diff --git a/docs/project-memory/shared-memory/decision-log.md b/docs/project-memory/shared-memory/decision-log.md index b908fffdd..0f022fc1b 100644 --- a/docs/project-memory/shared-memory/decision-log.md +++ b/docs/project-memory/shared-memory/decision-log.md @@ -16,6 +16,16 @@ --- +## 2026-07-23 BgFilter 失败审计使用硬上限与独立 tracking outbox + +- 背景:BgFilter worker 每个已发出的失败 provider attempt 都会启动 detached 审计任务;专用 worker 又关闭了 tracking outbox,使任务逐条等待 SpacetimeDB。`Q` 只约束内部 HTTP 请求生命周期,响应结束后无法限制仍在等待数据库的审计任务,部分失败、预算截短 timeout、重试恢复和熔断重置场景下可能持续堆积。 +- 决策:BgFilter worker 的失败审计在 `tokio::spawn` 前统一获取进程级 `1024` 个硬上限 permit,满载时直接丢弃并记录低基数指标,不创建等待任务。获准任务优先写入 worker 独立 tracking outbox,目录固定派生为共享 `GENARRATIVE_TRACKING_OUTBOX_DIR` 下的 `bgfilter-worker/` 子目录;worker 启动 outbox flush worker,退出时先排空已获准审计 enqueue,再封存并尽力 flush。BgFilter 专用策略在 outbox 缺失、容量拒绝或写盘失败时丢弃并观测,不回退同步直写 SpacetimeDB;其它外部 API 审计保持原有 fallback 语义。 +- 影响范围:`api-server` BgFilter worker、外部 API 失败审计策略、tracking outbox 进程接线、指标与测试、BgFilter 架构和开发运维文档;不修改 SpacetimeDB schema、procedure、bindings、前端或公开 DTO。 +- 验证方式:覆盖 spawn 前容量拒绝、flat / complex 共享总上限、permit 生命周期、独立 outbox 目录、outbox 满载 / 写盘失败不直写、SpacetimeDB 不可用时任务与磁盘保持有界,以及退出时 tracker drain 后再 flush;运行 api-server 定向测试、BgFilter fault smoke、Rust check、编码和 diff 检查。 +- 关联文档:`docs/technical/【后端架构】BgFilter受限资源调度方案-2026-07-21.md`、`docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md`、`docs/【开发运维】本地开发验证与生产运维-2026-05-15.md`。 + +--- + ## 2026-07-22 BgFilter flat 与 complex 使用独立熔断状态 - 背景:complex 请求在 provider 持续快速失败时仍会不断发起真实 provider attempt,并为每次已发出的失败生成异步审计;现有 flat 熔断不能约束 complex,且五分钟冷却会让短暂故障恢复后的等待过长。 diff --git a/docs/technical/【后端架构】BgFilter受限资源调度方案-2026-07-21.md b/docs/technical/【后端架构】BgFilter受限资源调度方案-2026-07-21.md index 6d6c23be3..948fdbdd5 100644 --- a/docs/technical/【后端架构】BgFilter受限资源调度方案-2026-07-21.md +++ b/docs/technical/【后端架构】BgFilter受限资源调度方案-2026-07-21.md @@ -89,7 +89,7 @@ flowchart LR 职责边界: - 父流程负责源对象已持久化、owner 校验、请求预算、flat fallback、Alpha / 尺寸恢复、动画 finalizer、最终 OSS / `asset_object` / 画布写回、计费和父终态。 -- `bgfilter-worker` 负责内部协议校验、OSS 签名、并发与排队上限、BgFilter 协议、两次顺序尝试、按模式隔离的熔断、provider 失败审计和结果图片校验。 +- `bgfilter-worker` 负责内部协议校验、OSS 签名、并发与排队上限、BgFilter 协议、两次顺序尝试、按模式隔离的熔断、provider 失败审计和结果图片校验;失败审计使用进程级 `1024` 个任务硬上限与 worker 独立 tracking outbox。 - SpacetimeDB 不参与本次内部调度;不新增表、reducer、procedure、facade 或生成 bindings。 `bgfilter-worker` 从实现形态看是只监听内部地址的同步 worker service,不是队列 consumer。父 worker 调另一个 worker 在这里是允许的:父进程明确选择保留调用栈和槽位,因此同步内部 HTTP 正是首版的最小交接方式。 @@ -423,7 +423,7 @@ BgFilter 成功二进制不是一份新的业务资产: 日志只写 `requestId`、父 job / request correlation、mode、attempt、排队耗时、provider 耗时、结果码和安全 object key;不得记录请求/响应图片 body。 -flat / complex 的每次 provider 失败审计都必须留在子 worker,保留“第一次失败、第二次成功”也可观察的事实。审计口径以“该次 attempt 是否已发出 provider HTTP”为界:已发出的失败一律写入 `external_api_call_failure`,包括被剩余预算截短后发生的 timeout 与 response 阶段超时(它们不计入熔断,但必须可审计);未发出的失败(预算不足未启动、签名失败)以及内部 admission、鉴权和本地配置错误不伪装成 BgFilter provider 失败。子 worker 进程角色不共享 api / extgen 的落盘 tracking outbox,失败审计由异步任务直写 SpacetimeDB,并纳入 shutdown tracker,优雅退出前排空;进程被强杀时可能丢失,属首版接受行为。 +flat / complex 的每次 provider 失败审计都必须留在子 worker,保留“第一次失败、第二次成功”也可观察的事实。审计口径以“该次 attempt 是否已发出 provider HTTP”为界:已发出的失败尝试写入 `external_api_call_failure`,包括被剩余预算截短后发生的 timeout 与 response 阶段超时(它们不计入熔断但仍属于审计候选);未发出的失败(预算不足未启动、签名失败)以及内部 admission、鉴权和本地配置错误不伪装成 BgFilter provider 失败。worker 在 `tokio::spawn` 前获取 flat / complex 共用的进程级 `1024` 个审计 permit,满载时直接丢弃并递增指标,不创建 semaphore waiter 或 detached task。获准任务写入共享 tracking outbox 根目录下独立的 `bgfilter-worker/` 子目录,由 worker 自己批量 flush;outbox 缺失、达到 `MAX_BYTES` 或写盘失败时只丢弃并观测,不同步直写 SpacetimeDB。审计任务纳入 shutdown tracker,优雅退出先排空 enqueue,再封存并尽力 flush;进程被强杀时,已写入 outbox 的记录可在下次启动重放,尚未 enqueue 的任务可能丢失。 ## 10. 实施与部署计划 @@ -478,7 +478,7 @@ flat / complex 的每次 provider 失败审计都必须留在子 worker,保留 - 父业务预算仍有效时,flat 两次失败、熔断、overload、内部 RPC deadline 或断连仍走“阿里云 → 本地”;complex 任意失败或自身熔断都直接失败,不接 flat fallback。 - flat / complex 分别按自身真实失败 attempt 计数且状态互不影响;由剩余业务预算截短的 timeout 不计入。两种模式都在 permit 前二次检查;已获准调用可完成第二次,后续同模式排队请求快速 `circuit_open`。 - `cancelled`(仅验证父侧映射,保留码首版不产生)、父 cancellation / 绝对 deadline、`invalid_request` 和 `unauthorized` 不启动 flat fallback;其它 flat 错误只在父业务预算仍有效时进入 fallback。 -- 已发出的 provider attempt 失败(含预算截短 timeout 与 response 阶段超时)全部落 `external_api_call_failure`;未发出与纯内部失败不落。审计任务由 shutdown tracker 排空后进程才退出。 +- 已发出的 provider attempt 失败(含预算截短 timeout 与 response 阶段超时)都是 `external_api_call_failure` 审计候选;未发出与纯内部失败不落。审计 task 在 spawn 前受全局 `1024` 硬上限约束,满载或独立 outbox 不可写时允许丢弃并上报指标;已获准任务由 shutdown tracker 排空并完成 enqueue 后,进程才封存和尽力 flush outbox。 - 客户端断连时,等待 permit 的请求最终由 deadline 收口;已开始 provider attempt 持有 permit 并排空。明确 cancellation 已被观察到后不再开始第二次,单纯 TCP 断连只作 best-effort 测试,不作为硬保证。 - 动画全部已提交帧继续 collect / drain,根因按现有稳定帧序号收口;首版不存在未实现的 group cancellation 承诺。 - 成功 body 为原始图片字节而非 Base64、JSON 或结果 object key;父侧 client 把该有界字节缓冲直接交给现有后处理,子 worker 不执行 raw OSS PUT。父、子两侧都拒绝空 body、MIME / 魔数不一致、chunked 超 `32 MiB` 和超过 `8192 × 8192` 的图片。 diff --git a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md index 8983970a6..bcf00d2e3 100644 --- a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md +++ b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md @@ -247,7 +247,7 @@ npm run check:server-rs-ddd - 抠图输入以私有 OSS 作为内存生命周期边界:生成原图和角色动作抽取帧上传时消费图片字节所有权,上传完成后不保留原图缓冲;手动去背景直接解析并校验已有 OSS object key,不下载原图。BgFilter 必须为 object key 签发 600 秒 GET URL 并通过 multipart `image_url` 提交,不用 `file` 重传;flat 链路进入阿里云 fallback 时由 `platform-matting` URL 接口单独下载并上传 `AuthorizeFileUpload` 临时对象,在推理前释放下载缓冲,继续 fallback 到本地键色时再单独下载一次原图,本地产出后释放本次原图下载缓冲。签名 URL 不得写入日志、审计或持久化。 - 角色动作抠图输入像素边界:仅图片画布角色动作链路在 FFmpeg 抽帧后、源帧上传 OSS 前,把帧解码为 RGB8,并按最终 `frameWidth × frameHeight` 的 contain 比例使用 `Triangle` 只缩放到内容尺寸;该阶段不得创建最终目标尺寸画布、不得引入 Alpha 通道,也不得插入任何 padding。BgFilter、阿里云通用抠图和本地键色降级共享这个无补边源帧 object key。抠图返回后才统一转为 RGBA8,按相同比例居中放入最终目标尺寸画布,并用 `RGBA(0,0,0,0)` 补齐透明 padding。以 `560×752 → 323×480` 为例,抠图输入固定为无 Alpha、无补边的 `323×434 RGB8 PNG`,最终输出为上下各 `23px` 透明补边的 `323×480 RGBA8 PNG`。旧 `/api/assets/character-animation/*` 动作发布链路继续保留原有帧 finalizer,不适用该输入规则。抽帧解码后若携带 Alpha 通道,必须先把像素按白底合成为不透明再转 RGB8,禁止直接丢弃 Alpha——全透明像素下未定义的 RGB 值会以杂色进入抠图输入,重新引入杂色边缘;共享 FFmpeg 抽帧命令保持不固定 `-pix_fmt`,白底合成只属于该链路的 BgFilter 输入准备阶段。 - 阿里云通用抠图的非上海地域输入不得使用 `viapiutils/GetOssStsToken`、固定 `viapi-customer-temp` 或 OSS V1 PUT。`platform-matting` 必须按官方新版 SDK Advance 协议调用 `AuthorizeFileUpload`,使用动态返回的单对象 Policy 执行 multipart POST,再把临时上海 OSS URL 交给 `SegmentCommonImage`;输入归一化、结果下载与原尺寸 Alpha 回贴继续留在同一适配器内。该协议仍上传图片字节,不等同于阿里云服务端直接抓取任意公网 URL,也不改变上层 BgFilter → 阿里云 → 本地降级顺序。 -- 编辑器抠图服务:手动 `POST /api/editor/images/background-removals` 与角色形象生成、图标 spritesheet 生成、UI 设计图素材提取、角色动作抽帧后的透明化统一通过唯一 loopback `bgfilter-worker` 调用 BgFilter provider。provider 配置继续使用 `GENARRATIVE_EDITOR_BGFILTER_BASE_URL` 与 `GENARRATIVE_EDITOR_BGFILTER_TOKEN`,默认 base URL 为 `http://58.87.105.82/bgfilter`;单次 provider attempt 上限不再独立配置,由公式 `N × est × 2` 运行时派生,其中 `est = GENARRATIVE_EDITOR_BGFILTER_SINGLE_IMAGE_ESTIMATE_MS`(默认 `5000`,依据为服务端高并发单图处理约 1-3s、网络约 3-5s),旧 `GENARRATIVE_EDITOR_BGFILTER_REQUEST_TIMEOUT_MS` 已删除;旧 `GENARRATIVE_EDITOR_BACKGROUND_REMOVAL_TOKEN` 只作为 provider token 的兼容回退别名,原手动去背景专用 base URL / timeout 配置已经删除。父流程先把候选 `objectKey`、`resourceId` 或 `assetId` 解析为当前 owner 已登记的私有 OSS object key;BFF 入队前统一拒绝 `data:` / `blob:`,底层 resolver 在解析引用前再次拒绝内联媒体并完成登记状态与 owner 校验。父流程只通过一次内部 HTTP RPC 传递 object key、排队预算 `maxQueueWaitMs`、调用预算 `callBudgetMs` 与模式参数,不传图片字节或签名 URL,并同步等待子 worker 返回的受限图片二进制 body。子 worker 在每次真实 provider attempt 前签发短期 OSS URL,承担 admission 保险丝 `Q`(默认 `2048`,仅防连接风暴)、provider 并发 `N`(生产 `16`);排队 deadline 从 `Q` admission 时刻起算,完成 JSON 校验并进入 provider permit 等待队列时再取得队长快照,按 `min((队长+5)×est×2, maxQueueWaitMs)` 约束排队等待。子 worker 还负责严格最多两次顺序 attempt、结果校验和按 flat / complex 隔离的进程级熔断;两种模式共享 `GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD=3` 和 `GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS=120` 默认值,但失败和成功只更新当前模式,且只由子 worker 读写。手动去背景固定使用 `background_mode=complex`、`seg_model=birefnet`、`cross_check=off`,不传 `file` 或 `screen_color`;complex provider 失败累计自身熔断,任意失败或自身熔断都直接返回父流程失败,不接 flat fallback,也不影响 flat 熔断。标准纯色背景四条链路固定使用 `background_mode=flat`、`screen_color=`、`seg_model=` 和 `cross_check=`,其中角色形象生成和角色动作逐帧去背传 `cross_check=on`,图标 spritesheet 生成和 UI 设计图素材提取传 `cross_check=off`。前端用户路径不展示抠图模型、模式或 cross-check,固定提交默认 `birefnet`,后端仍识别内部保留的 `anime-seg`;这些参数只属于后端内部供应商策略,不进入前端或外部 OpenAPI。父侧不重试已建立连接的内部 RPC,仅对 TCP 连接从未建立的失败按父预算有界退避重试(跨过 worker 重启与开机排序窗口,收到任何 HTTP 响应即停止);flat 两次 provider attempt 失败、熔断、overload、内部 deadline 或断连后,只要父业务预算仍有效,父流程才继续“阿里云通用抠图 → 本地 `editor_green_screen` 键色扣除”,熔断期不得直接退化到本地兜底。角色动作视频生成的背景色已与生图链路统一:`screenColor=auto` 时由视觉 LLM(`gpt-5-mini`,Responses 协议、low 推理档)读源角色图自动决策,并经硬过滤器剔除与前景 / 皮肤撞色的候选,手动 hex 则尊重用户选择;透明源角色图在提交 Ark 图生视频前先合成到选定背景色实色,使视频背景等于抠图键色;抽帧后每帧先上传私有 OSS 并释放原帧缓冲,再以 object key 固定使用 `seg_model=birefnet`、`cross_check=on` 进入上述三段式链路。阿里云通用抠图配置为 `GENARRATIVE_ALIYUN_MATTING_ENABLED`、`GENARRATIVE_ALIYUN_MATTING_ENDPOINT`、`GENARRATIVE_ALIYUN_MATTING_ACCESS_KEY_ID`、`GENARRATIVE_ALIYUN_MATTING_ACCESS_KEY_SECRET` 和 `GENARRATIVE_ALIYUN_MATTING_REQUEST_TIMEOUT_MS`;未配置专用 AK/SK 时可复用 `ALIBABA_CLOUD_ACCESS_KEY_ID` / `ALIBABA_CLOUD_ACCESS_KEY_SECRET`,默认 endpoint 为 `imageseg.cn-shanghai.aliyuncs.com`。标准纯色背景链路中,子 worker 已发出的 BgFilter provider 失败(含被剩余预算截短后发生的 timeout 与 response 阶段超时,这类失败不计入熔断但必须落审计)和父侧阿里云抠图链路已开始后的失败(包括源 OSS GET 成功后的解码、尺寸校验和归一化失败)都写入 `external_api_call_failure` 审计;真正开始外部调用前的本地预检不写该审计,并在 `failureStage` 中保留 `source_decode`、`source_validate` 等阶段。成功图片字节返回后,最终 Alpha / 尺寸恢复、OSS / asset object、画布写回、计费和父任务终态仍全部由父流程负责。 +- 编辑器抠图服务:手动 `POST /api/editor/images/background-removals` 与角色形象生成、图标 spritesheet 生成、UI 设计图素材提取、角色动作抽帧后的透明化统一通过唯一 loopback `bgfilter-worker` 调用 BgFilter provider。provider 配置继续使用 `GENARRATIVE_EDITOR_BGFILTER_BASE_URL` 与 `GENARRATIVE_EDITOR_BGFILTER_TOKEN`,默认 base URL 为 `http://58.87.105.82/bgfilter`;单次 provider attempt 上限不再独立配置,由公式 `N × est × 2` 运行时派生,其中 `est = GENARRATIVE_EDITOR_BGFILTER_SINGLE_IMAGE_ESTIMATE_MS`(默认 `5000`,依据为服务端高并发单图处理约 1-3s、网络约 3-5s),旧 `GENARRATIVE_EDITOR_BGFILTER_REQUEST_TIMEOUT_MS` 已删除;旧 `GENARRATIVE_EDITOR_BACKGROUND_REMOVAL_TOKEN` 只作为 provider token 的兼容回退别名,原手动去背景专用 base URL / timeout 配置已经删除。父流程先把候选 `objectKey`、`resourceId` 或 `assetId` 解析为当前 owner 已登记的私有 OSS object key;BFF 入队前统一拒绝 `data:` / `blob:`,底层 resolver 在解析引用前再次拒绝内联媒体并完成登记状态与 owner 校验。父流程只通过一次内部 HTTP RPC 传递 object key、排队预算 `maxQueueWaitMs`、调用预算 `callBudgetMs` 与模式参数,不传图片字节或签名 URL,并同步等待子 worker 返回的受限图片二进制 body。子 worker 在每次真实 provider attempt 前签发短期 OSS URL,承担 admission 保险丝 `Q`(默认 `2048`,仅防连接风暴)、provider 并发 `N`(生产 `16`);排队 deadline 从 `Q` admission 时刻起算,完成 JSON 校验并进入 provider permit 等待队列时再取得队长快照,按 `min((队长+5)×est×2, maxQueueWaitMs)` 约束排队等待。子 worker 还负责严格最多两次顺序 attempt、结果校验和按 flat / complex 隔离的进程级熔断;两种模式共享 `GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD=3` 和 `GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS=120` 默认值,但失败和成功只更新当前模式,且只由子 worker 读写。手动去背景固定使用 `background_mode=complex`、`seg_model=birefnet`、`cross_check=off`,不传 `file` 或 `screen_color`;complex provider 失败累计自身熔断,任意失败或自身熔断都直接返回父流程失败,不接 flat fallback,也不影响 flat 熔断。标准纯色背景四条链路固定使用 `background_mode=flat`、`screen_color=`、`seg_model=` 和 `cross_check=`,其中角色形象生成和角色动作逐帧去背传 `cross_check=on`,图标 spritesheet 生成和 UI 设计图素材提取传 `cross_check=off`。前端用户路径不展示抠图模型、模式或 cross-check,固定提交默认 `birefnet`,后端仍识别内部保留的 `anime-seg`;这些参数只属于后端内部供应商策略,不进入前端或外部 OpenAPI。父侧不重试已建立连接的内部 RPC,仅对 TCP 连接从未建立的失败按父预算有界退避重试(跨过 worker 重启与开机排序窗口,收到任何 HTTP 响应即停止);flat 两次 provider attempt 失败、熔断、overload、内部 deadline 或断连后,只要父业务预算仍有效,父流程才继续“阿里云通用抠图 → 本地 `editor_green_screen` 键色扣除”,熔断期不得直接退化到本地兜底。角色动作视频生成的背景色已与生图链路统一:`screenColor=auto` 时由视觉 LLM(`gpt-5-mini`,Responses 协议、low 推理档)读源角色图自动决策,并经硬过滤器剔除与前景 / 皮肤撞色的候选,手动 hex 则尊重用户选择;透明源角色图在提交 Ark 图生视频前先合成到选定背景色实色,使视频背景等于抠图键色;抽帧后每帧先上传私有 OSS 并释放原帧缓冲,再以 object key 固定使用 `seg_model=birefnet`、`cross_check=on` 进入上述三段式链路。阿里云通用抠图配置为 `GENARRATIVE_ALIYUN_MATTING_ENABLED`、`GENARRATIVE_ALIYUN_MATTING_ENDPOINT`、`GENARRATIVE_ALIYUN_MATTING_ACCESS_KEY_ID`、`GENARRATIVE_ALIYUN_MATTING_ACCESS_KEY_SECRET` 和 `GENARRATIVE_ALIYUN_MATTING_REQUEST_TIMEOUT_MS`;未配置专用 AK/SK 时可复用 `ALIBABA_CLOUD_ACCESS_KEY_ID` / `ALIBABA_CLOUD_ACCESS_KEY_SECRET`,默认 endpoint 为 `imageseg.cn-shanghai.aliyuncs.com`。标准纯色背景链路中,子 worker 已发出的 BgFilter provider 失败(含被剩余预算截短后发生的 timeout 与 response 阶段超时,这类失败不计入熔断但仍是审计候选)由进程级 `1024` 个审计任务硬上限保护,获准任务写入共享 tracking outbox 根目录下独立的 `bgfilter-worker/` 子目录并批量落库;满载、outbox 缺失、达到磁盘保护阈值或写盘失败时允许丢弃并记录指标,不回退逐条同步直写 SpacetimeDB。父侧阿里云抠图链路已开始后的失败(包括源 OSS GET 成功后的解码、尺寸校验和归一化失败)继续按通用外部 API 审计策略处理。真正开始外部调用前的本地预检不写该审计,并在 `failureStage` 中保留 `source_decode`、`source_validate` 等阶段。成功图片字节返回后,最终 Alpha / 尺寸恢复、OSS / asset object、画布写回、计费和父任务终态仍全部由父流程负责。 - BgFilter 连接复用、超时与动作帧流水线:`AppState` 分别复用父侧内部 worker HTTP Client 和子 worker 专用 BgFilter provider HTTP Client;父侧对一次逻辑调用至多让 worker 接收一次内部 RPC,不重试已建立连接后的失败;仅 TCP 连接从未建立时(worker 重启 / 开机排序窗口)按每轮重算 `maxQueueWaitMs` 的有界退避序列重连——增加的只是连接尝试次数,不产生第二次被接收的 RPC。重连配额按本进程是否已连通过 worker 分档:首连前(冷启动)flat 22.5s / complex 约 62.5s,首连后 flat ≤1.5s / complex 22.5s;每次重连计 `bgfilter_internal_connect_retry_total` 指标。子 worker 在同一个 `N` permit 内严格最多执行两次顺序 provider attempt。唯一子 worker 使用 `GENARRATIVE_BGFILTER_WORKER_CONCURRENCY=N`(生产 `16`)限制真实 provider 在途数;`GENARRATIVE_BGFILTER_WORKER_MAX_REQUESTS=Q` 降级为可选 admission 保险丝(默认 `2048`,仅防连接风暴,显式配置时必须 `>= N`)。超时全部由 `N` 与 `est` 运行时派生:单 attempt 上限 `N × est × 2`、调用预算 `callBudgetMs = 2 × attempt + 1s`(自取得 `N` permit 起算)、排队等待受 `min((provider 等待队列队长+5)×est×2, maxQueueWaitMs)` 双重上界(动态项充当自适应过载探测,超时带 `bound = estimate | parent` 标记),排队不侵蚀调用预算;`N` 与 `est` 必须同放共享 API 基础环境;请求携带的 `callBudgetMs` 只是父侧配置指纹,worker 比对后不一致只告警并计 `bgfilter_internal_call_budget_drift_total` 指标、始终以本进程公式值执行——发布调优 N / est 的新旧进程共存窗口不得误伤在途任务,持久漂移由部署脚本共享 env 对齐校验在启动前拦截。角色动作不再增加 `2000ms × 本次实际帧数`,`32 / 40 / 48` 帧使用相同公式。父侧按剩余绝对预算派生 `maxQueueWaitMs`(flat 扣除 `39s` 父侧预留(`37s` fallback + `2s` 传输窗),complex 只留 `2s` 传输窗;`<= 0` 时不发请求直接降级 / 失败),client timeout 取 `maxQueueWaitMs + callBudgetMs + 2s`;每次 attempt 前重新签发短期 OSS URL,剩余时间不足时不开始新的 attempt。父侧成功响应解码槽 `P = 8`。角色动作继续以 `buffer_unordered(frame_count.max(1))` 将全部单帧逻辑调用加入无序在途集合;返回结果携带原始帧序并在最终 collect / drain 全部已提交 Future 后排序,任一帧最终失败时必须先排空全部已启动 Future,再让整个动作任务失败退款,不能发布缺帧动画。单帧按“绿幕源图 owned 上传 OSS 并释放原帧 → 以 object key 调内部 worker / 按 object key 由父侧降级 → 父侧处理透明帧并落 OSS”流水化。角色动画源帧 PUT、透明帧 PUT 和最终帧 HEAD 仍统一复用 `AppState` 内初始化一次的 OSS HTTP Client(连接池参数为 connect 30 秒、request 60 秒、idle 300 秒、每 host 8 个 idle 连接、TCP keepalive 60 秒),并受进程级 8 路 OSS semaphore 限制;BgFilter provider 的 `N` 不占该 OSS permit,阿里云和本地处理既不占 OSS permit,也不受 `N / Q` 限制。每个 OSS 网络 attempt 单独获取 permit,退避期间释放;PUT/HEAD 动画帧请求最多 3 次(250ms、500ms 退避),只重试无 HTTP 响应的传输错误、timeout、OSS PutObject 的 `400 + RequestTimeout`、PUT `400` 错误体读取失败(未解析出 `Code`,按 timeout/transport 归类)、408、429 和 5xx。动作帧 PUT 只在 400 响应中有界读取最多 16 KiB OSS 错误 XML,并保留 `Code` 与响应头优先的 `x-oss-request-id`;错误体读取超时/断流时保留已读字节,已解析出的 `Code` 优先生效,未解析出 `Code` 则按 timeout/transport 归类重试;除 `RequestTimeout` 与该错误体读取失败情形外的其他 400、401/403/404、配置、URL/签名和空请求体错误不重试。最终帧 HEAD 失败只重试 HEAD,不重复 PUT。 - Match3D 物品 sheet:关卡整图完成后走 VectorEngine `/v1/images/edits` multipart `image`,模型为 `gpt-image-2`,`2K 1:1` 输出 `10*10` spritesheet;物品 sheet prompt 固定要求单一纯绿色 `#00FF00 / RGB(0,255,0)` 绿幕背景,后端上传 OSS 前必须把绿幕扣成透明 PNG,并把透明整图写入 `itemSpritesheetImageSrc/itemSpritesheetImageObjectKey`。后端优先按透明 alpha 连通域从该 sheet 识别真实素材矩形并持久化 20 个物品、每个 5 个形态;识别数量不足时才回退 `10*10` 固定网格。通用系列素材图集的行列索引按每行 2 个物品计算,必须落在 `1..=10`,难度只决定运行态加载 3 / 9 / 15 / 20 种。 - Match3D UI spritesheet 和背景派生图:关卡整图作为参考图并发生成 `1K 1:1` UI spritesheet 与 `1K 9:16` 背景图,模型均为 `gpt-image-2`。UI spritesheet prompt 固定要求单一纯绿色 `#00FF00 / RGB(0,255,0)` 绿幕背景,后端上传 OSS 前必须把绿幕扣成透明 PNG;背景图必须合成为全画幅不透明 PNG。 @@ -256,7 +256,7 @@ npm run check:server-rs-ddd - Hyper3D / Rodin:只保留后端安全代理和旧数据兼容;Rodin 提交、状态、下载和响应解析归属 `platform-hyper3d`,`api-server/src/hyper3d_generation.rs` 只做路由、配置和错误 envelope 映射;新 Match3D 草稿和批量新增不再生成 GLB。 - 音频:视觉小说专用音频路由保留;VectorEngine Suno/Vidu provider 协议、任务提交/查询、音频 URL 提取、下载、MIME/extension 归一和 OSS put 请求准备归属 `platform-audio`。`api-server/src/vector_engine_audio_generation.rs` 只做路由、配置、计费、asset object confirm、entity binding 和错误 envelope 映射;拼图、抓大鹅和敲木鱼提示词生成音效入口暂时关闭,通用 `/api/creation/audio/*` 对这些目标返回 `410 Gone`。敲木鱼创作只接收上传 / 录音音频资产;前端选择或录音阶段只在浏览器本地处理待提交音频,统一限制裁切后最长 1 秒、裁掉前后声音过小片段,并用浏览器端近似响度算法平衡到 `-15 LKFS` 后做峰值保护。点击生成时才直传 OSS 并确认 `asset_object`,创作 JSON 只提交轻量 `WoodenFishAudioAsset`,不得继续上传 Data URL 音频;未提供时由 `api-server` 写回内置默认木鱼音 `/wooden-fish/default-hit-sound.mp3`。 - OSS:私有 generated path 进入浏览器前必须通过 `/api/assets/read-url` 换签;不要裸请求 `/generated-*`。请求参数的安全语义不能混用:`legacyPublicPath` 是历史公开作品兼容口,只允许 `platform_oss::LEGACY_PUBLIC_PREFIXES` 中的 curated 前缀匿名换签;`objectKey` 是正式对象引用,绝不能复用该前缀旁路,必须查询 `asset_object` 并校验配置 bucket、精确 key、`PublicRead` 或当前 owner。External OpenAPI 的 `/api/external/v1/assets/read-url` 还必须有 `editor:asset` scope,并始终以 API Key 绑定的 `owner_user_id` 执行同一 owner 校验;后台跨账号预览只能走管理员鉴权后的 `/admin/api/assets/read-url`。`/api/assets/read-bytes` 与主站 read-url 共用完全相同的授权,默认仍应由浏览器使用 signed URL 直读,bytes 只作跨域字节读取 fallback。前端如果收到同一 OSS bucket 的完整 `https://*.oss-*.aliyuncs.com/generated-*` 地址,也必须先归一为 legacy path 后走同一换签链路,避免裸连私有 bucket 403 或绕过签名缓存。OSS 签名、读签名、HEAD 和 PUT 的结构化日志由 `platform-oss` 输出,排查资产写入 / 确认失败时优先按 `operation`、`object_key` / `key_prefix`、`status_class`、`error_kind` 和 `elapsed_ms` 下钻。新上传 generated 私有对象默认写入 `Cache-Control: public, max-age=31536000, immutable`;旧对象若缺该头,只能依赖 `ETag` / `Last-Modified` 协商缓存,应通过 OSS 元数据刷新或 CDN 配置补齐,不要恢复 api-server 静态代理。`editor-agent/` 前缀只用于服务端内部读写画布 Agent 会话消息文档,不属于浏览器直传 legacy public prefix;`/api/assets/direct-upload-tickets` 必须拒绝 `legacyPrefix=editor-agent`,内部读取只允许 `editor-agent/{conversationId}.json` 形态。 -- 外部 API 失败审计:外部供应商调用未成功时,`api-server` 必须发送 OTLP 失败事件并写入 `tracking_event`。VectorEngine 图片 provider 在 `platform-image` 内输出结构化日志和 `PlatformImageFailureAudit`,覆盖 `request_send`、`response_body`、`upstream_status`、`response_parse`、`missing_image` 和 `image_download` 阶段;编辑器 `screenColor=auto` 的 gpt-5-mini 背景色决策同样必须审计每次已发出的 LLM 调用失败,包括传输 / 超时、上游拒绝、响应体解析、空响应和返回候选外颜色;即使随后降级默认背景色并继续主流程也不得只记 warning。`api-server` 将这些失败映射成 `external_api_call_failure`,`scope_kind = module`、`scope_id = provider`、`module_key = external-api`。metadata 固定包含 provider、endpoint、operation、failureStage、statusCode、statusClass、timeout、retryable、errorMessage、latencyMs、promptChars、referenceImageCount、imageModel、rawExcerpt,以及在调用方可获得上下文时补充的 `userId`(触发者)和 `profileId`(草稿 / 作品 / 场景作用域)。图片生成入口应优先把 owner user id 和 profile id 透传到失败审计,不要只保留 provider 级聚合,否则很难按“谁触发、哪个作品触发”定位问题。入库优先复用 tracking outbox,outbox 不可写或保护阈值拒绝时回退同步写 SpacetimeDB;不得新增前端兜底或在 SpacetimeDB reducer 内做外部 I/O。唯一例外是 `bgfilter-worker` 进程角色:它不共享 api / extgen 的落盘 outbox 路径(避免多进程并发操作同一目录),provider 失败审计由 shutdown tracker 跟踪的异步任务直写 SpacetimeDB,优雅退出前排空,进程被强杀时可能丢失。 +- 外部 API 失败审计:外部供应商调用未成功时,`api-server` 必须发送 OTLP 失败事件并写入 `tracking_event`。VectorEngine 图片 provider 在 `platform-image` 内输出结构化日志和 `PlatformImageFailureAudit`,覆盖 `request_send`、`response_body`、`upstream_status`、`response_parse`、`missing_image` 和 `image_download` 阶段;编辑器 `screenColor=auto` 的 gpt-5-mini 背景色决策同样必须审计每次已发出的 LLM 调用失败,包括传输 / 超时、上游拒绝、响应体解析、空响应和返回候选外颜色;即使随后降级默认背景色并继续主流程也不得只记 warning。`api-server` 将这些失败映射成 `external_api_call_failure`,`scope_kind = module`、`scope_id = provider`、`module_key = external-api`。metadata 固定包含 provider、endpoint、operation、failureStage、statusCode、statusClass、timeout、retryable、errorMessage、latencyMs、promptChars、referenceImageCount、imageModel、rawExcerpt,以及在调用方可获得上下文时补充的 `userId`(触发者)和 `profileId`(草稿 / 作品 / 场景作用域)。图片生成入口应优先把 owner user id 和 profile id 透传到失败审计,不要只保留 provider 级聚合,否则很难按“谁触发、哪个作品触发”定位问题。普通调用入库优先复用 tracking outbox,outbox 不可写或保护阈值拒绝时回退同步写 SpacetimeDB;不得新增前端兜底或在 SpacetimeDB reducer 内做外部 I/O。`bgfilter-worker` 是受限资源例外:它使用共享 tracking outbox 基础目录下独立的 `bgfilter-worker/` 子目录,provider 失败审计在 spawn 前受进程级 `1024` 硬上限保护并由 shutdown tracker 跟踪;满载、outbox 缺失、保护阈值拒绝或写盘失败时直接丢弃并观测,不回退同步直写 SpacetimeDB。优雅退出先排空已获准任务的 enqueue,再封存并尽力 flush;进程被强杀时只有已 enqueue 记录可在下次启动重放。 - 外部生成运行记录:所有外部生成编排的完成态统一写入 `tracking_event`,`event_key = external_generation_run`,`scope_kind = module`,`scope_id = provider`,`module_key = external-generation`。metadata 固定包含 `runId`、`provider`、`operation`、`requestLabel`、`requestPayload`、`status`、`success`、`failureReason`、`providerRequestId`、`resultPayload`、`startedAtMicros`、`completedAtMicros` 和 `durationMs`。这类记录只用于运行审计和排障,不再走 `ai_task` 旧表。 ## SpacetimeDB 表目录 diff --git a/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md b/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md index 0e6451a5c..96f99c19b 100644 --- a/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md +++ b/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md @@ -734,7 +734,7 @@ cargo test -p platform-auth --manifest-path server-rs/Cargo.toml aliyun_send_sms 个人任务首版 scope 仅支持 `user`。每日登录任务按北京时间自然日 0 点重置;用户已登录并停留在“我的”页跨日时,前端需要先非阻断调用 refresh session 以写入新业务日 `daily_login`,再请求 `/api/profile/tasks` 刷新任务中心。认证成功后的 `daily_login` 必须通过 `SpacetimeClient::record_daily_login_tracking_event(...)` 调用 SpacetimeDB 专用 `record_daily_login_tracking_event_and_return` procedure,由数据库事务时间生成当日幂等事件并推进任务进度;不要改回普通 `record_tracking_event_after_success`、tracking outbox 或旧 `profile.login.daily` 事件键。后台、RPG、大鱼吃小鱼、Visual Novel、Story、Combat 等特定链路按 tracking 中间件排除规则处理;作品游玩统一使用 `work_play_start`。 -外部 API 失败审计复用 `tracking_event`,不新增表。失败事件优先写入本机 tracking outbox,再由后台 worker 批量落库;如果 outbox 因权限、磁盘或保护阈值不可写,会回退同步直写 SpacetimeDB。`metadata_json` 包含 endpoint、operation、failureStage、statusCode、statusClass、timeout、retryable、errorMessage、errorSource、latencyMs、promptChars、referenceImageCount、imageModel、rawExcerpt、userId、profileId 和 requestId;其中 `userId` 是触发生成的用户,`profileId` 是调用方传入的草稿 / 作品 / 场景作用域,`requestId` 用于回查同一次 HTTP 请求日志,入口拿不到上下文时允许为空。常用查询: +外部 API 失败审计复用 `tracking_event`,不新增表。普通 API / external-generation 调用的失败事件优先写入本机 tracking outbox,再由后台 worker 批量落库;如果 outbox 因权限、磁盘或保护阈值不可写,仍回退同步直写 SpacetimeDB。BgFilter worker 是受限资源例外:provider 失败审计在 spawn 前受进程级 `1024` 硬上限保护,获准任务写入 `GENARRATIVE_TRACKING_OUTBOX_DIR/bgfilter-worker/` 独立目录;任务满载、outbox 缺失、达到保护阈值或写盘失败时直接丢弃并记录指标,不同步直写。`metadata_json` 包含 endpoint、operation、failureStage、statusCode、statusClass、timeout、retryable、errorMessage、errorSource、latencyMs、promptChars、referenceImageCount、imageModel、rawExcerpt、userId、profileId 和 requestId;其中 `userId` 是触发生成的用户,`profileId` 是调用方传入的草稿 / 作品 / 场景作用域,`requestId` 用于回查同一次 HTTP 请求日志,入口拿不到上下文时允许为空。常用查询: ```sql SELECT event_id, scope_id AS provider, metadata_json, occurred_at @@ -774,7 +774,7 @@ GENARRATIVE_TRACKING_OUTBOX_MAX_BYTES=268435456 GENARRATIVE_API_SHUTDOWN_OUTBOX_FLUSH_TIMEOUT_MS=5000 ``` -outbox 采用 NDJSON 文件保存原始事件。达到 `BATCH_SIZE` 时会立刻把当前 active 文件原子封存为 sealed 文件,并马上切到新的 active 继续写入;后台 worker 异步 flush sealed 文件,HTTP 请求线程不等待 SpacetimeDB。`FLUSH_INTERVAL_MS` 只负责兜底封存长时间未满批的 active 文件。SpacetimeDB 批量 procedure 返回成功后删除 sealed 文件,失败则保留文件并重试。`MAX_BYTES` 是磁盘保护阈值,不是 flush 阈值;超过后低价值 route tracking 可以被丢弃并记录日志 / 指标,关键同步事件不进入该丢弃路径。sealed 文件若出现无法解析的坏行,会重命名为 `corrupt-*` 隔离并记录 `genarrative.tracking_outbox.files.corrupt` 指标,避免一个坏文件阻塞后续批量入库。api-server 收到退出信号后会在 `GENARRATIVE_API_SHUTDOWN_OUTBOX_FLUSH_TIMEOUT_MS` 窗口内封存 active 文件并尽力 flush sealed 文件,超时或 SpacetimeDB 暂不可用时保留本地文件给下次启动继续投递。该机制提供至少一次投递语义,依赖 `tracking_event.event_id` 幂等跳过重复事件。 +outbox 采用 NDJSON 文件保存原始事件。达到 `BATCH_SIZE` 时会立刻把当前 active 文件原子封存为 sealed 文件,并马上切到新的 active 继续写入;后台 worker 异步 flush sealed 文件,HTTP 请求线程不等待 SpacetimeDB。`FLUSH_INTERVAL_MS` 只负责兜底封存长时间未满批的 active 文件。SpacetimeDB 批量 procedure 返回成功后删除 sealed 文件,失败则保留文件并重试。`MAX_BYTES` 是每个 outbox 实例的磁盘保护阈值,不是 flush 阈值;超过后低价值 route tracking 和 BgFilter provider 失败审计可以被丢弃并记录日志 / 指标,关键同步事件不进入该丢弃路径。api-server 使用配置目录本身,BgFilter worker 固定使用其 `bgfilter-worker/` 子目录,两个进程不得操作同一个 active 文件。sealed 文件若出现无法解析的坏行,会重命名为 `corrupt-*` 隔离并记录 `genarrative.tracking_outbox.files.corrupt` 指标,避免一个坏文件阻塞后续批量入库。进程收到退出信号后会在 `GENARRATIVE_API_SHUTDOWN_OUTBOX_FLUSH_TIMEOUT_MS` 窗口内封存各自 active 文件并尽力 flush sealed 文件,超时或 SpacetimeDB 暂不可用时保留本地文件给下次同角色启动继续投递。该机制对已 enqueue 记录提供至少一次投递语义,依赖 `tracking_event.event_id` 幂等跳过重复事件;BgFilter 尚未 enqueue 或因硬上限 / 保护阈值被丢弃的审计不在该保证内。 release 机器如果日志每秒刷 `tracking outbox ... Permission denied (os error 13)`,先检查 `/etc/genarrative/api-server.env` 是否缺少 `GENARRATIVE_TRACKING_OUTBOX_DIR`。缺少时 `api-server` 会回退到本地开发默认相对路径 `server-rs/.data/tracking-outbox`,而 systemd 的工作目录是只读发布目录 `/opt/genarrative/releases/`,`genarrative` 用户无法在其中创建 `server-rs`。修复顺序: diff --git a/server-rs/crates/api-server/src/bgfilter_worker.rs b/server-rs/crates/api-server/src/bgfilter_worker.rs index 4cb64cbd6..b94bf9258 100644 --- a/server-rs/crates/api-server/src/bgfilter_worker.rs +++ b/server-rs/crates/api-server/src/bgfilter_worker.rs @@ -47,6 +47,9 @@ const BGFILTER_MAX_QUEUE_WAIT_MS: u64 = 21_600_000; const BGFILTER_QUEUE_ESTIMATE_HEADROOM: u64 = 5; const BGFILTER_QUEUE_ESTIMATE_SAFETY_FACTOR: u64 = 2; const BGFILTER_PROVIDER_MAX_ATTEMPTS: usize = 2; +/// provider 失败审计 task 的进程级硬上限。必须在 `tokio::spawn` 前获取 permit, +/// 避免 SpacetimeDB 或本机 outbox 变慢时形成无界 detached task / semaphore waiter。 +const BGFILTER_AUDIT_MAX_IN_FLIGHT: usize = 1_024; const BGFILTER_PROVIDER_ATTEMPT_RESERVE: Duration = Duration::from_secs(1); const BGFILTER_INTERNAL_CLIENT_RESPONSE_RESERVE: Duration = Duration::from_secs(2); const BGFILTER_SOURCE_URL_EXPIRE_SECONDS: u64 = 600; @@ -75,6 +78,8 @@ struct BgfilterMetrics { call_budget_drift_total: Counter, connect_retry_total: Counter, in_flight: UpDownCounter, + audit_in_flight: UpDownCounter, + audit_dropped_total: Counter, internal_request_seconds: Histogram, provider_http_seconds: Histogram, internal_request_total: Counter, @@ -136,6 +141,16 @@ fn bgfilter_metrics() -> &'static BgfilterMetrics { .with_unit("{request}") .with_description("Logical provider calls holding an N permit") .build(), + audit_in_flight: meter + .i64_up_down_counter("bgfilter_audit_in_flight") + .with_unit("{task}") + .with_description("BgFilter provider failure audit tasks holding a hard-limit permit") + .build(), + audit_dropped_total: meter + .u64_counter("bgfilter_audit_dropped_total") + .with_unit("{event}") + .with_description("BgFilter provider failure audits dropped before task spawn") + .build(), internal_request_seconds: meter .f64_histogram("bgfilter_internal_request_seconds") .with_unit("s") @@ -248,6 +263,7 @@ struct BgfilterWorkerRuntime { app_state: AppState, admission: Arc, provider: Arc, + audit: Arc, /// 正在等待 provider permit 的请求数;进入 provider permit 等待队列时的快照即 /// 「排在我前面的队长」,FIFO 语义下后来者不影响先到者的等待,快照可直接用于动态排队估时。 queue_depth: Arc, @@ -287,6 +303,7 @@ impl BgfilterWorkerRuntime { app_state, admission: Arc::new(Semaphore::new(max_requests)), provider: Arc::new(Semaphore::new(concurrency)), + audit: Arc::new(Semaphore::new(BGFILTER_AUDIT_MAX_IN_FLIGHT)), queue_depth: Arc::new(AtomicUsize::new(0)), internal_token: Arc::from(token), task_tracker, @@ -1154,6 +1171,7 @@ async fn execute_logical_request( audit_provider_attempt_failure( runtime.app_state.clone(), &runtime.task_tracker, + &runtime.audit, request.clone(), attempt, attempt_started, @@ -1860,20 +1878,76 @@ fn should_audit_provider_attempt_failure( && (error.transport || error.status_code.is_some() || error.invalid_result || error.timeout) } +struct BgfilterAuditPermit { + _permit: OwnedSemaphorePermit, +} + +impl BgfilterAuditPermit { + fn try_acquire(limiter: &Arc) -> Result { + let permit = limiter.clone().try_acquire_owned()?; + bgfilter_metrics().audit_in_flight.add(1, &[]); + Ok(Self { _permit: permit }) + } +} + +impl Drop for BgfilterAuditPermit { + fn drop(&mut self) { + bgfilter_metrics().audit_in_flight.add(-1, &[]); + } +} + +fn try_reserve_provider_failure_audit( + audit_limiter: &Arc, + attempt_started: bool, + error: &ProviderAttemptError, +) -> Option { + if !should_audit_provider_attempt_failure(attempt_started, error) { + return None; + } + + match BgfilterAuditPermit::try_acquire(audit_limiter) { + Ok(permit) => Some(permit), + Err(error) => { + let reason = match error { + TryAcquireError::NoPermits => "capacity", + TryAcquireError::Closed => "closed", + }; + bgfilter_metrics() + .audit_dropped_total + .add(1, &[KeyValue::new("reason", reason)]); + None + } + } +} + +fn try_start_provider_failure_audit( + task_tracker: &BgfilterTaskTracker, + audit_limiter: &Arc, + attempt_started: bool, + error: &ProviderAttemptError, +) -> Option<(BgfilterAuditPermit, BgfilterTaskGuard)> { + let audit_permit = try_reserve_provider_failure_audit(audit_limiter, attempt_started, error)?; + let task_guard = task_tracker.track_started(); + Some((audit_permit, task_guard)) +} + fn audit_provider_attempt_failure( state: AppState, task_tracker: &BgfilterTaskTracker, + audit_limiter: &Arc, request: BgfilterInternalRequest, attempt: usize, attempt_started: bool, error: ProviderAttemptError, ) { - if !should_audit_provider_attempt_failure(attempt_started, &error) { + let Some((audit_permit, task_guard)) = + try_start_provider_failure_audit(task_tracker, audit_limiter, attempt_started, &error) + else { return; - } - let task_guard = task_tracker.track_started(); + }; tokio::spawn(async move { let _task_guard = task_guard; + let _audit_permit = audit_permit; let audit = request .audit_context .unwrap_or(BgfilterInternalAuditContext { @@ -1887,7 +1961,7 @@ fn audit_provider_attempt_failure( request_id: audit.request_id, external_call_deadline: None, }; - crate::external_api_audit::record_matting_external_api_failure( + crate::external_api_audit::record_matting_external_api_failure_outbox_only( &state, &context, "bgfilter", @@ -2553,6 +2627,7 @@ mod tests { app_state: AppState::new(AppConfig::default()).expect("test state should build"), admission: Arc::new(Semaphore::new(admission)), provider: Arc::new(Semaphore::new(provider)), + audit: Arc::new(Semaphore::new(BGFILTER_AUDIT_MAX_IN_FLIGHT)), queue_depth: Arc::new(AtomicUsize::new(0)), internal_token: Arc::from("shared-token"), task_tracker: BgfilterTaskTracker::new(), @@ -2902,6 +2977,48 @@ mod tests { )); } + #[test] + fn provider_failure_audit_capacity_is_reserved_before_tracking_and_released() { + let limiter = Arc::new(Semaphore::new(1)); + let tracker = BgfilterTaskTracker::new(); + let error = ProviderAttemptError::upstream( + "upstream 500".to_string(), + 500, + None, + Duration::from_millis(10), + ); + + let reservation = try_start_provider_failure_audit(&tracker, &limiter, true, &error) + .expect("first audit should reserve the only permit"); + assert_eq!(limiter.available_permits(), 0); + assert_eq!( + tracker.inner.state.lock().unwrap().in_flight, + 1, + "accepted audit should be included in shutdown drain" + ); + + assert!(try_start_provider_failure_audit(&tracker, &limiter, true, &error).is_none()); + assert_eq!( + tracker.inner.state.lock().unwrap().in_flight, + 1, + "capacity rejection must not register another tracked task" + ); + + drop(reservation); + assert_eq!(limiter.available_permits(), 1); + assert_eq!(tracker.inner.state.lock().unwrap().in_flight, 0); + } + + #[test] + fn worker_runtime_uses_fixed_process_wide_audit_limit() { + let runtime = test_runtime(4, 2); + + assert_eq!( + runtime.audit.available_permits(), + BGFILTER_AUDIT_MAX_IN_FLIGHT + ); + } + #[test] fn provider_response_deadline_uses_earlier_attempt_or_rpc_deadline() { let now = Instant::now(); diff --git a/server-rs/crates/api-server/src/external_api_audit.rs b/server-rs/crates/api-server/src/external_api_audit.rs index f8eff93b9..5d2773b00 100644 --- a/server-rs/crates/api-server/src/external_api_audit.rs +++ b/server-rs/crates/api-server/src/external_api_audit.rs @@ -185,6 +185,45 @@ pub(crate) async fn record_matting_external_api_failure( record_external_api_failure(state, draft).await; } +/// BgFilter worker 专用入口:保留同一份 OTLP / tracking draft,但只允许写入本进程 +/// 独立 outbox。outbox 缺失、满载或写盘失败时丢弃,禁止在受限 worker 中逐条同步 +/// 直写 SpacetimeDB。 +#[allow(clippy::too_many_arguments)] +pub(crate) async fn record_matting_external_api_failure_outbox_only( + state: &AppState, + context: &ExternalApiAuditContext, + provider: &'static str, + endpoint: String, + operation: &'static str, + failure_stage: &'static str, + status_code: Option, + timeout: bool, + transport: bool, + latency_ms: Option, + error_message: String, + raw_excerpt: Option, +) { + let draft = build_matting_external_api_failure_draft( + provider, + endpoint, + operation, + failure_stage, + status_code, + timeout, + transport, + latency_ms, + error_message, + raw_excerpt, + context, + ); + record_external_api_failure_with_policy( + state, + draft, + ExternalApiAuditPersistencePolicy::RequireOutboxDropOnFailure, + ) + .await; +} + /// 构建抠图失败审计 draft。`transport` 必须由调用方从错误结构化字段读取, /// 不能用 `status_code.is_none()` 反推——本地处理失败同样没有上游 HTTP 状态。 #[allow(clippy::too_many_arguments)] @@ -352,8 +391,33 @@ pub(crate) fn app_error_status_class(status_code: StatusCode) -> &'static str { status_class(Some(status_code.as_u16())) } +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum ExternalApiAuditPersistencePolicy { + PreferOutboxThenSync, + RequireOutboxDropOnFailure, +} + +impl ExternalApiAuditPersistencePolicy { + fn allows_sync_fallback(self) -> bool { + matches!(self, Self::PreferOutboxThenSync) + } +} + /// 中文注释:外部供应商失败同时进入 OTLP 和 tracking_event;失败审计不能反向阻断主业务错误返回。 pub(crate) async fn record_external_api_failure(state: &AppState, draft: ExternalApiFailureDraft) { + record_external_api_failure_with_policy( + state, + draft, + ExternalApiAuditPersistencePolicy::PreferOutboxThenSync, + ) + .await; +} + +async fn record_external_api_failure_with_policy( + state: &AppState, + draft: ExternalApiFailureDraft, + persistence_policy: ExternalApiAuditPersistencePolicy, +) { record_external_api_failure_otlp(&draft); let tracking_event = build_external_api_failure_tracking_draft(&draft); @@ -366,6 +430,18 @@ pub(crate) async fn record_external_api_failure(state: &AppState, draft: Externa { Ok(crate::tracking_outbox::TrackingOutboxEnqueueOutcome::Enqueued) => {} Ok(crate::tracking_outbox::TrackingOutboxEnqueueOutcome::Dropped { reason }) => { + if !persistence_policy.allows_sync_fallback() { + crate::telemetry::record_external_api_audit_dropped(draft.provider, reason); + tracing::warn!( + provider = draft.provider, + endpoint = %draft.endpoint, + operation = %draft.operation, + failure_stage = draft.failure_stage, + reason, + "外部 API 失败审计写入专用 outbox 被保护阈值拒绝,已丢弃" + ); + return; + } tracing::warn!( provider = draft.provider, endpoint = %draft.endpoint, @@ -382,6 +458,21 @@ pub(crate) async fn record_external_api_failure(state: &AppState, draft: Externa .await; } Err(error) => { + if !persistence_policy.allows_sync_fallback() { + crate::telemetry::record_external_api_audit_dropped( + draft.provider, + "outbox_error", + ); + tracing::warn!( + provider = draft.provider, + endpoint = %draft.endpoint, + operation = %draft.operation, + failure_stage = draft.failure_stage, + error = %error, + "外部 API 失败审计写入专用 outbox 失败,已丢弃" + ); + return; + } tracing::warn!( provider = draft.provider, endpoint = %draft.endpoint, @@ -401,6 +492,18 @@ pub(crate) async fn record_external_api_failure(state: &AppState, draft: Externa return; } + if !persistence_policy.allows_sync_fallback() { + crate::telemetry::record_external_api_audit_dropped(draft.provider, "outbox_missing"); + tracing::warn!( + provider = draft.provider, + endpoint = %draft.endpoint, + operation = %draft.operation, + failure_stage = draft.failure_stage, + "外部 API 失败审计缺少专用 outbox,已丢弃" + ); + return; + } + crate::tracking::record_tracking_event_after_success( state, &audit_request_context(), @@ -572,6 +675,14 @@ mod tests { use super::*; + #[test] + fn bgfilter_outbox_only_policy_never_allows_sync_fallback() { + assert!(ExternalApiAuditPersistencePolicy::PreferOutboxThenSync.allows_sync_fallback()); + assert!( + !ExternalApiAuditPersistencePolicy::RequireOutboxDropOnFailure.allows_sync_fallback() + ); + } + #[test] fn external_api_failure_tracking_draft_uses_module_scope_and_safe_metadata() { let draft = build_external_api_failure_tracking_draft( diff --git a/server-rs/crates/api-server/src/main.rs b/server-rs/crates/api-server/src/main.rs index 0fe68f7b6..4f2b485dc 100644 --- a/server-rs/crates/api-server/src/main.rs +++ b/server-rs/crates/api-server/src/main.rs @@ -202,17 +202,18 @@ async fn run_bgfilter_worker_role(mut config: AppConfig) -> Result<(), io::Error let outbox_flush_timeout = config.shutdown_outbox_flush_timeout; let listener = build_tcp_listener(bind_address, listen_backlog)?; - // 专用 worker 不共享 api/extgen 进程的落盘 outbox,避免多个进程并发操作同一路径。 - // provider 失败审计仍通过 AppState 的无 outbox 路径 best-effort 写入 SpacetimeDB。 - config.tracking_outbox_enabled = false; - config.wallet_refund_outbox_enabled = false; + configure_bgfilter_worker_outboxes(&mut config); let state = AppState::new_with_empty_auth_store(config) .map_err(|error| io::Error::other(format!("初始化 bgfilter-worker 状态失败:{error}")))?; let (router, task_tracker) = build_bgfilter_worker_router(state.clone()) .map_err(|error| io::Error::other(format!("初始化 bgfilter-worker 路由失败:{error}")))?; + let tracking_outbox = state.tracking_outbox(); + if let Some(outbox) = tracking_outbox.clone() { + outbox.spawn_worker(); + } let shutdown_context = ShutdownContext { app_state: Some(state), - tracking_outbox: None, + tracking_outbox, wallet_refund_outbox: None, outbox_flush_timeout, }; @@ -236,6 +237,13 @@ async fn run_bgfilter_worker_role(mut config: AppConfig) -> Result<(), io::Error result } +fn configure_bgfilter_worker_outboxes(config: &mut AppConfig) { + // 多进程不能操作同一个 active 文件;worker 从共享基础目录派生自己的持久子目录。 + config.tracking_outbox_enabled = true; + config.tracking_outbox_dir = config.tracking_outbox_dir.join("bgfilter-worker"); + config.wallet_refund_outbox_enabled = false; +} + const DEFAULT_BGFILTER_WORKER_MAX_REQUESTS_FUSE: usize = 2_048; fn required_bgfilter_worker_capacity_from_env() -> Result<(usize, u64, usize), io::Error> { @@ -697,7 +705,7 @@ fn is_valid_env_key(key: &str) -> bool { #[cfg(test)] mod tests { use super::{ - AUTH_STORE_STARTUP_RETRY_INTERVAL, is_valid_env_key, + AUTH_STORE_STARTUP_RETRY_INTERVAL, configure_bgfilter_worker_outboxes, is_valid_env_key, parse_required_bgfilter_worker_capacity, protected_env_keys_from, should_initialize_editor_generation_pricing_for_startup, should_restore_auth_store_for_startup, should_start_profile_recharge_expiration_listener, @@ -744,6 +752,20 @@ mod tests { } } + #[test] + fn bgfilter_worker_uses_its_own_tracking_outbox_directory() { + let mut config = AppConfig::default(); + let base_dir = config.tracking_outbox_dir.clone(); + config.tracking_outbox_enabled = false; + config.wallet_refund_outbox_enabled = true; + + configure_bgfilter_worker_outboxes(&mut config); + + assert!(config.tracking_outbox_enabled); + assert_eq!(config.tracking_outbox_dir, base_dir.join("bgfilter-worker")); + assert!(!config.wallet_refund_outbox_enabled); + } + #[test] fn load_env_key_can_strip_utf8_bom_prefix() { let key = "\u{feff}SMS_AUTH_ENABLED" diff --git a/server-rs/crates/api-server/src/telemetry.rs b/server-rs/crates/api-server/src/telemetry.rs index f46e801b0..99ca82e41 100644 --- a/server-rs/crates/api-server/src/telemetry.rs +++ b/server-rs/crates/api-server/src/telemetry.rs @@ -180,6 +180,16 @@ pub(crate) fn record_external_api_failure( ); } +pub(crate) fn record_external_api_audit_dropped(provider: &'static str, reason: &'static str) { + external_api_metrics().audit_dropped.add( + 1, + &[ + KeyValue::new("provider", provider), + KeyValue::new("reason", reason), + ], + ); +} + fn track_response_body_in_flight(response: Response) -> Response { response.map(|body| { HTTP_RESPONSE_BODY_IN_FLIGHT.fetch_add(1, Ordering::Relaxed); @@ -219,6 +229,7 @@ struct TrackingOutboxMetrics { struct ExternalApiMetrics { failures: Counter, + audit_dropped: Counter, } struct HttpRequestPermitsAvailableGauges { @@ -363,6 +374,12 @@ fn external_api_metrics() -> &'static ExternalApiMetrics { "External API call failures grouped by provider and failure stage", ) .build(), + audit_dropped: meter + .u64_counter("genarrative.external_api.audit.dropped") + .with_description( + "External API failure audit records dropped when synchronous fallback is disabled", + ) + .build(), } }) }