feat: added reconstructed raw SSE capture for stream provider SDKs

- Added reconstructed SSE event emission for OpenAI, Azure, and Anthropic streams.
- Added raw SSE text to debug report bundles, including raw-sse.txt output.
- Included dropped-record metadata in raw SSE text when events were trimmed.
- Updated raw SSE and sse-debug tests for observer-based capture and safety checks.
This commit is contained in:
can1357
2026-06-08 02:35:08 +02:00
parent 71624ee47b
commit 00e6e14cb4
15 changed files with 575 additions and 556 deletions
+3 -2
View File
@@ -1,7 +1,6 @@
# Changelog
## [Unreleased]
### Added
- Added support for `impersonated_service_account` Application Default Credentials (ADC) in Vertex AI to enable chained impersonation without failing via 401 `invalid_client`.
@@ -9,6 +8,8 @@
### Changed
- Changed `onSseEvent` recording for OpenAI Responses, Azure OpenAI Responses, OpenAI Completions, and Anthropic stream providers to emit reconstructed SSE events from decoded SDK stream items instead of wrapping raw fetch responses
- Changed OpenAI Completions SSE diagnostics to include `event: "chat.completion.chunk"` in `onSseEvent` records for chunked responses
- Changed the default Anthropic model in `DEFAULT_MODEL_PER_PROVIDER` from `claude-sonnet-4-6` to `claude-opus-4-6`, so sessions that fall back to the provider default (no configured `default` role, no `--model`, no restored session) now start on Claude Opus 4.6.
### Fixed
@@ -3020,4 +3021,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.
+40 -19
View File
@@ -29,6 +29,7 @@ import type {
Message,
Model,
ProviderSessionState,
RawSseEvent,
RedactedThinkingContent,
ServiceTier,
SimpleStreamOptions,
@@ -62,7 +63,7 @@ import { isCopilotTransientModelError } from "../utils/retry";
import { COMBINATOR_KEYS, NO_STRICT, toolWireSchema } from "../utils/schema";
import { spillToDescription } from "../utils/schema/spill";
import { createSdkStreamRequestOptions } from "../utils/sdk-stream-timeout";
import { notifyRawSseEvent, wrapFetchForSseDebug } from "../utils/sse-debug";
import { notifyRawSseEvent } from "../utils/sse-debug";
import {
AnthropicConnectionTimeoutError,
type AnthropicFetchOptions,
@@ -863,7 +864,6 @@ export type AnthropicClientOptionsArgs = {
hasTools?: boolean;
thinkingEnabled?: boolean;
thinkingDisplay?: AnthropicThinkingDisplay;
onSseEvent?: AnthropicOptions["onSseEvent"];
fetch?: FetchImpl;
claudeCodeSessionId?: string;
};
@@ -1103,22 +1103,40 @@ async function getAnthropicStreamResponse(
request: unknown,
signal?: AbortSignal,
onSseEvent?: AnthropicOptions["onSseEvent"],
): Promise<{ events: AsyncIterable<RawMessageStreamEvent>; response: Response; requestId: string | null }> {
): Promise<{
events: AsyncIterable<RawMessageStreamEvent>;
response: Response;
requestId: string | null;
recordsRawSseEvents: boolean;
}> {
if (hasAnthropicRawResponseRequest(request)) {
const response = await request.asResponse();
return {
events: iterateAnthropicEvents(response, signal, onSseEvent),
response,
requestId: response.headers.get("request-id"),
recordsRawSseEvents: true,
};
}
if (hasAnthropicStreamWithResponseRequest(request)) {
const { data, response, request_id } = await request.withResponse();
return { events: data, response, requestId: request_id };
return { events: data, response, requestId: request_id, recordsRawSseEvents: false };
}
throw new Error("Anthropic SDK request did not expose a stream response");
}
async function* observeDecodedAnthropicSdkEvents(
events: AsyncIterable<RawMessageStreamEvent>,
observer: (event: RawSseEvent) => void,
): AsyncGenerator<RawMessageStreamEvent> {
for await (const event of events) {
const data = JSON.stringify(event);
// Reconstructed from decoded SDK event; not literal wire bytes.
notifyRawSseEvent(observer, { event: event.type, data, raw: [`event: ${event.type}`, `data: ${data}`] });
yield event;
}
}
function getAnthropicCompat(
model: Model<"anthropic-messages">,
): Required<NonNullable<Model<"anthropic-messages">["compat"]>> {
@@ -1285,6 +1303,9 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
let rawRequestDump: RawHttpRequestDump | undefined;
let activeAbortTracker = createAbortSourceTracker(options?.signal);
const onSseEvent = options?.onSseEvent;
const rawSseObserver = onSseEvent ? (event: RawSseEvent) => onSseEvent(event, model) : undefined;
try {
let client: AnthropicMessagesClientLike;
let isOAuthToken: boolean;
@@ -1319,7 +1340,6 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
hasTools: !!context.tools?.length,
thinkingEnabled: options?.thinkingEnabled,
thinkingDisplay: options?.thinkingDisplay,
onSseEvent: options?.onSseEvent,
fetch: options?.fetch,
claudeCodeSessionId: options?.sessionId ?? extractClaudeMetadataSessionId(options?.metadata?.user_id),
});
@@ -1398,16 +1418,14 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
let anthropicStream: AsyncIterable<RawMessageStreamEvent>;
let response: Response;
let requestId: string | null;
let recordsRawSseEvents: boolean;
try {
({
events: anthropicStream,
response,
requestId,
} = await getAnthropicStreamResponse(
anthropicRequest,
requestSignal,
options?.client ? event => options?.onSseEvent?.(event, model) : undefined,
));
recordsRawSseEvents,
} = await getAnthropicStreamResponse(anthropicRequest, requestSignal, rawSseObserver));
} catch (error) {
if (error instanceof AnthropicConnectionTimeoutError && !activeAbortTracker.wasCallerAbort()) {
throw firstEventTimeoutAbortError;
@@ -1421,7 +1439,7 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
let sawMessageStart = false;
let sawTerminalEnvelope = false;
for await (const event of iterateWithIdleTimeout(anthropicStream, {
const timedAnthropicStream = iterateWithIdleTimeout(anthropicStream, {
idleTimeoutMs,
firstItemTimeoutMs: firstEventTimeoutMs,
errorMessage: idleTimeoutAbortError.message,
@@ -1429,7 +1447,12 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
onIdle: () => activeAbortTracker.abortLocally(idleTimeoutAbortError),
onFirstItemTimeout: () => activeAbortTracker.abortLocally(firstEventTimeoutAbortError),
abortSignal: options?.signal,
})) {
});
const observedAnthropicStream =
rawSseObserver && !recordsRawSseEvents
? observeDecodedAnthropicSdkEvents(timedAnthropicStream, rawSseObserver)
: timedAnthropicStream;
for await (const event of observedAnthropicStream) {
sawEvent = true;
if (event.type === "message_start") {
@@ -1848,7 +1871,6 @@ export function buildAnthropicClientOptions(args: AnthropicClientOptionsArgs): A
thinkingEnabled = false,
thinkingDisplay,
isOAuth,
onSseEvent,
claudeCodeSessionId,
} = args;
const compat = getAnthropicCompat(model);
@@ -1862,7 +1884,6 @@ export function buildAnthropicClientOptions(args: AnthropicClientOptionsArgs): A
// Only OAuth requests inject the CC billing header; no API-key request can ever
// contain it, so there is no need to install the rewriter for those.
const cchFetch = oauthToken ? wrapFetchForCch(baseFetch) : baseFetch;
const debugFetch = onSseEvent ? wrapFetchForSseDebug(cchFetch, event => onSseEvent(event, model)) : cchFetch;
if (model.provider === "github-copilot") {
const copilotApiKey = parseGitHubCopilotApiKey(apiKey).accessToken;
const betaFeatures = [...extraBetas];
@@ -1888,7 +1909,7 @@ export function buildAnthropicClientOptions(args: AnthropicClientOptionsArgs): A
baseURL: baseUrl,
maxRetries: 5,
defaultHeaders,
fetch: debugFetch,
fetch: cchFetch,
...(tlsFetchOptions ? { fetchOptions: tlsFetchOptions } : {}),
};
}
@@ -1923,7 +1944,7 @@ export function buildAnthropicClientOptions(args: AnthropicClientOptionsArgs): A
baseURL: baseUrl,
maxRetries: 5,
defaultHeaders,
fetch: debugFetch,
fetch: cchFetch,
};
}
@@ -1939,7 +1960,7 @@ export function buildAnthropicClientOptions(args: AnthropicClientOptionsArgs): A
baseURL: baseUrl,
maxRetries: 5,
defaultHeaders,
...(debugFetch ? { fetch: debugFetch } : {}),
fetch: cchFetch,
...(tlsFetchOptions ? { fetchOptions: tlsFetchOptions } : {}),
};
}
@@ -1954,7 +1975,7 @@ export function buildAnthropicClientOptions(args: AnthropicClientOptionsArgs): A
baseURL: baseUrl,
maxRetries: 5,
defaultHeaders,
...(debugFetch ? { fetch: debugFetch } : {}),
fetch: cchFetch,
...(tlsFetchOptions ? { fetchOptions: tlsFetchOptions } : {}),
};
}
@@ -1966,7 +1987,7 @@ export function buildAnthropicClientOptions(args: AnthropicClientOptionsArgs): A
baseURL: baseUrl,
maxRetries: 5,
defaultHeaders,
fetch: debugFetch,
fetch: cchFetch,
...(tlsFetchOptions ? { fetchOptions: tlsFetchOptions } : {}),
};
}
@@ -11,6 +11,7 @@ import type {
AssistantMessage,
Context,
Model,
RawSseEvent,
ServiceTier,
StreamFunction,
StreamOptions,
@@ -27,7 +28,7 @@ import {
iterateWithIdleTimeout,
} from "../utils/idle-iterator";
import { sanitizeSchemaForOpenAIResponses, toolWireSchema } from "../utils/schema";
import { wrapFetchForSseDebug } from "../utils/sse-debug";
import { notifyRawSseEvent } from "../utils/sse-debug";
import { mapToOpenAIResponsesToolChoice } from "../utils/tool-choice";
import { normalizeOpenAIResponsesPromptCacheKey, supportsDeveloperRole } from "./openai-responses";
import {
@@ -89,6 +90,18 @@ type AzureOpenAIResponsesSamplingParams = ResponseCreateParamsStreaming & {
repetition_penalty?: number;
};
async function* observeDecodedAzureResponsesEvents(
events: AsyncIterable<ResponseStreamEvent>,
observer: (event: RawSseEvent) => void,
): AsyncGenerator<ResponseStreamEvent> {
for await (const event of events) {
const data = JSON.stringify(event);
// Reconstructed from decoded SDK event; not literal wire bytes.
notifyRawSseEvent(observer, { event: event.type, data, raw: [`event: ${event.type}`, `data: ${data}`] });
yield event;
}
}
/**
* Generate function for Azure OpenAI Responses API
*/
@@ -114,6 +127,8 @@ export const streamAzureOpenAIResponses: StreamFunction<"azure-openai-responses"
const abortTracker = createAbortSourceTracker(options?.signal);
const firstEventTimeoutAbortError = new Error(AZURE_OPENAI_RESPONSES_FIRST_EVENT_TIMEOUT_MESSAGE);
const { requestAbortController, requestSignal } = abortTracker;
const onSseEvent = options?.onSseEvent;
const rawSseObserver = onSseEvent ? (event: RawSseEvent) => onSseEvent(event, model) : undefined;
try {
// Create Azure OpenAI client
@@ -156,26 +171,24 @@ export const streamAzureOpenAIResponses: StreamFunction<"azure-openai-responses"
}
stream.push({ type: "start", partial: output });
await processResponsesStream(
iterateWithIdleTimeout(openaiStream, {
idleTimeoutMs,
firstItemTimeoutMs: firstEventTimeoutMs,
firstItemErrorMessage: AZURE_OPENAI_RESPONSES_FIRST_EVENT_TIMEOUT_MESSAGE,
errorMessage: "Azure OpenAI responses stream stalled while waiting for the next event",
onIdle: () => requestAbortController.abort(),
onFirstItemTimeout: () => abortTracker.abortLocally(firstEventTimeoutAbortError),
abortSignal: options?.signal,
isProgressItem: isOpenAIResponsesProgressEvent,
}),
output,
stream,
model,
{
onFirstToken: () => {
if (!firstTokenTime) firstTokenTime = Date.now();
},
const timedOpenaiStream = iterateWithIdleTimeout(openaiStream, {
idleTimeoutMs,
firstItemTimeoutMs: firstEventTimeoutMs,
firstItemErrorMessage: AZURE_OPENAI_RESPONSES_FIRST_EVENT_TIMEOUT_MESSAGE,
errorMessage: "Azure OpenAI responses stream stalled while waiting for the next event",
onIdle: () => requestAbortController.abort(),
onFirstItemTimeout: () => abortTracker.abortLocally(firstEventTimeoutAbortError),
abortSignal: options?.signal,
isProgressItem: isOpenAIResponsesProgressEvent,
});
const observedOpenaiStream = rawSseObserver
? observeDecodedAzureResponsesEvents(timedOpenaiStream, rawSseObserver)
: timedOpenaiStream;
await processResponsesStream(observedOpenaiStream, output, stream, model, {
onFirstToken: () => {
if (!firstTokenTime) firstTokenTime = Date.now();
},
);
});
const firstEventTimeoutError = abortTracker.getLocalAbortReason();
if (firstEventTimeoutError) {
@@ -269,7 +282,6 @@ function createClient(model: Model<"azure-openai-responses">, apiKey: string, op
const { baseUrl, apiVersion } = resolveAzureConfig(model, options);
const baseFetch = options?.fetch ?? fetch;
const onSseEvent = options?.onSseEvent;
return new AzureOpenAI({
apiKey,
apiVersion,
@@ -277,7 +289,7 @@ function createClient(model: Model<"azure-openai-responses">, apiKey: string, op
maxRetries: 5,
defaultHeaders: headers,
baseURL: baseUrl,
fetch: onSseEvent ? wrapFetchForSseDebug(baseFetch, event => onSseEvent(event, model)) : baseFetch,
fetch: baseFetch,
});
}
+26 -15
View File
@@ -23,6 +23,7 @@ import {
type Model,
type OpenAICompat,
type ProviderSessionState,
type RawSseEvent,
resolveServiceTier,
type ServiceTier,
type StopReason,
@@ -57,7 +58,7 @@ import { getKimiCommonHeaders } from "../utils/oauth/kimi";
import { notifyProviderResponse } from "../utils/provider-response";
import { callWithCopilotModelRetry } from "../utils/retry";
import { adaptSchemaForStrict, NO_STRICT, toolWireSchema } from "../utils/schema";
import { wrapFetchForSseDebug } from "../utils/sse-debug";
import { notifyRawSseEvent } from "../utils/sse-debug";
import {
getStreamMarkupHealingPattern,
type HealedToolCall,
@@ -406,6 +407,20 @@ export function getOpenAICompletionsStreamIdleTimeoutFallbackMs(
return undefined;
}
async function* observeDecodedOpenAICompletionChunks(
chunks: AsyncIterable<ChatCompletionChunk>,
observer: (event: RawSseEvent) => void,
): AsyncGenerator<ChatCompletionChunk> {
for await (const chunk of chunks) {
const data = JSON.stringify(chunk);
const event = typeof chunk.object === "string" ? chunk.object : null;
const raw = event === null ? [`data: ${data}`] : [`event: ${event}`, `data: ${data}`];
// Reconstructed from decoded SDK event; not literal wire bytes.
notifyRawSseEvent(observer, { event, data, raw });
yield chunk;
}
}
export const streamOpenAICompletions: StreamFunction<"openai-completions"> = (
model: Model<"openai-completions">,
context: Context,
@@ -423,6 +438,8 @@ export const streamOpenAICompletions: StreamFunction<"openai-completions"> = (
const abortTracker = createAbortSourceTracker(options?.signal);
const firstEventTimeoutAbortError = new Error(OPENAI_COMPLETIONS_FIRST_EVENT_TIMEOUT_MESSAGE);
const { requestAbortController, requestSignal } = abortTracker;
const onSseEvent = options?.onSseEvent;
const rawSseObserver = onSseEvent ? (event: RawSseEvent) => onSseEvent(event, model) : undefined;
try {
const apiKey = options?.apiKey || getEnvApiKey(model.provider) || "";
@@ -439,15 +456,7 @@ export const streamOpenAICompletions: StreamFunction<"openai-completions"> = (
requestHeaders,
getCapturedErrorResponse: captureErrorResponse,
clearCapturedErrorResponse,
} = await createClient(
model,
context,
apiKey,
options?.headers,
options?.initiatorOverride,
options?.onSseEvent,
options?.fetch,
);
} = await createClient(model, context, apiKey, options?.headers, options?.initiatorOverride, options?.fetch);
const premiumRequestsTotal = copilotPremiumRequests;
getCapturedErrorResponse = captureErrorResponse;
let appliedToolStrictMode: AppliedToolStrictMode = "mixed";
@@ -720,7 +729,7 @@ export const streamOpenAICompletions: StreamFunction<"openai-completions"> = (
for (const call of calls) emitHealedToolCall(call);
};
for await (const chunk of iterateWithIdleTimeout(openaiStream, {
const timedOpenaiStream = iterateWithIdleTimeout(openaiStream, {
idleTimeoutMs,
firstItemTimeoutMs: firstEventTimeoutMs,
firstItemErrorMessage: OPENAI_COMPLETIONS_FIRST_EVENT_TIMEOUT_MESSAGE,
@@ -729,7 +738,11 @@ export const streamOpenAICompletions: StreamFunction<"openai-completions"> = (
onFirstItemTimeout: () => abortTracker.abortLocally(firstEventTimeoutAbortError),
abortSignal: options?.signal,
isProgressItem: isOpenAICompletionsProgressChunk,
})) {
});
const observedOpenaiStream = rawSseObserver
? observeDecodedOpenAICompletionChunks(timedOpenaiStream, rawSseObserver)
: timedOpenaiStream;
for await (const chunk of observedOpenaiStream) {
if (!chunk || typeof chunk !== "object") continue;
// OpenAI documents ChatCompletionChunk.id as the unique chat completion identifier,
@@ -987,7 +1000,6 @@ async function createClient(
apiKey?: string,
extraHeaders?: Record<string, string>,
initiatorOverride?: MessageAttribution,
onSseEvent?: OpenAICompletionsOptions["onSseEvent"],
fetchOverride?: FetchImpl,
): Promise<{
client: OpenAI;
@@ -1086,7 +1098,6 @@ async function createClient(
},
baseFetch.preconnect ? { preconnect: baseFetch.preconnect } : {},
);
const debugFetch = onSseEvent ? wrapFetchForSseDebug(wrappedFetch, event => onSseEvent(event, model)) : wrappedFetch;
return {
client: new OpenAI({
apiKey,
@@ -1095,7 +1106,7 @@ async function createClient(
maxRetries: 5,
defaultHeaders: headers,
defaultQuery: azureDefaultQuery,
fetch: debugFetch,
fetch: wrappedFetch,
}),
copilotPremiumRequests,
baseUrl,
+38 -26
View File
@@ -4,6 +4,7 @@ import type {
Tool as OpenAITool,
ResponseCreateParamsStreaming,
ResponseInput,
ResponseStreamEvent,
} from "openai/resources/responses/responses";
import { getEnvApiKey } from "../stream";
import type {
@@ -15,6 +16,7 @@ import type {
Model,
OpenAICompat,
ProviderSessionState,
RawSseEvent,
ServiceTier,
StreamFunction,
StreamOptions,
@@ -42,7 +44,7 @@ import { notifyProviderResponse } from "../utils/provider-response";
import { callWithCopilotModelRetry } from "../utils/retry";
import { adaptSchemaForStrict, NO_STRICT, sanitizeSchemaForOpenAIResponses, toolWireSchema } from "../utils/schema";
import { createSdkStreamRequestOptions } from "../utils/sdk-stream-timeout";
import { wrapFetchForSseDebug } from "../utils/sse-debug";
import { notifyRawSseEvent } from "../utils/sse-debug";
import { mapToOpenAIResponsesToolChoice, type OpenAIResponsesToolChoice } from "../utils/tool-choice";
import {
buildCopilotDynamicHeaders,
@@ -184,6 +186,18 @@ type OpenAIResponsesSamplingParams = ResponseCreateParamsStreaming & {
stream_options?: { include_obfuscation?: boolean };
};
async function* observeDecodedOpenAIResponsesEvents(
events: AsyncIterable<ResponseStreamEvent>,
observer: (event: RawSseEvent) => void,
): AsyncGenerator<ResponseStreamEvent> {
for await (const event of events) {
const data = JSON.stringify(event);
// Reconstructed from decoded SDK event; not literal wire bytes.
notifyRawSseEvent(observer, { event: event.type, data, raw: [`event: ${event.type}`, `data: ${data}`] });
yield event;
}
}
/**
* Generate function for OpenAI Responses API
*/
@@ -208,6 +222,8 @@ export const streamOpenAIResponses: StreamFunction<"openai-responses"> = (
const abortTracker = createAbortSourceTracker(options?.signal);
const firstEventTimeoutAbortError = new Error(OPENAI_RESPONSES_FIRST_EVENT_TIMEOUT_MESSAGE);
const { requestAbortController, requestSignal } = abortTracker;
const onSseEvent = options?.onSseEvent;
const rawSseObserver = onSseEvent ? (event: RawSseEvent) => onSseEvent(event, model) : undefined;
try {
// Keep request routing on `sessionId` while allowing callers to pin a
@@ -222,7 +238,6 @@ export const streamOpenAIResponses: StreamFunction<"openai-responses"> = (
options?.headers,
options?.initiatorOverride,
routingSessionId,
options?.onSseEvent,
options?.fetch,
);
const premiumRequestsTotal = copilotPremiumRequests;
@@ -273,29 +288,27 @@ export const streamOpenAIResponses: StreamFunction<"openai-responses"> = (
stream.push({ type: "start", partial: output });
const nativeOutputItems: Array<Record<string, unknown>> = [];
await processResponsesStream(
iterateWithIdleTimeout(openaiStream, {
idleTimeoutMs,
firstItemTimeoutMs: firstEventTimeoutMs,
firstItemErrorMessage: OPENAI_RESPONSES_FIRST_EVENT_TIMEOUT_MESSAGE,
errorMessage: "OpenAI responses stream stalled while waiting for the next event",
onFirstItemTimeout: () => abortTracker.abortLocally(firstEventTimeoutAbortError),
onIdle: () => requestAbortController.abort(),
abortSignal: options?.signal,
isProgressItem: isOpenAIResponsesProgressEvent,
}),
output,
stream,
model,
{
onFirstToken: () => {
if (!firstTokenTime) firstTokenTime = Date.now();
},
onOutputItemDone: item => {
nativeOutputItems.push(structuredCloneJSON<unknown>(item) as unknown as Record<string, unknown>);
},
const timedOpenaiStream = iterateWithIdleTimeout(openaiStream, {
idleTimeoutMs,
firstItemTimeoutMs: firstEventTimeoutMs,
firstItemErrorMessage: OPENAI_RESPONSES_FIRST_EVENT_TIMEOUT_MESSAGE,
errorMessage: "OpenAI responses stream stalled while waiting for the next event",
onFirstItemTimeout: () => abortTracker.abortLocally(firstEventTimeoutAbortError),
onIdle: () => requestAbortController.abort(),
abortSignal: options?.signal,
isProgressItem: isOpenAIResponsesProgressEvent,
});
const observedOpenaiStream = rawSseObserver
? observeDecodedOpenAIResponsesEvents(timedOpenaiStream, rawSseObserver)
: timedOpenaiStream;
await processResponsesStream(observedOpenaiStream, output, stream, model, {
onFirstToken: () => {
if (!firstTokenTime) firstTokenTime = Date.now();
},
);
onOutputItemDone: item => {
nativeOutputItems.push(structuredCloneJSON<unknown>(item) as unknown as Record<string, unknown>);
},
});
if (premiumRequestsTotal !== undefined) output.usage.premiumRequests = premiumRequestsTotal;
const firstEventTimeoutError = abortTracker.getLocalAbortReason();
@@ -341,7 +354,6 @@ function createClient(
extraHeaders?: Record<string, string>,
initiatorOverride?: MessageAttribution,
sessionId?: string,
onSseEvent?: OpenAIResponsesOptions["onSseEvent"],
fetchOverride?: FetchImpl,
): {
client: OpenAI;
@@ -388,7 +400,7 @@ function createClient(
dangerouslyAllowBrowser: true,
maxRetries: 5,
defaultHeaders: headers,
fetch: onSseEvent ? wrapFetchForSseDebug(baseFetch, event => onSseEvent(event, model)) : baseFetch,
fetch: baseFetch,
}),
copilotPremiumRequests,
baseUrl,
-271
View File
@@ -1,9 +1,6 @@
import type { ServerSentEvent } from "@oh-my-pi/pi-utils";
import type { RawSseEvent } from "../types";
type FetchFunction = (input: string | URL | Request, init?: RequestInit) => Promise<Response>;
type FetchWithPreconnect = FetchFunction & { preconnect?: typeof fetch.preconnect };
type RawSseObserver = (event: RawSseEvent) => void;
export function notifyRawSseEvent(observer: RawSseObserver | undefined, event: ServerSentEvent | RawSseEvent): void {
@@ -19,271 +16,3 @@ export function notifyRawSseEvent(observer: RawSseObserver | undefined, event: S
// Raw stream observers are diagnostic only and must not affect generation.
}
}
function isSseResponse(response: Response): boolean {
// `response.body` is non-null for any fetch Response with a body, but we
// still guard because user-supplied `fetch` mocks may return `{ body: null }`
// for empty responses and we don't want to wrap those.
if (!response.ok || !response.body) return false;
const contentType = response.headers.get("content-type");
// All providers in this repo emit lowercase `text/event-stream` (verified
// against anthropic, openai-completions, openai-responses, azure-openai-responses,
// google-shared, google-gemini-cli, openai-codex-responses, pi-native-client,
// and the auth-gateway server). A canonical `includes` check is sufficient;
// if a future provider sends mixed case it will fall back to the unwrapped
// fetch — observably safe, just no debug tee for that response.
return contentType?.includes("text/event-stream") ?? false;
}
// Reused for every UTF-8 line decode. Safe because lines are split on LF
// (0x0a), which is single-byte ASCII and never appears inside a UTF-8
// multi-byte sequence — each line is a complete UTF-8 run, so the decoder
// carries no state across calls.
const SSE_LINE_DECODER = new TextDecoder("utf-8");
// Decode bytes [start, end) of an SSE line.
//
// A previous revision added an ASCII fast-path using `String.fromCharCode.apply`
// over chunked subarrays, on the theory that skipping `TextDecoder` would save
// the ~9.7% `decode` self-time the profile reported. In practice the swap
// *regressed* total wall time: `fromCharCode` became a new 7.8% hotspot,
// `Uint8Array` allocations grew 5.3%, and `subarray` rose from 11.5% to 18.3%
// — net loss of ~10pp. Bun's `TextDecoder.decode` has a fast C++ ASCII path
// that beats chunked `fromCharCode.apply` for the typical sub-1KB SSE line,
// so we keep the decoder. The line is bounded by LF (0x0a, single-byte
// ASCII), so each [start, end) slice is a complete UTF-8 run and the shared
// stateless decoder is safe to reuse.
function decodeSseLine(buf: Uint8Array, start: number, end: number): string {
if (start === 0 && end === buf.length) return SSE_LINE_DECODER.decode(buf);
return SSE_LINE_DECODER.decode(buf.subarray(start, end));
}
/**
* Inline SSE event splitter. Walks the byte stream as it flows through a
* `TransformStream`, dispatching parsed events to the debug observer while
* the bytes are forwarded unchanged to the response consumer. Replaces the
* previous `body.tee()` + `readSseEvents` re-parse pipeline so the byte
* stream is parsed exactly once when a debug observer is attached.
*
* Field parsing intentionally mirrors `readSseEvents` in `@oh-my-pi/pi-utils`
* (only `event` and `data` are observed; `id`/`retry` ignored; CR stripped
* before LF dispatch; leading space after `:` trimmed; `data:` lines join
* with `\n`). Reusing `readSseEvents` directly would require a second stream
* pipeline, which is exactly what this class avoids.
*/
class SseTeeParser {
#observer: RawSseObserver;
// Trailing bytes from the previous chunk that did not end with LF.
#partial: Uint8Array | null = null;
#event: string | null = null;
#data: string | null = null;
#raw: string[] = [];
constructor(observer: RawSseObserver) {
this.#observer = observer;
}
push(chunk: Uint8Array): void {
// Carry-forward path: concat the partial line with the new chunk so the
// LF scan walks a single contiguous buffer. The common case (partial is
// null) skips the allocation entirely.
let buf: Uint8Array;
if (this.#partial) {
buf = new Uint8Array(this.#partial.length + chunk.length);
buf.set(this.#partial, 0);
buf.set(chunk, this.#partial.length);
this.#partial = null;
} else {
buf = chunk;
}
const len = buf.length;
let i = 0;
while (i < len) {
const lf = buf.indexOf(0x0a, i);
if (lf === -1) {
// Retain the tail as a partial line for the next chunk. Copy
// because the source `chunk` buffer may be reused upstream.
this.#partial = buf.subarray(i).slice();
return;
}
let end = lf;
if (end > i && buf[end - 1] === 0x0d) end--;
this.#consumeLine(buf, i, end);
i = lf + 1;
}
}
flush(): void {
// Treat any trailing partial line (no terminating LF) as a complete line.
if (this.#partial) {
const tail = this.#partial;
this.#partial = null;
let end = tail.length;
if (end > 0 && tail[end - 1] === 0x0d) end--;
if (end > 0) this.#consumeLine(tail, 0, end);
}
// Real services don't always close on a blank line — flush any pending event.
this.#dispatch();
}
#consumeLine(buf: Uint8Array, start: number, end: number): void {
if (end === start) {
this.#dispatch();
return;
}
// Comment line: keep verbatim in `raw` for diagnostic context, skip parsing.
// SSE spec § 9.2.6: lines beginning with ':' are heartbeats/comments and
// MUST NOT contribute to the event dispatch state. Heartbeats are the
// single most common line type on long-poll provider streams, so the
// early-return here directly avoids ~half the field-parse work.
if (buf[start] === 0x3a /* ':' */) {
this.#raw.push(decodeSseLine(buf, start, end));
return;
}
// Byte-level field parse. We avoid `text.indexOf(':')` + two `String.slice`
// calls (~6% of CPU pre-optimization) by scanning bytes for the field
// delimiter and matching the field name byte-for-byte. Field-name bytes
// are ASCII per SSE spec, so byte offsets equal char offsets in the
// decoded string and we can `slice` the value directly off `text` without
// re-decoding.
//
// ASCII signatures (verified against SSE spec):
// "event" = 0x65 0x76 0x65 0x6e 0x74 (5 bytes)
// "data" = 0x64 0x61 0x74 0x61 (4 bytes)
let colon = -1;
for (let k = start; k < end; k++) {
if (buf[k] === 0x3a) {
colon = k;
break;
}
}
const fieldEnd = colon === -1 ? end : colon;
let valueStart = colon === -1 ? end : colon + 1;
// Per SSE spec, a single leading SP after the colon is stripped.
if (valueStart < end && buf[valueStart] === 0x20 /* ' ' */) valueStart++;
const fieldLen = fieldEnd - start;
const isEvent =
fieldLen === 5 &&
buf[start] === 0x65 &&
buf[start + 1] === 0x76 &&
buf[start + 2] === 0x65 &&
buf[start + 3] === 0x6e &&
buf[start + 4] === 0x74;
const isData =
!isEvent &&
fieldLen === 4 &&
buf[start] === 0x64 &&
buf[start + 1] === 0x61 &&
buf[start + 2] === 0x74 &&
buf[start + 3] === 0x61;
// Decode the line exactly once. Raw observers (debug buffer) want it
// regardless of field kind; `id`/`retry`/unknown lines pay only the
// decode cost, not any extra slicing.
const text = decodeSseLine(buf, start, end);
this.#raw.push(text);
if (isEvent) {
// `valueStart - start` is a byte offset into the line; since the
// "event:" prefix (and the optional SP) are pure ASCII, that byte
// offset equals the char offset in the decoded `text`.
this.#event = valueStart === end ? "" : text.slice(valueStart - start);
} else if (isData) {
const value = valueStart === end ? "" : text.slice(valueStart - start);
if (this.#data === null) this.#data = value;
else this.#data = `${this.#data}\n${value}`;
}
// `id` and `retry` are intentionally ignored — providers don't use them
// and reconnects are handled by the underlying transport.
}
// Hands ownership of the accumulated `raw` array to the observer. The
// observer (currently only `RawSseDebugBuffer.recordEvent`) MAY retain the
// array; we install a fresh `#raw = []` for the next event before invoking
// the observer so there is no aliasing across dispatches. This contract is
// mirrored in `notifyRawSseEvent` (no defensive clone) — see its comment.
//
// TODO(BufferOpt): once the buffer-side audit confirms it never mutates
// `event.raw`, the defensive `[...event.raw]` clone in older call paths
// (search for `notifyRawSseEvent`) can be dropped repository-wide.
#dispatch(): void {
if (this.#event === null && this.#data === null) return;
const event: RawSseEvent = {
event: this.#event,
data: this.#data ?? "",
raw: this.#raw,
};
this.#event = null;
this.#data = null;
this.#raw = [];
try {
this.#observer(event);
} catch {
// Raw stream observers are diagnostic only and must not affect generation.
}
}
}
export function wrapFetchForSseDebug(
fetchImpl: FetchWithPreconnect,
observer: RawSseObserver | undefined,
): FetchWithPreconnect {
if (!observer) return fetchImpl;
const wrapped = Object.assign(
async (input: string | URL | Request, init?: RequestInit): Promise<Response> => {
const response = await fetchImpl(input, init);
if (!isSseResponse(response)) {
return response;
}
const body = response.body;
if (!body) return response;
// Single-pass interception. Previously implemented as
// `body.pipeThrough(new TransformStream({...}))`, but the WHATWG
// TransformStream machinery imposes a per-chunk Promise boundary
// (`#handleNumberResult` showed at 8.8% self-time in CPU profile).
// A manual ReadableStream pulling directly from `body.getReader()`
// skips that hop: every `read()` immediately feeds both the parser
// and the controller in the same microtask.
const parser = new SseTeeParser(observer);
const reader = body.getReader();
const teed = new ReadableStream<Uint8Array>({
async pull(controller) {
try {
const { done, value } = await reader.read();
if (done) {
parser.flush();
controller.close();
return;
}
// Enqueue first so the consumer sees bytes ASAP; parser
// dispatch is best-effort diagnostic and runs after.
controller.enqueue(value);
parser.push(value);
} catch (err) {
// Mirror TransformStream semantics: surface upstream
// errors to the consumer; do not flush a partial event.
controller.error(err);
}
},
cancel(reason) {
// Propagate downstream cancellation to the source body so the
// underlying connection is released. Matches `pipeThrough`'s
// cancel-propagation behavior; `flush()` is intentionally NOT
// called (TransformStream skips `flush` on abort too).
return reader.cancel(reason);
},
});
return new Response(teed, {
status: response.status,
statusText: response.statusText,
headers: response.headers,
});
},
fetchImpl.preconnect ? { preconnect: fetchImpl.preconnect } : {},
);
return wrapped;
}
@@ -0,0 +1,283 @@
import { afterEach, describe, expect, it, vi } from "bun:test";
import { getBundledModel } from "../src/models";
import { streamAnthropic } from "../src/providers/anthropic";
import type { AnthropicMessagesClientLike } from "../src/providers/anthropic-client";
import type { RawMessageStreamEvent } from "../src/providers/anthropic-wire";
import { streamAzureOpenAIResponses } from "../src/providers/azure-openai-responses";
import { streamOpenAICompletions } from "../src/providers/openai-completions";
import { streamOpenAIResponses } from "../src/providers/openai-responses";
import type { Context, Model, RawSseEvent } from "../src/types";
const originalFetch = global.fetch;
const context: Context = {
messages: [{ role: "user", content: "Say hello", timestamp: Date.now() }],
};
const openAIResponsesModel = getBundledModel("openai", "gpt-5-mini") as Model<"openai-responses">;
const openAICompletionsModel = {
...(getBundledModel("openai", "gpt-4o-mini") as Model<"openai-completions">),
api: "openai-completions",
} satisfies Model<"openai-completions">;
const azureOpenAIResponsesModel: Model<"azure-openai-responses"> = {
id: "gpt-5-mini",
name: "GPT-5 Mini",
api: "azure-openai-responses",
provider: "azure",
baseUrl: "https://example.openai.azure.com/openai/v1",
reasoning: false,
input: ["text"],
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
contextWindow: 400_000,
maxTokens: 128_000,
};
const anthropicModel: Model<"anthropic-messages"> = {
id: "claude-sonnet-4-5",
name: "Claude Sonnet 4.5",
api: "anthropic-messages",
provider: "anthropic",
baseUrl: "https://api.anthropic.com",
reasoning: true,
input: ["text"],
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
contextWindow: 200_000,
maxTokens: 8_192,
};
const openAIResponsesEvents = [
{ type: "response.created", response: { id: "resp_raw_sse", status: "in_progress" } },
{
type: "response.output_item.added",
item: { type: "message", id: "msg_raw_sse", role: "assistant", status: "in_progress", content: [] },
},
{ type: "response.content_part.added", part: { type: "output_text", text: "" } },
{ type: "response.output_text.delta", delta: "Hello" },
{
type: "response.output_item.done",
item: {
type: "message",
id: "msg_raw_sse",
role: "assistant",
status: "completed",
content: [{ type: "output_text", text: "Hello" }],
},
},
{
type: "response.completed",
response: {
id: "resp_raw_sse",
status: "completed",
usage: {
input_tokens: 5,
output_tokens: 1,
total_tokens: 6,
input_tokens_details: { cached_tokens: 0 },
},
},
},
];
const anthropicEvents: RawMessageStreamEvent[] = [
{
type: "message_start",
message: {
id: "msg_raw_sse",
usage: {
input_tokens: 5,
output_tokens: 0,
cache_read_input_tokens: 0,
cache_creation_input_tokens: 0,
},
},
},
{ type: "content_block_start", index: 0, content_block: { type: "text", text: "" } },
{ type: "content_block_delta", index: 0, delta: { type: "text_delta", text: "Hello" } },
{ type: "content_block_stop", index: 0 },
{
type: "message_delta",
delta: { stop_reason: "end_turn" },
usage: {
input_tokens: 5,
output_tokens: 1,
cache_read_input_tokens: 0,
cache_creation_input_tokens: 0,
},
},
{ type: "message_stop" },
];
function createSseResponse(events: unknown[]): Response {
const payload = `${events
.map(event => `data: ${typeof event === "string" ? event : JSON.stringify(event)}`)
.join("\n\n")}\n\n`;
return new Response(payload, {
status: 200,
headers: { "content-type": "text/event-stream" },
});
}
function installFetchResponse(events: unknown[]) {
const fetchMock = vi.fn(async () => createSseResponse(events));
global.fetch = Object.assign(fetchMock, { preconnect: originalFetch.preconnect }) as typeof fetch;
return fetchMock;
}
function recordEvent(events: RawSseEvent[]): (event: RawSseEvent) => void {
return event => {
events.push({ event: event.event, data: event.data, raw: [...event.raw] });
};
}
async function* asyncEvents(events: RawMessageStreamEvent[]): AsyncGenerator<RawMessageStreamEvent> {
for (const event of events) yield event;
}
function createAnthropicSdkClient(events: RawMessageStreamEvent[]): AnthropicMessagesClientLike {
return {
messages: {
create: () => ({
async withResponse() {
return {
data: asyncEvents(events),
response: new Response(null, { status: 200, headers: { "request-id": "req_sdk" } }),
request_id: "req_sdk",
};
},
}),
},
};
}
function sseFrame(event: string, data: unknown): string {
return `event: ${event}\ndata: ${JSON.stringify(data)}\n\n`;
}
function createAnthropicRawClient(events: RawMessageStreamEvent[]): AnthropicMessagesClientLike {
return {
messages: {
create: () => ({
async asResponse() {
return new Response(events.map(event => sseFrame(event.type, event)).join(""), {
status: 200,
headers: { "content-type": "text/event-stream", "request-id": "req_raw" },
});
},
}),
},
};
}
afterEach(() => {
global.fetch = originalFetch;
vi.restoreAllMocks();
});
describe("SDK raw SSE capture", () => {
it("records OpenAI Responses SDK events from the decoded stream", async () => {
const fetchMock = installFetchResponse(openAIResponsesEvents);
const observed: RawSseEvent[] = [];
const result = await streamOpenAIResponses(openAIResponsesModel, context, {
apiKey: "test-key",
onSseEvent: recordEvent(observed),
}).result();
expect(result.stopReason).toBe("stop");
expect(fetchMock).toHaveBeenCalledTimes(1);
expect(observed.map(event => event.event)).toEqual(openAIResponsesEvents.map(event => event.type));
expect(JSON.parse(observed[0]!.data)).toEqual(openAIResponsesEvents[0]);
expect(observed[0]!.raw).toEqual([
"event: response.created",
`data: ${JSON.stringify(openAIResponsesEvents[0])}`,
]);
});
it("records OpenAI Chat Completions SDK events from the decoded stream", async () => {
const chunks = [
{
id: "chatcmpl_raw_sse",
object: "chat.completion.chunk",
created: 0,
model: openAICompletionsModel.id,
choices: [{ index: 0, delta: { content: "Hello" } }],
},
{
id: "chatcmpl_raw_sse",
object: "chat.completion.chunk",
created: 0,
model: openAICompletionsModel.id,
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
usage: {
prompt_tokens: 5,
completion_tokens: 1,
total_tokens: 6,
prompt_tokens_details: { cached_tokens: 0 },
},
},
"[DONE]",
];
installFetchResponse(chunks);
const observed: RawSseEvent[] = [];
const result = await streamOpenAICompletions(openAICompletionsModel, context, {
apiKey: "test-key",
onSseEvent: recordEvent(observed),
}).result();
expect(result.stopReason).toBe("stop");
expect(observed.map(event => event.event)).toEqual(["chat.completion.chunk", "chat.completion.chunk"]);
expect(JSON.parse(observed[0]!.data)).toEqual(chunks[0]);
expect(observed[0]!.raw).toEqual(["event: chat.completion.chunk", `data: ${JSON.stringify(chunks[0])}`]);
});
it("records Azure OpenAI Responses SDK events from the decoded stream", async () => {
installFetchResponse(openAIResponsesEvents);
const observed: RawSseEvent[] = [];
const result = await streamAzureOpenAIResponses(azureOpenAIResponsesModel, context, {
apiKey: "test-key",
azureBaseUrl: azureOpenAIResponsesModel.baseUrl,
azureApiVersion: "v1",
onSseEvent: recordEvent(observed),
}).result();
expect(result.stopReason).toBe("stop");
expect(observed.map(event => event.event)).toEqual(openAIResponsesEvents.map(event => event.type));
expect(JSON.parse(observed.at(-1)!.data)).toEqual(openAIResponsesEvents.at(-1));
});
it("records Anthropic SDK events from the decoded stream", async () => {
const observed: RawSseEvent[] = [];
const result = await streamAnthropic(anthropicModel, context, {
client: createAnthropicSdkClient(anthropicEvents),
onSseEvent: recordEvent(observed),
}).result();
expect(result.stopReason).toBe("stop");
expect(observed.map(event => event.event)).toEqual(anthropicEvents.map(event => event.type));
expect(JSON.parse(observed[0]!.data)).toEqual(anthropicEvents[0]);
expect(observed[0]!.raw).toEqual(["event: message_start", `data: ${JSON.stringify(anthropicEvents[0])}`]);
});
it("does not synthesize raw SSE records when no observer is installed", async () => {
installFetchResponse(openAIResponsesEvents);
const result = await streamOpenAIResponses(openAIResponsesModel, context, { apiKey: "test-key" }).result();
expect(result.stopReason).toBe("stop");
});
it("keeps Anthropic direct SSE parsing wired to the raw observer", async () => {
const observed: RawSseEvent[] = [];
const result = await streamAnthropic(anthropicModel, context, {
client: createAnthropicRawClient(anthropicEvents),
onSseEvent: recordEvent(observed),
}).result();
expect(result.stopReason).toBe("stop");
expect(observed.map(event => event.event)).toEqual(anthropicEvents.map(event => event.type));
expect(observed[0]!.raw).toEqual(["event: message_start", `data: ${JSON.stringify(anthropicEvents[0])}`]);
});
});
+24 -194
View File
@@ -1,205 +1,35 @@
import { describe, expect, it } from "bun:test";
import type { RawSseEvent } from "../src/types";
import { wrapFetchForSseDebug } from "../src/utils/sse-debug";
import { notifyRawSseEvent } from "../src/utils/sse-debug";
/**
* Exercises the inline SSE tee + parser in `sse-debug.ts`. There is no direct
* export for `SseTeeParser`; we drive it through `wrapFetchForSseDebug`, which
* is the only production caller. Each test:
* 1. Builds a mock `fetch` that returns a `text/event-stream` Response whose
* body emits a caller-controlled sequence of byte chunks (so we can
* exercise partial-line carry-forward and CR-LF handling deterministically).
* 2. Calls the wrapped fetch.
* 3. Reads the response body to completion so the `TransformStream` `flush`
* runs.
* 4. Asserts the events the observer received exactly match expectations.
*
* The point is to lock in behavior across the ASCII-fast-path / byte-level-
* field-parse rewrite: the observer MUST receive the same `{ event, data, raw }`
* shape it received with the prior decode-then-string-slice implementation.
*/
describe("notifyRawSseEvent", () => {
it("dispatches diagnostic events without cloning raw lines", () => {
const raw = ["event: message", "data: hello"];
let observed: RawSseEvent | undefined;
function chunkedStream(chunks: Uint8Array[]): ReadableStream<Uint8Array> {
let i = 0;
return new ReadableStream<Uint8Array>({
pull(controller) {
if (i >= chunks.length) {
controller.close();
return;
}
controller.enqueue(chunks[i++]);
},
});
}
notifyRawSseEvent(
event => {
observed = event;
},
{ event: "message", data: "hello", raw },
);
function sseResponse(chunks: Uint8Array[]): Response {
return new Response(chunkedStream(chunks), {
status: 200,
headers: { "content-type": "text/event-stream" },
});
}
const enc = new TextEncoder();
const b = (s: string): Uint8Array => enc.encode(s);
async function drain(response: Response): Promise<void> {
const reader = response.body!.getReader();
for (;;) {
const { done } = await reader.read();
if (done) return;
}
}
async function collect(chunks: Uint8Array[]): Promise<RawSseEvent[]> {
const events: RawSseEvent[] = [];
const fetchImpl = async () => sseResponse(chunks);
const wrapped = wrapFetchForSseDebug(fetchImpl, event => {
events.push(event);
});
const response = await wrapped("https://example.test/stream");
await drain(response);
return events;
}
describe("sse-debug parser", () => {
it("parses a single event terminated by blank line", async () => {
const events = await collect([b("event: message\ndata: hello\n\n")]);
expect(events).toEqual([{ event: "message", data: "hello", raw: ["event: message", "data: hello"] }]);
expect(observed).toEqual({ event: "message", data: "hello", raw });
expect(observed?.raw).toBe(raw);
});
it("joins multi-line data fields with newlines", async () => {
const events = await collect([b("data: line1\ndata: line2\ndata: line3\n\n")]);
expect(events).toHaveLength(1);
expect(events[0]!.event).toBe(null);
expect(events[0]!.data).toBe("line1\nline2\nline3");
expect(events[0]!.raw).toEqual(["data: line1", "data: line2", "data: line3"]);
it("keeps observer failures diagnostic-only", () => {
expect(() =>
notifyRawSseEvent(
() => {
throw new Error("observer failed");
},
{ event: "message", data: "hello", raw: ["event: message", "data: hello"] },
),
).not.toThrow();
});
it("strips a single leading SP after the colon but preserves further spaces", async () => {
const events = await collect([b("data: two-leading-spaces\n\n")]);
expect(events[0]!.data).toBe(" two-leading-spaces");
});
it("retains comment (`:`-prefixed) lines in raw but does not parse them", async () => {
const events = await collect([b(": heartbeat\ndata: payload\n\n")]);
expect(events).toHaveLength(1);
expect(events[0]!.data).toBe("payload");
expect(events[0]!.raw).toEqual([": heartbeat", "data: payload"]);
});
it("does not dispatch on a blank line if no event/data accumulated (pure heartbeats)", async () => {
const events = await collect([b(": ping\n\n: ping\n\n")]);
expect(events).toHaveLength(0);
});
it("handles CR-LF line endings and strips the CR before dispatch", async () => {
const events = await collect([b("event: ping\r\ndata: pong\r\n\r\n")]);
expect(events).toEqual([{ event: "ping", data: "pong", raw: ["event: ping", "data: pong"] }]);
});
it("ignores unknown fields (`id`, `retry`, gibberish) but keeps them in raw", async () => {
const events = await collect([b("id: 42\nretry: 1000\nfoo: bar\ndata: ok\n\n")]);
expect(events).toHaveLength(1);
expect(events[0]!.event).toBe(null);
expect(events[0]!.data).toBe("ok");
expect(events[0]!.raw).toEqual(["id: 42", "retry: 1000", "foo: bar", "data: ok"]);
});
it("treats a line with no colon as field-with-empty-value (data line still recorded)", async () => {
// Per SSE spec a bare `data` line is treated as `data:` with empty value.
const events = await collect([b("data\ndata: x\n\n")]);
expect(events).toHaveLength(1);
expect(events[0]!.data).toBe("\nx");
});
it("reassembles events split across arbitrary chunk boundaries", async () => {
// Split a single event across chunks: mid-field-name, mid-value, mid-LF-CRLF.
const events = await collect([b("eve"), b("nt: x\r"), b("\ndata: a"), b("bc\r\n\r"), b("\n")]);
expect(events).toEqual([{ event: "x", data: "abc", raw: ["event: x", "data: abc"] }]);
});
it("handles a chunk that ends exactly on LF (no partial carried)", async () => {
const events = await collect([b("data: a\n"), b("data: b\n"), b("\n")]);
expect(events).toHaveLength(1);
expect(events[0]!.data).toBe("a\nb");
});
it("flushes a trailing event with no terminating blank line", async () => {
// Stream closes without a final "\n\n". Parser must dispatch on flush.
const events = await collect([b("event: end\ndata: bye\n")]);
expect(events).toEqual([{ event: "end", data: "bye", raw: ["event: end", "data: bye"] }]);
});
it("flushes a trailing event with no terminating newline at all", async () => {
const events = await collect([b("event: end\ndata: bye")]);
expect(events).toEqual([{ event: "end", data: "bye", raw: ["event: end", "data: bye"] }]);
});
it("preserves UTF-8 multibyte characters via decoder fallback", async () => {
// Non-ASCII bytes (emoji, accented chars, CJK) must round-trip identically.
const events = await collect([b("data: caf\u00e9 \u2014 \u4f60\u597d \ud83d\ude00\n\n")]);
expect(events[0]!.data).toBe("café — 你好 😀");
});
it("handles a UTF-8 multibyte sequence split across chunk boundary", async () => {
// The 4-byte emoji U+1F600 ("😀") = F0 9F 98 80. Split it between chunks.
const full = b("data: \ud83d\ude00\n\n");
const split = full.indexOf(0xf0) + 2;
const events = await collect([full.subarray(0, split), full.subarray(split)]);
expect(events[0]!.data).toBe("😀");
});
it("emits multiple events in stream order", async () => {
const events = await collect([b("event: a\ndata: 1\n\nevent: b\ndata: 2\n\nevent: c\ndata: 3\n\n")]);
expect(events.map(e => [e.event, e.data])).toEqual([
["a", "1"],
["b", "2"],
["c", "3"],
]);
});
it("hands a fresh `raw` array to each observer call (no aliasing)", async () => {
const events = await collect([b("data: a\n\ndata: b\n\n")]);
expect(events).toHaveLength(2);
expect(events[0]!.raw).not.toBe(events[1]!.raw);
// Observer-side mutation of the first `raw` must not leak into the second.
events[0]!.raw.push("MUTATED");
expect(events[1]!.raw).toEqual(["data: b"]);
});
it("treats `data:` with no value as empty string and merges further data lines", async () => {
const events = await collect([b("data:\ndata: x\n\n")]);
expect(events[0]!.data).toBe("\nx");
});
it("returns the unwrapped fetch when observer is undefined", async () => {
const fetchImpl = async () => sseResponse([b("data: x\n\n")]);
const wrapped = wrapFetchForSseDebug(fetchImpl, undefined);
// Identity, not a wrapper: caller relies on this fast path.
expect(wrapped).toBe(fetchImpl as unknown as typeof wrapped);
});
it("passes through non-SSE responses untouched", async () => {
const events: RawSseEvent[] = [];
const fetchImpl = async () =>
new Response(b("not sse"), { status: 200, headers: { "content-type": "text/plain" } });
const wrapped = wrapFetchForSseDebug(fetchImpl, e => events.push(e));
const response = await wrapped("https://example.test/plain");
expect(await response.text()).toBe("not sse");
expect(events).toHaveLength(0);
});
it("forwards the byte stream byte-identically to the consumer", async () => {
// Critical invariant: tee must not mutate or re-shape bytes for the
// downstream consumer. Use a payload with UTF-8 + CR-LF + heartbeats to
// stress the parser without corrupting forwarded bytes.
const payload = b(": heartbeat\r\nevent: msg\r\ndata: caf\u00e9 \u4f60\u597d\r\n\r\ndata: tail\n\n");
// Chunk the input awkwardly so the TransformStream sees several chunks.
const chunks = [payload.subarray(0, 5), payload.subarray(5, 17), payload.subarray(17)];
const fetchImpl = async () => sseResponse(chunks);
const wrapped = wrapFetchForSseDebug(fetchImpl, () => {});
const response = await wrapped("https://example.test/stream");
const forwarded = new Uint8Array(await response.arrayBuffer());
expect(Array.from(forwarded)).toEqual(Array.from(payload));
it("is a no-op when no observer is installed", () => {
expect(() => notifyRawSseEvent(undefined, { event: null, data: "{}", raw: ["data: {}"] })).not.toThrow();
});
});
+2 -1
View File
@@ -1,14 +1,15 @@
# Changelog
## [Unreleased]
### Added
- Added `raw-sse.txt` to debug report bundles, exporting recent raw provider SSE diagnostics when captured
- Added `/model` visibility for auto-selected role defaults: inferred `pi/smol`/`pi/slow`/designer choices now show as compact `[ROLE auto]` badges, while explicitly configured roles keep the existing solid badges and thinking labels.
- Added credential provenance to the `/login` and `/logout` provider picker: each authenticated provider now shows where its credential comes from — `(login)`, `(api key)`, `(env: VAR_NAME)`, `(config)`, `(--api-key)`, or `(custom provider)` — so a real OAuth login is distinguishable from an env var that merely aliases the provider (e.g. `COPILOT_GITHUB_TOKEN`). The origin is also matched by the picker's type-to-search filter.
### Changed
- Changed raw SSE debug export output to prepend dropped-record metadata so truncated sessions in debug bundles now report dropped record and character counts
- Changed settings reads to cache pre-split schema paths and resolved values, with coarse invalidation on source/cwd changes.
- Changed status-line rendering to cache merged effective settings until `updateSettings()` changes the configuration.
- Changed `CustomEditor` app shortcut dispatch to parse each input packet once and match against precomputed canonical key sets, preserving the existing shortcut precedence while avoiding repeated key reparses.
+8
View File
@@ -195,6 +195,7 @@ export class DebugSelectorComponent extends Container {
const result = await createReportBundle({
sessionFile: this.ctx.sessionManager.getSessionFile(),
settings: this.#getResolvedSettings(),
rawSseText: this.#getRawSseText(),
cpuProfile,
workProfile,
});
@@ -253,6 +254,7 @@ export class DebugSelectorComponent extends Container {
const result = await createReportBundle({
sessionFile: this.ctx.sessionManager.getSessionFile(),
settings: this.#getResolvedSettings(),
rawSseText: this.#getRawSseText(),
});
loader.stop();
@@ -288,6 +290,7 @@ export class DebugSelectorComponent extends Container {
const result = await createReportBundle({
sessionFile: this.ctx.sessionManager.getSessionFile(),
settings: this.#getResolvedSettings(),
rawSseText: this.#getRawSseText(),
heapSnapshot,
});
@@ -490,6 +493,11 @@ export class DebugSelectorComponent extends Container {
}
}
#getRawSseText(): string | undefined {
const rawSseText = resolveRawSseDebugBuffer(this.ctx.session).toRawText();
return rawSseText.trim().length > 0 ? rawSseText : undefined;
}
#getResolvedSettings(): Record<string, unknown> {
// Extract key settings for the report
return {
@@ -152,9 +152,9 @@ export class RawSseDebugBuffer {
}
// Ownership contract for `event.raw`:
// The caller (either `notifyRawSseEvent` in `packages/ai/src/utils/sse-debug.ts`
// or `SseTeeParser.#dispatch` directly) hands us a freshly-allocated
// `string[]` per event and never retains, mutates, or re-dispatches it.
// The caller (`notifyRawSseEvent` in `packages/ai/src/utils/sse-debug.ts`)
// hands us a freshly-allocated `string[]` per event and never retains,
// mutates, or re-dispatches it.
// That lets `trimRawLines` keep the array by reference instead of
// cloning on every chunk — a measurable savings on the streaming hot
// path. If a future observer-chain mutates the array, restore the
@@ -192,7 +192,10 @@ export class RawSseDebugBuffer {
toRawText(): string {
// Reads the live array directly: `rawRecordText` only computes a string
// from each record, so no caller-visible mutation is possible.
return this.#records.map(rawRecordText).join("\n");
const body = this.#records.map(rawRecordText).join("\n");
if (this.#droppedRecords === 0) return body;
const dropped = `: omp-debug-dropped records=${this.#droppedRecords} chars=${this.#droppedChars}\n\n`;
return body.length > 0 ? `${dropped}${body}` : dropped;
}
#append(record: RawSseDebugRecord, chars: number): void {
@@ -45,6 +45,8 @@ export interface ReportBundleOptions {
heapSnapshot?: HeapSnapshot;
/** Work profile (for work scheduling reports) */
workProfile?: WorkProfile;
/** Raw provider SSE diagnostics captured by the session buffer */
rawSseText?: string;
}
export interface ReportBundleResult {
@@ -70,6 +72,7 @@ export interface DebugLogSource {
* - env.json: Sanitized environment variables
* - config.json: Resolved settings
* - profile.cpuprofile: CPU profile (performance report only)
* - raw-sse.txt: Recent raw provider SSE diagnostics (when captured)
* - profile.md: Markdown CPU profile (performance report only)
* - heap.heapsnapshot: Heap snapshot (memory report only)
* - work.folded: Work profile folded stacks (work report only)
@@ -109,6 +112,12 @@ export async function createReportBundle(options: ReportBundleOptions): Promise<
files.push("logs.txt");
}
// Recent raw provider SSE diagnostics
if (options.rawSseText && options.rawSseText.trim().length > 0) {
data["raw-sse.txt"] = options.rawSseText;
files.push("raw-sse.txt");
}
// Session file
if (options.sessionFile) {
try {
@@ -68,4 +68,27 @@ describe("RawSseDebugBuffer", () => {
expect(resolveRawSseDebugBuffer(owner)).toBe(buffer);
expect(buffer.snapshot().totalEvents).toBe(1);
});
it("keeps session-owned records captured before the viewer resolves the buffer", () => {
const session = { rawSseDebugBuffer: new RawSseDebugBuffer() };
session.rawSseDebugBuffer.recordResponse(
{ status: 200, requestId: "req_pre_viewer", headers: {}, metadata: { lastTransport: "sse" } },
model,
);
session.rawSseDebugBuffer.recordEvent(
{ event: "message_start", data: "{}", raw: ["event: message_start", "data: {}"] },
model,
);
session.rawSseDebugBuffer.recordEvent(
{ event: "message_stop", data: "{}", raw: ["event: message_stop", "data: {}"] },
model,
);
const buffer = resolveRawSseDebugBuffer(session);
expect(buffer).toBe(session.rawSseDebugBuffer);
expect(buffer.snapshot().totalEvents).toBe(2);
expect(buffer.toRawText()).toContain("requestId=req_pre_viewer");
expect(buffer.toRawText()).toContain("event: message_stop");
});
});
@@ -0,0 +1,76 @@
import { afterEach, describe, expect, it } from "bun:test";
import * as fs from "node:fs/promises";
import * as os from "node:os";
import * as path from "node:path";
import type { Model } from "@oh-my-pi/pi-ai";
import { getConfigRootDir, setAgentDir } from "@oh-my-pi/pi-utils";
import { RawSseDebugBuffer } from "../../src/debug/raw-sse-buffer";
import { createReportBundle } from "../../src/debug/report-bundle";
const model: Model<"anthropic-messages"> = {
id: "claude-test",
name: "Claude Test",
api: "anthropic-messages",
provider: "anthropic",
baseUrl: "https://api.anthropic.com",
reasoning: true,
input: ["text"],
cost: { input: 1, output: 1, cacheRead: 0, cacheWrite: 0 },
contextWindow: 200_000,
maxTokens: 8_192,
};
const originalAgentDir = process.env.PI_CODING_AGENT_DIR;
const originalXdgStateHome = process.env.XDG_STATE_HOME;
const fallbackAgentDir = path.join(getConfigRootDir(), "agent");
let cleanupRoot: string | undefined;
afterEach(async () => {
if (originalXdgStateHome === undefined) {
delete process.env.XDG_STATE_HOME;
} else {
process.env.XDG_STATE_HOME = originalXdgStateHome;
}
if (originalAgentDir) {
setAgentDir(originalAgentDir);
} else {
setAgentDir(fallbackAgentDir);
delete process.env.PI_CODING_AGENT_DIR;
}
if (cleanupRoot) {
await fs.rm(cleanupRoot, { recursive: true, force: true });
cleanupRoot = undefined;
}
});
describe("raw SSE report bundle", () => {
it("includes captured raw SSE text and dropped-record disclosure", async () => {
cleanupRoot = await fs.mkdtemp(path.join(os.tmpdir(), "omp-raw-sse-report-"));
const xdgStateHome = path.join(cleanupRoot, "state");
await fs.mkdir(path.join(xdgStateHome, "omp"), { recursive: true });
process.env.XDG_STATE_HOME = xdgStateHome;
setAgentDir(fallbackAgentDir);
const buffer = new RawSseDebugBuffer();
buffer.recordResponse(
{ status: 200, requestId: "req_report", headers: {}, metadata: { lastTransport: "sse" } },
model,
);
for (let i = 0; i < 1_001; i++) {
buffer.recordEvent(
{ event: "message_delta", data: `{"i":${i}}`, raw: ["event: message_delta", `data: {"i":${i}}`] },
model,
);
}
const rawSseText = buffer.toRawText();
expect(rawSseText).toContain(": omp-debug-dropped records=");
expect(rawSseText).toContain("event: message_delta");
const result = await createReportBundle({ sessionFile: undefined, rawSseText });
expect(result.files).toContain("raw-sse.txt");
const archive = new Bun.Archive(await Bun.file(result.path).bytes());
const files = await archive.files();
expect(await files.get("raw-sse.txt")?.text()).toBe(rawSseText);
});
});
+2 -2
View File
@@ -1,4 +1,4 @@
import { afterEach, describe, expect, it, spyOn, vi } from "bun:test";
import { afterEach, describe, expect, it, setSystemTime, spyOn, vi } from "bun:test";
import { TUI } from "@oh-my-pi/pi-tui";
import { Loader, type LoaderMessageColorFn } from "@oh-my-pi/pi-tui/components/loader";
import { visibleWidth } from "@oh-my-pi/pi-tui/utils";
@@ -91,7 +91,7 @@ describe("Loader component", () => {
it("requests render when animated message bytes change between spinner frames", () => {
vi.useFakeTimers();
vi.setSystemTime(1_000);
setSystemTime(new Date(1_000));
const ui = { requestRender: vi.fn() } as unknown as TUI;
const colorMessage = ((text: string) => `${text}-${Date.now()}`) as LoaderMessageColorFn & { animated: true };
colorMessage.animated = true;