feat(ai/providers): enhanced SSE parsing and propagated timeout headers
- Augmented streaming SSE observers to infer event names from payload fields when the event field is missing. - Added `X-Stainless-Timeout` request headers derived from `requestTimeoutMs` for OpenAI and Azure response/completions streaming calls.
This commit is contained in:
@@ -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<ResponseStreamEvent>;
|
||||
try {
|
||||
const headersWithTimeout = { ...headers };
|
||||
if (requestTimeoutMs !== undefined) {
|
||||
headersWithTimeout["X-Stainless-Timeout"] = Math.floor(requestTimeoutMs / 1000).toString();
|
||||
}
|
||||
const handle = await postOpenAIStream<ResponseStreamEvent>({
|
||||
url,
|
||||
headers,
|
||||
headers: headersWithTimeout,
|
||||
body: params,
|
||||
signal: requestSignal,
|
||||
fetch: options?.fetch,
|
||||
|
||||
@@ -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<ChatCompletionChunk>({
|
||||
url: completionsUrl,
|
||||
headers,
|
||||
headers: headersWithTimeout,
|
||||
body: params,
|
||||
signal: requestSignal,
|
||||
fetch: options?.fetch,
|
||||
|
||||
@@ -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<ResponseStreamEvent>({
|
||||
url: requestUrl,
|
||||
headers: requestHeaders,
|
||||
headers,
|
||||
body: requestParams,
|
||||
signal: requestSignal,
|
||||
fetch: options?.fetch,
|
||||
|
||||
Reference in New Issue
Block a user