Skip to main content

vertigo_cli/serve/
serve_run.rs

1use actix_proxy::IntoHttpResponse;
2use actix_web::{
3    App, HttpRequest, HttpResponse, HttpServer,
4    dev::{ServiceFactory, ServiceRequest},
5    http::header,
6    middleware::{Compress, Condition},
7    rt::System,
8    web,
9};
10use std::{num::NonZeroUsize, time::Duration};
11
12use crate::commons::{
13    ErrorCode,
14    spawn::{ServerOwner, term_signal},
15};
16use crate::serve::mount_path::MountConfig;
17
18use super::{
19    ServeOpts, ServeOptsInner, server_state::ServerState, vertigo_install::vertigo_install,
20};
21
22async fn wait_for_port(addr: &str, port: u16) {
23    for i in 0..20 {
24        match std::net::TcpListener::bind(addr) {
25            Ok(listener) => {
26                // Drop the listener to free the port for actix
27                drop(listener);
28                break;
29            }
30            Err(_) => {
31                log::warn!(
32                    "Port {} is still in use, waiting 1s... ({}/20)",
33                    port,
34                    i + 1
35                );
36                tokio::time::sleep(Duration::from_secs(1)).await;
37            }
38        }
39    }
40}
41
42pub async fn run(opts: ServeOpts, port_watch: Option<u16>) -> Result<(), ErrorCode> {
43    log::info!("serve params => {opts:#?}");
44
45    let ServeOptsInner {
46        host,
47        port,
48        proxy,
49        mount_point,
50        env,
51        wasm_preload,
52        disable_hydration,
53        disable_compression,
54        threads,
55    } = opts.inner;
56
57    let mount_config = MountConfig::new(
58        mount_point,
59        opts.common.dest_dir,
60        env,
61        wasm_preload,
62        disable_hydration,
63    )?;
64
65    ServerState::init_with_watch(&mount_config, port_watch)?;
66
67    let app = move || {
68        // Already-encoded responses (what `install_proxy` forwards) carry a `Content-Encoding`
69        // and `Compress` leaves those alone.
70        let mut app = App::new().wrap(Condition::new(!disable_compression, Compress::default()));
71
72        for (path, target) in &proxy {
73            app = install_proxy(app, path.clone(), target.clone());
74        }
75
76        app.configure(|cfg| {
77            vertigo_install(cfg, &mount_config);
78        })
79    };
80
81    let addr = format!("{host}:{port}");
82
83    wait_for_port(&addr, port).await;
84
85    let server =
86        HttpServer::new(app)
87            .workers(threads.unwrap_or_else(|| {
88                std::thread::available_parallelism().map_or(2, NonZeroUsize::get)
89            }))
90            .bind(addr.clone())
91            .map_err(|err| {
92                log::error!("Can't bind/serve on {addr}: {err}");
93                ErrorCode::ServeCantOpenPort
94            })?;
95
96    let server = server
97        .disable_signals()
98        .client_disconnect_timeout(Duration::from_secs(2))
99        .shutdown_timeout(5)
100        .run();
101
102    let handle = server.handle();
103    let handle2 = server.handle();
104
105    std::thread::spawn(move || System::new().block_on(server));
106
107    tokio::select! {
108        _ = ServerOwner { handle } => {},
109        msg = term_signal() => {
110            log::info!("{msg} received, shutting down");
111            handle2.stop(false).await;
112        }
113    }
114
115    Ok(())
116}
117
118fn install_proxy<T>(app: App<T>, path: String, target: String) -> App<T>
119where
120    T: ServiceFactory<ServiceRequest, Config = (), Error = actix_web::Error, InitError = ()>,
121{
122    app.service(web::scope(&path).default_service(web::to({
123        move |req: HttpRequest, body: web::Bytes| {
124            let path = path.clone();
125            let target = target.clone();
126            async move {
127                let method = req.method();
128                let uri = req.uri();
129                let current_path = uri.path();
130                let tail = current_path.strip_prefix(&path).unwrap_or(current_path);
131                let query = uri.query().map(|q| format!("?{}", q)).unwrap_or_default();
132
133                let target_url = format!("{target}{tail}{query}");
134
135                log::info!("proxy {method} {path}{tail} -> {target_url}");
136
137                let request = awc::Client::new().request_from(&target_url, req.head());
138
139                let response = if !body.is_empty() {
140                    request.send_body(body)
141                } else {
142                    request.send()
143                };
144
145                match response.await {
146                    Ok(response) => {
147                        let mut response = response.into_http_response();
148
149                        let headers = response.headers_mut();
150
151                        // `awc` decompresses the upstream body transparently, but
152                        // `into_http_response` copies the upstream headers across verbatim
153                        // so remove appropriate headers.
154                        if headers.remove(header::CONTENT_ENCODING).next().is_some() {
155                            headers.remove(header::CONTENT_LENGTH);
156                        }
157
158                        response
159                    }
160                    Err(error) => {
161                        let message = format!("Error fetching from url={target_url} error={error}");
162                        HttpResponse::InternalServerError().body(message)
163                    }
164                }
165            }
166        }
167    })))
168}