限制BgFilter失败审计堆积
增加进程级1024审计任务硬上限并在spawn前拒绝超额任务。 启用BgFilter-worker独立tracking-outbox并接入启动与退出flush生命周期。 增加outbox-only审计策略、丢弃指标及相关单元和故障验证。 同步后端架构、开发运维和项目决策文档。
This commit is contained in:
@@ -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,且五分钟冷却会让短暂故障恢复后的等待过长。
|
||||
|
||||
@@ -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` 的图片。
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -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/<version>`,`genarrative` 用户无法在其中创建 `server-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<u64>,
|
||||
connect_retry_total: Counter<u64>,
|
||||
in_flight: UpDownCounter<i64>,
|
||||
audit_in_flight: UpDownCounter<i64>,
|
||||
audit_dropped_total: Counter<u64>,
|
||||
internal_request_seconds: Histogram<f64>,
|
||||
provider_http_seconds: Histogram<f64>,
|
||||
internal_request_total: Counter<u64>,
|
||||
@@ -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<Semaphore>,
|
||||
provider: Arc<Semaphore>,
|
||||
audit: Arc<Semaphore>,
|
||||
/// 正在等待 provider permit 的请求数;进入 provider permit 等待队列时的快照即
|
||||
/// 「排在我前面的队长」,FIFO 语义下后来者不影响先到者的等待,快照可直接用于动态排队估时。
|
||||
queue_depth: Arc<AtomicUsize>,
|
||||
@@ -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<Semaphore>) -> Result<Self, TryAcquireError> {
|
||||
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<Semaphore>,
|
||||
attempt_started: bool,
|
||||
error: &ProviderAttemptError,
|
||||
) -> Option<BgfilterAuditPermit> {
|
||||
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<Semaphore>,
|
||||
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<Semaphore>,
|
||||
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();
|
||||
|
||||
@@ -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<u16>,
|
||||
timeout: bool,
|
||||
transport: bool,
|
||||
latency_ms: Option<u64>,
|
||||
error_message: String,
|
||||
raw_excerpt: Option<String>,
|
||||
) {
|
||||
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(
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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<Body>) -> Response<Body> {
|
||||
response.map(|body| {
|
||||
HTTP_RESPONSE_BODY_IN_FLIGHT.fetch_add(1, Ordering::Relaxed);
|
||||
@@ -219,6 +229,7 @@ struct TrackingOutboxMetrics {
|
||||
|
||||
struct ExternalApiMetrics {
|
||||
failures: Counter<u64>,
|
||||
audit_dropped: Counter<u64>,
|
||||
}
|
||||
|
||||
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(),
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user