Files
Genarrative/server-rs/crates/api-server/src/external_api_audit.rs
T

769 lines
27 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#[cfg(test)]
use axum::http::StatusCode;
use module_runtime::RuntimeTrackingScopeKind;
use platform_image::PlatformImageFailureAudit;
use serde_json::{Value, json};
use time::OffsetDateTime;
use uuid::Uuid;
use crate::{http_error::AppError, state::AppState, tracking::TrackingEventDraft};
pub(crate) const EXTERNAL_API_FAILURE_EVENT_KEY: &str = "external_api_call_failure";
pub(crate) const EXTERNAL_API_AUDIT_MODULE_KEY: &str = "external-api";
#[derive(Clone, Debug)]
pub(crate) struct ExternalApiFailureDraft {
pub(crate) provider: &'static str,
pub(crate) endpoint: String,
pub(crate) operation: String,
pub(crate) failure_stage: &'static str,
pub(crate) status_code: Option<u16>,
pub(crate) status_class: Option<&'static str>,
pub(crate) timeout: bool,
pub(crate) retryable: bool,
pub(crate) error_message: String,
pub(crate) error_source: Option<String>,
pub(crate) raw_excerpt: Option<String>,
pub(crate) latency_ms: Option<u64>,
pub(crate) prompt_chars: Option<usize>,
pub(crate) reference_image_count: Option<usize>,
pub(crate) image_model: Option<&'static str>,
pub(crate) user_id: Option<String>,
pub(crate) profile_id: Option<String>,
pub(crate) request_id: Option<String>,
}
impl ExternalApiFailureDraft {
pub(crate) fn new(
provider: &'static str,
endpoint: impl Into<String>,
operation: impl Into<String>,
failure_stage: &'static str,
error_message: impl Into<String>,
) -> Self {
Self {
provider,
endpoint: endpoint.into(),
operation: operation.into(),
failure_stage,
status_code: None,
status_class: None,
timeout: false,
retryable: false,
error_message: error_message.into(),
error_source: None,
raw_excerpt: None,
latency_ms: None,
prompt_chars: None,
reference_image_count: None,
image_model: None,
user_id: None,
profile_id: None,
request_id: None,
}
}
pub(crate) fn with_status_code(mut self, status_code: Option<u16>) -> Self {
self.status_code = status_code;
self
}
pub(crate) fn with_optional_status_class(mut self, status_class: Option<&'static str>) -> Self {
self.status_class = status_class;
self
}
pub(crate) fn with_timeout(mut self, timeout: bool) -> Self {
self.timeout = timeout;
self
}
pub(crate) fn with_retryable(mut self, retryable: bool) -> Self {
self.retryable = retryable;
self
}
pub(crate) fn with_error_source(mut self, error_source: Option<String>) -> Self {
self.error_source = error_source;
self
}
pub(crate) fn with_raw_excerpt(mut self, raw_excerpt: Option<String>) -> Self {
self.raw_excerpt = raw_excerpt;
self
}
pub(crate) fn with_latency_ms(mut self, latency_ms: Option<u64>) -> Self {
self.latency_ms = latency_ms;
self
}
pub(crate) fn with_prompt_chars(mut self, prompt_chars: Option<usize>) -> Self {
self.prompt_chars = prompt_chars;
self
}
pub(crate) fn with_reference_image_count(
mut self,
reference_image_count: Option<usize>,
) -> Self {
self.reference_image_count = reference_image_count;
self
}
pub(crate) fn with_image_model(mut self, image_model: Option<&'static str>) -> Self {
self.image_model = image_model;
self
}
pub(crate) fn with_user_id(mut self, user_id: Option<String>) -> Self {
self.user_id = user_id;
self
}
pub(crate) fn with_profile_id(mut self, profile_id: Option<String>) -> Self {
self.profile_id = profile_id;
self
}
pub(crate) fn with_request_id(mut self, request_id: Option<String>) -> Self {
self.request_id = request_id;
self
}
pub(crate) fn with_audit_context(mut self, context: &ExternalApiAuditContext) -> Self {
self.user_id = context.user_id.clone();
self.profile_id = context.profile_id.clone();
self.request_id = context.request_id.clone();
self
}
}
/// 外部 API 失败审计的调用方上下文(用户 / 档案 / 请求 id)。
#[derive(Clone, Debug, Default)]
pub(crate) struct ExternalApiAuditContext {
pub(crate) user_id: Option<String>,
pub(crate) profile_id: Option<String>,
pub(crate) request_id: Option<String>,
}
/// 抠图供应商(BgFilter / 阿里云通用抠图)调用失败的统一失败审计入口。
/// 即使随后兜底成功,供应商故障也必须进入 OTLP + tracking_event,不能只 warn! 后静默。
#[allow(clippy::too_many_arguments)]
pub(crate) async fn record_matting_external_api_failure(
state: &AppState,
context: &ExternalApiAuditContext,
provider: &'static str,
endpoint: String,
operation: &'static str,
failure_stage: &'static str,
status_code: Option<u16>,
timeout: 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,
latency_ms,
error_message,
raw_excerpt,
context,
);
record_external_api_failure(state, draft).await;
}
#[allow(clippy::too_many_arguments)]
fn build_matting_external_api_failure_draft(
provider: &'static str,
endpoint: String,
operation: &'static str,
failure_stage: &'static str,
status_code: Option<u16>,
timeout: bool,
latency_ms: Option<u64>,
error_message: String,
raw_excerpt: Option<String>,
context: &ExternalApiAuditContext,
) -> ExternalApiFailureDraft {
// status_code=None ⟺ statusClass=transport(见 status_class):DNS / 连接重置 / 读体中断 / 超时
// 这类传输层失败没有上游 HTTP 状态。它们必须与 "transport failures actionable" 语义一致,
// 记为 retryable=true,否则 statusClass=transport 却 retryable=false 会误导告警 / 重试分析。
let has_transport_error = status_code.is_none();
ExternalApiFailureDraft::new(provider, endpoint, operation, failure_stage, error_message)
.with_status_code(status_code)
.with_optional_status_class(Some(status_class(status_code)))
.with_timeout(timeout)
.with_retryable(is_retryable_external_api_failure(
status_code,
timeout,
has_transport_error,
))
.with_latency_ms(latency_ms)
.with_raw_excerpt(raw_excerpt)
.with_audit_context(context)
}
pub(crate) fn matting_failure_audit_status_code(error: &AppError) -> Option<u16> {
error
.details()
.and_then(|details| details.get("upstreamStatus"))
.and_then(Value::as_u64)
.and_then(|value| u16::try_from(value).ok())
.or_else(|| {
if matting_failure_audit_is_transport(error) {
None
} else {
Some(error.status_code().as_u16())
}
})
}
fn matting_failure_audit_is_transport(error: &AppError) -> bool {
matting_failure_audit_timeout(error)
|| error
.details()
.and_then(|details| details.get("transport"))
.and_then(Value::as_bool)
.unwrap_or(false)
}
pub(crate) fn matting_failure_audit_timeout(error: &AppError) -> bool {
error
.details()
.and_then(|details| details.get("timeout"))
.and_then(Value::as_bool)
.unwrap_or(false)
}
pub(crate) fn matting_failure_audit_latency_ms(error: &AppError) -> Option<u64> {
error
.details()
.and_then(|details| details.get("latencyMs"))
.and_then(Value::as_u64)
}
pub(crate) fn matting_failure_external_call_attempted(error: &AppError) -> bool {
error
.details()
.and_then(|details| details.get("externalCallAttempted"))
.and_then(Value::as_bool)
.unwrap_or(true)
}
pub(crate) fn matting_failure_audit_raw_excerpt(error: &AppError) -> Option<String> {
error
.details()
.and_then(|details| {
details
.get("upstreamMessage")
.or_else(|| details.get("rawExcerpt"))
})
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.map(|value| value.chars().take(800).collect())
}
pub(crate) fn build_external_api_failure_draft_from_platform_image_audit(
audit: &PlatformImageFailureAudit,
) -> ExternalApiFailureDraft {
ExternalApiFailureDraft::new(
audit.provider,
audit.endpoint.clone(),
audit.operation.clone(),
audit.failure_stage,
audit.error_message.clone(),
)
.with_status_code(audit.status_code)
.with_optional_status_class(audit.status_class)
.with_timeout(audit.timeout)
.with_retryable(audit.retryable)
.with_error_source(audit.error_source.clone())
.with_raw_excerpt(audit.raw_excerpt.clone())
.with_latency_ms(audit.latency_ms)
.with_prompt_chars(audit.prompt_chars)
.with_reference_image_count(audit.reference_image_count)
.with_image_model(audit.image_model)
.with_user_id(None)
.with_profile_id(None)
.with_request_id(None)
}
/// 中文注释:下载图片、OSS 读写等非标准 HTTP 状态统一显式归类,避免 OTLP 低基数 label 误落到 `transport`。
#[cfg(test)]
pub(crate) fn app_error_status_class(status_code: StatusCode) -> &'static str {
status_class(Some(status_code.as_u16()))
}
/// 中文注释:外部供应商失败同时进入 OTLP 和 tracking_event;失败审计不能反向阻断主业务错误返回。
pub(crate) async fn record_external_api_failure(state: &AppState, draft: ExternalApiFailureDraft) {
record_external_api_failure_otlp(&draft);
let tracking_event = build_external_api_failure_tracking_draft(&draft);
if let Some(outbox) = state.tracking_outbox() {
match outbox
.enqueue(crate::tracking::build_tracking_event_input(
tracking_event.clone(),
))
.await
{
Ok(crate::tracking_outbox::TrackingOutboxEnqueueOutcome::Enqueued) => {}
Ok(crate::tracking_outbox::TrackingOutboxEnqueueOutcome::Dropped { reason }) => {
tracing::warn!(
provider = draft.provider,
endpoint = %draft.endpoint,
operation = %draft.operation,
failure_stage = draft.failure_stage,
reason,
"外部 API 失败审计写入 outbox 被保护阈值拒绝,回退同步直写 SpacetimeDB"
);
crate::tracking::record_tracking_event_after_success(
state,
&audit_request_context(),
tracking_event,
)
.await;
}
Err(error) => {
tracing::warn!(
provider = draft.provider,
endpoint = %draft.endpoint,
operation = %draft.operation,
failure_stage = draft.failure_stage,
error = %error,
"外部 API 失败审计写入 outbox 失败,回退同步直写 SpacetimeDB"
);
crate::tracking::record_tracking_event_after_success(
state,
&audit_request_context(),
tracking_event,
)
.await;
}
}
return;
}
crate::tracking::record_tracking_event_after_success(
state,
&audit_request_context(),
tracking_event,
)
.await;
}
pub(crate) fn build_external_api_failure_tracking_draft(
failure: &ExternalApiFailureDraft,
) -> TrackingEventDraft {
let mut draft = TrackingEventDraft::new(
EXTERNAL_API_FAILURE_EVENT_KEY,
EXTERNAL_API_AUDIT_MODULE_KEY,
);
draft.scope_kind = RuntimeTrackingScopeKind::Module;
draft.scope_id = failure.provider.to_string();
draft.user_id = failure.user_id.clone();
draft.owner_user_id = failure.user_id.clone();
draft.profile_id = failure.profile_id.clone();
draft.metadata = build_external_api_failure_metadata(failure);
draft
}
fn build_external_api_failure_metadata(failure: &ExternalApiFailureDraft) -> Value {
let mut metadata = json!({
"provider": failure.provider,
"endpoint": failure.endpoint,
"operation": failure.operation,
"failureStage": failure.failure_stage,
"statusCode": failure.status_code,
"statusClass": failure.status_class.unwrap_or_else(|| status_class(failure.status_code)),
"timeout": failure.timeout,
"retryable": failure.retryable,
"errorMessage": truncate_field(failure.error_message.as_str(), 1_000),
"occurredAt": current_utc_iso_text(),
});
if let Some(latency_ms) = failure.latency_ms {
metadata["latencyMs"] = json!(latency_ms);
}
if let Some(prompt_chars) = failure.prompt_chars {
metadata["promptChars"] = json!(prompt_chars);
}
if let Some(reference_image_count) = failure.reference_image_count {
metadata["referenceImageCount"] = json!(reference_image_count);
}
if let Some(image_model) = failure.image_model {
metadata["imageModel"] = json!(image_model);
}
if let Some(user_id) = failure
.user_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
{
metadata["userId"] = json!(truncate_field(user_id, 1_000));
}
if let Some(profile_id) = failure
.profile_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
{
metadata["profileId"] = json!(truncate_field(profile_id, 1_000));
}
if let Some(request_id) = failure
.request_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
{
metadata["requestId"] = json!(truncate_field(request_id, 1_000));
}
if let Some(source) = failure
.error_source
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
{
metadata["errorSource"] = json!(truncate_field(source, 1_000));
}
if let Some(excerpt) = failure
.raw_excerpt
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
{
metadata["rawExcerpt"] = json!(truncate_field(excerpt, 800));
}
metadata
}
pub(crate) fn is_retryable_external_api_failure(
status_code: Option<u16>,
timeout: bool,
connect: bool,
) -> bool {
// 429 Too Many Requests / 408 Request Timeout / 5xx 视为可重试。
timeout
|| connect
|| status_code.is_some_and(|status| status == 429 || status == 408 || status >= 500)
}
fn record_external_api_failure_otlp(failure: &ExternalApiFailureDraft) {
crate::telemetry::record_external_api_failure(
failure.provider,
failure.failure_stage,
failure
.status_class
.unwrap_or_else(|| status_class(failure.status_code)),
failure.retryable,
);
tracing::error!(
provider = failure.provider,
endpoint = %failure.endpoint,
operation = %failure.operation,
failure_stage = failure.failure_stage,
status_code = failure.status_code,
status_class = failure.status_class.unwrap_or_else(|| status_class(failure.status_code)),
timeout = failure.timeout,
retryable = failure.retryable,
latency_ms = failure.latency_ms,
prompt_chars = failure.prompt_chars,
reference_image_count = failure.reference_image_count,
image_model = failure.image_model,
request_id = %failure.request_id.as_deref().unwrap_or_default(),
error_source = %failure.error_source.as_deref().unwrap_or_default(),
error = %failure.error_message,
"外部 API 调用失败"
);
}
fn status_class(status_code: Option<u16>) -> &'static str {
match status_code {
Some(100..=199) => "1xx",
Some(200..=299) => "2xx",
Some(300..=399) => "3xx",
Some(400..=499) => "4xx",
Some(500..=599) => "5xx",
Some(_) => "unknown",
None => "transport",
}
}
fn audit_request_context() -> crate::request_context::RequestContext {
crate::request_context::RequestContext::new(
format!("external-api-audit-{}", Uuid::new_v4()),
"external-api audit".to_string(),
std::time::Duration::ZERO,
false,
)
}
fn truncate_field(value: &str, max_chars: usize) -> String {
value.chars().take(max_chars).collect()
}
fn current_utc_iso_text() -> String {
shared_kernel::format_rfc3339(OffsetDateTime::now_utc())
.unwrap_or_else(|_| "1970-01-01T00:00:00Z".to_string())
}
#[cfg(test)]
mod tests {
use serde_json::Value;
use super::*;
#[test]
fn external_api_failure_tracking_draft_uses_module_scope_and_safe_metadata() {
let draft = build_external_api_failure_tracking_draft(
&ExternalApiFailureDraft::new(
"vector-engine",
"https://vector.example/v1/images/generations",
"拼图 UI 背景图生成失败",
"upstream_status",
"上游 429",
)
.with_status_code(Some(429))
.with_retryable(true)
.with_error_source(Some(
"client error (SendRequest) -> connection closed before message completed"
.to_string(),
))
.with_latency_ms(Some(1234))
.with_prompt_chars(Some(88))
.with_reference_image_count(Some(2))
.with_image_model(Some("gpt-image-2-all")),
);
assert_eq!(draft.event_key, EXTERNAL_API_FAILURE_EVENT_KEY);
assert_eq!(draft.scope_kind, RuntimeTrackingScopeKind::Module);
assert_eq!(draft.scope_id, "vector-engine");
assert_eq!(draft.module_key, Some(EXTERNAL_API_AUDIT_MODULE_KEY));
let metadata = draft.metadata;
assert_eq!(metadata["provider"], "vector-engine");
assert_eq!(metadata["statusCode"], 429);
assert_eq!(metadata["statusClass"], "4xx");
assert_eq!(metadata["retryable"], true);
assert_eq!(metadata["latencyMs"], 1234);
assert_eq!(metadata["promptChars"], 88);
assert_eq!(metadata["referenceImageCount"], 2);
assert_eq!(metadata["imageModel"], "gpt-image-2-all");
assert_eq!(
metadata["errorSource"],
"client error (SendRequest) -> connection closed before message completed"
);
assert!(matches!(metadata["occurredAt"], Value::String(_)));
}
#[test]
fn retryable_classification_keeps_transport_and_overload_failures_actionable() {
assert!(is_retryable_external_api_failure(None, true, false));
assert!(is_retryable_external_api_failure(None, false, true));
assert!(is_retryable_external_api_failure(Some(429), false, false));
assert!(is_retryable_external_api_failure(Some(502), false, false));
assert!(!is_retryable_external_api_failure(Some(400), false, false));
}
#[test]
fn matting_failure_draft_marks_non_timeout_transport_retryable() {
// status_code=None、timeout=false 的非超时 transport 失败(DNS / 连接重置 / 读体中断):
// statusClass 必须是 transport 且 retryable=true,否则与 "transport failures actionable" 冲突。
let draft = build_matting_external_api_failure_draft(
"bgfilter",
"https://bgfilter.example/remove-background".to_string(),
"editor-screen-background-removal",
"bgfilter_segment",
None,
false,
Some(67),
"请求 BgFilter 服务失败:dns error".to_string(),
Some("dns error".to_string()),
&ExternalApiAuditContext::default(),
);
assert_eq!(draft.status_code, None);
assert_eq!(draft.status_class, Some("transport"));
assert!(!draft.timeout);
assert!(
draft.retryable,
"非超时 transport 失败应当 retryable=true,与 statusClass=transport 保持一致"
);
}
#[test]
fn matting_failure_draft_keeps_client_error_non_retryable() {
// 真实上游 4xx(非 429 / 408)不是传输故障,仍应 retryable=false,不能被误判为可重试。
let draft = build_matting_external_api_failure_draft(
"aliyun-matting",
"imageseg.example".to_string(),
"editor-screen-background-removal",
"aliyun_segment",
Some(400),
false,
Some(12),
"通用抠图接口返回失败(HTTP 400,Code=InvalidImage):bad image".to_string(),
None,
&ExternalApiAuditContext::default(),
);
assert_eq!(draft.status_code, Some(400));
assert_eq!(draft.status_class, Some("4xx"));
assert!(!draft.retryable);
}
#[test]
fn app_error_status_class_can_override_successful_upstream_status() {
let draft = build_external_api_failure_tracking_draft(
&ExternalApiFailureDraft::new(
"vector-engine",
"https://cdn.example/generated.png",
"下载生成图片",
"image_download",
"下载生成图片失败",
)
.with_status_code(Some(200))
.with_optional_status_class(Some(app_error_status_class(StatusCode::BAD_GATEWAY))),
);
assert_eq!(draft.metadata["statusCode"], 200);
assert_eq!(draft.metadata["statusClass"], "5xx");
}
#[test]
fn matting_failure_audit_uses_upstream_status_instead_of_wrapped_status() {
let error = AppError::from_status(StatusCode::BAD_GATEWAY).with_details(json!({
"provider": "bgfilter",
"message": "BgFilter 服务返回非成功状态",
"upstreamStatus": 429,
"upstreamMessage": "too many requests",
"latencyMs": 345,
}));
let status_code = matting_failure_audit_status_code(&error);
let timeout = matting_failure_audit_timeout(&error);
let tracking = build_external_api_failure_tracking_draft(
&ExternalApiFailureDraft::new(
"bgfilter",
"https://bgfilter.example/remove-background",
"editor-screen-background-removal",
"bgfilter_segment",
error.message(),
)
.with_status_code(status_code)
.with_optional_status_class(Some(status_class(status_code)))
.with_timeout(timeout)
.with_latency_ms(matting_failure_audit_latency_ms(&error))
.with_retryable(is_retryable_external_api_failure(
status_code,
timeout,
false,
))
.with_raw_excerpt(matting_failure_audit_raw_excerpt(&error)),
);
assert_eq!(tracking.metadata["statusCode"], 429);
assert_eq!(tracking.metadata["statusClass"], "4xx");
assert_eq!(tracking.metadata["timeout"], false);
assert_eq!(tracking.metadata["retryable"], true);
assert_eq!(tracking.metadata["latencyMs"], 345);
assert_eq!(tracking.metadata["rawExcerpt"], "too many requests");
}
#[test]
fn matting_failure_audit_keeps_transport_timeout_classification() {
let error = AppError::from_status(StatusCode::GATEWAY_TIMEOUT).with_details(json!({
"provider": "bgfilter",
"message": "请求 BgFilter 服务失败:operation timed out",
"timeout": true,
}));
let status_code = matting_failure_audit_status_code(&error);
let timeout = matting_failure_audit_timeout(&error);
let tracking = build_external_api_failure_tracking_draft(
&ExternalApiFailureDraft::new(
"bgfilter",
"https://bgfilter.example/remove-background",
"editor-screen-background-removal",
"bgfilter_segment",
error.message(),
)
.with_status_code(status_code)
.with_optional_status_class(Some(status_class(status_code)))
.with_timeout(timeout)
.with_retryable(is_retryable_external_api_failure(
status_code,
timeout,
false,
)),
);
assert_eq!(tracking.metadata["statusCode"], Value::Null);
assert_eq!(tracking.metadata["statusClass"], "transport");
assert_eq!(tracking.metadata["timeout"], true);
assert_eq!(tracking.metadata["retryable"], true);
}
#[test]
fn matting_failure_audit_keeps_non_timeout_transport_classification() {
let error = AppError::from_status(StatusCode::BAD_GATEWAY).with_details(json!({
"provider": "bgfilter",
"message": "请求 BgFilter 服务失败:dns error",
"timeout": false,
"transport": true,
"latencyMs": 67,
}));
let status_code = matting_failure_audit_status_code(&error);
let timeout = matting_failure_audit_timeout(&error);
let tracking = build_external_api_failure_tracking_draft(
&ExternalApiFailureDraft::new(
"bgfilter",
"https://bgfilter.example/remove-background",
"editor-screen-background-removal",
"bgfilter_segment",
error.message(),
)
.with_status_code(status_code)
.with_optional_status_class(Some(status_class(status_code)))
.with_timeout(timeout)
.with_latency_ms(matting_failure_audit_latency_ms(&error))
.with_retryable(is_retryable_external_api_failure(
status_code,
timeout,
false,
)),
);
assert_eq!(tracking.metadata["statusCode"], Value::Null);
assert_eq!(tracking.metadata["statusClass"], "transport");
assert_eq!(tracking.metadata["timeout"], false);
assert_eq!(tracking.metadata["retryable"], false);
assert_eq!(tracking.metadata["latencyMs"], 67);
}
#[test]
fn matting_failure_external_call_attempted_defaults_to_true() {
let attempted = AppError::from_status(StatusCode::BAD_GATEWAY).with_details(json!({
"provider": "aliyun-matting",
"message": "真实上游失败",
}));
let skipped = AppError::from_status(StatusCode::SERVICE_UNAVAILABLE).with_details(json!({
"provider": "aliyun-matting",
"message": "阿里云抠图客户端未配置或未启用。",
"externalCallAttempted": false,
}));
assert!(matting_failure_external_call_attempted(&attempted));
assert!(!matting_failure_external_call_attempted(&skipped));
}
}