Some little fixes for epg_view and stream task yielding for better async operations

This commit is contained in:
euzu
2025-10-10 16:40:09 +02:00
parent 3dea155475
commit db0070def2
6 changed files with 34 additions and 26 deletions
@@ -30,20 +30,20 @@ impl BufferedStream {
mut stream: BoxedProviderStream,
client_close_signal: Arc<AtomicOnceFlag>,
) {
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:?}");
@@ -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}");
@@ -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,
}
}
}
}
impl<S: Unpin> Unpin for ThrottledStream<S> {}