fix(storage): retain checkpoint-referenced stages

This commit is contained in:
Gergő Magyar 2026-10-03 15:55:40 +01:00 • committed by GitHub
parent 047af140f4
commit 49062f48fc
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 138 additions and 3 deletions

View file

@ -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<string, unknown> =>
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;
}
};

View file

@ -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/<slug>/` 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.<id>`, `lbug.staging.<id>.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++;

View file

@ -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