Files
OmniRoute/open-sse/handlers/chatCore/streamingSemanticCacheStore.ts
Praveen K Palaniswamy 65e81158ab fix(ollama): route models by advertised capability (#11088)
Landed with the design call resolved per the owner's pick — **option 1**: the synced store is now endpoint-agnostic (persistDiscoveredModels and managedModelImport no longer drop non-chat models at write time), and chat selectability moved to read time (auto-pool expansion in autoStrategy applies filterChatSelectableModels; the models-route projection already had its chatOnly filter). Your discovery test now passes end-to-end (3/3): /api/show capabilities persist per connection and image/embedding requests route through the advertising host.

Reconciliation notes: conflicted areas merged onto the current tip (adobe discovery import, requestedModel preflight signature, resolvedProvider fast-path coexists with the synced-route override — explicit resolution wins); carried base-red drains (#10055 memoization, #11071 test variants) dropped as already-landed; the managed-model-import exclusion test was propagated to the new contract (image/video models persist; the read filter still hides them from chat pickers — pinned by a new assertion). Full battery: 205/206 focused (the one red is a confirmed periodic-timer timing flake on the loaded devbox — 20/20 isolated), autoCombo vitest 30/30, combo suites 46/46, gates + typecheck clean.

Thank you @yourspraveen — the capability probe + routing design was right; it just needed the store contract opened up. Fixes #11087.
2026-08-23 11:45:01 -03:00

99 lines
3.3 KiB
TypeScript

/**
* chatCore streaming semantic-cache store (Quality Gate v2 / Fase 9 — chatCore god-file
* decomposition, #3501).
*
* Extracted from handleChatCore's onStreamComplete callback: after a 200 streaming response is
* assembled, store it under its signature so a future temp=0 request can be served from cache.
* Side-effect only (cache write + debug log), wrapped in fail-open try/catch. Behaviour is
* byte-identical to the previous inline block — including the `_streamed` strip, the early
* skip-on-too-large, and the `Number(...) || 0` token accounting. The early return was the last
* statement of the callback, so returning from this helper is equivalent.
*/
import {
generateSignature as defaultGenerateSignature,
setCachedResponse as defaultSetCachedResponse,
isCacheableForWrite as defaultIsCacheableForWrite,
} from "@/lib/semanticCache";
import { isSmallEnoughForSemanticCache as defaultIsSmallEnough } from "../../utils/estimateSize.ts";
type LoggerLike = { debug?: (...args: unknown[]) => void } | null | undefined;
type CacheBody = {
messages?: unknown;
input?: unknown;
temperature?: number;
top_p?: number;
};
export interface StreamingSemanticCacheStoreDeps {
isCacheableForWrite: typeof defaultIsCacheableForWrite;
isSmallEnoughForSemanticCache: typeof defaultIsSmallEnough;
generateSignature: typeof defaultGenerateSignature;
setCachedResponse: typeof defaultSetCachedResponse;
}
const DEFAULT_DEPS: StreamingSemanticCacheStoreDeps = {
isCacheableForWrite: defaultIsCacheableForWrite,
isSmallEnoughForSemanticCache: defaultIsSmallEnough,
generateSignature: defaultGenerateSignature,
setCachedResponse: defaultSetCachedResponse,
};
interface StreamingCacheArgs {
enabled: boolean;
streamStatus: number;
streamResponseBody: Record<string, unknown> | null | undefined;
body: CacheBody;
headers: unknown;
model: string;
apiKeyId?: string;
streamUsage?: Record<string, unknown> | null;
log?: LoggerLike;
}
function streamTokensSaved(streamUsage: Record<string, unknown> | null | undefined): number {
const u = streamUsage as Record<string, unknown> | null;
return (Number(u?.prompt_tokens ?? 0) || 0) + (Number(u?.completion_tokens ?? 0) || 0);
}
function writeStreamingCacheEntry(
args: StreamingCacheArgs,
deps: StreamingSemanticCacheStoreDeps
): void {
try {
const cleanBody = { ...(args.streamResponseBody as Record<string, unknown>) };
delete cleanBody._streamed;
if (!deps.isSmallEnoughForSemanticCache(cleanBody)) return;
const sig = deps.generateSignature(
args.model,
args.body.messages ?? args.body.input,
args.body.temperature,
args.body.top_p,
args.apiKeyId ?? undefined
);
const tokensSaved = streamTokensSaved(args.streamUsage);
deps.setCachedResponse(sig, args.model, cleanBody, tokensSaved);
args.log?.debug?.(
"CACHE",
`Stored streaming response for ${args.model} (${tokensSaved} tokens)`
);
} catch {
// Cache write failed — non-critical
}
}
export function storeStreamingSemanticCacheResponse(
args: StreamingCacheArgs,
deps: StreamingSemanticCacheStoreDeps = DEFAULT_DEPS
): void {
if (
!args.enabled ||
args.streamStatus !== 200 ||
!args.streamResponseBody ||
!deps.isCacheableForWrite(args.body, args.headers)
) {
return;
}
writeStreamingCacheEntry(args, deps);
}