From db0070def239f0caf300033b6bcfd2affbef6f1c Mon Sep 17 00:00:00 2001 From: euzu Date: Fri, 10 Oct 2025 16:40:09 +0200 Subject: [PATCH] Some little fixes for epg_view and stream task yielding for better async operations --- Cargo.lock | 1 + .../src/api/model/streams/buffered_stream.rs | 22 ++++++++--------- .../model/streams/shared_stream_manager.rs | 1 + .../src/api/model/streams/throttled_stream.rs | 7 ++++-- frontend/Cargo.toml | 5 ++-- .../src/app/components/playlist/epg_view.rs | 24 ++++++++++--------- 6 files changed, 34 insertions(+), 26 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 690c39567..c5e7934e5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1197,6 +1197,7 @@ dependencies = [ "cron", "futures", "futures-signals", + "gloo-render 0.2.0", "gloo-storage 0.3.0", "gloo-timers 0.3.0", "gloo-utils 0.2.0", diff --git a/backend/src/api/model/streams/buffered_stream.rs b/backend/src/api/model/streams/buffered_stream.rs index ae9880ec5..fbc7ebbf3 100644 --- a/backend/src/api/model/streams/buffered_stream.rs +++ b/backend/src/api/model/streams/buffered_stream.rs @@ -30,20 +30,20 @@ impl BufferedStream { mut stream: BoxedProviderStream, client_close_signal: Arc, ) { - loop { - if !client_close_signal.is_active() { - break; - } + while client_close_signal.is_active() { match stream.next().await { Some(Ok(chunk)) => { - match tx.reserve().await { - Ok(permit) => permit.send(Ok(chunk)), - Err(_err) => { - // Receiver dropped, notify and exit - client_close_signal.notify(); - break; - } + match tx.reserve().await { + Ok(permit) => { + permit.send(Ok(chunk)); + tokio::task::yield_now().await; + }, + Err(_err) => { + // Receiver dropped, notify and exit + client_close_signal.notify(); + break; } + } } Some(Err(err)) => { //trace!("Buffered Stream Error: {err:?}"); diff --git a/backend/src/api/model/streams/shared_stream_manager.rs b/backend/src/api/model/streams/shared_stream_manager.rs index 56186cd59..7dcffec4d 100644 --- a/backend/src/api/model/streams/shared_stream_manager.rs +++ b/backend/src/api/model/streams/shared_stream_manager.rs @@ -107,6 +107,7 @@ impl SharedStreamState { debug!("Shared stream client send error: {address} {err}"); break; } + tokio::task::yield_now().await; } Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => { trace!("Client lagged behind. Skipped {skipped} messages. {address}"); diff --git a/backend/src/api/model/streams/throttled_stream.rs b/backend/src/api/model/streams/throttled_stream.rs index 39a905fde..8c7227140 100644 --- a/backend/src/api/model/streams/throttled_stream.rs +++ b/backend/src/api/model/streams/throttled_stream.rs @@ -56,7 +56,8 @@ where match Pin::new(&mut this.inner).poll_next(cx) { Poll::Ready(Some(Ok(bytes))) => { let len = bytes.len() as f64; - let delay_duration = Duration::from_secs_f64(len / this.rate_bytes_per_sec); + let delay_secs = (len / this.rate_bytes_per_sec).max(0.001); + let delay_duration = Duration::from_secs_f64(delay_secs); // Schedule the next delay this.next_delay = Some(Box::pin(sleep(delay_duration))); @@ -71,4 +72,6 @@ where Poll::Pending => Poll::Pending, } } -} \ No newline at end of file +} + +impl Unpin for ThrottledStream {} \ No newline at end of file diff --git a/frontend/Cargo.toml b/frontend/Cargo.toml index 0992392ee..e415da369 100644 --- a/frontend/Cargo.toml +++ b/frontend/Cargo.toml @@ -18,8 +18,9 @@ thiserror = "2" reqwasm = {version = "0.5", features = ["json"]} log = "0.4" wasm-logger = "0.2" -gloo-storage = "0" -gloo-utils = "0" +gloo-storage = "0.3" +gloo-utils = "0.2" +gloo-render = "0.2" gloo-timers = { version = "0", features = ["futures"] } implicit-clone = "0" futures-signals = "0.3" diff --git a/frontend/src/app/components/playlist/epg_view.rs b/frontend/src/app/components/playlist/epg_view.rs index 0d1de37ee..be6b66535 100644 --- a/frontend/src/app/components/playlist/epg_view.rs +++ b/frontend/src/app/components/playlist/epg_view.rs @@ -119,20 +119,22 @@ pub fn EpgView() -> Html { if let Some(prev) = debounce_handle_clone.borrow_mut().take() { prev.cancel(); } - if let Some(div) = container_ref.cast::() { - let scroll_top = div.scroll_top(); - let client_height = div.client_height(); // Calculate which channel rows are visible - let vr = visible_range.clone(); + // Schedule a new update after X ms (debounce) + let container_ref = container_ref.clone(); + let vr = visible_range.clone(); + let handle = Timeout::new(16, move || { + if let Some(div) = container_ref.cast::() { + let scroll_top = div.scroll_top(); + let client_height = div.client_height(); // Calculate which channel rows are visible - // Schedule a new update after 80ms (debounce) - let handle = Timeout::new(80, move || { - let start_index = (scroll_top / (channel_row_height as i32)).max(0); - let end_index = ((scroll_top + client_height) / (channel_row_height as i32) + 1).max(0); + // render 10 + 10 more lines + let start_index = (scroll_top / (channel_row_height as i32) - 10).max(0); + let end_index = ((scroll_top + client_height) / (channel_row_height as i32) + 10).max(0); vr.set((start_index as usize, end_index as usize)); - }); + } + }); - *debounce_handle_clone.borrow_mut() = Some(handle); - } + *debounce_handle_clone.borrow_mut() = Some(handle); }) as Box); div.add_event_listener_with_callback("scroll", onscroll.as_ref().unchecked_ref()).unwrap();