diff --git a/CHANGELOG.md b/CHANGELOG.md index d9706922b..6dfabe8e1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -17,6 +17,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Feature verification catalog now records the exact CLI and MCP surfaces, adds previously unlisted SDK and plugin integrations, and maps every category to an end-to-end procedure with prerequisites, assertions, cleanup, and automation limits. +### Fixed + +- `agent-relay-broker` now orders fleet `agent.deregister` before acknowledging a local worker release, so a restarted node can immediately recover the same agent name without an active-location collision, and fails fast instead of stalling the runtime API when fleet-control delivery is backpressured. + ## [10.6.6] - 2026-07-19 ### Fixed diff --git a/crates/broker/src/node_control.rs b/crates/broker/src/node_control.rs index 21366d9db..ee7477139 100644 --- a/crates/broker/src/node_control.rs +++ b/crates/broker/src/node_control.rs @@ -681,6 +681,17 @@ impl FleetDeliveryBook { self.bind_identity(&agent, &agent_id, true); } + /// Return the immutable identity currently bound to an agent name. + /// + /// Release paths need this before pruning the delivery book so they can + /// send an ordered `agent.deregister` frame to the fleet control plane. + pub(crate) fn active_agent_id(&self, agent: &str) -> Option<&str> { + self.active_agent_bindings_by_name + .get(agent) + .filter(|binding| binding.authoritative) + .map(|binding| binding.agent_id.as_str()) + } + /// Seed Relaycast's cumulative cursor after identity authority is bound. /// /// The immutable `agent_id` is the key: a later agent reusing the same name diff --git a/crates/broker/src/runtime/api.rs b/crates/broker/src/runtime/api.rs index 80beef668..c5b826e13 100644 --- a/crates/broker/src/runtime/api.rs +++ b/crates/broker/src/runtime/api.rs @@ -757,6 +757,20 @@ impl BrokerRuntime { workers.metrics.on_release(&name); match workers.release(&name).await { Ok(()) => { + let fleet_deregistration_error = super::fleet::deregister_fleet_agent( + fleet_control_tx, + fleet_delivery_book, + &name, + ) + .await + .err(); + if let Some(error) = &fleet_deregistration_error { + tracing::warn!( + worker = %name, + error = %error, + "released worker fleet deregistration was not queued; retaining its identity for retry" + ); + } if let Err(error) = relaycast_http.mark_agent_offline(&name).await { tracing::warn!( worker = %name, @@ -786,13 +800,26 @@ impl BrokerRuntime { if paths.persist { let _ = state.save(&paths.state); } - super::fleet::prune_fleet_agent_state( - fleet_control_tx, - fleet_inventory, - fleet_delivery_book, - &name, - ) - .await; + if fleet_deregistration_error.is_none() { + super::fleet::prune_fleet_agent_state( + fleet_control_tx, + fleet_inventory, + fleet_delivery_book, + &name, + ) + .await; + } else { + // Do not advertise the gone worker in the next + // inventory sync. Keep only its authoritative + // delivery-book binding so an idempotent release + // retry can emit the missing deregistration. + super::fleet::prune_fleet_inventory_entry( + fleet_control_tx, + fleet_inventory, + &name, + ) + .await; + } super::fleet::publish_fleet_load_snapshot( fleet_control_tx, u32::try_from(workers.workers.len()).unwrap_or(u32::MAX), @@ -811,23 +838,52 @@ impl BrokerRuntime { Some("http_api_release"), ) .await; - let _ = reply.send(Ok(json!({ "success": true, "name": name }))); + let response = match fleet_deregistration_error { + Some(error) => Err(format!( + "failed to deregister released worker from fleet control: {error}" + )), + None => Ok(json!({ "success": true, "name": name })), + }; + let _ = reply.send(response); } Err(e) => { let message = e.to_string(); if is_unknown_worker_error_message(&message) { + let fleet_deregistration_error = super::fleet::deregister_fleet_agent( + fleet_control_tx, + fleet_delivery_book, + &name, + ) + .await + .err(); + if let Some(error) = &fleet_deregistration_error { + tracing::warn!( + worker = %name, + error = %error, + "already-exited worker fleet deregistration was not queued; retaining its identity for retry" + ); + } relaycast_http.forget_agent_registration(&name); state.agents.remove(&name); if paths.persist { let _ = state.save(&paths.state); } - super::fleet::prune_fleet_agent_state( - fleet_control_tx, - fleet_inventory, - fleet_delivery_book, - &name, - ) - .await; + if fleet_deregistration_error.is_none() { + super::fleet::prune_fleet_agent_state( + fleet_control_tx, + fleet_inventory, + fleet_delivery_book, + &name, + ) + .await; + } else { + super::fleet::prune_fleet_inventory_entry( + fleet_control_tx, + fleet_inventory, + &name, + ) + .await; + } super::fleet::publish_fleet_load_snapshot( fleet_control_tx, u32::try_from(workers.workers.len()).unwrap_or(u32::MAX), @@ -840,7 +896,13 @@ impl BrokerRuntime { worker = %name, "ignoring duplicate HTTP API release for already exited worker" ); - let _ = reply.send(Ok(json!({ "success": true, "name": name }))); + let response = match fleet_deregistration_error { + Some(error) => Err(format!( + "failed to deregister released worker from fleet control: {error}" + )), + None => Ok(json!({ "success": true, "name": name })), + }; + let _ = reply.send(response); } else { eprintln!( "[agent-relay] HTTP API: failed to release '{}': {}", diff --git a/crates/broker/src/runtime/fleet.rs b/crates/broker/src/runtime/fleet.rs index 2d6466e3d..1195ed473 100644 --- a/crates/broker/src/runtime/fleet.rs +++ b/crates/broker/src/runtime/fleet.rs @@ -2,8 +2,8 @@ use super::*; use crate::{ fleet_wire::{ ActionInvoke, ActionResult, ActionResultError, ActionResultOutput, ActionResultPayload, - AgentRegister, BrokerToRelaycast, Deliver, DeliveryMode, RelaycastToBroker, - FLEET_WIRE_VERSION, + AgentDeregister, AgentRegister, BrokerToRelaycast, Deliver, DeliveryMode, + RelaycastToBroker, FLEET_WIRE_VERSION, }, node_control::{delivery_ack, handler_unavailable_result, DeliveryDecision}, }; @@ -693,29 +693,70 @@ pub(super) async fn publish_fleet_load_snapshot( handlers_live: bool, heartbeat_now: bool, ) { - let _ = fleet_control_tx - .send(FleetControlCommand::UpdateLoad(FleetLoadSnapshot { + if let Err(error) = + fleet_control_tx.try_send(FleetControlCommand::UpdateLoad(FleetLoadSnapshot { active_agents, max_agents, handlers_live, })) - .await; + { + tracing::warn!(error = %error, "fleet load update queue is unavailable; periodic heartbeat will retry"); + } if heartbeat_now { - let _ = fleet_control_tx - .send(FleetControlCommand::HeartbeatNow) - .await; + if let Err(error) = fleet_control_tx.try_send(FleetControlCommand::HeartbeatNow) { + tracing::warn!(error = %error, "fleet heartbeat queue is unavailable; periodic heartbeat will retry"); + } } } +/// Queue an `agent.deregister` frame before a released name can be reused. +/// +/// HTTP release and a subsequent same-name spawn are separate broker API +/// requests, but both converge on this single FIFO fleet-control channel. A +/// successful synchronous enqueue here precedes the release reply, so the +/// control plane observes deregistration before any later `agent.register`, +/// including when a restarted broker has a new node id. Backpressure fails the +/// release promptly and retains the authoritative identity for retry instead +/// of blocking the broker's single runtime API actor. Agents registered only +/// through the legacy HTTP fallback have no authoritative fleet identity and +/// remain covered by the REST offline call. +pub(super) async fn deregister_fleet_agent( + fleet_control_tx: &mpsc::Sender, + fleet_delivery_book: &FleetDeliveryBook, + name: &WorkerName, +) -> Result { + let Some(agent_id) = fleet_delivery_book.active_agent_id(name.as_str()) else { + return Ok(false); + }; + fleet_control_tx + .try_send(FleetControlCommand::Send( + BrokerToRelaycast::AgentDeregister(AgentDeregister { + v: FLEET_WIRE_VERSION, + id: None, + agent_id: agent_id.to_string(), + name: Some(name.as_str().to_string()), + }), + )) + .map_err(|error| match error { + tokio::sync::mpsc::error::TrySendError::Full(_) => { + "fleet_control_backpressure".to_string() + } + tokio::sync::mpsc::error::TrySendError::Closed(_) => { + "fleet_control_unavailable".to_string() + } + })?; + Ok(true) +} + pub(super) async fn publish_fleet_inventory_snapshot( fleet_control_tx: &mpsc::Sender, fleet_inventory: &HashMap, ) { - let _ = fleet_control_tx - .send(FleetControlCommand::UpdateInventory( - fleet_inventory.values().cloned().collect(), - )) - .await; + if let Err(error) = fleet_control_tx.try_send(FleetControlCommand::UpdateInventory( + fleet_inventory.values().cloned().collect(), + )) { + tracing::warn!(error = %error, "fleet inventory queue is unavailable; periodic heartbeat will retry"); + } } pub(super) async fn refresh_fleet_inventory_session_ref( @@ -1499,6 +1540,106 @@ mod tests { ); } + #[tokio::test] + async fn release_queues_fleet_deregister_with_authoritative_identity() { + let (tx, mut rx) = mpsc::channel::(1); + let mut delivery_book = FleetDeliveryBook::default(); + delivery_book.bind_authoritative_identity("agent-a", "agent-a-id"); + + assert!( + deregister_fleet_agent(&tx, &delivery_book, &WorkerName::from("agent-a")) + .await + .expect("deregister should enqueue") + ); + + let command = rx.recv().await.expect("deregister command emitted"); + let FleetControlCommand::Send(BrokerToRelaycast::AgentDeregister(request)) = command else { + panic!("expected AgentDeregister command"); + }; + assert_eq!(request.agent_id, "agent-a-id"); + assert_eq!(request.name.as_deref(), Some("agent-a")); + } + + #[tokio::test] + async fn release_without_fleet_identity_does_not_emit_deregister() { + let (tx, mut rx) = mpsc::channel::(1); + let delivery_book = FleetDeliveryBook::default(); + + assert!( + !deregister_fleet_agent(&tx, &delivery_book, &WorkerName::from("http-only")) + .await + .expect("missing identity should be a no-op") + ); + assert!(rx.try_recv().is_err()); + } + + #[tokio::test] + async fn release_does_not_deregister_nonauthoritative_http_identity() { + let (tx, mut rx) = mpsc::channel::(1); + let mut delivery_book = FleetDeliveryBook::default(); + let delivery = test_deliver( + "http-only", + "delivery-http-only", + "message-http-only", + json!({"text": "legacy delivery"}), + ); + delivery_book.commit_received(&delivery); + + assert!( + !deregister_fleet_agent(&tx, &delivery_book, &WorkerName::from("http-only")) + .await + .expect("non-authoritative identity should be a no-op") + ); + assert!(rx.try_recv().is_err()); + } + + #[tokio::test] + async fn failed_release_deregister_retains_identity_for_retry() { + let (tx, rx) = mpsc::channel::(1); + drop(rx); + let mut delivery_book = FleetDeliveryBook::default(); + delivery_book.bind_authoritative_identity("agent-a", "agent-a-id"); + + assert_eq!( + deregister_fleet_agent(&tx, &delivery_book, &WorkerName::from("agent-a")).await, + Err("fleet_control_unavailable".to_string()) + ); + assert_eq!(delivery_book.active_agent_id("agent-a"), Some("agent-a-id")); + } + + #[tokio::test] + async fn release_deregister_fails_fast_when_fleet_control_is_backpressured() { + let (tx, _rx) = mpsc::channel::(1); + tx.try_send(FleetControlCommand::HeartbeatNow) + .expect("fill fleet control queue"); + let mut delivery_book = FleetDeliveryBook::default(); + delivery_book.bind_authoritative_identity("agent-a", "agent-a-id"); + + let result = tokio::time::timeout( + Duration::from_millis(50), + deregister_fleet_agent(&tx, &delivery_book, &WorkerName::from("agent-a")), + ) + .await + .expect("backpressured deregister must not stall the runtime API actor"); + + assert_eq!(result, Err("fleet_control_backpressure".to_string())); + assert_eq!(delivery_book.active_agent_id("agent-a"), Some("agent-a-id")); + } + + #[tokio::test] + async fn load_publication_does_not_wait_for_fleet_control_capacity() { + let (tx, _rx) = mpsc::channel::(1); + tx.try_send(FleetControlCommand::HeartbeatNow) + .expect("fill fleet control queue"); + + tokio::time::timeout( + Duration::from_millis(50), + publish_fleet_load_snapshot(&tx, 1, 4, true, true), + ) + .await + .expect("load publication must not stall the runtime API actor"); + } + #[test] fn relaycast_spawn_session_ref_is_none_without_harness_session() { // A spawn with no harnessConfig session id yields None — the spawn is a diff --git a/crates/broker/src/worker.rs b/crates/broker/src/worker.rs index fafde5f06..ff300bdaf 100644 --- a/crates/broker/src/worker.rs +++ b/crates/broker/src/worker.rs @@ -907,6 +907,11 @@ impl WorkerRegistry { pub(crate) async fn release(&mut self, name: &str) -> Result<()> { tracing::info!(target = "broker::release", name = %name, "releasing worker"); self.initial_tasks.remove(name); + // An explicit release is terminal even when the process already exited + // and disappeared from `workers`. Cancel any pending restart before + // looking up the handle so maintenance cannot resurrect the released + // name after the API has acknowledged teardown. + self.supervisor.unregister(name); let mut handle = self .workers .remove(name) @@ -1668,6 +1673,57 @@ mod tests { assert!(!reg.has_worker_in_workspace("nonexistent", &workspace)); } + #[tokio::test] + async fn release_cancels_pending_restart_for_an_already_exited_worker() { + let mut reg = make_registry(vec![]); + let name = "released-stale-worker"; + let restart_policy = crate::supervisor::RestartPolicy { + cooldown_ms: 0, + ..crate::supervisor::RestartPolicy::default() + }; + let spec = AgentSpec { + name: WorkerName::from(name), + runtime: AgentRuntime::Headless, + provider: None, + cli: None, + session_id: None, + harness_config: None, + model: None, + cwd: None, + team: None, + shadow_of: None, + shadow_mode: None, + args: Vec::new(), + channels: Vec::new(), + restart_policy: Some(restart_policy.clone()), + }; + reg.supervisor.register( + name, + crate::supervisor::SupervisedAgent { + spec, + parent: None, + initial_task: None, + skip_relay_prompt: false, + agent_result: None, + }, + restart_policy, + ); + assert!(reg.supervisor.is_supervised(name)); + assert!(matches!( + reg.supervisor.on_exit(name, Some(1), None), + Some(crate::supervisor::RestartDecision::Restart { .. }) + )); + assert!(!reg.supervisor.pending_restarts().is_empty()); + + let error = reg + .release(name) + .await + .expect_err("missing process is still reported"); + assert!(error.to_string().contains("unknown worker")); + assert!(!reg.supervisor.is_supervised(name)); + assert!(reg.supervisor.pending_restarts().is_empty()); + } + #[test] fn worker_log_path_rejects_path_traversal() { let reg = make_registry(vec![]);