use super::{Chunk, ExitInfo};
use std::process::Child;
use std::sync::mpsc::{Receiver, TryRecvError};
use std::thread;
use std::time::{Duration, Instant};
const TERM_GRACE: Duration = Duration::from_millis(500);
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum StreamPoll {
Ready(Chunk),
Pending,
}
pub struct Stream {
child: Option<Child>,
rx: Receiver<Chunk>,
exit_emitted: bool,
}
impl Stream {
pub(super) fn new(child: Child, rx: Receiver<Chunk>) -> Self {
Self {
child: Some(child),
rx,
exit_emitted: false,
}
}
pub fn pid(&self) -> Option<u32> {
self.child.as_ref().map(std::process::Child::id)
}
pub fn try_next(&mut self) -> StreamPoll {
if self.exit_emitted {
return StreamPoll::Pending;
}
match self.rx.try_recv() {
Ok(chunk) => StreamPoll::Ready(chunk),
Err(TryRecvError::Empty) => StreamPoll::Pending,
Err(TryRecvError::Disconnected) => {
self.exit_emitted = true;
let status = self.child.take().and_then(|mut c| c.wait().ok());
StreamPoll::Ready(Chunk::Exited(exit_info(status)))
}
}
}
#[cfg(test)]
pub(super) fn from_rx(rx: Receiver<Chunk>) -> Self {
Self {
child: None,
rx,
exit_emitted: false,
}
}
}
impl Iterator for Stream {
type Item = Chunk;
fn next(&mut self) -> Option<Chunk> {
if self.exit_emitted {
return None;
}
if let Ok(chunk) = self.rx.recv() {
Some(chunk)
} else {
self.exit_emitted = true;
let status = self.child.take().and_then(|mut c| c.wait().ok());
Some(Chunk::Exited(exit_info(status)))
}
}
}
impl Drop for Stream {
fn drop(&mut self) {
let Some(mut child) = self.child.take() else {
return;
};
let pid = child.id() as i32;
super::sys::sigterm(pid);
let deadline = Instant::now() + TERM_GRACE;
while Instant::now() < deadline {
if let Ok(Some(_)) = child.try_wait() {
return;
}
thread::sleep(Duration::from_millis(25));
}
let _ = child.kill();
let _ = child.wait();
}
}
pub(super) fn exit_info(status: Option<std::process::ExitStatus>) -> ExitInfo {
use std::os::unix::process::ExitStatusExt;
status.map_or(ExitInfo::Unknown, |s| match (s.code(), s.signal()) {
(Some(c), _) => ExitInfo::Code(c),
(_, Some(sig)) => ExitInfo::Signal(sig),
_ => ExitInfo::Unknown,
})
}