diff --git a/gitnexus/src/core/embeddings/staged-embedding-recovery-child.ts b/gitnexus/src/core/embeddings/staged-embedding-recovery-child.ts index ada7c4297..b1fb2cb73 100644 --- a/gitnexus/src/core/embeddings/staged-embedding-recovery-child.ts +++ b/gitnexus/src/core/embeddings/staged-embedding-recovery-child.ts @@ -3,6 +3,7 @@ import fs from 'node:fs'; import path from 'node:path'; import lbug from '@ladybugdb/core'; import { createLbugDatabase, toNativeSafePath } from '../lbug/lbug-config.js'; +import { FAMILY_SUFFIXES } from '../../storage/embedding-recovery.js'; import { abortCachedEmbeddingsBuilder, createCachedEmbeddingsBuilder, @@ -11,15 +12,75 @@ import { } from './embedding-restore-spill.js'; import type { StagedEmbeddingExport } from './staged-embedding-recovery.js'; +/** Native replay and close may checkpoint; only give them disposable copies. */ +function copyRecoveryFamily(dbPath: string, exportDir: string): string { + const replayDir = fs.mkdtempSync(path.join(exportDir, 'replay-')); + const replayPath = path.join(replayDir, path.basename(dbPath)); + const noFollow = fs.constants.O_NOFOLLOW ?? 0; + const flags = fs.constants.O_RDONLY | noFollow | (fs.constants.O_NONBLOCK ?? 0); + const buffer = Buffer.allocUnsafe(1024 * 1024); + for (const suffix of FAMILY_SUFFIXES) { + const sourcePath = dbPath + suffix; + let entry: fs.BigIntStats; + try { + entry = fs.lstatSync(sourcePath, { bigint: true }); + } catch (error) { + if (suffix && (error as NodeJS.ErrnoException).code === 'ENOENT') continue; + throw error; + } + if (!entry.isFile()) throw new Error('staged embedding family is not a regular file'); + const source = fs.openSync(sourcePath, flags); + let destination: number | undefined; + try { + const opened = fs.fstatSync(source, { bigint: true }); + const current = fs.lstatSync(sourcePath, { bigint: true }); + if ( + !opened.isFile() || + !current.isFile() || + (noFollow === 0 && opened.ino === 0n) || + opened.dev !== entry.dev || + opened.ino !== entry.ino || + opened.dev !== current.dev || + opened.ino !== current.ino + ) { + throw new Error('staged embedding family changed while opening'); + } + destination = fs.openSync(replayPath + suffix, 'wx', 0o600); + let copied = 0n; + for (;;) { + const bytesRead = fs.readSync(source, buffer, 0, buffer.length, null); + if (bytesRead === 0) break; + fs.writeFileSync(destination, buffer.subarray(0, bytesRead)); + copied += BigInt(bytesRead); + } + const after = fs.fstatSync(source, { bigint: true }); + if ( + copied !== opened.size || + after.size !== opened.size || + after.mtimeNs !== opened.mtimeNs || + after.ctimeNs !== opened.ctimeNs + ) { + throw new Error('staged embedding family changed while copying'); + } + } finally { + try { + if (destination !== undefined) fs.closeSync(destination); + } finally { + fs.closeSync(source); + } + } + } + return replayPath; +} + async function extract(): Promise { const [dbPath, exportDir, dimensionsArg] = process.argv.slice(2); const dimensions = Number(dimensionsArg); if (!dbPath || !exportDir || !Number.isInteger(dimensions) || dimensions <= 0) { throw new Error('invalid staged embedding extraction arguments'); } - const stat = fs.lstatSync(dbPath); - if (!stat.isFile() || stat.isSymbolicLink()) - throw new Error('staged embedding DB is not a regular file'); + // The parent owns exportDir and reclaims it even after killing this child. + const replayPath = copyRecoveryFamily(dbPath, exportDir); const builder = createCachedEmbeddingsBuilder({ inMemoryRowLimit: 0, spillDir: exportDir }); const rejectedNodeIds = new Set(); let db: lbug.Database | undefined; @@ -27,7 +88,7 @@ async function extract(): Promise { try { // Avoid openLbugConnection's test-fixture lock sweep: a recovery source must // never have its WAL removed, even when an external slot resembles a fixture. - db = createLbugDatabase(lbug, toNativeSafePath(dbPath), { throwOnWalReplayFailure: true }); + db = createLbugDatabase(lbug, toNativeSafePath(replayPath), { throwOnWalReplayFailure: true }); conn = new lbug.Connection(db); const queried = await conn.query( 'MATCH (e:CodeEmbedding) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex, e.startLine AS startLine, e.endLine AS endLine, e.embedding AS embedding, e.contentHash AS contentHash',