diff --git a/crates/op-host-desktop/src/app_handler.rs b/crates/op-host-desktop/src/app_handler.rs index 0087ff16e..be23c2435 100644 --- a/crates/op-host-desktop/src/app_handler.rs +++ b/crates/op-host-desktop/src/app_handler.rs @@ -272,6 +272,13 @@ impl ApplicationHandler for DesktopApp { WindowEvent::CursorMoved { .. } => None, _ => Some(settings_io::fingerprint(self.host.editor_state())), }; + let mcp_cli_before = match &event { + WindowEvent::CursorMoved { .. } => None, + _ => { + let settings = &self.host.editor_state().editor_ui.agent_settings; + Some((settings.mcp_cli_enabled, settings.mcp_server.port)) + } + }; match event { WindowEvent::CloseRequested => { // The unsaved-changes prompt can abort the close. @@ -482,6 +489,12 @@ impl ApplicationHandler for DesktopApp { if self.poll_git_clone_job() { self.redraw_dirty = true; } + // Drain live MCP requests. Write tools must apply on the + // UI-owned EditorState so canvas state, history and + // selection stay canonical. + if self.poll_mcp_server() { + self.redraw_dirty = true; + } // Keep an open Git panel fresh against external repo // changes — re-request a snapshot at most every 2 s. // The query runs on a worker thread, so this never @@ -529,7 +542,7 @@ impl ApplicationHandler for DesktopApp { event_loop.set_control_flow(ControlFlow::WaitUntil( Instant::now() + Duration::from_millis(33), )); - } else if self.current_figma_import.is_some() { + } else if self.current_figma_import.is_some() || self.mcp_server_active() { event_loop.set_control_flow(ControlFlow::WaitUntil( Instant::now() + Duration::from_millis(100), )); @@ -887,6 +900,12 @@ impl ApplicationHandler for DesktopApp { } _ => {} } + if self.reconcile_mcp_server_from_settings() { + self.request_redraw(true); + } + if self.reconcile_mcp_cli_integrations(mcp_cli_before) { + self.request_redraw(true); + } if let Some(before) = settings_before { settings_io::save_if_changed(self.host.editor_state(), before); } @@ -912,6 +931,9 @@ impl ApplicationHandler for DesktopApp { // a focused-but-uncommitted edit isn't silently dropped. self.host.flush_settings_input(); settings_io::save(self.host.editor_state()); + if let Some(mut server) = self.mcp_server.take() { + server.stop(); + } // Save window geometry for next launch. Guarded on a window // having existed — a failed startup reaches `exiting` with // unseeded geometry and would clobber the previous good save. diff --git a/crates/op-host-desktop/src/main.rs b/crates/op-host-desktop/src/main.rs index 520f0e88a..6c7219425 100644 --- a/crates/op-host-desktop/src/main.rs +++ b/crates/op-host-desktop/src/main.rs @@ -29,6 +29,9 @@ mod iconify_host; mod image_search_session; mod keyboard_input; mod macos_app; +mod mcp_integrations; +mod mcp_live; +mod mcp_runtime; mod mcp_serve; mod menu; mod model_discovery; @@ -167,6 +170,8 @@ struct DesktopApp { /// periodic refresh that keeps an open panel current against /// external repository changes. last_git_refresh: Instant, + /// Live in-process MCP HTTP server, started from Settings -> MCP. + mcp_server: Option, } impl DesktopApp { @@ -232,6 +237,7 @@ impl DesktopApp { git_clone_job: None, git_clone_origin: None, last_git_refresh: Instant::now(), + mcp_server: None, } } diff --git a/crates/op-host-desktop/src/main_tests.rs b/crates/op-host-desktop/src/main_tests.rs index 84e6bd397..51fae21be 100644 --- a/crates/op-host-desktop/src/main_tests.rs +++ b/crates/op-host-desktop/src/main_tests.rs @@ -66,3 +66,63 @@ fn fresh_app_refits_blank_frame_to_actual_window_size_once() { let unchanged = app.host.editor_state().viewport; assert_eq!(v, unchanged); } + +#[test] +fn live_mcp_http_server_applies_write_requests_to_editor_state() { + use std::io::{Read, Write}; + use std::net::{TcpListener, TcpStream}; + use std::sync::mpsc; + use std::time::{Duration, Instant}; + + use op_editor_core::PenNodeExt; + + fn unused_port() -> u16 { + let listener = TcpListener::bind(("127.0.0.1", 0)).expect("bind ephemeral port"); + listener.local_addr().expect("local addr").port() + } + + fn post_json(port: u16, body: &str) -> String { + let mut stream = TcpStream::connect(("127.0.0.1", port)).expect("connect MCP server"); + let req = format!( + "POST /mcp HTTP/1.1\r\nHost: 127.0.0.1\r\nContent-Length: {}\r\n\r\n{}", + body.len(), + body + ); + stream.write_all(req.as_bytes()).expect("write request"); + let mut out = String::new(); + stream.read_to_string(&mut out).expect("read response"); + out + } + + let port = unused_port(); + let mut server = mcp_live::McpLiveServer::start(port).expect("start MCP server"); + let mut state = op_editor_core::EditorState::new(); + let body = r##"{"jsonrpc":"2.0","id":1,"method":"insert_node","params":{"kind":"rect","name":"From MCP","x":"10","y":"20","width":"100","height":"50","fill_hex":"#00ff00"}}"##; + let (tx, rx) = mpsc::channel(); + std::thread::spawn(move || { + let _ = tx.send(post_json(port, body)); + }); + + let started = Instant::now(); + let response = loop { + server.pump(&mut state); + if let Ok(response) = rx.try_recv() { + break response; + } + assert!( + started.elapsed() < Duration::from_secs(2), + "MCP request timed out" + ); + std::thread::sleep(Duration::from_millis(10)); + }; + + assert!(response.starts_with("HTTP/1.1 200 OK"), "{response}"); + assert!(response.contains(r#""wrote":"true""#), "{response}"); + assert!( + state + .active_children() + .iter() + .any(|node| node.base().name.as_deref() == Some("From MCP")), + "MCP write should mutate the live editor state" + ); +} diff --git a/crates/op-host-desktop/src/mcp_integrations.rs b/crates/op-host-desktop/src/mcp_integrations.rs new file mode 100644 index 000000000..b01881649 --- /dev/null +++ b/crates/op-host-desktop/src/mcp_integrations.rs @@ -0,0 +1,258 @@ +//! Terminal-side MCP client configuration for the desktop settings panel. +//! +//! The live server is owned by the GUI process. CLI integrations point +//! at that process over streamable HTTP so terminal agents can reach the +//! same canvas state the user is editing. + +use std::fs; +use std::path::{Path, PathBuf}; + +use op_editor_core::agent_settings::McpCli; +use serde_json::{Map, Value}; + +const SERVER_NAME: &str = "openpencil"; + +pub(crate) fn set_cli_enabled(cli: McpCli, enabled: bool, port: u16) -> Result { + let home = dirs::home_dir().ok_or_else(|| "home directory not available".to_string())?; + let path = config_path(cli, &home, true); + set_cli_enabled_at_path(cli, enabled, port, path) +} + +#[cfg(test)] +fn set_cli_enabled_at_home( + cli: McpCli, + enabled: bool, + port: u16, + home: &Path, +) -> Result { + let path = config_path(cli, home, false); + set_cli_enabled_at_path(cli, enabled, port, path) +} + +fn set_cli_enabled_at_path( + cli: McpCli, + enabled: bool, + port: u16, + path: PathBuf, +) -> Result { + match cli { + McpCli::Codex => update_codex_config(&path, enabled, port)?, + McpCli::ClaudeCode + | McpCli::Gemini + | McpCli::OpenCode + | McpCli::Kiro + | McpCli::GithubCopilot => update_json_config(&path, enabled, port)?, + } + Ok(path) +} + +fn config_path(cli: McpCli, home: &Path, use_env: bool) -> PathBuf { + match cli { + McpCli::ClaudeCode => home.join(".claude.json"), + McpCli::Codex => { + if use_env { + std::env::var_os("CODEX_HOME") + .map(PathBuf::from) + .unwrap_or_else(|| home.join(".codex")) + .join("config.toml") + } else { + home.join(".codex").join("config.toml") + } + } + McpCli::Gemini => home.join(".gemini").join("settings.json"), + McpCli::OpenCode => home.join(".opencode").join("config.json"), + McpCli::Kiro => home.join(".kiro").join("settings.json"), + McpCli::GithubCopilot => home.join(".config").join("github-copilot").join("mcp.json"), + } +} + +fn update_json_config(path: &Path, enabled: bool, port: u16) -> Result<(), String> { + let mut root = read_json_object(path)?; + if enabled { + let servers = root + .entry("mcpServers") + .or_insert_with(|| Value::Object(Map::new())); + if !servers.is_object() { + *servers = Value::Object(Map::new()); + } + let Some(servers) = servers.as_object_mut() else { + return Err("mcpServers is not an object".into()); + }; + servers.insert( + SERVER_NAME.into(), + serde_json::json!({ + "type": "http", + "url": endpoint(port), + }), + ); + } else if let Some(servers) = root.get_mut("mcpServers").and_then(Value::as_object_mut) { + servers.remove(SERVER_NAME); + if servers.is_empty() { + root.remove("mcpServers"); + } + } + write_json_object(path, &root) +} + +fn read_json_object(path: &Path) -> Result, String> { + let text = match fs::read_to_string(path) { + Ok(text) => text, + Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Map::new()), + Err(e) => return Err(format!("read {}: {e}", path.display())), + }; + if text.trim().is_empty() { + return Ok(Map::new()); + } + let value: Value = + serde_json::from_str(&text).map_err(|e| format!("parse {}: {e}", path.display()))?; + value + .as_object() + .cloned() + .ok_or_else(|| format!("{} must contain a JSON object", path.display())) +} + +fn write_json_object(path: &Path, root: &Map) -> Result<(), String> { + if let Some(parent) = path.parent() { + fs::create_dir_all(parent).map_err(|e| format!("create {}: {e}", parent.display()))?; + } + let text = serde_json::to_string_pretty(root) + .map_err(|e| format!("serialize {}: {e}", path.display()))?; + fs::write(path, format!("{text}\n")).map_err(|e| format!("write {}: {e}", path.display())) +} + +fn update_codex_config(path: &Path, enabled: bool, port: u16) -> Result<(), String> { + let existing = match fs::read_to_string(path) { + Ok(text) => text, + Err(e) if e.kind() == std::io::ErrorKind::NotFound => String::new(), + Err(e) => return Err(format!("read {}: {e}", path.display())), + }; + let mut text = remove_codex_server_block(&existing); + if enabled { + let prefix = text.trim_end(); + text = String::from(prefix); + if !text.is_empty() { + text.push_str("\n\n"); + } + text.push_str("[mcp_servers.openpencil]\n"); + text.push_str(&format!( + "url = \"{}\"\n", + toml_basic_string_escape(&endpoint(port)) + )); + } + if let Some(parent) = path.parent() { + fs::create_dir_all(parent).map_err(|e| format!("create {}: {e}", parent.display()))?; + } + fs::write(path, text).map_err(|e| format!("write {}: {e}", path.display())) +} + +fn remove_codex_server_block(input: &str) -> String { + let mut out = String::new(); + let mut skipping = false; + for line in input.split_inclusive('\n') { + let trimmed = line.trim(); + if is_codex_openpencil_table(trimmed) { + skipping = true; + continue; + } + if skipping && trimmed.starts_with('[') { + skipping = false; + } + if !skipping { + out.push_str(line); + } + } + out +} + +fn is_codex_openpencil_table(line: &str) -> bool { + matches!( + line, + "[mcp_servers.openpencil]" + | "[mcp_servers.\"openpencil\"]" + | "[\"mcp_servers\".\"openpencil\"]" + ) +} + +fn endpoint(port: u16) -> String { + format!("http://127.0.0.1:{port}/mcp") +} + +fn toml_basic_string_escape(s: &str) -> String { + s.replace('\\', "\\\\").replace('"', "\\\"") +} + +#[cfg(test)] +mod tests { + use super::*; + + fn temp_home(name: &str) -> PathBuf { + let path = std::env::temp_dir().join(format!( + "openpencil-mcp-{name}-{}-{}", + std::process::id(), + std::thread::current().name().unwrap_or("test") + )); + let _ = fs::remove_dir_all(&path); + fs::create_dir_all(&path).expect("create temp home"); + path + } + + #[test] + fn mcp_json_config_install_and_uninstall_preserves_other_servers() { + let home = temp_home("json"); + let path = home.join(".claude.json"); + fs::write( + &path, + r#"{"theme":"dark","mcpServers":{"other":{"type":"http","url":"http://x"}}}"#, + ) + .expect("seed config"); + + set_cli_enabled_at_home(McpCli::ClaudeCode, true, 3101, &home).expect("install"); + let text = fs::read_to_string(&path).expect("read installed"); + assert!(text.contains(r#""openpencil""#), "{text}"); + assert!( + text.contains(r#""url": "http://127.0.0.1:3101/mcp""#), + "{text}" + ); + assert!(text.contains(r#""other""#), "{text}"); + + set_cli_enabled_at_home(McpCli::ClaudeCode, false, 3101, &home).expect("uninstall"); + let text = fs::read_to_string(&path).expect("read uninstalled"); + assert!(!text.contains(r#""openpencil""#), "{text}"); + assert!(text.contains(r#""other""#), "{text}"); + + let _ = fs::remove_dir_all(home); + } + + #[test] + fn mcp_codex_config_replaces_existing_openpencil_block() { + let home = temp_home("codex"); + let path = home.join(".codex").join("config.toml"); + fs::create_dir_all(path.parent().expect("parent")).expect("create codex dir"); + fs::write( + &path, + "model = \"gpt-5\"\n\n[mcp_servers.openpencil]\nurl = \"http://old\"\n\n[profiles.dev]\nmodel = \"gpt-5-codex\"\n", + ) + .expect("seed config"); + + set_cli_enabled_at_home(McpCli::Codex, true, 3200, &home).expect("install"); + let text = fs::read_to_string(&path).expect("read installed"); + assert_eq!( + text.matches("[mcp_servers.openpencil]").count(), + 1, + "{text}" + ); + assert!( + text.contains("url = \"http://127.0.0.1:3200/mcp\""), + "{text}" + ); + assert!(text.contains("[profiles.dev]"), "{text}"); + + set_cli_enabled_at_home(McpCli::Codex, false, 3200, &home).expect("uninstall"); + let text = fs::read_to_string(&path).expect("read uninstalled"); + assert!(!text.contains("[mcp_servers.openpencil]"), "{text}"); + assert!(text.contains("model = \"gpt-5\""), "{text}"); + assert!(text.contains("[profiles.dev]"), "{text}"); + + let _ = fs::remove_dir_all(home); + } +} diff --git a/crates/op-host-desktop/src/mcp_live.rs b/crates/op-host-desktop/src/mcp_live.rs new file mode 100644 index 000000000..e0069cf00 --- /dev/null +++ b/crates/op-host-desktop/src/mcp_live.rs @@ -0,0 +1,206 @@ +//! In-process MCP HTTP server for the live desktop editor. +//! +//! The CLI-only `mcp_serve` path owns a file-backed `EditorState`. +//! This module keeps the GUI as the source of truth: the server thread +//! requests a fresh snapshot from the UI thread for each HTTP request, +//! then sends write commands back for the UI thread to apply. + +use std::io::Write; +use std::net::TcpListener; +use std::sync::mpsc::{self, Receiver, RecvTimeoutError, Sender, SyncSender, TryRecvError}; +use std::thread; +use std::time::Duration; + +use op_editor_core::{EditorCommand, EditorState}; + +const UI_ACK_TIMEOUT: Duration = Duration::from_secs(5); +const ACCEPT_IDLE_SLEEP: Duration = Duration::from_millis(25); + +pub(crate) struct McpLiveServer { + port: u16, + req_rx: Receiver, + stop_tx: Sender<()>, +} + +struct ApplyAck { + applied: bool, + state: EditorState, +} + +enum UiRequest { + Snapshot { + ack: SyncSender, + }, + Apply { + cmd: EditorCommand, + ack: SyncSender, + }, +} + +impl McpLiveServer { + pub(crate) fn start(port: u16) -> Result { + let listener = TcpListener::bind(("127.0.0.1", port)) + .map_err(|e| format!("bind 127.0.0.1:{port}: {e}"))?; + listener + .set_nonblocking(true) + .map_err(|e| format!("set nonblocking: {e}"))?; + let (req_tx, req_rx) = mpsc::channel(); + let (stop_tx, stop_rx) = mpsc::channel(); + thread::Builder::new() + .name("op-mcp-live-http".into()) + .spawn(move || server_loop(listener, req_tx, stop_rx)) + .map_err(|e| format!("spawn MCP live server: {e}"))?; + eprintln!("openpencil-desktop mcp: listening on 127.0.0.1:{port}/mcp"); + Ok(Self { + port, + req_rx, + stop_tx, + }) + } + + pub(crate) fn port(&self) -> u16 { + self.port + } + + pub(crate) fn pump(&mut self, state: &mut EditorState) -> bool { + let mut any_applied = false; + loop { + match self.req_rx.try_recv() { + Ok(UiRequest::Snapshot { ack }) => { + let _ = ack.send(state.clone()); + } + Ok(UiRequest::Apply { cmd, ack }) => { + let applied = state.apply(cmd); + let _ = ack.send(ApplyAck { + applied, + state: state.clone(), + }); + any_applied |= applied; + } + Err(TryRecvError::Empty) => break, + Err(TryRecvError::Disconnected) => break, + } + } + any_applied + } + + pub(crate) fn stop(&mut self) { + let _ = self.stop_tx.send(()); + } +} + +impl Drop for McpLiveServer { + fn drop(&mut self) { + self.stop(); + } +} + +fn server_loop(listener: TcpListener, req_tx: Sender, stop_rx: Receiver<()>) { + loop { + match stop_rx.try_recv() { + Ok(()) | Err(TryRecvError::Disconnected) => break, + Err(TryRecvError::Empty) => {} + } + match listener.accept() { + Ok((mut stream, _addr)) => { + let _ = stream.set_read_timeout(Some(UI_ACK_TIMEOUT)); + let _ = stream.set_write_timeout(Some(UI_ACK_TIMEOUT)); + if let Err(e) = serve_connection(&mut stream, &req_tx) { + eprintln!("openpencil-desktop mcp: {e}"); + let _ = write_http_response( + &mut stream, + "500 Internal Server Error", + &error_json(&e), + ); + } + } + Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => { + thread::sleep(ACCEPT_IDLE_SLEEP); + } + Err(e) => { + eprintln!("openpencil-desktop mcp: accept: {e}"); + thread::sleep(ACCEPT_IDLE_SLEEP); + } + } + } +} + +fn serve_connection( + stream: &mut S, + req_tx: &Sender, +) -> Result<(), String> { + let body = crate::mcp_serve::read_http_request_body(stream)?; + let mut state = request_snapshot(req_tx)?; + let response = + crate::mcp_serve::process_message_with_applier(&mut state, &body, |local_state, cmd| { + match request_apply(req_tx, cmd.clone()) { + Ok(ack) => { + *local_state = ack.state; + ack.applied + } + Err(e) => { + eprintln!("openpencil-desktop mcp: apply failed: {e}"); + false + } + } + })? + .unwrap_or_default(); + write_http_response(stream, "200 OK", &response) +} + +fn request_snapshot(req_tx: &Sender) -> Result { + let (ack_tx, ack_rx) = mpsc::sync_channel(1); + req_tx + .send(UiRequest::Snapshot { ack: ack_tx }) + .map_err(|_| "UI thread is not accepting MCP snapshot requests".to_string())?; + recv_with_timeout(ack_rx.recv_timeout(UI_ACK_TIMEOUT), "snapshot") +} + +fn request_apply(req_tx: &Sender, cmd: EditorCommand) -> Result { + let (ack_tx, ack_rx) = mpsc::sync_channel(1); + req_tx + .send(UiRequest::Apply { cmd, ack: ack_tx }) + .map_err(|_| "UI thread is not accepting MCP apply requests".to_string())?; + recv_with_timeout(ack_rx.recv_timeout(UI_ACK_TIMEOUT), "apply") +} + +fn recv_with_timeout(result: Result, label: &str) -> Result { + match result { + Ok(v) => Ok(v), + Err(RecvTimeoutError::Timeout) => Err(format!("timed out waiting for UI {label} ack")), + Err(RecvTimeoutError::Disconnected) => Err(format!("UI {label} ack channel closed")), + } +} + +fn write_http_response(stream: &mut S, status: &str, body: &str) -> Result<(), String> { + let http = format!( + "HTTP/1.1 {status}\r\nContent-Type: application/json\r\n\ + Content-Length: {}\r\nConnection: close\r\n\r\n{}", + body.len(), + body + ); + stream + .write_all(http.as_bytes()) + .map_err(|e| format!("http write: {e}"))?; + stream.flush().map_err(|e| format!("http flush: {e}")) +} + +fn error_json(message: &str) -> String { + format!(r#"{{"error":"{}"}}"#, json_escape(message)) +} + +fn json_escape(s: &str) -> String { + let mut out = String::with_capacity(s.len()); + for c in s.chars() { + match c { + '"' => out.push_str("\\\""), + '\\' => out.push_str("\\\\"), + '\n' => out.push_str("\\n"), + '\r' => out.push_str("\\r"), + '\t' => out.push_str("\\t"), + c if c.is_control() => out.push(' '), + c => out.push(c), + } + } + out +} diff --git a/crates/op-host-desktop/src/mcp_runtime.rs b/crates/op-host-desktop/src/mcp_runtime.rs new file mode 100644 index 000000000..ea2546425 --- /dev/null +++ b/crates/op-host-desktop/src/mcp_runtime.rs @@ -0,0 +1,114 @@ +//! Desktop-app glue for the live MCP server and terminal integrations. + +use super::{mcp_integrations, mcp_live, DesktopApp}; + +impl DesktopApp { + pub(crate) fn reconcile_mcp_server_from_settings(&mut self) -> bool { + let desired = self + .host + .editor_state() + .editor_ui + .agent_settings + .mcp_server + .running; + let port = self + .host + .editor_state() + .editor_ui + .agent_settings + .mcp_server + .port; + if !desired { + if let Some(mut server) = self.mcp_server.take() { + server.stop(); + } + return false; + } + if self + .mcp_server + .as_ref() + .is_some_and(|server| server.port() == port) + { + return false; + } + if let Some(mut server) = self.mcp_server.take() { + server.stop(); + } + match mcp_live::McpLiveServer::start(port) { + Ok(server) => { + self.mcp_server = Some(server); + false + } + Err(err) => { + eprintln!("openpencil-desktop mcp: failed to start: {err}"); + self.host + .editor_state_mut() + .editor_ui + .agent_settings + .mcp_server + .running = false; + self.host.mark_editor_state_dirty(); + true + } + } + } + + pub(crate) fn poll_mcp_server(&mut self) -> bool { + let Some(server) = self.mcp_server.as_mut() else { + return false; + }; + let changed = server.pump(self.host.editor_state_mut()); + if changed { + self.host.mark_editor_state_dirty(); + } + changed + } + + pub(crate) fn mcp_server_active(&self) -> bool { + self.mcp_server.is_some() + } + + pub(crate) fn reconcile_mcp_cli_integrations( + &mut self, + before: Option<([bool; 6], u16)>, + ) -> bool { + let Some((before_flags, before_port)) = before else { + return false; + }; + let settings = &self.host.editor_state().editor_ui.agent_settings; + let after_flags = settings.mcp_cli_enabled; + let port = settings.mcp_server.port; + if before_flags == after_flags && before_port == port { + return false; + } + + let mut reverted = false; + for (idx, cli) in op_editor_core::agent_settings::McpCli::ALL + .iter() + .copied() + .enumerate() + { + let flag_changed = before_flags[idx] != after_flags[idx]; + let enabled_port_changed = before_port != port && after_flags[idx]; + if !flag_changed && !enabled_port_changed { + continue; + } + if let Err(err) = mcp_integrations::set_cli_enabled(cli, after_flags[idx], port) { + eprintln!( + "openpencil-desktop mcp: failed to update {} integration: {err}", + cli.label() + ); + if flag_changed { + self.host + .editor_state_mut() + .editor_ui + .agent_settings + .mcp_cli_enabled[idx] = before_flags[idx]; + self.host.mark_editor_state_dirty(); + reverted = true; + } + } + } + reverted + } +} diff --git a/crates/op-host-desktop/src/mcp_serve.rs b/crates/op-host-desktop/src/mcp_serve.rs index be7c1bbc9..51d4993e1 100644 --- a/crates/op-host-desktop/src/mcp_serve.rs +++ b/crates/op-host-desktop/src/mcp_serve.rs @@ -28,7 +28,7 @@ use std::io::{BufRead, BufReader, BufWriter, Write}; use std::path::PathBuf; -use op_editor_core::EditorState; +use op_editor_core::{EditorCommand, EditorState}; use op_mcp::{ add_node_effect_snapshot, add_page_snapshot, align_selected_snapshot, batch_design_snapshot, clear_selection_snapshot, copy_node_snapshot, copy_selected_snapshot, count_nodes_snapshot, @@ -86,6 +86,34 @@ fn process_message( path: &std::path::Path, line: &str, ) -> Result, String> { + let mut applier_failed: Option = None; + let response = process_message_with_applier(state, line, |state, cmd| { + // `EditorState::apply` runs the pre-validate-then-mutate + // discipline; `false` means the command rejected and the + // document was NOT changed. + if !state.apply(cmd.clone()) { + return false; + } + if let Err(e) = save_editor_state(state, path) { + applier_failed = Some(format!("save failed: {e}")); + return false; + } + true + })?; + if let Some(msg) = applier_failed { + eprintln!("openpencil-desktop mcp: {msg}"); + } + Ok(response) +} + +pub(crate) fn process_message_with_applier( + state: &mut EditorState, + line: &str, + mut apply: F, +) -> Result, String> +where + F: FnMut(&mut EditorState, &EditorCommand) -> bool, +{ let trimmed = line.trim(); if trimmed.is_empty() { return Ok(None); @@ -112,27 +140,11 @@ fn process_message( // registry snapshots `state` at build time, so it no longer // borrows it once the applier closure mutates it. let registry = rebuild_registry(state); - let mut applier_failed: Option = None; let mut out: Vec = Vec::new(); { let mut input = std::io::Cursor::new(line.as_bytes()); - run_stdio_with_applier(®istry, &mut input, &mut out, |cmd| { - // `EditorState::apply` runs the pre-validate-then-mutate - // discipline; `false` means the command rejected and the - // document was NOT changed. - if !state.apply(cmd.clone()) { - return false; - } - if let Err(e) = save_editor_state(state, path) { - applier_failed = Some(format!("save failed: {e}")); - return false; - } - true - }) - .map_err(|e| format!("dispatch: {e}"))?; - } - if let Some(msg) = applier_failed { - eprintln!("openpencil-desktop mcp: {msg}"); + run_stdio_with_applier(®istry, &mut input, &mut out, |cmd| apply(state, cmd)) + .map_err(|e| format!("dispatch: {e}"))?; } let resp = String::from_utf8_lossy(&out).trim().to_string(); Ok((!resp.is_empty()).then_some(resp)) @@ -262,7 +274,7 @@ fn serve_http_connection( /// the `\r\n\r\n` header terminator, parses `Content-Length`, then /// reads exactly that many body bytes. The header block is capped so /// a malformed peer can't exhaust memory. -fn read_http_request_body(stream: &mut S) -> Result { +pub(crate) fn read_http_request_body(stream: &mut S) -> Result { const MAX_HEADER: usize = 64 * 1024; const MAX_BODY: usize = 8 * 1024 * 1024; let mut head: Vec = Vec::new();