226 lines
7.3 KiB
TypeScript
226 lines
7.3 KiB
TypeScript
// 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<LlmProviderTimeouts> = 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<typeof setTimeout> | 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<void>
|
|
}
|
|
|
|
/**
|
|
* 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<Uint8Array>,
|
|
options: StreamRelayOptions,
|
|
): ReadableStream<Uint8Array> {
|
|
const reader = upstream.getReader()
|
|
const tracker = createSseCompletionTracker()
|
|
let settled = false
|
|
let totalTimer: ReturnType<typeof setTimeout> | null = null
|
|
let idleTimer: ReturnType<typeof setTimeout> | null = null
|
|
let rejectIdle: ((reason: unknown) => void) | null = null
|
|
let downstream: ReadableStreamDefaultController<Uint8Array> | 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<boolean> => {
|
|
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<void> => {
|
|
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<ReadableStreamReadResult<Uint8Array>> =>
|
|
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<Uint8Array>({
|
|
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<Uint8Array>
|
|
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')
|
|
},
|
|
})
|
|
}
|