diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index d408a8b72..1a997916d 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -11,6 +11,7 @@ - Fixed concise `history://` transcript rendering for `find` and `search` so scoped `paths` arguments are visible instead of being hidden behind JSON fallback output or omitted when a search `pattern` is present. ([#3482](https://github.com/can1357/oh-my-pi/issues/3482)) - Fixed manual `/compact` leaving `session.isCompacting` false while active-turn abort teardown awaited, so the first steer/follow-up typed during compaction startup now routes through the compaction queue instead of being lost. ([#3485](https://github.com/can1357/oh-my-pi/issues/3485)) - Fixed ollama-cloud task/subagent fan-out exceeding the provider's three-request concurrency cap by adding a provider-specific subagent limiter, and let configured task/smol/advisor model roles inherit the default retry fallback chain when they do not define their own chain. ([#3464](https://github.com/can1357/oh-my-pi/issues/3464)) +- Fixed the per-provider subagent concurrency limiter (e.g. `providers.ollama-cloud.maxConcurrency`) being replaced with a fresh semaphore whenever the configured limit changed, which orphaned the in-flight slots on the old instance and let a runtime or mixed limit value exceed the cap. The limiter now resizes a single shared semaphore in place — raising the ceiling admits queued waiters immediately, lowering it drains in-flight holders without admitting past the new cap. ([#3464](https://github.com/can1357/oh-my-pi/issues/3464)) ## [16.1.19] - 2026-06-25 diff --git a/packages/coding-agent/src/task/executor.ts b/packages/coding-agent/src/task/executor.ts index de686e6b4..565e95080 100644 --- a/packages/coding-agent/src/task/executor.ts +++ b/packages/coding-agent/src/task/executor.ts @@ -205,19 +205,35 @@ interface ProviderSemaphoreEntry { const providerSemaphores = new Map(); -function getProviderConcurrencyLimit(settings: Settings, provider: string): number { +/** + * Resolve the configured concurrency ceiling for a provider, or `undefined` + * when the provider has no cap concept at all. A configured value `<= 0` means + * "unlimited" and maps to `Infinity` — still a tracked ceiling, so every run + * holds a slot and a later finite resize counts work started while unlimited. + */ +function getProviderConcurrencyLimit(settings: Settings, provider: string): number | undefined { const settingPath = PROVIDER_MAX_CONCURRENCY_SETTINGS[provider]; - if (!settingPath) return 0; + if (!settingPath) return undefined; const raw = settings.get(settingPath); const limit = Number.isFinite(raw) ? Math.trunc(raw) : 0; - return limit > 0 ? limit : 0; + return limit > 0 ? limit : Number.POSITIVE_INFINITY; } function getProviderSemaphore(settings: Settings, provider: string): Semaphore | undefined { const limit = getProviderConcurrencyLimit(settings, provider); - if (limit <= 0) return undefined; + if (limit === undefined) return undefined; + // Always hand out (and acquire on) the single shared limiter, even when + // unlimited (Infinity). Resizing it in place — rather than replacing it — + // keeps every in-flight slot counted, so a runtime or mixed limit change can + // never push concurrency past the cap (issue #3464 review feedback). const existing = providerSemaphores.get(provider); - if (existing?.limit === limit) return existing.semaphore; + if (existing) { + if (existing.limit !== limit) { + existing.limit = limit; + existing.semaphore.resize(limit); + } + return existing.semaphore; + } const semaphore = new Semaphore(limit); providerSemaphores.set(provider, { limit, semaphore }); return semaphore; diff --git a/packages/coding-agent/src/task/parallel.ts b/packages/coding-agent/src/task/parallel.ts index a07dd38c0..86fa4670d 100644 --- a/packages/coding-agent/src/task/parallel.ts +++ b/packages/coding-agent/src/task/parallel.ts @@ -135,11 +135,33 @@ export class Semaphore { } release(): void { - const next = this.#queue.shift(); - if (next) { + if (this.#current > 0) this.#current--; + // Admit the next waiter only if we are under the (possibly just-lowered) ceiling. + if (this.#current < this.#max) { + const next = this.#queue.shift(); + if (next) { + this.#current++; + next(); + } + } + } + + /** + * Adjust the maximum concurrency in place. Raising the ceiling immediately + * admits queued waiters that now fit; lowering it lets in-flight holders + * drain naturally (new acquires keep blocking until `#current` falls below + * the new max). Resizing the single shared instance — instead of replacing + * it — keeps in-flight slots counted, so a runtime or mixed limit change can + * never push concurrency past the cap (issue #3464 review feedback). + */ + resize(max: number): void { + const normalizedMax = Number.isFinite(max) ? Math.trunc(max) : 0; + this.#max = normalizedMax > 0 ? normalizedMax : Number.POSITIVE_INFINITY; + while (this.#current < this.#max) { + const next = this.#queue.shift(); + if (!next) break; + this.#current++; next(); - } else { - this.#current--; } } } 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 0c923fd13..5932dcb14 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 @@ -242,4 +242,74 @@ describe("issue #3464: ollama-cloud task backoff", () => { semaphore.release(); await next; }); + + it("raises the ceiling in place and admits queued waiters without a release", async () => { + const semaphore = new Semaphore(1); + await semaphore.acquire(); + const admitted: number[] = []; + const w1 = (async () => { + await semaphore.acquire(); + admitted.push(1); + })(); + const w2 = (async () => { + await semaphore.acquire(); + admitted.push(2); + })(); + await Bun.sleep(0); + expect(admitted).toEqual([]); + + semaphore.resize(3); + await Bun.sleep(0); + expect(admitted).toEqual([1, 2]); + await Promise.all([w1, w2]); + }); + + it("lowers the ceiling without admitting waiters past the new cap", async () => { + const semaphore = new Semaphore(3); + await semaphore.acquire(); + await semaphore.acquire(); + await semaphore.acquire(); + let admitted = false; + const waiter = (async () => { + await semaphore.acquire(); + admitted = true; + })(); + await Bun.sleep(0); + expect(admitted).toBe(false); + + semaphore.resize(1); + semaphore.release(); + await Bun.sleep(0); + expect(admitted).toBe(false); + semaphore.release(); + await Bun.sleep(0); + expect(admitted).toBe(false); + semaphore.release(); + await Bun.sleep(0); + expect(admitted).toBe(true); + await waiter; + semaphore.release(); + }); + + it("counts holders acquired while unlimited after a finite cap is re-enabled", async () => { + const semaphore = new Semaphore(0); // unlimited + await semaphore.acquire(); + await semaphore.acquire(); // two holders counted despite being unlimited + let admitted = false; + semaphore.resize(1); // re-enable a finite cap below the in-flight count + const waiter = (async () => { + await semaphore.acquire(); + admitted = true; + })(); + await Bun.sleep(0); + expect(admitted).toBe(false); + semaphore.release(); + await Bun.sleep(0); + expect(admitted).toBe(false); + semaphore.release(); + await Bun.sleep(0); + expect(admitted).toBe(true); + await waiter; + semaphore.release(); + }); });