Skip to main content

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;