Skip to main content

horto_os_ui_mcp/
server.rs

1//! MCP server (`rmcp`) for Horto (stdio or Streamable HTTP).
2
3use std::net::SocketAddr;
4use std::sync::Arc;
5
6use axum::extract::{ConnectInfo, Request, State};
7use axum::http::{header, StatusCode};
8use axum::middleware::{from_fn_with_state, Next};
9use axum::response::{IntoResponse, Response};
10use axum::Router;
11use rmcp::{
12    handler::server::wrapper::Parameters,
13    model::{CallToolResult, ContentBlock, ServerCapabilities, ServerConfig},
14    tool, tool_handler, tool_router, ErrorData as McpError, ServerHandler,
15};
16use subtle::ConstantTimeEq;
17use tokio::task::spawn_blocking;
18
19use crate::api_client::StatusApiClient;
20use crate::config::{McpMode, McpSettings};
21use crate::embedded_ops;
22use crate::net::is_local_network_ip;
23use crate::remote_ops;
24use crate::tool_args::{BackupEtcArgs, DockerRebuildArgs, SetupRunArgs, SetupStepArgs};
25
26/// Default HTTP listen when `HORTO_MCP_MODE=pc`.
27pub const DEFAULT_HTTP_LISTEN_PC: &str = "127.0.0.1:8790";
28/// Default HTTP listen when `HORTO_MCP_MODE=box`.
29pub const DEFAULT_HTTP_LISTEN_BOX: &str = "0.0.0.0:8790";
30
31/// MCP server handle.
32#[derive(Clone)]
33pub struct HortoMcp {
34    settings: Arc<McpSettings>,
35}
36
37impl HortoMcp {
38    /// Build from settings.
39    #[must_use]
40    pub fn new(settings: McpSettings) -> Self {
41        Self {
42            settings: Arc::new(settings),
43        }
44    }
45
46    fn settings(&self) -> Arc<McpSettings> {
47        Arc::clone(&self.settings)
48    }
49}
50
51/// Successful MCP tool result with a single text content block.
52pub(crate) fn text_ok(text: impl Into<String>) -> CallToolResult {
53    CallToolResult::success(vec![ContentBlock::text(text.into())])
54}
55
56/// MCP invalid-params error with a human message.
57pub(crate) fn mcp_err(msg: impl Into<String>) -> McpError {
58    McpError::invalid_params(msg.into(), None)
59}
60
61fn to_json_text<T: serde::Serialize>(value: &T) -> Result<String, McpError> {
62    serde_json::to_string_pretty(value).map_err(|err| mcp_err(err.to_string()))
63}
64
65async fn blocking_str<F>(f: F) -> Result<CallToolResult, McpError>
66where
67    F: FnOnce() -> anyhow::Result<String> + Send + 'static,
68{
69    spawn_blocking(f)
70        .await
71        .map_err(|err| mcp_err(format!("join error: {err}")))?
72        .map(text_ok)
73        .map_err(|err| mcp_err(err.to_string()))
74}
75
76#[tool_router]
77impl HortoMcp {
78    /// `GET /health` on the status API.
79    #[tool(description = "Horto status-api health (GET /health)")]
80    async fn health(&self) -> Result<CallToolResult, McpError> {
81        let client =
82            StatusApiClient::from_settings(&self.settings).map_err(|e| mcp_err(e.to_string()))?;
83        let v = client.health().await.map_err(|e| mcp_err(e.to_string()))?;
84        Ok(text_ok(to_json_text(&v)?))
85    }
86
87    /// `GET /v1/status` (bearer when token set).
88    #[tool(description = "Horto box status snapshot (GET /v1/status)")]
89    async fn get_status(&self) -> Result<CallToolResult, McpError> {
90        let client =
91            StatusApiClient::from_settings(&self.settings).map_err(|e| mcp_err(e.to_string()))?;
92        let v = client.status().await.map_err(|e| mcp_err(e.to_string()))?;
93        Ok(text_ok(to_json_text(&v)?))
94    }
95
96    /// Confirmed timestamped `/etc` backup via status-api.
97    #[tool(
98        description = "Timestamped /etc backup via status-api (confirm must be backup-etc; needs HORTO_API_TOKEN)"
99    )]
100    async fn backup_etc(
101        &self,
102        Parameters(args): Parameters<BackupEtcArgs>,
103    ) -> Result<CallToolResult, McpError> {
104        let client =
105            StatusApiClient::from_settings(&self.settings).map_err(|e| mcp_err(e.to_string()))?;
106        let v = client
107            .backup_etc(&args.confirm)
108            .await
109            .map_err(|e| mcp_err(e.to_string()))?;
110        Ok(text_ok(to_json_text(&v)?))
111    }
112
113    /// Probe remote box arch (PC mode only).
114    #[tool(description = "Probe box arch over SSH (HORTO_MCP_MODE=pc; needs HORTO_REMOTE_HOST)")]
115    async fn remote_probe(&self) -> Result<CallToolResult, McpError> {
116        if self.settings.mode != McpMode::Pc {
117            return Err(mcp_err(
118                "remote_probe is only available in HORTO_MCP_MODE=pc",
119            ));
120        }
121        let settings = self.settings();
122        blocking_str(move || remote_ops::remote_probe(&settings)).await
123    }
124
125    /// Setup status (remote SSH or embedded).
126    #[tool(description = "Setup pipeline status (remote SSH or embedded on box)")]
127    async fn setup_status(
128        &self,
129        Parameters(args): Parameters<SetupRunArgs>,
130    ) -> Result<CallToolResult, McpError> {
131        let settings = self.settings();
132        match settings.mode {
133            McpMode::Pc => {
134                blocking_str(move || remote_ops::setup_status(&settings, args.apply, args.full))
135                    .await
136            }
137            McpMode::Box => {
138                blocking_str(move || embedded_ops::setup_status_report(args.apply, args.full)).await
139            }
140        }
141    }
142
143    /// Full or minimal setup run.
144    #[tool(
145        description = "Run setup pipeline (apply default false = plan only). PC=SSH remote runner; box=embedded engine"
146    )]
147    async fn setup_run(
148        &self,
149        Parameters(args): Parameters<SetupRunArgs>,
150    ) -> Result<CallToolResult, McpError> {
151        let settings = self.settings();
152        match settings.mode {
153            McpMode::Pc => {
154                blocking_str(move || {
155                    remote_ops::setup_run(&settings, args.apply, args.full, args.skip_piper)
156                })
157                .await
158            }
159            McpMode::Box => {
160                blocking_str(move || {
161                    embedded_ops::setup_run_embedded(args.apply, args.full, args.skip_piper)
162                })
163                .await
164            }
165        }
166    }
167
168    /// Single setup step.
169    #[tool(description = "Run one setup step by id (s1, m1, d1, …)")]
170    async fn setup_step(
171        &self,
172        Parameters(args): Parameters<SetupStepArgs>,
173    ) -> Result<CallToolResult, McpError> {
174        let settings = self.settings();
175        let step_id = args.step_id.clone();
176        match settings.mode {
177            McpMode::Pc => {
178                blocking_str(move || {
179                    remote_ops::setup_step(&settings, &step_id, args.apply, args.full)
180                })
181                .await
182            }
183            McpMode::Box => {
184                blocking_str(move || {
185                    embedded_ops::setup_step_embedded(&step_id, args.apply, args.full)
186                })
187                .await
188            }
189        }
190    }
191
192    /// Doctor report.
193    #[tool(description = "Doctor / readiness report (remote or embedded)")]
194    async fn doctor(&self) -> Result<CallToolResult, McpError> {
195        let settings = self.settings();
196        match settings.mode {
197            McpMode::Pc => blocking_str(move || remote_ops::doctor(&settings)).await,
198            McpMode::Box => blocking_str(embedded_ops::doctor_embedded).await,
199        }
200    }
201
202    /// Docker container list / status.
203    #[tool(description = "Docker status (remote or embedded)")]
204    async fn docker_status(&self) -> Result<CallToolResult, McpError> {
205        let settings = self.settings();
206        match settings.mode {
207            McpMode::Pc => blocking_str(move || remote_ops::docker_status(&settings)).await,
208            McpMode::Box => blocking_str(embedded_ops::docker_status_embedded).await,
209        }
210    }
211
212    /// Docker compose rebuild (confirm `docker-rebuild`).
213    #[tool(description = "Docker compose rebuild (confirm must be docker-rebuild)")]
214    async fn docker_rebuild(
215        &self,
216        Parameters(args): Parameters<DockerRebuildArgs>,
217    ) -> Result<CallToolResult, McpError> {
218        let settings = self.settings();
219        let confirm = args.confirm.clone();
220        match settings.mode {
221            McpMode::Pc => {
222                blocking_str(move || remote_ops::docker_rebuild(&settings, &confirm)).await
223            }
224            McpMode::Box => {
225                blocking_str(move || embedded_ops::docker_rebuild_embedded(&confirm)).await
226            }
227        }
228    }
229
230    /// List timestamped `/etc` backups.
231    #[tool(description = "List timestamped /etc backups (remote or embedded)")]
232    async fn backup_list(&self) -> Result<CallToolResult, McpError> {
233        let settings = self.settings();
234        match settings.mode {
235            McpMode::Pc => blocking_str(move || remote_ops::backup_list(&settings)).await,
236            McpMode::Box => blocking_str(embedded_ops::backup_list_embedded).await,
237        }
238    }
239
240    /// Disk backup probe / status.
241    #[tool(description = "Disk backup readiness / status (remote or embedded)")]
242    async fn backup_disk_status(&self) -> Result<CallToolResult, McpError> {
243        let settings = self.settings();
244        match settings.mode {
245            McpMode::Pc => blocking_str(move || remote_ops::backup_disk_status(&settings)).await,
246            McpMode::Box => blocking_str(embedded_ops::backup_disk_status_embedded).await,
247        }
248    }
249}
250
251#[tool_handler]
252#[allow(clippy::unused_async_trait_impl)]
253impl ServerHandler for HortoMcp {
254    fn get_info(&self) -> ServerConfig {
255        ServerConfig::new(ServerCapabilities::builder().enable_tools().build())
256            .with_server_info(rmcp::model::Implementation::new(
257                "horto-os-ui",
258                env!("CARGO_PKG_VERSION"),
259            ))
260            .with_instructions(
261                "Horto tools: health, get_status, backup_etc (HTTP day-2); setup_*, doctor, docker_*, backup_list, backup_disk_status (SSH on pc / embedded on box). Set HORTO_MCP_MODE, HORTO_STATUS_API_URL, HORTO_API_TOKEN; PC needs HORTO_REMOTE_HOST. HTTP mode requires HORTO_MCP_TOKEN (or API token) and LAN peers only.",
262            )
263    }
264}
265
266#[derive(Clone)]
267struct HttpGate {
268    token: String,
269}
270
271fn bearer_ok(expected: &str, header: Option<&str>) -> bool {
272    let Some(provided) = header.and_then(|v| v.strip_prefix("Bearer ")) else {
273        return false;
274    };
275    if provided.len() != expected.len() {
276        return false;
277    }
278    bool::from(provided.as_bytes().ct_eq(expected.as_bytes()))
279}
280
281async fn gate_middleware(
282    State(gate): State<HttpGate>,
283    ConnectInfo(addr): ConnectInfo<SocketAddr>,
284    req: Request,
285    next: Next,
286) -> Response {
287    if !is_local_network_ip(addr.ip()) {
288        return (StatusCode::FORBIDDEN, "forbidden: non-LAN peer").into_response();
289    }
290    let auth = req
291        .headers()
292        .get(header::AUTHORIZATION)
293        .and_then(|v| v.to_str().ok());
294    if !bearer_ok(&gate.token, auth) {
295        return (StatusCode::UNAUTHORIZED, "unauthorized").into_response();
296    }
297    next.run(req).await
298}
299
300fn http_router(settings: McpSettings) -> Result<Router, String> {
301    let token = settings
302        .http_bearer()
303        .filter(|t| !t.is_empty())
304        .ok_or_else(|| "HTTP mode requires HORTO_MCP_TOKEN or HORTO_API_TOKEN".to_owned())?
305        .to_owned();
306    let mcp = HortoMcp::new(settings);
307    let config =
308        rmcp::transport::streamable_http_server::tower::StreamableHttpServerConfig::default();
309    let service = rmcp::transport::streamable_http_server::tower::StreamableHttpService::new(
310        move || Ok(mcp.clone()),
311        Arc::new(
312            rmcp::transport::streamable_http_server::session::local::LocalSessionManager::default(),
313        ),
314        config,
315    );
316    let method_router = axum::routing::any_service(service);
317    let gate = HttpGate { token };
318    Ok(Router::new()
319        .route("/mcp", method_router.clone())
320        .route("/mcp/", method_router)
321        .layer(from_fn_with_state(gate, gate_middleware)))
322}
323
324/// Serve Streamable HTTP until stopped.
325///
326/// # Errors
327///
328/// Returns bind/serve I/O errors, or missing token configuration.
329pub async fn run_http(addr: &str, settings: McpSettings) -> Result<(), String> {
330    let router = http_router(settings)?;
331    let listener = tokio::net::TcpListener::bind(addr)
332        .await
333        .map_err(|e| e.to_string())?;
334    let local = listener.local_addr().map_err(|e| e.to_string())?;
335    tracing::info!(%local, "horto-os-ui-mcp HTTP listening");
336    axum::serve(
337        listener,
338        router.into_make_service_with_connect_info::<SocketAddr>(),
339    )
340    .await
341    .map_err(|e| e.to_string())
342}
343
344/// Default listen string for a mode.
345#[must_use]
346pub const fn default_listen_for(mode: McpMode) -> &'static str {
347    match mode {
348        McpMode::Pc => DEFAULT_HTTP_LISTEN_PC,
349        McpMode::Box => DEFAULT_HTTP_LISTEN_BOX,
350    }
351}
352
353#[cfg(test)]
354mod tests {
355    use super::*;
356
357    #[test]
358    fn bearer_constant_time_compare() {
359        assert!(bearer_ok("secret", Some("Bearer secret")));
360        assert!(!bearer_ok("secret", Some("Bearer wrong")));
361        assert!(!bearer_ok("secret", Some("secret")));
362        assert!(!bearer_ok("secret", None));
363    }
364
365    #[tokio::test]
366    async fn http_requires_token_config() {
367        let settings = McpSettings {
368            mode: McpMode::Pc,
369            status_api_url: "http://127.0.0.1:8787".into(),
370            api_token: None,
371            mcp_token: None,
372            remote_host: None,
373            release_tag: None,
374            bin_dir: None,
375            install_ssh_key: false,
376        };
377        let err = http_router(settings).expect_err("token");
378        assert!(err.contains("HORTO_MCP_TOKEN") || err.contains("HORTO_API_TOKEN"));
379    }
380
381    #[tokio::test]
382    async fn http_rejects_missing_bearer() {
383        let settings = McpSettings {
384            mode: McpMode::Pc,
385            status_api_url: "http://127.0.0.1:8787".into(),
386            api_token: Some("tok".into()),
387            mcp_token: Some("tok".into()),
388            remote_host: None,
389            release_tag: None,
390            bin_dir: None,
391            install_ssh_key: false,
392        };
393        let router = http_router(settings).expect("router");
394        let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
395            .await
396            .expect("bind");
397        let addr = listener.local_addr().expect("addr");
398        let server = tokio::spawn(async move {
399            axum::serve(
400                listener,
401                router.into_make_service_with_connect_info::<SocketAddr>(),
402            )
403            .await
404        });
405
406        let res = reqwest::Client::new()
407            .post(format!("http://{addr}/mcp"))
408            .header("content-type", "application/json")
409            .header("accept", "application/json, text/event-stream")
410            .body("{}")
411            .send()
412            .await
413            .expect("post");
414        assert_eq!(res.status(), StatusCode::UNAUTHORIZED);
415
416        server.abort();
417        let _ = server.await;
418    }
419}