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;