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