fix(llm-proxy): reserve quota before the provider call and stop cutting off long streams
This commit is contained in:
parent
2f94d24c99
commit
b306034bfc
7 changed files with 1489 additions and 307 deletions
226
server/supabase/functions/llm-proxy/provider-deadline.ts
Normal file
226
server/supabase/functions/llm-proxy/provider-deadline.ts
Normal file
|
|
@ -0,0 +1,226 @@
|
|||
// 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')
|
||||
},
|
||||
})
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue