Skip to main content

datui_lib/loading/
measurements.rs

1//! What opening a dataset cost, measured rather than guessed: every figure is one datui
2//! produced itself (its own listing times, request counts, byte totals). Where Polars
3//! reads, datui cannot count, so nothing is recorded; hence the optional fields (see
4//! `docs/user-guide/dataset-info.md`). Written atomically by reading threads, read by
5//! the render without blocking.
6
7use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
8use std::time::Duration;
9
10/// One stretch of work, and what it cost.
11#[derive(Debug, Clone, Copy, PartialEq, Eq)]
12pub struct Cost {
13    /// How long it took, from the first request to the last answer.
14    pub took: Duration,
15    /// Data files found or footers read; `None` where the stretch never learned a count
16    /// (only the unreachable two-ends walk, see `schema_from_one_cloud_hive`, so it never
17    /// reports two files as the dataset's size).
18    pub files: Option<usize>,
19    /// Requests datui made and counted, with bytes returned. `None` on every listing (a
20    /// directory is read, not requested; a remote prefix pages inside the object store), so
21    /// zero never reads as "no data moved".
22    pub over_the_wire: Option<OverTheWire>,
23}
24
25/// Two stretches' wire figures as one: requests add; bytes add only when both weighed
26/// theirs.
27fn combine_wire(a: OverTheWire, b: OverTheWire) -> OverTheWire {
28    OverTheWire {
29        requests: a.requests + b.requests,
30        bytes: a.bytes.zip(b.bytes).map(|(x, y)| x + y),
31    }
32}
33
34/// Every stretch added up: how long the open has taken, and what it asked for.
35#[derive(Debug, Clone, Copy, PartialEq, Eq)]
36pub struct Total {
37    /// The stretches' times added together.
38    pub took: Duration,
39    /// The requests and bytes of the stretches that counted any, or `None` if none did.
40    pub over_the_wire: Option<OverTheWire>,
41}
42
43/// What datui asked for over a network, exactly.
44#[derive(Debug, Clone, Copy, PartialEq, Eq)]
45pub struct OverTheWire {
46    /// Requests datui issued.
47    pub requests: usize,
48    /// Bytes those requests returned, where counted; optional so a stretch can never claim
49    /// zero for unweighed requests.
50    pub bytes: Option<u64>,
51}
52
53/// A running tally of one kind of work. `ran` matters: an empty listing still took time
54/// and is a measurement.
55#[derive(Debug, Default)]
56struct Tally {
57    ran: AtomicBool,
58    nanos: AtomicU64,
59    files: AtomicUsize,
60    /// Counted as the requests are made, and still climbing while a pass runs.
61    live_requests: AtomicUsize,
62    live_bytes: AtomicU64,
63    /// Whether anything recorded here counted bytes at all.
64    counted_bytes: AtomicBool,
65    /// The counters at the last finished pass, which are shown: time and file count move
66    /// only at a pass's end, so live counters would pair one pass's requests with the
67    /// previous pass's time.
68    counted: AtomicBool,
69    requests: AtomicUsize,
70    bytes: AtomicU64,
71    /// Whether any stretch recorded here knew a count at all.
72    counted_files: AtomicBool,
73}
74
75impl Tally {
76    /// Record a finished stretch, added to the tally: one kind of work can span stretches
77    /// (a staged open's footer pass, a recount), so the count is footers read, not files.
78    /// Each dataset starts from an empty tally.
79    fn record(&self, took: Duration, files: Option<usize>, over_the_wire: bool) {
80        self.nanos.fetch_add(
81            took.as_nanos().min(u64::MAX as u128) as u64,
82            Ordering::Relaxed,
83        );
84        if let Some(files) = files {
85            self.files.fetch_add(files, Ordering::Relaxed);
86            self.counted_files.store(true, Ordering::Relaxed);
87        }
88        if over_the_wire {
89            // Read from the live counters here, so a request landing meanwhile is not lost.
90            // `fetch_max`: concurrent passes cannot finish together today, but the counters only
91            // climb, so the larger is right regardless.
92            self.requests.fetch_max(
93                self.live_requests.load(Ordering::Relaxed),
94                Ordering::Relaxed,
95            );
96            self.bytes
97                .fetch_max(self.live_bytes.load(Ordering::Relaxed), Ordering::Relaxed);
98            self.counted.store(true, Ordering::Relaxed);
99        }
100        // Last, and released, so a render that sees `ran` sees every field behind it.
101        self.ran.store(true, Ordering::Release);
102    }
103
104    /// Record a stretch in place of whatever this tally held, rather than adding to it.
105    fn replace(&self, took: Duration, files: Option<usize>) {
106        self.nanos.store(
107            took.as_nanos().min(u64::MAX as u128) as u64,
108            Ordering::Relaxed,
109        );
110        match files {
111            Some(files) => {
112                self.files.store(files, Ordering::Relaxed);
113                self.counted_files.store(true, Ordering::Relaxed);
114            }
115            None => self.counted_files.store(false, Ordering::Relaxed),
116        }
117        self.ran.store(true, Ordering::Release);
118    }
119
120    /// One more request, and the bytes it returned. Called from the threads doing the
121    /// reading, once per request, before the stretch is recorded.
122    fn request(&self, bytes: u64) {
123        self.live_requests.fetch_add(1, Ordering::Relaxed);
124        self.live_bytes.fetch_add(bytes, Ordering::Relaxed);
125        self.counted_bytes.store(true, Ordering::Relaxed);
126    }
127
128    /// A whole pass's worth of requests at once, for work counted against a meter of
129    /// its own and then folded in here.
130    fn add_requests(&self, wire: OverTheWire) {
131        self.live_requests
132            .fetch_add(wire.requests, Ordering::Relaxed);
133        if let Some(bytes) = wire.bytes {
134            self.live_bytes.fetch_add(bytes, Ordering::Relaxed);
135            self.counted_bytes.store(true, Ordering::Relaxed);
136        }
137    }
138
139    fn cost(&self) -> Option<Cost> {
140        if !self.ran.load(Ordering::Acquire) {
141            return None;
142        }
143        Some(Cost {
144            took: Duration::from_nanos(self.nanos.load(Ordering::Relaxed)),
145            files: self
146                .counted_files
147                .load(Ordering::Relaxed)
148                .then(|| self.files.load(Ordering::Relaxed)),
149            over_the_wire: self.counted.load(Ordering::Relaxed).then(|| OverTheWire {
150                requests: self.requests.load(Ordering::Relaxed),
151                bytes: self
152                    .counted_bytes
153                    .load(Ordering::Relaxed)
154                    .then(|| self.bytes.load(Ordering::Relaxed)),
155            }),
156        })
157    }
158}
159
160/// What an open reports as it goes: footer progress and cost so far, handed down every
161/// route together and owned by the open that made them (see [`Meter`]).
162#[derive(Debug, Clone, Default)]
163pub struct OpenReport {
164    /// How far the footer pass has got, for the loading screen.
165    pub progress: std::sync::Arc<crate::formats::schema_union::FooterProgress>,
166    /// What the work has cost, for the Info panel.
167    pub meter: std::sync::Arc<Meter>,
168    /// Where to find what a previous open learned and leave what this one learns; `None`
169    /// for routes with nowhere to keep it, and for tests.
170    pub remembered: Option<crate::cache::CacheManager>,
171    /// Where what is remembered for the home screen is written, off the open's path:
172    /// the app's, which the home listing settles before it reads.
173    pub(crate) writes: crate::app::background::CacheWrites,
174}
175
176/// What the open on screen cost, as measured; shared with reading threads (`Arc`,
177/// written through `&self`). Never cleared: each open builds a new one, since an
178/// abandoned load's reads still run and would add to the next dataset's figures.
179#[derive(Debug, Default)]
180pub struct Meter {
181    listing: Tally,
182    footers: Tally,
183    /// The most recent page of rows, replacing rather than adding: this one says what
184    /// the page on screen cost, not what every page since the open came to.
185    last_page: Tally,
186    /// Whether a pass to settle the row count has already been counted.
187    counted_rows: AtomicBool,
188}
189
190impl Meter {
191    /// Finding the files took `took` and found `files`. `over_the_wire` is always false
192    /// today (local walks make no requests; remote listings page inside the object store),
193    /// kept so both stretches record alike.
194    pub fn listed(&self, took: Duration, files: Option<usize>, over_the_wire: bool) {
195        self.listing.record(took, files, over_the_wire);
196    }
197
198    /// A pass over `files` footers took `took`. `over_the_wire` publishes requests counted
199    /// by [`Self::footer_request`]; a disk pass passes false (no requests, not zero).
200    pub fn read_footers(&self, took: Duration, files: Option<usize>, over_the_wire: bool) {
201        self.footers.record(took, files, over_the_wire);
202    }
203
204    /// A count pass over `files` footers took `took`: recorded once only (the return says
205    /// whether), since counts rerun on every invalidation and would inflate the open's
206    /// cost.
207    pub fn counted_rows(
208        &self,
209        took: Duration,
210        files: Option<usize>,
211        wire: Option<OverTheWire>,
212    ) -> bool {
213        // Only for an open this meter measured (a route reporting nothing, like Polars-handed
214        // directories, must not gain a section of after-the-fact work), checked before the
215        // one-shot.
216        if self.listing.cost().is_none() {
217            return false;
218        }
219        if self
220            .counted_rows
221            .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
222            .is_err()
223        {
224            return false;
225        }
226        // Folded in only after winning the one-shot, so a declined pass leaves counters alone.
227        if let Some(w) = wire {
228            self.footers.add_requests(w);
229        }
230        self.footers.record(took, files, wire.is_some());
231        true
232    }
233
234    /// The page on screen took `took` and read `files` files (`None` when Polars chose).
235    /// Replaces rather than adds: the page being viewed. No byte figure: Polars does not
236    /// report what it fetched, and row-group sizes would overstate it.
237    pub fn read_page(&self, took: Duration, files: Option<usize>) {
238        self.last_page.replace(took, files);
239    }
240
241    /// What the page now on screen cost, or `None` before one has been read.
242    pub fn last_page(&self) -> Option<Cost> {
243        self.last_page.cost()
244    }
245
246    /// One footer request, and the bytes it returned. Counted as the reads happen; what
247    /// they come to is published when the pass ends.
248    pub fn footer_request(&self, bytes: u64) {
249        self.footers.request(bytes);
250    }
251
252    /// What finding the files cost, or `None` if no listing has been measured.
253    pub fn listing(&self) -> Option<Cost> {
254        self.listing.cost()
255    }
256
257    /// What reading the footers cost, or `None` if no footer pass has been measured.
258    pub fn footers(&self) -> Option<Cost> {
259        self.footers.cost()
260    }
261
262    /// Everything measured, added; `None` until something is. No file count (listing counts
263    /// files, footer passes count footers). Wire figures sum the stretches that counted
264    /// them, `None` when none did.
265    pub fn total(&self) -> Option<Total> {
266        // Listing and footers only. The page is not part of opening the dataset — it is
267        // what looking at one costs, and it changes every time the view moves.
268        let parts: Vec<Cost> = [self.listing(), self.footers()]
269            .into_iter()
270            .flatten()
271            .collect();
272        // Two stretches or none: one stretch's total would repeat it under two labels.
273        if parts.len() < 2 {
274            return None;
275        }
276        Some(Total {
277            took: parts.iter().map(|p| p.took).sum(),
278            over_the_wire: parts
279                .iter()
280                .filter_map(|p| p.over_the_wire)
281                .reduce(combine_wire),
282        })
283    }
284}
285
286/// How long the event loop's own work takes: each frame drawn and each event handled,
287/// for the debug overlay and the log.
288#[derive(Debug, Default)]
289pub struct LoopTimes {
290    pub frames: Durations,
291    pub handlers: Durations,
292}
293
294impl LoopTimes {
295    /// One frame drawn. Every [`Durations::WINDOW`] frames the log hears both
296    /// summaries at debug level.
297    pub fn frame(&mut self, took: Duration) {
298        self.frames.record(took);
299        if self.frames.count.is_multiple_of(Durations::WINDOW as u64) {
300            log::debug!(
301                target: "datui",
302                "frames {}; handlers {}",
303                self.frames.summary(),
304                self.handlers.summary()
305            );
306        }
307    }
308
309    /// One event handled, its key included.
310    pub fn handler(&mut self, took: Duration) {
311        self.handlers.record(took);
312    }
313}
314
315/// The most recent durations of one kind of work, for its median, 99th percentile and
316/// worst.
317#[derive(Debug, Default)]
318pub struct Durations {
319    recent: std::collections::VecDeque<Duration>,
320    count: u64,
321}
322
323impl Durations {
324    /// How many of the latest durations the figures are over.
325    pub const WINDOW: usize = 240;
326
327    pub fn record(&mut self, took: Duration) {
328        if self.recent.len() == Self::WINDOW {
329            self.recent.pop_front();
330        }
331        self.recent.push_back(took);
332        self.count += 1;
333    }
334
335    /// How many were ever recorded.
336    #[cfg(test)]
337    pub fn count(&self) -> u64 {
338        self.count
339    }
340
341    /// The duration `q` of the way up the recent ones (0.5 the median, 1.0 the worst).
342    pub fn quantile(&self, q: f64) -> Option<Duration> {
343        let mut sorted: Vec<Duration> = self.recent.iter().copied().collect();
344        sorted.sort_unstable();
345        let last = sorted.len().checked_sub(1)?;
346        sorted.get(((last as f64) * q).round() as usize).copied()
347    }
348
349    /// `p50 1.2ms p99 3.4ms max 5.0ms`, or `-` before the first.
350    pub fn summary(&self) -> String {
351        let ms = |d: Duration| format!("{:.1}ms", d.as_secs_f64() * 1000.0);
352        match (self.quantile(0.5), self.quantile(0.99), self.quantile(1.0)) {
353            (Some(p50), Some(p99), Some(max)) => {
354                format!("p50 {} p99 {} max {}", ms(p50), ms(p99), ms(max))
355            }
356            _ => "-".to_string(),
357        }
358    }
359}
360
361#[cfg(test)]
362mod tests;