use super::{PoolConfig, PoolEvent};
use crate::{CancelToken, Merger, SampleCursor, SampleRange};
use indicatrix_net::SceneState;
use std::{
sync::{Condvar, Mutex, MutexGuard, PoisonError},
time::{Duration, Instant},
};
const POLL: Duration = Duration::from_millis(25);
pub(super) struct Sched {
in_flight: Mutex<u32>,
wake: Condvar,
}
impl Sched {
pub(super) const fn new() -> Self {
Self {
in_flight: Mutex::new(0),
wake: Condvar::new(),
}
}
fn lock(&self) -> MutexGuard<'_, u32> {
self.in_flight
.lock()
.unwrap_or_else(PoisonError::into_inner)
}
fn wait<'g>(&self, guard: MutexGuard<'g, u32>, timeout: Duration) -> MutexGuard<'g, u32> {
self.wake
.wait_timeout(guard, timeout)
.unwrap_or_else(PoisonError::into_inner)
.0
}
}
pub(super) struct Epoch<'a> {
pub(super) scene: &'a SceneState,
pub(super) cursor: SampleCursor,
pub(super) merger: &'a Merger,
pub(super) cancel: &'a CancelToken,
pub(super) events: &'a (dyn Fn(PoolEvent) + Sync),
pub(super) config: &'a PoolConfig,
pub(super) target: u32,
pub(super) pixels: usize,
pub(super) sched: Sched,
}
impl Epoch<'_> {
pub(super) fn emit(&self, event: PoolEvent) {
(self.events)(event);
}
pub(super) fn claim(&self, want: u32) -> Option<SampleRange> {
let mut in_flight = self.sched.lock();
loop {
if self.cancel.is_cancelled() {
return None;
}
if let Some((first, count)) = self.cursor.claim_any(want) {
*in_flight += 1;
return Some(SampleRange::new(first, count));
}
if *in_flight == 0 {
return None;
}
in_flight = self.sched.wait(in_flight, POLL);
}
}
pub(super) fn settle(&self, range: SampleRange, done: u32) {
let tail = range.after_prefix(done);
let mut in_flight = self.sched.lock();
self.cursor.requeue(tail.first_sample, tail.samples);
*in_flight -= 1;
drop(in_flight);
self.sched.wake.notify_all();
}
pub(super) fn pause(&self, total: Duration) -> bool {
let deadline = Instant::now() + total;
let mut in_flight = self.sched.lock();
loop {
if self.cancel.is_cancelled() || (*in_flight == 0 && self.cursor.fully_claimed()) {
return false;
}
let now = Instant::now();
if now >= deadline {
return true;
}
in_flight = self.sched.wait(in_flight, (deadline - now).min(POLL));
}
}
}