#[derive(Debug, Default)]
pub struct Utf8StreamDecoder {
pending: Vec<u8>,
}
impl Utf8StreamDecoder {
pub(crate) fn new() -> Self {
Self::default()
}
pub(crate) fn push(&mut self, bytes: &[u8]) -> String {
self.pending.extend_from_slice(bytes);
let mut out = String::with_capacity(bytes.len());
loop {
match std::str::from_utf8(&self.pending) {
Ok(text) => {
out.push_str(text);
self.pending.clear();
break;
}
Err(err) => {
let valid = err.valid_up_to();
if let Some(valid_bytes) = self.pending.get(..valid) {
out.push_str(&String::from_utf8_lossy(valid_bytes));
}
match err.error_len() {
Some(invalid_len) => {
out.push('\u{FFFD}');
self.pending.drain(..valid + invalid_len);
}
None => {
self.pending.drain(..valid);
break;
}
}
}
}
}
out
}
pub(crate) fn push_bytes(&mut self, bytes: &[u8], out: &mut Vec<u8>) {
self.pending.extend_from_slice(bytes);
loop {
match std::str::from_utf8(&self.pending) {
Ok(text) => {
out.extend_from_slice(text.as_bytes());
self.pending.clear();
break;
}
Err(err) => {
let valid = err.valid_up_to();
if let Some(valid_bytes) = self.pending.get(..valid) {
out.extend_from_slice(valid_bytes);
}
match err.error_len() {
Some(invalid_len) => {
out.extend_from_slice("\u{FFFD}".as_bytes());
self.pending.drain(..valid + invalid_len);
}
None => {
self.pending.drain(..valid);
break;
}
}
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn utf8_stream_decoder_push_bytes_matches_push() {
let full = "data: {\"hello\":\"world\"}\n\n".as_bytes();
let mut dec_str = Utf8StreamDecoder::new();
let mut dec_bytes = Utf8StreamDecoder::new();
for (i, &byte) in full.iter().enumerate() {
let mut out = Vec::new();
dec_bytes.push_bytes(std::slice::from_ref(&byte), &mut out);
let s = dec_str.push(std::slice::from_ref(&byte));
assert_eq!(out, s.into_bytes(), "byte {i}: push_bytes must match push");
}
let mut tail = Vec::new();
dec_bytes.push_bytes(&[], &mut tail);
assert!(tail.is_empty());
}
#[test]
fn utf8_stream_decoder_push_bytes_handles_split_multibyte() {
let mut dec = Utf8StreamDecoder::new();
let mut out = Vec::new();
dec.push_bytes(&[0xC3], &mut out);
assert!(out.is_empty(), "incomplete multibyte should produce no output");
dec.push_bytes(&[0xA9, b'h', b'i'], &mut out);
assert_eq!(std::str::from_utf8(&out).unwrap(), "\u{00E9}hi");
}
#[test]
fn utf8_stream_decoder_push_bytes_emits_replacement_for_invalid() {
let mut dec = Utf8StreamDecoder::new();
let mut out = Vec::new();
dec.push_bytes(&[b'a', 0xFF, b'b'], &mut out);
assert_eq!(std::str::from_utf8(&out).unwrap(), "a\u{FFFD}b");
}
}