Skip to content
22 changes: 10 additions & 12 deletions integrationTests/cluster/auditReplayYieldExcludedTables.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -56,10 +56,10 @@
* load-bearing oracle.
*
* The load-bearing oracle is qa522's cursor-progress sample: `skipAuditRecord()` (the
* fixed path) sends a sequence-position update via its trailing `return new
* Promise(setImmediate)` yield while the reverted path (`logger.debug?.(...)`) sends
* none until the whole synchronous walk finally completes. Confirmed against this exact
* mechanism in this environment: reverting #536 and re-running (3 cold reruns)
* fixed path) sends a sequence-position update within the probe window even on a long
* skip-run; the reverted path (`logger.debug?.(...)`) sends
* none until the whole synchronous walk finally completes. Before time-budget pacing,
* reverting #536 and re-running in this environment (3 cold reruns)
* consistently reproduces qa522's negative-control failure (`DEFECT-SHAPE: B's
* lastReceivedVersion only showed 1 distinct value(s)`); re-applying the fix consistently
* passes (3 cold reruns, >=2 distinct, monotonic).
Expand Down Expand Up @@ -105,10 +105,8 @@ const GHOST_COUNT = 70000; // larger than qa522's single 50k run
// near-baseline (sub-second) if the fix is doing its job; anything crossing this means
// the loop is not yielding like it should. Not the load-bearing oracle -- see file header.
const PING_STALL_THRESHOLD_MS = 4000;
// Fixed window for the ping/cursor probes, matching qa522's design. The whole run
// completes in a couple of seconds on a local machine regardless of fix state at this
// scale, so the window just needs to comfortably span that to catch qa522's proven
// pass/fail signal (see header).
// Fixed window for the ping/cursor probes, matching qa522's design; it includes receiver
// setup as well as audit replay. The header records the original pass/fail evidence.
const PROBE_WINDOW_MS = 8000;
const CURSOR_SAMPLE_INTERVAL_MS = 150;

Expand Down Expand Up @@ -370,15 +368,15 @@ suite(
cursorSamples.map((s) => `${s.t}:${s.version}`).join(' ')
);

// --- Assertion 1 (load-bearing): B's resume cursor advances during the run, not just
// once at the end. This is the signal that actually distinguishes fixed vs
// --- Assertion 1 (load-bearing): B's resume cursor advances within the probe window,
// not just once at the end. This is the signal that actually distinguishes fixed vs
// pre-#536 code -- see the file header for confirmed pass/fail evidence. ---
const distinctVersions = new Set(cursorSamples.map((s) => s.version));
ok(
distinctVersions.size >= 2,
`DEFECT-SHAPE: B's lastReceivedVersion only showed ${distinctVersions.size} distinct value(s) ` +
`across ${cursorSamples.length} samples during the skip-run — periodic sequence updates from ` +
`skipAuditRecord() are not firing during the run, so a reconnect mid-run would rescan from the start`
`across ${cursorSamples.length} samples in the probe window — a periodic sequence update from ` +
`skipAuditRecord() did not land, so a reconnect mid-run would rescan from the start`
);
const versionsInOrder = cursorSamples.map((s) => s.version);
let monotonic = true;
Expand Down
4 changes: 3 additions & 1 deletion replication/DESIGN.md
Original file line number Diff line number Diff line change
Expand Up @@ -334,7 +334,7 @@ discipline paces _retries_, and jittering a durability deadline would be activel

1. **Auth failures don't send DISCONNECT.** When the `authorization` promise rejects in `replicateOverWS`, the connection closes with "Unauthorized" but no DISCONNECT frame is sent — the client is expected to retry.

2. **A skipped record still advances the peer's cursor.** `sendAuditRecord` skips a record the subscriber must not receive (`LOCAL_ONLY`, a gated lock control entry, an unsubscribed or excluded table, an origin the subscription does not cover, a residency omission). Every skip branch returns `skipAuditRecord()`, which yields (#536) and, at most once per `SKIPPED_MESSAGE_SEQUENCE_UPDATE_DELAY` (300 ms), sends a `SEQUENCE_ID_UPDATE` when the last sent sequence id has fallen behind, so a long skipped run still moves the peer's resume cursor.
2. **A skipped record still advances the peer's cursor.** `sendAuditRecord` skips a record the subscriber must not receive (`LOCAL_ONLY`, a gated lock control entry, an unsubscribed or excluded table, an origin the subscription does not cover, a residency omission). Every skip branch returns `skipAuditRecord()`, which consults the shared send budget (note 28) and, at most once per `SKIPPED_MESSAGE_SEQUENCE_UPDATE_DELAY` (300 ms), sends a `SEQUENCE_ID_UPDATE` when the last sent sequence id has fallen behind, so a long skipped run still moves the peer's resume cursor.

3. **Blob back-pressure, transfer identity & timeout.** Blobs time out after `blobTimeout` (default 900s); concurrent sends are capped at `MAX_OUTSTANDING_BLOBS_BEING_SENT` (`replication.blobConcurrency`, default 5); back-pressure ratio (computed every `BACK_PRESSURE_INTERVAL`) tells senders to pause. Copy-mode `BLOB_CHUNK` headers carry an optional `transferId` so adjacent records that reuse one file id do not share skip/apply receive state. A current receiver keys those streams by `transferId`; an older sender omits it and the receiver falls back to `fileId`. Distinct records intentionally get distinct transfers, even when they reference the same file, because a skipped record must not consume an applied record's only transfer; those sends are serialized by file id for older receivers, and receive-in-flight markers are refcounted by file id. The identity-tie fast path proves only uncompressed blobs whose header and file length agree; compressed blobs fail closed because their completeness requires core's writer lock. Transfer ids come from a module-lifetime counter rather than resetting per connection, and their enumerable blob tag exists only inside the synchronous record encode's `try`/`finally`; it must never remain on the process-cached stored record. If you're seeing large-data replication hangs, look here first.

Expand Down Expand Up @@ -492,6 +492,8 @@ Regressions: `unitTests/replication/replicateFalseSendPaths.test.mjs` (frames),

---

28. **Audit sender fairness is shared per worker.** `yieldSendLoop` shares one pending macrotask and restarts its monotonic budget on resume: 2 ms on dedicated replication workers, 0.5 ms on HTTP workers and the main-thread fallback, selected by this thread's `workerData.name`, pinned by `sendLoopYield.test.mjs`. Normal sends (including base-copy records) consult it when neither drain nor blob saturation requires a wait; skips consult it so their sequence-update timer can fire. The copy-flush pacer paces on its own separate cadence (`copyFlushPacer.due`), including copy-only skipped rows.

## Tests

**Integration tests** live in `../integrationTests/cluster/`:
Expand Down
42 changes: 26 additions & 16 deletions replication/replicationConnection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,7 @@ import { verifyLegacyCopyBaseline } from './legacyCopy.ts';
import { getThisNodeName } from '../core/server/nodeName.ts';
import { isDedicatedPoolWorker } from '../core/server/threads/workerPools.ts';
import * as env from '../core/utility/environment/environmentManager.js';
import { CONFIG_PARAMS } from '../core/utility/hdbTerms.ts';
import { CONFIG_PARAMS, THREAD_TYPES } from '../core/utility/hdbTerms.ts';
import { registerBlobSend, noteBlobSendProgress, endBlobSend, isDrainingBlobSends } from './blobSendDrain.ts';
import {
PARK_WARN_MS,
Expand Down Expand Up @@ -105,7 +105,7 @@ import {
import { decode, encode, Packr } from 'msgpackr';
import { createStructon } from 'structon';
import { WebSocket } from 'ws';
import { threadId } from 'worker_threads';
import { threadId, workerData } from 'node:worker_threads';
import harperLogger from '../core/utility/logging/harper_logger.js';
const { forComponent, errorToString } = harperLogger;
import { disconnectedFromNode, connectedToNode, ensureNode } from './subscriptionManager.ts';
Expand Down Expand Up @@ -471,6 +471,25 @@ export function holdFailedFrame(
// (MAX_EVENT_DELAY_TIME = 3 s). Yield the event loop at least this often (ms) while decoding so the
// worker stays responsive during a bulk copy/clone.
const RECEIVE_YIELD_INTERVAL = env.get('replication_receiveYieldInterval') ?? 100;

// One pending turn is shared across senders so the yield cost does not multiply with peer count.
const SEND_YIELD_INTERVAL = workerData?.name === THREAD_TYPES.REPLICATION ? 2 : 0.5;
let lastSendYieldTime = 0;
let pendingSendYield: Promise<void> | undefined;

export function yieldSendLoop(): Promise<void> | undefined {
if (pendingSendYield) return pendingSendYield;
if (performance.now() - lastSendYieldTime >= SEND_YIELD_INTERVAL) {
return (pendingSendYield = new Promise<void>((resolve) => {
setImmediate(() => {
lastSendYieldTime = performance.now();
pendingSendYield = undefined;
resolve();
});
}));
}
}

// A queued frame is usually a subarray of the socket read chunk it arrived in, so it pins the whole
// chunk, not its own length. Both the frame ceiling and the honest retention bound are stated in these
// units rather than in frame lengths.
Expand Down Expand Up @@ -3097,7 +3116,7 @@ export function createReceiveWatchdog(opts: {
* Wall-clock pacer for the bulk-copy send loop. The copy normally flushes to the socket on a
* record-count checkpoint, but reading a large cold table dominates copy cost, so a single
* count-batch can exceed the receive watchdog window with no bytes on the wire — and the LOCAL_ONLY
* skip path bypasses the per-record flush+yield entirely. Either starves the watchdog into killing
* skip path bypasses normal record sending and its budgeted yield. Either starves the watchdog into killing
* the connection mid-copy. This bounds the wall-clock gap between flushes/yields: `due(now)` reports
* whether at least `intervalMs` has elapsed since the last one, and callers `mark(now)` after each
* flush or yield (whether triggered by this pacer or the count checkpoint) so the window restarts.
Expand Down Expand Up @@ -6112,13 +6131,6 @@ export function replicateOverWS(ws: ReplicationWebSocket, options: any, authoriz
if (!tableEntry) {
tableEntry = tableById[tableId] = tableToTableEntry(tableSubscriptionToReplicator.tableById[tableId]);
if (!tableEntry) {
// Must yield like every other skip path: a contiguous run of entries for a
// table this peer doesn't subscribe to (or a dropped table, or corrupt-entry
// sentinels with tableId undefined) otherwise iterates with await undefined,
// which never leaves the microtask queue. Timers, I/O, and watchdogs starve
// for the whole run, and the periodic sequence updates skipAuditRecord sends
// never go out, so the peer's cursor can't advance past the run and every
// reconnect rescans it from the start.
logger.debug?.('Not subscribed to table', tableId);
return skipAuditRecord();
}
Expand Down Expand Up @@ -6242,9 +6254,7 @@ export function replicateOverWS(ws: ReplicationWebSocket, options: any, authoriz
// entry is encoded, send it after checks for new structure and residency
}

// when we can skip an audit record, we still need to occasionally send a sequence update:
// every skip branch in sendAuditRecord must return this call — its contract is the
// trailing yield below, not the logging (see the !tableEntry skip's rationale above, #536).
// Every skip branch must share the yield budget so the sequence-update timer can fire.
function skipAuditRecord() {
logger.trace?.(connectionId, 'skipping audit record', auditRecord.recordId);
if (!skippedMessageSequenceUpdateTimer) {
Expand All @@ -6258,7 +6268,7 @@ export function replicateOverWS(ws: ReplicationWebSocket, options: any, authoriz
}
}, SKIPPED_MESSAGE_SEQUENCE_UPDATE_DELAY).unref();
}
return new Promise(setImmediate); // we still need to yield (otherwise we might never send a sequence id update)
return yieldSendLoop();
}
if (!sentNodeIds.has(auditRecord.nodeId)) {
sentNodeIds.add(auditRecord.nodeId);
Expand Down Expand Up @@ -6384,7 +6394,7 @@ export function replicateOverWS(ws: ReplicationWebSocket, options: any, authoriz
return new Promise((resolve) => {
blobSentCallbacks.push(resolve);
});
} else return new Promise(setImmediate); // yield on each turn for fairness and letting other things run
} else return yieldSendLoop();
};
const sendQueuedData = () => {
if (frame.position - frame.encodingStart > 8) {
Expand Down Expand Up @@ -6901,7 +6911,7 @@ export function replicateOverWS(ws: ReplicationWebSocket, options: any, authoriz
// Bound the wall-clock gap between socket flushes and event-loop yields,
// independent of record count. The count checkpoint below alone can let a cold
// batch run past the watchdog window with no bytes flushed (reads dominate cost),
// and the LOCAL_ONLY `continue` below skips the normal per-record flush+yield
// and the LOCAL_ONLY `continue` below skips normal record sending and its budgeted yield
// entirely — a contiguous skipped run would then never reach the timers phase, so
// the ping timer and receive side starve. Flush any pending batch (plain flush, NOT
// an end_txn — see the watermark note below) and yield a macrotask on this cadence
Expand Down
Loading
Loading