diff --git a/scripts/eval-retrieval.ts b/scripts/eval-retrieval.ts index f990898da..30c7f7e5c 100644 --- a/scripts/eval-retrieval.ts +++ b/scripts/eval-retrieval.ts @@ -72,6 +72,8 @@ export type GoldenRetrievalResult = { irrelevantSourceRateAt10: number; requiredSignalCoverageAt10: number; latencyMs: number; + searchTotalLatencyMs?: number; + retrievalPhaseLatenciesMs?: Record; retrievalStrategy: string | null; retrievalPlan: string | null; embeddingSkipped: boolean; @@ -572,6 +574,8 @@ export function evaluateGoldenRetrievalCase(args: { source_image_required?: boolean | null; source_image_satisfied?: boolean | null; second_stage_rerank_used?: boolean | null; + search_total_latency_ms?: number | null; + retrieval_phase_latencies_ms?: Record | null; }; latencyMs: number; timedOut?: boolean; @@ -650,6 +654,8 @@ export function evaluateGoldenRetrievalCase(args: { contentReciprocalRankAt10: contentReciprocalRankAt10(args.testCase.expectedContentTerms, args.results), ...signalMetrics, latencyMs: args.latencyMs, + searchTotalLatencyMs: args.telemetry.search_total_latency_ms ?? undefined, + retrievalPhaseLatenciesMs: args.telemetry.retrieval_phase_latencies_ms ?? undefined, retrievalStrategy: args.telemetry.retrieval_strategy ?? null, retrievalPlan: args.telemetry.retrieval_plan ?? null, embeddingSkipped: args.telemetry.embedding_skipped ?? false, @@ -773,10 +779,12 @@ export function summarizeGoldenRetrievalResults(results: GoldenRetrievalResult[] } function latencyFromTelemetry(telemetry: { + search_total_latency_ms?: number; supabase_rpc_latency_ms?: number; embedding_latency_ms?: number; rerank_latency_ms?: number; }) { + if (typeof telemetry.search_total_latency_ms === "number") return telemetry.search_total_latency_ms; return ( (telemetry.supabase_rpc_latency_ms ?? 0) + (telemetry.embedding_latency_ms ?? 0) + diff --git a/src/app/api/search/route.ts b/src/app/api/search/route.ts index 8316e88dc..08b55f9b5 100644 --- a/src/app/api/search/route.ts +++ b/src/app/api/search/route.ts @@ -478,6 +478,8 @@ function logWeakSearch(args: { } function telemetryLatencyMs(telemetry: Record) { + const total = telemetry.search_total_latency_ms; + if (typeof total === "number" && Number.isFinite(total)) return total; const keys = ["text_fast_path_latency_ms", "embedding_latency_ms", "supabase_rpc_latency_ms", "rerank_latency_ms"]; return keys.reduce((sum, key) => { const value = telemetry[key]; @@ -514,6 +516,8 @@ function telemetryRecord(telemetry: Record, key: string) { function retrievalDecisionTelemetry(telemetry: Record) { return { + search_total_latency_ms: telemetryNumber(telemetry, "search_total_latency_ms"), + retrieval_phase_latencies_ms: telemetryRecord(telemetry, "retrieval_phase_latencies_ms"), retrieval_plan: telemetryString(telemetry, "retrieval_plan"), retrieval_query_variant_count: telemetryNumber(telemetry, "retrieval_query_variant_count"), text_candidate_budget: telemetryNumber(telemetry, "text_candidate_budget"), @@ -823,6 +827,8 @@ async function buildScopedSearchPayload( smart_api_display_mode: smartApiPlan.displayMode, smart_api_source_link_count: smartApiPlan.sourceLinkCount, search_cache_hit: search.telemetry.search_cache_hit, + search_total_latency_ms: search.telemetry.search_total_latency_ms, + retrieval_phase_latencies_ms: search.telemetry.retrieval_phase_latencies_ms, shared_cache_hit: search.telemetry.shared_cache_hit, shared_cache_status: search.telemetry.shared_cache_status, shared_cache_miss_reason: search.telemetry.shared_cache_miss_reason, diff --git a/src/lib/rag-cache.ts b/src/lib/rag-cache.ts index a32ceae06..a9e2f350f 100644 --- a/src/lib/rag-cache.ts +++ b/src/lib/rag-cache.ts @@ -144,6 +144,7 @@ export async function getCachedAnswer( "query" | "documentId" | "documentIds" | "ownerId" | "accessScope" | "skipCache" | "queryMode" >, startedAt: number, + options?: { indexingVersionAtRequestStart?: string | null }, ): Promise { if (!answerCacheAllowedForOwner(args.ownerId) || args.skipCache) return null; if (env.RAG_ANSWER_CACHE_TTL_MS <= 0 || env.RAG_ANSWER_CACHE_SIZE <= 0) return null; @@ -155,7 +156,7 @@ export async function getCachedAnswer( answerCache.delete(key); return null; } - const indexingVersion = await cacheIndexingVersion(args); + const indexingVersion = options?.indexingVersionAtRequestStart ?? (await cacheIndexingVersion(args)); if (cached.indexingVersion !== indexingVersion) { answerCache.delete(key); return null; @@ -181,12 +182,8 @@ export async function setCachedAnswer( if (!answerCacheAllowedForOwner(args.ownerId) || args.skipCache) return; if (env.RAG_ANSWER_CACHE_TTL_MS <= 0 || env.RAG_ANSWER_CACHE_SIZE <= 0) return; - if (options?.indexingVersionAtRetrievalStart) { - const currentIndexingVersion = await cacheIndexingVersion(args); - if (currentIndexingVersion !== options.indexingVersionAtRetrievalStart) return; - } - - const indexingVersion = await cacheIndexingVersion(args); + const indexingVersion = await cacheIndexingVersion(args, { forceRefresh: true }); + if (options?.indexingVersionAtRetrievalStart && indexingVersion !== options.indexingVersionAtRetrievalStart) return; const key = scopedAnswerCacheKey(args); answerCache.set(key, { expiresAt: Date.now() + env.RAG_ANSWER_CACHE_TTL_MS, @@ -199,7 +196,7 @@ export async function setCachedAnswer( if (!oldestKey) break; answerCache.delete(oldestKey); } - setSharedCachedAnswer(args, answer); + setSharedCachedAnswer(args, answer, indexingVersion); } function stableHash(value: string) { @@ -256,6 +253,8 @@ function normalizeCacheStorageTelemetry(telemetry: SearchTelemetry): SearchTelem shared_cache_hit: false, shared_cache_status: undefined, shared_cache_miss_reason: null, + search_total_latency_ms: undefined, + retrieval_phase_latencies_ms: undefined, }; } @@ -267,10 +266,15 @@ export function isSearchCacheEnabled(args: Pick): return !args.skipCache && env.RAG_SEARCH_CACHE_TTL_MS > 0 && env.RAG_SEARCH_CACHE_SIZE > 0; } +export function isSearchCacheLookupEnabled(args: Pick): boolean { + return !args.skipCache && env.RAG_SEARCH_CACHE_TTL_MS > 0; +} + export async function getCachedSearch( args: SearchChunksArgs, queryClass?: RagQueryClass, queryVariants: string[] = [], + options?: { indexingVersionAtRequestStart?: string | null }, ): Promise<{ results: SearchResult[]; telemetry: SearchTelemetry } | null> { throwIfAborted(args.signal); if (!isSearchCacheEnabled(args)) return null; @@ -282,7 +286,7 @@ export async function getCachedSearch( searchCache.delete(key); return null; } - const indexingVersion = await cacheIndexingVersion(args); + const indexingVersion = options?.indexingVersionAtRequestStart ?? (await cacheIndexingVersion(args)); throwIfAborted(args.signal); if (cached.indexingVersion !== indexingVersion) { searchCache.delete(key); @@ -333,7 +337,7 @@ export async function setCachedSearch( if (!oldestKey) break; searchCache.delete(oldestKey); } - setSharedCachedSearch(args, results, cacheTelemetry, queryVariants); + setSharedCachedSearch(args, results, cacheTelemetry, indexingVersion, queryVariants); } type SharedCacheKind = "search" | "answer"; @@ -427,6 +431,7 @@ export async function getSharedCachedSearch( args: SearchChunksArgs, queryClass?: RagQueryClass, queryVariants: string[] = [], + options?: { indexingVersionAtRequestStart?: string | null }, ): Promise< | { kind: "hit"; results: SearchResult[]; telemetry: SearchTelemetry } | { kind: "miss"; reason: SharedCacheMissReason } @@ -435,7 +440,7 @@ export async function getSharedCachedSearch( throwIfAborted(args.signal); if (args.skipCache || env.RAG_SEARCH_CACHE_TTL_MS <= 0) return null; const normalizedQuery = retrievalPlanCacheQuery(args, queryClass, queryVariants); - const indexingVersion = await cacheIndexingVersion(args); + const indexingVersion = options?.indexingVersionAtRequestStart ?? (await cacheIndexingVersion(args)); try { let query = sharedCacheSelector(createAdminClient(), "search", args, indexingVersion, normalizedQuery); if (args.signal) query = query.abortSignal(args.signal); @@ -508,10 +513,11 @@ export async function getSharedCachedAnswer( "query" | "documentId" | "documentIds" | "ownerId" | "skipCache" | "queryMode" | "forceEmbedding" >, startedAt: number, + options?: { indexingVersionAtRequestStart?: string | null }, ) { if (!answerCacheAllowedForOwner(args.ownerId) || args.skipCache || env.RAG_ANSWER_CACHE_TTL_MS <= 0) return null; try { - const indexingVersion = await cacheIndexingVersion(args); + const indexingVersion = options?.indexingVersionAtRequestStart ?? (await cacheIndexingVersion(args)); const { data, error } = await sharedCacheSelector( createAdminClient(), "answer", @@ -546,13 +552,13 @@ async function replaceSharedCacheRow( >, payload: unknown, ttlMs: number, + indexingVersion: string, normalizedQuery: string = queryCacheKeyForStorage(normalizedCacheQuery(`${modeKey(args)} ${args.query}`)), ) { if (ttlMs <= 0) return; try { if (args.signal?.aborted) return; const supabase = createAdminClient(); - const indexingVersion = await cacheIndexingVersion(args); let deleteQuery = supabase .from("rag_response_cache") .delete() @@ -587,6 +593,7 @@ function setSharedCachedSearch( args: SearchChunksArgs, results: SearchResult[], telemetry: SearchTelemetry, + indexingVersion: string, queryVariants: string[] = [], ) { if (args.skipCache || env.RAG_SEARCH_CACHE_TTL_MS <= 0) return; @@ -595,6 +602,7 @@ function setSharedCachedSearch( args, { results: cloneSearchResults(results), telemetry }, env.RAG_SEARCH_CACHE_TTL_MS, + indexingVersion, retrievalPlanCacheQuery(args, telemetry.query_class, queryVariants), ); } @@ -605,6 +613,7 @@ function setSharedCachedAnswer( "query" | "documentId" | "documentIds" | "ownerId" | "skipCache" | "queryMode" | "forceEmbedding" >, answer: RagAnswer, + indexingVersion: string, ) { if (!answerCacheAllowedForOwner(args.ownerId) || args.skipCache || env.RAG_ANSWER_CACHE_TTL_MS <= 0) return; void replaceSharedCacheRow( @@ -612,6 +621,7 @@ function setSharedCachedAnswer( args, { answer: cloneAnswer(answer) }, env.RAG_ANSWER_CACHE_TTL_MS, + indexingVersion, sharedAnswerNormalizedQuery(args), ); } diff --git a/src/lib/rag-contracts.ts b/src/lib/rag-contracts.ts index f2bb3e097..345602d4a 100644 --- a/src/lib/rag-contracts.ts +++ b/src/lib/rag-contracts.ts @@ -27,10 +27,16 @@ export type SearchChunksArgs = { forceEmbedding?: boolean; // Lightweight-preview only: return lexical/trigram candidates without an embedding call. lexicalOnly?: boolean; + /** Internal: shares one request-start indexing-version resolution across nested cache reads. */ + cacheContext?: { + indexingVersionAtRequestStart?: Promise; + }; }; export type SearchTelemetry = { search_cache_hit: boolean; + search_total_latency_ms?: number; + retrieval_phase_latencies_ms?: Record; shared_cache_hit?: boolean; shared_cache_status?: "hit" | "miss"; shared_cache_miss_reason?: string | null; diff --git a/src/lib/rag.ts b/src/lib/rag.ts index 91b30e91c..6a91fe824 100644 --- a/src/lib/rag.ts +++ b/src/lib/rag.ts @@ -107,6 +107,7 @@ import { getSharedCachedAnswer, getSharedCachedSearch, isSearchCacheEnabled, + isSearchCacheLookupEnabled, packAdjacentSourceContext, packedContextCacheKey, scopedAnswerCacheKey, @@ -2257,24 +2258,24 @@ async function prepareCoverageGateResults(args: { queryClass: RagQueryClass; telemetry: SearchTelemetry; metadataCache: DocumentRankingMetadataCache; + timing: SearchTiming; }) { const startedAt = Date.now(); - const candidates = await attachDocumentRankingMetadata( - args.supabase, - args.candidates, - args.ownerId, - args.metadataCache, + const candidates = await measureSearchPhase(args.timing, "metadata_hydration", () => + attachDocumentRankingMetadata(args.supabase, args.candidates, args.ownerId, args.metadataCache), ); - let results = await attachPageVisualEvidence( - args.supabase, - selectRankedRetrievalResults({ - query: args.query, - queryClass: args.queryClass, - candidates, - topK: args.topK, - maxResultsPerDocument: args.maxResultsPerDocument, - telemetry: args.telemetry, - }), + let results = await measureSearchPhase(args.timing, "visual_hydration", () => + attachPageVisualEvidence( + args.supabase, + selectRankedRetrievalResults({ + query: args.query, + queryClass: args.queryClass, + candidates, + topK: args.topK, + maxResultsPerDocument: args.maxResultsPerDocument, + telemetry: args.telemetry, + }), + ), ); results = applySecondStageRerankIfNeeded({ queryClass: args.queryClass, @@ -2369,13 +2370,36 @@ function createSearchTelemetry(query: string, queryClass: RagQueryClass): Search }; } +type SearchTiming = { + startedAt: number; + phases: Record; +}; + +async function measureSearchPhase(timing: SearchTiming, phase: string, operation: () => Promise): Promise { + const startedAt = Date.now(); + try { + return await operation(); + } finally { + timing.phases[phase] = (timing.phases[phase] ?? 0) + (Date.now() - startedAt); + } +} + +function finishSearch(timing: SearchTiming, search: T): T { + search.telemetry.retrieval_phase_latencies_ms = { ...timing.phases }; + search.telemetry.search_total_latency_ms = Date.now() - timing.startedAt; + return search; +} + /** * Retrieves and ranks document chunks using lexical, structured, memory, and embedding-based evidence, while recording retrieval telemetry. * * @param args - Retrieval options, including the query, scope, search mode, and embedding preferences. * @returns The ranked search results and telemetry describing the retrieval process. */ -export async function searchChunksWithTelemetry(args: SearchChunksArgs) { +export async function searchChunksWithTelemetry( + args: SearchChunksArgs, +): Promise<{ results: SearchResult[]; telemetry: SearchTelemetry }> { + const searchTiming: SearchTiming = { startedAt: Date.now(), phases: {} }; args = { ...args, accessScope: retrievalAccessScopeForArgs(args) }; assertGlobalSearchAllowed(args); throwIfAborted(args.signal); @@ -2390,9 +2414,8 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { telemetry.embedding_skip_reason = "adversarial_manipulation_refused"; telemetry.retrieval_strategy = "unsupported_short_circuit"; recordSearchScoreTelemetry(telemetry, []); - return { results: [] as SearchResult[], telemetry }; + return finishSearch(searchTiming, { results: [] as SearchResult[], telemetry }); } - const indexingVersionAtRetrievalStart = await cacheIndexingVersion(args, { forceRefresh: true }); const supabase = createAdminClient(); // When the provider is source-only (offline mode, or auto mode without a usable key) we must // never call OpenAI for embeddings; retrieval falls back to the lexical text-fast-path only. @@ -2427,11 +2450,28 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { return undefined; } })(); - const queryAnalysis = await analyzeQueryWithClassifierFallback(retrievalQuery, analyzeClinicalQuery(retrievalQuery), { - corpusGrounding: corpusGroundingScope, - ownerId: args.ownerId, - signal: args.signal, - }); + const cacheContext = args.cacheContext ?? {}; + const indexingVersionPromise = isSearchCacheLookupEnabled(args) + ? (cacheContext.indexingVersionAtRequestStart ??= cacheIndexingVersion(args, { forceRefresh: true })) + : undefined; + const indexingVersionAtRetrievalStartPromise = indexingVersionPromise + ? measureSearchPhase(searchTiming, "index_version", () => indexingVersionPromise) + : Promise.resolve(null); + const queryAnalysisPromise = measureSearchPhase(searchTiming, "query_classification", () => + analyzeQueryWithClassifierFallback(retrievalQuery, analyzeClinicalQuery(retrievalQuery), { + corpusGrounding: corpusGroundingScope, + ownerId: args.ownerId, + signal: args.signal, + }), + ); + const ragAliasesPromise = measureSearchPhase(searchTiming, "alias_load", () => + fetchEnabledRagAliases(supabase, args.ownerId, args.accessScope, args.signal), + ); + const [indexingVersionAtRetrievalStart, queryAnalysis, ragAliases] = await Promise.all([ + indexingVersionAtRetrievalStartPromise, + queryAnalysisPromise, + ragAliasesPromise, + ]); throwIfAborted(args.signal); if (modeQueryClass) queryAnalysis.queryClass = modeQueryClass; const queryClassification = { @@ -2442,29 +2482,38 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { const telemetry = createSearchTelemetry(retrievalQuery, queryClassification.queryClass); if (queryAnalysis.corpusGrounding) telemetry.corpus_grounding = queryAnalysis.corpusGrounding; - const ragAliases = await fetchEnabledRagAliases(supabase, args.ownerId, args.accessScope); const ragAliasExpansions = selectRagAliasExpansions(retrievalQuery, ragAliases); telemetry.rag_alias_count = ragAliases.length; telemetry.rag_alias_expansion_count = ragAliasExpansions.length; const queryVariants = buildRetrievalQueryVariants(retrievalQuery, queryAnalysis, ragAliases); telemetry.retrieval_query_variant_count = queryVariants.length; - const cached = await getCachedSearch(args, queryClassification.queryClass, queryVariants); + const cached = await measureSearchPhase(searchTiming, "local_cache_lookup", () => + getCachedSearch(args, queryClassification.queryClass, queryVariants, { + indexingVersionAtRequestStart: indexingVersionAtRetrievalStart, + }), + ); // Only consult the shared cache when the process-local cache missed (preserves // the original short-circuit), then record the hit-rate counter ONCE with full // knowledge of both layers: a request served by either cache is a hit, so a // cold process that falls through to a warm shared cache is not miscounted as a // miss (deep /api/health cache hit-rate — docs/observability-slos.md §4). - const sharedCached = cached ? null : await getSharedCachedSearch(args, queryClassification.queryClass, queryVariants); + const sharedCached = cached + ? null + : await measureSearchPhase(searchTiming, "shared_cache_lookup", () => + getSharedCachedSearch(args, queryClassification.queryClass, queryVariants, { + indexingVersionAtRequestStart: indexingVersionAtRetrievalStart, + }), + ); const cacheOutcome = classifySearchCacheOutcome(isSearchCacheEnabled(args), Boolean(cached), sharedCached); if (cacheOutcome !== "skip") recordCacheLookup(cacheOutcome === "hit"); - if (cached) return cached; + if (cached) return finishSearch(searchTiming, cached); if (sharedCached?.kind === "hit") { await setCachedSearch(args, sharedCached.results, sharedCached.telemetry, queryVariants, { indexingVersionAtRetrievalStart, }); - return { results: sharedCached.results, telemetry: sharedCached.telemetry }; + return finishSearch(searchTiming, { results: sharedCached.results, telemetry: sharedCached.telemetry }); } if (sharedCached?.kind === "miss") { telemetry.shared_cache_status = "miss"; @@ -2487,7 +2536,16 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { min_sim: 0.45, }); if (typeof corrected === "string" && corrected && corrected.toLowerCase() !== retrievalQuery.toLowerCase()) { - return searchChunksWithTelemetry({ ...args, query: corrected, typoCorrected: true }); + const correctedSearch = await searchChunksWithTelemetry({ + ...args, + cacheContext, + query: corrected, + typoCorrected: true, + }); + for (const [phase, latency] of Object.entries(correctedSearch.telemetry.retrieval_phase_latencies_ms ?? {})) { + searchTiming.phases[phase] = (searchTiming.phases[phase] ?? 0) + latency; + } + return finishSearch(searchTiming, correctedSearch); } } telemetry.embedding_skipped = true; @@ -2495,7 +2553,7 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { telemetry.retrieval_strategy = "unsupported_short_circuit"; recordSearchScoreTelemetry(telemetry, []); await setCachedSearch(args, [], telemetry, queryVariants, { indexingVersionAtRetrievalStart }); - return { results: [] as SearchResult[], telemetry }; + return finishSearch(searchTiming, { results: [] as SearchResult[], telemetry }); } let expandedQuery = normalizeRetrievalVariant([expandClinicalQuery(retrievalQuery), ...ragAliasExpansions].join(" ")); @@ -2533,11 +2591,8 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { if (textData.length) { const rerankStartedAt = Date.now(); - const textCandidates = await attachDocumentRankingMetadata( - supabase, - textData as SearchResult[], - args.ownerId, - documentRankingMetadataCache, + const textCandidates = await measureSearchPhase(searchTiming, "metadata_hydration", () => + attachDocumentRankingMetadata(supabase, textData as SearchResult[], args.ownerId, documentRankingMetadataCache), ); expandedQuery = expandClinicalQueryWithCandidateMetadata(args.query, expandedQuery, textCandidates); const baseTextResults = selectRankedRetrievalResults({ @@ -2551,7 +2606,9 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { const baseTextFastPath = decideTextFastPath(args.query, baseTextResults, queryClassification.queryClass); if (!args.forceEmbedding && shouldReturnBeforeMemory(queryClassification.queryClass, baseTextFastPath)) { - textFastResults = await attachPageVisualEvidence(supabase, baseTextResults); + textFastResults = await measureSearchPhase(searchTiming, "visual_hydration", () => + attachPageVisualEvidence(supabase, baseTextResults), + ); textFastResults = applySecondStageRerankIfNeeded({ queryClass: queryClassification.queryClass, results: textFastResults, @@ -2563,19 +2620,21 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { telemetry.retrieval_strategy = "text_fast_path"; recordSearchScoreTelemetry(telemetry, textFastResults); await setCachedSearch(args, textFastResults, telemetry, queryVariants, { indexingVersionAtRetrievalStart }); - return { results: textFastResults, telemetry }; + return finishSearch(searchTiming, { results: textFastResults, telemetry }); } - const memoryBoost = await withMemoryBoostedCandidates({ - supabase, - query: retrievalQuery, - candidates: textCandidates, - ownerId: args.ownerId, - accessScope: args.accessScope, - documentIds: documentFilterList, - matchCount: candidateCount, - cardCache: memoryCardCache, - }); + const memoryBoost = await measureSearchPhase(searchTiming, "memory_hydration", () => + withMemoryBoostedCandidates({ + supabase, + query: retrievalQuery, + candidates: textCandidates, + ownerId: args.ownerId, + accessScope: args.accessScope, + documentIds: documentFilterList, + matchCount: candidateCount, + cardCache: memoryCardCache, + }), + ); telemetry.memory_card_count = Math.max(telemetry.memory_card_count ?? 0, memoryBoost.cards.length); telemetry.memory_top_score = Math.max( telemetry.memory_top_score ?? 0, @@ -2592,7 +2651,9 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { maxResultsPerDocument, telemetry, }); - textFastResults = await attachPageVisualEvidence(supabase, textFastResults); + textFastResults = await measureSearchPhase(searchTiming, "visual_hydration", () => + attachPageVisualEvidence(supabase, textFastResults), + ); textFastResults = applySecondStageRerankIfNeeded({ queryClass: queryClassification.queryClass, results: textFastResults, @@ -2607,7 +2668,7 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { telemetry.retrieval_strategy = "text_fast_path"; recordSearchScoreTelemetry(telemetry, textFastResults); await setCachedSearch(args, textFastResults, telemetry, queryVariants, { indexingVersionAtRetrievalStart }); - return { results: textFastResults, telemetry }; + return finishSearch(searchTiming, { results: textFastResults, telemetry }); } } @@ -2664,23 +2725,27 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { if (documentLookupData.length > 0) { const rerankStartedAt = Date.now(); - const documentLookupCandidates = await attachDocumentRankingMetadata( - supabase, - mergeSearchResults(documentLookupData, textFastResults), - args.ownerId, - documentRankingMetadataCache, + const documentLookupCandidates = await measureSearchPhase(searchTiming, "metadata_hydration", () => + attachDocumentRankingMetadata( + supabase, + mergeSearchResults(documentLookupData, textFastResults), + args.ownerId, + documentRankingMetadataCache, + ), ); expandedQuery = expandClinicalQueryWithCandidateMetadata(args.query, expandedQuery, documentLookupCandidates); - const memoryBoost = await withMemoryBoostedCandidates({ - supabase, - query: args.query, - candidates: documentLookupCandidates, - ownerId: args.ownerId, - accessScope: args.accessScope, - documentIds: documentFilterList, - matchCount: candidateCount, - cardCache: memoryCardCache, - }); + const memoryBoost = await measureSearchPhase(searchTiming, "memory_hydration", () => + withMemoryBoostedCandidates({ + supabase, + query: args.query, + candidates: documentLookupCandidates, + ownerId: args.ownerId, + accessScope: args.accessScope, + documentIds: documentFilterList, + matchCount: candidateCount, + cardCache: memoryCardCache, + }), + ); telemetry.memory_card_count = Math.max(telemetry.memory_card_count ?? 0, memoryBoost.cards.length); telemetry.memory_top_score = Math.max( telemetry.memory_top_score ?? 0, @@ -2694,16 +2759,18 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { topScore: Math.max(telemetry.memory_top_score ?? 0, ...memoryBoost.cards.map(memoryCardChunkScore)), }, ); - let documentLookupResults = await attachPageVisualEvidence( - supabase, - selectRankedRetrievalResults({ - query: retrievalQuery, - queryClass: queryClassification.queryClass, - candidates: memoryBoost.results, - topK: args.topK ?? 8, - maxResultsPerDocument, - telemetry, - }), + let documentLookupResults = await measureSearchPhase(searchTiming, "visual_hydration", () => + attachPageVisualEvidence( + supabase, + selectRankedRetrievalResults({ + query: retrievalQuery, + queryClass: queryClassification.queryClass, + candidates: memoryBoost.results, + topK: args.topK ?? 8, + maxResultsPerDocument, + telemetry, + }), + ), ); documentLookupResults = applySecondStageRerankIfNeeded({ queryClass: queryClassification.queryClass, @@ -2728,7 +2795,7 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { await setCachedSearch(args, documentLookupResults, telemetry, queryVariants, { indexingVersionAtRetrievalStart, }); - return { results: documentLookupResults, telemetry }; + return finishSearch(searchTiming, { results: documentLookupResults, telemetry }); } textFastResults = mergeSearchResults(documentLookupResults, textFastResults); } @@ -2745,6 +2812,7 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { queryClass: queryClassification.queryClass, telemetry, metadataCache: documentRankingMetadataCache, + timing: searchTiming, }); const coverageGate = evaluateEvidenceCoverageGate(args.query, coverageGateResults, queryClassification.queryClass); applyCoverageGateTelemetry(telemetry, coverageGate, !args.forceEmbedding && coverageGate.accepted); @@ -2752,7 +2820,7 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { telemetry.retrieval_strategy = coverageGate.strategy; recordSearchScoreTelemetry(telemetry, coverageGateResults); await setCachedSearch(args, coverageGateResults, telemetry, queryVariants, { indexingVersionAtRetrievalStart }); - return { results: coverageGateResults, telemetry }; + return finishSearch(searchTiming, { results: coverageGateResults, telemetry }); } textFastResults = mergeSearchResults(coverageGateResults, textFastResults); } @@ -2765,7 +2833,7 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { telemetry.embedding_skip_reason = sourceOnlyRetrieval ? SOURCE_ONLY_EMBEDDING_SKIP_REASON : "lexical_only"; telemetry.retrieval_strategy = telemetry.retrieval_strategy ?? "text_fast_path"; recordSearchScoreTelemetry(telemetry, textFastResults); - return { results: textFastResults, telemetry }; + return finishSearch(searchTiming, { results: textFastResults, telemetry }); } throwIfAborted(args.signal); @@ -2782,7 +2850,7 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { telemetry.vector_skipped_reason = classifyProviderFailure(error); telemetry.retrieval_strategy = telemetry.retrieval_strategy ?? "text_fast_path"; recordSearchScoreTelemetry(telemetry, textFastResults); - return { results: textFastResults, telemetry }; + return finishSearch(searchTiming, { results: textFastResults, telemetry }); } const { embedding, cacheHit } = embeddingResult; telemetry.embedding_latency_ms = Date.now() - embeddingStartedAt; @@ -2900,38 +2968,39 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { if (!hybridError) { const rerankStartedAt = Date.now(); const merged = args.forceEmbedding ? vectorCandidates : mergeSearchResults(vectorCandidates, textFastResults); - const mergedWithMetadata = await attachDocumentRankingMetadata( - supabase, - merged, - args.ownerId, - documentRankingMetadataCache, + const mergedWithMetadata = await measureSearchPhase(searchTiming, "metadata_hydration", () => + attachDocumentRankingMetadata(supabase, merged, args.ownerId, documentRankingMetadataCache), + ); + const memoryBoost = await measureSearchPhase(searchTiming, "memory_hydration", () => + withMemoryBoostedCandidates({ + supabase, + query: retrievalQuery, + candidates: mergedWithMetadata, + queryEmbedding: embedding, + ownerId: args.ownerId, + accessScope: args.accessScope, + documentIds: documentFilterList, + matchCount: candidateCount, + cardCache: memoryCardCache, + }), ); - const memoryBoost = await withMemoryBoostedCandidates({ - supabase, - query: retrievalQuery, - candidates: mergedWithMetadata, - queryEmbedding: embedding, - ownerId: args.ownerId, - accessScope: args.accessScope, - documentIds: documentFilterList, - matchCount: candidateCount, - cardCache: memoryCardCache, - }); telemetry.memory_card_count = Math.max(telemetry.memory_card_count ?? 0, memoryBoost.cards.length); telemetry.memory_top_score = Math.max( telemetry.memory_top_score ?? 0, ...memoryBoost.cards.map(memoryCardChunkScore), ); - let results = await attachPageVisualEvidence( - supabase, - selectRankedRetrievalResults({ - query: retrievalQuery, - queryClass: queryClassification.queryClass, - candidates: memoryBoost.results, - topK: args.topK ?? 8, - maxResultsPerDocument, - telemetry, - }), + let results = await measureSearchPhase(searchTiming, "visual_hydration", () => + attachPageVisualEvidence( + supabase, + selectRankedRetrievalResults({ + query: retrievalQuery, + queryClass: queryClassification.queryClass, + candidates: memoryBoost.results, + topK: args.topK ?? 8, + maxResultsPerDocument, + telemetry, + }), + ), ); results = applySecondStageRerankIfNeeded({ queryClass: queryClassification.queryClass, @@ -2943,7 +3012,7 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { telemetry.retrieval_strategy = "hybrid"; recordSearchScoreTelemetry(telemetry, results); await setCachedSearch(args, results, telemetry, queryVariants, { indexingVersionAtRetrievalStart }); - return { results, telemetry }; + return finishSearch(searchTiming, { results, telemetry }); } const vectorFilters = documentFilterList?.length ? documentFilterList : [null]; @@ -2986,38 +3055,44 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { mergeSearchResults(resultSets.flat(), embeddingFieldCandidates), indexUnitCandidates, ); - const mergedWithMetadata = await attachDocumentRankingMetadata( - supabase, - args.forceEmbedding ? fallbackVectorCandidates : mergeSearchResults(fallbackVectorCandidates, textFastResults), - args.ownerId, - documentRankingMetadataCache, + const mergedWithMetadata = await measureSearchPhase(searchTiming, "metadata_hydration", () => + attachDocumentRankingMetadata( + supabase, + args.forceEmbedding ? fallbackVectorCandidates : mergeSearchResults(fallbackVectorCandidates, textFastResults), + args.ownerId, + documentRankingMetadataCache, + ), + ); + const memoryBoost = await measureSearchPhase(searchTiming, "memory_hydration", () => + withMemoryBoostedCandidates({ + supabase, + query: retrievalQuery, + candidates: mergedWithMetadata, + queryEmbedding: embedding, + ownerId: args.ownerId, + accessScope: args.accessScope, + documentIds: documentFilterList, + matchCount: candidateCount, + cardCache: memoryCardCache, + }), ); - const memoryBoost = await withMemoryBoostedCandidates({ - supabase, - query: retrievalQuery, - candidates: mergedWithMetadata, - queryEmbedding: embedding, - ownerId: args.ownerId, - accessScope: args.accessScope, - documentIds: documentFilterList, - matchCount: candidateCount, - cardCache: memoryCardCache, - }); telemetry.memory_card_count = Math.max(telemetry.memory_card_count ?? 0, memoryBoost.cards.length); telemetry.memory_top_score = Math.max( telemetry.memory_top_score ?? 0, ...memoryBoost.cards.map(memoryCardChunkScore), ); - let results = await attachPageVisualEvidence( - supabase, - selectRankedRetrievalResults({ - query: retrievalQuery, - queryClass: queryClassification.queryClass, - candidates: memoryBoost.results, - topK: args.topK ?? 8, - maxResultsPerDocument, - telemetry, - }), + let results = await measureSearchPhase(searchTiming, "visual_hydration", () => + attachPageVisualEvidence( + supabase, + selectRankedRetrievalResults({ + query: retrievalQuery, + queryClass: queryClassification.queryClass, + candidates: memoryBoost.results, + topK: args.topK ?? 8, + maxResultsPerDocument, + telemetry, + }), + ), ); results = applySecondStageRerankIfNeeded({ queryClass: queryClassification.queryClass, @@ -3029,7 +3104,7 @@ export async function searchChunksWithTelemetry(args: SearchChunksArgs) { telemetry.retrieval_strategy = "vector_fallback"; recordSearchScoreTelemetry(telemetry, results); await setCachedSearch(args, results, telemetry, queryVariants, { indexingVersionAtRetrievalStart }); - return { results, telemetry }; + return finishSearch(searchTiming, { results, telemetry }); } /** Build related documents safe. */ @@ -3232,8 +3307,16 @@ async function answerQuestionWithScopeUncoalesced( // unchanged cache version) would bypass chooseAnswerRoute's refusal. Skipping the // cache lets the query flow to routing, which fails it closed to "unsupported". const adversarialQuery = hasAdversarialManipulationIntent(answerFocusQuery); - const indexingVersionAtRetrievalStart = adversarialQuery || args.skipCache ? null : await cacheIndexingVersion(args); - const cachedAnswer = adversarialQuery ? null : await getCachedAnswer(args, startedAt); + const cacheContext = args.cacheContext ?? {}; + const answerCacheLookupEnabled = + !adversarialQuery && answerCacheAllowedForOwner(args.ownerId) && !args.skipCache && env.RAG_ANSWER_CACHE_TTL_MS > 0; + const indexingVersionPromise = answerCacheLookupEnabled + ? (cacheContext.indexingVersionAtRequestStart ??= cacheIndexingVersion(args, { forceRefresh: true })) + : undefined; + const indexingVersionAtRetrievalStart = indexingVersionPromise ? await indexingVersionPromise : null; + const cachedAnswer = adversarialQuery + ? null + : await getCachedAnswer(args, startedAt, { indexingVersionAtRequestStart: indexingVersionAtRetrievalStart }); if (cachedAnswer) { const cachedSources = annotateSearchResults(answerFocusQuery, cachedAnswer.sources ?? []); const cachedRelevance = cachedAnswer.relevance ?? buildEvidenceRelevance(answerFocusQuery, cachedSources); @@ -3258,7 +3341,9 @@ async function answerQuestionWithScopeUncoalesced( : cachedAnswer.smartPanel, }); } - const sharedCachedAnswer = adversarialQuery ? null : await getSharedCachedAnswer(args, startedAt); + const sharedCachedAnswer = adversarialQuery + ? null + : await getSharedCachedAnswer(args, startedAt, { indexingVersionAtRequestStart: indexingVersionAtRetrievalStart }); if (sharedCachedAnswer) { await setCachedAnswer(args, sharedCachedAnswer, { indexingVersionAtRetrievalStart }); const cachedSources = annotateSearchResults(answerFocusQuery, sharedCachedAnswer.sources ?? []); @@ -3305,6 +3390,7 @@ async function answerQuestionWithScopeUncoalesced( skipCache: args.skipCache, queryMode: args.queryMode, signal: retrievalDeadline.signal, + cacheContext, }), ); } finally { diff --git a/tests/eval-retrieval.test.ts b/tests/eval-retrieval.test.ts index 68d89d7ad..ac323defd 100644 --- a/tests/eval-retrieval.test.ts +++ b/tests/eval-retrieval.test.ts @@ -60,6 +60,30 @@ describe("golden retrieval eval helpers", () => { expect(evaluated.failures).toEqual([]); }); + it("exposes total and phase latency telemetry in retrieval evaluation output", () => { + const evaluated = evaluateGoldenRetrievalCase({ + testCase: { + id: "latency-telemetry", + query: "What ANC threshold applies?", + expectedQueryClass: "table_threshold", + expectedDocumentSubstrings: [], + expectedContentTerms: [], + topK: 8, + expectTableEvidence: false, + }, + results: [result()], + telemetry: { + query_class: "table_threshold", + search_total_latency_ms: 321, + retrieval_phase_latencies_ms: { query_classification: 12, metadata_hydration: 34 }, + }, + latencyMs: 321, + }); + + expect(evaluated.searchTotalLatencyMs).toBe(321); + expect(evaluated.retrievalPhaseLatenciesMs).toEqual({ query_classification: 12, metadata_hydration: 34 }); + }); + it("scores ideal graded signal ranking, coverage, and irrelevant sources", () => { const evaluated = evaluateGoldenRetrievalCase({ testCase: { diff --git a/tests/rag-tail-latency.test.ts b/tests/rag-tail-latency.test.ts new file mode 100644 index 000000000..e202d0694 --- /dev/null +++ b/tests/rag-tail-latency.test.ts @@ -0,0 +1,568 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import type { SearchTelemetry } from "../src/lib/rag-contracts"; +import type { RagAnswer } from "../src/lib/types"; + +const ownerId = "aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa"; +const ragVersion = "test-rag-version"; + +function indexingVersion(stamp: string) { + return `${ragVersion}:doc-1:${stamp}:`; +} + +function baseTelemetry(overrides: Partial = {}): SearchTelemetry { + return { + search_cache_hit: false, + text_fast_path_latency_ms: 0, + embedding_skipped: true, + embedding_latency_ms: 0, + embedding_cache_hit: false, + supabase_rpc_latency_ms: 0, + rerank_latency_ms: 0, + ...overrides, + }; +} + +type CacheHarness = ReturnType; + +function createCacheHarness() { + let documentStamp = "request-start"; + let documentReads = 0; + let clientCreations = 0; + let advanceClockOnSecondClient = false; + const sharedReadVersions: string[] = []; + const inserts: Array> = []; + + function chain(table: string) { + const filters = new Map(); + const builder = { + select: () => builder, + delete: () => builder, + eq: (column: string, value: unknown) => { + filters.set(column, value); + return builder; + }, + is: (column: string, value: unknown) => { + filters.set(column, value); + return builder; + }, + in: () => builder, + or: () => builder, + gt: () => builder, + order: () => builder, + limit: () => builder, + maybeSingle: async () => { + sharedReadVersions.push(String(filters.get("indexing_version"))); + return { data: null, error: null }; + }, + then: ( + resolve: (value: { data: unknown[] | null; error: null }) => unknown, + reject?: (reason: unknown) => unknown, + ) => { + if (table === "documents") { + documentReads += 1; + return Promise.resolve({ + data: [{ id: "doc-1", updated_at: documentStamp, metadata: {} }], + error: null, + }).then(resolve, reject); + } + return Promise.resolve({ data: [], error: null }).then(resolve, reject); + }, + }; + return builder; + } + + const createAdminClient = vi.fn(() => { + clientCreations += 1; + if (advanceClockOnSecondClient && clientCreations === 2) { + vi.setSystemTime(new Date(Date.now() + 6_000)); + } + return { + from: (table: string) => ({ + ...chain(table), + insert: async (row: Record) => { + inserts.push(row); + return { data: null, error: null }; + }, + }), + }; + }); + + return { + createAdminClient, + get documentReads() { + return documentReads; + }, + get sharedReadVersions() { + return sharedReadVersions; + }, + get inserts() { + return inserts; + }, + setDocumentStamp(stamp: string) { + documentStamp = stamp; + }, + expireVersionBeforeSharedWrite() { + advanceClockOnSecondClient = true; + }, + }; +} + +async function loadCache(harness: CacheHarness) { + vi.doMock("@/lib/env", () => ({ + env: { + RAG_SEARCH_CACHE_TTL_MS: 60_000, + RAG_SEARCH_CACHE_SIZE: 200, + RAG_ANSWER_CACHE_TTL_MS: 60_000, + RAG_ANSWER_CACHE_SIZE: 200, + RAG_PERSIST_RAW_QUERY_TEXT: false, + RAG_QUERY_HASH_SECRET: "test-query-hash-secret", + }, + isDemoMode: () => false, + isLocalNoAuthMode: () => false, + })); + vi.doMock("@/lib/deep-memory", () => ({ ragDeepMemoryVersion: ragVersion })); + vi.doMock("@/lib/clinical-search", () => ({ + buildClinicalTextSearchQuery: (query: string) => query.trim(), + })); + vi.doMock("@/lib/supabase/admin", () => ({ createAdminClient: harness.createAdminClient })); + return import("../src/lib/rag-cache"); +} + +afterEach(() => { + vi.useRealTimers(); + vi.restoreAllMocks(); + vi.resetModules(); + vi.unstubAllEnvs(); +}); + +describe("RAG cache request indexing version", () => { + it("reuses one request-start version across shared answer and search cache reads", async () => { + const harness = createCacheHarness(); + const cache = await loadCache(harness); + const requestVersion = indexingVersion("request-start"); + const args = { query: "lithium monitoring", ownerId }; + + await cache.getSharedCachedAnswer(args, Date.now(), { indexingVersionAtRequestStart: requestVersion }); + await cache.getSharedCachedSearch(args, "medication_dose_risk", [], { + indexingVersionAtRequestStart: requestVersion, + }); + + expect(harness.documentReads).toBe(0); + expect(harness.sharedReadVersions).toEqual([requestVersion, requestVersion]); + }); + + it("uses the request-start version for a process-local search cache read", async () => { + const harness = createCacheHarness(); + const cache = await loadCache(harness); + const args = { query: "lithium monitoring", ownerId }; + const requestVersion = indexingVersion("request-start"); + + await cache.setCachedSearch(args, [], baseTelemetry({ query_class: "medication_dose_risk" }), [], { + indexingVersionAtRetrievalStart: requestVersion, + }); + harness.setDocumentStamp("changed-after-request-start"); + await cache.cacheIndexingVersion(args, { forceRefresh: true }); + + const cached = await cache.getCachedSearch(args, "medication_dose_risk", [], { + indexingVersionAtRequestStart: requestVersion, + }); + + expect(cached?.telemetry.search_cache_hit).toBe(true); + }); + + it("keeps one forced write refresh, rejects changed indexing, and passes the refresh to the shared writer", async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-07-14T00:00:00.000Z")); + const harness = createCacheHarness(); + harness.setDocumentStamp("changed"); + const cache = await loadCache(harness); + const args = { query: "lithium monitoring", ownerId }; + + await cache.setCachedSearch(args, [], baseTelemetry({ query_class: "medication_dose_risk" }), [], { + indexingVersionAtRetrievalStart: indexingVersion("request-start"), + }); + await Promise.resolve(); + + expect(harness.documentReads).toBe(1); + expect(harness.inserts).toHaveLength(0); + + vi.resetModules(); + const matchingHarness = createCacheHarness(); + matchingHarness.expireVersionBeforeSharedWrite(); + const matchingCache = await loadCache(matchingHarness); + await matchingCache.setCachedSearch(args, [], baseTelemetry({ query_class: "medication_dose_risk" }), [], { + indexingVersionAtRetrievalStart: indexingVersion("request-start"), + }); + await vi.waitFor(() => expect(matchingHarness.inserts).toHaveLength(1)); + + expect(matchingHarness.documentReads).toBe(1); + expect(matchingHarness.inserts[0]?.indexing_version).toBe(indexingVersion("request-start")); + }); +}); + +type Deferred = { + promise: Promise; + resolve: (value: T) => void; +}; + +function deferred(): Deferred { + let resolve!: (value: T) => void; + const promise = new Promise((done) => { + resolve = done; + }); + return { promise, resolve }; +} + +class EmptyQuery implements PromiseLike<{ data: unknown[]; error: null }> { + select() { + return this; + } + in() { + return this; + } + eq() { + return this; + } + is() { + return this; + } + neq() { + return this; + } + or() { + return this; + } + order() { + return this; + } + limit() { + return Promise.resolve({ data: [], error: null }); + } + then( + onfulfilled?: ((value: { data: unknown[]; error: null }) => TResult1 | PromiseLike) | null, + onrejected?: ((reason: unknown) => TResult2 | PromiseLike) | null, + ): PromiseLike { + return Promise.resolve({ data: [], error: null }).then(onfulfilled, onrejected); + } +} + +async function loadSearchWithCacheOutcome( + outcome: "cold" | "local" | "shared", + bootstrap?: { + version: Deferred; + classification: Deferred<{ verdict: "out_of_corpus" }>; + aliases: Deferred; + }, +) { + vi.doUnmock("@/lib/env"); + vi.doUnmock("@/lib/deep-memory"); + vi.doUnmock("@/lib/clinical-search"); + const staleTelemetry = baseTelemetry({ + search_cache_hit: true, + search_total_latency_ms: 99_999, + retrieval_phase_latencies_ms: { stale_phase: 99_999 }, + retrieval_strategy: "search_cache", + }); + const getCachedSearch = vi.fn(async () => + outcome === "local" ? { results: [], telemetry: { ...staleTelemetry } } : null, + ); + const getSharedCachedSearch = vi.fn(async () => + outcome === "shared" ? { kind: "hit" as const, results: [], telemetry: { ...staleTelemetry } } : null, + ); + let versionStarted = false; + let aliasesStarted = false; + let classificationStarted = false; + const fetchEnabledRagAliases = vi.fn(() => { + aliasesStarted = true; + return bootstrap?.aliases.promise ?? Promise.resolve([]); + }); + + vi.doMock("@/lib/rag-cache", async () => { + const actual = await vi.importActual("@/lib/rag-cache"); + return { + ...actual, + cacheIndexingVersion: vi.fn(() => { + versionStarted = true; + return bootstrap?.version.promise ?? Promise.resolve(indexingVersion("request-start")); + }), + getCachedSearch, + getSharedCachedSearch, + isSearchCacheEnabled: () => true, + setCachedSearch: vi.fn(async () => undefined), + }; + }); + vi.doMock("@/lib/rag-retrieval-variants", async () => { + const actual = + await vi.importActual("@/lib/rag-retrieval-variants"); + return { + ...actual, + fetchEnabledRagAliases, + }; + }); + vi.doMock("@/lib/corpus-grounding", () => ({ + classifyCorpusGrounding: vi.fn(() => { + classificationStarted = true; + return bootstrap?.classification.promise ?? Promise.resolve({ verdict: "out_of_corpus" }); + }), + })); + vi.doMock("@/lib/rag-provider", () => ({ + isSourceOnlyMode: () => true, + allowsAutoDegrade: () => true, + sourceOnlyReason: () => "source_only", + classifyProviderFailure: () => "provider_failure", + SOURCE_ONLY_EMBEDDING_SKIP_REASON: "source_only", + })); + vi.doMock("@/lib/supabase/admin", () => ({ + createAdminClient: () => ({ + rpc: vi.fn(async () => ({ data: [], error: null })), + from: vi.fn(() => new EmptyQuery()), + }), + })); + vi.stubEnv("RAG_SEARCH_CACHE_TTL_MS", "60000"); + vi.stubEnv("RAG_SEARCH_CACHE_SIZE", "200"); + vi.stubEnv("OPENAI_API_KEY", ""); + + const { searchChunksWithTelemetry } = await import("../src/lib/rag"); + return { + searchChunksWithTelemetry, + get started() { + return { versionStarted, aliasesStarted, classificationStarted }; + }, + fetchEnabledRagAliases, + }; +} + +describe("RAG search bootstrap and latency telemetry", () => { + it("starts indexing-version, query-classification, and alias bootstrap work concurrently", async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-07-14T00:00:00.000Z")); + const bootstrap = { + version: deferred(), + classification: deferred<{ verdict: "out_of_corpus" }>(), + aliases: deferred(), + }; + const loaded = await loadSearchWithCacheOutcome("cold", bootstrap); + const controller = new AbortController(); + + const pending = loaded.searchChunksWithTelemetry({ + query: "bipolar disorder", + ownerId, + lexicalOnly: true, + signal: controller.signal, + }); + await Promise.resolve(); + + expect(loaded.started).toEqual({ + versionStarted: true, + aliasesStarted: true, + classificationStarted: true, + }); + expect(loaded.fetchEnabledRagAliases).toHaveBeenCalledWith( + expect.anything(), + ownerId, + expect.objectContaining({ ownerId, includePublic: true }), + controller.signal, + ); + + vi.setSystemTime(new Date("2026-07-14T00:00:00.100Z")); + bootstrap.version.resolve(indexingVersion("request-start")); + bootstrap.classification.resolve({ verdict: "out_of_corpus" }); + bootstrap.aliases.resolve([]); + const result = await pending; + + expect(result.telemetry.retrieval_phase_latencies_ms).toMatchObject({ + index_version: 100, + query_classification: 100, + alias_load: 100, + }); + expect(result.telemetry.search_total_latency_ms).toBe(100); + expect(result.telemetry.search_total_latency_ms).toBeLessThan( + Object.values(result.telemetry.retrieval_phase_latencies_ms ?? {}).reduce((sum, value) => sum + value, 0), + ); + }); + + it.each(["cold", "local", "shared"] as const)( + "records current phase and total telemetry for the %s path", + async (outcome) => { + const loaded = await loadSearchWithCacheOutcome(outcome); + const result = await loaded.searchChunksWithTelemetry({ + query: outcome === "cold" ? "bipolar disorder" : "clozapine monitoring", + ownerId, + lexicalOnly: true, + }); + + expect(result.telemetry.search_total_latency_ms).toEqual(expect.any(Number)); + expect(result.telemetry.search_total_latency_ms).not.toBe(99_999); + expect(result.telemetry.retrieval_phase_latencies_ms).toEqual( + expect.objectContaining({ + index_version: expect.any(Number), + query_classification: expect.any(Number), + alias_load: expect.any(Number), + local_cache_lookup: expect.any(Number), + }), + ); + expect(result.telemetry.retrieval_phase_latencies_ms).not.toHaveProperty("stale_phase"); + if (outcome !== "local") { + expect(result.telemetry.retrieval_phase_latencies_ms).toHaveProperty("shared_cache_lookup"); + } + }, + ); +}); + +type OrchestrationEvent = + | { kind: "index_version"; value: string } + | { kind: "shared_read"; cacheKind: string; indexingVersion: string } + | { kind: "shared_write"; cacheKind: string; indexingVersion: string }; + +function createOrchestrationHarness() { + let documentStamp = "stale"; + let events: OrchestrationEvent[] = []; + + class Query implements PromiseLike<{ data: unknown[]; error: null }> { + private readonly filters = new Map(); + + constructor(private readonly table: string) {} + + select() { + return this; + } + delete() { + return this; + } + eq(column: string, value: unknown) { + this.filters.set(column, value); + return this; + } + is(column: string, value: unknown) { + this.filters.set(column, value); + return this; + } + in() { + return this; + } + abortSignal() { + return this; + } + or() { + return this; + } + gt() { + return this; + } + neq() { + return this; + } + order() { + return this; + } + limit() { + return this; + } + async maybeSingle() { + events.push({ + kind: "shared_read", + cacheKind: String(this.filters.get("cache_kind")), + indexingVersion: String(this.filters.get("indexing_version")), + }); + return { data: null, error: null }; + } + then( + onfulfilled?: ((value: { data: unknown[]; error: null }) => TResult1 | PromiseLike) | null, + onrejected?: ((reason: unknown) => TResult2 | PromiseLike) | null, + ): PromiseLike { + if (this.table === "documents") { + events.push({ kind: "index_version", value: documentStamp }); + return Promise.resolve({ + data: [{ id: "doc-1", updated_at: documentStamp, metadata: {} }], + error: null, + }).then(onfulfilled, onrejected); + } + return Promise.resolve({ data: [], error: null }).then(onfulfilled, onrejected); + } + } + + return { + createAdminClient: () => ({ + rpc: vi.fn(async () => ({ data: [], error: null })), + from: (table: string) => + Object.assign(new Query(table), { + insert: async (row: Record) => { + events.push({ + kind: "shared_write", + cacheKind: String(row.cache_kind), + indexingVersion: String(row.indexing_version), + }); + return { data: null, error: null }; + }, + }), + }), + get events() { + return events; + }, + setDocumentStamp(stamp: string) { + documentStamp = stamp; + }, + clearEvents() { + events = []; + }, + }; +} + +describe("answer-to-search cache request context", () => { + it("reuses one request-start indexing version across real answer and nested search cache reads", async () => { + vi.doUnmock("@/lib/env"); + vi.doUnmock("@/lib/deep-memory"); + vi.doUnmock("@/lib/clinical-search"); + vi.doUnmock("@/lib/rag-cache"); + vi.doUnmock("@/lib/rag-retrieval-variants"); + vi.doUnmock("@/lib/corpus-grounding"); + vi.doUnmock("@/lib/rag-provider"); + vi.stubEnv("RAG_PROVIDER_MODE", "offline"); + vi.stubEnv("RAG_SEARCH_CACHE_TTL_MS", "60000"); + vi.stubEnv("RAG_SEARCH_CACHE_SIZE", "200"); + vi.stubEnv("RAG_ANSWER_CACHE_TTL_MS", "60000"); + vi.stubEnv("RAG_ANSWER_CACHE_SIZE", "200"); + vi.stubEnv("OPENAI_API_KEY", ""); + + const harness = createOrchestrationHarness(); + vi.doMock("@/lib/supabase/admin", () => ({ createAdminClient: harness.createAdminClient })); + + const rag = await import("../src/lib/rag"); + const cache = await import("../src/lib/rag-cache"); + const args = { query: "coffee machine policy", ownerId, logQuery: false }; + + // Prime stale process-local answer and search entries. The orchestration request below + // must reject both local entries, continue through both shared reads, and carry the one + // request-start version into the nested search instead of resolving it again. + await rag.searchChunksWithTelemetry(args); + await cache.setCachedAnswer(args, {} as RagAnswer); + await vi.waitFor(() => expect(harness.events.some((event) => event.kind === "shared_write")).toBe(true)); + + harness.setDocumentStamp("request-start"); + harness.clearEvents(); + + const answer = await rag.answerQuestionWithScope(args); + await vi.waitFor(() => + expect(harness.events.filter((event) => event.kind === "shared_read").map((event) => event.cacheKind)).toEqual([ + "answer", + "search", + ]), + ); + + const sharedReads = harness.events.filter( + (event): event is Extract => event.kind === "shared_read", + ); + const requestVersion = sharedReads[0]?.indexingVersion; + const searchReadIndex = harness.events.findIndex( + (event) => event.kind === "shared_read" && event.cacheKind === "search", + ); + const versionsResolvedBeforeNestedSearch = harness.events + .slice(0, searchReadIndex) + .filter((event) => event.kind === "index_version"); + + expect(answer.routingReason).not.toContain("cache_hit"); + expect(sharedReads.map((event) => event.indexingVersion)).toEqual([requestVersion, requestVersion]); + expect(requestVersion).toMatch(/:doc-1:request-start:$/); + expect(versionsResolvedBeforeNestedSearch).toEqual([{ kind: "index_version", value: "request-start" }]); + }); +}); diff --git a/tests/rag-variant-early-exit.test.ts b/tests/rag-variant-early-exit.test.ts index fc5395117..147bd8187 100644 --- a/tests/rag-variant-early-exit.test.ts +++ b/tests/rag-variant-early-exit.test.ts @@ -125,7 +125,9 @@ describe("lexical variant early-exit (PT-02)", () => { expect(chunkTextCalls).toHaveLength(1); expect(telemetry.text_variant_early_exit).toBe(true); expect(telemetry.text_variant_rpc_calls?.match_document_chunks_text).toBe(1); - expect(from).toHaveBeenCalledWith("documents"); + // Cache TTL is disabled for this fixture, so retrieval must not pay the + // indexing-version `documents` preflight before issuing lexical RPCs. + expect(from).not.toHaveBeenCalledWith("documents"); expect(from).not.toHaveBeenCalledWith("document_pages"); });