Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions crates/connetto-client/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down Expand Up @@ -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
Expand Down
97 changes: 75 additions & 22 deletions crates/connetto-client/src/live.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -654,10 +654,12 @@ fn run_probe(conn: &mut SqliteConnection, probe: &str) -> Result<String, ClientE
}

/// Subscription ids and the wake signal shared with every [`LiveQuery`], so a
/// synchronous `Drop` can queue its unsubscribe for the async pump.
/// synchronous `Drop` can queue its unsubscribe for the async pump, and the
/// fault that stopped the pump, if one did.
struct Reaper {
pending: StdMutex<Vec<String>>,
wake: Arc<Notify>,
stopped: OnceLock<String>,
}

/// What [`LiveQuery`] and [`LiveValue`] both are underneath: a subscription
Expand Down Expand Up @@ -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()),
)
})
}
}

Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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<T, F, S>(
shared: Arc<Shared<T>>,
alive: Weak<ClientToken>,
Expand Down Expand Up @@ -2783,7 +2791,15 @@ async fn pump<T, F, S>(
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
Expand Down Expand Up @@ -2824,11 +2840,12 @@ async fn pump<T, F, S>(
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;
}
}
Expand Down Expand Up @@ -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<T>(state: &mut State<T>, shared: &Shared<T>)
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<T>(shared: &Arc<Shared<T>>, 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
Expand Down Expand Up @@ -2930,17 +2988,15 @@ 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.
if let Err(err) = state.conn.flush().await {
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,
Expand Down Expand Up @@ -2973,7 +3029,6 @@ where
PumpFlow::Recover
} else {
let _ = shared.events.send(ClientEvent::Closed);
pump_finished(shared);
PumpFlow::Exit
}
}
Expand All @@ -2982,7 +3037,6 @@ where
PumpFlow::Recover
} else {
let _ = shared.events.send(ClientEvent::Closed);
pump_finished(shared);
PumpFlow::Exit
}
}
Expand All @@ -3008,8 +3062,7 @@ where
if is_disconnect(&err) {
PumpFlow::Recover
} else {
pump_finished(shared);
PumpFlow::Exit
PumpFlow::Fail(err)
}
}
}
Expand Down
3 changes: 3 additions & 0 deletions crates/connetto-client/tests/it/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,9 @@ mod never_synced;

mod offline_start;

#[cfg(feature = "native-auth")]
mod pump_fault;

mod reconnect_live;

mod residual;
Expand Down
Loading
Loading