diff --git a/CHANGELOG.md b/CHANGELOG.md index 4a2a95de5e2..b11764d6573 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -15,6 +15,7 @@ ### 🐛 Bug Fixes +- **combo/streaming: fix intermittent 500s, corrupted SSE, and Gemini malformed-response handling ([#5976](https://github.com/diegosouzapw/OmniRoute/issues/5976)).** Five related fixes on the combo streaming path: (1) the quality check now reads a **clone** of the response so the original stream stays unlocked (kills the random `ERR_INVALID_STATE: ReadableStream is locked` 500s), and the abandoned clone branch is cancelled so it no longer buffers the whole body per request; (2) `withEarlyStreamKeepalive` only emits an in-band error frame when **no** bytes were forwarded yet, so a mid-flight upstream drop no longer corrupts a partially-delivered SSE stream; (3) Gemini `finishReason: "MALFORMED_RESPONSE"` maps to OpenAI `content_filter` and triggers combo failover instead of returning broken-but-"successful" output; (4) `/api/usage/call-logs?correlationId=` uses parameterized substring (`LIKE`) matching so partial IDs resolve; (5) per-model 500s skip model-lockout/cooldown for per-model-quota providers. Plus request-logger UI detail improvements. Regression guards: `tests/unit/{earlyStreamKeepalive,finishReason,combo-provider-cooldown-sibling,combo-context-relay,call-logs-correlation-substring,save-call-log-persistence,validate-response-quality}.test.ts`. ([#6216](https://github.com/diegosouzapw/OmniRoute/pull/6216) — thanks @hartmark) - **chatcore (tools): stop the default 128-tool cap from silently dropping opencode's `task`/MCP tools.** opencode (used as an MCP/agent host) sends a large tool list; when it exceeds the speculative `MAX_TOOLS_LIMIT` (128) default, `truncateToolList` did a blind `tools.slice(0, 128)`, dropping every tool past index 128 — including opencode's built-in `task` tool (subagent launch) and many MCP tools, so models routed through OmniRoute could no longer spawn subagents or reach part of their tools. The cap exists to avoid upstream `400`s for providers with real hard limits (e.g. grok-cli 200), so it is kept for those: detection of the opencode client (`isOpencodeClient` — any `x-opencode-*` header, or `opencode` in the user-agent) now only bypasses the **speculative 128 default**, never a known provider ceiling. Precedence is explicit — a proactive/detected provider limit always truncates (even for opencode); otherwise opencode forwards its full tool list; otherwise the unchanged 128 default applies to every other client. Refactors `getEffectiveToolLimit` into `getKnownToolLimit(provider) ?? DEFAULT_LIMIT` (byte-identical for existing callers) and fixes a cosmetic debug-log that reported the truncated count instead of the original. Regression guard: `tests/unit/tool-limit-detector.test.ts`. - **fix(mitm):** the macOS MITM-cert install check now matches the system keychain again. `security find-certificate -a -Z` prints the SHA-1 as a colon-less hex string, but the installed-check compared it against `getCertFingerprint()`'s colon-separated form, so the substring match never hit — the cert was reported as not-installed and re-prompted for the sudo install on every run. Fingerprints are now normalized (colons stripped, upper-cased) on both sides via the extracted `macCertOutputHasFingerprint` helper. Regression guard: `tests/unit/mitm-cert-mac-fingerprint.test.ts`. ([#6204](https://github.com/diegosouzapw/OmniRoute/pull/6204), closes [#6134](https://github.com/diegosouzapw/OmniRoute/issues/6134) — thanks @rianonehub) diff --git a/open-sse/handlers/chatCore.ts b/open-sse/handlers/chatCore.ts index 33c33d9ffb1..3031483c0e4 100644 --- a/open-sse/handlers/chatCore.ts +++ b/open-sse/handlers/chatCore.ts @@ -374,6 +374,7 @@ export async function handleChatCore({ skipUpstreamRetry = false, createPiiTransform = null, correlationId = null, + modelPinned = false, }) { let { provider, model, extendedContext } = modelInfo; // ── Memory pressure guard ──────────────────────────────────────────── @@ -806,6 +807,7 @@ export async function handleChatCore({ apiKeyInfo, noLogEnabled, correlationId, + modelPinned, }); // Primary path: merge client model id + alias target so config on either key applies; resolved diff --git a/open-sse/handlers/chatCore/attemptLogging.ts b/open-sse/handlers/chatCore/attemptLogging.ts index 8c7151e5679..9db3efbc74f 100644 --- a/open-sse/handlers/chatCore/attemptLogging.ts +++ b/open-sse/handlers/chatCore/attemptLogging.ts @@ -51,6 +51,7 @@ export type PersistAttemptLogsContext = { apiKeyInfo: { id?: string | null; name?: string | null } | null | undefined; noLogEnabled: unknown; correlationId?: string | null; + modelPinned?: boolean; }; function toConnectionId(value: unknown): string | null { @@ -108,6 +109,7 @@ export function persistAttemptLogs(args: PersistAttemptLogsArgs, ctx: PersistAtt apiKeyInfo, noLogEnabled, correlationId, + modelPinned, } = ctx; const initialConnectionId = toConnectionId(connectionId); const finalConnectionId = toConnectionId(credentials?.connectionId) || initialConnectionId; @@ -203,5 +205,6 @@ export function persistAttemptLogs(args: PersistAttemptLogsArgs, ctx: PersistAtt noLog: noLogEnabled, pipelinePayloads, correlationId, + modelPinned: modelPinned || false, }).catch(() => {}); } diff --git a/open-sse/services/combo.ts b/open-sse/services/combo.ts index 82543d0e3be..c78970360d5 100644 --- a/open-sse/services/combo.ts +++ b/open-sse/services/combo.ts @@ -12,6 +12,7 @@ import { formatRetryAfter, getModelLockoutInfo, getRuntimeProviderProfile, + hasPerModelQuota, isModelLocked, recordModelLockoutFailure, recordProviderFailure, @@ -107,7 +108,11 @@ import { getStickyWeightedExecutionKey, recordStickyWeightedSuccess, } from "./combo/rrState.ts"; -import { validateResponseQuality, toRetryAfterDisplayValue } from "./combo/validateQuality.ts"; +import { + validateResponseQuality, + releaseQualityClone, + toRetryAfterDisplayValue, +} from "./combo/validateQuality.ts"; import { resolveComboCooldownWaitDecision } from "./combo/comboCooldownRetry.ts"; import { computeClosestRetryAfter, @@ -705,7 +710,50 @@ export async function handleComboChat({ "COMBO", `Bypassing strategy — routing directly to pinned context model: ${pinnedModel}` ); - return handleSingleModelWithTimeout(body, pinnedModel); + let pinnedResult: Response | null = null; + try { + pinnedResult = await handleSingleModelWithTimeout(body, pinnedModel, { + modelPinned: true, + } as SingleModelTarget); + } catch (pinErr) { + log.warn( + "COMBO", + `Pinned model ${pinnedModel} threw error: ${pinErr instanceof Error ? pinErr.message : String(pinErr)}, falling through to combo retry/fallback` + ); + } + if (pinnedResult) { + if (pinnedResult.ok) { + let pinnedClone: Response; + try { + pinnedClone = pinnedResult.clone(); + } catch { + pinnedClone = pinnedResult; + } + const pinnedQuality = await validateResponseQuality( + pinnedClone, + clientRequestedStream, + log, + config.responseValidation + ); + releaseQualityClone(pinnedClone, pinnedResult, pinnedQuality); + if (pinnedQuality.valid) return pinnedResult; + log.warn( + "COMBO", + `Pinned model ${pinnedModel} returned 200 but failed quality check: ${pinnedQuality.reason}, falling through to combo retry/fallback` + ); + } else { + const pinnedStatus = pinnedResult.status || 500; + if (![408, 429, 500, 502, 503, 504].includes(pinnedStatus)) { + return pinnedResult; + } + log.warn( + "COMBO", + `Pinned model ${pinnedModel} failed (${pinnedStatus}), falling through to combo retry/fallback` + ); + } + } + // Fall through to the target iteration loop below — retries and sibling + // models will be tried via the normal combo machinery. } log.warn( "COMBO", @@ -713,7 +761,8 @@ export async function handleComboChat({ ? `Context-cache pin "${pinnedModel}" provider durably unhealthy — dropping pin, using strategy` : `Stale context-cache pin "${pinnedModel}" not in combo "${combo.name}" targets — dropping pin, using strategy` ); - return handleSingleModelWithTimeout(body, pinnedModel); + // Fall through to the normal target iteration loop below — the pin is + // dropped, so the combo strategy picks the best available target. } // Fusion strategy: parallel panel + judge synthesis. Handled in a separate module @@ -1504,12 +1553,22 @@ export async function handleComboChat({ undefined; const effectiveConnectionId = selectedConnectionId || target.connectionId || ""; + // Clone BEFORE quality check — validateResponseQuality reads the body + // via getReader() which locks the stream. The clone's body is consumed + // by the quality check; the original stays unlocked for piping. + let qualityClone: Response; + try { + qualityClone = result.clone(); + } catch { + qualityClone = result; + } const quality = await validateResponseQuality( - result, + qualityClone, clientRequestedStream, log, config.responseValidation ); + releaseQualityClone(qualityClone, result, quality); if (!quality.valid) { log.warn( "COMBO", @@ -1744,7 +1803,7 @@ export async function handleComboChat({ })(); } - return { ok: true, response: quality.clonedResponse ?? result }; + return { ok: true, response: result }; } // Extract error info from response @@ -2023,7 +2082,16 @@ export async function handleComboChat({ } log.warn("COMBO", `Model ${modelStr} failed, trying next`, { status: result.status }); - if (resilienceSettings.providerCooldown.enabled && provider && provider !== "unknown") { + // #5976: per-model-quota providers (Gemini, GitHub, etc.) multiplex models + // behind one connection. A model-level 500 must NOT cool down the entire + // provider — sibling models may still succeed. Skip cooldown recording for + // these providers on 500 errors so the next target can try. + if ( + resilienceSettings.providerCooldown.enabled && + provider && + provider !== "unknown" && + !(result.status === 500 && hasPerModelQuota(provider, rawModel)) + ) { recordProviderCooldown( provider, targetWithConnection.connectionId ?? undefined, @@ -2572,12 +2640,19 @@ async function handleRoundRobinCombo({ // Success — validate response quality before returning if (result.ok) { + let rrClone: Response; + try { + rrClone = result.clone(); + } catch { + rrClone = result; + } const quality = await validateResponseQuality( - result, + rrClone, clientRequestedStream, log, config.responseValidation ); + releaseQualityClone(rrClone, result, quality); if (!quality.valid) { log.warn( "COMBO-RR", @@ -2665,12 +2740,8 @@ async function handleRoundRobinCombo({ } })(); } - // validateResponseQuality peeks streaming bodies via getReader(), - // which locks `result.body`. It returns a clonedResponse that replays - // the buffered prefix and forwards the rest. Returning the original - // (now-locked) `result` makes Next.js throw "ReadableStream is locked" - // → 500. Mirror the priority strategy and return the replay response. - return quality.clonedResponse ?? result; + // Clone is consumed by quality check; original stays unlocked. + return result; } // Extract error info @@ -2847,7 +2918,12 @@ async function handleRoundRobinCombo({ if (offset > 0) fallbackCount++; log.warn("COMBO-RR", `${modelStr} failed, trying next model`, { status: result.status }); - if (resilienceSettings.providerCooldown.enabled && provider && provider !== "unknown") { + if ( + resilienceSettings.providerCooldown.enabled && + provider && + provider !== "unknown" && + !(result.status === 500 && hasPerModelQuota(provider, parseModel(modelStr).model || modelStr)) + ) { recordProviderCooldown( provider, targetWithConnection.connectionId ?? undefined, diff --git a/open-sse/services/combo/runtimeUnits.ts b/open-sse/services/combo/runtimeUnits.ts index 96951b81a2b..d106ade3366 100644 --- a/open-sse/services/combo/runtimeUnits.ts +++ b/open-sse/services/combo/runtimeUnits.ts @@ -2,7 +2,7 @@ import { errorResponse } from "../../utils/error.ts"; import { recordComboRequest } from "../comboMetrics.ts"; import { resolveDelayMs } from "./comboPredicates.ts"; -import { validateResponseQuality } from "./validateQuality.ts"; +import { validateResponseQuality, releaseQualityClone } from "./validateQuality.ts"; import type { ResponseValidationConfig } from "./responseValidation.ts"; import type { ComboCollectionLike, @@ -228,12 +228,19 @@ export async function executeRuntimeUnitCombo(args: { }); return { response, unit }; } + let unitClone: Response; + try { + unitClone = response.clone(); + } catch { + unitClone = response; + } const quality = await validateResponseQuality( - response, + unitClone, clientRequestedStream, args.log, args.config.responseValidation as ResponseValidationConfig | undefined ); + releaseQualityClone(unitClone, response, quality); if (quality.valid) { recordComboRequest(args.combo.name, unit.modelStr, { success: true, @@ -242,7 +249,7 @@ export async function executeRuntimeUnitCombo(args: { strategy: effectiveStrategy, target: { executionKey: unit.executionKey, stepId: unit.stepId, label: unit.label }, }); - return { response: quality.clonedResponse ?? response, unit }; + return { response, unit }; } } if (![408, 429, 500, 502, 503, 504].includes(response.status)) break; diff --git a/open-sse/services/combo/types.ts b/open-sse/services/combo/types.ts index 766d6f71d95..290b0964d78 100644 --- a/open-sse/services/combo/types.ts +++ b/open-sse/services/combo/types.ts @@ -46,6 +46,8 @@ export type SingleModelTarget = allowRateLimitedConnection?: boolean; effectiveComboStrategy?: string | null; modelAbortSignal?: AbortSignal | null; + /** True when this target was selected via context-cache session pinning. */ + modelPinned?: boolean; }) | { modelAbortSignal: AbortSignal }; diff --git a/open-sse/services/combo/validateQuality.ts b/open-sse/services/combo/validateQuality.ts index fe7b00313ec..d752a789897 100644 --- a/open-sse/services/combo/validateQuality.ts +++ b/open-sse/services/combo/validateQuality.ts @@ -10,10 +10,7 @@ import { createSSEDataLineNormalizer, isKnownNonClaudeStreamPayload, } from "../../utils/streamHelpers.ts"; -import { - evaluateResponseValidation, - type ResponseValidationConfig, -} from "./responseValidation.ts"; +import { evaluateResponseValidation, type ResponseValidationConfig } from "./responseValidation.ts"; import { getReasoningTokens } from "../../../src/lib/usage/tokenAccounting.ts"; import type { ComboRetryAfter } from "./types.ts"; @@ -99,6 +96,7 @@ export async function validateResponseQuality( let hasMessageStart = false; let hasContentBlock = false; let hasLifecycleEnd = false; + let anyContentFound = false; const sseLineNormalizer = createSSEDataLineNormalizer(); let pendingEventType = ""; @@ -236,6 +234,18 @@ export async function validateResponseQuality( return { valid: false, reason: "streaming empty content block" }; } + // Non-Claude stream with no recognizable content at all — the stream + // ended without any content deltas (e.g. Gemini returning HTTP 200 + // with an empty body or only metadata chunks). Mark as invalid for + // combo failover so the sibling model gets tried. + if (!anyContentFound && !hasContentBlock) { + log.warn?.( + "COMBO", + "Streaming response ended with no recognized content — marking as invalid for combo failover" + ); + return { valid: false, reason: "streaming no recognized content" }; + } + // Incomplete lifecycle or non-Claude stream — replay all buffered // bytes. The reader is exhausted so the forwarding reader will // immediately signal done. @@ -251,6 +261,7 @@ export async function validateResponseQuality( const foundContent = parseAccumulatedSse(); if (foundContent) { + anyContentFound = true; // A content_block_* event was found — stop peeking. Return a // clonedResponse that replays all buffered bytes (the current chunk // is already in bufferedChunks) and then forwards the remainder of @@ -259,9 +270,23 @@ export async function validateResponseQuality( return { valid: true, clonedResponse }; } } - } catch { - // If reading the stream fails, pass through — other mechanisms - // (stream readiness timeout) will catch truly broken streams. + } catch (streamErr) { + // If reading the stream fails due to a locked stream or pipe error, + // the content cannot be verified — mark as invalid for combo failover. + // A locked ReadableStream means the response body is already consumed + // or corrupted (e.g. "Invalid state: The ReadableStream is locked"). + // Broad match: Chrome/V8 throws "body used already", Firefox throws + // "ReadableStream is locked", etc. + const errMsg = streamErr instanceof Error ? streamErr.message : String(streamErr); + if ( + streamErr instanceof TypeError && + (errMsg.includes("locked") || + errMsg.includes("disturbed") || + errMsg.includes("used already")) + ) { + return { valid: false, reason: "stream locked or disturbed" }; + } + // Other read errors — pass through (stream readiness timeout will catch truly broken streams) return { valid: true }; } } @@ -308,7 +333,8 @@ export async function validateResponseQuality( const choices = json?.choices; if (json?.object === "response") { - if (!responsesApiOutputHasContent(json.output)) return { valid: false, reason: "empty_choices" }; + if (!responsesApiOutputHasContent(json.output)) + return { valid: false, reason: "empty_choices" }; const status = typeof json.status === "string" ? json.status : ""; if (status && !["completed", "done"].includes(status)) { return { valid: false, reason: "no_terminal" }; @@ -389,3 +415,27 @@ export async function validateResponseQuality( }), }; } + +/** + * Release the peek-and-abandon clone used by {@link validateResponseQuality}. + * + * The quality check clones the upstream response, reads the clone only until the + * first content block, then hands back a `clonedResponse` that callers on the + * streaming path DISCARD (they forward the original, untouched response). Because + * a `Response.clone()` tees the body, that abandoned branch would otherwise buffer + * the entire remaining body in memory until the original finishes streaming. + * + * Cancelling the abandoned branch releases that buffer. Per the ReadableStream tee + * contract, cancelling one branch does NOT cancel the shared source while the other + * branch (the original response being streamed to the client) is still active, so + * this is safe. No-op when the clone fell back to the original (clone unsupported) + * or when quality reading already exhausted the body (no `clonedResponse`). + */ +export function releaseQualityClone( + clone: Response, + original: Response, + quality: { clonedResponse?: Response } +): void { + if (clone === original) return; + void quality.clonedResponse?.body?.cancel().catch(() => {}); +} diff --git a/open-sse/utils/earlyStreamKeepalive.ts b/open-sse/utils/earlyStreamKeepalive.ts index db2b98f1946..c70180f927d 100644 --- a/open-sse/utils/earlyStreamKeepalive.ts +++ b/open-sse/utils/earlyStreamKeepalive.ts @@ -54,6 +54,8 @@ export type EarlyStreamKeepaliveOptions = { * for their stream watchdog and only a real `event: ping` keeps them from aborting. */ keepaliveFrame?: Uint8Array; + /** Extra headers to include in the keepalive response (e.g. X-Correlation-Id). */ + extraHeaders?: Record; }; type SettledHandler = { ok: true; response: Response } | { ok: false; error: unknown }; @@ -66,6 +68,7 @@ export async function withEarlyStreamKeepalive( const intervalMs = Math.max(250, options.intervalMs ?? 2_500); const signal = options.signal ?? null; const keepaliveFrame = options.keepaliveFrame ?? KEEPALIVE_FRAME; + const extraHeaders = options.extraHeaders ?? {}; // Settle into a tagged result so neither race branch leaves an unhandled // rejection when the threshold timer wins. @@ -154,10 +157,25 @@ export async function withEarlyStreamKeepalive( if (response.body && isSse) { // Real SSE stream — forward it verbatim. upstreamReader = response.body.getReader(); - while (true) { - const { done, value } = await upstreamReader.read(); - if (done) break; - if (value) controller.enqueue(value); + let bytesForwarded = 0; + try { + while (true) { + const { done, value } = await upstreamReader.read(); + if (done) break; + if (value) { + controller.enqueue(value); + bytesForwarded += value.byteLength; + } + } + } catch (readErr) { + // Upstream stream failed mid-flight. Only emit an error frame if + // NO content was forwarded yet — otherwise the client already + // received partial content and a late error frame would corrupt + // the SSE stream. Silently close instead; the client will see + // the stream end naturally. + if (bytesForwarded === 0) { + controller.enqueue(ERROR_FRAME); + } } } else { // Non-SSE response (e.g. a JSON error) reached us after we already @@ -204,6 +222,7 @@ export async function withEarlyStreamKeepalive( "Content-Type": "text/event-stream; charset=utf-8", "Cache-Control": "no-cache, no-transform", Connection: "keep-alive", + ...extraHeaders, }, }); } diff --git a/open-sse/utils/finishReason.ts b/open-sse/utils/finishReason.ts index 79221ad07f4..5d8ab33667d 100644 --- a/open-sse/utils/finishReason.ts +++ b/open-sse/utils/finishReason.ts @@ -13,6 +13,7 @@ const SAFETY_FINISH_REASONS = new Set([ "prohibited_content", "content_filtered", "policy_violation", + "malformed_response", ]); export function normalizeOpenAICompatibleFinishReason(value: unknown): unknown { diff --git a/src/app/api/usage/call-logs/route.ts b/src/app/api/usage/call-logs/route.ts index e59ab37e0b0..8a3540f2540 100644 --- a/src/app/api/usage/call-logs/route.ts +++ b/src/app/api/usage/call-logs/route.ts @@ -147,8 +147,10 @@ export async function GET(request: Request) { // (active + completed) that don't match — getCallLogs already filters // the DB rows but activeEntries/completedEntries bypass it. if (filter.correlationId) { - const cid = filter.correlationId; - return NextResponse.json(rows.filter((r: any) => r.correlationId === cid)); + const cid = filter.correlationId.toLowerCase(); + return NextResponse.json( + rows.filter((r: any) => (r.correlationId || "").toLowerCase().includes(cid)) + ); } return NextResponse.json(rows); diff --git a/src/app/api/v1/chat/completions/route.ts b/src/app/api/v1/chat/completions/route.ts index 9fd7c5452cf..0bc62035d75 100644 --- a/src/app/api/v1/chat/completions/route.ts +++ b/src/app/api/v1/chat/completions/route.ts @@ -1,6 +1,7 @@ import { CORS_HEADERS, handleCorsOptions } from "@/shared/utils/cors"; import { callCloudWithMachineId } from "@/shared/utils/cloud"; import { handleChat } from "@/sse/handlers/chat"; +import { generateRequestId } from "@/shared/utils/requestId"; import { initTranslators } from "@omniroute/open-sse/translator/index.ts"; import { createInjectionGuard } from "@/middleware/promptInjectionGuard"; import { acceptHeaderForcesStream } from "@omniroute/open-sse/utils/aiSdkCompat.ts"; @@ -99,9 +100,11 @@ export async function POST(request) { const wantsStreaming = (parsedBodyIsRecord && parsedBody.stream === true) || acceptForcesStream; if (wantsStreaming) { - return await withEarlyStreamKeepalive(handleChat(request, null, parsedBody), { + const reqId = generateRequestId(); + return await withEarlyStreamKeepalive(handleChat(request, null, parsedBody, reqId), { signal: request.signal, thresholdMs: resolveKeepaliveThreshold(parsedBody?.model), + extraHeaders: { "X-Correlation-Id": reqId }, }); } diff --git a/src/lib/db/combos.ts b/src/lib/db/combos.ts index b24f6e7e9b3..d582e6184f7 100644 --- a/src/lib/db/combos.ts +++ b/src/lib/db/combos.ts @@ -7,6 +7,7 @@ import { getDbInstance } from "./core"; import { backupDbFile } from "./backup"; import { invalidateDbCache } from "./readCache"; import { normalizeComboRecord } from "@/lib/combos/steps"; +import { clearSessionModelHistoryForCombo } from "./contextHandoffs"; type JsonRecord = Record; @@ -93,7 +94,9 @@ function getNextSortOrder() { export async function getCombos() { const db = getDbInstance(); const rawCombos = db - .prepare("SELECT data, sort_order, context_cache_protection FROM combos ORDER BY sort_order ASC, name COLLATE NOCASE ASC") + .prepare( + "SELECT data, sort_order, context_cache_protection FROM combos ORDER BY sort_order ASC, name COLLATE NOCASE ASC" + ) .all() .map((row) => parseComboRow(row)) .filter((row): row is JsonRecord => row !== null); @@ -111,7 +114,9 @@ export async function getCombos() { export async function getComboById(id: string) { const db = getDbInstance(); - const row = db.prepare("SELECT data, sort_order, context_cache_protection FROM combos WHERE id = ?").get(id); + const row = db + .prepare("SELECT data, sort_order, context_cache_protection FROM combos WHERE id = ?") + .get(id); const combo = parseComboRow(row); if (!combo) return null; return normalizeStoredCombo(combo, db, typeof combo.name === "string" ? [combo.name] : []); @@ -119,7 +124,9 @@ export async function getComboById(id: string) { export async function getComboByName(name: string) { const db = getDbInstance(); - const row = db.prepare("SELECT data, sort_order, context_cache_protection FROM combos WHERE name = ?").get(name); + const row = db + .prepare("SELECT data, sort_order, context_cache_protection FROM combos WHERE name = ?") + .get(name); const combo = parseComboRow(row); if (!combo) return null; return normalizeStoredCombo(combo, db, [name]); @@ -177,7 +184,9 @@ export async function createCombo(data: JsonRecord) { export async function updateCombo(id: string, data: JsonRecord) { const db = getDbInstance(); - const existing = db.prepare("SELECT data, sort_order, context_cache_protection FROM combos WHERE id = ?").get(id); + const existing = db + .prepare("SELECT data, sort_order, context_cache_protection FROM combos WHERE id = ?") + .get(id); if (!existing) return null; const current = parseComboRow(existing); @@ -210,7 +219,26 @@ export async function updateCombo(id: string, data: JsonRecord) { db.prepare( "UPDATE combos SET name = ?, data = ?, sort_order = ?, updated_at = ?, context_cache_protection = ? WHERE id = ?" - ).run(nextName, JSON.stringify(normalizedMerged), sortOrder, normalizedMerged.updatedAt, contextCacheProtection, id); + ).run( + nextName, + JSON.stringify(normalizedMerged), + sortOrder, + normalizedMerged.updatedAt, + contextCacheProtection, + id + ); + + // Invalidate stale context-cache pins when combo targets change. + // Without this, sessions pinned to removed models keep routing there forever. + if (data.models !== undefined) { + const cleared = clearSessionModelHistoryForCombo(currentName); + if (cleared > 0) { + // Also clear under the new name if the combo was renamed + if (nextName !== currentName) { + clearSessionModelHistoryForCombo(nextName); + } + } + } invalidateDbCache("combos"); backupDbFile("pre-write"); diff --git a/src/lib/db/contextHandoffs.ts b/src/lib/db/contextHandoffs.ts index ea7f2f36b01..56575d55309 100644 --- a/src/lib/db/contextHandoffs.ts +++ b/src/lib/db/contextHandoffs.ts @@ -216,3 +216,18 @@ export function getLastSessionModel(sessionId: string, comboName: string): strin return row?.model_str ?? null; } + +/** + * Clear all session model history entries for a given combo name. + * Called when a combo's model targets are changed to invalidate stale context-cache pins. + * + * @param comboName - The combo name whose pins should be cleared. + * @returns The number of deleted entries. + */ +export function clearSessionModelHistoryForCombo(comboName: string): number { + const db = getDbInstance() as unknown as DbLike; + const result = db + .prepare("DELETE FROM session_model_history WHERE combo_name = ?") + .run(comboName); + return result.changes ?? 0; +} diff --git a/src/lib/db/schemaColumns.ts b/src/lib/db/schemaColumns.ts index 762c1b9084b..ab8ec16525f 100644 --- a/src/lib/db/schemaColumns.ts +++ b/src/lib/db/schemaColumns.ts @@ -183,6 +183,10 @@ export function ensureCallLogsColumns(db: SqliteDatabase) { db.exec("ALTER TABLE call_logs ADD COLUMN correlation_id TEXT DEFAULT NULL"); console.log("[DB] Added call_logs.correlation_id column"); } + if (!columnNames.has("model_pinned")) { + db.exec("ALTER TABLE call_logs ADD COLUMN model_pinned INTEGER DEFAULT 0"); + console.log("[DB] Added call_logs.model_pinned column"); + } db.exec( "CREATE INDEX IF NOT EXISTS idx_call_logs_requested_model ON call_logs(requested_model)" diff --git a/src/lib/usage/callLogs.ts b/src/lib/usage/callLogs.ts index 961cd0f8c2b..f5fdb7cca4d 100644 --- a/src/lib/usage/callLogs.ts +++ b/src/lib/usage/callLogs.ts @@ -91,6 +91,7 @@ type CallLogSummaryRow = { provider_node_prefix?: string | null; resolved_account?: string | null; correlation_id?: string | null; + model_pinned?: number | null; }; const RESOLVED_ACCOUNT_SQL = "COALESCE(NULLIF(pc.name, ''), NULLIF(pc.email, ''), cl.account)"; @@ -513,6 +514,7 @@ function mapSummaryRow(row: CallLogSummaryRow) { hasResponseBody: toNumber(row.has_response_body) === 1, hasPipelineDetails: toNumber(row.has_pipeline_details) === 1, correlationId: row.correlation_id || null, + modelPinned: toNumber(row.model_pinned) === 1, }; } @@ -601,6 +603,7 @@ export async function saveCallLog(entry: any) { comboExecutionKey: toStringOrNull(entry.comboExecutionKey) || toStringOrNull(entry.comboStepId), correlationId: entry.correlationId || null, + modelPinned: entry.modelPinned ? 1 : 0, }; const requestSummary = noLogEnabled @@ -649,7 +652,7 @@ export async function saveCallLog(entry: any) { combo_name, combo_step_id, combo_execution_key, error_summary, detail_state, artifact_relpath, artifact_size_bytes, artifact_sha256, has_request_body, has_response_body, has_pipeline_details, request_summary, - correlation_id + correlation_id, model_pinned ) VALUES ( @id, @timestamp, @method, @path, @status, @model, @requestedModel, @provider, @@ -660,7 +663,7 @@ export async function saveCallLog(entry: any) { @comboName, @comboStepId, @comboExecutionKey, @errorSummary, @detailState, @artifactRelPath, @artifactSizeBytes, @artifactSha256, @hasRequestBody, @hasResponseBody, @hasPipelineDetails, @requestSummary, - @correlationId + @correlationId, @modelPinned ) ` ).run({ @@ -778,8 +781,8 @@ export async function getCallLogs(filter: any = {}) { params.apiKeyQ = `%${filter.apiKey}%`; } if (filter.correlationId) { - conditions.push("cl.correlation_id = @correlationId"); - params.correlationId = filter.correlationId; + conditions.push("cl.correlation_id LIKE @correlationId"); + params.correlationId = `%${filter.correlationId}%`; } if (filter.combo) { conditions.push("cl.combo_name IS NOT NULL"); diff --git a/src/shared/components/RequestLoggerDetail.tsx b/src/shared/components/RequestLoggerDetail.tsx index d4d42b0f313..8bdd43481e1 100644 --- a/src/shared/components/RequestLoggerDetail.tsx +++ b/src/shared/components/RequestLoggerDetail.tsx @@ -547,6 +547,19 @@ export default function RequestLoggerDetail({ {cacheSourceLabel} + {(detail?.modelPinned || log.modelPinned) && ( +
+
+ Model Pinning +
+ + + + + Active — model selected via session pinning + +
+ )}
Account diff --git a/src/shared/components/RequestLoggerV2.tsx b/src/shared/components/RequestLoggerV2.tsx index 25eb736b60e..3c7b266e96f 100644 --- a/src/shared/components/RequestLoggerV2.tsx +++ b/src/shared/components/RequestLoggerV2.tsx @@ -140,6 +140,8 @@ const RequestLoggerV2 = forwardRef(null); + const [groupedView, setGroupedView] = useState(false); const [detailLoading, setDetailLoading] = useState(false); // Column sort toggle: clicking a column header toggles asc/desc @@ -372,7 +374,15 @@ const RequestLoggerV2 = forwardRef { setLimit((prev) => prev + PAGE_SIZE); @@ -432,8 +442,31 @@ const RequestLoggerV2 = forwardRef { + if (!groupedView) return filteredLogs; + const byCid = new Map(); + const noCid: typeof filteredLogs = []; + for (const log of filteredLogs) { + const cid = log.correlationId; + if (!cid) { + noCid.push(log); + continue; + } + const existing = byCid.get(cid); + if ( + !existing || + new Date(log.timestamp).getTime() > new Date(existing.timestamp).getTime() + ) { + byCid.set(cid, log); + } + } + return [...noCid, ...byCid.values()]; + }, [filteredLogs, groupedView]); + const sortedLogs = useMemo(() => { - const arr = [...filteredLogs]; + const arr = [...dedupedLogs]; arr.sort((a, b) => { switch (sortBy) { @@ -466,7 +499,7 @@ const RequestLoggerV2 = forwardRef
+ {/* Group by CID toggle */} + + {/* Provider Dropdown */}