Merge branch 'master' into editor-agent-refactored

This commit is contained in:
2026-07-16 19:35:54 +08:00
4 changed files with 161 additions and 14 deletions
+2
View File
@@ -161,6 +161,8 @@ GENARRATIVE_DATABASE_BACKUP_KEEP_LOCAL=false
GENARRATIVE_DATABASE_BACKUP_STORAGE_FORMAT=archive
# files 模式并行 PUT/HEAD 数量,必须为 1-64;默认 16。
GENARRATIVE_DATABASE_BACKUP_FILES_CONCURRENCY=16
# 可选:仅限制备份进程上传总带宽,单位 bytes/s;为空或 0 表示不限速。
GENARRATIVE_DATABASE_BACKUP_UPLOAD_MAX_BYTES_PER_SECOND=
# 可选:显式要求备份工作目录所在文件系统至少保留的可用空间;为空时按数据目录大小 + 安全余量估算。
GENARRATIVE_DATABASE_BACKUP_MIN_FREE_BYTES=
# archive history 模式持久化已验真 full baseline 与追加批次;files 模式改用 work-dir 下的 files state。
@@ -331,6 +331,7 @@ GENARRATIVE_DATABASE_BACKUP_OSS_PREFIX=database-backups
GENARRATIVE_DATABASE_BACKUP_KEEP_LOCAL=false
GENARRATIVE_DATABASE_BACKUP_STORAGE_FORMAT=archive
GENARRATIVE_DATABASE_BACKUP_FILES_CONCURRENCY=16
GENARRATIVE_DATABASE_BACKUP_UPLOAD_MAX_BYTES_PER_SECOND=
GENARRATIVE_DATABASE_BACKUP_MIN_FREE_BYTES=
GENARRATIVE_DATABASE_BACKUP_BASELINE_STATE=/var/lib/genarrative/database-backups/genarrative-prod-history-state.json
# 仅 archive history 首次从一份 uploadStatus=uploaded 的全量 manifest 初始化 state 时设置或传 --baseline-manifest。
@@ -343,7 +344,7 @@ GENARRATIVE_DATABASE_BACKUP_OSS_ACCESS_KEY_SECRET=
`files-history` 使用仓库模板 `deploy/systemd/genarrative-database-backup-files-history.conf` 覆盖主 service 的 `ExecStart`,从 `/etc/genarrative/api-server.env` 读取 data-dir、database、bucket、prefix 与 OSS 凭据,不在 unit 写死环境目标,也不传 `--stop-service`。Server-Provision 在改动 drop-in 前,先用 current release 的同一脚本、同一 env 和 `DATABASE_BACKUP_FILES_HISTORY_WORK_DIR` 执行一次 history `--dry-run`;缺少已发布 full catalog 的 files state、current 脚本过旧或配置不匹配都会在安装 drop-in 和 `daemon-reload` 前失败。选择 `archive-full` 会主动删除仓库托管的 `10-files-history.conf` 与 dev 试点遗留的 `10-dev-files.conf`,防止 systemd 继续合并旧覆盖。dev 可继续指定已有 `/var/lib/genarrative/database-backups/dev-files`release 建议先在 `/var/lib/genarrative/database-backups/release-files` 建立自己的 full baseline;两台机器不得复用或互传本地 state 目录冒充本机基线。启用时通过 Server-Provision Job 选择目标、`DATABASE_BACKUP_PROFILE=files-history` 和对应 work-dir,先保持 `DRY_RUN=true` 核对,再以同参数正式 provision。不要直接在 `/etc/systemd/system` 手写第二份 drop-in。
files full 会递归扫描 data-dir,保留空目录、每个普通文件的相对路径,以及目标仍位于 data-dir 内部的相对符号链接;绝对链接或解析后越界的链接直接拒绝。文件按 SHA-256 上传到不可变对象 key,catalog 记录目录、路径、长度、SHA、对象 key 和相对链接目标,不写 staging 主机的绝对路径。相同 catalog 重跑不重复 PUT;新增或变化文件先 HEAD CAS 对象,存在且长度/SHA 元数据一致就复用,否则上传。16 MiB 及以下对象使用单次 PUT 后 HEAD 验真,大对象继续使用 multipart;对象操作默认以 16 路并行执行,可用 `GENARRATIVE_DATABASE_BACKUP_FILES_CONCURRENCY=1..64` 调整。并发和单次 PUT 只缩短传输与验真时间,不改变“全部对象、catalog 与 latest pointer 成功后才推进 state/清理”的顺序。full 基线必须来自停库后的 data-dir 或已通过恢复验证的冻结副本;源文件上传前后 stat 虽会复核,但在线扫描不能保证大量文件属于同一跨文件一致时点。catalog 验真后,脚本把最新 full/history 引用发布到固定 `<prefix>/<database>/latest.json`,全新机器不需要本地 state 即可自动发现恢复入口。
files full 会递归扫描 data-dir,保留空目录、每个普通文件的相对路径,以及目标仍位于 data-dir 内部的相对符号链接;绝对链接或解析后越界的链接直接拒绝。文件按 SHA-256 上传到不可变对象 key,catalog 记录目录、路径、长度、SHA、对象 key 和相对链接目标,不写 staging 主机的绝对路径。相同 catalog 重跑不重复 PUT;新增或变化文件先 HEAD CAS 对象,存在且长度/SHA 元数据一致就复用,否则上传。16 MiB 及以下对象使用单次 PUT 后 HEAD 验真,大对象继续使用 multipart;对象操作默认以 16 路并行执行,可用 `GENARRATIVE_DATABASE_BACKUP_FILES_CONCURRENCY=1..64` 调整。需要给线上入口留带宽时设置 `GENARRATIVE_DATABASE_BACKUP_UPLOAD_MAX_BYTES_PER_SECOND=<bytes/s>`,该共享限速器只包裹备份上传流,空值或 `0` 表示不限速,不修改主机全局 qdisc。并发、限速和单次 PUT 都不改变“全部对象、catalog 与 latest pointer 成功后才推进 state/清理”的顺序。full 基线必须来自停库后的 data-dir 或已通过恢复验证的冻结副本;源文件上传前后 stat 虽会复核,但在线扫描不能保证大量文件属于同一跨文件一致时点。catalog 验真后,脚本把最新 full/history 引用发布到固定 `<prefix>/<database>/latest.json`,全新机器不需要本地 state 即可自动发现恢复入口。
history 的安全边界按每个 replica 独立计算。设最新完整且未锁定的 snapshot offset 为 `S`;数字更大但缺少同 offset `.snapshot_bsatn`、仍存在同名 `.lock` 的目录不能参与边界计算。脚本必须保留起始 offset 小于等于 `S` 的最后一个 commitlog segment,以及它之后的全部 segment;只处理更早的 `.stdb.log` / `.stdb.ofs`,snapshot 只处理最新目录之前的旧目录。files history 会递归展开候选目录,逐对象复用或上传,随后依次验真候选对象、history catalog 与 full baseline catalog,再重新扫描边界和 stat fingerprint,最后覆盖发布并验真 `latest.json`;任何一步失败都不推进 state 或删除源文件。脚本在 work-dir 使用 PID lock 拒绝同库并发上传,SSH 超时后必须先检查原进程,不能直接重跑。
+65 -1
View File
@@ -5,12 +5,14 @@ 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,
@@ -47,6 +49,7 @@ async function main() {
assertInsufficientSpaceStopsBeforeServiceChanges();
assertArchiveFailureStillRestoresDependentServices();
await assertMultipartUploadRetriesAndVerifiesRemoteLength();
await assertUploadBandwidthLimiterSharesBudgetAndPropagatesErrors();
await assertDirectSmallFileUsesSinglePut();
await assertMissingPartEtagAbortsMultipartUpload();
await assertCompleteResponseAmbiguityUsesHeadVerification();
@@ -109,9 +112,13 @@ async function assertDirectSmallFileUsesSinglePut() {
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') {
@@ -131,9 +138,51 @@ async function assertDirectSmallFileUsesSinglePut() {
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() {
@@ -736,11 +785,13 @@ async function assertHeadShaMismatchAbortsMultipartUpload() {
async function assertManifestUploadUsesShaAndHeadVerification() {
const root = path.join(tmpRoot, 'manifest-upload');
const manifestPath = path.join(root, 'backup.manifest.json');
const body = Buffer.from('{"uploadStatus":"uploaded"}\n');
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',
@@ -752,6 +803,17 @@ async function assertManifestUploadUsesShaAndHeadVerification() {
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});
@@ -769,6 +831,8 @@ async function assertManifestUploadUsesShaAndHeadVerification() {
});
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'],
+92 -12
View File
@@ -288,6 +288,47 @@ function parseDirectFilesConcurrency(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);
}
@@ -2142,6 +2183,7 @@ export async function uploadArchive({
archiveSha256 = '',
contentType = 'application/gzip',
allowEmpty = false,
bandwidthLimiter = null,
}) {
const fileStat = statSync(archivePath);
if (!fileStat.isFile() || (!allowEmpty && fileStat.size <= 0)) {
@@ -2219,7 +2261,10 @@ export async function uploadArchive({
queries: {partNumber, uploadId},
headers: {'content-type': 'application/octet-stream'},
contentLength,
bodyFactory: () => createReadStream(archivePath, {start, end}),
bodyFactory: () => {
const stream = createReadStream(archivePath, {start, end});
return bandwidthLimiter ? bandwidthLimiter.wrap(stream) : stream;
},
operation: `UploadPart ${partNumber}/${partCount}`,
});
const etag = response.headers.get('etag');
@@ -2298,6 +2343,7 @@ export async function uploadDirectFile({
archiveSha256 = '',
contentType = 'application/octet-stream',
allowEmpty = true,
bandwidthLimiter = null,
}) {
const fileStat = statSync(archivePath);
if (!fileStat.isFile() || (!allowEmpty && fileStat.size <= 0)) {
@@ -2322,6 +2368,7 @@ export async function uploadDirectFile({
archiveSha256,
contentType,
allowEmpty,
bandwidthLimiter,
});
}
const verifiedArchiveSha256 = archiveSha256 || await sha256FileHex(archivePath);
@@ -2353,7 +2400,13 @@ export async function uploadDirectFile({
'x-oss-meta-backup-kind': backupKind,
},
contentLength: fileStat.size,
bodyFactory: () => fileStat.size === 0 ? Buffer.alloc(0) : createReadStream(archivePath),
bodyFactory: () => {
if (fileStat.size === 0) {
return Buffer.alloc(0);
}
const stream = createReadStream(archivePath);
return bandwidthLimiter ? bandwidthLimiter.wrap(stream) : stream;
},
operation: '上传逐文件对象',
});
const verification = await verifyUploadedObject({
@@ -2386,6 +2439,7 @@ export async function uploadManifestFile({
sleepImpl = sleep,
randomFn = Math.random,
maxAttempts = DEFAULT_OSS_REQUEST_MAX_ATTEMPTS,
bandwidthLimiter = null,
}) {
const body = readFileSync(manifestPath);
if (body.length === 0) {
@@ -2416,7 +2470,10 @@ export async function uploadManifestFile({
'x-oss-meta-backup-kind': 'spacetimedb-backup-manifest',
},
contentLength: body.length,
bodyFactory: () => body,
bodyFactory: () => {
const stream = createBufferReadStream(body);
return bandwidthLimiter ? bandwidthLimiter.wrap(stream) : stream;
},
operation: '上传 manifest',
});
const verification = await verifyUploadedObject({
@@ -2517,7 +2574,7 @@ export async function uploadHistoryArchiveWithCleanup({
return {result, uploadedManifest, cleanup, state};
}
async function uploadExistingArchive({args, env, bucket, endpoint, accessKeyId, accessKeySecret, objectPrefix}) {
async function uploadExistingArchive({args, env, bucket, endpoint, accessKeyId, accessKeySecret, objectPrefix, bandwidthLimiter}) {
const archivePath = resolvePath(args.uploadArchive);
if (!existsSync(archivePath)) {
throw new Error(`待上传备份文件不存在: ${archivePath}`);
@@ -2556,13 +2613,13 @@ async function uploadExistingArchive({args, env, bucket, endpoint, accessKeyId,
manifestPath,
manifest,
statePath,
uploadOptions: {bucket, endpoint, objectKey, accessKeyId, accessKeySecret},
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});
result = await uploadArchive({archivePath, bucket, endpoint, objectKey, accessKeyId, accessKeySecret, bandwidthLimiter});
const uploadedManifest = uploadedManifestPayload({manifest, database, result});
uploadedAt = uploadedManifest.uploadedAt;
writeManifest({manifestPath, payload: uploadedManifest});
@@ -2573,6 +2630,7 @@ async function uploadExistingArchive({args, env, bucket, endpoint, accessKeyId,
objectKey: uploadedManifest.manifestObjectKey,
accessKeyId,
accessKeySecret,
bandwidthLimiter,
});
uploadedManifest.manifestVerifiedAt = manifestUpload.verifiedAt;
uploadedManifest.manifestContentLength = manifestUpload.contentLength;
@@ -2602,7 +2660,7 @@ async function uploadExistingArchive({args, env, bucket, endpoint, accessKeyId,
}
}
async function publishExistingManifest({args, bucket, endpoint, accessKeyId, accessKeySecret}) {
async function publishExistingManifest({args, bucket, endpoint, accessKeyId, accessKeySecret, bandwidthLimiter}) {
const manifestPath = resolvePath(args.publishManifest);
const manifest = readManifest(manifestPath);
if (manifest.uploadStatus !== 'uploaded' || !manifest.objectKey) {
@@ -2617,6 +2675,7 @@ async function publishExistingManifest({args, bucket, endpoint, accessKeyId, acc
objectKey: manifest.manifestObjectKey,
accessKeyId,
accessKeySecret,
bandwidthLimiter,
});
manifest.manifestVerifiedAt = result.verifiedAt;
manifest.manifestContentLength = result.contentLength;
@@ -2686,6 +2745,7 @@ async function runHistoryBackup({
accessKeySecret,
objectPrefix,
keepLocal,
bandwidthLimiter,
}) {
const statePath = historyStatePath({args, env, workDir, database});
let state = loadOrImportHistoryState({args, env, statePath, database, dataDir});
@@ -2780,7 +2840,7 @@ async function runHistoryBackup({
manifestPath,
manifest,
statePath,
uploadOptions: {bucket, endpoint, objectKey, accessKeyId, accessKeySecret},
uploadOptions: {bucket, endpoint, objectKey, accessKeyId, accessKeySecret, bandwidthLimiter},
});
console.log(`[database-backup] history 上传并清理完成: ${JSON.stringify(historyResult.cleanup)}`);
if (args.resultFile) {
@@ -2823,6 +2883,7 @@ async function main() {
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}`);
@@ -2880,12 +2941,21 @@ async function main() {
}
if (args.publishManifest) {
await publishExistingManifest({args, bucket, endpoint, accessKeyId, accessKeySecret});
await publishExistingManifest({args, bucket, endpoint, accessKeyId, accessKeySecret, bandwidthLimiter: uploadBandwidthLimiter});
return;
}
if (args.uploadArchive) {
await uploadExistingArchive({args, env, bucket, endpoint, accessKeyId, accessKeySecret, objectPrefix});
await uploadExistingArchive({
args,
env,
bucket,
endpoint,
accessKeyId,
accessKeySecret,
objectPrefix,
bandwidthLimiter: uploadBandwidthLimiter,
});
return;
}
@@ -2911,7 +2981,7 @@ async function main() {
objectPrefix,
dryRun: args.dryRun,
resultFile: args.resultFile,
uploadOptions: {bucket, endpoint, accessKeyId, accessKeySecret},
uploadOptions: {bucket, endpoint, accessKeyId, accessKeySecret, bandwidthLimiter: uploadBandwidthLimiter},
concurrency: directFilesConcurrency,
});
} catch (error) {
@@ -2952,6 +3022,7 @@ async function main() {
accessKeySecret,
objectPrefix,
keepLocal,
bandwidthLimiter: uploadBandwidthLimiter,
});
return;
}
@@ -3026,7 +3097,15 @@ async function main() {
return;
}
const result = await uploadArchive({archivePath, bucket, endpoint, objectKey, accessKeyId, accessKeySecret});
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});
@@ -3037,6 +3116,7 @@ async function main() {
objectKey: uploadedManifest.manifestObjectKey,
accessKeyId,
accessKeySecret,
bandwidthLimiter: uploadBandwidthLimiter,
});
uploadedManifest.manifestVerifiedAt = manifestUpload.verifiedAt;
uploadedManifest.manifestContentLength = manifestUpload.contentLength;