fix(compaction): routed summary oneshots through provider cap

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
This commit is contained in:
roboomp
2026-06-28 21:11:48 +00:00
parent edaaec398c
commit 9599689af3
4 changed files with 245 additions and 5 deletions
@@ -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?: <TApi extends Api>(
model: Model<TApi>,
ctx: Context,
options: SimpleStreamOptions,
) => Promise<AssistantMessage>;
}
// ============================================================================
@@ -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
+18 -3
View File
@@ -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?: <TApi extends Api>(
model: Model<TApi>,
ctx: Context,
options: SimpleStreamOptions,
) => Promise<AssistantMessage>;
}
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") {
@@ -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) {
@@ -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<void>;
resolve: () => void;
}
function deferred(): Deferred {
const { promise, resolve } = Promise.withResolvers<void>();
return { promise, resolve };
}
function requireModel(provider: string, id: string): Model {
const model = getBundledModel(provider as Parameters<typeof getBundledModel>[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<void> {
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<AssistantMessage> => {
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<AssistantMessage> => {
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();
});
});