diff --git a/DESIGN.md b/DESIGN.md index 6610c97638..f67c4a8d8d 100644 --- a/DESIGN.md +++ b/DESIGN.md @@ -17,7 +17,8 @@ Index of the design notes for the harper core: one line per note, grouped by the - [`Table.ts` — section map](resources/DESIGN.md#tablets--section-map) — Section markers for the 4.7K-line `makeTable()` factory. - ["Where is X" cheat sheet](resources/DESIGN.md#where-is-x-cheat-sheet) — Symbol lookup for the read/write path, audit, subscriptions and schema. - [Full-text declarations and reader snapshots](resources/DESIGN.md#full-text-declarations-and-reader-snapshots) — Declaration names stay separate from stored attributes; a query retains one native reader across every page. -- [Audit retention floor](resources/DESIGN.md#audit-retention-floor) — A saved audit cursor below the floor must resync; the floor is internal, and `Table.commit` skips the out-of-order walk below it. +- [Audit retention floor](resources/DESIGN.md#audit-retention-floor) — The floor records what pruning removed; it is internal, and `Table.commit` skips the out-of-order walk below it. +- [Database generation and resumable positions](resources/DESIGN.md#database-generation-and-resumable-positions) — Every copy path stamps a new generation before the copy is readable; a position resumes only if it names it and sits at or above its resume floor. - [Path routing & parameterised routes](resources/DESIGN.md#path-routing--parameterised-routes) — How resource paths and route parameters resolve. - [Persisted relationship catalog](resources/DESIGN.md#persisted-relationship-catalog) — Where relationship definitions are stored and rebuilt. - [Typed, discoverable resources (code-first schema + request contract)](resources/DESIGN.md#typed-discoverable-resources-code-first-schema--request-contract) — Declaring schema and request contracts from code. diff --git a/bin/copyDb.ts b/bin/copyDb.ts index 7f1d35a3e8..1105296c5b 100644 --- a/bin/copyDb.ts +++ b/bin/copyDb.ts @@ -15,7 +15,7 @@ import OpenEnvironmentObject from '../utility/lmdb/OpenEnvironmentObject.ts'; import { OpenDBIObject } from '../utility/lmdb/OpenDBIObject.ts'; import { INTERNAL_DBIS_NAME, AUDIT_STORE_NAME } from '../utility/lmdb/terms.ts'; import { CONFIG_PARAMS, DATABASES_DIR_NAME, MIGRATING_DIR_SUFFIX } from '../utility/hdbTerms.ts'; -import { AUDIT_STORE_OPTIONS, auditRetention } from '../resources/auditStore.ts'; +import { AUDIT_STORE_OPTIONS, auditRetention, stampDatabaseGeneration } from '../resources/auditStore.ts'; import { blobsReadmeContent, copyBlobRootsByIndex } from '../dataLayer/blobBackup.ts'; import { describeSchema } from '../dataLayer/schemaDescribe.ts'; import { updateConfigValue } from '../config/configUtils.ts'; @@ -983,6 +983,10 @@ export async function copyDbToRocks(sourceRootStore, sourceDatabase: string, tar targetRootStore.putSync(REMOTE_NODE_IDS_KEY, asBinary(idMappingBytes)); } + // flushed so the stamp is durable before the caller publishes the staging directory + stampDatabaseGeneration(targetRootStore, { carriesLog: false }); + await targetRootStore.flush({ allowWriteStall: true }); + console.log('migrated database ' + sourceDatabase + ' to RocksDB'); } finally { endPendingMigrationBlobSaves(); diff --git a/dataLayer/DESIGN.md b/dataLayer/DESIGN.md index dc37e7bedf..c62fc43d1b 100644 --- a/dataLayer/DESIGN.md +++ b/dataLayer/DESIGN.md @@ -149,6 +149,11 @@ 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. - **`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 "successful" restore (or vice versa). `dropDatabase` takes the restore lock for every RocksDB or diff --git a/dataLayer/rocksdbBackup.ts b/dataLayer/rocksdbBackup.ts index da32308795..f677696b74 100644 --- a/dataLayer/rocksdbBackup.ts +++ b/dataLayer/rocksdbBackup.ts @@ -11,6 +11,7 @@ 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 { getDatabases, resolveDatabasePath } from '../resources/databases.ts'; +import { stampDatabaseDirectory } from '../resources/auditStore.ts'; import { type BlobCaptureDisposition, classifyBlobFileForCapture, @@ -601,6 +602,7 @@ export async function restoreBackup(request: any) { if (manifest.blobs) { await restoreBlobSnapshot(backupDir, backupId, databaseName, getBlobPathsForDatabaseName(databaseName)); } + await stampDatabaseDirectory(databaseDir, { carriesLog: true }); } catch (error: any) { // 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 @@ -1101,6 +1103,7 @@ export async function restoreBackupOffline( getBlobPathsForDatabaseName(targetDatabase ?? databaseName) ); } + await stampDatabaseDirectory(databaseDir, { carriesLog: true }); } catch (error: any) { // 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, diff --git a/resources/DESIGN.md b/resources/DESIGN.md index 758c786186..e276367c40 100644 --- a/resources/DESIGN.md +++ b/resources/DESIGN.md @@ -188,11 +188,10 @@ Index waits use the ordinary transaction timeout; they do not renew it. The adap ## Audit retention floor `Table.subscribe`'s `startTime` replay just begins wherever the audit log now begins, so a consumer -resuming below the retention horizon is silently handed a short replay. The floor is the primitive -that makes that detectable (harper#2447). It is internal, with deliberately no public accessor, and **no -resume path consumes it yet**: harper#2448 is to put the check inside `Table.subscribe` itself — the same shape as -replication's `shouldForceBaseCopyForRetention`, and the only one where the floor cannot move between -being read and being acted on. Until then the short replay above is unchanged. +resuming below the retention horizon is silently handed a short replay. The floor records what +pruning removed (harper#2447); it is internal, with deliberately no public accessor. **A resume is +not checked against it** but against the database generation (next section): this floor also steers +reconciliation, survives a restore, and absorbs at `Infinity`. **The one consumer today is not a resume**: `Table.commit`'s out-of-order reconciliation reads the floor before entering the audit walk (harper#2642). The walk terminates at the incoming write only by @@ -226,50 +225,31 @@ Things that are easy to get wrong here: surviving entry would do, because they prune a database-wide time prefix. `Table.deleteHistory` removes one table's entries from a database-scoped log, so a sibling's entry survives _below_ the newest entry it removed, and a floor taken from that survivor certifies cursors over removed history. -- **The record's presence is the trust marker.** `Symbol.for('audit-floor')` is a different key from - `last-removed`, which is still live and still maintained by the LMDB retention loop (#2338 hardened - its write path and added tests for the retry-carry — do not remove it). They coexist because they - answer different questions: `last-removed` records where the LMDB loop got to, after the fact, - while the floor is written ahead of every one of the five prune paths and its commit is verified. - A value found under `last-removed` therefore cannot be told apart from one carrying those - guarantees, which is why the floor needs its own key rather than reusing it. -- **A store with no floor record is a store whose retention history we cannot account for.** That - includes the empty audit store an LMDB→RocksDB migration leaves behind, since `bin/copyDb.ts` - deliberately does not migrate it, and the audit-DBI-less result of a table-scoped backup taken - without `include_audit` — so `openAuditStore` stamps `max(Date.now(), newest retained key)` as a - one-time resync epoch. There is no permissive-baseline case: creating the audit DBI proves the - DBI was absent, not that the database is new. -- **That epoch is a guess, and it is recorded as one.** Its bound is surviving state, which cannot see - history a selective prune already removed: a legacy `deleteHistory` takes one table's entries out of - the shared log, so a table that held the newest entries can leave the newest _survivor_ older than - entries that are gone, and a clock rolled back between the two stamps a floor below them (#2458). - Refusing to stamp is worse — `AUDIT_FLOOR_UNKNOWN` is absorbing (`raiseAuditFloor` cannot lift it, - `establishAuditFloor` skips any existing record), so it would make every upgraded deployment fail - closed forever. So `establishAuditFloor` writes the epoch under `Symbol.for('audit-floor-bootstrap')` - first, then stamps the floor from what that record holds. - - **The record's presence is the signal; comparing it against the floor is not.** A store carrying one - has an unverified pre-tracking window for as long as the record exists, however far the floor has - since moved — a prune raising the floor above the epoch certifies only what that prune removed, and - says nothing about history removed before tracking began, which may sit _above_ the epoch, since that - is precisely what the guess could not see. Worked example: a v4-era `deleteHistory` removes tableA up - to t=1000 while sibling tableB's newest survivor is 900; a rolled-back clock stamps bootstrap=900 and - floor=900; a later retention pass raises the floor to 950. A repair keyed on `floor > bootstrap` would - read 950 > 900, call it earned, and leave a consumer at cursor 970 certified over tableA's missing - 950–1000. So the mark is retired by a database generation (#2451), never by a floor that climbed past - it; what the recorded _value_ is for is telling that repair how far the guess reached. - - Two properties it does depend on. **Ordering:** the record is written first, so a crash between the - two writes leaves a record with no floor, which the next open retries because the early return tests - the _floor_. **Undecodable bytes are overwritten** rather than kept — unlike the floor, where a - present record may be a deliberate `AUDIT_FLOOR_UNKNOWN` and rewriting it would lower a floor. - Keeping torn bytes pinned the store to unknown _forever_: the resolver skipped the write because a - record existed, the read back failed identically on every later open, and no retry could succeed. - -- **`getHistory` is not in the floor's time domain.** The floor is an audit-log key, which is what - `subscribe`'s events carry as `localTime`; `getHistory` reports each entry's origin `version` under - that same name, and a backdated or replicated write makes the two differ. A cursor saved from - `getHistory` cannot be compared against the floor. +- **The record's presence is the trust marker.** `Symbol.for('audit-floor')` is not `last-removed`, + which the LMDB retention loop still maintains (#2338 — do not remove it): that one records where the + loop got to, after the fact, so a value there cannot be told apart from a write-ahead, verified floor. +- **A store with no floor record is one whose retention history we cannot account for** (the empty + audit store a migration leaves, a table-scoped backup without `include_audit`), so `openAuditStore` + stamps `max(Date.now(), newest retained key)` as a one-time resync epoch. There is no + permissive-baseline case: creating the audit DBI proves only that it was absent. +- **That epoch is a guess, and it is recorded as one.** Surviving state cannot see history a legacy + selective prune removed, so a rolled-back clock can stamp it below entries that are gone (#2458). + Refusing to stamp is worse — `AUDIT_FLOOR_UNKNOWN` is absorbing — so `establishAuditFloor` writes the + epoch under `Symbol.for('audit-floor-bootstrap')` first, then stamps the floor from that record. + + **The record's presence is the signal; comparing it against the floor is not.** Worked example: a + v4-era `deleteHistory` removes tableA up to t=1000 while tableB's newest survivor is 900; a + rolled-back clock stamps bootstrap=floor=900; a later pass raises the floor to 950 — and a cursor at + 970 still sits over tableA's missing 950–1000. No timestamp can close that window, which is why + resumable positions are bound to a generation instead: a position naming one postdates tracking. + + **Ordering:** the record is written first, so a crash between the two writes leaves a record with no + floor, which the next open retries. **Undecodable bytes are overwritten** (unlike the floor, where a + present record may be a deliberate unknown): keeping them pinned the store to unknown forever. + +- **`getHistory` is not in the floor's time domain.** The floor is an audit-log key (`subscribe`'s + `localTime`); `getHistory` reports each entry's origin `version` under that name, which a backdated + or replicated write makes differ, so its cursors cannot be compared against the floor. - **On RocksDB the floor tracks the configured retention horizon, not retained reality.** Whole-log-file purge granularity means the branch cannot know which entries a purge will drop, and the floor is written first, so each pass advances it to `Date.now() - auditRetention/(1+priority²)` whether a @@ -284,11 +264,22 @@ Things that are easy to get wrong here: - **Untrustworthy metadata resolves to `Infinity`, not to a number.** A wrong-length record, or eight bytes decoding to NaN/negative, must not become a floor: `cursor < NaN` is false, so a consumer spelling the check that way would read corrupt metadata as safe. -- **A restore is outside what the floor can see.** `restore_backup` reinstalls the backup's floor - along with everything else, so a cursor from after the backup point reads as safe against it. The - audit floor is one of three carriers of resumable state a restore rolls back (record versions and - per-node `Symbol.for('seq')` records are the others), so this wants a database-level generation - rather than a fix in this one field — harper#2451. +- **A copy keeps this floor honest for the log it carried.** A restore carries its log, so the floor + stands; a branch or migration carries none, so its stamp raises a finite floor to the copy's epoch. + Which history the database is belongs to the generation. + +--- + +## Database generation and resumable positions + +Every copy path gives the copy a generation of its own before anything can read it (harper#2451), so no live subscription or persisted position carries over from the source. **Invariant: no readable copy carries its source's generation; a position resumes only if it names the current generation, its cursor is finite, and no prune within the generation reached above it** (`isResumablePosition`). RocksDB only: an LMDB database carries neither a generation nor a resume floor, so no position resumes against it, and a caller applies the check to RocksDB databases alone. + +- **Two records beside the floor:** the generation (a 16-byte random id and a float64 epoch, 0 for genesis) and the resume floor (highest prune cutoff in the generation). `raiseAuditFloor` raises both in one verified transaction; the resume floor is not absorbed by an unknown floor. +- **Stamped into the copy before publication:** restore (before `completeRestore`), branch and migration (before their renames), via a private open plus an engine flush — a root-store write is not power-loss durable and directory fsync is best-effort. `harper copydb` and `compactOnStart` copy LMDB databases, so they stamp nothing. +- **An ordinary open never repairs:** it adopts the record, mints genesis as a compare-and-set on absence, or leaves the handle without one (every resume refused). The resume floor is read as never below a finite audit floor, since an older binary's prunes raise only the audit floor. +- **No scalar mode.** A position without an id is never resumable; bind an id only to a position established within the generation. Cursors must be progress-based: a snapshot's newest in-scope key can sit below a floor that retention advances on a quiet database, and would be refused forever. +- **Live subscriptions:** the per-path registry outlives the store, but a subscription from before a reopen can never deliver again (its commit listener and table stores belong to the closed handle). Every open first records its audit store as the only handle a registration on the path may use, then detaches the registry and ends each subscription: `DatabaseGenerationChangedError` for another or an unknown generation on RocksDB, and otherwise (the same generation, or any LMDB reopen) the retryable `DatabaseClosingError`. A registration through an earlier or closed handle, even from inside that close, is refused with the same pair of errors. +- **Not covered:** copies no generation-aware code made, a pre-generation binary pruning while the audit floor is unknown, keys reissued below a cursor after a clock rollback across a restart, cross-node identity (an id is per database per node), and subscription teardown and the handle check on a legacy LMDB `auditPath` root, which is reopened on every metadata read. --- diff --git a/resources/auditStore.ts b/resources/auditStore.ts index bf01a8988c..38ee852f8b 100644 --- a/resources/auditStore.ts +++ b/resources/auditStore.ts @@ -1,3 +1,4 @@ +import { randomBytes } from 'node:crypto'; import { readKey, writeKey } from 'ordered-binary'; import { initSync, get as envGet } from '../utility/environment/environmentManager.ts'; import { AUDIT_STORE_NAME } from '../utility/lmdb/terms.ts'; @@ -12,7 +13,8 @@ import { onStorageReclamation } from '../server/storageReclamation.ts'; import { RocksDatabase } from '@harperfast/rocksdb-js'; import { asBinary } from 'lmdb'; import { RocksTransactionLogStore } from './RocksTransactionLogStore.ts'; -import { isReadOnlyMode } from './databases.ts'; +import { endSubscriptionsFromEarlierHandles } from './transactionBroadcast.ts'; +import { isReadOnlyMode, openRocksDatabase } from './databases.ts'; /** * This module is responsible for the binary representation of audit records in an efficient form. @@ -138,9 +140,10 @@ const FLOAT_BUFFER = new Uint8Array(FLOAT_TARGET.buffer); const AUDIT_FLOOR_KEY = Symbol.for('audit-floor'); // The epoch `establishAuditFloor` stamped, never raised or removed. Its PRESENCE marks the floor as // unverified provenance — a guess bounded by what survived, blind to history a legacy prune removed -// before tracking began — until a database generation (#2451) retires the mark. No comparison does: -// a later prune certifies only what it removed, so `floor > bootstrap` says nothing about the older -// gap. The value records how far the guess reached, for that repair. See `establishAuditFloor`. +// before tracking began. No comparison retires the mark: a later prune certifies only what it +// removed, so `floor > bootstrap` says nothing about the older gap. Resume never relies on it, since a +// position is bound to a database generation and everything after its cursor was written after +// tracking began. The value records how far the guess reached. See `establishAuditFloor`. const AUDIT_FLOOR_BOOTSTRAP_KEY = Symbol.for('audit-floor-bootstrap'); /** * The floor's own eight bytes, deliberately NOT the FLOAT_TARGET/FLOAT_BUFFER pair the `last-removed` @@ -152,6 +155,25 @@ const FLOOR_TARGET = new Float64Array(1); const FLOOR_BUFFER = new Uint8Array(FLOOR_TARGET.buffer); /** No trustworthy floor: the highest possible floor, so every cursor compares as stale. */ const AUDIT_FLOOR_UNKNOWN = Infinity; +/** + * Which copy of this database's history this is: a sixteen-byte random id, then the float64 time the + * generation began (0 for genesis). Every path that publishes a copy stamps a fresh one first + * (`stampDatabaseGeneration`), so a position naming an older id is refused (harper#2451). RocksDB only, + * like the resume floor: an LMDB database carries neither, so no position is resumable against it. + */ +const DATABASE_GENERATION_KEY = Symbol.for('database-generation'); +const GENERATION_ID_BYTES = 16; +const GENERATION_RECORD_BYTES = GENERATION_ID_BYTES + 8; +/** + * The highest prune cutoff within the current generation, encoded like the floor. Kept apart from the + * floor, which also steers `Table.commit`'s reconciliation: resume validity must restart at a copy and + * must not be absorbed by the floor's unknown sentinel, and doing either to the floor changes merges. + */ +const AUDIT_RESUME_FLOOR_KEY = Symbol.for('audit-resume-floor'); + +function isRocksStore(store: any): boolean { + return store instanceof RocksTransactionLogStore || store instanceof RocksDatabase; +} /** Last resort on a detached path: a failing log sink must not itself become an unhandled rejection. */ function warnContained(message: string, error: unknown) { @@ -191,6 +213,8 @@ export function openAuditStore(rootStore) { rootStore.auditStore = auditStore; auditStore.rootStore = rootStore; establishAuditFloor(auditStore); + establishDatabaseGeneration(auditStore); + endSubscriptionsFromEarlierHandles(auditStore, isRocksStore(auditStore)); auditStore.tableStores = []; const deleteCallbacks = []; auditStore.addDeleteRemovalCallback = function (tableId, table, callback) { @@ -508,13 +532,13 @@ function decodeAuditFloor(stored: any): number { } /** - * Did the floor write land? A record has to be PRESENT, not merely decode to the value we wrote: - * `decodeAuditFloor(undefined)` is the unknown sentinel too, so on a floorless store — where the - * resolver writes exactly that sentinel — comparing decoded values alone reported a commit for a - * write that never happened, and the caller pruned with nothing persisted. + * Did a metadata write land? A record has to be PRESENT with exactly the bytes written, not merely + * decode to the same value: `decodeAuditFloor(undefined)` is the unknown sentinel too, so on a + * floorless store — where the resolver writes exactly that sentinel — comparing decoded values alone + * reported a commit for a write that never happened, and the caller pruned with nothing persisted. */ -function floorWriteLanded(stored: any, floor: number): boolean { - return stored !== undefined && decodeAuditFloor(stored) === floor; +function writeLanded(stored: any, written: Uint8Array): boolean { + return stored !== undefined && stored.byteLength === written.byteLength && Buffer.compare(stored, written) === 0; } /** Own eight bytes per write: the store must never be handed a live view of the reused module buffer. */ @@ -538,50 +562,66 @@ function updateAuditFloor( resolve: (current: number, recorded: boolean) => number | undefined, key: symbol = AUDIT_FLOOR_KEY ): void { - // A legacy `auditPath` layout is opened as its own standalone LMDB root (databases.ts) and has no - // `.rootStore`, so it owns the transaction itself. - const transactionOwner = auditStore?.rootStore ?? auditStore; + commitAuditMetadata( + auditStore, + (read) => { + const stored = read(key); + const floor = resolve(decodeAuditFloor(stored), stored !== undefined); + return floor === undefined ? undefined : [[key, encodeAuditFloor(floor)]]; + }, + 'audit retention floor' + ); +} + +/** + * Write audit metadata records under one store transaction, all or none: every write is read back, + * and a mismatch throws inside the transaction so it aborts — returning false commits on LMDB. + */ +function commitAuditMetadata( + store: any, + plan: (read: (key: symbol) => any) => Array<[symbol, Uint8Array]> | undefined, + what: string +): void { + // A copy being stamped passes its RocksDB root directly. A legacy `auditPath` layout is opened as its + // own standalone LMDB root (databases.ts) and has no `.rootStore`, so it owns the transaction itself. + const onRocksDB = isRocksStore(store); + const transactionOwner = store instanceof RocksDatabase ? store : (store?.rootStore ?? store); if (!transactionOwner?.transactionSync) - throw new Error('Cannot record the audit retention floor: this database has no audit store'); - // Both branches read their own write back and report `false` on mismatch, and the caller demands an - // explicit `true`. Both halves are load-bearing: a RocksDB transactionSync returns undefined for a - // swallowed abort rather than throwing (see RecordEncoder.saveStructures), and a write that fails - // without throwing is otherwise indistinguishable from one that landed — a caller pruning against a - // floor never recorded. Reads inside a write transaction see their own writes on both engines, so - // the read-back observes what commit will make durable. - const committed = - auditStore instanceof RocksTransactionLogStore - ? transactionOwner.transactionSync( - (txn) => { - const stored = txn.getBinarySync(key); - const floor = resolve(decodeAuditFloor(stored), stored !== undefined); - if (floor !== undefined) { - txn.putSync(key, asBinary(encodeAuditFloor(floor))); - if (!floorWriteLanded(txn.getBinarySync(key), floor)) return false; - } - return true; - }, - { retryOnBusy: true } - ) - : transactionOwner.transactionSync(() => { - const stored = auditStore.getBinary(key); - const floor = resolve(decodeAuditFloor(stored), stored !== undefined); - // `put` rather than `putSync`, and inside the transaction: lmdb's putSync is - // `put(...) === SYNC_PROMISE_SUCCESS`, so it drops whatever put returns, and a rejected put - // would leak with no owner. Within a write transaction put writes synchronously and returns - // an already-resolved sentinel, so the value is visible immediately either way and this - // only takes ownership of the failure case. - // asBinary: a legacy standalone audit root's encoder has no Uint8Array passthrough, so raw - // bytes would reach createAuditEntry and throw. This bypasses both encoders. - if (floor !== undefined) { - auditStore - .put(key, asBinary(encodeAuditFloor(floor))) - ?.catch?.((error) => warnContained('Error writing the audit retention floor', error)); - if (!floorWriteLanded(auditStore.getBinary(key), floor)) return false; + throw new Error(`Cannot record the ${what}: this database has no audit store`); + const notCommitted = () => new Error(`The ${what} transaction did not commit`); + // The caller still demands an explicit `true`: a RocksDB transactionSync returns undefined for a + // swallowed abort rather than throwing (see RecordEncoder.saveStructures). Reads inside a write + // transaction see their own writes on both engines, so the read-back observes what commit will make + // durable. + const committed = onRocksDB + ? transactionOwner.transactionSync( + (txn) => { + const writes = plan((key) => txn.getBinarySync(key)); + if (writes) { + for (const [key, bytes] of writes) txn.putSync(key, asBinary(bytes)); + for (const [key, bytes] of writes) if (!writeLanded(txn.getBinarySync(key), bytes)) throw notCommitted(); } return true; - }); - if (committed !== true) throw new Error('The audit retention floor transaction did not commit'); + }, + { retryOnBusy: true } + ) + : transactionOwner.transactionSync(() => { + const writes = plan((key) => store.getBinary(key)); + // `put` rather than `putSync`, and inside the transaction: lmdb's putSync is + // `put(...) === SYNC_PROMISE_SUCCESS`, so it drops whatever put returns, and a rejected put + // would leak with no owner. Within a write transaction put writes synchronously and returns + // an already-resolved sentinel, so the value is visible immediately either way and this + // only takes ownership of the failure case. + // asBinary: a legacy standalone audit root's encoder has no Uint8Array passthrough, so raw + // bytes would reach createAuditEntry and throw. This bypasses both encoders. + if (writes) { + for (const [key, bytes] of writes) + store.put(key, asBinary(bytes))?.catch?.((error) => warnContained(`Error writing the ${what}`, error)); + for (const [key, bytes] of writes) if (!writeLanded(store.getBinary(key), bytes)) throw notCommitted(); + } + return true; + }); + if (committed !== true) throw notCommitted(); } /** @@ -666,25 +706,42 @@ export function raiseAuditFloor(auditStore: any, cutoff: number): void { // database, a cutoff below one a wider prune already set), and taking the env write lock to // discover that serializes every worker's boot and reclamation on it. The in-transaction guards // below stay authoritative. - // Skips only the case it can prove is a no-op: a record that exists and already sits at or above - // the cutoff. An absent record is NOT decided here — the presence question is settled inside the - // transaction below, because another worker's establishAuditFloor can land between this read and - // that write. + // Skips only the case it can prove is a no-op: the floor, and on RocksDB the resume record, exist and + // already sit at or above the cutoff. An absent record is NOT decided here — the presence question is + // settled inside the transaction below, because another worker's establishAuditFloor can land between + // this read and that write. + const resumeTracked = isRocksStore(auditStore); if (auditStore?.getBinary) { const stored = auditStore.getBinary(AUDIT_FLOOR_KEY); - if (stored !== undefined && !(cutoff > decodeAuditFloor(stored))) return; + const resume = resumeTracked ? auditStore.getBinary(AUDIT_RESUME_FLOOR_KEY) : undefined; + if ( + stored !== undefined && + !(cutoff > decodeAuditFloor(stored)) && + (!resumeTracked || (resume !== undefined && !(cutoff > decodeAuditFloor(resume)))) + ) + return; } - updateAuditFloor(auditStore, (current, recorded) => { - // Still no record, and we are about to prune: persist the unknown sentinel. Leaving no marker - // lets the next open stamp a FINITE epoch, and a prune bound above that epoch (a future - // `endTime`, or a rolled-back clock) then certifies cursors whose history this prune deleted. - // Unknown is the honest value, because a store with no record may have been pruned before this - // run too. - if (!recorded) return AUDIT_FLOOR_UNKNOWN; - // A record appeared while we were getting here, so this is an ordinary monotonic raise: pruning - // to `cutoff` against a floor left below it is exactly the silent gap this function prevents. - return cutoff > current ? cutoff : undefined; - }); + commitAuditMetadata( + auditStore, + (read) => { + const writes: Array<[symbol, Uint8Array]> = []; + const stored = read(AUDIT_FLOOR_KEY); + // Still no record, and we are about to prune: persist the unknown sentinel. Leaving no marker + // lets the next open stamp a FINITE epoch, and a prune bound above that epoch (a future + // `endTime`, or a rolled-back clock) then certifies cursors whose history this prune deleted. + // Unknown is the honest value, because a store with no record may have been pruned before this + // run too. + if (stored === undefined) writes.push([AUDIT_FLOOR_KEY, encodeAuditFloor(AUDIT_FLOOR_UNKNOWN)]); + else if (cutoff > decodeAuditFloor(stored)) writes.push([AUDIT_FLOOR_KEY, encodeAuditFloor(cutoff)]); + if (resumeTracked) { + const resume = read(AUDIT_RESUME_FLOOR_KEY); + if (resume === undefined || cutoff > decodeAuditFloor(resume)) + writes.push([AUDIT_RESUME_FLOOR_KEY, encodeAuditFloor(cutoff)]); + } + return writes; + }, + 'audit retention floor' + ); } /** @@ -715,8 +772,9 @@ export function raiseAuditFloor(auditStore: any, cutoff: number): void { * unverified pre-tracking window for as long as it exists, however far the floor has since moved: a * prune raising the floor above the epoch certifies only what that prune removed, and says nothing * about history removed before tracking began — which may sit above the epoch, since that is exactly - * the case the guess cannot see. Retiring the mark takes a database generation (harper#2451), not a - * floor that has climbed past it. + * the case the guess cannot see. A floor that has climbed past it does not retire it. Resumable + * positions never depend on it: each is bound to a database generation minted after tracking began, + * so the pre-tracking window lies below every one of them (`getDatabaseGeneration`). * * Ordering. The provenance record is written **first**, so a crash between the two writes leaves a * record with no floor — which the next open retries, since the early return above tests the floor. @@ -787,10 +845,10 @@ export function establishAuditFloor(auditStore: any): void { /** * The floor of this database's retained audit history: every audit entry at or after the returned - * time is still retained, so a consumer whose last-processed audit-log cursor is `>=` it can resume - * incrementally, and one below it must resync — it may have lost nothing, but the floor cannot - * certify it either way. Returns `Infinity` when - * the floor is unknown, which fails closed — no cursor compares as safe. + * time is still retained, and below it the log may have lost anything. Returns `Infinity` when the + * floor is unknown, which fails closed — no cursor compares as safe. A resume is decided by the + * database generation and its resume floor (`getDatabaseGeneration`) instead: this floor survives a + * restore, and an unknown one is absorbing. * * **One exception, and it is the only one: history removed before the floor existed.** Every prune * that runs with a floor recorded is covered — it raises the floor first, so it cannot remove an entry @@ -798,7 +856,9 @@ export function establishAuditFloor(auditStore: any): void { * (`establishAuditFloor`), and a legacy `Table.deleteHistory` that removed a table's newest entries, * followed by a clock rollback, leaves that stamp below history that is gone; a cursor in the window * is then certified over the gap. So the guarantee is one-directional for tracked prunes and silent - * about untracked ones. Only a database generation can close that (harper#2451). + * about untracked ones. Resumable positions do not face it: they are checked against the database + * generation and its resume floor (`getDatabaseGeneration`), not against this floor, which also steers + * `Table.commit`'s reconciliation and so cannot move in the direction a cursor check would want. * * The time domain is the audit-log key: what `subscribe`'s events carry as `localTime` and what MQTT * durable sessions persist as `startTime`, so those compare against the floor directly. @@ -817,8 +877,9 @@ export function establishAuditFloor(auditStore: any): void { * * **A moment-in-time observation.** Retention can advance between this call and whatever the caller * does with the answer, so a check-then-resume sequence has a window where the floor moves under it. - * Closing that requires validating the cursor inside the resume itself (harper#2448); until then a - * lost race degrades to the truncation that happens today, never to anything worse. + * Closing that requires validating the position inside the resume itself (harper#2448), against the + * generation and the resume floor; until then a lost race degrades to the truncation that happens + * today, never to anything worse. * * **On RocksDB the floor tracks the configured horizon, not retained reality.** That branch purges at * whole-log-file granularity, so it cannot know before the fact which entries a purge will drop, and @@ -828,18 +889,179 @@ export function establishAuditFloor(auditStore: any): void { * resync. Conservative in the one safe direction, and the reason the LMDB branch (which can see a * single eligible entry) instead raises off the first one it finds. * - * **Not covered: copying a database's state without its history.** `restore_backup` replaces a - * database with the backup's, floor and all, and a RocksDB checkpoint (a branch database) copies the - * floor record but no transaction logs — so in both cases a cursor from after the copy point sits - * above a floor that is present, and therefore trusted, for history that is not there. Nothing here - * can detect that on its own, and it is not the audit log's problem alone: the same copy rolls back - * record versions and per-node replication sequence state, so making this one field honest while - * those stay stale would not give a consumer a coherent answer. It needs a database-level epoch — - * harper#2451. + * **Copies of a database's state.** `restore_backup` copies this record with everything else, and it + * stays accurate for the log the restore carried; a branch checkpoint copies it with no log at all. + * What a copy changes is which history the database is, and that is the database generation's to + * answer: every copy path stamps a new one (`stampDatabaseGeneration`), and a copy that carried no log + * raises this floor to the generation's epoch. */ export function getAuditFloor(auditStore: any): number { return decodeAuditFloor(auditStore.getBinary(AUDIT_FLOOR_KEY)); } + +export interface DatabaseGeneration { + /** Thirty-two lowercase hex characters. */ + id: string; + /** When the generation began, in the audit-log key domain; 0 for a genesis generation. */ + epoch: number; +} + +/** + * The generation this handle's open established, or undefined when it could not establish one, in + * which case nothing may resume against it. Per database per node, never cluster-wide. + */ +export function getDatabaseGeneration(auditStore: any): DatabaseGeneration | undefined { + return auditStore?.databaseGeneration; +} + +/** + * The bound a resumed position is checked against, or `Infinity` when unknown: the resume floor, but + * never below a finite audit floor, since a binary that predates the resume floor raised only the + * audit floor when it pruned. Such a binary pruning while the audit floor was unknown records nothing + * either record can show. + */ +export function getAuditResumeFloor(auditStore: any): number { + // The audit floor first: both records only rise, and an unknown audit floor stays unknown, so a prune + // committing between the two reads is still seen through the resume floor read second. + const auditFloor = getAuditFloor(auditStore); + const resumeFloor = decodeAuditFloor(auditStore.getBinary(AUDIT_RESUME_FLOOR_KEY)); + return Number.isFinite(auditFloor) && auditFloor > resumeFloor ? auditFloor : resumeFloor; +} + +/** + * Whether a position may resume here: it names this handle's generation and carries a finite cursor + * at or above the resume floor. The cursor must be progress-based — no lower than the log position + * observed when the position was established — or a quiet scope's snapshot key falls below a floor + * that retention keeps advancing and is refused on every resume. The answer holds as of the read: a + * prune that commits after it returns is the caller's to order against its replay. Always false on + * LMDB, which has no generation: a caller applies the check to RocksDB databases only. + */ +export function isResumablePosition(auditStore: any, generationId: string | undefined, cursor: number): boolean { + const generation = getDatabaseGeneration(auditStore); + return ( + generation !== undefined && + generationId === generation.id && + typeof cursor === 'number' && + Number.isFinite(cursor) && + cursor >= getAuditResumeFloor(auditStore) + ); +} + +function newGeneration(epoch: number): DatabaseGeneration { + return { id: randomBytes(GENERATION_ID_BYTES).toString('hex'), epoch }; +} + +function encodeGeneration(generation: DatabaseGeneration): Uint8Array { + const bytes = new Uint8Array(GENERATION_RECORD_BYTES); + bytes.set(Buffer.from(generation.id, 'hex')); + new DataView(bytes.buffer).setFloat64(GENERATION_ID_BYTES, generation.epoch, true); + return bytes; +} + +function decodeGeneration(stored: any): DatabaseGeneration | undefined { + if (stored?.byteLength !== GENERATION_RECORD_BYTES) return undefined; + const epoch = new DataView(stored.buffer, stored.byteOffset, GENERATION_RECORD_BYTES).getFloat64( + GENERATION_ID_BYTES, + true + ); + if (!Number.isFinite(epoch) || epoch < 0 || Object.is(epoch, -0)) return undefined; + return { id: Buffer.from(stored.buffer, stored.byteOffset, GENERATION_ID_BYTES).toString('hex'), epoch }; +} + +/** + * Establish this handle's generation at open: adopt the recorded one, or mint genesis for a RocksDB + * store that has none (an LMDB store gets none). An undecodable record is never replaced here — an + * ordinary open cannot know whether another worker already serves it; only a copy's stamp replaces a + * generation. Never fails the open. + */ +export function establishDatabaseGeneration(auditStore: any): void { + auditStore.databaseGeneration = undefined; + if (!isRocksStore(auditStore)) return; + try { + let stored = auditStore.getBinary(DATABASE_GENERATION_KEY); + if (stored === undefined) { + if (isReadOnlyMode()) return; + const genesis = encodeGeneration(newGeneration(0)); + commitAuditMetadata( + auditStore, + (read) => { + // compare-and-set on absence, so racing workers converge on the first one's id + if (read(DATABASE_GENERATION_KEY) !== undefined) return undefined; + const writes: Array<[symbol, Uint8Array]> = [[DATABASE_GENERATION_KEY, genesis]]; + if (read(AUDIT_RESUME_FLOOR_KEY) === undefined) { + // starting at the audit floor rather than 0 keeps the next prune's lock-free skip effective + const auditFloor = decodeAuditFloor(read(AUDIT_FLOOR_KEY)); + writes.push([AUDIT_RESUME_FLOOR_KEY, encodeAuditFloor(Number.isFinite(auditFloor) ? auditFloor : 0)]); + } + return writes; + }, + 'database generation' + ); + stored = auditStore.getBinary(DATABASE_GENERATION_KEY); + } + auditStore.databaseGeneration = decodeGeneration(stored); + if (!auditStore.databaseGeneration) + warnContained( + 'The database generation record is unreadable, so no subscription can resume against this database', + new Error(`database generation record of ${stored?.byteLength} bytes`) + ); + } catch (error) { + warnContained('Error establishing the database generation', error); + } +} + +/** + * Give a copy of a database a generation of its own, before anything can read it; the caller holds + * the copy exclusively and makes the stamp durable before publishing. A copy that kept no + * transaction log has nothing below its epoch, so a finite floor is raised to it — raising it over + * history a copy carried would change `Table.commit`'s reconciliation, and an unknown floor stays + * unknown. `generation` lets a caller replay a value it recorded durably first. + */ +export function stampDatabaseGeneration( + store: any, + { carriesLog, generation = newGeneration(Date.now()) }: { carriesLog: boolean; generation?: DatabaseGeneration } +): DatabaseGeneration { + if ( + !/^[0-9a-f]{32}$/.test(generation?.id) || + decodeGeneration(encodeGeneration(generation))?.epoch !== generation.epoch + ) + throw new Error(`Invalid database generation: ${JSON.stringify(generation)}`); + const stamp: Array<[symbol, Uint8Array]> = [ + [DATABASE_GENERATION_KEY, encodeGeneration(generation)], + [AUDIT_RESUME_FLOOR_KEY, encodeAuditFloor(0)], + ]; + commitAuditMetadata( + store, + (read) => { + if (carriesLog) return stamp; + const stored = read(AUDIT_FLOOR_KEY); + const floor = decodeAuditFloor(stored); + return stored !== undefined && Number.isFinite(floor) && floor < generation.epoch + ? [...stamp, [AUDIT_FLOOR_KEY, encodeAuditFloor(generation.epoch)]] + : stamp; + }, + 'database generation' + ); + return generation; +} + +/** + * Stamp the RocksDB database at `path`, which nothing else may have open, and flush it: a root-store + * write alone is not power-loss durable, and directory fsync is best-effort where unsupported. + */ +export async function stampDatabaseDirectory( + path: string, + options: { carriesLog: boolean; generation?: DatabaseGeneration } +): Promise { + const database = openRocksDatabase(path, { disableWAL: false }); + try { + const generation = stampDatabaseGeneration(database, options); + await database.flush({ allowWriteStall: true }); + return generation; + } finally { + database.close(); + } +} export function setAuditRetention(retentionTime, defaultDelay = DEFAULT_AUDIT_CLEANUP_DELAY) { auditRetention = retentionTime; DEFAULT_AUDIT_CLEANUP_DELAY = defaultDelay; diff --git a/resources/branchDatabase.ts b/resources/branchDatabase.ts index 426d0eb686..45351c5c68 100644 --- a/resources/branchDatabase.ts +++ b/resources/branchDatabase.ts @@ -26,6 +26,7 @@ import { retakeBranchIdentity, } from './databases.ts'; import { replayLogs, replayTimeBudgetMs } from './replayLogs.ts'; +import { stampDatabaseDirectory } from './auditStore.ts'; /** * Private per-application forks of a database, for running several variants of an application against @@ -345,6 +346,7 @@ async function materializeBranch( await base.createCheckpoint(staging); report.progress(); await cloneBlobRoots(baseName, baseRoots, blobRoots, report.progress); + await stampDatabaseDirectory(staging, { carriesLog: false }); await writeFile(join(staging, COMPLETION_MARKER), JSON.stringify({ blobRoots } satisfies BranchCompletion)); await rename(staging, branchPath); return blobRoots; diff --git a/resources/transactionBroadcast.ts b/resources/transactionBroadcast.ts index 365a396ecf..9bc0237383 100644 --- a/resources/transactionBroadcast.ts +++ b/resources/transactionBroadcast.ts @@ -1,10 +1,33 @@ +import { basename } from 'node:path'; import { warn } from '../utility/logging/harper_logger.js'; +import { DatabaseClosingError, DatabaseGenerationChangedError } from '../utility/errors/hdbError.ts'; import { IterableEventQueue } from './IterableEventQueue.ts'; import { keyArrayToString } from './Resources.ts'; import type { Id } from './ResourceInterface.ts'; const allSubscriptions = Object.create(null); // using it as a map that doesn't change much const allSameThreadSubscriptions = Object.create(null); // using it as a map that doesn't change much +// The handle each database path was last opened with on this thread, kept as a token so a retired database's closed +// store graph is not retained. A registration through any other handle would join the entry the current store's +// commits drive while reading through that handle's closed stores. +const HANDLE_TOKEN = Symbol('subscription-handle'); +const currentHandles = new Map< + string, + { token: symbol; generationId: string | undefined; tracksGeneration: boolean } +>(); + +// A store with no generation (LMDB) is never replaced by a copy, so its reopen is only ever a close. +function generationChanged( + generationId: string | null | undefined, + currentId: string | null | undefined, + tracksGeneration: boolean +): boolean { + return tracksGeneration && !(generationId != null && generationId === currentId); +} + +function closingError(auditStore: any, path: string): DatabaseClosingError { + return new DatabaseClosingError(auditStore?.rootStore?.databaseName ?? basename(path)); +} /** * This module/function is responsible for the main work of tracking subscriptions and listening for new transactions * that have occurred on any thread, and then reading through the transaction log to notify listeners. This is @@ -24,6 +47,15 @@ export function addSubscription(table, key, listener?: (key) => any, startTime?: if (!path) { throw new Error('No path for table primary store'); } + const generationId = table.auditStore?.databaseGeneration?.id ?? null; + const current = currentHandles.get(path); + const replaced = current !== undefined && table.auditStore?.[HANDLE_TOKEN] !== current.token; + if (replaced || table.auditStore?.rootStore?.status === 'closed') { + if (options?.scope === 'full-database') return; + throw replaced && generationChanged(generationId, current.generationId, current.tracksGeneration) + ? new DatabaseGenerationChangedError() + : closingError(table.auditStore, path); + } if (options?.crossThreads === false) { // we are only listening for commits on our own thread, so we use a separate subscriber and sequencer tracker databaseSubscriptions = allSameThreadSubscriptions[path] || (allSameThreadSubscriptions[path] = []); @@ -60,6 +92,7 @@ export function addSubscription(table, key, listener?: (key) => any, startTime?: } } databaseSubscriptions.auditStore = table.auditStore; + databaseSubscriptions.generationId ??= generationId; if (databaseSubscriptions.lastTxnTime == null) { databaseSubscriptions.lastTxnTime = Date.now(); } @@ -90,6 +123,49 @@ export function addSubscription(table, key, listener?: (key) => any, startTime?: return subscription; } +/** + * End every subscription this thread registered on the database before `auditStore` reopened it: its commit + * listener and its table's stores belong to the closed handle, so it can never deliver again. On a store that + * tracks generations, another or an unknown one must resynchronize; otherwise the retryable + * `DatabaseClosingError`. `auditStore` becomes the only handle `addSubscription` accepts on the path before + * any listener runs. + */ +export function endSubscriptionsFromEarlierHandles(auditStore: any, tracksGeneration: boolean): void { + const path = auditStore.rootStore.path; + const generationId = auditStore.databaseGeneration?.id; + const token = Symbol(basename(path)); + auditStore[HANDLE_TOKEN] = token; + currentHandles.set(path, { token, generationId, tracksGeneration }); + for (const registry of [allSubscriptions, allSameThreadSubscriptions]) { + const databaseSubscriptions = registry[path]; + if (!databaseSubscriptions) continue; + delete registry[path]; + const changed = generationChanged(databaseSubscriptions.generationId, generationId, tracksGeneration); + for (const tableId in databaseSubscriptions) { + const tableSubscriptions = databaseSubscriptions[tableId]; + if (!(tableSubscriptions instanceof Map)) continue; + for (const keySubscriptions of tableSubscriptions.values()) { + for (const subscription of [...keySubscriptions]) { + try { + subscription.close(changed ? new DatabaseGenerationChangedError() : closingError(auditStore, path)); + } catch (error) { + try { + warn(error); + } catch {} + } finally { + // a listener that threw on the final event left the queue open; a bare close sends nothing + if (!subscription.closed) { + try { + subscription.close(); + } catch {} + } + } + } + } + } + } +} + /** * This is the class that is returned from subscribe calls and provide the interface to set a callback, end the * subscription and get the initial state. diff --git a/unitTests/bin/migrationStagingRecovery.test.js b/unitTests/bin/migrationStagingRecovery.test.js index 600a29d43c..8b78271e31 100644 --- a/unitTests/bin/migrationStagingRecovery.test.js +++ b/unitTests/bin/migrationStagingRecovery.test.js @@ -92,6 +92,9 @@ describe('migration: staging directory recovery (#2012)', function () { assert(entry.value instanceof RecordObject, `record ${id} lost its record prototype`); assert(entry.version > 0, `record ${id} lost its version`); } + const generation = root.getBinarySync(Symbol.for('database-generation')); + assert.strictEqual(generation?.byteLength, 24, 'the migrated store must carry a generation'); + assert(Buffer.from(generation).readDoubleLE(16) > 0, 'a migrated generation is a copy, not genesis'); } finally { cf.close(); root.close(); diff --git a/unitTests/dataLayer/restoreGeneration.test.js b/unitTests/dataLayer/restoreGeneration.test.js new file mode 100644 index 0000000000..a26032f233 --- /dev/null +++ b/unitTests/dataLayer/restoreGeneration.test.js @@ -0,0 +1,106 @@ +const assert = require('node:assert'); +const { rmSync } = require('node:fs'); +const { setupTestDBPath } = require('../testUtils'); +const { table, closeDatabase } = require('#src/resources/databases'); +const { createBackupOffline, restoreBackupOffline, backupDirForDatabase } = require('#src/dataLayer/rocksdbBackup'); +const { getAuditFloor, getDatabaseGeneration, isResumablePosition } = require('#src/resources/auditStore'); +const { DatabaseGenerationChangedError } = require('#src/utility/errors/hdbError'); +const { setAnalyticsEnabled } = require('#src/resources/analytics/write'); +const { setMainIsWorker } = require('#js/server/threads/manageThreads'); +const { waitFor } = require('../waitFor'); +require('#src/server/serverHelpers/serverUtilities'); + +describe('A restore starts a new database generation', function () { + if (process.env.HARPER_STORAGE_ENGINE === 'lmdb') return; // restore_backup is RocksDB-only + let sequence = 0; + const databases = []; + before(() => { + setupTestDBPath(); + setMainIsWorker(true); + // An analytics flush declares its table, and the schema rescan that follows opens every database + // under STORAGE_PATH and keeps it open: one this suite is about to restore, or, once the flush + // outlives the suite, the next suite's. + setAnalyticsEnabled(false); + }); + after(() => { + setAnalyticsEnabled(true); + for (const database of databases) rmSync(backupDirForDatabase(database), { recursive: true, force: true }); + }); + + async function backedUpDatabase() { + const database = `restore_generation_${++sequence}`; + databases.push(database); + const open = () => + table({ + table: 'Restored', + database, + audit: true, + attributes: [{ name: 'id', isPrimaryKey: true }, { name: 'value' }], + }); + // Backed up while open, so the database is closed exactly once before the restore: a second + // close after a reopen in the same process leaves the root handle registered. Flushed first, + // because `table()`'s on-demand reopen does not replay the transaction log the way a server + // reload does. + const T = open(); + await T.put('A', { value: 1 }); + await T.auditStore.rootStore.flush({ allowWriteStall: true }); + const { backup_id: backupId } = await createBackupOffline(database); + await T.put('A', { value: 2 }); + return { database, open, T, backupId }; + } + + it('ends a live subscription from before the restore instead of resuming it against the restored data', async () => { + const { database, open, T, backupId } = await backedUpDatabase(); + const events = []; + const subscription = await T.subscribe({ id: 'A', listener: (event) => events.push(event) }); + assert.ok(await closeDatabase(database)); + await restoreBackupOffline(database, backupId); + const restored = open(); + assert.strictEqual((await restored.get('A')).value, 1, 'precondition: the restore took effect'); + const current = []; + await restored.subscribe({ id: 'A', listener: (event) => current.push(event) }); + await restored.put('A', { value: 3 }); + await waitFor(() => current.some((event) => event.value?.value === 3)); + assert.ok( + !events.some((event) => event.value?.value === 3), + 'the old subscriber must not silently receive the restored database’s writes' + ); + assert.strictEqual(subscription.closed, true); + }); + + it('refuses a position minted after the backup point, and accepts one minted after the restore', async () => { + const { database, open, T, backupId } = await backedUpDatabase(); + const before = getDatabaseGeneration(T.auditStore); + const cursor = Date.now(); + assert.strictEqual(isResumablePosition(T.auditStore, before.id, cursor), true, 'precondition'); + const floor = getAuditFloor(T.auditStore); + const events = []; + await T.subscribe({ id: 'A', listener: (event) => events.push(event) }); + assert.ok(await closeDatabase(database)); + await restoreBackupOffline(database, backupId); + const restored = open(); + const after = getDatabaseGeneration(restored.auditStore); + assert.notStrictEqual(after.id, before.id); + assert.ok(after.epoch >= cursor, 'the generation began at the restore'); + assert.strictEqual(isResumablePosition(restored.auditStore, before.id, cursor), false); + assert.strictEqual(isResumablePosition(restored.auditStore, after.id, Date.now()), true); + assert.strictEqual(getAuditFloor(restored.auditStore), floor, 'the restore carried its log, so its floor stands'); + assert.ok(events.at(-1) instanceof DatabaseGenerationChangedError, 'the old subscriber is told to resync'); + }); + + it('gives a restore into a new database a generation of its own', async () => { + const { database, T, backupId } = await backedUpDatabase(); + const source = getDatabaseGeneration(T.auditStore); + const target = `${database}_target`; + databases.push(target); + assert.ok(await closeDatabase(database)); + await restoreBackupOffline(database, backupId, target); + const copy = table({ + table: 'Restored', + database: target, + audit: true, + attributes: [{ name: 'id', isPrimaryKey: true }, { name: 'value' }], + }); + assert.notStrictEqual(getDatabaseGeneration(copy.auditStore).id, source.id); + }); +}); diff --git a/unitTests/resources/auditPurge.test.js b/unitTests/resources/auditPurge.test.js index e0f82593f7..66ad20b34d 100644 --- a/unitTests/resources/auditPurge.test.js +++ b/unitTests/resources/auditPurge.test.js @@ -36,24 +36,28 @@ describe('purgeAgedLogs', () => { return purgedFiles; }, }; - // A REAL RocksTransactionLogStore, because `updateAuditFloor` branches on `instanceof` and + // A REAL RocksTransactionLogStore, because the metadata write branches on `instanceof` and // purgeAgedLogs is typed `rootStore: RocksDatabase` — a plain object would silently route this // test through the LMDB branch and leave the one production path that always uses the RocksDB // branch untested. const auditStandIn = new RocksTransactionLogStore({ useLog: () => ({}) }); - // Holds what it is given, because updateAuditFloor reads its own write back inside the - // transaction to catch a put that failed without saying so. Starts at the baseline a database - // that has provably pruned nothing carries; an absent floor would instead mean "unknown", which - // a prune deliberately leaves alone. - let stored = new Uint8Array(Float64Array.of(1).buffer); + // Holds what it is given, per record, because the write reads itself back inside the transaction + // to catch a put that failed without saying so. Starts at the baseline a database that has + // provably pruned nothing carries: a floor, and the genesis resume floor. An absent floor would + // instead mean "unknown", which a prune deliberately leaves alone. + const FLOOR_KEY = Symbol.for('audit-floor'); + const records = new Map([ + [FLOOR_KEY, new Uint8Array(Float64Array.of(1).buffer)], + [Symbol.for('audit-resume-floor'), new Uint8Array(Float64Array.of(0).buffer)], + ]); const txn = { - getBinarySync: () => stored, - putSync(_key, value) { + getBinarySync: (key) => records.get(key), + putSync(key, value) { // the real write wraps the bytes in lmdb's asBinary(), which bypasses both engines' encoders const wrapped = Object.values(value)[0] ?? value; const bytes = Uint8Array.from(Object.values(wrapped)); - stored = failFloorWrite ? stored : bytes; - store.floorWrites.push(new Float64Array(bytes.slice().buffer)[0]); + if (!failFloorWrite) records.set(key, bytes); + if (key === FLOOR_KEY) store.floorWrites.push(new Float64Array(bytes.slice().buffer)[0]); }, }; let failFloorWrite = false; @@ -63,13 +67,13 @@ describe('purgeAgedLogs', () => { // A floorless store is the case where the resolver writes the unknown sentinel, and where a // read-back comparing decoded values alone cannot tell that write from no record at all. store.startFloorless = () => { - stored = undefined; + records.delete(FLOOR_KEY); }; auditStandIn.rootStore = { // RocksTransactionLogStore.getBinary delegates here, which is how raiseAuditFloor's lock-free // pre-check reads the floor - getBinarySync: () => stored, - // returns the callback's value, as both real engines do: updateAuditFloor requires an + getBinarySync: (key) => records.get(key), + // returns the callback's value, as both real engines do: the metadata write requires an // explicit `true` because RocksDB swallows an aborted transaction and returns undefined. transactionSync(callback) { order.push('floor'); diff --git a/unitTests/resources/branchDatabase.test.js b/unitTests/resources/branchDatabase.test.js index ab6e7adf8b..0a5d7b2c5a 100644 --- a/unitTests/resources/branchDatabase.test.js +++ b/unitTests/resources/branchDatabase.test.js @@ -6,7 +6,8 @@ const { join } = require('node:path'); const { RocksDatabase } = require('@harperfast/rocksdb-js'); const { setupTestDBPath } = require('../testUtils'); const { table, databases, database, BRANCH_ROOT_DIR, resolveBranchPath } = require('#src/resources/databases'); -const { getOrCreateBranch, removeBranches } = require('#src/resources/branchDatabase'); +const { getOrCreateBranch, removeBranches, closeBranchAt } = require('#src/resources/branchDatabase'); +const { getAuditFloor, getDatabaseGeneration } = require('#src/resources/auditStore'); const { replayLogs } = require('#src/resources/replayLogs'); const { setMainIsWorker } = require('#js/server/threads/manageThreads'); @@ -58,6 +59,20 @@ describeUnlessLmdb('branch lifecycle (harper#643)', () => { assert.ok(expected.includes(BRANCH_ROOT_DIR), 'branches belong under the reserved root'); }); + it('gives a new branch a generation of its own, and keeps it when the branch is adopted again', async function () { + const branch = await getOrCreateBranch('lifebase', 'appGeneration'); + const forked = getDatabaseGeneration(branch.tables.LifecycleSource.auditStore); + assert.notStrictEqual(forked.id, getDatabaseGeneration(Source.auditStore).id); + assert.ok(forked.epoch > 0, 'a branch is a copy, not genesis'); + assert.ok( + getAuditFloor(branch.tables.LifecycleSource.auditStore) >= forked.epoch, + 'the checkpoint carried no log below the fork' + ); + await closeBranchAt(resolveBranchPath('lifebase', 'appGeneration')); + const adopted = await getOrCreateBranch('lifebase', 'appGeneration'); + assert.deepStrictEqual(getDatabaseGeneration(adopted.tables.LifecycleSource.auditStore), forked); + }); + it('gives concurrent callers the same branch rather than racing two checkpoints', async function () { // Every worker thread loads the same applications, so this is the ordinary case, not an edge one. const [first, second, third] = await Promise.all([ diff --git a/unitTests/resources/databaseGeneration.test.js b/unitTests/resources/databaseGeneration.test.js new file mode 100644 index 0000000000..c9d9b4326a --- /dev/null +++ b/unitTests/resources/databaseGeneration.test.js @@ -0,0 +1,415 @@ +const assert = require('node:assert'); +const { setupTestDBPath } = require('../testUtils'); +const { table, closeDatabase, __setReadOnlyModeForTest } = require('#src/resources/databases'); +const { + getAuditFloor, + getAuditResumeFloor, + getDatabaseGeneration, + establishDatabaseGeneration, + isResumablePosition, + raiseAuditFloor, + stampDatabaseGeneration, + stampDatabaseDirectory, +} = require('#src/resources/auditStore'); +const { DatabaseClosingError, DatabaseGenerationChangedError } = require('#src/utility/errors/hdbError'); +const { setMainIsWorker } = require('#js/server/threads/manageThreads'); +const { waitFor } = require('../waitFor'); +require('#src/server/serverHelpers/serverUtilities'); + +const GENERATION_KEY = Symbol.for('database-generation'); +const RESUME_FLOOR_KEY = Symbol.for('audit-resume-floor'); +const FLOOR_KEY = Symbol.for('audit-floor'); +const isRocksDB = process.env.HARPER_STORAGE_ENGINE !== 'lmdb'; + +function floorBytes(value) { + return new Uint8Array(new Float64Array([value]).buffer); +} + +/** Write raw metadata the way an older binary, or corruption, would leave it. */ +function putRecord(auditStore, key, bytes) { + auditStore.putSync(key, bytes); +} + +async function clearRecord(auditStore, key) { + const root = auditStore.rootStore; + if (typeof root?.removeSync === 'function') root.removeSync(key); + await auditStore.remove(key); +} + +let sequence = 0; +function tableInOwnDatabase(name = `Generation${++sequence}`) { + return table({ + table: name, + database: `generation_${name}`, + audit: true, + attributes: [{ name: 'id', isPrimaryKey: true }, { name: 'value' }], + }); +} + +describe('Database generation', () => { + if (!isRocksDB) return; // RocksDB only; the LMDB behavior is below + before(() => { + setupTestDBPath(); + setMainIsWorker(true); + }); + + describe('genesis', () => { + it('gives a new database a generation of its own, with a finite resume floor', () => { + const store = tableInOwnDatabase().auditStore; + const generation = getDatabaseGeneration(store); + assert.match(generation.id, /^[0-9a-f]{32}$/); + assert.strictEqual(generation.epoch, 0, 'no copy produced it'); + assert.ok(Number.isFinite(getAuditResumeFloor(store))); + }); + + it('gives two databases different generations', () => { + assert.notStrictEqual( + getDatabaseGeneration(tableInOwnDatabase().auditStore).id, + getDatabaseGeneration(tableInOwnDatabase().auditStore).id + ); + }); + + it('keeps the generation across a close and reopen', async () => { + const name = 'Reopened'; + const before = getDatabaseGeneration(tableInOwnDatabase(name).auditStore); + assert.ok(await closeDatabase(`generation_${name}`)); + const reopened = tableInOwnDatabase(name).auditStore; + assert.deepStrictEqual(getDatabaseGeneration(reopened), before); + }); + + it("adopts another worker's genesis instead of minting a second one", async () => { + const store = tableInOwnDatabase().auditStore; + const winner = getDatabaseGeneration(store); + // This worker read "no generation" before the winner's commit landed; the compare-and-set + // inside the transaction sees the record and writes nothing. + const realGetBinary = store.getBinary.bind(store); + let hidden = false; + store.getBinary = (key) => { + if (key === GENERATION_KEY && !hidden) { + hidden = true; + return undefined; + } + return realGetBinary(key); + }; + try { + establishDatabaseGeneration(store); + } finally { + store.getBinary = realGetBinary; + } + assert.ok(hidden, 'precondition: the outer read was the stale one'); + assert.deepStrictEqual(getDatabaseGeneration(store), winner); + }); + + it('refuses to resume against an unreadable generation record, and does not repair it', async () => { + const store = tableInOwnDatabase().auditStore; + const garbage = new Uint8Array([1, 2, 3]); + putRecord(store, GENERATION_KEY, garbage); + establishDatabaseGeneration(store); + assert.strictEqual(getDatabaseGeneration(store), undefined); + assert.deepStrictEqual(new Uint8Array(store.getBinary(GENERATION_KEY)), garbage, 'an open must not rewrite it'); + assert.strictEqual(isResumablePosition(store, undefined, Date.now()), false); + await clearRecord(store, GENERATION_KEY); + }); + + it('mints nothing in read-only mode', async () => { + const store = tableInOwnDatabase().auditStore; + await clearRecord(store, GENERATION_KEY); + __setReadOnlyModeForTest(true); + try { + establishDatabaseGeneration(store); + } finally { + __setReadOnlyModeForTest(undefined); + } + assert.strictEqual(getDatabaseGeneration(store), undefined); + assert.strictEqual(store.getBinary(GENERATION_KEY), undefined); + }); + }); + + describe('the resume floor', () => { + it('sees a prune that commits between its two metadata reads', () => { + const store = tableInOwnDatabase().auditStore; + putRecord(store, FLOOR_KEY, floorBytes(Infinity)); + putRecord(store, RESUME_FLOOR_KEY, floorBytes(100)); + const { id } = getDatabaseGeneration(store); + const getBinary = store.getBinary; + let pruned = false; + store.getBinary = function (key) { + const value = getBinary.call(this, key); + if (!pruned) { + pruned = true; + raiseAuditFloor(store, 200); + } + return value; + }; + try { + assert.strictEqual(isResumablePosition(store, id, 150), false); + } finally { + store.getBinary = getBinary; + } + }); + + it('is raised by a prune together with the audit floor', () => { + const store = tableInOwnDatabase().auditStore; + const cutoff = Date.now() + 1000; + raiseAuditFloor(store, cutoff); + assert.strictEqual(getAuditFloor(store), cutoff); + assert.strictEqual(getAuditResumeFloor(store), cutoff); + }); + + it('is not absorbed by an unknown audit floor, which stays unknown', () => { + const store = tableInOwnDatabase().auditStore; + putRecord(store, FLOOR_KEY, floorBytes(Infinity)); + const cutoff = Date.now() + 1000; + raiseAuditFloor(store, cutoff); + assert.strictEqual(getAuditFloor(store), Infinity, "reconciliation's unknown floor is never rewritten"); + assert.strictEqual(getAuditResumeFloor(store), cutoff); + const { id } = getDatabaseGeneration(store); + assert.strictEqual(isResumablePosition(store, id, cutoff), true, 'an unknown floor no longer blocks resume'); + assert.strictEqual(isResumablePosition(store, id, cutoff - 1), false); + }); + + it('is raised even when the audit floor already covers the cutoff', async () => { + const store = tableInOwnDatabase().auditStore; + const cutoff = Date.now() + 1000; + raiseAuditFloor(store, cutoff + 1000); + await clearRecord(store, RESUME_FLOOR_KEY); + raiseAuditFloor(store, cutoff); + putRecord(store, FLOOR_KEY, floorBytes(Infinity)); + assert.strictEqual(getAuditResumeFloor(store), cutoff, 'the lock-free skip must check both records'); + }); + + it('is never below a finite audit floor, which a generation-unaware binary raises alone when it prunes', () => { + const store = tableInOwnDatabase().auditStore; + const { id } = getDatabaseGeneration(store); + const cursor = Date.now(); + assert.strictEqual(isResumablePosition(store, id, cursor), true, 'precondition'); + putRecord(store, FLOOR_KEY, floorBytes(cursor + 5000)); + assert.strictEqual(getAuditResumeFloor(store), cursor + 5000); + assert.strictEqual(isResumablePosition(store, id, cursor), false); + __setReadOnlyModeForTest(true); + try { + assert.strictEqual(isResumablePosition(store, id, cursor), false, 'read-only opens read the same bound'); + } finally { + __setReadOnlyModeForTest(undefined); + } + }); + }); + + describe('stamping a copy', () => { + it('starts a new generation, and leaves the floor of a copy that carried its log alone', () => { + const store = tableInOwnDatabase().auditStore; + const before = getDatabaseGeneration(store); + raiseAuditFloor(store, Date.now() - 60_000); + const floor = getAuditFloor(store); + const stamped = stampDatabaseGeneration(store, { carriesLog: true }); + assert.notStrictEqual(stamped.id, before.id); + assert.ok(stamped.epoch > 0); + assert.strictEqual(getAuditFloor(store), floor); + assert.strictEqual(getAuditResumeFloor(store), floor, 'no prune has run in the new generation'); + establishDatabaseGeneration(store); + assert.deepStrictEqual(getDatabaseGeneration(store), stamped); + assert.strictEqual(isResumablePosition(store, before.id, Date.now()), false, 'the old generation is refused'); + }); + + it('raises the floor of a copy that carried no log to its epoch, but never an unknown one', () => { + const store = tableInOwnDatabase().auditStore; + raiseAuditFloor(store, Date.now() - 60_000); + const stamped = stampDatabaseGeneration(store, { carriesLog: false }); + assert.strictEqual(getAuditFloor(store), stamped.epoch); + + const unknown = tableInOwnDatabase().auditStore; + putRecord(unknown, FLOOR_KEY, floorBytes(Infinity)); + stampDatabaseGeneration(unknown, { carriesLog: false }); + assert.strictEqual(getAuditFloor(unknown), Infinity); + }); + + it('replays a generation the caller recorded first', () => { + const store = tableInOwnDatabase().auditStore; + const recorded = { id: 'ab'.repeat(16), epoch: Date.now() }; + assert.deepStrictEqual(stampDatabaseGeneration(store, { carriesLog: true, generation: recorded }), recorded); + establishDatabaseGeneration(store); + assert.deepStrictEqual(getDatabaseGeneration(store), recorded); + }); + }); + + describe('resumable positions', () => { + it('require the current generation id and a finite cursor at or above the resume floor', () => { + const store = tableInOwnDatabase().auditStore; + const { id } = getDatabaseGeneration(store); + const floor = Date.now(); + raiseAuditFloor(store, floor); + assert.strictEqual(isResumablePosition(store, id, floor), true); + assert.strictEqual(isResumablePosition(store, id, floor - 1), false); + assert.strictEqual(isResumablePosition(store, 'cd'.repeat(16), floor), false); + assert.strictEqual(isResumablePosition(store, undefined, floor), false, 'no scalar mode'); + assert.strictEqual(isResumablePosition(store, id, Infinity), false); + assert.strictEqual(isResumablePosition(store, id, NaN), false); + }); + + it('refuses a cursor from before tracking began, however it compares with the bootstrap floor', () => { + // resources/DESIGN.md's worked example: a legacy prune removed tableA through 1000, the newest + // survivor was 900, and a rolled-back clock stamped the floor at 920. A cursor at 970 compares + // above every floor; it carries no generation id, so it cannot be certified. + const store = tableInOwnDatabase().auditStore; + putRecord(store, FLOOR_KEY, floorBytes(920)); + assert.strictEqual(isResumablePosition(store, undefined, 970), false); + }); + }); + + describe('live subscriptions across a generation change', function () { + async function reopenAsCopy(name, stampedTable) { + const path = stampedTable.auditStore.rootStore.path; + assert.ok(await closeDatabase(`generation_${name}`)); + await stampDatabaseDirectory(path, { carriesLog: true }); + return tableInOwnDatabase(name); + } + + it('ends a subscription registered before the database was replaced by a copy', async () => { + const name = 'LiveReplaced'; + const T = tableInOwnDatabase(name); + const events = []; + const subscription = await T.subscribe({ id: 'A', listener: (event) => events.push(event) }); + await T.put('A', { value: 1 }); + await waitFor(() => events.some((event) => event.value?.value === 1)); + const copy = await reopenAsCopy(name, T); + assert.strictEqual(subscription.closed, true); + assert.ok(events.at(-1) instanceof DatabaseGenerationChangedError, 'the consumer is told to resync'); + const after = []; + await copy.subscribe({ id: 'A', listener: (event) => after.push(event) }); + await copy.put('A', { value: 2 }); + await waitFor(() => after.some((event) => event.value?.value === 2)); + assert.ok(!events.some((event) => event.value?.value === 2), 'the old subscription must not see the copy'); + }); + + it('still closes a subscription whose listener throws on the terminal event', async () => { + const name = 'LiveThrowing'; + const T = tableInOwnDatabase(name); + const subscription = await T.subscribe({ + id: 'A', + listener: (event) => { + if (event instanceof Error) throw new Error('listener failed'); + }, + }); + const sibling = await T.subscribe({ id: 'A' }); + await reopenAsCopy(name, T); + assert.strictEqual(subscription.closed, true); + assert.strictEqual(sibling.closed, true, 'one throwing listener must not strand the others'); + }); + + async function subscribeAndWrite(copy) { + const events = []; + const subscription = await copy.subscribe({ id: 'A', listener: (event) => events.push(event) }); + await copy.put('A', { value: 2 }); + await waitFor(() => events.some((event) => event.value?.value === 2)); + return { subscription, events }; + } + + it('refuses a resubscribe that a listener makes through the replaced database while it is being ended', async () => { + const name = 'LiveReentrant'; + const T = tableInOwnDatabase(name); + await T.put('A', { value: 1 }); + const resource = await T.getResource('A', {}); + const stale = []; + let resubscribe; + await T.subscribe({ + id: 'A', + listener: (event) => { + if (event instanceof DatabaseGenerationChangedError && !resubscribe) { + resubscribe = resource.subscribe({ listener: (staleEvent) => stale.push(staleEvent) }); + resubscribe.catch(() => {}); + } + }, + }); + const copy = await reopenAsCopy(name, T); + assert.ok(resubscribe, 'precondition: the listener resubscribed during the teardown'); + await assert.rejects(resubscribe, DatabaseGenerationChangedError); + const current = await subscribeAndWrite(copy); + assert.deepStrictEqual(stale, []); + assert.ok(await closeDatabase(`generation_${name}`)); + tableInOwnDatabase(name); + assert.ok( + current.events.at(-1) instanceof DatabaseClosingError, + 'the copy’s subscriber ends as a same-generation reopen, not as a replaced database' + ); + }); + + it('refuses a subscription through a resource loaded before the database was replaced', async () => { + const name = 'LiveStaleResource'; + const T = tableInOwnDatabase(name); + await T.put('A', { value: 1 }); + const resource = await T.getResource('A', {}); + const copy = await reopenAsCopy(name, T); + const stale = []; + await assert.rejects( + resource.subscribe({ listener: (event) => stale.push(event) }), + DatabaseGenerationChangedError + ); + (await subscribeAndWrite(copy)).subscription.end(); + assert.deepStrictEqual(stale, []); + }); + + it('ends a subscription at a reopen of the same generation with a retryable error, and a resubscribe resumes', async () => { + const name = 'LiveSame'; + const T = tableInOwnDatabase(name); + const events = []; + const subscription = await T.subscribe({ id: 'A', listener: (event) => events.push(event) }); + assert.ok(await closeDatabase(`generation_${name}`)); + const reopened = tableInOwnDatabase(name); + await reopened.put('A', { value: 1 }); + await waitFor(() => events.some((event) => event instanceof Error || event.value?.value === 1)); + assert.ok(events.at(-1) instanceof DatabaseClosingError, 'the consumer is told to resubscribe'); + assert.strictEqual(subscription.closed, true); + (await subscribeAndWrite(reopened)).subscription.end(); + }); + + it('refuses a retry through the handle that a same-generation reopen closed', async () => { + const name = 'LiveSameRetry'; + const T = tableInOwnDatabase(name); + await T.put('A', { value: 1 }); + const resource = await T.getResource('A', {}); + const first = []; + await resource.subscribe({ listener: (event) => first.push(event) }); + assert.ok(await closeDatabase(`generation_${name}`)); + await assert.rejects(resource.subscribe({}), DatabaseClosingError, 'its store is closed'); + const reopened = tableInOwnDatabase(name); + assert.ok(first.at(-1) instanceof DatabaseClosingError, 'precondition: the reopen ended the first subscription'); + await assert.rejects(resource.subscribe({}), DatabaseClosingError, 'it is not the reopened handle'); + (await subscribeAndWrite(reopened)).subscription.end(); + }); + }); +}); + +describe('Database generation on LMDB', () => { + if (isRocksDB) return; + before(() => { + setupTestDBPath(); + setMainIsWorker(true); + }); + + it('records neither a generation nor a resume floor, so no position resumes', () => { + const store = tableInOwnDatabase().auditStore; + assert.strictEqual(getDatabaseGeneration(store), undefined); + const cutoff = Date.now() + 1000; + raiseAuditFloor(store, cutoff); + assert.strictEqual(getAuditFloor(store), cutoff); + assert.strictEqual(store.getBinary(GENERATION_KEY), undefined); + assert.strictEqual(store.getBinary(RESUME_FLOOR_KEY), undefined); + assert.strictEqual(isResumablePosition(store, undefined, cutoff + 1), false); + }); + + it('ends a subscription at a reopen with the retryable DatabaseClosingError, never as a replaced database', async () => { + const name = 'LmdbReopen'; + const T = tableInOwnDatabase(name); + const events = []; + const subscription = await T.subscribe({ id: 'A', listener: (event) => events.push(event) }); + assert.ok(await closeDatabase(`generation_${name}`)); + const reopened = tableInOwnDatabase(name); + assert.ok(events.at(-1) instanceof DatabaseClosingError); + assert.strictEqual(subscription.closed, true); + const current = []; + const resubscribed = await reopened.subscribe({ id: 'A', listener: (event) => current.push(event) }); + await reopened.put('A', { value: 1 }); + await waitFor(() => current.some((event) => event.value?.value === 1)); + resubscribed.end(); + }); +}); diff --git a/utility/errors/hdbError.ts b/utility/errors/hdbError.ts index 3e72f65b48..b6b0bfe498 100644 --- a/utility/errors/hdbError.ts +++ b/utility/errors/hdbError.ts @@ -117,6 +117,15 @@ export class DatabaseClosingError extends ServerError { } } +export class DatabaseGenerationChangedError extends ClientError { + code: string; + constructor() { + super('The database was replaced by a restored or copied state; resubscribe to resynchronize', 409); + this.name = 'DatabaseGenerationChangedError'; + this.code = 'DATABASE_GENERATION_CHANGED'; + } +} + export class DatabaseDrainTimeoutError extends DatabaseClosingError { constructor(databaseName: string, timeoutMilliseconds: number) { super(databaseName);