pub struct WebSocketClient { /* private fields */ }Expand description
WebSocket client with automatic reconnection.
Handles connection state, callbacks, and rate limiting. See module docs for architecture details.
Implementations§
Source§impl WebSocketClient
impl WebSocketClient
Sourcepub async fn connect_stream(
config: WebSocketConfig,
keyed_quotas: Vec<(String, Quota)>,
default_quota: Option<Quota>,
post_reconnect: Option<Arc<dyn Fn() + Send + Sync>>,
) -> Result<(MessageReader, Self), TransportError>
pub async fn connect_stream( config: WebSocketConfig, keyed_quotas: Vec<(String, Quota)>, default_quota: Option<Quota>, post_reconnect: Option<Arc<dyn Fn() + Send + Sync>>, ) -> Result<(MessageReader, Self), TransportError>
Creates a websocket client in stream mode that returns a MessageReader.
Returns a stream that the caller owns and reads from directly. Automatic reconnection
is disabled because the reader cannot be replaced internally. On disconnection, the
client transitions to CLOSED state and the caller must manually reconnect by calling
connect_stream again.
Use stream mode when you need custom reconnection logic, direct control over message reading, or fine-grained backpressure handling.
See WebSocketConfig documentation for comparison with handler mode.
§Errors
Returns an error if the connection cannot be established.
Sourcepub async fn connect(
config: WebSocketConfig,
message_handler: Option<MessageHandler>,
ping_handler: Option<PingHandler>,
post_reconnection: Option<Arc<dyn Fn() + Send + Sync>>,
keyed_quotas: Vec<(String, Quota)>,
default_quota: Option<Quota>,
) -> Result<Self, TransportError>
pub async fn connect( config: WebSocketConfig, message_handler: Option<MessageHandler>, ping_handler: Option<PingHandler>, post_reconnection: Option<Arc<dyn Fn() + Send + Sync>>, keyed_quotas: Vec<(String, Quota)>, default_quota: Option<Quota>, ) -> Result<Self, TransportError>
Creates a websocket client in handler mode with automatic reconnection.
The handler is called for each incoming message on an internal task. Automatic reconnection is enabled with exponential backoff. On disconnection, the client automatically attempts to reconnect and replaces the internal reader (the handler continues working seamlessly).
Use handler mode for simplified connection management, automatic reconnection, Python bindings, or callback-based message handling.
See WebSocketConfig documentation for comparison with stream mode.
§Errors
Returns an error if:
- The connection cannot be established.
message_handlerisNone(useconnect_streaminstead).
Sourcepub async fn connect_with_rate_limiter(
config: WebSocketConfig,
message_handler: Option<MessageHandler>,
ping_handler: Option<PingHandler>,
post_reconnection: Option<Arc<dyn Fn() + Send + Sync>>,
rate_limiter: Arc<RateLimiter<Ustr, MonotonicClock>>,
) -> Result<Self, TransportError>
pub async fn connect_with_rate_limiter( config: WebSocketConfig, message_handler: Option<MessageHandler>, ping_handler: Option<PingHandler>, post_reconnection: Option<Arc<dyn Fn() + Send + Sync>>, rate_limiter: Arc<RateLimiter<Ustr, MonotonicClock>>, ) -> Result<Self, TransportError>
Creates a websocket client in handler mode sharing an externally-owned rate limiter.
Use this constructor to share a single RateLimiter across multiple
WebSocketClient instances (for example, the WebSocket clients owned
by an exchange adapter’s data and execution clients). All quota state
lives inside the limiter, so passing the same Arc produces a single
shared bucket - the only way to honour a venue’s per-IP / per-account
WS message cap when more than one connection is opened in-process.
Behavior otherwise matches Self::connect.
§Errors
Returns an error if:
- The connection cannot be established.
message_handlerisNone(useconnect_streaminstead).
Sourcepub async fn connect_with_rate_limiter_and_epoch_handler(
config: WebSocketConfig,
epoch_handler: EpochMessageHandler,
ping_handler: Option<PingHandler>,
post_reconnection: Option<Arc<dyn Fn() + Send + Sync>>,
rate_limiter: Arc<RateLimiter<Ustr, MonotonicClock>>,
) -> Result<Self, TransportError>
pub async fn connect_with_rate_limiter_and_epoch_handler( config: WebSocketConfig, epoch_handler: EpochMessageHandler, ping_handler: Option<PingHandler>, post_reconnection: Option<Arc<dyn Fn() + Send + Sync>>, rate_limiter: Arc<RateLimiter<Ustr, MonotonicClock>>, ) -> Result<Self, TransportError>
Creates a handler-mode client whose incoming messages carry connection ownership.
The initial connection has epoch 0. Each replacement connection increments the epoch,
and both its incoming messages and RECONNECTED notification carry that new value. Use
Self::send_text_on_connection to bind an outgoing message to one of those epochs.
§Errors
Returns an error if the connection cannot be established.
Sourcepub fn reconnect_headers(&self) -> ReconnectHeaders
pub fn reconnect_headers(&self) -> ReconnectHeaders
Returns shared headers used by future automatic reconnect attempts.
Sourcepub fn connection_mode(&self) -> ConnectionMode
pub fn connection_mode(&self) -> ConnectionMode
Returns the current connection mode.
Sourcepub fn connection_epoch(&self) -> u64
pub fn connection_epoch(&self) -> u64
Returns the ownership epoch of the current connection.
A connection keeps one epoch for its lifetime. The value increments when the writer swaps to a replacement connection.
Sourcepub fn connection_mode_atomic(&self) -> Arc<AtomicU8> ⓘ
pub fn connection_mode_atomic(&self) -> Arc<AtomicU8> ⓘ
Returns a clone of the connection mode atomic for external state tracking.
This allows adapter clients to track connection state across reconnections without message-passing delays.
Sourcepub fn connection_epoch_atomic(&self) -> Arc<AtomicU64> ⓘ
pub fn connection_epoch_atomic(&self) -> Arc<AtomicU64> ⓘ
Returns shared read access to the current connection epoch.
Callers must treat the returned atomic as read-only. The WebSocket writer task exclusively advances it when installing a replacement connection.
Sourcepub fn is_active(&self) -> bool
pub fn is_active(&self) -> bool
Check if the client connection is active.
Returns true if the client is connected and has not been signalled to disconnect.
The client will automatically retry connection based on its configuration.
Sourcepub fn is_disconnected(&self) -> bool
pub fn is_disconnected(&self) -> bool
Check if the client is disconnected.
Sourcepub fn is_reconnecting(&self) -> bool
pub fn is_reconnecting(&self) -> bool
Check if the client is reconnecting.
Returns true if the client lost connection and is attempting to reestablish it.
The client will automatically retry connection based on its configuration.
Sourcepub fn set_auth_tracker(
&self,
tracker: AuthTracker,
reconnect_buffer_waits_for_auth: bool,
)
pub fn set_auth_tracker( &self, tracker: AuthTracker, reconnect_buffer_waits_for_auth: bool, )
Registers an AuthTracker with the client.
When the controller detects a dead connection and transitions to
Reconnect, it calls invalidate() on the tracker so that any
pending authenticated sends see the state change immediately. Terminal
transitions fail the tracker so pending auth waits can terminate.
Set reconnect_buffer_waits_for_auth for clients that must not replay
buffered messages until the next session authenticates.
Call this once after construction, before any authenticated sends.
Sourcepub fn is_disconnecting(&self) -> bool
pub fn is_disconnecting(&self) -> bool
Check if the client is disconnecting.
Returns true if the client is in disconnect mode.
Sourcepub fn is_closed(&self) -> bool
pub fn is_closed(&self) -> bool
Check if the client is closed.
Returns true if the client has been explicitly disconnected or reached
maximum reconnection attempts. In this state, the client cannot be reused
and a new client must be created for further connections.
Sourcepub fn notify_closed(&self)
pub fn notify_closed(&self)
Signals that the caller’s reader has observed EOF or a fatal error.
In stream mode the controller has no visibility into the caller-owned reader.
Call this method when reader.next().await returns None or an unrecoverable
error so the controller transitions to Closed and dependent tasks shut down.
For peer-initiated close frames (Message::Close), use disconnect
instead so the writer can send the close reply before shutting down.
If an AuthTracker is registered, this fails pending auth waits.
This is a no-op if the connection is already closed or disconnecting.
Sourcepub async fn disconnect(&self)
pub async fn disconnect(&self)
Set disconnect mode to true.
Controller task will periodically check the disconnect mode and shutdown the client if it is alive
If an AuthTracker is registered, this fails pending auth waits.
Sourcepub async fn send_text(
&self,
data: String,
keys: Option<&[Ustr]>,
) -> Result<(), SendError>
pub async fn send_text( &self, data: String, keys: Option<&[Ustr]>, ) -> Result<(), SendError>
Sends the given text data to the server.
Returns Ok(()) when the message is enqueued to the writer channel. This does NOT
guarantee delivery: if a disconnect occurs concurrently, the writer task may drop the
message. During reconnection, messages are buffered and replayed on the new connection.
§Errors
Returns a websocket error if unable to send.
Sourcepub async fn send_text_on_connection(
&self,
data: String,
keys: Option<&[Ustr]>,
connection_epoch: u64,
) -> Result<(), SendError>
pub async fn send_text_on_connection( &self, data: String, keys: Option<&[Ustr]>, connection_epoch: u64, ) -> Result<(), SendError>
Sends text once if the active connection matches connection_epoch.
The writer rejects the send if reconnection changes ownership before the write. The message is never replayed on another connection.
§Errors
Returns:
SendError::Closedif the client closes.SendError::Timeoutif the active wait times out before the write starts.SendError::ConnectionChangedif the expected connection no longer owns the writer.SendError::BrokenPipeif the command or transport write fails.SendError::WriteTimeoutif the write starts but does not complete within the write deadline. Delivery is undetermined in that case: the message is not replayed, and it must not be resent blindly.
Sourcepub async fn send_bytes(
&self,
data: Vec<u8>,
keys: Option<&[Ustr]>,
) -> Result<(), SendError>
pub async fn send_bytes( &self, data: Vec<u8>, keys: Option<&[Ustr]>, ) -> Result<(), SendError>
Sends the given bytes data to the server.
Returns Ok(()) when the message is enqueued to the writer channel. This does NOT
guarantee delivery: if a disconnect occurs concurrently, the writer task may drop the
message. During reconnection, messages are buffered and replayed on the new connection.
§Errors
Returns a websocket error if unable to send.