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
156 changes: 156 additions & 0 deletions integrationTests/cluster/replicationWedgeRecovery.test.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,156 @@
/**
* Open-but-idle replication wedge recovery (harper-pro#420).
*
* Field incident (JJill preprod, 5.1.5): after a cluster-wide simultaneous restart a per-DB receive
* socket settled "open but idle" — connected at the transport level, no bytes flowing, and no `close`
* event ever fired. The receive watchdog's `ws.terminate()` did not lead to recovery, the connection's
* node entry stayed `connected:true`, and the wedged `(peer, db)` pair made zero further connection
* attempts for over an hour, blocking replicated deploys. Only a manual staggered restart cleared it.
*
* The fix drives recovery from the watchdog through `NodeReplicationConnection.forceReconnect()`, which
* tears the socket down and schedules a fresh `connect()` independent of whether `close` ever fires.
*
* This test reproduces the exact no-`close` condition deterministically via the env-gated, one-shot
* `armReplicationWedgeForTest` hook: the first receive connection for the named DB is forced into the
* open-but-idle wedge (the watchdog observes frozen bytes, and the socket's terminate/close are
* neutralized so no `close` arrives). On the pre-#420 code this stays wedged forever (the test would time
* out); with the fix the connection reconnects on its own and replication resumes — with no restart.
*
* Proof of recovery is end-to-end: a record written on the source AFTER the wedge fires must replicate to
* the wedged subscriber, and `cluster_status` must return to `connected:true`.
*/
import { suite, test, before, after } from 'node:test';
import { ok } from 'node:assert';
import { setTimeout as delay } from 'node:timers/promises';
import { startHarper, teardownHarper, getNextAvailableLoopbackAddress } from '@harperfast/integration-testing';
import { join } from 'node:path';
import { sendOperation } from './clusterShared.mjs';

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

const NODE_COUNT = 2;
const WEDGE_DB = 'data';
// Short ping/watchdog windows so the receive watchdog fires within the test instead of the 60s default.
// Healthy connections pong every pingInterval (1s) so their watchdogs never false-fire; only the
// frozen-bytes wedged connection trips at pingTimeout.
const PING_INTERVAL_MS = 1000;
const PING_TIMEOUT_MS = 3000;
const RECOVERY_TIMEOUT_MS = 30000;
const POLL_INTERVAL_MS = 250;

function nodeStartOptions(node, { wedge = false } = {}) {
return {
config: {
analytics: { aggregatePeriod: -1 },
logging: { colors: false, stdStreams: true, console: true },
replication: {
securePort: node.hostname + ':9933',
databases: [WEDGE_DB],
pingInterval: PING_INTERVAL_MS,
pingTimeout: PING_TIMEOUT_MS,
},
},
// The wedge hook is per-process and one-shot; arming it only on the subscriber pins which
// (peer, db) receive socket gets forced open-but-idle.
env: wedge ? { HARPER_TEST_REPLICATION_WEDGE_DB: WEDGE_DB } : undefined,
};
}

async function dataSocketConnected(node) {
const status = await sendOperation(node, { operation: 'cluster_status' });
return status.connections.some(
(conn) => conn.database_sockets?.length > 0 && conn.database_sockets.every((socket) => socket.connected === true)
);
}

suite('Replication open-but-idle wedge recovery', { timeout: 120000 }, (ctx) => {
before(async () => {
// node[0] is the source, node[1] the subscriber whose receive socket is wedged.
ctx.nodes = [];
for (let i = 0; i < NODE_COUNT; i++) {
const nodeCtx = { name: ctx.name, harper: { hostname: await getNextAvailableLoopbackAddress() } };
const wedge = i === 1; // only the subscriber arms the one-shot wedge hook
ctx.nodes[i] = (await startHarper(nodeCtx, nodeStartOptions(nodeCtx.harper, { wedge }))).harper;
}
await Promise.all(
ctx.nodes.map((node) =>
sendOperation(node, {
operation: 'create_table',
database: WEDGE_DB,
table: 'test',
primary_key: 'id',
attributes: [
{ name: 'id', type: 'ID' },
{ name: 'name', type: 'String' },
],
})
)
);
});

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

test('a wedged open-but-idle receive socket reconnects on its own (no restart)', async () => {
// node1 subscribes to node0 for `data`. The first receive connection arms the wedge hook.
await sendOperation(ctx.nodes[1], {
operation: 'add_node',
rejectUnauthorized: false,
hostname: ctx.nodes[0].hostname,
authorization: ctx.nodes[1].admin,
});

// Let the wedge fire (watchdog at ~PING_TIMEOUT) and forceReconnect re-establish. The window is
// past the watchdog threshold plus a reconnect/backoff margin, so the record written next can only
// arrive over the RECOVERED socket — not the original pre-wedge one.
await delay(PING_TIMEOUT_MS + 5000);

const recordId = 'after-wedge-1';
await sendOperation(ctx.nodes[0], {
operation: 'insert',
database: WEDGE_DB,
table: 'test',
records: [{ id: recordId, name: 'recovered' }],
});

// Poll the wedged subscriber until the post-wedge write lands — recovery without a restart.
const deadline = Date.now() + RECOVERY_TIMEOUT_MS;
let replicated = false;
while (Date.now() < deadline) {
const result = await sendOperation(ctx.nodes[1], {
operation: 'search_by_id',
database: WEDGE_DB,
table: 'test',
ids: [recordId],
get_attributes: ['*'],
});
if (Array.isArray(result) && result.some((r) => r?.id === recordId)) {
replicated = true;
break;
}
await delay(POLL_INTERVAL_MS);
}
ok(replicated, 'record written after the wedge must replicate to the recovered subscriber (no restart)');

// And the socket-level view should be back to connected.
const deadlineConnected = Date.now() + RECOVERY_TIMEOUT_MS;
let connected = false;
while (Date.now() < deadlineConnected) {
if (await dataSocketConnected(ctx.nodes[1])) {
connected = true;
break;
}
await delay(POLL_INTERVAL_MS);
}
ok(connected, 'cluster_status should report the recovered data socket as connected');
});
});
107 changes: 104 additions & 3 deletions replication/replicationConnection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -296,6 +296,31 @@
return failureCount >= threshold && !alreadyLogged;
}

// Test-only fault injection for the harper-pro#420 regression test. When HARPER_TEST_REPLICATION_WEDGE_DB
// names a database, the FIRST receive (subscription) connection for it is forced into the open-but-idle
// wedge: the socket is paused (no frames arrive, so nothing re-arms the watchdog) and its terminate/close
// are neutralized so no 'close' event ever fires. The caller also freezes the watchdog's observed byte
// count for this connection — `socket.bytesRead` counts OS-buffered bytes (incl. pong frames) even while
// paused, so without freezing it the watchdog would keep seeing "activity" and never trip. Together these
// reproduce the field condition where recovery cannot come from the close handler and must come from
// forceReconnect's close-independent reconnect — so the test stays wedged on the pre-#420 bare-terminate()
// code and only recovers with the fix. One-shot per worker thread, so the reconnect's fresh socket
// recovers normally. Never arms in production: the env var is set only by the regression test.
let replicationWedgeForTestArmed = false;
export function armReplicationWedgeForTest(connection: any, ws: WebSocket, databaseName?: string): boolean {
// Guard the env var first: an unset var is undefined, and `undefined !== undefined` is false, so a
// connection with an undefined databaseName would otherwise arm the wedge in production.
if (!process.env.HARPER_TEST_REPLICATION_WEDGE_DB) return false;
if (!connection || replicationWedgeForTestArmed || process.env.HARPER_TEST_REPLICATION_WEDGE_DB !== databaseName)
return false;
replicationWedgeForTestArmed = true;
logger.warn?.(`[test] forcing open-but-idle replication wedge for db "${databaseName}" (harper-pro#420)`);
ws.terminate = () => {};
ws.close = () => {};
ws._socket?.pause();

Check failure on line 320 in replication/replicationConnection.ts

View workflow job for this annotation

GitHub Actions / Unit Tests (Node.js v22)

Property '_socket' does not exist on type 'WebSocket'.

Check failure on line 320 in replication/replicationConnection.ts

View workflow job for this annotation

GitHub Actions / Unit Tests (Node.js v26)

Property '_socket' does not exist on type 'WebSocket'.

Check failure on line 320 in replication/replicationConnection.ts

View workflow job for this annotation

GitHub Actions / Unit Tests (Node.js v24)

Property '_socket' does not exist on type 'WebSocket'.

Check failure on line 320 in replication/replicationConnection.ts

View workflow job for this annotation

GitHub Actions / Build Harper Pro (Node.js v24)

Property '_socket' does not exist on type 'WebSocket'.

Check failure on line 320 in replication/replicationConnection.ts

View workflow job for this annotation

GitHub Actions / Build Harper Pro (Node.js v22)

Property '_socket' does not exist on type 'WebSocket'.

Check failure on line 320 in replication/replicationConnection.ts

View workflow job for this annotation

GitHub Actions / Build Harper Pro (Node.js v26)

Property '_socket' does not exist on type 'WebSocket'.
return true; // tell the watchdog to treat this connection's byte count as frozen
}

/**
* Mark an error as a *source-reported* blob unavailability: the sender told us (via a BLOB_CHUNK
* `error` marker) that it cannot provide this blob — classically `ENOENT` because the blob was
Expand Down Expand Up @@ -587,6 +612,9 @@
// which is the post-close terminal marker. Anything else (protocol errors, peer
// DISCONNECT, etc.) leaves this false so the close handler schedules a retry.
intentionallyUnsubscribed = false;
// Set while a reconnect has already been scheduled (by forceReconnect or the close handler) so the
// two paths never both arm a connect() for the same drop — see forceReconnect / harper-pro#420.
reconnectScheduled = false;
nodeSubscriptions?: NodeSubscription[];
latency = 0;
replicateTablesByDefault: boolean;
Expand All @@ -612,7 +640,22 @@
if (this.intentionallyUnsubscribed) return;
if (!this.session) this.resetSession();
// TODO: Need to do this specifically for each node
this.socket = await createWebSocket(this.url, { serverName: this.nodeName, authorization: this.authorization });
try {
this.socket = await createWebSocket(this.url, {
serverName: this.nodeName,
authorization: this.authorization,
});
} finally {
// A forceReconnect-scheduled reconnect is now realized (or has failed) — stop suppressing the
// close handler's own retry. Clearing this only after this.socket is reassigned (rather than in
// the scheduling timer) keeps a late close from the superseded socket from arming a second
// connect() during the createWebSocket await window. See forceReconnect / harper-pro#420.
this.reconnectScheduled = false;
}
// Capture this attempt's socket so its close handler can tell whether it is still the live socket:
// forceReconnect can schedule a fresh connect() that replaces this.socket before this socket's
// (possibly delayed) terminate() finally fires close. See the close handler / harper-pro#420.
const socket = this.socket;

let session;
logger.debug?.(`Connecting to ${this.url}, db: ${this.databaseName}, process ${process.pid}`);
Expand Down Expand Up @@ -678,6 +721,11 @@
this.sessionReject(error);
});
this.socket.on('close', (code, reasonBuffer) => {
// Ignore a late close from a socket we have already replaced — forceReconnect may have scheduled
// a fresh connect() that swapped in a new this.socket before this superseded socket's terminate()
// finally propagated its close. Acting on it would wrongly tear down the live connection (mark it
// disconnected, drop its subscription listener, reset its session). See harper-pro#420.
if (this.socket !== socket) return;
// Only treat the close as terminal when something explicitly marked it as a deliberate
// teardown (user unsubscribe or the empty-subscription delayed close). Protocol-level
// closes — peer DISCONNECT, unauthorized after open, node-name-mismatch, invalid
Expand Down Expand Up @@ -713,6 +761,9 @@
);
}
session = null;
// forceReconnect (receive-watchdog path) may have already scheduled the reconnect before
// this close fired — don't arm a second connect() for the same drop.
if (this.reconnectScheduled) return;
this.resetSession();
// try to reconnect
setTimeout(() => {
Expand Down Expand Up @@ -743,6 +794,50 @@
this.intentionallyUnsubscribed = true;
this.socket?.close(1008, 'No longer subscribed');
}
// Drive recovery when the receive watchdog detects a silent connection. The normal recovery path is
// the close handler's retry, but an open-but-idle socket (copy stalls, no transport close) may never
// emit 'close' — and even when terminate() does fire one, the wedge reconciler skips a (peer, db)
// whose node entry is still connected:true, and the cached connection is still "reusable" so a
// re-subscribe hands back the same dead socket (replicator.isReusableConnection). So drive recovery
// here, independent of the close event: notify the main thread the node is disconnected (flipping the
// entry to connected:false for the reconciler backstop), tear the socket down, and schedule one fresh
// connect. See harper-pro#420.
forceReconnect() {
if (this.intentionallyUnsubscribed || this.isFinished || this.reconnectScheduled) return;
if (this.isConnected) {
if (this.nodeSubscriptions) {
disconnectedFromNode({
name: this.nodeName,
database: this.databaseName,
url: this.url,
finished: false,
});
}
this.isConnected = false;
}
this.reconnectScheduled = true;
// Drop this connection's stale subscription listener before reconnecting. The close handler
// normally does this (removeAllListeners), but its socket-identity guard early-returns for a
// superseded socket, so an open-but-idle wedge — whose old socket never fires a timely close —
// would otherwise leak one 'subscriptions-updated' listener per recovery cycle (eventually a
// MaxListenersExceededWarning). The fresh connect re-registers its own. See harper-pro#420.
this.removeAllListeners('subscriptions-updated');
this.resetSession();
const socket = this.socket;
// connect() clears reconnectScheduled once the new socket is installed (its finally), which keeps
// a late close from the old socket from double-scheduling during the createWebSocket await.
setTimeout(() => {
this.connect();
}, this.retryTime).unref();
// Match the close-handler backoff so a repeatedly-wedging peer backs off the same way (#339).
this.retryTime = Math.min(this.retryTime << 1, 30_000);
// Best-effort teardown of the dead socket; recovery above does not depend on the close it fires.
try {
socket?.terminate();
} catch {
// already destroyed — the scheduled connect still runs
}
}

getRecord(request) {
return this.session.then((session) => {
Expand Down Expand Up @@ -876,9 +971,10 @@
// sendPing tick above: if that tick is missed or its ws.terminate() does not propagate a
// 'close' event, the watchdog forces the reconnect path. See harper-pro#233 for the failure
// modes observed in the field.
const wedgedForTest = armReplicationWedgeForTest(options.connection, ws, databaseName);
receiveWatchdog = createReceiveWatchdog({
intervalMs: RECEIVE_SILENCE_THRESHOLD_MS,
getBytesRead: () => ws._socket?.bytesRead ?? 0,
getBytesRead: () => (wedgedForTest ? 0 : (ws._socket?.bytesRead ?? 0)),
onSilence: () => {
// Warn-level: if the active sendPing was healthy this watchdog should not have fired,
// so it is a signal that something is wrong upstream (event-loop stall, keepalive timer
Expand All @@ -889,7 +985,12 @@
logger.warn?.(
`Receive watchdog: ${direction} ${remoteNodeName}${dbContext} for ${RECEIVE_SILENCE_THRESHOLD_MS}ms — terminating connection and reconnecting`
);
ws.terminate();
// On the client (subscription) side drive recovery through the connection so it does not depend
// on terminate() propagating a 'close' (an open-but-idle socket may never emit one). A
// server-accepted connection has no connection object to reconnect — the remote client
// reconnects — so just terminate. See harper-pro#420.
if (options.connection) options.connection.forceReconnect();
else ws.terminate();
},
});
const resetPingTimer = receiveWatchdog.reset;
Expand Down
Loading
Loading