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}