mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
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>
This commit is contained in:
@@ -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<RelaySession>, mpsc::Receiver<MatchedEvent>) {
|
||||
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<Mutex<SessionState>>,
|
||||
requests: Arc<Mutex<HashMap<String, PendingRequest>>>,
|
||||
events: broadcast::Sender<MatchedEvent>,
|
||||
/// 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<Mutex<Option<mpsc::Sender<MatchedEvent>>>>,
|
||||
wake: mpsc::Sender<()>,
|
||||
cancel: CancellationToken,
|
||||
}
|
||||
@@ -202,8 +187,10 @@ impl SessionState {
|
||||
}
|
||||
|
||||
impl RelaySession {
|
||||
pub(crate) fn subscribe(&self) -> broadcast::Receiver<MatchedEvent> {
|
||||
self.events.subscribe()
|
||||
async fn attach_archive(&self) -> mpsc::Receiver<MatchedEvent> {
|
||||
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<nostr::Tag>,
|
||||
) -> (Arc<RelaySession>, mpsc::Receiver<MatchedEvent>) {
|
||||
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<nostr::Tag>) -> Arc<RelaySession> {
|
||||
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(),
|
||||
|
||||
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user