Skip to main content

truffle_core/synced_store/
sync.rs

1//! Background sync task for SyncedStore.
2//!
3//! Handles three event sources via `tokio::select!`:
4//! 1. Incoming sync messages from peers (namespace `"ss:{store_id}"`)
5//! 2. Peer join/leave events (for sync-on-join and cleanup-on-leave)
6//! 3. Outbound broadcast requests from `SyncedStore::set()`
7
8use 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
21/// Spawn the background sync task.
22///
23/// Returns a `JoinHandle` that can be aborted to stop the task.
24pub(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                // ── Source 1: incoming sync messages from peers ──
44                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                // ── Source 2: peer join/leave events ──
64                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(_) => {} // Updated, WsConnected, WsDisconnected, AuthRequired — ignore
73                        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                // ── Source 3: outbound broadcast requests from set() ──
88                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
106/// Handle an incoming sync message from a peer.
107async 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            // Peer is requesting our current slice.
147            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
192/// Apply a remote slice if the version is newer.
193async 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    // Don't apply our own data back.
203    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    // Check version — only accept if newer.
222    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    // Persist remote slice to backend.
245    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
263/// Handle a peer joining: send them a Request for their data.
264async 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    // Send Request to the new peer.
274    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    // Also send our current data to the new peer (proactive full sync).
286    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
303/// Handle a peer leaving: remove their slice and emit PeerRemoved.
304async 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        // Remove persisted slice for departed peer.
311        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}