From 8fdcf141c52df8825c893be3f922f8e51f6ec06d Mon Sep 17 00:00:00 2001 From: Nicolas Dreno Date: Fri, 18 Sep 2026 17:50:12 +0200 Subject: [PATCH 1/4] fix: stop dropping tokio runtimes inside the async context Three sites, all in the released binary, all the same panic: Cannot drop a runtime in a context where blocking is not allowed. **Shutdown.** `Gateway::load` builds a NATS publisher, a Kafka publisher and an LDAP client whatever the artifact contains, each carrying its own runtime for the synchronous calls WASM host functions make. Letting the gateway fall out of scope at the end of `run_serve` dropped all three inside the runtime. Every clean shutdown exited 101, which is how a supervisor decides a process crashed, so restart backoff, crash-loop counters and alerts fired on an ordinary stop. The drop moves to a plain thread, which has no async context. **Startup failure.** The same drop on every early return after the artifact loads: an invalid address, a taken port, an unreadable TLS config, a bad admin bind. The real cause was printed and then buried under a panic, and a configuration error and a crash became indistinguishable by exit code. **Compiling a URL-sourced plugin.** `reqwest::blocking` refuses to run inside a runtime on construction and on every request, so a manifest with a remote plugin could not be compiled at all from a dev-profile build. The download runs on a scoped thread, which can still borrow the URL. The lifecycle test spawns the binary, signals it and reads the exit code, which nothing did before: every other test drives the gateway in-process, which is exactly why this shipped. It fails against the old code with the panic in stderr. Closes #208. --- CHANGELOG.md | 5 + Cargo.lock | 1 + crates/barbacane-compiler/src/download.rs | 62 +++++- crates/barbacane-test/Cargo.toml | 2 + .../barbacane-test/tests/process_lifecycle.rs | 186 ++++++++++++++++++ crates/barbacane/src/main.rs | 29 +++ 6 files changed, 276 insertions(+), 9 deletions(-) create mode 100644 crates/barbacane-test/tests/process_lifecycle.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index 6540de5f..99067012 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed + +- **data plane**: `barbacane serve` exits 0 on `SIGTERM` instead of panicking and exiting 101. The plugin host builds a NATS publisher, a Kafka publisher and an LDAP client whatever the artifact contains, each carrying its own tokio runtime, and dropping a runtime inside an async context is a panic. Every clean shutdown therefore looked like a crash to a supervisor: restart backoff, crash-loop counters and alerts fired on an ordinary stop. The same panic fired when startup failed after the artifact loaded, for instance on a taken port, burying the real cause and making a configuration error indistinguishable from a crash by exit code. +- **compiler**: `compile` resolves a URL-sourced plugin instead of panicking. `reqwest::blocking` refuses to run inside a tokio runtime, on construction and on every request, and `compile` runs inside one, so a manifest with a remote plugin could not be compiled at all from a dev-profile build. The download now runs on a thread of its own. + ## [0.11.0] - 2026-09-17 Headline: a request now carries only the headers its operation admits, so the document describes what an upstream receives and not only what a client may send. Specs that could not be compiled at all now can, including any with a recursive schema. diff --git a/Cargo.lock b/Cargo.lock index cdeb1ba4..19a1d8b4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -527,6 +527,7 @@ dependencies = [ "base64", "flate2", "futures-util", + "libc", "p256", "predicates", "rcgen", diff --git a/crates/barbacane-compiler/src/download.rs b/crates/barbacane-compiler/src/download.rs index 5bd81b5a..903a78a3 100644 --- a/crates/barbacane-compiler/src/download.rs +++ b/crates/barbacane-compiler/src/download.rs @@ -23,14 +23,44 @@ pub struct DownloadResult { pub plugin_toml: Option, } -/// Build a blocking HTTP client with appropriate defaults. -fn build_client() -> Result { - reqwest::blocking::Client::builder() - .connect_timeout(Duration::from_secs(30)) - .timeout(Duration::from_secs(120)) - .user_agent(format!("barbacane-compiler/{}", env!("CARGO_PKG_VERSION"))) - .build() - .map_err(|e| CompileError::PluginResolution(format!("failed to build HTTP client: {e}"))) +/// The blocking HTTP client used for every plugin download. +/// +/// Built once and never dropped. A blocking client owns a tokio runtime, and +/// `compile` runs inside one, where dropping a runtime is a panic in tokio. +/// Keeping it in a `OnceLock` means the process never drops it, and downloads +/// reuse one connection pool rather than building a client each time. +static CLIENT: std::sync::OnceLock> = + std::sync::OnceLock::new(); + +/// The shared blocking HTTP client. +/// +/// Built on a plain thread, once, and never dropped. `compile` runs inside a +/// tokio runtime, and a blocking client both creates and drops a temporary +/// runtime while being constructed, which tokio refuses to do inside an async +/// context. A thread of our own has no such context. Storing the result means +/// the cost is paid once and downloads share one connection pool. +fn client() -> Result<&'static reqwest::blocking::Client, CompileError> { + CLIENT + .get_or_init(|| { + std::thread::Builder::new() + .name("barbacane-http-init".into()) + .spawn(|| { + reqwest::blocking::Client::builder() + .connect_timeout(Duration::from_secs(30)) + .timeout(Duration::from_secs(120)) + .user_agent(format!("barbacane-compiler/{}", env!("CARGO_PKG_VERSION"))) + .build() + .map_err(|e| format!("failed to build HTTP client: {e}")) + }) + .map_err(|e| format!("failed to start the HTTP client thread: {e}")) + .and_then(|h| { + h.join() + .map_err(|_| "the HTTP client thread panicked".to_string()) + }) + .and_then(|r| r) + }) + .as_ref() + .map_err(|e| CompileError::PluginResolution(e.clone())) } /// Derive candidate plugin.toml URLs from a .wasm URL. @@ -56,13 +86,27 @@ fn derive_plugin_toml_urls(wasm_url: &str) -> Vec { /// Fetches the .wasm binary and attempts to fetch `plugin.toml` from /// the same directory (best-effort — 404 is fine). pub fn download_plugin(url: &str) -> Result { + // `reqwest::blocking` refuses to run inside a tokio runtime, on construction + // and on every request, and `compile` runs inside one. A scoped thread has + // no async context and can still borrow `url`. + std::thread::scope( + |scope| match scope.spawn(|| download_plugin_blocking(url)).join() { + Ok(result) => result, + Err(_) => Err(CompileError::PluginResolution(format!( + "the download thread panicked fetching {url}" + ))), + }, + ) +} + +fn download_plugin_blocking(url: &str) -> Result { if !url.starts_with("https://") { return Err(CompileError::PluginResolution(format!( "plugin URL must use HTTPS: {url}" ))); } - let client = build_client()?; + let client = client()?; tracing::info!(url, "downloading remote plugin"); diff --git a/crates/barbacane-test/Cargo.toml b/crates/barbacane-test/Cargo.toml index 30e1fbe4..daa19721 100644 --- a/crates/barbacane-test/Cargo.toml +++ b/crates/barbacane-test/Cargo.toml @@ -8,6 +8,8 @@ license.workspace = true publish = false [dependencies] +# Sending a real SIGTERM in the process-lifecycle test. +libc = "0.2" barbacane-compiler = { workspace = true } tokio = { workspace = true } reqwest = { workspace = true } diff --git a/crates/barbacane-test/tests/process_lifecycle.rs b/crates/barbacane-test/tests/process_lifecycle.rs new file mode 100644 index 00000000..34b3ece7 --- /dev/null +++ b/crates/barbacane-test/tests/process_lifecycle.rs @@ -0,0 +1,186 @@ +//! The gateway as a process: does it start, serve, and stop cleanly? +//! +//! Every other test drives the gateway in-process, so nothing started the real +//! binary, signalled it and looked at its exit code. That is how a panic on +//! every graceful shutdown reached a release: the drain was correct, only the +//! exit code lied about it, and a supervisor reads a non-zero exit as a crash. +//! +//! The cause was a runtime dropped inside the async context. The plugin host +//! owns broker and directory clients that each carry one, so the panic fired +//! whatever the artifact contained. + +use std::process::{Command, Stdio}; +use std::time::{Duration, Instant}; + +/// The gateway binary built alongside these tests. +fn gateway_binary() -> std::path::PathBuf { + let mut dir = std::env::current_exe().expect("test binary path"); + dir.pop(); // deps/ + dir.pop(); // debug/ or release/ + dir.join("barbacane") +} + +/// The smallest artifact the repository can build: one mock route. +fn build_artifact(dir: &std::path::Path) -> Option { + let repo = std::path::Path::new(env!("CARGO_MANIFEST_DIR")) + .parent()? + .parent()? + .to_path_buf(); + let mock = repo.join("plugins/mock/mock.wasm"); + if !mock.exists() { + return None; + } + + let manifest = dir.join("barbacane.yaml"); + std::fs::write( + &manifest, + format!("plugins:\n mock:\n path: {}\n", mock.display()), + ) + .ok()?; + + let spec = dir.join("api.yaml"); + std::fs::write( + &spec, + r#"openapi: "3.0.3" +info: { title: lifecycle, version: "1.0.0" } +paths: + /ping: + get: + operationId: ping + x-barbacane-dispatch: { name: mock, config: { status: 200, body: "pong" } } + responses: { "200": { description: ok } } +"#, + ) + .ok()?; + + let out = dir.join("api.bca"); + let status = Command::new(gateway_binary()) + .args(["compile", "-s"]) + .arg(&spec) + .arg("-m") + .arg(&manifest) + .arg("-o") + .arg(&out) + .output() + .ok()?; + if !status.status.success() { + eprintln!( + "skipping: could not compile the fixture: {}", + String::from_utf8_lossy(&status.stderr) + ); + return None; + } + Some(out) +} + +fn wait_for_port(port: u16, limit: Duration) -> bool { + let deadline = Instant::now() + limit; + while Instant::now() < deadline { + if std::net::TcpStream::connect(("127.0.0.1", port)).is_ok() { + return true; + } + std::thread::sleep(Duration::from_millis(100)); + } + false +} + +/// SIGTERM on a healthy gateway must drain and exit 0. A non-zero exit is how +/// a supervisor decides the process crashed, so a clean stop that exits 101 +/// produces restart backoff, crash-loop counters and alerts on every deploy. +#[test] +fn sigterm_on_a_serving_gateway_exits_zero() { + let binary = gateway_binary(); + if !binary.exists() { + eprintln!("skipping: {} is not built", binary.display()); + return; + } + let dir = tempfile::tempdir().expect("temp dir"); + let Some(artifact) = build_artifact(dir.path()) else { + return; + }; + + let mut child = Command::new(&binary) + .arg("serve") + .arg("--artifact") + .arg(&artifact) + .args(["--listen", "127.0.0.1:34201"]) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn() + .expect("spawn the gateway"); + + assert!( + wait_for_port(34201, Duration::from_secs(30)), + "the gateway never accepted a connection" + ); + + unsafe { + libc::kill(child.id() as i32, libc::SIGTERM); + } + + let deadline = Instant::now() + Duration::from_secs(30); + let status = loop { + match child.try_wait().expect("wait") { + Some(status) => break status, + None if Instant::now() > deadline => { + let _ = child.kill(); + panic!("the gateway did not exit within 30s of SIGTERM"); + } + None => std::thread::sleep(Duration::from_millis(100)), + } + }; + + let mut stderr = String::new(); + if let Some(mut e) = child.stderr.take() { + use std::io::Read; + let _ = e.read_to_string(&mut stderr); + } + + assert!( + !stderr.contains("panicked"), + "the gateway panicked while shutting down:\n{stderr}" + ); + assert_eq!( + status.code(), + Some(0), + "SIGTERM must exit 0, got {status:?}\nstderr:\n{stderr}" + ); +} + +/// A taken port is a configuration error. It must report the cause and exit +/// non-zero without panicking, so a misconfiguration stays distinguishable +/// from a crash. +#[test] +fn a_taken_port_reports_the_cause_without_panicking() { + let binary = gateway_binary(); + if !binary.exists() { + return; + } + let dir = tempfile::tempdir().expect("temp dir"); + let Some(artifact) = build_artifact(dir.path()) else { + return; + }; + + let held = std::net::TcpListener::bind("127.0.0.1:34202").expect("hold the port"); + + let output = Command::new(&binary) + .arg("serve") + .arg("--artifact") + .arg(&artifact) + .args(["--listen", "127.0.0.1:34202"]) + .output() + .expect("run the gateway"); + + drop(held); + + let stderr = String::from_utf8_lossy(&output.stderr); + assert!( + stderr.contains("failed to bind"), + "the real cause must be reported:\n{stderr}" + ); + assert!( + !stderr.contains("panicked"), + "a taken port is a configuration error, not a panic:\n{stderr}" + ); + assert_eq!(output.status.code(), Some(1), "stderr:\n{stderr}"); +} diff --git a/crates/barbacane/src/main.rs b/crates/barbacane/src/main.rs index 37e1a692..00c3a43c 100644 --- a/crates/barbacane/src/main.rs +++ b/crates/barbacane/src/main.rs @@ -4779,6 +4779,7 @@ async fn run_serve( Ok(a) => a, Err(_) => { eprintln!("error: invalid listen address: {}", listen); + drop_off_runtime(gateway); return ExitCode::from(1); } }; @@ -4788,6 +4789,7 @@ async fn run_serve( Ok(l) => l, Err(e) => { eprintln!("error: failed to bind to {}: {}", addr, e); + drop_off_runtime(gateway); return ExitCode::from(1); } }; @@ -4799,6 +4801,7 @@ async fn run_serve( Ok(c) => c, Err(e) => { eprintln!("error: {}", e); + drop_off_runtime(gateway); return ExitCode::from(1); } }; @@ -4871,6 +4874,7 @@ async fn run_serve( "error: invalid --admin-bind address '{}': {}", admin_bind, e ); + drop_off_runtime(gateway); return ExitCode::from(1); } }; @@ -5173,9 +5177,34 @@ async fn run_serve( tokio::time::sleep(Duration::from_millis(500)).await; } + // The gateway holds runtimes; letting it fall out of scope here panics. + drop_off_runtime(gateway); + ExitCode::SUCCESS } +/// Drop a value outside the async runtime. +/// +/// The plugin host owns broker and directory clients, each carrying its own +/// tokio runtime for the synchronous host calls WASM makes. Dropping a runtime +/// inside an async context is a panic in tokio, so the last reference must not +/// fall out of scope on a runtime thread. A plain thread has no async context. +/// +/// Silent when the value is not the last reference: the drop is then cheap and +/// touches no runtime. +fn drop_off_runtime(value: T) { + if let Err(e) = std::thread::Builder::new() + .name("barbacane-drop".into()) + .spawn(move || drop(value)) + .and_then(|h| { + h.join() + .map_err(|_| std::io::Error::other("drop thread panicked")) + }) + { + tracing::debug!(error = %e, "could not drop off the runtime; leaking instead of panicking"); + } +} + /// Wait for shutdown signal (SIGTERM or SIGINT). async fn wait_for_shutdown_signal() { #[cfg(unix)] From 946c7597f2bb9292cc336cc3549c6821030ecd76 Mon Sep 17 00:00:00 2001 From: Nicolas Dreno Date: Mon, 21 Sep 2026 10:22:10 +0200 Subject: [PATCH 2/4] fix(wasm): shut the broker and directory runtimes down in the background Moving the gateway's drop off the runtime fixed the shutdown path only because of ordering. `SharedGateway` is an `Arc` that the MCP eviction and hot-reload tasks clone, and neither is joined, so either could hold the last reference and drop it on a worker thread after the local one is gone. The panic would have come back, rarely, which is worse than reliably. Fixed where the runtimes live instead of where they happen to be dropped. `NatsPublisher`, `KafkaPublisher` and `LdapClient` now hand theirs to `shutdown_background` rather than waiting for workers to stop, so the last reference is safe to fall anywhere. The tests drop one inside `block_on` and inside a spawned task, and both fail with the original panic if the plain drop is restored. The download path gains the regression test it lacked: the lifecycle test uses a local plugin path, so nothing reached `download_plugin`. Calling it from inside a runtime and from a spawned task must return an error rather than take the process down. Both fail without the scoped thread. --- Cargo.lock | 1 + crates/barbacane-compiler/Cargo.toml | 2 + crates/barbacane-compiler/src/download.rs | 51 +++++++++++++++ crates/barbacane-wasm/src/kafka_client.rs | 31 ++++++++- crates/barbacane-wasm/src/ldap_client.rs | 33 ++++++++-- crates/barbacane-wasm/src/nats_client.rs | 76 ++++++++++++++++++++++- 6 files changed, 184 insertions(+), 10 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 19a1d8b4..d378305a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -427,6 +427,7 @@ dependencies = [ "tar", "tempfile", "thiserror 2.0.18", + "tokio", "toml 0.8.23", "tracing", ] diff --git a/crates/barbacane-compiler/Cargo.toml b/crates/barbacane-compiler/Cargo.toml index 9674f7ad..b11312cc 100644 --- a/crates/barbacane-compiler/Cargo.toml +++ b/crates/barbacane-compiler/Cargo.toml @@ -32,6 +32,8 @@ workspace = true [dev-dependencies] tempfile = { workspace = true } +# `compile` runs inside a runtime; the download tests reproduce that context. +tokio = { workspace = true } criterion = { workspace = true } serde_json = { workspace = true } # Proves a parsed schema is usable by the data plane's validator, which builds diff --git a/crates/barbacane-compiler/src/download.rs b/crates/barbacane-compiler/src/download.rs index 903a78a3..fc76ba12 100644 --- a/crates/barbacane-compiler/src/download.rs +++ b/crates/barbacane-compiler/src/download.rs @@ -208,3 +208,54 @@ mod tests { ); } } + +#[cfg(test)] +mod runtime_context_tests { + use super::*; + + /// `compile` runs inside a tokio runtime, and `reqwest::blocking` refuses + /// to run inside one, on construction and on every request alike. A + /// manifest with a remote plugin could not be compiled at all because of + /// it. The call must reach the network and come back with an error, not + /// take the process down. + #[test] + fn downloading_from_inside_a_runtime_returns_an_error_rather_than_panicking() { + let rt = tokio::runtime::Builder::new_multi_thread() + .worker_threads(1) + .enable_all() + .build() + .expect("runtime"); + + let result = rt.block_on(async { + // A reserved domain that resolves nowhere, so the request fails + // fast without depending on the network being reachable. + download_plugin("https://barbacane-does-not-resolve.invalid/plugin.wasm") + }); + + assert!( + result.is_err(), + "the download should fail, not succeed against an invalid host" + ); + } + + /// The same from a spawned task, which is a worker thread rather than the + /// one `block_on` runs on. + #[test] + fn downloading_from_a_spawned_task_returns_an_error_rather_than_panicking() { + let rt = tokio::runtime::Builder::new_multi_thread() + .worker_threads(2) + .enable_all() + .build() + .expect("runtime"); + + let result = rt.block_on(async { + tokio::spawn(async { + download_plugin("https://barbacane-does-not-resolve.invalid/plugin.wasm") + }) + .await + .expect("the task must not panic") + }); + + assert!(result.is_err()); + } +} diff --git a/crates/barbacane-wasm/src/kafka_client.rs b/crates/barbacane-wasm/src/kafka_client.rs index 9831012f..b93d9cd0 100644 --- a/crates/barbacane-wasm/src/kafka_client.rs +++ b/crates/barbacane-wasm/src/kafka_client.rs @@ -31,7 +31,7 @@ const PRODUCE_TIMEOUT: Duration = Duration::from_secs(10); /// stay alive between publish calls. Connections are created lazily on first /// publish and reused for subsequent messages to the same broker. pub struct KafkaPublisher { - runtime: tokio::runtime::Runtime, + runtime: Option, clients: Mutex>>, /// When false, broker addresses resolving to internal/metadata ranges are /// rejected (SSRF guard). Operators opt in for trusted internal brokers. @@ -48,7 +48,7 @@ impl KafkaPublisher { .build() .map_err(|e| BrokerError::ConnectionFailed(format!("failed to create runtime: {e}")))?; Ok(Self { - runtime, + runtime: Some(runtime), clients: Mutex::new(HashMap::new()), allow_internal_egress, }) @@ -66,7 +66,7 @@ impl KafkaPublisher { payload: &str, headers: BTreeMap, ) -> Result { - self.runtime + self.runtime() .block_on(self.publish(brokers, topic, key, payload, headers)) } @@ -278,3 +278,28 @@ mod tests { assert!(matches!(result, Err(BrokerError::ConnectionFailed(_)))); } } + +impl KafkaPublisher { + /// The runtime, which is present for the whole life of the value and taken + /// only by `Drop`. + fn runtime(&self) -> &tokio::runtime::Runtime { + self.runtime + .as_ref() + .expect("the runtime is taken only while dropping") + } +} + +impl Drop for KafkaPublisher { + /// Hand the runtime to tokio's background shutdown instead of waiting for + /// it here. + /// + /// Dropping a runtime blocks until its workers stop, which tokio refuses + /// inside an async context. This type is reachable from tasks running on + /// the gateway's runtime, so the last reference can fall anywhere, and + /// `shutdown_background` is safe wherever that happens. + fn drop(&mut self) { + if let Some(runtime) = self.runtime.take() { + runtime.shutdown_background(); + } + } +} diff --git a/crates/barbacane-wasm/src/ldap_client.rs b/crates/barbacane-wasm/src/ldap_client.rs index 0ee4b941..02370b06 100644 --- a/crates/barbacane-wasm/src/ldap_client.rs +++ b/crates/barbacane-wasm/src/ldap_client.rs @@ -89,7 +89,7 @@ struct CachedConn { /// between calls. Search connections are created lazily and reused while they /// stay open; a closed connection is evicted and re-established on next use. pub struct LdapClient { - runtime: tokio::runtime::Runtime, + runtime: Option, connections: Mutex>, /// When false, directory addresses resolving to internal/metadata ranges are /// rejected (SSRF guard). Operators opt in for trusted internal directories. @@ -106,7 +106,7 @@ impl LdapClient { .build() .map_err(|e| LdapError::ConnectionFailed(format!("failed to create runtime: {e}")))?; Ok(Self { - runtime, + runtime: Some(runtime), connections: Mutex::new(HashMap::new()), allow_internal_egress, }) @@ -117,7 +117,7 @@ impl LdapClient { /// Must be called from a thread that is NOT inside a tokio runtime context /// (e.g. from within `std::thread::scope`). pub fn bind_blocking(&self, req: &LdapBindRequest) -> Result { - self.runtime.block_on(self.bind(req)) + self.runtime().block_on(self.bind(req)) } /// Blocking search for use from sync WASM host functions. `plugin` names the @@ -130,7 +130,7 @@ impl LdapClient { plugin: &str, req: &LdapSearchRequest, ) -> Result { - self.runtime.block_on(self.search(plugin, req)) + self.runtime().block_on(self.search(plugin, req)) } /// Verify credentials with a simple bind on a fresh connection. @@ -744,3 +744,28 @@ mod tests { assert_eq!(victim, Some(keys[0].clone())); } } + +impl LdapClient { + /// The runtime, which is present for the whole life of the value and taken + /// only by `Drop`. + fn runtime(&self) -> &tokio::runtime::Runtime { + self.runtime + .as_ref() + .expect("the runtime is taken only while dropping") + } +} + +impl Drop for LdapClient { + /// Hand the runtime to tokio's background shutdown instead of waiting for + /// it here. + /// + /// Dropping a runtime blocks until its workers stop, which tokio refuses + /// inside an async context. This type is reachable from tasks running on + /// the gateway's runtime, so the last reference can fall anywhere, and + /// `shutdown_background` is safe wherever that happens. + fn drop(&mut self) { + if let Some(runtime) = self.runtime.take() { + runtime.shutdown_background(); + } + } +} diff --git a/crates/barbacane-wasm/src/nats_client.rs b/crates/barbacane-wasm/src/nats_client.rs index b5fd9bc4..ba74ff0c 100644 --- a/crates/barbacane-wasm/src/nats_client.rs +++ b/crates/barbacane-wasm/src/nats_client.rs @@ -29,7 +29,7 @@ const PUBLISH_TIMEOUT: Duration = Duration::from_secs(10); /// (heartbeats, reconnection) stay alive between publish calls. Connections are /// created lazily on first publish and reused for subsequent messages to the same server. pub struct NatsPublisher { - runtime: tokio::runtime::Runtime, + runtime: Option, connections: Mutex>, /// When false, server addresses resolving to internal/metadata ranges are /// rejected (SSRF guard). Operators opt in for trusted internal servers. @@ -46,7 +46,7 @@ impl NatsPublisher { .build() .map_err(|e| BrokerError::ConnectionFailed(format!("failed to create runtime: {e}")))?; Ok(Self { - runtime, + runtime: Some(runtime), connections: Mutex::new(HashMap::new()), allow_internal_egress, }) @@ -63,7 +63,7 @@ impl NatsPublisher { payload: Bytes, headers: BTreeMap, ) -> Result { - self.runtime + self.runtime() .block_on(self.publish(url, subject, payload, headers)) } @@ -348,3 +348,73 @@ mod tests { assert!(matches!(result, Err(BrokerError::ConnectionFailed(_)))); } } + +impl NatsPublisher { + /// The runtime, which is present for the whole life of the value and taken + /// only by `Drop`. + fn runtime(&self) -> &tokio::runtime::Runtime { + self.runtime + .as_ref() + .expect("the runtime is taken only while dropping") + } +} + +impl Drop for NatsPublisher { + /// Hand the runtime to tokio's background shutdown instead of waiting for + /// it here. + /// + /// Dropping a runtime blocks until its workers stop, which tokio refuses + /// inside an async context. This type is reachable from tasks running on + /// the gateway's runtime, so the last reference can fall anywhere, and + /// `shutdown_background` is safe wherever that happens. + fn drop(&mut self) { + if let Some(runtime) = self.runtime.take() { + runtime.shutdown_background(); + } + } +} + +#[cfg(test)] +mod drop_safety_tests { + use super::*; + + /// The gateway holds this behind an `Arc` that background tasks clone, so + /// the last reference can fall on a tokio worker. Dropping a runtime there + /// blocks, which tokio refuses, and the process died on an ordinary + /// shutdown because of it. + #[test] + fn dropping_inside_a_runtime_does_not_panic() { + let outer = tokio::runtime::Builder::new_multi_thread() + .worker_threads(1) + .enable_all() + .build() + .expect("outer runtime"); + + outer.block_on(async { + let publisher = NatsPublisher::new(true).expect("publisher"); + drop(publisher); + }); + } + + /// And from a spawned task, which is where the gateway's eviction and + /// hot-reload tasks would drop it. + #[test] + fn dropping_inside_a_spawned_task_does_not_panic() { + let outer = tokio::runtime::Builder::new_multi_thread() + .worker_threads(2) + .enable_all() + .build() + .expect("outer runtime"); + + outer.block_on(async { + let shared = std::sync::Arc::new(NatsPublisher::new(true).expect("publisher")); + let held = shared.clone(); + let task = tokio::spawn(async move { + // The task outlives the local reference, so its drop is last. + drop(held); + }); + drop(shared); + task.await.expect("task"); + }); + } +} From c9d7da0f865317b73164dc8706cd614774a21010 Mon Sep 17 00:00:00 2001 From: Nicolas Dreno Date: Mon, 21 Sep 2026 10:32:24 +0200 Subject: [PATCH 3/4] style: keep the impls above the test module Clippy's items_after_test_module rejected the accessor and Drop blocks, which landed at the end of each file after the tests. Found by CI and not locally because I ran clippy as --lib --bins while CI runs --all-targets, so the lib-test target was never linted here. --- crates/barbacane-wasm/src/kafka_client.rs | 50 ++++---- crates/barbacane-wasm/src/ldap_client.rs | 50 ++++---- crates/barbacane-wasm/src/nats_client.rs | 140 +++++++++++----------- 3 files changed, 120 insertions(+), 120 deletions(-) diff --git a/crates/barbacane-wasm/src/kafka_client.rs b/crates/barbacane-wasm/src/kafka_client.rs index b93d9cd0..2c299fe9 100644 --- a/crates/barbacane-wasm/src/kafka_client.rs +++ b/crates/barbacane-wasm/src/kafka_client.rs @@ -189,6 +189,31 @@ impl KafkaPublisher { } } +impl KafkaPublisher { + /// The runtime, which is present for the whole life of the value and taken + /// only by `Drop`. + fn runtime(&self) -> &tokio::runtime::Runtime { + self.runtime + .as_ref() + .expect("the runtime is taken only while dropping") + } +} + +impl Drop for KafkaPublisher { + /// Hand the runtime to tokio's background shutdown instead of waiting for + /// it here. + /// + /// Dropping a runtime blocks until its workers stop, which tokio refuses + /// inside an async context. This type is reachable from tasks running on + /// the gateway's runtime, so the last reference can fall anywhere, and + /// `shutdown_background` is safe wherever that happens. + fn drop(&mut self) { + if let Some(runtime) = self.runtime.take() { + runtime.shutdown_background(); + } + } +} + #[cfg(test)] mod tests { use super::*; @@ -278,28 +303,3 @@ mod tests { assert!(matches!(result, Err(BrokerError::ConnectionFailed(_)))); } } - -impl KafkaPublisher { - /// The runtime, which is present for the whole life of the value and taken - /// only by `Drop`. - fn runtime(&self) -> &tokio::runtime::Runtime { - self.runtime - .as_ref() - .expect("the runtime is taken only while dropping") - } -} - -impl Drop for KafkaPublisher { - /// Hand the runtime to tokio's background shutdown instead of waiting for - /// it here. - /// - /// Dropping a runtime blocks until its workers stop, which tokio refuses - /// inside an async context. This type is reachable from tasks running on - /// the gateway's runtime, so the last reference can fall anywhere, and - /// `shutdown_background` is safe wherever that happens. - fn drop(&mut self) { - if let Some(runtime) = self.runtime.take() { - runtime.shutdown_background(); - } - } -} diff --git a/crates/barbacane-wasm/src/ldap_client.rs b/crates/barbacane-wasm/src/ldap_client.rs index 02370b06..c1b1fc0a 100644 --- a/crates/barbacane-wasm/src/ldap_client.rs +++ b/crates/barbacane-wasm/src/ldap_client.rs @@ -491,6 +491,31 @@ fn map_search_error(e: ldap3::LdapError) -> LdapError { } } +impl LdapClient { + /// The runtime, which is present for the whole life of the value and taken + /// only by `Drop`. + fn runtime(&self) -> &tokio::runtime::Runtime { + self.runtime + .as_ref() + .expect("the runtime is taken only while dropping") + } +} + +impl Drop for LdapClient { + /// Hand the runtime to tokio's background shutdown instead of waiting for + /// it here. + /// + /// Dropping a runtime blocks until its workers stop, which tokio refuses + /// inside an async context. This type is reachable from tasks running on + /// the gateway's runtime, so the last reference can fall anywhere, and + /// `shutdown_background` is safe wherever that happens. + fn drop(&mut self) { + if let Some(runtime) = self.runtime.take() { + runtime.shutdown_background(); + } + } +} + #[cfg(test)] mod tests { use super::*; @@ -744,28 +769,3 @@ mod tests { assert_eq!(victim, Some(keys[0].clone())); } } - -impl LdapClient { - /// The runtime, which is present for the whole life of the value and taken - /// only by `Drop`. - fn runtime(&self) -> &tokio::runtime::Runtime { - self.runtime - .as_ref() - .expect("the runtime is taken only while dropping") - } -} - -impl Drop for LdapClient { - /// Hand the runtime to tokio's background shutdown instead of waiting for - /// it here. - /// - /// Dropping a runtime blocks until its workers stop, which tokio refuses - /// inside an async context. This type is reachable from tasks running on - /// the gateway's runtime, so the last reference can fall anywhere, and - /// `shutdown_background` is safe wherever that happens. - fn drop(&mut self) { - if let Some(runtime) = self.runtime.take() { - runtime.shutdown_background(); - } - } -} diff --git a/crates/barbacane-wasm/src/nats_client.rs b/crates/barbacane-wasm/src/nats_client.rs index ba74ff0c..52985154 100644 --- a/crates/barbacane-wasm/src/nats_client.rs +++ b/crates/barbacane-wasm/src/nats_client.rs @@ -202,6 +202,76 @@ fn pinned_server_addrs( .map_err(|e| BrokerError::ConnectionFailed(format!("invalid pinned NATS address: {e}"))) } +impl NatsPublisher { + /// The runtime, which is present for the whole life of the value and taken + /// only by `Drop`. + fn runtime(&self) -> &tokio::runtime::Runtime { + self.runtime + .as_ref() + .expect("the runtime is taken only while dropping") + } +} + +impl Drop for NatsPublisher { + /// Hand the runtime to tokio's background shutdown instead of waiting for + /// it here. + /// + /// Dropping a runtime blocks until its workers stop, which tokio refuses + /// inside an async context. This type is reachable from tasks running on + /// the gateway's runtime, so the last reference can fall anywhere, and + /// `shutdown_background` is safe wherever that happens. + fn drop(&mut self) { + if let Some(runtime) = self.runtime.take() { + runtime.shutdown_background(); + } + } +} + +#[cfg(test)] +mod drop_safety_tests { + use super::*; + + /// The gateway holds this behind an `Arc` that background tasks clone, so + /// the last reference can fall on a tokio worker. Dropping a runtime there + /// blocks, which tokio refuses, and the process died on an ordinary + /// shutdown because of it. + #[test] + fn dropping_inside_a_runtime_does_not_panic() { + let outer = tokio::runtime::Builder::new_multi_thread() + .worker_threads(1) + .enable_all() + .build() + .expect("outer runtime"); + + outer.block_on(async { + let publisher = NatsPublisher::new(true).expect("publisher"); + drop(publisher); + }); + } + + /// And from a spawned task, which is where the gateway's eviction and + /// hot-reload tasks would drop it. + #[test] + fn dropping_inside_a_spawned_task_does_not_panic() { + let outer = tokio::runtime::Builder::new_multi_thread() + .worker_threads(2) + .enable_all() + .build() + .expect("outer runtime"); + + outer.block_on(async { + let shared = std::sync::Arc::new(NatsPublisher::new(true).expect("publisher")); + let held = shared.clone(); + let task = tokio::spawn(async move { + // The task outlives the local reference, so its drop is last. + drop(held); + }); + drop(shared); + task.await.expect("task"); + }); + } +} + #[cfg(test)] mod tests { use super::*; @@ -348,73 +418,3 @@ mod tests { assert!(matches!(result, Err(BrokerError::ConnectionFailed(_)))); } } - -impl NatsPublisher { - /// The runtime, which is present for the whole life of the value and taken - /// only by `Drop`. - fn runtime(&self) -> &tokio::runtime::Runtime { - self.runtime - .as_ref() - .expect("the runtime is taken only while dropping") - } -} - -impl Drop for NatsPublisher { - /// Hand the runtime to tokio's background shutdown instead of waiting for - /// it here. - /// - /// Dropping a runtime blocks until its workers stop, which tokio refuses - /// inside an async context. This type is reachable from tasks running on - /// the gateway's runtime, so the last reference can fall anywhere, and - /// `shutdown_background` is safe wherever that happens. - fn drop(&mut self) { - if let Some(runtime) = self.runtime.take() { - runtime.shutdown_background(); - } - } -} - -#[cfg(test)] -mod drop_safety_tests { - use super::*; - - /// The gateway holds this behind an `Arc` that background tasks clone, so - /// the last reference can fall on a tokio worker. Dropping a runtime there - /// blocks, which tokio refuses, and the process died on an ordinary - /// shutdown because of it. - #[test] - fn dropping_inside_a_runtime_does_not_panic() { - let outer = tokio::runtime::Builder::new_multi_thread() - .worker_threads(1) - .enable_all() - .build() - .expect("outer runtime"); - - outer.block_on(async { - let publisher = NatsPublisher::new(true).expect("publisher"); - drop(publisher); - }); - } - - /// And from a spawned task, which is where the gateway's eviction and - /// hot-reload tasks would drop it. - #[test] - fn dropping_inside_a_spawned_task_does_not_panic() { - let outer = tokio::runtime::Builder::new_multi_thread() - .worker_threads(2) - .enable_all() - .build() - .expect("outer runtime"); - - outer.block_on(async { - let shared = std::sync::Arc::new(NatsPublisher::new(true).expect("publisher")); - let held = shared.clone(); - let task = tokio::spawn(async move { - // The task outlives the local reference, so its drop is last. - drop(held); - }); - drop(shared); - task.await.expect("task"); - }); - } -} From 57dd7036659fe231d8fda0a59794d97a3a38418c Mon Sep 17 00:00:00 2001 From: Nicolas Dreno Date: Mon, 21 Sep 2026 11:45:17 +0200 Subject: [PATCH 4/4] test: cover the drop-in-async invariant for every runtime-owning client Review asked whether a background task holding the gateway could drop it on a runtime thread after `drop_off_runtime` has released the local reference. It can, and that is safe because each runtime-owning client shuts its runtime down in the background rather than waiting for it. Only `NatsPublisher` proved that; the same two tests now cover `KafkaPublisher` and `LdapClient`. Removing the `shutdown_background` call turns all six red, so they assert the invariant rather than describing it. `gateway_binary` resolves the same path it always did. The comments labelled the two `pop()` calls by the directory each left behind, which reads as though the file name is never removed; they now say what each call drops. The lifecycle tests no longer skip in silence. A missing binary or an unbuildable fixture returned early, which `cargo test` renders as a pass with the reason hidden in captured output. `build_artifact` returns the reason and both tests fail with it, matching every other test in the crate, which already refuses to run without the binary. --- .../barbacane-test/tests/process_lifecycle.rs | 54 +++++++++++-------- crates/barbacane-wasm/src/kafka_client.rs | 44 +++++++++++++++ crates/barbacane-wasm/src/ldap_client.rs | 44 +++++++++++++++ 3 files changed, 119 insertions(+), 23 deletions(-) diff --git a/crates/barbacane-test/tests/process_lifecycle.rs b/crates/barbacane-test/tests/process_lifecycle.rs index 34b3ece7..d7e7839a 100644 --- a/crates/barbacane-test/tests/process_lifecycle.rs +++ b/crates/barbacane-test/tests/process_lifecycle.rs @@ -13,22 +13,33 @@ use std::process::{Command, Stdio}; use std::time::{Duration, Instant}; /// The gateway binary built alongside these tests. +/// +/// A test executable lives at `target//deps/-`, so the +/// gateway is two levels up: one pop drops the file name, the second drops +/// `deps`. fn gateway_binary() -> std::path::PathBuf { let mut dir = std::env::current_exe().expect("test binary path"); + dir.pop(); // the test executable's own file name dir.pop(); // deps/ - dir.pop(); // debug/ or release/ dir.join("barbacane") } /// The smallest artifact the repository can build: one mock route. -fn build_artifact(dir: &std::path::Path) -> Option { +/// +/// Returns the reason it could not, never a bare `None`: a test that skips +/// without saying so reads as a pass, and `cargo test` hides the output of one. +fn build_artifact(dir: &std::path::Path) -> Result { let repo = std::path::Path::new(env!("CARGO_MANIFEST_DIR")) - .parent()? - .parent()? + .parent() + .and_then(std::path::Path::parent) + .ok_or("the repository root is not two levels above the crate")? .to_path_buf(); let mock = repo.join("plugins/mock/mock.wasm"); if !mock.exists() { - return None; + return Err(format!( + "{} is not built; run `make plugins`", + mock.display() + )); } let manifest = dir.join("barbacane.yaml"); @@ -36,7 +47,7 @@ fn build_artifact(dir: &std::path::Path) -> Option { &manifest, format!("plugins:\n mock:\n path: {}\n", mock.display()), ) - .ok()?; + .map_err(|e| format!("could not write the manifest: {e}"))?; let spec = dir.join("api.yaml"); std::fs::write( @@ -51,7 +62,7 @@ paths: responses: { "200": { description: ok } } "#, ) - .ok()?; + .map_err(|e| format!("could not write the spec: {e}"))?; let out = dir.join("api.bca"); let status = Command::new(gateway_binary()) @@ -62,15 +73,14 @@ paths: .arg("-o") .arg(&out) .output() - .ok()?; + .map_err(|e| format!("could not run the compiler: {e}"))?; if !status.status.success() { - eprintln!( - "skipping: could not compile the fixture: {}", + return Err(format!( + "compiling the fixture failed: {}", String::from_utf8_lossy(&status.stderr) - ); - return None; + )); } - Some(out) + Ok(out) } fn wait_for_port(port: u16, limit: Duration) -> bool { @@ -90,14 +100,14 @@ fn wait_for_port(port: u16, limit: Duration) -> bool { #[test] fn sigterm_on_a_serving_gateway_exits_zero() { let binary = gateway_binary(); - if !binary.exists() { - eprintln!("skipping: {} is not built", binary.display()); - return; - } + assert!( + binary.exists(), + "{} is not built. Every test in this crate drives the real binary, so a \ + missing one is a broken run, not a reason to pass quietly", + binary.display() + ); let dir = tempfile::tempdir().expect("temp dir"); - let Some(artifact) = build_artifact(dir.path()) else { - return; - }; + let artifact = build_artifact(dir.path()).expect("build the fixture artifact"); let mut child = Command::new(&binary) .arg("serve") @@ -157,9 +167,7 @@ fn a_taken_port_reports_the_cause_without_panicking() { return; } let dir = tempfile::tempdir().expect("temp dir"); - let Some(artifact) = build_artifact(dir.path()) else { - return; - }; + let artifact = build_artifact(dir.path()).expect("build the fixture artifact"); let held = std::net::TcpListener::bind("127.0.0.1:34202").expect("hold the port"); diff --git a/crates/barbacane-wasm/src/kafka_client.rs b/crates/barbacane-wasm/src/kafka_client.rs index 2c299fe9..586c9e85 100644 --- a/crates/barbacane-wasm/src/kafka_client.rs +++ b/crates/barbacane-wasm/src/kafka_client.rs @@ -214,6 +214,50 @@ impl Drop for KafkaPublisher { } } +#[cfg(test)] +mod drop_safety_tests { + use super::*; + + /// The last reference can fall on a runtime thread, since the gateway holds + /// this type behind an `Arc` that background tasks clone. Dropping a + /// runtime there is what tokio refuses, so the drop must not do it. + #[test] + fn dropping_inside_a_runtime_does_not_panic() { + let runtime = tokio::runtime::Builder::new_multi_thread() + .worker_threads(1) + .enable_all() + .build() + .expect("runtime"); + + runtime.block_on(async { + let client = KafkaPublisher::new(true).expect("kafka publisher"); + drop(client); + }); + } + + /// And from a spawned task, which is where the gateway's eviction and + /// hot-reload tasks would drop it. + #[test] + fn dropping_inside_a_spawned_task_does_not_panic() { + let outer = tokio::runtime::Builder::new_multi_thread() + .worker_threads(2) + .enable_all() + .build() + .expect("outer runtime"); + + outer.block_on(async { + let shared = std::sync::Arc::new(KafkaPublisher::new(true).expect("kafka publisher")); + let held = shared.clone(); + let task = tokio::spawn(async move { + // The task outlives the local reference, so its drop is last. + drop(held); + }); + drop(shared); + task.await.expect("task"); + }); + } +} + #[cfg(test)] mod tests { use super::*; diff --git a/crates/barbacane-wasm/src/ldap_client.rs b/crates/barbacane-wasm/src/ldap_client.rs index c1b1fc0a..4c331ffe 100644 --- a/crates/barbacane-wasm/src/ldap_client.rs +++ b/crates/barbacane-wasm/src/ldap_client.rs @@ -516,6 +516,50 @@ impl Drop for LdapClient { } } +#[cfg(test)] +mod drop_safety_tests { + use super::*; + + /// The last reference can fall on a runtime thread, since the gateway holds + /// this type behind an `Arc` that background tasks clone. Dropping a + /// runtime there is what tokio refuses, so the drop must not do it. + #[test] + fn dropping_inside_a_runtime_does_not_panic() { + let runtime = tokio::runtime::Builder::new_multi_thread() + .worker_threads(1) + .enable_all() + .build() + .expect("runtime"); + + runtime.block_on(async { + let client = LdapClient::new(true).expect("ldap client"); + drop(client); + }); + } + + /// And from a spawned task, which is where the gateway's eviction and + /// hot-reload tasks would drop it. + #[test] + fn dropping_inside_a_spawned_task_does_not_panic() { + let outer = tokio::runtime::Builder::new_multi_thread() + .worker_threads(2) + .enable_all() + .build() + .expect("outer runtime"); + + outer.block_on(async { + let shared = std::sync::Arc::new(LdapClient::new(true).expect("ldap client")); + let held = shared.clone(); + let task = tokio::spawn(async move { + // The task outlives the local reference, so its drop is last. + drop(held); + }); + drop(shared); + task.await.expect("task"); + }); + } +} + #[cfg(test)] mod tests { use super::*;