use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::mpsc::Sender;
use std::sync::{Arc, Mutex};
use std::time::Instant;
use polars::prelude::DataFrame;
use crate::loading::{LoadAnswer, LoadId};
use crate::{AppEvent, OpenOptions, logging};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum JobKind {
Load,
OpenNamed,
LookAtDirectory,
Classify,
Rows,
Analysis,
SampleRows,
SampleDraw,
Pivot,
ViewPivot,
ReshapePreview,
DrillRow,
InspectRow,
InspectJson,
InspectPretty,
InspectUnpack,
OpenValue,
Export,
Copy,
QualityReport,
FileFacts,
ChartExport,
ChartPrepare,
Find,
ValueCounts,
HexOpen,
HexFind,
UnfitCount,
FootersJoin,
JournalDetail,
IndexLines,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct Ticket {
id: u64,
generation: u64,
kind: JobKind,
}
impl Ticket {
pub fn kind(self) -> JobKind {
self.kind
}
}
#[derive(Debug, Clone)]
pub(crate) enum Job {
Load(LoadId),
OpenNamed(LoadId),
LookAtDirectory { load: LoadId, path: PathBuf },
Classify(Classify),
Rows(crate::InflightCollect),
OwedRows { dataset: u64, status: String },
Analysis(AnalysisRun),
SampleRows,
SampleDraw(Box<SampleDraw>),
Pivot,
ViewPivot(Box<(crate::view::SavedView, Option<crate::view::MatchReason>)>),
ReshapePreview { epoch: u64, token: u64 },
DrillRow,
InspectRow { frame: u64, row: usize },
InspectJson { token: u64 },
InspectPretty { token: u64 },
InspectUnpack { token: u64 },
OpenValue,
Export,
Copy,
QualityReport,
FileFacts { dataset: u64 },
ChartExport {
path: PathBuf,
format: crate::chart::chart_export::ChartExportFormat,
},
ChartPrepare(Box<ChartPrep>),
Find(crate::find::FindRun),
ValueCounts,
HexOpen {
origin: crate::app::hex_view::Origin,
fallback: bool,
record_size: Option<usize>,
},
HexFind(crate::app::hex_view::HexFindRun),
UnfitCount { dataset: u64, version: Option<u64> },
FootersJoin { dataset: u64 },
JournalDetail { dataset: u64 },
IndexLines { dataset: u64 },
}
#[derive(Debug, Clone)]
pub(crate) struct Classify {
pub(crate) path: PathBuf,
pub(crate) browsing: Option<PathBuf>,
pub(crate) jump: bool,
}
#[derive(Debug, Clone)]
pub(crate) struct ChartPrep {
pub(crate) request: crate::ChartRequest,
pub(crate) dataset: Option<u64>,
pub(crate) cancel: Arc<std::sync::atomic::AtomicBool>,
}
#[derive(Clone)]
pub(crate) struct SampleDraw {
pub(crate) sample: crate::analysis::sampling::Sample,
pub(crate) rows: Arc<crate::analysis::table_sample::SampleRows>,
pub(crate) watch: crate::analysis::sampling::ReadWatch,
pub(crate) through: bool,
pub(crate) replay: Option<crate::view::ViewSettings>,
pub(crate) then_analyze: bool,
pub(crate) path: Option<crate::analysis::table_sample::DrawPath>,
pub(crate) path_key: String,
pub(crate) schema: Option<polars::prelude::SchemaRef>,
}
impl std::fmt::Debug for SampleDraw {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SampleDraw")
.field("sample", &self.sample)
.field("through", &self.through)
.field("then_analyze", &self.then_analyze)
.finish_non_exhaustive()
}
}
#[derive(Debug, Clone, Default)]
pub(crate) struct AnalysisRun {
pub(crate) watch: Option<crate::analysis::data_quality::QualityWatch>,
pub(crate) runs_out: bool,
}
impl Job {
pub(crate) fn kind(&self) -> JobKind {
match self {
Job::Load(_) => JobKind::Load,
Job::OpenNamed(_) => JobKind::OpenNamed,
Job::LookAtDirectory { .. } => JobKind::LookAtDirectory,
Job::Classify(_) => JobKind::Classify,
Job::Rows(_) | Job::OwedRows { .. } => JobKind::Rows,
Job::Analysis(_) => JobKind::Analysis,
Job::SampleRows => JobKind::SampleRows,
Job::SampleDraw(_) => JobKind::SampleDraw,
Job::Pivot => JobKind::Pivot,
Job::ViewPivot(_) => JobKind::ViewPivot,
Job::ReshapePreview { .. } => JobKind::ReshapePreview,
Job::DrillRow => JobKind::DrillRow,
Job::InspectRow { .. } => JobKind::InspectRow,
Job::InspectJson { .. } => JobKind::InspectJson,
Job::InspectPretty { .. } => JobKind::InspectPretty,
Job::InspectUnpack { .. } => JobKind::InspectUnpack,
Job::OpenValue => JobKind::OpenValue,
Job::Export => JobKind::Export,
Job::Copy => JobKind::Copy,
Job::QualityReport => JobKind::QualityReport,
Job::FileFacts { .. } => JobKind::FileFacts,
Job::ChartExport { .. } => JobKind::ChartExport,
Job::ChartPrepare(_) => JobKind::ChartPrepare,
Job::Find(_) => JobKind::Find,
Job::ValueCounts => JobKind::ValueCounts,
Job::HexOpen { .. } => JobKind::HexOpen,
Job::HexFind(_) => JobKind::HexFind,
Job::UnfitCount { .. } => JobKind::UnfitCount,
Job::FootersJoin { .. } => JobKind::FootersJoin,
Job::JournalDetail { .. } => JobKind::JournalDetail,
Job::IndexLines { .. } => JobKind::IndexLines,
}
}
pub(crate) fn load(&self) -> Option<LoadId> {
match self {
Job::Load(load) | Job::OpenNamed(load) | Job::LookAtDirectory { load, .. } => {
Some(*load)
}
_ => None,
}
}
fn leased(&self) -> bool {
!matches!(
self,
Job::Rows(_)
| Job::SampleDraw(_)
| Job::OwedRows { .. }
| Job::OpenNamed(_)
| Job::LookAtDirectory { .. }
| Job::FileFacts { .. }
| Job::UnfitCount { .. }
| Job::FootersJoin { .. }
| Job::JournalDetail { .. }
| Job::IndexLines { .. }
| Job::ReshapePreview { .. }
| Job::ChartPrepare(_)
)
}
fn follows_the_generation(&self) -> bool {
!matches!(
self,
Job::FileFacts { .. }
| Job::SampleDraw(_)
| Job::UnfitCount { .. }
| Job::FootersJoin { .. }
| Job::JournalDetail { .. }
| Job::IndexLines { .. }
| Job::ChartExport { .. }
| Job::ChartPrepare(_)
| Job::OwedRows { .. }
| Job::ReshapePreview { .. }
)
}
}
pub(crate) enum Answer {
Load(Box<LoadAnswer>),
NamedPaths {
paths: Vec<PathBuf>,
options: Box<OpenOptions>,
directory: Option<PathBuf>,
},
NamedPathMissing(PathBuf),
LookedAt {
kind: crate::home::discover::EntryKind,
holds: Option<Box<crate::home::discover::Holds>>,
options: Box<OpenOptions>,
},
Kind(Option<crate::home::discover::EntryKind>),
Rows(crate::table::CollectResult),
RowsFailed {
message: String,
conversion: Option<Box<crate::error_display::ConversionFailure>>,
},
Analysis(
fn(
&mut crate::analysis::analysis_modal::AnalysisModal,
crate::analysis::statistics::AnalysisResults,
),
crate::analysis::statistics::AnalysisResults,
),
DataQuality {
results: Box<crate::analysis::data_quality::DataQualityResults>,
kept: Option<crate::KeptQualitySample>,
plan: Box<crate::analysis::data_quality::DataQualityPlan>,
},
Sample { df: DataFrame, label: String },
SampleDrawn(crate::analysis::table_sample::Drawn),
Pivoted {
spec: crate::app::modals::pivot_melt_modal::PivotSpec,
pivoted: DataFrame,
},
ViewPivoted(DataFrame),
ReshapePreviewed {
input: Option<crate::app::modals::pivot_melt_modal::PreviewInput>,
result: Result<crate::app::modals::pivot_melt_modal::PreviewFrame, String>,
},
DrillRow { group_index: usize, row: DataFrame },
FieldsRead(DataFrame),
JsonParsed(std::sync::Arc<serde_json::Value>),
Indented(std::sync::Arc<str>),
Unpacked(crate::inspector::inspector_bytes::Decoded),
ValueWritten(crate::inspector::external_open::ExternalOpen),
Exported(PathBuf),
Copied {
payload: crate::clipboard::Payload,
message: String,
},
QualityReportWritten(PathBuf),
ChartExported,
ChartPrepared(
Box<(
crate::chart::chart_plot::PlotData,
Option<crate::chart::chart_modal::ColorCounts>,
)>,
),
FileFacts(crate::widgets::info::FileFacts),
Found(Option<crate::find::Found>),
ValueCounts(Box<crate::analysis::value_counts::ValueCounts>),
HexOpened(Box<crate::app::hex_view::HexSource>),
HexFound(crate::app::hex_view::HexHit),
UnfitCounted(Vec<crate::formats::column_types::Unfit>),
#[cfg(test)]
Probe(Arc<()>),
FootersJoined(Option<Box<crate::table::FootersFound>>),
JournalDescribed(Box<crate::formats::text_formats::Detail>),
LinesIndexed(usize),
}
impl Answer {
pub(crate) fn then(self, after: impl FnOnce() + Send + 'static) -> Answered {
Answered {
answer: self,
then: Some(Box::new(after)),
}
}
}
pub(crate) struct Answered {
answer: Answer,
then: Option<Box<dyn FnOnce() + Send>>,
}
impl From<Answer> for Answered {
fn from(answer: Answer) -> Self {
Self { answer, then: None }
}
}
pub(crate) enum Outcome {
Answered(Box<Answer>),
Failed {
message: String,
panicked: bool,
},
}
impl Outcome {
#[cfg(test)]
pub(crate) fn answered(answer: Answer) -> Self {
Self::Answered(Box::new(answer))
}
}
#[derive(Debug, Clone)]
pub enum Progress {
ExportWriting { phase: &'static str, bytes: u64 },
QualityPhase(crate::analysis::data_quality::QualityPhase),
Finding { rows: usize },
HexFinding { read: u64, total: u64 },
SampleBegun(polars::prelude::SchemaRef),
SampleGrew,
}
pub(crate) struct Ended {
pub(crate) job: Job,
pub(crate) current: bool,
pub(crate) keys: Option<String>,
pub(crate) outcome: Outcome,
}
#[cfg(test)]
pub(crate) type WorkerDies = Box<dyn FnMut(&Job) -> bool + Send>;
#[cfg(test)]
pub(crate) type WorkerWaits = Box<dyn FnMut(&Job) -> Option<std::sync::mpsc::Receiver<()>> + Send>;
type Slot = Arc<Mutex<Option<Outcome>>>;
struct Record {
ticket: Ticket,
job: Job,
slot: Option<Slot>,
keys: Option<String>,
superseded: Option<Instant>,
cancelled: bool,
stale: Arc<std::sync::atomic::AtomicBool>,
}
std::thread_local! {
static RUNNING: std::cell::RefCell<Option<Arc<std::sync::atomic::AtomicBool>>> =
const { std::cell::RefCell::new(None) };
}
pub(crate) fn superseded() -> bool {
RUNNING.with(|running| {
running
.borrow()
.as_ref()
.is_some_and(|stale| stale.load(std::sync::atomic::Ordering::Relaxed))
})
}
impl Record {
fn running(&self) -> bool {
self.slot.is_some()
}
fn current(&self) -> bool {
self.running() && self.superseded.is_none()
}
fn holds(&self, generation: u64) -> bool {
self.current() && self.job.leased() && self.ticket.generation == generation
}
fn supersede(&mut self, now: Instant) {
self.superseded = Some(now);
self.keys = None;
self.stale.store(true, std::sync::atomic::Ordering::Relaxed);
}
}
pub(crate) struct Started {
ticket: Ticket,
slot: Slot,
events: Sender<AppEvent>,
ended: bool,
stale: Arc<std::sync::atomic::AtomicBool>,
#[cfg(test)]
dies: bool,
#[cfg(test)]
waits: Option<std::sync::mpsc::Receiver<()>>,
}
impl Started {
pub(crate) fn ticket(&self) -> Ticket {
self.ticket
}
pub(crate) fn run<F, R>(self, runtime: &tokio::runtime::Handle, work: F)
where
F: FnOnce(&Worker) -> Result<R, String> + Send + 'static,
R: Into<Answered>,
{
let mut started = self;
runtime.spawn_blocking(move || {
let worker = Worker {
ticket: started.ticket,
events: started.events.clone(),
};
#[cfg(test)]
let dies = started.dies;
#[cfg(test)]
if let Some(gate) = started.waits.take() {
let _ = gate.recv();
}
RUNNING.with(|running| *running.borrow_mut() = Some(started.stale.clone()));
let ran = logging::catch_panic(|| {
#[cfg(test)]
if dies {
panic!("worker died");
}
work(&worker)
});
RUNNING.with(|running| *running.borrow_mut() = None);
match ran {
Ok(Ok(answered)) => {
let Answered { answer, then } = answered.into();
started.finish(Outcome::Answered(Box::new(answer)));
if let Some(then) = then {
then();
}
}
Ok(Err(message)) => started.finish(Outcome::Failed {
message,
panicked: false,
}),
Err(message) => started.finish(Outcome::Failed {
message,
panicked: true,
}),
}
});
}
#[cfg(test)]
pub(crate) fn end(mut self, outcome: Outcome) {
self.finish(outcome);
}
fn finish(&mut self, outcome: Outcome) {
if std::mem::replace(&mut self.ended, true) {
return;
}
*self.slot.lock().unwrap_or_else(|e| e.into_inner()) = Some(outcome);
if self.events.send(AppEvent::JobEnded(self.ticket)).is_err() {
drop(self.slot.lock().unwrap_or_else(|e| e.into_inner()).take());
}
}
}
impl Drop for Started {
fn drop(&mut self) {
self.finish(Outcome::Failed {
message: "The background task stopped before it answered".to_string(),
panicked: true,
});
}
}
pub(crate) struct Worker {
ticket: Ticket,
events: Sender<AppEvent>,
}
impl Worker {
pub(crate) fn reporter(&self) -> impl Fn(Progress) + Send + Sync + 'static {
let (ticket, events) = (self.ticket, Mutex::new(self.events.clone()));
move |progress| {
let events = events.lock().unwrap_or_else(|e| e.into_inner());
let _ = events.send(AppEvent::JobProgress { ticket, progress });
}
}
pub(crate) fn send(&self, event: AppEvent) {
let _ = self.events.send(event);
}
}
type HoldCounts = Arc<Mutex<HashMap<u64, usize>>>;
#[must_use]
pub(crate) struct Hold {
generation: u64,
counts: HoldCounts,
}
impl Drop for Hold {
fn drop(&mut self) {
let mut counts = self.counts.lock().unwrap_or_else(|e| e.into_inner());
if let Some(n) = counts.get_mut(&self.generation) {
*n = n.saturating_sub(1);
if *n == 0 {
counts.remove(&self.generation);
}
}
}
}
pub(crate) struct Jobs {
events: Sender<AppEvent>,
generation: u64,
next_id: u64,
records: Vec<Record>,
holds: HoldCounts,
#[cfg(test)]
pub(crate) worker_dies: Option<WorkerDies>,
#[cfg(test)]
pub(crate) worker_waits: Option<WorkerWaits>,
}
impl Jobs {
pub(crate) fn new(events: Sender<AppEvent>) -> Self {
Self {
events,
generation: 0,
next_id: 0,
records: Vec::new(),
holds: Arc::default(),
#[cfg(test)]
worker_dies: None,
#[cfg(test)]
worker_waits: None,
}
}
pub(crate) fn generation(&self) -> u64 {
self.generation
}
fn ticket(&mut self, job: &Job) -> Ticket {
self.next_id = self.next_id.wrapping_add(1);
Ticket {
id: self.next_id,
generation: self.generation,
kind: job.kind(),
}
}
pub(crate) fn start(&mut self, job: Job, keys: Option<&str>) -> Started {
let ticket = self.ticket(&job);
#[cfg(test)]
let dies = self.worker_dies.as_mut().is_some_and(|dies| dies(&job));
#[cfg(test)]
let waits = self.worker_waits.as_mut().and_then(|waits| waits(&job));
let slot = Slot::default();
let stale = Arc::<std::sync::atomic::AtomicBool>::default();
self.records.push(Record {
ticket,
job,
slot: Some(slot.clone()),
keys: keys.map(str::to_string),
superseded: None,
cancelled: false,
stale: stale.clone(),
});
Started {
ticket,
slot,
events: self.events.clone(),
ended: false,
stale,
#[cfg(test)]
dies,
#[cfg(test)]
waits,
}
}
pub(crate) fn owe(&mut self, job: Job, keys: Option<&str>) {
let ticket = self.ticket(&job);
self.records.push(Record {
ticket,
job,
slot: None,
keys: keys.map(str::to_string),
superseded: None,
cancelled: false,
stale: Arc::default(),
});
}
pub(crate) fn owed(&self, which: impl Fn(&Job) -> bool) -> Option<&Job> {
self.records
.iter()
.find(|r| !r.running() && which(&r.job))
.map(|r| &r.job)
}
pub(crate) fn take_owed(&mut self, which: impl Fn(&Job) -> bool) -> Option<Job> {
let at = self
.records
.iter()
.position(|r| !r.running() && which(&r.job))?;
Some(self.records.remove(at).job)
}
pub(crate) fn end(&mut self, ticket: Ticket) -> Option<Ended> {
let at = self.records.iter().position(|r| r.ticket == ticket)?;
let outcome = self.records[at]
.slot
.as_ref()?
.lock()
.unwrap_or_else(|e| e.into_inner())
.take()?;
let record = self.records.remove(at);
Some(Ended {
job: record.job,
current: record.superseded.is_none(),
keys: record.keys,
outcome,
})
}
pub(crate) fn is_current(&self, ticket: Ticket) -> bool {
self.records
.iter()
.any(|r| r.ticket == ticket && r.current())
}
pub(crate) fn current(&self, which: impl Fn(&Job) -> bool) -> Option<(Ticket, &Job)> {
self.records
.iter()
.rev()
.find(|r| r.current() && which(&r.job))
.map(|r| (r.ticket, &r.job))
}
pub(crate) fn current_mut(&mut self, which: impl Fn(&Job) -> bool) -> Option<&mut Job> {
self.records
.iter_mut()
.rev()
.find(|r| r.current() && which(&r.job))
.map(|r| &mut r.job)
}
pub(crate) fn job_mut(&mut self, ticket: Ticket) -> Option<&mut Job> {
self.records
.iter_mut()
.find(|r| r.ticket == ticket)
.map(|r| &mut r.job)
}
pub(crate) fn cancelled_running(
&self,
which: impl Fn(&Job) -> bool,
) -> Option<(Instant, &Job)> {
self.records
.iter()
.filter(|r| r.running() && r.cancelled && which(&r.job))
.filter_map(|r| r.superseded.map(|since| (since, &r.job)))
.max_by_key(|(since, _)| *since)
}
pub(crate) fn running(&self, which: impl Fn(&Job) -> bool) -> bool {
self.records.iter().any(|r| r.running() && which(&r.job))
}
pub(crate) fn shows(&self, status: &str) -> bool {
self.records
.iter()
.any(|r| r.keys.as_deref() == Some(status))
}
pub(crate) fn would_strand(&self) -> bool {
self.records.iter().any(|r| r.holds(self.generation)) || self.held(self.generation)
}
fn held(&self, generation: u64) -> bool {
self.holds
.lock()
.unwrap_or_else(|e| e.into_inner())
.get(&generation)
.is_some_and(|n| *n > 0)
}
pub(crate) fn holds_keys(&self) -> bool {
self.records.iter().any(|r| r.keys.is_some())
}
pub(crate) fn keys_held_only_by(&self, which: impl Fn(&Job) -> bool) -> bool {
self.holds_keys()
&& self
.records
.iter()
.filter(|r| r.keys.is_some())
.all(|r| which(&r.job))
}
pub(crate) fn waited_on(&self, which: impl Fn(&Job) -> bool) -> bool {
self.records
.iter()
.rev()
.find(|r| r.current() && which(&r.job))
.is_some_and(|r| r.keys.is_some())
}
pub(crate) fn waiting_status(&self, which: impl Fn(&Job) -> bool) -> Option<&str> {
self.records
.iter()
.rev()
.filter(|r| r.superseded.is_none() && which(&r.job))
.find_map(|r| r.keys.as_deref())
}
pub(crate) fn wait_on(&mut self, which: impl Fn(&Job) -> bool, status: &str) -> bool {
let Some(record) = self
.records
.iter_mut()
.rev()
.find(|r| r.current() && which(&r.job))
else {
return false;
};
record.keys = Some(status.to_string());
true
}
pub(crate) fn quiet(&mut self, which: impl Fn(&Job) -> bool) -> Vec<String> {
self.records
.iter_mut()
.filter(|r| which(&r.job))
.filter_map(|r| r.keys.take())
.collect()
}
pub(crate) fn try_advance(&mut self) -> bool {
if self.would_strand() {
return false;
}
self.advance();
true
}
pub(crate) fn advance(&mut self) {
self.generation = self.generation.wrapping_add(1);
let now = Instant::now();
for record in &mut self.records {
if record.current() && record.job.follows_the_generation() {
record.supersede(now);
}
}
}
pub(crate) fn cancel(&mut self, which: impl Fn(&Job) -> bool) -> bool {
for record in &mut self.records {
if record.current() && which(&record.job) {
record.cancelled = true;
}
}
self.supersede(which)
}
pub(crate) fn supersede(&mut self, which: impl Fn(&Job) -> bool) -> bool {
let now = Instant::now();
let before = self.records.len();
self.records.retain(|r| r.running() || !which(&r.job));
let mut any = self.records.len() != before;
for record in &mut self.records {
if record.current() && which(&record.job) {
record.supersede(now);
any = true;
}
}
any
}
pub(crate) fn hold(&self) -> Hold {
*self
.holds
.lock()
.unwrap_or_else(|e| e.into_inner())
.entry(self.generation)
.or_default() += 1;
Hold {
generation: self.generation,
counts: self.holds.clone(),
}
}
pub(crate) fn in_flight(&self) -> bool {
self.records.iter().any(|r| r.running() && r.job.leased())
|| self
.holds
.lock()
.unwrap_or_else(|e| e.into_inner())
.values()
.any(|n| *n > 0)
}
pub(crate) fn running_behind(&self) -> bool {
self.records
.iter()
.any(|r| r.running() && r.job.leased() && r.ticket.generation != self.generation)
|| self
.holds
.lock()
.unwrap_or_else(|e| e.into_inner())
.iter()
.any(|(generation, n)| *generation != self.generation && *n > 0)
}
#[cfg(test)]
pub(crate) fn backdate_supersessions(&mut self, by: std::time::Duration) {
for record in &mut self.records {
if let Some(since) = record.superseded.as_mut() {
*since -= by;
}
}
}
}
#[cfg(test)]
mod tests;