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}