From d2a29613ff770c23494fe45ebc413f575ab17885 Mon Sep 17 00:00:00 2001 From: LucaCappelletti94 Date: Mon, 5 Oct 2026 14:06:17 +0200 Subject: [PATCH] Announce a pump fault and release the handles waiting on it --- crates/connetto-client/src/lib.rs | 10 + crates/connetto-client/src/live.rs | 97 +++++++-- crates/connetto-client/tests/it/main.rs | 3 + crates/connetto-client/tests/it/pump_fault.rs | 205 ++++++++++++++++++ docs/architecture/13-client-connection.md | 2 + examples/dioxus-desktop-demo/src/main.rs | 1 + examples/dioxus-web-demo/src/main.rs | 1 + examples/yew-web-demo/src/main.rs | 1 + 8 files changed, 298 insertions(+), 22 deletions(-) create mode 100644 crates/connetto-client/tests/it/pump_fault.rs diff --git a/crates/connetto-client/src/lib.rs b/crates/connetto-client/src/lib.rs index 5e32979d..57b2a1f3 100644 --- a/crates/connetto-client/src/lib.rs +++ b/crates/connetto-client/src/lib.rs @@ -205,6 +205,9 @@ pub enum ClientError { /// meets this keeps working and retries when a transport arrives. #[error("not connected: this operation needs a server")] NotConnected, + /// The client's pump stopped on a fault no reconnect can cure, named here. + #[error("the client stopped: {0}")] + Stopped(String), /// Acquiring or refreshing the access token failed. #[error("authentication error: {0}")] Auth(String), @@ -1066,6 +1069,13 @@ pub enum ClientEvent { /// Why the server closed the session. reason: FatalErrorReason, }, + /// The pump stopped on a fault no reconnect can cure, and said which. + /// [`Closed`](Self::Closed) follows, and every live handle's `changed` + /// returns [`ClientError::Stopped`] from then on. + Stopped { + /// The fault, as text. + detail: String, + }, /// The connection closed. Closed, /// The credential was rejected and refresh could not recover it, so the diff --git a/crates/connetto-client/src/live.rs b/crates/connetto-client/src/live.rs index 3117a1d8..d10eb6b1 100644 --- a/crates/connetto-client/src/live.rs +++ b/crates/connetto-client/src/live.rs @@ -25,7 +25,7 @@ use core::task::Poll; use core::time::Duration; use std::collections::{HashMap, HashSet}; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; -use std::sync::{Arc, Mutex as StdMutex, RwLock, Weak}; +use std::sync::{Arc, Mutex as StdMutex, OnceLock, RwLock, Weak}; use connetto_core::messages::{BindValue, SubscriptionSpec}; use connetto_core::quote_ident; @@ -654,10 +654,12 @@ fn run_probe(conn: &mut SqliteConnection, probe: &str) -> Result>, wake: Arc, + stopped: OnceLock, } /// What [`LiveQuery`] and [`LiveValue`] both are underneath: a subscription @@ -686,10 +688,15 @@ impl LiveHandleCore { } async fn changed(&mut self) -> Result<(), ClientError> { - self.changed - .changed() - .await - .map_err(|_| ClientError::Transport("live query driver stopped".to_owned())) + if let Some(detail) = self.reaper.stopped.get() { + return Err(ClientError::Stopped(detail.clone())); + } + self.changed.changed().await.map_err(|_| { + self.reaper.stopped.get().map_or_else( + || ClientError::Transport("live query driver stopped".to_owned()), + |detail| ClientError::Stopped(detail.clone()), + ) + }) } } @@ -1389,6 +1396,7 @@ where reaper: Arc::new(Reaper { pending: StdMutex::new(Vec::new()), wake: Arc::clone(&wake), + stopped: OnceLock::new(), }), events, next_live: AtomicU64::new(1), @@ -2727,7 +2735,7 @@ where /// A transport drop ends the pump, unless a reconnect driver is present, in /// which case the pump recovers: backoff, fresh transport, session resume, /// re-declared subscriptions. Local faults (session, apply, protocol) stay -/// terminal either way. +/// terminal either way and are announced as [`ClientEvent::Stopped`]. async fn pump( shared: Arc>, alive: Weak, @@ -2783,7 +2791,15 @@ async fn pump( needs_recovery = true; continue; } - PumpFlow::Exit => return, + PumpFlow::Exit => { + end_unresumed(&mut state, &shared).await; + return; + } + PumpFlow::Fail(err) => { + drop(state); + pump_fail(&shared, err).await; + return; + } } // One cancellable pump step. A wake interrupts the idle wait so lock @@ -2824,11 +2840,12 @@ async fn pump( continue; } PumpFlow::Exit => { - // The server ended the session with no driver to resume it, so - // the writes made meanwhile are queued for the next run. - if let Err(err) = state.conn.flush().await { - tracing::warn!(error = %err, "queueing the last writes at exit failed"); - } + end_unresumed(&mut state, &shared).await; + return; + } + PumpFlow::Fail(err) => { + drop(state); + pump_fail(&shared, err).await; return; } } @@ -2856,8 +2873,49 @@ enum PumpFlow { Proceed, /// The transport is gone. Recover at the top of the next iteration. Recover, - /// The pump has ended. + /// The session ended with no driver to resume it. Exit, + /// A fault no reconnect can cure ended the pump. + Fail(ClientError), +} + +/// Queue the writes made meanwhile for the next run, then end the pump, when +/// the session ended with no driver to resume it. +async fn end_unresumed(state: &mut State, shared: &Shared) +where + T: Transport + MaybeSend + 'static, + T::Error: core::fmt::Display, +{ + if let Err(err) = state.conn.flush().await { + tracing::warn!(error = %err, "queueing the last writes at exit failed"); + } + pump_finished(shared); +} + +/// End the pump on a fault, queueing what can be queued, closing the +/// transport, releasing every live handle with the fault, then announcing it +/// and `Closed`. +async fn pump_fail(shared: &Arc>, err: ClientError) +where + T: Transport + MaybeSend + 'static, + T::Error: core::fmt::Display, +{ + let detail = err.to_string(); + tracing::warn!(error = %detail, "the client pump stopped on a fault"); + // Set before the senders drop, so a woken handle reads the fault. + let _ = shared.reaper.stopped.set(detail.clone()); + let mut state = shared.state.lock().await; + if let Err(err) = state.conn.flush().await { + tracing::warn!(error = %err, "queueing the last writes at the fault failed"); + } + let _ = state.conn.close().await; + state.registry.clear(); + state.values.clear(); + state.computed.clear(); + drop(state); + let _ = shared.events.send(ClientEvent::Stopped { detail }); + let _ = shared.events.send(ClientEvent::Closed); + pump_finished(shared); } /// Queue the last writes, close the transport, announce `Closed`, and end @@ -2930,8 +2988,7 @@ where if is_disconnect(&err) { return PumpFlow::Recover; } - pump_finished(shared); - return PumpFlow::Exit; + return PumpFlow::Fail(err); } // Auto-submit local writes committed since the last step. With no socket // this only queues them durably, which must not wait for a server. @@ -2939,8 +2996,7 @@ where if is_disconnect(&err) { return PumpFlow::Recover; } - pump_finished(shared); - return PumpFlow::Exit; + return PumpFlow::Fail(err); } // No socket and a driver to find one, so go and find it. With no driver // there is nothing to recover to, so the pump falls through and parks, @@ -2973,7 +3029,6 @@ where PumpFlow::Recover } else { let _ = shared.events.send(ClientEvent::Closed); - pump_finished(shared); PumpFlow::Exit } } @@ -2982,7 +3037,6 @@ where PumpFlow::Recover } else { let _ = shared.events.send(ClientEvent::Closed); - pump_finished(shared); PumpFlow::Exit } } @@ -3008,8 +3062,7 @@ where if is_disconnect(&err) { PumpFlow::Recover } else { - pump_finished(shared); - PumpFlow::Exit + PumpFlow::Fail(err) } } } diff --git a/crates/connetto-client/tests/it/main.rs b/crates/connetto-client/tests/it/main.rs index 4e6dcec6..5054e254 100644 --- a/crates/connetto-client/tests/it/main.rs +++ b/crates/connetto-client/tests/it/main.rs @@ -59,6 +59,9 @@ mod never_synced; mod offline_start; +#[cfg(feature = "native-auth")] +mod pump_fault; + mod reconnect_live; mod residual; diff --git a/crates/connetto-client/tests/it/pump_fault.rs b/crates/connetto-client/tests/it/pump_fault.rs new file mode 100644 index 00000000..63c34061 --- /dev/null +++ b/crates/connetto-client/tests/it/pump_fault.rs @@ -0,0 +1,205 @@ +//! A pump that stops on a fault it cannot recover from says why and releases +//! everything waiting on it. + +use connetto_client::{ClientBuilder, ClientError, ClientEvent, ConnettoClient, DataDir}; +use connetto_core::messages::{BulkMessage, ControlMessage, HandshakeAck, LivePatch}; +use connetto_core::traits::{IncomingFrame, Transport}; +use diesel::prelude::*; +use std::collections::VecDeque; +use std::future::ready; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; +use tempfile::tempdir; + +const DDL: &str = "CREATE TABLE items (id INTEGER PRIMARY KEY, label TEXT)"; +const BOUND: Duration = Duration::from_secs(5); + +diesel::table! { + /// Synced test table. + items (id) { + /// Item identifier, the primary key + id -> Integer, + /// Optional item label + label -> Nullable, + } +} + +#[derive(Debug, Clone, PartialEq, Queryable)] +struct Item { + id: i32, + label: Option, +} + +/// A transport handing out queued frames, idling when none are queued, and +/// recording the subscriptions it was sent and whether it was closed. +#[derive(Clone, Default)] +struct Script { + frames: Arc>>, + subscribed: Arc>>, + closed: Arc, +} + +impl Script { + fn say(&self, frame: IncomingFrame) { + self.frames.lock().expect("script lock").push_back(frame); + } + + fn last_subscribed(&self) -> String { + self.subscribed + .lock() + .expect("script lock") + .last() + .cloned() + .expect("the client subscribed") + } +} + +impl Transport for Script { + type Error = std::convert::Infallible; + + fn send_control( + &mut self, + message: ControlMessage, + ) -> impl Future> { + if let ControlMessage::Subscribe(subscribe) = message { + self.subscribed + .lock() + .expect("script lock") + .push(subscribe.sub_id); + } + ready(Ok(())) + } + + fn send_bulk( + &mut self, + _message: BulkMessage, + ) -> impl Future> { + ready(Ok(())) + } + + fn recv(&mut self) -> impl Future, Self::Error>> { + let frames = Arc::clone(&self.frames); + async move { + loop { + if let Some(frame) = frames.lock().expect("script lock").pop_front() { + return Ok(Some(frame)); + } + tokio::task::yield_now().await; + } + } + } + + fn close(&mut self) -> impl Future> { + self.closed.store(true, Ordering::Relaxed); + ready(Ok(())) + } +} + +fn ack() -> IncomingFrame { + IncomingFrame::Control(ControlMessage::HandshakeAck(HandshakeAck { + connection_id: "script".to_owned(), + session_token: "script".to_owned(), + resume_token: "script".to_owned(), + current_cursor: connetto_core::Cursor::from(Vec::new()), + schema_version: None, + initial_credits: 64, + last_applied_seq: None, + })) +} + +async fn client() -> (ConnettoClient