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
137 changes: 118 additions & 19 deletions cloneNode/cloneNode.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { parseArgs } from 'node:util';
import { accessSync, readFileSync, writeFileSync, mkdirSync } from 'node:fs';
import { accessSync, readFileSync, writeFileSync, mkdirSync, renameSync, unlinkSync } from 'node:fs';
import { join, dirname } from 'node:path';
import { homedir } from 'node:os';
import { randomBytes } from 'node:crypto';
Expand Down Expand Up @@ -30,6 +30,10 @@ import {
} from '../core/utility/hdbTerms.ts';
import { fetchJWTKeyWithRetry } from './jwtKeyClone.ts';
import { monitorSyncLoop } from './syncMonitor.ts';
import {
isExplicitDatabaseSubscription,
isReplicatedDatabase as isReplicatedDatabaseUnder,
} from '../replication/replicatedDatabases.ts';

/**
* Environment Variables:
Expand Down Expand Up @@ -83,6 +87,8 @@ import { monitorSyncLoop } from './syncMonitor.ts';
const DEFAULT_SYNC_TIMEOUT_MS = 300000;
const DEFAULT_SYNC_CHECK_INTERVAL_MS = 3000;
const DEFAULT_REPLICATION_PORT = '9933';
const CLONE_ATTEMPT_FILE = '.cloneAttempt.json';
const CLONE_ATTEMPT_ENV = 'HARPER_CLONE_ATTEMPT';

const CONFIG_TO_EXCLUDE_FROM_CLONE = {
clustering_nodename: true,
Expand Down Expand Up @@ -236,6 +242,7 @@ export async function cloneNode(): Promise<void> {
const { main } = await import('../core/bin/run.js');
return main();
}
startCloneAttempt();

if (!usingCertAuth) {
// Request to leader to verify connectivity and credentials before proceeding with clone
Expand Down Expand Up @@ -308,6 +315,7 @@ export async function cloneNode(): Promise<void> {
operation: string;
verify_tls: boolean;
url: string;
isLeader: true;
authorization?:
| {
username: string;
Expand All @@ -320,6 +328,7 @@ export async function cloneNode(): Promise<void> {
operation: OPERATIONS_ENUM.ADD_NODE,
verify_tls: false, // set node cross-signs the cluster with harper self-signed certs
url: leaderReplicationURL,
isLeader: true,
};

if (!usingCertAuth) {
Expand Down Expand Up @@ -432,6 +441,7 @@ export async function cloneNode(): Promise<void> {

// Set a config value to indicate that this node has been cloned, which can be used by other processes to check clone status and prevent duplicate cloning
updateConfigValue(CONFIG_PARAMS.CLONED, true);
clearCloneAttempt();

log(`Clone from leader node ${leaderURL} complete`);
}
Expand Down Expand Up @@ -512,13 +522,47 @@ async function monitorSync(): Promise<SyncOutcome> {
`Starting to monitor sync status. Will check every ${DEFAULT_SYNC_CHECK_INTERVAL_MS}ms and fail if no replication data arrives for ${Math.round(stallTimeoutMs / 1000)}s`
);

// Whether the system database's socket is required has two independent gates.
//
// Local: this node only ever opens a system socket if its own `replication.databases` covers
// `system` — `shouldReplicateFromNode` runs every database, system included, through that
// filter. A node configured with e.g. `databases: ['data']` never subscribes to system, so
// requiring that socket would wedge the clone Unavailable forever.
//
// Leader capability: a legacy (v4) leader never replicates the system database either, while a
// v5+ leader must have it required up front — otherwise a small user database completing before
// the system subscription registers could finish the clone with the system copy unverified.
// registration_info is the version probe present on every leader version (see
// core/bin/cliOperations.ts). Fail CLOSED on that probe: only a positively-read legacy major
// version exempts system; a missing/unparseable version or a persistently failing probe requires
// it, so a transient probe error against a v5 leader cannot reopen the premature-Available race.
let systemSocketRequired = isReplicatedDatabase(SYSTEM_SCHEMA_NAME);
Comment thread
kriszyp marked this conversation as resolved.
if (!systemSocketRequired) {
log(`'${SYSTEM_SCHEMA_NAME}' is not in this node's replication.databases; not requiring its socket`, 'debug');
}
for (let attempt = 1; systemSocketRequired && attempt <= 3; attempt++) {
try {
const registration: any = await leaderRequest({ operation: 'registration_info' });
// First digit run tolerates prefixed version strings (e.g. "v4.3.7"), which parseInt would NaN.
const leaderMajorVersion = Number(String(registration?.version ?? '').match(/\d+/)?.[0] ?? NaN);
systemSocketRequired = !(leaderMajorVersion >= 1 && leaderMajorVersion < 5);
break;
} catch (err) {
log(`Leader version probe failed (attempt ${attempt}/3): ${err}`);
if (attempt < 3) await sleep(1000);
}
}

const outcome = await monitorSyncLoop({
targetTimestamps,
clusterStatus,
leaderReplicationURL,
stallTimeoutMs,
checkIntervalMs: DEFAULT_SYNC_CHECK_INTERVAL_MS,
log,
requiredSocketDatabases: Object.keys(targetTimestamps).filter(
(database) => database !== 'system' || systemSocketRequired
Comment thread
kriszyp marked this conversation as resolved.
),
});

if (outcome === 'synced') {
Expand Down Expand Up @@ -549,16 +593,49 @@ async function monitorSync(): Promise<SyncOutcome> {
* and record the most recent timestamp for each database in a JSON file.
* @returns {Promise<void>}
*/
// A database the clone doesn't subscribe to must not be pre-created (cloneSchemas) or become a
// sync target (getLastUpdatedRecord: its socket never exists, so a target would wedge the sync
// monitor). A sharded entry replicates only from a same-shard leader (`shouldReplicateFromNode`);
// the leader's shard comes from its configuration. Fail closed on an unreadable configuration by
// treating sharded entries as replicated: a wrong inclusion stalls the clone visibly, a wrong
// exclusion would skip verifying a database that is being copied.
function isReplicatedDatabase(dbName: string, shardedReplicates?: (entry: any) => boolean): boolean {
return isReplicatedDatabaseUnder(envMgr.get(CONFIG_PARAMS.REPLICATION_DATABASES), dbName, shardedReplicates);
}

async function leaderShardedReplicates(): Promise<(entry: any) => boolean> {
try {
const leaderConfiguration: any = await leaderRequest({ operation: 'get_configuration' });
const leaderShard = leaderConfiguration?.replication?.shard;
const localShard = envMgr.get(CONFIG_PARAMS.REPLICATION_SHARD);
return () => leaderShard === localShard;
} catch (err) {
log(`Could not read the leader configuration for shard matching (${err}); keeping sharded sync targets`);
return () => true;
}
}

async function getLastUpdatedRecord(): Promise<Record<string, number>> {
log('Getting last updated record timestamp for all database', 'debug');
const lastUpdated: Record<string, number> = {};
const systemDb: Record<string, any> = await leaderRequest({ operation: 'describe_database', database: 'system' });
lastUpdated['system'] = findMostRecentTimestamp(systemDb);

const shardedReplicates = await leaderShardedReplicates();
const { getHDBNodeTable } = await import('../replication/knownNodes.ts');
let leaderNode: any;
for (const node of getHDBNodeTable().search([])) {
if (node?.isLeader || node?.url === leaderReplicationURL) {
leaderNode = node;
break;
}
}
const allDb: Record<string, any> = await leaderRequest({ operation: 'describe_all' });
for (const db in allDb) {
// requestId is part of the describe response so we ignore it
if (typeof allDb[db] !== 'object') continue;
if (!isReplicatedDatabase(db, shardedReplicates) && !isExplicitDatabaseSubscription(leaderNode?.subscriptions, db))
continue;
lastUpdated[db] = findMostRecentTimestamp(allDb[db]);
}

Expand Down Expand Up @@ -912,27 +989,10 @@ async function cloneSchemas(): Promise<void> {
const { createSchema, createTable } = await import('../core/dataLayer/schema.js');
const { databases } = await import('../core/resources/databases.js');

// Filter by this node's `replication.databases` so we don't materialize empty databases the
// clone isn't even subscribing to. Matches the gating used by `shouldReplicateFromNode` in
// `replication/knownNodes.ts`: `undefined` or `'*'` accept everything; an array accepts only
// the names it lists (objects with `.name` are sharded-database entries).
const databaseReplications = envMgr.get(CONFIG_PARAMS.REPLICATION_DATABASES);
const isReplicatedDatabase = (dbName: string): boolean => {
if (!databaseReplications || databaseReplications === '*') return true;
if (!Array.isArray(databaseReplications)) return true;
return databaseReplications.some((entry: any) =>
typeof entry === 'string' ? entry === dbName : entry?.name === dbName
);
};

for (const dbName of Object.keys(allDb)) {
const dbDescribe = allDb[dbName];
if (!dbDescribe || typeof dbDescribe !== 'object' || dbName === SYSTEM_SCHEMA_NAME) continue;
if (!isReplicatedDatabase(dbName)) {
log(`Skipping schema pre-create for '${dbName}' (not in replication.databases)`, 'debug');
continue;
}

if (!isReplicatedDatabase(dbName)) continue;
if (!databases[dbName]) {
try {
await createSchema({ database: dbName, operation: OPERATIONS_ENUM.CREATE_DATABASE });
Expand Down Expand Up @@ -1139,6 +1199,45 @@ function writeJsonSync(path: string, data: any): void {
}
}

function cloneAttemptPath(): string {
return join(rootPath, CLONE_ATTEMPT_FILE);
}

function startCloneAttempt(): void {
const path = cloneAttemptPath();
let attemptId: string | undefined;
if (pathExists(path)) {
try {
const persisted = JSON.parse(readFileSync(path, 'utf8'));
if (typeof persisted?.attemptId === 'string') attemptId = persisted.attemptId;
} catch (error) {
log(`Could not read persisted clone attempt at ${path}: ${error}`, 'error');
}
}
if (!attemptId) {
attemptId = randomBytes(16).toString('hex');
try {
const temporaryPath = `${path}.${process.pid}.tmp`;
mkdirSync(dirname(path), { recursive: true });
writeFileSync(temporaryPath, JSON.stringify({ attemptId }), { encoding: 'utf8', mode: 0o600 });
renameSync(temporaryPath, path);
} catch (error) {
log(`Could not persist clone attempt at ${path}: ${error}`, 'error');
return;
}
}
process.env[CLONE_ATTEMPT_ENV] = attemptId;
}

function clearCloneAttempt(): void {
delete process.env[CLONE_ATTEMPT_ENV];
try {
unlinkSync(cloneAttemptPath());
} catch (error: any) {
if (error?.code !== 'ENOENT') log(`Could not remove clone attempt marker: ${error}`, 'error');
}
}

/**
* Get the database path for a given database name
* Checks to see if there is any custom DB pathing else uses the default storage path
Expand Down
62 changes: 47 additions & 15 deletions cloneNode/syncMonitor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,13 @@ export type SyncCheckResult = {
syncComplete: boolean;
/** Most recent arrival stamp (ms epoch) among databases still below their target; 0 if none. */
latestReceivedMs: number;
/** Databases that had a replication socket to the leader in this check. */
socketDatabases: Set<string>;
};

/**
* Check whether every database with a target timestamp has caught up, and report the freshest
* arrival stamp among the databases that have not.
* Check whether every database has caught up with the leader, and report the freshest arrival
* stamp among the databases that have not.
*
* Completion and liveness are deliberately separate signals: `lastReceivedVersion` is frozen for the
* whole bulk copy (it jumps to copyStartTime only on the post-copy end_txn), so it can only answer
Expand All @@ -21,42 +23,47 @@ export async function checkSyncStatus(
targetTimestamps: Record<string, number>,
clusterStatus: () => Promise<any>,
leaderReplicationURL: string,
log: SyncMonitorLog
log: SyncMonitorLog,
requiredSocketDatabases: string[] = Object.keys(targetTimestamps)
): Promise<SyncCheckResult> {
const clusterResponse = await clusterStatus();
log(`clone sync check cluster status response: ${JSON.stringify(clusterResponse)}`, 'debug');

if (!clusterResponse) {
log('No cluster status response received for clone, will wait and retry');
return { syncComplete: false, latestReceivedMs: 0 };
return { syncComplete: false, latestReceivedMs: 0, socketDatabases: new Set() };
}

if (!clusterResponse.connections?.length) {
log('No connections found in cluster status response for clone, will wait and retry');
return { syncComplete: false, latestReceivedMs: 0 };
return { syncComplete: false, latestReceivedMs: 0, socketDatabases: new Set() };
}

const leaderConnection = clusterResponse.connections.find((conn) => conn.url === leaderReplicationURL);

if (!leaderConnection) {
log('No connection found matching leader replication URL, will wait and retry');
return { syncComplete: false, latestReceivedMs: 0 };
return { syncComplete: false, latestReceivedMs: 0, socketDatabases: new Set() };
}

if (!leaderConnection.database_sockets?.length) {
log(`No database sockets found for connection leader ${leaderConnection.name}`, 'debug');
return { syncComplete: false, latestReceivedMs: 0 };
return { syncComplete: false, latestReceivedMs: 0, socketDatabases: new Set() };
}

let syncComplete = true;
let latestReceivedMs = 0;
const socketDatabases = new Set<string>();
for (const socket of leaderConnection.database_sockets) {
const dbName = socket.database;
const targetTime = targetTimestamps[dbName];
if (!targetTime) {
log(`Database ${dbName}: No target timestamp, skipping sync check`, 'debug');
continue;
}
socketDatabases.add(dbName);
// A missing target — an empty database, or a leader whose describe cannot report
// last_updated_record (RocksDB, harper#2091) — must not skip verification, or the check
// passes vacuously when every target is absent (#655). The received-version watermark is
// held at 0 for the whole bulk copy and only becomes positive via the final end_txn the
// sender emits at copyStartTime, so a positive watermark is the copy's own completion
// signal, independent of the leader's describe support.
const targetTime = targetTimestamps[dbName] || 1;
Comment thread
kriszyp marked this conversation as resolved.

// Raw version (high-precision float64) preserves the sub-millisecond precision needed for
// an accurate comparison against the leader's last_updated_record targets.
Expand All @@ -83,7 +90,20 @@ export async function checkSyncStatus(
if (Number.isFinite(receivedAt) && receivedAt > latestReceivedMs) latestReceivedMs = receivedAt;
}

return { syncComplete, latestReceivedMs };
// A required database with no socket yet (its subscription is still registering with the
// main thread) is pending, not verified — otherwise a lone early socket (e.g. the system DB,
// whose small copy finishes in seconds) could complete the check before the data databases'
// sockets even appear. Only databases the clone actually subscribes to are required: a legacy
// (v4) leader never replicates the system database, so demanding its socket would wedge the
// clone; when the socket does exist it is still verified by the loop above.
for (const dbName of requiredSocketDatabases) {
if (!socketDatabases.has(dbName)) {
log(`Database ${dbName}: no replication socket to the leader yet`, 'debug');
syncComplete = false;
}
}

return { syncComplete, latestReceivedMs, socketDatabases };
}

export type MonitorSyncLoopOptions = {
Expand All @@ -93,6 +113,8 @@ export type MonitorSyncLoopOptions = {
stallTimeoutMs: number;
checkIntervalMs: number;
log: SyncMonitorLog;
/** Databases whose replication socket must exist before sync can complete (default: every target). */
requiredSocketDatabases?: string[];
/** Test hooks: injectable clock and delay. */
now?: () => number;
delay?: (ms: number) => Promise<unknown>;
Expand All @@ -109,6 +131,12 @@ export async function monitorSyncLoop(options: MonitorSyncLoopOptions): Promise<
const delay = options.delay ?? sleep;
let lastProgressAt = now();
let loopCount = 0;
const baseRequired = options.requiredSocketDatabases ?? Object.keys(options.targetTimestamps);
// Ratchet: a target database whose socket has been seen once stays required even if the socket
// later drops, and a non-required one (e.g. `system`, optional because v4 leaders never
// replicate it) becomes required as soon as its socket appears — so on a v5 leader a small user
// database finishing first cannot complete the clone while the system copy is still pending.
const seenTargetSockets = new Set<string>();

while (now() - lastProgressAt < options.stallTimeoutMs) {
try {
Expand All @@ -121,7 +149,8 @@ export async function monitorSyncLoop(options: MonitorSyncLoopOptions): Promise<
options.targetTimestamps,
options.clusterStatus,
options.leaderReplicationURL,
options.log
options.log,
[...new Set([...baseRequired, ...seenTargetSockets])]
);
checkPromise.catch(() => {});
const result = await Promise.race([
Expand All @@ -133,7 +162,10 @@ export async function monitorSyncLoop(options: MonitorSyncLoopOptions): Promise<
await delay(options.checkIntervalMs);
continue;
}
const { syncComplete, latestReceivedMs } = result;
const { syncComplete, latestReceivedMs, socketDatabases } = result;
for (const dbName of socketDatabases) {
if (dbName in options.targetTimestamps) seenTargetSockets.add(dbName);
}

if (syncComplete) return 'synced';

Expand Down
Loading
Loading