diff --git a/.agentworkforce/trajectories/active/traj_a5b8spueklgc/trajectory.json b/.agentworkforce/trajectories/active/traj_a5b8spueklgc/trajectory.json new file mode 100644 index 0000000000..e6f46e1603 --- /dev/null +++ b/.agentworkforce/trajectories/active/traj_a5b8spueklgc/trajectory.json @@ -0,0 +1,125 @@ +{ + "id": "traj_a5b8spueklgc", + "version": 1, + "task": { + "title": "Resolve all active review and CI findings for Relay PR #1851", + "source": { + "system": "plain", + "id": "AgentWorkforce/relay#1851" + } + }, + "status": "active", + "startedAt": "2026-09-25T01:22:49.708Z", + "agents": [ + { + "name": "default", + "role": "lead", + "joinedAt": "2026-09-25T02:50:08.242Z" + } + ], + "chapters": [ + { + "id": "chap_b91llr1ynosb", + "title": "Work", + "agentName": "default", + "startedAt": "2026-09-25T02:50:08.242Z", + "events": [ + { + "ts": 1790304608243, + "type": "decision", + "content": "Use one atomically published receipt file per delivery ID: Use one atomically published receipt file per delivery ID", + "raw": { + "question": "Use one atomically published receipt file per delivery ID", + "chosen": "Use one atomically published receipt file per delivery ID", + "alternatives": [], + "reasoning": "Removes the fixed-capacity outage and whole-ledger rewrites while preserving indefinite idempotency without unsafe pruning" + }, + "significance": "high" + }, + { + "ts": 1790304615110, + "type": "decision", + "content": "Persist deferred native messages with explicit queued, in-doubt, and accepted states: Persist deferred native messages with explicit queued, in-doubt, and accepted states", + "raw": { + "question": "Persist deferred native messages with explicit queued, in-doubt, and accepted states", + "chosen": "Persist deferred native messages with explicit queued, in-doubt, and accepted states", + "alternatives": [], + "reasoning": "Restarts can recover queued messages, never replay an ambiguous injection, and safely discard a message already accepted before queue cleanup" + }, + "significance": "high" + }, + { + "ts": 1790307646867, + "type": "decision", + "content": "Classify a dropped runtime reply by consulting the durable receipt outside the Tokio event loop: Classify a dropped runtime reply by consulting the durable receipt outside the Tokio event loop", + "raw": { + "question": "Classify a dropped runtime reply by consulting the durable receipt outside the Tokio event loop", + "chosen": "Classify a dropped runtime reply by consulting the durable receipt outside the Tokio event loop", + "alternatives": [ + { + "option": "Always claim committed=true", + "reason": "" + }, + { + "option": "Always claim committed=false", + "reason": "" + } + ], + "reasoning": "The oneshot can close both before and after reservation. An exact receipt distinguishes safely retryable uncommitted delivery from queued or in-doubt delivery without weakening fail-closed semantics." + }, + "significance": "high" + }, + { + "ts": 1790307660847, + "type": "decision", + "content": "Require stable native sidecar session and runtime roots and place launch state under the Agent Relay user data directory: Require stable native sidecar session and runtime roots and place launch state under the Agent Relay user data directory", + "raw": { + "question": "Require stable native sidecar session and runtime roots and place launch state under the Agent Relay user data directory", + "chosen": "Require stable native sidecar session and runtime roots and place launch state under the Agent Relay user data directory", + "alternatives": [ + { + "option": "Keep the OS temp default with a suppression", + "reason": "" + }, + { + "option": "Use a random root that loses restart discovery", + "reason": "" + } + ], + "reasoning": "Deferred receipts must survive restart, and a predictable shared OS-temp root permits path attacks. The stable user data root preserves durability while removing the CodeQL temp-path flow." + }, + "significance": "high" + }, + { + "ts": 1790308221641, + "type": "decision", + "content": "Report post-publish receipt directory sync failures as in-doubt: Report post-publish receipt directory sync failures as in-doubt", + "raw": { + "question": "Report post-publish receipt directory sync failures as in-doubt", + "chosen": "Report post-publish receipt directory sync failures as in-doubt", + "alternatives": [ + { + "option": "Return ReceiptUnavailable after publish", + "reason": "" + }, + { + "option": "Ignore directory fsync failures", + "reason": "" + } + ], + "reasoning": "persist_noclobber has already made the reservation visible, so committed=false would invite a retry that can conflict with at-most-once semantics. Strict fsync remains required, but failure is classified conservatively." + }, + "significance": "high" + } + ] + } + ], + "commits": [], + "filesChanged": [], + "projectId": "AgentWorkforce/relay", + "tags": [], + "_trace": { + "startRef": "98cfa72090642b975ad2ac7cfe3c6d82b115f4b1", + "endRef": "98cfa72090642b975ad2ac7cfe3c6d82b115f4b1" + } +} \ No newline at end of file diff --git a/CHANGELOG.md b/CHANGELOG.md index 87cdd0bbd5..2502854ff2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,6 +10,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added - `teams.json` agents accept a per-agent `model` field when `up --spawn` starts them; an explicit `--model` or `-m` inside `cli` still wins. +- Persistent brokers expose authenticated, versioned native existing-session delivery and reconciliation for Cloud Babysitter. They durably reserve each `deliveryId` before the sole worker write, return stable receipts for exact duplicates, and reject session substitution or unsupported native input without sending. ### Changed diff --git a/crates/broker/src/lib.rs b/crates/broker/src/lib.rs index 04855823bc..ca30af129b 100644 --- a/crates/broker/src/lib.rs +++ b/crates/broker/src/lib.rs @@ -27,6 +27,7 @@ pub(crate) mod events; pub(crate) mod listen_api; #[allow(dead_code)] pub(crate) mod metrics; +pub(crate) mod native_delivery; pub(crate) mod node_control; pub(crate) mod node_delivery_probe; #[allow(dead_code)] diff --git a/crates/broker/src/listen_api.rs b/crates/broker/src/listen_api.rs index 88b6bfbe40..f485e3f807 100644 --- a/crates/broker/src/listen_api.rs +++ b/crates/broker/src/listen_api.rs @@ -6,6 +6,7 @@ use std::{ collections::HashMap, + path::PathBuf, sync::Arc, time::{Duration, Instant}, }; @@ -15,6 +16,9 @@ use crate::{ ids::{ ChannelName, DeliveryId, MessageTarget, ThreadId, WorkerName, WorkspaceAlias, WorkspaceId, }, + native_delivery::{ + NativeDeliveryError, NativeExistingSessionDelivery, NativeExistingSessionReconcile, + }, protocol::{MessageInjectionMode, ResolvedHarnessConfig}, relaycast::WorkspaceMembershipSummary, replay_buffer::ReplayBuffer, @@ -82,6 +86,14 @@ pub enum ListenApiRequest { List { reply: tokio::sync::oneshot::Sender>, }, + DeliverNativeExistingSession { + delivery: NativeExistingSessionDelivery, + reply: tokio::sync::oneshot::Sender>, + }, + ReconcileNativeExistingSession { + delivery: NativeExistingSessionReconcile, + reply: tokio::sync::oneshot::Sender>, + }, /// `GET /api/fleet-inventory` — snapshot of the in-process `fleet_inventory` /// map (what the broker last published to the engine via `inventory.sync`). /// Callers use this alongside `List` to detect the workers-vs-inventory @@ -415,6 +427,10 @@ struct ListenApiState { node_token: std::sync::Arc>>, /// Whether the broker is in persist mode persist: bool, + /// Direct receipt lookup used only when the runtime reply channel drops, + /// so the HTTP contract can distinguish pre-reservation failure from an + /// already committed or in-doubt delivery. + native_delivery_receipts: PathBuf, /// Node-control inbound introspection. Held directly (rather than reached /// through `tx`) so `GET /api/node-delivery` answers even when the runtime /// event loop is wedged — the case the endpoint exists to diagnose. @@ -455,6 +471,7 @@ pub struct ListenApiConfig { pub node_name: String, pub node_token: std::sync::Arc>>, pub persist: bool, + pub native_delivery_receipts: PathBuf, /// Node-control inbound introspection, read directly by /// `GET /api/node-delivery`. See [`crate::node_delivery_probe`]. pub node_delivery_probe: std::sync::Arc, @@ -500,6 +517,7 @@ pub(crate) fn listen_api_router_with_auth( node_name: config.node_name, node_token: config.node_token, persist: config.persist, + native_delivery_receipts: config.native_delivery_receipts, node_delivery_probe: config.node_delivery_probe, started_at: std::time::Instant::now(), input_serializers: Arc::new(tokio::sync::Mutex::new(HashMap::new())), @@ -515,6 +533,14 @@ pub(crate) fn listen_api_router_with_auth( .route("/api/session", routing::get(listen_api_session)) .route("/api/session/renew", routing::post(listen_api_renew_lease)) .route("/api/spawn", routing::post(listen_api_spawn)) + .route( + "/api/native-delivery/existing-session", + routing::post(listen_api_deliver_native_existing_session), + ) + .route( + "/api/native-delivery/existing-session/reconcile", + routing::post(listen_api_reconcile_native_existing_session), + ) .route("/api/spawned", routing::get(listen_api_list)) .route( "/api/fleet-inventory", @@ -1268,6 +1294,112 @@ async fn listen_api_list( } } +async fn listen_api_deliver_native_existing_session( + axum::extract::State(state): axum::extract::State, + axum::Json(delivery): axum::Json, +) -> (axum::http::StatusCode, axum::Json) { + let (reply_tx, reply_rx) = tokio::sync::oneshot::channel(); + let dropped_reply_delivery = delivery.clone(); + if state + .tx + .send(ListenApiRequest::DeliverNativeExistingSession { + delivery, + reply: reply_tx, + }) + .await + .is_err() + { + return internal_error(); + } + match reply_rx.await { + Ok(Ok(value)) => (axum::http::StatusCode::OK, axum::Json(value)), + Ok(Err(error)) => native_delivery_error_to_response(&error), + Err(_) if !state.persist => native_delivery_error_to_response( + &NativeDeliveryError::ReceiptUnavailable( + "runtime reply dropped before a durable receipt could be confirmed; retry the same deliveryId" + .to_string(), + ), + ), + Err(_) => match tokio::task::spawn_blocking({ + let receipt_root = state.native_delivery_receipts.clone(); + move || crate::native_delivery::existing_outcome(&receipt_root, &dropped_reply_delivery) + }) + .await + { + Err(error) => native_delivery_error_to_response(&NativeDeliveryError::InDoubt( + format!("runtime reply dropped and receipt lookup task failed: {error}"), + )), + Ok(Ok(Some(outcome))) => ( + axum::http::StatusCode::OK, + axum::Json(json!({ + "receiptId": outcome.receipt_id, + "status": "duplicate", + "state": outcome.state.as_str(), + })), + ), + Ok(Ok(None)) => native_delivery_error_to_response( + &NativeDeliveryError::ReceiptUnavailable( + "runtime reply dropped before durable reservation; retry the same deliveryId" + .to_string(), + ), + ), + Ok(Err(NativeDeliveryError::ReceiptUnavailable(error))) => { + native_delivery_error_to_response(&NativeDeliveryError::InDoubt(format!( + "runtime reply dropped and receipt lookup is unavailable: {error}" + ))) + } + Ok(Err(error)) => native_delivery_error_to_response(&error), + }, + } +} + +async fn listen_api_reconcile_native_existing_session( + axum::extract::State(state): axum::extract::State, + axum::Json(delivery): axum::Json, +) -> (axum::http::StatusCode, axum::Json) { + let (reply_tx, reply_rx) = tokio::sync::oneshot::channel(); + if state + .tx + .send(ListenApiRequest::ReconcileNativeExistingSession { + delivery, + reply: reply_tx, + }) + .await + .is_err() + { + return internal_error(); + } + match reply_rx.await { + Ok(Ok(value)) => (axum::http::StatusCode::OK, axum::Json(value)), + Ok(Err(error)) => native_delivery_error_to_response(&error), + Err(_) => internal_error(), + } +} + +fn native_delivery_error_to_response( + error: &NativeDeliveryError, +) -> (axum::http::StatusCode, axum::Json) { + use axum::http::StatusCode; + let status = match error { + NativeDeliveryError::Invalid(_) => StatusCode::BAD_REQUEST, + NativeDeliveryError::Conflict | NativeDeliveryError::Unauthorized(_) => { + StatusCode::CONFLICT + } + NativeDeliveryError::ReceiptUnavailable(_) | NativeDeliveryError::InDoubt(_) => { + StatusCode::SERVICE_UNAVAILABLE + } + }; + ( + status, + axum::Json(json!({ + "success": false, + "code": error.code(), + "error": error.to_string(), + "committed": error.committed(), + })), + ) +} + async fn listen_api_fleet_inventory( axum::extract::State(state): axum::extract::State, ) -> axum::Json { @@ -4029,6 +4161,10 @@ mod auth_tests { node_name: "test-node".to_string(), node_token: std::sync::Arc::new(std::sync::RwLock::new(None)), persist: false, + native_delivery_receipts: std::env::temp_dir().join(format!( + "agent-relay-listen-api-test-receipts-{}", + uuid::Uuid::new_v4() + )), node_delivery_probe: node_delivery_probe.clone(), }, broker_api_key.map(ToString::to_string), @@ -4235,6 +4371,237 @@ mod auth_tests { list_replier.await.expect("list replier should complete"); } + #[tokio::test] + async fn native_existing_session_delivery_requires_the_api_key() { + let (router, _rx) = test_router(Some("secret")); + let response = router + .oneshot( + Request::builder() + .uri("/api/native-delivery/existing-session") + .method("POST") + .header("content-type", "application/json") + .body(Body::from( + json!({ + "relayAgentName": "worker-a", + "sessionId": "native-session-1", + "deliveryId": "delivery-1", + "lineageId": "lineage-1", + "headSha": "a".repeat(40), + "message": "continue" + }) + .to_string(), + )) + .expect("request should build"), + ) + .await + .expect("request should succeed"); + + assert_eq!(response.status(), StatusCode::UNAUTHORIZED); + } + + #[tokio::test] + async fn native_existing_session_delivery_dispatches_typed_request_and_receipt() { + let (router, mut rx) = test_router(Some("secret")); + let replier = tokio::spawn(async move { + let Some(ListenApiRequest::DeliverNativeExistingSession { delivery, reply }) = + rx.recv().await + else { + panic!("expected native existing-session delivery"); + }; + assert_eq!(delivery.relay_agent_name, "worker-a"); + assert_eq!(delivery.session_id, "native-session-1"); + assert_eq!(delivery.delivery_id, "delivery-1"); + let _ = reply.send(Ok(json!({ + "receiptId": "ndr_receipt", + "status": "queued" + }))); + }); + + let response = router + .oneshot( + Request::builder() + .uri("/api/native-delivery/existing-session") + .method("POST") + .header("content-type", "application/json") + .header("x-api-key", "secret") + .body(Body::from( + json!({ + "relayAgentName": "worker-a", + "sessionId": "native-session-1", + "deliveryId": "delivery-1", + "lineageId": "lineage-1", + "headSha": "a".repeat(40), + "message": "continue" + }) + .to_string(), + )) + .expect("request should build"), + ) + .await + .expect("request should succeed"); + + assert_eq!(response.status(), StatusCode::OK); + assert_eq!(response_json(response).await["receiptId"], "ndr_receipt"); + replier.await.expect("delivery replier should complete"); + } + + #[tokio::test] + async fn dropped_native_delivery_reply_without_a_receipt_is_not_reported_committed() { + let (router, mut rx) = test_router(Some("secret")); + let replier = tokio::spawn(async move { + let Some(ListenApiRequest::DeliverNativeExistingSession { reply, .. }) = + rx.recv().await + else { + panic!("expected native existing-session delivery"); + }; + drop(reply); + }); + + let response = router + .oneshot( + Request::builder() + .uri("/api/native-delivery/existing-session") + .method("POST") + .header("content-type", "application/json") + .header("x-api-key", "secret") + .body(Body::from( + json!({ + "relayAgentName": "worker-a", + "sessionId": "native-session-1", + "deliveryId": "delivery-dropped-before-reservation", + "lineageId": "lineage-1", + "headSha": "a".repeat(40), + "message": "continue" + }) + .to_string(), + )) + .expect("request should build"), + ) + .await + .expect("request should succeed"); + + assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE); + let body = response_json(response).await; + assert_eq!(body["committed"], false); + assert_eq!(body["code"], "native_delivery_receipt_unavailable"); + replier.await.expect("delivery replier should complete"); + } + + #[tokio::test] + async fn dropped_native_delivery_reply_reconciles_an_in_doubt_receipt() { + let receipt_dir = tempfile::tempdir().expect("receipt root"); + let receipt_path = receipt_dir.path().join("native-delivery-receipts"); + let (tx, mut rx) = mpsc::channel(8); + let (events_tx, _events_rx) = broadcast::channel(8); + let router = listen_api_router_with_auth( + ListenApiConfig { + local_only: false, + tx, + events_tx, + replay_buffer: ReplayBuffer::new(DEFAULT_REPLAY_CAPACITY), + workspace_key: None, + relay_base_url: Some("https://relay.test".to_string()), + memberships: vec![], + default_workspace_id: None, + node_id: "node_test".to_string(), + node_name: "test-node".to_string(), + node_token: std::sync::Arc::new(std::sync::RwLock::new(None)), + persist: true, + native_delivery_receipts: receipt_path.clone(), + node_delivery_probe: std::sync::Arc::new( + crate::node_delivery_probe::NodeDeliveryProbe::new(), + ), + }, + Some("secret".to_string()), + ); + let replier = tokio::spawn(async move { + let Some(ListenApiRequest::DeliverNativeExistingSession { delivery, reply }) = + rx.recv().await + else { + panic!("expected native existing-session delivery"); + }; + let result = + crate::native_delivery::reserve_and_deliver(&receipt_path, &delivery, |_| async { + anyhow::bail!("fixture write failed") + }) + .await; + assert!(matches!( + result, + Err(crate::native_delivery::NativeDeliveryError::InDoubt(_)) + )); + drop(reply); + }); + + let response = router + .oneshot( + Request::builder() + .uri("/api/native-delivery/existing-session") + .method("POST") + .header("content-type", "application/json") + .header("x-api-key", "secret") + .body(Body::from( + json!({ + "relayAgentName": "worker-a", + "sessionId": "native-session-1", + "deliveryId": "delivery-dropped-after-reservation", + "lineageId": "lineage-1", + "headSha": "a".repeat(40), + "message": "continue" + }) + .to_string(), + )) + .expect("request should build"), + ) + .await + .expect("request should succeed"); + + assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE); + let body = response_json(response).await; + assert_eq!(body["committed"], true); + assert_eq!(body["code"], "native_delivery_in_doubt"); + replier.await.expect("delivery replier should complete"); + } + + #[tokio::test] + async fn native_existing_session_reconcile_dispatches_without_a_message() { + let (router, mut rx) = test_router(Some("secret")); + let replier = tokio::spawn(async move { + let Some(ListenApiRequest::ReconcileNativeExistingSession { delivery, reply }) = + rx.recv().await + else { + panic!("expected native existing-session reconciliation"); + }; + assert_eq!(delivery.delivery_id, "delivery-1"); + let _ = reply.send(Ok(json!({ "receiptId": null }))); + }); + + let response = router + .oneshot( + Request::builder() + .uri("/api/native-delivery/existing-session/reconcile") + .method("POST") + .header("content-type", "application/json") + .header("x-api-key", "secret") + .body(Body::from( + json!({ + "relayAgentName": "worker-a", + "sessionId": "native-session-1", + "deliveryId": "delivery-1", + "lineageId": "lineage-1", + "headSha": "a".repeat(40) + }) + .to_string(), + )) + .expect("request should build"), + ) + .await + .expect("request should succeed"); + + assert_eq!(response.status(), StatusCode::OK); + assert!(response_json(response).await["receiptId"].is_null()); + replier.await.expect("reconcile replier should complete"); + } + #[tokio::test] async fn api_route_accepts_lowercase_bearer_scheme() { let (router, mut rx) = test_router(Some("secret")); diff --git a/crates/broker/src/native_delivery.rs b/crates/broker/src/native_delivery.rs new file mode 100644 index 0000000000..9dee31b78f --- /dev/null +++ b/crates/broker/src/native_delivery.rs @@ -0,0 +1,759 @@ +//! Durable delivery into an already-running native harness session. +//! +//! The receipt is written before the broker writes to the worker. Once that +//! boundary is crossed, every retry with the same `delivery_id` is a read-only +//! lookup: an interrupted or failed write is in doubt and is never replayed. + +use std::{ + future::Future, + io::Write, + path::{Path, PathBuf}, +}; + +use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; +use thiserror::Error; + +use crate::{ + ids::{DeliveryId, EventId, MessageTarget, WorkerName}, + protocol::{MessageInjectionMode, RelayDelivery}, +}; + +pub(crate) const NATIVE_EXISTING_SESSION_CAPABILITY: &str = "relay:native-existing-session:v1"; +pub(crate) const NATIVE_EXISTING_SESSION_RECONCILE_CAPABILITY: &str = + "relay:native-existing-session-reconcile:v1"; +const MAX_MESSAGE_BYTES: usize = 128 * 1024; + +#[derive(Debug, Clone, Deserialize)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] +pub(crate) struct NativeExistingSessionDelivery { + /// Broker-owned worker identity used for live-session authorization. + pub(crate) relay_agent_name: String, + /// Broker-owned native session identity used for live-session authorization. + pub(crate) session_id: String, + /// Caller-assigned idempotency key, durably bound to the complete request. + pub(crate) delivery_id: String, + /// Caller assertion recorded for exact duplicate and reconciliation matching. + /// It is not an independent worker-authorization claim. + pub(crate) lineage_id: String, + /// Caller assertion recorded for exact duplicate and reconciliation matching. + /// It is not an independent worker-authorization claim. + pub(crate) head_sha: String, + /// Prompt content delivered to the authorized native session. + pub(crate) message: String, +} + +#[derive(Debug, Clone, Deserialize)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] +pub(crate) struct NativeExistingSessionReconcile { + /// Worker identity bound into the original durable receipt. + pub(crate) relay_agent_name: String, + /// Native session identity bound into the original durable receipt. + pub(crate) session_id: String, + /// Idempotency key whose durable receipt is being queried. + pub(crate) delivery_id: String, + /// Caller assertion that must exactly match the original receipt. + pub(crate) lineage_id: String, + /// Caller assertion that must exactly match the original receipt. + pub(crate) head_sha: String, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub(crate) enum NativeReceiptState { + /// The durable cancellation boundary was crossed. The worker write may or + /// may not have completed, so replay is forbidden. + InDoubt, + /// The sidecar confirmed durable queued custody or immediate acceptance of + /// the complete protocol frame. + Queued, +} + +impl NativeReceiptState { + pub(crate) fn as_str(self) -> &'static str { + match self { + Self::InDoubt => "in_doubt", + Self::Queued => "queued", + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub(crate) struct NativeDeliveryReceipt { + receipt_id: String, + delivery_id: String, + relay_agent_name: String, + session_id: String, + lineage_id: String, + head_sha: String, + request_digest: String, + state: NativeReceiptState, + recorded_at_ms: u64, +} + +impl NativeDeliveryReceipt { + pub(crate) fn receipt_id(&self) -> &str { + &self.receipt_id + } + + pub(crate) fn state(&self) -> NativeReceiptState { + self.state + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum NativeDeliveryDisposition { + Queued, + Duplicate, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct NativeDeliveryOutcome { + pub(crate) receipt_id: String, + pub(crate) disposition: NativeDeliveryDisposition, + pub(crate) state: NativeReceiptState, +} + +impl NativeDeliveryOutcome { + pub(crate) fn to_json(&self) -> serde_json::Value { + let status = match self.disposition { + NativeDeliveryDisposition::Queued => "queued", + NativeDeliveryDisposition::Duplicate => "duplicate", + }; + serde_json::json!({ + "receiptId": self.receipt_id, + "status": status, + "state": self.state.as_str(), + }) + } +} + +#[derive(Debug, Error)] +pub(crate) enum NativeDeliveryError { + #[error("invalid_native_delivery: {0}")] + Invalid(String), + #[error("native_delivery_conflict: delivery id is already bound to different input")] + Conflict, + #[error("native_session_unauthorized: {0}")] + Unauthorized(String), + #[error("native_delivery_receipt_unavailable: {0}")] + ReceiptUnavailable(String), + #[error("native_delivery_in_doubt: {0}")] + InDoubt(String), +} + +impl NativeDeliveryError { + pub(crate) fn committed(&self) -> bool { + matches!(self, Self::InDoubt(_)) + } + + pub(crate) fn code(&self) -> &'static str { + match self { + Self::Invalid(_) => "invalid_native_delivery", + Self::Conflict => "native_delivery_conflict", + Self::Unauthorized(_) => "native_session_unauthorized", + Self::ReceiptUnavailable(_) => "native_delivery_receipt_unavailable", + Self::InDoubt(_) => "native_delivery_in_doubt", + } + } +} + +impl NativeExistingSessionDelivery { + pub(crate) fn validate(&self) -> Result<(), NativeDeliveryError> { + validate_identifier(&self.relay_agent_name, "relayAgentName")?; + validate_identifier(&self.session_id, "sessionId")?; + validate_identifier(&self.delivery_id, "deliveryId")?; + validate_identifier(&self.lineage_id, "lineageId")?; + if self.head_sha.len() != 40 || !self.head_sha.bytes().all(|byte| byte.is_ascii_hexdigit()) + { + return Err(NativeDeliveryError::Invalid( + "headSha must be a 40-character hexadecimal commit id".to_string(), + )); + } + if self.message.trim().is_empty() { + return Err(NativeDeliveryError::Invalid( + "message must not be empty".to_string(), + )); + } + if self.message.len() > MAX_MESSAGE_BYTES { + return Err(NativeDeliveryError::Invalid(format!( + "message exceeds the {MAX_MESSAGE_BYTES}-byte limit" + ))); + } + Ok(()) + } + + /// Broker-namespaced identity for this delivery. It is the receipt id and + /// the worker-wire `delivery_id`, so custody can only be confirmed by the + /// sidecar's acknowledgement of this exact native delivery, never by an + /// ordinary delivery whose externally chosen id equals the caller's key. + fn receipt_id(&self) -> String { + format!("ndr_{}", digest(&self.delivery_id)) + } + + pub(crate) fn relay_delivery(&self) -> RelayDelivery { + RelayDelivery { + delivery_id: DeliveryId::new(self.receipt_id()), + event_id: EventId::new(self.delivery_id.clone()), + workspace_id: None, + workspace_alias: None, + from: "cloud-babysitter".to_string(), + target: MessageTarget::new(self.relay_agent_name.clone()), + body: self.message.clone(), + thread_id: None, + priority: None, + // `wait` maps to the native session's `on-idle` mode. A Babysitter + // wake-up must not interrupt an active coding turn. + injection_mode: MessageInjectionMode::Wait, + } + } + + fn request_digest(&self) -> String { + digest(&serde_json::json!([ + self.relay_agent_name, + self.session_id, + self.delivery_id, + self.lineage_id, + self.head_sha, + self.message, + ])) + } + + fn receipt(&self) -> NativeDeliveryReceipt { + NativeDeliveryReceipt { + receipt_id: self.receipt_id(), + delivery_id: self.delivery_id.clone(), + relay_agent_name: self.relay_agent_name.clone(), + session_id: self.session_id.clone(), + lineage_id: self.lineage_id.clone(), + head_sha: self.head_sha.clone(), + request_digest: self.request_digest(), + state: NativeReceiptState::InDoubt, + recorded_at_ms: chrono::Utc::now().timestamp_millis().max(0) as u64, + } + } +} + +impl NativeExistingSessionReconcile { + pub(crate) fn validate(&self) -> Result<(), NativeDeliveryError> { + validate_identifier(&self.relay_agent_name, "relayAgentName")?; + validate_identifier(&self.session_id, "sessionId")?; + validate_identifier(&self.delivery_id, "deliveryId")?; + validate_identifier(&self.lineage_id, "lineageId")?; + if self.head_sha.len() != 40 || !self.head_sha.bytes().all(|byte| byte.is_ascii_hexdigit()) + { + return Err(NativeDeliveryError::Invalid( + "headSha must be a 40-character hexadecimal commit id".to_string(), + )); + } + Ok(()) + } +} + +fn validate_identifier(value: &str, field: &str) -> Result<(), NativeDeliveryError> { + if value.trim().is_empty() || value.len() > 512 || value.chars().any(char::is_control) { + return Err(NativeDeliveryError::Invalid(format!( + "{field} must be a non-empty bounded identifier" + ))); + } + Ok(()) +} + +fn digest(value: &impl Serialize) -> String { + let bytes = serde_json::to_vec(value).expect("native delivery digest input is serializable"); + format!("{:x}", Sha256::digest(bytes)) +} + +fn receipt_path(root: &Path, delivery_id: &str) -> PathBuf { + root.join(format!("{}.json", digest(&delivery_id))) +} + +fn load_receipt( + root: &Path, + delivery_id: &str, +) -> Result, NativeDeliveryError> { + let path = receipt_path(root, delivery_id); + match std::fs::read(&path) { + Ok(bytes) => serde_json::from_slice(&bytes).map(Some).map_err(|error| { + NativeDeliveryError::ReceiptUnavailable(format!( + "could not parse {}: {error}", + path.display() + )) + }), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None), + Err(error) => Err(NativeDeliveryError::ReceiptUnavailable(format!( + "could not read {}: {error}", + path.display() + ))), + } +} + +fn save_receipt(root: &Path, receipt: &NativeDeliveryReceipt) -> Result<(), NativeDeliveryError> { + let path = receipt_path(root, &receipt.delivery_id); + crate::util::fs::write_json_atomic(&path, receipt).map_err(|error| { + NativeDeliveryError::ReceiptUnavailable(format!( + "could not persist {}: {error}", + path.display() + )) + })?; + sync_receipt_directories(root) +} + +#[cfg(unix)] +fn sync_receipt_directories(root: &Path) -> Result<(), NativeDeliveryError> { + let sync_dir = |path: &Path| { + std::fs::File::open(path) + .and_then(|directory| directory.sync_all()) + .map_err(|error| { + NativeDeliveryError::ReceiptUnavailable(format!( + "could not sync receipt directory {}: {error}", + path.display() + )) + }) + }; + sync_dir(root)?; + if let Some(parent) = root.parent().filter(|path| !path.as_os_str().is_empty()) { + sync_dir(parent)?; + } + Ok(()) +} + +#[cfg(not(unix))] +fn sync_receipt_directories(_root: &Path) -> Result<(), NativeDeliveryError> { + Ok(()) +} + +/// Create a durable per-delivery reservation without replacing an existing +/// receipt. One immutable file per delivery id avoids whole-ledger rewrites +/// and preserves idempotency history without a fixed lifetime capacity cliff. +fn create_receipt( + root: &Path, + receipt: &NativeDeliveryReceipt, +) -> Result { + std::fs::create_dir_all(root).map_err(|error| { + NativeDeliveryError::ReceiptUnavailable(format!( + "could not create {}: {error}", + root.display() + )) + })?; + let path = receipt_path(root, &receipt.delivery_id); + let mut file = tempfile::NamedTempFile::new_in(root).map_err(|error| { + NativeDeliveryError::ReceiptUnavailable(format!( + "could not create receipt temporary file in {}: {error}", + root.display() + )) + })?; + let bytes = serde_json::to_vec(receipt).expect("native receipt is serializable"); + file.write_all(&bytes) + .and_then(|()| file.as_file().sync_all()) + .map_err(|error| { + NativeDeliveryError::ReceiptUnavailable(format!( + "could not prepare {}: {error}", + path.display() + )) + })?; + match file.persist_noclobber(&path) { + Ok(_) => {} + Err(error) if error.error.kind() == std::io::ErrorKind::AlreadyExists => return Ok(false), + Err(error) => { + return Err(NativeDeliveryError::ReceiptUnavailable(format!( + "could not reserve {}: {}", + path.display(), + error.error + ))) + } + } + // The reservation file is already visible after persist_noclobber. A + // directory-sync failure cannot be reported as uncommitted: retrying the + // worker write would violate at-most-once delivery if the entry survives. + sync_receipt_directories(root) + .map_err(|error| NativeDeliveryError::InDoubt(error.to_string()))?; + Ok(true) +} + +/// Return the existing receipt for an exact duplicate. A delivery id reused +/// with any different field is a conflict, never an authorization to resend. +pub(crate) fn existing_receipt( + path: &Path, + input: &NativeExistingSessionDelivery, +) -> Result, NativeDeliveryError> { + input.validate()?; + let Some(receipt) = load_receipt(path, &input.delivery_id)? else { + return Ok(None); + }; + if receipt.request_digest != input.request_digest() { + return Err(NativeDeliveryError::Conflict); + } + Ok(Some(receipt)) +} + +fn duplicate_outcome( + receipt: NativeDeliveryReceipt, +) -> Result { + if receipt.state == NativeReceiptState::InDoubt { + return Err(NativeDeliveryError::InDoubt(format!( + "receipt {} remains in doubt; reconcile before proceeding", + receipt.receipt_id + ))); + } + Ok(NativeDeliveryOutcome { + receipt_id: receipt.receipt_id, + disposition: NativeDeliveryDisposition::Duplicate, + state: receipt.state, + }) +} + +pub(crate) fn existing_outcome( + path: &Path, + input: &NativeExistingSessionDelivery, +) -> Result, NativeDeliveryError> { + existing_receipt(path, input)? + .map(duplicate_outcome) + .transpose() +} + +/// Run receipt filesystem work on the blocking pool so directory and file +/// fsyncs never stall the broker runtime actor or a tokio worker thread. A +/// join failure after a reservation may have been published is in doubt. +async fn run_blocking(committed: bool, work: F) -> Result +where + T: Send + 'static, + F: FnOnce() -> Result + Send + 'static, +{ + tokio::task::spawn_blocking(work).await.map_err(|error| { + let message = format!("native receipt task failed: {error}"); + if committed { + NativeDeliveryError::InDoubt(message) + } else { + NativeDeliveryError::ReceiptUnavailable(message) + } + })? +} + +/// Look up an exact duplicate without blocking the async runtime. +pub(crate) async fn existing_outcome_async( + path: PathBuf, + input: NativeExistingSessionDelivery, +) -> Result, NativeDeliveryError> { + run_blocking(false, move || existing_outcome(&path, &input)).await +} + +/// Reserve an exact delivery durably, then perform its one permitted worker +/// write. The reservation remains `in_doubt` after every post-reservation +/// failure so a caller can reconcile but can never cause a second write. +pub(crate) async fn reserve_and_deliver( + path: &Path, + input: &NativeExistingSessionDelivery, + send: F, +) -> Result +where + F: FnOnce(RelayDelivery) -> Fut, + Fut: Future>, +{ + if let Some(outcome) = existing_outcome_async(path.to_path_buf(), input.clone()).await? { + return Ok(outcome); + } + + let receipt = input.receipt(); + // This is the cancellation boundary. No worker write may occur unless this + // exact reservation is durable on disk first. + let reserved = run_blocking(true, { + let root = path.to_path_buf(); + let receipt = receipt.clone(); + move || create_receipt(&root, &receipt) + }) + .await?; + if !reserved { + let receipt = run_blocking(false, { + let root = path.to_path_buf(); + let input = input.clone(); + move || existing_receipt(&root, &input) + }) + .await? + .ok_or_else(|| { + NativeDeliveryError::ReceiptUnavailable( + "concurrent receipt reservation was not readable".to_string(), + ) + })?; + return duplicate_outcome(receipt); + } + + if let Err(error) = send(input.relay_delivery()).await { + return Err(NativeDeliveryError::InDoubt(error.to_string())); + } + + let mut receipt = receipt; + receipt.state = NativeReceiptState::Queued; + let receipt_id = receipt.receipt_id.clone(); + // Failure here is still in doubt: the durable write-ahead record remains + // authoritative and a retry will return it without another worker write. + run_blocking(true, { + let root = path.to_path_buf(); + move || save_receipt(&root, &receipt) + }) + .await + .map_err(|error| match error { + NativeDeliveryError::InDoubt(_) => error, + other => NativeDeliveryError::InDoubt(other.to_string()), + })?; + Ok(NativeDeliveryOutcome { + receipt_id, + disposition: NativeDeliveryDisposition::Queued, + state: NativeReceiptState::Queued, + }) +} + +pub(crate) fn reconcile_receipt( + path: &Path, + input: &NativeExistingSessionReconcile, +) -> Result, NativeDeliveryError> { + input.validate()?; + let Some(receipt) = load_receipt(path, &input.delivery_id)? else { + return Ok(None); + }; + if receipt.relay_agent_name != input.relay_agent_name + || receipt.session_id != input.session_id + || receipt.lineage_id != input.lineage_id + || receipt.head_sha != input.head_sha + { + return Err(NativeDeliveryError::Conflict); + } + Ok(Some(receipt)) +} + +pub(crate) fn worker_name(input: &NativeExistingSessionDelivery) -> WorkerName { + WorkerName::new(input.relay_agent_name.clone()) +} + +/// Deliver with a live-session authorization captured by the runtime actor. +/// An exact duplicate is answered from its durable receipt even when the +/// worker has since exited; only a new reservation requires authorization. +pub(crate) async fn deliver_authorized( + path: PathBuf, + input: NativeExistingSessionDelivery, + authorization: Result, +) -> Result { + if let Some(outcome) = existing_outcome_async(path.clone(), input.clone()).await? { + return Ok(outcome); + } + let sender = authorization.map_err(NativeDeliveryError::Unauthorized)?; + reserve_and_deliver(&path, &input, |delivery| async move { + sender.deliver(delivery).await + }) + .await +} + +/// Reconcile an exact receipt without blocking the async runtime. +pub(crate) async fn reconcile_receipt_async( + path: PathBuf, + input: NativeExistingSessionReconcile, +) -> Result, NativeDeliveryError> { + run_blocking(false, move || reconcile_receipt(&path, &input)).await +} + +/// Response body for a reconciliation lookup. +pub(crate) fn reconcile_json(receipt: Option) -> serde_json::Value { + match receipt { + Some(receipt) => serde_json::json!({ + "receiptId": receipt.receipt_id(), + "state": receipt.state().as_str(), + }), + None => serde_json::json!({ "receiptId": null }), + } +} + +#[cfg(test)] +mod tests { + use std::sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, + }; + + use super::*; + + fn delivery(id: &str) -> NativeExistingSessionDelivery { + NativeExistingSessionDelivery { + relay_agent_name: "garden-coder".to_string(), + session_id: "native-session-1".to_string(), + delivery_id: id.to_string(), + lineage_id: "lineage-1".to_string(), + head_sha: "a".repeat(40), + message: "Review the exact live head".to_string(), + } + } + + #[tokio::test] + async fn worker_wire_delivery_id_is_namespaced_by_the_receipt() { + // The custody waiter is keyed by the wire delivery id. A caller key + // that equals an ordinary broker/engine delivery id (for example + // `del_42`) must not let that delivery's ACK confirm native custody. + let dir = tempfile::tempdir().expect("receipt dir"); + let path = dir.path().join("receipts"); + let wire = Arc::new(std::sync::Mutex::new(None)); + let captured = Arc::clone(&wire); + let outcome = reserve_and_deliver(&path, &delivery("del_42"), move |delivery| { + *captured.lock().expect("capture") = Some(delivery); + async { Ok(()) } + }) + .await + .expect("delivery"); + let wire = wire.lock().expect("capture").take().expect("worker write"); + assert_ne!(wire.delivery_id.as_str(), "del_42"); + assert_eq!(wire.delivery_id.as_str(), outcome.receipt_id); + assert!(outcome.receipt_id.starts_with("ndr_")); + assert_eq!( + wire.event_id.as_str(), + "del_42", + "the caller key remains the correlated event id" + ); + } + + #[tokio::test] + async fn writes_ahead_and_never_resends_an_ambiguous_delivery() { + let dir = tempfile::tempdir().expect("receipt dir"); + let path = dir.path().join("receipts.json"); + let writes = Arc::new(AtomicUsize::new(0)); + let first_writes = Arc::clone(&writes); + let reserved_path = path.clone(); + + let error = reserve_and_deliver(&path, &delivery("bst_1"), move |_| { + assert!( + reserved_path.exists(), + "write-ahead receipt must exist before the worker callback" + ); + first_writes.fetch_add(1, Ordering::SeqCst); + async { anyhow::bail!("writer outcome unknown") } + }) + .await + .expect_err("ambiguous write should fail"); + assert!(error.committed()); + + let retry_writes = Arc::clone(&writes); + let retry = reserve_and_deliver(&path, &delivery("bst_1"), move |_| { + retry_writes.fetch_add(1, Ordering::SeqCst); + async { Ok(()) } + }) + .await + .expect_err("retry must preserve the in-doubt signal"); + assert!(matches!(retry, NativeDeliveryError::InDoubt(_))); + assert_eq!(writes.load(Ordering::SeqCst), 1); + + let reconcile = NativeExistingSessionReconcile { + relay_agent_name: "garden-coder".to_string(), + session_id: "native-session-1".to_string(), + delivery_id: "bst_1".to_string(), + lineage_id: "lineage-1".to_string(), + head_sha: "a".repeat(40), + }; + assert_eq!( + reconcile_receipt(&path, &reconcile) + .expect("reconcile") + .expect("receipt") + .state(), + NativeReceiptState::InDoubt + ); + } + + #[tokio::test] + async fn exact_duplicate_is_idempotent_but_changed_input_conflicts() { + let dir = tempfile::tempdir().expect("receipt dir"); + let path = dir.path().join("receipts.json"); + let writes = Arc::new(AtomicUsize::new(0)); + let first_writes = Arc::clone(&writes); + let first = reserve_and_deliver(&path, &delivery("bst_2"), move |_| { + first_writes.fetch_add(1, Ordering::SeqCst); + async { Ok(()) } + }) + .await + .expect("first delivery"); + assert_eq!(first.disposition, NativeDeliveryDisposition::Queued); + assert_eq!(first.state, NativeReceiptState::Queued); + + let duplicate = reserve_and_deliver(&path, &delivery("bst_2"), |_| async { Ok(()) }) + .await + .expect("duplicate delivery"); + assert_eq!(duplicate.disposition, NativeDeliveryDisposition::Duplicate); + assert_eq!(duplicate.state, NativeReceiptState::Queued); + + let mut changed = delivery("bst_2"); + changed.message = "different authority".to_string(); + assert!(matches!( + reserve_and_deliver(&path, &changed, |_| async { Ok(()) }).await, + Err(NativeDeliveryError::Conflict) + )); + assert_eq!(writes.load(Ordering::SeqCst), 1); + } + + #[tokio::test] + async fn restart_reloads_the_durable_receipt_and_reconciliation_is_exact() { + let dir = tempfile::tempdir().expect("receipt dir"); + let path = dir.path().join("receipts.json"); + let request = delivery("bst_3"); + let queued = reserve_and_deliver(&path, &request, |_| async { Ok(()) }) + .await + .expect("first delivery"); + + let reconcile = NativeExistingSessionReconcile { + relay_agent_name: request.relay_agent_name.clone(), + session_id: request.session_id.clone(), + delivery_id: request.delivery_id.clone(), + lineage_id: request.lineage_id.clone(), + head_sha: request.head_sha.clone(), + }; + let restored = reconcile_receipt(&path, &reconcile) + .expect("reconcile") + .expect("durable receipt"); + assert_eq!(restored.receipt_id(), queued.receipt_id); + + let mismatched = NativeExistingSessionReconcile { + session_id: "replacement-session".to_string(), + ..reconcile + }; + assert!(matches!( + reconcile_receipt(&path, &mismatched), + Err(NativeDeliveryError::Conflict) + )); + } + + #[test] + fn invalid_inputs_fail_before_receipt_creation() { + let dir = tempfile::tempdir().expect("receipt dir"); + let path = dir.path().join("receipts.json"); + let mut input = delivery("bst_bad"); + input.head_sha = "not-a-sha".to_string(); + assert!(matches!( + existing_receipt(&path, &input), + Err(NativeDeliveryError::Invalid(_)) + )); + assert!(!path.exists()); + } + + #[tokio::test] + async fn concurrent_reservations_permit_only_one_worker_write() { + let dir = tempfile::tempdir().expect("receipt dir"); + let path = Arc::new(dir.path().join("receipts")); + let input = Arc::new(delivery("bst_race")); + let writes = Arc::new(AtomicUsize::new(0)); + + let attempt = |path: Arc, + input: Arc, + writes: Arc| async move { + reserve_and_deliver(path.as_path(), input.as_ref(), move |_| { + let writes = Arc::clone(&writes); + async move { + writes.fetch_add(1, Ordering::SeqCst); + tokio::task::yield_now().await; + Ok(()) + } + }) + .await + }; + + let (left, right) = tokio::join!( + attempt(Arc::clone(&path), Arc::clone(&input), Arc::clone(&writes)), + attempt(path, input, Arc::clone(&writes)), + ); + assert!(left.is_ok() || right.is_ok()); + assert_eq!(writes.load(Ordering::SeqCst), 1); + } +} diff --git a/crates/broker/src/runtime/api.rs b/crates/broker/src/runtime/api.rs index e1e82a1235..9c027dbc5b 100644 --- a/crates/broker/src/runtime/api.rs +++ b/crates/broker/src/runtime/api.rs @@ -1748,6 +1748,53 @@ impl BrokerRuntime { super::delivery::pending_message_counts(delivery_states, pending_deliveries); let _ = reply.send(Ok(json!({ "agents": workers.list(&counts) }))); } + ListenApiRequest::DeliverNativeExistingSession { delivery, reply } => { + if !paths.persist { + let _ = reply.send(Err( + crate::native_delivery::NativeDeliveryError::ReceiptUnavailable( + "persistent broker state is required for native delivery".to_string(), + ), + )); + return; + } + // Authorization reads only in-memory worker state. Receipt + // lookup, reservation, and fsyncs run in the spawned task on + // the blocking pool so they never stall this runtime actor. + let name = crate::native_delivery::worker_name(&delivery); + let authorization = workers + .authorize_native_existing_session(&name, &delivery.session_id) + .map_err(|error| error.to_string()); + let receipt_path = paths.native_delivery_receipts.clone(); + tokio::spawn(async move { + let result = crate::native_delivery::deliver_authorized( + receipt_path, + delivery, + authorization, + ) + .await + .map(|outcome| outcome.to_json()); + let _ = reply.send(result); + }); + } + ListenApiRequest::ReconcileNativeExistingSession { delivery, reply } => { + if !paths.persist { + let _ = reply.send(Err( + crate::native_delivery::NativeDeliveryError::ReceiptUnavailable( + "persistent broker state is required for native reconciliation" + .to_string(), + ), + )); + return; + } + let receipt_path = paths.native_delivery_receipts.clone(); + tokio::spawn(async move { + let result = + crate::native_delivery::reconcile_receipt_async(receipt_path, delivery) + .await + .map(crate::native_delivery::reconcile_json); + let _ = reply.send(result); + }); + } ListenApiRequest::FleetInventory { reply } => { // Report the in-process `fleet_inventory` map: the same // snapshot the broker publishes to the engine via diff --git a/crates/broker/src/runtime/fleet.rs b/crates/broker/src/runtime/fleet.rs index 388115db7e..42ff384bef 100644 --- a/crates/broker/src/runtime/fleet.rs +++ b/crates/broker/src/runtime/fleet.rs @@ -1405,6 +1405,15 @@ impl BrokerRuntime { self.handle_task_invoke(invoke).await; return; } + if invoke.action == crate::native_delivery::NATIVE_EXISTING_SESSION_CAPABILITY { + self.handle_native_existing_session_invoke(invoke).await; + return; + } + if invoke.action == crate::native_delivery::NATIVE_EXISTING_SESSION_RECONCILE_CAPABILITY { + self.handle_native_existing_session_reconcile_invoke(invoke) + .await; + return; + } let action = invoke.action.as_str(); if action == "spawn" || action.starts_with("spawn:") { self.handle_fleet_action_spawn(invoke).await; @@ -1428,6 +1437,77 @@ impl BrokerRuntime { .await; } + async fn handle_native_existing_session_invoke(&mut self, invoke: ActionInvoke) { + if !self.paths.persist { + self.reply_action_error( + &invoke.invocation_id, + "native_delivery_receipt_unavailable: persistent broker state is required; committed=false", + ) + .await; + return; + } + let delivery = match serde_json::from_value(invoke.input) { + Ok(delivery) => delivery, + Err(error) => { + self.reply_action_error( + &invoke.invocation_id, + &format!("invalid_native_delivery: {error}"), + ) + .await; + return; + } + }; + // Authorization reads only in-memory worker state. Receipt lookup, + // reservation, and fsyncs run in the spawned task on the blocking pool + // so they never stall the fleet actor. + let name = crate::native_delivery::worker_name(&delivery); + let authorization = self + .workers + .authorize_native_existing_session(&name, &delivery.session_id) + .map_err(|error| error.to_string()); + let receipt_path = self.paths.native_delivery_receipts.clone(); + let invocation_id = invoke.invocation_id; + let control_tx = self.fleet_control_tx.clone(); + tokio::spawn(async move { + let result = + crate::native_delivery::deliver_authorized(receipt_path, delivery, authorization) + .await + .map(|outcome| outcome.to_json()); + send_native_action_result(&control_tx, invocation_id, result).await; + }); + } + + async fn handle_native_existing_session_reconcile_invoke(&self, invoke: ActionInvoke) { + if !self.paths.persist { + self.reply_action_error( + &invoke.invocation_id, + "native_delivery_receipt_unavailable: persistent broker state is required; committed=false", + ) + .await; + return; + } + let delivery = match serde_json::from_value(invoke.input) { + Ok(delivery) => delivery, + Err(error) => { + self.reply_action_error( + &invoke.invocation_id, + &format!("invalid_native_delivery: {error}"), + ) + .await; + return; + } + }; + let receipt_path = self.paths.native_delivery_receipts.clone(); + let invocation_id = invoke.invocation_id; + let control_tx = self.fleet_control_tx.clone(); + tokio::spawn(async move { + let result = crate::native_delivery::reconcile_receipt_async(receipt_path, delivery) + .await + .map(crate::native_delivery::reconcile_json); + send_native_action_result(&control_tx, invocation_id, result).await; + }); + } + /// Run a `spawn` / `spawn:` node action by parsing the invoke input /// into spawn fields and calling the local spawn fn (which binds the agent /// to this node). Replies with `action.result { output }` on success or @@ -1980,6 +2060,31 @@ pub(super) struct FlushPendingRelayResult { pub(super) blocked_agent_id: Option, } +/// Reply to a native existing-session Fleet action from a spawned task. +async fn send_native_action_result( + control_tx: &mpsc::Sender, + invocation_id: String, + result: Result, +) { + let result = match result { + Ok(output) => ActionResultPayload::Output(ActionResultOutput { output }), + Err(error) => ActionResultPayload::Error(ActionResultError { + error: format!("{}; committed={}", error, error.committed()), + }), + }; + let _ = control_tx + .send(FleetControlCommand::Send(BrokerToRelaycast::ActionResult( + ActionResult { + task: None, + v: FLEET_WIRE_VERSION, + id: None, + invocation_id, + result, + }, + ))) + .await; +} + /// Inject a worker's held queue in FIFO order. A failed item and every item /// behind it remain queued. Relaycast ACKs advance only after the corresponding /// PTY write succeeds, so the emitted cursor is always an injected prefix. diff --git a/crates/broker/src/runtime/init.rs b/crates/broker/src/runtime/init.rs index f1f1eb9027..571f57d93e 100644 --- a/crates/broker/src/runtime/init.rs +++ b/crates/broker/src/runtime/init.rs @@ -330,6 +330,11 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re metadata: None, }); } + // A unique ephemeral state directory cannot restore receipts after a + // broker restart, so it must never advertise the durable Cloud contract. + if paths.persist { + append_native_delivery_capabilities(&mut node_manifest); + } // Retain the node name for the runtime: the HTTP `bind_agent_to_node` // fallback (used when node-control `agent.register` is unavailable) binds // spawned agents to this node so they become `via_node` and node delivery @@ -500,6 +505,7 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re node_name: session_node_name, node_token: session_node_token, persist: paths.persist, + native_delivery_receipts: paths.native_delivery_receipts.clone(), node_delivery_probe: node_delivery_probe.clone(), }); { @@ -1005,6 +1011,39 @@ fn bootstrap_node_manifest(node_name: &str, node_id: &str, broker_version: &str) } } +fn append_native_delivery_capabilities(manifest: &mut NodeManifest) { + manifest + .capabilities + .push(crate::protocol::NodeCapabilityManifest { + name: crate::native_delivery::NATIVE_EXISTING_SESSION_CAPABILITY.to_owned(), + kind: Some("action".to_owned()), + metadata: Some(HashMap::from([ + ("contract".to_owned(), json!("deliverNativeExistingSession")), + ("contractVersion".to_owned(), json!(1)), + ("durableReceipts".to_owned(), json!(true)), + ("idempotencyField".to_owned(), json!("deliveryId")), + ("sessionAuthorization".to_owned(), json!("exact")), + ( + "reconcileAction".to_owned(), + json!(crate::native_delivery::NATIVE_EXISTING_SESSION_RECONCILE_CAPABILITY), + ), + ])), + }); + manifest + .capabilities + .push(crate::protocol::NodeCapabilityManifest { + name: crate::native_delivery::NATIVE_EXISTING_SESSION_RECONCILE_CAPABILITY.to_owned(), + kind: Some("action".to_owned()), + metadata: Some(HashMap::from([ + ( + "contract".to_owned(), + json!("reconcileNativeExistingSession"), + ), + ("contractVersion".to_owned(), json!(1)), + ])), + }); +} + /// The harness names this broker can spawn, from `AGENT_RELAY_NODE_HARNESSES` /// (comma-separated, order-preserving, de-duplicated) or the built-in default. fn node_capacity_harnesses() -> Vec { @@ -1200,6 +1239,38 @@ mod tests { assert_eq!(manifest.version.as_deref(), Some("relay-broker/9.1.1")); } + #[test] + fn native_delivery_manifest_publishes_versioned_fail_closed_contract() { + let mut manifest = bootstrap_node_manifest("node-a", "node_a", "relay-broker/9.1.1"); + append_native_delivery_capabilities(&mut manifest); + + let delivery = manifest + .capabilities + .iter() + .find(|cap| cap.name == crate::native_delivery::NATIVE_EXISTING_SESSION_CAPABILITY) + .expect("native delivery action must be advertised"); + assert_eq!(delivery.kind.as_deref(), Some("action")); + let metadata = delivery.metadata.as_ref().expect("contract metadata"); + assert_eq!( + metadata.get("contract"), + Some(&json!("deliverNativeExistingSession")) + ); + assert_eq!(metadata.get("contractVersion"), Some(&json!(1))); + assert_eq!(metadata.get("durableReceipts"), Some(&json!(true))); + assert_eq!(metadata.get("idempotencyField"), Some(&json!("deliveryId"))); + assert_eq!(metadata.get("sessionAuthorization"), Some(&json!("exact"))); + assert_eq!( + metadata.get("reconcileAction"), + Some(&json!( + crate::native_delivery::NATIVE_EXISTING_SESSION_RECONCILE_CAPABILITY + )) + ); + assert!(manifest.capabilities.iter().any(|cap| { + cap.name == crate::native_delivery::NATIVE_EXISTING_SESSION_RECONCILE_CAPABILITY + && cap.kind.as_deref() == Some("action") + })); + } + #[test] fn broker_node_id_prefers_explicit_relay_node_id_env() { let _env_guard = clear_node_id_env(); diff --git a/crates/broker/src/runtime/paths.rs b/crates/broker/src/runtime/paths.rs index 275ee8eccb..13b118accc 100644 --- a/crates/broker/src/runtime/paths.rs +++ b/crates/broker/src/runtime/paths.rs @@ -10,6 +10,12 @@ pub(crate) struct RuntimePaths { /// Inbound-event dedup cache snapshot, so a restart that replays the /// persisted pending file cannot re-inject duplicates. pub(super) dedup: PathBuf, + /// Directory of per-delivery write-ahead receipts for the native + /// existing-session delivery lane. + /// These are persisted independently from ordinary inbox deliveries because + /// a transport timeout after a native session write must never authorize a + /// resend on broker restart. + pub(super) native_delivery_receipts: PathBuf, /// Held for process lifetime to prevent concurrent broker instances (persist mode only). #[allow(dead_code)] pub(super) _lock: Option, @@ -82,6 +88,7 @@ pub(crate) fn ensure_ephemeral_paths(_cwd: &Path, broker_name: &str) -> Result crate::native_delivery::NativeExistingSessionDelivery { + crate::native_delivery::NativeExistingSessionDelivery { + relay_agent_name: "missing-worker".to_string(), + session_id: "native-session-1".to_string(), + delivery_id: "delivery-1".to_string(), + lineage_id: "lineage-1".to_string(), + head_sha: "a".repeat(40), + message: "continue".to_string(), + } +} + +fn empty_worker_registry() -> WorkerRegistry { + let (events, _event_rx) = mpsc::channel(4); + WorkerRegistry::new(events, vec![], std::env::temp_dir(), Instant::now()) +} + +#[tokio::test] +async fn native_delivery_runtime_requires_persistence_and_authorizes_before_receipt() { + use tokio::sync::oneshot; + + let mut ephemeral = worker_event_runtime_fixture(empty_worker_registry(), HashMap::new()); + let ephemeral_receipts = ephemeral.runtime.paths.native_delivery_receipts.clone(); + let (reply, result) = oneshot::channel(); + ephemeral + .runtime + .handle_api_request(ListenApiRequest::DeliverNativeExistingSession { + delivery: native_existing_session_request(), + reply, + }) + .await; + assert!(matches!( + result.await.expect("runtime reply"), + Err(crate::native_delivery::NativeDeliveryError::ReceiptUnavailable(_)) + )); + assert!(!ephemeral_receipts.exists()); + + let mut persistent = worker_event_runtime_fixture(empty_worker_registry(), HashMap::new()); + persistent.runtime.paths.persist = true; + let persistent_receipts = persistent.runtime.paths.native_delivery_receipts.clone(); + let (reply, result) = oneshot::channel(); + persistent + .runtime + .handle_api_request(ListenApiRequest::DeliverNativeExistingSession { + delivery: native_existing_session_request(), + reply, + }) + .await; + assert!(matches!( + result.await.expect("runtime reply"), + Err(crate::native_delivery::NativeDeliveryError::Unauthorized(_)) + )); + assert!( + !persistent_receipts.exists(), + "authorization must precede the write-ahead reservation" + ); +} + +#[tokio::test] +async fn native_delivery_runtime_answers_exact_duplicate_after_worker_exit() { + use tokio::sync::oneshot; + + let mut fixture = worker_event_runtime_fixture(empty_worker_registry(), HashMap::new()); + fixture.runtime.paths.persist = true; + let receipts = fixture.runtime.paths.native_delivery_receipts.clone(); + let queued = crate::native_delivery::reserve_and_deliver( + &receipts, + &native_existing_session_request(), + |_| async { Ok(()) }, + ) + .await + .expect("original delivery"); + + // The worker no longer exists, so authorization fails; the durable + // receipt must still answer the exact retry without a worker write. + let (reply, result) = oneshot::channel(); + fixture + .runtime + .handle_api_request(ListenApiRequest::DeliverNativeExistingSession { + delivery: native_existing_session_request(), + reply, + }) + .await; + let duplicate = result + .await + .expect("runtime reply") + .expect("exact duplicate"); + assert_eq!(duplicate["status"], "duplicate"); + assert_eq!(duplicate["state"], "queued"); + assert_eq!(duplicate["receiptId"], queued.receipt_id); + + let (reply, result) = oneshot::channel(); + let mut reconcile = native_existing_session_request(); + reconcile.message = "changed".to_string(); + fixture + .runtime + .handle_api_request(ListenApiRequest::DeliverNativeExistingSession { + delivery: reconcile, + reply, + }) + .await; + assert!(matches!( + result.await.expect("runtime reply"), + Err(crate::native_delivery::NativeDeliveryError::Conflict) + )); +} + fn delivery_lifecycle_worker_event( name: &str, generation: Uuid, @@ -4714,6 +4821,7 @@ async fn api_spawn_retries_overload_and_only_safe_mode_falls_back() { node_name: "test-node".to_string(), node_token: std::sync::Arc::new(std::sync::RwLock::new(None)), persist: false, + native_delivery_receipts: fixture.runtime.paths.native_delivery_receipts.clone(), node_delivery_probe: std::sync::Arc::new( crate::node_delivery_probe::NodeDeliveryProbe::new(), ), diff --git a/crates/broker/src/runtime/worker_events.rs b/crates/broker/src/runtime/worker_events.rs index 3e4ee83416..1ea5d1ff36 100644 --- a/crates/broker/src/runtime/worker_events.rs +++ b/crates/broker/src/runtime/worker_events.rs @@ -651,6 +651,11 @@ impl BrokerRuntime { error = %error, "worker command writer failed; closing attached terminals and resetting worker" ); + workers.fail_native_delivery_custody_generation( + &name, + generation, + &format!("worker command writer failed: {error}"), + ); let session_ids: Vec = terminal_sessions .iter() .filter(|(_, session)| session.agent == name) @@ -693,6 +698,36 @@ impl BrokerRuntime { return; } if let Some(msg_type) = value.get("type").and_then(Value::as_str) { + if let Some(payload) = value.get("payload") { + let delivery_id = payload + .get("delivery_id") + .and_then(Value::as_str) + .unwrap_or(""); + if !delivery_id.is_empty() { + match msg_type { + "delivery_queued" | "delivery_ack" => { + workers.confirm_native_delivery_custody( + &name, + generation, + delivery_id, + ); + } + "delivery_failed" => { + let reason = payload + .get("reason") + .and_then(Value::as_str) + .unwrap_or("native sidecar rejected delivery"); + workers.fail_native_delivery_custody( + &name, + generation, + delivery_id, + reason, + ); + } + _ => {} + } + } + } if msg_type == "delivery_ack" { if let Some(payload) = value.get("payload") { let delivery_id = payload @@ -1832,6 +1867,11 @@ impl BrokerRuntime { .and_then(|p| p.get("signal")) .and_then(Value::as_str) .map(String::from); + workers.fail_native_delivery_custody_generation( + &name, + generation, + "native worker reported exit before confirming delivery custody", + ); tracing::info!( agent = %name, code = ?code, diff --git a/crates/broker/src/worker.rs b/crates/broker/src/worker.rs index 0bd3496e6c..4ca9ea17cd 100644 --- a/crates/broker/src/worker.rs +++ b/crates/broker/src/worker.rs @@ -2,11 +2,12 @@ use std::{ collections::{HashMap, HashSet, VecDeque}, path::{Path, PathBuf}, process::Stdio, + sync::{Arc, Mutex}, time::{Duration, Instant}, }; use crate::{ - ids::{RequestId, WorkerName}, + ids::{DeliveryId, RequestId, WorkerName}, metrics::MetricsCollector, protocol::{ AgentRuntime, AgentSpec, AppServerAuthType, AppServerHostOwnership, HarnessReleasePolicy, @@ -76,6 +77,11 @@ const WORKER_COMMAND_QUEUE_TIMEOUT: Duration = Duration::from_millis(250); /// slow provider response. The sole stdin writer must still eventually fault /// rather than wedge the worker lane, but should tolerate that short stall. const WORKER_WRITE_TIMEOUT: Duration = Duration::from_secs(5); +/// A native existing-session delivery is not durably queued merely because its +/// frame reached the worker pipe. Wait for the sidecar's correlated +/// `delivery_queued` or `delivery_ack`, which is emitted only after durable +/// custody or immediate acceptance. +const NATIVE_DELIVERY_CUSTODY_TIMEOUT: Duration = Duration::from_secs(10); /// A complete newline-delimited worker protocol frame. A dedicated task owns /// each worker's stdin and writes these frames in order, so cancelling a @@ -265,6 +271,189 @@ pub(crate) struct WorkerRegistry { pub(crate) owned_cleanup_journal: Option, pub(crate) supervisor: Supervisor, pub(crate) metrics: MetricsCollector, + native_delivery_custody: NativeDeliveryCustodyHub, +} + +#[derive(Debug, Clone, PartialEq, Eq, Hash)] +struct NativeDeliveryCustodyKey { + name: WorkerName, + generation: Uuid, + delivery_id: DeliveryId, +} + +type NativeDeliveryCustodyResult = std::result::Result<(), String>; + +/// Correlates the native delivery HTTP task with the sidecar protocol event +/// that proves durable custody. The registry owns the hub and authorized +/// senders hold clones, so a same-name replacement cannot satisfy a waiter for +/// an older process generation. +#[derive(Clone, Default)] +struct NativeDeliveryCustodyHub { + waiters: + Arc>>>, +} + +impl NativeDeliveryCustodyHub { + fn register( + &self, + key: NativeDeliveryCustodyKey, + ) -> Result> { + let (sender, receiver) = oneshot::channel(); + let mut waiters = self + .waiters + .lock() + .map_err(|_| anyhow::anyhow!("native delivery custody registry is unavailable"))?; + anyhow::ensure!( + !waiters.contains_key(&key), + "native delivery custody is already pending for '{}'", + key.delivery_id + ); + waiters.insert(key, sender); + Ok(receiver) + } + + fn resolve(&self, key: &NativeDeliveryCustodyKey, result: NativeDeliveryCustodyResult) -> bool { + let sender = self + .waiters + .lock() + .ok() + .and_then(|mut waiters| waiters.remove(key)); + sender.is_some_and(|sender| sender.send(result).is_ok()) + } + + fn cancel(&self, key: &NativeDeliveryCustodyKey) { + if let Ok(mut waiters) = self.waiters.lock() { + waiters.remove(key); + } + } + + fn fail_generation(&self, name: &WorkerName, generation: Uuid, error: &str) { + let senders = if let Ok(mut waiters) = self.waiters.lock() { + let keys: Vec<_> = waiters + .keys() + .filter(|key| key.name == *name && key.generation == generation) + .cloned() + .collect(); + keys.into_iter() + .filter_map(|key| waiters.remove(&key)) + .collect::>() + } else { + Vec::new() + }; + for sender in senders { + let _ = sender.send(Err(error.to_string())); + } + } +} + +struct NativeDeliveryCustodyRegistration { + custody: NativeDeliveryCustodyHub, + key: Option, +} + +impl NativeDeliveryCustodyRegistration { + fn new(custody: NativeDeliveryCustodyHub, key: NativeDeliveryCustodyKey) -> Self { + Self { + custody, + key: Some(key), + } + } + + fn disarm(&mut self) { + self.key = None; + } +} + +impl Drop for NativeDeliveryCustodyRegistration { + fn drop(&mut self) { + if let Some(key) = self.key.take() { + self.custody.cancel(&key); + } + } +} + +/// Cloneable handle to one authorized worker generation's sole stdin writer. +/// Native delivery tasks use this handle after leaving the broker actor so a +/// stalled pipe cannot block unrelated runtime events. +#[derive(Clone)] +pub(crate) struct WorkerDeliverySender { + name: WorkerName, + generation: Uuid, + command_tx: mpsc::Sender, + custody: NativeDeliveryCustodyHub, +} + +impl WorkerDeliverySender { + pub(crate) async fn deliver(&self, delivery: RelayDelivery) -> Result<()> { + tracing::debug!( + target = "broker::deliver", + worker = %self.name, + generation = %self.generation, + from = %delivery.from, + target = %delivery.target, + event_id = %delivery.event_id, + "delivering event to authorized worker generation" + ); + let delivery_id = delivery.delivery_id.clone(); + let frame = encode_worker_frame("deliver_relay", None, serde_json::to_value(delivery)?)?; + let custody_key = NativeDeliveryCustodyKey { + name: self.name.clone(), + generation: self.generation, + delivery_id, + }; + let custody_rx = self.custody.register(custody_key.clone())?; + let mut custody_registration = + NativeDeliveryCustodyRegistration::new(self.custody.clone(), custody_key); + let (completion_tx, completion_rx) = oneshot::channel(); + let write_result: Result<()> = async { + timeout( + WORKER_COMMAND_QUEUE_TIMEOUT, + self.command_tx.send(WorkerWriteCommand { + frame, + completion: Some(completion_tx), + }), + ) + .await + .map_err(|_| anyhow::anyhow!("worker command queue timed out for '{}'", self.name))? + .map_err(|_| { + anyhow::anyhow!("worker command writer is unavailable for '{}'", self.name) + })?; + completion_rx + .await + .map_err(|_| { + anyhow::anyhow!( + "worker command writer stopped before completing '{}'", + self.name + ) + })? + .map_err(anyhow::Error::msg) + .with_context(|| format!("failed writing frame to worker '{}'", self.name)) + } + .await; + write_result?; + + let custody_result = timeout(NATIVE_DELIVERY_CUSTODY_TIMEOUT, custody_rx).await; + if custody_result.is_ok() { + // Resolution removes the waiter from the hub. A timeout or a + // cancelled delivery future leaves the guard armed so Drop removes + // the registration and an exact retry can register immediately. + custody_registration.disarm(); + } + custody_result + .map_err(|_| { + anyhow::anyhow!( + "native delivery custody confirmation timed out for '{}'", + self.name + ) + })? + .map_err(|_| { + anyhow::anyhow!( + "native delivery custody waiter stopped before confirmation for '{}'", + self.name + ) + })? + .map_err(anyhow::Error::msg) + } } fn encode_worker_frame( @@ -375,6 +564,7 @@ impl WorkerRegistry { owned_cleanup_journal: None, supervisor: Supervisor::new(), metrics: MetricsCollector::new(broker_start), + native_delivery_custody: NativeDeliveryCustodyHub::default(), } } @@ -484,6 +674,94 @@ impl WorkerRegistry { self.workers.contains_key(name) } + /// Authorize the exact broker-owned Codex native session used by the + /// Cloud Babysitter delivery lane. Name-only liveness is insufficient: a + /// released worker can be replaced under the same Relay identity, so the + /// session id and native active-input capability are checked together at + /// the final local hop. + pub(crate) fn authorize_native_existing_session( + &mut self, + name: &WorkerName, + session_id: &str, + ) -> Result { + let handle = self + .workers + .get_mut(name) + .with_context(|| format!("native_session_not_found: no live worker named '{name}'"))?; + let live = match handle.child.try_wait() { + Ok(Some(_)) | Err(_) => false, + Ok(None) => { + #[cfg(unix)] + { + handle.child.id().is_some_and(|pid| !pid_is_gone(pid)) + } + #[cfg(not(unix))] + { + handle.child.id().is_some() + } + } + }; + anyhow::ensure!(live, "native_session_not_live: worker '{name}' is not live"); + anyhow::ensure!( + handle.ready_at.is_some(), + "native_session_not_ready: worker '{name}' has not proved readiness" + ); + authorize_native_existing_session_spec(&handle.spec, session_id)?; + anyhow::ensure!( + !self.initial_tasks.contains_key(name), + "native_session_not_ready: worker initial task has not been queued" + ); + Ok(WorkerDeliverySender { + name: name.clone(), + generation: handle.generation, + command_tx: handle.command_tx.clone(), + custody: self.native_delivery_custody.clone(), + }) + } + + pub(crate) fn confirm_native_delivery_custody( + &self, + name: &WorkerName, + generation: Uuid, + delivery_id: &str, + ) -> bool { + self.native_delivery_custody.resolve( + &NativeDeliveryCustodyKey { + name: name.clone(), + generation, + delivery_id: DeliveryId::from(delivery_id), + }, + Ok(()), + ) + } + + pub(crate) fn fail_native_delivery_custody( + &self, + name: &WorkerName, + generation: Uuid, + delivery_id: &str, + error: &str, + ) -> bool { + self.native_delivery_custody.resolve( + &NativeDeliveryCustodyKey { + name: name.clone(), + generation, + delivery_id: DeliveryId::from(delivery_id), + }, + Err(error.to_string()), + ) + } + + pub(crate) fn fail_native_delivery_custody_generation( + &self, + name: &WorkerName, + generation: Uuid, + error: &str, + ) { + self.native_delivery_custody + .fail_generation(name, generation, error); + } + /// True when a worker is registered AND its child process is still alive. /// Registration alone (`has_worker`) can lag a dead child until the periodic /// `reap_exited` sweep removes it, so callers that must not act on a @@ -1657,10 +1935,20 @@ impl WorkerRegistry { // looking up the handle so maintenance cannot resurrect the released // name after the API has acknowledged teardown. self.supervisor.unregister(name); + let generation = self + .workers + .get(name) + .with_context(|| format!("unknown worker '{name}'"))? + .generation; + self.fail_native_delivery_custody_generation( + &WorkerName::from(name), + generation, + "native worker was released before confirming delivery custody", + ); let mut handle = self .workers .remove(name) - .with_context(|| format!("unknown worker '{name}'"))?; + .expect("worker generation was checked before release"); let release_grace = release_grace_for_spec(&handle.spec); let shutdown_frame = ProtocolEnvelope { @@ -1826,6 +2114,11 @@ impl WorkerRegistry { ), } } + self.fail_native_delivery_custody_generation( + &name, + generation, + "native worker exited before confirming delivery custody", + ); self.workers.remove(&name); self.initial_tasks.remove(&name); self.argv_initial_tasks.remove(&name); @@ -1850,6 +2143,11 @@ impl WorkerRegistry { .workers .get(&name) .and_then(|handle| handle.exit_reason.clone()); + self.fail_native_delivery_custody_generation( + &name, + generation, + "native worker exited before confirming delivery custody", + ); self.workers.remove(&name); self.initial_tasks.remove(&name); self.argv_initial_tasks.remove(&name); @@ -1864,6 +2162,11 @@ impl WorkerRegistry { .workers .get(&name) .and_then(|handle| handle.exit_reason.clone()); + self.fail_native_delivery_custody_generation( + &name, + generation, + "native worker exited before confirming delivery custody", + ); self.workers.remove(&name); self.initial_tasks.remove(&name); self.argv_initial_tasks.remove(&name); @@ -1921,6 +2224,53 @@ pub(crate) fn native_harness_metadata(spec: &AgentSpec) -> Option<(u64, Option Result<()> { + let cli = spec.cli.as_deref().map_or_else( + || Ok(String::new()), + |raw| { + let (command, _) = + parse_cli_command(raw).with_context(|| format!("invalid CLI command '{raw}'"))?; + Ok::<_, anyhow::Error>(normalize_cli_name(&command).to_ascii_lowercase()) + }, + )?; + anyhow::ensure!( + cli == "codex" || cli == "codex.exe", + "native_session_unsupported_harness: only Codex native sessions are supported" + ); + let actual_session = spec.session_id.as_deref().or_else(|| { + spec.harness_config + .as_ref() + .and_then(ResolvedHarnessConfig::session_id) + }); + anyhow::ensure!( + actual_session == Some(session_id), + "native_session_mismatch: requested session does not match the live worker" + ); + let (version, capabilities) = native_harness_metadata(spec).ok_or_else(|| { + anyhow::anyhow!( + "native_session_unsupported_transport: worker is not using the native harness protocol" + ) + })?; + anyhow::ensure!( + version == 1, + "native_session_unsupported_protocol: expected native harness protocol version 1" + ); + let active_input = capabilities + .as_ref() + .and_then(|value| { + value + .get("activeInput") + .or_else(|| value.get("active_input")) + }) + .and_then(Value::as_bool) + .unwrap_or(false); + anyhow::ensure!( + active_input, + "native_session_input_unavailable: native session does not advertise active input" + ); + Ok(()) +} + fn release_policy_arg(policy: Option<&HarnessReleasePolicy>) -> &'static str { match policy { Some(HarnessReleasePolicy::Abort) => "abort", @@ -2859,6 +3209,139 @@ mod tests { WorkerRegistry::new(tx, env, PathBuf::from("/tmp/worker-tests"), Instant::now()) } + #[tokio::test] + async fn native_delivery_waits_for_correlated_sidecar_custody() { + let name = WorkerName::from("native"); + let generation = Uuid::new_v4(); + let custody = NativeDeliveryCustodyHub::default(); + let (command_tx, mut command_rx) = mpsc::channel(1); + let sender = WorkerDeliverySender { + name: name.clone(), + generation, + command_tx, + custody: custody.clone(), + }; + let delivery_id = DeliveryId::from("delivery-1"); + let delivery = RelayDelivery { + delivery_id: delivery_id.clone(), + event_id: "event-1".into(), + workspace_id: None, + workspace_alias: None, + from: "reviewer".to_string(), + target: "native".into(), + body: "status".to_string(), + thread_id: None, + priority: None, + injection_mode: Default::default(), + }; + + let delivery_task = tokio::spawn(async move { sender.deliver(delivery).await }); + let mut command = command_rx.recv().await.expect("worker command"); + command + .completion + .take() + .expect("write completion") + .send(Ok(())) + .expect("delivery task should await write completion"); + tokio::task::yield_now().await; + assert!( + !delivery_task.is_finished(), + "pipe write alone must not claim durable custody" + ); + + assert!(custody.resolve( + &NativeDeliveryCustodyKey { + name, + generation, + delivery_id, + }, + Ok(()), + )); + delivery_task + .await + .expect("delivery task should join") + .expect("correlated sidecar confirmation should complete delivery"); + } + + #[tokio::test] + async fn cancelled_native_delivery_removes_its_custody_waiter() { + let name = WorkerName::from("native"); + let generation = Uuid::new_v4(); + let custody = NativeDeliveryCustodyHub::default(); + let (command_tx, mut command_rx) = mpsc::channel(1); + let sender = WorkerDeliverySender { + name: name.clone(), + generation, + command_tx, + custody: custody.clone(), + }; + let delivery_id = DeliveryId::from("delivery-cancelled"); + let custody_key = NativeDeliveryCustodyKey { + name, + generation, + delivery_id: delivery_id.clone(), + }; + let delivery = RelayDelivery { + delivery_id, + event_id: "event-cancelled".into(), + workspace_id: None, + workspace_alias: None, + from: "reviewer".to_string(), + target: "native".into(), + body: "status".to_string(), + thread_id: None, + priority: None, + injection_mode: Default::default(), + }; + + let delivery_task = tokio::spawn(async move { sender.deliver(delivery).await }); + let mut command = command_rx.recv().await.expect("worker command"); + command + .completion + .take() + .expect("write completion") + .send(Ok(())) + .expect("delivery task should await custody"); + tokio::task::yield_now().await; + delivery_task.abort(); + assert!(delivery_task + .await + .expect_err("delivery should be cancelled") + .is_cancelled()); + + let replacement = custody + .register(custody_key.clone()) + .expect("cancelled delivery must not block an exact retry"); + custody.cancel(&custody_key); + assert!( + replacement.await.is_err(), + "cancelling the replacement waiter should close its receiver" + ); + } + + #[tokio::test] + async fn native_delivery_custody_fails_when_worker_generation_exits() { + let name = WorkerName::from("native"); + let generation = Uuid::new_v4(); + let custody = NativeDeliveryCustodyHub::default(); + let receiver = custody + .register(NativeDeliveryCustodyKey { + name: name.clone(), + generation, + delivery_id: DeliveryId::from("delivery-exit"), + }) + .expect("custody waiter should register"); + + custody.fail_generation(&name, generation, "worker exited"); + assert_eq!( + receiver + .await + .expect("exit should resolve custody waiter") + .expect_err("exit cannot confirm custody"), + "worker exited" + ); + } + #[cfg(unix)] fn git(repo: &Path, args: &[&str]) -> std::process::Output { std::process::Command::new("git") @@ -3684,6 +4167,49 @@ sleep 30 assert!(reg.supervisor.pending_restarts().is_empty()); } + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn release_fails_native_delivery_custody_waiters_before_removal() { + let mut registry = make_registry(Vec::new()); + let name = "released-native-worker"; + registry + .spawn( + sleeping_native_worker(name, None), + None, + None, + None, + true, + None, + None, + None, + None, + ) + .await + .expect("native worker should spawn"); + let generation = registry + .workers + .get(name) + .expect("spawned worker") + .generation; + let receiver = registry + .native_delivery_custody + .register(NativeDeliveryCustodyKey { + name: WorkerName::from(name), + generation, + delivery_id: DeliveryId::from("delivery-release"), + }) + .expect("custody waiter should register"); + + registry.release(name).await.expect("worker should release"); + assert_eq!( + receiver + .await + .expect("release should resolve custody waiter") + .expect_err("release cannot confirm custody"), + "native worker was released before confirming delivery custody" + ); + } + #[test] fn worker_log_path_rejects_path_traversal() { let reg = make_registry(vec![]); @@ -3780,6 +4306,108 @@ sleep 30 ); } + fn native_codex_authorization_spec() -> AgentSpec { + serde_json::from_value(json!({ + "name": "garden-coder", + "runtime": "headless", + "cli": "codex", + "sessionId": "native-1", + "args": [], + "channels": [], + "harnessConfig": { + "runtime": "native", + "command": "node", + "args": ["/tmp/sidecar.js"], + "sessionId": "native-1", + "metadata": { + "runtimeKind": "native", + "nativeHarnessProtocolVersion": 1, + "nativeHarnessCapabilities": {"activeInput": true} + } + } + })) + .expect("native Codex spec") + } + + #[test] + fn native_existing_session_authorization_requires_exact_session_and_active_input() { + let mut spec = native_codex_authorization_spec(); + authorize_native_existing_session_spec(&spec, "native-1") + .expect("exact native Codex session should authorize"); + + spec.cli = Some("/usr/local/bin/codex --model o3".to_string()); + authorize_native_existing_session_spec(&spec, "native-1") + .expect("inline Codex command should authorize by executable"); + + let mismatch = authorize_native_existing_session_spec(&spec, "replacement-session") + .expect_err("session substitution must fail") + .to_string(); + assert!(mismatch.contains("native_session_mismatch"), "{mismatch}"); + + let mut no_input = spec.clone(); + if let Some(ResolvedHarnessConfig::Native(config)) = no_input.harness_config.as_mut() { + config.metadata.as_mut().expect("metadata").insert( + "nativeHarnessCapabilities".to_string(), + json!({"activeInput": false}), + ); + } + let unavailable = authorize_native_existing_session_spec(&no_input, "native-1") + .expect_err("inactive input must fail") + .to_string(); + assert!( + unavailable.contains("native_session_input_unavailable"), + "{unavailable}" + ); + } + + #[test] + fn native_existing_session_authorization_rejects_non_codex_harnesses() { + let mut spec = native_codex_authorization_spec(); + spec.cli = Some("claude".to_string()); + let error = authorize_native_existing_session_spec(&spec, "native-1") + .expect_err("non-Codex native session must fail closed") + .to_string(); + assert!( + error.contains("native_session_unsupported_harness"), + "{error}" + ); + } + + #[cfg(unix)] + #[tokio::test] + async fn native_existing_session_authorization_reaps_an_exited_worker() { + let mut registry = make_registry(vec![]); + let name = WorkerName::from("exited-native-worker"); + let child = Command::new("true").spawn().expect("spawn exiting child"); + let generation = Uuid::new_v4(); + let (command_tx, _command_rx) = mpsc::channel(WORKER_WRITE_QUEUE_CAPACITY); + registry.workers.insert( + name.clone(), + WorkerHandle { + generation, + spec: native_codex_authorization_spec(), + parent: None, + workspace_id: None, + child, + command_tx, + harness_pid: None, + spawned_at: Instant::now(), + ready_at: Some(Instant::now()), + last_activity_at: Instant::now(), + context_budget_pct: None, + state: AgentWorkState::Idle, + exit_reason: None, + }, + ); + tokio::time::sleep(Duration::from_millis(50)).await; + + let error = match registry.authorize_native_existing_session(&name, "native-1") { + Ok(_) => panic!("an exited child must fail before receipt reservation"), + Err(error) => error.to_string(), + }; + assert!(error.contains("native_session_not_live"), "{error}"); + } + #[test] fn app_server_config_validation_rejects_missing_bearer_token() { let mut config = make_app_server_config(); diff --git a/packages/harnesses/src/ai-sdk/relay-session.test.ts b/packages/harnesses/src/ai-sdk/relay-session.test.ts index 92ae2e34e7..1f37de7325 100644 --- a/packages/harnesses/src/ai-sdk/relay-session.test.ts +++ b/packages/harnesses/src/ai-sdk/relay-session.test.ts @@ -1,4 +1,8 @@ import type { AgentIdentity, MessageContext, RelayMessage } from '@agent-relay/sdk'; +import { createHash } from 'node:crypto'; +import { mkdir, mkdtemp, readdir, readFile, rm, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { resolve } from 'node:path'; import { describe, expect, it, vi } from 'vitest'; import { formatInboundRelayPrompt, RelayHarnessSession } from './relay-session.js'; @@ -17,6 +21,36 @@ function context(id: string, mode: MessageContext['mode'] = 'immediate'): Messag return { id, mode, reason: 'message' }; } +async function readEntries(path: string): Promise>> { + const entries: Array> = []; + for (const directory of ['queue', 'receipts']) { + const entryPath = resolve(path, directory); + let files: string[]; + try { + files = (await readdir(entryPath)).filter((file) => file.endsWith('.json')); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') continue; + throw error; + } + for (const file of files) { + const envelope = JSON.parse(await readFile(resolve(entryPath, file), 'utf8')) as { + entry: Record; + }; + entries.push(envelope.entry); + } + } + return entries; +} + +async function writeEntry(path: string, entry: Record): Promise { + const directory = entry.state === 'queued' ? 'queue' : 'receipts'; + const entryPath = resolve(path, directory); + await mkdir(entryPath, { recursive: true }); + const key = String(entry.key); + const file = `${createHash('sha256').update(key).digest('hex')}.json`; + await writeFile(resolve(entryPath, file), JSON.stringify({ version: 2, entry })); +} + function fakeHost() { const listeners = new Set<(event: never) => void>(); let active = false; @@ -128,6 +162,248 @@ describe('RelayHarnessSession', () => { expect(fixture.host.startTurn).toHaveBeenLastCalledWith(expect.stringContaining('queued-2'), 'queued-2'); }); + it('restores a durably deferred on-idle message after a sidecar restart', async () => { + const root = await mkdtemp(resolve(tmpdir(), 'relay-deferred-')); + const queuePath = resolve(root, 'queue'); + const first = fakeHost(); + const firstSession = new RelayHarnessSession({ + identity, + host: first.host as never, + deferredQueuePath: queuePath, + }); + await firstSession.receiveMessage(message('active'), context('active')); + await expect( + firstSession.receiveMessage(message('durable'), context('durable', 'on-idle')) + ).resolves.toMatchObject({ status: 'deferred' }); + expect(await readEntries(queuePath)).toMatchObject([ + expect.objectContaining({ key: 'durable', state: 'queued' }), + ]); + + const restarted = fakeHost(); + const restartedSession = new RelayHarnessSession({ + identity, + host: restarted.host as never, + deferredQueuePath: queuePath, + }); + let stateAtAcceptance: string | undefined; + restartedSession.onEvent?.(async (event) => { + if (event.type !== 'delivery.accepted') return; + const persisted = await readEntries(queuePath); + stateAtAcceptance = persisted[0]?.state as string | undefined; + }); + await restartedSession.restoreDeferredMessages(); + expect(restarted.host.startTurn).toHaveBeenCalledWith( + expect.stringContaining('"messageId":"durable"'), + 'durable' + ); + expect(await readEntries(queuePath)).toMatchObject([ + expect.objectContaining({ key: 'durable', state: 'accepted' }), + ]); + expect(stateAtAcceptance).toBe('accepted'); + }); + + it('fails closed instead of replaying an in-flight deferred message after restart', async () => { + const root = await mkdtemp(resolve(tmpdir(), 'relay-deferred-indoubt-')); + const queuePath = resolve(root, 'queue'); + await writeEntry(queuePath, { + key: 'ambiguous', + deliveryId: 'ambiguous', + messageId: 'ambiguous', + state: 'in_doubt', + }); + await writeEntry(queuePath, { + key: 'already-accepted', + deliveryId: 'already-accepted', + messageId: 'already-accepted', + state: 'accepted', + }); + const restarted = fakeHost(); + const restartedSession = new RelayHarnessSession({ + identity, + host: restarted.host as never, + deferredQueuePath: queuePath, + }); + await restartedSession.restoreDeferredMessages(); + expect(restarted.host.startTurn).not.toHaveBeenCalled(); + expect(await readEntries(queuePath)).toEqual( + expect.arrayContaining([ + expect.objectContaining({ key: 'ambiguous', state: 'in_doubt' }), + expect.objectContaining({ key: 'already-accepted', state: 'accepted' }), + ]) + ); + await expect( + restartedSession.receiveMessage(message('ambiguous'), { + ...context('ambiguous', 'on-idle'), + idempotencyKey: 'ambiguous', + }) + ).resolves.toMatchObject({ status: 'failed', retryable: false }); + await expect( + restartedSession.receiveMessage(message('already-accepted'), { + ...context('already-accepted', 'on-idle'), + idempotencyKey: 'already-accepted', + }) + ).resolves.toMatchObject({ status: 'accepted' }); + expect(restarted.host.startTurn).not.toHaveBeenCalled(); + + const restartedAgain = fakeHost(); + const restartedAgainSession = new RelayHarnessSession({ + identity, + host: restartedAgain.host as never, + deferredQueuePath: queuePath, + maxDedupeEntries: 1, + }); + await restartedAgainSession.restoreDeferredMessages(); + await expect( + restartedAgainSession.receiveMessage(message('ambiguous'), { + ...context('ambiguous', 'on-idle'), + idempotencyKey: 'ambiguous', + }) + ).resolves.toMatchObject({ status: 'failed', retryable: false }); + await expect( + restartedAgainSession.receiveMessage(message('already-accepted'), { + ...context('already-accepted', 'on-idle'), + idempotencyKey: 'already-accepted', + }) + ).resolves.toMatchObject({ status: 'accepted' }); + expect(restartedAgain.host.startTurn).not.toHaveBeenCalled(); + }); + + it('retains a failed in-flight deferred message as a non-retryable tombstone', async () => { + const root = await mkdtemp(resolve(tmpdir(), 'relay-deferred-failed-')); + const queuePath = resolve(root, 'queue'); + const fixture = fakeHost(); + const session = new RelayHarnessSession({ + identity, + host: fixture.host as never, + deferredQueuePath: queuePath, + }); + await session.receiveMessage(message('active'), context('active')); + await session.receiveMessage(message('ambiguous'), { + ...context('ambiguous', 'on-idle'), + idempotencyKey: 'ambiguous', + }); + const failed = new Promise((resolveFailed) => { + session.onEvent?.((event) => { + if (event.type === 'delivery.failed' && event.deliveryId === 'ambiguous') resolveFailed(); + }); + }); + vi.mocked(fixture.host.startTurn).mockRejectedValueOnce(new Error('acceptance failed')); + fixture.settle(); + await failed; + + const tombstones = await readEntries(queuePath); + expect(tombstones).toMatchObject([expect.objectContaining({ key: 'ambiguous', state: 'in_doubt' })]); + expect(tombstones[0]).not.toHaveProperty('message'); + expect(tombstones[0]).not.toHaveProperty('context'); + await expect( + session.receiveMessage(message('ambiguous'), { + ...context('ambiguous', 'on-idle'), + idempotencyKey: 'ambiguous', + }) + ).resolves.toMatchObject({ status: 'failed', retryable: false }); + expect(fixture.host.startTurn).toHaveBeenCalledTimes(2); + await session.release?.('retired'); + expect(await readEntries(queuePath)).toEqual([]); + }); + + it('does not report failure after the host accepts when the terminal receipt cannot persist', async () => { + const root = await mkdtemp(resolve(tmpdir(), 'relay-deferred-accepted-persist-failure-')); + const queuePath = resolve(root, 'queue'); + const fixture = fakeHost(); + const session = new RelayHarnessSession({ + identity, + host: fixture.host as never, + deferredQueuePath: queuePath, + }); + await session.receiveMessage(message('active'), context('active')); + await session.receiveMessage(message('accepted'), { + ...context('accepted', 'on-idle'), + idempotencyKey: 'accepted', + }); + + const events: string[] = []; + let resolveAccepted!: () => void; + const accepted = new Promise((resolveEvent) => { + resolveAccepted = resolveEvent; + }); + session.onEvent?.(async (event) => { + events.push(event.type); + if (event.type === 'message.received' && event.message.id === 'accepted') { + const receipts = resolve(queuePath, 'receipts'); + await rm(receipts, { recursive: true, force: true }); + await writeFile(receipts, 'blocked'); + } + if (event.type === 'delivery.accepted' && event.deliveryId === 'accepted') resolveAccepted(); + }); + + fixture.settle(); + await accepted; + expect(events).toContain('delivery.accepted'); + expect(events).not.toContain('delivery.failed'); + await expect( + session.receiveMessage(message('accepted'), { + ...context('accepted', 'on-idle'), + idempotencyKey: 'accepted', + }) + ).resolves.toMatchObject({ status: 'accepted' }); + expect(fixture.host.startTurn).toHaveBeenCalledTimes(2); + await session.release?.('retired'); + }); + + it('keeps a durable queued delivery live when the in-doubt tombstone cannot persist', async () => { + const root = await mkdtemp(resolve(tmpdir(), 'relay-deferred-reservation-failure-')); + const queuePath = resolve(root, 'queue'); + const fixture = fakeHost(); + const session = new RelayHarnessSession({ + identity, + host: fixture.host as never, + deferredQueuePath: queuePath, + }); + await session.receiveMessage(message('active'), context('active')); + await session.receiveMessage(message('retry-after-recovery'), { + ...context('retry-after-recovery', 'on-idle'), + idempotencyKey: 'retry-after-recovery', + }); + const receipts = resolve(queuePath, 'receipts'); + await writeFile(receipts, 'blocked'); + const events: string[] = []; + let resolveAccepted!: () => void; + const accepted = new Promise((resolveEvent) => { + resolveAccepted = resolveEvent; + }); + session.onEvent?.((event) => { + events.push(event.type); + if (event.type === 'delivery.accepted' && event.deliveryId === 'retry-after-recovery') { + resolveAccepted(); + } + }); + + fixture.settle(); + await new Promise((resolveWait) => setTimeout(resolveWait, 0)); + expect(events).not.toContain('delivery.failed'); + expect(fixture.host.startTurn).toHaveBeenCalledTimes(1); + const queuedFiles = (await readdir(resolve(queuePath, 'queue'))).filter((file) => file.endsWith('.json')); + expect(queuedFiles).toHaveLength(1); + expect( + JSON.parse(await readFile(resolve(queuePath, 'queue', queuedFiles[0]!), 'utf8')).entry + ).toMatchObject({ key: 'retry-after-recovery', state: 'queued' }); + await expect( + session.receiveMessage(message('retry-after-recovery'), { + ...context('retry-after-recovery', 'on-idle'), + idempotencyKey: 'retry-after-recovery', + }) + ).resolves.toMatchObject({ status: 'deferred' }); + + await rm(receipts, { force: true }); + await accepted; + expect(fixture.host.startTurn).toHaveBeenCalledWith( + expect.stringContaining('"messageId":"retry-after-recovery"'), + 'retry-after-recovery' + ); + expect(events).toContain('delivery.accepted'); + expect(events).not.toContain('delivery.failed'); + }); + it('publishes capabilities and releases once', async () => { const fixture = fakeHost(); const session = new RelayHarnessSession({ identity, host: fixture.host as never }); @@ -137,4 +413,30 @@ describe('RelayHarnessSession', () => { await Promise.all([session.release?.('done'), session.release?.('again')]); expect(fixture.host.destroy).toHaveBeenCalledTimes(1); }); + + it('destroys the host even when durable queue cleanup fails', async () => { + const root = await mkdtemp(resolve(tmpdir(), 'relay-release-failure-')); + const queuePath = resolve(root, 'queue'); + const fixture = fakeHost(); + const session = new RelayHarnessSession({ + identity, + host: fixture.host as never, + deferredQueuePath: queuePath, + }); + const events: string[] = []; + session.onEvent?.((event) => events.push(event.type)); + await session.receiveMessage(message('active'), context('active')); + await session.receiveMessage(message('queued'), context('queued', 'on-idle')); + await rm(queuePath, { recursive: true }); + await writeFile(queuePath, 'blocked'); + + await expect(session.release?.('done')).rejects.toBeDefined(); + expect(fixture.host.destroy).toHaveBeenCalledTimes(1); + expect(events).toContain('session.released'); + await rm(queuePath); + await mkdir(resolve(queuePath, 'queue'), { recursive: true }); + await expect(session.release?.('retry')).resolves.toBeUndefined(); + expect(fixture.host.destroy).toHaveBeenCalledTimes(1); + expect(events.filter((event) => event === 'session.released')).toHaveLength(1); + }); }); diff --git a/packages/harnesses/src/ai-sdk/relay-session.ts b/packages/harnesses/src/ai-sdk/relay-session.ts index e1854724c7..a6adab6d23 100644 --- a/packages/harnesses/src/ai-sdk/relay-session.ts +++ b/packages/harnesses/src/ai-sdk/relay-session.ts @@ -1,4 +1,7 @@ import { createAgentActivityState, reduceAgentActivity } from '@agent-relay/sdk'; +import { createHash } from 'node:crypto'; +import { mkdir, mkdtemp, open, readdir, readFile, rename, rm, stat } from 'node:fs/promises'; +import { dirname, resolve } from 'node:path'; import type { AgentIdentity, AgentActivityState, @@ -18,12 +21,29 @@ export interface RelayHarnessSessionOptions { host: HarnessHost; maxQueueSize?: number; maxDedupeEntries?: number; + /** Durable queue used for on-idle messages accepted by a native sidecar. */ + deferredQueuePath?: string; } interface QueuedMessage { message: RelayMessage; context: MessageContext; key: string; + state: 'queued'; + order: number; +} + +interface DeliveryTombstone { + key: string; + deliveryId: string; + messageId: string; + state: 'in_doubt' | 'accepted'; +} + +type DurableDeliveryEntry = QueuedMessage | DeliveryTombstone; + +function durableEntryId(key: string): string { + return createHash('sha256').update(key).digest('hex'); } function senderName(message: RelayMessage): string { @@ -302,11 +322,18 @@ export class RelayHarnessSession implements AgentSession { readonly host: HarnessHost; readonly #maxQueueSize: number; readonly #maxDedupeEntries: number; + readonly #deferredQueuePath?: string; readonly #queue: QueuedMessage[] = []; readonly #receipts = new Map(); readonly #listeners = new Set<(event: AgentSessionEvent) => void | Promise>(); #operation = Promise.resolve(); + #nextQueueOrder = 0; #released = false; + #releaseCompleted = false; + #hostDestroyed = false; + #releaseEmitted = false; + #drainRetry?: ReturnType; + #drainRetryDelayMs = 100; #activity: AgentActivityState = createAgentActivityState(); constructor(options: RelayHarnessSessionOptions) { @@ -314,6 +341,7 @@ export class RelayHarnessSession implements AgentSession { this.host = options.host; this.#maxQueueSize = options.maxQueueSize ?? 100; this.#maxDedupeEntries = options.maxDedupeEntries ?? 10_000; + this.#deferredQueuePath = options.deferredQueuePath; this.capabilities = { messaging: { receive: true }, delivery: { @@ -376,13 +404,225 @@ export class RelayHarnessSession implements AgentSession { return receipt; } + #scheduleDrainRetry(): void { + if (this.#released || this.#drainRetry) return; + const delay = this.#drainRetryDelayMs; + this.#drainRetryDelayMs = Math.min(delay * 2, 5_000); + this.#drainRetry = setTimeout(() => { + this.#drainRetry = undefined; + void this.#serialized(() => this.#drain()).catch(() => undefined); + }, delay); + this.#drainRetry.unref?.(); + } + + #resetDrainRetry(): void { + if (this.#drainRetry) clearTimeout(this.#drainRetry); + this.#drainRetry = undefined; + this.#drainRetryDelayMs = 100; + } + + #durableReceipt(entry: DurableDeliveryEntry): MessageReceipt { + if (entry.state === 'queued') { + return { + status: 'deferred', + deliveryId: entry.context.id, + availableAt: new Date(Date.now() + 100).toISOString(), + reason: 'queued_until_idle', + metadata: { queued: true, restored: true }, + }; + } + if (entry.state === 'accepted') { + return { + status: 'accepted', + deliveryId: entry.deliveryId, + metadata: { restored: true }, + }; + } + return { + status: 'failed', + deliveryId: entry.deliveryId, + reason: 'Deferred delivery was in progress when the native sidecar stopped', + retryable: false, + }; + } + + #entryDirectory(state: DurableDeliveryEntry['state']): string | undefined { + if (!this.#deferredQueuePath) return undefined; + return resolve(this.#deferredQueuePath, state === 'queued' ? 'queue' : 'receipts'); + } + + #entryPath(key: string, state: DurableDeliveryEntry['state']): string | undefined { + const directory = this.#entryDirectory(state); + return directory ? resolve(directory, `${durableEntryId(key)}.json`) : undefined; + } + + async #syncEntryDirectories(directory: string): Promise { + if (!this.#deferredQueuePath || process.platform === 'win32') return; + const parent = dirname(this.#deferredQueuePath); + // The sidecar path is runtimeRoot/deferred-relay/. Syncing + // every newly-created directory level makes a first-use persist durable as + // well as the entry rename itself. + for (const path of [directory, this.#deferredQueuePath, parent, dirname(parent)]) { + const directory = await open(path, 'r'); + try { + await directory.sync(); + } finally { + await directory.close(); + } + } + } + + async #persistEntry(entry: DurableDeliveryEntry): Promise { + const directory = this.#entryDirectory(entry.state); + const destination = this.#entryPath(entry.key, entry.state); + if (!directory || !destination || !this.#deferredQueuePath) return; + await mkdir(directory, { recursive: true }); + const temporaryDirectory = await mkdtemp(resolve(directory, '.relay-deferred-')); + const temporary = resolve(temporaryDirectory, 'entry.json'); + try { + const file = await open(temporary, 'wx', 0o600); + try { + await file.writeFile(JSON.stringify({ version: 2, entry }), 'utf8'); + await file.sync(); + } finally { + await file.close(); + } + await rename(temporary, destination); + await this.#syncEntryDirectories(directory); + if (entry.state !== 'queued') { + const queued = this.#entryPath(entry.key, 'queued'); + if (queued) { + // The compact terminal receipt is authoritative once published. + // Cleanup cannot revoke it, and restore checks receipts before live + // queue entries, so a cleanup failure must not reverse the result. + await rm(queued, { force: true }).catch(() => undefined); + const queueDirectory = this.#entryDirectory('queued'); + if (queueDirectory) { + await this.#syncEntryDirectories(queueDirectory).catch(() => undefined); + } + } + } + } finally { + await rm(temporaryDirectory, { recursive: true, force: true }).catch(() => undefined); + } + } + + #parseEntry(value: unknown): DurableDeliveryEntry { + if (!value || typeof value !== 'object') throw new Error('Invalid deferred Relay delivery entry'); + const envelope = value as { version?: unknown; entry?: unknown }; + if (envelope.version !== 2 || !envelope.entry || typeof envelope.entry !== 'object') { + throw new Error('Invalid deferred Relay delivery entry'); + } + const entry = envelope.entry as { + key?: unknown; + state?: unknown; + message?: unknown; + context?: unknown; + order?: unknown; + deliveryId?: unknown; + messageId?: unknown; + }; + if (typeof entry.key !== 'string') throw new Error('Invalid deferred Relay delivery entry'); + if (entry.state === 'queued') { + if (!entry.message || !entry.context || typeof entry.order !== 'number') { + throw new Error('Invalid deferred Relay delivery entry'); + } + return entry as QueuedMessage; + } + if ( + (entry.state === 'accepted' || entry.state === 'in_doubt') && + typeof entry.deliveryId === 'string' && + typeof entry.messageId === 'string' + ) { + return entry as DeliveryTombstone; + } + throw new Error('Invalid deferred Relay delivery entry'); + } + + async #loadEntry(key: string): Promise { + for (const state of ['accepted', 'queued'] as const) { + const path = this.#entryPath(key, state); + if (!path) return undefined; + try { + const entry = this.#parseEntry(JSON.parse(await readFile(path, 'utf8'))); + if (entry.key !== key) throw new Error('Deferred Relay delivery key does not match its index'); + return entry; + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error; + } + } + return undefined; + } + + async #removeEntry(key: string): Promise { + const path = this.#entryPath(key, 'queued'); + const directory = this.#entryDirectory('queued'); + if (!path || !directory) return; + await rm(path, { force: true }); + await this.#syncEntryDirectories(directory); + } + + async #removeDurableSession(): Promise { + if (!this.#deferredQueuePath) return; + try { + await stat(this.#deferredQueuePath); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return; + throw error; + } + await rm(this.#deferredQueuePath, { recursive: true, force: true }); + if (process.platform === 'win32') return; + const parent = dirname(this.#deferredQueuePath); + for (const path of [parent, dirname(parent)]) { + const directory = await open(path, 'r'); + try { + await directory.sync(); + } finally { + await directory.close(); + } + } + } + + /** Restore deferred messages after a sidecar restart before reading stdin. */ + async restoreDeferredMessages(): Promise { + if (!this.#deferredQueuePath) return; + const queueDirectory = this.#entryDirectory('queued'); + if (!queueDirectory) return; + let files: string[]; + try { + files = await readdir(queueDirectory); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return; + throw error; + } + let cleanedStaleEntry = false; + for (const file of files.filter((candidate) => candidate.endsWith('.json'))) { + const entry = this.#parseEntry(JSON.parse(await readFile(resolve(queueDirectory, file), 'utf8'))); + if (file !== `${durableEntryId(entry.key)}.json`) { + throw new Error('Deferred Relay delivery key does not match its index'); + } + if (entry.state !== 'queued') throw new Error('Invalid queued Relay delivery entry'); + const current = await this.#loadEntry(entry.key); + if (current?.state !== 'queued') { + await rm(resolve(queueDirectory, file), { force: true }).catch(() => undefined); + cleanedStaleEntry = true; + continue; + } + this.#queue.push(current); + this.#nextQueueOrder = Math.max(this.#nextQueueOrder, current.order + 1); + this.#remember(current.key, this.#durableReceipt(current)); + } + if (cleanedStaleEntry) await this.#syncEntryDirectories(queueDirectory).catch(() => undefined); + this.#queue.sort((left, right) => left.order - right.order || left.key.localeCompare(right.key)); + await this.#serialized(() => this.#drain()); + } + async #accept(message: RelayMessage, context: MessageContext): Promise { const prompt = formatInboundRelayPrompt(message); if (this.host.hasActiveTurn) await this.host.submitUserMessage(prompt); else await this.host.startTurn(prompt, context.id); const receipt: MessageReceipt = { status: 'accepted', deliveryId: context.id }; await this.#emit({ type: 'message.received', message }); - await this.#emit({ type: 'delivery.accepted', messageId: message.id, deliveryId: context.id }); return receipt; } @@ -390,15 +630,34 @@ export class RelayHarnessSession implements AgentSession { if (this.#released || this.host.hasActiveTurn) return; const queued = this.#queue.shift(); if (!queued) return; + const inDoubt: DeliveryTombstone = { + key: queued.key, + deliveryId: queued.context.id, + messageId: queued.message.id, + state: 'in_doubt', + }; try { - const receipt = await this.#accept(queued.message, queued.context); - this.#remember(queued.key, receipt); + await this.#persistEntry(inDoubt); + } catch { + // The durable queued entry is still authoritative. Without a durable + // in-doubt marker it is unsafe to publish a terminal failure: a restart + // could restore and accept the queued entry. Keep it live and retry the + // transition with bounded backoff or when the broker redelivers it. + this.#queue.unshift(queued); + this.#remember(queued.key, this.#durableReceipt(queued)); + this.#scheduleDrainRetry(); + return; + } + this.#resetDrainRetry(); + let receipt: MessageReceipt; + try { + receipt = await this.#accept(queued.message, queued.context); } catch (error) { const receipt: MessageReceipt = { status: 'failed', deliveryId: queued.context.id, reason: error instanceof Error ? error.message : String(error), - retryable: true, + retryable: false, }; this.#remember(queued.key, receipt); await this.#emit({ @@ -406,10 +665,27 @@ export class RelayHarnessSession implements AgentSession { messageId: queued.message.id, deliveryId: queued.context.id, reason: receipt.reason, - retryable: true, + retryable: false, }); await this.#drain(); + return; + } + + const accepted: DeliveryTombstone = { ...inDoubt, state: 'accepted' }; + try { + await this.#persistEntry(accepted); + } catch { + // The host already accepted the message. The durable in-doubt marker is + // sufficient to prevent replay if the accepted transition cannot be + // published; do not report a false delivery failure for a turn that has + // started successfully. } + this.#remember(queued.key, receipt); + await this.#emit({ + type: 'delivery.accepted', + messageId: queued.message.id, + deliveryId: queued.context.id, + }); } receiveMessage(message: RelayMessage, context: MessageContext): Promise { @@ -419,7 +695,18 @@ export class RelayHarnessSession implements AgentSession { } const key = context.idempotencyKey ?? message.id ?? context.id; const previous = this.#receipts.get(key); - if (previous) return previous; + if (previous) { + if (previous.status === 'deferred' && !this.host.hasActiveTurn) await this.#drain(); + return this.#receipts.get(key) ?? previous; + } + const active = this.#queue.find((entry) => entry.key === key); + if (active) { + const receipt = this.#remember(key, this.#durableReceipt(active)); + if (!this.host.hasActiveTurn) await this.#drain(); + return this.#receipts.get(key) ?? receipt; + } + const durable = await this.#loadEntry(key); + if (durable) return this.#remember(key, this.#durableReceipt(durable)); const shouldQueue = this.host.hasActiveTurn && (context.mode === 'next-message' || context.mode === 'on-idle'); @@ -432,7 +719,32 @@ export class RelayHarnessSession implements AgentSession { retryable: true, }); } - this.#queue.push({ message, context, key }); + const queued: QueuedMessage = { + message, + context, + key, + state: 'queued', + order: this.#nextQueueOrder++, + }; + this.#queue.push(queued); + try { + await this.#persistEntry(queued); + } catch (error) { + this.#queue.pop(); + const inDoubt: DeliveryTombstone = { + key, + deliveryId: context.id, + messageId: message.id, + state: 'in_doubt', + }; + await this.#persistEntry(inDoubt).catch(() => undefined); + return this.#remember(key, { + status: 'failed', + deliveryId: context.id, + reason: `Could not durably queue Relay delivery: ${error instanceof Error ? error.message : String(error)}`, + retryable: false, + }); + } return this.#remember(key, { status: 'deferred', deliveryId: context.id, @@ -443,7 +755,13 @@ export class RelayHarnessSession implements AgentSession { } try { - return this.#remember(key, await this.#accept(message, context)); + const receipt = this.#remember(key, await this.#accept(message, context)); + await this.#emit({ + type: 'delivery.accepted', + messageId: message.id, + deliveryId: context.id, + }); + return receipt; } catch (error) { return this.#remember(key, { status: 'failed', @@ -457,11 +775,48 @@ export class RelayHarnessSession implements AgentSession { async release(reason?: string): Promise { return this.#serialized(async () => { - if (this.#released) return; + if (this.#releaseCompleted) return; this.#released = true; - this.#queue.length = 0; - await this.host.destroy(); - await this.#emit({ type: 'session.released', reason }); + this.#resetDrainRetry(); + const queued = this.#queue.splice(0); + let persistenceError: unknown; + for (const entry of queued) { + try { + await this.#removeEntry(entry.key); + } catch (error) { + persistenceError ??= error; + this.#queue.push(entry); + } + } + if (!persistenceError) { + try { + await this.#removeDurableSession(); + } catch (error) { + persistenceError = error; + } + } + let destroyError: unknown; + if (!this.#hostDestroyed) { + try { + await this.host.destroy(); + this.#hostDestroyed = true; + } catch (error) { + destroyError = error; + } + } + let emitError: unknown; + if (!this.#releaseEmitted) { + try { + await this.#emit({ type: 'session.released', reason }); + this.#releaseEmitted = true; + } catch (error) { + emitError = error; + } + } + if (persistenceError) throw persistenceError; + if (destroyError) throw destroyError; + if (emitError) throw emitError; + this.#releaseCompleted = true; }); } } diff --git a/packages/harnesses/src/ai-sdk/sidecar.test.ts b/packages/harnesses/src/ai-sdk/sidecar.test.ts index f55ba6434d..ded75f2605 100644 --- a/packages/harnesses/src/ai-sdk/sidecar.test.ts +++ b/packages/harnesses/src/ai-sdk/sidecar.test.ts @@ -1,4 +1,5 @@ -import { mkdtemp } from 'node:fs/promises'; +import { createHash } from 'node:crypto'; +import { mkdtemp, rm, stat, writeFile } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import { resolve } from 'node:path'; import { PassThrough } from 'node:stream'; @@ -9,6 +10,7 @@ import { runAiSdkSidecar } from './sidecar.js'; function fakeHarness() { const submitUserMessage = vi.fn(async () => undefined); + const turnResolvers: Array<() => void> = []; const session: HarnessV1Session = { sessionId: 'sidecar-session', isResume: false, @@ -19,7 +21,7 @@ function fakeHarness() { emit({ type: 'tool-call', toolCallId: 'tool-1', toolName: 'read', input: '{}' }); emit({ type: 'tool-approval-request', approvalId: 'approval-1', toolCallId: 'tool-1' }); return { - done: new Promise(() => undefined), + done: new Promise((resolveTurn) => turnResolvers.push(resolveTurn)), submitToolResult: vi.fn(async () => undefined), submitUserMessage, submitToolApproval: vi.fn(async () => undefined), @@ -38,11 +40,26 @@ function fakeHarness() { builtinTools: {}, doStart: vi.fn(async () => session), }; - return { harness, session, submitUserMessage }; + return { + harness, + session, + submitUserMessage, + settle() { + turnResolvers.shift()?.(); + }, + }; +} + +function hasDeliveryFrame(output: Array>, type: string, deliveryId: string): boolean { + return output.some( + (frame) => + frame.type === type && + (frame.payload as Record | undefined)?.delivery_id === deliveryId + ); } async function waitFor(predicate: () => boolean) { - for (let attempts = 0; attempts < 100; attempts += 1) { + for (let attempts = 0; attempts < 400; attempts += 1) { if (predicate()) return; await new Promise((resolveWait) => setTimeout(resolveWait, 5)); } @@ -55,6 +72,361 @@ afterEach(() => { }); describe('AI SDK native harness sidecar', () => { + it('requires a stable session id for durable deferred delivery', async () => { + const root = await mkdtemp(resolve(tmpdir(), 'relay-sidecar-')); + await expect( + runAiSdkSidecar( + { + name: 'Worker', + harness: 'fake', + workspace: resolve(root, 'workspace'), + runtimeRoot: resolve(root, 'runtime'), + }, + { input: new PassThrough(), write: () => undefined } + ) + ).rejects.toThrow('A stable sessionId is required for deferred delivery persistence'); + }); + + it('requires a stable runtime root for durable deferred delivery', async () => { + const root = await mkdtemp(resolve(tmpdir(), 'relay-sidecar-')); + await expect( + runAiSdkSidecar( + { + name: 'Worker', + harness: 'fake', + workspace: resolve(root, 'workspace'), + sessionId: 'sidecar-session', + }, + { input: new PassThrough(), write: () => undefined } + ) + ).rejects.toThrow('A stable runtimeRoot is required for deferred delivery persistence'); + }); + + it('keeps wait deliveries pending until the deferred message is accepted', async () => { + vi.stubEnv('RELAY_AGENT_TOKEN', 'at_live_native'); + vi.stubEnv('RELAY_WORKSPACE_KEY', 'rk_live_native'); + const fixture = fakeHarness(); + vi.spyOn(aiSdkAdapterRegistry, 'require').mockReturnValue({ + ...aiSdkAdapterRegistry.require('codex'), + createHarness: async () => fixture.harness, + }); + const root = await mkdtemp(resolve(tmpdir(), 'relay-sidecar-deferred-')); + const input = new PassThrough(); + const output: Array> = []; + const running = runAiSdkSidecar( + { + name: 'Worker', + harness: 'fake', + workspace: resolve(root, 'workspace'), + runtimeRoot: resolve(root, 'runtime'), + sessionId: 'sidecar-session', + }, + { input, write: (line) => output.push(JSON.parse(line)) } + ); + await waitFor(() => output.some((frame) => frame.type === 'agent_event')); + input.write(`${JSON.stringify({ v: 2, type: 'init_worker', payload: { agent: {} } })}\n`); + input.write( + `${JSON.stringify({ + v: 2, + type: 'deliver_relay', + payload: { + delivery_id: 'active', + event_id: 'event-active', + from: 'Human', + target: 'Worker', + body: 'active', + }, + })}\n` + ); + await waitFor(() => hasDeliveryFrame(output, 'delivery_ack', 'active')); + input.write( + `${JSON.stringify({ + v: 2, + type: 'deliver_relay', + payload: { + delivery_id: 'deferred', + event_id: 'event-deferred', + from: 'Human', + target: 'Worker', + body: 'later', + injection_mode: 'wait', + }, + })}\n` + ); + await waitFor(() => hasDeliveryFrame(output, 'delivery_queued', 'deferred')); + expect(hasDeliveryFrame(output, 'delivery_ack', 'deferred')).toBe(false); + + fixture.settle(); + await waitFor(() => hasDeliveryFrame(output, 'delivery_ack', 'deferred')); + + input.write( + `${JSON.stringify({ + v: 2, + type: 'deliver_relay', + payload: { + delivery_id: 'deferred-failed', + event_id: 'event-deferred-failed', + from: 'Human', + target: 'Worker', + body: 'later failure', + injection_mode: 'wait', + }, + })}\n` + ); + await waitFor(() => hasDeliveryFrame(output, 'delivery_queued', 'deferred-failed')); + vi.mocked(fixture.session.doPromptTurn).mockRejectedValueOnce(new Error('acceptance failed')); + fixture.settle(); + await waitFor(() => hasDeliveryFrame(output, 'delivery_failed', 'deferred-failed')); + expect(hasDeliveryFrame(output, 'delivery_ack', 'deferred-failed')).toBe(false); + input.end(); + await running; + }, 15_000); + + it('reports the final outcome of a deferred delivery restored after restart', async () => { + vi.stubEnv('RELAY_AGENT_TOKEN', 'at_live_native'); + vi.stubEnv('RELAY_WORKSPACE_KEY', 'rk_live_native'); + const original = fakeHarness(); + const restarted = fakeHarness(); + const createHarness = vi + .fn() + .mockResolvedValueOnce(original.harness) + .mockResolvedValueOnce(restarted.harness); + vi.spyOn(aiSdkAdapterRegistry, 'require').mockReturnValue({ + ...aiSdkAdapterRegistry.require('codex'), + createHarness, + }); + const root = await mkdtemp(resolve(tmpdir(), 'relay-sidecar-restored-outcome-')); + const config = { + name: 'Worker', + harness: 'fake', + workspace: resolve(root, 'workspace'), + runtimeRoot: resolve(root, 'runtime'), + sessionId: 'sidecar-session', + }; + + const firstInput = new PassThrough(); + const firstOutput: Array> = []; + const firstRun = runAiSdkSidecar(config, { + input: firstInput, + write: (line) => firstOutput.push(JSON.parse(line)), + }); + await waitFor(() => firstOutput.some((frame) => frame.type === 'agent_event')); + firstInput.write( + `${JSON.stringify({ + v: 2, + type: 'deliver_relay', + payload: { + delivery_id: 'active', + event_id: 'event-active', + from: 'Human', + target: 'Worker', + body: 'active', + }, + })}\n` + ); + await waitFor(() => hasDeliveryFrame(firstOutput, 'delivery_ack', 'active')); + firstInput.write( + `${JSON.stringify({ + v: 2, + type: 'deliver_relay', + payload: { + delivery_id: 'restored', + event_id: 'event-restored', + from: 'Human', + target: 'Worker', + body: 'after restart', + injection_mode: 'wait', + }, + })}\n` + ); + await waitFor(() => hasDeliveryFrame(firstOutput, 'delivery_queued', 'restored')); + firstInput.write( + `${JSON.stringify({ + v: 2, + type: 'deliver_relay', + payload: { + delivery_id: 'restored-failed', + event_id: 'event-restored-failed', + from: 'Human', + target: 'Worker', + body: 'fail after restart', + injection_mode: 'wait', + }, + })}\n` + ); + await waitFor(() => hasDeliveryFrame(firstOutput, 'delivery_queued', 'restored-failed')); + firstInput.end(); + await firstRun; + + const restartedInput = new PassThrough(); + const restartedOutput: Array> = []; + const restartedRun = runAiSdkSidecar(config, { + input: restartedInput, + write: (line) => restartedOutput.push(JSON.parse(line)), + }); + await waitFor(() => hasDeliveryFrame(restartedOutput, 'delivery_ack', 'restored')); + expect(restartedOutput).toContainEqual({ + v: 2, + type: 'delivery_ack', + payload: { + delivery_id: 'restored', + event_id: 'event-restored', + state: 'queued', + }, + }); + expect(restarted.session.doPromptTurn).toHaveBeenCalledTimes(1); + vi.mocked(restarted.session.doPromptTurn).mockRejectedValueOnce(new Error('restored acceptance failed')); + restarted.settle(); + await waitFor(() => hasDeliveryFrame(restartedOutput, 'delivery_failed', 'restored-failed')); + expect(restartedOutput).toContainEqual({ + v: 2, + type: 'delivery_failed', + payload: { + delivery_id: 'restored-failed', + event_id: 'event-restored-failed', + reason: 'restored acceptance failed', + }, + }); + restartedInput.end(); + await restartedRun; + }, 15_000); + + it('retires durable delivery state on broker worker shutdown', async () => { + vi.stubEnv('RELAY_AGENT_TOKEN', 'at_live_native'); + vi.stubEnv('RELAY_WORKSPACE_KEY', 'rk_live_native'); + const fixture = fakeHarness(); + vi.spyOn(aiSdkAdapterRegistry, 'require').mockReturnValue({ + ...aiSdkAdapterRegistry.require('codex'), + createHarness: async () => fixture.harness, + }); + const root = await mkdtemp(resolve(tmpdir(), 'relay-sidecar-shutdown-')); + const runtimeRoot = resolve(root, 'runtime'); + const sessionId = 'sidecar-session'; + const sessionStore = resolve( + runtimeRoot, + 'deferred-relay', + createHash('sha256').update(sessionId).digest('hex') + ); + const input = new PassThrough(); + const output: Array> = []; + const running = runAiSdkSidecar( + { + name: 'Worker', + harness: 'fake', + workspace: resolve(root, 'workspace'), + runtimeRoot, + sessionId, + }, + { input, write: (line) => output.push(JSON.parse(line)) } + ); + await waitFor(() => output.some((frame) => frame.type === 'agent_event')); + input.write( + `${JSON.stringify({ + v: 2, + type: 'deliver_relay', + payload: { + delivery_id: 'active', + event_id: 'event-active', + from: 'Human', + target: 'Worker', + body: 'active', + }, + })}\n` + ); + await waitFor(() => hasDeliveryFrame(output, 'delivery_ack', 'active')); + input.write( + `${JSON.stringify({ + v: 2, + type: 'deliver_relay', + payload: { + delivery_id: 'deferred', + event_id: 'event-deferred', + from: 'Human', + target: 'Worker', + body: 'later', + injection_mode: 'wait', + }, + })}\n` + ); + await waitFor(() => hasDeliveryFrame(output, 'delivery_queued', 'deferred')); + await expect(stat(sessionStore)).resolves.toBeDefined(); + + input.write(`${JSON.stringify({ v: 2, type: 'shutdown_worker', payload: {} })}\n`); + await running; + + expect(output.some((frame) => frame.type === 'worker_exited')).toBe(true); + expect(fixture.session.doDestroy).toHaveBeenCalledTimes(1); + await expect(stat(sessionStore)).rejects.toMatchObject({ code: 'ENOENT' }); + }); + + it('reports worker exit when broker shutdown cleanup fails', async () => { + vi.stubEnv('RELAY_AGENT_TOKEN', 'at_live_native'); + vi.stubEnv('RELAY_WORKSPACE_KEY', 'rk_live_native'); + const fixture = fakeHarness(); + vi.spyOn(aiSdkAdapterRegistry, 'require').mockReturnValue({ + ...aiSdkAdapterRegistry.require('codex'), + createHarness: async () => fixture.harness, + }); + const root = await mkdtemp(resolve(tmpdir(), 'relay-sidecar-shutdown-failure-')); + const runtimeRoot = resolve(root, 'runtime'); + const sessionId = 'sidecar-session'; + const sessionStore = resolve( + runtimeRoot, + 'deferred-relay', + createHash('sha256').update(sessionId).digest('hex') + ); + const input = new PassThrough(); + const output: Array> = []; + const running = runAiSdkSidecar( + { + name: 'Worker', + harness: 'fake', + workspace: resolve(root, 'workspace'), + runtimeRoot, + sessionId, + }, + { input, write: (line) => output.push(JSON.parse(line)) } + ); + await waitFor(() => output.some((frame) => frame.type === 'agent_event')); + input.write( + `${JSON.stringify({ + v: 2, + type: 'deliver_relay', + payload: { + delivery_id: 'active', + event_id: 'event-active', + from: 'Human', + target: 'Worker', + body: 'active', + }, + })}\n` + ); + await waitFor(() => hasDeliveryFrame(output, 'delivery_ack', 'active')); + input.write( + `${JSON.stringify({ + v: 2, + type: 'deliver_relay', + payload: { + delivery_id: 'deferred', + event_id: 'event-deferred', + from: 'Human', + target: 'Worker', + body: 'later', + injection_mode: 'wait', + }, + })}\n` + ); + await waitFor(() => hasDeliveryFrame(output, 'delivery_queued', 'deferred')); + await rm(sessionStore, { recursive: true }); + await writeFile(sessionStore, 'blocked'); + + input.write(`${JSON.stringify({ v: 2, type: 'shutdown_worker', payload: {} })}\n`); + await expect(running).rejects.toBeDefined(); + + expect(output).toContainEqual({ v: 2, type: 'worker_exited', payload: { code: 1 } }); + expect(fixture.session.doDestroy).toHaveBeenCalledTimes(1); + }); + it('speaks worker and native harness protocols with command deduplication', async () => { vi.stubEnv('RELAY_AGENT_TOKEN', 'at_live_native'); vi.stubEnv('RELAY_WORKSPACE_KEY', 'rk_live_native'); diff --git a/packages/harnesses/src/ai-sdk/sidecar.ts b/packages/harnesses/src/ai-sdk/sidecar.ts index 605b00ee39..5dd77af728 100644 --- a/packages/harnesses/src/ai-sdk/sidecar.ts +++ b/packages/harnesses/src/ai-sdk/sidecar.ts @@ -1,5 +1,6 @@ import { createHash } from 'node:crypto'; import { readFile } from 'node:fs/promises'; +import { resolve } from 'node:path'; import { createInterface } from 'node:readline'; import { NATIVE_HARNESS_PROTOCOL_VERSION, @@ -107,6 +108,12 @@ function diagnosticPayload(diagnostic: HarnessV1Diagnostic) { } export async function runAiSdkSidecar(config: AiSdkSidecarConfig, io: AiSdkSidecarIo): Promise { + if (!config.sessionId) { + throw new Error('A stable sessionId is required for deferred delivery persistence'); + } + if (!config.runtimeRoot) { + throw new Error('A stable runtimeRoot is required for deferred delivery persistence'); + } const entry = aiSdkAdapterRegistry.require(config.harness); const harness = await entry.createHarness(config.settings); const provider = new LocalHostSandboxProvider({ @@ -160,22 +167,46 @@ export async function runAiSdkSidecar(config: AiSdkSidecarConfig, io: AiSdkSidec handle: config.name.toLowerCase().replaceAll(/[^a-z0-9_-]/g, '-'), } satisfies AgentIdentity, host, + deferredQueuePath: resolve( + config.runtimeRoot, + 'deferred-relay', + createHash('sha256').update(host.sessionId).digest('hex') + ), }); const relayDeliveries = new Map(); relaySession.onEvent?.(async (event) => { await publishEvent(event); if (event.type === 'delivery.accepted' && event.deliveryId) { const delivery = relayDeliveries.get(event.deliveryId); - if (!delivery) return; relayDeliveries.delete(event.deliveryId); await write({ v: 2, type: 'delivery_ack', - payload: { delivery_id: event.deliveryId, event_id: delivery.eventId }, + payload: { + delivery_id: event.deliveryId, + // Restored deferred entries complete before stdin can rebuild the + // volatile delivery map. The durable entry retains the originating + // message id, so it remains sufficient to report the final outcome. + event_id: delivery?.eventId ?? event.messageId, + state: 'queued', + }, + }); + } else if (event.type === 'delivery.failed' && event.deliveryId) { + const delivery = relayDeliveries.get(event.deliveryId); + relayDeliveries.delete(event.deliveryId); + await write({ + v: 2, + type: 'delivery_failed', + payload: { + delivery_id: event.deliveryId, + event_id: delivery?.eventId ?? event.messageId, + reason: event.reason, + }, }); } }); await host.start(); + await relaySession.restoreDeferredMessages(); const maxCommandDedupeEntries = Math.max(1, config.maxCommandDedupeEntries ?? 10_000); const acknowledgements = new Map(); @@ -238,14 +269,27 @@ export async function runAiSdkSidecar(config: AiSdkSidecarConfig, io: AiSdkSidec } ); if (receipt.status === 'failed') throw new Error(receipt.reason); - if (receipt.status === 'accepted') { + if (receipt.status === 'deferred') { + const pending = relayDeliveries.get(delivery.delivery_id); + if (pending) { + await write({ + v: 2, + type: 'delivery_queued', + payload: { delivery_id: delivery.delivery_id, event_id: pending.eventId }, + }); + } + } else if (receipt.status === 'accepted') { const pending = relayDeliveries.get(delivery.delivery_id); if (pending) { relayDeliveries.delete(delivery.delivery_id); await write({ v: 2, type: 'delivery_ack', - payload: { delivery_id: delivery.delivery_id, event_id: pending.eventId }, + payload: { + delivery_id: delivery.delivery_id, + event_id: pending.eventId, + state: 'queued', + }, }); } } @@ -264,8 +308,13 @@ export async function runAiSdkSidecar(config: AiSdkSidecarConfig, io: AiSdkSidec continue; } if (frame.type === 'shutdown_worker') { - await host.destroy(); - await write({ v: 2, type: 'worker_exited', payload: { code: 0 } }); + try { + await relaySession.release('shutdown_worker'); + await write({ v: 2, type: 'worker_exited', payload: { code: 0 } }); + } catch (error) { + await write({ v: 2, type: 'worker_exited', payload: { code: 1 } }); + throw error; + } break; } if (!isCommandFrame(frame)) continue; diff --git a/packages/harnesses/src/define.ts b/packages/harnesses/src/define.ts index 96922da5b8..4a68b5a23e 100644 --- a/packages/harnesses/src/define.ts +++ b/packages/harnesses/src/define.ts @@ -1,4 +1,5 @@ import { randomUUID } from 'node:crypto'; +import { homedir } from 'node:os'; import path from 'node:path'; import { fileURLToPath } from 'node:url'; @@ -321,12 +322,17 @@ export function createNativeHarnessLaunch( const name = nextHarnessName(adapter.name, input.name); const cwd = path.resolve(input.cwd ?? process.cwd()); const sessionId = `native-${randomUUID()}`; + const runtimeRoot = path.resolve( + process.env.AGENT_RELAY_DATA_DIR?.trim() || path.join(homedir(), '.agentworkforce', 'relay'), + 'harness' + ); const sidecarEntry = fileURLToPath(new URL('./ai-sdk/sidecar-main.js', import.meta.url)); const sidecarConfig = { name, harness: adapter.name, workspace: cwd, sessionId, + runtimeRoot, settings: { ...(input.model ? { model: input.model } : {}) }, }; return { diff --git a/tests/e2e/fleet/fleet-e2e.test.ts b/tests/e2e/fleet/fleet-e2e.test.ts index 02346fab17..acc5e449f1 100644 --- a/tests/e2e/fleet/fleet-e2e.test.ts +++ b/tests/e2e/fleet/fleet-e2e.test.ts @@ -257,6 +257,8 @@ describe.skipIf(!pre.ok)('two-node fleet scenario matrix', () => { 'echo', 'relay:delivery-cursor-v1', 'relay:live-agents:v1', + 'relay:native-existing-session-reconcile:v1', + 'relay:native-existing-session:v1', 'release', 'spawn:claude', 'spawn:pool', @@ -266,6 +268,8 @@ describe.skipIf(!pre.ok)('two-node fleet scenario matrix', () => { 'ping', 'relay:delivery-cursor-v1', 'relay:live-agents:v1', + 'relay:native-existing-session-reconcile:v1', + 'relay:native-existing-session:v1', 'release', 'spawn:codex', 'spawn:pool', diff --git a/tests/relayflows/cases/1851-native-existing-session-delivery/case.json b/tests/relayflows/cases/1851-native-existing-session-delivery/case.json new file mode 100644 index 0000000000..fc9ddd51ea --- /dev/null +++ b/tests/relayflows/cases/1851-native-existing-session-delivery/case.json @@ -0,0 +1,21 @@ +{ + "version": 1, + "id": "1851-native-existing-session-delivery", + "kind": "feature", + "title": "Persistent brokers advertise durable existing-session delivery and reconciliation", + "requirements": ["broker-linux-x64"], + "runner": { + "command": ["node", "tests/relayflows/cases/1851-native-existing-session-delivery/run.mjs"] + }, + "timeoutSeconds": 120, + "expected": { + "base": { + "outcome": "absent", + "signature": "native_existing_session_capabilities_absent" + }, + "head": { + "outcome": "fixed", + "signature": "native_existing_session_capabilities_advertised" + } + } +} diff --git a/tests/relayflows/cases/1851-native-existing-session-delivery/run.mjs b/tests/relayflows/cases/1851-native-existing-session-delivery/run.mjs new file mode 100644 index 0000000000..4707d6c72a --- /dev/null +++ b/tests/relayflows/cases/1851-native-existing-session-delivery/run.mjs @@ -0,0 +1,144 @@ +import assert from 'node:assert/strict'; +import { execFileSync, spawn } from 'node:child_process'; +import { access, mkdir, mkdtemp, rm, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import path from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { engineFixture } from '../1766-durable-task-receipt/engine-fixture.mjs'; + +const caseId = '1851-native-existing-session-delivery'; +const deliverCapability = 'relay:native-existing-session:v1'; +const reconcileCapability = 'relay:native-existing-session-reconcile:v1'; +const required = (key) => { + assert(process.env[key], `Missing ${key}`); + return process.env[key]; +}; +const arm = required('RELAY_PR_PROOF_ARM'); +assert(['base', 'head'].includes(arm)); +const binary = path.resolve(required('RELAY_PR_PROOF_BROKER_BINARY')); +const target = required('RELAY_PR_PROOF_TARGET_DIR'); +const harness = required('RELAY_PR_PROOF_HARNESS_DIR'); +const resultPath = required('RELAY_PR_PROOF_RESULT_PATH'); +assert.equal( + execFileSync('git', ['-C', target, 'rev-parse', 'HEAD'], { encoding: 'utf8' }).trim(), + required(arm === 'base' ? 'RELAY_PR_PROOF_BASE_SHA' : 'RELAY_PR_PROOF_HEAD_SHA') +); +const relative = path.relative(path.resolve(harness), fileURLToPath(import.meta.url)); +assert(relative && !relative.startsWith('..') && !path.isAbsolute(relative)); +await access(binary, 1); + +const directory = await mkdtemp(path.join(tmpdir(), 'relayflow-native-delivery-')); +const state = path.join(directory, 'state'); +await mkdir(state); +const engine = await engineFixture(); +let broker; +let logs = ''; + +async function waitForRegistration(milliseconds = 15_000) { + const deadline = performance.now() + milliseconds; + while (performance.now() < deadline) { + engine.check(); + const registration = engine.frames.find((frame) => frame.type === 'node.register'); + if (registration) return registration; + if (broker?.exitCode !== null && broker?.exitCode !== undefined) { + throw new Error(`Broker exited before registration: ${logs}`); + } + await new Promise((resolve) => setTimeout(resolve, 20)); + } + throw new Error(`Timed out waiting for node.register: ${logs}`); +} + +async function stop() { + if (!broker || broker.exitCode !== null || broker.signalCode !== null) return; + const child = broker; + const exited = new Promise((resolve) => child.once('exit', resolve)); + child.kill('SIGTERM'); + const timer = setTimeout(() => child.kill('SIGKILL'), 2_000); + await exited; + clearTimeout(timer); + broker = undefined; +} + +try { + broker = spawn( + binary, + [ + 'init', + '--persist', + '--instance-name', + 'native-delivery-proof-node', + '--workspace-key', + 'rk_fixture_native_delivery_proof', + '--state-dir', + state, + '--api-port', + '0', + '--channels', + '', + ], + { + cwd: directory, + env: { + PATH: process.env.PATH, + TMPDIR: directory, + AGENT_RELAY_BROKER_LOG: 'stderr', + RELAYCAST_BASE_URL: engine.baseUrl, + RELAY_BASE_URL: engine.baseUrl, + RELAY_BROKER_API_KEY: 'br_fixture_native_delivery_proof', + RELAY_NODE_ID: 'node-fixture-native-delivery-proof', + // The shared loopback engine fixture deliberately accepts one fixed + // credential so a case cannot weaken its authentication behavior. + RELAY_NODE_TOKEN: 'nt_fixture_task_proof', + AGENT_RELAY_TELEMETRY_DISABLED: '1', + AGENT_RELAY_NO_DEBUG_FILES: '1', + }, + stdio: ['ignore', 'pipe', 'pipe'], + } + ); + broker.stdout.on('data', (chunk) => { + logs = (logs + chunk).slice(-10_000); + }); + broker.stderr.on('data', (chunk) => { + logs = (logs + chunk).slice(-10_000); + }); + + const registration = await waitForRegistration(); + const capabilities = new Map(registration.capabilities.map((value) => [value.name, value])); + if (arm === 'base') { + assert.equal(capabilities.has(deliverCapability), false); + assert.equal(capabilities.has(reconcileCapability), false); + } else { + const delivery = capabilities.get(deliverCapability); + const reconciliation = capabilities.get(reconcileCapability); + assert.equal(delivery?.kind, 'action'); + assert.equal(delivery?.metadata?.contract, 'deliverNativeExistingSession'); + assert.equal(delivery?.metadata?.durableReceipts, true); + assert.equal(delivery?.metadata?.idempotencyField, 'deliveryId'); + assert.equal(delivery?.metadata?.reconcileAction, reconcileCapability); + assert.equal(reconciliation?.kind, 'action'); + assert.equal(reconciliation?.metadata?.contract, 'reconcileNativeExistingSession'); + } + engine.check(); + await mkdir(path.dirname(resultPath), { recursive: true }); + await writeFile( + resultPath, + `${JSON.stringify({ + version: 1, + caseId, + arm, + outcome: arm === 'base' ? 'absent' : 'fixed', + signature: + arm === 'base' + ? 'native_existing_session_capabilities_absent' + : 'native_existing_session_capabilities_advertised', + details: + arm === 'base' + ? 'The exact base broker registered without either native existing-session action.' + : 'The exact head broker registered both versioned actions and advertised the durable receipt and reconciliation contract metadata.', + })}\n` + ); +} finally { + await stop(); + await engine.close(); + await rm(directory, { recursive: true, force: true }); +}