1use std::collections::HashMap;
48use std::path::PathBuf;
49use std::sync::mpsc::Sender;
50use std::sync::{Arc, Mutex};
51use std::time::Instant;
52
53use polars::prelude::DataFrame;
54
55use crate::loading::{LoadAnswer, LoadId};
56use crate::{AppEvent, OpenOptions, logging};
57
58#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
60pub enum JobKind {
61 Load,
62 OpenNamed,
63 LookAtDirectory,
64 Classify,
65 Rows,
66 Analysis,
67 SampleRows,
68 SampleDraw,
69 Pivot,
70 ViewPivot,
71 ReshapePreview,
72 DrillRow,
73 InspectRow,
74 InspectJson,
75 InspectPretty,
76 InspectUnpack,
77 OpenValue,
78 Export,
79 Copy,
80 QualityReport,
81 FileFacts,
82 ChartExport,
83 Find,
84 ValueCounts,
85 HexOpen,
86 HexFind,
87 UnfitCount,
88}
89
90#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
93pub struct Ticket {
94 id: u64,
95 generation: u64,
96 kind: JobKind,
97}
98
99impl Ticket {
100 pub fn kind(self) -> JobKind {
101 self.kind
102 }
103
104 pub fn generation(self) -> u64 {
106 self.generation
107 }
108}
109
110#[derive(Debug, Clone)]
115pub(crate) enum Job {
116 Load(LoadId),
120 OpenNamed(LoadId),
123 LookAtDirectory { load: LoadId, path: PathBuf },
125 Classify(Classify),
127 Rows(crate::InflightCollect),
129 OwedRows { dataset: u64, status: String },
132 Analysis(AnalysisRun),
134 SampleRows,
136 SampleDraw(Box<SampleDraw>),
138 Pivot,
140 ViewPivot(Box<(crate::view::SavedView, Option<crate::view::MatchReason>)>),
143 ReshapePreview { epoch: u64, token: u64 },
146 DrillRow,
148 InspectRow { frame: u64, row: usize },
151 InspectJson { token: u64 },
154 InspectPretty { token: u64 },
157 InspectUnpack { token: u64 },
160 OpenValue,
162 Export,
164 Copy,
166 QualityReport,
168 FileFacts { dataset: u64 },
172 ChartExport {
174 path: PathBuf,
175 format: crate::chart_export::ChartExportFormat,
176 },
177 Find(crate::find::FindRun),
179 ValueCounts,
181 HexOpen {
183 origin: crate::hex_view::Origin,
184 fallback: bool,
185 record_size: Option<usize>,
186 },
187 HexFind(crate::hex_view::HexFindRun),
189 UnfitCount { dataset: u64, version: Option<u64> },
193}
194
195#[derive(Debug, Clone)]
199pub(crate) struct Classify {
200 pub(crate) path: PathBuf,
201 pub(crate) browsing: Option<PathBuf>,
209 pub(crate) jump: bool,
211}
212
213#[derive(Clone)]
217pub(crate) struct SampleDraw {
218 pub(crate) sample: crate::sampling::Sample,
219 pub(crate) rows: Arc<crate::table_sample::SampleRows>,
220 pub(crate) watch: crate::sampling::ReadWatch,
222 pub(crate) through: bool,
224 pub(crate) replay: Option<crate::view::ViewSettings>,
227 pub(crate) then_analyze: bool,
229 pub(crate) path: Option<crate::table_sample::DrawPath>,
231 pub(crate) path_key: String,
233 pub(crate) schema: Option<polars::prelude::SchemaRef>,
237}
238
239impl std::fmt::Debug for SampleDraw {
240 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
241 f.debug_struct("SampleDraw")
242 .field("sample", &self.sample)
243 .field("through", &self.through)
244 .field("then_analyze", &self.then_analyze)
245 .finish_non_exhaustive()
246 }
247}
248
249#[derive(Debug, Clone, Default)]
251pub(crate) struct AnalysisRun {
252 pub(crate) watch: Option<crate::data_quality::QualityWatch>,
254 pub(crate) runs_out: bool,
256}
257
258impl Job {
259 pub(crate) fn kind(&self) -> JobKind {
260 match self {
261 Job::Load(_) => JobKind::Load,
262 Job::OpenNamed(_) => JobKind::OpenNamed,
263 Job::LookAtDirectory { .. } => JobKind::LookAtDirectory,
264 Job::Classify(_) => JobKind::Classify,
265 Job::Rows(_) | Job::OwedRows { .. } => JobKind::Rows,
266 Job::Analysis(_) => JobKind::Analysis,
267 Job::SampleRows => JobKind::SampleRows,
268 Job::SampleDraw(_) => JobKind::SampleDraw,
269 Job::Pivot => JobKind::Pivot,
270 Job::ViewPivot(_) => JobKind::ViewPivot,
271 Job::ReshapePreview { .. } => JobKind::ReshapePreview,
272 Job::DrillRow => JobKind::DrillRow,
273 Job::InspectRow { .. } => JobKind::InspectRow,
274 Job::InspectJson { .. } => JobKind::InspectJson,
275 Job::InspectPretty { .. } => JobKind::InspectPretty,
276 Job::InspectUnpack { .. } => JobKind::InspectUnpack,
277 Job::OpenValue => JobKind::OpenValue,
278 Job::Export => JobKind::Export,
279 Job::Copy => JobKind::Copy,
280 Job::QualityReport => JobKind::QualityReport,
281 Job::FileFacts { .. } => JobKind::FileFacts,
282 Job::ChartExport { .. } => JobKind::ChartExport,
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 }
289 }
290
291 pub(crate) fn load(&self) -> Option<LoadId> {
293 match self {
294 Job::Load(load) | Job::OpenNamed(load) | Job::LookAtDirectory { load, .. } => {
295 Some(*load)
296 }
297 _ => None,
298 }
299 }
300
301 fn leased(&self) -> bool {
317 !matches!(
318 self,
319 Job::Rows(_)
320 | Job::SampleDraw(_)
321 | Job::OwedRows { .. }
322 | Job::OpenNamed(_)
323 | Job::LookAtDirectory { .. }
324 | Job::FileFacts { .. }
325 | Job::UnfitCount { .. }
326 | Job::ReshapePreview { .. }
327 )
328 }
329
330 fn follows_the_generation(&self) -> bool {
335 !matches!(
336 self,
337 Job::FileFacts { .. }
338 | Job::SampleDraw(_)
339 | Job::UnfitCount { .. }
340 | Job::ChartExport { .. }
341 | Job::OwedRows { .. }
342 | Job::ReshapePreview { .. }
343 )
344 }
345}
346
347pub(crate) enum Answer {
349 Load(Box<LoadAnswer>),
352 NamedPaths {
354 paths: Vec<PathBuf>,
355 options: Box<OpenOptions>,
356 directory: Option<PathBuf>,
357 },
358 NamedPathMissing(PathBuf),
360 LookedAt {
363 kind: crate::discover::EntryKind,
364 holds: Option<Box<crate::discover::Holds>>,
365 options: Box<OpenOptions>,
366 },
367 Kind(Option<crate::discover::EntryKind>),
370 Rows(crate::widgets::datatable::CollectResult),
372 RowsFailed {
375 message: String,
376 conversion: Option<Box<crate::error_display::ConversionFailure>>,
377 },
378 Described(crate::statistics::AnalysisResults),
380 Distributions(crate::statistics::AnalysisResults),
382 Correlations(crate::statistics::AnalysisResults),
384 DataQuality {
387 results: Box<crate::data_quality::DataQualityResults>,
388 kept: Option<crate::KeptQualitySample>,
389 plan: Box<crate::data_quality::DataQualityPlan>,
390 },
391 Sample { df: DataFrame, label: String },
393 SampleDrawn(crate::table_sample::Drawn),
395 Pivoted {
397 spec: crate::pivot_melt_modal::PivotSpec,
398 pivoted: DataFrame,
399 },
400 ViewPivoted(DataFrame),
402 ReshapePreviewed {
405 input: Option<crate::pivot_melt_modal::PreviewInput>,
406 result: Result<crate::pivot_melt_modal::PreviewFrame, String>,
407 },
408 DrillRow { group_index: usize, row: DataFrame },
410 FieldsRead(DataFrame),
412 JsonParsed(std::sync::Arc<serde_json::Value>),
414 Indented(std::sync::Arc<str>),
416 Unpacked(crate::inspector_bytes::Decoded),
418 ValueWritten(crate::external_open::ExternalOpen),
420 Exported(PathBuf),
422 Copied {
425 payload: crate::clipboard::Payload,
426 message: String,
427 },
428 QualityReportWritten(PathBuf),
430 ChartExported,
432 FileFacts(crate::widgets::info::FileFacts),
434 Found(Option<crate::find::Found>),
436 ValueCounts(Box<crate::value_counts::ValueCounts>),
438 HexOpened(Box<crate::hex_view::HexSource>),
440 HexFound(crate::hex_view::HexHit),
442 UnfitCounted(Vec<crate::column_types::Unfit>),
444 #[cfg(test)]
446 Probe(Arc<()>),
447}
448
449impl Answer {
450 pub(crate) fn then(self, after: impl FnOnce() + Send + 'static) -> Answered {
454 Answered {
455 answer: self,
456 then: Some(Box::new(after)),
457 }
458 }
459}
460
461pub(crate) struct Answered {
463 answer: Answer,
464 then: Option<Box<dyn FnOnce() + Send>>,
465}
466
467impl From<Answer> for Answered {
468 fn from(answer: Answer) -> Self {
469 Self { answer, then: None }
470 }
471}
472
473pub(crate) enum Outcome {
475 Answered(Box<Answer>),
476 Failed {
479 message: String,
480 panicked: bool,
481 },
482}
483
484impl Outcome {
485 #[cfg(test)]
487 pub(crate) fn answered(answer: Answer) -> Self {
488 Self::Answered(Box::new(answer))
489 }
490}
491
492#[derive(Debug, Clone)]
494pub enum Progress {
495 ExportWriting { phase: &'static str, bytes: u64 },
497 QualityPhase(crate::data_quality::QualityPhase),
499 Finding { rows: usize },
501 HexFinding { read: u64, total: u64 },
503 SampleBegun(polars::prelude::SchemaRef),
505 SampleGrew,
507}
508
509pub(crate) struct Ended {
512 pub(crate) job: Job,
513 pub(crate) current: bool,
515 pub(crate) keys: Option<String>,
518 pub(crate) outcome: Outcome,
519}
520
521#[cfg(test)]
523pub(crate) type WorkerDies = Box<dyn FnMut(&Job) -> bool + Send>;
524
525#[cfg(test)]
527pub(crate) type WorkerWaits = Box<dyn FnMut(&Job) -> Option<std::sync::mpsc::Receiver<()>> + Send>;
528
529type Slot = Arc<Mutex<Option<Outcome>>>;
530
531struct Record {
533 ticket: Ticket,
534 job: Job,
535 slot: Option<Slot>,
538 keys: Option<String>,
540 superseded: Option<Instant>,
543 cancelled: bool,
546 stale: Arc<std::sync::atomic::AtomicBool>,
548}
549
550std::thread_local! {
551 static RUNNING: std::cell::RefCell<Option<Arc<std::sync::atomic::AtomicBool>>> =
553 const { std::cell::RefCell::new(None) };
554}
555
556pub(crate) fn superseded() -> bool {
559 RUNNING.with(|running| {
560 running
561 .borrow()
562 .as_ref()
563 .is_some_and(|stale| stale.load(std::sync::atomic::Ordering::Relaxed))
564 })
565}
566
567impl Record {
568 fn running(&self) -> bool {
569 self.slot.is_some()
570 }
571
572 fn current(&self) -> bool {
573 self.running() && self.superseded.is_none()
574 }
575
576 fn holds(&self, generation: u64) -> bool {
578 self.current() && self.job.leased() && self.ticket.generation == generation
579 }
580
581 fn supersede(&mut self, now: Instant) {
582 self.superseded = Some(now);
583 self.keys = None;
584 self.stale.store(true, std::sync::atomic::Ordering::Relaxed);
585 }
586}
587
588pub(crate) struct Started {
593 ticket: Ticket,
594 slot: Slot,
595 events: Sender<AppEvent>,
596 ended: bool,
597 stale: Arc<std::sync::atomic::AtomicBool>,
598 #[cfg(test)]
599 dies: bool,
600 #[cfg(test)]
601 waits: Option<std::sync::mpsc::Receiver<()>>,
602}
603
604impl Started {
605 pub(crate) fn ticket(&self) -> Ticket {
606 self.ticket
607 }
608
609 pub(crate) fn run<F, R>(self, runtime: &tokio::runtime::Handle, work: F)
614 where
615 F: FnOnce(&Worker) -> Result<R, String> + Send + 'static,
616 R: Into<Answered>,
617 {
618 let mut started = self;
619 runtime.spawn_blocking(move || {
620 let worker = Worker {
621 ticket: started.ticket,
622 events: started.events.clone(),
623 };
624 #[cfg(test)]
625 let dies = started.dies;
626 #[cfg(test)]
628 if let Some(gate) = started.waits.take() {
629 let _ = gate.recv();
630 }
631 RUNNING.with(|running| *running.borrow_mut() = Some(started.stale.clone()));
632 let ran = logging::catch_panic(|| {
633 #[cfg(test)]
634 if dies {
635 panic!("worker died");
636 }
637 work(&worker)
638 });
639 RUNNING.with(|running| *running.borrow_mut() = None);
640 match ran {
641 Ok(Ok(answered)) => {
642 let Answered { answer, then } = answered.into();
643 started.finish(Outcome::Answered(Box::new(answer)));
644 if let Some(then) = then {
645 then();
646 }
647 }
648 Ok(Err(message)) => started.finish(Outcome::Failed {
649 message,
650 panicked: false,
651 }),
652 Err(message) => started.finish(Outcome::Failed {
653 message,
654 panicked: true,
655 }),
656 }
657 });
658 }
659
660 #[cfg(test)]
663 pub(crate) fn end(mut self, outcome: Outcome) {
664 self.finish(outcome);
665 }
666
667 fn finish(&mut self, outcome: Outcome) {
668 if std::mem::replace(&mut self.ended, true) {
669 return;
670 }
671 *self.slot.lock().unwrap_or_else(|e| e.into_inner()) = Some(outcome);
672 if self.events.send(AppEvent::JobEnded(self.ticket)).is_err() {
675 drop(self.slot.lock().unwrap_or_else(|e| e.into_inner()).take());
676 }
677 }
678}
679
680impl Drop for Started {
681 fn drop(&mut self) {
682 self.finish(Outcome::Failed {
683 message: "The background task stopped before it answered".to_string(),
684 panicked: true,
685 });
686 }
687}
688
689pub(crate) struct Worker {
691 ticket: Ticket,
692 events: Sender<AppEvent>,
693}
694
695impl Worker {
696 pub(crate) fn reporter(&self) -> impl Fn(Progress) + Send + Sync + 'static {
699 let (ticket, events) = (self.ticket, Mutex::new(self.events.clone()));
700 move |progress| {
701 let events = events.lock().unwrap_or_else(|e| e.into_inner());
702 let _ = events.send(AppEvent::JobProgress { ticket, progress });
703 }
704 }
705
706 pub(crate) fn send(&self, event: AppEvent) {
709 let _ = self.events.send(event);
710 }
711}
712
713type HoldCounts = Arc<Mutex<HashMap<u64, usize>>>;
714
715#[must_use]
719pub(crate) struct Hold {
720 generation: u64,
721 counts: HoldCounts,
722}
723
724impl Drop for Hold {
725 fn drop(&mut self) {
726 let mut counts = self.counts.lock().unwrap_or_else(|e| e.into_inner());
727 if let Some(n) = counts.get_mut(&self.generation) {
728 *n = n.saturating_sub(1);
729 if *n == 0 {
730 counts.remove(&self.generation);
731 }
732 }
733 }
734}
735
736pub(crate) struct Jobs {
740 events: Sender<AppEvent>,
741 generation: u64,
744 next_id: u64,
745 records: Vec<Record>,
746 holds: HoldCounts,
747 #[cfg(test)]
750 pub(crate) worker_dies: Option<WorkerDies>,
751 #[cfg(test)]
754 pub(crate) worker_waits: Option<WorkerWaits>,
755}
756
757impl Jobs {
758 pub(crate) fn new(events: Sender<AppEvent>) -> Self {
759 Self {
760 events,
761 generation: 0,
762 next_id: 0,
763 records: Vec::new(),
764 holds: Arc::default(),
765 #[cfg(test)]
766 worker_dies: None,
767 #[cfg(test)]
768 worker_waits: None,
769 }
770 }
771
772 pub(crate) fn generation(&self) -> u64 {
773 self.generation
774 }
775
776 fn ticket(&mut self, job: &Job) -> Ticket {
777 self.next_id = self.next_id.wrapping_add(1);
778 Ticket {
779 id: self.next_id,
780 generation: self.generation,
781 kind: job.kind(),
782 }
783 }
784
785 pub(crate) fn start(&mut self, job: Job, keys: Option<&str>) -> Started {
789 let ticket = self.ticket(&job);
790 #[cfg(test)]
791 let dies = self.worker_dies.as_mut().is_some_and(|dies| dies(&job));
792 #[cfg(test)]
793 let waits = self.worker_waits.as_mut().and_then(|waits| waits(&job));
794 let slot = Slot::default();
795 let stale = Arc::<std::sync::atomic::AtomicBool>::default();
796 self.records.push(Record {
797 ticket,
798 job,
799 slot: Some(slot.clone()),
800 keys: keys.map(str::to_string),
801 superseded: None,
802 cancelled: false,
803 stale: stale.clone(),
804 });
805 Started {
806 ticket,
807 slot,
808 events: self.events.clone(),
809 ended: false,
810 stale,
811 #[cfg(test)]
812 dies,
813 #[cfg(test)]
814 waits,
815 }
816 }
817
818 pub(crate) fn owe(&mut self, job: Job, keys: Option<&str>) {
822 let ticket = self.ticket(&job);
823 self.records.push(Record {
824 ticket,
825 job,
826 slot: None,
827 keys: keys.map(str::to_string),
828 superseded: None,
829 cancelled: false,
830 stale: Arc::default(),
831 });
832 }
833
834 pub(crate) fn owed(&self, which: impl Fn(&Job) -> bool) -> Option<&Job> {
836 self.records
837 .iter()
838 .find(|r| !r.running() && which(&r.job))
839 .map(|r| &r.job)
840 }
841
842 pub(crate) fn take_owed(&mut self, which: impl Fn(&Job) -> bool) -> Option<Job> {
844 let at = self
845 .records
846 .iter()
847 .position(|r| !r.running() && which(&r.job))?;
848 Some(self.records.remove(at).job)
849 }
850
851 pub(crate) fn end(&mut self, ticket: Ticket) -> Option<Ended> {
858 let at = self.records.iter().position(|r| r.ticket == ticket)?;
859 let outcome = self.records[at]
860 .slot
861 .as_ref()?
862 .lock()
863 .unwrap_or_else(|e| e.into_inner())
864 .take()?;
865 let record = self.records.remove(at);
866 Some(Ended {
867 job: record.job,
868 current: record.superseded.is_none(),
869 keys: record.keys,
870 outcome,
871 })
872 }
873
874 pub(crate) fn is_current(&self, ticket: Ticket) -> bool {
876 self.records
877 .iter()
878 .any(|r| r.ticket == ticket && r.current())
879 }
880
881 pub(crate) fn current(&self, which: impl Fn(&Job) -> bool) -> Option<(Ticket, &Job)> {
883 self.records
884 .iter()
885 .rev()
886 .find(|r| r.current() && which(&r.job))
887 .map(|r| (r.ticket, &r.job))
888 }
889
890 pub(crate) fn current_mut(&mut self, which: impl Fn(&Job) -> bool) -> Option<&mut Job> {
892 self.records
893 .iter_mut()
894 .rev()
895 .find(|r| r.current() && which(&r.job))
896 .map(|r| &mut r.job)
897 }
898
899 pub(crate) fn job_mut(&mut self, ticket: Ticket) -> Option<&mut Job> {
901 self.records
902 .iter_mut()
903 .find(|r| r.ticket == ticket)
904 .map(|r| &mut r.job)
905 }
906
907 pub(crate) fn cancelled_running(
910 &self,
911 which: impl Fn(&Job) -> bool,
912 ) -> Option<(Instant, &Job)> {
913 self.records
914 .iter()
915 .filter(|r| r.running() && r.cancelled && which(&r.job))
916 .filter_map(|r| r.superseded.map(|since| (since, &r.job)))
917 .max_by_key(|(since, _)| *since)
918 }
919
920 pub(crate) fn shows(&self, status: &str) -> bool {
922 self.records
923 .iter()
924 .any(|r| r.keys.as_deref() == Some(status))
925 }
926
927 pub(crate) fn would_strand(&self) -> bool {
933 self.records.iter().any(|r| r.holds(self.generation)) || self.held(self.generation)
934 }
935
936 fn held(&self, generation: u64) -> bool {
937 self.holds
938 .lock()
939 .unwrap_or_else(|e| e.into_inner())
940 .get(&generation)
941 .is_some_and(|n| *n > 0)
942 }
943
944 pub(crate) fn holds_keys(&self) -> bool {
946 self.records.iter().any(|r| r.keys.is_some())
947 }
948
949 pub(crate) fn keys_held_only_by(&self, which: impl Fn(&Job) -> bool) -> bool {
951 self.holds_keys()
952 && self
953 .records
954 .iter()
955 .filter(|r| r.keys.is_some())
956 .all(|r| which(&r.job))
957 }
958
959 pub(crate) fn waited_on(&self, which: impl Fn(&Job) -> bool) -> bool {
961 self.records
962 .iter()
963 .rev()
964 .find(|r| r.current() && which(&r.job))
965 .is_some_and(|r| r.keys.is_some())
966 }
967
968 pub(crate) fn waiting_status(&self, which: impl Fn(&Job) -> bool) -> Option<&str> {
971 self.records
972 .iter()
973 .rev()
974 .filter(|r| r.superseded.is_none() && which(&r.job))
975 .find_map(|r| r.keys.as_deref())
976 }
977
978 pub(crate) fn wait_on(&mut self, which: impl Fn(&Job) -> bool, status: &str) -> bool {
982 let Some(record) = self
983 .records
984 .iter_mut()
985 .rev()
986 .find(|r| r.current() && which(&r.job))
987 else {
988 return false;
989 };
990 record.keys = Some(status.to_string());
991 true
992 }
993
994 pub(crate) fn quiet(&mut self, which: impl Fn(&Job) -> bool) -> Vec<String> {
998 self.records
999 .iter_mut()
1000 .filter(|r| which(&r.job))
1001 .filter_map(|r| r.keys.take())
1002 .collect()
1003 }
1004
1005 pub(crate) fn try_advance(&mut self) -> bool {
1007 if self.would_strand() {
1008 return false;
1009 }
1010 self.advance();
1011 true
1012 }
1013
1014 pub(crate) fn advance(&mut self) {
1018 self.generation = self.generation.wrapping_add(1);
1019 let now = Instant::now();
1020 for record in &mut self.records {
1021 if record.current() && record.job.follows_the_generation() {
1022 record.supersede(now);
1023 }
1024 }
1025 }
1026
1027 pub(crate) fn cancel(&mut self, which: impl Fn(&Job) -> bool) -> bool {
1030 for record in &mut self.records {
1031 if record.current() && which(&record.job) {
1032 record.cancelled = true;
1033 }
1034 }
1035 self.supersede(which)
1036 }
1037
1038 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 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 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 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 #[cfg(test)]
1098 pub(crate) fn backdate_supersessions(&mut self, by: std::time::Duration) {
1099 for record in &mut self.records {
1100 if let Some(since) = record.superseded.as_mut() {
1101 *since -= by;
1102 }
1103 }
1104 }
1105}
1106
1107#[cfg(test)]
1108mod tests {
1109 use super::*;
1110 use std::sync::mpsc::{Receiver, channel};
1111 use std::time::Duration;
1112
1113 fn jobs() -> (Jobs, Receiver<AppEvent>) {
1114 let (tx, rx) = channel();
1115 (Jobs::new(tx), rx)
1116 }
1117
1118 fn spawn<F, R>(jobs: &mut Jobs, job: Job, work: F) -> Ticket
1119 where
1120 F: FnOnce(&Worker) -> Result<R, String> + Send + 'static,
1121 R: Into<Answered>,
1122 {
1123 let started = jobs.start(job, None);
1124 let ticket = started.ticket();
1125 started.run(&crate::tests::test_runtime(), work);
1126 ticket
1127 }
1128
1129 fn ended(rx: &Receiver<AppEvent>) -> Ticket {
1131 match rx.recv_timeout(Duration::from_secs(30)) {
1132 Ok(AppEvent::JobEnded(ticket)) => ticket,
1133 Ok(_) => panic!("another event"),
1134 Err(e) => panic!("the job never ended: {e}"),
1135 }
1136 }
1137
1138 fn quiet(rx: &Receiver<AppEvent>) {
1139 assert!(
1140 rx.recv_timeout(Duration::from_millis(50)).is_err(),
1141 "nothing more"
1142 );
1143 }
1144
1145 #[test]
1148 fn an_answer_ends_its_job_once() {
1149 let (mut jobs, rx) = jobs();
1150 let (go, wait) = channel::<()>();
1151 let ticket = spawn(&mut jobs, Job::Export, move |_| {
1152 wait.recv().ok();
1153 Ok(Answer::Exported(PathBuf::from("out.csv")))
1154 });
1155 assert_eq!(ticket.kind(), JobKind::Export);
1156 assert!(jobs.would_strand(), "a bump would strand it");
1157 assert!(!jobs.try_advance(), "so the generation stays");
1158 assert!(jobs.is_current(ticket));
1159
1160 go.send(()).unwrap();
1161 assert_eq!(ended(&rx), ticket);
1162 assert!(
1163 jobs.would_strand(),
1164 "its record stands until the app has the answer"
1165 );
1166 let ended = jobs.end(ticket).expect("the outcome is in");
1167 assert!(ended.current);
1168 assert!(matches!(
1169 ended.outcome,
1170 Outcome::Answered(ref answer) if matches!(**answer, Answer::Exported(_))
1171 ));
1172 assert!(!jobs.would_strand(), "and goes with it");
1173 assert!(jobs.end(ticket).is_none(), "once");
1174 quiet(&rx);
1175 assert!(jobs.try_advance());
1176 }
1177
1178 #[test]
1180 fn an_error_and_a_panic_each_end_a_job_once() {
1181 let (mut jobs, rx) = jobs();
1182 for panics in [false, true] {
1183 let ticket = spawn(&mut jobs, Job::QualityReport, move |_| {
1184 assert!(!panics, "worker died");
1185 Err::<Answer, _>("disk full".to_string())
1186 });
1187 assert_eq!(ended(&rx), ticket);
1188 let ended = jobs.end(ticket).expect("the outcome is in");
1189 match ended.outcome {
1190 Outcome::Failed { message, panicked } => {
1191 assert_eq!(panicked, panics);
1192 assert!(message.contains(if panics { "worker died" } else { "disk full" }));
1193 }
1194 Outcome::Answered(_) => panic!("it failed"),
1195 }
1196 quiet(&rx);
1197 assert!(!jobs.would_strand());
1198 }
1199 }
1200
1201 #[test]
1203 fn a_job_dropped_unrun_ends_failed() {
1204 let (mut jobs, rx) = jobs();
1205 let started = jobs.start(Job::Pivot, None);
1206 let ticket = started.ticket();
1207 assert!(jobs.would_strand());
1208 drop(started);
1209 assert_eq!(ended(&rx), ticket);
1210 assert!(matches!(
1211 jobs.end(ticket).map(|e| e.outcome),
1212 Some(Outcome::Failed { panicked: true, .. })
1213 ));
1214 assert!(!jobs.would_strand());
1215 }
1216
1217 #[test]
1221 fn a_superseded_job_ends_stale_and_lets_go() {
1222 let (mut jobs, rx) = jobs();
1223 let held = Arc::new(());
1224 let analysis = || Job::Analysis(AnalysisRun::default());
1225 let is_analysis = |job: &Job| matches!(job, Job::Analysis(_));
1226 let stale = jobs.start(analysis(), Some("Running analysis..."));
1227 let own = jobs.start(
1228 Job::ChartExport {
1229 path: PathBuf::from("chart.png"),
1230 format: crate::chart_export::ChartExportFormat::Png,
1231 },
1232 None,
1233 );
1234 let (stale_ticket, own_ticket) = (stale.ticket(), own.ticket());
1235 assert!(jobs.would_strand());
1236 assert!(jobs.holds_keys());
1237
1238 jobs.advance();
1239 assert!(!jobs.is_current(stale_ticket), "the analysis is superseded");
1240 assert!(jobs.is_current(own_ticket), "the chart export is not");
1241 assert!(!jobs.would_strand(), "and neither holds the new generation");
1242 assert!(!jobs.holds_keys(), "nor the keys");
1243 assert!(jobs.running_behind(), "though it runs on");
1244 assert!(
1245 jobs.cancelled_running(is_analysis).is_none(),
1246 "superseded, not cancelled by the user"
1247 );
1248
1249 stale.end(Outcome::answered(Answer::Probe(held.clone())));
1250 assert_eq!(ended(&rx), stale_ticket);
1251 let stale_end = jobs.end(stale_ticket).expect("its outcome still arrives");
1252 assert!(!stale_end.current, "as stale");
1253 drop(stale_end);
1254 assert_eq!(Arc::strong_count(&held), 1, "and what it carried is let go");
1255 drop(own);
1256 assert_eq!(ended(&rx), own_ticket);
1257 assert!(jobs.end(own_ticket).is_some_and(|e| e.current));
1258 assert!(!jobs.running_behind());
1259 }
1260
1261 #[test]
1264 fn a_cancel_supersedes_only_what_it_picks() {
1265 let (mut jobs, rx) = jobs();
1266 let look = |path: &str| {
1267 Job::Classify(Classify {
1268 path: PathBuf::from(path),
1269 browsing: None,
1270 jump: false,
1271 })
1272 };
1273 let older = jobs.start(look("/a"), Some("Looking..."));
1274 let pivot = jobs.start(Job::Pivot, Some("Computing pivot..."));
1275 assert!(jobs.supersede(|job| matches!(job, Job::Classify(_))));
1276 let newer = jobs.start(look("/b"), Some("Looking..."));
1277 assert!(
1278 jobs.is_current(pivot.ticket()),
1279 "the cancel picked looks only"
1280 );
1281 assert!(
1282 jobs.current(|job| matches!(job, Job::Classify(c) if c.path.ends_with("b")))
1283 .is_some(),
1284 "the newer look is the one waited on"
1285 );
1286
1287 let ticket = older.ticket();
1288 older.end(Outcome::Failed {
1289 message: "gone".to_string(),
1290 panicked: false,
1291 });
1292 assert_eq!(ended(&rx), ticket);
1293 let ended_older = jobs.end(ticket).expect("its outcome arrives");
1294 assert!(!ended_older.current);
1295 assert!(ended_older.keys.is_none(), "it holds no keys to give back");
1296 assert!(jobs.is_current(newer.ticket()));
1297 assert!(jobs.holds_keys());
1298 drop((newer, pivot));
1299 }
1300
1301 #[test]
1305 fn keys_belong_to_the_job_the_user_waits_on() {
1306 let (mut jobs, rx) = jobs();
1307 let rows = || Job::Rows(crate::InflightCollect::for_tests(0, 100));
1308 let is_rows = |job: &Job| matches!(job, Job::Rows(_));
1309 let ahead = jobs.start(rows(), None);
1310 assert!(!jobs.holds_keys(), "a load-ahead holds no keys");
1311 assert!(jobs.wait_on(is_rows, "Loading buffer..."));
1312 assert!(
1313 jobs.holds_keys(),
1314 "a scroll that caught up with it waits on it"
1315 );
1316 jobs.quiet(is_rows);
1317 assert!(!jobs.holds_keys());
1318 assert!(jobs.wait_on(is_rows, "Loading buffer..."));
1319
1320 let ticket = ahead.ticket();
1321 ahead.end(Outcome::Failed {
1322 message: "no rows".to_string(),
1323 panicked: false,
1324 });
1325 assert_eq!(ended(&rx), ticket);
1326 assert!(jobs.holds_keys(), "until the app has the answer");
1327 let ended = jobs.end(ticket).expect("the outcome is in");
1328 assert_eq!(ended.keys.as_deref(), Some("Loading buffer..."));
1329 assert!(!jobs.holds_keys());
1330 }
1331
1332 #[test]
1335 fn an_owed_page_waits_for_the_generation() {
1336 let (mut jobs, _rx) = jobs();
1337 let hold = jobs.hold();
1338 let owed = |job: &Job| matches!(job, Job::OwedRows { .. });
1339 jobs.owe(
1340 Job::OwedRows {
1341 dataset: 3,
1342 status: "Loading buffer...".to_string(),
1343 },
1344 Some("Loading buffer..."),
1345 );
1346 assert!(jobs.holds_keys(), "the user waits on it");
1347 drop(hold);
1348 assert!(!jobs.would_strand(), "it holds no generation of its own");
1349 assert!(!jobs.in_flight(), "and nothing runs for it");
1350 jobs.advance();
1351 assert!(
1352 matches!(jobs.owed(owed), Some(Job::OwedRows { dataset: 3, .. })),
1353 "an advance does not put it down: it belongs to its dataset"
1354 );
1355 assert!(jobs.take_owed(owed).is_some());
1356 assert!(jobs.take_owed(owed).is_none(), "taken once");
1357 assert!(!jobs.holds_keys());
1358
1359 jobs.owe(
1360 Job::OwedRows {
1361 dataset: 4,
1362 status: String::new(),
1363 },
1364 Some("Loading buffer..."),
1365 );
1366 assert!(jobs.supersede(owed));
1367 assert!(jobs.owed(owed).is_none(), "a cancel drops it");
1368 assert!(!jobs.holds_keys());
1369 }
1370
1371 #[test]
1373 fn a_cancelled_job_says_since_when() {
1374 let (mut jobs, _rx) = jobs();
1375 let run = jobs.start(
1376 Job::Analysis(AnalysisRun {
1377 watch: None,
1378 runs_out: true,
1379 }),
1380 None,
1381 );
1382 assert!(jobs.cancel(|job| matches!(job, Job::Analysis(_))));
1383 assert!(!jobs.is_current(run.ticket()));
1384 let (since, job) = jobs
1385 .cancelled_running(|job| matches!(job, Job::Analysis(_)))
1386 .expect("still running");
1387 assert!(matches!(
1388 job,
1389 Job::Analysis(AnalysisRun { runs_out: true, .. })
1390 ));
1391 let before = since;
1392 jobs.backdate_supersessions(std::time::Duration::from_secs(5));
1393 let (since, _) = jobs
1394 .cancelled_running(|job| matches!(job, Job::Analysis(_)))
1395 .expect("still running");
1396 assert!(since < before);
1397 drop(run);
1398 }
1399
1400 #[test]
1402 fn a_stale_end_leaves_the_newer_job_alone() {
1403 let (mut jobs, rx) = jobs();
1404 let older = jobs.start(Job::Pivot, None);
1405 jobs.advance();
1406 let newer = jobs.start(Job::Pivot, None);
1407 let (older_ticket, newer_ticket) = (older.ticket(), newer.ticket());
1408
1409 older.end(Outcome::Failed {
1410 message: "gone".to_string(),
1411 panicked: false,
1412 });
1413 assert_eq!(ended(&rx), older_ticket);
1414 assert!(jobs.end(older_ticket).is_some_and(|e| !e.current));
1415 assert!(
1416 jobs.is_current(newer_ticket),
1417 "the newer pivot is still waited on"
1418 );
1419 assert!(jobs.would_strand());
1420 drop(newer);
1421 }
1422
1423 #[test]
1426 fn an_outcome_nobody_takes_is_dropped_with_the_owner() {
1427 let (mut jobs, rx) = jobs();
1428 let held = Arc::new(());
1429 let started = jobs.start(Job::Copy, None);
1430 started.end(Outcome::answered(Answer::Probe(held.clone())));
1431 assert!(matches!(rx.try_recv(), Ok(AppEvent::JobEnded(_))));
1432 assert_eq!(Arc::strong_count(&held), 2, "held by the record");
1433 drop(jobs);
1434 assert_eq!(Arc::strong_count(&held), 1);
1435 }
1436
1437 #[test]
1440 fn a_hold_keeps_the_generation_until_dropped() {
1441 let (mut jobs, _rx) = jobs();
1442 let hold = jobs.hold();
1443 assert!(jobs.would_strand());
1444 assert!(jobs.in_flight());
1445 assert!(!jobs.try_advance());
1446 drop(hold);
1447 assert!(!jobs.would_strand());
1448 assert!(!jobs.in_flight());
1449
1450 let passed = jobs.hold();
1451 jobs.advance();
1452 assert!(
1453 !jobs.would_strand(),
1454 "a passed generation's hold holds nothing up"
1455 );
1456 assert!(jobs.running_behind());
1457 drop(passed);
1458 assert!(!jobs.running_behind());
1459 }
1460
1461 #[test]
1464 fn replaceable_jobs_hold_no_lease() {
1465 let (mut jobs, _rx) = jobs();
1466 let rows = jobs.start(Job::Rows(crate::InflightCollect::for_tests(0, 100)), None);
1467 let load = crate::loading::LoadId::for_tests(1);
1468 let path = PathBuf::from("/data");
1469 let look = jobs.start(Job::LookAtDirectory { load, path }, None);
1470 let named = jobs.start(Job::OpenNamed(load), None);
1471 let facts = jobs.start(Job::FileFacts { dataset: 1 }, None);
1472 assert!(!jobs.would_strand());
1473 assert!(jobs.try_advance());
1474 assert!(!jobs.is_current(rows.ticket()));
1475 assert!(!jobs.is_current(look.ticket()));
1476 assert!(!jobs.is_current(named.ticket()));
1477 assert!(
1478 jobs.is_current(facts.ticket()),
1479 "file facts follow the dataset"
1480 );
1481 }
1482
1483 #[test]
1486 fn a_worker_reports_and_runs_on_after_its_answer() {
1487 let (mut jobs, rx) = jobs();
1488 let ticket = spawn(&mut jobs, Job::Export, |worker| {
1489 let report = worker.reporter();
1490 report(Progress::ExportWriting {
1491 phase: "Writing",
1492 bytes: 7,
1493 });
1494 let events = worker.events.clone();
1495 Ok(Answer::Exported(PathBuf::from("x")).then(move || {
1496 let _ = events.send(AppEvent::Wake);
1497 }))
1498 });
1499 let mut seen = Vec::new();
1500 for _ in 0..3 {
1501 seen.push(rx.recv_timeout(Duration::from_secs(30)).expect("an event"));
1502 }
1503 assert!(matches!(
1504 seen[0],
1505 AppEvent::JobProgress {
1506 ticket: t,
1507 progress: Progress::ExportWriting { bytes: 7, .. }
1508 } if t == ticket
1509 ));
1510 assert!(matches!(seen[1], AppEvent::JobEnded(t) if t == ticket));
1511 assert!(matches!(seen[2], AppEvent::Wake), "after the answer");
1512 }
1513
1514 #[test]
1516 fn a_worker_picked_to_die_fails_its_job() {
1517 let (mut jobs, rx) = jobs();
1518 jobs.worker_dies = Some(Box::new(|job| matches!(job, Job::DrillRow)));
1519 let ticket = spawn(&mut jobs, Job::DrillRow, |_| {
1520 Ok(Answer::Exported(PathBuf::from("never")))
1521 });
1522 assert_eq!(ended(&rx), ticket);
1523 assert!(matches!(
1524 jobs.end(ticket).map(|e| e.outcome),
1525 Some(Outcome::Failed { panicked: true, .. })
1526 ));
1527 }
1528
1529 #[test]
1531 fn a_worker_picked_to_wait_runs_once_let_go() {
1532 let (mut jobs, rx) = jobs();
1533 let (waits, release) = crate::tests::worker_waits_once(|job| matches!(job, Job::Pivot));
1534 jobs.worker_waits = waits;
1535 let ticket = spawn(&mut jobs, Job::Pivot, |_| {
1536 Ok(Answer::Exported(PathBuf::from("done")))
1537 });
1538 assert!(rx.try_recv().is_err(), "it waits");
1539 assert!(jobs.is_current(ticket));
1540 release.send(()).unwrap();
1541 assert_eq!(ended(&rx), ticket);
1542 assert!(matches!(
1543 jobs.end(ticket).map(|e| e.outcome),
1544 Some(Outcome::Answered(_))
1545 ));
1546 }
1547}