truffle_core/synced_store/
sync.rs1use std::sync::Arc;
9
10use serde::de::DeserializeOwned;
11use serde::Serialize;
12use tokio::sync::{broadcast, mpsc};
13
14use crate::network::NetworkProvider;
15use crate::node::Node;
16use crate::session::PeerEvent;
17
18use super::types::{Slice, StoreEvent, SyncMessage};
19use super::StoreInner;
20
21pub(super) fn spawn_sync_task<N, T>(
25 node: Arc<Node<N>>,
26 inner: Arc<StoreInner<T>>,
27 mut broadcast_rx: mpsc::UnboundedReceiver<SyncMessage>,
28) -> tokio::task::JoinHandle<()>
29where
30 N: NetworkProvider + 'static,
31 T: Serialize + DeserializeOwned + Clone + Send + Sync + 'static,
32{
33 let namespace = format!("ss:{}", inner.store_id);
34
35 tokio::spawn(async move {
36 let mut msg_rx = node.subscribe(&namespace);
37 let mut peer_rx = node.on_peer_change();
38
39 tracing::info!(store = inner.store_id.as_str(), "synced_store: sync task started");
40
41 loop {
42 tokio::select! {
43 result = msg_rx.recv() => {
45 match result {
46 Ok(msg) => {
47 handle_incoming_message(&node, &inner, &namespace, &msg.from, msg.payload).await;
48 }
49 Err(broadcast::error::RecvError::Lagged(n)) => {
50 tracing::warn!(
51 store = inner.store_id.as_str(),
52 missed = n,
53 "synced_store: message receiver lagged"
54 );
55 }
56 Err(broadcast::error::RecvError::Closed) => {
57 tracing::debug!(store = inner.store_id.as_str(), "synced_store: message channel closed");
58 break;
59 }
60 }
61 }
62
63 result = peer_rx.recv() => {
65 match result {
66 Ok(PeerEvent::Joined(state)) => {
67 handle_peer_joined(&node, &inner, &namespace, &state.id).await;
68 }
69 Ok(PeerEvent::Left(peer_id)) => {
70 handle_peer_left(&inner, &peer_id).await;
71 }
72 Ok(_) => {} Err(broadcast::error::RecvError::Lagged(n)) => {
74 tracing::warn!(
75 store = inner.store_id.as_str(),
76 missed = n,
77 "synced_store: peer event receiver lagged"
78 );
79 }
80 Err(broadcast::error::RecvError::Closed) => {
81 tracing::debug!(store = inner.store_id.as_str(), "synced_store: peer event channel closed");
82 break;
83 }
84 }
85 }
86
87 Some(sync_msg) = broadcast_rx.recv() => {
89 let msg_type = match &sync_msg {
90 SyncMessage::Update { .. } => "update",
91 SyncMessage::Full { .. } => "full",
92 SyncMessage::Request { .. } => "request",
93 SyncMessage::Clear { .. } => "clear",
94 };
95 if let Ok(payload) = serde_json::to_value(&sync_msg) {
96 node.broadcast_typed(&namespace, msg_type, &payload).await;
97 }
98 }
99 }
100 }
101
102 tracing::info!(store = inner.store_id.as_str(), "synced_store: sync task stopped");
103 })
104}
105
106async fn handle_incoming_message<N, T>(
108 node: &Node<N>,
109 inner: &StoreInner<T>,
110 namespace: &str,
111 from: &str,
112 payload: serde_json::Value,
113) where
114 N: NetworkProvider + 'static,
115 T: Serialize + DeserializeOwned + Clone + Send + Sync + 'static,
116{
117 let sync_msg: SyncMessage = match serde_json::from_value(payload) {
118 Ok(m) => m,
119 Err(e) => {
120 tracing::warn!(
121 store = inner.store_id.as_str(),
122 from = from,
123 "synced_store: failed to parse sync message: {e}"
124 );
125 return;
126 }
127 };
128
129 match sync_msg {
130 SyncMessage::Update {
131 device_id,
132 data,
133 version,
134 updated_at,
135 }
136 | SyncMessage::Full {
137 device_id,
138 data,
139 version,
140 updated_at,
141 } => {
142 apply_remote_slice(inner, &device_id, data, version, updated_at).await;
143 }
144
145 SyncMessage::Request {} => {
146 let local = inner.local.read().await;
148 if let Some(slice) = local.as_ref() {
149 let full = SyncMessage::Full {
150 device_id: inner.device_id.clone(),
151 data: match serde_json::to_value(&slice.data) {
152 Ok(v) => v,
153 Err(e) => {
154 tracing::error!(
155 store = inner.store_id.as_str(),
156 "synced_store: failed to serialize local data: {e}"
157 );
158 return;
159 }
160 },
161 version: slice.version,
162 updated_at: slice.updated_at,
163 };
164 if let Ok(payload) = serde_json::to_value(&full) {
165 if let Err(e) = node.send_typed(from, namespace, "full", &payload).await {
166 tracing::warn!(
167 store = inner.store_id.as_str(),
168 peer = from,
169 "synced_store: failed to send Full response: {e}"
170 );
171 }
172 }
173 }
174 }
175
176 SyncMessage::Clear { device_id } => {
177 let mut remotes = inner.remotes.write().await;
178 if remotes.remove(&device_id).is_some() {
179 let _ = inner.event_tx.send(StoreEvent::PeerRemoved {
180 device_id: device_id.clone(),
181 });
182 tracing::debug!(
183 store = inner.store_id.as_str(),
184 device = device_id.as_str(),
185 "synced_store: cleared remote slice (Clear message)"
186 );
187 }
188 }
189 }
190}
191
192async fn apply_remote_slice<T>(
194 inner: &StoreInner<T>,
195 device_id: &str,
196 data: serde_json::Value,
197 version: u64,
198 updated_at: u64,
199) where
200 T: Serialize + DeserializeOwned + Clone + Send + Sync + 'static,
201{
202 if device_id == inner.device_id {
204 return;
205 }
206
207 let typed_data: T = match serde_json::from_value(data) {
208 Ok(d) => d,
209 Err(e) => {
210 tracing::warn!(
211 store = inner.store_id.as_str(),
212 device = device_id,
213 "synced_store: failed to deserialize remote data: {e}"
214 );
215 return;
216 }
217 };
218
219 let mut remotes = inner.remotes.write().await;
220
221 if let Some(existing) = remotes.get(device_id) {
223 if version <= existing.version {
224 tracing::debug!(
225 store = inner.store_id.as_str(),
226 device = device_id,
227 existing_v = existing.version,
228 incoming_v = version,
229 "synced_store: stale update rejected"
230 );
231 return;
232 }
233 }
234
235 let slice = Slice {
236 device_id: device_id.to_string(),
237 data: typed_data.clone(),
238 version,
239 updated_at,
240 };
241
242 remotes.insert(device_id.to_string(), slice);
243
244 if let Ok(serialized) = serde_json::to_vec(&typed_data) {
246 inner.backend.save(&inner.store_id, device_id, &serialized, version);
247 }
248
249 let _ = inner.event_tx.send(StoreEvent::PeerUpdated {
250 device_id: device_id.to_string(),
251 data: typed_data,
252 version,
253 });
254
255 tracing::debug!(
256 store = inner.store_id.as_str(),
257 device = device_id,
258 version = version,
259 "synced_store: applied remote slice"
260 );
261}
262
263async fn handle_peer_joined<N, T>(
265 node: &Node<N>,
266 inner: &StoreInner<T>,
267 namespace: &str,
268 peer_id: &str,
269) where
270 N: NetworkProvider + 'static,
271 T: Serialize + DeserializeOwned + Clone + Send + Sync + 'static,
272{
273 let request = SyncMessage::Request {};
275 if let Ok(payload) = serde_json::to_value(&request) {
276 if let Err(e) = node.send_typed(peer_id, namespace, "request", &payload).await {
277 tracing::warn!(
278 store = inner.store_id.as_str(),
279 peer = peer_id,
280 "synced_store: failed to send Request to new peer: {e}"
281 );
282 }
283 }
284
285 let local = inner.local.read().await;
287 if let Some(slice) = local.as_ref() {
288 let full = SyncMessage::Full {
289 device_id: inner.device_id.clone(),
290 data: match serde_json::to_value(&slice.data) {
291 Ok(v) => v,
292 Err(_) => return,
293 },
294 version: slice.version,
295 updated_at: slice.updated_at,
296 };
297 if let Ok(payload) = serde_json::to_value(&full) {
298 let _ = node.send_typed(peer_id, namespace, "full", &payload).await;
299 }
300 }
301}
302
303async fn handle_peer_left<T>(inner: &StoreInner<T>, peer_id: &str)
305where
306 T: Clone + Send + Sync + 'static,
307{
308 let mut remotes = inner.remotes.write().await;
309 if remotes.remove(peer_id).is_some() {
310 inner.backend.remove(&inner.store_id, peer_id);
312
313 let _ = inner.event_tx.send(StoreEvent::PeerRemoved {
314 device_id: peer_id.to_string(),
315 });
316 tracing::debug!(
317 store = inner.store_id.as_str(),
318 device = peer_id,
319 "synced_store: removed peer slice (peer left)"
320 );
321 }
322}