Skip to content
Merged
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
34 changes: 34 additions & 0 deletions src/memory/compaction.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
});
});
209 changes: 136 additions & 73 deletions src/memory/compaction.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, true> = {
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') {
Expand Down Expand Up @@ -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<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) {
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<CognitiveMemoryRecord>;
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<{
Expand All @@ -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;
Expand All @@ -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})`);

Expand Down Expand Up @@ -218,43 +272,52 @@ 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<boolean> {
try {
await initDatabase();
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<CognitiveMemoryRecord>;
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;
Expand Down
1 change: 1 addition & 0 deletions src/memory/memoryOps.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) => ({
Expand Down
Loading
Loading