From 9599689af3d75f64bd535adcd1e5dad3867069ae Mon Sep 17 00:00:00 2001 From: roboomp Date: Sun, 28 Jun 2026 21:11:48 +0000 Subject: [PATCH] fix(compaction): routed summary oneshots through provider cap MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Compaction issued summarization HTTP requests via the default `completeSimple` transport, bypassing `wrapStreamFnWithProviderConcurrency` which was only wired into `Agent.streamFn` / `sideStreamFn`. With the per-LLM-turn bracket introduced in this PR, multiple ollama-cloud subagents that auto- or manually compact could issue uncapped summary requests in parallel and exceed `providers.ollama-cloud.maxConcurrency` (chatgpt-codex review on #3751). Added an optional `completeImpl` transport override to `SummaryOptions` and `GenerateBranchSummaryOptions` and threaded it into every `instrumentedCompleteSimple` call in compaction + branch-summarization. Wired `AgentSession.#compactWithFallbackModel` and the `generateBranchSummary` caller to route through `#sideStreamFn` — the same limiter-wrapped transport the handoff path already uses. Pinned with a coding-agent regression that drives `compact()` end to end against the wrapped sideStreamFn at maxConcurrency=1 and asserts peak in-flight stays at 1 across a concurrent unrelated side request. Fixes #3749 --- .../src/compaction/branch-summarization.ts | 15 +- packages/agent/src/compaction/compaction.ts | 21 +- .../coding-agent/src/session/agent-session.ts | 18 ++ ...issue-3751-compaction-provider-cap.test.ts | 196 ++++++++++++++++++ 4 files changed, 245 insertions(+), 5 deletions(-) create mode 100644 packages/coding-agent/test/issue-3751-compaction-provider-cap.test.ts diff --git a/packages/agent/src/compaction/branch-summarization.ts b/packages/agent/src/compaction/branch-summarization.ts index 11b04ebd4..015e2fd54 100644 --- a/packages/agent/src/compaction/branch-summarization.ts +++ b/packages/agent/src/compaction/branch-summarization.ts @@ -5,7 +5,7 @@ * a summary of the branch being left so context isn't lost. */ -import type { ApiKey, Model } from "@oh-my-pi/pi-ai"; +import type { Api, ApiKey, AssistantMessage, Context, Model, SimpleStreamOptions } from "@oh-my-pi/pi-ai"; import { preferredDialect } from "@oh-my-pi/pi-catalog/identity"; import { prompt } from "@oh-my-pi/pi-utils"; import { type AgentTelemetry, instrumentedCompleteSimple } from "../telemetry"; @@ -88,6 +88,17 @@ export interface GenerateBranchSummaryOptions { * wrapped in an OTEL chat span tagged with `pi.gen_ai.oneshot.kind = "branch_summary"`. */ telemetry?: AgentTelemetry; + /** + * Optional completion transport override (same contract as + * {@link SummaryOptions.completeImpl}). Lets the host route the branch + * summary HTTP request through its provider-concurrency limiter instead + * of the default `completeSimple` transport. + */ + completeImpl?: ( + model: Model, + ctx: Context, + options: SimpleStreamOptions, + ) => Promise; } // ============================================================================ @@ -310,7 +321,7 @@ export async function generateBranchSummary( model, { systemPrompt: [SUMMARIZATION_SYSTEM_PROMPT], messages: summarizationMessages }, { apiKey, signal, maxTokens: 2048, metadata }, - { telemetry: options.telemetry, oneshotKind: "branch_summary" }, + { telemetry: options.telemetry, oneshotKind: "branch_summary", completeImpl: options.completeImpl }, ); // Check if aborted or errored diff --git a/packages/agent/src/compaction/compaction.ts b/packages/agent/src/compaction/compaction.ts index 4bac74dfb..da31a02a7 100644 --- a/packages/agent/src/compaction/compaction.ts +++ b/packages/agent/src/compaction/compaction.ts @@ -680,6 +680,19 @@ export interface SummaryOptions { tools?: Tool[]; /** Optional fetch implementation threaded into remote compaction calls. */ fetch?: FetchImpl; + /** + * Optional completion transport override for host-level request wrappers + * (e.g. the coding-agent provider-concurrency limiter). When provided, + * every local summarization oneshot (`generateSummary`, + * `generateTurnPrefixSummary`, `generateShortSummary`) routes through it + * instead of the default `completeSimple`, so cap policies enforced on + * the live agent turn also bracket compaction HTTP requests. + */ + completeImpl?: ( + model: Model, + ctx: Context, + options: SimpleStreamOptions, + ) => Promise; } function formatPreviousSnapcompactArchive(archiveText: string): string { @@ -769,7 +782,7 @@ export async function generateSummary( initiatorOverride: options?.initiatorOverride, metadata: options?.metadata, }, - { telemetry: options?.telemetry, oneshotKind: "compaction_summary" }, + { telemetry: options?.telemetry, oneshotKind: "compaction_summary", completeImpl: options?.completeImpl }, ); if (response.stopReason === "error") { @@ -960,7 +973,7 @@ async function generateShortSummary( initiatorOverride: options?.initiatorOverride, metadata: options?.metadata, }, - { telemetry: options?.telemetry, oneshotKind: "compaction_short_summary" }, + { telemetry: options?.telemetry, oneshotKind: "compaction_short_summary", completeImpl: options?.completeImpl }, ); if (response.stopReason === "error") { @@ -1229,6 +1242,7 @@ export async function compact( promptCacheKey: options?.promptCacheKey, tools: options?.tools, fetch: options?.fetch, + completeImpl: options?.completeImpl, }; const previousSnapcompactArchive = snapcompact.getPreservedArchive(previousPreserveData); @@ -1413,6 +1427,7 @@ export async function compact( // resolves its own reasoning via resolveCompactionEffort. thinkingLevel: options?.thinkingLevel, fetch: summaryOptions.fetch, + completeImpl: summaryOptions.completeImpl, }); // Compute file lists and append to summary @@ -1476,7 +1491,7 @@ async function generateTurnPrefixSummary( initiatorOverride: options?.initiatorOverride, metadata: options?.metadata, }, - { telemetry: options?.telemetry, oneshotKind: "compaction_turn_prefix" }, + { telemetry: options?.telemetry, oneshotKind: "compaction_turn_prefix", completeImpl: options?.completeImpl }, ); if (response.stopReason === "error") { diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index da99737a9..c5ba40ab3 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -10726,6 +10726,18 @@ export class AgentSession { tools: this.agent.state.tools, sessionId: this.sessionId, promptCacheKey: this.sessionId, + // Route every summarization HTTP request through the + // session's side-stream transport so the provider + // concurrency cap (e.g. providers.ollama-cloud.maxConcurrency) + // brackets compaction the same way it brackets the live + // agent turn — without this, multiple ollama-cloud + // subagents auto/manually compacting issued uncapped + // summary requests in parallel (chatgpt-codex review on + // #3751). + completeImpl: async (requestModel, requestContext, requestOptions) => { + const stream = await this.#sideStreamFn(requestModel, requestContext, requestOptions); + return stream.result(); + }, }, ); } catch (error) { @@ -13554,6 +13566,12 @@ export class AgentSession { metadata: this.agent.metadataForProvider(model.provider), convertToLlm: messages => this.#convertToLlmForSideRequest(messages), telemetry: resolveTelemetry(this.agent.telemetry, this.sessionId), + // Same per-provider concurrency cap rationale as the compaction + // path above (chatgpt-codex review on #3751). + completeImpl: async (requestModel, requestContext, requestOptions) => { + const stream = await this.#sideStreamFn(requestModel, requestContext, requestOptions); + return stream.result(); + }, }); this.#branchSummaryAbortController = undefined; if (result.aborted) { diff --git a/packages/coding-agent/test/issue-3751-compaction-provider-cap.test.ts b/packages/coding-agent/test/issue-3751-compaction-provider-cap.test.ts new file mode 100644 index 000000000..428bdb1ba --- /dev/null +++ b/packages/coding-agent/test/issue-3751-compaction-provider-cap.test.ts @@ -0,0 +1,196 @@ +/** + * Regression for the [#3751](https://github.com/can1357/oh-my-pi/pull/3751) + * chatgpt-codex follow-up: the per-LLM-turn provider concurrency wrapper + * (`wrapStreamFnWithProviderConcurrency`) was only attached to + * `Agent.streamFn` / `Agent.sideStreamFn`, so direct compaction oneshots + * (`#compactWithFallbackModel` → `compact()` → `instrumentedCompleteSimple`) + * still went through the default `completeSimple` transport and bypassed + * `providers.ollama-cloud.maxConcurrency`. + * + * The fix threads a `completeImpl` through `SummaryOptions` / + * `GenerateBranchSummaryOptions`, and agent-session wires it to the same + * `#sideStreamFn` the handoff path already uses. This test asserts both + * halves of the contract: + * 1. `SummaryOptions.completeImpl` is honored by every fan-out + * summarizer (history + turn-prefix + short), so the default + * `completeSimple` is never reached. + * 2. When that override is built from the limiter-wrapped sideStreamFn — + * the exact shape agent-session installs — peak in-flight HTTP + * requests respect the configured cap even with a concurrent + * unrelated provider call. + */ +import { afterEach, describe, expect, it, vi } from "bun:test"; +import type { StreamFn } from "@oh-my-pi/pi-agent-core"; +import { compact, type CompactionPreparation, createFileOps, DEFAULT_COMPACTION_SETTINGS } from "@oh-my-pi/pi-agent-core/compaction"; +import type { AgentMessage } from "@oh-my-pi/pi-agent-core/types"; +import * as ai from "@oh-my-pi/pi-ai"; +import type { AssistantMessage, Model } from "@oh-my-pi/pi-ai"; +import { AssistantMessageEventStream } from "@oh-my-pi/pi-ai/utils/event-stream"; +import { getBundledModel } from "@oh-my-pi/pi-catalog/models"; +import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; +import { wrapStreamFnWithProviderConcurrency } from "@oh-my-pi/pi-coding-agent/task/provider-concurrency"; + +interface Deferred { + promise: Promise; + resolve: () => void; +} + +function deferred(): Deferred { + const { promise, resolve } = Promise.withResolvers(); + return { promise, resolve }; +} + +function requireModel(provider: string, id: string): Model { + const model = getBundledModel(provider as Parameters[0], id); + if (!model) throw new Error(`Expected bundled model ${provider}/${id}`); + return model; +} + +function makeAssistantMessage(text: string, model: Model): AssistantMessage { + return { + role: "assistant", + content: [{ type: "text", text }], + api: model.api, + provider: model.provider, + model: model.id, + usage: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 0, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }, + stopReason: "stop", + timestamp: Date.now(), + }; +} + +function makePreparation(): CompactionPreparation { + return { + firstKeptEntryId: "kept-1", + messagesToSummarize: [ + { role: "user", content: "history msg", timestamp: 1 } satisfies AgentMessage, + { + role: "assistant", + content: [{ type: "text", text: "history reply" }], + timestamp: 2, + } satisfies AgentMessage, + ], + turnPrefixMessages: [{ role: "user", content: "turn prefix msg", timestamp: 3 } satisfies AgentMessage], + recentMessages: [{ role: "user", content: "recent msg", timestamp: 4 } satisfies AgentMessage], + isSplitTurn: true, + tokensBefore: 12_345, + fileOps: createFileOps(), + settings: { ...DEFAULT_COMPACTION_SETTINGS, remoteEnabled: false }, + }; +} + +async function waitFor(check: () => boolean, label: string): Promise { + for (let i = 0; i < 1000 && !check(); i++) { + await Promise.resolve(); + } + if (!check()) throw new Error(`Timed out waiting for: ${label}`); +} + +afterEach(() => { + vi.restoreAllMocks(); +}); + +describe("issue #3751: compaction summaries respect provider concurrency cap", () => { + it("compact() routes every fan-out summarizer through SummaryOptions.completeImpl", async () => { + const model = requireModel("ollama-cloud", "gpt-oss:120b"); + const defaultSpy = vi + .spyOn(ai, "completeSimple") + .mockResolvedValue(makeAssistantMessage("default transport must not run", model)); + + const overrideCalls: { model: Model; system: string | undefined }[] = []; + const override = async (requestModel: Model, ctx: ai.Context): Promise => { + overrideCalls.push({ model: requestModel, system: ctx.systemPrompt?.[0] }); + return makeAssistantMessage("summary text", requestModel); + }; + + await compact(makePreparation(), model, "test-key", undefined, undefined, { + completeImpl: override, + }); + + // Split-turn preparation fans out into history + turn-prefix + short. + expect(overrideCalls).toHaveLength(3); + expect(defaultSpy).not.toHaveBeenCalled(); + // Each summarizer ships the system prompt that documents the compaction + // contract; we just sanity-check the override actually saw the request + // instead of being short-circuited by remote compaction or hook paths. + for (const call of overrideCalls) { + expect(call.model.provider).toBe(model.provider); + expect(typeof call.system).toBe("string"); + } + }); + + it("limiter-wrapped sideStreamFn caps compaction HTTP requests at maxConcurrency=1", async () => { + const model = requireModel("ollama-cloud", "gpt-oss:120b"); + const settings = Settings.isolated({ "providers.ollama-cloud.maxConcurrency": 1 }); + + let inFlight = 0; + let peakInFlight = 0; + const gates: Deferred[] = []; + const base: StreamFn = streamModel => { + const gate = deferred(); + gates.push(gate); + inFlight++; + peakInFlight = Math.max(peakInFlight, inFlight); + const events = new AssistantMessageEventStream(); + void gate.promise.then(() => { + inFlight--; + events.push({ type: "done", reason: "stop", message: makeAssistantMessage("ok", streamModel) }); + events.end(); + }); + return events; + }; + const sideStreamFn = wrapStreamFnWithProviderConcurrency(settings, base); + + // Exactly the shape agent-session.ts:#compactWithFallbackModel installs. + const completeImpl = async ( + requestModel: Model, + requestContext: ai.Context, + requestOptions: ai.SimpleStreamOptions, + ): Promise => { + const stream = await sideStreamFn(requestModel, requestContext, requestOptions); + return stream.result(); + }; + + const defaultSpy = vi + .spyOn(ai, "completeSimple") + .mockResolvedValue(makeAssistantMessage("default transport must not run", model)); + + // Kick off compaction; it must immediately try to acquire the slot. + const compaction = compact(makePreparation(), model, "test-key", undefined, undefined, { + completeImpl, + }); + await waitFor(() => gates.length === 1, "first compaction summarizer in flight"); + expect(inFlight).toBe(1); + + // A direct sideStreamFn call (e.g. /btw, IRC reply) issued while + // compaction is mid-flight must queue behind the held slot — proving + // the cap is shared. Under the pre-fix wiring this would bypass the + // limiter and run immediately, pushing peakInFlight to 2. + const concurrent = sideStreamFn(model, { messages: [] }, {}); + await Promise.resolve(); + await Promise.resolve(); + expect(gates).toHaveLength(1); + + // Drain in submission order: 4 HTTP requests total = 3 compaction + // summarizers + 1 concurrent side request. We intentionally do NOT + // distinguish which gate is which — that's the whole point of the + // shared cap. Whatever order the limiter admits them in, peak + // in-flight must never exceed 1. + for (let i = 0; i < 4; i++) { + await waitFor(() => gates.length > i, `request ${i + 1} admitted`); + expect(inFlight).toBe(1); + gates[i]!.resolve(); + } + await Promise.all([compaction, concurrent]); + + expect(peakInFlight).toBe(1); + expect(defaultSpy).not.toHaveBeenCalled(); + }); +});