mirror of
https://github.com/NanmiCoder/claude-code-haha.git
synced 2026-10-10 11:53:10 +08:00
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.
This commit is contained in:
@@ -1,5 +1,7 @@
|
||||
import { getComposerViewForTesting } from './MentionComposer'
|
||||
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 { useSideChatStore } from '@/stores/sideChatStore'
|
||||
import { fireEvent, render, screen, waitFor, within } from '@testing-library/react'
|
||||
@@ -248,6 +250,7 @@ describe('ChatInput file mentions', () => {
|
||||
vi.clearAllMocks()
|
||||
useSideChatStore.setState({ entries: {} })
|
||||
useTeamPlanStore.setState({ bySession: {} })
|
||||
useTeamStore.setState({ workbenchesBySession: {} })
|
||||
mocks.voiceSupported.mockReturnValue(false)
|
||||
useVoiceInputStore.setState({ catalog: null, loading: false, error: null })
|
||||
mocks.sideOpen.mockResolvedValue('side-tab')
|
||||
@@ -1357,6 +1360,44 @@ describe('ChatInput file mentions', () => {
|
||||
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) => {
|
||||
useChatStore.setState({
|
||||
sessions: {
|
||||
|
||||
@@ -19,6 +19,7 @@ import { useUIStore } from '../../stores/uiStore'
|
||||
import { useSessionStore } from '../../stores/sessionStore'
|
||||
import { useSessionRuntimeStore } from '../../stores/sessionRuntimeStore'
|
||||
import { useTeamStore } from '../../stores/teamStore'
|
||||
import { getMemberWorkState, resolveTeamMemberIdentity } from '../agentTeams/agentTeamsModel'
|
||||
import { useTeamPlanStore } from '@/stores/teamPlanStore'
|
||||
import { useSettingsStore } from '../../stores/settingsStore'
|
||||
import {
|
||||
@@ -300,10 +301,22 @@ export function ChatInput({ variant = 'default', compact = false, sessionId, vis
|
||||
const hasRunningSubagents = hasRunningSubagentTasks(sessionState?.backgroundAgentTasks)
|
||||
// Approved team processes are tracked by their plan, not background-agent
|
||||
// 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
|
||||
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 isWorkspaceMissing = workspaceState !== 'available'
|
||||
// Both composer branches (hero and inline) and the drop handler share this:
|
||||
|
||||
@@ -107,7 +107,7 @@ Claude 每完成一轮并修改文件,对话里会出现一张「{n} 个文件
|
||||
|
||||
后台跑的子 Agent 的工具活动也会冒泡到这里,不用等它跑完才知道它在干什么。
|
||||
|
||||
团队成员遇到模型服务断流、限流、5xx 这类临时故障会自己重试,成员行显示「自动重试 2/5」,不用你或主控介入;重试用完,或者遇到需要你处理的错误(比如 API Key 失效),成员行显示「出错」并通知主控。按停止按钮会让整组成员一起停下,但进度不会丢:你下一次给主控发消息时,它会知道哪些成员停在了哪个任务上,再按你的意思决定要不要让它们继续;你也可以直接给某个成员发消息,它会从保存的对话接着干,显示「已停止」的成员也一样。切换模型或权限模式不会打断团队;重启应用后团队还在,给成员发消息即可继续。删除会话或执行 `/clear` 才会结束团队。
|
||||
团队成员遇到模型服务断流、限流、5xx 这类临时故障会自己重试,成员行显示「自动重试 2/5」,不用你或主控介入;重试用完,或者遇到需要你处理的错误(比如 API Key 失效、余额不足、订阅额度用完),成员行显示「出错」并通知主控。问题解决后给主控发一句「继续」就行:它会收到出错成员的清单,逐个唤醒它们从保存的对话接着干,不会重做已完成的部分。按停止按钮会让整组成员一起停下,主控也不会再自己处理成员发来的汇报,但进度不会丢:你下一次给主控发消息时,它会知道哪些成员停在了哪个任务上,再按你的意思决定要不要让它们继续;你也可以直接给某个成员发消息,它会从保存的对话接着干,显示「已停止」的成员也一样。切换模型或权限模式不会打断团队;重启应用后团队还在,给成员发消息即可继续。删除会话或执行 `/clear` 才会结束团队。
|
||||
|
||||
## 轨迹:每一步到底发生了什么
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
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
|
||||
|
||||
|
||||
@@ -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
@@ -437,14 +437,20 @@ export type LeadMailboxPollState = {
|
||||
finalDrainDone: boolean
|
||||
/** Consecutive polls whose batch could not be acknowledged. */
|
||||
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 {
|
||||
return { finalDrainDone: false, consecutiveAckFailures: 0 }
|
||||
return { finalDrainDone: false, consecutiveAckFailures: 0, held: false }
|
||||
}
|
||||
|
||||
export type LeadMailboxPollStep =
|
||||
| { kind: 'stop' }
|
||||
| { kind: 'held' }
|
||||
| { kind: 'idle' }
|
||||
| { kind: 'retry' }
|
||||
| { kind: 'batch'; messages: TeammateMessage[] }
|
||||
@@ -461,6 +467,7 @@ export async function takeLeadMailboxBatch(
|
||||
options: { teamName: string | undefined; hasActiveTeammates: boolean },
|
||||
): Promise<LeadMailboxPollStep> {
|
||||
const { teamName, hasActiveTeammates } = options
|
||||
if (state.held) return { kind: 'held' }
|
||||
const unread = await readUnreadMessages(TEAM_LEAD_NAME, teamName)
|
||||
|
||||
if (!hasActiveTeammates) {
|
||||
@@ -2693,6 +2700,13 @@ function runHeadlessStreaming(
|
||||
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') {
|
||||
await sleep(POLL_INTERVAL_MS)
|
||||
continue
|
||||
@@ -3020,6 +3034,7 @@ function runHeadlessStreaming(
|
||||
}
|
||||
} else if (message.request.subtype === 'interrupt' || message.request.subtype === 'team_plan_pause') {
|
||||
sessionMessageInbox.cancelQueued(dequeueAllMatching)
|
||||
if (message.request.subtype === 'team_plan_pause') leadMailboxPoll.held = true
|
||||
// Track escapes for attribution (ant-only feature)
|
||||
if (feature('COMMIT_ATTRIBUTION')) {
|
||||
setAppState(prev => ({
|
||||
@@ -4452,6 +4467,9 @@ function runHeadlessStreaming(
|
||||
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({
|
||||
mode: 'prompt' as const,
|
||||
// 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' })
|
||||
})
|
||||
|
||||
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 () => {
|
||||
const service = new ConversationService() as any
|
||||
const killed: string[] = []
|
||||
|
||||
@@ -335,8 +335,8 @@ export type TeamWorkerStart = {
|
||||
* own: a lead restart keeps them, and Stop lets the runtime pause them first.
|
||||
*/
|
||||
export type TeamRuntimeListener = {
|
||||
/** Synchronously before Stop kills a lead's workers. */
|
||||
leadInterrupted?: (parentSessionId: string) => void
|
||||
/** Synchronously before Stop kills a lead's workers; true when it paused a running team. */
|
||||
leadInterrupted?: (parentSessionId: string) => boolean | void
|
||||
/** After any CLI session finished starting. */
|
||||
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
|
||||
// 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.
|
||||
let pausedTeam = false
|
||||
for (const listener of this.teamRuntimeListeners) {
|
||||
try {
|
||||
listener.leadInterrupted?.(sessionId)
|
||||
if (listener.leadInterrupted?.(sessionId) === true) pausedTeam = true
|
||||
} catch (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)
|
||||
})
|
||||
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, {
|
||||
type: 'control_request',
|
||||
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) {
|
||||
const { writeToMailbox } = await import('../../utils/teammateMailbox.js')
|
||||
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.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 h.waitFor(() => h.requests.length > before, 'resumed member turn')
|
||||
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)
|
||||
|
||||
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'])
|
||||
test("a user's direct message resumes a stopped member without waiting for the lead", async () => {
|
||||
const h = await startHarness(['reader'])
|
||||
try {
|
||||
const { deliverTeamPauseNotice } = await import('./teamPlanRuntime.js')
|
||||
await h.waitFor(async () => (await h.leadNotifications()).length === 2, 'initial turn reports')
|
||||
const sent: Array<[string, string]> = []
|
||||
const send = h.service.sendMessage.bind(h.service)
|
||||
const spy = spyOn(h.service, 'sendMessage').mockImplementation(async (...args: Parameters<typeof send>) => {
|
||||
sent.push([args[0], String(args[1])])
|
||||
return send(...args)
|
||||
await h.waitFor(async () => (await h.leadNotifications()).length === 1, 'initial turn report')
|
||||
h.service.sendInterrupt(h.parentId)
|
||||
await h.service.waitForTeamWorkersStopped(h.parentId)
|
||||
const before = h.requests.length
|
||||
await messageMember(h, 'reader', 'user', 'Keep going on your own')
|
||||
await h.waitFor(() => h.requests.length > before, 'resumed member turn')
|
||||
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 {
|
||||
// Without a Stop there is nothing to tell.
|
||||
await deliverTeamPauseNotice(h.parentId)
|
||||
expect(sent).toEqual([])
|
||||
|
||||
h.service.sendInterrupt(h.parentId)
|
||||
await messageMember(h, 'reader', 'team-lead', 'Check one more file')
|
||||
await h.waitFor(() => stopped, 'restart to begin')
|
||||
await h.service.waitForTeamWorkersStopped(h.parentId)
|
||||
await deliverTeamPauseNotice(h.parentId)
|
||||
await deliverTeamPauseNotice(h.parentId)
|
||||
|
||||
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')
|
||||
await new Promise(resolve => setTimeout(resolve, 800))
|
||||
expect(h.service.hasSession(sessionId)).toBe(false)
|
||||
expect((await h.teamMember('reader'))?.terminated).toBe(true)
|
||||
} finally {
|
||||
spy.mockRestore()
|
||||
}
|
||||
@@ -318,6 +360,87 @@ test("after a Stop, the lead's next message tells it which members stopped and h
|
||||
}
|
||||
}, 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 () => {
|
||||
const h = await startHarness(['reader'])
|
||||
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.waitFor(async () => (await h.teamMember('reader'))?.terminated === true, `stop ${cycle} to be recorded`)
|
||||
const before = h.requests.length
|
||||
await userMessagesLead(h)
|
||||
await messageMember(h, 'reader', 'team-lead', `Continue after stop ${cycle}`)
|
||||
await h.waitFor(() => h.requests.length > before, `resume ${cycle}`)
|
||||
expect(h.requests.at(-1)?.prompt).toContain(`Continue after stop ${cycle}`)
|
||||
|
||||
@@ -52,6 +52,10 @@ type WorkerRuntime = {
|
||||
restartsExhausted: boolean
|
||||
/** The user's Stop ended this process; resuming it is not a crash restart. */
|
||||
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>
|
||||
lastResultAt?: number
|
||||
}
|
||||
@@ -70,8 +74,12 @@ type TeamLaunch = {
|
||||
stopped: boolean
|
||||
/** Set by the user's Stop: hold automatic work until the lead or the user acts again. */
|
||||
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
|
||||
}
|
||||
|
||||
@@ -281,7 +289,7 @@ async function restartWorker(launch: TeamLaunch, worker: WorkerRuntime, reason:
|
||||
if (!worker.restartsExhausted) {
|
||||
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 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
|
||||
}
|
||||
@@ -290,7 +298,9 @@ async function restartWorker(launch: TeamLaunch, worker: WorkerRuntime, reason:
|
||||
const start = (async () => {
|
||||
try {
|
||||
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)
|
||||
return false
|
||||
}
|
||||
@@ -335,7 +345,7 @@ async function autoContinueWorker(launch: TeamLaunch, worker: WorkerRuntime, rea
|
||||
const { autoRetry: _autoRetry, ...rest } = current as MemberEntry & { autoRetry?: unknown }
|
||||
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> {
|
||||
@@ -352,6 +362,7 @@ async function handleWorkerResult(launch: TeamLaunch, worker: WorkerRuntime, mes
|
||||
}
|
||||
if (!failed) {
|
||||
worker.autoContinueAttempts = 0
|
||||
clearFailure(worker)
|
||||
cancelAutoContinue(worker)
|
||||
await updateWorkerEntry(launch, worker, entry => withoutFailure({ ...entry, isActive: false }))
|
||||
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 }
|
||||
return { ...rest, isActive: false, lastError: reason, ...(alive ? {} : { terminated: true }) }
|
||||
})
|
||||
await notifyLead(launch, worker, {
|
||||
idleReason: 'failed',
|
||||
failureReason: alive
|
||||
? 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.`,
|
||||
})
|
||||
await notifyFailure(launch, worker, alive
|
||||
? 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 ──────────────────────────────────────────────────────────────
|
||||
@@ -397,6 +421,7 @@ async function deliverToWorker(launch: TeamLaunch, worker: WorkerRuntime, messag
|
||||
await updateWorkerEntry(launch, worker, entry => ({ ...entry, isActive: false }))
|
||||
return
|
||||
}
|
||||
clearFailure(worker)
|
||||
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])))
|
||||
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
|
||||
}
|
||||
if (launch.pausedAt) {
|
||||
// Only a new instruction from the lead or the user resumes a paused team.
|
||||
const resume = messages.some(message => (message.from === TEAM_LEAD_NAME || message.from === 'user') && Date.parse(message.timestamp) >= launch.pausedAt!)
|
||||
// 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.
|
||||
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
|
||||
launch.pausedAt = undefined
|
||||
launch.resumeAllowedAt = undefined
|
||||
}
|
||||
if (!conversationService.hasSession(worker.sessionId)) {
|
||||
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. */
|
||||
function pauseLaunchesForLead(parentSessionId: string): void {
|
||||
function pauseLaunchesForLead(parentSessionId: string): boolean {
|
||||
let paused = false
|
||||
for (const launch of launches.values()) {
|
||||
if (launch.parentId !== parentSessionId || launch.stopped || !launch.running) continue
|
||||
paused = true
|
||||
launch.pausedAt = Date.now()
|
||||
launch.pauseNoticePending = true
|
||||
launch.resumeAllowedAt = undefined
|
||||
for (const worker of launch.workers.values()) {
|
||||
cancelAutoContinue(worker)
|
||||
worker.stoppedByUser = true
|
||||
worker.wakeNoticePending = true
|
||||
}
|
||||
// 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
|
||||
// restarts it, so the lead needs no special knowledge to continue later.
|
||||
// turn right after the user stopped everything. The user's next message
|
||||
// carries it instead (noteLeadUserMessage).
|
||||
void migrationMaintenance.track(mutateTeamFileAsync(launch.plan.teamName, team => {
|
||||
if (team.createdAt !== launch.createdAt) return
|
||||
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) }
|
||||
})).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
|
||||
* did to its team. Nothing is said at the Stop itself, which would start a lead
|
||||
* turn the user just stopped, but without this a lead told to "continue" waits
|
||||
* for members that are no longer running. The lead decides from the user's
|
||||
* words whether the members go on.
|
||||
* The user just sent the lead a message. After a Stop this is what allows the
|
||||
* lead to resume the team. And whether the team was stopped or members ran
|
||||
* into errors (a usage limit, billing, the network, a crash), the lead is told
|
||||
* once which members are not running: they keep their saved conversations and
|
||||
* 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()) {
|
||||
if (launch.parentId !== parentSessionId || launch.stopped || !launch.pauseNoticePending) continue
|
||||
launch.pauseNoticePending = false
|
||||
const stopped = [...launch.workers.values()].filter(worker => worker.released && worker.stoppedByUser)
|
||||
if (stopped.length === 0) continue
|
||||
if (launch.parentId !== parentSessionId || launch.stopped) continue
|
||||
if (launch.pausedAt) launch.resumeAllowedAt = Date.now()
|
||||
const waiting = [...launch.workers.values()].filter(worker => worker.released && worker.wakeNoticePending)
|
||||
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 lines = stopped.map(worker => {
|
||||
const lines = waiting.map(worker => {
|
||||
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
|
||||
? `- ${worker.member.name}: ${open.map(task => `#${task.id} ${task.subject} (${task.status})`).join('; ')}`
|
||||
: `- ${worker.member.name}: no unfinished task`
|
||||
? `- ${label}: ${open.map(task => `#${task.id} ${task.subject} (${task.status})`).join('; ')}`
|
||||
: `- ${label}: no unfinished task`
|
||||
})
|
||||
const stoppedByUser = waiting.some(worker => worker.stoppedByUser)
|
||||
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,
|
||||
'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'))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -32,7 +32,7 @@ import {
|
||||
ConversationStartupError,
|
||||
conversationService,
|
||||
} 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 {
|
||||
sessionService,
|
||||
@@ -1088,8 +1088,8 @@ async function handleUserMessage(
|
||||
if (!collaboration) emitSessionTurnEvent({ type: 'input-committed', sessionId })
|
||||
// After the user's own words, so the lead weighs them first.
|
||||
if (!collaboration) {
|
||||
void deliverTeamPauseNotice(sessionId).catch(error =>
|
||||
console.error('[WS] cannot tell the lead about its stopped team', error),
|
||||
void noteLeadUserMessage(sessionId).catch(error =>
|
||||
console.error('[WS] cannot tell the lead about its stopped team members', error),
|
||||
)
|
||||
}
|
||||
} finally {
|
||||
|
||||
Reference in New Issue
Block a user