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}