Skip to content
Closed
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
53 changes: 43 additions & 10 deletions replication/knownNodes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -268,10 +268,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 +289,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 +343,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,8 +376,10 @@ 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);
// 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);
}
const shards = new Map();
Expand Down Expand Up @@ -421,7 +428,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 +447,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 +576,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
15 changes: 10 additions & 5 deletions unitTests/replication/readNodeForAuth.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -37,11 +37,15 @@ import { resolveNodeForAuth, isValidNodeRecord } from '#src/replication/knownNod
function fakeStore({ data = {}, present = null, rangeKeys = null, throwOnGet = new Set() } = {}) {
const presentSet = present ?? new Set(Object.keys(data));
const rangeSet = rangeKeys ?? presentSet;
const pointGet = (key) => {
if (throwOnGet.has(key)) throw new Error('Record id is not defined for 0');
return data[key];
};
return {
get(key) {
if (throwOnGet.has(key)) throw new Error('Record id is not defined for 0');
return data[key];
},
get: pointGet,
// resolveNodeForAuth uses the SYNCHRONOUS point read (getSync) because the real system store is
// RocksDB, whose get() is a MaybePromise. getSync mirrors the synchronous get here.
getSync: pointGet,
doesExist(key) {
return presentSet.has(key);
},
Expand Down Expand Up @@ -124,7 +128,8 @@ describe('readNodeForAuth precondition (real msgpackr shared-structure decode)',

// Late-flipped node: the structure table was never persisted on this node, so decode fails.
const readerMissing = new Packr({ useRecords: true, maxSharedStructures: 32, getStructures: () => undefined });
let decodedMissing, threw = false;
let decodedMissing,
threw = false;
try {
decodedMissing = readerMissing.unpack(bytes);
} catch {
Expand Down
24 changes: 15 additions & 9 deletions unitTests/replication/scanNodesForSubscription.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -44,17 +44,21 @@ const DECODE_THROWS = Symbol('decode-throws');
* null (genuine tombstone, skip); any object ⇒ get() returns it.
*/
function fakeStore(rows) {
// probeNodeRow uses the SYNCHRONOUS point read (getSync) because the real system store is RocksDB,
// whose get() is a MaybePromise; getSync mirrors the synchronous point lookup here.
const pointGet = (key) => {
const row = rows[key];
if (row === DECODE_THROWS) throw new Error('missing shared structure');
return row ?? null;
};
return {
getRange() {
return Object.keys(rows)
.sort()
.map((key) => ({ key, value: rows[key] === DECODE_THROWS ? null : rows[key] }));
},
get(key) {
const row = rows[key];
if (row === DECODE_THROWS) throw new Error('missing shared structure');
return row ?? null;
},
get: pointGet,
getSync: pointGet,
};
}

Expand Down Expand Up @@ -107,28 +111,30 @@ describe('resolveScannedNode (harper-pro#460)', () => {
});

describe('probeNodeRow (harper-pro#460 review: tombstone vs decode failure)', () => {
// probeNodeRow reads via getSync (RocksDB get() is a MaybePromise — the synchronous read is required
// so a cache-miss tombstone is a clean null, not a truthy Promise<null> misclassified as a decode failure).
it('classifies a point lookup that THROWS as a decode failure', () => {
const store = {
get: () => {
getSync: () => {
throw new Error('missing shared structure');
},
};
expect(probeNodeRow(store, 'peer-a')).to.deep.equal({ outcome: 'decode-failure' });
});

it('classifies a clean null point lookup as a genuine tombstone (deleted)', () => {
const store = { get: () => null };
const store = { getSync: () => null };
expect(probeNodeRow(store, 'peer-a')).to.deep.equal({ outcome: 'deleted' });
});

it('classifies undefined (physically absent) as deleted', () => {
const store = { get: () => undefined };
const store = { getSync: () => undefined };
expect(probeNodeRow(store, 'peer-a')).to.deep.equal({ outcome: 'deleted' });
});

it('returns the recovered record when the point lookup succeeds where the range value was null', () => {
const rec = { name: 'peer-a', replicates: true };
const store = { get: () => rec };
const store = { getSync: () => rec };
expect(probeNodeRow(store, 'peer-a')).to.deep.equal({ outcome: 'decode-failure', record: rec });
});
});
Expand Down
86 changes: 86 additions & 0 deletions unitTests/replication/selfNodeReplicates.test.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
/**
* Coverage for selfNodeReplicates — the "does THIS node participate in replication" self-record read at
* the end of shouldReplicateFromNode (replication/knownNodes.ts).
*
* Root cause of the observed preprod wedge (CDP-confirmed on the live cluster): 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. The replication self-gate did
* an un-awaited `primaryStore.get(getThisNodeName())?.replicates`, so once `system` grew past block-cache
* size the self row was evicted, `get()` returned a Promise, and `(promise)?.replicates` was `undefined`
* — making the WHOLE shouldReplicateFromNode predicate falsy. That predicate is the `isDesired` gate for
* BOTH the wedge backstop (findWedgedNodeUrls -> reconcileWorkers) AND the onDatabase re-subscribe path,
* so a still-desired peer was silently excluded from all recovery (connected:false for hours, retries:0).
*
* The fix: selfNodeReplicates uses the SYNCHRONOUS point read (`getSync`), which forces the inline read
* regardless of cache state. A deleted/absent self-record yields null/undefined (falsy, correct); a
* genuine replicates:false is preserved.
*/

import { expect } from 'chai';
import { selfNodeReplicates } from '#src/replication/knownNodes';

/**
* Store stub modeling the RocksDB MaybePromise contract.
* - `get(key)` returns the value synchronously when the key is "cached" (block cache / memtable), and a
* resolved Promise on a cache miss — exactly the path that silently broke the old un-awaited gate.
* - `getSync(key)` always returns the value synchronously (forcing the inline disk read on a miss), and
* `undefined` for an absent key.
*/
function rocksLikeStore(disk = {}, cachedKeys = []) {
const cached = new Set(cachedKeys);
const has = (key) => Object.prototype.hasOwnProperty.call(disk, key);
return {
get(key) {
if (!has(key)) return undefined;
return cached.has(key) ? disk[key] : Promise.resolve(disk[key]); // cache miss -> Promise
},
getSync(key) {
return has(key) ? disk[key] : undefined;
},
};
}

describe('selfNodeReplicates', () => {
const SELF = 'node-a';

// THE bug-closing case: the self row is NOT in the block cache, so RocksDB get() returns a Promise.
// The old un-awaited `get(SELF)?.replicates` read undefined off that Promise; getSync reads the value.
it('returns the real replicates value on a block-cache MISS (get() returns a Promise)', () => {
const store = rocksLikeStore({ [SELF]: { name: SELF, replicates: true } } /* cachedKeys: none */);
// Demonstrate the hazard the fix removes: the async get() yields a Promise, not the record.
const asyncResult = store.get(SELF);
expect(typeof asyncResult.then).to.equal('function');
expect(asyncResult?.replicates).to.equal(undefined); // the old code path -> falsy gate -> wedge
// The fix:
expect(selfNodeReplicates(store, SELF)).to.equal(true);
});

it('returns the real replicates value on a cache HIT (get() is synchronous) too', () => {
const store = rocksLikeStore({ [SELF]: { name: SELF, replicates: true } }, [SELF]);
expect(selfNodeReplicates(store, SELF)).to.equal(true);
});

it('returns a replicates OBJECT verbatim (sends/sendsTo form)', () => {
const replicates = { sends: true, sendsTo: ['node-b'] };
const store = rocksLikeStore({ [SELF]: { name: SELF, replicates } });
expect(selfNodeReplicates(store, SELF)).to.deep.equal(replicates);
});

// False-positive guard: a genuine, decodable replicates:false must be preserved, not coerced truthy.
it('preserves a genuine replicates: false (does not force true)', () => {
const store = rocksLikeStore({ [SELF]: { name: SELF, replicates: false } });
expect(selfNodeReplicates(store, SELF)).to.equal(false);
});

// A genuinely-absent self-record stays undefined (falsy) — we do not invent a self record.
it('returns undefined when there is no self-record at all', () => {
const store = rocksLikeStore({});
expect(selfNodeReplicates(store, SELF)).to.equal(undefined);
});

// A removed-node tombstone is a clean null; `?.replicates` is undefined (falsy) — do NOT revive it.
it('returns undefined for a clean-null tombstone (does NOT revive a removed node)', () => {
const store = rocksLikeStore({ [SELF]: null });
expect(selfNodeReplicates(store, SELF)).to.equal(undefined);
});
});
Loading