datui_lib/app/jobs.rs
1//! Background operations, and their one owner.
2//!
3//! [`Jobs`] starts every general background operation and keeps one record per job
4//! until its outcome is handled: the [`Job`], its [`Ticket`], whether a generation
5//! bump would strand it, whether the user waits on it, and whether its answer is
6//! still wanted. The app keeps no job markers of its own.
7//!
8//! - **One outcome.** A worker returns its [`Answer`] or an error; a panic is
9//! caught. The outcome goes into the record and [`AppEvent::JobEnded`] says so. A
10//! [`Started`] job dropped without running ends as failed.
11//! - **Acceptance and release in one step.** [`Jobs::end`] hands over the outcome
12//! and removes the record in one call, so whatever the app starts while handling
13//! the answer takes the generation before anything else looks: it never reads
14//! free between two phases of one errand, and no finished job still holds it.
15//! - **Supersession.** Advancing the generation supersedes the jobs it scopes;
16//! [`Jobs::supersede`] cancels others. A superseded job stops holding the
17//! generation at once; its worker runs on and its stale outcome is dropped.
18//! - **Keys.** A waited-on job holds the keys and its footer line until it ends or
19//! is superseded; [`Jobs::quiet`] releases them and [`Jobs::wait_on`] takes them
20//! for a running job. A page asked for while the generation is held is owed
21//! ([`Jobs::owe`]): it holds the keys, not the generation, until run.
22//! - **Holds.** A continuation the pump has not dispatched, or a download waiting
23//! on the user, also holds the generation ([`Hold`]).
24//!
25//! Not owned here: the row count (`OwedCount`) and the home screen's workers, each
26//! keyed by its own marker. `reread_after_the_footers_joined` sends its jump straight
27//! to the channel, safe only because its callers checked nothing would be stranded.
28//!
29//! A worker that never returns (a `hard` NFS mount, a wedged object-store read)
30//! keeps its keys and lease until superseded.
31
32use std::collections::HashMap;
33use std::path::PathBuf;
34use std::sync::mpsc::Sender;
35use std::sync::{Arc, Mutex};
36use std::time::Instant;
37
38use polars::prelude::DataFrame;
39
40use crate::loading::{LoadAnswer, LoadId};
41use crate::{AppEvent, OpenOptions, logging};
42
43/// Which kind of operation a [`Ticket`] names.
44#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
45pub enum JobKind {
46 Load,
47 OpenNamed,
48 LookAtDirectory,
49 Classify,
50 Rows,
51 Analysis,
52 SampleRows,
53 SampleDraw,
54 Pivot,
55 ViewPivot,
56 ReshapePreview,
57 DrillRow,
58 InspectRow,
59 InspectJson,
60 InspectPretty,
61 InspectUnpack,
62 OpenValue,
63 Export,
64 Copy,
65 QualityReport,
66 FileFacts,
67 ChartExport,
68 ChartPrepare,
69 Find,
70 ValueCounts,
71 HexOpen,
72 HexFind,
73 UnfitCount,
74 FootersJoin,
75 JournalDetail,
76 IndexLines,
77}
78
79/// One started operation. Issued when it starts, carried by its worker, and handed
80/// back with its outcome.
81#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
82pub struct Ticket {
83 id: u64,
84 generation: u64,
85 kind: JobKind,
86}
87
88impl Ticket {
89 pub fn kind(self) -> JobKind {
90 self.kind
91 }
92}
93
94/// A background operation. Its fields are what the app needs of it; the record
95/// holding it is its only marker.
96#[derive(Debug, Clone)]
97pub(crate) enum Job {
98 /// A phase of the open it names, before its first rows. Judged by the open
99 /// ([`crate::loading::Loader`]), not the generation.
100 Load(LoadId),
101 /// The first phase of an open: whether the named paths exist and which is a
102 /// directory.
103 OpenNamed(LoadId),
104 /// The look at a directory named on the command line, before `load` opens it.
105 LookAtDirectory { load: LoadId, path: PathBuf },
106 /// A look at a path chosen on the home screen.
107 Classify(Classify),
108 /// The table's rows: a page the table waits on, or a load-ahead.
109 Rows(crate::InflightCollect),
110 /// A page asked for while the generation was held: no worker yet; read once
111 /// nothing would be stranded.
112 OwedRows { dataset: u64, status: String },
113 /// An Analysis tool's computation.
114 Analysis(AnalysisRun),
115 /// The sample, or the rows behind a finding, read to show as a table.
116 SampleRows,
117 /// The view's sample, drawn into memory as the table shows it.
118 SampleDraw(Box<SampleDraw>),
119 /// A pivot from the Pivot & Melt builder.
120 Pivot,
121 /// A view's pivot, read before its rows: the view, and why it is applied.
122 ViewPivot(Box<(crate::view::SavedView, crate::view::view_apply::Applying)>),
123 /// The Pivot & Melt preview, request `token` of opening `epoch`; judged by those,
124 /// not the generation. Nobody waits on it.
125 ReshapePreview { epoch: u64, token: u64 },
126 /// The group row Enter drills into, when the buffer did not hold it.
127 DrillRow,
128 /// The group a reopened dataset was drilled into, found again by its keys.
129 Regroup(Box<crate::table::DrillPlace>),
130 /// The inspector's fields of row `row` of frame `frame`, not in the buffer.
131 InspectRow { frame: u64, row: usize },
132 /// The inspector's text parsed as JSON to drill into; answers
133 /// [`crate::inspector::inspector_drill::JsonWait`].
134 InspectJson { token: u64 },
135 /// The inspector's long JSON indented for its JSON view; answers
136 /// [`crate::inspector::inspector_modal::Pretty`].
137 InspectPretty { token: u64 },
138 /// The inspector's gzip or zstd bytes decompressed for its Text view; answers
139 /// [`crate::inspector::inspector_modal::Unpack`].
140 InspectUnpack { token: u64 },
141 /// The inspector's value written to a file for another program to open.
142 OpenValue,
143 /// An export, from plan to committed file.
144 Export,
145 /// Collecting and formatting the view for a copy.
146 Copy,
147 /// Writing the Data Quality report.
148 QualityReport,
149 /// The open file's size and footer for the Info panel, for `dataset_generation`
150 /// `dataset` (the generation does not tell datasets apart).
151 FileFacts { dataset: u64 },
152 /// Writing a chart. The path and format reopen the form on a failure.
153 ChartExport {
154 path: PathBuf,
155 format: crate::chart::chart_export::ChartExportFormat,
156 },
157 /// A chart's data for the selection on screen; judged by itself (see
158 /// [`ChartPrep`]).
159 ChartPrepare(Box<ChartPrep>),
160 /// A find reading the view for its next match.
161 Find(crate::find::FindRun),
162 /// Counting a column's values for the Value Counts screen.
163 ValueCounts,
164 /// Mapping a file for the hex view, opened from `origin`.
165 HexOpen {
166 origin: crate::app::hex_view::Origin,
167 fallback: bool,
168 record_size: Option<usize>,
169 },
170 /// A find reading the hex view's file.
171 HexFind(crate::app::hex_view::HexFindRun),
172 /// Counting values the read's column types made null, for the Notes; judged by
173 /// `dataset`. `version` is the view's column changes counted; `None` for the read's
174 /// types.
175 UnfitCount { dataset: u64, version: Option<u64> },
176 /// The pass reading every footer behind a staged open; judged by `dataset`, since
177 /// it outlives several collects.
178 FootersJoin { dataset: u64 },
179 /// A piped journal's Info tab re-read once it ended, for `dataset`.
180 JournalDetail { dataset: u64 },
181 /// Indexing the rest of a text file's lines behind its first rows, for the
182 /// `dataset_generation` it was started for. Stopped by its flag, it fails.
183 IndexLines { dataset: u64 },
184}
185
186/// A look at a path chosen on the home screen. Home keys act while busy, so a
187/// second Enter supersedes the first.
188#[derive(Debug, Clone)]
189pub(crate) struct Classify {
190 pub(crate) path: PathBuf,
191 /// Where home was browsing when the look was asked; an answer for a place left
192 /// opens nothing. Not `HomeApp::generation`: listings rebuild for unrelated reasons
193 /// (another root's probe answering).
194 pub(crate) browsing: Option<PathBuf>,
195 /// The text typed at `~`, when the path came from there rather than a listed row.
196 pub(crate) typed: Option<String>,
197}
198
199/// A chart's data being prepared. One at a time, so a burst of selection changes
200/// does not fan out into a collect per column; the newest selection is prepared
201/// next. Its answer is kept only while current and only for its own dataset.
202#[derive(Debug, Clone)]
203pub(crate) struct ChartPrep {
204 pub(crate) request: crate::ChartRequest,
205 /// `len_generation` of the dataset read, so an answer is never installed into
206 /// another with the same column names.
207 pub(crate) dataset: Option<u64>,
208 /// Set when the selection or its view moves on. A streamed count or group-by stops
209 /// at its next batch; a sampled read runs out.
210 pub(crate) cancel: Arc<std::sync::atomic::AtomicBool>,
211}
212
213/// A view's sample being drawn. Judged by its rows, not the generation: paging,
214/// finds and inspection move the generation meanwhile. It lands in the view whose
215/// sample holds `rows`.
216#[derive(Clone)]
217pub(crate) struct SampleDraw {
218 pub(crate) sample: crate::analysis::sampling::Sample,
219 pub(crate) rows: Arc<crate::analysis::table_sample::SampleRows>,
220 /// Stops the draw; the rows so far stay.
221 pub(crate) watch: crate::analysis::sampling::ReadWatch,
222 /// Drawn from the view's query or filters rather than the source under them.
223 pub(crate) through: bool,
224 /// The steps laid on the built sample: query, filters, sort and columns of the view
225 /// drawn from, or being applied.
226 pub(crate) replay: Option<crate::view::ViewSettings>,
227 /// Analysis asked for it, and runs its tool once it is drawn.
228 pub(crate) then_analyze: bool,
229 /// How a random sample of a stream is drawn, decided before it starts.
230 pub(crate) path: Option<crate::analysis::table_sample::DrawPath>,
231 /// What the draw is remembered by, for drawing it the same way again.
232 pub(crate) path_key: String,
233 /// The drawn rows' columns, once cut to scope. The view becomes the sample's when
234 /// its first rows land; until then, and if none come, the old view stays.
235 pub(crate) schema: Option<polars::prelude::SchemaRef>,
236}
237
238impl std::fmt::Debug for SampleDraw {
239 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
240 f.debug_struct("SampleDraw")
241 .field("sample", &self.sample)
242 .field("through", &self.through)
243 .field("then_analyze", &self.then_analyze)
244 .finish_non_exhaustive()
245 }
246}
247
248/// An Analysis tool's run.
249#[derive(Debug, Clone, Default)]
250pub(crate) struct AnalysisRun {
251 /// A Data Quality run's watch: how it is told to stop, and what it has read.
252 pub(crate) watch: Option<crate::analysis::data_quality::QualityWatch>,
253 /// Cancelled during a read that cannot stop part way: it runs to its end.
254 pub(crate) runs_out: bool,
255}
256
257impl Job {
258 pub(crate) fn kind(&self) -> JobKind {
259 match self {
260 Job::Load(_) => JobKind::Load,
261 Job::OpenNamed(_) => JobKind::OpenNamed,
262 Job::LookAtDirectory { .. } => JobKind::LookAtDirectory,
263 Job::Classify(_) => JobKind::Classify,
264 Job::Rows(_) | Job::OwedRows { .. } => JobKind::Rows,
265 Job::Analysis(_) => JobKind::Analysis,
266 Job::SampleRows => JobKind::SampleRows,
267 Job::SampleDraw(_) => JobKind::SampleDraw,
268 Job::Pivot => JobKind::Pivot,
269 Job::ViewPivot(_) => JobKind::ViewPivot,
270 Job::ReshapePreview { .. } => JobKind::ReshapePreview,
271 Job::DrillRow | Job::Regroup(_) => JobKind::DrillRow,
272 Job::InspectRow { .. } => JobKind::InspectRow,
273 Job::InspectJson { .. } => JobKind::InspectJson,
274 Job::InspectPretty { .. } => JobKind::InspectPretty,
275 Job::InspectUnpack { .. } => JobKind::InspectUnpack,
276 Job::OpenValue => JobKind::OpenValue,
277 Job::Export => JobKind::Export,
278 Job::Copy => JobKind::Copy,
279 Job::QualityReport => JobKind::QualityReport,
280 Job::FileFacts { .. } => JobKind::FileFacts,
281 Job::ChartExport { .. } => JobKind::ChartExport,
282 Job::ChartPrepare(_) => JobKind::ChartPrepare,
283 Job::Find(_) => JobKind::Find,
284 Job::ValueCounts => JobKind::ValueCounts,
285 Job::HexOpen { .. } => JobKind::HexOpen,
286 Job::HexFind(_) => JobKind::HexFind,
287 Job::UnfitCount { .. } => JobKind::UnfitCount,
288 Job::FootersJoin { .. } => JobKind::FootersJoin,
289 Job::JournalDetail { .. } => JobKind::JournalDetail,
290 Job::IndexLines { .. } => JobKind::IndexLines,
291 }
292 }
293
294 /// Whether this job reads the open dataset's files, so that its failure saying a
295 /// file is gone since the open is one a reopen fixes and the Reopen question is
296 /// asked. No for the open's own phases (the loader says those), home's looks, the
297 /// raw bytes the hex view maps, work on a value or a file already read, writes of
298 /// what is already in hand, side reads for the Info panel and the Notes, and the
299 /// Pivot & Melt preview, which runs while typing. No wildcard: a new job decides.
300 pub(crate) fn reads_dataset(&self) -> bool {
301 match self {
302 Job::Rows(_)
303 | Job::OwedRows { .. }
304 | Job::Analysis(_)
305 | Job::SampleRows
306 | Job::SampleDraw(_)
307 | Job::Pivot
308 | Job::ViewPivot(_)
309 | Job::DrillRow
310 | Job::Regroup(_)
311 | Job::InspectRow { .. }
312 | Job::Export
313 | Job::Copy
314 | Job::ChartPrepare(_)
315 | Job::Find(_)
316 | Job::ValueCounts => true,
317 Job::Load(_)
318 | Job::OpenNamed(_)
319 | Job::LookAtDirectory { .. }
320 | Job::Classify(_)
321 | Job::ReshapePreview { .. }
322 | Job::InspectJson { .. }
323 | Job::InspectPretty { .. }
324 | Job::InspectUnpack { .. }
325 | Job::OpenValue
326 | Job::QualityReport
327 | Job::FileFacts { .. }
328 | Job::ChartExport { .. }
329 | Job::HexOpen { .. }
330 | Job::HexFind(_)
331 | Job::UnfitCount { .. }
332 | Job::FootersJoin { .. }
333 | Job::JournalDetail { .. }
334 | Job::IndexLines { .. } => false,
335 }
336 }
337
338 /// The open this job is a phase of, if it is one.
339 pub(crate) fn load(&self) -> Option<LoadId> {
340 match self {
341 Job::Load(load) | Job::OpenNamed(load) | Job::LookAtDirectory { load, .. } => {
342 Some(*load)
343 }
344 _ => None,
345 }
346 }
347
348 /// Whether a generation bump waits for this job's answer: stale answers are never
349 /// asked for again, so most jobs are leased. Not leased: buffer rows (re-asked, and
350 /// leasing would chain pages); looks at named paths (meant to be dropped on
351 /// Ctrl+O); file facts, journal detail and the footer pass (judged by dataset); the
352 /// Pivot & Melt preview (by request); chart data (by dataset, re-asked when
353 /// missing). An owed page has nothing running.
354 fn leased(&self) -> bool {
355 !matches!(
356 self,
357 Job::Rows(_)
358 | Job::SampleDraw(_)
359 | Job::OwedRows { .. }
360 | Job::OpenNamed(_)
361 | Job::LookAtDirectory { .. }
362 | Job::FileFacts { .. }
363 | Job::UnfitCount { .. }
364 | Job::FootersJoin { .. }
365 | Job::JournalDetail { .. }
366 | Job::IndexLines { .. }
367 | Job::ReshapePreview { .. }
368 | Job::ChartPrepare(_)
369 )
370 }
371
372 /// Whether advancing the generation makes this answer stale. Excluded jobs are put
373 /// down by their own owner: file facts (dataset), chart export and data (chart
374 /// view), an owed page (dataset), a preview (builder request).
375 fn follows_the_generation(&self) -> bool {
376 !matches!(
377 self,
378 Job::FileFacts { .. }
379 | Job::SampleDraw(_)
380 | Job::UnfitCount { .. }
381 | Job::FootersJoin { .. }
382 | Job::JournalDetail { .. }
383 | Job::IndexLines { .. }
384 | Job::ChartExport { .. }
385 | Job::ChartPrepare(_)
386 | Job::OwedRows { .. }
387 | Job::ReshapePreview { .. }
388 )
389 }
390}
391
392/// What a job's worker sends back when it succeeds. Each belongs to one [`Job`].
393pub(crate) enum Answer {
394 /// [`Job::Load`]: what the phase found. A download dropped unhandled removes its
395 /// file.
396 Load(Box<LoadAnswer>),
397 /// [`Job::OpenNamed`]: the paths are there; `directory` is one to look at first.
398 NamedPaths {
399 paths: Vec<PathBuf>,
400 options: Box<OpenOptions>,
401 directory: Option<PathBuf>,
402 },
403 /// [`Job::OpenNamed`]: a path named is not there.
404 NamedPathMissing(PathBuf),
405 /// [`Job::LookAtDirectory`]: what the look found; `holds` (a cloud listing) picks
406 /// the reader.
407 LookedAt {
408 kind: crate::home::discover::EntryKind,
409 holds: Option<Box<crate::home::discover::Holds>>,
410 options: Box<OpenOptions>,
411 },
412 /// [`Job::Classify`]: what the path turned out to be; `None` is a path that is not
413 /// there.
414 Kind(Option<crate::home::discover::EntryKind>),
415 /// [`Job::Rows`]: the rows read.
416 Rows(crate::table::CollectResult),
417 /// [`Job::Rows`]: the read failed. `conversion` is a value that would not convert,
418 /// for the SQL prompt to word itself.
419 RowsFailed {
420 message: String,
421 conversion: Option<Box<crate::error_display::ConversionFailure>>,
422 },
423 /// [`Job::Analysis`]: a Describe, Distributions or Correlations result, and
424 /// where on the modal it goes.
425 Analysis(
426 fn(
427 &mut crate::analysis::analysis_modal::AnalysisModal,
428 crate::analysis::statistics::AnalysisResults,
429 ),
430 crate::analysis::statistics::AnalysisResults,
431 ),
432 /// [`Job::Analysis`]: a Data Quality report, the rows a sampled run read, and the
433 /// plan it ran with.
434 DataQuality {
435 results: Box<crate::analysis::data_quality::DataQualityResults>,
436 kept: Option<crate::KeptQualitySample>,
437 plan: Box<crate::analysis::data_quality::DataQualityPlan>,
438 },
439 /// [`Job::SampleRows`]: rows to show as a table.
440 Sample { df: DataFrame, label: String },
441 /// [`Job::SampleDraw`]: the draw ended; its rows are in the job's chunks.
442 SampleDrawn(crate::analysis::table_sample::Drawn),
443 /// [`Job::Pivot`]: the pivot.
444 Pivoted {
445 spec: crate::app::modals::pivot_melt_modal::PivotSpec,
446 pivoted: DataFrame,
447 },
448 /// [`Job::ViewPivot`]: the view's pivot.
449 ViewPivoted(DataFrame),
450 /// [`Job::ReshapePreview`]: the head read, if any, and the preview or the reshape's
451 /// error.
452 ReshapePreviewed {
453 input: Option<crate::app::modals::pivot_melt_modal::PreviewInput>,
454 result: Result<crate::app::modals::pivot_melt_modal::PreviewFrame, String>,
455 },
456 /// [`Job::DrillRow`]: the group row.
457 DrillRow { group_index: usize, row: DataFrame },
458 /// [`Job::Regroup`]: the group's view row and row, or none when no group has its
459 /// keys now.
460 Regrouped(Option<(usize, DataFrame)>),
461 /// [`Job::InspectRow`]: the fields read.
462 FieldsRead(DataFrame),
463 /// [`Job::InspectJson`]: the document.
464 JsonParsed(std::sync::Arc<serde_json::Value>),
465 /// [`Job::InspectPretty`]: the text, indented.
466 Indented(std::sync::Arc<str>),
467 /// [`Job::InspectUnpack`]: the text, as far as it was decompressed.
468 Unpacked(crate::inspector::inspector_bytes::Decoded),
469 /// [`Job::OpenValue`]: the file, written.
470 ValueWritten(crate::inspector::external_open::ExternalOpen),
471 /// [`Job::Export`]: the file, committed.
472 Exported(PathBuf),
473 /// [`Job::Copy`]: the formatted view or field and its flash. The clipboard is
474 /// written on the event thread, which owns its handle.
475 Copied {
476 payload: crate::clipboard::Payload,
477 message: String,
478 },
479 /// [`Job::QualityReport`]: the report, written.
480 QualityReportWritten(PathBuf),
481 /// [`Job::ChartExport`]: the chart, written.
482 ChartExported,
483 /// [`Job::ChartPrepare`]: the chart's data, and the Color column's values if
484 /// counted.
485 ChartPrepared(
486 Box<(
487 crate::chart::chart_plot::PlotData,
488 Option<crate::chart::chart_modal::ColorCounts>,
489 )>,
490 ),
491 /// [`Job::FileFacts`]: what the file is.
492 FileFacts(crate::widgets::info::FileFacts),
493 /// [`Job::Find`]: the cell found, or `None` when nothing in the view matches.
494 Found(Option<crate::find::Found>),
495 /// [`Job::ValueCounts`]: the column's values, counted.
496 ValueCounts(Box<crate::analysis::value_counts::ValueCounts>),
497 /// [`Job::HexOpen`]: the file, mapped.
498 HexOpened(Box<crate::app::hex_view::HexSource>),
499 /// [`Job::HexFind`]: where the pattern is, if anywhere.
500 HexFound(crate::app::hex_view::HexHit),
501 /// [`Job::UnfitCount`]: the columns whose types made values null.
502 UnfitCounted(Vec<crate::formats::column_types::Unfit>),
503 /// A test's answer, which says when it is dropped.
504 #[cfg(test)]
505 Probe(Arc<()>),
506 /// [`Job::FootersJoin`]: what the footers say; `None` when unreadable, so the
507 /// dataset stops waiting.
508 FootersJoined(Option<Box<crate::table::FootersFound>>),
509 /// [`Job::JournalDetail`]: the journal's Info tab.
510 JournalDescribed(Box<crate::formats::text_formats::Detail>),
511 /// [`Job::IndexLines`]: every line is indexed, this many rows of them.
512 LinesIndexed(usize),
513}
514
515impl Answer {
516 /// This answer, then `after` on the worker's thread once sent: for work riding in a
517 /// job, like the row count a page read may answer, which must not delay the page.
518 pub(crate) fn then(self, after: impl FnOnce() + Send + 'static) -> Answered {
519 Answered {
520 answer: self,
521 then: Some(Box::new(after)),
522 }
523 }
524}
525
526/// An answer, and what the worker does once it has been sent.
527pub(crate) struct Answered {
528 answer: Answer,
529 then: Option<Box<dyn FnOnce() + Send>>,
530}
531
532impl From<Answer> for Answered {
533 fn from(answer: Answer) -> Self {
534 Self { answer, then: None }
535 }
536}
537
538/// How a job ended.
539pub(crate) enum Outcome {
540 Answered(Box<Answer>),
541 /// The worker returned an error or panicked; `panicked` means `message` is an
542 /// internal error naming the log.
543 Failed {
544 message: String,
545 panicked: bool,
546 },
547}
548
549impl Outcome {
550 /// A job's answer, for tests that end a job by hand.
551 #[cfg(test)]
552 pub(crate) fn answered(answer: Answer) -> Self {
553 Self::Answered(Box::new(answer))
554 }
555}
556
557/// A report from a job still running.
558#[derive(Debug, Clone)]
559pub enum Progress {
560 /// An export has written `bytes` of its file.
561 ExportWriting { phase: &'static str, bytes: u64 },
562 /// A Data Quality run entered a stage.
563 QualityPhase(crate::analysis::data_quality::QualityPhase),
564 /// A find has read `rows` of the view.
565 Finding { rows: usize },
566 /// A find in the hex view has read `read` of the file's `total` bytes.
567 HexFinding { read: u64, total: u64 },
568 /// A sample's rows are cut to its scope, with these columns: its view can be built.
569 SampleBegun(polars::prelude::SchemaRef),
570 /// A sample kept another chunk.
571 SampleGrew,
572}
573
574/// A job whose outcome has been taken: what it was, whether its answer is still
575/// wanted, and the outcome.
576pub(crate) struct Ended {
577 pub(crate) job: Job,
578 /// Not superseded: the answer is the one the app is waiting for.
579 pub(crate) current: bool,
580 /// The footer's line while the user waited, if they still do.
581 pub(crate) keys: Option<String>,
582 pub(crate) outcome: Outcome,
583}
584
585/// Picks the jobs that panic before their work starts; see [`Jobs::worker_dies`].
586#[cfg(test)]
587pub(crate) type WorkerDies = Box<dyn FnMut(&Job) -> bool + Send>;
588
589/// Picks the jobs whose worker waits before its work starts; see [`Jobs::worker_waits`].
590#[cfg(test)]
591pub(crate) type WorkerWaits = Box<dyn FnMut(&Job) -> Option<std::sync::mpsc::Receiver<()>> + Send>;
592
593type Slot = Arc<Mutex<Option<Outcome>>>;
594
595/// The record of one job: the only marker the app keeps for it.
596struct Record {
597 ticket: Ticket,
598 job: Job,
599 /// Where the worker puts the outcome; `None` for an owed job (no worker yet).
600 slot: Option<Slot>,
601 /// The footer's line while the user waits on it; set, keys wait.
602 keys: Option<String>,
603 /// When it was superseded: its answer is stale, and it holds neither generation
604 /// nor keys.
605 superseded: Option<Instant>,
606 /// Superseded by the user's cancel, not replacement: a read still going that a new
607 /// one should not start beside.
608 cancelled: bool,
609 /// Set when it is superseded, for its worker to see ([`superseded`]).
610 stale: Arc<std::sync::atomic::AtomicBool>,
611}
612
613std::thread_local! {
614 /// The stale flag of the job running on this thread, if a job runs on it.
615 static RUNNING: std::cell::RefCell<Option<Arc<std::sync::atomic::AtomicBool>>> =
616 const { std::cell::RefCell::new(None) };
617}
618
619/// Whether this thread's job was superseded: its answer will be dropped, so a wait
620/// inside may give up. False off a job's thread.
621pub(crate) fn superseded() -> bool {
622 RUNNING.with(|running| {
623 running
624 .borrow()
625 .as_ref()
626 .is_some_and(|stale| stale.load(std::sync::atomic::Ordering::Relaxed))
627 })
628}
629
630impl Record {
631 fn running(&self) -> bool {
632 self.slot.is_some()
633 }
634
635 fn current(&self) -> bool {
636 self.running() && self.superseded.is_none()
637 }
638
639 /// Whether a bump of the generation would strand its answer.
640 fn holds(&self, generation: u64) -> bool {
641 self.current() && self.job.leased() && self.ticket.generation == generation
642 }
643
644 fn supersede(&mut self, now: Instant) {
645 self.superseded = Some(now);
646 self.keys = None;
647 self.stale.store(true, std::sync::atomic::Ordering::Relaxed);
648 }
649}
650
651/// A started job until its outcome is in. [`Started::run`] hands it to a worker;
652/// dropped without ending, it ends as failed, so no record waits forever.
653pub(crate) struct Started {
654 ticket: Ticket,
655 slot: Slot,
656 events: Sender<AppEvent>,
657 ended: bool,
658 stale: Arc<std::sync::atomic::AtomicBool>,
659 #[cfg(test)]
660 dies: bool,
661 #[cfg(test)]
662 waits: Option<std::sync::mpsc::Receiver<()>>,
663}
664
665impl Started {
666 pub(crate) fn ticket(&self) -> Ticket {
667 self.ticket
668 }
669
670 /// Run `work` on a blocking thread of `runtime`; its answer, error or panic is the
671 /// outcome. What `work` holds is dropped before the outcome is sent.
672 pub(crate) fn run<F, R>(self, runtime: &tokio::runtime::Handle, work: F)
673 where
674 F: FnOnce(&Worker) -> Result<R, String> + Send + 'static,
675 R: Into<Answered>,
676 {
677 let mut started = self;
678 runtime.spawn_blocking(move || {
679 let worker = Worker {
680 ticket: started.ticket,
681 events: started.events.clone(),
682 };
683 #[cfg(test)]
684 let dies = started.dies;
685 // A send or a dropped sender both let it go, so a failing test frees it.
686 #[cfg(test)]
687 if let Some(gate) = started.waits.take() {
688 let _ = gate.recv();
689 }
690 RUNNING.with(|running| *running.borrow_mut() = Some(started.stale.clone()));
691 let ran = logging::catch_panic(|| {
692 #[cfg(test)]
693 if dies {
694 panic!("worker died");
695 }
696 work(&worker)
697 });
698 RUNNING.with(|running| *running.borrow_mut() = None);
699 match ran {
700 Ok(Ok(answered)) => {
701 let Answered { answer, then } = answered.into();
702 started.finish(Outcome::Answered(Box::new(answer)));
703 if let Some(then) = then {
704 then();
705 }
706 }
707 Ok(Err(message)) => started.finish(Outcome::Failed {
708 message,
709 panicked: false,
710 }),
711 Err(message) => started.finish(Outcome::Failed {
712 message,
713 panicked: true,
714 }),
715 }
716 });
717 }
718
719 /// End the job with `outcome` here, with no worker.
720 #[cfg(test)]
721 pub(crate) fn end(mut self, outcome: Outcome) {
722 self.finish(outcome);
723 }
724
725 fn finish(&mut self, outcome: Outcome) {
726 if std::mem::replace(&mut self.ended, true) {
727 return;
728 }
729 *self.slot.lock().unwrap_or_else(|e| e.into_inner()) = Some(outcome);
730 // Nobody to tell: the app is going, so drop what the outcome holds (a downloaded
731 // file) now.
732 if self.events.send(AppEvent::JobEnded(self.ticket)).is_err() {
733 drop(self.slot.lock().unwrap_or_else(|e| e.into_inner()).take());
734 }
735 }
736}
737
738impl Drop for Started {
739 fn drop(&mut self) {
740 self.finish(Outcome::Failed {
741 message: "The background task stopped before it answered".to_string(),
742 panicked: true,
743 });
744 }
745}
746
747/// What a running job's worker has: its ticket, and a way to report progress.
748pub(crate) struct Worker {
749 ticket: Ticket,
750 events: Sender<AppEvent>,
751}
752
753impl Worker {
754 /// A progress reporter; the app takes reports only while the job is current.
755 pub(crate) fn reporter(&self) -> impl Fn(Progress) + Send + Sync + 'static {
756 let (ticket, events) = (self.ticket, Mutex::new(self.events.clone()));
757 move |progress| {
758 let events = events.lock().unwrap_or_else(|e| e.into_inner());
759 let _ = events.send(AppEvent::JobProgress { ticket, progress });
760 }
761 }
762
763 /// Send an event not this job's: something read that is worth keeping whatever
764 /// becomes of the job.
765 pub(crate) fn send(&self, event: AppEvent) {
766 let _ = self.events.send(event);
767 }
768}
769
770type HoldCounts = Arc<Mutex<HashMap<u64, usize>>>;
771
772/// A hold on the generation by work that is not a running job: a continuation
773/// [`crate::app::event_pump::EventPump`] has not dispatched, the gap between an errand's
774/// phases, a download waiting on the user. Released on drop.
775#[must_use]
776pub(crate) struct Hold {
777 generation: u64,
778 counts: HoldCounts,
779}
780
781impl Drop for Hold {
782 fn drop(&mut self) {
783 let mut counts = self.counts.lock().unwrap_or_else(|e| e.into_inner());
784 if let Some(n) = counts.get_mut(&self.generation) {
785 *n = n.saturating_sub(1);
786 if *n == 0 {
787 counts.remove(&self.generation);
788 }
789 }
790 }
791}
792
793/// The owner of every general background operation: the generation, one record per
794/// job until handled, and the non-job holds.
795pub(crate) struct Jobs {
796 events: Sender<AppEvent>,
797 /// The generation answers are judged by; advancing it makes following jobs' answers
798 /// stale.
799 generation: u64,
800 next_id: u64,
801 records: Vec<Record>,
802 holds: HoldCounts,
803 /// Which jobs panic before starting, for tests.
804 #[cfg(test)]
805 pub(crate) worker_dies: Option<WorkerDies>,
806 /// Which jobs' workers wait on the returned receiver before starting, so tests can
807 /// catch a job mid-step without a race.
808 #[cfg(test)]
809 pub(crate) worker_waits: Option<WorkerWaits>,
810}
811
812impl Jobs {
813 pub(crate) fn new(events: Sender<AppEvent>) -> Self {
814 Self {
815 events,
816 generation: 0,
817 next_id: 0,
818 records: Vec::new(),
819 holds: Arc::default(),
820 #[cfg(test)]
821 worker_dies: None,
822 #[cfg(test)]
823 worker_waits: None,
824 }
825 }
826
827 pub(crate) fn generation(&self) -> u64 {
828 self.generation
829 }
830
831 fn ticket(&mut self, job: &Job) -> Ticket {
832 self.next_id = self.next_id.wrapping_add(1);
833 Ticket {
834 id: self.next_id,
835 generation: self.generation,
836 kind: job.kind(),
837 }
838 }
839
840 /// Record `job` as started on the current generation. With `keys` the user waits:
841 /// keys are held until it ends or is superseded, the footer showing `keys`.
842 pub(crate) fn start(&mut self, job: Job, keys: Option<&str>) -> Started {
843 let ticket = self.ticket(&job);
844 #[cfg(test)]
845 let dies = self.worker_dies.as_mut().is_some_and(|dies| dies(&job));
846 #[cfg(test)]
847 let waits = self.worker_waits.as_mut().and_then(|waits| waits(&job));
848 let slot = Slot::default();
849 let stale = Arc::<std::sync::atomic::AtomicBool>::default();
850 self.records.push(Record {
851 ticket,
852 job,
853 slot: Some(slot.clone()),
854 keys: keys.map(str::to_string),
855 superseded: None,
856 cancelled: false,
857 stale: stale.clone(),
858 });
859 Started {
860 ticket,
861 slot,
862 events: self.events.clone(),
863 ended: false,
864 stale,
865 #[cfg(test)]
866 dies,
867 #[cfg(test)]
868 waits,
869 }
870 }
871
872 /// Record `job` as owed until the generation is free (waited on with `keys`).
873 /// Nothing runs until the app takes it back with [`Self::take_owed`].
874 pub(crate) fn owe(&mut self, job: Job, keys: Option<&str>) {
875 let ticket = self.ticket(&job);
876 self.records.push(Record {
877 ticket,
878 job,
879 slot: None,
880 keys: keys.map(str::to_string),
881 superseded: None,
882 cancelled: false,
883 stale: Arc::default(),
884 });
885 }
886
887 /// The owed job `which` picks, if there is one.
888 pub(crate) fn owed(&self, which: impl Fn(&Job) -> bool) -> Option<&Job> {
889 self.records
890 .iter()
891 .find(|r| !r.running() && which(&r.job))
892 .map(|r| &r.job)
893 }
894
895 /// Take back the owed job `which` picks, to run or to put down.
896 pub(crate) fn take_owed(&mut self, which: impl Fn(&Job) -> bool) -> Option<Job> {
897 let at = self
898 .records
899 .iter()
900 .position(|r| !r.running() && which(&r.job))?;
901 Some(self.records.remove(at).job)
902 }
903
904 /// Take the outcome of `ticket`'s job and remove its record; `None` if no record or
905 /// no outcome yet. Removing it here means the job holds generation and keys until
906 /// the app has its answer, and no longer.
907 pub(crate) fn end(&mut self, ticket: Ticket) -> Option<Ended> {
908 let at = self.records.iter().position(|r| r.ticket == ticket)?;
909 let outcome = self.records[at]
910 .slot
911 .as_ref()?
912 .lock()
913 .unwrap_or_else(|e| e.into_inner())
914 .take()?;
915 let record = self.records.remove(at);
916 Some(Ended {
917 job: record.job,
918 current: record.superseded.is_none(),
919 keys: record.keys,
920 outcome,
921 })
922 }
923
924 /// Whether the job `ticket` names is running and its answer still wanted.
925 pub(crate) fn is_current(&self, ticket: Ticket) -> bool {
926 self.records
927 .iter()
928 .any(|r| r.ticket == ticket && r.current())
929 }
930
931 /// The newest running job `which` picks whose answer is still wanted.
932 pub(crate) fn current(&self, which: impl Fn(&Job) -> bool) -> Option<(Ticket, &Job)> {
933 self.records
934 .iter()
935 .rev()
936 .find(|r| r.current() && which(&r.job))
937 .map(|r| (r.ticket, &r.job))
938 }
939
940 /// As [`Self::current`], to change what the record says about the job.
941 pub(crate) fn current_mut(&mut self, which: impl Fn(&Job) -> bool) -> Option<&mut Job> {
942 self.records
943 .iter_mut()
944 .rev()
945 .find(|r| r.current() && which(&r.job))
946 .map(|r| &mut r.job)
947 }
948
949 /// What the record of the job `ticket` names says about it, to change it.
950 pub(crate) fn job_mut(&mut self, ticket: Ticket) -> Option<&mut Job> {
951 self.records
952 .iter_mut()
953 .find(|r| r.ticket == ticket)
954 .map(|r| &mut r.job)
955 }
956
957 /// The newest job `which` picks that the user cancelled and whose worker still runs,
958 /// with when it was cancelled.
959 pub(crate) fn cancelled_running(
960 &self,
961 which: impl Fn(&Job) -> bool,
962 ) -> Option<(Instant, &Job)> {
963 self.records
964 .iter()
965 .filter(|r| r.running() && r.cancelled && which(&r.job))
966 .filter_map(|r| r.superseded.map(|since| (since, &r.job)))
967 .max_by_key(|(since, _)| *since)
968 }
969
970 /// Whether a job `which` picks is running, its answer wanted or not.
971 pub(crate) fn running(&self, which: impl Fn(&Job) -> bool) -> bool {
972 self.records.iter().any(|r| r.running() && which(&r.job))
973 }
974
975 /// Whether a job still holding the keys shows `status` on the footer.
976 pub(crate) fn shows(&self, status: &str) -> bool {
977 self.records
978 .iter()
979 .any(|r| r.keys.as_deref() == Some(status))
980 }
981
982 /// Whether advancing the generation now would throw away an answer nothing will ask
983 /// for again: a stranded job or a hold. Asked of the records, not of a list of
984 /// kinds of work.
985 pub(crate) fn would_strand(&self) -> bool {
986 self.records.iter().any(|r| r.holds(self.generation)) || self.held(self.generation)
987 }
988
989 fn held(&self, generation: u64) -> bool {
990 self.holds
991 .lock()
992 .unwrap_or_else(|e| e.into_inner())
993 .get(&generation)
994 .is_some_and(|n| *n > 0)
995 }
996
997 /// Whether the user is waiting on a job: one running or owed that holds the keys.
998 pub(crate) fn holds_keys(&self) -> bool {
999 self.records.iter().any(|r| r.keys.is_some())
1000 }
1001
1002 /// Whether the user waits on some job, and on none but those `which` picks.
1003 pub(crate) fn keys_held_only_by(&self, which: impl Fn(&Job) -> bool) -> bool {
1004 self.holds_keys()
1005 && self
1006 .records
1007 .iter()
1008 .filter(|r| r.keys.is_some())
1009 .all(|r| which(&r.job))
1010 }
1011
1012 /// Whether the user waits on the newest current job `which` picks.
1013 pub(crate) fn waited_on(&self, which: impl Fn(&Job) -> bool) -> bool {
1014 self.records
1015 .iter()
1016 .rev()
1017 .find(|r| r.current() && which(&r.job))
1018 .is_some_and(|r| r.keys.is_some())
1019 }
1020
1021 /// The footer line of the newest waited-on job `which` picks, running or owed.
1022 pub(crate) fn waiting_status(&self, which: impl Fn(&Job) -> bool) -> Option<&str> {
1023 self.records
1024 .iter()
1025 .rev()
1026 .filter(|r| r.superseded.is_none() && which(&r.job))
1027 .find_map(|r| r.keys.as_deref())
1028 }
1029
1030 /// The user waits on the running job `which` picks from now on, with `status` on
1031 /// the footer (a load-ahead a scroll caught up with). Whether there was one.
1032 pub(crate) fn wait_on(&mut self, which: impl Fn(&Job) -> bool, status: &str) -> bool {
1033 let Some(record) = self
1034 .records
1035 .iter_mut()
1036 .rev()
1037 .find(|r| r.current() && which(&r.job))
1038 else {
1039 return false;
1040 };
1041 record.keys = Some(status.to_string());
1042 true
1043 }
1044
1045 /// Nobody waits on the jobs `which` picks any more, though their answers are still
1046 /// wanted: keys return to the user. Returns their footer lines.
1047 pub(crate) fn quiet(&mut self, which: impl Fn(&Job) -> bool) -> Vec<String> {
1048 self.records
1049 .iter_mut()
1050 .filter(|r| which(&r.job))
1051 .filter_map(|r| r.keys.take())
1052 .collect()
1053 }
1054
1055 /// Advance the generation if that strands nothing. Returns whether it did.
1056 pub(crate) fn try_advance(&mut self) -> bool {
1057 if self.would_strand() {
1058 return false;
1059 }
1060 self.advance();
1061 true
1062 }
1063
1064 /// Advance the generation whatever runs (a cancel, a new dataset): jobs that follow
1065 /// it are superseded.
1066 pub(crate) fn advance(&mut self) {
1067 self.generation = self.generation.wrapping_add(1);
1068 let now = Instant::now();
1069 for record in &mut self.records {
1070 if record.current() && record.job.follows_the_generation() {
1071 record.supersede(now);
1072 }
1073 }
1074 }
1075
1076 /// Supersede the running jobs `which` picks as the user's cancel: as
1077 /// [`Self::supersede`], and [`Self::cancelled_running`] says so while they run on.
1078 pub(crate) fn cancel(&mut self, which: impl Fn(&Job) -> bool) -> bool {
1079 for record in &mut self.records {
1080 if record.current() && which(&record.job) {
1081 record.cancelled = true;
1082 }
1083 }
1084 self.supersede(which)
1085 }
1086
1087 /// Supersede the running jobs `which` picks and drop the owed ones: stale, holding
1088 /// neither generation nor keys. Whether there were any.
1089 pub(crate) fn supersede(&mut self, which: impl Fn(&Job) -> bool) -> bool {
1090 let now = Instant::now();
1091 let before = self.records.len();
1092 self.records.retain(|r| r.running() || !which(&r.job));
1093 let mut any = self.records.len() != before;
1094 for record in &mut self.records {
1095 if record.current() && which(&record.job) {
1096 record.supersede(now);
1097 any = true;
1098 }
1099 }
1100 any
1101 }
1102
1103 /// Hold the current generation until the hold is dropped.
1104 pub(crate) fn hold(&self) -> Hold {
1105 *self
1106 .holds
1107 .lock()
1108 .unwrap_or_else(|e| e.into_inner())
1109 .entry(self.generation)
1110 .or_default() += 1;
1111 Hold {
1112 generation: self.generation,
1113 counts: self.holds.clone(),
1114 }
1115 }
1116
1117 /// Whether any leased job is still running, superseded or not, or anything holds a
1118 /// generation.
1119 pub(crate) fn in_flight(&self) -> bool {
1120 self.records.iter().any(|r| r.running() && r.job.leased())
1121 || self
1122 .holds
1123 .lock()
1124 .unwrap_or_else(|e| e.into_inner())
1125 .values()
1126 .any(|n| *n > 0)
1127 }
1128
1129 /// Whether leased work a cancel passed still runs (started on a generation since
1130 /// left).
1131 pub(crate) fn running_behind(&self) -> bool {
1132 self.records
1133 .iter()
1134 .any(|r| r.running() && r.job.leased() && r.ticket.generation != self.generation)
1135 || self
1136 .holds
1137 .lock()
1138 .unwrap_or_else(|e| e.into_inner())
1139 .iter()
1140 .any(|(generation, n)| *generation != self.generation && *n > 0)
1141 }
1142
1143 /// Move every supersession back by `by`, for tests of cancel timing.
1144 #[cfg(test)]
1145 pub(crate) fn backdate_supersessions(&mut self, by: std::time::Duration) {
1146 for record in &mut self.records {
1147 if let Some(since) = record.superseded.as_mut() {
1148 *since -= by;
1149 }
1150 }
1151 }
1152}
1153
1154#[cfg(test)]
1155mod tests;