|
|
@@ -11,7 +11,7 @@ import { SessionLogOffset, SessionId, type Session, type SessionEvent } from '@d
|
|
|
import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
|
|
|
import JsonlSessionPersistence from '@deepseek-ai/dsh-session-persistence-jsonl'
|
|
|
import SubagentService from '@deepseek-ai/dsh-subagent'
|
|
|
-import { queueSubagentPrompt, type HostPromptQueue } from '@deepseek-ai/dsh-subagent/internal'
|
|
|
+import { deliverSubagentPrompt, type HostPromptDeliverer } from '@deepseek-ai/dsh-subagent/internal'
|
|
|
import * as SubagentFork from '@deepseek-ai/dsh-subagent-fork-in-process'
|
|
|
import * as SubagentSpawn from '@deepseek-ai/dsh-subagent-spawn-in-process'
|
|
|
import { MockAdapter, textResponse } from '../../../core/agent-loop/tests/mock-adapter.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,
|
|
|
})
|
|
|
@@ -346,7 +346,7 @@ describe('Team identity and provisioning', () => {
|
|
|
diagnostics: ['string provider failure'],
|
|
|
})
|
|
|
await expect(first.ctx.agentTeams.sendMessage(first.lead, {
|
|
|
- target: 'string-failure', content: content('cannot deliver'), delivery: 'quiet', signal: SIGNAL,
|
|
|
+ target: 'string-failure', content: content('cannot deliver'), signal: SIGNAL,
|
|
|
})).rejects.toMatchObject({ code: 'TEAM_MEMBER_NOT_FOUND' })
|
|
|
|
|
|
const second = await setup([])
|
|
|
@@ -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,
|
|
|
@@ -594,7 +594,7 @@ describe('Team shared task DAG', () => {
|
|
|
})
|
|
|
|
|
|
it('enforces CAS, ownership, dependencies, transitions, and write-scope warnings', async () => {
|
|
|
- const { ctx, lead } = await setup(['hang', 'hang'])
|
|
|
+ const { ctx, lead } = await setup(['hang', 'hang', textResponse('beta integrated update')])
|
|
|
const firstMember = await spawn(ctx, lead, 'alpha')
|
|
|
const alpha = await waitRunning(ctx, firstMember.member.id)
|
|
|
const secondMember = await spawn(ctx, lead, 'beta')
|
|
|
@@ -938,29 +938,31 @@ describe('Team Remote API', () => {
|
|
|
})
|
|
|
|
|
|
describe('Team mailbox and waiting', () => {
|
|
|
- it('injects a quiet message addressed to the Lead and checkpoints its receipt', async () => {
|
|
|
- const { ctx, lead } = await setup([])
|
|
|
+ it('steers a message addressed to the Lead and checkpoints its receipt', async () => {
|
|
|
+ const { ctx, lead } = await setup(['hang'])
|
|
|
const message: TeamMessageSnapshot = {
|
|
|
- id: TeamMessageId('quiet-lead-message'),
|
|
|
+ id: TeamMessageId('steer-lead-message'),
|
|
|
senderId: SessionId('team-worker'),
|
|
|
senderName: 'worker',
|
|
|
targetId: lead.id,
|
|
|
- delivery: 'quiet',
|
|
|
- content: content('quiet report'),
|
|
|
+ content: content('progress report'),
|
|
|
}
|
|
|
lead.session.append('team/message/queued', {
|
|
|
- version: 1,
|
|
|
+ version: 2,
|
|
|
teamId: TeamId(lead.id),
|
|
|
message,
|
|
|
})
|
|
|
|
|
|
await expect(teamInternals(ctx).mailbox.tryDispatch(lead, message, SIGNAL)).resolves.toBe(true)
|
|
|
- expect(lead.inbox.nextStep.some(input => input.source.kind === 'team-message'
|
|
|
- && input.source.messageId === message.id)).toBe(true)
|
|
|
+ expect(lead.session.snapshotEvents().some(event => event.type === 'agent/inbox/spliced'
|
|
|
+ && event.data.inserted.some(input => input.source.kind === 'team-message'
|
|
|
+ && input.source.messageId === message.id))).toBe(true)
|
|
|
expect(durable(lead).pendingMessages).toEqual([])
|
|
|
+ lead.cancel({ kind: 'parent' })
|
|
|
+ await lead.whenIdle()
|
|
|
})
|
|
|
|
|
|
- it('acknowledges waking messages persisted by a busy Lead before model claim', async () => {
|
|
|
+ it('acknowledges steered messages persisted by a busy Lead before model claim', async () => {
|
|
|
const { ctx, lead, teamFiber } = await setup(['hang', 'hang'], { maxPendingMessagesPerMember: 1 })
|
|
|
const started = await spawn(ctx, lead, 'lead-reporter')
|
|
|
const reporter = await waitRunning(ctx, started.member.id)
|
|
|
@@ -968,10 +970,10 @@ describe('Team mailbox and waiting', () => {
|
|
|
await waitRunning(ctx, lead.id)
|
|
|
|
|
|
const first = await ctx.agentTeams.sendMessage(reporter, {
|
|
|
- target: 'lead', content: content('first wakeup report'), delivery: 'wakeup', signal: SIGNAL,
|
|
|
+ target: 'lead', content: content('first progress report'), signal: SIGNAL,
|
|
|
})
|
|
|
const second = await ctx.agentTeams.sendMessage(reporter, {
|
|
|
- target: 'lead', content: content('second wakeup report'), delivery: 'wakeup', signal: SIGNAL,
|
|
|
+ target: 'lead', content: content('second progress report'), signal: SIGNAL,
|
|
|
})
|
|
|
expect([first.status, second.status]).toEqual(['accepted', 'accepted'])
|
|
|
expect(lead.status).toBe('running')
|
|
|
@@ -1016,8 +1018,7 @@ describe('Team mailbox and waiting', () => {
|
|
|
const target = await waitRunning(ctx, started.member.id)
|
|
|
const immediate = await ctx.agentTeams.sendMessage(lead, {
|
|
|
target: 'pending-target',
|
|
|
- content: content('live quiet receipt'),
|
|
|
- delivery: 'quiet',
|
|
|
+ content: content('live steer receipt'),
|
|
|
signal: SIGNAL,
|
|
|
})
|
|
|
expect(immediate.status).toBe('accepted')
|
|
|
@@ -1030,11 +1031,10 @@ describe('Team mailbox and waiting', () => {
|
|
|
senderId: lead.id,
|
|
|
senderName: 'lead',
|
|
|
targetId: target.id,
|
|
|
- delivery: 'quiet',
|
|
|
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,
|
|
|
})
|
|
|
@@ -1098,7 +1098,7 @@ describe('Team mailbox and waiting', () => {
|
|
|
await waitNoAgent(ctx, target.id)
|
|
|
})
|
|
|
|
|
|
- it('acknowledges waking messages accepted by a busy target inbox', async () => {
|
|
|
+ it('acknowledges steered messages accepted by a busy target inbox', async () => {
|
|
|
const { ctx, lead } = await setup(['hang'], { maxPendingMessagesPerMember: 1 })
|
|
|
const started = await spawn(ctx, lead, 'busy-target')
|
|
|
const target = await waitRunning(ctx, started.member.id)
|
|
|
@@ -1110,24 +1110,24 @@ describe('Team mailbox and waiting', () => {
|
|
|
})
|
|
|
|
|
|
const first = await ctx.agentTeams.sendMessage(lead, {
|
|
|
- target: 'busy-target', content: content('first waking message'), delivery: 'wakeup', signal: SIGNAL,
|
|
|
+ target: 'busy-target', content: content('first steered message'), signal: SIGNAL,
|
|
|
})
|
|
|
|
|
|
expect(first.status).toBe('accepted')
|
|
|
expect(flushed).toEqual([lead.id, target.id, lead.id])
|
|
|
expect(durable(lead).pendingMessages).toEqual([])
|
|
|
- expect(target.inbox.nextTurn.some(message => message.source.kind === 'team-message'
|
|
|
+ expect(target.inbox.nextStep.some(message => message.source.kind === 'team-message'
|
|
|
&& message.source.messageId === first.messageId)).toBe(true)
|
|
|
|
|
|
flushed.length = 0
|
|
|
const second = await ctx.agentTeams.sendMessage(lead, {
|
|
|
- target: 'busy-target', content: content('second waking message'), delivery: 'wakeup', signal: SIGNAL,
|
|
|
+ target: 'busy-target', content: content('second steered message'), signal: SIGNAL,
|
|
|
})
|
|
|
|
|
|
expect(second.status).toBe('accepted')
|
|
|
expect(flushed).toEqual([lead.id, target.id, lead.id])
|
|
|
expect(durable(lead).pendingMessages).toEqual([])
|
|
|
- expect(target.inbox.nextTurn.filter(message => message.source.kind === 'team-message'
|
|
|
+ expect(target.inbox.nextStep.filter(message => message.source.kind === 'team-message'
|
|
|
&& (message.source.messageId === first.messageId || message.source.messageId === second.messageId)))
|
|
|
.toHaveLength(2)
|
|
|
|
|
|
@@ -1136,35 +1136,37 @@ describe('Team mailbox and waiting', () => {
|
|
|
await waitNoAgent(ctx, target.id)
|
|
|
})
|
|
|
|
|
|
- it('serializes concurrent waking 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)
|
|
|
+ it('serializes concurrent Steer delivery admission for one target', async () => {
|
|
|
+ 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 HostPromptQueue, queueSubagentPrompt)
|
|
|
- .mockImplementation(async (_parent, _childId, blocks) => {
|
|
|
+ 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 : ''
|
|
|
admitted.push(text)
|
|
|
- if (text === 'first waking') {
|
|
|
+ if (text === 'first steer') {
|
|
|
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, {
|
|
|
- target: 'ordered-target', content: content('first waking'), delivery: 'wakeup', signal: SIGNAL,
|
|
|
+ target: 'ordered-target', content: content('first steer'), signal: SIGNAL,
|
|
|
})
|
|
|
await entered.promise
|
|
|
let secondSettled = false
|
|
|
const second = ctx.agentTeams.sendMessage(lead, {
|
|
|
- target: 'ordered-target', content: content('second waking'), delivery: 'wakeup', signal: SIGNAL,
|
|
|
+ target: 'ordered-target', content: content('second steer'), signal: SIGNAL,
|
|
|
}).finally(() => { secondSettled = true })
|
|
|
- await new Promise<void>((resolve) => { setTimeout(resolve, 0) })
|
|
|
- expect(admitted).toEqual(['first waking'])
|
|
|
+ await vi.waitFor(() => { expect(durable(lead).pendingMessages).toHaveLength(2) })
|
|
|
+ expect(admitted).toEqual(['first steer'])
|
|
|
expect(secondSettled).toBe(false)
|
|
|
|
|
|
release.resolve(undefined)
|
|
|
@@ -1172,7 +1174,48 @@ describe('Team mailbox and waiting', () => {
|
|
|
{ status: 'accepted' },
|
|
|
{ status: 'accepted' },
|
|
|
])
|
|
|
- expect(admitted).toEqual(['first waking', 'second waking'])
|
|
|
+ expect(admitted).toEqual(['first steer', 'second steer'])
|
|
|
+
|
|
|
+ ctx.agentTeams.interrupt(lead, 'ordered-target')
|
|
|
+ target.cancel({ kind: 'parent' })
|
|
|
+ await waitNoAgent(ctx, target.id)
|
|
|
+ })
|
|
|
+
|
|
|
+ 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')
|
|
|
+ await waitNoAgent(ctx, started.member.id)
|
|
|
+ const earlier: TeamMessageSnapshot = {
|
|
|
+ id: TeamMessageId('earlier-message'),
|
|
|
+ senderId: lead.id,
|
|
|
+ senderName: 'lead',
|
|
|
+ targetId: started.member.id,
|
|
|
+ content: content('earlier steer'),
|
|
|
+ }
|
|
|
+ lead.session.append('team/message/queued', {
|
|
|
+ version: 2,
|
|
|
+ teamId: TeamId(lead.id),
|
|
|
+ message: earlier,
|
|
|
+ })
|
|
|
+ await ctx.sessions.flush(lead.session)
|
|
|
+
|
|
|
+ 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' })
|
|
|
+ await waitNoAgent(ctx, target.id)
|
|
|
})
|
|
|
|
|
|
it('deduplicates live target history and contains inspection and delivery failures', async () => {
|
|
|
@@ -1185,11 +1228,10 @@ describe('Team mailbox and waiting', () => {
|
|
|
senderId: lead.id,
|
|
|
senderName: 'lead',
|
|
|
targetId: live.id,
|
|
|
- delivery: 'wakeup',
|
|
|
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({
|
|
|
@@ -1214,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'))
|
|
|
@@ -1261,15 +1304,15 @@ describe('Team mailbox and waiting', () => {
|
|
|
await waitNoAgent(ctx, inactiveStarted.member.id)
|
|
|
const openRead = vi.spyOn(ctx.sessionPersistence, 'open').mockRejectedValueOnce(new Error('read unavailable'))
|
|
|
const uncertain = await ctx.agentTeams.sendMessage(lead, {
|
|
|
- target: 'inactive-target', content: content('inspection failure'), delivery: 'wakeup', signal: SIGNAL,
|
|
|
+ target: 'inactive-target', content: content('inspection failure'), signal: SIGNAL,
|
|
|
})
|
|
|
expect(uncertain.status).toBe('queued')
|
|
|
openRead.mockRestore()
|
|
|
|
|
|
- vi.spyOn(ctx.subagents as unknown as HostPromptQueue, queueSubagentPrompt)
|
|
|
+ vi.spyOn(ctx.subagents as unknown as HostPromptDeliverer, deliverSubagentPrompt)
|
|
|
.mockRejectedValueOnce(new Error('delivery unavailable'))
|
|
|
const failed = await ctx.agentTeams.sendMessage(lead, {
|
|
|
- target: 'inactive-target', content: content('delivery failure'), delivery: 'wakeup', signal: SIGNAL,
|
|
|
+ target: 'inactive-target', content: content('delivery failure'), signal: SIGNAL,
|
|
|
})
|
|
|
expect(failed.status).toBe('queued')
|
|
|
expect(warnings.some(warning => warning.includes('read unavailable'))).toBe(true)
|
|
|
@@ -1279,22 +1322,19 @@ describe('Team mailbox and waiting', () => {
|
|
|
await waitNoAgent(ctx, live.id)
|
|
|
})
|
|
|
|
|
|
- it('keeps quiet mail dormant, wakes on follow-up, preserves FIFO, and de-duplicates delivery', async () => {
|
|
|
- const { ctx, lead } = await setup(['hang', textResponse('beta first'), textResponse('beta resumed')])
|
|
|
+ it('cold-resumes an inactive sibling with sender attribution', async () => {
|
|
|
+ const { ctx, lead } = await setup(['hang', 'hang'])
|
|
|
const alphaStarted = await spawn(ctx, lead, 'alpha')
|
|
|
const alpha = await waitRunning(ctx, alphaStarted.member.id)
|
|
|
const betaStarted = await spawn(ctx, lead, 'beta')
|
|
|
- await waitNoAgent(ctx, betaStarted.member.id)
|
|
|
+ const beta = await waitRunning(ctx, betaStarted.member.id)
|
|
|
+ ctx.agentTeams.interrupt(lead, 'beta')
|
|
|
+ await waitNoAgent(ctx, beta.id)
|
|
|
|
|
|
- const quiet = await ctx.agentTeams.sendMessage(alpha, {
|
|
|
- target: 'beta', content: content('quiet info'), delivery: 'quiet', signal: SIGNAL,
|
|
|
- })
|
|
|
- expect(quiet.status).toBe('queued')
|
|
|
- expect(ctx.agents.get(betaStarted.member.id)).toBeUndefined()
|
|
|
- const waking = await ctx.agentTeams.sendMessage(alpha, {
|
|
|
- target: 'beta', content: content('do another turn'), delivery: 'wakeup', signal: SIGNAL,
|
|
|
+ const first = await ctx.agentTeams.sendMessage(alpha, {
|
|
|
+ target: 'beta', content: content('first update'), signal: SIGNAL,
|
|
|
})
|
|
|
- expect(waking.status).toBe('accepted')
|
|
|
+ expect(first.status).toBe('accepted')
|
|
|
await waitNoAgent(ctx, betaStarted.member.id)
|
|
|
await vi.waitFor(() => { expect(durable(lead).pendingMessages).toEqual([]) })
|
|
|
|
|
|
@@ -1305,18 +1345,16 @@ describe('Team mailbox and waiting', () => {
|
|
|
if (event.type !== 'user/message') return undefined
|
|
|
const block = event.data.content.at(-1)
|
|
|
return block?.type === 'text' ? block.text : undefined
|
|
|
- })).toEqual(['quiet info', 'do another turn'])
|
|
|
+ })).toEqual(['first update'])
|
|
|
expect(peerMessages.map(event => event.type === 'user/message'
|
|
|
? event.data.content[0]?.type === 'text' && event.data.content[0].text
|
|
|
: undefined)).toEqual([
|
|
|
expect.stringMatching(/^Team message .* from alpha:$/u),
|
|
|
- expect.stringMatching(/^Team message .* from alpha:$/u),
|
|
|
])
|
|
|
expect(peerMessages.map(event => event.type === 'user/message' && event.data.source.kind === 'team-message'
|
|
|
? [event.data.source.messageId, event.data.source.senderName]
|
|
|
: undefined)).toEqual([
|
|
|
- [quiet.messageId, 'alpha'],
|
|
|
- [waking.messageId, 'alpha'],
|
|
|
+ [first.messageId, 'alpha'],
|
|
|
])
|
|
|
|
|
|
ctx.agentTeams.interrupt(lead, 'alpha')
|
|
|
@@ -1331,25 +1369,26 @@ describe('Team mailbox and waiting', () => {
|
|
|
const target = await spawn(ctx, lead, 'target')
|
|
|
await waitNoAgent(ctx, target.member.id)
|
|
|
await expect(ctx.agentTeams.sendMessage(lead, {
|
|
|
- target: 'target', content: content('x'.repeat(300)), delivery: 'quiet', signal: SIGNAL,
|
|
|
+ target: 'target', content: content('x'.repeat(300)), signal: SIGNAL,
|
|
|
})).rejects.toMatchObject({ code: 'TEAM_MESSAGE_TOO_LARGE' })
|
|
|
+ vi.spyOn(ctx.sessionPersistence, 'open').mockRejectedValueOnce(new Error('temporary read failure'))
|
|
|
const queued = await ctx.agentTeams.sendMessage(lead, {
|
|
|
- target: 'target', content: content('one'), delivery: 'quiet', signal: SIGNAL,
|
|
|
+ target: 'target', content: content('one'), signal: SIGNAL,
|
|
|
})
|
|
|
expect(queued.status).toBe('queued')
|
|
|
await expect(ctx.agentTeams.sendMessage(lead, {
|
|
|
- target: 'target', content: content('two'), delivery: 'quiet', signal: SIGNAL,
|
|
|
+ target: 'target', content: content('two'), signal: SIGNAL,
|
|
|
})).rejects.toMatchObject({ code: 'TEAM_MAILBOX_FULL' })
|
|
|
await expect(ctx.agentTeams.sendMessage(lead, {
|
|
|
- target: 'lead', content: content('self'), delivery: 'quiet', signal: SIGNAL,
|
|
|
+ target: 'lead', content: content('self'), signal: SIGNAL,
|
|
|
})).rejects.toMatchObject({ code: 'TEAM_SELF_MESSAGE' })
|
|
|
await expect(ctx.agentTeams.sendMessage(lead, {
|
|
|
- target: 'missing', content: content('unknown target'), delivery: 'quiet', signal: SIGNAL,
|
|
|
+ target: 'missing', content: content('unknown target'), signal: SIGNAL,
|
|
|
})).rejects.toMatchObject({ code: 'TEAM_MEMBER_NOT_FOUND' })
|
|
|
const controller = new AbortController()
|
|
|
controller.abort(new TeamError('cancelled before queue', 'TEST_CANCELLED'))
|
|
|
await expect(ctx.agentTeams.sendMessage(lead, {
|
|
|
- target: 'target', content: content('cancelled'), delivery: 'quiet', signal: controller.signal,
|
|
|
+ target: 'target', content: content('cancelled'), signal: controller.signal,
|
|
|
})).rejects.toMatchObject({ code: 'TEST_CANCELLED' })
|
|
|
})
|
|
|
|
|
|
@@ -1358,12 +1397,12 @@ describe('Team mailbox and waiting', () => {
|
|
|
const started = await spawn(ctx, lead, 'worker')
|
|
|
const worker = await waitRunning(ctx, started.member.id)
|
|
|
const followup = await ctx.agentTeams.sendMessage(lead, {
|
|
|
- target: 'worker', content: content('retained follow-up'), delivery: 'wakeup', signal: SIGNAL,
|
|
|
+ target: 'worker', content: content('retained follow-up'), signal: SIGNAL,
|
|
|
})
|
|
|
expect(followup.status).toBe('accepted')
|
|
|
expect(ctx.agentTeams.interrupt(lead, 'worker')).toEqual({ previousStatus: 'running' })
|
|
|
await vi.waitFor(() => { expect(worker.status).toBe('idle') })
|
|
|
- expect(worker.inbox.nextTurn.some(message => message.source.kind === 'team-message'
|
|
|
+ expect(worker.inbox.nextStep.some(message => message.source.kind === 'team-message'
|
|
|
&& message.source.messageId === followup.messageId)).toBe(true)
|
|
|
worker.cancel({ kind: 'parent' })
|
|
|
await waitNoAgent(ctx, worker.id)
|
|
|
@@ -1538,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,
|
|
|
})
|
|
|
@@ -1551,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,
|
|
|
@@ -1574,7 +1613,7 @@ describe('Team mailbox and waiting', () => {
|
|
|
const entered = Promise.withResolvers<undefined>()
|
|
|
const aborted = Promise.withResolvers<undefined>()
|
|
|
const release = Promise.withResolvers<undefined>()
|
|
|
- vi.spyOn(ctx.subagents as unknown as HostPromptQueue, queueSubagentPrompt)
|
|
|
+ vi.spyOn(ctx.subagents as unknown as HostPromptDeliverer, deliverSubagentPrompt)
|
|
|
.mockImplementation(async (_parent, _childId, _content, _source, signal) => {
|
|
|
entered.resolve(undefined)
|
|
|
return await new Promise<never>((_resolve, reject) => {
|
|
|
@@ -1591,7 +1630,6 @@ describe('Team mailbox and waiting', () => {
|
|
|
const sending = ctx.agentTeams.sendMessage(lead, {
|
|
|
target: 'mailbox-worker',
|
|
|
content: content('resume during disposal'),
|
|
|
- delivery: 'wakeup',
|
|
|
signal: SIGNAL,
|
|
|
})
|
|
|
await entered.promise
|
|
|
@@ -1616,11 +1654,10 @@ describe('Team mailbox and waiting', () => {
|
|
|
senderId: SessionId('sender'),
|
|
|
senderName: 'sender',
|
|
|
targetId: lead.id,
|
|
|
- delivery: 'wakeup',
|
|
|
content: content('acknowledge before disposal'),
|
|
|
}
|
|
|
lead.session.append('team/message/queued', {
|
|
|
- version: 1,
|
|
|
+ version: 2,
|
|
|
teamId: TeamId(lead.id),
|
|
|
message,
|
|
|
})
|
|
|
@@ -1695,14 +1732,13 @@ describe('Team mailbox and waiting', () => {
|
|
|
signal: SIGNAL,
|
|
|
})).rejects.toMatchObject({ code: 'TEAM_DISPOSED' })
|
|
|
await expect(ctx.agentTeams.sendMessage(lead, {
|
|
|
- target: 'nobody', content: content('must reject'), delivery: 'quiet', signal: SIGNAL,
|
|
|
+ target: 'nobody', content: content('must reject'), signal: SIGNAL,
|
|
|
})).rejects.toMatchObject({ code: 'TEAM_DISPOSED' })
|
|
|
await expect(internal.mailbox.tryDispatch(lead, {
|
|
|
id: TeamMessageId('post-disposal-message'),
|
|
|
senderId: lead.id,
|
|
|
senderName: 'lead',
|
|
|
targetId: lead.id,
|
|
|
- delivery: 'quiet',
|
|
|
content: content('must not dispatch'),
|
|
|
}, SIGNAL)).resolves.toBe(false)
|
|
|
})
|
|
|
@@ -1774,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)
|
|
|
@@ -1794,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>()
|
|
|
@@ -1807,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' },
|
|
|
})
|