From 29dcbf55e05ad3ae8faf292c8dfe479981ca284c Mon Sep 17 00:00:00 2001 From: MattJackson <1085847+MattJackson@users.noreply.github.com> Date: Wed, 15 Apr 2026 03:41:03 +0000 Subject: [PATCH] M2TS implements Stream (read PES frames via TsDemux) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit M2tsStream::read() → TsDemux → CodecParser → PesFrame. M2TS input now works through PES pipeline. MKV input deferred (needs EBML → PES extraction). --- src/mux/m2ts.rs | 117 ++++++++++++++++++++++++++++++++++++++++++++- src/mux/resolve.rs | 11 +++-- 2 files changed, 122 insertions(+), 6 deletions(-) diff --git a/src/mux/m2ts.rs b/src/mux/m2ts.rs index 7cc1222..fb5f7fb 100644 --- a/src/mux/m2ts.rs +++ b/src/mux/m2ts.rs @@ -27,6 +27,13 @@ pub struct M2tsStream { finished: bool, /// Content size in bytes (file size minus header), set for read mode. content_size: Option, + // PES support + demuxer: Option, + parsers: Vec<(u16, Box)>, + pending_frames: std::collections::VecDeque, + pid_to_track: Vec<(u16, usize)>, + pes_buf: Vec, + pes_eof: bool, } impl M2tsStream { @@ -40,9 +47,32 @@ impl M2tsStream { }, finished: false, content_size: None, + demuxer: None, + parsers: Vec::new(), + pending_frames: std::collections::VecDeque::new(), + pid_to_track: Vec::new(), + pes_buf: Vec::new(), + pes_eof: false, } } + fn setup_pes(streams: &[DiscStream]) -> (Vec, Vec<(u16, Box)>, Vec<(u16, usize)>) { + let mut pids = Vec::new(); + let mut parsers: Vec<(u16, Box)> = Vec::new(); + let mut pid_to_track = Vec::new(); + for (i, s) in streams.iter().enumerate() { + let (pid, codec) = match s { + DiscStream::Video(v) => (v.pid, v.codec), + DiscStream::Audio(a) => (a.pid, a.codec), + DiscStream::Subtitle(s) => (s.pid, s.codec), + }; + pids.push(pid); + pid_to_track.push((pid, i)); + parsers.push((pid, super::codec::parser_for_codec(codec))); + } + (pids, parsers, pid_to_track) + } + /// Set stream metadata. Returns self for chaining. pub fn meta(mut self, dt: &DiscTitle) -> Self { self.disc_title = dt.clone(); @@ -61,13 +91,21 @@ impl M2tsStream { if let Ok(Some(m)) = meta::read_header(&mut reader) { let header_end = reader.stream_position()?; let content_size = file_size.saturating_sub(header_end); + let title = m.to_title(); + let (pids, parsers, pid_to_track) = Self::setup_pes(&title.streams); return Ok(Self { - disc_title: m.to_title(), + disc_title: title, mode: Mode::Read { reader: Box::new(reader), }, finished: false, content_size: Some(content_size), + demuxer: if pids.is_empty() { None } else { Some(ts::TsDemuxer::new(&pids)) }, + parsers, + pending_frames: std::collections::VecDeque::new(), + pid_to_track, + pes_buf: vec![0u8; 192 * 1024], + pes_eof: false, }); } @@ -89,6 +127,8 @@ impl M2tsStream { reader.seek(SeekFrom::Start(0))?; + let (pids, parsers, pid_to_track) = Self::setup_pes(&streams); + Ok(Self { disc_title: DiscTitle { duration_secs: duration, @@ -100,10 +140,85 @@ impl M2tsStream { }, finished: false, content_size: Some(file_size), + demuxer: if pids.is_empty() { None } else { Some(ts::TsDemuxer::new(&pids)) }, + parsers, + pending_frames: std::collections::VecDeque::new(), + pid_to_track, + pes_buf: vec![0u8; 192 * 1024], + pes_eof: false, }) } } +impl crate::pes::Stream for M2tsStream { + fn read(&mut self) -> io::Result> { + if let Some(frame) = self.pending_frames.pop_front() { + return Ok(Some(frame)); + } + if self.pes_eof { return Ok(None); } + + loop { + let reader = match &mut self.mode { + Mode::Read { reader } => reader, + _ => return Err(io::Error::new(io::ErrorKind::Unsupported, "not in read mode")), + }; + let mut buf = vec![0u8; 192 * 1024]; + let n = reader.read(&mut buf)?; + if n == 0 { + self.pes_eof = true; + return Ok(None); + } + + if let Some(ref mut demuxer) = self.demuxer { + let packets = demuxer.feed(&buf[..n]); + for pes in &packets { + if let Some((_, track)) = self.pid_to_track.iter().find(|(pid, _)| *pid == pes.pid) { + if let Some((_, parser)) = self.parsers.iter_mut().find(|(pid, _)| *pid == pes.pid) { + for frame in parser.parse(pes) { + self.pending_frames.push_back( + crate::pes::PesFrame::from_codec_frame(*track, frame) + ); + } + } + } + } + } + + if let Some(frame) = self.pending_frames.pop_front() { + return Ok(Some(frame)); + } + } + } + + fn write(&mut self, _frame: &crate::pes::PesFrame) -> io::Result<()> { + Err(io::Error::new(io::ErrorKind::Unsupported, "use M2tsOutputStream for writing")) + } + + fn finish(&mut self) -> io::Result<()> { Ok(()) } + + fn info(&self) -> &crate::disc::DiscTitle { &self.disc_title } + + fn codec_private(&self, track: usize) -> Option> { + let pid = self.pid_to_track.iter() + .find(|(_, idx)| *idx == track) + .map(|(pid, _)| *pid)?; + self.parsers.iter() + .find(|(p, _)| *p == pid) + .and_then(|(_, parser)| parser.codec_private()) + } + + fn headers_ready(&self) -> bool { + for (idx, s) in self.disc_title.streams.iter().enumerate() { + if let crate::disc::Stream::Video(v) = s { + if !v.secondary && self.codec_private(idx).is_none() { + return false; + } + } + } + true + } +} + impl IOStream for M2tsStream { fn info(&self) -> &DiscTitle { &self.disc_title diff --git a/src/mux/resolve.rs b/src/mux/resolve.rs index b51fc3c..ff7a729 100644 --- a/src/mux/resolve.rs +++ b/src/mux/resolve.rs @@ -340,15 +340,16 @@ pub fn input(url: &str, opts: &InputOptions) -> io::Result { validate_file_path(path, "m2ts")?; - // TODO: M2tsStream InputStream (TS demux → PES) - Err(io::Error::new(io::ErrorKind::Unsupported, - "m2ts:// PES input not yet implemented")) + let file = std::fs::File::open(path) + .map_err(|e| io::Error::new(e.kind(), format!("m2ts://{}: {}", path.display(), e)))?; + let reader = std::io::BufReader::with_capacity(IO_BUF_SIZE, file); + Ok(Box::new(M2tsStream::open(reader)?)) } StreamUrl::Mkv { ref path } => { validate_file_path(path, "mkv")?; - // TODO: MkvStream InputStream (MKV demux → PES) + // MKV as PES input requires EBML → PES frame extraction (TODO) Err(io::Error::new(io::ErrorKind::Unsupported, - "mkv:// PES input not yet implemented")) + "mkv:// as input not yet supported — use m2ts:// or iso:// as source")) } StreamUrl::Network { ref addr } => { validate_network_addr(addr)?;