mod epoch;
mod lane_loop;
#[cfg(test)]
mod tests;
use crate::{CancelToken, ChunkPolicy, Merger, RateModel, SampleCursor, SampleRange, WorkerLane};
use epoch::Epoch;
use glam::Vec3;
use indicatrix_net::SceneState;
use std::{
fmt,
sync::{Arc, Mutex, PoisonError},
thread,
time::Duration,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct PoolConfig {
pub policy: ChunkPolicy,
pub pause_after_failures: u32,
pub retire_after_failures: u32,
pub backoff_initial: Duration,
pub backoff_max: Duration,
}
impl PoolConfig {
pub const EXPORT: Self = Self {
policy: ChunkPolicy::EXPORT,
pause_after_failures: 2,
retire_after_failures: 5,
backoff_initial: Duration::from_secs(15),
backoff_max: Duration::from_secs(120),
};
pub const INTERACTIVE: Self = Self {
policy: ChunkPolicy::INTERACTIVE,
pause_after_failures: 1,
retire_after_failures: 3,
backoff_initial: Duration::from_millis(250),
backoff_max: Duration::from_secs(2),
};
#[must_use]
pub fn backoff(&self, consecutive_failures: u32) -> Duration {
if consecutive_failures < self.pause_after_failures.max(1) {
return Duration::ZERO;
}
let doublings = consecutive_failures - self.pause_after_failures.max(1);
self.backoff_initial
.saturating_mul(1u32 << doublings.min(16))
.min(self.backoff_max)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PoolEvent {
LaneStarted {
lane: usize,
},
ChunkMerged {
lane: usize,
range: SampleRange,
done: u32,
total_done: u32,
target: u32,
},
LaneFailed {
lane: usize,
error: String,
returned: SampleRange,
consecutive_failures: u32,
pause: Option<Duration>,
},
LaneRetired {
lane: usize,
consecutive_failures: u32,
},
LaneFinished {
lane: usize,
chunks: u32,
samples: u32,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PoolStatus {
Complete,
Cancelled {
missing: u32,
},
LanesExhausted {
missing: u32,
},
}
#[derive(Debug, Clone, PartialEq)]
pub struct PoolOutcome {
pub sum: Vec<Vec3>,
pub count: u32,
pub status: PoolStatus,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct MergerMismatch {
pub expected_pixels: usize,
pub expected_first_sample: u32,
pub merger_pixels: usize,
pub merger_first_sample: u32,
pub merger_samples: u32,
}
impl fmt::Display for MergerMismatch {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(
f,
"merger built for {} px from sample {} holding {} samples, run needs a fresh \
one for {} px from sample {}",
self.merger_pixels,
self.merger_first_sample,
self.merger_samples,
self.expected_pixels,
self.expected_first_sample
)
}
}
impl std::error::Error for MergerMismatch {}
struct LaneSlot {
lane: Arc<dyn WorkerLane>,
rate: Mutex<RateModel>,
}
pub struct LanePool {
config: PoolConfig,
lanes: Vec<LaneSlot>,
}
impl LanePool {
#[must_use]
pub const fn new(config: PoolConfig) -> Self {
Self {
config,
lanes: Vec::new(),
}
}
pub fn add_lane(&mut self, lane: Arc<dyn WorkerLane>, rate: RateModel) -> usize {
self.lanes.push(LaneSlot {
lane,
rate: Mutex::new(rate),
});
self.lanes.len() - 1
}
#[must_use]
pub const fn config(&self) -> &PoolConfig {
&self.config
}
#[must_use]
pub const fn lane_count(&self) -> usize {
self.lanes.len()
}
#[must_use]
pub fn lane_name(&self, lane: usize) -> Option<&str> {
self.lanes.get(lane).map(|slot| slot.lane.name())
}
#[must_use]
pub fn rate(&self, lane: usize) -> Option<RateModel> {
self.lanes
.get(lane)
.map(|slot| *slot.rate.lock().unwrap_or_else(PoisonError::into_inner))
}
pub fn run(
&self,
scene: &SceneState,
range: SampleRange,
cancel: &CancelToken,
events: &(dyn Fn(PoolEvent) + Sync),
) -> PoolOutcome {
let merger = Merger::new(pixel_count(scene), range.first_sample);
let status = self.run_unchecked(scene, range, &merger, cancel, events);
let (sum, count) = merger.into_parts();
PoolOutcome { sum, count, status }
}
pub fn run_into(
&self,
scene: &SceneState,
range: SampleRange,
merger: &Merger,
cancel: &CancelToken,
events: &(dyn Fn(PoolEvent) + Sync),
) -> Result<PoolStatus, MergerMismatch> {
let expected_pixels = pixel_count(scene);
let merger_samples = merger.total();
if merger.pixel_count() != expected_pixels
|| merger.first_sample() != range.first_sample
|| merger_samples != 0
{
return Err(MergerMismatch {
expected_pixels,
expected_first_sample: range.first_sample,
merger_pixels: merger.pixel_count(),
merger_first_sample: merger.first_sample(),
merger_samples,
});
}
Ok(self.run_unchecked(scene, range, merger, cancel, events))
}
fn run_unchecked(
&self,
scene: &SceneState,
range: SampleRange,
merger: &Merger,
cancel: &CancelToken,
events: &(dyn Fn(PoolEvent) + Sync),
) -> PoolStatus {
let epoch = Epoch {
scene,
cursor: SampleCursor::new(range.first_sample, range.end()),
merger,
cancel,
events,
config: &self.config,
target: range.samples,
pixels: merger.pixel_count(),
sched: epoch::Sched::new(),
};
if !range.is_empty() {
thread::scope(|scope| {
for (index, slot) in self.lanes.iter().enumerate() {
let epoch = &epoch;
scope.spawn(move || {
lane_loop::run_lane(epoch, index, &*slot.lane, &slot.rate);
});
}
});
}
let missing = range.samples.saturating_sub(merger.total());
if missing == 0 {
PoolStatus::Complete
} else if cancel.is_cancelled() {
PoolStatus::Cancelled { missing }
} else {
PoolStatus::LanesExhausted { missing }
}
}
}
const fn pixel_count(scene: &SceneState) -> usize {
scene.width as usize * scene.height as usize
}