diff --git a/crates/op-cli/src/app_control_cli.rs b/crates/op-cli/src/app_control_cli.rs index d7560a129..41d88ed94 100644 --- a/crates/op-cli/src/app_control_cli.rs +++ b/crates/op-cli/src/app_control_cli.rs @@ -456,10 +456,37 @@ fn reachable_live_port_file() -> Option<(u16, u32)> { /// Both candidates are confirmed with a JSON-RPC `ping`. Returns `None` /// when nothing is reachable. pub(crate) fn discover_running_port() -> Option { - if let Some((port, _)) = reachable_live_port_file() { - return Some(port); + discover_running_endpoint().map(|(port, _)| port) +} + +/// The reachable MCP endpoint and the instance token that authenticates it. +/// +/// The live endpoint authenticates every stateful call, so the port alone is +/// not enough to drive it — the token has to travel with it. +pub(crate) fn discover_running_endpoint() -> Option<(u16, String)> { + if let Some((port, _, token)) = read_live_port_file() { + if crate::mcp_http_cli::mcp_ping_live(port, &token) { + return Some((port, token)); + } } - reachable_headless_server().map(|info| info.port) + reachable_headless_server().map(|info| (info.port, info.token)) +} + +/// The instance token for `port`, or empty when this process cannot find one. +/// +/// Used when the caller pinned `--port` explicitly: the token still has to be +/// looked up, and an empty result simply means the request goes out unauthenticated +/// and the endpoint decides. +pub(crate) fn token_for_port(port: u16) -> String { + if let Some((live_port, _, token)) = read_live_port_file() { + if live_port == port { + return token; + } + } + running_mcp_from_pid_file() + .filter(|info| info.port == port) + .map(|info| info.token) + .unwrap_or_default() } pub(crate) fn ensure_document_file(path: &Path) -> Result<(), CliError> { diff --git a/crates/op-cli/src/export_cli.rs b/crates/op-cli/src/export_cli.rs index b892fe37e..739860de1 100644 --- a/crates/op-cli/src/export_cli.rs +++ b/crates/op-cli/src/export_cli.rs @@ -50,6 +50,7 @@ pub(crate) fn map_export(flags: &Flags) -> Result { pub(crate) fn run_export( port: u16, + token: &str, item_id: Option<&str>, output: &str, format: &str, @@ -68,6 +69,7 @@ pub(crate) fn run_export( } let response = post( port, + token, &tool_call_body("export_item", &Value::Object(arguments).to_string()), )?; write_export_response(&response, Path::new(output)) diff --git a/crates/op-cli/src/main.rs b/crates/op-cli/src/main.rs index c374fb5bd..f2f8a4551 100644 --- a/crates/op-cli/src/main.rs +++ b/crates/op-cli/src/main.rs @@ -73,10 +73,13 @@ fn run(args: &[String]) -> Result { | Command::ToolCallJson { .. } | Command::Export { .. } ); - let target_port = if port_explicit || !needs_server { - port + // The live MCP endpoint authenticates every stateful call, so the token + // has to be resolved alongside the port rather than after it. + let (target_port, target_token) = if port_explicit || !needs_server { + (port, app_control_cli::token_for_port(port)) } else { - app_control_cli::discover_running_port().unwrap_or(port) + app_control_cli::discover_running_endpoint() + .unwrap_or_else(|| (port, app_control_cli::token_for_port(port))) }; let out = match command { Command::Help => USAGE.to_string(), @@ -100,7 +103,7 @@ fn run(args: &[String]) -> Result { } Command::InstallSkill { target } => skill_install_cli::run_install(target.as_deref())?, Command::UninstallSkill { target } => skill_install_cli::run_uninstall(target.as_deref())?, - Command::ToolsList => post(target_port, &tools_list_body())?, + Command::ToolsList => post(target_port, &target_token, &tools_list_body())?, Command::ImportFigma { fig_path, out_path } => { figma_cli::run_import_figma(&fig_path, &out_path)? } @@ -113,10 +116,10 @@ fn run(args: &[String]) -> Result { out_path, } => html_cli::run_import_snapshot(&json_path, &out_path)?, Command::ToolCall { tool, args } => { - post(target_port, &tool_call_body(&tool, &args_to_json(&args)))? + post(target_port, &target_token, &tool_call_body(&tool, &args_to_json(&args)))? } Command::ToolCallJson { tool, args_json } => { - post(target_port, &tool_call_body(&tool, &args_json))? + post(target_port, &target_token, &tool_call_body(&tool, &args_json))? } Command::Export { item_id, @@ -126,6 +129,7 @@ fn run(args: &[String]) -> Result { scale, } => export_cli::run_export( target_port, + &target_token, item_id.as_deref(), &output, &format, diff --git a/crates/op-cli/src/mcp_http_cli.rs b/crates/op-cli/src/mcp_http_cli.rs index b1929b9ef..03245c7be 100644 --- a/crates/op-cli/src/mcp_http_cli.rs +++ b/crates/op-cli/src/mcp_http_cli.rs @@ -102,6 +102,7 @@ pub(crate) fn http_request(body: &str) -> String { /// service accepts but never replies from hanging the CLI indefinitely. fn post_raw( port: u16, + token: &str, body: &str, timeout: std::time::Duration, ) -> Result<(u16, String), CliError> { @@ -109,6 +110,7 @@ fn post_raw( // failure domain; flatten it into the CLI's `Transport` variant, whose // payload is the exact sentence `op:` prints. TcpJsonRpc::local_mcp(port) + .with_token(token) .post_raw(body, timeout) .map(|reply| (reply.status, reply.body)) .map_err(|e| CliError::Transport(e.to_string())) @@ -120,8 +122,8 @@ fn post_raw( /// `tools/call` content envelope is unwrapped to the raw tool result, and an /// `isError` tool result (or a JSON-RPC transport error) is surfaced as an /// error — so `op` exits non-zero on a failed tool, like the TS CLI. -pub(crate) fn post(port: u16, body: &str) -> Result { - let (status, reply) = post_raw(port, body, POST_TIMEOUT)?; +pub(crate) fn post(port: u16, token: &str, body: &str) -> Result { + let (status, reply) = post_raw(port, token, body, POST_TIMEOUT)?; if !(200..300).contains(&status) { return Err(CliError::Transport(format!( "MCP server on 127.0.0.1:{port} returned HTTP {status}: {reply}" @@ -204,7 +206,7 @@ fn is_web_canvas_health(body: &str) -> bool { /// `/api/mcp/*`, not `/mcp`) — so discovery never mis-routes tool calls. fn ping_result(port: u16) -> Option { let body = serde_json::to_string(&JsonRpcRequest::new(0, "ping", Value::Null)).ok()?; - let (status, body) = post_raw(port, &body, PING_TIMEOUT).ok()?; + let (status, body) = post_raw(port, "", &body, PING_TIMEOUT).ok()?; if !(200..300).contains(&status) { return None; } @@ -256,7 +258,7 @@ pub(crate) fn request_shutdown(port: u16, token: &str) -> bool { serde_json::json!({ "token": token }), )) .expect("shutdown JSON-RPC request serializes"); - match post_raw(port, &body, SHUTDOWN_TIMEOUT) { + match post_raw(port, "", &body, SHUTDOWN_TIMEOUT) { Ok((status, body)) if (200..300).contains(&status) => serde_json::from_str::(&body) .ok() .and_then(|v| { diff --git a/crates/op-host-desktop/src/collab_runtime/relay_bootstrap.rs b/crates/op-host-desktop/src/collab_runtime/relay_bootstrap.rs index 07dd95274..c100c1388 100644 --- a/crates/op-host-desktop/src/collab_runtime/relay_bootstrap.rs +++ b/crates/op-host-desktop/src/collab_runtime/relay_bootstrap.rs @@ -75,8 +75,11 @@ pub(super) struct RelayBootstrap { generation: u64, signed_payload: Box<[u8]>, regions: Vec, - /// Whether this run persisted the document that arms the anti-rollback - /// generation floor for the next start. + /// Whether the anti-rollback generation floor is intact around this load. + /// + /// False when the cache could not be read (this run had no floor to check + /// the fetched generation against) or could not be written (the next start + /// will have none). /// /// A failed cache write degrades rather than failing the bootstrap: the /// document was still fully verified, and an unwritable configuration @@ -174,14 +177,20 @@ impl EnvironmentRelayBootstrapProvider { .get_or_init(|| Mutex::new(())) .lock() .map_err(|_| BootstrapError::Cache)?; - // An absent cache is an empty floor and reads as `Ok(None)`, so the - // unwritable-configuration-directory case never reaches this error - // path — only a cache that exists and is corrupt, unreadable, or not a - // regular file does. That is an anomaly rather than a first run, and - // continuing past it would let the anomaly itself disarm the - // generation floor, so resolution stops here. The write side degrades - // instead; see the persist step below. - let cached = read_cache(&self.cache_path, self.endpoint.as_str())?; + // A cache that exists but is corrupt, unreadable, or not a regular file + // leaves this start with no generation floor — the same position an + // absent cache leaves it in, which reads as `Ok(None)` and has always + // been allowed. The threat model already accepts a missing floor: + // deleting the file achieves it, and corrupting or chmod-ing the file + // needs no more privilege than deleting it, so refusing here buys no + // security while turning a damaged cache into a hard inability to + // collaborate. Resolution therefore continues without a floor, and the + // degradation is recorded rather than swallowed — it travels with the + // returned document as `rollback_floor_armed`. + let (cached, cache_readable) = match read_cache(&self.cache_path, self.endpoint.as_str()) { + Ok(cached) => (cached, true), + Err(_) => (None, false), + }; let cached_verified = cached.as_ref().and_then(|cached| { verify_bootstrap( cached.body.as_bytes(), @@ -260,7 +269,7 @@ impl EnvironmentRelayBootstrapProvider { // and `reject_rollback` have already accepted the document, so the // verifier stays exactly as fail-closed as before. let mut verified = verified; - verified.rollback_floor_armed = match String::from_utf8(body) { + let persisted = match String::from_utf8(body) { Ok(body) => write_cache( &self.cache_path, &BootstrapCache { @@ -272,6 +281,9 @@ impl EnvironmentRelayBootstrapProvider { .is_ok(), Err(_) => false, }; + // Both halves matter: an unreadable cache cost this run its floor, and + // a failed write costs the next start its floor. + verified.rollback_floor_armed = cache_readable && persisted; Ok(Arc::new(verified)) } diff --git a/crates/op-host-desktop/src/collab_runtime/relay_bootstrap_tests.rs b/crates/op-host-desktop/src/collab_runtime/relay_bootstrap_tests.rs index 305bd5da8..380af1244 100644 --- a/crates/op-host-desktop/src/collab_runtime/relay_bootstrap_tests.rs +++ b/crates/op-host-desktop/src/collab_runtime/relay_bootstrap_tests.rs @@ -618,7 +618,16 @@ fn provider_sends_etag_and_accepts_only_matching_not_modified() { let _ = std::fs::remove_dir_all(cache_root); } -fn assert_bad_cache_cannot_disarm_rollback_floor( +/// A damaged cache degrades to "no floor" instead of blocking bootstrap. +/// +/// This is a deliberate weakening relative to failing closed, and the test +/// states its cost plainly: with the floor gone, the lower-generation document +/// IS accepted. It is accepted because refusing bought no security — deleting +/// the cache file already achieves a missing floor and needs no more privilege +/// than corrupting it — while it did cost the ability to collaborate at all on +/// a machine with a damaged or unreadable cache. What must not happen is the +/// loss passing silently, so `rollback_floor_armed` has to report it. +fn assert_bad_cache_degrades_to_an_unarmed_rollback_floor( label: &str, damage_cache: impl FnOnce(&std::path::Path), ) { @@ -680,25 +689,36 @@ fn assert_bad_cache_cannot_disarm_rollback_floor( development_http: true, cache_path, }; - assert_eq!(provider.load_inner(NOW).unwrap_err(), BootstrapError::Cache); + let loaded = provider + .load_inner(NOW) + .expect("a damaged cache degrades instead of blocking bootstrap"); + assert_eq!( + loaded.generation, + valid_payload().generation, + "with no floor the lower-generation document is accepted — this is the cost of degrading" + ); + assert!( + !loaded.rollback_floor_armed, + "a damaged cache must report the anti-rollback floor as unarmed" + ); stop_sender.send(()).ok(); assert!( - !server.join().unwrap(), - "an unsafe cache must fail before a lower-generation response is fetched" + server.join().unwrap(), + "degrading means the fetch proceeds rather than stopping at the damaged cache" ); let _ = std::fs::remove_dir_all(cache_root); } #[test] -fn corrupt_cache_cannot_disarm_the_rollback_floor() { - assert_bad_cache_cannot_disarm_rollback_floor("corrupt", |path| { +fn corrupt_cache_degrades_to_an_unarmed_rollback_floor() { + assert_bad_cache_degrades_to_an_unarmed_rollback_floor("corrupt", |path| { std::fs::write(path, b"not json").unwrap(); }); } #[cfg(unix)] #[test] -fn unreadable_cache_cannot_disarm_the_rollback_floor() { +fn unreadable_cache_degrades_to_an_unarmed_rollback_floor() { use std::os::unix::fs::PermissionsExt as _; let probe = std::env::temp_dir().join(format!( @@ -716,7 +736,7 @@ fn unreadable_cache_cannot_disarm_the_rollback_floor() { if can_still_read { return; } - assert_bad_cache_cannot_disarm_rollback_floor("unreadable", |path| { + assert_bad_cache_degrades_to_an_unarmed_rollback_floor("unreadable", |path| { std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o000)).unwrap(); }); } diff --git a/crates/op-host-services/src/mcp_live.rs b/crates/op-host-services/src/mcp_live.rs index 6e7d77952..69a20fe43 100644 --- a/crates/op-host-services/src/mcp_live.rs +++ b/crates/op-host-services/src/mcp_live.rs @@ -175,7 +175,12 @@ impl McpLiveServer { let (req_tx, req_rx) = mpsc::channel(); let (stop_tx, stop_rx) = mpsc::channel(); let token = make_live_token(); - let server_token = token.clone(); + // The connection threads need both halves of the admission material: + // the token they must see on every state-touching request, and the + // port they actually bound (so `Host` can be pinned to it — the + // DNS-rebinding check). Bundled so `server_loop`'s argument list + // stays the same width. + let admission = Arc::new(LiveAdmission::new(token.clone(), bound_port)); let quit_flag = Arc::new(AtomicBool::new(false)); let server_quit = Arc::clone(&quit_flag); let wake_ui: UiWake = Arc::new(wake_ui); @@ -188,7 +193,7 @@ impl McpLiveServer { listener, req_tx, stop_rx, - server_token, + admission, server_quit, wake_ui, server_identity, @@ -458,10 +463,12 @@ fn command_invalidates_layout(cmd: &EditorCommand) -> bool { } } +mod admission; mod connection; mod doc_sync; mod ui_requests; +use admission::*; use connection::*; use doc_sync::*; use ui_requests::*; diff --git a/crates/op-host-services/src/mcp_live/admission.rs b/crates/op-host-services/src/mcp_live/admission.rs new file mode 100644 index 000000000..54b462e4f --- /dev/null +++ b/crates/op-host-services/src/mcp_live/admission.rs @@ -0,0 +1,292 @@ +//! Admission control for the live-GUI MCP endpoint (`127.0.0.1:/mcp`). +//! +//! The live endpoint drives the on-screen document, and during a +//! collaboration session that document is the SHARED one. Before this +//! module existed the per-instance token authenticated exactly two +//! messages — the `ping` identity probe's reply and `openpencil/shutdown` +//! — so every read and every write tool call was reachable by any local +//! process, and (with no `Origin` / `Host` screening) by any web page on +//! the machine via DNS rebinding. That is a straight bypass of the +//! collaboration admission model: `CollabGatePolicy` can refuse an MCP +//! *mutation* mid-session, but it never saw reads, and it is not a +//! boundary — it is the last line behind one. +//! +//! Two independent gates, in this order: +//! +//! 1. [`check_boundary`] — browser screening. A `Host` that is not a +//! numeric loopback literal on the bound port, or ANY `Origin` other +//! than this instance's own loopback origin, is refused. Both headers +//! are browser-controlled but not page-forgeable, which is what +//! actually closes DNS rebinding: a rebound `evil.com` page still +//! sends `Host: evil.com:` and `Origin: http://evil.com`. +//! Requests with no `Origin` at all are normal non-browser clients +//! (the `op` CLI, the VS Code MCP proxy) and pass. +//! 2. [`check_token`] — the per-instance token from the +//! `X-OpenPencil-Token` header, compared in constant time. This is the +//! same credential the server already published to its clients (the +//! discovery file `~/.openpencil/.op-mcp-port` and the `ping` reply) +//! and the same header the managed web daemon uses +//! (`web_canvas_server::RequestAuth`) — no new credential system. +//! +//! Deliberately still tokenless: `OPTIONS` preflight, and the stateless +//! `initialize` / `notifications/initialized` / `ping` probes, which carry +//! no document data and are how a client discovers this instance in the +//! first place. `openpencil/shutdown` keeps its own body-carried token +//! check (`mcp_serve::shutdown_request_id`), unchanged. +//! +//! KNOWN RESIDUAL (documented, not introduced here): the `ping` reply +//! hands the token to any local caller, so this gate raises the bar for a +//! local process rather than closing it outright. Closing that needs a +//! change to the `op` CLI's discovery handshake, which lives in another +//! crate. + +use std::fmt; + +/// JSON-RPC error code for a refused request. Server-defined range +/// (-32000..=-32099), one step away from the -32000 this endpoint already +/// uses for "server busy" / "Invalid or missing session ID". +const DENIED_CODE: i32 = -32001; + +/// Per-instance admission material for the live endpoint: the token the +/// server published and the port it actually bound. Shared by every +/// connection thread (`Arc`), immutable for the life of the server. +pub(super) struct LiveAdmission { + token: String, + port: u16, +} + +impl LiveAdmission { + pub(super) fn new(token: String, port: u16) -> Self { + Self { token, port } + } + + /// The per-instance identity token — also what the `ping` reply and + /// the `openpencil/shutdown` check use, so the wire contract the CLI + /// already knows is preserved verbatim. + pub(super) fn token(&self) -> &str { + &self.token + } + + pub(super) fn port(&self) -> u16 { + self.port + } +} + +/// Why a request was refused. A typed enum rather than a `String` (the +/// workspace rule) — and deliberately NOT an `McpLiveError` variant: every +/// value here is a *client* fault answered on the wire with a 401/403 and +/// a JSON-RPC error body, never a server fault the accept loop logs and +/// turns into a 500. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(super) enum AdmissionDenial { + /// `Host` absent, not a numeric loopback literal, or carrying a port + /// other than the one this server bound. + ForeignHost, + /// An `Origin` header that is not this instance's own loopback origin. + ForeignOrigin, + /// No `X-OpenPencil-Token` header (or an empty one). + MissingToken, + /// A token that is not this instance's token. + BadToken, +} + +impl AdmissionDenial { + /// HTTP status line. Boundary failures are 403 (the caller is not + /// allowed to talk to this endpoint at all); token failures are 401 + /// (the caller may retry with the credential it was given). + pub(super) fn http_status(self) -> &'static str { + match self { + AdmissionDenial::ForeignHost | AdmissionDenial::ForeignOrigin => "403 Forbidden", + AdmissionDenial::MissingToken | AdmissionDenial::BadToken => "401 Unauthorized", + } + } + + /// Client-facing reason. Intentionally coarse: it names the gate, not + /// which byte of the token differed. + pub(super) fn message(self) -> &'static str { + match self { + AdmissionDenial::ForeignHost => { + "live MCP endpoint accepts loopback requests only (bad Host header)" + } + AdmissionDenial::ForeignOrigin => { + "live MCP endpoint refuses cross-origin requests (bad Origin header)" + } + AdmissionDenial::MissingToken => { + "live MCP endpoint requires the X-OpenPencil-Token header" + } + AdmissionDenial::BadToken => "live MCP endpoint token mismatch", + } + } +} + +impl fmt::Display for AdmissionDenial { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str(self.message()) + } +} + +/// Gate 1 — browser screening, applied to EVERY request (including the +/// stateless probes and the REST document-sync route) before any routing. +pub(super) fn check_boundary( + req: &crate::mcp_serve::HttpRequest, + admission: &LiveAdmission, +) -> Result<(), AdmissionDenial> { + if !host_allowed(req.host.as_deref(), admission.port()) { + return Err(AdmissionDenial::ForeignHost); + } + if !origin_allowed(req.origin.as_deref(), admission.port()) { + return Err(AdmissionDenial::ForeignOrigin); + } + Ok(()) +} + +/// Gate 2 — per-instance token, required by every request that can +/// observe or mutate the live document. +pub(super) fn check_token( + req: &crate::mcp_serve::HttpRequest, + admission: &LiveAdmission, +) -> Result<(), AdmissionDenial> { + if admission.token().is_empty() { + // A server with no token can authenticate nobody. Refusing is the + // only safe reading — the alternative ("empty means open") is the + // bug this module exists to remove. + return Err(AdmissionDenial::MissingToken); + } + match req.token.as_deref() { + // No header at all, or a present-but-blank one: the same "brought + // no credential" case, reported as such rather than as a mismatch. + None | Some("") => Err(AdmissionDenial::MissingToken), + Some(presented) if constant_time_eq(presented, admission.token()) => Ok(()), + Some(_) => Err(AdmissionDenial::BadToken), + } +} + +/// Non-early-exit byte compare. +/// +/// `subtle` is NOT a dependency of `op-host-services` (only +/// `op-host-desktop` and the two relay crates carry it), and adding one +/// would mean editing a manifest a concurrent session owns — so this is +/// the hand-rolled equivalent: fold every byte difference into a single +/// accumulator, no `return` inside the loop, no data-dependent branch, one +/// comparison at the end. Length is compared up front on purpose: the +/// token's shape is public (`{pid:x}-{nanos:x}`, see `make_live_token`), +/// so its length is not a secret; what must not leak is HOW MANY leading +/// bytes of a same-length guess were right. +pub(super) fn constant_time_eq(presented: &str, expected: &str) -> bool { + let (presented, expected) = (presented.as_bytes(), expected.as_bytes()); + if presented.len() != expected.len() { + return false; + } + let mut diff = 0u8; + for (a, b) in presented.iter().zip(expected.iter()) { + diff |= a ^ b; + } + diff == 0 +} + +/// JSON-RPC error body for a refused `/mcp` request, echoing the caller's +/// request id so a client correlates the refusal with its call instead of +/// hanging (same discipline as `op_mcp::parser`'s parse-failure path). +pub(super) fn denial_json_rpc(request_body: &str, denial: AdmissionDenial) -> String { + format!( + r#"{{"jsonrpc":"2.0","error":{{"code":{DENIED_CODE},"message":"{}"}},"id":{}}}"#, + crate::mcp_serve::json_escape(denial.message()), + request_id_raw(request_body) + ) +} + +/// REST-shaped denial body for the `POST /api/mcp/document` route, which +/// speaks `{ok,error}` rather than JSON-RPC. +pub(super) fn denial_rest(denial: AdmissionDenial) -> String { + crate::mcp_serve::rest_error_body(denial.message()) +} + +/// The caller's top-level JSON-RPC `id`, verbatim (so a string id stays a +/// string id), or `null` when the body is not a JSON object / has no id. +fn request_id_raw(request_body: &str) -> String { + serde_json::from_str::(request_body) + .ok() + .and_then(|value| value.get("id").map(|id| id.to_string())) + .unwrap_or_else(|| "null".to_string()) +} + +/// `Host` must name a numeric loopback address. A DNS name — including +/// `localhost` — is refused: rebinding attacks work precisely by pointing +/// a name at 127.0.0.1, and a browser writes the name it was given into +/// `Host`, so accepting names would leave the hole open. +/// +/// The port is checked when present. It may be absent: `op`'s own +/// transport sends a bare `Host: 127.0.0.1` +/// (`op_rpc_transport::TcpJsonRpc::http_post_request`), and a browser can +/// never produce that against a non-80 port — it always writes the target +/// port it dialled. So "no port" identifies a non-browser client rather +/// than widening the browser surface. +fn host_allowed(host: Option<&str>, expected_port: u16) -> bool { + // HTTP/1.1 requires `Host`; every browser and every client in this + // repo sends it. Absent ⇒ refuse rather than guess. + let Some(raw) = host else { + return false; + }; + let Some((host, port)) = split_authority(raw.trim()) else { + return false; + }; + is_numeric_loopback(host) && port.is_none_or(|port| port == expected_port) +} + +/// Any `Origin` other than this instance's own loopback origin is refused. +/// `None` is the normal non-browser case (CLI / proxy) and passes — a page +/// cannot suppress the header on a cross-origin request, so "no Origin" +/// is not something an attacker page can claim. +fn origin_allowed(origin: Option<&str>, expected_port: u16) -> bool { + let Some(origin) = origin else { + return true; + }; + let origin = origin.trim(); + // The live endpoint is plain HTTP on loopback, so only an `http://` + // origin can possibly be it; `null` (sandboxed iframe / `file://`) + // and any `https://` page fall through to a refusal. + let Some(authority) = origin.strip_prefix("http://") else { + return false; + }; + // A real serialized origin is scheme + authority and nothing else. + if authority.contains(['/', '@', '?', '#']) { + return false; + } + let Some((host, port)) = split_authority(authority) else { + return false; + }; + is_numeric_loopback(host) && port.unwrap_or(80) == expected_port +} + +/// Split an HTTP authority (`127.0.0.1:3100`, `127.0.0.1`, `[::1]:3100`) +/// into host and optional port. A malformed port refuses the whole value. +fn split_authority(value: &str) -> Option<(&str, Option)> { + if let Some(rest) = value.strip_prefix('[') { + let (inside, tail) = rest.split_once(']')?; + return match tail { + "" => Some((inside, None)), + tail => { + let port = tail.strip_prefix(':')?.parse::().ok()?; + Some((inside, Some(port))) + } + }; + } + match value.rsplit_once(':') { + // An unbracketed IPv6 literal lands here with a nonsense split; + // `is_numeric_loopback` then rejects the truncated host, which is + // correct — unbracketed IPv6 in an authority is malformed anyway. + Some((host, port)) => Some((host, Some(port.parse::().ok()?))), + None => Some((value, None)), + } +} + +/// A numeric IP literal in a loopback range (127.0.0.0/8 or `::1`). +/// Names never qualify. +fn is_numeric_loopback(host: &str) -> bool { + host.parse::() + .is_ok_and(|address| address.is_loopback()) +} + +#[cfg(test)] +#[path = "admission_tests.rs"] +mod tests; diff --git a/crates/op-host-services/src/mcp_live/admission_tests.rs b/crates/op-host-services/src/mcp_live/admission_tests.rs new file mode 100644 index 000000000..ed74d3974 --- /dev/null +++ b/crates/op-host-services/src/mcp_live/admission_tests.rs @@ -0,0 +1,278 @@ +//! Admission-gate tests for the live-MCP endpoint. Driven through the real +//! `serve_connection` router (generic over the stream, like the web-canvas +//! daemon's connection tests) so they cover the wiring, not just the +//! predicates: a refused request must never reach the UI-request channel. + +use super::super::*; +use super::*; + +const PORT: u16 = 51234; +const TOKEN: &str = "d34db33f-cafe"; +/// A read-only tool call. Reads were entirely unauthenticated before this +/// gate existed, so the read path is exactly what these tests must pin. +const LIST_PAGES_CALL: &str = r#"{"jsonrpc":"2.0","id":11,"method":"tools/call","params":{"name":"list_pages","arguments":{}}}"#; + +struct MockStream { + input: std::io::Cursor>, + output: Vec, +} + +impl std::io::Read for MockStream { + fn read(&mut self, buf: &mut [u8]) -> std::io::Result { + std::io::Read::read(&mut self.input, buf) + } +} + +impl std::io::Write for MockStream { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.output.extend_from_slice(buf); + Ok(buf.len()) + } + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } +} + +/// Run one raw HTTP request through the live router and return the raw +/// response. +fn drive(request: &str, req_tx: &Sender) -> String { + let admission = LiveAdmission::new(TOKEN.to_string(), PORT); + let stateful_lock = Mutex::new(()); + let quit_flag = AtomicBool::new(false); + let wake_ui: UiWake = Arc::new(|| {}); + let client_identity = Mutex::new(None); + let mut stream = MockStream { + input: std::io::Cursor::new(request.as_bytes().to_vec()), + output: Vec::new(), + }; + serve_connection( + &mut stream, + req_tx, + &admission, + &stateful_lock, + &quit_flag, + &wake_ui, + &client_identity, + ) + .expect("a refused request is answered on the wire, never a server error"); + String::from_utf8_lossy(&stream.output).into_owned() +} + +fn request(path: &str, headers: &str, body: &str) -> String { + format!( + "POST {path} HTTP/1.1\r\n{headers}Content-Length: {}\r\nConnection: close\r\n\r\n{body}", + body.len() + ) +} + +/// Headers a token-carrying non-browser client sends (VS Code MCP proxy +/// shape: explicit port, no `Origin`). +fn authed_headers() -> String { + format!("Host: 127.0.0.1:{PORT}\r\nX-OpenPencil-Token: {TOKEN}\r\n") +} + +#[test] +fn unauthenticated_tool_call_is_refused() { + let (req_tx, req_rx) = mpsc::channel(); + let headers = format!("Host: 127.0.0.1:{PORT}\r\n"); + let response = drive(&request("/mcp", &headers, LIST_PAGES_CALL), &req_tx); + + assert!( + response.starts_with("HTTP/1.1 401 Unauthorized"), + "{response}" + ); + assert!(response.contains(r#""code":-32001"#), "{response}"); + // The caller's id is echoed so a client fails fast instead of hanging. + assert!(response.contains(r#""id":11"#), "{response}"); + assert!( + req_rx.try_recv().is_err(), + "a refused tool call must never reach the UI thread" + ); +} + +#[test] +fn authenticated_tool_call_is_served() { + let (req_tx, req_rx) = mpsc::channel(); + let responder = thread::spawn(move || match req_rx.recv_timeout(Duration::from_secs(5)) { + Ok(UiRequest::ListPages { ack }) => ack + .send(op_mcp::ListPages { + page_count: 3, + active_page_index: 1, + pages: vec![("p1".to_string(), "One".to_string())], + }) + .is_ok(), + _ => false, + }); + let response = drive( + &request("/mcp", &authed_headers(), LIST_PAGES_CALL), + &req_tx, + ); + + assert!( + responder.join().expect("responder thread"), + "an authenticated tool call must reach the UI thread" + ); + assert!(response.starts_with("HTTP/1.1 200 OK"), "{response}"); + assert!(response.contains("pageCount"), "{response}"); + assert!(!response.contains("-32001"), "{response}"); +} + +/// The `op` CLI's exact wire shape: no `Origin` (not a browser) and a bare +/// `Host: 127.0.0.1` with no port +/// (`op_rpc_transport::TcpJsonRpc::http_post_request`). It must keep +/// working with the token it is handed by the discovery file. +#[test] +fn cli_request_without_origin_or_host_port_is_served() { + let (req_tx, req_rx) = mpsc::channel(); + let responder = thread::spawn(move || match req_rx.recv_timeout(Duration::from_secs(5)) { + Ok(UiRequest::ListPages { ack }) => ack + .send(op_mcp::ListPages { + page_count: 1, + active_page_index: 0, + pages: vec![("p1".to_string(), "One".to_string())], + }) + .is_ok(), + _ => false, + }); + let headers = format!("Host: 127.0.0.1\r\nX-OpenPencil-Token: {TOKEN}\r\n"); + let response = drive(&request("/mcp", &headers, LIST_PAGES_CALL), &req_tx); + + assert!( + responder.join().expect("responder thread"), + "a portless-Host CLI request must still reach the UI thread" + ); + assert!(response.starts_with("HTTP/1.1 200 OK"), "{response}"); +} + +#[test] +fn foreign_origin_tool_call_is_refused() { + let (req_tx, req_rx) = mpsc::channel(); + // Even WITH the right token: a browser page that somehow learned the + // token is still not an allowed caller. + let headers = format!( + "Host: 127.0.0.1:{PORT}\r\nOrigin: http://evil.example\r\nX-OpenPencil-Token: {TOKEN}\r\n" + ); + let response = drive(&request("/mcp", &headers, LIST_PAGES_CALL), &req_tx); + + assert!(response.starts_with("HTTP/1.1 403 Forbidden"), "{response}"); + assert!(req_rx.try_recv().is_err(), "refused before the UI thread"); + + // Its own loopback origin on its own port is the one allowed value. + assert!(origin_allowed( + Some(&format!("http://127.0.0.1:{PORT}")), + PORT + )); + // Right host, wrong port — a different local server's page. + assert!(!origin_allowed(Some("http://127.0.0.1:1"), PORT)); + // `localhost` is a NAME, and names are the rebinding vector. + assert!(!origin_allowed( + Some(&format!("http://localhost:{PORT}")), + PORT + )); + assert!(!origin_allowed(Some("null"), PORT)); + assert!(origin_allowed(None, PORT), "non-browser clients pass"); +} + +#[test] +fn non_loopback_or_wrong_port_host_is_refused() { + let (req_tx, req_rx) = mpsc::channel(); + // The DNS-rebinding shape: the browser resolved `evil.example` to + // 127.0.0.1 but still writes the NAME into `Host`. + let headers = + format!("Host: evil.example:{PORT}\r\nX-OpenPencil-Token: {TOKEN}\r\nOrigin: http://evil.example\r\n"); + let response = drive(&request("/mcp", &headers, LIST_PAGES_CALL), &req_tx); + assert!(response.starts_with("HTTP/1.1 403 Forbidden"), "{response}"); + assert!(req_rx.try_recv().is_err(), "refused before the UI thread"); + + // A loopback literal, but a port this server never bound. + let headers = format!( + "Host: 127.0.0.1:{}\r\nX-OpenPencil-Token: {TOKEN}\r\n", + PORT + 1 + ); + let response = drive(&request("/mcp", &headers, LIST_PAGES_CALL), &req_tx); + assert!(response.starts_with("HTTP/1.1 403 Forbidden"), "{response}"); + assert!(req_rx.try_recv().is_err(), "refused before the UI thread"); + + assert!(host_allowed(Some(&format!("127.0.0.1:{PORT}")), PORT)); + assert!(host_allowed(Some(&format!("[::1]:{PORT}")), PORT)); + assert!(!host_allowed(Some(&format!("localhost:{PORT}")), PORT)); + assert!(!host_allowed(Some(&format!("10.0.0.4:{PORT}")), PORT)); + assert!(!host_allowed(None, PORT), "a missing Host is refused"); +} + +#[test] +fn constant_time_compare_rejects_same_length_wrong_token() { + // Same length, differing only in the last byte — the case an + // early-exit `==` would answer faster than a first-byte mismatch. + assert_eq!(TOKEN.len(), "d34db33f-caff".len()); + assert!(!constant_time_eq("d34db33f-caff", TOKEN)); + assert!(!constant_time_eq("e34db33f-cafe", TOKEN)); + assert!(constant_time_eq(TOKEN, TOKEN)); + assert!(!constant_time_eq("", TOKEN)); + assert!(!constant_time_eq(&format!("{TOKEN}x"), TOKEN)); + + let admission = LiveAdmission::new(TOKEN.to_string(), PORT); + let same_length_guess = format!( + "Host: 127.0.0.1:{PORT}\r\nX-OpenPencil-Token: d34db33f-caff\r\nContent-Length: 0\r\n\r\n" + ); + let mut cursor = std::io::Cursor::new(format!("POST /mcp HTTP/1.1\r\n{same_length_guess}")); + let req = crate::mcp_serve::read_http_request(&mut cursor).expect("request parses"); + assert_eq!( + check_token(&req, &admission), + Err(AdmissionDenial::BadToken) + ); + + // An instance with no token authenticates nobody. + let tokenless = LiveAdmission::new(String::new(), PORT); + assert_eq!( + check_token(&req, &tokenless), + Err(AdmissionDenial::MissingToken) + ); +} + +/// The identity probes stay tokenless on purpose: `op` discovers this +/// instance by pinging it and matching the reply's token against the +/// discovery file. Gating `ping` would break discovery for every CLI. +#[test] +fn ping_probe_stays_tokenless_for_cli_discovery() { + let (req_tx, req_rx) = mpsc::channel(); + let headers = "Host: 127.0.0.1\r\n".to_string(); + let ping = r#"{"jsonrpc":"2.0","id":2,"method":"ping"}"#; + let response = drive(&request("/mcp", &headers, ping), &req_tx); + + assert!(response.starts_with("HTTP/1.1 200 OK"), "{response}"); + assert!(response.contains(r#""mode":"live""#), "{response}"); + assert!(response.contains(TOKEN), "{response}"); + assert!( + req_rx.try_recv().is_err(), + "ping never touches the UI thread" + ); + + // …but a ping from a foreign origin is still refused, so a web page + // cannot use the probe to harvest the token. + let headers = format!("Host: 127.0.0.1:{PORT}\r\nOrigin: http://evil.example\r\n"); + let response = drive(&request("/mcp", &headers, ping), &req_tx); + assert!(response.starts_with("HTTP/1.1 403 Forbidden"), "{response}"); + assert!(!response.contains(TOKEN), "{response}"); +} + +/// Whole-document sync replaces the LIVE (possibly shared) document — the +/// REST twin of the JSON-RPC write path, and equally unauthenticated +/// before this gate. +#[test] +fn document_sync_route_requires_the_instance_token() { + let (req_tx, req_rx) = mpsc::channel(); + let headers = format!("Host: 127.0.0.1:{PORT}\r\n"); + let body = r#"{"document":{"version":"1.0","children":[],"pages":[]}}"#; + let response = drive(&request("/api/mcp/document", &headers, body), &req_tx); + + assert!( + response.starts_with("HTTP/1.1 401 Unauthorized"), + "{response}" + ); + assert!(response.contains(r#""ok":false"#), "{response}"); + assert!( + req_rx.try_recv().is_err(), + "an unauthenticated whole-document sync must never reach the UI thread" + ); +} diff --git a/crates/op-host-services/src/mcp_live/connection.rs b/crates/op-host-services/src/mcp_live/connection.rs index 8f66dfe14..bd71a5878 100644 --- a/crates/op-host-services/src/mcp_live/connection.rs +++ b/crates/op-host-services/src/mcp_live/connection.rs @@ -20,7 +20,7 @@ pub(super) fn server_loop( listener: TcpListener, req_tx: Sender, stop_rx: Receiver<()>, - token: String, + admission: Arc, quit_flag: Arc, wake_ui: UiWake, client_identity: Arc>>, @@ -57,7 +57,7 @@ pub(super) fn server_loop( // pings. Threads are short-lived and detached. conn_count.fetch_add(1, Ordering::AcqRel); let req_tx = req_tx.clone(); - let token = token.clone(); + let admission = Arc::clone(&admission); let lock = Arc::clone(&stateful_lock); let quit = Arc::clone(&quit_flag); let conns = Arc::clone(&conn_count); @@ -78,7 +78,7 @@ pub(super) fn server_loop( if let Err(e) = serve_connection( &mut stream, &req_tx, - &token, + &admission, &lock, &quit, &wake, @@ -113,13 +113,25 @@ pub(super) fn server_loop( pub(super) fn serve_connection( stream: &mut S, req_tx: &Sender, - token: &str, + admission: &LiveAdmission, stateful_lock: &Mutex<()>, quit_flag: &AtomicBool, wake_ui: &UiWake, client_identity: &Mutex>, ) -> Result<(), McpLiveError> { let req = crate::mcp_serve::read_http_request(stream)?; + let token = admission.token(); + // Gate 1 (see `admission.rs`): browser screening, ahead of ALL routing — + // including the preflight and the stateless probes, so a foreign page + // cannot even fingerprint this endpoint. `Host`/`Origin` are not + // page-forgeable, which is what closes DNS rebinding. + if let Err(denial) = admission::check_boundary(&req, admission) { + return write_http( + stream, + denial.http_status(), + &admission::denial_json_rpc(&req.body, denial), + ); + } if req.method == "OPTIONS" { return write_http(stream, "204 No Content", ""); } @@ -128,6 +140,17 @@ pub(super) fn serve_connection( // client (`setSyncDocument` → POST `{document}`) drive THIS editor's // on-screen canvas, mirroring `apps/web/server/api/mcp/document.post.ts`. if crate::mcp_serve::is_document_sync_route(&req.method, &req.path) { + // Whole-document replacement of the live (possibly SHARED) document + // — the single most destructive thing this endpoint can do, so it + // takes the token gate before the body is even parsed. The REST + // route answers `{ok,error}`, not JSON-RPC. + if let Err(denial) = admission::check_token(&req, admission) { + return write_http( + stream, + denial.http_status(), + &admission::denial_rest(denial), + ); + } return serve_document_sync(stream, req_tx, wake_ui, stateful_lock, &req.body); } if req.path != "/mcp" && req.path != "/" { @@ -175,6 +198,20 @@ pub(super) fn serve_connection( } crate::mcp_serve::Stateless::NeedsState => {} } + // Gate 2 (see `admission.rs`): everything from here down reads or writes + // the live document — every `tools/call`, not just the write ones — so it + // requires the per-instance token. This sits IN FRONT of + // `CollabGatePolicy` (which still runs on the UI thread for each apply) + // and does not replace it: the policy decides what a session permits, this + // decides who is allowed to ask at all. Refusals are a JSON-RPC error + // carrying the caller's id, never a silent pass. + if let Err(denial) = admission::check_token(&req, admission) { + return write_http( + stream, + denial.http_status(), + &admission::denial_json_rpc(&req.body, denial), + ); + } // Everything below observes or mutates shared state — the live // `EditorState` OR a `--file` document on disk (a read-modify-write). // Serialize it under one lock so snapshots stay coherent, concurrent live diff --git a/crates/op-rpc-transport/src/lib.rs b/crates/op-rpc-transport/src/lib.rs index d6205b17b..507c3c07d 100644 --- a/crates/op-rpc-transport/src/lib.rs +++ b/crates/op-rpc-transport/src/lib.rs @@ -188,6 +188,10 @@ impl From for String { pub struct TcpJsonRpc { addr: SocketAddr, path: String, + /// Per-instance token for the live MCP endpoint, sent as + /// `X-OpenPencil-Token`. Empty means "send no header", which is what the + /// endpoint's tokenless probes (`initialize`, `ping`) still accept. + token: String, } impl TcpJsonRpc { @@ -195,6 +199,7 @@ impl TcpJsonRpc { Self { addr: SocketAddr::from(([127, 0, 0, 1], port)), path: path.to_string(), + token: String::new(), } } @@ -202,11 +207,27 @@ impl TcpJsonRpc { Self::localhost(port, MCP_PATH) } + /// Authenticates subsequent requests with the endpoint's instance token. + /// + /// The live MCP endpoint authenticates every stateful call, so a client + /// that omits this gets `401` on anything beyond the discovery probes. + #[must_use] + pub fn with_token(mut self, token: &str) -> Self { + self.token = token.to_string(); + self + } + pub fn http_post_request(&self, body: &str) -> String { + let authorization = if self.token.is_empty() { + String::new() + } else { + format!("X-OpenPencil-Token: {}\r\n", self.token) + }; format!( "POST {} HTTP/1.1\r\nHost: 127.0.0.1\r\nContent-Type: application/json\r\n\ - Content-Length: {}\r\nConnection: close\r\n\r\n{}", + {}Content-Length: {}\r\nConnection: close\r\n\r\n{}", self.path, + authorization, body.len(), body ) diff --git a/docs/security/p2p-collaboration-threat-model.md b/docs/security/p2p-collaboration-threat-model.md index 6d244cc28..ed8175489 100644 --- a/docs/security/p2p-collaboration-threat-model.md +++ b/docs/security/p2p-collaboration-threat-model.md @@ -301,8 +301,20 @@ than "routing metadata" in the narrow sense, and is stated here explicitly: - **The full admission ticket**, presented as the WSS `Authorization: Bearer` credential and verified by the relay. Its claims carry the global account subject, the device id, and — when present — the display name and avatar - URL. A relay operator can therefore reconstruct which accounts collaborate - with each other, from which devices, and when. + URL. Scope this precisely: peer admission today requires the remote ticket's + subject to equal the *local* account (`expected_subject` is the local + account on both sides, and a mismatch is rejected as `WrongSubject`), so the + product currently pairs only devices of the same account. What a relay + operator reconstructs is therefore **which devices of a given account sync, + from where, and when** — not a cross-account collaboration graph. That graph + becomes possible the moment cross-account collaboration ships, which is why + the credential is worth minimizing before then. + + Note also that the relay reads exactly one field out of this ticket, the + expiry it clamps the session deadline to. The subject, device id, `jti`, + display name, and avatar URL are format-checked and then discarded — they + are never compared, stored, or returned by the authorization path. The + disclosure is therefore gratuitous rather than load-bearing. - **The client hello**: route id, role, caller device X25519 public key, possession proof, and the embedded signed locator (owner Noise static key, region, validity window, discovery id). @@ -314,10 +326,36 @@ than "routing metadata" in the narrow sense, and is stated here explicitly: Residual risk: this is a disclosure to the relay operator, not to the network, and it is inherent to operating a relay that authenticates before forwarding. -It is accepted for the first-party relay. Two mitigations are open work: moving -to a claim-minimized relay credential that proves route authorization without -carrying account identity, and replacing the cleartext `session_id` in the -prologue with a binding commitment. Until then, the deployment requirement is +It is accepted for the first-party relay. Two mitigations were investigated. A claim-minimized relay credential that +carries no account identity is open work, and is designed: the relay's +authorization output is `(route, role, expiry)`, none of which comes from the +ticket's identity claims, so the credential can be reduced to an +audience-scoped token carrying only issuer, audience, version, scope, the +channel binding to the caller's X25519 key, and the time bounds. Two honest +limits on what that buys. First, it de-identifies but does not make the relay's +view unlinkable: the device's X25519 static key is persistent and travels in +cleartext in every hello, so the operator still builds a permanent *device* +graph — it simply can no longer name the nodes or join them to the account +namespace used by support, billing, and every other first-party service. +Second, an operator who runs both the issuer and the relay can re-identify a +connection by joining on the channel-binding key and the issuance timestamp, so +against the first party the change is close to symbolic; its real value is +against relay compromise, log leakage, a third-party or regional relay +operator, and relay-scoped lawful-access requests. Implementation is blocked on +the private provider ABI, which mints the credential. +Replacing the cleartext `session_id` in the prologue with a binding commitment +was prototyped and **rejected**: the prologue is not merely a binding, it is the +only channel by which a guest learns the session id at all. The LAN join path +reads it straight out of the prelude, the mDNS record deliberately carries no +session identifier, and the relay invite carries only the locator and route +capability. A commitment would therefore require publishing the session id +somewhere a guest can reach before connecting — on LAN that means the mDNS +record, which trades a handle visible to one relay operator for a handle +broadcast to the whole local network. The marginal gain was also narrower than +it first appeared: the relay already holds the `discovery_id` inside the signed +locator for the whole session, so only cross-epoch correlation would have been +closed. Closing this properly needs the session id added to the invite and a +different LAN join handshake — a product change, not a protocol tweak. Until then, the deployment requirement is that relay and locator operators must not persist bearer headers, hello bytes, or client addresses beyond what an incident needs, and must not log them at all: no relay or locator code path interpolates ticket, key, or identity