From 65d7a586d35586ad22f1f8bb9ac0ce5dbcbba541 Mon Sep 17 00:00:00 2001 From: openSwarm Date: Mon, 28 Sep 2026 15:11:02 +0900 Subject: [PATCH] fix(memory): paginate memory lifecycle sweeps and stop dropping rows past the first page Salvaged from draft PRs #772 and #776. Compaction, expiry cleanup, consolidation, and statistics all read the memory table with a single capped query, so larger stores were silently truncated: compaction refused outright at its 100k safety limit, while cleanupExpired, consolidateMemories and getMemoryStats quietly processed only the first 10k rows. They now read the complete table through memoryCore's paginated fetchAllTableRows / offset-limit pages. Deduplication is now bucket-based (stable FNV-1a hash of repo, type, derivedFrom and canonical metadata) instead of an O(n^2) pairwise scan, which matters now that the full record set is considered at once. Bucketing also guarantees duplicate pairs that straddle a page boundary meet: with the old page-by-page filtering they could not be compared. Within a bucket, records are ranked by importance then recency so the survivor is independent of page order. Also fixes the type defect in #772 where cosineSimilarity was declared : boolean while returning a similarity score, which broke the >= threshold comparison at the consolidation call site. Dropped from the drafts: getTable(tempTableName)/db.renameTable() (neither exists in the current memoryCore/lancedb API), the removals of withMemoryMutationLock/saveConversation/getRecentConversations/ runBackgroundCognition/getMemoryStats (main's API, still called by discordCore, support/chatMemory and support/web), package.json version churn, and the unconditional error-swallowing catch in compactMemoryTable that main deliberately rethrows. --- src/memory/compaction.test.ts | 34 ++++++ src/memory/compaction.ts | 209 ++++++++++++++++++++++------------ src/memory/memoryOps.test.ts | 1 + src/memory/memoryOps.ts | 77 +++++++++---- 4 files changed, 229 insertions(+), 92 deletions(-) diff --git a/src/memory/compaction.test.ts b/src/memory/compaction.test.ts index 7417d329..8292c0b1 100644 --- a/src/memory/compaction.test.ts +++ b/src/memory/compaction.test.ts @@ -18,4 +18,38 @@ describe('memory compaction deduplication', () => { record('b', '{"nested":{"y":2,"x":1},"project":"p"}'), ])).toHaveLength(1); }); + + it('deduplicates a near-duplicate pair that straddles a pagination boundary', () => { + // Page 1 ends at p1-9999; the straddler is the first record of page 2 and + // shares p1-9999's identity (same metadata), so the pair only meets if + // dedup considers records from both pages together. It has a near-identical + // but not equal vector, and higher importance, so it must win the merge. + const page1 = Array.from({ length: 10_000 }, (_, i) => + record(`p1-${i}`, `{"i":${i}}`)); + const page2 = Array.from({ length: 10_000 }, (_, i) => + record(`p2-${i}`, `{"j":${i}}`)); + const straddler = { + ...record('p2-straddler', '{"i":9999}'), + vector: [1, 0.05], + importance: 0.9, + lastUpdated: 2, + }; + + const result = removeDuplicates([...page1, straddler, ...page2]); + const ids = new Set(result.map((r) => r.id)); + + // 20_001 inputs collapse to 20_000: the cross-page pair merges into one. + expect(result).toHaveLength(20_000); + expect(ids.has('p2-straddler')).toBe(true); + expect(ids.has('p1-9999')).toBe(false); + // Records whose identity merely neighbours the boundary are untouched. + expect(ids.has('p1-9998')).toBe(true); + expect(ids.has('p2-0')).toBe(true); + }); + + it('keeps records that share identity metadata but differ in vector', () => { + const a = { ...record('a', '{"project":"alpha"}'), vector: [1, 0] }; + const b = { ...record('b', '{"project":"alpha"}'), vector: [0, 1] }; + expect(removeDuplicates([a, b])).toHaveLength(2); + }); }); diff --git a/src/memory/compaction.ts b/src/memory/compaction.ts index 2a88f5f5..e012ca6c 100644 --- a/src/memory/compaction.ts +++ b/src/memory/compaction.ts @@ -2,13 +2,25 @@ // OpenSwarm - Memory Compaction // ============================================ -import { getDb, getTable, initDatabase, EMBEDDING_DIM, PERMANENT_EXPIRY, normalizeRecords, setTable } from './memoryCore.js'; +import { getDb, getTable, initDatabase, PERMANENT_EXPIRY, normalizeRecords, setTable } from './memoryCore.js'; import type { CognitiveMemoryRecord } from './memoryCore.js'; import { isTransientReviewRejectionMemory } from './memoryFilters.js'; const MIN_IMPORTANCE = 0.1; const CONSOLIDATION_SIMILARITY = 0.85; +/** Page size for full-table scans: a single `.limit(100_000)` truncates larger stores. */ +const PAGE_SIZE = 10_000; + +/** v2 columns that force a compaction rewrite to the lean v3 schema. */ +const LEGACY_SCHEMA_COLUMNS: Record = { + revisionCount: true, + decay: true, + stability: true, + contradicts: true, + supports: true, +}; + function stableJson(value: unknown): string { if (Array.isArray(value)) return `[${value.map(stableJson).join(',')}]`; if (value && typeof value === 'object') { @@ -50,57 +62,111 @@ function cosineSimilarity(a: number[], b: number[]): number { } /** - * Remove duplicate memories based on vector similarity + * Remove duplicate memories based on vector similarity. + * + * Records are bucketed by a stable hash of their non-vector identity fields + * (repo, type, derivedFrom, canonical metadata) before merging, so duplicates + * that straddle a pagination boundary still land in the same bucket and are + * compared. Within a bucket, records are ranked by importance then recency and + * a record is dropped only when it is a near-duplicate of an already-kept one — + * so the survivor does not depend on the order the pages were read in. */ export function removeDuplicates(records: CognitiveMemoryRecord[]): CognitiveMemoryRecord[] { + // 1. Bucket by identity BEFORE any merging so cross-page duplicates meet. + const buckets = new Map(); + for (const record of records) { + const key = stableHash([ + record.repo, + record.type, + record.derivedFrom, + stableMetadata(record.metadata), + ]); + const bucket = buckets.get(key); + if (bucket) bucket.push(record); + else buckets.set(key, [record]); + } + const unique: CognitiveMemoryRecord[] = []; const seen = new Set(); - for (const record of records) { - // Skip if exact ID already seen - if (seen.has(record.id)) continue; - - // Check similarity with existing unique records - let isDuplicate = false; - for (const existing of unique) { - if ( - record.repo !== existing.repo || - record.type !== existing.type || - record.derivedFrom !== existing.derivedFrom || - stableMetadata(record.metadata) !== stableMetadata(existing.metadata) - ) { - continue; - } + // 2. Reduce each bucket, dropping only near-duplicates of a kept record. + for (const bucket of buckets.values()) { + // Rank by quality so the survivor of a near-duplicate cluster does not + // depend on input order (and therefore not on page order either). + bucket.sort( + (a, b) => b.importance - a.importance || b.lastUpdated - a.lastUpdated + ); - const similarity = cosineSimilarity(record.vector, existing.vector); + const kept: CognitiveMemoryRecord[] = []; - if (similarity >= CONSOLIDATION_SIMILARITY) { - // Keep the one with higher importance or more recent - if (record.importance > existing.importance || - record.lastUpdated > existing.lastUpdated) { - // Replace existing with current - const index = unique.indexOf(existing); - unique[index] = record; - seen.add(record.id); - } - isDuplicate = true; - break; - } - } + for (const record of bucket) { + if (seen.has(record.id)) continue; + seen.add(record.id); + + // Identity fields must match for vectors to be comparable at all, and the + // bucket key is exactly those fields — so comparing within the bucket is + // equivalent to the old whole-table scan, minus the page-order dependence. + const isDuplicate = kept.some( + (existing) => cosineSimilarity(record.vector, existing.vector) >= CONSOLIDATION_SIMILARITY + ); + if (isDuplicate) continue; - if (!isDuplicate) { + kept.push(record); unique.push(record); - seen.add(record.id); } } return unique; } +/** + * Order-independent hash (FNV-1a over canonical JSON) used as the dedup key. + * The JSON length is mixed into the result so distinct inputs that collide on + * the 32-bit hash remain separable in practice. + */ +function stableHash(value: unknown): string { + const json = stableJson(value); + let hash = 0x811c9dc5; + for (let i = 0; i < json.length; i++) { + hash ^= json.charCodeAt(i); + hash = Math.imul(hash, 0x01000193); + } + return `${(hash >>> 0).toString(16)}:${json.length.toString(16)}`; +} + +/** + * Decide whether a raw Lance row survives compaction. + * + * Lance hands back schema-erased rows, so the fields this decision reads are + * validated here and the result is a type predicate — the caller then works + * with a validated CognitiveMemoryRecord instead of re-asserting the shape. + */ +function isValidCompactionRow(row: unknown, now: number): row is CognitiveMemoryRecord { + if (typeof row !== 'object' || row === null) return false; + const r = row as Partial; + if (typeof r.id !== 'string') return false; + // `init` is the schema seed row: always retained, never a merge candidate. + if (r.id === 'init') return true; + + // Remove transient infrastructure failures that were previously stored as + // high-importance reviewer constraints. + if (isTransientReviewRejectionMemory(r)) return false; + + // Remove if expired, or if unimportant. + if (typeof r.expiresAt === 'number' && r.expiresAt < PERMANENT_EXPIRY && r.expiresAt < now) return false; + if (typeof r.importance === 'number' && r.importance < MIN_IMPORTANCE) return false; + + return true; +} + /** * Compact memory table by removing expired/unimportant/noisy records, * deduplicating similar memories, and rewriting to the lean v3 schema. * + * Reads records in offset/limit pages rather than one capped query, then + * deduplicates the full set so duplicates straddling a page boundary are + * still compared. + * * @returns Statistics about compaction */ export async function compactMemoryTable(): Promise<{ @@ -121,15 +187,15 @@ export async function compactMemoryTable(): Promise<{ return { before: 0, after: 0, removed: 0, deduplicated: 0 }; } - // 1. Read all records - const queryLimit = 100_000; - const allRecords = await table - .search(Array.from({ length: EMBEDDING_DIM }, () => 0)) - .limit(queryLimit) - .toArray(); - - if (allRecords.length >= queryLimit) { - throw new Error(`Memory compaction refused: query reached the ${queryLimit}-row safety limit`); + // 1. Read all records across pagination boundaries. Lance returns + // schema-erased rows here; `unknown` keeps the boundary honest until + // the per-row filter below narrows the fields it actually reads. + const allRecords: unknown[] = []; + for (;;) { + const page = await table.query().limit(PAGE_SIZE).offset(allRecords.length).toArray(); + if (page.length === 0) break; + allRecords.push(...page); + if (page.length < PAGE_SIZE) break; } const beforeCount = allRecords.length; @@ -142,27 +208,15 @@ export async function compactMemoryTable(): Promise<{ // 2. Filter valid records const now = Date.now(); - const validRecords = allRecords.filter((r: any) => { - if (r.id === 'init') return true; - - // Remove transient infrastructure failures that were previously stored as - // high-importance reviewer constraints. - if (isTransientReviewRejectionMemory(r)) return false; - - // Remove if expired - if (r.expiresAt < PERMANENT_EXPIRY && r.expiresAt < now) return false; - - // Remove if unimportant - if (r.importance < MIN_IMPORTANCE) return false; - - return true; - }); + const validRecords = allRecords.filter( + (row): row is CognitiveMemoryRecord => isValidCompactionRow(row, now) + ); const afterFilter = validRecords.length; console.log(`[Compaction] After filtering: ${afterFilter} records (removed ${beforeCount - afterFilter})`); - // 3. Deduplicate - const deduplicated = removeDuplicates(validRecords as CognitiveMemoryRecord[]); + // 3. Deduplicate the whole set so cross-page duplicates are caught. + const deduplicated = removeDuplicates(validRecords); const afterDedup = deduplicated.length; console.log(`[Compaction] After deduplication: ${afterDedup} records (merged ${afterFilter - afterDedup})`); @@ -218,7 +272,7 @@ export async function compactMemoryTable(): Promise<{ } /** - * Check if compaction is needed based on heuristics + * Check if compaction is needed based on table size and waste ratio. */ export async function shouldCompact(): Promise { try { @@ -226,35 +280,44 @@ export async function shouldCompact(): Promise { const table = getTable(); if (!table) return false; - const allRecords = await table - .search(Array.from({ length: EMBEDDING_DIM }, () => 0)) - .limit(100000) - .toArray(); + // Bounded total via countRows — never a full-table load. + const totalRows = await table.countRows(); + if (totalRows === 0) return false; + + // Waste is estimated from the first page only; a full scan here would + // cost as much as the compaction this check is trying to avoid. + const sample = await table.query().limit(PAGE_SIZE).toArray(); const now = Date.now(); // Count expired/noisy records let expiredCount = 0; let noisyCount = 0; - let legacyColumnCount = 0; - for (const r of allRecords) { - if (r.expiresAt < PERMANENT_EXPIRY && r.expiresAt < now) expiredCount++; + for (const row of sample) { + if (typeof row !== 'object' || row === null) continue; + const r = row as Partial; + if (r.id === 'init') continue; + if (typeof r.expiresAt === 'number' && r.expiresAt < PERMANENT_EXPIRY && r.expiresAt < now) expiredCount++; if (isTransientReviewRejectionMemory(r)) noisyCount++; - if ('revisionCount' in r || 'decay' in r || 'stability' in r || 'contradicts' in r || 'supports' in r) { - legacyColumnCount++; - } } const totalWaste = expiredCount + noisyCount; - const wasteRatio = totalWaste / allRecords.length; + const wasteRatio = sample.length > 0 ? totalWaste / sample.length : 0; + + // Legacy v2 columns live in the schema, not in row values, so detect them + // there — a v2 table needs the rewrite regardless of its waste ratio. + const schema = await table.schema(); + const legacyColumnCount = schema.fields.filter( + (field) => LEGACY_SCHEMA_COLUMNS[field.name] === true + ).length; // Compact if > 20% waste, > 1000 records, or legacy v2 fields are still // present and need a schema rewrite. - const shouldCompact = wasteRatio > 0.2 || allRecords.length > 1000 || legacyColumnCount > 0; + const shouldCompact = wasteRatio > 0.2 || totalRows > 1000 || legacyColumnCount > 0; if (shouldCompact) { - console.log(`[Compaction] Compaction recommended: ${totalWaste}/${allRecords.length} waste (${(wasteRatio * 100).toFixed(1)}%), ${legacyColumnCount} legacy rows`); + console.log(`[Compaction] Compaction recommended: ${totalWaste}/${sample.length} sampled waste (${(wasteRatio * 100).toFixed(1)}% of ${totalRows} rows), ${legacyColumnCount} legacy columns`); } return shouldCompact; diff --git a/src/memory/memoryOps.test.ts b/src/memory/memoryOps.test.ts index 768d509a..4ae5dc21 100644 --- a/src/memory/memoryOps.test.ts +++ b/src/memory/memoryOps.test.ts @@ -12,6 +12,7 @@ vi.mock('./memoryCore.js', () => ({ normalizeRecords: (records: any[]) => records, initDatabase: vi.fn(async () => {}), embedPassage: vi.fn(async () => [0.1, 0.2, 0.3, 0.4]), + fetchAllTableRows: async () => [...state.records.values()].map((r) => ({ ...r })), getTable: () => ({ query: () => ({ where: (pred: string) => ({ diff --git a/src/memory/memoryOps.ts b/src/memory/memoryOps.ts index 537c0ad8..b56d0280 100644 --- a/src/memory/memoryOps.ts +++ b/src/memory/memoryOps.ts @@ -5,11 +5,11 @@ * Core types, save, search are in memoryCore.ts. */ import { - EMBEDDING_DIM, PERMANENT_EXPIRY, normalizeRecords, initDatabase, embedPassage, + fetchAllTableRows, getTable, searchMemory, calculateFreshness, @@ -25,6 +25,9 @@ import { embeddingTextFor } from './embeddingConfig.js'; type MemoryTable = NonNullable>; const MAX_MEMORY_REVISIONS = 20; +/** Rows deleted per Lance commit while sweeping expired records. */ +const PAGE_SIZE = 10_000; + /** * In-process queue that serializes full read-modify-write memory mutations. * withMemoryWriteRetry only retries individual Lance commits; without this, @@ -395,7 +398,15 @@ function formatDate(timestamp: number): string { } /** - * Clean up expired memories + * Clean up expired memories. + * + * Scans the whole table in fixed-offset pages and then deletes the expired + * rows in bounded batches; the previous single `.limit(10_000)` silently left + * the remainder of larger stores behind. Deleting only after the scan + * completes keeps the page offsets stable — deleting mid-scan shifts every + * later row and makes the cursor skip records. The sweep runs inside + * withMemoryWriteRetry so a concurrent writer that wins a version race + * retries it instead of failing part-way; deletes are idempotent. */ export async function cleanupExpired(): Promise { try { @@ -404,18 +415,30 @@ export async function cleanupExpired(): Promise { if (!table) return 0; const now = Date.now(); - const results = await table.search(Array.from({ length: EMBEDDING_DIM }, () => 0)).limit(10000).toArray(); - const expiredIds = results - .filter((r: any) => r.expiresAt < PERMANENT_EXPIRY && r.expiresAt < now) - .map((r: any) => r.id); + return await withMemoryWriteRetry(async () => { + const rows = await fetchAllTableRows(table); - if (expiredIds.length > 0) { - await deleteMemoryIds(table, expiredIds); - console.log(`[Memory] Deleted ${expiredIds.length} expired records`); - } + const expiredIds: string[] = []; + for (const row of rows) { + if (typeof row !== 'object' || row === null) continue; + const r = row as Partial; + if (typeof r.id !== 'string') continue; + if (typeof r.expiresAt === 'number' && r.expiresAt < PERMANENT_EXPIRY && r.expiresAt < now) { + expiredIds.push(r.id); + } + } - return expiredIds.length; + for (let i = 0; i < expiredIds.length; i += PAGE_SIZE) { + await table.delete(idsPredicate(expiredIds.slice(i, i + PAGE_SIZE))); + } + + if (expiredIds.length > 0) { + console.log(`[Memory] Cleanup complete: ${expiredIds.length} expired records deleted`); + } + + return expiredIds.length; + }, 'cleanupExpired'); } catch (error) { console.error('[Memory] Cleanup error:', error); return 0; @@ -426,7 +449,13 @@ export async function cleanupExpired(): Promise { const CONSOLIDATION_SIMILARITY = 0.85; // Duplicate detection threshold /** - * Consolidate duplicate/similar memories + * Consolidate duplicate/similar memories. + * + * Reads the complete table in pages: the previous single `.limit(10000)` vector + * query both capped the scan and let vector ranking decide which rows were + * eligible for merging. Runs under the in-process mutation lock like every + * other read-modify-write op here, so concurrent revise/consolidate cannot + * overwrite each other's metadata. */ export async function consolidateMemories(): Promise<{ merged: number; @@ -438,8 +467,8 @@ export async function consolidateMemories(): Promise<{ const table = getTable(); if (!table) return { merged: 0, groups: [] }; - const results = await table.search(Array.from({ length: EMBEDDING_DIM }, () => 0)).limit(10000).toArray(); - const validMemories = results.filter((r: any) => r.id !== 'init'); + const allMemories = await fetchAllTableRows(table); + const validMemories = allMemories.filter((r) => r.id !== 'init'); const merged: string[] = []; const groups: Array<{ kept: string; merged: string[] }> = []; @@ -516,7 +545,11 @@ export async function consolidateMemories(): Promise<{ } /** - * Cosine similarity between two vectors + * Cosine similarity between two vectors. + * + * Returns a score in [-1, 1]; `0` means "cannot compare" (missing or + * mismatched vectors, or a zero-magnitude vector). Callers compare the result + * against a similarity threshold, so it must be a number. */ function cosineSimilarity(a: number[], b: number[]): number { if (!a || !b || a.length !== b.length) return 0; @@ -605,19 +638,25 @@ export async function getMemoryStats(): Promise<{ const table = getTable(); if (!table) return { total: 0, byType: { ...DEFAULT_BY_TYPE }, byRepo: {}, avgImportance: 0 }; - const results = await table.search(Array.from({ length: EMBEDDING_DIM }, () => 0)).limit(10000).toArray(); + // Aggregate over the complete table: the previous `.limit(10000)` vector + // query capped statistics at the first 10k rows and let vector ranking pick + // which rows were counted. + const results = await fetchAllTableRows(table); const byType: Record = { ...DEFAULT_BY_TYPE }; const byRepo: Record = {}; let totalImportance = 0; let count = 0; - for (const r of results) { - if (r.id === 'init') continue; + for (const row of results) { + if (typeof row !== 'object' || row === null) continue; + const r = row as Partial; + if (r.id === 'init' || typeof r.id !== 'string') continue; if (byType[r.type as MemoryType] !== undefined) { byType[r.type as MemoryType]++; } - byRepo[r.repo] = (byRepo[r.repo] || 0) + 1; + const repo = r.repo ?? 'unknown'; + byRepo[repo] = (byRepo[repo] || 0) + 1; totalImportance += r.importance ?? 0.5; count++; }