mirror of
https://github.com/jxxghp/MoviePilot.git
synced 2026-08-29 03:56:43 +08:00
refactor: inject plugin interaction event publisher
This commit is contained in:
@@ -2,14 +2,21 @@ import uuid
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timedelta
|
||||
from threading import Lock
|
||||
from typing import Any, Dict, List, Optional, Tuple, Union
|
||||
from typing import Any, Dict, List, Optional, Protocol, Tuple, Union
|
||||
|
||||
from app.application.messaging.interaction import InteractionContext, MessageGateway
|
||||
from app.runtime.events import EventManager
|
||||
from app.schemas.message import Message
|
||||
from app.schemas.types import EventType, NotificationChannel
|
||||
|
||||
|
||||
class PluginEventPublisher(Protocol):
|
||||
"""声明插件输入交互所需的同步事件发布端口。"""
|
||||
|
||||
def send_event(self, event_type: EventType, payload: dict) -> Any:
|
||||
"""同步发布插件输入事件,并返回宿主事件分发结果。"""
|
||||
...
|
||||
|
||||
|
||||
@dataclass
|
||||
class PendingPluginInputInteraction:
|
||||
"""
|
||||
@@ -370,9 +377,14 @@ plugin_input_interaction_manager = PluginInputInteractionManager()
|
||||
class PluginInputInteractionHandler:
|
||||
"""消费插件申请接管的下一条用户文本输入。"""
|
||||
|
||||
def __init__(self, messenger: MessageGateway):
|
||||
"""保存消息投递接口。"""
|
||||
def __init__(
|
||||
self,
|
||||
messenger: MessageGateway,
|
||||
event_publisher: PluginEventPublisher,
|
||||
):
|
||||
"""保存消息投递接口和宿主注入的事件发布端口。"""
|
||||
self._messenger = messenger
|
||||
self._event_publisher = event_publisher
|
||||
|
||||
def handle_text(
|
||||
self,
|
||||
@@ -410,8 +422,7 @@ class PluginInputInteractionHandler:
|
||||
return False
|
||||
|
||||
if status == "expired":
|
||||
# 调用时解析单例,避免模块级绑定在单例注册表被重置后与宿主脱钩
|
||||
EventManager().send_event(
|
||||
self._event_publisher.send_event(
|
||||
EventType.MessageAction,
|
||||
{
|
||||
"plugin_id": request.plugin_id,
|
||||
@@ -442,7 +453,7 @@ class PluginInputInteractionHandler:
|
||||
return not text.strip().startswith("/")
|
||||
|
||||
if is_cancel_text:
|
||||
EventManager().send_event(
|
||||
self._event_publisher.send_event(
|
||||
EventType.MessageAction,
|
||||
{
|
||||
"plugin_id": request.plugin_id,
|
||||
@@ -472,7 +483,7 @@ class PluginInputInteractionHandler:
|
||||
)
|
||||
return True
|
||||
|
||||
EventManager().send_event(
|
||||
self._event_publisher.send_event(
|
||||
EventType.MessageAction,
|
||||
{
|
||||
"plugin_id": request.plugin_id,
|
||||
|
||||
@@ -253,7 +253,10 @@ class MessageChain(ChainBase):
|
||||
is_channel_admin=is_channel_admin,
|
||||
)
|
||||
|
||||
if PluginInputInteractionHandler(messenger=self).handle_text(
|
||||
if PluginInputInteractionHandler(
|
||||
messenger=self,
|
||||
event_publisher=self.eventmanager,
|
||||
).handle_text(
|
||||
context=interaction_context,
|
||||
text=text,
|
||||
reply_to_message_id=reply_to_message_id,
|
||||
@@ -415,7 +418,10 @@ class MessageChain(ChainBase):
|
||||
)
|
||||
return False
|
||||
|
||||
if PluginInputInteractionHandler(messenger=self).handle_text(
|
||||
if PluginInputInteractionHandler(
|
||||
messenger=self,
|
||||
event_publisher=self.eventmanager,
|
||||
).handle_text(
|
||||
context=context,
|
||||
text=text,
|
||||
reply_to_message_id=reply_to_message_id,
|
||||
|
||||
@@ -6,7 +6,7 @@
|
||||
> 审计范围:宿主后端;排除 `app/plugins/**` 运行时插件副本
|
||||
> 规范优先级:`AGENTS.md` 与 `docs/rules/` 高于本文
|
||||
> 相关文档:`docs/architecture-overview.md`、`docs/refactor/backend-architecture-governance.md`、`docs/refactor/backend-module-refactor-compatibility.md`
|
||||
> 实施进度:阶段 0~6 的宿主架构能力已完成收口;API/Application 公共复杂度基线已清零,启动组合根的 SystemConfigOper 构造点已由 14 降至 1;API 进程内后台任务已完成首批统一登记,插件仓适配和 Outbox 外围扩展仍按风险切片推进。Model/Base 查询与写装饰器、legacy 隐式会话外壳均已清零,插件 SDK 也不再导出宿主 Model。2026-08-23 的长期整改阶段 0 已恢复宿主、启动性能、官方插件和 SDK 契约门禁的可信基线;阶段 1a 已补齐 TaskRegistry owner 零债务门禁和诚实的关停超时语义;阶段 1b1 已收口整理 worker、pending 回放、失败通知、进程内 AI 重试、插件监控与事件投递的生命周期所有权;2026-08-24 的阶段 2 已将 212 个已观察宿主模块方法的 legacy aggregation 清零,并补齐可执行 fanout 与下载器文件 DTO 边界;阶段 3 已将消息交互和远程命令的订阅删除统一到 Application/UoW/outbox,宿主不再调用裸线程统计入口;阶段 4 已统一七种消息渠道的宿主回环与后台执行边界;阶段 5 已补齐事件窗口聚合任务的生命周期所有权;阶段 6 已统一插件文件操作的取消完成语义;阶段 7 已统一插件协程补偿的终态等待;阶段 8 已统一宿主同步函数的异步线程池入口;阶段 9 已统一工作流运行时的宿主获取路径;阶段 10 已统一模块、插件与调度运行时的显式 getter 调用;阶段 11 已清除系统配置 getter 的 Oper 形别名;阶段 12 已完成工作流域的显式 Chain 数据端口迁移;阶段 13 已收口用户、交互与消息链的数据端口;阶段 14 已收口音乐订阅数据端口;阶段 15 已收口站点数据端口;阶段 16 已收口媒体服务器数据端口;阶段 17 已收口下载数据端口;阶段 18 已收口主订阅数据端口;阶段 19 已收口整理数据端口;阶段 20 已收口 Agent 数据端口;阶段 21 已收口监控历史端口;阶段 22 已统一服务配置应用边界;阶段 23 已补齐媒体服务器 API 遗留的类形配置读取路径;阶段 24 已清除 Scheduler 内部无 owner 的协程提交双轨;阶段 25 已补齐 TaskRegistry 跨线程 owner 并迁移整理 AI 接管;阶段 26 已统一 Agent 会话清理提交;阶段 27 已统一历史 AI 进度 owner;阶段 28 已托管旧插件订阅统计线程;阶段 29 已统一 Emby 系条目转换并清零重复代码白名单;阶段 30 已收口插件市场请求级子任务;阶段 31 已托管搜索 AI 推荐任务;阶段 32 已清除事件调度器绕过生命周期 owner 的投递回退;阶段 33 已统一宿主 Agent 运行时的获取路径;阶段 34 已统一 durable-required 事件与 Outbox topic 事实源;阶段 35 已统一 LLM provider 管理 API 的运行时解析路径;阶段 36 已统一 WebAgent 音频能力访问边界。
|
||||
> 实施进度:阶段 0~6 的宿主架构能力已完成收口;API/Application 公共复杂度基线已清零,启动组合根的 SystemConfigOper 构造点已由 14 降至 1;API 进程内后台任务已完成首批统一登记,插件仓适配和 Outbox 外围扩展仍按风险切片推进。Model/Base 查询与写装饰器、legacy 隐式会话外壳均已清零,插件 SDK 也不再导出宿主 Model。2026-08-23 的长期整改阶段 0 已恢复宿主、启动性能、官方插件和 SDK 契约门禁的可信基线;阶段 1a 已补齐 TaskRegistry owner 零债务门禁和诚实的关停超时语义;阶段 1b1 已收口整理 worker、pending 回放、失败通知、进程内 AI 重试、插件监控与事件投递的生命周期所有权;2026-08-24 的阶段 2 已将 212 个已观察宿主模块方法的 legacy aggregation 清零,并补齐可执行 fanout 与下载器文件 DTO 边界;阶段 3 已将消息交互和远程命令的订阅删除统一到 Application/UoW/outbox,宿主不再调用裸线程统计入口;阶段 4 已统一七种消息渠道的宿主回环与后台执行边界;阶段 5 已补齐事件窗口聚合任务的生命周期所有权;阶段 6 已统一插件文件操作的取消完成语义;阶段 7 已统一插件协程补偿的终态等待;阶段 8 已统一宿主同步函数的异步线程池入口;阶段 9 已统一工作流运行时的宿主获取路径;阶段 10 已统一模块、插件与调度运行时的显式 getter 调用;阶段 11 已清除系统配置 getter 的 Oper 形别名;阶段 12 已完成工作流域的显式 Chain 数据端口迁移;阶段 13 已收口用户、交互与消息链的数据端口;阶段 14 已收口音乐订阅数据端口;阶段 15 已收口站点数据端口;阶段 16 已收口媒体服务器数据端口;阶段 17 已收口下载数据端口;阶段 18 已收口主订阅数据端口;阶段 19 已收口整理数据端口;阶段 20 已收口 Agent 数据端口;阶段 21 已收口监控历史端口;阶段 22 已统一服务配置应用边界;阶段 23 已补齐媒体服务器 API 遗留的类形配置读取路径;阶段 24 已清除 Scheduler 内部无 owner 的协程提交双轨;阶段 25 已补齐 TaskRegistry 跨线程 owner 并迁移整理 AI 接管;阶段 26 已统一 Agent 会话清理提交;阶段 27 已统一历史 AI 进度 owner;阶段 28 已托管旧插件订阅统计线程;阶段 29 已统一 Emby 系条目转换并清零重复代码白名单;阶段 30 已收口插件市场请求级子任务;阶段 31 已托管搜索 AI 推荐任务;阶段 32 已清除事件调度器绕过生命周期 owner 的投递回退;阶段 33 已统一宿主 Agent 运行时的获取路径;阶段 34 已统一 durable-required 事件与 Outbox topic 事实源;阶段 35 已统一 LLM provider 管理 API 的运行时解析路径;阶段 36 已统一 WebAgent 音频能力访问边界;阶段 37 已统一插件输入事件发布路径。
|
||||
|
||||
## 当前复核结论(2026-08-24)
|
||||
|
||||
@@ -388,13 +388,23 @@
|
||||
导入 `AgentCapabilityManager`,但保留 `app.agent.llm` 的公开惰性导出供 Agent 内部与既有消费者使用。
|
||||
- 宿主依赖图模块仍为 `806`,内部边由 `6549` 降至 `6547`;12 组禁止边与唯一隔离 TMDB SCC 均未变化。
|
||||
|
||||
### 长期整改阶段 37:插件输入事件发布边界收口(2026-08-24)
|
||||
|
||||
- `PluginInputInteractionHandler` 原先位于 Application 层,却在超时、取消和正常输入三条分支重新构造
|
||||
`EventManager`;两个生产调用方所在的 MessageChain 已持有同一事件管理器,形成注入端口与全局定位双轨。
|
||||
- handler 现在必须接收显式事件发布端口,MessageChain 统一注入自身 `eventmanager`。事件类型、
|
||||
payload、`__mp_target_plugin_id` 定向字段、消息提示及消费返回值均未变化;架构门禁把 Application 到
|
||||
`app.runtime.events` 的直接依赖锁为零,插件公开事件合同和 V1/V2/V3 加载路径保持不变。
|
||||
- 宿主依赖图模块仍为 `806`,内部边由 `6547` 降至 `6545`;事件 producer 合同仍为 `78`,12 组禁止边
|
||||
与唯一隔离 TMDB SCC 均未变化。
|
||||
|
||||
### 总体判断
|
||||
|
||||
当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**:
|
||||
|
||||
- 继续采用单进程控制面是正确选择,不建议现在拆成微服务;插件、调度器、工作流、事件和数据库共享进程内状态,拆分会放大部署、事务和兼容成本。
|
||||
- `foundation/domain/runtime/adapters/application/chain/api/startup` 的职责方向基本成立;宿主架构基线、复杂度 ratchet、异步阻塞 ratchet 当前均通过。
|
||||
- 依赖图当前为 `806` 个 Python 模块、`6549` 条内部导入边;唯一非平凡 SCC 位于隔离的 TMDB 第三方移植包内部,不应为了指标归零重写。
|
||||
- 依赖图当前为 `806` 个 Python 模块、`6545` 条内部导入边;唯一非平凡 SCC 位于隔离的 TMDB 第三方移植包内部,不应为了指标归零重写。
|
||||
- 当前主要风险已经从“目录和依赖失控”转移到运行时协议、后台副作用的可靠性和遗留兼容面。换言之,下一阶段重点应是**语义收口和可验证性**,而不是继续搬文件或机械拆大文件。
|
||||
|
||||
综合评价:架构方向可持续,生产可用性较高;可演进性仍处于中等水平。现阶段没有静态审计发现必须立即推倒重来的 P0 架构问题,但存在需要按 P1/P2 计划治理的真实债务。
|
||||
|
||||
+2
-4
@@ -13,8 +13,8 @@
|
||||
"runtime_to_db": [],
|
||||
"workflow_to_db": []
|
||||
},
|
||||
"edge_count": 6547,
|
||||
"edge_sha256": "66726be779aae9069c32e8ccba1e7371eb99cfcc459666bc9805ddc25061ac63",
|
||||
"edge_count": 6545,
|
||||
"edge_sha256": "e74daad54d04652df3946334297eac0dd9a4560ca87f84365b25dc9f553436da",
|
||||
"edges": [
|
||||
"app -> app.runtime",
|
||||
"app -> app.runtime.compat",
|
||||
@@ -2651,8 +2651,6 @@
|
||||
"app.application.messaging.plugin -> app.application",
|
||||
"app.application.messaging.plugin -> app.application.messaging",
|
||||
"app.application.messaging.plugin -> app.application.messaging.interaction",
|
||||
"app.application.messaging.plugin -> app.runtime",
|
||||
"app.application.messaging.plugin -> app.runtime.events",
|
||||
"app.application.messaging.plugin -> app.schemas",
|
||||
"app.application.messaging.plugin -> app.schemas.message",
|
||||
"app.application.messaging.plugin -> app.schemas.types",
|
||||
|
||||
@@ -1221,6 +1221,17 @@ def test_host_consumers_use_agent_audio_capability_application_port():
|
||||
assert violations == {}
|
||||
|
||||
|
||||
def test_application_services_do_not_resolve_event_manager_singleton():
|
||||
"""Application 服务必须接收事件端口,不得自行定位进程级事件单例。"""
|
||||
violations = {
|
||||
module_name: dependencies & {"app.runtime.events"}
|
||||
for module_name, dependencies in _build_module_graph().items()
|
||||
if module_name.startswith("app.application")
|
||||
and "app.runtime.events" in dependencies
|
||||
}
|
||||
assert violations == {}
|
||||
|
||||
|
||||
def test_agent_tools_do_not_import_entrypoint_internals():
|
||||
"""Agent 工具不得穿透导入 HTTP 端点、调度器与命令注册表内部实现。
|
||||
|
||||
|
||||
@@ -15,7 +15,7 @@ from app.application.messaging.plugin import (
|
||||
PluginInputInteractionHandler,
|
||||
plugin_input_interaction_manager,
|
||||
)
|
||||
from app.schemas import IncomingMessage, TransferDirectoryConf
|
||||
from app.schemas import IncomingMessage, TransferDirectoryConf # pylint: disable=no-name-in-module
|
||||
from app.schemas.types import EventType, MediaSource, MediaType, NotificationChannel
|
||||
|
||||
|
||||
@@ -531,7 +531,10 @@ def test_plugin_input_session_ignores_none_text_messages():
|
||||
)
|
||||
image = IncomingMessage.MessageImage(ref="https://example.invalid/image.jpg")
|
||||
|
||||
handled = PluginInputInteractionHandler(messenger=chain).handle_text(
|
||||
handled = PluginInputInteractionHandler(
|
||||
messenger=chain,
|
||||
event_publisher=chain.eventmanager,
|
||||
).handle_text(
|
||||
context=InteractionContext(
|
||||
channel=NotificationChannel.Telegram,
|
||||
source="telegram-test",
|
||||
|
||||
@@ -7,10 +7,14 @@ from types import SimpleNamespace
|
||||
from unittest.mock import AsyncMock, Mock, patch
|
||||
|
||||
|
||||
from app.agent import AgentManager, _MessageTask, _async_start_processing_status
|
||||
from app.agent import ( # pylint: disable=no-name-in-module
|
||||
AgentManager,
|
||||
_MessageTask,
|
||||
_async_start_processing_status,
|
||||
)
|
||||
from app.chain.message import MessageChain
|
||||
from app.command import Command, _finish_command_processing_status
|
||||
from app.modules.telegram import TelegramModule
|
||||
from app.modules.telegram import TelegramModule # pylint: disable=no-name-in-module
|
||||
from app.modules.telegram.telegram import Telegram
|
||||
from app.runtime.config import global_vars
|
||||
from app.schemas.types import NotificationChannel
|
||||
@@ -259,6 +263,7 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
|
||||
|
||||
def test_async_agent_leaves_processing_status_to_worker(self):
|
||||
chain = MessageChain.__new__(MessageChain)
|
||||
chain.eventmanager = Mock()
|
||||
chain.runtime_config = replace(
|
||||
chain.runtime_config,
|
||||
ai_agent_enable=True,
|
||||
@@ -373,6 +378,7 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
|
||||
|
||||
def test_callback_stops_typing_when_message_handler_returns(self):
|
||||
chain = MessageChain.__new__(MessageChain)
|
||||
chain.eventmanager = Mock()
|
||||
status = MessageChain._ProcessingStatus(
|
||||
channel=NotificationChannel.Telegram,
|
||||
source="telegram-test",
|
||||
|
||||
Reference in New Issue
Block a user