Activepieces 架构决策解读:Engine 将运行时回调直连 App,移除 Worker 中转链路

📅 发布时间:2026/9/12 15:53:40
Activepieces 架构决策解读:Engine 将运行时回调直连 App,移除 Worker 中转链路
Activepieces 架构决策解读Engine 将运行时回调直连 App移除 Worker 中转链路【免费下载链接】activepiecesAI Agents MCPs AI Workflow Automation • (~400 MCP servers for AI agents) • AI Automation / AI Agent with MCPs • AI Workflows AI Agents • MCPs for AI Agents项目地址: https://gitcode.com/GitHub_Trending/ac/activepieces本文依据仓库决策记录 brain/knowledge/decisions/000003-engine-posts-run-time-callbacks-directly-to-the-app.md 展开结合 Engine 侧客户端、App 侧控制器 与 回调服务 等源码实现解读 Activepieces 中引擎Engine与 App 之间的运行时回调通道设计四个运行时回调如何从Engine → Worker → App的双跳链路收敛为Engine → App的直连 HTTP 通道以及uploadRunLog为何刻意保留双源。决策摘要四条回调直连 App决策记录status: accepted明确了 Activepieces 引擎与主应用之间的数据通路收敛方案Engine 将全部四个运行时回调updateRunProgress、updateStepProgress、sendFlowResponse、uploadRunLog通过 HTTP 直接发送给 App请求路径为internalApiUrlengineToken鉴权下的POST /v1/engine/*使用 ENGINE 主体principal。Worker 从数据通路中被移除Engine → Worker 的中转链路被删除。这一决策的核心是减少一跳此前 Engine 通过 Socket.IO 将回调发给 WorkerWorker 原封不动地再转发给 App形成纯转发、无增值的 1:1 中继a 1:1 relay adding no value。而 Engine 本来就通过同一 HTTP 通道与 App 交互 storeKV 存储、files文件与 connections连接运行时回调并入该通道顺理成章。背景从双跳中继到单一通道决策文档的 Context 部分还原了历史形态历史链路Engine 通过 Socket.IO 将回调发送给 WorkerWorker 再把每条回调逐字转发给 App——两跳之间是 verbatim 转发中继不附加任何业务价值已有先例Engine 与 App 之间已经存在直接 HTTP 通信store / files / connections 均走该通道回调并入同一通道无需引入新的传输机制。从当前仓库代码可以验证同一通道的说法Engine 侧对 store、文件、连接的回调分别封装在 store.ts、file-uploader.ts、connection-resolver.ts 中而运行时回调集中在 engine-run-api.ts——它们是并行存在的同一批 HTTP 客户端工具。动机统一本地与远程运行时不新增 Socket 鉴权面决策文档的 Why 部分给出了两条关键理由统一本地与远程运行时unifies local and remote runtimes无论 Engine 运行在本地sandbox还是远端 worker poolEngine 都以相同方式触达 App。直连 HTTP 消除了对Engine 与 Worker 同处一个 Socket.IO 会话这一隐含假设的依赖使运行时的部署形态与回调路径解耦保持 HTTP而非新建 Engine → App 的 Socket四个回调中三个是低频事件为低频通道引入一条常驻 Socket 连接不划算更重要的是Socket 层目前只接受 USER / WORKER 两类主体principal若让 ENGINE 走 Socket 就等于为 ENGINE 新增一个鉴权面auth surfacefor no present benefit——即在没有当下收益的情况下扩大了安全攻击面。因此该决策本质上是一次安全面收敛 链路简化把回调从需要为 ENGINE 开放 Socket 鉴权的潜在需求中摘除统一收编进已经为 ENGINE 开放的 HTTP 通道。实现印证App 侧的四个路由与服务控制器注册POST /v1/engine/*App 侧入口位于 packages/server/api/src/app/workers/engine-controller.ts。该 Fastify 插件注册了四个回调路由全部使用securityAccess.engine()安全配置HTTP 路由回调控制器实现POST /v1/engine/run-progressupdateRunProgressengineRunCallbackService.updateRunProgressPOST /v1/engine/step-progressupdateStepProgressengineRunCallbackService.updateStepProgressPOST /v1/engine/run-logsuploadRunLogengineRunCallbackService.uploadRunLogPOST /v1/engine/flow-responsesendFlowResponseengineRunCallbackService.sendFlowResponse控制器还注册了GET /populated-flows、GET /flows、GET /pieces/bundle等 Engine 拉取流版本与 piece 包的只读接口其中/pieces/bundle会显式校验request.principal.type ! PrincipalType.ENGINE并返回 401印证了 ENGINE 主体在 HTTP 层的独立鉴权身份。安全配置ENGINE 主体在 packages/server/api/src/app/core/security/authorization/fastify-security.ts 中可以看到engine()的实现/** * Creates a security configuration for routes that are only accessible to Engine principal. * * Effects: * - projectId field is available on the request.principal because the engine principal token contains projectId */ engine: () { return securityAccess.unscoped([PrincipalType.ENGINE]) },也就是说POST /v1/engine/*只接受携带合法 engineTokenBearer的 ENGINE 主体请求且请求主体principal中直接携带projectIdApp 侧据此把回调归入正确的项目作用域——这正是决策中ENGINE principal在代码层面的落点。统一回调服务一个服务支撑两个入口决策文档的 Consequences 特别强调One app-side service backs both entry points即 Engine 的 HTTP 回调与 Worker 的 Socket 回调共用同一个 App 侧服务。该服务位于 packages/server/api/src/app/flows/flow-run/engine-run-callback-service.ts各回调在 App 内的实际行为如下updateRunProgress通过websocketService.to(projectId)向项目广播UPDATE_RUN_PROGRESS事件把运行进度实时推送到前端运行详情页updateStepProgress向项目广播TEST_STEP_PROGRESS事件用于步骤级进度尤其测试执行场景的实时展示sendFlowResponse通过pubsub.publish(engine-run:sync:${request.workerHandlerId}, ...)发布消息把同步 HTTP 请求的响应httpRequestIdrunResponse送回等待中的调用方uploadRunLog核心元数据落库逻辑——校验运行状态是否终态isFlowRunStateTerminal终态时通过ensureLogsFileExists以 ZSTD 压缩写入运行日志文件FileType.FLOW_RUN_LOG随后把status、logsFileId、failedStep、startTime/finishTime、provisionMs/bootMs/runMs等元数据投递到runsMetadataQueue若状态命中FAILED_RUN_SYNC_STATUSESFAILED、INTERNAL_ERROR、TIMEOUT、MEMORY_LIMIT_EXCEEDED、LOG_SIZE_EXCEEDED且携带workerHandlerId/httpRequestId还会代为回送一个 500 响应The flow has failed and there is no response returned。从源码结构看uploadRunLog承担了运行元数据 日志文件 失败同步响应三合一职责是四个回调中业务最重、也因此被保留双源的一个。实现印证Engine 侧的 HTTP 客户端Engine 侧的调用集中在 packages/server/engine/src/lib/api/engine-run-api.ts四个方法一一对应四个路径async updateRunProgress({ apiUrl, engineToken, request }) { await post({ apiUrl, engineToken, path: run-progress, body: request }) }, async updateStepProgress({ apiUrl, engineToken, request }) { await post({ apiUrl, engineToken, path: step-progress, body: request, fetcher: global.fetch }) }, async uploadRunLog({ apiUrl, engineToken, request }) { await post({ apiUrl, engineToken, path: run-logs, body: request }) }, async sendFlowResponse({ apiUrl, engineToken, request }) { await post({ apiUrl, engineToken, path: flow-response, body: request }) }底层post统一拼装请求const response await fetcher(${apiUrl}v1/engine/${path}, { method: POST, headers: { Content-Type: application/json, Authorization: Bearer ${engineToken}, }, body: JSON.stringify(body), }) if (!response.ok) { throw new EngineGenericError( EngineRunCallbackError, Failed to POST ${path}: ${response.status} ${response.statusText}, ) }值得注意的实现细节默认使用retryFetch作为传输函数updateStepProgress例外地使用global.fetch失败时抛出EngineGenericError(EngineRunCallbackError, ...)——回调通道具备重试与可观测的失败语义apiUrl与engineToken的来源在 engine-constants.ts 中internalApiUrl必须以斜杠结尾否则构造期即抛出InternalApiUrlNotEndsWithSlashErrorEngine 还用它拼出GET /v1/worker/project等端点回调的实际触发点可在 flow-run-progress-reporter.ts 中看到Engine 在运行过程中调用updateStepProgress、updateRunProgress、uploadRunLog与决策文档中四个回调的名称一一对应。刻意的不对称uploadRunLog 为何双源这是该决策中最值得读者注意的设计取舍。决策文档的 Consequences 部分明确指出uploadRunLog是双源的dual-sourced同时保留在 Worker 的 Socket 通道上Worker 仍然以 WORKER 主体发起它用于记录 Engine 自己无法上报的终态crash、OOM、INTERNAL_ERROR。App 侧由同一个服务支撑两个入口。这种不对称是刻意为之——否则未来的读者会疑惑四个回调中为什么有一个还活在 Worker Socket 上。代码印证了这一点。在 packages/server/worker/src/lib/execute/jobs/execute-flow.ts 中正常路径下Worker 在 Engine 执行结束后会追加一次无状态的uploadRunLog仅携带provisionMs / bootMs / runMs计时数据代码注释说明这是 runs 页面的延迟分解复用 run-log 元数据上传通道合并 Engine 自身报告无法测量的计时异常路径下Worker 通过reportFlowStatus以uploadRunLog上报TIMEOUT、MEMORY_LIMIT_EXCEEDED、LOG_SIZE_EXCEEDED、INTERNAL_ERROR含RunInternalErrorSource.ENGINE的内部错误详情并携带logsFileId以便 App 把 internalError 持久化进日志文件命中INTERNAL_ERROR且为专用 Worker 时还会触发 on-call 分页告警。换言之Engine 直连 HTTP 覆盖运行中的高频/常规回调进度、步骤、同步响应、常规日志Worker 的 Socket 通道只兜底Engine 已死、无法自报的终态。Sandbox 崩溃、内存超限OOM、超时等场景下 Engine 进程可能已经无法发出 HTTP 请求此时唯一可靠的观测者是宿主的 Worker——这正是双源存在的理由。权衡与启示链路缩短带来的一致性与可观测性回调从两跳变一跳减少了消息丢失面与转发延迟也让回调是否送达的责任边界更清晰——失败即抛出EngineRunCallbackError可由 Engine 侧重试鉴权面最小化不为低频回调引入 ENGINE-on-Socket 的鉴权路径全部收敛到已有的 ENGINE HTTP 通道符合最小化攻击面的原则失败场景必须由外层兜底任何直连设计都需要回答调用方死了谁上报——本决策用 Worker 保留uploadRunLog单点兜底给出了答案且刻意在代码与文档中标注这种不对称避免后人修复它时引入回归一个服务、两个入口App 侧不做双份实现engineRunCallbackService同时服务 Engine HTTP 回调与 Worker Socket 回调保证两条链路的落库行为完全一致。相关代码索引决策原文brain/knowledge/decisions/000003-engine-posts-run-time-callbacks-directly-to-the-app.mdEngine 侧回调客户端packages/server/engine/src/lib/api/engine-run-api.tsEngine 侧回调触发点packages/server/engine/src/lib/helper/flow-run-progress-reporter.tsApp 侧控制器packages/server/api/src/app/workers/engine-controller.tsApp 侧统一回调服务packages/server/api/src/app/flows/flow-run/engine-run-callback-service.tsENGINE 鉴权配置packages/server/api/src/app/core/security/authorization/fastify-security.tsWorker 侧终态兜底上报packages/server/worker/src/lib/execute/jobs/execute-flow.tsEngine 常量与 URL 校验packages/server/engine/src/lib/handler/context/engine-constants.ts【免费下载链接】activepiecesAI Agents MCPs AI Workflow Automation • (~400 MCP servers for AI agents) • AI Automation / AI Agent with MCPs • AI Workflows AI Agents • MCPs for AI Agents项目地址: https://gitcode.com/GitHub_Trending/ac/activepieces创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考