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
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
11 changes: 11 additions & 0 deletions crates/broker/src/node_control.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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())
}

Comment thread
coderabbitai[bot] marked this conversation as resolved.
/// Seed Relaycast's cumulative cursor after identity authority is bound.
///
/// The immutable `agent_id` is the key: a later agent reusing the same name
Expand Down
94 changes: 78 additions & 16 deletions crates/broker/src/runtime/api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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),
Expand All @@ -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),
Expand All @@ -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 '{}': {}",
Expand Down
167 changes: 154 additions & 13 deletions crates/broker/src/runtime/fleet.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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},
};
Expand Down Expand Up @@ -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<FleetControlCommand>,
fleet_delivery_book: &FleetDeliveryBook,
name: &WorkerName,
) -> Result<bool, String> {
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<FleetControlCommand>,
fleet_inventory: &HashMap<WorkerName, InventoryAgent>,
) {
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(
Expand Down Expand Up @@ -1499,6 +1540,106 @@ mod tests {
);
}

#[tokio::test]
async fn release_queues_fleet_deregister_with_authoritative_identity() {
let (tx, mut rx) = mpsc::channel::<FleetControlCommand>(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::<FleetControlCommand>(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::<FleetControlCommand>(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::<FleetControlCommand>(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::<FleetControlCommand>(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::<FleetControlCommand>(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
Expand Down
Loading
Loading