From 56cf16c4c50cfb098fdf332659003a92f7bd9e93 Mon Sep 17 00:00:00 2001 From: YeonGyu-Kim Date: Tue, 31 Mar 2026 15:11:05 -0700 Subject: [PATCH] refactor(background-task): stop cancelling launched tasks during session wait --- src/features/background-agent/manager.test.ts | 85 ++++++++++++++++ .../create-background-task.test.ts | 91 ++++++++++++++++- .../background-task/create-background-task.ts | 16 +-- .../background-agent-executor.test.ts | 90 ++++++++++++++++- .../background-agent-executor.ts | 11 ++- .../background-executor.test.ts | 99 ++++++++++++++++++- .../call-omo-agent/background-executor.ts | 11 ++- 7 files changed, 384 insertions(+), 19 deletions(-) diff --git a/src/features/background-agent/manager.test.ts b/src/features/background-agent/manager.test.ts index 0d812329c..7050025ba 100644 --- a/src/features/background-agent/manager.test.ts +++ b/src/features/background-agent/manager.test.ts @@ -2414,6 +2414,91 @@ describe("BackgroundManager - Non-blocking Queue Integration", () => { expect(manager.getTask(secondTask.id)?.sessionID).toBe(secondSessionID) }) + test("should keep sibling launch running when concurrent launches share a parent and the first is cancelled during session creation", async () => { + // given + const firstSessionID = "ses-first-concurrent-cancelled" + const secondSessionID = "ses-second-concurrent-survives" + 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", + } + + // when + const [firstTask, secondTask] = await Promise.all([ + manager.launch(input), + manager.launch(input), + ]) + await firstCreateStarted + + 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" diff --git a/src/tools/background-task/create-background-task.test.ts b/src/tools/background-task/create-background-task.test.ts index 2afc20a0f..a7c108ca6 100644 --- a/src/tools/background-task/create-background-task.test.ts +++ b/src/tools/background-task/create-background-task.test.ts @@ -6,7 +6,13 @@ import type { PluginInput } from "@opencode-ai/plugin" import { createBackgroundTask } from "./create-background-task" describe("createBackgroundTask", () => { - const launchMock = mock(() => Promise.resolve({ + const launchMock = mock(async (): Promise<{ + id: string + sessionID: string | null + description: string + agent: string + status: string + }> => ({ id: "test-task-id", sessionID: null, description: "Test task", @@ -32,7 +38,11 @@ describe("createBackgroundTask", () => { sessionID: "test-session", messageID: "test-message", agent: "test-agent", + directory: "/Users/yeongyu/local-workspaces/omo", + worktree: "/Users/yeongyu/local-workspaces/omo", abort: new AbortController().signal, + metadata: () => {}, + ask: async () => {}, } const testArgs = { @@ -65,4 +75,83 @@ describe("createBackgroundTask", () => { expect(result).toContain("Task entered error state") expect(result).toContain("test-task-id") }) + + test("keeps launched background task alive when parent aborts before session id resolves", async () => { + //#given - background launch should survive parent abort during session-id wait + const abortController = new AbortController() + launchMock.mockResolvedValueOnce({ + id: "test-task-id", + sessionID: null, + description: "Test task", + agent: "test-agent", + status: "pending", + }) + getTaskMock.mockImplementationOnce(() => { + abortController.abort() + return { + id: "test-task-id", + sessionID: null, + description: "Test task", + agent: "test-agent", + status: "pending", + } + }) + + //#when + const result = await tool.execute(testArgs, { + ...testContext, + abort: abortController.signal, + }) + + //#then - tool should still report successful launch instead of cancelling child task + expect(result).toContain("Background task launched successfully.") + expect(result).toContain("Task ID: test-task-id") + expect(result).not.toContain("Task aborted and cancelled while waiting for session to start") + }) + + test("keeps sibling background task alive when two tasks start concurrently", async () => { + //#given - one aborted parent call should not interrupt a sibling launch from the same parent session + const firstAbortController = new AbortController() + const secondAbortController = new AbortController() + const states = new Map([ + ["task-1", { reads: 0, abortOnFirstRead: true, sessionID: "ses-1" }], + ["task-2", { reads: 0, abortOnFirstRead: false, sessionID: "ses-2" }], + ]) + let launchCount = 0 + launchMock.mockImplementation(async () => { + launchCount += 1 + return launchCount === 1 + ? { id: "task-1", sessionID: null, description: "Task 1", agent: "test-agent", status: "pending" } + : { id: "task-2", sessionID: null, description: "Task 2", agent: "test-agent", status: "pending" } + }) + getTaskMock.mockImplementation((taskID: string) => { + const state = states.get(taskID) + if (!state) return undefined + state.reads += 1 + if (state.abortOnFirstRead && state.reads === 1) { + firstAbortController.abort() + } + return state.reads >= 2 + ? { id: taskID, sessionID: state.sessionID, description: "Task", agent: "test-agent", status: "pending" } + : { id: taskID, sessionID: null, description: "Task", agent: "test-agent", status: "pending" } + }) + + //#when + const [firstResult, secondResult] = await Promise.all([ + tool.execute(testArgs, { + ...testContext, + abort: firstAbortController.signal, + }), + tool.execute(testArgs, { + ...testContext, + abort: secondAbortController.signal, + }), + ]) + + //#then - both launches still succeed and the sibling is not marked interrupted + expect(firstResult).toContain("Background task launched successfully.") + expect(secondResult).toContain("Background task launched successfully.") + expect(secondResult).toContain("Task ID: task-2") + expect(secondResult).not.toContain("interrupt") + }) }) diff --git a/src/tools/background-task/create-background-task.ts b/src/tools/background-task/create-background-task.ts index 8f57ed763..679bc533e 100644 --- a/src/tools/background-task/create-background-task.ts +++ b/src/tools/background-task/create-background-task.ts @@ -80,16 +80,18 @@ export function createBackgroundTask( const waitStart = Date.now() let sessionId = task.sessionID while (!sessionId && Date.now() - waitStart < WAIT_FOR_SESSION_TIMEOUT_MS) { - if (ctx.abort?.aborted) { - await manager.cancelTask(task.id) - return `Task aborted and cancelled while waiting for session to start.\n\nTask ID: ${task.id}` - } - await delay(WAIT_FOR_SESSION_INTERVAL_MS) const updated = manager.getTask(task.id) - if (!updated || updated.status === "error" || updated.status === "cancelled" || updated.status === "interrupt") { - return `Task ${!updated ? "was deleted" : `entered error state`}\.\n\nTask ID: ${task.id}` + if (updated?.status === "error" || updated?.status === "cancelled" || updated?.status === "interrupt") { + return `Task ${`entered error state`}\.\n\nTask ID: ${task.id}` } sessionId = updated?.sessionID + if (sessionId) { + break + } + if (ctx.abort?.aborted) { + break + } + await delay(WAIT_FOR_SESSION_INTERVAL_MS) } const bgMeta = { diff --git a/src/tools/call-omo-agent/background-agent-executor.test.ts b/src/tools/call-omo-agent/background-agent-executor.test.ts index d27575c15..ea74b2140 100644 --- a/src/tools/call-omo-agent/background-agent-executor.test.ts +++ b/src/tools/call-omo-agent/background-agent-executor.test.ts @@ -5,7 +5,13 @@ import type { PluginInput } from "@opencode-ai/plugin" import { executeBackgroundAgent } from "./background-agent-executor" describe("executeBackgroundAgent", () => { - const launchMock = mock(() => Promise.resolve({ + const launchMock = mock(async (): Promise<{ + id: string + sessionID: string | null + description: string + agent: string + status: string + }> => ({ id: "test-task-id", sessionID: null, description: "Test task", @@ -64,4 +70,86 @@ describe("executeBackgroundAgent", () => { expect(result).toContain("interrupt") expect(result).toContain("test-task-id") }) + + test("keeps launched background task alive when parent aborts before session id resolves", async () => { + //#given - parent abort after launch should stop waiting, not fail the background task + const abortController = new AbortController() + launchMock.mockResolvedValueOnce({ + id: "test-task-id", + sessionID: null, + description: "Test task", + agent: "test-agent", + status: "pending", + }) + getTaskMock.mockImplementationOnce(() => { + abortController.abort() + return { id: "test-task-id", sessionID: null, description: "Test task", agent: "test-agent", status: "pending" } + }) + + //#when + const result = await executeBackgroundAgent( + testArgs, + { + ...testContext, + abort: abortController.signal, + }, + mockManager, + mockClient + ) + + //#then - background launch should still be reported as launched + expect(result).toContain("Background agent task launched successfully") + expect(result).toContain("Task ID: test-task-id") + expect(result).not.toContain("Task aborted while waiting for session to start") + }) + + test("keeps sibling background agent launch alive when two tasks start concurrently", async () => { + //#given - one aborted parent call should not interrupt a sibling launch from the same parent session + const firstAbortController = new AbortController() + const secondAbortController = new AbortController() + const states = new Map([ + ["task-1", { reads: 0, abortOnFirstRead: true, sessionID: "ses-1" }], + ["task-2", { reads: 0, abortOnFirstRead: false, sessionID: "ses-2" }], + ]) + let launchCount = 0 + launchMock.mockImplementation(async () => { + launchCount += 1 + return launchCount === 1 + ? { id: "task-1", sessionID: null, description: "Task 1", agent: "test-agent", status: "pending" } + : { id: "task-2", sessionID: null, description: "Task 2", agent: "test-agent", status: "pending" } + }) + getTaskMock.mockImplementation((taskID: string) => { + const state = states.get(taskID) + if (!state) return undefined + state.reads += 1 + if (state.abortOnFirstRead && state.reads === 1) { + firstAbortController.abort() + } + return state.reads >= 2 + ? { id: taskID, sessionID: state.sessionID, description: "Task", agent: "test-agent", status: "pending" } + : { id: taskID, sessionID: null, description: "Task", agent: "test-agent", status: "pending" } + }) + + //#when + const [firstResult, secondResult] = await Promise.all([ + executeBackgroundAgent( + testArgs, + { ...testContext, abort: firstAbortController.signal }, + mockManager, + mockClient, + ), + executeBackgroundAgent( + testArgs, + { ...testContext, abort: secondAbortController.signal }, + mockManager, + mockClient, + ), + ]) + + //#then - both launches still succeed and the sibling is not marked interrupted + expect(firstResult).toContain("Background agent task launched successfully") + expect(secondResult).toContain("Background agent task launched successfully") + expect(secondResult).toContain("Task ID: task-2") + expect(secondResult).not.toContain("interrupt") + }) }) diff --git a/src/tools/call-omo-agent/background-agent-executor.ts b/src/tools/call-omo-agent/background-agent-executor.ts index c09f78df3..9b82e59a3 100644 --- a/src/tools/call-omo-agent/background-agent-executor.ts +++ b/src/tools/call-omo-agent/background-agent-executor.ts @@ -52,17 +52,20 @@ export async function executeBackgroundAgent( let sessionId = task.sessionID while (!sessionId && Date.now() - waitStart < waitTimeoutMs) { - if (toolContext.abort?.aborted) { - return `Task aborted while waiting for session to start.\n\nTask ID: ${task.id}` - } const updated = manager.getTask(task.id) if (updated?.status === "error" || updated?.status === "cancelled" || updated?.status === "interrupt") { return `Task failed to start (status: ${updated.status}).\n\nTask ID: ${task.id}` } + sessionId = updated?.sessionID + if (sessionId) { + break + } + if (toolContext.abort?.aborted) { + break + } await new Promise((resolve) => { setTimeout(resolve, waitIntervalMs) }) - sessionId = manager.getTask(task.id)?.sessionID } await toolContext.metadata?.({ diff --git a/src/tools/call-omo-agent/background-executor.test.ts b/src/tools/call-omo-agent/background-executor.test.ts index 53ea45d44..da8284059 100644 --- a/src/tools/call-omo-agent/background-executor.test.ts +++ b/src/tools/call-omo-agent/background-executor.test.ts @@ -5,7 +5,13 @@ import type { PluginInput } from "@opencode-ai/plugin" import { executeBackground } from "./background-executor" describe("executeBackground", () => { - const launchMock = mock(() => Promise.resolve({ + const launchMock = mock(async (_input?: { fallbackChain?: unknown }): Promise<{ + id: string + sessionID: string | null + description: string + agent: string + status: string + }> => ({ id: "test-task-id", sessionID: null, description: "Test task", @@ -83,7 +89,96 @@ describe("executeBackground", () => { await executeBackground(testArgs, testContext, mockManager, mockClient, fallbackChain) //#then - const launchArgs = launchMock.mock.calls.at(-1)?.[0] + const latestCall = [...launchMock.mock.calls].pop() + if (!latestCall) { + throw new Error("Expected background manager launch to be called") + } + const launchArgs = latestCall[0] + if (!launchArgs) { + throw new Error("Expected launch arguments") + } expect(launchArgs.fallbackChain).toEqual(fallbackChain) }) + + test("keeps launched background task alive when parent aborts before session id resolves", async () => { + //#given - parent abort after launch should stop waiting, not fail the background task + const abortController = new AbortController() + launchMock.mockResolvedValueOnce({ + id: "test-task-id", + sessionID: null, + description: "Test task", + agent: "test-agent", + status: "pending", + }) + getTaskMock.mockImplementationOnce(() => { + abortController.abort() + return { id: "test-task-id", sessionID: null, description: "Test task", agent: "test-agent", status: "pending" } + }) + + //#when + const result = await executeBackground( + testArgs, + { + ...testContext, + abort: abortController.signal, + }, + mockManager, + mockClient + ) + + //#then - background launch should still be reported as launched + expect(result).toContain("Background agent task launched successfully") + expect(result).toContain("Task ID: test-task-id") + expect(result).not.toContain("Task aborted while waiting for session to start") + }) + + test("keeps sibling background launch alive when two tasks start concurrently", async () => { + //#given - one aborted parent call should not interrupt a sibling launch from the same parent session + const firstAbortController = new AbortController() + const secondAbortController = new AbortController() + const states = new Map([ + ["task-1", { reads: 0, abortOnFirstRead: true, sessionID: "ses-1" }], + ["task-2", { reads: 0, abortOnFirstRead: false, sessionID: "ses-2" }], + ]) + let launchCount = 0 + launchMock.mockImplementation(async () => { + launchCount += 1 + return launchCount === 1 + ? { id: "task-1", sessionID: null, description: "Task 1", agent: "test-agent", status: "pending" } + : { id: "task-2", sessionID: null, description: "Task 2", agent: "test-agent", status: "pending" } + }) + getTaskMock.mockImplementation((taskID: string) => { + const state = states.get(taskID) + if (!state) return undefined + state.reads += 1 + if (state.abortOnFirstRead && state.reads === 1) { + firstAbortController.abort() + } + return state.reads >= 2 + ? { id: taskID, sessionID: state.sessionID, description: "Task", agent: "test-agent", status: "pending" } + : { id: taskID, sessionID: null, description: "Task", agent: "test-agent", status: "pending" } + }) + + //#when + const [firstResult, secondResult] = await Promise.all([ + executeBackground( + testArgs, + { ...testContext, abort: firstAbortController.signal }, + mockManager, + mockClient, + ), + executeBackground( + testArgs, + { ...testContext, abort: secondAbortController.signal }, + mockManager, + mockClient, + ), + ]) + + //#then - both launches still succeed and the sibling is not marked interrupted + expect(firstResult).toContain("Background agent task launched successfully") + expect(secondResult).toContain("Background agent task launched successfully") + expect(secondResult).toContain("Task ID: task-2") + expect(secondResult).not.toContain("interrupt") + }) }) diff --git a/src/tools/call-omo-agent/background-executor.ts b/src/tools/call-omo-agent/background-executor.ts index 13f6f6d21..fac45e1ce 100644 --- a/src/tools/call-omo-agent/background-executor.ts +++ b/src/tools/call-omo-agent/background-executor.ts @@ -61,15 +61,18 @@ export async function executeBackground( const waitStart = Date.now() let sessionId = task.sessionID while (!sessionId && Date.now() - waitStart < WAIT_FOR_SESSION_TIMEOUT_MS) { - if (toolContext.abort?.aborted) { - return `Task aborted while waiting for session to start.\n\nTask ID: ${task.id}` - } const updated = manager.getTask(task.id) if (updated?.status === "error" || updated?.status === "cancelled" || updated?.status === "interrupt") { return `Task failed to start (status: ${updated.status}).\n\nTask ID: ${task.id}` } + sessionId = updated?.sessionID + if (sessionId) { + break + } + if (toolContext.abort?.aborted) { + break + } await new Promise(resolve => setTimeout(resolve, WAIT_FOR_SESSION_INTERVAL_MS)) - sessionId = manager.getTask(task.id)?.sessionID } await toolContext.metadata?.({