fix(mcp): make live server startup usable

This commit is contained in:
Kayshen-X 2026-05-31 21:05:27 +08:00
parent a0a1f4da04
commit 1a2e97c2bd
5 changed files with 150 additions and 48 deletions

View file

@ -5,7 +5,6 @@
//! 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;
@ -107,7 +106,7 @@ fn server_loop(listener: TcpListener, req_tx: Sender<UiRequest>, stop_rx: Receiv
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(
let _ = crate::mcp_serve::write_mcp_http_response(
&mut stream,
"500 Internal Server Error",
&error_json(&e),
@ -129,23 +128,45 @@ 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 req = crate::mcp_serve::read_http_request(stream)?;
if req.method == "OPTIONS" {
return crate::mcp_serve::write_mcp_http_response(stream, "204 No Content", "");
}
if req.path != "/mcp" && req.path != "/" {
return crate::mcp_serve::write_mcp_http_response(
stream,
"404 Not Found",
r#"{"error":"Not found"}"#,
);
}
if req.method != "POST" {
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}"#,
);
}
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
}
let response = crate::mcp_serve::process_message_with_applier(
&mut state,
&req.body,
|local_state, cmd| match request_apply(req_tx, cmd.clone()) {
Ok(ack) => {
*local_state = ack.state;
ack.applied
}
})?
.unwrap_or_default();
write_http_response(stream, "200 OK", &response)
Err(e) => {
eprintln!("openpencil-desktop mcp: apply failed: {e}");
false
}
},
)?
.unwrap_or_default();
if response.is_empty() {
crate::mcp_serve::write_mcp_http_response(stream, "202 Accepted", "")
} else {
crate::mcp_serve::write_mcp_http_response(stream, "200 OK", &response)
}
}
fn request_snapshot(req_tx: &Sender<UiRequest>) -> Result<EditorState, String> {
@ -172,19 +193,6 @@ fn recv_with_timeout<T>(result: Result<T, RecvTimeoutError>, label: &str) -> Res
}
}
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))
}

View file

@ -256,25 +256,37 @@ fn serve_http_connection<S: std::io::Read + std::io::Write>(
state: &mut EditorState,
path: &std::path::Path,
) -> Result<(), String> {
let body = read_http_request_body(stream)?;
let response = process_message(state, path, &body)?.unwrap_or_default();
let http = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\n\
Content-Length: {}\r\nConnection: close\r\n\r\n{}",
response.len(),
response
);
stream
.write_all(http.as_bytes())
.map_err(|e| format!("http write: {e}"))?;
stream.flush().map_err(|e| format!("http flush: {e}"))
let req = read_http_request(stream)?;
if req.method == "OPTIONS" {
return write_mcp_http_response(stream, "204 No Content", "");
}
if req.path != "/mcp" && req.path != "/" {
return write_mcp_http_response(stream, "404 Not Found", r#"{"error":"Not found"}"#);
}
if req.method != "POST" {
return write_mcp_http_response(
stream,
"400 Bad Request",
r#"{"jsonrpc":"2.0","error":{"code":-32000,"message":"Invalid or missing session ID"},"id":null}"#,
);
}
match process_message(state, path, &req.body)? {
Some(response) => write_mcp_http_response(stream, "200 OK", &response),
None => write_mcp_http_response(stream, "202 Accepted", ""),
}
}
/// Read an HTTP request off `stream` and return its body. Reads to
/// 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.
pub(crate) fn read_http_request_body<S: std::io::Read>(stream: &mut S) -> Result<String, String> {
pub(crate) struct HttpRequest {
pub method: String,
pub path: String,
pub body: String,
}
/// Read an HTTP request off `stream`. Reads to the `\r\n\r\n` header
/// terminator, parses the request line + `Content-Length`, then reads
/// exactly that many body bytes. The header block is capped so a
/// malformed peer can't exhaust memory.
pub(crate) fn read_http_request<S: std::io::Read>(stream: &mut S) -> Result<HttpRequest, String> {
const MAX_HEADER: usize = 64 * 1024;
const MAX_BODY: usize = 8 * 1024 * 1024;
let mut head: Vec<u8> = Vec::new();
@ -295,6 +307,19 @@ pub(crate) fn read_http_request_body<S: std::io::Read>(stream: &mut S) -> Result
}
}
let headers = String::from_utf8_lossy(&head);
let mut lines = headers.lines();
let request_line = lines
.next()
.ok_or_else(|| "request line missing".to_string())?;
let mut request_parts = request_line.split_whitespace();
let method = request_parts
.next()
.ok_or_else(|| "request method missing".to_string())?
.to_ascii_uppercase();
let path = request_parts
.next()
.ok_or_else(|| "request path missing".to_string())?
.to_string();
let content_length = headers
.lines()
.find_map(|l| {
@ -311,7 +336,41 @@ pub(crate) fn read_http_request_body<S: std::io::Read>(stream: &mut S) -> Result
stream
.read_exact(&mut body)
.map_err(|e| format!("http body read: {e}"))?;
Ok(String::from_utf8_lossy(&body).into_owned())
Ok(HttpRequest {
method,
path,
body: String::from_utf8_lossy(&body).into_owned(),
})
}
/// Compatibility wrapper for older tests/callers that only care about
/// the JSON-RPC body.
#[cfg(test)]
pub(crate) fn read_http_request_body<S: std::io::Read>(stream: &mut S) -> Result<String, String> {
read_http_request(stream).map(|req| req.body)
}
pub(crate) fn write_mcp_http_response<S: std::io::Write>(
stream: &mut S,
status: &str,
body: &str,
) -> Result<(), String> {
let http = format!(
"HTTP/1.1 {status}\r\n\
Access-Control-Allow-Origin: *\r\n\
Access-Control-Allow-Methods: GET, POST, DELETE, OPTIONS\r\n\
Access-Control-Allow-Headers: Content-Type, mcp-session-id\r\n\
Access-Control-Expose-Headers: mcp-session-id\r\n\
mcp-session-id: openpencil\r\n\
Content-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}"))
}
/// Re-build the registry against the latest editor state so read-tool

View file

@ -239,8 +239,29 @@ fn http_transport_serves_initialize() {
let resp = String::from_utf8(stream.output).unwrap();
assert!(resp.starts_with("HTTP/1.1 200 OK"), "status line: {resp}");
assert!(resp.contains("Content-Type: application/json"));
assert!(resp.contains("mcp-session-id: openpencil"));
assert!(resp.contains("Access-Control-Allow-Origin: *"));
// The JSON-RPC initialize reply carries the protocol handshake +
// the request id, proving the body round-tripped over HTTP.
assert!(resp.contains(r#""protocolVersion""#), "body: {resp}");
assert!(resp.contains(r#""id":7"#), "body: {resp}");
}
#[test]
fn http_transport_serves_options_preflight() {
let request = "OPTIONS /mcp HTTP/1.1\r\nHost: x\r\nContent-Length: 0\r\n\r\n";
let mut stream = MockStream {
input: std::io::Cursor::new(request.as_bytes().to_vec()),
output: Vec::new(),
};
let mut state = EditorState::new();
serve_http_connection(
&mut stream,
&mut state,
std::path::Path::new("/tmp/unused.op"),
)
.expect("serve_http_connection");
let resp = String::from_utf8(stream.output).unwrap();
assert!(resp.starts_with("HTTP/1.1 204 No Content"), "{resp}");
assert!(resp.contains("Access-Control-Allow-Methods"));
}

View file

@ -94,6 +94,13 @@ impl WidgetHostNative {
.unwrap_or(0);
let v = &mut self.editor_state.editor_ui.agent_settings.mcp_cli_enabled[idx];
*v = !*v;
if *v {
self.editor_state
.editor_ui
.agent_settings
.mcp_server
.running = true;
}
}
AgentSettingsHit::ToggleImagesAdvanced => {
let v = &mut self

View file

@ -52,6 +52,13 @@ impl WidgetHost {
.position(|candidate| *candidate == cli)
.unwrap_or(0);
self.editor_state.editor_ui.agent_settings.mcp_cli_enabled[idx] ^= true;
if self.editor_state.editor_ui.agent_settings.mcp_cli_enabled[idx] {
self.editor_state
.editor_ui
.agent_settings
.mcp_server
.running = true;
}
}
AgentSettingsHit::ToggleImagesAdvanced => {
self.editor_state