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.
This commit is contained in:
Kayshen-X 2026-06-19 21:27:35 +08:00
parent 7acd125ccd
commit d5727d322d
7 changed files with 66 additions and 66 deletions

View file

@ -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_live::McpLiveServer>,
mcp_server: Option<op_web_daemon::mcp_live::McpLiveServer>,
/// 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

View file

@ -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);
}
}

View file

@ -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!(

View file

@ -1372,11 +1372,11 @@ fn serve_one<S: Read + Write>(
#[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)
},
)

View file

@ -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;

View file

@ -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<dyn Fn() + Send + Sync + 'static>;
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<Result<op_web_daemon::export::screenshot::ScreenshotPng, String>>,
spec: crate::export::screenshot::CaptureSpec,
ack: SyncSender<Result<crate::export::screenshot::ScreenshotPng, String>>,
},
}
impl McpLiveServer {
#[cfg(test)]
pub(crate) fn start(port: u16) -> Result<Self, String> {
#[doc(hidden)]
pub fn start(port: u16) -> Result<Self, String> {
Self::start_with_wake(port, || {})
}
pub(crate) fn start_with_wake<F>(port: u16, wake_ui: F) -> Result<Self, String>
pub fn start_with_wake<F>(port: u16, wake_ui: F) -> Result<Self, String>
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<S: std::io::Read + std::io::Write>(
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<S: std::io::Read + std::io::Write>(
}
// 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<S: std::io::Read + std::io::Write>(
// 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<S: std::io::Read + std::io::Write>(
// 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<S: std::io::Read + std::io::Write>(
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<S: std::io::Read + std::io::Write>(
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<S: std::io::Read + std::io::Write>(
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<S: std::io::Read + std::io::Write>(
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<S: std::io::Write>(
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<UiRequest>,
wake_ui: &UiWake,
req: op_mcp::ScreenshotRequest,
) -> Result<op_web_daemon::export::screenshot::ScreenshotPng, String> {
) -> Result<crate::export::screenshot::ScreenshotPng, String> {
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,

View file

@ -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<F>(body: &str, debug_enabled: bool, fulfill: F) -> Option<String>
pub fn maybe_serve<F>(body: &str, debug_enabled: bool, fulfill: F) -> Option<String>
where
F: FnOnce(ScreenshotRequest) -> Result<ScreenshotPng, String>,
{
@ -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,