From 362846c95f2aaf6d5824d487453b683b3e94290e Mon Sep 17 00:00:00 2001 From: KC <79471844+wolfyy970@users.noreply.github.com> Date: Sat, 1 Aug 2026 15:27:33 -0400 Subject: [PATCH 1/6] buzz-acp: add structured MCP server configuration Signed-off-by: KC <79471844+wolfyy970@users.noreply.github.com> --- .env.example | 5 + crates/buzz-acp/README.md | 55 ++ crates/buzz-acp/src/acp.rs | 428 +++++++++-- crates/buzz-acp/src/config.rs | 687 +++++++++++++++++- crates/buzz-acp/src/lib.rs | 326 +++++++-- crates/buzz-acp/tests/config_env.rs | 21 + .../src-tauri/src/managed_agents/env_vars.rs | 1 + .../src/managed_agents/env_vars/tests.rs | 1 + 8 files changed, 1415 insertions(+), 109 deletions(-) create mode 100644 crates/buzz-acp/tests/config_env.rs diff --git a/.env.example b/.env.example index b9bfcada0e..c65a8c19e0 100644 --- a/.env.example +++ b/.env.example @@ -160,6 +160,11 @@ RUST_LOG=buzz_relay=debug,buzz_datastore=info,buzz_db=debug,buzz_auth=debug,buzz # Binary for an optional MCP server sidecar (e.g. buzz-dev-mcp for buzz-agent). # BUZZ_ACP_MCP_COMMAND= +# Path to an optional version 1 JSON file defining additional stdio MCP servers. +# This file may contain credentials. Keep it out of Git and restrict it to +# its owner. +# BUZZ_ACP_MCP_CONFIG=/absolute/path/to/mcp-servers.json + # Number of parallel agent subprocesses (1–32). # BUZZ_ACP_AGENTS=1 diff --git a/crates/buzz-acp/README.md b/crates/buzz-acp/README.md index e6164b02dd..d5384df7e2 100644 --- a/crates/buzz-acp/README.md +++ b/crates/buzz-acp/README.md @@ -111,6 +111,7 @@ All configuration is via environment variables (or CLI flags — every env var h | `BUZZ_ACP_AGENT_COMMAND` | no | `goose` | Agent binary to spawn. | | `BUZZ_ACP_AGENT_ARGS` | no | `acp` | Agent arguments (comma-separated). | | `BUZZ_ACP_MCP_COMMAND` | no | `""` (empty) | Path to an optional MCP server binary to provide to the agent subprocess. | +| `BUZZ_ACP_MCP_CONFIG` | no | `""` (empty) | Path to a version 1 JSON file defining additional stdio MCP servers. | | `BUZZ_ACP_IDLE_TIMEOUT` | no | `620` | Idle timeout: max seconds of silence before cancelling a turn. Resets on any agent stdout activity. | | `BUZZ_ACP_MAX_TURN_DURATION` | no | `7200` | Absolute wall-clock cap per turn (safety valve). | | `BUZZ_API_TOKEN` | no | — | API token (required if relay enforces token auth). | @@ -119,6 +120,60 @@ All configuration is via environment variables (or CLI flags — every env var h **Legacy env vars:** `BUZZ_ACP_PRIVATE_KEY`, `BUZZ_ACP_API_TOKEN`, and `BUZZ_ACP_TURN_TIMEOUT` (replaced by `BUZZ_ACP_IDLE_TIMEOUT`) are still accepted as fallbacks. +### Multiple MCP servers + +Use `--mcp-config ` or `BUZZ_ACP_MCP_CONFIG` to add named stdio MCP +servers: + +```json +{ + "version": 1, + "servers": [ + { + "name": "analytics", + "command": "/opt/mcp/analytics-server", + "args": ["--stdio"], + "env": { + "ANALYTICS_TOKEN": "replace-me" + } + } + ] +} +``` + +The JSON is strict. The only top-level fields are `version` and `servers`. +Each server has `name`, `command`, `args`, and `env`. Server names must be +unique, contain 1 to 128 ASCII bytes using only letters, digits, `_`, or `-`, +and cannot contain `__`. Names are checked across both structured entries and +the legacy server. + +The config file is limited to 64 KiB. A harness can have at most 16 MCP +servers in total, including the server from `BUZZ_ACP_MCP_COMMAND`. An +unreadable file, malformed JSON, an unsupported version, an unknown field, or +an invalid server entry stops startup. Buzz does not silently drop a server. + +`BUZZ_ACP_MCP_COMMAND` keeps its current behavior. It defines one privileged +Buzz companion and receives the relay URL and Buzz identity credentials. +For a structured server, Buzz puts only the values listed in its `env` object +into the ACP `env` list. Protected Buzz identity and authentication keys are +rejected. Buzz sends the list to the ACP adapter in `session/new`, and the +adapter controls the MCP processes. Treat the adapter as a credential broker +and use one you trust. + +The adapter still inherits the harness environment so its shell tools can use +the `buzz` CLI. Some adapters may propagate inherited variables to MCP child +processes. Per-server `env` entries are explicit configuration, not a process +isolation boundary. Use a separate account, container, or credential-brokered +service when the MCP process must not inherit adapter credentials. + +If the JSON contains secrets, keep it outside Git and restrict the file to its +owner. On Unix: + +```bash +chmod 600 /absolute/path/to/mcp-servers.json +buzz-acp --mcp-config /absolute/path/to/mcp-servers.json +``` + ### Parallel Agents & Heartbeat | Flag | Env Var | Default | Description | diff --git a/crates/buzz-acp/src/acp.rs b/crates/buzz-acp/src/acp.rs index 700d5e8dcf..0ecdf6cd21 100644 --- a/crates/buzz-acp/src/acp.rs +++ b/crates/buzz-acp/src/acp.rs @@ -20,6 +20,8 @@ use crate::usage::{TurnUsage, UsageTracker}; /// Lines exceeding this limit are rejected to prevent OOM from rogue agents. const MAX_LINE_SIZE: usize = 10_000_000; // 10 MB +const REDACTED_ENV_VALUE: &str = "[REDACTED]"; + /// An MCP server configuration passed to `session/new`. /// /// Corresponds to the `McpServerStdio` variant in the ACP schema. @@ -114,13 +116,56 @@ pub enum AcpError { /// detail (e.g. a `data` field) is not lost. fn agent_error_from_json(error: &serde_json::Value) -> AcpError { let code = error.get("code").and_then(|c| c.as_i64()).unwrap_or(-32000); - let message = match error.get("message").and_then(|m| m.as_str()) { + let redacted_error = redact_wire_value(error); + let message = match redacted_error.get("message").and_then(|m| m.as_str()) { Some(m) => m.to_string(), - None => error.to_string(), + None => redacted_error.to_string(), }; AcpError::AgentError { code, message } } +/// Return a logging-safe copy of an ACP wire value. +/// +/// MCP environment values are needed by the adapter on the real wire, but +/// must not reach tracing or observer frames. +fn redact_wire_value(value: &serde_json::Value) -> serde_json::Value { + fn redact_in_place(value: &mut serde_json::Value) { + match value { + serde_json::Value::Array(values) => { + for value in values { + redact_in_place(value); + } + } + serde_json::Value::Object(fields) => { + for (key, value) in fields { + if key == "env" { + if let serde_json::Value::Array(entries) = value { + for entry in entries { + if let serde_json::Value::Object(env_var) = entry { + if env_var.contains_key("value") { + env_var.insert( + "value".to_string(), + serde_json::Value::String( + REDACTED_ENV_VALUE.to_string(), + ), + ); + } + } + } + } + } + redact_in_place(value); + } + } + _ => {} + } + } + + let mut redacted = value.clone(); + redact_in_place(&mut redacted); + redacted +} + fn build_initialize_params() -> serde_json::Value { serde_json::json!({ "protocolVersion": 2, @@ -788,7 +833,6 @@ impl AcpClient { "params": params, }); - tracing::debug!(target: "acp::wire", "→ {}", &serde_json::to_string(&msg).unwrap_or_default()); if let Err(e) = self.write_ndjson(&msg).await { self.last_prompt_id = None; self.current_hard_deadline = None; @@ -1047,6 +1091,8 @@ impl AcpClient { /// (e.g., it's stuck or dead), the write would otherwise block forever. async fn write_ndjson(&mut self, value: &serde_json::Value) -> Result<(), AcpError> { const WRITE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30); + let observed_value = redact_wire_value(value); + tracing::debug!(target: "acp::wire", "→ {observed_value}"); let line = serde_json::to_string(value)?; tokio::time::timeout(WRITE_TIMEOUT, async { self.stdin.write_all(line.as_bytes()).await?; @@ -1057,10 +1103,43 @@ impl AcpClient { .await .map_err(|_| AcpError::WriteTimeout(WRITE_TIMEOUT))? .map_err(AcpError::Io)?; - self.observe("acp_write", value.clone()); + self.observe("acp_write", observed_value); Ok(()) } + /// Parse one non-empty agent stdout line and emit only a safe copy. + /// + /// Parse failures expose the line length and parser error, never the raw + /// line. Successful messages retain their raw value for protocol handling. + fn parse_inbound_line(&self, line: &str) -> Option { + match serde_json::from_str(line) { + Ok(msg) => { + let observed_value = redact_wire_value(&msg); + tracing::debug!(target: "acp::wire", "← {observed_value}"); + self.observe("acp_read", observed_value); + Some(msg) + } + Err(error) => { + let line_length = line.len(); + let error = error.to_string(); + self.observe( + "acp_parse_error", + serde_json::json!({ + "lineLength": line_length, + "error": error, + }), + ); + tracing::warn!( + target: "acp::wire", + line_length, + error = %error, + "failed to parse agent stdout as JSON; skipping" + ); + None + } + } + } + /// Default timeout for non-prompt RPCs (initialize, session/new, etc.). const REQUEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60); @@ -1088,8 +1167,6 @@ impl AcpClient { "params": params, }); - tracing::debug!(target: "acp::wire", "→ {}", &serde_json::to_string(&msg).unwrap_or_default()); - // Wrap write + read in a single timeout so a hung agent can't block forever. // We cannot use an async block that borrows `self` mutably across two awaits // inside timeout(), so we sequence them with early-return on timeout. @@ -1153,7 +1230,6 @@ impl AcpClient { "params": params, }); - tracing::debug!(target: "acp::wire", "→ (notification) {}", &serde_json::to_string(&msg).unwrap_or_default()); self.write_ndjson(&msg).await?; Ok(()) } @@ -1194,27 +1270,10 @@ impl AcpClient { continue; } - // Only log and reset idle after we have a valid non-empty line. - tracing::debug!(target: "acp::wire", "← {trimmed}"); - - let msg: serde_json::Value = match serde_json::from_str(trimmed) { - Ok(v) => v, - Err(e) => { - self.observe( - "acp_parse_error", - serde_json::json!({ - "line": trimmed, - "error": e.to_string(), - }), - ); - tracing::warn!( - target: "acp::wire", - "failed to parse line as JSON: {e} — skipping" - ); - continue; - } + let msg = match self.parse_inbound_line(trimmed) { + Some(msg) => msg, + None => continue, }; - self.observe("acp_read", msg.clone()); // Check if this is a response to our expected request (has matching id // AND no `method` field — a `method` field means it's an agent-initiated @@ -1439,11 +1498,6 @@ impl AcpClient { "method": method, "params": params, }); - tracing::debug!( - target: "acp::wire", - "→ {}", - serde_json::to_string(&msg).unwrap_or_default() - ); match self.write_ndjson(&msg).await { Ok(()) => { pending_steer = Some((id, transport, req.ack_tx)); @@ -1518,26 +1572,10 @@ impl AcpClient { continue; } - tracing::debug!(target: "acp::wire", "← {trimmed}"); - - let msg: serde_json::Value = match serde_json::from_str(trimmed) { - Ok(v) => v, - Err(e) => { - self.observe( - "acp_parse_error", - serde_json::json!({ - "line": trimmed, - "error": e.to_string(), - }), - ); - tracing::warn!( - target: "acp::wire", - "failed to parse line as JSON: {e} — skipping" - ); - continue; - } + let msg = match self.parse_inbound_line(trimmed) { + Some(msg) => msg, + None => continue, }; - self.observe("acp_read", msg.clone()); let activity_now = Instant::now(); idle_deadline = activity_now + idle_timeout; @@ -1563,7 +1601,7 @@ impl AcpClient { .get("code") .and_then(|c| c.as_i64()) .unwrap_or(-1); - let message = error.to_string(); + let message = redact_wire_value(error).to_string(); crate::pool::SteerAck::Err( crate::pool::SteerError::AgentError { code, message }, ) @@ -2459,6 +2497,49 @@ mod tests { ); } + #[test] + fn wire_redaction_covers_nested_env_values_without_changing_source() { + let source = serde_json::json!({ + "method": "session/new", + "params": { + "mcpServers": [{ + "name": "analytics", + "env": [ + {"name": "ANALYTICS_TOKEN", "value": "secret-one"}, + {"name": "EMPTY_VALUE", "value": ""} + ] + }], + "nested": { + "env": [{"name": "OTHER_TOKEN", "value": "secret-two"}] + }, + "ordinary": {"value": "keep-me"} + } + }); + + let redacted = redact_wire_value(&source); + + assert_eq!( + source["params"]["mcpServers"][0]["env"][0]["value"], "secret-one", + "the source value sent on the wire must stay unchanged" + ); + assert_eq!( + redacted["params"]["mcpServers"][0]["env"][0]["value"], + REDACTED_ENV_VALUE + ); + assert_eq!( + redacted["params"]["mcpServers"][0]["env"][1]["value"], + REDACTED_ENV_VALUE + ); + assert_eq!( + redacted["params"]["nested"]["env"][0]["value"], + REDACTED_ENV_VALUE + ); + assert_eq!( + redacted["params"]["ordinary"]["value"], "keep-me", + "value fields outside env arrays must remain visible" + ); + } + #[test] fn session_prompt_request_format() { let prompt_text = "[Buzz @mention]\nChannel: test\nFrom: npub1...\nMessage: hello"; @@ -2885,6 +2966,74 @@ mod tests { .expect("failed to spawn test script") } + fn assert_safe_parse_error_event(observer: &ObserverHandle, raw_line: &str) { + let event = observer + .snapshot() + .into_iter() + .find(|event| event.kind == "acp_parse_error") + .expect("parse error observer event"); + assert_eq!( + event.payload["lineLength"].as_u64(), + Some(raw_line.len() as u64) + ); + assert!(event.payload["error"].is_string()); + assert!( + event.payload.get("line").is_none(), + "the raw line field must not exist" + ); + assert!( + !event.payload.to_string().contains(raw_line), + "the raw malformed line must not reach the observer" + ); + } + + #[tokio::test] + async fn regular_read_loop_reports_malformed_json_without_raw_line() { + let raw_line = "regular-loop-sensitive-malformed-json"; + let script = format!( + "read -t 2 _REQ\nprintf '%s\\n' '{raw_line}'\nprintf '%s\\n' \ + '{{\"jsonrpc\":\"2.0\",\"id\":0,\"result\":{{\"ok\":true}}}}'\nsleep 1" + ); + let mut client = spawn_script(&script).await; + let observer = ObserverHandle::in_process(); + client.set_observer(Some(observer.clone()), 0); + + let result = client + .send_request("test/request", serde_json::json!({})) + .await + .expect("valid response after malformed line"); + assert_eq!(result["ok"], true); + assert_safe_parse_error_event(&observer, raw_line); + client.shutdown().await; + } + + #[tokio::test] + async fn idle_read_loop_reports_malformed_json_without_raw_line() { + let raw_line = "idle-loop-sensitive-malformed-json"; + let script = format!( + "printf '%s\\n' '{raw_line}'\nprintf '%s\\n' \ + '{{\"jsonrpc\":\"2.0\",\"id\":999,\"result\":{{\"ok\":true}}}}'\nsleep 1" + ); + let mut client = spawn_script(&script).await; + let observer = ObserverHandle::in_process(); + client.set_observer(Some(observer.clone()), 0); + let max_duration = std::time::Duration::from_secs(5); + + let result = client + .read_until_response_with_idle_timeout( + "test", + 999, + std::time::Duration::from_secs(1), + tokio::time::Instant::now() + max_duration, + max_duration, + ) + .await + .expect("valid idle-loop response after malformed line"); + assert_eq!(result["ok"], true); + assert_safe_parse_error_event(&observer, raw_line); + client.shutdown().await; + } + /// Spawn a probe script whose file name carries a runtime identity (e.g. /// `hermes-acp`) and return the value of `var` as the child observed it. /// `` means the child did not receive the var. @@ -3322,6 +3471,161 @@ mod tests { ); } + #[tokio::test] + async fn session_new_sends_real_mcp_env_but_observer_only_sees_redacted_values() { + let script = r#" + read -t 2 _init + echo '{"jsonrpc":"2.0","id":0,"result":{"protocolVersion":2,"agentCapabilities":{}}}' + read -t 2 REQ + echo '{"jsonrpc":"2.0","id":1,"result":{"sessionId":"ses_secret_test","_receivedRequest":'"$REQ"'}}' + sleep 1 + "#; + let mut client = spawn_script(script).await; + let observer = ObserverHandle::in_process(); + client.set_observer(Some(observer.clone()), 0); + client + .initialize() + .await + .expect("initialize should succeed"); + + let secret = "mcp-secret-must-not-reach-observer"; + let response = client + .session_new_full( + "/tmp", + vec![McpServer { + name: "analytics".into(), + command: "analytics-mcp".into(), + args: vec!["--stdio".into()], + env: vec![EnvVar { + name: "ANALYTICS_TOKEN".into(), + value: secret.into(), + }], + }], + None, + None, + ) + .await + .expect("session/new should succeed"); + + assert_eq!( + response.raw["_receivedRequest"]["params"]["mcpServers"][0]["env"][0]["value"], secret, + "the adapter must receive the real MCP environment value" + ); + + let events = observer.snapshot(); + let session_write = events + .iter() + .find(|event| event.kind == "acp_write" && event.payload["method"] == "session/new") + .expect("session/new write observer event"); + assert_eq!( + session_write.payload["params"]["mcpServers"][0]["env"][0]["value"], + REDACTED_ENV_VALUE + ); + + let echoed_read = events + .iter() + .find(|event| { + event.kind == "acp_read" + && event.payload["result"]["sessionId"] == "ses_secret_test" + }) + .expect("session/new response observer event"); + assert_eq!( + echoed_read.payload["result"]["_receivedRequest"]["params"]["mcpServers"][0]["env"][0] + ["value"], + REDACTED_ENV_VALUE + ); + + let serialized_events = + serde_json::to_string(&events).expect("serialize observer snapshot"); + assert!( + !serialized_events.contains(secret), + "no observer frame may contain the real MCP environment value" + ); + client.shutdown().await; + } + + #[tokio::test] + async fn structured_mcp_servers_survive_repeated_sessions_and_adapter_restart() { + let script = r#" + read -t 2 _init + echo '{"jsonrpc":"2.0","id":0,"result":{"protocolVersion":2,"agentCapabilities":{}}}' + read -t 2 FIRST + echo '{"jsonrpc":"2.0","id":1,"result":{"sessionId":"ses_first","_receivedRequest":'"$FIRST"'}}' + read -t 2 SECOND + echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"ses_second","_receivedRequest":'"$SECOND"'}}' + sleep 1 + "#; + let servers = vec![ + McpServer { + name: "analytics".into(), + command: "/opt/MCP Servers/analytics,prod".into(), + args: vec!["--stdio".into(), "literal value".into()], + env: vec![EnvVar { + name: "ANALYTICS_ENDPOINT".into(), + value: "https://example.test/a=b".into(), + }], + }, + McpServer { + name: "search".into(), + command: "/opt/search-mcp".into(), + args: vec![], + env: vec![], + }, + ]; + + for _restart in 0..2 { + let mut client = spawn_script(script).await; + let observer = ObserverHandle::in_process(); + client.set_observer(Some(observer.clone()), 0); + client + .initialize() + .await + .expect("initialize should succeed"); + + for expected_session_id in ["ses_first", "ses_second"] { + let response = client + .session_new_full("/tmp", servers.clone(), None, None) + .await + .expect("session/new should succeed"); + assert_eq!(response.session_id, expected_session_id); + let received = &response.raw["_receivedRequest"]["params"]["mcpServers"]; + assert_eq!(received[0]["name"], "analytics"); + assert_eq!(received[0]["command"], "/opt/MCP Servers/analytics,prod"); + assert_eq!( + received[0]["args"], + serde_json::json!(["--stdio", "literal value"]) + ); + assert_eq!(received[0]["env"][0]["name"], "ANALYTICS_ENDPOINT"); + assert_eq!(received[0]["env"][0]["value"], "https://example.test/a=b"); + assert_eq!(received[1]["name"], "search"); + assert_eq!(received[1]["command"], "/opt/search-mcp"); + } + + let writes = observer + .snapshot() + .into_iter() + .filter(|event| { + event.kind == "acp_write" && event.payload["method"] == "session/new" + }) + .collect::>(); + assert_eq!(writes.len(), 2); + for write in writes { + let sent = &write.payload["params"]["mcpServers"]; + assert_eq!(sent[0]["name"], "analytics"); + assert_eq!(sent[0]["command"], "/opt/MCP Servers/analytics,prod"); + assert_eq!( + sent[0]["args"], + serde_json::json!(["--stdio", "literal value"]) + ); + assert_eq!(sent[0]["env"][0]["name"], "ANALYTICS_ENDPOINT"); + assert_eq!(sent[0]["env"][0]["value"], REDACTED_ENV_VALUE); + assert_eq!(sent[1]["name"], "search"); + assert_eq!(sent[1]["command"], "/opt/search-mcp"); + } + client.shutdown().await; + } + } + #[tokio::test] async fn goose_system_prompt_request_uses_append_contract() { let script = r#" @@ -4368,6 +4672,26 @@ mod tests { } } + #[test] + fn agent_error_from_json_redacts_mcp_env_before_display_or_turn_error() { + let secret = "agent-error-secret"; + let error = serde_json::json!({ + "code": -32002, + "data": { + "env": [{ + "name": "ANALYTICS_TOKEN", + "value": secret + }] + } + }); + + let rendered = super::agent_error_from_json(&error).to_string(); + + assert!(!rendered.contains(secret)); + assert!(rendered.contains(REDACTED_ENV_VALUE)); + assert_eq!(error["data"]["env"][0]["value"], secret); + } + #[test] fn agent_error_from_json_uses_message_field_when_present() { let error = serde_json::json!({"code": -32001, "message": "auth denied"}); diff --git a/crates/buzz-acp/src/config.rs b/crates/buzz-acp/src/config.rs index 35aaec188d..f23ddeca06 100644 --- a/crates/buzz-acp/src/config.rs +++ b/crates/buzz-acp/src/config.rs @@ -3,7 +3,8 @@ //! CLI-first: every option is a CLI flag with env var fallback. //! Config file (TOML) for complex subscription rules. -use std::collections::{HashMap, HashSet}; +use std::collections::{BTreeMap, HashMap, HashSet}; +use std::io::Read; use std::path::PathBuf; use clap::Parser; @@ -47,6 +48,250 @@ pub enum ConfigError { ConfigFile(String), } +const MCP_CONFIG_VERSION: u32 = 1; +const MCP_CONFIG_MAX_BYTES: u64 = 64 * 1024; +const MCP_SERVER_MAX_COUNT: usize = 16; +const MCP_SERVER_MAX_ARGS: usize = 128; +const MCP_SERVER_MAX_ENV: usize = 128; +const MCP_SERVER_NAME_MAX_BYTES: usize = 128; +const PROTECTED_MCP_ENV_NAMES: [&str; 6] = [ + "BUZZ_PRIVATE_KEY", + "NOSTR_PRIVATE_KEY", + "BUZZ_AUTH_TAG", + "BUZZ_API_TOKEN", + "BUZZ_ACP_PRIVATE_KEY", + "BUZZ_ACP_API_TOKEN", +]; + +/// One local stdio MCP server loaded from the structured MCP configuration. +#[derive(Debug, Clone, PartialEq, Eq, serde::Deserialize)] +#[serde(deny_unknown_fields)] +pub struct ConfiguredMcpServer { + /// Stable ACP identifier for this server. + pub name: String, + /// Executable to invoke, passed directly without shell parsing. + pub command: String, + /// Arguments passed to the executable in their configured order. + pub args: Vec, + /// Server-specific environment in deterministic key order. + #[serde(deserialize_with = "deserialize_mcp_env")] + pub env: BTreeMap, +} + +#[derive(Debug, serde::Deserialize)] +#[serde(deny_unknown_fields)] +struct McpConfigDocument { + version: u32, + servers: Vec, +} + +fn deserialize_mcp_env<'de, D>(deserializer: D) -> Result, D::Error> +where + D: serde::Deserializer<'de>, +{ + struct EnvVisitor; + + impl<'de> serde::de::Visitor<'de> for EnvVisitor { + type Value = BTreeMap; + + fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter.write_str("an object containing unique environment variable names") + } + + fn visit_map(self, mut map: A) -> Result + where + A: serde::de::MapAccess<'de>, + { + let mut env = BTreeMap::new(); + let mut normalized_names = HashSet::new(); + while let Some((key, value)) = map.next_entry::()? { + if !normalized_names.insert(key.to_ascii_uppercase()) { + return Err(serde::de::Error::custom(format!( + "duplicate environment key '{key}'" + ))); + } + env.insert(key, value); + } + Ok(env) + } + } + + deserializer.deserialize_map(EnvVisitor) +} + +/// Derive the ACP name used by the legacy single-command MCP configuration. +/// +/// This preserves the existing `build_mcp_servers` behavior so collision +/// validation and runtime construction use the same name. +pub fn legacy_mcp_server_name(command: &str) -> String { + std::path::Path::new(command) + .file_stem() + .and_then(|stem| stem.to_str()) + .unwrap_or("mcp") + .to_string() +} + +fn read_mcp_config(path: &std::path::Path) -> Result, ConfigError> { + let mut options = std::fs::OpenOptions::new(); + options.read(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.custom_flags(nix::libc::O_NONBLOCK); + } + let file = options.open(path).map_err(|error| { + ConfigError::ConfigFile(format!( + "failed to open MCP config {}: {error}", + path.display() + )) + })?; + let metadata = file.metadata().map_err(|error| { + ConfigError::ConfigFile(format!( + "failed to inspect MCP config {}: {error}", + path.display() + )) + })?; + if !metadata.is_file() { + return Err(ConfigError::ConfigFile(format!( + "MCP config {} must be a regular file", + path.display() + ))); + } + let mut content = Vec::new(); + file.take(MCP_CONFIG_MAX_BYTES + 1) + .read_to_end(&mut content) + .map_err(|error| { + ConfigError::ConfigFile(format!( + "failed to read MCP config {}: {error}", + path.display() + )) + })?; + if content.len() as u64 > MCP_CONFIG_MAX_BYTES { + return Err(ConfigError::ConfigFile(format!( + "MCP config {} exceeds the {} byte limit", + path.display(), + MCP_CONFIG_MAX_BYTES + ))); + } + Ok(content) +} + +fn valid_mcp_server_name(name: &str) -> bool { + !name.is_empty() + && name.len() <= MCP_SERVER_NAME_MAX_BYTES + && !name.contains("__") + && name + .bytes() + .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-')) +} + +fn valid_mcp_env_name(name: &str) -> bool { + let mut bytes = name.bytes(); + matches!(bytes.next(), Some(byte) if byte.is_ascii_alphabetic() || byte == b'_') + && bytes.all(|byte| byte.is_ascii_alphanumeric() || byte == b'_') +} + +fn load_mcp_config( + path: &std::path::Path, + legacy_mcp_command: &str, +) -> Result, ConfigError> { + let content = read_mcp_config(path)?; + let document: McpConfigDocument = serde_json::from_slice(&content).map_err(|error| { + ConfigError::ConfigFile(format!("invalid MCP config {}: {error}", path.display())) + })?; + + if document.version != MCP_CONFIG_VERSION { + return Err(ConfigError::ConfigFile(format!( + "unsupported MCP config version {} (expected {})", + document.version, MCP_CONFIG_VERSION + ))); + } + + let legacy_count = usize::from(!legacy_mcp_command.is_empty()); + if document.servers.len() + legacy_count > MCP_SERVER_MAX_COUNT { + return Err(ConfigError::ConfigFile(format!( + "too many MCP servers ({} structured + {legacy_count} legacy, max {MCP_SERVER_MAX_COUNT})", + document.servers.len() + ))); + } + + let legacy_name = + (!legacy_mcp_command.is_empty()).then(|| legacy_mcp_server_name(legacy_mcp_command)); + let mut names = HashSet::with_capacity(document.servers.len()); + for (index, server) in document.servers.iter().enumerate() { + if !valid_mcp_server_name(&server.name) { + return Err(ConfigError::ConfigFile(format!( + "MCP server {} has invalid name '{}': use 1 to {MCP_SERVER_NAME_MAX_BYTES} ASCII letters, digits, underscores, or hyphens, without '__'", + index + 1, + server.name + ))); + } + if !names.insert(server.name.as_str()) { + return Err(ConfigError::ConfigFile(format!( + "duplicate MCP server name '{}'", + server.name + ))); + } + if legacy_name.as_deref() == Some(server.name.as_str()) { + return Err(ConfigError::ConfigFile(format!( + "MCP server name '{}' collides with the legacy --mcp-command server", + server.name + ))); + } + if server.command.is_empty() || server.command.contains('\0') { + return Err(ConfigError::ConfigFile(format!( + "MCP server '{}' command must be nonempty and contain no NUL bytes", + server.name + ))); + } + if server.args.len() > MCP_SERVER_MAX_ARGS { + return Err(ConfigError::ConfigFile(format!( + "MCP server '{}' has too many arguments ({}, max {MCP_SERVER_MAX_ARGS})", + server.name, + server.args.len() + ))); + } + if server.args.iter().any(|argument| argument.contains('\0')) { + return Err(ConfigError::ConfigFile(format!( + "MCP server '{}' arguments must contain no NUL bytes", + server.name + ))); + } + if server.env.len() > MCP_SERVER_MAX_ENV { + return Err(ConfigError::ConfigFile(format!( + "MCP server '{}' has too many environment entries ({}, max {MCP_SERVER_MAX_ENV})", + server.name, + server.env.len() + ))); + } + for (key, value) in &server.env { + if !valid_mcp_env_name(key) { + return Err(ConfigError::ConfigFile(format!( + "MCP server '{}' has invalid environment key '{key}'", + server.name + ))); + } + if PROTECTED_MCP_ENV_NAMES + .iter() + .any(|protected| key.eq_ignore_ascii_case(protected)) + { + return Err(ConfigError::ConfigFile(format!( + "MCP server '{}' may not configure protected environment key '{key}'", + server.name + ))); + } + if value.contains('\0') { + return Err(ConfigError::ConfigFile(format!( + "MCP server '{}' environment value for '{key}' contains a NUL byte", + server.name + ))); + } + } + } + + Ok(document.servers) +} + #[derive(Debug, Clone, PartialEq, clap::ValueEnum)] pub enum SubscribeMode { Mentions, @@ -261,6 +506,10 @@ pub struct CliArgs { #[arg(long, env = "BUZZ_ACP_MCP_COMMAND", default_value = "")] pub mcp_command: String, + /// Path to a versioned JSON document defining additional local MCP servers. + #[arg(long, env = "BUZZ_ACP_MCP_CONFIG")] + pub mcp_config: Option, + /// Idle timeout: max seconds of silence before killing a turn. /// Resets on any agent stdout activity. #[arg(long, env = "BUZZ_ACP_IDLE_TIMEOUT")] @@ -500,6 +749,8 @@ pub struct Config { pub agent_command: String, pub agent_args: Vec, pub mcp_command: String, + /// Additional local MCP servers loaded once from `--mcp-config`. + pub configured_mcp_servers: Vec, pub idle_timeout_secs: u64, pub max_turn_duration_secs: u64, pub agents: u32, @@ -827,10 +1078,23 @@ pub fn propagate_legacy_env_vars() { } } +/// Prepare environment fallbacks before Clap and Tokio read process state. +/// +/// Deployment templates commonly render optional values as empty strings. +/// Clap treats an empty value for `Option` as an invalid supplied +/// value, so normalize this one optional path to the same state as an unset +/// variable before argument parsing starts. +pub fn prepare_process_env() { + propagate_legacy_env_vars(); + if std::env::var_os("BUZZ_ACP_MCP_CONFIG").is_some_and(|value| value.is_empty()) { + std::env::remove_var("BUZZ_ACP_MCP_CONFIG"); + } +} + impl Config { pub fn from_cli() -> Result { // Legacy env-var propagation is intentionally NOT done here. - // Call `propagate_legacy_env_vars()` before the tokio runtime starts + // Call `prepare_process_env()` before the tokio runtime starts // (in the sync `fn main()` wrapper) — see Rust 2024 edition safety. let args = CliArgs::parse(); Self::from_args(args) @@ -913,6 +1177,10 @@ impl Config { } let agent_args = normalize_agent_args(&agent_command, args.agent_args); + let configured_mcp_servers = match args.mcp_config.as_deref() { + Some(path) => load_mcp_config(path, &args.mcp_command)?, + None => Vec::new(), + }; if let Some(ref channels) = args.channels { for ch in channels { @@ -1066,6 +1334,7 @@ impl Config { agent_command, agent_args, mcp_command: args.mcp_command, + configured_mcp_servers, idle_timeout_secs, max_turn_duration_secs, agents: args.agents, @@ -1131,12 +1400,13 @@ impl Config { format!(" allowed_respond_to=[{}]", modes.join(",")) }; format!( - "relay={} pubkey={} agent_cmd={} {} mcp_cmd={} idle_timeout={}s max_turn={}s agents={} heartbeat={}s subscribe={:?} dedup={:?} meh={:?} ignore_self={} context_limit={} max_turns_per_session={} presence={} typing={} memory={} model={} permission_mode={} {}{}", + "relay={} pubkey={} agent_cmd={} {} legacy_mcp_server={} structured_mcp_servers={} idle_timeout={}s max_turn={}s agents={} heartbeat={}s subscribe={:?} dedup={:?} meh={:?} ignore_self={} context_limit={} max_turns_per_session={} presence={} typing={} memory={} model={} permission_mode={} {}{}", self.relay_url, self.keys.public_key().to_hex(), self.agent_command, self.agent_args.join(" "), - self.mcp_command, + !self.mcp_command.is_empty(), + self.configured_mcp_servers.len(), self.idle_timeout_secs, self.max_turn_duration_secs, self.agents, @@ -1445,6 +1715,7 @@ mod tests { agent_command: "goose".into(), agent_args: vec!["acp".into()], mcp_command: "".into(), + configured_mcp_servers: Vec::new(), idle_timeout_secs: DEFAULT_IDLE_TIMEOUT_SECS, max_turn_duration_secs: DEFAULT_MAX_TURN_DURATION_SECS, agents: 1, @@ -2924,6 +3195,414 @@ channels = "ALL" assert_eq!(compose_session_title(&agent, Some("buzz-dev")), agent); } + struct TempMcpConfig { + path: PathBuf, + } + + impl TempMcpConfig { + fn write(content: &[u8]) -> Self { + let path = + std::env::temp_dir().join(format!("buzz-acp-mcp-config-{}.json", Uuid::new_v4())); + std::fs::write(&path, content).expect("write temporary MCP config"); + Self { path } + } + } + + impl Drop for TempMcpConfig { + fn drop(&mut self) { + let _ = std::fs::remove_file(&self.path); + } + } + + fn config_from_mcp_file( + file: &TempMcpConfig, + legacy_command: Option<&str>, + ) -> Result { + let mut argv = vec![ + "buzz-acp".to_string(), + "--private-key".to_string(), + TEST_PRIVATE_KEY.to_string(), + "--mcp-config".to_string(), + file.path.display().to_string(), + ]; + if let Some(command) = legacy_command { + argv.push("--mcp-command".to_string()); + argv.push(command.to_string()); + } + let args = CliArgs::try_parse_from(argv).expect("clap should parse MCP config arguments"); + Config::from_args(args) + } + + fn server_json(name: &str) -> serde_json::Value { + serde_json::json!({ + "name": name, + "command": "mcp", + "args": [], + "env": {} + }) + } + + fn document_json(servers: Vec) -> Vec { + serde_json::to_vec(&serde_json::json!({ + "version": MCP_CONFIG_VERSION, + "servers": servers + })) + .expect("serialize MCP test document") + } + + #[test] + fn structured_mcp_config_preserves_order_and_exact_values() { + let posix_command = "/Applications/Tool Suite/工具 mcp"; + let windows_command = r"C:\Program Files\Agent Tools\server.exe"; + let literal_metacharacters = r#"$HOME;$(echo nope)|&<>*?`literal`"#; + let file = TempMcpConfig::write(&document_json(vec![ + serde_json::json!({ + "name": "analytics-primary", + "command": posix_command, + "args": [ + "", + "with spaces", + "comma,value", + "quote\"value", + r"C:\data\reports", + literal_metacharacters, + "雪" + ], + "env": { + "Z_LAST": "backslash\\quote\"雪", + "A_FIRST": "" + } + }), + serde_json::json!({ + "name": "windows_server", + "command": windows_command, + "args": ["--stdio"], + "env": {} + }), + ])); + + let config = config_from_mcp_file(&file, None).expect("structured config should load"); + assert_eq!(config.configured_mcp_servers.len(), 2); + assert_eq!(config.configured_mcp_servers[0].name, "analytics-primary"); + assert_eq!(config.configured_mcp_servers[0].command, posix_command); + assert_eq!( + config.configured_mcp_servers[0].args, + vec![ + "", + "with spaces", + "comma,value", + "quote\"value", + r"C:\data\reports", + literal_metacharacters, + "雪" + ] + ); + assert_eq!( + config.configured_mcp_servers[0] + .env + .keys() + .map(String::as_str) + .collect::>(), + vec!["A_FIRST", "Z_LAST"] + ); + assert_eq!( + config.configured_mcp_servers[0].env["Z_LAST"], + "backslash\\quote\"雪" + ); + assert_eq!(config.configured_mcp_servers[1].command, windows_command); + } + + #[test] + fn summary_reports_only_structured_mcp_count() { + let secret_value = "value-that-must-not-be-logged"; + let command = "/private/path/tool"; + let file = TempMcpConfig::write(&document_json(vec![serde_json::json!({ + "name": "safe", + "command": command, + "args": [], + "env": {"DOMAIN_TOKEN": secret_value} + })])); + let config = + config_from_mcp_file(&file, Some("/legacy/private/tool")).expect("config should load"); + + let summary = config.summary(); + assert!(summary.contains("legacy_mcp_server=true")); + assert!(summary.contains("structured_mcp_servers=1")); + assert!(!summary.contains(secret_value)); + assert!(!summary.contains(command)); + assert!(!summary.contains("/legacy/private/tool")); + assert!(!summary.contains("DOMAIN_TOKEN")); + } + + #[test] + fn legacy_mcp_name_matches_existing_file_stem_behavior() { + assert_eq!( + legacy_mcp_server_name("/opt/bin/my-mcp-server"), + "my-mcp-server" + ); + assert_eq!(legacy_mcp_server_name("."), "mcp"); + assert_eq!(legacy_mcp_server_name(""), "mcp"); + } + + #[test] + fn mcp_config_rejects_unreadable_file() { + let path = std::env::temp_dir().join(format!( + "buzz-acp-missing-mcp-config-{}.json", + Uuid::new_v4() + )); + let argv = vec![ + "buzz-acp".to_string(), + "--private-key".to_string(), + TEST_PRIVATE_KEY.to_string(), + "--mcp-config".to_string(), + path.display().to_string(), + ]; + let args = CliArgs::try_parse_from(argv).expect("clap should parse arguments"); + let error = Config::from_args(args).expect_err("missing MCP config must fail"); + assert!(error.to_string().contains("failed to open MCP config")); + } + + #[test] + fn mcp_config_rejects_non_regular_files() { + let path = std::env::temp_dir().join(format!("buzz-acp-mcp-config-dir-{}", Uuid::new_v4())); + std::fs::create_dir(&path).expect("create temporary MCP config directory"); + + let error = read_mcp_config(&path).expect_err("directories must not be read as MCP config"); + let _ = std::fs::remove_dir(&path); + assert!(error.to_string().contains("must be a regular file")); + } + + #[test] + fn mcp_config_enforces_file_size_boundary() { + let base = document_json(Vec::new()); + let mut at_limit = base.clone(); + at_limit.resize(MCP_CONFIG_MAX_BYTES as usize, b' '); + let file = TempMcpConfig::write(&at_limit); + config_from_mcp_file(&file, None).expect("64 KiB MCP config should be accepted"); + + let mut over_limit = base; + over_limit.resize(MCP_CONFIG_MAX_BYTES as usize + 1, b' '); + let file = TempMcpConfig::write(&over_limit); + let error = + config_from_mcp_file(&file, None).expect_err("MCP config over 64 KiB must fail"); + assert!(error.to_string().contains("65536 byte limit")); + } + + #[test] + fn mcp_config_rejects_malformed_wrong_version_and_unknown_fields() { + let cases: Vec<(&str, Vec)> = vec![ + ("malformed", br#"{"version":1,"servers":["#.to_vec()), + ( + "wrong version", + br#"{"version":2,"servers":[]}"#.to_vec(), + ), + ( + "unknown document field", + br#"{"version":1,"servers":[],"extra":true}"#.to_vec(), + ), + ( + "unknown server field", + br#"{"version":1,"servers":[{"name":"one","command":"mcp","args":[],"env":{},"extra":true}]}"#.to_vec(), + ), + ( + "missing required field", + br#"{"version":1,"servers":[{"name":"one","command":"mcp","env":{}}]}"#.to_vec(), + ), + ]; + + for (label, content) in cases { + let file = TempMcpConfig::write(&content); + assert!( + config_from_mcp_file(&file, None).is_err(), + "{label} should be rejected" + ); + } + } + + #[test] + fn mcp_config_validates_server_names_and_collisions() { + let invalid_names = vec![ + String::new(), + "contains space".to_string(), + "contains.dot".to_string(), + "double__underscore".to_string(), + "unicodé".to_string(), + "a".repeat(MCP_SERVER_NAME_MAX_BYTES + 1), + ]; + for invalid_name in invalid_names { + let file = TempMcpConfig::write(&document_json(vec![server_json(&invalid_name)])); + assert!( + config_from_mcp_file(&file, None).is_err(), + "invalid name {invalid_name:?} should fail" + ); + } + + let file = TempMcpConfig::write(&document_json(vec![ + server_json("same"), + server_json("same"), + ])); + let error = + config_from_mcp_file(&file, None).expect_err("duplicate server name should fail"); + assert!(error.to_string().contains("duplicate MCP server name")); + + let file = TempMcpConfig::write(&document_json(vec![server_json("my-mcp-server")])); + let error = config_from_mcp_file(&file, Some("/opt/bin/my-mcp-server")) + .expect_err("legacy name collision should fail"); + assert!(error.to_string().contains("collides")); + } + + #[test] + fn mcp_config_limits_total_servers_including_legacy() { + let sixteen = (0..MCP_SERVER_MAX_COUNT) + .map(|index| server_json(&format!("server-{index}"))) + .collect::>(); + let file = TempMcpConfig::write(&document_json(sixteen)); + config_from_mcp_file(&file, None).expect("16 structured servers should be accepted"); + assert!( + config_from_mcp_file(&file, Some("legacy-mcp")).is_err(), + "16 structured plus one legacy server should fail" + ); + + let seventeen = (0..=MCP_SERVER_MAX_COUNT) + .map(|index| server_json(&format!("server-{index}"))) + .collect::>(); + let file = TempMcpConfig::write(&document_json(seventeen)); + assert!( + config_from_mcp_file(&file, None).is_err(), + "17 structured servers should fail" + ); + } + + #[test] + fn mcp_config_enforces_argument_limit() { + let args = (0..MCP_SERVER_MAX_ARGS) + .map(|index| format!("arg-{index}")) + .collect::>(); + let file = TempMcpConfig::write(&document_json(vec![serde_json::json!({ + "name": "limit", + "command": "mcp", + "args": args, + "env": {} + })])); + config_from_mcp_file(&file, None).expect("128 arguments should be accepted"); + + let args = (0..=MCP_SERVER_MAX_ARGS) + .map(|index| format!("arg-{index}")) + .collect::>(); + let file = TempMcpConfig::write(&document_json(vec![serde_json::json!({ + "name": "over-limit", + "command": "mcp", + "args": args, + "env": {} + })])); + let error = + config_from_mcp_file(&file, None).expect_err("129 arguments should be rejected"); + assert!(error.to_string().contains("too many arguments")); + } + + #[test] + fn mcp_config_enforces_environment_limit() { + let env = (0..MCP_SERVER_MAX_ENV) + .map(|index| (format!("KEY_{index}"), format!("value-{index}"))) + .collect::>(); + let file = TempMcpConfig::write(&document_json(vec![serde_json::json!({ + "name": "limit", + "command": "mcp", + "args": [], + "env": env + })])); + config_from_mcp_file(&file, None).expect("128 environment entries should be accepted"); + + let env = (0..=MCP_SERVER_MAX_ENV) + .map(|index| (format!("KEY_{index}"), format!("value-{index}"))) + .collect::>(); + let file = TempMcpConfig::write(&document_json(vec![serde_json::json!({ + "name": "over-limit", + "command": "mcp", + "args": [], + "env": env + })])); + let error = config_from_mcp_file(&file, None) + .expect_err("129 environment entries should be rejected"); + assert!(error.to_string().contains("too many environment entries")); + } + + #[test] + fn mcp_config_rejects_invalid_duplicate_and_protected_env_names() { + for invalid_key in ["", "1STARTS_WITH_DIGIT", "BAD-NAME", "UNICODÉ"] { + let content = format!( + r#"{{"version":1,"servers":[{{"name":"one","command":"mcp","args":[],"env":{{"{invalid_key}":"value"}}}}]}}"# + ); + let file = TempMcpConfig::write(content.as_bytes()); + assert!( + config_from_mcp_file(&file, None).is_err(), + "invalid environment key {invalid_key:?} should fail" + ); + } + + for protected in PROTECTED_MCP_ENV_NAMES { + let lowercase = protected.to_ascii_lowercase(); + let content = format!( + r#"{{"version":1,"servers":[{{"name":"one","command":"mcp","args":[],"env":{{"{lowercase}":"value"}}}}]}}"# + ); + let file = TempMcpConfig::write(content.as_bytes()); + let error = config_from_mcp_file(&file, None) + .expect_err("protected environment key should fail case-insensitively"); + assert!(error.to_string().contains("protected environment key")); + } + + for duplicate_env in [ + r#"{"KEY":"one","KEY":"two"}"#, + r#"{"KEY":"one","key":"two"}"#, + ] { + let content = format!( + r#"{{"version":1,"servers":[{{"name":"one","command":"mcp","args":[],"env":{duplicate_env}}}]}}"# + ); + let file = TempMcpConfig::write(content.as_bytes()); + let error = + config_from_mcp_file(&file, None).expect_err("duplicate env key should fail"); + assert!(error.to_string().contains("duplicate environment key")); + } + } + + #[test] + fn mcp_config_rejects_empty_or_nul_process_values() { + let cases = vec![ + serde_json::json!({ + "name": "empty-command", + "command": "", + "args": [], + "env": {} + }), + serde_json::json!({ + "name": "nul-command", + "command": "mc\u{0}p", + "args": [], + "env": {} + }), + serde_json::json!({ + "name": "nul-arg", + "command": "mcp", + "args": ["ok", "bad\u{0}arg"], + "env": {} + }), + serde_json::json!({ + "name": "nul-value", + "command": "mcp", + "args": [], + "env": {"DOMAIN_KEY": "bad\u{0}value"} + }), + ]; + + for server in cases { + let file = TempMcpConfig::write(&document_json(vec![server])); + assert!( + config_from_mcp_file(&file, None).is_err(), + "invalid process value should fail" + ); + } + } + /// Every arg whose env var name contains KEY/SECRET/TOKEN/PASSWORD/CRED/AUTH /// must set `hide_env_values = true` to prevent credential leakage in --help. #[test] diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 811253e4ac..7f12d2a125 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -1282,7 +1282,7 @@ mod inactivity_tests { } pub fn run() -> Result<()> { - config::propagate_legacy_env_vars(); + config::prepare_process_env(); tokio_main() } @@ -4278,60 +4278,77 @@ async fn run_models(args: ModelsArgs) -> Result<()> { } fn build_mcp_servers(config: &Config) -> Vec { - if config.mcp_command.is_empty() { - return vec![]; - } - vec![McpServer { - name: std::path::Path::new(&config.mcp_command) - .file_stem() - .and_then(|s| s.to_str()) - .unwrap_or("mcp") - .to_string(), - command: config.mcp_command.clone(), - args: vec![], - env: { - let mut env = vec![ - EnvVar { - name: "BUZZ_RELAY_URL".into(), - value: config.relay_url.clone(), - }, - EnvVar { - name: "BUZZ_PRIVATE_KEY".into(), - // bech32 encoding of a valid secret key is infallible. - // Panic here is correct: injecting a bogus secret would cause - // delayed, hard-to-diagnose agent failures downstream. - value: config - .keys - .secret_key() - .to_bech32() - .expect("secret key bech32 encoding should never fail"), - }, - ]; - // Forward BUZZ_AUTH_TAG (NIP-OA owner attestation credential) - // so the MCP server can attach it to every signed event. - if let Ok(auth_tag) = std::env::var("BUZZ_AUTH_TAG") { - if !auth_tag.is_empty() { - env.push(EnvVar { - name: "BUZZ_AUTH_TAG".into(), - value: auth_tag, - }); + let mut servers = Vec::with_capacity( + usize::from(!config.mcp_command.is_empty()) + config.configured_mcp_servers.len(), + ); + + if !config.mcp_command.is_empty() { + servers.push(McpServer { + name: config::legacy_mcp_server_name(&config.mcp_command), + command: config.mcp_command.clone(), + args: vec![], + env: { + let mut env = vec![ + EnvVar { + name: "BUZZ_RELAY_URL".into(), + value: config.relay_url.clone(), + }, + EnvVar { + name: "BUZZ_PRIVATE_KEY".into(), + // bech32 encoding of a valid secret key is infallible. + // Panic here is correct: injecting a bogus secret would cause + // delayed, hard-to-diagnose agent failures downstream. + value: config + .keys + .secret_key() + .to_bech32() + .expect("secret key bech32 encoding should never fail"), + }, + ]; + // Forward BUZZ_AUTH_TAG (NIP-OA owner attestation credential) + // so the MCP server can attach it to every signed event. + if let Ok(auth_tag) = std::env::var("BUZZ_AUTH_TAG") { + if !auth_tag.is_empty() { + env.push(EnvVar { + name: "BUZZ_AUTH_TAG".into(), + value: auth_tag, + }); + } } - } - // Forward the agent's display name so dev-mcp can use it as the git - // author name instead of the raw npub. Read from the process env - // rather than Config: this is a pass-through of a contract owned - // upstream, and absent simply means dev-mcp falls back to the npub. - if let Ok(display_name) = std::env::var("BUZZ_ACP_DISPLAY_NAME") { - if !display_name.is_empty() { - env.push(EnvVar { - name: "BUZZ_ACP_DISPLAY_NAME".into(), - value: display_name, - }); + // Forward the agent's display name so dev-mcp can use it as the git + // author name instead of the raw npub. Read from the process env + // rather than Config: this is a pass-through of a contract owned + // upstream, and absent simply means dev-mcp falls back to the npub. + if let Ok(display_name) = std::env::var("BUZZ_ACP_DISPLAY_NAME") { + if !display_name.is_empty() { + env.push(EnvVar { + name: "BUZZ_ACP_DISPLAY_NAME".into(), + value: display_name, + }); + } } - } - env - }, - }] + env + }, + }); + } + + servers.extend(config.configured_mcp_servers.iter().map(|configured| { + McpServer { + name: configured.name.clone(), + command: configured.command.clone(), + args: configured.args.clone(), + env: configured + .env + .iter() + .map(|(name, value)| EnvVar { + name: name.clone(), + value: value.clone(), + }) + .collect(), + } + })); + + servers } #[cfg(test)] @@ -5089,6 +5106,8 @@ mod observer_chunk_coalescer_tests { #[cfg(test)] mod build_mcp_servers_tests { use super::*; + use clap::Parser; + use std::collections::BTreeMap; use std::sync::Mutex; /// Env-var-touching tests must run serially — env vars are process-global. @@ -5101,6 +5120,7 @@ mod build_mcp_servers_tests { agent_command: "goose".into(), agent_args: vec!["acp".into()], mcp_command: "test-mcp-server".into(), + configured_mcp_servers: Vec::new(), idle_timeout_secs: config::DEFAULT_IDLE_TIMEOUT_SECS, max_turn_duration_secs: config::DEFAULT_MAX_TURN_DURATION_SECS, agents: 1, @@ -5140,6 +5160,44 @@ mod build_mcp_servers_tests { } } + fn configured_server( + name: &str, + command: &str, + args: &[&str], + env: &[(&str, &str)], + ) -> config::ConfiguredMcpServer { + config::ConfiguredMcpServer { + name: name.into(), + command: command.into(), + args: args.iter().map(|arg| (*arg).to_string()).collect(), + env: env + .iter() + .map(|(name, value)| ((*name).to_string(), (*value).to_string())) + .collect::>(), + } + } + + struct TempMcpConfig { + path: std::path::PathBuf, + } + + impl TempMcpConfig { + fn write(content: &[u8]) -> Self { + let path = std::env::temp_dir().join(format!( + "buzz-acp-mcp-build-test-{}.json", + uuid::Uuid::new_v4() + )); + std::fs::write(&path, content).expect("write temporary MCP config"); + Self { path } + } + } + + impl Drop for TempMcpConfig { + fn drop(&mut self) { + let _ = std::fs::remove_file(&self.path); + } + } + #[test] fn session_new_mcp_server_has_required_fields() { let config = test_config(); @@ -5254,6 +5312,167 @@ mod build_mcp_servers_tests { ); } + #[test] + fn structured_servers_preserve_order_and_literal_values_without_legacy_credentials() { + let mut config = test_config(); + config.mcp_command.clear(); + config.configured_mcp_servers = vec![ + configured_server( + "analytics", + "/opt/MCP Servers/analytics,prod", + &[ + "--stdio", + "two words", + "comma,value", + "\"quoted\"", + r"C:\Program Files\MCP\server.exe", + "雪", + "$(literal)", + "`literal`", + "a|b;c", + ], + &[ + ("ANALYTICS_ENDPOINT", "https://example.test/a=b"), + ("LITERAL_VALUE", "$HOME;`id`|雪"), + ], + ), + configured_server("search", "/opt/search-mcp", &[], &[]), + ]; + + let servers = build_mcp_servers(&config); + + assert_eq!(servers.len(), 2); + assert_eq!(servers[0].name, "analytics"); + assert_eq!(servers[1].name, "search"); + assert_eq!(servers[0].command, "/opt/MCP Servers/analytics,prod"); + assert_eq!( + servers[0].args, + vec![ + "--stdio", + "two words", + "comma,value", + "\"quoted\"", + r"C:\Program Files\MCP\server.exe", + "雪", + "$(literal)", + "`literal`", + "a|b;c", + ] + ); + assert_eq!( + servers[0] + .env + .iter() + .map(|entry| (entry.name.as_str(), entry.value.as_str())) + .collect::>(), + vec![ + ("ANALYTICS_ENDPOINT", "https://example.test/a=b"), + ("LITERAL_VALUE", "$HOME;`id`|雪"), + ] + ); + assert!( + servers[0].env.iter().all(|entry| !matches!( + entry.name.as_str(), + "BUZZ_PRIVATE_KEY" | "BUZZ_AUTH_TAG" | "BUZZ_RELAY_URL" + )), + "structured servers must receive only their declared environment" + ); + } + + #[test] + fn legacy_server_remains_first_when_structured_servers_are_present() { + let mut config = test_config(); + config.configured_mcp_servers = vec![configured_server( + "analytics", + "/opt/analytics-mcp", + &["--stdio"], + &[("ANALYTICS_TOKEN", "opaque-test-token")], + )]; + + let servers = build_mcp_servers(&config); + + assert_eq!(servers.len(), 2); + assert_eq!(servers[0].name, "test-mcp-server"); + assert_eq!(servers[1].name, "analytics"); + assert!(servers[0] + .env + .iter() + .any(|entry| entry.name == "BUZZ_PRIVATE_KEY")); + assert_eq!( + servers[1] + .env + .iter() + .map(|entry| (entry.name.as_str(), entry.value.as_str())) + .collect::>(), + vec![("ANALYTICS_TOKEN", "opaque-test-token")] + ); + } + + #[test] + fn structured_json_reaches_initial_repeated_and_respawn_session_lists() { + let document = serde_json::json!({ + "version": 1, + "servers": [ + { + "name": "analytics", + "command": "/opt/MCP Servers/analytics,prod", + "args": ["--stdio", "literal value"], + "env": { + "ANALYTICS_ENDPOINT": "https://example.test/a=b" + } + }, + { + "name": "search", + "command": "/opt/search-mcp", + "args": [], + "env": {} + } + ] + }); + let file = TempMcpConfig::write( + &serde_json::to_vec(&document).expect("serialize MCP config fixture"), + ); + let private_key = nostr::Keys::generate() + .secret_key() + .to_bech32() + .expect("encode temporary private key"); + let args = config::CliArgs::try_parse_from([ + "buzz-acp", + "--private-key", + private_key.as_str(), + "--mcp-command", + "/opt/buzz-dev-mcp", + "--mcp-config", + file.path.to_str().expect("temporary path is UTF-8"), + ]) + .expect("parse MCP arguments"); + let config = Config::from_args(args).expect("load structured MCP config"); + + let initial_session = build_mcp_servers(&config); + let repeated_session = initial_session.clone(); + let respawned_session = repeated_session.clone(); + let initial_json = + serde_json::to_value(&initial_session).expect("serialize initial MCP list"); + let repeated_json = + serde_json::to_value(&repeated_session).expect("serialize repeated MCP list"); + let respawned_json = + serde_json::to_value(&respawned_session).expect("serialize respawned MCP list"); + + assert_eq!(initial_json, repeated_json); + assert_eq!(initial_json, respawned_json); + assert_eq!(initial_json.as_array().map(Vec::len), Some(3)); + assert_eq!(initial_json[0]["name"], "buzz-dev-mcp"); + assert_eq!(initial_json[1]["name"], "analytics"); + assert_eq!(initial_json[2]["name"], "search"); + assert_eq!( + initial_json[1]["env"], + serde_json::json!([{ + "name": "ANALYTICS_ENDPOINT", + "value": "https://example.test/a=b" + }]) + ); + } + #[test] fn absolute_path_mcp_command_uses_file_stem_as_name() { let mut config = test_config(); @@ -5323,6 +5542,7 @@ mod error_outcome_emission_tests { agent_command: "true".into(), agent_args: vec![], mcp_command: "test-mcp-server".into(), + configured_mcp_servers: Vec::new(), idle_timeout_secs: config::DEFAULT_IDLE_TIMEOUT_SECS, max_turn_duration_secs: config::DEFAULT_MAX_TURN_DURATION_SECS, agents: 1, diff --git a/crates/buzz-acp/tests/config_env.rs b/crates/buzz-acp/tests/config_env.rs new file mode 100644 index 0000000000..177ac3fea5 --- /dev/null +++ b/crates/buzz-acp/tests/config_env.rs @@ -0,0 +1,21 @@ +use std::process::Command; + +#[test] +fn empty_mcp_config_environment_value_is_treated_as_unset() { + let output = Command::new(env!("CARGO_BIN_EXE_buzz-acp")) + .env("BUZZ_ACP_MCP_CONFIG", "") + .args(["--private-key", "not-a-valid-nostr-key"]) + .output() + .expect("run buzz-acp"); + + let stderr = String::from_utf8_lossy(&output.stderr); + assert!(!output.status.success()); + assert!( + stderr.contains("configuration error: failed to parse nostr keys"), + "empty optional MCP config should reach normal configuration validation: {stderr}" + ); + assert!( + !stderr.contains("a value is required for '--mcp-config"), + "empty optional MCP config must not fail Clap parsing: {stderr}" + ); +} diff --git a/desktop/src-tauri/src/managed_agents/env_vars.rs b/desktop/src-tauri/src/managed_agents/env_vars.rs index 1653371e7f..965c173cc0 100644 --- a/desktop/src-tauri/src/managed_agents/env_vars.rs +++ b/desktop/src-tauri/src/managed_agents/env_vars.rs @@ -71,6 +71,7 @@ pub(crate) const RESERVED_ENV_KEYS: &[&str] = &[ "BUZZ_ACP_AGENT_COMMAND", "BUZZ_ACP_AGENT_ARGS", "BUZZ_ACP_MCP_COMMAND", + "BUZZ_ACP_MCP_CONFIG", // Security gates: respond-to mode + allowlist + legacy owner-only // fallback. Overriding would make the running agent's gate diverge // from the saved/UI-visible settings. diff --git a/desktop/src-tauri/src/managed_agents/env_vars/tests.rs b/desktop/src-tauri/src/managed_agents/env_vars/tests.rs index 534c2e0835..117dc5187c 100644 --- a/desktop/src-tauri/src/managed_agents/env_vars/tests.rs +++ b/desktop/src-tauri/src/managed_agents/env_vars/tests.rs @@ -175,6 +175,7 @@ fn reserved_keys_include_code_execution_surface() { "BUZZ_ACP_AGENT_COMMAND", "BUZZ_ACP_AGENT_ARGS", "BUZZ_ACP_MCP_COMMAND", + "BUZZ_ACP_MCP_CONFIG", ] { assert!(is_reserved_env_key(key), "{key} should be reserved"); } From cf126c9d47c56f7e81d8d5a2a5d36d2cad2177eb Mon Sep 17 00:00:00 2001 From: KC <79471844+wolfyy970@users.noreply.github.com> Date: Sat, 1 Aug 2026 17:28:10 -0400 Subject: [PATCH 2/6] buzz-acp: suppress serialized MCP config echoes Signed-off-by: KC <79471844+wolfyy970@users.noreply.github.com> --- crates/buzz-acp/src/acp.rs | 305 ++++++++++++++++++++++++++++++++++++- 1 file changed, 299 insertions(+), 6 deletions(-) diff --git a/crates/buzz-acp/src/acp.rs b/crates/buzz-acp/src/acp.rs index 0ecdf6cd21..45b7f3d597 100644 --- a/crates/buzz-acp/src/acp.rs +++ b/crates/buzz-acp/src/acp.rs @@ -124,6 +124,42 @@ fn agent_error_from_json(error: &serde_json::Value) -> AcpError { AcpError::AgentError { code, message } } +fn contains_serialized_json_key(text: &str, key: &str) -> bool { + text.match_indices(key).any(|(start, _)| { + let bytes = text.as_bytes(); + if start == 0 || bytes[start - 1] != b'"' { + return false; + } + + let mut cursor = start + key.len(); + while bytes.get(cursor) == Some(&b'\\') { + cursor += 1; + } + if bytes.get(cursor) != Some(&b'"') { + return false; + } + cursor += 1; + while bytes + .get(cursor) + .is_some_and(|byte| byte.is_ascii_whitespace()) + { + cursor += 1; + } + bytes.get(cursor) == Some(&b':') + }) +} + +fn redact_wire_text(text: &str) -> &str { + let looks_like_serialized_mcp_config = contains_serialized_json_key(text, "mcpServers") + && contains_serialized_json_key(text, "env") + && contains_serialized_json_key(text, "value"); + if looks_like_serialized_mcp_config { + REDACTED_ENV_VALUE + } else { + text + } +} + /// Return a logging-safe copy of an ACP wire value. /// /// MCP environment values are needed by the adapter on the real wire, but @@ -157,6 +193,12 @@ fn redact_wire_value(value: &serde_json::Value) -> serde_json::Value { redact_in_place(value); } } + serde_json::Value::String(text) => { + let observed = redact_wire_text(text); + if observed != text.as_str() { + *text = observed.to_string(); + } + } _ => {} } } @@ -621,6 +663,14 @@ impl AcpClient { /// Emit a semantic event to the local observer feed, if enabled. pub fn observe(&self, kind: impl Into, payload: serde_json::Value) { + if self.observer.is_none() { + return; + } + self.emit_observer(kind, redact_wire_value(&payload)); + } + + /// Emit an event whose payload is already a logging-safe copy. + fn emit_observer(&self, kind: impl Into, payload: serde_json::Value) { if let Some(observer) = &self.observer { observer.emit( kind, @@ -1103,7 +1153,7 @@ impl AcpClient { .await .map_err(|_| AcpError::WriteTimeout(WRITE_TIMEOUT))? .map_err(AcpError::Io)?; - self.observe("acp_write", observed_value); + self.emit_observer("acp_write", observed_value); Ok(()) } @@ -1116,7 +1166,7 @@ impl AcpClient { Ok(msg) => { let observed_value = redact_wire_value(&msg); tracing::debug!(target: "acp::wire", "← {observed_value}"); - self.observe("acp_read", observed_value); + self.emit_observer("acp_read", observed_value); Some(msg) } Err(error) => { @@ -1769,6 +1819,7 @@ impl AcpClient { match update_type { "agent_message_chunk" => { if let Some(text) = update["content"]["text"].as_str() { + let text = redact_wire_text(text); tracing::info!(target: "acp::stream", "{text}"); } false @@ -1782,6 +1833,8 @@ impl AcpClient { .get("kind") .and_then(|v| v.as_str()) .unwrap_or("unknown"); + let title = redact_wire_text(title); + let kind = redact_wire_text(kind); tracing::info!(target: "acp::tool", "tool_call: {title} ({kind})"); true } @@ -1791,6 +1844,8 @@ impl AcpClient { .and_then(|v| v.as_str()) .unwrap_or("?"); let status = update.get("status").and_then(|v| v.as_str()).unwrap_or("?"); + let tool_id = redact_wire_text(tool_id); + let status = redact_wire_text(status); tracing::info!(target: "acp::tool", "tool_call_update: {tool_id} → {status}"); false } @@ -1800,6 +1855,7 @@ impl AcpClient { } "agent_thought_chunk" => { if let Some(text) = update["content"]["text"].as_str() { + let text = redact_wire_text(text); tracing::debug!(target: "acp::thought", "{text}"); } false @@ -1809,7 +1865,12 @@ impl AcpClient { // Logged for observability; UI surfacing is a follow-up. let names: Vec<&str> = update["availableCommands"] .as_array() - .map(|cmds| cmds.iter().filter_map(|c| c["name"].as_str()).collect()) + .map(|cmds| { + cmds.iter() + .filter_map(|c| c["name"].as_str()) + .map(redact_wire_text) + .collect() + }) .unwrap_or_default(); tracing::info!( target: "acp::update", @@ -1836,9 +1897,10 @@ impl AcpClient { if let Some(goose_meta) = meta { match goose_meta.get("activeRunId") { Some(serde_json::Value::String(run_id)) => { + let observed_run_id = redact_wire_text(run_id); tracing::debug!( target: "acp::update", - "session_info_update: activeRunId={run_id}" + "session_info_update: activeRunId={observed_run_id}" ); self.active_run_id = Some(run_id.clone()); } @@ -1857,6 +1919,7 @@ impl AcpClient { } "keepalive" => false, other => { + let other = redact_wire_text(other); tracing::debug!(target: "acp::update", "session/update: {other}"); false } @@ -1884,9 +1947,10 @@ impl AcpClient { match serde_json::from_value::(params.clone()) { Ok(notif) => { if let GooseSessionUpdateVariant::UsageUpdate(payload) = ¬if.update { + let observed_session_id = redact_wire_text(¬if.session_id); tracing::debug!( target: "acp::usage", - session_id = %notif.session_id, + session_id = %observed_session_id, input = payload.accumulated_input_tokens, output = payload.accumulated_output_tokens, // A subset of `input`, logged so downstream accounting can @@ -1900,9 +1964,11 @@ impl AcpClient { } } Err(e) => { + let error = e.to_string(); + let observed_error = redact_wire_text(&error); tracing::debug!( target: "acp::usage", - "_goose/unstable/session/update: deserialization error: {e}" + "_goose/unstable/session/update: deserialization error: {observed_error}" ); } } @@ -2540,6 +2606,58 @@ mod tests { ); } + #[test] + fn wire_redaction_suppresses_serialized_mcp_configs_without_broad_string_matching() { + let secret = "quote\" slash\\ newline\n snowman \u{2603}"; + let request = serde_json::json!({ + "jsonrpc": "2.0", + "id": 1, + "method": "session/new", + "params": { + "mcpServers": [{ + "name": "analytics", + "command": "analytics-mcp", + "args": [], + "env": [{"name": "ANALYTICS_TOKEN", "value": secret}] + }] + } + }) + .to_string(); + let mut serialized_levels = vec![request]; + for _ in 0..3 { + let next = serde_json::to_string( + serialized_levels + .last() + .expect("at least one serialized request"), + ) + .expect("serialize request again"); + serialized_levels.push(next); + } + let source = serde_json::json!({ + "echoes": serialized_levels + .iter() + .map(|request| format!("adapter rejected {request}; check configuration")) + .collect::>(), + "ordinary": r#"invalid {"environment":"prod","value":"x"}"#, + "nearMiss": r#"invalid {"mcpServers":[],"envValue":"prod","value":"x"}"#, + "unrelated": "ordinary adapter error" + }); + let original_source = source.clone(); + + let redacted = redact_wire_value(&source); + + for echo in redacted["echoes"].as_array().expect("redacted echoes") { + assert_eq!(echo, REDACTED_ENV_VALUE); + } + assert_eq!(redacted["ordinary"], source["ordinary"]); + assert_eq!(redacted["nearMiss"], source["nearMiss"]); + assert_eq!(redacted["unrelated"], source["unrelated"]); + assert_eq!( + source, original_source, + "the protocol value must stay unchanged" + ); + } + #[test] fn session_prompt_request_format() { let prompt_text = "[Buzz @mention]\nChannel: test\nFrom: npub1...\nMessage: hello"; @@ -3544,6 +3662,181 @@ mod tests { client.shutdown().await; } + #[tokio::test] + async fn session_new_error_cannot_echo_serialized_mcp_config() { + let secret = "adapter-echo-secret"; + let embedded_request = serde_json::json!({ + "jsonrpc": "2.0", + "id": 1, + "method": "session/new", + "params": { + "mcpServers": [{ + "name": "analytics", + "command": "analytics-mcp", + "args": [], + "env": [{"name": "ANALYTICS_TOKEN", "value": secret}] + }] + } + }) + .to_string(); + let error_response = serde_json::json!({ + "jsonrpc": "2.0", + "id": 1, + "error": { + "code": -32055, + "message": format!("adapter rejected {embedded_request}") + } + }) + .to_string(); + let script = format!( + "read -t 2 _init\n\ + printf '%s\\n' '{{\"jsonrpc\":\"2.0\",\"id\":0,\"result\":{{\"protocolVersion\":2,\"agentCapabilities\":{{}}}}}}'\n\ + read -t 2 _session\n\ + printf '%s\\n' '{error_response}'\n\ + sleep 1" + ); + let mut client = spawn_script(&script).await; + let observer = ObserverHandle::in_process(); + client.set_observer(Some(observer.clone()), 0); + client + .initialize() + .await + .expect("initialize should succeed"); + + let result = client + .session_new_full( + "/tmp", + vec![McpServer { + name: "analytics".into(), + command: "analytics-mcp".into(), + args: vec![], + env: vec![EnvVar { + name: "ANALYTICS_TOKEN".into(), + value: secret.into(), + }], + }], + None, + None, + ) + .await; + + match result { + Err(AcpError::AgentError { code, message }) => { + assert_eq!(code, -32055); + assert_eq!(message, REDACTED_ENV_VALUE); + assert!(!message.contains(secret)); + } + Err(other) => panic!("expected redacted AgentError, got {other:?}"), + Ok(_) => panic!("expected session/new to return an error"), + } + + client.observe( + "adapter_diagnostic", + serde_json::json!({"message": embedded_request}), + ); + let serialized_events = + serde_json::to_string(&observer.snapshot()).expect("serialize observer snapshot"); + assert!(!serialized_events.contains(secret)); + assert!(serialized_events.contains(REDACTED_ENV_VALUE)); + client.shutdown().await; + } + + #[tokio::test] + async fn semantic_traces_cannot_echo_serialized_mcp_config() { + use std::io::Write; + use std::sync::{Arc, Mutex}; + + #[derive(Clone, Default)] + struct TraceCapture(Arc>>); + + struct TraceWriter(Arc>>); + + impl Write for TraceWriter { + fn write(&mut self, bytes: &[u8]) -> std::io::Result { + self.0 + .lock() + .expect("trace buffer lock") + .extend_from_slice(bytes); + Ok(bytes.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + + impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for TraceCapture { + type Writer = TraceWriter; + + fn make_writer(&'a self) -> Self::Writer { + TraceWriter(self.0.clone()) + } + } + + let secret = "semantic-trace-secret"; + let embedded_request = serde_json::json!({ + "method": "session/new", + "params": { + "mcpServers": [{ + "env": [{"name": "ANALYTICS_TOKEN", "value": secret}] + }] + } + }) + .to_string(); + let update = serde_json::json!({ + "jsonrpc": "2.0", + "method": "session/update", + "params": { + "update": { + "sessionUpdate": "agent_message_chunk", + "content": {"type": "text", "text": embedded_request} + } + } + }); + let mut client = spawn_inert_client().await; + let trace = TraceCapture::default(); + let subscriber = tracing_subscriber::fmt() + .without_time() + .with_ansi(false) + .with_max_level(tracing::Level::DEBUG) + .with_writer(trace.clone()) + .finish(); + + tracing::subscriber::with_default(subscriber, || { + let _ = client.handle_session_update(&update); + client.handle_goose_usage_update(&serde_json::json!({ + "params": { + "sessionId": embedded_request, + "update": { + "sessionUpdate": "usage_update", + "accumulatedInputTokens": 10, + "accumulatedOutputTokens": 5, + "accumulatedCachedInputTokens": null, + "accumulatedCost": null + } + } + })); + client.handle_goose_usage_update(&serde_json::json!({ + "params": { + "sessionId": "ordinary-session", + "update": { + "sessionUpdate": "usage_update", + "accumulatedInputTokens": embedded_request, + "accumulatedOutputTokens": 5, + "accumulatedCachedInputTokens": null, + "accumulatedCost": null + } + } + })); + }); + + let output = String::from_utf8(trace.0.lock().expect("trace buffer lock").clone()) + .expect("trace output should be UTF-8"); + assert!(!output.contains(secret)); + assert!(output.contains(REDACTED_ENV_VALUE)); + client.shutdown().await; + } + #[tokio::test] async fn structured_mcp_servers_survive_repeated_sessions_and_adapter_restart() { let script = r#" From a94b456cebbe679bcc2e0ed30a931d5f4c2fa3fa Mon Sep 17 00:00:00 2001 From: KC <79471844+wolfyy970@users.noreply.github.com> Date: Sun, 2 Aug 2026 16:47:34 -0400 Subject: [PATCH 3/6] refactor(acp): tag structured MCP transports Signed-off-by: KC <79471844+wolfyy970@users.noreply.github.com> --- crates/buzz-acp/src/config.rs | 135 +++++++++++++++++++++------------- crates/buzz-acp/src/lib.rs | 34 +++++---- 2 files changed, 105 insertions(+), 64 deletions(-) diff --git a/crates/buzz-acp/src/config.rs b/crates/buzz-acp/src/config.rs index f23ddeca06..f268dc842e 100644 --- a/crates/buzz-acp/src/config.rs +++ b/crates/buzz-acp/src/config.rs @@ -63,19 +63,26 @@ const PROTECTED_MCP_ENV_NAMES: [&str; 6] = [ "BUZZ_ACP_API_TOKEN", ]; -/// One local stdio MCP server loaded from the structured MCP configuration. +/// One MCP server loaded from the structured MCP configuration. +/// +/// The transport tag is part of the version-1 document even though this PR +/// implements only stdio. Additional transports can extend the same ordered +/// server list without introducing a parallel configuration format. #[derive(Debug, Clone, PartialEq, Eq, serde::Deserialize)] -#[serde(deny_unknown_fields)] -pub struct ConfiguredMcpServer { - /// Stable ACP identifier for this server. - pub name: String, - /// Executable to invoke, passed directly without shell parsing. - pub command: String, - /// Arguments passed to the executable in their configured order. - pub args: Vec, - /// Server-specific environment in deterministic key order. - #[serde(deserialize_with = "deserialize_mcp_env")] - pub env: BTreeMap, +#[serde(tag = "transport", rename_all = "snake_case", deny_unknown_fields)] +pub enum ConfiguredMcpServer { + /// A local MCP child process connected over stdio. + Stdio { + /// Stable ACP identifier for this server. + name: String, + /// Executable to invoke, passed directly without shell parsing. + command: String, + /// Arguments passed to the executable in their configured order. + args: Vec, + /// Server-specific environment in deterministic key order. + #[serde(deserialize_with = "deserialize_mcp_env")] + env: BTreeMap, + }, } #[derive(Debug, serde::Deserialize)] @@ -219,56 +226,62 @@ fn load_mcp_config( (!legacy_mcp_command.is_empty()).then(|| legacy_mcp_server_name(legacy_mcp_command)); let mut names = HashSet::with_capacity(document.servers.len()); for (index, server) in document.servers.iter().enumerate() { - if !valid_mcp_server_name(&server.name) { + let ConfiguredMcpServer::Stdio { + name, + command, + args, + env, + } = server; + if !valid_mcp_server_name(name) { return Err(ConfigError::ConfigFile(format!( "MCP server {} has invalid name '{}': use 1 to {MCP_SERVER_NAME_MAX_BYTES} ASCII letters, digits, underscores, or hyphens, without '__'", index + 1, - server.name + name ))); } - if !names.insert(server.name.as_str()) { + if !names.insert(name.as_str()) { return Err(ConfigError::ConfigFile(format!( "duplicate MCP server name '{}'", - server.name + name ))); } - if legacy_name.as_deref() == Some(server.name.as_str()) { + if legacy_name.as_deref() == Some(name.as_str()) { return Err(ConfigError::ConfigFile(format!( "MCP server name '{}' collides with the legacy --mcp-command server", - server.name + name ))); } - if server.command.is_empty() || server.command.contains('\0') { + if command.is_empty() || command.contains('\0') { return Err(ConfigError::ConfigFile(format!( "MCP server '{}' command must be nonempty and contain no NUL bytes", - server.name + name ))); } - if server.args.len() > MCP_SERVER_MAX_ARGS { + if args.len() > MCP_SERVER_MAX_ARGS { return Err(ConfigError::ConfigFile(format!( "MCP server '{}' has too many arguments ({}, max {MCP_SERVER_MAX_ARGS})", - server.name, - server.args.len() + name, + args.len() ))); } - if server.args.iter().any(|argument| argument.contains('\0')) { + if args.iter().any(|argument| argument.contains('\0')) { return Err(ConfigError::ConfigFile(format!( "MCP server '{}' arguments must contain no NUL bytes", - server.name + name ))); } - if server.env.len() > MCP_SERVER_MAX_ENV { + if env.len() > MCP_SERVER_MAX_ENV { return Err(ConfigError::ConfigFile(format!( "MCP server '{}' has too many environment entries ({}, max {MCP_SERVER_MAX_ENV})", - server.name, - server.env.len() + name, + env.len() ))); } - for (key, value) in &server.env { + for (key, value) in env { if !valid_mcp_env_name(key) { return Err(ConfigError::ConfigFile(format!( "MCP server '{}' has invalid environment key '{key}'", - server.name + name ))); } if PROTECTED_MCP_ENV_NAMES @@ -277,13 +290,13 @@ fn load_mcp_config( { return Err(ConfigError::ConfigFile(format!( "MCP server '{}' may not configure protected environment key '{key}'", - server.name + name ))); } if value.contains('\0') { return Err(ConfigError::ConfigFile(format!( "MCP server '{}' environment value for '{key}' contains a NUL byte", - server.name + name ))); } } @@ -3236,6 +3249,7 @@ channels = "ALL" fn server_json(name: &str) -> serde_json::Value { serde_json::json!({ "name": name, + "transport": "stdio", "command": "mcp", "args": [], "env": {} @@ -3258,6 +3272,7 @@ channels = "ALL" let file = TempMcpConfig::write(&document_json(vec![ serde_json::json!({ "name": "analytics-primary", + "transport": "stdio", "command": posix_command, "args": [ "", @@ -3275,6 +3290,7 @@ channels = "ALL" }), serde_json::json!({ "name": "windows_server", + "transport": "stdio", "command": windows_command, "args": ["--stdio"], "env": {} @@ -3283,11 +3299,17 @@ channels = "ALL" let config = config_from_mcp_file(&file, None).expect("structured config should load"); assert_eq!(config.configured_mcp_servers.len(), 2); - assert_eq!(config.configured_mcp_servers[0].name, "analytics-primary"); - assert_eq!(config.configured_mcp_servers[0].command, posix_command); + let ConfiguredMcpServer::Stdio { + name, + command, + args, + env, + } = &config.configured_mcp_servers[0]; + assert_eq!(name, "analytics-primary"); + assert_eq!(command, posix_command); assert_eq!( - config.configured_mcp_servers[0].args, - vec![ + args, + &vec![ "", "with spaces", "comma,value", @@ -3298,18 +3320,12 @@ channels = "ALL" ] ); assert_eq!( - config.configured_mcp_servers[0] - .env - .keys() - .map(String::as_str) - .collect::>(), + env.keys().map(String::as_str).collect::>(), vec!["A_FIRST", "Z_LAST"] ); - assert_eq!( - config.configured_mcp_servers[0].env["Z_LAST"], - "backslash\\quote\"雪" - ); - assert_eq!(config.configured_mcp_servers[1].command, windows_command); + assert_eq!(env["Z_LAST"], "backslash\\quote\"雪"); + let ConfiguredMcpServer::Stdio { command, .. } = &config.configured_mcp_servers[1]; + assert_eq!(command, windows_command); } #[test] @@ -3318,6 +3334,7 @@ channels = "ALL" let command = "/private/path/tool"; let file = TempMcpConfig::write(&document_json(vec![serde_json::json!({ "name": "safe", + "transport": "stdio", "command": command, "args": [], "env": {"DOMAIN_TOKEN": secret_value} @@ -3402,11 +3419,19 @@ channels = "ALL" ), ( "unknown server field", - br#"{"version":1,"servers":[{"name":"one","command":"mcp","args":[],"env":{},"extra":true}]}"#.to_vec(), + br#"{"version":1,"servers":[{"name":"one","transport":"stdio","command":"mcp","args":[],"env":{},"extra":true}]}"#.to_vec(), ), ( "missing required field", - br#"{"version":1,"servers":[{"name":"one","command":"mcp","env":{}}]}"#.to_vec(), + br#"{"version":1,"servers":[{"name":"one","transport":"stdio","command":"mcp","env":{}}]}"#.to_vec(), + ), + ( + "missing transport", + br#"{"version":1,"servers":[{"name":"one","command":"mcp","args":[],"env":{}}]}"#.to_vec(), + ), + ( + "unsupported transport", + br#"{"version":1,"servers":[{"name":"one","transport":"http","url":"https://example.test/mcp","headers":{}}]}"#.to_vec(), ), ]; @@ -3480,6 +3505,7 @@ channels = "ALL" .collect::>(); let file = TempMcpConfig::write(&document_json(vec![serde_json::json!({ "name": "limit", + "transport": "stdio", "command": "mcp", "args": args, "env": {} @@ -3491,6 +3517,7 @@ channels = "ALL" .collect::>(); let file = TempMcpConfig::write(&document_json(vec![serde_json::json!({ "name": "over-limit", + "transport": "stdio", "command": "mcp", "args": args, "env": {} @@ -3507,6 +3534,7 @@ channels = "ALL" .collect::>(); let file = TempMcpConfig::write(&document_json(vec![serde_json::json!({ "name": "limit", + "transport": "stdio", "command": "mcp", "args": [], "env": env @@ -3518,6 +3546,7 @@ channels = "ALL" .collect::>(); let file = TempMcpConfig::write(&document_json(vec![serde_json::json!({ "name": "over-limit", + "transport": "stdio", "command": "mcp", "args": [], "env": env @@ -3531,7 +3560,7 @@ channels = "ALL" fn mcp_config_rejects_invalid_duplicate_and_protected_env_names() { for invalid_key in ["", "1STARTS_WITH_DIGIT", "BAD-NAME", "UNICODÉ"] { let content = format!( - r#"{{"version":1,"servers":[{{"name":"one","command":"mcp","args":[],"env":{{"{invalid_key}":"value"}}}}]}}"# + r#"{{"version":1,"servers":[{{"name":"one","transport":"stdio","command":"mcp","args":[],"env":{{"{invalid_key}":"value"}}}}]}}"# ); let file = TempMcpConfig::write(content.as_bytes()); assert!( @@ -3543,7 +3572,7 @@ channels = "ALL" for protected in PROTECTED_MCP_ENV_NAMES { let lowercase = protected.to_ascii_lowercase(); let content = format!( - r#"{{"version":1,"servers":[{{"name":"one","command":"mcp","args":[],"env":{{"{lowercase}":"value"}}}}]}}"# + r#"{{"version":1,"servers":[{{"name":"one","transport":"stdio","command":"mcp","args":[],"env":{{"{lowercase}":"value"}}}}]}}"# ); let file = TempMcpConfig::write(content.as_bytes()); let error = config_from_mcp_file(&file, None) @@ -3556,7 +3585,7 @@ channels = "ALL" r#"{"KEY":"one","key":"two"}"#, ] { let content = format!( - r#"{{"version":1,"servers":[{{"name":"one","command":"mcp","args":[],"env":{duplicate_env}}}]}}"# + r#"{{"version":1,"servers":[{{"name":"one","transport":"stdio","command":"mcp","args":[],"env":{duplicate_env}}}]}}"# ); let file = TempMcpConfig::write(content.as_bytes()); let error = @@ -3570,24 +3599,28 @@ channels = "ALL" let cases = vec![ serde_json::json!({ "name": "empty-command", + "transport": "stdio", "command": "", "args": [], "env": {} }), serde_json::json!({ "name": "nul-command", + "transport": "stdio", "command": "mc\u{0}p", "args": [], "env": {} }), serde_json::json!({ "name": "nul-arg", + "transport": "stdio", "command": "mcp", "args": ["ok", "bad\u{0}arg"], "env": {} }), serde_json::json!({ "name": "nul-value", + "transport": "stdio", "command": "mcp", "args": [], "env": {"DOMAIN_KEY": "bad\u{0}value"} diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 7f12d2a125..8162927baa 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -4333,18 +4333,24 @@ fn build_mcp_servers(config: &Config) -> Vec { } servers.extend(config.configured_mcp_servers.iter().map(|configured| { - McpServer { - name: configured.name.clone(), - command: configured.command.clone(), - args: configured.args.clone(), - env: configured - .env - .iter() - .map(|(name, value)| EnvVar { - name: name.clone(), - value: value.clone(), - }) - .collect(), + match configured { + config::ConfiguredMcpServer::Stdio { + name, + command, + args, + env, + } => McpServer { + name: name.clone(), + command: command.clone(), + args: args.clone(), + env: env + .iter() + .map(|(name, value)| EnvVar { + name: name.clone(), + value: value.clone(), + }) + .collect(), + }, } })); @@ -5166,7 +5172,7 @@ mod build_mcp_servers_tests { args: &[&str], env: &[(&str, &str)], ) -> config::ConfiguredMcpServer { - config::ConfiguredMcpServer { + config::ConfiguredMcpServer::Stdio { name: name.into(), command: command.into(), args: args.iter().map(|arg| (*arg).to_string()).collect(), @@ -5415,6 +5421,7 @@ mod build_mcp_servers_tests { "servers": [ { "name": "analytics", + "transport": "stdio", "command": "/opt/MCP Servers/analytics,prod", "args": ["--stdio", "literal value"], "env": { @@ -5423,6 +5430,7 @@ mod build_mcp_servers_tests { }, { "name": "search", + "transport": "stdio", "command": "/opt/search-mcp", "args": [], "env": {} From 693e20d0f345d99b5e69638d1a9bb255459e0a1a Mon Sep 17 00:00:00 2001 From: KC <79471844+wolfyy970@users.noreply.github.com> Date: Mon, 3 Aug 2026 10:52:42 -0400 Subject: [PATCH 4/6] docs(acp): include MCP transport in example Signed-off-by: KC <79471844+wolfyy970@users.noreply.github.com> --- crates/buzz-acp/README.md | 1 + 1 file changed, 1 insertion(+) diff --git a/crates/buzz-acp/README.md b/crates/buzz-acp/README.md index d5384df7e2..d672d94b66 100644 --- a/crates/buzz-acp/README.md +++ b/crates/buzz-acp/README.md @@ -131,6 +131,7 @@ servers: "servers": [ { "name": "analytics", + "transport": "stdio", "command": "/opt/mcp/analytics-server", "args": ["--stdio"], "env": { From cd13211788d90bbb88727f7112ec09eb648e81e9 Mon Sep 17 00:00:00 2001 From: KC <79471844+wolfyy970@users.noreply.github.com> Date: Mon, 3 Aug 2026 11:34:39 -0400 Subject: [PATCH 5/6] fix(acp): redact MCP credentials from diagnostics Signed-off-by: KC <79471844+wolfyy970@users.noreply.github.com> --- Cargo.lock | 1 + crates/buzz-acp/Cargo.toml | 1 + crates/buzz-acp/src/acp.rs | 390 ++++++++++++++++++++++++++++------ crates/buzz-acp/src/config.rs | 54 ++++- 4 files changed, 383 insertions(+), 63 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 937ead564a..d75c96bd87 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -800,6 +800,7 @@ checksum = "5d20789868f4b01b2f2caec9f5c4e0213b41e3e5702a50157d699ae31ced2fcb" name = "buzz-acp" version = "0.1.0" dependencies = [ + "aho-corasick", "anyhow", "base64 0.22.1", "buzz-core", diff --git a/crates/buzz-acp/Cargo.toml b/crates/buzz-acp/Cargo.toml index d047849806..37c6a4146d 100644 --- a/crates/buzz-acp/Cargo.toml +++ b/crates/buzz-acp/Cargo.toml @@ -41,6 +41,7 @@ reqwest = { workspace = true } # Serialization serde = { workspace = true } serde_json = { workspace = true } +aho-corasick = "1.1" # IDs uuid = { workspace = true } diff --git a/crates/buzz-acp/src/acp.rs b/crates/buzz-acp/src/acp.rs index 45b7f3d597..8bee669345 100644 --- a/crates/buzz-acp/src/acp.rs +++ b/crates/buzz-acp/src/acp.rs @@ -8,7 +8,9 @@ //! 4. [`AcpClient::session_prompt_with_idle_timeout`] — send prompt with idle/hard deadline, return stop reason //! 5. [`AcpClient::session_cancel`] / [`AcpClient::cancel_with_cleanup`] — cancel in-flight turn +use aho_corasick::AhoCorasick; use futures_util::StreamExt; +use std::borrow::Cow; use tokio::io::AsyncWriteExt; use tokio::process::{Child, ChildStdin, ChildStdout}; use tokio_util::codec::{FramedRead, LinesCodec, LinesCodecError}; @@ -26,7 +28,7 @@ const REDACTED_ENV_VALUE: &str = "[REDACTED]"; /// /// Corresponds to the `McpServerStdio` variant in the ACP schema. /// All four fields are **required** by the schema (`args` and `env` may be empty arrays). -#[derive(Debug, Clone, serde::Serialize)] +#[derive(Clone, serde::Serialize)] pub struct McpServer { pub name: String, pub command: String, @@ -34,13 +36,35 @@ pub struct McpServer { pub env: Vec, } +impl std::fmt::Debug for McpServer { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter + .debug_struct("McpServer") + .field("name", &self.name) + .field("command", &self.command) + .field("arg_count", &self.args.len()) + .field("env", &self.env) + .finish() + } +} + /// A single environment variable for an MCP server. -#[derive(Debug, Clone, serde::Serialize)] +#[derive(Clone, serde::Serialize)] pub struct EnvVar { pub name: String, pub value: String, } +impl std::fmt::Debug for EnvVar { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter + .debug_struct("EnvVar") + .field("name", &self.name) + .field("value", &REDACTED_ENV_VALUE) + .finish() + } +} + /// Stop reason returned by `session/prompt` when the agent finishes a turn. /// /// Maps to the `stopReason` field in the `SessionPromptResponse`. @@ -114,9 +138,12 @@ pub enum AcpError { /// preserving the numeric code. When the `message` field is missing or /// non-string, fall back to the full JSON object so provider-specific /// detail (e.g. a `data` field) is not lost. -fn agent_error_from_json(error: &serde_json::Value) -> AcpError { +fn agent_error_from_json( + error: &serde_json::Value, + sensitive_value_matcher: Option<&AhoCorasick>, +) -> AcpError { let code = error.get("code").and_then(|c| c.as_i64()).unwrap_or(-32000); - let redacted_error = redact_wire_value(error); + let redacted_error = redact_wire_value(error, sensitive_value_matcher); let message = match redacted_error.get("message").and_then(|m| m.as_str()) { Some(m) => m.to_string(), None => redacted_error.to_string(), @@ -149,27 +176,45 @@ fn contains_serialized_json_key(text: &str, key: &str) -> bool { }) } -fn redact_wire_text(text: &str) -> &str { +fn redact_wire_text<'a>( + text: &'a str, + sensitive_value_matcher: Option<&AhoCorasick>, +) -> Cow<'a, str> { let looks_like_serialized_mcp_config = contains_serialized_json_key(text, "mcpServers") && contains_serialized_json_key(text, "env") && contains_serialized_json_key(text, "value"); if looks_like_serialized_mcp_config { - REDACTED_ENV_VALUE - } else { - text + return Cow::Borrowed(REDACTED_ENV_VALUE); + } + + let Some(matcher) = sensitive_value_matcher else { + return Cow::Borrowed(text); + }; + if matcher.find(text).is_some() { + // Suppress the whole diagnostic. Partial replacement can expand a + // bounded 10 MB adapter line many times over when a configured value + // is short or common. + return Cow::Borrowed(REDACTED_ENV_VALUE); } + Cow::Borrowed(text) } /// Return a logging-safe copy of an ACP wire value. /// /// MCP environment values are needed by the adapter on the real wire, but /// must not reach tracing or observer frames. -fn redact_wire_value(value: &serde_json::Value) -> serde_json::Value { - fn redact_in_place(value: &mut serde_json::Value) { +fn redact_wire_value( + value: &serde_json::Value, + sensitive_value_matcher: Option<&AhoCorasick>, +) -> serde_json::Value { + fn redact_in_place( + value: &mut serde_json::Value, + sensitive_value_matcher: Option<&AhoCorasick>, + ) { match value { serde_json::Value::Array(values) => { for value in values { - redact_in_place(value); + redact_in_place(value, sensitive_value_matcher); } } serde_json::Value::Object(fields) => { @@ -190,13 +235,13 @@ fn redact_wire_value(value: &serde_json::Value) -> serde_json::Value { } } } - redact_in_place(value); + redact_in_place(value, sensitive_value_matcher); } } serde_json::Value::String(text) => { - let observed = redact_wire_text(text); + let observed = redact_wire_text(text, sensitive_value_matcher); if observed != text.as_str() { - *text = observed.to_string(); + *text = observed.into_owned(); } } _ => {} @@ -204,7 +249,7 @@ fn redact_wire_value(value: &serde_json::Value) -> serde_json::Value { } let mut redacted = value.clone(); - redact_in_place(&mut redacted); + redact_in_place(&mut redacted, sensitive_value_matcher); redacted } @@ -259,6 +304,13 @@ pub struct AcpClient { observer_agent_index: Option, /// Best-effort context attached to raw ACP wire events. observer_context: ObserverContext, + /// Non-empty MCP environment values ever sent to this adapter. + /// + /// Values remain registered across sessions so a delayed adapter diagnostic + /// cannot disclose a credential from an earlier `session/new`. + sensitive_mcp_env_values: Vec, + /// Single-pass matcher rebuilt only when a session introduces a new value. + sensitive_mcp_env_matcher: Option, /// Most recently observed `_meta.goose.activeRunId` from a /// `session/update` notification of kind `session_info_update`. /// @@ -633,6 +685,8 @@ impl AcpClient { observer: None, observer_agent_index: None, observer_context: ObserverContext::default(), + sensitive_mcp_env_values: Vec::new(), + sensitive_mcp_env_matcher: None, active_run_id: None, steering_supported: false, steer_rx: None, @@ -661,12 +715,54 @@ impl AcpClient { self.observer_agent_index } + fn register_sensitive_mcp_env_values(&mut self, servers: &[McpServer]) -> Result<(), AcpError> { + let mut changed = false; + for value in servers + .iter() + .flat_map(|server| server.env.iter().map(|entry| &entry.value)) + .filter(|value| !value.is_empty()) + { + if !self + .sensitive_mcp_env_values + .iter() + .any(|registered| registered == value) + { + self.sensitive_mcp_env_values.push(value.clone()); + changed = true; + } + } + if !changed { + return Ok(()); + } + + self.sensitive_mcp_env_values + .sort_by(|left, right| right.len().cmp(&left.len()).then_with(|| left.cmp(right))); + self.sensitive_mcp_env_matcher = Some( + AhoCorasick::new(&self.sensitive_mcp_env_values).map_err(|_| { + AcpError::Protocol("failed to initialize MCP credential redaction".to_string()) + })?, + ); + Ok(()) + } + + fn redact_wire_text<'a>(&self, text: &'a str) -> Cow<'a, str> { + redact_wire_text(text, self.sensitive_mcp_env_matcher.as_ref()) + } + + fn redact_wire_value(&self, value: &serde_json::Value) -> serde_json::Value { + redact_wire_value(value, self.sensitive_mcp_env_matcher.as_ref()) + } + + fn agent_error_from_json(&self, error: &serde_json::Value) -> AcpError { + agent_error_from_json(error, self.sensitive_mcp_env_matcher.as_ref()) + } + /// Emit a semantic event to the local observer feed, if enabled. pub fn observe(&self, kind: impl Into, payload: serde_json::Value) { if self.observer.is_none() { return; } - self.emit_observer(kind, redact_wire_value(&payload)); + self.emit_observer(kind, self.redact_wire_value(&payload)); } /// Emit an event whose payload is already a logging-safe copy. @@ -699,7 +795,8 @@ impl AcpClient { .pointer("/_meta/steering/supported") .and_then(|v| v.as_bool()) .unwrap_or(false); - tracing::debug!(target: "acp::init", "initialize response: {result}"); + let observed_result = self.redact_wire_value(&result); + tracing::debug!(target: "acp::init", "initialize response: {observed_result}"); Ok(result) } @@ -737,6 +834,7 @@ impl AcpClient { system_prompt: Option>, session_title: Option<&str>, ) -> Result { + self.register_sensitive_mcp_env_values(&mcp_servers)?; let mut params = serde_json::json!({ "cwd": cwd, "mcpServers": mcp_servers, @@ -760,7 +858,8 @@ impl AcpClient { .as_str() .ok_or_else(|| AcpError::Protocol("session/new response missing sessionId".into()))? .to_owned(); - tracing::info!(target: "acp::session", "session created: {session_id}"); + let observed_session_id = self.redact_wire_text(&session_id); + tracing::info!(target: "acp::session", "session created: {observed_session_id}"); Ok(SessionNewResponse { session_id, raw: result, @@ -1103,9 +1202,10 @@ impl AcpClient { if !self.permission_responded { let response = permission_response_cancelled(&perm_id); self.write_ndjson(&response).await?; + let observed_permission_id = self.redact_wire_value(&perm_id); tracing::debug!( target: "acp::cancel", - "responded cancelled to pending permission id={perm_id}" + "responded cancelled to pending permission id={observed_permission_id}" ); } self.pending_permission_id = None; @@ -1114,7 +1214,8 @@ impl AcpClient { // Step 2: send session/cancel notification (no id) self.session_cancel(session_id).await?; - tracing::info!(target: "acp::cancel", "sent session/cancel for {session_id}"); + let observed_session_id = self.redact_wire_text(session_id); + tracing::info!(target: "acp::cancel", "sent session/cancel for {observed_session_id}"); // Use a fixed 30s idle timeout during cleanup — the cancel notification // needs time to propagate and the agent may go silent while winding down. // The separate hard_deadline bounds agents that keep producing output @@ -1141,7 +1242,7 @@ impl AcpClient { /// (e.g., it's stuck or dead), the write would otherwise block forever. async fn write_ndjson(&mut self, value: &serde_json::Value) -> Result<(), AcpError> { const WRITE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30); - let observed_value = redact_wire_value(value); + let observed_value = self.redact_wire_value(value); tracing::debug!(target: "acp::wire", "→ {observed_value}"); let line = serde_json::to_string(value)?; tokio::time::timeout(WRITE_TIMEOUT, async { @@ -1164,7 +1265,7 @@ impl AcpClient { fn parse_inbound_line(&self, line: &str) -> Option { match serde_json::from_str(line) { Ok(msg) => { - let observed_value = redact_wire_value(&msg); + let observed_value = self.redact_wire_value(&msg); tracing::debug!(target: "acp::wire", "← {observed_value}"); self.emit_observer("acp_read", observed_value); Some(msg) @@ -1331,7 +1432,7 @@ impl AcpClient { if let Some(id) = msg.get("id") { if *id == serde_json::json!(expected_id) && msg.get("method").is_none() { if let Some(error) = msg.get("error") { - return Err(agent_error_from_json(error)); + return Err(self.agent_error_from_json(error)); } return Ok(msg["result"].clone()); } @@ -1363,7 +1464,11 @@ impl AcpClient { // agent process is dead and continuing would hang. self.write_ndjson(&err_resp).await?; } - tracing::debug!(target: "acp::wire", "ignoring unknown method: {other}"); + let observed_method = self.redact_wire_text(other); + tracing::debug!( + target: "acp::wire", + "ignoring unknown method: {observed_method}" + ); } } } @@ -1651,7 +1756,7 @@ impl AcpClient { .get("code") .and_then(|c| c.as_i64()) .unwrap_or(-1); - let message = redact_wire_value(error).to_string(); + let message = self.redact_wire_value(error).to_string(); crate::pool::SteerAck::Err( crate::pool::SteerError::AgentError { code, message }, ) @@ -1721,6 +1826,8 @@ impl AcpClient { Some(serde_json::Value::String(s)) => s.clone(), Some(other) => other.to_string(), }; + let reported = + self.redact_wire_text(&reported).into_owned(); tracing::warn!( "steer rejected: {ACP_STEER_METHOD} returned \ unrecognized outcome {reported} — releasing \ @@ -1744,7 +1851,7 @@ impl AcpClient { let _ = ack_tx .send(crate::pool::SteerAck::PromptCompletedNeutral); } - return Err(agent_error_from_json(error)); + return Err(self.agent_error_from_json(error)); } if let Some((_, _, ack_tx)) = pending_steer.take() { let _ = @@ -1786,7 +1893,11 @@ impl AcpClient { // agent process is dead and continuing would hang. self.write_ndjson(&err_resp).await?; } - tracing::debug!(target: "acp::wire", "ignoring unknown method: {other}"); + let observed_method = self.redact_wire_text(other); + tracing::debug!( + target: "acp::wire", + "ignoring unknown method: {observed_method}" + ); } } } @@ -1819,7 +1930,7 @@ impl AcpClient { match update_type { "agent_message_chunk" => { if let Some(text) = update["content"]["text"].as_str() { - let text = redact_wire_text(text); + let text = self.redact_wire_text(text); tracing::info!(target: "acp::stream", "{text}"); } false @@ -1833,8 +1944,8 @@ impl AcpClient { .get("kind") .and_then(|v| v.as_str()) .unwrap_or("unknown"); - let title = redact_wire_text(title); - let kind = redact_wire_text(kind); + let title = self.redact_wire_text(title); + let kind = self.redact_wire_text(kind); tracing::info!(target: "acp::tool", "tool_call: {title} ({kind})"); true } @@ -1844,8 +1955,8 @@ impl AcpClient { .and_then(|v| v.as_str()) .unwrap_or("?"); let status = update.get("status").and_then(|v| v.as_str()).unwrap_or("?"); - let tool_id = redact_wire_text(tool_id); - let status = redact_wire_text(status); + let tool_id = self.redact_wire_text(tool_id); + let status = self.redact_wire_text(status); tracing::info!(target: "acp::tool", "tool_call_update: {tool_id} → {status}"); false } @@ -1855,7 +1966,7 @@ impl AcpClient { } "agent_thought_chunk" => { if let Some(text) = update["content"]["text"].as_str() { - let text = redact_wire_text(text); + let text = self.redact_wire_text(text); tracing::debug!(target: "acp::thought", "{text}"); } false @@ -1863,12 +1974,12 @@ impl AcpClient { "available_commands_update" => { // Advertised slash commands (ACP slash-commands extension). // Logged for observability; UI surfacing is a follow-up. - let names: Vec<&str> = update["availableCommands"] + let names: Vec> = update["availableCommands"] .as_array() .map(|cmds| { cmds.iter() .filter_map(|c| c["name"].as_str()) - .map(redact_wire_text) + .map(|name| self.redact_wire_text(name)) .collect() }) .unwrap_or_default(); @@ -1897,7 +2008,7 @@ impl AcpClient { if let Some(goose_meta) = meta { match goose_meta.get("activeRunId") { Some(serde_json::Value::String(run_id)) => { - let observed_run_id = redact_wire_text(run_id); + let observed_run_id = self.redact_wire_text(run_id); tracing::debug!( target: "acp::update", "session_info_update: activeRunId={observed_run_id}" @@ -1919,7 +2030,7 @@ impl AcpClient { } "keepalive" => false, other => { - let other = redact_wire_text(other); + let other = self.redact_wire_text(other); tracing::debug!(target: "acp::update", "session/update: {other}"); false } @@ -1947,7 +2058,7 @@ impl AcpClient { match serde_json::from_value::(params.clone()) { Ok(notif) => { if let GooseSessionUpdateVariant::UsageUpdate(payload) = ¬if.update { - let observed_session_id = redact_wire_text(¬if.session_id); + let observed_session_id = self.redact_wire_text(¬if.session_id); tracing::debug!( target: "acp::usage", session_id = %observed_session_id, @@ -1965,7 +2076,7 @@ impl AcpClient { } Err(e) => { let error = e.to_string(); - let observed_error = redact_wire_text(&error); + let observed_error = self.redact_wire_text(&error); tracing::debug!( target: "acp::usage", "_goose/unstable/session/update: deserialization error: {observed_error}" @@ -1999,9 +2110,10 @@ impl AcpClient { .as_array() .ok_or_else(|| AcpError::Protocol("permission request missing options".into()))?; + let observed_id = self.redact_wire_value(&id); tracing::debug!( target: "acp::permission", - "session/request_permission id={id}, {} options", + "session/request_permission id={observed_id}, {} options", options.len() ); @@ -2014,16 +2126,17 @@ impl AcpClient { let option_id = opt["optionId"] .as_str() .ok_or_else(|| AcpError::Protocol("allow_once option missing optionId".into()))?; + let observed_option_id = self.redact_wire_text(option_id); tracing::info!( target: "acp::permission", - "auto-approving permission id={id} with allow_once optionId={option_id:?}" + "auto-approving permission id={observed_id} with allow_once optionId={observed_option_id:?}" ); permission_response_selected(&id, option_id) } else { // No allow_once — fall back to reject_once. tracing::warn!( target: "acp::permission", - "no allow_once option found in permission request id={id}, falling back to reject_once" + "no allow_once option found in permission request id={observed_id}, falling back to reject_once" ); let reject = options .iter() @@ -2065,8 +2178,10 @@ impl AcpClient { let raw = result["stopReason"].as_str().ok_or_else(|| { AcpError::Protocol("session/prompt response missing stopReason".into()) })?; - StopReason::from_str(raw) - .ok_or_else(|| AcpError::Protocol(format!("unknown stopReason: {raw:?}"))) + StopReason::from_str(raw).ok_or_else(|| { + let observed = self.redact_wire_text(raw); + AcpError::Protocol(format!("unknown stopReason: {observed:?}")) + }) } } @@ -2360,6 +2475,17 @@ fn configure_no_window(cmd: &mut tokio::process::Command) { mod tests { use super::*; + fn sensitive_matcher(values: &[String]) -> Option { + let mut patterns = values + .iter() + .map(String::as_str) + .filter(|value| !value.is_empty()) + .collect::>(); + patterns.sort_by(|left, right| right.len().cmp(&left.len()).then_with(|| left.cmp(right))); + patterns.dedup(); + (!patterns.is_empty()).then(|| AhoCorasick::new(patterns).unwrap()) + } + #[test] fn stop_reason_parses_all_known_values() { assert_eq!(StopReason::from_str("end_turn"), Some(StopReason::EndTurn)); @@ -2563,6 +2689,30 @@ mod tests { ); } + #[test] + fn mcp_debug_redacts_environment_and_argument_values() { + let env_secret = "debug-env-secret-must-not-render"; + let arg_secret = "debug-arg-secret-must-not-render"; + let rendered = format!( + "{:?}", + McpServer { + name: "analytics".into(), + command: "analytics-mcp".into(), + args: vec!["--token".into(), arg_secret.into()], + env: vec![EnvVar { + name: "ANALYTICS_TOKEN".into(), + value: env_secret.into(), + }], + } + ); + + assert!(rendered.contains("ANALYTICS_TOKEN")); + assert!(rendered.contains("arg_count: 2")); + assert!(rendered.contains(REDACTED_ENV_VALUE)); + assert!(!rendered.contains(env_secret)); + assert!(!rendered.contains(arg_secret)); + } + #[test] fn wire_redaction_covers_nested_env_values_without_changing_source() { let source = serde_json::json!({ @@ -2582,7 +2732,7 @@ mod tests { } }); - let redacted = redact_wire_value(&source); + let redacted = redact_wire_value(&source, None); assert_eq!( source["params"]["mcpServers"][0]["env"][0]["value"], "secret-one", @@ -2644,7 +2794,7 @@ mod tests { }); let original_source = source.clone(); - let redacted = redact_wire_value(&source); + let redacted = redact_wire_value(&source, None); for echo in redacted["echoes"].as_array().expect("redacted echoes") { assert_eq!(echo, REDACTED_ENV_VALUE); @@ -2658,6 +2808,41 @@ mod tests { ); } + #[test] + fn wire_redaction_suppresses_plain_sensitive_echoes_without_changing_source() { + let shorter = "token"; + let longer = "token-with-suffix"; + let source = serde_json::json!({ + "error": { + "message": format!("[adapter] rejected {longer}; retry with {shorter}") + }, + "ordinary": "safe diagnostic" + }); + let original_source = source.clone(); + let sensitive_values = vec![ + longer.to_string(), + shorter.to_string(), + "[".to_string(), + String::new(), + longer.to_string(), + ]; + + let matcher = sensitive_matcher(&sensitive_values); + let redacted = redact_wire_value(&source, matcher.as_ref()); + + let message = redacted["error"]["message"] + .as_str() + .expect("redacted error message"); + assert!(!message.contains(longer)); + assert!(!message.contains(shorter)); + assert_eq!(message, REDACTED_ENV_VALUE); + assert_eq!(redacted["ordinary"], source["ordinary"]); + assert_eq!( + source, original_source, + "the protocol value must stay unchanged" + ); + } + #[test] fn session_prompt_request_format() { let prompt_text = "[Buzz @mention]\nChannel: test\nFrom: npub1...\nMessage: hello"; @@ -3742,7 +3927,85 @@ mod tests { } #[tokio::test] - async fn semantic_traces_cannot_echo_serialized_mcp_config() { + async fn plain_mcp_secret_echo_is_redacted_and_old_session_values_are_retained() { + let first_secret = "first-session-adapter-secret"; + let second_secret = "second-session-adapter-secret"; + let error_response = serde_json::json!({ + "jsonrpc": "2.0", + "id": 2, + "error": { + "code": -32056, + "message": format!("adapter failed while using {first_secret}") + } + }) + .to_string(); + let script = format!( + "read -t 2 _init\n\ + printf '%s\\n' '{{\"jsonrpc\":\"2.0\",\"id\":0,\"result\":{{\"protocolVersion\":2,\"agentCapabilities\":{{}}}}}}'\n\ + read -t 2 _first_session\n\ + printf '%s\\n' '{{\"jsonrpc\":\"2.0\",\"id\":1,\"result\":{{\"sessionId\":\"ses_first\"}}}}'\n\ + read -t 2 _second_session\n\ + printf '%s\\n' '{error_response}'\n\ + sleep 1" + ); + let mut client = spawn_script(&script).await; + let observer = ObserverHandle::in_process(); + client.set_observer(Some(observer.clone()), 0); + client + .initialize() + .await + .expect("initialize should succeed"); + + let server_with_secret = |value: &str| McpServer { + name: "analytics".into(), + command: "analytics-mcp".into(), + args: vec![], + env: vec![EnvVar { + name: "ANALYTICS_TOKEN".into(), + value: value.into(), + }], + }; + + client + .session_new_full("/tmp", vec![server_with_secret(first_secret)], None, None) + .await + .expect("first session should succeed"); + let result = client + .session_new_full("/tmp", vec![server_with_secret(second_secret)], None, None) + .await; + + match result { + Err(AcpError::AgentError { code, message }) => { + assert_eq!(code, -32056); + assert_eq!(message, REDACTED_ENV_VALUE); + assert!(!message.contains(first_secret)); + } + Err(other) => panic!("expected redacted AgentError, got {other:?}"), + Ok(_) => panic!("expected second session/new to return an error"), + } + + assert_eq!(client.sensitive_mcp_env_values.len(), 2); + assert!( + client + .sensitive_mcp_env_values + .iter() + .any(|value| value == first_secret), + "credentials from earlier sessions must remain registered" + ); + assert!(client + .sensitive_mcp_env_values + .iter() + .any(|value| value == second_secret)); + let serialized_events = + serde_json::to_string(&observer.snapshot()).expect("serialize observer snapshot"); + assert!(!serialized_events.contains(first_secret)); + assert!(!serialized_events.contains(second_secret)); + assert!(serialized_events.contains(REDACTED_ENV_VALUE)); + client.shutdown().await; + } + + #[tokio::test] + async fn semantic_traces_redact_plain_mcp_secret_echoes() { use std::io::Write; use std::sync::{Arc, Mutex}; @@ -3774,26 +4037,29 @@ mod tests { } let secret = "semantic-trace-secret"; - let embedded_request = serde_json::json!({ - "method": "session/new", - "params": { - "mcpServers": [{ - "env": [{"name": "ANALYTICS_TOKEN", "value": secret}] - }] - } - }) - .to_string(); + let plain_echo = format!("adapter diagnostic included {secret}"); let update = serde_json::json!({ "jsonrpc": "2.0", "method": "session/update", "params": { "update": { "sessionUpdate": "agent_message_chunk", - "content": {"type": "text", "text": embedded_request} + "content": {"type": "text", "text": plain_echo} } } }); let mut client = spawn_inert_client().await; + client + .register_sensitive_mcp_env_values(&[McpServer { + name: "analytics".into(), + command: "analytics-mcp".into(), + args: vec![], + env: vec![EnvVar { + name: "ANALYTICS_TOKEN".into(), + value: secret.into(), + }], + }]) + .expect("credential redactor should initialize"); let trace = TraceCapture::default(); let subscriber = tracing_subscriber::fmt() .without_time() @@ -3806,7 +4072,7 @@ mod tests { let _ = client.handle_session_update(&update); client.handle_goose_usage_update(&serde_json::json!({ "params": { - "sessionId": embedded_request, + "sessionId": plain_echo, "update": { "sessionUpdate": "usage_update", "accumulatedInputTokens": 10, @@ -3821,7 +4087,7 @@ mod tests { "sessionId": "ordinary-session", "update": { "sessionUpdate": "usage_update", - "accumulatedInputTokens": embedded_request, + "accumulatedInputTokens": plain_echo, "accumulatedOutputTokens": 5, "accumulatedCachedInputTokens": null, "accumulatedCost": null @@ -4953,7 +5219,7 @@ mod tests { // Errors without a string `message` field (e.g. only a `data` field) must // not be silently truncated to "unknown error" — the full JSON is preserved. let error = serde_json::json!({"code": -32000, "data": "quota exceeded"}); - match super::agent_error_from_json(&error) { + match super::agent_error_from_json(&error, None) { AcpError::AgentError { code, message } => { assert_eq!(code, -32000); assert!( @@ -4978,7 +5244,7 @@ mod tests { } }); - let rendered = super::agent_error_from_json(&error).to_string(); + let rendered = super::agent_error_from_json(&error, None).to_string(); assert!(!rendered.contains(secret)); assert!(rendered.contains(REDACTED_ENV_VALUE)); @@ -4988,7 +5254,7 @@ mod tests { #[test] fn agent_error_from_json_uses_message_field_when_present() { let error = serde_json::json!({"code": -32001, "message": "auth denied"}); - match super::agent_error_from_json(&error) { + match super::agent_error_from_json(&error, None) { AcpError::AgentError { code, message } => { assert_eq!(code, -32001); assert_eq!(message, "auth denied"); diff --git a/crates/buzz-acp/src/config.rs b/crates/buzz-acp/src/config.rs index f268dc842e..491f2b04ec 100644 --- a/crates/buzz-acp/src/config.rs +++ b/crates/buzz-acp/src/config.rs @@ -68,7 +68,7 @@ const PROTECTED_MCP_ENV_NAMES: [&str; 6] = [ /// The transport tag is part of the version-1 document even though this PR /// implements only stdio. Additional transports can extend the same ordered /// server list without introducing a parallel configuration format. -#[derive(Debug, Clone, PartialEq, Eq, serde::Deserialize)] +#[derive(Clone, PartialEq, Eq, serde::Deserialize)] #[serde(tag = "transport", rename_all = "snake_case", deny_unknown_fields)] pub enum ConfiguredMcpServer { /// A local MCP child process connected over stdio. @@ -85,6 +85,37 @@ pub enum ConfiguredMcpServer { }, } +struct RedactedMcpEnv<'a>(&'a BTreeMap); + +impl std::fmt::Debug for RedactedMcpEnv<'_> { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let mut map = formatter.debug_map(); + for key in self.0.keys() { + map.entry(key, &"[REDACTED]"); + } + map.finish() + } +} + +impl std::fmt::Debug for ConfiguredMcpServer { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Stdio { + name, + command, + args, + env, + } => formatter + .debug_struct("Stdio") + .field("name", name) + .field("command", command) + .field("arg_count", &args.len()) + .field("env", &RedactedMcpEnv(env)) + .finish(), + } + } +} + #[derive(Debug, serde::Deserialize)] #[serde(deny_unknown_fields)] struct McpConfigDocument { @@ -3351,6 +3382,27 @@ channels = "ALL" assert!(!summary.contains("DOMAIN_TOKEN")); } + #[test] + fn debug_output_redacts_structured_mcp_environment_values() { + let secret_value = "structured-debug-secret"; + let argument_secret = "structured-argument-secret"; + let file = TempMcpConfig::write(&document_json(vec![serde_json::json!({ + "name": "safe", + "transport": "stdio", + "command": "/opt/mcp", + "args": ["--token", argument_secret], + "env": {"DOMAIN_TOKEN": secret_value} + })])); + let config = config_from_mcp_file(&file, None).expect("config should load"); + + let rendered = format!("{config:?}"); + + assert!(rendered.contains("DOMAIN_TOKEN")); + assert!(rendered.contains("[REDACTED]")); + assert!(!rendered.contains(secret_value)); + assert!(!rendered.contains(argument_secret)); + } + #[test] fn legacy_mcp_name_matches_existing_file_stem_behavior() { assert_eq!( From dbf1cc4923d33cc7098761f02467fd51547bbcce Mon Sep 17 00:00:00 2001 From: KC <79471844+wolfyy970@users.noreply.github.com> Date: Mon, 3 Aug 2026 11:34:39 -0400 Subject: [PATCH 6/6] docs(acp): describe the MCP transport field Signed-off-by: KC <79471844+wolfyy970@users.noreply.github.com> --- crates/buzz-acp/README.md | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/crates/buzz-acp/README.md b/crates/buzz-acp/README.md index d672d94b66..198645f67c 100644 --- a/crates/buzz-acp/README.md +++ b/crates/buzz-acp/README.md @@ -143,10 +143,10 @@ servers: ``` The JSON is strict. The only top-level fields are `version` and `servers`. -Each server has `name`, `command`, `args`, and `env`. Server names must be -unique, contain 1 to 128 ASCII bytes using only letters, digits, `_`, or `-`, -and cannot contain `__`. Names are checked across both structured entries and -the legacy server. +Each server has `name`, `transport`, `command`, `args`, and `env`. Version 1 +supports the `stdio` transport. Server names must be unique, contain 1 to 128 +ASCII bytes using only letters, digits, `_`, or `-`, and cannot contain `__`. +Names are checked across both structured entries and the legacy server. The config file is limited to 64 KiB. A harness can have at most 16 MCP servers in total, including the server from `BUZZ_ACP_MCP_COMMAND`. An