saluki_io/net/resource.rs
1//! Registry-managed network listeners.
2//!
3//! A listening network socket is the canonical scarce resource: only one subsystem in the process can hold it,
4//! releasing it hands it back to the operating system, and re-binding can fail because something else took it in the
5//! meantime. The specifications here let a [`ResourceRegistry`] own listeners on behalf of the whole process and lend
6//! them out, so a component can be torn down and rebuilt without its sockets ever being released.
7
8use std::num::NonZeroUsize;
9
10use async_trait::async_trait;
11use saluki_core::runtime::state::{ResourceKind, ResourceSpecification, Subleases};
12use saluki_error::GenericError;
13use stringtheory::MetaString;
14
15use super::{
16 listener::{ConnectionOrientedListener, Listener},
17 util::retry::ExponentialBackoff,
18 ListenAddress,
19};
20
21/// Specification for a general-purpose network listener.
22///
23/// Creates a [`Listener`], which handles both connection-oriented and connectionless address families.
24///
25/// Only the listen address identifies the resource. Tuning -- the UDP socket count, the receive buffer size, the accept
26/// backoff -- is applied when the listener is first created and has no effect on a later acquisition of the same
27/// address, since the sockets are already bound. The registry logs a warning if a subsequent specification disagrees.
28#[derive(Clone, Debug)]
29pub struct SocketSpecification {
30 address: ListenAddress,
31 udp_streams: Option<NonZeroUsize>,
32 receive_buffer_size: Option<usize>,
33 accept_backoff: Option<ExponentialBackoff>,
34}
35
36impl SocketSpecification {
37 /// Creates a specification for the given listen address.
38 pub fn new(address: ListenAddress) -> Self {
39 Self {
40 address,
41 udp_streams: None,
42 receive_buffer_size: None,
43 accept_backoff: None,
44 }
45 }
46
47 /// Sets how many UDP sockets to bind to the address.
48 ///
49 /// Each socket is bound independently with `SO_REUSEPORT` so the kernel load-balances incoming datagrams across
50 /// them, and the listener yields one stream per socket. Ignored for non-UDP addresses. See
51 /// [`Listener::from_listen_address`] for the platform caveats.
52 pub fn with_udp_streams(mut self, udp_streams: Option<NonZeroUsize>) -> Self {
53 self.udp_streams = udp_streams;
54 self
55 }
56
57 /// Sets the receive buffer size applied to streams accepted from this listener.
58 ///
59 /// `None` keeps the operating system default.
60 pub fn with_receive_buffer_size(mut self, receive_buffer_size: Option<usize>) -> Self {
61 self.receive_buffer_size = receive_buffer_size;
62 self
63 }
64
65 /// Sets how long the listener waits between accepts when the system is out of resources.
66 ///
67 /// `None` keeps the listener's default. See [`Listener::with_accept_backoff`].
68 pub fn with_accept_backoff(mut self, accept_backoff: Option<ExponentialBackoff>) -> Self {
69 self.accept_backoff = accept_backoff;
70 self
71 }
72
73 /// Returns the listen address this specification names.
74 pub fn listen_address(&self) -> &ListenAddress {
75 &self.address
76 }
77}
78
79impl From<ListenAddress> for SocketSpecification {
80 fn from(address: ListenAddress) -> Self {
81 Self::new(address)
82 }
83}
84
85#[async_trait]
86impl ResourceSpecification for SocketSpecification {
87 type Resource = Listener;
88
89 const KIND: ResourceKind = ResourceKind::Socket;
90
91 fn key(&self) -> MetaString {
92 MetaString::from(self.address.to_string())
93 }
94
95 async fn create(&self, subleases: Subleases) -> Result<Self::Resource, GenericError> {
96 let listener = Listener::from_listen_address(self.address.clone(), self.udp_streams)
97 .await?
98 .with_receive_buffer_size(self.receive_buffer_size)
99 // For a connectionless family the bound socket *is* the stream, so every stream this listener yields
100 // shares its socket and can outlive the lease it was acquired under. Each takes a sublease, which keeps
101 // the registry from handing this listener on while a stream is still reading from the socket.
102 .with_subleases(subleases);
103
104 Ok(match self.accept_backoff.clone() {
105 Some(accept_backoff) => listener.with_accept_backoff(accept_backoff),
106 None => listener,
107 })
108 }
109
110 fn reset(listener: &mut Self::Resource) {
111 listener.rearm();
112 }
113}
114
115/// Specification for a connection-oriented network listener.
116///
117/// Creates a [`ConnectionOrientedListener`], which accepts connections on TCP or Unix stream addresses only.
118///
119/// Only the listen address identifies the resource. The accept backoff is applied when the listener is first created
120/// and has no effect on a later acquisition of the same address, since the socket is already bound. The registry logs
121/// a warning if a subsequent specification disagrees.
122#[derive(Clone, Debug)]
123pub struct ConnectionOrientedSocketSpecification {
124 address: ListenAddress,
125 accept_backoff: Option<ExponentialBackoff>,
126}
127
128impl ConnectionOrientedSocketSpecification {
129 /// Creates a specification for the given listen address.
130 pub fn new(address: ListenAddress) -> Self {
131 Self {
132 address,
133 accept_backoff: None,
134 }
135 }
136
137 /// Sets the backoff the listener applies when accepting fails for want of a system resource.
138 ///
139 /// `None` keeps the listener's default. See [`ConnectionOrientedListener::with_accept_backoff`].
140 pub fn with_accept_backoff(mut self, accept_backoff: Option<ExponentialBackoff>) -> Self {
141 self.accept_backoff = accept_backoff;
142 self
143 }
144
145 /// Returns the listen address this specification names.
146 pub fn listen_address(&self) -> &ListenAddress {
147 &self.address
148 }
149}
150
151impl From<ListenAddress> for ConnectionOrientedSocketSpecification {
152 fn from(address: ListenAddress) -> Self {
153 Self::new(address)
154 }
155}
156
157#[async_trait]
158impl ResourceSpecification for ConnectionOrientedSocketSpecification {
159 type Resource = ConnectionOrientedListener;
160
161 const KIND: ResourceKind = ResourceKind::Socket;
162
163 fn key(&self) -> MetaString {
164 MetaString::from(self.address.to_string())
165 }
166
167 async fn create(&self, _subleases: Subleases) -> Result<Self::Resource, GenericError> {
168 // Never subdivided, so there is nothing to sublet: `accept` mints a new socket per connection, and the
169 // listener keeps nothing that outlives the lease it was acquired under.
170 let listener = ConnectionOrientedListener::from_listen_address(self.address.clone()).await?;
171
172 Ok(match self.accept_backoff.clone() {
173 Some(accept_backoff) => listener.with_accept_backoff(accept_backoff),
174 None => listener,
175 })
176 }
177
178 fn reset(listener: &mut Self::Resource) {
179 listener.rearm();
180 }
181}
182
183#[cfg(test)]
184mod tests;