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

eidetica/sync/
peer.rs

1//! Peer management, sync relationships, and address handling for the sync system.
2
3use std::time::{Duration, SystemTime, UNIX_EPOCH};
4
5use tokio::sync::oneshot;
6use tracing::info;
7
8use super::{
9    Address, ConnectionState, PeerId, PeerInfo, PeerStatus, Sync, SyncError, SyncHandle,
10    SyncPeerInfo, SyncStatus, background::SyncCommand, peer_manager::PeerManager,
11};
12use crate::{Result, auth::crypto::PublicKey, entry::ID};
13
14impl Sync {
15    // === Peer Management Methods ===
16
17    /// Register a new remote peer in the sync network.
18    ///
19    /// # Arguments
20    /// * `pubkey` - The peer's public key
21    /// * `display_name` - Optional human-readable name for the peer
22    ///
23    /// # Returns
24    /// A Result indicating success or an error.
25    pub async fn register_peer(
26        &self,
27        pubkey: &PublicKey,
28        display_name: Option<&str>,
29    ) -> Result<()> {
30        // Store in sync tree via PeerManager
31        let txn = self.sync_tree.new_transaction().await?;
32        PeerManager::new(&txn)
33            .register_peer(pubkey, display_name)
34            .await?;
35        txn.commit().await?;
36
37        // Background sync will read peer info directly from sync tree when needed
38        Ok(())
39    }
40
41    /// Update the status of a registered peer.
42    ///
43    /// # Arguments
44    /// * `pubkey` - The peer's public key
45    /// * `status` - The new status for the peer
46    ///
47    /// # Returns
48    /// A Result indicating success or an error.
49    pub async fn update_peer_status(&self, pubkey: &PublicKey, status: PeerStatus) -> Result<()> {
50        let txn = self.sync_tree.new_transaction().await?;
51        PeerManager::new(&txn)
52            .update_peer_status(pubkey, status)
53            .await?;
54        txn.commit().await?;
55        Ok(())
56    }
57
58    /// Get information about a registered peer.
59    ///
60    /// # Arguments
61    /// * `pubkey` - The peer's public key
62    ///
63    /// # Returns
64    /// The peer information if found, None otherwise.
65    pub async fn get_peer_info(&self, pubkey: &PublicKey) -> Result<Option<PeerInfo>> {
66        let txn = self.sync_tree.new_transaction().await?;
67        PeerManager::new(&txn).get_peer_info(pubkey).await
68        // No commit - just reading
69    }
70
71    /// List all registered peers.
72    ///
73    /// # Returns
74    /// A vector of all registered peer information.
75    pub async fn list_peers(&self) -> Result<Vec<PeerInfo>> {
76        let txn = self.sync_tree.new_transaction().await?;
77        PeerManager::new(&txn).list_peers().await
78        // No commit - just reading
79    }
80
81    /// Remove a peer from the sync network.
82    ///
83    /// This removes the peer entry and all associated sync relationships and transport info.
84    ///
85    /// # Arguments
86    /// * `pubkey` - The peer's public key
87    ///
88    /// # Returns
89    /// A Result indicating success or an error.
90    pub async fn remove_peer(&self, pubkey: &PublicKey) -> Result<()> {
91        let txn = self.sync_tree.new_transaction().await?;
92        PeerManager::new(&txn).remove_peer(pubkey).await?;
93        txn.commit().await?;
94        Ok(())
95    }
96
97    // === Declarative Sync API ===
98
99    /// Register a peer for syncing (declarative API).
100    ///
101    /// This is the recommended way to set up syncing. It immediately registers
102    /// the peer and tree/peer relationship, then the background sync engine
103    /// handles the actual data synchronization.
104    ///
105    /// # Arguments
106    /// * `info` - Information about the peer and sync configuration
107    ///
108    /// # Returns
109    /// A handle for tracking sync status and adding more address hints.
110    ///
111    /// # Example
112    /// ```no_run
113    /// # use eidetica::*;
114    /// # use eidetica::sync::{SyncPeerInfo, Address, AuthParams};
115    /// # use eidetica::auth::PublicKey;
116    /// # async fn example(sync: sync::Sync, peer_pubkey: PublicKey, tree_id: entry::ID) -> Result<()> {
117    /// // Register peer for syncing
118    /// let handle = sync.register_sync_peer(SyncPeerInfo {
119    ///     peer_pubkey,
120    ///     tree_id,
121    ///     addresses: vec![Address {
122    ///         transport_type: "http".to_string(),
123    ///         address: "http://localhost:8080".to_string(),
124    ///     }],
125    ///     auth: None,
126    ///     display_name: Some("My Peer".to_string()),
127    /// }).await?;
128    ///
129    /// // Optionally wait for initial sync
130    /// handle.wait_for_initial_sync().await?;
131    ///
132    /// // Check status anytime
133    /// let status = handle.status().await?;
134    /// println!("Has local data: {}", status.has_local_data);
135    /// # Ok(())
136    /// # }
137    /// ```
138    pub async fn register_sync_peer(&self, info: SyncPeerInfo) -> Result<SyncHandle> {
139        let txn = self.sync_tree.new_transaction().await?;
140        let peer_mgr = PeerManager::new(&txn);
141
142        // Register peer if it doesn't exist
143        if peer_mgr.get_peer_info(&info.peer_pubkey).await?.is_none() {
144            peer_mgr
145                .register_peer(&info.peer_pubkey, info.display_name.as_deref())
146                .await?;
147        }
148
149        // Add all address hints
150        for addr in &info.addresses {
151            peer_mgr
152                .add_address(&info.peer_pubkey, addr.clone())
153                .await?;
154        }
155
156        // Register the tree/peer relationship
157        peer_mgr
158            .add_tree_sync(&info.peer_pubkey, &info.tree_id)
159            .await?;
160
161        // TODO: Store auth params if provided for bootstrap
162        // For now, auth is passed during the actual sync handshake via on_local_write callback
163
164        txn.commit().await?;
165
166        info!(
167            peer = %info.peer_pubkey,
168            tree = %info.tree_id,
169            address_count = info.addresses.len(),
170            "Registered peer for syncing"
171        );
172
173        Ok(SyncHandle {
174            tree_id: info.tree_id,
175            peer_pubkey: info.peer_pubkey,
176            sync: self.clone(),
177        })
178    }
179
180    /// Get the current sync status for a tree/peer pair.
181    ///
182    /// # Arguments
183    /// * `tree_id` - The tree to check
184    /// * `peer_pubkey` - The peer public key
185    ///
186    /// # Returns
187    /// Current sync status including whether we have local data.
188    pub async fn get_sync_status(
189        &self,
190        tree_id: &ID,
191        peer_pubkey: &PublicKey,
192    ) -> Result<SyncStatus> {
193        // Check if we have local data for this tree
194        let backend = self.backend()?;
195        let our_snapshot = backend.snapshot(tree_id).await.unwrap_or_default();
196
197        // TODO: Track last_error in sync tree
198        Ok(SyncStatus {
199            has_local_data: !our_snapshot.is_empty(),
200            last_sync: self.peer_last_sync(peer_pubkey),
201            last_error: None,
202        })
203    }
204
205    /// When the background engine last synced with a peer.
206    ///
207    /// Reads the shared liveness state directly rather than going through the
208    /// command channel, so it does not wait on an in-flight sync round. `None`
209    /// when no successful round is on record, including when the engine has not
210    /// started — both mean the same thing to a caller.
211    fn peer_last_sync(&self, peer_pubkey: &PublicKey) -> Option<SystemTime> {
212        let millis = self
213            .peer_state
214            .last_success_ms(&PeerId::from(peer_pubkey))?;
215        Some(UNIX_EPOCH + Duration::from_millis(millis))
216    }
217
218    // === Database Sync Relationship Methods ===
219
220    /// Add a tree to the sync relationship with a peer.
221    ///
222    /// # Arguments
223    /// * `peer_pubkey` - The peer's public key
224    /// * `tree_root_id` - The root ID of the tree to sync
225    ///
226    /// # Returns
227    /// A Result indicating success or an error.
228    pub async fn add_tree_sync(&self, peer_pubkey: &PublicKey, tree_root_id: &ID) -> Result<()> {
229        let txn = self.sync_tree.new_transaction().await?;
230        PeerManager::new(&txn)
231            .add_tree_sync(peer_pubkey, tree_root_id)
232            .await?;
233        txn.commit().await?;
234        Ok(())
235    }
236
237    /// Remove a tree from the sync relationship with a peer.
238    ///
239    /// # Arguments
240    /// * `peer_pubkey` - The peer's public key
241    /// * `tree_root_id` - The root ID of the tree to stop syncing
242    ///
243    /// # Returns
244    /// A Result indicating success or an error.
245    pub async fn remove_tree_sync(&self, peer_pubkey: &PublicKey, tree_root_id: &ID) -> Result<()> {
246        let txn = self.sync_tree.new_transaction().await?;
247        PeerManager::new(&txn)
248            .remove_tree_sync(peer_pubkey, tree_root_id)
249            .await?;
250        txn.commit().await?;
251        Ok(())
252    }
253
254    /// Get the list of trees synced with a peer.
255    ///
256    /// # Arguments
257    /// * `peer_pubkey` - The peer's public key
258    ///
259    /// # Returns
260    /// A vector of tree root IDs synced with this peer.
261    pub async fn get_peer_trees(&self, peer_pubkey: &PublicKey) -> Result<Vec<ID>> {
262        let txn = self.sync_tree.new_transaction().await?;
263        PeerManager::new(&txn).get_peer_trees(peer_pubkey).await
264        // No commit - just reading
265    }
266
267    /// Get all peers that sync a specific tree.
268    ///
269    /// # Arguments
270    /// * `tree_root_id` - The root ID of the tree
271    ///
272    /// # Returns
273    /// A vector of peer IDs that sync this tree.
274    pub async fn get_tree_peers(&self, tree_root_id: &ID) -> Result<Vec<PeerId>> {
275        let txn = self.sync_tree.new_transaction().await?;
276        PeerManager::new(&txn).get_tree_peers(tree_root_id).await
277        // No commit - just reading
278    }
279
280    /// Connect to a remote peer and perform handshake.
281    ///
282    /// This method initiates a connection to a peer, performs the handshake protocol,
283    /// and automatically registers the peer if successful.
284    ///
285    /// # Arguments
286    /// * `address` - The address of the peer to connect to
287    ///
288    /// # Returns
289    /// A Result containing the peer's public key if successful.
290    pub async fn connect_to_peer(&self, address: &Address) -> Result<PublicKey> {
291        let (tx, rx) = oneshot::channel();
292
293        self.background_tx
294            .get()
295            .ok_or(SyncError::NoTransportEnabled)?
296            .send(SyncCommand::ConnectToPeer {
297                address: address.clone(),
298                response: tx,
299            })
300            .await
301            .map_err(|e| SyncError::CommandSendError(e.to_string()))?;
302
303        rx.await
304            .map_err(|e| SyncError::Network(format!("Response channel error: {e}")))?
305    }
306
307    /// Update the connection state of a peer.
308    ///
309    /// # Arguments
310    /// * `pubkey` - The peer's public key
311    /// * `state` - The new connection state
312    ///
313    /// # Returns
314    /// A Result indicating success or an error.
315    pub async fn update_peer_connection_state(
316        &self,
317        pubkey: &PublicKey,
318        state: ConnectionState,
319    ) -> Result<()> {
320        let txn = self.sync_tree.new_transaction().await?;
321        let peer_manager = PeerManager::new(&txn);
322
323        // Get current peer info
324        let mut peer_info = match peer_manager.get_peer_info(pubkey).await? {
325            Some(info) => info,
326            None => return Err(SyncError::PeerNotFound(pubkey.to_string()).into()),
327        };
328
329        // Update connection state
330        peer_info.connection_state = state;
331        peer_info.touch_at(txn.now_rfc3339()?);
332
333        // Save updated peer info
334        peer_manager.update_peer_info(pubkey, peer_info).await?;
335        txn.commit().await?;
336        Ok(())
337    }
338
339    /// Check if a tree is synced with a specific peer.
340    ///
341    /// # Arguments
342    /// * `peer_pubkey` - The peer's public key
343    /// * `tree_root_id` - The root ID of the tree
344    ///
345    /// # Returns
346    /// True if the tree is synced with the peer, false otherwise.
347    pub async fn is_tree_synced_with_peer(
348        &self,
349        peer_pubkey: &PublicKey,
350        tree_root_id: &ID,
351    ) -> Result<bool> {
352        let txn = self.sync_tree.new_transaction().await?;
353        PeerManager::new(&txn)
354            .is_tree_synced_with_peer(peer_pubkey, tree_root_id)
355            .await
356        // No commit - just reading
357    }
358
359    // === Address Management Methods ===
360
361    /// Add an address to a peer.
362    ///
363    /// # Arguments
364    /// * `peer_pubkey` - The peer's public key
365    /// * `address` - The address to add
366    ///
367    /// # Returns
368    /// A Result indicating success or an error.
369    pub async fn add_peer_address(&self, peer_pubkey: &PublicKey, address: Address) -> Result<()> {
370        // Update sync tree via PeerManager
371        let txn = self.sync_tree.new_transaction().await?;
372        PeerManager::new(&txn)
373            .add_address(peer_pubkey, address)
374            .await?;
375        txn.commit().await?;
376
377        // Background sync will read updated peer info directly from sync tree when needed
378        Ok(())
379    }
380
381    /// Remove a specific address from a peer.
382    ///
383    /// # Arguments
384    /// * `peer_pubkey` - The peer's public key
385    /// * `address` - The address to remove
386    ///
387    /// # Returns
388    /// A Result indicating success or an error (true if removed, false if not found).
389    pub async fn remove_peer_address(
390        &self,
391        peer_pubkey: &PublicKey,
392        address: &Address,
393    ) -> Result<bool> {
394        let txn = self.sync_tree.new_transaction().await?;
395        let result = PeerManager::new(&txn)
396            .remove_address(peer_pubkey, address)
397            .await?;
398        txn.commit().await?;
399        Ok(result)
400    }
401
402    /// Get addresses for a peer, optionally filtered by transport type.
403    ///
404    /// # Arguments
405    /// * `peer_pubkey` - The peer's public key
406    /// * `transport_type` - Optional transport type filter
407    ///
408    /// # Returns
409    /// A vector of addresses matching the criteria.
410    pub async fn get_peer_addresses(
411        &self,
412        peer_pubkey: &PublicKey,
413        transport_type: Option<&str>,
414    ) -> Result<Vec<Address>> {
415        let txn = self.sync_tree.new_transaction().await?;
416        PeerManager::new(&txn)
417            .get_addresses(peer_pubkey, transport_type)
418            .await
419        // No commit - just reading
420    }
421}