diff --git a/docs/persistence-catalog.i18n.yaml b/docs/persistence-catalog.i18n.yaml index 84d0af864b..b6e2f8ac7f 100644 --- a/docs/persistence-catalog.i18n.yaml +++ b/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 diff --git a/docs/persistence-catalog.md b/docs/persistence-catalog.md index 5b8ae5748c..1c0c691998 100644 --- a/docs/persistence-catalog.md +++ b/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) diff --git a/docs/persistence-catalog.zh.md b/docs/persistence-catalog.zh.md index 9a5d1b4020..4cc8ba5b7f 100644 --- a/docs/persistence-catalog.zh.md +++ b/docs/persistence-catalog.zh.md @@ -778,7 +778,7 @@ export type SessionEvent = { ```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 = { ```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 = { ```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 = { ```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) diff --git a/packages/experimental/agent-team/src/mailbox.ts b/packages/experimental/agent-team/src/mailbox.ts index 8af111cd76..acd842308a 100644 --- a/packages/experimental/agent-team/src/mailbox.ts +++ b/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>() - private readonly activeDispatches = new Map() private readonly inFlightMessages = new Set() private readonly inFlightDispatches = new Set>() @@ -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 { - 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 { const targetId = message.targetId const prior = this.dispatchTails.get(targetId) ?? Promise.resolve() - const dispatch = async (): Promise => { - 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 { + 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 { 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, diff --git a/packages/experimental/agent-team/src/projection.ts b/packages/experimental/agent-team/src/projection.ts index 2d60102197..85fdfd0e0d 100644 --- a/packages/experimental/agent-team/src/projection.ts +++ b/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 const teamTaskEventSchema = z.object({ - version: z.literal(1), + version: z.literal(2), teamId: teamIdSchema, task: teamTaskSnapshotSchema, }).strict() as z.ZodType const teamMessageQueuedEventSchema = z.object({ - version: z.literal(1), + version: z.literal(2), teamId: teamIdSchema, message: teamMessageSnapshotSchema, }).strict() as z.ZodType 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) => { diff --git a/packages/experimental/agent-team/src/roster.ts b/packages/experimental/agent-team/src/roster.ts index 7ff934b150..3d0ca8372e 100644 --- a/packages/experimental/agent-team/src/roster.ts +++ b/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, }) diff --git a/packages/experimental/agent-team/src/task-board.ts b/packages/experimental/agent-team/src/task-board.ts index 5dd235f940..cc7064ceb6 100644 --- a/packages/experimental/agent-team/src/task-board.ts +++ b/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) }) } diff --git a/packages/experimental/agent-team/src/types.ts b/packages/experimental/agent-team/src/types.ts index 5818a6bc5f..f7b8ca414f 100644 --- a/packages/experimental/agent-team/src/types.ts +++ b/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 diff --git a/packages/experimental/agent-team/tests/invariant.spec.ts b/packages/experimental/agent-team/tests/invariant.spec.ts index 04d99d75a5..5dd8dba256 100644 --- a/packages/experimental/agent-team/tests/invariant.spec.ts +++ b/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'), diff --git a/packages/experimental/agent-team/tests/persistence.spec.ts b/packages/experimental/agent-team/tests/persistence.spec.ts index 12f3e1b9f3..0f550c591c 100644 --- a/packages/experimental/agent-team/tests/persistence.spec.ts +++ b/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, }) diff --git a/packages/experimental/agent-team/tests/projection-events.spec.ts b/packages/experimental/agent-team/tests/projection-events.spec.ts index d78b160ff0..97187e4eb6 100644 --- a/packages/experimental/agent-team/tests/projection-events.spec.ts +++ b/packages/experimental/agent-team/tests/projection-events.spec.ts @@ -79,15 +79,15 @@ function message(overrides: Partial = {}): 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 }, }, diff --git a/packages/experimental/agent-team/tests/team.spec.ts b/packages/experimental/agent-team/tests/team.spec.ts index 1a11e8c06c..4f099a0916 100644 --- a/packages/experimental/agent-team/tests/team.spec.ts +++ b/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() const release = Promise.withResolvers() 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((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, - }) - } + lead.session.append('team/message/queued', { + version: 2, + teamId: TeamId(lead.id), + message: earlier, + }) + await ctx.sessions.flush(lead.session) - const laterEntered = Promise.withResolvers() - const releaseLater = Promise.withResolvers() - 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']) - - 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() const release = Promise.withResolvers() @@ -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' }, }) diff --git a/packages/subagent/subagent/src/continuation.ts b/packages/subagent/subagent/src/continuation.ts index 6fabb6ea33..f127bc51de 100644 --- a/packages/subagent/subagent/src/continuation.ts +++ b/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 } /**