From b381460e602d4bc3d803dbd8a101804220199dc3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20Magyar?= Date: Sat, 3 Oct 2026 15:52:58 +0100 Subject: [PATCH] fix(embeddings): recover durable vectors after interrupted analysis --- gitnexus/src/core/run-analyze.ts | 184 ++++++++++++++++++++++++++++--- 1 file changed, 167 insertions(+), 17 deletions(-) diff --git a/gitnexus/src/core/run-analyze.ts b/gitnexus/src/core/run-analyze.ts index 92925de71..dc0dc88e7 100644 --- a/gitnexus/src/core/run-analyze.ts +++ b/gitnexus/src/core/run-analyze.ts @@ -36,7 +36,12 @@ import fs from 'fs/promises'; import { constants as fsConstants, existsSync } from 'node:fs'; import { randomUUID } from 'node:crypto'; import { retryRename } from '../storage/fs-atomic.js'; -import { acquireIndexLock, requireExclusiveIndexLock } from '../storage/index-lock.js'; +import { + acquireIndexLock, + requireExclusiveIndexLock, + sweepStagingArtifacts, +} from '../storage/index-lock.js'; +import { resolveEmbeddingRecovery } from '../storage/embedding-recovery.js'; import { invalidateNodeWorkspacePackages } from './ingestion/import-resolvers/node-workspace-packages.js'; import { logNameFallbackSummary, @@ -1287,6 +1292,8 @@ export async function runFullAnalysis( const log = (msg: string) => callbacks.onLog?.(stripControlCharacters(msg)); const acquireOpts = { log, + // Resolve and validate the canonical slot under the lock before cleanup. + sweep: false, onWaitStart: () => callbacks.onProgress('lock', 0, 'Waiting for another analyze to finish on this index…'), }; @@ -1345,6 +1352,7 @@ export async function runFullAnalysis( } const flatShared = writeTarget.placement.branch ? undefined : writeTarget.sharedStore; if (flatShared) await seedSharedSlot(flatShared, repoPath, log); + sweepStagingArtifacts(writeTarget.metaDir, log); const slotToLeave = options.noShare ? await optedInSlotToLeave(repoPath) : undefined; const result = await runFullAnalysisInner( repoPath, @@ -1842,7 +1850,13 @@ async function runFullAnalysisInner( decision = decideEmbeddingResume(checkpoint, embeddingIdentityForRun, resumeOptions); } if (decision.action === 'abort') throw new Error(decision.error); - log(decision.log); + log( + decision.action === 'resume' && checkpoint.recovery + ? `Previous analyze recorded an embedding checkpoint (${checkpoint.nodesProcessed}/` + + `${checkpoint.totalNodes} nodes); validating staged vectors before retrying ` + + `${decision.pendingNodeIds.size} pending node(s).` + : decision.log, + ); if (options.dropEmbeddings) { // --drop-embeddings has always implied a rebuild here; the decision only // covers the marker. @@ -2671,6 +2685,74 @@ async function runFullAnalysisInner( } } + // A checkpoint's pending decision and its paid, complete vectors are + // independent: --force discards the former but can still reuse the latter. + // Select only the explicitly referenced generation, never an orphan by age. + const stagedRecovery = resolveEmbeddingRecovery(metaDir, existingMeta?.embeddingCheckpoint); + const stagedCheckpoint = existingMeta?.embeddingCheckpoint; + const inheritedUnsafeNodeIds = new Set([ + ...pendingEmbeddingNodeIds, + ...(stagedRecovery?.unsafeNodeIds ?? []), + ]); + if (shouldLoadCache && !options.dropEmbeddings && stagedRecovery && stagedCheckpoint) { + if (!embeddingIdentityForRun) { + const { resolveEmbeddingIdentity } = await import('./embeddings/embedding-identity.js'); + embeddingIdentityForRun = resolveEmbeddingIdentity(); + } + const marker = stagedCheckpoint; + const matchesIdentity = + marker.model === embeddingIdentityForRun.model && + marker.dimensions === embeddingIdentityForRun.dimensions && + marker.provider === embeddingIdentityForRun.provider; + if (matchesIdentity && stagedRecovery.schemaFingerprint === SCHEMA_FINGERPRINT) { + // A force-discarded pending decision must not turn a known incomplete + // inherited group into a reusable cache merely because its hash matches. + if (inheritedUnsafeNodeIds.size > 0) { + const rows = cachedSnapshot.rows.filter((row) => !inheritedUnsafeNodeIds.has(row.nodeId)); + cachedSnapshot = { + ...cachedSnapshot, + rows, + embeddingNodeIds: new Set(rows.map((row) => row.nodeId)), + }; + } + let recovered: CachedEmbeddingsSnapshot | undefined; + try { + const { recoverStagedEmbeddings, mergeRecoveredEmbeddings } = + await import('./embeddings/staged-embedding-recovery.js'); + recovered = await recoverStagedEmbeddings(stagedRecovery.dbPath, { + dimensions: embeddingIdentityForRun.dimensions, + excludedNodeIds: stagedRecovery.unsafeNodeIds, + }); + const liveCacheDims = snapshotEmbeddingDims(cachedSnapshot); + if (liveCacheDims !== undefined && liveCacheDims !== embeddingIdentityForRun.dimensions) { + log( + `Embedding dimensions changed (${liveCacheDims}d -> ` + + `${embeddingIdentityForRun.dimensions}d), discarding published cache`, + ); + discardCachedEmbeddings(); + } + if (recovered.rows.length > 0) { + const merged = mergeRecoveredEmbeddings(cachedSnapshot, recovered); + disposeEmbeddingSpill(cachedSnapshot.spill); + adoptCachedEmbeddings(merged); + } + log( + `Recovered ${recovered.rows.length} complete staged embedding chunk(s) ` + + `for ${recovered.embeddingNodeIds.size} node(s); unchanged content can reuse them.`, + ); + } catch (err) { + log( + `Warning: could not recover staged embeddings (${(err as Error).message}); ` + + 'the retry will regenerate missing chunks.', + ); + } finally { + disposeEmbeddingSpill(recovered?.spill); + } + } else { + log('Staged embedding identity or schema changed; its vectors will not be reused.'); + } + } + // ── Load incremental parse cache ────────────────────────────────── // Content-addressed: `--force` reuses parser shards; `useParseCache: false` // stages a new generation under a run-unique parse-rebuild.* dir and publishes @@ -4411,6 +4493,9 @@ async function runFullAnalysisInner( // (the clean-run contract). Built inside Phase 4 so it carries the identity // of the run that actually wrote it — see the assignment below (#2790). let pendingEmbeddingCheckpoint: RepoMeta['embeddingCheckpoint']; + // A staged measurement is useful at publication, but cannot describe the + // live database while the replacement is still unpublished (#3456). + let lastMeasuredCheckpointEmbeddingCount: number | undefined; if (shouldGenerateEmbeddings) { const { skipForCap, capDisabled, nodeLimit } = deriveEmbeddingCap( @@ -4486,6 +4571,16 @@ async function runFullAnalysisInner( embeddingIdentityForRun = resolveEmbeddingIdentity(); } const embeddingIdentity = embeddingIdentityForRun; + const stagedRecoveryEnabled = useAtomicSwap && isManualCheckpointEnabled(); + if (useAtomicSwap && !stagedRecoveryEnabled) { + log( + 'Manual WAL checkpoints are disabled; new staged work cannot be recovered after interruption. ' + + 'Any previous durable recovery source is retained until publication.', + ); + } + const unsafeRecoveryNodeIds = new Set([...inheritedUnsafeNodeIds, ...restoreFailedNodeIds]); + let activeWindowNodeIds: string[] = []; + let recoveryGenerationDurable = false; // Build a Map from cached embeddings for incremental mode let existingEmbeddings: Map | undefined; if (cachedSnapshot.embeddingNodeIds.size > 0) { @@ -4521,12 +4616,9 @@ async function runFullAnalysisInner( // re-read the on-disk meta immediately before writing (the shape the // /api/embed checkpoint writer in server/api.ts already uses, which also // keeps a concurrent writer's update from being reverted by a stale - // snapshot) and replace ONLY `embeddingCheckpoint` — plus - // `stats.embeddings` when the caller actually MEASURED the live count - // (the post-window `onCheckpoint`). The window-start callback passes - // nothing: restating the previous run's count there both re-published a - // stale number and clobbered the live count a preceding `onCheckpoint` - // had just written. + // snapshot) and replace ONLY `embeddingCheckpoint`. In-place writers + // can also publish a measured live count; staged writers keep it in + // memory until the replacement is published. const saveEmbeddingCheckpoint = async ( checkpoint: { nodesProcessed: number; @@ -4537,6 +4629,16 @@ async function runFullAnalysisInner( embeddings?: number, ): Promise => { const latestMeta = (await loadMeta(metaDir)) ?? existingMeta; + // The supported manual-checkpoint opt-out cannot prove durability of + // this stage. Keep an earlier proven generation and its original + // identity/progress intact until the replacement is published. + if ( + useAtomicSwap && + !stagedRecoveryEnabled && + resolveEmbeddingRecovery(metaDir, latestMeta?.embeddingCheckpoint) + ) { + return; + } // First-ever analyze of this repo: no meta exists on disk yet (the // pre-wipe dirty stamp only fires when one does). Mint the minimum // RepoMeta requires, with `lastCommit: ''` — never `currentCommit` — @@ -4547,16 +4649,30 @@ async function runFullAnalysisInner( lastCommit: '', indexedAt: new Date().toISOString(), }; + const interrupted = mintInterruptedCheckpoint( + embeddingIdentity, + checkpoint, + pendingNodeIds, + ); await saveMeta(metaDir, { ...base, - ...(embeddings === undefined ? {} : { stats: { ...base.stats, embeddings } }), + ...(embeddings === undefined || useAtomicSwap + ? {} + : { stats: { ...base.stats, embeddings } }), // Written by a run that is still IN FLIGHT — see the `kind` doc in // repo-manager.ts. - embeddingCheckpoint: mintInterruptedCheckpoint( - embeddingIdentity, - checkpoint, - pendingNodeIds, - ), + embeddingCheckpoint: { + ...interrupted, + ...(stagedRecoveryEnabled + ? { + recovery: { + stagingFile: path.basename(buildPath), + schemaFingerprint: SCHEMA_FINGERPRINT, + unsafeNodeIds: [...unsafeRecoveryNodeIds], + }, + } + : {}), + }, }); }; @@ -4579,7 +4695,21 @@ async function runFullAnalysisInner( { forceReembedNodeIds: pendingEmbeddingNodeIds, onCheckpointWindowStart: async ({ nodeIds, ...checkpoint }) => { + const handoff = stagedRecoveryEnabled && !recoveryGenerationDurable; + if (handoff) { + if (!(await checkpointOnce())) { + throw new Error( + 'Could not checkpoint restored embeddings before recovery handoff.', + ); + } + recoveryGenerationDurable = true; + } + activeWindowNodeIds = nodeIds; + for (const id of nodeIds) unsafeRecoveryNodeIds.add(id); await saveEmbeddingCheckpoint(checkpoint, nodeIds); + // Reclaim the old source only after the new durable generation's + // reference is saved. Later windows retain this same generation. + if (handoff) sweepStagingArtifacts(metaDir, log); }, // ── The mid-run count is a DIAGNOSTIC, not a gate (#2790) ────── // This used to run the count query bare. THIS callback's rejection @@ -4593,8 +4723,15 @@ async function runFullAnalysisInner( // touch stats.embeddings" signal — so the checkpoint still lands, // with whatever count is already on disk left alone. onCheckpoint: async (checkpoint) => { - await checkpointOnce(); + const durable = await checkpointOnce(); + if (stagedRecoveryEnabled && !durable) { + throw new Error('Could not checkpoint the completed embedding window for recovery.'); + } + for (const id of activeWindowNodeIds) unsafeRecoveryNodeIds.delete(id); + activeWindowNodeIds = []; const measured = await measurePersistedEmbeddingCount(executeQuery); + const measuredCount = persistedEmbeddingCountOrUndefined(measured); + if (measuredCount !== undefined) lastMeasuredCheckpointEmbeddingCount = measuredCount; if (measured.kind === 'unknown') { log( `Warning: could not measure persisted embeddings at the embedding checkpoint ` + @@ -4759,7 +4896,7 @@ async function runFullAnalysisInner( embeddingCount === undefined ? ((await loadMeta(metaDir)) ?? existingMeta) : undefined; const persistedEmbeddingCount = resolvePersistedEmbeddingCount( measuredEmbeddingCount, - latestMetaForCount?.stats?.embeddings, + lastMeasuredCheckpointEmbeddingCount ?? latestMetaForCount?.stats?.embeddings, ); const { getRuntimeCapabilities } = await import('./platform/capabilities.js'); @@ -5081,6 +5218,7 @@ async function runFullAnalysisInner( // is a crash-safety improvement: a failed swap leaves the previous index // live and the next run recovers via the full-rebuild path. await saveMeta(metaDir, meta); + sweepStagingArtifacts(metaDir, log); // Registry freshness is published only after the graph and its metadata. // A failed close, swap, or metadata save must leave the previous registry @@ -5290,7 +5428,13 @@ async function runFullAnalysisInner( // rethrow below is the surface, and the lock's sweep remains the backstop. if (useAtomicSwap && buildPath !== lbugPath) { try { - await wipeLbugDbFiles(buildPath); + const recovery = resolveEmbeddingRecovery( + metaDir, + (await loadMeta(metaDir))?.embeddingCheckpoint, + ); + // Both paths belong to this locked slot. The validated generation + // basename identifies the same file even through a directory alias. + if (recovery?.stagingFile !== path.basename(buildPath)) await wipeLbugDbFiles(buildPath); } catch { /* swallow — orphan reclamation must never mask the real failure */ } @@ -5302,6 +5446,12 @@ async function runFullAnalysisInner( // IndexLockTimeoutError and other domain failures with `instanceof`. recordLiveIndexMutationRisk(err); } + if (/max(?:imum)?(?: database| db)? size|database size limit|maxDBSize/i.test(String(err))) { + log( + 'The database size limit was reached. Set GITNEXUS_LBUG_MAX_DB_SIZE to a larger ' + + 'byte limit before retrying analyze; retained complete embeddings can be reused.', + ); + } throw err; } }