diff --git a/app/chain/workflow.py b/app/chain/workflow.py index 705d1c196..0b4d0d6b1 100644 --- a/app/chain/workflow.py +++ b/app/chain/workflow.py @@ -18,7 +18,7 @@ from app.chain import ChainBase from app.runtime.config import global_vars from app.runtime.events import Event, eventmanager from app.application.workflow import get_workflow_manager -from app.application.chain.data import WorkflowPortProxy as WorkflowOper +from app.application.chain.data import get_chain_workflow_port from app.runtime.log import logger from app.schemas.workflow import ActionContext from app.schemas.workflow import ActionFlow @@ -1185,13 +1185,13 @@ class WorkflowChain(ChainBase): :param from_begin: 是否从头开始,默认为True :param progress_callback: 定时服务进度更新回调 """ - workflowoper = WorkflowOper() + workflowoper = get_chain_workflow_port() def save_step(action: Action, context: ActionContext, execution_state: dict, completed: bool): """ 保存上下文到数据库 """ - WorkflowOper().step( + get_chain_workflow_port().step( workflow_id, action_id=action.id if completed else "", context=_serialize_workflow_context(context), @@ -1263,18 +1263,18 @@ class WorkflowChain(ChainBase): """ 获取工作流列表 """ - return WorkflowOper().list_enabled() + return get_chain_workflow_port().list_enabled() @staticmethod def get_timer_workflows() -> List[Workflow]: """ 获取定时触发的工作流列表 """ - return WorkflowOper().get_timer_triggered_workflows() + return get_chain_workflow_port().get_timer_triggered_workflows() @staticmethod def get_event_workflows() -> List[Workflow]: """ 获取事件触发的工作流列表 """ - return WorkflowOper().get_event_triggered_workflows() + return get_chain_workflow_port().get_event_triggered_workflows() diff --git a/app/workflow/__init__.py b/app/workflow/__init__.py index 09e3e9e6c..05412b9ac 100644 --- a/app/workflow/__init__.py +++ b/app/workflow/__init__.py @@ -6,7 +6,7 @@ from pydantic import BaseModel from app.runtime.config import global_vars from app.runtime.events import eventmanager, Event -from app.application.chain.data import WorkflowPortProxy as WorkflowOper +from app.application.chain.data import get_chain_workflow_port from app.foundation.reflection import ModuleHelper from app.runtime.log import logger from app.schemas.workflow import ActionContext @@ -282,11 +282,11 @@ class WorkFlowManager(metaclass=Singleton): """ workflows = [] if workflow_id: - workflow = WorkflowOper().get(workflow_id) + workflow = get_chain_workflow_port().get(workflow_id) if workflow: workflows = [workflow] else: - workflows = WorkflowOper().get_event_triggered_workflows() + workflows = get_chain_workflow_port().get_event_triggered_workflows() try: for workflow in workflows: self.update_workflow_event(workflow) @@ -359,7 +359,7 @@ class WorkFlowManager(metaclass=Singleton): """ try: # 检查工作流是否存在且启用 - workflow = WorkflowOper().get(workflow_id) + workflow = get_chain_workflow_port().get(workflow_id) if not workflow or workflow.state == 'P': return diff --git a/app/workflow/actions/add_subscribe.py b/app/workflow/actions/add_subscribe.py index be77beca3..a6219b632 100644 --- a/app/workflow/actions/add_subscribe.py +++ b/app/workflow/actions/add_subscribe.py @@ -3,7 +3,7 @@ from app.chain.subscribe import SubscribeChain from app.application.configuration import get_chain_runtime_config_snapshot from app.runtime.config import global_vars from app.domain.context import MediaInfo -from app.application.chain.data import SubscribePortProxy as SubscribeOper +from app.application.chain.data import get_chain_subscribe_port from app.runtime.log import logger from app.schemas.workflow import ActionParams from app.schemas.workflow import ActionContext @@ -75,7 +75,7 @@ class AddSubscribeAction(BaseAction): if self._added_subscribes: logger.info(f"已添加 {len(self._added_subscribes)} 个订阅") for sid in self._added_subscribes: - context.subscribes.append(SubscribeOper().get(sid)) + context.subscribes.append(get_chain_subscribe_port().get(sid)) elif _started: self._has_error = True diff --git a/app/workflow/actions/transfer_file.py b/app/workflow/actions/transfer_file.py index 20041e930..0663523d9 100644 --- a/app/workflow/actions/transfer_file.py +++ b/app/workflow/actions/transfer_file.py @@ -6,7 +6,7 @@ from pydantic import Field from app.workflow.actions import BaseAction from app.runtime.config import global_vars -from app.application.chain.data import TransferHistoryPortProxy as TransferHistoryOper +from app.application.chain.data import get_chain_transfer_history_port from app.schemas.workflow import ActionParams from app.schemas.workflow import ActionContext from app.chain.storage import StorageChain @@ -67,7 +67,7 @@ class TransferFileAction(BaseAction): _failed_count = 0 storagechain = StorageChain() transferchain = TransferChain() - transferhis = TransferHistoryOper() + transferhis = get_chain_transfer_history_port() if params.source == "downloads": # 从下载任务中整理文件 for download in context.downloads: diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index ce8f06ccc..6b393dbad 100644 --- a/docs/refactor/backend-architecture-next-stage.md +++ b/docs/refactor/backend-architecture-next-stage.md @@ -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 形别名。 +> 实施进度:阶段 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 数据端口迁移。 ## 当前复核结论(2026-08-24) @@ -140,6 +140,16 @@ - 兼容边界不变:DB Oper 类、`app.db.oper` 懒导出、SDK/Compat 旧路径和 V1/V2/V3 插件加载均未改动; 已注入 `SystemConfigReader/SystemConfigService` 的 Agent 工具构造合同保持原样。 +### 长期整改阶段 12:工作流域 Chain 数据端口收口(2026-08-24) + +- `app.application.chain.data` 已提供命名 `get_chain_*_port()`,但 `WorkFlowManager`、`WorkflowChain`、 + 添加订阅和整理文件动作仍把 `*PortProxy` 别名为数据库 Oper,形成 getter 与代理双轨。四个工作流消费者 + 现统一使用显式 getter,测试替换同一个组合根接缝。 +- 架构门禁禁止 `app/workflow/**` 和 `app/chain/workflow.py` 重新导入任何 `*PortProxy`。其它 Chain 域仍按 + 独立风险切片迁移,不能因本阶段通过而宣称全部 Chain 数据代理已清零。 +- 兼容边界不变:全部 `*PortProxy` 类、`WorkFlowManager` 类路径/Singleton identity、动作类型与参数、 + 工作流事件和数据库 `WorkflowOper/SubscribeOper/TransferHistoryOper` 均保留,插件无需迁移。 + ### 总体判断 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: diff --git a/docs/rules/05-architecture.md b/docs/rules/05-architecture.md index f936ca3a5..1d6b13958 100644 --- a/docs/rules/05-architecture.md +++ b/docs/rules/05-architecture.md @@ -131,6 +131,10 @@ Session. `app/db/adapters/` is the concrete persistence-adapter layer: it may depend on Application-owned Protocols, UoW/Session and Oper implementations. This deliberate dependency inversion is the only `DB implementation -> Application contract` direction; Application must remain free of DB imports. +Workflow-domain consumers use the named `get_chain_*_port()` functions from +`app/application/chain/data.py`; they must not alias migration-time `*PortProxy` +classes back to database Oper names. Those proxy classes remain compatibility +boundaries while the other established Chain domains migrate independently. ### Adapter boundaries diff --git a/tests/test_architecture_dependencies.py b/tests/test_architecture_dependencies.py index 13f426b58..70b4dcd20 100644 --- a/tests/test_architecture_dependencies.py +++ b/tests/test_architecture_dependencies.py @@ -407,6 +407,27 @@ def test_host_code_uses_explicit_runtime_facade_getters(): assert violations == [] +def test_workflow_domain_uses_explicit_chain_data_port_getters(): + """工作流域不得重新使用迁移期 PortProxy 冒充数据库 Oper。""" + paths = [APP_ROOT / "chain" / "workflow.py"] + paths.extend((APP_ROOT / "workflow").rglob("*.py")) + violations: list[str] = [] + for path in paths: + tree = ast.parse(path.read_text(encoding="utf-8-sig"), filename=str(path)) + for node in ast.walk(tree): + if not isinstance(node, ast.ImportFrom): + continue + if node.module != "app.application.chain.data": + continue + for alias in node.names: + if alias.name.endswith("PortProxy"): + violations.append( + f"{path.relative_to(PROJECT_ROOT).as_posix()}:{node.lineno}:{alias.name}" + ) + + assert violations == [] + + def test_plugin_components_do_not_reexport_legacy_abi_names(): """新插件组件只提供 canonical 能力,不得复制旧 Helper、Manager 或 Oper 导出。""" violations: list[str] = [] diff --git a/tests/test_workflow_execution.py b/tests/test_workflow_execution.py index 2d23bff8b..a7992a812 100644 --- a/tests/test_workflow_execution.py +++ b/tests/test_workflow_execution.py @@ -635,7 +635,7 @@ def test_workflow_chain_process_serializes_circular_context(monkeypatch): fake_oper = _FakeWorkflowOper(workflow) monkeypatch.setattr(workflow_module, "get_workflow_manager", lambda: fake_manager) - monkeypatch.setattr(workflow_module, "WorkflowOper", lambda: fake_oper) + monkeypatch.setattr(workflow_module, "get_chain_workflow_port", lambda: fake_oper) monkeypatch.setattr(workflow_module.global_vars, "workflow_resume", lambda workflow_id: None) monkeypatch.setattr(workflow_module.global_vars, "is_workflow_stopped", lambda workflow_id: False) diff --git a/tests/test_workflow_runtime_config.py b/tests/test_workflow_runtime_config.py index a0497b207..564773350 100644 --- a/tests/test_workflow_runtime_config.py +++ b/tests/test_workflow_runtime_config.py @@ -123,7 +123,7 @@ def test_add_subscribe_uses_superuser_from_chain_snapshot(monkeypatch): monkeypatch.setattr(add_subscribe_module.global_vars, "is_workflow_stopped", lambda _: False) monkeypatch.setattr( add_subscribe_module, - "SubscribeOper", + "get_chain_subscribe_port", lambda: SimpleNamespace(get=lambda sid: sid), )