From db0070def239f0caf300033b6bcfd2affbef6f1c Mon Sep 17 00:00:00 2001 From: euzu Date: Fri, 10 Oct 2025 16:40:09 +0200 Subject: [PATCH 1/2] 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(); From c286a88f50e1c0c7a2eec17983af814d1fc2b95d Mon Sep 17 00:00:00 2001 From: euzu Date: Fri, 10 Oct 2025 17:46:26 +0200 Subject: [PATCH 2/2] Some little fixes for epg_view and stream task yielding for better async operations --- .../src/app/components/playlist/epg_view.rs | 25 ++++++++++++------- 1 file changed, 16 insertions(+), 9 deletions(-) diff --git a/frontend/src/app/components/playlist/epg_view.rs b/frontend/src/app/components/playlist/epg_view.rs index be6b66535..92b1a6fd0 100644 --- a/frontend/src/app/components/playlist/epg_view.rs +++ b/frontend/src/app/components/playlist/epg_view.rs @@ -99,7 +99,7 @@ pub fn EpgView() -> Html { .unwrap_or_else(|_| String::new()); // fallback if not set - row_height.trim_end_matches("px").parse::().unwrap_or(60) + row_height.trim_end_matches("px").parse::().unwrap_or(60).max(1) }); // Add scroll listener to calculate visible channels @@ -108,12 +108,13 @@ pub fn EpgView() -> Html { let visible_range = visible_range.clone(); let channel_row_height = *row_height; use_effect_with((), move |_| { + let debounce_handle: Rc>> = Rc::new(RefCell::new(None)); + let onscroll_handle: Rc>>> = Rc::new(RefCell::new(None)); if let Some(div) = container_ref.cast::() { let visible_range = visible_range.clone(); - // Store debounce timer in Rc - let debounce_handle: Rc>> = Rc::new(RefCell::new(None)); let debounce_handle_clone = debounce_handle.clone(); + let onscroll_handle_clone = onscroll_handle.clone(); let onscroll = Closure::wrap(Box::new(move |_event: web_sys::Event| { // Cancel previous scheduled update if let Some(prev) = debounce_handle_clone.borrow_mut().take() { @@ -128,7 +129,7 @@ pub fn EpgView() -> Html { let client_height = div.client_height(); // Calculate which channel rows are visible // render 10 + 10 more lines - let start_index = (scroll_top / (channel_row_height as i32) - 10).max(0); + 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)); } @@ -136,11 +137,17 @@ pub fn EpgView() -> Html { *debounce_handle_clone.borrow_mut() = Some(handle); }) as Box); - div.add_event_listener_with_callback("scroll", onscroll.as_ref().unchecked_ref()).unwrap(); - onscroll.forget(); + *onscroll_handle_clone.borrow_mut() = Some(onscroll); + } + move || { + if let Some(prev) = debounce_handle.borrow_mut().take() { + prev.cancel(); + } + if let Some(onscroll) = onscroll_handle.borrow_mut().take() { + drop(onscroll); + } } - || {} }); } @@ -176,7 +183,7 @@ pub fn EpgView() -> Html {
{ for tv.channels.iter().enumerate().skip(start_index).take(end_index - start_index).map(|(_i, ch)| { html! { -
+
{ if let Some(icon) = &ch.icon { html! { {ch.title.clone()} } @@ -215,7 +222,7 @@ pub fn EpgView() -> Html {
{ for tv.channels.iter().enumerate().skip(start_index).take(end_index - start_index).map(|(_i, ch)| { html! { -
+
{ for ch.programmes.iter().map(|p| { let is_active = now >= p.start && now < p.stop; let left = get_pos(p.start, start_window);