Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 27 additions & 3 deletions crates/video-streamer/src/normalizer/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -363,12 +363,19 @@ impl ClipNormalizer {
}

fn caught_up(&mut self) -> anyhow::Result<()> {
let mut decode_frame = Self::decode_frame;
let mut process_frame = Self::process_frame;
self.caught_up_with(&mut process_frame)
self.caught_up_with(&mut decode_frame, &mut process_frame)
}

fn caught_up_with<F>(&mut self, process_frame: &mut F) -> anyhow::Result<()>
/// Ends the history phase.
///
/// For a live-edge viewer, the frames of the latest group of pictures are only decoded, to rebuild the current
/// picture, and only the last one is processed. Its picture becomes the first output key frame, so the viewer
/// starts at the live picture instead of replaying the group.
fn caught_up_with<D, F>(&mut self, decode_frame: &mut D, process_frame: &mut F) -> anyhow::Result<()>
where
D: FnMut(&mut Self, PendingFrame) -> anyhow::Result<()>,
F: FnMut(&mut Self, PendingFrame) -> anyhow::Result<()>,
{
let history = match std::mem::replace(&mut self.phase, ClipPhase::Live) {
Expand All @@ -378,7 +385,15 @@ impl ClipNormalizer {
if matches!(history, HistoryPolicy::KeepLatestGop)
&& let Some(replay_point) = self.replay_point
{
self.replay_latest_gop_with(replay_point, process_frame)?;
let mut latest = None;
let mut keep_latest = |this: &mut Self, frame: PendingFrame| match latest.replace(frame) {
Some(previous) => decode_frame(this, previous),
None => Ok(()),
};
self.replay_latest_gop_with(replay_point, &mut keep_latest)?;
if let Some(latest) = latest {
process_frame(self, latest)?;
}
}
Ok(())
}
Expand Down Expand Up @@ -571,6 +586,15 @@ impl ClipNormalizer {
}))
}

/// Decodes a frame without showing it, to advance the decoder to a later picture.
fn decode_frame(&mut self, frame: PendingFrame) -> anyhow::Result<()> {
let input_decoder = self
.input_decoder
.get_or_insert_with(|| InputDecoder::new(frame.codec, self.config.encoder_threads));
input_decoder.decode(&frame.data)?;
Ok(())
}

fn process_frame(&mut self, frame: PendingFrame) -> anyhow::Result<()> {
let should_skip_encode = self.should_skip_encode(frame.timestamp);
let processing_started = Instant::now();
Expand Down
38 changes: 29 additions & 9 deletions crates/video-streamer/src/normalizer/tests/replay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,17 +23,24 @@ fn replay_uses_latest_gop_and_restores_the_same_reader() {
normalizer.scan_available().expect("scan available history");
let original_head = normalizer.reader_head;
let replay_point = normalizer.replay_point.expect("latest replay point");
let mut decoded = Vec::new();
let mut frames = Vec::new();
let mut decode_frame = |_: &mut ClipNormalizer, frame: PendingFrame| {
decoded.push((frame.timestamp, frame.key_frame));
Ok(())
};
let mut process_frame = |_: &mut ClipNormalizer, frame: PendingFrame| {
frames.push((frame.timestamp, frame.key_frame));
Ok(())
};

normalizer
.caught_up_with(&mut process_frame)
.caught_up_with(&mut decode_frame, &mut process_frame)
.expect("replay latest GOP");

assert_eq!(frames, vec![(90, true), (120, false)]);
// The viewer starts at the live picture: earlier frames of the group are decoded only.
assert_eq!(decoded, vec![(90, true)]);
assert_eq!(frames, vec![(120, false)]);
assert_eq!(normalizer.reader_head, original_head);

let stats = stats.lock().expect("reader stats lock");
Expand All @@ -44,7 +51,7 @@ fn replay_uses_latest_gop_and_restores_the_same_reader() {

#[test]
fn replay_error_restores_the_original_parser_and_reader_state() {
let data = video_clip_bytes(&[true, false, true], false);
let data = video_clip_bytes(&[true, false, true, false], false);
let visible = Arc::new(AtomicUsize::new(data.len()));
let stats = Arc::new(Mutex::new(ReaderStats::default()));
let reader = GrowingReader {
Expand All @@ -59,10 +66,11 @@ fn replay_error_restores_the_original_parser_and_reader_state() {
let original_head = normalizer.reader_head;
let original_decoder_position = normalizer.decoder.position();
let original_input = normalizer.input.clone();
let mut process_frame = |_: &mut ClipNormalizer, _: PendingFrame| Err(anyhow::anyhow!("replay callback failed"));
let mut decode_frame = |_: &mut ClipNormalizer, _: PendingFrame| Err(anyhow::anyhow!("replay callback failed"));
let mut process_frame = |_: &mut ClipNormalizer, _: PendingFrame| Ok(());

let error = normalizer
.caught_up_with(&mut process_frame)
.caught_up_with(&mut decode_frame, &mut process_frame)
.expect_err("replay callback failure");

assert!(format!("{error:#}").contains("replay callback failed"));
Expand Down Expand Up @@ -91,12 +99,15 @@ fn known_size_group_replay_excludes_a_partial_next_group_until_growth() {
normalizer.scan_available().expect("scan partial group");
assert!(normalizer.pending_block_group.is_some());
{
let mut decode_frame = |_: &mut ClipNormalizer, _: PendingFrame| -> anyhow::Result<()> {
panic!("a single-frame group has nothing to decode only")
};
let mut process_frame = |_: &mut ClipNormalizer, frame: PendingFrame| {
frames.push((frame.timestamp, frame.key_frame));
Ok(())
};
normalizer
.caught_up_with(&mut process_frame)
.caught_up_with(&mut decode_frame, &mut process_frame)
.expect("replay complete group");
}
assert_eq!(frames, vec![(0, true)]);
Expand Down Expand Up @@ -129,12 +140,15 @@ fn unknown_size_group_replay_uses_the_following_sibling_as_its_boundary() {
normalizer.scan_available().expect("scan unknown-sized group");
assert!(normalizer.pending_block_group.is_some());
{
let mut decode_frame = |_: &mut ClipNormalizer, _: PendingFrame| -> anyhow::Result<()> {
panic!("a single-frame group has nothing to decode only")
};
let mut process_frame = |_: &mut ClipNormalizer, frame: PendingFrame| {
frames.push((frame.timestamp, frame.key_frame));
Ok(())
};
normalizer
.caught_up_with(&mut process_frame)
.caught_up_with(&mut decode_frame, &mut process_frame)
.expect("replay through the complete boundary");
}
assert_eq!(frames, vec![(0, true)]);
Expand All @@ -161,18 +175,24 @@ fn replay_stops_at_a_completed_frame_inside_a_known_cluster() {
stats: Arc::new(Mutex::new(ReaderStats::default())),
};
let (mut normalizer, _receiver) = live_edge_normalizer(reader);
let mut decoded = Vec::new();
let mut frames = Vec::new();
let mut decode_frame = |_: &mut ClipNormalizer, frame: PendingFrame| {
decoded.push((frame.timestamp, frame.key_frame));
Ok(())
};
let mut process_frame = |_: &mut ClipNormalizer, frame: PendingFrame| {
frames.push((frame.timestamp, frame.key_frame));
Ok(())
};

normalizer.scan_available().expect("scan cross-cluster history");
normalizer
.caught_up_with(&mut process_frame)
.caught_up_with(&mut decode_frame, &mut process_frame)
.expect("replay completed frames without finalizing the document");

assert_eq!(frames, vec![(0, true), (30, false)]);
assert_eq!(decoded, vec![(0, true)]);
assert_eq!(frames, vec![(30, false)]);
}

#[test]
Expand Down
Loading
Loading