Skip to content
Closed
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
35 changes: 35 additions & 0 deletions src/memory/compaction.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
});
});
207 changes: 117 additions & 90 deletions src/memory/compaction.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, CognitiveMemoryRecord[]>();
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<string>();

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<{
Expand All @@ -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`);
Expand Down Expand Up @@ -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})`);
Expand All @@ -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<boolean> {
try {
Expand All @@ -228,20 +252,24 @@ export async function shouldCompact(): Promise<boolean> {

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++;
}
}
Expand All @@ -260,24 +288,23 @@ export async function shouldCompact(): Promise<boolean> {
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<number> {
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
Expand Down Expand Up @@ -308,4 +335,4 @@ export async function cleanupBackupFiles(): Promise<number> {
console.error('[Cleanup] Failed to clean backup files:', error);
return 0;
}
}
}
Loading