Remove Seek/File dependencies from stream readers
Streams are streams — they take impl Read, not Read+Seek or File. - MkvStream::open takes impl Read (was Read+Seek) - EBML element skipping uses skip_bytes() instead of seek(Current) - Byte position tracking uses remaining-bytes counter, not stream_position() - Removed file_size from open (progress is CLI concern) - M2tsStream::open takes impl Read (was Read+Seek) - Buffers first 1MB for FMKV header / PMT scan - Uses chain reader (buffered head + rest) for sequential reading - Duration unknown without seeking (0.0) — CLI can set from metadata - Removed ReadSeek trait (no longer needed) - WriteSeek kept (MKV muxer container format requires seeking internally)
This commit is contained in:
+44
-42
@@ -3,9 +3,9 @@
|
|||||||
//! Write: prepends FMKV metadata header, then passes through BD-TS bytes.
|
//! Write: prepends FMKV metadata header, then passes through BD-TS bytes.
|
||||||
//! Read: extracts metadata header (or scans PMT), then yields BD-TS bytes.
|
//! Read: extracts metadata header (or scans PMT), then yields BD-TS bytes.
|
||||||
|
|
||||||
use super::{meta, ts, IOStream, ReadSeek};
|
use super::{meta, ts, IOStream};
|
||||||
use crate::disc::{DiscTitle, Stream as DiscStream};
|
use crate::disc::{DiscTitle, Stream as DiscStream};
|
||||||
use std::io::{self, Read, Seek, SeekFrom, Write};
|
use std::io::{self, Read, Write};
|
||||||
|
|
||||||
type PesSetup = (Vec<u16>, Vec<(u16, Box<dyn super::codec::CodecParser>)>, Vec<(u16, usize)>);
|
type PesSetup = (Vec<u16>, Vec<(u16, Box<dyn super::codec::CodecParser>)>, Vec<(u16, usize)>);
|
||||||
|
|
||||||
@@ -18,10 +18,25 @@ enum Mode {
|
|||||||
header_written: bool,
|
header_written: bool,
|
||||||
},
|
},
|
||||||
Read {
|
Read {
|
||||||
reader: Box<dyn ReadSeek>,
|
reader: Box<dyn Read>,
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Read as many bytes as possible into buf (multiple read calls if needed).
|
||||||
|
/// Bounded by buf.len() — caller controls max bytes read.
|
||||||
|
fn read_fill(r: &mut impl Read, buf: &mut [u8]) -> io::Result<usize> {
|
||||||
|
let mut total = 0;
|
||||||
|
while total < buf.len() {
|
||||||
|
match r.read(&mut buf[total..]) {
|
||||||
|
Ok(0) => break,
|
||||||
|
Ok(n) => total += n,
|
||||||
|
Err(e) if e.kind() == io::ErrorKind::Interrupted => continue,
|
||||||
|
Err(e) => return Err(e),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Ok(total)
|
||||||
|
}
|
||||||
|
|
||||||
/// BD transport stream with embedded metadata.
|
/// BD transport stream with embedded metadata.
|
||||||
pub struct M2tsStream {
|
pub struct M2tsStream {
|
||||||
disc_title: DiscTitle,
|
disc_title: DiscTitle,
|
||||||
@@ -82,69 +97,56 @@ impl M2tsStream {
|
|||||||
self
|
self
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Open an m2ts file for reading.
|
/// Open an M2TS stream for reading. Takes any Read source — file, pipe, socket.
|
||||||
///
|
///
|
||||||
/// Tries FMKV metadata header first. Falls back to PMT scan + PTS duration.
|
/// Tries FMKV metadata header first. Falls back to PMT scan of first 1 MB.
|
||||||
pub fn open(mut reader: impl Read + Seek + 'static) -> io::Result<Self> {
|
pub fn open(mut reader: impl Read + 'static) -> io::Result<Self> {
|
||||||
// Get total file size for progress tracking
|
// Read first chunk — enough for FMKV header or PMT scan
|
||||||
let file_size = reader.seek(SeekFrom::End(0))?;
|
let mut head = vec![0u8; SCAN_SIZE];
|
||||||
reader.seek(SeekFrom::Start(0))?;
|
let head_len = read_fill(&mut reader, &mut head)?;
|
||||||
|
head.truncate(head_len);
|
||||||
|
|
||||||
// Try FMKV metadata header — save position so we can seek back on failure
|
// Try FMKV metadata header from the buffered head
|
||||||
let start = reader.stream_position()?;
|
let mut cursor = io::Cursor::new(&head);
|
||||||
if let Ok(Some(m)) = meta::read_header(&mut reader) {
|
if let Ok(Some(m)) = meta::read_header(&mut cursor) {
|
||||||
let header_end = reader.stream_position()?;
|
let header_end = cursor.position() as usize;
|
||||||
let content_size = file_size.saturating_sub(header_end);
|
|
||||||
let codec_privates = m.codec_privates();
|
|
||||||
let title = m.to_title();
|
let title = m.to_title();
|
||||||
let (pids, parsers, pid_to_track) = Self::setup_pes(&title.streams);
|
let (pids, parsers, pid_to_track) = Self::setup_pes(&title.streams);
|
||||||
|
// Chain: remaining head bytes + rest of reader
|
||||||
|
let remaining_head = &head[header_end..];
|
||||||
|
let chain: Box<dyn Read> = Box::new(io::Cursor::new(remaining_head.to_vec()).chain(reader));
|
||||||
return Ok(Self {
|
return Ok(Self {
|
||||||
disc_title: title,
|
disc_title: title.clone(),
|
||||||
mode: Mode::Read {
|
mode: Mode::Read { reader: chain },
|
||||||
reader: Box::new(reader),
|
|
||||||
},
|
|
||||||
finished: false,
|
finished: false,
|
||||||
content_size: Some(content_size),
|
content_size: None,
|
||||||
demuxer: if pids.is_empty() { None } else { Some(ts::TsDemuxer::new(&pids)) },
|
demuxer: if pids.is_empty() { None } else { Some(ts::TsDemuxer::new(&pids)) },
|
||||||
parsers,
|
parsers,
|
||||||
pending_frames: std::collections::VecDeque::new(),
|
pending_frames: std::collections::VecDeque::new(),
|
||||||
pid_to_track,
|
pid_to_track,
|
||||||
pes_eof: false,
|
pes_eof: false,
|
||||||
stored_codec_privates: codec_privates,
|
stored_codec_privates: title.codec_privates,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
// No FMKV header — seek back and try PMT scan
|
// No FMKV header — scan head for PMT
|
||||||
reader.seek(SeekFrom::Start(start))?;
|
let streams = ts::scan_streams(&head)
|
||||||
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"))?;
|
.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))?;
|
|
||||||
|
|
||||||
let (pids, parsers, pid_to_track) = Self::setup_pes(&streams);
|
let (pids, parsers, pid_to_track) = Self::setup_pes(&streams);
|
||||||
|
|
||||||
|
// Chain: full head (it's all TS data) + rest of reader
|
||||||
|
let chain: Box<dyn Read> = Box::new(io::Cursor::new(head).chain(reader));
|
||||||
|
|
||||||
Ok(Self {
|
Ok(Self {
|
||||||
disc_title: DiscTitle {
|
disc_title: DiscTitle {
|
||||||
duration_secs: duration,
|
duration_secs: 0.0, // unknown without seeking
|
||||||
streams,
|
streams,
|
||||||
..DiscTitle::empty()
|
..DiscTitle::empty()
|
||||||
},
|
},
|
||||||
mode: Mode::Read {
|
mode: Mode::Read { reader: chain },
|
||||||
reader: Box::new(reader),
|
|
||||||
},
|
|
||||||
finished: false,
|
finished: false,
|
||||||
content_size: Some(file_size),
|
content_size: None,
|
||||||
demuxer: if pids.is_empty() { None } else { Some(ts::TsDemuxer::new(&pids)) },
|
demuxer: if pids.is_empty() { None } else { Some(ts::TsDemuxer::new(&pids)) },
|
||||||
parsers,
|
parsers,
|
||||||
pending_frames: std::collections::VecDeque::new(),
|
pending_frames: std::collections::VecDeque::new(),
|
||||||
|
|||||||
+44
-41
@@ -7,11 +7,19 @@ use super::codec::{self, CodecParser};
|
|||||||
use super::lookahead::{LookaheadBuffer, LookaheadState, DEFAULT_LOOKAHEAD_SIZE};
|
use super::lookahead::{LookaheadBuffer, LookaheadState, DEFAULT_LOOKAHEAD_SIZE};
|
||||||
use super::mkv::{MkvMuxer, MkvTrack};
|
use super::mkv::{MkvMuxer, MkvTrack};
|
||||||
use super::ts::TsDemuxer;
|
use super::ts::TsDemuxer;
|
||||||
use super::{ebml, IOStream, ReadSeek, WriteSeek};
|
use super::{ebml, IOStream, WriteSeek};
|
||||||
|
use std::io::Seek;
|
||||||
|
|
||||||
type MkvHeaderResult = io::Result<(crate::disc::DiscTitle, Vec<(u16, Vec<u8>)>)>;
|
type MkvHeaderResult = io::Result<(crate::disc::DiscTitle, Vec<(u16, Vec<u8>)>)>;
|
||||||
|
|
||||||
|
/// Skip `n` bytes on a forward-only reader (no Seek required).
|
||||||
|
fn skip_bytes(r: &mut impl Read, n: u64) -> io::Result<()> {
|
||||||
|
io::copy(&mut r.take(n), &mut io::sink())?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
use crate::disc::*;
|
use crate::disc::*;
|
||||||
use std::io::{self, Read, Seek, SeekFrom, Write};
|
use std::io::{self, Read, Write};
|
||||||
|
|
||||||
/// Lookahead buffer for codec header detection (5 MB default).
|
/// Lookahead buffer for codec header detection (5 MB default).
|
||||||
const DEFAULT_MAX_BUFFER: usize = DEFAULT_LOOKAHEAD_SIZE;
|
const DEFAULT_MAX_BUFFER: usize = DEFAULT_LOOKAHEAD_SIZE;
|
||||||
@@ -35,7 +43,7 @@ struct WriteState {
|
|||||||
}
|
}
|
||||||
|
|
||||||
struct ReadState {
|
struct ReadState {
|
||||||
reader: Box<dyn ReadSeek>,
|
reader: Box<dyn Read>,
|
||||||
buf: Vec<u8>,
|
buf: Vec<u8>,
|
||||||
pos: usize,
|
pos: usize,
|
||||||
len: usize,
|
len: usize,
|
||||||
@@ -129,9 +137,7 @@ impl MkvStream {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Open an MKV file for reading.
|
/// Open an MKV file for reading.
|
||||||
pub fn open(mut reader: impl Read + Seek + 'static) -> io::Result<Self> {
|
pub fn open(mut reader: impl Read + 'static) -> io::Result<Self> {
|
||||||
let file_size = reader.seek(SeekFrom::End(0))?;
|
|
||||||
reader.seek(SeekFrom::Start(0))?;
|
|
||||||
let (disc_title, codec_privates) = parse_mkv_header(&mut reader)?;
|
let (disc_title, codec_privates) = parse_mkv_header(&mut reader)?;
|
||||||
Ok(Self {
|
Ok(Self {
|
||||||
disc_title,
|
disc_title,
|
||||||
@@ -146,7 +152,7 @@ impl MkvStream {
|
|||||||
}),
|
}),
|
||||||
max_buffer: 0,
|
max_buffer: 0,
|
||||||
finished: false,
|
finished: false,
|
||||||
file_size: Some(file_size),
|
file_size: None,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -196,9 +202,7 @@ impl crate::pes::Stream for MkvStream {
|
|||||||
}));
|
}));
|
||||||
}
|
}
|
||||||
_ => {
|
_ => {
|
||||||
// Skip unknown element
|
skip_bytes(&mut rs.reader, size)?;
|
||||||
let mut skip = vec![0u8; size as usize];
|
|
||||||
let _ = rs.reader.read_exact(&mut skip);
|
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -396,7 +400,7 @@ impl Read for MkvStream {
|
|||||||
}
|
}
|
||||||
_ => {
|
_ => {
|
||||||
if size != u64::MAX && size > 0 {
|
if size != u64::MAX && size > 0 {
|
||||||
rs.reader.seek(SeekFrom::Current(size as i64))?;
|
skip_bytes(&mut rs.reader, size)?;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -484,7 +488,7 @@ fn write_pes(
|
|||||||
|
|
||||||
/// Returns (DiscTitle, codec_privates: Vec<(track_number, codec_private_bytes)>)
|
/// Returns (DiscTitle, codec_privates: Vec<(track_number, codec_private_bytes)>)
|
||||||
fn parse_mkv_header(
|
fn parse_mkv_header(
|
||||||
r: &mut (impl Read + Seek),
|
r: &mut impl Read,
|
||||||
) -> MkvHeaderResult {
|
) -> MkvHeaderResult {
|
||||||
let mut title = String::new();
|
let mut title = String::new();
|
||||||
let mut duration_ms = 0.0f64;
|
let mut duration_ms = 0.0f64;
|
||||||
@@ -502,7 +506,7 @@ fn parse_mkv_header(
|
|||||||
"EBML header too large",
|
"EBML header too large",
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
r.seek(SeekFrom::Current(size as i64))?;
|
skip_bytes(r, size as u64)?;
|
||||||
|
|
||||||
let (id, _, _) = ebml::read_element_header(r)?;
|
let (id, _, _) = ebml::read_element_header(r)?;
|
||||||
if id != ebml::SEGMENT {
|
if id != ebml::SEGMENT {
|
||||||
@@ -522,24 +526,24 @@ fn parse_mkv_header(
|
|||||||
|
|
||||||
match id {
|
match id {
|
||||||
ebml::INFO => {
|
ebml::INFO => {
|
||||||
let end = r.stream_position()? + size;
|
let mut remaining = size;
|
||||||
while r.stream_position()? < end {
|
while remaining > 0 {
|
||||||
let (cid, cs, _) = ebml::read_element_header(r)?;
|
let (cid, cs, hlen) = ebml::read_element_header(r)?;
|
||||||
|
remaining = remaining.saturating_sub(hlen as u64 + cs);
|
||||||
match cid {
|
match cid {
|
||||||
ebml::TIMESTAMP_SCALE => ts_scale = ebml::read_uint_val(r, cs as usize)?,
|
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::DURATION => duration_ms = ebml::read_float_val(r, cs as usize)?,
|
||||||
ebml::TITLE => title = ebml::read_string_val(r, cs as usize)?,
|
ebml::TITLE => title = ebml::read_string_val(r, cs as usize)?,
|
||||||
_ => {
|
_ => { skip_bytes(r, cs)?; }
|
||||||
r.seek(SeekFrom::Current(cs as i64))?;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
got_info = true;
|
got_info = true;
|
||||||
}
|
}
|
||||||
ebml::TRACKS => {
|
ebml::TRACKS => {
|
||||||
let end = r.stream_position()? + size;
|
let mut remaining = size;
|
||||||
while r.stream_position()? < end {
|
while remaining > 0 {
|
||||||
let (cid, cs, _) = ebml::read_element_header(r)?;
|
let (cid, cs, hlen) = ebml::read_element_header(r)?;
|
||||||
|
remaining = remaining.saturating_sub(hlen as u64 + cs);
|
||||||
if cid == ebml::TRACK_ENTRY {
|
if cid == ebml::TRACK_ENTRY {
|
||||||
let (stream, tnum, cp) = parse_track(r, cs)?;
|
let (stream, tnum, cp) = parse_track(r, cs)?;
|
||||||
if let Some(s) = stream {
|
if let Some(s) = stream {
|
||||||
@@ -549,14 +553,14 @@ fn parse_mkv_header(
|
|||||||
codec_privates.push((tnum, cp));
|
codec_privates.push((tnum, cp));
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
r.seek(SeekFrom::Current(cs as i64))?;
|
skip_bytes(r, cs)?;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
got_tracks = true;
|
got_tracks = true;
|
||||||
}
|
}
|
||||||
ebml::CLUSTER => break,
|
ebml::CLUSTER => break,
|
||||||
_ if size != u64::MAX => {
|
_ if size != u64::MAX => {
|
||||||
r.seek(SeekFrom::Current(size as i64))?;
|
skip_bytes(r, size as u64)?;
|
||||||
}
|
}
|
||||||
_ => break,
|
_ => break,
|
||||||
}
|
}
|
||||||
@@ -573,17 +577,18 @@ fn parse_mkv_header(
|
|||||||
|
|
||||||
/// Returns (stream, track_number, codec_private_bytes)
|
/// Returns (stream, track_number, codec_private_bytes)
|
||||||
fn parse_track(
|
fn parse_track(
|
||||||
r: &mut (impl Read + Seek),
|
r: &mut impl Read,
|
||||||
size: u64,
|
size: u64,
|
||||||
) -> io::Result<(Option<crate::disc::Stream>, u16, Option<Vec<u8>>)> {
|
) -> io::Result<(Option<crate::disc::Stream>, u16, Option<Vec<u8>>)> {
|
||||||
let end = r.stream_position()? + size;
|
|
||||||
let (mut ttype, mut tnum) = (0u64, 0u16);
|
let (mut ttype, mut tnum) = (0u64, 0u16);
|
||||||
let (mut codec_id, mut lang, mut name) = (String::new(), String::from("und"), String::new());
|
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);
|
let (mut ph, mut sr, mut ch, mut forced) = (0u32, 0.0f64, 0u8, false);
|
||||||
let mut codec_priv: Option<Vec<u8>> = None;
|
let mut codec_priv: Option<Vec<u8>> = None;
|
||||||
|
|
||||||
while r.stream_position()? < end {
|
let mut remaining = size;
|
||||||
let (cid, cs, _) = ebml::read_element_header(r)?;
|
while remaining > 0 {
|
||||||
|
let (cid, cs, hlen) = ebml::read_element_header(r)?;
|
||||||
|
remaining = remaining.saturating_sub(hlen as u64 + cs);
|
||||||
match cid {
|
match cid {
|
||||||
ebml::TRACK_NUMBER => tnum = ebml::read_uint_val(r, cs as usize)? as u16,
|
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::TRACK_TYPE => ttype = ebml::read_uint_val(r, cs as usize)?,
|
||||||
@@ -593,32 +598,30 @@ fn parse_track(
|
|||||||
ebml::TRACK_NAME => name = 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::FLAG_FORCED => forced = ebml::read_uint_val(r, cs as usize)? != 0,
|
||||||
ebml::VIDEO => {
|
ebml::VIDEO => {
|
||||||
let ve = r.stream_position()? + cs;
|
let mut vrem = cs;
|
||||||
while r.stream_position()? < ve {
|
while vrem > 0 {
|
||||||
let (vid, vs, _) = ebml::read_element_header(r)?;
|
let (vid, vs, vhlen) = ebml::read_element_header(r)?;
|
||||||
|
vrem = vrem.saturating_sub(vhlen as u64 + vs);
|
||||||
if vid == ebml::PIXEL_HEIGHT {
|
if vid == ebml::PIXEL_HEIGHT {
|
||||||
ph = ebml::read_uint_val(r, vs as usize)? as u32;
|
ph = ebml::read_uint_val(r, vs as usize)? as u32;
|
||||||
} else {
|
} else {
|
||||||
r.seek(SeekFrom::Current(vs as i64))?;
|
skip_bytes(r, vs)?;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
ebml::AUDIO => {
|
ebml::AUDIO => {
|
||||||
let ae = r.stream_position()? + cs;
|
let mut arem = cs;
|
||||||
while r.stream_position()? < ae {
|
while arem > 0 {
|
||||||
let (aid, as_, _) = ebml::read_element_header(r)?;
|
let (aid, as_, ahlen) = ebml::read_element_header(r)?;
|
||||||
|
arem = arem.saturating_sub(ahlen as u64 + as_);
|
||||||
match aid {
|
match aid {
|
||||||
ebml::SAMPLING_FREQUENCY => sr = ebml::read_float_val(r, as_ as usize)?,
|
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,
|
ebml::CHANNELS => ch = ebml::read_uint_val(r, as_ as usize)? as u8,
|
||||||
_ => {
|
_ => { skip_bytes(r, as_)?; }
|
||||||
r.seek(SeekFrom::Current(as_ as i64))?;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
_ => { skip_bytes(r, cs)?; }
|
||||||
_ => {
|
|
||||||
r.seek(SeekFrom::Current(cs as i64))?;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+1
-4
@@ -75,9 +75,6 @@ pub trait IOStream: Read + Write {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Combined traits for internal trait objects.
|
// WriteSeek — used internally by MKV muxer (container format requires seeking).
|
||||||
pub(crate) trait ReadSeek: Read + Seek {}
|
|
||||||
impl<T: Read + Seek> ReadSeek for T {}
|
|
||||||
|
|
||||||
pub trait WriteSeek: Write + Seek {}
|
pub trait WriteSeek: Write + Seek {}
|
||||||
impl<T: Write + Seek> WriteSeek for T {}
|
impl<T: Write + Seek> WriteSeek for T {}
|
||||||
|
|||||||
Reference in New Issue
Block a user