Skip to main content

mesh_llm_console_server/
lib.rs

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}