Spaces:
Runtime error
Runtime error
File size: 4,746 Bytes
077865a | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 | /** Parse an HTTP `Retry-After` header (delta-seconds or an HTTP-date) into a
* millisecond delay. Returns undefined when absent or unparseable. */
export function parseRetryAfterMs(value) {
if (!value)
return undefined;
const trimmed = value.trim();
if (/^\d+$/.test(trimmed))
return Number(trimmed) * 1000;
const when = Date.parse(trimmed);
if (!Number.isNaN(when))
return Math.max(0, when - Date.now());
return undefined;
}
/** Build an error for a non-OK upstream response, capturing the status and any
* Retry-After hint. Used by every provider adapter so the proxy can honor a
* provider's explicit back-off when it sets the cooldown. */
export function providerHttpError(res, message) {
const err = new Error(message);
err.status = res.status;
const retryAfterMs = parseRetryAfterMs(res.headers?.get('retry-after'));
if (retryAfterMs !== undefined)
err.retryAfterMs = retryAfterMs;
return err;
}
export class BaseProvider {
/** Providers whose free tier needs no API key (e.g. Kilo's anonymous gateway).
* When true, the gateway stores a sentinel key row so routing still considers
* the platform "configured", and the provider omits the Authorization header
* on outgoing requests. Defaults to false; set by subclasses. */
keyless = false;
async fetchWithTimeout(url, init, timeoutMs = 15000) {
const controller = new AbortController();
const timeout = setTimeout(() => controller.abort(), timeoutMs);
try {
return await fetch(url, { ...init, signal: controller.signal });
}
finally {
clearTimeout(timeout);
}
}
makeId() {
return `chatcmpl-${Date.now()}-${Math.random().toString(36).slice(2, 8)}`;
}
/**
* Shared SSE reader for OpenAI-wire streaming endpoints (#231 audit).
*
* Hardened against the upstream failure modes observed live:
* - Inactivity timeout: fetchWithTimeout's abort timer dies the moment
* response HEADERS arrive, so a provider that stalls mid-body used to
* hang the client forever. Each read now has its own deadline.
* - Abrupt EOF: a stream that ends without `[DONE]` AND without any
* `finish_reason` is a truncated generation, not a completion. It used
* to end the generator silently (truncation logged as success); it now
* throws a retryable error so the proxy can fail over or report it.
* Providers that skip `[DONE]` but do send a terminal finish_reason
* (several compat shims) still complete normally.
*
* Malformed data lines are skipped, matching previous behavior.
*/
async *readSseStream(res, inactivityTimeoutMs = 90000) {
const reader = res.body?.getReader();
if (!reader)
throw new Error('No response body');
const decoder = new TextDecoder();
let buffer = '';
let sawFinishReason = false;
try {
while (true) {
let timer;
const result = await Promise.race([
reader.read(),
new Promise((_, reject) => {
timer = setTimeout(() => reject(new Error(`${this.name} stream stalled: no data for ${inactivityTimeoutMs}ms (timeout)`)), inactivityTimeoutMs);
}),
]).finally(() => clearTimeout(timer));
const { done, value } = result;
if (done)
break;
buffer += decoder.decode(value, { stream: true });
const lines = buffer.split('\n');
buffer = lines.pop() ?? '';
for (const line of lines) {
const trimmed = line.trim();
if (!trimmed || !trimmed.startsWith('data: '))
continue;
const data = trimmed.slice(6);
if (data === '[DONE]')
return;
try {
const chunk = JSON.parse(data);
if (chunk.choices?.some(c => c.finish_reason != null))
sawFinishReason = true;
yield chunk;
}
catch {
// Skip malformed chunks
}
}
}
}
finally {
reader.cancel().catch(() => { });
}
if (!sawFinishReason) {
throw new Error(`${this.name} stream ended unexpectedly (no [DONE], no finish_reason) — connection reset or truncated upstream`);
}
}
}
//# sourceMappingURL=base.js.map |