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

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}