#!/usr/bin/env node import { spawnSync } from 'node:child_process'; import { createHash, createHmac } from 'node:crypto'; import { chmodSync, closeSync, createReadStream, createWriteStream, existsSync, lstatSync, mkdirSync, openSync, readdirSync, readFileSync, readlinkSync, realpathSync, renameSync, rmSync, statfsSync, statSync, symlinkSync, writeFileSync, } from 'node:fs'; import { basename, dirname, isAbsolute, join, relative, resolve, sep, } from 'node:path'; import { Readable } from 'node:stream'; import { pipeline } from 'node:stream/promises'; import { setTimeout as sleep } from 'node:timers/promises'; import { fileURLToPath } from 'node:url'; import { gunzipSync, gzipSync } from 'node:zlib'; const __filename = fileURLToPath(import.meta.url); const __dirname = dirname(__filename); const REPO_ROOT = resolve(__dirname, '..'); const DEFAULT_LOCAL_DATA_DIR = resolve( REPO_ROOT, 'server-rs/.spacetimedb/local/data', ); const DEFAULT_LOCAL_WORK_DIR = resolve( REPO_ROOT, 'server-rs/.data/database-backups', ); const DEFAULT_PRODUCTION_DATA_DIR = '/stdb'; const DEFAULT_PRODUCTION_WORK_DIR = '/var/lib/genarrative/database-backups'; const DEFAULT_DATABASE_BACKUP_STOP_MARKER = join( DEFAULT_PRODUCTION_WORK_DIR, '.spacetimedb-stopped', ); const DEFAULT_SPACE_SAFETY_RATIO = 1.1; const DEFAULT_EXTRA_FREE_BYTES = 512 * 1024 * 1024; const OSS_ALGORITHM = 'OSS4-HMAC-SHA256'; const OSS_SERVICE = 'oss'; const OSS_REQUEST = 'aliyun_v4_request'; const UNSIGNED_PAYLOAD = 'UNSIGNED-PAYLOAD'; const DEFAULT_OSS_MULTIPART_PART_SIZE_BYTES = 128 * 1024 * 1024; const OSS_MIN_MULTIPART_PART_SIZE_BYTES = 100 * 1024; const OSS_MAX_MULTIPART_PART_SIZE_BYTES = 5 * 1024 * 1024 * 1024; const OSS_MAX_MULTIPART_PARTS = 10_000; const DEFAULT_OSS_REQUEST_MAX_ATTEMPTS = 5; const DEFAULT_OSS_RETRY_BASE_DELAY_MS = 1_000; const DEFAULT_OSS_RETRY_MAX_DELAY_MS = 30_000; const DEFAULT_DIRECT_FILES_CONCURRENCY = 16; const MAX_DIRECT_FILES_CONCURRENCY = 64; const DIRECT_FILES_SINGLE_PUT_MAX_BYTES = 16 * 1024 * 1024; const RETRYABLE_OSS_HTTP_STATUSES = new Set([408, 429, 500, 502, 503, 504]); const HISTORY_STATE_SCHEMA_VERSION = 1; const HISTORY_MANIFEST_SCHEMA_VERSION = 1; const DIRECT_FILES_STATE_SCHEMA_VERSION = 2; const LEGACY_DIRECT_FILES_STATE_SCHEMA_VERSION = 1; const DIRECT_FILES_CATALOG_SCHEMA_VERSION = 1; const DIRECT_FILES_LATEST_SCHEMA_VERSION = 1; function usage() { console.log(`用法: npm run database:backup:oss -- [--mode full|history] [--storage-format archive|files] [--data-dir ] [--work-dir ] [--bucket ] [--object-prefix ] [--keep-local] node -- scripts/database-backup-to-oss.mjs [--stop-service spacetimedb.service] [--restart-service-after genarrative-api.service] [--defer-upload] node -- scripts/database-backup-to-oss.mjs --upload-archive node -- scripts/database-backup-to-oss.mjs --upload-deferred-dir node -- scripts/database-backup-to-oss.mjs --publish-manifest node -- scripts/database-backup-to-oss.mjs --restore-files-state --restore-dir node -- scripts/database-backup-to-oss.mjs --restore-files-latest --restore-dir [--dry-run] 说明: 将 SpacetimeDB 数据目录以 .tar.gz 或逐文件 catalog 形式上传到阿里云 OSS 指定 bucket。 默认 full 模式保持原有全量冷备行为;history 模式只归档已被最新 snapshot 覆盖的历史 commitlog 与旧 snapshot。 --storage-format files 不打包:按原相对路径建立 catalog,文件内容以 SHA-256 不可变对象上传;重复运行只上传新增或变化内容。 archive history 必须有已验真的 full baseline state;files history 必须复用同一 work-dir 中已发布的 full catalog state。 --defer-upload 只生成本地冷备份和 manifest,不上传;后续用 --upload-archive 异步上传。 --upload-deferred-dir 串行收敛目录内 deferred/pending 及已上传未清理归档;只有 OSS 上传与验真完成后才按 keep-local 规则删除。 默认读取 .env / .env.local / .env.secrets.local;生产服务可传 --env-file /etc/genarrative/api-server.env。 shell 环境变量优先级最高,不会被 env 文件覆盖。 常用环境变量: GENARRATIVE_DATABASE_BACKUP_DATA_DIR 数据目录;生产建议 /stdb GENARRATIVE_DATABASE_BACKUP_WORK_DIR 本地临时备份目录;生产建议 /var/lib/genarrative/database-backups GENARRATIVE_DATABASE_BACKUP_OSS_BUCKET 备份 bucket;未设置时回退 ALIYUN_OSS_BUCKET GENARRATIVE_DATABASE_BACKUP_OSS_PREFIX 对象前缀,默认 database-backups GENARRATIVE_DATABASE_BACKUP_OSS_ENDPOINT OSS endpoint;未设置时回退 ALIYUN_OSS_ENDPOINT GENARRATIVE_DATABASE_BACKUP_STORAGE_FORMAT archive(默认)或 files GENARRATIVE_DATABASE_BACKUP_KEEP_LOCAL true 时保留本地 tar.gz GENARRATIVE_DATABASE_BACKUP_MIN_FREE_BYTES 备份前要求 work dir 所在文件系统至少有这些可用字节;未设置时按数据目录大小估算 GENARRATIVE_DATABASE_BACKUP_BASELINE_STATE history 使用的 full baseline 与追加批次状态文件 GENARRATIVE_DATABASE_BACKUP_BASELINE_MANIFEST 首次初始化 history state 的 uploaded full manifest ALIYUN_OSS_ACCESS_KEY_ID / ALIYUN_OSS_ACCESS_KEY_SECRET `); } function loadEnvFile(filePath, target, protectedKeys) { if (!existsSync(filePath)) { return; } const rawText = readFileSync(filePath, 'utf8'); for (const rawLine of rawText.split(/\r?\n/u)) { const line = rawLine.trim(); if (!line || line.startsWith('#')) { continue; } const match = line.match(/^([A-Za-z_][A-Za-z0-9_]*)=(.*)$/u); if (!match) { continue; } const [, key, rawValue] = match; if (protectedKeys.has(key)) { continue; } target[key] = rawValue.replace(/^['"]|['"]$/gu, ''); } } function loadRepoEnv() { const env = { ...process.env }; const protectedKeys = new Set( Object.entries(process.env) .filter(([, value]) => String(value ?? '').trim()) .map(([key]) => key), ); for (const fileName of ['.env', '.env.local', '.env.secrets.local']) { loadEnvFile(resolve(REPO_ROOT, fileName), env, protectedKeys); } return env; } function loadEffectiveEnv(envFiles) { const env = loadRepoEnv(); const protectedKeys = new Set( Object.entries(process.env) .filter(([, value]) => String(value ?? '').trim()) .map(([key]) => key), ); for (const filePath of envFiles) { loadEnvFile(resolvePath(filePath), env, protectedKeys); } return env; } function parseArgs(argv) { const options = { dataDir: '', workDir: '', bucket: '', endpoint: '', objectPrefix: '', accessKeyId: '', accessKeySecret: '', envFiles: [], keepLocal: false, stopService: '', restartServicesAfter: [], database: '', dryRun: false, deferUpload: false, uploadArchive: '', uploadDeferredDir: '', manifestFile: '', objectKey: '', resultFile: '', minFreeBytes: '', mode: 'full', baselineState: '', baselineManifest: '', publishManifest: '', storageFormat: '', restoreFilesState: '', restoreFilesLatest: false, restoreDir: '', }; for (let index = 0; index < argv.length; index += 1) { const arg = argv[index]; const readValue = () => { const value = argv[index + 1]; if (!value || value.startsWith('--')) { throw new Error(`${arg} 缺少参数值`); } index += 1; return value; }; switch (arg) { case '--help': case '-h': usage(); process.exit(0); break; case '--data-dir': options.dataDir = readValue(); break; case '--work-dir': options.workDir = readValue(); break; case '--bucket': options.bucket = readValue(); break; case '--endpoint': options.endpoint = readValue(); break; case '--object-prefix': options.objectPrefix = readValue(); break; case '--object-key': options.objectKey = readValue(); break; case '--access-key-id': options.accessKeyId = readValue(); break; case '--access-key-secret': options.accessKeySecret = readValue(); break; case '--env-file': options.envFiles.push(readValue()); break; case '--database': options.database = readValue(); break; case '--stop-service': options.stopService = readValue(); break; case '--restart-service-after': options.restartServicesAfter.push(readValue()); break; case '--keep-local': options.keepLocal = true; break; case '--dry-run': options.dryRun = true; break; case '--defer-upload': options.deferUpload = true; options.keepLocal = true; break; case '--upload-archive': options.uploadArchive = readValue(); break; case '--upload-deferred-dir': options.uploadDeferredDir = readValue(); break; case '--manifest-file': options.manifestFile = readValue(); break; case '--result-file': options.resultFile = readValue(); break; case '--min-free-bytes': options.minFreeBytes = readValue(); break; case '--mode': options.mode = readValue(); break; case '--baseline-state': options.baselineState = readValue(); break; case '--baseline-manifest': options.baselineManifest = readValue(); break; case '--publish-manifest': options.publishManifest = readValue(); break; case '--storage-format': options.storageFormat = readValue(); break; case '--restore-files-state': options.restoreFilesState = readValue(); break; case '--restore-files-latest': options.restoreFilesLatest = true; break; case '--restore-dir': options.restoreDir = readValue(); break; default: throw new Error(`未知参数: ${arg}`); } } return options; } function firstNonEmpty(...values) { return values.map((value) => String(value ?? '').trim()).find(Boolean) ?? ''; } function parseDirectFilesConcurrency(rawValue) { const value = Number( String(rawValue ?? DEFAULT_DIRECT_FILES_CONCURRENCY).trim(), ); if ( !Number.isSafeInteger(value) || value < 1 || value > MAX_DIRECT_FILES_CONCURRENCY ) { throw new Error( `GENARRATIVE_DATABASE_BACKUP_FILES_CONCURRENCY 必须是 1-${MAX_DIRECT_FILES_CONCURRENCY} 的整数,实际: ${rawValue}`, ); } return value; } export function createUploadBandwidthLimiter( rawValue, { nowFn = Date.now, sleepImpl = sleep } = {}, ) { const normalized = String(rawValue ?? '').trim(); if (!normalized || normalized === '0') { return null; } const maxBytesPerSecond = Number(normalized); if (!Number.isSafeInteger(maxBytesPerSecond) || maxBytesPerSecond < 1024) { throw new Error( `GENARRATIVE_DATABASE_BACKUP_UPLOAD_MAX_BYTES_PER_SECOND 必须为空、0 或 >= 1024 的整数,实际: ${rawValue}`, ); } let nextAvailableAtMs = 0; const waitForChunk = async (sizeBytes) => { const now = nowFn(); const startAt = Math.max(now, nextAvailableAtMs); const finishAt = startAt + (sizeBytes / maxBytesPerSecond) * 1000; nextAvailableAtMs = finishAt; const delayMs = Math.max(0, finishAt - now); if (delayMs > 0) { await sleepImpl(delayMs); } }; return { maxBytesPerSecond, wrap(readable) { return Readable.from( (async function* throttleUpload() { for await (const chunk of readable) { await waitForChunk(chunk.length); yield chunk; } })(), { objectMode: false }, ); }, }; } function createBufferReadStream(buffer, chunkSizeBytes = 64 * 1024) { return Readable.from( (function* readChunks() { for (let offset = 0; offset < buffer.length; offset += chunkSizeBytes) { yield buffer.subarray( offset, Math.min(offset + chunkSizeBytes, buffer.length), ); } })(), { objectMode: false }, ); } function resolvePath(value) { return isAbsolute(value) ? value : resolve(REPO_ROOT, value); } function normalizeEndpoint(raw) { return String(raw ?? '') .trim() .replace(/^https?:\/\//u, '') .replace(/\/+$/u, ''); } function sanitizeObjectPart(value, fallback) { const sanitized = String(value ?? '') .trim() .toLowerCase() .replace(/[^a-z0-9._-]+/gu, '-') .replace(/-+/gu, '-') .replace(/^-|-$/gu, ''); return sanitized || fallback; } function timestampForFile(date = new Date()) { const pad = (value) => String(value).padStart(2, '0'); return `${date.getUTCFullYear()}${pad(date.getUTCMonth() + 1)}${pad(date.getUTCDate())}T${pad(date.getUTCHours())}${pad(date.getUTCMinutes())}${pad(date.getUTCSeconds())}Z`; } function buildBackupNames({ database, dataDir, objectPrefix }) { const timestamp = timestampForFile(); const databasePart = sanitizeObjectPart( database || basename(dataDir), 'spacetimedb', ); const fileName = `${databasePart}-${timestamp}.tar.gz`; const prefix = String(objectPrefix || 'database-backups') .trim() .replace(/^\/+|\/+$/gu, '') .split('/') .filter(Boolean) .map((part) => sanitizeObjectPart(part, 'backup')) .join('/'); const objectKey = [prefix, databasePart, fileName].filter(Boolean).join('/'); return { fileName, objectKey }; } function atomicWriteBuffer(filePath, body) { mkdirSync(dirname(filePath), { recursive: true }); const tempPath = `${filePath}.${process.pid}.${Date.now()}.tmp`; writeFileSync(tempPath, body, { mode: 0o600 }); chmodSync(tempPath, 0o600); renameSync(tempPath, filePath); } function atomicWriteJson(filePath, payload) { atomicWriteBuffer( filePath, Buffer.from(`${JSON.stringify(payload, null, 2)}\n`, 'utf8'), ); } function atomicWriteGzipJson(filePath, payload) { const body = Buffer.from(`${JSON.stringify(payload)}\n`, 'utf8'); atomicWriteBuffer(filePath, gzipSync(body, { level: 9 })); } function processIsAlive(pid) { try { process.kill(pid, 0); return true; } catch (error) { return error?.code === 'EPERM'; } } function acquireBackupLock({ workDir, database }) { mkdirSync(workDir, { recursive: true }); const lockPath = join( workDir, `${sanitizeObjectPart(database, 'spacetimedb')}.backup.lock`, ); try { const fd = openSync(lockPath, 'wx', 0o600); writeFileSync(fd, `${process.pid}\n`, 'utf8'); closeSync(fd); const release = () => { try { const ownerPid = Number(String(readFileSync(lockPath, 'utf8')).trim()); if (ownerPid === process.pid) { rmSync(lockPath, { force: true }); } } catch { // The lock may already have been removed by the normal exit path. } }; process.once('exit', release); for (const signal of ['SIGINT', 'SIGTERM']) { process.once(signal, () => { release(); process.exit(signal === 'SIGINT' ? 130 : 143); }); } return lockPath; } catch (error) { if (error?.code !== 'EEXIST') { throw error; } } const ownerPid = Number(String(readFileSync(lockPath, 'utf8')).trim()); if ( Number.isSafeInteger(ownerPid) && ownerPid > 0 && processIsAlive(ownerPid) ) { throw new Error(`已有数据库备份进程持有锁: ${lockPath} pid=${ownerPid}`); } throw new Error( `发现失效数据库备份锁,拒绝自动抢锁;请核对 OSS multipart 与进程后手工删除: ${lockPath} pid=${ownerPid || ''}`, ); } function historyStatePath({ args, env, workDir, database }) { return resolvePath( firstNonEmpty( args.baselineState, env.GENARRATIVE_DATABASE_BACKUP_BASELINE_STATE, join( workDir, `${sanitizeObjectPart(database, 'spacetimedb')}-history-state.json`, ), ), ); } function baselineIdFor(baseline) { return sha256Hex( [ baseline.bucket, baseline.objectKey, baseline.verifiedAt, baseline.contentLength, baseline.archiveSha256, ].join('\0'), ).slice(0, 24); } function normalizeUploadedBaselineManifest(manifest, { database, dataDir }) { if (manifest.uploadStatus !== 'uploaded') { throw new Error( `baseline manifest 必须是 uploaded,实际: ${manifest.uploadStatus ?? ''}`, ); } if (manifest.backupKind !== 'spacetimedb-data-dir') { throw new Error( `baseline manifest backupKind 必须是 spacetimedb-data-dir,实际: ${manifest.backupKind ?? ''}`, ); } const baseline = { backupKind: 'spacetimedb-data-dir', database: firstNonEmpty(manifest.database, database), dataDir: resolvePath(firstNonEmpty(manifest.dataDir, dataDir)), bucket: String(manifest.bucket ?? '').trim(), objectKey: String(manifest.objectKey ?? '').trim(), verifiedAt: String(manifest.verifiedAt ?? '').trim(), uploadedAt: String(manifest.uploadedAt ?? '').trim(), contentLength: Number(manifest.contentLength), archiveSha256: String(manifest.archiveSha256 ?? '') .trim() .toLowerCase(), manifestObjectKey: String(manifest.manifestObjectKey ?? '').trim(), manifestContentLength: Number(manifest.manifestContentLength), manifestArchiveSha256: String(manifest.manifestArchiveSha256 ?? '') .trim() .toLowerCase(), manifestVerifiedAt: String(manifest.manifestVerifiedAt ?? '').trim(), }; if ( !baseline.bucket || !baseline.objectKey || !baseline.verifiedAt || !Number.isSafeInteger(baseline.contentLength) || baseline.contentLength <= 0 || !/^[a-f0-9]{64}$/u.test(baseline.archiveSha256) || !baseline.manifestObjectKey || !Number.isSafeInteger(baseline.manifestContentLength) || baseline.manifestContentLength <= 0 || !/^[a-f0-9]{64}$/u.test(baseline.manifestArchiveSha256) || !baseline.manifestVerifiedAt ) { throw new Error('baseline manifest 缺少已验真 OSS 归档或 sidecar 信息。'); } baseline.id = baselineIdFor(baseline); return baseline; } function validateHistoryState(state, { database, dataDir }) { if (state.schemaVersion !== HISTORY_STATE_SCHEMA_VERSION || !state.baseline) { throw new Error('history state schemaVersion 或 baseline 无效。'); } const baseline = normalizeUploadedBaselineManifest( { ...state.baseline, uploadStatus: 'uploaded' }, { database, dataDir }, ); if (baseline.database !== database) { throw new Error( `history state database 不匹配: expected=${database}, actual=${baseline.database}`, ); } if (resolvePath(baseline.dataDir) !== resolvePath(dataDir)) { throw new Error( `history state dataDir 不匹配: expected=${resolvePath(dataDir)}, actual=${resolvePath(baseline.dataDir)}`, ); } return { ...state, baseline, batches: Array.isArray(state.batches) ? state.batches : [], }; } function writeBaselineState({ statePath, baseline, previousState = null }) { const state = { schemaVersion: HISTORY_STATE_SCHEMA_VERSION, updatedAt: new Date().toISOString(), baseline, batches: previousState?.baseline?.id === baseline.id && Array.isArray(previousState.batches) ? previousState.batches : [], }; atomicWriteJson(statePath, state); return state; } function loadOrImportHistoryState({ args, env, statePath, database, dataDir }) { if (existsSync(statePath)) { return validateHistoryState(readManifest(statePath), { database, dataDir }); } const importPath = firstNonEmpty( args.baselineManifest, env.GENARRATIVE_DATABASE_BACKUP_BASELINE_MANIFEST, ); if (!importPath) { throw new Error( `history 模式缺少已验真 baseline state: ${statePath};可用 --baseline-manifest 导入已有 uploaded baseline manifest。`, ); } const baseline = normalizeUploadedBaselineManifest( readManifest(resolvePath(importPath)), { database, dataDir }, ); return writeBaselineState({ statePath, baseline }); } function assertSafeRelativePath(dataDir, absolutePath) { const relativePath = relative( resolvePath(dataDir), resolvePath(absolutePath), ); if ( !relativePath || relativePath === '..' || relativePath.startsWith(`..${sep}`) || isAbsolute(relativePath) ) { throw new Error(`history 候选路径越界或等于数据目录: ${absolutePath}`); } return relativePath.split(sep).join('/'); } function statFingerprint(absolutePath, rootPath = absolutePath) { // 候选 snapshot 可能包含数十万条目录项;增量更新摘要,避免把每条 // fingerprint 字符串同时保存在 entries[] 后再 join,造成一次性内存峰值。 const fingerprintHash = createHash('sha256'); let isFirstEntry = true; let entryCount = 0; let totalSize = 0n; const visit = (currentPath) => { const stat = lstatSync(currentPath, { bigint: true }); if (stat.isSymbolicLink()) { throw new Error(`history 候选不得包含符号链接: ${currentPath}`); } const entryPath = currentPath === rootPath ? '.' : relative(rootPath, currentPath).split(sep).join('/'); const kind = stat.isDirectory() ? 'directory' : stat.isFile() ? 'file' : 'other'; if (kind === 'other') { throw new Error(`history 候选只允许普通文件或目录: ${currentPath}`); } const entry = [ entryPath, kind, stat.dev.toString(), stat.ino.toString(), stat.mode.toString(), stat.size.toString(), stat.mtimeNs.toString(), ].join('\0'); if (!isFirstEntry) { fingerprintHash.update('\n'); } fingerprintHash.update(entry); isFirstEntry = false; entryCount += 1; if (stat.isFile()) { totalSize += stat.size; } else { for (const name of readdirSync(currentPath).sort()) { visit(join(currentPath, name)); } } }; visit(rootPath); return { fingerprint: fingerprintHash.digest('hex'), sizeBytes: totalSize.toString(), entryCount, }; } function findReplicasDir(dataDir) { const candidates = [ resolve(dataDir, 'replicas'), resolve(dataDir, 'data', 'replicas'), ].filter( (candidate) => existsSync(candidate) && lstatSync(candidate).isDirectory(), ); if (candidates.length !== 1) { throw new Error( `无法唯一确定 replicas 目录: ${candidates.length === 0 ? '' : candidates.join(', ')}`, ); } return candidates[0]; } function historyCandidate({ dataDir, absolutePath, kind, replicaId, transaction, }) { const stat = statFingerprint(absolutePath); return { path: assertSafeRelativePath(dataDir, absolutePath), kind, replicaId, transaction: transaction.toString(), ...stat, }; } export function discoverHistoryPlan({ dataDir }) { const resolvedDataDir = resolvePath(dataDir); const replicasDir = findReplicasDir(resolvedDataDir); const replicaEntries = readdirSync(replicasDir, { withFileTypes: true }); const replicas = []; const candidates = []; for (const replicaEntry of replicaEntries.sort((left, right) => left.name.localeCompare(right.name), )) { if (!replicaEntry.isDirectory()) { continue; } if (!/^\d+$/u.test(replicaEntry.name)) { throw new Error(`replica 目录名不符合预期: ${replicaEntry.name}`); } const replicaId = replicaEntry.name; const replicaDir = join(replicasDir, replicaId); const snapshotsDir = join(replicaDir, 'snapshots'); const clogDir = join(replicaDir, 'clog'); if (!existsSync(snapshotsDir) || !lstatSync(snapshotsDir).isDirectory()) { replicas.push({ replicaId, status: 'skipped', reason: 'no-snapshots-directory', }); continue; } const snapshotEntries = readdirSync(snapshotsDir, { withFileTypes: true }); const snapshots = snapshotEntries .flatMap((entry) => { const match = /^(\d{20})\.snapshot_dir$/u.exec(entry.name); if (!match) { return []; } const transaction = BigInt(match[1]); if (transaction > 0xffff_ffff_ffff_ffffn) { throw new Error(`snapshot transaction 超出 u64: ${entry.name}`); } if (!entry.isDirectory()) { throw new Error( `snapshot 候选必须是目录: ${join(snapshotsDir, entry.name)}`, ); } const snapshotDir = join(snapshotsDir, entry.name); const lockPath = join(snapshotsDir, `${match[1]}.lock`); const snapshotFile = join(snapshotDir, `${match[1]}.snapshot_bsatn`); if ( existsSync(lockPath) || !existsSync(snapshotFile) || !lstatSync(snapshotFile).isFile() ) { return []; } return [{ name: entry.name, transaction }]; }) .sort((left, right) => left.transaction < right.transaction ? -1 : left.transaction > right.transaction ? 1 : 0, ); if (snapshots.length === 0) { replicas.push({ replicaId, status: 'skipped', reason: 'no-snapshot' }); continue; } if (!existsSync(clogDir) || !lstatSync(clogDir).isDirectory()) { throw new Error(`replica ${replicaId} 缺少 clog 目录。`); } const segmentFiles = new Map(); for (const entry of readdirSync(clogDir, { withFileTypes: true })) { const match = /^(\d{20})\.stdb\.(log|ofs)$/u.exec(entry.name); if (!match) { throw new Error(`commitlog 文件名不符合预期: ${entry.name}`); } if (!entry.isFile()) { throw new Error( `commitlog 候选必须是普通文件: ${join(clogDir, entry.name)}`, ); } const transaction = BigInt(match[1]); if (transaction > 0xffff_ffff_ffff_ffffn) { throw new Error(`commitlog transaction 超出 u64: ${entry.name}`); } const key = transaction.toString(); const group = segmentFiles.get(key) ?? { transaction }; group[match[2]] = entry.name; segmentFiles.set(key, group); } for (const group of segmentFiles.values()) { if (group.ofs && !group.log) { throw new Error( `commitlog offset 缺少对应 log: replica=${replicaId}, transaction=${group.transaction}`, ); } } const segments = [...segmentFiles.values()] .filter((group) => group.log) .sort((left, right) => left.transaction < right.transaction ? -1 : left.transaction > right.transaction ? 1 : 0, ); const latestSnapshot = snapshots.at(-1).transaction; const boundarySegment = segments .filter((segment) => segment.transaction <= latestSnapshot) .at(-1); if (!boundarySegment) { throw new Error( `replica ${replicaId} 无法找到覆盖 latest snapshot ${latestSnapshot} 的 commitlog 边界。`, ); } for (const snapshot of snapshots.slice(0, -1)) { candidates.push( historyCandidate({ dataDir: resolvedDataDir, absolutePath: join(snapshotsDir, snapshot.name), kind: 'snapshot', replicaId, transaction: snapshot.transaction, }), ); } for (const segment of segments.filter( (item) => item.transaction < boundarySegment.transaction, )) { candidates.push( historyCandidate({ dataDir: resolvedDataDir, absolutePath: join(clogDir, segment.log), kind: 'commitlog', replicaId, transaction: segment.transaction, }), ); if (segment.ofs) { candidates.push( historyCandidate({ dataDir: resolvedDataDir, absolutePath: join(clogDir, segment.ofs), kind: 'commitlog-offset', replicaId, transaction: segment.transaction, }), ); } } replicas.push({ replicaId, status: 'ready', latestSnapshot: latestSnapshot.toString(), boundarySegment: boundarySegment.transaction.toString(), }); } candidates.sort((left, right) => left.path.localeCompare(right.path)); return { dataDir: resolvedDataDir, replicasDir: assertSafeRelativePath(resolvedDataDir, replicasDir), replicas, candidates, totalSizeBytes: candidates .reduce((sum, item) => sum + BigInt(item.sizeBytes), 0n) .toString(), }; } function runCommand(command, args, options = {}) { const result = spawnSync(command, args, { cwd: options.cwd ?? REPO_ROOT, env: options.env ?? process.env, encoding: 'utf8', stdio: options.stdio ?? 'pipe', shell: process.platform === 'win32', }); if (result.error) { throw new Error(`${command} 启动失败: ${result.error.message}`); } if (result.status !== 0) { const output = `${result.stdout ?? ''}\n${result.stderr ?? ''}`.trim(); throw new Error(`${command} 退出码 ${result.status}: ${output}`); } return result; } function parseByteSize(rawValue, label) { const value = String(rawValue ?? '').trim(); if (!value) { return null; } const match = /^(\d+)(?:\s*([KMGTPE]?)(?:I?B?)?)?$/iu.exec(value); if (!match) { throw new Error( `${label} 必须是字节数或 K/M/G/T/P/E 后缀大小,实际: ${rawValue}`, ); } const [, amountText, unitText = ''] = match; const multipliers = { '': 1n, K: 1024n, M: 1024n ** 2n, G: 1024n ** 3n, T: 1024n ** 4n, P: 1024n ** 5n, E: 1024n ** 6n, }; return BigInt(amountText) * multipliers[unitText.toUpperCase()]; } function formatBytes(bytes) { const value = BigInt(bytes); const gib = Number(value) / 1024 ** 3; if (gib >= 1) { return `${gib.toFixed(1)}GiB`; } const mib = Number(value) / 1024 ** 2; if (mib >= 1) { return `${mib.toFixed(1)}MiB`; } return `${value}B`; } function getDirectorySizeBytes(dataDir) { const result = runCommand('du', ['-sk', dataDir]); const [sizeKbText] = String(result.stdout ?? '') .trim() .split(/\s+/u); if (!sizeKbText || !/^\d+$/u.test(sizeKbText)) { throw new Error(`无法解析数据目录大小: ${result.stdout}`); } return BigInt(sizeKbText) * 1024n; } function getAvailableBytes(fileSystemPath) { const stat = statfsSync(fileSystemPath, { bigint: true }); return stat.bavail * stat.bsize; } function parseSafetyRatio(rawValue) { const value = String(rawValue ?? '').trim(); if (!value) { return DEFAULT_SPACE_SAFETY_RATIO; } const ratio = Number(value); if (!Number.isFinite(ratio) || ratio < 1) { throw new Error( `GENARRATIVE_DATABASE_BACKUP_SPACE_SAFETY_RATIO 必须是 >= 1 的数字,实际: ${rawValue}`, ); } return ratio; } function calculateRequiredFreeBytes({ dataSizeBytes, args, env }) { const explicitMinFreeBytes = parseByteSize( firstNonEmpty( args.minFreeBytes, env.GENARRATIVE_DATABASE_BACKUP_MIN_FREE_BYTES, ), 'GENARRATIVE_DATABASE_BACKUP_MIN_FREE_BYTES', ); if (explicitMinFreeBytes !== null) { return explicitMinFreeBytes; } const ratio = parseSafetyRatio( env.GENARRATIVE_DATABASE_BACKUP_SPACE_SAFETY_RATIO, ); const ratioBasisPoints = BigInt(Math.ceil(ratio * 10000)); const ratioRequirement = (dataSizeBytes * ratioBasisPoints + 9999n) / 10000n; const extraFreeBytes = parseByteSize( firstNonEmpty( env.GENARRATIVE_DATABASE_BACKUP_EXTRA_FREE_BYTES, String(DEFAULT_EXTRA_FREE_BYTES), ), 'GENARRATIVE_DATABASE_BACKUP_EXTRA_FREE_BYTES', ); const extraRequirement = dataSizeBytes + extraFreeBytes; return ratioRequirement > extraRequirement ? ratioRequirement : extraRequirement; } function assertSufficientWorkDirSpace({ dataDir, workDir, args, env }) { mkdirSync(workDir, { recursive: true }); const dataSizeBytes = getDirectorySizeBytes(dataDir); const availableBytes = getAvailableBytes(workDir); const requiredFreeBytes = calculateRequiredFreeBytes({ dataSizeBytes, args, env, }); console.log( `[database-backup] 备份空间预检: data=${formatBytes(dataSizeBytes)}, available=${formatBytes(availableBytes)}, required=${formatBytes(requiredFreeBytes)}`, ); if (availableBytes < requiredFreeBytes) { throw new Error( [ `备份工作目录所在文件系统剩余空间不足: ${workDir}`, `available=${formatBytes(availableBytes)}`, `required=${formatBytes(requiredFreeBytes)}`, `dataDir=${dataDir}`, `dataSize=${formatBytes(dataSizeBytes)}`, '为避免停库后写满磁盘,本次备份已在停止服务前中止。', ].join(';'), ); } } function assertSufficientHistoryWorkDirSpace({ historySizeBytes, workDir, args, env, }) { mkdirSync(workDir, { recursive: true }); const availableBytes = getAvailableBytes(workDir); const requiredFreeBytes = calculateRequiredFreeBytes({ dataSizeBytes: BigInt(historySizeBytes), args, env, }); console.log( `[database-backup] history 空间预检: candidates=${formatBytes(historySizeBytes)}, available=${formatBytes(availableBytes)}, required=${formatBytes(requiredFreeBytes)}`, ); if (availableBytes < requiredFreeBytes) { throw new Error( `history 工作目录剩余空间不足: available=${formatBytes(availableBytes)};required=${formatBytes(requiredFreeBytes)}`, ); } } function collectRestartServicesAfterBackup({ args, env }) { const serviceNames = [ ...String(env.GENARRATIVE_DATABASE_BACKUP_RESTART_SERVICE_AFTER ?? '') .split(',') .map((value) => value.trim()) .filter(Boolean), ...args.restartServicesAfter, ]; return [...new Set(serviceNames.filter(Boolean))]; } function databaseBackupStopMarkerPath(workDir) { return resolvePath( firstNonEmpty( process.env.GENARRATIVE_DATABASE_BACKUP_STOP_MARKER, workDir === DEFAULT_PRODUCTION_WORK_DIR ? DEFAULT_DATABASE_BACKUP_STOP_MARKER : join(workDir, '.spacetimedb-stopped'), ), ); } function writeDatabaseBackupStopMarker(markerPath, serviceName) { atomicWriteJson(markerPath, { serviceName, pid: process.pid, stoppedAt: new Date().toISOString(), }); } function clearDatabaseBackupStopMarker(markerPath) { if (markerPath) { rmSync(markerPath, { force: true }); } } function stopServiceIfNeeded(serviceName, stopMarkerPath) { if (!serviceName) { return false; } console.log(`[database-backup] 停止服务以获取冷备份: ${serviceName}`); writeDatabaseBackupStopMarker(stopMarkerPath, serviceName); // stop 命令失败时仍保留 marker:systemd 的 ExecStopPost 需要它判断是否要 // 兜底恢复,不能因为当前进程还能捕获异常就抹掉上一次停库证据。 runCommand('systemctl', ['stop', serviceName], { stdio: 'inherit' }); return true; } function startServiceIfNeeded(serviceName, wasStopped) { if (!serviceName || !wasStopped) { return; } console.log(`[database-backup] 恢复服务: ${serviceName}`); runCommand('systemctl', ['start', serviceName], { stdio: 'inherit' }); } function restartServicesAfterBackup(serviceNames) { const errors = []; for (const serviceName of serviceNames) { if (!serviceName) { continue; } console.log(`[database-backup] 冷备份后重启依赖服务: ${serviceName}`); try { runCommand('systemctl', ['restart', serviceName], { stdio: 'inherit' }); } catch (error) { errors.push(error); } } if (errors.length > 0) { throw new AggregateError( errors, `冷备份后重启依赖服务失败: ${errors.map((error) => error.message).join('; ')}`, ); } } function restoreServicesAfterBackup({ stopService, serviceStopped, restartServicesAfter, stopMarkerPath, }) { const errors = []; try { startServiceIfNeeded(stopService, serviceStopped); } catch (error) { errors.push(error); } try { restartServicesAfterBackup(restartServicesAfter); } catch (error) { errors.push(error); } if (errors.length > 0) { throw new AggregateError( errors, `恢复冷备份相关服务失败: ${errors.map((error) => error.message).join('; ')}`, ); } clearDatabaseBackupStopMarker(stopMarkerPath); } function createArchive({ dataDir, workDir, fileName }) { if (!existsSync(dataDir)) { throw new Error(`数据库数据目录不存在: ${dataDir}`); } const stat = statSync(dataDir); if (!stat.isDirectory()) { throw new Error(`数据库数据路径不是目录: ${dataDir}`); } mkdirSync(workDir, { recursive: true }); const archivePath = resolve(workDir, fileName); const parentDir = dirname(dataDir); const entryName = basename(dataDir); console.log(`[database-backup] 打包: ${dataDir} -> ${archivePath}`); runCommand('tar', ['-czf', archivePath, '-C', parentDir, entryName], { stdio: 'inherit', }); verifyArchive(archivePath); return archivePath; } function verifyArchive(archivePath) { console.log(`[database-backup] 校验归档: ${archivePath}`); runCommand('tar', ['-tzf', archivePath], { stdio: 'ignore' }); } function historyBatchId({ baselineId, plan }) { const identity = plan.candidates .map((candidate) => [ candidate.path, candidate.kind, candidate.transaction, candidate.fingerprint, ].join('\0'), ) .join('\n'); return sha256Hex(`${baselineId}\0${identity}`).slice(0, 32); } function buildHistoryNames({ database, objectPrefix, baselineId, batchId }) { const databasePart = sanitizeObjectPart(database, 'spacetimedb'); const prefix = String(objectPrefix || 'database-backups') .trim() .replace(/^\/+|\/+$/gu, '') .split('/') .filter(Boolean) .map((part) => sanitizeObjectPart(part, 'backup')) .join('/'); const fileName = `${databasePart}-history-${batchId}.tar.gz`; return { fileName, objectKey: [prefix, databasePart, 'history', baselineId, fileName] .filter(Boolean) .join('/'), }; } function createHistoryArchive({ dataDir, workDir, fileName, manifestPath, candidates, }) { mkdirSync(workDir, { recursive: true }); const archivePath = resolve(workDir, fileName); const candidatePaths = candidates.map((candidate) => candidate.path); console.log( `[database-backup] 打包 history: ${candidatePaths.length} 个候选 -> ${archivePath}`, ); runCommand( 'tar', [ '-czf', archivePath, '-C', dataDir, ...candidatePaths, '-C', dirname(manifestPath), basename(manifestPath), ], { stdio: 'inherit' }, ); verifyArchive(archivePath); return archivePath; } function recordHistoryBatch({ statePath, state, manifest, uploadResult, manifestUpload, status, cleanedAt = '', }) { const batch = { batchId: manifest.batchId, objectKey: uploadResult.objectKey, contentLength: uploadResult.contentLength, archiveSha256: uploadResult.archiveSha256, verifiedAt: uploadResult.verifiedAt, manifestObjectKey: manifestUpload.objectKey, manifestContentLength: manifestUpload.contentLength, manifestArchiveSha256: manifestUpload.archiveSha256, manifestVerifiedAt: manifestUpload.verifiedAt, uploadedAt: manifest.uploadedAt, status, cleanedAt, candidates: manifest.candidates, }; const batches = state.batches.filter( (item) => item.batchId !== batch.batchId, ); batches.push(batch); const nextState = { ...state, updatedAt: new Date().toISOString(), batches }; atomicWriteJson(statePath, nextState); return nextState; } function candidateKey(candidate) { return `${candidate.kind}\0${candidate.path}`; } export function cleanupHistoryCandidates({ dataDir, candidates }) { const currentPlan = discoverHistoryPlan({ dataDir }); const eligible = new Map( currentPlan.candidates.map((candidate) => [ candidateKey(candidate), candidate, ]), ); const existing = []; for (const candidate of candidates) { const absolutePath = resolve(dataDir, candidate.path); assertSafeRelativePath(dataDir, absolutePath); if (!existsSync(absolutePath)) { continue; } const current = eligible.get(candidateKey(candidate)); if (!current) { throw new Error( `history 候选已不在当前安全边界内,拒绝删除: ${candidate.path}`, ); } const currentStat = statFingerprint(absolutePath); if ( currentStat.fingerprint !== candidate.fingerprint || currentStat.sizeBytes !== candidate.sizeBytes ) { throw new Error(`history 候选 stat 漂移,拒绝删除: ${candidate.path}`); } existing.push({ candidate, absolutePath }); } existing.sort((left, right) => { const priority = { 'commitlog-offset': 0, commitlog: 1, snapshot: 2 }; return ( (priority[left.candidate.kind] ?? 3) - (priority[right.candidate.kind] ?? 3) || left.candidate.path.localeCompare(right.candidate.path) ); }); for (const { candidate, absolutePath } of existing) { rmSync(absolutePath, { recursive: candidate.kind === 'snapshot', force: false, }); console.log(`[database-backup] 已清理 history 源文件: ${candidate.path}`); } return { deletedCount: existing.length, alreadyMissingCount: candidates.length - existing.length, }; } function writeManifest({ manifestPath, payload }) { writeFileSync(manifestPath, `${JSON.stringify(payload, null, 2)}\n`, 'utf8'); } function readManifest(manifestPath) { if (!existsSync(manifestPath)) { throw new Error(`备份清单不存在: ${manifestPath}`); } return JSON.parse(readFileSync(manifestPath, 'utf8')); } function hmac(key, content, encoding) { return createHmac('sha256', key).update(content).digest(encoding); } function sha256Hex(content) { return createHash('sha256').update(content).digest('hex'); } async function sha256FileHex(filePath) { const hash = createHash('sha256'); for await (const chunk of createReadStream(filePath)) { hash.update(chunk); } return hash.digest('hex'); } function directFilesStatePath({ workDir, database }) { return join( workDir, `${sanitizeObjectPart(database, 'spacetimedb')}-files-state.json.gz`, ); } function legacyDirectFilesStatePath({ workDir, database }) { return join( workDir, `${sanitizeObjectPart(database, 'spacetimedb')}-files-state.json`, ); } function directCatalogLocalPaths({ workDir, database, catalog }) { const baseName = `${sanitizeObjectPart(database, 'spacetimedb')}-${catalog.mode}-${catalog.catalogId}.catalog.json`; return { jsonPath: join(workDir, baseName), gzipPath: join(workDir, `${baseName}.gz`), }; } function readLocalDirectCatalog({ workDir, database, catalog }) { const { jsonPath, gzipPath } = directCatalogLocalPaths({ workDir, database, catalog, }); let body = null; if (existsSync(jsonPath)) { body = readFileSync(jsonPath); } else if (existsSync(gzipPath)) { body = gunzipSync(readFileSync(gzipPath)); } if (!body) { return null; } if ( body.length !== catalog.contentLength || sha256Hex(body) !== catalog.sha256 ) { throw new Error( `本地 files catalog 长度或 SHA 与 state 引用不匹配: ${jsonPath}`, ); } const payload = JSON.parse(body.toString('utf8')); if ( payload.schemaVersion !== DIRECT_FILES_CATALOG_SCHEMA_VERSION || payload.database !== database || payload.mode !== catalog.mode || payload.catalogId !== catalog.catalogId || !Array.isArray(payload.files) ) { throw new Error(`本地 files catalog 与 state 引用不匹配: ${jsonPath}`); } return payload; } function readJsonOrGzip(filePath) { const body = readFileSync(filePath); const decoded = filePath.endsWith('.gz') || (body[0] === 0x1f && body[1] === 0x8b) ? gunzipSync(body) : body; return JSON.parse(decoded.toString('utf8')); } function compactDirectCatalogFile({ workDir, database, catalog }) { const { jsonPath, gzipPath } = directCatalogLocalPaths({ workDir, database, catalog, }); if (!existsSync(jsonPath)) { return existsSync(gzipPath) ? { compressed: false, gzipPath } : null; } const body = readFileSync(jsonPath); atomicWriteBuffer(gzipPath, gzipSync(body, { level: 9 })); rmSync(jsonPath, { force: false }); return { compressed: true, gzipPath }; } function compactDirectFilesLocalMetadata({ workDir, database, nextState, transientCatalogPaths = [], }) { const keepCatalog = nextState.latestCatalog; const keepCatalogIds = new Set([keepCatalog?.catalogId].filter(Boolean)); const databasePart = sanitizeObjectPart(database, 'spacetimedb'); const catalogPattern = new RegExp( `^${databasePart.replace(/[.*+?^${}()|[\]\\]/gu, '\\$&')}-(full|history)-([a-f0-9]{64})\\.catalog\\.json(?:\\.gz)?$`, 'u', ); let compressedCatalogCount = 0; let deletedCatalogCount = 0; let compactedResultCount = 0; for (const entry of readdirSync(workDir, { withFileTypes: true })) { if (!entry.isFile()) { continue; } const match = catalogPattern.exec(entry.name); if (!match) { continue; } const catalog = { mode: match[1], catalogId: match[2] }; if (keepCatalogIds.has(catalog.catalogId) && catalog.mode === 'full') { readLocalDirectCatalog({ workDir, database, catalog: keepCatalog }); if ( !entry.name.endsWith('.gz') && compactDirectCatalogFile({ workDir, database, catalog })?.compressed ) { compressedCatalogCount += 1; } continue; } rmSync(join(workDir, entry.name), { force: false }); deletedCatalogCount += 1; } for (const filePath of transientCatalogPaths) { if (existsSync(filePath)) { rmSync(filePath, { force: false }); deletedCatalogCount += 1; } } for (const entry of readdirSync(workDir, { withFileTypes: true })) { if ( !entry.isFile() || !entry.name.endsWith('.json') || entry.name.endsWith('.catalog.json') ) { continue; } const filePath = join(workDir, entry.name); let payload = null; try { payload = JSON.parse(readFileSync(filePath, 'utf8')); } catch { continue; } if ( payload?.catalog?.schemaVersion !== DIRECT_FILES_CATALOG_SCHEMA_VERSION || payload.catalog.database !== database || payload.catalog.bucket !== nextState.bucket || !Array.isArray(payload.catalog.files) ) { continue; } atomicWriteJson(filePath, compactDirectFilesResult(payload)); compactedResultCount += 1; } const result = { compressedCatalogCount, deletedCatalogCount, compactedResultCount, }; console.log( `[database-backup] files 本地元数据清理: ${JSON.stringify(result)}`, ); return result; } function normalizeObjectPrefix(objectPrefix, database) { const prefix = String(objectPrefix || 'database-backups') .trim() .replace(/^\/+|\/+$/gu, '') .split('/') .filter(Boolean) .map((part) => sanitizeObjectPart(part, 'backup')) .join('/'); return [prefix, sanitizeObjectPart(database, 'spacetimedb')] .filter(Boolean) .join('/'); } function directFileIdentity(filePath) { const stat = lstatSync(filePath, { bigint: true }); if (!stat.isFile() || stat.isSymbolicLink()) { throw new Error(`files 模式只允许普通文件: ${filePath}`); } return { dev: stat.dev.toString(), ino: stat.ino.toString(), size: stat.size.toString(), mtimeNs: stat.mtimeNs.toString(), mode: Number(stat.mode & 0o7777n), }; } function sameDirectFileIdentity(left, right) { return ( left.dev === right.dev && left.ino === right.ino && left.size === right.size && left.mtimeNs === right.mtimeNs && left.mode === right.mode ); } export async function collectDirectFileEntries({ dataDir, candidates = null, objectPrefix, database, }) { const resolvedDataDir = resolvePath(dataDir); if ( !existsSync(resolvedDataDir) || !lstatSync(resolvedDataDir).isDirectory() ) { throw new Error(`files 数据目录不存在或不是目录: ${resolvedDataDir}`); } const files = new Map(); const symlinks = new Map(); const directories = new Set(['.']); const roots = candidates === null ? [{ absolutePath: resolvedDataDir, relativePath: '.' }] : candidates.map((candidate) => ({ absolutePath: resolve(resolvedDataDir, candidate.path), relativePath: assertSafeRelativePath( resolvedDataDir, resolve(resolvedDataDir, candidate.path), ), })); const visit = async (absolutePath, relativePath) => { const stat = lstatSync(absolutePath); if (stat.isSymbolicLink()) { const target = readlinkSync(absolutePath, 'utf8'); if (!target || isAbsolute(target)) { throw new Error( `files 模式只允许 data-dir 内部的相对符号链接: ${absolutePath} -> ${target}`, ); } assertSafeRelativePath( resolvedDataDir, resolve(dirname(absolutePath), target), ); symlinks.set(relativePath, { path: relativePath, target }); return; } if (stat.isDirectory()) { directories.add(relativePath); for (const name of readdirSync(absolutePath).sort()) { const childRelative = relativePath === '.' ? name : `${relativePath}/${name}`; await visit(join(absolutePath, name), childRelative); } return; } if (!stat.isFile()) { throw new Error(`files 模式只允许普通文件或目录: ${absolutePath}`); } const before = directFileIdentity(absolutePath); const sha256 = await sha256FileHex(absolutePath); const after = directFileIdentity(absolutePath); if (!sameDirectFileIdentity(before, after)) { throw new Error(`files 扫描期间源文件发生变化: ${relativePath}`); } const basePrefix = normalizeObjectPrefix(objectPrefix, database); const file = { path: relativePath, sizeBytes: Number(after.size), sha256, mode: after.mode, objectKey: `${basePrefix}/files/sha256/${sha256.slice(0, 2)}/${sha256}`, }; // 上传前后的 inode/stat 仍用于防止在线扫描漂移,但设为不可枚举,避免 // 把仅供本地校验的副本再次写入 catalog 或 result JSON。 Object.defineProperty(file, 'sourceStat', { value: after, enumerable: false, }); files.set(relativePath, file); }; for (const root of roots.sort((left, right) => left.relativePath.localeCompare(right.relativePath), )) { if (!existsSync(root.absolutePath)) { throw new Error(`files 候选在扫描前消失: ${root.relativePath}`); } await visit(root.absolutePath, root.relativePath); } return { directories: [...directories].sort(), files: [...files.values()].sort((left, right) => left.path.localeCompare(right.path), ), symlinks: [...symlinks.values()].sort((left, right) => left.path.localeCompare(right.path), ), }; } function directCatalogIdentity({ mode, baselineCatalogId, rootName, directories, files, symlinks, }) { // 不把数十万条文件元数据先拼成一个巨型 JSON 字符串;分段写入 hash // 保持与 JSON.stringify 同样的字段顺序和转义结果,同时把峰值降到单条记录。 const hash = createHash('sha256'); hash.update('{"mode":'); hash.update(JSON.stringify(mode)); hash.update(',"baselineCatalogId":'); hash.update(JSON.stringify(baselineCatalogId || '')); hash.update(',"rootName":'); hash.update(JSON.stringify(rootName)); hash.update(',"directories":'); updateJsonArrayHash(hash, directories, (directory) => JSON.stringify(directory), ); hash.update(',"files":'); updateJsonArrayHash(hash, files, (file) => JSON.stringify({ path: file.path, sizeBytes: file.sizeBytes, sha256: file.sha256, mode: file.mode, objectKey: file.objectKey, }), ); hash.update(',"symlinks":'); updateJsonArrayHash(hash, symlinks, (symlink) => JSON.stringify(symlink)); hash.update('}'); return hash.digest('hex'); } function updateJsonArrayHash(hash, values, serialize) { hash.update('['); values.forEach((value, index) => { if (index > 0) { hash.update(','); } hash.update(serialize(value)); }); hash.update(']'); } function readDirectFilesState(statePath, { database, bucket }) { const candidates = statePath.endsWith('.gz') ? [statePath, statePath.slice(0, -3)] : [statePath, `${statePath}.gz`]; const existingPath = candidates.find((candidate) => existsSync(candidate)); if (!existingPath) { return null; } const state = readJsonOrGzip(existingPath); if ( ![ LEGACY_DIRECT_FILES_STATE_SCHEMA_VERSION, DIRECT_FILES_STATE_SCHEMA_VERSION, ].includes(state.schemaVersion) || state.backupKind !== 'spacetimedb-direct-files-state' || state.database !== database || state.bucket !== bucket ) { throw new Error(`files state 与本次数据源或 bucket 不匹配: ${statePath}`); } return state; } function directPreviousFiles({ state, workDir, database }) { if (!state?.latestCatalog) { return []; } if (Array.isArray(state?.latestCatalog?.files)) { return state.latestCatalog.files; } return ( readLocalDirectCatalog({ workDir, database, catalog: state?.latestCatalog }) ?.files ?? [] ); } async function ensureDirectObject({ file, dataDir, uploadOptions, previousFile, verifyCatalogReuse = false, uploadFn, verifyFn, }) { const absolutePath = resolve(dataDir, file.path); assertSafeRelativePath(dataDir, absolutePath); if ( previousFile?.sha256 === file.sha256 && previousFile?.sizeBytes === file.sizeBytes && previousFile?.objectKey === file.objectKey ) { if (verifyCatalogReuse) { await verifyFn({ ...uploadOptions, objectKey: file.objectKey, contentLength: file.sizeBytes, archiveSha256: file.sha256, }); } return { status: verifyCatalogReuse ? 'catalog-reused-verified' : 'catalog-reused', objectKey: file.objectKey, }; } try { await verifyFn({ ...uploadOptions, objectKey: file.objectKey, contentLength: file.sizeBytes, archiveSha256: file.sha256, }); return { status: 'oss-reused', objectKey: file.objectKey }; } catch (error) { if (error?.status !== 404) { throw error; } } const beforeUpload = directFileIdentity(absolutePath); if (!sameDirectFileIdentity(beforeUpload, file.sourceStat)) { throw new Error(`files 上传前源文件 stat 漂移: ${file.path}`); } await uploadFn({ archivePath: absolutePath, ...uploadOptions, objectKey: file.objectKey, archiveSha256: file.sha256, backupKind: 'spacetimedb-direct-file', contentType: 'application/octet-stream', allowEmpty: true, }); const afterUpload = directFileIdentity(absolutePath); if (!sameDirectFileIdentity(afterUpload, file.sourceStat)) { throw new Error(`files 上传期间源文件 stat 漂移: ${file.path}`); } return { status: 'uploaded', objectKey: file.objectKey }; } async function ensureDirectManifest({ manifestPath, objectKey, uploadOptions, uploadManifestFn, verifyFn, }) { const body = readFileSync(manifestPath); const archiveSha256 = sha256Hex(body); try { const verification = await verifyFn({ ...uploadOptions, objectKey, contentLength: body.length, archiveSha256, }); return { objectKey, contentLength: body.length, archiveSha256, verifiedAt: verification.verifiedAt, reused: true, }; } catch (error) { if (error?.status !== 404) { throw error; } } return uploadManifestFn({ manifestPath, ...uploadOptions, objectKey }); } function directCatalogRef(catalog) { return { mode: catalog.mode, catalogId: catalog.catalogId, objectKey: catalog.objectKey, contentLength: catalog.contentLength, sha256: catalog.sha256, verifiedAt: catalog.verifiedAt, }; } function normalizeDirectFilesState({ state, dataDir, database, bucket }) { return { schemaVersion: DIRECT_FILES_STATE_SCHEMA_VERSION, backupKind: 'spacetimedb-direct-files-state', database, dataDir, bucket, updatedAt: new Date().toISOString(), baselineCatalog: assertDirectCatalogRef( state?.baselineCatalog, 'full', 'baseline full', ), latestCatalog: assertDirectCatalogRef( state?.latestCatalog, 'full', 'latest full', ), historyCatalogs: (state?.historyCatalogs ?? []).map((catalog) => assertDirectCatalogRef(catalog, 'history', 'history'), ), }; } function persistDirectFilesState({ statePath, legacyStatePath, state }) { atomicWriteGzipJson(statePath, state); if (legacyStatePath !== statePath && existsSync(legacyStatePath)) { rmSync(legacyStatePath, { force: false }); } } function compactDirectFilesResult(result) { if (!result.catalog) { return result; } const { catalog, ...rest } = result; return { ...rest, catalog: { schemaVersion: catalog.schemaVersion, backupKind: catalog.backupKind, database: catalog.database, bucket: catalog.bucket, mode: catalog.mode, catalogId: catalog.catalogId, catalogObjectKey: catalog.catalogObjectKey, baselineCatalogId: catalog.baselineCatalogId, rootName: catalog.rootName, fileCount: Array.isArray(catalog.files) ? catalog.files.length : 0, symlinkCount: Array.isArray(catalog.symlinks) ? catalog.symlinks.length : 0, }, }; } function assertDirectCatalogRef(catalog, expectedMode, label) { if ( !catalog || catalog.mode !== expectedMode || !/^[a-f0-9]{64}$/u.test(catalog.catalogId) || typeof catalog.objectKey !== 'string' || !catalog.objectKey || !Number.isSafeInteger(catalog.contentLength) || catalog.contentLength <= 0 || !/^[a-f0-9]{64}$/u.test(catalog.sha256) || typeof catalog.verifiedAt !== 'string' || !catalog.verifiedAt ) { throw new Error(`files ${label} catalog ref 无效。`); } return directCatalogRef(catalog); } function buildDirectFilesLatest({ database, bucket, state }) { const latestFullCatalog = assertDirectCatalogRef( state?.latestCatalog, 'full', 'latest full', ); const historyCatalogs = (state?.historyCatalogs ?? []).map((catalog) => assertDirectCatalogRef(catalog, 'history', 'history'), ); return { schemaVersion: DIRECT_FILES_LATEST_SCHEMA_VERSION, backupKind: 'spacetimedb-direct-files-latest', database, bucket, updatedAt: new Date().toISOString(), latestFullCatalog, historyCatalogs, }; } function validateDirectFilesLatest(latest, { database, bucket }) { if ( latest?.schemaVersion !== DIRECT_FILES_LATEST_SCHEMA_VERSION || latest.backupKind !== 'spacetimedb-direct-files-latest' || latest.database !== database || latest.bucket !== bucket || !Array.isArray(latest.historyCatalogs) ) { throw new Error('files latest pointer 契约无效。'); } return { ...latest, latestFullCatalog: assertDirectCatalogRef( latest.latestFullCatalog, 'full', 'latest full', ), historyCatalogs: latest.historyCatalogs.map((catalog) => assertDirectCatalogRef(catalog, 'history', 'history'), ), }; } async function publishDirectFilesLatest({ workDir, database, bucket, objectPrefix, state, uploadOptions, uploadManifestFn, verifyFn, }) { const latest = buildDirectFilesLatest({ database, bucket, state }); for (const catalogRef of [ latest.latestFullCatalog, ...latest.historyCatalogs, ]) { await verifyFn({ ...uploadOptions, objectKey: catalogRef.objectKey, contentLength: catalogRef.contentLength, archiveSha256: catalogRef.sha256, }); } const latestPath = join( workDir, `${sanitizeObjectPart(database, 'spacetimedb')}-latest.json`, ); const latestObjectKey = `${normalizeObjectPrefix(objectPrefix, database)}/latest.json`; writeManifest({ manifestPath: latestPath, payload: latest }); const uploaded = await uploadManifestFn({ manifestPath: latestPath, ...uploadOptions, objectKey: latestObjectKey, }); const verification = await verifyFn({ ...uploadOptions, objectKey: latestObjectKey, contentLength: uploaded.contentLength, archiveSha256: uploaded.archiveSha256, }); return { latest, latestPath, latestObjectKey, contentLength: uploaded.contentLength, sha256: uploaded.archiveSha256, verifiedAt: verification.verifiedAt, }; } export async function runDirectFilesBackup({ mode, dataDir, workDir, database, bucket, objectPrefix, dryRun = false, resultFile = '', uploadOptions, uploadFn = uploadDirectFile, uploadManifestFn = uploadManifestFile, verifyFn = verifyOssObject, concurrency = 1, }) { mkdirSync(workDir, { recursive: true }); const statePath = directFilesStatePath({ workDir, database }); const legacyStatePath = legacyDirectFilesStatePath({ workDir, database }); const state = readDirectFilesState(statePath, { database, bucket }); if ( mode === 'history' && (!state?.baselineCatalog || state?.latestCatalog?.mode !== 'full') ) { throw new Error( `files history 模式缺少已发布 full baseline catalog: ${statePath}`, ); } const plan = mode === 'history' ? discoverHistoryPlan({ dataDir }) : null; const collected = await collectDirectFileEntries({ dataDir, candidates: plan?.candidates ?? null, objectPrefix, database, }); const baselineCatalogId = mode === 'history' ? (state?.baselineCatalog?.catalogId ?? '') : ''; const rootName = basename(dataDir); const catalogId = directCatalogIdentity({ mode, baselineCatalogId, rootName, ...collected, }); const basePrefix = normalizeObjectPrefix(objectPrefix, database); const catalogObjectKey = `${basePrefix}/catalogs/${mode}/${catalogId}.json`; const catalogPath = join( workDir, `${sanitizeObjectPart(database, 'spacetimedb')}-${mode}-${catalogId}.catalog.json`, ); const catalog = { schemaVersion: DIRECT_FILES_CATALOG_SCHEMA_VERSION, backupKind: mode === 'full' ? 'spacetimedb-data-dir-files' : 'spacetimedb-history-files', database, bucket, mode, catalogId, catalogObjectKey, baselineCatalogId, rootName, directories: collected.directories, files: collected.files, symlinks: collected.symlinks, }; writeManifest({ manifestPath: catalogPath, payload: catalog }); const summary = { statePath, catalogPath, catalogObjectKey, catalogId, fileCount: collected.files.length, symlinkCount: collected.symlinks.length, totalSizeBytes: collected.files .reduce((sum, file) => sum + BigInt(file.sizeBytes), 0n) .toString(), candidateCount: plan?.candidates.length ?? 0, }; if (resultFile) { atomicWriteJson(resolvePath(resultFile), { ...summary, dryRun }); } console.log( `[database-backup] files ${mode}: files=${summary.fileCount}, symlinks=${summary.symlinkCount}, size=${formatBytes(summary.totalSizeBytes)}, catalog=${catalogId}`, ); if (dryRun) { console.log( '[database-backup] files dry-run,仅扫描并生成本地 catalog,不上传或删除。', ); return { ...summary, catalog, uploadedCount: 0, reusedCount: 0 }; } if (mode === 'history' && plan.candidates.length === 0) { await verifyFn({ ...uploadOptions, objectKey: state.latestCatalog.objectKey, contentLength: state.latestCatalog.contentLength, archiveSha256: state.latestCatalog.sha256, }); const latestPointer = await publishDirectFilesLatest({ workDir, database, bucket, objectPrefix, state, uploadOptions, uploadManifestFn, verifyFn, }); const compactedState = normalizeDirectFilesState({ state, dataDir, database, bucket, }); persistDirectFilesState({ statePath, legacyStatePath, state: compactedState, }); const metadataCleanup = compactDirectFilesLocalMetadata({ workDir, database, nextState: compactedState, transientCatalogPaths: [catalogPath], }); console.log('[database-backup] files history 没有可归档候选。'); const emptyResult = { ...summary, catalog, latestPointer, uploadedCount: 0, reusedCount: 0, cleanup: null, metadataCleanup, }; if (resultFile) { atomicWriteJson( resolvePath(resultFile), compactDirectFilesResult(emptyResult), ); } return emptyResult; } if ( state?.latestCatalog?.catalogId === catalogId && state.latestCatalog.mode === mode ) { await verifyFn({ ...uploadOptions, objectKey: state.latestCatalog.objectKey, contentLength: state.latestCatalog.contentLength, archiveSha256: state.latestCatalog.sha256, }); if (mode === 'full') { const latestPointer = await publishDirectFilesLatest({ workDir, database, bucket, objectPrefix, state, uploadOptions, uploadManifestFn, verifyFn, }); const compactedState = normalizeDirectFilesState({ state, dataDir, database, bucket, }); persistDirectFilesState({ statePath, legacyStatePath, state: compactedState, }); const metadataCleanup = compactDirectFilesLocalMetadata({ workDir, database, nextState: compactedState, }); console.log('[database-backup] files catalog 未变化,无文件需要上传。'); return { ...summary, catalog, latestPointer, uploadedCount: 0, reusedCount: collected.files.length, unchanged: true, metadataCleanup, }; } } if (state?.latestCatalog) { await verifyFn({ ...uploadOptions, objectKey: state.latestCatalog.objectKey, contentLength: state.latestCatalog.contentLength, archiveSha256: state.latestCatalog.sha256, }); } const previousFiles = new Map( directPreviousFiles({ state, workDir, database }).map((file) => [ file.path, file, ]), ); let uploadedCount = 0; let reusedCount = 0; let nextIndex = 0; let completedCount = 0; const workerCount = Math.min(concurrency, collected.files.length); const workers = Array.from({ length: workerCount }, async () => { while (nextIndex < collected.files.length) { const index = nextIndex; nextIndex += 1; const file = collected.files[index]; const result = await ensureDirectObject({ file, dataDir, uploadOptions, previousFile: previousFiles.get(file.path), verifyCatalogReuse: mode === 'history', uploadFn, verifyFn, }); if (result.status === 'uploaded') { uploadedCount += 1; } else { reusedCount += 1; } completedCount += 1; if ( collected.files.length <= 100 || completedCount % 1000 === 0 || completedCount === collected.files.length ) { console.log( `[database-backup] files 进度: ${completedCount}/${collected.files.length} (${result.status}) ${file.path}`, ); } } }); await Promise.all(workers); const catalogUpload = await ensureDirectManifest({ manifestPath: catalogPath, objectKey: catalogObjectKey, uploadOptions, uploadManifestFn, verifyFn, }); await verifyFn({ ...uploadOptions, objectKey: catalogObjectKey, contentLength: catalogUpload.contentLength, archiveSha256: catalogUpload.archiveSha256, }); if (mode === 'history') { await verifyFn({ ...uploadOptions, objectKey: state.baselineCatalog.objectKey, contentLength: state.baselineCatalog.contentLength, archiveSha256: state.baselineCatalog.sha256, }); } const catalogRef = { mode, catalogId, objectKey: catalogObjectKey, contentLength: catalogUpload.contentLength, sha256: catalogUpload.archiveSha256, verifiedAt: catalogUpload.verifiedAt, }; const nextState = { schemaVersion: DIRECT_FILES_STATE_SCHEMA_VERSION, backupKind: 'spacetimedb-direct-files-state', database, dataDir, bucket, updatedAt: new Date().toISOString(), baselineCatalog: directCatalogRef(state?.baselineCatalog ?? catalogRef), latestCatalog: directCatalogRef( mode === 'full' ? catalogRef : state.latestCatalog, ), historyCatalogs: mode === 'history' ? [ ...(state.historyCatalogs ?? []) .filter((item) => item.catalogId !== catalogId) .map((item) => directCatalogRef(item)), directCatalogRef(catalogRef), ] : (state?.historyCatalogs ?? []).map((item) => directCatalogRef(item)), }; const latestPointer = await publishDirectFilesLatest({ workDir, database, bucket, objectPrefix, state: nextState, uploadOptions, uploadManifestFn, verifyFn, }); persistDirectFilesState({ statePath, legacyStatePath, state: nextState }); const metadataCleanup = compactDirectFilesLocalMetadata({ workDir, database, nextState, }); let cleanup = null; if (mode === 'history') { cleanup = cleanupHistoryCandidates({ dataDir, candidates: plan.candidates, }); } const finalResult = { ...summary, catalog, latestPointer, uploadedCount, reusedCount, cleanup, metadataCleanup, }; if (resultFile) { atomicWriteJson( resolvePath(resultFile), compactDirectFilesResult(finalResult), ); } return finalResult; } async function downloadOssBuffer({ objectKey, uploadOptions }) { const response = await signedOssRequest({ ...ossRequestDefaults(uploadOptions), method: 'GET', objectKey, operation: '下载对象', }); return Buffer.from(await response.arrayBuffer()); } async function downloadOssFile({ objectKey, destinationPath, uploadOptions }) { const response = await signedOssRequest({ ...ossRequestDefaults(uploadOptions), method: 'GET', objectKey, operation: '下载对象', }); const tempPath = `${destinationPath}.partial-${process.pid}`; rmSync(tempPath, { force: true }); try { if (response.body) { await pipeline( Readable.fromWeb(response.body), createWriteStream(tempPath, { mode: 0o600 }), ); } else { writeFileSync(tempPath, Buffer.alloc(0), { mode: 0o600 }); } renameSync(tempPath, destinationPath); } catch (error) { rmSync(tempPath, { force: true }); throw error; } } async function loadDirectFilesCatalog({ catalogRef, database, bucket, uploadOptions, downloadBufferFn, }) { const catalogBody = await downloadBufferFn({ objectKey: catalogRef.objectKey, uploadOptions, }); if ( catalogBody.length !== catalogRef.contentLength || sha256Hex(catalogBody) !== catalogRef.sha256 ) { throw new Error( `files restore catalog 长度或 SHA-256 不一致: ${catalogRef.objectKey}`, ); } const catalog = JSON.parse(catalogBody.toString('utf8')); if ( catalog.schemaVersion !== DIRECT_FILES_CATALOG_SCHEMA_VERSION || catalog.backupKind !== 'spacetimedb-data-dir-files' || catalog.database !== database || catalog.bucket !== bucket || catalog.catalogId !== catalogRef.catalogId || !Array.isArray(catalog.directories) || !Array.isArray(catalog.files) ) { throw new Error(`files restore catalog 契约无效: ${catalogRef.objectKey}`); } return { ...catalog, symlinks: catalog.symlinks ?? [] }; } function assertDirectCatalogFile(file, index) { if ( !file || typeof file.path !== 'string' || !Number.isSafeInteger(file.sizeBytes) || file.sizeBytes < 0 || !/^[a-f0-9]{64}$/u.test(file.sha256) || typeof file.objectKey !== 'string' || !file.objectKey || !Number.isSafeInteger(file.mode) ) { throw new Error(`files restore catalog 文件项无效: index=${index}`); } } function assertDirectCatalogSymlink(symlink, index, restoreDir) { if ( !symlink || typeof symlink.path !== 'string' || !symlink.path || typeof symlink.target !== 'string' || !symlink.target || isAbsolute(symlink.target) ) { throw new Error(`files restore catalog 符号链接项无效: index=${index}`); } const destinationPath = resolve(restoreDir, symlink.path); assertSafeRelativePath(restoreDir, destinationPath); assertSafeRelativePath( restoreDir, resolve(dirname(destinationPath), symlink.target), ); } async function restoreDirectFilesCatalog({ catalog, restoreDir, uploadOptions, resultFile = '', dryRun = false, downloadFileFn = downloadOssFile, }) { const resolvedRestoreDir = resolvePath(restoreDir); catalog.files.forEach(assertDirectCatalogFile); catalog.symlinks.forEach((symlink, index) => assertDirectCatalogSymlink(symlink, index, resolvedRestoreDir), ); const totalSizeBytes = catalog.files .reduce((sum, file) => sum + BigInt(file.sizeBytes), 0n) .toString(); if (dryRun) { const result = { restoreDir: resolvedRestoreDir, catalogId: catalog.catalogId, fileCount: catalog.files.length, symlinkCount: catalog.symlinks.length, totalSizeBytes, downloadedCount: 0, reusedCount: 0, dryRun: true, }; if (resultFile) { atomicWriteJson(resolvePath(resultFile), result); } return result; } mkdirSync(resolvedRestoreDir, { recursive: true, mode: 0o700 }); for (const directoryPath of catalog.directories) { if (directoryPath === '.') { continue; } const absolutePath = resolve(resolvedRestoreDir, directoryPath); assertSafeRelativePath(resolvedRestoreDir, absolutePath); mkdirSync(absolutePath, { recursive: true }); } let downloadedCount = 0; let reusedCount = 0; for (const [index, file] of catalog.files.entries()) { const destinationPath = resolve(resolvedRestoreDir, file.path); assertSafeRelativePath(resolvedRestoreDir, destinationPath); mkdirSync(dirname(destinationPath), { recursive: true }); let reusable = false; if (existsSync(destinationPath) && lstatSync(destinationPath).isFile()) { const stat = statSync(destinationPath); reusable = stat.size === file.sizeBytes && (await sha256FileHex(destinationPath)) === file.sha256; } if (reusable) { reusedCount += 1; } else { rmSync(destinationPath, { force: true }); await downloadFileFn({ objectKey: file.objectKey, destinationPath, uploadOptions, }); const stat = statSync(destinationPath); const sha256 = await sha256FileHex(destinationPath); if (stat.size !== file.sizeBytes || sha256 !== file.sha256) { rmSync(destinationPath, { force: true }); throw new Error( `files restore 对象长度或 SHA-256 不一致: ${file.path}`, ); } downloadedCount += 1; } chmodSync(destinationPath, file.mode & 0o7777); console.log( `[database-backup] files restore: ${index + 1}/${catalog.files.length} (${reusable ? 'reused' : 'downloaded'}) ${file.path}`, ); } for (const symlink of catalog.symlinks) { const destinationPath = resolve(resolvedRestoreDir, symlink.path); assertSafeRelativePath(resolvedRestoreDir, destinationPath); mkdirSync(dirname(destinationPath), { recursive: true }); rmSync(destinationPath, { recursive: true, force: true }); symlinkSync(symlink.target, destinationPath); console.log( `[database-backup] files restore: symlink ${symlink.path} -> ${symlink.target}`, ); } const result = { restoreDir: resolvedRestoreDir, catalogId: catalog.catalogId, fileCount: catalog.files.length, symlinkCount: catalog.symlinks.length, totalSizeBytes, downloadedCount, reusedCount, }; if (resultFile) { atomicWriteJson(resolvePath(resultFile), result); } return result; } export async function restoreDirectFilesBackup({ statePath, restoreDir, database, bucket, uploadOptions, resultFile = '', dryRun = false, downloadBufferFn = downloadOssBuffer, downloadFileFn = downloadOssFile, }) { const state = readDirectFilesState(resolvePath(statePath), { database, bucket, }); if (!state?.latestCatalog || state.latestCatalog.mode !== 'full') { throw new Error(`files restore 缺少 full baseline catalog: ${statePath}`); } const catalogRef = assertDirectCatalogRef( state.latestCatalog, 'full', 'latest full', ); const catalog = await loadDirectFilesCatalog({ catalogRef, database, bucket, uploadOptions, downloadBufferFn, }); return restoreDirectFilesCatalog({ catalog, restoreDir, uploadOptions, resultFile, dryRun, downloadFileFn, }); } export async function restoreDirectFilesLatest({ restoreDir, database, bucket, objectPrefix, uploadOptions, resultFile = '', dryRun = false, downloadBufferFn = downloadOssBuffer, downloadFileFn = downloadOssFile, verifyFn = verifyOssObject, }) { const latestObjectKey = `${normalizeObjectPrefix(objectPrefix, database)}/latest.json`; const latestBody = await downloadBufferFn({ objectKey: latestObjectKey, uploadOptions, }); const latestSha256 = sha256Hex(latestBody); await verifyFn({ ...uploadOptions, objectKey: latestObjectKey, contentLength: latestBody.length, archiveSha256: latestSha256, }); const latest = validateDirectFilesLatest( JSON.parse(latestBody.toString('utf8')), { database, bucket }, ); const catalogRef = latest.latestFullCatalog; await verifyFn({ ...uploadOptions, objectKey: catalogRef.objectKey, contentLength: catalogRef.contentLength, archiveSha256: catalogRef.sha256, }); const catalog = await loadDirectFilesCatalog({ catalogRef, database, bucket, uploadOptions, downloadBufferFn, }); return restoreDirectFilesCatalog({ catalog, restoreDir, uploadOptions, resultFile, dryRun, downloadFileFn, }); } function regionFromEndpoint(endpoint) { const match = /^oss-([a-z0-9-]+)\./u.exec(endpoint); if (!match) { throw new Error(`无法从 OSS endpoint 推断 region: ${endpoint}`); } return match[1]; } function formatScopeDate(date) { return timestampForFile(date).slice(0, 8); } function formatOssDate(date) { return timestampForFile(date).replace(/[-:]/gu, ''); } function encodePath(path) { return path .split('/') .map((segment) => encodeURIComponent(segment).replace( /[!'()*]/gu, (char) => `%${char.charCodeAt(0).toString(16).toUpperCase()}`, ), ) .join('/'); } function encodeQueryComponent(value) { return encodeURIComponent(String(value)).replace( /[!'()*]/gu, (char) => `%${char.charCodeAt(0).toString(16).toUpperCase()}`, ); } export function buildCanonicalQuery(queries = {}) { return Object.entries(queries) .map(([key, value]) => [ encodeQueryComponent(key), value === null ? null : encodeQueryComponent(value), ]) .sort(([leftKey, leftValue], [rightKey, rightValue]) => { if (leftKey !== rightKey) { return leftKey < rightKey ? -1 : 1; } const left = leftValue ?? ''; const right = rightValue ?? ''; return left === right ? 0 : left < right ? -1 : 1; }) .map(([key, value]) => (value === null ? key : `${key}=${value}`)) .join('&'); } function canonicalHeaderValue(value) { return String(value).trim().replace(/\s+/gu, ' '); } export function buildAuthorization({ method, bucket, endpoint, objectKey, accessKeyId, accessKeySecret, headers, date, queries = {}, }) { const region = regionFromEndpoint(endpoint); const scopeDate = formatScopeDate(date); const scope = `${scopeDate}/${region}/${OSS_SERVICE}/${OSS_REQUEST}`; const canonicalUri = `/${encodeURIComponent(bucket)}/${encodePath(objectKey)}`; const signedHeaders = Object.fromEntries( Object.entries(headers).map(([key, value]) => [ key.toLowerCase(), canonicalHeaderValue(value), ]), ); const canonicalHeaders = Object.entries(signedHeaders) .sort(([left], [right]) => left.localeCompare(right)) .map(([key, value]) => `${key}:${value}\n`) .join(''); const additionalHeaders = 'host'; const canonicalRequest = [ method, canonicalUri, buildCanonicalQuery(queries), canonicalHeaders, additionalHeaders, UNSIGNED_PAYLOAD, ].join('\n'); const stringToSign = [ OSS_ALGORITHM, headers['x-oss-date'], scope, sha256Hex(canonicalRequest), ].join('\n'); const signature = hmac( Buffer.from(`aliyun_v4${accessKeySecret}`, 'utf8'), scopeDate, ); const regionKey = hmac(signature, region); const serviceKey = hmac(regionKey, OSS_SERVICE); const signingKey = hmac(serviceKey, OSS_REQUEST); const finalSignature = hmac(signingKey, stringToSign, 'hex'); return `${OSS_ALGORITHM} Credential=${accessKeyId}/${scope},AdditionalHeaders=${additionalHeaders},Signature=${finalSignature}`; } function buildOssUrl({ bucket, endpoint, objectKey, queries = {} }) { const canonicalQuery = buildCanonicalQuery(queries); return `https://${bucket}.${endpoint}/${encodePath(objectKey)}${canonicalQuery ? `?${canonicalQuery}` : ''}`; } function isRetryableOssStatus(status) { return RETRYABLE_OSS_HTTP_STATUSES.has(status); } function retryDelayMs({ attempt, baseDelayMs, maxDelayMs, randomFn }) { const ceiling = Math.min( maxDelayMs, baseDelayMs * 2 ** Math.max(0, attempt - 1), ); return Math.floor(randomFn() * ceiling); } function ossRequestDefaults({ bucket, endpoint, accessKeyId, accessKeySecret, fetchImpl = globalThis.fetch, nowFn = () => new Date(), sleepImpl = sleep, randomFn = Math.random, maxAttempts = DEFAULT_OSS_REQUEST_MAX_ATTEMPTS, retryBaseDelayMs = DEFAULT_OSS_RETRY_BASE_DELAY_MS, retryMaxDelayMs = DEFAULT_OSS_RETRY_MAX_DELAY_MS, }) { return { bucket, endpoint, accessKeyId, accessKeySecret, fetchImpl, nowFn, sleepImpl, randomFn, maxAttempts, retryBaseDelayMs, retryMaxDelayMs, }; } async function signedOssRequest({ method, bucket, endpoint, objectKey, accessKeyId, accessKeySecret, queries = {}, headers = {}, bodyFactory, contentLength, operation, fetchImpl, nowFn, sleepImpl, randomFn, maxAttempts, retryBaseDelayMs, retryMaxDelayMs, }) { const targetUrl = buildOssUrl({ bucket, endpoint, objectKey, queries }); let lastError = null; for (let attempt = 1; attempt <= maxAttempts; attempt += 1) { const now = nowFn(); const signedHeaders = { host: `${bucket}.${endpoint}`, ...headers, 'x-oss-content-sha256': UNSIGNED_PAYLOAD, 'x-oss-date': formatOssDate(now), }; const authorization = buildAuthorization({ method, bucket, endpoint, objectKey, accessKeyId, accessKeySecret, headers: signedHeaders, date: now, queries, }); const requestHeaders = { ...signedHeaders, authorization }; if (contentLength !== undefined) { requestHeaders['content-length'] = String(contentLength); } const body = bodyFactory ? bodyFactory() : undefined; const requestOptions = { method, headers: requestHeaders }; if (body !== undefined) { requestOptions.body = body; requestOptions.duplex = 'half'; } let response; try { response = await fetchImpl(targetUrl, requestOptions); } catch (error) { lastError = new Error( `OSS ${operation}请求失败: oss://${bucket}/${objectKey}`, { cause: error }, ); } if (response?.ok) { return response; } if (response) { const responseText = await response.text(); const requestId = response.headers.get('x-oss-request-id'); lastError = new Error( `OSS ${operation}失败 HTTP ${response.status}${requestId ? ` requestId=${requestId}` : ''}: ${responseText.slice(0, 500)}`, ); lastError.status = response.status; } const retryable = !response || isRetryableOssStatus(response.status); if (!retryable || attempt >= maxAttempts) { throw lastError; } const delayMs = retryDelayMs({ attempt, baseDelayMs: retryBaseDelayMs, maxDelayMs: retryMaxDelayMs, randomFn, }); console.warn( `[database-backup] OSS ${operation}失败,${delayMs}ms 后重试 (${attempt}/${maxAttempts})`, ); await sleepImpl(delayMs); } throw ( lastError ?? new Error(`OSS ${operation}失败: oss://${bucket}/${objectKey}`) ); } function readXmlTag(xml, tagName) { const match = new RegExp(`<${tagName}>([\\s\\S]*?)`, 'u').exec( xml, ); if (!match) { return ''; } return match[1] .replace(/</gu, '<') .replace(/>/gu, '>') .replace(/"/gu, '"') .replace(/'/gu, "'") .replace(/&/gu, '&') .trim(); } function escapeXml(value) { return String(value) .replace(/&/gu, '&') .replace(//gu, '>') .replace(/"/gu, '"') .replace(/'/gu, '''); } function buildCompleteMultipartBody(parts) { const partXml = parts .map(({ partNumber, etag }) => [ '', `${partNumber}`, `${escapeXml(etag)}`, '', ].join(''), ) .join(''); return `${partXml}`; } function resolveMultipartPartSize(fileSize, configuredPartSize) { if ( !Number.isSafeInteger(configuredPartSize) || configuredPartSize < OSS_MIN_MULTIPART_PART_SIZE_BYTES ) { throw new Error( `OSS multipart part size 必须是 >= ${OSS_MIN_MULTIPART_PART_SIZE_BYTES} 的安全整数`, ); } const minimumForPartLimit = Math.ceil(fileSize / OSS_MAX_MULTIPART_PARTS); const partSize = Math.max(configuredPartSize, minimumForPartLimit); if (partSize > OSS_MAX_MULTIPART_PART_SIZE_BYTES) { throw new Error(`OSS multipart part size 超过 5GiB: ${partSize}`); } return partSize; } async function verifyUploadedObject({ requestOptions, expectedContentLength, expectedArchiveSha256, }) { const response = await signedOssRequest({ ...requestOptions, method: 'HEAD', operation: 'HEAD 验证', }); const contentLengthHeader = response.headers.get('content-length'); const metadataLengthHeader = response.headers.get('x-oss-meta-file-size'); const effectiveLengthHeader = /^\d+$/u.test(contentLengthHeader ?? '') ? contentLengthHeader : metadataLengthHeader; if (!effectiveLengthHeader || !/^\d+$/u.test(effectiveLengthHeader)) { throw new Error( `OSS HEAD 验证缺少有效 content-length/file-size: content-length=${contentLengthHeader ?? ''}, file-size=${metadataLengthHeader ?? ''}`, ); } const remoteContentLength = Number(effectiveLengthHeader); if (remoteContentLength !== expectedContentLength) { throw new Error( `OSS HEAD 验证长度不一致: local=${expectedContentLength}, remote=${remoteContentLength}`, ); } const remoteArchiveSha256 = String( response.headers.get('x-oss-meta-file-sha256') ?? response.headers.get('x-oss-meta-archive-sha256') ?? '', ) .trim() .toLowerCase(); if (remoteArchiveSha256 !== expectedArchiveSha256) { throw new Error( `OSS HEAD 验证 SHA-256 不一致: local=${expectedArchiveSha256}, remote=${remoteArchiveSha256 || ''}`, ); } return { verifiedAt: new Date().toISOString(), remoteContentLength, remoteArchiveSha256, }; } export async function verifyOssObject({ bucket, endpoint, objectKey, accessKeyId, accessKeySecret, contentLength, archiveSha256, fetchImpl = globalThis.fetch, nowFn = () => new Date(), sleepImpl = sleep, randomFn = Math.random, maxAttempts = DEFAULT_OSS_REQUEST_MAX_ATTEMPTS, }) { const requestOptions = { bucket, endpoint, objectKey, accessKeyId, accessKeySecret, fetchImpl, nowFn, sleepImpl, randomFn, maxAttempts, retryBaseDelayMs: DEFAULT_OSS_RETRY_BASE_DELAY_MS, retryMaxDelayMs: DEFAULT_OSS_RETRY_MAX_DELAY_MS, }; return verifyUploadedObject({ requestOptions, expectedContentLength: Number(contentLength), expectedArchiveSha256: String(archiveSha256 ?? '') .trim() .toLowerCase(), }); } async function abortMultipartUpload({ requestOptions, uploadId }) { try { await signedOssRequest({ ...requestOptions, method: 'DELETE', queries: { uploadId }, operation: 'AbortMultipartUpload', maxAttempts: Math.min(2, requestOptions.maxAttempts), }); console.warn( `[database-backup] 已清理失败的 multipart upload: ${uploadId}`, ); } catch (error) { console.warn( `[database-backup] 清理 multipart upload 失败: ${error.message}`, ); } } export async function uploadArchive({ archivePath, bucket, endpoint, objectKey, accessKeyId, accessKeySecret, partSizeBytes = DEFAULT_OSS_MULTIPART_PART_SIZE_BYTES, maxAttempts = DEFAULT_OSS_REQUEST_MAX_ATTEMPTS, retryBaseDelayMs = DEFAULT_OSS_RETRY_BASE_DELAY_MS, retryMaxDelayMs = DEFAULT_OSS_RETRY_MAX_DELAY_MS, fetchImpl = globalThis.fetch, nowFn = () => new Date(), sleepImpl = sleep, randomFn = Math.random, backupKind = 'spacetimedb-data-dir', archiveSha256 = '', contentType = 'application/gzip', allowEmpty = false, bandwidthLimiter = null, }) { const fileStat = statSync(archivePath); if (!fileStat.isFile() || (!allowEmpty && fileStat.size <= 0)) { throw new Error( `待上传备份必须是${allowEmpty ? '' : '非空'}普通文件: ${archivePath}`, ); } const verifiedArchiveSha256 = archiveSha256 || (await sha256FileHex(archivePath)); if (!/^[a-f0-9]{64}$/u.test(verifiedArchiveSha256)) { throw new Error(`归档 SHA-256 无效: ${verifiedArchiveSha256}`); } const requestOptions = { bucket, endpoint, objectKey, accessKeyId, accessKeySecret, fetchImpl, nowFn, sleepImpl, randomFn, maxAttempts, retryBaseDelayMs, retryMaxDelayMs, }; if (fileStat.size === 0) { await signedOssRequest({ ...requestOptions, method: 'PUT', headers: { 'content-type': contentType, 'x-oss-meta-archive-sha256': verifiedArchiveSha256, 'x-oss-meta-file-sha256': verifiedArchiveSha256, 'x-oss-meta-file-size': '0', 'x-oss-meta-backup-kind': backupKind, }, contentLength: 0, bodyFactory: () => Buffer.alloc(0), operation: '上传空文件', }); const verification = await verifyUploadedObject({ requestOptions, expectedContentLength: 0, expectedArchiveSha256: verifiedArchiveSha256, }); return { bucket, objectKey, contentLength: 0, archiveSha256: verifiedArchiveSha256, etag: '', uploadMode: 'single', partCount: 1, partSizeBytes: 0, verifiedAt: verification.verifiedAt, }; } const partSize = resolveMultipartPartSize(fileStat.size, partSizeBytes); const partCount = Math.ceil(fileStat.size / partSize); let uploadId = ''; let uploadCompleted = false; console.log( `[database-backup] multipart 上传 OSS: oss://${bucket}/${objectKey} (${partCount} parts)`, ); try { const initiateResponse = await signedOssRequest({ ...requestOptions, method: 'POST', queries: { uploads: null }, headers: { 'content-type': contentType, 'x-oss-meta-archive-sha256': verifiedArchiveSha256, 'x-oss-meta-file-sha256': verifiedArchiveSha256, 'x-oss-meta-file-size': String(fileStat.size), 'x-oss-meta-backup-kind': backupKind, }, operation: 'InitiateMultipartUpload', }); uploadId = readXmlTag(await initiateResponse.text(), 'UploadId'); if (!uploadId) { throw new Error('OSS InitiateMultipartUpload 响应缺少 UploadId'); } const parts = []; for (let partNumber = 1; partNumber <= partCount; partNumber += 1) { const start = (partNumber - 1) * partSize; const end = Math.min(fileStat.size, start + partSize) - 1; const contentLength = end - start + 1; const response = await signedOssRequest({ ...requestOptions, method: 'PUT', queries: { partNumber, uploadId }, headers: { 'content-type': 'application/octet-stream' }, contentLength, bodyFactory: () => { const stream = createReadStream(archivePath, { start, end }); return bandwidthLimiter ? bandwidthLimiter.wrap(stream) : stream; }, operation: `UploadPart ${partNumber}/${partCount}`, }); const etag = response.headers.get('etag'); if (!etag) { throw new Error( `OSS UploadPart ${partNumber}/${partCount} 响应缺少 ETag`, ); } parts.push({ partNumber, etag }); console.log( `[database-backup] multipart 进度: ${partNumber}/${partCount}`, ); } const completeBody = buildCompleteMultipartBody(parts); let completeResponse; try { completeResponse = await signedOssRequest({ ...requestOptions, method: 'POST', queries: { uploadId }, headers: { 'content-type': 'application/xml' }, contentLength: Buffer.byteLength(completeBody), bodyFactory: () => completeBody, operation: 'CompleteMultipartUpload', }); const completeResponseText = await completeResponse.text(); if (/)/u.test(completeResponseText)) { throw new Error( `OSS CompleteMultipartUpload 返回错误: ${completeResponseText.slice(0, 500)}`, ); } } catch (completeError) { try { await verifyUploadedObject({ requestOptions, expectedContentLength: fileStat.size, expectedArchiveSha256: verifiedArchiveSha256, }); completeResponse = null; } catch { throw completeError; } } const verification = await verifyUploadedObject({ requestOptions, expectedContentLength: fileStat.size, expectedArchiveSha256: verifiedArchiveSha256, }); uploadCompleted = true; return { bucket, objectKey, contentLength: fileStat.size, archiveSha256: verifiedArchiveSha256, etag: completeResponse?.headers.get('etag')?.replace(/^"|"$/gu, '') ?? '', uploadMode: 'multipart', partCount, partSizeBytes: partSize, verifiedAt: verification.verifiedAt, }; } catch (error) { if (uploadId && !uploadCompleted) { await abortMultipartUpload({ requestOptions, uploadId }); } throw error; } } export async function uploadDirectFile({ archivePath, bucket, endpoint, objectKey, accessKeyId, accessKeySecret, fetchImpl = globalThis.fetch, nowFn = () => new Date(), sleepImpl = sleep, randomFn = Math.random, maxAttempts = DEFAULT_OSS_REQUEST_MAX_ATTEMPTS, retryBaseDelayMs = DEFAULT_OSS_RETRY_BASE_DELAY_MS, retryMaxDelayMs = DEFAULT_OSS_RETRY_MAX_DELAY_MS, backupKind = 'spacetimedb-direct-file', archiveSha256 = '', contentType = 'application/octet-stream', allowEmpty = true, bandwidthLimiter = null, }) { const fileStat = statSync(archivePath); if (!fileStat.isFile() || (!allowEmpty && fileStat.size <= 0)) { throw new Error( `待上传备份必须是${allowEmpty ? '' : '非空'}普通文件: ${archivePath}`, ); } if (fileStat.size > DIRECT_FILES_SINGLE_PUT_MAX_BYTES) { return uploadArchive({ archivePath, bucket, endpoint, objectKey, accessKeyId, accessKeySecret, fetchImpl, nowFn, sleepImpl, randomFn, maxAttempts, retryBaseDelayMs, retryMaxDelayMs, backupKind, archiveSha256, contentType, allowEmpty, bandwidthLimiter, }); } const verifiedArchiveSha256 = archiveSha256 || (await sha256FileHex(archivePath)); if (!/^[a-f0-9]{64}$/u.test(verifiedArchiveSha256)) { throw new Error(`归档 SHA-256 无效: ${verifiedArchiveSha256}`); } const requestOptions = { bucket, endpoint, objectKey, accessKeyId, accessKeySecret, fetchImpl, nowFn, sleepImpl, randomFn, maxAttempts, retryBaseDelayMs, retryMaxDelayMs, }; await signedOssRequest({ ...requestOptions, method: 'PUT', headers: { 'content-type': contentType, 'x-oss-meta-archive-sha256': verifiedArchiveSha256, 'x-oss-meta-file-sha256': verifiedArchiveSha256, 'x-oss-meta-file-size': String(fileStat.size), 'x-oss-meta-backup-kind': backupKind, }, contentLength: fileStat.size, bodyFactory: () => { if (fileStat.size === 0) { return Buffer.alloc(0); } const stream = createReadStream(archivePath); return bandwidthLimiter ? bandwidthLimiter.wrap(stream) : stream; }, operation: '上传逐文件对象', }); const verification = await verifyUploadedObject({ requestOptions, expectedContentLength: fileStat.size, expectedArchiveSha256: verifiedArchiveSha256, }); return { bucket, objectKey, contentLength: fileStat.size, archiveSha256: verifiedArchiveSha256, etag: '', uploadMode: 'single', partCount: 1, partSizeBytes: fileStat.size, verifiedAt: verification.verifiedAt, }; } export async function uploadManifestFile({ manifestPath, bucket, endpoint, objectKey, accessKeyId, accessKeySecret, fetchImpl = globalThis.fetch, nowFn = () => new Date(), sleepImpl = sleep, randomFn = Math.random, maxAttempts = DEFAULT_OSS_REQUEST_MAX_ATTEMPTS, bandwidthLimiter = null, }) { const body = readFileSync(manifestPath); if (body.length === 0) { throw new Error(`待上传 manifest 不能为空: ${manifestPath}`); } const archiveSha256 = sha256Hex(body); const requestOptions = { bucket, endpoint, objectKey, accessKeyId, accessKeySecret, fetchImpl, nowFn, sleepImpl, randomFn, maxAttempts, retryBaseDelayMs: DEFAULT_OSS_RETRY_BASE_DELAY_MS, retryMaxDelayMs: DEFAULT_OSS_RETRY_MAX_DELAY_MS, }; await signedOssRequest({ ...requestOptions, method: 'PUT', headers: { 'content-type': 'application/json', 'x-oss-meta-archive-sha256': archiveSha256, 'x-oss-meta-file-size': String(body.length), 'x-oss-meta-backup-kind': 'spacetimedb-backup-manifest', }, contentLength: body.length, bodyFactory: () => { const stream = createBufferReadStream(body); return bandwidthLimiter ? bandwidthLimiter.wrap(stream) : stream; }, operation: '上传 manifest', }); const verification = await verifyUploadedObject({ requestOptions, expectedContentLength: body.length, expectedArchiveSha256: archiveSha256, }); return { objectKey, contentLength: body.length, archiveSha256, verifiedAt: verification.verifiedAt, }; } function uploadedManifestPayload({ manifest, database, result }) { return { ...manifest, database, bucket: result.bucket, objectKey: result.objectKey, manifestObjectKey: `${result.objectKey}.manifest.json`, contentLength: result.contentLength, archiveSha256: result.archiveSha256, etag: result.etag, uploadMode: result.uploadMode, partCount: result.partCount, partSizeBytes: result.partSizeBytes, verifiedAt: result.verifiedAt, uploadedAt: new Date().toISOString(), uploadStatus: 'uploaded', }; } export async function uploadHistoryArchiveWithCleanup({ archivePath, manifestPath, manifest, statePath, uploadOptions, uploadFn = uploadArchive, manifestUploadFn = uploadManifestFile, verifyFn = verifyOssObject, }) { const result = await uploadFn({ archivePath, ...uploadOptions, backupKind: 'spacetimedb-history', }); const uploadedManifest = uploadedManifestPayload({ manifest, database: manifest.database, result, }); writeManifest({ manifestPath, payload: uploadedManifest }); const manifestUpload = await manifestUploadFn({ manifestPath, ...uploadOptions, objectKey: uploadedManifest.manifestObjectKey, }); uploadedManifest.manifestVerifiedAt = manifestUpload.verifiedAt; uploadedManifest.manifestContentLength = manifestUpload.contentLength; uploadedManifest.manifestArchiveSha256 = manifestUpload.archiveSha256; writeManifest({ manifestPath, payload: uploadedManifest }); let state = validateHistoryState(readManifest(statePath), { database: uploadedManifest.database, dataDir: uploadedManifest.dataDir, }); if (state.baseline.id !== uploadedManifest.baselineId) { throw new Error( `history manifest baselineId 与 state 不匹配: manifest=${uploadedManifest.baselineId}, state=${state.baseline.id}`, ); } await verifyFn({ ...uploadOptions, bucket: state.baseline.bucket, objectKey: state.baseline.objectKey, contentLength: state.baseline.contentLength, archiveSha256: state.baseline.archiveSha256, }); await verifyFn({ ...uploadOptions, bucket: state.baseline.bucket, objectKey: state.baseline.manifestObjectKey, contentLength: state.baseline.manifestContentLength, archiveSha256: state.baseline.manifestArchiveSha256, }); state = recordHistoryBatch({ statePath, state, manifest: uploadedManifest, uploadResult: result, manifestUpload, status: 'uploaded', }); const cleanup = cleanupHistoryCandidates({ dataDir: uploadedManifest.dataDir, candidates: uploadedManifest.candidates, }); state = recordHistoryBatch({ statePath, state, manifest: uploadedManifest, uploadResult: result, manifestUpload, status: 'cleaned', cleanedAt: new Date().toISOString(), }); return { result, uploadedManifest, cleanup, state }; } export function discoverDeferredArchiveUploads({ workDir, database, includeUploaded = false, }) { const resolvedWorkDir = resolvePath(workDir); if (!existsSync(resolvedWorkDir)) { return { archives: [], missingArchives: [] }; } const archives = []; const missingArchives = []; const manifestSuffix = '.tar.gz.manifest.json'; const expectedDatabase = String(database || '').trim(); const entries = readdirSync(resolvedWorkDir, { withFileTypes: true }) .filter( (candidate) => candidate.isFile() && candidate.name.endsWith(manifestSuffix), ) .sort((left, right) => left.name.localeCompare(right.name, 'en')); for (const entry of entries) { const manifestPath = join(resolvedWorkDir, entry.name); const manifest = readManifest(manifestPath); const uploadStatus = String(manifest.uploadStatus || '').trim(); if ( !['deferred', 'pending'].includes(uploadStatus) && !(includeUploaded && uploadStatus === 'uploaded') ) { continue; } if ( expectedDatabase && String(manifest.database || '').trim() !== expectedDatabase ) { continue; } if (!manifest.archivePath) { throw new Error(`deferred 备份清单缺少 archivePath: ${manifestPath}`); } const archivePath = resolvePath(manifest.archivePath); if ( dirname(archivePath) !== resolvedWorkDir || manifestPath !== `${archivePath}.manifest.json` ) { throw new Error(`deferred 备份路径与清单不匹配: ${manifestPath}`); } const candidate = { archivePath, manifestPath, manifest }; if (existsSync(archivePath)) { const archiveStat = lstatSync(archivePath); if (!archiveStat.isFile() || archiveStat.isSymbolicLink()) { throw new Error( `deferred 备份归档必须是非符号链接的普通文件: ${archivePath}`, ); } archives.push(candidate); } else { missingArchives.push(candidate); } } return { archives, missingArchives }; } async function uploadExistingArchive({ args, env, bucket, endpoint, accessKeyId, accessKeySecret, objectPrefix, bandwidthLimiter, }) { const archivePath = resolvePath(args.uploadArchive); if (!existsSync(archivePath)) { throw new Error(`待上传备份文件不存在: ${archivePath}`); } const manifestPath = resolvePath( args.manifestFile || `${archivePath}.manifest.json`, ); const manifest = existsSync(manifestPath) ? readManifest(manifestPath) : {}; const dataDir = firstNonEmpty( manifest.dataDir, env.GENARRATIVE_DATABASE_BACKUP_DATA_DIR, DEFAULT_PRODUCTION_DATA_DIR, ); const database = firstNonEmpty( args.database, manifest.database, env.GENARRATIVE_SPACETIME_DATABASE, basename(dataDir), ); const objectKey = firstNonEmpty( args.objectKey, manifest.objectKey, buildBackupNames({ database, dataDir, objectPrefix }).objectKey, ); if (manifest.backupKind !== 'spacetimedb-history') { manifest.backupKind = 'spacetimedb-data-dir'; manifest.baselineStatePath = firstNonEmpty( manifest.baselineStatePath, historyStatePath({ args, env, workDir: dirname(archivePath), database }), ); } console.log(`[database-backup] 上传已有备份: ${archivePath}`); console.log(`[database-backup] 目标对象: oss://${bucket}/${objectKey}`); if (args.dryRun) { console.log('[database-backup] dry-run,仅校验上传配置。'); return; } const statePath = resolvePath( firstNonEmpty( manifest.baselineStatePath, historyStatePath({ args, env, workDir: dirname(archivePath), database }), ), ); let result; let uploadedAt; if (manifest.backupKind === 'spacetimedb-history') { const historyResult = await uploadHistoryArchiveWithCleanup({ archivePath, manifestPath, manifest, statePath, uploadOptions: { bucket, endpoint, objectKey, accessKeyId, accessKeySecret, bandwidthLimiter, }, }); result = historyResult.result; uploadedAt = historyResult.uploadedManifest.uploadedAt; console.log( `[database-backup] history 上传并清理完成: ${JSON.stringify(historyResult.cleanup)}`, ); } else { result = await uploadArchive({ archivePath, bucket, endpoint, objectKey, accessKeyId, accessKeySecret, bandwidthLimiter, }); const uploadedManifest = uploadedManifestPayload({ manifest, database, result, }); uploadedAt = uploadedManifest.uploadedAt; writeManifest({ manifestPath, payload: uploadedManifest }); const manifestUpload = await uploadManifestFile({ manifestPath, bucket, endpoint, objectKey: uploadedManifest.manifestObjectKey, accessKeyId, accessKeySecret, bandwidthLimiter, }); uploadedManifest.manifestVerifiedAt = manifestUpload.verifiedAt; uploadedManifest.manifestContentLength = manifestUpload.contentLength; uploadedManifest.manifestArchiveSha256 = manifestUpload.archiveSha256; writeManifest({ manifestPath, payload: uploadedManifest }); const previousState = existsSync(statePath) ? validateHistoryState(readManifest(statePath), { database, dataDir }) : null; const baseline = normalizeUploadedBaselineManifest(uploadedManifest, { database, dataDir, }); writeBaselineState({ statePath, baseline, previousState }); console.log(`[database-backup] 已写入 baseline state: ${statePath}`); } console.log(`[database-backup] 上传完成: ${JSON.stringify(result)}`); if (args.resultFile) { writeFileSync( resolvePath(args.resultFile), `${JSON.stringify({ archivePath, manifestPath, statePath, ...result, uploadedAt }, null, 2)}\n`, 'utf8', ); } const keepLocal = args.keepLocal || String(env.GENARRATIVE_DATABASE_BACKUP_KEEP_LOCAL ?? '') .trim() .toLowerCase() === 'true'; if (!keepLocal) { rmSync(archivePath, { force: true }); rmSync(manifestPath, { force: true }); console.log( '[database-backup] 已删除本地临时备份文件;如需保留请设置 --keep-local。', ); } else { console.log(`[database-backup] 已保留本地备份: ${archivePath}`); console.log(`[database-backup] 已保留备份清单: ${manifestPath}`); } } async function uploadDeferredArchives({ args, env, bucket, endpoint, accessKeyId, accessKeySecret, objectPrefix, database, bandwidthLimiter, }) { const workDir = resolvePath(args.uploadDeferredDir); const keepLocal = args.keepLocal || String(env.GENARRATIVE_DATABASE_BACKUP_KEEP_LOCAL ?? '') .trim() .toLowerCase() === 'true'; const { archives, missingArchives } = discoverDeferredArchiveUploads({ workDir, database, includeUploaded: !keepLocal, }); for (const { manifestPath } of missingArchives) { console.warn( `[database-backup] deferred 清单对应的本地归档不存在,跳过: ${manifestPath}`, ); } if (archives.length === 0) { console.log(`[database-backup] 没有可补偿的本地归档: ${workDir}`); return; } console.log( `[database-backup] 开始串行上传待补偿本地归档: count=${archives.length}`, ); for (const { archivePath, manifestPath } of archives) { await uploadExistingArchive({ args: { ...args, uploadArchive: archivePath, manifestFile: manifestPath }, env, bucket, endpoint, accessKeyId, accessKeySecret, objectPrefix, bandwidthLimiter, }); } console.log( `[database-backup] 待补偿本地归档上传完成: count=${archives.length}`, ); } async function publishExistingManifest({ args, bucket, endpoint, accessKeyId, accessKeySecret, bandwidthLimiter, }) { const manifestPath = resolvePath(args.publishManifest); const manifest = readManifest(manifestPath); if (manifest.uploadStatus !== 'uploaded' || !manifest.objectKey) { throw new Error( '只允许发布 uploadStatus=uploaded 且包含 objectKey 的备份 manifest。', ); } manifest.manifestObjectKey = manifest.manifestObjectKey || `${manifest.objectKey}.manifest.json`; writeManifest({ manifestPath, payload: manifest }); const result = await uploadManifestFile({ manifestPath, bucket, endpoint, objectKey: manifest.manifestObjectKey, accessKeyId, accessKeySecret, bandwidthLimiter, }); manifest.manifestVerifiedAt = result.verifiedAt; manifest.manifestContentLength = result.contentLength; manifest.manifestArchiveSha256 = result.archiveSha256; writeManifest({ manifestPath, payload: manifest }); console.log( `[database-backup] manifest 上传并验真完成: ${JSON.stringify(result)}`, ); } export async function resumeUploadedHistoryBatch({ statePath, state, dataDir, verificationOptions, verifyFn = verifyOssObject, }) { const pendingBatch = state.batches.find( (batch) => batch.status === 'uploaded', ); if (!pendingBatch) { return state; } console.log( `[database-backup] 重试已上传 history 批次的本地清理: ${pendingBatch.batchId}`, ); await verifyFn({ ...verificationOptions, objectKey: pendingBatch.objectKey, contentLength: pendingBatch.contentLength, archiveSha256: pendingBatch.archiveSha256, }); await verifyFn({ ...verificationOptions, objectKey: pendingBatch.manifestObjectKey, contentLength: pendingBatch.manifestContentLength, archiveSha256: pendingBatch.manifestArchiveSha256, }); const cleanup = cleanupHistoryCandidates({ dataDir, candidates: pendingBatch.candidates, }); const manifest = { batchId: pendingBatch.batchId, uploadedAt: pendingBatch.uploadedAt, candidates: pendingBatch.candidates, }; const uploadResult = { objectKey: pendingBatch.objectKey, contentLength: pendingBatch.contentLength, archiveSha256: pendingBatch.archiveSha256, verifiedAt: pendingBatch.verifiedAt, }; const manifestUpload = { objectKey: pendingBatch.manifestObjectKey, contentLength: pendingBatch.manifestContentLength, archiveSha256: pendingBatch.manifestArchiveSha256, verifiedAt: pendingBatch.manifestVerifiedAt, }; const nextState = recordHistoryBatch({ statePath, state, manifest, uploadResult, manifestUpload, status: 'cleaned', cleanedAt: new Date().toISOString(), }); console.log( `[database-backup] 已完成 history 清理重试: ${JSON.stringify(cleanup)}`, ); return nextState; } async function runHistoryBackup({ args, env, dataDir, workDir, database, bucket, endpoint, accessKeyId, accessKeySecret, objectPrefix, keepLocal, bandwidthLimiter, }) { const statePath = historyStatePath({ args, env, workDir, database }); let state = loadOrImportHistoryState({ args, env, statePath, database, dataDir, }); if (!args.dryRun && !args.deferUpload) { console.log( `[database-backup] 重新验真 full baseline: oss://${state.baseline.bucket}/${state.baseline.objectKey}`, ); await verifyOssObject({ bucket: state.baseline.bucket, endpoint, objectKey: state.baseline.objectKey, accessKeyId, accessKeySecret, contentLength: state.baseline.contentLength, archiveSha256: state.baseline.archiveSha256, }); await verifyOssObject({ bucket: state.baseline.bucket, endpoint, objectKey: state.baseline.manifestObjectKey, accessKeyId, accessKeySecret, contentLength: state.baseline.manifestContentLength, archiveSha256: state.baseline.manifestArchiveSha256, }); state = await resumeUploadedHistoryBatch({ statePath, state, dataDir, verificationOptions: { bucket, endpoint, accessKeyId, accessKeySecret }, }); } const plan = discoverHistoryPlan({ dataDir }); console.log( `[database-backup] history replicas: ${JSON.stringify(plan.replicas)}`, ); console.log( `[database-backup] history 候选: count=${plan.candidates.length}, size=${formatBytes(plan.totalSizeBytes)}`, ); if (args.resultFile) { writeFileSync( resolvePath(args.resultFile), `${JSON.stringify({ statePath, baseline: state.baseline, ...plan }, null, 2)}\n`, 'utf8', ); } if (args.dryRun) { console.log( '[database-backup] history dry-run,仅输出安全候选,不打包、上传或删除。', ); return; } if (plan.candidates.length === 0) { console.log('[database-backup] 没有可归档的 history 候选。'); return; } assertSufficientHistoryWorkDirSpace({ historySizeBytes: plan.totalSizeBytes, workDir, args, env, }); const batchId = historyBatchId({ baselineId: state.baseline.id, plan }); const { fileName, objectKey } = buildHistoryNames({ database, objectPrefix, baselineId: state.baseline.id, batchId, }); const archivePath = resolve(workDir, fileName); const manifestPath = `${archivePath}.manifest.json`; const manifest = { schemaVersion: HISTORY_MANIFEST_SCHEMA_VERSION, backupKind: 'spacetimedb-history', createdAt: new Date().toISOString(), database, dataDir, bucket, objectKey, archivePath, baselineId: state.baseline.id, baselineStatePath: statePath, batchId, replicas: plan.replicas, candidates: plan.candidates, totalSizeBytes: plan.totalSizeBytes, uploadStatus: args.deferUpload ? 'deferred' : 'pending', }; writeManifest({ manifestPath, payload: manifest }); createHistoryArchive({ dataDir, workDir, fileName, manifestPath, candidates: plan.candidates, }); if (args.deferUpload) { console.log( `[database-backup] 已生成 history 归档,延后上传且未清理源文件: ${archivePath}`, ); if (args.resultFile) { writeFileSync( resolvePath(args.resultFile), `${JSON.stringify({ archivePath, manifestPath, statePath, bucket, objectKey, batchId }, null, 2)}\n`, 'utf8', ); } return; } const historyResult = await uploadHistoryArchiveWithCleanup({ archivePath, manifestPath, manifest, statePath, uploadOptions: { bucket, endpoint, objectKey, accessKeyId, accessKeySecret, bandwidthLimiter, }, }); console.log( `[database-backup] history 上传并清理完成: ${JSON.stringify(historyResult.cleanup)}`, ); if (args.resultFile) { writeFileSync( resolvePath(args.resultFile), `${JSON.stringify( { archivePath, manifestPath, statePath, batchId, ...historyResult.result, uploadedAt: historyResult.uploadedManifest.uploadedAt, }, null, 2, )}\n`, 'utf8', ); } if (!keepLocal) { rmSync(archivePath, { force: true }); rmSync(manifestPath, { force: true }); console.log('[database-backup] 已删除本地 history 临时归档和清单。'); } } async function main() { const args = parseArgs(process.argv.slice(2)); const env = loadEffectiveEnv(args.envFiles); const isProductionLike = existsSync(DEFAULT_PRODUCTION_DATA_DIR) && process.platform !== 'win32'; const dataDir = resolvePath( firstNonEmpty( args.dataDir, env.GENARRATIVE_DATABASE_BACKUP_DATA_DIR, isProductionLike ? DEFAULT_PRODUCTION_DATA_DIR : DEFAULT_LOCAL_DATA_DIR, ), ); const workDir = resolvePath( firstNonEmpty( args.workDir, args.uploadDeferredDir, env.GENARRATIVE_DATABASE_BACKUP_WORK_DIR, isProductionLike ? DEFAULT_PRODUCTION_WORK_DIR : DEFAULT_LOCAL_WORK_DIR, ), ); const bucket = firstNonEmpty( args.bucket, env.GENARRATIVE_DATABASE_BACKUP_OSS_BUCKET, env.ALIYUN_OSS_BUCKET, ); const endpoint = normalizeEndpoint( firstNonEmpty( args.endpoint, env.GENARRATIVE_DATABASE_BACKUP_OSS_ENDPOINT, env.ALIYUN_OSS_ENDPOINT, ), ); const accessKeyId = firstNonEmpty( args.accessKeyId, env.GENARRATIVE_DATABASE_BACKUP_OSS_ACCESS_KEY_ID, env.ALIYUN_OSS_ACCESS_KEY_ID, ); const accessKeySecret = firstNonEmpty( args.accessKeySecret, env.GENARRATIVE_DATABASE_BACKUP_OSS_ACCESS_KEY_SECRET, env.ALIYUN_OSS_ACCESS_KEY_SECRET, ); const objectPrefix = firstNonEmpty( args.objectPrefix, env.GENARRATIVE_DATABASE_BACKUP_OSS_PREFIX, 'database-backups', ); const database = firstNonEmpty( args.database, env.GENARRATIVE_SPACETIME_DATABASE, basename(dataDir), ); const keepLocal = args.keepLocal || String(env.GENARRATIVE_DATABASE_BACKUP_KEEP_LOCAL ?? '') .trim() .toLowerCase() === 'true'; const storageFormat = firstNonEmpty( args.storageFormat, env.GENARRATIVE_DATABASE_BACKUP_STORAGE_FORMAT, 'archive', ); const directFilesConcurrency = parseDirectFilesConcurrency( env.GENARRATIVE_DATABASE_BACKUP_FILES_CONCURRENCY, ); const uploadBandwidthLimiter = createUploadBandwidthLimiter( env.GENARRATIVE_DATABASE_BACKUP_UPLOAD_MAX_BYTES_PER_SECOND, ); if (!['full', 'history'].includes(args.mode)) { throw new Error(`--mode 只能是 full 或 history,实际: ${args.mode}`); } if (!['archive', 'files'].includes(storageFormat)) { throw new Error( `--storage-format 只能是 archive 或 files,实际: ${storageFormat}`, ); } for (const [label, value] of Object.entries({ bucket, endpoint, accessKeyId, accessKeySecret, })) { if (!value) { throw new Error(`缺少 ${label} 配置`); } } if (args.restoreFilesState && args.restoreFilesLatest) { throw new Error( '--restore-files-state 与 --restore-files-latest 不能同时使用。', ); } if (args.restoreFilesState) { if (!args.restoreDir) { throw new Error('--restore-files-state 必须同时传 --restore-dir。'); } await restoreDirectFilesBackup({ statePath: args.restoreFilesState, restoreDir: args.restoreDir, database, bucket, uploadOptions: { bucket, endpoint, accessKeyId, accessKeySecret }, resultFile: args.resultFile, dryRun: args.dryRun, }); return; } if (args.restoreFilesLatest) { if (!args.restoreDir) { throw new Error('--restore-files-latest 必须同时传 --restore-dir。'); } await restoreDirectFilesLatest({ restoreDir: args.restoreDir, database, bucket, objectPrefix, uploadOptions: { bucket, endpoint, accessKeyId, accessKeySecret }, resultFile: args.resultFile, dryRun: args.dryRun, }); return; } if (args.restoreDir) { throw new Error( '--restore-dir 只能与 --restore-files-state 或 --restore-files-latest 一起使用。', ); } if (args.uploadArchive && args.uploadDeferredDir) { throw new Error('--upload-archive 与 --upload-deferred-dir 不能同时使用。'); } if (!args.dryRun) { const lockPath = acquireBackupLock({ workDir, database }); console.log(`[database-backup] 已获取进程锁: ${lockPath}`); } if (args.publishManifest) { await publishExistingManifest({ args, bucket, endpoint, accessKeyId, accessKeySecret, bandwidthLimiter: uploadBandwidthLimiter, }); return; } if (args.uploadArchive) { await uploadExistingArchive({ args, env, bucket, endpoint, accessKeyId, accessKeySecret, objectPrefix, bandwidthLimiter: uploadBandwidthLimiter, }); return; } if (args.uploadDeferredDir) { await uploadDeferredArchives({ args, env, bucket, endpoint, accessKeyId, accessKeySecret, objectPrefix, database, bandwidthLimiter: uploadBandwidthLimiter, }); return; } if (storageFormat === 'files') { if (args.deferUpload) { throw new Error( 'files 模式无需本地归档且不支持 --defer-upload;失败后使用同一 work-dir 重跑即可续传。', ); } const stopService = args.stopService || firstNonEmpty(env.GENARRATIVE_DATABASE_BACKUP_STOP_SERVICE); const restartServicesAfter = collectRestartServicesAfterBackup({ args, env, }); const stopMarkerPath = databaseBackupStopMarkerPath(workDir); let serviceStopped = false; let backupError = null; let restoreError = null; try { if (args.mode === 'full' && !args.dryRun) { serviceStopped = stopServiceIfNeeded(stopService, stopMarkerPath); } await runDirectFilesBackup({ mode: args.mode, dataDir, workDir, database, bucket, objectPrefix, dryRun: args.dryRun, resultFile: args.resultFile, uploadOptions: { bucket, endpoint, accessKeyId, accessKeySecret, bandwidthLimiter: uploadBandwidthLimiter, }, concurrency: directFilesConcurrency, }); } catch (error) { backupError = error; } finally { try { if (serviceStopped) { restoreServicesAfterBackup({ stopService, serviceStopped, restartServicesAfter, stopMarkerPath, }); } else if (!backupError && args.mode === 'full' && !args.dryRun) { restartServicesAfterBackup(restartServicesAfter); } } catch (error) { restoreError = error; } } if (backupError && restoreError) { throw new AggregateError( [backupError, restoreError], `files 备份失败,且恢复依赖服务时也失败: ${backupError.message}; ${restoreError.message}`, ); } if (backupError) { throw backupError; } if (restoreError) { throw restoreError; } return; } if (args.mode === 'history') { await runHistoryBackup({ args, env, dataDir, workDir, database, bucket, endpoint, accessKeyId, accessKeySecret, objectPrefix, keepLocal, bandwidthLimiter: uploadBandwidthLimiter, }); return; } const { fileName, objectKey } = buildBackupNames({ database, dataDir, objectPrefix, }); console.log(`[database-backup] 数据目录: ${dataDir}`); console.log(`[database-backup] 本地临时目录: ${workDir}`); console.log(`[database-backup] 目标对象: oss://${bucket}/${objectKey}`); if (args.dryRun) { console.log('[database-backup] dry-run,仅校验配置,不打包上传。'); return; } let archivePath = ''; let serviceStopped = false; let backupError = null; let restoreError = null; const stopService = args.stopService || firstNonEmpty(env.GENARRATIVE_DATABASE_BACKUP_STOP_SERVICE); const restartServicesAfter = collectRestartServicesAfterBackup({ args, env }); const stopMarkerPath = databaseBackupStopMarkerPath(workDir); try { assertSufficientWorkDirSpace({ dataDir, workDir, args, env }); serviceStopped = stopServiceIfNeeded(stopService, stopMarkerPath); archivePath = createArchive({ dataDir, workDir, fileName }); } catch (error) { backupError = error; } finally { try { if (serviceStopped) { restoreServicesAfterBackup({ stopService, serviceStopped, restartServicesAfter, stopMarkerPath, }); } else if (!backupError) { restartServicesAfterBackup(restartServicesAfter); } } catch (error) { restoreError = error; } } if (backupError) { if (restoreError) { throw new AggregateError( [backupError, restoreError], `数据库备份失败,且恢复依赖服务时也失败: ${backupError.message}; ${restoreError.message}`, ); } throw backupError; } if (restoreError) { throw restoreError; } const manifestPath = `${archivePath}.manifest.json`; const baselineStatePath = historyStatePath({ args, env, workDir, database }); const fullManifest = { backupKind: 'spacetimedb-data-dir', createdAt: new Date().toISOString(), database, dataDir, bucket, objectKey, archivePath, baselineStatePath, uploadStatus: args.deferUpload ? 'deferred' : 'pending', }; writeManifest({ manifestPath, payload: fullManifest, }); if (args.deferUpload) { console.log(`[database-backup] 已生成本地冷备份,延后上传: ${archivePath}`); console.log(`[database-backup] 已写入备份清单: ${manifestPath}`); if (args.resultFile) { writeFileSync( resolvePath(args.resultFile), `${JSON.stringify({ archivePath, manifestPath, baselineStatePath, bucket, objectKey }, null, 2)}\n`, 'utf8', ); } return; } const result = await uploadArchive({ archivePath, bucket, endpoint, objectKey, accessKeyId, accessKeySecret, bandwidthLimiter: uploadBandwidthLimiter, }); console.log(`[database-backup] 上传完成: ${JSON.stringify(result)}`); const uploadedManifest = uploadedManifestPayload({ manifest: fullManifest, database, result, }); writeManifest({ manifestPath, payload: uploadedManifest }); const manifestUpload = await uploadManifestFile({ manifestPath, bucket, endpoint, objectKey: uploadedManifest.manifestObjectKey, accessKeyId, accessKeySecret, bandwidthLimiter: uploadBandwidthLimiter, }); uploadedManifest.manifestVerifiedAt = manifestUpload.verifiedAt; uploadedManifest.manifestContentLength = manifestUpload.contentLength; uploadedManifest.manifestArchiveSha256 = manifestUpload.archiveSha256; writeManifest({ manifestPath, payload: uploadedManifest }); const previousState = existsSync(baselineStatePath) ? validateHistoryState(readManifest(baselineStatePath), { database, dataDir, }) : null; const baseline = normalizeUploadedBaselineManifest(uploadedManifest, { database, dataDir, }); writeBaselineState({ statePath: baselineStatePath, baseline, previousState }); console.log(`[database-backup] 已写入 baseline state: ${baselineStatePath}`); if (!keepLocal) { rmSync(archivePath, { force: true }); rmSync(manifestPath, { force: true }); console.log( '[database-backup] 已删除本地临时备份文件;如需保留请设置 --keep-local。', ); } else { console.log(`[database-backup] 已保留本地备份: ${archivePath}`); console.log(`[database-backup] 已保留备份清单: ${manifestPath}`); } } function formatErrorDetails(error) { if (!error || typeof error !== 'object') { return ''; } return ['code', 'errno', 'syscall', 'hostname', 'host', 'port', 'address'] .map((field) => { const value = error[field]; return value === undefined || value === null || value === '' ? '' : `${field}=${String(value)}`; }) .filter(Boolean) .join(' '); } function describeError(error) { const lines = []; let current = error; for (let depth = 0; current && depth < 5; depth += 1) { const label = depth === 0 ? 'error' : `cause[${depth}]`; if (!(current instanceof Error)) { lines.push(`${label}: ${String(current)}`); break; } lines.push(`${label}: ${current.name}: ${current.message}`); const details = formatErrorDetails(current); if (details) { lines.push(`${label} details: ${details}`); } if (current instanceof AggregateError) { current.errors.slice(0, 3).forEach((item, index) => { const itemText = item instanceof Error ? `${item.name}: ${item.message}` : String(item); const itemDetails = formatErrorDetails(item); lines.push( `${label}.errors[${index}]: ${itemText}${itemDetails ? ` (${itemDetails})` : ''}`, ); }); } current = current.cause; } return lines; } if ( process.argv[1] && realpathSync(resolve(process.argv[1])) === realpathSync(__filename) ) { main().catch((error) => { for (const line of describeError(error)) { console.error(`[database-backup] ${line}`); } process.exit(1); }); }