use std::collections::HashMap;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{Receiver, TryRecvError};
use std::sync::{Arc, Mutex};
use std::thread;
use crate::hurl::{HurlEntry, RunOutput};
use crate::report::context::ReportRunInputs;
use crate::report::model::{ReportResult, ReportRow};
use crate::report::run::{
DryRunner, EntryRunner, LiveRunner, RowEvent, RunContext, finalize, run_flow_raw,
};
pub enum RunUpdate {
Skeleton(ReportResult),
RowStarted(Vec<(usize, usize)>),
Row(Box<ReportRow>),
Done(ReportResult),
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum RowState {
Scheduled,
Running,
Finished,
}
pub struct RunProgress {
pub states: Vec<RowState>,
pub index: HashMap<Vec<(usize, usize)>, usize>,
pub done: usize,
pub total: usize,
}
#[derive(Clone, PartialEq, Eq, Hash, Debug)]
pub enum RunKey {
Path(std::path::PathBuf),
Id(u64),
}
impl RunKey {
pub fn of(report: &crate::report::Report) -> Self {
match &report.path {
Some(p) => RunKey::Path(p.clone()),
None => RunKey::Id(report.id),
}
}
}
#[derive(Default)]
pub struct ParkedRun {
pub result: Option<ReportResult>,
pub progress: Option<RunProgress>,
pub run: Option<RunHandle>,
pub results_exported: bool,
pub last_export: Option<String>,
pub params: crate::report::params::ParamValues,
pub view: Option<super::report_editor::EditorView>,
pub had_result: bool,
}
impl ParkedRun {
pub fn is_worth_keeping(&self) -> bool {
self.run.is_some() || self.result.is_some()
}
pub fn pump(&mut self) -> bool {
let Some(handle) = self.run.as_mut() else {
return false;
};
if matches!(
drain(handle, &mut self.result, &mut self.progress),
Drained::Disconnected
) {
self.run = None;
return false;
}
if handle.finished() {
self.run = None;
return false;
}
true
}
}
pub struct RunHandle {
cancel: Arc<AtomicBool>,
rx: Receiver<RunUpdate>,
finished: bool,
}
impl RunHandle {
pub fn cancel(&self) {
self.cancel.store(true, Ordering::Relaxed);
}
#[cfg(test)]
pub(crate) fn cancel_flag_for_test(&self) -> Arc<AtomicBool> {
self.cancel.clone()
}
pub fn cancelled(&self) -> bool {
self.cancel.load(Ordering::Relaxed)
}
pub fn finished(&self) -> bool {
self.finished
}
}
impl Drop for RunHandle {
fn drop(&mut self) {
self.cancel.store(true, Ordering::Relaxed);
}
}
#[cfg(test)]
pub(crate) fn test_handle() -> (RunHandle, std::sync::mpsc::Sender<RunUpdate>) {
let (tx, rx) = std::sync::mpsc::channel();
let handle = RunHandle {
cancel: Arc::new(AtomicBool::new(false)),
rx,
finished: false,
};
(handle, tx)
}
struct CancellableRunner<R: EntryRunner> {
inner: R,
cancel: Arc<AtomicBool>,
}
impl<R: EntryRunner> EntryRunner for CancellableRunner<R> {
fn run(&self, base: &HurlEntry, vars: &HashMap<String, String>) -> RunOutput {
if self.cancel.load(Ordering::Relaxed) {
return RunOutput {
entries: Vec::new(),
error: Some("cancelled".to_string()),
};
}
self.inner.run(base, vars)
}
}
pub fn spawn(inputs: ReportRunInputs) -> RunHandle {
let cancel = Arc::new(AtomicBool::new(false));
let cancel_worker = cancel.clone();
let (tx, rx) = std::sync::mpsc::channel();
thread::spawn(move || {
let ReportRunInputs {
flow,
entries,
helpers,
base_vars,
named_envs,
root,
file_root,
language,
params,
} = inputs;
let strings = crate::i18n::Strings::for_language(&language);
let skeleton = {
let dry_ctx = RunContext {
entries: &entries,
helpers: &helpers,
base_vars: base_vars.clone(),
named_envs: named_envs.clone(),
root: root.clone(),
runner: &DryRunner,
strings: &strings,
params: params.clone(),
sink: None,
};
run_flow_raw(&flow, &dry_ctx)
};
if tx.send(RunUpdate::Skeleton(skeleton)).is_err() {
return; }
let runner = CancellableRunner {
inner: LiveRunner { file_root },
cancel: cancel_worker,
};
let row_tx = Mutex::new(tx.clone());
let sink = move |ev: RowEvent| {
if let Ok(tx) = row_tx.lock() {
let msg = match ev {
RowEvent::Started(path) => RunUpdate::RowStarted(path.to_vec()),
RowEvent::Completed(row) => RunUpdate::Row(Box::new(row.clone())),
};
let _ = tx.send(msg);
}
};
let ctx = RunContext {
entries: &entries,
helpers: &helpers,
base_vars,
named_envs,
root,
runner: &runner,
strings: &strings,
params: params.clone(),
sink: Some(&sink),
};
let mut result = run_flow_raw(&flow, &ctx);
finalize(&mut result, &flow, &ctx);
let _ = tx.send(RunUpdate::Done(result));
});
RunHandle {
cancel,
rx,
finished: false,
}
}
pub enum Drained {
Idle,
Progress { done: usize, total: usize },
Done { rows: usize, errors: usize },
Disconnected,
}
pub fn drain(
handle: &mut RunHandle,
result: &mut Option<ReportResult>,
progress: &mut Option<RunProgress>,
) -> Drained {
let mut outcome = Drained::Idle;
loop {
match handle.rx.try_recv() {
Ok(update) => {
if let Some(step) = apply(handle, update, result, progress) {
outcome = step;
}
}
Err(TryRecvError::Empty) => break,
Err(TryRecvError::Disconnected) => {
*progress = None;
if !handle.finished {
return Drained::Disconnected;
}
break;
}
}
}
outcome
}
fn apply(
handle: &mut RunHandle,
update: RunUpdate,
result: &mut Option<ReportResult>,
progress: &mut Option<RunProgress>,
) -> Option<Drained> {
let cancelled = handle.cancelled();
match update {
RunUpdate::Skeleton(skeleton) => {
if cancelled {
return None;
}
let mut skeleton = skeleton;
let total = skeleton.rows.len();
skeleton.pending = (0..total).collect();
let index = skeleton
.rows
.iter()
.enumerate()
.map(|(i, row)| (row.path.clone(), i))
.collect();
*result = Some(skeleton);
*progress = Some(RunProgress {
states: vec![RowState::Scheduled; total],
index,
done: 0,
total,
});
Some(Drained::Progress { done: 0, total })
}
RunUpdate::RowStarted(path) => {
if cancelled {
return None;
}
if let Some(prog) = progress.as_mut()
&& let Some(&ri) = prog.index.get(&path)
&& prog.states.get(ri) == Some(&RowState::Scheduled)
{
prog.states[ri] = RowState::Running;
}
None
}
RunUpdate::Row(row) => {
if cancelled {
return None;
}
let (Some(res), Some(prog)) = (result.as_mut(), progress.as_mut()) else {
return None;
};
if let Some(&ri) = prog.index.get(&row.path)
&& ri < res.rows.len()
{
res.rows[ri] = *row;
res.pending.remove(&ri);
if prog.states[ri] != RowState::Finished {
prog.states[ri] = RowState::Finished;
prog.done += 1;
}
}
Some(Drained::Progress {
done: prog.done,
total: prog.total,
})
}
RunUpdate::Done(finalized) => {
handle.finished = true;
*progress = None;
if cancelled {
return Some(Drained::Done {
rows: result.as_ref().map(|r| r.rows.len()).unwrap_or(0),
errors: result.as_ref().map(|r| r.errors.len()).unwrap_or(0),
});
}
let rows = finalized.rows.len();
let errors = finalized.errors.len();
*result = Some(finalized);
Some(Drained::Done { rows, errors })
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::report::model::{ReportResult, ReportRow};
fn row(path: Vec<(usize, usize)>, status: &str) -> ReportRow {
let mut r = ReportRow::default();
r.path = path;
r.cells.insert("status".to_string(), status.to_string());
r
}
fn skeleton(n: usize) -> ReportResult {
let mut res = ReportResult::default();
res.column_order = vec!["status".to_string()];
res.rows = (0..n).map(|i| row(vec![(0, i)], "")).collect();
res
}
#[test]
fn streamed_rows_fill_their_slots_by_path_and_advance_progress() {
let (mut handle, tx) = test_handle();
let mut result = None;
let mut progress = None;
tx.send(RunUpdate::Skeleton(skeleton(3))).unwrap();
assert!(matches!(
drain(&mut handle, &mut result, &mut progress),
Drained::Progress { done: 0, total: 3 }
));
assert_eq!(progress.as_ref().unwrap().states.len(), 3);
assert!(
progress
.as_ref()
.unwrap()
.states
.iter()
.all(|s| *s == RowState::Scheduled)
);
tx.send(RunUpdate::RowStarted(vec![(0, 2)])).unwrap();
tx.send(RunUpdate::Row(Box::new(row(vec![(0, 2)], "200"))))
.unwrap();
assert!(matches!(
drain(&mut handle, &mut result, &mut progress),
Drained::Progress { done: 1, total: 3 }
));
let res = result.as_ref().unwrap();
assert_eq!(
res.rows[2].cells.get("status").map(String::as_str),
Some("200")
);
assert_eq!(progress.as_ref().unwrap().states[2], RowState::Finished);
assert_eq!(progress.as_ref().unwrap().states[0], RowState::Scheduled);
}
#[test]
fn done_replaces_the_grid_and_clears_progress() {
let (mut handle, tx) = test_handle();
let mut result = Some(skeleton(2));
let mut progress = Some(RunProgress {
states: vec![RowState::Scheduled; 2],
index: Default::default(),
done: 0,
total: 2,
});
let mut finalized = ReportResult::default();
finalized.column_order = vec!["status".to_string()];
finalized.rows = vec![row(vec![], "OK")];
tx.send(RunUpdate::Done(finalized)).unwrap();
assert!(matches!(
drain(&mut handle, &mut result, &mut progress),
Drained::Done { rows: 1, errors: 0 }
));
assert!(progress.is_none());
assert!(handle.finished());
assert_eq!(result.as_ref().unwrap().rows.len(), 1);
}
#[test]
fn a_cancelled_run_keeps_its_partial_grid_and_ignores_late_rows() {
let (mut handle, tx) = test_handle();
let mut result = None;
let mut progress = None;
tx.send(RunUpdate::Skeleton(skeleton(2))).unwrap();
drain(&mut handle, &mut result, &mut progress);
tx.send(RunUpdate::Row(Box::new(row(vec![(0, 0)], "200"))))
.unwrap();
drain(&mut handle, &mut result, &mut progress);
handle.cancel();
tx.send(RunUpdate::Row(Box::new(row(vec![(0, 1)], "500"))))
.unwrap();
tx.send(RunUpdate::Done(ReportResult::default())).unwrap();
drain(&mut handle, &mut result, &mut progress);
let res = result.as_ref().unwrap();
assert_eq!(
res.rows[0].cells.get("status").map(String::as_str),
Some("200")
);
assert_eq!(
res.rows[1].cells.get("status").map(String::as_str),
Some("")
);
assert!(handle.finished());
}
}