Просмотр исходного кода

fix(agent-team): preserve mailbox order on cold resume

Dudu-0223 2 недель назад
Родитель
Сommit
eeddd457cd

+ 2 - 2
docs/persistence-catalog.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 docs/persistence-catalog.md
-persistence-catalog.md: 5b8ae5748c997b829ce715534bee4390c0aac1f5
-persistence-catalog.zh.md: 9a5d1b40201a26d20839902f109d8c1ab8eb5301
+persistence-catalog.md: 1c0c6919987c691b82c4639aff0f779c95dca83c
+persistence-catalog.zh.md: 4cc8ba5b7fc76708a80285013ebcbb03fe3e8e4a

+ 4 - 4
docs/persistence-catalog.md

@@ -776,7 +776,7 @@ Source: [`packages/subagent/tool-subagent/src/model-selection-state.ts:17`](../p
 
 ```ts persistence-catalog
 /** Whole teammate lifecycle value, stored only in the Team Lead Session. */
-'team/member': { version: 1; teamId: TeamId; member: TeamMemberSnapshot }
+'team/member': { version: 2; teamId: TeamId; member: TeamMemberSnapshot }
 ```
 
 Types: [TeamId](subsystems/agent-team.md) · [TeamMemberSnapshot](subsystems/agent-team.md)
@@ -790,7 +790,7 @@ Source: [`packages/experimental/agent-team/src/types.ts:221`](../packages/experi
 ```ts persistence-catalog
 /** Durable acknowledgement that the target Session recorded the message. */
 'team/message/delivered': {
-  version: 1
+  version: 2
   teamId: TeamId
   messageId: TeamMessageId
   targetId: SessionId
@@ -807,7 +807,7 @@ Source: [`packages/experimental/agent-team/src/types.ts:227`](../packages/experi
 
 ```ts persistence-catalog
 /** Durable mailbox enqueue, stored before delivery is attempted. */
-'team/message/queued': { version: 1; teamId: TeamId; message: TeamMessageSnapshot }
+'team/message/queued': { version: 2; teamId: TeamId; message: TeamMessageSnapshot }
 ```
 
 Types: [TeamId](subsystems/agent-team.md) · [TeamMessageSnapshot](subsystems/agent-team.md)
@@ -820,7 +820,7 @@ Source: [`packages/experimental/agent-team/src/types.ts:225`](../packages/experi
 
 ```ts persistence-catalog
 /** Whole shared-task value, stored only in the Team Lead Session. */
-'team/task': { version: 1; teamId: TeamId; task: TeamTaskSnapshot }
+'team/task': { version: 2; teamId: TeamId; task: TeamTaskSnapshot }
 ```
 
 Types: [TeamId](subsystems/agent-team.md) · [TeamTaskSnapshot](subsystems/agent-team.md)

+ 4 - 4
docs/persistence-catalog.zh.md

@@ -778,7 +778,7 @@ export type SessionEvent<T extends SessionEventType = SessionEventType> = {
 
 ```ts persistence-catalog
 /** Whole teammate lifecycle value, stored only in the Team Lead Session. */
-'team/member': { version: 1; teamId: TeamId; member: TeamMemberSnapshot }
+'team/member': { version: 2; teamId: TeamId; member: TeamMemberSnapshot }
 ```
 
 类型:[TeamId](subsystems/agent-team.zh.md) · [TeamMemberSnapshot](subsystems/agent-team.zh.md)
@@ -792,7 +792,7 @@ export type SessionEvent<T extends SessionEventType = SessionEventType> = {
 ```ts persistence-catalog
 /** Durable acknowledgement that the target Session recorded the message. */
 'team/message/delivered': {
-  version: 1
+  version: 2
   teamId: TeamId
   messageId: TeamMessageId
   targetId: SessionId
@@ -809,7 +809,7 @@ export type SessionEvent<T extends SessionEventType = SessionEventType> = {
 
 ```ts persistence-catalog
 /** Durable mailbox enqueue, stored before delivery is attempted. */
-'team/message/queued': { version: 1; teamId: TeamId; message: TeamMessageSnapshot }
+'team/message/queued': { version: 2; teamId: TeamId; message: TeamMessageSnapshot }
 ```
 
 类型:[TeamId](subsystems/agent-team.zh.md) · [TeamMessageSnapshot](subsystems/agent-team.zh.md)
@@ -822,7 +822,7 @@ export type SessionEvent<T extends SessionEventType = SessionEventType> = {
 
 ```ts persistence-catalog
 /** Whole shared-task value, stored only in the Team Lead Session. */
-'team/task': { version: 1; teamId: TeamId; task: TeamTaskSnapshot }
+'team/task': { version: 2; teamId: TeamId; task: TeamTaskSnapshot }
 ```
 
 类型:[TeamId](subsystems/agent-team.zh.md) · [TeamTaskSnapshot](subsystems/agent-team.zh.md)

+ 27 - 24
packages/experimental/agent-team/src/mailbox.ts

@@ -26,7 +26,6 @@ import type {
 /** Owns every process-local state transition for the durable Team mailbox. */
 export class TeamMailbox {
   private readonly dispatchTails = new Map<SessionId, Promise<void>>()
-  private readonly activeDispatches = new Map<SessionId, TeamMessageSnapshot>()
   private readonly inFlightMessages = new Set<TeamMessageId>()
   private readonly inFlightDispatches = new Set<Promise<unknown>>()
 
@@ -139,7 +138,7 @@ export class TeamMailbox {
         throw new TeamError(`team message exceeds ${this.maxMessageBytes} bytes`, 'TEAM_MESSAGE_TOO_LARGE')
       }
       await this.journal.appendAndFlush(root, 'team/message/queued', {
-        version: 1,
+        version: 2,
         teamId: TeamId(root.id),
         message: queued,
       })
@@ -187,12 +186,7 @@ export class TeamMailbox {
     message: TeamMessageSnapshot,
     signal: AbortSignal,
   ): Promise<boolean> {
-    const active = this.activeDispatches.get(message.targetId)
-    const live = message.targetId === root.id ? root : this.ctx.agents.get(message.targetId)
-    if (active !== undefined && live !== undefined && this.messagePrecedes(root, message.id, active.id)) {
-      return await this.dispatchOnce(root, message, signal)
-    }
-    return await this.serializeDispatch(message, () => this.dispatchOnce(root, message, signal))
+    return await this.serializeDispatch(message, () => this.dispatchThrough(root, message, signal))
   }
 
   /** Serialize delivery admission for one durable target in queued order. */
@@ -202,16 +196,8 @@ export class TeamMailbox {
   ): Promise<boolean> {
     const targetId = message.targetId
     const prior = this.dispatchTails.get(targetId) ?? Promise.resolve()
-    const dispatch = async (): Promise<boolean> => {
-      this.activeDispatches.set(targetId, message)
-      try {
-        return await operation()
-      } finally {
-        this.activeDispatches.delete(targetId)
-      }
-    }
     /* v8 ignore next -- dispatch tails absorb rejection, so the recovery callback is a fail-safe backstop. */
-    const run = prior.then(dispatch, dispatch)
+    const run = prior.then(operation, operation)
     /* v8 ignore next -- dispatchOnce contains delivery failures and serializeDispatch itself does not throw. */
     const tail = run.then(() => undefined, () => undefined)
     this.dispatchTails.set(targetId, tail)
@@ -222,6 +208,29 @@ export class TeamMailbox {
     }
   }
 
+  /** Deliver every pending target message through `message` in durable queue order. */
+  private async dispatchThrough(
+    root: Agent,
+    message: TeamMessageSnapshot,
+    signal: AbortSignal,
+  ): Promise<boolean> {
+    const state = this.journal.state(root)
+    const pending = state.messages.filter(candidate =>
+      candidate.targetId === message.targetId && !state.delivered.includes(candidate.id))
+    const requested = pending.findIndex(candidate => candidate.id === message.id)
+    if (requested < 0) return state.delivered.includes(message.id)
+    for (const candidate of pending.slice(0, requested + 1)) {
+      const ownsInFlight = !this.inFlightMessages.has(candidate.id)
+      if (ownsInFlight) this.inFlightMessages.add(candidate.id)
+      try {
+        if (!await this.dispatchOnce(root, candidate, signal)) return false
+      } finally {
+        if (ownsInFlight) this.inFlightMessages.delete(candidate.id)
+      }
+    }
+    return true
+  }
+
   /** Attempt one queued delivery after target-local ordering admits it. */
   private async dispatchOnce(root: Agent, message: TeamMessageSnapshot, signal: AbortSignal): Promise<boolean> {
     try {
@@ -260,12 +269,6 @@ export class TeamMailbox {
     }
   }
 
-  /** Whether `left` was durably queued before `right` in one Lead log. */
-  private messagePrecedes(root: Agent, left: TeamMessageId, right: TeamMessageId): boolean {
-    const ids = this.journal.state(root).messages.map(message => message.id)
-    return ids.indexOf(left) < ids.indexOf(right)
-  }
-
   /** Flush one live target receipt before the Lead records its delivered edge. */
   private async checkpointDelivered(
     root: Agent,
@@ -286,7 +289,7 @@ export class TeamMailbox {
       const queued = state.messages.find(message => message.id === messageId)
       if (queued === undefined || queued.targetId !== targetId) return
       await this.journal.appendAndFlush(root, 'team/message/delivered', {
-        version: 1,
+        version: 2,
         teamId: TeamId(root.id),
         messageId,
         targetId,

+ 6 - 6
packages/experimental/agent-team/src/projection.ts

@@ -99,25 +99,25 @@ const teamEventSelectorSchema = z.object({
 }).loose()
 
 const teamMemberEventSchema = z.object({
-  version: z.literal(1),
+  version: z.literal(2),
   teamId: teamIdSchema,
   member: teamMemberSnapshotSchema,
 }).strict() as z.ZodType<SessionEventMap['team/member']>
 
 const teamTaskEventSchema = z.object({
-  version: z.literal(1),
+  version: z.literal(2),
   teamId: teamIdSchema,
   task: teamTaskSnapshotSchema,
 }).strict() as z.ZodType<SessionEventMap['team/task']>
 
 const teamMessageQueuedEventSchema = z.object({
-  version: z.literal(1),
+  version: z.literal(2),
   teamId: teamIdSchema,
   message: teamMessageSnapshotSchema,
 }).strict() as z.ZodType<SessionEventMap['team/message/queued']>
 
 const teamMessageDeliveredEventSchema = z.object({
-  version: z.literal(1),
+  version: z.literal(2),
   teamId: teamIdSchema,
   messageId: teamMessageIdSchema,
   targetId: sessionIdSchema,
@@ -224,7 +224,7 @@ function applyProjectionEvent(state: TeamProjectionState, event: SessionEvent):
   try {
     const selector = parsePersisted(event.type, teamEventSelectorSchema, event.data)
     if (selector.teamId !== state.id) return
-    if (selector.version !== 1) {
+    if (selector.version !== 2) {
       throw new Error(`unsupported Agent Teams event version ${String(selector.version)}`)
     }
     applyCurrentTeamEvent(state, parseCurrentTeamEvent(event))
@@ -306,7 +306,7 @@ function applyCurrentTeamEvent(state: TeamState, event: TeamSessionEvent): void
 /** Host-only Team projection selected by the projected Session identity. */
 export const teamProjectionDefinition = {
   key: 'agentTeam',
-  stateVersion: 2,
+  stateVersion: 3,
   stateSchema: teamProjectionEntrySchema,
   init: header => emptyTeamState(header.id),
   apply: (state, event) => {

+ 3 - 3
packages/experimental/agent-team/src/roster.ts

@@ -274,7 +274,7 @@ export class TeamRoster {
       if (state.members.length >= this.maxMembers) {
         throw new TeamError(`Team member limit ${this.maxMembers} reached`, 'TEAM_MEMBER_LIMIT')
       }
-      await this.journal.appendAndFlush(root, 'team/member', { version: 1, teamId: TeamId(root.id), member })
+      await this.journal.appendAndFlush(root, 'team/member', { version: 2, teamId: TeamId(root.id), member })
     })
 
     let started: ContinuableStart
@@ -424,7 +424,7 @@ export class TeamRoster {
           ...phase === 'failed' ? { error: failure } : {},
         }
         await this.journal.appendAndFlush(root, 'team/member', {
-          version: 1,
+          version: 2,
           teamId: TeamId(root.id),
           member: settled,
         })
@@ -472,7 +472,7 @@ export class TeamRoster {
       }
       if (current.phase !== 'provisioning') return current.phase
       await this.journal.appendAndFlush(root, 'team/member', {
-        version: 1,
+        version: 2,
         teamId: TeamId(root.id),
         member: terminal,
       })

+ 2 - 2
packages/experimental/agent-team/src/task-board.ts

@@ -67,7 +67,7 @@ export class TeamTaskBoard {
         writeScopes: this.writeScopes(request.writeScopes ?? []),
       }
       this.assertTaskGraph(state, task)
-      await this.journal.appendAndFlush(root, 'team/task', { version: 1, teamId: TeamId(root.id), task })
+      await this.journal.appendAndFlush(root, 'team/task', { version: 2, teamId: TeamId(root.id), task })
       return this.taskView(root, state, task)
     })
   }
@@ -209,7 +209,7 @@ export class TeamTaskBoard {
         revision: current.revision + 1,
       }
       this.assertTaskGraph(state, task)
-      await this.journal.appendAndFlush(root, 'team/task', { version: 1, teamId: TeamId(root.id), task })
+      await this.journal.appendAndFlush(root, 'team/task', { version: 2, teamId: TeamId(root.id), task })
       return this.taskView(root, state, task)
     })
   }

+ 4 - 4
packages/experimental/agent-team/src/types.ts

@@ -218,14 +218,14 @@ export interface TeamWaitResult {
 declare module '@deepseek-ai/dsh-session/types' {
   interface SessionEventMap {
     /** Whole teammate lifecycle value, stored only in the Team Lead Session. */
-    'team/member': { version: 1; teamId: TeamId; member: TeamMemberSnapshot }
+    'team/member': { version: 2; teamId: TeamId; member: TeamMemberSnapshot }
     /** Whole shared-task value, stored only in the Team Lead Session. */
-    'team/task': { version: 1; teamId: TeamId; task: TeamTaskSnapshot }
+    'team/task': { version: 2; teamId: TeamId; task: TeamTaskSnapshot }
     /** Durable mailbox enqueue, stored before delivery is attempted. */
-    'team/message/queued': { version: 1; teamId: TeamId; message: TeamMessageSnapshot }
+    'team/message/queued': { version: 2; teamId: TeamId; message: TeamMessageSnapshot }
     /** Durable acknowledgement that the target Session recorded the message. */
     'team/message/delivered': {
-      version: 1
+      version: 2
       teamId: TeamId
       messageId: TeamMessageId
       targetId: SessionId

+ 3 - 3
packages/experimental/agent-team/tests/invariant.spec.ts

@@ -30,13 +30,13 @@ describe('Agent Teams stream invariant', () => {
       phase: 'provisioning' as const,
     }
     expect(() => {
-      session.append('team/member', { version: 1, teamId: TeamId(session.id), member })
+      session.append('team/member', { version: 2, teamId: TeamId(session.id), member })
     }).not.toThrow()
 
     const invalid = ctx.sessions.create(SessionId('team-invariant-invalid'))
     expect(() => {
       invalid.append('team/member', {
-        version: 1,
+        version: 2,
         teamId: TeamId(invalid.id),
         member: { ...member, phase: 'active' },
       })
@@ -53,7 +53,7 @@ describe('Agent Teams stream invariant', () => {
 
     expect(() => {
       session.append('team/task', {
-        version: 1,
+        version: 2,
         teamId: TeamId(session.id),
         task: {
           id: TeamTaskId('task-1'),

+ 9 - 9
packages/experimental/agent-team/tests/persistence.spec.ts

@@ -178,12 +178,12 @@ for (const backend of backends) {
       await Promise.resolve()
 
       activeRoot.session.append('team/member', {
-        version: 1,
+        version: 2,
         teamId: TeamId(activeRoot.id),
         member: provisioning(childId, 'recoverable'),
       })
       failedRoot.session.append('team/member', {
-        version: 1,
+        version: 2,
         teamId: TeamId(failedRoot.id),
         member: provisioning(SessionId(`${backend.name}-missing`), 'missing'),
       })
@@ -248,7 +248,7 @@ for (const backend of backends) {
       await Promise.resolve()
       await Promise.resolve()
       root.session.append('team/member', {
-        version: 1,
+        version: 2,
         teamId: TeamId(root.id),
         member: provisioning(childId, 'pending-worker'),
       })
@@ -295,8 +295,8 @@ for (const backend of backends) {
         signal: SIGNAL,
       })
       await vi.waitFor(() => { expect(first.ctx.agents.get(started.member.id)).toBeUndefined() }, { timeout: 5_000 })
-      vi.spyOn(first.ctx.sessionPersistence, 'inspect')
-        .mockRejectedValueOnce(new Error('temporary target inspection failure'))
+      vi.spyOn(first.ctx.sessionPersistence, 'open')
+        .mockRejectedValueOnce(new Error('temporary target read failure'))
       const queued = await first.ctx.agentTeams.sendMessage(firstLead, {
         target: 'mail-worker',
         content: [{ type: 'text', text: 'durable retry context' }],
@@ -376,7 +376,7 @@ for (const backend of backends) {
         content: [{ type: 'text', text: 'already recorded before acknowledgement' }],
       }
       firstLead.session.append('team/message/queued', {
-        version: 1,
+        version: 2,
         teamId: TeamId(rootId),
         message: queued,
       })
@@ -428,17 +428,17 @@ for (const backend of backends) {
         content: [{ type: 'text', text: 'already durable in target inbox' }],
       }
       root.session.append('team/member', {
-        version: 1,
+        version: 2,
         teamId: TeamId(root.id),
         member: provisioned,
       })
       root.session.append('team/member', {
-        version: 1,
+        version: 2,
         teamId: TeamId(root.id),
         member: active,
       })
       root.session.append('team/message/queued', {
-        version: 1,
+        version: 2,
         teamId: TeamId(root.id),
         message: queued,
       })

+ 43 - 43
packages/experimental/agent-team/tests/projection-events.spec.ts

@@ -79,15 +79,15 @@ function message(overrides: Partial<TeamMessageSnapshot> = {}): TeamMessageSnaps
 describe('Agent Teams projection events', () => {
   it('projects current-team records independently from inherited records', () => {
     const records: SessionEvent[] = [
-      event('team/member', { version: 1, teamId: TeamId('ancestor'), member: member() }, SessionSeq(0)),
-      event('team/member', { version: 1, teamId: TEAM, member: member() }, SessionSeq(1)),
+      event('team/member', { version: 2, teamId: TeamId('ancestor'), member: member() }, SessionSeq(0)),
+      event('team/member', { version: 2, teamId: TEAM, member: member() }, SessionSeq(1)),
       event('team/member', {
-        version: 1,
+        version: 2,
         teamId: TEAM,
         member: member({ phase: 'active' }),
       }, SessionSeq(2)),
-      event('team/task', { version: 1, teamId: TEAM, task: task({ id: TeamTaskId('task-7') }) }, SessionSeq(3)),
-      event('team/message/queued', { version: 1, teamId: TEAM, message: message() }, SessionSeq(4)),
+      event('team/task', { version: 2, teamId: TEAM, task: task({ id: TeamTaskId('task-7') }) }, SessionSeq(3)),
+      event('team/message/queued', { version: 2, teamId: TEAM, message: message() }, SessionSeq(4)),
     ]
     const projected = project(ROOT, records)
     const state = teamState(projected)
@@ -103,53 +103,53 @@ describe('Agent Teams projection events', () => {
   })
 
   it('enforces teammate identity and lifecycle', () => {
-    const base = event('team/member', { version: 1, teamId: TEAM, member: member() }, SessionSeq(0))
+    const base = event('team/member', { version: 2, teamId: TEAM, member: member() }, SessionSeq(0))
     expect(() => projectTeam(ROOT, [event('team/member', {
-      version: 1,
+      version: 2,
       teamId: TEAM,
       member: member({ phase: 'active' }),
     }, SessionSeq(0))])).toThrow(/must begin provisioning/)
     expect(() => projectTeam(ROOT, [base, event('team/member', {
-      version: 1,
+      version: 2,
       teamId: TEAM,
       member: member({ name: 'renamed', phase: 'active' }),
     }, SessionSeq(1))])).toThrow(/immutable identity/)
     expect(() => projectTeam(ROOT, [base, event('team/member', {
-      version: 1,
+      version: 2,
       teamId: TEAM,
       member: member({ phase: 'active' }),
     }, SessionSeq(1)), event('team/member', {
-      version: 1,
+      version: 2,
       teamId: TEAM,
       member: member({ phase: 'failed' }),
     }, SessionSeq(2))])).toThrow(/invalid active -> failed/)
 
     const duplicateName = member({ id: SessionId('child-b') })
     expect(() => projectTeam(ROOT, [base, event('team/member', {
-      version: 1,
+      version: 2,
       teamId: TEAM,
       member: duplicateName,
     }, SessionSeq(1))])).toThrow(/name .* reused/)
   })
 
   it('enforces task revision continuity', () => {
-    const first = event('team/task', { version: 1, teamId: TEAM, task: task() }, SessionSeq(0))
+    const first = event('team/task', { version: 2, teamId: TEAM, task: task() }, SessionSeq(0))
     expect(() => projectTeam(ROOT, [event('team/task', {
-      version: 1,
+      version: 2,
       teamId: TEAM,
       task: task({ revision: 2 }),
     }, SessionSeq(0))])).toThrow(/begin at revision 1/)
     expect(() => projectTeam(ROOT, [first, event('team/task', {
-      version: 1,
+      version: 2,
       teamId: TEAM,
       task: task({ revision: 3 }),
     }, SessionSeq(1))])).toThrow(/revision is not contiguous/)
   })
 
   it('rejects every invalid persisted task dependency relation', () => {
-    const first = event('team/task', { version: 1, teamId: TEAM, task: task() }, SessionSeq(0))
+    const first = event('team/task', { version: 2, teamId: TEAM, task: task() }, SessionSeq(0))
     const second = event('team/task', {
-      version: 1,
+      version: 2,
       teamId: TEAM,
       task: task({
         id: TeamTaskId('task-2'),
@@ -159,7 +159,7 @@ describe('Agent Teams projection events', () => {
     const invalid: Array<{ records: SessionEvent[]; message: RegExp }> = [
       {
         records: [event('team/task', {
-          version: 1,
+          version: 2,
           teamId: TEAM,
           task: task({ blockedBy: [TeamTaskId('missing')] }),
         }, SessionSeq(0))],
@@ -167,7 +167,7 @@ describe('Agent Teams projection events', () => {
       },
       {
         records: [event('team/task', {
-          version: 1,
+          version: 2,
           teamId: TEAM,
           task: task({ blockedBy: [TeamTaskId('task-1')] }),
         }, SessionSeq(0))],
@@ -182,7 +182,7 @@ describe('Agent Teams projection events', () => {
       },
       {
         records: [first, second, event('team/task', {
-          version: 1,
+          version: 2,
           teamId: TEAM,
           task: task({ revision: 2, blockedBy: [TeamTaskId('task-2')] }),
         }, SessionSeq(2))],
@@ -190,7 +190,7 @@ describe('Agent Teams projection events', () => {
       },
       {
         records: [first, second, event('team/task', {
-          version: 1,
+          version: 2,
           teamId: TEAM,
           task: task({ revision: 2, status: 'deleted' }),
         }, SessionSeq(2))],
@@ -205,7 +205,7 @@ describe('Agent Teams projection events', () => {
 
   it('leaves numeric allocation unchanged for a branded nonstandard task id', () => {
     const state = projectTeam(ROOT, [event('team/task', {
-      version: 1,
+      version: 2,
       teamId: TEAM,
       task: task({ id: TeamTaskId('external-task') }),
     }, SessionSeq(0))])
@@ -214,16 +214,16 @@ describe('Agent Teams projection events', () => {
 
   it('rejects a persisted numeric task id outside the safe integer range', () => {
     expect(() => projectTeam(ROOT, [event('team/task', {
-      version: 1,
+      version: 2,
       teamId: TEAM,
       task: task({ id: TeamTaskId('task-9007199254740992') }),
     }, SessionSeq(0))])).toThrow(/persisted Agent Teams team\/task payload is invalid/)
   })
 
   it('enforces mailbox queue and acknowledgement relations', () => {
-    const queued = event('team/message/queued', { version: 1, teamId: TEAM, message: message() }, SessionSeq(0))
+    const queued = event('team/message/queued', { version: 2, teamId: TEAM, message: message() }, SessionSeq(0))
     const delivered = event('team/message/delivered', {
-      version: 1,
+      version: 2,
       teamId: TEAM,
       messageId: TeamMessageId('message-1'),
       targetId: CHILD,
@@ -241,42 +241,42 @@ describe('Agent Teams projection events', () => {
   it('validates every current-version persisted payload before projecting it', () => {
     const malformed = [
       {
-        ...event('team/member', { version: 1, teamId: TEAM, member: member() }, SessionSeq(0)),
-        data: { version: 1, teamId: TEAM, member: { ...member(), name: 42 } },
+        ...event('team/member', { version: 2, teamId: TEAM, member: member() }, SessionSeq(0)),
+        data: { version: 2, teamId: TEAM, member: { ...member(), name: 42 } },
       },
       {
-        ...event('team/task', { version: 1, teamId: TEAM, task: task() }, SessionSeq(0)),
-        data: { version: 1, teamId: TEAM, task: { ...task(), blockedBy: [42] } },
+        ...event('team/task', { version: 2, teamId: TEAM, task: task() }, SessionSeq(0)),
+        data: { version: 2, teamId: TEAM, task: { ...task(), blockedBy: [42] } },
       },
       {
-        ...event('team/message/queued', { version: 1, teamId: TEAM, message: message() }, SessionSeq(0)),
+        ...event('team/message/queued', { version: 2, teamId: TEAM, message: message() }, SessionSeq(0)),
         data: {
-          version: 1,
+          version: 2,
           teamId: TEAM,
           message: { ...message(), content: [{ type: 'text', text: 42 }] },
         },
       },
       {
         ...event('team/message/delivered', {
-          version: 1,
+          version: 2,
           teamId: TEAM,
           messageId: TeamMessageId('message-1'),
           targetId: CHILD,
         }, SessionSeq(0)),
         data: {
-          version: 1,
+          version: 2,
           teamId: TEAM,
           messageId: TeamMessageId('message-1'),
           targetId: 42,
         },
       },
       {
-        ...event('team/member', { version: 1, teamId: TEAM, member: member() }, SessionSeq(0)),
-        data: { version: 1, teamId: TEAM, member: member(), unexpected: true },
+        ...event('team/member', { version: 2, teamId: TEAM, member: member() }, SessionSeq(0)),
+        data: { version: 2, teamId: TEAM, member: member(), unexpected: true },
       },
       {
-        ...event('team/task', { version: 1, teamId: TEAM, task: task() }, SessionSeq(0)),
-        data: { version: 1, teamId: 42, task: task() },
+        ...event('team/task', { version: 2, teamId: TEAM, task: task() }, SessionSeq(0)),
+        data: { version: 2, teamId: 42, task: task() },
       },
     ] as unknown as SessionEvent[]
 
@@ -289,7 +289,7 @@ describe('Agent Teams projection events', () => {
   it('retains merge-extensible content blocks while rejecting malformed core variants', () => {
     const extension = { type: 'plugin/custom', payload: { value: 1 } } as never
     const state = projectTeam(ROOT, [event('team/message/queued', {
-      version: 1,
+      version: 2,
       teamId: TEAM,
       message: message({ content: [extension] }),
     }, SessionSeq(0))])
@@ -298,23 +298,23 @@ describe('Agent Teams projection events', () => {
 
   it('records unsupported event versions without applying them', () => {
     const invalid = event('team/task', {
-      version: 2 as 1,
+      version: 1 as 2,
       teamId: TEAM,
       task: task(),
     }, SessionSeq(0))
     const later = event('team/task', {
-      version: 1,
+      version: 2,
       teamId: TEAM,
       task: task(),
     }, SessionSeq(1))
     const state = project(ROOT, [invalid, later])
-    expect(state.failure).toMatch(/unsupported Agent Teams event version 2/)
+    expect(state.failure).toMatch(/unsupported Agent Teams event version 1/)
     expect(isEmptyState(state)).toBe(true)
   })
 
   it('isolates unsupported inherited Team records from the current Team', () => {
     const inherited = event('team/task', {
-      version: 2 as 1,
+      version: 1 as 2,
       teamId: TeamId('ancestor'),
       task: task(),
     }, SessionSeq(0))
@@ -326,12 +326,12 @@ describe('Agent Teams projection events', () => {
   it('ignores malformed current-version records inherited from another Team', () => {
     const inherited = {
       ...event('team/task', {
-        version: 1,
+        version: 2,
         teamId: TeamId('ancestor'),
         task: task(),
       }, SessionSeq(0)),
       data: {
-        version: 1,
+        version: 2,
         teamId: TeamId('ancestor'),
         task: { ...task(), subject: 42 },
       },

+ 51 - 63
packages/experimental/agent-team/tests/team.spec.ts

@@ -193,7 +193,7 @@ describe('Team identity and provisioning', () => {
       phase: 'provisioning' as const,
     }
     lead.session.append('team/member', {
-      version: 1,
+      version: 2,
       teamId: TeamId(lead.id),
       member: provisioning,
     })
@@ -354,7 +354,7 @@ describe('Team identity and provisioning', () => {
       const provisioning = durable(second.lead).members[0]
       if (provisioning === undefined) throw new Error('missing provisioning edge')
       second.lead.session.append('team/member', {
-        version: 1,
+        version: 2,
         teamId: TeamId(second.lead.id),
         member: { ...provisioning, phase: 'active' },
       })
@@ -555,7 +555,7 @@ describe('Team shared task DAG', () => {
     const { ctx, lead } = await setup([])
     const id = TeamTaskId(`task-${Number.MAX_SAFE_INTEGER}`)
     lead.session.append('team/task', {
-      version: 1,
+      version: 2,
       teamId: TeamId(lead.id),
       task: {
         id,
@@ -948,7 +948,7 @@ describe('Team mailbox and waiting', () => {
       content: content('progress report'),
     }
     lead.session.append('team/message/queued', {
-      version: 1,
+      version: 2,
       teamId: TeamId(lead.id),
       message,
     })
@@ -1034,7 +1034,7 @@ describe('Team mailbox and waiting', () => {
       content: content('durable pending receipt'),
     }
     lead.session.append('team/message/queued', {
-      version: 1,
+      version: 2,
       teamId: TeamId(lead.id),
       message,
     })
@@ -1070,7 +1070,7 @@ describe('Team mailbox and waiting', () => {
       content: content('canceled before checkpoint'),
     }
     lead.session.append('team/message/queued', {
-      version: 1,
+      version: 2,
       teamId: TeamId(lead.id),
       message: disappearing,
     })
@@ -1137,14 +1137,14 @@ describe('Team mailbox and waiting', () => {
   })
 
   it('serializes concurrent Steer delivery admission for one target', async () => {
-    const { ctx, lead } = await setup([textResponse('target initial')])
-    const target = await spawn(ctx, lead, 'ordered-target')
-    await waitNoAgent(ctx, target.member.id)
+    const { ctx, lead } = await setup(['hang'])
+    const started = await spawn(ctx, lead, 'ordered-target')
+    const target = await waitRunning(ctx, started.member.id)
     const entered = Promise.withResolvers<undefined>()
     const release = Promise.withResolvers<undefined>()
     const admitted: string[] = []
     vi.spyOn(ctx.subagents as unknown as HostPromptDeliverer, deliverSubagentPrompt)
-      .mockImplementation(async (_parent, _childId, blocks) => {
+      .mockImplementation(async (_parent, _childId, blocks, source) => {
         const last = blocks.at(-1)
         const text = last?.type === 'text' ? last.text : ''
         admitted.push(text)
@@ -1152,7 +1152,9 @@ describe('Team mailbox and waiting', () => {
           entered.resolve(undefined)
           await release.promise
         }
-        return createUserMessage({ content: blocks, source: { kind: 'user' } }).id
+        const input = createUserMessage({ content: blocks, source })
+        target.inject(input)
+        return input.id
       })
 
     const first = ctx.agentTeams.sendMessage(lead, {
@@ -1163,7 +1165,7 @@ describe('Team mailbox and waiting', () => {
     const second = ctx.agentTeams.sendMessage(lead, {
       target: 'ordered-target', content: content('second steer'), signal: SIGNAL,
     }).finally(() => { secondSettled = true })
-    await new Promise<void>((resolve) => { setTimeout(resolve, 0) })
+    await vi.waitFor(() => { expect(durable(lead).pendingMessages).toHaveLength(2) })
     expect(admitted).toEqual(['first steer'])
     expect(secondSettled).toBe(false)
 
@@ -1173,58 +1175,43 @@ describe('Team mailbox and waiting', () => {
       { status: 'accepted' },
     ])
     expect(admitted).toEqual(['first steer', 'second steer'])
+
+    ctx.agentTeams.interrupt(lead, 'ordered-target')
+    target.cancel({ kind: 'parent' })
+    await waitNoAgent(ctx, target.id)
   })
 
-  it('admits an earlier durable message ahead of a later in-flight resume', async () => {
-    const { ctx, lead } = await setup(['hang'])
+  it('delivers persisted mail before the later message that cold-resumes its target', async () => {
+    const { ctx, lead } = await setup([textResponse('target initial'), 'hang', 'hang'])
     const started = await spawn(ctx, lead, 'reordered-target')
-    const target = await waitRunning(ctx, started.member.id)
+    await waitNoAgent(ctx, started.member.id)
     const earlier: TeamMessageSnapshot = {
       id: TeamMessageId('earlier-message'),
       senderId: lead.id,
       senderName: 'lead',
-      targetId: target.id,
+      targetId: started.member.id,
       content: content('earlier steer'),
     }
-    const later: TeamMessageSnapshot = {
-      ...earlier,
-      id: TeamMessageId('later-message'),
-      content: content('later steer'),
-    }
-    for (const message of [earlier, later]) {
-      lead.session.append('team/message/queued', {
-        version: 1,
-        teamId: TeamId(lead.id),
-        message,
-      })
-    }
-
-    const laterEntered = Promise.withResolvers<undefined>()
-    const releaseLater = Promise.withResolvers<undefined>()
-    const admitted: string[] = []
-    vi.spyOn(ctx.subagents as unknown as HostPromptDeliverer, deliverSubagentPrompt)
-      .mockImplementation(async (_parent, _childId, blocks, source) => {
-        const last = blocks.at(-1)
-        const text = last?.type === 'text' ? last.text : ''
-        if (text === 'later steer') {
-          laterEntered.resolve(undefined)
-          await releaseLater.promise
-        }
-        const input = createUserMessage({ content: blocks, source })
-        target.inject(input)
-        admitted.push(text)
-        return input.id
-      })
-
-    const laterDispatch = teamInternals(ctx).mailbox.tryDispatch(lead, later, SIGNAL)
-    await laterEntered.promise
-    await expect(teamInternals(ctx).mailbox.tryDispatch(lead, earlier, SIGNAL)).resolves.toBe(true)
-    expect(admitted).toEqual(['earlier steer'])
+    lead.session.append('team/message/queued', {
+      version: 2,
+      teamId: TeamId(lead.id),
+      message: earlier,
+    })
+    await ctx.sessions.flush(lead.session)
 
-    releaseLater.resolve(undefined)
-    await expect(laterDispatch).resolves.toBe(true)
-    expect(admitted).toEqual(['earlier steer', 'later steer'])
-    expect(durable(lead).pendingMessages).toEqual([])
+    const later = await ctx.agentTeams.sendMessage(lead, {
+      target: 'reordered-target', content: content('later steer'), signal: SIGNAL,
+    })
+    expect(later.status).toBe('accepted')
+    const target = await waitRunning(ctx, started.member.id)
+    await vi.waitFor(() => {
+      const accepted = target.session.snapshotEvents().flatMap(event => event.type === 'agent/inbox/spliced'
+        ? event.data.inserted.flatMap(message => message.source.kind === 'team-message'
+          ? [message.source.messageId]
+          : [])
+        : [])
+      expect(accepted).toEqual([earlier.id, later.messageId])
+    })
 
     ctx.agentTeams.interrupt(lead, 'reordered-target')
     target.cancel({ kind: 'parent' })
@@ -1244,7 +1231,7 @@ describe('Team mailbox and waiting', () => {
       content: content('already in live history'),
     }
     lead.session.append('team/message/queued', {
-      version: 1, teamId: TeamId(lead.id), message,
+      version: 2, teamId: TeamId(lead.id), message,
     })
     await ctx.sessions.flush(lead.session)
     live.session.append('user/message', createUserMessage({
@@ -1269,13 +1256,14 @@ describe('Team mailbox and waiting', () => {
     }), { surfaceOp: 'append' })
     await expect(internal.tryDispatch(lead, message, SIGNAL)).resolves.toBe(true)
     await internal.markDelivered(lead, message.id, live.id)
+    await expect(internal.tryDispatch(lead, message, SIGNAL)).resolves.toBe(true)
 
     const wrongTarget: TeamMessageSnapshot = {
       ...message,
       id: TeamMessageId('wrong-target-message'),
     }
     lead.session.append('team/message/queued', {
-      version: 1, teamId: TeamId(lead.id), message: wrongTarget,
+      version: 2, teamId: TeamId(lead.id), message: wrongTarget,
     })
     await ctx.sessions.flush(lead.session)
     await internal.markDelivered(lead, wrongTarget.id, SessionId('wrong-target'))
@@ -1383,7 +1371,7 @@ describe('Team mailbox and waiting', () => {
     await expect(ctx.agentTeams.sendMessage(lead, {
       target: 'target', content: content('x'.repeat(300)), signal: SIGNAL,
     })).rejects.toMatchObject({ code: 'TEAM_MESSAGE_TOO_LARGE' })
-    vi.spyOn(ctx.sessionPersistence, 'inspect').mockRejectedValueOnce(new Error('temporary inspection failure'))
+    vi.spyOn(ctx.sessionPersistence, 'open').mockRejectedValueOnce(new Error('temporary read failure'))
     const queued = await ctx.agentTeams.sendMessage(lead, {
       target: 'target', content: content('one'), signal: SIGNAL,
     })
@@ -1589,7 +1577,7 @@ describe('Team mailbox and waiting', () => {
       phase: 'provisioning' as const,
     }
     lead.session.append('team/member', {
-      version: 1,
+      version: 2,
       teamId: TeamId(lead.id),
       member,
     })
@@ -1602,7 +1590,7 @@ describe('Team mailbox and waiting', () => {
     })
     await waitRunning(ctx, childId)
     lead.session.append('team/member', {
-      version: 1,
+      version: 2,
       teamId: TeamId(lead.id),
       member: {
         ...member,
@@ -1669,7 +1657,7 @@ describe('Team mailbox and waiting', () => {
       content: content('acknowledge before disposal'),
     }
     lead.session.append('team/message/queued', {
-      version: 1,
+      version: 2,
       teamId: TeamId(lead.id),
       message,
     })
@@ -1822,7 +1810,7 @@ describe('Team mailbox and waiting', () => {
       phase: 'provisioning' as const,
     }
     first.lead.session.append('team/member', {
-      version: 1, teamId: TeamId(first.lead.id), member: provisioning,
+      version: 2, teamId: TeamId(first.lead.id), member: provisioning,
     })
     const reconcileFirst = teamInternals(first.ctx).roster
     await reconcileFirst.reconcileProvisioning(first.lead, SIGNAL)
@@ -1842,7 +1830,7 @@ describe('Team mailbox and waiting', () => {
     const childId = SessionId('concurrently-settled-child')
     const member = { ...provisioning, id: childId, name: 'concurrent-child' }
     second.lead.session.append('team/member', {
-      version: 1, teamId: TeamId(second.lead.id), member,
+      version: 2, teamId: TeamId(second.lead.id), member,
     })
     const entered = Promise.withResolvers<undefined>()
     const release = Promise.withResolvers<undefined>()
@@ -1855,7 +1843,7 @@ describe('Team mailbox and waiting', () => {
     const reconciling = reconcileSecond.reconcileProvisioning(second.lead, SIGNAL)
     await entered.promise
     second.lead.session.append('team/member', {
-      version: 1,
+      version: 2,
       teamId: TeamId(second.lead.id),
       member: { ...member, phase: 'failed', error: 'settled elsewhere' },
     })

+ 9 - 1
packages/subagent/subagent/src/continuation.ts

@@ -134,7 +134,15 @@ export interface SubagentSendMessageOptions {
 
 /** Inputs shared by model steering and the human Queue adapter. */
 type ChildDeliveryOptions =
-  | { readonly delivery: 'steer'; readonly source?: MessageSource; readonly signal: AbortSignal }
+  | {
+    readonly delivery: 'steer'
+    /**
+     * A provided host source is preserved on the user message; omission attributes
+     * an adjacent-Agent message to the parent.
+     */
+    readonly source?: MessageSource
+    readonly signal: AbortSignal
+  }
   | { readonly delivery: 'queue'; readonly source: MessageSource; readonly signal: AbortSignal }
 
 /**