From a49b32c917a97500245b0e54dc5cdb67a0f7aa70 Mon Sep 17 00:00:00 2001 From: dawn <90008@gaze.systems> Date: Tue, 27 Jan 2026 21:21:12 +0300 Subject: [PATCH] rewrite to not ack events on drop but manually ack, dont use tokio tasks and buffers --- README.md | 20 ++-- standard-site-sync/src/main.rs | 16 +-- tapped/src/channel.rs | 188 ++++++++++----------------------- tapped/src/client.rs | 13 +-- tapped/src/handle.rs | 5 +- tapped/src/lib.rs | 13 +-- tapped/tests/integration.rs | 2 +- 7 files changed, 94 insertions(+), 163 deletions(-) diff --git a/README.md b/README.md index 3f468b5..86f85e4 100644 --- a/README.md +++ b/README.md @@ -51,10 +51,10 @@ async fn main() -> tapped::Result<()> { let handle = TapHandle::spawn_default(config).await?; // Subscribe to events - let mut channel = handle.channel().await?; + let (mut receiver, mut ack_sender) = handle.channel().await?; - while let Ok(received) = channel.recv().await { - match &received.event { + while let Ok((event, ack_id)) = receiver.recv().await { + match event { Event::Record(record) => { println!("[{:?}] {}/{}", record.action, @@ -66,7 +66,8 @@ async fn main() -> tapped::Result<()> { println!("Identity: {} -> {}", identity.did, identity.handle); } } - // Event is auto-acknowledged when `received` is dropped + // Manual acknowledgment required + ack_sender.ack(ack_id).await?; } Ok(()) @@ -164,15 +165,15 @@ let config = TapConfig::builder() ### Working with Events -Events are automatically acknowledged when dropped: +Events must be manually acknowledged: ```rust use tapped::{Event, RecordAction}; -let mut channel = client.channel().await?; +let (mut receiver, mut ack_sender) = client.channel().await?; -while let Ok(received) = channel.recv().await { - match &received.event { +while let Ok((event, ack_id)) = receiver.recv().await { + match event { Event::Record(record) => { match record.action { RecordAction::Create => { @@ -193,7 +194,8 @@ while let Ok(received) = channel.recv().await { println!("{} is now @{}", identity.did, identity.handle); } } - // Ack sent automatically here when `received` goes out of scope + // Ack must be sent manually + ack_sender.ack(ack_id).await?; } ``` diff --git a/standard-site-sync/src/main.rs b/standard-site-sync/src/main.rs index 05e80d2..3a5ab4d 100644 --- a/standard-site-sync/src/main.rs +++ b/standard-site-sync/src/main.rs @@ -67,7 +67,7 @@ async fn main() -> Result<(), Box> { client.health().await?; info!("Tap is healthy!"); - let mut receiver = client.channel().await?; + let (mut receiver, mut ack_sender) = client.channel().await?; info!("Connected! Waiting for events..."); // In-memory cache - load from disk if available @@ -83,8 +83,8 @@ async fn main() -> Result<(), Box> { loop { match receiver.recv().await { - Ok(received) => { - if let Event::Record(ref record_event) = *received { + Ok((event, ack_id)) => { + if let Event::Record(ref record_event) = event { // Track live vs backfill if record_event.live { live_count += 1; @@ -96,19 +96,23 @@ async fn main() -> Result<(), Box> { write_output_files(&cache)?; // Periodically show event source breakdown - if (live_count + backfill_count).is_multiple_of(10) { + if (live_count + backfill_count) % 10 == 0 { info!( "[Stats] Live events: {}, Backfill events: {}", live_count, backfill_count ); } - } else if let Event::Identity(ref identity_event) = *received { + } else if let Event::Identity(ref identity_event) = event { info!( "[IDENTITY] {} -> {} (active: {})", identity_event.did, identity_event.handle, identity_event.is_active ); } - // Event is automatically acked when `received` is dropped here + + if let Err(e) = ack_sender.ack(ack_id).await { + error!("Failed to ack event: {}", e); + break; + } } Err(e) => { eprintln!("Error receiving event: {}", e); diff --git a/tapped/src/channel.rs b/tapped/src/channel.rs index 6a9b020..0f3742f 100644 --- a/tapped/src/channel.rs +++ b/tapped/src/channel.rs @@ -1,10 +1,8 @@ //! WebSocket event channel and receiver. use serde::Serialize; -use tokio::sync::mpsc; use sockudo_ws::{Message, Http1, Config, Stream as WsTransportStream, SplitWriter, SplitReader}; use sockudo_ws::client::WebSocketClient; -use bytes::Bytes; use url::Url; use crate::types::RawEvent; @@ -13,58 +11,47 @@ use crate::{Error, Event, Result}; type WsSink = SplitWriter>; type WsSource = SplitReader>; +/// Opaque identifier for an event to be acknowledged. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub struct AckId(u64); + +/// Sender for event acknowledgments. +pub struct AckSender { + write: WsSink, +} + +impl AckSender { + /// Send an acknowledgment for an event. + pub async fn ack(&mut self, id: AckId) -> Result<()> { + #[derive(Serialize)] + struct AckMessage { + #[serde(rename = "type")] + type_: &'static str, + id: u64, + } + + let msg = AckMessage { type_: "ack", id: id.0 }; + let json = serde_json::to_string(&msg)?; + + self.write.send(Message::text(json)).await.map_err(|e| Error::WebSocket(Box::new(e))) + } +} + /// Receiver for events from a tap WebSocket channel. /// /// Events are received via the [`recv`](EventReceiver::recv) method. -/// Acknowledgments are sent automatically when events are dropped. +/// Acknowledgments must be sent manually using the [`AckSender`] returned by [`TapClient::channel()`](crate::TapClient::channel). /// /// This type does not implement auto-reconnection. If the connection /// closes, `recv()` will return an error and you must create a new /// `EventReceiver` via [`TapClient::channel()`](crate::TapClient::channel). pub struct EventReceiver { - event_rx: mpsc::Receiver, - _ack_tx: mpsc::Sender, -} - -struct EventWithAck { - event: Bytes, - ack_tx: mpsc::Sender, -} - -struct AckGuard { - id: u64, - ack_tx: Option>, -} - -impl Drop for AckGuard { - fn drop(&mut self) { - if let Some(tx) = self.ack_tx.take() { - // Fire and forget - if the channel is closed, we can't ack anyway - let id = self.id; - tokio::spawn(async move { - let _ = tx.send(id).await; - }); - } - } -} - -/// Wrapper around Event that includes the ack trigger. -pub struct ReceivedEvent { - pub event: Event, - _ack_guard: AckGuard, -} - -impl std::ops::Deref for ReceivedEvent { - type Target = Event; - - fn deref(&self) -> &Self::Target { - &self.event - } + read: WsSource, } impl EventReceiver { /// Connect to a tap WebSocket channel. - pub(crate) async fn connect(base_url: &Url, admin_password: Option<&str>) -> Result { + pub(crate) async fn connect(base_url: &Url, admin_password: Option<&str>) -> Result<(Self, AckSender)> { let mut ws_url = base_url.clone(); match ws_url.scheme() { "http" => ws_url.set_scheme("ws").unwrap(), @@ -122,116 +109,51 @@ impl EventReceiver { let (read, write) = ws_stream.split(); - // Buffer size increased to 2048 to prevent TCP window clamping during brief processing spikes - let (event_tx, event_rx) = mpsc::channel(2048); - let (ack_tx, ack_rx) = mpsc::channel(1000); - - let ack_tx_clone = ack_tx.clone(); - tokio::spawn(async move { - Self::writer_task(write, ack_rx).await; - }); - - tokio::spawn(async move { - Self::reader_task(read, event_tx, ack_tx_clone).await; - }); - - Ok(Self { - event_rx, - _ack_tx: ack_tx, - }) + Ok(( + Self { read }, + AckSender { write }, + )) } /// Receive the next event. /// - /// Returns the event wrapped in a [`ReceivedEvent`] that automatically - /// sends an acknowledgment when dropped. + /// Returns the event and an opaque ID that must be passed to [`AckSender::ack`] + /// to acknowledge the event. /// /// # Errors /// /// Returns [`Error::ChannelClosed`] if the WebSocket connection closes. - pub async fn recv(&mut self) -> Result { + pub async fn recv(&mut self) -> Result<(Event, AckId)> { loop { - match self.event_rx.recv().await { - Some(event_with_ack) => { - let json = event_with_ack.event; - let json_str = std::str::from_utf8(&json).expect("must be utf8"); + match self.read.next().await { + Some(Ok(Message::Text(event_bytes))) => { + let event_bytes = bytes::Bytes::from(event_bytes); + let json_str = match std::str::from_utf8(&event_bytes) { + Ok(s) => s, + Err(e) => { + tracing::warn!("Failed to parse event as utf8: {}", e); + continue; + } + }; + let raw = match serde_json::from_str::(json_str) { Ok(raw) => raw, Err(e) => { - tracing::warn!("Failed to parse event: {}", e); + tracing::warn!("Failed to parse event json: {}", e); continue; } }; - if let Some(event) = raw.into_event(json.clone()) { + if let Some(event) = raw.into_event(event_bytes.clone()) { let id = event.id(); - break Ok(ReceivedEvent { - event, - _ack_guard: AckGuard { - id, - ack_tx: Some(event_with_ack.ack_tx), - }, - }); + return Ok((event, AckId(id))); } } - None => break Err(Error::ChannelClosed), + Some(Ok(Message::Close(_))) => return Err(Error::ChannelClosed), + Some(Ok(_)) => continue, // Ping/Pong/Binary + Some(Err(_)) => return Err(Error::ChannelClosed), + None => return Err(Error::ChannelClosed), } } } - - /// Writer task: sends ack messages to the WebSocket. - async fn writer_task(mut write: WsSink, mut ack_rx: mpsc::Receiver) { - #[derive(Serialize)] - struct AckMessage { - #[serde(rename = "type")] - type_: &'static str, - id: u64, - } - - while let Some(id) = ack_rx.recv().await { - let msg = AckMessage { type_: "ack", id }; - let json = match serde_json::to_string(&msg) { - Ok(j) => j, - Err(e) => { - tracing::warn!("Failed to serialize ack: {}", e); - continue; - } - }; - - if let Err(e) = write.send(Message::text(json)).await { - tracing::warn!("Failed to send ack: {}", e); - break; - } - } - } - - /// Reader task: reads events from WebSocket and sends to channel. - async fn reader_task( - mut read: WsSource, - event_tx: mpsc::Sender, - ack_tx: mpsc::Sender, - ) { - while let Some(msg_result) = read.next().await { - match msg_result { - Ok(Message::Text(event)) => { - let event_with_ack = EventWithAck { - event, - ack_tx: ack_tx.clone(), - }; - if event_tx.send(event_with_ack).await.is_err() { - break; - } - } - Ok(Message::Close(_)) => { - break; - } - Ok(_) => { - // Ignore ping/pong/binary - } - Err(_) => { - break; - } - } - } - } -} +} \ No newline at end of file diff --git a/tapped/src/client.rs b/tapped/src/client.rs index f05edf7..40e29c8 100644 --- a/tapped/src/client.rs +++ b/tapped/src/client.rs @@ -337,8 +337,8 @@ impl TapClient { /// Connect to the WebSocket event channel. /// - /// Returns an [`EventReceiver`] for receiving events. Events are - /// automatically acknowledged when dropped. + /// Returns an [`EventReceiver`] for receiving events and an [`AckSender`] for + /// acknowledging them. Events must be manually acknowledged using [`AckSender::ack`]. /// /// # Example /// @@ -347,15 +347,16 @@ impl TapClient { /// use tapped::TapClient; /// /// let client = TapClient::new("http://localhost:2480")?; - /// let mut receiver = client.channel().await?; + /// let (mut receiver, mut ack_sender) = client.channel().await?; /// - /// while let Ok(event) = receiver.recv().await { - /// // Event is automatically acknowledged when dropped + /// while let Ok((event, ack_id)) = receiver.recv().await { + /// // Process event... + /// ack_sender.ack(ack_id).await?; /// } /// # Ok(()) /// # } /// ``` - pub async fn channel(&self) -> Result { + pub async fn channel(&self) -> Result<(EventReceiver, crate::channel::AckSender)> { EventReceiver::connect(&self.base_url, self.admin_password.as_deref()).await } } diff --git a/tapped/src/handle.rs b/tapped/src/handle.rs index 47d7fb1..9522d9a 100644 --- a/tapped/src/handle.rs +++ b/tapped/src/handle.rs @@ -32,9 +32,10 @@ use crate::Result; /// // Use the client methods directly on the handle /// handle.health().await?; /// -/// let mut channel = handle.channel().await?; -/// while let Ok(event) = channel.recv().await { +/// let (mut receiver, mut ack_sender) = handle.channel().await?; +/// while let Ok((event, ack_id)) = receiver.recv().await { /// // Handle event +/// ack_sender.ack(ack_id).await?; /// } /// /// Ok(()) diff --git a/tapped/src/lib.rs b/tapped/src/lib.rs index f2288b0..fa7555f 100644 --- a/tapped/src/lib.rs +++ b/tapped/src/lib.rs @@ -10,7 +10,7 @@ //! //! - Connect to an existing tap instance or spawn one as a subprocess //! - Strongly-typed configuration with builder pattern -//! - Async event streaming with automatic acknowledgment +//! - Async event streaming with manual acknowledgment //! - Full HTTP API coverage for repo management and statistics //! //! ## Example @@ -30,9 +30,10 @@ //! client.add_repos(&["did:plc:example1234567890abc"]).await?; //! //! // Stream events -//! let mut receiver = client.channel().await?; -//! while let Ok(event) = receiver.recv().await { -//! // Event is automatically acknowledged when dropped +//! let (mut receiver, mut ack_sender) = client.channel().await?; +//! while let Ok((event, ack_id)) = receiver.recv().await { +//! // Process event... +//! ack_sender.ack(ack_id).await?; //! } //! //! Ok(()) @@ -47,7 +48,7 @@ mod handle; mod process; mod types; -pub use channel::{EventReceiver, ReceivedEvent}; +pub use channel::{EventReceiver, AckSender, AckId}; pub use client::TapClient; pub use config::{LogLevel, TapConfig, TapConfigBuilder}; pub use error::Error; @@ -59,4 +60,4 @@ pub use types::{ }; /// A specialised Result type for tapped operations. -pub type Result = std::result::Result; +pub type Result = std::result::Result; \ No newline at end of file diff --git a/tapped/tests/integration.rs b/tapped/tests/integration.rs index dd32466..61af20c 100644 --- a/tapped/tests/integration.rs +++ b/tapped/tests/integration.rs @@ -149,7 +149,7 @@ async fn test_channel_connection() { .await .expect("Failed to spawn tap"); - let _channel = handle.channel().await.expect("channel connection failed"); + let (_receiver, _ack_sender) = handle.channel().await.expect("channel connection failed"); } #[tokio::test] -- 2.51.2