#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Framing {
Complete,
Failed,
Killed,
}
pub(crate) fn framing(bytes: &[u8]) -> Framing {
let terminated = match bytes.iter().rposition(|&b| b == b'\n') {
Some(idx) => bytes.get(..=idx).unwrap_or(bytes),
None => return Framing::Killed,
};
let lines: Vec<&[u8]> = terminated
.split(|&b| b == b'\n')
.filter(|l| !l.is_empty())
.collect();
let Some((last, rest)) = lines.split_last() else {
return Framing::Killed;
};
if event_type(last) != Some("end") {
return Framing::Killed;
}
let mut saw_finish = false;
for line in rest.iter().rev() {
match event_type(line) {
Some("end") => break,
Some("error") => return Framing::Failed,
Some("finish") => saw_finish = true,
_ => {}
}
}
if saw_finish {
Framing::Complete
} else {
Framing::Killed
}
}
pub(crate) fn segment_count(bytes: &[u8]) -> usize {
bytes
.split(|&b| b == b'\n')
.filter(|line| event_type(line) == Some("end"))
.count()
}
pub(crate) fn error_text(bytes: &[u8]) -> Option<String> {
let terminated = match bytes.iter().rposition(|&b| b == b'\n') {
Some(idx) => bytes.get(..=idx).unwrap_or(bytes),
None => return None,
};
let lines: Vec<&[u8]> = terminated
.split(|&b| b == b'\n')
.filter(|l| !l.is_empty())
.collect();
let (last, rest) = lines.split_last()?;
if event_type(last) != Some("end") {
return None;
}
for line in rest.iter().rev() {
match event_type(line) {
Some("end") => return None,
Some("error") => return Some(String::from_utf8_lossy(line).into_owned()),
_ => {}
}
}
None
}
pub(super) fn last_segment_complete(bytes: &[u8]) -> bool {
framing(bytes) == Framing::Complete
}
fn event_type(line: &[u8]) -> Option<&'static str> {
let value: serde_json::Value = serde_json::from_slice(line).ok()?;
match value.get("type")?.as_str()? {
"end" => Some("end"),
"finish" => Some("finish"),
"error" => Some("error"),
_ => None,
}
}
#[cfg(test)]
mod tests;