Skip to main content

ntex_server/
server.rs

1#![allow(clippy::missing_panics_doc)]
2use std::sync::{Arc, atomic::AtomicBool, atomic::Ordering};
3use std::task::{Context, Poll, ready};
4use std::{future::Future, io, pin::Pin};
5
6use async_channel::Sender;
7use ntex_rt::signals::Signal;
8
9use crate::manager::ServerCommand;
10
11#[derive(Debug)]
12pub(crate) struct ServerShared {
13    pub(crate) paused: AtomicBool,
14}
15
16/// Controller and completion future for a running server.
17///
18/// Clones can pause, resume, stop, or submit items to the server. Awaiting a
19/// `Server` resolves when the server has stopped. It always resolves to
20/// `Ok(())`.
21#[derive(Debug)]
22pub struct Server<T> {
23    shared: Arc<ServerShared>,
24    cmd: Sender<ServerCommand<T>>,
25    stop: Option<oneshot::AsyncReceiver<()>>,
26}
27
28impl Server<crate::net::Connection> {
29    /// Creates a network server builder with no application configuration.
30    pub fn builder() -> crate::net::ServerBuilder {
31        crate::net::ServerBuilder::default()
32    }
33}
34
35impl<T> Server<T> {
36    pub(crate) fn new(cmd: Sender<ServerCommand<T>>, shared: Arc<ServerShared>) -> Self {
37        Server {
38            cmd,
39            shared,
40            stop: None,
41        }
42    }
43
44    pub(crate) fn signal(&self, sig: Signal) {
45        let _ = self.cmd.try_send(ServerCommand::Signal(sig));
46    }
47
48    /// Submits an item to the worker pool.
49    ///
50    /// Returns the item unchanged if the server is paused or cannot accept it.
51    ///
52    /// The server starts paused and resumes once the first worker is ready.
53    /// It also pauses itself while no worker is available, so items submitted
54    /// right after start or during worker restarts may be rejected.
55    pub fn process(&self, item: T) -> Result<(), T> {
56        if self.shared.paused.load(Ordering::Acquire) {
57            Err(item)
58        } else if let Err(e) = self.cmd.try_send(ServerCommand::Item(item)) {
59            if let ServerCommand::Item(item) = e.into_inner() {
60                Err(item)
61            } else {
62                panic!()
63            }
64        } else {
65            Ok(())
66        }
67    }
68
69    /// Pauses processing new items.
70    ///
71    /// For network servers, listeners stop accepting connections, so new
72    /// connections wait in the kernel listen backlog. Existing connections
73    /// remain active.
74    pub fn pause(&self) -> impl Future<Output = ()> + use<T> {
75        let (tx, rx) = oneshot::channel();
76        let _ = self.cmd.try_send(ServerCommand::Pause(tx));
77        async move {
78            let _ = rx.await;
79        }
80    }
81
82    /// Resumes processing new items.
83    pub fn resume(&self) -> impl Future<Output = ()> + use<T> {
84        let (tx, rx) = oneshot::channel();
85        let _ = self.cmd.try_send(ServerCommand::Resume(tx));
86        async move {
87            let _ = rx.await;
88        }
89    }
90
91    /// Stops processing new items and shuts down all workers.
92    ///
93    /// If `graceful` is `true`, workers are given up to the graceful shutdown
94    /// timeout to finish active work before they are stopped. A zero timeout
95    /// makes the stop non-graceful.
96    ///
97    /// The returned future resolves once the stop has completed. If the
98    /// server has already stopped, it resolves immediately.
99    pub fn stop(&self, graceful: bool) -> impl Future<Output = ()> + use<T> {
100        let (tx, rx) = oneshot::channel();
101        let _ = self.cmd.try_send(ServerCommand::Stop {
102            graceful,
103            completion: Some(tx),
104        });
105        async move {
106            let _ = rx.await;
107        }
108    }
109}
110
111impl<T> Clone for Server<T> {
112    fn clone(&self) -> Self {
113        Self {
114            cmd: self.cmd.clone(),
115            shared: self.shared.clone(),
116            stop: None,
117        }
118    }
119}
120
121impl<T> Future for Server<T> {
122    type Output = io::Result<()>;
123
124    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
125        let this = self.get_mut();
126
127        if this.stop.is_none() {
128            let (tx, rx) = oneshot::async_channel();
129            if this.cmd.try_send(ServerCommand::NotifyStopped(tx)).is_err() {
130                return Poll::Ready(Ok(()));
131            }
132            this.stop = Some(rx);
133        }
134
135        let _ = ready!(Pin::new(this.stop.as_mut().unwrap()).poll(cx));
136
137        Poll::Ready(Ok(()))
138    }
139}