任务编排
一轮只被 produce 一次,写进按 taskId 键的 Redis stream buffer——live 响应和断线重连读的是同一个 buffer,不是两条代码路径,而是同一份流的两个读者。
编排器只 produce 一次,所有人读回
编排层其实是两个文件,按关注点切开。传输壳 apps/server/src/services/task-orchestrator.ts 很薄,只管一轮的流怎么送达、断了怎么恢复;真正的执行核心是 apps/server/src/services/task-runner.ts(startTaskUiStream),它和传输、和宿主进程都解耦——正因如此,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>,封帧交给调用方。
编排器从不碰 LLM 机制;agent 流内部的一切见 Agent Engine 及其子系统页。
请求路径:先 preflight,再入队
execute(taskId, user, body) 要尽快返回一个 HTTP Response,但不亲自跑这一轮——真正的执行留给入队后的 run。所以它只做传输形状的活,四步:
- 清掉陈旧 abort flag——
clearAbortRequest(taskId),放在任何 DB 往返之前的最前面,让本轮 Stop(在客户端“提交中”阶段按下的)落在清除之后,从而存活到 run 启动时被兑现。 - 在请求路径上 preflight——
preflightTask(deps, …):getTaskForExecution(归属 / not-found → HTTP 404)、creditService.checkQuota(→ HTTP 402),以及关键的一步——saveReceivedMessage在这里把入站消息落库。一个入了队却永远没跑的 job 绝不能把用户这轮悄悄弄丢,所以这次持久写留在请求路径上,不放进 run 里。 - 把 run 入队——
deps.jobQueue.enqueue("task.run", …, () => runTaskToStore(deps, …))。BullMQ 下交给 worker,inline 队列下在本进程跑;无论哪种,run 都是 produce 进 buffer。 - 读 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 响应)。
急(流之前):
- 拿锁——
acquireTaskLock(taskId)(Redis,已被占 → 409)+startLockHeartbeat续约 TTL,进程死亡时锁自然过期,不会把任务卡死。 - 绑定 task 行——再取一次
getTaskForExecution(model / approval / planning 都是任务绑定、跨轮不可变的,因为它们塑造缓存的 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、provider keys、agent config、subagent defs、沙箱就绪。- 记忆服务 +
loadIndex。 createRuntimeContext(...)——在完整工具 loadout 确定后一次性构建。- 组装
toolServices(记忆、browser bridge、kanban、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
流一旦彻底 drain 完,才轮到收尾。五步:
- 盖终态判决——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;失败只记日志,绝不致命。(这是服务端的名字;桌面端agent-handler.ts把同一步叫finalizeTurnCompaction。)saveExecutionRecords——inline 写,不入队,因为 admin UI 要立即用。- 入队三个后台 job——
credit.consume、resource.index、memory.extraction(经 JobQueue fire-and-forget;BullMQ 下在 worker,inline 兜底下在本进程)。
cleanup——onEnd 的 finally
无条件,成功失败都跑(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 点中断。
闭包隔离与及时 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 返回那刻一起落地。
桌面对照
agent-handler.ts 镜像核心的形状,但去掉 HTTP 相关的关注点:
| 关注点 | 服务端 | 桌面 |
|---|---|---|
| 分布式锁 | Redis SET NX + heartbeat | 不需要(单进程) |
| 额度检查 | preflight + run 里查 | 不需要 |
| produce / 读回 | Redis stream buffer | 不需要(IPC 一直在) |
| 跨进程 abort | Redis abort-bus | 只用进程内 abortManager |
| job 队列 | BullMQ on Redis | inline fire-and-forget |
| 轮末锚 | saveTurnAnchor | finalizeTurnCompaction |
阶段形状与闭包隔离纪律完全一致——GC 成本在桌面端最显眼,因为 job 是 inline 跑的。
相关文档
- Agent Engine——核心承载的那个循环
- 上下文压缩——part 寻址 checkpoint,以及
saveTurnAnchor写的跨轮锚 - 流式架构——可恢复 SSE buffer、produce/consume 拆分、WS 传输
- 后台任务队列——BullMQ vs inline,以及三个轮末 job
- 记忆系统——记忆沙箱与
memory.extractionjob