diff --git a/apps/aether-gateway/src/executor/candidate_loop.rs b/apps/aether-gateway/src/executor/candidate_loop.rs index 8806e66e3..86495df4d 100644 --- a/apps/aether-gateway/src/executor/candidate_loop.rs +++ b/apps/aether-gateway/src/executor/candidate_loop.rs @@ -522,7 +522,6 @@ where decision, plan_kind, transfer_tracker, - request_first_byte_started_at: Instant::now(), }; let loop_result = run_ai_attempt_loop(&port, plan_and_reports).await; if loop_result.is_err() { @@ -603,7 +602,6 @@ where decision, plan_kind, transfer_tracker, - request_first_byte_started_at: Instant::now(), }; let loop_result = run_dynamic_attempt_loop( &port, @@ -1121,10 +1119,6 @@ struct StreamAttemptLoopPort<'a> { decision: &'a GatewayControlDecision, plan_kind: &'a str, transfer_tracker: &'a ProviderTransferTracker, - /// All candidates in one downstream stream request share this origin. - /// Without it every retry receives a fresh full first-byte timeout and a - /// 30-second provider timeout can accumulate into a 60-120 second stall. - request_first_byte_started_at: Instant, } #[async_trait] @@ -1254,7 +1248,6 @@ where self.plan_kind, plan, watchdog_report_context, - self.request_first_byte_started_at, stop_on_transport_errors, move || async move { if let Some(response) = execution_plan_cost_capacity_response( @@ -1308,7 +1301,7 @@ where http::StatusCode::GATEWAY_TIMEOUT.as_u16(), "local_stream_candidate_watchdog_timeout", stream_candidate_watchdog_timeout_message(), - self.request_first_byte_started_at.elapsed().as_millis() as u64, + watchdog_started_at.elapsed().as_millis() as u64, ) .await?, ) @@ -1758,7 +1751,6 @@ async fn execute_stream_candidate_with_watchdog( plan_kind: &str, plan: &aether_contracts::ExecutionPlan, report_context: Option<&serde_json::Value>, - request_first_byte_started_at: Instant, stop_on_transport_errors: bool, execute: impl FnOnce() -> Fut, ) -> Result @@ -1768,7 +1760,6 @@ where > + Send, { let timeout_duration = resolve_stream_candidate_watchdog_timeout(plan, report_context); - let request_first_byte_deadline = request_first_byte_started_at + timeout_duration; let candidate_started_at = std::time::Instant::now(); let candidate_started_unix_ms = current_unix_ms(); let permit = match acquire_upstream_execution_gate(state, trace_id).await { @@ -1794,14 +1785,7 @@ where let watchdog_progress = StreamCandidateWatchdogProgress::shared(); let execution = watchdog_progress.clone().scope(execute()); tokio::pin!(execution); - // This is an absolute request-level deadline, not a new timeout for this - // candidate. Retries therefore consume only the budget left by earlier - // candidates instead of resetting the full provider timeout. - let candidate_budget_ms = request_first_byte_deadline - .saturating_duration_since(Instant::now()) - .as_millis() - .min(u128::from(u64::MAX)) as u64; - let deadline = tokio::time::sleep_until(request_first_byte_deadline); + let deadline = tokio::time::sleep(timeout_duration); tokio::pin!(deadline); let execution_result = tokio::select! { biased; @@ -1830,10 +1814,6 @@ where .map(|value| value.to_string()) .unwrap_or_else(|| "-".to_string()); let timeout_ms = u64::try_from(timeout_duration.as_millis()).unwrap_or(u64::MAX); - let request_elapsed_ms = request_first_byte_started_at - .elapsed() - .as_millis() - .min(u128::from(u64::MAX)) as u64; record_local_request_candidate_status( state, plan, @@ -1862,8 +1842,6 @@ where model_name, candidate_index = candidate_index.as_str(), timeout_ms, - candidate_budget_ms, - request_elapsed_ms, "gateway local stream candidate watchdog timed out" ); if stop_on_transport_errors { @@ -3150,7 +3128,6 @@ mod tests { "claude_cli_stream", &plan, Some(&report_context), - Instant::now(), false, || { std::future::pending::< @@ -3189,39 +3166,32 @@ mod tests { assert_eq!(record.candidate_index, 2); } - #[tokio::test] - async fn stream_candidate_retry_does_not_reset_an_expired_request_first_byte_budget() { - let writer = Arc::new(TestRequestCandidateWriter::default()); + async fn assert_stream_candidate_retry_gets_fresh_first_byte_budget( + provider_id: &str, + key_id: &str, + first_byte_ms: u64, + ) { + let writer = TestRequestCandidateWriter::default(); let plan = test_plan(Some(ExecutionTimeouts { - first_byte_ms: Some(250), + first_byte_ms: Some(100), ..ExecutionTimeouts::default() })); let report_context = test_report_context(); - // Stand in for earlier candidates having already consumed the request's - // complete first-byte budget. A per-candidate watchdog would wait a new - // 250 ms here; the shared absolute deadline must settle immediately. - let request_first_byte_started_at = Instant::now() - Duration::from_millis(300); - let result = tokio::time::timeout( - Duration::from_millis(100), - execute_stream_candidate_with_watchdog( - writer.as_ref(), - "trace_watchdog_shared_budget", - "claude_cli_stream", - &plan, - Some(&report_context), - request_first_byte_started_at, - false, - || { - std::future::pending::< - Result>, GatewayError>, - >() - }, - ), + let result = execute_stream_candidate_with_watchdog( + &writer, + "trace_watchdog_retry_budget", + "claude_cli_stream", + &plan, + Some(&report_context), + false, + || { + std::future::pending::< + Result>, GatewayError>, + >() + }, ) - .await - .expect("an expired request-level first-byte budget must not restart per candidate"); - + .await; assert!(matches!( result, Ok(StreamCandidateWatchdogOutcome::Executed( @@ -3231,14 +3201,119 @@ mod tests { } )) )); + + let mut next_plan = plan.clone(); + next_plan.candidate_id = Some("cand_watchdog_retry".to_string()); + next_plan.provider_id = provider_id.to_string(); + next_plan.key_id = key_id.to_string(); + next_plan.timeouts = Some(ExecutionTimeouts { + first_byte_ms: Some(first_byte_ms), + ..ExecutionTimeouts::default() + }); + let mut next_report_context = report_context.clone(); + next_report_context["candidate_id"] = json!("cand_watchdog_retry"); + next_report_context["candidate_index"] = json!(3); + + let result = execute_stream_candidate_with_watchdog( + &writer, + "trace_watchdog_retry_budget", + "claude_cli_stream", + &next_plan, + Some(&next_report_context), + false, + || async { + tokio::time::sleep(Duration::from_millis(60)).await; + Ok(AiAttemptExecutionOutcome::Responded(Response::new( + Body::from("retry succeeded"), + ))) + }, + ) + .await; + assert!( + matches!( + result, + Ok(StreamCandidateWatchdogOutcome::Executed( + AiAttemptExecutionOutcome::Responded(_) + )) + ), + "candidate {provider_id}/{key_id} must receive its own {first_byte_ms} ms budget" + ); + let records = writer.records.lock().await; assert_eq!(records.len(), 1); + assert_eq!(records[0].id, plan.candidate_id.as_deref().unwrap()); + assert_eq!(records[0].status, RequestCandidateStatus::Failed); assert_eq!( records[0].error_type.as_deref(), Some("local_stream_candidate_watchdog_timeout") ); } + #[tokio::test] + async fn stream_candidate_watchdog_failover_gets_fresh_first_byte_budget() { + for first_byte_ms in [100, 75, 150] { + assert_stream_candidate_retry_gets_fresh_first_byte_budget( + "provider_next", + "key_next", + first_byte_ms, + ) + .await; + } + } + + #[tokio::test] + async fn stream_candidate_watchdog_same_provider_retries_get_fresh_first_byte_budget() { + for key_id in ["key_next", "key_id"] { + assert_stream_candidate_retry_gets_fresh_first_byte_budget("provider_id", key_id, 100) + .await; + } + } + + #[tokio::test] + async fn stream_candidate_watchdog_starts_first_byte_budget_after_admission() { + let writer = TestRequestCandidateWriter::with_upstream_gate(1, Duration::from_secs(1)); + let held_permit = writer + .upstream_gate + .as_ref() + .expect("test gate should exist") + .try_acquire() + .expect("test gate permit should acquire"); + let plan = test_plan(Some(ExecutionTimeouts { + first_byte_ms: Some(50), + ..ExecutionTimeouts::default() + })); + let report_context = test_report_context(); + + let (result, ()) = tokio::join!( + execute_stream_candidate_with_watchdog( + &writer, + "trace_watchdog_admission_budget", + "claude_cli_stream", + &plan, + Some(&report_context), + false, + || async { + tokio::time::sleep(Duration::from_millis(20)).await; + Ok(AiAttemptExecutionOutcome::Responded(Response::new( + Body::from("admitted candidate succeeded"), + ))) + }, + ), + async move { + tokio::time::sleep(Duration::from_millis(100)).await; + drop(held_permit); + }, + ); + + assert!(matches!( + result, + Ok(StreamCandidateWatchdogOutcome::Executed( + AiAttemptExecutionOutcome::Responded(_) + )) + )); + assert!(writer.records.lock().await.is_empty()); + } + #[tokio::test] async fn stream_candidate_watchdog_can_stop_on_transport_error() { let writer = Arc::new(TestRequestCandidateWriter::default()); @@ -3254,7 +3329,6 @@ mod tests { "claude_cli_stream", &plan, Some(&report_context), - Instant::now(), true, || { std::future::pending::< @@ -3292,7 +3366,6 @@ mod tests { "claude_cli_stream", &plan, Some(&report_context), - Instant::now(), true, || async { mark_stream_candidate_watchdog_terminal_started(); @@ -3325,7 +3398,6 @@ mod tests { "claude_cli_stream", &plan, Some(&report_context), - Instant::now(), true, || async { Err(GatewayError::UpstreamUnavailable { @@ -3365,7 +3437,6 @@ mod tests { "claude_cli_stream", &plan, Some(&report_context), - Instant::now(), false, || async { panic!("execute future should not run while upstream execution gate is saturated") @@ -3413,7 +3484,6 @@ mod tests { "claude_cli_stream", &plan, Some(&report_context), - Instant::now(), false, || async { Err(GatewayError::AdmissionTimeout {