mirror of
https://github.com/jxxghp/MoviePilot.git
synced 2026-09-04 23:17:20 +08:00
refactor: unify web agent notice event routing
This commit is contained in:
+6
-150
@@ -45,7 +45,6 @@ from app.application.agent import (
|
||||
from app.chain.message import MessageChain
|
||||
from app.application.commands import get_command, get_commands
|
||||
from app.runtime.config import global_vars
|
||||
from app.runtime.events import Event, EventManager
|
||||
from app.api.principal import ApiPrincipal
|
||||
from app.api.dependencies.agent import (
|
||||
get_agent_chat_persistence,
|
||||
@@ -62,9 +61,12 @@ from app.application.messaging.chat import (
|
||||
from app.application.security.user import get_configured_user_id_lookup
|
||||
from app.application.configuration import get_api_runtime_config_snapshot
|
||||
from app.application.messaging.agent import (
|
||||
attach_web_agent_message_queue,
|
||||
attach_web_agent_edit_queue,
|
||||
create_web_agent_background_task,
|
||||
detach_web_agent_message_queue,
|
||||
detach_web_agent_edit_queue,
|
||||
is_web_agent_message_for_user,
|
||||
)
|
||||
from app.application.messaging.agent import agent_interaction_manager
|
||||
from app.application.messaging.agent import (
|
||||
@@ -94,9 +96,6 @@ WEB_AGENT_STREAM_COALESCE_MAX_CHARS = 256
|
||||
WEB_AGENT_STREAM_HEARTBEAT_SECONDS = 15.0
|
||||
WEB_AGENT_STREAM_QUEUE_MAX_SIZE = 64
|
||||
_WEB_AGENT_FILE_REGISTRY: dict[str, dict[str, Any]] = {}
|
||||
_WEB_AGENT_MESSAGE_QUEUES: dict[str, list[Queue[_SchemaMessage]]] = {}
|
||||
_WEB_AGENT_MESSAGE_LOCK = Lock()
|
||||
_WEB_AGENT_MESSAGE_LISTENER_REGISTERED = False
|
||||
|
||||
|
||||
class _WebAgentEventPublisher:
|
||||
@@ -1400,149 +1399,6 @@ def _has_web_agent_traditional_interaction(user_id: str) -> bool:
|
||||
return has_pending_interaction(user_id)
|
||||
|
||||
|
||||
def _extract_web_agent_message_from_event_data(
|
||||
data: dict,
|
||||
) -> Optional[_SchemaMessage]:
|
||||
"""
|
||||
从 NoticeMessage 事件数据中提取 WebAgent 通知。
|
||||
|
||||
:param data: NoticeMessage 事件数据,兼容扁平字段和 message 包装格式
|
||||
:return: WebAgent 通知,不属于 WebAgent 或数据无效时返回 None
|
||||
"""
|
||||
if not isinstance(data, dict):
|
||||
return None
|
||||
|
||||
try:
|
||||
message = data.get("message")
|
||||
if isinstance(message, _SchemaMessage):
|
||||
message = message
|
||||
elif isinstance(message, dict):
|
||||
message_data = copy.deepcopy(message)
|
||||
message_data.pop("type", None)
|
||||
message = _SchemaMessage(**message_data)
|
||||
else:
|
||||
message_data = copy.deepcopy(data)
|
||||
message_data.pop("type", None)
|
||||
message_data.pop("current_time", None)
|
||||
message = _SchemaMessage(**message_data)
|
||||
except Exception as err:
|
||||
logger.debug(f"解析WebAgent通知事件失败: {err}")
|
||||
return None
|
||||
|
||||
channel = message.channel
|
||||
channel_value = channel.value if isinstance(channel, NotificationChannel) else channel
|
||||
if channel_value != NotificationChannel.WebAgent.value:
|
||||
return None
|
||||
return message
|
||||
|
||||
|
||||
def _is_web_agent_message_for_user(
|
||||
message: _SchemaMessage,
|
||||
user_id: str,
|
||||
) -> bool:
|
||||
"""
|
||||
判断 NoticeMessage 事件是否属于当前 WebAgent 用户。
|
||||
|
||||
:param message: NoticeMessage 中的通知消息
|
||||
:param user_id: 当前登录用户 ID
|
||||
:return: 可被本次 WebAgent 请求消费时返回 True
|
||||
"""
|
||||
try:
|
||||
target_user = message.userid
|
||||
return target_user is None or str(target_user) == str(user_id)
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
|
||||
def _get_web_agent_message_user_id(message: _SchemaMessage) -> Optional[str]:
|
||||
"""
|
||||
从 NoticeMessage 事件中解析 WebAgent 目标用户。
|
||||
|
||||
:param message: NoticeMessage 中的通知消息
|
||||
:return: 用户 ID 字符串,事件不属于 WebAgent 时返回 None
|
||||
"""
|
||||
try:
|
||||
channel = message.channel
|
||||
channel_value = channel.value if isinstance(channel, NotificationChannel) else channel
|
||||
if channel_value != NotificationChannel.WebAgent.value:
|
||||
return None
|
||||
user_id = message.userid
|
||||
return str(user_id) if user_id is not None else None
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
def _dispatch_web_agent_message_event(event: Event) -> None:
|
||||
"""
|
||||
将 WebAgent NoticeMessage 分发给正在等待的请求队列。
|
||||
|
||||
:param event: NoticeMessage 广播事件
|
||||
"""
|
||||
data = event.event_data if isinstance(event.event_data, dict) else {}
|
||||
message = _extract_web_agent_message_from_event_data(data)
|
||||
if not message:
|
||||
return
|
||||
with _WEB_AGENT_MESSAGE_LOCK:
|
||||
user_id = _get_web_agent_message_user_id(message)
|
||||
if user_id is None:
|
||||
queues = [
|
||||
message_queue
|
||||
for user_queues in _WEB_AGENT_MESSAGE_QUEUES.values()
|
||||
for message_queue in user_queues
|
||||
]
|
||||
else:
|
||||
queues = list(_WEB_AGENT_MESSAGE_QUEUES.get(user_id) or [])
|
||||
for message_queue in queues:
|
||||
message_queue.put(message)
|
||||
|
||||
|
||||
def _ensure_web_agent_message_listener() -> None:
|
||||
"""
|
||||
确保 WebAgent NoticeMessage 全局监听器已注册。
|
||||
"""
|
||||
global _WEB_AGENT_MESSAGE_LISTENER_REGISTERED
|
||||
if _WEB_AGENT_MESSAGE_LISTENER_REGISTERED:
|
||||
return
|
||||
with _WEB_AGENT_MESSAGE_LOCK:
|
||||
if _WEB_AGENT_MESSAGE_LISTENER_REGISTERED:
|
||||
return
|
||||
EventManager().add_event_listener(
|
||||
EventType.NoticeMessage,
|
||||
_dispatch_web_agent_message_event,
|
||||
)
|
||||
_WEB_AGENT_MESSAGE_LISTENER_REGISTERED = True
|
||||
|
||||
|
||||
def _attach_web_agent_message_queue(user_id: str, message_queue: Queue[_SchemaMessage]) -> None:
|
||||
"""
|
||||
为当前 WebAgent 请求挂载通知收集队列。
|
||||
|
||||
:param user_id: 当前用户 ID
|
||||
:param message_queue: 用于接收通知事件的队列
|
||||
"""
|
||||
_ensure_web_agent_message_listener()
|
||||
with _WEB_AGENT_MESSAGE_LOCK:
|
||||
_WEB_AGENT_MESSAGE_QUEUES.setdefault(str(user_id), []).append(message_queue)
|
||||
|
||||
|
||||
def _detach_web_agent_message_queue(user_id: str, message_queue: Queue[_SchemaMessage]) -> None:
|
||||
"""
|
||||
移除当前 WebAgent 请求的通知收集队列。
|
||||
|
||||
:param user_id: 当前用户 ID
|
||||
:param message_queue: 需要移除的队列
|
||||
"""
|
||||
with _WEB_AGENT_MESSAGE_LOCK:
|
||||
queues = _WEB_AGENT_MESSAGE_QUEUES.get(str(user_id))
|
||||
if not queues:
|
||||
return
|
||||
_WEB_AGENT_MESSAGE_QUEUES[str(user_id)] = [
|
||||
item for item in queues if item is not message_queue
|
||||
]
|
||||
if not _WEB_AGENT_MESSAGE_QUEUES[str(user_id)]:
|
||||
_WEB_AGENT_MESSAGE_QUEUES.pop(str(user_id), None)
|
||||
|
||||
|
||||
def _build_web_agent_command_items() -> list[dict]:
|
||||
"""
|
||||
读取当前可用斜杠命令并转换为前端建议列表。
|
||||
@@ -1629,7 +1485,7 @@ async def _collect_web_agent_traditional_events(
|
||||
edit_queue: Queue[dict] = Queue()
|
||||
user_id = str(current_user.id)
|
||||
|
||||
_attach_web_agent_message_queue(user_id, message_queue)
|
||||
attach_web_agent_message_queue(user_id, message_queue)
|
||||
attach_web_agent_edit_queue(user_id, edit_queue)
|
||||
try:
|
||||
await run_in_threadpool(
|
||||
@@ -1668,13 +1524,13 @@ async def _collect_web_agent_traditional_events(
|
||||
break
|
||||
continue
|
||||
|
||||
if not _is_web_agent_message_for_user(message, user_id):
|
||||
if not is_web_agent_message_for_user(message, user_id):
|
||||
continue
|
||||
events.extend(await _build_web_agent_message_events_async(message))
|
||||
idle_deadline = time.monotonic() + WEB_AGENT_TRADITIONAL_IDLE_TIMEOUT_SECONDS
|
||||
return events
|
||||
finally:
|
||||
_detach_web_agent_message_queue(user_id, message_queue)
|
||||
detach_web_agent_message_queue(user_id, message_queue)
|
||||
detach_web_agent_edit_queue(user_id, edit_queue)
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user