eidetica/sync/background/mod.rs
1//! Background sync engine implementation.
2//!
3//! This module provides the BackgroundSync struct that handles all sync operations
4//! in a single background thread, removing circular dependency issues and providing
5//! automatic retry, periodic sync, and reconnection handling.
6
7use std::{sync::Arc, time::Duration};
8
9use tokio::{
10 sync::{mpsc, oneshot},
11 time::interval,
12};
13use tracing::{Instrument, debug, info, info_span, trace, warn};
14
15use super::{
16 error::SyncError,
17 handler::SyncHandlerImpl,
18 peer_manager::PeerManager,
19 peer_state::PeerStates,
20 peer_types::{Address, PeerId, PeerStatus},
21 protocol::{SyncRequest, SyncRequestAuth, SyncResponse, SyncTreeRequest},
22 queue::SyncQueue,
23 transport_manager::TransportManager,
24 transports::SyncTransport,
25};
26use crate::{
27 Database, Error, Instance, Result, WeakInstance,
28 auth::crypto::PublicKey,
29 entry::{Entry, ID},
30 store::DocStore,
31};
32
33mod conn;
34
35/// Commands that can be sent to the background sync engine
36#[allow(clippy::large_enum_variant)]
37pub enum SyncCommand {
38 /// Send entries to a specific peer
39 SendEntries { peer: PeerId, entries: Vec<Entry> },
40 /// Trigger immediate sync with a peer
41 SyncWithPeer { peer: PeerId },
42 /// Shutdown the background engine
43 Shutdown,
44
45 // Transport management
46 /// Add a named transport to the transport manager
47 AddTransport {
48 name: String,
49 transport: Box<dyn super::transports::SyncTransport>,
50 response: oneshot::Sender<Result<()>>,
51 },
52
53 // Server management commands
54 /// Start the sync server on specified or all transports
55 StartServer {
56 /// Transport name to start, or None for all transports
57 name: Option<String>,
58 response: oneshot::Sender<Result<()>>,
59 },
60 /// Stop the sync server on specified or all transports
61 StopServer {
62 /// Transport name to stop, or None for all transports
63 name: Option<String>,
64 response: oneshot::Sender<Result<()>>,
65 },
66 /// Get the server's listening address for a specific transport
67 GetServerAddress {
68 name: String,
69 response: oneshot::Sender<Result<String>>,
70 },
71 /// Get all server addresses for running servers
72 GetAllServerAddresses {
73 response: oneshot::Sender<Result<Vec<(String, String)>>>,
74 },
75
76 // Peer connection
77 /// Connect to a peer and perform handshake
78 ConnectToPeer {
79 address: Address,
80 response: oneshot::Sender<Result<PublicKey>>, // Returns peer pubkey
81 },
82
83 // Request/Response operations
84 /// Send a sync request and get response
85 SendRequest {
86 address: Address,
87 request: Box<SyncRequest>,
88 response: oneshot::Sender<Result<SyncResponse>>,
89 },
90
91 /// Flush: process all queued entries and retry queue, then respond.
92 Flush {
93 response: oneshot::Sender<Result<()>>,
94 },
95}
96
97// Manual Debug impl required because:
98// - `Box<dyn SyncTransport>` doesn't implement Debug (trait object)
99// - `oneshot::Sender` doesn't implement Debug (channel internals)
100// - Transports may contain secrets (e.g., Iroh's cryptographic keys)
101// This impl provides safe, useful debug output for logging.
102impl std::fmt::Debug for SyncCommand {
103 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
104 match self {
105 Self::SendEntries { peer, entries } => f
106 .debug_struct("SendEntries")
107 .field("peer", peer)
108 .field("entries_count", &entries.len())
109 .finish(),
110 Self::SyncWithPeer { peer } => {
111 f.debug_struct("SyncWithPeer").field("peer", peer).finish()
112 }
113 Self::Shutdown => write!(f, "Shutdown"),
114 Self::AddTransport {
115 name, transport, ..
116 } => f
117 .debug_struct("AddTransport")
118 .field("name", name)
119 .field("transport_type", &transport.transport_type())
120 .finish(),
121 Self::StartServer { name, .. } => {
122 f.debug_struct("StartServer").field("name", name).finish()
123 }
124 Self::StopServer { name, .. } => {
125 f.debug_struct("StopServer").field("name", name).finish()
126 }
127 Self::GetServerAddress { name, .. } => f
128 .debug_struct("GetServerAddress")
129 .field("name", name)
130 .finish(),
131 Self::GetAllServerAddresses { .. } => write!(f, "GetAllServerAddresses"),
132 Self::ConnectToPeer { address, .. } => f
133 .debug_struct("ConnectToPeer")
134 .field("address", address)
135 .finish(),
136 Self::SendRequest {
137 address, request, ..
138 } => f
139 .debug_struct("SendRequest")
140 .field("address", address)
141 .field("request", request)
142 .finish(),
143 Self::Flush { .. } => write!(f, "Flush"),
144 }
145 }
146}
147
148/// Entry in the retry queue for failed sends
149#[derive(Debug, Clone)]
150struct RetryEntry {
151 peer: PeerId,
152 entries: Vec<Entry>,
153 attempts: u32,
154 /// Timestamp of last attempt in milliseconds since Unix epoch
155 last_attempt_ms: u64,
156}
157
158/// Background sync engine that owns all sync state and handles operations
159pub struct BackgroundSync {
160 // Core components - owns everything
161 pub(super) transport_manager: TransportManager,
162 instance: WeakInstance,
163 pub(super) sync_tree_id: ID,
164
165 // Queue for entries pending synchronization (shared with Sync frontend)
166 queue: Arc<SyncQueue>,
167
168 // Per-peer liveness state (shared with Sync frontend, written only here)
169 peer_state: Arc<PeerStates>,
170
171 // Retry queue for failed sends
172 retry_queue: Vec<RetryEntry>,
173
174 // Communication
175 command_rx: mpsc::Receiver<SyncCommand>,
176}
177
178/// How many peers may be synced at once, and so how many outbound connections a
179/// periodic sync can hold open. Bounds a large peer set; well above the point
180/// where the work stops being dominated by any single unreachable peer.
181const MAX_CONCURRENT_PEER_SYNCS: usize = 16;
182
183/// Bound on a single registration handshake while a route is being selected.
184///
185/// It covers the handshake only. The transfer that follows runs outside it, so
186/// a legitimately large exchange is never cut short by a connect-shaped deadline
187/// — the transports impose their own read-silence timeouts for that.
188const ADDRESS_ATTEMPT_TIMEOUT: Duration = Duration::from_secs(30);
189
190impl BackgroundSync {
191 /// Start the background sync engine and return a command sender.
192 ///
193 /// The engine starts with no transports registered. Use `AddTransport`
194 /// commands to add transports after starting.
195 pub fn start(
196 instance: Instance,
197 sync_tree_id: ID,
198 queue: Arc<SyncQueue>,
199 peer_state: Arc<PeerStates>,
200 ) -> mpsc::Sender<SyncCommand> {
201 let (tx, rx) = mpsc::channel(100);
202
203 let background = Self {
204 transport_manager: TransportManager::new(),
205 instance: instance.downgrade(),
206 sync_tree_id,
207 queue,
208 peer_state,
209 retry_queue: Vec::new(),
210 command_rx: rx,
211 };
212
213 // Spawn background sync as a regular tokio task
214 // (Transaction is now Send since it uses Arc<Mutex>)
215 tokio::spawn(background.run());
216 tx
217 }
218
219 /// Upgrade the weak instance reference to a strong reference.
220 pub(super) fn instance(&self) -> Result<Instance> {
221 self.instance
222 .upgrade()
223 .ok_or_else(|| SyncError::InstanceDropped.into())
224 }
225
226 /// Collect what a handshake needs so it can run off the command loop.
227 ///
228 /// Everything gathered here is either a cheap handle or local engine state
229 /// (the listen addresses), so the borrow ends before the handshake starts.
230 fn handshake_ctx(&self, address: &Address) -> Result<conn::HandshakeCtx> {
231 let transport = self
232 .transport_manager
233 .handle_for_address(address)
234 .ok_or_else(|| SyncError::NoTransportForAddress {
235 address: address.clone(),
236 })?;
237 Ok(conn::HandshakeCtx {
238 transport,
239 instance: self.instance.clone(),
240 sync_tree_id: self.sync_tree_id.clone(),
241 listen_addresses: self
242 .transport_manager
243 .get_all_server_addresses()
244 .into_iter()
245 .map(|(transport_type, addr)| Address {
246 transport_type,
247 address: addr,
248 })
249 .collect(),
250 })
251 }
252
253 /// Get the sync tree for accessing peer data
254 pub(super) async fn get_sync_tree(&self) -> Result<Database> {
255 // Load sync tree with the device key
256 let instance = self.instance()?;
257 let signing_key = instance.signing_key()?.clone();
258 Ok(Database::open(&instance, &self.sync_tree_id)
259 .await?
260 .with_key(signing_key))
261 }
262
263 /// Get the minimum sync interval from all tracked databases
264 /// Returns None if no databases are tracked or no intervals are set
265 async fn get_min_sync_interval(&self) -> Option<u64> {
266 let sync_tree = match self.get_sync_tree().await {
267 Ok(tree) => tree,
268 Err(_) => return None,
269 };
270
271 let txn = match sync_tree.new_transaction().await {
272 Ok(txn) => txn,
273 Err(_) => return None,
274 };
275
276 let user_mgr = super::user_sync_manager::UserSyncManager::new(&txn);
277
278 // Get all tracked database IDs from the DATABASE_USERS_SUBTREE
279 let database_users = match txn
280 .get_store::<DocStore>(super::user_sync_manager::DATABASE_USERS_SUBTREE)
281 .await
282 {
283 Ok(store) => store,
284 Err(_) => return None,
285 };
286
287 let all_dbs = match database_users.get_all().await {
288 Ok(doc) => doc,
289 Err(_) => return None,
290 };
291
292 // Find the minimum interval across all databases
293 let mut min_interval: Option<u64> = None;
294 for db_id_str in all_dbs.keys() {
295 if let Ok(db_id) = ID::parse(db_id_str)
296 && let Ok(Some(settings)) = user_mgr.get_combined_settings(&db_id).await
297 && let Some(interval) = settings.interval_seconds
298 {
299 min_interval = Some(match min_interval {
300 Some(current_min) => current_min.min(interval),
301 None => interval,
302 });
303 }
304 }
305
306 min_interval
307 }
308
309 /// Main event loop that handles all sync operations
310 async fn run(mut self) {
311 async move {
312 info!("Starting background sync engine");
313
314 // Get initial sync interval from settings (default to 300 seconds if none set)
315 let mut current_interval_secs = self.get_min_sync_interval().await.unwrap_or(300);
316 info!("Initial periodic sync interval: {} seconds", current_interval_secs);
317
318 // Set up timers
319 let mut periodic_sync = interval(Duration::from_secs(current_interval_secs));
320 let mut queue_check = interval(Duration::from_secs(5)); // 5 seconds - batches local writes
321 let mut retry_check = interval(Duration::from_secs(30)); // 30 seconds
322 let mut connection_check = interval(Duration::from_secs(60)); // 1 minute
323 let mut settings_check = interval(Duration::from_secs(60)); // Check for settings changes every minute
324
325 // Skip initial tick to avoid immediate execution
326 periodic_sync.tick().await;
327 queue_check.tick().await;
328 retry_check.tick().await;
329 connection_check.tick().await;
330 settings_check.tick().await;
331
332 loop {
333 tokio::select! {
334 // Handle commands from frontend
335 Some(cmd) = self.command_rx.recv() => {
336 if let Err(e) = self.handle_command(cmd).await {
337 // Log errors but continue running - background sync should be resilient
338 tracing::error!("Background sync command error: {e}");
339 }
340 }
341
342 // Drain sync queue (batched entries)
343 _ = queue_check.tick() => {
344 self.process_queue().await;
345 }
346
347 // Periodic sync with all peers
348 _ = periodic_sync.tick() => {
349 self.periodic_sync_all_peers().await;
350 }
351
352 // Process retry queue
353 _ = retry_check.tick() => {
354 self.process_retry_queue().await;
355 }
356
357 // Check and reconnect disconnected peers
358 _ = connection_check.tick() => {
359 self.check_peer_connections().await;
360 }
361
362 // Check if sync interval settings have changed
363 _ = settings_check.tick() => {
364 if let Some(new_interval) = self.get_min_sync_interval().await
365 && new_interval != current_interval_secs {
366 info!("Sync interval changed from {} to {} seconds", current_interval_secs, new_interval);
367 current_interval_secs = new_interval;
368 // Recreate the periodic sync timer with new interval
369 periodic_sync = interval(Duration::from_secs(new_interval));
370 periodic_sync.tick().await; // Skip initial tick
371 }
372 }
373
374 // Channel closed, shutdown
375 else => {
376 // Normal shutdown when channel closes
377 info!("Background sync engine shutting down");
378 break;
379 }
380 }
381 }
382 }
383 .instrument(info_span!("background_sync"))
384 .await
385 }
386
387 /// Handle a single command from the frontend
388 async fn handle_command(&mut self, command: SyncCommand) -> Result<()> {
389 match command {
390 SyncCommand::SendEntries { peer, entries } => {
391 if let Err(e) = self.send_to_peer(&peer, entries.clone()).await {
392 let now_ms = self.instance().map(|i| i.clock().now_millis()).unwrap_or(0);
393 self.add_to_retry_queue(peer, entries, e, now_ms);
394 }
395 }
396
397 SyncCommand::SyncWithPeer { peer } => {
398 if let Err(e) = self.sync_with_peer(&peer).await {
399 // Log sync failure but don't crash the background engine
400 tracing::error!("Failed to sync with peer {peer}: {e}");
401 }
402 }
403
404 SyncCommand::AddTransport {
405 name,
406 transport,
407 response,
408 } => {
409 // Stop server on existing transport if running
410 if let Some(old) = self.transport_manager.get_mut(&name)
411 && old.is_server_running()
412 {
413 let _ = old.stop_server().await;
414 }
415 self.transport_manager.add(&name, Arc::from(transport));
416 tracing::debug!("Added transport: {}", name);
417 let _ = response.send(Ok(()));
418 }
419
420 SyncCommand::StartServer { name, response } => {
421 let result = self.start_server(name.as_deref()).await;
422 let _ = response.send(result);
423 }
424
425 SyncCommand::StopServer { name, response } => {
426 let result = self.stop_server(name.as_deref()).await;
427 let _ = response.send(result);
428 }
429
430 SyncCommand::GetServerAddress { name, response } => {
431 let result = self.transport_manager.get_server_address(&name);
432 let _ = response.send(result);
433 }
434
435 SyncCommand::GetAllServerAddresses { response } => {
436 let addresses = self.transport_manager.get_all_server_addresses();
437 let _ = response.send(Ok(addresses));
438 }
439
440 SyncCommand::ConnectToPeer { address, response } => {
441 // Served off the loop, for the same reason as SendRequest
442 // below: a handshake is aimed at a peer this engine has never
443 // reached, which is precisely the peer most likely not to be
444 // there. Awaiting it inline hands the whole engine to a
445 // stranger for a full connect deadline.
446 match self.handshake_ctx(&address) {
447 Ok(ctx) => {
448 tokio::spawn(async move {
449 let _ = response.send(conn::run_handshake(ctx, address).await);
450 });
451 }
452 Err(e) => {
453 let _ = response.send(Err(e));
454 }
455 }
456 }
457
458 SyncCommand::SendRequest {
459 address,
460 request,
461 response,
462 } => {
463 // Serve this from a task of its own rather than inline.
464 //
465 // Commands are handled one at a time, so awaiting a request
466 // here holds the engine for as long as the peer takes to
467 // answer — and a peer that accepts a connection and then goes
468 // quiet takes the full transport deadline. Every other
469 // request, the queue drain and the retry queue all wait behind
470 // it, so a deployment carrying a few retired peers starves the
471 // live ones. That presents as "everything times out", which
472 // reads as a local fault rather than as one absent peer.
473 //
474 // The transport handle is owned, so the task borrows nothing
475 // from the engine and the loop is free immediately.
476 match self.transport_manager.handle_for_address(&address) {
477 Some(transport) => {
478 tokio::spawn(async move {
479 let result = transport.send_request(&address, &request).await;
480 let _ = response.send(result);
481 });
482 }
483 None => {
484 let _ =
485 response.send(Err(SyncError::NoTransportForAddress { address }.into()));
486 }
487 }
488 }
489
490 SyncCommand::Flush { response } => {
491 // Process retry queue first (old failures), then main queue (new entries)
492 // This avoids double-trying entries that fail in process_queue
493 let retry_failures = self.flush_retry_queue().await;
494 let queue_failures = self.process_queue().await;
495
496 // Report error if any failures occurred
497 let result = match (retry_failures, queue_failures) {
498 (0, 0) => Ok(()),
499 (r, q) => Err(SyncError::Network(format!(
500 "Flush had failures: {r} from retry queue, {q} from new entries"
501 ))
502 .into()),
503 };
504 let _ = response.send(result);
505 }
506
507 SyncCommand::Shutdown => {
508 // Shutdown command received - exit cleanly
509 return Err(SyncError::Network("Shutdown requested".to_string()).into());
510 }
511 }
512 Ok(())
513 }
514
515 /// Race bounded registration handshakes and return the first usable route
516 /// **to `expected`**.
517 ///
518 /// For each address a transport handle is obtained via
519 /// [`TransportManager::handle_for_address`]. One task is spawned per
520 /// address, each subject to [`ADDRESS_ATTEMPT_TIMEOUT`]. Remaining
521 /// handshakes are **not** cancelled — they continue running so that
522 /// additional addresses can be registered by the remote peer. The caller
523 /// performs the real operation once on the selected route, outside this
524 /// timeout.
525 ///
526 /// A handshake identifies who actually answered, and only a route that
527 /// answers as `expected` is selected. A peer's address list only ever grows,
528 /// so a stale entry can be reoccupied by an unrelated node; that node
529 /// completes a handshake perfectly well. Taking it as the route would send
530 /// this peer's traffic to a stranger and, worse, stop the working address
531 /// behind it from ever being tried — the exact failure this whole path
532 /// exists to remove.
533 ///
534 /// If no address has a matching transport an error is returned. If all
535 /// tasks fail the last error is returned.
536 async fn select_route(
537 &self,
538 addresses: &[Address],
539 expected: &PublicKey,
540 ) -> Result<(std::sync::Arc<dyn SyncTransport>, Address)> {
541 // Collect transport handles upfront (before any spawn).
542 let mut tasks: Vec<(std::sync::Arc<dyn SyncTransport>, Address)> = Vec::new();
543 for addr in addresses {
544 if let Some(transport) = self.transport_manager.handle_for_address(addr) {
545 tasks.push((transport, addr.clone()));
546 }
547 }
548
549 if tasks.is_empty() {
550 return Err(
551 SyncError::InvalidAddress("No matching transport for any address".into()).into(),
552 );
553 }
554
555 let (tx, mut rx) = mpsc::channel(tasks.len());
556
557 for (transport, addr) in tasks {
558 let tx = tx.clone();
559 let mut ctx = self.handshake_ctx(&addr)?;
560 ctx.transport = Arc::clone(&transport);
561 let addr_info = addr.clone();
562 tokio::spawn(async move {
563 let result = tokio::time::timeout(
564 ADDRESS_ATTEMPT_TIMEOUT,
565 conn::run_handshake(ctx, addr.clone()),
566 )
567 .await;
568 let result = match result {
569 Ok(Ok(answered)) => Ok((transport, addr, answered)),
570 Ok(Err(error)) => Err(error),
571 Err(_) => {
572 warn!(
573 address = ?addr_info,
574 timeout = ?ADDRESS_ATTEMPT_TIMEOUT,
575 "Address attempt timed out",
576 );
577 Err(SyncError::Network(format!(
578 "Address attempt timed out after {ADDRESS_ATTEMPT_TIMEOUT:?}"
579 ))
580 .into())
581 }
582 };
583 let _ = tx.send(result).await;
584 });
585 }
586 drop(tx);
587
588 let mut last_err = None;
589 while let Some(result) = rx.recv().await {
590 match result {
591 Ok((transport, addr, answered)) if &answered == expected => {
592 return Ok((transport, addr));
593 }
594 Ok((_, addr, answered)) => {
595 warn!(
596 address = ?addr,
597 expected = %expected,
598 answered = %answered,
599 "Address answered as a different peer; not selecting it as the route"
600 );
601 last_err = Some(
602 SyncError::HandshakeFailed(format!(
603 "{addr:?} answered as {answered}, not {expected}"
604 ))
605 .into(),
606 );
607 }
608 Err(e) => last_err = Some(e),
609 }
610 }
611
612 Err(last_err.expect("at least one task was spawned"))
613 }
614
615 /// Send specific entries to a peer without duplicate filtering.
616 ///
617 /// This method performs direct entry transmission and is used by:
618 /// - `SendEntries` commands from the frontend (caller handles filtering)
619 /// - `sync_tree_with_peer()` after smart duplicate prevention analysis
620 ///
621 /// # Design Note
622 ///
623 /// This method does NOT perform duplicate prevention - that responsibility
624 /// lies with the caller. The background sync's smart duplicate prevention
625 /// happens in `sync_tree_with_peer()` via tip comparison, while direct
626 /// `SendEntries` commands trust the caller to send appropriate entries.
627 ///
628 /// # Error Handling
629 ///
630 /// Failed sends are automatically added to the retry queue with exponential backoff.
631 async fn send_to_peer(&self, peer: &PeerId, entries: Vec<Entry>) -> Result<()> {
632 // Get peer addresses from sync tree (extract and drop transaction before await)
633 let (addresses, request) = {
634 let sync_tree = self.get_sync_tree().await?;
635 let txn = sync_tree.new_transaction().await?;
636 let peer_info = PeerManager::new(&txn)
637 .get_peer_info(peer.public_key())
638 .await?
639 .ok_or_else(|| SyncError::PeerNotFound(peer.to_string()))?;
640
641 let addresses = peer_info.addresses.clone();
642 let request = SyncRequest::SendEntries(entries);
643 (addresses, request)
644 }; // Transaction is dropped here
645
646 let (transport, address) = self.select_route(&addresses, peer.public_key()).await?;
647 let response = transport.send_request(&address, &request).await?;
648
649 match response {
650 SyncResponse::Ack | SyncResponse::Count(_) => Ok(()),
651 SyncResponse::Error(msg) => Err(SyncError::SyncProtocolError(format!(
652 "Peer {peer} returned error: {msg}"
653 ))
654 .into()),
655 _ => Err(SyncError::UnexpectedResponse {
656 expected: "Ack or Count",
657 actual: format!("{response:?}"),
658 }
659 .into()),
660 }
661 }
662
663 /// Add failed send to retry queue
664 fn add_to_retry_queue(&mut self, peer: PeerId, entries: Vec<Entry>, error: Error, now_ms: u64) {
665 // Log send failure and add to retry queue
666 tracing::warn!("Failed to send to {peer}: {error}. Adding to retry queue.");
667 self.retry_queue.push(RetryEntry {
668 peer,
669 entries,
670 attempts: 1,
671 last_attempt_ms: now_ms,
672 });
673 }
674
675 /// Process entries from the sync queue, batching by peer.
676 ///
677 /// Drains the queue and sends entries to each peer. Failed sends
678 /// are added to the retry queue with exponential backoff.
679 /// Returns the number of peers that failed to receive entries.
680 async fn process_queue(&mut self) -> usize {
681 let batches = self.queue.drain();
682 if batches.is_empty() {
683 return 0;
684 }
685
686 let instance = match self.instance() {
687 Ok(i) => i,
688 Err(e) => {
689 tracing::warn!("Failed to get instance for queue processing: {e}");
690 return batches.len(); // All batches failed
691 }
692 };
693
694 let mut failures = 0;
695 for (peer, entry_ids) in batches {
696 // Fetch entries from backend
697 let mut entries = Vec::with_capacity(entry_ids.len());
698 for (entry_id, _tree_id) in &entry_ids {
699 match instance.backend().get(entry_id).await {
700 Ok(entry) => entries.push(entry),
701 Err(e) => {
702 tracing::warn!("Failed to fetch entry {entry_id} for peer {peer}: {e}");
703 }
704 }
705 }
706
707 if entries.is_empty() {
708 continue;
709 }
710
711 // Send batched entries to peer
712 if let Err(e) = self.send_to_peer(&peer, entries.clone()).await {
713 let now_ms = instance.clock().now_millis();
714 self.add_to_retry_queue(peer, entries, e, now_ms);
715 failures += 1;
716 }
717 }
718 failures
719 }
720
721 /// Process retry queue with exponential backoff
722 async fn process_retry_queue(&mut self) {
723 let now_ms = self.instance().map(|i| i.clock().now_millis()).unwrap_or(0);
724 let mut still_failed = Vec::new();
725
726 // Take the retry queue to avoid borrowing issues
727 let retry_queue = std::mem::take(&mut self.retry_queue);
728
729 // Process entries that are ready for retry
730 for mut entry in retry_queue {
731 // Backoff in milliseconds: 2^attempts * 1000ms, max 64 seconds
732 let backoff_ms = 2u64.pow(entry.attempts.min(6)) * 1000;
733 let elapsed_ms = now_ms.saturating_sub(entry.last_attempt_ms);
734
735 if elapsed_ms >= backoff_ms {
736 // Try sending again
737 if let Err(_e) = self.send_to_peer(&entry.peer, entry.entries.clone()).await {
738 entry.attempts += 1;
739 entry.last_attempt_ms = now_ms;
740
741 if entry.attempts < 10 {
742 // Max 10 attempts
743 still_failed.push(entry);
744 } else {
745 // Max retries exceeded - give up on this batch
746 tracing::error!("Giving up on sending to {} after 10 attempts", entry.peer);
747 }
748 } else {
749 // Successfully retried after failure
750 }
751 } else {
752 // Not ready for retry yet
753 still_failed.push(entry);
754 }
755 }
756
757 self.retry_queue = still_failed;
758 }
759
760 /// Flush retry queue immediately, ignoring backoff timers.
761 /// Returns the number of entries that still failed after retry.
762 async fn flush_retry_queue(&mut self) -> usize {
763 let mut still_failed = Vec::new();
764 let now_ms = self.instance().map(|i| i.clock().now_millis()).unwrap_or(0);
765
766 // Take the retry queue to process
767 let retry_queue = std::mem::take(&mut self.retry_queue);
768
769 // Try sending each entry immediately (ignore backoff)
770 for mut retry_entry in retry_queue {
771 if let Err(_e) = self
772 .send_to_peer(&retry_entry.peer, retry_entry.entries.clone())
773 .await
774 {
775 retry_entry.attempts += 1;
776 retry_entry.last_attempt_ms = now_ms;
777
778 if retry_entry.attempts < 10 {
779 still_failed.push(retry_entry);
780 } else {
781 tracing::error!(
782 "Giving up on sending to {} after 10 attempts",
783 retry_entry.peer
784 );
785 }
786 }
787 }
788
789 let failed_count = still_failed.len();
790 self.retry_queue = still_failed;
791 failed_count
792 }
793
794 /// Perform periodic sync with all active peers
795 async fn periodic_sync_all_peers(&self) {
796 // Periodic sync triggered
797
798 // Get all peers from sync tree
799 let peers = match self.get_sync_tree().await {
800 Ok(sync_tree) => match sync_tree.new_transaction().await {
801 Ok(txn) => match PeerManager::new(&txn).list_peers().await {
802 Ok(peers) => {
803 // Extract peer list and drop the operation before awaiting
804 peers
805 }
806 Err(_) => {
807 // Skip sync if we can't list peers
808 return;
809 }
810 },
811 Err(_) => {
812 // Skip sync if we can't create transaction
813 return;
814 }
815 },
816 Err(_) => {
817 // Skip sync if we can't get sync tree
818 return;
819 }
820 };
821
822 // Sync peers concurrently (the transaction is dropped, so no Send
823 // issues). A round is I/O-bound and dominated by peers that are down:
824 // done serially, one unreachable peer delays every peer behind it by a
825 // full connect timeout, so the round degrades with the number of dead
826 // peers rather than with the work there is to do. Concurrently, a round
827 // costs about as long as its slowest peer.
828 //
829 // Bounded so a large peer set can't open an unbounded number of
830 // connections at once. These futures are polled on this task rather
831 // than spawned, so they need no `Send`/`'static` bounds.
832 use futures_util::stream::StreamExt;
833 futures_util::stream::iter(
834 peers
835 .into_iter()
836 .filter(|p| p.status == PeerStatus::Active)
837 .map(|peer_info| async move {
838 if let Err(e) = self.sync_with_peer(&peer_info.id).await {
839 // Log individual peer sync failure but continue with others
840 tracing::error!("Periodic sync failed with {}: {e}", peer_info.id);
841 }
842 }),
843 )
844 .buffer_unordered(MAX_CONCURRENT_PEER_SYNCS)
845 .collect::<()>()
846 .await;
847 }
848
849 /// Sync with a specific peer (bidirectional)
850 async fn sync_with_peer(&self, peer_id: &PeerId) -> Result<()> {
851 async move {
852 info!(peer = %peer_id, "Starting peer synchronization");
853
854 // Get peer addresses and tree list from sync tree (extract and
855 // drop transaction before any network I/O).
856 let (addresses, sync_trees) = {
857 let sync_tree = self.get_sync_tree().await?;
858 let txn = sync_tree.new_transaction().await?;
859 let peer_manager = PeerManager::new(&txn);
860
861 let peer_info = peer_manager
862 .get_peer_info(peer_id.public_key())
863 .await?
864 .ok_or_else(|| SyncError::PeerNotFound(peer_id.to_string()))?;
865
866 let addresses = peer_info.addresses.clone();
867
868 // Find all trees that sync with this peer from sync tree
869 let sync_trees = peer_manager.get_peer_trees(peer_id.public_key()).await?;
870
871 (addresses, sync_trees)
872 }; // Transaction is dropped here
873
874 if sync_trees.is_empty() {
875 debug!(peer = %peer_id, "No trees configured for sync with peer");
876 return Ok(()); // No trees to sync
877 }
878
879 // Race bounded handshakes to select one route. Other successful
880 // handshakes may finish registration, but only the selected route
881 // continues into the tree exchange.
882 let (_transport, address) = self.select_route(&addresses, peer_id.public_key()).await?;
883
884 info!(peer = %peer_id, tree_count = sync_trees.len(), "Synchronizing trees with peer");
885
886 let tree_count = sync_trees.len();
887 let mut synced_any = false;
888 for (index, tree_id) in sync_trees.iter().enumerate() {
889 let Err(e) = self.sync_tree_with_peer(peer_id, tree_id, &address).await else {
890 synced_any = true;
891 continue;
892 };
893
894 tracing::error!("Failed to sync tree {tree_id} with peer {peer_id}: {e}");
895
896 // A failure that is about the peer rather than about this tree
897 // ends the walk. Every remaining tree would fail the same way,
898 // each paying a full deadline to establish what this one just
899 // established, so a peer that is down costs one round trip
900 // rather than one per tree registered against it.
901 //
902 // Any other failure is about this tree alone: the peer answered,
903 // so the trees behind it are still worth attempting. Keeping
904 // that distinction is what lets the walk stop early without a
905 // single broken tree stalling every tree behind it.
906 if e.is_network_error() {
907 debug!(
908 peer = %peer_id,
909 skipped = tree_count - index - 1,
910 "Peer stopped answering; ending the tree walk"
911 );
912 break;
913 }
914 }
915
916 // A round that moved at least one tree is the peer answering, which
917 // is the fact `SyncStatus.last_sync` reports. A round in which every
918 // tree failed is not, even though the walk itself returns `Ok`.
919 if synced_any && let Ok(instance) = self.instance() {
920 self.peer_state
921 .record_success(peer_id, instance.clock().now_millis());
922 }
923
924 info!(peer = %peer_id, "Completed peer synchronization");
925 Ok(())
926 }
927 .instrument(info_span!("sync_with_peer", peer = %peer_id))
928 .await
929 }
930
931 /// Sync a specific tree with a peer using smart duplicate prevention.
932 ///
933 /// This method implements Eidetica's core synchronization algorithm based on
934 /// Merkle-CRDT tip comparison. It eliminates duplicate sends by understanding
935 /// the semantic state of both peers' trees.
936 ///
937 /// # Algorithm
938 ///
939 /// 1. **Tip Exchange**: Get local tips and request peer's tips
940 /// 2. **Gap Analysis**: Compare tips to identify missing entries on both sides
941 /// 3. **Smart Transfer**: Only send/receive entries that are genuinely missing
942 /// 4. **DAG Completion**: Include all necessary ancestor entries
943 ///
944 /// # Benefits
945 ///
946 /// - **No duplicates**: Tips comparison guarantees no redundant network transfers
947 /// - **Complete data**: DAG traversal ensures all dependencies are satisfied
948 /// - **Bidirectional**: Both peers sync simultaneously for efficiency
949 /// - **Self-correcting**: Any missed entries are caught in subsequent syncs
950 ///
951 /// # Performance
952 ///
953 /// - **O(tip_count)** network requests for discovery
954 /// - **O(missing_entries)** data transfer (optimal)
955 /// - **Stateless**: No persistent tracking of individual sends needed
956 async fn sync_tree_with_peer(
957 &self,
958 peer_id: &PeerId,
959 tree_id: &ID,
960 address: &Address,
961 ) -> Result<()> {
962 async move {
963 trace!(peer = %peer_id, tree = %tree_id, "Starting unified tree synchronization");
964
965 // Get our tips for this tree (empty if tree doesn't exist)
966 let instance = self.instance()?;
967 let our_tips = instance
968 .backend()
969 .snapshot(tree_id)
970 .await
971 .map_err(|e| SyncError::BackendError(format!("Failed to get local tips: {e}")))?;
972
973 // Get our device public key for automatic peer tracking
974 let our_device_pubkey = Some(instance.id());
975
976 debug!(peer = %peer_id, tree = %tree_id, our_tips = our_tips.len(), "Sending sync tree request");
977
978 // Send unified sync request, signed so the peer can authorize the pull
979 let auth = SyncRequestAuth::sign(
980 instance.signing_key()?,
981 peer_id.public_key(),
982 tree_id,
983 &our_tips,
984 instance.clock().now_millis(),
985 );
986 let request = SyncRequest::SyncTree(SyncTreeRequest {
987 tree_id: tree_id.clone(),
988 our_tips,
989 peer_pubkey: our_device_pubkey,
990 requesting_key: None,
991 requesting_key_name: None,
992 requested_permission: None,
993 metadata: None,
994 auth: Some(auth),
995 });
996
997 let response = self.transport_manager.send_request(address, &request).await?;
998
999 match response {
1000 SyncResponse::Bootstrap(bootstrap_response) => {
1001 info!(peer = %peer_id, tree = %tree_id, entry_count = bootstrap_response.all_entries.len() + 1, "Received bootstrap response");
1002 self.handle_bootstrap_response(bootstrap_response).await?;
1003 }
1004 SyncResponse::Incremental(incremental_response) => {
1005 debug!(peer = %peer_id, tree = %tree_id,
1006 their_tips = incremental_response.their_tips.len(),
1007 missing_count = incremental_response.missing_entries.len(),
1008 "Received incremental sync response");
1009 self.handle_incremental_response(incremental_response).await?;
1010 }
1011 SyncResponse::Error(msg) => {
1012 return Err(SyncError::SyncProtocolError(format!("Sync error: {msg}")).into());
1013 }
1014 _ => {
1015 return Err(SyncError::UnexpectedResponse {
1016 expected: "Bootstrap or Incremental",
1017 actual: format!("{response:?}"),
1018 }.into());
1019 }
1020 }
1021
1022 trace!(peer = %peer_id, tree = %tree_id, "Completed unified tree synchronization");
1023 Ok(())
1024 }
1025 .instrument(info_span!("sync_tree", peer = %peer_id, tree = %tree_id))
1026 .await
1027 }
1028
1029 /// Check peer connections and attempt reconnection
1030 async fn check_peer_connections(&mut self) {
1031 // For now, this is a placeholder
1032 // In the future, we could implement connection health checks
1033 // and automatic reconnection logic here
1034 }
1035
1036 /// Start the sync server on specified or all transports
1037 async fn start_server(&mut self, name: Option<&str>) -> Result<()> {
1038 // Create a sync handler with instance access and sync tree ID
1039 let handler = Arc::new(SyncHandlerImpl::new(
1040 self.instance()?,
1041 self.sync_tree_id.clone(),
1042 ));
1043
1044 match name {
1045 Some(name) => {
1046 // Start server on specific transport
1047 self.transport_manager.start_server(name, handler).await?;
1048 tracing::info!("Sync server started for transport {name}");
1049 }
1050 None => {
1051 // Start servers on all transports
1052 self.transport_manager.start_all_servers(handler).await?;
1053 tracing::info!("Sync servers started for all transports");
1054 }
1055 }
1056
1057 Ok(())
1058 }
1059
1060 /// Stop the sync server on specified or all transports
1061 async fn stop_server(&mut self, name: Option<&str>) -> Result<()> {
1062 match name {
1063 Some(name) => {
1064 self.transport_manager.stop_server(name).await?;
1065 tracing::info!("Sync server stopped for transport {name}");
1066 }
1067 None => {
1068 self.transport_manager.stop_all_servers().await?;
1069 tracing::info!("All sync servers stopped");
1070 }
1071 }
1072 Ok(())
1073 }
1074}
1075
1076#[cfg(test)]
1077mod tests {
1078 use std::time::Duration;
1079
1080 use crate::{
1081 Error,
1082 sync::error::{SyncError, TimeoutPhase},
1083 };
1084
1085 /// The rule the tree walk applies to a failed tree: end the peer's round, or
1086 /// carry on to the trees behind it.
1087 fn ends_the_walk(err: SyncError) -> bool {
1088 Error::from(err).is_network_error()
1089 }
1090
1091 /// The peer itself stopped answering, so every remaining tree would pay a
1092 /// full deadline to establish what this one just established.
1093 #[test]
1094 fn a_peer_that_stopped_answering_ends_the_walk() {
1095 assert!(ends_the_walk(SyncError::ConnectionFailed {
1096 address: "peer:8080".to_string(),
1097 reason: "connection refused".to_string(),
1098 }));
1099 assert!(ends_the_walk(SyncError::Timeout {
1100 address: "peer:8080".to_string(),
1101 phase: TimeoutPhase::Connect,
1102 elapsed: Duration::from_secs(10),
1103 }));
1104 assert!(ends_the_walk(SyncError::Network(
1105 "connection reset by peer".to_string()
1106 )));
1107 }
1108
1109 /// A peer that connected and then went quiet is still a peer that stopped
1110 /// answering. The phase says whether it proved itself reachable, which is
1111 /// worth reporting, but it does not make the remaining trees worth
1112 /// attempting this round.
1113 #[test]
1114 fn a_peer_that_went_quiet_after_connecting_also_ends_the_walk() {
1115 assert!(ends_the_walk(SyncError::Timeout {
1116 address: "peer:8080".to_string(),
1117 phase: TimeoutPhase::Request,
1118 elapsed: Duration::from_secs(30),
1119 }));
1120 }
1121
1122 /// The peer answered in every one of these: the failure is about the one
1123 /// tree, so the trees behind it are still worth attempting. Ending the walk
1124 /// on these would let a single misconfigured tree stall every tree
1125 /// registered behind it against a peer that is perfectly healthy.
1126 #[test]
1127 fn a_failure_about_one_tree_does_not_end_the_walk() {
1128 assert!(!ends_the_walk(SyncError::PermissionDenied(
1129 "key has no write access".to_string()
1130 )));
1131 assert!(!ends_the_walk(SyncError::AuthenticationFailed(
1132 "signature did not verify".to_string()
1133 )));
1134 assert!(!ends_the_walk(SyncError::UnexpectedResponse {
1135 expected: "SyncResponse",
1136 actual: "Error".to_string(),
1137 }));
1138 assert!(!ends_the_walk(SyncError::SyncProtocolError(
1139 "malformed tip set".to_string()
1140 )));
1141 }
1142}