Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 9 additions & 7 deletions src/httpserver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -109,13 +109,15 @@ pub(crate) fn serve(snare: Arc<Snare>) -> Result<(), Box<dyn Error>> {
if let Ok(p) = std::env::var("SNARE_DEBUG_PORT_PATH") {
let value = match &listener {
Listener::Tcp(listener) => listener.local_addr().unwrap().port().to_string(),
Listener::Unix(listener) => listener
.local_addr()
.unwrap()
.as_pathname()
.unwrap()
.to_string_lossy()
.into_owned(),
Listener::Unix(_) => {
// On OpenBSD, Rust's stdlib (at least on 1.98.1) clips the last character
// from the path that Listener::Unix returns.
let lk = snare.conf.lock().unwrap();
let ListenAddr::Unix(path) = &lk.listen else {
panic!()
};
path.to_string_lossy().into_owned()
}
};
std::fs::write(p, value).unwrap();
}
Expand Down
51 changes: 42 additions & 9 deletions src/jobrunner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ use std::{
env,
error::Error,
fs::{self, remove_file},
io::{Read, Write},
io::{self, Read, Write},
os::unix::{
io::{AsRawFd, RawFd},
process::ExitStatusExt,
Expand All @@ -24,7 +24,7 @@ use std::{
time::{Duration, Instant},
};

use libc::c_int;
use libc::{c_int, ioctl};
use nix::{
fcntl::{fcntl, FcntlArg, OFlag},
poll::{poll, PollFd, PollFlags},
Expand Down Expand Up @@ -187,8 +187,10 @@ impl JobRunner {

// Has the HTTP server told us that we should check for new jobs and/or SIGCHLD/SIGHUP
// has been received?
let mut check_exit = false;
match self.pollfds[self.maxjobs * 2].revents() {
Some(flags) if flags == PollFlags::POLLIN => {
check_exit = true;
check_queue = true;
// It's fine for us to drain the event pipe completely: we'll process all the
// events it contains below.
Expand Down Expand Up @@ -218,19 +220,22 @@ impl JobRunner {
}
}

if let Some(Job {
stderr_hup: true,
stdout_hup: true,
..
}) = self.running[i]
{
// Ideally we'd fold this into the `if` below but that requires the 2024 edition of
// Rust.
if !check_exit {
continue;
}

// If a child process has exited, then we check all child processes to see if they
// are the one that has exited.
if let Some(job) = self.running[i].as_mut() {
// In the below, we know from the `let Some(_)` that `self.running[i]` is
// `Some(_)` and the unwrap thus safe.
let mut exited = false;
let mut exited_success = false;
let mut exit_type = "";
let mut exit_code = String::new();
match self.running[i].as_mut().unwrap().child.try_wait() {
match job.child.try_wait() {
Ok(Some(status)) => {
exited = true;
exited_success = status.success();
Expand All @@ -254,7 +259,25 @@ impl JobRunner {
Ok(None) => (),
}
if exited {
// When a child process exits, we might not yet have read all the pending
// data. However, there is a complication: while a child process will, on
// exit, have had all its file handles closed by the kernel, it may have
// passed those handles onto grandchildren which may be keeping them open.
// That means that simply reading until EOF might never finish.
// Fortunately, POSiX guarantees that the in "normal exit" situation all
// the bytes written before exit are available for us to read. We therefore
// read everything that's readily available in the pipe knowing that deals
// well with the "normal exit" situation and doesn't stall us in the
// "grandchildren are alive" situation.
if let Some(mut stderr) = job.child.stderr.take() {
drain_pipe(&mut stderr, job.stderrout.as_file_mut());
}
if let Some(mut stdout) = job.child.stdout.take() {
drain_pipe(&mut stdout, job.stderrout.as_file_mut());
}
if !exited_success {
job.stderr_hup = true;
job.stdout_hup = true;
let job = &self.running[i].as_ref().unwrap();
if job.is_errorcmd {
self.snare.error(&format!(
Expand All @@ -276,6 +299,7 @@ impl JobRunner {
job.child = errorchild;
job.is_errorcmd = true;
job.finish_by = finish_by;
self.update_pollfds();
continue;
}
}
Expand Down Expand Up @@ -681,6 +705,15 @@ struct Job {
rconf: RepoConfig,
}

/// Drain all readily available data from `pipe` into `out`. Note: this function deliberately
/// swallows errors.
fn drain_pipe<P: Read + AsRawFd>(pipe: &mut P, out: &mut impl Write) {
let mut available: c_int = 0;
if unsafe { ioctl(pipe.as_raw_fd(), libc::FIONREAD, &mut available) } != -1 {
io::copy(&mut pipe.take(available as u64), out).ok();
}
}

fn set_nonblock(fd: RawFd) -> Result<(), Box<dyn Error>> {
let mut flags = fcntl(fd, FcntlArg::F_GETFL)?;
flags |= OFlag::O_NONBLOCK.bits();
Expand Down
66 changes: 64 additions & 2 deletions tests/queue.rs
Original file line number Diff line number Diff line change
@@ -1,9 +1,13 @@
use nix::unistd::Uid;
use std::{error::Error, fs::read_dir, thread::sleep};
use std::{
error::Error,
fs::{read_dir, write},
thread::sleep,
};
use tempfile::Builder;

mod common;
use common::run_success;
use common::{run_success, SNARE_PAUSE};

// Note that `sleep_s` has to be a fairly high value as we really hope that snare has finished
// processing all the jobs we've thrown at it. There's no easy way to do that other than waiting
Expand Down Expand Up @@ -117,3 +121,61 @@ fn parallel() {

assert_eq!(run_queue("parallel", 20, "", 1,).unwrap(), 20);
}

#[test]
fn grandchildren_cant_cause_a_stall() {
if Uid::current().is_root() {
println!("test skipped: cannot run as root");
return;
}

// This test is for a bug where when a child process handed off its file descriptors to
// grandchildren, snare didn't check whether the child had exited or not.

let td = Builder::new()
.tempdir_in(env!("CARGO_TARGET_TMPDIR"))
.unwrap();
let tds = td.path().to_str().unwrap();
let mut reqs = Vec::new();
for i in 0..2 {
let path = td.path().to_owned();
reqs.push((
move |port| {
let body = r#"{"repository":{"owner":{"login":"testuser"},"name":"testrepo"}}"#;
Ok(format!(
"POST /payload HTTP/1.1\r\nHost: 127.0.0.1:{port}\r\nContent-Length: {}\r\nContent-Type: application/json\r\nX-GitHub-Event: issues\r\n\r\n{body}",
body.len()
))
},
move |response: String| {
assert!(response.starts_with("HTTP/1.1 200 OK"), "{}", response);
if i == 0 {
sleep(SNARE_PAUSE);
assert!(path.join("started").is_file());
} else {
// The second request is queued while the first command is still running.
write(path.join("release"), "")?;
sleep(SNARE_PAUSE);
assert!(path.join("second").is_file());
}
Ok(())
},
));
}

// The first command exits after "release" appears, leaving a background process holding
// stdout/stderr open. Both loops stop when the temporary directory is removed, even on failure.
run_success(
&format!(
r#"listen = "127.0.0.1:0";
maxjobs = 2;
github {{
match ".*" {{
cmd = "if [ -f {tds}/started ]; then touch {tds}/second; else touch {tds}/started; while [ -d {tds} ] && [ ! -f {tds}/release ]; do sleep 0.01; done; (while [ -d {tds} ]; do sleep 0.1; done) & fi";
}}
}}"#
),
&reqs,
)
.unwrap();
}
Loading