moonpool_transport/rpc/net_transport.rs
1//! `NetTransport`: Central transport coordinator (FDB pattern).
2//!
3//! Manages peer connections and dispatches incoming packets to endpoints.
4//! Provides synchronous send API (FDB pattern: never await on send).
5//!
6//! # FDB Reference
7//! From NetTransport.actor.cpp:300-600, NetTransport.h:195-314
8//!
9//! # Usage
10//!
11//! Use [`NetTransportBuilder`] to create a properly configured transport:
12//!
13//! ```rust,ignore
14//! // For servers and clients that need RPC responses:
15//! let transport = NetTransportBuilder::new(network, time, task)
16//! .local_address(addr)
17//! .build_listening()
18//! .await?;
19//!
20//! // For fire-and-forget senders only (no listening):
21//! let transport = NetTransportBuilder::new(network, time, task)
22//! .local_address(addr)
23//! .build();
24//! ```
25//!
26//! The builder automatically handles `Arc` wrapping and `set_weak_self()`,
27//! eliminating the most common footgun in `NetTransport` usage.
28
29use std::collections::BTreeMap;
30use std::sync::{Arc, RwLock, Weak};
31
32use super::failure_monitor::FailureMonitor;
33use super::net_notified_queue::ReplyQueueCloser;
34use super::reply_error::ReplyError;
35use crate::{
36 Endpoint, NetworkAddress, NetworkProvider, Peer, PeerConfig, Providers, TaskProvider,
37 TcpListenerTrait, UID, WellKnownToken,
38};
39use tokio_util::sync::CancellationToken;
40
41use super::endpoint_map::{EndpointMap, MessageReceiver};
42use super::request_stream::RequestStream;
43use crate::MessageCodec;
44use crate::error::MessagingError;
45use serde::Serialize;
46use serde::de::DeserializeOwned;
47
48/// Type alias for shared peer reference.
49type SharedPeer<P> = Arc<RwLock<Peer<P>>>;
50
51/// Pending reply entry: token + weak reference to the queue closer.
52type PendingReplyEntry = (UID, Weak<dyn ReplyQueueCloser>);
53
54/// Internal transport data (FDB `TransportData` equivalent).
55///
56/// Separates internal mutable state from public API, matching FDB's pattern.
57/// See NetTransport.actor.cpp:300-350 for the original `TransportData` struct.
58///
59/// # FDB Reference
60/// ```cpp
61/// struct TransportData {
62/// std::unordered_map<NetworkAddress, Reference<struct Peer>> peers;
63/// std::unordered_map<NetworkAddress, std::pair<double, double>> closedPeers;
64/// // ... endpoints, stats, etc.
65/// };
66/// ```
67struct TransportData<P: Providers> {
68 /// Endpoint map for routing incoming packets.
69 endpoints: EndpointMap,
70
71 /// Peer connections keyed by destination address (outgoing).
72 /// FDB: `std::unordered_map`<`NetworkAddress`, Reference<struct Peer>> peers;
73 peers: BTreeMap<String, SharedPeer<P>>,
74
75 /// Incoming peer connections (from accepted connections).
76 /// Separate from outgoing peers to avoid conflicts.
77 incoming_peers: BTreeMap<String, SharedPeer<P>>,
78
79 /// Failure monitor for address/endpoint failure tracking.
80 /// FDB: `IFailureMonitor` (FailureMonitor.h)
81 failure_monitor: Arc<FailureMonitor>,
82
83 /// Pending reply queues per remote address.
84 ///
85 /// Uses Weak refs so entries are automatically invalidated when `ReplyFuture` drops.
86 /// Cleaned lazily during `close_pending_replies`.
87 ///
88 /// FDB: endStreamOnDisconnect pattern (genericactors.actor.h:332)
89 pending_replies: BTreeMap<String, Vec<PendingReplyEntry>>,
90
91 /// Statistics.
92 stats: TransportStats,
93}
94
95impl<P: Providers> TransportData<P> {
96 fn new(providers: &P) -> Self {
97 Self {
98 endpoints: EndpointMap::new(),
99 peers: BTreeMap::new(),
100 incoming_peers: BTreeMap::new(),
101 failure_monitor: Arc::new(FailureMonitor::new(providers.time().clone())),
102 pending_replies: BTreeMap::new(),
103 stats: TransportStats::default(),
104 }
105 }
106}
107
108/// Central transport coordinator (FDB `NetTransport` equivalent).
109///
110/// # Design
111///
112/// - Manages peer connections lazily (created on first send)
113/// - Routes incoming packets to registered endpoints
114/// - Synchronous send API - queues immediately, returns (FDB pattern)
115/// - Explicit `Arc<NetTransport>` passing (testable, simulation-friendly)
116/// - Internal state in `TransportData` (matches FDB: `TransportData* self`)
117///
118/// # FDB Reference
119/// From NetTransport.h:195-314
120///
121/// # Multi-Node Support (Phase 12 Step 7d)
122///
123/// For multi-node operation, wrap in `Arc` and call `set_weak_self()`:
124/// ```ignore
125/// let transport = Arc::new(NetTransport::new(...));
126/// transport.set_weak_self(Arc::downgrade(&transport));
127/// transport.listen().await?; // Start accepting connections
128/// ```
129pub struct NetTransport<P: Providers, C: MessageCodec = crate::JsonCodec> {
130 /// Internal transport data (FDB: `TransportData`* self).
131 data: RwLock<TransportData<P>>,
132
133 /// Local address for this transport.
134 local_address: NetworkAddress,
135
136 /// Providers bundle for network, time, task, and random.
137 providers: P,
138
139 /// Message codec for serialization/deserialization.
140 codec: C,
141
142 /// Peer configuration.
143 peer_config: PeerConfig,
144
145 /// Weak self-reference for spawning background tasks.
146 /// Required for `connection_reader` tasks to dispatch back to this transport.
147 /// Set via `set_weak_self()` after wrapping in `Arc`.
148 weak_self: RwLock<Option<Weak<Self>>>,
149
150 /// Shutdown signal. Cancelled on drop so background tasks (`listen_task`)
151 /// exit. `CancellationToken` has deterministic FIFO wake order, unlike
152 /// `tokio::sync::watch` whose shard pick draws from tokio's runtime RNG
153 /// (entropy-seeded outside a tokio runtime).
154 shutdown: CancellationToken,
155
156 /// Weak reference as `dyn TransportHandle` for codec-erased types.
157 /// Set by the builder alongside `weak_self`.
158 weak_handle: RwLock<Option<std::sync::Weak<dyn super::transport_handle::TransportHandle>>>,
159}
160
161/// Statistics for the transport.
162#[derive(Debug, Default)]
163struct TransportStats {
164 /// Number of packets sent.
165 packets_sent: u64,
166 /// Number of packets dispatched to endpoints.
167 packets_dispatched: u64,
168 /// Number of packets that couldn't be delivered (endpoint not found).
169 packets_undelivered: u64,
170 /// Number of peers created.
171 peers_created: u64,
172}
173
174impl<P: Providers + Send + Sync, C: MessageCodec> NetTransport<P, C> {
175 /// Create a new `NetTransport`.
176 ///
177 /// # Arguments
178 ///
179 /// * `local_address` - Address of this transport (for local delivery checks)
180 /// * `providers` - Providers bundle for network, time, task, and random
181 /// * `codec` - Message codec for serialization/deserialization
182 ///
183 /// # Multi-Node Usage
184 ///
185 /// For multi-node operation, wrap in `Arc` and call `set_weak_self()`:
186 /// ```ignore
187 /// let transport = Arc::new(NetTransport::new(...));
188 /// transport.set_weak_self(Arc::downgrade(&transport));
189 /// ```
190 pub fn new(local_address: NetworkAddress, providers: P, codec: C) -> Self {
191 let data = TransportData::new(&providers);
192 Self {
193 data: RwLock::new(data),
194 local_address,
195 providers,
196 codec,
197 peer_config: PeerConfig::default(),
198 weak_self: RwLock::new(None),
199 shutdown: CancellationToken::new(),
200 weak_handle: RwLock::new(None),
201 }
202 }
203
204 /// Get a reference to the transport's codec.
205 pub fn codec(&self) -> &C {
206 &self.codec
207 }
208
209 /// Get the random provider for generating unique IDs.
210 pub fn random(&self) -> &P::Random {
211 self.providers.random()
212 }
213
214 /// Get the providers bundle.
215 pub fn providers(&self) -> &P {
216 &self.providers
217 }
218
219 /// Allocate a random base token for a dynamic interface.
220 ///
221 /// Each method in the interface uses `base.adjusted(method_index)`.
222 /// Method indices start at 1 (0 is the base identity).
223 pub fn allocate_interface_token(&self) -> UID {
224 use crate::RandomProvider as _;
225 UID::new(
226 self.providers.random().random::<u64>(),
227 self.providers.random().random::<u64>(),
228 )
229 }
230
231 /// Set the weak self-reference for background task spawning.
232 ///
233 /// Required for multi-node operation where `connection_reader` tasks need
234 /// to dispatch incoming packets back to this transport.
235 ///
236 /// # FDB Pattern
237 /// Similar to how FDB actors reference `TransportData* self`.
238 ///
239 /// # Example
240 /// ```ignore
241 /// let transport = Arc::new(NetTransport::new(...));
242 /// transport.set_weak_self(Arc::downgrade(&transport));
243 /// ```
244 ///
245 /// # Panics
246 ///
247 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
248 /// task panicked while holding the lock).
249 pub fn set_weak_self(&self, weak: Weak<Self>) {
250 *self
251 .weak_self
252 .write()
253 .expect("RwLock poisoned: prior task panicked") = Some(weak);
254 }
255
256 /// Get weak self-reference, panics if not set.
257 fn weak_self(&self) -> Weak<Self> {
258 self.weak_self
259 .read()
260 .expect("RwLock poisoned: prior task panicked")
261 .clone()
262 .expect("weak_self not set - call set_weak_self() after wrapping in Arc")
263 }
264
265 /// Create with custom peer configuration.
266 #[must_use]
267 pub fn with_peer_config(mut self, config: PeerConfig) -> Self {
268 self.peer_config = config;
269 self
270 }
271
272 /// Get the local address.
273 pub fn local_address(&self) -> &NetworkAddress {
274 &self.local_address
275 }
276
277 /// Get the failure monitor for this transport.
278 ///
279 /// Used by delivery mode functions to race replies against disconnect signals.
280 ///
281 /// # FDB Reference
282 /// `IFailureMonitor::failureMonitor()` (FailureMonitor.h)
283 ///
284 /// # Panics
285 ///
286 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
287 /// task panicked while holding the lock).
288 pub fn failure_monitor(&self) -> Arc<FailureMonitor> {
289 Arc::clone(
290 &self
291 .data
292 .read()
293 .expect("RwLock poisoned: prior task panicked")
294 .failure_monitor,
295 )
296 }
297
298 /// Register a well-known endpoint.
299 ///
300 /// Well-known endpoints have deterministic tokens for O(1) lookup.
301 ///
302 /// # Errors
303 ///
304 /// Returns [`MessagingError`] if a receiver is already registered under `token`.
305 ///
306 /// # Panics
307 ///
308 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
309 /// task panicked while holding the lock).
310 pub fn register_well_known(
311 &self,
312 token: WellKnownToken,
313 receiver: Arc<dyn MessageReceiver>,
314 ) -> Result<(), MessagingError> {
315 self.data
316 .write()
317 .expect("RwLock poisoned: prior task panicked")
318 .endpoints
319 .insert_well_known(token, receiver)
320 }
321
322 /// Register a dynamic endpoint with the given UID.
323 ///
324 /// Returns the endpoint that senders should use to address this receiver.
325 ///
326 /// # Panics
327 ///
328 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
329 /// task panicked while holding the lock).
330 pub fn register(&self, token: UID, receiver: Arc<dyn MessageReceiver>) -> Endpoint {
331 self.data
332 .write()
333 .expect("RwLock poisoned: prior task panicked")
334 .endpoints
335 .insert(token, receiver);
336 Endpoint::new(self.local_address.clone(), token)
337 }
338
339 /// Unregister a dynamic endpoint.
340 ///
341 /// # Panics
342 ///
343 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
344 /// task panicked while holding the lock).
345 pub fn unregister(&self, token: &UID) -> Option<Arc<dyn MessageReceiver>> {
346 self.data
347 .write()
348 .expect("RwLock poisoned: prior task panicked")
349 .endpoints
350 .remove(token)
351 }
352
353 /// Register a pending reply queue for a remote address.
354 ///
355 /// Called by `send_request` to track which reply queues should be closed
356 /// when a peer disconnects. Uses Weak refs so stale entries are cleaned lazily.
357 pub(crate) fn register_pending_reply(
358 &self,
359 addr: &str,
360 token: UID,
361 closer: Weak<dyn ReplyQueueCloser>,
362 ) {
363 let mut data = self
364 .data
365 .write()
366 .expect("RwLock poisoned: prior task panicked");
367 data.pending_replies
368 .entry(addr.to_string())
369 .or_default()
370 .push((token, closer));
371 }
372
373 /// Close all pending reply queues for a disconnected address.
374 ///
375 /// Upgrades Weak refs and calls `close_with_error`, then removes the
376 /// corresponding endpoints from the `EndpointMap`.
377 ///
378 /// FDB: endStreamOnDisconnect pattern (genericactors.actor.h:332)
379 pub(crate) fn close_pending_replies(&self, addr: &str, reason: &ReplyError) {
380 let entries = {
381 let mut data = self
382 .data
383 .write()
384 .expect("RwLock poisoned: prior task panicked");
385 data.pending_replies.remove(addr).unwrap_or_default()
386 };
387
388 if entries.is_empty() {
389 return;
390 }
391
392 tracing::debug!(
393 "close_pending_replies: closing {} reply queues for {}",
394 entries.len(),
395 addr,
396 );
397
398 for (token, weak_closer) in entries {
399 if let Some(closer) = weak_closer.upgrade() {
400 closer.close_with_error(reason.clone());
401 }
402 // Remove from endpoint map regardless (either closed or already dropped)
403 self.data
404 .write()
405 .expect("RwLock poisoned: prior task panicked")
406 .endpoints
407 .remove(&token);
408 }
409 }
410
411 /// Register a typed request handler in a single step.
412 ///
413 /// This is the preferred method for registering RPC handlers. It combines:
414 /// - Creating an endpoint from the local address and token
415 /// - Creating a `RequestStream` for the request type
416 /// - Registering the stream's queue with the transport
417 ///
418 /// The codec is taken from the transport (set at builder time).
419 ///
420 /// # Example
421 ///
422 /// ```rust,ignore
423 /// let stream = NetTransport::register_handler::<PingRequest, PingResponse>(
424 /// &transport, ping_token(),
425 /// );
426 ///
427 /// loop {
428 /// if let Some((req, reply)) = stream.recv().await {
429 /// reply.send(PingResponse { pong: req.ping });
430 /// }
431 /// }
432 /// ```
433 ///
434 /// # Panics
435 ///
436 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
437 /// task panicked while holding the lock).
438 pub fn register_handler<Req, Resp>(
439 transport: &Arc<Self>,
440 token: UID,
441 ) -> RequestStream<Req, Resp>
442 where
443 Req: DeserializeOwned + Send + Sync + 'static,
444 Resp: Serialize + Send + Sync + 'static,
445 {
446 let endpoint = Endpoint::new(transport.local_address.clone(), token);
447 let handle: Arc<dyn super::transport_handle::TransportHandle> =
448 Arc::clone(transport) as Arc<dyn super::transport_handle::TransportHandle>;
449 let stream = RequestStream::new(endpoint, transport.codec.clone(), handle);
450 transport
451 .data
452 .write()
453 .expect("RwLock poisoned: prior task panicked")
454 .endpoints
455 .insert(token, stream.queue() as Arc<dyn MessageReceiver>);
456 stream
457 }
458
459 /// Register a handler for a multi-method interface.
460 ///
461 /// The token is computed deterministically from `interface_id` and `method_index`.
462 /// The codec is taken from the transport.
463 ///
464 /// # Example
465 ///
466 /// ```rust,ignore
467 /// let (add_stream, _) = NetTransport::register_handler_at::<AddRequest, AddResponse>(
468 /// &transport, CALC_INTERFACE, METHOD_ADD,
469 /// );
470 ///
471 /// loop {
472 /// moonpool_core::select! {
473 /// Some((req, reply)) = add_stream.recv() => {
474 /// reply.send(AddResponse { result: req.a + req.b });
475 /// }
476 /// }
477 /// }
478 /// ```
479 pub fn register_handler_at<Req, Resp>(
480 transport: &Arc<Self>,
481 interface_id: u64,
482 method_index: u64,
483 ) -> (RequestStream<Req, Resp>, UID)
484 where
485 Req: DeserializeOwned + Send + Sync + 'static,
486 Resp: Serialize + Send + Sync + 'static,
487 {
488 let token = UID::new(interface_id, method_index);
489 let stream = Self::register_handler(transport, token);
490 (stream, token)
491 }
492
493 /// Send packet unreliably (best-effort, dropped on failure).
494 ///
495 /// This is a synchronous operation - it queues the packet and returns immediately.
496 /// The packet may be dropped if the connection fails (FDB pattern).
497 ///
498 /// # Arguments
499 ///
500 /// * `endpoint` - Destination endpoint
501 /// * `payload` - Message bytes (already serialized)
502 ///
503 /// # Errors
504 ///
505 /// Returns [`MessagingError`] if local delivery fails or the peer rejects the packet.
506 ///
507 /// # Panics
508 ///
509 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
510 /// task panicked while holding the lock).
511 pub fn send_unreliable(
512 &self,
513 endpoint: &Endpoint,
514 payload: &[u8],
515 ) -> Result<(), MessagingError> {
516 // Check for local delivery
517 if self.is_local_address(&endpoint.address) {
518 return self.deliver_local(&endpoint.token, payload);
519 }
520
521 // Get or create peer for remote address
522 let peer = self.get_or_open_peer(&endpoint.address);
523 peer.write()
524 .expect("RwLock poisoned: prior task panicked")
525 .send_unreliable(endpoint.token, payload)
526 .map_err(|e| MessagingError::PeerError {
527 message: e.to_string(),
528 })?;
529
530 self.data
531 .write()
532 .expect("RwLock poisoned: prior task panicked")
533 .stats
534 .packets_sent += 1;
535 Ok(())
536 }
537
538 /// Send packet reliably (queued, will retry on reconnect).
539 ///
540 /// This is a synchronous operation - it queues the packet and returns immediately.
541 /// The packet will be retried on connection failure (FDB pattern).
542 ///
543 /// # Arguments
544 ///
545 /// * `endpoint` - Destination endpoint
546 /// * `payload` - Message bytes (already serialized)
547 ///
548 /// # Errors
549 ///
550 /// Returns [`MessagingError`] if local delivery fails or the peer rejects the packet.
551 ///
552 /// # Panics
553 ///
554 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
555 /// task panicked while holding the lock).
556 pub fn send_reliable(&self, endpoint: &Endpoint, payload: &[u8]) -> Result<(), MessagingError> {
557 // Check for local delivery
558 if self.is_local_address(&endpoint.address) {
559 return self.deliver_local(&endpoint.token, payload);
560 }
561
562 // Get or create peer for remote address
563 let peer = self.get_or_open_peer(&endpoint.address);
564 peer.write()
565 .expect("RwLock poisoned: prior task panicked")
566 .send_reliable(endpoint.token, payload)
567 .map_err(|e| MessagingError::PeerError {
568 message: e.to_string(),
569 })?;
570
571 self.data
572 .write()
573 .expect("RwLock poisoned: prior task panicked")
574 .stats
575 .packets_sent += 1;
576 Ok(())
577 }
578
579 /// Check if address is local (same as this transport).
580 pub(crate) fn is_local_address(&self, address: &NetworkAddress) -> bool {
581 self.local_address == *address
582 }
583
584 /// Deliver packet locally (same process, no network).
585 fn deliver_local(&self, token: &UID, payload: &[u8]) -> Result<(), MessagingError> {
586 let data = self
587 .data
588 .read()
589 .expect("RwLock poisoned: prior task panicked");
590 if let Some(receiver) = data.endpoints.get(token) {
591 receiver.receive(payload);
592 drop(data); // Release borrow before mutating stats
593 self.data
594 .write()
595 .expect("RwLock poisoned: prior task panicked")
596 .stats
597 .packets_dispatched += 1;
598 Ok(())
599 } else {
600 drop(data); // Release borrow before mutating stats
601 self.data
602 .write()
603 .expect("RwLock poisoned: prior task panicked")
604 .stats
605 .packets_undelivered += 1;
606 Err(MessagingError::EndpointNotFound { token: *token })
607 }
608 }
609
610 /// Get or create a peer for the given address.
611 ///
612 /// Peers are created lazily on first send (FDB connectionKeeper pattern).
613 /// When a new peer is created, a `connection_reader` task is spawned to
614 /// handle incoming packets (FDB: connectionKeeper spawns connectionReader at line 843).
615 ///
616 /// Also checks incoming peers (from accepted connections) as a fallback.
617 /// This enables response routing over existing connections — when a server
618 /// receives a request from a client, it can send the response back on the
619 /// same connection without requiring the client to listen for new connections.
620 ///
621 /// # FDB Reference
622 /// `Reference<struct Peer> getOrOpenPeer(NetworkAddress const& address);`
623 fn get_or_open_peer(&self, address: &NetworkAddress) -> SharedPeer<P> {
624 let addr_str = address.to_string();
625
626 // Check if outgoing peer already exists
627 if let Some(peer) = self
628 .data
629 .read()
630 .expect("RwLock poisoned: prior task panicked")
631 .peers
632 .get(&addr_str)
633 {
634 return Arc::clone(peer);
635 }
636
637 // Check incoming peers — reuse accepted connection for responses
638 // (FDB pattern: responses flow back on the same connection)
639 if let Some(peer) = self
640 .data
641 .read()
642 .expect("RwLock poisoned: prior task panicked")
643 .incoming_peers
644 .get(&addr_str)
645 {
646 return Arc::clone(peer);
647 }
648
649 // Create new peer with failure monitor
650 let fm = Some(Arc::clone(
651 &self
652 .data
653 .read()
654 .expect("RwLock poisoned: prior task panicked")
655 .failure_monitor,
656 ));
657 let peer = Peer::new(
658 &self.providers,
659 addr_str.clone(),
660 self.peer_config.clone(),
661 fm,
662 );
663 let peer = Arc::new(RwLock::new(peer));
664
665 // Store in peers map
666 {
667 let mut data = self
668 .data
669 .write()
670 .expect("RwLock poisoned: prior task panicked");
671 data.peers.insert(addr_str.clone(), Arc::clone(&peer));
672 data.stats.peers_created += 1;
673 }
674
675 // Spawn connection_reader for incoming packets (FDB pattern: connectionKeeper spawns connectionReader)
676 // This handles responses for outgoing requests
677 self.spawn_connection_reader(Arc::clone(&peer), addr_str);
678
679 peer
680 }
681
682 /// Dispatch an incoming packet to the appropriate endpoint.
683 ///
684 /// Called by the transport loop when a packet is received.
685 ///
686 /// # Returns
687 ///
688 /// Ok(()) if delivered, Err if endpoint not found.
689 ///
690 /// # Errors
691 ///
692 /// Returns [`MessagingError::EndpointNotFound`] if no receiver is registered
693 /// under `token`.
694 ///
695 /// # Panics
696 ///
697 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
698 /// task panicked while holding the lock).
699 pub fn dispatch(&self, token: &UID, payload: &[u8]) -> Result<(), MessagingError> {
700 let data = self
701 .data
702 .read()
703 .expect("RwLock poisoned: prior task panicked");
704 if let Some(receiver) = data.endpoints.get(token) {
705 receiver.receive(payload);
706 drop(data); // Release borrow before mutating stats
707 self.data
708 .write()
709 .expect("RwLock poisoned: prior task panicked")
710 .stats
711 .packets_dispatched += 1;
712 Ok(())
713 } else {
714 drop(data); // Release borrow before mutating stats
715 self.data
716 .write()
717 .expect("RwLock poisoned: prior task panicked")
718 .stats
719 .packets_undelivered += 1;
720 Err(MessagingError::EndpointNotFound { token: *token })
721 }
722 }
723
724 /// Get statistics.
725 ///
726 /// # Panics
727 ///
728 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
729 /// task panicked while holding the lock).
730 pub fn packets_sent(&self) -> u64 {
731 self.data
732 .read()
733 .expect("RwLock poisoned: prior task panicked")
734 .stats
735 .packets_sent
736 }
737
738 /// Get the number of packets dispatched to endpoints.
739 ///
740 /// # Panics
741 ///
742 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
743 /// task panicked while holding the lock).
744 pub fn packets_dispatched(&self) -> u64 {
745 self.data
746 .read()
747 .expect("RwLock poisoned: prior task panicked")
748 .stats
749 .packets_dispatched
750 }
751
752 /// Get the number of packets that couldn't be delivered.
753 ///
754 /// # Panics
755 ///
756 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
757 /// task panicked while holding the lock).
758 pub fn packets_undelivered(&self) -> u64 {
759 self.data
760 .read()
761 .expect("RwLock poisoned: prior task panicked")
762 .stats
763 .packets_undelivered
764 }
765
766 /// Get the number of peers created.
767 ///
768 /// # Panics
769 ///
770 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
771 /// task panicked while holding the lock).
772 pub fn peers_created(&self) -> u64 {
773 self.data
774 .read()
775 .expect("RwLock poisoned: prior task panicked")
776 .stats
777 .peers_created
778 }
779
780 /// Get number of active peers.
781 ///
782 /// # Panics
783 ///
784 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
785 /// task panicked while holding the lock).
786 pub fn peer_count(&self) -> usize {
787 self.data
788 .read()
789 .expect("RwLock poisoned: prior task panicked")
790 .peers
791 .len()
792 }
793
794 /// Get number of registered endpoints.
795 ///
796 /// # Panics
797 ///
798 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
799 /// task panicked while holding the lock).
800 pub fn endpoint_count(&self) -> usize {
801 let data = self
802 .data
803 .read()
804 .expect("RwLock poisoned: prior task panicked");
805 data.endpoints.well_known_count() + data.endpoints.dynamic_count()
806 }
807
808 /// Spawn a `connection_reader` for a peer.
809 ///
810 /// FDB Pattern: connectionKeeper spawns connectionReader (line 843).
811 /// The `connection_reader` reads from the peer and dispatches to endpoints.
812 ///
813 /// # Arguments
814 ///
815 /// * `peer` - The peer to read from
816 /// * `peer_addr` - Address string for logging
817 fn spawn_connection_reader(&self, peer: SharedPeer<P>, peer_addr: String) {
818 // Only spawn if weak_self is set (multi-node mode)
819 if self
820 .weak_self
821 .read()
822 .expect("RwLock poisoned: prior task panicked")
823 .is_none()
824 {
825 return;
826 }
827
828 let transport_weak = self.weak_self();
829 drop(self.providers.task().spawn_task(
830 "connection_reader",
831 connection_reader(transport_weak, peer, peer_addr),
832 ));
833 }
834
835 // =========================================================================
836 // Server Listener Support (FDB: listen + connectionIncoming)
837 // =========================================================================
838
839 /// Start listening for incoming connections.
840 ///
841 /// FDB Pattern: `listen()` (NetTransport.actor.cpp:1646-1676)
842 /// Binds to the local address and spawns an accept loop that handles
843 /// incoming connections via `connection_incoming()`.
844 ///
845 /// # Requirements
846 ///
847 /// Must call `set_weak_self()` before calling this method.
848 ///
849 /// # Example
850 ///
851 /// ```ignore
852 /// let transport = Arc::new(NetTransport::new(...));
853 /// transport.set_weak_self(Arc::downgrade(&transport));
854 /// transport.listen().await?;
855 /// ```
856 ///
857 /// # Errors
858 ///
859 /// Returns [`MessagingError`] if `set_weak_self()` was not called, if binding
860 /// the local address fails, or if accepting incoming connections fails.
861 ///
862 /// # Panics
863 ///
864 /// Panics if the internal `RwLock` is poisoned (only possible if a prior
865 /// task panicked while holding the lock).
866 pub async fn listen(&self) -> Result<(), MessagingError> {
867 // Verify weak_self is set
868 if self
869 .weak_self
870 .read()
871 .expect("RwLock poisoned: prior task panicked")
872 .is_none()
873 {
874 return Err(MessagingError::InvalidState {
875 message: "weak_self not set - call set_weak_self() before listen()".to_string(),
876 });
877 }
878
879 // Bind to local address
880 let addr_str = self.local_address.to_string();
881 let listener = self
882 .providers
883 .network()
884 .bind(&addr_str)
885 .await
886 .map_err(|e| MessagingError::NetworkError {
887 message: format!("Failed to bind to {addr_str}: {e}"),
888 })?;
889
890 tracing::info!("NetTransport: listening on {}", addr_str);
891
892 // Spawn listen task (FDB: listen() actor)
893 // Pass a shutdown token clone so the task can exit when transport is dropped
894 let transport_weak = self.weak_self();
895 let shutdown = self.shutdown.clone();
896 drop(self.providers.task().spawn_task(
897 "listen",
898 listen_task(transport_weak, listener, addr_str, shutdown),
899 ));
900
901 Ok(())
902 }
903}
904
905/// Implement Drop to signal shutdown to background tasks.
906///
907/// When `NetTransport` is dropped, we signal all background tasks (`listen_task`)
908/// to exit gracefully. This prevents tasks from being stuck on `accept()` forever.
909impl<P: Providers, C: MessageCodec> Drop for NetTransport<P, C> {
910 fn drop(&mut self) {
911 tracing::debug!("NetTransport: signaling shutdown to background tasks");
912 self.shutdown.cancel();
913 }
914}
915
916// =============================================================================
917// TransportHandle implementation
918// =============================================================================
919
920impl<P: Providers + Send + Sync, C: MessageCodec> super::transport_handle::TransportHandle
921 for NetTransport<P, C>
922{
923 fn send_unreliable(
924 &self,
925 endpoint: &crate::Endpoint,
926 payload: &[u8],
927 ) -> Result<(), crate::error::MessagingError> {
928 NetTransport::send_unreliable(self, endpoint, payload)
929 }
930
931 fn send_reliable(
932 &self,
933 endpoint: &crate::Endpoint,
934 payload: &[u8],
935 ) -> Result<(), crate::error::MessagingError> {
936 NetTransport::send_reliable(self, endpoint, payload)
937 }
938
939 fn register(
940 &self,
941 token: crate::UID,
942 receiver: Arc<dyn super::endpoint_map::MessageReceiver>,
943 ) -> crate::Endpoint {
944 NetTransport::register(self, token, receiver)
945 }
946
947 fn unregister(
948 &self,
949 token: &crate::UID,
950 ) -> Option<Arc<dyn super::endpoint_map::MessageReceiver>> {
951 NetTransport::unregister(self, token)
952 }
953
954 fn register_pending_reply(
955 &self,
956 addr: &str,
957 token: crate::UID,
958 closer: std::sync::Weak<dyn super::net_notified_queue::ReplyQueueCloser>,
959 ) {
960 NetTransport::register_pending_reply(self, addr, token, closer);
961 }
962
963 fn failure_monitor(&self) -> Arc<super::failure_monitor::FailureMonitor> {
964 NetTransport::failure_monitor(self)
965 }
966
967 fn local_address(&self) -> &crate::NetworkAddress {
968 NetTransport::local_address(self)
969 }
970
971 fn allocate_interface_token(&self) -> crate::UID {
972 NetTransport::allocate_interface_token(self)
973 }
974
975 fn random_uid(&self) -> crate::UID {
976 use crate::RandomProvider as _;
977 crate::UID::new(
978 self.providers.random().random::<u64>(),
979 self.providers.random().random::<u64>(),
980 )
981 }
982
983 fn is_local_address(&self, address: &crate::NetworkAddress) -> bool {
984 NetTransport::is_local_address(self, address)
985 }
986
987 fn weak_for_cleanup(&self) -> std::sync::Weak<dyn super::transport_handle::TransportHandle> {
988 self.weak_handle
989 .read()
990 .expect("RwLock poisoned: prior task panicked")
991 .clone()
992 .unwrap_or_else(|| std::sync::Weak::<Self>::new())
993 }
994}
995
996impl<P: Providers + Send + Sync, C: MessageCodec> NetTransport<P, C> {
997 /// Set the weak handle reference for `dyn TransportHandle` usage.
998 ///
999 /// Called by the builder after wrapping in `Arc`.
1000 pub(crate) fn set_weak_handle(
1001 &self,
1002 weak: std::sync::Weak<dyn super::transport_handle::TransportHandle>,
1003 ) {
1004 *self
1005 .weak_handle
1006 .write()
1007 .expect("RwLock poisoned: prior task panicked") = Some(weak);
1008 }
1009}
1010
1011// =============================================================================
1012// NetTransportBuilder
1013// =============================================================================
1014
1015/// Builder for `NetTransport` that eliminates common footguns.
1016///
1017/// The manual `Arc` wrapping and `set_weak_self()` pattern is error-prone:
1018/// forgetting `set_weak_self()` causes a runtime panic. This builder handles
1019/// both automatically.
1020///
1021/// # Examples
1022///
1023/// ```rust,ignore
1024/// // Standard usage - server or client that needs RPC responses:
1025/// let transport = NetTransportBuilder::new(providers)
1026/// .local_address(addr)
1027/// .build_listening()
1028/// .await?;
1029///
1030/// // Fire-and-forget sender (no listening needed):
1031/// let transport = NetTransportBuilder::new(providers)
1032/// .local_address(addr)
1033/// .build();
1034///
1035/// // With custom peer config:
1036/// let transport = NetTransportBuilder::new(providers)
1037/// .local_address(addr)
1038/// .peer_config(config)
1039/// .build_listening()
1040/// .await?;
1041/// ```
1042///
1043/// # Why Build vs Build Listening?
1044///
1045/// - **`build_listening()`**: For most RPC use cases. Both servers AND clients
1046/// need this because responses are sent to the client's listening address.
1047///
1048/// - **`build()`**: For fire-and-forget messaging where you don't expect responses.
1049/// Also useful for testing where you control message flow manually.
1050pub struct NetTransportBuilder<P: Providers, C: MessageCodec = crate::JsonCodec> {
1051 providers: P,
1052 codec: C,
1053 local_address: Option<NetworkAddress>,
1054 peer_config: Option<PeerConfig>,
1055}
1056
1057impl<P: Providers> NetTransportBuilder<P, crate::JsonCodec> {
1058 /// Create a new builder with the providers bundle.
1059 ///
1060 /// Uses [`JsonCodec`](crate::JsonCodec) by default. Call [`.codec()`](Self::codec)
1061 /// to use a different codec.
1062 ///
1063 /// # Arguments
1064 ///
1065 /// * `providers` - Providers bundle for network, time, task, and random
1066 pub fn new(providers: P) -> Self {
1067 Self {
1068 providers,
1069 codec: crate::JsonCodec,
1070 local_address: None,
1071 peer_config: None,
1072 }
1073 }
1074}
1075
1076/// A [`NetTransport`] backed by the production [`TokioProviders`](moonpool_core::TokioProviders)
1077/// bundle (and the default [`JsonCodec`](crate::JsonCodec)) — the typical
1078/// production transport type.
1079#[cfg(feature = "tokio")]
1080pub type TokioTransport = NetTransport<moonpool_core::TokioProviders>;
1081
1082#[cfg(feature = "tokio")]
1083impl NetTransportBuilder<moonpool_core::TokioProviders, crate::JsonCodec> {
1084 /// Create a builder backed by the production
1085 /// [`TokioProviders`](moonpool_core::TokioProviders) bundle.
1086 ///
1087 /// Shorthand for `NetTransportBuilder::new(TokioProviders::new())` — the
1088 /// usual production entry point.
1089 #[must_use]
1090 pub fn tokio() -> Self {
1091 Self::new(moonpool_core::TokioProviders::new())
1092 }
1093}
1094
1095impl<P: Providers + Send + Sync, C: MessageCodec> NetTransportBuilder<P, C> {
1096 /// Set the message codec for this transport.
1097 ///
1098 /// The codec is used for all serialization/deserialization in RPC handlers
1099 /// and client endpoints created from this transport.
1100 pub fn codec<NewC: MessageCodec>(self, codec: NewC) -> NetTransportBuilder<P, NewC> {
1101 NetTransportBuilder {
1102 providers: self.providers,
1103 codec,
1104 local_address: self.local_address,
1105 peer_config: self.peer_config,
1106 }
1107 }
1108
1109 /// Set the local address for this transport.
1110 ///
1111 /// This is required before calling `build()` or `build_listening()`.
1112 ///
1113 /// # Arguments
1114 ///
1115 /// * `address` - The network address to bind to
1116 #[must_use]
1117 pub fn local_address(mut self, address: NetworkAddress) -> Self {
1118 self.local_address = Some(address);
1119 self
1120 }
1121
1122 /// Set custom peer configuration.
1123 ///
1124 /// If not set, uses `PeerConfig::default()`.
1125 #[must_use]
1126 pub fn peer_config(mut self, config: PeerConfig) -> Self {
1127 self.peer_config = Some(config);
1128 self
1129 }
1130
1131 /// Build the transport without starting the listener.
1132 ///
1133 /// Returns `Arc<NetTransport>` with `set_weak_self()` already called.
1134 /// Use this for fire-and-forget messaging or testing.
1135 ///
1136 /// For RPC (request/response), use `build_listening()` instead.
1137 ///
1138 /// # Errors
1139 ///
1140 /// Returns `MessagingError::MissingLocalAddress` if `local_address()` was not called.
1141 pub fn build(self) -> Result<Arc<NetTransport<P, C>>, MessagingError> {
1142 let address = self
1143 .local_address
1144 .ok_or(MessagingError::MissingLocalAddress)?;
1145
1146 let mut transport = NetTransport::new(address, self.providers, self.codec);
1147
1148 if let Some(config) = self.peer_config {
1149 transport = transport.with_peer_config(config);
1150 }
1151
1152 let transport = Arc::new(transport);
1153 transport.set_weak_self(Arc::downgrade(&transport));
1154 let handle: Arc<dyn super::transport_handle::TransportHandle> =
1155 transport.clone() as Arc<dyn super::transport_handle::TransportHandle>;
1156 transport.set_weak_handle(Arc::downgrade(&handle));
1157 Ok(transport)
1158 }
1159
1160 /// Build the transport and start listening for incoming connections.
1161 ///
1162 /// Returns `Arc<NetTransport>` with `set_weak_self()` already called
1163 /// and the listener started.
1164 ///
1165 /// Use this for typical RPC usage where you need to receive responses.
1166 /// Both servers AND clients need this in a request/response pattern.
1167 ///
1168 /// # Errors
1169 ///
1170 /// Returns `MessagingError::MissingLocalAddress` if `local_address()` was not called,
1171 /// or a network error if binding to the local address fails.
1172 pub async fn build_listening(self) -> Result<Arc<NetTransport<P, C>>, MessagingError> {
1173 let transport = self.build()?;
1174 transport.listen().await?;
1175 Ok(transport)
1176 }
1177}
1178
1179/// FDB: `connectionReader()` - reads from connection and dispatches to endpoints.
1180///
1181/// This is a background task that:
1182/// 1. Takes ownership of the peer's receiver channel at startup
1183/// 2. Reads incoming packets from the channel (no peer lock held during await)
1184/// 3. Dispatches them to the appropriate endpoint via `transport.dispatch()`
1185///
1186/// # FDB Reference
1187/// From NetTransport.actor.cpp:1401-1602 connectionReader
1188///
1189/// The FDB version is more complex (handles `ConnectPacket`, protocol negotiation),
1190/// but the core loop is: read packets → scanPackets → deliver.
1191async fn connection_reader<P: Providers + Send + Sync, C: MessageCodec>(
1192 transport: Weak<NetTransport<P, C>>,
1193 peer: SharedPeer<P>,
1194 peer_addr: String,
1195) {
1196 tracing::debug!("connection_reader: started for peer {}", peer_addr);
1197
1198 // Take ownership of the receiver at startup.
1199 // This avoids holding the peer lock across await points (critical safety fix).
1200 let mut receiver = {
1201 if let Some(rx) = peer
1202 .write()
1203 .expect("RwLock poisoned: prior task panicked")
1204 .take_receiver()
1205 {
1206 rx
1207 } else {
1208 tracing::error!(
1209 "connection_reader: receiver already taken for peer {}",
1210 peer_addr
1211 );
1212 return;
1213 }
1214 }; // peer lock released here
1215
1216 loop {
1217 // Await directly on the owned receiver - safe, no peer lock involved!
1218 if let Some((token, payload)) = receiver.recv().await {
1219 // Try to get transport reference
1220 let Some(transport) = transport.upgrade() else {
1221 tracing::debug!(
1222 "connection_reader: transport dropped, exiting for peer {}",
1223 peer_addr
1224 );
1225 break;
1226 };
1227
1228 // FDB: deliver() - looks up endpoint and dispatches
1229 if let Err(e) = transport.dispatch(&token, &payload) {
1230 tracing::debug!(
1231 "connection_reader: dispatch failed for token {}: {:?}",
1232 token,
1233 e
1234 );
1235
1236 // FDB: send WLTOKEN_ENDPOINT_NOT_FOUND back to sender
1237 // (FlowTransport.cpp:1244-1225)
1238 if !token.is_well_known() {
1239 let mut notification = [0u8; 16];
1240 notification[0..8].copy_from_slice(&token.first.to_le_bytes());
1241 notification[8..16].copy_from_slice(&token.second.to_le_bytes());
1242 let _ = peer
1243 .write()
1244 .expect("RwLock poisoned: prior task panicked")
1245 .send_unreliable(
1246 crate::peer::core::ENDPOINT_NOT_FOUND_TOKEN,
1247 ¬ification,
1248 );
1249 }
1250 }
1251 } else {
1252 // Channel closed - peer disconnected or shutdown
1253 tracing::debug!("connection_reader: peer {} receiver closed", peer_addr);
1254 // Close all pending reply queues for this peer with MaybeDelivered
1255 if let Some(transport) = transport.upgrade() {
1256 transport.close_pending_replies(&peer_addr, &ReplyError::MaybeDelivered);
1257 }
1258 break;
1259 }
1260 }
1261
1262 tracing::debug!("connection_reader: exiting for peer {}", peer_addr);
1263}
1264
1265/// FDB: `listen()` - accept loop spawning connectionIncoming per connection.
1266///
1267/// This is a background task that:
1268/// 1. Accepts incoming connections
1269/// 2. Spawns `connection_incoming` for each accepted connection
1270/// 3. Exits gracefully when shutdown signal is received
1271///
1272/// # FDB Reference
1273/// From NetTransport.actor.cpp:1646-1676 listen
1274async fn listen_task<P: Providers + Send + Sync, C: MessageCodec>(
1275 transport: Weak<NetTransport<P, C>>,
1276 listener: <P::Network as NetworkProvider>::TcpListener,
1277 listen_addr: String,
1278 shutdown: CancellationToken,
1279) {
1280 tracing::debug!("listen_task: started on {}", listen_addr);
1281
1282 loop {
1283 // Use select! to race between accept and shutdown signal
1284 moonpool_core::select! {
1285 // Bias toward shutdown to ensure timely exit
1286 biased;
1287
1288 // Check shutdown signal first
1289 () = shutdown.cancelled() => {
1290 tracing::debug!("listen_task: shutdown signal received, exiting for {}", listen_addr);
1291 break;
1292 }
1293
1294 // Accept next connection
1295 accept_result = listener.accept() => {
1296 match accept_result {
1297 Ok((stream, peer_addr)) => {
1298 tracing::debug!(
1299 "listen_task: accepted connection from {} on {}",
1300 peer_addr,
1301 listen_addr
1302 );
1303
1304 // Get transport reference
1305 let Some(transport_rc) = transport.upgrade() else {
1306 tracing::debug!("listen_task: transport dropped, exiting");
1307 break;
1308 };
1309
1310 // Handle the incoming connection (FDB: connectionIncoming)
1311 connection_incoming(
1312 Arc::downgrade(&transport_rc),
1313 stream,
1314 peer_addr,
1315 &transport_rc,
1316 );
1317
1318 }
1319 Err(e) => {
1320 tracing::warn!("listen_task: accept error on {}: {:?}", listen_addr, e);
1321 // Continue accepting - transient errors are expected
1322 }
1323 }
1324 }
1325 }
1326 }
1327
1328 tracing::debug!("listen_task: exiting for {}", listen_addr);
1329}
1330
1331/// FDB: `connectionIncoming()` - handles accepted connection.
1332///
1333/// This function:
1334/// 1. Creates a peer for the incoming connection
1335/// 2. Spawns a `connection_reader` for the peer
1336///
1337/// # FDB Reference
1338/// From NetTransport.actor.cpp:1604-1644 connectionIncoming
1339///
1340/// Note: FDB's connectionIncoming waits for `ConnectPacket` to identify the peer.
1341/// We simplify by using the peer address directly since `SimNetworkProvider`
1342/// already provides the peer address.
1343///
1344/// FDB Pattern: Use `Peer::new_incoming()` with the accepted stream, not `Peer::new()`.
1345/// This uses the already-established connection rather than trying to connect back.
1346fn connection_incoming<P: Providers + Send + Sync, C: MessageCodec>(
1347 transport_weak: Weak<NetTransport<P, C>>,
1348 stream: <P::Network as NetworkProvider>::TcpStream,
1349 peer_addr: String,
1350 transport: &NetTransport<P, C>,
1351) {
1352 tracing::debug!(
1353 "connection_incoming: handling connection from {}",
1354 peer_addr
1355 );
1356
1357 // Check if we already have a peer for this address (FDB: getOrOpenPeer in connectionReader:1555)
1358 // For incoming connections, we store in incoming_peers to avoid conflicts with outgoing peers
1359 //
1360 // Note: Unlike outgoing peers which can be reused, incoming peers with new streams
1361 // should replace the old one since the old connection is stale.
1362 let peer = {
1363 let data = transport
1364 .data
1365 .read()
1366 .expect("RwLock poisoned: prior task panicked");
1367 if data.incoming_peers.contains_key(&peer_addr) {
1368 tracing::debug!(
1369 "connection_incoming: replacing stale incoming peer for {}",
1370 peer_addr
1371 );
1372 }
1373 drop(data); // Release borrow
1374
1375 // FDB Pattern: Use Peer::new_incoming() with the accepted stream
1376 // (NetTransport.actor.cpp:1123 Peer::onIncomingConnection)
1377 // This uses the already-established connection rather than trying to connect back.
1378 let fm = Some(Arc::clone(
1379 &transport
1380 .data
1381 .read()
1382 .expect("RwLock poisoned: prior task panicked")
1383 .failure_monitor,
1384 ));
1385 let peer = Peer::new_incoming(
1386 &transport.providers,
1387 peer_addr.clone(),
1388 stream,
1389 transport.peer_config.clone(),
1390 fm,
1391 );
1392 let peer = Arc::new(RwLock::new(peer));
1393
1394 // Store in incoming_peers (replaces any existing stale peer)
1395 transport
1396 .data
1397 .write()
1398 .expect("RwLock poisoned: prior task panicked")
1399 .incoming_peers
1400 .insert(peer_addr.clone(), Arc::clone(&peer));
1401
1402 tracing::debug!(
1403 "connection_incoming: created new incoming peer for {}",
1404 peer_addr
1405 );
1406 peer
1407 };
1408
1409 // Spawn connection_reader to handle incoming packets
1410 drop(transport.providers.task().spawn_task(
1411 "connection_reader",
1412 connection_reader(transport_weak, peer, peer_addr),
1413 ));
1414}
1415
1416#[cfg(test)]
1417mod tests {
1418 use std::net::{IpAddr, Ipv4Addr};
1419
1420 use super::*;
1421 use crate::{
1422 JsonCodec, NetNotifiedQueue, TokioRandomProvider, TokioStorageProvider, TokioTaskProvider,
1423 TokioTimeProvider,
1424 };
1425
1426 // Simple mock network provider that fails all connections
1427 // (we only test local delivery, so connections are never actually made)
1428 #[derive(Clone)]
1429 struct MockNetworkProvider;
1430
1431 // Dummy stream type for the mock
1432 struct DummyStream;
1433
1434 impl futures::io::AsyncRead for DummyStream {
1435 fn poll_read(
1436 self: std::pin::Pin<&mut Self>,
1437 _cx: &mut std::task::Context<'_>,
1438 _buf: &mut [u8],
1439 ) -> std::task::Poll<std::io::Result<usize>> {
1440 std::task::Poll::Ready(Err(std::io::Error::other("dummy stream")))
1441 }
1442 }
1443
1444 impl futures::io::AsyncWrite for DummyStream {
1445 fn poll_write(
1446 self: std::pin::Pin<&mut Self>,
1447 _cx: &mut std::task::Context<'_>,
1448 _buf: &[u8],
1449 ) -> std::task::Poll<std::io::Result<usize>> {
1450 std::task::Poll::Ready(Err(std::io::Error::other("dummy stream")))
1451 }
1452
1453 fn poll_flush(
1454 self: std::pin::Pin<&mut Self>,
1455 _cx: &mut std::task::Context<'_>,
1456 ) -> std::task::Poll<std::io::Result<()>> {
1457 std::task::Poll::Ready(Err(std::io::Error::other("dummy stream")))
1458 }
1459
1460 fn poll_close(
1461 self: std::pin::Pin<&mut Self>,
1462 _cx: &mut std::task::Context<'_>,
1463 ) -> std::task::Poll<std::io::Result<()>> {
1464 std::task::Poll::Ready(Err(std::io::Error::other("dummy stream")))
1465 }
1466 }
1467
1468 impl std::marker::Unpin for DummyStream {}
1469
1470 // Dummy listener type for the mock
1471 struct DummyListener;
1472
1473 impl crate::TcpListenerTrait for DummyListener {
1474 type TcpStream = DummyStream;
1475
1476 async fn accept(&self) -> std::io::Result<(Self::TcpStream, String)> {
1477 Err(std::io::Error::other("dummy listener"))
1478 }
1479
1480 fn local_addr(&self) -> std::io::Result<String> {
1481 Err(std::io::Error::other("dummy listener"))
1482 }
1483 }
1484
1485 impl NetworkProvider for MockNetworkProvider {
1486 type TcpStream = DummyStream;
1487 type TcpListener = DummyListener;
1488
1489 async fn bind(&self, _addr: &str) -> std::io::Result<Self::TcpListener> {
1490 Err(std::io::Error::other("mock bind"))
1491 }
1492
1493 async fn connect(&self, _addr: &str) -> std::io::Result<Self::TcpStream> {
1494 Err(std::io::Error::other("mock connection"))
1495 }
1496 }
1497
1498 /// Mock providers bundle for testing
1499 #[derive(Clone)]
1500 struct MockProviders {
1501 network: MockNetworkProvider,
1502 time: TokioTimeProvider,
1503 task: TokioTaskProvider,
1504 random: TokioRandomProvider,
1505 storage: TokioStorageProvider,
1506 }
1507
1508 impl MockProviders {
1509 fn new() -> Self {
1510 Self {
1511 network: MockNetworkProvider,
1512 time: TokioTimeProvider::new(),
1513 task: TokioTaskProvider,
1514 random: TokioRandomProvider::new(),
1515 storage: TokioStorageProvider::new(),
1516 }
1517 }
1518 }
1519
1520 impl Providers for MockProviders {
1521 type Network = MockNetworkProvider;
1522 type Time = TokioTimeProvider;
1523 type Task = TokioTaskProvider;
1524 type Random = TokioRandomProvider;
1525 type Storage = TokioStorageProvider;
1526
1527 fn network(&self) -> &Self::Network {
1528 &self.network
1529 }
1530 fn time(&self) -> &Self::Time {
1531 &self.time
1532 }
1533 fn task(&self) -> &Self::Task {
1534 &self.task
1535 }
1536 fn random(&self) -> &Self::Random {
1537 &self.random
1538 }
1539 fn storage(&self) -> &Self::Storage {
1540 &self.storage
1541 }
1542 }
1543
1544 fn test_address() -> NetworkAddress {
1545 NetworkAddress::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 4500)
1546 }
1547
1548 fn create_test_transport() -> NetTransport<MockProviders> {
1549 NetTransport::new(test_address(), MockProviders::new(), JsonCodec)
1550 }
1551
1552 #[test]
1553 fn test_new_transport() {
1554 let transport = create_test_transport();
1555
1556 assert_eq!(transport.peer_count(), 0);
1557 assert_eq!(transport.endpoint_count(), 0);
1558 assert_eq!(transport.packets_sent(), 0);
1559 }
1560
1561 #[test]
1562 fn test_register_well_known() {
1563 let transport = create_test_transport();
1564
1565 let queue: Arc<NetNotifiedQueue<String>> = Arc::new(NetNotifiedQueue::with_codec(
1566 Endpoint::new(test_address(), UID::new(1, 1)),
1567 JsonCodec,
1568 ));
1569
1570 transport
1571 .register_well_known(WellKnownToken::Ping, queue)
1572 .expect("register should succeed");
1573
1574 assert_eq!(transport.endpoint_count(), 1);
1575 }
1576
1577 #[test]
1578 fn test_register_dynamic() {
1579 let transport = create_test_transport();
1580
1581 let token = UID::new(0x1234, 0x5678);
1582 let queue: Arc<NetNotifiedQueue<String>> = Arc::new(NetNotifiedQueue::with_codec(
1583 Endpoint::new(test_address(), token),
1584 JsonCodec,
1585 ));
1586
1587 let endpoint = transport.register(token, queue);
1588
1589 assert_eq!(endpoint.token, token);
1590 assert_eq!(transport.endpoint_count(), 1);
1591 }
1592
1593 #[test]
1594 fn test_local_delivery() {
1595 let transport = create_test_transport();
1596
1597 let token = UID::new(0x1234, 0x5678);
1598 let queue: Arc<NetNotifiedQueue<String>> = Arc::new(NetNotifiedQueue::with_codec(
1599 Endpoint::new(test_address(), token),
1600 JsonCodec,
1601 ));
1602
1603 let endpoint = transport.register(token, Arc::clone(&queue) as Arc<dyn MessageReceiver>);
1604
1605 // Send to local endpoint
1606 let payload = br#""hello local""#;
1607 transport
1608 .send_unreliable(&endpoint, payload)
1609 .expect("send should succeed");
1610
1611 assert_eq!(transport.packets_dispatched(), 1);
1612 assert_eq!(queue.try_recv(), Some("hello local".to_string()));
1613 }
1614
1615 #[test]
1616 fn test_endpoint_not_found() {
1617 let transport = create_test_transport();
1618
1619 let endpoint = Endpoint::new(test_address(), UID::new(999, 999));
1620 let payload = br#""test""#;
1621
1622 let result = transport.send_unreliable(&endpoint, payload);
1623 assert!(matches!(
1624 result,
1625 Err(MessagingError::EndpointNotFound { .. })
1626 ));
1627 assert_eq!(transport.packets_undelivered(), 1);
1628 }
1629
1630 #[test]
1631 fn test_unregister() {
1632 let transport = create_test_transport();
1633
1634 let token = UID::new(0x1234, 0x5678);
1635 let queue: Arc<NetNotifiedQueue<String>> = Arc::new(NetNotifiedQueue::with_codec(
1636 Endpoint::new(test_address(), token),
1637 JsonCodec,
1638 ));
1639
1640 transport.register(token, queue as Arc<dyn MessageReceiver>);
1641 assert_eq!(transport.endpoint_count(), 1);
1642
1643 let removed = transport.unregister(&token);
1644 assert!(removed.is_some());
1645 assert_eq!(transport.endpoint_count(), 0);
1646 }
1647
1648 #[test]
1649 fn test_dispatch() {
1650 let transport = create_test_transport();
1651
1652 let token = UID::new(0x1234, 0x5678);
1653 let queue: Arc<NetNotifiedQueue<String>> = Arc::new(NetNotifiedQueue::with_codec(
1654 Endpoint::new(test_address(), token),
1655 JsonCodec,
1656 ));
1657
1658 transport.register(token, Arc::clone(&queue) as Arc<dyn MessageReceiver>);
1659
1660 // Dispatch directly
1661 let payload = br#""dispatched""#;
1662 transport
1663 .dispatch(&token, payload)
1664 .expect("dispatch should succeed");
1665
1666 assert_eq!(queue.try_recv(), Some("dispatched".to_string()));
1667 }
1668
1669 // =========================================================================
1670 // NetTransportBuilder tests
1671 // =========================================================================
1672
1673 #[test]
1674 fn test_builder_build() {
1675 let transport = NetTransportBuilder::new(MockProviders::new())
1676 .local_address(test_address())
1677 .build()
1678 .expect("build should succeed");
1679
1680 // Should be properly initialized
1681 assert_eq!(transport.peer_count(), 0);
1682 assert_eq!(transport.endpoint_count(), 0);
1683 assert_eq!(transport.packets_sent(), 0);
1684 assert_eq!(*transport.local_address(), test_address());
1685 }
1686
1687 #[test]
1688 fn test_builder_weak_self_set() {
1689 let transport = NetTransportBuilder::new(MockProviders::new())
1690 .local_address(test_address())
1691 .build()
1692 .expect("build should succeed");
1693
1694 // weak_self should be set - verify by checking it doesn't panic
1695 // This verifies the builder correctly calls set_weak_self()
1696 assert!(
1697 transport
1698 .weak_self
1699 .read()
1700 .expect("RwLock poisoned: prior task panicked")
1701 .is_some()
1702 );
1703 }
1704
1705 #[test]
1706 fn test_builder_with_peer_config() {
1707 let config = PeerConfig::default();
1708 let transport = NetTransportBuilder::new(MockProviders::new())
1709 .local_address(test_address())
1710 .peer_config(config)
1711 .build()
1712 .expect("build should succeed");
1713
1714 // Transport should be created (peer_config is internal, but creation succeeds)
1715 assert_eq!(transport.peer_count(), 0);
1716 }
1717
1718 #[test]
1719 fn test_builder_missing_address_returns_error() {
1720 let result = NetTransportBuilder::new(MockProviders::new()).build();
1721
1722 assert!(matches!(result, Err(MessagingError::MissingLocalAddress)));
1723 }
1724
1725 #[test]
1726 fn test_builder_local_delivery_works() {
1727 // Verify the built transport functions correctly
1728 let transport = NetTransportBuilder::new(MockProviders::new())
1729 .local_address(test_address())
1730 .build()
1731 .expect("build should succeed");
1732
1733 let token = UID::new(0x1234, 0x5678);
1734 let queue: Arc<NetNotifiedQueue<String>> = Arc::new(NetNotifiedQueue::with_codec(
1735 Endpoint::new(test_address(), token),
1736 JsonCodec,
1737 ));
1738
1739 let endpoint = transport.register(token, Arc::clone(&queue) as Arc<dyn MessageReceiver>);
1740
1741 // Send to local endpoint
1742 let payload = br#""builder test""#;
1743 transport
1744 .send_unreliable(&endpoint, payload)
1745 .expect("send should succeed");
1746
1747 assert_eq!(transport.packets_dispatched(), 1);
1748 assert_eq!(queue.try_recv(), Some("builder test".to_string()));
1749 }
1750
1751 #[test]
1752 fn test_register_handler() {
1753 use serde::{Deserialize, Serialize};
1754
1755 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1756 struct TestRequest {
1757 value: i32,
1758 }
1759
1760 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1761 struct TestResponse {
1762 value: i32,
1763 }
1764
1765 let transport = NetTransportBuilder::new(MockProviders::new())
1766 .local_address(test_address())
1767 .build()
1768 .expect("build should succeed");
1769
1770 let token = UID::new(0xDEAD, 0xBEEF);
1771 let stream = NetTransport::register_handler::<TestRequest, TestResponse>(&transport, token);
1772
1773 // Handler should be registered
1774 assert_eq!(transport.endpoint_count(), 1);
1775
1776 // Endpoint should match
1777 assert_eq!(stream.endpoint().token, token);
1778 assert_eq!(stream.endpoint().address, test_address());
1779 }
1780
1781 #[test]
1782 fn test_register_handler_at_multi_method() {
1783 use serde::{Deserialize, Serialize};
1784
1785 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1786 struct AddRequest {
1787 a: i32,
1788 b: i32,
1789 }
1790
1791 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1792 struct AddResponse {
1793 result: i32,
1794 }
1795
1796 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1797 struct SubRequest {
1798 a: i32,
1799 b: i32,
1800 }
1801
1802 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1803 struct SubResponse {
1804 result: i32,
1805 }
1806
1807 const CALC_INTERFACE: u64 = 0xCA1C;
1808 const METHOD_ADD: u64 = 0;
1809 const METHOD_SUB: u64 = 1;
1810
1811 let transport = NetTransportBuilder::new(MockProviders::new())
1812 .local_address(test_address())
1813 .build()
1814 .expect("build should succeed");
1815
1816 // Register multiple handlers for the same interface
1817 let (add_stream, add_token) = NetTransport::register_handler_at::<AddRequest, AddResponse>(
1818 &transport,
1819 CALC_INTERFACE,
1820 METHOD_ADD,
1821 );
1822 let (sub_stream, sub_token) = NetTransport::register_handler_at::<SubRequest, SubResponse>(
1823 &transport,
1824 CALC_INTERFACE,
1825 METHOD_SUB,
1826 );
1827
1828 // Both handlers should be registered
1829 assert_eq!(transport.endpoint_count(), 2);
1830
1831 // Tokens should be deterministic
1832 assert_eq!(add_token, UID::new(CALC_INTERFACE, METHOD_ADD));
1833 assert_eq!(sub_token, UID::new(CALC_INTERFACE, METHOD_SUB));
1834
1835 // Streams should have correct endpoints
1836 assert_eq!(add_stream.endpoint().token, add_token);
1837 assert_eq!(sub_stream.endpoint().token, sub_token);
1838 }
1839
1840 // =========================================================================
1841 // Builder error path tests (Phase 12D Step 18)
1842 // =========================================================================
1843
1844 #[tokio::test]
1845 async fn test_build_listening_bind_error() {
1846 // MockNetworkProvider returns error on bind(), simulating port already in use
1847 let result = NetTransportBuilder::new(MockProviders::new())
1848 .local_address(test_address())
1849 .build_listening()
1850 .await;
1851
1852 // Should return NetworkError when bind fails
1853 assert!(matches!(
1854 result,
1855 Err(crate::error::MessagingError::NetworkError { .. })
1856 ));
1857 }
1858
1859 // =========================================================================
1860 // EndpointNotFound notification tests
1861 // =========================================================================
1862
1863 #[test]
1864 fn test_dispatch_not_found_for_well_known_token() {
1865 let transport = create_test_transport();
1866
1867 // Dispatch to a well-known token that's not registered
1868 let well_known = UID::well_known(42);
1869 let payload = b"test";
1870
1871 let result = transport.dispatch(&well_known, payload);
1872 assert!(matches!(
1873 result,
1874 Err(MessagingError::EndpointNotFound { .. })
1875 ));
1876 // The guard in connection_reader checks is_well_known() before
1877 // sending notification — well-known tokens should NOT trigger notifications.
1878 assert!(well_known.is_well_known());
1879 }
1880
1881 #[test]
1882 fn test_dispatch_not_found_for_dynamic_token() {
1883 let transport = create_test_transport();
1884
1885 // Dispatch to a dynamic (non-well-known) token that's not registered
1886 let dynamic_token = UID::new(0xCAFE, 0xBABE);
1887 let payload = b"test";
1888
1889 let result = transport.dispatch(&dynamic_token, payload);
1890 assert!(matches!(
1891 result,
1892 Err(MessagingError::EndpointNotFound { .. })
1893 ));
1894 // Dynamic tokens SHOULD trigger the notification in connection_reader.
1895 assert!(!dynamic_token.is_well_known());
1896 }
1897}