diff --git a/app/agent/tools/impl/query_workflows.py b/app/agent/tools/impl/query_workflows.py index ee67fea83..00d7f3412 100644 --- a/app/agent/tools/impl/query_workflows.py +++ b/app/agent/tools/impl/query_workflows.py @@ -7,7 +7,7 @@ from pydantic import BaseModel, Field from app.agent.tools.base import MoviePilotTool from app.agent.tools.tags import ToolTag -from app.application.agentdata import get_agent_workflow_port +from app.application.workflow import get_configured_workflow_query from app.runtime.log import logger @@ -54,8 +54,7 @@ class QueryWorkflowsTool(MoviePilotTool): logger.info(f"执行工具: {self.name}, 参数: state={state}, name={name}, trigger_type={trigger_type}") try: - workflow_oper = get_agent_workflow_port() - workflows = await workflow_oper.async_list() + workflows = await get_configured_workflow_query().list() # 过滤工作流 filtered_workflows = [] @@ -101,7 +100,11 @@ class QueryWorkflowsTool(MoviePilotTool): "event": "事件触发", "manual": "手动触发" } - trigger_type_desc = trigger_type_map.get(wf.trigger_type, wf.trigger_type or "定时触发") + trigger_type_key = wf.trigger_type or "timer" + trigger_type_desc = trigger_type_map.get( + trigger_type_key, + trigger_type_key, + ) simplified = { "id": wf.id, diff --git a/app/api/dependencies/workflow.py b/app/api/dependencies/workflow.py index 995e1aa3b..4837ae51c 100644 --- a/app/api/dependencies/workflow.py +++ b/app/api/dependencies/workflow.py @@ -61,8 +61,7 @@ def get_workflow_definition_command( def get_workflow_query_service( - db: AsyncSession = Depends(get_async_session), runtime: HostRuntime = Depends(get_host_runtime), ) -> WorkflowQueryService: - """组装工作流只读查询用例,避免端点直接持有数据库操作器。""" - return WorkflowQueryService(repository=runtime.workflow.repository(db)) + """返回组合根装配的工作流只读查询用例。""" + return runtime.workflow.query diff --git a/app/application/agentdata.py b/app/application/agentdata.py index 6b475d2e5..1843b167f 100644 --- a/app/application/agentdata.py +++ b/app/application/agentdata.py @@ -5,7 +5,6 @@ from __future__ import annotations from collections.abc import Callable from typing import Any - AgentDataFactory = Callable[[], Any] @@ -76,12 +75,6 @@ class DownloadHistoryPort(_PortProxy): port_name = "download_history" -class WorkflowPort(_PortProxy): - """工作流数据端口代理。""" - - port_name = "workflow" - - class PluginDataPort(_PortProxy): """插件数据端口代理。""" @@ -110,7 +103,6 @@ def configure_agent_data_ports(**factories: AgentDataFactory) -> None: "subscribe_history", "transfer_history", "download_history", - "workflow", "plugin_data", } missing = sorted(required - factories.keys()) @@ -167,11 +159,6 @@ def get_agent_download_history_port() -> Any: return get_agent_data_ports().download_history() -def get_agent_workflow_port() -> Any: - """创建 Agent 工作流数据端口实例。""" - return get_agent_data_ports().workflow() - - def get_agent_plugin_data_port() -> Any: """创建 Agent 插件数据端口实例。""" return get_agent_data_ports().plugin_data() diff --git a/app/application/chain/context.py b/app/application/chain/context.py index 465e0dcc6..20a9e29ad 100644 --- a/app/application/chain/context.py +++ b/app/application/chain/context.py @@ -6,7 +6,6 @@ from collections.abc import Callable from dataclasses import dataclass, field from typing import Any, Optional -from app.application.chain.data import ChainDataPorts from app.application.chain.events import ChainDurableEventWriter from app.application.configuration import ChainRuntimeConfig from app.runtime.stop import StopState, runtime_stop_state @@ -31,7 +30,6 @@ class ChainRuntimeContext: message_queue_factory: MessageQueueFactory module_dispatcher_factory: ModuleDispatcherFactory legacy_transfer_command: Optional[LegacyTransferCommand] = None - data_ports: Optional[ChainDataPorts] = None durable_event_writer: Optional[ChainDurableEventWriter] = None configuration: ChainRuntimeConfig = field( default_factory=lambda: ChainRuntimeConfig(media_extensions=()) diff --git a/app/application/chain/data.py b/app/application/chain/data.py index 92c92ea75..2714285f8 100644 --- a/app/application/chain/data.py +++ b/app/application/chain/data.py @@ -24,7 +24,6 @@ class ChainDataPorts: site: OperFactory subscribe: OperFactory - workflow: OperFactory download_history: OperFactory transfer_history: OperFactory transfer_pending: TransferAdmissionRepositoryFactory @@ -34,94 +33,34 @@ class ChainDataPorts: user: OperFactory -class _PortProxyMeta(type): - """让迁移期的 Oper 名称支持按方法打桩,同时仍转发到组合根端口。""" - - def __getattr__(cls, name: str) -> Any: - """把类级方法访问转发到一个新的端口实例。""" - return getattr(cls(), name) - - -class _ChainDataPortProxy(metaclass=_PortProxyMeta): - """将旧的 Oper 调用形态转发到 Chain 数据端口的内部代理。""" - - port_name: str - - def __getattr__(self, name: str) -> Any: - """转发未被测试替换的数据操作。""" - return getattr(getattr(get_chain_data_ports(), self.port_name)(), name) - - -class SitePortProxy(_ChainDataPortProxy): - """站点数据端口代理。""" - - port_name = "site" - - -class SubscribePortProxy(_ChainDataPortProxy): - """订阅数据端口代理。""" - - port_name = "subscribe" - - -class WorkflowPortProxy(_ChainDataPortProxy): - """工作流数据端口代理。""" - - port_name = "workflow" - - -class DownloadHistoryPortProxy(_ChainDataPortProxy): - """下载历史数据端口代理。""" - - port_name = "download_history" - - -class TransferHistoryPortProxy(_ChainDataPortProxy): - """整理历史数据端口代理。""" - - port_name = "transfer_history" - - -class MediaServerPortProxy(_ChainDataPortProxy): - """媒体服务器数据端口代理。""" - - port_name = "media_server" - - -class DownloadFailurePortProxy(_ChainDataPortProxy): - """下载失败数据端口代理。""" - - port_name = "download_failure" - - -class UserPortProxy(_ChainDataPortProxy): - """用户数据端口代理。""" - - port_name = "user" - - _ports: Optional[ChainDataPorts] = None -def configure_chain_data_ports(**factories: OperFactory) -> None: - """由启动组合根登记 Chain 的数据端口实现。""" - required = { - "site", - "subscribe", - "workflow", - "download_history", - "transfer_history", - "transfer_pending", - "transfer_execution", - "media_server", - "download_failure", - "user", - } - missing = sorted(required - factories.keys()) - if missing: - raise ValueError(f"Chain 数据端口缺少实现: {', '.join(missing)}") +def configure_chain_data_ports( + *, + site: OperFactory, + subscribe: OperFactory, + download_history: OperFactory, + transfer_history: OperFactory, + transfer_pending: TransferAdmissionRepositoryFactory, + transfer_execution: TransferExecutionRepositoryFactory, + media_server: OperFactory, + download_failure: OperFactory, + user: OperFactory, +) -> None: + """由启动组合根登记显式命名的 Chain 数据端口实现。""" global _ports - _ports = ChainDataPorts(**{name: factories[name] for name in required}) + _ports = ChainDataPorts( + site=site, + subscribe=subscribe, + download_history=download_history, + transfer_history=transfer_history, + transfer_pending=transfer_pending, + transfer_execution=transfer_execution, + media_server=media_server, + download_failure=download_failure, + user=user, + ) def get_chain_data_ports() -> ChainDataPorts: @@ -141,11 +80,6 @@ def get_chain_subscribe_port() -> Any: return get_chain_data_ports().subscribe() -def get_chain_workflow_port() -> Any: - """创建工作流数据端口实例。""" - return get_chain_data_ports().workflow() - - def get_chain_download_history_port() -> Any: """创建下载历史数据端口实例。""" return get_chain_data_ports().download_history() diff --git a/app/application/server/share.py b/app/application/server/share.py index d1f468a57..19f840e9d 100644 --- a/app/application/server/share.py +++ b/app/application/server/share.py @@ -4,8 +4,10 @@ from __future__ import annotations import json from collections.abc import Awaitable, Callable +from dataclasses import asdict from typing import Any, Optional +from app.application.workflow import WorkflowSnapshot from app.schemas.media import resolve_media_identity @@ -26,8 +28,10 @@ class ServerSharingService: *, subscribe_provider: Callable[[int], Any], async_subscribe_provider: Callable[[int], Awaitable[Any]], - workflow_provider: Callable[[int], Any], - async_workflow_provider: Callable[[int], Awaitable[Any]], + workflow_provider: Callable[[int], Optional[WorkflowSnapshot]], + async_workflow_provider: Callable[ + [int], Awaitable[Optional[WorkflowSnapshot]] + ], user_uuid_provider: Callable[[], str], subscribe_sender: Callable[[dict], Any], async_subscribe_sender: Callable[[dict], Awaitable[Any]], @@ -68,9 +72,9 @@ class ServerSharingService: return payload @staticmethod - def prepare_workflow(workflow: Any) -> dict: + def prepare_workflow(workflow: WorkflowSnapshot) -> dict: """移除本地字段并把动作和流程编码为中心服务兼容格式。""" - workflow_dict = workflow.to_dict() + workflow_dict = asdict(workflow) workflow_dict.pop("id", None) workflow_dict.pop("context", None) workflow_dict["actions"] = json.dumps(workflow_dict["actions"] or []) @@ -78,7 +82,9 @@ class ServerSharingService: return workflow_dict @staticmethod - def validate_workflow(workflow: Any) -> tuple[bool, str]: + def validate_workflow( + workflow: Optional[WorkflowSnapshot], + ) -> tuple[bool, str]: """验证工作流存在且同时包含动作与流程。""" if not workflow: return False, "工作流不存在" @@ -160,6 +166,8 @@ class ServerSharingService: valid, message = self.validate_workflow(workflow) if not valid: return False, message + if workflow is None: + return False, "工作流不存在" payload = { "share_title": share_title, "share_comment": share_comment, @@ -188,6 +196,8 @@ class ServerSharingService: valid, message = self.validate_workflow(workflow) if not valid: return False, message + if workflow is None: + return False, "工作流不存在" payload = { "share_title": share_title, "share_comment": share_comment, diff --git a/app/application/workflow.py b/app/application/workflow.py index 962ae8a69..c39ed0a6d 100644 --- a/app/application/workflow.py +++ b/app/application/workflow.py @@ -1,11 +1,12 @@ """工作流状态与定义写操作应用用例。""" -from dataclasses import dataclass import json from collections.abc import Awaitable +from dataclasses import dataclass from datetime import datetime -from typing import Any, Callable, Mapping, Optional, Protocol, TypeVar +from typing import Any, Callable, List, Mapping, Optional, Protocol, TypeVar +from app.schemas.common import JsonData WORKFLOW_TRIGGER_TIMER = "timer" WORKFLOW_TRIGGER_EVENT = "event" @@ -17,6 +18,30 @@ SUPPORTED_WORKFLOW_TRIGGERS = { } +@dataclass(frozen=True, slots=True) +class WorkflowSnapshot: + """工作流查询返回的脱离数据库会话的冻结快照。""" + + id: int + name: str + description: Optional[str] + timer: Optional[str] + trigger_type: Optional[str] + event_type: Optional[str] + event_conditions: Mapping[str, JsonData] + state: str + current_action: Optional[str] + result: Optional[str] + run_count: Optional[int] + actions: tuple[Mapping[str, JsonData], ...] + flows: tuple[Mapping[str, JsonData], ...] + context: Mapping[str, JsonData] + execution_config: Mapping[str, JsonData] + execution_state: Mapping[str, JsonData] + add_time: Optional[str] + last_time: Optional[str] + + class WorkflowRuntime(Protocol): """声明宿主入口与 Chain 消费的工作流运行时能力。""" @@ -40,7 +65,7 @@ class WorkflowRuntime(Protocol): """移除全部或指定工作流的事件触发器。""" ... - def update_workflow_event(self, workflow: Any) -> None: + def update_workflow_event(self, workflow: WorkflowSnapshot) -> None: """按最新定义刷新工作流事件触发器。""" ... @@ -87,15 +112,31 @@ def get_workflow_manager() -> WorkflowRuntime: return _workflow_runtime_provider() -class AsyncWorkflowQueryRepository(Protocol): - """工作流查询用例需要的异步读取端口。""" +class WorkflowQueryRepository(Protocol): + """工作流查询用例需要的同步与异步快照端口。""" - async def async_list(self) -> list[Any]: - """读取全部工作流。""" + def get(self, workflow_id: int) -> Optional[WorkflowSnapshot]: + """按 ID 读取工作流快照。""" ... - async def async_get(self, workflow_id: int) -> Optional[Any]: - """按 ID 读取工作流。""" + def list_enabled(self) -> List[WorkflowSnapshot]: + """读取全部启用的工作流快照。""" + ... + + def list_timer_enabled(self) -> List[WorkflowSnapshot]: + """读取启用的定时工作流快照。""" + ... + + def list_event_enabled(self) -> List[WorkflowSnapshot]: + """读取启用的事件工作流快照。""" + ... + + async def async_list(self) -> List[WorkflowSnapshot]: + """异步读取全部工作流快照。""" + ... + + async def async_get(self, workflow_id: int) -> Optional[WorkflowSnapshot]: + """异步按 ID 读取工作流快照。""" ... @@ -114,18 +155,34 @@ class WorkflowCachePort(Protocol): class WorkflowQueryService: """提供工作流列表和详情查询,隔离 API 与数据库会话。""" - def __init__(self, repository: AsyncWorkflowQueryRepository) -> None: - """保存请求级异步查询端口。""" + def __init__(self, repository: WorkflowQueryRepository) -> None: + """保存可返回脱离会话快照的查询端口。""" self._repository = repository - async def list(self) -> list[Any]: - """返回全部工作流。""" + async def list(self) -> List[WorkflowSnapshot]: + """返回全部工作流快照。""" return await self._repository.async_list() - async def get(self, workflow_id: int) -> Optional[Any]: - """返回指定工作流。""" + async def get(self, workflow_id: int) -> Optional[WorkflowSnapshot]: + """返回指定工作流快照。""" return await self._repository.async_get(workflow_id) + def get_sync(self, workflow_id: int) -> Optional[WorkflowSnapshot]: + """同步返回指定工作流快照。""" + return self._repository.get(workflow_id) + + def list_enabled(self) -> List[WorkflowSnapshot]: + """同步返回全部启用的工作流快照。""" + return self._repository.list_enabled() + + def list_timer_enabled(self) -> List[WorkflowSnapshot]: + """同步返回启用的定时工作流快照。""" + return self._repository.list_timer_enabled() + + def list_event_enabled(self) -> List[WorkflowSnapshot]: + """同步返回启用的事件工作流快照。""" + return self._repository.list_event_enabled() + _configured_workflow_query: WorkflowQueryService | None = None @@ -183,6 +240,56 @@ class UnitOfWork(Protocol): ... +class WorkflowExecutionPort(Protocol): + """工作流 Chain 提交执行状态所需的类型化事务端口。""" + + def start(self, workflow_id: int) -> bool: + """提交工作流运行中状态。""" + ... + + def success( + self, + workflow_id: int, + result: Optional[str] = None, + ) -> bool: + """提交工作流成功状态。""" + ... + + def fail(self, workflow_id: int, result: str) -> bool: + """提交工作流失败状态。""" + ... + + def step( + self, + workflow_id: int, + action_id: str, + context: dict[str, Any], + execution_state: Optional[dict[str, Any]] = None, + ) -> bool: + """提交工作流动作进度。""" + ... + + def reset(self, workflow_id: int, reset_count: bool = False) -> bool: + """提交工作流执行状态重置。""" + ... + + +_configured_workflow_execution: Optional[WorkflowExecutionPort] = None + + +def configure_workflow_execution(service: WorkflowExecutionPort) -> None: + """由启动组合根登记唯一工作流执行状态事务服务。""" + global _configured_workflow_execution + _configured_workflow_execution = service + + +def get_configured_workflow_execution() -> WorkflowExecutionPort: + """返回启动阶段登记的工作流执行状态事务服务。""" + if _configured_workflow_execution is None: + raise RuntimeError("工作流执行状态事务服务尚未配置") + return _configured_workflow_execution + + class WorkflowExecutionRepository(Protocol): """工作流执行状态写入所需的最小暂存端口。""" diff --git a/app/chain/__init__.py b/app/chain/__init__.py index b6ccd5315..4e7770c21 100644 --- a/app/chain/__init__.py +++ b/app/chain/__init__.py @@ -8,7 +8,6 @@ from pathlib import Path from typing import TYPE_CHECKING, Any, Dict, List, Optional, Set, Tuple, Union, cast from app.application.chain.context import ChainRuntimeContext, get_chain_runtime_context -from app.application.chain.data import get_chain_data_ports from app.application.configuration import ( ChainRuntimeConfig, get_chain_runtime_config_snapshot, @@ -59,7 +58,6 @@ class ChainBase(RecognitionMixin, MessageProcessingMixin, NotificationMixin, self.async_filecache = context.async_file_cache self.runtime_config = context.configuration self.stop_state = context.stop_state - self.data_ports = context.data_ports or get_chain_data_ports() self.durable_event_writer = context.durable_event_writer self._module_dispatcher = context.module_dispatcher_factory( module_catalog=self.modulemanager, diff --git a/app/chain/workflow.py b/app/chain/workflow.py index 2d40b4d50..e3c14606a 100644 --- a/app/chain/workflow.py +++ b/app/chain/workflow.py @@ -13,8 +13,12 @@ from typing import Any, Callable, List, Optional, Tuple from pydantic import BaseModel -from app.application.chain.data import get_chain_workflow_port -from app.application.workflow import get_workflow_manager +from app.application.workflow import ( + WorkflowSnapshot, + get_configured_workflow_execution, + get_configured_workflow_query, + get_workflow_manager, +) from app.chain import ChainBase from app.runtime.events import Event, eventmanager from app.runtime.execution import OwnedThreadPoolExecutor @@ -27,9 +31,6 @@ ARTIFACT_FIELDS = {"torrents", "medias", "fileitems", "downloads", "sites", "sub DEFAULT_WORKFLOW_MAX_WORKERS = 4 WORKFLOW_EXECUTOR_STOP_TIMEOUT_SECONDS = 10.0 CIRCULAR_REFERENCE_PLACEHOLDER = "[Circular]" -Workflow = Any - - def _serialize_workflow_key(key: Any) -> Any: """将映射键转换为 JSON 安全值。""" if key is None or isinstance(key, (str, int, float, bool)): @@ -118,7 +119,11 @@ class WorkflowExecutor: 工作流执行器 """ - def __init__(self, workflow: Workflow, step_callback: Callable = None): + def __init__( + self, + workflow: WorkflowSnapshot, + step_callback: Callable = None, + ): """ 初始化工作流执行器 :param workflow: 工作流对象 @@ -132,8 +137,13 @@ class WorkflowExecutor: if step_callback else False ) - self.actions = {action['id']: Action(**action) for action in workflow.actions} - self.flows = [ActionFlow(**flow) for flow in workflow.flows] + self.actions: dict[str, Action] = {} + for action_data in workflow.actions: + action = Action(**dict(action_data)) + if not action.id: + raise ValueError("工作流动作缺少 ID") + self.actions[action.id] = action + self.flows = [ActionFlow(**dict(flow)) for flow in workflow.flows] execution_config = getattr(workflow, "execution_config", None) or {} execution_state = getattr(workflow, "execution_state", None) or {} self.execution_config = ( @@ -652,7 +662,8 @@ class WorkflowExecutor: self.flow_satisfied.add(flow_key) if not source_success and self.node_states.get(source_id) == "failed": self.flow_failed.add(flow_key) - self.evaluate_target_state(flow.target) + if flow.target: + self.evaluate_target_state(flow.target) def evaluate_target_state(self, target_id: str) -> None: """ @@ -1275,15 +1286,35 @@ class WorkflowChain(ChainBase): :param from_begin: 是否从头开始,默认为True :param progress_callback: 定时服务进度更新回调 """ - workflowoper = get_chain_workflow_port() + workflow_execution = get_configured_workflow_execution() - def save_step(action: Action, context: ActionContext, execution_state: dict, completed: bool): - """ - 保存上下文到数据库 - """ - get_chain_workflow_port().step( + # 重置工作流 + if from_begin: + workflow_execution.reset(workflow_id) + + # 查询工作流数据 + workflow = get_configured_workflow_query().get_sync(workflow_id) + if not workflow: + logger.warn(f"工作流 {workflow_id} 不存在") + return False, "工作流不存在" + if not workflow.actions: + logger.warn(f"工作流 {workflow.name} 无动作") + return False, "工作流无动作" + if not workflow.flows: + logger.warn(f"工作流 {workflow.name} 无流程") + return False, "工作流无流程" + + def save_step( + action: Action, + context: ActionContext, + execution_state: dict, + completed: bool, + ) -> None: + """保存动作上下文和结构化执行状态。""" + persisted_action_id = (action.id or "") if completed else "" + workflow_execution.step( workflow_id, - action_id=action.id if completed else "", + action_id=persisted_action_id, context=_serialize_workflow_context(context), execution_state=_serialize_workflow_value(execution_state) ) @@ -1305,22 +1336,6 @@ class WorkflowChain(ChainBase): }, ) - # 重置工作流 - if from_begin: - workflowoper.reset(workflow_id) - - # 查询工作流数据 - workflow = workflowoper.get(workflow_id) - if not workflow: - logger.warn(f"工作流 {workflow_id} 不存在") - return False, "工作流不存在" - if not workflow.actions: - logger.warn(f"工作流 {workflow.name} 无动作") - return False, "工作流无动作" - if not workflow.flows: - logger.warn(f"工作流 {workflow.name} 无流程") - return False, "工作流无流程" - logger.info(f"开始执行工作流 {workflow.name},共 {len(workflow.actions)} 个动作 ...") if progress_callback: progress_callback( @@ -1334,7 +1349,7 @@ class WorkflowChain(ChainBase): logger.warning("工作流服务正在停止,拒绝执行 %s", workflow.name) return False, executor.errmsg try: - workflowoper.start(workflow_id) + workflow_execution.start(workflow_id) except Exception: executor.abort_before_execute() raise @@ -1346,31 +1361,31 @@ class WorkflowChain(ChainBase): if not executor.success or executor.has_failure: logger.info(f"工作流 {workflow.name} 执行失败:{executor.errmsg}") - workflowoper.fail(workflow_id, result=executor.errmsg) + workflow_execution.fail(workflow_id, result=executor.errmsg) return False, executor.errmsg logger.info(f"工作流 {workflow.name} 执行完成") - workflowoper.success(workflow_id) + workflow_execution.success(workflow_id) if progress_callback: progress_callback(value=100, text=f"工作流 {workflow.name} 执行完成") return True, "" @staticmethod - def get_workflows() -> List[Workflow]: + def get_workflows() -> List[WorkflowSnapshot]: """ 获取工作流列表 """ - return get_chain_workflow_port().list_enabled() + return get_configured_workflow_query().list_enabled() @staticmethod - def get_timer_workflows() -> List[Workflow]: + def get_timer_workflows() -> List[WorkflowSnapshot]: """ 获取定时触发的工作流列表 """ - return get_chain_workflow_port().get_timer_triggered_workflows() + return get_configured_workflow_query().list_timer_enabled() @staticmethod - def get_event_workflows() -> List[Workflow]: + def get_event_workflows() -> List[WorkflowSnapshot]: """ 获取事件触发的工作流列表 """ - return get_chain_workflow_port().get_event_triggered_workflows() + return get_configured_workflow_query().list_event_enabled() diff --git a/app/db/adapters/workflow.py b/app/db/adapters/workflow.py index 5f69b6dc2..15a90aba2 100644 --- a/app/db/adapters/workflow.py +++ b/app/db/adapters/workflow.py @@ -1,18 +1,144 @@ -"""工作流执行状态事务适配器。""" +"""工作流查询与执行状态事务适配器。""" -from collections.abc import Callable -from typing import Any, TypeVar +from collections.abc import Callable, Iterable +from contextlib import AbstractAsyncContextManager +from copy import deepcopy +from typing import Any, Optional, TypeVar, cast +from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.orm import Session -from app.application.workflow import WorkflowExecutionCommand +from app.application.workflow import WorkflowExecutionCommand, WorkflowSnapshot from app.db.oper.workflow import WorkflowOper from app.db.uow import SqlAlchemyUnitOfWork - +from app.schemas.common import JsonData _Result = TypeVar("_Result") +def _copy_json_mapping(value: object, field_name: str) -> dict[str, JsonData]: + """复制 ORM JSON 对象,拒绝把损坏结构带出 Session。""" + if value is None: + return {} + if not isinstance(value, dict): + raise ValueError(f"工作流 {field_name} 必须是 JSON 对象") + return cast(dict[str, JsonData], deepcopy(value)) + + +def _copy_json_sequence( + value: object, + field_name: str, +) -> tuple[dict[str, JsonData], ...]: + """复制 ORM JSON 对象序列,确保快照不共享可变容器。""" + if value is None: + return () + if not isinstance(value, list) or any(not isinstance(item, dict) for item in value): + raise ValueError(f"工作流 {field_name} 必须是 JSON 对象数组") + return tuple(cast(dict[str, JsonData], deepcopy(item)) for item in value) + + +def _project_workflow(record: object) -> WorkflowSnapshot: + """在持有数据库会话时把 ORM 记录投影为稳定快照。""" + workflow_id = getattr(record, "id", None) + name = getattr(record, "name", None) + state = getattr(record, "state", None) + if not isinstance(workflow_id, int) or not isinstance(name, str) or not isinstance(state, str): + raise ValueError("工作流记录缺少稳定身份或状态") + return WorkflowSnapshot( + id=workflow_id, + name=name, + description=getattr(record, "description", None), + timer=getattr(record, "timer", None), + trigger_type=getattr(record, "trigger_type", None), + event_type=getattr(record, "event_type", None), + event_conditions=_copy_json_mapping( + getattr(record, "event_conditions", None), + "event_conditions", + ), + state=state, + current_action=getattr(record, "current_action", None), + result=getattr(record, "result", None), + run_count=getattr(record, "run_count", None), + actions=_copy_json_sequence(getattr(record, "actions", None), "actions"), + flows=_copy_json_sequence(getattr(record, "flows", None), "flows"), + context=_copy_json_mapping(getattr(record, "context", None), "context"), + execution_config=_copy_json_mapping( + getattr(record, "execution_config", None), + "execution_config", + ), + execution_state=_copy_json_mapping( + getattr(record, "execution_state", None), + "execution_state", + ), + add_time=getattr(record, "add_time", None), + last_time=getattr(record, "last_time", None), + ) + + +class TransactionalWorkflowQueryRepository: + """在自有短 Session 内查询并投影工作流快照。""" + + def __init__( + self, + sync_session: Callable[[], Session], + async_session: Callable[[], AbstractAsyncContextManager[AsyncSession]], + ) -> None: + """保存同步 Session 工厂与异步 Session 作用域。""" + self._sync_session = sync_session + self._async_session = async_session + + def get(self, workflow_id: int) -> Optional[WorkflowSnapshot]: + """在同步短 Session 内读取并投影单条工作流。""" + session = self._sync_session() + try: + record = WorkflowOper(session).get(workflow_id) + return _project_workflow(record) if record else None + finally: + session.close() + + def list_enabled(self) -> list[WorkflowSnapshot]: + """在同步短 Session 内投影全部启用工作流。""" + return self._list_sync(lambda repository: repository.list_enabled()) + + def list_timer_enabled(self) -> list[WorkflowSnapshot]: + """在同步短 Session 内投影启用的定时工作流。""" + return self._list_sync( + lambda repository: repository.get_timer_triggered_workflows() + ) + + def list_event_enabled(self) -> list[WorkflowSnapshot]: + """在同步短 Session 内投影启用的事件工作流。""" + return self._list_sync( + lambda repository: repository.get_event_triggered_workflows() + ) + + async def async_list(self) -> list[WorkflowSnapshot]: + """在异步短 Session 内投影全部工作流。""" + async with self._async_session() as session: + records = await WorkflowOper(session).async_list() + return [_project_workflow(record) for record in records] + + async def async_get(self, workflow_id: int) -> Optional[WorkflowSnapshot]: + """在异步短 Session 内读取并投影单条工作流。""" + async with self._async_session() as session: + record = await WorkflowOper(session).async_get(workflow_id) + return _project_workflow(record) if record else None + + def _list_sync( + self, + operation: Callable[[WorkflowOper], Iterable[object]], + ) -> list[WorkflowSnapshot]: + """在同步短 Session 内执行列表查询并完成投影。""" + session = self._sync_session() + try: + return [ + _project_workflow(record) + for record in operation(WorkflowOper(session)) + ] + finally: + session.close() + + class TransactionalWorkflowExecutionService: """为每次工作流执行状态写入创建独立短会话和 UnitOfWork。""" diff --git a/app/db/oper/__init__.py b/app/db/oper/__init__.py index 07794a1cf..b72cb87ac 100644 --- a/app/db/oper/__init__.py +++ b/app/db/oper/__init__.py @@ -36,7 +36,6 @@ if TYPE_CHECKING: from app.db.oper.transfersettlementreceipt import TransferSettlementReceiptOper from app.db.oper.user import UserOper from app.db.oper.userconfig import UserConfigOper - from app.db.oper.workflow import WorkflowOper # 类名 -> 所在子模块。子模块名即实体名,与 app/db/models 对齐。 _OPER_MODULES = { @@ -56,7 +55,6 @@ _OPER_MODULES = { "TransferSettlementReceiptOper": "transfersettlementreceipt", "UserConfigOper": "userconfig", "UserOper": "user", - "WorkflowOper": "workflow", } @@ -99,5 +97,4 @@ __all__ = [ "TransferSettlementReceiptOper", "UserConfigOper", "UserOper", - "WorkflowOper", ] diff --git a/app/db/oper/workflow.py b/app/db/oper/workflow.py index 51334a411..8c2a6725b 100644 --- a/app/db/oper/workflow.py +++ b/app/db/oper/workflow.py @@ -1,4 +1,4 @@ -from typing import List, Mapping, Tuple, Optional, Any, Protocol +from typing import Any, List, Mapping, Optional, Tuple from sqlalchemy import delete as sqlalchemy_delete from sqlalchemy.orm import Session @@ -7,56 +7,6 @@ from app.db.base import DbOper from app.db.models.workflow import Workflow -class WorkflowLegacyWriter(Protocol): - """无显式 Session 的旧 Oper 写入口所需事务服务。""" - - def start(self, workflow_id: int) -> bool: - """提交工作流运行中状态。""" - ... - - def success( - self, - workflow_id: int, - result: Optional[str] = None, - ) -> bool: - """提交工作流成功状态。""" - ... - - def fail(self, workflow_id: int, result: str) -> bool: - """提交工作流失败状态。""" - ... - - def step( - self, - workflow_id: int, - action_id: str, - context: dict[str, Any], - execution_state: Optional[dict[str, Any]] = None, - ) -> bool: - """提交工作流动作进度。""" - ... - - def reset(self, workflow_id: int, reset_count: bool = False) -> bool: - """提交工作流执行状态重置。""" - ... - - -_legacy_writer: Optional[WorkflowLegacyWriter] = None - - -def configure_workflow_legacy_writer(writer: WorkflowLegacyWriter) -> None: - """由启动组合根为旧的无 Session Oper 写入口注入事务服务。""" - global _legacy_writer - _legacy_writer = writer - - -def _get_workflow_legacy_writer() -> WorkflowLegacyWriter: - """返回已装配的兼容事务服务,避免 Oper 自行创建会话。""" - if _legacy_writer is None: - raise RuntimeError("工作流兼容写服务尚未配置") - return _legacy_writer - - class WorkflowOper(DbOper): """ 工作流管理 @@ -193,67 +143,24 @@ class WorkflowOper(DbOper): workflow.run_count = 0 return workflow - def start(self, wid: int) -> bool: - """ - 启动 - """ - if self._db is None: - return _get_workflow_legacy_writer().start(wid) - return self.stage_start(wid) - def stage_start(self, wid: int) -> bool: """在调用方持有的会话中暂存运行中状态。""" if not isinstance(self._db, Session): raise RuntimeError("工作流暂存写入需要调用方提供同步 Session") return Workflow.start(self._db, wid) - def success(self, wid: int, result: Optional[str] = None) -> bool: - """ - 成功 - """ - if self._db is None: - return _get_workflow_legacy_writer().success(wid, result) - return self.stage_success(wid, result) - def stage_success(self, wid: int, result: Optional[str] = None) -> bool: """在调用方持有的会话中暂存成功状态。""" if not isinstance(self._db, Session): raise RuntimeError("工作流暂存写入需要调用方提供同步 Session") return Workflow.success(self._db, wid, result) - def fail(self, wid: int, result: str) -> bool: - """ - 失败 - """ - if self._db is None: - return _get_workflow_legacy_writer().fail(wid, result) - return self.stage_fail(wid, result) - def stage_fail(self, wid: int, result: str) -> bool: """在调用方持有的会话中暂存失败状态。""" if not isinstance(self._db, Session): raise RuntimeError("工作流暂存写入需要调用方提供同步 Session") return Workflow.fail(self._db, wid, result) - def step( - self, - wid: int, - action_id: str, - context: dict[str, Any], - execution_state: Optional[dict[str, Any]] = None, - ) -> bool: - """ - 步进 - """ - if self._db is None: - return _get_workflow_legacy_writer().step( - wid, - action_id, - context, - execution_state, - ) - return self.stage_step(wid, action_id, context, execution_state) - def stage_step( self, wid: int, @@ -272,14 +179,6 @@ class WorkflowOper(DbOper): execution_state ) - def reset(self, wid: int, reset_count: bool = False) -> bool: - """ - 重置 - """ - if self._db is None: - return _get_workflow_legacy_writer().reset(wid, reset_count) - return self.stage_execution_reset(wid, reset_count) - def stage_execution_reset( self, wid: int, diff --git a/app/runtime/compat/manifest.py b/app/runtime/compat/manifest.py index 2d05c3d79..549779d17 100644 --- a/app/runtime/compat/manifest.py +++ b/app/runtime/compat/manifest.py @@ -171,10 +171,10 @@ MODULE_ALIASES: Dict[str, ModuleAlias] = { owner="db", ), "app.db.workflow_oper": ModuleAlias( - target="app.db.oper.workflow", - replacement="app.db.oper.workflow", + target="app.sdk._legacy.workflow", + replacement="app.application.workflow.WorkflowExecutionPort", introduced="v3.0.0", - owner="db", + owner="sdk", ), "app.utils.crypto": ModuleAlias( target="app.foundation.crypto", @@ -774,6 +774,20 @@ _MESSAGE_NOTIFICATION_SYMBOL_ALIASES: Dict[str, SymbolAlias] = { } SYMBOL_ALIASES: Dict[str, Dict[str, SymbolAlias]] = { + "app.db.oper": { + "WorkflowOper": SymbolAlias( + target_module="app.sdk._legacy.workflow", + target_name="WorkflowOper", + replacement="app.application.workflow.WorkflowExecutionPort", + ), + }, + "app.workflow": { + "WorkFlowManager": SymbolAlias( + target_module="app.workflow", + target_name="WorkflowManager", + replacement="app.workflow.WorkflowManager", + ), + }, "app.application.transfer": { name: SymbolAlias( target_module="app.sdk._legacy.transfer", diff --git a/app/scheduler.py b/app/scheduler.py index 01811e3e2..85ac60140 100644 --- a/app/scheduler.py +++ b/app/scheduler.py @@ -33,6 +33,7 @@ from app.application.outbox import dispatch_pending_outbox from app.application.plugin.routes import register_plugin_api from app.application.plugin.runtime import get_plugin_manager from app.application.site.sites import SitesHelper # pylint: disable=import-error,no-name-in-module +from app.application.workflow import WorkflowSnapshot from app.chain import ChainBase from app.chain.mediaserver import MediaServerChain from app.chain.recommend import RecommendChain @@ -56,7 +57,6 @@ from app.schemas.dashboard import ScheduleProgress as _SchemaScheduleProgress from app.schemas.message import Message, MessageType from app.schemas.system import MediaServerConf as _SchemaMediaServerConf from app.schemas.types import EventType, SystemConfigKey -from app.schemas.workflow import Workflow lock = threading.Lock() SCHEDULER_PROGRESS_PREFIX = "scheduler" @@ -1703,7 +1703,7 @@ class Scheduler(ConfigReloadMixin, metaclass=SingletonClass): for workflow in WorkflowChain().get_timer_workflows() or []: self.update_workflow_job(workflow) - def remove_workflow_job(self, workflow: Workflow): + def remove_workflow_job(self, workflow: WorkflowSnapshot): """ 移除工作流服务 """ @@ -1787,7 +1787,7 @@ class Scheduler(ConfigReloadMixin, metaclass=SingletonClass): role="system", ) - def update_workflow_job(self, workflow: Workflow): + def update_workflow_job(self, workflow: WorkflowSnapshot): """ 更新工作流定时服务 """ diff --git a/app/sdk/_legacy/workflow.py b/app/sdk/_legacy/workflow.py new file mode 100644 index 000000000..e47e43dc1 --- /dev/null +++ b/app/sdk/_legacy/workflow.py @@ -0,0 +1,54 @@ +"""保留旧工作流 Oper 的无 Session 执行状态写入契约。""" + +from typing import Any, Optional + +from app.application.workflow import get_configured_workflow_execution +from app.db.oper.workflow import WorkflowOper as CanonicalWorkflowOper + + +class WorkflowOper(CanonicalWorkflowOper): + """继承显式 Session 查询,并保留旧执行状态写入方法。""" + + def start(self, wid: int) -> bool: + """按旧签名提交工作流运行中状态。""" + if self._db is None: + return get_configured_workflow_execution().start(wid) + return self.stage_start(wid) + + def success(self, wid: int, result: Optional[str] = None) -> bool: + """按旧签名提交工作流成功状态。""" + if self._db is None: + return get_configured_workflow_execution().success(wid, result) + return self.stage_success(wid, result) + + def fail(self, wid: int, result: str) -> bool: + """按旧签名提交工作流失败状态。""" + if self._db is None: + return get_configured_workflow_execution().fail(wid, result) + return self.stage_fail(wid, result) + + def step( + self, + wid: int, + action_id: str, + context: dict[str, Any], + execution_state: Optional[dict[str, Any]] = None, + ) -> bool: + """按旧签名提交工作流动作进度。""" + if self._db is None: + return get_configured_workflow_execution().step( + wid, + action_id, + context, + execution_state, + ) + return self.stage_step(wid, action_id, context, execution_state) + + def reset(self, wid: int, reset_count: bool = False) -> bool: + """按旧签名提交工作流执行状态重置。""" + if self._db is None: + return get_configured_workflow_execution().reset(wid, reset_count) + return self.stage_execution_reset(wid, reset_count) + + +__all__ = ["WorkflowOper"] diff --git a/app/startup/composition/context.py b/app/startup/composition/context.py index e2feb69fb..0ce9e2165 100644 --- a/app/startup/composition/context.py +++ b/app/startup/composition/context.py @@ -17,7 +17,7 @@ from app.application.subscription.mutation import ( SubscriptionHistoryMutationRepository, SubscriptionMutationRepository, ) -from app.application.workflow import WorkflowCachePort +from app.application.workflow import WorkflowCachePort, WorkflowQueryService from app.runtime.tasks import TaskRegistry @@ -165,6 +165,7 @@ class SiteRuntime: class WorkflowRuntime: """工作流定义、状态与缓存操作所需的数据工厂。""" + query: WorkflowQueryService repository: RepositoryFactory system_config: Callable[[], WorkflowCachePort] diff --git a/app/startup/initializers/modules.py b/app/startup/initializers/modules.py index 3b35a7071..02aad1b9b 100644 --- a/app/startup/initializers/modules.py +++ b/app/startup/initializers/modules.py @@ -43,7 +43,7 @@ from app.application.chain.context import ( ChainRuntimeContext, configure_chain_runtime_context_provider, ) -from app.application.chain.data import configure_chain_data_ports, get_chain_data_ports +from app.application.chain.data import configure_chain_data_ports from app.application.chain.events import ( restore_download_added, restore_transfer_result, @@ -106,7 +106,11 @@ from app.application.service import configure_service_directory from app.application.site.health import SiteHealthService, configure_site_health_service from app.application.site.query import SiteQueryService, configure_site_query_service from app.application.subscription.write import configure_subscribe_writer -from app.application.workflow import WorkflowQueryService, configure_workflow_query +from app.application.workflow import ( + WorkflowQueryService, + configure_workflow_execution, + configure_workflow_query, +) from app.command import CommandChain from app.db.adapters.chain import TransactionalChainDurableEventWriter from app.db.adapters.download import TransactionalDownloadFailureRepository @@ -119,7 +123,10 @@ from app.db.adapters.transfer.admission import TransactionalTransferAdmissionRep from app.db.adapters.transfer.execution import ( TransactionalTransferExecutionRepository, ) -from app.db.adapters.workflow import TransactionalWorkflowExecutionService +from app.db.adapters.workflow import ( + TransactionalWorkflowExecutionService, + TransactionalWorkflowQueryRepository, +) from app.db.oper.agentchat import AgentChatOper from app.db.oper.agenttask import AgentTaskOper from app.db.oper.downloadhistory import DownloadHistoryOper @@ -134,7 +141,7 @@ from app.db.oper.systemconfig import SystemConfigOper from app.db.oper.transferhistory import TransferHistoryOper from app.db.oper.user import UserOper from app.db.oper.userconfig import UserConfigOper -from app.db.oper.workflow import WorkflowOper, configure_workflow_legacy_writer +from app.db.oper.workflow import WorkflowOper from app.db.session import ( SessionFactory, async_session_scope, @@ -244,11 +251,6 @@ async def _async_get_subscribe(subscribe_id: int): return await SubscribeOper().async_get(subscribe_id) -async def _async_get_workflow(workflow_id: int): - """通过数据库操作器异步读取工作流,供服务端共享用例使用。""" - return await WorkflowOper().async_get(workflow_id) - - def _execute_legacy_transfer_command(**kwargs: Any) -> Any: """把旧 Chain ABI 延迟转入唯一 TransferChain durable command。""" from app.chain.transfer import TransferChain @@ -272,13 +274,12 @@ def _build_chain_runtime_context() -> ChainRuntimeContext: module_dispatcher_factory=ModuleInvocationDispatcher, legacy_transfer_command=_execute_legacy_transfer_command, configuration=build_chain_runtime_config(legacy_settings), - data_ports=get_chain_data_ports(), durable_event_writer=TransactionalChainDurableEventWriter(SessionFactory), stop_state=runtime_stop_state, ) -def configure_runtime_data_providers() -> None: +def configure_runtime_data_providers(workflow_query: WorkflowQueryService) -> None: """在启动组合层装配运行时和外部服务所需的数据库读取能力。""" configure_service_config_reader(lambda key: get_configured_system_config().get(key)) configure_module_runtime(lambda: ModuleManager()) @@ -314,8 +315,8 @@ def configure_runtime_data_providers() -> None: subscribe_id ), async_subscribe_provider=_async_get_subscribe, - workflow_provider=lambda workflow_id: WorkflowOper().get(workflow_id), - async_workflow_provider=_async_get_workflow, + workflow_provider=workflow_query.get_sync, + async_workflow_provider=workflow_query.get, user_uuid_provider=MoviePilotServerHelper.get_user_uuid, subscribe_sender=MoviePilotServerHelper.subscribe_share, async_subscribe_sender=MoviePilotServerHelper.async_subscribe_share, @@ -812,6 +813,13 @@ async def init_modules() -> HostRuntime: chain=lambda: build_chain_runtime_config(legacy_settings), ) runtime_settings = _build_runtime_settings_service() + workflow_query = WorkflowQueryService( + repository=TransactionalWorkflowQueryRepository( + sync_session=SessionFactory, + async_session=async_session_scope, + ) + ) + configure_workflow_query(workflow_query) agent_chat_persistence = AgentChatPersistenceService( repository=lambda session: AgentChatOper(session), async_executor=database_worker, @@ -852,6 +860,7 @@ async def init_modules() -> HostRuntime: outbox=SqlAlchemyAsyncOutboxStager, ), workflow=WorkflowRuntime( + query=workflow_query, repository=WorkflowOper, system_config=get_configured_system_config, ), @@ -866,16 +875,15 @@ async def init_modules() -> HostRuntime: configure_token_runtime_config(lambda: build_token_runtime_config(legacy_settings)) # 旧 app.api.data 导入只保留 ABI 转发,正式 API 依赖全部读取 HostRuntime。 configure_api_data_runtime(api_data) - configure_runtime_data_providers() + configure_runtime_data_providers(workflow_query) workflow_execution = TransactionalWorkflowExecutionService(SessionFactory) - configure_workflow_legacy_writer(workflow_execution) + configure_workflow_execution(workflow_execution) configure_chain_data_ports( site=lambda: TransactionalSiteRepository( sync_session=SessionFactory, async_session=async_session_scope, ), subscribe=lambda: SubscribeOper(), - workflow=lambda: WorkflowOper(), download_history=lambda: DownloadHistoryOper(), transfer_history=lambda: TransferHistoryOper(), transfer_pending=lambda: TransactionalTransferAdmissionRepository( @@ -921,7 +929,6 @@ async def init_modules() -> HostRuntime: sync_session=SessionFactory, async_session=async_session_scope, ))) - configure_workflow_query(WorkflowQueryService(repository=WorkflowOper())) configure_agent_data_ports( agent_chat=lambda: AgentChatOper(), agent_task=lambda: AgentTaskOper(), @@ -934,7 +941,6 @@ async def init_modules() -> HostRuntime: subscribe_history=lambda: SubscribeHistoryOper(), transfer_history=lambda: TransferHistoryOper(), download_history=lambda: DownloadHistoryOper(), - workflow=lambda: WorkflowOper(), plugin_data=lambda: PluginDataOper(), ) configure_agent_task_execution(AgentTaskExecutionService( diff --git a/app/startup/initializers/workflow.py b/app/startup/initializers/workflow.py index 291dd9712..028b8fc22 100644 --- a/app/startup/initializers/workflow.py +++ b/app/startup/initializers/workflow.py @@ -1,19 +1,19 @@ from app.application.workflow import configure_workflow_runtime -from app.workflow import WorkFlowManager +from app.workflow import WorkflowManager -# 启动模块是 concrete WorkFlowManager 的唯一宿主装配边界。 -configure_workflow_runtime(lambda: WorkFlowManager()) +# 启动模块是 concrete WorkflowManager 的唯一宿主装配边界。 +configure_workflow_runtime(lambda: WorkflowManager()) def init_workflow(): """ 初始化工作流 """ - WorkFlowManager() + WorkflowManager() def stop_workflow() -> bool: """ 停止工作流并返回全部活动执行是否收敛。 """ - return WorkFlowManager().stop() + return WorkflowManager().stop() diff --git a/app/workflow/__init__.py b/app/workflow/__init__.py index e2c694085..2d8d831e0 100644 --- a/app/workflow/__init__.py +++ b/app/workflow/__init__.py @@ -4,21 +4,23 @@ from typing import Any, Dict, List, Optional, Tuple from pydantic import BaseModel -from app.application.chain.data import get_chain_workflow_port -from app.application.workflow import WorkflowExecutionOwner +from app.application.workflow import ( + WorkflowExecutionOwner, + WorkflowSnapshot, + get_configured_workflow_query, +) from app.foundation.reflection import ModuleHelper from app.foundation.singleton import Singleton -from app.runtime.config import global_vars from app.runtime.events import Event, eventmanager from app.runtime.log import logger from app.runtime.stop import runtime_stop_state from app.schemas.types import EventType -from app.schemas.workflow import Action, ActionContext, ActionResult, Workflow +from app.schemas.workflow import Action, ActionContext, ActionResult _WORKFLOW_STOP_TIMEOUT_SECONDS = 10.0 -class WorkFlowManager(metaclass=Singleton): +class WorkflowManager(metaclass=Singleton): """ 工作流管理器 """ @@ -316,7 +318,7 @@ class WorkFlowManager(metaclass=Singleton): return {} return action.get_contract() - def update_workflow_event(self, workflow: Workflow): + def update_workflow_event(self, workflow: WorkflowSnapshot): """ 更新工作流事件触发器 """ @@ -333,11 +335,11 @@ class WorkFlowManager(metaclass=Singleton): """ workflows = [] if workflow_id: - workflow = get_chain_workflow_port().get(workflow_id) + workflow = get_configured_workflow_query().get_sync(workflow_id) if workflow: workflows = [workflow] else: - workflows = get_chain_workflow_port().get_event_triggered_workflows() + workflows = get_configured_workflow_query().list_event_enabled() try: for workflow in workflows: self.update_workflow_event(workflow) @@ -410,7 +412,7 @@ class WorkFlowManager(metaclass=Singleton): """ try: # 检查工作流是否存在且启用 - workflow = get_chain_workflow_port().get(workflow_id) + workflow = get_configured_workflow_query().get_sync(workflow_id) if not workflow or workflow.state == 'P': return diff --git a/docs/architecture-optimization-checklist.md b/docs/architecture-optimization-checklist.md index c860b174c..fe0bee642 100644 --- a/docs/architecture-optimization-checklist.md +++ b/docs/architecture-optimization-checklist.md @@ -69,7 +69,7 @@ MoviePilot V3 已经形成较清晰的模块化单体:`foundation`、`domain` | 指标 | 当前值 | 解释 | |---|---:|---| -| 宿主 Python 模块 / 内部依赖边 | 849 / 6,940 | `dependency-baseline.json` 当前快照 | +| 宿主 Python 模块 / 内部依赖边 | 850 / 6,944 | `dependency-baseline.json` 当前快照 | | 非平凡 SCC | 2 | 新增 Chain 包根环;另一个是隔离的 29 模块 TMDB 移植包环 | | 跨层 DB 边界债务 | 0 | Application、Chain、API、Agent、Runtime、Workflow 到 DB 的受控债务均为零 | | Model/Oper 事务债务 | 0 | 自建 Session、自动事务装饰器、直接 commit/rollback 等基线均为零 | @@ -77,9 +77,9 @@ MoviePilot V3 已经形成较清晰的模块化单体:`foundation`、`domain` | Event Contract | 53 | 均已有 payload model,但当前全部是 diagnostic enforcement | | Python 源码量 | 约 271,400 行 | 60 个文件超过 1,000 行,14 个超过 2,000 行 | | 长方法 | 281 个超过 80 行 | 67 个超过 150 行,23 个超过 250 行;大量是私有方法 | -| 全量 mypy 历史债务 | 11,820 / 596 文件 | strict frontier 当前覆盖 41 个文件,本批迁移路径的类型债务已清零 | -| Ruff 历史诊断 | 885 | 低水位门禁通过,但规则集只覆盖 `E4/E7/E9/F/I` | -| 覆盖率低水位 | Application 78.81%,Domain 79.29% | Chain、Runtime、Agent、Adapter、Startup 未进入包级覆盖率门禁 | +| 全量 mypy 历史债务 | 11,809 / 596 文件 | strict frontier 当前覆盖 41 个文件,本批迁移路径的类型债务已清零 | +| Ruff 历史诊断 | 875 | 低水位门禁通过,但规则集只覆盖 `E4/E7/E9/F/I` | +| 覆盖率低水位 | Application 78.89%,Domain 79.29% | Chain、Runtime、Agent、Adapter、Startup 未进入包级覆盖率门禁 | ### 3.3 热点文件 @@ -107,8 +107,8 @@ MoviePilot V3 已经形成较清晰的模块化单体:`foundation`、`domain` |---|---|---|---|---| | ARCH-001 | P0 | 已交付 | 恢复 mypy ratchet | `5df388719` 已推送,主线既有 CI gate 通过 | | ARCH-101 | P1 | 已交付 | 统一规则、总览、基线和语义门禁 | `113355784` 已推送,Unit Tests `33031697902`、Pylint `33031697785` 全绿,远端 `0/0` | -| ARCH-102 | P1 | 执行中 | 将 Transfer pending 升级为真实 E3 状态机 | `S1-L1.1` 至 `S1-L1.5` 全部交付后,崩溃窗口可判定恢复,结果未知时进入人工确认 | -| ARCH-103 | P1 | 待执行 | 类型化 Chain/Agent 数据 Port 与 DTO | 宿主主路径不再注入无 Session Oper,不向入口泄漏 ORM | +| ARCH-102 | P1 | 已交付 | 将 Transfer pending 升级为真实 E3 状态机 | `e9de149db`、`a2e249f20` 已推送;Unit Tests `33092427327`、Pylint `33092427348` 全绿,崩溃结果未知时进入人工确认 | +| ARCH-103 | P1 | 执行中 | 类型化 Chain/Agent 数据 Port 与 DTO | 宿主主路径不再注入无 Session Oper,不向入口泄漏 ORM | | ARCH-104 | P1 | 待执行 | 收口跨多次写入的业务事务 | 站点/规则引用清理可整体回滚或幂等恢复 | | ARCH-105 | P1 | 待执行 | 明确 post-commit 与 Outbox 完成语义 | “业务已提交、后置效果 pending”可被调用方正确识别 | | ARCH-106 | P1 | 待执行 | 让线程/队列/日志 writer 由 bootstrap/lifecycle 显式构造 | 导入或普通 Chain 构造不再启动进程资源 | @@ -196,8 +196,8 @@ MoviePilot V3 已经形成较清晰的模块化单体:`foundation`、`domain` claim/lease/heartbeat/attempt、过期接管、固定退避的唯一恢复入口和有界关闭 owner。 - `S1-L1.4 幂等执行与终态结算`:`VERIFIED`。已交付稳定 operation ledger、严格结果探测、 唯一 retry owner、`manual_review` 人工判定和 history/pending/outbox 同 UoW 终态结算。 -- `S1-L1.5 E3 全链收口`:`PLANNED`。完成崩溃矩阵、兼容验收与旧路径删除。此叶交付前, - ARCH-102 父项保持“执行中”,不得以局部绿色宣称 E3 完成。 +- `S1-L1.5 E3 全链收口`:`VERIFIED`。崩溃矩阵、3.0.17 升降级、重复回放、稳定计划身份、 + outcome/settlement 一致性和插件 ABI 已完成验收;旧 fail-open、重复状态与兼容层外旧入口已删除。 **问题与证据** @@ -227,7 +227,13 @@ MoviePilot V3 已经形成较清晰的模块化单体:`foundation`、`domain` - 管理员人工判定 API 只公开 `not_applied` 与带结果证据的 `applied`,并持久记录操作者、理由、结论 和 revision;无租约人工路径不能直接伪造失败终态。 - 以上实现已满足 `docs/adr/0007-background-action-reliability.md:123-139` 对 E3 稳定身份、步骤状态、 - lease/heartbeat 和人工恢复的阶段性要求;完整崩溃矩阵与兼容收口仍由 `S1-L1.5` 验收。 + lease/heartbeat 和人工恢复的要求。`RETRY_WAIT`、重放、双重失败、人工放弃和结算崩溃窗口均有 + 故障注入覆盖;计划指纹、步骤成员关系和所有状态写入使用精确 CAS,异常不再降级到旧执行路径。 +- canonical 模块已按职责聚合为 `app/application/chain/events.py`、`app/application/transfer/execution.py` + 和 `app/runtime/resources.py`;`durable_events.py`、`transfer_execution.py`、`managed_resources.py` + 等旧物理模块已退役,仅允许精确 Compat manifest 和兼容测试引用旧导入名,宿主不保留重复导出。 +- 交付提交为 `e9de149db`、`a2e249f20`;精确 head SHA 的 Unit Tests `33092427327` 与 Pylint + `33092427348` 全绿,覆盖率低水位同步提升至 Application `78.71%`。 **目标与步骤** @@ -266,16 +272,25 @@ MoviePilot V3 已经形成较清晰的模块化单体:`foundation`、`domain` `app/application/agentdata.py:91-120` 还通过 `__dict__.update()` 动态组装端口。 - `app/startup/initializers/modules.py:848-863,896-910` 仍向生产 Chain/Agent 注入多个无 Session Oper。 - 无 Session Oper 会为单次调用独立创建事务;一个业务操作的“查询后更新”可能被拆成多个事务。 -- Workflow query 和 Chain/Agent raw data port 仍返回 `Any`/ORM,Subscription mutation 内部也消费 ORM, - 因而存在 Session 生命周期外 detached/lazy-load 的潜在风险。公开 Subscription、Site、History - QueryService 已经投影 DTO,属于完成项,不应重做。 +- Workflow query 已在 S1-L2 迁入冻结 DTO 和 adapter-owned Session 投影;其余 Chain/Agent raw data port + 仍返回 `Any`/ORM,Subscription mutation 内部也消费 ORM,因而仍存在 Session 生命周期外 + detached/lazy-load 的潜在风险。公开 Subscription、Site、History QueryService 已经投影 DTO, + 属于完成项,不应重做。 +- S1-L2 由 `b4f873654`、`a01a35bcb` 交付;精确 head SHA 的 Unit Tests `33098869736` 与 + Pylint `33098869837` 全绿,Application 覆盖率低水位提升并固化至 `78.78%`。该证据只完成 + Workflow query 纵切面,不能替代 S1-L3 对其余 Chain/Agent raw data port 的清零。 **目标与步骤** - [ ] 按领域定义 Query/Command Protocol,不再使用通用 `OperFactory = Callable[[], Any]`。 - [ ] 写 Port 由 `db/adapters` 创建单操作 Session/UoW;Oper 的 canonical 写方法只 stage/flush, 读取方法仍可在调用方 Session 中查询。 -- [ ] 查询 Port 在 adapter Session 内映射为冻结 DTO/Projection,Application 和 API 不接收 ORM。 +- [x] Workflow 查询 Port 在 adapter Session 内映射为冻结 DTO/Projection,API、Agent、Chain、Scheduler、 + Workflow runtime 和中心服务分享均不接收 ORM。 +- [x] Workflow 执行写端由 Chain 直连 `WorkflowExecutionPort` 和短 Session/UoW 事务服务;canonical + `WorkflowOper` 只保留显式 Session query/stage,旧无 Session 五方法只存在于 SDK Legacy/Compat。 +- [x] 删除 Chain registry 中零消费者 `*PortProxy`/动态转发和 `ChainRuntimeContext.data_ports` + 伪注入;Workflow 执行服务只在 Application owner 配置一次,不再重复注册到 `ChainDataPorts`。 - [ ] `ChainDataPorts`/`AgentDataPorts` 可暂时保留为兼容聚合器,但字段必须显式、可类型检查。 - [ ] 以一个业务纵切面迁移并验证后,再迁移下一组,禁止一次替换所有 Oper。 - [ ] 增加 AST 门禁,禁止向 `ChainDataPorts`、`AgentDataPorts` 和新的 canonical use-case service diff --git a/docs/architecture-overview.md b/docs/architecture-overview.md index 1cbd4efc5..a366a3661 100644 --- a/docs/architecture-overview.md +++ b/docs/architecture-overview.md @@ -704,8 +704,8 @@ flowchart LR | 指标 | 当前值 | |---|---:| -| Python 模块 | 849 | -| 内部导入边 | 6,940 | +| Python 模块 | 850 | +| 内部导入边 | 6,944 | | 非平凡 SCC | 2(`ARCH-107` 临时 Chain 包根环;精确 containment 的 TMDB 移植包环) | | Direct egress | 66(12 条待迁移债务,54 条精确 containment) | | Module Contract V2 spec | 217(其中 215 个进入 `run_module` 观察面) | diff --git a/docs/architecture-refactor-roadmap.md b/docs/architecture-refactor-roadmap.md index eaa9254b3..db155a3d3 100644 --- a/docs/architecture-refactor-roadmap.md +++ b/docs/architecture-refactor-roadmap.md @@ -86,9 +86,9 @@ G-ARCH 只有在以下条件全部满足后才可完成: 退出条件:Transfer 达到 E3;正式数据 Port 类型化;跨表业务操作有单一 UoW;post-commit/Outbox 竞争、失败呈现、at-least-once 和幂等语义全部闭环。 -`S1-L1` 是 ARCH-102 的 Transfer E3 父项,当前状态为**执行中**。只有 `S1-L1.1` 至 -`S1-L1.5` 全部 `DELIVERED`,真实调用链完成迁移且旧 fail-open 路径退出 canonical 主程序后, -父项和 ARCH-102 才能标记已交付。 +`S1-L1` 是 ARCH-102 的 Transfer E3 父项,当前状态为 **DELIVERED**。`S1-L1.1` 至 +`S1-L1.5` 已全部交付,真实调用链已完成迁移,旧 fail-open、重复状态和兼容层外旧入口已退出 +canonical 主程序;兼容只经统一 Compat/SDK 门面提供。 | Leaf | 状态 | 依赖 | 完成定义 | |---|---|---|---| @@ -96,9 +96,18 @@ G-ARCH 只有在以下条件全部满足后才可完成: | S1-L1.2 Planning checkpoint | `VERIFIED` | S1-L1.1 | 版本化输入与指纹先持久化;无 legacy provider 时以 `accepted -> planned` CAS 提交完整计划,有 provider 时先提交 `provider_pending`,全部返回空后再以第二次 CAS 提交 `planned`;重放只执行冻结目标,所有文件副作用晚于对应 checkpoint commit | | S1-L1.3 Lease 与恢复调度 | `VERIFIED` | S1-L1.2 | claim/lease/heartbeat/attempt 与过期接管规则落地;启动回放和同进程恢复共用唯一调度入口,同一任务同时只有一个 worker owner | | S1-L1.4 幂等执行与终态结算 | `VERIFIED` | S1-L1.3 | 文件操作、历史提交和 checkpoint 可重放;唯一 retry owner 生效,未知外部结果进入 `manual_review`,仅完整终态删除 pending | -| S1-L1.5 E3 全链收口 | `PLANNED` | S1-L1.4 | 崩溃矩阵、升级/降级、重复回放和插件 ABI 验收完整;旧 fail-open、重复状态与兼容层外旧入口删除,ARCH-102 债务归零 | -| S1-L2 Workflow typed query | `PLANNED` | S0 | Workflow Application Port 不返回 `Any`/ORM,Session 内投影 DTO,正式调用方全部切换 | -| S1-L3 Chain/Agent typed data ports | `PLANNED` | S1-L2 | `ChainDataPorts`/`AgentDataPorts` 的 raw Oper/`Any` factory 全部清零,兼容调用进入 Legacy 层 | +| S1-L1.5 E3 全链收口 | `DELIVERED` | S1-L1.4 | `e9de149db`、`a2e249f20`:崩溃矩阵、3.0.17 升降级、重复回放和插件 ABI 验收完整;旧 fail-open、重复状态与兼容层外旧入口删除;Unit Tests `33092427327`、Pylint `33092427348` 全绿,ARCH-102 债务归零 | +| S1-L2 Workflow typed query | `DELIVERED` | S0 | `b4f873654`、`a01a35bcb`:Workflow Application Port 不返回 `Any`/ORM,Session 内投影冻结 DTO,正式调用方全部切换;Unit Tests `33098869736`、Pylint `33098869837` 全绿,覆盖率低水位提升至 Application `78.78%` | +| S1-L3 Chain/Agent typed data ports | `ACTIVE` | S1-L2 | `ChainDataPorts`/`AgentDataPorts` 的 raw Oper/`Any` factory 全部清零,兼容调用进入 Legacy 层 | +| S1-L3.1 Workflow typed execution | `DELIVERED` | S1-L2 | `17d8be2af`、`b33b29876`:Chain 直连类型化事务服务且单次执行只取一个 port;canonical Oper 删除旧 writer/无 Session 写方法,旧 ABI 只在 `_legacy/workflow.py` 与 Compat overlay;Unit Tests `33103913838`、Pylint `33103913935` 全绿,Application 覆盖率低水位提升至 `78.79%` | +| S1-L3.2 Chain registry/DI | `ACTIVE` | S1-L3.1 | 显式类型化 factory,删除 PortProxy 与失效的双重注入,构造器注入真实控制调用 | +| S1-L3.2.1 Registry hygiene | `DELIVERED` | S1-L3.1 | `ac7a20132`:删除零消费者 PortProxy/动态转发和 `ChainRuntimeContext.data_ports` 伪注入;Workflow 退出 Chain registry,只保留 Application owner 单一配置入口;Unit Tests `33120205586`、Pylint `33120205581` 全绿 | +| S1-L3.3 DownloadFailure/MediaServer | `PLANNED` | S1-L3.2 | 两组窄 DTO/Port/adapter 清零 raw Oper,不跨远端 I/O 持有 Session | +| S1-L3.4 User | `PLANNED` | S1-L3.3 | 认证、偏好与渠道绑定投影冻结快照,User Chain/Agent 不接收 ORM | +| S1-L3.5 History | `PLANNED` | S1-L3.4 | Download/Transfer history 统一 typed query/mutation,删除下载历史双事务 fail-open | +| S1-L3.6 Site | `PLANNED` | S1-L3.5 | 复用 Site query/health,补齐同步 typed command,Session 内完成 DTO 投影 | +| S1-L3.7 Subscription | `PLANNED` | S1-L3.6 | Chain/Workflow/interaction 全部消费 typed query/command;完成后进入 S1-L4 原子事务收口 | +| S1-L3.8 Agent/Transfer locator gate | `PLANNED` | S1-L3.7 | 删除 AgentDataPorts 与 Chain locator 跨层泄漏,AST 门禁确认 canonical 无 raw getter/Oper/Any | | S1-L4 Subscription mutation UoW | `PLANNED` | S1-L3 | Subscription mutation 不跨 Session 传 ORM,正式写路径一个 UoW,旧自动事务入口退出 canonical 路径 | | S1-L5 站点/规则引用原子清理 | `PLANNED` | S1-L4 | SystemConfig+Subscribe 同事务更新,commit 后快照原子发布,并发/故障注入无部分状态 | | S1-L6 Outbox 完成语义 | `PLANNED` | S0 | claim 竞争双发清零;业务提交与 effect pending 可区分;stager/store 分离;handler 幂等与崩溃测试完整 | @@ -144,7 +153,7 @@ G-ARCH 只有在以下条件全部满足后才可完成: | S4-L2 Event strict contract | `PLANNED` | S0-L2.6,S1-L6 | 宿主事件输入/输出按风险 strict,诊断例外只属于第三方插件兼容 | | S4-L3 Complexity v2 | `PLANNED` | S3 | 私有方法、class/file、圈复杂度进入门禁;所有超限通过职责拆分归零 | | S4-L4 全量 mypy 清零 | `PLANNED` | S3,S4-L1,S4-L2 | `mypy-baseline.json` 归零并删除债务接受路径,全宿主 strict 类型通过 | -| S4-L5 Ruff 治理债务清零 | `PLANNED` | S3 | 当前受控 885 条诊断归零,规则集扩展经过独立审查且新增诊断为零 | +| S4-L5 Ruff 治理债务清零 | `PLANNED` | S3 | 当前受控 875 条诊断归零,规则集扩展经过独立审查且新增诊断为零 | | S4-L6 Coverage/并发/质量证据 | `PLANNED` | S3,S4-L1,S4-L2 | 高风险包纳入 coverage;raw concurrency 分类清零;Module Quality 有真实 evidence test | ### S5:Plugin、Agent、Domain、Startup 与最终收口 @@ -263,7 +272,8 @@ git diff --check - 本叶不引入 claim、lease、heartbeat、attempt、执行步骤幂等或 `manual_review`;这些由 `S1-L1.3` 和 `S1-L1.4` 交付。 -- 文件操作成功后到历史结算前的未知结果仍未达到 E3,ARCH-102 父项继续保持执行中。 +- 本叶当时不单独承诺文件操作成功后到历史结算前的未知结果;该能力现已由 `S1-L1.4` 和 + `S1-L1.5` 的持久步骤账本、严格探测、`manual_review` 与 task-aware settlement 完整交付。 **Local verification (2026-08-27)** diff --git a/docs/rules/05-architecture.md b/docs/rules/05-architecture.md index 53cc82d30..ca0762deb 100644 --- a/docs/rules/05-architecture.md +++ b/docs/rules/05-architecture.md @@ -137,11 +137,12 @@ 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. -Migrated workflow, user, interaction, messaging, music, site, media-server, download, subscribe and transfer -Chain 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. +Migrated user, interaction, messaging, music, site, media-server, download, subscribe and transfer +Chain consumers temporarily use the named `get_chain_*_port()` functions from +`app/application/chain/data.py` while each owner establishes typed DTO/Port contracts. +The retired migration-time `*PortProxy` classes and dynamic `__getattr__` forwarding must not +be recreated; they had no host, SDK or plugin consumers. Workflow execution uses its owning +`app.application.workflow` service directly and must not be registered again in `ChainDataPorts`. Agent orchestration, memory and tool implementations follow the same rule via the named `get_agent_*_port()` functions from `app/application/agentdata.py`. The legacy Agent `*Port` proxy classes remain import-compatible boundaries and @@ -225,9 +226,19 @@ ModuleManager 与 startup 组合根继续关闭其余资源但必须向上返回 `app.runtime.execution.OwnedThreadPoolExecutor` 是进程级同步执行器有界收敛的唯一事实源;新的专用 线程池不得复制 Future 追踪、worker join 或重试关闭实现。DoH 查询线程池也必须复用该 owner:恢复系统 DNS 后有限等待,超时保留原 executor 并向 startup 返回 `False`,真实收敛前不得创建替代线程池或回填缓存。 -工作流节点线程池同样复用该 executor;所有 `WorkflowExecutor` 必须在 concrete `WorkFlowManager` 登记, +工作流节点线程池同样复用该 executor;所有 `WorkflowExecutor` 必须在 concrete `WorkflowManager` 登记, manager 停机先封口新执行并向活动 owner 发送本地取消,再有限等待执行线程和节点 worker。未收敛时必须 保留动作注册表和执行 owner,并让工作流生命周期 fail-fast,禁止继续释放仍被动作使用的插件或模块依赖。 +工作流读取统一使用 `app.application.workflow.WorkflowQueryService` 和冻结的 `WorkflowSnapshot`; +`app.db.adapters.workflow.TransactionalWorkflowQueryRepository` 必须在自有短 Session 内完成 ORM 投影与 +嵌套 JSON 深拷贝。API、Agent、Chain、Scheduler、`WorkflowManager` 和中心服务分享不得读取 raw +`WorkflowOper` 或把 ORM 带出 Session。旧 `WorkFlowManager` 拼写只由 Compat 符号覆盖承接,不进入 +canonical 模块定义或 `__all__`。 +工作流执行状态写入统一依赖 `app.application.workflow.WorkflowExecutionPort`;Chain 在一次执行中只获取 +一个事务端口,并由 `TransactionalWorkflowExecutionService` 为每次状态写入持有短 Session/UoW。 +canonical `app.db.oper.workflow.WorkflowOper` 只提供显式 Session 的 query/stage 方法;旧无 Session +`start/success/fail/step/reset` 仅由 `app.sdk._legacy.workflow` 和精确 Compat 映射承接,且不进入 +`app.db.oper.__all__`。 协程环境文件日志属于有界 E1 观测能力,只允许单一队列 writer;队列满时不得再以无界 executor 形成第二条异步写入路径。日志关闭必须有限等待 writer 与文件处理器,未收敛时 `LoggerManager` 保留原 owner 并让 lifespan 以关闭失败结束,不得先清空引用或用无界 `join()` 掩盖失败。 @@ -667,7 +678,8 @@ driven workflow registration. | `app/db/adapters/transfer/admission.py` | SQLAlchemy admission/checkpoint persistence, CAS state transition and detached snapshot adapter | | `app/application/scheduling.py` | Runtime scheduler facade for Agent tools and endpoints; `Scheduler` class registered by `app/startup/initializers/scheduler.py` | | `app/application/commands.py` | Command registry facade for Agent tools and endpoints; `Command` class registered by `app/startup/initializers/command.py` | -| `app/application/workflow.py` | Workflow use cases plus the runtime port consumed by API and Chain; `WorkFlowManager` is registered by `app/startup/initializers/workflow.py` | +| `app/application/workflow.py` | Workflow use cases, frozen query snapshot and typed runtime ports consumed by API, Agent, Chain and Scheduler; `WorkflowManager` is registered by `app/startup/initializers/workflow.py` | +| `app/db/adapters/workflow.py` | Short-session Workflow query projection and execution-state transaction adapters | | `app/db/adapters/` | SQLAlchemy repository/UoW implementations for Application-owned persistence Protocols | | `app/startup/composition/` | HostRuntime, configuration snapshots and cross-layer adapter wiring | | `app/startup/initializers/` | Domain-scoped initialization and shutdown hooks | diff --git a/tests/conftest.py b/tests/conftest.py index 81e3ea798..d7078c0cd 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -211,11 +211,12 @@ def configure_plugin_system_services(): from app.application.site.query import SiteQueryService, configure_site_query_service from app.application.workflow import ( WorkflowQueryService, + configure_workflow_execution, configure_workflow_query, configure_workflow_runtime, ) - from app.workflow import WorkFlowManager - configure_workflow_runtime(lambda: WorkFlowManager()) + from app.workflow import WorkflowManager + configure_workflow_runtime(lambda: WorkflowManager()) from app.application.agentdata import configure_agent_data_ports from app.application.agenttask import ( AgentTaskExecutionService, @@ -229,7 +230,10 @@ def configure_plugin_system_services(): from app.db.adapters.transfer.execution import ( TransactionalTransferExecutionRepository, ) - from app.db.adapters.workflow import TransactionalWorkflowExecutionService + from app.db.adapters.workflow import ( + TransactionalWorkflowExecutionService, + TransactionalWorkflowQueryRepository, + ) from app.db.oper.agentchat import AgentChatOper from app.db.oper.downloadhistory import DownloadHistoryOper from app.db.oper.mediaserver import MediaServerOper @@ -240,7 +244,7 @@ def configure_plugin_system_services(): from app.db.oper.subscribehistory import SubscribeHistoryOper from app.db.oper.transferhistory import TransferHistoryOper from app.db.oper.user import UserOper - from app.db.oper.workflow import WorkflowOper, configure_workflow_legacy_writer + from app.db.oper.workflow import WorkflowOper def create_sync_session() -> Session: """为无显式会话的 Oper 测试入口创建独占同步 Session。""" @@ -255,9 +259,8 @@ def configure_plugin_system_services(): async_=transaction_runner.async_, ) - configure_workflow_legacy_writer( - TransactionalWorkflowExecutionService(SessionFactory) - ) + workflow_execution = TransactionalWorkflowExecutionService(SessionFactory) + configure_workflow_execution(workflow_execution) configure_api_data_ports( sync_session=get_db, @@ -301,7 +304,6 @@ def configure_plugin_system_services(): configure_chain_data_ports( site=site_repository, subscribe=lambda: SubscribeOper(), - workflow=lambda: WorkflowOper(), download_history=lambda: DownloadHistoryOper(), transfer_history=lambda: TransferHistoryOper(), transfer_pending=lambda: TransactionalTransferAdmissionRepository( @@ -332,7 +334,12 @@ def configure_plugin_system_services(): )) configure_site_query_service(SiteQueryService(repository=site_repository())) configure_site_health_service(SiteHealthService(repository=site_repository())) - configure_workflow_query(WorkflowQueryService(repository=WorkflowOper())) + configure_workflow_query(WorkflowQueryService( + repository=TransactionalWorkflowQueryRepository( + sync_session=SessionFactory, + async_session=async_session_scope, + ) + )) from app.db.oper.agenttask import AgentTaskOper from app.db.oper.plugindata import PluginDataOper configure_agent_data_ports( @@ -344,7 +351,6 @@ def configure_plugin_system_services(): subscribe_history=lambda: SubscribeHistoryOper(), transfer_history=lambda: TransferHistoryOper(), download_history=lambda: DownloadHistoryOper(), - workflow=lambda: WorkflowOper(), plugin_data=lambda: PluginDataOper(), ) configure_agent_task_execution(AgentTaskExecutionService( diff --git a/tests/fixtures/architecture/coverage-baseline.json b/tests/fixtures/architecture/coverage-baseline.json index ee777a5a2..4114fab87 100644 --- a/tests/fixtures/architecture/coverage-baseline.json +++ b/tests/fixtures/architecture/coverage-baseline.json @@ -1,8 +1,8 @@ { "application": { - "covered_lines": 10068, - "percent": 78.81, - "statements": 12775 + "covered_lines": 10084, + "percent": 78.89, + "statements": 12782 }, "domain": { "covered_lines": 3392, diff --git a/tests/fixtures/architecture/dependency-baseline.json b/tests/fixtures/architecture/dependency-baseline.json index 5f2dcd3f2..d9e0c7111 100644 --- a/tests/fixtures/architecture/dependency-baseline.json +++ b/tests/fixtures/architecture/dependency-baseline.json @@ -1441,8 +1441,8 @@ "runtime_only": true } }, - "edge_count": 6940, - "edge_sha256": "9e3e8485c94c46a75eb577ccfc24708c80d4296b6422ec66ab9f23ef9964568f", + "edge_count": 6944, + "edge_sha256": "9f6be750c55150ef9e061f54ade997c328561b7950f6ec44bce6ada373d01f3e", "edges": [ "app -> app.runtime", "app -> app.runtime.compat", @@ -2587,7 +2587,7 @@ "app.agent.tools.impl.query_workflows -> app.agent.tools.base", "app.agent.tools.impl.query_workflows -> app.agent.tools.tags", "app.agent.tools.impl.query_workflows -> app.application", - "app.agent.tools.impl.query_workflows -> app.application.agentdata", + "app.agent.tools.impl.query_workflows -> app.application.workflow", "app.agent.tools.impl.query_workflows -> app.runtime", "app.agent.tools.impl.query_workflows -> app.runtime.log", "app.agent.tools.impl.read_file -> app.agent", @@ -3979,7 +3979,6 @@ "app.application.backup -> app.runtime.log", "app.application.chain.context -> app.application", "app.application.chain.context -> app.application.chain", - "app.application.chain.context -> app.application.chain.data", "app.application.chain.context -> app.application.chain.events", "app.application.chain.context -> app.application.configuration", "app.application.chain.context -> app.runtime", @@ -4344,6 +4343,8 @@ "app.application.servarr -> app.schemas.types", "app.application.server.report -> app.schemas", "app.application.server.report -> app.schemas.media", + "app.application.server.share -> app.application", + "app.application.server.share -> app.application.workflow", "app.application.server.share -> app.schemas", "app.application.server.share -> app.schemas.media", "app.application.service -> app.schemas", @@ -4476,10 +4477,11 @@ "app.application.transfer.workflow -> app.schemas.tmdb", "app.application.transfer.workflow -> app.schemas.transfer", "app.application.transfer.workflow -> app.schemas.types", + "app.application.workflow -> app.schemas", + "app.application.workflow -> app.schemas.common", "app.chain -> app.application", "app.chain -> app.application.chain", "app.chain -> app.application.chain.context", - "app.chain -> app.application.chain.data", "app.chain -> app.application.configuration", "app.chain -> app.chain._messaging", "app.chain -> app.chain._recognition", @@ -5064,8 +5066,6 @@ "app.chain.webhook -> app.schemas", "app.chain.webhook -> app.schemas.types", "app.chain.workflow -> app.application", - "app.chain.workflow -> app.application.chain", - "app.chain.workflow -> app.application.chain.data", "app.chain.workflow -> app.application.workflow", "app.chain.workflow -> app.chain", "app.chain.workflow -> app.runtime", @@ -5215,6 +5215,8 @@ "app.db.adapters.workflow -> app.db.oper", "app.db.adapters.workflow -> app.db.oper.workflow", "app.db.adapters.workflow -> app.db.uow", + "app.db.adapters.workflow -> app.schemas", + "app.db.adapters.workflow -> app.schemas.common", "app.db.base -> app.db", "app.db.base -> app.db.uow", "app.db.base -> app.runtime", @@ -7583,6 +7585,7 @@ "app.scheduler -> app.application.plugin.runtime", "app.scheduler -> app.application.scheduling", "app.scheduler -> app.application.site", + "app.scheduler -> app.application.workflow", "app.scheduler -> app.chain", "app.scheduler -> app.chain.mediaserver", "app.scheduler -> app.chain.recommend", @@ -7608,7 +7611,6 @@ "app.scheduler -> app.schemas.message", "app.scheduler -> app.schemas.system", "app.scheduler -> app.schemas.types", - "app.scheduler -> app.schemas.workflow", "app.schemas -> app.schemas.exports", "app.schemas.agent -> app.schemas", "app.schemas.agent -> app.schemas.common", @@ -7730,6 +7732,11 @@ "app.sdk._legacy.user -> app.db", "app.sdk._legacy.user -> app.db.oper", "app.sdk._legacy.user -> app.db.oper.user", + "app.sdk._legacy.workflow -> app.application", + "app.sdk._legacy.workflow -> app.application.workflow", + "app.sdk._legacy.workflow -> app.db", + "app.sdk._legacy.workflow -> app.db.oper", + "app.sdk._legacy.workflow -> app.db.oper.workflow", "app.sdk.browser -> app.adapters", "app.sdk.browser -> app.adapters.network", "app.sdk.browser -> app.adapters.network.browser", @@ -8209,14 +8216,11 @@ "app.testing.bootstrap -> app.startup.initializers.database", "app.testing.bootstrap -> app.startup.initializers.domain", "app.workflow -> app.application", - "app.workflow -> app.application.chain", - "app.workflow -> app.application.chain.data", "app.workflow -> app.application.workflow", "app.workflow -> app.foundation", "app.workflow -> app.foundation.reflection", "app.workflow -> app.foundation.singleton", "app.workflow -> app.runtime", - "app.workflow -> app.runtime.config", "app.workflow -> app.runtime.events", "app.workflow -> app.runtime.log", "app.workflow -> app.runtime.stop", @@ -8385,7 +8389,7 @@ "app.workflow.actions.transfer_file -> app.workflow", "app.workflow.actions.transfer_file -> app.workflow.actions" ], - "module_count": 849, + "module_count": 850, "modules": [ "app", "app.adapters", @@ -9179,6 +9183,7 @@ "app.sdk._legacy.transfer", "app.sdk._legacy.transferpending", "app.sdk._legacy.user", + "app.sdk._legacy.workflow", "app.sdk.browser", "app.sdk.cache", "app.sdk.config", diff --git a/tests/fixtures/architecture/mypy-baseline.json b/tests/fixtures/architecture/mypy-baseline.json index d90f7522b..fe44137ab 100644 --- a/tests/fixtures/architecture/mypy-baseline.json +++ b/tests/fixtures/architecture/mypy-baseline.json @@ -848,7 +848,7 @@ "arg-type": 7 }, "app/api/dependencies/workflow.py": { - "arg-type": 4, + "arg-type": 3, "redundant-cast": 2 }, "app/api/endpoints/agent.py": { @@ -1107,7 +1107,7 @@ "no-any-return": 4 }, "app/application/agentdata.py": { - "attr-defined": 10 + "attr-defined": 9 }, "app/application/agenttask.py": { "type-arg": 1 @@ -1299,7 +1299,6 @@ "type-arg": 5 }, "app/application/server/share.py": { - "no-any-return": 1, "type-arg": 6 }, "app/application/service.py": { @@ -1588,13 +1587,13 @@ "index": 1 }, "app/chain/workflow.py": { - "arg-type": 2, + "arg-type": 1, "assignment": 1, "call-arg": 1, "index": 4, "misc": 1, - "no-any-return": 10, - "no-untyped-def": 3, + "no-any-return": 7, + "no-untyped-def": 2, "truthy-function": 2, "type-arg": 17, "union-attr": 4, @@ -3232,7 +3231,7 @@ "unused-ignore": 1 }, "app/scheduler.py": { - "arg-type": 8, + "arg-type": 7, "attr-defined": 1, "import-untyped": 1, "misc": 1, @@ -3406,7 +3405,7 @@ "misc": 1, "no-any-return": 2, "no-untyped-call": 33, - "no-untyped-def": 13, + "no-untyped-def": 12, "return-value": 5 }, "app/startup/initializers/plugins.py": { @@ -3448,7 +3447,7 @@ "no-untyped-call": 2 }, "app/workflow/__init__.py": { - "arg-type": 2, + "arg-type": 1, "assignment": 2, "index": 1, "no-any-return": 1, diff --git a/tests/fixtures/architecture/ruff-baseline.json b/tests/fixtures/architecture/ruff-baseline.json index 438fbe92b..eb1c97429 100644 --- a/tests/fixtures/architecture/ruff-baseline.json +++ b/tests/fixtures/architecture/ruff-baseline.json @@ -291,9 +291,6 @@ "F401": 1, "I001": 1 }, - "app/application/agentdata.py": { - "I001": 1 - }, "app/application/agenttask.py": { "I001": 1 }, @@ -380,9 +377,6 @@ "app/application/torrent_cache.py": { "I001": 1 }, - "app/application/workflow.py": { - "I001": 1 - }, "app/chain/_music.py": { "E402": 5 }, @@ -401,9 +395,6 @@ "app/db/adapters/transaction.py": { "I001": 1 }, - "app/db/adapters/workflow.py": { - "I001": 1 - }, "app/db/base.py": { "I001": 1 }, @@ -482,9 +473,6 @@ "app/db/oper/userconfig.py": { "I001": 1 }, - "app/db/oper/workflow.py": { - "I001": 1 - }, "app/db/session.py": { "I001": 1 }, @@ -1050,9 +1038,6 @@ "E402": 27, "F401": 1 }, - "app/workflow/__init__.py": { - "F401": 1 - }, "scripts/architecture/task_ownership.py": { "I001": 1 }, @@ -1219,9 +1204,6 @@ "E402": 4, "I001": 1 }, - "tests/test_chain_runtime_context.py": { - "I001": 1 - }, "tests/test_cli_auto_update.py": { "I001": 1 }, @@ -1266,9 +1248,6 @@ "tests/test_db_lazy_engine.py": { "I001": 1 }, - "tests/test_db_oper_layer.py": { - "I001": 1 - }, "tests/test_db_oper_layer_extra.py": { "I001": 1 }, @@ -1353,10 +1332,6 @@ "tests/test_health_probes.py": { "I001": 1 }, - "tests/test_host_runtime_context.py": { - "F401": 1, - "I001": 1 - }, "tests/test_indexer_spider_search_url.py": { "I001": 1 }, @@ -1777,9 +1752,6 @@ "tests/test_workflow_authorization.py": { "I001": 1 }, - "tests/test_workflow_execution.py": { - "I001": 1 - }, "tests/test_workflow_runtime_config.py": { "I001": 1 } diff --git a/tests/fixtures/architecture/runtime-contract-baseline.json b/tests/fixtures/architecture/runtime-contract-baseline.json index fd884c692..a22e12b18 100644 --- a/tests/fixtures/architecture/runtime-contract-baseline.json +++ b/tests/fixtures/architecture/runtime-contract-baseline.json @@ -284,9 +284,9 @@ "app.db.workflow_oper": { "introduced": "v3.0.0", "is_package": false, - "owner": "db", - "replacement": "app.db.oper.workflow", - "target": "app.db.oper.workflow" + "owner": "sdk", + "replacement": "app.application.workflow.WorkflowExecutionPort", + "target": "app.sdk._legacy.workflow" }, "app.domain.string": { "introduced": "v3.0.0", @@ -926,6 +926,13 @@ "target_name": "MediaInteractionChain" } }, + "app.db.oper": { + "WorkflowOper": { + "replacement": "app.application.workflow.WorkflowExecutionPort", + "target_module": "app.sdk._legacy.workflow", + "target_name": "WorkflowOper" + } + }, "app.domain.media": { "MEDIA_SOURCE_ALIASES": { "replacement": "app.schemas.media.MEDIA_SOURCE_ALIASES", @@ -1184,6 +1191,13 @@ "target_module": "app.runtime.log", "target_name": "log_settings" } + }, + "app.workflow": { + "WorkFlowManager": { + "replacement": "app.workflow.WorkflowManager", + "target_module": "app.workflow", + "target_name": "WorkflowManager" + } } }, "virtual_packages": [ @@ -1440,12 +1454,12 @@ "caller": "app.workflow", "dynamic": true, "events": [], - "fingerprint": "042068d816db7e46ab4da6e96f8549af97b57bd9d75ba710cc1ff4fec7e5e188", + "fingerprint": "b99f557080a8ddb6dc9d2870d02c4cabc4b265b0daac086369c3bf1d51073c09", "handler": "self._handle_event", "invalid": false, "method": "add_event_listener", "priority": "", - "qualname": "WorkFlowManager.register_workflow_event", + "qualname": "WorkflowManager.register_workflow_event", "receiver_kind": "canonical_singleton", "registration_kind": "listener" } @@ -1836,7 +1850,7 @@ "50704edda70674af0932e769ddda40d2c21d50b134155c975d675035d7c933cd" ], "producer_fingerprints": [ - "611abaaf708f0c2d555d3b963ff540c015f2155ff493ed17b92d7c7f0bd45e96" + "b92d348c57dac3ba348079b4e71d02a0c13933bbf14b7e81188b042c3c5e2db3" ] } }, @@ -2980,10 +2994,10 @@ "events": [ "EventType.WorkflowExecute" ], - "fingerprint": "611abaaf708f0c2d555d3b963ff540c015f2155ff493ed17b92d7c7f0bd45e96", + "fingerprint": "b92d348c57dac3ba348079b4e71d02a0c13933bbf14b7e81188b042c3c5e2db3", "invalid": false, "method": "send_event", - "qualname": "WorkFlowManager._trigger_workflow", + "qualname": "WorkflowManager._trigger_workflow", "receiver_kind": "canonical_singleton" }, { diff --git a/tests/fixtures/architecture/runtime-contract-policy.json b/tests/fixtures/architecture/runtime-contract-policy.json index 4cd02d1ff..529463088 100644 --- a/tests/fixtures/architecture/runtime-contract-policy.json +++ b/tests/fixtures/architecture/runtime-contract-policy.json @@ -300,7 +300,7 @@ }, { "caller": "app.workflow", - "qualname": "WorkFlowManager.register_workflow_event", + "qualname": "WorkflowManager.register_workflow_event", "method": "add_event_listener", "receiver_kind": "canonical_singleton", "events": [], @@ -309,7 +309,7 @@ "handler": "self._handle_event", "registration_kind": "listener", "priority": "", - "fingerprint": "042068d816db7e46ab4da6e96f8549af97b57bd9d75ba710cc1ff4fec7e5e188", + "fingerprint": "b99f557080a8ddb6dc9d2870d02c4cabc4b265b0daac086369c3bf1d51073c09", "classification": "approved_dynamic_exception", "owner": "app.workflow", "reason": "工作流配置在运行期决定事件类型,receiver 与 handler 仍可静态证明。" diff --git a/tests/test_agent_data_ports.py b/tests/test_agent_data_ports.py index d698fa8bd..0be02b9b6 100644 --- a/tests/test_agent_data_ports.py +++ b/tests/test_agent_data_ports.py @@ -14,7 +14,6 @@ def test_named_agent_data_getters_use_registered_factories(monkeypatch) -> None: "subscribe_history": agentdata.get_agent_subscribe_history_port, "transfer_history": agentdata.get_agent_transfer_history_port, "download_history": agentdata.get_agent_download_history_port, - "workflow": agentdata.get_agent_workflow_port, "plugin_data": agentdata.get_agent_plugin_data_port, } factories = { diff --git a/tests/test_agent_query_workflows_tool.py b/tests/test_agent_query_workflows_tool.py index 17e372955..9ec539128 100644 --- a/tests/test_agent_query_workflows_tool.py +++ b/tests/test_agent_query_workflows_tool.py @@ -1,39 +1,49 @@ import asyncio import json -import unittest -from types import SimpleNamespace -from unittest.mock import AsyncMock, MagicMock, patch +from unittest.mock import AsyncMock, MagicMock from app.agent.tools.impl.query_workflows import QueryWorkflowsTool +from app.application.workflow import WorkflowSnapshot -class TestQueryWorkflowsTool(unittest.TestCase): - def test_query_workflows_omits_large_result_field(self): - tool = QueryWorkflowsTool(session_id="session-1", user_id="10001") - workflow = SimpleNamespace( - id=1, - name="demo", - description="demo workflow", - state="S", - trigger_type="manual", - run_count=1, - timer=None, - event_type=None, - add_time="2026-05-08 10:00:00", - last_time="2026-05-08 10:01:00", - current_action=None, - result="x" * 10000, - ) - workflow_oper = MagicMock() - workflow_oper.async_list = AsyncMock(return_value=[workflow]) +def _workflow() -> WorkflowSnapshot: + """构造 Agent 查询使用的真实工作流快照。""" + return WorkflowSnapshot( + id=1, + name="demo", + description="demo workflow", + timer=None, + trigger_type="manual", + event_type=None, + event_conditions={}, + state="S", + current_action=None, + result="x" * 10000, + run_count=1, + actions=(), + flows=(), + context={}, + execution_config={}, + execution_state={}, + add_time="2026-05-08 10:00:00", + last_time="2026-05-08 10:01:00", + ) - with patch( - "app.agent.tools.impl.query_workflows.get_agent_workflow_port", - return_value=workflow_oper, - ): - result = asyncio.run(tool.run()) - payload = json.loads(result) - self.assertEqual(len(payload), 1) - self.assertEqual(payload[0]["name"], "demo") - self.assertNotIn("result", payload[0]) +def test_query_workflows_omits_large_result_field(monkeypatch) -> None: + """Agent 列表查询使用统一快照服务且不返回大结果字段。""" + tool = QueryWorkflowsTool(session_id="session-1", user_id="10001") + query = MagicMock() + query.list = AsyncMock(return_value=[_workflow()]) + monkeypatch.setattr( + "app.agent.tools.impl.query_workflows.get_configured_workflow_query", + lambda: query, + ) + + result = asyncio.run(tool.run()) + + payload = json.loads(result) + assert len(payload) == 1 + assert payload[0]["name"] == "demo" + assert "result" not in payload[0] + query.list.assert_awaited_once_with() diff --git a/tests/test_architecture_dependencies.py b/tests/test_architecture_dependencies.py index 303d4825f..51b6ad743 100644 --- a/tests/test_architecture_dependencies.py +++ b/tests/test_architecture_dependencies.py @@ -283,6 +283,216 @@ def test_retired_canonical_filenames_do_not_return(): assert leftovers == [] +def test_workflow_query_contract_returns_only_typed_snapshots(): + """工作流正式查询端口不得退化为 Any 或 ORM 返回值。""" + path = APP_ROOT / "application" / "workflow.py" + tree = ast.parse(path.read_text(encoding="utf-8"), filename=str(path)) + query_classes = { + node.name: node + for node in tree.body + if isinstance(node, ast.ClassDef) + and node.name in {"WorkflowQueryRepository", "WorkflowQueryService"} + } + methods = [ + node + for query_class in query_classes.values() + for node in query_class.body + if isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef)) + and not node.name.startswith("__") + and node.returns is not None + ] + + assert set(query_classes) == {"WorkflowQueryRepository", "WorkflowQueryService"} + assert methods + for method in methods: + annotation = ast.unparse(method.returns) + assert "Any" not in annotation + assert "WorkflowSnapshot" in annotation + + +def test_workflow_query_consumers_do_not_reach_raw_oper(): + """API、Agent、共享服务和运行时管理器只消费统一快照查询服务。""" + consumer_paths = ( + "app/api/dependencies/workflow.py", + "app/agent/tools/impl/query_workflows.py", + "app/application/server/share.py", + "app/workflow/__init__.py", + ) + violations = {} + for relative_path in consumer_paths: + source = (PROJECT_ROOT / relative_path).read_text(encoding="utf-8") + forbidden = { + name + for name in ( + "WorkflowOper", + "get_agent_workflow_port", + "get_chain_workflow_port", + ) + if name in source + } + if forbidden: + violations[relative_path] = sorted(forbidden) + + assert violations == {} + + +def test_workflow_query_adapter_owns_projection_sessions(): + """唯一查询适配器必须在自有同步和异步 Session 内投影快照。""" + path = APP_ROOT / "db" / "adapters" / "workflow.py" + source = path.read_text(encoding="utf-8") + + assert "class TransactionalWorkflowQueryRepository" in source + assert "session.close()" in source + assert "async with self._async_session() as session" in source + assert "_project_workflow(record)" in source + + +def test_workflow_execution_chain_uses_single_application_owned_port(): + """工作流 Chain 写端必须只使用 Application owner 的唯一配置入口。""" + contract_path = APP_ROOT / "application" / "workflow.py" + contract_tree = ast.parse( + contract_path.read_text(encoding="utf-8"), + filename=str(contract_path), + ) + contract = next( + node + for node in contract_tree.body + if isinstance(node, ast.ClassDef) + and node.name == "WorkflowExecutionPort" + ) + methods = { + node.name: ast.unparse(node.returns) + for node in contract.body + if isinstance(node, ast.FunctionDef) + and node.returns is not None + } + assert methods == { + "start": "bool", + "success": "bool", + "fail": "bool", + "step": "bool", + "reset": "bool", + } + + data_path = APP_ROOT / "application" / "chain" / "data.py" + data_tree = ast.parse( + data_path.read_text(encoding="utf-8"), + filename=str(data_path), + ) + data_class = next( + node + for node in data_tree.body + if isinstance(node, ast.ClassDef) and node.name == "ChainDataPorts" + ) + data_fields = { + node.target.id + for node in data_class.body + if isinstance(node, ast.AnnAssign) + and isinstance(node.target, ast.Name) + } + data_functions = { + node.name + for node in data_tree.body + if isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef)) + } + assert "workflow" not in data_fields + assert "get_chain_workflow_port" not in data_functions + + chain_source = (APP_ROOT / "chain" / "workflow.py").read_text(encoding="utf-8") + startup_source = ( + APP_ROOT / "startup" / "initializers" / "modules.py" + ).read_text(encoding="utf-8") + assert chain_source.count("get_configured_workflow_execution()") == 1 + assert "get_chain_workflow_port" not in chain_source + assert "configure_workflow_execution(workflow_execution)" in startup_source + assert "workflow=lambda:" not in startup_source + + +def test_chain_registry_has_no_dynamic_proxies_or_dead_context_injection(): + """Chain registry 不得恢复零消费者动态代理或失效 data_ports 伪注入。""" + data_path = APP_ROOT / "application" / "chain" / "data.py" + data_tree = ast.parse( + data_path.read_text(encoding="utf-8"), + filename=str(data_path), + ) + proxy_classes = { + node.name + for node in data_tree.body + if isinstance(node, ast.ClassDef) + and (node.name.endswith("PortProxy") or node.name == "_PortProxyMeta") + } + dynamic_getters = { + node.name + for node in ast.walk(data_tree) + if isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef)) + and node.name == "__getattr__" + } + assert proxy_classes == set() + assert dynamic_getters == set() + + context_source = ( + APP_ROOT / "application" / "chain" / "context.py" + ).read_text(encoding="utf-8") + chain_base_source = (APP_ROOT / "chain" / "__init__.py").read_text( + encoding="utf-8" + ) + startup_source = ( + APP_ROOT / "startup" / "initializers" / "modules.py" + ).read_text(encoding="utf-8") + assert "data_ports" not in context_source + assert "self.data_ports" not in chain_base_source + assert "data_ports=" not in startup_source + + +def test_canonical_workflow_oper_has_no_legacy_writer_or_duplicate_exports(): + """工作流旧写入口只能存在于 SDK Legacy facade。""" + oper_path = APP_ROOT / "db" / "oper" / "workflow.py" + oper_tree = ast.parse( + oper_path.read_text(encoding="utf-8"), + filename=str(oper_path), + ) + oper_class = next( + node + for node in oper_tree.body + if isinstance(node, ast.ClassDef) and node.name == "WorkflowOper" + ) + method_names = { + node.name + for node in oper_class.body + if isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef)) + } + assert {"start", "success", "fail", "step", "reset"}.isdisjoint(method_names) + assert "legacy" not in oper_path.read_text(encoding="utf-8").lower() + + package_source = (APP_ROOT / "db" / "oper" / "__init__.py").read_text( + encoding="utf-8" + ) + assert '"WorkflowOper"' not in package_source + + +def test_agent_data_ports_do_not_duplicate_workflow_query_capability(): + """Agent 数据聚合器不得重新暴露无类型工作流读取入口。""" + source = (APP_ROOT / "application" / "agentdata.py").read_text( + encoding="utf-8" + ) + + assert "WorkflowPort" not in source + assert "get_agent_workflow_port" not in source + + +def test_host_uses_canonical_workflow_manager_name(): + """宿主代码不得继续定义或导入旧 WorkFlowManager 拼写。""" + violations = [] + for path in APP_ROOT.rglob("*.py"): + relative_path = path.relative_to(APP_ROOT) + if relative_path.parts[:2] == ("runtime", "compat"): + continue + if "WorkFlowManager" in path.read_text(encoding="utf-8"): + violations.append(str(path.relative_to(PROJECT_ROOT))) + + assert violations == [] + + def test_startup_root_contains_only_composition_packages(): """组合根顶层只保留稳定分区,禁止再次堆叠扁平实现文件。""" startup_root = APP_ROOT / "startup" diff --git a/tests/test_architecture_event_policy.py b/tests/test_architecture_event_policy.py index f25c924b5..78e22e2ce 100644 --- a/tests/test_architecture_event_policy.py +++ b/tests/test_architecture_event_policy.py @@ -137,7 +137,7 @@ def test_current_event_consumer_policy_matches_exact_reviewed_set() -> None: ] assert dynamic_entries == [{ "caller": "app.workflow", - "qualname": "WorkFlowManager.register_workflow_event", + "qualname": "WorkflowManager.register_workflow_event", "method": "add_event_listener", "receiver_kind": "canonical_singleton", "events": [], @@ -147,7 +147,7 @@ def test_current_event_consumer_policy_matches_exact_reviewed_set() -> None: "registration_kind": "listener", "priority": "", "fingerprint": ( - "042068d816db7e46ab4da6e96f8549af97b57bd9d75ba710cc1ff4fec7e5e188" + "b99f557080a8ddb6dc9d2870d02c4cabc4b265b0daac086369c3bf1d51073c09" ), "classification": DYNAMIC_CLASSIFICATION, "owner": "app.workflow", diff --git a/tests/test_chain_runtime_context.py b/tests/test_chain_runtime_context.py index 833d97d85..6d7cc410c 100644 --- a/tests/test_chain_runtime_context.py +++ b/tests/test_chain_runtime_context.py @@ -2,9 +2,12 @@ from unittest.mock import Mock +import pytest + +from app.application.chain import context as chain_context +from app.application.chain import data as chain_data from app.application.chain.context import ChainRuntimeContext from app.application.configuration import ChainRuntimeConfig -from app.application.chain import context as chain_context from app.chain import ChainBase from app.runtime.extensions.module.dispatcher import ModuleInvocationDispatcher @@ -51,3 +54,42 @@ def test_no_arg_chain_uses_compatibility_context_provider(monkeypatch) -> None: provider.assert_called_once_with() assert chain.modulemanager is context.module_manager assert chain.pluginmanager is context.plugin_manager + + +def test_chain_runtime_context_rejects_unconfigured_provider(monkeypatch) -> None: + """未由组合根配置运行上下文时必须显式拒绝无参 Chain。""" + monkeypatch.setattr( + chain_context, + "_context_provider", + chain_context._unconfigured_chain_runtime_context, + ) + + with pytest.raises(RuntimeError, match="Chain 运行上下文尚未由启动组合根配置"): + chain_context.get_chain_runtime_context() + + +def test_chain_data_registry_rejects_unconfigured_and_returns_factories( + monkeypatch, +) -> None: + """数据 registry 未配置时拒绝访问,配置后按字段返回工厂实例。""" + monkeypatch.setattr(chain_data, "_ports", None) + + with pytest.raises(RuntimeError, match="Chain 数据端口尚未配置"): + chain_data.get_chain_data_ports() + + media_server = Mock() + user = Mock() + chain_data.configure_chain_data_ports( + site=Mock, + subscribe=Mock, + download_history=Mock, + transfer_history=Mock, + transfer_pending=Mock, + transfer_execution=Mock, + media_server=lambda: media_server, + download_failure=Mock, + user=lambda: user, + ) + + assert chain_data.get_chain_media_server_port() is media_server + assert chain_data.get_chain_user_port() is user diff --git a/tests/test_db_oper_layer.py b/tests/test_db_oper_layer.py index 5010b4d1d..674019740 100644 --- a/tests/test_db_oper_layer.py +++ b/tests/test_db_oper_layer.py @@ -6,12 +6,11 @@ Oper 层大多是模型方法的薄封装,但薄封装恰恰是最容易出错 验证 Oper 的对外契约,而不是验证它调了哪个模型方法。 """ import asyncio +import importlib from unittest.mock import Mock import pytest -from app.db.oper.downloadhistory import DownloadHistoryOper -from app.db.oper.mediaserver import MediaServerOper from app.db.models.downloadhistory import DownloadFiles, DownloadHistory from app.db.models.mediaserver import MediaServerItem from app.db.models.plugindata import PluginData @@ -22,6 +21,8 @@ from app.db.models.siteuserdata import SiteUserData from app.db.models.user import User from app.db.models.userconfig import UserConfig from app.db.models.workflow import Workflow +from app.db.oper.downloadhistory import DownloadHistoryOper +from app.db.oper.mediaserver import MediaServerOper from app.db.oper.plugindata import PluginDataOper from app.db.oper.site import SiteOper from app.db.oper.user import UserOper @@ -374,7 +375,7 @@ def test_workflow_oper_add_rejects_duplicate_name(db): assert oper.add(**_workflow_kwargs("op-wf")) == (False, "工作流已存在") -def test_workflow_oper_exposes_lists_and_lifecycle(db): +def test_workflow_oper_exposes_lists_and_staged_lifecycle(db): """ 列表入口与生命周期方法都应透传到模型并落库。 """ @@ -387,23 +388,24 @@ def test_workflow_oper_exposes_lists_and_lifecycle(db): assert {w.name for w in oper.list_enabled()} >= {"op-wf-life"} assert {w.name for w in oper.get_timer_triggered_workflows()} >= {"op-wf-life"} - oper.start(flow.id) + oper.stage_start(flow.id) assert oper.get(flow.id).state == "R" - oper.step(flow.id, "a1", {"n": 1}) + oper.stage_step(flow.id, "a1", {"n": 1}) assert oper.get(flow.id).current_action == "a1" - oper.success(flow.id, "完成") + oper.stage_success(flow.id, "完成") assert oper.get(flow.id).state == "S" - oper.fail(flow.id, "出错") + oper.stage_fail(flow.id, "出错") assert oper.get(flow.id).state == "F" - oper.reset(flow.id, reset_count=True) + oper.stage_execution_reset(flow.id, reset_count=True) assert (oper.get(flow.id).state, oper.get(flow.id).run_count) == ("W", 0) def test_workflow_oper_no_session_uses_configured_uow_writer(db): """旧的无 Session Oper 写入口仍可用,但事务由组合根服务持有。""" flow = db.add(Workflow(**_workflow_kwargs("op-wf-legacy"))) + legacy = importlib.import_module("app.db.workflow_oper") - assert WorkflowOper().start(flow.id) is True + assert legacy.WorkflowOper().start(flow.id) is True db.session.expire_all() assert WorkflowOper(db=db.session).get(flow.id).state == "R" diff --git a/tests/test_db_workflow_queries.py b/tests/test_db_workflow_queries.py index 4b1164be2..37d1e6c16 100644 --- a/tests/test_db_workflow_queries.py +++ b/tests/test_db_workflow_queries.py @@ -6,13 +6,17 @@ `run_count` 的自增必须留在 SQL 侧,否则并发执行会丢计数。 """ import asyncio +from dataclasses import FrozenInstanceError import pytest +from app.application.workflow import WorkflowSnapshot from app.db import base as db_base +from app.db.adapters.workflow import TransactionalWorkflowQueryRepository from app.db.models.workflow import Workflow from app.db.oper.workflow import WorkflowOper -from app.db.session import async_session_scope +from app.db.session import SessionFactory, async_session_scope +from app.schemas.workflow import Workflow as WorkflowResponse @pytest.fixture(autouse=True) @@ -90,6 +94,63 @@ def test_workflow_oper_reuses_explicit_query_sessions(db, monkeypatch): asyncio.run(check()) +def test_query_repository_returns_detached_deep_copied_snapshot(db): + """查询仓储必须在关闭短 Session 前投影,且 JSON 不与 ORM 记录共享。""" + workflow = _flow("wf-snapshot") + workflow.actions = [{"id": "action-1", "config": {"value": 1}}] + workflow.flows = [{"source": "action-1", "target": "end"}] + workflow.context = {"nested": {"value": 1}} + created = db.add(workflow) + repository = TransactionalWorkflowQueryRepository( + sync_session=SessionFactory, + async_session=async_session_scope, + ) + + snapshot = repository.get(created.id) + + assert isinstance(snapshot, WorkflowSnapshot) + assert snapshot.name == "wf-snapshot" + with pytest.raises(FrozenInstanceError): + snapshot.name = "changed" + snapshot.actions[0]["config"]["value"] = 2 + snapshot.context["nested"]["value"] = 2 + refreshed = repository.get(created.id) + assert refreshed.actions[0]["config"]["value"] == 1 + assert refreshed.context["nested"]["value"] == 1 + + +def test_query_repository_async_projection_survives_session_close(db): + """异步查询返回值在仓储退出 Session 作用域后仍可完整序列化。""" + created = db.add(_flow("wf-async-snapshot")) + repository = TransactionalWorkflowQueryRepository( + sync_session=SessionFactory, + async_session=async_session_scope, + ) + + snapshot = asyncio.run(repository.async_get(created.id)) + listed = asyncio.run(repository.async_list()) + + assert isinstance(snapshot, WorkflowSnapshot) + assert snapshot.name == "wf-async-snapshot" + assert created.id in {item.id for item in listed} + + +def test_workflow_snapshot_validates_against_api_response_contract(db): + """冻结快照可直接序列化为 API 合同且不会暴露内部执行上下文。""" + created = db.add(_flow("wf-api-snapshot")) + repository = TransactionalWorkflowQueryRepository( + sync_session=SessionFactory, + async_session=async_session_scope, + ) + + response = WorkflowResponse.model_validate(repository.get(created.id)) + payload = response.model_dump() + + assert payload["id"] == created.id + assert payload["name"] == "wf-api-snapshot" + assert "context" not in payload + + def test_enabled_workflows_exclude_paused(db): """ 启用列表排除暂停状态。 diff --git a/tests/test_host_runtime_context.py b/tests/test_host_runtime_context.py index 14a9183ca..1706fccba 100644 --- a/tests/test_host_runtime_context.py +++ b/tests/test_host_runtime_context.py @@ -14,6 +14,13 @@ from app.api.context import ( get_agent_chat_transaction, ) from app.api.dependencies.agent import get_agent_chat_persistence +from app.application.configuration import ( + ApiRuntimeConfig, + ChainRuntimeConfig, + RuntimeConfiguration, + RuntimeSettingsService, + SchedulerRuntimeConfig, +) from app.startup import lifecycle from app.startup.composition.context import ( AgentChatRuntime, @@ -26,14 +33,6 @@ from app.startup.composition.context import ( SubscriptionRuntime, WorkflowRuntime, ) -from app.application.configuration import ( - ApiRuntimeConfig, - ChainRuntimeConfig, - RuntimeConfiguration, - RuntimeSettingsService, - SchedulerRuntimeConfig, -) - PROJECT_ROOT = Path(__file__).parents[1] @@ -153,6 +152,7 @@ def _runtime() -> HostRuntime: outbox=_Outbox, ), workflow=WorkflowRuntime( + query=SimpleNamespace(), repository=_Repository, system_config=lambda: _Repository(object()), ), diff --git a/tests/test_legacy_db_behavior_compat.py b/tests/test_legacy_db_behavior_compat.py index ccbedf474..a257d28af 100644 --- a/tests/test_legacy_db_behavior_compat.py +++ b/tests/test_legacy_db_behavior_compat.py @@ -7,6 +7,62 @@ from app.db.models.transferhistory import TransferHistory from app.schemas.file import FileItem +def test_legacy_workflow_writes_delegate_to_configured_execution_port(monkeypatch): + """旧 WorkflowOper 无 Session 写入必须完整委托类型化事务端口。""" + legacy = importlib.import_module("app.db.workflow_oper") + calls = [] + + class ExecutionPort: + """记录五种旧工作流写入调用。""" + + def start(self, workflow_id): + """记录启动。""" + calls.append(("start", workflow_id)) + return True + + def success(self, workflow_id, result=None): + """记录成功。""" + calls.append(("success", workflow_id, result)) + return True + + def fail(self, workflow_id, result): + """记录失败。""" + calls.append(("fail", workflow_id, result)) + return True + + def step(self, workflow_id, action_id, context, execution_state=None): + """记录步骤。""" + calls.append( + ("step", workflow_id, action_id, context, execution_state) + ) + return True + + def reset(self, workflow_id, reset_count=False): + """记录重置。""" + calls.append(("reset", workflow_id, reset_count)) + return True + + monkeypatch.setattr( + legacy, + "get_configured_workflow_execution", + lambda: ExecutionPort(), + ) + oper = legacy.WorkflowOper() + + assert oper.start(7) is True + assert oper.success(7, "done") is True + assert oper.fail(7, "failed") is True + assert oper.step(7, "A", {"value": 1}, {"runtime": {}}) is True + assert oper.reset(7, reset_count=True) is True + assert calls == [ + ("start", 7), + ("success", 7, "done"), + ("fail", 7, "failed"), + ("step", 7, "A", {"value": 1}, {"runtime": {}}), + ("reset", 7, True), + ] + + def test_legacy_subscribe_add_delegates_to_application_service(monkeypatch): """旧 SubscribeOper.add 应保留 mediainfo 写入签名。""" legacy = importlib.import_module("app.db.subscribe_oper") diff --git a/tests/test_legacy_import_compat.py b/tests/test_legacy_import_compat.py index 999014e14..8f588a47e 100644 --- a/tests/test_legacy_import_compat.py +++ b/tests/test_legacy_import_compat.py @@ -314,6 +314,19 @@ def test_db_refactor_legacy_modules_are_all_registered(): assert expected <= set(MODULE_ALIASES) +def test_workflow_oper_compatibility_is_only_exposed_by_overlay(): + """旧工作流写入口只由 Legacy facade 和精确符号映射提供。""" + legacy = importlib.import_module("app.db.workflow_oper") + canonical = importlib.import_module("app.db.oper.workflow") + oper_package = importlib.import_module("app.db.oper") + + assert MODULE_ALIASES["app.db.workflow_oper"].target == "app.sdk._legacy.workflow" + assert issubclass(legacy.WorkflowOper, canonical.WorkflowOper) + assert legacy.WorkflowOper is not canonical.WorkflowOper + assert oper_package.WorkflowOper is legacy.WorkflowOper + assert "WorkflowOper" not in oper_package.__all__ + + def test_split_user_oper_facade_exports_data_and_auth_contracts(): """旧 user_oper 同时提供 UserOper 与八个认证依赖。""" legacy = importlib.import_module("app.db.user_oper") @@ -410,6 +423,16 @@ def test_chain_media_legacy_scraping_symbols_resolve_to_scraping_chain(): assert legacy_media.ScrapingConfig is canonical_scraping.ScrapingConfig +def test_workflow_manager_legacy_name_resolves_only_through_symbol_overlay(): + """旧 WorkFlowManager 仍可显式导入,但不进入 canonical 模块公开面。""" + install_legacy_import_hook() + workflow_module = importlib.import_module("app.workflow") + + assert workflow_module.WorkFlowManager is workflow_module.WorkflowManager + assert "WorkFlowManager" not in vars(workflow_module) + assert "WorkFlowManager" not in workflow_module.__all__ + + def test_rules_domain_legacy_modules_resolve_to_rules(): """规则域收敛后,filter/filter_rules 旧路径应复用 rules 模块。""" canonical = importlib.import_module("app.application.rules") @@ -459,6 +482,8 @@ def test_plugin_scan_reports_moved_symbol_import(tmp_path: Path): def test_symbol_alias_manifest_covers_all_moved_public_symbols(): """符号级映射清单应覆盖媒体身份、整理工作项、刮削拆分与消息/通知命名统一的旧入口。""" + assert set(SYMBOL_ALIASES["app.db.oper"]) == {"WorkflowOper"} + assert set(SYMBOL_ALIASES["app.workflow"]) == {"WorkFlowManager"} assert set(SYMBOL_ALIASES["app.domain.media"]) == { "MEDIA_SOURCE_ALIASES", "MEDIA_SOURCE_PREFIXES", diff --git a/tests/test_server_sharing_service.py b/tests/test_server_sharing_service.py index e4a598442..51d503c32 100644 --- a/tests/test_server_sharing_service.py +++ b/tests/test_server_sharing_service.py @@ -1,8 +1,34 @@ import asyncio +import json from types import SimpleNamespace from unittest.mock import AsyncMock, Mock from app.application.server.share import ServerSharingService +from app.application.workflow import WorkflowSnapshot + + +def _workflow(*, actions=(), flows=()) -> WorkflowSnapshot: + """构造中心服务分享使用的真实工作流快照。""" + return WorkflowSnapshot( + id=1, + name="Demo Workflow", + description="demo", + timer=None, + trigger_type="manual", + event_type=None, + event_conditions={}, + state="W", + current_action=None, + result=None, + run_count=0, + actions=actions, + flows=flows, + context={"private": True}, + execution_config={}, + execution_state={}, + add_time=None, + last_time=None, + ) def _service(**overrides) -> ServerSharingService: @@ -66,7 +92,7 @@ def test_subscribe_share_builds_public_payload_and_clears_cache_after_success(): def test_workflow_validation_stops_before_transport(): """缺少动作或流程的工作流不会进入中心服务传输。""" sender = Mock() - workflow = SimpleNamespace(actions=[], flows=[{"id": 1}]) + workflow = _workflow(flows=({"id": 1},)) service = _service( workflow_provider=Mock(return_value=workflow), workflow_sender=sender, @@ -84,6 +110,36 @@ def test_workflow_validation_stops_before_transport(): sender.assert_not_called() +def test_workflow_share_serializes_snapshot_without_local_fields(): + """同步工作流分享从冻结快照生成兼容载荷并剔除本地上下文。""" + sender = Mock(return_value=SimpleNamespace(status_code=200)) + workflow = _workflow( + actions=({"id": "action-1"},), + flows=({"source": "action-1", "target": "end"},), + ) + service = _service( + workflow_provider=Mock(return_value=workflow), + workflow_sender=sender, + ) + + result = service.share_workflow( + enabled=True, + workflow_id=1, + share_title="Title", + share_comment="Comment", + share_user="User", + ) + + assert result == (True, "") + payload = sender.call_args.args[0] + assert "id" not in payload + assert "context" not in payload + assert json.loads(payload["actions"]) == [{"id": "action-1"}] + assert json.loads(payload["flows"]) == [ + {"source": "action-1", "target": "end"} + ] + + def test_async_subscribe_share_uses_async_reader_and_transport(): """异步分享路径不会回退到同步数据库或网络端口。""" subscribe = SimpleNamespace(to_dict=lambda: { @@ -110,3 +166,29 @@ def test_async_subscribe_share_uses_async_reader_and_transport(): assert result == (True, "") reader.assert_awaited_once_with(1) sender.assert_awaited_once() + + +def test_async_workflow_share_uses_snapshot_reader_and_transport(): + """异步工作流分享复用同一快照契约且不回退同步端口。""" + workflow = _workflow( + actions=({"id": "action-1"},), + flows=({"source": "action-1", "target": "end"},), + ) + reader = AsyncMock(return_value=workflow) + sender = AsyncMock(return_value=SimpleNamespace(status_code=200)) + service = _service( + async_workflow_provider=reader, + async_workflow_sender=sender, + ) + + result = asyncio.run(service.async_share_workflow( + enabled=True, + workflow_id=1, + share_title="Title", + share_comment="Comment", + share_user="User", + )) + + assert result == (True, "") + reader.assert_awaited_once_with(1) + sender.assert_awaited_once() diff --git a/tests/test_workflow_actions.py b/tests/test_workflow_actions.py index 31c8029f8..fa648db98 100644 --- a/tests/test_workflow_actions.py +++ b/tests/test_workflow_actions.py @@ -3,7 +3,7 @@ from types import SimpleNamespace from app.schemas.download import DownloadTask from app.schemas.file import FileItem from app.schemas.workflow import ActionContext, ActionResult -from app.workflow import WorkFlowManager +from app.workflow import WorkflowManager from app.workflow.actions import BaseAction from app.workflow.actions import fetch_downloads as fetch_downloads_module from app.workflow.actions import fetch_torrents as fetch_torrents_module @@ -239,7 +239,7 @@ def test_execute_with_inputs_maps_contract_inputs_outputs_and_runtime(monkeypatc def test_workflow_manager_list_actions_exposes_contract(): """动作列表应返回固定输入输出契约。""" - manager = object.__new__(WorkFlowManager) + manager = object.__new__(WorkflowManager) manager._actions = {"FetchRssAction": FetchRssAction} actions = manager.list_actions() diff --git a/tests/test_workflow_execution.py b/tests/test_workflow_execution.py index 53ebf48aa..0cfd69200 100644 --- a/tests/test_workflow_execution.py +++ b/tests/test_workflow_execution.py @@ -6,11 +6,11 @@ from types import SimpleNamespace import pytest +from app import workflow as workflow_package from app.chain import workflow as workflow_module from app.runtime.correlation import correlation_scope, get_correlation_id from app.schemas.types import EventType from app.schemas.workflow import Action, ActionContext, ActionResult -from app import workflow as workflow_package def _build_workflow(current_action=None, context=None, actions=None, flows=None, @@ -635,9 +635,24 @@ def test_workflow_chain_process_serializes_circular_context(monkeypatch): flows=[{"id": "flow-end", "source": "A", "target": "END", "animated": True}], ) fake_oper = _FakeWorkflowOper(workflow) + port_calls = [] + + def get_execution_port(): + """记录单次执行获取事务端口的次数。""" + port_calls.append(True) + return fake_oper monkeypatch.setattr(workflow_module, "get_workflow_manager", lambda: fake_manager) - monkeypatch.setattr(workflow_module, "get_chain_workflow_port", lambda: fake_oper) + monkeypatch.setattr( + workflow_module, + "get_configured_workflow_execution", + get_execution_port, + ) + monkeypatch.setattr( + workflow_module, + "get_configured_workflow_query", + lambda: SimpleNamespace(get_sync=lambda _workflow_id: workflow), + ) monkeypatch.setattr(workflow_module.runtime_stop_state, "resume_workflow", lambda workflow_id: None) monkeypatch.setattr(workflow_module.runtime_stop_state, "is_workflow_stopped", lambda workflow_id: False) @@ -645,6 +660,7 @@ def test_workflow_chain_process_serializes_circular_context(monkeypatch): assert success is True assert message == "" + assert port_calls == [True] assert fake_oper.succeeded is True saved_workflow_context = fake_oper.steps[-1]["context"]["workflow_context"] saved_self = saved_workflow_context["self"] @@ -825,7 +841,7 @@ def test_workflow_manager_shutdown_retains_blocked_execution_for_retry(monkeypat release.wait() return ActionResult(success=True, context=context) - manager = object.__new__(workflow_package.WorkFlowManager) + manager = object.__new__(workflow_package.WorkflowManager) manager._lock = threading.RLock() manager._actions = {"BlockingAction": BlockingAction} manager._event_workflows = {} @@ -904,7 +920,7 @@ def test_workflow_manager_shutdown_continues_across_owner_failures(): self.manager.unregister_execution(self) return True - manager = object.__new__(workflow_package.WorkFlowManager) + manager = object.__new__(workflow_package.WorkflowManager) manager._lock = threading.RLock() action_marker = object() manager._actions = {"FakeAction": action_marker} @@ -940,7 +956,16 @@ def test_workflow_chain_rejects_execution_before_persisting_running_state(monkey workflowoper = _FakeWorkflowOper(workflow) manager = RejectingWorkflowManager([]) monkeypatch.setattr(workflow_module, "get_workflow_manager", lambda: manager) - monkeypatch.setattr(workflow_module, "get_chain_workflow_port", lambda: workflowoper) + monkeypatch.setattr( + workflow_module, + "get_configured_workflow_execution", + lambda: workflowoper, + ) + monkeypatch.setattr( + workflow_module, + "get_configured_workflow_query", + lambda: SimpleNamespace(get_sync=lambda _workflow_id: workflow), + ) def unexpected_resume(_workflow_id: int) -> None: """拒绝准入时若仍恢复停止标记则立即暴露回归。""" @@ -970,7 +995,7 @@ def test_workflow_chain_releases_admitted_owner_when_start_fails(monkeypatch): _ = wid raise RuntimeError("start failed") - manager = object.__new__(workflow_package.WorkFlowManager) + manager = object.__new__(workflow_package.WorkflowManager) manager._lock = threading.RLock() manager._actions = {"FakeAction": object()} manager._event_workflows = {} @@ -978,7 +1003,18 @@ def test_workflow_chain_releases_admitted_owner_when_start_fails(monkeypatch): manager._executions = {} workflowoper = FailingWorkflowOper(_build_workflow()) monkeypatch.setattr(workflow_module, "get_workflow_manager", lambda: manager) - monkeypatch.setattr(workflow_module, "get_chain_workflow_port", lambda: workflowoper) + monkeypatch.setattr( + workflow_module, + "get_configured_workflow_execution", + lambda: workflowoper, + ) + monkeypatch.setattr( + workflow_module, + "get_configured_workflow_query", + lambda: SimpleNamespace( + get_sync=lambda _workflow_id: workflowoper.workflow + ), + ) monkeypatch.setattr(workflow_module.runtime_stop_state, "resume_workflow", lambda _workflow_id: None) with pytest.raises(RuntimeError, match="start failed"): @@ -1018,7 +1054,7 @@ class _FakeEventManager: def test_workflow_event_listener_keeps_shared_handler_until_last_workflow(monkeypatch): """同一事件下移除单个工作流时不应断开其他工作流监听。""" fake_eventmanager = _FakeEventManager() - manager = object.__new__(workflow_package.WorkFlowManager) + manager = object.__new__(workflow_package.WorkflowManager) manager._lock = threading.Lock() manager._event_workflows = {} @@ -1057,7 +1093,7 @@ def test_workflow_manager_retries_action_until_success(monkeypatch): return ActionResult(success=False, message="第一次失败", context=context) return ActionResult(success=True, message="第二次成功", context=context, outputs={"ok": True}) - manager = object.__new__(workflow_package.WorkFlowManager) + manager = object.__new__(workflow_package.WorkflowManager) manager._actions = {"RetryAction": RetryAction} monkeypatch.setattr(workflow_package.runtime_stop_state, "is_workflow_stopped", lambda workflow_id: False) diff --git a/tests/test_workflow_mutation_command.py b/tests/test_workflow_mutation_command.py index 0f587cb9a..2d97cfa41 100644 --- a/tests/test_workflow_mutation_command.py +++ b/tests/test_workflow_mutation_command.py @@ -3,11 +3,13 @@ from unittest.mock import AsyncMock, Mock import pytest +import app.application.workflow as workflow_application from app.application.workflow import ( WorkflowDefinitionCommand, WorkflowExecutionCommand, WorkflowMutationCommand, WorkflowQueryService, + WorkflowSnapshot, ) @@ -21,6 +23,30 @@ def _workflow(trigger_type="timer", timer="0 0 * * *", event_type="DownloadAdded ) +def _snapshot() -> WorkflowSnapshot: + """构造查询服务返回的冻结工作流快照。""" + return WorkflowSnapshot( + id=7, + name="query", + description=None, + timer="0 0 * * *", + trigger_type="timer", + event_type=None, + event_conditions={}, + state="W", + current_action=None, + result=None, + run_count=0, + actions=(), + flows=(), + context={}, + execution_config={}, + execution_state={}, + add_time=None, + last_time=None, + ) + + def _command(workflow=None, commit_error=None): """构造可观察工作流事务与运行时副作用的命令。""" repository = Mock() @@ -62,6 +88,23 @@ def _execution_command(commit_error=None): ), repository, unit_of_work +def test_workflow_execution_port_requires_explicit_configuration(monkeypatch): + """执行状态端口必须显式装配,并原样返回组合根登记的服务。""" + monkeypatch.setattr( + workflow_application, + "_configured_workflow_execution", + None, + ) + + with pytest.raises(RuntimeError, match="工作流执行状态事务服务尚未配置"): + workflow_application.get_configured_workflow_execution() + + service = Mock() + workflow_application.configure_workflow_execution(service) + + assert workflow_application.get_configured_workflow_execution() is service + + def test_execution_step_is_staged_before_unit_of_work_commit(): """工作流进度写入必须由应用命令暂存后统一提交。""" command, repository, unit_of_work = _execution_command() @@ -96,8 +139,9 @@ def test_execution_commit_failure_rolls_back(): async def test_workflow_query_service_delegates_list_and_get_to_repository(): """工作流查询服务只调用读取端口,不持有数据库会话或事务。""" repository = Mock() - repository.async_list = AsyncMock(return_value=[_workflow()]) - repository.async_get = AsyncMock(return_value=_workflow()) + snapshot = _snapshot() + repository.async_list = AsyncMock(return_value=[snapshot]) + repository.async_get = AsyncMock(return_value=snapshot) service = WorkflowQueryService(repository) listed = await service.list() @@ -105,6 +149,8 @@ async def test_workflow_query_service_delegates_list_and_get_to_repository(): assert listed == repository.async_list.return_value assert fetched == repository.async_get.return_value + assert all(isinstance(item, WorkflowSnapshot) for item in listed) + assert isinstance(fetched, WorkflowSnapshot) repository.async_list.assert_awaited_once_with() repository.async_get.assert_awaited_once_with(7) @@ -123,6 +169,18 @@ def test_start_timer_workflow_commits_before_registering_job(): dependencies["repository"].stage_state.assert_called_once_with(7, "W") +def test_start_rejects_missing_workflow_without_transaction(): + """工作流不存在时不得暂存状态或触发事务。""" + command, dependencies = _command() + + result = command.start(7) + + assert result.success is False + assert result.message == "工作流不存在" + dependencies["repository"].stage_state.assert_not_called() + dependencies["unit_of_work"].commit.assert_not_called() + + def test_start_rejects_invalid_trigger_without_transaction(): """未知触发类型不得更新数据库或注册运行时触发器。""" command, dependencies = _command(_workflow(trigger_type="unknown")) diff --git a/tests/test_workflow_runtime_facade.py b/tests/test_workflow_runtime_facade.py index a22ec8875..c45da9d41 100644 --- a/tests/test_workflow_runtime_facade.py +++ b/tests/test_workflow_runtime_facade.py @@ -18,7 +18,7 @@ def test_workflow_runtime_facade_preserves_registered_identity(monkeypatch) -> N def test_workflow_runtime_facade_fails_before_composition(monkeypatch) -> None: - """未装配时不得隐式创建第二个 WorkFlowManager Singleton。""" + """未装配时不得隐式创建第二个 WorkflowManager Singleton。""" monkeypatch.setattr( workflow_application, "_workflow_runtime_provider",