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 一个函数。

05 agent.ts:113 send → :119 inbox.splice → agent-loop/src/inbox.ts:166 splice → :198 mutate → (下一步) 08 preStep:245 inbox.claim

示例本次示例:我们的消息在 inbox 里的三个 splice

示例轨迹 07-1 · 消息进入 → 认领 → 清空
# ① 用户点发送 → 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)
来源:inbox.ts 行级解读 + 05 页 send/claim 路径;splice 记录字段来自 inbox.ts:227-232
示例轨迹 07-2 · resume 时这页如何重放
# 进程重启 → 新 ReactLoopInbox 构造(本页 75-81 行)
# 构造器只注册投影定义(key 'inbox'),不再自己回放
# 投影注册表按日志里的 agent/inbox/spliced 事件增量折叠
# 命中上面 ② 的认领记录 → apply() 折叠 → next-turn 队列空
# 命中一个进程崩溃前没折叠的插入 → 消息回到队列
# 这就是「待办不丢」:真相在日志,投影只是缓存视图。
来源:inbox.ts:32-39 + 20 页持久化章

§1投影定义与构造器

packages/core/agent-loop/src/inbox.ts投影定义与构造器21-77
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
21-24

inboxProjectionSchema:投影状态的 wire 校验 schema——两个只读数组。它是 viewSchema 的来源(远程视图共享同一 schema)。

27-33

投影定义:key: 'inbox'(注册表的键)、stateSchema、init(空双队列)、apply(state, event)——只对 agent/inbox/spliced 事件折叠,其他事件原样返回 state。

34-55

校验搬进了投影的 apply:坐标合法性(37-41)、插入后跨队列 id 唯一(42-49)、失败带 session seq 的 cause 重抛(54)。0.1.2 时这些在 Inbox 的私有 validate 里、由构造器回放时调用;现在任何读这条日志的路径(恢复、重放、远程视图)都自动校验。

64

stateVersion: 1:投影状态自身的版本号——与 SESSION_FORMAT_VERSION(日志格式)是两个独立维度。

74-79

0.1.6-alpha.2:构造器现在是空的() {})。旧版这里有一句 this.projections.register(inboxProjectionDefinition);现在注册移交给了 AgentLoop——它在该文件 416-417 行一次性注册 turnBoundary 与 inbox 两个投影。ReactLoopInbox 仍需持有 projections 引用(读状态用,见 §3),但不再负责注册。

§2mutate:先校验,再落盘

packages/core/agent-loop/src/inbox.ts所有变更的唯一出口198-243
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  }
166-172

公开的 splice(166 行)与内部 mutate 的差别只在 discardRemoved:公开 splice 传 true(被移除的消息是「取消」),claim 内部传 false(被移除的消息是「认领」,不算取消)。

205

状态从投影注册表读:this.current() → projections.stateOf(session, 'inbox')。不再有 this.state 字段——单一真相源是日志,投影是它的缓存视图。

207-217

归一化坐标:截断小数、NaN 归零、负 offset 从尾部算、越界钳到界内。任何输入都归一成「合法 splice」。零操作直接返回(不产生无意义事件)。

218-225

0.1.3-alpha.2 的顺序变化:先本地试算并校验,再落盘。toSpliced 造出 candidate、跨队列 id 唯一性检查全部在 append 之前——坏 splice 根本进不了日志(此前是「先 append,由投影 apply 校验,失败则抛」)。

226-233

splice 记录是归一化后的坐标(target/start/removedCount/inserted/outcome),不是操作原始输入——回放只需机械重放,无需重算。outcome: 'canceled' 只标在「真的删了东西」的公开 splice 上,是 foldConsumedWork 对账的依据(§4)。

234-235

唯一写入:removed 从当前状态切片(供通知用),然后 append。事件进日志后,投影注册表按 apply 增量折叠出新状态——下一个 current() 就能读到。

236-242

两个通知方向:被移除的消息发 discarded(仅公开 splice),插入的消息发 inserted。都是从 event.data 取,不从内存状态取——通知的是「这次日志追加了什么」。

§3claim 与 current()

packages/core/agent-loop/src/inbox.tsclaim 与 current()111-116, 189-197
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  }
109-111

claim 的语义:总是先清空 next-step(全部),仅当 target 是 next-turn 时再取一条排队 turn。两次 mutate 都传 discardRemoved=false——认领的删除不是取消,不落 outcome:'canceled',不发 discarded 通知。

112

每条认领消息发 agent/inbox/claimed(带 owning turn)——被拒 step 的消息终点就是这里。

186-195

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 上标记:取消是对账需要的事实,认领不是。