跳到主要内容

Scope

本文聚焦当前项目的 Workflow 执行系统,不覆盖编辑器画布 UI 细节。目标是说明工作流如何被触发、如何进入执行引擎、执行期间如何维护上下文、如何处理暂停/恢复/交互、以及最终如何落盘。

Main Responsibilities

  • 接收执行请求:WebSocket 手动运行、Hook HTTP 触发、Cron 定时触发。
  • 构建执行快照:从 workflow 模板或请求内 snapshot 派生本次执行的 nodes/edges/groups/variables。
  • 维护执行会话:保存状态、执行顺序、上下文、步骤日志、最近事件、断点/暂停信息。
  • 分派节点执行:内建节点、Agent 节点、SQLite/KB 节点、子工作流、循环节点、插件节点、客户端节点。
  • 对外发事件:节点开始/完成/报错、执行日志、上下文、暂停/恢复、完成/失败。
  • 处理阻塞交互:alert/prompt/form/table 节点通过客户端往返恢复执行。
  • 落盘历史:保存 execution log,并暴露短期恢复态给前端重连。

Entry Points

  • WebSocket 手动执行:packages/server/src/ws/execution-channels.ts
    • workflow:execute -> executionManager.execute()
    • workflow:pause -> executionManager.pause()
    • workflow:resume -> executionManager.resume()
    • workflow:stop -> executionManager.stop()
    • workflow:debug-node -> executionManager.debugNode()
    • workflow:get-execution-recovery -> executionManager.getExecutionRecovery()
  • Hook 触发:packages/server/src/routes/workflow-hook.ts
    • POST /api/workflows/hook/:hookName -> triggerService.getHookBindings() -> executionManager.execute()
  • Cron 触发:packages/server/src/services/workflow-trigger-service.ts
    • WorkflowTriggerService.start() 启动时扫描全部 workflow
    • registerCronJob()node-cron 注册定时任务,触发后直接 executionManager.execute({ workflowId }, '__cron__')
  • 服务启动 wiring:packages/server/src/app.ts
    • 创建 InteractionManagerClientNodeManagerExecutionManagerWorkflowTriggerService
    • 挂载 createWorkflowHookRouter(triggerService, executionManager)
    • 启动后调用 triggerService.start()

Execution Flow

packages/server/src/app.ts
-> new ExecutionManager(...)
-> registerExecutionChannels(executionManager)
-> triggerService.setExecutionManager(executionManager)
-> triggerService.start()

WebSocket client
-> packages/server/src/ws/execution-channels.ts: registerHandler('workflow:execute')
-> ExecutionManager.execute(request, clientId, eventSink, workspaceId)
-> workflowStore.getWorkflow(workflowId)
-> resolveExecutionSnapshot()
-> createSession()
-> run()
-> buildExecutionOrder()
-> runSafe()
-> runFromIndex()
-> executeNode()
-> dispatchNode()
-> emitEvent()/emitLog()/emitContext()
-> workflowStore.addExecutionLog()

HTTP hook
-> packages/server/src/routes/workflow-hook.ts: POST /hook/:hookName
-> triggerService.getHookBindings(hookName)
-> executionManager.execute(...)

Cron
-> WorkflowTriggerService.registerCronJob()
-> node-cron callback
-> executionManager.execute({ workflowId }, '__cron__')

Session State Model

ExecutionManager.createSession()packages/server/src/services/execution-manager.ts 构造执行会话,核心字段包括:

  • nodes / edges / groups / variables
    • 本次执行快照,优先使用请求里的 snapshot,否则克隆 workflow 当前定义。
  • context
    • __input__:启动输入。
    • __env__:变量输出对象与请求 env 合并后的环境。
    • __data__:额外上下文数据。
    • __config__:执行前加载的插件配置。
    • 以及每个节点执行后的 session.context[node.id] = step.output
  • executionOrder
    • buildExecutionOrder() 基于 runtime edges 拓扑排序得到。
  • steps
    • 逐节点执行日志,保存输入、输出、错误、分步 logs。
  • status
    • idle / running / paused / completed / error
  • currentIndexpauseRequestedstopRequestedpauseReasonpauseNodeId
    • 用于人工暂停、断点暂停、停止与恢复。
  • recentEvents
    • 供恢复接口返回最近执行事件。

Snapshot And Partial Execution

resolveExecutionSnapshot() 支持两种重要模式:

  • 全量执行
    • 直接使用 workflow 当前图,或请求显式传入的 snapshot。
  • 局部执行
    • 当请求带 startNodeId 时,先校验目标起点,再调用 buildReachableSnapshot() 裁剪可达节点子图。
    • 如果根层存在多个 start 节点且未显式指定 startNodeId,直接报错。

这意味着执行系统把“编辑中的未保存画布”与“已保存模板”统一成同一套 snapshot 执行模型。

Node Dispatch

真正的节点执行分两层:

  • executeNode()
    • 解析变量引用和 dry-run 覆盖。
    • 生成 step,发送 node:start
    • dispatchNode() 拿结果。
    • 写入 session.context 和节点执行数据,发送 node:completenode:error
  • dispatchNode()
    • switch (node.type) 分派到各类实现。

已验证的节点类别:

  • 基础流程节点:startendswitchlooploop_breaksub_workflow
  • 数据处理节点:parse_jsonstring_concatflatten_arrayvariable_aggregate、变量读写删
  • 执行节点:run_coderun_python
  • 数据源节点:sqlite_*kb_*
  • AI 节点:agent_run
  • 交互节点:alertpromptformtable_display
  • 展示/无操作节点:markdownsticky_notegallery_preview
  • 扩展节点:
    • isClientPluginNode(node) 为真时走 executeClientNode()
    • pluginService.canExecuteWorkflowNode(node.type) 为真时走服务端插件执行
    • 插件声明 requiresClientExecution() 时仍回落到客户端执行

这说明扩展点不是单一 registry,而是三层并存:

  • 内建 switch(node.type) 分派
  • 服务端插件运行时 pluginService
  • 客户端插件桥接 ClientNodeManager

Interaction And Client Round-Trips

阻塞式交互由 InteractionManager 统一管理,位置在 packages/server/src/services/interaction-manager.ts

调用链:

ExecutionManager.executeAlertDialog/executePromptDialog/executeFormDialog/executeTableDisplay
-> interactionManager.request(...)
-> sendToClient(clientId, payload[channel='workflow:interaction'])
-> 前端弹窗收集结果
-> WS 'workflow:interaction'
-> packages/server/src/ws/handler.ts
-> handleInteractionResponse(...)
-> InteractionManager.handleResponse(...)
-> 原 Promise resolve/reject
-> ExecutionManager 继续执行

关键机制:

  • 默认超时 5 分钟。
  • 客户端断线后保留 30 秒重连宽限期。
  • cancelExecution() 会中止某个 executionId 下所有待响应交互。
  • 重连时会重新发送 pending payload。

Pause / Resume / Recovery

暂停恢复有两层:

  • 运行态暂停
    • pause() 仅把 pauseRequested = true
    • 真正暂停发生在下一个节点边界:executeNodeAtIndex() 检测后把状态切到 paused,记录 pauseReason='manual'
    • resume() 清掉 pause 状态,重新 runSafe(session, session.currentIndex)
  • 重连恢复
    • getExecutionRecovery() 先找活跃 session,再找 finishedRecoveries
    • createRecoveryState() 返回 statuscurrentNodeIdpauseReasonlogcontextrecentEvents
    • 已结束执行不会永久常驻,只在 finishedRecoveries 中短期保留,随后 pruneFinishedRecoveries() 清理

这套设计偏向“前端刷新后恢复执行视图”,不是持久化级别的引擎续跑。

Persistence And External Boundaries

Workflow 存储采用每个 workflow 一个目录的结构,定义在 packages/server/src/storage/workflow-store.ts

目录模型:

workflows/{workflowId}/
workflow.json
versions/{versionId}.json
execution_history/{logId}.json
plugin_configs/{pluginId}/{schemeName}.json
staging.json
operation_history.json
chat.json
workflows/folders.json

稳定结论:

  • 模板本体:workflow.json
  • 执行历史:execution_history/*.json
    • ExecutionManager.persistAndCleanup() 最终调用 workflowStore.addExecutionLog()
    • 默认保留最近 100 条
  • 版本快照:versions/*.json
    • 默认保留最近 100 个
  • 插件配置:plugin_configs/{pluginId}/{scheme}.json
  • 编辑暂存:staging.json
  • 画布操作历史:operation_history.json
  • 工作流 Agent 对话:chat.json
  • 兼容旧格式
    • 首次访问时自动把旧的 workflows/<id>.json 扁平文件迁移到目录格式

外部边界:

  • WebSocket:手动执行、交互响应、客户端节点执行、恢复查询
  • HTTP:hook 触发
  • node-cron:定时调度
  • 文件系统:workflow 定义、版本、执行历史、插件配置
  • 插件运行时:pluginService
  • 客户端插件桥接:ClientNodeManager
  • AI/代码/数据库能力:通过节点分派调用到各自服务

Trigger Registration Lifecycle

WorkflowTriggerService 的职责仅是“注册触发器并把触发交给 ExecutionManager”,自己不保存执行态。

主要行为:

  • start()
    • 服务启动时扫描 store.listWorkflows(),批量注册全部 trigger
  • reloadWorkflow(workflowId)
    • workflow 变更后清理并重建该 workflow 的 trigger
  • removeWorkflow(workflowId)
    • workflow 删除后移除 trigger
  • getHookBindings(hookName)
    • 供 hook route 反查绑定关系
  • validateCron(cronExpr)
    • 用于校验 cron 配置并返回后续触发时间预览

Hook 触发器本质上只是 hookName -> {workflowId, triggerId}[] 的内存索引。

  • packages/server/src/services/execution-manager.ts
    • 执行主循环、节点分派、恢复态、落盘逻辑都在这里
  • packages/server/src/services/workflow-trigger-service.ts
    • cron/hook 触发注册和重载逻辑
  • packages/server/src/services/interaction-manager.ts
    • 阻塞式 UI 节点的暂停与恢复
  • packages/server/src/storage/workflow-store.ts
    • workflow 目录结构和执行历史保留策略
  • packages/server/src/ws/execution-channels.ts
    • 手动执行入口与控制命令
  • packages/server/src/routes/workflow-hook.ts
    • 外部系统通过 HTTP 驱动 workflow 的入口
  • packages/server/src/services/execution-node-helpers.ts
    • 拓扑排序、客户端插件节点识别、输入构造

Open Questions / Risks

  • 已修复的问题:ExecutionManager.emitEvent() 现在会把 session.workspaceId 传给 deps.emit()packages/server/src/app.ts 的全局广播只在存在显式 workspaceId 时才调用 broadcastToWorkspace()
    • 这修正了原先“把 workflowId 错传给 broadcastToWorkspace()”的问题。
    • 结果是:来自 WS 的执行会按真实连接 scope 广播;Hook/SSE 执行继续走 eventSink,不依赖全局广播。
  • 仍然存在的设计缺口:系统依旧缺少统一的“workflow execution event scope”模型。
    • workflow 模板本身没有 workspaceId 字段。
    • 前端同样存在两套归属约定:编辑器使用 workspaces[0]?.id 建立 WS,分享页固定使用 getWS('workflows')
    • Cron 触发默认没有前端受众 scope,因此不会进入全局广播。
    • 结论:这次修复解决了错参 bug,但没有解决 execution scope 的长期模型问题。
  • 当前恢复机制只保留内存态和最终 execution log,没有持久化 session checkpoint。
    • 进程重启后无法从中间节点继续跑,只能恢复展示或重新触发。
  • WorkflowTriggerService 的 hook/cron 注册状态是纯内存的。
    • 依赖服务启动时重新扫全量 workflow;如果未来要做分布式部署,需要改成集中调度或外部注册中心。