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}