Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
49 changes: 40 additions & 9 deletions crates/cli/src/daemon/broker/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -1962,6 +1989,7 @@ fn handle_worker_communication_failure(
state: &Arc<DaemonState>,
fingerprint: Fingerprint,
worker_id: &str,
reason: &'static str,
) {
if state
.registry
Expand All @@ -1971,14 +1999,17 @@ 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!(
target: "nemo_relay.daemon",
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"
);
Expand Down Expand Up @@ -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<DaemonState>, activation_id: &str) {
fn cancel_staged_activation(state: &Arc<DaemonState>, activation_id: &str, reason: &'static str) {
let workers = {
let mut sessions = lock(&state.worker_sessions);
let ids: Vec<_> = sessions
Expand All @@ -2196,7 +2227,7 @@ fn cancel_staged_activation(state: &Arc<DaemonState>, activation_id: &str) {
.collect::<Vec<_>>()
};
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,
Expand All @@ -2213,7 +2244,7 @@ fn expire_activation_routes(state: &Arc<DaemonState>, 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!(
Expand All @@ -2231,7 +2262,7 @@ fn handle_release_action(state: Arc<DaemonState>, 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",
Expand Down
6 changes: 3 additions & 3 deletions crates/cli/src/daemon/broker/server/socket.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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);
}
})
Expand Down
96 changes: 95 additions & 1 deletion crates/cli/tests/coverage/daemon/server_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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]);
Expand Down Expand Up @@ -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));
Expand Down
11 changes: 9 additions & 2 deletions docs/daemon/operations.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down
6 changes: 6 additions & 0 deletions docs/daemon/reference.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading