#![cfg(unix)]
use std::net::SocketAddr;
use std::process::Stdio;
use std::time::Duration;
use rmcp::service::RoleClient;
use rmcp::service::RunningService;
use rmcp::transport::streamable_http_client::StreamableHttpClientTransportConfig;
use rmcp::transport::StreamableHttpClientTransport;
use rmcp::ServiceExt as _;
use tokio::io::AsyncBufReadExt;
use tokio::io::AsyncWriteExt;
use tokio::io::BufReader;
use tokio::process::Child;
use tokio::process::ChildStdin;
use tokio::process::Command;
#[path = "common/deadline.rs"]
mod deadline;
#[path = "common/scrub_env.rs"]
mod scrub_env;
#[path = "common/startup_line.rs"]
mod startup_line;
use deadline::bounded;
use startup_line::HTTP_READY;
const EXIT_TIMEOUT: Duration = Duration::from_secs(5);
#[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,
);
}
const STDIO_READY: &str = "Starting Bugzilla MCP server on stdio";
const OVER_CAP_FRAME: usize = 5 * 1024 * 1024;
const OVER_CAP_LINE: &str =
"stdio frame exceeds the 4194304-byte request cap; closing the transport";
const RMCP_READ_ERROR_TARGET: &str = "rmcp::transport::async_rw";
const OVER_CAP_CLOSE: &str = "stdio transport closed: an inbound frame exceeded the request cap";
const OVER_CAP_EXIT_LINE: &str =
"Error: stdio transport closed: an inbound frame exceeded the request cap";
const SERVING_ERROR_PREFIX: &str = "serving error: ";
const EXIT_LINE_PREFIX: &str = "Error: stdio serving failed: ";
const HANDSHAKE_LINE_MAX: usize = 200;
const FILLER_CHARS: usize = 100_000;
const FILLER_RUN: usize = 64;
fn level_of(line: &str) -> Option<&str> {
line.split_whitespace().nth(1)
}
fn at_level(line: &str, level: &str) -> bool {
level_of(line) == Some(level)
}
fn from_bugwarden_at(line: &str, level: &str) -> bool {
const TARGETS: [&str; 4] = [
" bugwarden: ",
" bugwarden::",
" bugwarden_core: ",
" bugwarden_core::",
];
at_level(line, level) && TARGETS.iter().any(|target| line.contains(target))
}
struct Server {
child: Child,
stdin: Option<ChildStdin>,
stderr: startup_line::StderrLines,
log: String,
}
impl Server {
fn spawn(args: &[&str]) -> Self {
Self::spawn_with_stdout(args, Stdio::piped())
}
fn spawn_with_stdout(args: &[&str], stdout: Stdio) -> Self {
let mut cmd = Command::new(env!("CARGO_BIN_EXE_bugwarden"));
cmd.args(args)
.stdin(Stdio::piped())
.stdout(stdout)
.stderr(Stdio::piped());
for var in scrub_env::AMBIENT_VARS {
cmd.env_remove(var);
}
cmd.env("RUST_LOG", "info");
cmd.kill_on_drop(true);
let mut child = cmd.spawn().expect("the built binary must start");
let stdin = child.stdin.take();
let stderr = startup_line::stderr_lines(&mut child);
Self {
child,
stdin,
stderr,
log: String::new(),
}
}
async fn wait_for_stderr(&mut self, needle: &str) -> String {
startup_line::wait_for_line(&mut self.stderr, &mut self.log, needle, EXIT_TIMEOUT).await
}
async fn next_stderr_line(&mut self) -> Option<String> {
startup_line::next_logged_line(&mut self.stderr, &mut self.log).await
}
fn pid(&self) -> u32 {
self.child.id().expect("the child has a pid")
}
fn signal(&self, signal: &str) {
let pid = self.pid();
let status = std::process::Command::new("kill")
.args(["-s", signal, &pid.to_string()])
.status()
.expect("kill must be executable");
assert!(status.success(), "kill -s {signal} {pid} failed: {status}");
}
async fn assert_graceful_exit(mut self, what: &str) {
let status = tokio::time::timeout(EXIT_TIMEOUT, self.child.wait())
.await
.unwrap_or_else(|_| panic!("{what}: the process must exit within {EXIT_TIMEOUT:?}"))
.expect("the child must be waitable");
let _ = tokio::time::timeout(EXIT_TIMEOUT, async {
while self.next_stderr_line().await.is_some() {}
})
.await;
assert_eq!(
status.code(),
Some(0),
"{what}: SIGTERM/SIGINT must be a clean exit 0, not a signal-kill \
(code=None) or an error: status={status:?} stderr={log}",
log = self.log
);
assert!(
self.log.contains("received shutdown signal"),
"{what}: the shutdown path must log that it ran: {log}",
log = self.log
);
}
async fn assert_over_cap_exit(mut self, what: &str, pre_handshake: bool) {
let status = tokio::time::timeout(EXIT_TIMEOUT, self.child.wait())
.await
.unwrap_or_else(|_| {
panic!("{what}: an over-cap frame must end the process within {EXIT_TIMEOUT:?}")
})
.expect("the child must be waitable");
let _ = tokio::time::timeout(EXIT_TIMEOUT, async {
while self.next_stderr_line().await.is_some() {}
})
.await;
assert_eq!(
status.code(),
Some(1),
"{what}: a refused frame must be a failure exit, not the `0` rmcp's \
silent close would produce: status={status:?} stderr={log}",
log = self.log
);
let refusal = self
.log
.lines()
.find(|line| line.contains(OVER_CAP_LINE))
.unwrap_or_else(|| {
panic!(
"{what}: bugwarden must log the cap it enforced: {log}",
log = self.log
)
});
assert_eq!(
level_of(refusal),
Some("WARN"),
"{what}: the refusal is the cap doing its job and a client \
repeats it by reconnecting, so the reader's line is a WARN and \
not an ERROR: {refusal}"
);
let rmcp_echo = self
.log
.lines()
.find(|line| line.contains(RMCP_READ_ERROR_TARGET));
assert!(
rmcp_echo.is_some_and(|line| at_level(line, "ERROR")),
"{what}: rmcp still echoes the reader's error at ERROR — the \
line DESIGN.md documents, and the reason the directive that \
quiets it is documented with it: {log}",
log = self.log
);
let own_errors: Vec<&str> = self
.log
.lines()
.filter(|line| from_bugwarden_at(line, "ERROR"))
.collect();
if pre_handshake {
assert_eq!(
own_errors.len(),
1,
"{what}: one refusal earns bugwarden one ERROR line, the one \
that ends the process: {own_errors:?}"
);
assert!(
own_errors[0].contains(SERVING_ERROR_PREFIX),
"{what}: and that line is `serving error`, not the reader's \
refusal wearing its level: {line}",
line = own_errors[0]
);
} else {
assert!(
own_errors.is_empty(),
"{what}: after the handshake `main` bails with no line of its \
own, so bugwarden says nothing at ERROR here: {own_errors:?}"
);
}
assert!(
self.log.lines().any(|line| line == OVER_CAP_EXIT_LINE),
"{what}: the process must exit naming the cap, in the one wording \
both sides of the handshake share: {log}",
log = self.log
);
let serving_error = self
.log
.lines()
.find(|line| line.contains(SERVING_ERROR_PREFIX));
if pre_handshake {
let line = serving_error.unwrap_or_else(|| {
panic!(
"{what}: a refusal out of `serve` must log it: {log}",
log = self.log
)
});
assert!(
line.ends_with(OVER_CAP_CLOSE),
"{what}: the `serving error` line must name the cap, not a \
hangup the peer never chose: {line}"
);
assert!(
line.chars().count() <= HANDSHAKE_LINE_MAX,
"{what}: the `serving error` line must stay under \
{HANDSHAKE_LINE_MAX} chars: {len} chars",
len = line.chars().count()
);
} else {
assert!(
serving_error.is_none(),
"{what}: a refusal out of `waiting()` never entered `serve`'s \
error path: {log}",
log = self.log
);
}
}
async fn assert_bounded_handshake_failure(
mut self,
what: &str,
classification: &str,
forbidden: Option<&str>,
) {
let drained = tokio::time::timeout(EXIT_TIMEOUT, async {
while self.next_stderr_line().await.is_some() {}
})
.await;
assert!(
drained.is_ok(),
"{what}: a failed handshake must end the process within {EXIT_TIMEOUT:?}: {log}",
log = self.log
);
let status = tokio::time::timeout(EXIT_TIMEOUT, self.child.wait())
.await
.unwrap_or_else(|_| panic!("{what}: the process must be reaped after its stderr ends"))
.expect("the child must be waitable");
assert_eq!(
status.code(),
Some(1),
"{what}: a failed handshake must stay a failure exit: status={status:?} \
stderr={log}",
log = self.log
);
for prefix in [SERVING_ERROR_PREFIX, EXIT_LINE_PREFIX] {
let mut matched = self.log.lines().filter(|line| line.contains(prefix));
let line = matched.next().unwrap_or_else(|| {
panic!(
"{what}: stderr must carry a {prefix:?} line: {log}",
log = self.log
)
});
assert!(
matched.next().is_none(),
"{what}: {prefix:?} must be written once, not per attempt: {log}",
log = self.log
);
assert!(
line.ends_with(classification),
"{what}: {prefix:?} must be followed by the classification and \
nothing else: {line}"
);
if prefix == SERVING_ERROR_PREFIX {
assert_eq!(
level_of(line),
Some("ERROR"),
"{what}: a handshake failure ends the process, which is \
what this workspace keeps ERROR for: {line}"
);
}
assert!(
line.chars().count() <= HANDSHAKE_LINE_MAX,
"{what}: {prefix:?} must stay under {HANDSHAKE_LINE_MAX} chars, \
not grow with the frame: {len} chars",
len = line.chars().count()
);
}
if let Some(needle) = forbidden {
assert!(
!self.log.contains(needle),
"{what}: no {run}-char run of the client's own string may reach \
stderr: {log}",
run = needle.chars().count(),
log = self.log
);
}
}
async fn write_frame(&mut self, frame: &str) {
let stdin = self.stdin.as_mut().expect("stdin is piped");
let _ = tokio::time::timeout(EXIT_TIMEOUT, stdin.write_all(frame.as_bytes())).await;
let _ = tokio::time::timeout(EXIT_TIMEOUT, stdin.write_all(b"\n")).await;
}
fn close_stdin(&mut self) {
self.stdin = None;
}
}
fn stdout_with_no_reader() -> Stdio {
let (reader, mut writer) = std::io::pipe().expect("the test host must provide a pipe");
drop(reader);
for _ in 0..PROBE_ATTEMPTS {
match std::io::Write::write(&mut writer, b"\0") {
Err(e) if e.kind() == std::io::ErrorKind::BrokenPipe => return Stdio::from(writer),
_ => std::thread::sleep(PROBE_WAIT),
}
}
panic!("a pipe whose only reader was dropped never became BrokenPipe");
}
fn run_of(filler: char, run: usize) -> String {
std::iter::repeat_n(filler, run).collect()
}
const PROBE_ATTEMPTS: usize = 1000;
const PROBE_WAIT: Duration = Duration::from_millis(1);
fn spawn_stdio() -> Server {
Server::spawn(&STDIO_ARGS)
}
const STDIO_ARGS: [&str; 6] = [
"--transport",
"stdio",
"--bugzilla-server",
"https://bugzilla.example.invalid",
"--api-key",
"test-key",
];
async fn initialize_and_drain(server: &mut Server) -> tokio::task::JoinHandle<()> {
let stdin = server.stdin.as_mut().expect("stdin is piped");
let stdout = server.child.stdout.take().expect("stdout is piped");
let mut stdout = BufReader::new(stdout).lines();
stdin
.write_all(
br#"{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-11-25","capabilities":{},"clientInfo":{"name":"binary-shutdown-test","version":"0"}}}
"#,
)
.await
.expect("the child must accept the handshake");
let reply = tokio::time::timeout(EXIT_TIMEOUT, stdout.next_line())
.await
.expect("the handshake must not hang")
.expect("stdout must be readable")
.expect("the server must answer initialize");
assert!(
reply.contains("bugwarden"),
"initialize must complete before the test acts: {reply}"
);
tokio::spawn(async move { while stdout.next_line().await.ok().flatten().is_some() {} })
}
async fn probe_and_drain(server: &mut Server) -> tokio::task::JoinHandle<()> {
let stdin = server.stdin.as_mut().expect("stdin is piped");
let stdout = server.child.stdout.take().expect("stdout is piped");
let mut stdout = BufReader::new(stdout).lines();
stdin
.write_all(
br#"{"jsonrpc":"2.0","id":1,"method":"server/discover","params":{"_meta":{"io.modelcontextprotocol/protocolVersion":"2026-07-28","io.modelcontextprotocol/clientCapabilities":{}}}}
"#,
)
.await
.expect("the child must accept the probe");
let reply = tokio::time::timeout(EXIT_TIMEOUT, stdout.next_line())
.await
.expect("the probe must not hang")
.expect("stdout must be readable")
.expect("the server must answer the probe");
assert!(
reply.contains("supportedVersions"),
"the probe must be answered before the test acts: {reply}"
);
tokio::spawn(async move { while stdout.next_line().await.ok().flatten().is_some() {} })
}
async fn write_an_over_cap_frame(server: &mut Server) {
let stdin = server.stdin.as_mut().expect("stdin is piped");
let frame = vec![b'a'; OVER_CAP_FRAME];
let _ = tokio::time::timeout(EXIT_TIMEOUT, stdin.write_all(&frame)).await;
}
async fn spawn_http() -> (Server, SocketAddr) {
let mut server = Server::spawn(&[
"--bugzilla-server",
"https://bugzilla.example.invalid",
"--host",
"127.0.0.1",
"--port",
"0",
"--insecure-no-auth",
]);
let line = server.wait_for_stderr(HTTP_READY).await;
let addr = startup_line::parse_bound_addr(&line);
wait_for_tcp(addr).await;
(server, addr)
}
async fn wait_for_tcp(addr: SocketAddr) {
let ready = tokio::time::timeout(EXIT_TIMEOUT, async {
loop {
if tokio::net::TcpStream::connect(addr).await.is_ok() {
return;
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
})
.await;
assert!(ready.is_ok(), "the binary must start serving on {addr}");
}
async fn connect_insecure(addr: SocketAddr) -> RunningService<RoleClient, ()> {
let transport = StreamableHttpClientTransport::with_client(
reqwest::Client::new(),
StreamableHttpClientTransportConfig::with_uri(format!("http://{addr}/mcp")),
);
bounded("the MCP handshake", ().serve(transport))
.await
.expect("MCP handshake must succeed under --insecure-no-auth")
}
#[tokio::test]
async fn http_sigterm_exits_zero_on_an_idle_listener() {
let (server, _addr) = spawn_http().await;
server.signal("TERM");
server.assert_graceful_exit("http idle SIGTERM").await;
}
#[tokio::test]
async fn http_sigint_still_exits_zero() {
let (server, _addr) = spawn_http().await;
server.signal("INT");
server.assert_graceful_exit("http idle SIGINT").await;
}
#[tokio::test]
async fn http_sigterm_cancels_a_live_session() {
let (server, addr) = spawn_http().await;
let _session = connect_insecure(addr).await;
server.signal("TERM");
server
.assert_graceful_exit("http live-session SIGTERM")
.await;
}
#[tokio::test]
async fn stdio_sigterm_during_the_handshake_wait_exits_zero() {
let mut server = spawn_stdio();
server.wait_for_stderr(STDIO_READY).await;
server.signal("TERM");
server
.assert_graceful_exit("stdio pre-handshake SIGTERM")
.await;
}
#[tokio::test]
async fn stdio_an_over_cap_frame_before_the_handshake_exits_one() {
let mut server = spawn_stdio();
server.wait_for_stderr(STDIO_READY).await;
write_an_over_cap_frame(&mut server).await;
server
.assert_over_cap_exit("stdio pre-handshake over-cap", true)
.await;
}
#[tokio::test]
async fn stdio_a_probe_then_an_over_cap_frame_still_exits_one() {
let mut server = spawn_stdio();
server.wait_for_stderr(STDIO_READY).await;
let _drain = probe_and_drain(&mut server).await;
write_an_over_cap_frame(&mut server).await;
server
.assert_over_cap_exit("stdio probed pre-handshake over-cap", true)
.await;
}
#[tokio::test]
async fn stdio_sigterm_after_initialize_exits_zero() {
let mut server = spawn_stdio();
server.wait_for_stderr(STDIO_READY).await;
let _drain = initialize_and_drain(&mut server).await;
server.signal("TERM");
server
.assert_graceful_exit("stdio post-handshake SIGTERM")
.await;
}
#[tokio::test]
async fn stdio_an_over_cap_frame_after_initialize_exits_one() {
let mut server = spawn_stdio();
server.wait_for_stderr(STDIO_READY).await;
let _drain = initialize_and_drain(&mut server).await;
write_an_over_cap_frame(&mut server).await;
server
.assert_over_cap_exit("stdio post-handshake over-cap", false)
.await;
}
#[tokio::test]
async fn stdio_a_pre_initialize_tool_call_never_logs_the_frame() {
let mut server = spawn_stdio();
server.wait_for_stderr(STDIO_READY).await;
let query = run_of('a', FILLER_CHARS);
server
.write_frame(&format!(
r#"{{"jsonrpc":"2.0","id":1,"method":"tools/call","params":{{"name":"bug_info","arguments":{{"query":"{query}"}}}}}}"#
))
.await;
server
.assert_bounded_handshake_failure(
"stdio pre-handshake tools/call",
"the first frame was a request other than initialize",
Some(&run_of('a', FILLER_RUN)),
)
.await;
}
#[tokio::test]
async fn stdio_a_pre_initialize_notification_is_named_by_kind() {
let mut server = spawn_stdio();
server.wait_for_stderr(STDIO_READY).await;
server
.write_frame(r#"{"jsonrpc":"2.0","method":"notifications/initialized"}"#)
.await;
server
.assert_bounded_handshake_failure(
"stdio pre-handshake notification",
"the first frame was a notification, not initialize",
None,
)
.await;
}
#[tokio::test]
async fn stdio_a_pre_initialize_response_never_logs_its_id() {
let mut server = spawn_stdio();
server.wait_for_stderr(STDIO_READY).await;
let id = run_of('x', 4096);
server
.write_frame(&format!(r#"{{"jsonrpc":"2.0","id":"{id}","result":{{}}}}"#))
.await;
server
.assert_bounded_handshake_failure(
"stdio pre-handshake response",
"the first frame was a response, not initialize",
Some(&run_of('x', FILLER_RUN)),
)
.await;
}
#[tokio::test]
async fn stdio_a_pre_initialize_error_reply_never_logs_its_message() {
let mut server = spawn_stdio();
server.wait_for_stderr(STDIO_READY).await;
let message = run_of('e', FILLER_CHARS);
server
.write_frame(&format!(
r#"{{"jsonrpc":"2.0","id":1,"error":{{"code":-32000,"message":"{message}"}}}}"#
))
.await;
server
.assert_bounded_handshake_failure(
"stdio pre-handshake error reply",
"the first frame was an error reply, not initialize",
Some(&run_of('e', FILLER_RUN)),
)
.await;
}
#[tokio::test]
async fn stdio_a_hangup_before_initialize_is_classified_not_echoed() {
let mut server = spawn_stdio();
server.wait_for_stderr(STDIO_READY).await;
server.close_stdin();
server
.assert_bounded_handshake_failure(
"stdio pre-handshake hangup",
"the stream ended before initialize",
None,
)
.await;
}
#[tokio::test]
async fn stdio_a_broken_stdout_during_the_handshake_is_classified() {
let mut server = Server::spawn_with_stdout(&STDIO_ARGS, stdout_with_no_reader());
server.wait_for_stderr(STDIO_READY).await;
server
.write_frame(r#"{"jsonrpc":"2.0","id":1,"method":"tools/list"}"#)
.await;
server
.assert_bounded_handshake_failure(
"stdio pre-handshake broken stdout",
"the stdio transport failed during the handshake",
None,
)
.await;
}
#[cfg(target_os = "linux")]
#[tokio::test]
async fn stdio_a_refused_initialize_is_classified_not_echoed() {
let dir = tempfile::tempdir().expect("a temp dir");
let config = dir.path().join("audit.toml");
std::fs::write(
&config,
"path = \"/dev/full\"\nfail_mode = \"closed_all\"\n",
)
.expect("the audit config must be writable");
let mut server = Server::spawn(&[
"--transport",
"stdio",
"--bugzilla-server",
"https://bugzilla.example.invalid",
"--api-key",
"test-key",
"--audit-config",
config.to_str().expect("a utf-8 temp path"),
]);
server.wait_for_stderr(STDIO_READY).await;
server
.write_frame(
r#"{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-11-25","capabilities":{},"clientInfo":{"name":"binary-shutdown-test","version":"0"}}}"#,
)
.await;
server
.assert_bounded_handshake_failure(
"stdio initialize refused by the audit sink",
"this server refused initialize",
None,
)
.await;
}