use super::{Chunk, 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 struct StreamedLine {
pub text: String,
pub err: bool,
}
pub fn stdout_text(lines: &[StreamedLine]) -> String {
text_of(lines, false)
}
pub fn stderr_text(lines: &[StreamedLine]) -> String {
text_of(lines, true)
}
fn text_of(lines: &[StreamedLine], err: bool) -> String {
lines
.iter()
.filter(|l| l.err == err)
.map(|l| l.text.as_str())
.collect::<Vec<_>>()
.join("\n")
}
fn tag(texts: Vec<String>, err: bool) -> Vec<StreamedLine> {
texts
.into_iter()
.map(|text| StreamedLine { text, err })
.collect()
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum StreamedPoll {
Lines(Vec<StreamedLine>),
Pending,
Done(StreamedOutcome),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StreamedOutcome {
pub lines: Vec<StreamedLine>,
pub exit: i32,
}
pub struct Streamed {
stream: Stream,
out: LineBuf,
err: LineBuf,
}
impl Streamed {
pub fn new(stream: Stream) -> Self {
Self {
stream,
out: LineBuf::default(),
err: LineBuf::default(),
}
}
pub fn poll(&mut self) -> StreamedPoll {
let mut lines = Vec::new();
loop {
match self.stream.try_next() {
StreamPoll::Ready(Chunk::Stdout(b)) => lines.extend(tag(self.out.push(&b), false)),
StreamPoll::Ready(Chunk::Stderr(b)) => lines.extend(tag(self.err.push(&b), true)),
StreamPoll::Ready(Chunk::Exited(e)) => {
lines.extend(tag(self.out.finish().into_iter().collect(), false));
lines.extend(tag(self.err.finish().into_iter().collect(), true));
return StreamedPoll::Done(StreamedOutcome {
lines,
exit: e.shell_code(),
});
}
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))
}
}