mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-04 00:17:45 +08:00
fix(gateway): reset stream first-byte timeout per candidate
This commit is contained in:
@@ -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<Fut>(
|
||||
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<StreamCandidateWatchdogOutcome, GatewayError>
|
||||
@@ -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<AiAttemptExecutionOutcome<Response<Body>>, 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<AiAttemptExecutionOutcome<Response<Body>>, 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 {
|
||||
|
||||
Reference in New Issue
Block a user