#!/usr/bin/env node import {spawnSync} from 'node:child_process'; import {createHash} from 'node:crypto'; import {chmodSync, existsSync, lstatSync, mkdirSync, mkdtempSync, readFileSync, readlinkSync, rmSync, statSync, symlinkSync, writeFileSync} from 'node:fs'; import {tmpdir} from 'node:os'; import path from 'node:path'; import {Readable} from 'node:stream'; import { buildAuthorization, buildCanonicalQuery, cleanupHistoryCandidates, collectDirectFileEntries, createUploadBandwidthLimiter, discoverHistoryPlan, restoreDirectFilesBackup, restoreDirectFilesLatest, resumeUploadedHistoryBatch, runDirectFilesBackup, uploadArchive, uploadDirectFile, uploadHistoryArchiveWithCleanup, uploadManifestFile, } from './database-backup-to-oss.mjs'; const BACKUP_SCRIPT = path.resolve('scripts/database-backup-to-oss.mjs'); const tmpRoot = mkdtempSync(path.join(tmpdir(), 'genarrative-database-backup-check-')); const failures = []; try { await main(); } finally { rmSync(tmpRoot, {recursive: true, force: true}); } if (failures.length > 0) { console.error('[check:database-backup] FAILED'); for (const failure of failures) { console.error(`- ${failure}`); } process.exit(1); } console.log('[check:database-backup] OK'); async function main() { assertCanonicalQueryAndAuthorizationIncludeMultipartParameters(); assertInsufficientSpaceStopsBeforeServiceChanges(); assertArchiveFailureStillRestoresDependentServices(); await assertMultipartUploadRetriesAndVerifiesRemoteLength(); await assertUploadBandwidthLimiterSharesBudgetAndPropagatesErrors(); await assertDirectSmallFileUsesSinglePut(); await assertMissingPartEtagAbortsMultipartUpload(); await assertCompleteResponseAmbiguityUsesHeadVerification(); await assertHeadLengthMismatchAbortsMultipartUpload(); await assertHeadShaMismatchAbortsMultipartUpload(); await assertManifestUploadUsesShaAndHeadVerification(); assertHistoryDiscoversDevAndProductionLayoutsWithMultipleReplicas(); assertHistoryRequiresBaselineAndProducesDeterministicDeferredBatch(); assertHistoryBackupLockRejectsLiveAndStaleOwners(); assertHistorySkipsReplicaWithoutSnapshotAndRejectsMalformedNames(); assertHistoryStatDriftPreventsAnyCleanup(); await assertHistoryUploadFailureDoesNotDeleteSources(); await assertHistorySuccessfulUploadCleansAndIsIdempotent(); await assertHistoryResumeReverifiesArchiveAndManifest(); await assertDirectFilesPreservePathsAndIncrementWithoutDuplicateUpload(); await assertDirectFilesConcurrencyIsBounded(); await assertDirectHistoryPublishesCatalogBeforeCleanup(); await assertDirectHistoryWithoutCandidatesPublishesLatest(); await assertDirectFilesRestoreDownloadsCatalogAndObjects(); } function createDirectOssHarness() { const objects = new Map(); const uploadedKeys = []; const verifiedKeys = []; const uploadFn = async ({archivePath, objectKey, archiveSha256}) => { const body = readFileSync(archivePath); const sha256 = createHash('sha256').update(body).digest('hex'); assertEqual(sha256, archiveSha256, `direct file ${objectKey} 的上传 SHA 必须来自实际内容。`); objects.set(objectKey, {body, contentLength: body.length, sha256}); uploadedKeys.push(objectKey); return {objectKey, contentLength: body.length, archiveSha256: sha256, verifiedAt: '2026-07-16T01:00:00.000Z'}; }; const uploadManifestFn = async ({manifestPath, objectKey}) => { const body = readFileSync(manifestPath); const sha256 = createHash('sha256').update(body).digest('hex'); objects.set(objectKey, {body, contentLength: body.length, sha256}); uploadedKeys.push(objectKey); return {objectKey, contentLength: body.length, archiveSha256: sha256, verifiedAt: '2026-07-16T01:00:01.000Z'}; }; const verifyFn = async ({objectKey, contentLength, archiveSha256}) => { verifiedKeys.push(objectKey); const object = objects.get(objectKey); if (!object) { const error = new Error(`missing ${objectKey}`); error.status = 404; throw error; } if (object.contentLength !== contentLength || object.sha256 !== archiveSha256) { throw new Error(`mismatch ${objectKey}`); } return {verifiedAt: '2026-07-16T01:00:02.000Z'}; }; return {objects, uploadedKeys, verifiedKeys, uploadFn, uploadManifestFn, verifyFn}; } async function assertDirectSmallFileUsesSinglePut() { const filePath = path.join(tmpRoot, 'direct-small-file.bin'); const body = Buffer.from('small direct object'); const sha256 = createHash('sha256').update(body).digest('hex'); writeFileSync(filePath, body); const methods = []; let uploadedBytes = 0; const fetchImpl = async (_url, options) => { methods.push(options.method); if (options.method === 'PUT') { for await (const chunk of options.body) { uploadedBytes += chunk.length; } return new Response('', {status: 200, headers: {etag: '"single-etag"'}}); } if (options.method === 'HEAD') { return new Response(null, {status: 200, headers: { 'content-length': String(body.length), 'x-oss-meta-file-sha256': sha256, }}); } throw new Error(`unexpected method ${options.method}`); }; const result = await uploadDirectFile({ archivePath: filePath, bucket: 'backup-bucket', endpoint: 'oss-cn-shanghai.aliyuncs.com', objectKey: 'database-backups/test-db/files/small', accessKeyId: 'test-id', accessKeySecret: 'test-secret', archiveSha256: sha256, fetchImpl, bandwidthLimiter: createUploadBandwidthLimiter(64 * 1024), }); assertEqual(result.uploadMode, 'single', '小型逐文件对象必须使用单次 PUT。'); assertEqual(methods.join(','), 'PUT,HEAD', '小型逐文件对象只能执行 PUT 后 HEAD 验真,不得进入 multipart。'); assertEqual(uploadedBytes, body.length, '逐文件带宽限制流不得丢失上传内容。'); } async function assertUploadBandwidthLimiterSharesBudgetAndPropagatesErrors() { assertEqual(createUploadBandwidthLimiter('0'), null, '上传带宽限制为 0 时必须关闭。'); assertThrows( () => createUploadBandwidthLimiter('1023'), '必须为空、0 或 >= 1024 的整数', '上传带宽限制必须拒绝过小值。', ); const delays = []; const limiter = createUploadBandwidthLimiter(1024, { nowFn: () => 0, sleepImpl: async (delayMs) => { delays.push(delayMs); }, }); const consume = async (stream) => { let totalBytes = 0; for await (const chunk of stream) { totalBytes += chunk.length; } return totalBytes; }; const [firstBytes, secondBytes] = await Promise.all([ consume(limiter.wrap(Readable.from([Buffer.alloc(1024), Buffer.alloc(1024)], {objectMode: false}))), consume(limiter.wrap(Readable.from([Buffer.alloc(1024), Buffer.alloc(1024)], {objectMode: false}))), ]); assertEqual(firstBytes + secondBytes, 4096, '共享上传限速器不得丢失并发流内容。'); assertEqual(delays.join(','), '1000,2000,3000,4000', '两个并发上传流必须共享同一个累计带宽预算。'); let sourceError = null; try { await consume(limiter.wrap(Readable.from((async function* failingSource() { yield Buffer.alloc(1); throw new Error('source-read-failed'); })(), {objectMode: false}))); } catch (error) { sourceError = error; } assertIncludes(sourceError?.message, 'source-read-failed', '限速流必须向上传请求透传源读取错误。'); } async function assertDirectFilesPreservePathsAndIncrementWithoutDuplicateUpload() { const root = path.join(tmpRoot, 'direct-files-incremental'); const dataDir = path.join(root, 'stdb'); const workDir = path.join(root, 'work'); mkdirSync(path.join(dataDir, 'replicas', '1', 'snapshots', '00000000000000000010.snapshot_dir', 'objects'), {recursive: true}); mkdirSync(path.join(dataDir, 'empty-directory'), {recursive: true}); mkdirSync(path.join(dataDir, 'bin', '2.6.0'), {recursive: true}); symlinkSync('2.6.0', path.join(dataDir, 'bin', 'current')); writeFileSync(path.join(dataDir, 'control-db'), 'control'); writeFileSync( path.join(dataDir, 'replicas', '1', 'snapshots', '00000000000000000010.snapshot_dir', 'objects', 'object.bin'), 'snapshot object', ); const harness = createDirectOssHarness(); const options = { mode: 'full', dataDir, workDir, database: 'test-db', bucket: 'backup-bucket', objectPrefix: 'database-backups', uploadOptions: {}, uploadFn: harness.uploadFn, uploadManifestFn: harness.uploadManifestFn, verifyFn: harness.verifyFn, }; const collected = await collectDirectFileEntries({dataDir, database: 'test-db', objectPrefix: 'database-backups'}); assertTrue( collected.files.some(({path: filePath}) => filePath === 'replicas/1/snapshots/00000000000000000010.snapshot_dir/objects/object.bin'), 'files catalog 必须原样保留 snapshot 内文件的相对路径。', ); assertTrue(collected.directories.includes('empty-directory'), 'files catalog 必须保留空目录。'); assertTrue( collected.symlinks.some(({path: symlinkPath, target}) => symlinkPath === 'bin/current' && target === '2.6.0'), 'files catalog 必须保留指向 data-dir 内部的相对符号链接。', ); const first = await runDirectFilesBackup(options); assertEqual(first.uploadedCount, 2, '首次 files full 应上传全部普通文件。'); assertTrue(!Object.hasOwn(first.catalog, 'dataDir'), '远端 files catalog 不得绑定 staging 主机的绝对 data-dir。'); const latestObjectKey = 'database-backups/test-db/latest.json'; const latest = JSON.parse(harness.objects.get(latestObjectKey).body.toString('utf8')); assertEqual(latest.latestFullCatalog.catalogId, first.catalogId, 'latest pointer 必须指向已验真的最新 full catalog。'); assertTrue(!Object.hasOwn(latest.latestFullCatalog, 'files'), 'latest full ref 不得嵌入 files 数组。'); assertTrue(latest.historyCatalogs.every((catalog) => !Object.hasOwn(catalog, 'files')), 'latest history ref 不得嵌入 files 数组。'); const immutableUploadsAfterFirst = harness.uploadedKeys.filter((objectKey) => objectKey !== latestObjectKey).length; const repeated = await runDirectFilesBackup(options); assertEqual(repeated.uploadedCount, 0, '相同目录重复运行不得重复上传文件。'); assertEqual( harness.uploadedKeys.filter((objectKey) => objectKey !== latestObjectKey).length, immutableUploadsAfterFirst, '相同 catalog 重跑不得重复 PUT 文件或 catalog,但应覆盖验真 latest pointer。', ); writeFileSync(path.join(dataDir, 'control-db'), 'control changed'); writeFileSync(path.join(dataDir, 'new-program.bin'), 'new program'); const incremental = await runDirectFilesBackup(options); assertEqual(incremental.uploadedCount, 2, '增量 files full 只应上传新增和变化文件。'); assertEqual(incremental.reusedCount, 1, '增量 files full 应复用未变化 snapshot 文件。'); } async function assertDirectFilesConcurrencyIsBounded() { const root = path.join(tmpRoot, 'direct-files-concurrency'); const dataDir = path.join(root, 'stdb'); const workDir = path.join(root, 'work'); mkdirSync(dataDir, {recursive: true}); for (let index = 0; index < 12; index += 1) { writeFileSync(path.join(dataDir, `file-${index}.bin`), `content-${index}`); } const harness = createDirectOssHarness(); let activeUploads = 0; let maxActiveUploads = 0; const uploadFn = async (options) => { activeUploads += 1; maxActiveUploads = Math.max(maxActiveUploads, activeUploads); await new Promise((resolve) => setTimeout(resolve, 5)); try { return await harness.uploadFn(options); } finally { activeUploads -= 1; } }; await runDirectFilesBackup({ mode: 'full', dataDir, workDir, database: 'test-db', bucket: 'backup-bucket', objectPrefix: 'database-backups', uploadOptions: {}, uploadFn, uploadManifestFn: harness.uploadManifestFn, verifyFn: harness.verifyFn, concurrency: 3, }); assertTrue(maxActiveUploads > 1, 'files 备份应按配置并发处理多个对象。'); assertTrue(maxActiveUploads <= 3, 'files 备份对象并发数不得超过配置上限。'); } async function assertDirectHistoryPublishesCatalogBeforeCleanup() { const fixture = createHistoryFixture('direct-files-history-cleanup', {nestedData: false}); createReplicaHistory(fixture.replicasDir, '1', {snapshots: [0, 10], segments: [0, 1, 11]}); const harness = createDirectOssHarness(); const common = { dataDir: fixture.dataDir, workDir: fixture.workDir, database: 'test-db', bucket: 'backup-bucket', objectPrefix: 'database-backups', uploadOptions: {}, uploadFn: harness.uploadFn, verifyFn: harness.verifyFn, }; await runDirectFilesBackup({...common, mode: 'full', uploadManifestFn: harness.uploadManifestFn}); const plan = discoverHistoryPlan({dataDir: fixture.dataDir}); let failure = null; try { await runDirectFilesBackup({ ...common, mode: 'history', uploadManifestFn: async () => { throw new Error('synthetic direct catalog failure'); }, }); } catch (error) { failure = error; } assertIncludes(failure?.message ?? '', 'synthetic direct catalog failure', 'direct history catalog 发布失败必须向上返回。'); for (const candidate of plan.candidates) { assertTrue(existsSync(path.join(fixture.dataDir, candidate.path)), `direct history catalog 发布失败不得删除: ${candidate.path}`); } let pointerFailure = null; try { await runDirectFilesBackup({ ...common, mode: 'history', uploadManifestFn: harness.uploadManifestFn, verifyFn: async (options) => { if (options.objectKey.endsWith('/latest.json')) { throw new Error('synthetic latest pointer HEAD failure'); } return harness.verifyFn(options); }, }); } catch (error) { pointerFailure = error; } assertIncludes(pointerFailure?.message ?? '', 'synthetic latest pointer HEAD failure', 'latest pointer HEAD 验真失败必须向上返回。'); for (const candidate of plan.candidates) { assertTrue(existsSync(path.join(fixture.dataDir, candidate.path)), `latest pointer 发布失败不得删除: ${candidate.path}`); } const success = await runDirectFilesBackup({...common, mode: 'history', uploadManifestFn: harness.uploadManifestFn}); assertEqual(success.uploadedCount, 0, 'history 文件已在 full CAS baseline 时不应重复上传内容。'); for (const file of success.catalog.files) { assertTrue( harness.verifiedKeys.includes(file.objectKey), `history 清理前必须逐个验真 baseline 复用对象: ${file.path}`, ); } assertEqual(success.cleanup?.deletedCount, plan.candidates.length, 'catalog 和 baseline 验真后才应清理全部安全候选。'); const state = JSON.parse(readFileSync(success.statePath, 'utf8')); const historyCatalogObjectKey = state.historyCatalogs[0].objectKey; harness.objects.delete(historyCatalogObjectKey); let brokenHistoryFailure = null; try { await runDirectFilesBackup({...common, mode: 'history', uploadManifestFn: harness.uploadManifestFn}); } catch (error) { brokenHistoryFailure = error; } assertIncludes( brokenHistoryFailure?.message ?? '', `missing ${historyCatalogObjectKey}`, 'latest pointer 发布前必须重新验真所有 history catalog 引用。', ); } async function assertDirectHistoryWithoutCandidatesPublishesLatest() { const fixture = createHistoryFixture('direct-files-history-empty', {nestedData: false}); createReplicaHistory(fixture.replicasDir, '1', {snapshots: [10], segments: [0]}); const harness = createDirectOssHarness(); const common = { dataDir: fixture.dataDir, workDir: fixture.workDir, database: 'test-db', bucket: 'backup-bucket', objectPrefix: 'database-backups', uploadOptions: {}, uploadFn: harness.uploadFn, uploadManifestFn: harness.uploadManifestFn, verifyFn: harness.verifyFn, }; await runDirectFilesBackup({...common, mode: 'full'}); const latestObjectKey = 'database-backups/test-db/latest.json'; harness.objects.delete(latestObjectKey); const result = await runDirectFilesBackup({...common, mode: 'history'}); assertEqual(result.candidateCount, 0, 'fixture 应没有可归档 history 候选。'); assertTrue(harness.objects.has(latestObjectKey), 'history 无候选时仍必须从现有 state 发布 latest pointer。'); assertTrue(harness.verifiedKeys.includes(latestObjectKey), 'history 无候选时 latest pointer 仍必须 HEAD 验真。'); } async function assertDirectFilesRestoreDownloadsCatalogAndObjects() { const root = path.join(tmpRoot, 'direct-files-restore'); const dataDir = path.join(root, 'stdb'); const workDir = path.join(root, 'work'); const restoreDir = path.join(root, 'restore'); mkdirSync(path.join(dataDir, 'empty-directory'), {recursive: true}); mkdirSync(path.join(dataDir, 'config'), {recursive: true}); mkdirSync(path.join(dataDir, 'bin', '2.6.0'), {recursive: true}); symlinkSync('2.6.0', path.join(dataDir, 'bin', 'current')); const keyPath = path.join(dataDir, 'config', 'id_ecdsa'); writeFileSync(keyPath, 'private key fixture'); chmodSync(keyPath, 0o640); const harness = createDirectOssHarness(); await runDirectFilesBackup({ mode: 'full', dataDir, workDir, database: 'test-db', bucket: 'backup-bucket', objectPrefix: 'database-backups', uploadOptions: {}, uploadFn: harness.uploadFn, uploadManifestFn: harness.uploadManifestFn, verifyFn: harness.verifyFn, }); writeFileSync(keyPath, 'updated private key fixture'); chmodSync(keyPath, 0o640); const latestFull = await runDirectFilesBackup({ mode: 'full', dataDir, workDir, database: 'test-db', bucket: 'backup-bucket', objectPrefix: 'database-backups', uploadOptions: {}, uploadFn: harness.uploadFn, uploadManifestFn: harness.uploadManifestFn, verifyFn: harness.verifyFn, }); const restored = await restoreDirectFilesBackup({ statePath: latestFull.statePath, restoreDir, database: 'test-db', bucket: 'backup-bucket', uploadOptions: {}, downloadBufferFn: async ({objectKey}) => Buffer.from(harness.objects.get(objectKey)?.body ?? ''), downloadFileFn: async ({objectKey, destinationPath}) => { const object = harness.objects.get(objectKey); if (!object) { throw new Error(`missing ${objectKey}`); } writeFileSync(destinationPath, object.body); }, }); assertEqual(restored.downloadedCount, 1, 'files restore 必须从对象存储下载 catalog 中的普通文件。'); assertEqual(readFileSync(path.join(restoreDir, 'config', 'id_ecdsa'), 'utf8'), 'updated private key fixture', 'files restore 必须按最新 full catalog 的原相对路径恢复内容。'); assertTrue(existsSync(path.join(restoreDir, 'empty-directory')), 'files restore 必须重建空目录。'); assertEqual(statSync(path.join(restoreDir, 'config', 'id_ecdsa')).mode & 0o7777, 0o640, 'files restore 必须恢复文件权限。'); assertTrue(lstatSync(path.join(restoreDir, 'bin', 'current')).isSymbolicLink(), 'files restore 必须重建符号链接。'); assertEqual(readlinkSync(path.join(restoreDir, 'bin', 'current'), 'utf8'), '2.6.0', 'files restore 必须保留符号链接目标。'); rmSync(restoreDir, {recursive: true, force: true}); const downloadBufferFn = async ({objectKey}) => { const object = harness.objects.get(objectKey); if (!object) { throw new Error(`missing ${objectKey}`); } return Buffer.from(object.body); }; let objectDownloadCount = 0; const downloadFileFn = async ({objectKey, destinationPath}) => { objectDownloadCount += 1; const object = harness.objects.get(objectKey); if (!object) { throw new Error(`missing ${objectKey}`); } writeFileSync(destinationPath, object.body); }; const dryRun = await restoreDirectFilesLatest({ restoreDir, database: 'test-db', bucket: 'backup-bucket', objectPrefix: 'database-backups', uploadOptions: {}, dryRun: true, downloadBufferFn, downloadFileFn, verifyFn: harness.verifyFn, }); assertEqual(dryRun.catalogId, latestFull.catalogId, 'OSS-only dry-run 必须选择 latestFullCatalog。'); assertEqual(dryRun.fileCount, 1, 'OSS-only dry-run 应返回 full catalog 文件数。'); assertEqual(dryRun.symlinkCount, 1, 'OSS-only dry-run 应返回 full catalog 符号链接数。'); assertEqual(dryRun.totalSizeBytes, String(Buffer.byteLength('updated private key fixture')), 'OSS-only dry-run 应返回总字节数。'); assertEqual(objectDownloadCount, 0, 'OSS-only dry-run 不得下载数据对象。'); assertTrue(!existsSync(restoreDir), 'OSS-only dry-run 不得创建恢复目录。'); const latestRestored = await restoreDirectFilesLatest({ restoreDir, database: 'test-db', bucket: 'backup-bucket', objectPrefix: 'database-backups', uploadOptions: {}, downloadBufferFn, downloadFileFn, verifyFn: harness.verifyFn, }); assertEqual(latestRestored.catalogId, latestFull.catalogId, 'OSS-only restore 必须选择 latestFullCatalog。'); assertEqual(latestRestored.downloadedCount, 1, 'OSS-only restore 应下载 latest full catalog 的数据对象。'); assertEqual(readFileSync(path.join(restoreDir, 'config', 'id_ecdsa'), 'utf8'), 'updated private key fixture', 'OSS-only restore 应还原最新 full 内容。'); assertEqual(readlinkSync(path.join(restoreDir, 'bin', 'current'), 'utf8'), '2.6.0', 'OSS-only restore 应还原符号链接。'); } function assertCanonicalQueryAndAuthorizationIncludeMultipartParameters() { assertEqual(buildCanonicalQuery({uploads: null}), 'uploads', 'InitiateMultipartUpload 必须使用无等号的 uploads 参数。'); assertEqual( buildCanonicalQuery({uploadId: 'abc+/= xyz', partNumber: 12}), 'partNumber=12&uploadId=abc%2B%2F%3D%20xyz', 'multipart query 必须按 key 排序并使用 RFC3986 编码。', ); const date = new Date('2026-07-13T10:20:30.000Z'); const headers = { host: 'genarrative-test.oss-cn-shanghai.aliyuncs.com', 'content-type': 'application/octet-stream', 'x-oss-content-sha256': 'UNSIGNED-PAYLOAD', 'x-oss-date': '20260713T102030Z', }; const withoutQuery = buildAuthorization({ method: 'PUT', bucket: 'genarrative-test', endpoint: 'oss-cn-shanghai.aliyuncs.com', objectKey: 'database-backups/test/archive.tar.gz', accessKeyId: 'test-access-key', accessKeySecret: 'test-access-secret', headers, date, }); const withQuery = buildAuthorization({ method: 'PUT', bucket: 'genarrative-test', endpoint: 'oss-cn-shanghai.aliyuncs.com', objectKey: 'database-backups/test/archive.tar.gz', accessKeyId: 'test-access-key', accessKeySecret: 'test-access-secret', headers, date, queries: {partNumber: 12, uploadId: 'abc+/= xyz'}, }); assertNotEqual(withQuery, withoutQuery, 'multipart query 必须参与 V4 Authorization 计算。'); assertEqual( withQuery, 'OSS4-HMAC-SHA256 Credential=test-access-key/20260713/cn-shanghai/oss/aliyun_v4_request,AdditionalHeaders=host,Signature=9323dd3b7272b52f416c4d32115fcc00460eaccdcdaf011575c2502a63a27b1f', 'multipart V4 Authorization 必须保持固定签名向量。', ); } function assertInsufficientSpaceStopsBeforeServiceChanges() { const fixture = createFixture('insufficient-space'); const result = runBackup(fixture, [ '--stop-service', 'spacetimedb.service', '--restart-service-after', 'genarrative-api.service', '--min-free-bytes', '999999999999999999', ]); assertStatus(result, 1, '空间不足时必须失败。'); assertIncludes(result.stdout, '备份空间预检', '空间不足失败前应打印空间预检。'); assertIncludes(result.stderr, '剩余空间不足', '空间不足失败应说明剩余空间不足。'); assertFileMissing(fixture.systemctlLog, '空间不足时不能调用 systemctl。'); assertFileMissing(fixture.tarLog, '空间不足时不能调用 tar。'); } function assertArchiveFailureStillRestoresDependentServices() { const fixture = createFixture('tar-failure'); const result = runBackup(fixture, [ '--stop-service', 'spacetimedb.service', '--restart-service-after', 'genarrative-api.service', '--restart-service-after', 'genarrative-external-generation-worker@1.service', '--restart-service-after', 'genarrative-external-generation-controller.service', '--min-free-bytes', '1', ]); assertStatus(result, 1, 'tar 失败时备份脚本必须失败。'); assertIncludes(result.stderr, 'fake tar failure', 'tar 失败原因应保留在错误输出中。'); const systemctlLog = readFile(fixture.systemctlLog); const expectedCommands = [ 'systemctl stop spacetimedb.service', 'systemctl start spacetimedb.service', 'systemctl restart genarrative-api.service', 'systemctl restart genarrative-external-generation-worker@1.service', 'systemctl restart genarrative-external-generation-controller.service', ]; for (const command of expectedCommands) { assertIncludes(systemctlLog, command, `tar 失败后必须执行: ${command}`); } } async function assertMultipartUploadRetriesAndVerifiesRemoteLength() { const root = path.join(tmpRoot, 'multipart-success'); const archivePath = path.join(root, 'backup.tar.gz'); const partSizeBytes = 100 * 1024; const payload = Buffer.concat([ Buffer.alloc(partSizeBytes, 'a'), Buffer.alloc(partSizeBytes, 'b'), Buffer.alloc(17, 'c'), ]); const payloadSha256 = createHash('sha256').update(payload).digest('hex'); mkdirSync(root, {recursive: true}); writeFileSync(archivePath, payload); const requests = []; const retryDelays = []; let firstPartAttempts = 0; const uploadId = 'upload+/= id'; const fetchImpl = async (url, options) => { const body = await readRequestBody(options.body); requests.push({url, method: options.method, headers: options.headers, body}); const parsedUrl = new URL(url); if (options.method === 'POST' && parsedUrl.search === '?uploads') { return new Response(`${uploadId}`, {status: 200}); } if (options.method === 'PUT') { const partNumber = Number(parsedUrl.searchParams.get('partNumber')); if (partNumber === 1) { firstPartAttempts += 1; if (firstPartAttempts === 1) { return new Response('ServiceUnavailable', {status: 503}); } } return new Response('', {status: 200, headers: {etag: `"etag-${partNumber}"`}}); } if (options.method === 'POST' && parsedUrl.searchParams.has('uploadId')) { return new Response('', {status: 200, headers: {etag: '"complete-etag"'}}); } if (options.method === 'HEAD') { return new Response(null, {status: 200, headers: { 'content-length': String(payload.length), 'x-oss-meta-archive-sha256': payloadSha256, }}); } throw new Error(`unexpected request: ${options.method} ${url}`); }; const result = await uploadArchive({ archivePath, bucket: 'genarrative-test', endpoint: 'oss-cn-shanghai.aliyuncs.com', objectKey: 'database-backups/test/backup.tar.gz', accessKeyId: 'test-access-key', accessKeySecret: 'test-access-secret', partSizeBytes, maxAttempts: 3, retryBaseDelayMs: 1, retryMaxDelayMs: 1, fetchImpl, nowFn: () => new Date('2026-07-13T10:20:30.000Z'), sleepImpl: async (delayMs) => retryDelays.push(delayMs), randomFn: () => 0, }); assertEqual(result.uploadMode, 'multipart', '上传结果必须记录 multipart 模式。'); assertEqual(result.partCount, 3, 'multipart 应按配置大小切成三段。'); assertEqual(result.contentLength, payload.length, '上传结果应保留完整归档长度。'); assertEqual(result.etag, 'complete-etag', '上传结果应保留 CompleteMultipartUpload ETag。'); assertEqual(firstPartAttempts, 2, '503 后应仅重试失败的第一段。'); assertEqual(retryDelays.length, 1, '一次可重试失败应触发一次退避。'); const initiateRequest = requests[0]; assertTrue(initiateRequest.url.endsWith('?uploads'), 'InitiateMultipartUpload URL 必须使用裸 uploads 参数。'); assertTrue(!initiateRequest.url.endsWith('?uploads='), 'InitiateMultipartUpload URL 不能把裸参数写成 uploads=。'); assertEqual(initiateRequest.headers['x-oss-meta-archive-sha256'], payloadSha256, 'multipart 对象必须保存本地归档 SHA-256 元数据。'); const firstPartRequests = requests.filter(({method, url}) => method === 'PUT' && new URL(url).searchParams.get('partNumber') === '1'); assertEqual(firstPartRequests.length, 2, '第一段应产生原请求和一次重试。'); assertBufferEqual(firstPartRequests[0].body, payload.subarray(0, partSizeBytes), '第一段原请求内容必须完整。'); assertBufferEqual(firstPartRequests[1].body, payload.subarray(0, partSizeBytes), '第一段重试必须重新创建并完整读取 stream。'); assertTrue( firstPartRequests[0].url.includes('?partNumber=1&uploadId=upload%2B%2F%3D%20id'), 'UploadPart URL 必须使用排序并编码后的 canonical query。', ); const completeRequest = requests.find(({method, url}) => method === 'POST' && new URL(url).searchParams.has('uploadId')); assertIncludes(completeRequest?.body.toString('utf8') ?? '', '1"etag-1"', 'Complete XML 应包含第一段 ETag。'); assertIncludes(completeRequest?.body.toString('utf8') ?? '', '3"etag-3"', 'Complete XML 应包含最后一段 ETag。'); assertTrue(requests.some(({method}) => method === 'HEAD'), 'Complete 后必须执行签名 HEAD 验证。'); for (const request of requests) { assertTrue(String(request.headers.authorization ?? '').startsWith('OSS4-HMAC-SHA256 '), `${request.method} 请求必须携带 V4 Authorization。`); } } async function assertHeadLengthMismatchAbortsMultipartUpload() { const root = path.join(tmpRoot, 'multipart-head-mismatch'); const archivePath = path.join(root, 'backup.tar.gz'); const partSizeBytes = 100 * 1024; const payload = Buffer.alloc(partSizeBytes + 1, 'x'); const payloadSha256 = createHash('sha256').update(payload).digest('hex'); mkdirSync(root, {recursive: true}); writeFileSync(archivePath, payload); const requests = []; const fetchImpl = async (url, options) => { await readRequestBody(options.body); requests.push({url, method: options.method}); const parsedUrl = new URL(url); if (options.method === 'POST' && parsedUrl.search === '?uploads') { return new Response('mismatch-upload', {status: 200}); } if (options.method === 'PUT') { return new Response('', {status: 200, headers: {etag: `"etag-${parsedUrl.searchParams.get('partNumber')}"`}}); } if (options.method === 'POST') { return new Response('', {status: 200, headers: {etag: '"complete-etag"'}}); } if (options.method === 'HEAD') { return new Response(null, {status: 200, headers: { 'content-length': String(payload.length - 1), 'x-oss-meta-archive-sha256': payloadSha256, }}); } if (options.method === 'DELETE') { return new Response(null, {status: 204}); } throw new Error(`unexpected request: ${options.method} ${url}`); }; let uploadError = null; try { await uploadArchive({ archivePath, bucket: 'genarrative-test', endpoint: 'oss-cn-shanghai.aliyuncs.com', objectKey: 'database-backups/test/mismatch.tar.gz', accessKeyId: 'test-access-key', accessKeySecret: 'test-access-secret', partSizeBytes, maxAttempts: 2, retryBaseDelayMs: 1, retryMaxDelayMs: 1, fetchImpl, nowFn: () => new Date('2026-07-13T10:20:30.000Z'), sleepImpl: async () => {}, randomFn: () => 0, }); } catch (error) { uploadError = error; } assertTrue(uploadError instanceof Error, 'HEAD 长度不一致时上传必须失败。'); assertIncludes(uploadError?.message ?? '', 'HEAD 验证长度不一致', 'HEAD 长度不一致错误应保留本地和远端长度。'); const abortRequest = requests.find(({method}) => method === 'DELETE'); assertTrue(Boolean(abortRequest), 'HEAD 长度不一致后必须 best-effort AbortMultipartUpload。'); assertTrue(abortRequest?.url.endsWith('?uploadId=mismatch-upload'), 'AbortMultipartUpload 必须携带同一 uploadId。'); } async function assertHeadShaMismatchAbortsMultipartUpload() { const root = path.join(tmpRoot, 'multipart-head-sha-mismatch'); const archivePath = path.join(root, 'backup.tar.gz'); const payload = Buffer.alloc(100 * 1024, 's'); mkdirSync(root, {recursive: true}); writeFileSync(archivePath, payload); const requests = []; const fetchImpl = async (url, options) => { await readRequestBody(options.body); requests.push({url, method: options.method}); const parsedUrl = new URL(url); if (options.method === 'POST' && parsedUrl.search === '?uploads') { return new Response('sha-mismatch-upload', {status: 200}); } if (options.method === 'PUT') { return new Response('', {status: 200, headers: {etag: '"part-etag"'}}); } if (options.method === 'POST') { return new Response('', {status: 200}); } if (options.method === 'HEAD') { return new Response(null, {status: 200, headers: { 'content-length': String(payload.length), 'x-oss-meta-archive-sha256': '0'.repeat(64), }}); } if (options.method === 'DELETE') { return new Response(null, {status: 204}); } throw new Error(`unexpected request: ${options.method} ${url}`); }; let uploadError = null; try { await uploadArchive({ archivePath, bucket: 'genarrative-test', endpoint: 'oss-cn-shanghai.aliyuncs.com', objectKey: 'database-backups/test/sha-mismatch.tar.gz', accessKeyId: 'test-access-key', accessKeySecret: 'test-access-secret', partSizeBytes: 100 * 1024, maxAttempts: 1, fetchImpl, nowFn: () => new Date('2026-07-13T10:20:30.000Z'), sleepImpl: async () => {}, randomFn: () => 0, }); } catch (error) { uploadError = error; } assertIncludes(uploadError?.message ?? '', 'SHA-256 不一致', 'HEAD SHA-256 不一致时上传必须失败。'); assertTrue( requests.some(({method, url}) => method === 'DELETE' && url.endsWith('?uploadId=sha-mismatch-upload')), 'HEAD SHA-256 不一致后必须 best-effort AbortMultipartUpload。', ); } async function assertManifestUploadUsesShaAndHeadVerification() { const root = path.join(tmpRoot, 'manifest-upload'); const manifestPath = path.join(root, 'backup.manifest.json'); const body = Buffer.from(JSON.stringify({uploadStatus: 'uploaded', catalog: 'x'.repeat(150 * 1024)})); const bodySha256 = createHash('sha256').update(body).digest('hex'); mkdirSync(root, {recursive: true}); writeFileSync(manifestPath, body); const requests = []; let limitedChunkCount = 0; let limitedBytes = 0; const result = await uploadManifestFile({ manifestPath, bucket: 'genarrative-test', endpoint: 'oss-cn-shanghai.aliyuncs.com', objectKey: 'database-backups/test/backup.tar.gz.manifest.json', accessKeyId: 'test-access-key', accessKeySecret: 'test-access-secret', maxAttempts: 1, nowFn: () => new Date('2026-07-13T10:20:30.000Z'), sleepImpl: async () => {}, randomFn: () => 0, bandwidthLimiter: { wrap(readable) { return Readable.from((async function* observeLimitedManifest() { for await (const chunk of readable) { limitedChunkCount += 1; limitedBytes += chunk.length; yield chunk; } })(), {objectMode: false}); }, }, fetchImpl: async (url, options) => { const requestBody = await readRequestBody(options.body); requests.push({url, method: options.method, headers: options.headers, body: requestBody}); if (options.method === 'PUT') { return new Response(null, {status: 200}); } if (options.method === 'HEAD') { return new Response(null, {status: 200, headers: { 'x-oss-meta-file-size': String(body.length), 'x-oss-meta-archive-sha256': bodySha256, }}); } throw new Error(`unexpected request: ${options.method} ${url}`); }, }); assertEqual(result.archiveSha256, bodySha256, 'manifest 上传结果必须记录本地 SHA-256。'); assertBufferEqual(requests.find(({method}) => method === 'PUT')?.body, body, 'manifest PUT 必须上传完整 JSON。'); assertEqual(limitedBytes, body.length, 'manifest 必须完整经过上传带宽限制流。'); assertTrue(limitedChunkCount > 1, '大型 manifest 必须分块经过限速器,不能整块突发上传。'); assertTrue(requests.some(({method}) => method === 'HEAD'), 'manifest PUT 后必须执行 HEAD 验真。'); assertEqual( requests.find(({method}) => method === 'PUT')?.headers['x-oss-meta-file-size'], String(body.length), 'manifest PUT 必须记录原始字节数,供动态压缩 HEAD 缺少 content-length 时验真。', ); } async function assertMissingPartEtagAbortsMultipartUpload() { const root = path.join(tmpRoot, 'multipart-missing-etag'); const archivePath = path.join(root, 'backup.tar.gz'); mkdirSync(root, {recursive: true}); writeFileSync(archivePath, Buffer.alloc(100 * 1024, 'e')); const requests = []; const fetchImpl = async (url, options) => { await readRequestBody(options.body); requests.push({url, method: options.method}); const parsedUrl = new URL(url); if (options.method === 'POST' && parsedUrl.search === '?uploads') { return new Response('missing-etag-upload', {status: 200}); } if (options.method === 'PUT') { return new Response('', {status: 200}); } if (options.method === 'DELETE') { return new Response(null, {status: 204}); } throw new Error(`unexpected request: ${options.method} ${url}`); }; let uploadError = null; try { await uploadArchive({ archivePath, bucket: 'genarrative-test', endpoint: 'oss-cn-shanghai.aliyuncs.com', objectKey: 'database-backups/test/missing-etag.tar.gz', accessKeyId: 'test-access-key', accessKeySecret: 'test-access-secret', partSizeBytes: 100 * 1024, maxAttempts: 1, fetchImpl, nowFn: () => new Date('2026-07-13T10:20:30.000Z'), sleepImpl: async () => {}, randomFn: () => 0, }); } catch (error) { uploadError = error; } assertIncludes(uploadError?.message ?? '', '响应缺少 ETag', 'UploadPart 缺少 ETag 时必须失败。'); assertTrue( requests.some(({method, url}) => method === 'DELETE' && url.endsWith('?uploadId=missing-etag-upload')), 'UploadPart 缺少 ETag 后必须 AbortMultipartUpload。', ); } async function assertCompleteResponseAmbiguityUsesHeadVerification() { const root = path.join(tmpRoot, 'multipart-complete-ambiguity'); const archivePath = path.join(root, 'backup.tar.gz'); const payload = Buffer.alloc(100 * 1024, 'c'); const payloadSha256 = createHash('sha256').update(payload).digest('hex'); mkdirSync(root, {recursive: true}); writeFileSync(archivePath, payload); const requests = []; let completeAttempts = 0; const fetchImpl = async (url, options) => { await readRequestBody(options.body); requests.push({url, method: options.method}); const parsedUrl = new URL(url); if (options.method === 'POST' && parsedUrl.search === '?uploads') { return new Response('ambiguous-upload', {status: 200}); } if (options.method === 'PUT') { return new Response('', {status: 200, headers: {etag: '"part-etag"'}}); } if (options.method === 'POST' && parsedUrl.searchParams.has('uploadId')) { completeAttempts += 1; if (completeAttempts === 1) { throw new TypeError('socket closed after remote complete'); } return new Response('NoSuchUpload', {status: 404}); } if (options.method === 'HEAD') { return new Response(null, {status: 200, headers: { 'content-length': String(payload.length), 'x-oss-meta-archive-sha256': payloadSha256, }}); } if (options.method === 'DELETE') { return new Response(null, {status: 204}); } throw new Error(`unexpected request: ${options.method} ${url}`); }; const result = await uploadArchive({ archivePath, bucket: 'genarrative-test', endpoint: 'oss-cn-shanghai.aliyuncs.com', objectKey: 'database-backups/test/complete-ambiguity.tar.gz', accessKeyId: 'test-access-key', accessKeySecret: 'test-access-secret', partSizeBytes: 100 * 1024, maxAttempts: 2, retryBaseDelayMs: 1, retryMaxDelayMs: 1, fetchImpl, nowFn: () => new Date('2026-07-13T10:20:30.000Z'), sleepImpl: async () => {}, randomFn: () => 0, }); assertEqual(completeAttempts, 2, 'Complete 网络错误后应按策略重试。'); assertEqual(result.contentLength, payload.length, 'Complete 结果不确定时应以 HEAD 长度验真收口。'); assertTrue(requests.some(({method}) => method === 'HEAD'), 'Complete 结果不确定时必须执行 HEAD 验真。'); assertTrue(!requests.some(({method}) => method === 'DELETE'), 'HEAD 已证实对象完整时不得 Abort 已完成上传。'); } function assertHistoryDiscoversDevAndProductionLayoutsWithMultipleReplicas() { const dev = createHistoryFixture('history-dev-layout', {nestedData: false}); createReplicaHistory(dev.replicasDir, '1', { snapshots: [0, 187, 279], segments: [0, 188, 280], }); createReplicaHistory(dev.replicasDir, '2', { snapshots: [50, 99], segments: [0, 51, 100], }); const devPlan = discoverHistoryPlan({dataDir: dev.dataDir}); assertEqual(devPlan.replicas.length, 2, 'history 应逐 replica 计算安全边界。'); assertEqual(devPlan.candidates.length, 7, '多 replica history 候选数量必须符合 snapshot/segment 边界。'); assertTrue( devPlan.candidates.some(({path}) => path === 'replicas/1/clog/00000000000000000000.stdb.log'), 'dev 布局应识别边界 segment 之前的 commitlog。', ); assertTrue( !devPlan.candidates.some(({path}) => path.includes('00000000000000000188.stdb.log')), '跨越 latest snapshot 的边界 segment 必须保留。', ); assertTrue( !devPlan.candidates.some(({path}) => path.includes('00000000000000000279.snapshot_dir')), '每个 replica 的 latest snapshot 必须保留。', ); const production = createHistoryFixture('history-production-layout', {nestedData: true}); createReplicaHistory(production.replicasDir, '7', { snapshots: [10, 20], segments: [0, 11, 21], }); const productionPlan = discoverHistoryPlan({dataDir: production.dataDir}); assertEqual(productionPlan.replicasDir, 'data/replicas', 'history 必须兼容 /stdb/data/replicas 布局。'); assertEqual(productionPlan.candidates.length, 3, 'production 布局应识别一个旧 snapshot 与一对旧 commitlog 文件。'); const importResult = runHistoryDryRun(dev); assertStatus(importResult, 0, 'history dry-run 应能从已有 uploaded baseline manifest 导入 state。'); assertTrue(existsSync(dev.statePath), 'history dry-run 应持久化导入后的 baseline state。'); assertIncludes(importResult.stdout, 'history dry-run', 'history dry-run 应明确说明不会上传或删除。'); for (const candidate of devPlan.candidates) { assertTrue(existsSync(path.join(dev.dataDir, candidate.path)), `history dry-run 不得删除候选: ${candidate.path}`); } } function assertHistorySkipsReplicaWithoutSnapshotAndRejectsMalformedNames() { const noSnapshot = createHistoryFixture('history-no-snapshot', {nestedData: false}); createReplicaHistory(noSnapshot.replicasDir, '1', {snapshots: [], segments: [0]}); const plan = discoverHistoryPlan({dataDir: noSnapshot.dataDir}); assertEqual(plan.candidates.length, 0, '没有 snapshot 的 replica 不得产生可删除候选。'); assertEqual(plan.replicas[0]?.reason, 'no-snapshot', '没有 snapshot 时应记录明确跳过原因。'); const incompleteSnapshot = createHistoryFixture('history-incomplete-snapshot', {nestedData: false}); const incompleteReplica = createReplicaHistory(incompleteSnapshot.replicasDir, '1', {snapshots: [0, 10], segments: [0, 1, 11]}); mkdirSync(path.join(incompleteReplica.snapshotsDir, '00000000000000000020.snapshot_dir')); mkdirSync(path.join(incompleteReplica.snapshotsDir, '00000000000000000030.snapshot_dir')); writeFileSync(path.join(incompleteReplica.snapshotsDir, '00000000000000000030.snapshot_dir', '00000000000000000030.snapshot_bsatn'), 'locked'); writeFileSync(path.join(incompleteReplica.snapshotsDir, '00000000000000000030.lock'), `${process.pid}\n`); const incompletePlan = discoverHistoryPlan({dataDir: incompleteSnapshot.dataDir}); assertEqual(incompletePlan.replicas[0]?.latestSnapshot, '10', '缺少 snapshot_bsatn 或仍有 lockfile 的目录不得成为 latest snapshot。'); assertTrue( !incompletePlan.candidates.some(({path}) => path.includes('00000000000000000010.snapshot_dir')), '最后一个完整且未锁定的 snapshot 必须保留。', ); const malformedLog = createHistoryFixture('history-malformed-log', {nestedData: false}); const malformedLogReplica = createReplicaHistory(malformedLog.replicasDir, '1', {snapshots: [10], segments: [0, 11]}); writeFileSync(path.join(malformedLogReplica.clogDir, 'broken.stdb.log'), 'broken'); assertThrows( () => discoverHistoryPlan({dataDir: malformedLog.dataDir}), 'commitlog 文件名不符合预期', '异常 commitlog 名称必须阻断整个清理计划。', ); } function assertHistoryRequiresBaselineAndProducesDeterministicDeferredBatch() { const missingBaseline = createHistoryFixture('history-missing-baseline', {nestedData: false}); createReplicaHistory(missingBaseline.replicasDir, '1', {snapshots: [0, 10], segments: [0, 1, 11]}); rmSync(missingBaseline.baselineManifestPath, {force: true}); const missingResult = runHistoryCommand(missingBaseline, ['--dry-run'], {includeBaselineManifest: false}); assertStatus(missingResult, 1, 'history 没有 baseline state 或 imported manifest 时必须失败。'); assertIncludes(missingResult.stderr, '缺少已验真 baseline state', 'baseline 门禁失败应给出明确错误。'); const wrongKind = createHistoryFixture('history-wrong-baseline-kind', {nestedData: false}); createReplicaHistory(wrongKind.replicasDir, '1', {snapshots: [0, 10], segments: [0, 1, 11]}); const wrongKindManifest = JSON.parse(readFileSync(wrongKind.baselineManifestPath, 'utf8')); wrongKindManifest.backupKind = 'spacetimedb-history'; writeFileSync(wrongKind.baselineManifestPath, `${JSON.stringify(wrongKindManifest)}\n`); const wrongKindResult = runHistoryDryRun(wrongKind); assertStatus(wrongKindResult, 1, 'history archive manifest 不得被导入为 full baseline。'); assertIncludes(wrongKindResult.stderr, 'backupKind 必须是 spacetimedb-data-dir', 'baseline 类型不匹配应失败关闭。'); const deterministic = createHistoryFixture('history-deterministic-batch', {nestedData: false}); createReplicaHistory(deterministic.replicasDir, '1', {snapshots: [0, 10], segments: [0, 1, 11]}); importHistoryState(deterministic); const firstResultFile = path.join(deterministic.workDir, 'defer-first.json'); const secondResultFile = path.join(deterministic.workDir, 'defer-second.json'); const first = runHistoryCommand(deterministic, ['--defer-upload', '--result-file', firstResultFile]); const second = runHistoryCommand(deterministic, ['--defer-upload', '--result-file', secondResultFile]); assertStatus(first, 0, '第一次 history defer 应成功生成归档。'); assertStatus(second, 0, '相同候选重复 history defer 应幂等复用 batch identity。'); const firstPayload = JSON.parse(readFileSync(firstResultFile, 'utf8')); const secondPayload = JSON.parse(readFileSync(secondResultFile, 'utf8')); assertEqual(firstPayload.batchId, secondPayload.batchId, '相同 baseline 与候选必须生成确定性 batchId。'); assertEqual(firstPayload.objectKey, secondPayload.objectKey, '相同 batch 重跑不得制造新的 OSS object key。'); const archiveListing = spawnSync('tar', ['-tzf', firstPayload.archivePath], {encoding: 'utf8'}); assertStatus(archiveListing, 0, 'history 归档应可被 tar 正常读取。'); assertIncludes(archiveListing.stdout, path.basename(firstPayload.manifestPath), 'history 归档内部必须携带安全候选 manifest。'); const dryRunPending = createHistoryFixture('history-dry-run-pending-cleanup', {nestedData: false}); createReplicaHistory(dryRunPending.replicasDir, '1', {snapshots: [0, 10], segments: [0, 1, 11]}); const dryRunState = importHistoryState(dryRunPending); const dryRunPlan = discoverHistoryPlan({dataDir: dryRunPending.dataDir}); dryRunState.batches.push({ batchId: 'pending-cleanup', objectKey: 'database-backups/test-db/history/pending-cleanup.tar.gz', contentLength: 10, verifiedAt: '2026-07-16T00:20:00.000Z', uploadedAt: '2026-07-16T00:20:00.000Z', status: 'uploaded', candidates: dryRunPlan.candidates, }); writeFileSync(dryRunPending.statePath, `${JSON.stringify(dryRunState)}\n`); const pendingDryRunResult = runHistoryDryRun(dryRunPending); assertStatus(pendingDryRunResult, 0, '存在待清理 uploaded batch 时 history dry-run 仍应只读成功。'); for (const candidate of dryRunPlan.candidates) { assertTrue(existsSync(path.join(dryRunPending.dataDir, candidate.path)), `history dry-run 不得恢复执行待清理 batch: ${candidate.path}`); } } function assertHistoryBackupLockRejectsLiveAndStaleOwners() { const liveOwner = createHistoryFixture('history-live-lock', {nestedData: false}); createReplicaHistory(liveOwner.replicasDir, '1', {snapshots: [0, 10], segments: [0, 1, 11]}); importHistoryState(liveOwner); const liveLockPath = path.join(liveOwner.workDir, 'test-db.backup.lock'); writeFileSync(liveLockPath, `${process.pid}\n`); const liveResult = runHistoryCommand(liveOwner, ['--defer-upload']); assertStatus(liveResult, 1, '仍存活进程持有 backup lock 时必须拒绝并发备份。'); assertIncludes(liveResult.stderr, '已有数据库备份进程持有锁', '并发备份失败应报告 lock owner pid。'); const staleOwner = createHistoryFixture('history-stale-lock', {nestedData: false}); createReplicaHistory(staleOwner.replicasDir, '1', {snapshots: [0, 10], segments: [0, 1, 11]}); importHistoryState(staleOwner); const staleLockPath = path.join(staleOwner.workDir, 'test-db.backup.lock'); writeFileSync(staleLockPath, '2147483647\n'); const staleResult = runHistoryCommand(staleOwner, ['--defer-upload']); assertStatus(staleResult, 1, '失效 owner pid 的 backup lock 也必须失败关闭,避免并发抢锁。'); assertIncludes(staleResult.stderr, '拒绝自动抢锁', '失效 backup lock 应要求人工核对 multipart 与进程。'); assertTrue(existsSync(staleLockPath), '失效 backup lock 未经人工核对不得自动删除。'); } function assertHistoryStatDriftPreventsAnyCleanup() { const fixture = createHistoryFixture('history-stat-drift', {nestedData: false}); createReplicaHistory(fixture.replicasDir, '1', {snapshots: [0, 10], segments: [0, 1, 11]}); const plan = discoverHistoryPlan({dataDir: fixture.dataDir}); const driftCandidate = plan.candidates.find(({kind}) => kind === 'commitlog'); const untouchedCandidate = plan.candidates.find(({kind}) => kind === 'snapshot'); writeFileSync(path.join(fixture.dataDir, driftCandidate.path), 'changed-after-plan'); assertThrows( () => cleanupHistoryCandidates({dataDir: fixture.dataDir, candidates: plan.candidates}), 'stat 漂移', '任一候选 stat 漂移时必须在删除任何文件前失败。', ); assertTrue(existsSync(path.join(fixture.dataDir, untouchedCandidate.path)), 'stat 漂移失败时不得删除其他候选。'); } async function assertHistoryUploadFailureDoesNotDeleteSources() { const fixture = createHistoryFixture('history-upload-failure', {nestedData: false}); createReplicaHistory(fixture.replicasDir, '1', {snapshots: [0, 10], segments: [0, 1, 11]}); const state = importHistoryState(fixture); const plan = discoverHistoryPlan({dataDir: fixture.dataDir}); const archivePath = path.join(fixture.workDir, 'history.tar.gz'); const manifestPath = `${archivePath}.manifest.json`; writeFileSync(archivePath, 'history archive'); const manifest = createHistoryManifest({fixture, state, plan, archivePath}); writeFileSync(manifestPath, `${JSON.stringify(manifest)}\n`); let uploadError = null; try { await uploadHistoryArchiveWithCleanup({ archivePath, manifestPath, manifest, statePath: fixture.statePath, uploadOptions: {}, uploadFn: async () => { throw new Error('synthetic upload failure'); }, manifestUploadFn: async () => { throw new Error('manifest upload must not run after archive failure'); }, }); } catch (error) { uploadError = error; } assertIncludes(uploadError?.message ?? '', 'synthetic upload failure', 'history 应保留上传失败原因。'); for (const candidate of plan.candidates) { assertTrue(existsSync(path.join(fixture.dataDir, candidate.path)), `上传失败不得删除 history 源文件: ${candidate.path}`); } const stateAfterFailure = JSON.parse(readFileSync(fixture.statePath, 'utf8')); assertEqual(stateAfterFailure.batches.length, 0, '上传失败不得把 batch 标记为 uploaded。'); const manifestFailure = createHistoryFixture('history-manifest-upload-failure', {nestedData: false}); createReplicaHistory(manifestFailure.replicasDir, '1', {snapshots: [0, 10], segments: [0, 1, 11]}); const manifestFailureState = importHistoryState(manifestFailure); const manifestFailurePlan = discoverHistoryPlan({dataDir: manifestFailure.dataDir}); const manifestFailureArchive = path.join(manifestFailure.workDir, 'history.tar.gz'); const manifestFailurePath = `${manifestFailureArchive}.manifest.json`; writeFileSync(manifestFailureArchive, 'history archive'); const manifestFailurePayload = createHistoryManifest({ fixture: manifestFailure, state: manifestFailureState, plan: manifestFailurePlan, archivePath: manifestFailureArchive, }); writeFileSync(manifestFailurePath, `${JSON.stringify(manifestFailurePayload)}\n`); let manifestUploadError = null; try { await uploadHistoryArchiveWithCleanup({ archivePath: manifestFailureArchive, manifestPath: manifestFailurePath, manifest: manifestFailurePayload, statePath: manifestFailure.statePath, uploadOptions: {}, uploadFn: async () => ({ bucket: 'backup-bucket', objectKey: manifestFailurePayload.objectKey, contentLength: 15, archiveSha256: 'c'.repeat(64), verifiedAt: '2026-07-16T00:05:00.000Z', }), manifestUploadFn: async () => { throw new Error('synthetic manifest upload failure'); }, }); } catch (error) { manifestUploadError = error; } assertIncludes(manifestUploadError?.message ?? '', 'synthetic manifest upload failure', 'sidecar manifest 上传失败应阻断清理。'); for (const candidate of manifestFailurePlan.candidates) { assertTrue(existsSync(path.join(manifestFailure.dataDir, candidate.path)), `manifest 上传失败不得删除源文件: ${candidate.path}`); } } async function assertHistorySuccessfulUploadCleansAndIsIdempotent() { const fixture = createHistoryFixture('history-upload-success', {nestedData: false}); createReplicaHistory(fixture.replicasDir, '1', {snapshots: [0, 10], segments: [0, 1, 11]}); const state = importHistoryState(fixture); const plan = discoverHistoryPlan({dataDir: fixture.dataDir}); const archivePath = path.join(fixture.workDir, 'history.tar.gz'); const manifestPath = `${archivePath}.manifest.json`; writeFileSync(archivePath, 'history archive'); const manifest = createHistoryManifest({fixture, state, plan, archivePath}); writeFileSync(manifestPath, `${JSON.stringify(manifest)}\n`); let baselineVerifyCount = 0; const result = await uploadHistoryArchiveWithCleanup({ archivePath, manifestPath, manifest, statePath: fixture.statePath, uploadOptions: {}, uploadFn: async () => ({ bucket: 'backup-bucket', objectKey: manifest.objectKey, contentLength: 15, archiveSha256: 'b'.repeat(64), etag: 'test-etag', uploadMode: 'multipart', partCount: 1, partSizeBytes: 102400, verifiedAt: '2026-07-16T00:10:00.000Z', }), manifestUploadFn: async ({objectKey}) => ({ objectKey, contentLength: 512, archiveSha256: 'd'.repeat(64), verifiedAt: '2026-07-16T00:10:01.000Z', }), verifyFn: async () => { baselineVerifyCount += 1; return {verifiedAt: '2026-07-16T00:10:02.000Z'}; }, }); assertEqual(baselineVerifyCount, 2, 'history 删除源文件前必须重新验真 full baseline 与 sidecar。'); assertEqual(result.cleanup.deletedCount, plan.candidates.length, '验真上传成功后应删除全部安全候选。'); for (const candidate of plan.candidates) { assertTrue(!existsSync(path.join(fixture.dataDir, candidate.path)), `验真成功后应删除 history 源文件: ${candidate.path}`); } const repeatedCleanup = cleanupHistoryCandidates({dataDir: fixture.dataDir, candidates: plan.candidates}); assertEqual(repeatedCleanup.alreadyMissingCount, plan.candidates.length, '重复清理同一 uploaded batch 应幂等。'); } async function assertHistoryResumeReverifiesArchiveAndManifest() { const fixture = createHistoryFixture('history-resume-verification', {nestedData: false}); createReplicaHistory(fixture.replicasDir, '1', {snapshots: [0, 10], segments: [0, 1, 11]}); const state = importHistoryState(fixture); const plan = discoverHistoryPlan({dataDir: fixture.dataDir}); state.batches.push({ batchId: 'resume-batch', objectKey: 'database-backups/test-db/history/resume.tar.gz', contentLength: 100, archiveSha256: 'e'.repeat(64), verifiedAt: '2026-07-16T00:30:00.000Z', manifestObjectKey: 'database-backups/test-db/history/resume.tar.gz.manifest.json', manifestContentLength: 200, manifestArchiveSha256: 'f'.repeat(64), manifestVerifiedAt: '2026-07-16T00:30:01.000Z', uploadedAt: '2026-07-16T00:30:00.000Z', status: 'uploaded', candidates: plan.candidates, }); writeFileSync(fixture.statePath, `${JSON.stringify(state)}\n`); const verifiedKeys = []; await resumeUploadedHistoryBatch({ statePath: fixture.statePath, state, dataDir: fixture.dataDir, verificationOptions: {}, verifyFn: async ({objectKey}) => { verifiedKeys.push(objectKey); return {verifiedAt: '2026-07-16T00:31:00.000Z'}; }, }); assertEqual(verifiedKeys.length, 2, '续清理前必须重新验真 history archive 与 sidecar manifest。'); for (const candidate of plan.candidates) { assertTrue(!existsSync(path.join(fixture.dataDir, candidate.path)), `续清理验真后应删除候选: ${candidate.path}`); } const failure = createHistoryFixture('history-resume-verification-failure', {nestedData: false}); createReplicaHistory(failure.replicasDir, '1', {snapshots: [0, 10], segments: [0, 1, 11]}); const failureState = importHistoryState(failure); const failurePlan = discoverHistoryPlan({dataDir: failure.dataDir}); failureState.batches.push({...state.batches[0], candidates: failurePlan.candidates}); writeFileSync(failure.statePath, `${JSON.stringify(failureState)}\n`); let resumeError = null; try { await resumeUploadedHistoryBatch({ statePath: failure.statePath, state: failureState, dataDir: failure.dataDir, verificationOptions: {}, verifyFn: async () => { throw new Error('synthetic resume HEAD failure'); }, }); } catch (error) { resumeError = error; } assertIncludes(resumeError?.message ?? '', 'synthetic resume HEAD failure', '续清理 OSS 复核失败应保留错误。'); for (const candidate of failurePlan.candidates) { assertTrue(existsSync(path.join(failure.dataDir, candidate.path)), `续清理验真失败不得删除候选: ${candidate.path}`); } } function createHistoryFixture(name, {nestedData}) { const root = path.join(tmpRoot, name); const dataDir = path.join(root, 'stdb'); const replicasDir = nestedData ? path.join(dataDir, 'data', 'replicas') : path.join(dataDir, 'replicas'); const workDir = path.join(root, 'work'); const statePath = path.join(workDir, 'history-state.json'); const baselineManifestPath = path.join(workDir, 'baseline.manifest.json'); mkdirSync(replicasDir, {recursive: true}); mkdirSync(workDir, {recursive: true}); writeFileSync(baselineManifestPath, `${JSON.stringify({ backupKind: 'spacetimedb-data-dir', uploadStatus: 'uploaded', database: 'test-db', dataDir, bucket: 'backup-bucket', objectKey: 'database-backups/test-db/baseline.tar.gz', manifestObjectKey: 'database-backups/test-db/baseline.tar.gz.manifest.json', contentLength: 1234, archiveSha256: 'a'.repeat(64), manifestContentLength: 512, manifestArchiveSha256: '9'.repeat(64), manifestVerifiedAt: '2026-07-16T00:00:00.500Z', verifiedAt: '2026-07-16T00:00:00.000Z', uploadedAt: '2026-07-16T00:00:01.000Z', }, null, 2)}\n`); return {root, dataDir, replicasDir, workDir, statePath, baselineManifestPath}; } function createReplicaHistory(replicasDir, replicaId, {snapshots, segments}) { const replicaDir = path.join(replicasDir, replicaId); const snapshotsDir = path.join(replicaDir, 'snapshots'); const clogDir = path.join(replicaDir, 'clog'); mkdirSync(snapshotsDir, {recursive: true}); mkdirSync(clogDir, {recursive: true}); for (const transaction of snapshots) { const name = `${String(transaction).padStart(20, '0')}.snapshot_dir`; const snapshotDir = path.join(snapshotsDir, name); mkdirSync(path.join(snapshotDir, 'objects'), {recursive: true}); writeFileSync(path.join(snapshotDir, `${String(transaction).padStart(20, '0')}.snapshot_bsatn`), `snapshot-${transaction}`); writeFileSync(path.join(snapshotDir, 'objects', 'object.bin'), `object-${transaction}`); } for (const transaction of segments) { const prefix = String(transaction).padStart(20, '0'); writeFileSync(path.join(clogDir, `${prefix}.stdb.log`), `log-${transaction}`); writeFileSync(path.join(clogDir, `${prefix}.stdb.ofs`), `ofs-${transaction}`); } return {replicaDir, snapshotsDir, clogDir}; } function runHistoryDryRun(fixture) { const resultFile = path.join(fixture.workDir, 'dry-run-result.json'); return runHistoryCommand(fixture, [ '--result-file', resultFile, '--dry-run', ]); } function runHistoryCommand(fixture, extraArgs = [], {includeBaselineManifest = true} = {}) { const baselineManifestArgs = includeBaselineManifest ? ['--baseline-manifest', fixture.baselineManifestPath] : []; return spawnSync(process.execPath, [ BACKUP_SCRIPT, '--mode', 'history', '--data-dir', fixture.dataDir, '--work-dir', fixture.workDir, '--database', 'test-db', '--bucket', 'backup-bucket', '--endpoint', 'oss-cn-shanghai.aliyuncs.com', '--access-key-id', 'test-access-key', '--access-key-secret', 'test-access-secret', '--baseline-state', fixture.statePath, ...baselineManifestArgs, ...extraArgs, ], {encoding: 'utf8'}); } function importHistoryState(fixture) { const result = runHistoryDryRun(fixture); assertStatus(result, 0, '测试 fixture 应能导入 baseline state。'); return JSON.parse(readFileSync(fixture.statePath, 'utf8')); } function createHistoryManifest({fixture, state, plan, archivePath}) { return { schemaVersion: 1, backupKind: 'spacetimedb-history', database: 'test-db', dataDir: fixture.dataDir, bucket: 'backup-bucket', objectKey: `database-backups/test-db/history/${state.baseline.id}/test-batch.tar.gz`, archivePath, baselineId: state.baseline.id, baselineStatePath: fixture.statePath, batchId: 'test-batch', candidates: plan.candidates, uploadStatus: 'pending', }; } async function readRequestBody(body) { if (body === undefined || body === null) { return Buffer.alloc(0); } if (typeof body === 'string') { return Buffer.from(body); } if (Buffer.isBuffer(body)) { return body; } const chunks = []; for await (const chunk of body) { chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); } return Buffer.concat(chunks); } function createFixture(name) { const root = path.join(tmpRoot, name); const binDir = path.join(root, 'bin'); const dataDir = path.join(root, 'data'); const workDir = path.join(root, 'work'); const systemctlLog = path.join(root, 'systemctl.log'); const tarLog = path.join(root, 'tar.log'); mkdirSync(binDir, {recursive: true}); mkdirSync(dataDir, {recursive: true}); writeFileSync(path.join(dataDir, 'sample.bin'), 'sample backup payload\n', 'utf8'); writeExecutable( path.join(binDir, 'systemctl'), `#!/usr/bin/env bash printf 'systemctl %s\\n' "$*" >> "${systemctlLog}" exit 0 `, ); writeExecutable( path.join(binDir, 'tar'), `#!/usr/bin/env bash printf 'tar %s\\n' "$*" >> "${tarLog}" echo 'fake tar failure' >&2 exit 2 `, ); return {root, binDir, dataDir, workDir, systemctlLog, tarLog}; } function runBackup(fixture, extraArgs = []) { return spawnSync( process.execPath, [ BACKUP_SCRIPT, '--data-dir', fixture.dataDir, '--work-dir', fixture.workDir, '--bucket', 'genarrative-test', '--endpoint', 'oss-cn-shanghai.aliyuncs.com', '--access-key-id', 'test', '--access-key-secret', 'test', ...extraArgs, ], { cwd: process.cwd(), encoding: 'utf8', env: { ...process.env, PATH: `${fixture.binDir}${path.delimiter}${process.env.PATH ?? ''}`, }, }, ); } function writeExecutable(filePath, content) { writeFileSync(filePath, content, 'utf8'); spawnSync('chmod', ['0755', filePath], {encoding: 'utf8'}); } function readFile(filePath) { return existsSync(filePath) ? readFileSync(filePath, 'utf8') : ''; } function assertStatus(result, expected, reason) { const actual = result.status ?? 0; if (actual !== expected) { failures.push( `${reason} 预期退出码 ${expected},实际 ${actual}。\nstdout:\n${result.stdout}\nstderr:\n${result.stderr}`, ); } } function assertIncludes(content, expected, reason) { if (!String(content).includes(expected)) { failures.push(`${reason} 缺少: ${expected}`); } } function assertEqual(actual, expected, reason) { if (actual !== expected) { failures.push(`${reason} 预期: ${String(expected)},实际: ${String(actual)}`); } } function assertNotEqual(actual, expected, reason) { if (actual === expected) { failures.push(`${reason} 两者均为: ${String(actual)}`); } } function assertTrue(condition, reason) { if (!condition) { failures.push(reason); } } function assertThrows(callback, expectedMessage, reason) { let thrown = null; try { callback(); } catch (error) { thrown = error; } if (!(thrown instanceof Error)) { failures.push(`${reason} 预期抛出错误。`); return; } assertIncludes(thrown.message, expectedMessage, reason); } function assertBufferEqual(actual, expected, reason) { if (!Buffer.isBuffer(actual) || !actual.equals(expected)) { failures.push(`${reason} 预期 ${expected.length} bytes,实际 ${actual?.length ?? ''} bytes。`); } } function assertFileMissing(filePath, reason) { if (existsSync(filePath)) { failures.push(`${reason} 实际存在: ${filePath}\n${readFile(filePath)}`); } }