use std::collections::HashSet;
use std::process::{ExitStatus, Stdio};
use std::time::Duration;
use serde_json::{json, Value};
use tokio::io::{AsyncBufReadExt as _, AsyncWriteExt as _, BufReader};
use tokio::process::{Child, ChildStdin};
#[path = "common/scrub_env.rs"]
mod scrub_env;
const REPLY_BUDGET: Duration = Duration::from_secs(20);
const UNSERVED: &str = "2027-01-01";
#[test]
fn the_scrub_list_covers_every_environment_fallback() {
scrub_env::assert_the_scrub_list_covers_every_environment_fallback(
scrub_env::AMBIENT_VARS,
scrub_env::HTTP_TOKEN_VARS,
);
}
struct Server {
child: Child,
stdin: ChildStdin,
stdout: tokio::io::Lines<BufReader<tokio::process::ChildStdout>>,
}
impl Server {
fn spawn() -> Self {
Self::spawn_with(Stdio::null())
}
fn spawn_logged() -> Self {
Self::spawn_with(Stdio::piped())
}
fn spawn_with(stderr: Stdio) -> Self {
let mut cmd = tokio::process::Command::new(env!("CARGO_BIN_EXE_bugwarden"));
cmd.args(["--transport", "stdio"])
.args(["--bugzilla-server", "https://bugzilla.example.invalid"])
.args(["--api-key", "test-key"])
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(stderr)
.kill_on_drop(true);
for var in scrub_env::AMBIENT_VARS {
cmd.env_remove(var);
}
let mut child = cmd.spawn().expect("the built binary must start");
let stdin = child.stdin.take().expect("stdin is piped");
let stdout = BufReader::new(child.stdout.take().expect("stdout is piped")).lines();
Server {
child,
stdin,
stdout,
}
}
async fn send(&mut self, message: Value) {
self.stdin
.write_all(format!("{message}\n").as_bytes())
.await
.expect("the child must accept input");
}
async fn recv(&mut self) -> Value {
let line = tokio::time::timeout(REPLY_BUDGET, self.stdout.next_line())
.await
.expect("the child must answer within the budget")
.expect("the child's stdout must be readable")
.expect("the child must not close stdout before answering");
serde_json::from_str(&line).unwrap_or_else(|e| panic!("{line:?}: {e}"))
}
async fn call(&mut self, message: Value) -> Value {
self.send(message).await;
self.recv().await
}
async fn burst(&mut self, frames: &[Value]) {
let mut buffer = String::new();
for frame in frames {
buffer.push_str(&format!("{frame}\n"));
}
self.stdin
.write_all(buffer.as_bytes())
.await
.expect("the child must accept input");
}
async fn finish(self) -> (ExitStatus, String) {
drop(self.stdin);
let output = tokio::time::timeout(REPLY_BUDGET, self.child.wait_with_output())
.await
.expect("the child must exit when stdin closes")
.expect("the child must be waitable");
(
output.status,
String::from_utf8_lossy(&output.stderr).into_owned(),
)
}
async fn close(self) {
let (status, _) = self.finish().await;
assert!(
status.success(),
"a probe must refuse the request, never the process: {status}"
);
}
}
fn go_sdk_discover(id: u32, version: &str) -> Value {
json!({
"jsonrpc": "2.0", "id": id, "method": "server/discover",
"params": { "_meta": {
"io.modelcontextprotocol/protocolVersion": version,
"io.modelcontextprotocol/clientCapabilities": {},
"io.modelcontextprotocol/clientInfo": { "name": "agy", "version": "1.7.0" },
}}
})
}
fn legacy_initialize(id: u32) -> Value {
json!({
"jsonrpc": "2.0", "id": id, "method": "initialize",
"params": {
"protocolVersion": "2025-11-25",
"capabilities": {},
"clientInfo": { "name": "agy", "version": "1.7.0" },
}
})
}
fn bare_list(id: u32) -> Value {
json!({ "jsonrpc": "2.0", "id": id, "method": "tools/list", "params": {} })
}
#[tokio::test]
async fn a_refused_probe_leaves_the_legacy_handshake_intact() {
let mut server = Server::spawn();
let probe = server.call(go_sdk_discover(1, UNSERVED)).await;
assert_eq!(probe["error"]["code"], json!(-32022), "{probe}");
assert_eq!(probe["error"]["data"]["requested"], UNSERVED, "{probe}");
assert!(
probe["error"]["data"]["supported"]
.as_array()
.is_some_and(|served| served.contains(&json!("2026-07-28"))),
"the refusal must hand back the served list the client retries from: {probe}"
);
let init = server.call(legacy_initialize(2)).await;
assert_eq!(init["result"]["protocolVersion"], "2025-11-25", "{init}");
server
.send(json!({ "jsonrpc": "2.0", "method": "notifications/initialized" }))
.await;
let listed = server.call(bare_list(3)).await;
assert!(
listed["result"]["tools"]
.as_array()
.is_some_and(|tools| !tools.is_empty()),
"a probed-then-handshook session must be served its tools: {listed}"
);
server.close().await;
}
#[tokio::test]
async fn a_served_probe_leaves_the_legacy_handshake_intact() {
let mut server = Server::spawn();
let probe = server.call(go_sdk_discover(1, "2026-07-28")).await;
assert_eq!(
probe["result"]["_meta"]["io.modelcontextprotocol/serverInfo"],
json!({ "name": "bugwarden", "version": env!("CARGO_PKG_VERSION") }),
"the probe must name this build: {probe}"
);
server.call(legacy_initialize(2)).await;
server
.send(json!({ "jsonrpc": "2.0", "method": "notifications/initialized" }))
.await;
let listed = server.call(bare_list(3)).await;
assert!(
listed["result"]["tools"]
.as_array()
.is_some_and(|tools| !tools.is_empty()),
"a served probe must not commit the session either: {listed}"
);
server.close().await;
}
#[tokio::test]
async fn a_probe_with_no_meta_refuses_the_request_not_the_process() {
let mut server = Server::spawn();
let probe = server
.call(
json!({ "jsonrpc": "2.0", "id": 1, "method": "server/discover",
"params": {} }),
)
.await;
assert_eq!(probe["error"]["code"], json!(-32602), "{probe}");
assert_eq!(
probe["error"]["message"],
json!(
"request _meta is missing or has malformed required fields: \
io.modelcontextprotocol/protocolVersion, \
io.modelcontextprotocol/clientCapabilities"
),
"{probe}"
);
let init = server.call(legacy_initialize(2)).await;
assert_eq!(init["result"]["protocolVersion"], "2025-11-25", "{init}");
server
.send(json!({ "jsonrpc": "2.0", "method": "notifications/initialized" }))
.await;
let listed = server.call(bare_list(3)).await;
assert!(
listed["result"]["tools"]
.as_array()
.is_some_and(|tools| !tools.is_empty()),
"a malformed probe must not end the session: {listed}"
);
server.close().await;
}
#[tokio::test]
async fn a_probe_still_commits_nothing_that_a_declaring_request_does_not() {
let mut server = Server::spawn();
server.call(go_sdk_discover(1, "2026-07-28")).await;
let declared = server
.call(json!({
"jsonrpc": "2.0", "id": 2, "method": "tools/list",
"params": { "_meta": {
"io.modelcontextprotocol/protocolVersion": "2026-07-28",
"io.modelcontextprotocol/clientCapabilities": {},
}}
}))
.await;
assert!(
declared["result"]["tools"]
.as_array()
.is_some_and(|tools| !tools.is_empty()),
"{declared}"
);
let bare = server.call(bare_list(3)).await;
assert_eq!(
bare["error"]["code"],
json!(-32602),
"the handshake-free lifecycle must still demand per-request _meta: {bare}"
);
server.close().await;
}
#[tokio::test]
async fn a_probe_only_client_hangs_up_cleanly() {
for probe in [
go_sdk_discover(1, "2026-07-28"),
go_sdk_discover(1, UNSERVED),
json!({ "jsonrpc": "2.0", "id": 1, "method": "server/discover", "params": {} }),
] {
let mut server = Server::spawn_logged();
server.call(probe.clone()).await;
let (status, log) = server.finish().await;
assert!(
status.success(),
"a probe-only client hung up: {probe}: {status}"
);
assert!(
log.contains("peer hung up after a server/discover probe"),
"the hangup must be logged as one: {probe}: {log}"
);
assert!(
!log.contains("serving error"),
"a hangup is not a serving failure: {probe}: {log}"
);
}
}
#[tokio::test]
async fn a_pipelined_probe_is_answered_over_real_pipes() {
const FRAMES: u32 = 40;
let mut server = Server::spawn();
server.call(legacy_initialize(1)).await;
server
.send(json!({ "jsonrpc": "2.0", "method": "notifications/initialized" }))
.await;
let frames: Vec<Value> = (2..=FRAMES + 1)
.map(|id| {
if id % 2 == 0 {
go_sdk_discover(id, "2026-07-28")
} else {
bare_list(id)
}
})
.collect();
server.burst(&frames).await;
let mut answered = HashSet::new();
for _ in 0..FRAMES {
let reply = server.recv().await;
answered.insert(
reply["id"]
.as_u64()
.unwrap_or_else(|| panic!("every reply names its request: {reply}")),
);
}
let unanswered: Vec<u64> = (2..=u64::from(FRAMES) + 1)
.filter(|id| !answered.contains(id))
.collect();
assert!(unanswered.is_empty(), "unanswered ids: {unanswered:?}");
server.close().await;
}