import crypto from 'node:crypto'; import { fingerprintContent } from './context-budget.mjs'; import { pgvectorMemoryBackendInternals } from './memory-v2-pgvector.mjs'; import { isTemporalRecallQuery } from './temporal-recall-service/keyword-rules.mjs'; export const RECALL_FUSION_MODES = Object.freeze(['off', 'shadow', 'active']); export const DEFAULT_RECALL_FUSION_RRF_K = 60; function normalizeMode(value, fallback = 'off') { const raw = String(value ?? fallback).trim().toLowerCase(); return RECALL_FUSION_MODES.includes(raw) ? raw : fallback; } export function resolveRecallFusionMode(env = process.env) { return normalizeMode(env.MEMIND_RECALL_FUSION_MODE, 'off'); } export function resolveRecallFusionRrfK(env = process.env) { const raw = Number(env.MEMIND_RECALL_FUSION_RRF_K ?? DEFAULT_RECALL_FUSION_RRF_K); return Number.isFinite(raw) && raw > 0 ? Math.floor(raw) : DEFAULT_RECALL_FUSION_RRF_K; } export function resolveRecallFusionLimit(env = process.env, fallback = 3) { const raw = Number(env.MEMIND_RECALL_FUSION_LIMIT ?? fallback); return Number.isFinite(raw) && raw > 0 ? Math.floor(raw) : fallback; } function memoryItemText(item) { return String(item?.text ?? item?.content ?? item?.summary ?? '').trim(); } function memoryItemLabel(item) { return String(item?.label ?? item?.title ?? '').trim() || null; } export function parseRecallQuery(query, { now = new Date() } = {}) { const text = String(query ?? '').trim(); return { text, temporal: isTemporalRecallQuery(text), keywordTerms: pgvectorMemoryBackendInternals.extractKeywordTerms(text), now, }; } export function normalizeRecallCandidate(item, { source, rank = 0, query = '' } = {}) { if (!item || typeof item !== 'object') return null; const text = memoryItemText(item); if (!text) return null; const label = memoryItemLabel(item); const id = String( item?.id ?? item?.memoryId ?? item?.sessionId ?? item?.event_id ?? `${source}:${label ?? text.slice(0, 32)}`, ).trim(); const lexical = query ? pgvectorMemoryBackendInternals.lexicalQueryCoverage(query, text) : 0; return { id, source, label, text, rank, lexical, fingerprint: fingerprintContent(text), raw: item, }; } export function reciprocalRankFusion(lists, { k = DEFAULT_RECALL_FUSION_RRF_K } = {}) { const scores = new Map(); const meta = new Map(); for (const list of Array.isArray(lists) ? lists : []) { const source = String(list?.source ?? 'unknown'); const candidates = Array.isArray(list?.candidates) ? list.candidates : []; candidates.forEach((candidate, index) => { if (!candidate?.id) return; const rrf = 1 / (k + index + 1); const prev = scores.get(candidate.id) ?? 0; scores.set(candidate.id, prev + rrf); const existing = meta.get(candidate.id); if (!existing) { meta.set(candidate.id, { ...candidate, sources: [source], ranks: { [source]: index + 1 }, rrfScore: rrf, }); return; } existing.sources = [...new Set([...existing.sources, source])]; existing.ranks[source] = index + 1; existing.rrfScore = scores.get(candidate.id); if ((candidate.lexical ?? 0) > (existing.lexical ?? 0)) { existing.lexical = candidate.lexical; existing.label = candidate.label ?? existing.label; existing.text = candidate.text ?? existing.text; existing.raw = candidate.raw ?? existing.raw; } }); } return [...meta.values()] .sort((left, right) => { const scoreDelta = (right.rrfScore ?? 0) - (left.rrfScore ?? 0); if (scoreDelta !== 0) return scoreDelta; const lexicalDelta = (right.lexical ?? 0) - (left.lexical ?? 0); if (lexicalDelta !== 0) return lexicalDelta; return String(left.id).localeCompare(String(right.id)); }); } export function fuseRecallCandidates({ query = '', personalMemories = [], episodicMemories = [], temporalItems = [], limit = 3, env = process.env, } = {}) { const mode = resolveRecallFusionMode(env); const parsedQuery = parseRecallQuery(query); const lists = [ { source: 'personal', candidates: (Array.isArray(personalMemories) ? personalMemories : []) .map((item, index) => normalizeRecallCandidate(item, { source: 'personal', rank: index, query })) .filter(Boolean), }, { source: 'episodic', candidates: (Array.isArray(episodicMemories) ? episodicMemories : []) .map((item, index) => normalizeRecallCandidate(item, { source: 'episodic', rank: index, query })) .filter(Boolean), }, { source: 'temporal', candidates: (Array.isArray(temporalItems) ? temporalItems : []) .map((item, index) => normalizeRecallCandidate(item, { source: 'temporal', rank: index, query, })) .filter(Boolean), }, ]; const fused = reciprocalRankFusion(lists, { k: resolveRecallFusionRrfK(env) }); const seenFingerprints = new Set(); const memories = []; const duplicateItems = []; for (const item of fused) { if (item.fingerprint && seenFingerprints.has(item.fingerprint)) { duplicateItems.push({ id: item.id, sources: item.sources, reason: 'duplicate_fingerprint', }); continue; } if (item.fingerprint) seenFingerprints.add(item.fingerprint); memories.push(item.raw ?? { id: item.id, label: item.label, text: item.text, source: item.sources?.[0] ?? item.source, }); if (memories.length >= limit) break; } return { mode, query: parsedQuery, inputCounts: { personal: lists[0].candidates.length, episodic: lists[1].candidates.length, temporal: lists[2].candidates.length, }, fusedCount: fused.length, keptCount: memories.length, duplicateCount: duplicateItems.length, duplicateItems, memories, topScores: fused.slice(0, Math.min(5, fused.length)).map((item) => ({ id: item.id, sources: item.sources, rrfScore: Number(item.rrfScore?.toFixed?.(4) ?? item.rrfScore), lexical: item.lexical, })), }; } export function applyRecallFusionToMemoryMerge({ personalMemories = [], episodicMemories = [], temporalItems = [], query = '', limit = 3, legacyMerge, env = process.env, } = {}) { const mode = resolveRecallFusionMode(env); const fusion = fuseRecallCandidates({ query, personalMemories, episodicMemories, temporalItems, limit, env, }); if (mode === 'off') { return { memories: legacyMerge?.({ personalMemories, episodicMemories, query, limit, }) ?? [], fusion: null, applied: false, }; } if (mode === 'shadow') { return { memories: legacyMerge?.({ personalMemories, episodicMemories, query, limit, }) ?? fusion.memories, fusion, applied: false, }; } return { memories: fusion.memories, fusion, applied: true, }; } export function buildRecallFusionResolvedEvent(fusion) { if (!fusion) return null; return { mode: fusion.mode, temporalQuery: Boolean(fusion.query?.temporal), inputCounts: fusion.inputCounts, fusedCount: fusion.fusedCount, keptCount: fusion.keptCount, duplicateCount: fusion.duplicateCount, topScores: fusion.topScores, }; } export function hashRecallFusionPlan(fusion) { if (!fusion) return null; return crypto.createHash('sha256') .update(JSON.stringify({ keptCount: fusion.keptCount, topScores: fusion.topScores, inputCounts: fusion.inputCounts, })) .digest('hex') .slice(0, 16); }