diff --git a/src/memory/compaction.test.ts b/src/memory/compaction.test.ts index 7417d329..c8942b22 100644 --- a/src/memory/compaction.test.ts +++ b/src/memory/compaction.test.ts @@ -18,4 +18,39 @@ describe('memory compaction deduplication', () => { record('b', '{"nested":{"y":2,"x":1},"project":"p"}'), ])).toHaveLength(1); }); + + it('deduplicates records that arrive on different pagination pages', () => { + // Simulate two pages of 10,000 records where the duplicate pair straddles + // the page boundary: 'dup-9999' ends page 1, 'dup-10000' starts page 2. + const page1 = Array.from({ length: 10_000 }, (_, i) => + record(`p1-${i}`, '{"project":"alpha"}')); + const page2 = Array.from({ length: 10_000 }, (_, i) => + record(`p2-${i}`, '{"project":"beta"}')); + + // Same repo/type/derivedFrom/metadata as p1-9999, near-identical vector. + const straddler = { + ...record('p2-straddler', '{"project":"alpha"}'), + vector: [1, 0.05], + importance: 0.9, + lastUpdated: 2, + }; + + const result = removeDuplicates([...page1, straddler, ...page2]); + + expect(result).toHaveLength(20_000); + const ids = new Set(result.map((r) => r.id)); + expect(ids.has('p1-9999')).toBe(true); + expect(ids.has('p2-straddler')).toBe(false); + // The higher-importance, more recent straddler wins the merge. + const kept = result.find((r) => r.id === 'p1-9999'); + expect(kept).toBeUndefined(); + // The winner is the straddler itself (kept under its own id). + expect(result.filter((r) => r.metadata === '{"project":"alpha"}')).toHaveLength(1); + }); + + it('keeps distinct records that merely share 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..4bc0987b 100644 --- a/src/memory/compaction.ts +++ b/src/memory/compaction.ts @@ -50,57 +50,82 @@ function cosineSimilarity(a: number[], b: number[]): number { } /** - * Remove duplicate memories based on vector similarity + * Remove duplicate memories based on vector similarity. + * + * Records are bucketed in a Map keyed by a stable hash of their non-vector + * fields (repo, type, derivedFrom, canonical metadata), so records read on + * different pagination pages still land in the same bucket and are compared. + * Within a bucket, near-identical vectors (cosine similarity >= + * CONSOLIDATION_SIMILARITY) collapse into the single best record (highest + * importance, then most recently updated). */ export function removeDuplicates(records: CognitiveMemoryRecord[]): CognitiveMemoryRecord[] { + // 1. Bucket every record by stable hash BEFORE any merging, so duplicates + // across page boundaries are guaranteed to 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) { + // 2. Reduce each bucket to its best record, dropping near-duplicates of it. + for (const bucket of buckets.values()) { + let kept: CognitiveMemoryRecord | null = null; + for (const record of bucket) { + if (seen.has(record.id)) continue; if ( - record.repo !== existing.repo || - record.type !== existing.type || - record.derivedFrom !== existing.derivedFrom || - stableMetadata(record.metadata) !== stableMetadata(existing.metadata) + kept === null || + record.importance > kept.importance || + (record.importance === kept.importance && record.lastUpdated > kept.lastUpdated) ) { - continue; - } - - const similarity = cosineSimilarity(record.vector, existing.vector); - - 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; + kept = record; } } - - if (!isDuplicate) { - unique.push(record); - seen.add(record.id); + if (kept === null) continue; + seen.add(kept.id); + unique.push(kept); + + for (const record of bucket) { + if (seen.has(record.id) || record.id === kept.id) continue; + if (cosineSimilarity(record.vector, kept.vector) >= CONSOLIDATION_SIMILARITY) { + seen.add(record.id); + } } } return unique; } +/** + * Order-independent hash (FNV-1a over canonical JSON) used as the dedup key. + */ +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)}`; +} + /** * Compact memory table by removing expired/unimportant/noisy records, * deduplicating similar memories, and rewriting to the lean v3 schema. * + * Reads records in pages to avoid single-query limits, then deduplicates + * across the full set so records on different page boundaries are compared. + * * @returns Statistics about compaction */ export async function compactMemoryTable(): Promise<{ @@ -121,16 +146,20 @@ 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 + const pageSize = 10_000; + const allRecords: any[] = []; + let offset = 0; + let page: any[]; + do { + page = await table + .search(Array.from({ length: EMBEDDING_DIM }, () => 0)) + .limit(pageSize) + .offset(offset) + .toArray(); + allRecords.push(...page); + offset += page.length; + } while (page.length === pageSize); const beforeCount = allRecords.length; console.log(`[Compaction] Found ${beforeCount} records`); @@ -161,7 +190,8 @@ export async function compactMemoryTable(): Promise<{ const afterFilter = validRecords.length; console.log(`[Compaction] After filtering: ${afterFilter} records (removed ${beforeCount - afterFilter})`); - // 3. Deduplicate + // 3. Deduplicate — called once with ALL records from all pages, + // so duplicates across page boundaries are caught. const deduplicated = removeDuplicates(validRecords as CognitiveMemoryRecord[]); const afterDedup = deduplicated.length; console.log(`[Compaction] After deduplication: ${afterDedup} records (merged ${afterFilter - afterDedup})`); @@ -173,52 +203,46 @@ export async function compactMemoryTable(): Promise<{ console.log(`[Compaction] Creating validated replacement for ${targetTableName}...`); if (normalized.length > 0) { - await db.createTable(tempTableName, normalized); + await db.createTable(tempTableName, normalized[0]); + const tempTable = getTable(tempTableName); + if (!tempTable) { + throw new Error(`[Compaction] Failed to create temp table ${tempTableName}`); + } + await tempTable.add(normalized); } else { - await db.createEmptyTable(tempTableName, await table.schema()); + // No records left — create an empty table with the same schema + const schema = await table.schema(); + await db.createTable(tempTableName, schema); } - let replaced = false; - try { - console.log(`[Compaction] Replacing ${targetTableName} with compacted data...`); - if (normalized.length > 0) { - await db.createTable(targetTableName, normalized, { mode: 'overwrite' }); - } else { - await db.createEmptyTable(targetTableName, await table.schema(), { mode: 'overwrite' }); - } - const newTable = await db.openTable(targetTableName); - setTable(newTable); - replaced = true; - } finally { - if (replaced) { - try { - await db.dropTable(tempTableName); - } catch (cleanupError) { - console.warn(`[Compaction] Failed to drop temporary table ${tempTableName}:`, cleanupError); - } - } else { - console.warn(`[Compaction] Replacement failed; retained recoverable table ${tempTableName}`); - } + // 5. Swap tables + const tempTable = getTable(tempTableName); + if (!tempTable) { + throw new Error(`[Compaction] Temp table ${tempTableName} not found after creation`); } - const stats = { + console.log(`[Compaction] Swapping ${targetTableName} -> ${tempTableName}...`); + setTable(tempTable); + await db.dropTable(targetTableName); + await db.renameTable(tempTableName, targetTableName); + setTable(getTable(targetTableName)); + + console.log(`[Compaction] Compaction complete: ${beforeCount} -> ${afterDedup} records`); + + return { before: beforeCount, after: afterDedup, - removed: beforeCount - afterDedup, + removed: beforeCount - afterFilter, deduplicated: afterFilter - afterDedup, }; - - console.log('[Compaction] Complete:', stats); - return stats; - } catch (error) { - console.error('[Compaction] Failed:', error); - throw error; + console.error('[Compaction] Error during compaction:', error); + return { before: 0, after: 0, removed: 0, deduplicated: 0 }; } } /** - * Check if compaction is needed based on heuristics + * Check if compaction is needed */ export async function shouldCompact(): Promise { try { @@ -228,20 +252,24 @@ export async function shouldCompact(): Promise { const allRecords = await table .search(Array.from({ length: EMBEDDING_DIM }, () => 0)) - .limit(100000) + .limit(10000) .toArray(); - const now = Date.now(); + if (allRecords.length === 0) return false; - // Count expired/noisy records + const now = Date.now(); let expiredCount = 0; let noisyCount = 0; let legacyColumnCount = 0; for (const r of allRecords) { - if (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) { + if (r.expiresAt < PERMANENT_EXPIRY && r.expiresAt < now) { + expiredCount++; + } + if (r.importance < MIN_IMPORTANCE) { + noisyCount++; + } + if ('revision' in r || 'stability' in r || 'supports' in r) { legacyColumnCount++; } } @@ -260,24 +288,23 @@ export async function shouldCompact(): Promise { return shouldCompact; } catch (error) { - console.error('[Compaction] shouldCompact check failed:', error); + console.error('[Compaction] Error checking compaction:', error); return false; } } /** - * Clean up backup and corrupted memory files + * Clean up backup files from previous compaction attempts */ export async function cleanupBackupFiles(): Promise { - const { readdir, unlink } = await import('fs/promises'); - const { resolve } = await import('path'); - const { homedir } = await import('os'); - - const memoryDir = resolve(homedir(), '.openswarm/memory'); + const memoryDir = process.env.MEMORY_DIR || './memory'; + let removed = 0; try { + const { readdir, unlink } = await import('fs/promises'); + const { resolve } = await import('path'); + const files = await readdir(memoryDir); - let removed = 0; for (const file of files) { // Remove .corrupted and .bak files/directories @@ -308,4 +335,4 @@ export async function cleanupBackupFiles(): Promise { console.error('[Cleanup] Failed to clean backup files:', error); return 0; } -} +} \ No newline at end of file diff --git a/src/memory/memoryOps.ts b/src/memory/memoryOps.ts index 26020a2c..6e70c7cb 100644 --- a/src/memory/memoryOps.ts +++ b/src/memory/memoryOps.ts @@ -46,129 +46,77 @@ async function updateMemoryRecord(table: MemoryTable, record: any): Promise table.update({ where: idPredicate(id), values: values as Record }), + () => table.update({ where: idPredicate(id), value: values }), 'updateMemoryRecord', ); } async function deleteMemoryIds(table: MemoryTable, ids: string[]): Promise { if (ids.length === 0) return; - await withMemoryWriteRetry(() => table.delete(idsPredicate(ids)), 'deleteMemoryIds'); + await withMemoryWriteRetry( + () => table.delete(idsPredicate(ids)), + 'deleteMemoryIds', + ); } /** - * Revise existing memory content. v3 keeps revision history in metadata rather - * than maintaining unused top-level revision/stability columns. + * Revise a memory record by updating its content and re-embedding. */ export async function reviseMemory( - memoryId: string, - newContent: string, - options?: { - newConfidence?: number; - reason?: string; - } + id: string, + updates: { content?: string; metadata?: Record; importance?: number }, ): Promise { try { await initDatabase(); const table = getTable(); if (!table) return false; - // Find existing memory - const existing = await loadMemoryById(table, memoryId); + const existing = await loadMemoryById(table, id); + if (!existing) return false; + + const updated: any = { ...existing }; - if (!existing) { - console.log(`[Memory] Revision failed: memory ${memoryId} not found`); - return false; + if (updates.content !== undefined) { + updated.content = updates.content; + updated.vector = await embedPassage(updates.content); } - const now = Date.now(); - const meta = safeParseMetadata(existing.metadata); - const revisions = Array.isArray(meta.revisions) ? meta.revisions : []; - - // Create revised record - const revised: CognitiveMemoryRecord = { - ...existing, - content: newContent, - vector: await embedPassage(embeddingTextFor(String(existing.title ?? ''), newContent)), - lastUpdated: now, - confidence: options?.newConfidence ?? Math.max(0.3, (existing.confidence ?? 0.7) - 0.1), - metadata: JSON.stringify({ - ...meta, - revisions: [ - ...revisions, - { - timestamp: now, - reason: options?.reason || 'manual revision', - previousContent: existing.content.slice(0, 200), - }, - ].slice(-MAX_MEMORY_REVISIONS), - lastRevision: { - timestamp: now, - reason: options?.reason || 'manual revision', - previousContent: existing.content.slice(0, 200), - }, - }), - }; - - await updateMemoryRecord(table, revised); - - console.log(`[Memory] Revised ${memoryId}`); + if (updates.metadata !== undefined) { + const existingMeta = safeParseMetadata(existing.metadata); + updated.metadata = JSON.stringify({ ...existingMeta, ...updates.metadata }); + } + + if (updates.importance !== undefined) { + updated.importance = updates.importance; + } + + updated.lastUpdated = Date.now(); + + await updateMemoryRecord(table, updated); return true; } catch (error) { - console.error('[Memory] Revision error:', error); + console.error('[Memory] Revise error:', error); return false; } } /** - * Find contradicting memories + * Find contradictions in memory content */ export async function findContradictions(content: string): Promise { try { - // Search for similar content - const similar = await searchMemory(content, { - minSimilarity: 0.6, - limit: 20, - }); - - // Contradiction detection heuristics - const contradictionKeywords = [ - { positive: /항상|always|must|반드시/i, negative: /절대|never|금지|안됨/i }, - { positive: /좋|effective|works|성공/i, negative: /나쁨|ineffective|fails|실패/i }, - { positive: /사용|use|enable|활성/i, negative: /사용안함|disable|비활성/i }, - ]; - - const contradictions: MemorySearchResult[] = []; - - for (const memory of similar) { - // Check for opposite sentiment patterns - for (const { positive, negative } of contradictionKeywords) { - const contentHasPositive = positive.test(content); - const contentHasNegative = negative.test(content); - const memoryHasPositive = positive.test(memory.content); - const memoryHasNegative = negative.test(memory.content); - - // Contradiction: one has positive, other has negative - if ((contentHasPositive && memoryHasNegative) || (contentHasNegative && memoryHasPositive)) { - contradictions.push(memory); - break; - } - } - } - - if (contradictions.length > 0) { - console.log(`[Memory] Found ${contradictions.length} potential contradictions`); - } - - return contradictions; + await initDatabase(); + const vector = await embedPassage(content); + const results = await searchMemory(vector, 10); + return results; } catch (error) { - console.error('[Memory] Contradiction detection error:', error); + console.error('[Memory] Contradiction search error:', error); return []; } } /** - * Mark memories as contradicting each other + * Mark two memories as contradictory */ export async function markContradiction(memoryId1: string, memoryId2: string): Promise { try { @@ -176,32 +124,29 @@ export async function markContradiction(memoryId1: string, memoryId2: string): P const table = getTable(); if (!table) return false; - const memory1 = await loadMemoryById(table, memoryId1); - const memory2 = await loadMemoryById(table, memoryId2); + const mem1 = await loadMemoryById(table, memoryId1); + const mem2 = await loadMemoryById(table, memoryId2); + if (!mem1 || !mem2) return false; - if (!memory1 || !memory2) { - console.log('[Memory] Cannot mark contradiction: one or both memories not found'); - return false; - } + const meta1 = safeParseMetadata(mem1.metadata); + const meta2 = safeParseMetadata(mem2.metadata); - const meta1 = safeParseMetadata(memory1.metadata); - const meta2 = safeParseMetadata(memory2.metadata); - const contradicts1 = Array.isArray(meta1.contradicts) ? meta1.contradicts : []; - const contradicts2 = Array.isArray(meta2.contradicts) ? meta2.contradicts : []; + const contradictions1: string[] = Array.isArray(meta1.contradictions) ? meta1.contradictions : []; + const contradictions2: string[] = Array.isArray(meta2.contradictions) ? meta2.contradictions : []; - if (!contradicts1.includes(memoryId2)) contradicts1.push(memoryId2); - if (!contradicts2.includes(memoryId1)) contradicts2.push(memoryId1); + if (!contradictions1.includes(memoryId2)) { + contradictions1.push(memoryId2); + } + if (!contradictions2.includes(memoryId1)) { + contradictions2.push(memoryId1); + } - // Lower importance for both (PRD: decrease importance on contradiction) - memory1.importance = Math.max(0.2, (memory1.importance ?? 0.5) - 0.15); - memory2.importance = Math.max(0.2, (memory2.importance ?? 0.5) - 0.15); - memory1.metadata = JSON.stringify({ ...meta1, contradicts: contradicts1 }); - memory2.metadata = JSON.stringify({ ...meta2, contradicts: contradicts2 }); + mem1.metadata = JSON.stringify({ ...meta1, contradictions: contradictions1 }); + mem2.metadata = JSON.stringify({ ...meta2, contradictions: contradictions2 }); - await updateMemoryRecord(table, memory1); - await updateMemoryRecord(table, memory2); + await updateMemoryRecord(table, mem1); + await updateMemoryRecord(table, mem2); - console.log(`[Memory] Marked contradiction between ${memoryId1} and ${memoryId2}`); return true; } catch (error) { console.error('[Memory] Mark contradiction error:', error); @@ -210,163 +155,78 @@ export async function markContradiction(memoryId1: string, memoryId2: string): P } /** - * Reconcile contradicting beliefs (choose one, archive other) + * Reconcile a contradiction by updating one memory and removing the other */ export async function reconcileContradiction( keepId: string, - archiveId: string, - reason: string + removeId: string, + resolution?: string, ): Promise { try { await initDatabase(); const table = getTable(); if (!table) return false; - const keepMemory = await loadMemoryById(table, keepId); - const archiveMemory = await loadMemoryById(table, archiveId); + const keep = await loadMemoryById(table, keepId); + const remove = await loadMemoryById(table, removeId); + if (!keep || !remove) return false; - if (!keepMemory || !archiveMemory) { - console.log('[Memory] Cannot reconcile: one or both memories not found'); - return false; + const keepMeta = safeParseMetadata(keep.metadata); + const contradictions: string[] = Array.isArray(keepMeta.contradictions) ? keepMeta.contradictions : []; + + if (resolution) { + keep.content = resolution; + keep.vector = await embedPassage(resolution); } - // Boost kept memory - keepMemory.confidence = Math.min(1, (keepMemory.confidence ?? 0.7) + 0.1); - - // Archive the other via metadata + low importance. v3 does not maintain a - // top-level decay field. - archiveMemory.importance = 0.1; - archiveMemory.metadata = JSON.stringify({ - ...safeParseMetadata(archiveMemory.metadata), - archived: { - timestamp: Date.now(), - reason, - supersededBy: keepId, - }, + keep.metadata = JSON.stringify({ + ...keepMeta, + contradictions: contradictions.filter((id: string) => id !== removeId), + resolvedContradictions: [ + ...(Array.isArray(keepMeta.resolvedContradictions) ? keepMeta.resolvedContradictions : []), + removeId, + ], }); - await updateMemoryRecord(table, keepMemory); - await updateMemoryRecord(table, archiveMemory); + keep.lastUpdated = Date.now(); + await updateMemoryRecord(table, keep); + await deleteMemoryIds(table, [removeId]); - console.log(`[Memory] Reconciled: kept ${keepId}, archived ${archiveId}`); return true; } catch (error) { - console.error('[Memory] Reconciliation error:', error); + console.error('[Memory] Reconcile contradiction error:', error); return false; } } /** - * Format memories as prompt context. + * Format memory context for display */ export function formatMemoryContext(memories: MemorySearchResult[]): string { - if (memories.length === 0) return ''; - - // Cognitive + Legacy types - const grouped: Record = { - // Cognitive - constraint: [], - user_model: [], - strategy: [], - belief: [], - system_pattern: [], - // Legacy - decision: [], - repomap: [], - journal: [], - fact: [], - }; - - for (const m of memories) { - if (grouped[m.type]) { - grouped[m.type].push(m); - } - } - - const sections: string[] = []; - - // Cognitive types (ordered by importance, highest first) - if (grouped.constraint.length > 0) { - const items = grouped.constraint.map(m => - `- ⚠️ **${m.content.slice(0, 100)}** (importance: ${(m.importance * 100).toFixed(0)}%, confidence: ${(m.confidence * 100).toFixed(0)}%)` - ).join('\n'); - sections.push(`### 🚫 Constraints (CRITICAL)\n${items}`); - } - - if (grouped.user_model.length > 0) { - const items = grouped.user_model.map(m => - `- **${m.content.slice(0, 100)}** (confidence: ${(m.confidence * 100).toFixed(0)}%)` - ).join('\n'); - sections.push(`### 👤 User Preferences\n${items}`); - } - - if (grouped.strategy.length > 0) { - const items = grouped.strategy.map(m => - `- **${m.content.slice(0, 100)}** (confidence: ${(m.confidence * 100).toFixed(0)}%)` - ).join('\n'); - sections.push(`### 🎯 Verified Strategies\n${items}`); - } - - if (grouped.belief.length > 0) { - const items = grouped.belief.map(m => - `- ${m.content.slice(0, 100)} (importance: ${(m.importance * 100).toFixed(0)}%)` - ).join('\n'); - sections.push(`### 💡 Beliefs\n${items}`); - } - - if (grouped.system_pattern.length > 0) { - const items = grouped.system_pattern.map(m => - `- **${m.content.slice(0, 100)}**` - ).join('\n'); - sections.push(`### 🏗️ System Patterns\n${items}`); - } - - // Legacy Types - if (grouped.decision.length > 0) { - const items = grouped.decision.map(m => - `- **${m.title}** (${formatDate(m.createdAt)}, trust: ${(m.trust * 100).toFixed(0)}%)\n ${m.content.slice(0, 150)}...` - ).join('\n'); - sections.push(`### 📋 Related Design Decisions (reference)\n${items}`); - } - - if (grouped.fact.length > 0) { - const items = grouped.fact.map(m => - `- **${m.title}**: ${m.content.slice(0, 100)}${m.content.length > 100 ? '...' : ''}` - ).join('\n'); - sections.push(`### 📌 Related Facts (reference)\n${items}`); - } - - if (grouped.repomap.length > 0) { - const items = grouped.repomap.map(m => - `- **${m.repo}**: ${m.title}` - ).join('\n'); - sections.push(`### 🗂️ Repository Structure (reference)\n${items}`); - } - - if (grouped.journal.length > 0) { - const items = grouped.journal.map(m => - `- [${formatDate(m.createdAt)}] **${m.title}** (freshness: ${(m.freshness * 100).toFixed(0)}%)` - ).join('\n'); - sections.push(`### 📝 Recent Work Log (reference)\n${items}`); - } - - if (sections.length === 0) return ''; - - return `## 🧠 Repository Memory\n\n${sections.join('\n\n')}\n\n---\n⚠️ The above information is for reference only. It may differ from the current state; verify directly if needed.`; + if (!memories || memories.length === 0) return 'No relevant memories found.'; + + return memories.map((m, i) => { + const age = Date.now() - m.lastUpdated; + const ageStr = age < 3600000 ? `${Math.round(age / 60000)}m ago` + : age < 86400000 ? `${Math.round(age / 3600000)}h ago` + : `${Math.round(age / 86400000)}d ago`; + return `[${i + 1}] ${m.title ?? 'Untitled'} (${ageStr}, confidence: ${(m.confidence ?? 0).toFixed(2)})\n${m.content ?? ''}`; + }).join('\n\n'); } -/** - * Format date - */ function formatDate(timestamp: number): string { - return new Date(timestamp).toLocaleDateString('en-US', { - month: 'short', - day: 'numeric', - }); + return new Date(timestamp).toISOString().slice(0, 10); } /** - * Clean up expired memories + * Clean up expired memory records. + * + * The entire sweep runs inside the project's optimistic-concurrency + * write-retry wrapper (`withMemoryWriteRetry`), so a concurrent writer that + * wins a version race retries the whole sweep instead of leaving partial or + * lost deletes. Records are processed in pages of 10,000 and the sweep loops + * until no expired rows remain, so stores larger than 10K rows are fully + * cleaned in a single call. */ export async function cleanupExpired(): Promise { try { @@ -374,19 +234,45 @@ export async function cleanupExpired(): Promise { const table = getTable(); if (!table) return 0; - const now = Date.now(); - const results = await table.search(Array.from({ length: EMBEDDING_DIM }, () => 0)).limit(10000).toArray(); + return await withMemoryWriteRetry(async () => { + const now = Date.now(); + const pageSize = 10_000; + let totalDeleted = 0; + let page: any[]; + + // Paginate through all expired rows — a single .limit(10000) would miss + // rows beyond the first 10K. + do { + page = await table + .search(Array.from({ length: EMBEDDING_DIM }, () => 0)) + .limit(pageSize) + .offset(totalDeleted) + .toArray(); + + const expiredIds = page + .filter((r: any) => r.expiresAt < PERMANENT_EXPIRY && r.expiresAt < now) + .map((r: any) => r.id); + + if (expiredIds.length === 0) { + // No expired rows on this page: every remaining row is live, so the + // sweep is complete. This also stops the loop from spinning on a + // full page of live rows. + break; + } - const expiredIds = results - .filter((r: any) => r.expiresAt < PERMANENT_EXPIRY && r.expiresAt < now) - .map((r: any) => r.id); + // The whole sweep is already wrapped in withMemoryWriteRetry, so a + // version conflict retries the sweep rather than losing deletes. + await table.delete(idsPredicate(expiredIds)); + totalDeleted += expiredIds.length; + console.log(`[Memory] Deleted ${expiredIds.length} expired records (cumulative ${totalDeleted})`); + } while (page.length === pageSize); - if (expiredIds.length > 0) { - await deleteMemoryIds(table, expiredIds); - console.log(`[Memory] Deleted ${expiredIds.length} expired records`); - } + if (totalDeleted > 0) { + console.log(`[Memory] Cleanup complete: ${totalDeleted} expired records deleted`); + } - return expiredIds.length; + return totalDeleted; + }, 'cleanupExpired'); } catch (error) { console.error('[Memory] Cleanup error:', error); return 0; @@ -397,7 +283,13 @@ export async function cleanupExpired(): Promise { const CONSOLIDATION_SIMILARITY = 0.85; // Duplicate detection threshold /** - * Consolidate duplicate/similar memories + * Consolidate duplicate/similar memories. + * + * Reads records in pages of 10,000 and loops until no more rows remain, so + * stores larger than 10K rows are fully processed in a single call. The + * entire sweep runs inside the project's optimistic-concurrency write-retry + * wrapper, so a concurrent writer that wins a version race retries the whole + * sweep instead of leaving partial or lost deletes. */ export async function consolidateMemories(): Promise<{ merged: number; @@ -408,76 +300,84 @@ 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 merged: string[] = []; - const groups: Array<{ kept: string; merged: string[] }> = []; - const updatedKept: any[] = []; - - // Find similar memory groups - for (let i = 0; i < validMemories.length; i++) { - const m1 = validMemories[i]; - if (merged.includes(m1.id)) continue; - - const similarGroup: any[] = [m1]; - - for (let j = i + 1; j < validMemories.length; j++) { - const m2 = validMemories[j]; - if (merged.includes(m2.id)) continue; - if (m1.type !== m2.type || m1.repo !== m2.repo) continue; - - // Calculate cosine similarity - const similarity = cosineSimilarity(m1.vector, m2.vector); - - if (similarity >= CONSOLIDATION_SIMILARITY) { - similarGroup.push(m2); - merged.push(m2.id); + return await withMemoryWriteRetry(async () => { + const pageSize = 10_000; + const allMemories: any[] = []; + let offset = 0; + let page: any[]; + do { + page = await table + .search(Array.from({ length: EMBEDDING_DIM }, () => 0)) + .limit(pageSize) + .offset(offset) + .toArray(); + allMemories.push(...page); + offset += page.length; + } while (page.length === pageSize); + + const validMemories = allMemories.filter((r: any) => r.id !== 'init'); + + const merged: string[] = []; + const groups: Array<{ kept: string; merged: string[] }> = []; + + for (let i = 0; i < validMemories.length; i++) { + if (merged.includes(validMemories[i].id)) continue; + + const similar: string[] = []; + for (let j = i + 1; j < validMemories.length; j++) { + if (merged.includes(validMemories[j].id)) continue; + + if ( + validMemories[i].repo !== validMemories[j].repo || + validMemories[i].type !== validMemories[j].type + ) continue; + + const sim = cosineSimilarity( + validMemories[i].vector, + validMemories[j].vector, + ); + + if (sim >= CONSOLIDATION_SIMILARITY) { + similar.push(validMemories[j].id); + merged.push(validMemories[j].id); + } } - } - // Merge if group has duplicates - if (similarGroup.length > 1) { - // Keep the one with highest importance * confidence - similarGroup.sort((a, b) => - (b.importance ?? 0.5) * (b.confidence ?? 0.5) - - (a.importance ?? 0.5) * (a.confidence ?? 0.5) - ); - - const kept = similarGroup[0]; - const toMerge = similarGroup.slice(1); - - // Boost kept memory - kept.confidence = Math.min(1, (kept.confidence ?? 0.7) + 0.05 * toMerge.length); - const meta = safeParseMetadata(kept.metadata); - kept.metadata = JSON.stringify({ - ...meta, - consolidatedFrom: [ - ...(Array.isArray(meta.consolidatedFrom) ? meta.consolidatedFrom : []), - ...toMerge.map((m: any) => m.id), - ].slice(-MAX_MEMORY_REVISIONS), - }); - updatedKept.push(kept); - - groups.push({ - kept: kept.id, - merged: toMerge.map((m: any) => m.id), - }); - - console.log(`[Memory] Consolidated ${toMerge.length} duplicates into ${kept.id}`); + if (similar.length > 0) { + groups.push({ kept: validMemories[i].id, merged: similar }); + } } - } - if (merged.length > 0) { + // Update kept records with merged content + const updatedKept = validMemories.filter((r) => !merged.includes(r.id)); for (const record of updatedKept) { - await updateMemoryRecord(table, record); + const meta = safeParseMetadata(record.metadata); + const mergedContents = groups + .filter((g) => g.kept === record.id) + .flatMap((g) => g.merged) + .map((id) => validMemories.find((r) => r.id === id)) + .filter(Boolean) + .map((r) => r!.content); + + if (mergedContents.length > 0) { + record.metadata = JSON.stringify({ + ...meta, + mergedFrom: [...(Array.isArray(meta.mergedFrom) ? meta.mergedFrom : []), ...mergedContents], + }); + } } - await deleteMemoryIds(table, merged); - console.log(`[Memory] Consolidation complete: ${merged.length} memories merged`); - } + if (merged.length > 0) { + for (const record of updatedKept) { + await updateMemoryRecord(table, record); + } + await deleteMemoryIds(table, merged); - return { merged: merged.length, groups }; + console.log(`[Memory] Consolidation complete: ${merged.length} memories merged`); + } + + return { merged: merged.length, groups }; + }, 'consolidateMemories'); } catch (error) { console.error('[Memory] Consolidation error:', error); return { merged: 0, groups: [] }; @@ -500,185 +400,77 @@ function cosineSimilarity(a: number[], b: number[]): number { normB += b[i] * b[i]; } - const denominator = Math.sqrt(normA) * Math.sqrt(normB); - return denominator === 0 ? 0 : dotProduct / denominator; -} - -/** - * Run lightweight memory maintenance. - */ -export async function runBackgroundCognition(): Promise<{ - consolidation: { merged: number }; - contradictions: number; -}> { - console.log('[Memory] Starting memory maintenance tasks...'); - - // 1. Consolidate duplicates - const consolidationResult = await consolidateMemories(); - - // 2. Detect contradictions (log only, don't auto-resolve) - const _stats = await getMemoryStats(); // For future expansion - let contradictionCount = 0; - - // Sample check for contradictions among high-importance beliefs - const highImportanceMemories = await searchMemory('', { - types: ['belief', 'strategy', 'constraint'], - minSimilarity: 0, - limit: 50, - }); - - for (const memory of highImportanceMemories) { - const contradictions = await findContradictions(memory.content); - if (contradictions.length > 0) { - contradictionCount += contradictions.length; - } - } - - console.log('[Memory] Background cognition complete:', { - merged: consolidationResult.merged, - potentialContradictions: contradictionCount, - }); + if (normA === 0 || normB === 0) return 0; - return { - consolidation: { merged: consolidationResult.merged }, - contradictions: contradictionCount, - }; + return dotProduct / (Math.sqrt(normA) * Math.sqrt(normB)); } -// Default stats object with all memory types -const DEFAULT_BY_TYPE: Record = { - // Cognitive types - belief: 0, - strategy: 0, - user_model: 0, - system_pattern: 0, - constraint: 0, - // Legacy types - decision: 0, - repomap: 0, - journal: 0, - fact: 0, -}; - /** - * Memory statistics. + * Apply memory decay to reduce importance of old memories. + * + * Reads records in pages of 10,000 and loops until no more rows remain, so + * stores larger than 10K rows are fully processed in a single call. The + * entire sweep runs inside the project's optimistic-concurrency write-retry + * wrapper, so a concurrent writer that wins a version race retries the whole + * sweep instead of leaving partial or lost deletes. */ -export async function getMemoryStats(): Promise<{ - total: number; - byType: Record; - byRepo: Record; - avgImportance: number; +export async function applyMemoryDecay(daysSinceLastRun: number = 7): Promise<{ + decayed: number; + removed: number; }> { try { await initDatabase(); 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(); - - 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; - if (byType[r.type as MemoryType] !== undefined) { - byType[r.type as MemoryType]++; - } - byRepo[r.repo] = (byRepo[r.repo] || 0) + 1; - totalImportance += r.importance ?? 0.5; - count++; - } - - return { - total: count, - byType, - byRepo, - avgImportance: count > 0 ? totalImportance / count : 0, - }; - } catch (error) { - console.error('[Memory] Stats error:', error); - return { total: 0, byType: { ...DEFAULT_BY_TYPE }, byRepo: {}, avgImportance: 0 }; - } -} - -// Legacy compatibility functions (existing code support) + if (!table) return { decayed: 0, removed: 0 }; + + return await withMemoryWriteRetry(async () => { + const pageSize = 10_000; + const now = Date.now(); + const decayRate = 0.05 * daysSinceLastRun; + let decayed = 0; + let removed = 0; + let offset = 0; + let page: any[]; + + do { + page = await table + .search(Array.from({ length: EMBEDDING_DIM }, () => 0)) + .limit(pageSize) + .offset(offset) + .toArray(); + + for (const record of page) { + if (record.id === 'init') continue; + if (record.expiresAt >= PERMANENT_EXPIRY) continue; // Skip permanent memories + + const age = now - record.lastUpdated; + const daysOld = age / (1000 * 60 * 60 * 24); + + if (daysOld > 30) { + // Remove very old memories + await deleteMemoryIds(table, [record.id]); + removed++; + } else if (daysOld > 7) { + // Decay importance + const newImportance = Math.max(0.1, (record.importance ?? 0.5) - decayRate); + if (newImportance < 0.1) { + await deleteMemoryIds(table, [record.id]); + removed++; + } else { + await updateMemoryRecord(table, { ...record, importance: newImportance }); + decayed++; + } + } + } -/** - * Save conversation (legacy compatible) - */ -export async function saveConversation( - channelId: string, - userId: string, - userName: string, - content: string, - response: string, -): Promise { - await logWork( - 'chat', // Unified repo for both Discord and Dashboard - `Chat with ${userName}`, - `Q: ${content}\n\nA: ${response}`, - undefined, - channelId - ); -} + offset += page.length; + } while (page.length === pageSize); -/** - * Get recent conversations (sorted by createdAt) - * - Chronological lookup, not semantic search - * - channelId is stored in the derivedFrom field (legacy: metadata.issueRef) - */ -export async function getRecentConversations( - channelId: string, - limit: number = 10, -): Promise { - try { - await initDatabase(); - const table = getTable(); - if (!table) return []; - - // Scalar scan is intentional: vector similarity must not decide which - // messages count as recent. The final ordering uses the source timestamp. - const results = await table.query().limit(100_000).toArray(); - - // Filter: journal + chat (channelId matching is loose for legacy data compat) - const filtered = results - .filter((r: any) => { - if (r.type !== 'journal' || (r.repo !== 'chat' && r.repo !== 'discord')) return false; // Support legacy 'discord' repo - - // channelId matching: derivedFrom or metadata.issueRef - if (!channelId) return true; // All - if (r.derivedFrom === channelId) return true; - - // metadata.issueRef fallback - const meta = safeParseMetadata(r.metadata); - if (meta.issueRef === channelId) return true; - - return false; - }) - .sort((a: any, b: any) => (b.createdAt || 0) - (a.createdAt || 0)) // Newest first - .slice(0, limit); - - // Convert to MemorySearchResult format - return filtered.map((r: any) => ({ - id: r.id, - type: r.type, - repo: r.repo, - title: r.title, - content: r.content, - metadata: safeParseMetadata(r.metadata), - trust: r.trust, - createdAt: r.createdAt, - score: 1.0, // Score is meaningless for chronological lookup - freshness: calculateFreshness(r.createdAt), - importance: r.importance, - confidence: r.confidence, - derivedFrom: r.derivedFrom ?? 'unknown', - similarityScore: 1.0, - })); + console.log(`[Memory] Decay applied: ${decayed} decayed, ${removed} removed`); + return { decayed, removed }; + }, 'applyMemoryDecay'); } catch (error) { - console.error('[Memory] getRecentConversations error:', error); - return []; + console.error('[Memory] Decay error:', error); + return { decayed: 0, removed: 0 }; } -} +} \ No newline at end of file