From 62d338d9b2f04125e9968207a00bc4ff3f82e6a6 Mon Sep 17 00:00:00 2001 From: steebchen Date: Mon, 2 Feb 2026 17:19:09 +0000 Subject: [PATCH] feat(gateway): add pending log tracking for timeout detection Insert a pending log entry at the start of each /v1/chat/completions request to track requests that may timeout or be cancelled before completing normally. Changes: - Add PENDING status to UnifiedFinishReason enum - Add insertPendingLog() for synchronous DB insert at request start - Add updateLog() for async updates via LOG_UPDATE_QUEUE - Add processLogUpdateQueue() worker loop for handling updates - Convert all insertLog() calls to finalizeLog() helper that updates pending logs or falls back to insert This allows detecting timed out or cancelled requests by querying for logs with unifiedFinishReason = 'pending' that are older than expected. Co-Authored-By: Claude Opus 4.5 --- apps/gateway/src/chat/chat.ts | 65 ++++++++++++--- apps/gateway/src/lib/logs.ts | 121 +++++++++++++++++++++++++++- apps/worker/src/worker.ts | 144 ++++++++++++++++++++++++++++++++++ packages/cache/src/redis.ts | 1 + packages/db/src/schema.ts | 1 + 5 files changed, 317 insertions(+), 15 deletions(-) diff --git a/apps/gateway/src/chat/chat.ts b/apps/gateway/src/chat/chat.ts index f333cf1d4c..db998801cf 100644 --- a/apps/gateway/src/chat/chat.ts +++ b/apps/gateway/src/chat/chat.ts @@ -17,7 +17,12 @@ import { import { isCodingModel } from "@/lib/coding-models.js"; import { calculateCosts, shouldBillCancelledRequests } from "@/lib/costs.js"; import { throwIamException, validateModelAccess } from "@/lib/iam.js"; -import { calculateDataStorageCost, insertLog } from "@/lib/logs.js"; +import { + calculateDataStorageCost, + insertLog, + insertPendingLog, + updateLog, +} from "@/lib/logs.js"; import { createCombinedSignal, isTimeoutError } from "@/lib/timeout-config.js"; import { @@ -1654,6 +1659,42 @@ chat.openapi(completions, async (c) => { }); } + // Insert pending log to track requests that may timeout or be cancelled + // This is done before cache checks so we can detect issues with the request flow + let pendingLogId: string | null = null; + try { + pendingLogId = await insertPendingLog({ + requestId, + organizationId: project.organizationId, + projectId: apiKey.projectId, + apiKeyId: apiKey.id, + requestedModel: initialRequestedModel, + requestedProvider: requestedProvider || null, + mode: project.mode, + usedMode: providerKey?.id ? "api-keys" : "credits", + messages, + streamed: stream || false, + source: source || null, + userAgent: userAgent || null, + }); + } catch (error) { + // Log the error but don't fail the request - we can still process without tracking + logger.warn("Failed to insert pending log, continuing without tracking", { + requestId, + error: error instanceof Error ? error.message : String(error), + }); + } + + // Helper to finalize log - updates pending log if available, otherwise inserts new log + const finalizeLog = async ( + logData: Parameters[0], + ): Promise => { + if (pendingLogId) { + return await updateLog({ ...logData, logId: pendingLogId }); + } + return await insertLog(logData); + }; + // Check if caching is enabled for this project const { enabled: cachingEnabled, duration: cacheDuration } = await isCachingEnabled(project.id); @@ -1794,7 +1835,7 @@ chat.openapi(completions, async (c) => { inputImageCount, ); - await insertLog({ + await finalizeLog({ ...baseLogEntry, duration: 0, // No processing time for cached response timeToFirstToken: null, // Not applicable for cached response @@ -1918,7 +1959,7 @@ chat.openapi(completions, async (c) => { inputImageCount, ); - await insertLog({ + await finalizeLog({ ...baseLogEntry, duration, timeToFirstToken: null, // Not applicable for cached response @@ -2306,7 +2347,7 @@ chat.openapi(completions, async (c) => { undefined, // No plugin results for error case ); - await insertLog({ + await finalizeLog({ ...baseLogEntry, duration: Date.now() - startTime, timeToFirstToken: null, @@ -2425,7 +2466,7 @@ chat.openapi(completions, async (c) => { undefined, // No plugin results for canceled request ); - await insertLog({ + await finalizeLog({ ...baseLogEntry, duration: Date.now() - startTime, timeToFirstToken: null, // Not applicable for canceled request @@ -2533,7 +2574,7 @@ chat.openapi(completions, async (c) => { undefined, // No plugin results for error case ); - await insertLog({ + await finalizeLog({ ...baseLogEntry, duration: Date.now() - startTime, timeToFirstToken: null, // Not applicable for error case @@ -2692,7 +2733,7 @@ chat.openapi(completions, async (c) => { undefined, // No plugin results for error case ); - await insertLog({ + await finalizeLog({ ...baseLogEntry, duration: Date.now() - startTime, timeToFirstToken: null, // Not applicable for error case @@ -4057,7 +4098,7 @@ chat.openapi(completions, async (c) => { const shouldIncludeTokensForBilling = !canceled || (canceled && billCancelledRequests); - await insertLog({ + await finalizeLog({ ...baseLogEntry, duration, timeToFirstToken, @@ -4289,7 +4330,7 @@ chat.openapi(completions, async (c) => { undefined, // No plugin results for error case ); - await insertLog({ + await finalizeLog({ ...baseLogEntry, duration, timeToFirstToken: null, // Not applicable for error case @@ -4421,7 +4462,7 @@ chat.openapi(completions, async (c) => { undefined, // No plugin results for canceled request ); - await insertLog({ + await finalizeLog({ ...baseLogEntry, duration, timeToFirstToken: null, // Not applicable for canceled request @@ -4533,7 +4574,7 @@ chat.openapi(completions, async (c) => { undefined, // No plugin results for error case ); - await insertLog({ + await finalizeLog({ ...baseLogEntry, duration, timeToFirstToken: null, // Not applicable for error case @@ -4859,7 +4900,7 @@ chat.openapi(completions, async (c) => { }); } - await insertLog({ + await finalizeLog({ ...baseLogEntry, duration, timeToFirstToken: null, // Not applicable for non-streaming requests diff --git a/apps/gateway/src/lib/logs.ts b/apps/gateway/src/lib/logs.ts index e2c26be325..eab4006f12 100644 --- a/apps/gateway/src/lib/logs.ts +++ b/apps/gateway/src/lib/logs.ts @@ -1,9 +1,14 @@ -import { publishToQueue, LOG_QUEUE } from "@llmgateway/cache"; -import { UnifiedFinishReason, type LogInsertData } from "@llmgateway/db"; +import { publishToQueue, LOG_QUEUE, LOG_UPDATE_QUEUE } from "@llmgateway/cache"; +import { + UnifiedFinishReason, + type LogInsertData, + cdb as db, + log, + shortid, +} from "@llmgateway/db"; import { logger } from "@llmgateway/logger"; import type { InferInsertModel } from "@llmgateway/db"; -import type { log } from "@llmgateway/db"; /** * Check if a finish reason is expected to map to UNKNOWN @@ -178,3 +183,113 @@ export async function insertLog(logData: LogInsertData): Promise { await publishToQueue(LOG_QUEUE, logData); return 1; // Return 1 to match test expectations } + +/** + * Data required for creating a pending log entry. + * Only includes the minimal fields needed at request start. + */ +export interface PendingLogData { + requestId: string; + organizationId: string; + projectId: string; + apiKeyId: string; + requestedModel: string; + requestedProvider: string | null; + mode: "api-keys" | "credits" | "hybrid"; + usedMode: "api-keys" | "credits"; + messages: unknown; + streamed: boolean; + source: string | null; + userAgent: string | null; +} + +/** + * Insert a pending log entry synchronously to the database. + * This is called at the start of a request to track requests that may timeout or be cancelled. + * Returns the log ID that should be used in the subsequent updateLog() call. + */ +export async function insertPendingLog(data: PendingLogData): Promise { + const logId = shortid(); + + try { + await db.insert(log).values({ + id: logId, + requestId: data.requestId, + organizationId: data.organizationId, + projectId: data.projectId, + apiKeyId: data.apiKeyId, + requestedModel: data.requestedModel, + requestedProvider: data.requestedProvider, + // Set placeholder values for required fields - these will be updated later + usedModel: data.requestedModel, + usedProvider: data.requestedProvider || "unknown", + duration: 0, + responseSize: 0, + mode: data.mode, + usedMode: data.usedMode, + // Set pending status + unifiedFinishReason: UnifiedFinishReason.PENDING, + // Include available data + messages: data.messages, + streamed: data.streamed, + source: data.source, + userAgent: data.userAgent, + // Initialize other fields + hasError: false, + canceled: false, + cached: false, + dataStorageCost: "0", + }); + + return logId; + } catch (error) { + logger.error("Failed to insert pending log", { + requestId: data.requestId, + error: error instanceof Error ? error.message : String(error), + }); + throw error; + } +} + +/** + * Data for updating an existing log entry. + * Includes the log ID and all the fields to update. + */ +export interface LogUpdateData extends Omit { + logId: string; +} + +/** + * Update an existing log entry via the message queue. + * This is called at the end of a request to update the pending log with final data. + */ +export async function updateLog(logData: LogUpdateData): Promise { + if (logData.unifiedFinishReason === undefined) { + if (logData.canceled) { + logData.unifiedFinishReason = UnifiedFinishReason.CANCELED; + } else { + logData.unifiedFinishReason = getUnifiedFinishReason( + logData.finishReason, + logData.usedProvider, + ); + + if ( + logData.unifiedFinishReason === UnifiedFinishReason.UNKNOWN && + logData.finishReason && + !isExpectedUnknownFinishReason( + logData.finishReason, + logData.usedProvider, + ) + ) { + logger.error("Unknown finish reason encountered", { + requestId: logData.requestId, + finishReason: logData.finishReason, + provider: logData.usedProvider, + model: logData.usedModel, + }); + } + } + } + await publishToQueue(LOG_UPDATE_QUEUE, logData); + return 1; +} diff --git a/apps/worker/src/worker.ts b/apps/worker/src/worker.ts index ec071306fd..4f5080e8db 100644 --- a/apps/worker/src/worker.ts +++ b/apps/worker/src/worker.ts @@ -5,6 +5,7 @@ import { z } from "zod"; import { consumeFromQueue, LOG_QUEUE, + LOG_UPDATE_QUEUE, closeRedisClient, publishToQueue, } from "@llmgateway/cache"; @@ -918,6 +919,116 @@ export async function processLogQueue(): Promise { } } +/** + * Type for log update data that includes the logId field. + */ +interface LogUpdateData extends LogInsertData { + logId: string; +} + +export async function processLogUpdateQueue(): Promise { + const message = await consumeFromQueue(LOG_UPDATE_QUEUE); + + if (!message) { + return; + } + + const MAX_RETRIES = 5; + + try { + const logUpdateData = message.map((i) => JSON.parse(i) as LogUpdateData); + + const processedLogData: LogUpdateData[] = await Promise.all( + logUpdateData.map(async (data) => { + const org = await db.query.organization.findFirst({ + where: { + id: { + eq: data.organizationId, + }, + }, + }); + + if (org?.retentionLevel === "none") { + const { + messages: _messages, + content: _content, + reasoningContent: _reasoningContent, + tools: _tools, + toolChoice: _toolChoice, + toolResults: _toolResults, + ...metadataOnly + } = data; + return metadataOnly as LogUpdateData; + } + + return data; + }), + ); + + // Update logs with retry logic + let lastError: Error | undefined; + for (let attempt = 0; attempt <= MAX_RETRIES; attempt++) { + try { + // Process each update individually since we need to update by ID + for (const updateData of processedLogData) { + const { logId, ...fieldsToUpdate } = updateData; + await db + .update(log) + .set(fieldsToUpdate as Partial) + .where(eq(log.id, logId)); + } + return; // Success, exit function + } catch (updateError) { + lastError = + updateError instanceof Error + ? updateError + : new Error(String(updateError)); + + if (attempt < MAX_RETRIES) { + const delay = Math.pow(2, attempt) * 1000; // 1s, 2s, 4s, 8s, 16s, ... + logger.warn( + `Failed to update logs (attempt ${attempt + 1}/${MAX_RETRIES + 1}), retrying in ${delay}ms...`, + lastError, + ); + await new Promise((resolve) => { + setTimeout(resolve, delay); + }); + } + } + } + + // All retries exhausted, push messages back to queue for later processing + logger.error( + `Failed to update logs after ${MAX_RETRIES + 1} attempts, pushing back to queue`, + lastError, + ); + + // Re-add messages to queue + for (const msg of message) { + await publishToQueue(LOG_UPDATE_QUEUE, JSON.parse(msg)); + } + } catch (error) { + logger.error( + "Error processing log update message", + error instanceof Error ? error : new Error(String(error)), + ); + + // Re-add messages to queue on unexpected errors + try { + for (const msg of message) { + await publishToQueue(LOG_UPDATE_QUEUE, JSON.parse(msg)); + } + } catch (requeueError) { + logger.error( + "Failed to re-queue log update messages", + requeueError instanceof Error + ? requeueError + : new Error(String(requeueError)), + ); + } + } +} + let isWorkerRunning = false; let shouldStop = false; let minutelyIntervalId: NodeJS.Timeout | null = null; @@ -959,6 +1070,38 @@ async function runLogQueueLoop() { } } +async function runLogUpdateQueueLoop() { + activeLoops++; + logger.info("Starting log update queue processing loop..."); + try { + // eslint-disable-next-line no-unmodified-loop-condition + while (!shouldStop) { + try { + await processLogUpdateQueue(); + + if (!shouldStop) { + await new Promise((resolve) => { + setTimeout(resolve, 1000); + }); + } + } catch (error) { + logger.error( + "Error in log update queue loop", + error instanceof Error ? error : new Error(String(error)), + ); + if (!shouldStop) { + await new Promise((resolve) => { + setTimeout(resolve, 5000); + }); + } + } + } + } finally { + activeLoops--; + logger.info("Log update queue loop stopped"); + } +} + async function runAutoTopUpLoop() { activeLoops++; const interval = (process.env.NODE_ENV === "production" ? 120 : 5) * 1000; // 2 minutes in prod, 5 seconds in dev @@ -1243,6 +1386,7 @@ export async function startWorker() { // Start all parallel worker loops void runLogQueueLoop(); + void runLogUpdateQueueLoop(); void runAutoTopUpLoop(); void runBatchProcessLoop(); void runDataRetentionLoop(); diff --git a/packages/cache/src/redis.ts b/packages/cache/src/redis.ts index 6e99eb0767..6d1cabdbe1 100644 --- a/packages/cache/src/redis.ts +++ b/packages/cache/src/redis.ts @@ -16,6 +16,7 @@ redisClient.on("error", (err) => ); export const LOG_QUEUE = "log_queue_" + process.env.NODE_ENV; +export const LOG_UPDATE_QUEUE = "log_update_queue_" + process.env.NODE_ENV; export async function publishToQueue( queue: string, diff --git a/packages/db/src/schema.ts b/packages/db/src/schema.ts index f86bf0b1f5..ddd2f825d9 100644 --- a/packages/db/src/schema.ts +++ b/packages/db/src/schema.ts @@ -18,6 +18,7 @@ import type { errorDetails, tools, toolChoice, toolResults } from "./types.js"; import type z from "zod"; export const UnifiedFinishReason = { + PENDING: "pending", COMPLETED: "completed", LENGTH_LIMIT: "length_limit", CONTENT_FILTER: "content_filter",