feat(mcp): enable live desktop server
This commit is contained in:
parent
1610a01da0
commit
c21e9160fd
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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<mcp_live::McpLiveServer>,
|
||||
}
|
||||
|
||||
impl DesktopApp {
|
||||
|
|
@ -232,6 +237,7 @@ impl DesktopApp {
|
|||
git_clone_job: None,
|
||||
git_clone_origin: None,
|
||||
last_git_refresh: Instant::now(),
|
||||
mcp_server: None,
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
);
|
||||
}
|
||||
|
|
|
|||
258
crates/op-host-desktop/src/mcp_integrations.rs
Normal file
258
crates/op-host-desktop/src/mcp_integrations.rs
Normal file
|
|
@ -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<PathBuf, String> {
|
||||
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<PathBuf, String> {
|
||||
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<PathBuf, String> {
|
||||
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<Map<String, Value>, 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<String, Value>) -> 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);
|
||||
}
|
||||
}
|
||||
206
crates/op-host-desktop/src/mcp_live.rs
Normal file
206
crates/op-host-desktop/src/mcp_live.rs
Normal file
|
|
@ -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<UiRequest>,
|
||||
stop_tx: Sender<()>,
|
||||
}
|
||||
|
||||
struct ApplyAck {
|
||||
applied: bool,
|
||||
state: EditorState,
|
||||
}
|
||||
|
||||
enum UiRequest {
|
||||
Snapshot {
|
||||
ack: SyncSender<EditorState>,
|
||||
},
|
||||
Apply {
|
||||
cmd: EditorCommand,
|
||||
ack: SyncSender<ApplyAck>,
|
||||
},
|
||||
}
|
||||
|
||||
impl McpLiveServer {
|
||||
pub(crate) fn start(port: u16) -> Result<Self, String> {
|
||||
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<UiRequest>, 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<S: std::io::Read + std::io::Write>(
|
||||
stream: &mut S,
|
||||
req_tx: &Sender<UiRequest>,
|
||||
) -> 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<UiRequest>) -> Result<EditorState, String> {
|
||||
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<UiRequest>, cmd: EditorCommand) -> Result<ApplyAck, String> {
|
||||
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<T>(result: Result<T, RecvTimeoutError>, label: &str) -> Result<T, String> {
|
||||
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<S: Write>(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
|
||||
}
|
||||
114
crates/op-host-desktop/src/mcp_runtime.rs
Normal file
114
crates/op-host-desktop/src/mcp_runtime.rs
Normal file
|
|
@ -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
|
||||
}
|
||||
}
|
||||
|
|
@ -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<Option<String>, String> {
|
||||
let mut applier_failed: Option<String> = 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<F>(
|
||||
state: &mut EditorState,
|
||||
line: &str,
|
||||
mut apply: F,
|
||||
) -> Result<Option<String>, 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<String> = None;
|
||||
let mut out: Vec<u8> = 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<S: std::io::Read + std::io::Write>(
|
|||
/// 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<S: std::io::Read>(stream: &mut S) -> Result<String, String> {
|
||||
pub(crate) fn read_http_request_body<S: std::io::Read>(stream: &mut S) -> Result<String, String> {
|
||||
const MAX_HEADER: usize = 64 * 1024;
|
||||
const MAX_BODY: usize = 8 * 1024 * 1024;
|
||||
let mut head: Vec<u8> = Vec::new();
|
||||
|
|
|
|||
Loading…
Reference in a new issue