Skip to content
6 changes: 6 additions & 0 deletions replication/clusterStatus.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@ import {
BLOB_FAILURE_COUNT_POSITION,
LAST_BLOB_FAILURE_TIME_POSITION,
readConnectionTruth,
COPY_SHORTFALL_COUNT_POSITION,
LAST_COPY_SHORTFALL_TIME_POSITION,
} from './replicationConnection.ts';
import '../core/server/serverHelpers/serverUtilities.ts';

Expand Down Expand Up @@ -67,6 +69,10 @@ export async function clusterStatus() {
// `|| undefined` so a healthy link omits the field entirely (matching lastBlobFailure's asDate(0)).
socket.blobReplicationFailures = replicationSharedStatus[BLOB_FAILURE_COUNT_POSITION] || undefined;
socket.lastBlobFailure = asDate(replicationSharedStatus[LAST_BLOB_FAILURE_TIME_POSITION]);
// Copy delivery shortfall tripwire: records the sender reported sent that never arrived intact
// on this link (alert-only verification at COPY_COMPLETE). Omitted entirely when healthy.
socket.copyDeliveryShortfall = replicationSharedStatus[COPY_SHORTFALL_COUNT_POSITION] || undefined;
socket.lastCopyShortfall = asDate(replicationSharedStatus[LAST_COPY_SHORTFALL_TIME_POSITION]);
// W1 (harper-pro#431): the shared-memory connection truth is authoritative over the edge-triggered
// map mirror in requestClusterStatus, which can still read connected:true for an open-but-idle
// wedge that never delivered a disconnect (#289/#233). Also surface the last disconnect (#214).
Expand Down
111 changes: 102 additions & 9 deletions replication/replicationConnection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -139,6 +139,12 @@
export const LAST_LIVENESS_TIME_POSITION = 10; // wall-clock ms of last confirmed liveness (pong or received message)
export const LAST_ERROR_CODE_POSITION = 11; // close code of the most recent disconnect
export const LAST_ERROR_TIME_POSITION = 12; // wall-clock ms of the most recent disconnect
// Copy delivery shortfall (alert-only verification at COPY_COMPLETE): cumulative count of records
// the sender reported sent but this receiver never counted as delivered, and when it last grew.
// Non-zero means silent in-flight loss or undecodable drops occurred on this link; per-table detail
// is in the error log and on the connection's copyDeliveryShortfall marker.
export const COPY_SHORTFALL_COUNT_POSITION = 13;
export const LAST_COPY_SHORTFALL_TIME_POSITION = 14;
export const CONNECTION_STATE_DOWN = 0;
export const CONNECTION_STATE_CONNECTED = 2;
// LIVENESS_STALE_MS is defined below, after PING_TIMEOUT, so it can be derived from the configured
Expand Down Expand Up @@ -246,6 +252,28 @@
return requestedStartTime < Math.max(oldestRetainedTime ?? 0, retentionCutoffTime);
}

/**
* Compare the sender's per-table sent counts (carried on COPY_COMPLETE) against the records this
* receiver actually decoded from the copy stream. Both sides count the same stream for the same
* session, so the comparison is exact: it is immune to concurrent writes, eviction/expiration, and
* resume (a resumed copy counts only its own tail on both ends). A deficit means records were lost
* in flight or dropped undecodable on this end — loss the receive path otherwise silently skips
* past. A malformed sent count is ignored; an extra received table (sender predates the payload,
* or names diverge) is not an error. Pure so the comparison is unit-testable.
*/
export function computeCopyDeliveryShortfall(
sent: Record<string, number>,
received: Record<string, number>
): Array<{ table: string; sent: number; received: number }> {
const shortfalls: Array<{ table: string; sent: number; received: number }> = [];
for (const [table, sentCount] of Object.entries(sent)) {
if (typeof sentCount !== 'number' || !(sentCount >= 0)) continue;
const receivedCount = received[table] ?? 0;
if (receivedCount < sentCount) shortfalls.push({ table, sent: sentCount, received: receivedCount });
}
return shortfalls;
}

export const tableUpdateListeners = new Map();
// This a map of the database name to the subscription object, for the subscriptions from our tables to the replication module
// when we receive messages from other nodes, we then forward them on to as a notification on these subscriptions
Expand Down Expand Up @@ -607,7 +635,7 @@
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 638 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 638 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 638 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 638 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 638 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'.

Check failure on line 638 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 638 in replication/replicationConnection.ts

View workflow job for this annotation

GitHub Actions / YCSB single-node (Node.js v24)

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

Expand Down Expand Up @@ -1466,6 +1494,10 @@
let copyModeOrderVersion; // copy-order version the leader announced in COPY_START; persisted in the cursor (#421)
let copyFromNodeId; // local id of the node we are copying from — the key for the persisted cursor
let copyCompleteReceived = false;
// Per-table count of copy records this session delivered intact (decoded, or intentionally
// dropped by policy). Compared at COPY_COMPLETE against the sender's sent counts; see
// computeCopyDeliveryShortfall.
let copyReceivedCounts: Record<string, number> = {};
// Staged key-based copy resume cursor (#426). The copy cursor (`{currentTable, afterKey, ...}` = "fully
// copied through this key") is KEY-based and, exactly like the sequence watermark, must only be
// PERSISTED once the copied key's blob — and every earlier blob — is durable. We must NOT await blobs in
Expand Down Expand Up @@ -2377,6 +2409,7 @@
// the leader is (re)starting a bulk copy; track a resume cursor for it
inCopyMode = true;
pendingCopyCursor = null; // discard any cursor staged by a prior copy on this connection
copyReceivedCounts = {}; // fresh delivery tally for this copy session
copyBytesSinceFlush = 0; // reset the copy-apply flush gate for this (re)start (harper-pro#480)
lastCopyFlushTime = performance.now();
// The byte watchdog was already (re)armed for THIS frame back in the synchronous
Expand All @@ -2395,7 +2428,7 @@
copyFromNodeId = getIdOfRemoteNode(remoteNodeName, auditStore);
logger.debug?.(connectionId, 'bulk copy starting from', remoteNodeName, new Date(copyModeStartTime));
break;
case COPY_COMPLETE:
case COPY_COMPLETE: {
// Copy signalled complete. Stay in copy mode so batches still committing keep advancing the
// cursor; maybeFinishCopy exits copy mode and clears the cursor once those commits drain.
copyCompleteReceived = true;
Expand All @@ -2404,7 +2437,44 @@
copyProgressWatchdog?.stop();
maybeFinishCopy();
logger.debug?.(connectionId, 'bulk copy complete from', remoteNodeName);
// Alert-only delivery verification: the payload (absent from old senders) carries how
// many records the sender actually sent per table in this copy session, and
// copyReceivedCounts tallied what arrived intact. Both sides count the same stream, so
// any deficit is real loss — records dropped in transit or skipped undecodable on this
// end, which the decode path otherwise silently advances past. Deliberately no
// corrective action: a forced re-copy would loop forever against permanently
// undecodable records, and at-rest completeness belongs to a range-checksum check, not
// a delivery tally.
if (data && typeof data === 'object') {
const shortfalls = computeCopyDeliveryShortfall(data as Record<string, number>, copyReceivedCounts);
if (shortfalls.length > 0) {
logger.error?.(
`Copy delivery verification failed for database ${databaseName} from ${remoteNodeName}: ` +
shortfalls.map((s) => `${s.table} sent=${s.sent} received=${s.received}`).join(', ') +
'. The missing records were lost in transit or dropped as undecodable (see any decode errors above); not forcing a re-copy.'
);
if (options.connection) {
// surfaced for inspection/monitoring alongside the connection's other state
options.connection.copyDeliveryShortfall = {
time: Date.now(),
database: databaseName,
from: remoteNodeName,
shortfalls,
};
}
// cluster_status tripwire (mirrors the blob-failure slots, harper-pro#386)
const sharedStatus = getSharedStatus();
if (sharedStatus) {
sharedStatus[COPY_SHORTFALL_COUNT_POSITION] += shortfalls.reduce(
(total, shortfall) => total + (shortfall.sent - shortfall.received),
0
);
sharedStatus[LAST_COPY_SHORTFALL_TIME_POSITION] = Date.now();
}
}
}
break;
}
case SEQUENCE_ID_UPDATE:
// we need to record the sequence number that the remote node has received
lastSequenceIdReceived = data;
Expand Down Expand Up @@ -2773,13 +2843,15 @@
subscriptionToHdbNodes = subscription;
for await (const event of subscriptionToHdbNodes) {
const node = event.value;
if (!(
node?.replicates === true ||
node?.replicates?.receives ||
node?.replicates?.receivesFrom?.some(
(sub) => sub.source === getThisNodeName() && sub.database === databaseName
if (
!(
node?.replicates === true ||
node?.replicates?.receives ||
node?.replicates?.receivesFrom?.some(
(sub) => sub.source === getThisNodeName() && sub.database === databaseName
)
)
)) {
) {
closed = true;
close(1008, `Unauthorized database subscription to ${databaseName}`);
return;
Expand Down Expand Up @@ -3157,7 +3229,8 @@
});
// find the earliest start time of the subscriptions
let copyResume:
{ copyStartTime: number; currentTable: string; afterKey: any; copyOrder?: number } | undefined;
| { copyStartTime: number; currentTable: string; afterKey: any; copyOrder?: number }
| undefined;
for (const subscription of nodeSubscriptions) {
if (subscription.startTime < currentSequenceId) currentSequenceId = subscription.startTime;
// a follower resuming an interrupted bulk copy sends back where it left off. This keeps the
Expand Down Expand Up @@ -3287,6 +3360,11 @@
// Paces the flush/yield cadence inside the copy loop below (see
// COPY_CHECKPOINT_MAX_INTERVAL_MS). Marked on every in-loop flush/yield.
const copyFlushPacer = createCopyFlushPacer(COPY_CHECKPOINT_MAX_INTERVAL_MS, Date.now());
// Per-table tally of records actually sent in THIS copy session, shipped on
// COPY_COMPLETE so the follower can verify delivery (computeCopyDeliveryShortfall).
// Counting the stream keeps the comparison exact: no table scan, no race with
// concurrent writes or eviction, and a resumed copy counts only its own tail.
const copySentCounts: Record<string, number> = {};
// If resuming, the follower already committed every table before currentTable (records commit
// in stable iteration order), so skip to currentTable and continue after its last committed key.
let reachedResumeTable = !copyResume;
Expand Down Expand Up @@ -3414,6 +3492,7 @@
},
entry.localTime
);
copySentCounts[tableName] = (copySentCounts[tableName] ?? 0) + 1;
logger.debug?.(
'sent record from table',
entry.key,
Expand Down Expand Up @@ -3462,7 +3541,11 @@
);
// The full copy is done — tell the follower to clear its resume cursor and fall back to
// normal audit-log replication from the persisted seqId (which is copyStartTime).
ws.send(encode([COPY_COMPLETE]));
// The per-table sent tally rides along so the follower can verify every record it
// was sent arrived intact (computeCopyDeliveryShortfall); the record frames all
// precede this message on the socket, so the follower's tally is settled by the
// time it compares.
ws.send(encode([COPY_COMPLETE, copySentCounts]));
getSharedStatus()[SENDING_TIME_POSITION] = 0;
currentSequenceId = copyStartTime;
}
Expand Down Expand Up @@ -3580,6 +3663,9 @@
'from',
remoteNodeName
);
// a policy drop still counts as delivered for the copy tally: the record arrived intact
if (messageIsCopyFrame)
copyReceivedCounts[tableDecoder.name] = (copyReceivedCounts[tableDecoder.name] ?? 0) + 1;
decoder.position = start + eventLength;
continue;
}
Expand All @@ -3595,6 +3681,9 @@
'from',
remoteNodeName
);
// a policy drop still counts as delivered for the copy tally: the record arrived intact
if (messageIsCopyFrame && tableDecoder)
copyReceivedCounts[tableDecoder.name] = (copyReceivedCounts[tableDecoder.name] ?? 0) + 1;
decoder.position = start + eventLength;
continue;
}
Expand Down Expand Up @@ -3737,6 +3826,10 @@
// otherwise double-apply). Rows older than copyStartTime are never redelivered, so the
// snapshot is safe and carries no audit. Strict `<` keeps the boundary row audited.
event.isCopyApply = messageIsCopyFrame && copyApplyActive() && auditRecord.version < copyModeStartTime;
// Count only records that decoded successfully: a decode failure leaves `event` unset and
// falls through above, so its absence here is exactly the deficit the COPY_COMPLETE
// delivery verification reports.
if (messageIsCopyFrame) copyReceivedCounts[event.table] = (copyReceivedCounts[event.table] ?? 0) + 1;
tableSubscriptionToReplicator.send(event);
// Per-record backpressure: a single large WS message can synchronously decode
// thousands of records, each holding a decoded value object and a closure over
Expand Down
60 changes: 60 additions & 0 deletions unitTests/replication/copyDeliveryVerification.test.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
/**
* Coverage for computeCopyDeliveryShortfall — the COPY_COMPLETE delivery check (#537).
*
* The sender tallies records it actually sent per table during a copy session and ships the tally
* on COPY_COMPLETE; the receiver tallies records that arrived intact (decoded, or intentionally
* dropped by policy). Both sides count the same stream, so the comparison is exact and immune to
* concurrent writes, eviction, and resume — unlike a point-in-time table count, which races live
* writes and reads an estimate. A deficit means records were lost in transit or dropped
* undecodable on the receiver (the decode path logs and skips those, letting the cursor advance
* past them). Alert-only by design: a forced re-copy would loop forever against permanently
* undecodable records.
*/

import { expect } from 'chai';
import { computeCopyDeliveryShortfall } from '#src/replication/replicationConnection';

describe('computeCopyDeliveryShortfall', () => {
it('reports nothing when every table matches', () => {
expect(computeCopyDeliveryShortfall({ dogs: 100, cats: 0 }, { dogs: 100, cats: 0 })).to.deep.equal([]);
});

it('reports nothing for an empty sent tally (old sender or empty copy)', () => {
expect(computeCopyDeliveryShortfall({}, { dogs: 5 })).to.deep.equal([]);
});

it('reports a deficit with both counts', () => {
expect(computeCopyDeliveryShortfall({ dogs: 100 }, { dogs: 97 })).to.deep.equal([
{ table: 'dogs', sent: 100, received: 97 },
]);
});

it('treats a table absent from the received tally as zero delivered', () => {
expect(computeCopyDeliveryShortfall({ dogs: 3 }, {})).to.deep.equal([{ table: 'dogs', sent: 3, received: 0 }]);
});

it('does not report a table the sender sent nothing for', () => {
expect(computeCopyDeliveryShortfall({ dogs: 0 }, {})).to.deep.equal([]);
});

it('ignores tables only the receiver saw (renames, older senders)', () => {
expect(computeCopyDeliveryShortfall({ dogs: 2 }, { dogs: 2, cats: 9 })).to.deep.equal([]);
});

it('never reports a surplus as a shortfall', () => {
expect(computeCopyDeliveryShortfall({ dogs: 2 }, { dogs: 5 })).to.deep.equal([]);
});

it('ignores malformed sent counts', () => {
expect(computeCopyDeliveryShortfall({ dogs: -1, cats: NaN, fish: '7', birds: 2 }, { birds: 1 })).to.deep.equal([
{ table: 'birds', sent: 2, received: 1 },
]);
});

it('reports multiple shortfalls in sent-tally order', () => {
expect(computeCopyDeliveryShortfall({ a: 10, b: 5, c: 1 }, { a: 10, b: 0, c: 0 })).to.deep.equal([
{ table: 'b', sent: 5, received: 0 },
{ table: 'c', sent: 1, received: 0 },
]);
});
});
Loading