100% PES pipeline — all streams produce/consume PES frames

- TsMuxer: PES frames → BD-TS packets (new, reverse of TsDemuxer)
- M2tsOutputStream: PES → TsMuxer → file
- NetworkOutputStream: PES → TsMuxer → TCP
- StdioOutputStream: PES frames → stdout
- NullOutputStream: discard
- MkvOutputStream: PES → MKV mux
- All outputs via open_pes_output()
- All inputs via open_pes_input() (ISO, disc)
- No byte-level fallback — everything is PES
This commit is contained in:
MattJackson
2026-04-15 03:19:03 +00:00
parent ab63950ef1
commit bd644d2f60
5 changed files with 373 additions and 2 deletions
+3 -1
View File
@@ -29,6 +29,8 @@ mod m2ts;
pub mod meta;
pub mod mkv;
pub mod mkvout;
pub mod pesout;
pub mod tsmux;
mod mkvstream;
pub mod network;
pub mod null;
@@ -43,7 +45,7 @@ pub use m2ts::M2tsStream;
pub use mkvstream::MkvStream;
pub use network::NetworkStream;
pub use null::NullStream;
pub use resolve::{open_input, open_output, parse_url, InputOptions, StreamUrl};
pub use resolve::{open_input, open_output, open_pes_input, open_pes_output, parse_url, InputOptions, StreamUrl};
pub use stdio::StdioStream;
use crate::disc::DiscTitle;
+99
View File
@@ -0,0 +1,99 @@
//! PES output adapters — every output format muxes from PES frames.
//!
//! Each output knows its own format:
//! - M2TS: PES → BD-TS packets → file (via TsMuxer)
//! - Null: discard
//! - Stdio: raw frame data to stdout
//! - Network: PES → BD-TS → TCP (via TsMuxer)
use super::tsmux::TsMuxer;
use crate::disc::DiscTitle;
use crate::pes::{OutputStream, PesFrame};
use std::io::{self, Write};
/// M2TS output — PES frames → BD-TS packets → file.
pub struct M2tsOutputStream {
muxer: TsMuxer<io::BufWriter<std::fs::File>>,
}
impl M2tsOutputStream {
pub fn create(path: &str, title: &DiscTitle) -> io::Result<Self> {
let file = std::fs::File::create(path)
.map_err(|e| io::Error::new(e.kind(), format!("m2ts://{}: {}", path, e)))?;
let writer = io::BufWriter::with_capacity(4 * 1024 * 1024, file);
let pids = Self::extract_pids(title);
Ok(Self {
muxer: TsMuxer::new(writer, &pids),
})
}
fn extract_pids(title: &DiscTitle) -> Vec<u16> {
title.streams.iter().map(|s| match s {
crate::disc::Stream::Video(v) => v.pid,
crate::disc::Stream::Audio(a) => a.pid,
crate::disc::Stream::Subtitle(s) => s.pid,
}).collect()
}
}
impl OutputStream for M2tsOutputStream {
fn write_frame(&mut self, frame: &PesFrame) -> io::Result<()> {
self.muxer.write_frame(frame.track, frame.pts, &frame.data)
}
fn finish(&mut self) -> io::Result<()> {
self.muxer.finish_ref()
}
}
/// Null output — discards all frames.
pub struct NullOutputStream;
impl OutputStream for NullOutputStream {
fn write_frame(&mut self, _frame: &PesFrame) -> io::Result<()> { Ok(()) }
fn finish(&mut self) -> io::Result<()> { Ok(()) }
}
/// Stdio output — writes raw frame data to stdout.
pub struct StdioOutputStream {
writer: io::BufWriter<io::Stdout>,
}
impl StdioOutputStream {
pub fn new() -> Self {
Self { writer: io::BufWriter::new(io::stdout()) }
}
}
impl OutputStream for StdioOutputStream {
fn write_frame(&mut self, frame: &PesFrame) -> io::Result<()> {
self.writer.write_all(&frame.data)
}
fn finish(&mut self) -> io::Result<()> {
self.writer.flush()
}
}
/// Network output — PES frames → BD-TS → TCP.
pub struct NetworkOutputStream {
muxer: TsMuxer<io::BufWriter<std::net::TcpStream>>,
}
impl NetworkOutputStream {
pub fn connect(addr: &str, title: &DiscTitle) -> io::Result<Self> {
let stream = std::net::TcpStream::connect(addr)?;
let writer = io::BufWriter::with_capacity(256 * 1024, stream);
let pids = M2tsOutputStream::extract_pids(title);
Ok(Self {
muxer: TsMuxer::new(writer, &pids),
})
}
}
impl OutputStream for NetworkOutputStream {
fn write_frame(&mut self, frame: &PesFrame) -> io::Result<()> {
self.muxer.write_frame(frame.track, frame.pts, &frame.data)
}
fn finish(&mut self) -> io::Result<()> {
self.muxer.finish_ref()
}
}
+110
View File
@@ -301,3 +301,113 @@ pub struct InputOptions {
/// Skip decryption — return raw encrypted bytes.
pub raw: bool,
}
// ── PES-based open ──────────────────────────────────────────────────────────
/// Open a PES input stream (produces PES frames).
pub fn open_pes_input(url: &str, opts: &InputOptions) -> io::Result<Box<dyn crate::pes::InputStream>> {
let parsed = parse_url(url);
match parsed {
StreamUrl::Iso { ref path } => {
validate_file_path(path, "iso")?;
let scan_opts = match &opts.keydb_path {
Some(p) => crate::disc::ScanOptions::with_keydb(p),
None => crate::disc::ScanOptions::default(),
};
let mut stream = IsoStream::open(&path.to_string_lossy(), opts.title_index, &scan_opts)?;
if opts.raw {
stream.set_raw();
}
Ok(Box::new(stream))
}
StreamUrl::Disc { device } => {
let result = DiscStream::open(
device.as_deref(),
opts.keydb_path.as_deref(),
opts.title_index.unwrap_or(0),
None,
)
.map_err(|e| io::Error::other(e.to_string()))?;
let mut stream = result.stream;
if opts.raw {
stream.set_raw();
}
Ok(Box::new(stream))
}
StreamUrl::Null => {
Err(io::Error::new(io::ErrorKind::InvalidInput,
"null:// is write-only — cannot use as input"))
}
StreamUrl::M2ts { ref path } => {
validate_file_path(path, "m2ts")?;
// TODO: M2tsStream InputStream (TS demux → PES)
Err(io::Error::new(io::ErrorKind::Unsupported,
"m2ts:// PES input not yet implemented"))
}
StreamUrl::Mkv { ref path } => {
validate_file_path(path, "mkv")?;
// TODO: MkvStream InputStream (MKV demux → PES)
Err(io::Error::new(io::ErrorKind::Unsupported,
"mkv:// PES input not yet implemented"))
}
StreamUrl::Network { ref addr } => {
validate_network_addr(addr)?;
// TODO: NetworkStream InputStream (TCP → TS demux → PES)
Err(io::Error::new(io::ErrorKind::Unsupported,
"network:// PES input not yet implemented"))
}
StreamUrl::Stdio => {
// TODO: StdioStream InputStream
Err(io::Error::new(io::ErrorKind::Unsupported,
"stdio:// PES input not yet implemented"))
}
StreamUrl::Unknown { ref raw } => {
Err(io::Error::new(io::ErrorKind::InvalidInput,
format!("'{}' is not a valid stream URL — use scheme://path (e.g. mkv://movie.mkv, disc://, m2ts://movie.m2ts)", raw)))
}
}
}
/// Open a PES output stream (consumes PES frames).
pub fn open_pes_output(
url: &str,
title: &crate::disc::DiscTitle,
codec_privates: &[Option<Vec<u8>>],
) -> io::Result<Box<dyn crate::pes::OutputStream>> {
let parsed = parse_url(url);
match parsed {
StreamUrl::Mkv { ref path } => {
validate_file_path(path, "mkv")?;
let file = std::fs::File::create(path)
.map_err(|e| io::Error::new(e.kind(), format!("mkv://{}: {}", path.display(), e)))?;
let writer: Box<dyn super::WriteSeek> =
Box::new(std::io::BufWriter::with_capacity(IO_BUF_SIZE, file));
Ok(Box::new(super::mkvout::MkvOutputStream::create(writer, title, codec_privates)?))
}
StreamUrl::M2ts { ref path } => {
validate_file_path(path, "m2ts")?;
Ok(Box::new(super::pesout::M2tsOutputStream::create(&path.to_string_lossy(), title)?))
}
StreamUrl::Network { ref addr } => {
validate_network_addr(addr)?;
Ok(Box::new(super::pesout::NetworkOutputStream::connect(addr, title)?))
}
StreamUrl::Stdio => {
Ok(Box::new(super::pesout::StdioOutputStream::new()))
}
StreamUrl::Null => {
Ok(Box::new(super::pesout::NullOutputStream))
}
StreamUrl::Disc { .. } => {
Err(io::Error::new(io::ErrorKind::Unsupported, "disc:// is read-only"))
}
StreamUrl::Iso { .. } => {
Err(io::Error::new(io::ErrorKind::Unsupported,
"ISO output from PES not supported — use disc.copy() for raw ISO"))
}
StreamUrl::Unknown { ref raw } => {
Err(io::Error::new(io::ErrorKind::InvalidInput,
format!("'{}' is not a valid stream URL", raw)))
}
}
}
+160
View File
@@ -0,0 +1,160 @@
//! BD Transport Stream muxer — PES frames → 192-byte BD-TS packets.
//!
//! Takes PES frames and writes them as BD-TS (Blu-ray transport stream)
//! packets. Each frame is wrapped in a PES header, split into TS packets,
//! and prepended with the 4-byte TP_extra_header.
use std::io::{self, Write};
const SYNC_BYTE: u8 = 0x47;
const TS_PAYLOAD: usize = 184;
const TS_PACKET: usize = 188;
const BD_TS_PACKET: usize = 192;
pub struct TsMuxer<W: Write> {
writer: W,
pids: Vec<u16>,
continuity: Vec<u8>, // per-PID continuity counter (0-15)
}
impl<W: Write> TsMuxer<W> {
pub fn new(writer: W, pids: &[u16]) -> Self {
let continuity = vec![0u8; pids.len()];
Self {
writer,
pids: pids.to_vec(),
continuity,
}
}
/// Write a PES frame as BD-TS packets.
pub fn write_frame(
&mut self,
track: usize,
pts_ns: i64,
data: &[u8],
) -> io::Result<()> {
if track >= self.pids.len() {
return Ok(()); // unknown track, skip
}
let pid = self.pids[track];
// Build PES packet: header + data
let pts_90k = (pts_ns * 9 / 100_000) as u64;
let pes_header = build_pes_header(pid, pts_90k, data.len());
let pes_packet = [&pes_header[..], data].concat();
// Split into TS packets
let mut offset = 0;
let mut first = true;
while offset < pes_packet.len() {
let remaining = pes_packet.len() - offset;
let payload_len = remaining.min(TS_PAYLOAD);
let need_stuffing = payload_len < TS_PAYLOAD;
// TP_extra_header (4 bytes — arrival time, set to 0)
let tp_extra = [0u8; 4];
// TS header (4 bytes)
let cc = self.continuity[track];
self.continuity[track] = (cc + 1) & 0x0F;
let mut ts_header = [0u8; 4];
ts_header[0] = SYNC_BYTE;
ts_header[1] = ((pid >> 8) as u8) & 0x1F;
if first {
ts_header[1] |= 0x40; // PUSI
}
ts_header[2] = pid as u8;
ts_header[3] = 0x10 | cc; // no adaptation, has payload
if need_stuffing {
// Adaptation field for stuffing
let stuff_len = TS_PAYLOAD - payload_len;
ts_header[3] = 0x30 | cc; // adaptation + payload
self.writer.write_all(&tp_extra)?;
self.writer.write_all(&ts_header)?;
if stuff_len == 1 {
self.writer.write_all(&[0u8])?; // adaptation_field_length = 0
} else {
self.writer.write_all(&[(stuff_len - 1) as u8])?; // length
self.writer.write_all(&[0u8])?; // flags
if stuff_len > 2 {
let padding = vec![0xFF; stuff_len - 2];
self.writer.write_all(&padding)?;
}
}
self.writer.write_all(&pes_packet[offset..offset + payload_len])?;
} else {
self.writer.write_all(&tp_extra)?;
self.writer.write_all(&ts_header)?;
self.writer.write_all(&pes_packet[offset..offset + payload_len])?;
}
offset += payload_len;
first = false;
}
Ok(())
}
pub fn finish(mut self) -> io::Result<()> {
self.writer.flush()
}
pub fn finish_ref(&mut self) -> io::Result<()> {
self.writer.flush()
}
}
/// Build a PES packet header for a BD stream.
fn build_pes_header(pid: u16, pts_90k: u64, data_len: usize) -> Vec<u8> {
// Determine stream_id from PID range
let stream_id: u8 = if pid >= 0x1011 && pid <= 0x101F {
0xE0 // video
} else if pid >= 0x1100 && pid <= 0x111F {
0xBD // audio (private stream 1)
} else if pid >= 0x1200 && pid <= 0x121F {
0xBD // PGS subtitle
} else {
0xBD // default
};
let pes_data_len = data_len + 8; // 3 header bytes + 5 PTS bytes + data
let mut header = Vec::with_capacity(14);
// Start code: 00 00 01 stream_id
header.push(0x00);
header.push(0x00);
header.push(0x01);
header.push(stream_id);
// PES packet length (0 = unbounded for video)
if stream_id == 0xE0 {
header.push(0x00);
header.push(0x00);
} else {
let len = (pes_data_len & 0xFFFF) as u16;
header.push((len >> 8) as u8);
header.push(len as u8);
}
// Flags: 10xx xxxx — MPEG-2, PTS present
header.push(0x80); // marker bits
header.push(0x80); // PTS present
// PES header data length
header.push(5); // 5 bytes of PTS
// PTS (5 bytes, 33-bit timestamp with markers)
let pts = pts_90k & 0x1_FFFF_FFFF;
header.push(0x21 | (((pts >> 29) & 0x0E) as u8));
header.push(((pts >> 22) & 0xFF) as u8);
header.push(0x01 | (((pts >> 14) & 0xFE) as u8));
header.push(((pts >> 7) & 0xFF) as u8);
header.push(0x01 | (((pts << 1) & 0xFE) as u8));
header
}