From 529b65729bbcce6e6381669047b011458758f704 Mon Sep 17 00:00:00 2001 From: Shiju Date: Wed, 23 Sep 2026 16:03:24 +0530 Subject: [PATCH 1/2] fix(network): refuse protocol upgrades on JSON-RPC and MCP endpoints JSON-RPC and MCP rules apply to each HTTP request, but the proxy could forward a request that also carried upgrade headers. After an upstream answered 101, route selection and the forward proxy relayed the connection without inspection. Refuse any request that carries an Upgrade header on JSON-RPC-family endpoints before the L7 policy decision, in every enforcement mode. Share the check with the existing h2c refusal and call it from relay_jsonrpc as well. Record the refusal as a policy denial and answer with the unsupported_l7_protocol error, because no policy rule can allow the request. If a JSON-RPC-family endpoint still receives 101, close the connection instead of relaying raw bytes. Document the refusal and the WebSocket alternative. Signed-off-by: Shiju --- architecture/security-policy.md | 5 +- .../src/l7/relay.rs | 456 +++++++++++++++++- .../src/l7/rest.rs | 169 +++++++ .../openshell-supervisor-network/src/proxy.rs | 159 +++++- .../how-it-works/policies/manage-policies.mdx | 1 + docs/how-it-works/policies/network-rules.mdx | 10 + docs/how-it-works/policies/schema.mdx | 2 +- docs/observability/logging.mdx | 2 +- 8 files changed, 775 insertions(+), 29 deletions(-) diff --git a/architecture/security-policy.md b/architecture/security-policy.md index 36a919573a..29fa93b842 100644 --- a/architecture/security-policy.md +++ b/architecture/security-policy.md @@ -113,7 +113,10 @@ allowing them to continue under stale authorization. HTTP upgrades switch to raw relay by default. A `protocol: rest` endpoint can opt in to `websocket_credential_rewrite` for client-to-server WebSocket text messages after an allowed `101` upgrade; server-to-client traffic and all other upgraded -protocols remain raw passthrough. +protocols remain raw passthrough. JSON-RPC and MCP endpoints refuse every +request that carries an `Upgrade` header with `403` before forwarding, +whatever the enforcement mode, because their rules apply to individual HTTP +requests and a raw relay would bypass them. A `protocol: tcp` hostname is a connection-routing constraint, not an application-authority boundary. Transparent capture validates the approved DNS diff --git a/crates/openshell-supervisor-network/src/l7/relay.rs b/crates/openshell-supervisor-network/src/l7/relay.rs index 574ad3be8f..e08b3cd681 100644 --- a/crates/openshell-supervisor-network/src/l7/relay.rs +++ b/crates/openshell-supervisor-network/src/l7/relay.rs @@ -707,34 +707,45 @@ fn engine_type_for_protocol(protocol: L7Protocol) -> &'static str { } } -async fn deny_h2c_upgrade_if_requested( +/// Refuses an upgrade the endpoint cannot inspect and reports whether the +/// request was answered. +/// +/// Every L7 request loop calls this before the L7 policy decision. A refusal +/// records a policy denial for the endpoint, emits a parse-rejection event, +/// and answers `403` regardless of enforcement mode; see +/// `unsupported_upgrade_detail` for which upgrades each protocol refuses. The +/// response carries the `unsupported_l7_protocol` error rather than the +/// policy-denial body, because no policy rule can allow the request. +async fn deny_unsupported_upgrade_if_requested( req: &crate::l7::provider::L7Request, config: &L7EndpointConfig, ctx: &L7EvalContext, + observer: Option<&EndpointObserver>, client: &mut C, ) -> Result where C: AsyncRead + AsyncWrite + Unpin + Send, { - if !crate::l7::rest::request_is_h2c_upgrade(&req.raw_header) { + let Some(detail) = + crate::l7::rest::unsupported_upgrade_detail(&req.raw_header, config.protocol) + else { return Ok(false); - } + }; - emit_parse_rejection( - ctx, - crate::l7::rest::UNSUPPORTED_H2C_UPGRADE_DETAIL, - engine_type_for_protocol(config.protocol), - ); - crate::l7::rest::RestProvider::default() - .deny_with_redacted_target( - req, - &ctx.policy_name, - crate::l7::rest::UNSUPPORTED_H2C_UPGRADE_DETAIL, - client, - None, - Some(crate::l7::rest::DenyResponseContext::from_l7_context(ctx)), - ) - .await?; + if let Some(observer) = observer { + observer.observe(EndpointResult::PolicyDenied); + } + emit_parse_rejection(ctx, detail, engine_type_for_protocol(config.protocol)); + crate::l7::rest::send_json_response( + &ctx.policy_name, + serde_json::json!({ + "error": "unsupported_l7_protocol", + "detail": detail, + }), + client, + "403 Forbidden", + ) + .await?; Ok(true) } @@ -934,7 +945,9 @@ where .await?; return Ok(()); } - if deny_h2c_upgrade_if_requested(&req, config, ctx, client).await? { + if deny_unsupported_upgrade_if_requested(&req, config, ctx, observer.as_ref(), client) + .await? + { return Ok(()); } @@ -1283,6 +1296,26 @@ where websocket_permessage_deflate, websocket_subprotocol, } => { + // JSON-RPC and MCP rules apply to individual HTTP requests. + // No current path forwards upgrade headers for these + // protocols: the request-side refusal rejects them, and + // request middleware cannot add upgrade or connection + // headers. If a later change lets such a request reach an + // upstream that answers `101`, close instead of relaying + // frames that no rule would inspect. + if config.protocol.is_jsonrpc_family() { + warn!( + host = %ctx.host, + port = ctx.port, + "closing JSON-RPC connection after unexpected protocol upgrade" + ); + if let Some(session) = middleware_session.take() { + session + .end(openshell_core::proto::MiddlewareSessionEndReason::ProtocolError) + .await; + } + return Ok(()); + } let mut options = upgrade_options( config, ctx, @@ -1751,7 +1784,7 @@ where reject_request_authority_mismatch(client, ctx, &req.action).await?; return Ok(()); } - if deny_h2c_upgrade_if_requested(&req, config, ctx, client).await? { + if deny_unsupported_upgrade_if_requested(&req, config, ctx, None, client).await? { return Ok(()); } @@ -2218,6 +2251,11 @@ where reject_request_authority_mismatch(client, ctx, &req.action).await?; return Ok(()); } + if deny_unsupported_upgrade_if_requested(&req, config, ctx, observer.as_ref(), client) + .await? + { + return Ok(()); + } if close_if_stale(engine.generation_guard(), ctx) { return Ok(()); } @@ -2510,7 +2548,7 @@ where reject_request_authority_mismatch(client, ctx, &req.action).await?; return Ok(()); } - if deny_h2c_upgrade_if_requested(&req, config, ctx, client).await? { + if deny_unsupported_upgrade_if_requested(&req, config, ctx, None, client).await? { return Ok(()); } @@ -9028,6 +9066,382 @@ network_policies: let _ = tokio::time::timeout(std::time::Duration::from_secs(1), relay).await; } + /// A `tools/call` that no JSON-RPC-family fixture below allows. + const UNALLOWED_TOOL_CALL: &[u8] = + br#"{"jsonrpc":"2.0","id":2,"method":"tools/call","params":{"name":"delete_resource","arguments":{}}}"#; + + /// A receive-stream GET that also asks to switch to WebSocket. + const JSONRPC_WEBSOCKET_UPGRADE_REQUEST: &[u8] = b"GET /mcp HTTP/1.1\r\nHost: mcp.example.test:8000\r\nAccept: text/event-stream\r\nMCP-Protocol-Version: 2025-11-25\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\nSec-WebSocket-Version: 13\r\n\r\n"; + + /// Builds two endpoints on one host and port so every request goes + /// through per-request route selection: a JSON-RPC-family endpoint at + /// `/mcp` that allows only `initialize`, and a REST endpoint at `/api/**`. + fn jsonrpc_and_rest_route_configs( + protocol: &str, + enforcement: &str, + ) -> (Vec, TunnelPolicyEngine, L7EvalContext) { + let data = format!( + r#" +network_policies: + shared_api: + name: shared_api + endpoints: + - host: mcp.example.test + port: 8000 + path: "/mcp" + protocol: {protocol} + enforcement: {enforcement} + rules: + - allow: + method: initialize + - host: mcp.example.test + port: 8000 + path: "/api/**" + protocol: rest + enforcement: enforce + rules: + - allow: + method: GET + path: "/api/**" + binaries: + - {{ path: /usr/bin/python3 }} +"# + ); + let engine = OpaEngine::from_strings(TEST_POLICY, &data).unwrap(); + let input = NetworkInput { + host: "mcp.example.test".into(), + port: 8000, + binary_path: PathBuf::from("/usr/bin/python3"), + binary_sha256: "unused".into(), + ancestors: vec![], + cmdline_paths: vec![], + }; + let (endpoint_configs, generation) = engine + .query_endpoint_configs_with_generation(&input) + .unwrap(); + let configs: Vec = endpoint_configs + .iter() + .map(|config| crate::l7::parse_l7_config(config).unwrap()) + .collect(); + assert_eq!(configs.len(), 2, "both endpoints must share the route"); + let tunnel_engine = engine.clone_engine_for_tunnel(generation).unwrap(); + let ctx = L7EvalContext { + host: "mcp.example.test".into(), + port: 8000, + request_default_port: Some(8000), + policy_name: "shared_api".into(), + binary_path: "/usr/bin/python3".into(), + ancestors: vec![], + cmdline_paths: vec![], + secret_resolver: None, + ..Default::default() + }; + (configs, tunnel_engine, ctx) + } + + /// Result of sending one upgrade request through a relay whose upstream + /// accepts every upgrade it receives. + struct UpgradeScenario { + /// The response head the client received, or empty if none arrived. + response: String, + /// The response body of a non-`101` response. + body: String, + /// Every byte the upstream received. + upstream_seen: Vec, + } + + /// Sends `request`, answers any forwarded upgrade with a valid `101`, and + /// after a `101` writes `frame` as a WebSocket text message. The upstream + /// never refuses, so a relay that forwards the upgrade and then copies + /// bytes delivers `frame` to it. + async fn run_upgrade_scenario(request: &[u8], frame: &[u8], relay: F) -> UpgradeScenario + where + F: FnOnce( + tokio::io::DuplexStream, + tokio::io::DuplexStream, + ) -> tokio::task::JoinHandle>, + { + let (mut app, relay_client) = tokio::io::duplex(8192); + let (relay_upstream, mut upstream) = tokio::io::duplex(8192); + let relay = relay(relay_client, relay_upstream); + let upstream_task = tokio::spawn(async move { + let mut seen = Vec::new(); + let mut buf = [0u8; 4096]; + let mut answered = false; + loop { + let read = tokio::time::timeout( + std::time::Duration::from_secs(2), + upstream.read(&mut buf), + ) + .await; + let Ok(Ok(n)) = read else { break }; + if n == 0 { + break; + } + seen.extend_from_slice(&buf[..n]); + if !answered && seen.windows(4).any(|w| w == b"\r\n\r\n") { + answered = true; + let _ = upstream + .write_all( + b"HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=\r\n\r\n", + ) + .await; + } + } + seen + }); + + app.write_all(request).await.unwrap(); + let mut response = Vec::new(); + let mut byte = [0u8; 1]; + while !response.ends_with(b"\r\n\r\n") { + let read = + tokio::time::timeout(std::time::Duration::from_secs(2), app.read(&mut byte)).await; + let Ok(Ok(1)) = read else { break }; + response.push(byte[0]); + } + let response = String::from_utf8_lossy(&response).into_owned(); + let mut body = Vec::new(); + if response.starts_with("HTTP/1.1 101") { + let _ = app.write_all(&masked_text_frame(frame)).await; + } else { + // Refusals close the connection after the body. + let _ = tokio::time::timeout( + std::time::Duration::from_secs(2), + app.read_to_end(&mut body), + ) + .await; + } + drop(app); + let _ = tokio::time::timeout(std::time::Duration::from_secs(2), relay).await; + let upstream_seen = upstream_task.await.unwrap(); + UpgradeScenario { + response, + body: String::from_utf8_lossy(&body).into_owned(), + upstream_seen, + } + } + + fn contains_bytes(haystack: &[u8], needle: &[u8]) -> bool { + haystack + .windows(needle.len()) + .any(|window| window == needle) + } + + fn assert_upgrade_denied_before_forwarding(scenario: &UpgradeScenario) { + assert!( + !contains_bytes( + &scenario.upstream_seen, + &masked_text_frame(UNALLOWED_TOOL_CALL) + ), + "an uninspected tools/call frame reached the upstream" + ); + assert!( + scenario.upstream_seen.is_empty(), + "the upgrade request must not reach the upstream, got: {}", + String::from_utf8_lossy(&scenario.upstream_seen) + ); + assert!( + scenario.response.starts_with("HTTP/1.1 403"), + "expected a 403 denial, got: {}", + scenario.response + ); + assert!( + scenario.body.contains("\"unsupported_l7_protocol\"") + && scenario + .body + .contains(crate::l7::rest::UNSUPPORTED_JSONRPC_UPGRADE_DETAIL), + "expected the upgrade refusal, got: {}", + scenario.body + ); + } + + #[tokio::test] + async fn mcp_websocket_upgrade_refusal_records_policy_denied() { + use openshell_core::endpoint_status::EndpointStatusCommand; + + for route_selected in [false, true] { + let (mut config, tunnel_engine, mut ctx) = mcp_test_relay_context(); + let mut receiver = install_mcp_test_observation(&mut config, &mut ctx).await; + let scenario = run_upgrade_scenario( + JSONRPC_WEBSOCKET_UPGRADE_REQUEST, + UNALLOWED_TOOL_CALL, + move |mut client, mut upstream| { + tokio::spawn(async move { + if route_selected { + relay_with_route_selection( + &[config], + tunnel_engine, + &mut client, + &mut upstream, + &ctx, + ) + .await + } else { + relay_with_inspection( + &config, + tunnel_engine, + &mut client, + &mut upstream, + &ctx, + ) + .await + } + }) + }, + ) + .await; + assert_upgrade_denied_before_forwarding(&scenario); + assert!( + matches!( + receiver.try_recv(), + Ok(EndpointStatusCommand::Observe { + result: EndpointResult::PolicyDenied, + .. + }) + ), + "route_selected={route_selected}: refusal must record a policy denial" + ); + assert!( + receiver.try_recv().is_err(), + "route_selected={route_selected}: one result per exchange" + ); + } + } + + #[tokio::test] + async fn route_selected_mcp_websocket_upgrade_is_denied_before_forwarding() { + let (configs, tunnel_engine, ctx) = jsonrpc_and_rest_route_configs("mcp", "enforce"); + let scenario = run_upgrade_scenario( + JSONRPC_WEBSOCKET_UPGRADE_REQUEST, + UNALLOWED_TOOL_CALL, + move |mut client, mut upstream| { + tokio::spawn(async move { + relay_with_route_selection( + &configs, + tunnel_engine, + &mut client, + &mut upstream, + &ctx, + ) + .await + }) + }, + ) + .await; + assert_upgrade_denied_before_forwarding(&scenario); + } + + #[tokio::test] + async fn route_selected_audit_jsonrpc_websocket_upgrade_is_denied_before_forwarding() { + // Audit mode forwards requests that policy would deny, so the upgrade + // refusal must not depend on the policy decision. + let (configs, tunnel_engine, ctx) = jsonrpc_and_rest_route_configs("json-rpc", "audit"); + let scenario = run_upgrade_scenario( + JSONRPC_WEBSOCKET_UPGRADE_REQUEST, + UNALLOWED_TOOL_CALL, + move |mut client, mut upstream| { + tokio::spawn(async move { + relay_with_route_selection( + &configs, + tunnel_engine, + &mut client, + &mut upstream, + &ctx, + ) + .await + }) + }, + ) + .await; + assert_upgrade_denied_before_forwarding(&scenario); + } + + #[tokio::test] + async fn single_endpoint_mcp_websocket_upgrade_is_denied_before_forwarding() { + let (config, tunnel_engine, ctx) = mcp_test_relay_context(); + let scenario = run_upgrade_scenario( + JSONRPC_WEBSOCKET_UPGRADE_REQUEST, + UNALLOWED_TOOL_CALL, + move |mut client, mut upstream| { + tokio::spawn(async move { + relay_with_inspection(&config, tunnel_engine, &mut client, &mut upstream, &ctx) + .await + }) + }, + ) + .await; + assert_upgrade_denied_before_forwarding(&scenario); + } + + #[tokio::test] + async fn route_selected_rest_websocket_upgrade_still_relays_beside_mcp() { + // The refusal is scoped to JSON-RPC-family endpoints: a REST upgrade + // on the same host and port keeps its documented raw relay. + let (configs, tunnel_engine, ctx) = jsonrpc_and_rest_route_configs("mcp", "enforce"); + let frame = br#"{"type":"ping"}"#; + let scenario = run_upgrade_scenario( + b"GET /api/ws HTTP/1.1\r\nHost: mcp.example.test:8000\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\nSec-WebSocket-Version: 13\r\n\r\n", + frame, + move |mut client, mut upstream| { + tokio::spawn(async move { + relay_with_route_selection( + &configs, + tunnel_engine, + &mut client, + &mut upstream, + &ctx, + ) + .await + }) + }, + ) + .await; + assert!( + scenario.response.starts_with("HTTP/1.1 101"), + "REST upgrade should still switch protocols, got: {}", + scenario.response + ); + assert!(contains_bytes( + &scenario.upstream_seen, + &masked_text_frame(frame) + )); + } + + #[tokio::test] + async fn route_selected_mcp_receive_stream_without_upgrade_is_still_forwarded() { + let (configs, tunnel_engine, ctx) = jsonrpc_and_rest_route_configs("mcp", "enforce"); + let (mut app, mut relay_client) = tokio::io::duplex(8192); + let (mut relay_upstream, mut upstream) = tokio::io::duplex(8192); + let relay = tokio::spawn(async move { + relay_with_route_selection( + &configs, + tunnel_engine, + &mut relay_client, + &mut relay_upstream, + &ctx, + ) + .await + }); + + app.write_all( + b"GET /mcp HTTP/1.1\r\nHost: mcp.example.test:8000\r\nAccept: text/event-stream\r\nMCP-Protocol-Version: 2025-11-25\r\n\r\n", + ) + .await + .unwrap(); + let forwarded = tokio::time::timeout( + std::time::Duration::from_secs(2), + read_http_headers(&mut upstream), + ) + .await + .expect("receive-stream GET should reach the upstream"); + let forwarded = String::from_utf8_lossy(&forwarded); + assert!(forwarded.starts_with("GET /mcp HTTP/1.1\r\n")); + assert!(!forwarded.to_ascii_lowercase().contains("upgrade")); + relay.abort(); + let _ = relay.await; + } + fn masked_text_frame(payload: &[u8]) -> Vec { let mask = [0x11, 0x22, 0x33, 0x44]; assert!( diff --git a/crates/openshell-supervisor-network/src/l7/rest.rs b/crates/openshell-supervisor-network/src/l7/rest.rs index 4592600231..b7818142b9 100644 --- a/crates/openshell-supervisor-network/src/l7/rest.rs +++ b/crates/openshell-supervisor-network/src/l7/rest.rs @@ -84,6 +84,8 @@ const HTTP_METHOD_PREFIXES: &[&[u8]] = &[ pub(crate) const HTTP2_PRIOR_KNOWLEDGE_PREFACE: &[u8] = b"PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n"; pub(crate) const UNSUPPORTED_H2C_UPGRADE_DETAIL: &str = "HTTP/2 cleartext upgrade (h2c) is not supported for L7-inspected endpoints"; +pub(crate) const UNSUPPORTED_JSONRPC_UPGRADE_DETAIL: &str = + "HTTP upgrade is not supported for JSON-RPC or MCP endpoints"; const MIN_HTTP2_PREFACE_DETECTION_BYTES: usize = 8; /// Idle timeout for `relay_until_eof`. If no data arrives within this window @@ -2393,6 +2395,52 @@ pub(crate) fn request_is_h2c_upgrade(raw_header: &[u8]) -> bool { upgrade_h2c && connection_upgrade } +/// Returns why an L7 endpoint using `protocol` must refuse this request's +/// upgrade, or `None` when the request may continue. +/// +/// Every inspected protocol refuses h2c. JSON-RPC and MCP policy applies to +/// individual HTTP requests, so after any protocol switch no rule would see +/// the messages; those endpoints refuse every request that carries an +/// `Upgrade` header. Callers apply this before the L7 policy decision and +/// regardless of enforcement mode, because an upgrade would end inspection +/// rather than break a rule that audit mode could log. +pub(crate) fn unsupported_upgrade_detail( + raw_header: &[u8], + protocol: crate::l7::L7Protocol, +) -> Option<&'static str> { + if request_is_h2c_upgrade(raw_header) { + return Some(UNSUPPORTED_H2C_UPGRADE_DETAIL); + } + if protocol.is_jsonrpc_family() && request_has_upgrade_header(raw_header) { + return Some(UNSUPPORTED_JSONRPC_UPGRADE_DETAIL); + } + None +} + +/// Returns true when a request carries an `Upgrade` header, whatever its +/// value. +/// +/// Both relay checks that can lead to a protocol switch require that header: +/// `request_is_websocket_upgrade`, which decides whether upgrade headers are +/// forwarded, and `client_requested_upgrade`, which decides whether an +/// upstream `101` may reach the client. A refusal based on this test therefore +/// covers every request either one treats as an upgrade. Headers that are not +/// UTF-8 return false; the shared relay rejects them before forwarding. +fn request_has_upgrade_header(raw_header: &[u8]) -> bool { + let header_end = raw_header + .windows(4) + .position(|w| w == b"\r\n\r\n") + .map_or(raw_header.len(), |p| p + 4); + let Ok(header_str) = std::str::from_utf8(&raw_header[..header_end]) else { + return false; + }; + + header_str.lines().skip(1).any(|line| { + line.split_once(':') + .is_some_and(|(name, _)| name.trim().eq_ignore_ascii_case("upgrade")) + }) +} + fn rewrite_websocket_extensions_for_mode( raw_header: &[u8], mode: WebSocketExtensionMode, @@ -8897,6 +8945,127 @@ mod tests { assert!(!client_requested_upgrade(headers)); } + #[test] + fn unsupported_upgrade_detail_refuses_h2c_for_every_protocol() { + let raw = b"GET /api HTTP/1.1\r\nHost: example.com\r\nConnection: Upgrade, HTTP2-Settings\r\nUpgrade: h2c\r\nHTTP2-Settings: AAMAAABkAAQAAP__\r\n\r\n"; + for protocol in [ + crate::l7::L7Protocol::Rest, + crate::l7::L7Protocol::Websocket, + crate::l7::L7Protocol::Graphql, + crate::l7::L7Protocol::JsonRpc, + crate::l7::L7Protocol::Mcp, + ] { + assert_eq!( + unsupported_upgrade_detail(raw, protocol), + Some(UNSUPPORTED_H2C_UPGRADE_DETAIL), + "{protocol:?}" + ); + } + } + + #[test] + fn unsupported_upgrade_detail_refuses_any_upgrade_on_jsonrpc_family() { + let websocket = format!( + "GET /mcp HTTP/1.1\r\nHost: example.com\r\nAccept: text/event-stream\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Key: {VALID_WS_KEY}\r\nSec-WebSocket-Version: 13\r\n\r\n" + ); + // Any request that carries an `Upgrade` header is refused, including + // one without `Connection: upgrade` that the relay would not upgrade. + let requests: [&[u8]; 3] = [ + websocket.as_bytes(), + b"POST /mcp HTTP/1.1\r\nHost: example.com\r\nUpgrade: websocket\r\nConnection: keep-alive, upgrade\r\nContent-Length: 0\r\n\r\n", + b"GET /mcp HTTP/1.1\r\nHost: example.com\r\nUpgrade: custom\r\n\r\n", + ]; + for raw in requests { + for protocol in [crate::l7::L7Protocol::JsonRpc, crate::l7::L7Protocol::Mcp] { + assert_eq!( + unsupported_upgrade_detail(raw, protocol), + Some(UNSUPPORTED_JSONRPC_UPGRADE_DETAIL), + "{protocol:?}: {}", + String::from_utf8_lossy(raw) + ); + } + } + } + + #[test] + fn unsupported_upgrade_detail_allows_ordinary_and_non_jsonrpc_requests() { + let websocket = format!( + "GET /ws HTTP/1.1\r\nHost: example.com\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Key: {VALID_WS_KEY}\r\nSec-WebSocket-Version: 13\r\n\r\n" + ); + for protocol in [ + crate::l7::L7Protocol::Rest, + crate::l7::L7Protocol::Websocket, + crate::l7::L7Protocol::Graphql, + ] { + assert_eq!( + unsupported_upgrade_detail(websocket.as_bytes(), protocol), + None, + "{protocol:?}" + ); + } + + // Streamable HTTP requests, a look-alike header name, and a stray + // `Connection: upgrade` without an `Upgrade` header stay allowed; + // none of them can switch protocols. + let ordinary: [&[u8]; 4] = [ + b"POST /mcp HTTP/1.1\r\nHost: example.com\r\nContent-Type: application/json\r\nConnection: keep-alive\r\nContent-Length: 2\r\n\r\n{}", + b"GET /mcp HTTP/1.1\r\nHost: example.com\r\nAccept: text/event-stream\r\n\r\n", + b"GET /mcp HTTP/1.1\r\nHost: example.com\r\nUpgrade-Insecure-Requests: 1\r\n\r\n", + b"GET /mcp HTTP/1.1\r\nHost: example.com\r\nConnection: Upgrade\r\n\r\n", + ]; + for raw in ordinary { + for protocol in [crate::l7::L7Protocol::JsonRpc, crate::l7::L7Protocol::Mcp] { + assert_eq!( + unsupported_upgrade_detail(raw, protocol), + None, + "{protocol:?}: {}", + String::from_utf8_lossy(raw) + ); + } + } + } + + #[test] + fn unsupported_upgrade_detail_covers_every_relay_upgrade_check() { + // The refusal must be at least as broad as both relay checks that can + // lead to a protocol switch, across header spellings that pass + // ingress validation. + let valid_websocket = |head: &str| { + format!("{head}Sec-WebSocket-Key: {VALID_WS_KEY}\r\nSec-WebSocket-Version: 13\r\n\r\n") + }; + let requests = [ + valid_websocket( + "GET /mcp HTTP/1.1\r\nHost: example.com\r\nUPGRADE: WebSocket\r\nCONNECTION: UPGRADE\r\n", + ), + valid_websocket( + "GET /mcp HTTP/1.1\r\nHost: example.com\r\nConnection: keep-alive\r\nConnection: upgrade\r\nUpgrade: websocket\r\n", + ), + "GET /mcp HTTP/1.1\r\nHost: example.com\r\nUpgrade:\r\nConnection: upgrade\r\n\r\n" + .to_string(), + "GET /mcp HTTP/1.0\r\nUpgrade: websocket\r\nConnection:upgrade\r\n\r\n".to_string(), + "POST /mcp HTTP/1.1\r\nHost: example.com\r\nTransfer-Encoding: chunked\r\nUpgrade: custom\r\nConnection: close, Upgrade\r\n\r\n" + .to_string(), + "GET /mcp HTTP/1.1\r\nHost: example.com\r\nUpgrade: websocket\r\nUpgrade: h2c\r\nConnection: upgrade\r\n\r\n" + .to_string(), + ]; + for raw in &requests { + assert!( + validate_http_request_header_block(raw.as_bytes()).is_ok(), + "fixture must pass ingress validation: {raw}" + ); + assert!( + client_requested_upgrade(raw) || request_is_websocket_upgrade(raw.as_bytes()), + "fixture must be an upgrade to the relay: {raw}" + ); + for protocol in [crate::l7::L7Protocol::JsonRpc, crate::l7::L7Protocol::Mcp] { + assert!( + unsupported_upgrade_detail(raw.as_bytes(), protocol).is_some(), + "{protocol:?} must refuse: {raw}" + ); + } + } + } + #[test] fn client_requested_upgrade_handles_comma_separated_connection() { let headers = "GET /ws HTTP/1.1\r\nHost: example.com\r\nUpgrade: websocket\r\nConnection: keep-alive, Upgrade\r\n\r\n"; diff --git a/crates/openshell-supervisor-network/src/proxy.rs b/crates/openshell-supervisor-network/src/proxy.rs index 6ac6628cc4..0a2cc72f7b 100644 --- a/crates/openshell-supervisor-network/src/proxy.rs +++ b/crates/openshell-supervisor-network/src/proxy.rs @@ -5647,7 +5647,15 @@ async fn handle_forward_proxy( .await?; return Ok(()); } - if crate::l7::rest::request_is_h2c_upgrade(&forward_request_bytes) { + // Refuse upgrades this endpoint cannot inspect before the L7 policy + // decision; see `unsupported_upgrade_detail` for the per-protocol rule. + if let Some(upgrade_detail) = crate::l7::rest::unsupported_upgrade_detail( + &forward_request_bytes, + l7_config.config.protocol, + ) { + if let Some(observer) = endpoint_observer.as_ref() { + observer.observe(EndpointResult::PolicyDenied); + } let event = HttpActivityBuilder::new(openshell_ocsf::ctx::ctx()) .activity(ActivityId::Other) .action(ActionId::Denied) @@ -5666,9 +5674,9 @@ async fn handle_forward_proxy( ) .firewall_rule(policy_str, "l7") .message(format!( - "FORWARD_L7 denied unsupported h2c upgrade for {method} {host_lc}:{port}{telemetry_path}" + "FORWARD_L7 denied unsupported upgrade for {method} {host_lc}:{port}{telemetry_path}" )) - .status_detail(crate::l7::rest::UNSUPPORTED_H2C_UPGRADE_DETAIL) + .status_detail(upgrade_detail) .build(); ocsf_emit!(event); emit_activity_simple(activity_tx, true, "l7_parse_rejection"); @@ -5678,7 +5686,7 @@ async fn handle_forward_proxy( port, &binary_str, &decision, - crate::l7::rest::UNSUPPORTED_H2C_UPGRADE_DETAIL, + upgrade_detail, "forward-l7-parse-rejection", ); respond( @@ -5687,7 +5695,7 @@ async fn handle_forward_proxy( 403, "Forbidden", "unsupported_l7_protocol", - crate::l7::rest::UNSUPPORTED_H2C_UPGRADE_DETAIL, + upgrade_detail, ), ) .await?; @@ -6559,6 +6567,28 @@ async fn handle_forward_proxy( websocket_permessage_deflate, websocket_subprotocol, } => { + // JSON-RPC and MCP rules apply to individual HTTP requests. No + // current path forwards upgrade headers for these protocols: the + // request-side refusal rejects them, and request middleware cannot + // add upgrade or connection headers. If a later change lets such a + // request reach an upstream that answers `101`, close instead of + // relaying frames that no rule would inspect. + if forward_upgrade_config + .as_ref() + .is_some_and(|config| config.protocol.is_jsonrpc_family()) + { + warn!( + host = %host_lc, + port, + "closing forwarded JSON-RPC connection after unexpected protocol upgrade" + ); + if let Some(session) = middleware_session.take() { + session + .end(openshell_core::proto::MiddlewareSessionEndReason::ProtocolError) + .await; + } + return Ok(()); + } let mut upgrade_options = if let (Some(config), Some(engine)) = ( forward_upgrade_config.as_ref(), forward_tunnel_engine.as_ref(), @@ -8577,6 +8607,125 @@ network_policies: } } + #[tokio::test] + async fn forward_mcp_websocket_upgrade_is_denied_before_connecting_upstream() { + if !cfg!(target_os = "linux") { + eprintln!("skipping: handler identity binding requires /proc (Linux)"); + return; + } + let Some(upstream_ip) = non_loopback_test_ipv4() else { + eprintln!("skipping: no routable non-loopback IPv4 test address"); + return; + }; + + let upstream_listener = TcpListener::bind((upstream_ip, 0)) + .await + .expect("bind MCP upstream listener"); + let upstream_port = upstream_listener.local_addr().unwrap().port(); + let executable = std::env::current_exe().expect("current executable"); + let data = format!( + r#" +network_policies: + mcp-upstream: + name: mcp-upstream + endpoints: + - host: "{upstream_ip}" + port: {upstream_port} + path: /mcp + protocol: mcp + enforcement: enforce + rules: + - allow: + method: initialize + binaries: + - {{ path: "{executable}" }} +"#, + executable = executable.display(), + ); + let engine = Arc::new( + OpaEngine::from_strings(include_str!("../data/sandbox-policy.rego"), &data) + .expect("load MCP policy"), + ); + + let proxy_listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("bind proxy listener"); + let proxy_address = proxy_listener.local_addr().unwrap(); + let target = format!("http://{upstream_ip}:{upstream_port}/mcp"); + // A receive-stream GET that the MCP policy allows, plus WebSocket + // upgrade headers. + let request = format!( + "GET {target} HTTP/1.1\r\nHost: {upstream_ip}:{upstream_port}\r\nAccept: text/event-stream\r\nMCP-Protocol-Version: 2025-11-25\r\nConnection: Upgrade\r\nUpgrade: websocket\r\nSec-WebSocket-Version: 13\r\nSec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\n\r\n" + ); + let client = tokio::spawn(async move { + let mut socket = TcpStream::connect(proxy_address) + .await + .expect("connect proxy"); + let mut response = Vec::new(); + socket + .read_to_end(&mut response) + .await + .expect("read proxy response"); + response + }); + let (proxy_connection, _) = proxy_listener.accept().await.unwrap(); + let socket_addrs = proxy_connection + .peer_addr() + .ok() + .zip(proxy_connection.local_addr().ok()); + let stream: BoundaryDuplexStream = Box::new(proxy_connection); + let mut proxy_connection = tokio::io::BufReader::new(stream); + + tokio::time::timeout( + std::time::Duration::from_secs(30), + Box::pin(handle_forward_proxy( + "GET", + &target, + request.as_bytes(), + request.len(), + &mut proxy_connection, + None, + socket_addrs, + engine, + Arc::new(BinaryIdentityCache::new()), + Arc::new(AtomicU32::new(std::process::id())), + None, + AgentProposals::default(), + Arc::new(None), + Arc::new(None), + None, + None, + None, + None, + None, + None, + )), + ) + .await + .expect("refused upgrade must complete without an upstream response") + .expect("handle refused MCP WebSocket upgrade"); + drop(proxy_connection); + + let response = String::from_utf8(client.await.unwrap()).expect("UTF-8 response"); + assert!( + response.starts_with("HTTP/1.1 403"), + "the upgrade must be refused: {response}" + ); + assert!( + response.contains(crate::l7::rest::UNSUPPORTED_JSONRPC_UPGRADE_DETAIL), + "the refusal must name the unsupported upgrade: {response}" + ); + assert!( + tokio::time::timeout( + std::time::Duration::from_millis(100), + upstream_listener.accept() + ) + .await + .is_err(), + "a refused upgrade must not establish an upstream connection" + ); + } + #[tokio::test] async fn plaintext_websocket_preflight_denial_does_not_connect_upstream() { if !cfg!(target_os = "linux") { diff --git a/docs/how-it-works/policies/manage-policies.mdx b/docs/how-it-works/policies/manage-policies.mdx index 681f7f0f29..89c7d44ce2 100644 --- a/docs/how-it-works/policies/manage-policies.mdx +++ b/docs/how-it-works/policies/manage-policies.mdx @@ -340,6 +340,7 @@ code in the response to find the cause: | `request_authority_mismatch` | The HTTP request's host or port differs from the connection's destination. | The client's `Host` header, including any non-default port, matches the connection. | | `credential_endpoint_mismatch` | A network rule allowed the request, but the provider credential is not bound to this destination. | The provider's profile endpoints or credential binding. Do not widen the network rule. | | `credential_placeholder_in_request_body` | The request body contains an invalid or revoked credential placeholder. | Remove the stale placeholder or restore the provider. | +| `unsupported_l7_protocol` | The request used a protocol or upgrade that the endpoint cannot inspect, such as h2c or an `Upgrade` header sent to an MCP or JSON-RPC endpoint. No rule can allow it. | Send the request without the upgrade, or allow WebSocket traffic through a separate `protocol: websocket` endpoint. | ### A Change Fails to Load diff --git a/docs/how-it-works/policies/network-rules.mdx b/docs/how-it-works/policies/network-rules.mdx index bcaba59218..ac2a161a4f 100644 --- a/docs/how-it-works/policies/network-rules.mdx +++ b/docs/how-it-works/policies/network-rules.mdx @@ -538,6 +538,16 @@ responses and SSE messages are relayed without MCP policy parsing. Do not put an MCP endpoint on the same host and port as an endpoint that uses a different protocol. `openshell policy update` rejects this combination. +MCP and JSON-RPC endpoints carry only HTTP requests, because their rules apply +to each request. OpenShell answers a request to these endpoints that carries an +`Upgrade` header, such as a WebSocket upgrade, with `403 Forbidden` before it +reaches the server, in both `enforce` and `audit` mode. If the server also +accepts WebSocket connections, allow them with a separate +`protocol: websocket` endpoint, which applies `GET` and `WEBSOCKET_TEXT` rules, +not MCP method or tool rules. For an MCP server, put that endpoint on a +different host or port. A JSON-RPC endpoint can share its host and port with a +WebSocket endpoint that uses a different `path`. + ### Allow Native TCP Use `protocol: tcp` for a client that speaks a protocol other than HTTP, such diff --git a/docs/how-it-works/policies/schema.mdx b/docs/how-it-works/policies/schema.mdx index a4801b0b0a..f829633c69 100644 --- a/docs/how-it-works/policies/schema.mdx +++ b/docs/how-it-works/policies/schema.mdx @@ -152,7 +152,7 @@ the gateway host. | Field | Type | Default | Description | |---|---|---|---| -| `protocol` | string | None | `rest`, `websocket`, `graphql`, `mcp`, or `json-rpc` for request inspection, or `tcp` for a native TCP connection. Refer to [Connection and Request Checks](/how-it-works/policies/network-rules#connection-and-request-checks). | +| `protocol` | string | None | `rest`, `websocket`, `graphql`, `mcp`, or `json-rpc` for request inspection, or `tcp` for a native TCP connection. `mcp` and `json-rpc` endpoints refuse requests that carry an `Upgrade` header with `403`. Refer to [Connection and Request Checks](/how-it-works/policies/network-rules#connection-and-request-checks). | | `tls` | string | Automatic | `skip` relays traffic without terminating TLS, so OpenShell cannot inspect it. Do not use it with a request protocol. | | `enforcement` | string | `audit` | `enforce` blocks requests that break the endpoint's rules. `audit` logs them and allows the request. | | `access` | string | None | Access preset: `read-only`, `read-write`, or `full`. Refer to [Access Presets](#access-presets). | diff --git a/docs/observability/logging.mdx b/docs/observability/logging.mdx index e9b6a03845..5ac30d7438 100644 --- a/docs/observability/logging.mdx +++ b/docs/observability/logging.mdx @@ -228,7 +228,7 @@ An upstream that the proxy cannot reach returns `502 Bad Gateway`: } ``` -The `error` field is a short machine-readable code (`policy_denied`, `middleware_denied`, `middleware_failed`, `ssrf_denied`, `upstream_unreachable`). The `detail` field is a human-readable explanation suitable for display in an agent transcript. The optional `reason` field, when present, provides the specific denial cause from the policy engine (for example, which binary was not allowed or which rule was missing). +The `error` field is a short machine-readable code (`policy_denied`, `middleware_denied`, `middleware_failed`, `ssrf_denied`, `upstream_unreachable`, `unsupported_l7_protocol`). `unsupported_l7_protocol` means the request used a protocol or upgrade that the endpoint cannot inspect, such as h2c or an upgrade on an MCP or JSON-RPC endpoint; no policy rule can allow it. The `detail` field is a human-readable explanation suitable for display in an agent transcript. The optional `reason` field, when present, provides the specific denial cause from the policy engine (for example, which binary was not allowed or which rule was missing). For L7 REST policy denials, the body also includes structured policy fields such as `method`, `path`, `rule_missing`, and `next_steps`. When the policy advisor is enabled, the body also includes `agent_guidance`, a short plain-language instruction telling the agent to read `/etc/openshell/skills/policy_advisor.md`, propose the narrowest rule through `http://policy.local/v1/proposals`, wait for `policy_reloaded: true`, and retry. A middleware denial instead identifies the policy-local config in `middleware` and can include a validated `reason_code`. A fail-closed runtime failure uses `middleware_failed` with platform-owned text. Both middleware responses omit `rule_missing`, `next_steps`, and `agent_guidance` because no policy rule is missing. From 5b9a48234790bca27684c826b840ee7f8c9170a7 Mon Sep 17 00:00:00 2001 From: Shiju Date: Mon, 28 Sep 2026 21:24:28 +0530 Subject: [PATCH 2/2] docs(observability): remove duplicate protocol error definition Keep unsupported_l7_protocol in the response error-code list and retain its explanation in the policy troubleshooting table. Signed-off-by: Shiju --- docs/observability/logging.mdx | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/observability/logging.mdx b/docs/observability/logging.mdx index 5ac30d7438..a277992117 100644 --- a/docs/observability/logging.mdx +++ b/docs/observability/logging.mdx @@ -228,7 +228,7 @@ An upstream that the proxy cannot reach returns `502 Bad Gateway`: } ``` -The `error` field is a short machine-readable code (`policy_denied`, `middleware_denied`, `middleware_failed`, `ssrf_denied`, `upstream_unreachable`, `unsupported_l7_protocol`). `unsupported_l7_protocol` means the request used a protocol or upgrade that the endpoint cannot inspect, such as h2c or an upgrade on an MCP or JSON-RPC endpoint; no policy rule can allow it. The `detail` field is a human-readable explanation suitable for display in an agent transcript. The optional `reason` field, when present, provides the specific denial cause from the policy engine (for example, which binary was not allowed or which rule was missing). +The `error` field is a short machine-readable code (`policy_denied`, `middleware_denied`, `middleware_failed`, `ssrf_denied`, `upstream_unreachable`, `unsupported_l7_protocol`). The `detail` field is a human-readable explanation suitable for display in an agent transcript. The optional `reason` field, when present, provides the specific denial cause from the policy engine (for example, which binary was not allowed or which rule was missing). For L7 REST policy denials, the body also includes structured policy fields such as `method`, `path`, `rule_missing`, and `next_steps`. When the policy advisor is enabled, the body also includes `agent_guidance`, a short plain-language instruction telling the agent to read `/etc/openshell/skills/policy_advisor.md`, propose the narrowest rule through `http://policy.local/v1/proposals`, wait for `policy_reloaded: true`, and retry. A middleware denial instead identifies the policy-local config in `middleware` and can include a validated `reason_code`. A fail-closed runtime failure uses `middleware_failed` with platform-owned text. Both middleware responses omit `rule_missing`, `next_steps`, and `agent_guidance` because no policy rule is missing.