收口生产内存工作集
Project CI / Repository checks (pull_request) Failing after 8s
Project CI / Native shell tests (pull_request) Failing after 10m34s
Project CI / Backend tests (pull_request) Failing after 9s
Project CI / Frontend tests (pull_request) Successful in 3m38s

限制外部生成 worker 脱管任务的执行容量

限制 AI 任务文本与终态内存

清理过期认证会话并限制临时状态

降低备份扫描与 catalog 哈希峰值

禁止 release 使用高风险 files-history 并增加备份资源护栏

同步后端、运维与项目决策文档
This commit is contained in:
2026-08-27 17:12:34 +08:00
parent 5942ff25e7
commit 0ef3461df9
12 changed files with 550 additions and 70 deletions
@@ -11,6 +11,12 @@ WorkingDirectory=/opt/genarrative/current
EnvironmentFile=/etc/genarrative/api-server.env
ExecStart=/usr/bin/node -- /opt/genarrative/current/scripts/database-backup-to-oss.mjs --env-file /etc/genarrative/api-server.env --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
# 备份脚本必须受独立内存上限保护,不能因目录扫描异常拖垮整台 release 主机。
Environment=NODE_OPTIONS=--max-old-space-size=768
MemoryHigh=768M
MemoryMax=1G
OOMPolicy=stop
# 备份需要停止 / 启动 spacetimedb.service,并读取 /stdb、写入 /var/lib/genarrative/database-backups。
# 停止 SpacetimeDB 会连带停止 Requires 它的 API / worker / controller,冷备份后必须显式拉起。
PrivateTmp=true
@@ -1309,8 +1309,8 @@ CI 上 `background_agent_runtime_recovers_stale_running_before_pending_task` 在
- OSS 固定恢复入口为 `<prefix>/<database>/latest.json`。CAS 文件和 full/history catalog 保持不可变;latest pointer 只保存最新 full catalog 与已发布 history catalog 的 object key、长度和 SHA,不包含主机绝对路径或文件内容。每次 state 变化先验真全部引用 catalog,再覆盖上传并 HEAD 验真 latest pointer,成功后才落本地 statehistory 还必须在 pointer 成功后才允许删除源文件。全新机器可仅凭 bucket、database、prefix 与 OSS 凭据自动下载 pointer 和 full catalog。
- dev 带宽不足时,允许把已冻结的 dev 基线经 `10.2.0.10 -> 10.2.4.16` 内网 rsync 到 release 独立 staging,再用 release 出口上传 dev bucketstaging 不得指向 release `/stdb`,不得停止或修改 release 服务,传输凭据必须临时创建并在演练后移除。catalog 不记录 staging 绝对路径,files state 可回传 dev 继续 history。
- 恢复边界:恢复时默认从 OSS `latest.json` 自动定位 full catalog,创建目录并按相对路径下载每个对象、逐文件校验长度与 SHA;本地 state 只用于备份续跑,不再是异机恢复前置条件。远程 dev 已完成真实 OSS、清理、重启和异机隔离恢复演练;release timer 与 publish 前备份继续保持原行为。
- systemd 接线:主 service 保持 `archive-full`。Server-Provision 新增默认值为 `archive-full``DATABASE_BACKUP_PROFILE`dev 或 release 显式选择 `files-history`,必须为各自主机指定独立 work-dir并先用 current release 脚本执行 history dry-run,确认已有 full state 后才安装仓库托管 drop-in,并删除现场手写旧 drop-in。切回默认 profile 必须删除所有 history 覆盖
- 影响范围:`scripts/database-backup-to-oss.mjs`、备份门禁、生产 env 示例、systemd 模板、Server-Provision、SpacetimeDB 运维与恢复流程;release timer 可在独立 baseline 验证后显式选择 profile,publish 前备份是否切换仍需单独决策。
- systemd 接线:主 service 保持 `archive-full`。Server-Provision 新增默认值为 `archive-full``DATABASE_BACKUP_PROFILE`development 可显式选择 `files-history`,必须指定独立 work-dir 并先用 current release 脚本执行 history dry-run,确认已有 full state 后才安装仓库托管 drop-inrelease 拒绝 `files-history`,直到流式 catalog 改造完成,以免大目录扫描再次触发 Node 内存峰值。切回默认 profile 必须删除所有 history 覆盖;备份 unit 同时设置 Node heap 与 systemd memory 上限,避免备份异常拖垮业务主机
- 影响范围:`scripts/database-backup-to-oss.mjs`、备份门禁、生产 env 示例、systemd 模板、Server-Provision、SpacetimeDB 运维与恢复流程;release timer 固定使用 archive-full,publish 前备份是否切换仍需单独决策。
- 验证方式:`npm run check:database-backup``npm run check:production-ops``npm run check:encoding``git diff --check`;dev 现场必须完成逐文件 full catalog、重复 full 零 PUT、history dry-run、上传后清理、STDB 重启和按 catalog 隔离恢复 roundtrip。
- 关联:<https://github.com/clockworklabs/SpacetimeDB/issues/5542#issuecomment-4981566448>。
@@ -329,6 +329,7 @@ Responses 的终态载荷既是工具调用的恢复源,也是正文的恢复
- Rust 结构体:`AiTask`
- 源码:`server-rs/crates/spacetime-module/src/ai/tasks.rs`
- `module-ai` 的进程内热状态不是持久化真相:文本增量按阶段有序聚合并受单阶段 512 KiB 上限约束;terminal task 立即释放增量明细,内存工作集最多保留 1024 个任务。需要长期查询时必须读取 SpacetimeDB 的 `ai_task` / `ai_task_stage` 投影,不得依赖进程重启后仍存在的内存快照。
### `ai_task_event`
@@ -413,7 +414,7 @@ Responses 的终态载荷既是工具调用的恢复源,也是正文的恢复
- Rust 结构体:`AuthStoreProjectionMeta`
- 源码:`server-rs/crates/spacetime-module/src/auth/tables.rs`
认证恢复策略:`api-server` 启动时只从 SpacetimeDB 正式认证表(`user_account` / `auth_identity` / `refresh_session`)导出 typed `AuthStoreProjectionView`,再恢复 `module-auth` 的进程内认证工作集;运行中 Bearer `sid` 或 refresh cookie 在本进程工作集内未命中时直接按失效处理,不再从 SpacetimeDB 导出整包认证状态刷新内存,避免旧投影把重复手机号或旧会话重新灌回进程。`module-auth` 只保留内存工作集和 projection 导入 / 导出能力,不再保留 JSON 快照导入 / 导出能力,也不写本地持久化文件;`auth-store.json` / `GENARRATIVE_AUTH_STORE_PATH` 不再是兼容恢复源。认证创建、登录会话、刷新、退出、改密、重置密码、绑定和资料变更等写操作必须在返回客户端前通过 `sync_auth_store_projection` 成功同步 SpacetimeDB 正式认证表;同步失败时接口返回错误,不允许把只存在于当前进程内存的账号或会话当成成功结果。新用户注册奖励、邀请码绑定和登录埋点必须排在认证同步成功之后,避免认证没落库时先写出钱包或邀请关系。若启动恢复阶段 SpacetimeDB 不可连接或超时,`api-server` 会按固定间隔持续重试认证工作集恢复,恢复成功后才开始监听 HTTP,避免一次短超时让进程永久停留在依赖不可用状态。
认证恢复策略:`api-server` 启动时只从 SpacetimeDB 正式认证表(`user_account` / `auth_identity` / `refresh_session`)导出 typed `AuthStoreProjectionView`,再恢复 `module-auth` 的进程内认证工作集;运行中 Bearer `sid` 或 refresh cookie 在本进程工作集内未命中时直接按失效处理,不再从 SpacetimeDB 导出整包认证状态刷新内存,避免旧投影把重复手机号或旧会话重新灌回进程。`module-auth` 只保留内存工作集和 projection 导入 / 导出能力,不再保留 JSON 快照导入 / 导出能力,也不写本地持久化文件;`auth-store.json` / `GENARRATIVE_AUTH_STORE_PATH` 不再是兼容恢复源。认证创建、登录会话、刷新、退出、改密、重置密码、绑定和资料变更等写操作必须在返回客户端前通过 `sync_auth_store_projection` 成功同步 SpacetimeDB 正式认证表;同步失败时接口返回错误,不允许把只存在于当前进程内存的账号或会话当成成功结果。新用户注册奖励、邀请码绑定和登录埋点必须排在认证同步成功之后,避免认证没落库时先写出钱包或邀请关系。refresh session 工作集最多保留 8192 条,短信验证码和微信登录 state 分别最多保留 4096 条;写入前清理过期项,达到上限时拒绝新增而不继续膨胀。若启动恢复阶段 SpacetimeDB 不可连接或超时,`api-server` 会按固定间隔持续重试认证工作集恢复,恢复成功后才开始监听 HTTP,避免一次短超时让进程永久停留在依赖不可用状态。
`auth_store_snapshot` 表和旧 `import_auth_store_snapshot_json` / `export_auth_store_snapshot_from_tables` procedure 已删除。认证投影同步只读写 `user_account``auth_identity``refresh_session``auth_store_projection_meta``auth_identity` 不再写 `phone_e164``display_name``avatar_url`,这些账号资料只以 `user_account` 为准。
@@ -1232,6 +1233,7 @@ RPG 创作入口的配置 ID 是 `rpg`,当前 `visible=true`、`open=true`
- Rust 结构体:`RefreshSession`
- 源码:`server-rs/crates/spacetime-module/src/auth/tables.rs`
- 认证工作集只保留 active 会话以及最近 24 小时内的 revoked / expired 会话;超过宽限期的失效会话在 refresh session 写路径和 projection 导出前从内存索引移除,并随下一次 typed projection 同步从正式表清理。该清理不改变 active 多端登录、单端登出或全端登出语义。
### `runtime_setting`
@@ -99,7 +99,7 @@ HTTP 角色的 `GENARRATIVE_SPACETIME_POOL_SIZE` 只表示 procedure / reducer
生产拆分角色时,`external-generation-worker``external-generation-controller` 的专属 env 示例会把 `GENARRATIVE_SPACETIME_POOL_SIZE` 覆盖为 `1`;非 HTTP 角色不创建 API 缓存读连接,只保留 `external_generation_job` 队列窄订阅作为响应式唤醒信号,实际抢占和扩缩容判断仍走 SpacetimeDB procedure。worker / controller 不执行模型定价 seed,启动时先调用受 runtime writer 鉴权的 queue-stats procedure 做只读预检,身份不匹配时 fail-fast;当前正式 systemd unit 通过共同加载 API env 继承同一 `GENARRATIVE_SPACETIME_TOKEN`,默认路径为 `/etc/genarrative/api-server.env`,自定义部署由 provision 和 API deploy 按实际参数渲染,专属角色 env 示例不重复配置该 token。`GENARRATIVE_EXTERNAL_GENERATION_WORKER_POLL_INTERVAL_MS` 与 controller poll interval 只作为订阅失效、漏事件和 lease 过期这类时间条件的兜底,不作为正常领取任务的主路径。
生产 worker 默认 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LEASE_SECONDS=600`,只覆盖 worker 心跳抖动和短暂断连窗口,不再把 lease 当成完整任务时长;默认 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_JOB_TIMEOUT_SECONDS=900``editor_image_generation``editor_image_edit``editor_icon_spritesheet_generation``editor_ui_design_asset_extraction` 四类 VectorEngine 图片任务与角色动画 / 视频类长任务使用 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LONG_JOB_TIMEOUT_SECONDS=1800`,手动去背景、音效和背景音乐继续使用普通预算。worker 在单次尝试超过执行预算后会停止续租并释放 worker 槽位,但不会取消已启动的业务 future 或主动写入失败 / 重试状态;在途执行由 lease fencing 仲裁,有效租约内写回仍可完成,租约过期后任务才可重新领取,attempt 耗尽时由认领事务标记失败并结算退款。生产部署和 provision 脚本会给 `/etc/genarrative/api-server.env``/etc/genarrative/external-generation-worker.env` 补齐这些变量;已有自定义值不覆盖,只会把历史旧默认 `3600` 迁移为 `600`
生产 worker 默认 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LEASE_SECONDS=600`,只覆盖 worker 心跳抖动和短暂断连窗口,不再把 lease 当成完整任务时长;默认 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_JOB_TIMEOUT_SECONDS=900``editor_image_generation``editor_image_edit``editor_icon_spritesheet_generation``editor_ui_design_asset_extraction` 四类 VectorEngine 图片任务与角色动画 / 视频类长任务使用 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LONG_JOB_TIMEOUT_SECONDS=1800`,手动去背景、音效和背景音乐继续使用普通预算。worker 在单次尝试超过执行预算后会停止续租,但不会取消已启动的业务 future 或主动写入失败 / 重试状态;执行许可会一直绑定到 active 或 detached work 真正结束(或超过租约仲裁窗口被取消),避免超时任务脱管后立即补进新的高内存任务。在途执行由 lease fencing 仲裁,有效租约内写回仍可完成,租约过期后任务才可重新领取,attempt 耗尽时由认领事务标记失败并结算退款。生产部署和 provision 脚本会给 `/etc/genarrative/api-server.env``/etc/genarrative/external-generation-worker.env` 补齐这些变量;已有自定义值不覆盖,只会把历史旧默认 `3600` 迁移为 `600`
lease 过期后不代表任务一定再次执行:claim transaction 只有在 `attempt < max_attempts` 时才会递增 attempt 并返回 worker;如果过期的是最终 attempt,则直接把 job 收口为 `failed`、清理 lease,并按入队冻结价格为当前 attempt 原子退款或写 cancellation intent。该终态任务不会再次进入 provider executor,迟到 consume 会被 settlement intent 拒绝。
@@ -401,7 +401,7 @@ UI 相关修改要重点验证:
### SpacetimeDB 数据目录 OSS 备份
数据库备份不放进 `spacetime-module` reducer / procedure:备份属于文件系统与 OSS 外部副作用,必须由运维脚本在 SpacetimeDB 宿主外执行。当前统一脚本为 `scripts/database-backup-to-oss.mjs`npm 命令 `npm run database:backup:oss`)。默认 `--storage-format archive --mode full` 保持原有全量压缩包冷备行为;`--storage-format files` 不生成 tar.gz,而是把目录树映射成逐文件 CAS 对象与 catalog,full 重跑只上传新增或内容变化的文件,history 只处理已被最新 snapshot 完全覆盖的历史 commitlog 与旧 snapshot。`Genarrative-Server-Provision``DATABASE_BACKUP_PROFILE` 默认是 `archive-full`,继续安装每天 `03:20` 左右执行的全量冷备主 service;development 和 release 都可以显式选择 `files-history`指定 work-dir 必须已经有与本机 database/bucket 匹配且已发布的 full baseline state
数据库备份不放进 `spacetime-module` reducer / procedure:备份属于文件系统与 OSS 外部副作用,必须由运维脚本在 SpacetimeDB 宿主外执行。当前统一脚本为 `scripts/database-backup-to-oss.mjs`npm 命令 `npm run database:backup:oss`)。默认 `--storage-format archive --mode full` 保持原有全量压缩包冷备行为;`--storage-format files` 不生成 tar.gz,而是把目录树映射成逐文件 CAS 对象与 catalog,full 重跑只上传新增或内容变化的文件,history 只处理已被最新 snapshot 完全覆盖的历史 commitlog 与旧 snapshot。`Genarrative-Server-Provision``DATABASE_BACKUP_PROFILE` 默认是 `archive-full`,继续安装每天 `03:20` 左右执行的全量冷备主 service;当前 release 只允许 `archive-full`,避免 `files-history` 在大目录上构造全量 catalog 导致 Node 内存峰值;development 才可以显式选择 `files-history`指定 work-dir 必须已经有与本机 database/bucket 匹配且已发布的 full baseline state
```bash
npm run database:backup:oss -- --data-dir /stdb --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
@@ -437,7 +437,7 @@ GENARRATIVE_DATABASE_BACKUP_OSS_ACCESS_KEY_SECRET=
`GENARRATIVE_DATABASE_BACKUP_OSS_BUCKET` 为空时会回退 `ALIYUN_OSS_BUCKET`AccessKey 默认复用 `ALIYUN_OSS_ACCESS_KEY_ID` / `ALIYUN_OSS_ACCESS_KEY_SECRET`,也可用 `GENARRATIVE_DATABASE_BACKUP_OSS_ACCESS_KEY_ID` / `GENARRATIVE_DATABASE_BACKUP_OSS_ACCESS_KEY_SECRET` 为备份 bucket 单独配置最小权限账号。冷备脚本会在停止 SpacetimeDB 前检查 `GENARRATIVE_DATABASE_BACKUP_WORK_DIR` 所在文件系统剩余空间;未设置 `GENARRATIVE_DATABASE_BACKUP_MIN_FREE_BYTES` 时,按数据目录大小加安全余量估算,空间不足会在停库前失败,避免写满根分区。即使打包或上传前步骤失败,只要脚本已经停过 SpacetimeDB,也会先恢复 SpacetimeDB 并执行 `--restart-service-after` 指定的 API / worker / controller,再带着原始备份错误退出。`Genarrative-Server-Provision` 会创建 `/var/lib/genarrative/database-backups` 并归属 `genarrative:genarrative`,同时安装并启用 `genarrative-database-backup.timer`。手动检查定时器:`systemctl list-timers genarrative-database-backup.timer`;手动触发一次:`systemctl start genarrative-database-backup.service`。如果 timer 显示 `enabled``inactive/dead``NEXT` / `Trigger` 为空,先写入当前 stamp 避免 `Persistent=true` 在白天立刻补跑冷备份:`touch /var/lib/systemd/timers/stamp-genarrative-database-backup.timer && systemctl daemon-reload && systemctl start genarrative-database-backup.timer`,随后确认下一次触发时间约为次日 `03:20`
`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-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 在 development 改动 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 继续合并旧覆盖。`genarrative-database-backup.service` 还通过 `NODE_OPTIONS=--max-old-space-size=768``MemoryHigh=768M``MemoryMax=1G``OOMPolicy=stop` 给备份进程设置独立护栏;release 若现场残留 files-history drop-in,必须先按 archive-full 重新 provision 并确认 drop-in 已删除,再恢复定时器。dev 可继续指定已有 `/var/lib/genarrative/database-backups/dev-files`;两台机器不得复用或互传本地 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` 调整。需要给线上入口留带宽时设置 `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 即可自动发现恢复入口。
@@ -478,7 +478,7 @@ node -- scripts/database-backup-to-oss.mjs \
dev 出口过慢时,可以把冻结基线经内网 rsync 到 release 独立 staging,再由 release 上传 dev bucket。staging 必须位于 `/var/lib/genarrative/dev-database-backup-staging/` 一类隔离目录,命令显式传 staging `--data-dir`、独立 `--work-dir`、dev `--bucket`,且不得传 `--stop-service`;禁止指向或修改 release `/stdb`。中转 key 只为本次传输临时授权,结束后从 dev 私钥和 release `authorized_keys` 同时移除。上传完成后把整个 files work-dir/state 回传 devhistory 才能延续同一 baseline catalog。
完整恢复默认从 OSS 固定 `latest.json` 读取最新 full catalog:先创建 `directories`,再把每个 `files[].objectKey` 下载到 `<restore-root>/<files[].path>` 并逐项核对 `sizeBytes` / `sha256`history catalog 用于证明已清理历史仍有 OSS 对象,不需要把已被 full baseline 覆盖的旧文件叠回当前恢复目录。本地 state 仍可作为兼容入口,并同时支持旧 v1 JSON 与 v2 gzip,但不再是异机恢复的前置条件。随后用隔离 data-dir 启动同版本 standalone,验证 `/v1/ping`、日志中的 snapshot restore / commitlog replay / module launch、代表性 SQL 和 reducer。dev 已完成这轮 OSS-only 异机恢复与重启演练;release 使用独立 `/var/lib/genarrative/database-backups/release-files` full baseline 和 `files-history` profile,现场最终 `ExecStart`、timer 状态与最近备份结果仍须在变更时重新核对。
完整恢复默认从 OSS 固定 `latest.json` 读取最新 full catalog:先创建 `directories`,再把每个 `files[].objectKey` 下载到 `<restore-root>/<files[].path>` 并逐项核对 `sizeBytes` / `sha256`history catalog 用于证明已清理历史仍有 OSS 对象,不需要把已被 full baseline 覆盖的旧文件叠回当前恢复目录。本地 state 仍可作为兼容入口,并同时支持旧 v1 JSON 与 v2 gzip,但不再是异机恢复的前置条件。随后用隔离 data-dir 启动同版本 standalone,验证 `/v1/ping`、日志中的 snapshot restore / commitlog replay / module launch、代表性 SQL 和 reducer。dev 已完成这轮 OSS-only 异机恢复与重启演练;release 使用 archive-full 时,现场最终 `ExecStart`、timer 状态与最近备份结果仍须在变更时重新核对。
```bash
node -- scripts/database-backup-to-oss.mjs \
@@ -34,8 +34,8 @@ pipeline {
string(name: 'WEB_LINK', defaultValue: '/srv/genarrative/web', description: 'Nginx 静态站点目录或软链接')
string(name: 'API_ENV_FILE', defaultValue: '/etc/genarrative/api-server.env', description: 'api-server 环境文件')
string(name: 'API_PORT', defaultValue: '8082', description: 'api-server 本机监听端口')
choice(name: 'DATABASE_BACKUP_PROFILE', choices: ['archive-full', 'files-history'], description: '数据库定时备份 profile默认 archive-fullfiles-history 仅在指定 work-dir 已有完整 full baseline 后启用')
string(name: 'DATABASE_BACKUP_FILES_HISTORY_WORK_DIR', defaultValue: '/var/lib/genarrative/database-backups/files-history', description: 'files-history 的本地 state/catalog 目录;dev/release 必须使用各自已建立 full baseline 的独立目录')
choice(name: 'DATABASE_BACKUP_PROFILE', choices: ['archive-full', 'files-history'], description: '数据库定时备份 profilerelease 仅允许 archive-fullfiles-history 仅供 development 在指定 work-dir 已有完整 full baseline 后启用')
string(name: 'DATABASE_BACKUP_FILES_HISTORY_WORK_DIR', defaultValue: '/var/lib/genarrative/database-backups/files-history', description: 'development files-history 的本地 state/catalog 目录;必须使用已建立 full baseline 的独立目录')
choice(name: 'NGINX_CONFIG_MODE', choices: ['none', 'production-https', 'development-http'], description: 'Nginx 配置模式;开发服无域名时选 development-httprelease 正式入口选 production-https')
booleanParam(name: 'ENABLE_SERVICES', defaultValue: true, description: '启用并启动 spacetimedb 与 api-server systemd 服务')
booleanParam(name: 'ENABLE_OTELCOL', defaultValue: true, description: '安装并启用本机 OpenTelemetry Collectorapi-server 模板默认开启 OTLP,如需关闭请在 API_ENV_FILE 中将 GENARRATIVE_OTEL_ENABLED 改为 false')
@@ -113,6 +113,9 @@ pipeline {
if (!(databaseBackupProfile in ['archive-full', 'files-history'])) {
error("DATABASE_BACKUP_PROFILE 只能是 archive-full 或 files-history,当前值: ${params.DATABASE_BACKUP_PROFILE}")
}
if (params.DEPLOY_TARGET == 'release' && databaseBackupProfile == 'files-history') {
error('release 仅允许 archive-fullfiles-history 会把整棵历史目录加载到 Node 内存,需先完成流式 catalog 改造后才能重新启用。')
}
def databaseBackupFilesHistoryWorkDir = params.DATABASE_BACKUP_FILES_HISTORY_WORK_DIR?.trim()
if (!(databaseBackupFilesHistoryWorkDir ==~ /^\/var\/lib\/genarrative\/database-backups\/[A-Za-z0-9._\/-]+$/) || databaseBackupFilesHistoryWorkDir.contains('..')) {
error("DATABASE_BACKUP_FILES_HISTORY_WORK_DIR 必须是 /var/lib/genarrative/database-backups/ 下不含连续点号的绝对路径,当前值: ${params.DATABASE_BACKUP_FILES_HISTORY_WORK_DIR}")
+14 -2
View File
@@ -821,6 +821,18 @@ const checks = [
reason:
'生产冷备份 service 必须用 node -- 分隔脚本参数,避免 Node 22 抢占业务 --env-file。',
},
{
file: 'deploy/systemd/genarrative-database-backup.service',
includes: 'Environment=NODE_OPTIONS=--max-old-space-size=768',
reason:
'备份 Node 进程必须设置独立 heap 上限,避免目录扫描异常拖垮 release 主机。',
},
{
file: 'deploy/systemd/genarrative-database-backup.service',
includes: 'MemoryMax=1G',
reason:
'备份 service 必须设置 systemd 内存硬上限,避免异常进程消耗整机内存。',
},
{
file: 'deploy/systemd/genarrative-database-backup.service',
excludes: '--storage-format files',
@@ -901,10 +913,10 @@ const checks = [
},
{
file: 'jenkins/Jenkinsfile.production-server-provision',
excludes:
includes:
"params.DEPLOY_TARGET == 'release' && databaseBackupProfile == 'files-history'",
reason:
'release 必须能在显式选择 profile 且 baseline 预检通过后启用 files-history。',
'release 必须拒绝 files-history,避免逐文件 catalog 扫描再次触发生产内存峰值。',
},
{
file: 'scripts/database-backup-to-oss.mjs',
+55 -16
View File
@@ -555,7 +555,11 @@ function assertSafeRelativePath(dataDir, absolutePath) {
}
function statFingerprint(absolutePath, rootPath = absolutePath) {
const entries = [];
// 候选 snapshot 可能包含数十万条目录项;增量更新摘要,避免把每条
// fingerprint 字符串同时保存在 entries[] 后再 join,造成一次性内存峰值。
const fingerprintHash = createHash('sha256');
let isFirstEntry = true;
let entryCount = 0;
let totalSize = 0n;
const visit = (currentPath) => {
const stat = lstatSync(currentPath, {bigint: true});
@@ -567,7 +571,7 @@ function statFingerprint(absolutePath, rootPath = absolutePath) {
if (kind === 'other') {
throw new Error(`history 候选只允许普通文件或目录: ${currentPath}`);
}
entries.push([
const entry = [
entryPath,
kind,
stat.dev.toString(),
@@ -575,7 +579,13 @@ function statFingerprint(absolutePath, rootPath = absolutePath) {
stat.mode.toString(),
stat.size.toString(),
stat.mtimeNs.toString(),
].join('\0'));
].join('\0');
if (!isFirstEntry) {
fingerprintHash.update('\n');
}
fingerprintHash.update(entry);
isFirstEntry = false;
entryCount += 1;
if (stat.isFile()) {
totalSize += stat.size;
} else {
@@ -586,9 +596,9 @@ function statFingerprint(absolutePath, rootPath = absolutePath) {
};
visit(rootPath);
return {
fingerprint: sha256Hex(entries.join('\n')),
fingerprint: fingerprintHash.digest('hex'),
sizeBytes: totalSize.toString(),
entryCount: entries.length,
entryCount,
};
}
@@ -1285,14 +1295,17 @@ export async function collectDirectFileEntries({dataDir, candidates = null, obje
throw new Error(`files 扫描期间源文件发生变化: ${relativePath}`);
}
const basePrefix = normalizeObjectPrefix(objectPrefix, database);
files.set(relativePath, {
const file = {
path: relativePath,
sizeBytes: Number(after.size),
sha256,
mode: after.mode,
objectKey: `${basePrefix}/files/sha256/${sha256.slice(0, 2)}/${sha256}`,
sourceStat: after,
});
};
// 上传前后的 inode/stat 仍用于防止在线扫描漂移,但设为不可枚举,避免
// 把仅供本地校验的副本再次写入 catalog 或 result JSON。
Object.defineProperty(file, 'sourceStat', {value: after, enumerable: false});
files.set(relativePath, file);
};
for (const root of roots.sort((left, right) => left.relativePath.localeCompare(right.relativePath))) {
@@ -1309,14 +1322,40 @@ export async function collectDirectFileEntries({dataDir, candidates = null, obje
}
function directCatalogIdentity({mode, baselineCatalogId, rootName, directories, files, symlinks}) {
return sha256Hex(JSON.stringify({
mode,
baselineCatalogId: baselineCatalogId || '',
rootName,
directories,
files: files.map(({path, sizeBytes, sha256, mode, objectKey}) => ({path, sizeBytes, sha256, mode, objectKey})),
symlinks,
// 不把数十万条文件元数据先拼成一个巨型 JSON 字符串;分段写入 hash
// 保持与 JSON.stringify 同样的字段顺序和转义结果,同时把峰值降到单条记录。
const hash = createHash('sha256');
hash.update('{"mode":');
hash.update(JSON.stringify(mode));
hash.update(',"baselineCatalogId":');
hash.update(JSON.stringify(baselineCatalogId || ''));
hash.update(',"rootName":');
hash.update(JSON.stringify(rootName));
hash.update(',"directories":');
updateJsonArrayHash(hash, directories, (directory) => JSON.stringify(directory));
hash.update(',"files":');
updateJsonArrayHash(hash, files, (file) => JSON.stringify({
path: file.path,
sizeBytes: file.sizeBytes,
sha256: file.sha256,
mode: file.mode,
objectKey: file.objectKey,
}));
hash.update(',"symlinks":');
updateJsonArrayHash(hash, symlinks, (symlink) => JSON.stringify(symlink));
hash.update('}');
return hash.digest('hex');
}
function updateJsonArrayHash(hash, values, serialize) {
hash.update('[');
values.forEach((value, index) => {
if (index > 0) {
hash.update(',');
}
hash.update(serialize(value));
});
hash.update(']');
}
function readDirectFilesState(statePath, {database, bucket}) {
@@ -1621,7 +1660,7 @@ export async function runDirectFilesBackup({
baselineCatalogId,
rootName,
directories: collected.directories,
files: collected.files.map(({sourceStat: _sourceStat, ...file}) => file),
files: collected.files,
symlinks: collected.symlinks,
};
writeManifest({manifestPath: catalogPath, payload: catalog});
+4
View File
@@ -78,6 +78,10 @@ validate_database_backup_profile() {
exit 1
;;
esac
if [[ "${DEPLOY_TARGET}" == "release" && "${DATABASE_BACKUP_PROFILE}" == "files-history" ]]; then
echo "[server-provision] release 仅允许 archive-fullfiles-history 会把整棵历史目录加载到 Node 内存,需先完成流式 catalog 改造后才能重新启用。" >&2
exit 1
fi
if [[ ! "${DATABASE_BACKUP_FILES_HISTORY_WORK_DIR}" =~ ^/var/lib/genarrative/database-backups/[A-Za-z0-9._/-]+$ || "${DATABASE_BACKUP_FILES_HISTORY_WORK_DIR}" == *..* ]]; then
echo "[server-provision] DATABASE_BACKUP_FILES_HISTORY_WORK_DIR 必须是 /var/lib/genarrative/database-backups/ 下不含连续点号的绝对路径,当前值: ${DATABASE_BACKUP_FILES_HISTORY_WORK_DIR}" >&2
exit 1
@@ -14,6 +14,7 @@ use spacetime_client::{
ExternalGenerationJobRenewLeaseRecordInput, ExternalGenerationQueueWakeSubscription,
};
use tokio::{
sync::{OwnedSemaphorePermit, Semaphore},
task::{JoinHandle, JoinSet},
time::sleep,
};
@@ -92,6 +93,10 @@ pub(crate) async fn run_external_generation_worker(state: AppState) -> Result<()
let concurrency = state.config.external_generation_worker_concurrency.max(1);
let poll_interval = state.config.external_generation_worker_poll_interval;
let lease = state.config.external_generation_worker_lease;
// 超时任务不能立即取消(在途 procedure 仍可能写回),因此执行容量必须同时
// 约束 active 与 detached work;否则每次超时都会释放 tasks 槽位,实际内存占用
// 会超过配置并发。
let work_slots = std::sync::Arc::new(Semaphore::new(concurrency));
let mut tasks = JoinSet::new();
let mut shutdown = external_generation_worker_shutdown_signal();
let mut queue_wake = None;
@@ -115,14 +120,24 @@ pub(crate) async fn run_external_generation_worker(state: AppState) -> Result<()
loop {
ensure_external_generation_queue_wake_subscription(&state, &mut queue_wake).await;
while tasks.len() >= concurrency {
if await_worker_task_or_shutdown(&mut tasks, &mut shutdown).await {
drain_external_generation_worker_tasks(&mut tasks).await;
return Ok(());
while work_slots.available_permits() == 0 {
tokio::select! {
_ = shutdown.as_mut() => {
drain_external_generation_worker_tasks(&mut tasks).await;
return Ok(());
}
permit = work_slots.clone().acquire_owned() => {
if permit.is_err() {
drain_external_generation_worker_tasks(&mut tasks).await;
return Ok(());
}
// 只用 acquire 作为容量变化唤醒信号,许可立即归还;真正领取任务
// 时在下方按返回的 job 数量逐个 try_acquire。
}
}
}
let available = concurrency.saturating_sub(tasks.len()).max(1);
let available = work_slots.available_permits().max(1);
let now_micros = current_utc_micros();
let lease_expires_at_micros = now_micros.saturating_add(duration_micros_i64(lease));
@@ -178,9 +193,13 @@ pub(crate) async fn run_external_generation_worker(state: AppState) -> Result<()
for job in jobs {
let state = state.clone();
let worker_id = worker_id.clone();
let permit = work_slots
.clone()
.try_acquire_owned()
.expect("claimed job must have an execution capacity permit");
tasks.spawn(async move {
if let Err(error) =
process_external_generation_job(state, worker_id, lease, job).await
process_external_generation_job(state, worker_id, lease, job, permit).await
{
error!(error = %error, "external generation worker 执行任务失败");
}
@@ -255,16 +274,6 @@ async fn await_worker_task(tasks: &mut JoinSet<()>) {
}
}
async fn await_worker_task_or_shutdown(
tasks: &mut JoinSet<()>,
shutdown: &mut ExternalGenerationShutdownSignal,
) -> bool {
tokio::select! {
_ = shutdown.as_mut() => true,
_ = await_worker_task(tasks) => false,
}
}
async fn await_one_task_or_queue_wake_or_sleep_or_shutdown(
tasks: &mut JoinSet<()>,
queue_wake: &mut Option<ExternalGenerationQueueWakeSubscription>,
@@ -326,6 +335,7 @@ async fn process_external_generation_job(
worker_id: String,
lease: Duration,
job: ExternalGenerationJobRecord,
permit: OwnedSemaphorePermit,
) -> Result<(), String> {
let heartbeat_interval = external_generation_worker_heartbeat_interval(lease);
let job_timeout = external_generation_worker_job_timeout(&state.config, job.job_kind.as_str());
@@ -384,6 +394,7 @@ async fn process_external_generation_job(
&job,
lease,
"任务超过执行预算",
Some(permit),
);
Err(message)
}
@@ -393,6 +404,7 @@ async fn process_external_generation_job(
&job,
lease,
"任务租约续期失败",
Some(permit),
);
Err(error)
}
@@ -417,6 +429,7 @@ fn detach_external_generation_work_until_lease_expiry(
job: &ExternalGenerationJobRecord,
lease: Duration,
reason: &'static str,
permit: Option<OwnedSemaphorePermit>,
) {
let job_id = job.job_id.clone();
let job_kind = job.job_kind.clone();
@@ -454,6 +467,9 @@ fn detach_external_generation_work_until_lease_expiry(
);
}
}
// 保持执行许可直到 work 真正结束或被取消,避免超时任务脱管后继续
// 累积图片/音频响应占用。
drop(permit);
});
}
@@ -1972,6 +1988,7 @@ mod tests {
&job,
Duration::from_millis(200),
"任务超过执行预算",
None,
);
tokio::time::sleep(Duration::from_millis(100)).await;
@@ -1981,6 +1998,39 @@ mod tests {
);
}
#[tokio::test]
async fn worker_detached_work_keeps_execution_slot_until_finished() {
let slots = std::sync::Arc::new(tokio::sync::Semaphore::new(1));
let permit = slots
.clone()
.acquire_owned()
.await
.expect("the only execution slot should be available");
let work_handle = tokio::spawn(async {
tokio::time::sleep(Duration::from_millis(20)).await;
Ok(())
});
let job = external_generation_job_record_fixture(Some("lease-1"));
detach_external_generation_work_until_lease_expiry(
work_handle,
&job,
Duration::from_millis(200),
"任务超过执行预算",
Some(permit),
);
assert!(
slots.try_acquire().is_err(),
"脱管 work 完成前不得重新领取执行容量"
);
tokio::time::sleep(Duration::from_millis(100)).await;
assert!(
slots.try_acquire().is_ok(),
"脱管 work 完成后应归还执行容量"
);
}
#[tokio::test]
async fn worker_detached_work_is_aborted_after_lease_arbitration_window() {
let connection = std::sync::Arc::new(tokio::sync::Semaphore::new(1));
@@ -1999,6 +2049,7 @@ mod tests {
&job,
Duration::from_millis(10),
"任务超过执行预算",
None,
);
let reacquired =
@@ -1,5 +1,5 @@
use std::{
collections::HashMap,
collections::{BTreeMap, HashMap},
sync::{Arc, Mutex},
};
@@ -9,6 +9,10 @@ use crate::{
use super::ensure_task_is_not_terminal;
const MAX_RETAINED_TASKS: usize = 1024;
const MAX_TASK_TEXT_OUTPUT_BYTES: usize = 512 * 1024;
const MAX_RETAINED_TASK_TEXT_BYTES: usize = 64 * 1024 * 1024;
#[derive(Clone, Debug, Default)]
pub struct InMemoryAiTaskStore {
inner: Arc<Mutex<InMemoryAiTaskStoreState>>,
@@ -17,7 +21,9 @@ pub struct InMemoryAiTaskStore {
#[derive(Debug, Default)]
struct InMemoryAiTaskStoreState {
tasks: HashMap<String, AiTaskSnapshot>,
text_chunks: HashMap<String, Vec<AiTextChunkSnapshot>>,
// Keep only the ordered deltas needed to handle an out-of-order chunk.
// Completed tasks drop this map immediately; it is not a second durable log.
text_chunks: HashMap<String, HashMap<crate::AiTaskStageKind, BTreeMap<u32, String>>>,
}
impl InMemoryAiTaskStore {
@@ -34,7 +40,25 @@ impl InMemoryAiTaskStore {
return Err(AiTaskServiceError::TaskAlreadyExists);
}
state.text_chunks.insert(task.task_id.clone(), Vec::new());
if state.tasks.len() >= MAX_RETAINED_TASKS {
let oldest_terminal = state
.tasks
.values()
.filter(|value| value.status.is_terminal())
.min_by_key(|value| value.completed_at_micros.or(Some(value.updated_at_micros)))
.map(|value| value.task_id.clone());
let Some(task_id) = oldest_terminal else {
return Err(AiTaskServiceError::Store(
"AI 任务仓储已达到内存容量上限".to_string(),
));
};
state.tasks.remove(&task_id);
state.text_chunks.remove(&task_id);
}
state
.text_chunks
.insert(task.task_id.clone(), HashMap::new());
state.tasks.insert(task.task_id.clone(), task.clone());
Ok(task)
}
@@ -56,7 +80,11 @@ impl InMemoryAiTaskStore {
.get_mut(task_id.trim())
.ok_or(AiTaskServiceError::TaskNotFound)?;
apply(task)?;
Ok(task.clone())
let snapshot = task.clone();
if snapshot.status.is_terminal() {
state.text_chunks.remove(task_id.trim());
}
Ok(snapshot)
}
pub(super) fn append_text_chunk(
@@ -67,41 +95,62 @@ impl InMemoryAiTaskStore {
.inner
.lock()
.map_err(|_| AiTaskServiceError::Store("AI 任务仓储锁已中毒".to_string()))?;
{
if chunk.delta_text.len() > MAX_TASK_TEXT_OUTPUT_BYTES {
return Err(AiTaskServiceError::Store(
"AI 任务文本输出超过内存上限".to_string(),
));
}
let (previous_stage_output_bytes, previous_latest_output_bytes) = {
let task = state
.tasks
.get_mut(&chunk.task_id)
.get(&chunk.task_id)
.ok_or(AiTaskServiceError::TaskNotFound)?;
ensure_task_is_not_terminal(task.status)?;
let stage = task
.stages
.iter_mut()
.iter()
.find(|stage| stage.stage_kind == chunk.stage_kind)
.ok_or(AiTaskServiceError::StageNotFound)?;
if stage.status == AiTaskStageStatus::Pending {
stage.status = AiTaskStageStatus::Running;
stage.started_at_micros = Some(chunk.created_at_micros);
}
(
stage.text_output.as_ref().map_or(0, String::len),
task.latest_text_output.as_ref().map_or(0, String::len),
)
};
task.status = AiTaskStatus::Running;
task.started_at_micros
.get_or_insert(chunk.created_at_micros);
let (previous_chunk, aggregated_bytes, aggregated_text) = {
let chunks = state
.text_chunks
.get_mut(&chunk.task_id)
.ok_or(AiTaskServiceError::TaskNotFound)?;
let stage_chunks = chunks.entry(chunk.stage_kind).or_default();
let previous_chunk = stage_chunks.insert(chunk.sequence, chunk.delta_text.clone());
let aggregated_bytes = stage_chunks
.values()
.fold(0_usize, |total, delta| total.saturating_add(delta.len()));
let mut aggregated_text = String::with_capacity(aggregated_bytes);
for delta in stage_chunks.values() {
aggregated_text.push_str(delta);
}
(previous_chunk, aggregated_bytes, aggregated_text)
};
if aggregated_bytes > MAX_TASK_TEXT_OUTPUT_BYTES {
rollback_text_chunk(&mut state, &chunk, previous_chunk);
return Err(AiTaskServiceError::Store(
"AI 任务文本输出超过内存上限".to_string(),
));
}
let chunks = state
.text_chunks
.get_mut(&chunk.task_id)
.ok_or(AiTaskServiceError::TaskNotFound)?;
chunks.push(chunk.clone());
chunks.sort_by_key(|value| value.sequence);
let retained_text_bytes = retained_text_bytes(&state)
.saturating_sub(previous_stage_output_bytes)
.saturating_sub(previous_latest_output_bytes)
.saturating_add(aggregated_bytes.saturating_mul(2));
if retained_text_bytes > MAX_RETAINED_TASK_TEXT_BYTES {
rollback_text_chunk(&mut state, &chunk, previous_chunk);
return Err(AiTaskServiceError::Store(
"AI 任务仓储文本工作集超过内存上限".to_string(),
));
}
let aggregated_text = chunks
.iter()
.filter(|value| value.stage_kind == chunk.stage_kind)
.map(|value| value.delta_text.as_str())
.collect::<Vec<_>>()
.join("");
let normalized_output = if aggregated_text.trim().is_empty() {
None
} else {
@@ -112,11 +161,19 @@ impl InMemoryAiTaskStore {
.tasks
.get_mut(&chunk.task_id)
.ok_or(AiTaskServiceError::TaskNotFound)?;
ensure_task_is_not_terminal(task.status)?;
let stage = task
.stages
.iter_mut()
.find(|stage| stage.stage_kind == chunk.stage_kind)
.ok_or(AiTaskServiceError::StageNotFound)?;
if stage.status == AiTaskStageStatus::Pending {
stage.status = AiTaskStageStatus::Running;
stage.started_at_micros = Some(chunk.created_at_micros);
}
task.status = AiTaskStatus::Running;
task.started_at_micros
.get_or_insert(chunk.created_at_micros);
stage.text_output = normalized_output.clone();
task.latest_text_output = normalized_output;
task.updated_at_micros = chunk.created_at_micros;
@@ -136,3 +193,41 @@ impl InMemoryAiTaskStore {
.ok_or(AiTaskServiceError::TaskNotFound)
}
}
fn rollback_text_chunk(
state: &mut InMemoryAiTaskStoreState,
chunk: &AiTextChunkSnapshot,
previous_chunk: Option<String>,
) {
if let Some(stage_chunks) = state
.text_chunks
.get_mut(&chunk.task_id)
.and_then(|chunks| chunks.get_mut(&chunk.stage_kind))
{
if let Some(previous_chunk) = previous_chunk {
stage_chunks.insert(chunk.sequence, previous_chunk);
} else {
stage_chunks.remove(&chunk.sequence);
}
}
}
fn retained_text_bytes(state: &InMemoryAiTaskStoreState) -> usize {
let snapshot_bytes = state.tasks.values().fold(0_usize, |total, task| {
let latest = task.latest_text_output.as_ref().map_or(0, String::len);
let stages = task.stages.iter().fold(0_usize, |stage_total, stage| {
stage_total.saturating_add(stage.text_output.as_ref().map_or(0, String::len))
});
total.saturating_add(latest).saturating_add(stages)
});
state
.text_chunks
.values()
.fold(snapshot_bytes, |total, stages| {
stages.values().fold(total, |stage_total, chunks| {
chunks.values().fold(stage_total, |chunk_total, delta| {
chunk_total.saturating_add(delta.len())
})
})
})
}
+41
View File
@@ -112,6 +112,47 @@ fn append_text_chunk_aggregates_stream_output_by_stage() {
assert_eq!(second_chunk.sequence, 2);
}
#[test]
fn append_text_chunk_rejects_output_over_stage_memory_limit_without_mutating_task() {
let service = build_service();
let task = service
.create_task(build_create_input(AiTaskKind::CharacterChat))
.expect("task should create");
let max_output = "a".repeat(512 * 1024);
let (updated, _) = service
.append_text_chunk(
&task.task_id,
AiTaskStageKind::RequestModel,
1,
max_output.clone(),
task.created_at_micros + 1,
)
.expect("the stage limit itself should be accepted");
assert_eq!(
updated.latest_text_output.as_deref().map(str::len),
Some(max_output.len())
);
let error = service
.append_text_chunk(
&task.task_id,
AiTaskStageKind::RequestModel,
2,
"b".to_string(),
task.created_at_micros + 2,
)
.expect_err("output beyond the stage limit should fail");
assert!(matches!(error, AiTaskServiceError::Store(_)));
let after_rejection = service
.get_task(&task.task_id)
.expect("task should remain readable");
assert_eq!(
after_rejection.latest_text_output.as_deref().map(str::len),
Some(max_output.len())
);
}
#[test]
fn complete_stage_updates_latest_outputs() {
let service = build_service();
+227
View File
@@ -28,6 +28,11 @@ use shared_kernel::{
use time::{Duration, OffsetDateTime};
use tracing::{info, warn};
const REFRESH_SESSION_STALE_RETENTION: Duration = Duration::days(1);
const MAX_REFRESH_SESSIONS: usize = 8_192;
const MAX_PHONE_CODES: usize = 4_096;
const MAX_WECHAT_STATES: usize = 4_096;
#[derive(Clone, Debug)]
pub struct InMemoryAuthStore {
inner: Arc<Mutex<InMemoryAuthStoreState>>,
@@ -364,6 +369,7 @@ impl RefreshSessionService {
input: CreateRefreshSessionInput,
now: OffsetDateTime,
) -> Result<CreateRefreshSessionResult, RefreshSessionError> {
self.store.prune_stale_sessions(now)?;
self.store
.find_by_user_id(&input.user_id)
.map_err(map_password_store_error)?
@@ -400,6 +406,7 @@ impl RefreshSessionService {
input: RotateRefreshSessionInput,
now: OffsetDateTime,
) -> Result<RotateRefreshSessionResult, RefreshSessionError> {
self.store.prune_stale_sessions(now)?;
let Some(refresh_token_hash) = normalize_required_string(&input.refresh_token_hash) else {
return Err(RefreshSessionError::MissingToken);
};
@@ -454,6 +461,7 @@ impl RefreshSessionService {
user_id: &str,
now: OffsetDateTime,
) -> Result<ListActiveRefreshSessionsResult, RefreshSessionError> {
self.store.prune_stale_sessions(now)?;
self.store
.find_by_user_id(user_id)
.map_err(map_password_store_error)?
@@ -468,6 +476,7 @@ impl RefreshSessionService {
input: RevokeRefreshSessionByUserInput,
now: OffsetDateTime,
) -> Result<RevokeRefreshSessionResult, RefreshSessionError> {
self.store.prune_stale_sessions(now)?;
self.store
.find_by_user_id(&input.user_id)
.map_err(map_password_store_error)?
@@ -492,6 +501,7 @@ impl RefreshSessionService {
session_id: &str,
now: OffsetDateTime,
) -> Result<bool, RefreshSessionError> {
self.store.prune_stale_sessions(now)?;
self.store
.is_session_active_for_user(user_id, session_id.trim(), now)
}
@@ -511,6 +521,7 @@ impl PhoneAuthService {
input: SendPhoneCodeInput,
now: OffsetDateTime,
) -> Result<SendPhoneCodeResult, PhoneAuthError> {
self.store.prune_expired_phone_codes(now)?;
let scene = input.scene.clone();
validate_mainland_china_country_code(input.country_code.as_deref())?;
let normalized_phone = normalize_mainland_china_phone_number(&input.pure_phone_number)?;
@@ -525,6 +536,8 @@ impl PhoneAuthService {
);
self.store
.ensure_phone_code_not_cooling_down(&normalized_phone.e164, &scene, now)?;
self.store
.ensure_phone_code_capacity(&normalized_phone.e164, &scene)?;
let expires_at = now
.checked_add(Duration::minutes(SMS_CODE_TTL_MINUTES))
.ok_or_else(|| PhoneAuthError::Store("短信验证码过期时间计算溢出".to_string()))?;
@@ -787,6 +800,7 @@ impl WechatAuthStateService {
input: CreateWechatAuthStateInput,
now: OffsetDateTime,
) -> Result<CreateWechatAuthStateResult, WechatAuthError> {
self.store.prune_wechat_states(now)?;
let created_at = format_rfc3339(now).map_err(|message| {
WechatAuthError::Store(format!("微信 state 时间格式化失败:{message}"))
})?;
@@ -1063,10 +1077,18 @@ impl InMemoryAuthStoreState {
}
}
let now = OffsetDateTime::now_utc();
for session in view.refresh_sessions {
if !existing_user_ids.contains(&session.user_id) {
continue;
}
if should_prune_refresh_session_fields(
&session.expires_at,
session.revoked_at.as_deref(),
now,
) {
continue;
}
let client_info =
serde_json::from_str::<RefreshSessionClientInfo>(&session.client_info_json)
.map_err(|error| format!("解析 refresh session 客户端信息失败:{error}"))?;
@@ -1185,6 +1207,8 @@ impl InMemoryAuthStore {
&self,
updated_at_micros: i64,
) -> Result<AuthStoreProjectionView, String> {
self.prune_stale_sessions(OffsetDateTime::now_utc())
.map_err(|error| error.to_string())?;
let state = self
.inner
.lock()
@@ -1262,6 +1286,38 @@ impl InMemoryAuthStore {
})
}
fn prune_stale_sessions(&self, now: OffsetDateTime) -> Result<(), RefreshSessionError> {
let mut state = self
.inner
.lock()
.map_err(|_| RefreshSessionError::Store("会话仓储锁已中毒".to_string()))?;
let stale_session_ids = state
.sessions_by_id
.iter()
.filter(|(_, stored)| should_prune_refresh_session(&stored.session, now))
.map(|(session_id, _)| session_id.clone())
.collect::<Vec<_>>();
if stale_session_ids.is_empty() {
return Ok(());
}
for session_id in stale_session_ids {
let Some(stored) = state.sessions_by_id.remove(&session_id) else {
continue;
};
if state
.session_id_by_refresh_token_hash
.get(&stored.session.refresh_token_hash)
.is_some_and(|mapped_id| mapped_id == &session_id)
{
state
.session_id_by_refresh_token_hash
.remove(&stored.session.refresh_token_hash);
}
}
self.persist_refresh_state(&state)
}
fn persist_state(&self, state: &InMemoryAuthStoreState) -> Result<(), String> {
let _ = state;
Ok(())
@@ -1894,6 +1950,11 @@ impl InMemoryAuthStore {
"refresh token hash 已存在,无法重复创建会话".to_string(),
));
}
if state.sessions_by_id.len() >= MAX_REFRESH_SESSIONS {
return Err(RefreshSessionError::Store(
"refresh session 内存容量已达到上限".to_string(),
));
}
state.session_id_by_refresh_token_hash.insert(
session.refresh_token_hash.clone(),
@@ -1918,10 +1979,41 @@ impl InMemoryAuthStore {
.map_err(|_| PhoneAuthError::Store("短信验证码仓储锁已中毒".to_string()))?;
// 手机号和业务场景共同决定同一份验证码快照,重复发送时直接覆盖旧值。
let key = build_phone_code_key(&code.phone_number, &code.scene);
if !state.phone_codes_by_key.contains_key(&key)
&& state.phone_codes_by_key.len() >= MAX_PHONE_CODES
{
return Err(PhoneAuthError::Store(
"短信验证码内存容量已达到上限,请稍后重试".to_string(),
));
}
state.phone_codes_by_key.insert(key, code);
Ok(())
}
fn prune_expired_phone_codes(&self, now: OffsetDateTime) -> Result<(), PhoneAuthError> {
let mut state = self
.inner
.lock()
.map_err(|_| PhoneAuthError::Store("短信验证码仓储锁已中毒".to_string()))?;
let expired_keys = state
.phone_codes_by_key
.iter()
.filter_map(|(key, stored)| {
OffsetDateTime::parse(
&stored.expires_at,
&time::format_description::well_known::Rfc3339,
)
.ok()
.filter(|expires_at| *expires_at <= now)
.map(|_| key.clone())
})
.collect::<Vec<_>>();
for key in expired_keys {
state.phone_codes_by_key.remove(&key);
}
Ok(())
}
fn ensure_phone_code_not_cooling_down(
&self,
phone_number: &str,
@@ -1961,6 +2053,26 @@ impl InMemoryAuthStore {
})
}
fn ensure_phone_code_capacity(
&self,
phone_number: &str,
scene: &PhoneAuthScene,
) -> Result<(), PhoneAuthError> {
let state = self
.inner
.lock()
.map_err(|_| PhoneAuthError::Store("短信验证码仓储锁已中毒".to_string()))?;
let key = build_phone_code_key(phone_number, scene);
if state.phone_codes_by_key.contains_key(&key)
|| state.phone_codes_by_key.len() < MAX_PHONE_CODES
{
return Ok(());
}
Err(PhoneAuthError::Store(
"短信验证码内存容量已达到上限,请稍后重试".to_string(),
))
}
fn get_active_phone_code(
&self,
phone_number: &str,
@@ -2041,6 +2153,11 @@ impl InMemoryAuthStore {
{
return Err(WechatAuthError::Store("微信 state 已存在".to_string()));
}
if state.wechat_states_by_token.len() >= MAX_WECHAT_STATES {
return Err(WechatAuthError::Store(
"微信登录 state 内存容量已达到上限,请稍后重试".to_string(),
));
}
state.wechat_states_by_token.insert(
state_record.state_token.clone(),
StoredWechatAuthState {
@@ -2405,6 +2522,33 @@ impl InMemoryAuthStore {
Ok(())
}
fn prune_wechat_states(&self, now: OffsetDateTime) -> Result<(), WechatAuthError> {
let mut state = self
.inner
.lock()
.map_err(|_| WechatAuthError::Store("微信 state 仓储锁已中毒".to_string()))?;
let stale_tokens = state
.wechat_states_by_token
.iter()
.filter_map(|(token, stored)| {
if stored.state.consumed_at.is_some() {
return Some(token.clone());
}
OffsetDateTime::parse(
&stored.state.expires_at,
&time::format_description::well_known::Rfc3339,
)
.ok()
.filter(|expires_at| *expires_at <= now)
.map(|_| token.clone())
})
.collect::<Vec<_>>();
for token in stale_tokens {
state.wechat_states_by_token.remove(&token);
}
Ok(())
}
fn revoke_session_by_user_and_session_id(
&self,
user_id: &str,
@@ -2568,6 +2712,25 @@ impl InMemoryAuthStore {
}
}
fn should_prune_refresh_session(session: &RefreshSessionRecord, now: OffsetDateTime) -> bool {
should_prune_refresh_session_fields(&session.expires_at, session.revoked_at.as_deref(), now)
}
fn should_prune_refresh_session_fields(
expires_at: &str,
revoked_at: Option<&str>,
now: OffsetDateTime,
) -> bool {
let stale_before = now.saturating_sub(REFRESH_SESSION_STALE_RETENTION);
if let Some(revoked_at) = revoked_at {
return OffsetDateTime::parse(revoked_at, &time::format_description::well_known::Rfc3339)
.is_ok_and(|timestamp| timestamp <= stale_before);
}
OffsetDateTime::parse(expires_at, &time::format_description::well_known::Rfc3339)
.is_ok_and(|timestamp| timestamp <= stale_before)
}
fn map_sms_provider_error_to_phone_error(error: SmsProviderError) -> PhoneAuthError {
match error {
SmsProviderError::InvalidVerifyCode => PhoneAuthError::InvalidVerifyCode,
@@ -4012,6 +4175,70 @@ mod tests {
);
}
#[tokio::test]
async fn stale_refresh_sessions_are_pruned_from_both_indexes() {
let store = build_store();
let refresh_service = build_refresh_service(store.clone());
let user = create_phone_login_user(store.clone(), "13800138008").await;
let now = OffsetDateTime::now_utc();
refresh_service
.create_session(
CreateRefreshSessionInput {
user_id: user.id.clone(),
refresh_token_hash: hash_refresh_session_token("stale-revoked"),
issued_by_provider: AuthLoginMethod::Password,
client_info: build_client_info(),
},
now - Duration::days(2),
)
.expect("stale session should create");
store
.revoke_session_by_refresh_token_hash(
&hash_refresh_session_token("stale-revoked"),
now - Duration::days(2),
)
.expect("stale session should revoke");
refresh_service
.create_session(
CreateRefreshSessionInput {
user_id: user.id.clone(),
refresh_token_hash: hash_refresh_session_token("recent-revoked"),
issued_by_provider: AuthLoginMethod::Password,
client_info: build_client_info(),
},
now,
)
.expect("recent session should create");
store
.revoke_session_by_refresh_token_hash(
&hash_refresh_session_token("recent-revoked"),
now,
)
.expect("recent session should revoke");
let projection = store
.export_projection_view(now.unix_timestamp())
.expect("projection export should prune stale sessions");
assert_eq!(projection.refresh_sessions.len(), 1);
assert_eq!(
projection.refresh_sessions[0].refresh_token_hash,
hash_refresh_session_token("recent-revoked")
);
let stale_error = refresh_service
.rotate_session(
RotateRefreshSessionInput {
refresh_token_hash: hash_refresh_session_token("stale-revoked"),
next_refresh_token_hash: hash_refresh_session_token("stale-next"),
},
now,
)
.expect_err("pruned session should no longer be indexed");
assert_eq!(stale_error, RefreshSessionError::SessionNotFound);
}
#[tokio::test]
async fn wechat_login_hits_existing_user_by_union_id_before_openid() {
let store = build_store();