From 00e6e14cb49fb9c93941df018d9301b35cc8e52c Mon Sep 17 00:00:00 2001 From: can1357 Date: Mon, 8 Jun 2026 02:35:08 +0200 Subject: [PATCH] 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. --- packages/ai/CHANGELOG.md | 5 +- packages/ai/src/providers/anthropic.ts | 59 ++-- .../src/providers/azure-openai-responses.ts | 56 ++-- .../ai/src/providers/openai-completions.ts | 41 ++- packages/ai/src/providers/openai-responses.ts | 64 ++-- packages/ai/src/utils/sse-debug.ts | 271 ----------------- packages/ai/test/raw-sse-sdk-capture.test.ts | 283 ++++++++++++++++++ packages/ai/test/sse-debug.test.ts | 218 ++------------ packages/coding-agent/CHANGELOG.md | 3 +- packages/coding-agent/src/debug/index.ts | 8 + .../coding-agent/src/debug/raw-sse-buffer.ts | 11 +- .../coding-agent/src/debug/report-bundle.ts | 9 + .../test/debug/raw-sse-buffer.test.ts | 23 ++ .../test/debug/raw-sse-report-bundle.test.ts | 76 +++++ packages/tui/test/loader.test.ts | 4 +- 15 files changed, 575 insertions(+), 556 deletions(-) create mode 100644 packages/ai/test/raw-sse-sdk-capture.test.ts create mode 100644 packages/coding-agent/test/debug/raw-sse-report-bundle.test.ts diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index d4efedab2..c510ce7ee 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -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. \ No newline at end of file diff --git a/packages/ai/src/providers/anthropic.ts b/packages/ai/src/providers/anthropic.ts index c8e9f3c35..172033d4e 100644 --- a/packages/ai/src/providers/anthropic.ts +++ b/packages/ai/src/providers/anthropic.ts @@ -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; response: Response; requestId: string | null }> { +): Promise<{ + events: AsyncIterable; + 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, + observer: (event: RawSseEvent) => void, +): AsyncGenerator { + 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["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; 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 } : {}), }; } diff --git a/packages/ai/src/providers/azure-openai-responses.ts b/packages/ai/src/providers/azure-openai-responses.ts index 26b3f0a16..711f02b33 100644 --- a/packages/ai/src/providers/azure-openai-responses.ts +++ b/packages/ai/src/providers/azure-openai-responses.ts @@ -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, + observer: (event: RawSseEvent) => void, +): AsyncGenerator { + 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, }); } diff --git a/packages/ai/src/providers/openai-completions.ts b/packages/ai/src/providers/openai-completions.ts index ca8a78d33..42f3d19bd 100644 --- a/packages/ai/src/providers/openai-completions.ts +++ b/packages/ai/src/providers/openai-completions.ts @@ -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, + observer: (event: RawSseEvent) => void, +): AsyncGenerator { + 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, 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, diff --git a/packages/ai/src/providers/openai-responses.ts b/packages/ai/src/providers/openai-responses.ts index ed6de0481..f3f251099 100644 --- a/packages/ai/src/providers/openai-responses.ts +++ b/packages/ai/src/providers/openai-responses.ts @@ -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, + observer: (event: RawSseEvent) => void, +): AsyncGenerator { + 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> = []; - 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(item) as unknown as Record); - }, + 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(item) as unknown as Record); + }, + }); if (premiumRequestsTotal !== undefined) output.usage.premiumRequests = premiumRequestsTotal; const firstEventTimeoutError = abortTracker.getLocalAbortReason(); @@ -341,7 +354,6 @@ function createClient( extraHeaders?: Record, 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, diff --git a/packages/ai/src/utils/sse-debug.ts b/packages/ai/src/utils/sse-debug.ts index b42028a9f..63a83826f 100644 --- a/packages/ai/src/utils/sse-debug.ts +++ b/packages/ai/src/utils/sse-debug.ts @@ -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; -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 => { - 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({ - 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; -} diff --git a/packages/ai/test/raw-sse-sdk-capture.test.ts b/packages/ai/test/raw-sse-sdk-capture.test.ts new file mode 100644 index 000000000..02ab86e03 --- /dev/null +++ b/packages/ai/test/raw-sse-sdk-capture.test.ts @@ -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 { + 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])}`]); + }); +}); diff --git a/packages/ai/test/sse-debug.test.ts b/packages/ai/test/sse-debug.test.ts index 7260c67a5..bc839265b 100644 --- a/packages/ai/test/sse-debug.test.ts +++ b/packages/ai/test/sse-debug.test.ts @@ -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 { - let i = 0; - return new ReadableStream({ - 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 { - const reader = response.body!.getReader(); - for (;;) { - const { done } = await reader.read(); - if (done) return; - } -} - -async function collect(chunks: Uint8Array[]): Promise { - 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(); }); }); diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 2b65b40c5..4c2cc7671 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -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. diff --git a/packages/coding-agent/src/debug/index.ts b/packages/coding-agent/src/debug/index.ts index 17840cc4b..5def149a0 100644 --- a/packages/coding-agent/src/debug/index.ts +++ b/packages/coding-agent/src/debug/index.ts @@ -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 { // Extract key settings for the report return { diff --git a/packages/coding-agent/src/debug/raw-sse-buffer.ts b/packages/coding-agent/src/debug/raw-sse-buffer.ts index 9120637b6..1bb6d0b3e 100644 --- a/packages/coding-agent/src/debug/raw-sse-buffer.ts +++ b/packages/coding-agent/src/debug/raw-sse-buffer.ts @@ -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 { diff --git a/packages/coding-agent/src/debug/report-bundle.ts b/packages/coding-agent/src/debug/report-bundle.ts index 635babe57..0e7914c47 100644 --- a/packages/coding-agent/src/debug/report-bundle.ts +++ b/packages/coding-agent/src/debug/report-bundle.ts @@ -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 { diff --git a/packages/coding-agent/test/debug/raw-sse-buffer.test.ts b/packages/coding-agent/test/debug/raw-sse-buffer.test.ts index cd2e35b8b..56008eab7 100644 --- a/packages/coding-agent/test/debug/raw-sse-buffer.test.ts +++ b/packages/coding-agent/test/debug/raw-sse-buffer.test.ts @@ -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"); + }); }); diff --git a/packages/coding-agent/test/debug/raw-sse-report-bundle.test.ts b/packages/coding-agent/test/debug/raw-sse-report-bundle.test.ts new file mode 100644 index 000000000..bb43fdb58 --- /dev/null +++ b/packages/coding-agent/test/debug/raw-sse-report-bundle.test.ts @@ -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); + }); +}); diff --git a/packages/tui/test/loader.test.ts b/packages/tui/test/loader.test.ts index 9aa28131d..c2f26101f 100644 --- a/packages/tui/test/loader.test.ts +++ b/packages/tui/test/loader.test.ts @@ -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;