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

eidetica/sync/
mod.rs

1//! Synchronization module for Eidetica database.
2//!
3//! # Quick Start
4//!
5//! ```rust,ignore
6//! // Enable sync
7//! instance.enable_sync().await?;
8//! let sync = instance.sync().unwrap();
9//!
10//! // Register transports with their configurations
11//! sync.register_transport("http", HttpTransport::builder()
12//!     .bind("127.0.0.1:8080")
13//! ).await?;
14//! sync.register_transport("p2p", IrohTransport::builder()).await?;
15//!
16//! // Start accepting incoming connections
17//! sync.accept_connections().await?;
18//!
19//! // Outbound sync works via the registered transports
20//! sync.sync_with_peer(&Address::http("peer:8080"), Some(&tree_id)).await?;
21//! ```
22//!
23//! # Architecture
24//!
25//! The sync system uses a Background Sync architecture with command-pattern communication:
26//!
27//! - **[`Sync`]**: Thread-safe frontend using `Arc<Sync>` with interior mutability
28//!   (`OnceLock`). Provides the public API and sends commands to the background.
29//! - **[`background::BackgroundSync`]**: Single background thread handling all sync operations
30//!   via a command loop. Owns the transport and retry queue.
31//! - **Write Callbacks**: Automatically trigger sync when entries are committed via
32//!   database write callbacks.
33//!
34//! # Connection Model
35//!
36//! The sync system separates **outbound** and **inbound** connection handling:
37//!
38//! - **Outbound** ([`Sync::sync_with_peer`]): Works after registering transports via
39//!   [`Sync::register_transport`]. Each transport can be configured via its builder.
40//! - **Inbound** ([`Sync::accept_connections`]): Must be explicitly called each time
41//!   the instance starts. Starts servers on all registered transports.
42//!
43//! This design provides security by default. Nodes don't accept incoming connections
44//! unless explicitly opted in.
45//!
46//! # Bootstrap Protocol
47//!
48//! The sync protocol detects whether a client needs bootstrap or incremental sync:
49//!
50//! - **Empty tips** → Full bootstrap (complete database transfer)
51//! - **Has tips** → Incremental sync (only missing entries)
52//!
53//! Use [`Sync::sync_with_peer`] which handles both cases automatically.
54//!
55//! # Duplicate Prevention
56//!
57//! Uses Merkle-DAG tip comparison instead of tracking individual sent entries:
58//!
59//! 1. Exchange tips with peer
60//! 2. Compare DAGs to find missing entries
61//! 3. Send only what peer doesn't have
62//! 4. Receive only what we're missing
63//!
64//! This approach requires no extra storage apart from tracking relevant Merkle-DAG tips.
65//!
66//! # Transport Layer
67//!
68//! Two transport implementations are available:
69//!
70//! - **HTTP** ([`transports::http::HttpTransport`]): REST API at `/api/v0`, JSON serialization
71//! - **Iroh P2P** ([`transports::iroh::IrohTransport`]): QUIC-based with NAT traversal
72//!
73//! Transports are registered via [`Sync::register_transport`] with their builders.
74//! State (like Iroh node identity) is automatically persisted per named instance.
75//! Both implement the [`transports::SyncTransport`] trait.
76//!
77//! # Peer and State Management
78//!
79//! Peers and sync relationships are stored in a dedicated sync database (`_sync`):
80//!
81//! - `peer_manager::PeerManager`: Handles peer registration and relationships
82//!
83//! # Connection Behavior
84//!
85//! - **Lazy connections**: Established on-demand, not at peer registration
86//! - **Periodic sync**: Configurable interval (default 5 minutes)
87//! - **Retry queue**: Failed sends retried with exponential backoff
88
89use std::sync::{Arc, OnceLock};
90use std::time::SystemTime;
91use tokio::sync::mpsc;
92
93use crate::{
94    Database, Instance, Result, WeakInstance,
95    auth::{Permission, crypto::PublicKey},
96    crdt::{Doc, doc::Value},
97    entry::ID,
98    instance::backend::Backend,
99    store::{DocStore, Registry},
100};
101
102// Public submodules
103pub mod background;
104pub mod error;
105pub mod handler;
106mod handler_tree_ops;
107pub mod peer_manager;
108pub mod peer_types;
109pub mod protocol;
110pub mod ticket;
111pub mod transports;
112pub mod utils;
113
114// Private submodules
115mod bootstrap;
116mod bootstrap_request_manager;
117mod ops;
118mod peer;
119mod peer_state;
120mod queue;
121mod transport;
122mod transport_manager;
123mod user;
124mod user_sync_manager;
125
126// Re-exports
127use background::SyncCommand;
128pub use bootstrap_request_manager::{BootstrapRequest, RequestStatus};
129pub use error::{SyncError, TimeoutPhase};
130use peer_state::PeerStates;
131pub use peer_types::{Address, ConnectionState, PeerId, PeerInfo, PeerStatus};
132use queue::SyncQueue;
133pub use ticket::DatabaseTicket;
134use transports::TransportConfig;
135
136/// Private constant for the sync settings subtree name
137const SETTINGS_SUBTREE: &str = "settings_map";
138
139/// Private constant for the transports registry subtree name
140const TRANSPORTS_SUBTREE: &str = "transports";
141
142/// Private constant for the transport state store name (persisted identity/state per transport instance)
143const TRANSPORT_STATE_STORE: &str = "transport_state";
144
145/// Authentication parameters for sync operations.
146#[derive(Debug, Clone)]
147pub struct AuthParams {
148    /// The public key making the request
149    pub requesting_key: PublicKey,
150    /// The name/ID of the requesting key
151    pub requesting_key_name: String,
152    /// The permission level being requested
153    pub requested_permission: Permission,
154}
155
156/// Information needed to register a peer for syncing.
157///
158/// This is used with [`Sync::register_sync_peer()`] to declare sync intent.
159#[derive(Debug, Clone)]
160pub struct SyncPeerInfo {
161    /// The peer's public key
162    pub peer_pubkey: PublicKey,
163    /// The tree/database to sync
164    pub tree_id: ID,
165    /// Initial address hints where the peer might be found
166    pub addresses: Vec<Address>,
167    /// Optional authentication parameters for bootstrap
168    pub auth: Option<AuthParams>,
169    /// Optional display name for the peer
170    pub display_name: Option<String>,
171}
172
173/// Handle for tracking sync status with a specific peer.
174///
175/// Returned by [`Sync::register_sync_peer()`].
176#[derive(Debug, Clone)]
177pub struct SyncHandle {
178    tree_id: ID,
179    peer_pubkey: PublicKey,
180    sync: Sync,
181}
182
183impl SyncHandle {
184    /// Get the current sync status.
185    pub async fn status(&self) -> Result<SyncStatus> {
186        self.sync
187            .get_sync_status(&self.tree_id, &self.peer_pubkey)
188            .await
189    }
190
191    /// Add another address hint for this peer.
192    pub async fn add_address(&self, address: Address) -> Result<()> {
193        self.sync.add_peer_address(&self.peer_pubkey, address).await
194    }
195
196    /// Block until initial sync completes (has local data).
197    ///
198    /// This is a convenience method for backwards compatibility.
199    /// The sync happens in the background, this just polls until data arrives.
200    pub async fn wait_for_initial_sync(&self) -> Result<()> {
201        loop {
202            let status = self.status().await?;
203            if status.has_local_data {
204                return Ok(());
205            }
206            tokio::time::sleep(tokio::time::Duration::from_millis(100)).await;
207        }
208    }
209
210    /// Get the tree ID being synced.
211    pub fn tree_id(&self) -> &ID {
212        &self.tree_id
213    }
214
215    /// Get the peer public key.
216    pub fn peer_pubkey(&self) -> &PublicKey {
217        &self.peer_pubkey
218    }
219}
220
221/// Current sync status for a tree/peer pair.
222#[derive(Debug, Clone)]
223pub struct SyncStatus {
224    /// Whether we have local data for this tree
225    pub has_local_data: bool,
226    /// When at least one tree last synced with this peer.
227    ///
228    /// Peer-wide rather than specific to `has_local_data`'s tree, and held only
229    /// in memory by the running sync engine, so it reads `None` before the
230    /// first successful round and again after a restart.
231    pub last_sync: Option<SystemTime>,
232    /// Last error encountered (if any)
233    pub last_error: Option<String>,
234}
235
236/// Synchronization manager for the database.
237///
238/// The Sync module is a thin frontend that communicates with a background
239/// sync engine thread via command channels. All actual sync operations, transport
240/// communication, and state management happen in the background thread.
241///
242/// ## Multi-Transport Support
243///
244/// Multiple transports can be enabled simultaneously (e.g., HTTP + Iroh P2P),
245/// allowing peers to be reachable via different networks. Requests are automatically
246/// routed to the appropriate transport based on address type.
247///
248/// ```rust,ignore
249/// // Enable both HTTP and Iroh transports
250/// sync.enable_http_transport().await?;
251/// sync.enable_iroh_transport().await?;
252///
253/// // Start servers on all transports
254/// sync.start_server("127.0.0.1:0").await?;
255///
256/// // Get all server addresses
257/// let addresses = sync.get_all_server_addresses().await?;
258/// ```
259#[derive(Debug)]
260pub struct Sync {
261    /// Communication channel to the background sync engine.
262    /// Initialized when the first transport is enabled via `enable_*_transport()` or `add_transport()`.
263    background_tx: OnceLock<mpsc::Sender<SyncCommand>>,
264    /// The instance for read operations and tree management
265    instance: WeakInstance,
266    /// The tree containing synchronization settings
267    sync_tree: Database,
268    /// Queue for entries pending synchronization
269    queue: Arc<SyncQueue>,
270    /// Serialize user-preference reconciliation so concurrent service writes
271    /// cannot race updates to the persisted combined settings tree.
272    reconciliation: Arc<tokio::sync::Mutex<()>>,
273    /// Per-peer liveness state, written by the background engine.
274    peer_state: Arc<PeerStates>,
275}
276
277impl Clone for Sync {
278    fn clone(&self) -> Self {
279        let background_tx = OnceLock::new();
280        if let Some(tx) = self.background_tx.get() {
281            let _ = background_tx.set(tx.clone());
282        }
283        Self {
284            background_tx,
285            instance: self.instance.clone(),
286            sync_tree: self.sync_tree.clone(),
287            queue: Arc::clone(&self.queue),
288            reconciliation: Arc::clone(&self.reconciliation),
289            peer_state: Arc::clone(&self.peer_state),
290        }
291    }
292}
293
294impl Sync {
295    /// Create a new Sync instance with a dedicated settings tree.
296    ///
297    /// # Arguments
298    /// * `instance` - The database instance for tree operations
299    ///
300    /// # Returns
301    /// A new Sync instance with its own settings tree.
302    pub async fn new(instance: Instance) -> Result<Self> {
303        // Get device key from instance
304        let signing_key = instance.signing_key()?.clone();
305
306        let mut sync_settings = Doc::new();
307        sync_settings.set("name", "_sync");
308        sync_settings.set("type", "sync_settings");
309
310        let sync_tree = Database::create(&instance, signing_key, sync_settings).await?;
311
312        let sync = Self {
313            background_tx: OnceLock::new(),
314            instance: instance.downgrade(),
315            sync_tree,
316            queue: Arc::new(SyncQueue::new()),
317            reconciliation: Arc::new(tokio::sync::Mutex::new(())),
318            peer_state: Arc::new(PeerStates::default()),
319        };
320
321        // Initialize combined settings for all tracked users
322        sync.initialize_user_settings().await?;
323
324        Ok(sync)
325    }
326
327    /// Load an existing Sync instance from a sync tree root ID.
328    ///
329    /// # Arguments
330    /// * `instance` - The database instance
331    /// * `sync_tree_root_id` - The root ID of the existing sync tree
332    ///
333    /// # Returns
334    /// A Sync instance loaded from the existing tree.
335    pub async fn load(instance: Instance, sync_tree_root_id: &ID) -> Result<Self> {
336        let device_key = instance.signing_key()?.clone();
337
338        let sync_tree = Database::open(&instance, sync_tree_root_id)
339            .await?
340            .with_key(device_key);
341
342        let sync = Self {
343            background_tx: OnceLock::new(),
344            instance: instance.downgrade(),
345            sync_tree,
346            queue: Arc::new(SyncQueue::new()),
347            reconciliation: Arc::new(tokio::sync::Mutex::new(())),
348            peer_state: Arc::new(PeerStates::default()),
349        };
350
351        // Initialize combined settings for all tracked users
352        sync.initialize_user_settings().await?;
353
354        Ok(sync)
355    }
356
357    /// Get the root ID of the sync settings tree.
358    pub fn sync_tree_root_id(&self) -> &ID {
359        self.sync_tree.root_id()
360    }
361
362    /// Store a setting in the sync_settings subtree.
363    ///
364    /// # Arguments
365    /// * `key` - The setting key
366    /// * `value` - The setting value
367    pub async fn set_setting(
368        &self,
369        key: impl Into<String>,
370        value: impl Into<String>,
371    ) -> Result<()> {
372        let txn = self.sync_tree.new_transaction().await?;
373        let sync_settings = txn.get_store::<DocStore>(SETTINGS_SUBTREE).await?;
374        sync_settings.set(key, Value::Text(value.into())).await?;
375        txn.commit().await?;
376        Ok(())
377    }
378
379    /// Retrieve a setting from the settings_map subtree.
380    ///
381    /// # Arguments
382    /// * `key` - The setting key to retrieve
383    ///
384    /// # Returns
385    /// The setting value if found, None otherwise.
386    pub async fn get_setting(&self, key: impl AsRef<str>) -> Result<Option<String>> {
387        let sync_settings = self
388            .sync_tree
389            .get_store_viewer::<DocStore>(SETTINGS_SUBTREE)
390            .await?;
391        match sync_settings.get_string(key).await {
392            Ok(value) => Ok(Some(value)),
393            Err(e) if e.is_not_found() => Ok(None),
394            Err(e) => Err(e),
395        }
396    }
397
398    /// Load a transport configuration from the `_sync` database.
399    ///
400    /// Transport configurations are stored in the `transports` subtree,
401    /// keyed by their name. If no configuration exists for the transport,
402    /// returns the default configuration.
403    ///
404    /// # Type Parameters
405    /// * `T` - The transport configuration type implementing [`TransportConfig`]
406    ///
407    /// # Arguments
408    /// * `name` - The name of the transport instance (e.g., "iroh", "http")
409    ///
410    /// # Returns
411    /// The loaded configuration, or the default if not found.
412    ///
413    /// # Example
414    ///
415    /// ```ignore
416    /// use eidetica::sync::transports::iroh::IrohTransportConfig;
417    ///
418    /// let config: IrohTransportConfig = sync.load_transport_config("iroh")?;
419    /// ```
420    pub async fn load_transport_config<T: TransportConfig>(&self, name: &str) -> Result<T> {
421        let tx = self.sync_tree.new_transaction().await?;
422        let registry = Registry::new(&tx, TRANSPORTS_SUBTREE).await?;
423
424        match registry.get_entry(name).await {
425            Ok(entry) => {
426                // Verify the type matches
427                if entry.type_id != T::type_id() {
428                    return Err(SyncError::TransportTypeMismatch {
429                        name: name.to_string(),
430                        expected: T::type_id().to_string(),
431                        found: entry.type_id,
432                    }
433                    .into());
434                }
435                entry.config.get_json("data").map_err(|e| {
436                    SyncError::SerializationError(format!(
437                        "Failed to deserialize transport config '{name}': {e}"
438                    ))
439                    .into()
440                })
441            }
442            Err(e) if e.is_not_found() => Ok(T::default()),
443            Err(e) => Err(e),
444        }
445    }
446
447    /// Save a transport configuration to the `_sync` database.
448    ///
449    /// Transport configurations are stored in the `transports` subtree,
450    /// keyed by their name. This persists the configuration so it can
451    /// be loaded on subsequent startups.
452    ///
453    /// # Type Parameters
454    /// * `T` - The transport configuration type implementing [`TransportConfig`]
455    ///
456    /// # Arguments
457    /// * `name` - The name of the transport instance (e.g., "iroh", "http")
458    /// * `config` - The configuration to save
459    ///
460    /// # Example
461    ///
462    /// ```ignore
463    /// use eidetica::sync::transports::iroh::IrohTransportConfig;
464    ///
465    /// let mut config = IrohTransportConfig::default();
466    /// config.get_or_create_secret_key(); // Generate key
467    /// sync.save_transport_config("iroh", &config)?;
468    /// ```
469    pub async fn save_transport_config<T: TransportConfig>(
470        &self,
471        name: &str,
472        config: &T,
473    ) -> Result<()> {
474        let mut config_doc = crate::crdt::Doc::new();
475        config_doc.set_json("data", config).map_err(|e| {
476            SyncError::SerializationError(format!(
477                "Failed to serialize transport config '{name}': {e}"
478            ))
479        })?;
480        let tx = self.sync_tree.new_transaction().await?;
481        let registry = Registry::new(&tx, TRANSPORTS_SUBTREE).await?;
482        registry.set_entry(name, T::type_id(), config_doc).await?;
483        tx.commit().await?;
484        Ok(())
485    }
486
487    /// Get a reference to the Instance.
488    pub fn instance(&self) -> Result<Instance> {
489        self.instance
490            .upgrade()
491            .ok_or_else(|| SyncError::InstanceDropped.into())
492    }
493
494    /// Get a clone of the backend seam.
495    pub fn backend(&self) -> Result<Arc<dyn Backend>> {
496        Ok(self.instance()?.backend().clone())
497    }
498
499    /// Get the sync tree database.
500    pub fn sync_tree(&self) -> &Database {
501        &self.sync_tree
502    }
503
504    /// Get the device public key for this sync instance.
505    ///
506    /// # Returns
507    /// The device's public key.
508    pub fn get_device_pubkey(&self) -> Result<PublicKey> {
509        Ok(self.instance()?.id())
510    }
511}
512
513impl Sync {
514    // === Test Helpers ===
515}