use std::pin::pin;
use async_stream::try_stream;
use futures::StreamExt;
use crate::MiasmaStream;
pub fn with_bytes_counted(
stream: impl MiasmaStream,
record_size: impl AsyncFnOnce(usize),
) -> impl MiasmaStream {
let mut stream_size = 0;
try_stream! {
let mut stream = pin!(stream);
while let Some(chunk) = stream.next().await {
if let Ok(ref chunk) = chunk {
stream_size += chunk.len();
}
yield chunk?;
}
record_size(stream_size).await;
}
}
#[cfg(test)]
mod test {
use async_stream::try_stream;
use super::*;
async fn drain_stream(stream: impl MiasmaStream) -> String {
let mut out = String::new();
let mut stream = pin!(stream);
while let Some(chunk) = stream.next().await {
out.push_str(str::from_utf8(&chunk.unwrap()).unwrap());
}
out
}
#[tokio::test]
async fn counts_bytes_and_does_not_mutate() {
let stream = try_stream! {
yield "hello".into();
yield " ".into();
yield "world".into();
};
let mut bytes = 0;
let with_size = with_bytes_counted(stream, async |n| bytes = n);
let result = drain_stream(with_size).await;
#[allow(clippy::needless_as_bytes)]
let expected_size = "hello world".as_bytes().len();
assert_eq!(result, "hello world");
assert_eq!(bytes, expected_size);
}
}