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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions scripts/eval-retrieval.ts
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,8 @@ export type GoldenRetrievalResult = {
irrelevantSourceRateAt10: number;
requiredSignalCoverageAt10: number;
latencyMs: number;
searchTotalLatencyMs?: number;
retrievalPhaseLatenciesMs?: Record<string, number>;
retrievalStrategy: string | null;
retrievalPlan: string | null;
embeddingSkipped: boolean;
Expand Down Expand Up @@ -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<string, number> | null;
};
latencyMs: number;
timedOut?: boolean;
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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) +
Expand Down
6 changes: 6 additions & 0 deletions src/app/api/search/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -478,6 +478,8 @@ function logWeakSearch(args: {
}

function telemetryLatencyMs(telemetry: Record<string, unknown>) {
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];
Expand Down Expand Up @@ -514,6 +516,8 @@ function telemetryRecord(telemetry: Record<string, unknown>, key: string) {

function retrievalDecisionTelemetry(telemetry: Record<string, unknown>) {
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"),
Expand Down Expand Up @@ -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,
Expand Down
36 changes: 23 additions & 13 deletions src/lib/rag-cache.ts
Original file line number Diff line number Diff line change
Expand Up @@ -144,6 +144,7 @@ export async function getCachedAnswer(
"query" | "documentId" | "documentIds" | "ownerId" | "accessScope" | "skipCache" | "queryMode"
>,
startedAt: number,
options?: { indexingVersionAtRequestStart?: string | null },
): Promise<RagAnswer | null> {
if (!answerCacheAllowedForOwner(args.ownerId) || args.skipCache) return null;
if (env.RAG_ANSWER_CACHE_TTL_MS <= 0 || env.RAG_ANSWER_CACHE_SIZE <= 0) return null;
Expand All @@ -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;
Expand All @@ -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,
Expand All @@ -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) {
Expand Down Expand Up @@ -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,
};
}

Expand All @@ -267,10 +266,15 @@ export function isSearchCacheEnabled(args: Pick<SearchChunksArgs, "skipCache">):
return !args.skipCache && env.RAG_SEARCH_CACHE_TTL_MS > 0 && env.RAG_SEARCH_CACHE_SIZE > 0;
}

export function isSearchCacheLookupEnabled(args: Pick<SearchChunksArgs, "skipCache">): 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;
Expand All @@ -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);
Expand Down Expand Up @@ -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";
Expand Down Expand Up @@ -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 }
Expand All @@ -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);
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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;
Expand All @@ -595,6 +602,7 @@ function setSharedCachedSearch(
args,
{ results: cloneSearchResults(results), telemetry },
env.RAG_SEARCH_CACHE_TTL_MS,
indexingVersion,
retrievalPlanCacheQuery(args, telemetry.query_class, queryVariants),
);
}
Expand All @@ -605,13 +613,15 @@ 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(
"answer",
args,
{ answer: cloneAnswer(answer) },
env.RAG_ANSWER_CACHE_TTL_MS,
indexingVersion,
sharedAnswerNormalizedQuery(args),
);
}
Expand Down
6 changes: 6 additions & 0 deletions src/lib/rag-contracts.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string>;
};
};

export type SearchTelemetry = {
search_cache_hit: boolean;
search_total_latency_ms?: number;
retrieval_phase_latencies_ms?: Record<string, number>;
shared_cache_hit?: boolean;
shared_cache_status?: "hit" | "miss";
shared_cache_miss_reason?: string | null;
Expand Down
Loading