From 4cca14d3f201e48ff5eb6e9b86ef0f579686e4d6 Mon Sep 17 00:00:00 2001 From: Wren <5217c5c2f7bfb4333e46d17c98a9255a52dadee18dcd43a43536b95e6776dfa0@buzz.block.builderlab.xyz> Date: Sun, 16 Aug 2026 01:05:06 -0400 Subject: [PATCH] fix(desktop): preserve archive event backpressure Route persistent archive events directly through one bounded mpsc channel and await delivery in the socket loop. Archive subscriptions are live-only, so replay cannot repair broadcast eviction; throttling preserves B's no-loss contract while finite catalog requests remain isolated in the request map. The catalog intentionally verifies finite events twice: transport verification bounds memory against forged input, while catalog verification keeps its projection helper sound for every caller despite the measured cost. Co-authored-by: Wren <5217c5c2f7bfb4333e46d17c98a9255a52dadee18dcd43a43536b95e6776dfa0@buzz.block.builderlab.xyz> Signed-off-by: Wren <5217c5c2f7bfb4333e46d17c98a9255a52dadee18dcd43a43536b95e6776dfa0@buzz.block.builderlab.xyz> --- desktop/src-tauri/src/native_relay_client.rs | 85 +++++++------------ .../src/native_relay_client_tests.rs | 81 ++++++++++++++++-- 2 files changed, 103 insertions(+), 63 deletions(-) diff --git a/desktop/src-tauri/src/native_relay_client.rs b/desktop/src-tauri/src/native_relay_client.rs index 3f80dcb70..d63539883 100644 --- a/desktop/src-tauri/src/native_relay_client.rs +++ b/desktop/src-tauri/src/native_relay_client.rs @@ -27,7 +27,7 @@ use std::{ use buzz_ws_client_pkg::{NostrWsConnection, RelayMessage}; use nostr::{Event, Keys}; use tokio::{ - sync::{broadcast, mpsc, oneshot, Mutex}, + sync::{mpsc, oneshot, Mutex}, time::Instant, }; use tokio_util::sync::CancellationToken; @@ -118,26 +118,7 @@ impl NativeRelayClient { keys: Keys, ) -> (Arc, mpsc::Receiver) { let session = self.ensure_session(relay_url, keys).await; - let mut events = session.subscribe(); - let (event_tx, event_rx) = mpsc::channel(256); - tauri::async_runtime::spawn(async move { - loop { - match events.recv().await { - Ok(event) => { - if event_tx.send(event).await.is_err() { - return; - } - } - Err(broadcast::error::RecvError::Lagged(skipped)) => { - eprintln!( - "buzz-desktop: archive relay receiver lagged by {skipped} events; stopping sync rather than silently losing archive data" - ); - return; - } - Err(broadcast::error::RecvError::Closed) => return, - } - } - }); + let event_rx = session.attach_archive().await; (session, event_rx) } } @@ -145,7 +126,11 @@ impl NativeRelayClient { pub(crate) struct RelaySession { state: Arc>, requests: Arc>>, - events: broadcast::Sender, + /// The archive is the sole persistent-event consumer. Sending through its + /// bounded channel is awaited by the socket loop, preserving the + /// backpressure required by live-only (`limit: 0`) subscriptions: dropping + /// an event here cannot be repaired by replaying it later. + archive_events: Arc>>>, wake: mpsc::Sender<()>, cancel: CancellationToken, } @@ -202,8 +187,10 @@ impl SessionState { } impl RelaySession { - pub(crate) fn subscribe(&self) -> broadcast::Receiver { - self.events.subscribe() + async fn attach_archive(&self) -> mpsc::Receiver { + let (events, receiver) = mpsc::channel(256); + *self.archive_events.lock().await = Some(events); + receiver } /// Fetches one finite page over this session without disturbing persistent @@ -295,41 +282,22 @@ impl RelaySession { /// desired set — never a snapshot captured at connect time, so a subscription /// change during an outage is honored by the reconnect that follows. #[cfg(test)] -pub(crate) fn start( +pub(crate) async fn start( relay_url: String, keys: Keys, auth_tag: Option, ) -> (Arc, mpsc::Receiver) { let session = start_managed(relay_url, keys, auth_tag); - let mut events = session.subscribe(); - let (event_tx, event_rx) = mpsc::channel(256); - tauri::async_runtime::spawn(async move { - loop { - match events.recv().await { - Ok(event) => { - if event_tx.send(event).await.is_err() { - return; - } - } - Err(broadcast::error::RecvError::Lagged(skipped)) => { - eprintln!( - "buzz-desktop: native_relay_client: legacy receiver lagged by {skipped} events" - ); - } - Err(broadcast::error::RecvError::Closed) => return, - } - } - }); - (session, event_rx) + let events = session.attach_archive().await; + (session, events) } fn start_managed(relay_url: String, keys: Keys, auth_tag: Option) -> Arc { - let (events, _) = broadcast::channel(256); let (wake, wake_rx) = mpsc::channel(1); let session = Arc::new(RelaySession { state: Arc::new(Mutex::new(SessionState::default())), requests: Arc::new(Mutex::new(HashMap::new())), - events, + archive_events: Arc::new(Mutex::new(None)), wake, cancel: CancellationToken::new(), }); @@ -497,13 +465,20 @@ async fn run_connection( // accumulated backoff for it is stale. Mirrors the JS // port's per-event `closedRetryAttempt = 0`. retries.remove(&subscription_id); - // A broadcast session remains healthy when no feature - // currently observes persistent events (catalog fetches - // are fulfilled above through `requests`). - let _ = session.events.send(MatchedEvent { - subscription_id, - event, - }); + // Persistent archive subscriptions are live-only, so + // losing an event cannot be repaired with a later REQ. + // Await the bounded archive channel to push back on the + // socket read loop instead. Finite catalog requests are + // fulfilled above and never enter this channel. + let sender = session.archive_events.lock().await.clone(); + if let Some(sender) = sender { + let _ = sender + .send(MatchedEvent { + subscription_id, + event, + }) + .await; + } } Ok(RelayMessage::Closed { subscription_id, message }) => { // The relay dropped it; forget it so a reopen re-sends @@ -827,7 +802,7 @@ mod relay_backed_tests { // shape: that the `#p` tag key and the `limit: 0` live tail produce a // REQ a real relay accepts and answers. Scope demultiplexing on the // archive side is covered in `archive/sync_tests.rs`. - let (session, mut events) = start(relay_url.clone(), owner.clone(), None); + let (session, mut events) = start(relay_url.clone(), owner.clone(), None).await; session .set_subscriptions(vec![Subscription { id: "archive:owner_p:test".to_string(), diff --git a/desktop/src-tauri/src/native_relay_client_tests.rs b/desktop/src-tauri/src/native_relay_client_tests.rs index 25cc69f52..21851d5c0 100644 --- a/desktop/src-tauri/src/native_relay_client_tests.rs +++ b/desktop/src-tauri/src/native_relay_client_tests.rs @@ -159,7 +159,7 @@ async fn settle() { #[tokio::test] async fn finite_fetch_multiplexes_with_persistent_delivery_on_a_real_websocket() { let (relay_url, mut frames, commands) = stub_relay().await; - let (session, mut events) = start(relay_url, Keys::generate(), None); + let (session, mut events) = start(relay_url, Keys::generate(), None).await; session.set_subscriptions(vec![probe_subscription()]).await; assert_eq!(next_req(&mut frames, "the persistent REQ").await, PROBE_ID); @@ -229,6 +229,71 @@ async fn finite_fetch_multiplexes_with_persistent_delivery_on_a_real_websocket() session.shutdown(); } +async fn run_persistent_burst(drain_concurrently: bool) { + const BURST: usize = 1_200; + + let (relay_url, mut frames, commands) = stub_relay().await; + let (session, mut events) = start(relay_url, Keys::generate(), None).await; + session.set_subscriptions(vec![probe_subscription()]).await; + assert_eq!(next_req(&mut frames, "the burst REQ").await, PROBE_ID); + + let relay_keys = Keys::generate(); + let event = EventBuilder::text_note("persistent burst event") + .sign_with_keys(&relay_keys) + .unwrap(); + let send_burst = tokio::spawn({ + let commands = commands.clone(); + let event = serde_json::to_value(&event).unwrap(); + async move { + for _ in 0..BURST { + commands + .send(StubCommand::Event(PROBE_ID.into(), event.clone())) + .await + .unwrap(); + } + } + }); + + if !drain_concurrently { + // Let the bounded archive channel fill before draining. The socket loop + // must wait here rather than evicting live-only events. + tokio::time::sleep(Duration::from_millis(100)).await; + } + for _ in 0..BURST { + tokio::time::timeout(Duration::from_secs(60), events.recv()) + .await + .expect("timed out draining persistent burst") + .expect("archive receiver closed during persistent burst"); + } + send_burst.await.unwrap(); + + let after = EventBuilder::text_note("persistent event after burst") + .sign_with_keys(&relay_keys) + .unwrap(); + commands + .send(StubCommand::Event( + PROBE_ID.into(), + serde_json::to_value(&after).unwrap(), + )) + .await + .unwrap(); + let delivered = tokio::time::timeout(Duration::from_secs(60), events.recv()) + .await + .expect("timed out after persistent burst") + .expect("archive receiver closed after persistent burst"); + assert_eq!(*delivered.event, after); + session.shutdown(); +} + +/// Persistent archive subscriptions use `limit: 0`, so an event lost during a +/// slow-consumer burst cannot be replayed. Both a fast control and a receiver +/// that starts late must therefore get the whole burst and remain live after it. +#[tokio::test] +async fn persistent_delivery_applies_backpressure_without_losing_a_burst() { + run_persistent_burst(true).await; + run_persistent_burst(false).await; +} + /// The blocker: a CLOSED with the desired set never changing again must /// still reopen the subscription. /// @@ -239,7 +304,7 @@ async fn finite_fetch_multiplexes_with_persistent_delivery_on_a_real_websocket() #[tokio::test] async fn a_closed_subscription_reopens_without_a_desired_set_change() { let (relay_url, mut frames, closed) = stub_relay().await; - let (session, _events) = start(relay_url, Keys::generate(), None); + let (session, _events) = start(relay_url, Keys::generate(), None).await; session.set_subscriptions(vec![probe_subscription()]).await; assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); @@ -267,7 +332,7 @@ async fn a_closed_subscription_reopens_without_a_desired_set_change() { #[tokio::test] async fn a_terminal_closed_is_not_retried_on_the_same_socket() { let (relay_url, mut frames, closed) = stub_relay().await; - let (session, _events) = start(relay_url, Keys::generate(), None); + let (session, _events) = start(relay_url, Keys::generate(), None).await; session.set_subscriptions(vec![probe_subscription()]).await; assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); @@ -303,7 +368,7 @@ async fn a_terminal_closed_is_not_retried_on_the_same_socket() { #[tokio::test] async fn a_recreated_subscription_does_not_inherit_a_terminal_latch() { let (relay_url, mut frames, closed) = stub_relay().await; - let (session, _events) = start(relay_url, Keys::generate(), None); + let (session, _events) = start(relay_url, Keys::generate(), None).await; session.set_subscriptions(vec![probe_subscription()]).await; assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); @@ -347,7 +412,7 @@ async fn a_recreated_subscription_does_not_inherit_a_terminal_latch() { #[tokio::test] async fn a_recreated_subscription_is_not_suppressed_when_the_writes_coalesce() { let (relay_url, mut frames, closed) = stub_relay().await; - let (session, _events) = start(relay_url, Keys::generate(), None); + let (session, _events) = start(relay_url, Keys::generate(), None).await; session.set_subscriptions(vec![probe_subscription()]).await; assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); @@ -382,7 +447,7 @@ async fn a_recreated_subscription_is_not_suppressed_when_the_writes_coalesce() { #[tokio::test] async fn a_reconcile_preserves_the_backoff_of_a_still_desired_subscription() { let (relay_url, mut frames, closed) = stub_relay().await; - let (session, _events) = start(relay_url, Keys::generate(), None); + let (session, _events) = start(relay_url, Keys::generate(), None).await; session.set_subscriptions(vec![probe_subscription()]).await; assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); @@ -436,7 +501,7 @@ async fn a_reconcile_preserves_the_backoff_of_a_still_desired_subscription() { #[tokio::test] async fn a_closed_arriving_after_removal_does_not_mint_retry_state() { let (relay_url, mut frames, closed) = stub_relay().await; - let (session, _events) = start(relay_url, Keys::generate(), None); + let (session, _events) = start(relay_url, Keys::generate(), None).await; session.set_subscriptions(vec![probe_subscription()]).await; assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); @@ -485,7 +550,7 @@ async fn a_closed_arriving_after_removal_does_not_mint_retry_state() { #[tokio::test] async fn a_stale_terminal_closed_does_not_blackhole_a_recreated_subscription() { let (relay_url, mut frames, closed) = stub_relay().await; - let (session, mut events) = start(relay_url, Keys::generate(), None); + let (session, mut events) = start(relay_url, Keys::generate(), None).await; session.set_subscriptions(vec![probe_subscription()]).await; assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID);