From d5727d322d775bcca8e36d50acdbbc9bfa9836ce Mon Sep 17 00:00:00 2001 From: Kayshen-X Date: Fri, 19 Jun 2026 21:27:35 +0800 Subject: [PATCH] refactor(host): move mcp_live to op-web-daemon (Phase 5, Task 5.2) The live MCP HTTP server (McpLiveServer + mcp_live/screenshot + tests) moves to op_web_daemon::mcp_live; its op_web_daemon::{mcp_serve,export} refs flip to crate::. start_with_wake stays generic over F: Fn() (the winit wake closure is built at the mcp_runtime call site, not inside mcp_live). McpLiveServer methods + McpPumpOutcome fields promoted to pub for the desktop consumers; the test-convenience start(port) made a non-cfg(test) doc-hidden seam. main.rs field + mcp_runtime + main_tests + web_canvas_server(screenshot) repointed. op-web-daemon 253 + live-MCP 8 green; no dep/lock change. --- crates/op-host-desktop/src/main.rs | 3 +- crates/op-host-desktop/src/main_tests.rs | 16 ++-- crates/op-host-desktop/src/mcp_runtime.rs | 6 +- .../op-host-desktop/src/web_canvas_server.rs | 4 +- crates/op-web-daemon/src/lib.rs | 1 + .../src/mcp_live.rs | 96 +++++++++---------- .../src/mcp_live/screenshot.rs | 6 +- 7 files changed, 66 insertions(+), 66 deletions(-) rename crates/{op-host-desktop => op-web-daemon}/src/mcp_live.rs (90%) rename crates/{op-host-desktop => op-web-daemon}/src/mcp_live/screenshot.rs (97%) diff --git a/crates/op-host-desktop/src/main.rs b/crates/op-host-desktop/src/main.rs index 32e0ecdf3..bed2d9b14 100644 --- a/crates/op-host-desktop/src/main.rs +++ b/crates/op-host-desktop/src/main.rs @@ -35,7 +35,6 @@ mod kit_io; mod kit_persistence; mod macos_app; mod mcp_integrations; -mod mcp_live; mod mcp_port_file; mod mcp_runtime; mod mcp_serve; @@ -241,7 +240,7 @@ struct DesktopApp { /// external repository changes. last_git_refresh: Instant, /// Live in-process MCP HTTP server, started from Settings -> MCP. - mcp_server: Option, + mcp_server: Option, /// When set (via the `--live-mcp[=port]` launch flag used by /// `op start`), the editor force-enables the live MCP server on /// this port during `resumed()`, regardless of the persisted diff --git a/crates/op-host-desktop/src/main_tests.rs b/crates/op-host-desktop/src/main_tests.rs index bc173c9ca..63edce895 100644 --- a/crates/op-host-desktop/src/main_tests.rs +++ b/crates/op-host-desktop/src/main_tests.rs @@ -184,7 +184,7 @@ fn live_mcp_http_server_applies_write_requests_to_editor_state() { use op_editor_core::PenNodeExt; - fn start_live_server() -> (mcp_live::McpLiveServer, u16) { + fn start_live_server() -> (op_web_daemon::mcp_live::McpLiveServer, u16) { // `bind(0)` to grab an ephemeral port, then re-`start` on that port, // has a TOCTOU window where the OS can reassign the port between the // probe-listener drop and the server bind — so a single attempt @@ -195,7 +195,7 @@ fn live_mcp_http_server_applies_write_requests_to_editor_state() { let listener = TcpListener::bind(("127.0.0.1", 0)).expect("bind ephemeral port"); listener.local_addr().expect("local addr").port() }; - if let Ok(server) = mcp_live::McpLiveServer::start(port) { + if let Ok(server) = op_web_daemon::mcp_live::McpLiveServer::start(port) { return (server, port); } } @@ -257,13 +257,13 @@ fn live_mcp_http_server_waits_for_split_http_request() { use std::sync::mpsc; use std::time::{Duration, Instant}; - fn start_live_server() -> (mcp_live::McpLiveServer, u16) { + fn start_live_server() -> (op_web_daemon::mcp_live::McpLiveServer, u16) { for _ in 0..20 { let port = { let listener = TcpListener::bind(("127.0.0.1", 0)).expect("bind ephemeral port"); listener.local_addr().expect("local addr").port() }; - if let Ok(server) = mcp_live::McpLiveServer::start(port) { + if let Ok(server) = op_web_daemon::mcp_live::McpLiveServer::start(port) { return (server, port); } } @@ -320,7 +320,7 @@ fn live_mcp_http_server_routes_file_path_requests_to_target_file() { use std::sync::mpsc; use std::time::{Duration, Instant}; - fn start_live_server() -> (mcp_live::McpLiveServer, u16) { + fn start_live_server() -> (op_web_daemon::mcp_live::McpLiveServer, u16) { // `bind(0)` to grab an ephemeral port, then re-`start` on that port, // has a TOCTOU window where the OS can reassign the port between the // probe-listener drop and the server bind — so a single attempt @@ -331,7 +331,7 @@ fn live_mcp_http_server_routes_file_path_requests_to_target_file() { let listener = TcpListener::bind(("127.0.0.1", 0)).expect("bind ephemeral port"); listener.local_addr().expect("local addr").port() }; - if let Ok(server) = mcp_live::McpLiveServer::start(port) { + if let Ok(server) = op_web_daemon::mcp_live::McpLiveServer::start(port) { return (server, port); } } @@ -420,13 +420,13 @@ fn live_mcp_http_server_replaces_document_via_rest_document_sync() { use std::sync::mpsc; use std::time::{Duration, Instant}; - fn start_live_server() -> (mcp_live::McpLiveServer, u16) { + fn start_live_server() -> (op_web_daemon::mcp_live::McpLiveServer, u16) { for _ in 0..20 { let port = { let listener = TcpListener::bind(("127.0.0.1", 0)).expect("bind ephemeral port"); listener.local_addr().expect("local addr").port() }; - if let Ok(server) = mcp_live::McpLiveServer::start(port) { + if let Ok(server) = op_web_daemon::mcp_live::McpLiveServer::start(port) { return (server, port); } } diff --git a/crates/op-host-desktop/src/mcp_runtime.rs b/crates/op-host-desktop/src/mcp_runtime.rs index 3109c8a37..ee0a2e902 100644 --- a/crates/op-host-desktop/src/mcp_runtime.rs +++ b/crates/op-host-desktop/src/mcp_runtime.rs @@ -1,6 +1,6 @@ //! Desktop-app glue for the live MCP server and terminal integrations. -use super::{mcp_integrations, mcp_live, DesktopApp}; +use super::{mcp_integrations, DesktopApp}; impl DesktopApp { pub(crate) fn bootstrap_mcp_runtime_from_settings(&mut self) -> bool { @@ -78,7 +78,7 @@ impl DesktopApp { if let Some(mut server) = self.mcp_server.take() { server.stop(); } - match mcp_live::McpLiveServer::start_with_wake(port, self.mcp_wake_callback()) { + match op_web_daemon::mcp_live::McpLiveServer::start_with_wake(port, self.mcp_wake_callback()) { Ok(server) => { let bound_port = server.port(); self.mcp_server = Some(server); @@ -88,7 +88,7 @@ impl DesktopApp { Err(err) => { eprintln!("openpencil-desktop mcp: failed to start on {port}: {err}"); if port != 0 { - match mcp_live::McpLiveServer::start_with_wake(0, self.mcp_wake_callback()) { + match op_web_daemon::mcp_live::McpLiveServer::start_with_wake(0, self.mcp_wake_callback()) { Ok(server) => { let bound_port = server.port(); eprintln!( diff --git a/crates/op-host-desktop/src/web_canvas_server.rs b/crates/op-host-desktop/src/web_canvas_server.rs index 21bf833e8..65e2e5418 100644 --- a/crates/op-host-desktop/src/web_canvas_server.rs +++ b/crates/op-host-desktop/src/web_canvas_server.rs @@ -1372,11 +1372,11 @@ fn serve_one( #[cfg(feature = "mcp-debug-tools")] if let Some(response) = { let guard = state.lock().unwrap_or_else(|p| p.into_inner()); - crate::mcp_live::screenshot::maybe_serve( + op_web_daemon::mcp_live::screenshot::maybe_serve( &req.body, op_mcp::debug_tools_enabled(), |shot_req| { - let spec = crate::mcp_live::screenshot::capture_spec(&shot_req); + let spec = op_web_daemon::mcp_live::screenshot::capture_spec(&shot_req); op_web_daemon::export::screenshot::capture(&guard.editor, &spec) }, ) diff --git a/crates/op-web-daemon/src/lib.rs b/crates/op-web-daemon/src/lib.rs index e9a9c1052..371c93849 100644 --- a/crates/op-web-daemon/src/lib.rs +++ b/crates/op-web-daemon/src/lib.rs @@ -33,6 +33,7 @@ pub mod design_session; pub mod doc_io; pub mod export; pub mod export_pdf; +pub mod mcp_live; pub mod mcp_serve; pub mod model_discovery; pub mod pre_validator; diff --git a/crates/op-host-desktop/src/mcp_live.rs b/crates/op-web-daemon/src/mcp_live.rs similarity index 90% rename from crates/op-host-desktop/src/mcp_live.rs rename to crates/op-web-daemon/src/mcp_live.rs index 473817e48..1eb1a9d2e 100644 --- a/crates/op-host-desktop/src/mcp_live.rs +++ b/crates/op-web-daemon/src/mcp_live.rs @@ -15,7 +15,7 @@ use std::time::{Duration, SystemTime, UNIX_EPOCH}; use op_editor_core::{EditorCommand, EditorState}; #[cfg(feature = "mcp-debug-tools")] -pub(crate) mod screenshot; +pub mod screenshot; /// Per-request budget for a UI-thread snapshot/apply ack. Sized to cover a /// large single editor operation without the connection giving up; the CLI @@ -35,7 +35,7 @@ const LIVE_CONN_STACK_SIZE: usize = 16 * 1024 * 1024; const _: () = assert!(LIVE_CONN_STACK_SIZE <= 16 * 1024 * 1024); type UiWake = Arc; -pub(crate) struct McpLiveServer { +pub struct McpLiveServer { port: u16, /// Per-instance identity token, reported in the live `ping` reply and /// written into `~/.openpencil/.op-mcp-port`. Lets the `op` CLI confirm @@ -55,9 +55,9 @@ struct ApplyAck { } #[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] -pub(crate) struct McpPumpOutcome { - pub(crate) repaint: bool, - pub(crate) layout_dirty: bool, +pub struct McpPumpOutcome { + pub repaint: bool, + pub layout_dirty: bool, } enum UiRequest { @@ -84,18 +84,18 @@ enum UiRequest { /// bounded-wait discipline as the chat canvas tools. #[cfg(feature = "mcp-debug-tools")] Screenshot { - spec: op_web_daemon::export::screenshot::CaptureSpec, - ack: SyncSender>, + spec: crate::export::screenshot::CaptureSpec, + ack: SyncSender>, }, } impl McpLiveServer { - #[cfg(test)] - pub(crate) fn start(port: u16) -> Result { + #[doc(hidden)] + pub fn start(port: u16) -> Result { Self::start_with_wake(port, || {}) } - pub(crate) fn start_with_wake(port: u16, wake_ui: F) -> Result + pub fn start_with_wake(port: u16, wake_ui: F) -> Result where F: Fn() + Send + Sync + 'static, { @@ -138,22 +138,22 @@ impl McpLiveServer { }) } - pub(crate) fn port(&self) -> u16 { + pub fn port(&self) -> u16 { self.port } /// Identity token to publish in the discovery file. - pub(crate) fn token(&self) -> &str { + pub fn token(&self) -> &str { &self.token } /// Whether a token-authed `openpencil/shutdown` was received — the UI /// thread should exit the event loop. - pub(crate) fn shutdown_requested(&self) -> bool { + pub fn shutdown_requested(&self) -> bool { self.quit_flag.load(Ordering::Acquire) } - pub(crate) fn pump(&mut self, state: &mut EditorState) -> McpPumpOutcome { + pub fn pump(&mut self, state: &mut EditorState) -> McpPumpOutcome { let mut outcome = McpPumpOutcome::default(); for _ in 0..UI_PUMP_REQUEST_BUDGET { match self.req_rx.try_recv() { @@ -182,7 +182,7 @@ impl McpLiveServer { #[cfg(feature = "mcp-debug-tools")] Ok(UiRequest::Screenshot { spec, ack }) => { // Read-only render of the live state — no repaint needed. - let _ = ack.send(op_web_daemon::export::screenshot::capture(state, &spec)); + let _ = ack.send(crate::export::screenshot::capture(state, &spec)); } Err(TryRecvError::Empty) => break, Err(TryRecvError::Disconnected) => break, @@ -191,7 +191,7 @@ impl McpLiveServer { outcome } - pub(crate) fn stop(&mut self) { + pub fn stop(&mut self) { let _ = self.stop_tx.send(()); } } @@ -257,7 +257,7 @@ fn server_loop( if conn_count.load(Ordering::Acquire) >= MAX_LIVE_CONNS { // Shed load rather than spawn unbounded threads. let _ = stream.set_write_timeout(Some(ACCEPT_IDLE_SLEEP)); - let _ = op_web_daemon::mcp_serve::write_mcp_http_response( + let _ = crate::mcp_serve::write_mcp_http_response( &mut stream, "503 Service Unavailable", r#"{"jsonrpc":"2.0","error":{"code":-32000,"message":"server busy"},"id":null}"#, @@ -290,7 +290,7 @@ fn server_loop( serve_connection(&mut stream, &req_tx, &token, &lock, &quit, &wake) { eprintln!("openpencil-desktop mcp: {e}"); - let _ = op_web_daemon::mcp_serve::write_mcp_http_response( + let _ = crate::mcp_serve::write_mcp_http_response( &mut stream, "500 Internal Server Error", &error_json(&e), @@ -322,26 +322,26 @@ fn serve_connection( quit_flag: &AtomicBool, wake_ui: &UiWake, ) -> Result<(), String> { - let req = op_web_daemon::mcp_serve::read_http_request(stream)?; + let req = crate::mcp_serve::read_http_request(stream)?; if req.method == "OPTIONS" { - return op_web_daemon::mcp_serve::write_mcp_http_response(stream, "204 No Content", ""); + return crate::mcp_serve::write_mcp_http_response(stream, "204 No Content", ""); } // TS live-canvas whole-document sync (REST `POST /api/mcp/document`), // distinct from the JSON-RPC `/mcp` path below. Lets a TS whole-doc-sync // client (`setSyncDocument` → POST `{document}`) drive THIS editor's // on-screen canvas, mirroring `apps/web/server/api/mcp/document.post.ts`. - if op_web_daemon::mcp_serve::is_document_sync_route(&req.method, &req.path) { + if crate::mcp_serve::is_document_sync_route(&req.method, &req.path) { return serve_document_sync(stream, req_tx, wake_ui, stateful_lock, &req.body); } if req.path != "/mcp" && req.path != "/" { - return op_web_daemon::mcp_serve::write_mcp_http_response( + return crate::mcp_serve::write_mcp_http_response( stream, "404 Not Found", r#"{"error":"Not found"}"#, ); } if req.method != "POST" { - return op_web_daemon::mcp_serve::write_mcp_http_response( + return crate::mcp_serve::write_mcp_http_response( stream, "400 Bad Request", r#"{"jsonrpc":"2.0","error":{"code":-32000,"message":"Invalid or missing session ID"},"id":null}"#, @@ -349,8 +349,8 @@ fn serve_connection( } // Token-authed graceful shutdown: ack, then flag the UI thread to exit // the event loop. No pid-kill ⇒ no signal-the-wrong-process race. - if let Some(id) = op_web_daemon::mcp_serve::shutdown_request_id(&req.body, token) { - write_json_rpc_response(stream, &op_web_daemon::mcp_serve::shutdown_ok_response(&id))?; + if let Some(id) = crate::mcp_serve::shutdown_request_id(&req.body, token) { + write_json_rpc_response(stream, &crate::mcp_serve::shutdown_ok_response(&id))?; quit_flag.store(true, Ordering::Release); wake_ui(); return Ok(()); @@ -359,17 +359,17 @@ fn serve_connection( // take the stateful lock, so `op`'s `ping` probe stays fast and never // false-negatives a busy editor. The live `ping` reply carries our // identity token so the CLI can confirm THIS server published the file. - match op_web_daemon::mcp_serve::classify_stateless(&req.body) { - op_web_daemon::mcp_serve::Stateless::Respond(resp) => { + match crate::mcp_serve::classify_stateless(&req.body) { + crate::mcp_serve::Stateless::Respond(resp) => { return write_json_rpc_response(stream, &resp); } - op_web_daemon::mcp_serve::Stateless::Swallow => { + crate::mcp_serve::Stateless::Swallow => { return write_json_rpc_response(stream, ""); } - op_web_daemon::mcp_serve::Stateless::Ping(id) => { + crate::mcp_serve::Stateless::Ping(id) => { return write_json_rpc_response(stream, &live_ping_response(&id, token)); } - op_web_daemon::mcp_serve::Stateless::NeedsState => {} + crate::mcp_serve::Stateless::NeedsState => {} } // Everything below mutates shared state — the live `EditorState` OR a // `--file` document on disk (a read-modify-write). Serialize ALL of it @@ -382,7 +382,7 @@ fn serve_connection( // File-backed path (`--file` arg): handle the whole read-modify-write // while holding the lock. if let Some(response) = - op_web_daemon::mcp_serve::file_path::process_message_for_file_path_arg(None, &req.body)? + crate::mcp_serve::file_path::process_message_for_file_path_arg(None, &req.body)? { return write_json_rpc_response(stream, &response); } @@ -401,7 +401,7 @@ fn serve_connection( return write_json_rpc_response(stream, &response); } let mut state = request_snapshot(req_tx, wake_ui)?; - let response = op_web_daemon::mcp_serve::process_message_with_applier( + let response = crate::mcp_serve::process_message_with_applier( &mut state, &req.body, |local_state, cmd| match request_apply(req_tx, wake_ui, cmd.clone()) { @@ -437,13 +437,13 @@ fn serve_document_sync( stateful_lock: &Mutex<()>, body: &str, ) -> Result<(), String> { - let document_json = match op_web_daemon::mcp_serve::parse_document_sync_body(body) { + let document_json = match crate::mcp_serve::parse_document_sync_body(body) { Ok(json) => json, Err(message) => { - return op_web_daemon::mcp_serve::write_mcp_http_response( + return crate::mcp_serve::write_mcp_http_response( stream, "400 Bad Request", - &op_web_daemon::mcp_serve::rest_error_body(&message), + &crate::mcp_serve::rest_error_body(&message), ); } }; @@ -454,10 +454,10 @@ fn serve_document_sync( let loaded = match op_pen_loader::load_canonical(&document_json) { Ok(loaded) => loaded, Err(e) => { - return op_web_daemon::mcp_serve::write_mcp_http_response( + return crate::mcp_serve::write_mcp_http_response( stream, "400 Bad Request", - &op_web_daemon::mcp_serve::rest_error_body(&e.to_string()), + &crate::mcp_serve::rest_error_body(&e.to_string()), ); } }; @@ -474,18 +474,18 @@ fn serve_document_sync( match request_replace(req_tx, wake_ui, loaded.value) { Ok(()) => { let version = LIVE_SYNC_VERSION.fetch_add(1, Ordering::Relaxed) + 1; - op_web_daemon::mcp_serve::write_mcp_http_response( + crate::mcp_serve::write_mcp_http_response( stream, "200 OK", - &op_web_daemon::mcp_serve::document_sync_ok(version), + &crate::mcp_serve::document_sync_ok(version), ) } // The UI thread is gone or didn't ack in time — a server fault, mapped // to 500 (TS throws → 500 for server-side failures). - Err(transport_err) => op_web_daemon::mcp_serve::write_mcp_http_response( + Err(transport_err) => crate::mcp_serve::write_mcp_http_response( stream, "500 Internal Server Error", - &op_web_daemon::mcp_serve::rest_error_body(&transport_err), + &crate::mcp_serve::rest_error_body(&transport_err), ), } } @@ -495,9 +495,9 @@ fn write_json_rpc_response( response: &str, ) -> Result<(), String> { if response.is_empty() { - op_web_daemon::mcp_serve::write_mcp_http_response(stream, "202 Accepted", "") + crate::mcp_serve::write_mcp_http_response(stream, "202 Accepted", "") } else { - op_web_daemon::mcp_serve::write_mcp_http_response(stream, "200 OK", response) + crate::mcp_serve::write_mcp_http_response(stream, "200 OK", response) } } @@ -532,7 +532,7 @@ fn request_screenshot( req_tx: &Sender, wake_ui: &UiWake, req: op_mcp::ScreenshotRequest, -) -> Result { +) -> Result { let timeout = Duration::from_millis(req.timeout_ms.max(1)); let spec = screenshot::capture_spec(&req); let (ack_tx, ack_rx) = mpsc::sync_channel(1); @@ -591,15 +591,15 @@ fn make_live_token() -> String { fn live_ping_response(id_raw: &str, token: &str) -> String { format!( r#"{{"jsonrpc":"2.0","id":{id_raw},"result":{{"server":"{}","mode":"live","token":"{}"}}}}"#, - op_web_daemon::mcp_serve::MCP_SERVER_NAME, - op_web_daemon::mcp_serve::json_escape(token) + crate::mcp_serve::MCP_SERVER_NAME, + crate::mcp_serve::json_escape(token) ) } fn error_json(message: &str) -> String { format!( r#"{{"error":"{}"}}"#, - op_web_daemon::mcp_serve::json_escape(message) + crate::mcp_serve::json_escape(message) ) } @@ -724,7 +724,7 @@ mod tests { let (ack_tx, ack_rx) = mpsc::sync_channel(1); req_tx .send(UiRequest::Screenshot { - spec: op_web_daemon::export::screenshot::CaptureSpec { + spec: crate::export::screenshot::CaptureSpec { node_id: None, padding: 0.0, scale: 1.0, diff --git a/crates/op-host-desktop/src/mcp_live/screenshot.rs b/crates/op-web-daemon/src/mcp_live/screenshot.rs similarity index 97% rename from crates/op-host-desktop/src/mcp_live/screenshot.rs rename to crates/op-web-daemon/src/mcp_live/screenshot.rs index b275091c2..ab7eeaabf 100644 --- a/crates/op-host-desktop/src/mcp_live/screenshot.rs +++ b/crates/op-web-daemon/src/mcp_live/screenshot.rs @@ -12,13 +12,13 @@ use op_mcp::{RequestId, ScreenshotRequest, ScreenshotTarget, ToolErrorCode, ToolResponse}; use serde_json::{json, Value}; -use op_web_daemon::export::screenshot::{CaptureSpec, ScreenshotPng}; +use crate::export::screenshot::{CaptureSpec, ScreenshotPng}; /// If `body` is a `tools/call` for `debug_screenshot` AND the debug /// gate is open, produce the full JSON-RPC response via `fulfill`. /// `None` ⇒ not a screenshot call (or gate closed) — the caller falls /// through to the generic dispatch. -pub(crate) fn maybe_serve(body: &str, debug_enabled: bool, fulfill: F) -> Option +pub fn maybe_serve(body: &str, debug_enabled: bool, fulfill: F) -> Option where F: FnOnce(ScreenshotRequest) -> Result, { @@ -44,7 +44,7 @@ where } /// Convert validated wire args into the renderer-side capture spec. -pub(crate) fn capture_spec(req: &ScreenshotRequest) -> CaptureSpec { +pub fn capture_spec(req: &ScreenshotRequest) -> CaptureSpec { CaptureSpec { node_id: match &req.target { ScreenshotTarget::Root => None,