From e33b13a7ca0faa3d9a5d6d4da4b5af5dbfd4d45f Mon Sep 17 00:00:00 2001 From: MattJackson <1085847+MattJackson@users.noreply.github.com> Date: Wed, 15 Apr 2026 03:03:45 +0000 Subject: [PATCH] PES pipeline: InputStream on DiscStream + IsoStream, MkvOutputStream MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - DiscStream: next_frame() reads sectors → decrypts → demuxes → returns PesFrame - IsoStream: same pattern, reads from ISO file - MkvOutputStream: accepts PesFrame, writes MKV via MkvMuxer - InputStream trait: next_frame(), info(), codec_private(), headers_ready() - OutputStream trait: write_frame(), finish() - FileSectorReader: SectorReader backed by a file - PesFrame: track + pts + keyframe + data Old Read/Write IOStream impls preserved for backward compatibility. Next: replace pipe() in CLI to use PES pipeline. --- src/mux/disc.rs | 20 ++++++++++ src/mux/iso.rs | 98 +++++++++++++++++++++++++++++++++++++++++++++++ src/mux/mkvout.rs | 66 +++++++++++++++++++++++++++++++ src/mux/mod.rs | 1 + src/pes.rs | 7 ++++ 5 files changed, 192 insertions(+) create mode 100644 src/mux/mkvout.rs diff --git a/src/mux/disc.rs b/src/mux/disc.rs index cc4e05a..bbc4eba 100644 --- a/src/mux/disc.rs +++ b/src/mux/disc.rs @@ -402,6 +402,26 @@ impl crate::pes::InputStream for DiscStream { fn info(&self) -> &DiscTitle { &self.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.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 Read for DiscStream { diff --git a/src/mux/iso.rs b/src/mux/iso.rs index 3b83da1..5fd2c17 100644 --- a/src/mux/iso.rs +++ b/src/mux/iso.rs @@ -77,6 +77,11 @@ pub struct IsoStream { // Write side iso_writer: Option>>, write_started: bool, + // PES output (for InputStream impl) + demuxer: Option, + parsers: Vec<(u16, Box)>, + pending_frames: std::collections::VecDeque, + pid_to_track: Vec<(u16, usize)>, } impl IsoStream { @@ -116,6 +121,21 @@ impl IsoStream { .collect(); let sectors_remaining = extents.first().map(|e| e.1).unwrap_or(0); + // Set up PES demux from title stream PIDs + let mut pids = Vec::new(); + let mut parsers: Vec<(u16, Box)> = Vec::new(); + let mut pid_to_track = Vec::new(); + for (i, s) in disc_title.streams.iter().enumerate() { + let (pid, codec) = match s { + crate::disc::Stream::Video(v) => (v.pid, v.codec), + crate::disc::Stream::Audio(a) => (a.pid, a.codec), + crate::disc::Stream::Subtitle(s) => (s.pid, s.codec), + }; + pids.push(pid); + pid_to_track.push((pid, i)); + parsers.push((pid, super::codec::parser_for_codec(codec))); + } + Ok(IsoStream { disc_title, disc: Some(disc), @@ -130,6 +150,10 @@ impl IsoStream { decrypt_keys, iso_writer: None, write_started: false, + demuxer: if pids.is_empty() { None } else { Some(super::ts::TsDemuxer::new(&pids)) }, + parsers, + pending_frames: std::collections::VecDeque::new(), + pid_to_track, }) } @@ -154,6 +178,10 @@ impl IsoStream { eof: false, iso_writer: Some(iso_writer), write_started: false, + demuxer: None, + parsers: Vec::new(), + pending_frames: std::collections::VecDeque::new(), + pid_to_track: Vec::new(), }) } @@ -228,6 +256,76 @@ impl IsoStream { } } +impl crate::pes::InputStream for IsoStream { + fn next_frame(&mut self) -> io::Result> { + // Return buffered frame + if let Some(frame) = self.pending_frames.pop_front() { + return Ok(Some(frame)); + } + + if self.eof { + return Ok(None); + } + + // Read until we produce at least one frame + loop { + if !self.read_next_batch()? { + self.eof = true; + return Ok(None); + } + + // Data in batch_buf is already decrypted by fill_next_batch + if let Some(ref mut demuxer) = self.demuxer { + let packets = demuxer.feed(&self.batch_buf[..self.buf_len]); + for pes in &packets { + if let Some((pid_idx, _)) = self.pid_to_track.iter().enumerate() + .find(|(_, (pid, _))| *pid == pes.pid) + { + let track_idx = self.pid_to_track[pid_idx].1; + 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_idx, frame) + ); + } + } + } + } + } + + if let Some(frame) = self.pending_frames.pop_front() { + return Ok(Some(frame)); + } + } + } + + 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 IsoStream { fn info(&self) -> &DiscTitle { &self.disc_title diff --git a/src/mux/mkvout.rs b/src/mux/mkvout.rs new file mode 100644 index 0000000..33fbe6b --- /dev/null +++ b/src/mux/mkvout.rs @@ -0,0 +1,66 @@ +//! MKV output stream — accepts PES frames, writes Matroska container. +//! +//! Implements OutputStream. Takes PES frames directly — no TS demuxing needed. +//! Creates the MKV muxer once codec_private is provided for all tracks. + +use super::mkv::{MkvMuxer, MkvTrack}; +use super::WriteSeek; +use crate::disc::DiscTitle; +use crate::pes::{OutputStream, PesFrame}; +use std::io; + +pub struct MkvOutputStream { + muxer: Option>>, +} + +impl MkvOutputStream { + /// Create an MKV output stream. + /// `codec_privates` provides initialization data per track (from InputStream). + /// Tracks without codec_private get None. + pub fn create( + writer: Box, + title: &DiscTitle, + codec_privates: &[Option>], + ) -> io::Result { + let mut tracks = Vec::new(); + for (idx, s) in title.streams.iter().enumerate() { + let mut track = match s { + crate::disc::Stream::Video(v) => MkvTrack::video(v), + crate::disc::Stream::Audio(a) => MkvTrack::audio(a), + crate::disc::Stream::Subtitle(s) => MkvTrack::subtitle(s), + }; + if let Some(cp) = codec_privates.get(idx).and_then(|c| c.as_ref()) { + track.codec_private = Some(cp.clone()); + } + tracks.push(track); + } + + let muxer = MkvMuxer::new_with_chapters( + writer, + &tracks, + Some(&title.playlist), + title.duration_secs, + &title.chapters, + )?; + + Ok(Self { muxer: Some(muxer) }) + } +} + +impl OutputStream for MkvOutputStream { + fn write_frame(&mut self, frame: &PesFrame) -> io::Result<()> { + if let Some(ref mut muxer) = self.muxer { + muxer.write_frame(frame.track, frame.pts, frame.keyframe, &frame.data) + } else { + Ok(()) + } + } + + fn finish(&mut self) -> io::Result<()> { + if let Some(muxer) = self.muxer.take() { + muxer.finish() + } else { + Ok(()) + } + } +} diff --git a/src/mux/mod.rs b/src/mux/mod.rs index 67b65a1..aecec08 100644 --- a/src/mux/mod.rs +++ b/src/mux/mod.rs @@ -28,6 +28,7 @@ pub mod lookahead; mod m2ts; pub mod meta; pub mod mkv; +pub mod mkvout; mod mkvstream; pub mod network; pub mod null; diff --git a/src/pes.rs b/src/pes.rs index dcc5332..bd0f0ab 100644 --- a/src/pes.rs +++ b/src/pes.rs @@ -39,6 +39,13 @@ pub trait InputStream { /// Stream metadata (tracks, duration, etc). fn info(&self) -> &crate::disc::DiscTitle; + + /// Codec initialization data for a track (SPS/PPS for HEVC, etc). + /// Returns None until enough frames have been parsed. + fn codec_private(&self, track: usize) -> Option>; + + /// True when codec_private is available for all video tracks. + fn headers_ready(&self) -> bool; } /// Output stream — consumes PES frames to any destination.