任务编排

一轮只被 produce 一次,写进按 taskId 键的 Redis stream buffer——live 响应和断线重连读的是同一个 buffer,不是两条代码路径,而是同一份流的两个读者。

编排器只 produce 一次,所有人读回

编排层其实是两个文件,按关注点切开。传输壳 apps/server/src/services/task-orchestrator.ts 很薄,只管一轮的流怎么送达、断了怎么恢复;真正的执行核心是 apps/server/src/services/task-runner.tsstartTaskUiStream),它和传输、和宿主进程都解耦——正因如此,API 进程和 BullMQ worker 才能跑同一份代码。(桌面端在 apps/desktop/src/main/handlers/agent-handler.ts 里镜像这套核心。)

这套设计真正承重的地方,是只有一个读核。在 SSE 路径上,run 不把流直接推给客户端,而是把自己产出的 UIMessageChunk produce 进一个按 taskId 键的 Redis stream buffer,再由 execute() 读这个 buffer 作响应。客户端断线后调 resumeStream(taskId),读的还是同一个 buffer、同一种读法。所以 live 和 resume 从来不是两条代码路径,而是同一份已 produce 流的两个读者。run 自己则始终不碰传输——它只返回一个传输中立的 ReadableStream<UIMessageChunk>,封帧交给调用方。

HTTP 请求 Task Orchestrator Agent Engine Sandbox server-only @zapvol/backend 锁、额度、Redis 恢复 后台任务、MCP 生命周期 ReAct 循环、Prompt 组装 上下文压缩、工具执行 文件系统、Shell 代码执行 架构边界

编排器从不碰 LLM 机制;agent 流内部的一切见 Agent Engine 及其子系统页。

请求路径:先 preflight,再入队

execute(taskId, user, body) 要尽快返回一个 HTTP Response,但不亲自跑这一轮——真正的执行留给入队后的 run。所以它只做传输形状的活,四步:

  1. 清掉陈旧 abort flag——clearAbortRequest(taskId),放在任何 DB 往返之前的最前面,让本轮 Stop(在客户端“提交中”阶段按下的)落在清除之后,从而存活到 run 启动时被兑现。
  2. 在请求路径上 preflight——preflightTask(deps, …)getTaskForExecution(归属 / not-found → HTTP 404)、creditService.checkQuota(→ HTTP 402),以及关键的一步——saveReceivedMessage 在这里把入站消息落库。一个入了队却永远没跑的 job 绝不能把用户这轮悄悄弄丢,所以这次持久写留在请求路径上,不放进 run 里。
  3. 把 run 入队——deps.jobQueue.enqueue("task.run", …, () => runTaskToStore(deps, …))。BullMQ 下交给 worker,inline 队列下在本进程跑;无论哪种,run 都是 produce 进 buffer
  4. 读 buffer 作响应——streamBuffer.read(taskId),和重连走的是同一次读。

无 Redis 兜底(dev / 桌面)整个跳过 buffer:startTaskUiStream inline 跑、直接流回。executeWs(taskId, userId, body, send) 是 WebSocket 封装——驱动同一个 startTaskUiStream,把每个 chunk 作 task:stream:frame 发出;不走 Redis buffer,WS 客户端断线后靠重新拉取 DB 历史恢复。

run 本身:startTaskUiStream

这才是真正跑一轮的地方,返回一个传输中立的流。它的节奏是调用时急、drain 时懒——下面的 setup 在函数被调用的当下就跑;而 agent 循环与 finalize 要等返回的流被 drain 时才跑(drain 它的可能是 buffer 生产者、WS reader,或 inline SSE 响应)。

急(流之前):

  1. 拿锁——acquireTaskLock(taskId)(Redis,已被占 → 409)+ startLockHeartbeat 续约 TTL,进程死亡时锁自然过期,不会把任务卡死。
  2. 绑定 task 行——再取一次 getTaskForExecution(model / approval / planning 都是任务绑定、跨轮不可变的,因为它们塑造缓存的 prompt 前缀)+ creditService.checkQuota(run 要自洽,worker 会重查一遍)。
  3. 加载任务数据——loadTaskData(taskId) → 裸 uiMessages、一个新的 assistantMessageId,以及 oldAssistantMessage(resume 时)。一轮起手只需裸流;引擎的 buildTurnInput 每轮从它重建压缩前缀(见 压缩)。
  4. 沙箱 + 上传——createSandbox({ id: taskId }),再由 seedUploadsIntoSandbox 把用户本轮附的文件落进沙箱,好让 read_file / shell 够得到。
  5. session + abort 接线——建 ExecutionSessionabortManager.create(taskId) 管进程内 Stop,加 subscribeAbort(Redis abort-bus)让 Stop 在 run 处于 worker 时也能抵达,再加一次 isAbortRequested 复查以覆盖“订阅前”那道窗口。

drain 时(createUIMessageStreamexecute 回调): setup 边跑边流出进度状态(agent_buildingcontext_buildingmcp_connectingagent_running,见状态机):

  1. 并发发起 connectMcpTools(网络绑定,与其余 setup 重叠)。
  2. buildAgentSetup——并行解析 tier、provider keys、agent config、subagent defs、沙箱就绪。
  3. 记忆服务 + loadIndex
  4. createRuntimeContext(...)——在完整工具 loadout 确定后一次性构建。
  5. 组装 toolServices(记忆、browser bridge、kanban、delivery、wait_and_resume、subagents);剥掉 team 工具(task 的请求-响应流没有 team 需要的 idle-wake 生命周期)。
  6. 压缩 repos(toolCompaction / snapshot / taskBudget)+ buildCompactionDeps——挂在 session 上,好让 finalize 写跨轮锚。
  7. await MCP 工具 + tryAttachToolDiscoveryinjectUserSkills 处理 /skill-name 激活。
  8. runAgentLoop(...)——返回一个普通对象 AgentLoopResult { agentStream, stepUsages }。setup 阶段(模型吐字节之前)落下的 Stop 会在这里被捕获并吞掉,让流干净地收成一次 abort,而非错误。
  9. writer.merge(toAgentUIMessageStream(...))——chunk 开始流动。

finalize——onEndfinalizeExecution

流一旦彻底 drain 完,才轮到收尾。五步:

  1. 盖终态判决——abort 压过零星 error(迟到的 Stop 可能触发一次伪 onError);服务器关停引发的 abort 记为 interrupted 而非 aborted,把重启抖动挡在“用户中止”指标之外。在途的孤儿工具 part 被收成终态,免得前端永久 shimmer。
  2. taskService.finalizeTurn(...)——持久化 assistant 消息、累计用量、更新任务 metadata,返回 { roundUsage, lifecycle }。一条 task:event WS 消息把终态 kind 告诉客户端。
  3. saveTurnAnchor({ ctx, deps, repos })——把跨轮真实用量锚(meta:anchor)写进 task_compaction_snapshot,供下一轮 preflight 种入。它从不写 task_message、也不入队任何 job;失败只记日志,绝不致命。(这是服务端的名字;桌面端 agent-handler.ts 把同一步叫 finalizeTurnCompaction。)
  4. saveExecutionRecords——inline 写,不入队,因为 admin UI 要立即用。
  5. 入队三个后台 job——credit.consumeresource.indexmemory.extraction(经 JobQueue fire-and-forget;BullMQ 下在 worker,inline 兜底下在本进程)。

cleanup——onEndfinally

无条件,成功失败都跑(setup 抛错的 catch 里也镜像一份):unsubscribeAbortabortManager.remove(taskId)mcpClientManager.disconnect(taskId)stopLockHeartbeat()releaseTaskLock(taskId)

跨回调的实体

createUIMessageStream 由回调驱动,execute / onEnd / onError 之间不传返回值,所以一个可变的 ExecutionSession 是它们唯一的通道。

实体携带什么生命周期
ExecutionSessionabortControllerstreamStartedAtstepUsages,以及(在 execute 里写入的)resolvedModelIdcontextcompactionDepscompactionReposmemoryServiceerroronError急切创建,onEnd
RuntimeContext平台无关的 agent 环境——writersandboxmemorySandboxsubagentDefstodosreminders,以及 write() / writeTransient() 辅助方法execute 里一次性构建,用到 finalize
CompactionDeps共享的压缩 bundle(createModelcontextWindowcompactionModel)——挂 session 供 memory.extraction job 用;引擎内部经 buildCompactionDeps 另建自己的一份execute 里建,由后台 job 消费

session 刻意保持窄:只有跨回调数据才上 session(sandbox 留作 execute 的局部,onEnd 从不读它)。这是三条闭包隔离手法之一,让每请求对象图在 onEnd 返回那一刻就能被 GC——见下。

abort——跨进程的真 Stop

真正的 Stop 要能穿过进程边界——run 可能在 API 进程、也可能在 worker。taskOrchestrator.abort(userId, taskId)(来自 POST /:id/aborttask:stream:abort)按一个刻意的顺序跑:

  1. taskService.abort——归属校验 + isActive = false。未授权访问会抛错,所以攻击者无法枚举 taskId 去取消别人的 run。把它放第一步,未授权的调用者根本到不了第二步。
  2. abortManager.abort(taskId)(进程内)+ publishAbort(taskId)(Redis abort-bus)。run 可能在本进程、也可能在 worker,所以两个都发;abort 是幂等的。run 里(startTaskUiStream 装的)subscribeAbort 收到 bus 消息;isAbortRequested 复查覆盖“抢在订阅之前”的那次 Stop。信号被穿进 runAgentLoop,在下一个 await 点中断。

闭包隔离与及时 GC

一次 run 的对象图很重——消息历史、step runner、StreamTextResult 全挂在上面。三条手法让它在该释放的那一刻立刻释放、不被闭包钉住:

手法效果
ExecutionSession 字段保持最少onEnd 不读的都留作 execute 的局部(如 sandbox
返回普通结果,而非闭包住整个 runrunAgentLoop 返回 { agentStream, stepUsages };编排器不持有伸进 run 的引用,丢掉局部即释放 StreamTextResult + step-runner 链
长生命闭包经模块级工厂构建例如 title 的 .then 闭包只捕获 taskId + writer,够不到外层作用域

BullMQ 下 run 的闭包立即被丢(payload 序列化进 Redis,worker 重建自己的 deps);inline 兜底下,正是这套纪律让 runResult + step runner + 消息历史在 onEnd 返回那刻一起落地。

桌面对照

agent-handler.ts 镜像核心的形状,但去掉 HTTP 相关的关注点:

关注点服务端桌面
分布式锁Redis SET NX + heartbeat不需要(单进程)
额度检查preflight + run 里查不需要
produce / 读回Redis stream buffer不需要(IPC 一直在)
跨进程 abortRedis abort-bus只用进程内 abortManager
job 队列BullMQ on Redisinline fire-and-forget
轮末锚saveTurnAnchorfinalizeTurnCompaction

阶段形状与闭包隔离纪律完全一致——GC 成本在桌面端最显眼,因为 job 是 inline 跑的。

相关文档

这页有帮助吗?