fix(agent-teams): enforce task dependencies in approved teams

An approved team handed every member its instructions at the same moment,
so members whose tasks depended on unfinished work started anyway: the
second and third layers of a plan ran before the first, a member could mark
a blocked task in progress, and results only ever went to the lead. The
official CLI states dependencies in tool prompts and nothing more, so the
plan the user approved was not what ran.

- A member whose tasks all wait on other tasks gets its instructions only
  once one of them is ready. What the lead or a teammate sends it before
  that waits in its inbox and arrives with the instructions; only the user
  writing to it, or a shutdown request, reaches it earlier. This survives a
  Stop and a server restart.
- TaskUpdate refuses a teammate that starts or completes a task while a
  task it is blocked by is unfinished. The lead is not held to it.
- Each member is told who waits on its tasks and to send them its result
  before completing, so a dependent member starts with that result in
  hand. A member that starts without one is told whose is missing.
- A member that exits after approving a shutdown request is no longer
  recorded as failed, and the lead gets no failure notice for it. The
  shutdown-approval schema keeps leaving out the desktop backend on
  purpose, now documented and tested.
- TeamPlan tells the lead that a dependent member's prompt is delivered
  when its dependencies are done.
This commit is contained in:
程序员阿江(Relakkes)
2026-10-06 01:43:38 +08:00
parent 8b2f9102a5
commit ea8d7be6ca
11 changed files with 618 additions and 20 deletions
@@ -25,6 +25,10 @@ const checks: Check[] = [
title: 'Agent Teams legacy mailbox migration to unread inbox plus history', title: 'Agent Teams legacy mailbox migration to unread inbox plus history',
command: ['bun', 'test', './src/utils/teammateMailbox.test.ts', '--test-name-pattern', 'legacy inbox files'], command: ['bun', 'test', './src/utils/teammateMailbox.test.ts', '--test-name-pattern', 'legacy inbox files'],
}, },
{
title: 'Agent Teams members recorded before deferred instructions rehydrate as already instructed',
command: ['bun', 'test', './src/server/services/teamPlanRuntime.supervisor.test.ts', '--test-name-pattern', 'written before deferred instructions'],
},
{ {
title: 'Agent Teams teammate resume from agent metadata written before the teammate fields', title: 'Agent Teams teammate resume from agent metadata written before the teammate fields',
command: ['bun', 'test', './src/utils/swarm/inProcessRunner.resume.test.ts', '--test-name-pattern', 'metadata written before'], command: ['bun', 'test', './src/utils/swarm/inProcessRunner.resume.test.ts', '--test-name-pattern', 'metadata written before'],
@@ -19,6 +19,14 @@ sdk.onmessage = async event => {
sdk.send(JSON.stringify({ type: 'control_response', response: { subtype: 'success', request_id: message.request_id, response: {} } })) sdk.send(JSON.stringify({ type: 'control_response', response: { subtype: 'success', request_id: message.request_id, response: {} } }))
} else if (message.type === 'user') { } else if (message.type === 'user') {
const content = JSON.stringify(message.message?.content) const content = JSON.stringify(message.message?.content)
// A member that agrees to shut down tells the lead so, then leaves. The
// approval is what the real CLI writes for a desktop member, backend included.
if (content.includes('FIXTURE_APPROVE_SHUTDOWN')) {
const { writeToMailbox } = await import('../../../utils/teammateMailbox.js')
const from = arg('--agent-name')!
await writeToMailbox('team-lead', { from, timestamp: new Date().toISOString(), text: JSON.stringify({ type: 'shutdown_approved', requestId: 'fixture', from, timestamp: new Date().toISOString(), paneId: '', backendType: 'process' }) }, arg('--team-name'))
sdk.close(); setTimeout(() => process.exit(0), 5); continue
}
if (content.includes('FIXTURE_SHUTDOWN')) { sdk.close(); setTimeout(() => process.exit(0), 5); continue } if (content.includes('FIXTURE_SHUTDOWN')) { sdk.close(); setTimeout(() => process.exit(0), 5); continue }
const url = new URL(process.env.ANTHROPIC_BASE_URL!) const url = new URL(process.env.ANTHROPIC_BASE_URL!)
if (url.hostname !== '127.0.0.1') throw new Error('Fixture refuses non-loopback upstream') if (url.hostname !== '127.0.0.1') throw new Error('Fixture refuses non-loopback upstream')
@@ -53,7 +53,11 @@ type Harness = Awaited<ReturnType<typeof startHarness>>
* A real ConversationService with fixture worker CLIs over a loopback SDK * A real ConversationService with fixture worker CLIs over a loopback SDK
* bridge. `respond` decides how each upstream "model" call ends. * bridge. `respond` decides how each upstream "model" call ends.
*/ */
async function startHarness(memberNames: string[]) { async function startHarness(memberNames: string[], options: {
tasks?: Array<{ id: string; subject: string; ownerId: string; dependencies: string[] }>
/** Members expected to get their instructions at launch (default: all). */
instructedAtLaunch?: string[]
} = {}) {
const { conversationService: service } = await import('./conversationService.js') const { conversationService: service } = await import('./conversationService.js')
const { ProviderService } = await import('./providerService.js') const { ProviderService } = await import('./providerService.js')
const { launchTeamPlanRuntime } = await import('./teamPlanRuntime.js') const { launchTeamPlanRuntime } = await import('./teamPlanRuntime.js')
@@ -89,7 +93,7 @@ async function startHarness(memberNames: string[]) {
// The lead is a fixture process too; it only needs to answer control requests. // The lead is a fixture process too; it only needs to answer control requests.
await service.startSession(parentId, home, `ws://127.0.0.1:${bridge.port}/sdk/${parentId}?token=parent`, { providerId: provider.id, model: 'cheap', teamWorker: { parentSessionId: 'external', teamName: 'lead-fixture', memberId: 'lead', name: 'leader', systemPrompt: 'lead fixture' } }) await service.startSession(parentId, home, `ws://127.0.0.1:${bridge.port}/sdk/${parentId}?token=parent`, { providerId: provider.id, model: 'cheap', teamWorker: { parentSessionId: 'external', teamName: 'lead-fixture', memberId: 'lead', name: 'leader', systemPrompt: 'lead fixture' } })
const members = memberNames.map(name => ({ id: name, name, agentType: 'research', prompt: `Work as ${name}`, runtime: { providerId: provider.id, modelId: 'cheap' }, agentSnapshot: { systemPrompt: 'Read only', tools: ['Read'] } })) const members = memberNames.map(name => ({ id: name, name, agentType: 'research', prompt: `Work as ${name}`, runtime: { providerId: provider.id, modelId: 'cheap' }, agentSnapshot: { systemPrompt: 'Read only', tools: ['Read'] } }))
const tasks = memberNames.map(name => ({ id: `task-${name}`, subject: `Task for ${name}`, ownerId: name, dependencies: [] })) const tasks = options.tasks ?? memberNames.map(name => ({ id: `task-${name}`, subject: `Task for ${name}`, ownerId: name, dependencies: [] }))
const planId = crypto.randomUUID() const planId = crypto.randomUUID()
const plan = { schemaVersion: 1, planId, sessionId: parentId, teamName, incarnationId: createHash('sha256').update(JSON.stringify([teamName, parentId, createdAt])).digest('hex'), revision: 2, state: 'launching', workDir: home, members, tasks, leaderRuntime: members[0]!.runtime, createdAt, updatedAt: createdAt, approvedSnapshot: { revision: 1, members, tasks, leaderRuntime: members[0]!.runtime, approvedAt: Date.now(), requestId: 'approve' } } as any const plan = { schemaVersion: 1, planId, sessionId: parentId, teamName, incarnationId: createHash('sha256').update(JSON.stringify([teamName, parentId, createdAt])).digest('hex'), revision: 2, state: 'launching', workDir: home, members, tasks, leaderRuntime: members[0]!.runtime, createdAt, updatedAt: createdAt, approvedSnapshot: { revision: 1, members, tasks, leaderRuntime: members[0]!.runtime, approvedAt: Date.now(), requestId: 'approve' } } as any
await writeFile(join(getTeamDir(teamName), 'plan.json'), JSON.stringify(plan)) await writeFile(join(getTeamDir(teamName), 'plan.json'), JSON.stringify(plan))
@@ -101,7 +105,8 @@ async function startHarness(memberNames: string[]) {
while (!(await condition()) && Date.now() < end) await new Promise(resolve => setTimeout(resolve, 20)) while (!(await condition()) && Date.now() < end) await new Promise(resolve => setTimeout(resolve, 20))
if (!(await condition())) throw new Error(`Timed out waiting for ${label}`) if (!(await condition())) throw new Error(`Timed out waiting for ${label}`)
} }
await waitFor(() => requests.length >= memberNames.length, 'initial member turns') const instructed = options.instructedAtLaunch ?? memberNames
await waitFor(() => instructed.every(name => requests.some(request => request.prompt.includes(`Work as ${name}`))), 'initial member turns')
return { return {
service, teamName, parentId, planId, memberIds, requests, replies, waitFor, service, teamName, parentId, planId, memberIds, requests, replies, waitFor,
/** Keeps model replies (and so member turns) pending until released. */ /** Keeps model replies (and so member turns) pending until released. */
@@ -569,3 +574,291 @@ test('a newly unblocked task wakes its idle owner once', async () => {
await h.cleanup() await h.cleanup()
} }
}, 30_000) }, 30_000)
/** recon maps the release; analyst reviews it once the map exists. */
async function startDependentTeam() {
const h = await startHarness(['recon', 'analyst'], {
tasks: [
{ id: 'map', subject: 'Map the release', ownerId: 'recon', dependencies: [] },
{ id: 'review', subject: 'Review the release', ownerId: 'analyst', dependencies: ['map'] },
],
instructedAtLaunch: ['recon'],
})
const { listTasks, getCanonicalTeamTaskListId } = await import('../../utils/tasks.js')
const taskListId = getCanonicalTeamTaskListId(h.teamName)
const tasks = await listTasks(taskListId)
return {
...h,
taskListId,
map: tasks.find(task => task.subject === 'Map the release')!.id,
review: tasks.find(task => task.subject === 'Review the release')!.id,
analystBriefs: () => h.requests.filter(request => request.prompt.includes('Work as analyst')),
}
}
test('a member whose tasks all wait on other tasks gets its instructions only once one is ready', async () => {
const { setTeamRuntimeTimingForTests } = await import('./teamPlanRuntime.js')
// A redundant "task is ready" wake after the instructions would come at once.
setTeamRuntimeTimingForTests({ unblockedTaskWakeDelayMs: 0 })
const h = await startDependentTeam()
try {
const { updateTask } = await import('../../utils/tasks.js')
expect(await h.teamMember('analyst')).toMatchObject({ isActive: false, awaitingDependencies: true })
await new Promise(resolve => setTimeout(resolve, 600))
expect(h.analystBriefs()).toHaveLength(0)
await updateTask(h.taskListId, h.map, { status: 'completed' })
await h.waitFor(() => h.analystBriefs().length > 0, 'instructions once the review is ready')
expect(h.analystBriefs()[0]!.prompt).toContain(`#${h.review} Review the release`)
// recon completed without sending anything: the analyst is told, so it
// asks instead of working without the result it depends on.
expect(h.analystBriefs()[0]!.prompt).toContain('Nothing has reached you from recon yet')
await h.waitFor(async () => (await h.teamMember('analyst'))?.awaitingDependencies === undefined, 'instructions recorded as delivered')
await new Promise(resolve => setTimeout(resolve, 600))
expect(h.requests.filter(request => request.prompt.includes(`#${h.review}`))).toHaveLength(1)
} finally {
await h.cleanup()
}
}, 30_000)
test("only the user's own message starts a waiting member early; the lead's waits and arrives with its instructions", async () => {
const h = await startDependentTeam()
try {
// A lead relays results as they come in; that must not start the member.
await messageMember(h, 'analyst', 'team-lead', 'Forwarded: the inventory so far')
await new Promise(resolve => setTimeout(resolve, 700))
expect(h.analystBriefs()).toHaveLength(0)
expect(await h.teamMember('analyst')).toMatchObject({ isActive: false, awaitingDependencies: true })
await messageMember(h, 'analyst', 'user', 'Start now, on my word')
await h.waitFor(() => h.analystBriefs().length > 0, 'instructions with the user message')
expect(h.analystBriefs()[0]!.prompt).toContain('Start now, on my word')
expect(h.analystBriefs()[0]!.prompt).toContain('Forwarded: the inventory so far')
await h.waitFor(async () => (await h.teamMember('analyst'))?.awaitingDependencies === undefined, 'instructions recorded as delivered')
await messageMember(h, 'analyst', 'team-lead', 'One more thing')
await h.waitFor(() => h.requests.some(request => request.prompt.includes('One more thing')), 'second message')
expect(h.analystBriefs()).toHaveLength(1)
} finally {
await h.cleanup()
}
}, 30_000)
test('a shutdown request reaches a member that has not started yet, without its instructions', async () => {
const h = await startDependentTeam()
try {
const { updateTask } = await import('../../utils/tasks.js')
const { createShutdownRequestMessage } = await import('../../utils/teammateMailbox.js')
await messageMember(h, 'analyst', 'team-lead', JSON.stringify(createShutdownRequestMessage({ requestId: 'shutdown-1', from: 'team-lead' })))
await h.waitFor(() => h.requests.some(request => request.prompt.includes('shutdown-1')), 'shutdown request delivered')
expect(h.analystBriefs()).toHaveLength(0)
await h.waitFor(async () => (await h.teamMember('analyst'))?.isActive === false, 'member idle again')
expect((await h.teamMember('analyst'))?.awaitingDependencies).toBe(true)
// Declined or ignored, it still starts properly once its task is ready.
await updateTask(h.taskListId, h.map, { status: 'completed' })
await h.waitFor(() => h.analystBriefs().length > 0, 'instructions once the review is ready')
} finally {
await h.cleanup()
}
}, 30_000)
test('after a Stop, a member that has not started yet starts once the team runs again and its task is ready', async () => {
const h = await startDependentTeam()
try {
const { updateTask } = await import('../../utils/tasks.js')
await h.waitFor(async () => (await h.leadNotifications()).length === 1, 'recon report')
h.service.sendInterrupt(h.parentId)
await h.service.waitForTeamWorkersStopped(h.parentId)
const notices = await userMessagesLead(h)
expect(notices).toHaveLength(1)
expect(notices[0]).toMatch(/- recon: #\d+ Map the release \(pending\)/)
expect(notices[0]).toContain(`- analyst: not started yet; it starts by itself once #${h.map} is completed, and messages to it wait until then`)
await updateTask(h.taskListId, h.map, { status: 'completed' })
await new Promise(resolve => setTimeout(resolve, 700))
// A finished dependency does not undo the user's Stop.
expect(h.analystBriefs()).toHaveLength(0)
// The lead's instruction resumes the team even when it goes to the member
// that has not started; that member then starts because its task is ready.
await messageMember(h, 'analyst', 'team-lead', 'The user wants you to continue')
await h.waitFor(() => h.analystBriefs().length > 0, 'analyst starts once the team runs again')
expect(h.analystBriefs()[0]!.prompt).toContain('The user wants you to continue')
// It never had a turn, so its own session starts afresh rather than resuming.
expect(h.service.hasSession(h.memberIds.analyst!)).toBe(true)
} finally {
await h.cleanup()
}
}, 30_000)
test('a member still waiting after a server restart gets its approved instructions once a task is ready', async () => {
const h = await startDependentTeam()
try {
const runtime = await import('./teamPlanRuntime.js')
const { updateTask } = await import('../../utils/tasks.js')
await h.service.stopSessionAndWait(h.memberIds.recon!)
await h.service.stopSessionAndWait(h.memberIds.analyst!)
runtime.forgetTeamPlanRuntimesForTests()
expect(await runtime.rehydrateTeamPlanRuntimesForSession(h.parentId)).toBe(true)
await new Promise(resolve => setTimeout(resolve, 600))
expect(h.analystBriefs()).toHaveLength(0)
await updateTask(h.taskListId, h.map, { status: 'completed' })
await h.waitFor(() => h.analystBriefs().length > 0, 'instructions after the restart')
expect(h.service.hasSession(h.memberIds.analyst!)).toBe(true)
// The rebuilt instructions carry the live task ids.
expect(h.analystBriefs()[0]!.prompt).toContain(`\\"id\\":\\"${h.review}\\"`)
} finally {
await h.cleanup()
}
}, 30_000)
test('members written before deferred instructions rehydrate as already instructed', async () => {
const h = await startHarness(['reader'])
try {
const runtime = await import('./teamPlanRuntime.js')
// An older team file never carries awaitingDependencies.
expect(await h.teamMember('reader')).not.toHaveProperty('awaitingDependencies')
await h.service.stopSessionAndWait(h.memberIds.reader!)
runtime.forgetTeamPlanRuntimesForTests()
expect(await runtime.rehydrateTeamPlanRuntimesForSession(h.parentId)).toBe(true)
const before = h.requests.length
await messageMember(h, 'reader', 'team-lead', 'Resume after restart')
await h.waitFor(() => h.requests.length > before, 'restarted member turn')
expect(h.requests.at(-1)?.prompt).toContain('Resume after restart')
expect(h.requests.at(-1)?.prompt).not.toContain('Work as reader')
} finally {
await h.cleanup()
}
}, 30_000)
test("a teammate's note does not start a waiting member and arrives with its instructions", async () => {
const h = await startDependentTeam()
try {
const { updateTask } = await import('../../utils/tasks.js')
const { readUnreadMessages } = await import('../../utils/teammateMailbox.js')
// recon hands its result over before it closes the task.
await messageMember(h, 'analyst', 'recon', 'Release map: three feature groups')
await new Promise(resolve => setTimeout(resolve, 700))
expect(h.analystBriefs()).toHaveLength(0)
expect(await readUnreadMessages('analyst', h.teamName)).toHaveLength(1)
expect(await h.teamMember('analyst')).toMatchObject({ isActive: false, awaitingDependencies: true })
await updateTask(h.taskListId, h.map, { status: 'completed' })
await h.waitFor(() => h.analystBriefs().length > 0, 'instructions once the review is ready')
expect(h.analystBriefs()[0]!.prompt).toContain('Release map: three feature groups')
expect(h.analystBriefs()[0]!.prompt).toContain(`#${h.review} Review the release`)
// The result it depends on is in hand, so nothing is reported missing.
expect(h.analystBriefs()[0]!.prompt).not.toContain('Nothing has reached you')
await h.waitFor(async () => (await readUnreadMessages('analyst', h.teamName)).length === 0, 'note marked as read')
await new Promise(resolve => setTimeout(resolve, 600))
expect(h.requests.filter(request => request.prompt.includes('Release map: three feature groups'))).toHaveLength(1)
} finally {
await h.cleanup()
}
}, 30_000)
test('the release-review plan that started its second layer first now runs layer by layer', async () => {
const { setTeamRuntimeTimingForTests } = await import('./teamPlanRuntime.js')
setTeamRuntimeTimingForTests({ unblockedTaskWakeDelayMs: 0 })
// The approved plan of the reported run: one inventory, three analyses that
// depend on it, a verification of all three, and a report by one analyst.
const names = ['recon', 'perf-analyst', 'sec-analyst', 'ux-analyst', 'verifier']
const h = await startHarness(names, {
tasks: [
{ id: 't1', subject: 'Inventory', ownerId: 'recon', dependencies: [] },
{ id: 't2', subject: 'Performance', ownerId: 'perf-analyst', dependencies: ['t1'] },
{ id: 't3', subject: 'Security', ownerId: 'sec-analyst', dependencies: ['t1'] },
{ id: 't4', subject: 'Interaction', ownerId: 'ux-analyst', dependencies: ['t1'] },
{ id: 't5', subject: 'Verification', ownerId: 'verifier', dependencies: ['t2', 't3', 't4'] },
{ id: 't6', subject: 'Report', ownerId: 'perf-analyst', dependencies: ['t2', 't3', 't4', 't5'] },
],
instructedAtLaunch: ['recon'],
})
try {
const { listTasks, updateTask, getCanonicalTeamTaskListId } = await import('../../utils/tasks.js')
const taskListId = getCanonicalTeamTaskListId(h.teamName)
const id = Object.fromEntries((await listTasks(taskListId)).map(task => [task.subject, task.id]))
const complete = (subject: string) => updateTask(taskListId, id[subject]!, { status: 'completed' })
const briefs = (name: string) => h.requests.filter(request => request.prompt.includes(`Work as ${name}`))
const instructed = () => names.filter(name => briefs(name).length > 0)
const settle = () => new Promise(resolve => setTimeout(resolve, 600))
await settle()
expect(instructed()).toEqual(['recon'])
// Each member is told who waits on its task, and to hand its result over
// before completing it: the inventory goes to the three analysts.
expect(briefs('recon')[0]!.prompt).toContain(`\\"handOffTo\\":[\\"perf-analyst\\",\\"sec-analyst\\",\\"ux-analyst\\"]`)
expect(briefs('recon')[0]!.prompt).toContain('first SendMessage its result to each member listed there, then mark it completed')
await h.waitFor(() => h.requests.some(request => request.prompt.includes(`Approved team ${h.teamName} is running`)), 'launch notice to the lead')
expect(h.requests.find(request => request.prompt.includes(`Approved team ${h.teamName} is running`))!.prompt)
.toContain('Every task of perf-analyst, sec-analyst, ux-analyst, verifier waits on other tasks')
await complete('Inventory')
await h.waitFor(() => instructed().length === 4, 'the three analysts to start')
await settle()
expect(instructed()).toEqual(['recon', 'perf-analyst', 'sec-analyst', 'ux-analyst'])
await complete('Performance')
await complete('Security')
await settle()
expect(briefs('verifier')).toHaveLength(0)
await complete('Interaction')
await h.waitFor(() => briefs('verifier').length > 0, 'the verifier to start')
expect(briefs('verifier')[0]!.prompt).toContain(`#${id.Verification} Verification`)
// An analyst hands its analysis to the verifier; its own report task is not a hand-off.
expect(briefs('perf-analyst')[0]!.prompt).toContain(`\\"handOffTo\\":[\\"verifier\\"]`)
// The verification goes to the analyst who writes the report.
expect(briefs('verifier')[0]!.prompt).toContain(`\\"handOffTo\\":[\\"perf-analyst\\"]`)
// The analyst that also owns the report is woken for it, not instructed twice.
await h.waitFor(async () => (await h.teamMember('perf-analyst'))?.isActive === false, 'perf-analyst idle')
await complete('Verification')
await h.waitFor(() => h.requests.some(request => request.prompt.includes(`Task #${id.Report}`)), 'wake for the report')
await settle()
for (const name of names) expect(briefs(name)).toHaveLength(1)
} finally {
await h.cleanup()
}
}, 60_000)
test('a member that leaves after approving a shutdown request is not reported as a failure', async () => {
const h = await startHarness(['reader', 'writer'])
try {
const { createShutdownRequestMessage } = await import('../../utils/teammateMailbox.js')
await h.waitFor(async () => (await h.leadNotifications()).length === 2, 'initial turn reports')
const request = (marker: string) => JSON.stringify(createShutdownRequestMessage({ requestId: `stop-${marker}`, from: 'team-lead', reason: marker }))
// reader agrees and leaves; writer's process disappears without answering.
await messageMember(h, 'reader', 'team-lead', request('FIXTURE_APPROVE_SHUTDOWN'))
await messageMember(h, 'writer', 'team-lead', request('FIXTURE_SHUTDOWN'))
await h.waitFor(async () => (await h.teamMember('reader'))?.terminated === true && (await h.teamMember('writer'))?.terminated === true, 'both processes gone')
await h.waitFor(async () => (await h.leadNotifications()).some(item => item.from === 'writer' && item.idleReason === 'failed'), 'failure notice for the member that vanished')
// An exit that is not recognised as approved is recorded a moment later,
// so give that time to happen before asserting it did not.
await new Promise(resolve => setTimeout(resolve, 900))
expect((await h.teamMember('reader'))?.lastError).toBeUndefined()
expect((await h.teamMember('writer'))?.lastError).toContain('exited')
const notices = await h.leadNotifications()
expect(notices.some(item => item.from === 'reader' && item.type === 'shutdown_approved')).toBe(true)
expect(notices.filter(item => item.from === 'reader' && item.idleReason === 'failed')).toEqual([])
// With the user's next message the lead hears about the failure only.
const told = await userMessagesLead(h)
expect(told).toHaveLength(1)
expect(told[0]).toContain('- writer (stopped on an error')
expect(told[0]).not.toContain('- reader')
// The approval covers that one exit: a later crash is a failure again.
const sessionId = h.memberIds.reader!
await messageMember(h, 'reader', 'team-lead', 'One more check, please')
await h.waitFor(() => h.service.hasSession(sessionId), 'reader restarted by the message')
await h.waitFor(async () => (await h.teamMember('reader'))?.isActive === false, 'reader idle after the message')
h.replies.push('FIXTURE_CRASH')
await messageMember(h, 'reader', 'team-lead', 'And another')
await h.waitFor(async () => (await h.leadNotifications()).some(item => item.from === 'reader' && item.idleReason === 'failed'), 'failure notice for the later crash')
expect((await h.teamMember('reader'))?.lastError).toBeTruthy()
} finally {
await h.cleanup()
}
}, 30_000)
+182 -16
View File
@@ -2,7 +2,7 @@ import { createHash, randomUUID } from 'node:crypto'
import { migrationMaintenance } from '../migrationMaintenance.js' import { migrationMaintenance } from '../migrationMaintenance.js'
import { readdir, readFile, stat } from 'node:fs/promises' import { readdir, readFile, stat } from 'node:fs/promises'
import { join } from 'node:path' import { join } from 'node:path'
import { isValidTeamMemberName, teamPlanRecordSchema, type TeamPlanRecord, type TeamPlanMember } from '../../shared/teamPlan.js' import { isValidTeamMemberName, teamPlanRecordSchema, type TeamPlanRecord, type TeamPlanMember, type TeamPlanTask } from '../../shared/teamPlan.js'
import { conversationService } from './conversationService.js' import { conversationService } from './conversationService.js'
import { ProviderService } from './providerService.js' import { ProviderService } from './providerService.js'
import { CLAUDE_OFFICIAL_PROVIDER_ID } from '../types/provider.js' import { CLAUDE_OFFICIAL_PROVIDER_ID } from '../types/provider.js'
@@ -15,7 +15,7 @@ import { cleanupTeamDirectories, getTeamDir, mutateTeamFileAsync, readTeamFile,
import { getTeamsDir as getTeamsDirectory } from '../../utils/envUtils.js' import { getTeamsDir as getTeamsDirectory } from '../../utils/envUtils.js'
import { TEAM_LEAD_NAME } from '../../utils/swarm/constants.js' import { TEAM_LEAD_NAME } from '../../utils/swarm/constants.js'
import { createTask, listTasks, updateTask, withTaskListLifecycleLock, getCanonicalTeamTaskListId, type Task } from '../../utils/tasks.js' import { createTask, listTasks, updateTask, withTaskListLifecycleLock, getCanonicalTeamTaskListId, type Task } from '../../utils/tasks.js'
import { readUnreadMessages, markMessagesAsReadByPredicate, writeToMailbox, createIdleNotification, formatTeammateMessages, type TeammateMessage } from '../../utils/teammateMailbox.js' import { readUnreadMessages, readMailboxHistory, markMessagesAsReadByPredicate, writeToMailbox, createIdleNotification, formatTeammateMessages, isShutdownRequest, type TeammateMessage } from '../../utils/teammateMailbox.js'
import { formatTeammateAutoContinuePrompt, isTransientTurnFailure, summarizeTurnFailure, TEAMMATE_AUTO_CONTINUE_DELAYS_MS } from '../../utils/swarm/turnFailure.js' import { formatTeammateAutoContinuePrompt, isTransientTurnFailure, summarizeTurnFailure, TEAMMATE_AUTO_CONTINUE_DELAYS_MS } from '../../utils/swarm/turnFailure.js'
const providerService = new ProviderService() const providerService = new ProviderService()
@@ -58,6 +58,14 @@ type WorkerRuntime = {
wakeNoticePending?: boolean wakeNoticePending?: boolean
wokenForTaskIds: Set<string> wokenForTaskIds: Set<string>
lastResultAt?: number lastResultAt?: number
/**
* Approved instructions held back because every task the member owns waits
* on unfinished work. They go out when one of its tasks becomes ready, or
* when the user writes to the member, whichever comes first.
*/
brief?: string
/** When a shutdown request last reached the member; an approval after it makes its exit an ordinary one. */
shutdownRequestedAt?: number
} }
type TeamLaunch = { type TeamLaunch = {
@@ -177,6 +185,39 @@ async function materializeTasks(plan: TeamPlanRecord): Promise<Record<string, st
return mapping return mapping
} }
/**
* What an approved member is told to do: its preset's opening, its prompt and
* its shared tasks. Each task names who waits on it (`handOffTo`): a member
* cannot see that from its own tasks, and without it a result only ever went
* to the lead, reaching the member that needed it late or not at all.
*/
function memberBrief(member: TeamPlanMember, members: TeamPlanMember[], tasks: TeamPlanTask[], taskMapping: Record<string, string>): string {
const ownerName = (task: TeamPlanTask) => members.find(candidate => candidate.id === task.ownerId)?.name ?? TEAM_LEAD_NAME
const assigned = tasks.filter(task => task.ownerId === member.id).map(task => {
const handOffTo = [...new Set(tasks.filter(other => other.dependencies.includes(task.id)).map(ownerName))].filter(name => name !== member.name)
return { ...task, id: taskMapping[task.id], dependencies: task.dependencies.map(dep => taskMapping[dep]), ...(handOffTo.length > 0 ? { handOffTo } : {}) }
})
const handOff = assigned.some(task => 'handOffTo' in task)
? '\n\nWhen you finish a task that has handOffTo: first SendMessage its result to each member listed there, then mark it completed. Their work starts from your message the moment you complete the task.'
: ''
return `${member.agentSnapshot?.initialPrompt ? member.agentSnapshot.initialPrompt + "\n" : ""}${member.prompt}\n\nApproved shared tasks (use TaskGet/TaskUpdate; start a task only once all of its dependencies are completed):\n${JSON.stringify(assigned)}${handOff}`
}
function unresolvedTaskIds(tasks: Task[]): Set<string> {
return new Set(tasks.filter(task => task.status !== 'completed').map(task => task.id))
}
/** Unfinished tasks the member owns whose dependencies are all complete. */
function readyTasksOf(memberName: string, tasks: Task[]): Task[] {
const unresolved = unresolvedTaskIds(tasks)
return tasks.filter(task => task.owner === memberName && task.status !== 'completed' && task.blockedBy.every(id => !unresolved.has(id)))
}
/** The member owns unfinished work and all of it waits on other tasks. */
function waitsForDependencies(memberName: string, tasks: Task[]): boolean {
return tasks.some(task => task.owner === memberName && task.status !== 'completed') && readyTasksOf(memberName, tasks).length === 0
}
// ── Worker failure classification ─────────────────────────────────────────── // ── Worker failure classification ───────────────────────────────────────────
/** A worker's CLI `result` text after a failed turn (see utils/swarm/turnFailure). */ /** A worker's CLI `result` text after a failed turn (see utils/swarm/turnFailure). */
@@ -368,6 +409,16 @@ async function handleWorkerResult(launch: TeamLaunch, worker: WorkerRuntime, mes
await notifyLead(launch, worker, { idleReason: 'available', ...(text ? { result: text } : {}) }) await notifyLead(launch, worker, { idleReason: 'available', ...(text ? { result: text } : {}) })
return return
} }
if (!alive && await approvedShutdown(launch, worker)) {
// It left because the lead asked it to and it agreed. Every process exit
// arrives here as an error result; recording this one as a failure showed
// a finished team as broken and told the lead its members had crashed.
worker.autoContinueAttempts = 0
clearFailure(worker)
cancelAutoContinue(worker)
await updateWorkerEntry(launch, worker, entry => withoutFailure({ ...entry, isActive: false, terminated: true }))
return
}
const reason = firstLine(text || 'The member process exited') const reason = firstLine(text || 'The member process exited')
if (alive && isTransientWorkerFailure(text) && worker.autoContinueAttempts < timing.autoContinueDelaysMs.length) { if (alive && isTransientWorkerFailure(text) && worker.autoContinueAttempts < timing.autoContinueDelaysMs.length) {
const attempt = ++worker.autoContinueAttempts const attempt = ++worker.autoContinueAttempts
@@ -390,6 +441,33 @@ async function handleWorkerResult(launch: TeamLaunch, worker: WorkerRuntime, mes
: `${worker.member.name}'s process exited (${reason}). Messaging it restarts it from its saved conversation.`) : `${worker.member.name}'s process exited (${reason}). Messaging it restarts it from its saved conversation.`)
} }
/**
* The approval a member writes to the lead before it exits. The mailbox's own
* isShutdownApproved does not recognise a desktop member's approval, on
* purpose (see ShutdownApprovedMessageSchema), so the type is read here.
*/
function isShutdownApproval(text: string): boolean {
try {
return (JSON.parse(text) as { type?: unknown } | null)?.type === 'shutdown_approved'
} catch {
return false
}
}
/** Whether the member approved the shutdown request it was last given; asked once per request. */
async function approvedShutdown(launch: TeamLaunch, worker: WorkerRuntime): Promise<boolean> {
const since = worker.shutdownRequestedAt
if (since === undefined) return false
// A later exit without a new request is a crash again.
worker.shutdownRequestedAt = undefined
const approved = async () => (await readMailboxHistory(TEAM_LEAD_NAME, launch.plan.teamName).catch(() => []))
.some(message => message.from === worker.member.name && Date.parse(message.timestamp) >= since && isShutdownApproval(message.text))
if (await approved()) return true
// The lead may be archiving its inbox at this very moment.
await new Promise(resolve => setTimeout(resolve, 300))
return approved()
}
/** /**
* A failure nothing retries any more: the lead hears it now through its * A failure nothing retries any more: the lead hears it now through its
* mailbox, and again with the user's next message, which is when the cause * mailbox, and again with the user's next message, which is when the cause
@@ -416,12 +494,21 @@ async function deliverToWorker(launch: TeamLaunch, worker: WorkerRuntime, messag
// The message replaces any scheduled retry and answers the last failure; // The message replaces any scheduled retry and answers the last failure;
// the turn it starts records its own outcome. // the turn it starts records its own outcome.
await updateWorkerEntry(launch, worker, entry => withoutFailure({ ...entry, isActive: true, terminated: false })) await updateWorkerEntry(launch, worker, entry => withoutFailure({ ...entry, isActive: true, terminated: false }))
const accepted = await conversationService.sendMessage(worker.sessionId, formatTeammateMessages(messages)) // A member that has not started yet gets its approved instructions with the
// first thing it is told, unless that is only to shut down.
const starts = worker.brief !== undefined && !messages.every(message => isShutdownRequest(message.text))
if (messages.some(message => isShutdownRequest(message.text))) worker.shutdownRequestedAt = Date.now()
const text = formatTeammateMessages(messages)
const accepted = await conversationService.sendMessage(worker.sessionId, starts ? `${worker.brief}\n\n${text}` : text)
if (!accepted) { if (!accepted) {
await updateWorkerEntry(launch, worker, entry => ({ ...entry, isActive: false })) await updateWorkerEntry(launch, worker, entry => ({ ...entry, isActive: false }))
return return
} }
clearFailure(worker) clearFailure(worker)
if (starts) {
worker.brief = undefined
await updateWorkerEntry(launch, worker, ({ awaitingDependencies: _awaiting, ...entry }) => entry)
}
const ids = new Set(messages.map(message => message.id).filter(Boolean)) const ids = new Set(messages.map(message => message.id).filter(Boolean))
const legacy = new Set(messages.filter(message => !message.id).map(message => JSON.stringify([message.from, message.timestamp, message.text]))) const legacy = new Set(messages.filter(message => !message.id).map(message => JSON.stringify([message.from, message.timestamp, message.text])))
await markMessagesAsReadByPredicate(worker.member.name, message => message.id ? ids.has(message.id) : legacy.has(JSON.stringify([message.from, message.timestamp, message.text])), launch.plan.teamName) await markMessagesAsReadByPredicate(worker.member.name, message => message.id ? ids.has(message.id) : legacy.has(JSON.stringify([message.from, message.timestamp, message.text])), launch.plan.teamName)
@@ -435,7 +522,7 @@ async function deliverToWorker(launch: TeamLaunch, worker: WorkerRuntime, messag
async function wakeForReadyTask(launch: TeamLaunch, worker: WorkerRuntime, entry: MemberEntry | undefined, tasks: Task[]): Promise<void> { async function wakeForReadyTask(launch: TeamLaunch, worker: WorkerRuntime, entry: MemberEntry | undefined, tasks: Task[]): Promise<void> {
if (!entry || entry.isActive || worker.autoContinueTimer || !conversationService.hasSession(worker.sessionId)) return if (!entry || entry.isActive || worker.autoContinueTimer || !conversationService.hasSession(worker.sessionId)) return
if (worker.lastResultAt === undefined || Date.now() - worker.lastResultAt < timing.unblockedTaskWakeDelayMs) return if (worker.lastResultAt === undefined || Date.now() - worker.lastResultAt < timing.unblockedTaskWakeDelayMs) return
const unresolved = new Set(tasks.filter(task => task.status !== 'completed').map(task => task.id)) const unresolved = unresolvedTaskIds(tasks)
const ready = tasks.find(task => task.owner === worker.member.name && task.status === 'pending' && !worker.wokenForTaskIds.has(task.id) const ready = tasks.find(task => task.owner === worker.member.name && task.status === 'pending' && !worker.wokenForTaskIds.has(task.id)
&& task.blockedBy.length > 0 && task.blockedBy.every(id => !unresolved.has(id))) && task.blockedBy.length > 0 && task.blockedBy.every(id => !unresolved.has(id)))
if (!ready) return if (!ready) return
@@ -448,6 +535,41 @@ async function wakeForReadyTask(launch: TeamLaunch, worker: WorkerRuntime, entry
}]) }])
} }
/**
* Give a member that has not started yet its approved instructions once one
* of its tasks is ready, so it never works ahead of the dependencies the user
* approved. A Stop or a server restart may have ended its idle process; it
* starts again for this, as it would for a message.
*/
async function startWhenReady(launch: TeamLaunch, worker: WorkerRuntime, tasks: Task[], held: TeammateMessage[]): Promise<void> {
const ready = readyTasksOf(worker.member.name, tasks)
if (ready.length === 0) return
if (!conversationService.hasSession(worker.sessionId)) {
void migrationMaintenance.track(restartWorker(launch, worker, 'message'))
return
}
// The instructions already say what is ready; no separate wake follows.
for (const task of ready) worker.wokenForTaskIds.add(task.id)
// The members it depends on were told to send their results before
// completing. When one did not, say so now rather than let this member find
// out halfway through that it is working without its input.
const teammates = new Set([...launch.workers.values()].map(other => other.member.name))
const upstream = new Set(ready.flatMap(task => task.blockedBy).map(id => tasks.find(task => task.id === id)?.owner)
.filter((owner): owner is string => !!owner && owner !== worker.member.name && teammates.has(owner)))
const silent = [...upstream].filter(owner => !held.some(message => message.from === owner))
const missing = silent.length > 0
? `\nNothing has reached you from ${silent.join(', ')} yet. If you need ${silent.length === 1 ? 'its' : 'their'} result, ask with SendMessage before you rely on it.`
: ''
// What the lead and teammates wrote while it waited (the results it depends
// on) arrives with its instructions.
await deliverToWorker(launch, worker, [...held, {
from: 'task-list',
text: `Ready now (dependencies complete): ${ready.map(task => `#${task.id} ${task.subject}`).join('; ')}${missing}`,
timestamp: new Date().toISOString(),
read: false,
}])
}
async function superviseLaunch(launch: TeamLaunch): Promise<void> { async function superviseLaunch(launch: TeamLaunch): Promise<void> {
if (migrationMaintenance.isActive) return if (migrationMaintenance.isActive) return
const team = readTeamFile(launch.plan.teamName) const team = readTeamFile(launch.plan.teamName)
@@ -468,12 +590,6 @@ async function superviseLaunch(launch: TeamLaunch): Promise<void> {
await updateWorkerEntry(launch, worker, current => ({ ...current, isActive: false, terminated: true })) await updateWorkerEntry(launch, worker, current => ({ ...current, isActive: false, terminated: true }))
} }
const messages = await readUnreadMessages(worker.member.name, launch.plan.teamName) const messages = await readUnreadMessages(worker.member.name, launch.plan.teamName)
if (messages.length === 0) {
if (!leadAlive || launch.pausedAt) continue
tasks ??= await listTasks(getCanonicalTeamTaskListId(launch.plan.teamName)).catch(() => [])
await wakeForReadyTask(launch, worker, entry, tasks)
continue
}
if (launch.pausedAt) { if (launch.pausedAt) {
// Only a new instruction resumes a paused team: the user's own message to // Only a new instruction resumes a paused team: the user's own message to
// a member, or the lead's once the user has spoken to it again. // a member, or the lead's once the user has spoken to it again.
@@ -486,6 +602,18 @@ async function superviseLaunch(launch: TeamLaunch): Promise<void> {
launch.pausedAt = undefined launch.pausedAt = undefined
launch.resumeAllowedAt = undefined launch.resumeAllowedAt = undefined
} }
// A member that waits on dependencies is started by its tasks, or by the
// user writing to it. What the lead or a teammate sends it meanwhile (a
// forwarded result, a broadcast) waits in its inbox and arrives with its
// instructions; starting on such a message would undo the approved order.
const held = worker.brief !== undefined && !messages.some(message => message.from === 'user' || isShutdownRequest(message.text))
if (messages.length === 0 || held) {
if (!leadAlive) continue
tasks ??= await listTasks(getCanonicalTeamTaskListId(launch.plan.teamName)).catch(() => [])
if (worker.brief !== undefined) await startWhenReady(launch, worker, tasks, messages)
else await wakeForReadyTask(launch, worker, entry, tasks)
continue
}
if (!conversationService.hasSession(worker.sessionId)) { if (!conversationService.hasSession(worker.sessionId)) {
const fromUser = messages.some(message => message.from === 'user') const fromUser = messages.some(message => message.from === 'user')
// Without a lead nobody coordinates the restarted member; a direct user // Without a lead nobody coordinates the restarted member; a direct user
@@ -559,9 +687,18 @@ export async function noteLeadUserMessage(parentSessionId: string): Promise<void
if (waiting.length === 0) continue if (waiting.length === 0) continue
for (const worker of waiting) worker.wakeNoticePending = false for (const worker of waiting) worker.wakeNoticePending = false
const tasks = await listTasks(getCanonicalTeamTaskListId(launch.plan.teamName)).catch(() => [] as Task[]) const tasks = await listTasks(getCanonicalTeamTaskListId(launch.plan.teamName)).catch(() => [] as Task[])
const unresolved = unresolvedTaskIds(tasks)
const lines = waiting.map(worker => { const lines = waiting.map(worker => {
const open = tasks.filter(task => task.owner === worker.member.name && task.status !== 'completed') const open = tasks.filter(task => task.owner === worker.member.name && task.status !== 'completed')
const label = worker.failure ? `${worker.member.name} (stopped on an error: ${worker.failure})` : worker.member.name const label = worker.failure ? `${worker.member.name} (stopped on an error: ${worker.failure})` : worker.member.name
// A member whose start failed is reported like any other failure.
if (worker.brief !== undefined && !worker.failure) {
const ready = readyTasksOf(worker.member.name, tasks)
const blockers = [...new Set(open.flatMap(task => task.blockedBy.filter(id => unresolved.has(id))))]
if (ready.length > 0) return `- ${label}: not started yet; ${ready.map(task => `#${task.id} ${task.subject}`).join('; ')} is ready for it, so message it to start`
if (blockers.length > 0) return `- ${label}: not started yet; it starts by itself once ${blockers.map(id => `#${id}`).join(', ')} ${blockers.length === 1 ? 'is' : 'are'} completed, and messages to it wait until then`
return `- ${label}: not started yet; no unfinished task`
}
return open.length > 0 return open.length > 0
? `- ${label}: ${open.map(task => `#${task.id} ${task.subject} (${task.status})`).join('; ')}` ? `- ${label}: ${open.map(task => `#${task.id} ${task.subject} (${task.status})`).join('; ')}`
: `- ${label}: no unfinished task` : `- ${label}: no unfinished task`
@@ -655,6 +792,11 @@ export async function launchTeamPlanRuntime(plan: TeamPlanRecord): Promise<{ mem
await writeTeamFileAsync(plan.teamName, team) await writeTeamFileAsync(plan.teamName, team)
}) })
taskMapping = await materializeTasks(plan) taskMapping = await materializeTasks(plan)
// Every member's process starts now, but one whose tasks all wait on other
// tasks gets its instructions only once one of them is ready. Handing them
// out at once let members start before the work they depend on existed.
const startingTasks = await listTasks(getCanonicalTeamTaskListId(plan.teamName))
const waiting = new Set(members.filter(member => waitsForDependencies(member.name, startingTasks)).map(member => member.id))
await startTeamWorkersBarrier(members, async member => { await startTeamWorkersBarrier(members, async member => {
if (launch.stopped || !conversationService.hasSession(plan.sessionId)) throw new Error('Team launch cancelled') if (launch.stopped || !conversationService.hasSession(plan.sessionId)) throw new Error('Team launch cancelled')
const id = randomUUID() const id = randomUUID()
@@ -676,21 +818,30 @@ export async function launchTeamPlanRuntime(plan: TeamPlanRecord): Promise<{ mem
team.members.push({ agentId: `${entry.name}@${plan.teamName}`, name: entry.name, agentType: entry.agentType, team.members.push({ agentId: `${entry.name}@${plan.teamName}`, name: entry.name, agentType: entry.agentType,
model: entry.runtime.modelId, providerId: entry.runtime.providerId, providerName: typeof entry.providerName === 'string' ? entry.providerName : undefined, effortLevel: entry.runtime.effortLevel, model: entry.runtime.modelId, providerId: entry.runtime.providerId, providerName: typeof entry.providerName === 'string' ? entry.providerName : undefined, effortLevel: entry.runtime.effortLevel,
planMemberId: entry.id, joinedAt: Date.now(), tmuxPaneId: '', cwd: plan.workDir, subscriptions: [], planMemberId: entry.id, joinedAt: Date.now(), tmuxPaneId: '', cwd: plan.workDir, subscriptions: [],
sessionId: memberIds[entry.id], backendType: 'process', isActive: true }) sessionId: memberIds[entry.id], backendType: 'process',
...(waiting.has(entry.id) ? { isActive: false, awaitingDependencies: true } : { isActive: true }) })
} }
await writeTeamFileAsync(plan.teamName, team) await writeTeamFileAsync(plan.teamName, team)
}) })
await sendTeamSnapshot(plan.sessionId, plan.teamName, launch.createdAt) await sendTeamSnapshot(plan.sessionId, plan.teamName, launch.createdAt)
} }
const assigned = approved.tasks.filter(task => task.ownerId === member.id).map(task => ({ ...task, id: taskMapping[task.id], dependencies: task.dependencies.map(dep => taskMapping[dep]) })) const brief = memberBrief(member, members, approved.tasks, taskMapping)
await withTaskListLifecycleLock(getCanonicalTeamTaskListId(plan.teamName), async () => { await withTaskListLifecycleLock(getCanonicalTeamTaskListId(plan.teamName), async () => {
const latest = await readTeamPlan(plan.teamName) const latest = await readTeamPlan(plan.teamName)
if (launch.stopped || !latest || latest.planId !== plan.planId || latest.incarnationId !== plan.incarnationId || latest.state !== 'launching' || latest.revision !== plan.revision || latest.approvedSnapshot?.requestId !== approved.requestId) throw new Error('Team launch authorization has been revoked') if (launch.stopped || !latest || latest.planId !== plan.planId || latest.incarnationId !== plan.incarnationId || latest.state !== 'launching' || latest.revision !== plan.revision || latest.approvedSnapshot?.requestId !== approved.requestId) throw new Error('Team launch authorization has been revoked')
executionStarted = true executionStarted = true
launch.released = true launch.released = true
const sent = await conversationService.sendMessage(id, `${member.agentSnapshot?.initialPrompt ? member.agentSnapshot.initialPrompt + "\n" : ""}${member.prompt}\n\nApproved shared tasks (use TaskGet/TaskUpdate; respect dependencies):\n${JSON.stringify(assigned)}`, undefined, { canSend: () => !launch.stopped })
if (!sent) throw new Error(`Failed to release ${member.name}`)
const worker = launch.workers.get(member.id) const worker = launch.workers.get(member.id)
if (waiting.has(member.id)) {
// The supervisor delivers it when a task is ready (startWhenReady).
if (worker) {
worker.brief = brief
worker.released = true
}
return
}
const sent = await conversationService.sendMessage(id, brief, undefined, { canSend: () => !launch.stopped })
if (!sent) throw new Error(`Failed to release ${member.name}`)
if (worker) worker.released = true if (worker) worker.released = true
}) })
}, id => conversationService.stopSessionAndWait(id)) }, id => conversationService.stopSessionAndWait(id))
@@ -699,7 +850,11 @@ export async function launchTeamPlanRuntime(plan: TeamPlanRecord): Promise<{ mem
// remains connected while idle and wakes when another teammate writes. // remains connected while idle and wakes when another teammate writes.
launch.running = true launch.running = true
startSupervisor(launch) startSupervisor(launch)
await conversationService.sendMessage(plan.sessionId, `Approved team ${plan.teamName} is running. Members and their shared tasks are ready. Continue coordinating the approved team; do not spawn these members again. Members report back automatically when they finish or fail; you do not need to poll or wait in a loop.`) const waitingNames = members.filter(member => waiting.has(member.id)).map(member => member.name)
const waitingNote = waitingNames.length > 0
? ` Every task of ${waitingNames.join(', ')} waits on other tasks: such a member gets its instructions automatically once one of its tasks is ready. Messages sent to it before that wait in its inbox and arrive with those instructions. Each member is told to send its result to the members whose tasks depend on it before completing, so forward a result yourself only when a member says one is missing.`
: ''
await conversationService.sendMessage(plan.sessionId, `Approved team ${plan.teamName} is running. Members and their shared tasks are ready.${waitingNote} Continue coordinating the approved team; do not spawn these members again. Members report back automatically when they finish or fail; you do not need to poll or wait in a loop.`)
return { memberIds } return { memberIds }
} catch (error) { } catch (error) {
await stopTeamPlanRuntime(plan.planId) await stopTeamPlanRuntime(plan.planId)
@@ -740,14 +895,25 @@ export async function rehydrateTeamPlanRuntime(plan: TeamPlanRecord): Promise<bo
if (!team || team.leadSessionId !== plan.sessionId || incarnationOf(team) !== plan.incarnationId) return false if (!team || team.leadSessionId !== plan.sessionId || incarnationOf(team) !== plan.incarnationId) return false
const memberIds = plan.launch?.memberIds ?? {} const memberIds = plan.launch?.memberIds ?? {}
const workers = new Map<string, WorkerRuntime>() const workers = new Map<string, WorkerRuntime>()
let taskMapping: Record<string, string> | undefined
for (const member of plan.approvedSnapshot.members) { for (const member of plan.approvedSnapshot.members) {
const sessionId = memberIds[member.id] ?? team.members.find(entry => entry.planMemberId === member.id && entry.backendType === 'process')?.sessionId const sessionId = memberIds[member.id] ?? team.members.find(entry => entry.planMemberId === member.id && entry.backendType === 'process')?.sessionId
if (!sessionId) continue if (!sessionId) continue
const entry = team.members.find(candidate => candidate.sessionId === sessionId) const entry = team.members.find(candidate => candidate.sessionId === sessionId)
if (!entry) continue if (!entry) continue
workers.set(member.id, { member, sessionId, released: true, autoContinueAttempts: 0, restartHistory: [], restartsExhausted: false, wokenForTaskIds: new Set() }) const worker: WorkerRuntime = { member, sessionId, released: true, autoContinueAttempts: 0, restartHistory: [], restartsExhausted: false, wokenForTaskIds: new Set() }
// Still waiting for its first ready task: rebuild the instructions it was approved with.
if (entry.awaitingDependencies === true) {
taskMapping ??= Object.fromEntries((await listTasks(getCanonicalTeamTaskListId(plan.teamName)).catch(() => [] as Task[]))
.filter(task => task.metadata?.teamPlanId === plan.planId && typeof task.metadata?.teamPlanTaskId === 'string')
.map(task => [task.metadata!.teamPlanTaskId as string, task.id]))
worker.brief = memberBrief(member, plan.approvedSnapshot.members, plan.approvedSnapshot.tasks, taskMapping)
}
workers.set(member.id, worker)
} }
if (workers.size === 0) return false if (workers.size === 0) return false
// Another start may have re-owned the team while the task list was read.
if (launches.has(plan.planId)) return true
ensureRuntimeHooks() ensureRuntimeHooks()
const launch: TeamLaunch = { const launch: TeamLaunch = {
parentId: plan.sessionId, plan, createdAt: team.createdAt, released: true, running: true, parentId: plan.sessionId, plan, createdAt: team.createdAt, released: true, running: true,
+57
View File
@@ -24,6 +24,7 @@ import {
updateTask, updateTask,
withTaskListLifecycleLock, withTaskListLifecycleLock,
} from '../utils/tasks.js' } from '../utils/tasks.js'
import { clearDynamicTeamContext, setDynamicTeamContext } from '../utils/teammate.js'
import { TaskCreateTool } from './TaskCreateTool/TaskCreateTool.js' import { TaskCreateTool } from './TaskCreateTool/TaskCreateTool.js'
import { TaskGetTool } from './TaskGetTool/TaskGetTool.js' import { TaskGetTool } from './TaskGetTool/TaskGetTool.js'
import { TaskListTool } from './TaskListTool/TaskListTool.js' import { TaskListTool } from './TaskListTool/TaskListTool.js'
@@ -528,3 +529,59 @@ describe('Task tool execution ordering', () => {
} }
}) })
}) })
describe('Task dependencies', () => {
it('a teammate cannot start or finish a task ahead of the tasks it waits on, the lead can', async () => {
const configDir = await mkdtemp(join(tmpdir(), 'task-tool-dependencies-'))
const taskListId = 'task-tool-dependency-team'
const previousConfigDir = process.env.CLAUDE_CONFIG_DIR
const previousTaskListId = process.env.CLAUDE_CODE_TASK_LIST_ID
process.env.CLAUDE_CONFIG_DIR = configDir
process.env.CLAUDE_CODE_TASK_LIST_ID = taskListId
let appState: Record<string, unknown> = { expandedView: undefined, inbox: { messages: [] } }
const context = {
abortController: new AbortController(),
getAppState: () => appState,
setAppState: (update: (prev: Record<string, unknown>) => Record<string, unknown>) => {
appState = update(appState)
},
} as unknown as ToolUseContext
const asTeammate = () => setDynamicTeamContext({ agentId: `analyst@${taskListId}`, agentName: 'analyst', teamName: taskListId, planModeRequired: false })
const statusOf = async (taskId: string) => (await TaskGetTool.call({ taskId }, context)).data.task?.status
try {
const inventory = (await TaskCreateTool.call({ subject: 'Inventory', description: 'Map the release' }, context)).data.task.id
const review = (await TaskCreateTool.call({ subject: 'Review', description: 'Needs the inventory' }, context)).data.task.id
await TaskUpdateTool.call({ taskId: review, addBlockedBy: [inventory] }, context)
// The reported run: the analyst marked its review in progress while the
// inventory it depends on had not even been started.
asTeammate()
for (const status of ['in_progress', 'completed'] as const) {
const refused = await TaskUpdateTool.call({ taskId: review, status }, context)
expect(refused.data).toMatchObject({ success: false, updatedFields: [] })
expect(refused.data.error).toContain(`Task #${review} waits on #${inventory}`)
}
expect(await statusOf(review)).toBe('pending')
// Everything but the status still goes through.
expect((await TaskUpdateTool.call({ taskId: review, description: 'Needs the inventory first' }, context)).data.success).toBe(true)
// The lead coordinates and may decide the dependency no longer matters.
clearDynamicTeamContext()
expect((await TaskUpdateTool.call({ taskId: review, status: 'in_progress' }, context)).data.success).toBe(true)
await updateTask(taskListId, review, { status: 'pending' })
await updateTask(taskListId, inventory, { status: 'completed' })
asTeammate()
expect((await TaskUpdateTool.call({ taskId: review, status: 'in_progress' }, context)).data.success).toBe(true)
expect((await TaskUpdateTool.call({ taskId: review, status: 'completed' }, context)).data.success).toBe(true)
} finally {
clearDynamicTeamContext()
if (previousConfigDir === undefined) delete process.env.CLAUDE_CONFIG_DIR
else process.env.CLAUDE_CONFIG_DIR = previousConfigDir
if (previousTaskListId === undefined) delete process.env.CLAUDE_CODE_TASK_LIST_ID
else process.env.CLAUDE_CODE_TASK_LIST_ID = previousTaskListId
await rm(configDir, { recursive: true, force: true })
}
})
})
@@ -24,6 +24,7 @@ import {
getAgentName, getAgentName,
getTeammateColor, getTeammateColor,
getTeamName, getTeamName,
isTeammate,
} from '../../utils/teammate.js' } from '../../utils/teammate.js'
import { writeToMailbox } from '../../utils/teammateMailbox.js' import { writeToMailbox } from '../../utils/teammateMailbox.js'
import { VERIFICATION_AGENT_TYPE } from '../AgentTool/constants.js' import { VERIFICATION_AGENT_TYPE } from '../AgentTool/constants.js'
@@ -164,6 +165,33 @@ export const TaskUpdateTool = buildTool({
} }
} }
// A teammate cannot start or finish a task ahead of the tasks it waits on;
// without this the list's dependencies are only a suggestion, and a member
// marks its blocked task in progress while the work it needs is unwritten.
// The lead is not held to it: it coordinates, and may decide a dependency
// no longer matters.
if (
isTeammate() &&
(status === 'in_progress' || status === 'completed') &&
status !== existingTask.status &&
existingTask.blockedBy.length > 0
) {
const unresolved = new Set(
(await listTasks(taskListId)).filter(task => task.status !== 'completed').map(task => task.id),
)
const waitingOn = existingTask.blockedBy.filter(id => unresolved.has(id))
if (waitingOn.length > 0) {
return {
data: {
success: false,
taskId,
updatedFields: [],
error: `Task #${taskId} waits on ${waitingOn.map(id => `#${id}`).join(', ')}, not completed yet. Do not work on it or change its status until then; TaskList shows which of your tasks are ready.`,
},
}
}
}
const updatedFields: string[] = [] const updatedFields: string[] = []
if (status !== undefined) { if (status !== undefined) {
// Handle deletion - delete the task file and return early // Handle deletion - delete the task file and return early
@@ -87,6 +87,14 @@ describe('whole-team planning tools', () => {
delete process.env.CC_HAHA_TEAM_REVIEW_REQUIRED delete process.env.CC_HAHA_TEAM_REVIEW_REQUIRED
expect(getPrompt()).not.toContain('Human review before execution') expect(getPrompt()).not.toContain('Human review before execution')
}) })
test('planning instructions say a dependent member starts only once its dependencies are done', async () => {
// A lead once wrote "if recon has not reported yet, start on your own",
// and the members did, ahead of the dependency the user approved.
const prompt = await TeamPlanTool.prompt({} as never)
expect(prompt).toContain('receives its prompt only once one of them is ready')
expect(prompt).toContain('never tell a member to start without them')
})
test('legacy Agent spawns only stage members, then whole-plan submit waits for human review', async () => { test('legacy Agent spawns only stage members, then whole-plan submit waits for human review', async () => {
const created = await TeamCreateTool.call({ team_name: 'review-team' }, context) const created = await TeamCreateTool.call({ team_name: 'review-team' }, context)
expect(created.data.plan?.state).toBe('draft') expect(created.data.plan?.state).toBe('draft')
+1 -1
View File
@@ -28,7 +28,7 @@ export const TeamPlanTool = buildTool({
isEnabled() { return isAgentSwarmsEnabled() && isTeamReviewRequired() && !isTeammate() }, isEnabled() { return isAgentSwarmsEnabled() && isTeamReviewRequired() && !isTeammate() },
async description() { return 'Read, replace or submit a team draft for human review. Never starts members or approves a plan.' }, async description() { return 'Read, replace or submit a team draft for human review. Never starts members or approves a plan.' },
async prompt() { async prompt() {
return 'After TeamCreate, use TeamPlan to submit the complete roster and tasks together. Give each member a stable id, a launchable name using only letters, numbers, underscores or hyphens (team-lead is reserved), available agentType, task prompt, suggested runtime and a short reason. Task ownerId refers to a member id; dependencies refer to task ids. Use get to read the current revision. Replace and submit require expected_revision. Submit may include a complete replacement plan. After submit, end the planning turn and wait for human approval. Do not start work, claim tasks, poll the plan or call Agent to bypass review. Only the user can approve in the team panel.' return 'After TeamCreate, use TeamPlan to submit the complete roster and tasks together. Give each member a stable id, a launchable name using only letters, numbers, underscores or hyphens (team-lead is reserved), available agentType, task prompt, suggested runtime and a short reason. Task ownerId refers to a member id; dependencies refer to task ids. A member whose tasks all have dependencies receives its prompt only once one of them is ready, so write that prompt for the moment its dependencies are complete and never tell a member to start without them. Use get to read the current revision. Replace and submit require expected_revision. Submit may include a complete replacement plan. After submit, end the planning turn and wait for human approval. Do not start work, claim tasks, poll the plan or call Agent to bypass review. Only the user can approve in the team panel.'
}, },
toAutoClassifierInput(input) { return `${input.operation} ${input.team_name}` }, toAutoClassifierInput(input) { return `${input.operation} ${input.team_name}` },
renderToolUseMessage(input) { return `${input.operation} team plan: ${input.team_name}` }, renderToolUseMessage(input) { return `${input.operation} team plan: ${input.team_name}` },
+6
View File
@@ -110,6 +110,12 @@ export type TeamFile = {
lastError?: string lastError?: string
/** Pending automatic continuation after a transient provider failure. */ /** Pending automatic continuation after a transient provider failure. */
autoRetry?: { attempt: number; max: number; nextAt: number } autoRetry?: { attempt: number; max: number; nextAt: number }
/**
* The member has not been given its approved instructions yet: every task
* it owns waited on unfinished work when the team started. Absent on
* members written before this was recorded, which were all instructed.
*/
awaitingDependencies?: boolean
}> }>
} }
+20
View File
@@ -15,6 +15,7 @@ import {
getInboxPath, getInboxPath,
IDLE_RESULT_MAX_CHARS, IDLE_RESULT_MAX_CHARS,
isIdleNotification, isIdleNotification,
isShutdownApproved,
markMessagesAsReadByIdentity, markMessagesAsReadByIdentity,
markMessagesAsReadByPredicate, markMessagesAsReadByPredicate,
readMailbox, readMailbox,
@@ -533,3 +534,22 @@ describe('teammate message formatting', () => {
) )
}) })
}) })
describe('shutdown approvals', () => {
const approval = (extra: Record<string, unknown>) => JSON.stringify({
type: 'shutdown_approved', requestId: 'shutdown-1@reader', from: 'reader', timestamp: new Date().toISOString(), ...extra,
})
test('are recognised for teammates the lead runs itself', () => {
expect(isShutdownApproved(approval({ paneId: '%1', backendType: 'tmux' }))?.from).toBe('reader')
expect(isShutdownApproved(approval({ backendType: 'in-process' }))?.from).toBe('reader')
})
// Recognising it makes the lead drop the member from the team file and
// unassign its tasks. The desktop keeps a stopped member on the roster so a
// message can restart it, and wakes waiting members by the approved task
// owners; adding `process` to the schema would quietly break both.
test('leave a desktop member to the server runtime', () => {
expect(isShutdownApproved(approval({ paneId: '', backendType: 'process' }))).toBeNull()
})
})
+8
View File
@@ -1040,6 +1040,14 @@ export type ShutdownRequestMessage = z.infer<
/** /**
* Shutdown approved message sent from teammate to leader via mailbox * Shutdown approved message sent from teammate to leader via mailbox
*
* `backendType` deliberately leaves out `process`, the desktop's members.
* Whoever recognises an approval here removes the member from the team file
* and unassigns its unfinished tasks, which is right for a teammate the lead
* itself runs. A desktop member belongs to the server's team runtime instead:
* it stays on the roster as stopped so a message can restart it, and its tasks
* keep the owners the user approved, which is what wakes the members waiting
* on them. Its approval therefore reaches the lead as an ordinary message.
*/ */
export const ShutdownApprovedMessageSchema = lazySchema(() => export const ShutdownApprovedMessageSchema = lazySchema(() =>
z.strictObject({ z.strictObject({