diff --git a/app/application/workflow.py b/app/application/workflow.py index f547a72b5..962ae8a69 100644 --- a/app/application/workflow.py +++ b/app/application/workflow.py @@ -44,6 +44,26 @@ class WorkflowRuntime(Protocol): """按最新定义刷新工作流事件触发器。""" ... + def register_execution(self, owner: "WorkflowExecutionOwner") -> bool: + """登记活动工作流执行 owner;停机封口后返回 False。""" + ... + + def unregister_execution(self, owner: "WorkflowExecutionOwner") -> None: + """在工作流执行真实终止后释放 owner。""" + ... + + +class WorkflowExecutionOwner(Protocol): + """声明 concrete 工作流管理器需要持有的执行生命周期能力。""" + + def request_stop(self) -> None: + """请求停止继续调度,并通知支持取消的活动动作。""" + ... + + def wait_stopped(self, timeout: float) -> bool: + """有限等待执行及其节点线程池真实终止。""" + ... + WorkflowRuntimeProvider = Callable[[], WorkflowRuntime] diff --git a/app/chain/workflow.py b/app/chain/workflow.py index 0b4d0d6b1..f58c23345 100644 --- a/app/chain/workflow.py +++ b/app/chain/workflow.py @@ -5,11 +5,10 @@ import inspect import pickle import threading from collections import defaultdict, deque -from concurrent.futures import ThreadPoolExecutor from contextvars import Context, copy_context from datetime import date, datetime from functools import partial -from time import sleep +from time import monotonic, sleep from typing import Any, Callable, List, Optional, Tuple from pydantic import BaseModel @@ -19,6 +18,7 @@ from app.runtime.config import global_vars from app.runtime.events import Event, eventmanager from app.application.workflow import get_workflow_manager from app.application.chain.data import get_chain_workflow_port +from app.runtime.execution import OwnedThreadPoolExecutor from app.runtime.log import logger from app.schemas.workflow import ActionContext from app.schemas.workflow import ActionFlow @@ -29,6 +29,7 @@ from app.schemas.types import EventType ARTIFACT_FIELDS = {"torrents", "medias", "fileitems", "downloads", "sites", "subscribes"} DEFAULT_WORKFLOW_MAX_WORKERS = 4 +WORKFLOW_EXECUTOR_STOP_TIMEOUT_SECONDS = 10.0 CIRCULAR_REFERENCE_PLACEHOLDER = "[Circular]" Workflow = Any @@ -93,18 +94,27 @@ class WorkflowCancelToken: 工作流取消令牌。 """ - def __init__(self, workflow_id: int): + def __init__( + self, + workflow_id: int, + stop_event: Optional[threading.Event] = None, + ) -> None: """ 初始化取消令牌。 :param workflow_id: 工作流ID + :param stop_event: 单次执行 owner 的本地停止信号 """ self.workflow_id = workflow_id + self.stop_event = stop_event def is_cancelled(self) -> bool: """ 判断工作流是否已被取消。 """ - return global_vars.is_workflow_stopped(self.workflow_id) + return bool( + (self.stop_event and self.stop_event.is_set()) + or global_vars.is_workflow_stopped(self.workflow_id) + ) class WorkflowExecutor: @@ -163,6 +173,13 @@ class WorkflowExecutor: # 工作流管理器 self.workflowmanager = get_workflow_manager() + # 具体管理器登记活动执行;旧自定义 provider 不实现 owner 接口时保持原调用兼容。 + self._execution_lock = threading.RLock() + self._admission_state = "pending" + self._registered_execution = False + self._stop_event = threading.Event() + self._execute_returned = threading.Event() + self._stopped_event = threading.Event() # 线程安全队列 self.queue = deque() self.queued_actions = set() @@ -170,8 +187,8 @@ class WorkflowExecutor: # 锁用于保证线程安全 self.lock = threading.Lock() # 线程池 - self.executor = ThreadPoolExecutor(max_workers=self.get_workflow_max_workers()) - self.cancel_token = WorkflowCancelToken(self.workflow.id) + self.executor = OwnedThreadPoolExecutor(max_workers=self.get_workflow_max_workers()) + self.cancel_token = WorkflowCancelToken(self.workflow.id, self._stop_event) # 跟踪运行中的任务数 self.running_tasks = 0 @@ -188,8 +205,6 @@ class WorkflowExecutor: self.context = self.restore_context() self.ensure_context_partitions() - # 恢复工作流 - global_vars.workflow_resume(self.workflow.id) # 恢复时重新释放已终态节点的出边,使后继节点能继续执行或保持跳过传播。 for action_id, state in self.node_states.items(): if state == "success": @@ -201,6 +216,78 @@ class WorkflowExecutor: if self.node_states.get(action_id) == "pending" and not self.incoming_flows.get(action_id): self.enqueue_node(action_id) + def admit(self) -> bool: + """向 concrete manager 登记本次执行,并保持旧 provider 可直接运行。""" + with self._execution_lock: + if self._admission_state == "admitted": + return True + if self._admission_state == "rejected": + return False + register = getattr(self.workflowmanager, "register_execution", None) + if callable(register) and not register(self): + self._admission_state = "rejected" + self.success = False + self.stopped = True + self.errmsg = "工作流服务正在停止" + self.executor.shutdown_bounded(timeout=0.0) + self._execute_returned.set() + self._stopped_event.set() + return False + self._registered_execution = callable(register) + self._admission_state = "admitted" + # 只有获得执行准入后才能清除历史单工作流停止标记。 + global_vars.workflow_resume(self.workflow.id) + return True + + def request_stop(self) -> None: + """停止调度新节点,并通过本地令牌通知支持取消的活动动作。""" + self._stop_event.set() + + def abort_before_execute(self) -> None: + """执行状态启动失败时释放尚未使用的节点池和 manager owner。""" + self.request_stop() + converged = self.executor.shutdown_bounded(timeout=0.0) + self._execute_returned.set() + if converged: + self._release_execution() + + def wait_stopped(self, timeout: float) -> bool: + """有限等待 execute 返回,并在需要时重试节点线程池收敛。""" + deadline = monotonic() + max(0.0, timeout) + if self._stopped_event.is_set(): + return True + if not self._execute_returned.wait( + timeout=max(0.0, deadline - monotonic()), + ): + return False + if self._stopped_event.is_set(): + return True + if not self.executor.shutdown_bounded( + timeout=max(0.0, deadline - monotonic()), + ): + return False + self._release_execution() + return self._stopped_event.is_set() + + def _release_execution(self) -> None: + """从 concrete manager 释放已真实终止的执行 owner。""" + with self._execution_lock: + if self._stopped_event.is_set(): + return + if self._registered_execution: + unregister = getattr(self.workflowmanager, "unregister_execution", None) + if callable(unregister): + unregister(self) + self._registered_execution = False + self._stopped_event.set() + + def _stop_requested(self) -> bool: + """判断本次执行或全局工作流是否已收到停止请求。""" + return bool( + self._stop_event.is_set() + or global_vars.is_workflow_stopped(self.workflow.id) + ) + def get_workflow_max_workers(self) -> int: """ 获取工作流最大并发数。 @@ -339,12 +426,14 @@ class WorkflowExecutor: """ 执行工作流 """ + if not self.admit(): + return try: while True: should_sleep = False node_id = None with self.lock: - if global_vars.is_workflow_stopped(self.workflow.id): + if self._stop_requested(): self.success = False self.stopped = True self.errmsg = "工作流已停止" @@ -385,7 +474,12 @@ class WorkflowExecutor: ) future.add_done_callback(partial(context.run, self.on_node_complete)) finally: - self.executor.shutdown(wait=True, cancel_futures=True) + converged = self.executor.shutdown_bounded( + timeout=WORKFLOW_EXECUTOR_STOP_TIMEOUT_SECONDS, + ) + self._execute_returned.set() + if converged: + self._release_execution() def pop_dispatchable_node(self) -> Optional[str]: """ @@ -432,7 +526,7 @@ class WorkflowExecutor: try: action, action_result = future.result() with self.lock: - if global_vars.is_workflow_stopped(self.workflow.id): + if self._stop_requested(): self.success = False self.stopped = True self.errmsg = "工作流已停止" @@ -1238,10 +1332,16 @@ class WorkflowChain(ChainBase): text=f"开始执行工作流 {workflow.name} ...", data={"total": len(workflow.actions), "finished": 0}, ) - workflowoper.start(workflow_id) - # 执行工作流 executor = WorkflowExecutor(workflow, step_callback=save_step) + if not executor.admit(): + logger.warning("工作流服务正在停止,拒绝执行 %s", workflow.name) + return False, executor.errmsg + try: + workflowoper.start(workflow_id) + except Exception: + executor.abort_before_execute() + raise executor.execute() if executor.stopped: diff --git a/app/startup/initializers/workflow.py b/app/startup/initializers/workflow.py index 8a9e436ec..bd2a72015 100644 --- a/app/startup/initializers/workflow.py +++ b/app/startup/initializers/workflow.py @@ -13,8 +13,8 @@ def init_workflow(): WorkFlowManager() -def stop_workflow(): +def stop_workflow() -> bool: """ - 停止工作流 + 停止工作流并返回全部活动执行是否收敛。 """ - WorkFlowManager().stop() + return WorkFlowManager().stop() diff --git a/app/startup/lifecycle/__init__.py b/app/startup/lifecycle/__init__.py index 00796c9cc..9dadf6f22 100644 --- a/app/startup/lifecycle/__init__.py +++ b/app/startup/lifecycle/__init__.py @@ -482,6 +482,7 @@ def build_lifecycle_components(app: FastAPI) -> tuple[LifecycleComponent, ...]: stop_order=20, start_timeout_seconds=120, stop_timeout_seconds=120, + stop_failure=LifecycleFailurePolicy.FAIL_FAST, ), LifecycleComponent( name="插件备份", diff --git a/app/workflow/__init__.py b/app/workflow/__init__.py index 05412b9ac..ce0e182a5 100644 --- a/app/workflow/__init__.py +++ b/app/workflow/__init__.py @@ -7,6 +7,7 @@ from pydantic import BaseModel from app.runtime.config import global_vars from app.runtime.events import eventmanager, Event from app.application.chain.data import get_chain_workflow_port +from app.application.workflow import WorkflowExecutionOwner from app.foundation.reflection import ModuleHelper from app.runtime.log import logger from app.schemas.workflow import ActionContext @@ -16,17 +17,22 @@ from app.schemas.workflow import Workflow from app.schemas.types import EventType from app.foundation.singleton import Singleton +_WORKFLOW_STOP_TIMEOUT_SECONDS = 10.0 + class WorkFlowManager(metaclass=Singleton): """ 工作流管理器 """ - def __init__(self): + def __init__(self) -> None: + """创建动作、事件触发器和活动执行 owner 注册表。""" # 所有动作定义 - self._lock = threading.Lock() + self._lock = threading.RLock() self._actions: Dict[str, Any] = {} self._event_workflows: Dict[str, List[int]] = {} + self._accepting_executions = True + self._executions: Dict[int, WorkflowExecutionOwner] = {} self.init() def init(self): @@ -62,14 +68,61 @@ class WorkFlowManager(metaclass=Singleton): # 加载工作流事件触发器 self.load_workflow_events() - def stop(self): + def register_execution(self, owner: WorkflowExecutionOwner) -> bool: + """登记活动执行;生命周期封口后拒绝新的工作流。""" + with self._lock: + if not self._accepting_executions: + return False + self._executions[id(owner)] = owner + return True + + def unregister_execution(self, owner: WorkflowExecutionOwner) -> None: + """仅在执行及其节点线程池真实终止后释放 owner。""" + with self._lock: + self._executions.pop(id(owner), None) + + def stop( + self, + timeout: float = _WORKFLOW_STOP_TIMEOUT_SECONDS, + ) -> bool: """ - 停止 + 封口工作流入口并有限等待全部活动执行。 + + :param timeout: 等待活动执行真实终止的最长秒数 + :return: 全部执行终止并安全释放动作注册表时返回 True """ - for event_type_str in list(self._event_workflows.keys()): + with self._lock: + self._accepting_executions = False + event_type_values = tuple(self._event_workflows) + for event_type_str in event_type_values: self.remove_workflow_event(event_type_str=event_type_str) - self._actions = {} - self._event_workflows = {} + with self._lock: + self._event_workflows = {} + executions = tuple(self._executions.values()) + + converged = True + for execution in executions: + try: + execution.request_stop() + except Exception as err: + converged = False + logger.error("请求停止工作流执行失败:%s", err) + + deadline = monotonic() + max(0.0, timeout) + for execution in executions: + try: + if not execution.wait_stopped( + timeout=max(0.0, deadline - monotonic()), + ): + converged = False + except Exception as err: + converged = False + logger.error("等待工作流执行停止失败:%s", err) + with self._lock: + converged = converged and not self._executions + if converged: + self._actions = {} + return converged def execute(self, workflow_id: int, action: Action, context: ActionContext = None, inputs: Optional[dict] = None, runtime: Optional[dict] = None, diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index 5474366b6..67d04a15b 100644 --- a/docs/refactor/backend-architecture-next-stage.md +++ b/docs/refactor/backend-architecture-next-stage.md @@ -2,13 +2,13 @@ > 文档性质:当前架构复核、优秀 Python 后端实践对标、AI 可执行任务手册 > 适用仓库:`MoviePilot`,分支 `v3` -> 审计基线:`e9053a65`(2026-08-24) +> 审计基线:`41b1460b`(2026-08-24) > 审计范围:宿主后端;排除 `app/plugins/**` 运行时插件副本 > 规范优先级:`AGENTS.md` 与 `docs/rules/` 高于本文 > 相关文档:`docs/architecture-overview.md`、`docs/refactor/backend-architecture-governance.md`、`docs/refactor/backend-module-refactor-compatibility.md` -> 实施进度:阶段 0~6 的宿主架构能力已完成收口;API/Application 公共复杂度基线已清零,启动组合根的 SystemConfigOper 构造点已由 14 降至 1;API 进程内后台任务已完成首批统一登记,插件仓适配和 Outbox 外围扩展仍按风险切片推进。Model/Base 查询与写装饰器、legacy 隐式会话外壳均已清零,插件 SDK 也不再导出宿主 Model。2026-08-23 的长期整改阶段 0 已恢复宿主、启动性能、官方插件和 SDK 契约门禁的可信基线;阶段 1a 已补齐 TaskRegistry owner 零债务门禁和诚实的关停超时语义;阶段 1b1 已收口整理 worker、pending 回放、失败通知、进程内 AI 重试、插件监控与事件投递的生命周期所有权;2026-08-24 的阶段 2 已将 212 个已观察宿主模块方法的 legacy aggregation 清零,并补齐可执行 fanout 与下载器文件 DTO 边界;阶段 3 已将消息交互和远程命令的订阅删除统一到 Application/UoW/outbox,宿主不再调用裸线程统计入口;阶段 4 已统一七种消息渠道的宿主回环与后台执行边界;阶段 5 已补齐事件窗口聚合任务的生命周期所有权;阶段 6 已统一插件文件操作的取消完成语义;阶段 7 已统一插件协程补偿的终态等待;阶段 8 已统一宿主同步函数的异步线程池入口;阶段 9 已统一工作流运行时的宿主获取路径;阶段 10 已统一模块、插件与调度运行时的显式 getter 调用;阶段 11 已清除系统配置 getter 的 Oper 形别名;阶段 12 已完成工作流域的显式 Chain 数据端口迁移;阶段 13 已收口用户、交互与消息链的数据端口;阶段 14 已收口音乐订阅数据端口;阶段 15 已收口站点数据端口;阶段 16 已收口媒体服务器数据端口;阶段 17 已收口下载数据端口;阶段 18 已收口主订阅数据端口;阶段 19 已收口整理数据端口;阶段 20 已收口 Agent 数据端口;阶段 21 已收口监控历史端口;阶段 22 已统一服务配置应用边界;阶段 23 已补齐媒体服务器 API 遗留的类形配置读取路径;阶段 24 已清除 Scheduler 内部无 owner 的协程提交双轨;阶段 25 已补齐 TaskRegistry 跨线程 owner 并迁移整理 AI 接管;阶段 26 已统一 Agent 会话清理提交;阶段 27 已统一历史 AI 进度 owner;阶段 28 已托管旧插件订阅统计线程;阶段 29 已统一 Emby 系条目转换并清零重复代码白名单;阶段 30 已收口插件市场请求级子任务;阶段 31 已托管搜索 AI 推荐任务;阶段 32 已清除事件调度器绕过生命周期 owner 的投递回退;阶段 33 已统一宿主 Agent 运行时的获取路径;阶段 34 已统一 durable-required 事件与 Outbox topic 事实源;阶段 35 已统一 LLM provider 管理 API 的运行时解析路径;阶段 36 已统一 WebAgent 音频能力访问边界;阶段 37 已统一插件输入事件发布路径;阶段 38 已统一 WebAgent 通知事件监听与队列边界;阶段 39 已补齐搜索 SSE 断线时的上游任务清理;阶段 40 已补齐异步防抖取消的终态所有权;阶段 41 已统一优雅重启兜底线程的唯一所有权;阶段 42 已补齐 Telegram typing 的多实例隔离和终态 owner;阶段 43 已统一 Discord typing 的异步 owner 和 shutdown 收尾;阶段 44 已清除 WebAgent 测试临时事件循环提前关闭产生的 CI 红注解;阶段 45 已统一影视与字幕搜索的请求级逐页任务编排;阶段 46 已收口启动性能门禁的托管 runner 假失败与诊断输出;阶段 47 已补齐 Agent 渠道流式刷新任务的重入 owner;阶段 48 已统一工件上传 action 的 Node 24 主版本;阶段 49 已统一插件安装的同步/异步代际解析事实源;阶段 50 已统一插件市场 GitHub 请求降级策略;阶段 51 已统一插件索引请求与响应三态策略;阶段 52 已统一插件 Release 分页策略;阶段 53 已统一远端插件安装模式决策;阶段 54 已补齐同步安装成功后的临时回滚备份清理;阶段 55~56 已收口官方插件观察基线与报告保留策略;阶段 57 已统一进程级运行时 Facade 门禁并补齐 ModuleManager 边界;阶段 58 已消除 AgentTask 关闭回归的跨线程零时长等待竞态;阶段 59 已统一 Feishu 多实例长连接的 SDK 循环路由;阶段 60 已清除命令服务虚假的关停 owner 声明;阶段 61 已统一 Capability Runtime 同步/异步关闭的诚实收敛结果;阶段 62 已统一消息渠道长连接的多实例关闭收敛合同;阶段 63 已补齐应用消息队列线程的关闭收敛合同;阶段 64 已统一共享线程池的有界关闭 owner;阶段 65 已统一异步文件日志的单一有界写入和关闭 owner;阶段 66 已统一 DoH 与共享线程池的有界 executor owner。 +> 实施进度:阶段 0~6 的宿主架构能力已完成收口;API/Application 公共复杂度基线已清零,启动组合根的 SystemConfigOper 构造点已由 14 降至 1;API 进程内后台任务已完成首批统一登记,插件仓适配和 Outbox 外围扩展仍按风险切片推进。Model/Base 查询与写装饰器、legacy 隐式会话外壳均已清零,插件 SDK 也不再导出宿主 Model。2026-08-23 的长期整改阶段 0 已恢复宿主、启动性能、官方插件和 SDK 契约门禁的可信基线;阶段 1a 已补齐 TaskRegistry owner 零债务门禁和诚实的关停超时语义;阶段 1b1 已收口整理 worker、pending 回放、失败通知、进程内 AI 重试、插件监控与事件投递的生命周期所有权;2026-08-24 的阶段 2 已将 212 个已观察宿主模块方法的 legacy aggregation 清零,并补齐可执行 fanout 与下载器文件 DTO 边界;阶段 3 已将消息交互和远程命令的订阅删除统一到 Application/UoW/outbox,宿主不再调用裸线程统计入口;阶段 4 已统一七种消息渠道的宿主回环与后台执行边界;阶段 5 已补齐事件窗口聚合任务的生命周期所有权;阶段 6 已统一插件文件操作的取消完成语义;阶段 7 已统一插件协程补偿的终态等待;阶段 8 已统一宿主同步函数的异步线程池入口;阶段 9 已统一工作流运行时的宿主获取路径;阶段 10 已统一模块、插件与调度运行时的显式 getter 调用;阶段 11 已清除系统配置 getter 的 Oper 形别名;阶段 12 已完成工作流域的显式 Chain 数据端口迁移;阶段 13 已收口用户、交互与消息链的数据端口;阶段 14 已收口音乐订阅数据端口;阶段 15 已收口站点数据端口;阶段 16 已收口媒体服务器数据端口;阶段 17 已收口下载数据端口;阶段 18 已收口主订阅数据端口;阶段 19 已收口整理数据端口;阶段 20 已收口 Agent 数据端口;阶段 21 已收口监控历史端口;阶段 22 已统一服务配置应用边界;阶段 23 已补齐媒体服务器 API 遗留的类形配置读取路径;阶段 24 已清除 Scheduler 内部无 owner 的协程提交双轨;阶段 25 已补齐 TaskRegistry 跨线程 owner 并迁移整理 AI 接管;阶段 26 已统一 Agent 会话清理提交;阶段 27 已统一历史 AI 进度 owner;阶段 28 已托管旧插件订阅统计线程;阶段 29 已统一 Emby 系条目转换并清零重复代码白名单;阶段 30 已收口插件市场请求级子任务;阶段 31 已托管搜索 AI 推荐任务;阶段 32 已清除事件调度器绕过生命周期 owner 的投递回退;阶段 33 已统一宿主 Agent 运行时的获取路径;阶段 34 已统一 durable-required 事件与 Outbox topic 事实源;阶段 35 已统一 LLM provider 管理 API 的运行时解析路径;阶段 36 已统一 WebAgent 音频能力访问边界;阶段 37 已统一插件输入事件发布路径;阶段 38 已统一 WebAgent 通知事件监听与队列边界;阶段 39 已补齐搜索 SSE 断线时的上游任务清理;阶段 40 已补齐异步防抖取消的终态所有权;阶段 41 已统一优雅重启兜底线程的唯一所有权;阶段 42 已补齐 Telegram typing 的多实例隔离和终态 owner;阶段 43 已统一 Discord typing 的异步 owner 和 shutdown 收尾;阶段 44 已清除 WebAgent 测试临时事件循环提前关闭产生的 CI 红注解;阶段 45 已统一影视与字幕搜索的请求级逐页任务编排;阶段 46 已收口启动性能门禁的托管 runner 假失败与诊断输出;阶段 47 已补齐 Agent 渠道流式刷新任务的重入 owner;阶段 48 已统一工件上传 action 的 Node 24 主版本;阶段 49 已统一插件安装的同步/异步代际解析事实源;阶段 50 已统一插件市场 GitHub 请求降级策略;阶段 51 已统一插件索引请求与响应三态策略;阶段 52 已统一插件 Release 分页策略;阶段 53 已统一远端插件安装模式决策;阶段 54 已补齐同步安装成功后的临时回滚备份清理;阶段 55~56 已收口官方插件观察基线与报告保留策略;阶段 57 已统一进程级运行时 Facade 门禁并补齐 ModuleManager 边界;阶段 58 已消除 AgentTask 关闭回归的跨线程零时长等待竞态;阶段 59 已统一 Feishu 多实例长连接的 SDK 循环路由;阶段 60 已清除命令服务虚假的关停 owner 声明;阶段 61 已统一 Capability Runtime 同步/异步关闭的诚实收敛结果;阶段 62 已统一消息渠道长连接的多实例关闭收敛合同;阶段 63 已补齐应用消息队列线程的关闭收敛合同;阶段 64 已统一共享线程池的有界关闭 owner;阶段 65 已统一异步文件日志的单一有界写入和关闭 owner;阶段 66 已统一 DoH 与共享线程池的有界 executor owner;阶段 67 已补齐工作流活动执行的生命周期 owner。 > 当前 canonical 状态:API/Application 公共复杂度基线已清零,组合根外 `SystemConfigOper()` 构造和 Model/Oper 隐式事务均为 0;命名 Chain/Agent 数据端口、TaskRegistry owner、Module Contract V2、typed Event、Outbox durable intent、请求关联和插件运行时 getter 已形成当前路径。插件仓适配、未知第三方 fallback 和其它 E1/E3 副作用仍按风险持续治理。 -> 最新阶段:阶段 66 已统一 DoH 与共享线程池的有界 executor owner。 +> 最新阶段:阶段 67 已补齐工作流活动执行的生命周期 owner。 ## 当前复核结论(2026-08-24) @@ -17,7 +17,7 @@ ### 长期整改阶段 0:治理门禁恢复(2026-08-23) -- 宿主依赖基线已审查 TaskRegistry、有界后台 owner 与插件变更准入接入后的语义差异:当前为 `810` 个模块、`6562` 条内部导入边,12 组重点禁止边继续全部为 `0`,唯一非平凡 SCC 仍是隔离的 TMDB 移植包。 +- 宿主依赖基线已审查 TaskRegistry、有界后台 owner 与插件变更准入接入后的语义差异:当前为 `810` 个模块、`6564` 条内部导入边,12 组重点禁止边继续全部为 `0`,唯一非平凡 SCC 仍是隔离的 TMDB 移植包。 - 启动性能探针会在隔离生命周期中真实创建并释放 TaskRegistry;normal/safe 组件数分别为 `23`/`11`,CI 只读检查使用稳定的宿主模块集合和生命周期组件顺序,不再把 Python/平台模块数量当作硬合同。 - 官方插件快照覆盖 `plugins.v3`、`plugins.v2` 以及 V3 实际会从 `package.json` 回退加载的 31 个默认实现;`app/plugins/**` 仍只是宿主运行副本,不进入扫描。 - SDK 快照以各模块显式 `__all__` 为公开合同,能够记录赋值别名;`typing`、`__future__` 等实现期导入不再被误冻结,既有数据库备份门面已补精确导出清单。 @@ -727,13 +727,36 @@ Pylint 10.00/10、strict mypy 41 文件、全部宿主与质量 ratchet 及四分片全量 `5955 passed, 3 skipped` 均通过。 +### 长期整改阶段 67:工作流活动执行生命周期 owner 收口(2026-08-24) + +- 工作流生命周期组件原本先于 Scheduler、插件和模块关闭,但 concrete `WorkFlowManager.stop()` 只注销 + 事件监听并清空动作表,不登记或等待 `WorkflowExecutor`。阻塞动作探针确认旧 stop 立即返回 `None`, + 活动执行线程仍存活且动作表已经清空,后续组件会继续释放仍被节点使用的运行时依赖。 +- Application 工作流运行时协议现在声明最小 `WorkflowExecutionOwner`;concrete manager 是全部活动执行的 + 唯一注册表,停机先原子封口新准入,再向快照中的每个 owner 请求本地取消并使用 10 秒共享 deadline + 有限等待。单个 owner 抛错或超时不会跳过其它执行的停止尝试;未收敛 owner 和动作表均保留供同实例重试。 +- `WorkflowExecutor` 的节点池不再维护独立标准库实现,统一复用阶段 66 的 + `OwnedThreadPoolExecutor`;正常路径在节点完成后即时收敛,异常/阻塞路径保留执行 owner。manager 的本地 + 停止事件同时进入 `WorkflowCancelToken`,支持取消的重试和循环动作可尽快退出,不响应取消的第三方同步 + 动作则诚实返回未收敛。 +- `WorkflowChain.process()` 在写入运行中状态前取得执行准入,停机后的 API、事件、Agent 或 Scheduler + 新触发不会制造新的 `R` 状态。工作流组件改为 fail-fast owner;活动执行未终止时不得继续关闭监控器、 + Scheduler、插件、模块或 HTTP 依赖。 +- `WorkflowExecutor`、`WorkflowCancelToken(workflow_id)`、`WorkFlowManager()`、无参数 `stop()`、 + `WorkflowChain.process()`、旧自定义 runtime provider、事件/定时/API/Agent 调用路径均保留;新增可选取消 + 信号、准入和布尔收敛结果为兼容扩展。未修改 SDK/Compat 清单、插件 Hook 或插件仓,V1/V2/V3 边界不变。 +- 阻塞动作、有限返回、本地取消、owner/动作表保留、拒绝新执行、释放后重试、多 owner 异常隔离和生命周期 + fail-fast 均有故障注入;工作流/生命周期专项 91 项、工作流与旧导入联合专项 213 项、架构合同 94 项、 + Pylint 10.00/10、strict mypy 41 文件、全部宿主与质量 ratchet 及四分片全量 + `5959 passed, 3 skipped` 均通过。 + ### 总体判断 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: - 继续采用单进程控制面是正确选择,不建议现在拆成微服务;插件、调度器、工作流、事件和数据库共享进程内状态,拆分会放大部署、事务和兼容成本。 - `foundation/domain/runtime/adapters/application/chain/api/startup` 的职责方向基本成立;宿主架构基线、复杂度 ratchet、异步阻塞 ratchet 当前均通过。 -- 依赖图当前为 `810` 个 Python 模块、`6560` 条内部导入边;唯一非平凡 SCC 位于隔离的 TMDB 第三方移植包内部,不应为了指标归零重写。 +- 依赖图当前为 `810` 个 Python 模块、`6564` 条内部导入边;唯一非平凡 SCC 位于隔离的 TMDB 第三方移植包内部,不应为了指标归零重写。 - 当前主要风险已经从“目录和依赖失控”转移到运行时协议、后台副作用的可靠性和遗留兼容面。换言之,下一阶段重点应是**语义收口和可验证性**,而不是继续搬文件或机械拆大文件。 综合评价:架构方向可持续,生产可用性较高;可演进性仍处于中等水平。现阶段没有静态审计发现必须立即推倒重来的 P0 架构问题,但存在需要按 P1/P2 计划治理的真实债务。 diff --git a/docs/rules/05-architecture.md b/docs/rules/05-architecture.md index 345cf2677..69f3ba2fd 100644 --- a/docs/rules/05-architecture.md +++ b/docs/rules/05-architecture.md @@ -204,6 +204,9 @@ ModuleManager 与 startup 组合根继续关闭其余资源但必须向上返回 `app.runtime.execution.OwnedThreadPoolExecutor` 是进程级同步执行器有界收敛的唯一事实源;新的专用 线程池不得复制 Future 追踪、worker join 或重试关闭实现。DoH 查询线程池也必须复用该 owner:恢复系统 DNS 后有限等待,超时保留原 executor 并向 startup 返回 `False`,真实收敛前不得创建替代线程池或回填缓存。 +工作流节点线程池同样复用该 executor;所有 `WorkflowExecutor` 必须在 concrete `WorkFlowManager` 登记, +manager 停机先封口新执行并向活动 owner 发送本地取消,再有限等待执行线程和节点 worker。未收敛时必须 +保留动作注册表和执行 owner,并让工作流生命周期 fail-fast,禁止继续释放仍被动作使用的插件或模块依赖。 协程环境文件日志属于有界 E1 观测能力,只允许单一队列 writer;队列满时不得再以无界 executor 形成第二条异步写入路径。日志关闭必须有限等待 writer 与文件处理器,未收敛时 `LoggerManager` 保留原 owner 并让 lifespan 以关闭失败结束,不得先清空引用或用无界 `join()` 掩盖失败。 diff --git a/tests/fixtures/architecture/dependency-baseline.json b/tests/fixtures/architecture/dependency-baseline.json index d6a4790a9..1b5c7b0e7 100644 --- a/tests/fixtures/architecture/dependency-baseline.json +++ b/tests/fixtures/architecture/dependency-baseline.json @@ -13,8 +13,8 @@ "runtime_to_db": [], "workflow_to_db": [] }, - "edge_count": 6562, - "edge_sha256": "08b6e0005a3c316a1de675193f7749ff8a93e8c1c052eb48339aa58b308ad6bf", + "edge_count": 6564, + "edge_sha256": "333de3c1ef3f7e49be96dda288a0f382bcd19a853817708453d67557d138aa80", "edges": [ "app -> app.runtime", "app -> app.runtime.compat", @@ -3495,6 +3495,7 @@ "app.chain.workflow -> app.runtime", "app.chain.workflow -> app.runtime.config", "app.chain.workflow -> app.runtime.events", + "app.chain.workflow -> app.runtime.execution", "app.chain.workflow -> app.runtime.log", "app.chain.workflow -> app.schemas", "app.chain.workflow -> app.schemas.types", @@ -6407,6 +6408,7 @@ "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", diff --git a/tests/test_lifecycle_shutdown.py b/tests/test_lifecycle_shutdown.py index 03e3f6439..3f0977804 100644 --- a/tests/test_lifecycle_shutdown.py +++ b/tests/test_lifecycle_shutdown.py @@ -127,7 +127,6 @@ def _patch_lifespan(monkeypatch, *, failing_step: str | None = None) -> dict: "failing_step", [ "backup_plugins", - "stop_workflow", "stop_modules", "close_http", ], @@ -275,6 +274,7 @@ _ORDERED_SHUTDOWN_STEPS = ( "failing_step", ( "stop_plugin_monitor", + "stop_workflow", "stop_monitor", "stop_scheduler", "stop_agent", @@ -579,6 +579,7 @@ def test_lifecycle_manifest_declares_normal_and_safe_mode_order() -> None: } == { "插件变更监控", "后台任务登记器", + "工作流", "监控器", "定时器", "AI智能体会话", @@ -596,6 +597,7 @@ def test_lifecycle_manifest_declares_normal_and_safe_mode_order() -> None: not in { "插件变更监控", "后台任务登记器", + "工作流", "监控器", "定时器", "AI智能体会话", diff --git a/tests/test_workflow_execution.py b/tests/test_workflow_execution.py index a7992a812..e48d90256 100644 --- a/tests/test_workflow_execution.py +++ b/tests/test_workflow_execution.py @@ -4,6 +4,8 @@ import threading import time from types import SimpleNamespace +import pytest + from app.chain import workflow as workflow_module from app.runtime.correlation import correlation_scope, get_correlation_id from app.schemas.types import EventType @@ -802,6 +804,190 @@ def test_workflow_executor_stop_is_not_success(monkeypatch): assert executor.errmsg == "工作流已停止" +def test_workflow_manager_shutdown_retains_blocked_execution_for_retry(monkeypatch): + """阻塞动作超时时保留 manager owner,释放后同一执行可重试收敛。""" + entered = threading.Event() + release = threading.Event() + + class BlockingAction: + """模拟不响应取消令牌的第三方同步工作流动作。""" + + def __init__(self, action_id): + """保存动作标识。""" + self.action_id = action_id + self.success = True + self.message = "" + + def execute_with_inputs(self, workflow_id, params, inputs, runtime, context): + """阻塞到测试显式释放,并返回原工作流上下文。""" + _ = workflow_id, params, inputs, runtime + entered.set() + release.wait() + return ActionResult(success=True, context=context) + + manager = object.__new__(workflow_package.WorkFlowManager) + manager._lock = threading.RLock() + manager._actions = {"BlockingAction": BlockingAction} + manager._event_workflows = {} + manager._accepting_executions = True + manager._executions = {} + workflow = _build_workflow( + actions=[{ + "id": "A", + "type": "BlockingAction", + "name": "阻塞动作", + "data": {}, + }], + flows=[], + ) + monkeypatch.setattr(workflow_module, "get_workflow_manager", lambda: manager) + monkeypatch.setattr(workflow_module.global_vars, "workflow_resume", lambda _workflow_id: None) + monkeypatch.setattr( + workflow_module.global_vars, + "is_workflow_stopped", + lambda _workflow_id: False, + ) + + executor = workflow_module.WorkflowExecutor(workflow) + execution_thread = threading.Thread(target=executor.execute, daemon=True) + execution_thread.start() + try: + assert entered.wait(timeout=1) + started_at = time.monotonic() + + assert manager.stop(timeout=0.01) is False + + assert time.monotonic() - started_at < 1 + assert executor.cancel_token.is_cancelled() is True + assert manager._actions == {"BlockingAction": BlockingAction} + assert manager._executions == {id(executor): executor} + + rejected = workflow_module.WorkflowExecutor(workflow) + assert rejected.admit() is False + assert rejected.errmsg == "工作流服务正在停止" + finally: + release.set() + execution_thread.join(timeout=2) + + assert not execution_thread.is_alive() + assert executor.wait_stopped(timeout=1) is True + assert manager.stop(timeout=1) is True + assert manager._executions == {} + assert manager._actions == {} + + +def test_workflow_manager_shutdown_continues_across_owner_failures(): + """单个执行 owner 抛错不得跳过其它活动工作流的停止和等待。""" + + class ObservedOwner: + """记录 manager 对多个 owner 的关闭调用。""" + + def __init__(self, manager, *, fail: bool = False): + """保存管理器、失败开关和调用计数。""" + self.manager = manager + self.fail = fail + self.stop_calls = 0 + self.wait_calls = 0 + + def request_stop(self) -> None: + """记录停止请求,并按需模拟第三方 owner 异常。""" + self.stop_calls += 1 + if self.fail: + raise RuntimeError("stop failed") + + def wait_stopped(self, timeout: float) -> bool: + """记录等待;正常 owner 从 manager 注册表释放自身。""" + _ = timeout + self.wait_calls += 1 + if self.fail: + raise RuntimeError("wait failed") + self.manager.unregister_execution(self) + return True + + manager = object.__new__(workflow_package.WorkFlowManager) + manager._lock = threading.RLock() + action_marker = object() + manager._actions = {"FakeAction": action_marker} + manager._event_workflows = {} + manager._accepting_executions = True + manager._executions = {} + failing_owner = ObservedOwner(manager, fail=True) + healthy_owner = ObservedOwner(manager) + assert manager.register_execution(failing_owner) is True + assert manager.register_execution(healthy_owner) is True + + assert manager.stop(timeout=0.01) is False + + assert failing_owner.stop_calls == 1 + assert failing_owner.wait_calls == 1 + assert healthy_owner.stop_calls == 1 + assert healthy_owner.wait_calls == 1 + assert manager._executions == {id(failing_owner): failing_owner} + assert manager._actions == {"FakeAction": action_marker} + + +def test_workflow_chain_rejects_execution_before_persisting_running_state(monkeypatch): + """停机封口后的新执行不得先把数据库状态写成运行中。""" + + class RejectingWorkflowManager(_FakeWorkflowManager): + """模拟已经封口的 concrete 工作流运行时。""" + + def register_execution(self, _owner) -> bool: + """拒绝停机后的新执行 owner。""" + return False + + workflow = _build_workflow() + workflowoper = _FakeWorkflowOper(workflow) + manager = RejectingWorkflowManager([]) + monkeypatch.setattr(workflow_module, "get_workflow_manager", lambda: manager) + monkeypatch.setattr(workflow_module, "get_chain_workflow_port", lambda: workflowoper) + + def unexpected_resume(_workflow_id: int) -> None: + """拒绝准入时若仍恢复停止标记则立即暴露回归。""" + raise AssertionError("拒绝准入时不得恢复工作流") + + monkeypatch.setattr( + workflow_module.global_vars, + "workflow_resume", + unexpected_resume, + ) + + success, message = workflow_module.WorkflowChain.process(workflow_id=1) + + assert success is False + assert message == "工作流服务正在停止" + assert workflowoper.started is False + + +def test_workflow_chain_releases_admitted_owner_when_start_fails(monkeypatch): + """数据库启动状态写入异常时不得遗留尚未执行的 manager owner。""" + + class FailingWorkflowOper(_FakeWorkflowOper): + """模拟执行状态 start 事务失败。""" + + def start(self, wid): + """在 owner 已准入后抛出持久化异常。""" + _ = wid + raise RuntimeError("start failed") + + manager = object.__new__(workflow_package.WorkFlowManager) + manager._lock = threading.RLock() + manager._actions = {"FakeAction": object()} + manager._event_workflows = {} + manager._accepting_executions = True + 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.global_vars, "workflow_resume", lambda _workflow_id: None) + + with pytest.raises(RuntimeError, match="start failed"): + workflow_module.WorkflowChain.process(workflow_id=1) + + assert manager._executions == {} + assert manager._accepting_executions is True + + def test_workflow_context_merge_preserves_runtime_objects(): """合并上下文时应保留运行时对象,而不是转成字典。""" executor = object.__new__(workflow_module.WorkflowExecutor)