use std::time::Duration;
use dvb_si::demux::SiDemux;
use dvb_stream::SectionStream;
use futures_core::stream::Stream;
use std::pin::Pin;
use std::task::{Context, Poll};
use tokio::io::AsyncRead;
async fn collect_section_stream<R: AsyncRead + Unpin>(
mut stream: SectionStream<R>,
) -> Vec<dvb_si::demux::SectionEvent> {
tokio::time::timeout(Duration::from_secs(60), async {
let mut events = Vec::new();
while let Some(ev) = futures_util_poll_once(&mut stream).await {
events.push(ev);
}
events
})
.await
.expect("SectionStream stalled — hang guard (issue #807) fired after 60 s")
}
async fn futures_util_poll_once<S: Stream + Unpin>(stream: &mut S) -> Option<S::Item> {
std::future::poll_fn(|cx| Pin::new(&mut *stream).poll_next(cx)).await
}
fn sync_oracle(data: &[u8]) -> Vec<dvb_si::demux::SectionEvent> {
let mut demux = SiDemux::builder().build();
let mut events = Vec::new();
for pkt in data.chunks_exact(188) {
for ev in demux.feed(pkt) {
events.push(ev);
}
}
events
}
fn m6_fixture_path() -> &'static str {
concat!(env!("CARGO_MANIFEST_DIR"), "/../fixtures/ts/m6-single.ts")
}
#[tokio::test]
async fn section_stream_matches_sync_oracle_m6() {
let path = m6_fixture_path();
let data = std::fs::read(path).expect("m6-single.ts fixture not found");
let oracle = sync_oracle(&data);
assert!(
!oracle.is_empty(),
"oracle produced no events — fixture empty?"
);
let cursor = tokio::io::BufReader::new(std::io::Cursor::new(data));
let stream = SectionStream::new(cursor);
let async_events = collect_section_stream(stream).await;
assert_eq!(
async_events.len(),
oracle.len(),
"event count mismatch: async={} oracle={}",
async_events.len(),
oracle.len()
);
for (i, (got, want)) in async_events.iter().zip(oracle.iter()).enumerate() {
assert_eq!(
got.pid(),
want.pid(),
"event[{i}] pid mismatch: got={:?} want={:?}",
got.pid(),
want.pid()
);
assert_eq!(
got.table_id(),
want.table_id(),
"event[{i}] table_id mismatch"
);
assert_eq!(got.bytes(), want.bytes(), "event[{i}] bytes mismatch");
}
}
#[tokio::test]
async fn section_stream_in_memory_cursor_produces_events() {
let path = m6_fixture_path();
let data = std::fs::read(path).expect("m6-single.ts fixture not found");
let cursor = std::io::Cursor::new(data);
let stream = SectionStream::new(cursor);
let events = tokio::time::timeout(Duration::from_secs(60), async {
let mut events = Vec::new();
let mut stream = stream;
loop {
let item = std::future::poll_fn(|cx| Pin::new(&mut stream).poll_next(cx)).await;
match item {
Some(ev) => events.push(ev),
None => break,
}
}
events
})
.await
.expect("hang guard (issue #807) fired: stream never completed within 60s");
assert!(
!events.is_empty(),
"expected at least one section event from m6-single.ts"
);
}
struct OneByteAtATime {
data: Vec<u8>,
pos: usize,
}
impl AsyncRead for OneByteAtATime {
fn poll_read(
mut self: Pin<&mut Self>,
_cx: &mut Context<'_>,
buf: &mut tokio::io::ReadBuf<'_>,
) -> Poll<std::io::Result<()>> {
if self.pos >= self.data.len() {
return Poll::Ready(Ok(())); }
buf.put_slice(&self.data[self.pos..self.pos + 1]);
self.pos += 1;
Poll::Ready(Ok(()))
}
}
#[tokio::test]
async fn section_stream_one_byte_at_a_time_matches_oracle() {
let path = m6_fixture_path();
let data = std::fs::read(path).expect("m6-single.ts fixture not found");
let oracle = sync_oracle(&data);
let reader = OneByteAtATime {
data: data.clone(),
pos: 0,
};
let stream = SectionStream::new(reader);
let async_events = tokio::time::timeout(Duration::from_secs(60), async {
let mut events = Vec::new();
let mut stream = stream;
loop {
let item = std::future::poll_fn(|cx| Pin::new(&mut stream).poll_next(cx)).await;
match item {
Some(ev) => events.push(ev),
None => break,
}
}
events
})
.await
.expect("hang guard (issue #807) fired: one-byte-at-a-time stream never completed within 60s");
assert_eq!(
async_events.len(),
oracle.len(),
"one-byte reader: event count mismatch async={} oracle={}",
async_events.len(),
oracle.len()
);
for (i, (got, want)) in async_events.iter().zip(oracle.iter()).enumerate() {
assert_eq!(got.pid(), want.pid(), "event[{i}] pid");
assert_eq!(got.table_id(), want.table_id(), "event[{i}] table_id");
assert_eq!(got.bytes(), want.bytes(), "event[{i}] bytes");
}
}
#[tokio::test]
async fn section_stream_stats_after_completion() {
let path = m6_fixture_path();
let data = std::fs::read(path).expect("m6-single.ts fixture not found");
let data_len = data.len();
let cursor = std::io::Cursor::new(data);
let mut stream = SectionStream::new(cursor);
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("hang guard (issue #807) fired: stream never completed within 60s");
let stats = stream.stats();
assert!(
stats.packets >= (data_len / 188) as u64,
"expected at least {} packets, got {}",
data_len / 188,
stats.packets
);
assert!(stats.emitted > 0, "expected emitted > 0");
assert_eq!(stats.crc_failures, 0, "expected zero CRC failures");
}