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 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 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 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}