Skip to main content

atomic_websocket/
lib.rs

1//! # atomic_websocket
2//!
3//! `atomic_websocket` is a high-level WebSocket client and server implementation for Rust built on top of tokio-tungstenite.
4//! It provides resilient WebSocket connections with the following features:
5//!
6//! - Automatic connection recovery and ping/pong handling
7//! - Local network server auto-discovery
8//! - Connection status monitoring and events
9//! - Serialization/deserialization support
10//!
11//! ## Basic Usage
12//!
13//! ### Client Example
14//!
15//! ```rust,ignore
16//! use atomic_websocket::{
17//!     AtomicWebsocket,
18//!     connection_store::{ConnectionStore, NativeDbConnectionStore},
19//!     server_sender::{ClientOptions, SenderStatus},
20//!     schema::ServerConnectInfo,
21//! };
22//! use std::sync::Arc;
23//!
24//! #[tokio::main]
25//! async fn main() -> Result<(), Box<dyn std::error::Error>> {
26//!     // Configure client options
27//!     let mut client_options = ClientOptions::default();
28//!     client_options.retry_seconds = 2;
29//!     client_options.use_keep_ip = true;
30//!
31//!     // Initialize your own DB (implementation details omitted), then wrap
32//!     // it once in a ConnectionStore — the library never touches your DB
33//!     // handle directly, only through this trait.
34//!     let db = initialize_database().await?;
35//!     let connection_store: Arc<dyn ConnectionStore> =
36//!         Arc::new(NativeDbConnectionStore::new(db.clone()));
37//!     let server_sender = initialize_server_sender().await?;
38//!
39//!     // Create client
40//!     let atomic_client = AtomicWebsocket::get_internal_client_with_server_sender(
41//!         connection_store.clone(),
42//!         client_options,
43//!         server_sender.clone(),
44//!     ).await;
45//!
46//!     // Connect to server
47//!     let result = atomic_client
48//!         .get_internal_connect(
49//!             Some(ServerConnectInfo {
50//!                 server_ip: "",
51//!                 port: "9000",
52//!             }),
53//!             connection_store.clone(),
54//!         )
55//!         .await;
56//!
57//!     Ok(())
58//! }
59//! ```
60//!
61//! ### Server Example
62//!
63//! ```rust,ignore
64//! use atomic_websocket::{
65//!     AtomicWebsocket,
66//!     client_sender::ServerOptions,
67//! };
68//!
69//! #[tokio::main]
70//! async fn main() -> Result<(), Box<dyn std::error::Error>> {
71//!     // Configure server options
72//!     let options = ServerOptions::default();
73//!
74//!     // Initialize client senders for managing connections
75//!     let client_senders = initialize_client_senders().await?;
76//!
77//!     // Create and start server
78//!     let address = "0.0.0.0:9000";
79//!     let atomic_server = AtomicWebsocket::get_internal_server_with_client_senders(
80//!         address.to_string(),
81//!         options,
82//!         client_senders.clone(),
83//!     ).await?;
84//!
85//!     // Set up message handler
86//!     let handle_message_receiver = atomic_server.get_handle_message_receiver().await;
87//!     tokio::spawn(handle_messages(handle_message_receiver));
88//!
89//!     Ok(())
90//! }
91//! ```
92
93use std::sync::Arc;
94
95use helpers::{
96    connection_store::ConnectionStore,
97    internal_client::{AtomicClient, ClientOptions},
98    internal_server::{AtomicServer, ServerOptions},
99};
100#[cfg(feature = "native-db")]
101use native_db::{native_db, ToKey};
102#[cfg(feature = "native-db")]
103use native_model::{native_model, Model};
104use serde::{Deserialize, Serialize};
105
106/// Module that re-exports various external dependencies.
107pub mod external {
108    pub use async_trait;
109    #[cfg(feature = "bebop")]
110    pub use bebop;
111    pub use futures_util;
112    pub use nanoid;
113    #[cfg(feature = "native-db")]
114    pub use native_db;
115    #[cfg(feature = "native-db")]
116    pub use native_model;
117    #[cfg(feature = "rustls")]
118    pub use rustls;
119    pub use tokio;
120    pub use tokio_tungstenite;
121}
122
123/// Module containing message schema definitions.
124///
125/// This module defines the message formats exchanged between clients and servers.
126#[cfg(feature = "bebop")]
127pub mod schema {
128    pub use crate::generated::schema::*;
129}
130
131/// Module providing LAN server discovery (used by Android clients to
132/// locate the POS before the first connect).
133pub mod scan_manager {
134    pub use crate::helpers::scan_manager::*;
135}
136
137/// Module providing functionality for managing client connections on the server side.
138pub mod client_sender {
139    pub use crate::helpers::client_sender::*;
140    pub use crate::helpers::internal_server::{handle_upgraded_connection, ServerOptions};
141}
142
143/// Module providing functionality for managing server connections on the client side.
144pub mod server_sender {
145    pub use crate::helpers::internal_client::{
146        get_internal_connect, get_ip_address, ClientOptions,
147    };
148    pub use crate::helpers::server_sender::*;
149}
150
151/// Module providing the pluggable connection-identity persistence trait
152/// (client ID + last-known server connect info) used by `ServerSender`/
153/// `AtomicClient`, plus the library's default `Settings`-table-backed
154/// implementation.
155pub mod connection_store {
156    pub use crate::helpers::connection_store::{ConnectionStore, NativeDbConnectionStore};
157}
158
159/// Module providing common utility functions for WebSocket communication.
160pub mod common {
161    #[cfg(feature = "bebop")]
162    pub use crate::helpers::common::make_response_message;
163    pub use crate::helpers::common::{get_setting_by_key, make_atomic_message, set_setting};
164    pub use crate::helpers::get_internal_websocket::get_id;
165}
166
167/// Module containing common type definitions used throughout the library.
168pub mod types {
169    pub use crate::helpers::types::*;
170}
171
172/// Module providing builder patterns for configuration.
173pub mod builder {
174    pub use crate::helpers::builder::{ClientOptionsBuilder, ServerOptionsBuilder};
175}
176
177/// Module providing metrics and observability.
178pub mod metrics {
179    pub use crate::helpers::metrics::{Metrics, MetricsSnapshot};
180}
181
182/// Module providing middleware/interceptor pattern for WebSocket message handling.
183pub mod middleware {
184    pub use crate::helpers::middleware::{MessageMiddleware, MiddlewareResult};
185}
186
187use server_sender::{ServerSender, ServerSenderTrait};
188use tokio::sync::RwLock;
189use tokio_util::sync::CancellationToken;
190use types::{RwClientSenders, RwServerSender};
191
192#[cfg(feature = "bebop")]
193mod generated;
194mod helpers;
195
196/// Database model for storing client settings and state.
197#[cfg(feature = "native-db")]
198#[derive(Serialize, Deserialize, Debug, Clone)]
199#[native_model(id = 1004, version = 1)]
200#[native_db]
201pub struct Settings {
202    #[primary_key]
203    pub key: String,
204    pub value: Vec<u8>,
205}
206
207/// Settings struct for when native-db feature is disabled.
208/// Stores key-value pairs in memory only.
209#[cfg(not(feature = "native-db"))]
210#[derive(Serialize, Deserialize, Debug, Clone)]
211pub struct Settings {
212    pub key: String,
213    pub value: Vec<u8>,
214}
215
216/// Primary entry point for creating WebSocket clients and servers.
217///
218/// This struct provides static methods for creating client and server
219/// instances for both internal and external WebSocket connections.
220pub struct AtomicWebsocket {}
221
222/// Internal enum to distinguish WebSocket connection types.
223#[derive(Debug, Clone)]
224pub enum AtomicWebsocketType {
225    /// Connection to a server on the local network
226    Internal,
227    /// Connection to an external internet server
228    External,
229}
230
231impl AtomicWebsocket {
232    /// Creates a basic client instance for internal network use.
233    ///
234    /// # Arguments
235    ///
236    /// * `connection_store` - Persistence for connection-identity state (client ID, server connect info)
237    /// * `options` - Client connection options (auto-reconnect, ping intervals, etc.)
238    ///
239    /// # Returns
240    ///
241    /// A newly created `AtomicClient` instance
242    ///
243    /// # Examples
244    ///
245    /// ```rust,ignore
246    /// let client_options = ClientOptions::default();
247    /// let client = AtomicWebsocket::get_internal_client(connection_store.clone(), client_options).await;
248    /// ```
249    pub async fn get_internal_client(
250        connection_store: Arc<dyn ConnectionStore>,
251        mut options: ClientOptions,
252    ) -> AtomicClient {
253        options.atomic_websocket_type = AtomicWebsocketType::Internal;
254        get_client(connection_store, options, None).await
255    }
256
257    /// Creates a client instance for internal network use with an existing ServerSender.
258    ///
259    /// # Arguments
260    ///
261    /// * `connection_store` - Persistence for connection-identity state
262    /// * `options` - Client connection options
263    /// * `server_sender` - Existing ServerSender instance
264    ///
265    /// # Returns
266    ///
267    /// A newly created `AtomicClient` instance
268    pub async fn get_internal_client_with_server_sender(
269        connection_store: Arc<dyn ConnectionStore>,
270        mut options: ClientOptions,
271        server_sender: RwServerSender,
272    ) -> AtomicClient {
273        options.atomic_websocket_type = AtomicWebsocketType::Internal;
274        get_client(connection_store, options, Some(server_sender)).await
275    }
276
277    /// Creates a client instance for external servers.
278    ///
279    /// # Arguments
280    ///
281    /// * `connection_store` - Persistence for connection-identity state
282    /// * `options` - Client connection options (including URL, TLS settings, etc.)
283    ///
284    /// # Returns
285    ///
286    /// A newly created `AtomicClient` instance
287    pub async fn get_outer_client(
288        connection_store: Arc<dyn ConnectionStore>,
289        mut options: ClientOptions,
290    ) -> AtomicClient {
291        options.atomic_websocket_type = AtomicWebsocketType::External;
292        get_client(connection_store, options, None).await
293    }
294
295    /// Creates a client instance for external servers with an existing ServerSender.
296    ///
297    /// # Arguments
298    ///
299    /// * `connection_store` - Persistence for connection-identity state
300    /// * `options` - Client connection options
301    /// * `server_sender` - Existing ServerSender instance
302    ///
303    /// # Returns
304    ///
305    /// A newly created `AtomicClient` instance
306    pub async fn get_outer_client_with_server_sender(
307        connection_store: Arc<dyn ConnectionStore>,
308        mut options: ClientOptions,
309        server_sender: RwServerSender,
310    ) -> AtomicClient {
311        options.atomic_websocket_type = AtomicWebsocketType::External;
312        get_client(connection_store, options, Some(server_sender)).await
313    }
314
315    /// Creates a server instance for internal network use.
316    ///
317    /// # Arguments
318    ///
319    /// * `addr` - Address to bind the server to (e.g., "0.0.0.0:9000")
320    /// * `option` - Server configuration options
321    ///
322    /// # Returns
323    ///
324    /// A `Result` containing the new `AtomicServer` instance, or an IO error if binding fails
325    ///
326    /// # Examples
327    ///
328    /// ```rust,ignore
329    /// let server_options = ServerOptions::default();
330    /// let server = AtomicWebsocket::get_internal_server("127.0.0.1:9000".to_string(), server_options).await?;
331    /// ```
332    pub async fn get_internal_server(
333        addr: String,
334        option: ServerOptions,
335    ) -> std::io::Result<AtomicServer> {
336        AtomicServer::new(&addr, option, None).await
337    }
338
339    /// Creates a server instance for internal network use with existing ClientSenders.
340    ///
341    /// # Arguments
342    ///
343    /// * `addr` - Address to bind the server to (e.g., "0.0.0.0:9000")
344    /// * `option` - Server configuration options
345    /// * `client_senders` - Existing ClientSenders instance
346    ///
347    /// # Returns
348    ///
349    /// A `Result` containing the new `AtomicServer` instance, or an IO error if binding fails
350    pub async fn get_internal_server_with_client_senders(
351        addr: String,
352        option: ServerOptions,
353        client_senders: RwClientSenders,
354    ) -> std::io::Result<AtomicServer> {
355        AtomicServer::new(&addr, option, Some(client_senders)).await
356    }
357}
358
359/// Internal helper function: Creates a client instance.
360///
361/// # Arguments
362///
363/// * `connection_store` - Persistence for connection-identity state
364/// * `options` - Client options
365/// * `atomic_websocket_type` - Connection type (Internal or External)
366/// * `server_sender` - Optional existing ServerSender instance
367///
368/// # Returns
369///
370/// An initialized AtomicClient instance
371async fn get_client(
372    connection_store: Arc<dyn ConnectionStore>,
373    options: ClientOptions,
374    server_sender: Option<RwServerSender>,
375) -> AtomicClient {
376    let mut server_sender = match server_sender {
377        Some(server_sender) => {
378            let server_sender_clone = server_sender.clone();
379            let mut server_sender_clone = server_sender_clone.write().await;
380            server_sender_clone.server_ip = options.url.clone();
381            server_sender_clone.options = options.clone();
382            drop(server_sender_clone);
383            server_sender
384        }
385        None => Arc::new(RwLock::new(ServerSender::new(
386            connection_store.clone(),
387            options.url.clone(),
388            options.clone(),
389        ))),
390    };
391    server_sender.regist(server_sender.clone()).await;
392
393    let cancel_token = CancellationToken::new();
394    let atomic_websocket: AtomicClient = AtomicClient {
395        server_sender,
396        options,
397        cancel_token,
398    };
399    match atomic_websocket.options.atomic_websocket_type {
400        AtomicWebsocketType::Internal => {
401            atomic_websocket.internal_initialize(connection_store).await
402        }
403        AtomicWebsocketType::External => {
404            atomic_websocket.outer_initialize(connection_store).await
405        }
406    }
407    atomic_websocket
408}