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
4 changes: 0 additions & 4 deletions .gitmodules
Original file line number Diff line number Diff line change
@@ -1,4 +0,0 @@
[submodule "core"]
path = core
url = git@github.com:HarperFast/harper.git
branch = main
203 changes: 203 additions & 0 deletions integrationTests/cluster/blockCacheEviction.test.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,203 @@
/**
* Integration test: replication must survive RocksDB `store.get()` returning a Promise.
*
* Background — the bug class this guards against:
* On RocksDB, `store.get(id)` returns a `MaybePromise`: the record SYNCHRONOUSLY when it is in the
* block cache / memtable, but a *Promise* on a cache miss that needs a disk read. (LMDB is always
* synchronous, so this is RocksDB-only.) Several replication paths read `hdb_nodes` with an un-awaited
* `primaryStore.get(...)` and consume the result synchronously (`?.replicates`, `?.url`, truthiness).
* While the row is cached `get()` is synchronous, so it works. Once the system database grows past the
* block cache — or immediately after a restart, when the cache is COLD — `get()` returns a Promise;
* `Promise?.replicates` is `undefined`, which silently disables replication / drops a node from
* cluster_status / never opens a retrieval connection. The fix is `getSync(...)` at those sites.
*
* Why this test reproduces it:
* - Each node runs RocksDB with a SMALL (but viable) block cache, so system-table blocks do not stay
* resident under churn. (A sub-MB cache makes RocksDB hang on open, so "small" here is ~32 MB, far
* below the default ~25%-of-RAM — the cold restart below is what makes the miss DETERMINISTIC.)
* - Every node is RESTARTED, giving a COLD block cache + empty memtable, so the first post-restart read
* of each `hdb_nodes` record is a guaranteed cache miss (a Promise from `get()`). The startup
* replication paths (ensureThisNode / shouldReplicateFromNode) run exactly in that window.
* - We then drive the procedures that depend on those synchronous reads — a rolling restart — and assert replication stays healthy and converges.
*
* Pre-fix this fails: a post-restart cache miss makes shouldReplicateFromNode falsy (unsubscribe) and/or
* flips `isFullyReplicating = false` ("Disabling replication"), so the post-restart write never converges.
* Post-fix the synchronous reads return the real records regardless of cache state, so it converges.
*/
import { suite, test, before, after } from 'node:test';
import { ok, equal } from 'node:assert/strict';
import { setTimeout as delay } from 'node:timers/promises';
import { startHarper, teardownHarper, getNextAvailableLoopbackAddress } from '@harperfast/integration-testing';
import { join } from 'node:path';
import { sendOperation, readLog } from './clusterShared.mjs';

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

// Small but viable RocksDB block cache: far below the default (~25% of RAM) so blocks are evicted under
// churn, yet large enough that opening Harper's databases does not hang (a sub-MB cache does). The cold
// restart is what makes the cache-miss deterministic; this keeps misses from being papered over by a
// warm cache. We deliberately do NOT shrink the WriteBufferManager — a tiny WBM with allowStall stalls
// the schema writes during startup.
const SMALL_ROCKS = { blockCacheSize: 32 * 1024 * 1024 };

// Some padded records so the data table spans multiple SST blocks (cache pressure). Convergence is
// checked via sentinel records, not a full count, so the exact number/pagination doesn't matter.
const SEED_RECORD_COUNT = 500;
const PADDING = 'x'.repeat(1024);

async function pollHealth(node, { retries = 60, intervalMs = 2000 } = {}) {
for (let i = 0; i < retries; i++) {
try {
const status = await sendOperation(node, { operation: 'cluster_status' });
if (status) return status;
} catch {
/* transient during restart */
}
await delay(intervalMs);
}
throw new Error(`node ${node.hostname} did not come healthy in time`);
}

async function waitForRecord(node, id, { retries = 90, intervalMs = 1000 } = {}) {
for (let i = 0; i < retries; i++) {
const rows = await sendOperation(node, {
operation: 'search_by_value',
database: 'data',
table: 'cache_evict_test',
search_attribute: 'id',
search_value: id,
get_attributes: ['id', 'value'],
}).catch(() => []);
const found = (rows ?? [])[0];
if (found) return found;
await delay(intervalMs);
}
return null;
}

suite(
'replication survives RocksDB get() cache-miss Promises (small block cache + restart)',
{ timeout: 240000 },
(ctx) => {
before(async () => {
const hostnameA = await getNextAvailableLoopbackAddress();
const hostnameB = await getNextAvailableLoopbackAddress();

const makeNodeCtx = (hostname) => ({ name: ctx.name, harper: { hostname } });

// Plaintext replication of BOTH 'data' and 'system' so hdb_nodes replicates across the pair
// (the system table is the one that falls out of the block cache).
const nodeConfig = (hostname) => ({
config: {
analytics: { aggregatePeriod: -1 },
logging: { colors: false, stdStreams: false, console: true },
replication: { port: hostname + ':9933', securePort: null, databases: ['data', 'system'] },
storage: { engine: 'rocksdb', rocks: SMALL_ROCKS },
},
env: { HARPER_NO_FLUSH_ON_EXIT: true },
});

const ctxA = makeNodeCtx(hostnameA);
const ctxB = makeNodeCtx(hostnameB);
await Promise.all([startHarper(ctxA, nodeConfig(hostnameA)), startHarper(ctxB, nodeConfig(hostnameB))]);
ctx.nodeA = ctxA.harper;
ctx.nodeB = ctxB.harper;

// Seed table + data on A, ending with a sentinel we can wait on.
await sendOperation(ctx.nodeA, {
operation: 'create_table',
database: 'data',
table: 'cache_evict_test',
primary_key: 'id',
});
const records = Array.from({ length: SEED_RECORD_COUNT }, (_, i) => ({
id: `seed-${i}`,
value: `v${i}`,
pad: PADDING,
}));
records.push({ id: 'seed-sentinel', value: 'seeded', pad: PADDING });
for (let i = 0; i < records.length; i += 250) {
await sendOperation(ctx.nodeA, {
operation: 'upsert',
database: 'data',
table: 'cache_evict_test',
records: records.slice(i, i + 250),
});
}

// B joins A as leader and full-copies the seed data.
await sendOperation(ctx.nodeB, {
operation: 'add_node',
hostname: ctx.nodeA.hostname,
rejectUnauthorized: false,
isLeader: true,
authorization: ctx.nodeA.admin,
});

const onB = await waitForRecord(ctx.nodeB, 'seed-sentinel');
ok(onB, 'node B should have received the seeded data (sentinel) before we start perturbing it');
});

after(async () => {
await Promise.all([
ctx.nodeA && teardownHarper({ harper: ctx.nodeA }),
ctx.nodeB && teardownHarper({ harper: ctx.nodeB }),
]);
});

test('cold-cache restart does not silently disable replication; cluster reconverges', async () => {
const { nodeA, nodeB } = ctx;

// Rolling restart -> COLD block cache on each node. The startup replication paths
// (ensureThisNode / shouldReplicateFromNode / cluster bootstrap) now read hdb_nodes from a cold
// cache, which is the exact get()->Promise condition this test guards.
for (const node of [nodeA, nodeB]) {
await sendOperation(node, { operation: 'restart' }).catch(() => {});
await pollHealth(node);
}

// (1) Direct catch for the silent-disable bug: a cold-cache Promise self-row logs
// "Disabling replication". It must not appear.
for (const node of [nodeA, nodeB]) {
const log = await readLog(node);
ok(
!/Disabling replication/.test(log),
`node ${node.hostname} logged "Disabling replication" after a cold-cache restart (get() Promise self-row)`
);
}

// (2) cluster_status must still report this node's own record after the cold restart
// (clusterStatus reads hdb_nodes for the self record; a Promise there omits node_name).
for (const node of [nodeA, nodeB]) {
const status = await pollHealth(node);
ok(status.node_name, `cluster_status on ${node.hostname} is missing node_name after restart`);
}

// (3) A write made AFTER the cold restart must converge to B. If a cold-cache get() Promise
// silently disabled replication / unsubscribed the peer, this never arrives.
await sendOperation(nodeA, {
operation: 'upsert',
database: 'data',
table: 'cache_evict_test',
records: [{ id: 'post-restart-1', value: 'after-cold-cache', pad: PADDING }],
});
const found = await waitForRecord(nodeB, 'post-restart-1');
ok(found, 'post-restart write did not converge to node B (replication silently disabled?)');
equal(found.value, 'after-cold-cache', 'wrong value converged to node B');
});

// NOTE: a dedicated remove_node→add_node cycle test was dropped here. Removing a node's *leader*
// leaves it with a null self-record so it (correctly) disables replication and does not re-converge
// within the window — a remove_node re-subscription behavior orthogonal to the get() MaybePromise
// fix this suite guards (the cold-restart test above already exercises ensureThisNode /
// shouldReplicateFromNode / cluster_status on a cold cache). add_node itself is covered by the
// `before` hook and by replicationReconnect.test.mjs / replicationTopology.test.mjs.
}
);
4 changes: 3 additions & 1 deletion replication/clusterStatus.ts
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,9 @@ export async function clusterStatus() {
// Add node name and shard/url info for this node
response.node_name = getThisNodeName();
// If it doesn't exist and or needs to be updated.
const thisNode = getHDBNodeTable().primaryStore.get(response.node_name);
// getSync (not get): a get() Promise on a cache miss has no .shard/.url, so cluster_status would
// silently omit this node's shard/url once hdb_nodes grows past the block cache.
const thisNode = getHDBNodeTable().primaryStore.getSync(response.node_name);
if (thisNode?.shard) response.shard = thisNode.shard;
if (thisNode?.url) response.url = thisNode.url;
response.is_enabled = true; // if we have replication, replication is enabled
Expand Down
92 changes: 78 additions & 14 deletions replication/knownNodes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,10 +12,41 @@ import * as env from '../core/utility/environment/environmentManager.js';
import { CONFIG_PARAMS } from '../core/utility/hdbTerms.ts';
import { logger } from '../core/utility/logging/logger.ts';

let hdbNodeTable;
type MaybePromise<T> = T | Promise<T>;

/**
* A row in the system `hdb_nodes` table — the canonical {@link Node} shape plus an index signature for
* the incidental extra fields some paths read (e.g. `authorization`). Reusing `Node` keeps it assignable
* everywhere a `Node` is expected (`getNodeURL`, `server.nodes`, …).
*/
export interface NodeRecord extends Node {
[key: string]: any;
}

/**
* Typed view of the `hdb_nodes` primaryStore. It is RocksDB-capable, so `get()` returns a MaybePromise —
* the record synchronously on a block-cache / memtable hit, but a *Promise* on a cache miss that needs a
* disk read. Typing `get()` this way makes any SYNCHRONOUS consumption (e.g. `store.get(id)?.replicates`)
* a compile error: use `getSync(id)` (forces the inline read) or `await` / `when` the result. LMDB is
* always synchronous so its Promise arm is never taken at runtime, but the type keeps callers honest on
* both engines. Other members fall through the index signature — this intentionally constrains only the
* point-read path the MaybePromise hazard affects.
*/
export type NodeStore = {
get(id: string, options?: any): MaybePromise<NodeRecord | undefined>;
getSync(id: string, options?: any): NodeRecord | undefined;
[key: string]: any;
};

interface HdbNodeTable {
primaryStore: NodeStore;
[key: string]: any;
}

let hdbNodeTable: HdbNodeTable | undefined;
server.nodes = [];

export function getHDBNodeTable() {
export function getHDBNodeTable(): HdbNodeTable {
return (
hdbNodeTable ||
(hdbNodeTable = table({
Expand Down Expand Up @@ -57,7 +88,7 @@ export function getHDBNodeTable() {
attribute: '__updatedtime__',
},
],
}))
}) as unknown as HdbNodeTable)
);
}
export function getReplicationSharedStatus(
Expand Down Expand Up @@ -268,10 +299,13 @@ export function readNodeForAuth(name: string, routeRecord?: any): any {
export function resolveNodeForAuth(store: any, name: string, routeRecord?: any): any {
let record: any;
try {
record = store.get(name);
// Synchronous point read. RocksDB get() is a MaybePromise (a Promise on a block-cache miss);
// an un-awaited get() would return a pending Promise that isValidNodeRecord rejects, stranding
// the peer. getSync forces the inline read. See selfNodeReplicates.
record = store.getSync(name);
} catch {
// A present-but-undecodable row throws here (missing shared structure); fall through to
// the physical-existence check below rather than treating the peer as unknown.
// Defensive: a read error here falls through to the range-visibility check rather than
// treating the peer as unknown.
record = undefined;
}
if (isValidNodeRecord(record)) return record;
Expand All @@ -286,7 +320,7 @@ export function resolveNodeForAuth(store: any, name: string, routeRecord?: any):
logger.warn?.(
'hdb_nodes record for',
name,
'did not decode to a valid node descriptor on the point lookup but is range-visible (likely a v5-era shared-structure row transiently misreading at boot during an in-place upgrade; see harper-pro#352). Authorizing the peer by its certificate-validated name; the full record self-heals on the next decodable update.'
'did not resolve to a valid node descriptor on the (synchronous) point lookup but is range-visible. Authorizing the peer by its certificate-validated name (defense-in-depth; the connection is independently TLS-validated). This path should be rare now that the point read is synchronous (see harper-pro#470).'
);
return { name };
}
Expand Down Expand Up @@ -340,7 +374,9 @@ function storeRecordRangeVisible(store: any, name: string): boolean {
if (typeof store.doesExist === 'function' && store.doesExist(name)) return true;
if (typeof store.getBinaryFast === 'function' && store.getBinaryFast(name) != null) return true;
try {
return store.get(name) != null;
// Synchronous read: RocksDB get() is a MaybePromise; an un-awaited get() would return a
// truthy Promise on a cache miss, making an absent key look present. See selfNodeReplicates.
return store.getSync(name) != null;
} catch {
return true;
}
Expand Down Expand Up @@ -371,9 +407,11 @@ async function processNodeUpdateEvent(event: any, listener: (node: any, id: stri
}
} else if (event.type === 'patch' && node_name !== getThisNodeName() && event.value?.isLeader !== undefined) {
// add_node { isLeader: true } reaches us as a patch event; read the merged
// record from LMDB so server.nodes reflects the full record (including isLeader).
const fullRecord = getHDBNodeTable().primaryStore.get(node_name);
if (fullRecord) server.nodes.push(fullRecord);
// record so server.nodes reflects the full record (including isLeader). Sync read
// (getSync): RocksDB get() is a MaybePromise (Promise on a block-cache miss); an un-awaited
// get() here would push a pending Promise into server.nodes. See selfNodeReplicates.
const fullRecord = getHDBNodeTable().primaryStore.getSync(node_name);
if (fullRecord) server.nodes.push(fullRecord as any);
}
const shards = new Map();
for await (const node of getHDBNodeTable().search({})) {
Expand Down Expand Up @@ -421,7 +459,8 @@ async function processNodeUpdateEvent(event: any, listener: (node: any, id: stri
}
} else if (event.type === 'patch' && event.value?.isLeader !== undefined) {
// isLeader patches need to drive subscription bootstrap; pass the merged record.
const fullRecord = getHDBNodeTable().primaryStore.get(event.id);
// Sync read (getSync) — RocksDB get() is a MaybePromise; see selfNodeReplicates.
const fullRecord = getHDBNodeTable().primaryStore.getSync(event.id);
if (fullRecord) listener(fullRecord, event.id);
}
}
Expand All @@ -439,7 +478,10 @@ async function processNodeUpdateEvent(event: any, listener: (node: any, id: stri
export function probeNodeRow(store: any, key: unknown): { outcome: 'deleted' | 'decode-failure'; record?: any } {
let record: any;
try {
record = store.get(key);
// Synchronous read: RocksDB get() is a MaybePromise. An un-awaited get() returns Promise<null>
// for a tombstone on a cache miss, so `record == null` would be false and a removed node would
// be misclassified as a decode-failure and revived. getSync forces the inline read. See #470.
record = store.getSync(key);
} catch {
// present-but-undecodable row → decode failure, reconstruct.
return { outcome: 'decode-failure' };
Expand Down Expand Up @@ -565,11 +607,33 @@ export function shouldReplicateFromNode(node: Node, databaseName: string) {
: dbReplication.name === databaseName &&
(!dbReplication.sharded || node.shard === env.get(CONFIG_PARAMS.REPLICATION_SHARD));
}))) &&
getHDBNodeTable().primaryStore.get(getThisNodeName())?.replicates) ||
selfNodeReplicates(getHDBNodeTable().primaryStore, getThisNodeName())) ||
node.subscriptions?.some((sub) => (sub.database || sub.schema) === databaseName && sub.subscribe)
);
}

/**
* Read this node's own `replicates` flag from its hdb_nodes self-record. This is the "does THIS node
* participate in replication at all" gate at the end of {@link shouldReplicateFromNode}; it is consulted
* by BOTH the wedge backstop (findWedgedNodeUrls -> reconcileWorkers) and the onDatabase re-subscribe
* path, in synchronous contexts (the `isDesired` filter, the change-stream listener).
*
* It MUST use the synchronous point read. The system database is RocksDB, whose `get()` returns a
* `MaybePromise`: the value synchronously when the row is in the block cache / memtable, but a `Promise`
* on a cache miss that needs a disk read. An un-awaited `get()` therefore works only while `system` is
* small enough that this row stays cached; once `system` grows past block-cache size the row is evicted,
* `get()` returns a Promise, and `(promise)?.replicates` silently becomes `undefined` — disabling ALL
* replication recovery for a still-desired peer (the observed preprod wedge, where the self row read as
* an empty object that was actually a pending Promise). `getSync()` forces the synchronous read (doing
* the disk read inline on a cache miss), which is what these sync callers require.
*
* A deleted/absent self-record yields `null`/`undefined` (`?.` -> falsy), which is correct; a genuine
* `replicates: false` is preserved. Store is passed explicitly so it stays unit-testable.
*/
export function selfNodeReplicates(store: any, name: string): any {
return store.getSync(name)?.replicates;
}

const replicationConfirmationFloat64s = new Map<string, Map<string, Float64Array>>();
/** Ensure that the shared user buffers are instantiated so we can communicate through them
*/
Expand Down
4 changes: 3 additions & 1 deletion replication/replicator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -551,7 +551,9 @@ function getRetrievalConnectionByName(nodeName, subscription, dbName): NodeRepli
}
let connection = dbConnections.get(dbName);
if (isReusableConnection(connection)) return connection;
const node = getHDBNodeTable().primaryStore.get(nodeName);
// getSync (not get): get() returns a Promise on a RocksDB block-cache miss; `Promise?.url` is undefined,
// so the connection would silently never be created (no error) for a node that actually exists.
const node = getHDBNodeTable().primaryStore.getSync(nodeName);
if (node?.url) {
connection = new NodeReplicationConnection(getNodeURL(node), subscription, dbName, nodeName, node.authorization);
// cache the connection
Expand Down
Loading
Loading