(&self, timeout_duration: Duration)
| 198 | } |
| 199 | |
| 200 | pub async fn wait_for_keyframe(&self, timeout_duration: Duration) -> Option<SharedFrame> { |
| 201 | self.inner.native.set_client_foreground(true); |
| 202 | let deadline = Instant::now() + timeout_duration; |
| 203 | let baseline_sequence = self |
| 204 | .latest_keyframe() |
| 205 | .map_or(0, |frame| frame.frame_sequence); |
| 206 | let mut rx = self.inner.sender.subscribe(); |
| 207 | self.request_keyframe_immediate(); |
| 208 | |
| 209 | loop { |
| 210 | if let Some(frame) = self.latest_keyframe() { |
| 211 | if frame.frame_sequence > baseline_sequence { |
| 212 | return Some(frame); |
| 213 | } |
| 214 | } |
| 215 | |
| 216 | let now = Instant::now(); |
| 217 | if now >= deadline { |
| 218 | return None; |
| 219 | } |
| 220 | |
| 221 | let remaining = deadline - now; |
| 222 | match timeout(remaining, rx.recv()).await { |
| 223 | Ok(Ok(frame)) if frame.is_keyframe && frame.frame_sequence > baseline_sequence => { |
| 224 | return Some(frame) |
| 225 | } |
| 226 | Ok(Ok(_)) => self.request_keyframe(), |
| 227 | Ok(Err(broadcast::error::RecvError::Lagged(_))) => { |
| 228 | self.request_keyframe(); |
| 229 | } |
| 230 | Ok(Err(_)) | Err(_) => return None, |
| 231 | } |
| 232 | } |
| 233 | } |
| 234 | |
| 235 | pub fn request_refresh(&self) { |
| 236 | self.inner.request_refresh(); |
nothing calls this directly
no test coverage detected