1
0
mirror of https://github.com/fafhrd91/actix-net synced 2024-11-24 17:23:00 +01:00
actix-net/actix-server/src/builder.rs

453 lines
14 KiB
Rust
Raw Normal View History

2018-08-19 19:47:04 +02:00
use std::time::Duration;
use std::{io, mem, net};
2018-08-19 19:47:04 +02:00
2018-12-10 06:51:35 +01:00
use actix_rt::{spawn, Arbiter, System};
use futures::future::{lazy, ok};
use futures::stream::futures_unordered;
use futures::sync::mpsc::{unbounded, UnboundedReceiver};
use futures::{Async, Future, Poll, Stream};
2018-12-06 23:04:42 +01:00
use log::{error, info};
2018-08-19 19:47:04 +02:00
use net2::TcpBuilder;
use num_cpus;
2018-12-10 06:51:35 +01:00
use tokio_timer::sleep;
2018-08-19 19:47:04 +02:00
2018-12-11 06:06:54 +01:00
use crate::accept::{AcceptLoop, AcceptNotify, Command};
use crate::server::{Server, ServerCommand};
2019-03-09 01:26:30 +01:00
use crate::service_config::{ConfiguredService, ServiceConfig};
use crate::services::{InternalServiceFactory, ServiceFactory, StreamNewService};
2018-12-11 06:06:54 +01:00
use crate::signals::{Signal, Signals};
use crate::worker::{self, Worker, WorkerAvailability, WorkerClient};
2019-03-05 00:41:16 +01:00
use crate::{ssl, Token};
2018-12-10 05:30:04 +01:00
2018-12-10 06:51:35 +01:00
/// Server builder
2018-12-10 05:30:04 +01:00
pub struct ServerBuilder {
2018-08-19 19:47:04 +02:00
threads: usize,
token: Token,
2019-03-11 20:01:55 +01:00
backlog: i32,
workers: Vec<(usize, WorkerClient)>,
2018-09-27 05:40:45 +02:00
services: Vec<Box<InternalServiceFactory>>,
2018-08-19 19:47:04 +02:00
sockets: Vec<(Token, net::TcpListener)>,
accept: AcceptLoop,
exit: bool,
2018-09-27 05:40:45 +02:00
shutdown_timeout: Duration,
2018-08-19 19:47:04 +02:00
no_signals: bool,
2018-12-10 06:51:35 +01:00
cmd: UnboundedReceiver<ServerCommand>,
server: Server,
}
impl Default for ServerBuilder {
fn default() -> Self {
Self::new()
}
2018-08-19 19:47:04 +02:00
}
2018-12-10 05:30:04 +01:00
impl ServerBuilder {
2018-12-10 06:51:35 +01:00
/// Create new Server builder instance
pub fn new() -> ServerBuilder {
let (tx, rx) = unbounded();
let server = Server::new(tx);
2018-12-10 05:30:04 +01:00
ServerBuilder {
2018-08-19 19:47:04 +02:00
threads: num_cpus::get(),
token: Token(0),
2018-08-19 19:47:04 +02:00
workers: Vec::new(),
services: Vec::new(),
sockets: Vec::new(),
2018-12-10 06:51:35 +01:00
accept: AcceptLoop::new(server.clone()),
2019-03-11 20:01:55 +01:00
backlog: 2048,
2018-08-19 19:47:04 +02:00
exit: false,
2018-09-27 05:40:45 +02:00
shutdown_timeout: Duration::from_secs(30),
2018-08-19 19:47:04 +02:00
no_signals: false,
2018-12-10 06:51:35 +01:00
cmd: rx,
server,
2018-08-19 19:47:04 +02:00
}
}
/// Set number of workers to start.
///
/// By default server uses number of available logical cpu as workers
2018-08-19 19:47:04 +02:00
/// count.
pub fn workers(mut self, num: usize) -> Self {
self.threads = num;
self
}
2019-03-11 20:01:55 +01:00
/// Set 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.
///
/// This method should be called before `bind()` method call.
pub fn backlog(mut self, num: i32) -> Self {
self.backlog = num;
self
}
2018-08-19 19:47:04 +02:00
/// Sets the maximum per-worker number of concurrent connections.
///
/// All socket listeners will stop accepting connections when this limit is
/// reached for each worker.
///
2018-09-07 20:35:25 +02:00
/// By default max connections is set to a 25k per worker.
pub fn maxconn(self, num: usize) -> Self {
2018-09-08 18:36:38 +02:00
worker::max_concurrent_connections(num);
2018-08-19 19:47:04 +02:00
self
}
2019-03-05 00:41:16 +01:00
/// Sets the maximum per-worker concurrent connection establish process.
///
/// All listeners will stop accepting connections when this limit is reached. It
/// can be used to limit the global SSL CPU usage.
2018-08-19 19:47:04 +02:00
///
2019-03-05 00:41:16 +01:00
/// By default max connections is set to a 256.
2019-03-05 01:16:39 +01:00
pub fn maxconnrate(self, num: usize) -> Self {
2019-03-05 00:41:16 +01:00
ssl::max_concurrent_ssl_connect(num);
self
}
/// Stop actix system.
2018-08-19 19:47:04 +02:00
pub fn system_exit(mut self) -> Self {
self.exit = true;
self
}
/// Disable signal handling
pub fn disable_signals(mut self) -> Self {
self.no_signals = true;
self
}
2018-09-27 05:40:45 +02:00
/// Timeout for graceful workers shutdown in seconds.
2018-08-19 19:47:04 +02:00
///
/// After receiving a stop signal, workers have this much time to finish
/// serving requests. Workers still alive after the timeout are force
/// dropped.
///
/// By default shutdown timeout sets to 30 seconds.
pub fn shutdown_timeout(mut self, sec: u16) -> Self {
2018-09-27 05:40:45 +02:00
self.shutdown_timeout = Duration::from_secs(u64::from(sec));
2018-08-19 19:47:04 +02:00
self
}
/// Execute external configuration as part of the server building
/// process.
2018-08-22 20:36:56 +02:00
///
/// This function is useful for moving parts of configuration to a
/// different module or even library.
2018-12-10 06:51:35 +01:00
pub fn configure<F>(mut self, f: F) -> io::Result<ServerBuilder>
2018-08-22 20:36:56 +02:00
where
F: Fn(&mut ServiceConfig) -> io::Result<()>,
2018-08-22 20:36:56 +02:00
{
2019-03-11 20:01:55 +01:00
let mut cfg = ServiceConfig::new(self.threads, self.backlog);
f(&mut cfg)?;
if let Some(apply) = cfg.apply {
let mut srv = ConfiguredService::new(apply);
for (name, lst) in cfg.services {
let token = self.token.next();
srv.stream(token, name, lst.local_addr()?);
self.sockets.push((token, lst));
}
self.services.push(Box::new(srv));
}
self.threads = cfg.threads;
Ok(self)
2018-08-22 20:36:56 +02:00
}
2018-12-10 06:51:35 +01:00
/// Add new service to the server.
2018-09-18 05:19:48 +02:00
pub fn bind<F, U, N: AsRef<str>>(mut self, name: N, addr: U, factory: F) -> io::Result<Self>
2018-08-19 19:47:04 +02:00
where
F: ServiceFactory,
2018-08-19 19:47:04 +02:00
U: net::ToSocketAddrs,
{
2019-03-11 20:01:55 +01:00
let sockets = bind_addr(addr, self.backlog)?;
2018-08-19 19:47:04 +02:00
for lst in sockets {
2019-03-09 16:27:56 +01:00
let token = self.token.next();
self.services.push(StreamNewService::create(
name.as_ref().to_string(),
token,
factory.clone(),
lst.local_addr()?,
));
2018-11-14 23:20:33 +01:00
self.sockets.push((token, lst));
2018-08-19 19:47:04 +02:00
}
Ok(self)
}
2018-12-10 06:51:35 +01:00
/// Add new service to the server.
2018-09-18 05:19:48 +02:00
pub fn listen<F, N: AsRef<str>>(
2018-10-30 04:29:47 +01:00
mut self,
name: N,
lst: net::TcpListener,
factory: F,
) -> io::Result<Self>
2018-08-19 19:47:04 +02:00
where
F: ServiceFactory,
2018-08-19 19:47:04 +02:00
{
let token = self.token.next();
self.services.push(StreamNewService::create(
name.as_ref().to_string(),
token,
factory,
2019-03-09 16:27:56 +01:00
lst.local_addr()?,
));
2018-09-27 05:40:45 +02:00
self.sockets.push((token, lst));
Ok(self)
2018-09-27 05:40:45 +02:00
}
2018-08-19 19:47:04 +02:00
/// Spawn new thread and start listening for incoming connections.
///
/// This method spawns new thread and starts new actix system. Other than
/// that it is similar to `start()` method. This method blocks.
///
/// This methods panics if no socket addresses get bound.
///
/// ```rust,ignore
/// use actix_web::*;
///
2019-03-09 01:26:30 +01:00
/// fn main() -> std::io::Result<()> {
2018-08-19 19:47:04 +02:00
/// Server::new().
/// .service(
2019-03-09 01:26:30 +01:00
/// HttpServer::new(|| App::new().service(web::service("/").to(|| HttpResponse::Ok())))
2018-08-19 19:47:04 +02:00
/// .bind("127.0.0.1:0")
2019-03-09 01:26:30 +01:00
/// .run()
2018-08-19 19:47:04 +02:00
/// }
/// ```
2019-03-09 01:26:30 +01:00
pub fn run(self) -> io::Result<()> {
2018-08-19 19:47:04 +02:00
let sys = System::new("http-server");
self.start();
2019-03-09 01:26:30 +01:00
sys.run()
2018-08-19 19:47:04 +02:00
}
2018-12-10 06:51:35 +01:00
/// Starts processing incoming connections and return server controller.
2018-12-10 05:30:04 +01:00
pub fn start(mut self) -> Server {
2018-08-19 19:47:04 +02:00
if self.sockets.is_empty() {
2018-12-10 06:51:35 +01:00
panic!("Server should have at least one bound socket");
2018-08-19 19:47:04 +02:00
} else {
2018-09-18 05:19:48 +02:00
info!("Starting {} workers", self.threads);
2018-08-19 19:47:04 +02:00
// start workers
let mut workers = Vec::new();
for idx in 0..self.threads {
let worker = self.start_worker(idx, self.accept.get_notify());
workers.push(worker.clone());
self.workers.push((idx, worker));
2018-08-19 19:47:04 +02:00
}
// start accept thread
for sock in &self.sockets {
2018-09-18 05:19:48 +02:00
info!("Starting server on {}", sock.1.local_addr().ok().unwrap());
2018-08-19 19:47:04 +02:00
}
2018-12-10 06:51:35 +01:00
self.accept
2018-08-19 19:47:04 +02:00
.start(mem::replace(&mut self.sockets, Vec::new()), workers);
2018-12-11 06:06:54 +01:00
// handle signals
if !self.no_signals {
Signals::start(self.server.clone());
}
2018-08-19 19:47:04 +02:00
// start http server actor
2018-12-10 06:51:35 +01:00
let server = self.server.clone();
spawn(self);
server
2018-08-19 19:47:04 +02:00
}
}
fn start_worker(&self, idx: usize, notify: AcceptNotify) -> WorkerClient {
2018-11-01 23:33:35 +01:00
let (tx1, rx1) = unbounded();
let (tx2, rx2) = unbounded();
2018-09-27 05:40:45 +02:00
let timeout = self.shutdown_timeout;
2018-09-07 22:06:51 +02:00
let avail = WorkerAvailability::new(notify);
2018-11-01 23:33:35 +01:00
let worker = WorkerClient::new(idx, tx1, tx2, avail.clone());
2018-09-27 05:40:45 +02:00
let services: Vec<Box<InternalServiceFactory>> =
2018-08-19 19:47:04 +02:00
self.services.iter().map(|v| v.clone_factory()).collect();
2018-12-10 05:30:04 +01:00
Arbiter::new().send(lazy(move || {
2018-12-06 23:04:42 +01:00
Worker::start(rx1, rx2, services, avail, timeout);
Ok::<_, ()>(())
}));
2018-08-19 19:47:04 +02:00
worker
2018-08-19 19:47:04 +02:00
}
2018-12-10 06:51:35 +01:00
2018-12-11 06:06:54 +01:00
fn handle_cmd(&mut self, item: ServerCommand) {
match item {
ServerCommand::Pause(tx) => {
self.accept.send(Command::Pause);
let _ = tx.send(());
}
ServerCommand::Resume(tx) => {
self.accept.send(Command::Resume);
let _ = tx.send(());
}
ServerCommand::Signal(sig) => {
// Signals support
// Handle `SIGINT`, `SIGTERM`, `SIGQUIT` signals and stop actix system
match sig {
Signal::Int => {
info!("SIGINT received, exiting");
self.exit = true;
self.handle_cmd(ServerCommand::Stop {
graceful: false,
completion: None,
})
2018-12-10 06:51:35 +01:00
}
2018-12-11 06:06:54 +01:00
Signal::Term => {
info!("SIGTERM received, stopping");
self.exit = true;
self.handle_cmd(ServerCommand::Stop {
graceful: true,
completion: None,
})
2018-12-10 06:51:35 +01:00
}
2018-12-11 06:06:54 +01:00
Signal::Quit => {
info!("SIGQUIT received, exiting");
self.exit = true;
self.handle_cmd(ServerCommand::Stop {
graceful: false,
completion: None,
})
}
_ => (),
}
}
ServerCommand::Stop {
graceful,
completion,
} => {
let exit = self.exit;
// stop accept thread
self.accept.send(Command::Stop);
// stop workers
if !self.workers.is_empty() {
spawn(
futures_unordered(
self.workers
.iter()
.map(move |worker| worker.1.stop(graceful)),
)
.collect()
.then(move |_| {
if let Some(tx) = completion {
let _ = tx.send(());
}
if exit {
2018-12-10 06:51:35 +01:00
spawn(sleep(Duration::from_millis(300)).then(|_| {
2018-08-19 19:47:04 +02:00
System::current().stop();
2018-12-10 06:51:35 +01:00
ok(())
}));
2018-08-19 19:47:04 +02:00
}
2018-12-11 06:06:54 +01:00
ok(())
}),
)
} else {
// we need to stop system if server was spawned
if self.exit {
spawn(sleep(Duration::from_millis(300)).then(|_| {
System::current().stop();
ok(())
}));
2018-08-19 19:47:04 +02:00
}
2018-12-11 06:06:54 +01:00
if let Some(tx) = completion {
let _ = tx.send(());
}
}
}
ServerCommand::WorkerDied(idx) => {
let mut found = false;
for i in 0..self.workers.len() {
if self.workers[i].0 == idx {
self.workers.swap_remove(i);
found = true;
break;
}
}
if found {
error!("Worker has died {:?}, restarting", idx);
let mut new_idx = self.workers.len();
'found: loop {
2018-08-19 19:47:04 +02:00
for i in 0..self.workers.len() {
2018-12-11 06:06:54 +01:00
if self.workers[i].0 == new_idx {
new_idx += 1;
continue 'found;
2018-08-19 19:47:04 +02:00
}
}
2018-12-11 06:06:54 +01:00
break;
}
2018-08-19 19:47:04 +02:00
2018-12-11 06:06:54 +01:00
let worker = self.start_worker(new_idx, self.accept.get_notify());
self.workers.push((new_idx, worker.clone()));
self.accept.send(Command::Worker(worker));
}
}
}
}
}
2018-12-10 06:51:35 +01:00
2018-12-11 06:06:54 +01:00
impl Future for ServerBuilder {
type Item = ();
type Error = ();
fn poll(&mut self) -> Poll<Self::Item, Self::Error> {
loop {
match self.cmd.poll() {
Ok(Async::Ready(None)) | Err(_) => return Ok(Async::Ready(())),
Ok(Async::NotReady) => return Ok(Async::NotReady),
Ok(Async::Ready(Some(item))) => self.handle_cmd(item),
2018-08-19 19:47:04 +02:00
}
}
}
}
2019-03-11 20:01:55 +01:00
pub(super) fn bind_addr<S: net::ToSocketAddrs>(
addr: S,
backlog: i32,
) -> io::Result<Vec<net::TcpListener>> {
2018-08-19 19:47:04 +02:00
let mut err = None;
let mut succ = false;
let mut sockets = Vec::new();
for addr in addr.to_socket_addrs()? {
2019-03-11 20:01:55 +01:00
match create_tcp_listener(addr, backlog) {
2018-08-19 19:47:04 +02:00
Ok(lst) => {
succ = true;
sockets.push(lst);
}
Err(e) => err = Some(e),
}
}
if !succ {
if let Some(e) = err.take() {
Err(e)
} else {
Err(io::Error::new(
io::ErrorKind::Other,
"Can not bind to address.",
))
}
} else {
Ok(sockets)
}
}
2019-03-11 20:01:55 +01:00
fn create_tcp_listener(addr: net::SocketAddr, backlog: i32) -> io::Result<net::TcpListener> {
2018-08-19 19:47:04 +02:00
let builder = match addr {
net::SocketAddr::V4(_) => TcpBuilder::new_v4()?,
net::SocketAddr::V6(_) => TcpBuilder::new_v6()?,
};
builder.reuse_address(true)?;
builder.bind(addr)?;
2019-03-11 20:01:55 +01:00
Ok(builder.listen(backlog)?)
2018-08-19 19:47:04 +02:00
}