fusen_rs/protocol/http/
server.rs1use crate::{common::utils::shutdown::Shutdown, server::router::Router};
2use hyper_util::rt::TokioExecutor;
3use hyper_util::{rt::TokioIo, server::conn::auto::Builder};
4use std::sync::Arc;
5use tokio::net::TcpStream;
6use tokio::{
7 io,
8 net::TcpListener,
9 sync::{broadcast, mpsc},
10};
11
12#[derive(Clone)]
13pub struct TcpServer;
14
15impl TcpServer {
16 pub async fn run(port: u16, router: Router, mut shutdown: Shutdown) -> io::Result<()> {
17 let mut builder = hyper_util::server::conn::auto::Builder::new(TokioExecutor::new());
18 builder.http2().max_concurrent_streams(None);
19 builder.http1().keep_alive(true);
20 let builder = Arc::new(builder);
21 let tcp_listener = TcpListener::bind(format!("0.0.0.0:{port}")).await?;
22 let notify_shutdown: tokio::sync::broadcast::Sender<()> = broadcast::channel(1).0;
23 let (sender, mut recv) = mpsc::channel::<()>(1);
24 let router = Arc::new(router);
25 loop {
26 let (tcp_stream, _socketaddr) = tokio::select! {
27 stream = tcp_listener.accept() => stream?,
28 _ = shutdown.recv() => {
29 drop(notify_shutdown);
30 drop(sender);
31 let _ = recv.recv().await;
32 return Ok(());
33 }
34 };
35 let shutdown = Shutdown::new(notify_shutdown.subscribe());
36 let router = router.clone();
37 let sender = sender.clone();
38 let builder = builder.clone();
39 tokio::spawn(async move {
40 let handler = HttpStreamHandler {
41 router,
42 builder,
43 stream: tcp_stream,
44 shutdown,
45 _sender: sender,
46 };
47 handler.run().await;
48 });
49 }
50 }
51}
52
53pub struct HttpStreamHandler {
54 pub router: Arc<Router>,
55 pub builder: Arc<Builder<TokioExecutor>>,
56 pub stream: TcpStream,
57 pub shutdown: Shutdown,
58 pub _sender: mpsc::Sender<()>,
59}
60
61impl HttpStreamHandler {
62 pub async fn run(self) {
63 let HttpStreamHandler {
64 router,
65 builder,
66 stream,
67 mut shutdown,
68 _sender,
69 } = self;
70 let hyper_io = TokioIo::new(stream);
71 let server = builder.serve_connection_with_upgrades(hyper_io, router);
72 let _result = tokio::select! {
73 result = server => result,
74 _ = shutdown.recv() => Ok(())
75 };
76 }
77}