Skip to main content

indicatrix_dispatch/sample_cursor/
mod.rs

1//! [`SampleCursor`]: the shared claim point every backend contributing to one image
2//! draws disjoint absolute sample sub-ranges from -- the desktop export's local and
3//! remote lanes, the live viewport's epoch (`LiveEpoch` in `indicatrix-cut`), and every
4//! lane of a [`crate::LanePool`].
5//!
6//! Moved here from `indicatrix-cut`'s `bridge::sample_cursor`, which now re-exports it,
7//! so the coordinator can share the exact same claim semantics.
8//!
9//! # Why a shared claim point, not a one-shot split
10//!
11//! A fixed up-front split (measured once, dispatched once) can leave a fast engine's
12//! assigned slice finished early with nothing further to claim while the render keeps
13//! going, sitting idle for the rest of it. `SampleCursor` replaces that with a shared
14//! claim point every engine pulls fresh work from whenever free, for as long as any
15//! remains.
16//!
17//! # Disjointness by construction
18//!
19//! [`claim`](SampleCursor::claim) and [`claim_local`](SampleCursor::claim_local) hand
20//! out ranges via a single atomic `fetch_add`, so any number of callers on any number
21//! of threads can never receive overlapping indices, by construction, with no reliance
22//! on callers coordinating access any other way. Disjointness matters because both the
23//! CPU tracer and GPU megakernel derive each sample's pixel-jitter/RNG draws from its
24//! *absolute* index: overlapping ranges would redraw identical samples and silently
25//! bias the average toward whatever indices got traced twice -- wrong, but not
26//! obviously wrong.
27//!
28//! # Two retry piles
29//!
30//! A range a lane failed to finish must be traced again, exactly once, by someone.
31//!
32//! - [`return_to_local`](SampleCursor::return_to_local) (the desktop's single-remote
33//!   model) pushes it onto a pile only [`claim_local`](SampleCursor::claim_local) /
34//!   [`claim_local_bounded`](SampleCursor::claim_local_bounded) drain, never back onto
35//!   the shared pool [`claim`](SampleCursor::claim) draws from. Otherwise remote could
36//!   re-claim its own just-failed range on its very next iteration -- an infinite loop
37//!   for a persistent (not transient) failure. Routing it to a local-only pile
38//!   guarantees it is retried at most once more, locally.
39//! - [`requeue`](SampleCursor::requeue) (the [`crate::LanePool`] model, where every
40//!   lane is equal) pushes it onto a pile [`claim_any`](SampleCursor::claim_any) drains
41//!   before fresh work. The pool itself prevents the infinite loop: a failing lane
42//!   backs off and is retired after a bounded number of consecutive failures.
43//!
44//! The two models are never mixed on one cursor in practice; [`claim`] alone never
45//! sees either pile.
46//!
47//! [`claim`]: SampleCursor::claim
48
49use std::{
50    collections::VecDeque,
51    sync::{
52        Mutex, PoisonError,
53        atomic::{AtomicU32, Ordering},
54    },
55};
56
57#[cfg(test)]
58mod tests;
59
60/// See this module's doc comment. Every method takes `&self`, so one instance can be
61/// shared (by reference inside a `thread::scope`, or behind an `Arc`) across every
62/// lane's own thread.
63#[derive(Debug)]
64pub struct SampleCursor {
65    /// Absolute index of the next sample NO engine has claimed yet. Every claim is one
66    /// `fetch_add` against this -- see [`claim`](Self::claim).
67    next: AtomicU32,
68    /// One past the last sample in this budget (exclusive). A claim is never honoured
69    /// past this, however large a range was requested.
70    end: u32,
71    /// Ranges a remote chunk failed to finish, reserved for the local lane alone -- see
72    /// this module's doc comment.
73    local_retry: Mutex<VecDeque<(u32, u32)>>,
74    /// Ranges a pool lane failed to finish, open to every lane via
75    /// [`claim_any`](Self::claim_any).
76    requeued: Mutex<VecDeque<(u32, u32)>>,
77}
78
79impl SampleCursor {
80    /// Builds a cursor over `[start, end)`. For an export, `start` is wherever a prior
81    /// sequential phase (remote/hybrid calibration) left `samples_done` and `end` is
82    /// the export's total `samples_per_pixel`; for a live epoch it is `[0,
83    /// target_samples)`; for a coordinator job it is the request's
84    /// `[first_sample, first_sample + samples)`.
85    #[must_use]
86    pub const fn new(start: u32, end: u32) -> Self {
87        Self {
88            next: AtomicU32::new(start),
89            end,
90            local_retry: Mutex::new(VecDeque::new()),
91            requeued: Mutex::new(VecDeque::new()),
92        }
93    }
94
95    /// One past the last sample of this budget.
96    #[must_use]
97    pub const fn end(&self) -> u32 {
98        self.end
99    }
100
101    /// Claims up to `want` fresh samples for ANY engine, returning `Some((start,
102    /// count))` with `1 <= count <= want`, or `None` if the budget is already
103    /// exhausted. Never returns a range past `end`, however large `want` is.
104    ///
105    /// `want == 0` always returns `None` rather than a zero-length range.
106    ///
107    /// # Why `fetch_add` alone is enough
108    ///
109    /// One atomic read-modify-write, unconditionally advancing `next` by `want` and
110    /// then checking whether the range it was handed starts past `end`. `fetch_add` is
111    /// inherently exclusive, so every call receives a DISTINCT `start` with no CAS
112    /// retry loop needed. `next` can end up advanced PAST `end` (the last claim before
113    /// exhaustion typically requests more than remains) -- harmless, since the
114    /// returned `count` is clamped to `end - start` and every later claim sees
115    /// `start >= end` and reports `None`. `saturating` is not needed in practice
116    /// (budgets are far below `u32::MAX`), but `fetch_add` wraps, so a budget near
117    /// `u32::MAX` must not be used.
118    pub fn claim(&self, want: u32) -> Option<(u32, u32)> {
119        if want == 0 {
120            return None;
121        }
122        let start = self.next.fetch_add(want, Ordering::Relaxed);
123        if start >= self.end {
124            return None;
125        }
126        Some((start, want.min(self.end - start)))
127    }
128
129    /// Claims for the LOCAL lane only: a previously remote-failed range first (WHOLE,
130    /// however long it is), falling back to [`claim`](Self::claim) once that retry
131    /// pile is empty. See this module's doc comment for why a retried range never goes
132    /// back to remote. The export's local lane batches whatever it gets; the live
133    /// viewport uses [`claim_local_bounded`](Self::claim_local_bounded) instead.
134    pub fn claim_local(&self, want: u32) -> Option<(u32, u32)> {
135        let retried = self
136            .local_retry
137            .lock()
138            .unwrap_or_else(PoisonError::into_inner)
139            .pop_front();
140        if retried.is_some() {
141            return retried;
142        }
143        self.claim(want)
144    }
145
146    /// Like [`claim_local`](Self::claim_local), but never hands out more than `want`
147    /// samples even from the retry pile: a longer retried range is split, its head
148    /// returned now and its tail pushed back to the FRONT of the pile for the next
149    /// call. The live viewport traces one short frame per claim, so a whole failed
150    /// remote chunk (possibly hundreds of samples) must not land in a single frame.
151    /// `want == 0` returns `None` without touching anything.
152    pub fn claim_local_bounded(&self, want: u32) -> Option<(u32, u32)> {
153        if want == 0 {
154            return None;
155        }
156        pop_bounded(&self.local_retry, want).or_else(|| self.claim(want))
157    }
158
159    /// Claims for a [`crate::LanePool`] lane: up to `want` samples of a requeued range
160    /// first (split like [`claim_local_bounded`](Self::claim_local_bounded)), falling
161    /// back to fresh work from [`claim`](Self::claim). Never touches the local-only
162    /// retry pile. `want == 0` returns `None` without touching anything.
163    pub fn claim_any(&self, want: u32) -> Option<(u32, u32)> {
164        if want == 0 {
165            return None;
166        }
167        pop_bounded(&self.requeued, want).or_else(|| self.claim(want))
168    }
169
170    /// Whether the SHARED pool has nothing left to claim -- `true` once every sample up
171    /// to `end` has been handed out to some engine. Says nothing about the retry
172    /// piles. A cheap, lock-free read for callers deciding whether waiting around could
173    /// still yield work.
174    pub fn shared_pool_exhausted(&self) -> bool {
175        self.next.load(Ordering::Relaxed) >= self.end
176    }
177
178    /// Whether nothing at all is claimable right now: the shared pool is exhausted and
179    /// both retry piles are empty. A range currently being traced may still come back
180    /// through [`requeue`](Self::requeue) or [`return_to_local`](Self::return_to_local)
181    /// if its lane fails.
182    pub fn fully_claimed(&self) -> bool {
183        self.shared_pool_exhausted()
184            && self
185                .requeued
186                .lock()
187                .unwrap_or_else(PoisonError::into_inner)
188                .is_empty()
189            && self
190                .local_retry
191                .lock()
192                .unwrap_or_else(PoisonError::into_inner)
193                .is_empty()
194    }
195
196    /// Returns `[start, start + count)` to the queue for GUARANTEED local processing
197    /// after a remote chunk failed to finish it -- never re-offered to remote. A `count`
198    /// of `0` is a no-op (nothing to retry) rather than an empty queue entry.
199    pub fn return_to_local(&self, start: u32, count: u32) {
200        if count == 0 {
201            return;
202        }
203        self.local_retry
204            .lock()
205            .unwrap_or_else(PoisonError::into_inner)
206            .push_back((start, count));
207    }
208
209    /// Returns `[start, start + count)` to the pile every [`claim_any`](Self::claim_any)
210    /// caller drains first, after a pool lane failed to finish it. A `count` of `0` is
211    /// a no-op.
212    pub fn requeue(&self, start: u32, count: u32) {
213        if count == 0 {
214            return;
215        }
216        self.requeued
217            .lock()
218            .unwrap_or_else(PoisonError::into_inner)
219            .push_back((start, count));
220    }
221}
222
223/// Pops at most `want` samples off the front of `pile`, pushing a longer entry's tail
224/// back to the front for the next caller. `want` must be non-zero.
225fn pop_bounded(pile: &Mutex<VecDeque<(u32, u32)>>, want: u32) -> Option<(u32, u32)> {
226    let mut pile = pile.lock().unwrap_or_else(PoisonError::into_inner);
227    let head = pile.pop_front();
228    let claimed = match head {
229        Some((start, count)) if count > want => {
230            pile.push_front((start + want, count - want));
231            Some((start, want))
232        }
233        other => other,
234    };
235    drop(pile);
236    claimed
237}