use std::pin::Pin;
use std::time::Duration;
use dvb_stream::SectionStream;
use futures_core::stream::Stream;
fn make_si_ts_packet(pid: u16, _payload_byte: u8) -> [u8; 188] {
let mut pkt = [0xFFu8; 188];
pkt[0] = 0x47; pkt[1] = 0x40 | (((pid >> 8) & 0x1F) as u8); pkt[2] = (pid & 0xFF) as u8;
pkt[3] = 0x10;
let section: [u8; 15] = [
0x00, 0xB0, 0x0D, 0x00, 0x01, 0xC1, 0x00, 0x00, 0x00, 0xE0, 0x10, 0x00, 0x00, 0x00, 0x00, ];
pkt[4] = 0x00; pkt[5..5 + section.len()].copy_from_slice(§ion);
pkt
}
async fn drain_section_stream<R: tokio::io::AsyncRead + Unpin>(stream: &mut SectionStream<R>) {
tokio::time::timeout(Duration::from_secs(60), async {
loop {
let item = std::future::poll_fn(|cx| Pin::new(&mut *stream).poll_next(cx)).await;
if item.is_none() {
break;
}
}
})
.await
.expect("SectionStream stalled — hang guard (issue #807) fired after 60 s");
}
#[tokio::test]
async fn section_stream_clean_resync_stats() {
let pid = 0x0000u16;
let mut data = Vec::new();
for i in 0..3 {
let pkt = make_si_ts_packet(pid, i);
data.extend_from_slice(&pkt);
}
let cursor = std::io::Cursor::new(data);
let mut stream = SectionStream::new(cursor);
drain_section_stream(&mut stream).await;
let stats = stream.resync_stats();
assert_eq!(stats.resyncs, 1, "expected 1 resync, got {}", stats.resyncs);
assert_eq!(stats.bytes_discarded, 0, "expected 0 bytes discarded");
assert_eq!(stats.desyncs, 0, "expected 0 desyncs");
}
#[tokio::test]
async fn section_stream_leading_junk_increments_resync_counters() {
let pid = 0x0000u16;
let mut data = Vec::new();
data.extend_from_slice(&[0x00u8; 42]);
for i in 0..3 {
let pkt = make_si_ts_packet(pid, i);
data.extend_from_slice(&pkt);
}
let cursor = std::io::Cursor::new(data);
let mut stream = SectionStream::new(cursor);
drain_section_stream(&mut stream).await;
let stats = stream.resync_stats();
assert_eq!(stats.resyncs, 1, "expected 1 resync");
assert_eq!(stats.bytes_discarded, 42, "expected 42 bytes discarded");
assert_eq!(stats.desyncs, 0, "expected 0 desyncs");
}
#[tokio::test]
async fn section_stream_junk_only_discards_all_bytes() {
let data = vec![0x00u8; 500]; let cursor = std::io::Cursor::new(data);
let mut stream = SectionStream::new(cursor);
drain_section_stream(&mut stream).await;
let stats = stream.resync_stats();
assert_eq!(stats.resyncs, 0, "expected 0 resyncs");
assert_eq!(stats.bytes_discarded, 500, "expected 500 bytes discarded");
assert_eq!(stats.desyncs, 0, "expected 0 desyncs");
}
#[tokio::test]
async fn section_stream_mid_stream_corruption_detected() {
let pid = 0x0000u16;
let mut data = Vec::new();
for i in 0..2 {
data.extend_from_slice(&make_si_ts_packet(pid, i));
}
let mut corrupt = make_si_ts_packet(pid, 2);
corrupt[0] = 0x00;
data.extend_from_slice(&corrupt);
for i in 3..5 {
data.extend_from_slice(&make_si_ts_packet(pid, i));
}
let cursor = std::io::Cursor::new(data);
let mut stream = SectionStream::new(cursor);
drain_section_stream(&mut stream).await;
let stats = stream.resync_stats();
assert_eq!(stats.resyncs, 1, "expected 1 resync (initial only)");
assert_eq!(
stats.bytes_discarded, 564,
"expected 564 bytes discarded, got {}",
stats.bytes_discarded
);
assert_eq!(stats.desyncs, 1, "expected 1 desync, got {}", stats.desyncs);
}
use tokio::io::AsyncRead;
struct SlowReader {
data: Vec<u8>,
pos: usize,
}
impl AsyncRead for SlowReader {
fn poll_read(
mut self: Pin<&mut Self>,
_cx: &mut std::task::Context<'_>,
buf: &mut tokio::io::ReadBuf<'_>,
) -> std::task::Poll<std::io::Result<()>> {
if self.pos >= self.data.len() {
return std::task::Poll::Ready(Ok(())); }
let chunk_size = 18usize.min(self.data.len() - self.pos);
buf.put_slice(&self.data[self.pos..self.pos + chunk_size]);
self.pos += chunk_size;
std::task::Poll::Ready(Ok(()))
}
}
#[tokio::test]
async fn section_stream_slow_reader_corruption_recovers() {
let pid = 0x0000u16;
let mut data = Vec::new();
for i in 0..2 {
data.extend_from_slice(&make_si_ts_packet(pid, i));
}
let mut corrupt = make_si_ts_packet(pid, 2);
corrupt[0] = 0x00; data.extend_from_slice(&corrupt);
for i in 3..7 {
data.extend_from_slice(&make_si_ts_packet(pid, i));
}
let reader = SlowReader { data, pos: 0 };
let mut stream = SectionStream::new(reader);
drain_section_stream(&mut stream).await;
let stats = stream.resync_stats();
assert_eq!(
stats.resyncs, 2,
"expected 2 resyncs (initial + post-desync), got {}",
stats.resyncs
);
assert!(
stats.desyncs >= 1,
"expected at least 1 desync, got {}",
stats.desyncs
);
assert!(
stats.bytes_discarded >= 188,
"expected at least 188 bytes discarded, got {}",
stats.bytes_discarded
);
}