1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
use std::{fmt, io, net, sync::Arc};
use ntex_io::Io;
use ntex_rt::System;
use ntex_service::{IntoService, Service, cfg::SharedCfg};
use ntex_util::time::Millis;
use socket2::{Domain, SockAddr, Socket, Type};
use crate::{NoConfig, Server, ServerAppConfig, WorkerPool};
use super::accept::AcceptLoop;
use super::config::ServiceConfig;
use super::factory::{self, FactoryServiceType};
use super::{Connection, ServerStatus, StreamServer, Token, socket::Listener};
/// Builder for a network server.
///
/// Register listeners and their service factories, configure the worker pool,
/// and call [`run`](Self::run) to start the server.
pub struct ServerBuilder<Cfg = NoConfig> {
name: String,
token: Token,
backlog: i32,
state: Arc<Cfg>,
services: Vec<FactoryServiceType<Cfg>>,
sockets: Vec<(Token, String, Listener)>,
accept: AcceptLoop,
pool: WorkerPool,
}
impl Default for ServerBuilder {
fn default() -> Self {
Self::new(NoConfig)
}
}
impl<Cfg> ServerBuilder<Cfg>
where
Cfg: ServerAppConfig,
{
#[must_use]
/// Creates a server builder with the specified application configuration.
pub fn new(cfg: Cfg) -> ServerBuilder<Cfg> {
let sys = System::current();
let mut accept = AcceptLoop::default();
accept.name(sys.name());
if sys.testing() {
accept.testing();
}
ServerBuilder {
accept,
name: sys.name().to_string(),
token: Token(0),
state: Arc::new(cfg),
services: Vec::new(),
sockets: Vec::new(),
backlog: 2048,
pool: WorkerPool::default().name(sys.name()),
}
}
#[must_use]
/// Sets the server name.
///
/// The name is also used for the accept and worker thread names. It
/// defaults to the current system name.
pub fn name<T: AsRef<str>>(mut self, name: T) -> Self {
self.name = name.as_ref().to_string();
self.accept.name(self.name.as_str());
self.pool = self.pool.name(self.name.as_str());
self
}
#[must_use]
/// Sets the number of worker threads to start.
///
/// By default, the server uses the number of available logical CPUs.
pub fn workers(mut self, num: usize) -> Self {
self.pool = self.pool.workers(num);
self
}
#[must_use]
/// Sets the maximum number of pending connections.
///
/// This refers to the number of clients that can be waiting to be served.
/// Exceeding this number results in the client getting an error when
/// attempting to connect. It should only affect servers under significant
/// load.
///
/// Generally set in the 64-2048 range. Default value is 2048.
///
/// It applies to listeners created by later [`bind`](Self::bind) and
/// [`configure`](Self::configure) calls. It does not affect listeners
/// passed to [`listen`](Self::listen) or `listen_uds`.
pub fn backlog(mut self, num: i32) -> Self {
self.backlog = num;
self
}
#[must_use]
/// Sets the maximum per-worker number of concurrent connections.
///
/// A worker stops taking new connections while it is at this limit. When
/// no worker can take a connection, the listeners stop accepting.
///
/// The limit is a process-wide setting shared by every server in the
/// process. Set it before the server starts, because each worker reads
/// it when its first service is created.
///
/// The default is 25,600 connections per worker.
pub fn max_connections(self, num: usize) -> Self {
super::max_concurrent_connections(num);
self
}
#[must_use]
/// Stops the current ntex runtime after the server has stopped.
///
/// By default, "stop runtime" is disabled.
pub fn stop_runtime(mut self) -> Self {
self.pool = self.pool.stop_runtime();
self
}
#[must_use]
/// Stops the server when one of the workers fails.
///
/// A worker fails when it panics or its service cannot be created. The
/// stop is graceful only if [`graceful_shutdown`](Self::graceful_shutdown)
/// is enabled. Without this option, a failed worker is restarted.
///
/// By default, "stop on panic" is disabled.
pub fn stop_on_panic(mut self) -> Self {
self.pool = self.pool.stop_on_panic();
self
}
#[must_use]
/// Disables signal handling.
///
/// By default, the server stops on SIGINT, SIGTERM, and SIGQUIT.
pub fn disable_signals(mut self) -> Self {
self.pool = self.pool.disable_signals();
self
}
#[must_use]
/// Enables CPU affinity for worker threads.
///
/// By default, affinity is disabled.
pub fn enable_affinity(mut self) -> Self {
self.pool = self.pool.enable_affinity();
self
}
#[must_use]
/// Enables graceful shutdown on SIGQUIT, fatal signals, and panics.
///
/// When enabled, SIGQUIT, SIGSEGV, SIGABRT, application panics, and
/// worker failures with "stop on panic" stop the server gracefully.
/// SIGTERM always stops gracefully and SIGINT always stops immediately.
///
/// By default, these events stop the server immediately.
pub fn graceful_shutdown(mut self) -> Self {
self.pool = self.pool.graceful_shutdown();
self
}
#[must_use]
/// Timeout for graceful worker shutdown.
///
/// After receiving a stop signal, workers have this much time to finish
/// serving requests. Workers that are still alive after the timeout are
/// forcefully dropped.
///
/// This bounds the worker as a whole, not an individual connection. Each
/// connection is bound separately by `IoConfig::set_shutdown_timeout`, so
/// this value should leave room for the connections a worker is still
/// draining to shut down themselves.
///
/// By default, the timeout is set to 30 seconds.
pub fn graceful_shutdown_timeout<T: Into<Millis>>(mut self, timeout: T) -> Self {
self.pool = self.pool.graceful_shutdown_timeout(timeout);
self
}
#[must_use]
/// Sets the server status handler.
///
/// The handler runs on the accept thread. It receives
/// [`ServerStatus::Ready`] when the listeners resume accepting and
/// [`ServerStatus::NotReady`] when they pause. The same status may be
/// reported more than once.
pub fn status_handler<F>(mut self, handler: F) -> Self
where
F: FnMut(ServerStatus) + Send + 'static,
{
self.accept.set_status_handler(handler);
self
}
/// Runs asynchronous configuration as part of the server building
/// process.
///
/// Listeners registered on the [`ServiceConfig`] are added to the server.
/// Services for them are attached per worker in
/// [`ServiceConfig::on_worker_start`]. This is useful for moving parts of
/// the configuration to a different module or library.
pub async fn configure<F>(mut self, f: F) -> io::Result<Self>
where
F: AsyncFnOnce(ServiceConfig<Cfg>) -> io::Result<()>,
{
let cfg = ServiceConfig::new(self.token, self.backlog);
f(cfg.clone()).await?;
let (token, sockets, factory) = cfg.into_factory();
self.token = token;
self.sockets.extend(sockets);
self.services.push(factory);
Ok(self)
}
#[allow(clippy::needless_pass_by_value)]
/// Binds TCP listeners and registers a service factory.
///
/// A listener is created for every address resolved from `addr`. Binding
/// succeeds if at least one of them binds; addresses that fail to bind
/// are skipped.
///
/// `cfg` is the I/O configuration for accepted connections. `factory` is
/// called once per worker with that worker's application state and
/// returns the connection service.
pub fn bind<F, S, I>(
mut self,
name: impl AsRef<str>,
addr: impl net::ToSocketAddrs,
cfg: impl Into<SharedCfg>,
factory: F,
) -> io::Result<Self>
where
F: AsyncFn(&Cfg::State) -> I + Send + Clone + 'static,
S: Service<Cfg::State, Io> + 'static,
I: IntoService<S, Cfg::State, Io> + 'static,
{
let cfg = cfg.into();
let sockets = bind_addr(addr, self.backlog)?;
let mut tokens = Vec::new();
for lst in sockets {
let token = self.token.next();
self.sockets
.push((token, name.as_ref().to_string(), Listener::from_tcp(lst)));
tokens.push((token, cfg.clone()));
}
self.services.push(factory::create_factory_service(
name.as_ref().to_string(),
tokens,
factory,
));
Ok(self)
}
#[cfg(unix)]
/// Binds a Unix domain socket and registers a service factory.
///
/// Any existing file at `addr` is removed before binding. The socket file
/// is removed again when the server stops. See [`bind`](Self::bind) for
/// `cfg` and `factory`.
pub fn bind_uds<F, I, S>(
self,
name: impl AsRef<str>,
addr: impl AsRef<std::path::Path>,
cfg: impl Into<SharedCfg>,
factory: F,
) -> io::Result<Self>
where
F: AsyncFn(&Cfg::State) -> I + Send + Clone + 'static,
I: IntoService<S, Cfg::State, Io> + 'static,
S: Service<Cfg::State, Io> + 'static,
{
use std::os::unix::net::UnixListener;
// The path must not exist when we try to bind.
// Try to remove it to avoid bind error.
if let Err(e) = std::fs::remove_file(addr.as_ref()) {
// NotFound is expected and not an issue. Anything else is.
if e.kind() != std::io::ErrorKind::NotFound {
return Err(e);
}
}
let lst = UnixListener::bind(addr)?;
self.listen_uds(name, lst, cfg.into(), factory)
}
#[cfg(unix)]
/// Registers a service factory for an existing Unix domain listener.
///
/// This is useful for socket activation, including listeners acquired
/// through systemd. The listener is switched to non-blocking mode. See
/// [`bind`](Self::bind) for `cfg` and `factory`.
pub fn listen_uds<F, I, S>(
mut self,
name: impl AsRef<str>,
lst: std::os::unix::net::UnixListener,
cfg: impl Into<SharedCfg>,
factory: F,
) -> io::Result<Self>
where
F: AsyncFn(&Cfg::State) -> I + Send + Clone + 'static,
I: IntoService<S, Cfg::State, Io> + 'static,
S: Service<Cfg::State, Io> + 'static,
{
let token = self.token.next();
self.services.push(factory::create_factory_service(
name.as_ref().to_string(),
vec![(token, cfg.into())],
factory,
));
self.sockets
.push((token, name.as_ref().to_string(), Listener::from_uds(lst)));
Ok(self)
}
/// Registers a service factory for an existing TCP listener.
///
/// The listener is switched to non-blocking mode. See
/// [`bind`](Self::bind) for `cfg` and `factory`.
pub fn listen<F, S, I>(
mut self,
name: impl AsRef<str>,
lst: net::TcpListener,
cfg: impl Into<SharedCfg>,
factory: F,
) -> io::Result<Self>
where
F: AsyncFn(&Cfg::State) -> I + Send + Clone + 'static,
S: Service<Cfg::State, Io> + 'static,
I: IntoService<S, Cfg::State, Io> + 'static,
{
let token = self.token.next();
self.services.push(factory::create_factory_service(
name.as_ref().to_string(),
vec![(token, cfg.into())],
factory,
));
self.sockets
.push((token, name.as_ref().to_string(), Listener::from_tcp(lst)));
Ok(self)
}
/// Starts processing incoming connections and returns a server controller.
///
/// # Panics
///
/// Panics if no listener has been registered.
pub fn run(self) -> Server<Connection> {
assert!(
!self.sockets.is_empty(),
"Server should have at least one bound socket"
);
let srv = StreamServer::new(self.accept.notify(), self.state, self.services);
let svc = self.pool.run(srv);
let sockets = self
.sockets
.into_iter()
.map(|sock| {
log::info!("Starting \"{}\" service on {}", sock.1, sock.2);
(sock.0, sock.2)
})
.collect();
self.accept.start(sockets, svc.clone());
svc
}
}
impl<Cfg> fmt::Debug for ServerBuilder<Cfg> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ServerBuilder")
.field("name", &self.name)
.field("token", &self.token)
.field("backlog", &self.backlog)
.field("sockets", &self.sockets)
.field("accept", &self.accept)
.field("worker-pool", &self.pool)
.finish()
}
}
/// Binds TCP listeners for every address resolved from `addr`.
///
/// Succeeds if at least one address binds and returns only the listeners that
/// bound. Otherwise, returns the last bind error.
pub fn bind_addr<S: net::ToSocketAddrs>(
addr: S,
backlog: i32,
) -> io::Result<Vec<net::TcpListener>> {
let mut err = None;
let mut succ = false;
let mut sockets = Vec::new();
for addr in addr.to_socket_addrs()? {
match create_tcp_listener(addr, backlog) {
Ok(lst) => {
succ = true;
sockets.push(lst);
}
Err(e) => err = Some(e),
}
}
if succ {
Ok(sockets)
} else if let Some(e) = err.take() {
Err(e)
} else {
Err(io::Error::new(
io::ErrorKind::InvalidInput,
"Cannot bind to address.",
))
}
}
/// Creates and binds a TCP listener with the specified listen backlog.
pub fn create_tcp_listener(addr: net::SocketAddr, backlog: i32) -> io::Result<net::TcpListener> {
let builder = match addr {
net::SocketAddr::V4(_) => Socket::new(Domain::IPV4, Type::STREAM, None)?,
net::SocketAddr::V6(_) => Socket::new(Domain::IPV6, Type::STREAM, None)?,
};
// On Windows, this allows rebinding sockets which are actively in use,
// which allows “socket hijacking”, so we explicitly don't set it here.
// https://docs.microsoft.com/en-us/windows/win32/winsock/using-so-reuseaddr-and-so-exclusiveaddruse
#[cfg(not(windows))]
builder.set_reuse_address(true)?;
builder.bind(&SockAddr::from(addr))?;
builder.listen(backlog)?;
Ok(net::TcpListener::from(builder))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_bind_addr() {
let addrs: Vec<net::SocketAddr> = Vec::new();
assert!(bind_addr(&addrs[..], 10).is_err());
}
#[ntex::test]
async fn test_debug() {
let builder = ServerBuilder::default();
assert!(format!("{builder:?}").contains("ServerBuilder"));
}
}