Skip to content
Merged
6 changes: 5 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,11 @@ All notable changes to Agent Relay will be documented in this file.
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/),
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

## [Unreleased]
## [Unreleased - Minor]

### Added

- Opt-in persistent broker task providers preserve final results across reconnects and acknowledge callbacks only after durable Relaycast receipts.

## [12.2.2] - 2026-09-15

Expand Down
121 changes: 121 additions & 0 deletions crates/broker/src/fleet_wire.rs
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,8 @@ pub struct FleetCapability {
skip_serializing_if = "Option::is_none"
)]
pub queue: Option<bool>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub execution_mode: Option<String>,
#[serde(
default,
deserialize_with = "deserialize_optional_presence",
Expand Down Expand Up @@ -375,6 +377,7 @@ pub struct ActionResult {
pub id: Option<String>,
pub invocation_id: String,
pub result: ActionResultPayload,
pub task: Option<TaskResultFields>,
}

impl Serialize for ActionResult {
Expand Down Expand Up @@ -420,6 +423,14 @@ struct ActionResultWire {
skip_serializing_if = "Option::is_none"
)]
pub error: Option<String>,
#[serde(default, rename = "final", skip_serializing_if = "Option::is_none")]
pub final_result: Option<bool>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub execution_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub worker_generation: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub accounting: Option<BTreeMap<String, serde_json::Number>>,
}

fn deserialize_optional_presence<'de, D, T>(deserializer: D) -> Result<Option<T>, D::Error>
Expand Down Expand Up @@ -483,13 +494,27 @@ impl From<&ActionResult> for ActionResultWire {
v: value.v,
id: value.id.clone(),
invocation_id: value.invocation_id.clone(),
final_result: value.task.as_ref().map(|task| task.final_result),
execution_id: value.task.as_ref().map(|task| task.execution_id.clone()),
worker_generation: value
.task
.as_ref()
.map(|task| task.worker_generation.clone()),
accounting: value.task.as_ref().and_then(|task| task.accounting.clone()),
output: Some(output.output.clone()),
error: None,
},
ActionResultPayload::Error(error) => Self {
v: value.v,
id: value.id.clone(),
invocation_id: value.invocation_id.clone(),
final_result: value.task.as_ref().map(|task| task.final_result),
execution_id: value.task.as_ref().map(|task| task.execution_id.clone()),
worker_generation: value
.task
.as_ref()
.map(|task| task.worker_generation.clone()),
accounting: value.task.as_ref().and_then(|task| task.accounting.clone()),
output: None,
error: Some(error.error.clone()),
},
Expand All @@ -509,11 +534,44 @@ impl TryFrom<ActionResultWire> for ActionResult {
}
};

let task = if value.final_result.is_some()
|| value.execution_id.is_some()
|| value.worker_generation.is_some()
|| value.accounting.is_some()
{
if value.id.as_ref().is_none_or(|id| id.is_empty()) {
return Err("task result requires a request id".into());
}
let execution_id = value
.execution_id
.filter(|id| !id.is_empty())
.ok_or("task result requires execution_id")?;
let worker_generation = value
.worker_generation
.filter(|id| !id.is_empty() && id.len() <= 512)
.ok_or("task result requires worker_generation")?;
if value.accounting.as_ref().is_some_and(|values| {
values
.values()
.any(|n| n.as_f64().is_none_or(|n| !n.is_finite() || n < 0.0))
}) {
return Err("task accounting must be finite and nonnegative".into());
}
Some(TaskResultFields {
final_result: value.final_result.ok_or("task result requires final")?,
execution_id,
worker_generation,
accounting: value.accounting,
})
} else {
None
};
Ok(Self {
v: value.v,
id: value.id,
invocation_id: value.invocation_id,
result,
task,
})
}
}
Expand Down Expand Up @@ -605,6 +663,8 @@ pub struct ActionInvoke {
skip_serializing_if = "Option::is_none"
)]
pub agent_name: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub task_execution: Option<Box<TaskExecution>>,
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
Expand Down Expand Up @@ -720,6 +780,38 @@ where
}
}

// Nested inside the inbound ActionInvoke frame, so this must retain the same
// forward-compatibility rule: a future engine field cannot make the broker
// drop the entire invocation before acknowledging it.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct TaskExecution {
pub execution_id: String,
pub run_id: String,
pub step_id: String,
pub dispatch_id: String,
pub deadline: String,
}
Comment thread
cursor[bot] marked this conversation as resolved.

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct ActionAccept {
pub v: FleetWireVersion,
pub id: String,
pub invocation_id: String,
pub execution_id: String,
pub worker_generation: String,
}

#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct TaskResultFields {
#[serde(rename = "final")]
pub final_result: bool,
pub execution_id: String,
pub worker_generation: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub accounting: Option<BTreeMap<String, serde_json::Number>>,
}

#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "type")]
pub enum NodeToServer {
Expand All @@ -737,6 +829,8 @@ pub enum NodeToServer {
DeliveryAck(DeliveryAck),
#[serde(rename = "action.result")]
ActionResult(ActionResult),
#[serde(rename = "action.accept")]
ActionAccept(ActionAccept),
#[serde(rename = "inventory.sync")]
InventorySync(InventorySync),
}
Expand Down Expand Up @@ -863,6 +957,7 @@ mod tests {
#[test]
fn action_result_allows_error_payloads() {
let msg = BrokerToRelaycast::ActionResult(ActionResult {
task: None,
v: FLEET_WIRE_VERSION,
id: None,
invocation_id: "inv_2".to_string(),
Expand Down Expand Up @@ -993,6 +1088,7 @@ mod tests {
name: "builder-1".to_string(),
node_id: "node_1".to_string(),
capabilities: vec![FleetCapability {
execution_mode: None,
name: "spawn:codex".to_string(),
kind: Some("capacity".to_string()),
global: None,
Expand Down Expand Up @@ -1330,4 +1426,29 @@ mod tests {
let decoded: RelaycastToBroker = serde_json::from_value(value).unwrap();
assert_eq!(decoded, msg);
}

#[test]
fn action_invoke_accepts_future_nested_task_execution_fields() {
let invoke: RelaycastToBroker = serde_json::from_value(json!({
"type": "action.invoke",
"v": 1,
"invocation_id": "inv_task_1",
"action": "task.run",
"input": {},
"task_execution": {
"execution_id": "inv_task_1/1",
"run_id": "run_1",
"step_id": "step_1",
"dispatch_id": "dispatch_1",
"deadline": "2026-09-15T12:00:00.000Z",
"future_engine_field": { "value": 1 }
}
}))
.expect("inbound task execution metadata must be forward compatible");

let RelaycastToBroker::ActionInvoke(invoke) = invoke else {
panic!("action invoke")
};
assert_eq!(invoke.task_execution.unwrap().execution_id, "inv_task_1/1");
}
}
14 changes: 14 additions & 0 deletions crates/broker/src/listen_api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -280,12 +280,16 @@ impl std::error::Error for DeliveryRouteError {}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AgentResultRouteError {
InvalidToken,
Retryable,
Conflict,
}

impl std::fmt::Display for AgentResultRouteError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
AgentResultRouteError::InvalidToken => write!(f, "invalid_result_token"),
AgentResultRouteError::Retryable => write!(f, "task_receipt_pending"),
AgentResultRouteError::Conflict => write!(f, "task_result_conflict"),
}
}
}
Expand Down Expand Up @@ -1396,6 +1400,16 @@ async fn listen_api_agent_result(

match reply_rx.await {
Ok(Ok(value)) => (axum::http::StatusCode::OK, axum::Json(value)),
Ok(Err(AgentResultRouteError::Retryable)) => (
axum::http::StatusCode::SERVICE_UNAVAILABLE,
axum::Json(
json!({ "success": false, "error": "task_receipt_pending", "retryable": true }),
),
),
Ok(Err(AgentResultRouteError::Conflict)) => (
axum::http::StatusCode::CONFLICT,
axum::Json(json!({ "success": false, "error": "task_result_conflict" })),
),
Ok(Err(AgentResultRouteError::InvalidToken)) => (
axum::http::StatusCode::UNAUTHORIZED,
axum::Json(json!({ "success": false, "error": "invalid_result_token" })),
Expand Down
62 changes: 61 additions & 1 deletion crates/broker/src/node_control.rs
Original file line number Diff line number Diff line change
Expand Up @@ -519,6 +519,7 @@ impl FleetLoadSnapshot {
active_agent_names.sort();
active_agent_names.dedup();
capabilities.push(FleetCapability {
execution_mode: None,
name: LIVE_AGENT_CAPABILITY_NAME.to_string(),
kind: Some("capacity".to_string()),
global: None,
Expand Down Expand Up @@ -1306,6 +1307,7 @@ impl FleetDeliveryBook {

pub(crate) fn handler_unavailable_result(invocation_id: &str) -> ActionResult {
ActionResult {
task: None,
v: FLEET_WIRE_VERSION,
id: None,
invocation_id: invocation_id.to_string(),
Expand All @@ -1327,9 +1329,10 @@ pub(crate) fn build_node_register(
.iter()
.filter(|capability| capability.name != crate::fleet_wire::DELIVERY_CURSOR_CAPABILITY)
.map(|capability| FleetCapability {
execution_mode: (capability.name == "task.run").then(|| "task".to_owned()),
name: capability.name.clone(),
kind: capability.kind.clone(),
global: None,
global: (capability.name == "task.run").then_some(true),
queue: None,
metadata: capability.metadata.as_ref().map(|metadata| {
metadata
Expand All @@ -1340,6 +1343,7 @@ pub(crate) fn build_node_register(
})
.collect::<Vec<_>>();
capabilities.push(FleetCapability {
execution_mode: None,
name: crate::fleet_wire::DELIVERY_CURSOR_CAPABILITY.to_string(),
kind: Some("capacity".to_string()),
global: None,
Expand Down Expand Up @@ -2635,6 +2639,14 @@ where
}
match frame {
RelaycastToBroker::Reply(reply) => {
if reply.id.starts_with(crate::runtime::task_request_prefix()) {
return event_tx
.send(FleetControlEvent::Message(RelaycastToBroker::Reply(
reply,
)))
.await
.is_ok();
}
if let Some(pending) = pending_deregistrations.remove(&reply.id) {
let result = if reply.ok {
Ok(())
Expand Down Expand Up @@ -2677,6 +2689,14 @@ where
}
}
RelaycastToBroker::Error(error) => {
if error.id.starts_with(crate::runtime::task_request_prefix()) {
return event_tx
.send(FleetControlEvent::Message(RelaycastToBroker::Error(
error,
)))
.await
.is_ok();
}
if let Some(pending) = pending_deregistrations.remove(&error.id) {
let _ =
pending.send(Err(format!("{}: {}", error.code, error.message)));
Expand Down Expand Up @@ -4218,6 +4238,7 @@ mod tests {
assert_eq!(
register.capabilities.last(),
Some(&FleetCapability {
execution_mode: None,
name: crate::fleet_wire::DELIVERY_CURSOR_CAPABILITY.to_string(),
kind: Some("capacity".to_string()),
global: None,
Expand Down Expand Up @@ -4319,6 +4340,7 @@ mod tests {

ws.send(Message::Text(
serde_json::to_string(&RelaycastToBroker::ActionInvoke(ActionInvoke {
task_execution: None,
v: FLEET_WIRE_VERSION,
invocation_id: "inv-1".to_string(),
action: "run:test".to_string(),
Expand Down Expand Up @@ -4361,6 +4383,7 @@ mod tests {
command_tx
.send(FleetControlCommand::Send(BrokerToRelaycast::ActionResult(
ActionResult {
task: None,
v: FLEET_WIRE_VERSION,
id: None,
invocation_id: "inv-1".to_string(),
Expand Down Expand Up @@ -4827,6 +4850,7 @@ mod tests {
result: ActionResultPayload::Output(ActionResultOutput {
output: json!({"ok": true}),
}),
task: None,
},
)))
.await
Expand Down Expand Up @@ -6195,6 +6219,42 @@ mod tests {
consecutive_unauthorized.saturating_add(1)
));
}
#[tokio::test]
async fn task_receipts_are_forwarded_without_consuming_agent_registration_waiters() {
let (tx, mut rx) = mpsc::channel(4);
let mut registrations = HashMap::new();
let mut deregistrations = HashMap::new();
let mut liveness = ApplicationLiveness::new(Duration::from_secs(1));
let mut sink = futures_util::sink::drain();
for raw in [
serde_json::json!({"v":1,"type":"reply","id":"task_receipt_a","ok":true,"data":{"status":"running"}}),
serde_json::json!({"v":1,"type":"error","id":"task_receipt_b","ok":false,"code":"stale_task_execution","message":"stale"}),
] {
assert!(
handle_server_message(
Message::Text(raw.to_string()),
&tx,
&mut registrations,
&mut deregistrations,
&mut liveness,
"node-test",
&mut sink,
None,
)
.await
);
let event = rx.recv().await.unwrap();
match event {
FleetControlEvent::Message(RelaycastToBroker::Reply(reply)) => {
assert_eq!(reply.id, "task_receipt_a")
}
FleetControlEvent::Message(RelaycastToBroker::Error(error)) => {
assert_eq!(error.id, "task_receipt_b")
}
other => panic!("unexpected event {other:?}"),
}
}
}
}

#[cfg(test)]
Expand Down
Loading
Loading