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
2 changes: 1 addition & 1 deletion core
Submodule core updated 118 files
128 changes: 128 additions & 0 deletions replication/blobSendDrain.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
/**
* Per-worker tracker of in-flight replication blob *sends*, used to gracefully drain them before an
* http worker shuts down during a restart (deploy reload, rolling restart).
*
* A worker restart that tears down a blob send mid-stream leaves the peer receiver's copy diverged
* until it re-requests the blob (harper-pro#527 handles the receiver side). To avoid interrupting a
* transfer that is about to finish, the worker's shutdown path drains sends that are still making
* progress — bounded by an absolute deadline. `sendBlobs` registers each send, marks progress every
* time bytes are written to the socket, and unregisters on completion; core's shutdown drain
* ({@link ../core/components/shutdownDrain.ts}) polls {@link isDrainComplete} until every remaining
* send has either finished, stalled (no bytes for {@link STALL_MS}), or the deadline has passed.
*
* Reset-on-bytes is the point: a healthy large transfer that keeps flowing is never treated as
* stalled and is allowed to finish (up to the deadline); only a send that goes silent is abandoned.
*
* State is module-local, so it is naturally per-worker (each worker is a fresh module realm).
*/

/** A send that goes this long without writing bytes is considered stalled and no longer drained. */
export const STALL_MS = Math.max(0, Number(process.env.HARPER_BLOB_SEND_DRAIN_STALL_MS) || 0) || 5000;
/** How often the drain re-checks whether the remaining sends have finished or stalled. */
export const POLL_MS = 250;

export interface BlobSendProgress {
lastProgressAt: number;
}

const activeSends = new Set<BlobSendProgress>();

// Set once the worker begins draining for shutdown, so no NEW sends are started while we drain (a
// worker still listening could otherwise keep registering fresh sends and hold the drain open to the
// ceiling). In-flight sends continue; skipped ones are re-requested by the peer on reconnect (fix (1)).
let draining = false;

/** Whether the worker has begun draining sends for shutdown (new sends should not be started). */
export function isDrainingBlobSends(): boolean {
return draining;
}

/** Start tracking a blob send. Returns a handle to mark progress on and to end. */
export function registerBlobSend(): BlobSendProgress {
const entry: BlobSendProgress = { lastProgressAt: Date.now() };
activeSends.add(entry);
return entry;
}

/** Mark that bytes were just written for this send — must stay cheap, it runs per chunk. */
export function noteBlobSendProgress(entry: BlobSendProgress): void {
entry.lastProgressAt = Date.now();
}

/** Stop tracking a finished (or failed) blob send. */
export function endBlobSend(entry: BlobSendProgress): void {
activeSends.delete(entry);
}

/** Whether any blob send is currently in flight. */
export function hasActiveBlobSends(): boolean {
return activeSends.size > 0;
}

/**
* Whether any blob send is currently making progress (wrote bytes within the stall window). This —
* not merely "active" — gates the shutdown-deadline extension: an already-stalled send is worth
* nothing to drain, so it must not push the force-kill backstops out to the ceiling.
*/
export function hasProgressingBlobSends(stallMs: number = STALL_MS): boolean {
const now = Date.now();
for (const entry of activeSends) {
if (isSendProgressing(entry.lastProgressAt, now, stallMs)) return true;
}
return false;
}

/** A send is still worth waiting on if it wrote bytes within the stall window. */
export function isSendProgressing(lastProgressAt: number, now: number, stallMs: number): boolean {
return now - lastProgressAt < stallMs;
}

/**
* The drain is complete when the deadline has passed, or when no remaining send is still progressing
* (all have finished or stalled). An empty set is trivially complete.
*/
export function isDrainComplete(
lastProgressTimes: Iterable<number>,
now: number,
stallMs: number,
deadlineMs: number
): boolean {
if (now >= deadlineMs) return true;
for (const lastProgressAt of lastProgressTimes) {
if (isSendProgressing(lastProgressAt, now, stallMs)) return false;
}
return true;
}

/**
* Resolve once every in-flight blob send has finished, stalled, or the absolute `deadlineMs` (epoch
* timestamp) has passed — whichever comes first. Polls rather than waiting event-driven; this only
* runs on the shutdown path, so the small poll cost is irrelevant and the logic stays trivially
* testable.
*/
export function drainBlobSends(
deadlineMs: number,
stallMs: number = STALL_MS,
pollMs: number = POLL_MS
): Promise<void> {
draining = true; // quiesce: stop starting new sends for the rest of this worker's life
return new Promise<void>((resolve) => {
const check = () => {
const now = Date.now();
const times: number[] = [];
for (const entry of activeSends) times.push(entry.lastProgressAt);
if (isDrainComplete(times, now, stallMs, deadlineMs)) {
resolve();
return;
}
setTimeout(check, Math.min(pollMs, Math.max(0, deadlineMs - now))).unref();
};
check();
});
}

/** Test-only: drop all tracked sends so unit tests start from a clean per-module state. */
export function _resetBlobSendDrainForTest(): void {
activeSends.clear();
draining = false;
}
73 changes: 71 additions & 2 deletions replication/replicationConnection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
import { getThisNodeName } from '../core/server/nodeName.ts';
import * as env from '../core/utility/environment/environmentManager.js';
import { CONFIG_PARAMS } from '../core/utility/hdbTerms.ts';
import { registerBlobSend, noteBlobSendProgress, endBlobSend, isDrainingBlobSends } from './blobSendDrain.ts';
import { HAS_STRUCTURE_UPDATE, lastMetadata, lastValueEncoding, METADATA } from '../core/resources/RecordEncoder.ts';
import { decode, encode, Packr } from 'msgpackr';
import { createStructon } from 'structon';
Expand Down Expand Up @@ -376,6 +377,37 @@
return lastChunk + blobTimeout < now;
}

/**
* Race a backpressure `drain` wait against the connection going away, so a mid-flush peer disconnect
* always lets the waiter settle instead of parking forever on a `drain` that will never fire.
*
* `sendBlobs` awaits this whenever `ws._socket.writableNeedDrain` is true (both the mid-loop wait and
* the terminal-frame flush wait have this exact shape). Without racing `close`/`error`, a peer that
* closes the connection while backpressured leaves the await unsettled, so `sendBlobs`'s `finally`
* never runs: `endBlobSend` is skipped and the drain token leaks in `blobSendDrain`'s module-global
* `activeSends` for the rest of the worker's life (harper-pro#529 review, cb1kenobi). Listening on both
* the raw socket (`drain`/`close`/`error`) and the WebSocket wrapper (`close`) covers a close that
* surfaces on either emitter; all listeners are removed once one fires, so nothing is left registered
* after the promise settles.
*
* Exported so the race itself is unit-testable with plain EventEmitters standing in for the socket/ws.
*/
export function waitForDrainOrSocketEnd(socket: EventEmitter, ws: EventEmitter): Promise<void> {
return new Promise<void>((resolve) => {
const done = () => {
socket.off('drain', done);
socket.off('close', done);
socket.off('error', done);
ws.off('close', done);
resolve();
};
socket.once('drain', done);
socket.once('close', done);
socket.once('error', done);
ws.once('close', done);
});
}

/**
* Credit the back-pressure pause back to every in-flight blob stream at the moment the receiver resumes
* (the point where the pause-reason refcount drops to zero). Because reads are suspended while paused,
Expand Down Expand Up @@ -525,7 +557,7 @@

/**
* Create the PassThrough that receives a blob's bytes on the way to its file store. It carries a no-op
* `'error'` listener from creation so that destroying it with an error — most importantly the blobsTimer

Check failure on line 560 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 560 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 560 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 560 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 560 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 560 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'.
* sweep tearing down a blob whose save never wired up, leaving the stream orphaned in `blobsInFlight`
* after an app/source error (harper-pro#1337) — cannot promote a stream `'error'` into a process-level
* uncaughtException. saveBlob's pipeline still observes and reports real errors via its completion
Expand Down Expand Up @@ -3810,6 +3842,12 @@
return;
}
if (wsClosed) return;
if (isDrainingBlobSends()) {
// The worker is draining for shutdown; don't start a new send that we'd only tear down (or
// hold the drain open for). The peer re-requests this blob on reconnect (harper-pro#527).
logger.debug?.('Worker draining, not starting new blob send', id);
return;
}
blobsBeingSent.add(id);
// Acquire a send slot before opening the blob stream. Enforcing the cap only at the audit
// writer's backpressure check didn't bound concurrency: there it sat in an else-if behind the
Expand All @@ -3822,8 +3860,18 @@
blobsBeingSent.delete(id);
return;
}
if (isDrainingBlobSends()) {
// A shutdown drain can start while this send was queued behind the concurrency cap;
// don't let a freshly-dequeued send start after that point (same reasoning as the
// pre-queue check above).
blobsBeingSent.delete(id);
return;
}
}
const iterator = blob.stream()[Symbol.asyncIterator]();
// Track this send so a worker restart can gracefully drain it (finish it if it's still making
// progress) before shutting down, rather than tearing it down mid-stream. See blobSendDrain.ts.
const drainToken = registerBlobSend();
try {
let lastBuffer: Buffer;
outstandingBlobsBeingSent++;
Expand Down Expand Up @@ -3877,12 +3925,20 @@
);
}
lastBuffer = buffer;
if (ws._socket.writableNeedDrain) {
// Optional-chain the guard: the connection can close during an await, leaving `_socket` null.
if (ws._socket?.writableNeedDrain) {
logger.debug?.('draining', id);
await new Promise((resolve) => ws._socket.once('drain', resolve));
// Waiting on the socket to flush IS progress — mark it so a shutdown drain doesn't misread a
// slow-but-alive peer (a large chunk taking longer than the stall window to flush) as stalled.
noteBlobSendProgress(drainToken);
// Races against close/error so a mid-flush disconnect still lets `finally` below run
// (endBlobSend/outstandingBlobsBeingSent cleanup) instead of hanging on a `drain` that
// will never fire (harper-pro#529 review, cb1kenobi).
await waitForDrainOrSocketEnd(ws._socket, ws);
logger.debug?.('drained', id);
}
recordAction(buffer.length, 'bytes-sent', `${remoteNodeName}.${databaseName}`, 'replication', 'blob');
noteBlobSendProgress(drainToken);
}
logger.debug?.('Sending final blob chunk', id, 'length', lastBuffer.length);
if (checkExcessMessageSize(lastBuffer.length)) throw new Error('Blob chunk too large');
Expand All @@ -3897,6 +3953,18 @@
lastBuffer,
])
);
noteBlobSendProgress(drainToken);
// Keep this send "in flight" until the terminal frame has actually flushed to the socket, so a
// concurrent shutdown drain waits for the `finished:true` frame rather than exiting with it still
// buffered (which would leave the peer's blob diverged until it re-requests). Only waits under
// backpressure; if the peer stops reading, the drain's stall detection abandons it after the
// stall window and the receiver re-requests (harper-pro#527).
if (ws._socket?.writableNeedDrain) {
// Same close/error race as the mid-loop wait above — otherwise a peer that disconnects
// while this terminal-frame flush is parked on backpressure never lets `finally` run.
await waitForDrainOrSocketEnd(ws._socket, ws);
noteBlobSendProgress(drainToken);
}
} catch (error) {
try {
await iterator.return?.();
Expand Down Expand Up @@ -3934,6 +4002,7 @@
])
);
} finally {
endBlobSend(drainToken);
blobsBeingSent.delete(id);
outstandingBlobsBeingSent--;
while (outstandingBlobsBeingSent < MAX_OUTSTANDING_BLOBS_BEING_SENT && blobSentCallbacks.length > 0) {
Expand Down
11 changes: 11 additions & 0 deletions replication/replicator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,8 @@ import {
LATENCY_POSITION,
} from './replicationConnection.ts';
import { redactOperationForLog } from './logRedaction.ts';
import { registerShutdownDrain } from '../core/components/shutdownDrain.ts';
import { hasProgressingBlobSends, drainBlobSends } from './blobSendDrain.ts';
import { server } from '../core/server/Server.ts';
import * as env from '../core/utility/environment/environmentManager.js';
import * as logger from '../core/utility/logging/harper_logger.js';
Expand Down Expand Up @@ -720,6 +722,15 @@ export function forceReconnectToNode({ url, nodes, database }) {
exportIdMapping,
getIdOfRemoteNode,
};

// Gracefully drain in-flight replication blob sends before this worker shuts down during a restart,
// so a deploy reload doesn't tear a transfer down mid-stream and leave the peer's copy diverged
// (harper-pro#527 covers the receiver side). Only sends still making progress are waited on, bounded
// by an absolute deadline core supplies; see blobSendDrain.ts.
registerShutdownDrain({
hasWork: () => hasProgressingBlobSends(),
drain: (deadlineMs: number) => drainBlobSends(deadlineMs),
});
export function urlToNodeName(nodeUrl) {
if (nodeUrl) return new URL(nodeUrl).hostname; // this the part of the URL that is the node name, as we want it to match common name in the certificate
}
Expand Down
Loading
Loading