DOCS v0.1.0 BUILD: STABLE GitHub ↗
章节目录
S13

Stateful Agent

steering/followUp 双队列 + idle/running 状态机 + 事件串行派发。

前置概念 双队列模式steering vs followUp状态机事件串行派发状态快照隔离

s13: 有状态 Agent — 双队列与状态机

把 agentLoop 包进一个类,让循环”有记忆”——能被叫停、能被改方向、还能在被叫停之后再被推一把。

...s12s13


问题

s01 的 agentLoop 是个无状态纯函数:每次调用都从头传 messages 数组,跑完就还给调用方。s04 给它加了 AbortSignalgetSteeringMessages 回调,让循环可以被外部打断、被中途改方向——但回调只是一根”线”,没有”对象”承载状态。结果调用方 main 里堆的胶水越来越多:history 数组、controller 句柄、steering 闭包……全靠人肉维护。

真实 pi 需要一个”可复用的有状态 Agent”:

  1. 维护 messages:自己持有对话历史,调用方不需要传。
  2. 两个注入入口steer() 在循环跑着的时候塞指令,followUp() 在 agent 本该停下来的时候塞指令——两种注入时机完全不同。
  3. 支持 abort / waitForIdle:外界能硬停当前 run,也能等它跑完。
  4. 串行派发事件:listener 按订阅顺序 await,listener 自己的 promise 也要计入 run 的 settlement。

两个注入入口的差异是本章的核心认知:steering 是”循环内注入”(当前轮工具照常跑完,下一轮 LLM 调用前生效);followUp 是”循环外注入”(内层循环已经退出、agent 本该停止时,被它唤起再跑一轮)。


解决方案

把 s01 的 agentLoop 包进一个 MiniAgent 类,维护四样东西:

  1. 状态字段_messages / _systemPrompt / _listeners 集合
  2. 双队列_steeringQueue_followUpQueue,都是 PendingQueue 实例(mode 可选 "all" 批量排干或 "one-at-a-time" 逐条排干)
  3. 状态机三件套_activeRun = { promise, resolve, abortController }——run 在则 isStreaming=true,run 走完 finishRun() 清场
  4. 双层循环:内层处理 tool-call + steering,外层处理 followUp
DUAL QUEUE · steering vs followUp
STEERING · 循环内注入
1
user 提问
2
assistant 流式输出
3
tool_call
4
tool_result
5
turn_end
6
★ drain steering
7
assistant 基于新指令输出
8
stop
steeringQueue:(空)
FOLLOWUP · 循环外注入
1
user 提问
2
assistant 流式输出
3
tool_call
4
tool_result
5
turn_end
6
agent 即将停止
7
★ drain followUp
8
续跑一轮
9
stop
followUpQueue:(空)
状态机
idle
prompt()
running
finishRun()
aborted
activeRun 三件套
promisesettled
resolve已调用
abortControlleraborted=false

关键认知: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 调用。试这三个场景,观察两种注入点的时序差异:

  1. steering(循环内注入):输入”用工具列文件”,agent 开始流式输出 + 工具调用。运行中再输入 /steer 别动 test 目录,应看到紫色 [steering]turn_end 之后、下一轮 LLM 之前被注入,当前轮工具调用照常执行完毕。
  2. followUp(循环外注入):输入一个简单问题(不含”工具”),agent 一轮就停。在它停止前输入 /follow 顺便补一个 README,应看到紫色 [followUp] 在内层循环退出后、agent_end 之前被注入,agent 续跑一轮。
  3. abort:输入一个慢任务,运行中输入 /abort(或按 Ctrl+C),应看到 run 立刻终止,回到提示符。

前置概念清单

本章引入五个新概念:

  1. 双队列模式:steering 与 followUp 共享同一个 PendingQueue 容器,只是 drain 的位置不同。两个队列、两种语义、一套数据结构。
  2. steering vs followUp:steering 是循环内注入(下一轮 LLM 调用前生效,当前轮工具照常跑完);followUp 是循环外注入(内层循环退出、agent 本该停止时生效,让循环续跑一轮)。差异在 drain 位置,不在 API 形状。
  3. 状态机idle ⇄ running,由 _activeRun 是否存在决定。run 在则 isStreaming=true,run 走完 finishRun() 清场回到 idle。abort 走 abortController.abort(),是 running 的一种特殊退出路径。
  4. 事件串行派发:listener 按订阅顺序 for...of + await,不是 forEach 并发。listener 的 promise 计入 run settlement——waitForIdle() 等到最后一个 listener settle。
  5. 状态快照隔离:每次 prompt() 调用前 this._messages.slice() 传给 loop——loop 拿到的是快照,即使 run 中途状态被外部改动,loop 内部用的是它启动时的那份。

源码锚点

mini-pi 的 MiniAgent,在pi 里是 packages/agent/src/agent.tsAgent 类。配套双层循环在 packages/agent/src/agent-loop.ts。读完回答三个问题:

锚点行号看什么
Agent166-557整个有状态 Agent 的骨架,对照 MiniAgent
MutableAgentState59-93拷贝语义:tools/messages 的 getter/setter 都 .slice()
PendingMessageQueue118-152双队列共享容器,drain 两种 mode 实现
steer / followUp263-271两个入口只是 enqueue,差异在 drain 位置
runWithLifecycle451-474状态机三件套,try/finally finishRun
finishRun494-500清场:resolve promise + 清空 activeRun
processEvents509-556switch 更新 state + for...of await 串行派发
drain 位置(agent-loop.ts)167 / 253 / 257steering 在循环开始前 + 内层每轮末尾;followUp 在内层退出后
  1. agent.ts 第 59-93 行 MutableAgentState 的 tools 和 messages 都用 getter/setter 包了 .slice()为什么赋值时也要 copy?(提示:外部数组随时可能被原地改,agent 内部状态必须独立)
  2. agent.ts 第 451-474 行 runWithLifecycletry/catch/finally——为什么 catch 里调 handleRunFailure 模拟一个 failure message,而 mini-pi 直接让异常外抛?(提示:pi 要保证 listener 一定能收到 agent_end 事件,mini-pi 用 finally 清场就够了)
  3. 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 等)本章无动态改写需求
convertToLlmagent 内部消息格式 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,凭记忆补全四个函数体:

  1. PendingQueue.drain():两种 mode 的实现("all" 批量 vs "one-at-a-time" 逐条)
  2. runWithLifecycle():状态机三件套(AbortController + Promise + try/finally finishRun)
  3. loop():双层循环 + 两个 drain 位置(内层开头 drain steering、外层退出后 drain followUp)
  4. 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”的关键。


下一章:s14 会话分支树 — 一个 .jsonl 如何承载对话分叉