M2TS implements Stream (read PES frames via TsDemux)
M2tsStream::read() → TsDemux → CodecParser → PesFrame. M2TS input now works through PES pipeline. MKV input deferred (needs EBML → PES extraction).
This commit is contained in:
+116
-1
@@ -27,6 +27,13 @@ pub struct M2tsStream {
|
||||
finished: bool,
|
||||
/// Content size in bytes (file size minus header), set for read mode.
|
||||
content_size: Option<u64>,
|
||||
// PES support
|
||||
demuxer: Option<ts::TsDemuxer>,
|
||||
parsers: Vec<(u16, Box<dyn super::codec::CodecParser>)>,
|
||||
pending_frames: std::collections::VecDeque<crate::pes::PesFrame>,
|
||||
pid_to_track: Vec<(u16, usize)>,
|
||||
pes_buf: Vec<u8>,
|
||||
pes_eof: bool,
|
||||
}
|
||||
|
||||
impl M2tsStream {
|
||||
@@ -40,9 +47,32 @@ impl M2tsStream {
|
||||
},
|
||||
finished: false,
|
||||
content_size: None,
|
||||
demuxer: None,
|
||||
parsers: Vec::new(),
|
||||
pending_frames: std::collections::VecDeque::new(),
|
||||
pid_to_track: Vec::new(),
|
||||
pes_buf: Vec::new(),
|
||||
pes_eof: false,
|
||||
}
|
||||
}
|
||||
|
||||
fn setup_pes(streams: &[DiscStream]) -> (Vec<u16>, Vec<(u16, Box<dyn super::codec::CodecParser>)>, Vec<(u16, usize)>) {
|
||||
let mut pids = Vec::new();
|
||||
let mut parsers: Vec<(u16, Box<dyn super::codec::CodecParser>)> = Vec::new();
|
||||
let mut pid_to_track = Vec::new();
|
||||
for (i, s) in streams.iter().enumerate() {
|
||||
let (pid, codec) = match s {
|
||||
DiscStream::Video(v) => (v.pid, v.codec),
|
||||
DiscStream::Audio(a) => (a.pid, a.codec),
|
||||
DiscStream::Subtitle(s) => (s.pid, s.codec),
|
||||
};
|
||||
pids.push(pid);
|
||||
pid_to_track.push((pid, i));
|
||||
parsers.push((pid, super::codec::parser_for_codec(codec)));
|
||||
}
|
||||
(pids, parsers, pid_to_track)
|
||||
}
|
||||
|
||||
/// Set stream metadata. Returns self for chaining.
|
||||
pub fn meta(mut self, dt: &DiscTitle) -> Self {
|
||||
self.disc_title = dt.clone();
|
||||
@@ -61,13 +91,21 @@ impl M2tsStream {
|
||||
if let Ok(Some(m)) = meta::read_header(&mut reader) {
|
||||
let header_end = reader.stream_position()?;
|
||||
let content_size = file_size.saturating_sub(header_end);
|
||||
let title = m.to_title();
|
||||
let (pids, parsers, pid_to_track) = Self::setup_pes(&title.streams);
|
||||
return Ok(Self {
|
||||
disc_title: m.to_title(),
|
||||
disc_title: title,
|
||||
mode: Mode::Read {
|
||||
reader: Box::new(reader),
|
||||
},
|
||||
finished: false,
|
||||
content_size: Some(content_size),
|
||||
demuxer: if pids.is_empty() { None } else { Some(ts::TsDemuxer::new(&pids)) },
|
||||
parsers,
|
||||
pending_frames: std::collections::VecDeque::new(),
|
||||
pid_to_track,
|
||||
pes_buf: vec![0u8; 192 * 1024],
|
||||
pes_eof: false,
|
||||
});
|
||||
}
|
||||
|
||||
@@ -89,6 +127,8 @@ impl M2tsStream {
|
||||
|
||||
reader.seek(SeekFrom::Start(0))?;
|
||||
|
||||
let (pids, parsers, pid_to_track) = Self::setup_pes(&streams);
|
||||
|
||||
Ok(Self {
|
||||
disc_title: DiscTitle {
|
||||
duration_secs: duration,
|
||||
@@ -100,10 +140,85 @@ impl M2tsStream {
|
||||
},
|
||||
finished: false,
|
||||
content_size: Some(file_size),
|
||||
demuxer: if pids.is_empty() { None } else { Some(ts::TsDemuxer::new(&pids)) },
|
||||
parsers,
|
||||
pending_frames: std::collections::VecDeque::new(),
|
||||
pid_to_track,
|
||||
pes_buf: vec![0u8; 192 * 1024],
|
||||
pes_eof: false,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl crate::pes::Stream for M2tsStream {
|
||||
fn read(&mut self) -> io::Result<Option<crate::pes::PesFrame>> {
|
||||
if let Some(frame) = self.pending_frames.pop_front() {
|
||||
return Ok(Some(frame));
|
||||
}
|
||||
if self.pes_eof { return Ok(None); }
|
||||
|
||||
loop {
|
||||
let reader = match &mut self.mode {
|
||||
Mode::Read { reader } => reader,
|
||||
_ => return Err(io::Error::new(io::ErrorKind::Unsupported, "not in read mode")),
|
||||
};
|
||||
let mut buf = vec![0u8; 192 * 1024];
|
||||
let n = reader.read(&mut buf)?;
|
||||
if n == 0 {
|
||||
self.pes_eof = true;
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
if let Some(ref mut demuxer) = self.demuxer {
|
||||
let packets = demuxer.feed(&buf[..n]);
|
||||
for pes in &packets {
|
||||
if let Some((_, track)) = self.pid_to_track.iter().find(|(pid, _)| *pid == pes.pid) {
|
||||
if let Some((_, parser)) = self.parsers.iter_mut().find(|(pid, _)| *pid == pes.pid) {
|
||||
for frame in parser.parse(pes) {
|
||||
self.pending_frames.push_back(
|
||||
crate::pes::PesFrame::from_codec_frame(*track, frame)
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(frame) = self.pending_frames.pop_front() {
|
||||
return Ok(Some(frame));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn write(&mut self, _frame: &crate::pes::PesFrame) -> io::Result<()> {
|
||||
Err(io::Error::new(io::ErrorKind::Unsupported, "use M2tsOutputStream for writing"))
|
||||
}
|
||||
|
||||
fn finish(&mut self) -> io::Result<()> { Ok(()) }
|
||||
|
||||
fn info(&self) -> &crate::disc::DiscTitle { &self.disc_title }
|
||||
|
||||
fn codec_private(&self, track: usize) -> Option<Vec<u8>> {
|
||||
let pid = self.pid_to_track.iter()
|
||||
.find(|(_, idx)| *idx == track)
|
||||
.map(|(pid, _)| *pid)?;
|
||||
self.parsers.iter()
|
||||
.find(|(p, _)| *p == pid)
|
||||
.and_then(|(_, parser)| parser.codec_private())
|
||||
}
|
||||
|
||||
fn headers_ready(&self) -> bool {
|
||||
for (idx, s) in self.disc_title.streams.iter().enumerate() {
|
||||
if let crate::disc::Stream::Video(v) = s {
|
||||
if !v.secondary && self.codec_private(idx).is_none() {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
}
|
||||
true
|
||||
}
|
||||
}
|
||||
|
||||
impl IOStream for M2tsStream {
|
||||
fn info(&self) -> &DiscTitle {
|
||||
&self.disc_title
|
||||
|
||||
+6
-5
@@ -340,15 +340,16 @@ pub fn input(url: &str, opts: &InputOptions) -> io::Result<Box<dyn crate::pes::S
|
||||
}
|
||||
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"))
|
||||
let file = std::fs::File::open(path)
|
||||
.map_err(|e| io::Error::new(e.kind(), format!("m2ts://{}: {}", path.display(), e)))?;
|
||||
let reader = std::io::BufReader::with_capacity(IO_BUF_SIZE, file);
|
||||
Ok(Box::new(M2tsStream::open(reader)?))
|
||||
}
|
||||
StreamUrl::Mkv { ref path } => {
|
||||
validate_file_path(path, "mkv")?;
|
||||
// TODO: MkvStream InputStream (MKV demux → PES)
|
||||
// MKV as PES input requires EBML → PES frame extraction (TODO)
|
||||
Err(io::Error::new(io::ErrorKind::Unsupported,
|
||||
"mkv:// PES input not yet implemented"))
|
||||
"mkv:// as input not yet supported — use m2ts:// or iso:// as source"))
|
||||
}
|
||||
StreamUrl::Network { ref addr } => {
|
||||
validate_network_addr(addr)?;
|
||||
|
||||
Reference in New Issue
Block a user