feat(observability): 引入错误链路 error_flow 元数据并区分上游/客户端错误

- 网关在本地 failover 时构建 error_flow 元数据(分类/决策/传播策略),写入 report_context
- scheduler-core 解析并透传 error_flow 至候选 extra_data
- admin usage 详情拆分 request/upstream/client/failure_summary 错误域,敏感上游错误标记为 suppressed
- 前端 RequestDetailDrawer 拆出"返回客户端"与"上游响应"双错误卡片
- HorizontalRequestTimeline 节点详情展示真实请求错误及 error_flow 标签
This commit is contained in:
fawney19
2026-04-25 16:58:17 +08:00
parent bc97e383d3
commit 00744c0ce5
9 changed files with 997 additions and 44 deletions

View File

@@ -315,6 +315,307 @@ fn admin_usage_strip_trace_metadata(metadata: &mut serde_json::Map<String, Value
metadata.remove("trace_id");
}
fn admin_usage_string_field<'a>(value: &'a Value, field: &str) -> Option<&'a str> {
value
.get(field)
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
}
fn admin_usage_error_field<'a>(body: &'a Value, field: &str) -> Option<&'a str> {
body.get("error")
.and_then(|error| match error {
Value::Object(object) => object.get(field).and_then(Value::as_str),
_ => None,
})
.map(str::trim)
.filter(|value| !value.is_empty())
.or_else(|| admin_usage_string_field(body, field))
}
fn admin_usage_error_message_from_body(body: &Value) -> Option<String> {
body.get("error")
.and_then(|error| match error {
Value::Object(object) => object.get("message").and_then(Value::as_str),
Value::String(message) => Some(message.as_str()),
_ => None,
})
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
.or_else(|| {
admin_usage_string_field(body, "message")
.or_else(|| admin_usage_string_field(body, "detail"))
.map(ToOwned::to_owned)
})
}
fn admin_usage_error_type_from_body(body: &Value) -> Option<String> {
admin_usage_error_field(body, "type")
.or_else(|| admin_usage_error_field(body, "status"))
.or_else(|| admin_usage_error_field(body, "kind"))
.map(ToOwned::to_owned)
}
fn admin_usage_error_code_from_body(body: &Value) -> Option<Value> {
body.get("error")
.and_then(|error| match error {
Value::Object(object) => object.get("code").cloned(),
_ => None,
})
.or_else(|| body.get("code").cloned())
}
fn admin_usage_header_content_type(headers: Option<&Value>) -> Option<String> {
let object = headers?.as_object()?;
for (key, value) in object {
if key.eq_ignore_ascii_case("content-type") {
return value
.as_str()
.map(ToOwned::to_owned)
.or_else(|| (!value.is_null()).then(|| value.to_string()));
}
}
None
}
fn admin_usage_error_domain_json(
source: &str,
status_code: Option<u16>,
headers: Option<&Value>,
body: Option<&Value>,
fallback_type: Option<&str>,
fallback_message: Option<&str>,
) -> Value {
let error_type = body
.and_then(admin_usage_error_type_from_body)
.or_else(|| fallback_type.map(ToOwned::to_owned));
let message = body
.and_then(admin_usage_error_message_from_body)
.or_else(|| {
fallback_message
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
});
let code = body.and_then(admin_usage_error_code_from_body);
let content_type = admin_usage_header_content_type(headers);
if status_code.is_none()
&& error_type.is_none()
&& message.is_none()
&& code.is_none()
&& body.is_none()
{
return Value::Null;
}
json!({
"source": source,
"status_code": status_code,
"type": error_type,
"message": message,
"code": code,
"content_type": content_type,
"body": body.cloned().unwrap_or(Value::Null),
})
}
fn admin_usage_error_domain_message(domain: &Value) -> Option<String> {
domain
.get("message")
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
fn admin_usage_error_domain_type(domain: &Value) -> Option<String> {
domain
.get("type")
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
fn admin_usage_error_domain_source(domain: &Value) -> Option<String> {
domain
.get("source")
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
fn admin_usage_failure_summary_json(
item: &StoredRequestUsageAudit,
request_error: &Value,
upstream_error: &Value,
client_error: &Value,
) -> Value {
let selected = [client_error, upstream_error, request_error]
.into_iter()
.find(|domain| !domain.is_null() && admin_usage_error_domain_message(domain).is_some());
let message = selected
.and_then(admin_usage_error_domain_message)
.or_else(|| {
item.error_message
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
});
let Some(message) = message else {
return Value::Null;
};
let error_type = selected
.and_then(admin_usage_error_domain_type)
.or_else(|| item.error_category.clone());
let source = selected
.and_then(admin_usage_error_domain_source)
.unwrap_or_else(|| "usage_summary".to_string());
json!({
"source": source,
"status_code": item.status_code,
"type": error_type,
"message": message,
"category": item.error_category,
})
}
fn admin_usage_error_domains_json(item: &StoredRequestUsageAudit) -> Value {
let has_upstream_attempt = item.candidate_id.is_some()
|| item.provider_api_key_id.is_some()
|| item.provider_request_headers.is_some()
|| item.provider_request_body.is_some()
|| item.provider_request_body_ref.is_some();
let upstream_error = if has_upstream_attempt {
admin_usage_error_domain_json(
"upstream_response",
item.status_code,
item.response_headers.as_ref(),
item.response_body.as_ref(),
item.error_category.as_deref(),
item.error_message.as_deref(),
)
} else {
Value::Null
};
let client_error = admin_usage_error_domain_json(
"client_response",
item.status_code,
item.client_response_headers.as_ref(),
item.client_response_body.as_ref(),
item.error_category.as_deref(),
item.error_message.as_deref(),
);
let request_error = Value::Null;
let failure_summary =
admin_usage_failure_summary_json(item, &request_error, &upstream_error, &client_error);
json!({
"request_error": request_error,
"upstream_error": upstream_error,
"client_error": client_error,
"failure_summary": failure_summary,
})
}
fn admin_usage_error_domain_search_text(domain: &Value) -> String {
[
domain.get("type").and_then(Value::as_str),
domain.get("message").and_then(Value::as_str),
domain.get("code").and_then(Value::as_str),
]
.into_iter()
.flatten()
.map(str::to_ascii_lowercase)
.collect::<Vec<_>>()
.join(" ")
}
fn admin_usage_upstream_error_is_sensitive(domain: &Value) -> bool {
let text = admin_usage_error_domain_search_text(domain);
[
"insufficient_quota",
"insufficient quota",
"quota exhausted",
"credits exhausted",
"credit balance",
"credit limit",
"payment_required",
"payment required",
"account disabled",
"account_deactivated",
"subscription inactive",
"verification required",
]
.iter()
.any(|pattern| text.contains(pattern))
}
fn admin_usage_error_flow_json(item: &StoredRequestUsageAudit, error_domains: &Value) -> Value {
let request_error = &error_domains["request_error"];
let upstream_error = &error_domains["upstream_error"];
let client_error = &error_domains["client_error"];
let failure_summary = &error_domains["failure_summary"];
if request_error.is_null()
&& upstream_error.is_null()
&& client_error.is_null()
&& failure_summary.is_null()
{
return Value::Null;
}
let upstream_sensitive =
!upstream_error.is_null() && admin_usage_upstream_error_is_sensitive(upstream_error);
let propagation = if upstream_sensitive {
"suppressed"
} else if !client_error.is_null() && !upstream_error.is_null() {
let upstream_message = admin_usage_error_domain_message(upstream_error);
let client_message = admin_usage_error_domain_message(client_error);
if upstream_message.is_some() && upstream_message == client_message {
"passthrough"
} else {
"converted"
}
} else if !client_error.is_null() {
"local"
} else if !upstream_error.is_null() {
"captured"
} else {
"none"
};
let source = if !request_error.is_null() {
"request"
} else if !upstream_error.is_null() {
"upstream"
} else if !client_error.is_null() {
"gateway"
} else {
"summary"
};
let client_response_source = if client_error.is_null() {
Value::Null
} else if !upstream_error.is_null() {
Value::String("converted_or_sanitized_upstream".to_string())
} else {
Value::String("gateway_generated".to_string())
};
json!({
"source": source,
"status_code": item.status_code,
"propagation": propagation,
"client_response_source": client_response_source,
"safe_to_expose_upstream": !upstream_sensitive,
"summary_source": failure_summary.get("source").cloned().unwrap_or(Value::Null),
})
}
fn maybe_insert_number_field(
object: &mut serde_json::Map<String, Value>,
key: &str,
@@ -1790,6 +2091,14 @@ pub fn build_admin_usage_detail_payload(
payload["body_capture"] = admin_usage_body_capture_json(item);
payload["settlement"] = admin_usage_settlement_json(item);
payload["trace"] = admin_usage_trace_json(item);
let error_domains = admin_usage_error_domains_json(item);
let error_flow = admin_usage_error_flow_json(item, &error_domains);
payload["errors"] = error_domains.clone();
payload["request_error"] = error_domains["request_error"].clone();
payload["upstream_error"] = error_domains["upstream_error"].clone();
payload["client_error"] = error_domains["client_error"].clone();
payload["failure_summary"] = error_domains["failure_summary"].clone();
payload["error_flow"] = error_flow;
payload["has_request_body"] = json!(admin_usage_has_body_value(
item,
item.request_body.as_ref(),
@@ -2261,6 +2570,158 @@ mod tests {
assert_eq!(payload["has_client_response_body"], true);
}
#[test]
fn detail_payload_separates_upstream_client_and_summary_errors() {
let item = StoredRequestUsageAudit {
error_message: Some(
"execution runtime stream returned retryable status 400".to_string(),
),
error_category: Some("server_error".to_string()),
response_headers: Some(json!({
"content-type": "application/json"
})),
response_body: Some(json!({
"error": {
"type": "retryable_upstream_status",
"message": "execution runtime stream returned retryable status 400",
"code": 400
}
})),
client_response_headers: Some(json!({
"content-type": "application/json"
})),
client_response_body: Some(json!({
"error": {
"type": "http_error",
"message": "local execution runtime exhausted"
}
})),
..sample_usage(
"failed",
Some(503),
Some("execution runtime stream returned retryable status 400"),
)
};
let payload = build_admin_usage_detail_payload(
&item,
&BTreeMap::new(),
&BTreeMap::new(),
false,
false,
None,
true,
Some(json!({"model": "gpt-5.4"})),
&BTreeMap::new(),
);
assert_eq!(payload["upstream_error"]["source"], "upstream_response");
assert_eq!(
payload["errors"]["upstream_error"],
payload["upstream_error"]
);
assert_eq!(
payload["upstream_error"]["message"],
"execution runtime stream returned retryable status 400"
);
assert_eq!(payload["client_error"]["source"], "client_response");
assert_eq!(
payload["client_error"]["message"],
"local execution runtime exhausted"
);
assert_eq!(payload["failure_summary"]["source"], "client_response");
assert_eq!(
payload["failure_summary"]["message"],
"local execution runtime exhausted"
);
assert_eq!(payload["error_flow"]["propagation"], "converted");
assert!(payload["request_error"].is_null());
}
#[test]
fn detail_payload_marks_sensitive_upstream_account_errors_suppressed() {
let item = StoredRequestUsageAudit {
error_message: Some("credit balance exhausted".to_string()),
error_category: Some("server_error".to_string()),
response_body: Some(json!({
"error": {
"type": "insufficient_quota",
"message": "credit balance exhausted"
}
})),
client_response_body: Some(json!({
"error": {
"type": "http_error",
"message": "upstream provider unavailable"
}
})),
..sample_usage("failed", Some(503), Some("credit balance exhausted"))
};
let payload = build_admin_usage_detail_payload(
&item,
&BTreeMap::new(),
&BTreeMap::new(),
false,
false,
None,
true,
Some(json!({"model": "gpt-5.4"})),
&BTreeMap::new(),
);
assert_eq!(payload["error_flow"]["propagation"], "suppressed");
assert_eq!(payload["error_flow"]["safe_to_expose_upstream"], false);
assert_eq!(payload["upstream_error"]["type"], "insufficient_quota");
assert_eq!(
payload["client_error"]["message"],
"upstream provider unavailable"
);
}
#[test]
fn detail_payload_does_not_promote_local_client_error_to_upstream_error() {
let message = "没有可用提供商支持模型 gpt-5.4 的同步请求。请检查模型映射、端点启用状态和 API Key 权限(原因代码: candidate_list_empty";
let item = StoredRequestUsageAudit {
provider_api_key_id: None,
provider_request_headers: None,
provider_request_body: None,
provider_request_body_ref: None,
candidate_id: None,
error_category: Some("http_error".to_string()),
client_response_body: Some(json!({
"error": {
"type": "http_error",
"message": message
}
})),
response_body: Some(json!({
"error": {
"type": "http_error",
"message": message
}
})),
..sample_usage("failed", Some(503), Some(message))
};
let payload = build_admin_usage_detail_payload(
&item,
&BTreeMap::new(),
&BTreeMap::new(),
false,
false,
None,
true,
Some(json!({"model": "gpt-5.4"})),
&BTreeMap::new(),
);
assert!(payload["upstream_error"].is_null());
assert_eq!(payload["client_error"]["message"], message);
assert_eq!(payload["error_flow"]["source"], "gateway");
assert_eq!(payload["error_flow"]["propagation"], "local");
}
#[test]
fn detail_payload_preserves_legacy_body_capture_metadata_keys() {
let item = StoredRequestUsageAudit {

View File

@@ -23,6 +23,7 @@ pub struct SchedulerRequestCandidateReportContext {
pub header_rules: Option<Value>,
pub body_rules: Option<Value>,
pub proxy: Option<Value>,
pub error_flow: Option<Value>,
}
#[derive(Debug, Clone, PartialEq)]
@@ -47,6 +48,19 @@ pub struct SchedulerExecutionRequestCandidateSeed {
pub report_context: Value,
}
#[derive(Debug, Clone, Default)]
struct ReportCandidateExtraDataInput {
client_api_format: Option<String>,
provider_api_format: Option<String>,
upstream_url: Option<String>,
mapped_model: Option<String>,
key_name: Option<String>,
header_rules: Option<Value>,
body_rules: Option<Value>,
proxy: Option<Value>,
error_flow: Option<Value>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SchedulerRequestCandidateStatusUpdate {
pub status: RequestCandidateStatus,
@@ -123,6 +137,10 @@ pub fn parse_request_candidate_report_context(
.get("proxy")
.cloned()
.filter(|value| !value.is_null()),
error_flow: report_context
.get("error_flow")
.cloned()
.filter(|value| !value.is_null()),
})
}
@@ -151,6 +169,7 @@ pub fn resolve_report_request_candidate_slot(
header_rules,
body_rules,
proxy,
error_flow,
} = metadata;
let request_id = request_id?;
let synthesized_extra_data = build_report_candidate_extra_data(ReportCandidateExtraDataInput {
@@ -162,6 +181,7 @@ pub fn resolve_report_request_candidate_slot(
header_rules,
body_rules,
proxy,
error_flow,
});
let created_at_unix_ms = matched_candidate
.as_ref()
@@ -325,6 +345,7 @@ pub fn build_local_request_candidate_status_record(
header_rules: metadata.header_rules.clone(),
body_rules: metadata.body_rules.clone(),
proxy: metadata.proxy.clone(),
error_flow: metadata.error_flow.clone(),
});
let created_at_unix_ms = started_at_unix_ms.or(finished_at_unix_ms);
@@ -526,17 +547,6 @@ fn next_candidate_index(candidates: &[StoredRequestCandidate]) -> u32 {
.unwrap_or_default()
}
struct ReportCandidateExtraDataInput {
client_api_format: Option<String>,
provider_api_format: Option<String>,
upstream_url: Option<String>,
mapped_model: Option<String>,
key_name: Option<String>,
header_rules: Option<Value>,
body_rules: Option<Value>,
proxy: Option<Value>,
}
fn build_report_candidate_extra_data(input: ReportCandidateExtraDataInput) -> Option<Value> {
let ReportCandidateExtraDataInput {
client_api_format,
@@ -547,6 +557,7 @@ fn build_report_candidate_extra_data(input: ReportCandidateExtraDataInput) -> Op
header_rules,
body_rules,
proxy,
error_flow,
} = input;
let mut extra_data = Map::with_capacity(8);
extra_data.insert("gateway_execution_runtime".to_string(), Value::Bool(true));
@@ -581,6 +592,9 @@ fn build_report_candidate_extra_data(input: ReportCandidateExtraDataInput) -> Op
if let Some(proxy) = proxy {
extra_data.insert("proxy".to_string(), proxy);
}
if let Some(error_flow) = error_flow {
extra_data.insert("error_flow".to_string(), error_flow);
}
(!extra_data.is_empty()).then_some(Value::Object(extra_data))
}
@@ -729,6 +743,11 @@ mod tests {
"node_id": "proxy-node-1",
"node_name": "edge-1",
"source": "provider"
},
"error_flow": {
"classification": "retry_upstream_failure",
"decision": "retry_next_candidate",
"propagation": "suppressed"
}
})))
.expect("metadata");
@@ -777,6 +796,13 @@ mod tests {
.map(Vec::len),
Some(1)
);
assert_eq!(
slot.extra_data
.as_ref()
.and_then(|value| value.get("error_flow"))
.and_then(|value| value.get("propagation")),
Some(&json!("suppressed"))
);
}
#[test]