Skip to main content

app_forge_kit_http_server/
provider.rs

1use crate::config::Config;
2use crate::routing::Router;
3use app_forge_kit_service::{Error, Observable, Signal};
4use app_forge_kit_telemetry_tracing::types::Level;
5use async_trait::async_trait;
6use std::net::SocketAddr;
7use std::sync::Arc;
8use tokio::net::TcpListener;
9use tokio::sync::Mutex;
10use tower_http::trace::{DefaultOnRequest, DefaultOnResponse, TraceLayer};
11
12pub struct Provider {
13    done: Arc<Mutex<Option<tokio::sync::watch::Sender<Option<()>>>>>,
14    config: Option<Config>,
15    router: Router,
16}
17
18#[allow(clippy::new_without_default)]
19impl Provider {
20    pub fn new() -> Self {
21        Provider {
22            done: Arc::new(Mutex::new(None)),
23            config: None,
24            router: Router::new(),
25        }
26    }
27
28    pub fn with_config(self, config: Config) -> Self {
29        Self {
30            config: Some(config),
31            ..self
32        }
33    }
34
35    pub fn with_router(self, router: Router) -> Self {
36        Self { router, ..self }
37    }
38}
39
40const DEFAULT_LISTEN_ADDR: &str = "127.0.0.1:8080";
41
42#[async_trait]
43impl Observable for Provider {
44    async fn serve(&self) -> Result<(), Error> {
45        let config = self.config.clone().unwrap_or_default();
46
47        let tcp_listener = TcpListener::bind(
48            config.listen.unwrap_or(
49                DEFAULT_LISTEN_ADDR
50                    .parse::<SocketAddr>()
51                    .map_err(|err| Error::from(crate::Error::from(err)))?,
52            ),
53        )
54        .await?;
55
56        let (done_tx, mut done_rx) = tokio::sync::watch::channel(None);
57        self.done.lock().await.replace(done_tx);
58
59        let router = self.router.clone().layer(
60            TraceLayer::new_for_http()
61                .on_request(DefaultOnRequest::new().level(Level::INFO))
62                .on_response(DefaultOnResponse::new().level(Level::INFO)),
63        );
64
65        axum::serve(tcp_listener, router)
66            .with_graceful_shutdown(async move {
67                let _ = done_rx.changed().await;
68            })
69            .await?;
70
71        Ok(())
72    }
73
74    async fn signal(&self, signal: &Signal) -> Result<(), Error> {
75        if signal.is_terminate()
76            && let Some(done) = self.done.lock().await.take()
77        {
78            let _ = done.send(None);
79        }
80
81        Ok(())
82    }
83}