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
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
jsResource:
files: resources.js # entry module — runs at component load
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
// Test-only component: monkey-patches `fs.open` so that read-mode (`'r'`) opens of a DETERMINISTIC
// subset of blob files ALWAYS fail with `ENOENT` — modelling blobs that are GONE AT THE SOURCE (e.g.
// evicted/expired on an expiration cache table). When the replication SENDER tries to read such a blob
// to stream it (`sendBlobs` -> `blob.stream()` -> blob.ts read path's `fs.open(path, 'r', cb)`), the
// read fails; `sendBlobs` catches it and forwards a BLOB_CHUNK `error` marker carrying `errorCode:
// 'ENOENT'`, which is the trigger for the receive-side "source-reported permanent unavailable" branch
// under test (receiveBlobs -> markSourceBlobUnavailable/isUnrecoverableSourceBlobError -> advance the
// resume cursor past it instead of holding `hasBlobGap` forever). See harper-pro#403.
//
// Why a deterministic subset (by fileId) rather than "every Nth open": blob.ts's read path RETRIES on
// ENOENT (up to 1000×) while the blob might still be mid-write, so a counter that fails only one open
// per blob lets the retry succeed and the error never reaches `sendBlobs`. Keying the failure to the
// fileId (the path basename) makes EVERY open of those blobs fail — retries included — so the ENOENT
// deterministically propagates. The selected blobs are permanently unsendable; the rest replicate.
//
// Distinct from fixture-blob-fail-transient / fixture-blob-fail-injector, which patch the RECEIVER's
// `createWriteStream` to model a LOCAL save fault (which must keep HOLDING the cursor, not advance).
// Install on the SOURCE node. Toggle on with HARPER_TEST_BLOB_READ_FAIL_MODULUS=<positive int>: a blob
// is failed when parseInt(fileId, 16) % modulus === 0 (e.g. 5 fails ~1 in 5 blobs).
import { createRequire } from 'node:module';

const modulus = Number.parseInt(process.env.HARPER_TEST_BLOB_READ_FAIL_MODULUS || '0', 10);
if (Number.isFinite(modulus) && modulus > 0) {
const require = createRequire(import.meta.url);
const fs = require('node:fs');
const path = require('node:path');
const realOpen = fs.open;
let failedReads = 0;
const shouldFail = (p) => {
if (typeof p !== 'string' || !p.includes('/blobs/')) return false;
const fileId = path.basename(p);
const id = Number.parseInt(fileId, 16);
return Number.isFinite(id) && id % modulus === 0;
};
fs.open = function patchedOpen(p, flags, mode, cb) {
// fs.open signatures: (path, cb) | (path, flags, cb) | (path, flags, mode, cb). Normalize so we
// can inspect flags and find the callback regardless of arity.
const callback = typeof cb === 'function' ? cb : typeof mode === 'function' ? mode : flags;
const readMode = flags === 'r' || flags === undefined; // default flag is 'r'
if (readMode && typeof callback === 'function' && shouldFail(p)) {
const err = new Error("ENOENT: no such file or directory, open '" + p + "'");
err.code = 'ENOENT';
err.errno = -2;
err.syscall = 'open';
err.path = p;
console.log('[blob-fail-source-read] failing read open #' + ++failedReads + ' ' + p);
process.nextTick(() => callback(err));
return;
}
return realOpen.apply(this, arguments);
};
console.log(
'[blob-fail-source-read] installed; failing /blobs/ reads where parseInt(fileId,16) % ' + modulus + ' === 0'
);
}
184 changes: 184 additions & 0 deletions integrationTests/cluster/replicationBlobSourceUnavailable.test.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,184 @@
/**
* Regression guard for the permanent blob-replication WEDGE on an expiration/cache table (harper-pro#403).
*
* Mechanism (replication/replicationConnection.ts):
* - When the SENDER cannot read a blob to stream it (`sendBlobs` -> `blob.stream()` hits ENOENT —
* e.g. the blob was evicted/expired at the source), it forwards a BLOB_CHUNK `error` marker.
* - The receiver `stream.destroy()`s the blob stream with that error and `saveBlob` rejects.
* - BEFORE the fix, `receiveBlobs` treated every save failure identically: it set `hasBlobGap`, which
* pins the persisted resume cursor. On reconnect the source re-streams the SAME blob, hits the same
* ENOENT, and the cursor holds again — forever. `blobReplicationFailures` climbs and never drains.
* - The fix classifies a source-reported permanent absence (`isUnrecoverableSourceBlobError`) and
* advances the cursor PAST it (recorded loudly), leaving the record for backfill (harper-pro#388),
* while still HOLDING for genuinely local/transient save faults (the createWriteStream injectors).
*
* This is the end-to-end wiring the unit test (unitTests/replication/blobReplicationFailure.test.mjs)
* cannot cover: source read ENOENT -> `error` marker over the WS -> receiver takes the advance branch.
* The discriminating signal is B's log: FIXED code emits "advancing the resume cursor past it"; UNFIXED
* code emits only "Blob save failed for ..." and never advances. Heavy/stress-gated like its sibling
* blobGapDeadlock.test.mjs.
*/

import { suite, test, before, after } from 'node:test';
import { ok } from 'node:assert';
import { setTimeout as delay } from 'node:timers/promises';
import {
startHarper,
teardownHarper,
setupHarperWithFixture,
getNextAvailableLoopbackAddress,
targz,
} from '@harperfast/integration-testing';
import { join } from 'node:path';
import { sendOperation, fetchWithRetry, concurrent, readLog } from './clusterShared.mjs';

process.env.HARPER_INTEGRATION_TEST_INSTALL_SCRIPT = join(
import.meta.dirname ?? module.path,
'..',
'..',
'dist',
'bin',
'harper.js'
);

const RECORDS = Number.parseInt(process.env.HARPER_TEST_SRC_UNAVAIL_RECORDS || '120', 10);
const REQ_CONCURRENCY = Number.parseInt(process.env.HARPER_TEST_SRC_UNAVAIL_CONCURRENCY || '20', 10);
// Fail reads of a deterministic subset of blobs on the SOURCE (fileId % MODULUS === 0) so those blobs
// are PERMANENTLY unsendable (retries fail too), mirroring a scatter of evicted/expired blobs, while
// the rest replicate normally.
const READ_FAIL_MODULUS = process.env.HARPER_TEST_BLOB_READ_FAIL_MODULUS || '5';
const BLOB_CHUNKS = process.env.HARPER_TEST_BLOB_CHUNKS || '16'; // 16 * 4096 = 64 KB per blob (file-backed)
const CONVERGE_TIMEOUT_MS = Number.parseInt(process.env.HARPER_TEST_SRC_UNAVAIL_CONVERGE_MS || '60000', 10);

const STRESS = process.env.HARPER_RUN_STRESS_TESTS === '1';

suite('Source-unavailable blob does not permanently wedge replication', { skip: !STRESS, timeout: 300000 }, (ctx) => {

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.

Does this need to be explicitly wired up in the ci stress matrix?

before(async () => {
const nodeA = { name: ctx.name, harper: { hostname: await getNextAvailableLoopbackAddress() } };
const nodeB = { name: ctx.name, harper: { hostname: await getNextAvailableLoopbackAddress() } };
const sharedConfig = (host) => ({
analytics: { aggregatePeriod: -1 },
logging: { colors: false, console: true, level: 'debug' },
replication: { securePort: host + ':9933' },
});
// A is the SENDER: install the source read-fail injector so a subset of its blobs are unsendable.
await setupHarperWithFixture(nodeA, join(import.meta.dirname, 'fixture-blob-fail-source-read'), {
config: sharedConfig(nodeA.harper.hostname),
env: {
HARPER_NO_FLUSH_ON_EXIT: true,
HARPER_TEST_BLOB_CHUNKS: BLOB_CHUNKS,
HARPER_TEST_BLOB_READ_FAIL_MODULUS: READ_FAIL_MODULUS,
},
});
await startHarper(nodeB, {
config: sharedConfig(nodeB.harper.hostname),
env: { HARPER_NO_FLUSH_ON_EXIT: true, HARPER_TEST_BLOB_CHUNKS: BLOB_CHUNKS },
});
ctx.nodes = [nodeA.harper, nodeB.harper];

// Connect A↔B.
const tokenResp = await sendOperation(ctx.nodes[0], {
operation: 'create_authentication_tokens',
authorization: ctx.nodes[0].admin,
});
await sendOperation(ctx.nodes[1], {
operation: 'add_node',
rejectUnauthorized: false,
hostname: ctx.nodes[0].hostname,
authorization: 'Bearer ' + tokenResp.operation_token,
});
for (let retries = 0; retries < 15; retries++) {
const status = await Promise.all(ctx.nodes.map((n) => sendOperation(n, { operation: 'cluster_status' })));
if (status.every((r) => (r.connections ?? []).every((c) => (c.database_sockets ?? []).every((s) => s.connected))))
break;
await delay(200 * (retries + 1));
}

// Deploy the caching blob component to A; `replicated: true` installs it on B too.
const payload = await targz(join(import.meta.dirname, 'fixture-blob-gap-deadlock-source'));
await sendOperation(ctx.nodes[0], {
operation: 'deploy_component',
project: 'blob-gap-deadlock-source',
payload,
replicated: true,
restart: true,
});
await delay(35000);

const bootLog = await readLog(ctx.nodes[0]);
ok(
bootLog.includes('[blob-fail-source-read] installed'),
'source read-fail injector did not load on A — test would not exercise the source-unavailable path'
);

// Wait for replication to be LIVE before driving load (same race guard as blobGapDeadlock).
const probeId = 9_000_000 + Math.floor(Math.random() * 1000);
await fetchWithRetry(ctx.nodes[0].httpURL + '/Prerender/' + probeId);
let live = false;
for (let i = 0; i < 90; i++) {
const bDesc = await sendOperation(ctx.nodes[1], { operation: 'describe_table', table: 'Prerender' });
if ((bDesc.record_count ?? 0) > 0) {
live = true;
break;
}
await delay(1000);
}
ok(live, 'replication from A to B never went live within 90s — setup race, not a wedge');
});

after(async () => {
if (!ctx.nodes) return;
await Promise.all(ctx.nodes.map((n) => teardownHarper({ harper: n })));
});

test('receiver advances past a source-unavailable blob instead of wedging the cursor', async () => {
const [A, B] = ctx.nodes;

// Drive RECORDS distinct cache misses on A (each mints a new file-backed blob), continuously,
// so the source's replication send hits the read-fail injector on a subset of blobs.
let nextId = 0;
const { execute, finish } = concurrent(() => fetchWithRetry(A.httpURL + '/Prerender/' + nextId++), REQ_CONCURRENCY);
const loadPromise = (async () => {
for (let i = 0; i < RECORDS; i++) await execute();
await finish();
})();

// B should keep converging — the apply loop never blocks on blobs — and, critically, take the
// new advance branch for the source-unavailable blobs rather than holding the cursor forever.
const deadline = Date.now() + CONVERGE_TIMEOUT_MS;
let bCount = 0;
let aCount = 0;
let loadDone = false;
loadPromise.then(() => (loadDone = true)).catch(() => (loadDone = true));
while (Date.now() < deadline) {
aCount = (await sendOperation(A, { operation: 'describe_table', table: 'Prerender' })).record_count;
bCount = (await sendOperation(B, { operation: 'describe_table', table: 'Prerender' })).record_count;
if (loadDone && bCount >= aCount && aCount > 0) break;
await delay(1000);
}
await loadPromise.catch(() => {});

const aLog = await readLog(A);
const bLog = await readLog(B);
const sourceReadFailures = (aLog.match(/\[blob-fail-source-read\] failing read open/g) ?? []).length;
const advancedPast = (bLog.match(/is unrecoverable at source .* advancing the resume cursor past it/g) ?? [])
.length;
console.log(
`source-unavailable: A=${aCount} B=${bCount} sourceReadFailures=${sourceReadFailures} advancedPast=${advancedPast}`
);

ok(
sourceReadFailures > 0,
`expected the source read-fail injector to fire at least once, found ${sourceReadFailures}`
);
// The discriminating assertion: the receiver took the advance branch for a source-reported
// unavailable blob. On unfixed code this branch does not exist — B logs only "Blob save failed
// for ..." and pins the cursor — so `advancedPast` is 0 and this fails.
ok(
advancedPast > 0,
`receiver never advanced past a source-unavailable blob (advancedPast=${advancedPast}); the cursor would wedge (harper-pro#403)`
);
// Sanity: the apply loop kept up (records still flow; the wedge is about the cursor, not visibility).
ok(bCount >= aCount, `B did not converge on record count: ${bCount}/${aCount}`);
});
});
2 changes: 2 additions & 0 deletions replication/DESIGN.md
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,8 @@ Schema (defined in that function): `name` (PK), `subscriptions[]`, `system_info`

6. **The initial bulk clone copy is resumable (PK cursor).** When a follower requests a full copy (`startTime: 0`), the leader sends `COPY_START{copyStartTime}`, walks each table's primary store in key order, flushes a checkpoint transaction every `COPY_CHECKPOINT_RECORDS` (timed at `copyStartTime` so the persisted `seqId` stays pinned there, never a record's `localTime`), and sends `COPY_COMPLETE` at the end. The follower persists a cursor `{copyStartTime, currentTable, afterKey}` under `dbisDB` key `Symbol.for('copyCursor')` — but only in the `end_txn` `onCommit`, **after** the batch commits, so the cursor can never get ahead of committed data (a resume re-copies a few records idempotently but never skips). On reconnect, `sendSubscriptionRequestUpdate` reads the cursor and sends it as `copyResume` on the subscription request (overriding the persisted `seqId`, which alone would skip the un-copied tables); the leader skips tables before `currentTable` (stable iteration order ⇒ already committed) and resumes `currentTable` after `afterKey`. `COPY_COMPLETE` clears the cursor so subsequent connections resume normally from `seqId`. Before this, an interrupted copy restarted from zero and never converged for a large table (issue #241).

8. **Blob durability watermark: holds on a local/transient gap, advances past a source-missing (ENOENT) one.** Records commit (== become visible) without waiting on their blobs; the persisted resume cursor instead tracks `lastDurableSequenceId`, which only advances to a committed sequence once that sequence's blobs (and all earlier ones) are durably saved (the `.finally` watermark advance + the `onCommit` clamp, both gated on `!hasBlobGap`). A blob save failure in `receiveBlobs` is classified. A **local/transient** fault (receiver `createWriteStream` ENOENT, disk full, mid-stream timeout) sets `hasBlobGap`, pinning the watermark so a reconnect re-streams and re-saves the blob — no silent loss (#368/#386). A **source-reported PERMANENT** failure — the sender's `sendBlobs` catch forwarded a `BLOB_CHUNK` `error` marker with `errorCode: 'ENOENT'` because the blob is gone at the origin (evicted/expired) — is unrecoverable: re-streaming reproduces it, so holding would wedge the connection forever. The receiver instead logs it loudly (`cluster_status.blobReplicationFailures` + a per-blob "advancing the resume cursor past it" error), advances, and leaves the diverged record for proactive blob backfill (#388). Classification is deliberately narrow (`isPermanentSourceBlobErrorCode` = ENOENT only): a transient sender fault (EIO, EMFILE, timeout) or an older sender that doesn't forward `errorCode` stays unmarked, so it HOLDS like a local gap and a reconnect retries — never silently skipping a recoverable blob. The trigger is set on the destroy error via `markSourceBlobUnavailable`; the save `.catch` keys on `isUnrecoverableSourceBlobError`. See harper-pro#403.

---

## Tests
Expand Down
Loading
Loading