From 49062f48fcbbb0d51949b81c1cb9912aa2ae7c73 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20Magyar?= Date: Sat, 3 Oct 2026 15:55:40 +0100 Subject: [PATCH] fix(storage): retain checkpoint-referenced stages --- gitnexus/src/storage/embedding-recovery.ts | 118 +++++++++++++++++++++ gitnexus/src/storage/index-lock.ts | 9 +- gitnexus/src/storage/repo-meta.ts | 14 +++ 3 files changed, 138 insertions(+), 3 deletions(-) create mode 100644 gitnexus/src/storage/embedding-recovery.ts diff --git a/gitnexus/src/storage/embedding-recovery.ts b/gitnexus/src/storage/embedding-recovery.ts new file mode 100644 index 000000000..5cfcc0101 --- /dev/null +++ b/gitnexus/src/storage/embedding-recovery.ts @@ -0,0 +1,118 @@ +/** + * Filesystem-only staged embedding provenance. Keep this independent of native + * and model imports: every index-lock caller needs the retention decision. + */ +import { lstatSync, readFileSync, realpathSync } from 'node:fs'; +import path from 'node:path'; +import type { EmbeddingRecoveryReference } from './repo-meta.js'; +import { INDEX_METADATA_FILE, LEGACY_METADATA_FILE } from './storage-constants.js'; + +const STAGING_FILENAME = + /^lbug\.staging\.[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/; +const FAMILY_SUFFIXES = ['', '.wal', '.shadow', '.wal.checkpoint', '.lock'] as const; + +export interface ResolvedEmbeddingRecovery extends EmbeddingRecoveryReference { + dbPath: string; + /** Exact basenames, never a prefix match that can preserve another generation. */ + familyFiles: string[]; +} + +const isRecord = (value: unknown): value is Record => + value !== null && typeof value === 'object' && !Array.isArray(value); +const isNonemptyString = (value: unknown): value is string => + typeof value === 'string' && value.length > 0; +const isNodeIds = (value: unknown): value is string[] => + Array.isArray(value) && value.every(isNonemptyString); +const isCount = (value: unknown): value is number => + typeof value === 'number' && Number.isSafeInteger(value) && value >= 0; + +/** + * Resolve only an explicit interrupted-generation receipt in this canonical + * slot. This validates provenance, not native contents or current identity; + * those checks must pass separately before any rows can be reused. + */ +export const resolveEmbeddingRecovery = ( + lockDir: string, + checkpoint: unknown, +): ResolvedEmbeddingRecovery | undefined => { + if (!isRecord(checkpoint) || !isRecord(checkpoint.recovery)) return undefined; + if (checkpoint.kind !== undefined && checkpoint.kind !== 'interrupted') return undefined; + if ( + !isNonemptyString(checkpoint.at) || + !Number.isFinite(Date.parse(checkpoint.at)) || + !isCount(checkpoint.nodesProcessed) || + !isCount(checkpoint.totalNodes) || + checkpoint.nodesProcessed > checkpoint.totalNodes || + !isCount(checkpoint.chunksProcessed) || + !isNonemptyString(checkpoint.model) || + !isCount(checkpoint.dimensions) || + checkpoint.dimensions === 0 || + !isNonemptyString(checkpoint.provider) || + (checkpoint.pendingNodeIds !== undefined && !isNodeIds(checkpoint.pendingNodeIds)) + ) { + return undefined; + } + const recovery = checkpoint.recovery; + if ( + !isNonemptyString(recovery.stagingFile) || + !STAGING_FILENAME.test(recovery.stagingFile) || + !isNonemptyString(recovery.schemaFingerprint) || + !isNodeIds(recovery.unsafeNodeIds) + ) { + return undefined; + } + const unsafeNodeIds = new Set(recovery.unsafeNodeIds); + if ( + isNodeIds(checkpoint.pendingNodeIds) && + checkpoint.pendingNodeIds.some((nodeId) => !unsafeNodeIds.has(nodeId)) + ) { + return undefined; + } + + try { + const canonicalDir = realpathSync(lockDir); + if (!lstatSync(canonicalDir).isDirectory()) return undefined; + const familyFiles = FAMILY_SUFFIXES.map((suffix) => recovery.stagingFile + suffix); + for (const [index, filename] of familyFiles.entries()) { + try { + // lstat refuses both live and dangling symlinks, without following one + // to a database outside the slot. The base file must exist. + if (!lstatSync(path.join(canonicalDir, filename)).isFile()) return undefined; + } catch (error) { + if (index > 0 && (error as NodeJS.ErrnoException).code === 'ENOENT') continue; + return undefined; + } + } + return { + dbPath: path.join(canonicalDir, recovery.stagingFile), + stagingFile: recovery.stagingFile, + schemaFingerprint: recovery.schemaFingerprint, + unsafeNodeIds: [...unsafeNodeIds], + familyFiles, + }; + } catch { + return undefined; + } +}; + +/** Synchronous mirror of loadMeta's primary-first, absent-only fallback rule. */ +export const readEmbeddingRecovery = (lockDir: string): ResolvedEmbeddingRecovery | undefined => { + let metadataPath = path.join(lockDir, INDEX_METADATA_FILE); + try { + try { + if (!lstatSync(metadataPath).isFile()) return undefined; + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + if (code !== 'ENOENT' && code !== 'ENOTDIR') return undefined; + metadataPath = path.join(lockDir, LEGACY_METADATA_FILE); + if (!lstatSync(metadataPath).isFile()) return undefined; + } + const meta: unknown = JSON.parse(readFileSync(metadataPath, 'utf8')); + if (!isRecord(meta)) return undefined; + // storagePath describes the flat/cache root, including in branch-slot + // metadata. The current locked directory is the generation boundary. + return resolveEmbeddingRecovery(lockDir, meta.embeddingCheckpoint); + } catch { + return undefined; + } +}; diff --git a/gitnexus/src/storage/index-lock.ts b/gitnexus/src/storage/index-lock.ts index 5588127e6..5d3b6cc47 100644 --- a/gitnexus/src/storage/index-lock.ts +++ b/gitnexus/src/storage/index-lock.ts @@ -62,6 +62,7 @@ import path from 'node:path'; import os from 'node:os'; import { randomBytes, randomUUID, createHash } from 'node:crypto'; import { isProcessAlive } from '../utils/process-identity.js'; +import { readEmbeddingRecovery } from './embedding-recovery.js'; const LOCK_FILENAME = 'analyze.lock'; const LOCK_RECORD_VERSION = 1 as const; @@ -433,8 +434,9 @@ const deniedCreateHandle = ( /** * Delete orphaned build/staging artifacts left in the lock directory by a * crashed prior writer. Safe precisely because we hold the exclusive lock: no - * other writer can be creating these here right now, so anything present is a - * crash orphan. Matches this slot's staging files ONLY — never `lbug` itself, + * other writer can be creating these here right now. A checkpoint-referenced + * generation is retained for embedding recovery; all other stages are orphans. + * Matches this slot's staging files ONLY — never `lbug` itself, * never `lbug.wal`/`lbug.shadow` (the LIVE index's own sidecars), and never a * `branches//` sub-slot (which owns its own lock + sweep). Non-recursive. */ @@ -442,6 +444,7 @@ export const sweepStagingArtifacts = (lockDir: string, log?: (msg: string) => vo // Matches `lbug.new`, `lbug.new.wal`, `lbug.staging.`, `lbug.staging..wal`, … // Does NOT match `lbug`, `lbug.wal`, `lbug.shadow`. const stagingRe = /^lbug\.(staging\..+|new(\..+)?)$/; + const retained = new Set(readEmbeddingRecovery(lockDir)?.familyFiles ?? []); let removed = 0; let entries: string[]; try { @@ -450,7 +453,7 @@ export const sweepStagingArtifacts = (lockDir: string, log?: (msg: string) => vo return; } for (const name of entries) { - if (!stagingRe.test(name)) continue; + if (!stagingRe.test(name) || retained.has(name)) continue; try { unlinkSync(path.join(lockDir, name)); removed++; diff --git a/gitnexus/src/storage/repo-meta.ts b/gitnexus/src/storage/repo-meta.ts index 0c412dd7c..1562ee8c7 100644 --- a/gitnexus/src/storage/repo-meta.ts +++ b/gitnexus/src/storage/repo-meta.ts @@ -43,6 +43,15 @@ export type ContentRetention = 'full' | 'symbol' | 'none'; export type FtsProfile = 'full' | 'symbol-no-file-content' | 'name-only'; export const CONTENT_RETENTION_SCHEMA_VERSION = 1; +/** Exact staged generation whose completed embedding groups can survive a retry. */ +export interface EmbeddingRecoveryReference { + /** A run-minted basename within this metadata file's index slot. */ + stagingFile: string; + schemaFingerprint: string; + /** Active-window and inherited incomplete groups: never reusable until completed. */ + unsafeNodeIds: string[]; +} + /** * Versioned receipt for the analyzer process that produced an index. * @@ -521,6 +530,11 @@ export interface RepoMeta { * subset of their chunks; for `'partial'` they hold none. */ pendingNodeIds?: string[]; + /** + * Interrupted atomic builds only. This does not publish the staged graph + * or advance live embedding statistics; it identifies a recovery source. + */ + recovery?: EmbeddingRecoveryReference; }; /** * Name of the git branch this index represents (#2106). Absent for the