fix(gateway): restore failover and usage diagnostics

This commit is contained in:
elky
2026-07-30 09:12:11 +08:00
parent a04673a90d
commit 1ab4f079c9
4 changed files with 116 additions and 30 deletions
@@ -148,18 +148,17 @@ pub(super) fn build_stream_failure_from_execution_error(
error: &ExecutionError, error: &ExecutionError,
) -> StreamFailureReport { ) -> StreamFailureReport {
let transport_error = execution_error_is_transport(error); let transport_error = execution_error_is_transport(error);
let status_code = error.upstream_status.unwrap_or_else(|| { let fallback_status_code = if matches!(
if matches!( error.kind,
error.kind, ExecutionErrorKind::ConnectTimeout
ExecutionErrorKind::ConnectTimeout | ExecutionErrorKind::FirstByteTimeout
| ExecutionErrorKind::FirstByteTimeout | ExecutionErrorKind::ReadTimeout
| ExecutionErrorKind::ReadTimeout ) {
) { 504
504 } else {
} else { 502
502 };
} let status_code = error.upstream_status.unwrap_or(fallback_status_code);
});
let error_type = serde_json::to_value(&error.kind) let error_type = serde_json::to_value(&error.kind)
.ok() .ok()
.and_then(|value| value.as_str().map(ToOwned::to_owned)) .and_then(|value| value.as_str().map(ToOwned::to_owned))
@@ -639,6 +638,7 @@ pub(super) async fn handle_prefetch_stream_failure(
) )
.await; .await;
} }
let honor_local_failover = honor_http_failover && retry_scope_out.is_some();
let failure_analysis = record_stream_sync_failure( let failure_analysis = record_stream_sync_failure(
state, state,
plan, plan,
@@ -646,14 +646,14 @@ pub(super) async fn handle_prefetch_stream_failure(
&payload, &payload,
candidate_status_code, candidate_status_code,
None, None,
if honor_http_failover { if honor_local_failover {
StreamFailureHandling::HonorLocalFailover StreamFailureHandling::HonorLocalFailover
} else { } else {
StreamFailureHandling::Terminal StreamFailureHandling::Terminal
}, },
) )
.await; .await;
if honor_http_failover if honor_local_failover
&& matches!( && matches!(
failure_analysis.decision, failure_analysis.decision,
LocalFailoverDecision::RetryNextCandidate LocalFailoverDecision::RetryNextCandidate
@@ -722,8 +722,8 @@ async fn handle_prefetch_transport_stream_failure(
payload.report_context.as_ref(), payload.report_context.as_ref(),
) )
.await; .await;
let retrying_next_candidate = let retrying_next_candidate = retry_scope_out.is_some()
matches!(analysis.decision, LocalFailoverDecision::RetryNextCandidate); && matches!(analysis.decision, LocalFailoverDecision::RetryNextCandidate);
if !retrying_next_candidate { if !retrying_next_candidate {
crate::execution_runtime::mark_stream_candidate_watchdog_terminal_started(); crate::execution_runtime::mark_stream_candidate_watchdog_terminal_started();
let report_context_with_diagnostics = let report_context_with_diagnostics =
@@ -1577,6 +1577,8 @@ mod tests {
const TEST_OPENAI_IMAGE_SYNC_PLAN_KIND: &str = "openai_image_sync"; const TEST_OPENAI_IMAGE_SYNC_PLAN_KIND: &str = "openai_image_sync";
const TEST_STANDARD_TEXT_SYNC_PLAN_KIND: &str = "openai_responses_compact_sync"; const TEST_STANDARD_TEXT_SYNC_PLAN_KIND: &str = "openai_responses_compact_sync";
const HEARTBEAT_USAGE_POLL_INTERVAL: Duration = Duration::from_millis(10);
const HEARTBEAT_USAGE_SETTLE_TIMEOUT: Duration = Duration::from_secs(30);
struct TestSyncAttemptSource { struct TestSyncAttemptSource {
attempts: VecDeque<AiSyncAttempt>, attempts: VecDeque<AiSyncAttempt>,
@@ -1733,11 +1735,32 @@ mod tests {
usage_repository: &InMemoryUsageReadRepository, usage_repository: &InMemoryUsageReadRepository,
request_id: &str, request_id: &str,
) { ) {
let usage = usage_repository let deadline = Instant::now() + HEARTBEAT_USAGE_SETTLE_TIMEOUT;
.find_by_request_id(request_id) let usage = loop {
.await let usage = usage_repository
.expect("usage should read") .find_by_request_id(request_id)
.expect("terminal usage should be recorded"); .await
.expect("usage should read");
if usage.as_ref().is_some_and(|usage| {
matches!(usage.status.as_str(), "completed" | "failed" | "cancelled")
}) {
break usage.expect("terminal usage should be recorded");
}
let now = Instant::now();
let last_status = usage.as_ref().map(|usage| usage.status.as_str());
assert!(
now < deadline,
"terminal usage should be recorded within {HEARTBEAT_USAGE_SETTLE_TIMEOUT:?}; \
last status: {}",
last_status.unwrap_or("<missing>")
);
tokio::time::sleep(HEARTBEAT_USAGE_POLL_INTERVAL.min(deadline - now)).await;
};
assert_eq!(
usage.status, "completed",
"heartbeat usage should complete successfully"
);
let request_metadata = usage let request_metadata = usage
.request_metadata .request_metadata
.as_ref() .as_ref()
@@ -12692,7 +12692,6 @@ mod tests {
apply_usage_body_capture_policy_to_event( apply_usage_body_capture_policy_to_event(
UsageBodyCapturePolicy { UsageBodyCapturePolicy {
record_level: UsageRequestRecordLevel::Basic, record_level: UsageRequestRecordLevel::Basic,
..UsageBodyCapturePolicy::default()
}, },
&mut event, &mut event,
); );
+72 -8
View File
@@ -1728,14 +1728,17 @@ fn build_usage_event_data_seed_with_detail(
.as_deref() .as_deref()
.and_then(infer_endpoint_kind) .and_then(infer_endpoint_kind)
.map(ToOwned::to_owned); .map(ToOwned::to_owned);
let request_metadata = build_runtime_request_metadata_seed_from_parts( let request_metadata = merge_usage_request_metadata_owned(
plan, build_usage_request_metadata_seed(plan, context),
context, build_runtime_request_metadata_seed_from_parts(
request_capture.request_body.is_some(), plan,
request_capture.request_body_ref.as_deref(), context,
request_capture.provider_request.is_some(), request_capture.request_body.is_some(),
request_capture.provider_request_body_ref.as_deref(), request_capture.request_body_ref.as_deref(),
plan.body.body_bytes_b64.as_deref(), request_capture.provider_request.is_some(),
request_capture.provider_request_body_ref.as_deref(),
plan.body.body_bytes_b64.as_deref(),
),
); );
sanitize_usage_event_data(UsageEventData { sanitize_usage_event_data(UsageEventData {
user_id: context_string(context, "user_id"), user_id: context_string(context, "user_id"),
@@ -6729,6 +6732,67 @@ mod tests {
); );
} }
#[test]
fn usage_event_data_seed_preserves_timing_and_runtime_body_metadata() {
let plan = ExecutionPlan {
request_id: "req-seed-metadata-1".to_string(),
candidate_id: Some("cand-seed-metadata-1".to_string()),
provider_name: Some("OpenAI".to_string()),
provider_id: "provider-1".to_string(),
endpoint_id: "endpoint-1".to_string(),
key_id: "key-1".to_string(),
method: "POST".to_string(),
url: "https://example.com/v1/chat/completions".to_string(),
headers: BTreeMap::new(),
content_type: Some("application/json".to_string()),
content_encoding: None,
body: RequestBody::from_json(json!({"model": "gpt-5"})),
stream: false,
client_api_format: "openai:chat".to_string(),
provider_api_format: "openai:chat".to_string(),
model_name: Some("gpt-5".to_string()),
proxy: None,
transport_profile: None,
timeouts: None,
};
let data = build_usage_event_data_seed(
&plan,
Some(&json!({
"trace_id": "trace-seed-metadata-1",
"end_to_end_time_ms": 10_626,
"end_to_end_first_byte_time_ms": 10_120,
"db_timings_ms": {"query_count": 2},
"original_request_body": {"messages": []}
})),
);
let metadata = data
.request_metadata
.as_ref()
.and_then(Value::as_object)
.expect("request metadata should be preserved");
assert_eq!(metadata.get("end_to_end_time_ms"), Some(&json!(10_626)));
assert_eq!(
metadata.get("end_to_end_first_byte_time_ms"),
Some(&json!(10_120))
);
assert_eq!(
metadata.get("db_timings_ms"),
Some(&json!({"query_count": 2}))
);
assert_eq!(
metadata.get("trace_id"),
Some(&json!("trace-seed-metadata-1"))
);
let body_size = metadata
.get("body_size")
.and_then(Value::as_object)
.expect("runtime body size metadata should remain");
assert!(body_size.get("client_request_body").is_some());
assert!(body_size.get("provider_request_body").is_some());
}
#[test] #[test]
fn masks_known_sensitive_header_values() { fn masks_known_sensitive_header_values() {
let token = "Bearer eyJhbGciOiJSUzI1NiJ9.payload-here.signature-tail"; let token = "Bearer eyJhbGciOiJSUzI1NiJ9.payload-here.signature-tail";