s13: 有状态 Agent — 双队列与状态机
把 agentLoop 包进一个类,让循环”有记忆”——能被叫停、能被改方向、还能在被叫停之后再被推一把。
... → s12 → s13
问题
s01 的 agentLoop 是个无状态纯函数:每次调用都从头传 messages 数组,跑完就还给调用方。s04 给它加了 AbortSignal 和 getSteeringMessages 回调,让循环可以被外部打断、被中途改方向——但回调只是一根”线”,没有”对象”承载状态。结果调用方 main 里堆的胶水越来越多:history 数组、controller 句柄、steering 闭包……全靠人肉维护。
真实 pi 需要一个”可复用的有状态 Agent”:
- 维护 messages:自己持有对话历史,调用方不需要传。
- 两个注入入口:
steer()在循环跑着的时候塞指令,followUp()在 agent 本该停下来的时候塞指令——两种注入时机完全不同。 - 支持 abort / waitForIdle:外界能硬停当前 run,也能等它跑完。
- 串行派发事件:listener 按订阅顺序
await,listener 自己的 promise 也要计入 run 的 settlement。
两个注入入口的差异是本章的核心认知:steering 是”循环内注入”(当前轮工具照常跑完,下一轮 LLM 调用前生效);followUp 是”循环外注入”(内层循环已经退出、agent 本该停止时,被它唤起再跑一轮)。
解决方案
把 s01 的 agentLoop 包进一个 MiniAgent 类,维护四样东西:
- 状态字段:
_messages/_systemPrompt/_listeners集合 - 双队列:
_steeringQueue与_followUpQueue,都是PendingQueue实例(mode 可选"all"批量排干或"one-at-a-time"逐条排干) - 状态机三件套:
_activeRun = { promise, resolve, abortController }——run 在则isStreaming=true,run 走完finishRun()清场 - 双层循环:内层处理 tool-call + steering,外层处理 followUp
关键认知:steering 和 followUp 不是两个语义不同的 API,而是同一个队列模式在循环的两个位置被排干。drain 的位置决定了”注入时机”:内层 while 的开头 drain → steering(下一轮 LLM 前);外层 while 内层退出后 drain → followUp(agent 即将停止时)。
工作原理
打开 code.ts,重点看四块。
第 1 块:PendingQueue。 一个 ~20 行的小容器,enqueue / hasItems / drain / clear 四个方法。drain() 的行为由 mode 决定:
drain(): ChatMessage[] {
if (this.mode === "all") {
const drained = this.messages.slice();
this.messages = [];
return drained;
}
const first = this.messages[0];
if (!first) return [];
this.messages = this.messages.slice(1);
return [first];
}
"all" 是批量排干(适合”用户连续敲了三句话,下次 LLM 一次性全看到”);"one-at-a-time" 是逐条排干(pi 默认值,每轮只看一条,模型可以独立回应每条指令)。两种 mode 只是 drain 的切片粒度不同,队列入队逻辑完全一样。
第 2 块:MiniAgent 类骨架。 状态字段 + getter + 状态机三件套:
class MiniAgent {
private _messages: ChatMessage[] = [];
private _systemPrompt: string;
private readonly _listeners = new Set<EventListener>();
private readonly _steeringQueue = new PendingQueue("one-at-a-time");
private readonly _followUpQueue = new PendingQueue("one-at-a-time");
private _activeRun?: ActiveRun;
// getter: messages / isStreaming
// steer(msg) / followUp(msg) / clearAllQueues()
// subscribe(listener) 返回 unsubscribe
// abort() / waitForIdle()
}
_activeRun 是状态机的全部——它存在即 isStreaming === true,它被 finishRun() 清空即回到 idle。abort() 调 abortController.abort(),waitForIdle() 返回 activeRun.promise——三件套各司其职。
第 3 块:prompt() → runWithLifecycle() → loop()。 这是双层循环的入口链:
async prompt(input: string): Promise<void> {
if (this._activeRun) throw new Error("Agent 正在运行...");
this._messages.push({ role: "user", content: input });
await this.runWithLifecycle((signal) => this.loop(this._messages.slice(), signal));
}
private async runWithLifecycle(executor) {
// 1. new AbortController + new Promise,把 resolve 存到外部变量
// 2. this._activeRun = { promise, resolve, abortController }
// 3. try { await executor(signal) } finally { this.finishRun() }
}
runWithLifecycle 是状态机的”开关”——开 run 时建三件套,run 走完(无论成功失败)finally 必然执行 finishRun()。这就是”异常视为 run 结束”的语义:catch 不接住错误往外抛,但 finally 仍然清场。
第 4 块:事件订阅 + processEvents 串行派发。 subscribe 把 listener 塞进 Set,返回 unsubscribe 函数。processEvents 是关键:
private async processEvents(event: AgentEvent): Promise<void> {
const signal = this._activeRun?.abortController.signal;
if (!signal) throw new Error("listener 在 active run 之外被调用");
for (const listener of this._listeners) {
await listener(event, signal);
}
}
串行 await 不是 forEach——forEach 会并发触发所有 listener 不等返回;for...of + await 才是按订阅顺序逐个等。listener 的 promise 计入 run settlement:waitForIdle() 返回的 activeRun.promise 要等到 finishRun() 才 resolve,而 finishRun() 在 try/finally 的 finally 里执行——所以最后一个 listener 的 await 跑完,run 才算结束。
运行
node code.ts
无需 API key——本章用 mock 流式 assistant 演示,重点是时序而非真实 LLM 调用。试这三个场景,观察两种注入点的时序差异:
- steering(循环内注入):输入”用工具列文件”,agent 开始流式输出 + 工具调用。运行中再输入
/steer 别动 test 目录,应看到紫色[steering]在turn_end之后、下一轮 LLM 之前被注入,当前轮工具调用照常执行完毕。 - followUp(循环外注入):输入一个简单问题(不含”工具”),agent 一轮就停。在它停止前输入
/follow 顺便补一个 README,应看到紫色[followUp]在内层循环退出后、agent_end之前被注入,agent 续跑一轮。 - abort:输入一个慢任务,运行中输入
/abort(或按 Ctrl+C),应看到 run 立刻终止,回到提示符。
前置概念清单
本章引入五个新概念:
- 双队列模式:steering 与 followUp 共享同一个
PendingQueue容器,只是 drain 的位置不同。两个队列、两种语义、一套数据结构。 - steering vs followUp:steering 是循环内注入(下一轮 LLM 调用前生效,当前轮工具照常跑完);followUp 是循环外注入(内层循环退出、agent 本该停止时生效,让循环续跑一轮)。差异在 drain 位置,不在 API 形状。
- 状态机:
idle ⇄ running,由_activeRun是否存在决定。run 在则isStreaming=true,run 走完finishRun()清场回到 idle。abort 走abortController.abort(),是 running 的一种特殊退出路径。 - 事件串行派发:listener 按订阅顺序
for...of+await,不是forEach并发。listener 的 promise 计入 run settlement——waitForIdle()等到最后一个 listener settle。 - 状态快照隔离:每次
prompt()调用前this._messages.slice()传给 loop——loop 拿到的是快照,即使 run 中途状态被外部改动,loop 内部用的是它启动时的那份。
源码锚点
mini-pi 的 MiniAgent,在pi 里是 packages/agent/src/agent.ts 的 Agent 类。配套双层循环在 packages/agent/src/agent-loop.ts。读完回答三个问题:
| 锚点 | 行号 | 看什么 |
|---|---|---|
Agent 类 | 166-557 | 整个有状态 Agent 的骨架,对照 MiniAgent |
MutableAgentState | 59-93 | 拷贝语义:tools/messages 的 getter/setter 都 .slice() |
PendingMessageQueue | 118-152 | 双队列共享容器,drain 两种 mode 实现 |
steer / followUp | 263-271 | 两个入口只是 enqueue,差异在 drain 位置 |
runWithLifecycle | 451-474 | 状态机三件套,try/finally finishRun |
finishRun | 494-500 | 清场:resolve promise + 清空 activeRun |
processEvents | 509-556 | switch 更新 state + for...of await 串行派发 |
| drain 位置(agent-loop.ts) | 167 / 253 / 257 | steering 在循环开始前 + 内层每轮末尾;followUp 在内层退出后 |
agent.ts第 59-93 行MutableAgentState的 tools 和 messages 都用 getter/setter 包了.slice(),为什么赋值时也要 copy?(提示:外部数组随时可能被原地改,agent 内部状态必须独立)agent.ts第 451-474 行runWithLifecycle的try/catch/finally——为什么 catch 里调handleRunFailure模拟一个 failure message,而 mini-pi 直接让异常外抛?(提示:pi 要保证 listener 一定能收到agent_end事件,mini-pi 用 finally 清场就够了)agent-loop.ts第 167 行在循环开始前就调了一次getSteeringMessages,第 253 行在内层每轮末尾再调一次,第 257 行在外层退出前调getFollowUpMessages——这三个 drain 位置为什么不能合并成一个?(提示:用户在三个时间窗口都可能敲字:启动到首次 LLM 之间、两轮工具调用之间、agent 即将停止时)
动手任务(改 mini-pi):把 PendingQueue 的 mode 从 "one-at-a-time" 改成 "all",观察 /steer 连续输入三条消息时的行为差异——批量注入 vs 逐条注入。再想想:什么场景下 "all" 比 "one-at-a-time" 更合适?(提示:用户连珠炮式敲字 vs 每条指令需要独立回应)
妥协清单
mini-pi 比pi 少做了什么,以及为什么省略是安全:
| 省略项 | pi 的做法 | 为什么本章可以省 |
|---|---|---|
continue() 外部入口 | 从最后一条消息续跑,区分 user/toolResult/assistant 三种角色 | 本章 REPL 一次提问走完一个 run,没有”从中间续跑”的场景 |
prepareNextTurn 钩子 | 每轮后可换 model / 改 context | 本章 mock 不切模型,s11 已演示过模型切换 |
shouldStopAfterTurn 钩子 | 每轮后调回调决定是否停 | 本章用 maxTurns 硬编码表达同一个思想 |
transformContext | 每轮前可改 messages(注入 system reminder 等) | 本章无动态改写需求 |
convertToLlm | agent 内部消息格式 vs LLM API 消息格式的转换层 | 本章 mock 直接用 ChatMessage,无格式差异 |
handleRunFailure 四事件模拟 | 错误时 emit message_start/message_end/turn_end/agent_end 四个事件 | 本章异常直接终止,listener 收不到 agent_end 也无所谓 |
| 真实 SSE 流式 | streamFn 逐 token 推送,event 有 message_update | 本章用 setTimeout 一次性返回,演示时序足够 |
MutableAgentState 拷贝语义 | tools/messages 的 getter/setter 都 .slice() | 本章 messages 直接暴露,调用方不外部改动即可 |
| 多事件类型 | turn_start/message_start/message_update/message_end/turn_end/agent_end/tool_execution_* | 本章只 emit 5 种事件,足以演示串行派发 |
| steering drain 位置 | pi 在循环开始前(agent-loop.ts 第 167 行)调一次 + 每轮末尾(第 253 行)调一次 | 教学版在内层 while 开头 + 末尾各 drain 一次;pi 的”启动前 drain”处理用户在 agent 启动到首次调 LLM 之间敲的字 |
| processEvents 状态机 | pi processEvents 有 switch 语句更新内部 state(streamingMessage / pendingToolCalls / errorMessage) | 教学版 processEvents 只派发事件,不更新 state;listener 拿不到 streamingMessage 等运行时状态 |
这张表往后会逐行划掉:s14 之后会补上 convertToLlm,s15 会补上真实 SSE 流式。
默写验收
合上 code.ts 和本 README,打开 practice.ts,凭记忆补全四个函数体:
PendingQueue.drain():两种 mode 的实现("all"批量 vs"one-at-a-time"逐条)runWithLifecycle():状态机三件套(AbortController + Promise + try/finally finishRun)loop():双层循环 + 两个 drain 位置(内层开头 drain steering、外层退出后 drain followUp)processEvents():串行 await listeners(for...of不是forEach)
通过标准:node practice.ts 跑起来,能正常对话;用 /steer xxx 注入能在下一轮 LLM 前生效;用 /follow xxx 注入能让 agent 续跑一轮;/abort 能中止当前 run。
写不出 loop() 说明两个 drain 位置没进脑子,回到「工作原理」第 3 块重读,重点记”内层开头 drain steering(循环内)→ 内层退出后 drain followUp(循环外)“这个时序差异。写不出 processEvents() 的串行 await——把 for...of 写成 forEach——回到第 4 块重读,这是”listener 的 promise 计入 run settlement”的关键。