Skip to main content

moirai_async/net/
types.rs

1//! Network configuration and statistics types.
2
3#![expect(
4    clippy::unwrap_used,
5    reason = "ratchet MOIRAI-UNWRAP-1: pre-existing debt"
6)]
7
8use std::collections::HashMap;
9use std::sync::Mutex;
10use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
11use std::time::{Duration, Instant};
12
13/// Configuration for TCP server behavior
14#[derive(Debug, Clone)]
15pub struct TcpServerConfig {
16    /// Maximum number of concurrent connections
17    pub max_connections: Option<usize>,
18    /// Socket TCP_NODELAY setting, applied to every accepted stream.
19    pub nodelay: bool,
20    /// TCP keep-alive duration.
21    ///
22    /// Currently not applied: neither `std::net::TcpStream` nor
23    /// `moirai_pal::net::AsyncTcpStream` exposes a keep-alive setter (it
24    /// requires `SO_KEEPALIVE`/`TCP_KEEPIDLE` support in the PAL). Pending PAL
25    /// wiring; until then the value is configuration-only.
26    pub keep_alive: Option<Duration>,
27    /// Connection timeout.
28    ///
29    /// Currently not applied: accepted sockets are non-blocking (reactor
30    /// driven), where `SO_RCVTIMEO`/`SO_SNDTIMEO` have no effect. Async
31    /// deadlines are expressed by wrapping operations in
32    /// [`crate::timer::timeout()`]; automatic application of this value is
33    /// pending.
34    pub timeout: Option<Duration>,
35}
36
37impl Default for TcpServerConfig {
38    fn default() -> Self {
39        Self {
40            max_connections: Some(1000),
41            nodelay: true,
42            keep_alive: Some(Duration::from_secs(300)),
43            timeout: Some(Duration::from_secs(30)),
44        }
45    }
46}
47
48/// TCP server statistics for monitoring
49#[derive(Debug, Default)]
50pub struct ServerStats {
51    /// Connections accepted over the server's lifetime.
52    pub total_connections: AtomicU64,
53    /// Connections currently tracked as open.
54    pub active_connections: AtomicU64,
55    /// Bytes received across all connections.
56    pub bytes_received: AtomicU64,
57    /// Bytes sent across all connections.
58    pub bytes_sent: AtomicU64,
59}
60
61/// Unique, monotonically-assigned identifier for a tracked connection.
62///
63/// Connections are keyed by id rather than by peer [`std::net::SocketAddr`] so that
64/// (a) a stream can be removed from the pool at drop time without re-querying
65/// the (possibly already-reset) socket, and (b) two connections sharing a peer
66/// address (NAT, rapid address reuse) cannot collide in the tracking map.
67pub type ConnectionId = u64;
68
69/// Connection information tracking
70#[derive(Debug, Clone)]
71pub struct ConnectionInfo {
72    /// Peer address captured at accept/connect time (never re-queried).
73    pub peer_addr: std::net::SocketAddr,
74    /// Instant the connection entered the pool.
75    pub connected_at: Instant,
76    /// Bytes received on this connection.
77    pub bytes_received: u64,
78    /// Bytes sent on this connection.
79    pub bytes_sent: u64,
80    /// Instant of the most recent recorded activity.
81    pub last_activity: Instant,
82}
83
84/// Connection pool for managing active connections
85#[derive(Debug)]
86pub struct ConnectionPool {
87    active_connections: Mutex<HashMap<ConnectionId, ConnectionInfo>>,
88    reserved_connections: AtomicUsize,
89    next_connection_id: AtomicU64,
90    max_connections: Option<usize>,
91}
92
93impl ConnectionPool {
94    /// Create a pool bounded to `max_connections`, or unbounded on `None`.
95    #[must_use]
96    pub fn new(max_connections: Option<usize>) -> Self {
97        Self {
98            active_connections: Mutex::new(HashMap::new()),
99            reserved_connections: AtomicUsize::new(0),
100            next_connection_id: AtomicU64::new(0),
101            max_connections,
102        }
103    }
104
105    /// Reserve one connection slot ahead of an accept.
106    ///
107    /// Returns false when the pool (active plus reserved) is at capacity.
108    pub fn try_reserve(&self) -> bool {
109        let max = match self.max_connections {
110            Some(m) => m,
111            None => return true,
112        };
113
114        let connections = self.active_connections.lock().unwrap();
115        let current = connections.len();
116        // The mutex orders all admission increments. Releases may race this
117        // snapshot, but they only reduce the reservation count; a stale release
118        // can conservatively reject an admission, never over-admit one. The
119        // counter carries no payload, so Relaxed is sufficient for its atomic
120        // accounting.
121        let reserved = self.reserved_connections.load(Ordering::Relaxed);
122        if current + reserved < max {
123            self.reserved_connections.fetch_add(1, Ordering::Relaxed);
124            true
125        } else {
126            false
127        }
128    }
129
130    /// Release a slot taken by [`Self::try_reserve`] without admitting a
131    /// connection.
132    pub fn cancel_reservation(&self) {
133        if self.max_connections.is_some() {
134            self.reserved_connections.fetch_sub(1, Ordering::Relaxed);
135        }
136    }
137
138    /// Register a connection and return its unique id. The id is what the owning
139    /// stream stores and later passes to [`Self::remove_connection`].
140    pub fn add_connection(&self, addr: std::net::SocketAddr) -> ConnectionId {
141        let id = self.next_connection_id.fetch_add(1, Ordering::Relaxed);
142        let now = Instant::now();
143        self.active_connections.lock().unwrap().insert(
144            id,
145            ConnectionInfo {
146                peer_addr: addr,
147                connected_at: now,
148                bytes_received: 0,
149                bytes_sent: 0,
150                last_activity: now,
151            },
152        );
153        id
154    }
155
156    /// Convert a successful reservation into a tracked connection, returning the
157    /// new connection id. Releases exactly the one reservation taken by
158    /// [`Self::try_reserve`].
159    pub fn add_connection_reserved(&self, addr: std::net::SocketAddr) -> ConnectionId {
160        let id = self.add_connection(addr);
161        if self.max_connections.is_some() {
162            self.reserved_connections.fetch_sub(1, Ordering::Relaxed);
163        }
164        id
165    }
166
167    /// Record I/O activity on a tracked connection, updating its byte counters
168    /// and `last_activity` timestamp. No-op when the connection is no longer
169    /// tracked (already removed or never pool-tracked).
170    pub fn record_io(&self, id: ConnectionId, bytes_received: u64, bytes_sent: u64) {
171        let mut connections = self.active_connections.lock().unwrap();
172        if let Some(info) = connections.get_mut(&id) {
173            info.bytes_received += bytes_received;
174            info.bytes_sent += bytes_sent;
175            info.last_activity = Instant::now();
176        }
177    }
178
179    /// Remove a tracked connection; returns whether it was present.
180    pub fn remove_connection(&self, id: ConnectionId) -> bool {
181        self.active_connections
182            .lock()
183            .unwrap()
184            .remove(&id)
185            .is_some()
186    }
187
188    /// Return whether the pool can admit another connection.
189    pub fn has_capacity(&self) -> bool {
190        match self.max_connections {
191            Some(max) => {
192                let current = self.connection_count();
193                let reserved = self.reserved_connections.load(Ordering::Relaxed);
194                current + reserved < max
195            }
196            None => true,
197        }
198    }
199
200    /// Snapshot the tracked connections.
201    pub fn get_active_connections(&self) -> HashMap<ConnectionId, ConnectionInfo> {
202        self.active_connections.lock().unwrap().clone()
203    }
204
205    /// Count of connections currently tracked as open.
206    pub fn connection_count(&self) -> usize {
207        self.active_connections.lock().unwrap().len()
208    }
209}
210
211/// Statistics for individual TCP connections
212#[derive(Debug, Clone)]
213pub struct ConnectionStats {
214    /// Bytes read on the connection.
215    pub bytes_read: u64,
216    /// Bytes written on the connection.
217    pub bytes_written: u64,
218    /// Completed read operations.
219    pub read_ops: u64,
220    /// Completed write operations.
221    pub write_ops: u64,
222}
223
224/// Configuration for UDP socket behavior
225#[derive(Debug, Clone)]
226pub struct UdpConfig {
227    /// Socket buffer size.
228    ///
229    /// Currently not applied: `std::net::UdpSocket` exposes no
230    /// `SO_RCVBUF`/`SO_SNDBUF` setter (requires PAL/socket2-level support).
231    /// Pending PAL wiring; until then the value is configuration-only.
232    pub buffer_size: usize,
233    /// Broadcast support, applied at bind time via `SO_BROADCAST`.
234    pub broadcast: bool,
235    /// Multicast support.
236    ///
237    /// Currently not applied: joining a multicast group requires a group
238    /// address, which this flag cannot carry, and the PAL exposes no
239    /// `join_multicast` surface. Pending a typed multicast configuration
240    /// (group + interface) in place of this flag.
241    pub multicast: bool,
242}
243
244impl Default for UdpConfig {
245    fn default() -> Self {
246        Self {
247            buffer_size: 65536,
248            broadcast: false,
249            multicast: false,
250        }
251    }
252}
253
254/// Statistics for UDP socket operations
255#[derive(Debug, Default)]
256pub struct UdpStats {
257    /// Datagrams sent.
258    pub packets_sent: AtomicU64,
259    /// Datagrams received.
260    pub packets_received: AtomicU64,
261    /// Bytes sent.
262    pub bytes_sent: AtomicU64,
263    /// Bytes received.
264    pub bytes_received: AtomicU64,
265}
266
267/// Public TCP server statistics
268#[derive(Debug, Clone)]
269pub struct TcpServerStats {
270    /// Connections accepted over the server's lifetime.
271    pub total_connections: u64,
272    /// Connections currently open.
273    pub active_connections: u64,
274    /// Bytes received across all connections.
275    pub bytes_received: u64,
276    /// Bytes sent across all connections.
277    pub bytes_sent: u64,
278}
279
280/// Public UDP socket statistics  
281#[derive(Debug, Clone)]
282pub struct UdpSocketStats {
283    /// Datagrams sent.
284    pub packets_sent: u64,
285    /// Datagrams received.
286    pub packets_received: u64,
287    /// Bytes sent.
288    pub bytes_sent: u64,
289    /// Bytes received.
290    pub bytes_received: u64,
291}