fix(agent-teams): keep a stopped team stopped and wake members after any failure (#1446)

* fix(agent-teams): keep a stopped team stopped and wake members after any failure

Stop killed every member, but the lead kept polling its mailbox: reports
sent before the Stop started hidden lead turns, and a single SendMessage
from one of them unpaused the whole team. Stop now sends team_plan_pause
so the lead holds teammate mail until the user's next message, and the
server only lets lead instructions written after that message resume the
team. A restart in flight when Stop arrives is stopped as well, and the
idle Stop button hides once every member is stopped.

Members that stopped on a usage limit, billing, network or crash failure
were only resumed if the lead happened to message them. The user's next
message now lists every stopped member with its failure and open tasks,
once per failure, whether or not Stop was pressed.

* feat(desktop): show where each team member stands on its tasks

A member that stopped or failed left its task reading "In progress" with
an animated bar, and its card said only "Stopped" or "Error". Tasks whose
owner stopped, failed or waits to retry now show that state without
animation, and the card says which task the member stopped on.

The member card chips tell done, current and next tasks apart, with the
subject on hover. The member drawer groups tasks into now, up next and
done, names the unfinished dependencies of blocked work, shows done as
3/5, and no longer prints +0:00 for a duration polling could not measure.
This commit is contained in:
程序员阿江-Relakkes
2026-10-04 19:31:10 +08:00
committed by GitHub
parent 5dc59c5892
commit 9623064b15
22 changed files with 843 additions and 126 deletions
@@ -402,6 +402,32 @@ describe('AgentTeamsCanvas', () => {
expect(lead).toBeNull() expect(lead).toBeNull()
}) })
it('chips what each member has done, is doing and does next, with the subject on hover', () => {
render(<AgentTeamsCanvas {...props()} />)
const chips = (agentId: string) => Array.from(
screen.getByTestId(`agent-teams-canvas-member-${agentId}`).querySelectorAll('[data-member-task]'),
).map(chip => ({
id: chip.getAttribute('data-member-task'),
status: chip.getAttribute('data-task-status'),
next: chip.getAttribute('data-task-next') === 'true',
title: chip.getAttribute('title'),
}))
expect(chips('builder@canvas-team')).toEqual([
{ id: '1', status: 'completed', next: false, title: 'Task 1' },
{ id: '2', status: 'in_progress', next: false, title: 'Task 2' },
])
// Nothing started yet, but the member's next task is already known.
expect(chips('reviewer@canvas-team')).toEqual([
{ id: '3', status: 'pending', next: true, title: 'Task 3' },
])
// The connector line runs behind the chip row; a see-through chip lets
// it strike through the label.
const next = screen.getByTestId('agent-teams-canvas-member-reviewer@canvas-team').querySelector<HTMLElement>('[data-task-next]')!
expect(next.style.backgroundColor).not.toBe('transparent')
expect(next.style.backgroundColor).not.toBe('')
})
it('labels stopped, retrying and failed members and keeps their failure on hover', () => { it('labels stopped, retrying and failed members and keeps their failure on hover', () => {
const current = snapshot('current') const current = snapshot('current')
const recovering: TeamWorkbenchSnapshot = { const recovering: TeamWorkbenchSnapshot = {
@@ -436,7 +462,15 @@ describe('AgentTeamsCanvas', () => {
const builder = screen.getByTestId('agent-teams-canvas-member-builder@canvas-team') const builder = screen.getByTestId('agent-teams-canvas-member-builder@canvas-team')
expect(builder.getAttribute('data-member-state')).toBe('retrying') expect(builder.getAttribute('data-member-state')).toBe('retrying')
expect(screen.getByText('Auto-retry 2/5').getAttribute('title')).toBe('API Error: 529 overloaded') // It names the task it stopped on, not just that it stopped.
expect(screen.getByText('Auto-retry 2/5 · at #2').getAttribute('title')).toBe('API Error: 529 overloaded')
// The task list still says in progress; the card must not animate work
// nobody is doing.
const stalledTask = screen.getByTestId('agent-teams-canvas-task-2')
expect(stalledTask.getAttribute('data-stalled')).toBe('retrying')
expect(stalledTask.textContent).toContain('Auto-retry 2/5')
expect(stalledTask.textContent).not.toContain('In progress')
expect(stalledTask.querySelector('[data-progress="indeterminate"]')).toBeNull()
const reviewer = screen.getByTestId('agent-teams-canvas-member-reviewer@canvas-team') const reviewer = screen.getByTestId('agent-teams-canvas-member-reviewer@canvas-team')
expect(reviewer.getAttribute('data-member-state')).toBe('stopped') expect(reviewer.getAttribute('data-member-state')).toBe('stopped')
@@ -17,9 +17,11 @@ import {
parseWorkbenchMessageBody, parseWorkbenchMessageBody,
resolveMemberModel, resolveMemberModel,
resolveTeamMemberIdentity, resolveTeamMemberIdentity,
stalledTaskOwnerState,
taskOwnedByMember, taskOwnedByMember,
type MemberWorkState, type MemberWorkState,
type PositionedWorkbenchTask, type PositionedWorkbenchTask,
type StalledTaskOwnerState,
type WorkbenchTaskState, type WorkbenchTaskState,
} from './agentTeamsModel' } from './agentTeamsModel'
@@ -50,6 +52,8 @@ type MemberPosition = {
percent: number percent: number
inbox: number inbox: number
recentTasks: TeamWorkbenchTask[] recentTasks: TeamWorkbenchTask[]
/** The first of the member's tasks that has not started yet. */
nextTask?: TeamWorkbenchTask
} }
type OwnerVisual = { type OwnerVisual = {
@@ -138,6 +142,32 @@ function memberStateLabel(state: MemberWorkState, member: TeamMember, t: Transla
return t(`agentTeams.member.${state}` as TranslationKey) return t(`agentTeams.member.${state}` as TranslationKey)
} }
function isStalledState(state: MemberWorkState): state is StalledTaskOwnerState {
return state === 'stopped' || state === 'error' || state === 'retrying'
}
function stalledColors(state: StalledTaskOwnerState) {
if (state === 'error') {
return { background: 'var(--color-error-container)', foreground: 'var(--color-on-error-container)', border: 'var(--color-error)' }
}
if (state === 'retrying') {
return { background: 'var(--color-warning-container)', foreground: 'var(--color-on-warning-container)', border: 'var(--color-warning)' }
}
return { background: 'var(--color-surface-container-high)', foreground: 'var(--color-text-secondary)', border: 'var(--color-border-strong)' }
}
function memberTaskChipStyle(task: TeamWorkbenchTask, next: boolean, accent: string) {
if (next) {
// Opaque like the other chips: the formation's connector line runs
// behind this row and would cut through the label.
return { borderStyle: 'dashed', borderColor: 'var(--color-outline)', backgroundColor: 'var(--color-background)', color: 'var(--color-text-tertiary)' }
}
if (task.status === 'completed') {
return { borderColor: 'var(--color-border)', backgroundColor: 'var(--color-success-container)', color: 'var(--color-on-success-container)' }
}
return { borderColor: accent, backgroundColor: 'var(--color-surface-container-lowest)', color: 'var(--color-text-primary)' }
}
function leadStatusLabel(snapshot: TeamWorkbenchSnapshot, t: TranslationFn): string { function leadStatusLabel(snapshot: TeamWorkbenchSnapshot, t: TranslationFn): string {
if (snapshot.deletedAt) return t('agentTeams.lead.archived') if (snapshot.deletedAt) return t('agentTeams.lead.archived')
if (snapshot.tasks.length === 0) return t('agentTeams.lead.forming') if (snapshot.tasks.length === 0) return t('agentTeams.lead.forming')
@@ -254,6 +284,18 @@ function memberInboxCount(member: TeamMember, messages: TeamWorkbenchMessage[]):
)).length )).length
} }
/**
* The chips under a member card: what it has done and is doing, then the next
* task waiting for it, within the four chips the card has room for.
*/
function memberTaskChips(ownedTasks: TeamWorkbenchTask[]): Pick<MemberPosition, 'recentTasks' | 'nextTask'> {
const nextTask = ownedTasks
.filter(task => task.status === 'pending')
.sort((left, right) => left.id.localeCompare(right.id, undefined, { numeric: true }))[0]
const started = ownedTasks.filter(task => task.status !== 'pending')
return { recentTasks: started.slice(nextTask ? -3 : -4), nextTask }
}
function workState( function workState(
member: TeamMember, member: TeamMember,
snapshot: TeamWorkbenchSnapshot, snapshot: TeamWorkbenchSnapshot,
@@ -479,7 +521,9 @@ function MemberNode({
: waitingDependency : waitingDependency
? t('agentTeams.member.waitingForDependency', { task: waitingDependency }) ? t('agentTeams.member.waitingForDependency', { task: waitingDependency })
: t('agentTeams.member.waitingForTask') : t('agentTeams.member.waitingForTask')
: memberStateLabel(state, member, t) : isStalledState(state) && position.currentTask
? t('agentTeams.member.stalledOnTask', { state: memberStateLabel(state, member, t), task: position.currentTask.id })
: memberStateLabel(state, member, t)
const characterClass = state === 'working' const characterClass = state === 'working'
? 'agent-teams-character-working' ? 'agent-teams-character-working'
: state === 'idle' || state === 'retrying' : state === 'idle' || state === 'retrying'
@@ -582,14 +626,22 @@ function MemberNode({
<span className="block h-full rounded-full" style={{ width: `${position.percent}%`, backgroundColor: accent }} /> <span className="block h-full rounded-full" style={{ width: `${position.percent}%`, backgroundColor: accent }} />
</span> </span>
<span className="mt-1.5 flex max-w-[172px] items-center justify-center gap-1 overflow-hidden"> <span className="mt-1.5 flex max-w-[172px] items-center justify-center gap-1 overflow-hidden">
{position.recentTasks.map(task => ( {[...position.recentTasks, ...(position.nextTask ? [position.nextTask] : [])].map(task => {
<span const next = task === position.nextTask
key={task.id} return (
className="rounded-full border border-[var(--color-border)] bg-[var(--color-surface-container-high)] px-1.5 py-px font-mono text-[9px] font-extrabold text-[var(--color-text-secondary)]" <span
> key={task.id}
#{task.id} data-member-task={task.id}
</span> data-task-status={task.status}
))} data-task-next={next ? 'true' : undefined}
title={task.subject}
className="rounded-full border px-1.5 py-px font-mono text-[9px] font-extrabold"
style={memberTaskChipStyle(task, next, accent)}
>
#{task.id}
</span>
)
})}
</span> </span>
</> </>
) : null} ) : null}
@@ -628,6 +680,10 @@ function TaskCard({
const owner = taskOwnerVisual(task, snapshot, members, depth) const owner = taskOwnerVisual(task, snapshot, members, depth)
const accent = owner?.accent ?? 'var(--color-brand)' const accent = owner?.accent ?? 'var(--color-brand)'
const colors = taskStateColors(state, accent) const colors = taskStateColors(state, accent)
// Still in progress on the list, but its owner stopped, failed or waits to
// retry: say so instead of animating work nobody is doing.
const stalled = owner ? stalledTaskOwnerState(task, snapshot) : undefined
const stall = stalled ? stalledColors(stalled) : undefined
const progress = taskProgress(task) const progress = taskProgress(task)
const dependencies = task.blockedBy.map(id => `#${id}`).join(' ') const dependencies = task.blockedBy.map(id => `#${id}`).join(' ')
const ownerLabel = owner const ownerLabel = owner
@@ -645,7 +701,8 @@ function TaskCard({
data-state={state} data-state={state}
data-depth={depth} data-depth={depth}
data-chain-active={focused ? 'true' : 'false'} data-chain-active={focused ? 'true' : 'false'}
aria-label={`${task.subject}, ${taskStateLabel(state, t)}`} data-stalled={stalled}
aria-label={`${task.subject}, ${stalled && owner ? memberStateLabel(stalled, owner.member, t) : taskStateLabel(state, t)}`}
onClick={onSelect} onClick={onSelect}
onMouseEnter={onHover} onMouseEnter={onHover}
onMouseLeave={onHoverEnd} onMouseLeave={onHoverEnd}
@@ -656,7 +713,7 @@ function TaskCard({
left: x, left: x,
top: y, top: y,
backgroundColor: colors.background, backgroundColor: colors.background,
borderColor: justUnlocked ? 'var(--color-success)' : colors.border, borderColor: justUnlocked ? 'var(--color-success)' : stall?.border ?? colors.border,
opacity: dimmed ? 0.34 : state === 'blocked' ? 0.72 : 1, opacity: dimmed ? 0.34 : state === 'blocked' ? 0.72 : 1,
}} }}
> >
@@ -667,12 +724,16 @@ function TaskCard({
<span <span
className="shrink-0 rounded-full border px-1.5 py-px text-[9.5px] font-extrabold" className="shrink-0 rounded-full border px-1.5 py-px text-[9.5px] font-extrabold"
style={{ style={{
backgroundColor: justUnlocked ? 'var(--color-success-container)' : colors.pillBackground, backgroundColor: justUnlocked ? 'var(--color-success-container)' : stall?.background ?? colors.pillBackground,
borderColor: justUnlocked ? 'var(--color-success)' : colors.border, borderColor: justUnlocked ? 'var(--color-success)' : stall?.border ?? colors.border,
color: justUnlocked ? 'var(--color-on-success-container)' : colors.pillForeground, color: justUnlocked ? 'var(--color-on-success-container)' : stall?.foreground ?? colors.pillForeground,
}} }}
> >
{justUnlocked ? t('agentTeams.task.unlocked') : taskStateLabel(state, t)} {justUnlocked
? t('agentTeams.task.unlocked')
: stalled && owner
? memberStateLabel(stalled, owner.member, t)
: taskStateLabel(state, t)}
</span> </span>
</span> </span>
@@ -706,7 +767,7 @@ function TaskCard({
: 'var(--color-surface-container-high)', : 'var(--color-surface-container-high)',
}} }}
> >
{progress === null ? ( {stalled && progress === null ? null : progress === null ? (
<span <span
data-progress="indeterminate" data-progress="indeterminate"
className="agent-teams-task-running-fill block h-full rounded-full" className="agent-teams-task-running-fill block h-full rounded-full"
@@ -716,7 +777,7 @@ function TaskCard({
<span <span
data-progress={Math.round(progress)} data-progress={Math.round(progress)}
className="block h-full rounded-full" className="block h-full rounded-full"
style={{ width: `${progress}%`, backgroundColor: colors.progress }} style={{ width: `${progress}%`, backgroundColor: stall ? stall.border : colors.progress }}
/> />
)} )}
</span> </span>
@@ -779,7 +840,7 @@ export function AgentTeamsCanvas({
total: ownedTasks.length, total: ownedTasks.length,
percent: ownedTasks.length === 0 ? 0 : Math.round((completed / ownedTasks.length) * 100), percent: ownedTasks.length === 0 ? 0 : Math.round((completed / ownedTasks.length) * 100),
inbox: memberInboxCount(member, snapshot.messages), inbox: memberInboxCount(member, snapshot.messages),
recentTasks: ownedTasks.filter(task => task.status !== 'pending').slice(-4), ...memberTaskChips(ownedTasks),
}) })
} }
@@ -125,7 +125,59 @@ describe('AgentTeamsMemberInspector', () => {
expect(row.textContent).toContain('Completed') expect(row.textContent).toContain('Completed')
expect(row.textContent).toContain(`${expectedStart} +7:00`) expect(row.textContent).toContain(`${expectedStart} +7:00`)
expect(row.getAttribute('data-task-state')).toBe('completed') expect(row.getAttribute('data-task-state')).toBe('completed')
expect(screen.getByText('1', { selector: 'dd' })).toBeTruthy() expect(screen.getByText('1/1', { selector: 'dd' })).toBeTruthy()
})
it('groups what a member is doing, will do next and has done, and says where a stopped member stopped', () => {
const memberTask = (id: string, status: TeamWorkbenchTask['status'], blockedBy: string[] = []): TeamWorkbenchTask => ({
id, subject: `Task ${id}`, description: '', owner: 'builder', status, blocks: [], blockedBy, taskListId: 'team-a',
})
const stopped: TeamMember = { ...builder, activity: 'stopped' }
const frame = (generatedAt: string, tasks: TeamWorkbenchTask[]): TeamWorkbenchSnapshot => ({
version: 'v1',
generatedAt,
team: { name: 'team-a', leadAgentId: 'lead@team-a', leadSessionId: 'lead-session', members: [stopped, reviewer] },
tasks,
messages: [],
})
const snapshots = [
frame('2026-08-08T07:00:00.000Z', [memberTask('1', 'pending'), memberTask('2', 'pending'), memberTask('3', 'pending', ['2']), memberTask('4', 'pending')]),
// #1 started and finished between two polls: its duration is unknown,
// not zero.
frame('2026-08-08T07:05:00.000Z', [memberTask('1', 'completed'), memberTask('2', 'in_progress'), memberTask('3', 'pending', ['2']), memberTask('4', 'pending')]),
]
render(
<AgentTeamsMemberInspector
snapshots={snapshots}
selectedIndex={1}
snapshot={snapshots[1]!}
member={stopped}
isLead={false}
leadIsStreaming={false}
onBack={vi.fn()}
onClose={vi.fn()}
onOpenExecution={vi.fn()}
/>,
)
const group = (name: string) => screen.getByTestId(`agent-teams-member-task-group-${name}`)
const ids = (name: string) => Array.from(group(name).querySelectorAll('[data-task-state]')).map(row => row.getAttribute('data-testid'))
expect(screen.getAllByTestId(/^agent-teams-member-task-group-/).map(node => node.getAttribute('data-testid'))).toEqual([
'agent-teams-member-task-group-running',
'agent-teams-member-task-group-upcoming',
'agent-teams-member-task-group-completed',
])
expect(ids('running')).toEqual(['agent-teams-member-task-2'])
expect(ids('upcoming')).toEqual(['agent-teams-member-task-3', 'agent-teams-member-task-4'])
expect(ids('completed')).toEqual(['agent-teams-member-task-1'])
// The task list still says in progress, but nobody is working on it.
expect(screen.getByTestId('agent-teams-member-task-2-state').textContent).toBe('Stopped')
expect(screen.getByText('Stopped · at #2')).toBeTruthy()
expect(screen.getByTestId('agent-teams-member-task-3').textContent).toContain('Depends on #2')
expect(screen.getByTestId('agent-teams-member-task-1').textContent).toContain(formatWorkbenchMessageTime('2026-08-08T07:05:00.000Z'))
expect(screen.getByTestId('agent-teams-member-task-1').textContent).not.toContain('+0:00')
expect(screen.getByText('1/4', { selector: 'dd' })).toBeTruthy()
}) })
it('shows message direction, renders human Markdown, and narrates protocol payloads', () => { it('shows message direction, renders human Markdown, and narrates protocol payloads', () => {
@@ -282,7 +334,8 @@ describe('AgentTeamsMemberInspector', () => {
/>, />,
) )
expect(screen.getByText('Stopped')).toBeTruthy() // Where it stopped, not just that it stopped.
expect(screen.getByText('Stopped · at #7')).toBeTruthy()
const notice = screen.getByTestId('agent-teams-member-recovery') const notice = screen.getByTestId('agent-teams-member-recovery')
expect(notice.getAttribute('data-member-state')).toBe('stopped') expect(notice.getAttribute('data-member-state')).toBe('stopped')
expect(screen.getByTestId('agent-teams-member-recovery-hint').textContent) expect(screen.getByTestId('agent-teams-member-recovery-hint').textContent)
@@ -12,8 +12,10 @@ import {
parseWorkbenchMessageBody, parseWorkbenchMessageBody,
resolveMemberModel, resolveMemberModel,
resolveTeamMemberIdentity, resolveTeamMemberIdentity,
stalledTaskOwnerState,
taskOwnedByMember, taskOwnedByMember,
type MemberWorkState, type MemberWorkState,
type StalledTaskOwnerState,
type WorkbenchMessageBody, type WorkbenchMessageBody,
type WorkbenchTaskState, type WorkbenchTaskState,
} from '@/components/agentTeams/agentTeamsModel' } from '@/components/agentTeams/agentTeamsModel'
@@ -183,7 +185,30 @@ function formatDuration(durationMs: number): string {
function formatTaskSpan(entry: TaskHistoryEntry): string { function formatTaskSpan(entry: TaskHistoryEntry): string {
if (entry.startedAt === null || entry.durationMs === null) return '—' if (entry.startedAt === null || entry.durationMs === null) return '—'
return `${formatWorkbenchMessageTime(new Date(entry.startedAt).toISOString())} +${formatDuration(entry.durationMs)}` const start = formatWorkbenchMessageTime(new Date(entry.startedAt).toISOString())
// Times come from polled frames: a task that started and ended between two
// of them has an unknown duration, not a zero one.
return entry.durationMs < 1000 ? start : `${start} +${formatDuration(entry.durationMs)}`
}
type TaskGroup = 'running' | 'upcoming' | 'completed'
const TASK_GROUPS: Array<{ group: TaskGroup; label: TranslationKey }> = [
{ group: 'running', label: 'agentTeams.inspector.groupRunning' },
{ group: 'upcoming', label: 'agentTeams.inspector.groupUpcoming' },
{ group: 'completed', label: 'agentTeams.inspector.groupCompleted' },
]
function taskGroup(state: WorkbenchTaskState): TaskGroup {
if (state === 'running') return 'running'
if (state === 'completed') return 'completed'
return 'upcoming'
}
function stalledTone(state: StalledTaskOwnerState): Tone {
if (state === 'error') return 'danger'
if (state === 'retrying') return 'warning'
return 'neutral'
} }
function taskTone(state: WorkbenchTaskState): Tone { function taskTone(state: WorkbenchTaskState): Tone {
@@ -193,6 +218,15 @@ function taskTone(state: WorkbenchTaskState): Tone {
return 'neutral' return 'neutral'
} }
function memberStateLabel(state: MemberWorkState, member: TeamMember, t: TranslationFn): string {
return state === 'retrying'
? t('agentTeams.member.retrying', {
attempt: member.autoRetry?.attempt ?? '?',
max: member.autoRetry?.max ?? '?',
})
: t(`agentTeams.member.${state}` as TranslationKey)
}
function memberTone(state: MemberWorkState): Tone { function memberTone(state: MemberWorkState): Tone {
if (state === 'working') return 'brand' if (state === 'working') return 'brand'
if (state === 'error') return 'danger' if (state === 'error') return 'danger'
@@ -346,12 +380,9 @@ export function AgentTeamsMemberInspector({
? t('agentTeams.member.waitingForDependency', { task: waitingDependency }) ? t('agentTeams.member.waitingForDependency', { task: waitingDependency })
: workState === 'idle' : workState === 'idle'
? t('agentTeams.member.waitingForTask') ? t('agentTeams.member.waitingForTask')
: workState === 'retrying' : (workState === 'stopped' || workState === 'error' || workState === 'retrying') && runningTask
? t('agentTeams.member.retrying', { ? t('agentTeams.member.stalledOnTask', { state: memberStateLabel(workState, member, t), task: runningTask.id })
attempt: member.autoRetry?.attempt ?? '?', : memberStateLabel(workState, member, t)
max: member.autoRetry?.max ?? '?',
})
: t(`agentTeams.member.${workState}` as TranslationKey)
// Stopped, retrying and failed members all come back on their own or through // Stopped, retrying and failed members all come back on their own or through
// a message; say why they paused and what brings them back. // a message; say why they paused and what brings them back.
const awaitsRecovery = !isLead && ( const awaitsRecovery = !isLead && (
@@ -433,7 +464,7 @@ export function AgentTeamsMemberInspector({
<dt className="text-[10px] font-semibold text-[var(--color-text-tertiary)]"> <dt className="text-[10px] font-semibold text-[var(--color-text-tertiary)]">
{t('agentTeams.inspector.completedTasks')} {t('agentTeams.inspector.completedTasks')}
</dt> </dt>
<dd className="mt-0.5 font-extrabold tabular-nums">{completedTasks}</dd> <dd className="mt-0.5 font-extrabold tabular-nums">{completedTasks}/{taskHistory.length}</dd>
</div> </div>
<div> <div>
<dt className="text-[10px] font-semibold text-[var(--color-text-tertiary)]"> <dt className="text-[10px] font-semibold text-[var(--color-text-tertiary)]">
@@ -509,40 +540,66 @@ export function AgentTeamsMemberInspector({
</span> </span>
</h3> </h3>
{taskHistory.length > 0 ? ( {taskHistory.length > 0 ? (
<ol className="mt-2" data-testid="agent-teams-member-task-history"> <div data-testid="agent-teams-member-task-history">
{taskHistory.map(entry => { {TASK_GROUPS.map(({ group, label }) => {
const span = formatTaskSpan(entry) const entries = taskHistory.filter(entry => taskGroup(entry.state) === group)
if (entries.length === 0) return null
// Up next reads in the order the work unblocks, not by when
// (never) it started.
if (group === 'upcoming') entries.sort((left, right) => left.task.id.localeCompare(right.task.id, undefined, { numeric: true }))
return ( return (
<li <section key={group} data-testid={`agent-teams-member-task-group-${group}`} className="mt-2">
key={entry.task.id} <h4 className="text-[10px] font-semibold text-[var(--color-text-secondary)]">
data-testid={`agent-teams-member-task-${entry.task.id}`} {t(label)} · {entries.length}
data-task-state={entry.state} </h4>
className="flex min-w-0 items-center gap-2 border-b border-[var(--color-border)] py-1.5 last:border-b-0" <ol className="mt-0.5">
> {entries.map(entry => {
<span className="w-[26px] shrink-0 font-mono text-[10px] font-extrabold text-[var(--color-text-tertiary)]"> const stalled = stalledTaskOwnerState(entry.task, snapshot)
#{entry.task.id} const openDependencies = entry.state === 'blocked'
</span> ? entry.task.blockedBy.filter(id => snapshot.tasks.find(task => task.id === id)?.status !== 'completed')
<span className="min-w-0 flex-1 truncate text-[11.5px] leading-[1.3]" title={entry.task.subject}> : []
{entry.task.subject} const aside = openDependencies.length > 0
</span> ? `${t('agentTeams.task.dependsOn')} ${openDependencies.map(id => `#${id}`).join(' ')}`
<Badge : formatTaskSpan(entry)
data-testid={`agent-teams-member-task-${entry.task.id}-state`} return (
tone={taskTone(entry.state)} <li
size="xs" key={entry.task.id}
bordered data-testid={`agent-teams-member-task-${entry.task.id}`}
> data-task-state={entry.state}
{t(`agentTeams.task.${entry.state}` as TranslationKey)} data-task-stalled={stalled}
</Badge> className="flex min-w-0 items-center gap-2 border-b border-[var(--color-border)] py-1.5 last:border-b-0"
<time >
dateTime={entry.startedAt === null ? undefined : new Date(entry.startedAt).toISOString()} <span className="w-[26px] shrink-0 font-mono text-[10px] font-extrabold text-[var(--color-text-tertiary)]">
className="w-[82px] shrink-0 text-right font-mono text-[9.5px] tabular-nums text-[var(--color-text-tertiary)]" #{entry.task.id}
> </span>
{span} <span className="min-w-0 flex-1 truncate text-[11.5px] leading-[1.3]" title={entry.task.subject}>
</time> {entry.task.subject}
</li> </span>
<Badge
data-testid={`agent-teams-member-task-${entry.task.id}-state`}
tone={stalled ? stalledTone(stalled) : taskTone(entry.state)}
size="xs"
bordered
>
{stalled
? memberStateLabel(stalled, member, t)
: t(`agentTeams.task.${entry.state}` as TranslationKey)}
</Badge>
<time
dateTime={entry.startedAt === null ? undefined : new Date(entry.startedAt).toISOString()}
className="w-[82px] shrink-0 truncate text-right font-mono text-[9.5px] tabular-nums text-[var(--color-text-tertiary)]"
title={aside}
>
{aside}
</time>
</li>
)
})}
</ol>
</section>
) )
})} })}
</ol> </div>
) : ( ) : (
<p className="mt-2 text-[11px] text-[var(--color-text-tertiary)]"> <p className="mt-2 text-[11px] text-[var(--color-text-tertiary)]">
{t('agentTeams.noMemberTasks')} {t('agentTeams.noMemberTasks')}
@@ -15,6 +15,7 @@ import {
runningTaskForMember, runningTaskForMember,
shortModelLabel, shortModelLabel,
snapshotWithHistoricalMembers, snapshotWithHistoricalMembers,
stalledTaskOwnerState,
taskOwnedByMember, taskOwnedByMember,
WORKBENCH_TASK_WIDTH, WORKBENCH_TASK_WIDTH,
} from './agentTeamsModel' } from './agentTeamsModel'
@@ -446,6 +447,37 @@ describe('Agent Teams workbench model', () => {
}) })
}) })
describe('Agent Teams stalled tasks', () => {
const member = (name: string, extra: Partial<TeamMember>): TeamMember => ({
agentId: `${name}@team-a`, name, role: name, status: 'idle', ...extra,
})
const withMembers = (tasks: TeamWorkbenchTask[], members: TeamMember[]): TeamWorkbenchSnapshot => {
const base = snapshot(tasks)
return { ...base, team: { ...base.team, members } }
}
it("reports the owner's state for an in-progress task its owner is no longer working on", () => {
// The task list keeps `in_progress` until the member itself updates it,
// so a stopped or failed owner would otherwise look busy on the board.
const frame = withMembers(
[task('1', 'in_progress', [], 'stopper'), task('2', 'in_progress', [], 'failer'), task('3', 'in_progress', [], 'retrier'), task('4', 'in_progress', [], 'worker')],
[
member('stopper', { activity: 'stopped' }),
member('failer', { status: 'error', activity: 'idle', lastError: 'Credit balance is too low' }),
member('retrier', { activity: 'idle', lastError: 'overloaded', autoRetry: { attempt: 1, max: 5, nextAt: 1 } }),
member('worker', { status: 'running', activity: 'active' }),
],
)
expect(frame.tasks.map(item => stalledTaskOwnerState(item, frame))).toEqual(['stopped', 'error', 'retrying', undefined])
})
it('never marks finished, unstarted or unowned work as stalled', () => {
const owner = member('stopper', { activity: 'stopped' })
const frame = withMembers([task('1', 'completed', [], 'stopper'), task('2', 'pending', [], 'stopper'), task('3', 'in_progress')], [owner])
expect(frame.tasks.map(item => stalledTaskOwnerState(item, frame))).toEqual([undefined, undefined, undefined])
})
})
describe('Agent Teams member model display', () => { describe('Agent Teams member model display', () => {
const worker: TeamMember = { const worker: TeamMember = {
agentId: 'builder@team-a', agentId: 'builder@team-a',
@@ -314,6 +314,27 @@ export function getMemberWorkState(
return member.status === 'running' ? 'working' : 'idle' return member.status === 'running' ? 'working' : 'idle'
} }
export type StalledTaskOwnerState = Extract<MemberWorkState, 'stopped' | 'error' | 'retrying'>
/**
* Why a task still marked in progress is not moving: its owner was stopped,
* failed, or waits for an automatic retry. Only the member updates its own
* task, so the list keeps saying `in_progress`, and the board would animate
* work nobody is doing.
*/
export function stalledTaskOwnerState(
task: TeamWorkbenchTask,
snapshot: TeamWorkbenchSnapshot,
): StalledTaskOwnerState | undefined {
if (task.status !== 'in_progress' || snapshot.deletedAt) return undefined
const owner = inferTaskOwner(task, snapshot)
if (!owner) return undefined
const { member, isLead } = resolveTeamMemberIdentity(snapshot.team, owner.identity)
if (isLead || !snapshot.team.members.includes(member)) return undefined
const state = getMemberWorkState(member)
return state === 'stopped' || state === 'error' || state === 'retrying' ? state : undefined
}
export type TaskOwnerAttribution = { export type TaskOwnerAttribution = {
identity: string identity: string
/** True when the name was recovered from the mailbox rather than recorded. */ /** True when the name was recovered from the mailbox rather than recorded. */
@@ -1,5 +1,7 @@
import { getComposerViewForTesting } from './MentionComposer' import { getComposerViewForTesting } from './MentionComposer'
import { useTeamPlanStore } from '@/stores/teamPlanStore' import { useTeamPlanStore } from '@/stores/teamPlanStore'
import { useTeamStore } from '@/stores/teamStore'
import type { TeamDetail, TeamMember } from '@/types/team'
import type { TeamPlanRecord } from '../../../../src/shared/teamPlan' import type { TeamPlanRecord } from '../../../../src/shared/teamPlan'
import { useSideChatStore } from '@/stores/sideChatStore' import { useSideChatStore } from '@/stores/sideChatStore'
import { fireEvent, render, screen, waitFor, within } from '@testing-library/react' import { fireEvent, render, screen, waitFor, within } from '@testing-library/react'
@@ -248,6 +250,7 @@ describe('ChatInput file mentions', () => {
vi.clearAllMocks() vi.clearAllMocks()
useSideChatStore.setState({ entries: {} }) useSideChatStore.setState({ entries: {} })
useTeamPlanStore.setState({ bySession: {} }) useTeamPlanStore.setState({ bySession: {} })
useTeamStore.setState({ workbenchesBySession: {} })
mocks.voiceSupported.mockReturnValue(false) mocks.voiceSupported.mockReturnValue(false)
useVoiceInputStore.setState({ catalog: null, loading: false, error: null }) useVoiceInputStore.setState({ catalog: null, loading: false, error: null })
mocks.sideOpen.mockResolvedValue('side-tab') mocks.sideOpen.mockResolvedValue('side-tab')
@@ -1357,6 +1360,44 @@ describe('ChatInput file mentions', () => {
expect(screen.queryByRole('button', { name: 'Stop' })).not.toBeInTheDocument() expect(screen.queryByRole('button', { name: 'Stop' })).not.toBeInTheDocument()
}) })
it('stops offering Stop once Stop has paused every member of a running team', () => {
const plan: TeamPlanRecord = {
schemaVersion: 1, planId: 'approved-plan', sessionId, teamName: 'approved-team',
incarnationId: 'incarnation', revision: 2, state: 'running', workDir: '/repo', createdAt: 1, updatedAt: 2,
leaderRuntime: { providerId: 'fake-provider', modelId: 'fake-model' }, members: [], tasks: [],
}
useTeamPlanStore.setState({ bySession: { [sessionId]: { plan, loading: false, busy: false, error: null, conflict: false } } })
const showTeam = (readerActivity: TeamMember['activity']) => {
const team: TeamDetail = {
name: 'approved-team',
leadAgentId: 'team-lead@approved-team',
leadSessionId: sessionId,
members: [
// The lead's own row never says whether the team still runs.
{ agentId: 'team-lead@approved-team', name: 'team-lead', role: 'team-lead', status: 'idle' },
{ agentId: 'reader@approved-team', name: 'reader', role: 'reader', status: 'idle', activity: readerActivity },
{ agentId: 'writer@approved-team', name: 'writer', role: 'writer', status: 'completed', activity: 'exited' },
],
}
useTeamStore.setState({
workbenchesBySession: {
[sessionId]: { teamName: team.name, loading: false, error: null, snapshots: [{ version: '1', generatedAt: '', team, tasks: [], messages: [] }] },
},
})
}
act(() => showTeam('active'))
render(<ChatInput compact />)
expect(screen.getByRole('button', { name: 'Stop' })).toBeInTheDocument()
// The plan stays running while paused, so only the members tell that
// nothing is left to stop; a Stop button here invites endless clicking.
act(() => showTeam('stopped'))
expect(screen.queryByRole('button', { name: 'Stop' })).not.toBeInTheDocument()
act(() => showTeam('idle'))
expect(screen.getByRole('button', { name: 'Stop' })).toBeInTheDocument()
})
it.each(['local_bash', 'dream'])('does not turn Run into Stop for a running %s task', (taskType) => { it.each(['local_bash', 'dream'])('does not turn Run into Stop for a running %s task', (taskType) => {
useChatStore.setState({ useChatStore.setState({
sessions: { sessions: {
+14 -1
View File
@@ -19,6 +19,7 @@ import { useUIStore } from '../../stores/uiStore'
import { useSessionStore } from '../../stores/sessionStore' import { useSessionStore } from '../../stores/sessionStore'
import { useSessionRuntimeStore } from '../../stores/sessionRuntimeStore' import { useSessionRuntimeStore } from '../../stores/sessionRuntimeStore'
import { useTeamStore } from '../../stores/teamStore' import { useTeamStore } from '../../stores/teamStore'
import { getMemberWorkState, resolveTeamMemberIdentity } from '../agentTeams/agentTeamsModel'
import { useTeamPlanStore } from '@/stores/teamPlanStore' import { useTeamPlanStore } from '@/stores/teamPlanStore'
import { useSettingsStore } from '../../stores/settingsStore' import { useSettingsStore } from '../../stores/settingsStore'
import { import {
@@ -300,10 +301,22 @@ export function ChatInput({ variant = 'default', compact = false, sessionId, vis
const hasRunningSubagents = hasRunningSubagentTasks(sessionState?.backgroundAgentTasks) const hasRunningSubagents = hasRunningSubagentTasks(sessionState?.backgroundAgentTasks)
// Approved team processes are tracked by their plan, not background-agent // Approved team processes are tracked by their plan, not background-agent
// notifications. Keep Stop available after the review card is dismissed. // notifications. Keep Stop available after the review card is dismissed.
const hasRunningTeam = useTeamPlanStore(state => { const teamPlanRunning = useTeamPlanStore(state => {
const plan = activeTabId ? state.bySession[activeTabId]?.plan : undefined const plan = activeTabId ? state.bySession[activeTabId]?.plan : undefined
return plan?.state === 'launching' || plan?.state === 'running' return plan?.state === 'launching' || plan?.state === 'running'
}) })
// Stop pauses a running team without ending its plan, so the plan alone
// would keep offering Stop after every member is already stopped.
const teamMembersAllStopped = useTeamStore(state => {
const team = activeTabId ? state.workbenchesBySession[activeTabId]?.snapshots.at(-1)?.team : undefined
if (!team) return false
const members = team.members.filter(member => !resolveTeamMemberIdentity(team, member.agentId).isLead)
return members.length > 0 && members.every(member => {
const work = getMemberWorkState(member)
return work === 'stopped' || work === 'exited'
})
})
const hasRunningTeam = teamPlanRunning && !teamMembersAllStopped
const workspaceState = getSessionWorkspaceState(activeSession) const workspaceState = getSessionWorkspaceState(activeSession)
const isWorkspaceMissing = workspaceState !== 'available' const isWorkspaceMissing = workspaceState !== 'available'
// Both composer branches (hero and inline) and the drop handler share this: // Both composer branches (hero and inline) and the drop handler share this:
+4
View File
@@ -3343,6 +3343,10 @@ Row 9, all 8 cells: continuing from straight down, turning left through lower-le
'agentTeams.member.retryNow': 'Retrying automatically now', 'agentTeams.member.retryNow': 'Retrying automatically now',
'agentTeams.member.stoppedHint': 'Send a message to resume it from its saved conversation', 'agentTeams.member.stoppedHint': 'Send a message to resume it from its saved conversation',
'agentTeams.member.errorHint': 'Once the problem is fixed, send a message to let it continue', 'agentTeams.member.errorHint': 'Once the problem is fixed, send a message to let it continue',
'agentTeams.member.stalledOnTask': '{state} · at #{task}',
'agentTeams.inspector.groupRunning': 'Now',
'agentTeams.inspector.groupUpcoming': 'Up next',
'agentTeams.inspector.groupCompleted': 'Done',
'agentTeams.openMember': 'Open {name} details', 'agentTeams.openMember': 'Open {name} details',
'agentTeams.resizeCommunication': 'Resize communication panel', 'agentTeams.resizeCommunication': 'Resize communication panel',
'agentTeams.memberTranscriptLoading': 'Loading member transcript...', 'agentTeams.memberTranscriptLoading': 'Loading member transcript...',
+4
View File
@@ -3344,6 +3344,10 @@ export const jp: Record<TranslationKey, string> = {
'agentTeams.member.retryNow': 'まもなく自動で再試行', 'agentTeams.member.retryNow': 'まもなく自動で再試行',
'agentTeams.member.stoppedHint': 'メッセージを送ると、保存された会話から再開します', 'agentTeams.member.stoppedHint': 'メッセージを送ると、保存された会話から再開します',
'agentTeams.member.errorHint': '問題を解決してからメッセージを送ると、作業を再開します', 'agentTeams.member.errorHint': '問題を解決してからメッセージを送ると、作業を再開します',
'agentTeams.member.stalledOnTask': '{state} · #{task} で停止',
'agentTeams.inspector.groupRunning': '進行中',
'agentTeams.inspector.groupUpcoming': '次のタスク',
'agentTeams.inspector.groupCompleted': '完了',
'agentTeams.openMember': '{name} の詳細を開く', 'agentTeams.openMember': '{name} の詳細を開く',
'agentTeams.resizeCommunication': '通信パネルの幅を変更', 'agentTeams.resizeCommunication': '通信パネルの幅を変更',
'agentTeams.memberTranscriptLoading': 'メンバーの transcript を読み込み中...', 'agentTeams.memberTranscriptLoading': 'メンバーの transcript を読み込み中...',
+4
View File
@@ -3346,6 +3346,10 @@ export const kr: Record<TranslationKey, string> = {
'agentTeams.member.retryNow': '곧 자동 재시도', 'agentTeams.member.retryNow': '곧 자동 재시도',
'agentTeams.member.stoppedHint': '메시지를 보내면 저장된 대화에서 다시 시작합니다', 'agentTeams.member.stoppedHint': '메시지를 보내면 저장된 대화에서 다시 시작합니다',
'agentTeams.member.errorHint': '문제를 해결한 뒤 메시지를 보내면 작업을 이어갑니다', 'agentTeams.member.errorHint': '문제를 해결한 뒤 메시지를 보내면 작업을 이어갑니다',
'agentTeams.member.stalledOnTask': '{state} · #{task}에서 멈춤',
'agentTeams.inspector.groupRunning': '진행 중',
'agentTeams.inspector.groupUpcoming': '다음 작업',
'agentTeams.inspector.groupCompleted': '완료',
'agentTeams.openMember': '{name} 세부 정보 열기', 'agentTeams.openMember': '{name} 세부 정보 열기',
'agentTeams.resizeCommunication': '통신 패널 너비 조정', 'agentTeams.resizeCommunication': '통신 패널 너비 조정',
'agentTeams.memberTranscriptLoading': '멤버 transcript 불러오는 중...', 'agentTeams.memberTranscriptLoading': '멤버 transcript 불러오는 중...',
+4
View File
@@ -3343,6 +3343,10 @@ export const zh: Record<TranslationKey, string> = {
'agentTeams.member.retryNow': '即將自動重試', 'agentTeams.member.retryNow': '即將自動重試',
'agentTeams.member.stoppedHint': '發訊息即可從已儲存的對話恢復', 'agentTeams.member.stoppedHint': '發訊息即可從已儲存的對話恢復',
'agentTeams.member.errorHint': '解決問題後,發訊息即可讓它繼續', 'agentTeams.member.errorHint': '解決問題後,發訊息即可讓它繼續',
'agentTeams.member.stalledOnTask': '{state} · 停在 #{task}',
'agentTeams.inspector.groupRunning': '進行中',
'agentTeams.inspector.groupUpcoming': '接下來',
'agentTeams.inspector.groupCompleted': '已完成',
'agentTeams.openMember': '開啟 {name} 的詳細資料', 'agentTeams.openMember': '開啟 {name} 的詳細資料',
'agentTeams.resizeCommunication': '調整通訊面板寬度', 'agentTeams.resizeCommunication': '調整通訊面板寬度',
'agentTeams.memberTranscriptLoading': '正在載入成員 transcript...', 'agentTeams.memberTranscriptLoading': '正在載入成員 transcript...',
+4
View File
@@ -3342,6 +3342,10 @@ export const zh: Record<TranslationKey, string> = {
'agentTeams.member.retryNow': '即将自动重试', 'agentTeams.member.retryNow': '即将自动重试',
'agentTeams.member.stoppedHint': '发消息即可从保存的对话恢复', 'agentTeams.member.stoppedHint': '发消息即可从保存的对话恢复',
'agentTeams.member.errorHint': '解决问题后,发消息即可让它继续', 'agentTeams.member.errorHint': '解决问题后,发消息即可让它继续',
'agentTeams.member.stalledOnTask': '{state} · 停在 #{task}',
'agentTeams.inspector.groupRunning': '进行中',
'agentTeams.inspector.groupUpcoming': '接下来',
'agentTeams.inspector.groupCompleted': '已完成',
'agentTeams.openMember': '打开 {name} 的详情', 'agentTeams.openMember': '打开 {name} 的详情',
'agentTeams.resizeCommunication': '调整通讯面板宽度', 'agentTeams.resizeCommunication': '调整通讯面板宽度',
'agentTeams.memberTranscriptLoading': '正在加载成员 transcript...', 'agentTeams.memberTranscriptLoading': '正在加载成员 transcript...',
+1 -1
View File
@@ -107,7 +107,7 @@ Claude 每完成一轮并修改文件,对话里会出现一张「{n} 个文件
后台跑的子 Agent 的工具活动也会冒泡到这里,不用等它跑完才知道它在干什么。 后台跑的子 Agent 的工具活动也会冒泡到这里,不用等它跑完才知道它在干什么。
团队成员遇到模型服务断流、限流、5xx 这类临时故障会自己重试,成员行显示「自动重试 2/5」,不用你或主控介入;重试用完,或者遇到需要你处理的错误(比如 API Key 失效),成员行显示「出错」并通知主控。按停止按钮会让整组成员一起停下,但进度不会丢:你下一次给主控发消息时,它会知道哪些成员停在了哪个任务上,再按你的意思决定要不要让它们继续;你也可以直接给某个成员发消息,它会从保存的对话接着干,显示「已停止」的成员也一样。切换模型或权限模式不会打断团队;重启应用后团队还在,给成员发消息即可继续。删除会话或执行 `/clear` 才会结束团队。 团队成员遇到模型服务断流、限流、5xx 这类临时故障会自己重试,成员行显示「自动重试 2/5」,不用你或主控介入;重试用完,或者遇到需要你处理的错误(比如 API Key 失效、余额不足、订阅额度用完),成员行显示「出错」并通知主控。问题解决后给主控发一句「继续」就行:它会收到出错成员的清单,逐个唤醒它们从保存的对话接着干,不会重做已完成的部分。按停止按钮会让整组成员一起停下,主控也不会再自己处理成员发来的汇报,但进度不会丢:你下一次给主控发消息时,它会知道哪些成员停在了哪个任务上,再按你的意思决定要不要让它们继续;你也可以直接给某个成员发消息,它会从保存的对话接着干,显示「已停止」的成员也一样。切换模型或权限模式不会打断团队;重启应用后团队还在,给成员发消息即可继续。删除会话或执行 `/clear` 才会结束团队。
## 轨迹:每一步到底发生了什么 ## 轨迹:每一步到底发生了什么
+1 -1
View File
@@ -107,7 +107,7 @@ The first button on the right of the tab bar opens the Activity panel, which lis
Tool activity from background subagents bubbles up here too, so you don't have to wait for one to finish to see what it's doing. Tool activity from background subagents bubbles up here too, so you don't have to wait for one to finish to see what it's doing.
Team members retry on their own when the model service drops a stream, rate-limits, or returns a 5xx; the member row shows "Auto-retry 2/5" and neither you nor the lead has to step in. When the retries run out, or the error needs you (an expired API key, say), the row shows "Error" and the lead is told. The Stop button halts the whole team without losing progress: your next message to the lead tells it which members stopped on which tasks, and it decides from your words whether they carry on. You can also message a member directly; it picks up from its saved conversation, and the same goes for a member marked "Stopped". Switching the model or permission mode doesn't interrupt the team, and after an app restart the team is still there — message a member to continue. Deleting the session or running `/clear` ends the team. Team members retry on their own when the model service drops a stream, rate-limits, or returns a 5xx; the member row shows "Auto-retry 2/5" and neither you nor the lead has to step in. When the retries run out, or the error needs you (an expired API key, an empty balance, a used-up subscription limit), the row shows "Error" and the lead is told. Once the problem is fixed, tell the lead to continue: it gets the list of members that stopped on an error and wakes each one, which picks up from its saved conversation without redoing finished work. The Stop button halts the whole team, and the lead stops acting on its members' reports, without losing progress: your next message to the lead tells it which members stopped on which tasks, and it decides from your words whether they carry on. You can also message a member directly; it picks up from its saved conversation, and the same goes for a member marked "Stopped". Switching the model or permission mode doesn't interrupt the team, and after an app restart the team is still there — message a member to continue. Deleting the session or running `/clear` ends the team.
## Trajectory: what actually happened, step by step ## Trajectory: what actually happened, step by step
+179
View File
@@ -0,0 +1,179 @@
import { afterAll, beforeEach, expect, mock, test } from 'bun:test'
import { mkdtempSync, rmSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { createSandboxedTestEnvironment } from '../../scripts/pr/test-environment.js'
const home = mkdtempSync(join(tmpdir(), 'lead-mailbox-pause-'))
const originalEnv = { ...process.env }
for (const key of Object.keys(process.env)) delete process.env[key]
Object.assign(
process.env,
createSandboxedTestEnvironment(
home,
{ CLAUDE_CODE_SIMPLE: '1', CC_HAHA_AGENT_TEAMS_ENABLED: '1' },
originalEnv,
),
)
const LEAD_ID = 'team-lead@pause-team'
const WORKER_ID = 'worker@pause-team'
// Each lead turn records its prompt instead of calling a model.
const prompts: string[] = []
const queryEngine = await import('../QueryEngine.js')
const originalQueryEngine = { ...queryEngine }
mock.module('../QueryEngine.js', () => ({
...queryEngine,
ask: async function* (params: { prompt: unknown }) {
prompts.push(
typeof params.prompt === 'string'
? params.prompt
: JSON.stringify(params.prompt),
)
},
}))
const { __runHeadlessStreamingForTests } = await import('./print.js')
const { StructuredIO } = await import('./structuredIO.js')
const { Stream } = await import('../utils/stream.js')
const { getDefaultAppState } = await import('../state/AppStateStore.js')
const { clearCommandQueue } = await import('../utils/messageQueueManager.js')
const { setIsInteractive, getIsInteractive } = await import(
'../bootstrap/state.js'
)
const mailbox = await import('../utils/teammateMailbox.js')
const wasInteractive = getIsInteractive()
setIsInteractive(false)
afterAll(() => {
// mock.restore does not undo mock.module; later files in the same run
// must see the real QueryEngine.
mock.module('../QueryEngine.js', () => originalQueryEngine)
clearCommandQueue()
setIsInteractive(wasInteractive)
for (const key of Object.keys(process.env)) delete process.env[key]
Object.assign(process.env, originalEnv)
rmSync(home, { recursive: true, force: true })
})
beforeEach(() => {
prompts.length = 0
clearCommandQueue()
})
function startLead(team: string) {
const defaults = getDefaultAppState()
let state = {
...defaults,
teamContext: {
teamName: team,
teamFilePath: join(home, 'teams', team, 'config.json'),
leadAgentId: LEAD_ID,
teammates: {
[LEAD_ID]: { name: 'team-lead' },
[WORKER_ID]: { name: 'worker' },
},
},
} as unknown as typeof defaults
const input = new Stream<string>()
const responses = new Set<string>()
const output = __runHeadlessStreamingForTests(
new StructuredIO(input),
[],
[],
[],
[],
(() => undefined) as never,
{},
() => state,
update => {
state = update(state)
},
[],
{ outputFormat: 'stream-json' } as never,
)
const drained = (async () => {
for await (const message of output) {
const response = (message as { type?: string; response?: { request_id?: string } })
if (response.type === 'control_response' && response.response?.request_id) {
responses.add(response.response.request_id)
}
}
})()
const send = (message: unknown) => input.enqueue(`${JSON.stringify(message)}\n`)
return {
say: (content: string) => send({ type: 'user', message: { role: 'user', content } }),
async control(subtype: string) {
const id = `${subtype}-${Math.random()}`
send({ type: 'control_request', request_id: id, request: { subtype } })
await until(() => responses.has(id))
},
async close() {
// Without teammates the closing input ends the session instead of
// asking the lead to shut a team down.
state = { ...state, teamContext: { ...state.teamContext!, teammates: { [LEAD_ID]: { name: 'team-lead' } } } } as typeof state
input.done()
await drained
clearCommandQueue()
},
}
}
async function until(done: () => boolean, ms = 10_000): Promise<void> {
const deadline = Date.now() + ms
while (!done() && Date.now() < deadline) await Bun.sleep(20)
expect(done()).toBe(true)
}
async function report(team: string, text: string): Promise<void> {
expect(
await mailbox.writeToMailbox(
'team-lead',
{ from: 'worker', text, timestamp: new Date().toISOString() },
team,
),
).toBe(true)
}
const delivered = (text: string) =>
`<teammate-message teammate_id="worker">\n${text}\n</teammate-message>`
test('a team pause keeps teammate mail from starting lead turns until the user speaks again', async () => {
const team = 'pause-team'
const lead = startLead(team)
try {
lead.say('start')
await until(() => prompts.length === 1)
await lead.control('team_plan_pause')
// A report a member wrote before the Stop killed it.
await report(team, 'halfway done')
// Several poll intervals: the lead must not act on its own.
await Bun.sleep(1_500)
expect(prompts).toEqual(['start'])
expect((await mailbox.readUnreadMessages('team-lead', team)).map(m => m.text)).toEqual(['halfway done'])
lead.say('continue')
await until(() => prompts.length === 3)
expect(prompts).toEqual(['start', 'continue', delivered('halfway done')])
} finally {
await lead.close()
}
}, 30_000)
test('a plain interrupt keeps delivering teammate mail between turns', async () => {
const team = 'interrupt-team'
const lead = startLead(team)
try {
lead.say('start')
await until(() => prompts.length === 1)
await lead.control('interrupt')
await report(team, 'done')
await until(() => prompts.length === 2)
expect(prompts).toEqual(['start', delivered('done')])
} finally {
await lead.close()
}
}, 30_000)
+19 -1
View File
@@ -437,14 +437,20 @@ export type LeadMailboxPollState = {
finalDrainDone: boolean finalDrainDone: boolean
/** Consecutive polls whose batch could not be acknowledged. */ /** Consecutive polls whose batch could not be acknowledged. */
consecutiveAckFailures: number consecutiveAckFailures: number
/**
* The host paused the team (the desktop's Stop): teammate mail waits unread
* for the user's next message instead of starting lead turns nobody asked for.
*/
held: boolean
} }
export function createLeadMailboxPollState(): LeadMailboxPollState { export function createLeadMailboxPollState(): LeadMailboxPollState {
return { finalDrainDone: false, consecutiveAckFailures: 0 } return { finalDrainDone: false, consecutiveAckFailures: 0, held: false }
} }
export type LeadMailboxPollStep = export type LeadMailboxPollStep =
| { kind: 'stop' } | { kind: 'stop' }
| { kind: 'held' }
| { kind: 'idle' } | { kind: 'idle' }
| { kind: 'retry' } | { kind: 'retry' }
| { kind: 'batch'; messages: TeammateMessage[] } | { kind: 'batch'; messages: TeammateMessage[] }
@@ -461,6 +467,7 @@ export async function takeLeadMailboxBatch(
options: { teamName: string | undefined; hasActiveTeammates: boolean }, options: { teamName: string | undefined; hasActiveTeammates: boolean },
): Promise<LeadMailboxPollStep> { ): Promise<LeadMailboxPollStep> {
const { teamName, hasActiveTeammates } = options const { teamName, hasActiveTeammates } = options
if (state.held) return { kind: 'held' }
const unread = await readUnreadMessages(TEAM_LEAD_NAME, teamName) const unread = await readUnreadMessages(TEAM_LEAD_NAME, teamName)
if (!hasActiveTeammates) { if (!hasActiveTeammates) {
@@ -2693,6 +2700,13 @@ function runHeadlessStreaming(
break break
} }
if (step.kind === 'held') {
logForDebugging(
'[print.ts] Team paused by the host; teammate mail waits for the next user message',
)
break
}
if (step.kind === 'retry') { if (step.kind === 'retry') {
await sleep(POLL_INTERVAL_MS) await sleep(POLL_INTERVAL_MS)
continue continue
@@ -3020,6 +3034,7 @@ function runHeadlessStreaming(
} }
} else if (message.request.subtype === 'interrupt' || message.request.subtype === 'team_plan_pause') { } else if (message.request.subtype === 'interrupt' || message.request.subtype === 'team_plan_pause') {
sessionMessageInbox.cancelQueued(dequeueAllMatching) sessionMessageInbox.cancelQueued(dequeueAllMatching)
if (message.request.subtype === 'team_plan_pause') leadMailboxPoll.held = true
// Track escapes for attribution (ant-only feature) // Track escapes for attribution (ant-only feature)
if (feature('COMMIT_ATTRIBUTION')) { if (feature('COMMIT_ATTRIBUTION')) {
setAppState(prev => ({ setAppState(prev => ({
@@ -4452,6 +4467,9 @@ function runHeadlessStreaming(
trackReceivedMessageUuid(message.uuid) trackReceivedMessageUuid(message.uuid)
} }
// The user speaking again ends a host pause: mail held since then reaches
// the lead after this turn, alongside whatever the user now wants.
leadMailboxPoll.held = false
enqueue({ enqueue({
mode: 'prompt' as const, mode: 'prompt' as const,
// file_attachments rides the protobuf catchall from the web composer. // file_attachments rides the protobuf catchall from the web composer.
@@ -205,6 +205,23 @@ test('stopping a parent stops its workers and cancels relayed pending approvals'
expect(events).toContainEqual({ type: 'control_cancel_request', request_id: 'approval' }) expect(events).toContainEqual({ type: 'control_cancel_request', request_id: 'approval' })
}) })
test("Stop asks a lead whose team was paused to hold its teammates' mail, and interrupts any other session plainly", () => {
const service = new ConversationService() as any
const requests: Array<[string, string]> = []
service.sendSdkMessage = (id: string, message: any) => {
requests.push([id, message.request.subtype])
return true
}
const remove = service.addTeamRuntimeListener({ leadInterrupted: (id: string) => id === 'lead' })
try {
service.sendInterrupt('lead')
service.sendInterrupt('solo')
} finally {
remove()
}
expect(requests).toEqual([['lead', 'team_plan_pause'], ['solo', 'interrupt']])
})
test('a lead process that exits on its own or restarts keeps its approved workers', async () => { test('a lead process that exits on its own or restarts keeps its approved workers', async () => {
const service = new ConversationService() as any const service = new ConversationService() as any
const killed: string[] = [] const killed: string[] = []
+8 -4
View File
@@ -335,8 +335,8 @@ export type TeamWorkerStart = {
* own: a lead restart keeps them, and Stop lets the runtime pause them first. * own: a lead restart keeps them, and Stop lets the runtime pause them first.
*/ */
export type TeamRuntimeListener = { export type TeamRuntimeListener = {
/** Synchronously before Stop kills a lead's workers. */ /** Synchronously before Stop kills a lead's workers; true when it paused a running team. */
leadInterrupted?: (parentSessionId: string) => void leadInterrupted?: (parentSessionId: string) => boolean | void
/** After any CLI session finished starting. */ /** After any CLI session finished starting. */
sessionStarted?: (sessionId: string, info: { isTeamWorker: boolean }) => void sessionStarted?: (sessionId: string, info: { isTeamWorker: boolean }) => void
} }
@@ -1192,9 +1192,10 @@ export class ConversationService {
// already working is paused rather than destroyed: the runtime marks its // already working is paused rather than destroyed: the runtime marks its
// members as user-stopped before their processes die, and a later message // members as user-stopped before their processes die, and a later message
// from the lead or the user resumes each member from its own transcript. // from the lead or the user resumes each member from its own transcript.
let pausedTeam = false
for (const listener of this.teamRuntimeListeners) { for (const listener of this.teamRuntimeListeners) {
try { try {
listener.leadInterrupted?.(sessionId) if (listener.leadInterrupted?.(sessionId) === true) pausedTeam = true
} catch (error) { } catch (error) {
console.error('[ConversationService] Team runtime interrupt hook failed', error) console.error('[ConversationService] Team runtime interrupt hook failed', error)
} }
@@ -1211,10 +1212,13 @@ export class ConversationService {
if (this.teamStopOperations.get(sessionId) === stop) this.teamStopOperations.delete(sessionId) if (this.teamStopOperations.get(sessionId) === stop) this.teamStopOperations.delete(sessionId)
}) })
this.teamStopOperations.set(sessionId, stop) this.teamStopOperations.set(sessionId, stop)
// A paused team's lead also holds the mail its members sent before the
// Stop until the user speaks again; acting on it would start lead turns,
// and restart members, that the user just stopped.
return this.sendSdkMessage(sessionId, { return this.sendSdkMessage(sessionId, {
type: 'control_request', type: 'control_request',
request_id: crypto.randomUUID(), request_id: crypto.randomUUID(),
request: { subtype: 'interrupt' }, request: { subtype: pausedTeam ? 'team_plan_pause' : 'interrupt' },
}) })
} }
@@ -131,6 +131,27 @@ async function startHarness(memberNames: string[]) {
} }
} }
/**
* The user sends the lead a message. The lead here is a fixture whose turns
* would call the fake model too, so what the runtime tells it is recorded
* instead of delivered.
*/
async function userMessagesLead(h: Harness, toLead: string[] = []) {
const { noteLeadUserMessage } = await import('./teamPlanRuntime.js')
const send = h.service.sendMessage.bind(h.service)
const spy = spyOn(h.service, 'sendMessage').mockImplementation(async (...args: Parameters<typeof send>) => {
if (args[0] !== h.parentId) return send(...args)
toLead.push(String(args[1]))
return true
})
try {
await noteLeadUserMessage(h.parentId)
} finally {
spy.mockRestore()
}
return toLead
}
async function messageMember(h: Harness, to: string, from: string, text: string) { async function messageMember(h: Harness, to: string, from: string, text: string) {
const { writeToMailbox } = await import('../../utils/teammateMailbox.js') const { writeToMailbox } = await import('../../utils/teammateMailbox.js')
await writeToMailbox(to, { from, text, timestamp: new Date().toISOString() }, h.teamName) await writeToMailbox(to, { from, text, timestamp: new Date().toISOString() }, h.teamName)
@@ -273,6 +294,14 @@ test("the user's Stop pauses the team instead of ending it, and only an instruct
expect(h.requests.length).toBe(before) expect(h.requests.length).toBe(before)
expect(h.service.hasSession(h.memberIds.writer!)).toBe(false) expect(h.service.hasSession(h.memberIds.writer!)).toBe(false)
// A lead message the user did not ask for (one the lead was sending as the
// user pressed Stop, or a turn it started on its own) must not undo the Stop.
await messageMember(h, 'writer', 'team-lead', 'Lead message from before the user spoke again')
await new Promise(resolve => setTimeout(resolve, 700))
expect(h.requests.length).toBe(before)
expect(h.service.hasSession(h.memberIds.writer!)).toBe(false)
await userMessagesLead(h)
await messageMember(h, 'writer', 'team-lead', 'The user wants you to continue') await messageMember(h, 'writer', 'team-lead', 'The user wants you to continue')
await h.waitFor(() => h.requests.length > before, 'resumed member turn') await h.waitFor(() => h.requests.length > before, 'resumed member turn')
expect(h.requests.at(-1)?.resumed).toBe(h.memberIds.writer) expect(h.requests.at(-1)?.resumed).toBe(h.memberIds.writer)
@@ -283,33 +312,46 @@ test("the user's Stop pauses the team instead of ending it, and only an instruct
} }
}, 30_000) }, 30_000)
test("after a Stop, the lead's next message tells it which members stopped and how to resume them", async () => { test("a user's direct message resumes a stopped member without waiting for the lead", async () => {
const h = await startHarness(['reader', 'writer']) const h = await startHarness(['reader'])
try { try {
const { deliverTeamPauseNotice } = await import('./teamPlanRuntime.js') await h.waitFor(async () => (await h.leadNotifications()).length === 1, 'initial turn report')
await h.waitFor(async () => (await h.leadNotifications()).length === 2, 'initial turn reports') h.service.sendInterrupt(h.parentId)
const sent: Array<[string, string]> = [] await h.service.waitForTeamWorkersStopped(h.parentId)
const send = h.service.sendMessage.bind(h.service) const before = h.requests.length
const spy = spyOn(h.service, 'sendMessage').mockImplementation(async (...args: Parameters<typeof send>) => { await messageMember(h, 'reader', 'user', 'Keep going on your own')
sent.push([args[0], String(args[1])]) await h.waitFor(() => h.requests.length > before, 'resumed member turn')
return send(...args) expect(h.requests.at(-1)?.resumed).toBe(h.memberIds.reader)
expect(h.requests.at(-1)?.prompt).toContain('Keep going on your own')
} finally {
await h.cleanup()
}
}, 30_000)
test('a Stop while a member is restarting leaves that member stopped', async () => {
const h = await startHarness(['reader'])
try {
await h.waitFor(async () => (await h.leadNotifications()).length === 1, 'initial turn report')
const sessionId = h.memberIds.reader!
await h.service.stopSessionAndWait(sessionId)
const start = h.service.startSession.bind(h.service)
let stopped = false
const spy = spyOn(h.service, 'startSession').mockImplementation(async (...args: Parameters<typeof start>) => {
// The user presses Stop after the restart began but before its process
// exists, so Stop's kill pass cannot see it.
if (args[0] === sessionId && !stopped) {
stopped = true
h.service.sendInterrupt(h.parentId)
}
return start(...args)
}) })
try { try {
// Without a Stop there is nothing to tell. await messageMember(h, 'reader', 'team-lead', 'Check one more file')
await deliverTeamPauseNotice(h.parentId) await h.waitFor(() => stopped, 'restart to begin')
expect(sent).toEqual([])
h.service.sendInterrupt(h.parentId)
await h.service.waitForTeamWorkersStopped(h.parentId) await h.service.waitForTeamWorkersStopped(h.parentId)
await deliverTeamPauseNotice(h.parentId) await new Promise(resolve => setTimeout(resolve, 800))
await deliverTeamPauseNotice(h.parentId) expect(h.service.hasSession(sessionId)).toBe(false)
expect((await h.teamMember('reader'))?.terminated).toBe(true)
const notices = sent.filter(([sessionId]) => sessionId === h.parentId).map(([, text]) => text)
expect(notices).toHaveLength(1)
expect(notices[0]).toContain('[Team runtime notice]')
expect(notices[0]).toMatch(/- reader: #\d+ Task for reader \(pending\)/)
expect(notices[0]).toMatch(/- writer: #\d+ Task for writer \(pending\)/)
expect(notices[0]).toContain('SendMessage each member that still has work')
} finally { } finally {
spy.mockRestore() spy.mockRestore()
} }
@@ -318,6 +360,87 @@ test("after a Stop, the lead's next message tells it which members stopped and h
} }
}, 30_000) }, 30_000)
test("after a Stop, the lead's next message tells it which members stopped and how to resume them", async () => {
const h = await startHarness(['reader', 'writer'])
try {
await h.waitFor(async () => (await h.leadNotifications()).length === 2, 'initial turn reports')
// Without a Stop or a failure there is nothing to tell.
expect(await userMessagesLead(h)).toEqual([])
h.service.sendInterrupt(h.parentId)
await h.service.waitForTeamWorkersStopped(h.parentId)
const notices = await userMessagesLead(h)
await userMessagesLead(h, notices)
expect(notices).toHaveLength(1)
expect(notices[0]).toContain("[Team runtime notice] The user's Stop also stopped your team")
expect(notices[0]).toMatch(/- reader: #\d+ Task for reader \(pending\)/)
expect(notices[0]).toMatch(/- writer: #\d+ Task for writer \(pending\)/)
expect(notices[0]).toContain('SendMessage each member that still has work')
} finally {
await h.cleanup()
}
}, 30_000)
test("after any failure, the user's next message tells the lead which members to wake, once per failure", async () => {
const { setTeamRuntimeTimingForTests } = await import('./teamPlanRuntime.js')
setTeamRuntimeTimingForTests({ autoContinueDelaysMs: [40] })
const h = await startHarness(['billing', 'network', 'crash', 'healthy'])
try {
await h.waitFor(async () => (await h.leadNotifications()).length === 4, 'initial turn reports')
const notices: string[] = []
const failureCount = async () => (await h.leadNotifications()).filter(item => item.idleReason === 'failed').length
// Model calls are answered in arrival order, so members fail one at a
// time: billing is final, a dropped connection exhausts its one retry,
// and a crash ends the process.
for (const [name, replies] of [
['billing', ['FIXTURE_ERROR:Credit balance is too low']],
['network', ['FIXTURE_ERROR:API Error: Connection error.', 'FIXTURE_ERROR:API Error: Connection error.']],
['crash', ['FIXTURE_CRASH']],
] as const) {
const count = await failureCount()
h.replies.push(...replies)
await messageMember(h, name, 'team-lead', 'Keep working')
await h.waitFor(async () => await failureCount() > count, `${name} failure`)
}
// No Stop happened, yet the user's next message names every member
// that stopped on an error.
await userMessagesLead(h, notices)
expect(notices).toHaveLength(1)
const notice = notices[0]!
expect(notice).toContain('[Team runtime notice]')
expect(notice).toMatch(/- billing \(stopped on an error: Credit balance is too low\): #\d+ Task for billing \(pending\)/)
expect(notice).toMatch(/- network \(stopped on an error: [^)]*Connection error[^\n]*\): #\d+ Task for network/)
expect(notice).toMatch(/- crash \(stopped on an error: [^\n]*exited[^\n]*\): #\d+ Task for crash/)
expect(notice).not.toContain('healthy')
expect(notice).toContain('SendMessage each member that still has work')
// Each failure is announced once.
await userMessagesLead(h, notices)
expect(notices).toHaveLength(1)
// The lead wakes a member: it continues in its own conversation.
const before = h.requests.length
await messageMember(h, 'billing', 'team-lead', 'Billing is fixed, continue')
await h.waitFor(() => h.requests.length > before, 'woken member turn')
expect(h.requests.at(-1)?.prompt).toContain('Billing is fixed, continue')
await h.waitFor(async () => (await h.teamMember('billing'))?.isActive === false, 'woken turn to finish')
// A new failure is announced again, alone.
const count = await failureCount()
h.replies.push('FIXTURE_ERROR:Credit balance is too low')
await messageMember(h, 'billing', 'team-lead', 'One more thing')
await h.waitFor(async () => await failureCount() > count, 'second billing failure')
await userMessagesLead(h, notices)
expect(notices).toHaveLength(2)
expect(notices[1]).toContain('- billing (stopped on an error')
expect(notices[1]).not.toContain('- network')
} finally {
await h.cleanup()
}
}, 60_000)
test('stopping and resuming the team repeatedly is never mistaken for a crash loop', async () => { test('stopping and resuming the team repeatedly is never mistaken for a crash loop', async () => {
const h = await startHarness(['reader']) const h = await startHarness(['reader'])
try { try {
@@ -328,6 +451,7 @@ test('stopping and resuming the team repeatedly is never mistaken for a crash lo
await h.service.waitForTeamWorkersStopped(h.parentId) await h.service.waitForTeamWorkersStopped(h.parentId)
await h.waitFor(async () => (await h.teamMember('reader'))?.terminated === true, `stop ${cycle} to be recorded`) await h.waitFor(async () => (await h.teamMember('reader'))?.terminated === true, `stop ${cycle} to be recorded`)
const before = h.requests.length const before = h.requests.length
await userMessagesLead(h)
await messageMember(h, 'reader', 'team-lead', `Continue after stop ${cycle}`) await messageMember(h, 'reader', 'team-lead', `Continue after stop ${cycle}`)
await h.waitFor(() => h.requests.length > before, `resume ${cycle}`) await h.waitFor(() => h.requests.length > before, `resume ${cycle}`)
expect(h.requests.at(-1)?.prompt).toContain(`Continue after stop ${cycle}`) expect(h.requests.at(-1)?.prompt).toContain(`Continue after stop ${cycle}`)
+75 -32
View File
@@ -52,6 +52,10 @@ type WorkerRuntime = {
restartsExhausted: boolean restartsExhausted: boolean
/** The user's Stop ended this process; resuming it is not a crash restart. */ /** The user's Stop ended this process; resuming it is not a crash restart. */
stoppedByUser?: boolean stoppedByUser?: boolean
/** Why the member's last turn ended in a failure nothing retries any more. */
failure?: string
/** The lead has not yet been told, with the user's next message, that this member waits to be woken. */
wakeNoticePending?: boolean
wokenForTaskIds: Set<string> wokenForTaskIds: Set<string>
lastResultAt?: number lastResultAt?: number
} }
@@ -70,8 +74,12 @@ type TeamLaunch = {
stopped: boolean stopped: boolean
/** Set by the user's Stop: hold automatic work until the lead or the user acts again. */ /** Set by the user's Stop: hold automatic work until the lead or the user acts again. */
pausedAt?: number pausedAt?: number
/** The lead has not yet been told what the user's last Stop did to this team. */ /**
pauseNoticePending?: boolean * When the user next messaged the lead after a Stop. Only lead instructions
* written since then resume the team: earlier ones were sent before the
* user's Stop took effect, or in turns the user never asked for.
*/
resumeAllowedAt?: number
permissionMode: string permissionMode: string
} }
@@ -281,7 +289,7 @@ async function restartWorker(launch: TeamLaunch, worker: WorkerRuntime, reason:
if (!worker.restartsExhausted) { if (!worker.restartsExhausted) {
worker.restartsExhausted = true worker.restartsExhausted = true
await updateWorkerEntry(launch, worker, entry => ({ ...entry, isActive: false, terminated: true, lastError: 'Stopped restarting after repeated exits. A message from the user restarts it.' })) await updateWorkerEntry(launch, worker, entry => ({ ...entry, isActive: false, terminated: true, lastError: 'Stopped restarting after repeated exits. A message from the user restarts it.' }))
await notifyLead(launch, worker, { idleReason: 'failed', failureReason: `${worker.member.name} exited ${worker.restartHistory.length} times in ${WORKER_RESTART_WINDOW_MS / 60_000} minutes; automatic restarts are paused` }) await notifyFailure(launch, worker, `${worker.member.name} exited ${worker.restartHistory.length} times in ${WORKER_RESTART_WINDOW_MS / 60_000} minutes; automatic restarts are paused`)
} }
return false return false
} }
@@ -290,7 +298,9 @@ async function restartWorker(launch: TeamLaunch, worker: WorkerRuntime, reason:
const start = (async () => { const start = (async () => {
try { try {
await startWorkerProcess(launch, worker, true) await startWorkerProcess(launch, worker, true)
if (launch.stopped) { // The team was ended, or the user pressed Stop while this process was
// starting and Stop's kill pass could not see it yet.
if (launch.stopped || worker.stoppedByUser) {
await conversationService.stopSessionAndWait(worker.sessionId) await conversationService.stopSessionAndWait(worker.sessionId)
return false return false
} }
@@ -335,7 +345,7 @@ async function autoContinueWorker(launch: TeamLaunch, worker: WorkerRuntime, rea
const { autoRetry: _autoRetry, ...rest } = current as MemberEntry & { autoRetry?: unknown } const { autoRetry: _autoRetry, ...rest } = current as MemberEntry & { autoRetry?: unknown }
return { ...rest, isActive: false, lastError: reason } return { ...rest, isActive: false, lastError: reason }
}) })
await notifyLead(launch, worker, { idleReason: 'failed', failureReason: `${reason} (could not continue automatically; message ${worker.member.name} to retry)` }) await notifyFailure(launch, worker, `${reason} (could not continue automatically; message ${worker.member.name} to retry)`)
} }
async function handleWorkerResult(launch: TeamLaunch, worker: WorkerRuntime, message: { is_error?: boolean; result?: unknown }): Promise<void> { async function handleWorkerResult(launch: TeamLaunch, worker: WorkerRuntime, message: { is_error?: boolean; result?: unknown }): Promise<void> {
@@ -352,6 +362,7 @@ async function handleWorkerResult(launch: TeamLaunch, worker: WorkerRuntime, mes
} }
if (!failed) { if (!failed) {
worker.autoContinueAttempts = 0 worker.autoContinueAttempts = 0
clearFailure(worker)
cancelAutoContinue(worker) cancelAutoContinue(worker)
await updateWorkerEntry(launch, worker, entry => withoutFailure({ ...entry, isActive: false })) await updateWorkerEntry(launch, worker, entry => withoutFailure({ ...entry, isActive: false }))
await notifyLead(launch, worker, { idleReason: 'available', ...(text ? { result: text } : {}) }) await notifyLead(launch, worker, { idleReason: 'available', ...(text ? { result: text } : {}) })
@@ -374,12 +385,25 @@ async function handleWorkerResult(launch: TeamLaunch, worker: WorkerRuntime, mes
const { autoRetry: _autoRetry, ...rest } = entry as MemberEntry & { autoRetry?: unknown } const { autoRetry: _autoRetry, ...rest } = entry as MemberEntry & { autoRetry?: unknown }
return { ...rest, isActive: false, lastError: reason, ...(alive ? {} : { terminated: true }) } return { ...rest, isActive: false, lastError: reason, ...(alive ? {} : { terminated: true }) }
}) })
await notifyLead(launch, worker, { await notifyFailure(launch, worker, alive
idleReason: 'failed', ? exhausted ? `${reason} (automatic retries exhausted; message ${worker.member.name} to continue)` : reason
failureReason: alive : `${worker.member.name}'s process exited (${reason}). Messaging it restarts it from its saved conversation.`)
? exhausted ? `${reason} (automatic retries exhausted; message ${worker.member.name} to continue)` : reason }
: `${worker.member.name}'s process exited (${reason}). Messaging it restarts it from its saved conversation.`,
}) /**
* 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
* (a usage limit, billing, the network) has usually been dealt with.
*/
async function notifyFailure(launch: TeamLaunch, worker: WorkerRuntime, failureReason: string): Promise<void> {
worker.failure = failureReason
worker.wakeNoticePending = true
await notifyLead(launch, worker, { idleReason: 'failed', failureReason })
}
function clearFailure(worker: WorkerRuntime): void {
worker.failure = undefined
worker.wakeNoticePending = false
} }
// ── Supervisor ────────────────────────────────────────────────────────────── // ── Supervisor ──────────────────────────────────────────────────────────────
@@ -397,6 +421,7 @@ async function deliverToWorker(launch: TeamLaunch, worker: WorkerRuntime, messag
await updateWorkerEntry(launch, worker, entry => ({ ...entry, isActive: false })) await updateWorkerEntry(launch, worker, entry => ({ ...entry, isActive: false }))
return return
} }
clearFailure(worker)
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)
@@ -450,10 +475,16 @@ async function superviseLaunch(launch: TeamLaunch): Promise<void> {
continue continue
} }
if (launch.pausedAt) { if (launch.pausedAt) {
// Only a new instruction from the lead or the user resumes a paused team. // Only a new instruction resumes a paused team: the user's own message to
const resume = messages.some(message => (message.from === TEAM_LEAD_NAME || message.from === 'user') && Date.parse(message.timestamp) >= launch.pausedAt!) // a member, or the lead's once the user has spoken to it again.
const resume = messages.some(message => {
const at = Date.parse(message.timestamp)
if (message.from === 'user') return at >= launch.pausedAt!
return message.from === TEAM_LEAD_NAME && launch.resumeAllowedAt !== undefined && at >= launch.resumeAllowedAt
})
if (!resume) continue if (!resume) continue
launch.pausedAt = undefined launch.pausedAt = undefined
launch.resumeAllowedAt = undefined
} }
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')
@@ -486,50 +517,62 @@ async function sendTeamSnapshot(parentId: string, teamName: string, createdAt: n
} }
/** The user's Stop: every member stops now, nothing is lost, the next instruction resumes. */ /** The user's Stop: every member stops now, nothing is lost, the next instruction resumes. */
function pauseLaunchesForLead(parentSessionId: string): void { function pauseLaunchesForLead(parentSessionId: string): boolean {
let paused = false
for (const launch of launches.values()) { for (const launch of launches.values()) {
if (launch.parentId !== parentSessionId || launch.stopped || !launch.running) continue if (launch.parentId !== parentSessionId || launch.stopped || !launch.running) continue
paused = true
launch.pausedAt = Date.now() launch.pausedAt = Date.now()
launch.pauseNoticePending = true launch.resumeAllowedAt = undefined
for (const worker of launch.workers.values()) { for (const worker of launch.workers.values()) {
cancelAutoContinue(worker) cancelAutoContinue(worker)
worker.stoppedByUser = true worker.stoppedByUser = true
worker.wakeNoticePending = true
} }
// No notice goes to the lead's mailbox: delivering it would start a lead // No notice goes to the lead's mailbox: delivering it would start a lead
// turn right after the user stopped everything. Messaging a stopped member // turn right after the user stopped everything. The user's next message
// restarts it, so the lead needs no special knowledge to continue later. // carries it instead (noteLeadUserMessage).
void migrationMaintenance.track(mutateTeamFileAsync(launch.plan.teamName, team => { void migrationMaintenance.track(mutateTeamFileAsync(launch.plan.teamName, team => {
if (team.createdAt !== launch.createdAt) return if (team.createdAt !== launch.createdAt) return
const sessions = new Set([...launch.workers.values()].map(worker => worker.sessionId)) const sessions = new Set([...launch.workers.values()].map(worker => worker.sessionId))
return { ...team, members: team.members.map(entry => entry.sessionId && sessions.has(entry.sessionId) ? { ...entry, isActive: false, terminated: true } : entry) } return { ...team, members: team.members.map(entry => entry.sessionId && sessions.has(entry.sessionId) ? { ...entry, isActive: false, terminated: true } : entry) }
})).catch(error => console.error('[TeamPlanRuntime] cannot record the paused team', error)) })).catch(error => console.error('[TeamPlanRuntime] cannot record the paused team', error))
} }
return paused
} }
/** /**
* The lead's first message from the user after a Stop: tell it what the Stop * The user just sent the lead a message. After a Stop this is what allows the
* did to its team. Nothing is said at the Stop itself, which would start a lead * lead to resume the team. And whether the team was stopped or members ran
* turn the user just stopped, but without this a lead told to "continue" waits * into errors (a usage limit, billing, the network, a crash), the lead is told
* for members that are no longer running. The lead decides from the user's * once which members are not running: they keep their saved conversations and
* words whether the members go on. * do nothing until messaged, so a lead told to "continue" would otherwise wait
* for them, or redo their work. The lead decides from the user's words whether
* they go on; nothing is said at the Stop or failure itself, which would start
* a lead turn the user did not ask for.
*/ */
export async function deliverTeamPauseNotice(parentSessionId: string): Promise<void> { export async function noteLeadUserMessage(parentSessionId: string): Promise<void> {
for (const launch of launches.values()) { for (const launch of launches.values()) {
if (launch.parentId !== parentSessionId || launch.stopped || !launch.pauseNoticePending) continue if (launch.parentId !== parentSessionId || launch.stopped) continue
launch.pauseNoticePending = false if (launch.pausedAt) launch.resumeAllowedAt = Date.now()
const stopped = [...launch.workers.values()].filter(worker => worker.released && worker.stoppedByUser) const waiting = [...launch.workers.values()].filter(worker => worker.released && worker.wakeNoticePending)
if (stopped.length === 0) continue if (waiting.length === 0) continue
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 lines = stopped.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
return open.length > 0 return open.length > 0
? `- ${worker.member.name}: ${open.map(task => `#${task.id} ${task.subject} (${task.status})`).join('; ')}` ? `- ${label}: ${open.map(task => `#${task.id} ${task.subject} (${task.status})`).join('; ')}`
: `- ${worker.member.name}: no unfinished task` : `- ${label}: no unfinished task`
}) })
const stoppedByUser = waiting.some(worker => worker.stoppedByUser)
await conversationService.sendMessage(parentSessionId, [ await conversationService.sendMessage(parentSessionId, [
`[Team runtime notice] The user's Stop also stopped your team ${launch.plan.teamName}. These members are stopped and keep their work so far:`, stoppedByUser
? `[Team runtime notice] The user's Stop also stopped your team ${launch.plan.teamName}. These members are stopped and keep their work so far:`
: `[Team runtime notice] These members of your team ${launch.plan.teamName} stopped on an error and keep their work so far:`,
...lines, ...lines,
'A stopped member does nothing until it is messaged: if the work should go on, SendMessage each member that still has work, and it resumes from its saved conversation where it left off. Do not spawn replacements or redo their tasks. If the user wants the team to stay stopped, leave them.', 'A stopped member does nothing until it is messaged: if the work should go on, SendMessage each member that still has work, and it resumes from its saved conversation where it left off. Do not spawn replacements or redo their tasks. If a woken member fails again with the same error (a usage limit, billing, the network), tell the user instead of retrying. If the user wants the team to stay stopped, leave them.',
].join('\n')) ].join('\n'))
} }
} }
+3 -3
View File
@@ -32,7 +32,7 @@ import {
ConversationStartupError, ConversationStartupError,
conversationService, conversationService,
} from '../services/conversationService.js' } from '../services/conversationService.js'
import { deliverTeamPauseNotice, endTeamsForParent, hasActiveTeamWorkForParent } from '../services/teamPlanRuntime.js' import { endTeamsForParent, hasActiveTeamWorkForParent, noteLeadUserMessage } from '../services/teamPlanRuntime.js'
import { computerUseApprovalService } from '../services/computerUseApprovalService.js' import { computerUseApprovalService } from '../services/computerUseApprovalService.js'
import { import {
sessionService, sessionService,
@@ -1088,8 +1088,8 @@ async function handleUserMessage(
if (!collaboration) emitSessionTurnEvent({ type: 'input-committed', sessionId }) if (!collaboration) emitSessionTurnEvent({ type: 'input-committed', sessionId })
// After the user's own words, so the lead weighs them first. // After the user's own words, so the lead weighs them first.
if (!collaboration) { if (!collaboration) {
void deliverTeamPauseNotice(sessionId).catch(error => void noteLeadUserMessage(sessionId).catch(error =>
console.error('[WS] cannot tell the lead about its stopped team', error), console.error('[WS] cannot tell the lead about its stopped team members', error),
) )
} }
} finally { } finally {