vertigo_cli/serve/
serve_run.rs1use 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(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 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 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}