合并最新master到设计代理分支
Project CI / Repository checks (pull_request) Failing after 4m24s
Project CI / Frontend tests (pull_request) Successful in 6m0s
Project CI / Backend tests (pull_request) Successful in 9m12s
Project CI / Native shell tests (pull_request) Successful in 16m21s

同步 master 的退款 outbox、认证投影与外部生成历史清理改动

保留 GDD 审批修复及双方决策记录
This commit is contained in:
2026-08-28 03:00:51 +00:00
78 changed files with 8665 additions and 1503 deletions
@@ -53,6 +53,7 @@ services:
- "host.docker.internal:host-gateway"
volumes:
- api-tracking-outbox:/var/lib/genarrative/tracking-outbox
- api-wallet-refund-outbox:/var/lib/genarrative/wallet-refund-outbox
ulimits:
nofile:
soft: 4096
@@ -85,6 +86,9 @@ services:
OTEL_SERVICE_NAME: genarrative-external-generation-worker
extra_hosts:
- "host.docker.internal:host-gateway"
volumes:
- external-generation-tracking-outbox:/var/lib/genarrative/tracking-outbox-worker
- external-generation-wallet-refund-outbox:/var/lib/genarrative/wallet-refund-outbox
ulimits:
nofile:
soft: 4096
@@ -142,4 +146,7 @@ services:
volumes:
spacetime-data:
api-tracking-outbox:
api-wallet-refund-outbox:
external-generation-tracking-outbox:
external-generation-wallet-refund-outbox:
nginx-logs:
@@ -11,6 +11,17 @@ 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
Environment=GENARRATIVE_DATABASE_BACKUP_STOP_MARKER=/var/lib/genarrative/database-backups/.spacetimedb-stopped
MemoryHigh=768M
MemoryMax=1G
OOMPolicy=stop
# 主进程可能在停库后被 MemoryMax/OOMPolicy 强制终止,JS finally 无法执行;
# 仅当备份脚本留下停库 marker 且本次 service 非正常成功时,由 systemd 兜底恢复全部依赖服务。
ExecStopPost=/bin/sh -c 'if [ "${SERVICE_RESULT}" != "success" ] && [ -f "${GENARRATIVE_DATABASE_BACKUP_STOP_MARKER}" ]; then systemctl start spacetimedb.service; systemctl restart genarrative-api.service; systemctl restart genarrative-external-generation-worker@1.service; systemctl restart genarrative-external-generation-controller.service; if systemctl is-active --quiet spacetimedb.service && systemctl is-active --quiet genarrative-api.service && systemctl is-active --quiet genarrative-external-generation-worker@1.service && systemctl is-active --quiet genarrative-external-generation-controller.service; then rm -f "${GENARRATIVE_DATABASE_BACKUP_STOP_MARKER}"; fi; fi'
# 备份需要停止 / 启动 spacetimedb.service,并读取 /stdb、写入 /var/lib/genarrative/database-backups。
# 停止 SpacetimeDB 会连带停止 Requires 它的 API / worker / controller,冷备份后必须显式拉起。
PrivateTmp=true
@@ -24,6 +24,31 @@
- 验证方式:首次 collecting 带 invented-confirmation 必须拒绝且不落 GDDreject continuation 再交 `user_revision` 的 v2 仍成功。
- 关联文档:`docs/technical/【技术方案】立项策划AgentFast GDD-2026-08-10.md`
## 2026-08-27 退款 emergency spool 容量溢出保持可恢复
## 2026-08-27 退款 emergency spool 容量溢出保持可恢复
- 背景:本机 emergency spool 仅作为 SpacetimeDB 完全不可达时的最后恢复路径,原有 `MAX_BYTES` 分支会直接返回 `Dropped`,导致扣费已经完成但没有可重放记录。
- 决策:达到普通 outbox `MAX_BYTES` 时,将退款记录写入同一持久目录的 `refund-overflow-*` 文件;该文件与普通 pending 文件一样由启动恢复和后台 worker 重放到 SpacetimeDB,且按 refund ledger id 保持幂等。溢出文件不计入普通阈值,但必须触发容量告警;底层磁盘写入失败仍进入关键退款人工补偿流程。
- 影响范围:api-server wallet refund emergency spool、资产失败退款日志、loadtest / 预览 Compose 持久卷、后端架构与开发运维文档。
- 验证方式:运行 api-server `wallet_refund_outbox` 定向测试,确认超限写入并保留 overflow 文件;运行 SpacetimeDB profile 测试、Compose 配置校验、编码和 diff 门禁。
## 2026-08-27 短期认证状态进入共享 typed projection
- 背景:短信验证码和微信 OAuth state 仍只存在 API 进程内 HashMap,多节点请求或 API 重启会直接丢失,无法满足无粘性会话的鉴权恢复要求。
- 决策:`AuthStoreProjectionView` 增加 `phone_codes``wechat_states` typed 字段,由 `auth_store_projection_meta` 以 JSON 投影持久化;启动恢复、CAS 同步和失败后的权威刷新都覆盖这两类短期状态。验证码哈希使用部署级稳定盐(当前复用 `GENARRATIVE_JWT_SECRET`),各 API 节点必须一致;发码前先刷新权威投影并用占位验证码记录做一次 projection CAS,只有占用成功才调用短信 provider,避免跨节点冷却竞态;认证 handler 在发码、消费验证码、创建/消费微信 state 后都要完成 projection sync,失败即返回服务错误;所有会读取或变更本机认证工作集的认证主链路(登录、刷新、`/me`、会话管理、密码、绑定和微信 state)在领域操作前先从正式投影做一次受 CAS 保护的只读刷新,受保护 Bearer 中间件也会在进入业务 handler 前执行同样的刷新,刷新失败时 fail closed,不能依赖粘性会话;同步遇到 CAS 冲突时,若本次尝试期间没有新的本地变更则恢复正式快照,若仍有待同步 revision 则由后续认证请求重试,避免节点永久卡在 pending。微信 OAuth state 设置有界活动数量,避免单个 JSON 投影无界膨胀。短期状态仍由 `module-auth` 内存工作集执行领域校验,但不再把本机 HashMap 当作持久化或跨节点真相。
- 影响范围:`module-auth` projection、`spacetime-module` auth schema/procedure、`spacetime-client` bindings/facade、api-server 手机号 / 微信 handler、认证架构与运维文档。
- 验证方式:运行 module-auth projection roundtrip(验证码可跨恢复校验、微信 state 可跨恢复消费)、SpacetimeDB schema/runtime/DDD 门禁、api-server 定向测试、编码和 diff 检查。
---
## 2026-08-27 外部生成历史采用受控保留清理
- 背景:`external_generation_job``external_generation_job_summary``external_generation_job_event` 都是持久化表;摘要和 payload 边界收紧后,已确认的终态历史仍会继续占用 SpacetimeDB 常驻内存,且事件审计链会随任务数量增长。
- 决策:新增仅 migration operator 可调用的 `prune_external_generation_job_history_and_return`。默认按 `source_module=editor-canvas`、30 天保留期和 `job_id` 游标分批运行;只删除主任务与摘要状态一致、属于 completed / failed / cancelled、摘要已有 `notification_acknowledged_at` 且终态时间达到 cutoff 的任务。事件、摘要和主任务仍按同一事务顺序删除,但每次事务最多删除 256 条事件;事件未删完时保留任务与摘要并返回同一个 job cursor,维护脚本下一次继续,避免单个任务形成无界事务写集。默认 dry-run,必须固定 dry-run 返回的 cutoff 后再 applypending / running、未确认通知、摘要缺失或状态不一致的数据永不删除。其他 source module 必须显式指定并单独评估;资产对象和钱包流水不随任务历史删除;不新增自动定时器或 runtime 清理权限。
- 影响范围:`server-rs/crates/spacetime-module/src/external_generation.rs`、外部生成事件 job_id 单列索引、SpacetimeDB 生成 bindings、`scripts/spacetime-maintain-external-generation-jobs.mjs`、架构与生产运维文档。
- 验证方式:覆盖终态 / 活跃态 / 已确认与未确认摘要、状态或身份不一致、cutoff 边界测试;运行 SpacetimeDB module tests/check、bindings 生成、schema/encoding/diff 门禁,并在维护窗口先 dry-run 再 apply。
- 关联文档:`docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md``docs/【开发运维】本地开发验证与生产运维-2026-05-15.md`、PR #203
## 2026-08-27 SpacetimeDB 工具链统一升级到 2.8.3
- 背景:SpacetimeDB 2.8.0 引入 TypeScript submodule 与调度延迟观测,2.8.1 修复 v1 WebSocket 订阅移除死锁、TypeScript SDK `array<u8>` 读缓存别名和 Rust string 默认值支持,2.8.2 修复 table accessor 改名自动迁移,2.8.3 修复 scheduled function 从实际执行时间重排导致的长期漂移。仓库若继续锁定 2.7.0,会保留这些已知运行时与 SDK 问题。
@@ -1334,8 +1359,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>。
@@ -7793,3 +7818,17 @@ CI 上 `background_agent_runtime_recovers_stale_running_before_pending_task` 在
- 决策:`.codex/skills/` 下的 SpacetimeDB 指导收敛为单一 `.codex/skills/genarrative-spacetimedb/SKILL.md`。官方 `spacetimedb` 插件负责通用 concepts、Rust server、CLI、TypeScript client 和 MCP 知识;项目 skill 只保留 Genarrative 的架构边界、schema/migration 门禁、目标 server 安全规则、运行时排障和验证路径。
- 路由:涉及 SpacetimeDB 的任务统一先读取项目适配 skill,再按需读取 `spacetimedb:concepts``spacetimedb:rust-server``spacetimedb:cli``spacetimedb:typescript-client``spacetimedb:mcp`。插件通用示例不得覆盖项目禁止 `maincloud`、禁止人工 `spacetime --root-dir`、显式 server 和后端分层等约束。
- 安装:团队环境缺少插件时使用 `codex plugin marketplace add clockworklabs/SpacetimeDB --sparse .agents --sparse codex-plugin``codex plugin add spacetimedb\@spacetimedb-plugins`;个人配置、缓存和凭据不进入仓库。
## 2026-08-27 退款 outbox 主路径迁入 SpacetimeDB
- 决策:`profile_wallet_refund_outbox` 是跨 API 节点退款的正式持久化队列。扣费失败、外部生成 attempt 失败或最终 lease 过期时,在同一个 SpacetimeDB 事务内按 `refund_ledger_id` 幂等写入 pending 行;worker 从库内 pending 行批量处理,退款账本写入与 outbox 成功删除保持在同一事务内,失败由 `available_at` / `attempts` 驱动重试。`asset_operation_wallet_settlement` 继续负责退款先于 consume 可见时的取消 intent,阻止迟到扣费。只有 SpacetimeDB 完全不可达时,api-server 才写本机 `wallet-refund-outbox` emergency spool;本机文件不能替代库内队列,必须持久挂载、告警、恢复演练并支持人工补偿。
- 影响范围:`profile_wallet_refund_outbox` 表及 bindings、runtime enqueue/process procedure、外部生成失败事务、inline 资产退款、api-server 跨节点 worker 和 emergency spool、后端架构与开发运维文档。
- 验证方式:运行 `npm run spacetime:generate``npm run check:spacetime-schema``npm run check:server-rs-ddd``cargo check -p spacetime-module -p spacetime-client -p api-server --manifest-path server-rs/Cargo.toml`、退款 outbox / asset billing / external generation 定向测试、`npm run check:encoding``git diff --check`
## 2026-08-27 release 内存增长修复与备份 OOM 恢复兜底
- 决策:外部生成 worker 每轮主动 `try_join_next` 回收已完成 `JoinHandle`,避免持续有队列任务时只归还 semaphore permit 却让 `JoinSet` 句柄集合无界增长;超时脱管任务在 abort 后等待句柄结束,执行许可保持到 work 真正结束或被取消。
- 决策:`module-ai` 的阶段终态写入与流式增量统一受文本、结构化 JSON、warning 和全局 retained 工作集上限约束;认证投影恢复对过滤后的 refresh session 重新计数,超过 8192 条直接拒绝启动恢复,避免超限快照灌入内存。
- 复审补充:`spacetime-module` 的 AI procedure 复用同一组任务元数据 / payload / 输出 / 结果引用上限,流式文本聚合超过 512 KiB 或每阶段 8192 个 chunk 时回滚事务,terminal task 收口后按有界批次删除 `ai_text_chunk` 明细,避免真实持久化链路绕过内存边界。
- 决策:备份脚本停库前写入受保护 `.spacetimedb-stopped` marker,正常恢复完成后清理;systemd 备份 service 通过 `MemoryHigh/MemoryMax/OOMPolicy``ExecStopPost` 在 Node OOM kill、无法执行 JS finally 时兜底拉起 SpacetimeDB、API、worker、controller,恢复不完整则保留 marker。
- 验证:worker/module-auth/module-ai 定向 Rust 测试、database-backup/production-ops/encoding 门禁和 `git diff --check` 必须在提交前通过;release 现场需按 archive-full 重新 provision 并核验旧 files-history drop-in 已删除、四个服务 active。
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
@@ -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}")
File diff suppressed because it is too large Load Diff
+28 -4
View File
@@ -821,6 +821,31 @@ 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:
'Environment=GENARRATIVE_DATABASE_BACKUP_STOP_MARKER=/var/lib/genarrative/database-backups/.spacetimedb-stopped',
reason:
'备份停库 marker 必须固定在受保护的 release work-dir,供 OOM 后 systemd 兜底恢复服务。',
},
{
file: 'deploy/systemd/genarrative-database-backup.service',
includes: 'MemoryMax=1G',
reason:
'备份 service 必须设置 systemd 内存硬上限,避免异常进程消耗整机内存。',
},
{
file: 'deploy/systemd/genarrative-database-backup.service',
includes: 'ExecStopPost=/bin/sh -c',
reason:
'备份主进程被 OOM kill 后必须由 systemd 兜底恢复停掉的 SpacetimeDB、API、worker 和 controller。',
},
{
file: 'deploy/systemd/genarrative-database-backup.service',
excludes: '--storage-format files',
@@ -941,10 +966,9 @@ const checks = [
},
{
file: 'jenkins/Jenkinsfile.production-server-provision',
excludes:
"params.DEPLOY_TARGET == 'release' && databaseBackupProfile == 'files-history'",
includes: 'release 仅允许 archive-fullfiles-history',
reason:
'release 必须能在显式选择 profile 且 baseline 预检通过后启用 files-history。',
'release 必须拒绝 files-history,避免逐文件 catalog 扫描再次触发生产内存峰值。',
},
{
file: 'scripts/database-backup-to-oss.mjs',
@@ -954,7 +978,7 @@ const checks = [
{
file: 'scripts/database-backup-to-oss.mjs',
includes:
'restoreServicesAfterBackup({stopService, serviceStopped, restartServicesAfter})',
'restoreServicesAfterBackup({stopService, serviceStopped, restartServicesAfter, stopMarkerPath})',
reason: '生产冷备份打包失败时也必须恢复 SpacetimeDB 及依赖服务。',
},
{
File diff suppressed because it is too large Load Diff
+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
@@ -8,23 +8,29 @@ import {
} from './spacetime-migration-common.mjs';
const MAX_BATCH_SIZE = 25;
const DEFAULT_RETENTION_DAYS = 30;
const MICROS_PER_DAY = 86_400_000_000;
function usage() {
return `用法:
node scripts/spacetime-maintain-external-generation-jobs.mjs --database <name> [选项]
默认只 dry-run 一批历史终态任务 payload 压缩,不修改数据库。
使用 --prune-history 时改为清理已确认通知且超过保留期的历史任务、摘要与事件。
公共选项:
--database <name> 目标数据库(必填,也可用 GENARRATIVE_SPACETIME_DATABASE
--server <name-or-url> spacetime CLI server 名或 URL
--server-url <url> 显式 server URL
--limit <1-${MAX_BATCH_SIZE}> 单批任务数,默认 10
--limit <1-${MAX_BATCH_SIZE}> 单批任务数,默认 10
--cursor-job-id <jobId> 从上一批 next_cursor_job_id 继续
--apply 执行写入;省略时始终 dry-run
--backfill-summaries 改为回填轻量摘要投影
--prune-history 改为清理已确认通知的终态历史
--owner-user-id <userId> 仅摘要回填可选,限定 owner
--completed-before-micros <n> 仅 payload 压缩可选,限定终态完成时间
--source-module <module> 仅历史清理可选,默认 editor-canvas
--retention-days <n> 仅历史清理可选,默认 ${DEFAULT_RETENTION_DAYS}
--completed-before-micros <n> 限定终态完成时间;历史清理默认按 retention-days 计算
--help 显示帮助
必须使用已授权 migration operator 的 spacetime CLI 登录态。脚本每次只处理一批;
@@ -40,6 +46,9 @@ function parseOptions(argv) {
database: process.env.GENARRATIVE_SPACETIME_DATABASE || '',
limit: 10,
ownerUserId: '',
pruneHistory: false,
retentionDays: DEFAULT_RETENTION_DAYS,
sourceModule: 'editor-canvas',
passthrough: [],
server: process.env.GENARRATIVE_SPACETIME_SERVER || '',
serverUrl: process.env.GENARRATIVE_SPACETIME_SERVER_URL || '',
@@ -82,6 +91,15 @@ function parseOptions(argv) {
options.apply = true;
} else if (arg === '--backfill-summaries') {
options.backfillSummaries = true;
} else if (arg === '--prune-history') {
options.pruneHistory = true;
} else if (arg === '--source-module') {
options.sourceModule = readValue(arg).trim();
if (!options.sourceModule) {
throw new Error('--source-module 不能为空。');
}
} else if (arg === '--retention-days') {
options.retentionDays = parsePositiveInteger(readValue(arg), arg);
} else if (arg === '--help' || arg === '-h') {
options.help = true;
} else {
@@ -95,12 +113,49 @@ function parseOptions(argv) {
if (options.ownerUserId && !options.backfillSummaries) {
throw new Error('--owner-user-id 只能与 --backfill-summaries 一起使用。');
}
if (options.backfillSummaries && options.pruneHistory) {
throw new Error('--backfill-summaries 与 --prune-history 不能同时使用。');
}
if (options.sourceModule !== 'editor-canvas' && !options.pruneHistory) {
throw new Error('--source-module 只能与 --prune-history 一起使用。');
}
if (
options.retentionDays !== DEFAULT_RETENTION_DAYS &&
!options.pruneHistory
) {
throw new Error('--retention-days 只能与 --prune-history 一起使用。');
}
if (options.completedBeforeMicros !== null && options.backfillSummaries) {
throw new Error('--completed-before-micros 不能用于摘要回填。');
}
if (
options.completedBeforeMicros !== null &&
options.pruneHistory &&
options.retentionDays !== DEFAULT_RETENTION_DAYS
) {
throw new Error(
'--completed-before-micros 与 --retention-days 不能同时使用。',
);
}
return options;
}
function resolveRetentionCutoffMicros(options) {
if (!options.pruneHistory) {
return options.completedBeforeMicros;
}
if (options.completedBeforeMicros !== null) {
return options.completedBeforeMicros;
}
const cutoff = Date.now() * 1000 - options.retentionDays * MICROS_PER_DAY;
if (!Number.isSafeInteger(cutoff)) {
throw new Error(
'--retention-days 计算出的 completed_before_micros 超出安全整数范围。',
);
}
return cutoff;
}
try {
const options = parseOptions(process.argv.slice(2));
if (options.help) {
@@ -113,24 +168,36 @@ try {
);
}
const procedureName = options.backfillSummaries
? 'backfill_external_generation_job_summaries_and_return'
: 'compact_external_generation_job_payloads_and_return';
const input = options.backfillSummaries
const completedBeforeMicros = resolveRetentionCutoffMicros(options);
const procedureName = options.pruneHistory
? 'prune_external_generation_job_history_and_return'
: options.backfillSummaries
? 'backfill_external_generation_job_summaries_and_return'
: 'compact_external_generation_job_payloads_and_return';
const input = options.pruneHistory
? {
owner_user_id: encodeSpacetimeCliOption(options.ownerUserId || null),
source_module: options.sourceModule,
limit: options.limit,
cursor_job_id: encodeSpacetimeCliOption(options.cursorJobId || null),
completed_before_micros: completedBeforeMicros,
dry_run: !options.apply,
}
: {
dry_run: !options.apply,
limit: options.limit,
cursor_job_id: encodeSpacetimeCliOption(options.cursorJobId || null),
completed_before_micros: encodeSpacetimeCliOption(
options.completedBeforeMicros,
),
};
: options.backfillSummaries
? {
owner_user_id: encodeSpacetimeCliOption(options.ownerUserId || null),
limit: options.limit,
cursor_job_id: encodeSpacetimeCliOption(options.cursorJobId || null),
dry_run: !options.apply,
}
: {
dry_run: !options.apply,
limit: options.limit,
cursor_job_id: encodeSpacetimeCliOption(options.cursorJobId || null),
completed_before_micros: encodeSpacetimeCliOption(
completedBeforeMicros,
),
};
const result = await callSpacetimeProcedureViaCli(
options,
procedureName,
@@ -138,10 +205,29 @@ try {
);
ensureProcedureOk(result);
console.log(JSON.stringify({ procedure: procedureName, ...result }, null, 2));
const pendingApplyCount = options.backfillSummaries
? Number(result.selected_count ?? 0)
: Number(result.matched_count ?? 0);
console.log(
JSON.stringify(
{
procedure: procedureName,
...(options.pruneHistory
? {
source_module: options.sourceModule,
completed_before_micros: completedBeforeMicros,
...(options.completedBeforeMicros === null
? { retention_days: options.retentionDays }
: {}),
}
: {}),
...result,
},
null,
2,
),
);
const pendingApplyCount =
options.pruneHistory || options.backfillSummaries
? Number(result.selected_count ?? 0)
: Number(result.matched_count ?? 0);
if (result.has_more && options.apply) {
console.log(
`仍有后续批次;下一次追加 --cursor-job-id ${result.next_cursor_job_id ?? '<missing>'}`,
@@ -150,8 +236,11 @@ try {
const currentCursor = options.cursorJobId
? `保留 --cursor-job-id ${options.cursorJobId}`
: '仍从首批开始';
const cutoffHint = options.pruneHistory
? `并固定 --completed-before-micros ${completedBeforeMicros}`
: '';
console.log(
`当前仅 dry-run;请${currentCursor}并追加 --apply 重跑同一批。apply 成功后再使用其 next_cursor_job_id 进入下一批。`,
`当前仅 dry-run;请${currentCursor}${cutoffHint}并追加 --apply 重跑同一批。apply 成功后再使用其 next_cursor_job_id 进入下一批。`,
);
}
} catch (error) {
+59 -12
View File
@@ -9,11 +9,13 @@ export function parseArgs(argv) {
'GENARRATIVE_SPACETIME_MIGRATION_CHUNK_SIZE',
),
database: process.env.GENARRATIVE_SPACETIME_DATABASE || '',
bootstrapSecret: process.env.GENARRATIVE_SPACETIME_MIGRATION_BOOTSTRAP_SECRET || '',
bootstrapSecret:
process.env.GENARRATIVE_SPACETIME_MIGRATION_BOOTSTRAP_SECRET || '',
bootstrapSecretFile:
process.env.GENARRATIVE_SPACETIME_MIGRATION_BOOTSTRAP_SECRET_FILE || '',
includeTables: [],
operatorIdentity: process.env.GENARRATIVE_SPACETIME_MIGRATION_OPERATOR_IDENTITY || '',
operatorIdentity:
process.env.GENARRATIVE_SPACETIME_MIGRATION_OPERATOR_IDENTITY || '',
passthrough: [],
note: '',
server: process.env.GENARRATIVE_SPACETIME_SERVER || '',
@@ -142,7 +144,9 @@ export function buildSpacetimeCallArgs(options, procedureName, input) {
export async function callSpacetimeProcedure(options, procedureName, input) {
if (!options.database) {
throw new Error('必须传入 --database,或设置 GENARRATIVE_SPACETIME_DATABASE。');
throw new Error(
'必须传入 --database,或设置 GENARRATIVE_SPACETIME_DATABASE。',
);
}
validateSpacetimeDatabaseName(options.database);
@@ -196,7 +200,9 @@ export async function createSpacetimeWebIdentity(options) {
const text = await response.text();
if (!response.ok) {
throw new Error(`SpacetimeDB identity HTTP ${response.status}: ${trimPreview(text)}`);
throw new Error(
`SpacetimeDB identity HTTP ${response.status}: ${trimPreview(text)}`,
);
}
let payload;
@@ -209,16 +215,25 @@ export async function createSpacetimeWebIdentity(options) {
}
const identity =
payload.identity ?? payload.Identity ?? payload.identity_hex ?? payload.identityHex;
payload.identity ??
payload.Identity ??
payload.identity_hex ??
payload.identityHex;
const token = payload.token ?? payload.Token;
if (typeof identity !== 'string' || typeof token !== 'string') {
throw new Error(`SpacetimeDB identity 响应缺少 identity/token: ${trimPreview(text)}`);
throw new Error(
`SpacetimeDB identity 响应缺少 identity/token: ${trimPreview(text)}`,
);
}
return { identity, token };
}
export async function callSpacetimeProcedureAuto(options, procedureName, input) {
export async function callSpacetimeProcedureAuto(
options,
procedureName,
input,
) {
if (options.useHttp) {
return callSpacetimeProcedure(options, procedureName, input);
}
@@ -226,7 +241,11 @@ export async function callSpacetimeProcedureAuto(options, procedureName, input)
return callSpacetimeProcedureViaCli(options, procedureName, input);
}
export async function callSpacetimeProcedureViaCli(options, procedureName, input) {
export async function callSpacetimeProcedureViaCli(
options,
procedureName,
input,
) {
const args = buildSpacetimeCallArgs(options, procedureName, input);
const output = await runSpacetimeCli(args);
return parseProcedureResult(output, procedureName);
@@ -335,7 +354,8 @@ function normalizeSatsProduct(value, procedureName) {
}
if (
procedureName === 'normalize_editor_character_animation_metadata_and_return' &&
procedureName ===
'normalize_editor_character_animation_metadata_and_return' &&
value.length === 19
) {
return {
@@ -427,6 +447,24 @@ function normalizeSatsProduct(value, procedureName) {
};
}
if (
procedureName === 'prune_external_generation_job_history_and_return' &&
value.length === 10
) {
return {
ok: normalizeSatsValue(value[0]),
dry_run: normalizeSatsValue(value[1]),
scanned_count: normalizeSatsValue(value[2]),
selected_count: normalizeSatsValue(value[3]),
deleted_job_count: normalizeSatsValue(value[4]),
deleted_summary_count: normalizeSatsValue(value[5]),
deleted_event_count: normalizeSatsValue(value[6]),
next_cursor_job_id: normalizeSatsOption(value[7]),
has_more: normalizeSatsValue(value[8]),
error_message: normalizeSatsOption(value[9]),
};
}
if (value.length === 3) {
return {
ok: normalizeSatsValue(value[0]),
@@ -497,7 +535,10 @@ function normalizeSatsValue(value) {
if (value && typeof value === 'object') {
return Object.fromEntries(
Object.entries(value).map(([key, entry]) => [key, normalizeSatsValue(entry)]),
Object.entries(value).map(([key, entry]) => [
key,
normalizeSatsValue(entry),
]),
);
}
@@ -581,7 +622,9 @@ export function resolveServerUrl(options) {
return 'http://127.0.0.1:3101';
}
throw new Error(`未知 SpacetimeDB server: ${server}。请改用 --server-url 显式传入地址。`);
throw new Error(
`未知 SpacetimeDB server: ${server}。请改用 --server-url 显式传入地址。`,
);
}
function resolveCliServer(options) {
@@ -635,7 +678,11 @@ function runSpacetimeCli(args) {
return;
}
if (code !== 0) {
reject(new Error(`spacetime call 失败,退出码 ${code}: ${trimPreview(output)}`));
reject(
new Error(
`spacetime call 失败,退出码 ${code}: ${trimPreview(output)}`,
),
);
return;
}
@@ -223,4 +223,35 @@ describe('SpacetimeDB CLI SATS option encoding', () => {
expect(objectResult.batch_sha256).toBe('d'.repeat(64));
expect(objectResult).not.toHaveProperty('batch_sha_256');
});
it('normalizes external generation history prune tuple results', () => {
const result = parseProcedureResult(
JSON.stringify([
true,
false,
25,
2,
2,
2,
10,
[0, 'job-25'],
true,
[0, '清理失败'],
]),
'prune_external_generation_job_history_and_return',
);
expect(result).toEqual({
ok: true,
dry_run: false,
scanned_count: 25,
selected_count: 2,
deleted_job_count: 2,
deleted_summary_count: 2,
deleted_event_count: 10,
next_cursor_job_id: 'job-25',
has_more: true,
error_message: '清理失败',
});
});
});
+1
View File
@@ -5435,6 +5435,7 @@ dependencies = [
"shared-contracts",
"spacetimedb",
"spacetimedb-lib",
"time",
]
[[package]]
+4
View File
@@ -657,6 +657,10 @@ pub async fn admin_list_editor_assets(
Extension(_admin): Extension<AuthenticatedAdmin>,
Query(query): Query<AdminEditorAssetListQuery>,
) -> Result<Json<Value>, AppError> {
state
.refresh_auth_store_from_spacetime()
.await
.map_err(map_admin_spacetime_error)?;
let page_size = query
.limit
.unwrap_or(ADMIN_EDITOR_ASSET_DEFAULT_LIMIT)
@@ -64,6 +64,7 @@ pub async fn admin_list_recharge_orders(
Extension(_admin): Extension<AuthenticatedAdmin>,
Query(query): Query<AdminRechargeOrderListQuery>,
) -> Result<Json<Value>, Response> {
refresh_auth_projection(&state, &request_context).await?;
let user_id = resolve_optional_user_id(&state, query.user_id, query.public_user_code)
.map_err(|error| error_response(&request_context, error))?;
let input = build_runtime_profile_recharge_order_admin_list_input(
@@ -112,6 +113,7 @@ pub async fn admin_get_user_detail(
Extension(admin): Extension<AuthenticatedAdmin>,
Query(query): Query<AdminUserDetailQuery>,
) -> Result<Json<Value>, Response> {
refresh_auth_projection(&state, &request_context).await?;
let user = resolve_user(&state, query.user_id, query.public_user_code)
.map_err(|error| error_response(&request_context, error))?;
let wallet_detail = state
@@ -181,6 +183,7 @@ pub async fn admin_reconcile_user_consumption(
Extension(admin): Extension<AuthenticatedAdmin>,
Json(payload): Json<AdminUserConsumptionReconcileRequest>,
) -> Result<Json<Value>, Response> {
refresh_auth_projection(&state, &request_context).await?;
let user = resolve_user(&state, Some(payload.user_id), None)
.map_err(|error| error_response(&request_context, error))?;
let input = build_runtime_profile_wallet_consumption_reconcile_input(
@@ -589,6 +592,7 @@ pub async fn admin_update_wallet_restriction(
Extension(admin): Extension<AuthenticatedAdmin>,
Json(payload): Json<AdminWalletRestrictionRequest>,
) -> Result<Json<Value>, Response> {
refresh_auth_projection(&state, &request_context).await?;
let user = resolve_user(&state, Some(payload.user_id), None)
.map_err(|error| error_response(&request_context, error))?;
let input = build_runtime_profile_wallet_manual_restriction_upsert_input(
@@ -1040,6 +1044,22 @@ fn map_hold(hold: RuntimeProfileRechargeRefundHoldSnapshot) -> AdminRechargeRefu
}
}
async fn refresh_auth_projection(
state: &AppState,
request_context: &RequestContext,
) -> Result<(), Response> {
state
.refresh_auth_store_from_spacetime()
.await
.map_err(|error| {
error_response(
request_context,
AppError::from_status(StatusCode::BAD_GATEWAY)
.with_message(format!("刷新用户认证信息失败:{error}")),
)
})
}
fn resolve_optional_user_id(
state: &AppState,
user_id: Option<String>,
+143 -69
View File
@@ -499,87 +499,127 @@ async fn refund_asset_operation_points_with_job_id(
external_generation_claim_attempt: Option<u32>,
) -> Result<(), AppError> {
let created_at_micros = current_utc_micros();
let metadata_json = wallet_metadata_json(
external_generation_job_id.as_deref(),
let current_attempt_is_owned_by_failure_transaction =
external_generation_job_id.as_deref().is_some_and(|job_id| {
current_external_generation_billing_context().is_some_and(|context| {
context.job_id == job_id
&& Some(context.claim_attempt) == external_generation_claim_attempt
})
});
if current_attempt_is_owned_by_failure_transaction {
// 队列当前 attempt 的 refund 由 fail_external_generation_job transaction 原子写入
// SpacetimeDB outbox;这里不能先写另一笔独立退款,避免任务成功写回后被误退。
return Ok(());
}
let settlement_reason = if external_generation_job_id.is_some() {
"stale_attempt_recovery"
} else {
"asset_operation_failed"
};
let enqueue_input = module_runtime::build_runtime_profile_wallet_refund_outbox_enqueue_input(
owner_user_id.clone(),
points_cost,
ledger_id.clone(),
created_at_micros,
asset_kind.clone(),
asset_id.clone(),
settlement_reason.to_string(),
external_generation_job_id.clone(),
external_generation_claim_attempt,
);
let result = state
)
.map_err(|error| {
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR).with_details(json!({
"provider": "profile-wallet-refund-outbox",
"message": error.to_string(),
}))
})?;
let enqueue_result = state
.spacetime_client()
.refund_profile_wallet_points_with_metadata(
owner_user_id.clone(),
points_cost,
ledger_id.clone(),
created_at_micros,
metadata_json,
)
.enqueue_profile_wallet_refund_outbox(enqueue_input)
.await;
if let Err(error) = result {
let refund_error = error.to_string();
let app_error = map_asset_operation_wallet_error(error);
if let Some(outbox) = state.wallet_refund_outbox() {
match outbox
.enqueue(WalletRefundOutboxRecord {
owner_user_id: owner_user_id.clone(),
amount: points_cost,
ledger_id: ledger_id.clone(),
created_at_micros,
asset_kind: asset_kind.clone(),
asset_id: asset_id.clone(),
external_generation_job_id: external_generation_job_id.clone(),
})
.await
{
Ok(WalletRefundOutboxEnqueueOutcome::Enqueued) => {
tracing::warn!(
owner_user_id,
asset_kind,
asset_id,
external_generation_job_id,
ledger_id,
error = %refund_error,
"资产操作失败后的泥点退款立即执行失败,已写入 wallet refund outbox"
);
}
Ok(WalletRefundOutboxEnqueueOutcome::Dropped { reason }) => {
tracing::error!(
owner_user_id,
asset_kind,
asset_id,
external_generation_job_id,
ledger_id,
reason,
error = %refund_error,
"资产操作失败后的泥点退款立即执行失败,且 wallet refund outbox 因容量限制丢弃"
);
}
Err(outbox_error) => {
tracing::error!(
owner_user_id,
asset_kind,
asset_id,
external_generation_job_id,
ledger_id,
refund_error = %refund_error,
outbox_error = %outbox_error,
"资产操作失败后的泥点退款立即执行失败,且写入 wallet refund outbox 失败"
);
}
}
} else {
tracing::error!(
match enqueue_result {
Ok(_) => {
tracing::info!(
owner_user_id,
asset_kind,
asset_id,
external_generation_job_id,
external_generation_claim_attempt,
ledger_id,
error = %refund_error,
"资产操作失败后的泥点退款失败,且 wallet refund outbox 未启用"
"资产操作失败后的泥点退款已写入 SpacetimeDB refund outbox"
);
Ok(())
}
return Err(app_error);
Err(error) if should_use_wallet_refund_emergency_spool(&error) => {
let refund_error = error.to_string();
let app_error = map_asset_operation_wallet_error(error);
if let Some(outbox) = state.wallet_refund_outbox() {
match outbox
.enqueue(WalletRefundOutboxRecord {
owner_user_id: owner_user_id.clone(),
amount: points_cost,
ledger_id: ledger_id.clone(),
created_at_micros,
asset_kind: asset_kind.clone(),
asset_id: asset_id.clone(),
settlement_reason: settlement_reason.to_string(),
external_generation_job_id: external_generation_job_id.clone(),
external_generation_claim_attempt,
})
.await
{
Ok(WalletRefundOutboxEnqueueOutcome::Enqueued) => {
tracing::warn!(
owner_user_id,
asset_kind,
asset_id,
external_generation_job_id,
ledger_id,
error = %refund_error,
"SpacetimeDB refund outbox 不可达,已写入本机 emergency spool"
);
}
Ok(WalletRefundOutboxEnqueueOutcome::OverflowEnqueued { reason }) => {
tracing::error!(
owner_user_id,
asset_kind,
asset_id,
external_generation_job_id,
ledger_id,
reason,
error = %refund_error,
"SpacetimeDB refund outbox 不可达,退款已写入本机 emergency spool overflow 文件;需监控并尽快恢复库内队列"
);
}
Err(outbox_error) => {
tracing::error!(
owner_user_id,
asset_kind,
asset_id,
external_generation_job_id,
ledger_id,
refund_error = %refund_error,
outbox_error = %outbox_error,
"SpacetimeDB refund outbox 不可达,且写入本机 emergency spool 失败"
);
}
}
} else {
tracing::error!(
owner_user_id,
asset_kind,
asset_id,
external_generation_job_id,
external_generation_claim_attempt,
ledger_id,
error = %refund_error,
"SpacetimeDB refund outbox 不可达,且本机 emergency spool 未启用"
);
}
Err(app_error)
}
Err(error) => Err(map_asset_operation_wallet_error(error)),
}
Ok(())
}
fn current_external_generation_billing_context() -> Option<ExternalGenerationBillingContext> {
@@ -683,6 +723,22 @@ pub(crate) fn should_skip_asset_operation_billing_for_connectivity(
}
}
fn should_use_wallet_refund_emergency_spool(error: &SpacetimeClientError) -> bool {
match error {
SpacetimeClientError::ConnectDropped | SpacetimeClientError::Timeout(_) => true,
SpacetimeClientError::Build(message)
| SpacetimeClientError::Procedure(message)
| SpacetimeClientError::Runtime(message) => {
message.contains("503")
|| message.contains("Service Unavailable")
|| message.contains("Failed to connect")
|| message.contains("WebSocket")
|| message.contains("连接已断开")
|| message.contains("连接在返回结果前已断开")
}
}
}
fn current_utc_micros() -> i64 {
time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64 / 1_000
}
@@ -838,6 +894,24 @@ mod tests {
));
}
#[test]
fn wallet_refund_emergency_spool_requires_database_unavailability() {
assert!(should_use_wallet_refund_emergency_spool(
&SpacetimeClientError::ConnectDropped
));
assert!(should_use_wallet_refund_emergency_spool(
&SpacetimeClientError::Runtime("503 Service Unavailable".to_string())
));
assert!(!should_use_wallet_refund_emergency_spool(
&SpacetimeClientError::Procedure(
"No such procedure: enqueue_profile_wallet_refund_outbox_and_return".to_string(),
)
));
assert!(!should_use_wallet_refund_emergency_spool(
&SpacetimeClientError::Procedure("泥点余额不足".to_string())
));
}
#[test]
fn asset_operation_wallet_insufficient_balance_is_public_message() {
for domain_message in [
+72 -38
View File
@@ -19,6 +19,7 @@ use serde_json::{Value, json};
use shared_contracts::auth::RuntimeGuestTokenResponse;
#[cfg(any())]
use shared_kernel::{format_rfc3339, new_uuid_simple_string};
#[cfg(test)]
use time::OffsetDateTime;
use tracing::warn;
@@ -145,6 +146,16 @@ pub async fn require_bearer_auth(
let Some(authenticated) = authenticate_request(&state, headers, request_id).await? else {
return Err(AppError::from_status(StatusCode::UNAUTHORIZED));
};
// JWT 会话校验走 SpacetimeDB;随后刷新用户/身份投影,保证所有受保护路由在
// 读取进程内工作集时都不会依赖粘性会话命中创建或更新它的 API 节点。
state
.refresh_auth_store_from_spacetime()
.await
.map_err(|error| {
warn!(error = %error, "受保护请求刷新认证投影失败");
AppError::from_status(StatusCode::SERVICE_UNAVAILABLE)
.with_message("认证状态服务暂不可用")
})?;
request.extensions_mut().insert(authenticated.clone());
let mut response = next.run(request).await;
@@ -236,54 +247,77 @@ async fn authenticate_request(
);
AppError::from_status(StatusCode::UNAUTHORIZED)
})?;
let current_user = state
.auth_user_service()
.get_user_by_id(claims.user_id())
.map_err(|error| {
warn!(
%request_id,
error = %error,
"Bearer JWT 用户快照读取失败"
);
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
})?;
let Some(current_user) = current_user else {
warn!(
%request_id,
user_id = %claims.user_id(),
"Bearer JWT 对应用户不存在"
);
return Err(AppError::from_status(StatusCode::UNAUTHORIZED));
};
if current_user.token_version != claims.token_version() {
warn!(
%request_id,
user_id = %claims.user_id(),
token_version = claims.token_version(),
current_token_version = current_user.token_version,
"Bearer JWT 版本已失效"
);
return Err(AppError::from_status(StatusCode::UNAUTHORIZED)
.with_message("当前登录态已失效,请重新登录"));
}
#[cfg(not(test))]
let session_is_active = state
.refresh_session_service()
.is_session_active_for_user(
claims.user_id(),
claims.session_id(),
OffsetDateTime::now_utc(),
)
.spacetime_client()
.validate_auth_session(spacetime_client::AuthSessionValidationRecordInput {
user_id: claims.user_id().to_string(),
session_id: claims.session_id().to_string(),
token_version: claims.token_version(),
})
.await
.map_err(|error| {
warn!(
%request_id,
user_id = %claims.user_id(),
session_id = %claims.session_id(),
error = %error,
"Bearer JWT refresh session 状态读取失败"
"Bearer JWT SpacetimeDB 会话状态读取失败"
);
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
})?;
#[cfg(test)]
let session_is_active = {
let current_user = state
.auth_user_service()
.get_user_by_id(claims.user_id())
.map_err(|error| {
warn!(
%request_id,
error = %error,
"Bearer JWT 用户快照读取失败"
);
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
})?;
let Some(current_user) = current_user else {
warn!(
%request_id,
user_id = %claims.user_id(),
"Bearer JWT 对应用户不存在"
);
return Err(AppError::from_status(StatusCode::UNAUTHORIZED));
};
if current_user.token_version != claims.token_version() {
warn!(
%request_id,
user_id = %claims.user_id(),
token_version = claims.token_version(),
current_token_version = current_user.token_version,
"Bearer JWT 版本已失效"
);
return Err(AppError::from_status(StatusCode::UNAUTHORIZED)
.with_message("当前登录态已失效,请重新登录"));
}
state
.refresh_session_service()
.is_session_active_for_user(
claims.user_id(),
claims.session_id(),
OffsetDateTime::now_utc(),
)
.map_err(|error| {
warn!(
%request_id,
user_id = %claims.user_id(),
session_id = %claims.session_id(),
error = %error,
"Bearer JWT refresh session 状态读取失败"
);
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
})?
};
if !session_is_active {
warn!(
%request_id,
@@ -15,6 +15,13 @@ pub async fn get_public_user_by_code(
Extension(request_context): Extension<RequestContext>,
Path(code): Path<String>,
) -> Result<Json<serde_json::Value>, AppError> {
state
.refresh_auth_store_from_spacetime()
.await
.map_err(|error| {
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
.with_message(format!("刷新认证状态失败:{error}"))
})?;
let user = state
.password_entry_service()
.get_user_by_public_user_code(&code)
@@ -41,6 +48,14 @@ pub async fn get_public_user_by_id(
return Err(AppError::from_status(StatusCode::BAD_REQUEST).with_message("用户 ID 不能为空"));
}
state
.refresh_auth_store_from_spacetime()
.await
.map_err(|error| {
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
.with_message(format!("刷新认证状态失败:{error}"))
})?;
let user = state
.auth_user_service()
.get_user_by_id(user_id)
@@ -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;
@@ -113,16 +118,29 @@ pub(crate) async fn run_external_generation_worker(state: AppState) -> Result<()
);
loop {
// 持续有队列任务时不会进入等待分支,因此必须在每轮主动回收已完成的
// JoinHandle;否则 permit 虽已归还,JoinSet 仍会保留每个历史任务的句柄。
reap_finished_external_generation_worker_tasks(&mut tasks);
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 +196,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,13 +277,11 @@ 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,
fn reap_finished_external_generation_worker_tasks(tasks: &mut JoinSet<()>) {
while let Some(result) = tasks.try_join_next() {
if let Err(error) = result {
error!(error = %error, "external generation worker 子任务 panic");
}
}
}
@@ -326,6 +346,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());
@@ -377,13 +398,14 @@ async fn process_external_generation_job(
job_id = %job.job_id,
job_kind = %job.job_kind,
timeout_seconds = job_timeout.as_secs(),
"external generation worker 任务超过执行预算,停止续租并释放 worker 槽位,在途执行交由租约仲裁"
"external generation worker 任务超过执行预算,停止续租并保留 worker 槽位,在途执行交由租约仲裁"
);
detach_external_generation_work_until_lease_expiry(
work_handle,
&job,
lease,
"任务超过执行预算",
Some(permit),
);
Err(message)
}
@@ -393,6 +415,7 @@ async fn process_external_generation_job(
&job,
lease,
"任务租约续期失败",
Some(permit),
);
Err(error)
}
@@ -417,6 +440,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();
@@ -445,6 +469,9 @@ fn detach_external_generation_work_until_lease_expiry(
),
Err(_) => {
work_handle.abort();
// 仅调用 abort 不会从 JoinHandle/JoinSet 中消费完成结果;等待被取消
// 的 handle,确保 permit 与任务句柄在同一生命周期内一起释放。
let _ = work_handle.await;
warn!(
job_id = %job_id,
job_kind = %job_kind,
@@ -454,6 +481,9 @@ fn detach_external_generation_work_until_lease_expiry(
);
}
}
// 保持执行许可直到 work 真正结束或被取消,避免超时任务脱管后继续
// 累积图片/音频响应占用。
drop(permit);
});
}
@@ -1972,6 +2002,7 @@ mod tests {
&job,
Duration::from_millis(200),
"任务超过执行预算",
None,
);
tokio::time::sleep(Duration::from_millis(100)).await;
@@ -1981,6 +2012,61 @@ 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_reaps_completed_tasks_while_queue_remains_busy() {
let mut tasks = JoinSet::new();
let completed = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
const TASK_COUNT: usize = 128;
for _ in 0..TASK_COUNT {
let completed = completed.clone();
tasks.spawn(async move {
completed.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
});
}
while completed.load(std::sync::atomic::Ordering::SeqCst) < TASK_COUNT {
tokio::task::yield_now().await;
}
assert_eq!(tasks.len(), TASK_COUNT);
reap_finished_external_generation_worker_tasks(&mut tasks);
assert!(tasks.is_empty(), "已完成任务的 JoinHandle 应在每轮被回收");
}
#[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 +2085,7 @@ mod tests {
&job,
Duration::from_millis(10),
"任务超过执行预算",
None,
);
let reacquired =
+8 -4
View File
@@ -529,21 +529,24 @@ async fn finalize_shutdown(context: ShutdownContext) {
}
if let Some(outbox) = context.wallet_refund_outbox {
info!(timeout_ms, "api-server 退出前 flush wallet refund outbox");
info!(
timeout_ms,
"api-server 退出前 flush wallet refund emergency spool"
);
match timeout(context.outbox_flush_timeout, outbox.flush_for_shutdown()).await {
Ok(Ok(())) => {
info!("api-server 退出前 wallet refund outbox flush 完成");
info!("api-server 退出前 wallet refund emergency spool flush 完成");
}
Ok(Err(error)) => {
warn!(
error = %error,
"api-server 退出前 wallet refund outbox flush 未完成,已保留本地文件等待下次启动重试"
"api-server 退出前 wallet refund emergency spool flush 未完成,已保留本地文件等待下次启动重试"
);
}
Err(_) => {
warn!(
timeout_ms,
"api-server 退出前 wallet refund outbox flush 超时,已保留本地文件等待下次启动重试"
"api-server 退出前 wallet refund emergency spool flush 超时,已保留本地文件等待下次启动重试"
);
}
}
@@ -557,6 +560,7 @@ fn spawn_common_app_state_background_workers(state: &AppState) {
if let Some(outbox) = state.wallet_refund_outbox() {
outbox.spawn_worker();
}
state.profile_wallet_refund_outbox_worker().spawn_worker();
}
fn spawn_http_app_state_background_workers(state: &AppState, process_role: ProcessRole) {
@@ -27,6 +27,13 @@ pub async fn password_entry(
headers: HeaderMap,
Json(payload): Json<PasswordEntryRequest>,
) -> Result<impl IntoResponse, AppError> {
state
.refresh_auth_store_from_spacetime()
.await
.map_err(|error| {
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
.with_message(format!("刷新认证状态失败:{error}"))
})?;
let input = PasswordEntryInput {
country_code: payload.country_code,
pure_phone_number: payload.pure_phone_number,
@@ -9,6 +9,7 @@ use shared_contracts::auth::{
PasswordChangeRequest, PasswordChangeResponse, PasswordResetRequest, PasswordResetResponse,
};
use time::OffsetDateTime;
use tracing::warn;
use crate::{
api_response::json_success_body,
@@ -81,7 +82,17 @@ pub async fn reset_password(
);
}
let result = state
// reset_password 消费的是跨节点共享的短期验证码;先恢复正式投影,
// 避免发码节点与消费节点的本机工作集不一致。
state
.refresh_auth_store_from_spacetime()
.await
.map_err(|error| {
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
.with_message(format!("刷新短信验证码状态失败:{error}"))
})?;
let result = match state
.phone_auth_service()
.reset_password(
ResetPasswordInput {
@@ -93,7 +104,17 @@ pub async fn reset_password(
OffsetDateTime::now_utc(),
)
.await
.map_err(map_phone_auth_error)?;
{
Ok(result) => result,
Err(error) => {
if let Err(sync_error) = state.sync_auth_store_tables_to_spacetime().await {
warn!(error = %sync_error, "重置密码失败后的短信验证码状态同步失败");
return Err(AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
.with_message("同步短信验证码状态失败"));
}
return Err(map_phone_auth_error(error));
}
};
let session_client = resolve_session_client_context(&headers);
let signed_session = create_auth_session(
&state,
+80 -20
View File
@@ -50,16 +50,33 @@ pub async fn send_phone_code(
phone_input_masked = phone_input_masked.as_str(),
"收到手机号验证码发送请求"
);
let send_input = SendPhoneCodeInput {
country_code: payload.country_code,
pure_phone_number: payload.pure_phone_number,
scene: scene.clone(),
};
let send_now = OffsetDateTime::now_utc();
state
.refresh_auth_store_from_spacetime()
.await
.map_err(|error| {
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
.with_message(format!("刷新短信验证码状态失败:{error}"))
})?;
state
.phone_auth_service()
.reserve_code_send(&send_input, send_now)
.map_err(map_phone_auth_error)?;
state
.sync_auth_store_tables_to_spacetime()
.await
.map_err(|error| {
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
.with_message(format!("占用短信验证码发送窗口失败:{error}"))
})?;
let result = match state
.phone_auth_service()
.send_code(
SendPhoneCodeInput {
country_code: payload.country_code,
pure_phone_number: payload.pure_phone_number,
scene: scene.clone(),
},
OffsetDateTime::now_utc(),
)
.send_code_after_authoritative_reservation(send_input, send_now)
.await
{
Ok(result) => {
@@ -91,6 +108,14 @@ pub async fn send_phone_code(
}
};
state
.sync_auth_store_tables_to_spacetime()
.await
.map_err(|error| {
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
.with_message(format!("同步短信验证码状态失败:{error}"))
})?;
Ok(json_success_body(
Some(&request_context),
PhoneSendCodeResponse {
@@ -114,19 +139,44 @@ pub async fn phone_login(
AppError::from_status(StatusCode::BAD_REQUEST).with_message("手机号登录暂未启用")
);
}
let invite_code = payload.invite_code.clone();
let result = match state
.phone_auth_service()
.login(
PhoneLoginInput {
country_code: payload.country_code,
pure_phone_number: payload.pure_phone_number,
verify_code: payload.code,
},
OffsetDateTime::now_utc(),
)
state
.refresh_auth_store_from_spacetime()
.await
{
.map_err(|error| {
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
.with_message(format!("刷新短信验证码状态失败:{error}"))
})?;
let invite_code = payload.invite_code.clone();
let login_input = PhoneLoginInput {
country_code: payload.country_code,
pure_phone_number: payload.pure_phone_number,
verify_code: payload.code,
};
let login_now = OffsetDateTime::now_utc();
let login_result = state
.phone_auth_service()
.login(login_input.clone(), login_now)
.await;
let login_result = match login_result {
Err(PhoneAuthError::VerifyCodeNotFound) => {
if let Err(sync_error) = state.refresh_auth_store_from_spacetime().await {
warn!(
request_id = request_context.request_id(),
operation = request_context.operation(),
error = %sync_error,
"手机号验证码未命中后的认证投影刷新失败"
);
return Err(AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
.with_message("刷新短信验证码状态失败"));
}
state
.phone_auth_service()
.login(login_input, login_now)
.await
}
result => result,
};
let result = match login_result {
Ok(result) => {
info!(
request_id = request_context.request_id(),
@@ -142,6 +192,16 @@ pub async fn phone_login(
result
}
Err(error) => {
if let Err(sync_error) = state.sync_auth_store_tables_to_spacetime().await {
warn!(
request_id = request_context.request_id(),
operation = request_context.operation(),
error = %sync_error,
"手机号验证码登录失败后的状态同步失败"
);
return Err(AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
.with_message("同步短信验证码状态失败"));
}
warn!(
request_id = request_context.request_id(),
operation = request_context.operation(),
@@ -265,6 +265,14 @@ async fn process_expired_virtual_payment_order(
state: &AppState,
order: &RuntimeProfileRechargeOrderRecord,
) -> Result<(), ExpirationCompensationError> {
state
.refresh_auth_store_from_spacetime()
.await
.map_err(|error| {
ExpirationCompensationError::Runtime(format!(
"failed to refresh auth projection for virtual payment query: {error}"
))
})?;
let identity = state
.wechat_auth_service()
.get_identity_by_user_id(&order.user_id)

Some files were not shown because too many files have changed in this diff Show More