1use mesh_llm_ui::{ConsoleAssetProvider, FileSystemConsoleAssets};
2use std::{net::SocketAddr, path::PathBuf, sync::Arc, time::Duration};
3use tokio::{
4 io::{AsyncReadExt, AsyncWriteExt},
5 net::{TcpListener, TcpStream},
6 sync::oneshot,
7 task::JoinHandle,
8};
9
10#[derive(Clone, Debug)]
11pub struct ConsoleServerOptions {
12 pub asset_dir: PathBuf,
13 pub port: u16,
14 pub listen_all: bool,
15}
16
17#[derive(Debug)]
18pub struct ConsoleServerHandle {
19 url: String,
20 shutdown_tx: Option<oneshot::Sender<()>>,
21 task: JoinHandle<()>,
22}
23
24impl ConsoleServerHandle {
25 pub fn url(&self) -> &str {
26 &self.url
27 }
28
29 pub async fn stop(mut self) {
30 if let Some(tx) = self.shutdown_tx.take() {
31 let _ = tx.send(());
32 }
33 let _ = self.task.await;
34 }
35}
36
37pub async fn start_file_console(
38 options: ConsoleServerOptions,
39) -> anyhow::Result<ConsoleServerHandle> {
40 let assets = Arc::new(FileSystemConsoleAssets::new(options.asset_dir));
41 if assets.index().is_none() {
42 anyhow::bail!("console asset directory must contain index.html");
43 }
44 start_console(options.port, options.listen_all, assets).await
45}
46
47pub async fn start_console(
48 port: u16,
49 listen_all: bool,
50 assets: Arc<dyn ConsoleAssetProvider>,
51) -> anyhow::Result<ConsoleServerHandle> {
52 let bind_addr = if listen_all { "0.0.0.0" } else { "127.0.0.1" };
53 let listener = TcpListener::bind(format!("{bind_addr}:{port}")).await?;
54 let addr = listener.local_addr()?;
55 let url = console_url(addr, listen_all);
56 let (shutdown_tx, shutdown_rx) = oneshot::channel();
57 let task = tokio::spawn(run(listener, assets, shutdown_rx));
58 Ok(ConsoleServerHandle {
59 url,
60 shutdown_tx: Some(shutdown_tx),
61 task,
62 })
63}
64
65async fn run(
66 listener: TcpListener,
67 assets: Arc<dyn ConsoleAssetProvider>,
68 mut shutdown_rx: oneshot::Receiver<()>,
69) {
70 loop {
71 tokio::select! {
72 result = listener.accept() => {
73 let Ok((stream, _)) = result else {
74 continue;
75 };
76 let assets = assets.clone();
77 tokio::spawn(async move {
78 let _ = handle_connection(stream, assets).await;
79 });
80 }
81 _ = &mut shutdown_rx => break,
82 }
83 }
84}
85
86async fn handle_connection(
87 mut stream: TcpStream,
88 assets: Arc<dyn ConsoleAssetProvider>,
89) -> anyhow::Result<()> {
90 let Some(request) = read_request(&mut stream).await? else {
91 return Ok(());
92 };
93 let Some((method, path)) = parse_request_line(&request) else {
94 respond_text(&mut stream, 400, "Bad Request", "bad request").await?;
95 return Ok(());
96 };
97 let path_only = path.split('?').next().unwrap_or(path);
98 if method != "GET" {
99 respond_text(&mut stream, 405, "Method Not Allowed", "method not allowed").await?;
100 return Ok(());
101 }
102
103 if is_index_route(path_only) {
104 respond_asset(&mut stream, assets.index(), 500, "console bundle missing").await?;
105 } else if is_static_asset_route(path_only) {
106 respond_asset(&mut stream, assets.asset(path_only), 404, "not found").await?;
107 } else {
108 respond_text(&mut stream, 404, "Not Found", "not found").await?;
109 }
110 Ok(())
111}
112
113fn is_index_route(path: &str) -> bool {
114 matches!(
115 path,
116 "/" | "/dashboard"
117 | "/dashboard/"
118 | "/reserves"
119 | "/reserves/"
120 | "/chat"
121 | "/chat/"
122 | "/configuration"
123 | "/configuration/"
124 | "/__playground"
125 | "/__meshviz-perf"
126 ) || path.starts_with("/chat/")
127 || path.starts_with("/configuration/")
128 || path.starts_with("/plugins/")
129}
130
131fn is_static_asset_route(path: &str) -> bool {
132 path.starts_with("/assets/")
133 || matches!(path.rsplit('.').next(), Some("png" | "ico" | "webmanifest"))
134 || (path.ends_with(".json") && !path.starts_with("/api/"))
135}
136
137async fn read_request(stream: &mut TcpStream) -> anyhow::Result<Option<Vec<u8>>> {
138 tokio::time::timeout(Duration::from_secs(5), read_request_headers(stream))
139 .await
140 .unwrap_or(Ok(None))
141}
142
143async fn read_request_headers(stream: &mut TcpStream) -> anyhow::Result<Option<Vec<u8>>> {
144 const MAX_REQUEST_HEADER_BYTES: usize = 16 * 1024;
145
146 let mut buffer = Vec::with_capacity(1024);
147 loop {
148 if request_headers_complete(&buffer) {
149 return Ok(Some(buffer));
150 }
151 if buffer.len() >= MAX_REQUEST_HEADER_BYTES {
152 return Ok(Some(buffer));
153 }
154
155 let remaining = MAX_REQUEST_HEADER_BYTES - buffer.len();
156 let mut chunk = [0_u8; 1024];
157 let chunk_len = remaining.min(chunk.len());
158 let read = stream.read(&mut chunk[..chunk_len]).await?;
159 if read == 0 {
160 return if buffer.is_empty() {
161 Ok(None)
162 } else {
163 Ok(Some(buffer))
164 };
165 }
166 buffer.extend_from_slice(&chunk[..read]);
167 }
168}
169
170fn request_headers_complete(request: &[u8]) -> bool {
171 request.windows(4).any(|window| window == b"\r\n\r\n")
172}
173
174fn parse_request_line(request: &[u8]) -> Option<(&str, &str)> {
175 let line_end = request.windows(2).position(|window| window == b"\r\n")?;
176 let line = std::str::from_utf8(&request[..line_end]).ok()?;
177 let mut parts = line.split_whitespace();
178 Some((parts.next()?, parts.next()?))
179}
180
181async fn respond_asset(
182 stream: &mut TcpStream,
183 asset: Option<mesh_llm_ui::UiAsset>,
184 missing_code: u16,
185 missing_message: &str,
186) -> anyhow::Result<()> {
187 let Some(asset) = asset else {
188 return respond_text(
189 stream,
190 missing_code,
191 status_text(missing_code),
192 missing_message,
193 )
194 .await;
195 };
196 let header = format!(
197 "HTTP/1.1 200 OK\r\nContent-Type: {}\r\nContent-Length: {}\r\nCache-Control: {}\r\nConnection: close\r\n\r\n",
198 asset.content_type,
199 asset.contents.len(),
200 asset.cache_control
201 );
202 stream.write_all(header.as_bytes()).await?;
203 stream.write_all(asset.contents.as_ref()).await?;
204 stream.shutdown().await?;
205 Ok(())
206}
207
208async fn respond_text(
209 stream: &mut TcpStream,
210 code: u16,
211 status: &str,
212 body: &str,
213) -> anyhow::Result<()> {
214 let header = format!(
215 "HTTP/1.1 {code} {status}\r\nContent-Type: text/plain; charset=utf-8\r\nContent-Length: {}\r\nCache-Control: no-cache\r\nConnection: close\r\n\r\n",
216 body.len()
217 );
218 stream.write_all(header.as_bytes()).await?;
219 stream.write_all(body.as_bytes()).await?;
220 stream.shutdown().await?;
221 Ok(())
222}
223
224fn status_text(code: u16) -> &'static str {
225 match code {
226 404 => "Not Found",
227 405 => "Method Not Allowed",
228 400 => "Bad Request",
229 500 => "Internal Server Error",
230 _ => "OK",
231 }
232}
233
234fn console_url(addr: SocketAddr, listen_all: bool) -> String {
235 if listen_all && addr.ip().is_unspecified() {
236 format!("http://127.0.0.1:{}", addr.port())
237 } else {
238 format!("http://{addr}")
239 }
240}
241
242#[cfg(test)]
243mod tests {
244 use super::{start_file_console, ConsoleServerOptions};
245 use std::{fs, io::Write};
246
247 #[tokio::test]
248 async fn serves_index_and_assets_from_directory() {
249 let root =
250 std::env::temp_dir().join(format!("mesh-llm-console-server-{}", std::process::id()));
251 let _ = fs::remove_dir_all(&root);
252 fs::create_dir_all(root.join("assets")).expect("create asset root");
253 fs::write(root.join("index.html"), "<html>console</html>").expect("write index");
254 fs::write(root.join("assets/app.js"), "console.log('ok')").expect("write app");
255
256 let handle = start_file_console(ConsoleServerOptions {
257 asset_dir: root.clone(),
258 port: 0,
259 listen_all: false,
260 })
261 .await
262 .expect("start console");
263
264 let index = blocking_get(handle.url().to_string(), "/".to_string()).await;
265 assert!(index.contains("200 OK"));
266 assert!(index.contains("<html>console</html>"));
267
268 let asset = blocking_get(handle.url().to_string(), "/assets/app.js".to_string()).await;
269 assert!(asset.contains("200 OK"));
270 assert!(asset.contains("text/javascript"));
271
272 handle.stop().await;
273 let _ = fs::remove_dir_all(root);
274 }
275
276 #[tokio::test]
277 async fn serves_index_for_console_deep_links() {
278 let root = std::env::temp_dir().join(format!(
279 "mesh-llm-console-server-deep-link-{}",
280 std::process::id()
281 ));
282 let _ = fs::remove_dir_all(&root);
283 fs::create_dir_all(root.join("assets")).expect("create asset root");
284 fs::write(root.join("index.html"), "<html>console</html>").expect("write index");
285
286 let handle = start_file_console(ConsoleServerOptions {
287 asset_dir: root.clone(),
288 port: 0,
289 listen_all: false,
290 })
291 .await
292 .expect("start console");
293
294 for path in [
295 "/configuration",
296 "/configuration/defaults",
297 "/configuration/local-deployment",
298 "/plugins/web-ui-exemplar/overview",
299 "/reserves",
300 "/chat/thread",
301 ] {
302 let response = blocking_get(handle.url().to_string(), path.to_string()).await;
303 assert!(
304 response.contains("200 OK"),
305 "expected {path} to serve index, got {response}"
306 );
307 assert!(response.contains("<html>console</html>"));
308 }
309
310 handle.stop().await;
311 let _ = fs::remove_dir_all(root);
312 }
313
314 #[tokio::test]
315 async fn handles_request_line_split_across_reads() {
316 let root = std::env::temp_dir().join(format!(
317 "mesh-llm-console-server-split-request-{}",
318 std::process::id()
319 ));
320 let _ = fs::remove_dir_all(&root);
321 fs::create_dir_all(root.join("assets")).expect("create asset root");
322 fs::write(root.join("index.html"), "<html>console</html>").expect("write index");
323
324 let handle = start_file_console(ConsoleServerOptions {
325 asset_dir: root.clone(),
326 port: 0,
327 listen_all: false,
328 })
329 .await
330 .expect("start console");
331
332 let response =
333 blocking_split_get(handle.url().to_string(), "/configuration".to_string()).await;
334 assert!(response.contains("200 OK"), "got {response}");
335 assert!(response.contains("<html>console</html>"));
336
337 handle.stop().await;
338 let _ = fs::remove_dir_all(root);
339 }
340
341 async fn blocking_get(base: String, path: String) -> String {
342 tokio::task::spawn_blocking(move || {
343 let url = base.strip_prefix("http://").expect("test server uses http");
344 let mut stream = std::net::TcpStream::connect(url).expect("connect");
345 write!(stream, "GET {path} HTTP/1.1\r\nHost: {url}\r\n\r\n").expect("write request");
346 let mut response = String::new();
347 std::io::Read::read_to_string(&mut stream, &mut response).expect("read response");
348 response
349 })
350 .await
351 .expect("blocking get")
352 }
353
354 async fn blocking_split_get(base: String, path: String) -> String {
355 tokio::task::spawn_blocking(move || {
356 let url = base.strip_prefix("http://").expect("test server uses http");
357 let mut stream = std::net::TcpStream::connect(url).expect("connect");
358 write!(stream, "GET {path}").expect("write partial request");
359 std::thread::sleep(std::time::Duration::from_millis(50));
360 write!(stream, " HTTP/1.1\r\nHost: {url}\r\n\r\n").expect("write request end");
361 let mut response = String::new();
362 std::io::Read::read_to_string(&mut stream, &mut response).expect("read response");
363 response
364 })
365 .await
366 .expect("blocking split get")
367 }
368}