Development Documentation (main branch) - For stable release docs, see docs.rs/eidetica
Skip to main content

eidetica/sync/transports/
shared.rs

1//! Shared utilities for transport implementations.
2//!
3//! This module provides common functionality used across different transport
4//! implementations to reduce code duplication and ensure consistency.
5
6use std::sync::Mutex;
7
8use tokio::sync::oneshot;
9
10use crate::sync::{
11    error::SyncError,
12    protocol::{SyncRequest, SyncResponse},
13};
14
15/// Manages server state common to all transport implementations.
16///
17/// Transports are shared rather than owned exclusively, so that an outbound
18/// request can be served from a clone while the engine goes on handling other
19/// commands. That rules out `&mut self`, so the mutable parts live behind a
20/// lock. It is a `std` mutex deliberately: every critical section here is a
21/// few field assignments with no await inside, so an async mutex would buy
22/// nothing and cost a scheduling point.
23pub struct ServerState {
24    inner: Mutex<ServerStateInner>,
25}
26
27/// A cancellation-safe reservation for starting one server.
28pub(super) struct ServerStart<'a> {
29    state: &'a ServerState,
30    completed: bool,
31}
32
33struct ServerStateInner {
34    /// Whether the server is starting or running.
35    active: bool,
36    /// Whether startup completed and the server is accepting requests.
37    running: bool,
38    /// Shutdown signal for the server loop.
39    shutdown: Option<oneshot::Sender<()>>,
40    /// The server's address.
41    address: Option<String>,
42}
43
44impl Default for ServerState {
45    fn default() -> Self {
46        Self::new()
47    }
48}
49
50impl ServerState {
51    /// Create a new server state manager.
52    pub fn new() -> Self {
53        Self {
54            inner: Mutex::new(ServerStateInner {
55                active: false,
56                running: false,
57                shutdown: None,
58                address: None,
59            }),
60        }
61    }
62
63    /// Reserve this transport for startup.
64    ///
65    /// Dropping the returned guard before startup completes releases the
66    /// reservation, including when the startup future is cancelled.
67    pub(super) fn begin_start(&self, address: &str) -> Result<ServerStart<'_>, SyncError> {
68        let mut inner = self.lock();
69        if inner.active {
70            return Err(SyncError::ServerAlreadyRunning {
71                address: address.to_string(),
72            });
73        }
74        inner.active = true;
75        Ok(ServerStart {
76            state: self,
77            completed: false,
78        })
79    }
80
81    /// Check if the server is currently running.
82    pub fn is_running(&self) -> bool {
83        self.lock().running
84    }
85
86    /// Get the server address if available.
87    pub fn get_address(&self) -> Result<String, SyncError> {
88        self.lock()
89            .address
90            .clone()
91            .ok_or(SyncError::ServerNotRunning)
92    }
93
94    /// Release a startup reservation after startup fails.
95    fn startup_failed(&self) {
96        let mut inner = self.lock();
97        inner.active = false;
98        inner.running = false;
99        inner.address = None;
100        inner.shutdown = None;
101    }
102
103    /// Stop the server by triggering shutdown and clearing state.
104    /// This combines the commonly used pair: trigger_shutdown + set_stopped.
105    pub fn stop_server(&self) {
106        let mut inner = self.lock();
107        // First trigger shutdown if we have a sender
108        if let Some(tx) = inner.shutdown.take() {
109            let _ = tx.send(());
110        }
111        // Then mark as stopped and clear address
112        inner.active = false;
113        inner.running = false;
114        inner.address = None;
115    }
116
117    /// Take the lock, recovering from a poisoned one.
118    ///
119    /// A panic while holding this lock leaves the flags describing a server
120    /// that is no longer being driven. Refusing to serve from then on would
121    /// turn one panicked task into a permanently dead transport, so the state
122    /// is taken as-is and the next start/stop corrects it.
123    fn lock(&self) -> std::sync::MutexGuard<'_, ServerStateInner> {
124        self.inner.lock().unwrap_or_else(|e| e.into_inner())
125    }
126}
127
128impl ServerStart<'_> {
129    /// Mark startup complete and install the running server state.
130    pub(super) fn complete(mut self, address: String, shutdown_sender: oneshot::Sender<()>) {
131        let mut inner = self.state.lock();
132        debug_assert!(inner.active);
133        inner.running = true;
134        inner.address = Some(address);
135        inner.shutdown = Some(shutdown_sender);
136        self.completed = true;
137    }
138}
139
140impl Drop for ServerStart<'_> {
141    fn drop(&mut self) {
142        if !self.completed {
143            self.state.startup_failed();
144        }
145    }
146}
147
148/// Utilities for handling JSON serialization/deserialization in transports.
149pub struct JsonHandler;
150
151impl JsonHandler {
152    /// Serialize a SyncRequest to JSON bytes.
153    pub fn serialize_request(request: &SyncRequest) -> Result<Vec<u8>, SyncError> {
154        serde_json::to_vec(request)
155            .map_err(|e| SyncError::Network(format!("Failed to serialize request: {e}")))
156    }
157
158    /// Serialize a SyncResponse to JSON bytes.
159    pub fn serialize_response(response: &SyncResponse) -> Result<Vec<u8>, SyncError> {
160        serde_json::to_vec(response)
161            .map_err(|e| SyncError::Network(format!("Failed to serialize response: {e}")))
162    }
163
164    /// Deserialize JSON bytes to a SyncRequest.
165    pub fn deserialize_request(bytes: &[u8]) -> Result<SyncRequest, SyncError> {
166        serde_json::from_slice(bytes)
167            .map_err(|e| SyncError::Network(format!("Failed to deserialize request: {e}")))
168    }
169
170    /// Deserialize JSON bytes to a SyncResponse.
171    pub fn deserialize_response(bytes: &[u8]) -> Result<SyncResponse, SyncError> {
172        serde_json::from_slice(bytes)
173            .map_err(|e| SyncError::Network(format!("Failed to deserialize response: {e}")))
174    }
175}
176
177/// Waits for server ready signal and maps errors appropriately.
178pub async fn wait_for_ready(
179    ready_rx: oneshot::Receiver<()>,
180    address: &str,
181) -> Result<(), SyncError> {
182    ready_rx.await.map_err(|_| SyncError::ServerBind {
183        address: address.to_string(),
184        reason: "Server startup failed".to_string(),
185    })
186}