feat(sync): add sync-hosts store, NDJSON --json mode, and --check to subminer sync

This commit is contained in:
2026-07-11 17:45:27 -07:00
parent 8797719a09
commit 0a3f76c0a8
16 changed files with 979 additions and 27 deletions
+158 -11
View File
@@ -18,8 +18,11 @@ import {
} from '../sync/ssh.js';
import { resolvePathMaybe } from '../util.js';
import { canConnectUnixSocket } from '../mpv.js';
import { recordHostSyncResultToDisk } from '../sync/sync-hosts.js';
import type { LauncherCommandContext } from './context.js';
import type { RemoteRunResult } from '../sync/ssh.js';
import type { SyncMergeSummary, SyncProgressEvent } from '../../src/shared/sync/sync-events.js';
import type { SyncResultStatus } from '../../src/shared/sync/sync-hosts-store.js';
export interface SyncCommandDeps {
createDbSnapshot: typeof createDbSnapshot;
@@ -39,6 +42,8 @@ export interface SyncCommandDeps {
consoleLog: typeof console.log;
writeStdout: typeof process.stdout.write;
ensureTrackerQuiescent: (context: LauncherCommandContext, dbPath: string) => Promise<void>;
emitEvent: (event: SyncProgressEvent) => void;
recordHostSyncResult: (host: string, status: SyncResultStatus, detail: string | null) => void;
}
function resolveDbPath(context: LauncherCommandContext): string {
@@ -96,12 +101,26 @@ const defaultSyncCommandDeps: SyncCommandDeps = {
consoleLog: console.log,
writeStdout: process.stdout.write.bind(process.stdout),
ensureTrackerQuiescent: async (context, dbPath) => ensureTrackerQuiescent(context, dbPath),
emitEvent: () => {},
recordHostSyncResult: recordHostSyncResultToDisk,
};
function resolveSyncCommandDeps(inputDeps: Partial<SyncCommandDeps> = {}): SyncCommandDeps {
return { ...defaultSyncCommandDeps, ...inputDeps };
}
// In --json mode every line on stdout is an NDJSON event: human console output
// is silenced and events are written through the original console logger.
function withJsonEvents(deps: SyncCommandDeps): SyncCommandDeps {
const writeLine = deps.consoleLog;
return {
...deps,
consoleLog: () => {},
writeStdout: (() => true) as typeof process.stdout.write,
emitEvent: (event) => writeLine(JSON.stringify(event)),
};
}
export function runSnapshotMode(
context: LauncherCommandContext,
dbPath: string,
@@ -109,7 +128,13 @@ export function runSnapshotMode(
): void {
const deps = resolveSyncCommandDeps(inputDeps);
const outPath = resolvePathMaybe(context.args.syncSnapshotPath);
deps.emitEvent({
type: 'stage',
stage: 'snapshot-local',
message: `Snapshotting local database (${dbPath})`,
});
deps.createDbSnapshot(dbPath, outPath);
deps.emitEvent({ type: 'snapshot-created', path: outPath });
deps.consoleLog(outPath);
}
@@ -121,10 +146,73 @@ export async function runMergeMode(
const deps = resolveSyncCommandDeps(inputDeps);
await deps.ensureTrackerQuiescent(context, dbPath);
const snapshotPath = resolvePathMaybe(context.args.syncMergePath);
deps.emitEvent({
type: 'stage',
stage: 'merge-local',
message: `Merging ${snapshotPath} into the local database`,
});
const summary = deps.mergeSnapshotIntoDb(dbPath, snapshotPath);
deps.emitEvent({ type: 'merge-summary', target: 'local', summary });
deps.consoleLog(deps.formatMergeSummary(summary));
}
function formatHostSyncDetail(
direction: 'both' | 'push' | 'pull',
pulledSummary: SyncMergeSummary | null,
): string {
if (!pulledSummary) return direction === 'push' ? 'Pushed local stats' : 'Sync complete';
const merged = `${pulledSummary.sessionsMerged} session${pulledSummary.sessionsMerged === 1 ? '' : 's'} merged`;
return direction === 'pull' ? merged : `${merged}; pushed local stats`;
}
export async function runCheckMode(
context: LauncherCommandContext,
inputDeps: Partial<SyncCommandDeps> = {},
): Promise<void> {
const deps = resolveSyncCommandDeps(inputDeps);
const { args } = context;
const host = args.syncHost;
deps.assertSafeSshHost(host);
deps.consoleLog(`Checking SSH connection to ${host}...`);
let remoteCommand: string | null = null;
let remoteVersion: string | null = null;
let error: string | null = null;
const probe = deps.runSsh(host, 'echo subminer-check-ok');
const sshOk = probe.status === 0 && probe.stdout.includes('subminer-check-ok');
if (!sshOk) {
error = formatRemoteRunError(`Could not reach ${host} over SSH.`, probe);
} else {
deps.consoleLog('SSH connection: ok');
try {
remoteCommand = deps.resolveRemoteSubminerCommand(host, args.syncRemoteCmd || null);
const version = deps.runSsh(host, `${remoteCommand} --version`);
remoteVersion = version.status === 0 ? version.stdout.trim() || null : null;
deps.consoleLog(
`Remote subminer: ${remoteCommand}${remoteVersion ? ` (${remoteVersion})` : ''}`,
);
} catch (resolveError) {
error = resolveError instanceof Error ? resolveError.message : String(resolveError);
}
}
const ok = sshOk && remoteCommand !== null;
deps.emitEvent({
type: 'check-result',
host,
sshOk,
remoteCommand,
remoteVersion,
ok,
error,
});
if (!ok) {
throw new Error(error ?? `Connection check failed for ${host}.`);
}
deps.consoleLog('Check passed.');
}
function cleanupRemote(host: string, remoteTmpDir: string, deps: SyncCommandDeps): void {
if (!remoteTmpDir.startsWith('/tmp/')) return;
deps.runSsh(host, `rm -rf ${shellQuote(remoteTmpDir)}`);
@@ -155,6 +243,7 @@ export async function runHostSync(
const localTmpDir = deps.mkdtempSync(path.join(os.tmpdir(), 'subminer-sync-'));
let remoteTmpDir = '';
let pulledSummary: SyncMergeSummary | null = null;
try {
// Signal failures by throwing (not fail(), which exits synchronously and
// would skip the finally cleanup, leaking temp dirs holding snapshot data).
@@ -170,12 +259,18 @@ export async function runHostSync(
const localSnapshot = path.join(localTmpDir, 'local.sqlite');
if (shouldPush) {
deps.consoleLog(`Snapshotting local database (${dbPath})...`);
deps.emitEvent({
type: 'stage',
stage: 'snapshot-local',
message: `Snapshotting local database (${dbPath})`,
});
deps.createDbSnapshot(dbPath, localSnapshot);
}
const remoteSnapshot = `${remoteTmpDir}/snapshot.sqlite`;
if (shouldPull) {
deps.consoleLog(`Snapshotting ${host}...`);
deps.emitEvent({ type: 'stage', stage: 'snapshot-remote', message: `Snapshotting ${host}` });
const snapshotRun = deps.runSsh(
host,
`${remoteCmd} sync --snapshot ${shellQuote(remoteSnapshot)}${forceFlag}`,
@@ -186,25 +281,50 @@ export async function runHostSync(
}
const pulledSnapshot = path.join(localTmpDir, 'remote.sqlite');
if (shouldPull) deps.runScp(`${host}:${remoteSnapshot}`, pulledSnapshot);
if (shouldPull) {
deps.emitEvent({
type: 'stage',
stage: 'download',
message: `Copying snapshot from ${host}`,
});
deps.runScp(`${host}:${remoteSnapshot}`, pulledSnapshot);
}
const incomingSnapshot = `${remoteTmpDir}/incoming.sqlite`;
if (shouldPush) deps.runScp(localSnapshot, `${host}:${incomingSnapshot}`);
if (shouldPush) {
deps.emitEvent({ type: 'stage', stage: 'upload', message: `Copying snapshot to ${host}` });
deps.runScp(localSnapshot, `${host}:${incomingSnapshot}`);
}
if (shouldPull) {
deps.consoleLog(`\nMerging ${host} -> local:`);
deps.emitEvent({
type: 'stage',
stage: 'merge-local',
message: `Merging ${host} into the local database`,
});
await deps.ensureTrackerQuiescent(context, dbPath);
const summary = deps.mergeSnapshotIntoDb(dbPath, pulledSnapshot);
pulledSummary = summary;
deps.emitEvent({ type: 'merge-summary', target: 'local', summary });
deps.consoleLog(deps.formatMergeSummary(summary));
}
if (shouldPush) {
deps.consoleLog(`\nMerging local -> ${host}:`);
deps.emitEvent({
type: 'stage',
stage: 'merge-remote',
message: `Merging the local database into ${host}`,
});
await deps.ensureTrackerQuiescent(context, dbPath);
const mergeRun = deps.runSsh(
host,
`${remoteCmd} sync --merge ${shellQuote(incomingSnapshot)}${forceFlag}`,
);
deps.writeStdout(mergeRun.stdout);
if (mergeRun.stdout.trim()) {
deps.emitEvent({ type: 'remote-output', text: mergeRun.stdout });
}
if (mergeRun.status !== 0) {
const retryCommand =
direction === 'push' ? `subminer sync ${host} --push` : `subminer sync ${host}`;
@@ -219,6 +339,18 @@ export async function runHostSync(
}
deps.consoleLog('\nSync complete.');
deps.recordHostSyncResult(host, 'success', formatHostSyncDetail(direction, pulledSummary));
} catch (error) {
try {
deps.recordHostSyncResult(
host,
'error',
error instanceof Error ? error.message : String(error),
);
} catch {
// best effort
}
throw error;
} finally {
deps.rmSync(localTmpDir, { recursive: true, force: true });
if (remoteTmpDir) {
@@ -235,19 +367,34 @@ export async function runSyncCommand(
context: LauncherCommandContext,
inputDeps: Partial<SyncCommandDeps> = {},
): Promise<boolean> {
const deps = resolveSyncCommandDeps(inputDeps);
let deps = resolveSyncCommandDeps(inputDeps);
const { args } = context;
if (!args.sync) return false;
if (args.syncJson) deps = withJsonEvents(deps);
const dbPath = resolveDbPath(context);
if (args.syncSnapshotPath) {
runSnapshotMode(context, dbPath, deps);
} else if (args.syncMergePath) {
await runMergeMode(context, dbPath, deps);
} else if (args.syncHost) {
await runHostSync(context, dbPath, deps);
} else {
deps.fail('sync requires a host, --snapshot <file>, or --merge <file>.');
try {
if (args.syncCheck) {
await runCheckMode(context, deps);
} else if (args.syncSnapshotPath) {
runSnapshotMode(context, dbPath, deps);
} else if (args.syncMergePath) {
await runMergeMode(context, dbPath, deps);
} else if (args.syncHost) {
await runHostSync(context, dbPath, deps);
} else {
deps.fail('sync requires a host, --snapshot <file>, or --merge <file>.');
}
} catch (error) {
if (args.syncJson) {
deps.emitEvent({
type: 'result',
ok: false,
error: error instanceof Error ? error.message : String(error),
});
}
throw error;
}
if (args.syncJson) deps.emitEvent({ type: 'result', ok: true, error: null });
return true;
}