diff --git a/apps/aether-gateway/src/execution_runtime/stream/execution_failures.rs b/apps/aether-gateway/src/execution_runtime/stream/execution_failures.rs index 117bdb71d..0361035fd 100644 --- a/apps/aether-gateway/src/execution_runtime/stream/execution_failures.rs +++ b/apps/aether-gateway/src/execution_runtime/stream/execution_failures.rs @@ -148,18 +148,17 @@ pub(super) fn build_stream_failure_from_execution_error( error: &ExecutionError, ) -> StreamFailureReport { let transport_error = execution_error_is_transport(error); - let status_code = error.upstream_status.unwrap_or_else(|| { - if matches!( - error.kind, - ExecutionErrorKind::ConnectTimeout - | ExecutionErrorKind::FirstByteTimeout - | ExecutionErrorKind::ReadTimeout - ) { - 504 - } else { - 502 - } - }); + let fallback_status_code = if matches!( + error.kind, + ExecutionErrorKind::ConnectTimeout + | ExecutionErrorKind::FirstByteTimeout + | ExecutionErrorKind::ReadTimeout + ) { + 504 + } else { + 502 + }; + let status_code = error.upstream_status.unwrap_or(fallback_status_code); let error_type = serde_json::to_value(&error.kind) .ok() .and_then(|value| value.as_str().map(ToOwned::to_owned)) @@ -639,6 +638,7 @@ pub(super) async fn handle_prefetch_stream_failure( ) .await; } + let honor_local_failover = honor_http_failover && retry_scope_out.is_some(); let failure_analysis = record_stream_sync_failure( state, plan, @@ -646,14 +646,14 @@ pub(super) async fn handle_prefetch_stream_failure( &payload, candidate_status_code, None, - if honor_http_failover { + if honor_local_failover { StreamFailureHandling::HonorLocalFailover } else { StreamFailureHandling::Terminal }, ) .await; - if honor_http_failover + if honor_local_failover && matches!( failure_analysis.decision, LocalFailoverDecision::RetryNextCandidate @@ -722,8 +722,8 @@ async fn handle_prefetch_transport_stream_failure( payload.report_context.as_ref(), ) .await; - let retrying_next_candidate = - matches!(analysis.decision, LocalFailoverDecision::RetryNextCandidate); + let retrying_next_candidate = retry_scope_out.is_some() + && matches!(analysis.decision, LocalFailoverDecision::RetryNextCandidate); if !retrying_next_candidate { crate::execution_runtime::mark_stream_candidate_watchdog_terminal_started(); let report_context_with_diagnostics = diff --git a/apps/aether-gateway/src/executor/orchestration.rs b/apps/aether-gateway/src/executor/orchestration.rs index 1137a9943..eb7e59f03 100644 --- a/apps/aether-gateway/src/executor/orchestration.rs +++ b/apps/aether-gateway/src/executor/orchestration.rs @@ -1577,6 +1577,8 @@ mod tests { const TEST_OPENAI_IMAGE_SYNC_PLAN_KIND: &str = "openai_image_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 { attempts: VecDeque, @@ -1733,11 +1735,32 @@ mod tests { usage_repository: &InMemoryUsageReadRepository, request_id: &str, ) { - let usage = usage_repository - .find_by_request_id(request_id) - .await - .expect("usage should read") - .expect("terminal usage should be recorded"); + let deadline = Instant::now() + HEARTBEAT_USAGE_SETTLE_TIMEOUT; + let usage = loop { + let usage = usage_repository + .find_by_request_id(request_id) + .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("") + ); + 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 .request_metadata .as_ref() diff --git a/crates/aether-usage/runtime/src/runtime.rs b/crates/aether-usage/runtime/src/runtime.rs index 60d59217a..926ea1185 100644 --- a/crates/aether-usage/runtime/src/runtime.rs +++ b/crates/aether-usage/runtime/src/runtime.rs @@ -12692,7 +12692,6 @@ mod tests { apply_usage_body_capture_policy_to_event( UsageBodyCapturePolicy { record_level: UsageRequestRecordLevel::Basic, - ..UsageBodyCapturePolicy::default() }, &mut event, ); diff --git a/crates/aether-usage/runtime/src/write.rs b/crates/aether-usage/runtime/src/write.rs index 2875dd098..93d1d760a 100644 --- a/crates/aether-usage/runtime/src/write.rs +++ b/crates/aether-usage/runtime/src/write.rs @@ -1728,14 +1728,17 @@ fn build_usage_event_data_seed_with_detail( .as_deref() .and_then(infer_endpoint_kind) .map(ToOwned::to_owned); - let request_metadata = build_runtime_request_metadata_seed_from_parts( - plan, - context, - request_capture.request_body.is_some(), - request_capture.request_body_ref.as_deref(), - request_capture.provider_request.is_some(), - request_capture.provider_request_body_ref.as_deref(), - plan.body.body_bytes_b64.as_deref(), + let request_metadata = merge_usage_request_metadata_owned( + build_usage_request_metadata_seed(plan, context), + build_runtime_request_metadata_seed_from_parts( + plan, + context, + request_capture.request_body.is_some(), + request_capture.request_body_ref.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 { 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] fn masks_known_sensitive_header_values() { let token = "Bearer eyJhbGciOiJSUzI1NiJ9.payload-here.signature-tail";