初始化提交
This commit is contained in:
@@ -0,0 +1,11 @@
|
||||
# ///
|
||||
# __init__.py
|
||||
# 描述:服务模块导出
|
||||
# ///
|
||||
|
||||
from services.node_manager import node_manager
|
||||
from services.monitor_service import monitor_service
|
||||
from services.log_service import log_service
|
||||
from services.queue_service import queue_service
|
||||
|
||||
__all__ = ["node_manager", "monitor_service", "log_service", "queue_service"]
|
||||
@@ -0,0 +1,145 @@
|
||||
# ///
|
||||
# log_service.py
|
||||
# 描述:日志服务,支持Redis缓存和文件持久化
|
||||
# 作者:AI Generated
|
||||
# 创建日期:2026-04-05
|
||||
# ///
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
from datetime import datetime
|
||||
from typing import List, Dict, Any, Optional
|
||||
from collections import deque
|
||||
|
||||
from config import get_settings
|
||||
|
||||
|
||||
class LogService:
|
||||
def __init__(self, max_logs: int = 500, log_file: str = "logs/app.log"):
|
||||
self.max_logs = max_logs
|
||||
self.log_file = log_file
|
||||
self._logs: deque = deque(maxlen=max_logs)
|
||||
self._subscribers: List[asyncio.Queue] = []
|
||||
self._log_levels = ["INFO", "WARNING", "ERROR", "SUCCESS"]
|
||||
self._redis_client = None
|
||||
self._redis_available = False
|
||||
self._init_redis()
|
||||
|
||||
os.makedirs(os.path.dirname(log_file), exist_ok=True)
|
||||
self._load_from_file()
|
||||
|
||||
def _init_redis(self):
|
||||
try:
|
||||
import redis
|
||||
settings = get_settings()
|
||||
redis_config = settings.redis
|
||||
self._redis_client = redis.from_url(
|
||||
redis_config.url,
|
||||
db=redis_config.db,
|
||||
password=redis_config.password,
|
||||
decode_responses=True
|
||||
)
|
||||
self._redis_client.ping()
|
||||
self._redis_available = True
|
||||
except Exception:
|
||||
self._redis_available = False
|
||||
self._redis_client = None
|
||||
|
||||
def _load_from_file(self):
|
||||
if not os.path.exists(self.log_file):
|
||||
return
|
||||
try:
|
||||
with open(self.log_file, "r", encoding="utf-8") as f:
|
||||
for line in f:
|
||||
if line.strip():
|
||||
try:
|
||||
log_entry = json.loads(line.strip())
|
||||
if len(self._logs) < self.max_logs:
|
||||
self._logs.append(log_entry)
|
||||
except Exception:
|
||||
continue
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def _save_to_file(self, log_entry: Dict[str, Any]):
|
||||
try:
|
||||
with open(self.log_file, "a", encoding="utf-8") as f:
|
||||
f.write(json.dumps(log_entry, ensure_ascii=False) + "\n")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def _save_to_redis(self, log_entry: Dict[str, Any]):
|
||||
if not self._redis_available:
|
||||
return
|
||||
try:
|
||||
self._redis_client.lpush("wxauto:logs", json.dumps(log_entry, ensure_ascii=False))
|
||||
self._redis_client.ltrim("wxauto:logs", 0, self.max_logs - 1)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def add_log(self, message: str, level: str = "INFO", source: str = "System"):
|
||||
log_entry = {
|
||||
"timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
|
||||
"level": level,
|
||||
"source": source,
|
||||
"message": message
|
||||
}
|
||||
|
||||
self._logs.append(log_entry)
|
||||
|
||||
if self._redis_available:
|
||||
self._save_to_redis(log_entry)
|
||||
else:
|
||||
self._save_to_file(log_entry)
|
||||
|
||||
self._notify_subscribers(log_entry)
|
||||
|
||||
def info(self, message: str, source: str = "System"):
|
||||
self.add_log(message, "INFO", source)
|
||||
|
||||
def warning(self, message: str, source: str = "System"):
|
||||
self.add_log(message, "WARNING", source)
|
||||
|
||||
def error(self, message: str, source: str = "System"):
|
||||
self.add_log(message, "ERROR", source)
|
||||
|
||||
def success(self, message: str, source: str = "System"):
|
||||
self.add_log(message, "SUCCESS", source)
|
||||
|
||||
async def subscribe(self) -> asyncio.Queue:
|
||||
queue = asyncio.Queue(maxsize=100)
|
||||
self._subscribers.append(queue)
|
||||
return queue
|
||||
|
||||
def unsubscribe(self, queue: asyncio.Queue):
|
||||
if queue in self._subscribers:
|
||||
self._subscribers.remove(queue)
|
||||
|
||||
async def _notify_subscribers(self, log_entry: Dict[str, Any]):
|
||||
for queue in self._subscribers:
|
||||
try:
|
||||
await queue.put(log_entry)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def get_recent_logs(self, count: int = 50) -> List[Dict[str, Any]]:
|
||||
if self._redis_available:
|
||||
try:
|
||||
logs = self._redis_client.lrange("wxauto:logs", 0, count - 1)
|
||||
return [json.loads(log) for log in reversed(logs)]
|
||||
except Exception:
|
||||
pass
|
||||
return list(self._logs)[-count:]
|
||||
|
||||
def get_all_logs(self) -> List[Dict[str, Any]]:
|
||||
if self._redis_available:
|
||||
try:
|
||||
logs = self._redis_client.lrange("wxauto:logs", 0, self.max_logs - 1)
|
||||
return [json.loads(log) for log in reversed(logs)]
|
||||
except Exception:
|
||||
pass
|
||||
return list(self._logs)
|
||||
|
||||
|
||||
log_service = LogService()
|
||||
@@ -0,0 +1,179 @@
|
||||
# ///
|
||||
# monitor_service.py
|
||||
# 描述:节点状态监控服务,定时检测API和微信状态
|
||||
# 作者:AI Generated
|
||||
# 创建日期:2026-04-05
|
||||
# ///
|
||||
|
||||
import asyncio
|
||||
from datetime import datetime
|
||||
from typing import Dict, Any, Optional
|
||||
import httpx
|
||||
|
||||
from config import get_settings, WebhookConfig
|
||||
from services.node_manager import node_manager
|
||||
from services.log_service import log_service
|
||||
from database_service import get_db_service
|
||||
|
||||
|
||||
class MonitorService:
|
||||
def __init__(self):
|
||||
self._running = False
|
||||
self._task: Optional[asyncio.Task] = None
|
||||
self._last_status: Dict[str, Dict[str, str]] = {}
|
||||
|
||||
async def start(self):
|
||||
if self._running:
|
||||
return
|
||||
self._running = True
|
||||
self._task = asyncio.create_task(self._monitor_loop())
|
||||
|
||||
async def stop(self):
|
||||
self._running = False
|
||||
if self._task:
|
||||
self._task.cancel()
|
||||
try:
|
||||
await self._task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
|
||||
async def _monitor_loop(self):
|
||||
settings = get_settings()
|
||||
log_service.info("监控服务启动", "Monitor")
|
||||
while self._running:
|
||||
try:
|
||||
await self._check_all_nodes()
|
||||
except Exception as e:
|
||||
log_service.error(f"监控循环错误: {e}", "Monitor")
|
||||
await asyncio.sleep(settings.monitor.check_interval)
|
||||
|
||||
async def _check_all_nodes(self):
|
||||
settings = get_settings()
|
||||
if not settings.monitor.enabled:
|
||||
return
|
||||
|
||||
for node in await node_manager.get_all_nodes():
|
||||
await self._check_node(node.node_id)
|
||||
|
||||
async def _check_node(self, node_id: str):
|
||||
node = await node_manager.get_node(node_id)
|
||||
if not node:
|
||||
return
|
||||
|
||||
api_status_changed = False
|
||||
wechat_status_changed = False
|
||||
old_api_status = self._last_status.get(node_id, {}).get("api_status")
|
||||
old_wechat_status = self._last_status.get(node_id, {}).get("wechat_status")
|
||||
|
||||
api_status = await node_manager.check_node_status(node_id)
|
||||
if old_api_status is not None and old_api_status != api_status:
|
||||
api_status_changed = True
|
||||
self._last_status.setdefault(node_id, {})["api_status"] = api_status
|
||||
|
||||
wechat_status = node.wechat_status
|
||||
if old_wechat_status is not None and old_wechat_status != wechat_status:
|
||||
wechat_status_changed = True
|
||||
self._last_status.setdefault(node_id, {})["wechat_status"] = wechat_status
|
||||
|
||||
if api_status_changed:
|
||||
event = "node_online" if self._last_status[node_id]["api_status"] == "active" else "node_offline"
|
||||
msg = f"节点 {node.name} ({node_id}) API状态变为 {self._last_status[node_id]['api_status']}"
|
||||
log_service.info(msg, "Monitor")
|
||||
if event == "node_offline":
|
||||
log_service.error(msg, "Monitor")
|
||||
await self._send_webhook(event, {
|
||||
"node_id": node_id,
|
||||
"node_name": node.name,
|
||||
"api_status": self._last_status[node_id]["api_status"],
|
||||
"message": msg
|
||||
})
|
||||
|
||||
if wechat_status_changed:
|
||||
event = "wechat_online" if wechat_status == "online" else "wechat_offline"
|
||||
msg = f"节点 {node.name} ({node_id}) 微信状态变为 {wechat_status}"
|
||||
log_service.info(msg, "Monitor")
|
||||
if event == "wechat_offline":
|
||||
log_service.warning(msg, "Monitor")
|
||||
await self._send_webhook(event, {
|
||||
"node_id": node_id,
|
||||
"node_name": node.name,
|
||||
"wechat_status": wechat_status,
|
||||
"message": msg
|
||||
})
|
||||
|
||||
async def _send_webhook(self, event: str, data: Dict[str, Any]):
|
||||
settings = get_settings()
|
||||
webhook_config: WebhookConfig = settings.webhook
|
||||
|
||||
db_service = get_db_service()
|
||||
webhooks = []
|
||||
if db_service:
|
||||
webhooks = db_service.get_enabled_webhooks()
|
||||
|
||||
if not webhooks:
|
||||
if not webhook_config.enabled or not webhook_config.urls:
|
||||
return
|
||||
webhooks = [{"url": url, "format": webhook_config.format} for url in webhook_config.urls]
|
||||
|
||||
event_names = {
|
||||
"api_error": "API错误",
|
||||
"wechat_offline": "微信掉线",
|
||||
"wechat_online": "微信上线",
|
||||
"node_offline": "节点离线",
|
||||
"node_online": "节点在线"
|
||||
}
|
||||
|
||||
event_name = event_names.get(event, event)
|
||||
node_name = data.get('node_name', 'N/A')
|
||||
status_info = data.get('wechat_status') or data.get('api_status') or 'N/A'
|
||||
timestamp = datetime.now().strftime('%Y-%m-%d %H:%M:%S')
|
||||
|
||||
for webhook in webhooks:
|
||||
url = webhook.url if hasattr(webhook, 'url') else webhook.get('url')
|
||||
format_type = webhook.format if hasattr(webhook, 'format') else webhook.get('format', 'bark')
|
||||
event_types = webhook.event_types if hasattr(webhook, 'event_types') else []
|
||||
|
||||
if event_types and event not in event_types:
|
||||
continue
|
||||
|
||||
if format_type == "bark":
|
||||
title = f"WXAuto告警 - {event_name}"
|
||||
content = f"{node_name}\n状态: {status_info}\n时间: {timestamp}"
|
||||
payload = {
|
||||
"title": title,
|
||||
"body": content,
|
||||
"icon": "https://is3-ssl.mzstatic.com/image/thumb/Purple122/v4/8e/1a/b5/8e1ab5a2-09c7-2c06-ae03-2056b0348d2c/[email protected]/0x0ss.png"
|
||||
}
|
||||
else:
|
||||
content = f"""【WXAuto Center 告警】
|
||||
事件类型: {event_name}
|
||||
节点名称: {node_name}
|
||||
状态信息: {status_info}
|
||||
时间: {timestamp}"""
|
||||
payload = {
|
||||
"msgtype": "text",
|
||||
"text": {
|
||||
"content": content
|
||||
}
|
||||
}
|
||||
|
||||
for attempt in range(webhook_config.retry_times):
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=webhook_config.timeout) as client:
|
||||
response = await client.post(url, json=payload)
|
||||
if response.status_code == 200:
|
||||
log_service.success(f"Webhook通知发送成功: {event} -> {url}", "Monitor")
|
||||
break
|
||||
except Exception as e:
|
||||
log_service.warning(f"Webhook发送失败 (尝试 {attempt + 1}): {e}", "Monitor")
|
||||
else:
|
||||
log_service.error(f"Webhook通知发送失败 (已重试{webhook_config.retry_times}次): {event} -> {url}", "Monitor")
|
||||
|
||||
def get_last_status(self, node_id: str) -> Dict[str, str]:
|
||||
return self._last_status.get(node_id, {})
|
||||
|
||||
def get_all_status(self) -> Dict[str, Dict[str, str]]:
|
||||
return self._last_status.copy()
|
||||
|
||||
|
||||
monitor_service = MonitorService()
|
||||
@@ -0,0 +1,225 @@
|
||||
# ///
|
||||
# node_manager.py
|
||||
# 描述:节点管理服务,负责与WXAUTO-HTTP-API节点通信
|
||||
# 作者:AI Generated
|
||||
# 创建日期:2026-04-05
|
||||
# ///
|
||||
|
||||
import asyncio
|
||||
from typing import Dict, List, Optional, Any
|
||||
import httpx
|
||||
|
||||
from services.log_service import log_service
|
||||
from database_service import get_db_service
|
||||
|
||||
|
||||
class NodeStatus:
|
||||
ACTIVE = "active"
|
||||
INACTIVE = "inactive"
|
||||
ERROR = "error"
|
||||
|
||||
|
||||
class Node:
|
||||
def __init__(self, node_id: str, name: str, api_url: str, api_key: str, description: str = "", group: str = "default"):
|
||||
self.node_id = node_id
|
||||
self.name = name
|
||||
self.api_url = api_url.rstrip("/")
|
||||
self.api_key = api_key
|
||||
self.description = description
|
||||
self.group = group
|
||||
self.enabled = True
|
||||
self.status = NodeStatus.INACTIVE
|
||||
self.wechat_status = "online"
|
||||
self.is_healthy = False
|
||||
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
return {
|
||||
"node_id": self.node_id,
|
||||
"name": self.name,
|
||||
"api_url": self.api_url,
|
||||
"enabled": self.enabled,
|
||||
"description": self.description,
|
||||
"group": self.group,
|
||||
"status": self.status,
|
||||
"wechat_status": self.wechat_status,
|
||||
"is_healthy": self.is_healthy
|
||||
}
|
||||
|
||||
|
||||
class NodeManager:
|
||||
def __init__(self):
|
||||
self.nodes: Dict[str, Node] = {}
|
||||
|
||||
async def add_node(self, node_id: str, name: str, api_url: str, api_key: str, description: str = "", group: str = "default") -> bool:
|
||||
if node_id in self.nodes:
|
||||
return False
|
||||
self.nodes[node_id] = Node(node_id, name, api_url, api_key, description, group)
|
||||
await self.check_node_status(node_id)
|
||||
return True
|
||||
|
||||
async def remove_node(self, node_id: str) -> bool:
|
||||
if node_id not in self.nodes:
|
||||
return False
|
||||
del self.nodes[node_id]
|
||||
return True
|
||||
|
||||
async def get_node(self, node_id: str) -> Optional[Node]:
|
||||
return self.nodes.get(node_id)
|
||||
|
||||
async def get_all_nodes(self) -> List[Node]:
|
||||
return list(self.nodes.values())
|
||||
|
||||
async def get_enabled_nodes(self) -> List[Node]:
|
||||
return [n for n in self.nodes.values() if n.enabled]
|
||||
|
||||
async def get_nodes_by_group(self, group: str) -> List[Node]:
|
||||
return [n for n in self.nodes.values() if n.group == group]
|
||||
|
||||
def to_config(self) -> List[Dict[str, Any]]:
|
||||
return [{
|
||||
"node_id": n.node_id,
|
||||
"name": n.name,
|
||||
"api_url": n.api_url,
|
||||
"api_key": n.api_key,
|
||||
"enabled": n.enabled,
|
||||
"description": n.description,
|
||||
"group": n.group
|
||||
} for n in self.nodes.values()]
|
||||
|
||||
async def update_node(self, node_id: str, **kwargs) -> bool:
|
||||
node = self.nodes.get(node_id)
|
||||
if not node:
|
||||
return False
|
||||
for key, value in kwargs.items():
|
||||
if hasattr(node, key):
|
||||
setattr(node, key, value)
|
||||
if "api_url" in kwargs or "api_key" in kwargs:
|
||||
await self.check_node_status(node_id)
|
||||
return True
|
||||
|
||||
async def check_node_status(self, node_id: str) -> str:
|
||||
node = self.nodes.get(node_id)
|
||||
if not node or not node.enabled:
|
||||
return NodeStatus.INACTIVE
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=5.0) as client:
|
||||
response = await client.get(
|
||||
f"{node.api_url}/health",
|
||||
headers={"X-API-Key": node.api_key}
|
||||
)
|
||||
if response.status_code == 200:
|
||||
node.status = NodeStatus.ACTIVE
|
||||
else:
|
||||
node.status = NodeStatus.ERROR
|
||||
except Exception:
|
||||
node.status = NodeStatus.ERROR
|
||||
return node.status
|
||||
|
||||
async def is_node_healthy(self, node_id: str) -> tuple[bool, str]:
|
||||
node = self.nodes.get(node_id)
|
||||
if not node or not node.enabled:
|
||||
return False, "Node not found or disabled"
|
||||
|
||||
api_status = await self.check_node_status(node_id)
|
||||
if api_status != NodeStatus.ACTIVE:
|
||||
return False, f"API status: {api_status}"
|
||||
|
||||
return True, "OK"
|
||||
|
||||
async def call_node_api(self, node_id: str, endpoint: str, method: str = "GET", data: Optional[Dict] = None) -> Dict[str, Any]:
|
||||
node = self.nodes.get(node_id)
|
||||
if not node or not node.enabled:
|
||||
return {"success": False, "error": f"Node '{node_id}' not found or disabled"}
|
||||
|
||||
url = f"{node.api_url}{endpoint}"
|
||||
headers = {"X-API-Key": node.api_key}
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=10.0) as client:
|
||||
if method.upper() == "GET":
|
||||
response = await client.get(url, headers=headers)
|
||||
elif method.upper() == "POST":
|
||||
response = await client.post(url, headers=headers, json=data)
|
||||
else:
|
||||
return {"success": False, "error": f"Unsupported method: {method}"}
|
||||
|
||||
if response.status_code == 200:
|
||||
return {"success": True, "data": response.json()}
|
||||
else:
|
||||
return {"success": False, "error": f"HTTP {response.status_code}", "url": url}
|
||||
except Exception as e:
|
||||
return {"success": False, "error": str(e), "url": url}
|
||||
|
||||
async def send_message(self, node_id: str, who: str, msg: str, msg_type: str = "text") -> Dict[str, Any]:
|
||||
node = self.nodes.get(node_id)
|
||||
await self.call_node_api(node_id, "/api/wechat/initialize", method="POST")
|
||||
result = await self.call_node_api(
|
||||
node_id,
|
||||
"/api/message/send",
|
||||
method="POST",
|
||||
data={"receiver": who, "message": msg}
|
||||
)
|
||||
if node:
|
||||
node.wechat_status = "online" if result.get("success") else "offline"
|
||||
db_service = get_db_service()
|
||||
if db_service:
|
||||
db_service.update_node(node_id, wechat_status=node.wechat_status)
|
||||
return result
|
||||
|
||||
async def send_message_to_group(self, node_id: str, group_name: str, msg: str, msg_type: str = "text") -> Dict[str, Any]:
|
||||
node = self.nodes.get(node_id)
|
||||
await self.call_node_api(node_id, "/api/wechat/initialize", method="POST")
|
||||
result = await self.call_node_api(
|
||||
node_id,
|
||||
"/api/message/send",
|
||||
method="POST",
|
||||
data={"receiver": group_name, "message": msg}
|
||||
)
|
||||
if node:
|
||||
node.wechat_status = "online" if result.get("success") else "offline"
|
||||
db_service = get_db_service()
|
||||
if db_service:
|
||||
db_service.update_node(node_id, wechat_status=node.wechat_status)
|
||||
return result
|
||||
|
||||
async def get_friends_list(self, node_id: str) -> Dict[str, Any]:
|
||||
return await self.call_node_api(node_id, "/api/friend/list", method="GET")
|
||||
|
||||
async def get_groups_list(self, node_id: str) -> Dict[str, Any]:
|
||||
return await self.call_node_api(node_id, "/api/group/list", method="GET")
|
||||
|
||||
async def get_messages(self, node_id: str, chat_name: str, limit: int = 20) -> Dict[str, Any]:
|
||||
return await self.call_node_api(
|
||||
node_id,
|
||||
"/api/message/get-history",
|
||||
method="POST",
|
||||
data={"chat_name": chat_name, "save_pic": False, "save_video": False, "save_file": False, "save_voice": False}
|
||||
)
|
||||
|
||||
async def reset_wechat_status(self, node_id: str) -> Dict[str, Any]:
|
||||
node = self.nodes.get(node_id)
|
||||
if not node:
|
||||
return {"success": False, "error": f"Node '{node_id}' not found"}
|
||||
|
||||
await self.call_node_api(node_id, "/api/wechat/initialize", method="POST")
|
||||
result = await self.call_node_api(
|
||||
node_id,
|
||||
"/api/message/send",
|
||||
method="POST",
|
||||
data={"receiver": "文件传输助手", "message": "WXAuto Center 状态检测消息"}
|
||||
)
|
||||
|
||||
if result.get("success"):
|
||||
node.wechat_status = "online"
|
||||
else:
|
||||
node.wechat_status = "offline"
|
||||
|
||||
db_service = get_db_service()
|
||||
if db_service:
|
||||
db_service.update_node(node_id, wechat_status=node.wechat_status)
|
||||
|
||||
return result
|
||||
|
||||
|
||||
node_manager = NodeManager()
|
||||
@@ -0,0 +1,148 @@
|
||||
# ///
|
||||
# plugin_notification_service.py
|
||||
# 描述:插件通知服务,支持向配置的地址发送插件告警通知
|
||||
# 作者:User
|
||||
# 创建日期:2026-04-07
|
||||
# ///
|
||||
|
||||
import logging
|
||||
import httpx
|
||||
from typing import Any, Dict, List, Optional
|
||||
from datetime import datetime
|
||||
|
||||
logger = logging.getLogger("plugin_notification")
|
||||
|
||||
|
||||
class PluginNotificationService:
|
||||
def __init__(self):
|
||||
self._notification_url: Optional[str] = None
|
||||
self._enabled: bool = False
|
||||
self._webhook_configs: List[Dict[str, Any]] = []
|
||||
self._plugin_webhook_map: Dict[str, str] = {}
|
||||
|
||||
def initialize(self):
|
||||
from database_service import get_db_service
|
||||
db_service = get_db_service()
|
||||
if not db_service:
|
||||
logger.warning("Database service not available, plugin notification disabled")
|
||||
return
|
||||
|
||||
try:
|
||||
self._load_notification_config(db_service)
|
||||
self._load_webhook_configs(db_service)
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to initialize plugin notification service: {e}")
|
||||
|
||||
def _load_notification_config(self, db_service):
|
||||
try:
|
||||
from models import WebhookUrl
|
||||
webhooks = db_service.get_all_webhooks()
|
||||
for wh in webhooks:
|
||||
event_types = getattr(wh, 'event_types', '') or ''
|
||||
if 'plugin_notification' in event_types or event_types == '*':
|
||||
self._notification_url = wh.url
|
||||
self._enabled = wh.enabled
|
||||
logger.info(f"Plugin notification enabled: {wh.name} - {wh.url}")
|
||||
break
|
||||
except Exception as e:
|
||||
logger.warning(f"Failed to load notification config: {e}")
|
||||
|
||||
def _load_webhook_configs(self, db_service):
|
||||
try:
|
||||
from models import WebhookUrl
|
||||
webhooks = db_service.get_all_webhooks()
|
||||
self._webhook_configs = []
|
||||
for wh in webhooks:
|
||||
event_types = getattr(wh, 'event_types', '') or ''
|
||||
event_types_list = event_types.split(',') if isinstance(event_types, str) else []
|
||||
self._webhook_configs.append({
|
||||
'id': wh.id,
|
||||
'name': wh.name,
|
||||
'url': wh.url,
|
||||
'format': getattr(wh, 'format', 'bark'),
|
||||
'enabled': wh.enabled,
|
||||
'event_types': event_types_list
|
||||
})
|
||||
if 'plugin_notification' in event_types_list or '*' in event_types_list:
|
||||
self._notification_url = wh.url
|
||||
self._enabled = wh.enabled
|
||||
except Exception as e:
|
||||
logger.warning(f"Failed to load webhook configs: {e}")
|
||||
|
||||
def reload_config(self):
|
||||
self.initialize()
|
||||
|
||||
def send_notification(
|
||||
self,
|
||||
plugin_name: str,
|
||||
title: str,
|
||||
message: str,
|
||||
level: str = "info"
|
||||
) -> bool:
|
||||
logger.info(f"send_notification called: {plugin_name} - {title}, enabled={self._enabled}, url={self._notification_url}")
|
||||
|
||||
if not self._enabled or not self._notification_url:
|
||||
logger.warning(f"Plugin notification disabled or no URL: enabled={self._enabled}, url={self._notification_url}")
|
||||
logger.info(f"Available webhook configs: {self._webhook_configs}")
|
||||
return False
|
||||
|
||||
timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||||
|
||||
content = f"""【{plugin_name} 插件通知】
|
||||
级别: {level.upper()}
|
||||
标题: {title}
|
||||
消息: {message}
|
||||
时间: {timestamp}"""
|
||||
|
||||
payload = {
|
||||
"title": f"[{level.upper()}] {plugin_name}",
|
||||
"body": content,
|
||||
"icon": "https://is3-ssl.mzstatic.com/image/thumb/Purple122/v4/8e/1a/b5/8e1ab5a2-09c7-2c06-ae03-2056b0348d2c/[email protected]/0x0ss.png"
|
||||
}
|
||||
|
||||
for config in self._webhook_configs:
|
||||
if not config['enabled']:
|
||||
continue
|
||||
event_types = config['event_types']
|
||||
if 'plugin_notification' not in event_types and '*' not in event_types:
|
||||
continue
|
||||
|
||||
try:
|
||||
url = config['url']
|
||||
format_type = config.get('format', 'bark')
|
||||
|
||||
if format_type == 'wechat':
|
||||
payload = {
|
||||
"msgtype": "text",
|
||||
"text": {
|
||||
"content": content
|
||||
}
|
||||
}
|
||||
|
||||
result = self._send_webhook(url, payload)
|
||||
if result:
|
||||
logger.info(f"Plugin notification sent: {plugin_name} - {level} -> {config['name']}")
|
||||
return result
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to send plugin notification: {e}")
|
||||
|
||||
return False
|
||||
|
||||
def _send_webhook(self, url: str, payload: Dict[str, Any]) -> bool:
|
||||
try:
|
||||
with httpx.Client(timeout=10) as client:
|
||||
response = client.post(url, json=payload)
|
||||
return response.status_code == 200
|
||||
except Exception as e:
|
||||
logger.error(f"Webhook request failed: {e}")
|
||||
return False
|
||||
|
||||
def is_enabled(self) -> bool:
|
||||
return self._enabled
|
||||
|
||||
def get_notification_url(self) -> Optional[str]:
|
||||
return self._notification_url
|
||||
|
||||
|
||||
plugin_notification_service = PluginNotificationService()
|
||||
@@ -0,0 +1,119 @@
|
||||
# ///
|
||||
# queue_service.py
|
||||
# 描述:消息队列服务,支持Redis后端
|
||||
# 作者:AI Generated
|
||||
# 创建日期:2026-04-05
|
||||
# ///
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
from typing import Dict, Any, Optional, Callable
|
||||
from datetime import datetime
|
||||
from enum import Enum
|
||||
|
||||
from config import get_settings
|
||||
|
||||
|
||||
class QueueName(Enum):
|
||||
MESSAGE_SEND = "wxauto:queue:message_send"
|
||||
MESSAGE_BROADCAST = "wxauto:queue:message_broadcast"
|
||||
|
||||
|
||||
class QueueService:
|
||||
def __init__(self):
|
||||
self._redis_client = None
|
||||
self._redis_available = False
|
||||
self._local_queue = asyncio.Queue()
|
||||
self._init_redis()
|
||||
self._worker_task = None
|
||||
self._running = False
|
||||
|
||||
def _init_redis(self):
|
||||
try:
|
||||
import redis
|
||||
settings = get_settings()
|
||||
redis_config = settings.redis
|
||||
self._redis_client = redis.from_url(
|
||||
redis_config.url,
|
||||
db=redis_config.db,
|
||||
password=redis_config.password,
|
||||
decode_responses=True
|
||||
)
|
||||
self._redis_client.ping()
|
||||
self._redis_available = True
|
||||
except Exception:
|
||||
self._redis_available = False
|
||||
self._redis_client = None
|
||||
|
||||
def enqueue(self, queue_name: QueueName, data: Dict[str, Any]) -> bool:
|
||||
job_data = {
|
||||
"data": data,
|
||||
"enqueued_at": datetime.now().isoformat(),
|
||||
"status": "pending"
|
||||
}
|
||||
|
||||
if self._redis_available:
|
||||
try:
|
||||
self._redis_client.rpush(queue_name.value, json.dumps(job_data, ensure_ascii=False))
|
||||
return True
|
||||
except Exception:
|
||||
return False
|
||||
else:
|
||||
try:
|
||||
asyncio.create_task(self._local_enqueue(queue_name, job_data))
|
||||
return True
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
async def _local_enqueue(self, queue_name: QueueName, job_data: Dict[str, Any]):
|
||||
await self._local_queue.put((queue_name, job_data))
|
||||
|
||||
async def dequeue(self, queue_name: QueueName, timeout: int = 1) -> Optional[Dict[str, Any]]:
|
||||
if self._redis_available:
|
||||
try:
|
||||
result = self._redis_client.blpop(queue_name.value, timeout=timeout)
|
||||
if result:
|
||||
_, job_data = result
|
||||
return json.loads(job_data)
|
||||
except Exception:
|
||||
pass
|
||||
else:
|
||||
try:
|
||||
queue_item = await asyncio.wait_for(self._local_queue.get(), timeout=timeout)
|
||||
if queue_item and queue_item[0] == queue_name:
|
||||
return queue_item[1]
|
||||
except asyncio.TimeoutError:
|
||||
pass
|
||||
return None
|
||||
|
||||
def get_queue_length(self, queue_name: QueueName) -> int:
|
||||
if self._redis_available:
|
||||
try:
|
||||
return self._redis_client.llen(queue_name.value)
|
||||
except Exception:
|
||||
pass
|
||||
return 0
|
||||
|
||||
def start_worker(self, processor: Callable):
|
||||
if self._running:
|
||||
return
|
||||
self._running = True
|
||||
self._worker_task = asyncio.create_task(self._worker_loop(processor))
|
||||
|
||||
async def _worker_loop(self, processor: Callable):
|
||||
while self._running:
|
||||
for queue_name in [QueueName.MESSAGE_SEND]:
|
||||
job = await self.dequeue(queue_name, timeout=1)
|
||||
if job:
|
||||
try:
|
||||
await processor(job)
|
||||
except Exception as e:
|
||||
pass
|
||||
|
||||
def stop_worker(self):
|
||||
self._running = False
|
||||
if self._worker_task:
|
||||
self._worker_task.cancel()
|
||||
|
||||
|
||||
queue_service = QueueService()
|
||||
@@ -0,0 +1,66 @@
|
||||
# ///
|
||||
# scheduler_service.py
|
||||
# 描述:定时任务调度服务,支持 cron 表达式
|
||||
# 作者:AI Generated
|
||||
# 创建日期:2026-04-06
|
||||
# ///
|
||||
|
||||
import asyncio
|
||||
from typing import Dict, Any
|
||||
from croniter import croniter
|
||||
from plugins.base import PluginRegistry
|
||||
from services.log_service import log_service
|
||||
|
||||
|
||||
class SchedulerService:
|
||||
def __init__(self):
|
||||
self._running = False
|
||||
self._task = None
|
||||
self._last_run_times: Dict[str, float] = {}
|
||||
|
||||
async def start(self):
|
||||
if self._running:
|
||||
return
|
||||
self._running = True
|
||||
self._task = asyncio.create_task(self._scheduler_loop())
|
||||
log_service.info("调度服务启动", "Scheduler")
|
||||
|
||||
async def stop(self):
|
||||
self._running = False
|
||||
if self._task:
|
||||
self._task.cancel()
|
||||
try:
|
||||
await self._task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
log_service.info("调度服务已停止", "Scheduler")
|
||||
|
||||
async def _scheduler_loop(self):
|
||||
while self._running:
|
||||
try:
|
||||
await self._check_and_run_tasks()
|
||||
except Exception as e:
|
||||
log_service.error(f"调度循环错误: {e}", "Scheduler")
|
||||
await asyncio.sleep(10)
|
||||
|
||||
async def _check_and_run_tasks(self):
|
||||
tasks = PluginRegistry.get_scheduled_tasks()
|
||||
import time
|
||||
current_time = time.time()
|
||||
|
||||
for task in tasks:
|
||||
try:
|
||||
should_run, next_run = task.should_run_and_get_next()
|
||||
if should_run:
|
||||
log_service.info(f"执行定时任务: {task.plugin_name}", "Scheduler")
|
||||
result = task.run_task()
|
||||
if result.get("success"):
|
||||
log_service.success(f"定时任务完成: {task.plugin_name}", "Scheduler")
|
||||
else:
|
||||
log_service.error(f"定时任务失败: {task.plugin_name} - {result.get('error')}", "Scheduler")
|
||||
self._last_run_times[task.plugin_name] = current_time
|
||||
except Exception as e:
|
||||
log_service.error(f"执行定时任务 {task.plugin_name} 出错: {e}", "Scheduler")
|
||||
|
||||
|
||||
scheduler_service = SchedulerService()
|
||||
Reference in New Issue
Block a user