任务编排
编排器只 produce 一次,所有人读回
编排层按关注点切开。传输壳 apps/server/src/services/task-orchestrator.ts 很薄,只管一轮的流怎么送达、断了怎么恢复;server 适配器是 apps/server/src/services/task-runner.ts(startTaskUiStream),真正的执行核心是 packages/core/src/services/agent/run-task-turn.ts(runTaskTurn)。正因如此,API 进程、BullMQ worker 和 desktop 都能用不同端口驱动同一份 turn body。桌面端在 apps/desktop/src/main/handlers/agent-stream.ts 注入自己的本地子集。
这套设计真正承重的地方,是只有一个读核。在 SSE 路径上,run 不把流直接推给客户端,而是把自己产出的 UIMessageChunk produce 进一个按 taskId 键的 Redis stream buffer,再由 execute() 读这个 buffer 作响应。客户端断线后调 resumeStream(taskId),读的还是同一个 buffer、同一种读法。所以 live 和 resume 从来不是两条代码路径,而是同一份已 produce 流的两个读者。run 自己则始终不碰传输——它只返回一个传输中立的 ReadableStream<UIMessageChunk>,封帧交给调用方。
编排器从不碰 LLM 机制;agent 流内部的一切见 Agent Engine 及其子系统页。
设计判断
请求路径只做那些“job 可能还没开始就必须已经完成”的事:确认归属、检查额度、持久化入站消息、重置流 buffer。其余工作进入 run 本身,这样重试、worker 执行、desktop inline 执行都共享同一套语义。
stream buffer 也是产品判断。live delivery 和 resume 不是两套实现,而是同一份已产出 turn 的两个读者。断线恢复因此属于 task runtime 的能力,而不是浏览器客户端上的补丁。
请求路径:先 preflight,再入队
execute(taskId, user, body) 要尽快返回一个 HTTP Response,但不亲自跑这一轮——真正的执行留给入队后的 run。它只做传输形状的工作,五步:
- 清掉陈旧 abort flag——
clearAbortRequest(taskId),放在任何 DB 往返之前的最前面。这个POST在用户还来不及按 Stop 时就已发出,把清除搁在最顶端,本轮 Stop(在客户端“提交中”阶段按下的)的publishAbort才落在清除之后,从而存活到 run 启动时被兑现。 - 在请求路径上 preflight——
preflightTask(deps, …):getTaskForExecution(归属 / not-found → HTTP 404)、creditService.checkQuota(→ HTTP 402),以及关键的一步——saveReceivedMessage在这里把入站消息落库。一个入了队却永远没跑的 job 绝不能把用户这轮悄悄弄丢,所以这次持久写留在请求路径上,不进 run。 - 重置上一轮的 buffer——
streamBuffer.reset(taskId)。task 的流按taskId跨轮复用同一个 key,若不先重置,本轮新挂上的读者会把上一轮仍在 TTL 内的 chunk 回填出来——等于把上一条回答重播进这一轮。生产者自己那次 reset 跑得太晚(在 worker 里,读者早已读过),所以重置放在返回读者之前、请求路径上。 - 把 run 入队——
deps.jobQueue.enqueue("task.run", …, () => runTaskToStore(deps, …))。BullMQ 下交给 worker,inline 队列下在本进程跑;无论哪种,run 都是 produce 进 buffer。 - 读 buffer 作响应——
streamBuffer.read(taskId),和重连(resumeStream)读的是同一份 buffer。
无 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 响应)。
急(流之前):
- 拿锁——
acquireTaskLock(taskId)(Redis,已被占 → 409)+startLockHeartbeat续约 TTL,进程死亡时锁自然过期,不会把任务卡死。 - 绑定 task 行——再取一次
getTaskForExecution(model / approval / planning / knowledge 四项都是任务绑定、跨轮不可变的:它们塑造缓存的 prompt 前缀,所以创建时锁定、绝不逐消息读,否则跨轮就会破坏缓存与上下文窗口的假设)+creditService.checkQuota(run 要自洽,worker 会重查一遍)。 - 加载任务数据——
loadTaskData(taskId)→ 裸uiMessages、一个新的assistantMessageId,以及oldAssistantMessage(resume 时)。一轮起手只需裸流;引擎的buildTurnInput每轮从它重建压缩前缀(见 压缩)。 - 沙箱 + 上传——
createSandbox({ id: taskId }),再由seedUploadsIntoSandbox把用户本轮附的文件落进沙箱,好让read_file/ shell 够得到。 - session + abort 接线——建
ExecutionSession;abortManager.create(taskId)管进程内 Stop,加subscribeAbort(Redis abort-bus)让 Stop 在 run 处于 worker 时也能抵达,再加一次isAbortRequested复查以覆盖“订阅前”那道窗口。
drain 时(createUIMessageStream 的 execute 回调): setup 边跑边流出进度状态(agent_building → context_building → mcp_connecting → agent_running,见状态机):
- 并发发起
connectMcpTools(网络绑定,与其余 setup 重叠)。 buildAgentSetup——先取 tier,再并行取 agent config(含 provider key)、subagent defs、沙箱就绪。- 记忆服务 +
loadIndex。 createRuntimeContext(...)——在完整工具 loadout 确定后一次性构建。- 组装
toolServices(记忆、knowledge——经createScopedKnowledgeService只 scope 到本任务附加的文档、workbench、delivery、wait_and_resume、subagents);剥掉team工具(task 的请求-响应流没有 team 需要的 idle-wake 生命周期)。 - 压缩 repos(
toolCompaction/snapshot/taskBudget)+buildCompactionDeps——挂在 session 上,好让 finalize 写跨轮锚。 - await MCP 工具 +
tryAttachToolDiscovery;injectUserSkills处理/skill-name激活。 runAgentLoop(...)——返回一个普通对象AgentLoopResult { agentStream, stepUsages }。setup 阶段(模型吐字节之前)落下的 Stop 会在这里被捕获并吞掉,让流干净地收成一次 abort,而非错误。writer.merge(toAgentUIMessageStream(...))——chunk 开始流动。
finalize——onEnd → finalizeExecution
run 把 chunk 交给流之后,故事还没完——收尾要等,流一旦彻底 drain 完,onEnd 才触发 finalizeExecution。五步:
- 盖终态判决——abort 压过零星 error(迟到的 Stop 可能触发一次伪
onError);服务器关停引发的 abort 记为interrupted而非aborted,把重启抖动挡在“用户中止”指标之外。在途的孤儿工具 part 被收成终态,免得前端永久 shimmer。 taskService.finalizeTurn(...)——持久化 assistant 消息、累计用量、更新任务 metadata,返回{ roundUsage, lifecycle }。一条task:eventWS 消息把终态 kind 告诉客户端。saveTurnAnchor({ ctx, deps, repos })——把跨轮真实用量锚(meta:anchor)写进task_compaction_snapshot,供下一轮 preflight 种入。它从不写task_message、也不入队任何 job;失败只记日志,绝不致命。server 与 desktop 当前都通过finalizeTaskTurn走到这一步。saveExecutionRecords——inline await 写,不入队:它只是一次廉价的本地遥测落库(token 用量、每步task_step_info、子 agent 执行记录);重活与外部副作用才留给下面的 job。- 入队三个后台 job——
credit.consume、resource.index、memory.extraction(经 JobQueue fire-and-forget;BullMQ 下在 worker,inline 兜底下在本进程)。
cleanup——onEnd 的 finally
finalize 把这一轮的成果落库,cleanup 则把这一轮占用的资源交还——无条件,成功失败都跑(setup 抛错的 catch 里也镜像一份):unsubscribeAbort → abortManager.remove(taskId) → mcpClientManager.disconnect(taskId) → stopLockHeartbeat() → releaseTaskLock(taskId)。
跨回调的实体
createUIMessageStream 由回调驱动,execute / onEnd / onError 之间不传返回值,所以一个可变的 ExecutionSession 是它们唯一的通道。
| 实体 | 携带什么 | 生命周期 |
|---|---|---|
ExecutionSession | abortController、streamStartedAt、stepUsages,以及(在 execute 里写入的)resolvedModelId、context、compactionDeps、compactionRepos、memoryService;error 由 onError 写 | 急切创建,onEnd 读 |
RuntimeContext | 平台无关的 agent 环境——writer、sandbox、memorySandbox、subagentDefs、todos、reminders,以及 write() / writeTransient() 辅助方法 | 在 execute 里一次性构建,用到 finalize |
CompactionDeps | 共享的压缩 bundle(createModel、contextWindow、compactionModel)——挂 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/abort 与 task:stream:abort)按一个刻意的顺序跑:
taskService.abort——归属校验 +isActive = false。未授权访问会抛错,所以攻击者无法枚举taskId去取消别人的 run。把它放第一步,未授权的调用者根本到不了第二步。abortManager.abort(taskId)(进程内)+publishAbort(taskId)(Redis abort-bus)。run 可能在本进程、也可能在 worker,所以两个都发;abort 是幂等的。run 里(startTaskUiStream装的)subscribeAbort收到 bus 消息;isAbortRequested复查覆盖“抢在订阅之前”的那次 Stop。信号被穿进runAgentLoop,在下一个 await 点中断。
两步跑完,编排器不等 agent 慢慢收场,立即向用户推一条 task:event(kind: "aborted")WS 消息。finalize 时还会再推一次等价的终态,但那只是无副作用的重复:aborted 不弹 toast,客户端至多多 invalidate 一次列表。
闭包隔离与及时 GC
一次 run 的对象图很重——消息历史、step runner、StreamTextResult 全挂在上面。三条手法让它在该释放的那一刻立刻释放、不被闭包钉住:
| 手法 | 效果 |
|---|---|
ExecutionSession 字段保持最少 | onEnd 不读的都留作 execute 的局部(如 sandbox) |
| 返回普通结果,而非闭包住整个 run | runAgentLoop 返回 { agentStream, stepUsages };编排器不持有伸进 run 的引用,丢掉局部即释放 StreamTextResult + step-runner 链 |
| 长生命闭包经模块级工厂构建 | 例如 title 的 .then 闭包只捕获 taskId + writer,够不到外层作用域 |
BullMQ 下 run 的闭包立即被丢(payload 序列化进 Redis,worker 重建自己的 deps);inline 兜底下,正是这套纪律让 runResult + step runner + 消息历史在 onEnd 返回那刻一起落地。
桌面对照
apps/desktop/src/main/handlers/agent-stream.ts 是同一个 runTaskTurn core 的 desktop 宿主。它通过 loopback HTTP 路由 POST /api/tasks/:id/messages 提供 task streaming,因此 renderer 像 Web 一样消费 SSE,而不是旧的 agent:run IPC 路径。
| 关注点 | 服务端 | 桌面 |
|---|---|---|
| 分布式锁 | Redis SET NX + heartbeat | 不需要(单进程) |
| 额度检查 | preflight + run 里查 | 不需要 |
| produce / 读回 | Redis stream buffer | 不需要(loopback SSE 直连) |
| 跨进程 abort | Redis abort-bus | 只用进程内 abortManager |
| job 队列 | BullMQ on Redis | inline fire-and-forget |
| 轮末锚 | finalizeTaskTurn 内调用 saveTurnAnchor | finalizeTaskTurn 内调用 saveTurnAnchor |
阶段形状与闭包隔离纪律完全一致——GC 成本在桌面端最显眼,因为 job 是 inline 跑的。
相关文档
- Agent Engine——核心承载的那个循环
- 上下文压缩——part 寻址 checkpoint,以及
saveTurnAnchor写的跨轮锚 - 流式架构——可恢复 SSE buffer、produce/consume 拆分、WS 传输
- 后台任务队列——BullMQ vs inline,以及三个轮末 job
- 记忆系统——记忆沙箱与
memory.extractionjob