/** * Memory store - CRUD operations with prepared statements and caching */ import { getDbInstance } from "../db/core"; import { upsertSemanticMemoryPoint, deleteSemanticMemoryPoint } from "./qdrant"; import { Memory, MemoryType } from "./types"; import { logger } from "../../../open-sse/utils/logger.ts"; import { sanitizeErrorMessage } from "../../../open-sse/utils/error.ts"; import { resolveEmbeddingSource, embed } from "./embedding"; import { getVectorStore } from "./vectorStore"; import { getMemorySettings } from "./settings"; import { markMemoryNeedsReindex } from "@/lib/localDb"; const log = logger("MEMORY_STORE"); interface CacheEntry { value: T; timestamp: number; } interface MemoryRow { id: string; api_key_id: string; session_id: string | null; type: MemoryType; key: string | null; content: string; metadata: string | null; created_at: string; updated_at: string; expires_at: string | null; } // Memory cache configuration const MEMORY_CACHE_TTL = 60_000; // 1 minute const MEMORY_MAX_CACHE_SIZE = 500; // Cache for recently accessed memories const _memoryCache = new Map>(); // Helper function to safely parse JSON strings function parseJSON(value: unknown): Record { if (!value || typeof value !== "string" || value.trim() === "") { return {}; } try { const parsed = JSON.parse(value); return typeof parsed === "object" && parsed !== null ? parsed : {}; } catch { return {}; } } function invalidateMemoryCache(key: string) { _memoryCache.delete(key); } function evictIfNeeded(cache: Map) { if (cache.size > MEMORY_MAX_CACHE_SIZE) { // Remove oldest entries first const keysArray = Array.from(cache.keys()); const entriesToRemove = Math.floor(cache.size * 0.2); for (let i = 0; i < entriesToRemove; i++) { cache.delete(keysArray[i]); } } } function rowToMemory(row: MemoryRow): Memory { return { id: String(row.id), apiKeyId: String(row.api_key_id), sessionId: typeof row.session_id === "string" ? row.session_id : "", type: row.type as MemoryType, key: typeof row.key === "string" ? row.key : "", content: String(row.content), metadata: parseJSON(row.metadata), createdAt: new Date(String(row.created_at)), updatedAt: new Date(String(row.updated_at)), expiresAt: row.expires_at ? new Date(String(row.expires_at)) : null, }; } /** * Find existing memory by apiKeyId and key (for UPSERT logic) */ function findExistingMemory( db: ReturnType, apiKeyId: string, key: string ): MemoryRow | undefined { if (!key) return undefined; const stmt = db.prepare( "SELECT * FROM memories WHERE api_key_id = ? AND key = ? ORDER BY created_at DESC LIMIT 1" ); return stmt.get(apiKeyId, key) as MemoryRow | undefined; } /** * Fire-and-forget: generate embedding for a memory and upsert into sqlite-vec. * Errors are logged but never thrown — this must never block the SQLite write. */ /** * Best-effort: try to mark a memory needs_reindex. Swallows errors so that DB-closed * states (e.g. test teardown after the parent promise resolved) never escape as * unhandledRejection. Producing this side-effect is opportunistic by design. */ function safeMarkNeedsReindex(id: string, needs: boolean): void { try { markMemoryNeedsReindex(id, needs); } catch { // intentional swallow — DB may be closed (test teardown) or schema not yet ready } } function scheduleVectorUpsert(id: string, content: string): void { setImmediate(async () => { try { const settings = await getMemorySettings(); const resolution = resolveEmbeddingSource(settings); if (!resolution.source) return; const embeddingResult = await embed(content, settings); if (!("vector" in embeddingResult)) { log.warn("memory.vec.embed.fail", { id, reason: embeddingResult.reason, message: sanitizeErrorMessage(embeddingResult.message), }); safeMarkNeedsReindex(id, true); return; } const vec = getVectorStore(); if (!vec) { safeMarkNeedsReindex(id, true); return; } await vec.ensureReady(resolution); await vec.upsertVector(id, embeddingResult.vector); safeMarkNeedsReindex(id, false); } catch (err: unknown) { log.warn("memory.vec.upsert.fail", { id, error: sanitizeErrorMessage(err instanceof Error ? err.message : String(err)), }); safeMarkNeedsReindex(id, true); } }); } /** * Create a new memory entry (UPSERT: updates existing if same apiKeyId + key) */ export async function createMemory( memory: Omit ): Promise { const db = getDbInstance(); const now = new Date().toISOString(); // Check for existing memory with same apiKeyId + key (UPSERT logic) const existing = memory.key ? findExistingMemory(db, memory.apiKeyId, memory.key) : undefined; if (existing) { // UPDATE existing record const updatedMetadata = { ...parseJSON(existing.metadata), ...memory.metadata }; const stmt = db.prepare( "UPDATE memories SET content = ?, metadata = ?, updated_at = ?, session_id = ?, type = ?, expires_at = ? WHERE id = ?" ); stmt.run( memory.content, JSON.stringify(updatedMetadata), now, memory.sessionId, memory.type, memory.expiresAt ?? null, existing.id ); const updatedMemory: Memory = { id: String(existing.id), apiKeyId: memory.apiKeyId, sessionId: memory.sessionId, type: memory.type, key: memory.key, content: memory.content, metadata: updatedMetadata, createdAt: new Date(String(existing.created_at)), updatedAt: new Date(now), expiresAt: memory.expiresAt ?? null, }; // Invalidate and update cache invalidateMemoryCache(existing.id); evictIfNeeded(_memoryCache); _memoryCache.set(existing.id, { value: updatedMemory, timestamp: Date.now() }); log.info("memory.updated", { apiKeyId: memory.apiKeyId, type: memory.type, id: existing.id, key: memory.key, }); // Best-effort vector upsert (fire-and-forget — content changed so regenerate) scheduleVectorUpsert(String(existing.id), memory.content); // Best-effort re-sync to Qdrant after update upsertSemanticMemoryPoint({ id: String(existing.id), apiKeyId: memory.apiKeyId || "", sessionId: memory.sessionId || "", key: memory.key || "", content: memory.content, metadata: updatedMetadata || {}, createdAt: String(existing.created_at), expiresAt: memory.expiresAt ? memory.expiresAt.toISOString() : null, }) .then((r) => { if (r.ok) log.debug?.("qdrant.upsert.ok", { id: existing.id, latencyMs: r.latencyMs }); else if (r.error && r.error !== "not_configured") log.warn?.("qdrant.upsert.fail", { id: existing.id, error: r.error }); }) .catch((e) => log.warn?.("qdrant.upsert.error", { id: existing.id, error: String(e) })); return updatedMemory; } // INSERT new record if not exists const id = crypto.randomUUID(); const stmt = db.prepare( "INSERT INTO memories (id, api_key_id, session_id, type, key, content, metadata, created_at, updated_at, expires_at) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)" ); stmt.run( id, memory.apiKeyId, memory.sessionId, memory.type, memory.key, memory.content, JSON.stringify(memory.metadata ?? {}), now, now, memory.expiresAt?.toISOString() ?? null ); const createdMemory: Memory = { id, apiKeyId: memory.apiKeyId, sessionId: memory.sessionId, type: memory.type, key: memory.key, content: memory.content, metadata: memory.metadata, createdAt: new Date(now), updatedAt: new Date(now), expiresAt: memory.expiresAt ?? null, }; // Cache the newly created memory invalidateMemoryCache(id); evictIfNeeded(_memoryCache); _memoryCache.set(id, { value: createdMemory, timestamp: Date.now() }); log.info("memory.stored", { apiKeyId: memory.apiKeyId, type: memory.type, id }); // Best-effort vector upsert (fire-and-forget) scheduleVectorUpsert(id, memory.content); // Best-effort sync to semantic memory store (Qdrant). Failures do not block the SQLite write. upsertSemanticMemoryPoint({ id, apiKeyId: memory.apiKeyId || "", sessionId: memory.sessionId || "", key: memory.key || "", content: memory.content, metadata: memory.metadata || {}, createdAt: now, expiresAt: memory.expiresAt ? memory.expiresAt.toISOString() : null, }) .then((r) => { if (r.ok) log.debug?.("qdrant.upsert.ok", { id, latencyMs: r.latencyMs }); else if (r.error && r.error !== "not_configured") log.warn?.("qdrant.upsert.fail", { id, error: r.error }); }) .catch((e) => log.warn?.("qdrant.upsert.error", { id, error: String(e) })); return createdMemory; } /** * Get a memory by ID */ export async function getMemory(id: string): Promise { if (!id || typeof id !== "string") return null; // Check cache first const cached = _memoryCache.get(id); if (cached && Date.now() - cached.timestamp < MEMORY_CACHE_TTL) { return cached.value; } const db = getDbInstance(); const stmt = db.prepare("SELECT * FROM memories WHERE id = ?"); const row = stmt.get(id) as MemoryRow | undefined; if (!row) { // Cache negative result briefly to prevent repeated DB hits evictIfNeeded(_memoryCache); _memoryCache.set(id, { value: null, timestamp: Date.now() }); return null; } const memory = rowToMemory(row); // Cache the result evictIfNeeded(_memoryCache); _memoryCache.set(id, { value: memory, timestamp: Date.now() }); return memory; } /** * Update a memory entry */ export async function updateMemory( id: string, updates: Partial> ): Promise { if (!id || typeof id !== "string") return false; const db = getDbInstance(); const now = new Date().toISOString(); // Fetch current state to detect content/key change (needed for vector re-gen) const currentRow = db.prepare("SELECT content, key FROM memories WHERE id = ?").get(id) as | { content: string; key: string | null } | undefined; // Build dynamic update query const fields: string[] = []; const values: unknown[] = []; if (updates.type !== undefined) { fields.push("type = ?"); values.push(updates.type); } if (updates.key !== undefined) { fields.push("key = ?"); values.push(updates.key); } if (updates.content !== undefined) { fields.push("content = ?"); values.push(updates.content); } if (updates.metadata !== undefined) { fields.push("metadata = ?"); values.push(JSON.stringify(updates.metadata)); } if (updates.expiresAt !== undefined) { fields.push("expires_at = ?"); values.push(updates.expiresAt?.toISOString() ?? null); } // Always update the updatedAt timestamp fields.push("updated_at = ?"); values.push(now); values.push(id); // For WHERE clause const stmt = db.prepare(`UPDATE memories SET ${fields.join(", ")} WHERE id = ?`); const result = stmt.run(...values); if (result.changes === 0) { return false; } // Invalidate cache for this memory invalidateMemoryCache(id); // Regenerate vector if content or key changed (fire-and-forget) const contentChanged = updates.content !== undefined && updates.content !== currentRow?.content; const keyChanged = updates.key !== undefined && updates.key !== currentRow?.key; if (contentChanged || keyChanged) { const newContent = updates.content ?? currentRow?.content ?? ""; scheduleVectorUpsert(id, newContent); } return true; } /** * Delete a memory by ID. * D15 (bug #3): MUST call both vec.deleteVector AND deleteSemanticMemoryPoint * before the SQLite DELETE to keep all stores in sync. */ export async function deleteMemory(id: string): Promise { if (!id || typeof id !== "string") return false; // 1. Delete from sqlite-vec (best-effort — does not fail if vec not loaded) const vec = getVectorStore(); if (vec) { await vec.deleteVector(id).catch((e: unknown) => log.warn("memory.vec.delete.fail", { id, error: sanitizeErrorMessage(e instanceof Error ? e.message : String(e)), }) ); } // 2. Delete from Qdrant (best-effort — already existed before plan 21) await deleteSemanticMemoryPoint(id).catch((e: unknown) => log.warn("memory.qdrant.delete.fail", { id, error: sanitizeErrorMessage(e instanceof Error ? e.message : String(e)), }) ); // 3. Delete from SQLite const db = getDbInstance(); const stmt = db.prepare("DELETE FROM memories WHERE id = ?"); const result = stmt.run(id); if (result.changes === 0) { return false; } // Invalidate cache for this memory invalidateMemoryCache(id); log.info("memory.deleted", { id }); return true; } /** * List memories with optional filtering and pagination */ export async function listMemories(filters: { apiKeyId?: string; type?: MemoryType; sessionId?: string; query?: string; limit?: number; offset?: number; page?: number; }): Promise<{ data: Memory[]; total: number; byType: Record }> { const db = getDbInstance(); // Build dynamic query conditions const whereClauses: string[] = []; const whereParams: unknown[] = []; if (filters.apiKeyId) { whereClauses.push("api_key_id = ?"); whereParams.push(filters.apiKeyId); } if (filters.type) { whereClauses.push("type = ?"); whereParams.push(filters.type); } if (filters.sessionId) { whereClauses.push("session_id = ?"); whereParams.push(filters.sessionId); } if (typeof filters.query === "string" && filters.query.trim().length > 0) { const likeQuery = `%${filters.query.trim().toLowerCase()}%`; whereClauses.push("(LOWER(content) LIKE ? OR LOWER(key) LIKE ?)"); whereParams.push(likeQuery, likeQuery); } // Run COUNT query + byType aggregation in a single query let countQuery = "SELECT COUNT(*) as total FROM memories"; if (whereClauses.length > 0) { countQuery += " WHERE " + whereClauses.join(" AND "); } const countStmt = db.prepare(countQuery); const countRow = countStmt.get(...whereParams) as { total: number }; const total = countRow.total; // Build byType aggregation (counts ALL matching rows, not just the page) let byTypeQuery = "SELECT type, COUNT(*) as count FROM memories"; const byTypeParams: unknown[] = [...whereParams]; if (whereClauses.length > 0) { byTypeQuery += " WHERE " + whereClauses.join(" AND "); } byTypeQuery += " GROUP BY type"; const byTypeStmt = db.prepare(byTypeQuery); const byTypeRows = byTypeStmt.all(...byTypeParams) as { type: string; count: number }[]; const byType = Object.fromEntries(byTypeRows.map((r) => [r.type, r.count])) as Record< string, number >; // Calculate effective limit and offset const effectiveLimit = filters.limit ?? 50; const effectivePage = filters.page ?? 1; const effectiveOffset = filters.offset ?? (effectivePage - 1) * effectiveLimit; // Build SELECT query with pagination let query = "SELECT * FROM memories"; if (whereClauses.length > 0) { query += " WHERE " + whereClauses.join(" AND "); } // Add ordering and pagination query += " ORDER BY created_at DESC LIMIT ? OFFSET ?"; // Build params for SELECT query (WHERE params + pagination params) const params = [...whereParams, effectiveLimit, effectiveOffset]; const stmt = db.prepare(query); const rows = stmt.all(...params); return { data: (rows as MemoryRow[]).map(rowToMemory), total, byType, }; } /** * Total estimated tokens across stored memories (4 chars ≈ 1 token), computed in * SQL so we never load every memory's content into process memory. Scoped to a * single API key when `apiKeyId` is provided, otherwise counts all memories. */ export function getMemoryTokensUsed(apiKeyId?: string): number { const db = getDbInstance(); const stmt = db.prepare( "SELECT COALESCE(SUM((LENGTH(content) + 3) / 4), 0) as tokensUsed FROM memories" + (apiKeyId ? " WHERE api_key_id = ?" : "") ); const row = stmt.get(...(apiKeyId ? [apiKeyId] : [])) as { tokensUsed: number } | undefined; return row?.tokensUsed ?? 0; }