use super::{Chunk, ExitInfo, Stream, StreamPoll};
const MAX_LINE: usize = 4096;
const TRUNC_MARK: &str = "…[truncated]";
#[derive(Debug, Default)]
struct LineBuf {
partial: Vec<u8>,
truncated: bool,
}
impl LineBuf {
fn push(&mut self, bytes: &[u8]) -> Vec<String> {
let mut lines = Vec::new();
for &b in bytes {
if b == b'\n' {
lines.push(self.flush());
} else if self.partial.len() < MAX_LINE {
self.partial.push(b);
} else {
self.truncated = true;
}
}
lines
}
fn finish(&mut self) -> Option<String> {
(!self.partial.is_empty()).then(|| self.flush())
}
fn flush(&mut self) -> String {
let mut line = String::from_utf8_lossy(&self.partial).into_owned();
if self.truncated {
line.push_str(TRUNC_MARK);
}
self.partial.clear();
self.truncated = false;
line
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum StreamedPoll {
Lines(Vec<String>),
Pending,
Done(StreamedOutcome),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StreamedOutcome {
pub lines: Vec<String>,
pub exit: i32,
pub stderr: String,
}
pub struct Streamed {
stream: Stream,
buf: LineBuf,
stderr: Vec<u8>,
}
impl Streamed {
pub fn new(stream: Stream) -> Self {
Self {
stream,
buf: LineBuf::default(),
stderr: Vec::new(),
}
}
pub fn poll(&mut self) -> StreamedPoll {
let mut lines = Vec::new();
loop {
match self.stream.try_next() {
StreamPoll::Ready(Chunk::Stdout(b)) => lines.extend(self.buf.push(&b)),
StreamPoll::Ready(Chunk::Stderr(b)) => self.stderr.extend_from_slice(&b),
StreamPoll::Ready(Chunk::Exited(e)) => {
lines.extend(self.buf.finish());
return StreamedPoll::Done(self.outcome(lines, e));
}
StreamPoll::Pending => {
return if lines.is_empty() {
StreamedPoll::Pending
} else {
StreamedPoll::Lines(lines)
};
}
}
}
}
#[cfg(test)]
pub(crate) fn from_rx(rx: std::sync::mpsc::Receiver<Chunk>) -> Self {
Self::new(Stream::from_rx(rx))
}
fn outcome(&self, lines: Vec<String>, exit: ExitInfo) -> StreamedOutcome {
StreamedOutcome {
lines,
exit: exit.shell_code(),
stderr: String::from_utf8_lossy(&self.stderr).into_owned(),
}
}
}