Skip to main content

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}