Add IOStream trait and stream-based I/O architecture

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).
This commit is contained in:
MattJackson
2026-04-10 19:13:53 -07:00
parent 37c98e8826
commit 8c2f3898b8
20 changed files with 2188 additions and 268 deletions
+19 -6
View File
@@ -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<Title>,
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,
+2 -2
View File
@@ -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();
+8 -1
View File
@@ -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};
+16 -4
View File
@@ -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
}
+2 -1
View File
@@ -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;
+2 -1
View File
@@ -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;
+20 -13
View File
@@ -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(),
}]
}
+130
View File
@@ -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(()) }
}
+148 -4
View File
@@ -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
+139
View File
@@ -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")),
}
}
}
+293
View File
@@ -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,
}
}
+7
View File
@@ -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)?;
+512
View File
@@ -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,
]
}
+48 -13
View File
@@ -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.0.0.1: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 {}
+128
View File
@@ -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")),
}
}
}
+43
View File
@@ -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"))
}
}
+130
View File
@@ -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 }
}
}
-218
View File
@@ -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(())
}
}
+262 -5
View File
@@ -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::*;