eidetica/sync/handler.rs
1//! Sync request handler trait and implementation.
2//!
3//! This module contains transport-agnostic handlers that process
4//! sync requests and generate responses. These handlers can be
5//! used by any transport implementation through the SyncHandler trait.
6
7use std::{
8 collections::{BTreeMap, HashMap},
9 sync::Mutex,
10};
11
12use async_trait::async_trait;
13use tracing::{Instrument, debug, error, info, info_span, trace, warn};
14
15use super::{
16 bootstrap_request_manager::{BootstrapRequest, BootstrapRequestManager, RequestStatus},
17 peer_manager::PeerManager,
18 peer_types::Address,
19 protocol::{
20 BootstrapResponse, HandshakeRequest, HandshakeResponse, IncrementalResponse,
21 PROTOCOL_VERSION, RequestContext, SyncRequest, SyncResponse, SyncTreeRequest,
22 },
23 user_sync_manager::UserSyncManager,
24};
25use crate::{
26 Database, Entry, Error, Instance, Result, WeakInstance,
27 auth::{
28 Permission,
29 crypto::{PublicKey, create_challenge_response, generate_challenge},
30 },
31 entry::ID,
32 store::SettingsStore,
33 sync::error::SyncError,
34};
35
36/// Trait for handling sync requests with database access.
37///
38/// Implementations of this trait can process sync requests and generate
39/// appropriate responses, with full access to the database backend for
40/// storing and retrieving entries.
41#[async_trait]
42pub trait SyncHandler: Send + std::marker::Sync {
43 /// Handle a sync request and generate an appropriate response.
44 ///
45 /// This is the main entry point for processing sync messages,
46 /// regardless of which transport they arrived through.
47 ///
48 /// # Arguments
49 /// * `request` - The sync request to process
50 /// * `context` - Context about the request (remote address, etc.)
51 ///
52 /// # Returns
53 /// The appropriate response for the given request.
54 async fn handle_request(&self, request: &SyncRequest, context: &RequestContext)
55 -> SyncResponse;
56}
57
58/// How far a request's timestamp may sit from our clock, in either direction.
59///
60/// Bounds how long a captured signature stays useful. Peers with clocks further
61/// apart than this cannot sync — the window trades clock tolerance against
62/// replay exposure.
63const MAX_REQUEST_AGE_MS: u64 = 60_000;
64
65/// Default implementation of SyncHandler with database backend access.
66pub struct SyncHandlerImpl {
67 instance: WeakInstance,
68 sync_tree_id: ID,
69 /// Nonces spent within the freshness window, keyed by the claiming key.
70 ///
71 /// Makes each signature single-use: without this, anyone who observes a
72 /// signed request can replay it verbatim until it ages out.
73 spent_nonces: Mutex<HashMap<(PublicKey, Vec<u8>), u64>>,
74}
75
76impl SyncHandlerImpl {
77 /// Create a new SyncHandlerImpl with the given instance.
78 ///
79 /// # Arguments
80 /// * `instance` - Database instance for storing and retrieving entries
81 /// * `sync_tree_id` - Root ID of the sync database for storing bootstrap requests
82 pub fn new(instance: Instance, sync_tree_id: ID) -> Self {
83 Self {
84 instance: instance.downgrade(),
85 sync_tree_id,
86 spent_nonces: Mutex::new(HashMap::new()),
87 }
88 }
89
90 /// Authenticate a request that is about to be served data.
91 ///
92 /// Returns the key the caller has proven it holds. Checks, in order:
93 /// the signature covers *this* request to *us*, the request is fresh, and
94 /// its nonce has not been spent.
95 ///
96 /// This does not decide authorization — see [`Database::can_access`].
97 fn authenticate_request(&self, request: &SyncTreeRequest) -> Result<PublicKey> {
98 let key = self.verify_request_signature(request)?;
99 let auth = request.auth.as_ref().expect("signature was verified");
100 let now = self.instance()?.clock().now_millis();
101
102 // The map holds one entry per verified request in the last window, so
103 // it is bounded by request rate, not by uptime. Cap it if a peer can
104 // ever outrun that.
105 let mut spent = self
106 .spent_nonces
107 .lock()
108 .unwrap_or_else(|poisoned| poisoned.into_inner());
109 spent.retain(|_, spent_at| now.abs_diff(*spent_at) <= MAX_REQUEST_AGE_MS);
110 if spent
111 .insert((auth.key.clone(), auth.nonce.clone()), auth.timestamp_ms)
112 .is_some()
113 {
114 return Err(
115 SyncError::AuthenticationFailed("request nonce already spent".to_string()).into(),
116 );
117 }
118
119 Ok(key)
120 }
121
122 /// Verify possession without consuming the request nonce.
123 ///
124 /// Manual bootstrap retries prove who is asking but disclose no data. The
125 /// nonce is consumed only if the request later reaches a path that serves
126 /// entries.
127 fn verify_request_signature(&self, request: &SyncTreeRequest) -> Result<PublicKey> {
128 let auth = request
129 .auth
130 .as_ref()
131 .ok_or_else(|| SyncError::AuthenticationRequired(request.tree_id.to_string()))?;
132 let instance = self.instance()?;
133 auth.verify(&instance.id(), &request.tree_id, &request.our_tips)
134 .map_err(|_| {
135 SyncError::AuthenticationFailed("invalid request signature".to_string())
136 })?;
137 let now = instance.clock().now_millis();
138 if now.abs_diff(auth.timestamp_ms) > MAX_REQUEST_AGE_MS {
139 return Err(SyncError::AuthenticationFailed(
140 "request timestamp outside the freshness window".to_string(),
141 )
142 .into());
143 }
144 Ok(auth.key.clone())
145 }
146
147 /// Whether this request may be served entries from `tree_id`.
148 ///
149 /// A database with no auth configured, or with a global `*` grant, is
150 /// world-readable by design and needs no credentials. Otherwise the caller
151 /// must prove it holds a key that [`Database::can_access`] accepts —
152 /// directly, through the global grant, or through a delegated tree.
153 async fn authorize_read(&self, request: &SyncTreeRequest) -> Result<()> {
154 if !self.check_if_database_has_auth(&request.tree_id).await? {
155 return Ok(());
156 }
157
158 let key = self.authenticate_request(request)?;
159 if !Database::can_access(&self.instance()?, &request.tree_id, &key, &Permission::Read)
160 .await?
161 {
162 return Err(SyncError::PermissionDenied(format!(
163 "key {key} is not authorized to read {}",
164 request.tree_id
165 ))
166 .into());
167 }
168
169 Ok(())
170 }
171
172 /// Upgrade the weak instance reference to a strong reference.
173 pub(super) fn instance(&self) -> Result<Instance> {
174 self.instance
175 .upgrade()
176 .ok_or_else(|| SyncError::InstanceDropped.into())
177 }
178
179 /// Get access to the sync tree for bootstrap request management.
180 ///
181 /// # Returns
182 /// A Database instance for the sync tree with device key authentication.
183 async fn get_sync_tree(&self) -> Result<Database> {
184 // Load sync tree with the device key
185 let instance = self.instance()?;
186 let signing_key = instance.signing_key()?.clone();
187 Ok(Database::open(&instance, &self.sync_tree_id)
188 .await?
189 .with_key(signing_key))
190 }
191
192 /// Resolve an explicit bootstrap request under the lifecycle lock.
193 ///
194 /// Rejected history wins over current authority; otherwise live authority
195 /// approves immediately and a missing grant creates or resumes a pending
196 /// record.
197 async fn resolve_bootstrap_request(
198 &self,
199 sync_request: &SyncTreeRequest,
200 requesting_key_name: &str,
201 ) -> Result<BootstrapRequestOutcome> {
202 let tree_id = &sync_request.tree_id;
203 let requesting_key = sync_request
204 .requesting_key
205 .as_ref()
206 .expect("explicit bootstrap request has a requesting key");
207 let requested_permission = sync_request
208 .requested_permission
209 .expect("explicit bootstrap request has a requested permission");
210 let sync_tree = self.get_sync_tree().await?;
211 let instance = self.instance()?;
212 let lock = instance.tree_lock(&self.sync_tree_id);
213 let guard = lock.lock_owned().await;
214 let txn = sync_tree.new_transaction().await?;
215 let manager = BootstrapRequestManager::new(&txn);
216
217 if let Some((request_id, request)) = manager
218 .find_existing_request(tree_id, requesting_key, &requested_permission)
219 .await?
220 {
221 match request.status {
222 RequestStatus::Rejected { .. } => {
223 return Ok(BootstrapRequestOutcome::Rejected(request_id));
224 }
225 RequestStatus::Pending => {
226 return Ok(BootstrapRequestOutcome::Pending(request_id));
227 }
228 RequestStatus::Approved { .. } => {
229 if Database::can_access(
230 &instance,
231 tree_id,
232 requesting_key,
233 &requested_permission,
234 )
235 .await?
236 {
237 return Ok(BootstrapRequestOutcome::Approved);
238 }
239 }
240 }
241 }
242
243 if Database::can_access(&instance, tree_id, requesting_key, &requested_permission).await? {
244 return Ok(BootstrapRequestOutcome::Approved);
245 }
246
247 let request = BootstrapRequest {
248 tree_id: tree_id.clone(),
249 requesting_pubkey: requesting_key.clone(),
250 requesting_key_name: requesting_key_name.to_string(),
251 requested_permission,
252 timestamp: self.instance()?.clock().now_rfc3339(),
253 status: RequestStatus::Pending,
254 // TODO: We need to get the actual peer address from the transport layer
255 // For now, use a placeholder that will need to be fixed when implementing notifications
256 peer_address: Address {
257 transport_type: "unknown".to_string(),
258 address: "unknown".to_string(),
259 },
260 // TODO(bootstrap-metadata-bound): `metadata` is unbounded, remote-supplied
261 // data persisted to the system sync tree before any approval decision. A
262 // hostile peer can spam large/many pending requests (storage/DoS on the
263 // multi-tenant boundary). Bound the metadata size here, and cap pending
264 // requests per peer, before this is exposed to untrusted peers. Note the
265 // `Doc` is already deserialized at the protocol layer ahead of the
266 // sync-enabled gate, so the size cap ideally belongs there too.
267 metadata: sync_request.metadata.clone(),
268 };
269
270 let request_id = manager.store_request(request).await?;
271 txn.commit_under_tree_lock(guard).await?;
272
273 Ok(BootstrapRequestOutcome::Pending(request_id))
274 }
275}
276
277enum BootstrapRequestOutcome {
278 Pending(String),
279 Rejected(String),
280 Approved,
281}
282
283#[async_trait]
284impl SyncHandler for SyncHandlerImpl {
285 async fn handle_request(
286 &self,
287 request: &SyncRequest,
288 context: &RequestContext,
289 ) -> SyncResponse {
290 match request {
291 SyncRequest::Handshake(handshake_req) => {
292 debug!("Received handshake request");
293 self.handle_handshake(handshake_req, context).await
294 }
295 SyncRequest::SyncTree(sync_req) => {
296 debug!(tree_id = %sync_req.tree_id, tips_count = sync_req.our_tips.len(), "Received sync tree request");
297 self.handle_sync_tree(sync_req, context).await
298 }
299 SyncRequest::SendEntries(entries) => {
300 // Process and store the received entries
301 let count = entries.len();
302 info!(count = count, "Received entries for synchronization");
303
304 let instance = match self.instance() {
305 Ok(i) => i,
306 Err(e) => return SyncResponse::Error(format!("Instance dropped: {e}")),
307 };
308
309 // Group entries by tree_id so we can fire callbacks per-database.
310 // BTreeMap so iteration order is deterministic (sorted by id);
311 // sender order within a tree is preserved by per-tree push order.
312 //
313 // Root entries declare an empty `tree.root` and act as their
314 // own tree_id. Non-root entries always carry a tree_id;
315 // well-formed peers should never send a non-root entry with
316 // `root() == None`. If they do, the entry ends up filed under
317 // its own id and parent-existence checks downstream reject it.
318 let mut by_tree: BTreeMap<ID, Vec<Entry>> = BTreeMap::new();
319 for entry in entries {
320 let tree_id = entry.root().unwrap_or_else(|| entry.id());
321 by_tree.entry(tree_id).or_default().push(entry.clone());
322 }
323
324 let mut stored_count = 0usize;
325 for (tree_id, tree_entries) in by_tree {
326 let batch_size = tree_entries.len();
327 // Entries arrive over the wire without per-entry signature
328 // verification; `put_remote_entries` stores them
329 // Unverified so a future re-verification pass can promote
330 // them.
331 match instance.put_remote_entries(&tree_id, tree_entries).await {
332 Ok(n) => {
333 stored_count += n;
334 debug!(tree_id = %tree_id, requested = batch_size, stored = n, "Stored entries");
335 }
336 Err(e) => {
337 error!(tree_id = %tree_id, error = %e, "Failed to store entries batch");
338 }
339 }
340 }
341
342 debug!(
343 received = count,
344 stored = stored_count,
345 "Completed entry synchronization"
346 );
347 if count <= 1 {
348 SyncResponse::Ack
349 } else {
350 SyncResponse::Count(stored_count)
351 }
352 }
353 }
354 }
355}
356
357impl SyncHandlerImpl {
358 /// Get the highest permission level a key has in the database's auth settings.
359 ///
360 /// This looks up all permissions the key has (direct + global wildcard) and returns
361 /// the highest one. Used for auto-detecting permissions during bootstrap.
362 ///
363 /// # Arguments
364 /// * `tree_id` - The database/tree ID to check auth settings for
365 /// * `requesting_pubkey` - The public key to look up
366 ///
367 /// # Returns
368 /// - `Ok(Some(Permission))` if key has any permissions
369 /// - `Ok(None)` if key not found in auth settings
370 /// - `Err` if database access fails
371 async fn get_key_highest_permission(
372 &self,
373 tree_id: &ID,
374 requesting_pubkey: &PublicKey,
375 ) -> Result<Option<Permission>> {
376 let database = Database::open(&self.instance()?, tree_id).await?;
377 let transaction = database.new_transaction().await?;
378 let settings_store = SettingsStore::new(&transaction)?;
379 let auth_settings = settings_store.auth_snapshot().await?;
380
381 let results = auth_settings.find_all_sigkeys_for_pubkey(requesting_pubkey);
382
383 if results.is_empty() {
384 return Ok(None);
385 }
386
387 // Results are sorted highest first, so take the first one
388 Ok(Some(results[0].1))
389 }
390
391 /// Check that the caller signed this request with `claimed_key`.
392 ///
393 /// Named-key bootstrap requests enter a persistent lifecycle and may expose
394 /// its status, so possession is required even when the database is public.
395 /// This validation is non-consuming; the nonce is spent only if entries are
396 /// served.
397 fn prove_possession(&self, request: &SyncTreeRequest, claimed_key: &PublicKey) -> Result<()> {
398 let proven = self.verify_request_signature(request)?;
399 if proven != *claimed_key {
400 return Err(SyncError::AuthenticationFailed(format!(
401 "request is signed by {proven}, which is not the claimed key {claimed_key}"
402 ))
403 .into());
404 }
405 Ok(())
406 }
407
408 /// Check if a database requires authentication for unauthenticated requests.
409 ///
410 /// This method checks if the database requires authentication for bootstrap requests
411 /// that don't provide credentials. A database allows unauthenticated access if:
412 /// 1. It has no auth settings configured at all (empty auth), OR
413 /// 2. It has a global `*` permission configured that allows unauthenticated access
414 ///
415 /// # Arguments
416 /// * `tree_id` - The database/tree ID to check auth configuration for
417 ///
418 /// # Returns
419 /// - `Ok(true)` if database requires authentication (has auth but no global permission)
420 /// - `Ok(false)` if database allows unauthenticated access (no auth or has global permission)
421 /// - `Err` if the check fails
422 async fn check_if_database_has_auth(&self, tree_id: &ID) -> Result<bool> {
423 let database = Database::open(&self.instance()?, tree_id).await?;
424 let transaction = database.new_transaction().await?;
425 let settings_store = SettingsStore::new(&transaction)?;
426
427 let auth_settings = settings_store.auth_snapshot().await?;
428
429 // Check if auth settings is completely empty (no auth configured)
430 if auth_settings.as_doc().is_empty() {
431 debug!(
432 tree_id = %tree_id,
433 "Database has no auth configured - allowing unauthenticated access"
434 );
435 return Ok(false); // No auth required
436 }
437
438 // Auth is configured - check if there's an Active global permission
439 if let Ok(global_key) = auth_settings.get_global_key()
440 && global_key.is_active()
441 {
442 debug!(
443 tree_id = %tree_id,
444 global_permission = ?global_key.permissions(),
445 "Database has global permission - allowing unauthenticated access"
446 );
447 return Ok(false); // Global permission allows unauthenticated access
448 }
449
450 // Auth is configured but no global permission - require authentication
451 debug!(
452 tree_id = %tree_id,
453 "Database has auth configured without global permission - requiring authentication"
454 );
455 Ok(true) // Auth required
456 }
457
458 /// Check if a database has sync enabled by at least one user.
459 ///
460 /// This is a security-critical check that determines if a database should accept
461 /// any sync requests at all. A database is only eligible for sync if at least one
462 /// user has it in their preferences with `sync_enabled: true`.
463 ///
464 /// # Security
465 /// This method implements fail-closed behavior:
466 /// - Returns `false` on any error (no information leakage)
467 /// - Returns `false` if no users have the database in preferences
468 /// - Returns `false` if combined_settings.sync_enabled is false
469 /// - Only returns `true` if explicitly enabled
470 ///
471 /// # Arguments
472 /// * `tree_id` - The ID of the database to check
473 ///
474 /// # Returns
475 /// `true` if the database has sync enabled, `false` otherwise (including errors)
476 async fn is_database_sync_enabled(&self, tree_id: &ID) -> bool {
477 let instance = match self.instance() {
478 Ok(i) => i,
479 Err(_) => return false, // Fail closed
480 };
481
482 let signing_key = match instance.signing_key() {
483 Ok(k) => k.clone(),
484 Err(_) => return false, // Fail closed
485 };
486
487 let sync_database = match Database::open(&instance, &self.sync_tree_id).await {
488 Ok(db) => db.with_key(signing_key),
489 Err(_) => return false, // Fail closed
490 };
491
492 let transaction = match sync_database.new_transaction().await {
493 Ok(tx) => tx,
494 Err(_) => return false, // Fail closed
495 };
496
497 // Use UserSyncManager to get combined settings
498 let user_mgr = UserSyncManager::new(&transaction);
499 match user_mgr.get_combined_settings(tree_id).await {
500 Ok(Some(settings)) => settings.sync_enabled,
501 _ => false, // Fail closed: no settings or error
502 }
503 }
504
505 /// Register an incoming peer and add their addresses to the peer list.
506 ///
507 /// This method registers a peer that initiated a connection to us during handshake.
508 /// It adds both the peer-advertised addresses and the transport-provided remote address.
509 ///
510 /// # Arguments
511 /// * `peer_pubkey` - The peer's public key
512 /// * `display_name` - Optional display name for the peer
513 /// * `advertised_addresses` - Addresses the peer advertised in their handshake
514 /// * `remote_address` - The actual address from which the connection originated
515 ///
516 /// # Returns
517 /// Result indicating success or failure of registration
518 async fn register_incoming_peer(
519 &self,
520 peer_pubkey: &PublicKey,
521 display_name: Option<&str>,
522 advertised_addresses: &[Address],
523 remote_address: &Option<Address>,
524 ) -> Result<()> {
525 let sync_tree = self.get_sync_tree().await?;
526 let txn = sync_tree.new_transaction().await?;
527 let peer_manager = PeerManager::new(&txn);
528
529 // Try to register the peer (ignore if already exists)
530 match peer_manager.register_peer(peer_pubkey, display_name).await {
531 Ok(()) => {
532 info!(peer_pubkey = %peer_pubkey, "Registered new incoming peer");
533 }
534 Err(Error::Sync(ref e)) if matches!(**e, SyncError::PeerAlreadyExists(_)) => {
535 debug!(peer_pubkey = %peer_pubkey, "Peer already registered, updating addresses");
536 }
537 Err(e) => return Err(e),
538 }
539
540 // Add all advertised addresses
541 for addr in advertised_addresses {
542 if let Err(e) = peer_manager.add_address(peer_pubkey, addr.clone()).await {
543 warn!(peer_pubkey = %peer_pubkey, address = ?addr, error = %e, "Failed to add advertised address");
544 }
545 }
546
547 // Add the remote address from transport if available
548 if let Some(addr) = remote_address
549 && let Err(e) = peer_manager.add_address(peer_pubkey, addr.clone()).await
550 {
551 warn!(peer_pubkey = %peer_pubkey, address = ?addr, error = %e, "Failed to add remote address");
552 }
553
554 txn.commit().await?;
555 Ok(())
556 }
557
558 /// Track tree/peer sync relationship when a peer requests a tree.
559 ///
560 /// This method adds the tree to the peer's sync list, enabling bidirectional
561 /// sync for the requested tree. This is critical for `sync_on_commit` to work
562 /// in both directions.
563 ///
564 /// # Arguments
565 /// * `tree_id` - The ID of the tree being requested
566 /// * `peer_pubkey` - The public key of the peer requesting the tree (device key, not auth key)
567 ///
568 /// # Returns
569 /// Result indicating success or failure
570 async fn track_tree_sync_relationship(
571 &self,
572 tree_id: &ID,
573 peer_pubkey: &PublicKey,
574 ) -> Result<()> {
575 let sync_tree = self.get_sync_tree().await?;
576 let txn = sync_tree.new_transaction().await?;
577 let peer_manager = PeerManager::new(&txn);
578
579 // Add the tree sync relationship
580 peer_manager.add_tree_sync(peer_pubkey, tree_id).await?;
581 txn.commit().await?;
582
583 debug!(tree_id = %tree_id, peer_pubkey = %peer_pubkey, "Tracked tree/peer sync relationship");
584 Ok(())
585 }
586
587 /// Handle a handshake request from a peer.
588 async fn handle_handshake(
589 &self,
590 request: &HandshakeRequest,
591 context: &RequestContext,
592 ) -> SyncResponse {
593 async move {
594 debug!(
595 peer_device_id = %request.device_id,
596 peer_public_key = %request.public_key,
597 display_name = ?request.display_name,
598 protocol_version = request.protocol_version,
599 "Processing handshake request"
600 );
601
602 // Check protocol version compatibility
603 if request.protocol_version != PROTOCOL_VERSION {
604 warn!(
605 expected = PROTOCOL_VERSION,
606 received = request.protocol_version,
607 "Protocol version mismatch"
608 );
609 return SyncResponse::Error(format!(
610 "Protocol version mismatch: expected {}, got {}",
611 PROTOCOL_VERSION, request.protocol_version
612 ));
613 }
614
615 // Get device signing key from backend
616 let instance = match self.instance() {
617 Ok(i) => i,
618 Err(e) => {
619 error!(error = %e, "Failed to get instance");
620 return SyncResponse::Error(format!("Failed to get instance: {e}"));
621 }
622 };
623 let signing_key = match instance.signing_key() {
624 Ok(k) => k.clone(),
625 Err(e) => {
626 error!(error = %e, "Failed to get device key");
627 return SyncResponse::Error(format!("Failed to get device key: {e}"));
628 }
629 };
630
631 // Generate device ID and public key from signing key
632 let public_key = signing_key.public_key();
633 let device_id = public_key.clone(); // Device ID is the public key
634
635 // Sign the challenge with our device key to prove identity
636 let challenge_response = create_challenge_response(&request.challenge, &signing_key);
637
638 // Generate a new challenge for mutual authentication
639 let new_challenge = generate_challenge();
640
641 // Get available trees for discovery
642 let available_trees = self.get_available_trees().await;
643
644 // Register the peer and add their addresses to our peer list
645 match self.register_incoming_peer(&request.public_key, request.display_name.as_deref(), &request.listen_addresses, &context.remote_address).await {
646 Ok(()) => {
647 debug!(peer_pubkey = %request.public_key, "Successfully registered incoming peer");
648 }
649 Err(e) => {
650 // Log the error but don't fail the handshake - peer registration is best-effort
651 warn!(peer_pubkey = %request.public_key, error = %e, "Failed to register incoming peer");
652 }
653 }
654
655 info!(
656 our_device_id = %device_id,
657 peer_device_id = %request.device_id,
658 tree_count = available_trees.len(),
659 "Handshake completed successfully"
660 );
661
662 SyncResponse::Handshake(HandshakeResponse {
663 device_id,
664 public_key,
665 display_name: Some("Eidetica Peer".to_string()),
666 protocol_version: PROTOCOL_VERSION,
667 challenge_response,
668 new_challenge,
669 available_trees,
670 })
671 }
672 .instrument(info_span!("handle_handshake", peer = %request.device_id))
673 .await
674 }
675
676 /// Handle a unified sync tree request (bootstrap or incremental).
677 ///
678 /// This method routes between two sync modes:
679 /// 1. **Bootstrap**: When peer has no tips (empty database), sends complete tree
680 /// 2. **Incremental**: When peer has existing tips, sends only new entries
681 ///
682 /// # Bootstrap Authentication
683 /// During bootstrap, if the peer provides authentication credentials:
684 /// - `requesting_key`: Public key to add
685 /// - `requesting_key_name`: Name for the key
686 /// - `requested_permission`: Access level requested
687 /// - `auth`: Proof that the requester holds `requesting_key`
688 ///
689 /// The handler will evaluate the bootstrap policy and either:
690 /// - Auto-approve a proven key with existing authority
691 /// - Store a proven request for manual approval
692 /// - Proceed anonymously only when no named-key request is made and the
693 /// database is public
694 async fn handle_sync_tree(
695 &self,
696 request: &SyncTreeRequest,
697 context: &RequestContext,
698 ) -> SyncResponse {
699 async move {
700 trace!(tree_id = %request.tree_id, "Processing sync tree request");
701
702 // Track tree/peer sync relationship for bidirectional sync
703 // IMPORTANT: Only use context.peer_pubkey (device key from handshake)
704 // Do NOT use request.requesting_key (that's an auth key for database access)
705 if let Some(peer_pubkey) = &context.peer_pubkey {
706 if let Err(e) = self.track_tree_sync_relationship(&request.tree_id, peer_pubkey).await {
707 // Log the error but don't fail the sync - relationship tracking is best-effort
708 warn!(tree_id = %request.tree_id, peer_pubkey = %peer_pubkey, error = %e, "Failed to track tree/peer relationship");
709 }
710 } else {
711 debug!(tree_id = %request.tree_id, "No peer pubkey in context, skipping relationship tracking");
712 }
713
714 // Check if peer needs bootstrap (empty tips indicates no local data)
715 if request.our_tips.is_empty() {
716 debug!(tree_id = %request.tree_id, "Peer needs bootstrap - sending full tree");
717 return self.handle_bootstrap_request(request).await;
718 }
719
720 // Handle incremental sync (peer has existing data, needs updates)
721 debug!(tree_id = %request.tree_id, peer_tips = request.our_tips.len(), "Handling incremental sync");
722 self.handle_incremental_sync(request).await
723 }
724 .instrument(info_span!("handle_sync_tree", tree = %request.tree_id))
725 .await
726 }
727
728 /// Handle bootstrap request by sending complete tree state and optionally approving auth key.
729 ///
730 /// Bootstrap is the initial synchronization when a peer has no local data for a tree.
731 /// This method:
732 /// 1. Validates the tree exists and sync is enabled
733 /// 2. Processes authentication and permission resolution
734 /// 3. Sends all entries from the tree to the peer
735 ///
736 /// # Authentication Flow
737 ///
738 /// The bootstrap process handles three authentication scenarios:
739 ///
740 /// ## 1. Explicit Permission Request
741 /// When all three request parameters are provided (`requesting_key`, `requesting_key_name`, `requested_permission`):
742 /// - Verify `auth` proves possession of `requesting_key`
743 /// - Check if the key already has sufficient permissions
744 /// - If yes: Approve immediately without adding a key
745 /// - If no: Store the request for manual approval and return `BootstrapPending`
746 ///
747 /// ## 2. Auto-Detection
748 /// When key is provided but `requested_permission` is `None`:
749 /// - Look up key's existing permissions in database auth settings
750 /// - Uses `find_all_sigkeys_for_pubkey()` to find all permissions (direct + global wildcard)
751 /// - If key found: Use highest available permission and approve immediately
752 /// - If key not found: Reject with authentication error
753 ///
754 /// ## 3. Anonymous Access
755 /// When no named-key request is made:
756 /// - Only allowed if the database has no auth configured or has a global wildcard permission
757 /// - Otherwise rejected with authentication required error
758 ///
759 /// # Arguments
760 /// * `tree_id` - The database/tree to bootstrap
761 /// * `requesting_key` - Optional public key requesting access; possession must be proven
762 /// * `requesting_key_name` - Optional name/identifier for the key
763 /// * `requested_permission` - Optional permission level requested (if None, auto-detects from auth settings)
764 ///
765 /// # Returns
766 /// - `BootstrapResponse`: Contains entries and approval status (key_approved, granted_permission)
767 /// - `BootstrapPending`: Manual approval required (request queued)
768 /// - `Error`: Tree not found, auth required, key not authorized, or processing failure
769 async fn handle_bootstrap_request(&self, request: &SyncTreeRequest) -> SyncResponse {
770 let tree_id = &request.tree_id;
771 let requesting_key = request.requesting_key.as_ref();
772 let requesting_key_name = request.requesting_key_name.as_deref();
773 let requested_permission = request.requested_permission;
774
775 // SECURITY: Check if database has sync enabled (FIRST CHECK - before anything else)
776 // This prevents information leakage about database existence: the gate
777 // returns false both for databases that are absent and for databases that
778 // are present-but-not-tracked-for-sync, and we deliberately respond with
779 // the same opaque "Tree not found" to peers in either case.
780 if !self.is_database_sync_enabled(tree_id).await {
781 warn!(
782 tree_id = %tree_id,
783 requesting_key = ?requesting_key,
784 requesting_key_name = ?requesting_key_name,
785 "Bootstrap request rejected: database is absent or has no user with sync enabled (responding as not-found)"
786 );
787 return SyncResponse::Error(format!("Tree not found: {tree_id}"));
788 }
789
790 // Get the root entry (to verify tree exists)
791 let instance = match self.instance() {
792 Ok(i) => i,
793 Err(e) => return SyncResponse::Error(format!("Instance dropped: {e}")),
794 };
795 let _root_entry = match instance.backend().get(tree_id).await {
796 Ok(entry) => entry,
797 Err(e) if e.is_not_found() => {
798 warn!(
799 tree_id = %tree_id,
800 requesting_key = ?requesting_key,
801 requesting_key_name = ?requesting_key_name,
802 "Bootstrap request rejected: a user has this tree marked sync-enabled but the backend has no root entry for it"
803 );
804 return SyncResponse::Error(format!("Tree not found: {tree_id}"));
805 }
806 Err(e) => {
807 error!(tree_id = %tree_id, error = %e, "Failed to get root entry");
808 return SyncResponse::Error(format!("Failed to get tree root: {e}"));
809 }
810 };
811
812 // Check if database has authentication configured
813 let auth_configured = match self.check_if_database_has_auth(tree_id).await {
814 Ok(has_auth) => has_auth,
815 Err(e) => {
816 error!(tree_id = %tree_id, error = %e, "Failed to check if database has auth");
817 return SyncResponse::Error(format!("Failed to check database auth: {e}"));
818 }
819 };
820
821 // If auth is configured but no credentials provided, reject the request
822 if auth_configured && requesting_key.is_none() {
823 warn!(
824 tree_id = %tree_id,
825 "Unauthenticated bootstrap request rejected - database requires authentication"
826 );
827 return SyncResponse::Error(
828 "Authentication required: This database requires authenticated access. \
829 Please provide credentials (requesting_key, requesting_key_name, requested_permission) \
830 to bootstrap sync.".to_string()
831 );
832 }
833
834 // A named key must prove ownership before any approval lookup or
835 // lifecycle response, even when the database is public.
836 if let Some(key) = requesting_key
837 && let Err(e) = self.prove_possession(request, key)
838 {
839 warn!(
840 tree_id = %tree_id,
841 requesting_key = %key,
842 error = %e,
843 "Bootstrap request rejected: caller did not prove it holds the claimed key"
844 );
845 return SyncResponse::Error(e.to_string());
846 }
847
848 // Handle key approval for bootstrap requests FIRST
849 let (key_approved, granted_permission) = match (
850 requesting_key,
851 requesting_key_name,
852 requested_permission,
853 ) {
854 // Case 1: All three parameters provided - explicit permission request
855 (Some(key), Some(key_name), Some(permission)) => {
856 info!(
857 tree_id = %tree_id,
858 requesting_key = %key,
859 key_name = %key_name,
860 requested_permission = ?permission,
861 "Processing key approval request for bootstrap"
862 );
863
864 match self.resolve_bootstrap_request(request, key_name).await {
865 Ok(BootstrapRequestOutcome::Pending(request_id)) => {
866 info!(
867 tree_id = %tree_id,
868 request_id = %request_id,
869 "Bootstrap request stored for manual approval"
870 );
871 return SyncResponse::BootstrapPending {
872 request_id,
873 message: "Bootstrap request pending manual approval".to_string(),
874 };
875 }
876 Ok(BootstrapRequestOutcome::Rejected(request_id)) => {
877 return SyncResponse::BootstrapRejected {
878 request_id,
879 message: "Bootstrap request was rejected by an administrator"
880 .to_string(),
881 };
882 }
883 Ok(BootstrapRequestOutcome::Approved) => (true, Some(permission)),
884 Err(e) => {
885 error!(
886 tree_id = %tree_id,
887 error = %e,
888 "Failed to resolve bootstrap request"
889 );
890 return SyncResponse::Error(format!(
891 "Failed to resolve bootstrap request: {e}"
892 ));
893 }
894 }
895 }
896
897 // Case 2: Key provided but permission not specified - auto-detect from auth settings
898 (Some(key), Some(_key_name), None) => {
899 info!(
900 tree_id = %tree_id,
901 requesting_key = %key,
902 "Auto-detecting permission from auth settings for bootstrap request"
903 );
904
905 match self.get_key_highest_permission(tree_id, key).await {
906 Ok(Some(permission)) => {
907 info!(
908 tree_id = %tree_id,
909 requesting_key = %key,
910 detected_permission = ?permission,
911 "Approved bootstrap using auto-detected permission from auth settings"
912 );
913 (true, Some(permission))
914 }
915 Ok(None) => {
916 warn!(
917 tree_id = %tree_id,
918 requesting_key = %key,
919 "Key not found in auth settings - rejecting bootstrap request"
920 );
921 return SyncResponse::Error(
922 "Authentication required: provided key is not authorized for this database".to_string()
923 );
924 }
925 Err(e) => {
926 error!(
927 tree_id = %tree_id,
928 requesting_key = %key,
929 error = %e,
930 "Failed to lookup key permissions"
931 );
932 return SyncResponse::Error(format!("Failed to access auth settings: {e}"));
933 }
934 }
935 }
936
937 // Case 3: No key provided, or key provided without key_name - unauthenticated access
938 _ => {
939 debug!(
940 tree_id = %tree_id,
941 "No authentication credentials provided - proceeding with unauthenticated bootstrap"
942 );
943 (false, None)
944 }
945 };
946
947 // A database with auth configured serves entries only to a caller that
948 // proved it holds a key with read access. Cases that fall through
949 // without approval (no credentials, or a key with no key name) must not
950 // be served just because they reached this point.
951 if auth_configured && !key_approved {
952 warn!(
953 tree_id = %tree_id,
954 requesting_key = ?requesting_key,
955 "Bootstrap request rejected: no proven authority for a database that requires authentication"
956 );
957 return SyncResponse::Error(
958 SyncError::AuthenticationRequired(tree_id.to_string()).to_string(),
959 );
960 }
961
962 // Named-key requests validate non-consumingly before lifecycle access so
963 // legitimate pending retries can reuse the same proof. Once this path
964 // will serve entries, consume the nonce and reject replays.
965 if let Some(key) = requesting_key
966 && let Err(e) = self.authenticate_request(request).and_then(|proven| {
967 if proven == *key {
968 Ok(())
969 } else {
970 Err(SyncError::AuthenticationFailed(format!(
971 "request is signed by {proven}, which is not the claimed key {key}"
972 ))
973 .into())
974 }
975 })
976 {
977 return SyncResponse::Error(e.to_string());
978 }
979
980 // NOW collect all entries after key approval (so we get the updated database state)
981 let all_entries = match self.collect_all_entries_for_bootstrap(tree_id).await {
982 Ok(entries) => entries,
983 Err(e) => {
984 error!(tree_id = %tree_id, error = %e, "Failed to collect all entries for bootstrap after key approval");
985 return SyncResponse::Error(format!(
986 "Failed to collect all entries for bootstrap: {e}"
987 ));
988 }
989 };
990
991 // For bootstrap, we need to send the actual root entry (tree_id) as root_entry
992 // The root_entry should always be the tree's root, not a tip
993 let instance = match self.instance() {
994 Ok(i) => i,
995 Err(e) => return SyncResponse::Error(format!("Instance dropped: {e}")),
996 };
997 let root_entry = match instance.backend().get(tree_id).await {
998 Ok(entry) => entry,
999 Err(e) => {
1000 error!(tree_id = %tree_id, error = %e, "Failed to get root entry");
1001 return SyncResponse::Error(format!("Failed to get root entry: {e}"));
1002 }
1003 };
1004
1005 // Filter out the root from all_entries since we send it separately as root_entry
1006 let other_entries: Vec<_> = all_entries
1007 .into_iter()
1008 .filter(|entry| entry.id() != *tree_id)
1009 .collect();
1010
1011 info!(
1012 tree_id = %tree_id,
1013 entry_count = other_entries.len() + 1,
1014 key_approved = key_approved,
1015 "Sending bootstrap response"
1016 );
1017
1018 SyncResponse::Bootstrap(BootstrapResponse {
1019 tree_id: tree_id.clone(),
1020 root_entry,
1021 all_entries: other_entries,
1022 key_approved,
1023 granted_permission,
1024 })
1025 }
1026
1027 /// Handle incremental sync request.
1028 ///
1029 /// The caller selects this path by sending any non-empty tip list, so it
1030 /// enforces the same read policy bootstrap does. Without that, a single
1031 /// fabricated tip — which matches nothing in our DAG and therefore never
1032 /// stops the ancestor walk — returns the entire tree.
1033 async fn handle_incremental_sync(&self, request: &SyncTreeRequest) -> SyncResponse {
1034 let tree_id = &request.tree_id;
1035 let peer_tips = request.our_tips.tips();
1036
1037 // SECURITY: Check if database has sync enabled (FIRST CHECK - before anything else)
1038 // This prevents information leakage about database existence: the gate
1039 // returns false both for databases that are absent and for databases that
1040 // are present-but-not-tracked-for-sync, and we deliberately respond with
1041 // the same opaque "Tree not found" to peers in either case.
1042 if !self.is_database_sync_enabled(tree_id).await {
1043 warn!(
1044 tree_id = %tree_id,
1045 peer_tip_count = peer_tips.len(),
1046 "Incremental sync request rejected: database is absent or has no user with sync enabled (responding as not-found)"
1047 );
1048 return SyncResponse::Error(format!("Tree not found: {tree_id}"));
1049 }
1050
1051 if let Err(e) = self.authorize_read(request).await {
1052 warn!(
1053 tree_id = %tree_id,
1054 peer_tip_count = peer_tips.len(),
1055 error = %e,
1056 "Incremental sync request rejected: caller is not authorized to read this database"
1057 );
1058 return SyncResponse::Error(e.to_string());
1059 }
1060
1061 // Get our current tips
1062 let instance = match self.instance() {
1063 Ok(i) => i,
1064 Err(e) => return SyncResponse::Error(format!("Instance dropped: {e}")),
1065 };
1066 let our_tips: Vec<ID> = match instance.backend().snapshot(tree_id).await {
1067 Ok(snap) => snap.into_tips(),
1068 Err(e) => {
1069 error!(tree_id = %tree_id, error = %e, "Failed to get our tips");
1070 return SyncResponse::Error(format!("Failed to get tips: {e}"));
1071 }
1072 };
1073
1074 // Find entries peer is missing
1075 let missing_entries = match self
1076 .find_missing_entries_for_peer(&our_tips, peer_tips)
1077 .await
1078 {
1079 Ok(entries) => entries,
1080 Err(e) => {
1081 error!(tree_id = %tree_id, error = %e, "Failed to find missing entries");
1082 return SyncResponse::Error(format!("Failed to find missing entries: {e}"));
1083 }
1084 };
1085
1086 debug!(
1087 tree_id = %tree_id,
1088 our_tips = our_tips.len(),
1089 peer_tips = peer_tips.len(),
1090 missing_count = missing_entries.len(),
1091 "Sending incremental sync response"
1092 );
1093
1094 SyncResponse::Incremental(IncrementalResponse {
1095 tree_id: tree_id.clone(),
1096 their_tips: our_tips,
1097 missing_entries,
1098 })
1099 }
1100}