fix(background-agent): prevent queue item loss on concurrent cancel and guard against cancelled task resurrection
This commit is contained in:
@@ -2195,6 +2195,170 @@ describe("BackgroundManager - Non-blocking Queue Integration", () => {
|
|||||||
// then
|
// then
|
||||||
expect(retryTask.status).toBe("pending")
|
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<void>((resolve) => {
|
||||||
|
resolveFirstCreateStarted = resolve
|
||||||
|
})
|
||||||
|
const secondPromptAsyncStarted = new Promise<void>((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<never>((_, 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<void>((resolve) => {
|
||||||
|
resolveCreateStarted = resolve
|
||||||
|
})
|
||||||
|
const abortCalled = new Promise<void>((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<never>((_, 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", () => {
|
describe("pending task can be cancelled", () => {
|
||||||
|
|||||||
@@ -340,14 +340,16 @@ export class BackgroundManager {
|
|||||||
try {
|
try {
|
||||||
const queue = this.queuesByKey.get(key)
|
const queue = this.queuesByKey.get(key)
|
||||||
while (queue && queue.length > 0) {
|
while (queue && queue.length > 0) {
|
||||||
const item = queue[0]
|
const item = queue.shift()
|
||||||
|
if (!item) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
await this.concurrencyManager.acquire(key)
|
await this.concurrencyManager.acquire(key)
|
||||||
|
|
||||||
if (item.task.status === "cancelled" || item.task.status === "error" || item.task.status === "interrupt") {
|
if (item.task.status === "cancelled" || item.task.status === "error" || item.task.status === "interrupt") {
|
||||||
this.rollbackPreStartDescendantReservation(item.task)
|
this.rollbackPreStartDescendantReservation(item.task)
|
||||||
this.concurrencyManager.release(key)
|
this.concurrencyManager.release(key)
|
||||||
queue.shift()
|
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -363,8 +365,6 @@ export class BackgroundManager {
|
|||||||
this.concurrencyManager.release(key)
|
this.concurrencyManager.release(key)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
queue.shift()
|
|
||||||
}
|
}
|
||||||
} finally {
|
} finally {
|
||||||
this.processingKeys.delete(key)
|
this.processingKeys.delete(key)
|
||||||
@@ -411,6 +411,17 @@ export class BackgroundManager {
|
|||||||
}
|
}
|
||||||
|
|
||||||
const sessionID = createResult.data.id
|
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)
|
this.settlePreStartDescendantReservation(task)
|
||||||
subagentSessions.add(sessionID)
|
subagentSessions.add(sessionID)
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user