From 36bed4cb888369699274fdf3e07b133b1a4786e0 Mon Sep 17 00:00:00 2001 From: antra-tess Date: Fri, 25 Sep 2026 17:20:04 -0700 Subject: [PATCH] feat(strategy): fail loudly on invalid store topology; merge adjacency in store order MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Every load audits the summary archive for crossed ownership — a summary whose leaves are not contiguous among owned messages (chunk members ∪ chunk-linked L1 sourceIds) in store order. topologyPolicy 'reject' (default) throws StoreTopologyError from initialize, so ContextManager.open refuses the store until it is repaired; 'report' logs at error level and surfaces the count through getCompressionDebt().topologyViolations (state critical). A kv-unified config that opts into gap handling (preserveGapBearingSummaries / treeifyNonContiguousSummaries) defaults to 'report'. ContextManager.open releases a store it opened when initialize rejects. scripts/audit-topology.ts runs the same audit read-only. Merge adjacency is judged in store order instead of chunk-record order, so a chunk minted late over an early message (#122's head ratchet) no longer joins the frontier's run; demand-path merges (#95) are split into strictly adjacent runs; executeMerge refuses any group that is not one level below the target and strictly adjacent in store order — no model call, the entry moves to the merge quarantine with outcome topology_violation, and getCompressionDebt().topologyRefusals / state critical say so. A crossed node is never minted. Calibrated on Linn's store copy (the 7 summaries the kv-unified forest rejects) and on the synthetic corpus (7 crossed L2/L3s). Co-Authored-By: Claude Fable 5.1 --- changelog.d/store-topology-fail-loud.added.md | 19 + scripts/audit-topology.ts | 57 +++ src/context-manager.ts | 13 +- src/index.ts | 1 + src/strategies/autobiographical.ts | 377 +++++++++++++++++- src/types/strategy.ts | 17 + test/topology-guard.test.ts | 293 ++++++++++++++ 7 files changed, 759 insertions(+), 18 deletions(-) create mode 100644 changelog.d/store-topology-fail-loud.added.md create mode 100644 scripts/audit-topology.ts create mode 100644 test/topology-guard.test.ts diff --git a/changelog.d/store-topology-fail-loud.added.md b/changelog.d/store-topology-fail-loud.added.md new file mode 100644 index 0000000..5478b85 --- /dev/null +++ b/changelog.d/store-topology-fail-loud.added.md @@ -0,0 +1,19 @@ +- Store topology fails loudly. Every load audits the summary archive for + crossed ownership (a summary whose leaves are not contiguous among + chunk-owned messages in store order — issue #122's cross-era merges, + restore/branch interleavings, hand surgery). `topologyPolicy: 'reject'` + (default) throws `StoreTopologyError` from `initialize`, so + `ContextManager.open` refuses the store until it is repaired; + `'report'` logs the violations at error level and reports them through + `getCompressionDebt().topologyViolations` (state `critical`). A kv-unified + config that opts into gap handling (`preserveGapBearingSummaries` or + `treeifyNonContiguousSummaries`) defaults to `'report'`. + `scripts/audit-topology.ts` runs the same audit read-only on a store path. +- Merge adjacency is judged in store order, not chunk-record order, so a chunk + minted late over an early message can no longer join the frontier's merge run + (the second half of #122). Demand-path merges (`enqueueMergeForRange`, #95) + are split into strictly adjacent runs like the threshold path. `executeMerge` + refuses any group that is not one level below the target and strictly adjacent + in store order: no model call, the entry moves into the merge quarantine with + outcome `topology_violation`, and `getCompressionDebt().topologyRefusals` / + state `critical` say so. A crossed node is never minted. diff --git a/scripts/audit-topology.ts b/scripts/audit-topology.ts new file mode 100644 index 0000000..a5e90bb --- /dev/null +++ b/scripts/audit-topology.ts @@ -0,0 +1,57 @@ +/** + * Audit a store's summary topology: every summary whose leaves are not + * contiguous among chunk-owned messages in store order (crossed ownership — + * issue #122 cross-era merges, restore/branch interleavings, hand surgery). + * Read-only; the same audit `initialize` runs (topologyPolicy). + * + * Usage: + * node dist/scripts/audit-topology.js [--namespace ] [--json] + * + * Exit code 0 = clean, 2 = violations found, 1 = could not open. + * Run it on a stopped resident's store or a copy: opening a chronicle store + * takes its lock and may rewrite its state index. + */ + +import { ContextManager, AutobiographicalStrategy } from '../src/index.js'; + +async function main(): Promise { + const args = process.argv.slice(2); + const storePath = args.find((a) => !a.startsWith('--')); + if (!storePath) { + console.error('Usage: audit-topology [--namespace ] [--json]'); + process.exit(1); + } + const nsAt = args.indexOf('--namespace'); + const namespace = nsAt >= 0 ? args[nsAt + 1] : undefined; + const json = args.includes('--json'); + const strategy = new AutobiographicalStrategy({ + adaptiveResolution: true, + hierarchical: true, + autoTickOnNewMessage: false, + topologyPolicy: 'report', + }); + const membrane = { complete: async () => ({ content: [{ type: 'text', text: '[audit]' }] }) }; + const manager = await ContextManager.open({ + path: storePath, strategy, membrane: membrane as never, ...(namespace ? { namespace } : {}), + }); + try { + const violations = strategy.getTopologyViolations(); + const debt = strategy.getCompressionDebt(); + if (json) { + console.log(JSON.stringify({ store: storePath, namespace, violations, debt }, null, 2)); + } else { + console.log(`${storePath}${namespace ? ` (${namespace})` : ''}: ${manager.getMessageCount()} messages, ` + + `${violations.length} topology violation(s), compression debt ${debt.state}`); + for (const v of violations) { + console.log(` L${v.level} ${v.id}: ${v.leafCount} leaves ${v.span.first}..${v.span.last}, ` + + `${v.holes} hole(s) e.g. ${v.holeSample.join(',')}` + + `${v.holeOwners.length ? ` owned by ${v.holeOwners.join(',')}` : ''}`); + } + } + process.exitCode = violations.length > 0 ? 2 : 0; + } finally { + manager.close(); + } +} + +main().catch((error) => { console.error(error); process.exit(1); }); diff --git a/src/context-manager.ts b/src/context-manager.ts index d376151..6484af0 100644 --- a/src/context-manager.ts +++ b/src/context-manager.ts @@ -345,9 +345,18 @@ export class ContextManager { auxiliaryStores, ); - // Initialize strategy + // Initialize strategy. A strategy that refuses the store (e.g. + // StoreTopologyError) must not leave a store we opened locked behind a + // rejected promise: release it, then rethrow. const openingBranch = observeStoreBranch(store); - await manager.initializeStrategy(openingBranch); + try { + await manager.initializeStrategy(openingBranch); + } catch (error) { + if (ownsStore) { + try { store.close(); } catch { /* the initialize error is the one to report */ } + } + throw error; + } manager.initialized = true; return manager; diff --git a/src/index.ts b/src/index.ts index a64a494..6381af8 100644 --- a/src/index.ts +++ b/src/index.ts @@ -36,6 +36,7 @@ export type { ConfigLayer, ConfigResolutionSemantics, EffectiveConfigReport } fr // classifyInferenceError); exporting them from the root gives consumers a // real `instanceof` instead of stringly-typed `err.name` matching. export { OverBudgetError, UncoveredDropError } from './adaptive/picker.js'; +export { StoreTopologyError, type TopologyViolation } from './strategies/autobiographical.js'; export type { OverBudgetDiagnostics } from './adaptive/picker.js'; // Types diff --git a/src/strategies/autobiographical.ts b/src/strategies/autobiographical.ts index 73810a9..fa2844b 100644 --- a/src/strategies/autobiographical.ts +++ b/src/strategies/autobiographical.ts @@ -481,6 +481,8 @@ interface CompressionRefusalNormalizedConfig { type CompressionAttemptOutcome = | 'refusal' + /** The merge group is not contiguous in store order; refused before any model call. */ + | 'topology_violation' | 'unusable_empty' | 'provider_error' | 'admission_rejected' @@ -546,6 +548,42 @@ interface MergeQuarantineRecord { quarantinedAt: number; } +/** One summary whose leaves are not contiguous among chunk-owned messages. */ +export interface TopologyViolation { + id: string; + level: number; + leafCount: number; + /** First/last leaf message id in store order. */ + span: { first: string; last: string }; + /** Chunk-owned messages inside the span that the summary does not own. */ + holes: number; + /** Up to five hole message ids. */ + holeSample: string[]; + /** Up to five L1 summaries that own the holes (the interleaved representation). */ + holeOwners: string[]; +} + +/** + * Thrown by `initialize` (so by `ContextManager.open`) when the summary + * archive carries crossed ownership and `topologyPolicy` is `'reject'`. + * The store is intact; nothing was written. Repair it (or open with + * `topologyPolicy: 'report'` to inspect) before running a resident on it. + */ +export class StoreTopologyError extends Error { + constructor(readonly violations: TopologyViolation[]) { + super( + `store topology rejected: ${violations.length} summar${violations.length === 1 ? 'y owns' : 'ies own'} ` + + `non-contiguous leaves — ` + + violations.slice(0, 8).map((v) => + `L${v.level} ${v.id} (${v.leafCount} leaves ${v.span.first}..${v.span.last}, ${v.holes} hole(s)` + + `${v.holeOwners.length ? ` owned by ${v.holeOwners.join('/')}` : ''})`).join('; ') + + (violations.length > 8 ? `; +${violations.length - 8} more` : '') + + `. Repair the store or set topologyPolicy: 'report' to open it for inspection.`, + ); + this.name = 'StoreTopologyError'; + } +} + interface CompressionRefusalOutcomeRecord { curveLabel: string; requestHash: string; @@ -1667,6 +1705,8 @@ export class AutobiographicalStrategy implements ResettableStrategy { this.rebuildChunks(ctx.messageStore); if (abortIfStale()) return; this.sanitizePersistedMergeQueue(ctx.messageStore); + if (abortIfStale()) return; + this.assertStoreTopology(messages); // Kick the merge ladder for pre-existing unmerged summaries. Normally a // compression/merge completion does this, but a store that boots with a // backlog above threshold and an empty queue (e.g. after a pyramid @@ -3738,6 +3778,272 @@ export class AutobiographicalStrategy implements ResettableStrategy { return merge; } + // =========================================================================== + // Store topology: contiguity of ownership in STORE order. + // + // Merge adjacency used to be judged in chunk-record order. Records are + // appended when a chunk is minted, so a chunk minted late over an early + // message (issue #122: the head ratchet peeling opening messages into + // one-message L1s days later) sat next to the open frontier and merged with + // it, producing L2/L3s that mix the chronicle's opening with weeks-later + // material. The demand path (#95) never checked adjacency at all. Every + // grouping decision and the mint itself now use one index: chunk-owned + // messages, positioned by the store listing. Messages no chunk owns (head, + // never-chunked) occupy no position, so they never split a run; a hole is + // always another live representation. + // =========================================================================== + + /** Store position per message id, refreshed from every listing we see. */ + private _storeOrder = new Map(); + /** Load-time audit result (see assertStoreTopology). */ + private topologyViolations: TopologyViolation[] = []; + /** Merges refused by executeMerge because they would have minted a crossed node. */ + private topologyRefusals = 0; + + protected refreshStoreOrder(messages: ReadonlyArray<{ id: MessageId }>): void { + if (messages.length === 0) return; + const order = new Map(); + for (let i = 0; i < messages.length; i++) order.set(messages[i].id, i); + this._storeOrder = order; + } + + /** + * Position per owned message id (chunk records ∪ live L1 sourceIds): store + * order where known, then any owned id the last listing did not contain + * (messages appended since) after it, in record order. Without any listing + * (hosts that drive the grouping without a store view) this degrades to + * record order. + */ + protected mergePositionIndex(): Map { + // Ownership authority = chunk members ∪ live L1 sourceIds (an L1 is + // coverage even where its chunk record is missing; see ownedMessageIds). + const chunkMember = new Set(); + const members: MessageId[] = []; + for (const ch of this.chunks) for (const m of ch.messages) if (!chunkMember.has(m.id)) { chunkMember.add(m.id); members.push(m.id); } + const seen = new Set(chunkMember); + for (const id of this.liveL1Ids()) { + const s = this.summaryById(id); + if (!s) continue; + for (const leaf of s.sourceIds) if (!seen.has(leaf)) { seen.add(leaf); members.push(leaf); } + } + const order = this._storeOrder; + if (order.size === 0) return new Map(members.map((id, i) => [id, i] as const)); + const known = members.filter((id) => order.has(id)).sort((a, b) => order.get(a)! - order.get(b)!); + const index = new Map(); + let position = 0; + for (const id of known) index.set(id, position++); + // A chunk member the listing lacks was appended since that listing: it + // is newer than everything listed. An L1 source the listing lacks is + // hidden (viewFilter) or pruned: it occupies no position at all. + for (const id of members) if (!index.has(id) && chunkMember.has(id)) index.set(id, position++); + return index; + } + + private _summaryIndex: { source: readonly SummaryEntry[]; length: number; byId: Map } | null = null; + private summaryById(id: string): SummaryEntry | undefined { + const cached = this._summaryIndex; + if (!cached || cached.source !== this.summaries || cached.length !== this.summaries.length) { + this._summaryIndex = { source: this.summaries, length: this.summaries.length, byId: new Map(this.summaries.map((s) => [s.id, s] as const)) }; + } + return this._summaryIndex!.byId.get(id); + } + + /** + * The L1s that are ownership: those a chunk record points at. Legacy + * archives keep superseded L1 generations (prefix families the migration + * sweep skipped); they are prose, not coverage. Without any chunk→L1 link + * (a strategy driven without records) every L1 counts. + */ + protected liveL1Ids(): Set { + const linked = new Set(); + for (const ch of this.chunks) if (ch.summaryId) linked.add(ch.summaryId); + if (linked.size > 0) return linked; + const all = new Set(); + for (const s of this.summaries) if (s.level === 1) all.add(s.id); + return all; + } + + /** Split same-level sources into strictly adjacent runs (store order). Sources + * whose range endpoints have no position are left out of every run. */ + protected contiguousSourceRuns( + sources: readonly SummaryEntry[], + position: ReadonlyMap, + ): SummaryEntry[][] { + const placed: Array<{ s: SummaryEntry; first: number; last: number }> = []; + for (const s of sources) { + const a = position.get(s.sourceRange.first); + const b = position.get(s.sourceRange.last); + if (a === undefined || b === undefined) continue; + placed.push({ s, first: Math.min(a, b), last: Math.max(a, b) }); + } + placed.sort((x, y) => x.first - y.first); + const runs: SummaryEntry[][] = []; + let run: SummaryEntry[] = []; + let end = -Infinity; + for (const p of placed) { + if (run.length > 0 && p.first !== end + 1) { runs.push(run); run = []; } + run.push(p.s); + end = Math.max(end, p.last); + } + if (run.length > 0) runs.push(run); + return runs; + } + + /** Why this group must not be merged into an L{targetLevel}, or null when it may. */ + protected mergeTopologyProblem( + targetLevel: number, + sources: readonly SummaryEntry[], + position: ReadonlyMap, + ): string | null { + if (sources.length < 2) return `${sources.length} source(s); a merge needs at least two`; + const wrongLevel = sources.filter((s) => s.level !== targetLevel - 1); + if (wrongLevel.length > 0) { + return `source level mismatch: ${wrongLevel.map((s) => `${s.id}=L${s.level}`).join(',')} for an L${targetLevel} merge`; + } + const unplaced = sources.filter((s) => + position.get(s.sourceRange.first) === undefined || position.get(s.sourceRange.last) === undefined); + if (unplaced.length > 0) { + return `source position unresolved (pruned or unowned range endpoints): ${unplaced.map((s) => s.id).join(',')}`; + } + const runs = this.contiguousSourceRuns(sources, position); + if (runs.length !== 1) { + return `sources are not adjacent in store order: ` + + runs.map((r) => `[${r.map((s) => `${s.id}(${s.sourceRange.first}..${s.sourceRange.last})`).join(',')}]`).join(' | '); + } + return null; + } + + /** The merge is never executed: no model call, sources stay unmerged, the + * entry moves into the durable merge quarantine, health goes critical. */ + protected refuseCrossedMerge(targetLevel: SummaryLevel, sourceIds: string[], reason: string): void { + this.requireBranchMutation('refuseCrossedMerge'); + this.topologyRefusals++; + const head = this.mergeQueue[0]; + const entry = head && head.sourceIds === sourceIds ? head : { level: targetLevel, sourceIds, attempts: 0 }; + if (entry === head) this.dequeueMerge(); + const record: MergeQuarantineRecord = { + key: sha256Json(sourceIds), + level: targetLevel, + sourceIds: [...sourceIds], + attempts: entry.attempts ?? 0, + lastOutcome: 'topology_violation', + lastErrorType: reason, + quarantinedAt: Date.now(), + }; + this.mergeQuarantine.set(record.key, record); + this.persistMergeQuarantine(); + console.error( + `[merge-topology] ⛔ refused L${targetLevel} merge over ${sourceIds.length} source(s) ` + + `(${sourceIds.join(', ')}): ${reason}. The node was NOT minted; entry quarantined ` + + `(key=${record.key.slice(0, 12)}). Inspect the store, then clearMergeQuarantine.`, + ); + logCompressionCall({ + event: 'merge:topology-refused', + operation: `merge_l${targetLevel}`, + metadata: { ...record }, + }); + } + + /** + * Leaf-level ownership audit over the whole summary archive: every summary + * whose leaves are not contiguous among chunk-owned messages in store + * order. Read-only; safe to call on any loaded strategy. + */ + auditStoreTopology(messages: ReadonlyArray<{ id: MessageId }>): TopologyViolation[] { + this.refreshStoreOrder(messages); + const position = this.mergePositionIndex(); + if (position.size === 0) return []; + const byPosition: MessageId[] = []; + for (const [id, p] of position) byPosition[p] = id; + const byId = new Map(this.summaries.map((s) => [s.id, s] as const)); + const live = this.liveL1Ids(); + const l1Owner = new Map(); + for (const s of this.summaries) if (s.level === 1 && live.has(s.id)) for (const id of s.sourceIds) l1Owner.set(id, s.id); + const leaves = new Map(); + const collect = (s: SummaryEntry, trail: Set): MessageId[] => { + const cached = leaves.get(s.id); + if (cached) return cached; + if (trail.has(s.id)) return []; // cyclic sourceIds: reported by the pyramid checks, not here + trail.add(s.id); + let out: MessageId[]; + // A superseded L1 generation owns nothing (see liveL1Ids). + if (s.sourceLevel === 0 || s.level === 1) out = live.has(s.id) ? [...s.sourceIds] : []; + else { + out = []; + for (const childId of s.sourceIds) { + const child = byId.get(childId); + if (child) out.push(...collect(child, trail)); + } + } + trail.delete(s.id); + leaves.set(s.id, out); + return out; + }; + const violations: TopologyViolation[] = []; + for (const s of this.summaries) { + const owned = new Set(); + for (const id of collect(s, new Set())) { + const p = position.get(id); + if (p !== undefined) owned.add(p); + } + if (owned.size < 2) continue; + let min = Infinity, max = -Infinity; + for (const p of owned) { if (p < min) min = p; if (p > max) max = p; } + const holes = max - min + 1 - owned.size; + if (holes === 0) continue; + const holeSample: string[] = []; + const holeOwners = new Set(); + for (let p = min; p <= max && (holeSample.length < 5 || holeOwners.size < 5); p++) { + if (owned.has(p)) continue; + const id = byPosition[p]; + if (holeSample.length < 5) holeSample.push(id); + const owner = l1Owner.get(id); + if (owner && holeOwners.size < 5) holeOwners.add(owner); + } + violations.push({ + id: s.id, level: s.level, leafCount: owned.size, + span: { first: byPosition[min], last: byPosition[max] }, + holes, holeSample, holeOwners: [...holeOwners], + }); + } + return violations; + } + + /** Violations found by the load-time audit (empty under 'reject', which throws instead). */ + getTopologyViolations(): readonly TopologyViolation[] { + return this.topologyViolations; + } + + protected resolvedTopologyPolicy(): 'reject' | 'report' { + if (this.config.topologyPolicy) return this.config.topologyPolicy; + const kv = this.config.kvUnified; + if (kv && (kv.preserveGapBearingSummaries || kv.treeifyNonContiguousSummaries)) return 'report'; + return 'reject'; + } + + /** Load-time gate: audit, then throw (reject) or record + shout (report). */ + protected assertStoreTopology(messages: ReadonlyArray<{ id: MessageId }>): void { + const violations = this.auditStoreTopology(messages); + this.topologyViolations = violations; + if (violations.length === 0) return; + const policy = this.resolvedTopologyPolicy(); + const detail = violations.slice(0, 20).map((v) => + ` L${v.level} ${v.id}: ${v.leafCount} leaves ${v.span.first}..${v.span.last}, ${v.holes} hole(s)` + + ` e.g. ${v.holeSample.join(',')}${v.holeOwners.length ? ` owned by ${v.holeOwners.join(',')}` : ''}`).join('\n'); + console.error( + `[store-topology] ⛔ ${violations.length} summar${violations.length === 1 ? 'y owns' : 'ies own'} ` + + `non-contiguous leaves (policy=${policy}):\n${detail}` + + (violations.length > 20 ? `\n … +${violations.length - 20} more` : ''), + ); + logCompressionCall({ + event: 'store-topology-violation', + policy, + count: violations.length, + violations: violations.slice(0, 50), + }); + if (policy === 'reject') throw new StoreTopologyError(violations); + } + /** Drop persisted queue entries authored under an older grouping grammar. * A queue is intent, not memory: if its sources are now parented, missing, * out of order, or separated by another live representation, replaying it @@ -4070,15 +4376,9 @@ export class AutobiographicalStrategy implements ResettableStrategy { } } - // Sequence index per message id, for "within range" tests. Use the - // current chunk store as the ordering source. - const messageOrder = new Map(); - let seq = 0; - for (const ch of this.chunks) { - for (const m of ch.messages) { - messageOrder.set(m.id, seq++); - } - } + // Position per chunk-owned message id in store order (see + // mergePositionIndex), for "within range" and contiguity tests. + const messageOrder = this.mergePositionIndex(); const firstSeq = messageOrder.get(firstMsgId); const lastSeq = messageOrder.get(lastMsgId); @@ -4100,8 +4400,29 @@ export class AutobiographicalStrategy implements ResettableStrategy { } if (sources.length < 2) return; + // The demand path used to take whatever unmerged sources fell inside the + // range, holes included (issue #95): a source whose neighbour's L1 had + // not landed yet was folded across it, minting a crossed node. Group by + // the same strict-adjacency grammar as the threshold path and take the + // longest contiguous run; the rest waits for its neighbours. + const runs = this.contiguousSourceRuns(sources, messageOrder); + const run = runs.reduce((best, r) => (r.length > best.length ? r : best), []); + if (run.length < 2) { + console.warn( + `[autobiographical] demand L${targetLevel} merge over ${sources.length} source(s) in ` + + `${firstMsgId}..${lastMsgId} has no contiguous pair; nothing enqueued`, + ); + return; + } + if (run.length !== sources.length) { + const left = sources.filter((s) => !run.includes(s)).map((s) => s.id); + console.warn( + `[autobiographical] demand L${targetLevel} merge split by a hole: enqueuing ` + + `${run.map((s) => s.id).join(',')}; ${left.join(',')} wait for contiguous neighbours`, + ); + } const N = this.config.mergeThreshold ?? 6; - const toMerge = sources.slice(0, N); + const toMerge = run.slice(0, N); this.enqueueMerge({ level: targetLevel as SummaryLevel, sourceIds: toMerge.map((s) => s.id), @@ -4477,6 +4798,10 @@ export class AutobiographicalStrategy implements ResettableStrategy { compressionQuarantineCount: number; unmergedFrontier: { l1: number; l2: number; l3: number }; lastMintAt: number | null; + /** Summaries with crossed ownership found by the load-time audit (topologyPolicy 'report'). */ + topologyViolations: number; + /** Merges refused this process because they would have minted a crossed node. */ + topologyRefusals: number; } { const DEGRADED_AFTER_MS = 60 * 60 * 1000; const CRITICAL_AFTER_MS = 6 * 60 * 60 * 1000; @@ -4490,6 +4815,8 @@ export class AutobiographicalStrategy implements ResettableStrategy { compressionQuarantineCount: 0, unmergedFrontier: { l1: 0, l2: 0, l3: 0 }, lastMintAt: null, + topologyViolations: 0, + topologyRefusals: 0, }; try { // Exclude the trailing open chunk: it is life, not debt. @@ -4535,6 +4862,9 @@ export class AutobiographicalStrategy implements ResettableStrategy { ) { state = 'critical'; } + // Crossed ownership is never a matter of waiting: it is critical until + // an operator repairs the store (or the refused merge group). + if (this.topologyViolations.length > 0 || this.topologyRefusals > 0) state = 'critical'; return { state, pendingChunks: pending.length, @@ -4549,6 +4879,8 @@ export class AutobiographicalStrategy implements ResettableStrategy { l3: this.summaries.filter((s) => s.level === 3 && !s.mergedInto).length, }, lastMintAt, + topologyViolations: this.topologyViolations.length, + topologyRefusals: this.topologyRefusals, }; } catch { return empty; @@ -6803,11 +7135,10 @@ export class AutobiographicalStrategy implements ResettableStrategy { threshold: number, ): SummaryEntry[] | null { if (unmerged.length < threshold) return null; - const messageOrder = new Map(); - let seq = 0; - for (const ch of this.chunks) { - for (const m of ch.messages) messageOrder.set(m.id, seq++); - } + // Store order, not chunk-record order: a chunk minted late over an early + // message (issue #122's head ratchet) is adjacent to the frontier in + // record order and to the chronicle's opening in store order. + const messageOrder = this.mergePositionIndex(); const spanBase = this.config.mergeMaxSourceSpanMessages ?? 1500; const mergeK = this.config.mergeThreshold ?? 6; const withPos: Array<{ s: SummaryEntry; first: number; last: number }> = []; @@ -7012,6 +7343,17 @@ export class AutobiographicalStrategy implements ResettableStrategy { return; } + // Last line of defence: a crossed node is never minted. Whatever enqueued + // this group (threshold pass, demand path, a persisted queue from an older + // grammar, an operator), the sources must be one level below the target + // and strictly adjacent among chunk-owned messages in store order. + this.refreshStoreOrder(ctx.messageStore.getAll()); + const topology = this.mergeTopologyProblem(targetLevel, sources, this.mergePositionIndex()); + if (topology) { + this.refuseCrossedMerge(targetLevel, sourceIds, topology); + return; + } + const targetTokens = this.config.summaryTargetTokens ?? 2000; const participant = this.config.summaryParticipant ?? 'Claude'; const mergeAttempts = @@ -7713,6 +8055,7 @@ export class AutobiographicalStrategy implements ResettableStrategy { let _t = _diag ? Date.now() : 0; this.loadCalibration(store); const messages = store.getAll(); + this.refreshStoreOrder(messages); if (_diag) { console.error(`[cm-cache] selectAdaptive: calibration+getAll ${Date.now() - _t}ms`); _t = Date.now(); } const msgCap = this.config.maxMessageTokens; // Post-strip estimates (see postStripEstimates): every budgeting site in @@ -10167,7 +10510,9 @@ export class AutobiographicalStrategy implements ResettableStrategy { // ---- 1. Materialize persisted records (they OWN their messages). ---- const byId = new Map(); - for (const m of store.getAll()) byId.set(m.id, m); + const listing = store.getAll(); + this.refreshStoreOrder(listing); + for (const m of listing) byId.set(m.id, m); const consumed = new Set(); let orphaned = 0; diff --git a/src/types/strategy.ts b/src/types/strategy.ts index 0042345..0530a37 100644 --- a/src/types/strategy.ts +++ b/src/types/strategy.ts @@ -885,6 +885,23 @@ export interface AutobiographicalConfig { * (never silently retried in a loop, never canonized). Default: 5. */ mergeAttemptLimit?: number; + /** + * Store-topology policy. On every load the summary archive is audited for + * crossed ownership: a summary whose leaves are not contiguous among the + * chunk-owned messages in store order (an interleaved live representation + * sits inside its span — issue #122's cross-era merges, restore/branch + * interleavings, hand surgery). `'reject'` (default) throws + * `StoreTopologyError` from `initialize`, so `ContextManager.open` fails + * and the operator repairs the store before the resident runs on it. + * `'report'` logs the violations at error level and reports them through + * `getCompressionDebt().topologyViolations` (state `critical`). + * A kv-unified config that explicitly opts into gap handling + * (`preserveGapBearingSummaries` or `treeifyNonContiguousSummaries`) + * defaults to `'report'`: those stores are known to carry gaps. + * Independent of this policy, a merge that WOULD mint a crossed node is + * never executed: it is refused into the merge quarantine. + */ + topologyPolicy?: 'reject' | 'report'; /** Legacy first-choice target-only merge request. Default false. */ compressionMergeSourceOnly?: boolean; /** Use target-only merge request only on the final persisted merge attempt. Default false. */ diff --git a/test/topology-guard.test.ts b/test/topology-guard.test.ts new file mode 100644 index 0000000..f80ec2c --- /dev/null +++ b/test/topology-guard.test.ts @@ -0,0 +1,293 @@ +/** + * Store topology guard — merge adjacency in STORE order, refusal of crossed + * mints, and the load-time audit (issues #122 and #95). + */ + +import { test, describe, before, after } from 'node:test'; +import assert from 'node:assert/strict'; +import { rmSync, existsSync } from 'node:fs'; + +import { ContextManager, StoreTopologyError } from '../src/index.js'; +import { AutobiographicalStrategy } from '../src/strategies/autobiographical.js'; +import type { SummaryEntry, StrategyContext } from '../src/types/strategy.js'; +import type { ContentBlock } from '@animalabs/membrane'; + +class Probe extends AutobiographicalStrategy { + /** Chunk records in RECORD order (the order chunks were minted). */ + setChunks(groups: string[][]): void { + (this as unknown as { chunks: unknown[] }).chunks = groups.map((ids) => ({ messages: ids.map((id) => ({ id })) })); + } + /** The store listing (chronicle order). */ + setStore(messageIds: string[]): void { + this.refreshStoreOrder(messageIds.map((id) => ({ id }))); + } + setSummaries(summaries: SummaryEntry[]): void { + (this as unknown as { summaries: SummaryEntry[] }).summaries = summaries; + } + allSummaries(): SummaryEntry[] { + return (this as unknown as { summaries: SummaryEntry[] }).summaries; + } + pick(unmerged: SummaryEntry[], threshold: number): SummaryEntry[] | null { + return this.contiguousMergeCandidates(unmerged, threshold); + } + demand(level: number, first: string, last: string): void { + this.enqueueMergeForRange(level, first, last); + } + queue(): Array<{ level: number; sourceIds: string[] }> { + return (this as unknown as { mergeQueue: Array<{ level: number; sourceIds: string[] }> }).mergeQueue; + } + setQueue(queue: Array<{ level: number; sourceIds: string[] }>): void { + (this as unknown as { mergeQueue: unknown[] }).mergeQueue = queue; + } + quarantine(): Map { + return (this as unknown as { mergeQuarantine: Map }).mergeQuarantine; + } + audit(messageIds: string[]): ReturnType { + return this.auditStoreTopology(messageIds.map((id) => ({ id }))); + } + gate(messageIds: string[]): void { + this.assertStoreTopology(messageIds.map((id) => ({ id }))); + } + async merge(level: number, sourceIds: string[], ctx: StrategyContext): Promise { + await this.executeMerge(level as never, sourceIds, ctx); + } +} + +const ids = (from: number, to: number): string[] => Array.from({ length: to - from + 1 }, (_, i) => `m-${from + i}`); + +function l1(id: string, first: number, last: number, level = 1): SummaryEntry { + return { + id, level, content: `s ${id}`, tokens: 100, sourceLevel: level - 1, + sourceIds: level === 1 ? ids(first, last) : [], + sourceRange: { first: `m-${first}`, last: `m-${last}` }, + created: 1, + } as SummaryEntry; +} + +function probe(): Probe { + return new Probe({ adaptiveResolution: true, hierarchical: true, autoTickOnNewMessage: false, mergeThreshold: 6 }); +} + +describe('merge adjacency is judged in store order (issue #122)', () => { + test('a late chunk over an opening message does not join the frontier run', () => { + const p = probe(); + // Store: m-0..m-199. Chunk RECORDS: the frontier chunks first, then the + // stray minted last over the chronicle's opening message m-3. + p.setStore(ids(0, 199)); + // m-4..m-99 are owned by earlier (already merged) chunks: the stray's only + // neighbours in store order are other live representations. + const history = [ids(4, 49), ids(50, 99)]; + const frontier = [ids(100, 109), ids(110, 119), ids(120, 129), ids(130, 139), ids(140, 149)]; + p.setChunks([...history, ...frontier, ['m-3']]); + const stray = l1('L1-stray', 3, 3); + const run5 = frontier.map((g, i) => l1(`L1-f${i}`, 100 + i * 10, 109 + i * 10)); + // In record order the stray is adjacent to L1-f4 and the six would merge. + // In store order it is an interior singleton: the frontier run has five + // members and waits for its sixth. + assert.equal(p.pick([...run5, stray], 6), null); + // With a sixth frontier L1 the run merges WITHOUT the stray. + p.setChunks([...history, ...frontier, ids(150, 159), ['m-3']]); + const six = [...run5, l1('L1-f5', 150, 159)]; + const picked = p.pick([...six, stray], 6); + assert.deepEqual(picked?.map((s) => s.id), six.map((s) => s.id)); + }); + + test('without a store listing the index degrades to chunk-record order', () => { + const p = probe(); + p.setChunks([ids(0, 9), ids(10, 19)]); + const picked = p.pick([l1('a', 0, 9), l1('b', 10, 19)], 2); + assert.deepEqual(picked?.map((s) => s.id), ['a', 'b']); + }); +}); + +/** A Probe bound to a real store: branch-guarded entrypoints need one. */ +async function openProbe(path: string, membrane: unknown, count = 40): Promise<{ manager: ContextManager; p: Probe; ids: string[]; ctx: StrategyContext }> { + const p = new Probe({ adaptiveResolution: true, hierarchical: true, autoTickOnNewMessage: false, mergeThreshold: 6, compressionModel: 'test-compression-model' }); + const manager = await ContextManager.open({ path, strategy: p, membrane: membrane as never }); + const out: string[] = []; + for (let i = 0; i < count; i++) out.push(manager.addMessage(i % 2 ? 'agent' : 'user', [{ type: 'text', text: `message ${i} ` + 'w '.repeat(20) }])); + const ctx = (manager as unknown as { createStrategyContext(): StrategyContext }).createStrategyContext(); + return { manager, p, ids: out, ctx }; +} +const real = (all: string[]) => (a: number, b: number) => all.slice(a, b + 1); +function l1r(id: string, src: string[], level = 1): SummaryEntry { + return { id, level, content: `s ${id}`, tokens: 100, sourceLevel: level - 1, sourceIds: src, + sourceRange: { first: src[0], last: src[src.length - 1] }, created: 1 } as SummaryEntry; +} +const STORES = ['./test-topology-guard-demand', './test-topology-guard-refuse', './test-topology-guard-pass', './test-topology-guard-store']; +const wipe = () => { for (const s of STORES) if (existsSync(s)) rmSync(s, { recursive: true, force: true }); }; + +describe('demand-path merges respect holes (issue #95)', () => { + before(wipe); + after(wipe); + test('a source separated by an unlanded neighbour is left out; no contiguous pair enqueues nothing', async () => { + const membrane = { complete: async () => ({ stopReason: 'end_turn', content: [{ type: 'text', text: 'x' }] }) }; + const { manager, p, ids: all } = await openProbe(STORES[0], membrane); + const s = real(all); + p.setChunks([s(0, 9), s(10, 19), s(20, 29), s(30, 39)]); + // The L1 over 20..29 has not landed: folding a and d across it would cross. + p.setSummaries([l1r('a', s(0, 9)), l1r('b', s(10, 19)), l1r('d', s(30, 39))]); + p.demand(2, all[0], all[39]); + assert.deepEqual(p.queue().map((m) => m.sourceIds), [['a', 'b']]); + p.setQueue([]); + p.setSummaries([l1r('a', s(0, 9)), l1r('d', s(30, 39))]); + p.demand(2, all[0], all[39]); + assert.deepEqual(p.queue(), []); + manager.close(); + }); +}); + +describe('a crossed merge is refused at mint time', () => { + before(wipe); + after(wipe); + test('non-adjacent sources: no model call, entry quarantined, health critical', async () => { + let calls = 0; + const membrane = { complete: async () => { calls++; return { stopReason: 'end_turn', content: [{ type: 'text', text: 'merged' }] }; } }; + const { manager, p, ids: all, ctx } = await openProbe(STORES[1], membrane); + const s = real(all); + p.setChunks([s(0, 9), s(10, 19), s(20, 29)]); + p.setSummaries([l1r('a', s(0, 9)), l1r('b', s(10, 19)), l1r('c', s(20, 29))]); + const queued = { level: 2, sourceIds: ['a', 'c'] }; + p.setQueue([queued]); + await p.merge(2, queued.sourceIds, ctx); + assert.equal(calls, 0, 'the summarizer must not be called for a crossed group'); + assert.deepEqual(p.queue(), [], 'the entry left the queue'); + const [record] = [...p.quarantine().values()]; + assert.equal(record.lastOutcome, 'topology_violation'); + assert.match(record.lastErrorType ?? '', /not adjacent in store order/); + const debt = p.getCompressionDebt(); + assert.equal(debt.topologyRefusals, 1); + assert.equal(debt.state, 'critical'); + assert.equal(p.allSummaries().filter((x) => x.level === 2).length, 0, 'nothing was minted'); + manager.close(); + }); + + test('adjacent sources pass the gate (reaches the summarizer)', async () => { + let calls = 0; + const membrane = { complete: async () => { calls++; return { stopReason: 'end_turn', content: [{ type: 'text', text: 'merged '.repeat(10) }], usage: { input_tokens: 1, output_tokens: 1 } }; } }; + const { manager, p, ids: all, ctx } = await openProbe(STORES[2], membrane); + const s = real(all); + p.setChunks([s(0, 9), s(10, 19)]); + p.setSummaries([l1r('a', s(0, 9)), l1r('b', s(10, 19))]); + p.setQueue([{ level: 2, sourceIds: ['a', 'b'] }]); + await p.merge(2, ['a', 'b'], ctx); + assert.equal(calls, 1); + assert.equal(p.getCompressionDebt().topologyRefusals, 0); + manager.close(); + }); +}); + +describe('load-time topology audit', () => { + function crossedFixture(): { p: Probe; store: string[] } { + const p = probe(); + const store = ids(0, 49); + p.setChunks([['m-3'], ids(4, 13), ids(14, 23), ids(24, 33), ids(34, 43)]); + const stray = l1('L1-stray', 3, 3); + const era = [l1('L1-a', 4, 13), l1('L1-b', 14, 23), l1('L1-c', 24, 33), l1('L1-d', 34, 43)]; + // L2-x owns the stray opening message plus m-24..43: crossed over L1-a/L1-b. + const l2 = { ...l1('L2-x', 3, 43, 2), sourceIds: ['L1-stray', 'L1-c', 'L1-d'] } as SummaryEntry; + stray.mergedInto = 'L2-x'; era[2].mergedInto = 'L2-x'; era[3].mergedInto = 'L2-x'; + p.setSummaries([stray, ...era, l2]); + return { p, store }; + } + + test('finds the crossed summary and names the interleaved owners', () => { + const { p, store } = crossedFixture(); + const violations = p.audit(store); + assert.deepEqual(violations.map((v) => v.id), ['L2-x']); + const [v] = violations; + assert.equal(v.level, 2); + assert.equal(v.leafCount, 21); + assert.deepEqual(v.span, { first: 'm-3', last: 'm-43' }); + assert.equal(v.holes, 20); + assert.deepEqual(v.holeOwners, ['L1-a', 'L1-b']); + }); + + test('a clean pyramid has no violations', () => { + const p = probe(); + p.setChunks([ids(0, 9), ids(10, 19), ids(20, 29)]); + const l1s = [l1('a', 0, 9), l1('b', 10, 19), l1('c', 20, 29)]; + const l2 = { ...l1('L2', 0, 19, 2), sourceIds: ['a', 'b'] } as SummaryEntry; + p.setSummaries([...l1s, l2]); + assert.deepEqual(p.audit(ids(0, 40)), []); + }); + + test("never-chunked messages inside a span are not holes", () => { + const p = probe(); + // m-10 is unowned (e.g. a message the chunker skipped): it occupies no position. + p.setChunks([ids(0, 9), ids(11, 20)]); + const l2 = { ...l1('L2', 0, 20, 2), sourceIds: ['a', 'b'] } as SummaryEntry; + p.setSummaries([l1('a', 0, 9), l1('b', 11, 20), l2]); + assert.deepEqual(p.audit(ids(0, 20)), []); + }); + + test("'reject' throws StoreTopologyError; 'report' records and goes critical", () => { + const { p, store } = crossedFixture(); + assert.throws(() => p.gate(store), (e: unknown) => e instanceof StoreTopologyError && e.violations.length === 1 && /L2-x/.test(e.message)); + const r = new Probe({ adaptiveResolution: true, hierarchical: true, autoTickOnNewMessage: false, topologyPolicy: 'report' }); + const { p: fixture } = crossedFixture(); + r.setChunks((fixture as unknown as { chunks: Array<{ messages: Array<{ id: string }> }> }).chunks.map((c) => c.messages.map((m) => m.id))); + r.setSummaries(fixture.allSummaries()); + r.gate(store); + assert.equal(r.getTopologyViolations().length, 1); + assert.equal(r.getCompressionDebt().topologyViolations, 1); + assert.equal(r.getCompressionDebt().state, 'critical'); + }); + + test('explicit kv-unified gap handling defaults to report', () => { + const { p: fixture, store } = crossedFixture(); + const kv = new Probe({ + adaptiveResolution: true, hierarchical: true, autoTickOnNewMessage: false, + foldingStrategy: 'kv-unified', + kvUnified: { treeifyNonContiguousSummaries: false, preserveGapBearingSummaries: true } as never, + }); + kv.setChunks((fixture as unknown as { chunks: Array<{ messages: Array<{ id: string }> }> }).chunks.map((c) => c.messages.map((m) => m.id))); + kv.setSummaries(fixture.allSummaries()); + kv.gate(store); + assert.equal(kv.getTopologyViolations().length, 1); + }); +}); + +describe('ContextManager.open fails closed on a crossed store', () => { + const STORE = STORES[3]; + before(wipe); + after(wipe); + const t = (s: string): ContentBlock[] => [{ type: 'text', text: s }]; + + test('a persisted crossed L2 rejects open; report opens and surfaces it', async () => { + const membrane = { complete: async () => ({ stopReason: 'end_turn', content: [{ type: 'text', text: 'x' }] }) }; + const seed = new AutobiographicalStrategy({ adaptiveResolution: true, hierarchical: true, autoTickOnNewMessage: false }); + const manager = await ContextManager.open({ path: STORE, strategy: seed, membrane: membrane as never }); + const msgIds: string[] = []; + for (let i = 0; i < 40; i++) msgIds.push(manager.addMessage(i % 2 ? 'agent' : 'user', t(`message ${i} ` + 'w '.repeat(20)))); + const mk = (id: string, level: number, src: string[], first: string, last: string, mergedInto?: string): SummaryEntry => ({ + id, level, content: `s ${id}`, tokens: 50, sourceLevel: level - 1, sourceIds: src, + sourceRange: { first, last }, created: 1, ...(mergedInto ? { mergedInto } : {}), + } as SummaryEntry); + const slice = (a: number, b: number) => msgIds.slice(a, b + 1); + const summaries = [ + mk('L1-stray', 1, slice(3, 3), msgIds[3], msgIds[3], 'L2-x'), + mk('L1-a', 1, slice(4, 13), msgIds[4], msgIds[13]), + mk('L1-b', 1, slice(14, 23), msgIds[14], msgIds[23]), + mk('L1-c', 1, slice(24, 33), msgIds[24], msgIds[33], 'L2-x'), + mk('L2-x', 2, ['L1-stray', 'L1-c'], msgIds[3], msgIds[33]), + ]; + const internals = seed as unknown as { store: { setStateJson(id: string, v: unknown): void }; summariesStateId: string }; + internals.store.setStateJson(internals.summariesStateId, summaries); + manager.close(); + + await assert.rejects( + ContextManager.open({ path: STORE, strategy: new AutobiographicalStrategy({ adaptiveResolution: true, hierarchical: true, autoTickOnNewMessage: false }), membrane: membrane as never }), + (e: unknown) => e instanceof StoreTopologyError && e.violations.map((v) => v.id).join() === 'L2-x', + ); + + const report = new AutobiographicalStrategy({ adaptiveResolution: true, hierarchical: true, autoTickOnNewMessage: false, topologyPolicy: 'report' }); + const reopened = await ContextManager.open({ path: STORE, strategy: report, membrane: membrane as never }); + try { + assert.deepEqual(report.getTopologyViolations().map((v) => v.id), ['L2-x']); + assert.equal(report.getCompressionDebt().state, 'critical'); + } finally { + reopened.close(); + } + }); +});