(
&self,
timeout_duration: Duration,
)
| 1346 | } |
| 1347 | |
| 1348 | pub(crate) async fn wait_for_keyframe( |
| 1349 | &self, |
| 1350 | timeout_duration: Duration, |
| 1351 | ) -> Option<SharedFrame> { |
| 1352 | let deadline = Instant::now() + timeout_duration; |
| 1353 | let baseline_sequence = self |
| 1354 | .inner |
| 1355 | .latest_keyframe |
| 1356 | .read() |
| 1357 | .unwrap() |
| 1358 | .as_ref() |
| 1359 | .map_or(0, |frame| frame.frame_sequence); |
| 1360 | let mut receiver = self.inner.sender.subscribe(); |
| 1361 | self.request_keyframe(); |
| 1362 | |
| 1363 | loop { |
| 1364 | if let Some(frame) = self.inner.latest_keyframe.read().unwrap().clone() { |
| 1365 | if frame.frame_sequence > baseline_sequence { |
| 1366 | return Some(frame); |
| 1367 | } |
| 1368 | } |
| 1369 | let remaining = deadline.checked_duration_since(Instant::now())?; |
| 1370 | match time::timeout(remaining, receiver.recv()).await { |
| 1371 | Ok(Ok(frame)) if frame.is_keyframe && frame.frame_sequence > baseline_sequence => { |
| 1372 | return Some(frame) |
| 1373 | } |
| 1374 | Ok(Ok(_)) | Ok(Err(broadcast::error::RecvError::Lagged(_))) => { |
| 1375 | self.request_keyframe(); |
| 1376 | } |
| 1377 | Ok(Err(_)) | Err(_) => return None, |
| 1378 | } |
| 1379 | } |
| 1380 | } |
| 1381 | |
| 1382 | pub(crate) fn request_refresh(&self) {} |
| 1383 |
no test coverage detected