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}