diff --git a/src/cmd/src/bin/query_regression_runner/otlp.rs b/src/cmd/src/bin/query_regression_runner/otlp.rs index 5be02b0efe..32e2b06de4 100644 --- a/src/cmd/src/bin/query_regression_runner/otlp.rs +++ b/src/cmd/src/bin/query_regression_runner/otlp.rs @@ -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> { - let status = self.child.try_wait()?; - if status.is_some() { - self.reaped = true; - } - Ok(status) - } - - fn kill_and_wait(&mut self) -> std::io::Result { - 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 { +fn cleanup_child(child: &mut Child) -> std::io::Result { + 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 { 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 { + 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()