From ecd6f608e58593e209decfed719f85526319990d Mon Sep 17 00:00:00 2001 From: can1357 Date: Thu, 25 Jun 2026 19:04:50 +0200 Subject: [PATCH] fix(agent): resize provider concurrency limiter in place instead of replacing it The per-provider subagent limiter (providers.ollama-cloud.maxConcurrency) created a fresh Semaphore whenever the configured limit changed, orphaning in-flight slots on the old instance so a runtime or mixed limit value could exceed the cap. getProviderSemaphore now always hands out one shared limiter (Infinity when unlimited, so every run is still counted) and resizes it in place. Semaphore.release() decrements before admitting, and the new Semaphore.resize() raises the ceiling by admitting queued waiters while lowering it drains in-flight holders without admitting past the new cap. Refs #3464 --- packages/coding-agent/CHANGELOG.md | 1 + packages/coding-agent/src/task/executor.ts | 26 +++++-- packages/coding-agent/src/task/parallel.ts | 30 ++++++-- ...sue-3464-ollama-cloud-task-backoff.test.ts | 70 +++++++++++++++++++ 4 files changed, 118 insertions(+), 9 deletions(-) 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(); + }); });