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

eidetica/sync/
ops.rs

1//! Core sync operations for the sync system.
2
3use std::future::Future;
4use std::time::Duration;
5
6use tokio::sync::oneshot;
7use tracing::{debug, info, warn};
8
9use super::{
10    Address, DatabaseTicket, PeerId, Sync, SyncError,
11    background::SyncCommand,
12    peer_manager::PeerManager,
13    peer_types,
14    protocol::{self, SyncRequest, SyncRequestAuth, SyncResponse, SyncTreeRequest},
15    user_sync_manager::UserSyncManager,
16};
17use crate::{
18    Database, Entry, Result,
19    auth::Permission,
20    auth::crypto::{PrivateKey, PublicKey},
21    crdt::Doc,
22    entry::ID,
23    store::Table,
24    user::types::UserInfo,
25};
26
27use super::utils::collect_ancestors_to_send;
28
29impl Sync {
30    // === Core Sync Methods ===
31
32    /// Synchronize a specific tree with a peer using bidirectional sync.
33    ///
34    /// This is the main synchronization method that implements tip exchange
35    /// and bidirectional entry transfer to keep trees in sync between peers.
36    /// It performs both pull (fetch missing entries) and push (send our entries).
37    ///
38    /// # Arguments
39    /// * `peer_pubkey` - The public key of the peer to sync with
40    /// * `tree_id` - The ID of the tree to synchronize
41    ///
42    /// # Returns
43    /// A Result indicating success or failure of the sync operation.
44    pub async fn sync_tree_with_peer(&self, peer_pubkey: &PublicKey, tree_id: &ID) -> Result<()> {
45        self.sync_tree_with_peer_as(peer_pubkey, tree_id, None)
46            .await
47    }
48
49    /// Synchronize a tree with a peer, signing the request with `signing_key`.
50    ///
51    /// The peer authorizes the pull against the key that signed it, so this must
52    /// be a key with read access on `tree_id`. `None` uses this instance's device
53    /// key, which is right when the device itself holds the access (directly or
54    /// through a delegated tree). A database that instead granted a *user* key —
55    /// the usual outcome of bootstrapping with one — needs that key here.
56    pub async fn sync_tree_with_peer_as(
57        &self,
58        peer_pubkey: &PublicKey,
59        tree_id: &ID,
60        signing_key: Option<&PrivateKey>,
61    ) -> Result<()> {
62        let addresses = self.peer_addresses(peer_pubkey).await?;
63        let peer = peer_pubkey.clone();
64        let tree = tree_id.clone();
65        let key = signing_key.cloned();
66        self.with_selected_address(&addresses, peer_pubkey, move |sync, addr| {
67            let peer = peer.clone();
68            let tree = tree.clone();
69            let key = key.clone();
70            async move {
71                sync.sync_tree_with_peer_at(&addr, &peer, &tree, key.as_ref())
72                    .await
73            }
74        })
75        .await
76    }
77
78    /// One address's attempt at [`Self::sync_tree_with_peer_as`].
79    async fn sync_tree_with_peer_at(
80        &self,
81        address: &Address,
82        peer_pubkey: &PublicKey,
83        tree_id: &ID,
84        signing_key: Option<&PrivateKey>,
85    ) -> Result<()> {
86        // Get our current tips for this tree (empty if tree doesn't exist)
87        let backend = self.backend()?;
88        let our_tips = backend
89            .snapshot(tree_id)
90            .await
91            .map_err(|e| SyncError::BackendError(format!("Failed to get local tips: {e}")))?;
92
93        // Get our device public key for automatic peer tracking
94        let our_device_pubkey = self.get_device_pubkey().ok();
95
96        // Send unified sync request, signed so the peer can authorize the pull
97        let instance = self.instance()?;
98        let auth = SyncRequestAuth::sign(
99            signing_key.unwrap_or(instance.signing_key()?),
100            peer_pubkey,
101            tree_id,
102            &our_tips,
103            instance.clock().now_millis(),
104        );
105        let request = SyncRequest::SyncTree(SyncTreeRequest {
106            tree_id: tree_id.clone(),
107            our_tips,
108            peer_pubkey: our_device_pubkey,
109            requesting_key: None,
110            requesting_key_name: None,
111            requested_permission: None,
112            metadata: None,
113            auth: Some(auth),
114        });
115
116        // Send request via background sync command
117        let (tx, rx) = oneshot::channel();
118        self.background_tx
119            .get()
120            .ok_or(SyncError::NoTransportEnabled)?
121            .send(SyncCommand::SendRequest {
122                address: address.clone(),
123                request: Box::new(request),
124                response: tx,
125            })
126            .await
127            .map_err(|e| SyncError::CommandSendError(e.to_string()))?;
128
129        let response = rx
130            .await
131            .map_err(|e| SyncError::Network(format!("Response channel error: {e}")))?
132            .map_err(|e| SyncError::Network(format!("Request failed: {e}")))?;
133
134        match response {
135            SyncResponse::Bootstrap(bootstrap_response) => {
136                self.handle_bootstrap_response(bootstrap_response).await?;
137            }
138            SyncResponse::Incremental(incremental_response) => {
139                self.handle_incremental_response(incremental_response, address)
140                    .await?;
141            }
142            SyncResponse::Error(msg) => {
143                return Err(SyncError::SyncProtocolError(format!("Sync error: {msg}")).into());
144            }
145            _ => {
146                return Err(SyncError::UnexpectedResponse {
147                    expected: "Bootstrap or Incremental",
148                    actual: format!("{response:?}"),
149                }
150                .into());
151            }
152        }
153
154        // Track tree/peer relationship for sync_on_commit to work
155        // This allows on_local_write() to find this peer when queueing entries
156        self.add_tree_sync(peer_pubkey, tree_id).await?;
157
158        Ok(())
159    }
160
161    /// Handle bootstrap response by storing root and all entries
162    pub(super) async fn handle_bootstrap_response(
163        &self,
164        response: protocol::BootstrapResponse,
165    ) -> Result<()> {
166        tracing::info!(tree_id = %response.tree_id, "Processing bootstrap response");
167
168        // Integrity check: the root entry's content must hash to the declared
169        // tree_id. Rejects peers serving substituted content. A mismatch here
170        // also covers the cross-algorithm bootstrap case (e.g. a SHA-256 tree
171        // advertised to a BLAKE3-default node) — those are unsupported until
172        // the backend gains multi-CID-per-entry storage, so failing loudly is
173        // better than silently re-keying the DAG under the wrong algorithm.
174        let derived = response.root_entry.id();
175        if derived != response.tree_id {
176            return Err(SyncError::InvalidEntry(format!(
177                "root entry content hashes to {} but bootstrap response declares tree_id {}",
178                derived, response.tree_id
179            ))
180            .into());
181        }
182
183        // Combine root entry with all other entries into a single batch
184        let mut all_entries = Vec::with_capacity(1 + response.all_entries.len());
185        all_entries.push(response.root_entry);
186        all_entries.extend(response.all_entries);
187
188        // Store all entries and fire callbacks once
189        self.store_received_entries(&response.tree_id, all_entries)
190            .await?;
191
192        tracing::info!(tree_id = %response.tree_id, "Bootstrap completed successfully");
193        Ok(())
194    }
195
196    /// Handle incremental response by storing missing entries and sending back what server is missing
197    pub(super) async fn handle_incremental_response(
198        &self,
199        response: protocol::IncrementalResponse,
200        peer_address: &peer_types::Address,
201    ) -> Result<()> {
202        tracing::debug!(tree_id = %response.tree_id, "Processing incremental response");
203
204        // Step 1: Store missing entries
205        self.store_received_entries(&response.tree_id, response.missing_entries)
206            .await?;
207
208        // Step 2: Check if server is missing entries from us
209        let backend = self.backend()?;
210        let our_snapshot = backend.snapshot(&response.tree_id).await?;
211        let their_tips = &response.their_tips;
212
213        // Find tips they don't have
214        let missing_tip_ids: Vec<_> = our_snapshot
215            .tips()
216            .iter()
217            .filter(|tip_id| !their_tips.contains(tip_id))
218            .cloned()
219            .collect();
220
221        if !missing_tip_ids.is_empty() {
222            tracing::debug!(
223                tree_id = %response.tree_id,
224                missing_tips = missing_tip_ids.len(),
225                "Server is missing some of our entries, sending them back"
226            );
227
228            // Collect entries server is missing
229            let engine = self
230                .backend()?
231                .local_engine()
232                .expect("sync requires local backend");
233            let entries_for_server =
234                collect_ancestors_to_send(engine.as_ref(), &missing_tip_ids, their_tips).await?;
235
236            if !entries_for_server.is_empty() {
237                // Send these entries back to server
238                self.send_missing_entries_to_peer(
239                    peer_address,
240                    &response.tree_id,
241                    entries_for_server,
242                )
243                .await?;
244            }
245        }
246
247        tracing::debug!(tree_id = %response.tree_id, "Incremental sync completed");
248        Ok(())
249    }
250
251    /// Send entries that the server is missing back to complete bidirectional sync
252    async fn send_missing_entries_to_peer(
253        &self,
254        peer_address: &peer_types::Address,
255        tree_id: &ID,
256        entries: Vec<Entry>,
257    ) -> Result<()> {
258        if entries.is_empty() {
259            return Ok(());
260        }
261
262        tracing::debug!(
263            tree_id = %tree_id,
264            entry_count = entries.len(),
265            "Sending missing entries back to peer for bidirectional sync"
266        );
267
268        let request = protocol::SyncRequest::SendEntries(entries);
269
270        // Send via command channel
271        let (tx, rx) = tokio::sync::oneshot::channel();
272        self.background_tx
273            .get()
274            .ok_or(SyncError::NoTransportEnabled)?
275            .send(SyncCommand::SendRequest {
276                address: peer_address.clone(),
277                request: Box::new(request),
278                response: tx,
279            })
280            .await
281            .map_err(|e| SyncError::CommandSendError(e.to_string()))?;
282
283        // Wait for acknowledgment
284        let response = rx
285            .await
286            .map_err(|e| SyncError::Network(format!("Response channel error: {e}")))?
287            .map_err(|e| SyncError::Network(format!("Request failed: {e}")))?;
288
289        match response {
290            protocol::SyncResponse::Ack | protocol::SyncResponse::Count(_) => {
291                tracing::debug!(tree_id = %tree_id, "Server acknowledged receipt of missing entries");
292                Ok(())
293            }
294            protocol::SyncResponse::Error(e) => {
295                Err(SyncError::Network(format!("Server error receiving entries: {e}")).into())
296            }
297            _ => Err(SyncError::UnexpectedResponse {
298                expected: "Ack or Count",
299                actual: format!("{response:?}"),
300            }
301            .into()),
302        }
303    }
304
305    /// Validate and store received entries from a peer, firing remote write callbacks.
306    pub(super) async fn store_received_entries(
307        &self,
308        tree_id: &ID,
309        entries: Vec<Entry>,
310    ) -> Result<()> {
311        // These entries arrive without per-entry declared IDs — they were batched
312        // under a single tree_id by the sender. Content is stored under whatever
313        // ID our local `entry.id()` derives, so substitution attacks on individual
314        // entries would fail DAG connectivity checks via parent pointers rather
315        // than a per-entry hash check here. Root-level integrity is verified by
316        // the bootstrap handler against the declared tree_id.
317        //
318        // TODO: Add signature verification and parent-existence / DAG-connectivity
319        // checks before marking entries as verified.
320
321        // Store entries and fire callbacks via Instance::put_remote_entries.
322        // Stored Unverified: these arrived from a peer and have not been
323        // verified by this node.
324        let instance = self.instance()?;
325        instance
326            .put_remote_entries(tree_id, entries)
327            .await
328            .map_err(|e| SyncError::BackendError(format!("Failed to store entries: {e}")))?;
329
330        Ok(())
331    }
332
333    /// Send a batch of entries to a sync peer (async version).
334    ///
335    /// # Arguments
336    /// * `entries` - The entries to send
337    /// * `address` - The address of the peer to send to
338    ///
339    /// # Returns
340    /// A Result indicating whether the entries were successfully acknowledged.
341    pub async fn send_entries(
342        &self,
343        entries: impl AsRef<[Entry]>,
344        address: &Address,
345    ) -> Result<()> {
346        let entries_vec = entries.as_ref().to_vec();
347        let request = SyncRequest::SendEntries(entries_vec);
348        let response = self.send_request(&request, address).await?;
349
350        match response {
351            SyncResponse::Ack | SyncResponse::Count(_) => Ok(()),
352            SyncResponse::Error(msg) => Err(SyncError::SyncProtocolError(format!(
353                "Peer {} returned error: {}",
354                address.address, msg
355            ))
356            .into()),
357            _ => Err(SyncError::UnexpectedResponse {
358                expected: "Ack or Count",
359                actual: format!("{response:?}"),
360            }
361            .into()),
362        }
363    }
364
365    /// Send specific entries to a peer via the background sync engine.
366    ///
367    /// This method queues entries for direct transmission without duplicate filtering.
368    /// The caller is responsible for determining which entries should be sent.
369    ///
370    /// # Duplicate Prevention Architecture
371    ///
372    /// Eidetica uses **smart duplicate prevention** in the background sync engine:
373    /// - **Database sync** (`SyncWithPeer` command): Uses tip comparison for semantic filtering
374    /// - **Direct send** (this method): Trusts caller to provide appropriate entries
375    ///
376    /// For automatic duplicate prevention, use tree-based sync relationships instead
377    /// of calling this method directly.
378    ///
379    /// # Arguments
380    /// * `peer_id` - The peer ID to send to
381    /// * `entries` - The specific entries to send (no filtering applied)
382    ///
383    /// # Returns
384    /// A Result indicating whether the command was successfully queued for background processing.
385    pub async fn send_entries_to_peer(&self, peer_id: &PeerId, entries: Vec<Entry>) -> Result<()> {
386        self.background_tx
387            .get()
388            .ok_or(SyncError::NoTransportEnabled)?
389            .send(SyncCommand::SendEntries {
390                peer: peer_id.clone(),
391                entries,
392            })
393            .await
394            .map_err(|e| SyncError::CommandSendError(e.to_string()))?;
395        Ok(())
396    }
397
398    /// Queue an entry for sync to a peer (non-blocking, for use in callbacks).
399    ///
400    /// This method is designed for use in write callbacks where async operations
401    /// are not possible. It uses try_send to avoid blocking, and logs errors
402    /// rather than failing the callback.
403    ///
404    /// # Arguments
405    /// * `peer_pubkey` - The public key of the peer to sync with
406    /// * `entry_id` - The ID of the entry to queue
407    /// * `tree_id` - The tree ID where the entry belongs
408    ///
409    /// # Returns
410    /// Ok(()) if the entry was successfully queued.
411    /// Only returns Err if transport is not enabled.
412    pub fn queue_entry_for_sync(
413        &self,
414        peer_id: &PeerId,
415        entry_id: &ID,
416        tree_id: &ID,
417    ) -> Result<()> {
418        // Ensure background sync is running
419        if self.background_tx.get().is_none() {
420            return Err(SyncError::NoTransportEnabled.into());
421        }
422
423        // Add to queue - BackgroundSync will process and send
424        self.queue
425            .enqueue(peer_id, entry_id.clone(), tree_id.clone());
426
427        Ok(())
428    }
429
430    /// Handle local write events for automatic sync.
431    ///
432    /// This method is called by the Instance write callback system when entries
433    /// are committed locally. It looks up the combined sync settings for the database
434    /// and queues the entry for sync with all configured peers if sync is enabled.
435    ///
436    /// This is the core method that implements automatic sync-on-commit behavior.
437    ///
438    /// # Arguments
439    /// * `event` - The write event containing the newly committed entries
440    /// * `database` - The database where the entries were committed
441    ///
442    /// # Returns
443    /// Ok(()) on success, or an error if settings lookup fails
444    pub(crate) async fn on_local_write(
445        &self,
446        event: &crate::instance::WriteEvent,
447        database: &Database,
448    ) -> Result<()> {
449        // Early return if background sync not running
450        if self.background_tx.get().is_none() {
451            return Ok(());
452        }
453
454        // Look up combined settings for this database
455        let tx = self.sync_tree.new_transaction().await?;
456        let user_mgr = UserSyncManager::new(&tx);
457        let peer_mgr = PeerManager::new(&tx);
458
459        let combined_settings = match user_mgr.get_combined_settings(database.root_id()).await? {
460            Some(settings) => settings,
461            None => {
462                // No settings configured for this database - no sync needed
463                debug!(database_id = %database.root_id(), "No sync settings for database, skipping");
464                return Ok(());
465            }
466        };
467
468        // Check if sync is enabled and sync_on_commit is true
469        if !combined_settings.sync_enabled || !combined_settings.sync_on_commit {
470            debug!(
471                database_id = %database.root_id(),
472                sync_enabled = combined_settings.sync_enabled,
473                sync_on_commit = combined_settings.sync_on_commit,
474                "Sync not enabled for database"
475            );
476            return Ok(());
477        }
478
479        // Get list of peers for this database
480        let peers = peer_mgr.get_tree_peers(database.root_id()).await?;
481
482        if peers.is_empty() {
483            debug!(database_id = %database.root_id(), "No peers configured for database");
484            return Ok(());
485        }
486
487        // Queue each entry for sync with each peer. The event no longer
488        // carries entry payloads — expand the cursor advance back into a
489        // concrete set of IDs by walking the DAG diff. Cost is bounded
490        // by the cursor delta, not the full tree.
491        let tree_id = database.root_id();
492
493        let new_ids = database
494            .ids_added(event.previous_tips(), event.post_tips())
495            .await?;
496
497        for entry_id in &new_ids {
498            debug!(
499                database_id = %tree_id,
500                entry_id = %entry_id,
501                peer_count = peers.len(),
502                "Queueing entry for automatic sync"
503            );
504
505            for peer_id in &peers {
506                self.queue_entry_for_sync(peer_id, entry_id, tree_id)?;
507            }
508        }
509
510        Ok(())
511    }
512
513    /// Initialize combined settings for all users.
514    ///
515    /// This is called during Sync initialization. For new sync trees (just created),
516    /// it scans the _users database to register all existing users. For existing
517    /// sync trees (loaded), it updates combined settings for already-tracked users.
518    pub(super) async fn initialize_user_settings(&self) -> Result<()> {
519        self.reconcile_user_settings().await
520    }
521
522    /// Reconcile the persisted user directory and every tracked preferences tree.
523    ///
524    /// The daemon calls this at startup and whenever a service client changes
525    /// either source tree. It is intentionally idempotent: `sync_user` compares
526    /// the persisted preferences cursor before updating combined settings.
527    pub(crate) async fn reconcile_user_settings(&self) -> Result<()> {
528        let _guard = self.reconciliation.lock().await;
529
530        let instance = self.instance.upgrade().ok_or(SyncError::InstanceDropped)?;
531        let users_db = instance.users_db().await?;
532        let users_table = users_db
533            .get_store_viewer::<Table<UserInfo>>("users")
534            .await?;
535
536        // Always scan `_users`, not only when the sync tree is empty. A daemon
537        // may already track some users when another account is created through
538        // the service socket.
539        for (user_uuid, user_info) in users_table.search(|_| true).await? {
540            self.sync_user(&user_uuid, &user_info.user_database_id)
541                .await?;
542        }
543
544        Ok(())
545    }
546
547    /// Whether a local write can change the persisted user sync intent.
548    pub(crate) async fn is_reconciliation_source(&self, database_id: &ID) -> Result<bool> {
549        let instance = self.instance.upgrade().ok_or(SyncError::InstanceDropped)?;
550        if instance.users_db_id() == database_id {
551            return Ok(true);
552        }
553        let users_db = instance.users_db().await?;
554        let users_table = users_db
555            .get_store_viewer::<Table<UserInfo>>("users")
556            .await?;
557
558        Ok(users_table
559            .search(|user| &user.user_database_id == database_id)
560            .await?
561            .into_iter()
562            .next()
563            .is_some())
564    }
565
566    /// Send a sync request to a peer and get a response (async version).
567    ///
568    /// # Arguments
569    /// * `request` - The sync request to send
570    /// * `address` - The address of the peer
571    ///
572    /// # Returns
573    /// The sync response from the peer.
574    pub(super) async fn send_request(
575        &self,
576        request: &SyncRequest,
577        address: &Address,
578    ) -> Result<SyncResponse> {
579        let (tx, rx) = oneshot::channel();
580
581        self.background_tx
582            .get()
583            .ok_or(SyncError::NoTransportEnabled)?
584            .send(SyncCommand::SendRequest {
585                address: address.clone(),
586                request: Box::new(request.clone()),
587                response: tx,
588            })
589            .await
590            .map_err(|e| SyncError::CommandSendError(e.to_string()))?;
591
592        rx.await
593            .map_err(|e| SyncError::Network(format!("Response channel error: {e}")))?
594    }
595
596    /// Discover available trees from a peer (simplified API).
597    ///
598    /// This method connects to a peer and retrieves the list of trees they're willing to sync.
599    /// This is useful for discovering what can be synced before setting up sync relationships.
600    ///
601    /// # Arguments
602    /// * `address` - The transport address of the peer.
603    ///
604    /// # Returns
605    /// A vector of TreeInfo describing available trees, or an error.
606    pub async fn discover_peer_trees(&self, address: &Address) -> Result<Vec<protocol::TreeInfo>> {
607        // Connect and get handshake info
608        let _peer_pubkey = self.connect_to_peer(address).await?;
609
610        // The handshake already contains the tree list, but we need to get it again
611        // since connect_to_peer doesn't return it. For now, return empty list
612        // TODO: Enhance this to actually return the tree list from handshake
613
614        tracing::warn!(
615            "discover_peer_trees not fully implemented - handshake contains tree info but API needs enhancement"
616        );
617        Ok(vec![])
618    }
619
620    /// Sync with a peer at a given address.
621    ///
622    /// This is a blocking convenience method that:
623    /// 1. Connects to discover the peer's public key
624    /// 2. Registers the peer and performs immediate sync
625    /// 3. Returns after sync completes
626    ///
627    /// For new code, prefer using [`register_sync_peer()`](Self::register_sync_peer)
628    /// directly, which registers intent and lets background sync handle it.
629    ///
630    /// # Arguments
631    /// * `address` - The transport address of the peer.
632    /// * `tree_id` - Optional tree ID to sync (None = discover available trees)
633    ///
634    /// # Returns
635    /// Result indicating success or failure.
636    pub async fn sync_with_peer(&self, address: &Address, tree_id: Option<&ID>) -> Result<()> {
637        self.sync_with_peer_as(address, tree_id, None).await
638    }
639
640    /// Sync with a peer, signing requests with `signing_key`.
641    ///
642    /// See [`sync_tree_with_peer_as`](Self::sync_tree_with_peer_as) for which
643    /// key to pass.
644    pub async fn sync_with_peer_as(
645        &self,
646        address: &Address,
647        tree_id: Option<&ID>,
648        signing_key: Option<&PrivateKey>,
649    ) -> Result<()> {
650        // Connect to peer if not already connected
651        let peer_pubkey = self.connect_to_peer(address).await?;
652
653        // Store the address for this peer (needed for sync_tree_with_peer)
654        self.add_peer_address(&peer_pubkey, address.clone()).await?;
655
656        if let Some(tree_id) = tree_id {
657            // Sync specific tree
658            self.sync_tree_with_peer_as(&peer_pubkey, tree_id, signing_key)
659                .await?;
660        } else {
661            // TODO: Sync all available trees
662            tracing::warn!(
663                "Syncing all trees not yet implemented - need to enhance discover_peer_trees first"
664            );
665        }
666
667        Ok(())
668    }
669
670    /// Sync with a peer using a [`DatabaseTicket`].
671    ///
672    /// Races bounded handshakes against every address hint, then performs the
673    /// tree exchange once through the first usable route. Other successful
674    /// handshakes may finish peer registration in the background.
675    ///
676    /// # Arguments
677    /// * `ticket` - A ticket containing the database ID and address hints.
678    ///
679    /// # Errors
680    /// Returns [`SyncError::InvalidAddress`] if the ticket has no address hints.
681    /// Returns the last sync error if no address succeeded.
682    pub async fn sync_with_ticket(&self, ticket: &DatabaseTicket) -> Result<()> {
683        let database_id = ticket.database_id().clone();
684        let (address, peer_pubkey) = self.select_address(ticket.addresses(), None).await?;
685        self.add_peer_address(&peer_pubkey, address.clone()).await?;
686        self.sync_tree_with_peer_at(&address, &peer_pubkey, &database_id, None)
687            .await
688    }
689
690    /// Sync a specific tree with a peer, with optional authentication for bootstrap.
691    ///
692    /// This is a lower-level method that allows specifying authentication parameters
693    /// for bootstrap scenarios where access needs to be requested.
694    ///
695    /// # Arguments
696    /// * `peer_pubkey` - The public key of the peer to sync with
697    /// * `tree_id` - The ID of the tree to sync
698    /// * `requesting_key` - Optional private key to sign with and request access for
699    /// * `requesting_key_name` - Optional name/ID of the requesting key
700    /// * `requested_permission` - Optional permission level being requested
701    ///
702    /// # Returns
703    /// A Result indicating success or failure.
704    pub async fn sync_tree_with_peer_auth(
705        &self,
706        peer_pubkey: &PublicKey,
707        tree_id: &ID,
708        requesting_key: Option<&PrivateKey>,
709        requesting_key_name: Option<&str>,
710        requested_permission: Option<Permission>,
711        metadata: Option<Doc>,
712    ) -> Result<()> {
713        let addresses = self.peer_addresses(peer_pubkey).await?;
714        let peer = peer_pubkey.clone();
715        let tree = tree_id.clone();
716        let key = requesting_key.cloned();
717        let key_name = requesting_key_name.map(str::to_string);
718        self.with_selected_address(&addresses, peer_pubkey, move |sync, addr| {
719            let peer = peer.clone();
720            let tree = tree.clone();
721            let key = key.clone();
722            let key_name = key_name.clone();
723            let metadata = metadata.clone();
724            async move {
725                sync.sync_tree_with_peer_auth_at(
726                    &addr,
727                    &peer,
728                    &tree,
729                    key.as_ref(),
730                    key_name.as_deref(),
731                    requested_permission,
732                    metadata,
733                )
734                .await
735            }
736        })
737        .await
738    }
739
740    /// One address's attempt at [`Self::sync_tree_with_peer_auth`].
741    #[allow(clippy::too_many_arguments)]
742    pub(super) async fn sync_tree_with_peer_auth_at(
743        &self,
744        address: &Address,
745        peer_pubkey: &PublicKey,
746        tree_id: &ID,
747        requesting_key: Option<&PrivateKey>,
748        requesting_key_name: Option<&str>,
749        requested_permission: Option<Permission>,
750        metadata: Option<Doc>,
751    ) -> Result<()> {
752        // Get our current tips for this tree (empty if tree doesn't exist)
753        let backend = self.backend()?;
754        let our_tips = backend
755            .snapshot(tree_id)
756            .await
757            .map_err(|e| SyncError::BackendError(format!("Failed to get local tips: {e}")))?;
758
759        // Get our device public key for automatic peer tracking
760        let our_device_pubkey = self.get_device_pubkey().ok();
761
762        // Send the request with proof from the named requesting key. Calls
763        // without a named-key request still sign with the device key for
764        // authenticated data access where required.
765        let instance = self.instance()?;
766        let signing_key = requesting_key.unwrap_or(instance.signing_key()?);
767        let auth = SyncRequestAuth::sign(
768            signing_key,
769            peer_pubkey,
770            tree_id,
771            &our_tips,
772            instance.clock().now_millis(),
773        );
774        let request = SyncRequest::SyncTree(SyncTreeRequest {
775            tree_id: tree_id.clone(),
776            our_tips,
777            peer_pubkey: our_device_pubkey,
778            requesting_key: requesting_key.map(|k| k.public_key()),
779            requesting_key_name: requesting_key_name.map(|k| k.to_string()),
780            requested_permission,
781            metadata,
782            auth: Some(auth),
783        });
784
785        // Send request via background sync command
786        let (tx, rx) = oneshot::channel();
787        self.background_tx
788            .get()
789            .ok_or(SyncError::NoTransportEnabled)?
790            .send(SyncCommand::SendRequest {
791                address: address.clone(),
792                request: Box::new(request),
793                response: tx,
794            })
795            .await
796            .map_err(|_| {
797                SyncError::CommandSendError("Background sync command channel closed".to_string())
798            })?;
799
800        // Wait for response
801        let response = rx
802            .await
803            .map_err(|_| {
804                SyncError::CommandSendError("Background sync response channel closed".to_string())
805            })?
806            .map_err(|e| SyncError::Network(format!("Sync request failed: {e}")))?;
807
808        // Handle the response (same logic as existing sync_tree_with_peer)
809        match response {
810            SyncResponse::Bootstrap(bootstrap_response) => {
811                info!(peer = %peer_pubkey, tree = %tree_id, entry_count = bootstrap_response.all_entries.len() + 1, "Received bootstrap response");
812
813                // Store root + all entries as a single batch with callback dispatch
814                let mut all_entries = Vec::with_capacity(1 + bootstrap_response.all_entries.len());
815                all_entries.push(bootstrap_response.root_entry);
816                all_entries.extend(bootstrap_response.all_entries);
817
818                // Bootstrap entries come from a peer; stored Unverified.
819                let instance = self.instance()?;
820                instance.put_remote_entries(tree_id, all_entries).await?;
821
822                info!(peer = %peer_pubkey, tree = %tree_id, "Bootstrap sync completed successfully");
823            }
824            SyncResponse::Incremental(incremental_response) => {
825                info!(peer = %peer_pubkey, tree = %tree_id, missing_count = incremental_response.missing_entries.len(), "Received incremental sync response");
826
827                // Use the enhanced handler that supports bidirectional sync
828                self.handle_incremental_response(incremental_response, address)
829                    .await?;
830
831                debug!(peer = %peer_pubkey, tree = %tree_id, "Incremental sync completed");
832            }
833            SyncResponse::BootstrapPending {
834                request_id,
835                message,
836            } => {
837                info!(peer = %peer_pubkey, tree = %tree_id, request_id = %request_id, "Bootstrap request pending manual approval");
838                return Err(SyncError::BootstrapPending {
839                    request_id,
840                    message,
841                }
842                .into());
843            }
844            SyncResponse::BootstrapRejected {
845                request_id,
846                message,
847            } => {
848                info!(peer = %peer_pubkey, tree = %tree_id, request_id = %request_id, "Bootstrap request was rejected");
849                return Err(SyncError::BootstrapRejected {
850                    request_id,
851                    message,
852                }
853                .into());
854            }
855            SyncResponse::Error(err) => {
856                return Err(SyncError::Network(format!("Peer returned error: {err}")).into());
857            }
858            _ => {
859                return Err(SyncError::SyncProtocolError(
860                    "Unexpected response type for sync tree request".to_string(),
861                )
862                .into());
863            }
864        }
865
866        // Track tree/peer relationship for sync_on_commit to work
867        // This allows on_local_write() to find this peer when queueing entries
868        self.add_tree_sync(peer_pubkey, tree_id).await?;
869
870        Ok(())
871    }
872
873    // === Flush Operations ===
874
875    /// Process all queued entries and retry any failed sends.
876    ///
877    /// This method:
878    /// 1. Retries all entries in the retry queue (ignoring backoff timers)
879    /// 2. Processes all entries in the sync queue (batched by peer)
880    ///
881    /// When this method returns, all pending sync work has been attempted.
882    /// This is useful to eensuree that all pending pushes have completed.
883    ///
884    /// # Returns
885    /// `Ok(())` if all operations completed successfully, or an error
886    /// if the background sync engine is not running or sends failed.
887    pub async fn flush(&self) -> Result<()> {
888        let (tx, rx) = oneshot::channel();
889
890        self.background_tx
891            .get()
892            .ok_or(SyncError::NoTransportEnabled)?
893            .send(SyncCommand::Flush { response: tx })
894            .await
895            .map_err(|e| SyncError::CommandSendError(e.to_string()))?;
896
897        rx.await
898            .map_err(|e| SyncError::Network(format!("Response channel error: {e}")))?
899    }
900
901    /// Every address recorded for `peer_pubkey`.
902    async fn peer_addresses(&self, peer_pubkey: &PublicKey) -> Result<Vec<Address>> {
903        let peer_info = self
904            .get_peer_info(peer_pubkey)
905            .await?
906            .ok_or_else(|| SyncError::PeerNotFound(peer_pubkey.to_string()))?;
907        if peer_info.addresses.is_empty() {
908            return Err(
909                SyncError::Network(format!("No addresses found for peer {peer_pubkey}")).into(),
910            );
911        }
912        Ok(peer_info.addresses)
913    }
914
915    /// Select a route with bounded handshakes, then run the real operation once.
916    ///
917    /// A peer's address list only ever grows: anything that changes address on
918    /// restart appends a new entry and leaves the old one in place. Dialing just
919    /// one of them makes a peer that is up and reachable permanently unreachable
920    /// as soon as the entry that happens to be dialed goes stale — and the
921    /// failure presents as a timeout, which reads as "the peer is down", the one
922    /// diagnosis that leads away from the cause.
923    ///
924    /// Ticket bootstrap has always raced its address hints. This puts every
925    /// subsequent sync on the same footing, which makes an accumulated list
926    /// harmless rather than fatal. Pruning dead addresses is a separate concern
927    /// and is not needed for reachability once all of them are tried.
928    ///
929    /// On total failure the last error is returned unchanged — callers match on
930    /// specific variants (`BootstrapPending`, for one), so it must not be
931    /// flattened into a connectivity error. The attempted addresses are logged
932    /// instead, since a bare timeout naming no address is not actionable.
933    async fn with_selected_address<F, Fut>(
934        &self,
935        addresses: &[Address],
936        peer_pubkey: &PublicKey,
937        f: F,
938    ) -> Result<()>
939    where
940        F: Fn(Sync, Address) -> Fut,
941        Fut: Future<Output = Result<()>> + Send + 'static,
942    {
943        let result = match self.select_address(addresses, Some(peer_pubkey)).await {
944            Ok((address, _)) => f(self.clone(), address).await,
945            Err(error) => Err(error),
946        };
947        if let Err(e) = &result {
948            warn!(
949                peer = %peer_pubkey,
950                attempted = ?addresses,
951                error = %e,
952                "No address answered for this peer"
953            );
954        }
955        result
956    }
957
958    /// Timeout applied to each address attempt in
959    /// [`select_address`](Self::select_address).
960    const ADDRESS_ATTEMPT_TIMEOUT: Duration = Duration::from_secs(30);
961
962    /// Race registration handshakes and return the first usable address, along
963    /// with the identity that answered on it.
964    ///
965    /// Spawns one detached task per address via [`tokio::spawn`]. Remaining
966    /// tasks are **not** cancelled — they continue running in the background so
967    /// that additional addresses can be registered for future syncs. The real
968    /// operation runs separately and is therefore not bounded by this timeout.
969    ///
970    /// `expected` is the peer the caller means to reach, when it knows. A
971    /// handshake says who actually answered, and a route that answers as
972    /// somebody else is rejected rather than selected: a peer's address list
973    /// only ever grows, so a stale entry can be reoccupied by an unrelated node
974    /// that handshakes perfectly well, and selecting it would both misdirect the
975    /// exchange and stop the working address behind it from ever being tried.
976    /// Ticket paths pass `None` — a ticket's whole purpose is to learn an
977    /// identity the caller does not yet have.
978    ///
979    /// If all tasks fail the last error is returned. If `addresses` is empty
980    /// an [`SyncError::InvalidAddress`] error is returned.
981    pub(super) async fn select_address(
982        &self,
983        addresses: &[Address],
984        expected: Option<&PublicKey>,
985    ) -> Result<(Address, PublicKey)> {
986        if addresses.is_empty() {
987            return Err(SyncError::InvalidAddress("Ticket has no address hints".into()).into());
988        }
989
990        let (tx, mut rx) = tokio::sync::mpsc::channel(addresses.len());
991
992        for addr in addresses {
993            let tx = tx.clone();
994            let sync = self.clone();
995            let addr = addr.clone();
996            let addr_info = addr.clone();
997            // Detached spawn: the task keeps running even after we return.
998            tokio::spawn(async move {
999                let result = tokio::time::timeout(
1000                    Self::ADDRESS_ATTEMPT_TIMEOUT,
1001                    sync.connect_to_peer(&addr),
1002                )
1003                .await;
1004                let result = match result {
1005                    Ok(Ok(peer_pubkey)) => Ok((addr, peer_pubkey)),
1006                    Ok(Err(error)) => Err(error),
1007                    Err(_) => {
1008                        warn!(
1009                            address = ?addr_info,
1010                            "Address attempt timed out after {:?}",
1011                            Self::ADDRESS_ATTEMPT_TIMEOUT,
1012                        );
1013                        Err(SyncError::Network(format!(
1014                            "Address attempt timed out after {:?}",
1015                            Self::ADDRESS_ATTEMPT_TIMEOUT,
1016                        ))
1017                        .into())
1018                    }
1019                };
1020                match &result {
1021                    Ok(_) => debug!(address = ?addr_info, "Address attempt succeeded"),
1022                    Err(e) => debug!(address = ?addr_info, error = %e, "Address attempt failed"),
1023                }
1024                // Ignore send errors — the receiver is dropped on early success,
1025                // but the task still completes its work (peer registration, etc.).
1026                let _ = tx.send(result).await;
1027            });
1028        }
1029        // Drop our sender so the channel closes when all tasks finish.
1030        drop(tx);
1031
1032        let mut last_err = None;
1033        while let Some(result) = rx.recv().await {
1034            match result {
1035                Ok((addr, answered)) => match expected {
1036                    Some(want) if want != &answered => {
1037                        warn!(
1038                            address = ?addr,
1039                            expected = %want,
1040                            answered = %answered,
1041                            "Address answered as a different peer; not selecting it as the route"
1042                        );
1043                        last_err = Some(
1044                            SyncError::HandshakeFailed(format!(
1045                                "{addr:?} answered as {answered}, not {want}"
1046                            ))
1047                            .into(),
1048                        );
1049                    }
1050                    _ => return Ok((addr, answered)),
1051                },
1052                Err(e) => last_err = Some(e),
1053            }
1054        }
1055
1056        Err(last_err.expect("at least one task was spawned"))
1057    }
1058}