// server/supabase/functions/llm-proxy/provider-deadline.ts // Deadlines for the Anthropic call, split by phase. // // A single AbortSignal.timeout(45 s) on fetch also errors the response body, // so it used to cut off streamed answers that took longer than 45 s in total. // The phases now have their own limits: // * time to first byte (response headers): ttfbMs — covers the provider call. // For non-stream requests the whole body is the answer, so the deadline // stays armed until the JSON has been read. // * streamed body: an idle limit between chunks and a total limit for the // stream, sized for a max_tokens=4096 answer. // The client's request signal is tied in, so a client that disconnects cancels // the upstream generation instead of letting it run to completion. export interface LlmProviderTimeouts { /** Response headers must arrive within this time. */ ttfbMs: number /** Longest gap allowed between two streamed chunks. */ streamIdleMs: number /** Longest a streamed body may run after the headers arrived. */ streamTotalMs: number } export const DEFAULT_LLM_PROVIDER_TIMEOUTS: Readonly = Object.freeze({ ttfbMs: 45_000, streamIdleMs: 30_000, streamTotalMs: 180_000, }) export function providerTimeoutError(message: string): DOMException { return new DOMException(message, 'TimeoutError') } export interface ProviderCallDeadline { /** Pass to fetch; aborts on the first-byte deadline, a stream deadline or a client disconnect. */ readonly signal: AbortSignal /** Stop the first-byte timer (the response headers arrived). */ headersReceived(): void abort(reason: unknown): void /** Clear timers and listeners. Safe to call more than once. */ dispose(): void } export function createProviderCallDeadline( ttfbMs: number, clientSignal: AbortSignal | null, ): ProviderCallDeadline { const controller = new AbortController() let ttfbTimer: ReturnType | null = setTimeout(() => { ttfbTimer = null controller.abort(providerTimeoutError('provider response timeout')) }, ttfbMs) const onClientAbort = (): void => { controller.abort(clientSignal?.reason ?? new DOMException('client disconnected', 'AbortError')) } if (clientSignal?.aborted) onClientAbort() else clientSignal?.addEventListener('abort', onClientAbort, { once: true }) const clearTtfb = (): void => { if (ttfbTimer !== null) clearTimeout(ttfbTimer) ttfbTimer = null } return { signal: controller.signal, headersReceived: clearTtfb, abort(reason: unknown): void { controller.abort(reason) }, dispose(): void { clearTtfb() clientSignal?.removeEventListener('abort', onClientAbort) }, } } /** * Tracks whether an Anthropic SSE stream reached `message_stop`. * Chunks can split the marker, so a short tail of the previous text is kept. */ export function createSseCompletionTracker(): { push(chunk: Uint8Array): void readonly completed: boolean } { const MARKER = 'message_stop' const decoder = new TextDecoder() let tail = '' let completed = false return { push(chunk: Uint8Array): void { if (completed) return const text = tail + decoder.decode(chunk, { stream: true }) if (text.includes(MARKER)) { completed = true tail = '' return } tail = text.slice(-(MARKER.length - 1)) }, get completed(): boolean { return completed }, } } /** How a relayed stream ended. */ export type StreamRelayOutcome = /** Upstream ended after message_stop: the client received the full answer. */ | 'completed' /** Upstream errored, timed out, or ended without message_stop. */ | 'failed' /** The client stopped reading (disconnect/cancel). */ | 'cancelled' export interface StreamRelayOptions { idleMs: number totalMs: number /** Abort the upstream request (fetch signal). */ abortUpstream(reason: unknown): void /** Called exactly once when the relay ends. Awaited before the client stream closes or errors. */ onSettled(outcome: StreamRelayOutcome): Promise } /** * Relay a provider SSE body to the client with an idle deadline between * chunks and a total deadline for the whole stream. */ export function relayProviderStream( upstream: ReadableStream, options: StreamRelayOptions, ): ReadableStream { const reader = upstream.getReader() const tracker = createSseCompletionTracker() let settled = false let totalTimer: ReturnType | null = null let idleTimer: ReturnType | null = null let rejectIdle: ((reason: unknown) => void) | null = null let downstream: ReadableStreamDefaultController | null = null const clearTimers = (): void => { if (totalTimer !== null) clearTimeout(totalTimer) if (idleTimer !== null) clearTimeout(idleTimer) totalTimer = null idleTimer = null } const settle = async (outcome: StreamRelayOutcome): Promise => { if (settled) return false settled = true clearTimers() try { await options.onSettled(outcome) } catch { // Settlement is best effort; the caller logs its own failures. } return true } const terminate = async (reason: unknown): Promise => { if (settled) return options.abortUpstream(reason) reader.cancel(reason).catch(() => undefined) rejectIdle?.(reason) if (await settle('failed')) { try { downstream?.error(reason) } catch { // Already closed. } } } const readWithIdleDeadline = (): Promise> => new Promise((resolve, reject) => { rejectIdle = reject idleTimer = setTimeout(() => { idleTimer = null void terminate(providerTimeoutError('provider stream idle timeout')) }, options.idleMs) reader.read().then(resolve, reject).finally(() => { if (idleTimer !== null) clearTimeout(idleTimer) idleTimer = null rejectIdle = null }) }) return new ReadableStream({ start(controller) { downstream = controller // Also fires while the client applies backpressure and no read is pending. totalTimer = setTimeout(() => { totalTimer = null void terminate(providerTimeoutError('provider stream total timeout')) }, options.totalMs) }, async pull(controller) { if (settled) return let result: ReadableStreamReadResult try { result = await readWithIdleDeadline() } catch (err) { if (settled) return options.abortUpstream(err) if (await settle('failed')) controller.error(err) return } if (settled) return if (result.done) { const outcome: StreamRelayOutcome = tracker.completed ? 'completed' : 'failed' // A stream that ended without message_stop is still closed cleanly: the // client already treats a missing message_stop as incomplete, and any // in-band `event: error` the provider sent stays readable. if (await settle(outcome)) controller.close() return } tracker.push(result.value) controller.enqueue(result.value) }, async cancel(reason) { options.abortUpstream(reason) reader.cancel(reason).catch(() => undefined) await settle('cancelled') }, }) }