1use 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
26pub const DEFAULT_HTTP_LISTEN_PC: &str = "127.0.0.1:8790";
28pub const DEFAULT_HTTP_LISTEN_BOX: &str = "0.0.0.0:8790";
30
31#[derive(Clone)]
33pub struct HortoMcp {
34 settings: Arc<McpSettings>,
35}
36
37impl HortoMcp {
38 #[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
51pub(crate) fn text_ok(text: impl Into<String>) -> CallToolResult {
53 CallToolResult::success(vec![ContentBlock::text(text.into())])
54}
55
56pub(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 #[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 #[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 #[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 #[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 #[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 #[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 #[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 #[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 #[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 #[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 #[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 #[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
324pub 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#[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}