diff --git a/.github/workflows/review-coverage.yml b/.github/workflows/review-coverage.yml index 66cb77891d..41e78e6c52 100644 --- a/.github/workflows/review-coverage.yml +++ b/.github/workflows/review-coverage.yml @@ -123,6 +123,7 @@ jobs: dataLayer/harperBridge/ResourceBridge.ts dataLayer/harperBridge/lmdbBridge/lmdbMethods/lmdbGetBackup.js dataLayer/restoreMarker.ts + dataLayer/restoreStaging.ts dataLayer/rocksdbBackup.ts dataLayer/schemaDescribe.ts resources/DatabaseTransaction.ts diff --git a/dataLayer/DESIGN.md b/dataLayer/DESIGN.md index 3fe9700efd..97400ce767 100644 --- a/dataLayer/DESIGN.md +++ b/dataLayer/DESIGN.md @@ -56,9 +56,44 @@ The certificate verification tables (`hdb_certificate_cache`, `hdb_crl_cache`, ` ## RocksDB backup/restore: the restore lock + marker protocol (`dataLayer/restoreMarker.ts`, `dataLayer/rocksdbBackup.ts`) -The `restore_backup` operation restores a user database on a live server by closing it across all -worker threads, purging its directory (`backups.restore` with `purgeAllFiles`), and reloading it. -Three non-obvious mechanics keep that safe: +The `restore_backup` operation restores a user database on a live server by restoring the backup +into a staging directory, closing the database across all worker threads, swapping the staged copy +in, and reloading it. Several non-obvious mechanics keep that safe: + +- **Nothing is destroyed until this build has opened the backup (harper#2965, + `dataLayer/restoreStaging.ts`).** `backups.restore`'s `purgeAllFiles` clears its target _before_ + copying, and an unreadable backup only fails inside the copy (checksum) or at the next open + (an unsupported table `format_version`, or a transaction log with an unsupported header, which + `open` silently skips). So the restore lands in `` `restore`/.staging ``, every staged + transaction-log store must pass `validateTransactionLogStore`, and the engine is opened and stamped + there; only then is the database published by two renames — itself to `.replaced`, staging to + its name. The first rename is the one-way door the marker protocol guards, so a staging failure + (corrupt, unreadable, `ENOSPC`) is a pre-destruction failure. Readability, not completeness, is the + bar: validation is non-strict, since a torn log tail is something open-time recovery truncates and + the operator has no override to restore past a refusal. The open's proof is bounded: with the + default table cache RocksDB loads only part of a large database's tables at open, so an + unsupported table outside that set still surfaces later, on a cold read. Opening every table + (`maxOpenFiles: -1`) would close that gap at the cost of a descriptor per table, which a large + database can exhaust. Cost: disk for one extra engine copy while staging, written online while + every database on that filesystem keeps serving, so a copy that would not leave headroom (the + larger of 256 MiB and a tenth of the copy) is refused with a 507 before it starts, where the + purge it replaced freed the space first. Two limits: an in-place restore needs room for two + copies even offline, where nothing else is serving, and the check is per restore, so concurrent + restores of different databases on one filesystem are not reserved against each other. Blob + roots are not staged (they span filesystems and the archive capabilities already gate their + encodings), so a blob-copy failure after publication still requires a rerun. A database + directory that is a symlink is refused, since the swap would replace the link with a directory, + and so is one that is a mount point of another filesystem, since staging (beside it) would land + elsewhere and the rename could only fail after a full copy. The check compares devices, so a + same-device bind mount passes it and fails at the first rename, with the database untouched. + **`.replaced` outlives every attempt under a preexisting marker**: a crash between the renames leaves it as the only copy of the database, so + a rerun keeps it until its own replacement publishes (and drops whatever is at the database path + then, a disposable candidate, before the space check; online, only once nothing holds it open). That inference needs `.replaced` never to + outlive its restore, so a finished restore renames it to `.discarded` while its marker still + stands, and only the removal of that may fail quietly; a `.replaced` that no marker accounts for + is refused (409), never trusted. A failed second rename moves it back, and + only a rollback whose directories were fsynced counts as "nothing destroyed". Once staging is + published the marker stays on any later failure, even where nothing was displaced. - **Two files in an isolated `` `restore` `` directory beside (never inside) the database directory**, each keyed by `sha256(basename(dbPath)).slice(0,32)`: `.lock`, an OS-level exclusive flock @@ -89,15 +124,18 @@ Three non-obvious mechanics keep that safe: it and broadcasting a reload would surface the earlier attempt's partial/corrupt directory as healthy. Only a _fresh_ marker on a _previously healthy_ database that failed before destruction is safe to clear. -- **The ITC close broadcast is normally best-effort, so closure is verified before the purge.** A +- **The ITC close broadcast is normally best-effort, so closure is verified before publication.** A SCHEMA broadcast (`signalSchemaChange`) usually resolves after remote handlers complete but times out at 30s "best-effort" and swallows errors. The restore `close` phase is stricter: it waits until every eligible recipient acknowledges or its port closes, and aborts the restore with a retryable 409 if that does not happen within 30s, because proceeding past an unconfirmed blob-save barrier - would re-open the race the barrier exists to close. A destructive purge still + would re-open the race the barrier exists to close. Publication still verifies closure independently: `restoreBackup` polls rocksdb-js `registryStatus()` (process-global across worker threads) until the database path has no open instance, and aborts with a 409 — - _cleaning up the marker, since nothing was destroyed_ — if handles remain. + _cleaning up the marker, since nothing was destroyed_ — if handles remain. The close can only reach + a loaded database, and online staging holds the marker for the whole copy, so a rescan during it + keeps a marked root this thread already has open (`restoreBlocksLoad`) instead of unloading it and + orphaning the handle. - **The close acknowledgement fences blob saves, deferred reclamation and orphan cleanup, not just database handles.** A store handle can close while a `saveBlob` file pipeline it started is still pending, because blob roots @@ -128,7 +166,7 @@ Three non-obvious mechanics keep that safe: correct.** rocksdb-js's registry is process-global but records only a per-path refCount, with no attribution to a thread or component; Harper keeps no component→database ownership map. So when a loaded component holds its own handle on the target database, `registryStatus()` stays non-zero, - Harper can neither identify nor force-close that handle, and an in-place purge would corrupt a live + Harper can neither identify nor force-close that handle, and swapping its directory would corrupt a live instance. `verifyDatabaseClosed` therefore waits only a short grace period (`DATABASE_CLOSE_WAIT_MS`, for a just-finished job worker's own close to drain) and then fails fast with a 409 that points at @@ -150,13 +188,12 @@ Three non-obvious mechanics keep that safe: calls `closeLoadedDatabases()` (`resources/databases.ts`) in its `finally`, closing every loaded user database on that thread (the non-enumerable `system` DB is intentionally skipped), so an exited job worker leaves no residual handle to be mistaken for a live holder. -- **A restore stamps a new database generation before `completeRestore`.** The restored files carry - the backup's generation, so both paths open the restored directory privately, stamp it and flush - (`stampDatabaseDirectory`, [database generation](../resources/DESIGN.md#database-generation-and-resumable-positions)) - inside the destructive section: a failed stamp leaves the marker, and the rerun re-purges and - re-stamps. The stamp is as durable as the marker protocol it runs inside. +- **A restore stamps a new database generation before publishing.** The restored files carry + the backup's generation, so both paths stamp the staged copy and flush + (`stampDatabaseDirectory`, [database generation](../resources/DESIGN.md#database-generation-and-resumable-positions)); + that open is also the readability proof above, and a failed stamp destroys nothing. - **`dropDatabase` and `restore_backup` serialize on the same lock, not a check-then-act probe.** - A drop's `destroy()` interleaving with a restore's purge-and-copy on the same directory would gut + A drop's `destroy()` interleaving with a restore's publication on the same directory would gut a "successful" restore (or vice versa). `dropDatabase` takes the restore lock for every RocksDB or LMDB root and publishes a positional `.dropping` marker beside each root before deleting any of them. A restore in progress makes the acquire fail with 409; `beginRestore` likewise refuses a @@ -180,13 +217,13 @@ Three non-obvious mechanics keep that safe: cheap pre-check that answers a plainly blocked caller without taking a lock. - **The offline restore probes RocksDB's own `LOCK` file, and fails closed.** The offline path runs only when the CLI sees no server (a PID heuristic; the PID file is briefly absent - mid-`harper restart`), and `backups.restore`'s `purgeAllFiles` never takes RocksDB's lock — so - before purging, `restoreBackupOffline` opens the database to probe. It now takes the restore + mid-`harper restart`), and publishing by rename never takes RocksDB's lock — so before + staging and again before publishing, `restoreBackupOffline` opens the database to probe. It now takes the restore lock+marker _before_ probing (so a server that starts afterward sees the marker and refuses to load), and recognizes the rocksdb-js lock error by message (`isRocksDbLockError`; at 2.5.0, when this was written, a plain `Error` with no `code`) — `IO error: While lock file: /LOCK: Resource temporarily unavailable` — aborting with a 409 - rather than purging a database another process holds open. Any _other_ open failure + rather than replacing a database another process holds open. Any _other_ open failure (corrupt/half-restored) is exactly what restore recovers, so only a lock conflict aborts. Known limitation: the flock is process-owned; if the restore job's worker _thread_ dies without diff --git a/dataLayer/restoreMarker.ts b/dataLayer/restoreMarker.ts index 6222e0dafe..47ac836727 100644 --- a/dataLayer/restoreMarker.ts +++ b/dataLayer/restoreMarker.ts @@ -13,7 +13,7 @@ import logger from '../utility/logging/harper_logger.ts'; * mutate the same directory concurrently. * * Restore metadata lives in an isolated `` `restore` `` directory *beside* the database directory - * (never inside it, since a restore purges the destination). Each database's two files are keyed by + * (never inside it, since a restore replaces the destination). Each database's entries are keyed by * a hash of the database directory name rather than being suffixed onto the name itself. That keeps * them out of the database-name namespace — a legal database literally named `orders.restoring` * would otherwise be mistaken for the restore marker of `orders`, and a 250-character name plus a @@ -49,6 +49,8 @@ import logger from '../utility/logging/harper_logger.ts'; * roots, and retrying the same drop removes remaining roots and the blob directories recorded for * their physical store identities before clearing the markers. The recorded identity is covered by * a digest, so damaged marker content fails closed instead of redirecting blob deletion. + * - `/.staging/`, `.replaced/` and `.discarded/` — a restore's proven + * replacement, the database it displaces, and that database once the restore no longer needs it; see `restoreStaging.ts` for when each may be removed. */ // The backtick makes this an illegal database name (schemaRegex rejects `/` and backtick only), so @@ -57,6 +59,9 @@ export const RESTORE_META_DIR = '`restore`'; export const RESTORE_LOCK_SUFFIX = '.lock'; export const RESTORING_MARKER_SUFFIX = '.restoring'; export const DROPPING_MARKER_SUFFIX = '.dropping'; +const RESTORE_STAGING_SUFFIX = '.staging'; +const RESTORE_REPLACED_SUFFIX = '.replaced'; +const RESTORE_DISCARDED_SUFFIX = '.discarded'; // Deliberately not a `.restoring` suffix: `scanBlockedRestores` selects markers by that suffix, and // a half-written temp must never be mistaken for one. const MARKER_TEMP_SUFFIX = '.tmp'; @@ -90,6 +95,18 @@ export function restoringMarkerPath(dbPath: string): string { return join(restoreMetaDir(dbPath), restoreMetaKey(dbPath) + RESTORING_MARKER_SUFFIX); } +export function restoreStagingPath(dbPath: string): string { + return join(restoreMetaDir(dbPath), restoreMetaKey(dbPath) + RESTORE_STAGING_SUFFIX); +} + +export function restoreReplacedPath(dbPath: string): string { + return join(restoreMetaDir(dbPath), restoreMetaKey(dbPath) + RESTORE_REPLACED_SUFFIX); +} + +export function restoreDiscardedPath(dbPath: string): string { + return join(restoreMetaDir(dbPath), restoreMetaKey(dbPath) + RESTORE_DISCARDED_SUFFIX); +} + export function droppingMarkerPath(dbPath: string): string { return join(restoreMetaDir(dbPath), restoreMetaKey(dbPath) + DROPPING_MARKER_SUFFIX); } diff --git a/dataLayer/restoreStaging.ts b/dataLayer/restoreStaging.ts new file mode 100644 index 0000000000..6f87629695 --- /dev/null +++ b/dataLayer/restoreStaging.ts @@ -0,0 +1,231 @@ +import { chmodSync, lstatSync, mkdirSync, readdirSync, renameSync, rmSync, statfsSync, statSync } from 'node:fs'; +import { dirname, join } from 'node:path'; +import { backups, validateTransactionLogStore } from '@harperfast/rocksdb-js'; +import { stampDatabaseDirectory } from '../resources/auditStore.ts'; +import { ClientError, ServerError } from '../utility/errors/hdbError.ts'; +import { fsyncDirectory, pathPresent } from '../utility/durableFile.ts'; +import logger from '../utility/logging/harper_logger.ts'; +import { + restoreDiscardedPath, + restoreMetaDir, + restoreReplacedPath, + restoreStagingPath, + type RestoreLock, +} from './restoreMarker.ts'; + +// Stage → prove → publish for `restore_backup` (harper#2965); the protocol is in dataLayer/DESIGN.md. +// The caller holds the restore lock (and its marker) throughout. + +/** + * Clear what an earlier attempt left behind. Staging was never published, so it is always disposable. + * `.replaced` is the database as it was before an interrupted publication — the only copy of it if + * that publication died between its renames — so it survives every attempt that runs under a + * preexisting marker. It exists only while a publication is unfinished: a finished restore moves it to + * `.discarded` while it still holds its marker (`discardReplaced`). + */ +export function prepareRestoreStaging(lock: RestoreLock, state: PublishResult): void { + const databaseDir = lock.dbPath; + if (isSymbolicLink(databaseDir)) { + // Repointing moves the restore metadata with the path, so an earlier restore's marker stops + // guarding its half-restored directory; only an offline rerun before the next start covers that. + const remedy = lock.preexisting + ? 'An earlier restore of this database did not finish: with Harper stopped, point the configured database path at the real directory, then rerun this restore offline before starting Harper again, so the half-restored directory is never loaded' + : 'Point the configured database path at the real directory, then rerun the restore'; + throw new ClientError( + `Cannot restore into ${databaseDir}: it is a symbolic link, and a restore replaces the database directory itself. ${remedy}` + ); + } + // Staging lives beside the database, so a mount point would cost a full copy only for the rename to fail. + if (pathPresent(databaseDir) && statSync(databaseDir).dev !== statSync(restoreMetaDir(databaseDir)).dev) { + throw new ClientError( + `Cannot restore into ${databaseDir}: it is a mount point, on a different filesystem from ${dirname(databaseDir)}, and a restore replaces the database directory by renaming it. Mount the volume at the parent directory instead, or restore offline into a new target_database` + ); + } + const replacedDir = restoreReplacedPath(databaseDir); + // Only a marker proves `.replaced` is this database's unfinished publication; without one the + // database path may be the live database, not a candidate. + if (!lock.preexisting && pathPresent(replacedDir)) { + throw new ClientError( + `Cannot restore into ${databaseDir}: ${replacedDir} is left from an earlier restore that recorded no restore marker, so it cannot be told apart from the database. Inspect it and remove it, then rerun the restore`, + 409 + ); + } + rmSync(restoreStagingPath(databaseDir), { recursive: true, force: true }); + rmSync(restoreDiscardedPath(databaseDir), { recursive: true, force: true }); + if (pathPresent(replacedDir)) { + // A publication began and never finished, so the database path holds a candidate, never the + // database; dropped now rather than at publish so the space check does not count a third copy. + state.destroyed = true; + rmSync(databaseDir, { recursive: true, force: true }); + } +} + +/** + * Readability, not completeness, is the bar: a torn log tail is something open-time recovery + * truncates, so it is not a reason to refuse a backup the operator has no other way to restore. + */ +export async function stageRestore(backupDir: string, backupId: number, lock: RestoreLock): Promise { + const databaseDir = lock.dbPath; + const stagingDir = restoreStagingPath(databaseDir); + await assertRoomToStage(backupDir, backupId, lock); + try { + mkdirSync(stagingDir); + // After a crash between the publication renames, `.replaced` is the only record of the database's mode. + const modeSource = pathPresent(databaseDir) ? databaseDir : restoreReplacedPath(databaseDir); + if (pathPresent(modeSource)) chmodSync(stagingDir, lstatSync(modeSource).mode & 0o7777); + await backups.restore(backupDir, stagingDir, { backupId, mode: 'purgeAllFiles' }); + const logsDir = join(stagingDir, 'transaction_logs'); + if (pathPresent(logsDir)) { + for (const entry of readdirSync(logsDir, { withFileTypes: true })) { + if (!entry.isDirectory()) continue; + const result = await validateTransactionLogStore(join(logsDir, entry.name)); + if (!result.valid) { + const problems = [ + ...result.errors, + ...result.files.flatMap((file) => file.errors.map((e) => `${file.file}: ${e}`)), + ]; + throw new Error(`transaction log '${entry.name}' is unreadable: ${problems.join('; ')}`); + } + } + } + await stampDatabaseDirectory(stagingDir, { carriesLog: true }); + } catch (error) { + throw new Error( + `Backup ${backupId} could not be staged and verified, so ${untouched(lock)}: ${error instanceof Error ? error.message : String(error)}`, + { + cause: error, + } + ); + } +} + +// Free space a restore leaves on the shared filesystem for the databases still serving there. +const STAGING_HEADROOM_BYTES = 256 * 1024 ** 2; + +/** + * Staging needs a second engine copy where the purge it replaced freed the space first, and online it + * is written while every database on that filesystem keeps serving. Running out partway would fail + * their writes too, so a copy that will not fit is refused before it starts. + */ +async function assertRoomToStage(backupDir: string, backupId: number, lock: RestoreLock): Promise { + const databaseDir = lock.dbPath; + const engineBytes = (await backups.list(backupDir)).find((backup) => backup.backupId === backupId)?.size ?? 0; + const needed = engineBytes + directoryBytes(join(backupDir, 'transaction_logs', String(backupId))); + const headroom = Math.max(STAGING_HEADROOM_BYTES, needed / 10); + const { bavail, bsize } = statfsSync(restoreMetaDir(databaseDir)); + const available = Number(bavail) * Number(bsize); + if (available < needed + headroom) { + throw new ServerError( + `Cannot restore backup ${backupId}: staging it needs about ${formatBytes(needed)} beside ${databaseDir}, plus ${formatBytes(headroom)} left free for the databases still serving on that filesystem, but only ${formatBytes(available)} is available. Free space there and rerun the restore; ${untouched(lock)}`, + 507 + ); + } +} + +/** What a refusal before publication can truthfully say about the destination. */ +function untouched(lock: RestoreLock): string { + return lock.preexisting + ? `${lock.dbPath} keeps the marker of an earlier restore that did not finish, and may be incomplete; rerun restore_backup to recover` + : `${lock.dbPath} was not modified`; +} + +function directoryBytes(path: string): number { + if (!pathPresent(path)) return 0; + let total = 0; + for (const entry of readdirSync(path, { withFileTypes: true })) { + const entryPath = join(path, entry.name); + total += entry.isDirectory() ? directoryBytes(entryPath) : statSync(entryPath).size; + } + return total; +} + +function formatBytes(bytes: number): string { + return `${(bytes / 1024 ** 3).toFixed(2)} GiB`; +} + +export type PublishResult = { destroyed: boolean }; + +/** + * Swap the staged database in. Reports through `destroyed` whether the destination may have changed, + * which is what decides whether the caller may clear its marker; an error is rethrown either way. + */ +export function publishStagedRestore(lock: RestoreLock, state: PublishResult): void { + const databaseDir = lock.dbPath; + const stagingDir = restoreStagingPath(databaseDir); + const replacedDir = restoreReplacedPath(databaseDir); + const parentDir = dirname(databaseDir); + const metaDir = restoreMetaDir(databaseDir); + if (pathPresent(replacedDir)) { + // Whatever is at the database path is a candidate an earlier attempt published and never finished. + state.destroyed = true; + rmSync(databaseDir, { recursive: true, force: true }); + } else if (pathPresent(databaseDir)) { + renameSync(databaseDir, replacedDir); + state.destroyed = true; + fsyncDirectory(parentDir); + fsyncDirectory(metaDir); + } + try { + renameSync(stagingDir, databaseDir); + } catch (error) { + if (pathPresent(replacedDir)) rollBackPublication(replacedDir, databaseDir, state); + throw error; + } + // Even with nothing displaced, a published engine whose blobs never landed must keep its marker. + state.destroyed = true; + fsyncDirectory(metaDir); + fsyncDirectory(parentDir); +} + +/** + * Put the pre-restore database back. Only a rollback that is itself durable lets the caller clear its + * marker; otherwise the marker stays and a rerun finds `.replaced` where this left it. + */ +function rollBackPublication(replacedDir: string, databaseDir: string, state: PublishResult): void { + try { + renameSync(replacedDir, databaseDir); + fsyncDirectory(restoreMetaDir(databaseDir)); + fsyncDirectory(dirname(databaseDir)); + state.destroyed = false; + } catch (error) { + logger.error(`Could not move the pre-restore copy of ${databaseDir} back; it remains at ${replacedDir}`, error); + } +} + +/** + * Retire the pre-restore copy before the marker clears. The rename is atomic, so `.replaced` can never + * outlive its restore and later pass for an unfinished publication; it throws while the marker still + * stands. Only removing the renamed copy may fail quietly, since the next restore removes it. + */ +export function discardReplaced(lock: RestoreLock): void { + const replacedDir = restoreReplacedPath(lock.dbPath); + if (!pathPresent(replacedDir)) return; + const discardedDir = restoreDiscardedPath(lock.dbPath); + rmSync(discardedDir, { recursive: true, force: true }); + renameSync(replacedDir, discardedDir); + fsyncDirectory(restoreMetaDir(lock.dbPath)); + try { + rmSync(discardedDir, { recursive: true, force: true }); + } catch (error) { + logger.warn(`Could not remove the pre-restore copy of ${lock.dbPath}; the next restore removes it`, error); + } +} + +/** Staging is never published by a failed attempt, so it is always safe to drop; never fatal. */ +export function discardRestoreStaging(lock: RestoreLock): void { + const stagingDir = restoreStagingPath(lock.dbPath); + try { + rmSync(stagingDir, { recursive: true, force: true }); + } catch (error) { + logger.warn(`Could not remove restore staging at ${stagingDir}; the next restore removes it`, error); + } +} + +function isSymbolicLink(path: string): boolean { + try { + return lstatSync(path).isSymbolicLink(); + } catch (error) { + if ((error as any)?.code === 'ENOENT') return false; + throw error; + } +} diff --git a/dataLayer/rocksdbBackup.ts b/dataLayer/rocksdbBackup.ts index ed81817ce8..97d8ef354b 100644 --- a/dataLayer/rocksdbBackup.ts +++ b/dataLayer/rocksdbBackup.ts @@ -12,7 +12,6 @@ import { setTimeout as delay } from 'node:timers/promises'; import { pack as tarPack, type Pack } from 'tar-stream'; import { RocksDatabase, backups, registryStatus, type BackupInfo } from '@harperfast/rocksdb-js'; import { databases, getDatabases, resolveDatabasePath } from '../resources/databases.ts'; -import { stampDatabaseDirectory } from '../resources/auditStore.ts'; import { type BlobCaptureDisposition, classifyBlobFileForCapture, @@ -33,6 +32,7 @@ import { abandonRestore, checkRestoreState, releaseRestoreLock, + restoreReplacedPath, type RestoreLock, } from './restoreMarker.ts'; import { @@ -68,6 +68,14 @@ import { writeBackupManifest, type BackupManifest, } from './backupManifest.ts'; +import { + discardReplaced, + discardRestoreStaging, + prepareRestoreStaging, + publishStagedRestore, + stageRestore, + type PublishResult, +} from './restoreStaging.ts'; import logger from '../utility/logging/harper_logger.ts'; /** @@ -529,8 +537,8 @@ files by hand. ## Restore -Restore is destructive: it purges and rewrites the database directory — and every blob root — from -the backup (blobs are restored automatically). Restore the latest backup in place: +Restore is destructive: once a staged copy of the backup has opened, it replaces the database directory — and +rewrites every blob root — from the backup (blobs are restored automatically). Restore the latest backup in place: harper restore_backup database=${databaseName} @@ -616,8 +624,9 @@ export async function validateRestoreBackup(request: any) { /** * Online restore of a user database (see the design's restore lock + marker protocol): - * take the per-database restore lock, write the restoring marker, close the database across all - * worker threads, restore, delete the marker, release the lock, and reload everywhere. + * take the per-database restore lock, write the restoring marker, stage and verify the backup, close + * the database across all worker threads, publish the staged copy, delete the marker, release the + * lock, and reload everywhere. */ export async function restoreBackup(request: any) { const databaseName = getDatabaseName(request); @@ -654,11 +663,11 @@ export async function restoreBackup(request: any) { await assertBlobSnapshotRestorable(backupDir, backupId, blobRoots); const allowEngineOnly = requireBooleanOption(request.allow_engine_only, 'allow_engine_only'); // Once is enough: the decision reads the manifest and the opt-in, never the destination, so no - // concurrent writer can change the answer between here and the purge. + // concurrent writer can change the answer between here and publication. assertEngineOnlyRestoreAllowed(databaseName, { backupHasBlobs: manifest.blobs, allowEngineOnly }); const pinId = restorePinId(databaseDir); const restoreToken = randomUUID(); - let destructionStarted = false; + const publication: PublishResult = { destroyed: false }; // Re-check before replacing a previous attempt's claim; publish the new claim before the marker // under both locks, so a crash cannot leave a marked database with an unprotected source. const lock = await withBackupRepositoryLock(backupDir, databaseName, async () => { @@ -668,6 +677,13 @@ export async function restoreBackup(request: any) { ); }); try { + // Preparing a rerun of an unfinished publication drops the candidate it left, so nothing may hold it. + if (lock.preexisting && pathPresent(restoreReplacedPath(databaseDir))) { + await verifyDatabaseClosed(databaseDir, databaseName); + } + // Staged while the database is still open and serving, so the copy is not downtime. + prepareRestoreStaging(lock, publication); + await stageRestore(backupDir, backupId, lock); // Block new blob saves, drain in-flight saves, and close the database across all worker threads. // Each thread also rescans, and the restoring marker keeps it from reloading mid-restore. try { @@ -679,26 +695,26 @@ export async function restoreBackup(request: any) { } // A live component (or the system database) can hold its own handle on the database that // Harper does not track and cannot close, so verify actual process-wide closure before - // purging — restoring under an open instance would corrupt it. If handles remain, fail - // with a clear pointer to the offline CLI path rather than purging. + // publishing — replacing the directory under an open instance would corrupt it. If handles + // remain, fail with a clear pointer to the offline CLI path instead. await verifyDatabaseClosed(databaseDir, databaseName); - destructionStarted = true; - await backups.restore(backupDir, databaseDir, { backupId, mode: 'purgeAllFiles' }); + publishStagedRestore(lock, publication); // restore blobs only for a backup that captured them (an engine-only backup leaves the live // blob roots untouched); the manifest, not the mere presence of a snapshot dir, is the source // of truth so a mid-copy or absent snapshot can't be misread if (manifest.blobs) { await restoreBlobSnapshot(backupDir, backupId, databaseName, getBlobPathsForDatabaseName(databaseName)); } - await stampDatabaseDirectory(databaseDir, { carriesLog: true }); + discardReplaced(lock); } catch (error: any) { + discardRestoreStaging(lock); // Leave the marker (so startup/rescan detection reports an incomplete restore until a rerun - // succeeds) when either the destructive purge has begun, OR this attempt was itself a recovery - // over a pre-existing marker: in that case the directory may already be half-purged from an - // earlier failed restore, so clearing the marker and reloading it as healthy would surface - // partial/corrupt data. Only a *fresh* marker on a *previously healthy* database that failed - // before any destruction is safe to clear. - if (destructionStarted || lock.preexisting) { + // succeeds) when either publication has replaced the database, OR this attempt was itself a + // recovery over a pre-existing marker: in that case the directory may hold an unfinished + // publication from an earlier restore, so clearing the marker and reloading it as healthy would + // surface partial/corrupt data. Only a *fresh* marker on a *previously healthy* database that + // failed before any destruction is safe to clear. + if (publication.destroyed || lock.preexisting) { // The marker stays, so the database is unloadable until a rerun — and the rerun needs this // backup. The pin stays with it, and lapses on its own once the marker is gone. abandonRestore(lock); @@ -715,12 +731,7 @@ export async function restoreBackup(request: any) { // and a process restart clears it regardless. logger.error(`Could not release the blob fence after a failed restore of '${databaseName}'`, releaseError); } - // wrap rather than mutate error.message: a frozen/library error can have a non-writable - // message (assigning it throws TypeError under 'use strict') - throw new Error( - `Restore of database '${databaseName}' from backup ${backupId} failed (rerun restore_backup to recover): ${error.message}`, - { cause: error } - ); + throw rerunRequiredError(databaseName, backupId, error); } // nothing destructive happened and the marker was fresh — clear it and let every thread reload // the intact database @@ -845,7 +856,7 @@ function beginRestoreForDatabase( * surfaces it as a plain `Error` with no `code` and a message like * `IO error: While lock file: /LOCK: Resource temporarily unavailable`, so string-matching is * the only signal available (there is no typed error to key on — a native primitive is a rocksdb-js - * follow-on). We match conservatively and fail *closed* on a hit so the offline restore never purges + * follow-on). We match conservatively and fail *closed* on a hit so the offline restore never replaces * a database another process still has open. */ function isRocksDbLockError(error: any): boolean { @@ -1271,7 +1282,7 @@ export async function restoreBackupOffline( await assertBlobSnapshotRestorable(backupDir, backupId, blobRoots); assertEngineOnlyRestoreAllowed(targetDatabase ?? databaseName, { backupHasBlobs: manifest.blobs, allowEngineOnly }); const pinId = restorePinId(databaseDir); - let destructionStarted = false; + const publication: PublishResult = { destroyed: false }; // As online, claim before marking under both locks, and mark before probing the destination. const lock = await withBackupRepositoryLock(backupDir, databaseName, async () => { await findBackup(backupDir, backupId as number, databaseName); @@ -1303,29 +1314,11 @@ export async function restoreBackupOffline( }); }); try { - // The offline path is entered only when the CLI sees no running server (getHdbPid), but that is - // a heuristic: the PID file is briefly absent mid-`harper restart`, and backups.restore's - // purgeAllFiles never takes RocksDB's own lock. Probe that lock by opening the database — a live - // holder makes open throw its LOCK-file error (isRocksDbLockError) — so we fail closed rather - // than purge a database another process still has open. A directory that fails to open for any - // *other* reason (corrupt or half-restored) is exactly what restore recovers, so only a lock - // conflict aborts. - if (existsSync(join(databaseDir, 'CURRENT'))) { - let handle: RocksDatabase | undefined; - try { - handle = RocksDatabase.open(databaseDir); - } catch (error: any) { - if (isRocksDbLockError(error)) { - throw new BackupInProgressError( - `Cannot restore database '${databaseName}': it is open by a running Harper process — stop Harper before restoring offline` - ); - } - // otherwise corrupt/half-restored — fall through and let restore recover it - } - handle?.close(); - } - destructionStarted = true; - await backups.restore(backupDir, databaseDir, { backupId, mode: 'purgeAllFiles' }); + assertNotOpenElsewhere(databaseDir, databaseName); + prepareRestoreStaging(lock, publication); + await stageRestore(backupDir, backupId, lock); + assertNotOpenElsewhere(databaseDir, databaseName); + publishStagedRestore(lock, publication); // restore blobs only for a backup that captured them (per the manifest, not snapshot presence) if (manifest.blobs) { await restoreBlobSnapshot( @@ -1335,24 +1328,22 @@ export async function restoreBackupOffline( getBlobPathsForDatabaseName(targetDatabase ?? databaseName) ); } - await stampDatabaseDirectory(databaseDir, { carriesLog: true }); + discardReplaced(lock); } catch (error: any) { + discardRestoreStaging(lock); // Preserve the marker on a destructive failure or a recovery over a pre-existing marker (see // the online restoreBackup for the rationale); otherwise clear the fresh marker so an intact, // merely-locked database is not left flagged as an incomplete restore. // The pin stays exactly as long as the marker does: a retained marker means a rerun is required, // and the rerun needs this backup to still be there. - if (destructionStarted || lock.preexisting) abandonRestore(lock); + if (publication.destroyed || lock.preexisting) abandonRestore(lock); else { releaseRestoreClaim(backupDir, pinId, lock, databaseName); } // preserve typed client errors (e.g. the 409 lock probe) unwrapped; only wrap an opaque restore // failure after destruction has begun - if (destructionStarted && !(error instanceof ClientError)) { - throw new Error( - `Restore of database '${databaseName}' from backup ${backupId} failed (rerun restore_backup to recover): ${error.message}`, - { cause: error } - ); + if (publication.destroyed && !(error instanceof ClientError)) { + throw rerunRequiredError(databaseName, backupId, error); } throw error; } @@ -1365,6 +1356,43 @@ export async function restoreBackupOffline( }; } +/** + * The offline path is entered only when the CLI sees no running server (getHdbPid), but that is a + * heuristic: the PID file is briefly absent mid-`harper restart`, and publishing by rename never takes + * RocksDB's own lock. Probe that lock by opening the database — a live holder makes open throw its + * LOCK-file error (isRocksDbLockError) — so we fail closed rather than replace a database another + * process still has open. A directory that fails to open for any *other* reason (corrupt or + * half-restored) is exactly what restore recovers, so only a lock conflict aborts. + */ +function assertNotOpenElsewhere(databaseDir: string, databaseName: string): void { + if (!existsSync(join(databaseDir, 'CURRENT'))) return; + let handle: RocksDatabase | undefined; + try { + handle = RocksDatabase.open(databaseDir); + } catch (error: any) { + if (isRocksDbLockError(error)) { + throw new BackupInProgressError( + `Cannot restore database '${databaseName}': it is open by a running Harper process — stop Harper before restoring offline` + ); + } + } + handle?.close(); +} + +/** + * A restore that may have changed the destination needs a rerun. Wrapped rather than mutated, since a + * frozen or library error can have a non-writable message; the status code of a refusal such as a 507 + * is carried over. + */ +function rerunRequiredError(databaseName: string, backupId: number, error: any): Error { + const wrapped: any = new Error( + `Restore of database '${databaseName}' from backup ${backupId} failed (rerun restore_backup to recover): ${error instanceof Error ? error.message : String(error)}`, + { cause: error } + ); + if (typeof error?.statusCode === 'number') wrapped.statusCode = error.statusCode; + return wrapped; +} + function isMissingOrEmptyDir(path: string): boolean { try { return readdirSync(path).length === 0; diff --git a/integrationTests/database/restore-backup-staging.test.ts b/integrationTests/database/restore-backup-staging.test.ts new file mode 100644 index 0000000000..d2262fd254 --- /dev/null +++ b/integrationTests/database/restore-backup-staging.test.ts @@ -0,0 +1,197 @@ +import { suite, test, before, after } from 'node:test'; +import { ok, strictEqual, match } from 'node:assert'; +import { execFile } from 'node:child_process'; +import { closeSync, existsSync, openSync, readdirSync, writeSync } from 'node:fs'; +import { join, resolve } from 'node:path'; +import { setTimeout as sleep } from 'node:timers/promises'; +import { promisify } from 'node:util'; +import { startHarper, killHarper, teardownHarper, type ContextWithHarper } from '@harperfast/integration-testing'; + +const run = promisify(execFile); +const HARPER_BIN = resolve(import.meta.dirname, '../../dist/bin/harper.js'); +const DATABASE = 'restore_staging'; +const TABLE = 'Items'; +const skipSuite = process.env.HARPER_RUNTIME === 'bun' || process.platform === 'win32'; +const JOB_TIMEOUT_MS = 120_000; +const BASE_ROWS = 50; + +suite('restore_backup verifies before it replaces (harper#2965)', { skip: skipSuite }, (ctx: ContextWithHarper) => { + const baseIds = Array.from({ length: BASE_ROWS }, (_, i) => `base-${i}`); + let backupId = 0; + + async function op(operation: Record): Promise { + const { username, password } = ctx.harper.admin; + const res = await fetch(ctx.harper.operationsAPIURL, { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + 'Authorization': 'Basic ' + Buffer.from(`${username}:${password}`).toString('base64'), + }, + body: JSON.stringify(operation), + signal: AbortSignal.timeout(60_000), + }); + const text = await res.text(); + let body: any = text; + try { + body = JSON.parse(text); + } catch { + /* keep text */ + } + return { status: res.status, body, text }; + } + + async function runJob(operation: Record): Promise { + const started = await op(operation); + strictEqual(started.status, 200, `${operation.operation} rejected: ${started.text}`); + const jobId = /Starting job with id ([\w-]+)/.exec(started.body?.message ?? '')?.[1]; + ok(jobId, `no job id in ${started.text}`); + const deadline = Date.now() + JOB_TIMEOUT_MS; + let last = ''; + while (Date.now() < deadline) { + const r = await op({ operation: 'get_job', id: jobId }); + const record = Array.isArray(r.body) ? r.body[0] : r.body; + if (record?.status === 'COMPLETE' || record?.status === 'ERROR') return record; + last = r.text.slice(0, 300); + await sleep(500); + } + throw new Error(`job ${jobId} did not settle within ${JOB_TIMEOUT_MS}ms; last=${last}`); + } + + async function createBackup(): Promise { + const job = await runJob({ operation: 'create_backup', database: DATABASE }); + strictEqual(job.status, 'COMPLETE', `create_backup failed: ${JSON.stringify(job)}`); + const id = job.result?.backup_id ?? (typeof job.message === 'object' ? job.message?.backup_id : undefined); + ok(Number.isInteger(id), `no backup_id in job: ${JSON.stringify(job)}`); + return id; + } + + async function insert(ids: string[]): Promise { + const r = await op({ + operation: 'insert', + database: DATABASE, + table: TABLE, + records: ids.map((id) => ({ id, note: `note-${id}` })), + }); + strictEqual(r.status, 200, `insert failed: ${r.text}`); + } + + async function readIds(): Promise { + const r = await op({ + operation: 'search_by_value', + database: DATABASE, + table: TABLE, + search_attribute: 'id', + search_value: '*', + get_attributes: ['id', 'note'], + }); + strictEqual(r.status, 200, `search failed: ${r.text}`); + for (const row of r.body) strictEqual(row.note, `note-${row.id}`, `row ${row.id} not value-exact`); + return r.body.map((row: any) => row.id).sort(); + } + + async function expectIds(expected: string[], label: string): Promise { + deepEqualSorted(await readIds(), expected, label); + } + + function deepEqualSorted(actual: string[], expected: string[], label: string): void { + const want = [...expected].sort(); + strictEqual(actual.length, want.length, `${label}: expected ${want.length} rows, got ${actual.length}`); + strictEqual(actual.join(','), want.join(','), `${label}: row set differs`); + } + + // The backups root is configurable and the default dir name is an implementation detail, so + // locate the repository by its `private/` engine directory instead of hard-coding it. + function findBackupPrivateDir(id: number): string { + const root = ctx.harper.dataRootDir; + for (const top of readdirSync(root)) { + const candidate = join(root, top, DATABASE, 'private', String(id)); + if (existsSync(candidate)) return candidate; + } + throw new Error(`no backup repository for ${DATABASE}#${id} under ${root}`); + } + + function corruptManifest(id: number): void { + const dir = findBackupPrivateDir(id); + const manifest = readdirSync(dir).find((f) => f.startsWith('MANIFEST-')); + ok(manifest, `no MANIFEST-* in ${dir}`); + const fd = openSync(join(dir, manifest), 'r+'); + try { + writeSync(fd, Buffer.alloc(64, 0x5a), 0, 64, 0); + } finally { + closeSync(fd); + } + } + + const start = () => + startHarper(ctx, { config: { threads: { count: 3 }, logging: { console: true, level: 'error' } } }); + + async function cliRestore(id: number): Promise<{ exit: number; out: string }> { + try { + const r = await run(process.execPath, [HARPER_BIN, 'restore_backup', `database=${DATABASE}`, `backup_id=${id}`], { + env: { ...process.env, ROOTPATH: ctx.harper.dataRootDir }, + timeout: JOB_TIMEOUT_MS, + maxBuffer: 16 * 1024 * 1024, + }); + return { exit: 0, out: `${r.stdout}\n${r.stderr}` }; + } catch (e: any) { + return { exit: typeof e.code === 'number' ? e.code : 1, out: `${e.stdout ?? ''}\n${e.stderr ?? ''}` }; + } + } + + before(async () => { + await start(); + const t = await op({ operation: 'create_table', database: DATABASE, table: TABLE, primary_key: 'id' }); + ok(t.status === 200, `create_table failed: ${t.text}`); + await insert(baseIds); + backupId = await createBackup(); + }); + + after(async () => { + await teardownHarper(ctx); + }); + + test('online: corrupt backup is refused and the database is untouched', async () => { + const postIds = Array.from({ length: 10 }, (_, i) => `post-${i}`); + await insert(postIds); + corruptManifest(backupId); + + const job = await runJob({ operation: 'restore_backup', database: DATABASE, backup_id: backupId }); + strictEqual(job.status, 'ERROR', `restore of a corrupt backup must fail: ${JSON.stringify(job)}`); + match(String(job.message), /was not modified/); + + await expectIds([...baseIds, ...postIds], 'after refused restore'); + await insert(['after-refusal']); + await expectIds([...baseIds, ...postIds, 'after-refusal'], 'after write post-refusal'); + }); + + test('offline CLI: corrupt backup is refused and the database is intact after restart', async () => { + const corruptId = await createBackup(); + const before = await readIds(); + await killHarper(ctx); + corruptManifest(corruptId); + + const { exit, out } = await cliRestore(corruptId); + ok(exit !== 0, `CLI restore of a corrupt backup must fail; output:\n${out}`); + match(out, /was not modified/); + + await start(); + await expectIds(before, 'after offline refusal and restart'); + }); + + // The success path runs offline: online, the restore stages but is then refused at the closure check + // while leaked handles stay open (harper#3120). + test('offline CLI: a valid backup is staged, swapped in, and survives a restart', async () => { + const validId = await createBackup(); + const kept = await readIds(); + await insert(['after-valid-backup']); + await killHarper(ctx); + + const { exit, out } = await cliRestore(validId); + strictEqual(exit, 0, `CLI restore of a valid backup failed; output:\n${out}`); + + await start(); + await expectIds(kept, 'after restore and restart'); + await insert(['after-restore']); + await expectIds([...kept, 'after-restore'], 'writable after restore'); + }); +}); diff --git a/resources/databases.ts b/resources/databases.ts index 7c24d2db2e..51f78d7863 100644 --- a/resources/databases.ts +++ b/resources/databases.ts @@ -1068,7 +1068,7 @@ export function getDatabases(): Databases { blockedByDrop.databaseNames.has(dbName) ) continue; - if (blockedByRestore.has(dbName)) continue; + if (restoreBlocksLoad(blockedByRestore, dbName, dbPath)) continue; if (isOpenBranchPath(dbPath)) continue; if ( @@ -1136,8 +1136,8 @@ export function getDatabases(): Databases { if (databaseEntry.name.endsWith(MIGRATING_DIR_SUFFIX)) continue; // migration staging dir if (databaseEntry.name === RESTORE_META_DIR) continue; // reserved restore-metadata dir if (databaseEntry.name === BRANCH_ROOT_DIR) continue; // reserved branch root - if (blockedByRestore.has(basename(databaseEntry.name, '.mdb'))) continue; const dbPath = join(databasePath, databaseEntry.name); + if (restoreBlocksLoad(blockedByRestore, basename(databaseEntry.name, '.mdb'), dbPath)) continue; if (databaseRootUnavailable(dbPath)) continue; if (blockedByDrop.rootPaths.has(dbPath) || blockedByDrop.databaseNames.has(dbName)) continue; if (isOpenBranchPath(dbPath)) continue; @@ -1384,6 +1384,15 @@ function reportRelationshipError(key: string, message: string): void { logger.error(message); } +/** + * A marked root this thread already has open is the live database an online restore is staging + * beside: it keeps serving and stays loaded until the restore's close broadcast closes it. Dropping + * it here would orphan the handle, since `closeDatabase` only reaches loaded databases. + */ +function restoreBlocksLoad(blockedByRestore: Set, dbName: string, dbPath: string): boolean { + return blockedByRestore.has(dbName) && !rocksdbDatabaseEnvs.has(dbPath) && !lmdbDatabaseEnvs.has(dbPath); +} + /** * Scan a databases directory's entries for restore lock/marker files and return the names of * databases that must not be loaded: a held restore lock means a restore is in progress in some @@ -1395,7 +1404,11 @@ function databasesBlockedByRestore(databasePath: string): Set { const blocked = new Set(); for (const [dbName, state] of scanBlockedRestores(databasePath)) { if (state === 'in-progress') { - logger.warn(`A restore of database '${dbName}' is in progress; not loading it`); + // An online restore stages beside a database this thread keeps serving (`restoreBlocksLoad`). + const serving = + rocksdbDatabaseEnvs.has(join(databasePath, dbName)) || + lmdbDatabaseEnvs.has(join(databasePath, `${dbName}.mdb`)); + if (!serving) logger.warn(`A restore of database '${dbName}' is in progress; not loading it`); blocked.add(dbName); } else if (state === 'incomplete') { logger.error( diff --git a/server/itc/serverHandlers.js b/server/itc/serverHandlers.js index 1cbfc8afeb..e61b3f76f7 100644 --- a/server/itc/serverHandlers.js +++ b/server/itc/serverHandlers.js @@ -100,8 +100,8 @@ async function schemaHandler(event) { event.message.dropPreparationRootPaths ); } - // restore_backup: this thread must release its store handles so the restore can purge and - // rewrite the database directory. The rescan below (resetDatabases) skips reloading it while + // restore_backup: this thread must release its store handles so the restore can replace the + // database directory. The rescan below (resetDatabases) skips reloading it while // the restoring marker is present, and reloads it on the completion signal (marker gone). let resumeBlobSavesFor; if (event.message?.operation === hdbTerms.OPERATIONS_ENUM.RESTORE_BACKUP && event.message.schema) { diff --git a/unitTests/dataLayer/rocksdbBackup.test.js b/unitTests/dataLayer/rocksdbBackup.test.js index d7301ffd7c..19c1c7a558 100644 --- a/unitTests/dataLayer/rocksdbBackup.test.js +++ b/unitTests/dataLayer/rocksdbBackup.test.js @@ -1,10 +1,24 @@ 'use strict'; const assert = require('node:assert'); -const { chmodSync, existsSync, mkdirSync, mkdtempSync, readFileSync, rmSync, writeFileSync } = require('node:fs'); +const fs = require('node:fs'); +const { + chmodSync, + existsSync, + lstatSync, + mkdirSync, + mkdtempSync, + readdirSync, + readFileSync, + renameSync, + rmSync, + symlinkSync, + writeFileSync, +} = fs; const { dirname, join } = require('node:path'); const { tmpdir } = require('node:os'); const { spawn } = require('node:child_process'); +const { syncBuiltinESMExports } = require('node:module'); const { extract } = require('tar-stream'); const { RocksDatabase } = require('@harperfast/rocksdb-js'); const { @@ -39,10 +53,18 @@ const { getBlobPathsForDatabaseName } = require('#src/resources/blob'); // managed-backup ops self-enforce super_user (see requireSuperUser in rocksdbBackup.ts); requests // in the online-operation tests below must therefore carry a super_user role. const SU = { hdb_user: { role: { permission: { super_user: true } } } }; -const { abandonRestore, beginRestore, completeRestore, checkRestoreState } = require('#src/dataLayer/restoreMarker'); +const { + abandonRestore, + beginRestore, + completeRestore, + checkRestoreState, + restoreDiscardedPath, + restoreReplacedPath, + restoreStagingPath, +} = require('#src/dataLayer/restoreMarker'); const { pinBackup, readBackupPins, unpinBackup, withBackupRepositoryLock } = require('#src/dataLayer/backupRepository'); const { backups } = require('@harperfast/rocksdb-js'); -const { closeLoadedDatabases } = require('#src/resources/databases'); +const { closeLoadedDatabases, resetDatabases } = require('#src/resources/databases'); const { ARCHIVE_MANIFEST_ENTRY, ARCHIVE_SCHEMA_VERSION, @@ -322,6 +344,469 @@ describe('rocksdbBackup', function () { }); }); + // harper#2965 + describe('restore staging', function () { + const STAGED = `${DB_NAME}-staged`; + const stagedDir = () => join(storageDir, STAGED); + const RECORDS = Array.from({ length: 500 }, (_, i) => [`k${i}`, { i, pad: 'x'.repeat(100) }]); + + afterEach(async function () { + // a refused online restore reloads the database it left intact, in this process too + await closeLoadedDatabases(); + if (checkRestoreState(stagedDir()) !== 'clear') completeRestore(beginRestore(stagedDir())); + rmSync(stagedDir(), { recursive: true, force: true }); + rmSync(backupDirForDatabase(STAGED), { recursive: true, force: true }); + rmSync(restoreStagingPath(stagedDir()), { recursive: true, force: true }); + rmSync(restoreReplacedPath(stagedDir()), { recursive: true, force: true }); + }); + + async function seed() { + const database = RocksDatabase.open(stagedDir()); + try { + await database.transaction(async (transaction) => { + transaction.useLog('audit').addEntry(Buffer.from('entry')); + for (const [key, value] of RECORDS) transaction.putSync(key, value); + }); + await database.flush(); + } finally { + database.close(); + } + const { backup_id: backupId } = await createBackupOffline(STAGED); + // diverge from the backup, so "intact" cannot be satisfied by a restore that happened to succeed + const after = RocksDatabase.open(stagedDir()); + try { + after.putSync('after-backup', { kept: true }); + } finally { + after.close(); + } + return backupId; + } + + function assertNoDebris() { + for (const path of [ + restoreStagingPath(stagedDir()), + restoreReplacedPath(stagedDir()), + restoreDiscardedPath(stagedDir()), + ]) { + assert.ok(!existsSync(path), `${path} must not be left behind`); + } + } + + function assertDestinationIntact() { + assert.strictEqual(checkRestoreState(stagedDir()), 'clear', 'a refused restore must not leave a marker'); + const database = RocksDatabase.open(stagedDir()); + try { + assert.deepStrictEqual(database.getSync('k7'), RECORDS[7][1]); + assert.deepStrictEqual(database.getSync('after-backup'), { kept: true }, 'the pre-restore data must survive'); + } finally { + database.close(); + } + assertNoDebris(); + } + + function assertRestoredFromBackup() { + const database = RocksDatabase.open(stagedDir()); + try { + assert.deepStrictEqual(database.getSync('k7'), RECORDS[7][1]); + assert.strictEqual(database.getSync('after-backup'), undefined, 'the restore took effect'); + } finally { + database.close(); + } + } + + // RocksDB's backup meta records a crc32c per file, verified as the file is copied out. + function crc32c(buffer) { + let crc = 0xffffffff; + for (const byte of buffer) { + crc ^= byte; + for (let bit = 0; bit < 8; bit++) crc = crc & 1 ? (crc >>> 1) ^ 0x82f63b78 : crc >>> 1; + } + return (crc ^ 0xffffffff) >>> 0; + } + + // A table whose footer declares a format_version this build does not know — what a newer + // binding's table format looks like to this one — with its checksum re-recorded, so the copy + // succeeds and only an open by this build can refuse it. + function writeUnsupportedTableFormat(backupId) { + const backupDir = backupDirForDatabase(STAGED); + const sst = readdirSync(join(backupDir, 'shared_checksum')).find((name) => name.endsWith('.sst')); + assert.ok(sst, 'precondition: the backup holds a table file'); + const sstPath = join(backupDir, 'shared_checksum', sst); + const contents = readFileSync(sstPath); + contents.writeUInt32LE(99, contents.length - 12); + writeFileSync(sstPath, contents); + const metaPath = join(backupDir, 'meta', String(backupId)); + const meta = readFileSync(metaPath, 'utf8'); + const entry = new RegExp(`^(shared_checksum/${sst} crc32 )\\d+$`, 'm'); + assert.match(meta, entry); + writeFileSync(metaPath, meta.replace(entry, `$1${crc32c(contents)}`)); + } + + function writeUnsupportedLogVersion(backupId) { + const store = join(backupDirForDatabase(STAGED), 'transaction_logs', String(backupId), 'audit'); + const log = readdirSync(store).find((name) => name.endsWith('.txnlog')); + const contents = readFileSync(join(store, log)); + contents[4] = 99; + writeFileSync(join(store, log), contents); + } + + function corruptBackupFile(backupId) { + const privateDir = join(backupDirForDatabase(STAGED), 'private', String(backupId)); + const manifest = readdirSync(privateDir).find((name) => name.startsWith('MANIFEST-')); + const contents = readFileSync(join(privateDir, manifest)); + contents.fill(0x5a, 0, Math.min(contents.length, 64)); + writeFileSync(join(privateDir, manifest), contents); + } + + const restores = { + offline: (backupId) => restoreBackupOffline(STAGED, backupId), + online: (backupId) => restoreBackup({ ...SU, database: STAGED, backup_id: backupId }), + }; + const unreadable = [ + ['a corrupt backup', corruptBackupFile, /checksum/i], + ['a table format this build cannot open', writeUnsupportedTableFormat, /format_version 99/], + ['a transaction log format this build cannot read', writeUnsupportedLogVersion, /version: 99/], + ]; + for (const [mode, restore] of Object.entries(restores)) { + for (const [label, damage, signature] of unreadable) { + it(`${mode}: refuses ${label} without touching the destination`, async function () { + this.timeout(30000); + const backupId = await seed(); + damage(backupId); + await assert.rejects( + restore(backupId), + (error) => signature.test(error.message) && /was not modified/.test(error.message) + ); + assertDestinationIntact(); + }); + } + + it(`${mode}: puts the database back when publication fails after moving it aside`, async function () { + this.timeout(30000); + const backupId = await seed(); + const realRename = fs.renameSync; + fs.renameSync = (from, to) => { + if (from === restoreStagingPath(stagedDir())) + throw Object.assign(new Error('injected EXDEV'), { code: 'EXDEV' }); + return realRename(from, to); + }; + // under TypeStrip the module binds the ESM builtin export, which only this re-syncs + syncBuiltinESMExports(); + try { + await assert.rejects(restore(backupId), /injected EXDEV/); + } finally { + fs.renameSync = realRename; + syncBuiltinESMExports(); + } + assertDestinationIntact(); + }); + } + + it('online: restores a valid backup through the swap', async function () { + this.timeout(30000); + const backupId = await seed(); + const result = await restoreBackup({ ...SU, database: STAGED, backup_id: backupId }); + assert.strictEqual(result.backup_id, backupId); + assert.strictEqual(checkRestoreState(stagedDir()), 'clear'); + await closeLoadedDatabases(); + assertRestoredFromBackup(); + assertNoDebris(); + }); + + it('online: a rescan during staging keeps the live database closable', async function () { + this.timeout(30000); + const backupId = await seed(); + const realRestore = backups.restore; + let release; + let parked; + const staged = new Promise((resolve) => (parked = resolve)); + backups.restore = async (...args) => { + parked(); + await new Promise((resolve) => (release = resolve)); + return realRestore.apply(backups, args); + }; + try { + const restoring = restoreBackup({ ...SU, database: STAGED, backup_id: backupId }); + await staged; + resetDatabases(); + release(); + await restoring; + } finally { + backups.restore = realRestore; + } + await closeLoadedDatabases(); + assertRestoredFromBackup(); + assertNoDebris(); + }); + + it('refuses a pre-restore copy that no restore marker accounts for', async function () { + this.timeout(30000); + const backupId = await seed(); + const replacedDir = restoreReplacedPath(stagedDir()); + for (const restore of Object.values(restores)) { + mkdirSync(replacedDir, { recursive: true }); + try { + await assert.rejects( + restore(backupId), + (error) => error.statusCode === 409 && /recorded no restore marker/.test(error.message) + ); + assert.ok(existsSync(replacedDir)); + } finally { + rmSync(replacedDir, { recursive: true, force: true }); + } + assertDestinationIntact(); + } + }); + + it('keeps the database directory mode across the swap', async function () { + if (process.platform === 'win32') this.skip(); + this.timeout(30000); + const backupId = await seed(); + chmodSync(stagedDir(), 0o700); + await restoreBackupOffline(STAGED, backupId); + assertRestoredFromBackup(); + assert.strictEqual(lstatSync(stagedDir()).mode & 0o777, 0o700); + assertNoDebris(); + }); + + // Nothing was displaced, so only the published engine says the destination changed; the marker + // has to outlive a blob restore that fails after it. + it('keeps the marker when a restore into a new target fails after publication', async function () { + this.timeout(30000); + const backupId = await seed(); + const TARGET = `${STAGED}-target`; + const targetDir = join(storageDir, TARGET); + const realRestoreBlobs = blobBackupModule.restoreBlobSnapshot; + blobBackupModule.restoreBlobSnapshot = async () => { + throw new Error('injected blob failure'); + }; + try { + await assert.rejects(restoreBackupOffline(STAGED, backupId, TARGET), /injected blob failure/); + assert.strictEqual(checkRestoreState(targetDir), 'incomplete'); + } finally { + blobBackupModule.restoreBlobSnapshot = realRestoreBlobs; + completeRestore(beginRestore(targetDir)); + rmSync(targetDir, { recursive: true, force: true }); + rmSync(restoreStagingPath(targetDir), { recursive: true, force: true }); + for (const root of getBlobPathsForDatabaseName(TARGET)) rmSync(root, { recursive: true, force: true }); + } + }); + + it('refuses a backup that would not fit beside the database before staging anything', async function () { + this.timeout(30000); + const backupId = await seed(); + const realStatfs = fs.statfsSync; + fs.statfsSync = (path, ...rest) => ({ ...realStatfs(path, ...rest), bavail: 1, bsize: 4096 }); + syncBuiltinESMExports(); + try { + for (const restore of Object.values(restores)) { + await assert.rejects( + restore(backupId), + (error) => error.statusCode === 507 && /was not modified/.test(error.message) + ); + assertDestinationIntact(); + } + } finally { + fs.statfsSync = realStatfs; + syncBuiltinESMExports(); + } + }); + + it('drops the published candidate before measuring space when rerunning a failed publication', async function () { + this.timeout(30000); + const backupId = await seed(); + const realRestoreBlobs = blobBackupModule.restoreBlobSnapshot; + blobBackupModule.restoreBlobSnapshot = async () => { + throw new Error('injected blob failure'); + }; + try { + await assert.rejects(restoreBackupOffline(STAGED, backupId), /injected blob failure/); + } finally { + blobBackupModule.restoreBlobSnapshot = realRestoreBlobs; + } + assert.ok(existsSync(restoreReplacedPath(stagedDir())), 'precondition: publication began'); + assert.ok(existsSync(stagedDir()), 'precondition: a candidate is published'); + + const realStatfs = fs.statfsSync; + let candidateAtCheck; + fs.statfsSync = (path, ...rest) => { + candidateAtCheck = existsSync(stagedDir()); + return realStatfs(path, ...rest); + }; + syncBuiltinESMExports(); + try { + await restoreBackupOffline(STAGED, backupId); + } finally { + fs.statfsSync = realStatfs; + syncBuiltinESMExports(); + } + assert.strictEqual(candidateAtCheck, false); + assert.strictEqual(checkRestoreState(stagedDir()), 'clear'); + assertRestoredFromBackup(); + assertNoDebris(); + }); + + it('online: refuses to drop a published candidate something still holds open', async function () { + this.timeout(30000); + const backupId = await seed(); + const realRestoreBlobs = blobBackupModule.restoreBlobSnapshot; + blobBackupModule.restoreBlobSnapshot = async () => { + throw new Error('injected blob failure'); + }; + try { + await assert.rejects(restoreBackupOffline(STAGED, backupId), /injected blob failure/); + } finally { + blobBackupModule.restoreBlobSnapshot = realRestoreBlobs; + } + assert.ok(existsSync(restoreReplacedPath(stagedDir())), 'precondition: publication began'); + const holder = RocksDatabase.open(stagedDir()); + try { + await assert.rejects( + restoreBackup({ ...SU, database: STAGED, backup_id: backupId }), + (error) => error.statusCode === 409 + ); + assert.ok(existsSync(join(stagedDir(), 'CURRENT')), 'the held candidate was not removed'); + } finally { + holder.close(); + } + await restoreBackup({ ...SU, database: STAGED, backup_id: backupId }); + assert.strictEqual(checkRestoreState(stagedDir()), 'clear'); + await closeLoadedDatabases(); + assertRestoredFromBackup(); + assertNoDebris(); + }); + + // `.replaced` is what tells a rerun that the database path holds only a candidate, so it may not + // outlive the restore that made it. + it('keeps the marker when the pre-restore copy cannot be retired', async function () { + this.timeout(30000); + const backupId = await seed(); + const realRename = fs.renameSync; + fs.renameSync = (from, to) => { + if (to === restoreDiscardedPath(stagedDir())) throw Object.assign(new Error('injected EIO'), { code: 'EIO' }); + return realRename(from, to); + }; + syncBuiltinESMExports(); + try { + await assert.rejects(restoreBackupOffline(STAGED, backupId), /injected EIO/); + } finally { + fs.renameSync = realRename; + syncBuiltinESMExports(); + } + assert.strictEqual(checkRestoreState(stagedDir()), 'incomplete'); + await restoreBackupOffline(STAGED, backupId); + assert.strictEqual(checkRestoreState(stagedDir()), 'clear'); + assertRestoredFromBackup(); + assertNoDebris(); + }); + + it('refuses a database directory that is a mount point before staging anything', async function () { + this.timeout(30000); + const backupId = await seed(); + const realStat = fs.statSync; + fs.statSync = (path, ...rest) => { + const stats = realStat(path, ...rest); + return path === stagedDir() + ? Object.assign(Object.create(Object.getPrototypeOf(stats)), stats, { dev: stats.dev + 1 }) + : stats; + }; + syncBuiltinESMExports(); + try { + await assert.rejects( + restoreBackupOffline(STAGED, backupId), + (error) => error.statusCode === 400 && /mount point/.test(error.message) + ); + } finally { + fs.statSync = realStat; + syncBuiltinESMExports(); + } + assertDestinationIntact(); + }); + + it('refuses a symlinked database directory before staging anything', async function () { + this.timeout(30000); + const backupId = await seed(); + const realDir = `${stagedDir()}-real`; + renameSync(stagedDir(), realDir); + symlinkSync(realDir, stagedDir(), 'dir'); + try { + await assert.rejects( + restoreBackupOffline(STAGED, backupId), + (error) => error.statusCode === 400 && /symbolic link/.test(error.message) + ); + assert.ok(lstatSync(stagedDir()).isSymbolicLink(), 'the link is untouched'); + assert.strictEqual(checkRestoreState(stagedDir()), 'clear'); + assertNoDebris(); + } finally { + rmSync(stagedDir(), { force: true }); + renameSync(realDir, stagedDir()); + } + }); + + // Killed between moving the database aside and publishing its replacement: the database path is + // empty and `.replaced` is the only copy of it. + async function killBetweenRenames(backupId) { + const child = spawn( + process.execPath, + [ + '-e', + ` + const fs = require('node:fs'); + const env = require(${JSON.stringify(require.resolve('#src/utility/environment/environmentManager'))}); + env.initSync(); + env.setProperty('storage.backupPath', ${JSON.stringify(getBackupsRoot())}); + const rename = fs.renameSync; + fs.renameSync = (from, to) => { + rename(from, to); + if (to.endsWith('.replaced')) process.kill(process.pid, 'SIGKILL'); + }; + require('node:module').syncBuiltinESMExports(); + const backup = require(${JSON.stringify(require.resolve('#src/dataLayer/rocksdbBackup'))}); + backup.restoreBackupOffline(${JSON.stringify(STAGED)}, ${backupId}) + .then(() => process.exit(2), (error) => { console.error(error); process.exit(3); }); + `, + ], + { timeout: 20000, env: { ...process.env, STORAGE_PATH: storageDir } } + ); + let stderr = ''; + child.stderr.on('data', (data) => { + stderr += data; + }); + const signal = await new Promise((resolve, reject) => { + child.once('error', reject); + child.once('exit', (code, exitSignal) => + exitSignal ? resolve(exitSignal) : reject(new Error(`child exited ${code}: ${stderr}`)) + ); + }); + assert.strictEqual(signal, 'SIGKILL'); + assert.strictEqual(checkRestoreState(stagedDir()), 'incomplete'); + assert.ok(!existsSync(stagedDir()), 'precondition: killed with the database moved aside'); + assert.ok(existsSync(restoreReplacedPath(stagedDir()))); + } + + it('keeps the moved-aside database through a failed rerun', async function () { + this.timeout(60000); + const good = await seed(); + await killBetweenRenames(good); + corruptBackupFile(good); + await assert.rejects(restoreBackupOffline(STAGED, good), /rerun restore_backup to recover/); + assert.strictEqual(checkRestoreState(stagedDir()), 'incomplete', 'a failed recovery keeps the marker'); + assert.ok(existsSync(restoreReplacedPath(stagedDir())), 'and the only copy of the database'); + }); + + it('reruns a publication interrupted between its renames, keeping the moved-aside mode', async function () { + this.timeout(60000); + const backupId = await seed(); + chmodSync(stagedDir(), 0o700); + await killBetweenRenames(backupId); + await restoreBackupOffline(STAGED, backupId); + assert.strictEqual(checkRestoreState(stagedDir()), 'clear'); + assertRestoredFromBackup(); + if (process.platform !== 'win32') assert.strictEqual(lstatSync(stagedDir()).mode & 0o777, 0o700); + assertNoDebris(); + }); + }); + describe('online restore_backup validation', function () { it('rejects database=system with a pointer at running it offline', async function () { for (const fn of [validateRestoreBackup, restoreBackup]) { diff --git a/unitTests/resources/databases.test.js b/unitTests/resources/databases.test.js index 9af39ea707..0a7323bc2e 100644 --- a/unitTests/resources/databases.test.js +++ b/unitTests/resources/databases.test.js @@ -267,7 +267,7 @@ describe('dropDatabase restore serialization', () => { assert.ok(existsSync(reservedDir), 'the reserved dir itself is left in place (used for lifecycle metadata)'); }); - it('never loads a database whose pre-atomic restoring marker is empty', function () { + it('never loads a database whose pre-atomic restoring marker is empty', async function () { const databaseName = 'corrupt-marker-startup-test'; const Table = table({ table: 'CorruptMarker', @@ -279,6 +279,8 @@ describe('dropDatabase restore serialization', () => { abandonRestore(beginRestore(rootStore.path)); writeFileSync(restoringMarkerPath(rootStore.path), ''); + // Startup has nothing open; a root this thread still holds is kept loaded for the restore's close. + await closeDatabase(databaseName); try { resetDatabases(); assert.strictEqual(