From bc9bf022ea57057fcf9f192afa252990379b0769 Mon Sep 17 00:00:00 2001 From: MattJackson <1085847+MattJackson@users.noreply.github.com> Date: Fri, 10 Apr 2026 19:13:53 -0700 Subject: [PATCH] Add IOStream trait and stream-based I/O architecture MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Introduce IOStream trait for uniform read/write across disc, file, network, and null streams. Rename Title→DiscTitle, add stream URL resolver, split old stream.rs into focused modules (m2ts, mkvstream, network, disc, null, resolve, meta). --- src/disc.rs | 25 ++- src/labels/mod.rs | 4 +- src/lib.rs | 9 +- src/mux/codec/ac3.rs | 20 +- src/mux/codec/h264.rs | 3 +- src/mux/codec/hevc.rs | 3 +- src/mux/codec/vc1.rs | 33 +-- src/mux/disc.rs | 130 +++++++++++ src/mux/ebml.rs | 152 ++++++++++++- src/mux/m2ts.rs | 139 ++++++++++++ src/mux/meta.rs | 293 ++++++++++++++++++++++++ src/mux/mkv.rs | 7 + src/mux/mkvstream.rs | 512 ++++++++++++++++++++++++++++++++++++++++++ src/mux/mod.rs | 61 +++-- src/mux/network.rs | 128 +++++++++++ src/mux/null.rs | 43 ++++ src/mux/resolve.rs | 130 +++++++++++ src/mux/stream.rs | 218 ------------------ src/mux/ts.rs | 267 +++++++++++++++++++++- tests/streams.rs | 279 +++++++++++++++++++++++ 20 files changed, 2188 insertions(+), 268 deletions(-) create mode 100644 src/mux/disc.rs create mode 100644 src/mux/m2ts.rs create mode 100644 src/mux/meta.rs create mode 100644 src/mux/mkvstream.rs create mode 100644 src/mux/network.rs create mode 100644 src/mux/null.rs create mode 100644 src/mux/resolve.rs delete mode 100644 src/mux/stream.rs create mode 100644 tests/streams.rs diff --git a/src/disc.rs b/src/disc.rs index ecad532..92a1812 100644 --- a/src/disc.rs +++ b/src/disc.rs @@ -34,7 +34,7 @@ pub struct Disc { /// Number of layers (1 = single, 2 = dual) pub layers: u8, /// Titles sorted by duration (longest first), then playlist name - pub titles: Vec, + pub titles: Vec<DiscTitle>, /// Disc region pub region: DiscRegion, /// AACS state -- None if disc is unencrypted or keys unavailable @@ -80,7 +80,7 @@ pub enum BdRegion { /// A title (one MPLS playlist). #[derive(Debug, Clone)] -pub struct Title { +pub struct DiscTitle { /// Playlist filename (e.g. "00800.mpls") pub playlist: String, /// Playlist number (e.g. 800) @@ -279,7 +279,20 @@ impl ColorSpace { } } -impl Title { +impl DiscTitle { + /// Empty DiscTitle with no streams. + pub fn empty() -> Self { + Self { + playlist: String::new(), + playlist_id: 0, + duration_secs: 0.0, + size_bytes: 0, + clips: Vec::new(), + streams: Vec::new(), + extents: Vec::new(), + } + } + /// Duration formatted as "Xh Ym" pub fn duration_display(&self) -> String { let hrs = (self.duration_secs / 3600.0) as u32; @@ -627,7 +640,7 @@ impl Disc { // ── Internal helpers ──────────────────────────────────────────────────── /// Detect disc format from the main title's video streams. - fn detect_format(titles: &[Title]) -> DiscFormat { + fn detect_format(titles: &[DiscTitle]) -> DiscFormat { for title in titles.iter().take(3) { for stream in &title.streams { if let Stream::Video(v) = stream { @@ -695,7 +708,7 @@ impl Disc { udf_fs: &udf::UdfFs, filename: &str, data: &[u8], - ) -> Option<Title> { + ) -> Option<DiscTitle> { let parsed = mpls::parse(data).ok()?; // Calculate duration from play items @@ -811,7 +824,7 @@ impl Disc { let playlist_num = filename.trim_end_matches(".mpls").trim_end_matches(".MPLS"); let playlist_id = playlist_num.parse::<u16>().unwrap_or(0); - Some(Title { + Some(DiscTitle { playlist: filename.to_string(), playlist_id, duration_secs, diff --git a/src/labels/mod.rs b/src/labels/mod.rs index 8d46551..e59476f 100644 --- a/src/labels/mod.rs +++ b/src/labels/mod.rs @@ -15,7 +15,7 @@ pub mod vocab; use crate::drive::DriveSession; use crate::udf::UdfFs; -use crate::disc::{Title, Stream}; +use crate::disc::{DiscTitle, Stream}; /// A stream label extracted from disc config files. #[derive(Debug, Clone)] @@ -79,7 +79,7 @@ const PARSERS: &[(&str, DetectFn, ParseFn)] = &[ /// Search disc for config files, extract labels, apply to streams. /// This is 100% optional — if anything fails, streams are untouched. -pub fn apply(session: &mut DriveSession, udf: &UdfFs, titles: &mut [Title]) { +pub fn apply(session: &mut DriveSession, udf: &UdfFs, titles: &mut [DiscTitle]) { let labels = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { extract(session, udf) })).unwrap_or_default(); diff --git a/src/lib.rs b/src/lib.rs index a8614dd..cbd84c9 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -92,7 +92,14 @@ pub use profile::DriveProfile; // Platform trait is pub(crate) -- callers use DriveSession, not Platform directly pub use scsi::ScsiTransport; pub use speed::DriveSpeed; -pub use disc::{Disc, DiscFormat, Title, Clip, Stream, VideoStream, AudioStream, SubtitleStream, +pub use disc::{Disc, DiscFormat, DiscTitle, Clip, Stream, VideoStream, AudioStream, SubtitleStream, Codec, HdrFormat, ColorSpace, Extent, ContentReader, AacsState, KeySource, ScanOptions}; +pub use mux::IOStream; pub use mux::MkvStream; +pub use mux::M2tsStream; +pub use mux::NetworkStream; +pub use mux::DiscStream; +pub use mux::NullStream; +pub use mux::DiscOptions; +pub use mux::{open_input, open_output, parse_url, InputOptions}; diff --git a/src/mux/codec/ac3.rs b/src/mux/codec/ac3.rs index 71a2cd6..bbea0cf 100644 --- a/src/mux/codec/ac3.rs +++ b/src/mux/codec/ac3.rs @@ -16,22 +16,34 @@ impl Ac3Parser { impl CodecParser for Ac3Parser { fn parse(&mut self, pes: &PesPacket) -> Vec<Frame> { - if pes.data.is_empty() { + if pes.data.len() < 2 { return Vec::new(); } let pts_ns = pes.pts.map(pts_to_ns).unwrap_or(0); - // AC3: each PES = one frame, always a keyframe + // Find AC3 syncword (0x0B77) — skip any garbage before it + let data = &pes.data; + let start = find_ac3_sync(data).unwrap_or(0); + vec![Frame { pts_ns, keyframe: true, - data: pes.data.clone(), + data: data[start..].to_vec(), }] } fn codec_private(&self) -> Option<Vec<u8>> { - // AC3 doesn't need codecPrivate in MKV None } } + +/// Find AC3 syncword (0x0B77) in data. +fn find_ac3_sync(data: &[u8]) -> Option<usize> { + for i in 0..data.len().saturating_sub(1) { + if data[i] == 0x0B && data[i + 1] == 0x77 { + return Some(i); + } + } + None +} diff --git a/src/mux/codec/h264.rs b/src/mux/codec/h264.rs index 259dcdd..3a916de 100644 --- a/src/mux/codec/h264.rs +++ b/src/mux/codec/h264.rs @@ -29,7 +29,8 @@ impl CodecParser for H264Parser { return Vec::new(); } - let pts_ns = pes.pts.map(pts_to_ns).unwrap_or(0); + // Use DTS when available (monotonic for B-frame content), fall back to PTS + let pts_ns = pes.dts.or(pes.pts).map(pts_to_ns).unwrap_or(0); // Scan NAL units for SPS, PPS, and IDR detection let mut keyframe = false; diff --git a/src/mux/codec/hevc.rs b/src/mux/codec/hevc.rs index 738da9b..971312c 100644 --- a/src/mux/codec/hevc.rs +++ b/src/mux/codec/hevc.rs @@ -34,7 +34,8 @@ impl CodecParser for HevcParser { return Vec::new(); } - let pts_ns = pes.pts.map(pts_to_ns).unwrap_or(0); + // Use DTS when available (monotonic for B-frame content), fall back to PTS + let pts_ns = pes.dts.or(pes.pts).map(pts_to_ns).unwrap_or(0); let data = &pes.data; let mut keyframe = false; diff --git a/src/mux/codec/vc1.rs b/src/mux/codec/vc1.rs index db75e2d..984542c 100644 --- a/src/mux/codec/vc1.rs +++ b/src/mux/codec/vc1.rs @@ -28,8 +28,10 @@ impl CodecParser for Vc1Parser { return Vec::new(); } - let pts_ns = pes.pts.map(pts_to_ns).unwrap_or(0); - let mut keyframe = false; + // Use DTS when available (monotonic for B-frame content), fall back to PTS + let ts_ns = pes.dts.or(pes.pts).map(pts_to_ns).unwrap_or(0); + let mut has_seq_header = false; + let mut frame_start: Option<usize> = None; // Scan for start codes (00 00 01 XX) let data = &pes.data; @@ -39,23 +41,18 @@ impl CodecParser for Vc1Parser { let sc_type = data[i + 3]; match sc_type { SC_SEQUENCE_HEADER => { - // Capture everything from here to next start code let end = find_next_sc(data, i + 4).unwrap_or(data.len()); self.seq_header = Some(data[i..end].to_vec()); + has_seq_header = true; } SC_ENTRY_POINT => { let end = find_next_sc(data, i + 4).unwrap_or(data.len()); self.entry_point = Some(data[i..end].to_vec()); } SC_FRAME => { - // First frame after sequence header + entry point is a keyframe - if self.seq_header.is_some() && self.entry_point.is_some() { - keyframe = true; - } - // Also check frame type from bitstream (bit after start code) - if i + 4 < data.len() { - // For Advanced profile: first 2 bits of frame data indicate type - // But simpler: any frame preceded by seq+entry is I-frame + // Frame data starts at this start code + if frame_start.is_none() { + frame_start = Some(i); } } _ => {} @@ -66,10 +63,20 @@ impl CodecParser for Vc1Parser { } } + // Keyframe = this PES contains a sequence header (I-frame indicator in BD) + let keyframe = has_seq_header; + + // Strip sequence header + entry point from frame data — those are in codecPrivate. + // Only include data from the frame start code onwards. + let frame_data = match frame_start { + Some(start) => &data[start..], + None => data, // no frame start code found, pass through entire PES + }; + vec![Frame { - pts_ns, + pts_ns: ts_ns, keyframe, - data: pes.data.clone(), + data: frame_data.to_vec(), }] } diff --git a/src/mux/disc.rs b/src/mux/disc.rs new file mode 100644 index 0000000..83c32af --- /dev/null +++ b/src/mux/disc.rs @@ -0,0 +1,130 @@ +//! DiscStream — read BD-TS data from an optical disc drive. +//! +//! Read-only stream. Wraps DriveSession + Disc. +//! Handles drive init, AACS decryption, and sector reading. + +use std::io::{self, Read, Write}; +use std::path::Path; +use super::IOStream; +use crate::disc::{DiscTitle, Disc}; +use crate::drive::DriveSession; +use crate::error::Error; + +/// Options for opening a disc stream. +pub struct DiscOptions { + /// Device path (e.g. "/dev/sg4"). None = auto-detect. + pub device: Option<String>, + /// KEYDB.cfg path. None = search standard locations. + pub keydb_path: Option<String>, + /// Which title to read (0-based). None = longest title. + pub title_index: Option<usize>, +} + +impl Default for DiscOptions { + fn default() -> Self { + Self { device: None, keydb_path: None, title_index: None } + } +} + +/// Optical disc stream. Read-only — yields decrypted BD-TS bytes. +pub struct DiscStream { + disc_title: DiscTitle, + disc: Disc, + session: DriveSession, + title_index: usize, + // Read buffer: holds one batch from ContentReader + batch_buf: Vec<u8>, + batch_pos: usize, + started: bool, + eof: bool, +} + +impl DiscStream { + /// Open the disc drive and scan disc metadata. + pub fn open(opts: DiscOptions) -> Result<Self, Error> { + let device = match opts.device { + Some(ref d) => crate::drive::resolve_device(d)?.0, + None => crate::drive::find_drive() + .ok_or_else(|| Error::DeviceNotFound { path: String::new() })?, + }; + + let mut session = DriveSession::open(Path::new(&device))?; + session.wait_ready()?; + let _ = session.init(); + let _ = session.probe_disc(); + + let scan_opts = match opts.keydb_path { + Some(ref kp) => crate::disc::ScanOptions::with_keydb(kp), + None => crate::disc::ScanOptions::default(), + }; + let disc = Disc::scan(&mut session, &scan_opts)?; + + let title_index = opts.title_index.unwrap_or(0); + if title_index >= disc.titles.len() { + return Err(Error::DiscTitleRange { index: title_index, count: disc.titles.len() }); + } + let disc_title = disc.titles[title_index].clone(); + + Ok(Self { + disc_title, disc, session, title_index, + batch_buf: Vec::new(), batch_pos: 0, + started: false, eof: false, + }) + } + + /// Get the full Disc (for listing all titles, etc.) + pub fn disc(&self) -> &Disc { &self.disc } +} + +impl IOStream for DiscStream { + fn info(&self) -> &DiscTitle { &self.disc_title } + fn finish(&mut self) -> io::Result<()> { Ok(()) } +} + +impl Read for DiscStream { + fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> { + // Drain buffer first + if self.batch_pos < self.batch_buf.len() { + let n = (self.batch_buf.len() - self.batch_pos).min(buf.len()); + buf[..n].copy_from_slice(&self.batch_buf[self.batch_pos..self.batch_pos + n]); + self.batch_pos += n; + return Ok(n); + } + + if self.eof { return Ok(0); } + + // Open reader on first call + if !self.started { + self.started = true; + } + + // Read next batch via a temporary ContentReader + // ContentReader borrows session and disc, so we create it inline + let mut reader = self.disc.open_title(&mut self.session, self.title_index) + .map_err(|e| io::Error::new(io::ErrorKind::Other, e.to_string()))?; + + match reader.read_batch() { + Ok(Some(batch)) => { + let n = batch.len().min(buf.len()); + buf[..n].copy_from_slice(&batch[..n]); + if batch.len() > n { + self.batch_buf = batch.to_vec(); + self.batch_pos = n; + } else { + self.batch_buf.clear(); + self.batch_pos = 0; + } + Ok(n) + } + Ok(None) => { self.eof = true; Ok(0) } + Err(e) => Err(io::Error::new(io::ErrorKind::Other, e.to_string())), + } + } +} + +impl Write for DiscStream { + fn write(&mut self, _buf: &[u8]) -> io::Result<usize> { + Err(io::Error::new(io::ErrorKind::Unsupported, "disc is read-only")) + } + fn flush(&mut self) -> io::Result<()> { Ok(()) } +} diff --git a/src/mux/ebml.rs b/src/mux/ebml.rs index 082cdd8..d567918 100644 --- a/src/mux/ebml.rs +++ b/src/mux/ebml.rs @@ -3,16 +3,16 @@ //! EBML uses variable-length integers for element IDs and sizes. //! This module provides low-level writers for constructing MKV files. -use std::io::{self, Write, Seek, SeekFrom}; +use std::io::{self, Read, Write, Seek, SeekFrom}; /// Write an EBML element ID (1-4 bytes, already encoded). /// Element IDs are predefined constants — we write them verbatim. pub fn write_id(w: &mut impl Write, id: u32) -> io::Result<()> { - if id <= 0x7F { + if id <= 0xFF { w.write_all(&[id as u8]) - } else if id <= 0x7FFF { + } else if id <= 0xFFFF { w.write_all(&[(id >> 8) as u8, id as u8]) - } else if id <= 0x7F_FFFF { + } else if id <= 0xFF_FFFF { w.write_all(&[(id >> 16) as u8, (id >> 8) as u8, id as u8]) } else { w.write_all(&[(id >> 24) as u8, (id >> 16) as u8, (id >> 8) as u8, id as u8]) @@ -140,6 +140,149 @@ pub fn end_master<W: Write + Seek>(w: &mut W, size_pos: u64) -> io::Result<()> { Ok(()) } +// ============================================================ +// EBML Read primitives +// ============================================================ + +/// Read an EBML element ID. Returns (id, bytes_consumed). +pub fn read_id(r: &mut impl Read) -> io::Result<(u32, usize)> { + let mut first = [0u8; 1]; + r.read_exact(&mut first)?; + let b0 = first[0]; + + if b0 & 0x80 != 0 { + Ok((b0 as u32, 1)) + } else if b0 & 0x40 != 0 { + let mut b = [0u8; 1]; + r.read_exact(&mut b)?; + Ok((((b0 as u32) << 8) | b[0] as u32, 2)) + } else if b0 & 0x20 != 0 { + let mut b = [0u8; 2]; + r.read_exact(&mut b)?; + Ok((((b0 as u32) << 16) | (b[0] as u32) << 8 | b[1] as u32, 3)) + } else if b0 & 0x10 != 0 { + let mut b = [0u8; 3]; + r.read_exact(&mut b)?; + Ok((((b0 as u32) << 24) | (b[0] as u32) << 16 | (b[1] as u32) << 8 | b[2] as u32, 4)) + } else { + Err(io::Error::new(io::ErrorKind::InvalidData, "invalid EBML ID")) + } +} + +/// Read an EBML variable-length size. Returns (size, bytes_consumed). +/// Size of u64::MAX means "unknown size". +pub fn read_size(r: &mut impl Read) -> io::Result<(u64, usize)> { + let mut first = [0u8; 1]; + r.read_exact(&mut first)?; + let b0 = first[0]; + + if b0 & 0x80 != 0 { + let val = (b0 & 0x7F) as u64; + if val == 0x7F { return Ok((u64::MAX, 1)); } // unknown + Ok((val, 1)) + } else if b0 & 0x40 != 0 { + let mut b = [0u8; 1]; + r.read_exact(&mut b)?; + let val = (((b0 & 0x3F) as u64) << 8) | b[0] as u64; + if val == 0x3FFF { return Ok((u64::MAX, 2)); } + Ok((val, 2)) + } else if b0 & 0x20 != 0 { + let mut b = [0u8; 2]; + r.read_exact(&mut b)?; + let val = (((b0 & 0x1F) as u64) << 16) | (b[0] as u64) << 8 | b[1] as u64; + if val == 0x1FFFFF { return Ok((u64::MAX, 3)); } + Ok((val, 3)) + } else if b0 & 0x10 != 0 { + let mut b = [0u8; 3]; + r.read_exact(&mut b)?; + let val = (((b0 & 0x0F) as u64) << 24) | (b[0] as u64) << 16 | (b[1] as u64) << 8 | b[2] as u64; + if val == 0x0FFFFFFF { return Ok((u64::MAX, 4)); } + Ok((val, 4)) + } else if b0 & 0x08 != 0 { + let mut b = [0u8; 4]; + r.read_exact(&mut b)?; + let val = (((b0 & 0x07) as u64) << 32) | (b[0] as u64) << 24 | (b[1] as u64) << 16 | (b[2] as u64) << 8 | b[3] as u64; + Ok((val, 5)) + } else if b0 & 0x04 != 0 { + let mut b = [0u8; 5]; + r.read_exact(&mut b)?; + let val = (((b0 & 0x03) as u64) << 40) | (b[0] as u64) << 32 | (b[1] as u64) << 24 | (b[2] as u64) << 16 | (b[3] as u64) << 8 | b[4] as u64; + Ok((val, 6)) + } else if b0 & 0x02 != 0 { + let mut b = [0u8; 6]; + r.read_exact(&mut b)?; + let val = (((b0 & 0x01) as u64) << 48) | (b[0] as u64) << 40 | (b[1] as u64) << 32 | (b[2] as u64) << 24 | (b[3] as u64) << 16 | (b[4] as u64) << 8 | b[5] as u64; + Ok((val, 7)) + } else { + let mut b = [0u8; 7]; + r.read_exact(&mut b)?; + let val = (b[0] as u64) << 48 | (b[1] as u64) << 40 | (b[2] as u64) << 32 | (b[3] as u64) << 24 | (b[4] as u64) << 16 | (b[5] as u64) << 8 | b[6] as u64; + if val == 0x00FFFFFFFFFFFFFF { return Ok((u64::MAX, 8)); } + Ok((val, 8)) + } +} + +/// Read an EBML element header (ID + size). Returns (id, data_size, header_bytes). +pub fn read_element_header(r: &mut impl Read) -> io::Result<(u32, u64, usize)> { + let (id, id_len) = read_id(r)?; + let (size, size_len) = read_size(r)?; + Ok((id, size, id_len + size_len)) +} + +/// Read an unsigned integer value of `len` bytes. +pub fn read_uint_val(r: &mut impl Read, len: usize) -> io::Result<u64> { + let mut buf = [0u8; 8]; + r.read_exact(&mut buf[..len])?; + let mut val = 0u64; + for &b in &buf[..len] { + val = (val << 8) | b as u64; + } + Ok(val) +} + +/// Read a float value (4 or 8 bytes). +pub fn read_float_val(r: &mut impl Read, len: usize) -> io::Result<f64> { + if len == 4 { + let mut buf = [0u8; 4]; + r.read_exact(&mut buf)?; + Ok(f32::from_be_bytes(buf) as f64) + } else { + let mut buf = [0u8; 8]; + r.read_exact(&mut buf)?; + Ok(f64::from_be_bytes(buf)) + } +} + +/// Read a UTF-8 string value of `len` bytes. +pub fn read_string_val(r: &mut impl Read, len: usize) -> io::Result<String> { + let mut buf = vec![0u8; len]; + r.read_exact(&mut buf)?; + // Strip trailing nulls + while buf.last() == Some(&0) { buf.pop(); } + String::from_utf8(buf).map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e)) +} + +/// Read binary data of `len` bytes. +pub fn read_binary_val(r: &mut impl Read, len: usize) -> io::Result<Vec<u8>> { + let mut buf = vec![0u8; len]; + r.read_exact(&mut buf)?; + Ok(buf) +} + +/// Read a VINT (track number) from a SimpleBlock. Returns (value, bytes_consumed). +pub fn read_vint(r: &mut impl Read) -> io::Result<(u64, usize)> { + let mut first = [0u8; 1]; + r.read_exact(&mut first)?; + let b0 = first[0]; + if b0 & 0x80 != 0 { return Ok(((b0 & 0x7F) as u64, 1)); } + if b0 & 0x40 != 0 { + let mut b = [0u8; 1]; + r.read_exact(&mut b)?; + return Ok(((((b0 & 0x3F) as u64) << 8) | b[0] as u64, 2)); + } + Err(io::Error::new(io::ErrorKind::InvalidData, "unsupported VINT width")) +} + // ============================================================ // Matroska Element IDs // ============================================================ @@ -183,6 +326,7 @@ pub const FLAG_FORCED: u32 = 0x55AA; pub const LANGUAGE: u32 = 0x22B59C; pub const CODEC_ID: u32 = 0x86; pub const CODEC_PRIVATE: u32 = 0x63A2; +pub const TRACK_NAME: u32 = 0x536E; pub const DEFAULT_DURATION: u32 = 0x23E383; // Video diff --git a/src/mux/m2ts.rs b/src/mux/m2ts.rs new file mode 100644 index 0000000..f5d7050 --- /dev/null +++ b/src/mux/m2ts.rs @@ -0,0 +1,139 @@ +//! M2tsStream — BD transport stream with embedded metadata header. +//! +//! Write: prepends FMKV metadata header, then passes through BD-TS bytes. +//! Read: extracts metadata header (or scans PMT), then yields BD-TS bytes. + +use std::io::{self, Read, Write, Seek, SeekFrom}; +use super::{IOStream, ReadSeek, meta, ts}; +use crate::disc::{DiscTitle, Stream as DiscStream}; + +/// Size of initial scan buffer for PMT/stream detection. +const SCAN_SIZE: usize = 1024 * 1024; + +enum Mode { + Write { + writer: Box<dyn Write>, + header_written: bool, + }, + Read { + reader: Box<dyn ReadSeek>, + }, +} + +/// BD transport stream with embedded metadata. +pub struct M2tsStream { + disc_title: DiscTitle, + mode: Mode, + finished: bool, +} + +impl M2tsStream { + /// Create for writing. Metadata header is written on first write(). + pub fn new(writer: impl Write + 'static) -> Self { + Self { + disc_title: DiscTitle::empty(), + mode: Mode::Write { + writer: Box::new(writer), + header_written: false, + }, + finished: false, + } + } + + /// Set stream metadata. Returns self for chaining. + pub fn meta(mut self, dt: &DiscTitle) -> Self { + self.disc_title = dt.clone(); + self + } + + /// Open an m2ts file for reading. + /// + /// Tries FMKV metadata header first. Falls back to PMT scan + PTS duration. + pub fn open(mut reader: impl Read + Seek + 'static) -> io::Result<Self> { + // Try FMKV metadata header + if let Ok(Some(m)) = meta::read_header(&mut reader) { + return Ok(Self { + disc_title: m.to_title(), + mode: Mode::Read { reader: Box::new(reader) }, + finished: false, + }); + } + + // Fallback: scan PMT for streams, PTS for duration + reader.seek(SeekFrom::Start(0))?; + let mut buf = vec![0u8; SCAN_SIZE]; + let n = reader.read(&mut buf)?; + + let streams = ts::scan_streams(&buf[..n]) + .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidData, "no streams found"))?; + + let video_pid = streams.iter().find_map(|s| match s { + DiscStream::Video(v) => Some(v.pid), + _ => None, + }); + let duration = video_pid + .and_then(|pid| ts::scan_duration(&mut reader, pid)) + .unwrap_or(0.0); + + reader.seek(SeekFrom::Start(0))?; + + Ok(Self { + disc_title: DiscTitle { + duration_secs: duration, + streams, + ..DiscTitle::empty() + }, + mode: Mode::Read { reader: Box::new(reader) }, + finished: false, + }) + } +} + +impl IOStream for M2tsStream { + fn info(&self) -> &DiscTitle { &self.disc_title } + + fn finish(&mut self) -> io::Result<()> { + if self.finished { return Ok(()); } + self.finished = true; + if let Mode::Write { ref mut writer, .. } = self.mode { + writer.flush() + } else { + Ok(()) + } + } +} + +impl Write for M2tsStream { + fn write(&mut self, buf: &[u8]) -> io::Result<usize> { + match self.mode { + Mode::Write { ref mut writer, ref mut header_written } => { + if !*header_written { + if !self.disc_title.streams.is_empty() { + let m = meta::M2tsMeta::from_title(&self.disc_title); + meta::write_header(&mut *writer, &m)?; + } + *header_written = true; + } + writer.write(buf) + } + Mode::Read { .. } => Err(io::Error::new(io::ErrorKind::Unsupported, "stream opened for reading")), + } + } + + fn flush(&mut self) -> io::Result<()> { + if let Mode::Write { ref mut writer, .. } = self.mode { + writer.flush() + } else { + Ok(()) + } + } +} + +impl Read for M2tsStream { + fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> { + match self.mode { + Mode::Read { ref mut reader } => reader.read(buf), + Mode::Write { .. } => Err(io::Error::new(io::ErrorKind::Unsupported, "stream opened for writing")), + } + } +} diff --git a/src/mux/meta.rs b/src/mux/meta.rs new file mode 100644 index 0000000..48983c8 --- /dev/null +++ b/src/mux/meta.rs @@ -0,0 +1,293 @@ +//! M2TS metadata header — embeds title/stream info in raw m2ts files. +//! +//! Format: [8B magic] [4B json_len] [JSON] [padding to 192B boundary] [BD-TS data...] +//! Other tools skip the header during TS sync recovery (scan for 0x47). + +use std::io::{self, Read, Seek, SeekFrom, Write}; +use serde::{Serialize, Deserialize}; +use crate::disc::{DiscTitle, Stream, VideoStream, AudioStream, SubtitleStream, + Codec, HdrFormat, ColorSpace}; + +/// Magic bytes: "FMKV" + version 1 + 2 reserved bytes. +const MAGIC: [u8; 8] = [b'F', b'M', b'K', b'V', 0x00, 0x01, 0x00, 0x00]; + +/// BD-TS packet size (header must be padded to this boundary). +const PACKET_SIZE: usize = 192; + +/// Metadata embedded in an m2ts file. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct M2tsMeta { + /// Format version. + pub v: u8, + /// Title name (e.g. filename stem or disc title). + #[serde(default)] + pub title: String, + /// Duration in seconds. + #[serde(default)] + pub duration: f64, + /// Stream descriptors. + pub streams: Vec<MetaStream>, +} + +/// A single stream descriptor in the metadata. +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(tag = "type")] +pub enum MetaStream { + #[serde(rename = "video")] + Video { + pid: u16, + codec: String, + #[serde(default)] resolution: String, + #[serde(default)] frame_rate: String, + #[serde(default)] hdr: String, + #[serde(default)] label: String, + #[serde(default)] secondary: bool, + }, + #[serde(rename = "audio")] + Audio { + pid: u16, + codec: String, + #[serde(default)] channels: String, + #[serde(default)] language: String, + #[serde(default)] sample_rate: String, + #[serde(default)] label: String, + #[serde(default)] secondary: bool, + }, + #[serde(rename = "subtitle")] + Subtitle { + pid: u16, + codec: String, + #[serde(default)] language: String, + #[serde(default)] forced: bool, + }, +} + +impl M2tsMeta { + /// Build metadata from a disc Title. + pub fn from_title(title: &DiscTitle) -> Self { + let streams = title.streams.iter().map(|s| match s { + Stream::Video(v) => MetaStream::Video { + pid: v.pid, + codec: codec_to_str(v.codec), + resolution: v.resolution.clone(), + frame_rate: v.frame_rate.clone(), + hdr: hdr_to_str(v.hdr), + label: v.label.clone(), + secondary: v.secondary, + }, + Stream::Audio(a) => MetaStream::Audio { + pid: a.pid, + codec: codec_to_str(a.codec), + channels: a.channels.clone(), + language: a.language.clone(), + sample_rate: a.sample_rate.clone(), + label: a.label.clone(), + secondary: a.secondary, + }, + Stream::Subtitle(s) => MetaStream::Subtitle { + pid: s.pid, + codec: codec_to_str(s.codec), + language: s.language.clone(), + forced: s.forced, + }, + }).collect(); + + Self { + v: 1, + title: title.playlist.clone(), + duration: title.duration_secs, + streams, + } + } + + /// Convert back to a library Title (for remux). + pub fn to_title(&self) -> DiscTitle { + let streams = self.streams.iter().map(|s| match s { + MetaStream::Video { pid, codec, resolution, frame_rate, hdr, label, secondary } => { + Stream::Video(VideoStream { + pid: *pid, + codec: str_to_codec(codec), + resolution: resolution.clone(), + frame_rate: frame_rate.clone(), + hdr: str_to_hdr(hdr), + color_space: ColorSpace::Bt709, + secondary: *secondary, + label: label.clone(), + }) + } + MetaStream::Audio { pid, codec, channels, language, sample_rate, label, secondary } => { + Stream::Audio(AudioStream { + pid: *pid, + codec: str_to_codec(codec), + channels: channels.clone(), + language: language.clone(), + sample_rate: sample_rate.clone(), + secondary: *secondary, + label: label.clone(), + }) + } + MetaStream::Subtitle { pid, codec, language, forced } => { + Stream::Subtitle(SubtitleStream { + pid: *pid, + codec: str_to_codec(codec), + language: language.clone(), + forced: *forced, + }) + } + }).collect(); + + DiscTitle { + playlist: self.title.clone(), + playlist_id: 0, + duration_secs: self.duration, + size_bytes: 0, + clips: Vec::new(), + streams, + extents: Vec::new(), + } + } +} + +/// Write the metadata header to a writer. Padded to 192-byte boundary. +pub fn write_header(w: &mut impl Write, meta: &M2tsMeta) -> io::Result<()> { + let json = serde_json::to_vec(meta) + .map_err(|e| io::Error::new(io::ErrorKind::Other, e))?; + + let json_len = json.len() as u32; + let raw_len = 8 + 4 + json.len(); // magic + len + json + let padded_len = ((raw_len + PACKET_SIZE - 1) / PACKET_SIZE) * PACKET_SIZE; + let padding = padded_len - raw_len; + + w.write_all(&MAGIC)?; + w.write_all(&json_len.to_be_bytes())?; + w.write_all(&json)?; + if padding > 0 { + w.write_all(&vec![0u8; padding])?; + } + Ok(()) +} + +/// Try to read a metadata header from the start of an m2ts file. +/// Returns None for bare m2ts files (no header). +/// On success, leaves reader positioned at the first TS packet. +/// On failure, seeks back to the start. +pub fn read_header<R: Read + Seek>(r: &mut R) -> io::Result<Option<M2tsMeta>> { + let start = r.stream_position()?; + + let mut magic = [0u8; 8]; + if r.read_exact(&mut magic).is_err() { + r.seek(SeekFrom::Start(start))?; + return Ok(None); + } + + if magic[..4] != MAGIC[..4] { + // Not a freemkv m2ts — seek back + r.seek(SeekFrom::Start(start))?; + return Ok(None); + } + + let mut len_buf = [0u8; 4]; + r.read_exact(&mut len_buf)?; + let json_len = u32::from_be_bytes(len_buf) as usize; + + let mut json_buf = vec![0u8; json_len]; + r.read_exact(&mut json_buf)?; + + let meta: M2tsMeta = serde_json::from_slice(&json_buf) + .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; + + // Skip padding to next 192-byte boundary + let raw_len = 8 + 4 + json_len; + let padded_len = ((raw_len + PACKET_SIZE - 1) / PACKET_SIZE) * PACKET_SIZE; + let padding = padded_len - raw_len; + if padding > 0 { + r.seek(SeekFrom::Current(padding as i64))?; + } + + Ok(Some(meta)) +} + +/// Read a metadata header from a forward-only stream (no Seek required). +/// Returns None if the magic bytes don't match. Consumes the header bytes. +pub fn read_header_from_stream(r: &mut impl Read) -> io::Result<Option<M2tsMeta>> { + let mut magic = [0u8; 8]; + r.read_exact(&mut magic)?; + + if magic[..4] != MAGIC[..4] { + return Ok(None); + } + + let mut len_buf = [0u8; 4]; + r.read_exact(&mut len_buf)?; + let json_len = u32::from_be_bytes(len_buf) as usize; + + let mut json_buf = vec![0u8; json_len]; + r.read_exact(&mut json_buf)?; + + let meta: M2tsMeta = serde_json::from_slice(&json_buf) + .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; + + // Skip padding + let raw_len = 8 + 4 + json_len; + let padded_len = ((raw_len + PACKET_SIZE - 1) / PACKET_SIZE) * PACKET_SIZE; + let padding = padded_len - raw_len; + if padding > 0 { + let mut skip = vec![0u8; padding]; + r.read_exact(&mut skip)?; + } + + Ok(Some(meta)) +} + +// Codec string conversion (compact, no English — just codec identifiers) +fn codec_to_str(c: Codec) -> String { + match c { + Codec::Hevc => "hevc", + Codec::H264 => "h264", + Codec::Vc1 => "vc1", + Codec::Mpeg2 => "mpeg2", + Codec::TrueHd => "truehd", + Codec::DtsHdMa => "dtshd_ma", + Codec::DtsHdHr => "dtshd_hr", + Codec::Dts => "dts", + Codec::Ac3 => "ac3", + Codec::Ac3Plus => "eac3", + Codec::Lpcm => "lpcm", + Codec::Pgs => "pgs", + Codec::Unknown(_) => "unknown", + }.into() +} + +fn str_to_codec(s: &str) -> Codec { + match s { + "hevc" => Codec::Hevc, + "h264" => Codec::H264, + "vc1" => Codec::Vc1, + "mpeg2" => Codec::Mpeg2, + "truehd" => Codec::TrueHd, + "dtshd_ma" => Codec::DtsHdMa, + "dtshd_hr" => Codec::DtsHdHr, + "dts" => Codec::Dts, + "ac3" => Codec::Ac3, + "eac3" => Codec::Ac3Plus, + "lpcm" => Codec::Lpcm, + "pgs" => Codec::Pgs, + _ => Codec::Unknown(0), + } +} + +fn hdr_to_str(h: HdrFormat) -> String { + match h { + HdrFormat::Sdr => "sdr", + HdrFormat::Hdr10 => "hdr10", + HdrFormat::DolbyVision => "dv", + }.into() +} + +fn str_to_hdr(s: &str) -> HdrFormat { + match s { + "hdr10" => HdrFormat::Hdr10, + "dv" => HdrFormat::DolbyVision, + _ => HdrFormat::Sdr, + } +} diff --git a/src/mux/mkv.rs b/src/mux/mkv.rs index ad50627..b4fc801 100644 --- a/src/mux/mkv.rs +++ b/src/mux/mkv.rs @@ -13,6 +13,7 @@ pub struct MkvTrack { pub track_type: u64, // 1=video, 2=audio, 17=subtitle pub codec_id: &'static str, pub language: String, + pub name: String, // Track name / label (e.g. "English (Lossless)") pub codec_private: Option<Vec<u8>>, pub is_default: bool, pub is_forced: bool, @@ -39,6 +40,7 @@ impl MkvTrack { track_type: ebml::TRACK_TYPE_VIDEO, codec_id, language: "und".into(), + name: v.label.clone(), codec_private: None, // filled later by parser is_default: !v.secondary, is_forced: false, @@ -65,6 +67,7 @@ impl MkvTrack { track_type: ebml::TRACK_TYPE_AUDIO, codec_id, language: a.language.clone(), + name: a.label.clone(), codec_private: None, is_default: !a.secondary, is_forced: false, @@ -81,6 +84,7 @@ impl MkvTrack { track_type: ebml::TRACK_TYPE_SUBTITLE, codec_id: "S_HDMV/PGS", language: s.language.clone(), + name: String::new(), codec_private: None, is_default: false, is_forced: s.forced, @@ -163,6 +167,9 @@ impl<W: Write + Seek> MkvMuxer<W> { ebml::write_uint(&mut writer, ebml::FLAG_LACING, 0)?; ebml::write_string(&mut writer, ebml::CODEC_ID, track.codec_id)?; ebml::write_string(&mut writer, ebml::LANGUAGE, &track.language)?; + if !track.name.is_empty() { + ebml::write_string(&mut writer, ebml::TRACK_NAME, &track.name)?; + } if !track.is_default { ebml::write_uint(&mut writer, ebml::FLAG_DEFAULT, 0)?; diff --git a/src/mux/mkvstream.rs b/src/mux/mkvstream.rs new file mode 100644 index 0000000..e037671 --- /dev/null +++ b/src/mux/mkvstream.rs @@ -0,0 +1,512 @@ +//! MkvStream — Matroska container stream. +//! +//! Write: BD-TS bytes in → demux → codec parse → MKV container out. +//! Read: MKV container in → extract frames → wrap as BD-TS → bytes out. + +use std::io::{self, Read, Write, Seek, SeekFrom}; +use super::{IOStream, WriteSeek, ReadSeek, ebml}; +use super::ts::TsDemuxer; +use super::mkv::{MkvMuxer, MkvTrack}; +use super::codec::{self, CodecParser}; +use super::lookahead::{LookaheadBuffer, LookaheadState, DEFAULT_LOOKAHEAD_SIZE}; +use crate::disc::*; + +/// Lookahead buffer for codec header detection (5 MB default). +const DEFAULT_MAX_BUFFER: usize = DEFAULT_LOOKAHEAD_SIZE; + +#[derive(Debug, Clone, Copy, PartialEq)] +enum WritePhase { Scanning, Streaming } + +struct WriteState { + demuxer: TsDemuxer, + muxer: Option<MkvMuxer<Box<dyn WriteSeek>>>, + writer: Option<Box<dyn WriteSeek>>, + parsers: Vec<(u16, Box<dyn CodecParser>)>, + pid_to_track: Vec<(u16, usize)>, + tracks: Vec<MkvTrack>, + lookahead: LookaheadBuffer, + phase: WritePhase, + video_pending: usize, +} + +struct ReadState { + reader: Box<dyn ReadSeek>, + buf: Vec<u8>, + pos: usize, + len: usize, + cluster_ts_ms: i64, +} + +enum Mode { + Write(WriteState), + Read(ReadState), +} + +/// Matroska container stream. +pub struct MkvStream { + disc_title: DiscTitle, + mode: Mode, + max_buffer: usize, + finished: bool, +} + +impl MkvStream { + /// Create for writing. BD-TS bytes written to this stream produce MKV output. + pub fn new(writer: impl Write + Seek + 'static) -> Self { + Self { + disc_title: DiscTitle::empty(), + mode: Mode::Write(WriteState { + demuxer: TsDemuxer::new(&[]), + muxer: None, + writer: Some(Box::new(writer)), + parsers: Vec::new(), + pid_to_track: Vec::new(), + tracks: Vec::new(), + lookahead: LookaheadBuffer::new(DEFAULT_MAX_BUFFER), + phase: WritePhase::Scanning, + video_pending: 0, + }), + max_buffer: DEFAULT_MAX_BUFFER, + finished: false, + } + } + + /// Set stream metadata. Returns self for chaining. + pub fn meta(mut self, dt: &DiscTitle) -> Self { + if let Mode::Write(ref mut ws) = self.mode { + let mut pids = Vec::new(); + for s in &dt.streams { + let (pid, track, parser) = match s { + crate::disc::Stream::Video(v) => { + ws.video_pending += 1; + (v.pid, MkvTrack::video(v), codec::parser_for_codec(v.codec)) + } + crate::disc::Stream::Audio(a) => { + (a.pid, MkvTrack::audio(a), codec::parser_for_codec(a.codec)) + } + crate::disc::Stream::Subtitle(s) => { + (s.pid, MkvTrack::subtitle(s), codec::parser_for_codec(s.codec)) + } + }; + let idx = ws.tracks.len(); + pids.push(pid); + ws.pid_to_track.push((pid, idx)); + ws.parsers.push((pid, parser)); + ws.tracks.push(track); + } + ws.demuxer = TsDemuxer::new(&pids); + } + self.disc_title = dt.clone(); + self + } + + /// Set lookahead buffer size. Returns self. + pub fn max_buffer(mut self, size: usize) -> Self { + self.max_buffer = size; + if let Mode::Write(ref mut ws) = self.mode { + ws.lookahead = LookaheadBuffer::new(size); + } + self + } + + /// Open an MKV file for reading. + pub fn open(mut reader: impl Read + Seek + 'static) -> io::Result<Self> { + let disc_title = parse_mkv_header(&mut reader)?; + Ok(Self { + disc_title, + mode: Mode::Read(ReadState { + reader: Box::new(reader), + buf: Vec::new(), + pos: 0, + len: 0, + cluster_ts_ms: 0, + }), + max_buffer: 0, + finished: false, + }) + } +} + +impl IOStream for MkvStream { + fn info(&self) -> &DiscTitle { &self.disc_title } + + fn finish(&mut self) -> io::Result<()> { + if self.finished { return Ok(()); } + self.finished = true; + if let Mode::Write(ref mut ws) = self.mode { + // Flush remaining PES packets + if let Some(ref mut muxer) = ws.muxer { + for pes in &ws.demuxer.flush() { + write_pes(&ws.pid_to_track, &mut ws.parsers, muxer, pes)?; + } + } + // Write cues and finalize + if let Some(muxer) = ws.muxer.take() { + muxer.finish()?; + } + } + Ok(()) + } +} + +// ── Write ────────────────────────────────────────────────────── + +impl Write for MkvStream { + fn write(&mut self, buf: &[u8]) -> io::Result<usize> { + let dt = &self.disc_title; + let ws = match self.mode { + Mode::Write(ref mut ws) => ws, + Mode::Read(_) => return Err(io::Error::new(io::ErrorKind::Unsupported, "stream opened for reading")), + }; + + match ws.phase { + WritePhase::Scanning => { + // Feed demuxer for codec detection + let packets = ws.demuxer.feed(buf); + for pes in &packets { + if let Some((_, p)) = ws.parsers.iter_mut().find(|(pid, _)| *pid == pes.pid) { + let _ = p.parse(pes); + } + } + + let state = ws.lookahead.push(buf); + + // Check if all video codec headers found + if check_codec_private(ws) { + ws.lookahead.mark_ready(); + begin_streaming(ws, dt)?; + return Ok(buf.len()); + } + + match state { + LookaheadState::Collecting | LookaheadState::Ready => Ok(buf.len()), + LookaheadState::Overflow => Err(io::Error::new( + io::ErrorKind::OutOfMemory, + "no codec headers found within lookahead buffer", + )), + } + } + WritePhase::Streaming => { + let packets = ws.demuxer.feed(buf); + if let Some(ref mut muxer) = ws.muxer { + for pes in &packets { + write_pes(&ws.pid_to_track, &mut ws.parsers, muxer, pes)?; + } + } + Ok(buf.len()) + } + } + } + + fn flush(&mut self) -> io::Result<()> { Ok(()) } +} + +// ── Read ─────────────────────────────────────────────────────── + +impl Read for MkvStream { + fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> { + let rs = match self.mode { + Mode::Read(ref mut rs) => rs, + Mode::Write(_) => return Err(io::Error::new(io::ErrorKind::Unsupported, "stream opened for writing")), + }; + + // Drain internal buffer first + if rs.pos < rs.len { + let n = (rs.len - rs.pos).min(buf.len()); + buf[..n].copy_from_slice(&rs.buf[rs.pos..rs.pos + n]); + rs.pos += n; + return Ok(n); + } + + // Read next element from MKV + loop { + let (id, size, _) = match ebml::read_element_header(&mut rs.reader) { + Ok(h) => h, + Err(_) => return Ok(0), + }; + + match id { + ebml::CLUSTER => continue, + ebml::CLUSTER_TIMESTAMP => { + rs.cluster_ts_ms = ebml::read_uint_val(&mut rs.reader, size as usize)? as i64; + continue; + } + ebml::SIMPLE_BLOCK => { + let block = ebml::read_binary_val(&mut rs.reader, size as usize)?; + if block.len() < 4 { continue; } + + let (track, vl) = block_vint(&block); + if vl + 3 > block.len() { continue; } + + let rel_ts = i16::from_be_bytes([block[vl], block[vl + 1]]); + let frame = &block[vl + 3..]; + let pts_ms = rs.cluster_ts_ms + rel_ts as i64; + + rs.buf.clear(); + frame_to_ts(&mut rs.buf, track as u16, pts_ms, frame); + rs.pos = 0; + rs.len = rs.buf.len(); + + if rs.len > 0 { + let n = rs.len.min(buf.len()); + buf[..n].copy_from_slice(&rs.buf[..n]); + rs.pos = n; + return Ok(n); + } + } + _ => { + if size != u64::MAX && size > 0 { + rs.reader.seek(SeekFrom::Current(size as i64))?; + } + } + } + } + } +} + +// ── Write internals ──────────────────────────────────────────── + +fn check_codec_private(ws: &mut WriteState) -> bool { + if ws.video_pending == 0 { return true; } + for (pid, parser) in &ws.parsers { + if let Some(cp) = parser.codec_private() { + if let Some((_, idx)) = ws.pid_to_track.iter().find(|(p, _)| p == pid) { + if ws.tracks[*idx].codec_private.is_none() { + ws.tracks[*idx].codec_private = Some(cp); + ws.video_pending -= 1; + } + } + } + } + ws.video_pending == 0 +} + +fn begin_streaming(ws: &mut WriteState, dt: &DiscTitle) -> io::Result<()> { + let writer = ws.writer.take() + .ok_or_else(|| io::Error::new(io::ErrorKind::Other, "writer already consumed"))?; + + ws.muxer = Some(MkvMuxer::new(writer, &ws.tracks, Some(&dt.playlist), dt.duration_secs)?); + ws.phase = WritePhase::Streaming; + + // Re-parse buffered data through a fresh demuxer + let buffered = ws.lookahead.drain(); + if !buffered.is_empty() { + let pids: Vec<u16> = ws.pid_to_track.iter().map(|(pid, _)| *pid).collect(); + let mut temp = TsDemuxer::new(&pids); + let packets = temp.feed(&buffered); + if let Some(ref mut muxer) = ws.muxer { + for pes in &packets { + write_pes(&ws.pid_to_track, &mut ws.parsers, muxer, pes)?; + } + } + } + Ok(()) +} + +fn write_pes( + pid_to_track: &[(u16, usize)], + parsers: &mut [(u16, Box<dyn CodecParser>)], + muxer: &mut MkvMuxer<Box<dyn WriteSeek>>, + pes: &super::ts::PesPacket, +) -> io::Result<()> { + let idx = match pid_to_track.iter().find(|(pid, _)| *pid == pes.pid) { + Some((_, idx)) => *idx, + None => return Ok(()), + }; + let parser = match parsers.iter_mut().find(|(pid, _)| *pid == pes.pid) { + Some((_, p)) => p, + None => return Ok(()), + }; + for frame in parser.parse(pes) { + muxer.write_frame(idx, frame.pts_ns, frame.keyframe, &frame.data)?; + } + Ok(()) +} + +// ── MKV header parsing (read side) ──────────────────────────── + +fn parse_mkv_header(r: &mut (impl Read + Seek)) -> io::Result<DiscTitle> { + let mut title = String::new(); + let mut duration_ms = 0.0f64; + let mut ts_scale: u64 = 1_000_000; + let mut streams: Vec<crate::disc::Stream> = Vec::new(); + + let (id, size, _) = ebml::read_element_header(r)?; + if id != ebml::EBML { return Err(io::Error::new(io::ErrorKind::InvalidData, "not EBML")); } + r.seek(SeekFrom::Current(size as i64))?; + + let (id, _, _) = ebml::read_element_header(r)?; + if id != ebml::SEGMENT { return Err(io::Error::new(io::ErrorKind::InvalidData, "no Segment")); } + + let (mut got_info, mut got_tracks) = (false, false); + + loop { + if got_info && got_tracks { break; } + let (id, size, _) = match ebml::read_element_header(r) { Ok(h) => h, Err(_) => break }; + + match id { + ebml::INFO => { + let end = r.stream_position()? + size; + while r.stream_position()? < end { + let (cid, cs, _) = ebml::read_element_header(r)?; + match cid { + ebml::TIMESTAMP_SCALE => ts_scale = ebml::read_uint_val(r, cs as usize)?, + ebml::DURATION => duration_ms = ebml::read_float_val(r, cs as usize)?, + ebml::TITLE => title = ebml::read_string_val(r, cs as usize)?, + _ => { r.seek(SeekFrom::Current(cs as i64))?; } + } + } + got_info = true; + } + ebml::TRACKS => { + let end = r.stream_position()? + size; + while r.stream_position()? < end { + let (cid, cs, _) = ebml::read_element_header(r)?; + if cid == ebml::TRACK_ENTRY { + if let Some(s) = parse_track(r, cs)? { streams.push(s); } + } else { r.seek(SeekFrom::Current(cs as i64))?; } + } + got_tracks = true; + } + ebml::CLUSTER => break, + _ if size != u64::MAX => { r.seek(SeekFrom::Current(size as i64))?; } + _ => break, + } + } + + Ok(DiscTitle { + playlist: title, + duration_secs: duration_ms * (ts_scale as f64) / 1_000_000_000.0, + streams, + ..DiscTitle::empty() + }) +} + +fn parse_track(r: &mut (impl Read + Seek), size: u64) -> io::Result<Option<crate::disc::Stream>> { + let end = r.stream_position()? + size; + let (mut ttype, mut tnum) = (0u64, 0u16); + let (mut codec_id, mut lang, mut name) = (String::new(), String::from("und"), String::new()); + let (mut ph, mut sr, mut ch, mut forced) = (0u32, 0.0f64, 0u8, false); + + while r.stream_position()? < end { + let (cid, cs, _) = ebml::read_element_header(r)?; + match cid { + ebml::TRACK_NUMBER => tnum = ebml::read_uint_val(r, cs as usize)? as u16, + ebml::TRACK_TYPE => ttype = ebml::read_uint_val(r, cs as usize)?, + ebml::CODEC_ID => codec_id = ebml::read_string_val(r, cs as usize)?, + ebml::LANGUAGE => lang = ebml::read_string_val(r, cs as usize)?, + ebml::TRACK_NAME => name = ebml::read_string_val(r, cs as usize)?, + ebml::FLAG_FORCED => forced = ebml::read_uint_val(r, cs as usize)? != 0, + ebml::VIDEO => { + let ve = r.stream_position()? + cs; + while r.stream_position()? < ve { + let (vid, vs, _) = ebml::read_element_header(r)?; + if vid == ebml::PIXEL_HEIGHT { ph = ebml::read_uint_val(r, vs as usize)? as u32; } + else { r.seek(SeekFrom::Current(vs as i64))?; } + } + } + ebml::AUDIO => { + let ae = r.stream_position()? + cs; + while r.stream_position()? < ae { + let (aid, as_, _) = ebml::read_element_header(r)?; + match aid { + ebml::SAMPLING_FREQUENCY => sr = ebml::read_float_val(r, as_ as usize)?, + ebml::CHANNELS => ch = ebml::read_uint_val(r, as_ as usize)? as u8, + _ => { r.seek(SeekFrom::Current(as_ as i64))?; } + } + } + } + _ => { r.seek(SeekFrom::Current(cs as i64))?; } + } + } + + let codec = match codec_id.as_str() { + "V_MPEGH/ISO/HEVC" => Codec::Hevc, "V_MPEG4/ISO/AVC" => Codec::H264, + "V_MS/VFW/FOURCC" => Codec::Vc1, "V_MPEG2" => Codec::Mpeg2, + "A_AC3" => Codec::Ac3, "A_EAC3" => Codec::Ac3Plus, + "A_TRUEHD" => Codec::TrueHd, "A_DTS" => Codec::Dts, + "A_PCM/INT/BIG" => Codec::Lpcm, "S_HDMV/PGS" => Codec::Pgs, + _ => Codec::Unknown(0), + }; + let res = format!("{}p", ph); + let chs: String = match ch { 8 => "7.1", 6 => "5.1", 2 => "stereo", 1 => "mono", _ => "5.1" }.into(); + let srs: String = if sr >= 96000.0 { "96kHz" } else { "48kHz" }.into(); + + Ok(match ttype { + 1 => Some(crate::disc::Stream::Video(VideoStream { + pid: tnum, codec, resolution: res, frame_rate: String::new(), + hdr: HdrFormat::Sdr, color_space: ColorSpace::Bt709, secondary: false, label: name, + })), + 2 => Some(crate::disc::Stream::Audio(AudioStream { + pid: tnum, codec, channels: chs, language: lang, sample_rate: srs, + secondary: false, label: name, + })), + 17 => Some(crate::disc::Stream::Subtitle(SubtitleStream { + pid: tnum, codec, language: lang, forced, + })), + _ => None, + }) +} + +// ── BD-TS frame wrapping (read side) ────────────────────────── + +fn block_vint(d: &[u8]) -> (u64, usize) { + if d.is_empty() { return (0, 0); } + if d[0] & 0x80 != 0 { return ((d[0] & 0x7F) as u64, 1); } + if d[0] & 0x40 != 0 && d.len() >= 2 { + return ((((d[0] & 0x3F) as u64) << 8) | d[1] as u64, 2); + } + (0, 1) +} + +fn frame_to_ts(out: &mut Vec<u8>, track: u16, pts_ms: i64, data: &[u8]) { + let pid = if track == 1 { 0x1011 } else { 0x1100 + (track - 2) as u16 }; + let stream_id: u8 = if track == 1 { 0xE0 } else { 0xBD }; + let pts = encode_pts(pts_ms * 90); + let hdr = [0x00, 0x00, 0x01, stream_id, 0x00, 0x00, 0x80, 0x80, 0x05]; + + let mut pes = Vec::with_capacity(hdr.len() + pts.len() + data.len()); + pes.extend_from_slice(&hdr); + pes.extend_from_slice(&pts); + pes.extend_from_slice(data); + + let mut off = 0; + let mut pusi = true; + while off < pes.len() { + let mut pkt = [0u8; 192]; + pkt[4] = 0x47; + pkt[5] = (pid >> 8) as u8 & 0x1F; + if pusi { pkt[5] |= 0x40; pusi = false; } + pkt[6] = pid as u8; + + let space = 184; + let rem = pes.len() - off; + let n = rem.min(space); + + if n < space { + let pad = space - n; + pkt[7] = 0x30; // AF + payload + pkt[8] = pad as u8; + if pad > 1 { pkt[9] = 0x00; } + for i in 10..(8 + pad).min(192) { pkt[i] = 0xFF; } + pkt[8 + pad..8 + pad + n].copy_from_slice(&pes[off..off + n]); + } else { + pkt[7] = 0x10; // payload only + pkt[8..8 + n].copy_from_slice(&pes[off..off + n]); + } + + out.extend_from_slice(&pkt); + off += n; + } +} + +fn encode_pts(pts: i64) -> [u8; 5] { + let p = pts as u64; + [ + 0x21 | ((p >> 29) & 0x0E) as u8, + ((p >> 22) & 0xFF) as u8, + 0x01 | ((p >> 14) & 0xFE) as u8, + ((p >> 7) & 0xFF) as u8, + 0x01 | ((p << 1) & 0xFE) as u8, + ] +} diff --git a/src/mux/mod.rs b/src/mux/mod.rs index b8b603e..fc3bd90 100644 --- a/src/mux/mod.rs +++ b/src/mux/mod.rs @@ -1,25 +1,60 @@ -//! MKV muxing pipeline. +//! Stream-based I/O pipeline. //! -//! Provides BD transport stream → MKV remuxing via composable streams. -//! -//! The main type is `MkvStream` — wraps any `Write + Seek` output, -//! receives raw BD-TS bytes via `write()`, outputs MKV. +//! All formats are streams. Two URLs, left reads, right writes: //! //! ```text -//! disc.rip(title, MkvStream::new(file, &title)) +//! freemkv disc:// mkv://Dune.mkv +//! freemkv m2ts://Dune.m2ts mkv://Dune.mkv +//! freemkv disc:// network://10.1.7.11:9000 //! ``` //! -//! Components (for advanced use): -//! - `ts`: BD transport stream demuxer (192-byte packets → PES frames) -//! - `ebml`: EBML write primitives for Matroska container -//! - `mkv`: MKV muxer (tracks, clusters, blocks, cues) -//! - `codec`: Elementary stream parsers (frame boundaries, codec headers) +//! Streams implement `IOStream` for uniform handling: +//! +//! ```text +//! let mut input = open_input("disc://", &opts)?; +//! let mut output = open_output("mkv://Dune.mkv", input.info())?; +//! io::copy(&mut *input, &mut *output)?; +//! output.finish()?; +//! ``` pub mod ebml; pub mod ts; pub mod mkv; pub mod codec; pub mod lookahead; -pub mod stream; +pub mod meta; +mod m2ts; +mod mkvstream; +pub mod network; +pub mod disc; +pub mod null; +pub mod resolve; -pub use stream::MkvStream; +pub use m2ts::M2tsStream; +pub use mkvstream::MkvStream; +pub use network::NetworkStream; +pub use disc::{DiscStream, DiscOptions}; +pub use null::NullStream; +pub use resolve::{open_input, open_output, parse_url, InputOptions, StreamUrl}; + +use std::io::{self, Read, Write, Seek}; +use crate::disc::DiscTitle; + +/// Common interface for all stream types. +/// +/// A stream can be opened for reading or created for writing. +/// Calling the unsupported direction returns an error. +pub trait IOStream: Read + Write { + /// Get stream metadata. + fn info(&self) -> &DiscTitle; + + /// Finalize the stream (flush, write index/cues, close). + fn finish(&mut self) -> io::Result<()>; +} + +// Combined traits for internal trait objects. +pub(crate) trait ReadSeek: Read + Seek {} +impl<T: Read + Seek> ReadSeek for T {} + +pub(crate) trait WriteSeek: Write + Seek {} +impl<T: Write + Seek> WriteSeek for T {} diff --git a/src/mux/network.rs b/src/mux/network.rs new file mode 100644 index 0000000..946aefa --- /dev/null +++ b/src/mux/network.rs @@ -0,0 +1,128 @@ +//! NetworkStream — BD-TS over TCP with embedded metadata. +//! +//! Write side (sender): connects to a listener, sends FMKV header + BD-TS data. +//! Read side (receiver): listens for a connection, reads FMKV header + BD-TS data. +//! +//! The FMKV metadata header is the same format as M2tsStream uses — so a +//! NetworkStream reader can hand off to any output stream (MKV, M2TS, etc.) +//! with full metadata (labels, languages, duration). + +use std::io::{self, Read, Write, BufReader, BufWriter}; +use std::net::{TcpListener, TcpStream}; +use super::{IOStream, meta}; +use crate::disc::DiscTitle; + +/// I/O buffer size for network reads/writes. +const NET_BUF_SIZE: usize = 256 * 1024; + +enum Mode { + Write { + writer: BufWriter<TcpStream>, + header_written: bool, + }, + Read { + reader: BufReader<TcpStream>, + }, +} + +/// TCP network stream for distributed rip/remux. +pub struct NetworkStream { + disc_title: DiscTitle, + mode: Mode, + finished: bool, +} + +impl NetworkStream { + /// Connect to a remote listener for writing. + /// Sends FMKV metadata header on first write. + pub fn connect(addr: &str) -> io::Result<Self> { + let stream = TcpStream::connect(addr)?; + stream.set_nodelay(true)?; + Ok(Self { + disc_title: DiscTitle::empty(), + mode: Mode::Write { + writer: BufWriter::with_capacity(NET_BUF_SIZE, stream), + header_written: false, + }, + finished: false, + }) + } + + /// Set stream metadata (for write side). Returns self for chaining. + pub fn meta(mut self, dt: &DiscTitle) -> Self { + self.disc_title = dt.clone(); + self + } + + /// Listen for an incoming connection and read from it. + /// Extracts FMKV metadata header from the sender. + pub fn listen(addr: &str) -> io::Result<Self> { + let listener = TcpListener::bind(addr)?; + let (stream, _peer) = listener.accept()?; + stream.set_nodelay(true)?; + let mut reader = BufReader::with_capacity(NET_BUF_SIZE, stream); + + // Read FMKV metadata header (inline, since TcpStream doesn't impl Seek) + let disc_title = meta::read_header_from_stream(&mut reader)? + .ok_or_else(|| io::Error::new( + io::ErrorKind::InvalidData, + "no FMKV metadata header from sender", + ))? + .to_title(); + + Ok(Self { + disc_title, + mode: Mode::Read { reader }, + finished: false, + }) + } +} + +impl IOStream for NetworkStream { + fn info(&self) -> &DiscTitle { &self.disc_title } + + fn finish(&mut self) -> io::Result<()> { + if self.finished { return Ok(()); } + self.finished = true; + if let Mode::Write { ref mut writer, .. } = self.mode { + writer.flush()?; + writer.get_ref().shutdown(std::net::Shutdown::Write)?; + } + Ok(()) + } +} + +impl Write for NetworkStream { + fn write(&mut self, buf: &[u8]) -> io::Result<usize> { + match self.mode { + Mode::Write { ref mut writer, ref mut header_written } => { + if !*header_written { + if !self.disc_title.streams.is_empty() { + let m = meta::M2tsMeta::from_title(&self.disc_title); + meta::write_header(&mut *writer, &m)?; + } + *header_written = true; + } + writer.write(buf) + } + Mode::Read { .. } => Err(io::Error::new(io::ErrorKind::Unsupported, "stream opened for reading")), + } + } + + fn flush(&mut self) -> io::Result<()> { + if let Mode::Write { ref mut writer, .. } = self.mode { + writer.flush() + } else { + Ok(()) + } + } +} + +impl Read for NetworkStream { + fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> { + match self.mode { + Mode::Read { ref mut reader } => reader.read(buf), + Mode::Write { .. } => Err(io::Error::new(io::ErrorKind::Unsupported, "stream opened for writing")), + } + } +} diff --git a/src/mux/null.rs b/src/mux/null.rs new file mode 100644 index 0000000..7623f52 --- /dev/null +++ b/src/mux/null.rs @@ -0,0 +1,43 @@ +//! NullStream — discards all data. Write-only. For benchmarking. + +use std::io::{self, Read, Write}; +use super::IOStream; +use crate::disc::DiscTitle; + +/// Null stream — accepts writes, discards data. For benchmarking rip speed. +pub struct NullStream { + disc_title: DiscTitle, + bytes_written: u64, +} + +impl NullStream { + pub fn new() -> Self { + Self { disc_title: DiscTitle::empty(), bytes_written: 0 } + } + + pub fn meta(mut self, dt: &DiscTitle) -> Self { + self.disc_title = dt.clone(); + self + } + + pub fn bytes_written(&self) -> u64 { self.bytes_written } +} + +impl IOStream for NullStream { + fn info(&self) -> &DiscTitle { &self.disc_title } + fn finish(&mut self) -> io::Result<()> { Ok(()) } +} + +impl Write for NullStream { + fn write(&mut self, buf: &[u8]) -> io::Result<usize> { + self.bytes_written += buf.len() as u64; + Ok(buf.len()) + } + fn flush(&mut self) -> io::Result<()> { Ok(()) } +} + +impl Read for NullStream { + fn read(&mut self, _buf: &mut [u8]) -> io::Result<usize> { + Err(io::Error::new(io::ErrorKind::Unsupported, "null stream is write-only")) + } +} diff --git a/src/mux/resolve.rs b/src/mux/resolve.rs new file mode 100644 index 0000000..5655ae8 --- /dev/null +++ b/src/mux/resolve.rs @@ -0,0 +1,130 @@ +//! Stream URL resolver — parses URL strings into IOStream instances. +//! +//! Schemes: disc://, m2ts://, mkv://, network:// +//! Bare paths infer scheme from file extension. + +use std::io::{self, BufReader, BufWriter}; +use super::{IOStream, M2tsStream, MkvStream}; +use super::network::NetworkStream; +use super::null::NullStream; +use super::disc::{DiscStream, DiscOptions}; +use crate::disc::DiscTitle; + +/// I/O buffer size for file streams. +const IO_BUF_SIZE: usize = 4 * 1024 * 1024; + +/// MKV lookahead buffer size. +const MKV_LOOKAHEAD: usize = 10 * 1024 * 1024; + +/// Parsed stream URL. +pub struct StreamUrl { + pub scheme: String, + pub path: String, +} + +/// Parse a URL string into scheme + path. +/// +/// Supports: `disc://`, `disc:///dev/sg4`, `m2ts://path`, `mkv://path`, +/// `network://host:port`, or bare paths (infer from extension). +pub fn parse_url(url: &str) -> StreamUrl { + if let Some(rest) = url.strip_prefix("disc://") { + return StreamUrl { scheme: "disc".into(), path: rest.to_string() }; + } + if let Some(rest) = url.strip_prefix("m2ts://") { + return StreamUrl { scheme: "m2ts".into(), path: rest.to_string() }; + } + if let Some(rest) = url.strip_prefix("mkv://") { + return StreamUrl { scheme: "mkv".into(), path: rest.to_string() }; + } + if let Some(rest) = url.strip_prefix("network://") { + return StreamUrl { scheme: "network".into(), path: rest.to_string() }; + } + if url == "null://" || url.starts_with("null://") { + return StreamUrl { scheme: "null".into(), path: String::new() }; + } + + // Infer from extension + let scheme = if url.ends_with(".m2ts") { + "m2ts" + } else if url.ends_with(".mkv") { + "mkv" + } else { + "m2ts" // default + }; + + StreamUrl { scheme: scheme.into(), path: url.to_string() } +} + +/// Open a stream URL for reading (source). +pub fn open_input(url: &str, opts: &InputOptions) -> io::Result<Box<dyn IOStream>> { + let parsed = parse_url(url); + + match parsed.scheme.as_str() { + "disc" => { + let disc_opts = DiscOptions { + device: if parsed.path.is_empty() { None } else { Some(parsed.path) }, + keydb_path: opts.keydb_path.clone(), + title_index: opts.title_index, + }; + let stream = DiscStream::open(disc_opts) + .map_err(|e| io::Error::new(io::ErrorKind::Other, e.to_string()))?; + Ok(Box::new(stream)) + } + "m2ts" => { + let file = std::fs::File::open(&parsed.path)?; + let reader = BufReader::with_capacity(IO_BUF_SIZE, file); + Ok(Box::new(M2tsStream::open(reader)?)) + } + "mkv" => { + let file = std::fs::File::open(&parsed.path)?; + let reader = BufReader::with_capacity(IO_BUF_SIZE, file); + Ok(Box::new(MkvStream::open(reader)?)) + } + "network" => { + Ok(Box::new(NetworkStream::listen(&parsed.path)?)) + } + _ => Err(io::Error::new(io::ErrorKind::InvalidInput, + format!("unknown scheme: {}", parsed.scheme))), + } +} + +/// Open a stream URL for writing (destination). +pub fn open_output(url: &str, meta: &DiscTitle) -> io::Result<Box<dyn IOStream>> { + let parsed = parse_url(url); + + match parsed.scheme.as_str() { + "disc" => { + Err(io::Error::new(io::ErrorKind::Unsupported, "disc is read-only")) + } + "null" => { + Ok(Box::new(NullStream::new().meta(meta))) + } + "m2ts" => { + let file = std::fs::File::create(&parsed.path)?; + let writer = BufWriter::with_capacity(IO_BUF_SIZE, file); + Ok(Box::new(M2tsStream::new(writer).meta(meta))) + } + "mkv" => { + let file = std::fs::File::create(&parsed.path)?; + let writer = BufWriter::with_capacity(IO_BUF_SIZE, file); + Ok(Box::new(MkvStream::new(writer).meta(meta).max_buffer(MKV_LOOKAHEAD))) + } + "network" => { + Ok(Box::new(NetworkStream::connect(&parsed.path)?.meta(meta))) + } + _ => Err(io::Error::new(io::ErrorKind::InvalidInput, + format!("unknown scheme: {}", parsed.scheme))), + } +} + +/// Options for opening an input stream. +pub struct InputOptions { + pub keydb_path: Option<String>, + pub title_index: Option<usize>, +} + +impl Default for InputOptions { + fn default() -> Self { + Self { keydb_path: None, title_index: None } + } +} diff --git a/src/mux/stream.rs b/src/mux/stream.rs deleted file mode 100644 index dfe3d64..0000000 --- a/src/mux/stream.rs +++ /dev/null @@ -1,218 +0,0 @@ -//! MkvStream — a Write adapter that demuxes BD-TS and writes MKV. -//! -//! ```rust,ignore -//! let output = MkvStream::new(file) -//! .title(&disc.titles[0]) -//! .max_buffer(20 * 1024 * 1024); -//! -//! disc.rip(0, output)?; -//! ``` - -use std::io::{self, Write, Seek}; -use super::ts::{TsDemuxer, PesPacket}; -use super::mkv::{MkvMuxer, MkvTrack}; -use super::codec::{self, CodecParser}; -use super::lookahead::{LookaheadBuffer, LookaheadState, DEFAULT_LOOKAHEAD_SIZE}; -use crate::disc::{Stream, Title}; - -/// Phase of the MkvStream. -#[derive(Debug, Clone, Copy, PartialEq)] -enum Phase { - /// Collecting data in lookahead buffer, scanning for codec setup. - Scanning, - /// Header written, streaming directly to muxer. - Streaming, -} - -/// MKV output stream. Implements `Write`. -pub struct MkvStream<W: Write + Seek> { - demuxer: TsDemuxer, - muxer: Option<MkvMuxer<W>>, - writer: Option<W>, - parsers: Vec<(u16, Box<dyn CodecParser>)>, - pid_to_track: Vec<(u16, usize)>, - tracks: Vec<MkvTrack>, - title_name: String, - duration_secs: f64, - lookahead: LookaheadBuffer, - phase: Phase, - video_tracks_pending: usize, -} - -impl<W: Write + Seek> MkvStream<W> { - /// Create a new MkvStream wrapping an output writer. - pub fn new(writer: W) -> Self { - Self { - demuxer: TsDemuxer::new(&[]), - muxer: None, - writer: Some(writer), - parsers: Vec::new(), - pid_to_track: Vec::new(), - tracks: Vec::new(), - title_name: String::new(), - duration_secs: 0.0, - lookahead: LookaheadBuffer::new(DEFAULT_LOOKAHEAD_SIZE), - phase: Phase::Scanning, - video_tracks_pending: 0, - } - } - - /// Set the title metadata (streams, duration, name). Returns self. - pub fn title(mut self, title: &Title) -> Self { - let mut pids = Vec::new(); - - for stream in &title.streams { - let (pid, track, parser) = match stream { - Stream::Video(v) => { - self.video_tracks_pending += 1; - (v.pid, MkvTrack::video(v), codec::parser_for_codec(v.codec)) - } - Stream::Audio(a) => (a.pid, MkvTrack::audio(a), codec::parser_for_codec(a.codec)), - Stream::Subtitle(s) => (s.pid, MkvTrack::subtitle(s), codec::parser_for_codec(s.codec)), - }; - let track_idx = self.tracks.len(); - pids.push(pid); - self.pid_to_track.push((pid, track_idx)); - self.parsers.push((pid, parser)); - self.tracks.push(track); - } - - self.demuxer = TsDemuxer::new(&pids); - self.title_name = title.playlist.clone(); - self.duration_secs = title.duration_secs; - self - } - - /// Set the lookahead buffer size in bytes. Default 5 MB. Returns self. - pub fn max_buffer(mut self, size: usize) -> Self { - self.lookahead = LookaheadBuffer::new(size); - self - } - - /// Finalize the MKV file — close cluster, write cues. - pub fn finish(mut self) -> io::Result<()> { - if let Some(ref mut muxer) = self.muxer { - let remaining = self.demuxer.flush(); - for pes in &remaining { - Self::process_one_pes(&self.pid_to_track, &mut self.parsers, muxer, pes)?; - } - } - if let Some(muxer) = self.muxer { - muxer.finish()?; - } - Ok(()) - } - - fn check_codec_private(&mut self) -> bool { - if self.video_tracks_pending == 0 { - return true; - } - for (pid, parser) in &self.parsers { - if let Some(cp) = parser.codec_private() { - if let Some((_, track_idx)) = self.pid_to_track.iter().find(|(p, _)| p == pid) { - if self.tracks[*track_idx].codec_private.is_none() { - self.tracks[*track_idx].codec_private = Some(cp); - self.video_tracks_pending -= 1; - } - } - } - } - self.video_tracks_pending == 0 - } - - fn start_streaming(&mut self) -> io::Result<()> { - let writer = self.writer.take().ok_or_else(|| { - io::Error::new(io::ErrorKind::Other, "writer already consumed") - })?; - - let muxer = MkvMuxer::new(writer, &self.tracks, Some(&self.title_name), self.duration_secs)?; - self.muxer = Some(muxer); - self.phase = Phase::Streaming; - - // Re-parse and write buffered data - let buffered = self.lookahead.drain(); - if !buffered.is_empty() { - let pids: Vec<u16> = self.pid_to_track.iter().map(|(pid, _)| *pid).collect(); - let mut temp_demuxer = TsDemuxer::new(&pids); - let mut packets = temp_demuxer.feed(&buffered); - packets.extend(temp_demuxer.flush()); - - if let Some(ref mut muxer) = self.muxer { - for pes in &packets { - Self::process_one_pes(&self.pid_to_track, &mut self.parsers, muxer, pes)?; - } - } - } - - Ok(()) - } - - fn process_one_pes( - pid_to_track: &[(u16, usize)], - parsers: &mut [(u16, Box<dyn CodecParser>)], - muxer: &mut MkvMuxer<W>, - pes: &PesPacket, - ) -> io::Result<()> { - let track_idx = match pid_to_track.iter().find(|(pid, _)| *pid == pes.pid) { - Some((_, idx)) => *idx, - None => return Ok(()), - }; - let parser = match parsers.iter_mut().find(|(pid, _)| *pid == pes.pid) { - Some((_, p)) => p, - None => return Ok(()), - }; - let frames = parser.parse(pes); - for frame in frames { - muxer.write_frame(track_idx, frame.pts_ns, frame.keyframe, &frame.data)?; - } - Ok(()) - } -} - -impl<W: Write + Seek> Write for MkvStream<W> { - fn write(&mut self, buf: &[u8]) -> io::Result<usize> { - match self.phase { - Phase::Scanning => { - // Parse for codec info - let packets = self.demuxer.feed(buf); - for pes in &packets { - if let Some((_, parser)) = self.parsers.iter_mut().find(|(p, _)| *p == pes.pid) { - let _ = parser.parse(pes); - } - } - - // Try to buffer - let state = self.lookahead.push(buf); - - // Check if we have everything - if self.check_codec_private() { - self.lookahead.mark_ready(); - self.start_streaming()?; - return Ok(buf.len()); - } - - match state { - LookaheadState::Collecting => Ok(buf.len()), - LookaheadState::Overflow => Err(io::Error::new( - io::ErrorKind::OutOfMemory, - "MKV lookahead buffer overflow — no codec data found within buffer limit", - )), - LookaheadState::Ready => Ok(buf.len()), - } - } - Phase::Streaming => { - let packets = self.demuxer.feed(buf); - if let Some(ref mut muxer) = self.muxer { - for pes in &packets { - Self::process_one_pes(&self.pid_to_track, &mut self.parsers, muxer, pes)?; - } - } - Ok(buf.len()) - } - } - } - - fn flush(&mut self) -> io::Result<()> { - Ok(()) - } -} diff --git a/src/mux/ts.rs b/src/mux/ts.rs index 526a83e..3a7c583 100644 --- a/src/mux/ts.rs +++ b/src/mux/ts.rs @@ -15,6 +15,12 @@ const TS_PACKET_SIZE: usize = 188; /// TS sync byte. const SYNC_BYTE: u8 = 0x47; +/// Size of head buffer for PTS/stream scanning (1 MB). +const SCAN_HEAD_SIZE: usize = 1024 * 1024; + +/// Size of tail buffer for last-PTS scanning (2 MB). +const SCAN_TAIL_SIZE: usize = 2 * 1024 * 1024; + /// A reassembled PES packet with timestamp info. #[derive(Debug)] pub struct PesPacket { @@ -94,6 +100,7 @@ impl PesAssembler { pub struct TsDemuxer { assemblers: Vec<PesAssembler>, pid_index: [i16; 8192], // PID → index into assemblers, -1 = not tracked + remainder: Vec<u8>, // leftover bytes from previous feed() call } impl TsDemuxer { @@ -105,17 +112,31 @@ impl TsDemuxer { pid_index[pid as usize] = i as i16; assemblers.push(PesAssembler::new(pid)); } - Self { assemblers, pid_index } + Self { assemblers, pid_index, remainder: Vec::new() } } - /// Feed a chunk of BD transport stream data (must be aligned to 192-byte packets). - /// Returns completed PES packets. + /// Feed a chunk of BD transport stream data. Handles non-192-byte-aligned input + /// by buffering leftover bytes between calls. Returns completed PES packets. pub fn feed(&mut self, data: &[u8]) -> Vec<PesPacket> { let mut completed = Vec::new(); + + // Prepend any remainder from previous call + let work: &[u8]; + let mut combined: Vec<u8> = Vec::new(); + if !self.remainder.is_empty() { + combined.reserve(self.remainder.len() + data.len()); + combined.extend_from_slice(&self.remainder); + combined.extend_from_slice(data); + self.remainder.clear(); + work = &combined; + } else { + work = data; + } + let mut offset = 0; - while offset + BD_TS_PACKET_SIZE <= data.len() { - let packet = &data[offset..offset + BD_TS_PACKET_SIZE]; + while offset + BD_TS_PACKET_SIZE <= work.len() { + let packet = &work[offset..offset + BD_TS_PACKET_SIZE]; offset += BD_TS_PACKET_SIZE; // Skip 4-byte TP_extra_header, check sync byte @@ -172,6 +193,11 @@ impl TsDemuxer { } } + // Save leftover bytes for next call + if offset < work.len() { + self.remainder.extend_from_slice(&work[offset..]); + } + completed } @@ -242,6 +268,237 @@ fn parse_timestamp(data: &[u8]) -> i64 { | b4 >> 1 } +// ============================================================ +// Stream scanning (PAT/PMT → stream list) +// ============================================================ + +/// Scan BD-TS data for streams by parsing PAT and PMT tables. +/// Returns None if no valid program is found. +pub fn scan_streams(data: &[u8]) -> Option<Vec<crate::disc::Stream>> { + use crate::disc::*; + + // Pass 1: find PMT PID from PAT + let mut pat_pmt_pid: Option<u16> = None; + let mut offset = 0; + while offset + BD_TS_PACKET_SIZE <= data.len() { + if data[offset + 4] != SYNC_BYTE { offset += 1; continue; } + let pid = (((data[offset + 5] & 0x1F) as u16) << 8) | data[offset + 6] as u16; + let pusi = data[offset + 5] & 0x40 != 0; + + if pid == 0 && pusi { + let payload_start = offset + 4 + 4; + if payload_start + 12 < data.len() { + let pointer = data[payload_start] as usize; + let pat_start = payload_start + 1 + pointer; + if pat_start + 12 < data.len() && data[pat_start] == 0x00 { + let section_len = (((data[pat_start + 1] & 0x0F) as usize) << 8) | data[pat_start + 2] as usize; + let entries_start = pat_start + 8; + let entries_end = pat_start + 3 + section_len - 4; + let mut e = entries_start; + while e + 4 <= data.len() && e < entries_end { + let prog_num = ((data[e] as u16) << 8) | data[e + 1] as u16; + let p = (((data[e + 2] & 0x1F) as u16) << 8) | data[e + 3] as u16; + if prog_num != 0 { + pat_pmt_pid = Some(p); + break; + } + e += 4; + } + } + } + } + offset += BD_TS_PACKET_SIZE; + } + + let pmt_pid = pat_pmt_pid?; + + // Pass 2: parse PMT for stream entries + let mut streams = Vec::new(); + offset = 0; + while offset + BD_TS_PACKET_SIZE <= data.len() { + if data[offset + 4] != SYNC_BYTE { offset += 1; continue; } + let pid = (((data[offset + 5] & 0x1F) as u16) << 8) | data[offset + 6] as u16; + let pusi = data[offset + 5] & 0x40 != 0; + + if pid == pmt_pid && pusi { + let payload_start = offset + 4 + 4; + if payload_start + 1 >= data.len() { offset += BD_TS_PACKET_SIZE; continue; } + let pointer = data[payload_start] as usize; + let pmt_start = payload_start + 1 + pointer; + if pmt_start + 12 >= data.len() { offset += BD_TS_PACKET_SIZE; continue; } + if data[pmt_start] != 0x02 { offset += BD_TS_PACKET_SIZE; continue; } + + let section_len = (((data[pmt_start + 1] & 0x0F) as usize) << 8) | data[pmt_start + 2] as usize; + let prog_info_len = (((data[pmt_start + 10] & 0x0F) as usize) << 8) | data[pmt_start + 11] as usize; + let mut pos = pmt_start + 12 + prog_info_len; + let end = pmt_start + 3 + section_len - 4; + + while pos + 5 <= data.len() && pos < end { + let stream_type = data[pos]; + let es_pid = (((data[pos + 1] & 0x1F) as u16) << 8) | data[pos + 2] as u16; + let es_info_len = (((data[pos + 3] & 0x0F) as usize) << 8) | data[pos + 4] as usize; + + let stream = match stream_type { + 0x1B => Some(Stream::Video(VideoStream { + pid: es_pid, codec: Codec::H264, + resolution: "1080p".into(), frame_rate: String::new(), + hdr: HdrFormat::Sdr, color_space: ColorSpace::Bt709, + secondary: false, label: String::new(), + })), + 0x24 => Some(Stream::Video(VideoStream { + pid: es_pid, codec: Codec::Hevc, + resolution: "2160p".into(), frame_rate: String::new(), + hdr: HdrFormat::Sdr, color_space: ColorSpace::Bt709, + secondary: false, label: String::new(), + })), + 0xEA => Some(Stream::Video(VideoStream { + pid: es_pid, codec: Codec::Vc1, + resolution: "1080p".into(), frame_rate: String::new(), + hdr: HdrFormat::Sdr, color_space: ColorSpace::Bt709, + secondary: false, label: String::new(), + })), + 0x02 => Some(Stream::Video(VideoStream { + pid: es_pid, codec: Codec::Mpeg2, + resolution: "1080i".into(), frame_rate: String::new(), + hdr: HdrFormat::Sdr, color_space: ColorSpace::Bt709, + secondary: false, label: String::new(), + })), + 0x81 => Some(Stream::Audio(AudioStream { + pid: es_pid, codec: Codec::Ac3, + channels: "5.1".into(), language: "und".into(), + sample_rate: "48kHz".into(), secondary: false, label: String::new(), + })), + 0x83 => Some(Stream::Audio(AudioStream { + pid: es_pid, codec: Codec::TrueHd, + channels: "5.1".into(), language: "und".into(), + sample_rate: "48kHz".into(), secondary: false, label: String::new(), + })), + 0x84 | 0xA1 => Some(Stream::Audio(AudioStream { + pid: es_pid, codec: Codec::Ac3Plus, + channels: "5.1".into(), language: "und".into(), + sample_rate: "48kHz".into(), secondary: false, label: String::new(), + })), + 0x85 | 0x86 => Some(Stream::Audio(AudioStream { + pid: es_pid, codec: Codec::DtsHdMa, + channels: "5.1".into(), language: "und".into(), + sample_rate: "48kHz".into(), secondary: false, label: String::new(), + })), + 0x82 => Some(Stream::Audio(AudioStream { + pid: es_pid, codec: Codec::Dts, + channels: "5.1".into(), language: "und".into(), + sample_rate: "48kHz".into(), secondary: false, label: String::new(), + })), + 0x90 => Some(Stream::Subtitle(SubtitleStream { + pid: es_pid, codec: Codec::Pgs, + language: "und".into(), forced: false, + })), + _ => None, + }; + + if let Some(s) = stream { + streams.push(s); + } + pos += 5 + es_info_len; + } + break; + } + offset += BD_TS_PACKET_SIZE; + } + + if streams.is_empty() { None } else { Some(streams) } +} + +// ============================================================ +// PTS scanning utilities (for duration detection) +// ============================================================ + +/// Find the first PTS for a given PID in BD-TS data. +pub fn scan_first_pts(data: &[u8], target_pid: u16) -> Option<i64> { + let mut offset = 0; + while offset + BD_TS_PACKET_SIZE <= data.len() { + if data[offset + 4] != SYNC_BYTE { offset += 1; continue; } + let pid = (((data[offset + 5] & 0x1F) as u16) << 8) | data[offset + 6] as u16; + let pusi = data[offset + 5] & 0x40 != 0; + if pid == target_pid && pusi { + let ts = &data[offset + 4..offset + BD_TS_PACKET_SIZE]; + let afc = (ts[3] >> 4) & 0x03; + let payload_start = if afc == 3 { 5 + ts[4] as usize } else { 4 }; + if payload_start < TS_PACKET_SIZE { + let payload = &ts[payload_start..]; + if payload.len() >= 14 && payload[0] == 0 && payload[1] == 0 && payload[2] == 1 { + let pts_dts_flags = (payload[7] >> 6) & 0x03; + if pts_dts_flags >= 2 { + return Some(parse_timestamp(&payload[9..14])); + } + } + } + } + offset += BD_TS_PACKET_SIZE; + } + None +} + +/// Find the last PTS for a given PID in BD-TS data. +pub fn scan_last_pts(data: &[u8], target_pid: u16) -> Option<i64> { + let mut last_pts = None; + let mut offset = 0; + while offset + BD_TS_PACKET_SIZE <= data.len() { + if data[offset + 4] != SYNC_BYTE { offset += 1; continue; } + let pid = (((data[offset + 5] & 0x1F) as u16) << 8) | data[offset + 6] as u16; + let pusi = data[offset + 5] & 0x40 != 0; + if pid == target_pid && pusi { + let ts = &data[offset + 4..offset + BD_TS_PACKET_SIZE]; + let afc = (ts[3] >> 4) & 0x03; + let payload_start = if afc == 3 { 5 + ts[4] as usize } else { 4 }; + if payload_start < TS_PACKET_SIZE { + let payload = &ts[payload_start..]; + if payload.len() >= 14 && payload[0] == 0 && payload[1] == 0 && payload[2] == 1 { + let pts_dts_flags = (payload[7] >> 6) & 0x03; + if pts_dts_flags >= 2 { + last_pts = Some(parse_timestamp(&payload[9..14])); + } + } + } + } + offset += BD_TS_PACKET_SIZE; + } + last_pts +} + +/// Scan an m2ts file for duration by reading first and last PTS. +/// Returns duration in seconds, or None if PTS cannot be found. +/// The reader position is restored after scanning. +pub fn scan_duration<R: std::io::Read + std::io::Seek>(r: &mut R, video_pid: u16) -> Option<f64> { + use std::io::SeekFrom; + + let start_pos = r.stream_position().ok()?; + + // Read first 1MB for first PTS + let mut head_buf = vec![0u8; SCAN_HEAD_SIZE]; + r.seek(SeekFrom::Start(0)).ok()?; + let head_n = r.read(&mut head_buf).ok()?; + let first_pts = scan_first_pts(&head_buf[..head_n], video_pid)?; + + // Read last 2MB for last PTS (aligned to 192-byte boundary) + let file_size = r.seek(SeekFrom::End(0)).ok()?; + let tail_size: u64 = SCAN_TAIL_SIZE as u64; + let raw_pos = if file_size > tail_size { file_size - tail_size } else { 0 }; + let seek_pos = (raw_pos / BD_TS_PACKET_SIZE as u64) * BD_TS_PACKET_SIZE as u64; + r.seek(SeekFrom::Start(seek_pos)).ok()?; + let mut tail_buf = vec![0u8; tail_size as usize]; + let tail_n = r.read(&mut tail_buf).ok()?; + let last_pts = scan_last_pts(&tail_buf[..tail_n], video_pid)?; + + // Restore reader position + let _ = r.seek(SeekFrom::Start(start_pos)); + + if last_pts > first_pts { + Some((last_pts - first_pts) as f64 / 90000.0) + } else { + None + } +} + #[cfg(test)] mod tests { use super::*; diff --git a/tests/streams.rs b/tests/streams.rs new file mode 100644 index 0000000..c185b03 --- /dev/null +++ b/tests/streams.rs @@ -0,0 +1,279 @@ +//! Integration tests for the IOStream pipeline. + +use std::io::{Cursor, Read, Write, Seek, SeekFrom}; +use libfreemkv::*; +use libfreemkv::mux::meta::M2tsMeta; + +fn sample_disc_title() -> DiscTitle { + DiscTitle { + playlist: "Test Movie".into(), + playlist_id: 0, + duration_secs: 7200.0, + size_bytes: 0, + clips: Vec::new(), + streams: vec![ + Stream::Video(VideoStream { + pid: 0x1011, codec: Codec::Hevc, + resolution: "2160p".into(), frame_rate: "23.976".into(), + hdr: HdrFormat::Hdr10, color_space: ColorSpace::Bt709, + secondary: false, label: "Main".into(), + }), + Stream::Audio(AudioStream { + pid: 0x1100, codec: Codec::TrueHd, + channels: "7.1".into(), language: "eng".into(), + sample_rate: "48kHz".into(), secondary: false, + label: "English Atmos".into(), + }), + Stream::Audio(AudioStream { + pid: 0x1101, codec: Codec::Ac3, + channels: "5.1".into(), language: "fra".into(), + sample_rate: "48kHz".into(), secondary: false, + label: "French".into(), + }), + Stream::Subtitle(SubtitleStream { + pid: 0x1200, codec: Codec::Pgs, + language: "eng".into(), forced: false, + }), + ], + extents: Vec::new(), + } +} + +// ── URL parsing ──────────────────────────────────────────────── + +#[test] +fn parse_url_disc() { + let u = parse_url("disc://"); + assert_eq!(u.scheme, "disc"); + assert_eq!(u.path, ""); +} + +#[test] +fn parse_url_disc_device() { + let u = parse_url("disc:///dev/sg4"); + assert_eq!(u.scheme, "disc"); + assert_eq!(u.path, "/dev/sg4"); +} + +#[test] +fn parse_url_mkv() { + let u = parse_url("mkv://Dune.mkv"); + assert_eq!(u.scheme, "mkv"); + assert_eq!(u.path, "Dune.mkv"); +} + +#[test] +fn parse_url_network() { + let u = parse_url("network://10.1.7.11:9000"); + assert_eq!(u.scheme, "network"); + assert_eq!(u.path, "10.1.7.11:9000"); +} + +#[test] +fn parse_url_bare_mkv() { + let u = parse_url("Dune.mkv"); + assert_eq!(u.scheme, "mkv"); + assert_eq!(u.path, "Dune.mkv"); +} + +#[test] +fn parse_url_bare_m2ts() { + let u = parse_url("Dune.m2ts"); + assert_eq!(u.scheme, "m2ts"); + assert_eq!(u.path, "Dune.m2ts"); +} + +// ── M2TS metadata roundtrip ─────────────────────────────────── + +#[test] +fn m2ts_meta_roundtrip() { + let dt = sample_disc_title(); + let meta = M2tsMeta::from_title(&dt); + let restored = meta.to_title(); + + assert_eq!(restored.playlist, dt.playlist); + assert_eq!(restored.duration_secs, dt.duration_secs); + assert_eq!(restored.streams.len(), dt.streams.len()); + + // Check video + if let Stream::Video(v) = &restored.streams[0] { + assert_eq!(v.codec, Codec::Hevc); + assert_eq!(v.resolution, "2160p"); + assert_eq!(v.label, "Main"); + } else { panic!("expected video"); } + + // Check audio + if let Stream::Audio(a) = &restored.streams[1] { + assert_eq!(a.codec, Codec::TrueHd); + assert_eq!(a.language, "eng"); + assert_eq!(a.label, "English Atmos"); + } else { panic!("expected audio"); } + + // Check subtitle + if let Stream::Subtitle(s) = &restored.streams[3] { + assert_eq!(s.language, "eng"); + assert!(!s.forced); + } else { panic!("expected subtitle"); } +} + +// ── M2TS header write + read ────────────────────────────────── + +#[test] +fn m2ts_header_write_read() { + let dt = sample_disc_title(); + let meta = M2tsMeta::from_title(&dt); + + // Write header to buffer + let mut buf = Vec::new(); + libfreemkv::mux::meta::write_header(&mut buf, &meta).unwrap(); + + // Verify magic + assert_eq!(&buf[..4], b"FMKV"); + + // Verify 192-byte alignment + assert_eq!(buf.len() % 192, 0); + + // Read it back + let mut cursor = Cursor::new(&buf); + let read_back = libfreemkv::mux::meta::read_header(&mut cursor).unwrap().unwrap(); + assert_eq!(read_back.title, "Test Movie"); + assert_eq!(read_back.duration, 7200.0); + assert_eq!(read_back.streams.len(), 4); + + // Cursor should be at end of header (192-aligned) + assert_eq!(cursor.position() as usize, buf.len()); +} + +// ── M2tsStream write + read ─────────────────────────────────── + +#[test] +fn m2ts_stream_write_read() { + let dt = sample_disc_title(); + + // Build fake BD-TS packets + let mut ts_data = Vec::new(); + for i in 0..10u8 { + let mut pkt = [0u8; 192]; + pkt[4] = 0x47; + pkt[5] = 0x10; + pkt[6] = 0x11; + pkt[7] = 0x10; + pkt[8] = i; + ts_data.extend_from_slice(&pkt); + } + + // Write through M2tsStream to a Cursor + let output = Cursor::new(Vec::new()); + let mut stream = M2tsStream::new(output).meta(&dt); + stream.write_all(&ts_data).unwrap(); + stream.finish().unwrap(); + + // M2tsStream consumed the cursor — we need the inner data. + // For this test, write to a shared buffer instead. + // Use a second pass: write header + data manually to verify read side. + let mut encoded = Vec::new(); + let meta = M2tsMeta::from_title(&dt); + libfreemkv::mux::meta::write_header(&mut encoded, &meta).unwrap(); + encoded.extend_from_slice(&ts_data); + + // Read back + let cursor = Cursor::new(encoded); + let mut stream = M2tsStream::open(cursor).unwrap(); + let info = stream.info(); + assert_eq!(info.streams.len(), 4); + assert_eq!(info.duration_secs, 7200.0); + + // Read BD-TS data + let mut read_buf = vec![0u8; 192 * 10]; + let mut total = 0; + loop { + match stream.read(&mut read_buf[total..]) { + Ok(0) => break, + Ok(n) => total += n, + Err(_) => break, + } + } + + assert_eq!(total, 192 * 10); + for i in 0..10u8 { + assert_eq!(read_buf[i as usize * 192 + 4], 0x47); + assert_eq!(read_buf[i as usize * 192 + 8], i); + } +} + +// ── M2tsStream passthrough identity ─────────────────────────── + +#[test] +fn m2ts_passthrough_preserves_data() { + let dt = sample_disc_title(); + + // Create original BD-TS data + let mut original = Vec::new(); + for i in 0..100u8 { + let mut pkt = [0u8; 192]; + pkt[4] = 0x47; + pkt[5] = (i % 3) << 4; + pkt[6] = i; + for j in 8..192 { pkt[j] = i.wrapping_add(j as u8); } + original.extend_from_slice(&pkt); + } + + // Build encoded: header + original data + let mut encoded = Vec::new(); + let meta = M2tsMeta::from_title(&dt); + libfreemkv::mux::meta::write_header(&mut encoded, &meta).unwrap(); + encoded.extend_from_slice(&original); + + // Read back through M2tsStream + let cursor = Cursor::new(encoded); + let mut stream = M2tsStream::open(cursor).unwrap(); + let mut decoded = vec![0u8; original.len()]; + let mut total = 0; + loop { + match stream.read(&mut decoded[total..]) { + Ok(0) => break, + Ok(n) => total += n, + Err(_) => break, + } + } + + // BD-TS data must be byte-identical + assert_eq!(total, original.len()); + assert_eq!(decoded, original); +} + +// ── IOStream trait ──────────────────────────────────────────── + +#[test] +fn m2ts_implements_iostream() { + let dt = sample_disc_title(); + let output = Cursor::new(Vec::new()); + let stream = M2tsStream::new(output).meta(&dt); + + let mut boxed: Box<dyn IOStream> = Box::new(stream); + let meta = boxed.info(); + assert_eq!(meta.streams.len(), 4); + + let pkt = [0u8; 192]; + boxed.write_all(&pkt).unwrap(); + boxed.finish().unwrap(); +} + +#[test] +fn m2ts_read_returns_error_on_write_stream() { + let output = Cursor::new(Vec::new()); + let stream = M2tsStream::new(output); + let mut boxed: Box<dyn IOStream> = Box::new(stream); + let mut buf = [0u8; 10]; + assert!(boxed.read(&mut buf).is_err()); +} + +// ── DiscTitle::empty ────────────────────────────────────────── + +#[test] +fn disc_title_empty() { + let dt = DiscTitle::empty(); + assert_eq!(dt.streams.len(), 0); + assert_eq!(dt.duration_secs, 0.0); + assert!(dt.playlist.is_empty()); +}