Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
65 changes: 53 additions & 12 deletions apps/gateway/src/chat/chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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),
});
}
Comment on lines +1662 to +1686

Copilot AI Feb 2, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The pending log insertion happens at line 1666 (based on the diff), but any exceptions thrown after this point and before the first finalizeLog call will leave orphaned pending logs in the database. While the PR description mentions detecting stuck requests by querying for pending logs older than 5 minutes, there's no cleanup mechanism implemented. Without a cleanup job, these orphaned logs will accumulate indefinitely. Consider adding a periodic cleanup task in the worker to mark or remove truly stuck pending logs.

Copilot uses AI. Check for mistakes.

// Helper to finalize log - updates pending log if available, otherwise inserts new log
const finalizeLog = async (
logData: Parameters<typeof insertLog>[0],
): Promise<unknown> => {
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);
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -1918,7 +1959,7 @@ chat.openapi(completions, async (c) => {
inputImageCount,
);

await insertLog({
await finalizeLog({
...baseLogEntry,
duration,
timeToFirstToken: null, // Not applicable for cached response
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -4057,7 +4098,7 @@ chat.openapi(completions, async (c) => {
const shouldIncludeTokensForBilling =
!canceled || (canceled && billCancelledRequests);

await insertLog({
await finalizeLog({
...baseLogEntry,
duration,
timeToFirstToken,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -4859,7 +4900,7 @@ chat.openapi(completions, async (c) => {
});
}

await insertLog({
await finalizeLog({
...baseLogEntry,
duration,
timeToFirstToken: null, // Not applicable for non-streaming requests
Expand Down
121 changes: 118 additions & 3 deletions apps/gateway/src/lib/logs.ts
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -178,3 +183,113 @@ export async function insertLog(logData: LogInsertData): Promise<unknown> {
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<string> {
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",
});
Comment on lines +191 to +242

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟠 Major

Respect retentionLevel when inserting pending logs.

Pending logs are written directly to the DB, so for retentionLevel = "none" the request messages can persist indefinitely (the update path won’t clear them). This breaks retention guarantees and can leak sensitive content.

🛡️ Proposed fix (avoid storing messages when retention is "none")
 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;
+	retentionLevel?: "retain" | "none" | null;
 }

 export async function insertPendingLog(data: PendingLogData): Promise<string> {
 	const logId = shortid();
+	const storeMessages = data.retentionLevel !== "none";

 	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,
 			usedModel: data.requestedModel,
 			usedProvider: data.requestedProvider || "unknown",
 			duration: 0,
 			responseSize: 0,
 			mode: data.mode,
 			usedMode: data.usedMode,
 			unifiedFinishReason: UnifiedFinishReason.PENDING,
-			messages: data.messages,
+			messages: storeMessages ? data.messages : null,
 			streamed: data.streamed,
 			source: data.source,
 			userAgent: data.userAgent,
 			hasError: false,
 			canceled: false,
 			cached: false,
 			dataStorageCost: "0",
 		});

And pass the retention level at the call site:

 	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,
+		retentionLevel,
 	});
🤖 Prompt for AI Agents
In `@apps/gateway/src/lib/logs.ts` around lines 191 - 242, The pending log
currently always writes request messages to the DB; update insertPendingLog and
PendingLogData to accept a retentionLevel ("none" | ...) and, when
retentionLevel === "none", avoid storing messages (e.g., set messages: null or
omit it) so sensitive content isn't persisted; update the call sites to pass the
retention level into insertPendingLog and adjust any types/DB schema usage
accordingly (refer to PendingLogData and insertPendingLog to locate changes and
add the conditional use of data.retentionLevel when assigning the messages
field).


return logId;
} catch (error) {
logger.error("Failed to insert pending log", {
requestId: data.requestId,
error: error instanceof Error ? error.message : String(error),
});
throw error;
}
Comment on lines +214 to +251

Copilot AI Feb 2, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The insertPendingLog function inserts a log synchronously to the database, which could impact request latency. Since this operation happens on every request before any caching or processing begins, any database slowness or connection issues will directly affect the user-facing API response time. Consider moving this to an async fire-and-forget pattern or making it optional based on a feature flag to reduce the performance impact on the critical path.

Suggested change
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;
}
// Perform the database insert in a fire-and-forget manner so it does not
// block the request's critical path. Any errors are logged but do not
// affect the main request flow.
void (async () => {
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",
});
} catch (error) {
logger.error("Failed to insert pending log", {
requestId: data.requestId,
error: error instanceof Error ? error.message : String(error),
});
}
})();
return logId;

Copilot uses AI. Check for mistakes.
}

/**
* Data for updating an existing log entry.
* Includes the log ID and all the fields to update.
*/
export interface LogUpdateData extends Omit<LogInsertData, "id"> {
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<unknown> {
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;
}
Comment on lines +211 to +295

Copilot AI Feb 2, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The new functions insertPendingLog and updateLog lack test coverage. The existing logs.spec.ts file has tests for other log-related functions, so these new functions should also have tests. Tests should cover successful insertion, error handling, queue publishing, and the unified finish reason logic in updateLog.

Copilot uses AI. Check for mistakes.
Loading
Loading