#[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 {
use super::*;
const FINISH_END: &[u8] = br#"{"type":"message_start","v":1,"role":"assistant"}
{"type":"finish","reason":"stop"}
{"type":"end"}
"#;
const ERROR_END: &[u8] = br#"{"type":"message_start","v":1,"role":"assistant"}
{"type":"error","kind":"transport","message":"reset"}
{"type":"end"}
"#;
#[test]
fn finish_then_end_is_complete() {
assert!(last_segment_complete(FINISH_END));
}
#[test]
fn error_then_end_is_not_complete() {
assert!(!last_segment_complete(ERROR_END));
}
#[test]
fn no_trailing_end_is_not_complete() {
let jsonl = br#"{"type":"content_delta","index":0,"delta":{"text_delta":"hi"}}
"#;
assert!(!last_segment_complete(jsonl));
}
#[test]
fn empty_or_newline_only_is_not_complete() {
assert!(!last_segment_complete(b""));
assert!(!last_segment_complete(b"\n\n"));
}
#[test]
fn trailing_partial_line_after_end_is_ignored() {
let jsonl = b"{\"type\":\"finish\",\"reason\":\"stop\"}\n{\"type\":\"end\"}\n{partial";
assert!(last_segment_complete(jsonl));
}
#[test]
fn only_latest_segment_decides_complete() {
let jsonl = br#"{"type":"error","kind":"x"}
{"type":"end"}
{"type":"message_start","v":1}
{"type":"finish","reason":"stop"}
{"type":"end"}
"#;
assert!(last_segment_complete(jsonl));
}
#[test]
fn latest_segment_error_after_earlier_finish_is_not_complete() {
let jsonl = br#"{"type":"finish","reason":"stop"}
{"type":"end"}
{"type":"error","kind":"x"}
{"type":"end"}
"#;
assert!(!last_segment_complete(jsonl));
}
#[test]
fn end_without_finish_or_error_is_not_complete() {
assert!(!last_segment_complete(
b"{\"type\":\"message_start\"}\n{\"type\":\"end\"}\n"
));
}
#[test]
fn malformed_last_line_is_not_complete() {
assert!(!last_segment_complete(b"{\"type\":\"finish\"}\nnot json\n"));
}
#[test]
fn framing_classifies_the_three_outcomes() {
assert_eq!(framing(FINISH_END), Framing::Complete);
assert_eq!(framing(ERROR_END), Framing::Failed);
assert_eq!(
framing(b"{\"type\":\"content_delta\",\"index\":0}\n"),
Framing::Killed
);
assert_eq!(
framing(b"{\"type\":\"message_start\"}\n{\"type\":\"end\"}\n"),
Framing::Killed
);
}
#[test]
fn error_text_returns_the_error_line_iff_framing_is_failed() {
let failed = br#"{"type":"error","kind":"http","status":401}
{"type":"end"}
"#;
assert_eq!(
error_text(failed).as_deref(),
Some(r#"{"type":"error","kind":"http","status":401}"#)
);
assert_eq!(framing(failed), Framing::Failed);
assert_eq!(error_text(FINISH_END), None); assert_eq!(error_text(b"{\"type\":\"content_delta\"}\n"), None); assert_eq!(error_text(b""), None); assert_eq!(error_text(b"\n\n"), None); assert_eq!(error_text(b"no trailing newline"), None); }
#[test]
fn error_text_reads_only_the_latest_segment() {
let retried = br#"{"type":"error","kind":"x"}
{"type":"end"}
{"type":"finish","reason":"stop"}
{"type":"end"}
"#;
assert_eq!(error_text(retried), None);
assert_eq!(framing(retried), Framing::Complete);
}
#[test]
fn segment_count_counts_end_events() {
let jsonl = br#"{"type":"error","kind":"x"}
{"type":"end"}
{"type":"finish","reason":"stop"}
{"type":"end"}
"#;
assert_eq!(segment_count(jsonl), 2);
assert_eq!(
segment_count(b"{\"type\":\"content_delta\",\"index\":0}\n"),
0
);
assert_eq!(segment_count(b""), 0);
}
}