Skip to content

Commit 143a13c

Browse files
committed
fix(delegation): settle stopped subagents reliably
1 parent ca9b404 commit 143a13c

12 files changed

Lines changed: 1698 additions & 122 deletions

src/main/delegation/acp-execution.test.ts

Lines changed: 148 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -631,6 +631,154 @@ describe('ACP delegate execution production adapter', () => {
631631
await expect(execution.reserve(1)).resolves.toHaveProperty('slotIds')
632632
})
633633

634+
it('settles cancellation when the terminated provider prompt never returns', async () => {
635+
const prompt = deferred<PromptResponse>()
636+
const cleanup: string[] = []
637+
let callbacks!: AcpDelegateExecutionCallbacks
638+
const execution = createAcpDelegateExecution({
639+
capacity: 1,
640+
prepare: async (input) => ({
641+
executionId: input.attemptId,
642+
provenance: {
643+
projectId: input.session.projectId,
644+
sessionId: input.session.sessionId,
645+
agentFrameId: input.frameId,
646+
runtimeSegmentId: input.runtimeSegmentId
647+
},
648+
workspace: { cwd: '/workspace/terminated-provider' },
649+
runtimeHome: '/runtime/terminated-provider',
650+
frameworkId: 'certified-test',
651+
capability: {
652+
revoke: async () => {
653+
cleanup.push('revoke')
654+
}
655+
},
656+
disposeResources: async () => {
657+
cleanup.push('resources')
658+
}
659+
}),
660+
assertFrameworkNativeDelegationDisabled: async () => undefined,
661+
createRuntime: (_scope, runtimeCallbacks) => {
662+
callbacks = runtimeCallbacks
663+
return {
664+
createSession: async () => ({ sessionId: 'provider-terminated' }),
665+
sendAppContinuation: () => {
666+
callbacks.onProviderPromptAccepted('provider-terminated')
667+
return prompt.promise
668+
},
669+
cancelPrompt: async () => undefined,
670+
respondToPermission: async () => undefined,
671+
setPermissionProfile: async () => undefined,
672+
deleteSession: async () => {
673+
cleanup.push('delete')
674+
},
675+
shutdownForQuit: async () => {
676+
cleanup.push('shutdown')
677+
return { reaped: true }
678+
}
679+
}
680+
}
681+
})
682+
const reservation = await execution.reserve(1)
683+
const running = execution.run(makeInput('terminated-provider'), reservation.slotIds[0])
684+
await running.accepted
685+
686+
await expect(running.cancel()).resolves.toBeUndefined()
687+
await expect(running.completion).resolves.toEqual({ status: 'cancelled' })
688+
await vi.waitFor(() => expect(cleanup).toEqual(['revoke', 'delete', 'shutdown', 'resources']))
689+
await expect(execution.reserve(1)).resolves.toHaveProperty('slotIds')
690+
})
691+
692+
it('does not start a provider prompt when cancellation lands during deferred Turn begin', async () => {
693+
const { execution, controls } = makeHarness(1)
694+
const begin = deferred<void>()
695+
const reservation = await execution.reserve(1)
696+
const running = execution.run(
697+
{
698+
...makeInput('deferred-begin'),
699+
turn: {
700+
promptMessageId: 'prompt-deferred-begin',
701+
messageBranchId: 'branch-deferred-begin',
702+
runtimeSegmentId: 'segment-deferred-begin',
703+
begin: () => begin.promise
704+
}
705+
},
706+
reservation.slotIds[0]
707+
)
708+
await vi.waitFor(() => expect(controls.has('deferred-begin')).toBe(true))
709+
710+
await running.cancel()
711+
begin.resolve()
712+
713+
await expect(running.accepted).rejects.toMatchObject({
714+
name: 'DelegateMessagePreAcceptanceError'
715+
})
716+
await expect(running.completion).resolves.toEqual({ status: 'cancelled' })
717+
await vi.waitFor(() => expect(controls.get('deferred-begin')?.prompts).toEqual([]))
718+
})
719+
720+
it('reports cleanup failure and retries only the unfinished cleanup step', async () => {
721+
let revokeAttempts = 0
722+
const cleanup: string[] = []
723+
let callbacks!: AcpDelegateExecutionCallbacks
724+
const execution = createAcpDelegateExecution({
725+
capacity: 1,
726+
prepare: async (input) => ({
727+
executionId: input.attemptId,
728+
provenance: {
729+
projectId: input.session.projectId,
730+
sessionId: input.session.sessionId,
731+
agentFrameId: input.frameId,
732+
runtimeSegmentId: input.runtimeSegmentId
733+
},
734+
workspace: { cwd: '/workspace/retry-cleanup' },
735+
runtimeHome: '/runtime/retry-cleanup',
736+
frameworkId: 'certified-test',
737+
capability: {
738+
async revoke() {
739+
revokeAttempts += 1
740+
if (revokeAttempts === 1) throw new Error('revoke failed once')
741+
cleanup.push('revoke')
742+
}
743+
},
744+
disposeResources: async () => {
745+
cleanup.push('resources')
746+
}
747+
}),
748+
assertFrameworkNativeDelegationDisabled: async () => undefined,
749+
createRuntime: (_scope, runtimeCallbacks) => {
750+
callbacks = runtimeCallbacks
751+
return {
752+
createSession: async () => ({ sessionId: 'provider-retry-cleanup' }),
753+
sendAppContinuation: () => {
754+
callbacks.onProviderPromptAccepted('provider-retry-cleanup')
755+
return new Promise<PromptResponse>(() => undefined)
756+
},
757+
cancelPrompt: async () => undefined,
758+
respondToPermission: async () => undefined,
759+
setPermissionProfile: async () => undefined,
760+
deleteSession: async () => {
761+
cleanup.push('delete')
762+
},
763+
shutdownForQuit: async () => {
764+
cleanup.push('shutdown')
765+
return { reaped: true }
766+
}
767+
}
768+
}
769+
})
770+
const reservation = await execution.reserve(1)
771+
const running = execution.run(makeInput('retry-cleanup'), reservation.slotIds[0])
772+
await running.accepted
773+
774+
await expect(running.cancel()).rejects.toThrow('revoke failed once')
775+
await expect(running.cancel()).resolves.toBeUndefined()
776+
777+
expect(revokeAttempts).toBe(2)
778+
expect(cleanup).toEqual(['revoke', 'delete', 'shutdown', 'resources'])
779+
await expect(execution.reserve(1)).resolves.toHaveProperty('slotIds')
780+
})
781+
634782
it('reserves an entire batch atomically and releases terminal slots', async () => {
635783
const { execution, controls } = makeHarness(2)
636784
const reservation = await execution.reserve(2)

src/main/delegation/acp-execution.ts

Lines changed: 45 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -285,6 +285,10 @@ const createAcpDelegateExecution = (options: AcpDelegateExecutionOptions): Deleg
285285
let ownsWorkspace = false
286286
let writable = true
287287
let capabilityRevoked = false
288+
let providerSessionDeleted = false
289+
let runtimeShutdown = false
290+
let resourcesDisposed = false
291+
let slotReleased = false
288292
let acceptedSettled = false
289293
let terminalSettled = false
290294
let cancelRequested = false
@@ -396,27 +400,29 @@ const createAcpDelegateExecution = (options: AcpDelegateExecutionOptions): Deleg
396400
activeMessage = undefined
397401
for (const pending of queuedPrompts.splice(0)) pending.acceptance.reject(deliveryError)
398402
if (scope && !capabilityRevoked) {
399-
capabilityRevoked = true
400403
await scope.capability.revoke()
404+
capabilityRevoked = true
401405
}
402406
}
403-
const cleanup = async (): Promise<void> => {
407+
const cleanupOnce = async (): Promise<void> => {
404408
let firstError: unknown
405409
try {
406410
await revokeWrites()
407411
} catch (error) {
408412
firstError = error
409413
}
410-
if (runtime && providerSessionId) {
414+
if (runtime && providerSessionId && !providerSessionDeleted) {
411415
try {
412416
await runtime.deleteSession({ sessionId: providerSessionId })
417+
providerSessionDeleted = true
413418
} catch (error) {
414419
firstError ??= error
415420
}
416421
}
417-
if (runtime) {
422+
if (runtime && !runtimeShutdown) {
418423
try {
419424
await runtime.shutdownForQuit()
425+
runtimeShutdown = true
420426
} catch (error) {
421427
firstError ??= error
422428
}
@@ -430,16 +436,28 @@ const createAcpDelegateExecution = (options: AcpDelegateExecutionOptions): Deleg
430436
activeWorkspaces.delete(scope.workspace.cwd)
431437
ownsWorkspace = false
432438
}
433-
try {
434-
await scope.disposeResources?.()
435-
} catch (error) {
436-
firstError ??= error
439+
if (!resourcesDisposed) {
440+
try {
441+
await scope.disposeResources?.()
442+
resourcesDisposed = true
443+
} catch (error) {
444+
firstError ??= error
445+
}
437446
}
438447
}
439448
listeners.clear()
440-
releaseSlot(slotId)
449+
if (!slotReleased) {
450+
slotReleased = true
451+
releaseSlot(slotId)
452+
}
441453
if (firstError !== undefined) throw firstError
442454
}
455+
let cleanupTail = Promise.resolve()
456+
const cleanup = (): Promise<void> => {
457+
const next = cleanupTail.then(cleanupOnce, cleanupOnce)
458+
cleanupTail = next.catch(() => undefined)
459+
return next
460+
}
443461

444462
const promptRequest = (
445463
text: string
@@ -533,6 +551,7 @@ const createAcpDelegateExecution = (options: AcpDelegateExecutionOptions): Deleg
533551
let response = ''
534552
while (!cancelRequested) {
535553
await activeTurn?.begin?.()
554+
if (cancelRequested) break
536555
currentResponse = []
537556
providerPromptStarted = true
538557
const outcome = await runtime.sendAppContinuation(promptRequest(nextPrompt))
@@ -653,13 +672,27 @@ const createAcpDelegateExecution = (options: AcpDelegateExecutionOptions): Deleg
653672
publish({ kind: 'permission', awaiting: false, requestId: response.requestId })
654673
},
655674
async cancel() {
656-
if (terminalSettled) return
675+
if (terminalSettled && !cancelRequested) return
657676
cancelRequested = true
658-
await revokeWrites().catch(() => undefined)
677+
await revokeWrites()
659678
if (runtime && providerSessionId) {
660679
await runtime.cancelPrompt({ sessionId: providerSessionId }).catch(() => undefined)
661680
}
662-
await work.catch(() => undefined)
681+
// A terminated provider may never settle its in-flight prompt. Cancellation owns the
682+
// terminal signal; cleanup is idempotent and remains observed if transport teardown stalls
683+
// or the provider task resumes later.
684+
settleAccepted(
685+
'provider_prompt_completed',
686+
new DelegateMessagePreAcceptanceError(
687+
'delegate execution was cancelled before provider acceptance'
688+
)
689+
)
690+
if (!terminalSettled) {
691+
terminalSettled = true
692+
terminal.resolve({ status: 'cancelled' })
693+
}
694+
await cleanup()
695+
void work.catch(() => undefined)
663696
}
664697
})
665698
}

src/main/delegation/delegated-turn-lifecycle.ts

Lines changed: 45 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -64,27 +64,44 @@ const createDelegatedTurnLifecycle = (options: {
6464
}> => {
6565
let currentArtifact: DelegatedArtifactHandle | undefined
6666
const artifactHandles: DelegatedArtifactHandle[] = []
67+
const disposedArtifacts = new Set<DelegatedArtifactHandle>()
68+
const pendingArtifactOpens = new Set<Promise<void>>()
6769
let artifactHandoffFile: string | undefined
6870
let stagedRuntimeUpdateCount = 0
6971
let completedTurnMessage: DurableMessage | undefined
72+
let disposeRequested = false
7073

71-
const openArtifact = async (context: TurnContext, executionId: string): Promise<void> => {
72-
const artifact = await options.artifactEvidence?.open({
73-
session: options.session,
74-
executionId,
75-
attemptId: options.attemptId,
76-
rootFrameId: context.rootFrameId,
77-
agentFrameId: options.agentFrameId,
78-
messageBranchId: context.messageBranchId,
79-
runtimeSegmentId: context.runtimeSegmentId,
80-
promptMessageId: context.promptMessageId,
81-
agentName: options.agentName
82-
})
83-
if (!artifact) return
84-
currentArtifact = artifact
85-
artifactHandles.push(artifact)
86-
if (artifactHandoffFile) await artifact.activateAt?.(artifactHandoffFile)
87-
else artifactHandoffFile = artifact.execution?.currentRunFile
74+
const openArtifact = (context: TurnContext, executionId: string): Promise<void> => {
75+
if (disposeRequested) return Promise.resolve()
76+
const opening = (async () => {
77+
const artifact = await options.artifactEvidence?.open({
78+
session: options.session,
79+
executionId,
80+
attemptId: options.attemptId,
81+
rootFrameId: context.rootFrameId,
82+
agentFrameId: options.agentFrameId,
83+
messageBranchId: context.messageBranchId,
84+
runtimeSegmentId: context.runtimeSegmentId,
85+
promptMessageId: context.promptMessageId,
86+
agentName: options.agentName
87+
})
88+
if (!artifact) return
89+
artifactHandles.push(artifact)
90+
if (disposeRequested) {
91+
await artifact.dispose()
92+
disposedArtifacts.add(artifact)
93+
return
94+
}
95+
currentArtifact = artifact
96+
if (artifactHandoffFile) await artifact.activateAt?.(artifactHandoffFile)
97+
else artifactHandoffFile = artifact.execution?.currentRunFile
98+
})()
99+
pendingArtifactOpens.add(opening)
100+
void opening.then(
101+
() => pendingArtifactOpens.delete(opening),
102+
() => pendingArtifactOpens.delete(opening)
103+
)
104+
return opening
88105
}
89106

90107
return {
@@ -140,7 +157,17 @@ const createDelegatedTurnLifecycle = (options: {
140157
await currentArtifact?.finalize(terminalMessageId)
141158
},
142159
async dispose() {
143-
await Promise.allSettled(artifactHandles.map((artifact) => artifact.dispose()))
160+
disposeRequested = true
161+
const openingAtDisposal = [...pendingArtifactOpens]
162+
const disposalAtStart = artifactHandles
163+
.filter((artifact) => !disposedArtifacts.has(artifact))
164+
.map(async (artifact) => {
165+
await artifact.dispose()
166+
disposedArtifacts.add(artifact)
167+
})
168+
// An open that was already in flight when disposal began owns disposal of its late handle.
169+
// Await it alongside current handles so failure keeps the outer finalizer retryable.
170+
await Promise.all([...disposalAtStart, ...openingAtDisposal])
144171
}
145172
}
146173
}

src/main/delegation/durable-delegated-work-contract.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -298,6 +298,10 @@ type CreateDurableDelegatedWorkOptions = Readonly<{
298298
reviewEvidence?: DelegatedReviewEvidence
299299
onRootPermissionEvent?(event: RootDelegatePermissionEvent): void
300300
onAgentRuntimeUpdate?(update: AcpAgentRuntimeUpdate): void
301+
onCleanupError?(
302+
scope: Readonly<{ session: SessionKey; frameId: string; attemptId: string }>,
303+
error: unknown
304+
): void
301305
now?: () => number
302306
createId?: (kind: 'frame' | 'attempt' | 'message' | 'runtime' | 'question') => string
303307
collectPollIntervalMs?: number

0 commit comments

Comments
 (0)