feat(ai): implemented bounded retries for empty assistant completions
- Introduced `withEmptyCompletionRetry` utility to manage bounded retries with exponential backoff for empty streaming responses. - Integrated completion validation across Anthropic and OpenAI providers to ensure robust handling of intermittent empty outputs. - Updated `omp bench` and `omp dry-balance` to correctly resolve extension-contributed providers and report content-less runs as failures. - Added comprehensive unit and regression tests to validate retry logic, provider resolution, and failure reporting metrics.
This commit is contained in:
@@ -1,6 +1,10 @@
|
||||
# Changelog
|
||||
|
||||
## [Unreleased]
|
||||
### Added
|
||||
|
||||
- Added bounded auto-retry for empty assistant completions specifically to the OpenAI Responses provider
|
||||
- Added bounded auto-retry for empty assistant completions across the OpenAI Chat Completions, OpenAI Responses, and Anthropic Messages providers. A benign terminal stop that streamed no content and billed no output tokens — the signature of a flaky OpenAI-/Anthropic-compatible gateway that intermittently 200s with an empty body — is now retried up to twice with exponential backoff (honoring `providerRetryWait`) before being surfaced, instead of silently stalling the agent loop. Retries fire only before any content streams, so live streaming (including thinking) is never delayed, retried, or duplicated.
|
||||
|
||||
## [16.1.3] - 2026-06-19
|
||||
|
||||
@@ -3950,4 +3954,4 @@ _Dedicated to Peter's shoulder ([@steipete](https://twitter.com/steipete))_
|
||||
|
||||
## [0.9.4] - 2025-11-26
|
||||
|
||||
Initial release with multi-provider LLM support.
|
||||
Initial release with multi-provider LLM support.
|
||||
@@ -46,6 +46,7 @@ import type {
|
||||
import { resolveServiceTier } from "../types";
|
||||
import { isRecord, normalizeSystemPrompts, normalizeToolCallId, resolveCacheRetention } from "../utils";
|
||||
import { createAbortSourceTracker } from "../utils/abort";
|
||||
import { withEmptyCompletionRetry } from "../utils/empty-completion-retry";
|
||||
import { AssistantMessageEventStream } from "../utils/event-stream";
|
||||
import { isFoundryEnabled } from "../utils/foundry";
|
||||
import { finalizeErrorMessage, type RawHttpRequestDump, rewriteCopilotError } from "../utils/http-inspector";
|
||||
@@ -1575,7 +1576,7 @@ export function applyAnthropicUsageExtras(usage: Usage, source: AnthropicUsageLi
|
||||
}
|
||||
}
|
||||
|
||||
export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
|
||||
const streamAnthropicOnce = (
|
||||
model: Model<"anthropic-messages">,
|
||||
context: Context,
|
||||
options?: AnthropicOptions,
|
||||
@@ -2309,6 +2310,16 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
|
||||
return stream;
|
||||
};
|
||||
|
||||
/**
|
||||
* Public entry: wrap the single-attempt streamer with bounded empty-completion
|
||||
* retries (a benign terminal stop carrying no content/usage would otherwise
|
||||
* stall the agent loop). The inner attempt keeps its own provider-failure retry
|
||||
* loop; this layer only re-issues a fresh request on an empty success. Shared
|
||||
* with the OpenAI-completions provider via `withEmptyCompletionRetry`.
|
||||
*/
|
||||
export const streamAnthropic: StreamFunction<"anthropic-messages"> = (model, context, options) =>
|
||||
withEmptyCompletionRetry(model, context, options, streamAnthropicOnce);
|
||||
|
||||
export type AnthropicSystemBlock = {
|
||||
type: "text";
|
||||
text: string;
|
||||
|
||||
@@ -27,6 +27,7 @@ import type {
|
||||
} from "../types";
|
||||
import { normalizeSystemPrompts } from "../utils";
|
||||
import { createAbortSourceTracker } from "../utils/abort";
|
||||
import { hasVisibleAssistantContent, withEmptyCompletionRetry } from "../utils/empty-completion-retry";
|
||||
import { AssistantMessageEventStream } from "../utils/event-stream";
|
||||
import { finalizeErrorMessage, type RawHttpRequestDump, rewriteCopilotError } from "../utils/http-inspector";
|
||||
import {
|
||||
@@ -537,7 +538,7 @@ const OPENAI_COMPLETIONS_FIRST_EVENT_TIMEOUT_MESSAGE =
|
||||
// converts the already-successful response into a timeout error.
|
||||
const OPENAI_COMPLETIONS_POST_FINISH_GRACE_MS = 2_500;
|
||||
|
||||
export const streamOpenAICompletions: StreamFunction<"openai-completions"> = (
|
||||
const streamOpenAICompletionsOnce = (
|
||||
model: Model<"openai-completions">,
|
||||
context: Context,
|
||||
options?: OpenAICompletionsOptions,
|
||||
@@ -1234,7 +1235,7 @@ export const streamOpenAICompletions: StreamFunction<"openai-completions"> = (
|
||||
if (
|
||||
policy.stream.emptyLengthFinishIsContextError &&
|
||||
output.stopReason === "length" &&
|
||||
!hasVisibleCompletionContent(output)
|
||||
!hasVisibleAssistantContent(output)
|
||||
) {
|
||||
output.stopReason = "error";
|
||||
output.errorMessage = EMPTY_OLLAMA_LENGTH_COMPLETION_MESSAGE;
|
||||
@@ -1288,6 +1289,15 @@ export const streamOpenAICompletions: StreamFunction<"openai-completions"> = (
|
||||
return stream;
|
||||
};
|
||||
|
||||
/**
|
||||
* Public entry: wrap the single-attempt streamer with bounded empty-completion
|
||||
* retries — flaky gateways occasionally 200 with `delta: {}` + `finish_reason:
|
||||
* "stop"` and no usage, which would otherwise stall the agent loop. Shared with
|
||||
* the Anthropic provider via `withEmptyCompletionRetry`.
|
||||
*/
|
||||
export const streamOpenAICompletions: StreamFunction<"openai-completions"> = (model, context, options) =>
|
||||
withEmptyCompletionRetry(model, context, options, streamOpenAICompletionsOnce);
|
||||
|
||||
function createRequestSetup(
|
||||
model: Model<"openai-completions">,
|
||||
context: Context,
|
||||
@@ -2061,16 +2071,6 @@ function convertTools(
|
||||
};
|
||||
}
|
||||
|
||||
const NON_WHITESPACE_RE = /\S/;
|
||||
|
||||
function hasVisibleCompletionContent(message: AssistantMessage): boolean {
|
||||
for (const block of message.content) {
|
||||
if (block.type === "toolCall") return true;
|
||||
if (block.type === "text" && NON_WHITESPACE_RE.test(block.text)) return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
const EMPTY_OLLAMA_LENGTH_COMPLETION_MESSAGE =
|
||||
"Model returned no content: prompt filled the context window; raise Ollama num_ctx or shorten the prompt.";
|
||||
|
||||
|
||||
@@ -21,6 +21,7 @@ import {
|
||||
sanitizeOpenAIResponsesHistoryItemsForReplay,
|
||||
} from "../utils";
|
||||
import { createAbortSourceTracker } from "../utils/abort";
|
||||
import { withEmptyCompletionRetry } from "../utils/empty-completion-retry";
|
||||
import { AssistantMessageEventStream } from "../utils/event-stream";
|
||||
import { finalizeErrorMessage, type RawHttpRequestDump, rewriteCopilotError } from "../utils/http-inspector";
|
||||
import {
|
||||
@@ -338,7 +339,7 @@ type OpenAIResponsesSamplingParams = ResponseCreateParamsStreaming & {
|
||||
/**
|
||||
* Generate function for OpenAI Responses API
|
||||
*/
|
||||
export const streamOpenAIResponses: StreamFunction<"openai-responses"> = (
|
||||
const streamOpenAIResponsesOnce = (
|
||||
model: Model<"openai-responses">,
|
||||
context: Context,
|
||||
options?: OpenAIResponsesOptions,
|
||||
@@ -737,6 +738,15 @@ export const streamOpenAIResponses: StreamFunction<"openai-responses"> = (
|
||||
return stream;
|
||||
};
|
||||
|
||||
/**
|
||||
* Public entry: wrap the single-attempt Responses streamer with bounded
|
||||
* empty-completion retries — a `response.completed` carrying no content/usage
|
||||
* would otherwise stall the agent loop. Shared with the OpenAI-completions and
|
||||
* Anthropic providers via `withEmptyCompletionRetry`.
|
||||
*/
|
||||
export const streamOpenAIResponses: StreamFunction<"openai-responses"> = (model, context, options) =>
|
||||
withEmptyCompletionRetry(model, context, options, streamOpenAIResponsesOnce);
|
||||
|
||||
export function buildParams(
|
||||
model: Model<"openai-responses">,
|
||||
context: Context,
|
||||
|
||||
@@ -0,0 +1,159 @@
|
||||
/**
|
||||
* Bounded retries for an empty assistant completion.
|
||||
*
|
||||
* Some providers — and especially flaky OpenAI-/Anthropic-compatible gateways —
|
||||
* intermittently return a benign terminal stop carrying no content and no usage
|
||||
* (e.g. a single OpenAI `delta: {}` + `finish_reason: "stop"` chunk). Delivered
|
||||
* as-is the agent loop has nothing to act on and silently halts mid-task, so the
|
||||
* request must be retried instead of surfaced.
|
||||
*
|
||||
* This wraps a single-attempt provider stream and re-invokes it (a fresh request
|
||||
* with its own message state) when an attempt produces no meaningful content.
|
||||
* Only a stream that streamed nothing meaningful is retried: the moment any
|
||||
* text/thinking/tool delta is forwarded the attempt is committed, so live
|
||||
* streaming (including thinking) is never delayed, retried, or duplicated.
|
||||
*
|
||||
* Mirrors the Gemini empty-response policy in `google-shared` (which keeps its
|
||||
* own integrated loop) and is shared by the OpenAI-completions and
|
||||
* Anthropic-messages providers.
|
||||
*/
|
||||
import { scheduler } from "node:timers/promises";
|
||||
import type { AssistantMessage, AssistantMessageEvent, Context } from "../types";
|
||||
import { AssistantMessageEventStream } from "./event-stream";
|
||||
|
||||
export const MAX_EMPTY_COMPLETION_RETRIES = 2;
|
||||
export const EMPTY_COMPLETION_BASE_DELAY_MS = 500;
|
||||
|
||||
const NON_WHITESPACE_RE = /\S/;
|
||||
|
||||
/**
|
||||
* Whether a completed assistant message carries content worth delivering: a tool
|
||||
* call or any non-whitespace text. An empty/whitespace-only message — or one
|
||||
* that only ever produced thinking — is the "empty response" failure.
|
||||
*/
|
||||
export function hasVisibleAssistantContent(message: AssistantMessage): boolean {
|
||||
for (const block of message.content) {
|
||||
if (block.type === "toolCall") return true;
|
||||
if (block.type === "text" && NON_WHITESPACE_RE.test(block.text)) return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/** A streamed event that delivers content worth committing the attempt for. */
|
||||
function isMeaningfulCompletionEvent(event: AssistantMessageEvent): boolean {
|
||||
switch (event.type) {
|
||||
case "text_delta":
|
||||
case "thinking_delta":
|
||||
case "toolcall_delta":
|
||||
return event.delta.length > 0;
|
||||
case "text_end":
|
||||
case "thinking_end":
|
||||
return event.content.length > 0;
|
||||
case "toolcall_start":
|
||||
case "toolcall_end":
|
||||
return true;
|
||||
default:
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
interface EmptyCompletionRetryOptions {
|
||||
signal?: AbortSignal;
|
||||
providerRetryWait?: (delayMs: number, signal?: AbortSignal) => Promise<void>;
|
||||
}
|
||||
|
||||
/**
|
||||
* Wrap a single-attempt provider stream with bounded empty-completion retries.
|
||||
* `attempt` MUST create a fresh request (and its own output message) on each
|
||||
* call so a retry never inherits stale metadata from an empty attempt.
|
||||
*/
|
||||
export function withEmptyCompletionRetry<M, O extends EmptyCompletionRetryOptions>(
|
||||
model: M,
|
||||
context: Context,
|
||||
options: O | undefined,
|
||||
attempt: (model: M, context: Context, options?: O) => AssistantMessageEventStream,
|
||||
): AssistantMessageEventStream {
|
||||
const outer = new AssistantMessageEventStream();
|
||||
const signal = options?.signal;
|
||||
void (async () => {
|
||||
for (let emptyAttempt = 0; ; emptyAttempt++) {
|
||||
const inner = attempt(model, context, options);
|
||||
const buffered: AssistantMessageEvent[] = [];
|
||||
let committed = false;
|
||||
let terminal: AssistantMessageEvent | undefined;
|
||||
const flush = (): void => {
|
||||
for (const event of buffered) outer.push(event);
|
||||
buffered.length = 0;
|
||||
};
|
||||
try {
|
||||
for await (const event of inner) {
|
||||
if (event.type === "done" || event.type === "error") {
|
||||
terminal = event;
|
||||
break;
|
||||
}
|
||||
// Buffer pre-content events (start/*_start) so an empty attempt can
|
||||
// be discarded; commit the moment real content streams.
|
||||
if (!committed && !isMeaningfulCompletionEvent(event)) {
|
||||
buffered.push(event);
|
||||
continue;
|
||||
}
|
||||
committed = true;
|
||||
flush();
|
||||
outer.push(event);
|
||||
if (outer.done) return;
|
||||
}
|
||||
} catch (error) {
|
||||
flush();
|
||||
outer.fail(error);
|
||||
return;
|
||||
}
|
||||
|
||||
// Retry only a genuinely degenerate completion: a normal stop that
|
||||
// produced no visible content AND billed no output tokens (the flaky
|
||||
// gateway signature — charged nothing, returned nothing). A stop that
|
||||
// reports output tokens spent its budget somewhere (e.g. thinking) and
|
||||
// is left alone.
|
||||
const message = terminal?.type === "done" ? terminal.message : undefined;
|
||||
const isRetryableEmpty =
|
||||
!committed &&
|
||||
message !== undefined &&
|
||||
message.stopReason === "stop" &&
|
||||
!message.errorMessage &&
|
||||
(message.usage?.output ?? 0) <= 0 &&
|
||||
!hasVisibleAssistantContent(message);
|
||||
|
||||
if (isRetryableEmpty && emptyAttempt < MAX_EMPTY_COMPLETION_RETRIES && !signal?.aborted) {
|
||||
const delayMs = EMPTY_COMPLETION_BASE_DELAY_MS * 2 ** emptyAttempt;
|
||||
try {
|
||||
if (options?.providerRetryWait) await options.providerRetryWait(delayMs, signal);
|
||||
else await scheduler.wait(delayMs, { signal });
|
||||
} catch (waitError) {
|
||||
// Aborted during backoff: deliver the empty result rather than hang.
|
||||
// Any other wait failure is a real error and must surface.
|
||||
flush();
|
||||
if (signal?.aborted) {
|
||||
if (terminal) outer.push(terminal);
|
||||
} else {
|
||||
outer.fail(waitError);
|
||||
}
|
||||
return;
|
||||
}
|
||||
// Discard the buffered `start` from this empty attempt and retry.
|
||||
continue;
|
||||
}
|
||||
|
||||
flush();
|
||||
if (terminal) {
|
||||
outer.push(terminal);
|
||||
} else if (!outer.done) {
|
||||
try {
|
||||
outer.end(await inner.result());
|
||||
} catch (error) {
|
||||
outer.fail(error);
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
})();
|
||||
return outer;
|
||||
}
|
||||
@@ -0,0 +1,258 @@
|
||||
/**
|
||||
* Contracts for `withEmptyCompletionRetry` (shared by the OpenAI-completions and
|
||||
* Anthropic-messages providers): a benign terminal stop with no content/usage is
|
||||
* retried a bounded number of times; once any content streams the attempt is
|
||||
* committed (no retry, no duplicate `start`); the cap delivers the empty result;
|
||||
* backoff failures surface unless the caller aborted.
|
||||
*/
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import type { AssistantMessage, AssistantMessageEvent, Context, Usage } from "@oh-my-pi/pi-ai/types";
|
||||
import { MAX_EMPTY_COMPLETION_RETRIES, withEmptyCompletionRetry } from "@oh-my-pi/pi-ai/utils/empty-completion-retry";
|
||||
import { AssistantMessageEventStream } from "@oh-my-pi/pi-ai/utils/event-stream";
|
||||
|
||||
const CTX = {} as Context;
|
||||
|
||||
function usage(): Usage {
|
||||
return {
|
||||
input: 0,
|
||||
output: 0,
|
||||
cacheRead: 0,
|
||||
cacheWrite: 0,
|
||||
totalTokens: 0,
|
||||
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
||||
};
|
||||
}
|
||||
|
||||
function assistant(texts: string[] = []): AssistantMessage {
|
||||
return {
|
||||
role: "assistant",
|
||||
content: texts.map(text => ({ type: "text" as const, text })),
|
||||
api: "openai-completions",
|
||||
provider: "test",
|
||||
model: "test-model",
|
||||
timestamp: 1,
|
||||
stopReason: "stop",
|
||||
usage: usage(),
|
||||
};
|
||||
}
|
||||
|
||||
function streamFromEvents(events: AssistantMessageEvent[]): AssistantMessageEventStream {
|
||||
const stream = new AssistantMessageEventStream();
|
||||
for (const event of events) stream.push(event);
|
||||
return stream;
|
||||
}
|
||||
|
||||
/** start + stop with no content/usage — the flaky-gateway empty completion. */
|
||||
function emptyAttempt(): AssistantMessageEventStream {
|
||||
const message = assistant();
|
||||
return streamFromEvents([
|
||||
{ type: "start", partial: message },
|
||||
{ type: "done", reason: "stop", message },
|
||||
] as unknown as AssistantMessageEvent[]);
|
||||
}
|
||||
|
||||
function contentAttempt(): AssistantMessageEventStream {
|
||||
const message = assistant(["hello"]);
|
||||
return streamFromEvents([
|
||||
{ type: "start", partial: message },
|
||||
{ type: "text_start", contentIndex: 0, partial: message },
|
||||
{ type: "text_delta", contentIndex: 0, delta: "hello", partial: message },
|
||||
{ type: "text_end", contentIndex: 0, content: "hello", partial: message },
|
||||
{ type: "done", reason: "stop", message },
|
||||
] as unknown as AssistantMessageEvent[]);
|
||||
}
|
||||
|
||||
async function drain(stream: AssistantMessageEventStream): Promise<AssistantMessageEvent[]> {
|
||||
const events: AssistantMessageEvent[] = [];
|
||||
for await (const event of stream) events.push(event);
|
||||
return events;
|
||||
}
|
||||
|
||||
describe("withEmptyCompletionRetry", () => {
|
||||
it("retries past empty attempts and delivers the first non-empty one", async () => {
|
||||
let attempts = 0;
|
||||
const waits: number[] = [];
|
||||
const stream = withEmptyCompletionRetry({}, CTX, { providerRetryWait: async ms => void waits.push(ms) }, () => {
|
||||
attempts++;
|
||||
return attempts <= MAX_EMPTY_COMPLETION_RETRIES ? emptyAttempt() : contentAttempt();
|
||||
});
|
||||
|
||||
const events = await drain(stream);
|
||||
const result = await stream.result();
|
||||
|
||||
expect(attempts).toBe(MAX_EMPTY_COMPLETION_RETRIES + 1);
|
||||
expect(waits).toHaveLength(MAX_EMPTY_COMPLETION_RETRIES);
|
||||
// Discarded attempts' `start` events must not leak — exactly one survives.
|
||||
expect(events.filter(e => e.type === "start")).toHaveLength(1);
|
||||
expect(events.some(e => e.type === "text_delta")).toBe(true);
|
||||
expect(events.at(-1)?.type).toBe("done");
|
||||
expect(result.content).toEqual([{ type: "text", text: "hello" }]);
|
||||
});
|
||||
|
||||
it("delivers the empty result after exhausting the retry cap", async () => {
|
||||
let attempts = 0;
|
||||
const waits: number[] = [];
|
||||
const stream = withEmptyCompletionRetry({}, CTX, { providerRetryWait: async ms => void waits.push(ms) }, () => {
|
||||
attempts++;
|
||||
return emptyAttempt();
|
||||
});
|
||||
|
||||
const events = await drain(stream);
|
||||
const result = await stream.result();
|
||||
|
||||
expect(attempts).toBe(MAX_EMPTY_COMPLETION_RETRIES + 1);
|
||||
expect(waits).toHaveLength(MAX_EMPTY_COMPLETION_RETRIES);
|
||||
expect(events.filter(e => e.type === "start")).toHaveLength(1);
|
||||
expect(events.at(-1)?.type).toBe("done");
|
||||
expect(result.content).toEqual([]);
|
||||
});
|
||||
|
||||
it("does not retry when the first attempt streams content", async () => {
|
||||
let attempts = 0;
|
||||
let waited = false;
|
||||
const stream = withEmptyCompletionRetry(
|
||||
{},
|
||||
CTX,
|
||||
{
|
||||
providerRetryWait: async () => {
|
||||
waited = true;
|
||||
},
|
||||
},
|
||||
() => {
|
||||
attempts++;
|
||||
return contentAttempt();
|
||||
},
|
||||
);
|
||||
|
||||
await drain(stream);
|
||||
|
||||
expect(attempts).toBe(1);
|
||||
expect(waited).toBe(false);
|
||||
});
|
||||
|
||||
it("commits on streamed thinking and does not retry a thinking-only stop", async () => {
|
||||
let attempts = 0;
|
||||
const stream = withEmptyCompletionRetry({}, CTX, {}, () => {
|
||||
attempts++;
|
||||
const message = assistant(); // no visible content; only thinking streams
|
||||
return streamFromEvents([
|
||||
{ type: "start", partial: message },
|
||||
{ type: "thinking_delta", contentIndex: 0, delta: "pondering", partial: message },
|
||||
{ type: "done", reason: "stop", message },
|
||||
] as unknown as AssistantMessageEvent[]);
|
||||
});
|
||||
|
||||
const events = await drain(stream);
|
||||
|
||||
expect(attempts).toBe(1);
|
||||
expect(events.some(e => e.type === "thinking_delta")).toBe(true);
|
||||
});
|
||||
|
||||
it("propagates a non-abort backoff failure instead of masking the empty result", async () => {
|
||||
const stream = withEmptyCompletionRetry(
|
||||
{},
|
||||
CTX,
|
||||
{
|
||||
providerRetryWait: async () => {
|
||||
throw new Error("wait boom");
|
||||
},
|
||||
},
|
||||
() => emptyAttempt(),
|
||||
);
|
||||
|
||||
let caught: unknown;
|
||||
try {
|
||||
await drain(stream);
|
||||
} catch (error) {
|
||||
caught = error;
|
||||
}
|
||||
expect((caught as Error | undefined)?.message).toBe("wait boom");
|
||||
});
|
||||
|
||||
it("delivers the empty result when aborted during backoff", async () => {
|
||||
const controller = new AbortController();
|
||||
let attempts = 0;
|
||||
const stream = withEmptyCompletionRetry(
|
||||
{},
|
||||
CTX,
|
||||
{
|
||||
signal: controller.signal,
|
||||
providerRetryWait: async () => {
|
||||
controller.abort();
|
||||
throw new Error("aborted");
|
||||
},
|
||||
},
|
||||
() => {
|
||||
attempts++;
|
||||
return emptyAttempt();
|
||||
},
|
||||
);
|
||||
|
||||
const events = await drain(stream);
|
||||
const result = await stream.result();
|
||||
|
||||
expect(attempts).toBe(1);
|
||||
expect(events.at(-1)?.type).toBe("done");
|
||||
expect(result.content).toEqual([]);
|
||||
});
|
||||
|
||||
it("discards buffered pre-content markers from a retried empty attempt", async () => {
|
||||
let attempts = 0;
|
||||
const stream = withEmptyCompletionRetry({}, CTX, { providerRetryWait: async () => {} }, () => {
|
||||
attempts++;
|
||||
if (attempts === 1) {
|
||||
const message = assistant();
|
||||
return streamFromEvents([
|
||||
{ type: "start", partial: message },
|
||||
{ type: "thinking_start", contentIndex: 0, partial: message },
|
||||
{ type: "done", reason: "stop", message },
|
||||
] as unknown as AssistantMessageEvent[]);
|
||||
}
|
||||
return contentAttempt();
|
||||
});
|
||||
|
||||
const events = await drain(stream);
|
||||
|
||||
expect(attempts).toBe(2);
|
||||
// The empty attempt's start + thinking_start were discarded; only the
|
||||
// successful attempt's events reach the consumer.
|
||||
expect(events.filter(e => e.type === "start")).toHaveLength(1);
|
||||
expect(events.some(e => e.type === "thinking_start")).toBe(false);
|
||||
expect(events.some(e => e.type === "text_delta")).toBe(true);
|
||||
});
|
||||
|
||||
it("streams content as it arrives without waiting for the terminal event", async () => {
|
||||
let waited = false;
|
||||
const message = assistant(["streamed"]);
|
||||
const inner = new AssistantMessageEventStream();
|
||||
const stream = withEmptyCompletionRetry(
|
||||
{},
|
||||
CTX,
|
||||
{
|
||||
providerRetryWait: async () => {
|
||||
waited = true;
|
||||
},
|
||||
},
|
||||
() => inner,
|
||||
);
|
||||
|
||||
const iterator = stream[Symbol.asyncIterator]();
|
||||
// Push content with no terminal yet: the buffered start then the delta must
|
||||
// surface before any `done` exists, proving the wrapper does not buffer
|
||||
// meaningful content until completion.
|
||||
inner.push({ type: "start", partial: message } as unknown as AssistantMessageEvent);
|
||||
inner.push({
|
||||
type: "text_delta",
|
||||
contentIndex: 0,
|
||||
delta: "streamed",
|
||||
partial: message,
|
||||
} as unknown as AssistantMessageEvent);
|
||||
|
||||
expect((await iterator.next()).value?.type).toBe("start");
|
||||
expect((await iterator.next()).value?.type).toBe("text_delta");
|
||||
|
||||
inner.push({ type: "done", reason: "stop", message } as unknown as AssistantMessageEvent);
|
||||
expect((await iterator.next()).value?.type).toBe("done");
|
||||
expect(waited).toBe(false);
|
||||
});
|
||||
});
|
||||
@@ -21,18 +21,30 @@ const context: Context = {
|
||||
|
||||
function createSseResponse(): Response {
|
||||
return new Response(
|
||||
`data: ${JSON.stringify({
|
||||
type: "response.completed",
|
||||
response: {
|
||||
status: "completed",
|
||||
usage: { input_tokens: 1, output_tokens: 1, total_tokens: 2, input_tokens_details: { cached_tokens: 0 } },
|
||||
},
|
||||
})}\n\n`,
|
||||
`data: ${JSON.stringify({ type: "response.content_part.added", part: { type: "output_text", text: "" } })}\n\n` +
|
||||
`data: ${JSON.stringify({ type: "response.output_text.delta", delta: "ok" })}\n\n` +
|
||||
`data: ${JSON.stringify({
|
||||
type: "response.completed",
|
||||
response: {
|
||||
status: "completed",
|
||||
usage: {
|
||||
input_tokens: 1,
|
||||
output_tokens: 1,
|
||||
total_tokens: 2,
|
||||
input_tokens_details: { cached_tokens: 0 },
|
||||
},
|
||||
},
|
||||
})}\n\n`,
|
||||
{ status: 200, headers: { "content-type": "text/event-stream" } },
|
||||
);
|
||||
}
|
||||
function createChatDoneResponse(): Response {
|
||||
return new Response("data: [DONE]\n\n", { status: 200, headers: { "content-type": "text/event-stream" } });
|
||||
return new Response(
|
||||
`data: ${JSON.stringify({ choices: [{ index: 0, delta: { content: "ok" }, finish_reason: null }] })}\n\n` +
|
||||
`data: ${JSON.stringify({ choices: [{ index: 0, delta: {}, finish_reason: "stop" }] })}\n\n` +
|
||||
`data: [DONE]\n\n`,
|
||||
{ status: 200, headers: { "content-type": "text/event-stream" } },
|
||||
);
|
||||
}
|
||||
|
||||
function buildOpenRouterModel(
|
||||
|
||||
@@ -1,10 +1,13 @@
|
||||
# Changelog
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed `omp bench` and `omp dry-balance` failing to resolve models from extension providers
|
||||
- Improved error reporting for `omp bench` runs that return no output or tokens
|
||||
- Cache-miss marker no longer fires on a cold turn whose predecessor only *wrote* the prompt cache (never read it back). The session's opening request always writes the prefix with `cacheRead 0`, so a long-running first tool call (e.g. `gh run watch`) that outlived the provider's cache TTL surfaced a spurious `⊘ cache miss` divider right under the opening message. The marker now requires the previous turn to have actually read a warm prefix, so it flags only a demonstrably working cache going cold — and collapses a run of consecutive cold turns to a single marker at the moment the cache broke.
|
||||
- Fixed `omp bench` and `omp dry-balance` failing to resolve models from extension-registered providers (`pi.registerProvider(...)`). Both commands build a one-shot `ModelRegistry` that previously only knew built-in catalog providers, so a `provider/model` selector for a provider contributed by an extension under `~/.omp/agent/extensions/` errored with "Model not found". A new `loadCliExtensionProviders` helper loads the session's extensions, drains their provider registrations into the registry, and discovers dynamic provider catalogs before resolving selectors — mirroring the interactive session and `omp models` paths.
|
||||
- `omp bench` now reports a run that streamed no content and measured no output tokens as a failure ("provider returned no output") instead of a misleading green check with `tokens 0 / TPS 0.0`.
|
||||
|
||||
## [16.1.3] - 2026-06-19
|
||||
|
||||
@@ -12132,4 +12135,4 @@ Initial public release.
|
||||
|
||||
## [0.7.6] - 2025-11-13
|
||||
|
||||
Previous releases did not maintain a changelog.
|
||||
Previous releases did not maintain a changelog.
|
||||
@@ -39,6 +39,21 @@ function fakeStream(): AssistantMessageEventStream {
|
||||
return Object.assign(iterator, { result: async () => message }) as unknown as AssistantMessageEventStream;
|
||||
}
|
||||
|
||||
function emptyStream(): AssistantMessageEventStream {
|
||||
const message = {
|
||||
role: "assistant",
|
||||
content: [],
|
||||
stopReason: "stop",
|
||||
usage: { input: 5, output: 0 },
|
||||
duration: 120,
|
||||
} as unknown as AssistantMessage;
|
||||
const events = [{ type: "done", message }] as unknown as AssistantMessageEvent[];
|
||||
const iterator = (async function* () {
|
||||
for (const event of events) yield event;
|
||||
})();
|
||||
return Object.assign(iterator, { result: async () => message }) as unknown as AssistantMessageEventStream;
|
||||
}
|
||||
|
||||
interface FakeRegistryOptions {
|
||||
models: Model<Api>[];
|
||||
authedProviders: string[];
|
||||
@@ -66,7 +81,11 @@ function fakeRegistry(opts: FakeRegistryOptions): BenchModelRegistry {
|
||||
};
|
||||
}
|
||||
|
||||
async function runBench(selector: string, registry: BenchModelRegistry) {
|
||||
async function runBench(
|
||||
selector: string,
|
||||
registry: BenchModelRegistry,
|
||||
streamFactory: () => AssistantMessageEventStream = fakeStream,
|
||||
) {
|
||||
const stderr: string[] = [];
|
||||
const summary = await runBenchCommand(
|
||||
{ models: [selector], flags: { runs: 1, maxTokens: 64, json: false } },
|
||||
@@ -76,7 +95,7 @@ async function runBench(selector: string, registry: BenchModelRegistry) {
|
||||
writeStdout: () => {},
|
||||
writeStderr: text => stderr.push(text),
|
||||
setExitCode: () => {},
|
||||
streamSimple: () => fakeStream(),
|
||||
streamSimple: () => streamFactory(),
|
||||
now: () => 0,
|
||||
stdoutIsTTY: false,
|
||||
},
|
||||
@@ -135,3 +154,17 @@ describe("bench credential-aware provider selection", () => {
|
||||
expect(stderr).not.toContain("benchmarking");
|
||||
});
|
||||
});
|
||||
|
||||
describe("bench empty-output guard", () => {
|
||||
it("reports a run with no streamed content and no tokens as a failure", async () => {
|
||||
const registry = fakeRegistry({ models: [fakeModel("acme", "model-x")], authedProviders: ["acme"] });
|
||||
|
||||
const { summary } = await runBench("acme/model-x", registry, emptyStream);
|
||||
|
||||
expect(summary.failures).toBe(1);
|
||||
const run = summary.models[0].results[0];
|
||||
expect(run.ok).toBe(false);
|
||||
if (!run.ok) expect(run.error).toContain("no output");
|
||||
expect(summary.models[0].average).toBeNull();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -0,0 +1,89 @@
|
||||
/**
|
||||
* Regression test for `loadCliExtensionProviders`.
|
||||
*
|
||||
* One-shot CLIs (`omp bench`, dry-balance) build a bare `ModelRegistry` that
|
||||
* only knows built-in catalog providers. Before the helper existed they never
|
||||
* loaded extensions, so a provider contributed by an extension
|
||||
* (`pi.registerProvider(...)`, e.g. a custom OpenAI-compatible gateway under
|
||||
* `~/.omp/agent/extensions/`) was invisible to model resolution and
|
||||
* `omp bench <provider>/<model>` failed with "Model not found".
|
||||
*
|
||||
* Contract under test: after `loadCliExtensionProviders` drains the extension's
|
||||
* provider registrations into the registry, a `provider/id` selector for that
|
||||
* extension provider resolves. Discovery is disabled and the extension path is
|
||||
* passed explicitly so the test never touches the developer's real `~/.omp`.
|
||||
*/
|
||||
|
||||
import { afterAll, beforeAll, expect, test } from "bun:test";
|
||||
import * as fs from "node:fs/promises";
|
||||
import { AuthStorage } from "@oh-my-pi/pi-ai";
|
||||
import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry";
|
||||
import { getModelMatchPreferences, resolveCliModel } from "@oh-my-pi/pi-coding-agent/config/model-resolver";
|
||||
import { resetSettingsForTest, Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
|
||||
import { loadCliExtensionProviders } from "@oh-my-pi/pi-coding-agent/sdk";
|
||||
import { TempDir } from "@oh-my-pi/pi-utils";
|
||||
|
||||
let tmp: TempDir;
|
||||
let extPath: string;
|
||||
let dbPath: string;
|
||||
|
||||
beforeAll(async () => {
|
||||
tmp = await TempDir.create("@cli-ext-providers-");
|
||||
extPath = tmp.join("ext.ts");
|
||||
dbPath = tmp.join("auth.db");
|
||||
await fs.writeFile(
|
||||
extPath,
|
||||
`export default function (pi) {
|
||||
pi.registerProvider("bench-gw", {
|
||||
baseUrl: "https://example.com/v1",
|
||||
apiKey: "literal-test-key",
|
||||
api: "openai-completions",
|
||||
models: [{
|
||||
id: "bench-model",
|
||||
name: "Bench Model",
|
||||
reasoning: false,
|
||||
input: ["text"],
|
||||
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
|
||||
contextWindow: 128000,
|
||||
maxTokens: 4096,
|
||||
}],
|
||||
});
|
||||
}
|
||||
`,
|
||||
);
|
||||
});
|
||||
|
||||
afterAll(async () => {
|
||||
resetSettingsForTest();
|
||||
await tmp.remove();
|
||||
});
|
||||
|
||||
test("loadCliExtensionProviders makes extension providers resolvable by selector", async () => {
|
||||
const authStorage = await AuthStorage.create(dbPath);
|
||||
try {
|
||||
const settings = await Settings.init({
|
||||
inMemory: true,
|
||||
cwd: tmp.path(),
|
||||
overrides: { extensions: [extPath], disabledExtensions: [] },
|
||||
});
|
||||
const modelRegistry = new ModelRegistry(authStorage);
|
||||
const preferences = getModelMatchPreferences(settings);
|
||||
|
||||
// Before the drain the extension provider is unknown: resolution fails.
|
||||
const before = resolveCliModel({ cliModel: "bench-gw/bench-model", modelRegistry, preferences });
|
||||
expect(before.model).toBeUndefined();
|
||||
|
||||
await loadCliExtensionProviders(modelRegistry, settings, tmp.path(), {
|
||||
disableExtensionDiscovery: true,
|
||||
additionalExtensionPaths: [extPath],
|
||||
});
|
||||
|
||||
// After the drain the same selector resolves to the extension provider.
|
||||
const after = resolveCliModel({ cliModel: "bench-gw/bench-model", modelRegistry, preferences });
|
||||
expect(after.error).toBeUndefined();
|
||||
expect(after.model?.provider).toBe("bench-gw");
|
||||
expect(after.model?.id).toBe("bench-model");
|
||||
} finally {
|
||||
authStorage.close();
|
||||
}
|
||||
});
|
||||
Reference in New Issue
Block a user