diff --git a/cli/src/native/actions.rs b/cli/src/native/actions.rs index 5ba4e19..2887853 100644 --- a/cli/src/native/actions.rs +++ b/cli/src/native/actions.rs @@ -1903,11 +1903,17 @@ async fn handle_reload(state: &mut DaemonState) -> Result { let mut rx = mgr.client.subscribe(); let _ = tokio::time::timeout(tokio::time::Duration::from_secs(10), async { - while let Ok(event) = rx.recv().await { - if event.method == "Page.loadEventFired" - && event.session_id.as_deref() == Some(&session_id) - { - return; + loop { + match rx.recv().await { + Ok(event) => { + if event.method == "Page.loadEventFired" + && event.session_id.as_deref() == Some(&session_id) + { + return; + } + } + Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue, + Err(_) => break, } } }) @@ -4165,6 +4171,7 @@ async fn handle_responsebody(cmd: &Value, state: &DaemonState) -> Result continue, Ok(Err(_)) => return Err("Event stream closed".to_string()), Err(_) => { return Err(format!( @@ -4203,6 +4210,7 @@ async fn handle_waitfordownload(cmd: &Value, state: &DaemonState) -> Result continue, Ok(Err(_)) => return Err("Event stream closed".to_string()), Err(_) => return Err("Timeout waiting for download".to_string()), } diff --git a/cli/src/native/browser.rs b/cli/src/native/browser.rs index 7a84177..c41bc37 100644 --- a/cli/src/native/browser.rs +++ b/cli/src/native/browser.rs @@ -469,9 +469,17 @@ impl BrowserManager { let timeout = tokio::time::Duration::from_millis(self.default_timeout_ms); tokio::time::timeout(timeout, async { - while let Ok(event) = rx.recv().await { - if event.method == event_name && event.session_id.as_deref() == Some(session_id) { - return Ok(()); + loop { + match rx.recv().await { + Ok(event) => { + if event.method == event_name + && event.session_id.as_deref() == Some(session_id) + { + return Ok(()); + } + } + Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue, + Err(tokio::sync::broadcast::error::RecvError::Closed) => break, } } Err("Event stream closed".to_string()) @@ -526,6 +534,7 @@ impl BrowserManager { } } Ok(Ok(_)) => {} + Ok(Err(tokio::sync::broadcast::error::RecvError::Lagged(_))) => continue, Ok(Err(_)) => break, Err(_) => { // Timeout on recv -- check if idle long enough