Skip to main content

ntex/web/
server.rs

1use std::{io, marker::PhantomData, net, sync::Arc, sync::Mutex};
2
3#[cfg(feature = "openssl")]
4use tls_openssl::ssl::{AlpnError, SslAcceptor, SslAcceptorBuilder};
5#[cfg(feature = "rustls")]
6use tls_rustls::ServerConfig as RustlsServerConfig;
7
8use crate::error::IntoFailure;
9use crate::http::{self, Request, Response, ResponseError};
10use crate::server::{NoConfig, Server, ServerAppConfig, ServerBuilder};
11use crate::service::{IntoServiceFactory, Service, ServiceFactory, State, fn_service};
12use crate::{SharedCfg, time::Seconds};
13
14struct Config {
15    host: Option<String>,
16}
17
18/// An HTTP Server.
19///
20/// Create new http server with application factory.
21///
22/// ```rust,no_run
23/// use ntex::web::{self, App, HttpResponse, HttpServer};
24///
25/// #[ntex::main]
26/// async fn main() -> std::io::Result<()> {
27///     HttpServer::new(
28///         async |_| App::new()
29///             .service(web::resource("/").to(async || { HttpResponse::Ok() })))
30///         .bind("127.0.0.1:59090", ntex::SharedCfg::default())?
31///         .run()
32///         .await
33/// }
34/// ```
35#[derive(derive_more::Debug)]
36#[debug("HttpServer")]
37pub struct HttpServer<Cfg, F, I, Sf>
38where
39    Cfg: ServerAppConfig,
40    Cfg::State: Clone,
41    F: AsyncFn(&Cfg::State) -> I + Send + Clone + 'static,
42    I: IntoServiceFactory<Sf, Cfg::State, Request>,
43    Sf: ServiceFactory<Cfg::State, Request>,
44    Sf::Res: Into<Response>,
45    Sf::Error: ResponseError,
46    Sf::InitError: IntoFailure,
47{
48    factory: F,
49    config: Arc<Mutex<Config>>,
50    backlog: i32,
51    builder: ServerBuilder<Cfg>,
52    _t: PhantomData<Sf>,
53}
54
55impl<F, I, Sf> HttpServer<NoConfig, F, I, Sf>
56where
57    F: AsyncFn(&()) -> I + Send + Clone + 'static,
58    I: IntoServiceFactory<Sf, (), Request>,
59    Sf: ServiceFactory<(), Request> + 'static,
60    Sf::Res: Into<Response>,
61    Sf::Error: ResponseError,
62    Sf::InitError: IntoFailure,
63{
64    #[must_use]
65    /// Create new http server with application factory
66    pub fn new(factory: F) -> Self {
67        HttpServer {
68            factory,
69            config: Arc::new(Mutex::new(Config { host: None })),
70            backlog: 1024,
71            builder: ServerBuilder::default(),
72            _t: PhantomData,
73        }
74    }
75}
76
77impl<Cfg, F, I, Sf> HttpServer<Cfg, F, I, Sf>
78where
79    Cfg: ServerAppConfig,
80    Cfg::State: Clone,
81    F: AsyncFn(&Cfg::State) -> I + Send + Clone + 'static,
82    I: IntoServiceFactory<Sf, Cfg::State, Request>,
83    Sf: ServiceFactory<Cfg::State, Request> + 'static,
84    Sf::Res: Into<Response>,
85    Sf::Error: ResponseError,
86    Sf::InitError: IntoFailure,
87{
88    #[must_use]
89    /// Create new http server with application factory and state mapping
90    pub fn with_config(cfg: Cfg, factory: F) -> Self
91    where
92        Cfg: ServerAppConfig,
93    {
94        HttpServer {
95            factory,
96            config: Arc::new(Mutex::new(Config { host: None })),
97            backlog: 1024,
98            builder: ServerBuilder::new(cfg),
99            _t: PhantomData,
100        }
101    }
102
103    #[must_use]
104    /// Set number of workers to start.
105    ///
106    /// By default http server uses number of available logical cpu as threads
107    /// count.
108    pub fn workers(mut self, num: usize) -> Self {
109        self.builder = self.builder.workers(num);
110        self
111    }
112
113    #[must_use]
114    /// Set the maximum number of pending connections.
115    ///
116    /// This refers to the number of clients that can be waiting to be served.
117    /// Exceeding this number results in the client getting an error when
118    /// attempting to connect. It should only affect servers under significant
119    /// load.
120    ///
121    /// Generally set in the 64-2048 range. Default value is 1024.
122    ///
123    /// This method should be called before `bind()` method call.
124    pub fn backlog(mut self, backlog: i32) -> Self {
125        self.backlog = backlog;
126        self.builder = self.builder.backlog(backlog);
127        self
128    }
129
130    #[must_use]
131    /// Sets the maximum per-worker number of concurrent connections.
132    ///
133    /// All socket listeners will stop accepting connections when this limit is reached
134    /// for each worker.
135    ///
136    /// By default max connections is set to a 25k.
137    pub fn max_connections(mut self, num: usize) -> Self {
138        self.builder = self.builder.max_connections(num);
139        self
140    }
141
142    #[must_use]
143    /// Sets the maximum per-worker number of concurrent TLS handshakes.
144    ///
145    /// This limit applies only to TLS acceptors (OpenSSL and rustls). When it is
146    /// reached, TLS acceptors stop accepting new handshakes until one completes.
147    /// Plain TCP listeners are not affected. It can be used to limit the CPU
148    /// usage of TLS handshakes.
149    ///
150    /// By default the limit is 256.
151    pub fn max_tls_handshakes(self, num: usize) -> Self {
152        ntex_tls::max_concurrent_ssl_accept(num);
153        self
154    }
155
156    #[must_use]
157    /// Set server host name.
158    ///
159    /// Host name is used by application router as a hostname for url generation.
160    /// Check [`ConnectionInfo::host()`](crate::web::dev::ConnectionInfo::host)
161    /// documentation for more information.
162    ///
163    /// By default host name is set to a "localhost:8080" value.
164    pub fn server_hostname<T: AsRef<str>>(self, val: T) -> Self {
165        self.config.lock().unwrap().host = Some(val.as_ref().to_owned());
166        self
167    }
168
169    #[must_use]
170    /// Stop ntex runtime when server get dropped.
171    ///
172    /// By default "stop runtime" is disabled.
173    pub fn stop_runtime(mut self) -> Self {
174        self.builder = self.builder.stop_runtime();
175        self
176    }
177
178    #[must_use]
179    /// Stops the server when one of the workers panics.
180    ///
181    /// By default, "stop on panic" is disabled.
182    pub fn stop_on_panic(mut self) -> Self {
183        self.builder = self.builder.stop_on_panic();
184        self
185    }
186
187    #[must_use]
188    /// Disable signal handling.
189    ///
190    /// By default, signal handling is enabled.
191    pub fn disable_signals(mut self) -> Self {
192        self.builder = self.builder.disable_signals();
193        self
194    }
195
196    #[must_use]
197    /// Timeout for graceful worker shutdown.
198    ///
199    /// After receiving a stop signal, workers have this much time to finish
200    /// serving requests. Workers that are still alive after the timeout are
201    /// forcefully dropped.
202    ///
203    /// This bounds the worker as a whole, not an individual connection. Each
204    /// connection is bound separately by `IoConfig::set_shutdown_timeout`, so
205    /// this value should leave room for the connections a worker is still
206    /// draining to shut down themselves.
207    ///
208    /// By default, the timeout is set to 30 seconds.
209    pub fn graceful_shutdown_timeout(mut self, sec: Seconds) -> Self {
210        self.builder = self.builder.graceful_shutdown_timeout(sec);
211        self
212    }
213
214    #[must_use]
215    /// Enable cpu affinity.
216    ///
217    /// By default, affinity is disabled.
218    pub fn enable_affinity(mut self) -> Self {
219        self.builder = self.builder.enable_affinity();
220        self
221    }
222
223    #[must_use]
224    /// Graceful shutdown.
225    ///
226    /// When enabled, SIGQUIT, SIGSEGV, SIGABRT, and application panics stop the
227    /// server gracefully. SIGTERM always stops the server gracefully, and SIGINT
228    /// always stops it immediately.
229    ///
230    /// By default, these events stop the server immediately.
231    pub fn graceful_shutdown(mut self) -> Self {
232        self.builder = self.builder.graceful_shutdown();
233        self
234    }
235
236    /// Use listener for accepting incoming connection requests
237    ///
238    /// `HttpServer` does not change any configuration for `TcpListener`,
239    /// it needs to be configured before passing it to `listen()` method.
240    pub fn listen(mut self, lst: net::TcpListener, cfg: impl Into<SharedCfg>) -> io::Result<Self> {
241        let factory = self.factory.clone();
242        let addr = lst.local_addr().unwrap();
243
244        self.builder = self.builder.listen(
245            format!("ntex-web-service-{addr}"),
246            lst,
247            cfg.into(),
248            async move |st| {
249                let state = st.clone();
250                fn_service(async move |req| {
251                    Ok(State {
252                        req,
253                        state: state.clone(),
254                    })
255                })
256                .and_then(http::HttpService::new(factory(st).await))
257            },
258        )?;
259        Ok(self)
260    }
261
262    #[cfg(feature = "openssl")]
263    /// Use listener for accepting incoming tls connection requests.
264    ///
265    /// This method sets alpn protocols to "h2" and "http/1.1"
266    pub fn listen_openssl(
267        self,
268        lst: net::TcpListener,
269        cfg: impl Into<SharedCfg>,
270        builder: SslAcceptorBuilder,
271    ) -> io::Result<Self> {
272        self.listen_openssl_inner(lst, cfg.into(), openssl_acceptor(builder)?)
273    }
274
275    #[cfg(feature = "openssl")]
276    fn listen_openssl_inner(
277        mut self,
278        lst: net::TcpListener,
279        cfg: SharedCfg,
280        acceptor: SslAcceptor,
281    ) -> io::Result<Self> {
282        let factory = self.factory.clone();
283        let addr = lst.local_addr().unwrap();
284
285        self.builder = self.builder.listen(
286            format!("ntex-web-service-{addr}"),
287            lst,
288            cfg,
289            async move |st| {
290                let state = st.clone();
291
292                http::openssl(
293                    acceptor.clone(),
294                    fn_service(async move |req| {
295                        Ok(State {
296                            req,
297                            state: state.clone(),
298                        })
299                    })
300                    .and_then(http::HttpService::new(factory(st).await)),
301                )
302            },
303        )?;
304        Ok(self)
305    }
306
307    #[cfg(feature = "rustls")]
308    /// Use listener for accepting incoming tls connection requests.
309    ///
310    /// This method sets alpn protocols to "h2" and "http/1.1"
311    pub fn listen_rustls(
312        self,
313        lst: net::TcpListener,
314        cfg: impl Into<SharedCfg>,
315        config: RustlsServerConfig,
316    ) -> io::Result<Self> {
317        self.listen_rustls_inner(lst, cfg.into(), config)
318    }
319
320    #[cfg(feature = "rustls")]
321    fn listen_rustls_inner(
322        mut self,
323        lst: net::TcpListener,
324        cfg: SharedCfg,
325        config: RustlsServerConfig,
326    ) -> io::Result<Self> {
327        let factory = self.factory.clone();
328        let addr = lst.local_addr().unwrap();
329
330        self.builder = self.builder.listen(
331            format!("ntex-web-rustls-service-{addr}"),
332            lst,
333            cfg,
334            async move |st| {
335                let state = st.clone();
336                http::rustls(
337                    config.clone(),
338                    http::ALPN_PROTOS,
339                    fn_service(async move |req| {
340                        Ok(State {
341                            req,
342                            state: state.clone(),
343                        })
344                    })
345                    .and_then(http::HttpService::new(factory(st).await)),
346                )
347            },
348        )?;
349        Ok(self)
350    }
351
352    /// The socket address to bind.
353    ///
354    /// To bind multiple addresses this method can be called multiple times.
355    pub fn bind<A: net::ToSocketAddrs>(
356        mut self,
357        addr: A,
358        cfg: impl Into<SharedCfg>,
359    ) -> io::Result<Self> {
360        let cfg = cfg.into();
361        for lst in self.bind2(addr)? {
362            self = self.listen(lst, cfg.clone())?;
363        }
364
365        Ok(self)
366    }
367
368    fn bind2<A: net::ToSocketAddrs>(&self, addr: A) -> io::Result<Vec<net::TcpListener>> {
369        let mut err = None;
370        let mut succ = false;
371        let mut sockets = Vec::new();
372        for addr in addr.to_socket_addrs()? {
373            match crate::server::create_tcp_listener(addr, self.backlog) {
374                Ok(lst) => {
375                    succ = true;
376                    sockets.push(lst);
377                }
378                Err(e) => err = Some(e),
379            }
380        }
381
382        if succ {
383            Ok(sockets)
384        } else if let Some(e) = err.take() {
385            Err(e)
386        } else {
387            Err(io::Error::new(
388                io::ErrorKind::InvalidInput,
389                "Cannot bind to address.",
390            ))
391        }
392    }
393
394    #[cfg(feature = "openssl")]
395    /// Start listening for incoming tls connections.
396    ///
397    /// This method sets alpn protocols to "h2" and "http/1.1"
398    pub fn bind_openssl<A>(
399        mut self,
400        addr: A,
401        builder: SslAcceptorBuilder,
402        cfg: impl Into<SharedCfg>,
403    ) -> io::Result<Self>
404    where
405        A: net::ToSocketAddrs,
406    {
407        let cfg = cfg.into();
408        let sockets = self.bind2(addr)?;
409        let acceptor = openssl_acceptor(builder)?;
410
411        for lst in sockets {
412            self = self.listen_openssl_inner(lst, cfg.clone(), acceptor.clone())?;
413        }
414
415        Ok(self)
416    }
417
418    #[cfg(feature = "rustls")]
419    /// Start listening for incoming tls connections.
420    ///
421    /// This method sets alpn protocols to "h2" and "http/1.1"
422    pub fn bind_rustls<A: net::ToSocketAddrs>(
423        mut self,
424        addr: A,
425        config: &RustlsServerConfig,
426        cfg: impl Into<SharedCfg>,
427    ) -> io::Result<Self> {
428        let cfg = cfg.into();
429        let sockets = self.bind2(addr)?;
430        for lst in sockets {
431            self = self.listen_rustls_inner(lst, cfg.clone(), config.clone())?;
432        }
433        Ok(self)
434    }
435
436    #[cfg(unix)]
437    /// Start listening for unix domain connections on existing listener.
438    ///
439    /// This method is available only on Unix.
440    pub fn listen_uds(
441        mut self,
442        lst: std::os::unix::net::UnixListener,
443        cfg: impl Into<SharedCfg>,
444    ) -> io::Result<Self> {
445        let factory = self.factory.clone();
446        let addr = format!("ntex-web-service-{:?}", lst.local_addr()?);
447
448        self.builder = self
449            .builder
450            .listen_uds(addr, lst, cfg.into(), async move |st| {
451                let state = st.clone();
452                fn_service(async move |req| {
453                    Ok(State {
454                        req,
455                        state: state.clone(),
456                    })
457                })
458                .and_then(http::HttpService::new(factory(st).await))
459            })?;
460        Ok(self)
461    }
462
463    #[cfg(unix)]
464    /// Start listening for incoming unix domain connections.
465    ///
466    /// This method is available only on Unix.
467    pub fn bind_uds<A>(mut self, addr: A, cfg: impl Into<SharedCfg>) -> io::Result<Self>
468    where
469        A: AsRef<std::path::Path>,
470    {
471        let factory = self.factory.clone();
472
473        self.builder = self.builder.bind_uds(
474            format!("ntex-web-service-{:?}", addr.as_ref().display()),
475            addr,
476            cfg.into(),
477            async move |st| {
478                let state = st.clone();
479                fn_service(async move |req| {
480                    Ok(State {
481                        req,
482                        state: state.clone(),
483                    })
484                })
485                .and_then(http::HttpService::new(factory(st).await))
486            },
487        )?;
488        Ok(self)
489    }
490}
491
492impl<Cfg, F, I, Sf> HttpServer<Cfg, F, I, Sf>
493where
494    Cfg: ServerAppConfig,
495    Cfg::State: Clone,
496    F: AsyncFn(&Cfg::State) -> I + Send + Clone + 'static,
497    I: IntoServiceFactory<Sf, Cfg::State, Request>,
498    Sf: ServiceFactory<Cfg::State, Request> + 'static,
499    Sf::Res: Into<Response>,
500    Sf::Error: ResponseError,
501    Sf::InitError: IntoFailure,
502{
503    /// Start listening for incoming connections.
504    ///
505    /// This method starts number of http workers in separate threads.
506    /// For each address this method starts separate thread which does
507    /// `accept()` in a loop.
508    ///
509    /// # Panics
510    ///
511    /// Panics if no listener was registered before calling this method.
512    ///
513    /// ```rust,no_run
514    /// use ntex::web::{self, App, HttpResponse, HttpServer};
515    ///
516    /// #[ntex::main]
517    /// async fn main() -> std::io::Result<()> {
518    ///     HttpServer::new(
519    ///         async |_| App::new().service(web::resource("/").to(async || { HttpResponse::Ok() }))
520    ///     )
521    ///         .bind("127.0.0.1:0", ntex::SharedCfg::default())?
522    ///         .run()
523    ///         .await
524    /// }
525    /// ```
526    pub fn run(self) -> Server {
527        self.builder.run()
528    }
529}
530
531#[cfg(feature = "openssl")]
532/// Configure `SslAcceptorBuilder` with custom server flags.
533fn openssl_acceptor(mut builder: SslAcceptorBuilder) -> io::Result<SslAcceptor> {
534    builder.set_alpn_select_callback(|_, protos| {
535        const H2: &[u8] = b"\x02h2";
536        const H11: &[u8] = b"\x08http/1.1";
537        if protos.windows(3).any(|window| window == H2) {
538            Ok(b"h2")
539        } else if protos.windows(9).any(|window| window == H11) {
540            Ok(b"http/1.1")
541        } else {
542            Err(AlpnError::NOACK)
543        }
544    });
545    builder.set_alpn_protos(b"\x08http/1.1\x02h2")?;
546
547    Ok(builder.build())
548}