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/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..0e2aea90f 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 diff --git a/app/chain/workflow.py b/app/chain/workflow.py index 2d40b4d50..dd4193ffb 100644 --- a/app/chain/workflow.py +++ b/app/chain/workflow.py @@ -14,7 +14,11 @@ 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_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: """ @@ -1277,10 +1288,29 @@ class WorkflowChain(ChainBase): """ workflowoper = get_chain_workflow_port() - def save_step(action: Action, context: ActionContext, execution_state: dict, completed: bool): - """ - 保存上下文到数据库 - """ + # 重置工作流 + if from_begin: + workflowoper.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: + """保存动作上下文和结构化执行状态。""" get_chain_workflow_port().step( workflow_id, action_id=action.id if completed else "", @@ -1305,22 +1335,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( @@ -1355,22 +1369,22 @@ class WorkflowChain(ChainBase): 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/runtime/compat/manifest.py b/app/runtime/compat/manifest.py index 2d05c3d79..01a4791cd 100644 --- a/app/runtime/compat/manifest.py +++ b/app/runtime/compat/manifest.py @@ -774,6 +774,13 @@ _MESSAGE_NOTIFICATION_SYMBOL_ALIASES: Dict[str, SymbolAlias] = { } SYMBOL_ALIASES: Dict[str, Dict[str, SymbolAlias]] = { + "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/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 9ebe1256b..bf5d74509 100644 --- a/app/startup/initializers/modules.py +++ b/app/startup/initializers/modules.py @@ -114,7 +114,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 @@ -239,11 +242,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 @@ -273,7 +271,7 @@ def _build_chain_runtime_context() -> ChainRuntimeContext: ) -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()) @@ -309,8 +307,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, @@ -799,6 +797,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, @@ -839,6 +844,7 @@ async def init_modules() -> HostRuntime: outbox=SqlAlchemyAsyncOutboxStager, ), workflow=WorkflowRuntime( + query=workflow_query, repository=WorkflowOper, system_config=get_configured_system_config, ), @@ -853,7 +859,7 @@ 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_chain_data_ports( @@ -908,7 +914,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(), @@ -921,7 +926,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 cc649e83f..2d4d76008 100644 --- a/docs/architecture-optimization-checklist.md +++ b/docs/architecture-optimization-checklist.md @@ -69,7 +69,7 @@ MoviePilot V3 已经形成较清晰的模块化单体:`foundation`、`domain` | 指标 | 当前值 | 解释 | |---|---:|---| -| 宿主 Python 模块 / 内部依赖边 | 844 / 6,898 | `dependency-baseline.json` 当前快照 | +| 宿主 Python 模块 / 内部依赖边 | 844 / 6,901 | `dependency-baseline.json` 当前快照 | | 非平凡 SCC | 2 | 新增 Chain 包根环;另一个是隔离的 29 模块 TMDB 移植包环 | | 跨层 DB 边界债务 | 0 | Application、Chain、API、Agent、Runtime、Workflow 到 DB 的受控债务均为零 | | Model/Oper 事务债务 | 0 | 自建 Session、自动事务装饰器、直接 commit/rollback 等基线均为零 | @@ -77,8 +77,8 @@ 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 历史诊断 | 888 | 低水位门禁通过,但规则集只覆盖 `E4/E7/E9/F/I` | +| 全量 mypy 历史债务 | 11,809 / 596 文件 | strict frontier 当前覆盖 41 个文件,本批迁移路径的类型债务已清零 | +| Ruff 历史诊断 | 881 | 低水位门禁通过,但规则集只覆盖 `E4/E7/E9/F/I` | | 覆盖率低水位 | Application 78.71%,Domain 79.29% | Chain、Runtime、Agent、Adapter、Startup 未进入包级覆盖率门禁 | ### 3.3 热点文件 @@ -272,16 +272,18 @@ 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, + 属于完成项,不应重做。 **目标与步骤** - [ ] 按领域定义 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。 - [ ] `ChainDataPorts`/`AgentDataPorts` 可暂时保留为兼容聚合器,但字段必须显式、可类型检查。 - [ ] 以一个业务纵切面迁移并验证后,再迁移下一组,禁止一次替换所有 Oper。 - [ ] 增加 AST 门禁,禁止向 `ChainDataPorts`、`AgentDataPorts` 和新的 canonical use-case service diff --git a/docs/architecture-overview.md b/docs/architecture-overview.md index 3a9265edd..f505fc971 100644 --- a/docs/architecture-overview.md +++ b/docs/architecture-overview.md @@ -705,7 +705,7 @@ flowchart LR | 指标 | 当前值 | |---|---:| | Python 模块 | 844 | -| 内部导入边 | 6,898 | +| 内部导入边 | 6,901 | | 非平凡 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 f15670f6f..5a0ed93cf 100644 --- a/docs/architecture-refactor-roadmap.md +++ b/docs/architecture-refactor-roadmap.md @@ -97,7 +97,7 @@ canonical 主程序;兼容只经统一 Compat/SDK 门面提供。 | 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 全链收口 | `DELIVERED` | S1-L1.4 | `e9de149db`、`a2e249f20`:崩溃矩阵、3.0.17 升降级、重复回放和插件 ABI 验收完整;旧 fail-open、重复状态与兼容层外旧入口删除;Unit Tests `33092427327`、Pylint `33092427348` 全绿,ARCH-102 债务归零 | -| S1-L2 Workflow typed query | `PLANNED` | S0 | Workflow Application Port 不返回 `Any`/ORM,Session 内投影 DTO,正式调用方全部切换 | +| S1-L2 Workflow typed query | `VERIFIED` | 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-L4 Subscription mutation UoW | `PLANNED` | S1-L3 | Subscription mutation 不跨 Session 传 ORM,正式写路径一个 UoW,旧自动事务入口退出 canonical 路径 | | S1-L5 站点/规则引用原子清理 | `PLANNED` | S1-L4 | SystemConfig+Subscribe 同事务更新,commit 后快照原子发布,并发/故障注入无部分状态 | @@ -144,7 +144,7 @@ canonical 主程序;兼容只经统一 Compat/SDK 门面提供。 | 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 | 当前受控 888 条诊断归零,规则集扩展经过独立审查且新增诊断为零 | +| S4-L5 Ruff 治理债务清零 | `PLANNED` | S3 | 当前受控 881 条诊断归零,规则集扩展经过独立审查且新增诊断为零 | | S4-L6 Coverage/并发/质量证据 | `PLANNED` | S3,S4-L1,S4-L2 | 高风险包纳入 coverage;raw concurrency 分类清零;Module Quality 有真实 evidence test | ### S5:Plugin、Agent、Domain、Startup 与最终收口 diff --git a/docs/rules/05-architecture.md b/docs/rules/05-architecture.md index 53cc82d30..e37dc6722 100644 --- a/docs/rules/05-architecture.md +++ b/docs/rules/05-architecture.md @@ -225,9 +225,14 @@ 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__`。 协程环境文件日志属于有界 E1 观测能力,只允许单一队列 writer;队列满时不得再以无界 executor 形成第二条异步写入路径。日志关闭必须有限等待 writer 与文件处理器,未收敛时 `LoggerManager` 保留原 owner 并让 lifespan 以关闭失败结束,不得先清空引用或用无界 `join()` 掩盖失败。 @@ -667,7 +672,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..0d498972a 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -214,8 +214,8 @@ def configure_plugin_system_services(): 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 +229,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 @@ -332,7 +335,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 +352,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/dependency-baseline.json b/tests/fixtures/architecture/dependency-baseline.json index 3e75a6bf2..2fac978e0 100644 --- a/tests/fixtures/architecture/dependency-baseline.json +++ b/tests/fixtures/architecture/dependency-baseline.json @@ -1441,8 +1441,8 @@ "runtime_only": true } }, - "edge_count": 6898, - "edge_sha256": "9bceb8aaa1bb5857ba6ee075c50214455cfb4f59030cee49b3e5af751cc95972", + "edge_count": 6901, + "edge_sha256": "860590c25bd889096c9faa04f35ad3e3312e2c28ab415e57dce1d669d490e0a1", "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", @@ -4340,6 +4340,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", @@ -4472,6 +4474,8 @@ "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", @@ -5201,6 +5205,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", @@ -7551,6 +7557,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", @@ -7576,7 +7583,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", @@ -8167,14 +8173,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", 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 6758b3e1d..d0311fd57 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 }, @@ -1056,9 +1047,6 @@ "E402": 27, "F401": 1 }, - "app/workflow/__init__.py": { - "F401": 1 - }, "scripts/architecture/task_ownership.py": { "I001": 1 }, @@ -1362,10 +1350,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 }, @@ -1786,9 +1770,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 14942cc72..b42dfb21c 100644 --- a/tests/fixtures/architecture/runtime-contract-baseline.json +++ b/tests/fixtures/architecture/runtime-contract-baseline.json @@ -1184,6 +1184,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 +1447,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 +1843,7 @@ "50704edda70674af0932e769ddda40d2c21d50b134155c975d675035d7c933cd" ], "producer_fingerprints": [ - "611abaaf708f0c2d555d3b963ff540c015f2155ff493ed17b92d7c7f0bd45e96" + "b92d348c57dac3ba348079b4e71d02a0c13933bbf14b7e81188b042c3c5e2db3" ] } }, @@ -2980,10 +2987,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..9e9d7c9c4 100644 --- a/tests/test_architecture_dependencies.py +++ b/tests/test_architecture_dependencies.py @@ -283,6 +283,93 @@ 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_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_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_import_compat.py b/tests/test_legacy_import_compat.py index 999014e14..3c05b33ef 100644 --- a/tests/test_legacy_import_compat.py +++ b/tests/test_legacy_import_compat.py @@ -410,6 +410,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 +469,7 @@ 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.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..c65b06259 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, @@ -638,6 +638,11 @@ def test_workflow_chain_process_serializes_circular_context(monkeypatch): 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_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) @@ -825,7 +830,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 +909,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} @@ -941,6 +946,11 @@ def test_workflow_chain_rejects_execution_before_persisting_running_state(monkey 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_query", + lambda: SimpleNamespace(get_sync=lambda _workflow_id: workflow), + ) def unexpected_resume(_workflow_id: int) -> None: """拒绝准入时若仍恢复停止标记则立即暴露回归。""" @@ -970,7 +980,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 = {} @@ -979,6 +989,13 @@ def test_workflow_chain_releases_admitted_owner_when_start_fails(monkeypatch): 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_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 +1035,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 +1074,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..aa09d2dac 100644 --- a/tests/test_workflow_mutation_command.py +++ b/tests/test_workflow_mutation_command.py @@ -8,6 +8,7 @@ from app.application.workflow import ( WorkflowExecutionCommand, WorkflowMutationCommand, WorkflowQueryService, + WorkflowSnapshot, ) @@ -21,6 +22,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() @@ -96,8 +121,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 +131,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) 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",