diff --git a/packages/coding-agent/src/async/job-manager.ts b/packages/coding-agent/src/async/job-manager.ts index fa098f9bb..6e36f22b8 100644 --- a/packages/coding-agent/src/async/job-manager.ts +++ b/packages/coding-agent/src/async/job-manager.ts @@ -220,9 +220,7 @@ export class AsyncJobManager { } getDeliveryState(filter?: AsyncJobFilter): AsyncJobDeliveryState { - const deliveries = filter?.ownerId - ? this.#deliveries.filter(delivery => this.#jobs.get(delivery.jobId)?.ownerId === filter.ownerId) - : this.#deliveries; + const deliveries = this.#filterDeliveries(filter); const nextRetryAt = deliveries.reduce((next, delivery) => { if (next === undefined) return delivery.nextAttemptAt; return Math.min(next, delivery.nextAttemptAt); @@ -300,6 +298,12 @@ export class AsyncJobManager { const deadline = hasDeadline ? Date.now() + Math.max(timeoutMs, 0) : Number.POSITIVE_INFINITY; while (this.hasPendingDeliveries(filter)) { + if (filter?.ownerId) { + const delivered = await this.#deliverNextFiltered(filter, deadline); + if (delivered) continue; + return false; + } + this.#ensureDeliveryLoop(); const loop = this.#deliveryLoop; if (!loop) { @@ -392,6 +396,53 @@ export class AsyncJobManager { this.#evictionTimers.clear(); } + #filterDeliveries(filter?: AsyncJobFilter): AsyncJobDelivery[] { + const ownerId = filter?.ownerId; + if (!ownerId) return this.#deliveries; + return this.#deliveries.filter(delivery => this.#jobs.get(delivery.jobId)?.ownerId === ownerId); + } + + async #deliverNextFiltered(filter: AsyncJobFilter, deadline: number): Promise { + let selected: AsyncJobDelivery | undefined; + for (const delivery of this.#deliveries) { + if (this.#jobs.get(delivery.jobId)?.ownerId !== filter.ownerId) continue; + if (this.isDeliverySuppressed(delivery.jobId)) continue; + if (!selected || delivery.nextAttemptAt < selected.nextAttemptAt) { + selected = delivery; + } + } + if (!selected) return true; + + const now = Date.now(); + if (selected.nextAttemptAt > now) { + if (selected.nextAttemptAt > deadline) return false; + await Bun.sleep(selected.nextAttemptAt - now); + } + + const index = this.#deliveries.indexOf(selected); + if (index === -1) return true; + + try { + await this.#onJobComplete(selected.jobId, selected.text, this.#jobs.get(selected.jobId)); + this.#deliveries.splice(index, 1); + } catch (error) { + selected.attempt += 1; + selected.lastError = error instanceof Error ? error.message : String(error); + selected.nextAttemptAt = Date.now() + this.#getRetryDelay(selected.attempt); + this.#deliveries.splice(index, 1); + if (!this.isDeliverySuppressed(selected.jobId)) { + this.#deliveries.push(selected); + } + logger.warn("Async job completion delivery failed", { + jobId: selected.jobId, + attempt: selected.attempt, + nextRetryAt: selected.nextAttemptAt, + error: selected.lastError, + }); + } + return true; + } + isDeliverySuppressed(jobId: string): boolean { return this.#suppressedDeliveries.has(jobId) || this.#watchedJobs.has(jobId); } diff --git a/packages/coding-agent/test/agent-session-concurrent.test.ts b/packages/coding-agent/test/agent-session-concurrent.test.ts index 060dfcd4a..3bea6ef5a 100644 --- a/packages/coding-agent/test/agent-session-concurrent.test.ts +++ b/packages/coding-agent/test/agent-session-concurrent.test.ts @@ -383,248 +383,6 @@ describe("AgentSession concurrent prompt guard", () => { }), ).toBe(true); }); - const model = getBundledModel("anthropic", "claude-sonnet-4-5")!; - const mock = createMockModel({ handler: () => ({ content: ["Done"] }) }); - const agent = new Agent({ - getApiKey: () => "test-key", - // Regression: a subscriber that fires the next prompt synchronously from the - // agent_end listener (the shape every wire transport ends up in — rpc-mode - // stdout subscriber, ACP bridge, Cursor exec) must not collide with the - // outgoing turn's still-unwinding in-flight bookkeeping. Before the wire-level - // agent_end was deferred until #promptInFlightCount drops to 0, the - // subscriber observed agent_end while Session.isStreaming was still true (the - // agent's own `isStreaming` had flipped, but #promptWithMessage's finally had - // not yet decremented the prompt-in-flight counter), and the next prompt - // threw AgentBusyError. Surfaced as `RpcCommandError: prompt: Agent is - // already processing` from omp-rpc clients (robomp triage reminder path). - it("subscriber may prompt() synchronously from agent_end without AgentBusyError", async () => { - const model = getBundledModel("anthropic", "claude-sonnet-4-5")!; - const mock = createMockModel({ handler: () => ({ content: ["Done"] }) }); - const agent = new Agent({ - getApiKey: () => "test-key", - initialState: { model, systemPrompt: ["Test"], tools: [] }, - streamFn: mock.stream, - }); - - const sessionManager = SessionManager.inMemory(); - const settings = Settings.isolated(); - const authStorage = await AuthStorage.create(path.join(tempDir, "testauth.db")); - authStorages.push(authStorage); - const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models.yml")); - authStorage.setRuntimeApiKey("anthropic", "test-key"); - - session = new AgentSession({ agent, sessionManager, settings, modelRegistry }); - - const observedIsStreamingAtAgentEnd: boolean[] = []; - const reentrantPromptResults: Array<"resolved" | { error: string }> = []; - let reentrantPrompted = false; - - session.subscribe(event => { - if (event.type !== "agent_end") return; - observedIsStreamingAtAgentEnd.push(session.isStreaming); - if (reentrantPrompted) return; - reentrantPrompted = true; - void session - .prompt("Second message") - .then(() => reentrantPromptResults.push("resolved")) - .catch((err: Error) => reentrantPromptResults.push({ error: err.message })); - }); - - await session.prompt("First message"); - await waitFor(() => reentrantPromptResults.length > 0, 2000); - await session.waitForIdle(); - - expect(observedIsStreamingAtAgentEnd).not.toContain(true); - expect(reentrantPromptResults).toEqual(["resolved"]); - }); - - it("queues idle ACP client-triggered custom messages instead of starting an ownerless turn", async () => { - const model = getBundledModel("anthropic", "claude-sonnet-4-5")!; - const mock = createMockModel({ handler: () => ({ content: ["Done"] }) }); - const agent = new Agent({ - getApiKey: () => "test-key", - initialState: { - model, - systemPrompt: ["Test"], - tools: [], - }, - convertToLlm, - streamFn: mock.stream, - }); - - const sessionManager = SessionManager.inMemory(); - const settings = Settings.isolated(); - const authStorage = await AuthStorage.create(":memory:"); - authStorages.push(authStorage); - const modelRegistry = new ModelRegistry(authStorage); - authStorage.setRuntimeApiKey("anthropic", "test-key"); - - session = new AgentSession({ - agent, - sessionManager, - settings, - modelRegistry, - }); - session.setClientBridge({ - capabilities: {}, - deferAgentInitiatedTurns: true, - }); - - await session.prompt("First message"); - expect(session.isStreaming).toBe(false); - const callsAfterFirstPrompt = mock.calls.length; - - await session.sendCustomMessage( - { - customType: "async-result", - content: "Background result", - display: true, - attribution: "agent", - }, - { deliverAs: "followUp", triggerTurn: true }, - ); - - expect(mock.calls).toHaveLength(callsAfterFirstPrompt); - expect(session.isStreaming).toBe(false); - - await session.prompt("Next user prompt"); - await session.dispose(); - session = undefined as unknown as AgentSession; - expect(mock.calls).toHaveLength(callsAfterFirstPrompt + 1); - expect( - mock.calls.at(-1)?.context.messages.some(message => { - if (typeof message.content === "string") { - return message.content.includes("Background result"); - } - - return message.content.some( - content => content.type === "text" && content.text.includes("Background result"), - ); - }), - ).toBe(true); - }); - streamFn: mock.stream, - }); - - const sessionManager = SessionManager.inMemory(); - const settings = Settings.isolated(); - // Regression: a subscriber that fires the next prompt synchronously from the - // agent_end listener (the shape every wire transport ends up in — rpc-mode - // stdout subscriber, ACP bridge, Cursor exec) must not collide with the - // outgoing turn's still-unwinding in-flight bookkeeping. Before the wire-level - // agent_end was deferred until #promptInFlightCount drops to 0, the - // subscriber observed agent_end while Session.isStreaming was still true (the - // agent's own `isStreaming` had flipped, but #promptWithMessage's finally had - // not yet decremented the prompt-in-flight counter), and the next prompt - // threw AgentBusyError. Surfaced as `RpcCommandError: prompt: Agent is - // already processing` from omp-rpc clients (robomp triage reminder path). - it("subscriber may prompt() synchronously from agent_end without AgentBusyError", async () => { - const model = getBundledModel("anthropic", "claude-sonnet-4-5")!; - const mock = createMockModel({ handler: () => ({ content: ["Done"] }) }); - const agent = new Agent({ - getApiKey: () => "test-key", - initialState: { model, systemPrompt: ["Test"], tools: [] }, - streamFn: mock.stream, - }); - - const sessionManager = SessionManager.inMemory(); - const settings = Settings.isolated(); - const authStorage = await AuthStorage.create(path.join(tempDir, "testauth.db")); - authStorages.push(authStorage); - const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models.yml")); - authStorage.setRuntimeApiKey("anthropic", "test-key"); - - session = new AgentSession({ agent, sessionManager, settings, modelRegistry }); - - const observedIsStreamingAtAgentEnd: boolean[] = []; - const reentrantPromptResults: Array<"resolved" | { error: string }> = []; - let reentrantPrompted = false; - - session.subscribe(event => { - if (event.type !== "agent_end") return; - observedIsStreamingAtAgentEnd.push(session.isStreaming); - if (reentrantPrompted) return; - reentrantPrompted = true; - void session - .prompt("Second message") - .then(() => reentrantPromptResults.push("resolved")) - .catch((err: Error) => reentrantPromptResults.push({ error: err.message })); - }); - - await session.prompt("First message"); - await waitFor(() => reentrantPromptResults.length > 0, 2000); - await session.waitForIdle(); - - expect(observedIsStreamingAtAgentEnd).not.toContain(true); - expect(reentrantPromptResults).toEqual(["resolved"]); - }); - - it("queues idle ACP client-triggered custom messages instead of starting an ownerless turn", async () => { - const model = getBundledModel("anthropic", "claude-sonnet-4-5")!; - const mock = createMockModel({ handler: () => ({ content: ["Done"] }) }); - const agent = new Agent({ - getApiKey: () => "test-key", - initialState: { - model, - systemPrompt: ["Test"], - tools: [], - }, - convertToLlm, - streamFn: mock.stream, - }); - - const sessionManager = SessionManager.inMemory(); - const settings = Settings.isolated(); - const authStorage = await AuthStorage.create(":memory:"); - authStorages.push(authStorage); - const modelRegistry = new ModelRegistry(authStorage); - authStorage.setRuntimeApiKey("anthropic", "test-key"); - - session = new AgentSession({ - agent, - sessionManager, - settings, - modelRegistry, - }); - session.setClientBridge({ - capabilities: {}, - deferAgentInitiatedTurns: true, - }); - - await session.prompt("First message"); - expect(session.isStreaming).toBe(false); - const callsAfterFirstPrompt = mock.calls.length; - - await session.sendCustomMessage( - { - customType: "async-result", - content: "Background result", - display: true, - attribution: "agent", - }, - { deliverAs: "followUp", triggerTurn: true }, - ); - - expect(mock.calls).toHaveLength(callsAfterFirstPrompt); - expect(session.isStreaming).toBe(false); - - await session.prompt("Next user prompt"); - await session.dispose(); - session = undefined as unknown as AgentSession; - expect(mock.calls).toHaveLength(callsAfterFirstPrompt + 1); - expect( - mock.calls.at(-1)?.context.messages.some(message => { - if (typeof message.content === "string") { - return message.content.includes("Background result"); - } - - return message.content.some( - content => content.type === "text" && content.text.includes("Background result"), - ); - }), - ).toBe(true); - }); - }); }); describe("AgentSession TTSR resume gate", () => { diff --git a/packages/coding-agent/test/async-job-manager.test.ts b/packages/coding-agent/test/async-job-manager.test.ts index ac1b6d9eb..b8a4a7c88 100644 --- a/packages/coding-agent/test/async-job-manager.test.ts +++ b/packages/coding-agent/test/async-job-manager.test.ts @@ -229,6 +229,38 @@ describe("AsyncJobManager", () => { expect(manager.hasPendingDeliveries()).toBe(false); }); + test("scoped delivery drain returns once matching owner deliveries finish", async () => { + let firstOwnerAttempts = 0; + const secondOwnerCompletions: Array<{ jobId: string; text: string }> = []; + const manager = new AsyncJobManager({ + onJobComplete: async (jobId, text, job) => { + if (job?.ownerId === "0-Main") { + firstOwnerAttempts++; + throw new Error("first owner delivery retry"); + } + secondOwnerCompletions.push({ jobId, text }); + }, + }); + + manager.register("task", "main job", async () => "main result", { ownerId: "0-Main" }); + const targetJobId = manager.register("task", "subagent job", async () => "subagent result", { + ownerId: "3-AuthLoader", + }); + await manager.waitForAll(); + const firstAttemptDeadline = Date.now() + 2_000; + while (firstOwnerAttempts === 0) { + if (Date.now() >= firstAttemptDeadline) throw new Error("Timed out waiting for first owner delivery attempt"); + await Bun.sleep(5); + } + + const drained = await manager.drainDeliveries({ timeoutMs: 50, filter: { ownerId: "3-AuthLoader" } }); + + expect(drained).toBe(true); + expect(secondOwnerCompletions).toEqual([{ jobId: targetJobId, text: "subagent result" }]); + expect(manager.hasPendingDeliveries({ ownerId: "3-AuthLoader" })).toBe(false); + expect(manager.hasPendingDeliveries({ ownerId: "0-Main" })).toBe(true); + }); + test("cancelAll with ownerId only cancels matching jobs", async () => { const manager = new AsyncJobManager({ onJobComplete: async () => {},