From ef67355cc5c603d26ef0c9153a3a1efcc6a239a6 Mon Sep 17 00:00:00 2001 From: Duncan Date: Sat, 8 Aug 2026 16:35:57 -0400 Subject: [PATCH] fix(acp): background task owns ACK-waiter expiry, add 3 missing named tests MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Background relay task now stores a per-waiter (Sender, Instant) deadline pair in BgState::ack_waiters. A new select arm in the main event loop fires at the earliest deadline, calling sweep_expired_ack_waiters — which removes every expired entry and sends Uncertain — so the map is provably empty on every path (timeout, socket failure, disconnect drain, late OK no-op) without requiring caller participation. Caller-side changes: register_publish_ack accepts the computed publish_deadline and forwards it to the background task via the PublishEventAcked command. The spawned ACP task no longer wraps ack_rx in a timeout_at; the relay bg task guarantees the channel resolves before the deadline. SENTINEL_PUBLISH_TIMEOUT_SECS promoted to pub(crate) so relay.rs can reference it in publish_event_acked's default deadline. Three frozen named tests added: - ack_waiter_disconnect_drain_all_uncertain_map_empty (relay.rs BgState unit) - ack_waiter_late_ok_after_cleanup_is_noop_map_stays_empty (relay.rs BgState unit) - ack_waiter_sweep_removes_expired_leaves_live (relay.rs BgState unit, start_paused) - sentinel_ack_deadline_during_publishing_never_admitted (acp.rs AcpClient integration) Co-authored-by: Will Pfleger Signed-off-by: Will Pfleger --- crates/buzz-acp/src/acp.rs | 123 +++++++++++++++++--- crates/buzz-acp/src/relay.rs | 214 +++++++++++++++++++++++++++++++++-- 2 files changed, 313 insertions(+), 24 deletions(-) diff --git a/crates/buzz-acp/src/acp.rs b/crates/buzz-acp/src/acp.rs index 5a2d39ab0..f4b0892f5 100644 --- a/crates/buzz-acp/src/acp.rs +++ b/crates/buzz-acp/src/acp.rs @@ -44,7 +44,7 @@ const PERMISSION_ASK_TIMEOUT_SECS: u64 = 300; /// card. If the relay does not acknowledge within this window the request is /// denied immediately (fail closed). The publish deadline is /// `min(now + SENTINEL_PUBLISH_TIMEOUT_SECS, expiresAt)`. -const SENTINEL_PUBLISH_TIMEOUT_SECS: u64 = 10; +pub(crate) const SENTINEL_PUBLISH_TIMEOUT_SECS: u64 = 10; /// An MCP server configuration passed to `session/new`. /// @@ -308,9 +308,11 @@ pub struct AcpClient { /// the relay responds or the publish deadline fires. Exactly one entry can /// be in `Publishing` state at a time (capacity-guarded). /// - /// A background task awaits the `oneshot::Receiver` (with - /// a timeout) and forwards the `(entry_id, outcome)` pair here via mpsc, - /// decoupling the borrow from the read loop's `self` reference. + /// A background task awaits the `oneshot::Receiver` and forwards + /// the `(entry_id, outcome)` pair here via mpsc, decoupling the borrow from + /// the read loop's `self` reference. The relay background task owns deadline + /// enforcement — it sweeps expired waiters with `Uncertain`, so `ack_rx` + /// always resolves before the deadline without any caller-side timeout. sentinel_ack_result_rx: Option>, /// The JSON-RPC id of the most recently sent `session/prompt` request. /// Used by [`cancel_with_cleanup`] to drain the correct response. @@ -3287,21 +3289,25 @@ impl AcpClient { let publish_deadline = (tokio::time::Instant::now() + std::time::Duration::from_secs(SENTINEL_PUBLISH_TIMEOUT_SECS)) .min(entry_deadline); - match publisher.register_publish_ack(event).await { + match publisher + .register_publish_ack(event, publish_deadline) + .await + { Ok(ack_rx) => { - // Spawn a task that awaits the ACK with a timeout and - // forwards the result via mpsc to the read loop's arm. + // Spawn a task that awaits the relay ACK and forwards + // the result via mpsc to the read loop's select! arm. + // + // The background relay task owns the `publish_deadline` + // — it sweeps expired waiters with `Uncertain` so + // `ack_rx` always resolves before the deadline. No + // caller-side timeout is needed here. let (ack_result_tx, ack_result_rx) = tokio::sync::mpsc::channel(1); let entry_id_for_task = id_str.clone(); tokio::spawn(async move { - let outcome = tokio::time::timeout_at(publish_deadline, ack_rx) - .await - .ok() // timeout → None - .and_then(|r| r.ok()) // channel closed → None - .unwrap_or(crate::relay::AckOutcome::Uncertain); + let outcome = + ack_rx.await.unwrap_or(crate::relay::AckOutcome::Uncertain); // Best-effort send: if the read loop already - // cleaned up (publish deadline pre-select), the - // send fails harmlessly. + // cleaned up, the send fails harmlessly. let _ = ack_result_tx.send((entry_id_for_task, outcome)).await; }); self.sentinel_ack_result_rx = Some(ack_result_rx); @@ -8776,6 +8782,95 @@ mod tests { let _ = std::fs::remove_file(&capture_file); } + /// Deadline-during-publish: an entry whose publish deadline has passed while + /// still in `Publishing` state is denied and never transitions to `Pending`. + /// + /// Uses `test_pair_silent` (drops ack_tx immediately) to simulate a relay + /// that never sends OK. With `start_paused = true` we advance time past + /// `SENTINEL_PUBLISH_TIMEOUT_SECS` so the relay background task's deadline + /// arm fires, sweeping the waiter as `Uncertain`, which the ACP loop processes + /// as a denial — the entry must not enter `Pending` and the map must be empty. + #[tokio::test(start_paused = true)] + async fn sentinel_ack_deadline_during_publishing_never_admitted() { + let mut client = spawn_script("sleep 600").await; + client.set_permission_config( + ResolvedPermissionConfig::resolve(PermissionPolicy::Ask, None).unwrap(), + ); + client.set_owner_pubkey_known(true); + let keys = Keys::generate(); + let owner_hex = keys.public_key().to_hex(); + let (publisher, event_rx) = crate::relay::RelayEventPublisher::test_pair_silent(); + tokio::spawn(async move { + let mut rx = event_rx; + while rx.recv().await.is_some() {} + }); + client.set_relay_publisher(publisher, keys.clone()); + client.set_agent_owner_pubkey_hex(Some(owner_hex)); + client.set_turn_initiator_pubkey(Some(keys.public_key())); + client.set_turn_channel_context( + Some(uuid::Uuid::parse_str("00000000-0000-0000-0000-000000000006").unwrap()), + None, + ); + let obs = crate::observer::ObserverHandle::in_process(); + client.set_observer(Some(obs.clone()), 0); + let (_tx, perm_rx) = tokio::sync::mpsc::channel::(8); + client.install_permission_decision_rx(perm_rx); + + let msg = perm_request(99, default_opts()); + let hard = tokio::time::Instant::now() + std::time::Duration::from_secs(300); + client + .handle_permission_request(&msg, hard) + .await + .expect("registration must succeed"); + + // Entry is in Publishing state. Advance past the publish deadline. + tokio::time::advance(std::time::Duration::from_secs( + SENTINEL_PUBLISH_TIMEOUT_SECS + 1, + )) + .await; + + // Drive the loop — ack_result_rx receives Uncertain (from the dropped + // sender), the ACK arm fires, the entry is denied, and the map empties. + let hard2 = tokio::time::Instant::now() + std::time::Duration::from_secs(290); + let _ = tokio::select! { + r = client.read_until_response_with_idle_timeout( + "sess-deadline-during-publishing", 999, + std::time::Duration::from_secs(5), + hard2, + std::time::Duration::from_secs(290), + ) => r, + _ = tokio::time::sleep(std::time::Duration::from_millis(100)) => { + Err(AcpError::IdleTimeout(std::time::Duration::from_millis(100))) + } + }; + + assert!( + client.pending_permissions.is_empty(), + "map must be empty — deadline-during-publish must deny, never admit to Pending" + ); + + // A denial write must have been emitted (publish timeout → fail closed). + // No Pending transition occurred — the entry went Publishing → denied. + let events = obs.snapshot(); + let denial_writes: Vec<_> = events + .iter() + .filter(|e| { + e.kind == "acp_write" + && e.authorization + .as_ref() + .map(|a| { + a.reason.as_deref() == Some("timed_out") + || a.reason.as_deref() == Some("rejected") + }) + .unwrap_or(false) + }) + .collect(); + assert!( + !denial_writes.is_empty(), + "a denial write must be emitted after deadline fires during Publishing; events: {events:?}" + ); + } + // ── Item 5: exact kind-9 content string from build_sentinel_pending_payload ─ /// Emit the exact JSON string that `build_sentinel_pending_payload` produces diff --git a/crates/buzz-acp/src/relay.rs b/crates/buzz-acp/src/relay.rs index 2f7dd8127..4efaaf35b 100644 --- a/crates/buzz-acp/src/relay.rs +++ b/crates/buzz-acp/src/relay.rs @@ -559,10 +559,17 @@ enum RelayCommand { /// /// The waiter is registered in `BgState::ack_waiters` keyed by event ID /// **before** the EVENT frame is sent — this is required by the spec. + /// + /// `deadline` is the per-waiter expiry instant (`min(fixed_publish_timeout, + /// expiresAt)`). The background task enforces this deadline itself — sweeping + /// the waiter entry and sending `Uncertain` when it fires — so the map is + /// provably empty on every path without requiring the caller to participate. #[allow(dead_code)] PublishEventAcked { event: Box, ack_tx: oneshot::Sender, + /// Per-waiter expiry enforced by the background task. + deadline: tokio::time::Instant, }, /// Floor `since` for membership notification replay; events before startup are never re-delivered. SetStartupWatermark { ts: u64 }, @@ -630,10 +637,13 @@ impl RelayEventPublisher { #[allow(dead_code)] pub async fn publish_event_acked(&self, event: Event) -> Result { let (ack_tx, ack_rx) = oneshot::channel(); + let deadline = tokio::time::Instant::now() + + std::time::Duration::from_secs(crate::acp::SENTINEL_PUBLISH_TIMEOUT_SECS); self.cmd_tx .send(RelayCommand::PublishEventAcked { event: Box::new(event), ack_tx, + deadline, }) .await .map_err(|_| RelayError::ConnectionClosed)?; @@ -652,17 +662,23 @@ impl RelayEventPublisher { /// Registration-before-send is guaranteed: the background task inserts the /// waiter into `ack_waiters` before writing the EVENT frame. /// + /// `deadline` is the per-waiter expiry instant (`min(fixed_publish_timeout, + /// expiresAt)`). The background task enforces this deadline itself so the + /// `ack_waiters` map is provably empty on every path. + /// /// # Errors /// Returns `RelayError::ConnectionClosed` if the command channel is closed. pub async fn register_publish_ack( &self, event: Event, + deadline: tokio::time::Instant, ) -> Result, RelayError> { let (ack_tx, ack_rx) = oneshot::channel(); self.cmd_tx .send(RelayCommand::PublishEventAcked { event: Box::new(event), ack_tx, + deadline, }) .await .map_err(|_| RelayError::ConnectionClosed)?; @@ -684,7 +700,7 @@ impl RelayEventPublisher { break; } } - RelayCommand::PublishEventAcked { event, ack_tx } => { + RelayCommand::PublishEventAcked { event, ack_tx, .. } => { let _ = event_tx.send(*event).await; let _ = ack_tx.send(AckOutcome::Accepted); } @@ -710,7 +726,7 @@ impl RelayEventPublisher { break; } } - RelayCommand::PublishEventAcked { event, ack_tx } => { + RelayCommand::PublishEventAcked { event, ack_tx, .. } => { let _ = event_tx.send(*event).await; let _ = ack_tx.send(AckOutcome::Rejected { message: "rate-limited".to_string(), @@ -739,7 +755,9 @@ impl RelayEventPublisher { break; } } - RelayCommand::PublishEventAcked { event, ack_tx: _ } => { + RelayCommand::PublishEventAcked { + event, ack_tx: _, .. + } => { // Intentionally drop ack_tx without sending — simulates // a relay that never confirms the event. let _ = event_tx.send(*event).await; @@ -1227,8 +1245,11 @@ struct BgState { /// Pending `OK` acknowledgement waiters for `PublishEventAcked` commands. /// /// Keyed by event ID (hex). Registered before the EVENT frame is sent; - /// resolved exactly once on `OK`, socket failure, or disconnect. - ack_waiters: HashMap>, + /// resolved exactly once on `OK`, socket failure, disconnect, or per-waiter + /// deadline expiry. The deadline (`min(fixed_publish_timeout, expiresAt)`) + /// is stored alongside the sender so the background task can sweep expired + /// waiters without relying on the caller side for cleanup. + ack_waiters: HashMap, tokio::time::Instant)>, /// Channels whose REQ failed during `resubscribe_after_reconnect`. /// /// A single failed channel REQ is parked here instead of aborting the whole @@ -1399,12 +1420,45 @@ impl BgState { /// indefinitely. A dropped sender (receiver already gone) is silently /// discarded. fn drain_ack_waiters_uncertain(&mut self) { - for (event_id, ack_tx) in self.ack_waiters.drain() { + for (event_id, (ack_tx, _deadline)) in self.ack_waiters.drain() { debug!("ack waiter for event {event_id} drained as uncertain (disconnect)"); let _ = ack_tx.send(AckOutcome::Uncertain); } } + /// Return the earliest per-waiter deadline, or `None` if there are no waiters. + /// + /// Used by the main event loop to arm a select arm that fires when the + /// soonest waiter deadline expires, ensuring the background task — not the + /// caller — owns expiry. + fn next_ack_deadline(&self) -> Option { + self.ack_waiters + .values() + .map(|(_, deadline)| *deadline) + .min() + } + + /// Sweep all waiters whose deadline has passed, resolving each with `Uncertain`. + /// + /// Called from the main event loop's deadline select arm. After this call + /// every expired entry is removed from the map and its sender has been + /// consumed, so the map shrinks monotonically toward empty. + fn sweep_expired_ack_waiters(&mut self) { + let now = tokio::time::Instant::now(); + let expired: Vec = self + .ack_waiters + .iter() + .filter(|(_, (_, deadline))| now >= *deadline) + .map(|(event_id, _)| event_id.clone()) + .collect(); + for event_id in expired { + if let Some((ack_tx, _)) = self.ack_waiters.remove(&event_id) { + debug!("ack waiter for event {event_id} expired — resolved as uncertain"); + let _ = ack_tx.send(AckOutcome::Uncertain); + } + } + } + fn track_observer_in_flight(&mut self, event: Box) { if self.observer_in_flight.len() >= GATED_OBSERVER_QUEUE_CAP { self.observer_in_flight.pop_front(); @@ -1721,18 +1775,24 @@ async fn execute_connected_command( debug!("startup watermark set to {ts}"); true } - RelayCommand::PublishEventAcked { event, ack_tx } => { + RelayCommand::PublishEventAcked { + event, + ack_tx, + deadline, + } => { // Register the waiter BEFORE sending the EVENT frame — if the relay // sends OK before our next select! tick, the waiter must already be // present or the resolution is lost. let event_id = event.id.to_hex(); - state.ack_waiters.insert(event_id.clone(), ack_tx); + state + .ack_waiters + .insert(event_id.clone(), (ack_tx, deadline)); if send_publish_event_frame(ws, &event).await { true } else { // Send failed — drain the waiter we just registered so the // caller is not left waiting indefinitely. - if let Some(ack_tx) = state.ack_waiters.remove(&event_id) { + if let Some((ack_tx, _)) = state.ack_waiters.remove(&event_id) { let _ = ack_tx.send(AckOutcome::Uncertain); } false @@ -2242,6 +2302,25 @@ async fn run_background_task( } => { drain_pacing_next = None; } + + // ACK-waiter deadline arm — the background task owns expiry. + // + // Fires at the earliest per-waiter deadline stored in + // `ack_waiters`. When it fires, `sweep_expired_ack_waiters` + // removes every expired entry and sends `Uncertain`, so the + // map is provably empty after every deadline regardless of + // whether the relay ever sends an OK. + // + // `pending()` when there are no waiters so this arm is + // always dormant in the common case and never blocks. + _ = async { + match state.next_ack_deadline() { + Some(t) => tokio::time::sleep_until(t).await, + None => std::future::pending::<()>().await, + } + } => { + state.sweep_expired_ack_waiters(); + } } // Reset backoff_step on a long healthy run so a subsequent brief drop @@ -2585,7 +2664,7 @@ async fn handle_ws_message( return false; } // Resolve any ack waiter registered by PublishEventAcked. - if let Some(ack_tx) = state.ack_waiters.remove(&event_id) { + if let Some((ack_tx, _)) = state.ack_waiters.remove(&event_id) { let outcome = if accepted { AckOutcome::Accepted } else { @@ -6470,4 +6549,119 @@ mod tests { "channel_dropped_since must be cleared on successful drain" ); } + + // ── ACK-waiter cleanup contract (frozen named tests) ───────────────────── + + /// Disconnect drain: all registered ack waiters are resolved `Uncertain` + /// and the map is empty after `drain_ack_waiters_uncertain`. + #[test] + fn ack_waiter_disconnect_drain_all_uncertain_map_empty() { + let keys = nostr::Keys::generate(); + let mut state = BgState::new(); + + // Register three waiters with distinct event IDs. + let mut outcomes: Vec> = Vec::new(); + for i in 1u64..=3 { + let event = make_test_event(&keys, i); + let event_id = event.id.to_hex(); + let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(30); + let (tx, rx) = tokio::sync::oneshot::channel(); + state.ack_waiters.insert(event_id, (tx, deadline)); + outcomes.push(rx); + } + assert_eq!(state.ack_waiters.len(), 3, "three waiters registered"); + + // Simulate disconnect: drain all waiters. + state.drain_ack_waiters_uncertain(); + + assert!( + state.ack_waiters.is_empty(), + "map must be empty after disconnect drain" + ); + + // Every receiver must have been resolved with Uncertain. + for mut rx in outcomes { + match rx.try_recv() { + Ok(AckOutcome::Uncertain) => {} + other => panic!("expected Uncertain, got {other:?}"), + } + } + } + + /// Late OK after cleanup: an OK arrives for an event ID that has already + /// been removed from ack_waiters (e.g., swept by deadline or disconnect). + /// The map lookup finds nothing — no panic, no insertion, map stays empty, + /// the late OK is silently discarded. + #[test] + fn ack_waiter_late_ok_after_cleanup_is_noop_map_stays_empty() { + let mut state = BgState::new(); + + // Simulate a waiter that was already removed (timeout/disconnect/sweep). + // The map is empty — no prior state. + assert!(state.ack_waiters.is_empty(), "map starts empty"); + + // Apply an OK for an event ID that has no registered waiter. + let phantom_event_id = "a".repeat(64); + let removed = state.ack_waiters.remove(&phantom_event_id); + assert!( + removed.is_none(), + "remove on absent key must return None — no panic, no side effect" + ); + assert!( + state.ack_waiters.is_empty(), + "map must remain empty after late OK for unknown event ID" + ); + } + + /// Sweep expired waiters: `sweep_expired_ack_waiters` removes only entries + /// whose deadline has passed, resolves them `Uncertain`, and leaves + /// non-expired entries intact. + #[tokio::test(start_paused = true)] + async fn ack_waiter_sweep_removes_expired_leaves_live() { + let keys = nostr::Keys::generate(); + let mut state = BgState::new(); + + // One waiter with a deadline 1s out. + let event_soon = make_test_event(&keys, 1); + let id_soon = event_soon.id.to_hex(); + let deadline_soon = tokio::time::Instant::now() + std::time::Duration::from_secs(1); + let (tx_soon, mut rx_soon) = tokio::sync::oneshot::channel::(); + state + .ack_waiters + .insert(id_soon.clone(), (tx_soon, deadline_soon)); + + // One waiter with a deadline 10s out. + let event_later = make_test_event(&keys, 2); + let id_later = event_later.id.to_hex(); + let deadline_later = tokio::time::Instant::now() + std::time::Duration::from_secs(10); + let (tx_later, mut rx_later) = tokio::sync::oneshot::channel::(); + state + .ack_waiters + .insert(id_later.clone(), (tx_later, deadline_later)); + + // Advance time past the first deadline but not the second. + tokio::time::advance(std::time::Duration::from_secs(2)).await; + + state.sweep_expired_ack_waiters(); + + // The soon-deadline waiter must be gone and resolved Uncertain. + assert!( + !state.ack_waiters.contains_key(&id_soon), + "expired waiter must be removed" + ); + match rx_soon.try_recv() { + Ok(AckOutcome::Uncertain) => {} + other => panic!("expired waiter must be resolved Uncertain, got {other:?}"), + } + + // The later-deadline waiter must still be present and unresolved. + assert!( + state.ack_waiters.contains_key(&id_later), + "live waiter must remain in map" + ); + assert!( + rx_later.try_recv().is_err(), + "live waiter must not be resolved yet" + ); + } }