Sfoglia il codice sorgente

fix(client): pin inbox claim semantics

imccyu 1 mese fa
parent
commit
b4527bedc7

+ 2 - 2
.agents/notes/implemented/architecture/2026-08-09-client-conversation-node-assembly.i18n.yaml

@@ -2,5 +2,5 @@
 # side as of the last confirmed-consistent state. Both languages carry equal authority;
 # after editing either side, bring the other along and re-record with:
 #   pnpm run verify-translation-pairing --write .agents/notes/implemented/architecture/2026-08-09-client-conversation-node-assembly.md
-2026-08-09-client-conversation-node-assembly.md: 29ee0821638e6a594e53ce1410426c34f099b93d
-2026-08-09-client-conversation-node-assembly.zh.md: 18c3df86c970cdf1d5d51e131956271ae049b628
+2026-08-09-client-conversation-node-assembly.md: 4866f228052633af32a882dd4c96dd28e06851b5
+2026-08-09-client-conversation-node-assembly.zh.md: b50762b70d14c764f8ccf5d803993f5801ae755f

+ 1 - 1
.agents/notes/implemented/architecture/2026-08-09-client-conversation-node-assembly.md

@@ -18,7 +18,7 @@ Client Runtime provides a target-neutral Conversation Node assembly engine. Busi
 
 This Note retains the derivation, business-by-business validation, responsibilities, algorithms, and trade-offs that remain relevant after implementation.
 
-Chat registers an Inbox Definition only for `next-step`, because message classification is its sole consumer; `next-turn` splices remain durable Session inputs but create no Chat Context. Chat and Trajectory each keep target-owned next-step state. Every insertion stores only message IDs in an immutable splice node. A successful claim materializes the pending chain once, replaces the previous claimed set with that batch, and lets later Contexts share the set until another claim. Historical Contexts therefore retain linear ID state instead of cumulative array and Set snapshots.
+Chat registers an Inbox Definition only for `next-step`, because message classification is its sole consumer; `next-turn` splices remain durable Session inputs but create no Chat Context. Chat and Trajectory each keep target-owned next-step state. Every insertion stores only message IDs in an immutable splice node. A successful claim materializes the pending chain once, replaces the previous claimed set with that batch, and lets later Contexts share the set until another claim. The AgentLoop appends every message admitted from that claim before it can claim another batch; a rejected claim appends no `user/message`, so later classification needs only the current batch. Historical Contexts therefore retain linear ID state instead of cumulative array and Set snapshots.
 
 ### Responsibility layers
 

+ 1 - 1
.agents/notes/implemented/architecture/2026-08-09-client-conversation-node-assembly.zh.md

@@ -18,7 +18,7 @@ Client Runtime 提供 target-neutral 的 Conversation Node 组装引擎,业务
 
 本 Note 保留实现后仍有价值的方案推导、逐业务适配、职责、算法和取舍。
 
-Chat 只注册 `next-step` Inbox Definition,因为消息分类是其唯一消费方;`next-turn` splice 仍是持久 Session input,但不会创建 Chat Context。Chat 与 Trajectory 各自维护 target 专属 next-step state。每次插入只把消息 ID 写入不可变 splice 节点。成功 claim 时只 materialize 一次 pending 链,以当前批次替换上一个 claimed Set,并让后续 Context 共享该 Set,直到下一次 claim。历史 Context 因而只保留线性 ID state,不再保留累计数组和 Set 快照。
+Chat 只注册 `next-step` Inbox Definition,因为消息分类是其唯一消费方;`next-turn` splice 仍是持久 Session input,但不会创建 Chat Context。Chat 与 Trajectory 各自维护 target 专属 next-step state。每次插入只把消息 ID 写入不可变 splice 节点。成功 claim 时只 materialize 一次 pending 链,以当前批次替换上一个 claimed Set,并让后续 Context 共享该 Set,直到下一次 claim。AgentLoop 会在领取下一批消息之前追加当前 claim 接纳的全部消息;被拒绝的 claim 不追加 `user/message`,因此后续分类只需当前批次。历史 Context 因而只保留线性 ID state,不再保留累计数组和 Set 快照。
 
 ### 责任分层
 

+ 2 - 2
.agents/notes/implemented/architecture/2026-08-11-trajectory-conversation-context-assembly.i18n.yaml

@@ -2,5 +2,5 @@
 # side as of the last confirmed-consistent state. Both languages carry equal authority;
 # after editing either side, bring the other along and re-record with:
 #   pnpm run verify-translation-pairing --write .agents/notes/implemented/architecture/2026-08-11-trajectory-conversation-context-assembly.md
-2026-08-11-trajectory-conversation-context-assembly.md: 7d0aea2fc09f0f04bd5de923bee15c42a773489a
-2026-08-11-trajectory-conversation-context-assembly.zh.md: 1909243da02093de1a5f1d0857a11d3742258492
+2026-08-11-trajectory-conversation-context-assembly.md: ad4bff77b87649a6e8b208aa6a55575c1a4e23f0
+2026-08-11-trajectory-conversation-context-assembly.zh.md: e7bd90dce00a8bd8a23eed8d076c5350df5d7c15

+ 1 - 1
.agents/notes/implemented/architecture/2026-08-11-trajectory-conversation-context-assembly.md

@@ -40,7 +40,7 @@ Assistant chunks update only their `turn:step` Context. Content-bearing chunks r
 
 Trajectory reconstructs steering from durable inbox history, using the same identity rule as the [Chat steering decision](../feature/2026-08-04-web-context-source-and-steer-marks.md) without sharing Chat's final Node.
 
-Each `agent/inbox/spliced` Event targeting `next-step` starts an invisible Context identified by its Event seq. Its `start()` reads the nearest earlier inbox Context, applies the splice, and stores the pending identities plus the cumulative set of claimed message IDs. A later user-origin `user/message` reads the nearest earlier inbox Context: a claimed ID produces a Steering Node, while every other user-origin message produces an ordinary User Node.
+Each `agent/inbox/spliced` Event targeting `next-step` starts an invisible Context identified by its Event seq. Its `start()` reads the nearest earlier inbox Context, appends the splice to persistent pending-ID state, and materializes that state only when a claim replaces the current claimed batch. The AgentLoop appends every admitted message from one claim before it can claim another batch; a rejected claim appends no `user/message`. A later user-origin `user/message` reads the nearest earlier inbox Context: an ID in the current claim produces a Steering Node, while every other user-origin message produces an ordinary User Node.
 
 A Reader miss while older history remains records a window-gap dependency. When prepend supplies the missing predecessor, the Assembler replays the affected inbox chain and message Contexts in forward Event order. Historical page direction therefore cannot permanently misclassify a message.
 

+ 1 - 1
.agents/notes/implemented/architecture/2026-08-11-trajectory-conversation-context-assembly.zh.md

@@ -40,7 +40,7 @@ Assistant chunk 只更新对应的 `turn:step` Context。带内容的 chunk 请
 
 Trajectory 从持久 inbox 历史恢复 steering,使用与 [Chat steering 决策](../feature/2026-08-04-web-context-source-and-steer-marks.zh.md)相同的标识规则,但不共享 Chat 的最终 Node。
 
-每条目标为 `next-step` 的 `agent/inbox/spliced` Event 都会启动一个以 Event seq 标识的不可见 Context。它的 `start()` 读取最近的前序 inbox Context,应用 splice,并存储待处理标识以及累计的已领取 message ID 集合。后续用户来源的 `user/message` 读取最近的前序 inbox Context:已领取的 ID 生成 Steering Node,其余用户来源消息生成普通 User Node。
+每条目标为 `next-step` 的 `agent/inbox/spliced` Event 都会启动一个以 Event seq 标识的不可见 Context。它的 `start()` 读取最近的前序 inbox Context,把 splice 追加到持久的 pending ID state,并只在 claim 时 materialize 该 state、替换当前 claimed batch。AgentLoop 会在领取下一批消息之前追加当前 claim 接纳的全部消息;被拒绝的 claim 不追加 `user/message`。后续用户来源的 `user/message` 读取最近的前序 inbox Context:ID 属于当前 claim 时生成 Steering Node,其余用户来源消息生成普通 User Node。
 
 仍有更早历史时,Reader miss 会记录 window-gap 依赖。prepend 补齐缺失的前驱后,Assembler 按 Event 正序重放受影响的 inbox chain 与 message Context。因此,历史分页方向不会永久错误分类消息。
 

+ 17 - 9
packages/client/ui-chat/src/client/conversation-nodes/inbox.ts

@@ -29,16 +29,16 @@ interface PendingSplice {
 
 type PendingState = PendingSnapshot | PendingSplice
 
-/** Cumulative state after one durable inbox splice. */
+/** Persistent next-step state after one durable Inbox splice. */
 export interface InboxState {
   /** Persistent splice chain materialized only when a next-step batch is claimed. */
   readonly pending: PendingState
-  /** Message ids in the current claimed batch, shared until the next claim. */
-  readonly claimed: ReadonlySet<string>
+  /** Message ids in the current claim, shared until the next claim. */
+  readonly currentClaimed: ReadonlySet<string>
 }
 
 const EMPTY_PENDING: PendingState = { kind: 'snapshot', ids: [] }
-const EMPTY_CLAIMED: ReadonlySet<string> = new Set()
+const EMPTY_CURRENT_CLAIMED: ReadonlySet<string> = new Set()
 
 function materializePending(state: PendingState): string[] {
   const splices: PendingSplice[] = []
@@ -48,8 +48,7 @@ function materializePending(state: PendingState): string[] {
     current = current.previous
   }
   const pending = [...current.ids]
-  for (let index = splices.length - 1; index >= 0; index--) {
-    const splice = splices[index] as PendingSplice
+  for (const splice of splices.reverse()) {
     pending.splice(splice.start, splice.removedCount, ...splice.inserted)
   }
   return pending
@@ -68,6 +67,12 @@ function withoutInserted(
   return next ?? claimed
 }
 
+/**
+ * Apply one next-step splice under the AgentLoop's durable event ordering.
+ * An entered claim logs its complete message batch before another claim; a
+ * rejected claim logs no messages, so only the current claim can classify a
+ * later `user/message`.
+ */
 function applySplice(
   previous: ConversationPreviousContext<InboxState> | undefined,
   splice: InboxSplice,
@@ -75,15 +80,18 @@ function applySplice(
   const priorPending = previous?.state.pending ?? EMPTY_PENDING
   const inserted = splice.inserted.map(identity => identity.id)
   const removedCount = splice.removedCount ?? 0
-  const priorClaimed = withoutInserted(previous?.state.claimed ?? EMPTY_CLAIMED, inserted)
   if (removedCount > 0 && splice.outcome !== 'canceled') {
     const pending = materializePending(priorPending)
     const removed = pending.splice(splice.start, removedCount, ...inserted)
     return {
       pending: { kind: 'snapshot', ids: pending },
-      claimed: new Set(removed),
+      currentClaimed: new Set(removed),
     }
   }
+  const currentClaimed = withoutInserted(
+    previous?.state.currentClaimed ?? EMPTY_CURRENT_CLAIMED,
+    inserted,
+  )
   return {
     pending: {
       kind: 'splice',
@@ -92,7 +100,7 @@ function applySplice(
       removedCount,
       inserted,
     },
-    claimed: priorClaimed,
+    currentClaimed,
   }
 }
 

+ 2 - 1
packages/client/ui-chat/src/client/conversation-nodes/message.ts

@@ -58,7 +58,8 @@ export const messageDefinition: ConversationNodeDefinition<MessageNode> = {
         form: contextForm(event.data.source),
       }
     }
-    const claimed = reader.previous<InboxState>('inbox-next-step')?.state.claimed.has(String(event.data.id)) === true
+    const claimed = reader.previous<InboxState>('inbox-next-step')
+      ?.state.currentClaimed.has(String(event.data.id)) === true
     return claimed
       ? {
         kind: 'steering',

+ 57 - 0
packages/client/ui-chat/tests/conversation-node-definitions.client.spec.ts

@@ -418,6 +418,63 @@ describe('built-in conversation node Definitions', () => {
     ])
   })
 
+  it('replays pending splice chains and scopes steering to the current claim', () => {
+    const first = textMessage('claim-first', 'first')
+    const second = textMessage('claim-second', 'second')
+    const canceled = textMessage('claim-canceled', 'canceled')
+    const requeued = textMessage('claim-requeued', 'requeued')
+    const later = textMessage('claim-later', 'later')
+    const current = snapshot(assembler([
+      at(1, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, inserted: [first],
+      }),
+      at(2, 'agent/inbox/spliced', {
+        target: 'next-step', start: 1, inserted: [canceled],
+      }),
+      at(3, 'agent/inbox/spliced', {
+        target: 'next-step', start: 1, inserted: [second],
+      }),
+      at(4, 'agent/inbox/spliced', {
+        target: 'next-step', start: 2, removedCount: 1, inserted: [], outcome: 'canceled',
+      }),
+      at(5, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, removedCount: 2, inserted: [],
+      }),
+      at(6, 'user/message', first, { surfaceOp: 'append' }),
+      at(7, 'user/message', second, { surfaceOp: 'append' }),
+      at(8, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, inserted: [requeued],
+      }),
+      at(9, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, removedCount: 1, inserted: [],
+      }),
+      at(10, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, inserted: [requeued],
+      }),
+      at(11, 'user/message', requeued, { surfaceOp: 'append' }),
+      at(12, 'user/message', canceled, { surfaceOp: 'append' }),
+      at(13, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, removedCount: 1, inserted: [], outcome: 'canceled',
+      }),
+      at(14, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, inserted: [later],
+      }),
+      at(15, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, removedCount: 1, inserted: [],
+      }),
+      at(16, 'user/message', later, { surfaceOp: 'append' }),
+    ]))
+
+    expect(current.order.map(key => current.nodes.get(key)).filter(node =>
+      node?.kind === 'user' || node?.kind === 'steering')).toMatchObject([
+      { kind: 'steering', data: { seq: 6 } },
+      { kind: 'steering', data: { seq: 7 } },
+      { kind: 'user', data: { seq: 11 } },
+      { kind: 'user', data: { seq: 12 } },
+      { kind: 'steering', data: { seq: 16 } },
+    ])
+  })
+
   it('orders a command-started Turn first steering before its process control', () => {
     const steering = textMessage('command-task', 'plan this change')
     const value = assembler([

+ 13 - 3
packages/client/ui-tool/tests/tool-row.client.spec.tsx

@@ -260,14 +260,24 @@ describe('ToolRow', () => {
 
   it('formats the argument body only while expanding it', () => {
     const stringify = vi.spyOn(JSON, 'stringify')
+    const bodyFormatCalls = () => stringify.mock.calls.filter(
+      ([value, replacer, space]) => typeof value === 'object'
+        && value !== null
+        && 'a' in value
+        && value.a === 1
+        && replacer === null
+        && space === 2,
+    ).length
     const view = render(<ToolRow {...rowProps} />)
-    expect(stringify).not.toHaveBeenCalled()
+    expect(bodyFormatCalls()).toBe(0)
 
     fireEvent.click(view.getByRole('button'))
-    expect(stringify).toHaveBeenCalledTimes(1)
+    expect(bodyFormatCalls()).toBe(1)
+    expect(view.getByText(/"a": 1/)).toBeTruthy()
 
     fireEvent.click(view.getByRole('button'))
-    expect(stringify).toHaveBeenCalledTimes(1)
+    expect(bodyFormatCalls()).toBe(1)
+    expect(view.queryByText(/"a": 1/)).toBeNull()
   })
 
   it('running keeps the icon (row sweep carries the signal); error swaps in a StateDot', () => {

+ 17 - 9
packages/client/ui-trajectory/src/client/trajectory-message-definitions.ts

@@ -39,14 +39,14 @@ type PendingState = PendingSnapshot | PendingSplice
 interface InboxState {
   /** Persistent splice chain materialized only when a next-step batch is claimed. */
   readonly pending: PendingState
-  /** Message ids in the current claimed batch, shared until the next claim. */
-  readonly claimed: ReadonlySet<string>
+  /** Message ids in the current claim, shared until the next claim. */
+  readonly currentClaimed: ReadonlySet<string>
 }
 
 type MessageNode = UserMessageNode | SteeringMessageNode | ContextMessageNode
 
 const EMPTY_PENDING: PendingState = { kind: 'snapshot', ids: [] }
-const EMPTY_CLAIMED: ReadonlySet<string> = new Set()
+const EMPTY_CURRENT_CLAIMED: ReadonlySet<string> = new Set()
 
 function materializePending(state: PendingState): string[] {
   const splices: PendingSplice[] = []
@@ -56,8 +56,7 @@ function materializePending(state: PendingState): string[] {
     current = current.previous
   }
   const pending = [...current.ids]
-  for (let index = splices.length - 1; index >= 0; index--) {
-    const splice = splices[index] as PendingSplice
+  for (const splice of splices.reverse()) {
     pending.splice(splice.start, splice.removedCount, ...splice.inserted)
   }
   return pending
@@ -76,6 +75,12 @@ function withoutInserted(
   return next ?? claimed
 }
 
+/**
+ * Apply one next-step splice under the AgentLoop's durable event ordering.
+ * An entered claim logs its complete message batch before another claim; a
+ * rejected claim logs no messages, so only the current claim can classify a
+ * later `user/message`.
+ */
 function applySplice(
   previous: ConversationPreviousContext<InboxState> | undefined,
   splice: InboxSplice,
@@ -83,15 +88,18 @@ function applySplice(
   const priorPending = previous?.state.pending ?? EMPTY_PENDING
   const inserted = splice.inserted.map(identity => identity.id)
   const removedCount = splice.removedCount ?? 0
-  const priorClaimed = withoutInserted(previous?.state.claimed ?? EMPTY_CLAIMED, inserted)
   if (removedCount > 0 && splice.outcome !== 'canceled') {
     const pending = materializePending(priorPending)
     const removed = pending.splice(splice.start, removedCount, ...inserted)
     return {
       pending: { kind: 'snapshot', ids: pending },
-      claimed: new Set(removed),
+      currentClaimed: new Set(removed),
     }
   }
+  const currentClaimed = withoutInserted(
+    previous?.state.currentClaimed ?? EMPTY_CURRENT_CLAIMED,
+    inserted,
+  )
   return {
     pending: {
       kind: 'splice',
@@ -100,7 +108,7 @@ function applySplice(
       removedCount,
       inserted,
     },
-    claimed: priorClaimed,
+    currentClaimed,
   }
 }
 
@@ -148,7 +156,7 @@ const trajectoryMessageDefinition: ConversationNodeDefinition<MessageNode> = {
       }
     }
     const claimed = reader.previous<InboxState>('trajectory-inbox-next-step')
-      ?.state.claimed.has(String(event.data.id)) === true
+      ?.state.currentClaimed.has(String(event.data.id)) === true
     return claimed
       ? {
         kind: 'steering',

+ 66 - 0
packages/client/ui-trajectory/tests/conversation-definitions.client.spec.ts

@@ -508,4 +508,70 @@ describe('Trajectory conversation Definitions', () => {
       ? request.promptChange?.kind
       : undefined)).toEqual(['initial', undefined])
   })
+
+  it('replays pending splice chains and scopes steering to the current claim', () => {
+    const message = (id: string, text: string) => ({
+      id,
+      role: 'user',
+      content: [{ type: 'text', text }],
+      source: { kind: 'user' },
+    })
+    const first = message('claim-first', 'first')
+    const second = message('claim-second', 'second')
+    const canceled = message('claim-canceled', 'canceled')
+    const requeued = message('claim-requeued', 'requeued')
+    const later = message('claim-later', 'later')
+    const current = snapshot(assembler([
+      at(1, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, inserted: [first],
+      }),
+      at(2, 'agent/inbox/spliced', {
+        target: 'next-step', start: 1, inserted: [canceled],
+      }),
+      at(3, 'agent/inbox/spliced', {
+        target: 'next-step', start: 1, inserted: [second],
+      }),
+      at(4, 'agent/inbox/spliced', {
+        target: 'next-step', start: 2, removedCount: 1, inserted: [], outcome: 'canceled',
+      }),
+      at(5, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, removedCount: 2, inserted: [],
+      }),
+      at(6, 'user/message', first),
+      at(7, 'user/message', second),
+      at(8, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, inserted: [requeued],
+      }),
+      at(9, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, removedCount: 1, inserted: [],
+      }),
+      at(10, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, inserted: [requeued],
+      }),
+      at(11, 'user/message', requeued),
+      at(12, 'user/message', canceled),
+      at(13, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, removedCount: 1, inserted: [], outcome: 'canceled',
+      }),
+      at(14, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, inserted: [later],
+      }),
+      at(15, 'agent/inbox/spliced', {
+        target: 'next-step', start: 0, removedCount: 1, inserted: [],
+      }),
+      at(16, 'user/message', later),
+    ]))
+
+    expect(current.eventNodes.filter(node =>
+      node.kind === 'user' || node.kind === 'steering').map(node => ({
+      kind: node.kind,
+      seq: node.seq,
+    }))).toEqual([
+      { kind: 'steering', seq: 6 },
+      { kind: 'steering', seq: 7 },
+      { kind: 'user', seq: 11 },
+      { kind: 'user', seq: 12 },
+      { kind: 'steering', seq: 16 },
+    ])
+  })
 })

+ 2 - 2
packages/core/agent-loop/README.i18n.yaml

@@ -2,5 +2,5 @@
 # side as of the last confirmed-consistent state. Both languages carry equal authority;
 # after editing either side, bring the other along and re-record with:
 #   pnpm run verify-translation-pairing --write packages/core/agent-loop/README.md
-README.md: a37d75545b4c1eccb4363c7ff2f09ae102f480db
-README.zh.md: a209e6a99d96cb8790a19b0ebb043332d619026f
+README.md: 60d1e75a5038500a42244a45fb077e1eb71da092
+README.zh.md: 166ba6ff366f9a8c8c687b4521bb04b591814851

+ 1 - 1
packages/core/agent-loop/README.md

@@ -107,7 +107,7 @@ Creation is one rollback-covered transaction: construct a private session, concr
 
 ### Turn and step flow
 
-The driver owns one agent for its lifetime and runs inside `ctx.agents.withInitiator(agent, ...)`. At a turn boundary it opens the durable turn, then atomically claims pending next-step input plus one queued prompt; between steps it claims only next-step input. `agent/pre-step` decides what enters the step; each successful model call appends one `assistant/message` anchor citing its chunk seqs, and a cancelled stream appends an `interrupted: true` anchor with the delivered prefix so the next request contains what the user saw. Within a step, exclusive calls form barriers and parallel-safe calls use the bounded rolling pool; policy, durable results, and result context remain model-ordered.
+The driver owns one agent for its lifetime and runs inside `ctx.agents.withInitiator(agent, ...)`. At a turn boundary it opens the durable turn, then atomically claims pending next-step input plus one queued prompt; between steps it claims only next-step input. `agent/pre-step` decides what enters the step. An entered decision appends its complete `user/message` batch before the driver can claim again, while a rejected decision appends none. Each successful model call appends one `assistant/message` anchor citing its chunk seqs, and a cancelled stream appends an `interrupted: true` anchor with the delivered prefix so the next request contains what the user saw. Within a step, exclusive calls form barriers and parallel-safe calls use the bounded rolling pool; policy, durable results, and result context remain model-ordered.
 
 ### Failure and cancellation
 

+ 1 - 1
packages/core/agent-loop/README.zh.md

@@ -107,7 +107,7 @@ const handle = await ctx.agents.create({
 
 ### 轮次与步骤流程
 
-驱动器在其整个生命周期内拥有一个 agent,并在 `ctx.agents.withInitiator(agent, ...)` 内运行。在轮次边界,它先打开持久轮次,再原子领取待处理的 next-step 输入与一条排队提示词;在步骤之间则只领取 next-step 输入。`agent/pre-step` 决定什么进入该步骤;每次成功的模型调用都恰好追加一个引用其分片 seq 的 `assistant/message` 锚点,被取消的流则追加带 `interrupted: true` 的锚点并携带已交付前缀,使下一次请求包含用户看到的内容。在步骤内,独占调用形成屏障,并行安全调用使用有界滚动池;策略、持久结果与结果上下文保持模型顺序。
+驱动器在其整个生命周期内拥有一个 agent,并在 `ctx.agents.withInitiator(agent, ...)` 内运行。在轮次边界,它先打开持久轮次,再原子领取待处理的 next-step 输入与一条排队提示词;在步骤之间则只领取 next-step 输入。`agent/pre-step` 决定什么进入该步骤。进入步骤的决定会在驱动器再次领取消息前追加完整的 `user/message` 批次,被拒绝的决定则不追加任何消息。每次成功的模型调用都恰好追加一个引用其分片 seq 的 `assistant/message` 锚点,被取消的流则追加带 `interrupted: true` 的锚点并携带已交付前缀,使下一次请求包含用户看到的内容。在步骤内,独占调用形成屏障,并行安全调用使用有界滚动池;策略、持久结果与结果上下文保持模型顺序。
 
 ### 失败与取消
 

+ 49 - 0
packages/core/agent-loop/tests/contract-regressions.spec.ts

@@ -512,6 +512,55 @@ describe('adapter registration, routing, and accepted-input ownership', () => {
     expect(steeringSources).toEqual([{ kind: 'plugin', plugin: 'goal' }])
   })
 
+  it('records each admitted next-step batch before the following claim', async () => {
+    const adapter = new MockAdapter([
+      toolCallResponse('c1', 'steer_next', {}),
+      toolCallResponse('c2', 'steer_next', {}),
+      textResponse('done'),
+    ])
+    const ctx = await harness(adapter)
+    const steering = [
+      createUserMessage({ content: [{ type: 'text', text: 'first steer' }], source: { kind: 'user' } }),
+      createUserMessage({ content: [{ type: 'text', text: 'second steer' }], source: { kind: 'user' } }),
+    ]
+    const agent = ctx.agentLoop.create(SessionId('claim-order'), { provider: 'mock', model: 'mock' })
+    let execution = 0
+    ctx.tools.register(defineContentToolFixture({
+      name: 'steer_next',
+      description: '',
+      parameters: {},
+      async execute() {
+        const message = steering[execution]
+        execution += 1
+        if (message !== undefined) agent.steer(message)
+        return []
+      },
+    }))
+
+    send(agent, 'go')
+    await waitForIdle(ctx, agent)
+
+    const events = agent.session.events
+    const claims = events.flatMap(event => event.type === 'agent/inbox/spliced'
+      && event.data.target === 'next-step'
+      && event.data.outcome !== 'canceled'
+      && (event.data.removedCount ?? 0) > 0
+      ? [event]
+      : [])
+    expect(claims).toHaveLength(2)
+    for (const [index, message] of steering.entries()) {
+      const claim = claims[index]
+      const admitted = events.find(event =>
+        event.type === 'user/message' && event.data.id === message.id)
+      expect(claim).toBeDefined()
+      expect(admitted).toBeDefined()
+      if (claim === undefined || admitted === undefined) continue
+      expect(admitted.seq).toBeGreaterThan(claim.seq)
+      const nextClaim = claims[index + 1]
+      if (nextClaim !== undefined) expect(admitted.seq).toBeLessThan(nextClaim.seq)
+    }
+  })
+
 })
 
 describe('turn numbering continues across seeded sessions', () => {