From 786c7a84d00c515d7b0b2c8ed698a6ec7e2c83a6 Mon Sep 17 00:00:00 2001 From: YeonGyu-Kim Date: Fri, 13 Mar 2026 13:12:59 +0900 Subject: [PATCH] fix(background-agent): prevent queue item loss on concurrent cancel and guard against cancelled task resurrection --- src/features/background-agent/manager.test.ts | 164 ++++++++++++++++++ src/features/background-agent/manager.ts | 19 +- 2 files changed, 179 insertions(+), 4 deletions(-) diff --git a/src/features/background-agent/manager.test.ts b/src/features/background-agent/manager.test.ts index 78a7b64f0..8fd3e2133 100644 --- a/src/features/background-agent/manager.test.ts +++ b/src/features/background-agent/manager.test.ts @@ -2195,6 +2195,170 @@ describe("BackgroundManager - Non-blocking Queue Integration", () => { // then expect(retryTask.status).toBe("pending") }) + + test("should keep the next queued task when the first task is cancelled during session creation", async () => { + // given + const firstSessionID = "ses-first-cancelled-during-create" + const secondSessionID = "ses-second-survives-queue" + let createCallCount = 0 + let resolveFirstCreate: ((value: { data: { id: string } }) => void) | undefined + let resolveFirstCreateStarted: (() => void) | undefined + let resolveSecondPromptAsync: (() => void) | undefined + const firstCreateStarted = new Promise((resolve) => { + resolveFirstCreateStarted = resolve + }) + const secondPromptAsyncStarted = new Promise((resolve) => { + resolveSecondPromptAsync = resolve + }) + + manager.shutdown() + manager = new BackgroundManager( + { + client: { + session: { + create: async () => { + createCallCount += 1 + if (createCallCount === 1) { + resolveFirstCreateStarted?.() + return await new Promise<{ data: { id: string } }>((resolve) => { + resolveFirstCreate = resolve + }) + } + + return { data: { id: secondSessionID } } + }, + get: async () => ({ data: { directory: "/test/dir" } }), + prompt: async () => ({}), + promptAsync: async ({ path }: { path: { id: string } }) => { + if (path.id === secondSessionID) { + resolveSecondPromptAsync?.() + } + + return {} + }, + messages: async () => ({ data: [] }), + todo: async () => ({ data: [] }), + status: async () => ({ data: {} }), + abort: async () => ({}), + }, + }, + directory: tmpdir(), + } as unknown as PluginInput, + { defaultConcurrency: 1 } + ) + + const input = { + description: "Test task", + prompt: "Do something", + agent: "test-agent", + parentSessionID: "parent-session", + parentMessageID: "parent-message", + } + + const firstTask = await manager.launch(input) + const secondTask = await manager.launch(input) + await firstCreateStarted + + // when + const cancelled = await manager.cancelTask(firstTask.id, { + source: "test", + abortSession: false, + }) + resolveFirstCreate?.({ data: { id: firstSessionID } }) + + await Promise.race([ + secondPromptAsyncStarted, + new Promise((_, reject) => setTimeout(() => reject(new Error("timeout")), 100)), + ]) + + // then + expect(cancelled).toBe(true) + expect(createCallCount).toBe(2) + expect(manager.getTask(firstTask.id)?.status).toBe("cancelled") + expect(manager.getTask(secondTask.id)?.status).toBe("running") + expect(manager.getTask(secondTask.id)?.sessionID).toBe(secondSessionID) + }) + + test("should keep task cancelled and abort the session when cancellation wins during session creation", async () => { + // given + const createdSessionID = "ses-cancelled-during-create" + let resolveCreate: ((value: { data: { id: string } }) => void) | undefined + let resolveCreateStarted: (() => void) | undefined + let resolveAbortCalled: (() => void) | undefined + const createStarted = new Promise((resolve) => { + resolveCreateStarted = resolve + }) + const abortCalled = new Promise((resolve) => { + resolveAbortCalled = resolve + }) + const abortCalls: string[] = [] + const promptAsyncSessionIDs: string[] = [] + + manager.shutdown() + manager = new BackgroundManager( + { + client: { + session: { + create: async () => { + resolveCreateStarted?.() + return await new Promise<{ data: { id: string } }>((resolve) => { + resolveCreate = resolve + }) + }, + get: async () => ({ data: { directory: "/test/dir" } }), + prompt: async () => ({}), + promptAsync: async ({ path }: { path: { id: string } }) => { + promptAsyncSessionIDs.push(path.id) + return {} + }, + messages: async () => ({ data: [] }), + todo: async () => ({ data: [] }), + status: async () => ({ data: {} }), + abort: async ({ path }: { path: { id: string } }) => { + abortCalls.push(path.id) + resolveAbortCalled?.() + return {} + }, + }, + }, + directory: tmpdir(), + } as unknown as PluginInput, + { defaultConcurrency: 1 } + ) + + const input = { + description: "Test task", + prompt: "Do something", + agent: "test-agent", + parentSessionID: "parent-session", + parentMessageID: "parent-message", + } + + const task = await manager.launch(input) + await createStarted + + // when + const cancelled = await manager.cancelTask(task.id, { + source: "test", + abortSession: false, + }) + resolveCreate?.({ data: { id: createdSessionID } }) + + await Promise.race([ + abortCalled, + new Promise((_, reject) => setTimeout(() => reject(new Error("timeout")), 100)), + ]) + await Promise.resolve() + + // then + const updatedTask = manager.getTask(task.id) + expect(cancelled).toBe(true) + expect(updatedTask?.status).toBe("cancelled") + expect(updatedTask?.sessionID).toBeUndefined() + expect(promptAsyncSessionIDs).not.toContain(createdSessionID) + expect(abortCalls).toEqual([createdSessionID]) + expect(getConcurrencyManager(manager).getCount("test-agent")).toBe(0) + }) }) describe("pending task can be cancelled", () => { diff --git a/src/features/background-agent/manager.ts b/src/features/background-agent/manager.ts index 137c64850..0947b2d17 100644 --- a/src/features/background-agent/manager.ts +++ b/src/features/background-agent/manager.ts @@ -340,14 +340,16 @@ export class BackgroundManager { try { const queue = this.queuesByKey.get(key) while (queue && queue.length > 0) { - const item = queue[0] + const item = queue.shift() + if (!item) { + continue + } await this.concurrencyManager.acquire(key) if (item.task.status === "cancelled" || item.task.status === "error" || item.task.status === "interrupt") { this.rollbackPreStartDescendantReservation(item.task) this.concurrencyManager.release(key) - queue.shift() continue } @@ -363,8 +365,6 @@ export class BackgroundManager { this.concurrencyManager.release(key) } } - - queue.shift() } } finally { this.processingKeys.delete(key) @@ -411,6 +411,17 @@ export class BackgroundManager { } const sessionID = createResult.data.id + + if (task.status === "cancelled") { + await this.client.session.abort({ + path: { id: sessionID }, + }).catch((error) => { + log("[background-agent] Failed to abort cancelled pre-start session:", error) + }) + this.concurrencyManager.release(concurrencyKey) + return + } + this.settlePreStartDescendantReservation(task) subagentSessions.add(sessionID)