diff --git a/packages/ai/src/providers/azure-openai-responses.ts b/packages/ai/src/providers/azure-openai-responses.ts index f48c805d4..60dca0355 100644 --- a/packages/ai/src/providers/azure-openai-responses.ts +++ b/packages/ai/src/providers/azure-openai-responses.ts @@ -115,7 +115,26 @@ export const streamAzureOpenAIResponses: StreamFunction<"azure-openai-responses" 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; + const rawSseObserver = onSseEvent + ? (event: RawSseEvent) => { + if (!event.event && event.data && event.data !== "[DONE]") { + try { + const parsed = JSON.parse(event.data); + const resolvedEvent = + typeof parsed.type === "string" + ? parsed.type + : typeof parsed.object === "string" + ? parsed.object + : null; + if (resolvedEvent) { + event.event = resolvedEvent; + event.raw = [`event: ${resolvedEvent}`, ...event.raw]; + } + } catch {} + } + onSseEvent(event, model); + } + : undefined; try { const apiKey = options?.apiKey || getEnvApiKey(model.provider) || ""; @@ -141,9 +160,13 @@ export const streamAzureOpenAIResponses: StreamFunction<"azure-openai-responses" } let openaiStream: AsyncIterable; try { + const headersWithTimeout = { ...headers }; + if (requestTimeoutMs !== undefined) { + headersWithTimeout["X-Stainless-Timeout"] = Math.floor(requestTimeoutMs / 1000).toString(); + } const handle = await postOpenAIStream({ url, - headers, + headers: headersWithTimeout, body: params, signal: requestSignal, fetch: options?.fetch, diff --git a/packages/ai/src/providers/openai-completions.ts b/packages/ai/src/providers/openai-completions.ts index 71885c138..87cf8787d 100644 --- a/packages/ai/src/providers/openai-completions.ts +++ b/packages/ai/src/providers/openai-completions.ts @@ -424,7 +424,26 @@ export const streamOpenAICompletions: StreamFunction<"openai-completions"> = ( 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; + const rawSseObserver = onSseEvent + ? (event: RawSseEvent) => { + if (!event.event && event.data && event.data !== "[DONE]") { + try { + const parsed = JSON.parse(event.data); + const resolvedEvent = + typeof parsed.type === "string" + ? parsed.type + : typeof parsed.object === "string" + ? parsed.object + : null; + if (resolvedEvent) { + event.event = resolvedEvent; + event.raw = [`event: ${resolvedEvent}`, ...event.raw]; + } + } catch {} + } + onSseEvent(event, model); + } + : undefined; // Assigned once the block helpers exist (they are scoped to the `try`); // the catch handler uses it to close any open blocks before emitting the // terminal error so both exit paths obey the same block lifecycle. @@ -485,9 +504,13 @@ export const streamOpenAICompletions: StreamFunction<"openai-completions"> = ( ); } try { + const headersWithTimeout = { ...headers }; + if (requestTimeoutMs !== undefined) { + headersWithTimeout["X-Stainless-Timeout"] = Math.floor(requestTimeoutMs / 1000).toString(); + } const { events, response, requestId } = await postOpenAIStream({ url: completionsUrl, - headers, + headers: headersWithTimeout, body: params, signal: requestSignal, fetch: options?.fetch, diff --git a/packages/ai/src/providers/openai-responses.ts b/packages/ai/src/providers/openai-responses.ts index 7b932f2e7..881598386 100644 --- a/packages/ai/src/providers/openai-responses.ts +++ b/packages/ai/src/providers/openai-responses.ts @@ -364,7 +364,26 @@ export const streamOpenAIResponses: StreamFunction<"openai-responses"> = ( 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; + const rawSseObserver = onSseEvent + ? (event: RawSseEvent) => { + if (!event.event && event.data && event.data !== "[DONE]") { + try { + const parsed = JSON.parse(event.data); + const resolvedEvent = + typeof parsed.type === "string" + ? parsed.type + : typeof parsed.object === "string" + ? parsed.object + : null; + if (resolvedEvent) { + event.event = resolvedEvent; + event.raw = [`event: ${resolvedEvent}`, ...event.raw]; + } + } catch {} + } + onSseEvent(event, model); + } + : undefined; try { // Keep request routing on `sessionId` while allowing callers to pin a @@ -418,9 +437,13 @@ export const streamOpenAIResponses: StreamFunction<"openai-responses"> = ( ); } try { + const headers = { ...requestHeaders }; + if (requestTimeoutMs !== undefined) { + headers["X-Stainless-Timeout"] = Math.floor(requestTimeoutMs / 1000).toString(); + } const { events, response, requestId } = await postOpenAIStream({ url: requestUrl, - headers: requestHeaders, + headers, body: requestParams, signal: requestSignal, fetch: options?.fetch,