From 7e90e4d081bee9c380ff0ab0a476bcf48edfbef8 Mon Sep 17 00:00:00 2001 From: roboomp Date: Thu, 25 Jun 2026 12:14:45 +0000 Subject: [PATCH] fix(agent): released ollama-cloud semaphore slot when waiter aborts Semaphore.acquire now accepts an AbortSignal so a queued waiter that is cancelled (parent task abort, wall-clock budget elapsing) removes itself from the wait queue instead of being resolved by the next release. The provider semaphore in runSubprocess passes the run's abortSignal through, preventing aborted ollama-cloud subagents from permanently draining the provider concurrency budget. Fixes #3464 --- packages/coding-agent/src/task/executor.ts | 3 +- packages/coding-agent/src/task/parallel.ts | 36 +++++++++++++++++-- ...sue-3464-ollama-cloud-task-backoff.test.ts | 25 +++++++++++++ 3 files changed, 59 insertions(+), 5 deletions(-) diff --git a/packages/coding-agent/src/task/executor.ts b/packages/coding-agent/src/task/executor.ts index 41cc0a2e4..de686e6b4 100644 --- a/packages/coding-agent/src/task/executor.ts +++ b/packages/coding-agent/src/task/executor.ts @@ -1991,9 +1991,8 @@ export async function runSubprocess(options: ExecutorOptions): Promise 0 ? normalizedMax : Number.POSITIVE_INFINITY; } - async acquire(): Promise { + /** + * Resolves when a slot is available. Pass an `AbortSignal` so callers that + * stop waiting (parent task cancelled, wall-clock budget elapsed) also stop + * occupying a queue slot — otherwise a later `release()` would resolve the + * abandoned waiter, permanently shrinking effective concurrency for the + * remaining lifetime of the process (issue #3464 review feedback). + */ + async acquire(signal?: AbortSignal): Promise { + if (signal?.aborted) { + throw semaphoreAbortReason(signal); + } if (this.#current < this.#max) { this.#current++; return; } - const { promise, resolve } = Promise.withResolvers(); - this.#queue.push(resolve); + const { promise, resolve, reject } = Promise.withResolvers(); + const queue = this.#queue; + let waiter: () => void = resolve; + if (signal) { + const onAbort = () => { + const index = queue.indexOf(waiter); + if (index >= 0) queue.splice(index, 1); + reject(semaphoreAbortReason(signal)); + }; + waiter = () => { + signal.removeEventListener("abort", onAbort); + resolve(); + }; + signal.addEventListener("abort", onAbort, { once: true }); + } + queue.push(waiter); return promise; } @@ -119,3 +143,9 @@ export class Semaphore { } } } + +function semaphoreAbortReason(signal: AbortSignal): unknown { + const reason = signal.reason; + if (reason !== undefined) return reason; + return new Error("Semaphore acquire aborted"); +} diff --git a/packages/coding-agent/test/issue-3464-ollama-cloud-task-backoff.test.ts b/packages/coding-agent/test/issue-3464-ollama-cloud-task-backoff.test.ts index 06e27adf2..0c923fd13 100644 --- a/packages/coding-agent/test/issue-3464-ollama-cloud-task-backoff.test.ts +++ b/packages/coding-agent/test/issue-3464-ollama-cloud-task-backoff.test.ts @@ -17,6 +17,7 @@ import { import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage"; import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager"; import { runSubprocess } from "@oh-my-pi/pi-coding-agent/task/executor"; +import { Semaphore } from "@oh-my-pi/pi-coding-agent/task/parallel"; import type { AgentDefinition } from "@oh-my-pi/pi-coding-agent/task/types"; import { EventBus } from "@oh-my-pi/pi-coding-agent/utils/event-bus"; import { TempDir } from "@oh-my-pi/pi-utils"; @@ -217,4 +218,28 @@ describe("issue #3464: ollama-cloud task backoff", () => { gates.get("CloudTwo")?.resolve(); await second; }); + + it("frees a queued slot when its acquire waiter is aborted", async () => { + const semaphore = new Semaphore(1); + await semaphore.acquire(); + const controller = new AbortController(); + const aborted = semaphore.acquire(controller.signal); + controller.abort(); + await aborted.then( + () => { + throw new Error("Aborted semaphore.acquire should reject"); + }, + () => {}, + ); + + const nextStarted = deferred(); + const next = (async () => { + await semaphore.acquire(); + nextStarted.resolve(); + })(); + semaphore.release(); + await nextStarted.promise; + semaphore.release(); + await next; + }); });