Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions resources/DESIGN.md
Original file line number Diff line number Diff line change
Expand Up @@ -480,6 +480,7 @@ Opt-in deflate compression for file-backed blobs (harper#2443) has three load-be
- **Write-bearing request transactions are aborted and poisoned** (issue #1407). The monitor calls `abortDueToTimeout()`, which sets `timedOut`, forces `open = CLOSED` (so `doneReadTxn` takes the discard path instead of re-entering `commit()` via the `LINGERING` branch — which would now throw), then `abort()`s. `addWrite`/`commit` both guard on `timedOut` and throw `transactionOpenTooLongError` (503), so the in-flight request rolls back cleanly rather than the monitor silently force-committing a partial write set (atomicity violation + orphaned secondary-index entries that only a full rebuild repairs). The old behavior `commit()`d and reused the still-open transaction.
- **`hasPendingWrites()` walks the `next` chain.** Writes to a second database live on `transaction.next` (see `txnForContext`), so a transaction that reads database A (head, tracked via its read snapshot, empty `writes`) and writes database B (`next`) is still write-bearing. Without the walk the head looks read-only and the monitor's force-commit path would cascade-commit B. `abortDueToTimeout()` poisons + aborts the whole chain.
- **Read-only, `sourceApply`, and `isReplay` transactions keep the prior force-commit behavior.** Read-only long transactions (large scans/exports) have no atomicity/index risk and must not have their ongoing reads poisoned. Canonical-source applies (replication peer / external caching source) and crash-recovery replay have no resubscribe/resume path: aborting a write would drop it while the resume cursor advances past it — a permanent divergence (harper-pro#348). `sourceApply` is propagated down the `next` chain in `txnForContext`, so gating on the head suffices. (Replay is additionally synchronous, so the async monitor can't fire mid-replay anyway.)
- **The owner's commit waits for the monitor's force-commit.** That commit claims every staged write, marks the transaction `CLOSED`, and detaches its handle, so the owner's later `commit()` finds nothing to do. The monitor therefore stores the promise it got back as `monitorCommit`: `commit()` chains on it and the monitor clears it on success. A failure is kept until the owner's final (`doneWriting`) commit rejects with it and releases the context, since the monitor itself only logs it. That release is deliberately not `abort()`: a multi-store commit can fail after its head store landed, and `abort()`'s blob cleanup would unlink files the head's audit entries still reference. A landed store always finishes its own bookkeeping (writes cleared, record locks released) before a chained store's failure surfaces, including a synchronous throw from `next.commit()`.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion (non-blocking): This bullet reads as a universal invariant, but the mechanism it describes (monitorCommit) is RocksDB-only — LMDBTransaction.commit()/startMonitoringTxns() overrides this path entirely and still force-commits fire-and-forget (only harperLogger.debug?.-logged on failure), so LMDB retains a narrower version of the same window. The PR body already discloses this ("the LMDB monitor has a narrower version of the same window..."), so worth making that scope explicit here too, mirroring the existing "RocksDB-write-path only" callout at DatabaseTransaction.ts:374-376.


- **A transaction parked in its commit phase is spared, not poisoned** (issue #2062). `commit()` sets `committing` around its pre-commit await (the `before`/`beforeIntermediate` completions — in practice a blob's durable file write) and the monitor logs instead of aborting while it is set. The limit polices an _application_ holding a transaction open with an unfinished write set; once `commit()` is entered the write set is sealed and the caller is awaiting the commit, so the time is core's own I/O, and a multi-tens-of-MB deploy payload legitimately outruns the limit. Poisoning there was actively destructive: `abort()` cleared the write set and unlinked the write's pre-saved blobs, and the resumed commit then found nothing to write and resolved as **success** — the caller was told its write landed, and was left holding a blob whose file was gone but whose `fileId` was still set, so its next `put` silently minted a reference to a destroyed file (the deploy-payload case: `Blob file not found` on the peer, unrecoverably). The grace is bounded — `COMMIT_PHASE_GRACE` over-limit ticks, ~10 min at the 30s default, since sparing re-arms `timeout` — because the transaction still pins a read snapshot; a source that stalls rather than finishing falls through to the normal abort. `sourceApply`/`isReplay` are spared without a bound: they may be neither aborted (harper-pro#348) nor force-committed mid-write (that would durably commit a replica record whose blob file is still being written), and their blob sources are bounded by the receive-side idle watchdog instead.
- **Resuming from that await re-checks that the transaction is still alive.** `timedOut` (monitor poison, including via the `next` chain) throws `transactionOpenTooLongError`; a write set cleared with the handle released — a plain `abort()` in the same window — throws `Transaction was aborted while its commit was waiting on pre-commit work`. Without both, either path resolves as a phantom commit. `LMDBTransaction.commit` carries the same pair around its own `before` phase.
Expand Down
50 changes: 40 additions & 10 deletions resources/DatabaseTransaction.ts
Original file line number Diff line number Diff line change
Expand Up @@ -737,6 +737,9 @@ export class DatabaseTransaction implements Transaction {
// open-transaction limit. Once poisoned, any further addWrite/commit throws transactionOpenTooLongError
// so the request rolls back cleanly instead of silently committing a partial write set (issue #1407).
declare timedOut?: boolean;
// The monitor's force-commit, which the owner's commit joins (resources/DESIGN.md). Kept after a failure
// until a final commit reports it: the monitor only logs it.
declare monitorCommit?: Promise<CommitResolution>;
// Set once the retained read handle's write intents have been released (see commit()'s
// outstanding-iterators branch), so a retry round cannot re-fire the release.
declare writesAbandoned?: boolean;
Expand Down Expand Up @@ -1420,6 +1423,19 @@ export class DatabaseTransaction implements Transaction {
*/
commit(options: CommitOptions = {}): MaybePromise<CommitResolution> {
if (this.timedOut) throw transactionOpenTooLongError();
if (this.monitorCommit && !options.transaction) {
return this.monitorCommit.then(
() => this.commit(options),
(error) => {
// The failed commit ran its terminal cleanup, but as a non-final commit it kept the context.
if (options.doneWriting) {
this.monitorCommit = undefined;
this.releaseContext(true);
}
throw error;
}
);
}
// reused across retries — the native layer resets it in place (fresh snapshot) on IsBusy/TryAgain —
// but reassigned to a fresh replay transaction when outstanding read iterators retain this.transaction
let transaction = options.transaction ?? this.transaction;
Expand Down Expand Up @@ -1669,9 +1685,15 @@ export class DatabaseTransaction implements Transaction {
if (this.next) {
// never forward options.transaction (a retry/replay round's HEAD-store handle) to
// the next store — it must commit its own writes through its own transaction
completions.push(
this.next.commit(options.transaction ? { ...options, transaction: undefined } : options)
);
let nextCommit: MaybePromise<CommitResolution>;
try {
nextCommit = this.next.commit(options.transaction ? { ...options, transaction: undefined } : options);
} catch (error) {
// This store has landed; its bookkeeping below must run before the failure surfaces.
nextCommit = Promise.reject(error);
nextCommit.catch(() => {}); // still rejects Promise.all, even if bookkeeping throws first
}
completions.push(nextCommit);
}
if (options?.flush) {
completions.push(this.writes[0].store.flushed);
Expand Down Expand Up @@ -2483,15 +2505,23 @@ function startMonitoringTxns() {
// Read-only long transaction (no atomicity/index risk — e.g. a large scan or export), or a
// canonical-source apply/replay that must never drop a write: preserve the prior behavior of
// committing to close out the snapshot without poisoning the transaction.
let result: MaybePromise<CommitResolution>;
try {
const result = txn.commit();
if ((result as any)?.then) {
(result as any).catch((error) => {
harperLogger.debug?.(`Error committing timed out transaction: ${error.message}`);
});
}
result = txn.commit();
} catch (error) {
harperLogger.debug?.(`Error committing timed out transaction: ${error.message}`);
result = Promise.reject(error);
}
if ((result as any)?.then) {
const monitorCommit = result as Promise<CommitResolution>;
txn.monitorCommit = monitorCommit;
monitorCommit.then(
() => {
if (txn.monitorCommit === monitorCommit) txn.monitorCommit = undefined;
},
(error) => {
harperLogger.debug?.(`Error committing timed out transaction: ${error.message}`);
}
);
}
txn.timeout = Math.max(txnExpiration, txn.timeoutBudget ?? 0);
}
Expand Down
97 changes: 97 additions & 0 deletions unitTests/resources/txn-tracking.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -462,6 +462,83 @@ describe('Write txn timeout', () => {
}
});

describe("the owner's commit waits for the monitor's force-commit", () => {
function gateNativeCommit(link) {
const nativeTxn = link.transaction;
const commit = nativeTxn.commit;
let enter;
const entered = new Promise((resolve) => (enter = resolve));
let release;
const released = new Promise((resolve) => (release = resolve));
nativeTxn.commit = function (...args) {
nativeTxn.commit = commit;
enter();
return released.then((failure) => (failure ? Promise.reject(failure) : commit.apply(this, args)));
};
return { link, entered, release };
}

function runSourceApply(id, value, whileHeld) {
const context = { sourceApply: true };
let settled = false;
let onGated;
const gated = new Promise((resolve) => (onGated = resolve));
const committed = transaction(context, async () => {
await IndexedResource.put(id, { t: value }, context);
const gate = gateNativeCommit(databaseTxns(context)[0]);
await gate.entered;
onGated(gate);
await whileHeld?.(gate);
});
committed.then(
() => (settled = true),
() => (settled = true)
);
return { context, committed, gated, isSettled: () => settled };
}

beforeEach(function () {
if (isLMDB) this.skip();
setExpiration(20);
});
afterEach(() => setExpiration(30000));

it('resolves only after the force-commit lands', async function () {
await IndexedResource.put(403, { t: 1 });
const { committed, gated, isSettled } = runSourceApply(403, 2);
const gate = await gated;
await delay(50);
assert.equal(isSettled(), false, "transaction() must not resolve while the monitor's commit is in flight");
gate.release();
await committed;
assert.strictEqual((await IndexedResource.get(403))?.t, 2);
});

it('rejects when the force-commit fails while the owner is waiting on it', async function () {
await IndexedResource.put(404, { t: 1 });
const { context, committed, gated } = runSourceApply(404, 2);
const gate = await gated;
await delay(50);
gate.release(new Error('injected native commit failure'));
await assert.rejects(committed, /injected native commit failure/);
assert.strictEqual((await IndexedResource.get(404))?.t, 1);
assert.notStrictEqual(context.transaction, gate.link, 'the failed transaction must release its context');
await transaction.commit(context);
});

it('rejects when the force-commit failed before the handler returned', async function () {
await IndexedResource.put(405, { t: 1 });
const { committed } = runSourceApply(405, 2, async (gate) => {
const monitorCommit = gate.link.monitorCommit;
assert.ok(monitorCommit, 'test setup: the monitor must have submitted its commit');
gate.release(new Error('injected native commit failure'));
await monitorCommit.catch(() => {});
});
await assert.rejects(committed, /injected native commit failure/);
assert.strictEqual((await IndexedResource.get(405))?.t, 1);
});
});

describe('abort releases the native handle', () => {
// A write-first link (save() built the handle with no prior read) has no readTxnsUsed, so the
// refcount loop never runs and the handle was stranded — permanently, since rocksdb-js's
Expand Down Expand Up @@ -627,6 +704,26 @@ describe('Commit-phase pre-commit work is not poisoned by the monitor (#2062)',
}
}

it("finishes a landed store's bookkeeping when the next store's commit throws synchronously", async function () {
if (isLMDB) this.skip();
const context = {};
let links;
await assert.rejects(
transaction(context, async () => {
await SecondaryBlobResource.put({ id: 2080, value: 'head' }, context);
await ThirdResource.put({ id: 2080, value: 'next' }, context);
links = databaseTxns(context);
assert.equal(links.length, 2);
links[1].commit = () => {
throw new Error('injected synchronous next-store failure');
};
}),
/injected synchronous next-store failure/
);
assert.equal((await SecondaryBlobResource.get(2080))?.value, 'head', 'test setup: the head store must land');
assert.equal(links[0].writes.length, 0, "the landed head's write set must be cleared");
});

it('marks and clears the commit phase across an LMDB transaction chain', function () {
const head = new LMDBTransaction();
const next = new LMDBTransaction();
Expand Down
Loading