eidetica/sync/transports/
shared.rs1use std::sync::Mutex;
7
8use tokio::sync::oneshot;
9
10use crate::sync::{
11 error::SyncError,
12 protocol::{SyncRequest, SyncResponse},
13};
14
15pub struct ServerState {
24 inner: Mutex<ServerStateInner>,
25}
26
27pub(super) struct ServerStart<'a> {
29 state: &'a ServerState,
30 completed: bool,
31}
32
33struct ServerStateInner {
34 active: bool,
36 running: bool,
38 shutdown: Option<oneshot::Sender<()>>,
40 address: Option<String>,
42}
43
44impl Default for ServerState {
45 fn default() -> Self {
46 Self::new()
47 }
48}
49
50impl ServerState {
51 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 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 pub fn is_running(&self) -> bool {
83 self.lock().running
84 }
85
86 pub fn get_address(&self) -> Result<String, SyncError> {
88 self.lock()
89 .address
90 .clone()
91 .ok_or(SyncError::ServerNotRunning)
92 }
93
94 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 pub fn stop_server(&self) {
106 let mut inner = self.lock();
107 if let Some(tx) = inner.shutdown.take() {
109 let _ = tx.send(());
110 }
111 inner.active = false;
113 inner.running = false;
114 inner.address = None;
115 }
116
117 fn lock(&self) -> std::sync::MutexGuard<'_, ServerStateInner> {
124 self.inner.lock().unwrap_or_else(|e| e.into_inner())
125 }
126}
127
128impl ServerStart<'_> {
129 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
148pub struct JsonHandler;
150
151impl JsonHandler {
152 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 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 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 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
177pub 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}