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}