Main Chain B · Step 07
Inbox:双队列的 session projection
packages/core/agent-loop/src/inbox.ts(244 行)。0.1.3-alpha.2 起 Inbox 从 core/agent 迁到 core/agent-loop,改名 ReactLoopInbox,并重构为一个 session projection:状态不再由构造器自己回放日志维护,而是注册成 sessionProjections 的一个投影单元(key 'inbox'),由投影注册表按事件增量折叠。所有输入方法仍汇到 mutate 一个函数。
示例本次示例:我们的消息在 inbox 里的三个 splice
# ① 用户点发送 → followup → send('next-turn', wakeup=true)(05 页 122 行)
# → inbox.splice('next-turn', Infinity, 0, [msg]) → mutate(本页 198 行)
# 归一化:start=Infinity → 钳到队尾 0;插入 1 条 → 落盘:
{ "type":"agent/inbox/spliced", "seq":N, "data":{
"target":"next-turn", "start":0, "inserted":[{"id":"m1",...}] } }
# (无 removedCount → 纯插入;无 outcome)
# ② turn 1 的 preStep claim(08 页 245 行)→ claim('next-turn', 1)(本页 109 行)
# 先清 next-step(空)→ 再取 next-turn 队首 1 条 → 两条纯删除 splice 落盘:
{ "data":{ "target":"next-step", "start":0 } } ← 空操作不落盘(220 行)
{ "data":{ "target":"next-turn", "start":0, "removedCount":1 } } ← 认领,无 outcome
# → live 通知 agent/inbox/claimed {message, turn:1}(112 行)
# ③ 用户取消(假设)→ cancel → inbox.clear()(95 行)
# 清 next-step 再清 next-turn → 这次 discardRemoved=true:
{ "data":{ "target":"next-turn", "start":0, "removedCount":1, "outcome":"canceled" } }
# outcome:'canceled' 是对账依据(§4 foldConsumedWork)
# 进程重启 → 新 ReactLoopInbox 构造(本页 75-81 行) # 构造器只注册投影定义(key 'inbox'),不再自己回放 # 投影注册表按日志里的 agent/inbox/spliced 事件增量折叠 # 命中上面 ② 的认领记录 → apply() 折叠 → next-turn 队列空 # 命中一个进程崩溃前没折叠的插入 → 消息回到队列 # 这就是「待办不丢」:真相在日志,投影只是缓存视图。
§1投影定义与构造器
27export const inboxProjectionDefinition = { // 投影定义:0.1.3-alpha.2 的核心
64 stateVersion: 1, // 投影状态自身的版本号
74export class ReactLoopInbox implements InboxContract {
75 constructor(
76 private readonly projections: SessionProjectionRegistry,
77 private readonly session: Session,
78 private readonly dispatch: AgentEventDispatch,
79 ) {} // 0.1.6-alpha.2:构造器现在完全是空的注册已移交给 AgentLoop
inboxProjectionSchema:投影状态的 wire 校验 schema——两个只读数组。它是 viewSchema 的来源(远程视图共享同一 schema)。
投影定义:key: 'inbox'(注册表的键)、stateSchema、init(空双队列)、apply(state, event)——只对 agent/inbox/spliced 事件折叠,其他事件原样返回 state。
校验搬进了投影的 apply:坐标合法性(37-41)、插入后跨队列 id 唯一(42-49)、失败带 session seq 的 cause 重抛(54)。0.1.2 时这些在 Inbox 的私有 validate 里、由构造器回放时调用;现在任何读这条日志的路径(恢复、重放、远程视图)都自动校验。
stateVersion: 1:投影状态自身的版本号——与 SESSION_FORMAT_VERSION(日志格式)是两个独立维度。
0.1.6-alpha.2:构造器现在是空的() {})。旧版这里有一句 this.projections.register(inboxProjectionDefinition);现在注册移交给了 AgentLoop——它在该文件 416-417 行一次性注册 turnBoundary 与 inbox 两个投影。ReactLoopInbox 仍需持有 projections 引用(读状态用,见 §3),但不再负责注册。
§2mutate:先校验,再落盘
198 private mutate(
199 target: InboxTarget,
200 start: number,
201 deleteCount: number,
202 inserted: UserMessage[],
203 discardRemoved: boolean,
204 ): UserMessage[] {
205 const state = this.current()
206 const inbox = state[target]
207 const truncatedStart = Math.trunc(start)
208 const offset = Number.isNaN(truncatedStart) ? 0 : truncatedStart
209 const actualStart = offset < 0
210 ? Math.max(inbox.length + offset, 0)
211 : Math.min(offset, inbox.length)
212 const truncatedDeleteCount = Math.trunc(deleteCount)
213 const actualDeleteCount = Math.min(
214 Math.max(Number.isNaN(truncatedDeleteCount) ? 0 : truncatedDeleteCount, 0),
215 inbox.length - actualStart,
216 )
217 if (actualDeleteCount === 0 && inserted.length === 0) return []
218 const candidate = inbox.toSpliced(actualStart, actualDeleteCount, ...inserted)先试算
219 const ids = new Set<string>()
220 for (const message of target === 'next-turn'
221 ? [...candidate, ...state['next-step']]
222 : [...state['next-turn'], ...candidate]) {
223 if (ids.has(message.id)) throw new Error(`message "${message.id}" is already pending`)跨队列唯一性
224 ids.add(message.id)
225 }
226 const outcome = discardRemoved && actualDeleteCount > 0 ? 'canceled' as const : undefined
227 const splice: SessionEventMap['agent/inbox/spliced'] = {
228 target,
229 start: actualStart,
230 ...(actualDeleteCount === 0 ? {} : { removedCount: actualDeleteCount }),
231 inserted,
232 ...(outcome === undefined ? {} : { outcome }),
233 }
234 const removed = inbox.slice(actualStart, actualStart + actualDeleteCount)
235 const event = this.session.append('agent/inbox/spliced', splice)
236-242 // 落盘后:按 event.data 发 discarded / inserted 通知
243 }
公开的 splice(166 行)与内部 mutate 的差别只在 discardRemoved:公开 splice 传 true(被移除的消息是「取消」),claim 内部传 false(被移除的消息是「认领」,不算取消)。
状态从投影注册表读:this.current() → projections.stateOf(session, 'inbox')。不再有 this.state 字段——单一真相源是日志,投影是它的缓存视图。
归一化坐标:截断小数、NaN 归零、负 offset 从尾部算、越界钳到界内。任何输入都归一成「合法 splice」。零操作直接返回(不产生无意义事件)。
0.1.3-alpha.2 的顺序变化:先本地试算并校验,再落盘。toSpliced 造出 candidate、跨队列 id 唯一性检查全部在 append 之前——坏 splice 根本进不了日志(此前是「先 append,由投影 apply 校验,失败则抛」)。
splice 记录是归一化后的坐标(target/start/removedCount/inserted/outcome),不是操作原始输入——回放只需机械重放,无需重算。outcome: 'canceled' 只标在「真的删了东西」的公开 splice 上,是 foldConsumedWork 对账的依据(§4)。
唯一写入:removed 从当前状态切片(供通知用),然后 append。事件进日志后,投影注册表按 apply 增量折叠出新状态——下一个 current() 就能读到。
两个通知方向:被移除的消息发 discarded(仅公开 splice),插入的消息发 inserted。都是从 event.data 取,不从内存状态取——通知的是「这次日志追加了什么」。
§3claim 与 current()
109 claim(target: InboxTarget, turn: number): UserMessage[] {
110 const claimed = this.mutate('next-step', 0, this.nextStep.length, [], false)
111 if (target === 'next-turn') claimed.push(...this.mutate('next-turn', 0, 1, [], false))
112 for (const message of claimed) this.dispatch.emit('agent/inbox/claimed', { message, turn })
113 return claimed
114 }
187 private current(): InboxState {
188 const state = this.projections.stateOf(this.session, 'inbox')
189 /* v8 ignore next -- the constructor registers this key before any read */
190 if (state === undefined) {
191 throw new Error(
192 `agent "${this.session.id}" cannot read inbox state: its projection registration is not active`,
193 )
194 }
195 return state
196 }
claim 的语义:总是先清空 next-step(全部),仅当 target 是 next-turn 时再取一条排队 turn。两次 mutate 都传 discardRemoved=false——认领的删除不是取消,不落 outcome:'canceled',不发 discarded 通知。
每条认领消息发 agent/inbox/claimed(带 owning turn)——被拒 step 的消息终点就是这里。
current() 是唯一的读状态入口:查投影注册表,未注册直接抛(fail-loud,而不是静默返回空队列——那会悄悄丢输入)。
§4对账:foldConsumedWork
packages/core/agent/src/consumed-work.ts(108 行)的纯函数 foldConsumedWork 是「从日志可重建」原则的又一实例:只靠 session 日志(turn/step 边界 + inbox splice)就能判定一个 agent 到底消费/丢弃了哪些工作——不需要任何内存状态。它的对账依据就是本页落下的两类记录:
- turn 的开合(turn/start、turn/end)圈定「哪些工作属于哪个 turn」;
- splice 的
removedCount+outcome:'canceled'区分「被认领消费」与「被取消丢弃」。
这解释了 226 行为什么 outcome 只在公开 splice 上标记:取消是对账需要的事实,认领不是。