Skip to main content

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                            &notification,
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}