Skip to main content

WebSocketClient

Struct WebSocketClient 

Source
pub struct WebSocketClient { /* private fields */ }
Expand description

A WebSocket client with rate limiting, heartbeats, and automatic reconnection in handler mode.

Handler mode owns the reader and writer tasks, buffers sends during reconnection, and replays them against the replacement connection. Stream mode returns the reader to the caller and does not reconnect automatically. See crate::websocket for connection ownership, replay, and epoch guarantees.

Implementations§

Source§

impl WebSocketClient

Source

pub fn reconnect_headers(&self) -> ReconnectHeaders

Returns shared headers used by future automatic reconnect attempts.

Source

pub fn reconnect_handle(&self) -> WebSocketReconnectHandle

Returns a cloneable handle to this client’s reconnect controller.

Source

pub fn request_reconnect(&self) -> bool

Requests that the controller replace the active transport.

Returns true only when this call transitions a handler-mode client from active to reconnecting. Stream-mode, duplicate, disconnecting, and closed requests return false.

Source

pub fn connection_mode(&self) -> ConnectionMode

Returns the current connection mode.

Source

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.

Source

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.

Source

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.

Source

pub fn is_active(&self) -> bool

Returns whether 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.

Source

pub fn is_disconnected(&self) -> bool

Returns whether the controller task has stopped.

Source

pub fn is_reconnecting(&self) -> bool

Returns whether 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.

Source

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.

Source

pub fn is_disconnecting(&self) -> bool

Returns whether the client is disconnecting.

Returns true if the client is in disconnect mode.

Source

pub fn is_closed(&self) -> bool

Returns whether 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.

Source

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.

Source

pub async fn disconnect(&self)

Disconnects the client and waits for the controller task to stop.

If an AuthTracker is registered, this fails pending auth waits.

Source

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.

Source

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:

Source

pub async fn send_pong(&self, data: Vec<u8>) -> Result<(), SendError>

Sends a pong frame back to the server when the connection is active.

The pong is skipped silently if the connection is not active when called: a pong belongs to the connection whose ping caused it, so this method does not wait for reconnection before enqueueing it.

§Errors

Returns an error if:

  • The payload exceeds 125 bytes, the RFC 6455 control-frame limit.
  • The writer channel is broken.
Source

pub async fn send_pong_on_connection( &self, data: Vec<u8>, connection_epoch: u64, ) -> Result<(), SendError>

Sends a pong if the connection that received its ping is still active.

The pong is skipped silently if the connection is inactive or its epoch has changed.

§Errors

Returns an error if:

  • The payload exceeds 125 bytes, the RFC 6455 control-frame limit.
  • The writer channel is broken.
Source

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.

Source

pub async fn send_close_message(&self) -> Result<(), SendError>

Sends a close message to the server.

§Errors

Returns a websocket error if unable to send.

Source

pub fn stream_builder() -> WebSocketClientStreamBuilder

Returns a builder for a websocket client in stream mode.

Calling connect 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 create a new connection.

Use stream mode when you need custom reconnection logic, direct control over message reading, or fine-grained backpressure handling.

default_quota and keyed_quotas limit outgoing messages. state_sink reports transport availability changes.

See WebSocketConfig documentation for comparison with handler mode.

§Errors

Returns an error if the connection cannot be established.

Source

pub fn builder() -> WebSocketClientBuilder

Returns a builder for a websocket client in handler mode.

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, or callback-based message handling.

See WebSocketConfig documentation for comparison with stream mode.

Set rate_limiter to share message quota state across clients. Otherwise, the client creates one from default_quota and keyed_quotas. connection_rate_limiter gates the initial connection and reconnects using connection_rate_keys.

Without initial_connect_retry_policy the builder makes exactly one connection attempt. With one, failures classified as retryable are retried up to its max_attempts; see InitialConnectRetryPolicy for which failures return before that bound is reached.

cancellation_token aborts the initial connection only, and is observed during the connection rate-limit wait, the dial itself, and each backoff delay. It has no effect once this function returns a client: use WebSocketClient::disconnect to stop an established one, whose reconnect loop the token does not govern.

The message handler is required:

use nautilus_network::websocket::{WebSocketClient, WebSocketConfig};

let config: WebSocketConfig = unimplemented!();
let _ = WebSocketClient::builder().config(config).connect();
§Errors

Returns an error if:

  • The configuration is invalid or the connection cannot be established.
  • A shared rate limiter is combined with quota configuration.
  • The connection rate limiter and its keys are not configured together.
Source

pub fn epoch_builder() -> WebSocketClientEpochBuilder

Returns a builder for a handler-mode client whose 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. Rate-limit, state, initial-connect retry, and cancellation options match Self::builder. Set either ping_handler or epoch_ping_handler when custom ping handling is required.

The epoch handler is required:

use nautilus_network::websocket::{WebSocketClient, WebSocketConfig};

let config: WebSocketConfig = unimplemented!();
let _ = WebSocketClient::epoch_builder().config(config).connect();
§Errors

Returns an error if:

  • The configuration is invalid or the connection cannot be established.
  • A shared rate limiter is combined with quota configuration.
  • The connection rate limiter and its keys are not configured together.
  • Both ping_handler and epoch_ping_handler are configured.

Trait Implementations§

Source§

impl Debug for WebSocketClient

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl Drop for WebSocketClient

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

§

impl<T> Instrument for T

§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided [Span], returning an Instrumented wrapper. Read more
§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

§

impl<T> PolicyExt for T
where T: ?Sized,

§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns [Action::Follow] only if self and other return Action::Follow. Read more
§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns [Action::Follow] if either self or other returns Action::Follow. Read more
§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
§

impl<T> Ungil for T
where T: Send,

§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V

§

impl<T> WithSubscriber for T

§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a [WithDispatch] wrapper. Read more
§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a [WithDispatch] wrapper. Read more