为 complex 增加独立 BgFilter 熔断

为 flat 与 complex 分别维护熔断状态并补充双模式指标。

将统一熔断冷却默认值从五分钟降为两分钟。

为部署和 Provision 增加旧默认值定向迁移及门禁测试。

补充熔断状态隔离测试并同步架构、运维与项目决策文档。
This commit is contained in:
2026-07-22 14:39:41 +00:00
parent 835fe20c63
commit ec4fd552e3
13 changed files with 214 additions and 102 deletions
+1 -1
View File
@@ -33,7 +33,7 @@ GENARRATIVE_WALLET_REFUND_OUTBOX_MAX_BYTES=67108864
GENARRATIVE_BGFILTER_WORKER_CONCURRENCY=16
GENARRATIVE_EDITOR_BGFILTER_SINGLE_IMAGE_ESTIMATE_MS=5000
GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD=3
GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS=300
GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS=120
# BgFilter 失败后的中间兜底:阿里云通用抠图(SegmentCommonImage)。AccessKey 留空则跳过该层,
# BgFilter 失败直接本地 editor_green_screen 去背;填入后恢复 BgFilter→阿里云→本地三级兜底。
# AccessKey 也可复用标准 SDK 命名 ALIBABA_CLOUD_ACCESS_KEY_ID / ALIBABA_CLOUD_ACCESS_KEY_SECRET。
+2 -2
View File
@@ -7,9 +7,9 @@ GENARRATIVE_BGFILTER_WORKER_HOST=127.0.0.1
GENARRATIVE_BGFILTER_WORKER_PORT=8083
# Qadmission 保险丝,只防连接风暴;正常业务不应触达,显式配置时必须 >= N。
GENARRATIVE_BGFILTER_WORKER_MAX_REQUESTS=2048
# flat 熔断只由本进程维护;complex 不读写熔断状态
# flat / complex 熔断只由本进程维护;两种模式共享参数,但状态互相独立
GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD=3
GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS=300
GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS=120
GENARRATIVE_API_LOG=info,tower_http=info
GENARRATIVE_OTEL_ENABLED=true
@@ -16,6 +16,15 @@
---
## 2026-07-22 BgFilter flat 与 complex 使用独立熔断状态
- 背景:complex 请求在 provider 持续快速失败时仍会不断发起真实 provider attempt,并为每次已发出的失败生成异步审计;现有 flat 熔断不能约束 complex,且五分钟冷却会让短暂故障恢复后的等待过长。
- 决策:把现有 flat 熔断行为按原语义复用到 complex。flat / complex 共享 `GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD=3``GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS=120`,但在唯一 `bgfilter-worker` 内分别维护独立的连续失败数和打开截止时间;真实 provider attempt 的失败或成功只更新当前模式。两种模式都在排队前及取得 provider permit 后、第一次真实 HTTP 前检查自身熔断;已经通过第二次检查的逻辑调用仍可完成自己的第二次顺序 attempt。complex 熔断仍直接使父流程失败,不获得 flat 的阿里云 / 本地 fallback;本次不修改失败审计的异步处理流程。
- 部署边界:deploy / Provision 将 worker env 中历史模板默认 cooldown `300` 定向迁移到 `120`,其它显式自定义值保持不变;`bgfilter_circuit_state` 分别上报 `mode=flat``mode=complex`
- 影响范围:`api-server` BgFilter worker、熔断指标与测试、worker 环境模板、生产部署迁移门禁、BgFilter 架构和运维文档;不修改 SpacetimeDB schema、父业务 fallback、计费或失败审计流程。
- 验证方式:覆盖两种模式状态隔离、阈值、成功重置、cooldown 到期、permit 前二次检查和部署默认值迁移;运行 api-server BgFilter 定向测试、生产部署脚本门禁、编码检查与 diff 检查。
- 关联文档:`docs/technical/【后端架构】BgFilter受限资源调度方案-2026-07-21.md``docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md``docs/【开发运维】本地开发验证与生产运维-2026-05-15.md`
## 2026-07-21 BgFilter 首版采用单实例同步内部 HTTP 与父流程原地等待
- 背景:角色动画在单个 `external_generation_job` 内通过 `buffer_unordered(frame_count)` 可并发发射最多 `48` 次 BgFilter 请求;限制父 worker 并发不能限制单个父 job 内的实际 BgFilter 并发。父 job checkpoint / continuation 和 SpacetimeDB 持久子任务都会扩大父状态机、attempt、计费、恢复和清理改动,而当前 BgFilter 成功结果本来就是 HTTP 图片二进制。
File diff suppressed because one or more lines are too long
@@ -1,6 +1,6 @@
# BgFilter 受限资源调度方案(同步内部 HTTP 原地等待版)
更新时间:`2026-07-21`
更新时间:`2026-07-22`
状态:`已实施,待生产压测`
@@ -19,7 +19,7 @@
| 并发 | 进程内 `Semaphore(N)`,并增加有界 admission 上限 `Q` |
| 重试 | 子 worker 对一次逻辑调用最多做两次顺序 provider attempt;父侧不重试整次内部 RPC |
| 超时 | 双预算:父侧派生排队预算 `maxQueueWaitMs` 与调用预算 `callBudgetMs`attempt 上限由 `N × est × 2` 公式运行时派生(est 默认 `5s`),排队不侵蚀调用时间 |
| flat 熔断 | 迁到唯一子 worker 的进程内状态;连续失败达到阈值后暂时跳过 BgFilter,complex 完全不参与 |
| flat / complex 熔断 | 迁到唯一子 worker;两种模式共享阈值和 `120s` cooldown,但分别维护独立进程内状态 |
| 业务语义 | 父流程继续负责 Alpha / 尺寸恢复、flat fallback、最终 OSS、画布写回、计费和父终态 |
| 动画失败 | 首版保持当前“所有已提交帧都等待并排空”语义,不新增跨帧取消组 |
| 崩溃恢复 | 不查询、不恢复 BgFilter 结果;父 job 沿用现有 lease、失败和退款语义 |
@@ -89,14 +89,14 @@ flowchart LR
职责边界:
- 父流程负责源对象已持久化、owner 校验、请求预算、flat fallback、Alpha / 尺寸恢复、动画 finalizer、最终 OSS / `asset_object` / 画布写回、计费和父终态。
- `bgfilter-worker` 负责内部协议校验、OSS 签名、并发与排队上限、BgFilter 协议、两次顺序尝试、flat 熔断、provider 失败审计和结果图片校验。
- `bgfilter-worker` 负责内部协议校验、OSS 签名、并发与排队上限、BgFilter 协议、两次顺序尝试、按模式隔离的熔断、provider 失败审计和结果图片校验。
- SpacetimeDB 不参与本次内部调度;不新增表、reducer、procedure、facade 或生成 bindings。
`bgfilter-worker` 从实现形态看是只监听内部地址的同步 worker service,不是队列 consumer。父 worker 调另一个 worker 在这里是允许的:父进程明确选择保留调用栈和槽位,因此同步内部 HTTP 正是首版的最小交接方式。
这里的“同步等待”是控制流上的 request / response `await`:不会阻塞 OS 执行线程或整个父进程,但父 job future 仍留在通用 worker 的并发集合中,占用一个父 worker 槽,并由现有 heartbeat 继续续租。
首版进程角色仍复用现有完整 `AppState` 构造路径,以获得 OSS、BgFilter provider、SpacetimeDB 审计、HTTP client 和可观测性依赖;进程角色只阻止它挂载公共路由、claim 外部生成 job 或启动其它后台循环,并不等于它只需要几个调度环境变量。因此生产 unit 必须先加载 `/etc/genarrative/api-server.env` 中父子共享的 `N / est` 与基础配置,再加载 `/etc/genarrative/bgfilter-worker.env` 覆盖监听地址、`Q`、flat 熔断和 worker 独占参数。后续若拆出轻量专用 state,可再缩小共享配置依赖,首版不能假设该拆分已经存在。
首版进程角色仍复用现有完整 `AppState` 构造路径,以获得 OSS、BgFilter provider、SpacetimeDB 审计、HTTP client 和可观测性依赖;进程角色只阻止它挂载公共路由、claim 外部生成 job 或启动其它后台循环,并不等于它只需要几个调度环境变量。因此生产 unit 必须先加载 `/etc/genarrative/api-server.env` 中父子共享的 `N / est` 与基础配置,再加载 `/etc/genarrative/bgfilter-worker.env` 覆盖监听地址、`Q`、flat / complex 统一熔断参数和 worker 独占参数。后续若拆出轻量专用 state,可再缩小共享配置依赖,首版不能假设该拆分已经存在。
### 3.1 图片数据流口径
@@ -151,7 +151,7 @@ Authorization: Bearer <internal-token>
- 当前部署只有一个配置内私有 OSS bucket,因此请求只传 `sourceObjectKey`,子 worker 从自身 OSS 配置取 bucket 并生成短期签名 URL。
- 如果未来确实支持多个 bucket,新增字段也必须由服务端 allowlist 校验;不能接受调用方提供任意下载 URL。
- `backgroundMode` 只允许 `flat / complex``segModel` 继续沿用当前 `birefnet / anime-seg` allowlistcomplex 固定使用当前参数组合。
- `screenColor` 只对 flat 必填;complex 不得误接 flat 参数或熔断
- `screenColor` 只对 flat 必填;complex 不得误接 flat 参数,两种模式的熔断状态必须隔离
- `maxQueueWaitMs``callBudgetMs` 都是相对预算,不是跨机器绝对时间。前者从 admission 起约束排队阶段(worker 还会用 §5.2 的动态估计对其取 min);后者从取得 provider permit 起计时,覆盖签名、两次 attempt、结果校验和响应构造。`callBudgetMs` 是父侧按 `N / est` 公式算出的“配置指纹”,仅作核对:worker 始终以自己按同一公式派生的值执行,不一致时不拒绝请求,而是记录 warn 日志并递增漂移指标。发布调优 N / est 时新旧进程共存的瞬态漂移因此不会误伤在途任务;持久性漂移的硬拦截由部署脚本的共享 env 对齐校验承担。
- JSON body 设置很小的固定上限;源图字节不进入该 JSON。
@@ -177,7 +177,7 @@ Content-Type: image/png
失败返回有界 JSON,稳定错误码只保留:
- `provider_exhausted`:两次真实 provider attempt 都失败;
- `circuit_open`flat 熔断已打开,未发送 provider 请求;
- `circuit_open`当前请求模式的熔断已打开,未发送 provider 请求;
- `deadline_exceeded`:排队、provider 或响应阶段预算耗尽,使用 `phase = queue | provider | response``phase = queue` 时附带触发边界 `bound = estimate | parent`,区分动态过载探测与父上限;
- `overloaded`admission 保险丝 `Q` 触达(默认 `2048`,正常业务不应出现);
- `cancelled`:保留错误码,首版子 worker 不产生。首版没有显式取消信号通道,单纯 TCP 断连后 handler future 被 drop、也无法再返回响应;该码为第二阶段 group cancellation 预留,父侧已按“不启动 fallback”实现映射;
@@ -316,18 +316,18 @@ inline / External v1 当前没有显式 `RequestContext` deadline 时,内部 R
### 6.3 熔断
熔断是故障保护:flat 的真实 provider attempt 连续失败达到阈值后,在 cooldown 内暂时不再请求 BgFilter,而是快速返回 `circuit_open`,由父流程进入“阿里云 → 本地”fallback,避免故障 provider 持续占满并发和超时
熔断是故障保护:flat 或 complex 的真实 provider attempt 连续失败达到阈值后,当前模式在 cooldown 内暂时不再请求 BgFilter,而是快速返回 `circuit_open`避免故障 provider 持续占满并发和超时。flat 由父流程继续进入“阿里云 → 本地”fallback;complex 仍直接失败,不获得 flat fallback
保持当前语义:
- 只有 flat 读取和更新熔断;complex 完全不读写
- flat 在取得 permit、即将发送第一次 provider HTTP 前重新检查熔断,避免 48 个排队请求在熔断打开前全部通过旧检查。
- flat / complex 共享同一阈值与 cooldown 配置,但分别维护独立的 `consecutive_failures / open_until`;任一模式的失败或成功只更新自身状态,不影响另一模式
- 两种模式都在取得 permit、即将发送第一次 provider HTTP 前重新检查自身熔断,避免排队请求在熔断打开前全部通过旧检查。
- 已经获准执行的逻辑调用,即使第一次失败使熔断打开,也仍允许在预算内完成自己的第二次顺序 attempt;后续请求快速返回 `circuit_open`
- 每个真实失败 attempt 计一次失败,保持当前计数口径;flat 任一真实 attempt 成功后重置。
- 每个真实失败 attempt 计一次失败,保持当前计数口径;任一真实 attempt 成功后重置当前模式
- 只有拿到完整公式 attempt 上限(`N × est × 2`)后发生的 provider timeout,以及真实传输失败、非 2xx 和无效 / 超限图片计入。因 `callBudgetMs` 剩余不足而被截短的 timeout 返回 `deadline_exceeded`,不更新熔断;保险丝拒绝、排队超时、客户端取消、鉴权和本地配置错误同样不计入。
- 进程重启后熔断状态清零是首版接受行为。
flat 熔断`GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD``GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS` 属于 `bgfilter-worker` 运行参数;父 API / external-generation worker 不再读取或更新熔断。生产示例必须把这两个值放进 worker 专属环境,避免运维人员在父侧修改了一个实际不生效的配置
flat / complex 统一使用`GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD=3``GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS=120` 属于 `bgfilter-worker` 运行参数;父 API / external-generation worker 不再读取或更新熔断。生产示例必须把这两个值放进 worker 专属环境deploy / Provision 只把历史模板默认 cooldown `300` 定向迁移为 `120`,保留其它显式自定义值
首版不增加 QPS 限制。若 provider 以后要求 QPS,需要另加 token bucket;不能把并发 semaphore 当作 QPS。
@@ -364,7 +364,7 @@ flat 熔断的 `GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD` 与 `GENA
| flat | `cancelled`(保留码,首版子 worker 不产生),或父 job cancellation / 绝对 deadline 已生效 | 立即向上退出,不再启动阿里云或本地 fallback |
| flat | `invalid_request / unauthorized` | 作为内部契约或部署配置错误失败,不 fallback、不计入 BgFilter 熔断 |
| complex 手动去背景 | 图片二进制 | 父流程继续最终 OSS、资源和画布写回 |
| complex 手动去背景 | 任意非成功或断连 | 父流程直接失败;不得接 flat fallback,不得修改 flat 熔断 |
| complex 手动去背景 | 任意非成功或断连 | 父流程直接失败;provider 失败只累计 complex 熔断,不得接 flat fallback修改 flat 熔断 |
| 角色 / 图标 / UI 后处理最终失败 | BgFilter 与 fallback 都未得到可用结果 | 保留已持久化 provider 原图,以现有 `completed + warning` 收口 |
| 动画任一帧最终失败 | 该帧完整 fallback / finalizer / PUT 仍失败 | 排空其它已提交帧后,整项动画按现有语义失败退款 |
@@ -411,11 +411,11 @@ BgFilter 成功二进制不是一份新的业务资产:
- `bgfilter_provider_http_seconds{mode,attempt,outcome}`
- `bgfilter_internal_request_total{mode,outcome}`
- `bgfilter_internal_response_bytes`
- `bgfilter_circuit_state`
- `bgfilter_circuit_state{mode=flat|complex}`
日志只写 `requestId`、父 job / request correlation、mode、attempt、排队耗时、provider 耗时、结果码和安全 object key;不得记录请求/响应图片 body。
当前 flat 的每次 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 进程角色不共享 api / extgen 的落盘 tracking outbox,失败审计由异步任务直写 SpacetimeDB,并纳入 shutdown tracker,优雅退出前排空;进程被强杀时可能丢失,属首版接受行为。
## 10. 实施与部署计划
@@ -423,7 +423,7 @@ BgFilter 成功二进制不是一份新的业务资产:
1. 在现有 Rust 后端增加 `bgfilter-worker` 进程角色和独立 loopback Axum listener;它不启动用户 HTTP router,也不 claim `external_generation_job`
2. 增加内部 request / binary response / typed error 契约、Token 校验、JSON body 上限、object key allowlist 和健康检查;listener 在 body 解析前接入连接 / request concurrency limit、固定 backlog 和 load shedding。
3. 增加 admission `Q``Semaphore(N)`、两次顺序 attempt、预算检查、结果限长 / 解码校验、response-body permit guard 和 flat 进程级熔断。
3. 增加 admission `Q``Semaphore(N)`、两次顺序 attempt、预算检查、结果限长 / 解码校验、response-body permit guard 和 flat / complex 独立进程级熔断。
4. 增加父侧共享内部 HTTP client。该 client 不自动重试;对 `2xx` 读取并返回受限图片字节,对非 `2xx` 只解析有界类型化 JSON 错误;把父剩余预算显式转换为 `maxQueueWaitMs` 与公式 `callBudgetMs`client timeout 固定取两者之和加 `2s`。父绝对预算只在派生 `maxQueueWaitMs` 时扣除 callBudget 与父侧预留,不在发送阶段重新裁剪或挪用两笔相对预算。
5. 用内部 client 替换两个集中调用边界:
- flat`remove_editor_generated_screen_background_with_bgfilter`
@@ -467,8 +467,8 @@ BgFilter 成功二进制不是一份新的业务资产:
- `callBudgetMs` 与 worker 本进程公式值不一致时不拒绝:worker 以自身公式值执行,记 warn 并递增漂移指标;发布调优 N / est 的新旧进程共存窗口内,在途 flat 任务仍能正常执行或走既有 fallback,不得因瞬态漂移触发 `invalid_request`(该码禁止 fallback)。attempt、callBudget、client timeout 全部由 `N / est` 运行时派生,代码不存在硬编码结果值。
- parent client timeout 精确取 `maxQueueWaitMs + callBudgetMs + 2s`helper 保持 infallible;父绝对预算通过 `maxQueueWaitMs` 的派生公式预先约束,结果校验等待也必须 deadline-aware,不能只在校验完成后事后判超时。
- 第一次失败后预算不足时不开始第二次;父侧从不重试整次内部 RPC。
- 父业务预算仍有效时,flat 两次失败、熔断、overload、内部 RPC deadline 或断连仍走“阿里云 → 本地”;complex 任意失败直接失败且不读写熔断
- flat 熔断按真实失败 attempt 计数;由剩余业务预算截短的 timeout 不计入。已获准调用可完成第二次,后续排队请求快速 `circuit_open`
- 父业务预算仍有效时,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 排空后进程才退出。
- 客户端断连时,等待 permit 的请求最终由 deadline 收口;已开始 provider attempt 持有 permit 并排空。明确 cancellation 已被观察到后不再开始第二次,单纯 TCP 断连只作 best-effort 测试,不作为硬保证。
@@ -480,7 +480,7 @@ BgFilter 成功二进制不是一份新的业务资产:
- worker 重启 / RPC 丢失不查询、不恢复结果;父 job 的 heartbeat、lease、失败退款和 fencing 保持现状。
- External v1 / inline 不再直连 BgFilter;公共 router、BFF、账单和任务列表中没有内部 endpoint 或内部调用记录。
- 生产不存在两个同时运行的 `bgfilter-worker``N``est` 缺失、为 `0` 时 fail-closed`Q` 显式配置时必须 `>= N`
- worker unit 先加载共享 API env、再加载 worker 专属 env;父子有效 `N``est` 完全一致(两者都在共享 API env),flat 熔断参数只由子 worker 配置和执行。
- worker unit 先加载共享 API env、再加载 worker 专属 env;父子有效 `N``est` 完全一致(两者都在共享 API env),flat / complex 统一熔断参数只由子 worker 配置和执行cooldown 默认 `120s`
- `external-generation-worker.env` 后加载时不得把内部 base URL、Token / Token 文件、connect timeout、`N``est`、OSS bucket 或 endpoint 覆盖为与共享 API env 不同的有效值;父侧必须把源对象写到子 worker 将要签名读取的同一 OSS 位置。外部生成 worker 可使用同 bucket 下权限等价或更小的独立 AK,不要求凭据文本相同。
- 父进程 `GENARRATIVE_BGFILTER_WORKER_BASE_URL`、子 worker `HOST / PORT` 和部署 readiness URL 必须指向同一个 `127.0.0.1:<port>` endpoint;旧非空配置不能因为“无需补默认值”而绕过一致性检查。
- `genarrative-bgfilter-worker.service` 必须保持 `TimeoutStopSec=900`,覆盖 `callBudgetMs` 排空上界(约 `321s`)与停止收口余量;停止时排队请求立即类型化失败,不参与排空。
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
+45 -1
View File
@@ -57,6 +57,7 @@ function main() {
assertDeployRejectsNonLoopbackBgFilterListener();
assertDeployRejectsMissingBgFilterWorkerCapacity();
assertDeployMigratesOldDefaultBgFilterAdmissionLimit();
assertDeployMigratesOldDefaultBgFilterCircuitCooldown();
assertDeployRejectsInvalidBgFilterWorkerCapacity();
assertDeployRejectsInlineBgFilterInternalToken();
assertDeployRejectsEmptyBgFilterInternalToken();
@@ -395,7 +396,12 @@ function assertDeployCopiesPingoraDirectReleaseDependencies() {
assertIncludes(
bgfilterEnvExample,
'GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD=3',
'BgFilter 专属 env 必须提供 flat 熔断阈值。',
'BgFilter 专属 env 必须提供 flat / complex 统一熔断阈值。',
);
assertIncludes(
bgfilterEnvExample,
'GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS=120',
'BgFilter 专属 env 必须提供 flat / complex 统一的 120 秒熔断冷却时间。',
);
for (const sharedKey of [
'GENARRATIVE_BGFILTER_WORKER_CONCURRENCY=',
@@ -1292,6 +1298,44 @@ function assertDeployMigratesOldDefaultBgFilterAdmissionLimit() {
}
}
function assertDeployMigratesOldDefaultBgFilterCircuitCooldown() {
const fixture = prepareFixture('bgfilter-worker-migrate-old-default-circuit-cooldown');
writeFileSync(
fixture.bgfilterWorkerEnvFile,
`${readFileSync(fixture.bgfilterWorkerEnvFile, 'utf8')}GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS=300\n`,
'utf8',
);
const result = runDeploy(fixture);
if (result.status !== 0) {
failures.push(`历史默认熔断 cooldown=300 迁移到 120 时部署不应失败:${result.stderr}`);
return;
}
const migrated = readFileSync(fixture.bgfilterWorkerEnvFile, 'utf8');
if (!/^GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS=120$/mu.test(migrated)) {
failures.push('部署必须把历史模板默认熔断 cooldown=300 定向迁移为 120。');
}
if (/^GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS=300$/mu.test(migrated)) {
failures.push('部署完成后不得继续保留历史模板默认熔断 cooldown=300。');
}
const custom = prepareFixture('bgfilter-worker-preserve-custom-circuit-cooldown');
writeFileSync(
custom.bgfilterWorkerEnvFile,
`${readFileSync(custom.bgfilterWorkerEnvFile, 'utf8')}GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS=90\n`,
'utf8',
);
const customResult = runDeploy(custom);
if (customResult.status !== 0) {
failures.push(`自定义熔断 cooldown=90 时部署不应失败:${customResult.stderr}`);
return;
}
const preserved = readFileSync(custom.bgfilterWorkerEnvFile, 'utf8');
if (!/^GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS=90$/mu.test(preserved)) {
failures.push('部署只能迁移历史默认熔断 cooldown=300,必须保留其它显式自定义值。');
}
}
function assertMissingReleaseManifestFails() {
const fixture = prepareFixture('missing-release-manifest');
rmSync(path.join(fixture.sourceDir, 'release-manifest.json'));
@@ -1691,6 +1691,13 @@ const checks = [
reason:
'Server-Provision 迁移历史默认值时必须读取最后一次有效赋值,不能覆盖后写的自定义运行态值。',
},
{
file: 'scripts/jenkins-server-provision.sh',
includes:
'ensure_env_value_migrates_old_default "${BGFILTER_WORKER_ENV_FILE}" "GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS" "300" "120"',
reason:
'Server-Provision 必须把 BgFilter 熔断 cooldown 历史模板默认 300 定向迁移为 120,并保留其它显式定制值。',
},
{
file: 'scripts/jenkins-server-provision.sh',
includes: "root:genarrative:440",
+1 -1
View File
@@ -497,7 +497,7 @@ ensure_bgfilter_worker_runtime_env_defaults() {
# 其它显式定制值继续保留。
ensure_env_value_migrates_old_default "${bgfilter_env_file}" "GENARRATIVE_BGFILTER_WORKER_MAX_REQUESTS" "128" "2048"
ensure_env_value "${bgfilter_env_file}" "GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD" "3"
ensure_env_value "${bgfilter_env_file}" "GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS" "300"
ensure_env_value_migrates_old_default "${bgfilter_env_file}" "GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS" "300" "120"
# N 已迁入共享 API env;worker 专属文件中的旧值会与共享值形成双写风险,直接移除。
remove_env_key_if_present "${bgfilter_env_file}" "GENARRATIVE_BGFILTER_WORKER_CONCURRENCY"
}
+1 -1
View File
@@ -806,7 +806,7 @@ ensure_bgfilter_worker_runtime_env_defaults() {
# 其它显式定制值继续保留。
ensure_env_value_migrates_old_default "${BGFILTER_WORKER_ENV_FILE}" "GENARRATIVE_BGFILTER_WORKER_MAX_REQUESTS" "128" "2048"
ensure_env_value "${BGFILTER_WORKER_ENV_FILE}" "GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD" "3"
ensure_env_value "${BGFILTER_WORKER_ENV_FILE}" "GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS" "300"
ensure_env_value_migrates_old_default "${BGFILTER_WORKER_ENV_FILE}" "GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS" "300" "120"
}
bgfilter_internal_token_file_is_single_segment() {
@@ -53,6 +53,7 @@ const BGFILTER_SOURCE_URL_EXPIRE_SECONDS: u64 = 600;
const BGFILTER_PROVIDER_TOKEN_HEADER: &str = "X-Genarrative-Image-Token";
static BGFILTER_FLAT_CIRCUIT: OnceLock<Mutex<BgfilterCircuitState>> = OnceLock::new();
static BGFILTER_COMPLEX_CIRCUIT: OnceLock<Mutex<BgfilterCircuitState>> = OnceLock::new();
struct BgfilterMetrics {
_circuit_state: ObservableGauge<i64>,
@@ -72,14 +73,22 @@ fn bgfilter_metrics() -> &'static BgfilterMetrics {
let meter = global::meter("genarrative-bgfilter-worker");
let circuit_state = meter
.i64_observable_gauge("bgfilter_circuit_state")
.with_description("Flat BgFilter circuit state: 0 closed, 1 open")
.with_description("BgFilter circuit state by mode: 0 closed, 1 open")
.with_callback(|observer| {
let open = flat_circuit_state()
.lock()
.ok()
.and_then(|circuit| circuit.open_until)
.is_some_and(|open_until| open_until > Instant::now());
observer.observe(i64::from(open), &[KeyValue::new("mode", "flat")]);
for mode in [
BgfilterBackgroundMode::Flat,
BgfilterBackgroundMode::Complex,
] {
let open = circuit_state(mode)
.lock()
.ok()
.and_then(|circuit| circuit.open_until)
.is_some_and(|open_until| open_until > Instant::now());
observer.observe(
i64::from(open),
&[KeyValue::new("mode", mode.as_str())],
);
}
})
.build();
BgfilterMetrics {
@@ -144,10 +153,6 @@ impl BgfilterBackgroundMode {
Self::Complex => "complex",
}
}
fn uses_flat_circuit(self) -> bool {
matches!(self, Self::Flat)
}
}
#[derive(Clone, Debug)]
@@ -966,20 +971,19 @@ async fn execute_logical_request(
request: BgfilterInternalRequest,
admitted_at: Instant,
) -> WorkerOutcome {
if request.background_mode.uses_flat_circuit() {
if let Some(remaining) = flat_circuit_open_remaining(&runtime.app_state) {
return WorkerOutcome {
result: Err(WorkerFailure::new(
"circuit_open",
format!(
"BgFilter flat 熔断仍有 {}ms",
remaining.as_millis().min(u128::from(u64::MAX))
),
true,
)),
provider_permit: None,
};
}
if let Some(remaining) = circuit_open_remaining(&runtime.app_state, request.background_mode) {
return WorkerOutcome {
result: Err(WorkerFailure::new(
"circuit_open",
format!(
"BgFilter {} 熔断仍有 {}ms",
request.background_mode.as_str(),
remaining.as_millis().min(u128::from(u64::MAX))
),
true,
)),
provider_permit: None,
};
}
let (queue_len_ahead, queue_depth_guard) = QueueDepthGuard::enter(runtime.queue_depth.clone());
@@ -1056,20 +1060,19 @@ async fn execute_logical_request(
// 中文注释:排队期间熔断可能刚被其它请求打开。必须在取得 provider permit 后、
// 第一次真实 HTTP 前再检查一次;本调用一旦通过该检查,自己的第二次 attempt 不再重查。
if request.background_mode.uses_flat_circuit() {
if let Some(remaining) = flat_circuit_open_remaining(&runtime.app_state) {
return WorkerOutcome {
result: Err(WorkerFailure::new(
"circuit_open",
format!(
"BgFilter flat 熔断仍有 {}ms",
remaining.as_millis().min(u128::from(u64::MAX))
),
true,
)),
provider_permit: None,
};
}
if let Some(remaining) = circuit_open_remaining(&runtime.app_state, request.background_mode) {
return WorkerOutcome {
result: Err(WorkerFailure::new(
"circuit_open",
format!(
"BgFilter {} 熔断仍有 {}ms",
request.background_mode.as_str(),
remaining.as_millis().min(u128::from(u64::MAX))
),
true,
)),
provider_permit: None,
};
}
let attempt_state = &runtime.app_state;
@@ -1127,8 +1130,8 @@ async fn execute_logical_request(
attempt_started,
error.clone(),
);
if request.background_mode.uses_flat_circuit() && error.counts_for_circuit {
record_flat_failure(&runtime.app_state);
if error.counts_for_circuit {
record_circuit_failure(&runtime.app_state, request.background_mode);
}
},
)
@@ -1136,9 +1139,7 @@ async fn execute_logical_request(
match attempts {
SequentialAttemptOutcome::Success { value, .. } => {
if request.background_mode.uses_flat_circuit() {
record_flat_success();
}
record_circuit_success(request.background_mode);
WorkerOutcome {
result: Ok(value),
provider_permit: Some(provider_permit),
@@ -1889,41 +1890,66 @@ struct BgfilterCircuitState {
open_until: Option<Instant>,
}
fn flat_circuit_state() -> &'static Mutex<BgfilterCircuitState> {
BGFILTER_FLAT_CIRCUIT.get_or_init(|| Mutex::new(BgfilterCircuitState::default()))
impl BgfilterCircuitState {
fn open_remaining(&mut self, now: Instant) -> Option<Duration> {
let open_until = self.open_until?;
if open_until > now {
return Some(open_until.duration_since(now));
}
*self = Self::default();
None
}
fn record_success(&mut self) {
*self = Self::default();
}
fn record_failure(&mut self, threshold: u32, cooldown: Duration, now: Instant) {
self.consecutive_failures = self.consecutive_failures.saturating_add(1);
if self.consecutive_failures >= threshold {
self.open_until = Some(now + cooldown);
}
}
}
fn flat_circuit_open_remaining(state: &AppState) -> Option<Duration> {
fn circuit_state(mode: BgfilterBackgroundMode) -> &'static Mutex<BgfilterCircuitState> {
match mode {
BgfilterBackgroundMode::Flat => {
BGFILTER_FLAT_CIRCUIT.get_or_init(|| Mutex::new(BgfilterCircuitState::default()))
}
BgfilterBackgroundMode::Complex => {
BGFILTER_COMPLEX_CIRCUIT.get_or_init(|| Mutex::new(BgfilterCircuitState::default()))
}
}
}
fn circuit_open_remaining(state: &AppState, mode: BgfilterBackgroundMode) -> Option<Duration> {
if state.config.editor_bgfilter_circuit_failure_threshold == 0 {
return None;
}
let mut circuit = flat_circuit_state().lock().ok()?;
let open_until = circuit.open_until?;
let now = Instant::now();
if open_until > now {
return Some(open_until.duration_since(now));
}
*circuit = BgfilterCircuitState::default();
None
circuit_state(mode)
.lock()
.ok()?
.open_remaining(Instant::now())
}
fn record_flat_success() {
if let Ok(mut circuit) = flat_circuit_state().lock() {
*circuit = BgfilterCircuitState::default();
fn record_circuit_success(mode: BgfilterBackgroundMode) {
if let Ok(mut circuit) = circuit_state(mode).lock() {
circuit.record_success();
}
}
fn record_flat_failure(state: &AppState) {
fn record_circuit_failure(state: &AppState, mode: BgfilterBackgroundMode) {
let threshold = state.config.editor_bgfilter_circuit_failure_threshold;
if threshold == 0 {
return;
}
if let Ok(mut circuit) = flat_circuit_state().lock() {
circuit.consecutive_failures = circuit.consecutive_failures.saturating_add(1);
if circuit.consecutive_failures >= threshold {
circuit.open_until =
Some(Instant::now() + state.config.editor_bgfilter_circuit_cooldown);
}
if let Ok(mut circuit) = circuit_state(mode).lock() {
circuit.record_failure(
threshold,
state.config.editor_bgfilter_circuit_cooldown,
Instant::now(),
);
}
}
@@ -2862,9 +2888,34 @@ mod tests {
}
#[test]
fn complex_mode_never_uses_flat_circuit_and_internal_errors_do_not_retry() {
assert!(BgfilterBackgroundMode::Flat.uses_flat_circuit());
assert!(!BgfilterBackgroundMode::Complex.uses_flat_circuit());
fn flat_and_complex_circuits_are_independent_and_internal_errors_do_not_retry() {
assert!(!std::ptr::eq(
circuit_state(BgfilterBackgroundMode::Flat),
circuit_state(BgfilterBackgroundMode::Complex),
));
let now = Instant::now();
let cooldown = Duration::from_secs(120);
let mut flat = BgfilterCircuitState::default();
let mut complex = BgfilterCircuitState::default();
for _ in 0..2 {
flat.record_failure(3, cooldown, now);
}
assert_eq!(flat.open_remaining(now), None);
flat.record_failure(3, cooldown, now);
assert_eq!(flat.open_remaining(now), Some(cooldown));
assert_eq!(complex.open_remaining(now), None);
for _ in 0..3 {
complex.record_failure(3, cooldown, now);
}
assert_eq!(complex.open_remaining(now), Some(cooldown));
complex.record_success();
assert_eq!(complex.open_remaining(now), None);
assert_eq!(flat.open_remaining(now), Some(cooldown));
assert_eq!(flat.open_remaining(now + cooldown), None);
assert_eq!(flat.consecutive_failures, 0);
assert!(!ProviderAttemptError::internal("config error").should_retry());
assert_eq!(BGFILTER_PROVIDER_MAX_ATTEMPTS, 2);
}
@@ -3051,7 +3102,7 @@ mod tests {
}
#[test]
fn worker_rechecks_flat_circuit_and_signs_each_attempt_while_parent_sends_once() {
fn worker_rechecks_mode_circuit_and_signs_each_attempt_while_parent_sends_once() {
let source = include_str!("bgfilter_worker.rs");
let execute = source
.split_once("async fn execute_logical_request")
@@ -3061,7 +3112,7 @@ mod tests {
.expect("attempt budget boundary")
.0;
let circuit_checks = execute
.match_indices("flat_circuit_open_remaining")
.match_indices("circuit_open_remaining")
.map(|(index, _)| index)
.collect::<Vec<_>>();
assert_eq!(circuit_checks.len(), 2);
+2 -1
View File
@@ -22,7 +22,7 @@ const DEFAULT_EDITOR_BGFILTER_SINGLE_IMAGE_ESTIMATE_MS: u64 = 5_000;
const BGFILTER_ATTEMPT_SAFETY_FACTOR: u64 = 2;
const BGFILTER_WORKER_RESPONSE_WINDOW_MS: u64 = 1_000;
const DEFAULT_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD: u32 = 3;
const DEFAULT_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS: u64 = 300;
const DEFAULT_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS: u64 = 120;
const DEFAULT_ALIYUN_MATTING_ENDPOINT: &str = "imageseg.cn-shanghai.aliyuncs.com";
const DEFAULT_ALIYUN_MATTING_REQUEST_TIMEOUT_MS: u64 = 30_000;
@@ -1642,6 +1642,7 @@ mod tests {
config.editor_bgfilter_circuit_cooldown.as_secs(),
DEFAULT_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS
);
assert_eq!(DEFAULT_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS, 120);
assert!(config.editor_bgfilter_token.is_none());
assert_eq!(
config.external_generation_worker_lease.as_secs(),