Skip to main content

wallr_core/video/
playback.rs

1use crate::video::{
2    DecoderInfo, FrameScheduler, HwAccel, VideoDecoder, VideoError, VideoFrame, VideoMetadata,
3};
4use std::path::Path;
5use std::sync::Mutex;
6use std::sync::atomic::{AtomicU64, Ordering};
7use std::time::{Duration, Instant};
8
9pub struct VideoPlayback {
10    decoder: Mutex<Option<VideoDecoder>>,
11    scheduler: Mutex<Option<FrameScheduler>>,
12    pending_frame: Mutex<Option<VideoFrame>>,
13    generation: AtomicU64,
14}
15
16pub struct PreparedVideoPlayback {
17    decoder: VideoDecoder,
18    metadata: VideoMetadata,
19    first_frame: Option<VideoFrame>,
20}
21
22impl PreparedVideoPlayback {
23    pub fn metadata(&self) -> &VideoMetadata {
24        &self.metadata
25    }
26
27    pub fn take_first_frame(&mut self) -> Option<VideoFrame> {
28        self.first_frame.take()
29    }
30}
31
32impl VideoPlayback {
33    pub fn new() -> Self {
34        Self {
35            decoder: Mutex::new(None),
36            scheduler: Mutex::new(None),
37            pending_frame: Mutex::new(None),
38            generation: AtomicU64::new(u64::MAX),
39        }
40    }
41
42    pub fn start(
43        &self,
44        path: &Path,
45        hw_accel: HwAccel,
46        preload_frames: usize,
47        generation: u64,
48    ) -> Result<VideoMetadata, crate::video::error::VideoError> {
49        let prepared = Self::prepare(
50            path,
51            hw_accel,
52            preload_frames,
53            Duration::from_millis(0),
54            |_| Ok(()),
55        )?;
56        Ok(self.commit(prepared, generation))
57    }
58
59    pub fn prepare<F>(
60        path: &Path,
61        hw_accel: HwAccel,
62        preload_frames: usize,
63        first_frame_timeout: Duration,
64        validate: F,
65    ) -> Result<PreparedVideoPlayback, crate::video::error::VideoError>
66    where
67        F: FnOnce(&VideoMetadata) -> Result<(), crate::video::error::VideoError>,
68    {
69        let decoder = VideoDecoder::with_preload(path, hw_accel, preload_frames)?;
70        let metadata = decoder.metadata().clone();
71        validate(&metadata)?;
72        let deadline = Instant::now() + first_frame_timeout;
73        let first_frame = loop {
74            if let Some(frame) = decoder.next_frame() {
75                break Some(frame);
76            }
77            if Instant::now() >= deadline {
78                tracing::warn!("Timed out waiting for first video frame");
79                break None;
80            }
81            std::thread::sleep(Duration::from_millis(5));
82        };
83
84        Ok(PreparedVideoPlayback {
85            decoder,
86            metadata,
87            first_frame,
88        })
89    }
90
91    pub fn commit(&self, prepared: PreparedVideoPlayback, generation: u64) -> VideoMetadata {
92        let PreparedVideoPlayback {
93            decoder,
94            metadata,
95            first_frame,
96        } = prepared;
97        let scheduler = FrameScheduler::new(metadata.duration);
98        self.stop();
99        self.generation.store(generation, Ordering::Release);
100        *self.lock_decoder() = Some(decoder);
101        *self.lock_scheduler() = Some(scheduler);
102        *self.lock_pending() = first_frame;
103        tracing::info!(
104            "Video started: {}x{} @ {:.2} fps, {:?}",
105            metadata.width,
106            metadata.height,
107            metadata.fps,
108            metadata.duration
109        );
110        metadata
111    }
112
113    pub fn stop(&self) {
114        self.generation.store(u64::MAX, Ordering::Release);
115        *self.lock_decoder() = None;
116        *self.lock_scheduler() = None;
117        *self.lock_pending() = None;
118    }
119
120    /// Stops playback only when `generation` is still active.
121    /// Superseded callers cannot stop their successor's playback.
122    pub fn stop_generation(&self, generation: u64) {
123        let mut decoder = self.lock_decoder();
124        let mut scheduler = self.lock_scheduler();
125        let mut pending = self.lock_pending();
126        if self.generation.load(Ordering::Acquire) != generation {
127            return;
128        }
129        self.generation.store(u64::MAX, Ordering::Release);
130        *decoder = None;
131        *scheduler = None;
132        *pending = None;
133    }
134
135    pub fn pause(&self) {
136        if let Some(s) = self.lock_scheduler().as_mut() {
137            s.pause();
138        }
139        if let Some(d) = self.lock_decoder().as_ref() {
140            d.pause();
141        }
142    }
143
144    pub fn resume(&self) {
145        if let Some(s) = self.lock_scheduler().as_mut() {
146            s.resume();
147        }
148        if let Some(d) = self.lock_decoder().as_ref() {
149            d.resume();
150        }
151    }
152
153    pub fn seek(&self, timestamp: Duration) -> Result<(), crate::video::error::VideoError> {
154        let mut decoder = self.lock_decoder();
155        let mut scheduler = self.lock_scheduler();
156        let mut pending = self.lock_pending();
157        let scheduler = scheduler.as_mut().ok_or_else(|| {
158            VideoError::SeekFailed(timestamp, anyhow::anyhow!("no video is playing"))
159        })?;
160        scheduler.seek(timestamp)?;
161        if let Some(d) = decoder.as_mut() {
162            d.seek(timestamp);
163            d.drain();
164        }
165        *pending = None;
166        tracing::info!("Video seeked to {:?}", timestamp);
167        Ok(())
168    }
169
170    pub fn next_frame(&self) -> Option<VideoFrame> {
171        self.next_frame_for_generation(None)
172    }
173
174    /// Returns a frame only when `generation` is still active.
175    pub fn next_frame_in_generation(&self, generation: u64) -> Option<VideoFrame> {
176        self.next_frame_for_generation(Some(generation))
177    }
178
179    fn next_frame_for_generation(&self, generation: Option<u64>) -> Option<VideoFrame> {
180        let mut decoder = self.lock_decoder();
181        if generation.is_some_and(|expected| self.generation.load(Ordering::Acquire) != expected) {
182            return None;
183        }
184        let mut scheduler = self.lock_scheduler();
185        let mut pending = self.lock_pending();
186        let (Some(decoder), Some(scheduler)) = (decoder.as_mut(), scheduler.as_mut()) else {
187            return None;
188        };
189
190        take_due_frame(scheduler, &mut pending, || decoder.next_frame())
191    }
192
193    pub fn wait_first_frame(&self, timeout: Duration) -> Option<VideoFrame> {
194        let deadline = Instant::now() + timeout;
195        loop {
196            if let Some(frame) = self.next_frame() {
197                return Some(frame);
198            }
199            if Instant::now() >= deadline {
200                tracing::warn!("Timed out waiting for first video frame");
201                return None;
202            }
203            std::thread::sleep(Duration::from_millis(5));
204        }
205    }
206
207    pub fn time_until_next_frame(&self) -> Option<Duration> {
208        self.time_until_next_frame_for_generation(None)
209    }
210
211    /// Returns a deadline only when `generation` is still active.
212    pub fn time_until_next_frame_in_generation(&self, generation: u64) -> Option<Duration> {
213        self.time_until_next_frame_for_generation(Some(generation))
214    }
215
216    fn time_until_next_frame_for_generation(&self, generation: Option<u64>) -> Option<Duration> {
217        let scheduler = self.lock_scheduler();
218        let pending = self.lock_pending();
219        if generation.is_some_and(|expected| self.generation.load(Ordering::Acquire) != expected) {
220            return None;
221        }
222        let (Some(scheduler), Some(frame)) = (scheduler.as_ref(), pending.as_ref()) else {
223            return None;
224        };
225        scheduler.time_until_next_frame(frame.pts)
226    }
227
228    pub fn metadata(&self) -> Option<VideoMetadata> {
229        self.lock_decoder().as_ref().map(|d| d.metadata().clone())
230    }
231
232    pub fn decoder_info(&self) -> Option<DecoderInfo> {
233        self.lock_decoder().as_ref().map(|d| d.decoder_info())
234    }
235
236    pub fn hw_accel_in_use(&self) -> HwAccel {
237        self.lock_decoder()
238            .as_ref()
239            .map(VideoDecoder::hw_accel_in_use)
240            .unwrap_or(HwAccel::Software)
241    }
242
243    pub fn position(&self) -> Option<Duration> {
244        self.lock_scheduler().as_ref().map(|s| s.current_position())
245    }
246
247    pub fn is_paused(&self) -> bool {
248        self.lock_scheduler()
249            .as_ref()
250            .map(|s| s.is_paused())
251            .unwrap_or(false)
252    }
253
254    pub fn is_playing(&self) -> bool {
255        self.lock_decoder().is_some()
256    }
257
258    fn lock_decoder(&self) -> std::sync::MutexGuard<'_, Option<VideoDecoder>> {
259        self.decoder.lock().unwrap_or_else(|p| p.into_inner())
260    }
261
262    fn lock_scheduler(&self) -> std::sync::MutexGuard<'_, Option<FrameScheduler>> {
263        self.scheduler.lock().unwrap_or_else(|p| p.into_inner())
264    }
265
266    fn lock_pending(&self) -> std::sync::MutexGuard<'_, Option<VideoFrame>> {
267        self.pending_frame.lock().unwrap_or_else(|p| p.into_inner())
268    }
269}
270
271fn take_due_frame(
272    scheduler: &mut FrameScheduler,
273    pending: &mut Option<VideoFrame>,
274    mut next_frame: impl FnMut() -> Option<VideoFrame>,
275) -> Option<VideoFrame> {
276    let mut selected = None;
277
278    if let Some(frame) = pending.take() {
279        if !scheduler.should_display(frame.pts) {
280            *pending = Some(frame);
281            return None;
282        }
283        if scheduler.should_upload(frame.pts) {
284            selected = Some(frame);
285        }
286    }
287
288    while let Some(frame) = next_frame() {
289        if !scheduler.should_display(frame.pts) {
290            *pending = Some(frame);
291            break;
292        }
293        if scheduler.should_upload(frame.pts) {
294            selected = Some(frame);
295        }
296    }
297
298    selected
299}
300
301impl Default for VideoPlayback {
302    fn default() -> Self {
303        Self::new()
304    }
305}
306
307#[cfg(test)]
308mod tests {
309    use super::*;
310    use std::collections::VecDeque;
311
312    fn frame(pts_ms: u64) -> VideoFrame {
313        VideoFrame {
314            data: crate::video::VideoFrameData::Rgba(vec![pts_ms as u8]),
315            width: 1,
316            height: 1,
317            pts: Duration::from_millis(pts_ms),
318            index: pts_ms,
319        }
320    }
321
322    #[test]
323    fn retains_future_frame_until_due() {
324        let mut scheduler = FrameScheduler::new(Duration::from_secs(1));
325        let mut pending = None;
326        let mut frames = VecDeque::from([frame(0), frame(100)]);
327
328        let first = take_due_frame(&mut scheduler, &mut pending, || frames.pop_front()).unwrap();
329        assert_eq!(first.pts, Duration::ZERO);
330        assert_eq!(pending.as_ref().unwrap().pts, Duration::from_millis(100));
331
332        scheduler.seek(Duration::from_millis(110)).unwrap();
333        let second = take_due_frame(&mut scheduler, &mut pending, || None).unwrap();
334        assert_eq!(second.pts, Duration::from_millis(100));
335        assert!(pending.is_none());
336    }
337
338    #[test]
339    fn selects_latest_due_frame() {
340        let mut scheduler = FrameScheduler::new(Duration::from_secs(1));
341        scheduler.seek(Duration::from_millis(50)).unwrap();
342        let mut pending = None;
343        let mut frames = VecDeque::from([frame(0), frame(16), frame(32), frame(100)]);
344
345        let selected = take_due_frame(&mut scheduler, &mut pending, || frames.pop_front()).unwrap();
346        assert_eq!(selected.pts, Duration::from_millis(32));
347        assert_eq!(pending.as_ref().unwrap().pts, Duration::from_millis(100));
348    }
349}