fix(test): bound OTLP cleanup failure handling and regression test

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
(cherry picked from commit d57f27008cd5adf18d3cf4cb8e5f1334282592f0)
This commit is contained in:
discord9
2026-09-08 21:13:38 +08:00
parent 929e1e549e
commit fc65f23dc5
+91 -62
View File
@@ -186,25 +186,25 @@ async fn run_otelgen_load(
);
let started = Instant::now();
if load.warmup_seconds > 0 {
let _ = wait_for_child(&mut child, Duration::from_secs(load.warmup_seconds)).await?;
let _ = wait_for_child(&mut child.0, Duration::from_secs(load.warmup_seconds)).await?;
}
let warmed = fetch_otlp_metrics(client, http_port, &clock).await?;
let mut timed_out = false;
if child.try_wait()?.is_none() {
if child.0.try_wait()?.is_none() {
let remaining = load
.duration_seconds
.saturating_sub(load.warmup_seconds)
.saturating_add(60)
.max(60);
if !wait_for_child(&mut child, Duration::from_secs(remaining)).await? {
if !wait_for_child(&mut child.0, Duration::from_secs(remaining)).await? {
timed_out = true;
let _ = child.kill_and_wait()?;
let _ = cleanup_child(&mut child.0)?;
}
}
let final_snapshot = fetch_otlp_metrics(client, http_port, &clock).await?;
let returncode = match child.try_wait()? {
let returncode = match child.0.try_wait()? {
Some(status) => status.code(),
None => child.kill_and_wait()?.code(),
None => cleanup_child(&mut child.0)?.code(),
};
let elapsed_seconds = started.elapsed().as_secs_f64();
Ok(json!({
@@ -219,45 +219,36 @@ async fn run_otelgen_load(
}))
}
struct ChildGuard {
child: Child,
reaped: bool,
}
struct ChildGuard(Child);
impl ChildGuard {
fn new(child: Child) -> Self {
Self {
child,
reaped: false,
}
}
fn try_wait(&mut self) -> std::io::Result<Option<ExitStatus>> {
let status = self.child.try_wait()?;
if status.is_some() {
self.reaped = true;
}
Ok(status)
}
fn kill_and_wait(&mut self) -> std::io::Result<ExitStatus> {
self.child.kill()?;
let status = self.child.wait()?;
self.reaped = true;
Ok(status)
Self(child)
}
}
impl Drop for ChildGuard {
fn drop(&mut self) {
if !self.reaped {
let _ = self.child.kill();
let _ = self.child.wait();
if let Err(error) = cleanup_child(&mut self.0) {
eprintln!("failed to clean up otelgen child: {error}");
}
}
}
async fn wait_for_child(child: &mut ChildGuard, timeout: Duration) -> Result<bool> {
fn cleanup_child(child: &mut Child) -> std::io::Result<ExitStatus> {
match child.kill() {
Ok(()) => child.wait(),
Err(kill_error) => match child.try_wait() {
Ok(Some(status)) => Ok(status),
Ok(None) => Err(kill_error),
Err(status_error) => Err(std::io::Error::other(format!(
"failed to kill child: {kill_error}; failed to check child status: {status_error}"
))),
},
}
}
async fn wait_for_child(child: &mut Child, timeout: Duration) -> Result<bool> {
let deadline = Instant::now() + timeout;
loop {
if child.try_wait()?.is_some() {
@@ -524,22 +515,47 @@ mod tests {
#[cfg(unix)]
#[tokio::test]
async fn metrics_failure_after_spawn_reaps_otelgen() {
use std::io::{Read, Write};
use std::io::{self, ErrorKind, Read, Write};
use std::net::TcpListener;
use std::os::unix::fs::PermissionsExt;
use std::thread;
let temp_dir = tempfile::tempdir().unwrap();
fn accept_before(
listener: &TcpListener,
deadline: Instant,
) -> io::Result<std::net::TcpStream> {
loop {
match listener.accept() {
Ok((stream, _)) => return Ok(stream),
Err(error) if error.kind() == ErrorKind::WouldBlock => {
if Instant::now() >= deadline {
return Err(io::Error::new(
ErrorKind::TimedOut,
"timed out accepting request",
));
}
thread::sleep(Duration::from_millis(10));
}
Err(error) => return Err(error),
}
}
}
let temp_dir = tempfile::Builder::new()
.prefix("otlp child with spaces ")
.tempdir()
.unwrap();
let ready_path = temp_dir.path().join("ready");
let pid_path = temp_dir.path().join("pid");
let otelgen_path = temp_dir.path().join("otelgen");
fs::write(
&otelgen_path,
format!(
"#!/bin/sh\necho $$ > {}\ntouch {}\nwhile :; do :; done\n",
pid_path.display(),
ready_path.display(),
),
r#"#!/bin/sh
cd "$(dirname "$0")" || exit 1
echo "$$" > pid
: > ready
exec sleep 5
"#,
)
.unwrap();
let mut permissions = fs::metadata(&otelgen_path).unwrap().permissions();
@@ -547,38 +563,51 @@ mod tests {
fs::set_permissions(&otelgen_path, permissions).unwrap();
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
listener.set_nonblocking(true).unwrap();
let port = listener.local_addr().unwrap().port();
let server = thread::spawn(move || {
for status in ["200 OK", "500 Internal Server Error"] {
let (mut stream, _) = listener.accept().unwrap();
let mut request = [0; 1024];
let _ = stream.read(&mut request);
if status == "500 Internal Server Error" {
while !ready_path.exists() {
thread::yield_now();
}
let server = thread::spawn(move || -> io::Result<()> {
let mut initial = accept_before(&listener, Instant::now() + Duration::from_secs(2))?;
initial.set_nonblocking(false)?;
initial.set_read_timeout(Some(Duration::from_secs(1)))?;
initial.set_write_timeout(Some(Duration::from_secs(1)))?;
let mut request = [0; 1024];
let _ = initial.read(&mut request)?;
initial
.write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 0\r\nConnection: close\r\n\r\n")?;
let readiness_deadline = Instant::now() + Duration::from_secs(2);
while !ready_path.exists() {
if Instant::now() >= readiness_deadline {
return Err(io::Error::new(ErrorKind::TimedOut, "otelgen did not start"));
}
write!(
stream,
"HTTP/1.1 {status}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
)
.unwrap();
thread::sleep(Duration::from_millis(10));
}
let mut warmed = accept_before(&listener, Instant::now() + Duration::from_secs(2))?;
warmed.set_nonblocking(false)?;
warmed.set_read_timeout(Some(Duration::from_secs(1)))?;
warmed.set_write_timeout(Some(Duration::from_secs(1)))?;
let _ = warmed.read(&mut request)?;
warmed.write_all(b"HTTP/1.1 500 Internal Server Error\r\nContent-Length: 0\r\nConnection: close\r\n\r\n")?;
Ok(())
});
let client = Client::builder().build().unwrap();
let client = Client::builder()
.timeout(Duration::from_secs(3))
.build()
.unwrap();
let mut load = otlp_load_for_test();
load.warmup_seconds = 0;
assert!(
run_otelgen_load(&otelgen_path, port, temp_dir.path(), &load, &client)
.await
.is_err()
);
server.join().unwrap();
let error = run_otelgen_load(&otelgen_path, port, temp_dir.path(), &load, &client)
.await
.unwrap_err();
server.join().unwrap().unwrap();
assert!(error.to_string().contains("500 Internal Server Error"));
let pid = fs::read_to_string(pid_path).unwrap();
assert!(
!std::process::Command::new("sh")
.args(["-c", &format!("kill -0 {pid}")])
.args(["-c", &format!("kill -0 {} 2>/dev/null", pid.trim())])
.status()
.unwrap()
.success()