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;