diff --git a/crates/broker/src/protocol.rs b/crates/broker/src/protocol.rs index 31b9c8dc0..b65929c0f 100644 --- a/crates/broker/src/protocol.rs +++ b/crates/broker/src/protocol.rs @@ -337,6 +337,15 @@ pub struct ProtocolError { pub data: Option, } +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum DeliveryReadAckStatus { + Marked, + Failed, + SkippedSynthetic, + SuppressedDuplicate, +} + #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] #[serde(tag = "kind", rename_all = "snake_case")] pub enum BrokerEvent { @@ -412,6 +421,14 @@ pub enum BrokerEvent { from: String, to: MessageTarget, }, + DeliveryReadAck { + name: WorkerName, + delivery_id: DeliveryId, + event_id: EventId, + status: DeliveryReadAckStatus, + #[serde(default, skip_serializing_if = "Option::is_none")] + reason: Option, + }, MessageDeliveryFailed { name: WorkerName, #[serde(default)] diff --git a/crates/broker/src/relaycast/ws.rs b/crates/broker/src/relaycast/ws.rs index c2959d971..884dd0031 100644 --- a/crates/broker/src/relaycast/ws.rs +++ b/crates/broker/src/relaycast/ws.rs @@ -422,6 +422,37 @@ impl RelaycastHttpClient { .map_err(|error| anyhow::anyhow!("{error}")) } + async fn registered_agent_client_as( + &self, + agent_name: &str, + cli_hint: Option<&str>, + ) -> Result { + let registration = self + .registration + .as_ref() + .as_ref() + .context("SDK relay client not initialized")?; + registration + .registered_agent_client(agent_name, cli_hint.or(Some(self.default_cli.as_str()))) + .await + .map_err(|error| anyhow::anyhow!("{error}")) + } + + /// Impersonation by design: delivery read-acks must be attributed to the + /// recipient worker's agent identity, not the broker identity. + pub async fn mark_read_as_agent( + &self, + agent_name: &str, + cli_hint: Option<&str>, + message_id: &str, + ) -> Result { + self.registered_agent_client_as(agent_name, cli_hint) + .await? + .mark_read(message_id) + .await + .map_err(|error| anyhow::anyhow!("relaycast mark_read failed: {error}")) + } + /// Register an action whose handler is this broker's agent. Spawn/release /// are exposed as relaycast actions so other agents can invoke them as /// structured agent-to-agent RPC. @@ -843,6 +874,54 @@ mod tests { assert!(message.contains("pre-register")); } + #[tokio::test] + async fn mark_read_as_agent_uses_seeded_recipient_token_without_respawn() { + let server = MockServer::start(); + let read_mock = server.mock(|when, then| { + when.method(POST) + .path("/v1/messages/msg_1/read") + .header("authorization", "Bearer at_live_existing_recipient"); + then.status(200).json_body(json!({ + "ok": true, + "data": { + "message_id": "msg_1", + "agent_id": "agent_existing_recipient", + "read_at": "2026-06-08T10:00:00.000Z" + } + })); + }); + let spawn_mock = server.mock(|when, then| { + when.method(POST).path("/v1/agents/spawn"); + then.status(200).json_body(json!({ + "ok": true, + "data": { + "agent": { + "id": "agent_fresh_wrong", + "name": "recipient", + "type": "agent", + "status": "online", + "created_at": "2026-06-08T10:00:00.000Z", + "last_seen": "2026-06-08T10:00:00.000Z", + "metadata": {} + }, + "token": "at_live_fresh_wrong" + } + })); + }); + + let client = RelaycastHttpClient::new(server.base_url(), "rk_live_test", "broker", "codex"); + client.seed_agent_token("recipient", "at_live_existing_recipient"); + + let result = client + .mark_read_as_agent("recipient", Some("codex"), "msg_1") + .await + .expect("seeded recipient should mark read"); + + assert_eq!(result["agent_id"], "agent_existing_recipient"); + read_mock.assert_hits(1); + spawn_mock.assert_hits(0); + } + #[tokio::test] #[ignore = "relaycast API response fixture mismatch - needs investigation"] async fn send_with_mode_forwards_steer_for_relaycast_dm_targets() { diff --git a/crates/broker/src/runtime/api.rs b/crates/broker/src/runtime/api.rs index f5e2884f9..788cd970a 100644 --- a/crates/broker/src/runtime/api.rs +++ b/crates/broker/src/runtime/api.rs @@ -96,8 +96,15 @@ impl BrokerRuntime { } }; - // Caller-supplied agent_token overrides auto-registration - let worker_relay_key = agent_token.or(worker_relay_key); + // Caller-supplied agent_token overrides auto-registration. + // Seed it so broker-side read-acks later act as this exact + // recipient identity instead of minting a replacement token. + let worker_relay_key = if let Some(token) = agent_token { + seed_supplied_agent_token(relaycast_http, &name, &token); + Some(token) + } else { + worker_relay_key + }; let mut effective_task = normalize_initial_task(task); if let Some(ref continue_from) = continue_from { diff --git a/crates/broker/src/runtime/delivery.rs b/crates/broker/src/runtime/delivery.rs index 6ab763765..206601938 100644 --- a/crates/broker/src/runtime/delivery.rs +++ b/crates/broker/src/runtime/delivery.rs @@ -116,6 +116,197 @@ pub(crate) struct DeliveryAckPayload { pub(super) event_id: EventId, } +/// Classify delivery ids that are meaningful Relaycast message ids for +/// read-ack purposes. A read-ack means "delivered to the recipient location", +/// not proof that a model turn cognitively processed the message. +pub(crate) fn synthetic_delivery_read_ack_reason(event_id: &EventId) -> Option<&'static str> { + let event_id = event_id.as_str().trim(); + if event_id.is_empty() { + return Some("blank_event_id"); + } + if event_id.starts_with("http_") { + return Some("http_api_synthetic_event_id"); + } + if event_id.starts_with("init_") { + return Some("initial_task_synthetic_event_id"); + } + if event_id.starts_with("cont_load_") { + return Some("continuity_synthetic_event_id"); + } + if event_id.starts_with("flush_") { + return Some("manual_flush_synthetic_event_id"); + } + None +} + +#[cfg(test)] +pub(crate) fn delivery_read_ack_is_relaycast_message(event_id: &EventId) -> bool { + synthetic_delivery_read_ack_reason(event_id).is_none() +} + +pub(crate) fn seed_supplied_agent_token( + relaycast_http: &RelaycastHttpClient, + agent_name: &str, + token: &str, +) { + relaycast_http.seed_agent_token(agent_name, token); +} + +const DELIVERY_READ_ACK_TIMEOUT: Duration = Duration::from_secs(2); + +pub(crate) fn mark_delivery_read_ack( + relaycast_http: &RelaycastHttpClient, + sdk_out_tx: &mpsc::Sender>, + dedup: &mut DedupCache, + worker_name: &WorkerName, + cli_hint: Option<&str>, + delivery_id: &DeliveryId, + event_id: &EventId, +) { + mark_delivery_read_ack_with_timeout( + relaycast_http, + sdk_out_tx, + dedup, + worker_name, + cli_hint, + delivery_id, + event_id, + DELIVERY_READ_ACK_TIMEOUT, + ); +} + +#[allow(clippy::too_many_arguments)] +pub(crate) fn mark_delivery_read_ack_with_timeout( + relaycast_http: &RelaycastHttpClient, + sdk_out_tx: &mpsc::Sender>, + dedup: &mut DedupCache, + worker_name: &WorkerName, + cli_hint: Option<&str>, + delivery_id: &DeliveryId, + event_id: &EventId, + timeout_window: Duration, +) { + let dedup_key = format!("delivery_read_ack:{worker_name}:{event_id}"); + if !dedup.insert_if_new(&dedup_key, Instant::now()) { + emit_delivery_read_ack_telemetry( + sdk_out_tx.clone(), + BrokerEvent::DeliveryReadAck { + name: worker_name.clone(), + delivery_id: delivery_id.clone(), + event_id: event_id.clone(), + status: DeliveryReadAckStatus::SuppressedDuplicate, + reason: Some("duplicate_delivery_read_ack".to_string()), + }, + ); + return; + } + + if let Some(reason) = synthetic_delivery_read_ack_reason(event_id) { + emit_delivery_read_ack_telemetry( + sdk_out_tx.clone(), + BrokerEvent::DeliveryReadAck { + name: worker_name.clone(), + delivery_id: delivery_id.clone(), + event_id: event_id.clone(), + status: DeliveryReadAckStatus::SkippedSynthetic, + reason: Some(reason.to_string()), + }, + ); + return; + } + + let relaycast_http = relaycast_http.clone(); + let sdk_out_tx = sdk_out_tx.clone(); + let worker_name = worker_name.clone(); + let cli_hint = cli_hint.map(str::to_string); + let delivery_id = delivery_id.clone(); + let event_id = event_id.clone(); + + tokio::spawn(async move { + let result = timeout( + timeout_window, + relaycast_http.mark_read_as_agent( + worker_name.as_str(), + cli_hint.as_deref(), + event_id.as_str(), + ), + ) + .await; + + match result { + Ok(Ok(_)) => { + let _ = send_broker_event( + &sdk_out_tx, + BrokerEvent::DeliveryReadAck { + name: worker_name, + delivery_id, + event_id, + status: DeliveryReadAckStatus::Marked, + reason: None, + }, + ) + .await; + } + Ok(Err(error)) => { + let reason = error.to_string(); + tracing::warn!( + target = "agent_relay::broker", + worker = %worker_name, + delivery_id = %delivery_id, + event_id = %event_id, + error = %reason, + "failed to mark relaycast message read after delivery_ack" + ); + let _ = send_broker_event( + &sdk_out_tx, + BrokerEvent::DeliveryReadAck { + name: worker_name, + delivery_id, + event_id, + status: DeliveryReadAckStatus::Failed, + reason: Some(reason), + }, + ) + .await; + } + Err(_) => { + let reason = format!( + "relaycast mark_read timed out after {}ms", + timeout_window.as_millis() + ); + tracing::warn!( + target = "agent_relay::broker", + worker = %worker_name, + delivery_id = %delivery_id, + event_id = %event_id, + timeout_ms = %timeout_window.as_millis(), + "timed out marking relaycast message read after delivery_ack" + ); + let _ = send_broker_event( + &sdk_out_tx, + BrokerEvent::DeliveryReadAck { + name: worker_name, + delivery_id, + event_id, + status: DeliveryReadAckStatus::Failed, + reason: Some(reason), + }, + ) + .await; + } + } + }); +} + +fn emit_delivery_read_ack_telemetry( + sdk_out_tx: mpsc::Sender>, + event: BrokerEvent, +) { + tokio::spawn(async move { + let _ = send_broker_event(&sdk_out_tx, event).await; + }); +} + /// Outcome of [`queue_inbound_for_delivery_mode`]. Distinguishes the /// three cases broker call sites care about: the message is queued and /// should wait for an explicit flush, the queue should be drained now, diff --git a/crates/broker/src/runtime/mod.rs b/crates/broker/src/runtime/mod.rs index 258b4f060..7153731d0 100644 --- a/crates/broker/src/runtime/mod.rs +++ b/crates/broker/src/runtime/mod.rs @@ -31,9 +31,9 @@ use crate::{ WorkspaceAlias, WorkspaceId, }, protocol::{ - AgentRuntime, AgentSpec, BrokerEvent, HeadlessProvider as ProtocolHeadlessProvider, - MessageInjectionMode, ProtocolEnvelope, RelayDelivery, ResolvedHarnessConfig, - PROTOCOL_VERSION, + AgentRuntime, AgentSpec, BrokerEvent, DeliveryReadAckStatus, + HeadlessProvider as ProtocolHeadlessProvider, MessageInjectionMode, ProtocolEnvelope, + RelayDelivery, ResolvedHarnessConfig, PROTOCOL_VERSION, }, relaycast::{ agent_name_eq, format_worker_preregistration_error, is_self_name, map_ws_event, diff --git a/crates/broker/src/runtime/relaycast_events.rs b/crates/broker/src/runtime/relaycast_events.rs index b7c30c5b2..2f89e08ca 100644 --- a/crates/broker/src/runtime/relaycast_events.rs +++ b/crates/broker/src/runtime/relaycast_events.rs @@ -268,9 +268,9 @@ impl BrokerRuntime { // short timeout keeps spawn latency bounded while still // giving the registration call a real chance. let worker_relay_key = { - let ws_token = relaycast_ws_spawn_token(&ws_value); - if ws_token.is_some() { - ws_token + if let Some(token) = relaycast_ws_spawn_token(&ws_value) { + seed_supplied_agent_token(&workspace_http, &name, &token); + Some(token) } else { const REG_TIMEOUT: Duration = Duration::from_secs(3); match tokio::time::timeout( @@ -498,9 +498,9 @@ impl BrokerRuntime { // Pre-register (same logic as primary WS spawn path). let worker_relay_key = { - let ws_token = relaycast_ws_spawn_token(&ws_value); - if ws_token.is_some() { - ws_token + if let Some(token) = relaycast_ws_spawn_token(&ws_value) { + seed_supplied_agent_token(&workspace_http, &name, &token); + Some(token) } else { const REG_TIMEOUT: Duration = Duration::from_secs(3); match tokio::time::timeout( diff --git a/crates/broker/src/runtime/tests.rs b/crates/broker/src/runtime/tests.rs index e730e7cb4..678239154 100644 --- a/crates/broker/src/runtime/tests.rs +++ b/crates/broker/src/runtime/tests.rs @@ -10,8 +10,8 @@ use crate::ids::{ ChannelName, DeliveryId, EventId, MessageTarget, WorkerName, WorkspaceAlias, WorkspaceId, }; use crate::protocol::{ - AgentSpec, HarnessReleasePolicy, HeadlessHarnessConfig, HeadlessHarnessDriver, - MessageInjectionMode, RelayDelivery, ResolvedHarnessConfig, + AgentSpec, BrokerEvent, DeliveryReadAckStatus, HarnessReleasePolicy, HeadlessHarnessConfig, + HeadlessHarnessDriver, MessageInjectionMode, RelayDelivery, ResolvedHarnessConfig, }; use crate::worker::{AgentWorkState, WorkerEvent, WorkerHandle, WorkerRegistry}; use crate::{ @@ -30,20 +30,24 @@ use tokio::sync::mpsc; use super::{ build_agent_state_transition_event, build_http_api_spawn_spec, build_thread_infos, channels_from_csv, clear_pending_delivery_if_event_matches, continuity_dir, - delivery_retry_interval, derive_ws_base_url_from_http, display_target_for_dashboard, - drop_pending_for_worker, emit_delivery_attempt_outcome, emit_dropped_delivery_failures, - ensure_ephemeral_paths, extract_mcp_message_ids, http_api_event_emit_timeout, - http_api_local_delivery_timeout, http_api_relaycast_send_timeout, - is_relaycast_self_control_target, is_unknown_worker_error_message, normalize_channel, - normalize_initial_task, normalize_sender, parse_sort_key_from_raw_timestamp, - queue_inbound_for_delivery_mode, relaycast_spawn_control_dedup_key, - relaycast_ws_control_dedup_key, relaycast_ws_should_apply_local_spawn_echo_dedup, - relaycast_ws_spawn_token, retry_pending_delivery, sender_is_dashboard_label, - should_clear_pending_delivery_for_event, AgentRuntime, DeliveryAttemptOutcome, InboundContext, + delivery_read_ack_is_relaycast_message, delivery_retry_interval, derive_ws_base_url_from_http, + display_target_for_dashboard, drop_pending_for_worker, emit_delivery_attempt_outcome, + emit_dropped_delivery_failures, ensure_ephemeral_paths, extract_mcp_message_ids, + http_api_event_emit_timeout, http_api_local_delivery_timeout, http_api_relaycast_send_timeout, + is_relaycast_self_control_target, is_unknown_worker_error_message, mark_delivery_read_ack, + mark_delivery_read_ack_with_timeout, normalize_channel, normalize_initial_task, + normalize_sender, parse_sort_key_from_raw_timestamp, queue_inbound_for_delivery_mode, + relaycast_spawn_control_dedup_key, relaycast_ws_control_dedup_key, + relaycast_ws_should_apply_local_spawn_echo_dedup, relaycast_ws_spawn_token, + retry_pending_delivery, seed_supplied_agent_token, send_broker_event, + sender_is_dashboard_label, should_clear_pending_delivery_for_event, + synthetic_delivery_read_ack_reason, AgentRuntime, DeliveryAttemptOutcome, InboundContext, InboundQueueOutcome, PendingDelivery, ProtocolHeadlessProvider, MAX_DELIVERY_RETRIES, }; use crate::dedup::DedupCache; -use crate::relaycast::{format_worker_preregistration_error, RelaycastRegistrationError}; +use crate::relaycast::{ + format_worker_preregistration_error, RelaycastHttpClient, RelaycastRegistrationError, +}; use crate::types::{InboundDeliveryMode, InboundDeliveryState}; fn env_test_lock() -> &'static Mutex<()> { @@ -121,6 +125,28 @@ fn inbound_ctx<'a>(event_id: &'a str) -> InboundContext<'a> { } } +fn pending_delivery(worker_name: &str, delivery_id: &str, event_id: &str) -> PendingDelivery { + PendingDelivery { + worker_name: WorkerName::from(worker_name), + delivery: RelayDelivery { + delivery_id: DeliveryId::new(delivery_id), + event_id: EventId::new(event_id), + workspace_id: Some(WorkspaceId::new("ws_test")), + workspace_alias: Some(WorkspaceAlias::new("test")), + from: "sender".to_string(), + target: MessageTarget::new(worker_name), + body: "hello".to_string(), + thread_id: None, + priority: None, + injection_mode: MessageInjectionMode::Wait, + }, + attempts: 1, + next_retry_at: Instant::now(), + queued_at_ms: super::unix_timestamp_millis(), + last_error: None, + } +} + #[tokio::test] async fn inbound_queue_auto_inject_drains_immediately_with_full_context() { let worker_name = "worker-a"; @@ -1348,6 +1374,390 @@ fn clear_pending_delivery_returns_none_for_stale_event_id() { assert!(pending.contains_key("del_1")); } +#[test] +fn delivery_read_ack_classification_skips_synthetic_event_ids() { + let cases = [ + ("", Some("blank_event_id")), + (" ", Some("blank_event_id")), + ("http_123", Some("http_api_synthetic_event_id")), + ("init_123", Some("initial_task_synthetic_event_id")), + ("cont_load_123", Some("continuity_synthetic_event_id")), + ("flush_123", Some("manual_flush_synthetic_event_id")), + ("msg_123", None), + ("1780911342_317109", None), + ]; + + for (event_id, expected) in cases { + let event_id = EventId::new(event_id); + assert_eq!(synthetic_delivery_read_ack_reason(&event_id), expected); + assert_eq!( + delivery_read_ack_is_relaycast_message(&event_id), + expected.is_none() + ); + } +} + +#[test] +fn delivery_read_ack_event_shape_is_stable() { + let event = BrokerEvent::DeliveryReadAck { + name: WorkerName::new("Worker1"), + delivery_id: DeliveryId::new("del_1"), + event_id: EventId::new("msg_1"), + status: DeliveryReadAckStatus::SkippedSynthetic, + reason: Some("initial_task_synthetic_event_id".to_string()), + }; + + let encoded = serde_json::to_value(&event).expect("event serializes"); + assert_eq!(encoded["kind"], "delivery_read_ack"); + assert_eq!(encoded["name"], "Worker1"); + assert_eq!(encoded["delivery_id"], "del_1"); + assert_eq!(encoded["event_id"], "msg_1"); + assert_eq!(encoded["status"], "skipped_synthetic"); + assert_eq!(encoded["reason"], "initial_task_synthetic_event_id"); +} + +#[tokio::test] +async fn confirmed_delivery_read_ack_marks_relaycast_exactly_once() { + use httpmock::{Method::POST, MockServer}; + + let server = MockServer::start(); + let read_mock = server.mock(|when, then| { + when.method(POST) + .path("/v1/messages/msg_1/read") + .header("authorization", "Bearer at_live_supplied_recipient"); + then.status(200).json_body(json!({ + "ok": true, + "data": { + "message_id": "msg_1", + "agent_id": "agent_supplied_recipient", + "read_at": "2026-06-08T10:00:00.000Z" + } + })); + }); + let spawn_mock = server.mock(|when, then| { + when.method(POST).path("/v1/agents/spawn"); + then.status(200).json_body(json!({ + "ok": true, + "data": { + "agent": { + "id": "agent_fresh_wrong", + "name": "recipient", + "type": "agent", + "status": "online", + "created_at": "2026-06-08T10:00:00.000Z", + "last_seen": "2026-06-08T10:00:00.000Z", + "metadata": {} + }, + "token": "at_live_fresh_wrong" + } + })); + }); + let client = RelaycastHttpClient::new(server.base_url(), "rk_live_test", "broker", "codex"); + seed_supplied_agent_token(&client, "recipient", "at_live_supplied_recipient"); + let mut dedup = DedupCache::new(Duration::from_secs(300), 16); + let (tx, mut rx) = mpsc::channel(4); + let mut pending = HashMap::from([( + DeliveryId::new("del_1"), + pending_delivery("recipient", "del_1", "msg_1"), + )]); + + let confirmed = clear_pending_delivery_if_event_matches( + &mut pending, + "del_1", + Some("msg_1"), + "recipient", + "delivery_ack", + ) + .expect("matching delivery_ack confirms the pending delivery"); + + mark_delivery_read_ack( + &client, + &tx, + &mut dedup, + &WorkerName::new("recipient"), + Some("codex"), + &confirmed.delivery.delivery_id, + &confirmed.delivery.event_id, + ); + + let frame = tokio::time::timeout(Duration::from_secs(1), rx.recv()) + .await + .expect("delivery_read_ack telemetry should arrive") + .expect("delivery_read_ack event emitted"); + assert_eq!(frame.msg_type, "event"); + assert_eq!(frame.payload["kind"], "delivery_read_ack"); + assert_eq!(frame.payload["name"], "recipient"); + assert_eq!(frame.payload["delivery_id"], "del_1"); + assert_eq!(frame.payload["event_id"], "msg_1"); + assert_eq!(frame.payload["status"], "marked"); + assert!(frame.payload.get("reason").is_none()); + read_mock.assert_hits(1); + spawn_mock.assert_hits(0); +} + +#[tokio::test] +async fn duplicate_delivery_read_ack_suppresses_repeat_mark_read() { + use httpmock::{Method::POST, MockServer}; + + let server = MockServer::start(); + let read_mock = server.mock(|when, then| { + when.method(POST) + .path("/v1/messages/msg_dup/read") + .header("authorization", "Bearer at_live_recipient_dup"); + then.status(200).json_body(json!({ + "ok": true, + "data": { + "message_id": "msg_dup", + "agent_id": "agent_recipient_dup", + "read_at": "2026-06-08T10:00:00.000Z" + } + })); + }); + let spawn_mock = server.mock(|when, then| { + when.method(POST).path("/v1/agents/spawn"); + then.status(200).json_body(json!({ + "ok": true, + "data": { + "agent": { + "id": "agent_fresh_wrong", + "name": "recipient", + "type": "agent", + "status": "online", + "created_at": "2026-06-08T10:00:00.000Z", + "last_seen": "2026-06-08T10:00:00.000Z", + "metadata": {} + }, + "token": "at_live_fresh_wrong" + } + })); + }); + let client = RelaycastHttpClient::new(server.base_url(), "rk_live_test", "broker", "codex"); + seed_supplied_agent_token(&client, "recipient", "at_live_recipient_dup"); + let mut dedup = DedupCache::new(Duration::from_secs(300), 16); + let (tx, mut rx) = mpsc::channel(4); + + mark_delivery_read_ack( + &client, + &tx, + &mut dedup, + &WorkerName::new("recipient"), + Some("codex"), + &DeliveryId::new("del_dup_1"), + &EventId::new("msg_dup"), + ); + mark_delivery_read_ack( + &client, + &tx, + &mut dedup, + &WorkerName::new("recipient"), + Some("codex"), + &DeliveryId::new("del_dup_2"), + &EventId::new("msg_dup"), + ); + + let mut statuses = Vec::new(); + for _ in 0..2 { + let frame = tokio::time::timeout(Duration::from_secs(1), rx.recv()) + .await + .expect("delivery_read_ack telemetry should arrive") + .expect("delivery_read_ack event emitted"); + assert_eq!(frame.payload["kind"], "delivery_read_ack"); + statuses.push( + frame.payload["status"] + .as_str() + .unwrap_or_default() + .to_string(), + ); + } + + assert!(statuses.iter().any(|status| status == "marked")); + assert!(statuses + .iter() + .any(|status| status == "suppressed_duplicate")); + read_mock.assert_hits(1); + spawn_mock.assert_hits(0); +} + +#[tokio::test] +async fn stale_delivery_ack_event_id_does_not_mark_read() { + use httpmock::{Method::POST, MockServer}; + + let server = MockServer::start(); + let read_mock = server.mock(|when, then| { + when.method(POST).path("/v1/messages/msg_current/read"); + then.status(200).json_body(json!({"ok": true, "data": {}})); + }); + let client = RelaycastHttpClient::new(server.base_url(), "rk_live_test", "broker", "codex"); + seed_supplied_agent_token(&client, "recipient", "at_live_recipient"); + let mut dedup = DedupCache::new(Duration::from_secs(300), 16); + let (tx, mut rx) = mpsc::channel(4); + let mut pending = HashMap::from([( + DeliveryId::new("del_stale"), + pending_delivery("recipient", "del_stale", "msg_current"), + )]); + + let confirmed = clear_pending_delivery_if_event_matches( + &mut pending, + "del_stale", + Some("msg_stale"), + "recipient", + "delivery_ack", + ); + if let Some(confirmed) = confirmed { + mark_delivery_read_ack( + &client, + &tx, + &mut dedup, + &WorkerName::new("recipient"), + Some("codex"), + &confirmed.delivery.delivery_id, + &confirmed.delivery.event_id, + ); + } + + assert!(pending.contains_key("del_stale")); + read_mock.assert_hits(0); + assert!(tokio::time::timeout(Duration::from_millis(50), rx.recv()) + .await + .is_err()); +} + +#[tokio::test] +async fn synthetic_delivery_read_ack_skips_mark_read() { + use httpmock::{Method::POST, MockServer}; + + let server = MockServer::start(); + let read_mock = server.mock(|when, then| { + when.method(POST).path("/v1/messages/init_123/read"); + then.status(200).json_body(json!({"ok": true, "data": {}})); + }); + let client = RelaycastHttpClient::new(server.base_url(), "rk_live_test", "broker", "codex"); + let mut dedup = DedupCache::new(Duration::from_secs(300), 16); + let (tx, mut rx) = mpsc::channel(4); + + mark_delivery_read_ack( + &client, + &tx, + &mut dedup, + &WorkerName::new("recipient"), + Some("codex"), + &DeliveryId::new("del_init"), + &EventId::new("init_123"), + ); + mark_delivery_read_ack( + &client, + &tx, + &mut dedup, + &WorkerName::new("recipient"), + Some("codex"), + &DeliveryId::new("del_init_duplicate"), + &EventId::new("init_123"), + ); + + let first = tokio::time::timeout(Duration::from_secs(1), rx.recv()) + .await + .expect("synthetic skip telemetry should arrive") + .expect("delivery_read_ack event emitted"); + assert_eq!(first.payload["kind"], "delivery_read_ack"); + assert_eq!(first.payload["status"], "skipped_synthetic"); + assert_eq!(first.payload["reason"], "initial_task_synthetic_event_id"); + + let duplicate = tokio::time::timeout(Duration::from_secs(1), rx.recv()) + .await + .expect("duplicate synthetic telemetry should arrive") + .expect("delivery_read_ack event emitted"); + assert_eq!(duplicate.payload["kind"], "delivery_read_ack"); + assert_eq!(duplicate.payload["status"], "suppressed_duplicate"); + assert_eq!(duplicate.payload["reason"], "duplicate_delivery_read_ack"); + read_mock.assert_hits(0); +} + +#[tokio::test] +async fn slow_delivery_read_ack_does_not_block_confirmation_path() { + use httpmock::{Method::POST, MockServer}; + + let server = MockServer::start(); + let read_mock = server.mock(|when, then| { + when.method(POST) + .path("/v1/messages/msg_slow/read") + .header("authorization", "Bearer at_live_slow_recipient"); + then.status(200) + .delay(Duration::from_millis(200)) + .json_body(json!({ + "ok": true, + "data": { + "message_id": "msg_slow", + "agent_id": "agent_slow_recipient", + "read_at": "2026-06-08T10:00:00.000Z" + } + })); + }); + let client = RelaycastHttpClient::new(server.base_url(), "rk_live_test", "broker", "codex"); + seed_supplied_agent_token(&client, "recipient", "at_live_slow_recipient"); + let mut dedup = DedupCache::new(Duration::from_secs(300), 16); + let (tx, mut rx) = mpsc::channel(4); + let mut pending = HashMap::from([( + DeliveryId::new("del_slow"), + pending_delivery("recipient", "del_slow", "msg_slow"), + )]); + + let confirmed = clear_pending_delivery_if_event_matches( + &mut pending, + "del_slow", + Some("msg_slow"), + "recipient", + "delivery_ack", + ) + .expect("matching delivery_ack confirms the pending delivery"); + send_broker_event( + &tx, + BrokerEvent::MessageDeliveryConfirmed { + name: WorkerName::new("recipient"), + delivery_id: confirmed.delivery.delivery_id.clone(), + event_id: confirmed.delivery.event_id.clone(), + from: confirmed.delivery.from.clone(), + to: confirmed.delivery.target.clone(), + }, + ) + .await + .expect("confirmation event should enqueue before read-ack scheduling"); + + let start = Instant::now(); + mark_delivery_read_ack_with_timeout( + &client, + &tx, + &mut dedup, + &WorkerName::new("recipient"), + Some("codex"), + &confirmed.delivery.delivery_id, + &confirmed.delivery.event_id, + Duration::from_millis(20), + ); + assert!( + start.elapsed() < Duration::from_millis(50), + "read-ack scheduling must not wait for slow Relaycast mark_read" + ); + + let confirmation = tokio::time::timeout(Duration::from_millis(50), rx.recv()) + .await + .expect("delivery confirmation must not wait on mark_read") + .expect("confirmation event emitted"); + assert_eq!(confirmation.payload["kind"], "message_delivery_confirmed"); + assert_eq!(confirmation.payload["delivery_id"], "del_slow"); + + let read_ack = tokio::time::timeout(Duration::from_secs(1), rx.recv()) + .await + .expect("read-ack failure telemetry should arrive after timeout") + .expect("delivery_read_ack event emitted"); + assert_eq!(read_ack.payload["kind"], "delivery_read_ack"); + assert_eq!(read_ack.payload["status"], "failed"); + assert!(read_ack.payload["reason"] + .as_str() + .unwrap_or_default() + .contains("timed out")); + read_mock.assert_hits(1); +} + #[test] fn should_clear_pending_delivery_without_event_id_for_compatibility() { let pending = PendingDelivery { diff --git a/crates/broker/src/runtime/worker_events.rs b/crates/broker/src/runtime/worker_events.rs index dc576a086..df79021a2 100644 --- a/crates/broker/src/runtime/worker_events.rs +++ b/crates/broker/src/runtime/worker_events.rs @@ -6,8 +6,10 @@ impl BrokerRuntime { let paths = &self.paths; let state = &mut self.state; let sdk_out_tx = &self.sdk_out_tx; + let relaycast_http = self.relaycast_http.clone(); let ws_control_tx = &self.ws_control_tx; let workers = &mut self.workers; + let dedup = &mut self.dedup; let pending_deliveries = &mut self.pending_deliveries; let terminal_failed_deliveries = &mut self.terminal_failed_deliveries; let pending_requests = &mut self.pending_requests; @@ -65,6 +67,13 @@ impl BrokerRuntime { ) .await; if let Some(pending) = pending_for_confirmation { + let read_ack_delivery_id = pending.delivery.delivery_id.clone(); + let read_ack_event_id = pending.delivery.event_id.clone(); + let cli_hint = workers + .workers + .get(&name) + .and_then(|handle| handle.spec.cli.as_deref()) + .map(str::to_string); if let Some(handle) = workers.workers.get_mut(&name) { handle.last_activity_at = Instant::now(); handle.state = AgentWorkState::Working; @@ -80,6 +89,15 @@ impl BrokerRuntime { }, ) .await; + mark_delivery_read_ack( + &relaycast_http, + sdk_out_tx, + dedup, + &name, + cli_hint.as_deref(), + &read_ack_delivery_id, + &read_ack_event_id, + ); } } } else if msg_type == "delivery_queued" || msg_type == "delivery_injected" { diff --git a/packages/harness-driver/src/lifecycle-hooks.ts b/packages/harness-driver/src/lifecycle-hooks.ts index 88e43b2aa..a343f9b9d 100644 --- a/packages/harness-driver/src/lifecycle-hooks.ts +++ b/packages/harness-driver/src/lifecycle-hooks.ts @@ -159,6 +159,7 @@ export type DriverAgentActivityReason = | 'delivery_injected' | 'delivery_active' | 'delivery_ack' + | 'delivery_read_ack' | 'delivery_failed' | 'message_delivery_confirmed' | 'message_delivery_failed' diff --git a/packages/harness-driver/src/protocol.ts b/packages/harness-driver/src/protocol.ts index f6be68f71..9a01d2281 100644 --- a/packages/harness-driver/src/protocol.ts +++ b/packages/harness-driver/src/protocol.ts @@ -361,6 +361,14 @@ export type BrokerEvent = from: string; to: string; } + | { + kind: 'delivery_read_ack'; + name: string; + delivery_id: string; + event_id: string; + status: 'marked' | 'failed' | 'skipped_synthetic' | 'suppressed_duplicate'; + reason?: string; + } | { kind: 'message_delivery_failed'; name: string;