From 947f59c894daee147468cbd90e7a424c2be4504d Mon Sep 17 00:00:00 2001 From: kdletters Date: Thu, 16 Jul 2026 19:00:18 +0800 Subject: [PATCH 1/2] =?UTF-8?q?=E9=99=90=E5=88=B6=E6=95=B0=E6=8D=AE?= =?UTF-8?q?=E5=BA=93=E5=A4=87=E4=BB=BD=E4=B8=8A=E4=BC=A0=E5=B8=A6=E5=AE=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 新增共享上传带宽限制器并覆盖数据对象、分片、catalog 与 latest pointer 补齐并发聚合限速、源流错误传播和大型 manifest 分块测试 更新备份环境变量示例与生产运维文档 --- deploy/env/api-server.env.example | 2 + ...发运维】本地开发验证与生产运维-2026-05-15.md | 3 +- scripts/check-database-backup-to-oss.mjs | 69 +++++++++++- scripts/database-backup-to-oss.mjs | 103 ++++++++++++++++-- 4 files changed, 163 insertions(+), 14 deletions(-) diff --git a/deploy/env/api-server.env.example b/deploy/env/api-server.env.example index 11c056e69..b4a392fa7 100644 --- a/deploy/env/api-server.env.example +++ b/deploy/env/api-server.env.example @@ -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。 diff --git a/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md b/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md index 8bc9a31a6..dd76e5325 100644 --- a/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md +++ b/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md @@ -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 引用发布到固定 `//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=`,该共享限速器只包裹备份上传流,空值或 `0` 表示不限速,不修改主机全局 qdisc。并发、限速和单次 PUT 都不改变“全部对象、catalog 与 latest pointer 成功后才推进 state/清理”的顺序。full 基线必须来自停库后的 data-dir 或已通过恢复验证的冻结副本;源文件上传前后 stat 虽会复核,但在线扫描不能保证大量文件属于同一跨文件一致时点。catalog 验真后,脚本把最新 full/history 引用发布到固定 `//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 超时后必须先检查原进程,不能直接重跑。 diff --git a/scripts/check-database-backup-to-oss.mjs b/scripts/check-database-backup-to-oss.mjs index 31f685693..e92a1cce7 100644 --- a/scripts/check-database-backup-to-oss.mjs +++ b/scripts/check-database-backup-to-oss.mjs @@ -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,54 @@ 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 的整数', + '上传带宽限制必须拒绝过小值。', + ); + let nowMs = 0; + const delays = []; + const limiter = createUploadBandwidthLimiter(1024, { + nowFn: () => nowMs, + sleepImpl: async (delayMs) => { + delays.push(delayMs); + nowMs += 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(nowMs, 4000, '两个并发上传流必须共享同一个总带宽预算。'); + assertEqual(delays.reduce((sum, delayMs) => sum + delayMs, 0), 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 +788,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 +806,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 +834,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'], diff --git a/scripts/database-backup-to-oss.mjs b/scripts/database-backup-to-oss.mjs index 4c7a92aa8..c4a4b1ef6 100644 --- a/scripts/database-backup-to-oss.mjs +++ b/scripts/database-backup-to-oss.mjs @@ -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}); @@ -2602,7 +2659,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 +2674,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 +2744,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 +2839,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 +2882,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 +2940,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 +2980,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 +3021,7 @@ async function main() { accessKeySecret, objectPrefix, keepLocal, + bandwidthLimiter: uploadBandwidthLimiter, }); return; } @@ -3026,7 +3096,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 +3115,7 @@ async function main() { objectKey: uploadedManifest.manifestObjectKey, accessKeyId, accessKeySecret, + bandwidthLimiter: uploadBandwidthLimiter, }); uploadedManifest.manifestVerifiedAt = manifestUpload.verifiedAt; uploadedManifest.manifestContentLength = manifestUpload.contentLength; From 0253a35e72affa4f60012207a5d5959a0088522c Mon Sep 17 00:00:00 2001 From: kdletters Date: Thu, 16 Jul 2026 19:02:42 +0800 Subject: [PATCH 2/2] =?UTF-8?q?=E8=A1=A5=E9=BD=90=E5=BD=92=E6=A1=A3?= =?UTF-8?q?=E6=B8=85=E5=8D=95=E4=B8=8A=E4=BC=A0=E9=99=90=E9=80=9F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 让已有全量归档的 manifest 继承同一上传带宽限制器 收紧并发上传共享累计带宽预算测试 --- scripts/check-database-backup-to-oss.mjs | 7 ++----- scripts/database-backup-to-oss.mjs | 1 + 2 files changed, 3 insertions(+), 5 deletions(-) diff --git a/scripts/check-database-backup-to-oss.mjs b/scripts/check-database-backup-to-oss.mjs index e92a1cce7..809d54266 100644 --- a/scripts/check-database-backup-to-oss.mjs +++ b/scripts/check-database-backup-to-oss.mjs @@ -152,13 +152,11 @@ async function assertUploadBandwidthLimiterSharesBudgetAndPropagatesErrors() { '必须为空、0 或 >= 1024 的整数', '上传带宽限制必须拒绝过小值。', ); - let nowMs = 0; const delays = []; const limiter = createUploadBandwidthLimiter(1024, { - nowFn: () => nowMs, + nowFn: () => 0, sleepImpl: async (delayMs) => { delays.push(delayMs); - nowMs += delayMs; }, }); const consume = async (stream) => { @@ -173,8 +171,7 @@ async function assertUploadBandwidthLimiterSharesBudgetAndPropagatesErrors() { consume(limiter.wrap(Readable.from([Buffer.alloc(1024), Buffer.alloc(1024)], {objectMode: false}))), ]); assertEqual(firstBytes + secondBytes, 4096, '共享上传限速器不得丢失并发流内容。'); - assertEqual(nowMs, 4000, '两个并发上传流必须共享同一个总带宽预算。'); - assertEqual(delays.reduce((sum, delayMs) => sum + delayMs, 0), 4000, '共享上传限速器必须预约完整字节时长。'); + assertEqual(delays.join(','), '1000,2000,3000,4000', '两个并发上传流必须共享同一个累计带宽预算。'); let sourceError = null; try { diff --git a/scripts/database-backup-to-oss.mjs b/scripts/database-backup-to-oss.mjs index c4a4b1ef6..80fb5392c 100644 --- a/scripts/database-backup-to-oss.mjs +++ b/scripts/database-backup-to-oss.mjs @@ -2630,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;