Skip to main content

indicatrix_dispatch/rate/
mod.rs

1//! [`RateModel`]: one lane's throughput estimate, and [`ChunkPolicy`]: how a chunk is
2//! sized from it.
3//!
4//! Generalised from the desktop export's remote lane (`remote_chunk_samples`,
5//! `remote_marginal_rate`, `calibrate_remote_rate`, the per-video rate carry) and the
6//! hybrid split's 0.7-old/0.3-new moving average:
7//!
8//! 1. **Initial guess.** A lane that has never been measured starts from a guess and
9//!    is handed a small calibration chunk ([`ChunkPolicy::calibration_samples`]) rather
10//!    than a chunk sized from the guess, so a wildly wrong guess costs one short chunk.
11//! 2. **First-chunk calibration.** The first valid measurement REPLACES the guess
12//!    outright; a guess is not evidence worth averaging with.
13//! 3. **Blending.** Every later measurement is folded in with an exponential moving
14//!    average ([`DEFAULT_SMOOTHING`] weight on the new sample), so thermal drift and load
15//!    changes are followed without one noisy chunk swinging the next chunk's size.
16//!
17//! A measurement is the lane's own reported marginal rate when it has one (a remote
18//! worker's steady-state rate, excluding connection setup and upload), otherwise
19//! `done / wall time` of the whole chunk ([`marginal_rate`]).
20
21use std::time::Duration;
22
23#[cfg(test)]
24mod tests;
25
26/// Weight of a new measurement in the moving average: 0.3 new, 0.7 old, the same
27/// blend the hybrid CPU/GPU split uses.
28pub const DEFAULT_SMOOTHING: f64 = 0.3;
29
30/// The worker's hidden per-request cap (`MAX_SAMPLES_PER_REQUEST` in
31/// `indicatrix-worker`'s request validation): a chunk sent to a worker must never be
32/// larger, or the worker rejects the whole request.
33pub const DEFAULT_MAX_CHUNK_SAMPLES: u32 = 65_536;
34
35/// Throughput in samples per second: `delta_samples` traced over `elapsed`. `None`
36/// when nothing was traced or the span is zero, so a degenerate measurement never
37/// reaches a [`RateModel`].
38#[must_use]
39pub fn marginal_rate(delta_samples: u32, elapsed: Duration) -> Option<f64> {
40    if delta_samples == 0 {
41        return None;
42    }
43    let secs = elapsed.as_secs_f64();
44    if secs <= 0.0 {
45        return None;
46    }
47    Some(f64::from(delta_samples) / secs)
48}
49
50/// How chunks are sized for a [`crate::LanePool`] run.
51#[derive(Debug, Clone, Copy, PartialEq, Eq)]
52pub struct ChunkPolicy {
53    /// Target wall-clock duration of one chunk. Long enough that per-request overhead
54    /// (handshake, upload, the final transfer) is negligible against it; short enough
55    /// that a slow lane never holds the tail of the image hostage while every other
56    /// lane idles.
57    pub target: Duration,
58    /// Floor for a rate-sized chunk: a zero, tiny or non-finite rate must not collapse
59    /// the chunk to something overhead dominates.
60    pub min_samples: u32,
61    /// Ceiling for any chunk ([`DEFAULT_MAX_CHUNK_SAMPLES`] for remote workers).
62    pub max_samples: u32,
63    /// Size of an uncalibrated lane's first chunk (see the module doc).
64    pub calibration_samples: u32,
65}
66
67impl ChunkPolicy {
68    /// The still export and tilt video: about 22 s per chunk, a 32-sample floor
69    /// (below that a remote round trip costs more than it saves), an 8-sample
70    /// calibration probe -- the desktop export's existing numbers.
71    pub const EXPORT: Self = Self {
72        target: Duration::from_secs(22),
73        min_samples: 32,
74        max_samples: DEFAULT_MAX_CHUNK_SAMPLES,
75        calibration_samples: 8,
76    };
77
78    /// Interactive (live-view) requests: about 1.5 s per chunk so progress arrives
79    /// often, single-sample floor and calibration.
80    pub const INTERACTIVE: Self = Self {
81        target: Duration::from_millis(1500),
82        min_samples: 1,
83        max_samples: DEFAULT_MAX_CHUNK_SAMPLES,
84        calibration_samples: 1,
85    };
86
87    /// Every chunk (calibration included) is exactly `samples` long, whatever the rate.
88    /// The chunk partition of a run is then independent of timing (see
89    /// [`crate::Merger`]'s determinism notes). `0` is treated as `1`.
90    #[must_use]
91    pub const fn fixed(samples: u32) -> Self {
92        let samples = if samples == 0 { 1 } else { samples };
93        Self {
94            target: Duration::ZERO,
95            min_samples: samples,
96            max_samples: samples,
97            calibration_samples: samples,
98        }
99    }
100
101    /// The chunk size for a lane measured at `rate` samples per second: `rate *
102    /// target`, rounded, clamped to `[min_samples, max_samples]`; `min_samples` for a
103    /// non-finite or non-positive rate. Never `0`.
104    #[must_use]
105    pub fn samples_for_rate(&self, rate: f64) -> u32 {
106        let (min, max) = self.bounds();
107        let target = rate * self.target.as_secs_f64();
108        if !target.is_finite() || target < f64::from(min) {
109            return min;
110        }
111        if target >= f64::from(max) {
112            return max;
113        }
114        (target.round() as u32).clamp(min, max)
115    }
116
117    /// `(min, max)` with both at least `1` and `min <= max`.
118    fn bounds(&self) -> (u32, u32) {
119        let max = self.max_samples.max(1);
120        (self.min_samples.clamp(1, max), max)
121    }
122
123    /// The first chunk of an uncalibrated lane, clamped to `[1, max_samples]`.
124    #[must_use]
125    pub fn first_chunk_samples(&self) -> u32 {
126        self.calibration_samples.clamp(1, self.bounds().1)
127    }
128}
129
130/// One lane's throughput estimate in samples per second. See the module doc.
131#[derive(Debug, Clone, Copy, PartialEq)]
132pub struct RateModel {
133    /// Used until the first measurement.
134    guess: f64,
135    /// `None` until calibrated.
136    estimate: Option<f64>,
137    /// Weight of a new measurement, in `(0, 1]`.
138    smoothing: f64,
139    /// Measurements folded in so far.
140    observations: u32,
141}
142
143impl RateModel {
144    /// An uncalibrated model starting from `initial_guess` samples per second. The
145    /// guess only matters for reporting ([`Self::rate`]); an uncalibrated lane's first
146    /// chunk is [`ChunkPolicy::first_chunk_samples`] long regardless.
147    #[must_use]
148    pub const fn new(initial_guess: f64) -> Self {
149        Self {
150            guess: initial_guess,
151            estimate: None,
152            smoothing: DEFAULT_SMOOTHING,
153            observations: 0,
154        }
155    }
156
157    /// A model already calibrated at `rate` (carried over from an earlier epoch of the
158    /// same lane, e.g. the previous tilt-video frame), so no calibration chunk is spent.
159    /// A non-finite or non-positive `rate` yields an uncalibrated model instead.
160    #[must_use]
161    pub fn calibrated(rate: f64) -> Self {
162        let mut model = Self::new(rate);
163        if valid_rate(rate) {
164            model.estimate = Some(rate);
165        }
166        model
167    }
168
169    /// Replaces the moving-average weight of a new measurement; clamped to
170    /// `[0.01, 1.0]` (`1.0` means "always trust the latest chunk").
171    #[must_use]
172    pub const fn with_smoothing(mut self, weight_of_new: f64) -> Self {
173        self.smoothing = if weight_of_new.is_finite() {
174            weight_of_new.clamp(0.01, 1.0)
175        } else {
176            DEFAULT_SMOOTHING
177        };
178        self
179    }
180
181    /// Whether at least one real measurement has replaced the guess.
182    #[must_use]
183    pub const fn is_calibrated(&self) -> bool {
184        self.estimate.is_some()
185    }
186
187    /// The current estimate, or the initial guess while uncalibrated.
188    #[must_use]
189    pub fn rate(&self) -> f64 {
190        self.estimate.unwrap_or(self.guess)
191    }
192
193    /// The calibrated estimate only (`None` while uncalibrated) -- what a caller
194    /// carries into the next epoch with [`Self::calibrated`].
195    #[must_use]
196    pub const fn estimate(&self) -> Option<f64> {
197        self.estimate
198    }
199
200    /// Measurements folded in so far.
201    #[must_use]
202    pub const fn observations(&self) -> u32 {
203        self.observations
204    }
205
206    /// Folds in one finished chunk: `done` samples in `elapsed` wall time, with the
207    /// lane's own `reported` marginal rate preferred when it is valid. Returns the
208    /// measurement used, or `None` (model unchanged) when there was nothing to measure.
209    pub fn observe(&mut self, done: u32, elapsed: Duration, reported: Option<f64>) -> Option<f64> {
210        let measured = reported
211            .filter(|&rate| valid_rate(rate))
212            .or_else(|| marginal_rate(done, elapsed))
213            .filter(|&rate| valid_rate(rate))?;
214        self.estimate = Some(
215            self.estimate
216                .map_or(measured, |old| self.smoothing.mul_add(measured - old, old)),
217        );
218        self.observations += 1;
219        Some(measured)
220    }
221
222    /// The size of this lane's next chunk under `policy`: the calibration chunk while
223    /// uncalibrated, otherwise [`ChunkPolicy::samples_for_rate`] of the estimate.
224    #[must_use]
225    pub fn chunk_samples(&self, policy: &ChunkPolicy) -> u32 {
226        self.estimate.map_or_else(
227            || policy.first_chunk_samples(),
228            |rate| policy.samples_for_rate(rate),
229        )
230    }
231}
232
233/// A usable throughput: finite and strictly positive.
234fn valid_rate(rate: f64) -> bool {
235    rate.is_finite() && rate > 0.0
236}