use std::io::{Read, Seek, SeekFrom};
use std::ops::Range;
use itertools::Either;
use re_chunk::{Chunk, ChunkId, EntityPath, RowId, TimeColumn, TimePoint};
use re_log_types::{TimeType, Timeline, TimelineName};
use re_sdk_types::archetypes::VideoStream;
use re_sdk_types::components::VideoCodec;
use re_video::player::GetVideoSource;
use re_video::{
Mp4TranscodeOptions, SampleIndex, SampleMetadataState, VideoDataDescription, VideoSource,
};
use crate::Mp4Error;
trait ReadSeek: Read + Seek {}
impl<T: Read + Seek> ReadSeek for T {}
type ChunkIter = Box<dyn Iterator<Item = Result<Chunk, Mp4Error>>>;
pub(crate) enum StreamInput {
#[cfg(not(target_arch = "wasm32"))]
Path(std::path::PathBuf),
Bytes(Vec<u8>),
}
impl StreamInput {
#[cfg_attr(target_arch = "wasm32", expect(clippy::unnecessary_wraps))]
fn open(&self) -> Result<Box<dyn ReadSeek>, Mp4Error> {
Ok(match self {
#[cfg(not(target_arch = "wasm32"))]
Self::Path(path) => Box::new(std::io::BufReader::new(std::fs::File::open(path)?)),
Self::Bytes(bytes) => Box::new(std::io::Cursor::new(bytes.clone())),
})
}
}
pub(crate) fn iter_chunks(
input: StreamInput,
entity_path: &EntityPath,
timeline_name: TimelineName,
chunk_by_gop: bool,
timeline_type: TimeType,
transcode: &Mp4TranscodeOptions,
debug_name: &str,
) -> Result<ChunkIter, Mp4Error> {
re_tracing::profile_function!();
let mut reader = input.open()?;
let size = reader.seek(SeekFrom::End(0))?;
reader.seek(SeekFrom::Start(0))?;
let desc = VideoDataDescription::load_mp4_from_reader(&mut reader, size, debug_name)?;
VideoCodec::try_from(desc.codec.clone()).map_err(|_err| Mp4Error::ImageSequenceInStreamMode)?;
let target_codec = transcode
.output_codec
.clone()
.unwrap_or_else(|| desc.codec.clone());
let output_mapped_codec =
VideoCodec::try_from(target_codec).map_err(|_err| Mp4Error::ImageSequenceInStreamMode)?;
let needs_decoder_reordering = !desc.samples_statistics.dts_always_equal_pts;
let wants_transcode = needs_decoder_reordering || transcode.requests_transform(&desc.codec);
let segments: Box<dyn Iterator<Item = Result<Segment, Mp4Error>>> = if !wants_transcode {
Box::new(std::iter::once(Segment::new(reader, desc)))
} else {
drop(reader);
#[cfg(not(target_arch = "wasm32"))]
{
Box::new(transcoded_segments(
input,
desc.codec.clone(),
transcode,
debug_name,
)?)
}
#[cfg(target_arch = "wasm32")]
{
let _ = (input, transcode);
return Err(Mp4Error::TranscodeRequiresFfmpeg);
}
};
let entity_path = entity_path.clone();
let codec_chunk = build_codec_chunk(&entity_path, output_mapped_codec);
let gop_chunks = segments.flat_map(move |segment| {
gop_chunks(
segment,
entity_path.clone(),
timeline_name,
timeline_type,
chunk_by_gop,
)
});
Ok(Box::new(std::iter::chain(
std::iter::once(codec_chunk),
gop_chunks,
)))
}
struct Segment {
reader: Box<dyn ReadSeek>,
desc: VideoDataDescription,
timescale: re_video::Timescale,
}
impl Segment {
fn new(reader: Box<dyn ReadSeek>, desc: VideoDataDescription) -> Result<Self, Mp4Error> {
let timescale = desc.timescale.ok_or(Mp4Error::NoTimescale)?;
if !desc.samples.is_empty()
&& desc.keyframe_indices.first() != Some(&desc.samples.min_index())
{
return Err(Mp4Error::SamplesBeforeFirstKeyframe);
}
Ok(Self {
reader,
desc,
timescale,
})
}
}
fn gop_chunks(
segment: Result<Segment, Mp4Error>,
entity_path: EntityPath,
timeline_name: TimelineName,
timeline_type: TimeType,
chunk_by_gop: bool,
) -> impl Iterator<Item = Result<Chunk, Mp4Error>> {
let Segment {
mut reader,
desc,
timescale,
} = match segment {
Ok(segment) => segment,
Err(err) => return Either::Left(std::iter::once(Err(err))),
};
let ranges = sample_ranges(&desc, chunk_by_gop, &entity_path);
Either::Right(ranges.into_iter().filter_map(move |range| {
build_gop_chunk(
&mut *reader,
&desc,
timescale,
timeline_name,
timeline_type,
&entity_path,
range,
)
.transpose()
}))
}
#[cfg(not(target_arch = "wasm32"))]
fn transcoded_segments(
input: StreamInput,
source_codec: re_video::VideoCodec,
transcode: &Mp4TranscodeOptions,
debug_name: &str,
) -> Result<impl Iterator<Item = Result<Segment, Mp4Error>> + use<>, Mp4Error> {
let StreamInput::Path(path) = input else {
return Err(Mp4Error::TranscodeRequiresSeekableFile);
};
let chunks = re_video::transcode_mp4(&path, source_codec, transcode, debug_name)
.map_err(|err| map_ffmpeg_err(&err))?;
let scanner = FragmentScanner::new(ChunkReader::new(chunks), debug_name)?;
let debug_name = debug_name.to_owned();
Ok(scanner.map(move |mini_mp4| {
let mini_mp4 = mini_mp4?;
let size = mini_mp4.len() as u64;
let mut reader = std::io::Cursor::new(mini_mp4);
let desc = VideoDataDescription::load_mp4_from_reader(&mut reader, size, &debug_name)?;
Segment::new(Box::new(reader), desc)
}))
}
#[cfg(not(target_arch = "wasm32"))]
fn map_ffmpeg_err(err: &re_video::FFmpegError) -> Mp4Error {
if let re_video::FFmpegError::NoEncoderForCodec { codec } = err {
return Mp4Error::NoEncoderAvailable {
codec: codec.clone(),
};
}
let mut msg = err.to_string();
if matches!(err, re_video::FFmpegError::FFmpegNotInstalled)
&& let Some(url) = re_video::ffmpeg_download_url()
{
msg = format!("{msg} You can download a build of FFmpeg at {url}");
}
Mp4Error::Transcode(msg)
}
#[cfg(not(target_arch = "wasm32"))]
struct ChunkReader<I> {
chunks: I,
current: Vec<u8>,
pos: usize,
}
#[cfg(not(target_arch = "wasm32"))]
impl<I> ChunkReader<I> {
fn new(chunks: I) -> Self {
Self {
chunks,
current: Vec::new(),
pos: 0,
}
}
}
#[cfg(not(target_arch = "wasm32"))]
impl<I: Iterator<Item = Result<Vec<u8>, re_video::FFmpegError>>> Read for ChunkReader<I> {
fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
while self.pos >= self.current.len() {
match self.chunks.next() {
Some(Ok(chunk)) => {
self.current = chunk;
self.pos = 0;
}
Some(Err(err)) => return Err(std::io::Error::other(err.to_string())),
None => return Ok(0),
}
}
let n = (self.current.len() - self.pos).min(buf.len());
buf[..n].copy_from_slice(&self.current[self.pos..self.pos + n]);
self.pos += n;
Ok(n)
}
}
#[cfg(not(target_arch = "wasm32"))]
type Mp4Box = ([u8; 4], Vec<u8>);
#[cfg(not(target_arch = "wasm32"))]
struct FragmentScanner<R> {
reader: R,
init: Vec<u8>,
pending_moof: Option<Vec<u8>>,
done: bool,
}
#[cfg(not(target_arch = "wasm32"))]
impl<R: Read> FragmentScanner<R> {
fn new(mut reader: R, debug_name: &str) -> Result<Self, Mp4Error> {
let mut init = Vec::new();
let mut pending_moof = None;
loop {
match read_box(&mut reader)? {
None => break, Some((box_type, bytes)) => {
if &box_type == b"moof" {
pending_moof = Some(bytes);
break;
}
init.extend_from_slice(&bytes);
}
}
}
if init.is_empty() {
return Err(Mp4Error::Transcode(format!(
"ffmpeg produced no init segment for {debug_name}"
)));
}
Ok(Self {
reader,
init,
pending_moof,
done: false,
})
}
fn next_fragment(&mut self) -> Result<Option<Vec<u8>>, Mp4Error> {
if self.done {
return Ok(None);
}
let Some(mut fragment) = self.pending_moof.take() else {
self.done = true;
return Ok(None);
};
match read_box(&mut self.reader)? {
Some((box_type, mdat)) if &box_type == b"mdat" => fragment.extend_from_slice(&mdat),
_ => {
self.done = true;
return Err(Mp4Error::Transcode(
"ffmpeg fragmented mp4 has a `moof` without a following `mdat`".to_owned(),
));
}
}
match read_box(&mut self.reader)? {
Some((box_type, bytes)) if &box_type == b"moof" => self.pending_moof = Some(bytes),
_ => self.done = true,
}
Ok(Some(fragment))
}
}
#[cfg(not(target_arch = "wasm32"))]
impl<R: Read> Iterator for FragmentScanner<R> {
type Item = Result<Vec<u8>, Mp4Error>;
fn next(&mut self) -> Option<Self::Item> {
match self.next_fragment() {
Ok(Some(fragment)) => {
let mut mini_mp4 = Vec::with_capacity(self.init.len() + fragment.len());
mini_mp4.extend_from_slice(&self.init);
mini_mp4.extend_from_slice(&fragment);
Some(Ok(mini_mp4))
}
Ok(None) => None,
Err(err) => Some(Err(err)),
}
}
}
#[cfg(not(target_arch = "wasm32"))]
fn read_box<R: Read>(reader: &mut R) -> Result<Option<Mp4Box>, Mp4Error> {
let mut header = [0u8; 8];
if !read_exact_or_eof(reader, &mut header)? {
return Ok(None);
}
let size32 = u32::from_be_bytes([header[0], header[1], header[2], header[3]]);
let mut box_type = [0u8; 4];
box_type.copy_from_slice(&header[4..8]);
let mut bytes = Vec::new();
bytes.extend_from_slice(&header);
let total = if size32 == 1 {
let mut ext = [0u8; 8];
reader.read_exact(&mut ext)?;
bytes.extend_from_slice(&ext);
u64::from_be_bytes(ext) as usize
} else if size32 == 0 {
reader.read_to_end(&mut bytes)?;
return Ok(Some((box_type, bytes)));
} else {
size32 as usize
};
if total < bytes.len() {
return Err(Mp4Error::Transcode(format!(
"ffmpeg produced an mp4 box with an invalid size {total}"
)));
}
let already = bytes.len();
bytes.resize(total, 0);
reader.read_exact(&mut bytes[already..])?;
Ok(Some((box_type, bytes)))
}
#[cfg(not(target_arch = "wasm32"))]
fn read_exact_or_eof<R: Read>(reader: &mut R, buf: &mut [u8]) -> Result<bool, Mp4Error> {
let mut filled = 0;
while filled < buf.len() {
match reader.read(&mut buf[filled..])? {
0 => {
if filled == 0 {
return Ok(false);
}
return Err(Mp4Error::Transcode(
"ffmpeg output ended in the middle of an mp4 box".to_owned(),
));
}
n => filled += n,
}
}
Ok(true)
}
fn sample_ranges(
desc: &VideoDataDescription,
chunk_by_gop: bool,
entity_path: &EntityPath,
) -> Vec<Range<SampleIndex>> {
if chunk_by_gop {
(0..desc.keyframe_indices.len())
.filter_map(|i| desc.gop_sample_range_for_keyframe(i))
.collect()
} else {
desc.samples
.iter_indexed()
.filter_map(|(idx, sample)| match sample {
SampleMetadataState::Present(_) => Some(idx..idx + 1),
SampleMetadataState::Unloaded { .. } => {
re_log::warn_once!(
"Skipping unloaded sample {idx} in mp4 demux (entity path: {entity_path})"
);
None
}
})
.collect()
}
}
fn build_codec_chunk(entity_path: &EntityPath, codec: VideoCodec) -> Result<Chunk, Mp4Error> {
let chunk = Chunk::builder(entity_path.clone())
.with_archetype(
RowId::new(),
TimePoint::default(),
&VideoStream::update_fields().with_codec(codec),
)
.build()?;
Ok(chunk)
}
fn build_gop_chunk(
reader: &mut dyn ReadSeek,
desc: &VideoDataDescription,
timescale: re_video::Timescale,
timeline_name: TimelineName,
timeline_type: TimeType,
entity_path: &EntityPath,
range: Range<SampleIndex>,
) -> Result<Option<Chunk>, Mp4Error> {
let mut time_values: Vec<i64> = Vec::with_capacity(range.len());
let mut sample_blobs: Vec<Vec<u8>> = Vec::with_capacity(range.len());
let mut is_keyframe: Vec<bool> = Vec::with_capacity(range.len());
let mut sample_bytes = vec![];
for sample_idx in range {
let SampleMetadataState::Present(meta) = &desc.samples[sample_idx] else {
re_log::warn_once!(
"Skipping unloaded sample {sample_idx} in mp4 demux (entity path: {entity_path})"
);
continue;
};
let pts_ns = meta.presentation_timestamp.into_nanos(timescale);
time_values.push(pts_ns);
let VideoSource::Span(span) = meta.source else {
return Err(Mp4Error::SampleConversion(format!(
"sample {sample_idx} has a non-span source; mp4 demux only produces spans"
)));
};
struct FullSource<'a>(&'a [u8]);
impl GetVideoSource for FullSource<'_> {
fn get_video_chunk(&self, _source: VideoSource) -> &[u8] {
self.0
}
fn require_video_source(&self, _source: VideoSource) {}
fn indicate_video_source(&self, _source: VideoSource) {}
}
let byte_range = span.range_usize();
reader.seek(SeekFrom::Start(byte_range.start as u64))?;
sample_bytes.resize(byte_range.len(), 0);
reader.read_exact(&mut sample_bytes)?;
let chunk = meta
.get(&FullSource(sample_bytes.as_slice()), sample_idx)
.ok_or_else(|| {
Mp4Error::SampleConversion(format!(
"sample {sample_idx} could not be read from the mp4 buffer"
))
})?;
sample_blobs.push(
desc.sample_data_in_stream_format(&chunk)
.map_err(|err| Mp4Error::SampleConversion(err.to_string()))?,
);
is_keyframe.push(meta.is_sync);
}
if time_values.is_empty() {
return Ok(None);
}
let timeline = Timeline::new(timeline_name, timeline_type);
let time_column = TimeColumn::new(
Some(true),
timeline,
arrow::buffer::ScalarBuffer::from(time_values),
);
let components: Vec<_> = VideoStream::update_fields()
.with_many_sample(sample_blobs)
.with_many_is_keyframe(is_keyframe)
.columns_of_unit_batches()
.map_err(|err| Mp4Error::SampleConversion(err.to_string()))?
.collect();
let chunk = Chunk::from_auto_row_ids(
ChunkId::new(),
entity_path.clone(),
std::iter::once((*timeline.name(), time_column)).collect(),
components.into_iter().collect(),
)?;
Ok(Some(chunk))
}
#[cfg(all(test, not(target_arch = "wasm32")))]
mod tests {
use super::FragmentScanner;
fn mp4_box(box_type: &[u8; 4], body: &[u8]) -> Vec<u8> {
let mut bytes = ((8 + body.len()) as u32).to_be_bytes().to_vec();
bytes.extend_from_slice(box_type);
bytes.extend_from_slice(body);
bytes
}
#[test]
fn yields_one_init_plus_fragment_mini_mp4_per_gop() {
let ftyp = mp4_box(b"ftyp", b"isomiso2");
let moov = mp4_box(b"moov", b"fake-init-metadata");
let init: Vec<u8> = [ftyp, moov].concat();
let fragments: Vec<Vec<u8>> = (0..3u8)
.map(|i| [mp4_box(b"moof", &[i; 4]), mp4_box(b"mdat", &[0xAB; 16])].concat())
.collect();
let mut stream = init.clone();
for fragment in &fragments {
stream.extend_from_slice(fragment);
}
stream.extend_from_slice(&mp4_box(b"mfra", b"index"));
let scanner = FragmentScanner::new(std::io::Cursor::new(stream), "test").unwrap();
let got: Vec<Vec<u8>> = scanner.map(Result::unwrap).collect();
let expected: Vec<Vec<u8>> = fragments
.iter()
.map(|fragment| [init.clone(), fragment.clone()].concat())
.collect();
assert_eq!(
got, expected,
"one init+fragment mini-mp4 per moof+mdat pair"
);
}
#[test]
fn init_only_stream_yields_no_fragments() {
let stream = [mp4_box(b"ftyp", b"isom"), mp4_box(b"moov", b"meta")].concat();
let mut scanner = FragmentScanner::new(std::io::Cursor::new(stream), "test").unwrap();
assert!(scanner.next().is_none());
}
}