operon 0.7.0

A workflow engine for parallel, incremental scheduling of DAG-defined multiplex tasks.
Documentation
use std::io::{self, BufRead, BufReader};
use std::os::unix::io::{AsRawFd, FromRawFd, OwnedFd, RawFd};
use std::thread::JoinHandle;

type EmitFn = Box<dyn Fn(String) + Send + 'static>;

struct FdSpec {
    fd: RawFd,
    name: String,
    emit: EmitFn,
}

struct FdHandle {
    original: OwnedFd,
    target_fd: RawFd,
    thread: Option<JoinHandle<()>>,
}

#[derive(Default)]
pub(super) struct FdRedirect {
    specs: Vec<FdSpec>,
}

pub(super) struct FdRedirectHandle {
    handles: Vec<FdHandle>,
}

impl FdRedirectHandle {
    pub(super) fn original_fd(&self, target: &impl AsRawFd) -> Option<std::fs::File> {
        self.handles
            .iter()
            .find(|h| h.target_fd == target.as_raw_fd())
            .and_then(|h| {
                nix::unistd::dup(h.original.as_raw_fd())
                    .ok()
                    .map(|duped| unsafe { std::fs::File::from_raw_fd(duped) })
            })
    }
}

impl FdRedirect {
    pub(super) fn new() -> Self {
        Self::default()
    }

    pub(super) fn with_fd(
        mut self,
        target: &impl AsRawFd,
        name: &str,
        emit: impl Fn(String) + Send + 'static,
    ) -> Self {
        self.specs.push(FdSpec {
            fd: target.as_raw_fd(),
            name: name.to_owned(),
            emit: Box::new(emit),
        });
        self
    }

    pub(super) fn capture(self) -> io::Result<FdRedirectHandle> {
        let mut handles = Vec::with_capacity(self.specs.len());

        for spec in self.specs {
            match Self::capture_fd(spec) {
                Ok(handle) => handles.push(handle),
                Err(e) => {
                    drop(FdRedirectHandle { handles });
                    return Err(e);
                }
            }
        }

        Ok(FdRedirectHandle { handles })
    }

    fn capture_fd(spec: FdSpec) -> io::Result<FdHandle> {
        let original_raw = nix::unistd::dup(spec.fd).map_err(io::Error::from)?;
        let original = unsafe { OwnedFd::from_raw_fd(original_raw) };

        let (read_end, write_end) = nix::unistd::pipe().map_err(io::Error::from)?;

        let _ = nix::unistd::dup2(write_end.as_raw_fd(), spec.fd).map_err(io::Error::from)?;

        drop(write_end);

        let read_file: std::fs::File = read_end.into();
        let thread_name = format!("{}-capture", spec.name);
        let emit = spec.emit;

        let thread = std::thread::Builder::new()
            .name(thread_name)
            .spawn(move || {
                let reader = BufReader::new(read_file);
                for line in reader.lines() {
                    match line {
                        Ok(line) => emit(line),
                        Err(_) => break,
                    }
                }
            })?;

        Ok(FdHandle {
            original,
            target_fd: spec.fd,
            thread: Some(thread),
        })
    }
}

impl Drop for FdRedirectHandle {
    fn drop(&mut self) {
        for handle in &mut self.handles {
            let _ = nix::unistd::dup2(handle.original.as_raw_fd(), handle.target_fd);
            let Some(thread) = handle.thread.take() else {
                continue;
            };
            if let Err(panic) = thread.join() {
                tracing::error!(
                    "fd capture thread for fd {} panicked: {:?}",
                    handle.target_fd,
                    panic,
                );
            }
        }
    }
}

pub(super) fn capture_std_outputs() -> io::Result<FdRedirectHandle> {
    // Capture standard outputs at the Error level,
    // since the standard I/O shouldn't be filtered out by log level filters.
    FdRedirect::new()
        .with_fd(&io::stdout(), "stdout", |line| {
            tracing::error!(target: "stdio::stdout", "{}", line);
        })
        .with_fd(&io::stderr(), "stderr", |line| {
            tracing::error!(target: "stdio::stderr", "{}", line);
        })
        .capture()
}