From bf253bcdfefb64e7fa9bcc402268f8675fe0a084 Mon Sep 17 00:00:00 2001 From: James Lopez Date: Fri, 28 Aug 2026 10:59:22 -0500 Subject: [PATCH] chore(acp): trace channel delivery boundaries Signed-off-by: James Lopez --- crates/buzz-acp/src/relay.rs | 64 ++++++++++++++++++++++++++++++++---- 1 file changed, 57 insertions(+), 7 deletions(-) diff --git a/crates/buzz-acp/src/relay.rs b/crates/buzz-acp/src/relay.rs index 6188e57a11d..5d133300d68 100644 --- a/crates/buzz-acp/src/relay.rs +++ b/crates/buzz-acp/src/relay.rs @@ -2228,6 +2228,13 @@ async fn handle_ws_message( } else if let Some(channel_id) = channel_id_from_sub_id(&subscription_id) { let ts = event.created_at.as_secs(); let event_id_hex = event.id.to_hex(); + debug!( + channel_id = %channel_id, + subscription_id = %subscription_id, + event_id = %event_id_hex, + kind = event.kind.as_u16(), + "channel event received from relay" + ); if state.record_event(channel_id, &event) { let buzz_event = BuzzEvent { channel_id, @@ -2244,7 +2251,14 @@ async fn handle_ws_message( ); } match event_tx.try_send(Some(buzz_event)) { - Ok(()) => {} + Ok(()) => { + debug!( + channel_id = %channel_id, + subscription_id = %subscription_id, + event_id = %event_id_hex, + "channel event forwarded to harness" + ); + } Err(mpsc::error::TrySendError::Full(_)) => { // Remove from dedup set so the replayed event // won't be rejected as a duplicate after reconnect. @@ -3289,12 +3303,12 @@ async fn send_subscribe( match ws_send_timeout(ws, Message::Text(text.into()), WS_SEND_TIMEOUT_SECS).await { Ok(()) => { debug!( - "subscribed to channel {channel_id}{}", - if since.is_some() { - " (with since filter)" - } else { - " (since=now)" - } + channel_id = %channel_id, + subscription_id = %sub_id, + kinds = ?filter.kinds, + require_mention = filter.require_mention, + since = since_ts, + "channel subscription REQ sent" ); true } @@ -4483,6 +4497,42 @@ mod tests { (client, server.await.expect("join test websocket server")) } + #[tokio::test] + async fn channel_event_is_forwarded_from_matching_subscription() { + let (mut client, _server) = test_ws_pair().await; + let (event_tx, mut event_rx) = mpsc::channel(1); + let (observer_tx, _observer_rx) = mpsc::channel(1); + let mut state = BgState::new(); + let channel_id = Uuid::new_v4(); + let event = make_test_event(&nostr::Keys::generate(), 1_000_000); + let event_id = event.id; + let text = serde_json::to_string(&json!(["EVENT", channel_sub_id(channel_id), event,])) + .expect("serialize relay event"); + + assert!( + handle_ws_message( + Message::Text(text.into()), + &mut client, + &event_tx, + &observer_tx, + &mut state, + &nostr::Keys::generate(), + "ws://relay.example", + "agent-pubkey", + None, + ) + .await + ); + + let forwarded = timeout(Duration::from_secs(1), event_rx.recv()) + .await + .expect("timed out waiting for forwarded channel event") + .expect("event channel unexpectedly closed") + .expect("connection-loss sentinel instead of channel event"); + assert_eq!(forwarded.channel_id, channel_id); + assert_eq!(forwarded.event.id, event_id); + } + async fn next_test_frame( server: &mut WebSocketStream, ) -> serde_json::Value {