indicatrix_dispatch/pool/mod.rs
1//! [`LanePool`]: N lanes rendering one image epoch against one [`SampleCursor`].
2//!
3//! # One run
4//!
5//! [`LanePool::run`] (or [`LanePool::run_into`]) spawns one scoped thread per lane.
6//! Each lane loops:
7//!
8//! 1. size a chunk from its own [`RateModel`] under the pool's [`ChunkPolicy`]
9//! (a calibration chunk while uncalibrated);
10//! 2. claim it from the epoch's cursor ([`SampleCursor::claim_any`]: a failed lane's
11//! returned remainder first, fresh samples after);
12//! 3. trace it with [`WorkerLane::render_chunk`];
13//! 4. merge the valid prefix into the [`Merger`] (chunk-start order, see its
14//! determinism notes) and fold the measurement into its rate model;
15//! 5. on a short chunk: requeue the untraced tail for any lane, count a consecutive
16//! failure, back off, and retire after [`PoolConfig::retire_after_failures`].
17//!
18//! A lane with nothing to claim waits while any other lane still has a chunk in flight
19//! (that chunk could fail and come back); claims and in-flight bookkeeping share one
20//! lock, so no lane can conclude "all done" between another lane's claim and its
21//! bookkeeping. The run ends when every lane has exited: all samples traced exactly
22//! once ([`PoolStatus::Complete`]), cancelled, or every lane retired with samples
23//! left ([`PoolStatus::LanesExhausted`]).
24//!
25//! A lane that panics inside `render_chunk` is treated as a failed chunk with nothing
26//! traced; the panic does not take the pool down.
27
28mod epoch;
29mod lane_loop;
30#[cfg(test)]
31mod tests;
32
33use crate::{CancelToken, ChunkPolicy, Merger, RateModel, SampleCursor, SampleRange, WorkerLane};
34use epoch::Epoch;
35use glam::Vec3;
36use indicatrix_net::SceneState;
37use std::{
38 fmt,
39 sync::{Arc, Mutex, PoisonError},
40 thread,
41 time::Duration,
42};
43
44/// Failure handling and chunk sizing for a [`LanePool`].
45#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46pub struct PoolConfig {
47 /// How chunks are sized from each lane's rate.
48 pub policy: ChunkPolicy,
49 /// Consecutive failed chunks a lane may have before it starts pausing between
50 /// attempts (`1`: pause after every failure). A single dropped connection is common
51 /// enough that retrying at once is often right.
52 pub pause_after_failures: u32,
53 /// Consecutive failed chunks after which a lane is retired for the rest of the
54 /// epoch. `u32::MAX` never retires.
55 pub retire_after_failures: u32,
56 /// The first pause; each further consecutive failure doubles it.
57 pub backoff_initial: Duration,
58 /// The longest pause.
59 pub backoff_max: Duration,
60}
61
62impl PoolConfig {
63 /// Still export / tilt video: 22 s chunks, pause from the 2nd consecutive failure
64 /// (15 s doubling to 120 s, the desktop export's schedule), retire after 5.
65 pub const EXPORT: Self = Self {
66 policy: ChunkPolicy::EXPORT,
67 pause_after_failures: 2,
68 retire_after_failures: 5,
69 backoff_initial: Duration::from_secs(15),
70 backoff_max: Duration::from_secs(120),
71 };
72
73 /// Interactive requests: 1.5 s chunks, short pauses (250 ms to 2 s), retire after 3.
74 pub const INTERACTIVE: Self = Self {
75 policy: ChunkPolicy::INTERACTIVE,
76 pause_after_failures: 1,
77 retire_after_failures: 3,
78 backoff_initial: Duration::from_millis(250),
79 backoff_max: Duration::from_secs(2),
80 };
81
82 /// The pause after `consecutive_failures` failed chunks in a row: zero below
83 /// [`Self::pause_after_failures`], then `backoff_initial` doubling per further
84 /// failure, capped at `backoff_max`.
85 #[must_use]
86 pub fn backoff(&self, consecutive_failures: u32) -> Duration {
87 if consecutive_failures < self.pause_after_failures.max(1) {
88 return Duration::ZERO;
89 }
90 let doublings = consecutive_failures - self.pause_after_failures.max(1);
91 self.backoff_initial
92 .saturating_mul(1u32 << doublings.min(16))
93 .min(self.backoff_max)
94 }
95}
96
97/// Something that happened during a [`LanePool`] run. `lane` is the index
98/// [`LanePool::add_lane`] returned; [`LanePool::lane_name`] names it.
99///
100/// Events are delivered on the lane threads, concurrently, in no global order.
101#[derive(Debug, Clone, PartialEq, Eq)]
102pub enum PoolEvent {
103 /// The lane's thread started claiming work.
104 LaneStarted {
105 /// The lane.
106 lane: usize,
107 },
108 /// A chunk's traced prefix was merged: progress.
109 ChunkMerged {
110 /// The lane.
111 lane: usize,
112 /// The chunk as claimed.
113 range: SampleRange,
114 /// How many of its samples were traced and merged (a prefix).
115 done: u32,
116 /// The exact merged total across all lanes after this chunk.
117 total_done: u32,
118 /// The run's target sample count.
119 target: u32,
120 },
121 /// A chunk ended short; its untraced tail went back to the cursor.
122 LaneFailed {
123 /// The lane.
124 lane: usize,
125 /// Why (the lane's own error text, or a pool-side reason).
126 error: String,
127 /// The tail returned to the cursor for any lane.
128 returned: SampleRange,
129 /// Failed chunks in a row, this one included.
130 consecutive_failures: u32,
131 /// The pause the lane now takes; `None` when it is being retired instead.
132 pause: Option<Duration>,
133 },
134 /// The lane failed too often in a row and takes no more work this run.
135 LaneRetired {
136 /// The lane.
137 lane: usize,
138 /// Failed chunks in a row.
139 consecutive_failures: u32,
140 },
141 /// The lane stopped normally: nothing left to claim, or cancelled.
142 LaneFinished {
143 /// The lane.
144 lane: usize,
145 /// Chunks it merged at least one sample from.
146 chunks: u32,
147 /// Samples it contributed.
148 samples: u32,
149 },
150}
151
152/// How a run ended.
153#[derive(Debug, Clone, Copy, PartialEq, Eq)]
154pub enum PoolStatus {
155 /// Every sample of the range was traced exactly once.
156 Complete,
157 /// The cancel token was raised; `missing` samples were not traced.
158 Cancelled {
159 /// Samples of the range not merged.
160 missing: u32,
161 },
162 /// Every lane retired (or there were none) with `missing` samples untraced. A
163 /// coordinator reports this to the viewer as a lost job.
164 LanesExhausted {
165 /// Samples of the range not merged.
166 missing: u32,
167 },
168}
169
170/// The result of [`LanePool::run`].
171#[derive(Debug, Clone, PartialEq)]
172pub struct PoolOutcome {
173 /// Merged per-pixel radiance sum, `width * height` long.
174 pub sum: Vec<Vec3>,
175 /// Exact number of samples in `sum`.
176 pub count: u32,
177 /// How the run ended.
178 pub status: PoolStatus,
179}
180
181/// [`LanePool::run_into`] was handed a [`Merger`] built for a different image or range.
182#[derive(Debug, Clone, Copy, PartialEq, Eq)]
183pub struct MergerMismatch {
184 /// `scene.width * scene.height`.
185 pub expected_pixels: usize,
186 /// The range's first sample.
187 pub expected_first_sample: u32,
188 /// The merger's pixel count.
189 pub merger_pixels: usize,
190 /// The merger's first sample.
191 pub merger_first_sample: u32,
192 /// Samples the merger already held (a fresh merger holds none).
193 pub merger_samples: u32,
194}
195
196impl fmt::Display for MergerMismatch {
197 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
198 write!(
199 f,
200 "merger built for {} px from sample {} holding {} samples, run needs a fresh \
201 one for {} px from sample {}",
202 self.merger_pixels,
203 self.merger_first_sample,
204 self.merger_samples,
205 self.expected_pixels,
206 self.expected_first_sample
207 )
208 }
209}
210
211impl std::error::Error for MergerMismatch {}
212
213/// One registered lane and its rate model, which persists across runs.
214struct LaneSlot {
215 lane: Arc<dyn WorkerLane>,
216 rate: Mutex<RateModel>,
217}
218
219/// See the module doc. Lanes and their rate models persist across runs, so a caller
220/// rendering many images of one scene (a tilt video, successive coordinator jobs)
221/// calibrates each lane once.
222pub struct LanePool {
223 config: PoolConfig,
224 lanes: Vec<LaneSlot>,
225}
226
227impl LanePool {
228 /// An empty pool.
229 #[must_use]
230 pub const fn new(config: PoolConfig) -> Self {
231 Self {
232 config,
233 lanes: Vec::new(),
234 }
235 }
236
237 /// Registers `lane` starting from `rate`, returning its index.
238 pub fn add_lane(&mut self, lane: Arc<dyn WorkerLane>, rate: RateModel) -> usize {
239 self.lanes.push(LaneSlot {
240 lane,
241 rate: Mutex::new(rate),
242 });
243 self.lanes.len() - 1
244 }
245
246 /// The pool's configuration.
247 #[must_use]
248 pub const fn config(&self) -> &PoolConfig {
249 &self.config
250 }
251
252 /// How many lanes are registered.
253 #[must_use]
254 pub const fn lane_count(&self) -> usize {
255 self.lanes.len()
256 }
257
258 /// Lane `lane`'s [`WorkerLane::name`].
259 #[must_use]
260 pub fn lane_name(&self, lane: usize) -> Option<&str> {
261 self.lanes.get(lane).map(|slot| slot.lane.name())
262 }
263
264 /// Lane `lane`'s current rate model.
265 #[must_use]
266 pub fn rate(&self, lane: usize) -> Option<RateModel> {
267 self.lanes
268 .get(lane)
269 .map(|slot| *slot.rate.lock().unwrap_or_else(PoisonError::into_inner))
270 }
271
272 /// Renders `range` of `scene` with every lane and returns the merged result.
273 /// Blocks until the run ends (see [`PoolStatus`]); `events` is called from the lane
274 /// threads as things happen.
275 pub fn run(
276 &self,
277 scene: &SceneState,
278 range: SampleRange,
279 cancel: &CancelToken,
280 events: &(dyn Fn(PoolEvent) + Sync),
281 ) -> PoolOutcome {
282 let merger = Merger::new(pixel_count(scene), range.first_sample);
283 let status = self.run_unchecked(scene, range, &merger, cancel, events);
284 let (sum, count) = merger.into_parts();
285 PoolOutcome { sum, count, status }
286 }
287
288 /// Like [`Self::run`], merging into a caller-owned `merger` so another thread can
289 /// read progressive snapshots ([`Merger::snapshot_into`]) while the run is going.
290 /// `merger` must be fresh, built with `Merger::new(width * height,
291 /// range.first_sample)`.
292 ///
293 /// # Errors
294 ///
295 /// [`MergerMismatch`] when `merger` was built for another size or range; nothing
296 /// is rendered then.
297 pub fn run_into(
298 &self,
299 scene: &SceneState,
300 range: SampleRange,
301 merger: &Merger,
302 cancel: &CancelToken,
303 events: &(dyn Fn(PoolEvent) + Sync),
304 ) -> Result<PoolStatus, MergerMismatch> {
305 let expected_pixels = pixel_count(scene);
306 let merger_samples = merger.total();
307 if merger.pixel_count() != expected_pixels
308 || merger.first_sample() != range.first_sample
309 || merger_samples != 0
310 {
311 return Err(MergerMismatch {
312 expected_pixels,
313 expected_first_sample: range.first_sample,
314 merger_pixels: merger.pixel_count(),
315 merger_first_sample: merger.first_sample(),
316 merger_samples,
317 });
318 }
319 Ok(self.run_unchecked(scene, range, merger, cancel, events))
320 }
321
322 fn run_unchecked(
323 &self,
324 scene: &SceneState,
325 range: SampleRange,
326 merger: &Merger,
327 cancel: &CancelToken,
328 events: &(dyn Fn(PoolEvent) + Sync),
329 ) -> PoolStatus {
330 let epoch = Epoch {
331 scene,
332 cursor: SampleCursor::new(range.first_sample, range.end()),
333 merger,
334 cancel,
335 events,
336 config: &self.config,
337 target: range.samples,
338 pixels: merger.pixel_count(),
339 sched: epoch::Sched::new(),
340 };
341 if !range.is_empty() {
342 thread::scope(|scope| {
343 for (index, slot) in self.lanes.iter().enumerate() {
344 let epoch = &epoch;
345 scope.spawn(move || {
346 lane_loop::run_lane(epoch, index, &*slot.lane, &slot.rate);
347 });
348 }
349 });
350 }
351 let missing = range.samples.saturating_sub(merger.total());
352 if missing == 0 {
353 PoolStatus::Complete
354 } else if cancel.is_cancelled() {
355 PoolStatus::Cancelled { missing }
356 } else {
357 PoolStatus::LanesExhausted { missing }
358 }
359 }
360}
361
362/// `width * height` of `scene`.
363const fn pixel_count(scene: &SceneState) -> usize {
364 scene.width as usize * scene.height as usize
365}