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;