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

eidetica/sync/transports/
iroh.rs

1//! Iroh transport implementation for sync communication.
2//!
3//! This module provides peer-to-peer sync communication using
4//! Iroh's QUIC-based networking with hole punching and relay servers.
5
6use std::{sync::Arc, time::Duration};
7use tokio::{sync::Mutex, time::timeout};
8
9use async_trait::async_trait;
10use iroh::{
11    Endpoint, RelayMode, SecretKey,
12    endpoint::{Connection, RecvStream, SendStream, presets},
13};
14use iroh_tickets::{Ticket, endpoint::EndpointTicket};
15use serde::{Deserialize, Serialize};
16#[allow(unused_imports)] // Used by write_all method on streams
17use tokio::io::AsyncWriteExt;
18use tokio::sync::oneshot;
19
20use super::{SyncTransport, TransportBuilder, TransportConfig, shared::*};
21use crate::{
22    Result,
23    crdt::Doc,
24    store::Registered,
25    sync::{
26        error::{SyncError, TimeoutPhase},
27        handler::SyncHandler,
28        peer_types::Address,
29        protocol::{RequestContext, SyncRequest, SyncResponse},
30    },
31};
32
33const SYNC_ALPN: &[u8] = b"eidetica/v0";
34
35/// Maximum time spent establishing a connection to a peer.
36///
37/// Looser than the HTTP transport's equivalent, because it bounds something
38/// larger: hole punching and relay fallback, not a TCP handshake. Against a
39/// live but awkwardly-NATed peer that work can legitimately run past ten
40/// seconds, and cutting it short would report a reachable peer as unreachable
41/// rather than merely slow.
42const CONNECT_TIMEOUT: Duration = Duration::from_secs(30);
43
44/// Maximum time for the request/response exchange once connected, covering
45/// opening the stream, writing the request and reading the response.
46///
47/// Without it, a peer that completes the QUIC handshake and then stops
48/// answering holds the request open forever. Requests are awaited inline by the
49/// background sync engine, so an unbounded one stalls every other thing that
50/// engine does.
51const REQUEST_TIMEOUT: Duration = Duration::from_secs(30);
52
53/// Serializable relay mode setting for transport configuration.
54///
55/// This is a simplified version of Iroh's `RelayMode` that can be
56/// persisted to storage. Custom relay configurations are not supported
57/// in the persisted config (use the builder API for custom relays).
58#[derive(Debug, Clone, Serialize, Deserialize, Default, PartialEq, Eq)]
59pub enum RelayModeSetting {
60    /// Use n0's production relay servers (recommended for most deployments)
61    #[default]
62    Default,
63    /// Use n0's staging relay infrastructure (for testing)
64    Staging,
65    /// Disable relay servers entirely (local/direct connections only)
66    Disabled,
67}
68
69impl From<RelayModeSetting> for RelayMode {
70    fn from(setting: RelayModeSetting) -> Self {
71        match setting {
72            RelayModeSetting::Default => RelayMode::Default,
73            RelayModeSetting::Staging => RelayMode::Staging,
74            RelayModeSetting::Disabled => RelayMode::Disabled,
75        }
76    }
77}
78
79/// Persistable configuration for the Iroh transport.
80///
81/// This configuration is stored in the `_sync` database's `transport_configs`
82/// subtree and is automatically loaded when `enable_iroh_transport()` is called.
83///
84/// The most important field is `secret_key_hex`, which stores the node's
85/// cryptographic identity. When this is persisted, the node will have the
86/// same address across restarts.
87///
88/// # Example
89///
90/// ```ignore
91/// use eidetica::sync::transports::iroh::IrohTransportConfig;
92///
93/// // Create a default config (secret key will be generated on first use)
94/// let config = IrohTransportConfig::default();
95///
96/// // Or create with specific settings
97/// let config = IrohTransportConfig {
98///     relay_mode: RelayModeSetting::Disabled,
99///     ..Default::default()
100/// };
101/// ```
102#[derive(Debug, Clone, Serialize, Deserialize)]
103pub struct IrohTransportConfig {
104    /// Secret key bytes (hex encoded for JSON storage).
105    ///
106    /// When `None`, a new secret key will be generated on first use
107    /// and stored back to the config. Once set, this ensures the node
108    /// maintains the same identity (and thus address) across restarts.
109    #[serde(default, skip_serializing_if = "Option::is_none")]
110    pub secret_key_hex: Option<String>,
111
112    /// Relay mode setting for NAT traversal.
113    #[serde(default)]
114    pub relay_mode: RelayModeSetting,
115}
116
117impl Default for IrohTransportConfig {
118    fn default() -> Self {
119        Self {
120            secret_key_hex: None,
121            relay_mode: RelayModeSetting::Default,
122        }
123    }
124}
125
126impl Registered for IrohTransportConfig {
127    fn type_id() -> &'static str {
128        "iroh:v0"
129    }
130}
131
132impl TransportConfig for IrohTransportConfig {}
133
134impl IrohTransportConfig {
135    /// Get the secret key from config, or generate a new one.
136    ///
137    /// If a secret key is already stored in the config, it will be decoded
138    /// and returned. Otherwise, a new random secret key is generated,
139    /// stored in the config (as hex), and returned.
140    ///
141    /// This method mutates the config to store the newly generated key,
142    /// so the caller should persist the config after calling this.
143    pub fn get_or_create_secret_key(&mut self) -> SecretKey {
144        if let Some(hex) = &self.secret_key_hex {
145            let bytes = hex::decode(hex).expect("valid hex in stored secret key");
146            let bytes: [u8; 32] = bytes.try_into().expect("secret key should be 32 bytes");
147            SecretKey::from_bytes(&bytes)
148        } else {
149            // Generate new secret key
150            use rand::RngCore;
151            let mut secret_bytes = [0u8; 32];
152            rand::rngs::OsRng.fill_bytes(&mut secret_bytes);
153            let key = SecretKey::from_bytes(&secret_bytes);
154            self.secret_key_hex = Some(hex::encode(key.to_bytes()));
155            key
156        }
157    }
158
159    /// Check if a secret key has been set in this config.
160    pub fn has_secret_key(&self) -> bool {
161        self.secret_key_hex.is_some()
162    }
163}
164
165/// Builder for configuring IrohTransport with different relay modes and options.
166///
167/// # Examples
168///
169/// ## Production deployment (default)
170/// ```no_run
171/// use eidetica::sync::transports::iroh::IrohTransport;
172///
173/// # fn main() -> Result<(), Box<dyn std::error::Error>> {
174/// let transport = IrohTransport::builder()
175///     .build()?;
176/// // Uses n0's production relay servers by default
177/// # Ok(())
178/// # }
179/// ```
180///
181/// ## Local testing without internet
182/// ```no_run
183/// use eidetica::sync::transports::iroh::IrohTransport;
184/// use iroh::RelayMode;
185///
186/// # fn main() -> Result<(), Box<dyn std::error::Error>> {
187/// let transport = IrohTransport::builder()
188///     .relay_mode(RelayMode::Disabled)
189///     .build()?;
190/// // Direct P2P only, no relay servers
191/// # Ok(())
192/// # }
193/// ```
194///
195/// ## Enterprise deployment with custom relay
196/// ```no_run
197/// use eidetica::sync::transports::iroh::IrohTransport;
198/// use iroh::{RelayConfig, RelayMode, RelayMap, RelayUrl};
199///
200/// # fn main() -> Result<(), Box<dyn std::error::Error>> {
201/// let relay_url: RelayUrl = "https://relay.example.com".parse()?;
202/// let relay_config: RelayConfig = relay_url.into();
203/// let transport = IrohTransport::builder()
204///     .relay_mode(RelayMode::Custom(RelayMap::from_iter([relay_config])))
205///     .build()?;
206/// # Ok(())
207/// # }
208/// ```
209#[derive(Debug, Clone)]
210pub struct IrohTransportBuilder {
211    relay_mode: RelayMode,
212    secret_key: Option<SecretKey>,
213}
214
215impl IrohTransportBuilder {
216    /// Create a new builder with production defaults.
217    ///
218    /// By default, uses `RelayMode::Default` which connects to n0's
219    /// production relay infrastructure for NAT traversal.
220    pub fn new() -> Self {
221        Self {
222            relay_mode: RelayMode::Default, // Use n0's production relays by default
223            secret_key: None,
224        }
225    }
226
227    /// Set the relay mode for the transport.
228    ///
229    /// # Relay Modes
230    ///
231    /// - `RelayMode::Default` - Use n0's production relay servers (recommended)
232    /// - `RelayMode::Staging` - Use n0's staging infrastructure for testing
233    /// - `RelayMode::Disabled` - No relay servers, direct P2P only (local testing)
234    /// - `RelayMode::Custom(RelayMap)` - Use custom relay servers (enterprise deployments)
235    ///
236    /// # Example
237    ///
238    /// ```no_run
239    /// use eidetica::sync::transports::iroh::IrohTransport;
240    /// use iroh::RelayMode;
241    ///
242    /// # fn main() -> Result<(), Box<dyn std::error::Error>> {
243    /// let transport = IrohTransport::builder()
244    ///     .relay_mode(RelayMode::Disabled)
245    ///     .build()?;
246    /// # Ok(())
247    /// # }
248    /// ```
249    pub fn relay_mode(mut self, mode: RelayMode) -> Self {
250        self.relay_mode = mode;
251        self
252    }
253
254    /// Set the secret key for persistent node identity.
255    ///
256    /// When a secret key is provided, the node will have the same
257    /// cryptographic identity (and thus the same address) across restarts.
258    /// This is essential for maintaining stable peer connections.
259    ///
260    /// If not set, a random secret key will be generated on each startup,
261    /// resulting in a different node address each time.
262    ///
263    /// # Example
264    ///
265    /// ```no_run
266    /// use eidetica::sync::transports::iroh::IrohTransport;
267    /// use iroh::SecretKey;
268    /// use rand::RngCore;
269    ///
270    /// # fn main() -> Result<(), Box<dyn std::error::Error>> {
271    /// // Generate or load secret key from storage
272    /// let mut secret_bytes = [0u8; 32];
273    /// rand::rngs::OsRng.fill_bytes(&mut secret_bytes);
274    /// let secret_key = SecretKey::from_bytes(&secret_bytes);
275    ///
276    /// let transport = IrohTransport::builder()
277    ///     .secret_key(secret_key)
278    ///     .build()?;
279    /// # Ok(())
280    /// # }
281    /// ```
282    pub fn secret_key(mut self, key: SecretKey) -> Self {
283        self.secret_key = Some(key);
284        self
285    }
286
287    /// Build the IrohTransport with the configured options.
288    ///
289    /// Returns a configured `IrohTransport` ready to be used with
290    /// `SyncEngine::enable_iroh_transport_with_config()`.
291    pub fn build(self) -> Result<IrohTransport> {
292        Ok(IrohTransport {
293            endpoint: Arc::new(Mutex::new(None)),
294            server_state: ServerState::new(),
295            runtime_config: IrohRuntimeConfig {
296                relay_mode: self.relay_mode,
297                secret_key: self.secret_key,
298            },
299        })
300    }
301}
302
303impl Default for IrohTransportBuilder {
304    fn default() -> Self {
305        Self::new()
306    }
307}
308
309/// Key for storing secret key in persisted Doc
310const SECRET_KEY_FIELD: &str = "secret_key";
311
312#[async_trait]
313impl TransportBuilder for IrohTransportBuilder {
314    type Transport = IrohTransport;
315
316    /// Build the transport using persisted state for identity.
317    ///
318    /// The secret key (node identity) is loaded from the persisted Doc.
319    /// If no secret key exists, a new one is generated and the Doc is returned
320    /// for persistence. Any secret_key set via the builder is ignored -
321    /// persisted state takes precedence.
322    async fn build(self, mut persisted: Doc) -> Result<(Self::Transport, Option<Doc>)> {
323        use crate::crdt::doc::Value;
324
325        // Load or generate secret key from persisted state
326        let (secret_key, updated) = match persisted.get(SECRET_KEY_FIELD) {
327            Some(Value::Text(hex_str)) => {
328                // Decode existing secret key
329                let bytes = hex::decode(hex_str).map_err(|e| {
330                    SyncError::TransportInit(format!("Invalid secret key hex: {}", e))
331                })?;
332                let bytes: [u8; 32] = bytes.try_into().map_err(|_| {
333                    SyncError::TransportInit("Secret key must be 32 bytes".to_string())
334                })?;
335                (SecretKey::from_bytes(&bytes), None)
336            }
337            Some(_) => {
338                return Err(SyncError::TransportInit(
339                    "Secret key must be a text value".to_string(),
340                )
341                .into());
342            }
343            None => {
344                // Generate new secret key and store it
345                // FIXME: Use SecretKey::generate() here (and elsewhere)
346                // Waiting on a rand_core version mismatch to be resolved.
347                use rand::RngCore;
348                let mut secret_bytes = [0u8; 32];
349                rand::rngs::OsRng.fill_bytes(&mut secret_bytes);
350                let key = SecretKey::from_bytes(&secret_bytes);
351                persisted.set(SECRET_KEY_FIELD, hex::encode(key.to_bytes()));
352                (key, Some(persisted))
353            }
354        };
355
356        let transport = IrohTransport {
357            endpoint: Arc::new(Mutex::new(None)),
358            server_state: ServerState::new(),
359            runtime_config: IrohRuntimeConfig {
360                relay_mode: self.relay_mode,
361                secret_key: Some(secret_key),
362            },
363        };
364
365        Ok((transport, updated))
366    }
367}
368
369/// Runtime configuration for IrohTransport (internal, not persisted).
370///
371/// This holds the actual runtime values used by the transport,
372/// including the decoded secret key and relay mode.
373#[derive(Debug, Clone)]
374struct IrohRuntimeConfig {
375    relay_mode: RelayMode,
376    secret_key: Option<SecretKey>,
377}
378
379/// Iroh transport implementation using QUIC peer-to-peer networking.
380///
381/// Provides NAT traversal and direct peer-to-peer connectivity using the Iroh
382/// protocol. Supports both relay-assisted and direct connections.
383///
384/// # How It Works
385///
386/// 1. **Discovery**: Peers find each other via relay servers or direct addresses
387/// 2. **Connection**: Attempts direct connection through NAT hole-punching
388/// 3. **Fallback**: Uses relay servers if direct connection fails
389/// 4. **Upgrade**: Automatically upgrades to direct connection when possible
390///
391/// # Server Addresses
392///
393/// Addresses use iroh's standard `EndpointTicket` format (postcard + base32-lower
394/// with `endpoint` prefix) for both `get_server_address()` and `DatabaseTicket` URLs.
395///
396/// # Example
397///
398/// ```no_run
399/// use eidetica::sync::transports::iroh::IrohTransport;
400/// use iroh::RelayMode;
401///
402/// # fn main() -> Result<(), Box<dyn std::error::Error>> {
403/// // Create with defaults (production relay servers)
404/// let transport = IrohTransport::new()?;
405///
406/// // Or use the builder for custom configuration
407/// let transport = IrohTransport::builder()
408///     .relay_mode(RelayMode::Staging)
409///     .build()?;
410/// # Ok(())
411/// # }
412/// ```
413pub struct IrohTransport {
414    /// The Iroh endpoint for P2P communication (lazily initialized).
415    endpoint: Arc<Mutex<Option<Endpoint>>>,
416    /// Shared server state management.
417    server_state: ServerState,
418    /// Runtime configuration (relay mode, secret key, etc.)
419    runtime_config: IrohRuntimeConfig,
420}
421
422impl IrohTransport {
423    /// Transport type identifier for Iroh
424    pub const TRANSPORT_TYPE: &'static str = "iroh";
425
426    /// Create a new Iroh transport instance with production defaults.
427    ///
428    /// Uses `RelayMode::Default` which connects to n0's production relay
429    /// infrastructure. For custom configuration, use `IrohTransport::builder()`.
430    ///
431    /// The endpoint will be lazily initialized on first use.
432    ///
433    /// # Example
434    ///
435    /// ```no_run
436    /// use eidetica::sync::transports::iroh::IrohTransport;
437    ///
438    /// # fn main() -> Result<(), Box<dyn std::error::Error>> {
439    /// let transport = IrohTransport::new()?;
440    /// // Use with: sync.enable_iroh_transport_with_config(transport)?;
441    /// # Ok(())
442    /// # }
443    /// ```
444    pub fn new() -> Result<Self> {
445        IrohTransportBuilder::new().build()
446    }
447
448    /// Create a builder for configuring the transport.
449    ///
450    /// Allows customization of relay modes and other transport options.
451    ///
452    /// # Example
453    ///
454    /// ```no_run
455    /// use eidetica::sync::transports::iroh::IrohTransport;
456    /// use iroh::RelayMode;
457    ///
458    /// # fn main() -> Result<(), Box<dyn std::error::Error>> {
459    /// let transport = IrohTransport::builder()
460    ///     .relay_mode(RelayMode::Disabled)
461    ///     .build()?;
462    /// # Ok(())
463    /// # }
464    /// ```
465    pub fn builder() -> IrohTransportBuilder {
466        IrohTransportBuilder::new()
467    }
468
469    /// Initialize the Iroh endpoint if not already done.
470    async fn ensure_endpoint(&self) -> Result<Endpoint> {
471        let mut endpoint_lock = self.endpoint.lock().await;
472
473        if endpoint_lock.is_none() {
474            // Create a new Iroh endpoint with configured relay mode
475            let mut builder = Endpoint::builder(presets::N0)
476                .alpns(vec![SYNC_ALPN.to_vec()])
477                .relay_mode(self.runtime_config.relay_mode.clone());
478
479            // Use the provided secret key for persistent identity
480            if let Some(secret_key) = &self.runtime_config.secret_key {
481                builder = builder.secret_key(secret_key.clone());
482            }
483
484            let endpoint = builder.bind().await.map_err(|e| {
485                SyncError::TransportInit(format!("Failed to create Iroh endpoint: {e}"))
486            })?;
487
488            *endpoint_lock = Some(endpoint);
489        }
490
491        Ok(endpoint_lock.as_ref().unwrap().clone())
492    }
493
494    /// Start the server request handling loop.
495    async fn start_server_loop(
496        &self,
497        endpoint: Endpoint,
498        ready_tx: oneshot::Sender<()>,
499        shutdown_rx: oneshot::Receiver<()>,
500        handler: Arc<dyn SyncHandler>,
501    ) -> Result<()> {
502        let mut shutdown_rx = shutdown_rx;
503
504        // Signal that we're ready
505        let _ = ready_tx.send(());
506
507        // Accept incoming connections
508        tokio::spawn(async move {
509            loop {
510                tokio::select! {
511                    // Check for shutdown signal
512                    _ = &mut shutdown_rx => {
513                        break;
514                    }
515                    // Accept incoming connections
516                    connection_result = endpoint.accept() => {
517                        match connection_result {
518                            Some(connecting) => {
519                                let handler_clone = handler.clone();
520                                tokio::spawn(async move {
521                                    if let Ok(conn) = connecting.await {
522                                        Self::handle_connection(conn, handler_clone).await;
523                                    }
524                                });
525                            }
526                            None => break, // Endpoint closed
527                        }
528                    }
529                }
530            }
531            // Server loop has exited - the shutdown was triggered by stop_server()
532            // which already marked the server as stopped, so no additional cleanup needed here
533        });
534
535        Ok(())
536    }
537
538    /// Handle an incoming connection.
539    async fn handle_connection(conn: Connection, handler: Arc<dyn SyncHandler>) {
540        // Get the remote peer node ID for context
541        let remote_endpoint_id = conn.remote_id();
542        let remote_address = Address {
543            transport_type: Self::TRANSPORT_TYPE.to_string(),
544            address: remote_endpoint_id.to_string(),
545        };
546
547        // Accept incoming streams and process sequentially
548        // Note: We process streams sequentially because SyncHandler::handle_request
549        // returns non-Send futures (internal types use Rc/RefCell).
550        while let Ok((send_stream, recv_stream)) = conn.accept_bi().await {
551            Self::handle_stream(
552                send_stream,
553                recv_stream,
554                handler.clone(),
555                remote_address.clone(),
556            )
557            .await;
558        }
559    }
560
561    /// Handle an incoming bidirectional stream.
562    async fn handle_stream(
563        mut send_stream: SendStream,
564        mut recv_stream: RecvStream,
565        handler: Arc<dyn SyncHandler>,
566        remote_address: Address,
567    ) {
568        // Read the request with size limit (1MB)
569        let buffer: Vec<u8> = match recv_stream.read_to_end(1024 * 1024).await {
570            Ok(buffer) => buffer,
571            Err(e) => {
572                tracing::error!("Failed to read stream: {e}");
573                return;
574            }
575        };
576
577        // Deserialize the request using JsonHandler
578        let request: SyncRequest = match JsonHandler::deserialize_request(&buffer) {
579            Ok(req) => req,
580            Err(e) => {
581                tracing::error!("Failed to deserialize request: {e}");
582                return;
583            }
584        };
585
586        // Extract peer_pubkey from SyncTreeRequest if present
587        let peer_pubkey = match &request {
588            SyncRequest::SyncTree(sync_tree_request) => sync_tree_request.peer_pubkey.clone(),
589            _ => None,
590        };
591
592        // Create request context with remote address and peer pubkey
593        let context = RequestContext {
594            remote_address: Some(remote_address),
595            peer_pubkey,
596        };
597
598        // Handle the request using the SyncHandler
599        let response = handler.handle_request(&request, &context).await;
600
601        // Serialize and send response using JsonHandler
602        match JsonHandler::serialize_response(&response) {
603            Ok(response_bytes) => {
604                if let Err(e) = send_stream.write_all(&response_bytes).await {
605                    tracing::error!("Failed to write response: {e}");
606                    return;
607                }
608                if let Err(e) = send_stream.finish() {
609                    tracing::error!("Failed to finish stream: {e}");
610                }
611            }
612            Err(e) => {
613                tracing::error!("Failed to serialize response: {e}");
614            }
615        }
616    }
617}
618
619#[async_trait]
620impl SyncTransport for IrohTransport {
621    fn transport_type(&self) -> &'static str {
622        Self::TRANSPORT_TYPE
623    }
624
625    fn can_handle_address(&self, address: &Address) -> bool {
626        address.transport_type == Self::TRANSPORT_TYPE
627    }
628
629    async fn start_server(&self, handler: Arc<dyn SyncHandler>) -> Result<()> {
630        let start = self.server_state.begin_start("iroh-endpoint")?;
631
632        // Ensure we have an endpoint and get EndpointAddr with direct addresses
633        let endpoint = self.ensure_endpoint().await?;
634        let endpoint_clone = endpoint.clone();
635
636        // Get the EndpointAddr with direct addresses
637        // Note: We don't wait for online() - direct addresses are available immediately
638        // after bind(), and relay connections happen asynchronously in the background.
639        let endpoint_addr = endpoint.addr();
640        let endpoint_addr_str = EndpointTicket::new(endpoint_addr).encode_string();
641
642        // Create server coordination channels
643        let (ready_tx, ready_rx) = oneshot::channel();
644        let (shutdown_tx, shutdown_rx) = oneshot::channel();
645
646        // Start server loop
647        self.start_server_loop(endpoint_clone, ready_tx, shutdown_rx, handler)
648            .await?;
649
650        // Wait for server to be ready using shared utility
651        wait_for_ready(ready_rx, "iroh-endpoint").await?;
652
653        // Start server state with EndpointAddr string and shutdown sender
654        start.complete(endpoint_addr_str, shutdown_tx);
655
656        Ok(())
657    }
658
659    async fn stop_server(&self) -> Result<()> {
660        if !self.server_state.is_running() {
661            return Err(SyncError::ServerNotRunning.into());
662        }
663
664        self.server_state.stop_server();
665        if let Some(endpoint) = self.endpoint.lock().await.take() {
666            endpoint.close().await;
667        }
668
669        Ok(())
670    }
671
672    async fn send_request(&self, address: &Address, request: &SyncRequest) -> Result<SyncResponse> {
673        if !self.can_handle_address(address) {
674            return Err(SyncError::UnsupportedTransport {
675                transport_type: address.transport_type.clone(),
676            }
677            .into());
678        }
679
680        // Ensure we have an endpoint (lazy initialization)
681        let endpoint = self.ensure_endpoint().await?;
682
683        // Deserialize the EndpointTicket address
684        let endpoint_ticket =
685            <EndpointTicket as Ticket>::decode_string(&address.address).map_err(|e| {
686                SyncError::SerializationError(format!(
687                    "Failed to parse EndpointTicket '{}': {e}",
688                    address.address
689                ))
690            })?;
691        let endpoint_addr = endpoint_ticket.endpoint_addr().clone();
692
693        // Connect to the peer, bounding how long an unreachable one can cost.
694        let conn = timeout(CONNECT_TIMEOUT, endpoint.connect(endpoint_addr, SYNC_ALPN))
695            .await
696            .map_err(|_| SyncError::Timeout {
697                address: address.address.clone(),
698                phase: TimeoutPhase::Connect,
699                elapsed: CONNECT_TIMEOUT,
700            })?
701            .map_err(|e| SyncError::ConnectionFailed {
702                address: address.address.clone(),
703                reason: e.to_string(),
704            })?;
705
706        // Serialize the request before starting the exchange clock.
707        let request_bytes = JsonHandler::serialize_request(request)?;
708
709        // One deadline covers the whole exchange: a peer can otherwise stall
710        // any single step of it indefinitely.
711        let response_bytes: Vec<u8> = timeout(REQUEST_TIMEOUT, async {
712            // Open a bidirectional stream
713            let (mut send_stream, mut recv_stream) = conn
714                .open_bi()
715                .await
716                .map_err(|e| SyncError::Network(format!("Failed to open stream: {e}")))?;
717
718            send_stream
719                .write_all(&request_bytes)
720                .await
721                .map_err(|e| SyncError::Network(format!("Failed to write request: {e}")))?;
722
723            send_stream
724                .finish()
725                .map_err(|e| SyncError::Network(format!("Failed to finish send stream: {e}")))?;
726
727            // Read the response with size limit (1MB)
728            recv_stream
729                .read_to_end(1024 * 1024)
730                .await
731                .map_err(|e| SyncError::Network(format!("Failed to read response: {e}")))
732        })
733        .await
734        .map_err(|_| {
735            // The peer answered the connection attempt, so it is reachable:
736            // the phase records that, rather than the peer being reported as
737            // unreachable.
738            SyncError::Timeout {
739                address: address.address.clone(),
740                phase: TimeoutPhase::Request,
741                elapsed: REQUEST_TIMEOUT,
742            }
743        })??;
744
745        // Deserialize the response using JsonHandler
746        let response: SyncResponse = JsonHandler::deserialize_response(&response_bytes)?;
747
748        Ok(response)
749    }
750
751    fn is_server_running(&self) -> bool {
752        self.server_state.is_running()
753    }
754
755    fn get_server_address(&self) -> Result<String> {
756        self.server_state.get_address().map_err(|e| e.into())
757    }
758}
759
760#[cfg(test)]
761mod tests {
762    use super::*;
763    use iroh::{EndpointAddr, TransportAddr};
764
765    /// Round-trip: EndpointAddr → EndpointTicket string → EndpointAddr
766    #[test]
767    fn endpoint_ticket_round_trip() {
768        let secret_key = SecretKey::from_bytes(&[1u8; 32]);
769        let endpoint_addr = EndpointAddr::from_parts(
770            secret_key.public(),
771            vec![
772                TransportAddr::Ip("127.0.0.1:1234".parse().unwrap()),
773                TransportAddr::Ip("192.168.1.1:5678".parse().unwrap()),
774            ],
775        );
776
777        // Serialize to EndpointTicket string
778        let ticket_str = EndpointTicket::new(endpoint_addr.clone()).encode_string();
779        assert!(
780            ticket_str.starts_with("endpoint"),
781            "EndpointTicket should start with 'endpoint' prefix: {ticket_str}"
782        );
783
784        // Deserialize back
785        let ticket = <EndpointTicket as Ticket>::decode_string(&ticket_str).unwrap();
786        let round_tripped = ticket.endpoint_addr();
787
788        assert_eq!(endpoint_addr.id, round_tripped.id);
789        assert_eq!(
790            endpoint_addr.ip_addrs().collect::<Vec<_>>(),
791            round_tripped.ip_addrs().collect::<Vec<_>>()
792        );
793    }
794}