diff --git a/crates/cli/src/daemon/broker/server.rs b/crates/cli/src/daemon/broker/server.rs index 694780721..411a15250 100644 --- a/crates/cli/src/daemon/broker/server.rs +++ b/crates/cli/src/daemon/broker/server.rs @@ -772,7 +772,11 @@ fn cancel_activation_blocking( Ok(outcome) => { if outcome == crate::daemon::common::control::ActivationCancellation::Cancelled { revoke_activation(&state, &request.payload.activation_id); - cancel_staged_activation(&state, &request.payload.activation_id); + cancel_staged_activation( + &state, + &request.payload.activation_id, + "activation_cancelled", + ); } state.sockets.changed.notify_waiters(); Json(outcome).into_response() @@ -815,7 +819,7 @@ fn activation_failed_blocking( { Ok(()) => { revoke_activation(&state, &request.payload.activation_id); - cancel_staged_activation(&state, &request.payload.activation_id); + cancel_staged_activation(&state, &request.payload.activation_id, "activation_failed"); let fingerprint = authenticated.fingerprint.to_string(); log::error!( target: "nemo_relay.daemon", @@ -1761,15 +1765,38 @@ async fn forward_to_worker( ) .await; let route_failure = take_worker_route_failure(&mut outcome.response); + if matches!(outcome.failure, Some(ForwardFailure::ResponseHeadTimeout)) && !route_failure { + // Middleware and provider work can delay headers while worker control stays healthy. + // A request deadline alone does not establish loss of the assigned worker. + let fingerprint = fingerprint.to_string(); + log::warn!( + target: "nemo_relay.daemon", + event = "worker_request_failed", + fingerprint = fingerprint.as_str(), + worker_id = worker_id.as_str(), + reason = "response_head_timeout", + route_mode = "worker"; + "Worker request timed out waiting for response headers" + ); + return outcome.response; + } if outcome.failure.is_some() || route_failure { - handle_worker_communication_failure(&state, fingerprint, &worker_id); + let reason = outcome + .failure + .map_or("worker_route_rejected", ForwardFailure::as_str); + handle_worker_communication_failure(&state, fingerprint, &worker_id, reason); return outcome.response; } let (parts, body) = outcome.response.into_parts(); let observed = ErrorObservedBody { body, on_error: Some(move || { - handle_worker_communication_failure(&state, fingerprint, &worker_id); + handle_worker_communication_failure( + &state, + fingerprint, + &worker_id, + "response_body_error", + ); }), }; Response::from_parts(parts, Body::new(observed)) @@ -1962,6 +1989,7 @@ fn handle_worker_communication_failure( state: &Arc, fingerprint: Fingerprint, worker_id: &str, + reason: &'static str, ) { if state .registry @@ -1971,7 +1999,9 @@ fn handle_worker_communication_failure( return; } // Force control reconnect, retaining the assigned generation during its 30-second grace. - state.sockets.cancel_worker(worker_id); + state + .sockets + .cancel_worker(worker_id, "worker_communication_failed"); state.sockets.changed.notify_waiters(); let fingerprint = fingerprint.to_string(); log::error!( @@ -1979,6 +2009,7 @@ fn handle_worker_communication_failure( event = "worker_communication_failed", fingerprint = fingerprint.as_str(), worker_id = worker_id, + reason = reason, route_mode = "pass_through"; "Worker communication failed; route changed to pass-through" ); @@ -2180,7 +2211,7 @@ fn revoke_activation(state: &DaemonState, activation_id: &str) { } /// Release terminal activations immediately instead of retaining staged admission slots. -fn cancel_staged_activation(state: &Arc, activation_id: &str) { +fn cancel_staged_activation(state: &Arc, activation_id: &str, reason: &'static str) { let workers = { let mut sessions = lock(&state.worker_sessions); let ids: Vec<_> = sessions @@ -2196,7 +2227,7 @@ fn cancel_staged_activation(state: &Arc, activation_id: &str) { .collect::>() }; for (id, session) in workers { - state.sockets.cancel_worker(&id); + state.sockets.cancel_worker(&id, reason); tokio::spawn(socket::cleanup_worker_session( Arc::clone(state), id, @@ -2213,7 +2244,7 @@ fn expire_activation_routes(state: &Arc, now_unix_ms: u64) { } in state.registry.expire_activations(now_unix_ms) { revoke_activation(state, &activation_id); - cancel_staged_activation(state, &activation_id); + cancel_staged_activation(state, &activation_id, "activation_expired"); state.sockets.changed.notify_waiters(); let fingerprint = fingerprint.to_string(); log::error!( @@ -2231,7 +2262,7 @@ fn handle_release_action(state: Arc, fingerprint: Fingerprint, acti ReleaseAction::NoChange => {} ReleaseAction::CancelActivation { activation_id } => { revoke_activation(&state, &activation_id); - cancel_staged_activation(&state, &activation_id); + cancel_staged_activation(&state, &activation_id, "activation_cancelled"); let fingerprint = fingerprint.to_string(); log::info!( target: "nemo_relay.daemon", diff --git a/crates/cli/src/daemon/broker/server/socket.rs b/crates/cli/src/daemon/broker/server/socket.rs index e73fa8453..9282bd305 100644 --- a/crates/cli/src/daemon/broker/server/socket.rs +++ b/crates/cli/src/daemon/broker/server/socket.rs @@ -157,9 +157,9 @@ impl Hub { peer.cancel.cancel("outbound_queue_full"); } } - pub(super) fn cancel_worker(&self, worker_id: &str) { + pub(super) fn cancel_worker(&self, worker_id: &str, reason: &'static str) { if let Some(peer) = lock(&self.peers).get(&key(ComponentRole::Worker, worker_id)) { - peer.cancel.cancel("activation_expired"); + peer.cancel.cancel(reason); } } pub(super) fn draining(&self, worker_id: &str) -> bool { @@ -765,7 +765,7 @@ async fn disconnected( work.registry.mcp_disconnected(fingerprint, &session_id) { revoke_activation(&work, &activation_id); - cancel_staged_activation(&work, &activation_id); + cancel_staged_activation(&work, &activation_id, "activation_cancelled"); lock(&work.pending_directives).remove(&disconnected_id); } }) diff --git a/crates/cli/tests/coverage/daemon/server_tests.rs b/crates/cli/tests/coverage/daemon/server_tests.rs index 32da1e007..ab67a2135 100644 --- a/crates/cli/tests/coverage/daemon/server_tests.rs +++ b/crates/cli/tests/coverage/daemon/server_tests.rs @@ -841,6 +841,95 @@ async fn global_pass_through_authenticates_and_forwards_only_provider_headers() server.abort(); } +#[tokio::test] +async fn worker_response_head_timeout_preserves_the_route_and_next_request() { + for require_worker in [false, true] { + let token = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode([0x75_u8; 32]); + let credential = RouteCredential::parse(token.clone()).unwrap(); + let mut state = test_daemon_state(false, &token, GatewayConfig::default()); + Arc::get_mut(&mut state).unwrap().registry = + Registry::new(false).with_require_worker(require_worker); + let fingerprint = MachineIdentity::generate().unwrap().identity.fingerprint(); + let launch = fresh_launch(WorkerNetworkHint::new("127.0.0.1", None).unwrap()).unwrap(); + let activation_id = launch.activation_id.clone(); + state + .registry + .register_mcp( + McpRegistration { + fingerprint, + token_digest: credential.digest(), + session_id: McpSessionId::new("slow-worker-mcp").unwrap(), + lease_expires_at_unix_ms: u64::MAX, + }, + launch, + ) + .unwrap(); + + let started = Arc::new(tokio::sync::Notify::new()); + let first_request = Arc::new(std::sync::atomic::AtomicBool::new(true)); + let worker_started = Arc::clone(&started); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let endpoint = format!("http://{}", listener.local_addr().unwrap()); + let worker_task = tokio::spawn(async move { + let app = Router::new().fallback(move || { + let first_request = Arc::clone(&first_request); + let started = Arc::clone(&worker_started); + async move { + if first_request.swap(false, std::sync::atomic::Ordering::SeqCst) { + started.notify_one(); + std::future::pending::<()>().await; + } + StatusCode::NO_CONTENT + } + }); + axum::serve(listener, app).await.unwrap(); + }); + state + .registry + .mark_worker_ready( + fingerprint, + &activation_id, + Arc::new( + WorkerTarget::new( + "slow-worker", + endpoint, + SensitiveString::new("worker-token").unwrap(), + ) + .unwrap(), + ), + ) + .unwrap(); + let app = router(Arc::clone(&state)); + let request = || { + Request::post("/v1/responses") + .header(CLIENT_TOKEN_HEADER, &token) + .body(Body::empty()) + .unwrap() + }; + let response = tokio::spawn(app.clone().oneshot(request())); + tokio::time::timeout(Duration::from_secs(5), started.notified()) + .await + .unwrap(); + tokio::time::pause(); + tokio::time::advance(RESPONSE_HEAD_TIMEOUT).await; + let response = response.await.unwrap().unwrap(); + tokio::time::resume(); + let route_ready = matches!( + state.registry.resolve_target(&credential.digest()), + Ok(ResolvedTarget::Worker(_)) + ); + let next_status = if route_ready { + Some(app.oneshot(request()).await.unwrap().status()) + } else { + None + }; + worker_task.abort(); + assert_eq!(response.status(), StatusCode::GATEWAY_TIMEOUT); + assert!(route_ready); + assert_eq!(next_status, Some(StatusCode::NO_CONTENT)); + } +} + #[tokio::test] async fn unreachable_worker_marks_its_authenticated_route_pass_through() { let token = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode([0x22_u8; 32]); @@ -2593,7 +2682,12 @@ fn communication_failure_invalidates_route_without_waiting_for_durable_revocatio let (sent, received) = std::sync::mpsc::channel(); let failure_state = Arc::clone(&state); let task = runtime.spawn_blocking(move || { - handle_worker_communication_failure(&failure_state, fingerprint, "failed-worker"); + handle_worker_communication_failure( + &failure_state, + fingerprint, + "failed-worker", + "transport_error", + ); sent.send(()).unwrap(); }); let completed = received.recv_timeout(Duration::from_secs(5)); diff --git a/docs/daemon/operations.mdx b/docs/daemon/operations.mdx index a0471ff48..26b787bfb 100644 --- a/docs/daemon/operations.mdx +++ b/docs/daemon/operations.mdx @@ -160,7 +160,7 @@ logging. Failure events use warning or error levels. | Worker control | `worker_control_connected`, `worker_control_reconnected`, `worker_control_disconnected`, `worker_control_lost`, `worker_control_recovery_failed` | | Worker shutdown | `worker_drain_started`, `worker_removed`, `worker_stopped`, `worker_activation_cancelled` | | Worker degradation | `worker_registration_failed`, `worker_readiness_failed`, `worker_activation_failed`, `worker_activation_expired`, `worker_communication_failed`, `worker_failed`, `worker_recovery_failed`, `worker_recovery_expired`, `worker_relaunch_failed` | -| Provider forwarding | `upstream_request_failed` | +| Request forwarding | `upstream_request_failed`, `worker_request_failed` | Use `fingerprint` to match an MCP session to its worker route. Use `worker_id` to follow one worker generation. A `route_mode = "pass_through"` field reports @@ -171,8 +171,15 @@ Both `worker_launch_failed` and `worker_activation_failed` include a such as starting the worker, sending its startup permission, or waiting for it to be ready. It does not include activation data or the source error text. Read the worker's configured logs for the full error message. +`worker_communication_failed` includes a `reason` that identifies a transport +error, invalid destination or response, response-body error, or explicit worker +route rejection. `worker_request_failed` with `reason = "response_head_timeout"` +reports a timed-out request without disconnecting the worker or changing its +route. Check worker and provider logs to identify the delay. Control disconnect events include a `reason` such as `acknowledgement_timeout`, -`pong_timeout`, or `peer_closed`. Provider failure events are rate-limited and +`pong_timeout`, `peer_closed`, or `worker_communication_failed`. A staged worker +cancelled when its launcher exits reports `activation_cancelled`; this does not +mean that the startup deadline expired. Provider failure events are rate-limited and report the number of intervening failures in `suppressed_since_last_emit`. ## Diagnose Common Failures diff --git a/docs/daemon/reference.mdx b/docs/daemon/reference.mdx index 541b7e49c..cdd5b6990 100644 --- a/docs/daemon/reference.mdx +++ b/docs/daemon/reference.mdx @@ -439,6 +439,12 @@ disconnects, Relay cancels the corresponding provider request. Relay has separat deadlines for connections and response headers. It does not set a total time limit after streaming starts. +A request that exceeds the daemon's 60-second response-header deadline returns +`504`. This timeout does not disconnect the worker or change its route. The +worker can still be waiting for middleware or a provider. A worker transport +failure, response-body failure, or explicit route rejection triggers control +recovery before worker forwarding resumes. + The worker also sends request bodies directly through the pooled Hyper client when no LLM request guardrail, request interceptor, or request-sanitization guardrail is registered. If middleware needs the complete JSON request, the