import { existsSync } from 'node:fs'; import { join } from 'node:path'; import { setTimeout as sleep } from 'node:timers/promises'; import type { BackupCommand, MaintenanceEvent, RestoreStatus } from '@dorfteich/shared'; import { archiveFileName, dumpFileName } from './backup-set.js'; import type { RemoteLogger } from './remote.js'; import { writeRestoreStatus } from './status.js'; /** * The in-app restore, orchestrated by the sidecar (issue #103): * * 1. `restore-status.json` → `running` — the api's maintenance gate flips * to 503 for everything but health and the status endpoint. * 2. NOTIFY maintenance `enter` — the collab server persists + closes every * live session and refuses new ones, so no in-memory document writes * pre-restore content back afterwards. * 3. Grace period, then terminate all other database connections. * 4. Fetch the set (remote source: download + unpack the bundle), verify * both artifacts exist, `pg_restore --clean` + volume archive extract — * the exact path of the operator's `restore.sh`. * 5. `restore-status.json` → final state; NOTIFY `exit`. The api restarts * itself on `succeeded` (fresh caches, migrate-on-start for older dumps). * * Every failure lands in `restore-status.json` — the admin watches that * file through the exempt status endpoint, so it must always resolve. */ export interface RestoreOrchestratorDeps { backupsDir: string; now(): Date; /** NOTIFY on a short-lived connection (maintenance enter/exit). */ notifyMaintenance(event: MaintenanceEvent): Promise; /** Kills every other DB connection right before pg_restore. */ terminateOtherConnections(): Promise; /** Downloads + unpacks the remote bundle into backupsDir (remote.ts). */ fetchRemoteSet(backupId: string): Promise; /** pg_restore + volume extract — shared with the restore.js CLI. */ restoreSet(backupId: string): Promise; /** Milliseconds between maintenance enter and connection termination. */ graceMs?: number; log: RemoteLogger; } export async function orchestrateRestore( deps: RestoreOrchestratorDeps, command: Extract, ): Promise { const startedAt = deps.now().toISOString(); const base: Omit = { schemaVersion: 1, backupId: command.backupId, source: command.source, requestedBy: command.requestedBy, startedAt, }; await writeRestoreStatus(deps.backupsDir, { ...base, state: 'running', finishedAt: null }); deps.log.info( { backupId: command.backupId, source: command.source, requestedBy: command.requestedBy }, 'restore started — instance entering maintenance mode', ); let entered = false; try { await deps.notifyMaintenance({ phase: 'enter' }); entered = true; await sleep(deps.graceMs ?? 5000); if (command.source === 'remote') { await deps.fetchRemoteSet(command.backupId); } for (const file of [dumpFileName(command.backupId), archiveFileName(command.backupId)]) { if (!existsSync(join(deps.backupsDir, file))) { throw new Error(`restore set is incomplete — ${file} not found`); } } await deps.terminateOtherConnections(); await deps.restoreSet(command.backupId); const status: RestoreStatus = { ...base, state: 'succeeded', finishedAt: deps.now().toISOString(), }; await writeRestoreStatus(deps.backupsDir, status); deps.log.info({ backupId: command.backupId }, 'restore succeeded — api will restart'); return status; } catch (error) { const message = error instanceof Error ? error.message : String(error); const status: RestoreStatus = { ...base, state: 'failed', finishedAt: deps.now().toISOString(), error: message, }; // Best effort — if even this write fails the api's staleness bound // (RESTORE_STALE_MAX_AGE_MINUTES) unblocks the instance eventually. await writeRestoreStatus(deps.backupsDir, status).catch(() => undefined); deps.log.error({ backupId: command.backupId, error: message }, 'restore failed'); return status; } finally { if (entered) { await deps.notifyMaintenance({ phase: 'exit' }).catch((error: unknown) => { deps.log.warn({ error: String(error) }, 'maintenance exit notify failed'); }); } } }