feat: track thread ids and retry stream disconnects
This commit is contained in:
+163
-3
@@ -45,6 +45,7 @@ const DEFAULT_CONFIG = {
|
||||
retryable_status_codes: [429, 503],
|
||||
retryable_error_messages: [
|
||||
"Selected model is at capacity. Please try a different model.",
|
||||
"stream disconnected before completion: Concurrency limit exceeded for account, please retry later",
|
||||
],
|
||||
upstream_fetch_retry_attempts: 5,
|
||||
upstream_fetch_retry_backoff_ms: 350,
|
||||
@@ -60,6 +61,37 @@ const REASONING_POINTERS = [
|
||||
"/response/usage/output_tokens_details/reasoning_tokens",
|
||||
"/response/usage/completion_tokens_details/reasoning_tokens",
|
||||
];
|
||||
const REQUEST_THREAD_ID_POINTERS = [
|
||||
"/thread_id",
|
||||
"/thread",
|
||||
"/thread/id",
|
||||
"/conversation_id",
|
||||
"/conversation",
|
||||
"/conversation/id",
|
||||
];
|
||||
const RESPONSE_THREAD_ID_POINTERS = [
|
||||
"/thread_id",
|
||||
"/thread",
|
||||
"/thread/id",
|
||||
"/conversation_id",
|
||||
"/conversation",
|
||||
"/conversation/id",
|
||||
"/response/thread_id",
|
||||
"/response/thread",
|
||||
"/response/thread/id",
|
||||
"/response/conversation_id",
|
||||
"/response/conversation",
|
||||
"/response/conversation/id",
|
||||
];
|
||||
const NON_STREAM_RESPONSE_ID_POINTERS = [
|
||||
"/id",
|
||||
"/response_id",
|
||||
"/response/id",
|
||||
];
|
||||
const STREAM_RESPONSE_ID_POINTERS = [
|
||||
"/response_id",
|
||||
"/response/id",
|
||||
];
|
||||
|
||||
function parseArgs(argv) {
|
||||
const args = { config: null, log: null };
|
||||
@@ -150,6 +182,49 @@ function firstInteger(...values) {
|
||||
return null;
|
||||
}
|
||||
|
||||
function firstNonEmptyString(...values) {
|
||||
for (const value of values) {
|
||||
if (typeof value !== "string") {
|
||||
continue;
|
||||
}
|
||||
const trimmed = value.trim();
|
||||
if (trimmed) {
|
||||
return trimmed;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function extractStringByPointers(payload, pointers) {
|
||||
for (const pointer of pointers) {
|
||||
const raw = jsonPointerGet(payload, pointer);
|
||||
if (typeof raw !== "string") {
|
||||
continue;
|
||||
}
|
||||
const trimmed = raw.trim();
|
||||
if (trimmed) {
|
||||
return trimmed;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function extractRequestThreadId(payload) {
|
||||
return extractStringByPointers(payload, REQUEST_THREAD_ID_POINTERS);
|
||||
}
|
||||
|
||||
function extractResponseThreadId(payload) {
|
||||
return extractStringByPointers(payload, RESPONSE_THREAD_ID_POINTERS);
|
||||
}
|
||||
|
||||
function extractNonStreamingResponseId(payload) {
|
||||
return extractStringByPointers(payload, NON_STREAM_RESPONSE_ID_POINTERS);
|
||||
}
|
||||
|
||||
function extractStreamingResponseId(payload) {
|
||||
return extractStringByPointers(payload, STREAM_RESPONSE_ID_POINTERS);
|
||||
}
|
||||
|
||||
function normalizeUsageSnapshot(payload) {
|
||||
const usage = payload?.usage || payload?.response?.usage || null;
|
||||
if (!usage || typeof usage !== "object") {
|
||||
@@ -458,6 +533,8 @@ function openRequestsDatabase(dbPath) {
|
||||
CREATE TABLE IF NOT EXISTS requests (
|
||||
seq INTEGER PRIMARY KEY,
|
||||
request_id TEXT,
|
||||
response_id TEXT,
|
||||
thread_id TEXT,
|
||||
started_at TEXT,
|
||||
finished_at TEXT,
|
||||
duration_ms INTEGER,
|
||||
@@ -496,7 +573,15 @@ function openRequestsDatabase(dbPath) {
|
||||
if (!requestColumnNames.has("request_id")) {
|
||||
db.exec("ALTER TABLE requests ADD COLUMN request_id TEXT");
|
||||
}
|
||||
if (!requestColumnNames.has("response_id")) {
|
||||
db.exec("ALTER TABLE requests ADD COLUMN response_id TEXT");
|
||||
}
|
||||
if (!requestColumnNames.has("thread_id")) {
|
||||
db.exec("ALTER TABLE requests ADD COLUMN thread_id TEXT");
|
||||
}
|
||||
db.exec("CREATE INDEX IF NOT EXISTS idx_requests_request_id ON requests(request_id)");
|
||||
db.exec("CREATE INDEX IF NOT EXISTS idx_requests_response_id ON requests(response_id)");
|
||||
db.exec("CREATE INDEX IF NOT EXISTS idx_requests_thread_id ON requests(thread_id)");
|
||||
return db;
|
||||
}
|
||||
|
||||
@@ -528,6 +613,8 @@ function requestRowFromEntry(entry) {
|
||||
return {
|
||||
seq: entry.seq,
|
||||
request_id: entry.request_id || null,
|
||||
response_id: entry.response_id || null,
|
||||
thread_id: entry.thread_id || null,
|
||||
started_at: entry.started_at || null,
|
||||
finished_at: entry.finished_at || null,
|
||||
duration_ms: Number.isInteger(entry.duration_ms) ? entry.duration_ms : null,
|
||||
@@ -560,13 +647,13 @@ function requestRowFromEntry(entry) {
|
||||
function insertRequestRow(db, row) {
|
||||
db.prepare(`
|
||||
INSERT OR REPLACE INTO requests (
|
||||
seq, request_id, started_at, finished_at, duration_ms, profile_name, method, path, model,
|
||||
seq, request_id, response_id, thread_id, started_at, finished_at, duration_ms, profile_name, method, path, model,
|
||||
requested_model, forwarded_model, request_stream, response_stream, inspected, matched,
|
||||
status_code, upstream_status_code, reasoning_tokens, input_tokens, output_tokens,
|
||||
total_tokens, cached_tokens, error, upstream_origin, upstream_path,
|
||||
upstream_auth_mode, upstream_auth_source, payload_json
|
||||
) VALUES (
|
||||
@seq, @request_id, @started_at, @finished_at, @duration_ms, @profile_name, @method, @path, @model,
|
||||
@seq, @request_id, @response_id, @thread_id, @started_at, @finished_at, @duration_ms, @profile_name, @method, @path, @model,
|
||||
@requested_model, @forwarded_model, @request_stream, @response_stream, @inspected, @matched,
|
||||
@status_code, @upstream_status_code, @reasoning_tokens, @input_tokens, @output_tokens,
|
||||
@total_tokens, @cached_tokens, @error, @upstream_origin, @upstream_path,
|
||||
@@ -627,6 +714,8 @@ function buildRequestQueryFilters({ query, filter }) {
|
||||
lower(coalesce(method, '')) LIKE @query OR
|
||||
lower(coalesce(path, '')) LIKE @query OR
|
||||
lower(coalesce(request_id, '')) LIKE @query OR
|
||||
lower(coalesce(response_id, '')) LIKE @query OR
|
||||
lower(coalesce(thread_id, '')) LIKE @query OR
|
||||
lower(coalesce(model, '')) LIKE @query OR
|
||||
lower(coalesce(requested_model, '')) LIKE @query OR
|
||||
lower(coalesce(forwarded_model, '')) LIKE @query OR
|
||||
@@ -1222,6 +1311,8 @@ function buildRequestEntry({ seq, startedAt, startedMs, req, pathname, requestJs
|
||||
return {
|
||||
seq,
|
||||
request_id: null,
|
||||
response_id: null,
|
||||
thread_id: extractRequestThreadId(requestJson),
|
||||
lifecycle_state: "sent",
|
||||
started_at: startedAt.toISOString(),
|
||||
first_response_at: null,
|
||||
@@ -2562,14 +2653,26 @@ function findRetryableStreamErrorMatch(config, parsedBody, bodyText, eventName =
|
||||
return matchRetryableMessage(config, parsedBody, bodyText);
|
||||
}
|
||||
|
||||
function findRetryableStreamTerminationMatch(config, error) {
|
||||
const message = `${error?.message || error || ""}`.trim();
|
||||
if (!message) {
|
||||
return null;
|
||||
}
|
||||
return matchRetryableMessage(config, null, message);
|
||||
}
|
||||
|
||||
function isExpectedStreamTermination(error) {
|
||||
if (!error) {
|
||||
return false;
|
||||
}
|
||||
const message = `${error?.message || ""}`.trim().toLowerCase();
|
||||
if (error.name === "AbortError") {
|
||||
return true;
|
||||
}
|
||||
return error instanceof TypeError && error.message === "terminated";
|
||||
return error instanceof TypeError && (
|
||||
message === "terminated" ||
|
||||
message.includes("stream disconnected before completion")
|
||||
);
|
||||
}
|
||||
|
||||
function isRetryableUpstreamFetchError(error) {
|
||||
@@ -2701,6 +2804,8 @@ function inspectSseChunk(state, chunk, config) {
|
||||
const result = {
|
||||
reasoning: null,
|
||||
usage: null,
|
||||
response_id: null,
|
||||
thread_id: null,
|
||||
retryable_upstream_error: null,
|
||||
};
|
||||
|
||||
@@ -2735,6 +2840,8 @@ function inspectSseChunk(state, chunk, config) {
|
||||
result.reasoning = reasoning;
|
||||
}
|
||||
result.usage = mergeUsageSnapshots(result.usage, normalizeUsageSnapshot(parsed));
|
||||
result.response_id = result.response_id || extractStreamingResponseId(parsed);
|
||||
result.thread_id = result.thread_id || extractResponseThreadId(parsed);
|
||||
} catch {
|
||||
// ignore malformed SSE payloads
|
||||
}
|
||||
@@ -2770,6 +2877,11 @@ async function handleNonStreaming({
|
||||
: null;
|
||||
const reasoning = parsed ? extractReasoningTokens(parsed) : null;
|
||||
const usage = parsed ? normalizeUsageSnapshot(parsed) : null;
|
||||
const responseId = parsed ? extractNonStreamingResponseId(parsed) : null;
|
||||
const threadId = firstNonEmptyString(
|
||||
requestEntry.thread_id,
|
||||
parsed ? extractResponseThreadId(parsed) : null,
|
||||
);
|
||||
const matched = reasoningMatched(config, reasoning);
|
||||
const retryableUpstreamError = terminalRetryableUpstreamError || findRetryableUpstreamErrorMatch(
|
||||
config,
|
||||
@@ -2777,6 +2889,8 @@ async function handleNonStreaming({
|
||||
parsed,
|
||||
bodyText,
|
||||
);
|
||||
requestEntry.response_id = responseId || requestEntry.response_id || null;
|
||||
requestEntry.thread_id = threadId || requestEntry.thread_id || null;
|
||||
|
||||
recordInspectedResponse(monitor, reasoning, matched || Boolean(retryableUpstreamError));
|
||||
|
||||
@@ -2799,6 +2913,8 @@ async function handleNonStreaming({
|
||||
upstream_status_code: upstreamResponse.status,
|
||||
reasoning_tokens: reasoning,
|
||||
usage,
|
||||
response_id: requestEntry.response_id,
|
||||
thread_id: requestEntry.thread_id,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -2828,6 +2944,8 @@ async function handleNonStreaming({
|
||||
usage,
|
||||
error: `retryable upstream error: ${retryableUpstreamError.matched_pattern}`,
|
||||
match_reason: "retryable_upstream_error",
|
||||
response_id: requestEntry.response_id,
|
||||
thread_id: requestEntry.thread_id,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -2841,6 +2959,8 @@ async function handleNonStreaming({
|
||||
upstream_status_code: upstreamResponse.status,
|
||||
reasoning_tokens: reasoning,
|
||||
usage,
|
||||
response_id: requestEntry.response_id,
|
||||
thread_id: requestEntry.thread_id,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -2878,6 +2998,24 @@ async function handleStreaming({
|
||||
readResult = await reader.read();
|
||||
} catch (error) {
|
||||
if (isExpectedStreamTermination(error)) {
|
||||
const retryableTerminationError = findRetryableStreamTerminationMatch(config, error);
|
||||
if (retryableTerminationError) {
|
||||
return {
|
||||
inspected: true,
|
||||
matched: true,
|
||||
retry_requested: strict502Mode || !wroteAnyChunk,
|
||||
retryable_upstream_error: retryableTerminationError,
|
||||
upstream_status_code: upstreamResponse.status,
|
||||
reasoning_tokens: observedReasoning,
|
||||
usage: observedUsage,
|
||||
error: `retryable upstream error: ${retryableTerminationError.matched_pattern}`,
|
||||
match_reason: "retryable_upstream_error",
|
||||
response_id: requestEntry.response_id,
|
||||
thread_id: requestEntry.thread_id,
|
||||
response_bytes_received: requestEntry.response_bytes_received,
|
||||
stream_chunk_count: requestEntry.stream_chunk_count,
|
||||
};
|
||||
}
|
||||
recordInspectedResponse(monitor, observedReasoning, false);
|
||||
persistStreamingProgress(runtime, requestEntry, { force: true }, new Date());
|
||||
if (strict502Mode) {
|
||||
@@ -2892,6 +3030,8 @@ async function handleStreaming({
|
||||
reasoning_tokens: observedReasoning,
|
||||
usage: observedUsage,
|
||||
error: "upstream stream terminated before completion",
|
||||
response_id: requestEntry.response_id,
|
||||
thread_id: requestEntry.thread_id,
|
||||
response_bytes_received: requestEntry.response_bytes_received,
|
||||
stream_chunk_count: requestEntry.stream_chunk_count,
|
||||
};
|
||||
@@ -2905,6 +3045,8 @@ async function handleStreaming({
|
||||
reasoning_tokens: observedReasoning,
|
||||
usage: observedUsage,
|
||||
error: "upstream stream terminated before completion",
|
||||
response_id: requestEntry.response_id,
|
||||
thread_id: requestEntry.thread_id,
|
||||
response_bytes_received: requestEntry.response_bytes_received,
|
||||
stream_chunk_count: requestEntry.stream_chunk_count,
|
||||
};
|
||||
@@ -2931,6 +3073,8 @@ async function handleStreaming({
|
||||
upstream_status_code: upstreamResponse.status,
|
||||
reasoning_tokens: observedReasoning,
|
||||
usage: observedUsage,
|
||||
response_id: requestEntry.response_id,
|
||||
thread_id: requestEntry.thread_id,
|
||||
response_bytes_received: requestEntry.response_bytes_received,
|
||||
stream_chunk_count: requestEntry.stream_chunk_count,
|
||||
};
|
||||
@@ -2953,6 +3097,8 @@ async function handleStreaming({
|
||||
usage: observedUsage,
|
||||
error: `retryable upstream error: ${retryableUpstreamError.matched_pattern}`,
|
||||
match_reason: "retryable_upstream_error",
|
||||
response_id: requestEntry.response_id,
|
||||
thread_id: requestEntry.thread_id,
|
||||
response_bytes_received: requestEntry.response_bytes_received,
|
||||
stream_chunk_count: requestEntry.stream_chunk_count,
|
||||
};
|
||||
@@ -2965,6 +3111,12 @@ async function handleStreaming({
|
||||
if (Number.isInteger(reasoning)) {
|
||||
observedReasoning = reasoning;
|
||||
}
|
||||
if (inspection.response_id) {
|
||||
requestEntry.response_id = inspection.response_id;
|
||||
}
|
||||
if (inspection.thread_id) {
|
||||
requestEntry.thread_id = inspection.thread_id;
|
||||
}
|
||||
updateStreamingProgress(requestEntry, {
|
||||
chunkBytes: chunkBuffer.length,
|
||||
usage: inspection.usage,
|
||||
@@ -3006,6 +3158,8 @@ async function handleStreaming({
|
||||
upstream_status_code: upstreamResponse.status,
|
||||
reasoning_tokens: reasoning,
|
||||
usage: observedUsage,
|
||||
response_id: requestEntry.response_id,
|
||||
thread_id: requestEntry.thread_id,
|
||||
response_bytes_received: requestEntry.response_bytes_received,
|
||||
stream_chunk_count: requestEntry.stream_chunk_count,
|
||||
};
|
||||
@@ -3080,6 +3234,7 @@ async function proxyRequest(runtime, req, res) {
|
||||
let totalUpstreamAttempts = 0;
|
||||
requestEntry.request_body_bytes = rawRequestBody.length;
|
||||
requestEntry.request_id = computeRequestId(pathname, rawRequestBody);
|
||||
requestEntry.thread_id = extractRequestThreadId(parsedRequestJson);
|
||||
requestEntry.model = requestJson?.model || null;
|
||||
requestEntry.requested_model = parsedRequestJson?.model || null;
|
||||
requestEntry.forwarded_model = forwardedModel || parsedRequestJson?.model || null;
|
||||
@@ -3140,6 +3295,11 @@ async function proxyRequest(runtime, req, res) {
|
||||
copyHeadersToClient(upstreamResponse.headers, res);
|
||||
res.writeHead(upstreamResponse.status);
|
||||
const body = Buffer.from(await upstreamResponse.arrayBuffer());
|
||||
const parsed = isJsonContentType(upstreamResponse.headers.get("content-type"))
|
||||
? parseJsonSafely(body)
|
||||
: null;
|
||||
requestEntry.response_id = requestEntry.response_id || (parsed ? extractNonStreamingResponseId(parsed) : null);
|
||||
requestEntry.thread_id = requestEntry.thread_id || (parsed ? extractResponseThreadId(parsed) : null);
|
||||
res.end(body);
|
||||
recordRequestEntry(
|
||||
runtime,
|
||||
|
||||
Reference in New Issue
Block a user