use std::collections::HashSet;
use std::fmt;
use std::io::{self, IsTerminal, Write};
use std::sync::atomic::{AtomicBool, AtomicU8, AtomicU64, Ordering};
use std::sync::mpsc::{self, Receiver, RecvTimeoutError, SyncSender};
use std::sync::{Arc, Mutex};
use std::thread::{self, JoinHandle};
use std::time::Duration;
use crossterm::cursor::{Hide, MoveTo, Show};
use crossterm::event::{
self, DisableMouseCapture, EnableMouseCapture, Event, KeyCode, KeyEventKind, KeyModifiers,
};
use crossterm::execute;
use crossterm::terminal::{
Clear, ClearType, EnterAlternateScreen, LeaveAlternateScreen, disable_raw_mode, enable_raw_mode,
};
use indicatif::{MultiProgress, ProgressBar, ProgressDrawTarget, ProgressStyle};
use serde_json::Value;
use super::error::ReportingError;
use super::phase::{Phase, TaskDisplayKind, TaskKey};
const REFRESH_INTERVAL: Duration = Duration::from_millis(100);
const MESSAGE_CAPACITY: usize = 256;
static TERMINAL_OWNED: AtomicBool = AtomicBool::new(false);
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum TaskStatus {
Pending,
Running,
Completed,
Reused,
Failed,
Cancelled,
Skipped,
}
impl TaskStatus {
pub const fn as_str(self) -> &'static str {
match self {
Self::Pending => "pending",
Self::Running => "running",
Self::Completed => "completed",
Self::Failed => "failed",
Self::Reused => "reused",
Self::Cancelled => "cancelled",
Self::Skipped => "skipped",
}
}
fn encode(self) -> u8 {
match self {
Self::Pending => 0,
Self::Running => 1,
Self::Completed => 2,
Self::Failed => 3,
Self::Reused => 4,
Self::Cancelled => 5,
Self::Skipped => 6,
}
}
fn decode(value: u8) -> Self {
match value {
0 => Self::Pending,
1 => Self::Running,
2 => Self::Completed,
3 => Self::Failed,
4 => Self::Reused,
5 => Self::Cancelled,
6 => Self::Skipped,
_ => unreachable!("task status is written only through TaskStatus::encode"),
}
}
fn label(self) -> &'static str {
self.as_str()
}
}
#[derive(Clone, Debug)]
pub struct TaskIdentity {
label: Arc<str>,
key: TaskKey,
configuration: crate::configuration::TaskConfig,
}
impl TaskIdentity {
pub fn label(&self) -> &str {
&self.label
}
pub fn len(&self) -> usize {
self.configuration.parameters().len()
}
pub fn is_empty(&self) -> bool {
self.configuration.parameters().is_empty()
}
pub fn value(&self, key: &str) -> Option<&Value> {
self.configuration.value(key)
}
pub fn iter(&self) -> Box<dyn Iterator<Item = (&str, &Value)> + '_> {
Box::new(self.configuration.parameters().iter())
}
pub fn task_key(&self) -> &TaskKey {
&self.key
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct ProgressSummary {
total: u64,
pending: u64,
running: u64,
completed: u64,
reused: u64,
failed: u64,
cancelled: u64,
skipped: u64,
}
impl ProgressSummary {
pub fn total(&self) -> u64 {
self.total
}
pub fn pending(&self) -> u64 {
self.pending
}
pub fn running(&self) -> u64 {
self.running
}
pub fn completed(&self) -> u64 {
self.completed
}
pub fn reused(&self) -> u64 {
self.reused
}
pub fn failed(&self) -> u64 {
self.failed
}
pub fn cancelled(&self) -> u64 {
self.cancelled
}
pub fn skipped(&self) -> u64 {
self.skipped
}
pub fn is_success(&self) -> bool {
self.completed + self.reused == self.total
&& self.pending == 0
&& self.running == 0
&& self.failed == 0
&& self.cancelled == 0
&& self.skipped == 0
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum OutputMode {
Auto,
Terminal,
Plain,
Hidden,
}
pub(crate) struct RuntimeReporterBuilder {
slots: Arc<[Arc<ProgressSlot>]>,
output: OutputMode,
cancellation: Option<CancellationToken>,
}
impl RuntimeReporterBuilder {
pub(crate) fn cancellation_token(mut self, cancellation: CancellationToken) -> Self {
self.cancellation = Some(cancellation);
self
}
pub fn terminal(mut self) -> Self {
self.output = OutputMode::Terminal;
self
}
pub fn plain(mut self) -> Self {
self.output = OutputMode::Plain;
self
}
pub fn hidden(mut self) -> Self {
self.output = OutputMode::Hidden;
self
}
pub fn start(self) -> Result<RuntimeReporter, ReportingError> {
start_reporter(self.slots, self.output, self.cancellation)
}
}
impl fmt::Debug for RuntimeReporterBuilder {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("RuntimeReporterBuilder")
.field("tasks", &self.slots.len())
.field("output", &self.output)
.finish_non_exhaustive()
}
}
pub(crate) struct RuntimeReporter {
inner: Arc<ReporterInner>,
renderer: Option<JoinHandle<()>>,
finished: bool,
}
impl RuntimeReporter {
pub fn for_phase(
phase: &Phase,
heading: &str,
) -> Result<RuntimeReporterBuilder, ReportingError> {
Ok(RuntimeReporterBuilder {
slots: build_phase_slots(std::slice::from_ref(phase), Some(heading))?,
output: OutputMode::Auto,
cancellation: None,
})
}
pub fn start_progress(
&self,
key: &TaskKey,
initial_iteration: u64,
target_iteration: Option<u64>,
) -> Result<TaskProgress, ReportingError> {
let slot = managed_slot(&self.inner.slots, key)?;
if slot.display_kind != TaskDisplayKind::Progress {
return Err(kind_mismatch(&slot, "progress"));
}
start_slot(&self.inner, slot, initial_iteration, target_iteration)
}
pub fn start_activity(&self, key: &TaskKey) -> Result<ActivityTask, ReportingError> {
let slot = managed_slot(&self.inner.slots, key)?;
if slot.display_kind != TaskDisplayKind::Activity {
return Err(kind_mismatch(&slot, "activity"));
}
start_activity_slot(&self.inner, slot)
}
pub fn mark_reused(&self, key: &TaskKey) -> Result<(), ReportingError> {
let slot = managed_slot(&self.inner.slots, key)?;
mark_slot_reused(&slot)
}
pub(crate) fn mark_skipped(&self, key: &TaskKey) -> Result<(), ReportingError> {
let slot = managed_slot(&self.inner.slots, key)?;
mark_pending_terminal(&slot, TaskStatus::Skipped, "skipped")
}
pub(crate) fn mark_cancelled(&self, key: &TaskKey) -> Result<(), ReportingError> {
let slot = managed_slot(&self.inner.slots, key)?;
mark_pending_terminal(&slot, TaskStatus::Cancelled, "cancelled")
}
pub(crate) fn mark_delayed(&self, key: &TaskKey, rank: usize) -> Result<(), ReportingError> {
let slot = managed_slot(&self.inner.slots, key)?;
if TaskStatus::decode(slot.status.load(Ordering::Acquire)) != TaskStatus::Pending {
return Err(ReportingError::TaskAlreadyStarted {
identity: slot.identity.label().to_owned(),
});
}
*lock(&slot.detail) = format!("delayed start (rank {rank})").into_boxed_str();
Ok(())
}
pub(crate) fn request_cancellation(&self) {
self.inner.cancelled.store(true, Ordering::Release);
}
pub(crate) fn is_cancelled(&self) -> bool {
self.inner.cancelled.load(Ordering::Acquire)
}
pub(crate) fn cancellation_flag(&self) -> Arc<AtomicBool> {
Arc::clone(&self.inner.cancelled)
}
pub fn summary(&self) -> ProgressSummary {
summarize(&self.inner.slots)
}
pub(crate) fn task_execution_snapshots(&self) -> Vec<TaskExecutionSnapshot> {
self.inner
.slots
.iter()
.map(|slot| TaskExecutionSnapshot {
key: slot.identity.task_key().clone(),
status: TaskStatus::decode(slot.status.load(Ordering::Acquire)),
current_iteration: (slot.display_kind == TaskDisplayKind::Progress)
.then(|| slot.current.load(Ordering::Relaxed)),
target_iteration: (slot.display_kind == TaskDisplayKind::Progress
&& slot.target_known.load(Ordering::Acquire))
.then(|| slot.target.load(Ordering::Relaxed)),
})
.collect()
}
pub fn complete(
mut self,
message: impl Into<String>,
) -> Result<ProgressSummary, ReportingError> {
let summary = self.summary();
if !summary.is_success() {
self.stop(false, "workflow did not complete".to_owned())?;
return Err(ReportingError::IncompleteProgress {
pending: summary.pending,
running: summary.running,
failed: summary.failed,
});
}
self.stop(true, message.into())?;
Ok(summary)
}
pub fn fail(mut self, message: impl Into<String>) -> Result<ProgressSummary, ReportingError> {
let summary = self.summary();
self.stop(false, message.into())?;
Ok(summary)
}
fn stop(&mut self, success: bool, message: String) -> Result<(), ReportingError> {
self.inner
.events
.send(RenderEvent::Stop { success, message })
.map_err(|_| ReportingError::RendererUnavailable)?;
self.finished = true;
if self
.renderer
.take()
.expect("an unfinished reporter owns one renderer")
.join()
.is_err()
{
return Err(ReportingError::RendererPanicked);
}
Ok(())
}
}
pub(crate) struct TaskExecutionSnapshot {
pub(crate) key: TaskKey,
pub(crate) status: TaskStatus,
pub(crate) current_iteration: Option<u64>,
pub(crate) target_iteration: Option<u64>,
}
impl fmt::Debug for RuntimeReporter {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("RuntimeReporter")
.field("tasks", &self.inner.slots.len())
.field("summary", &self.summary())
.finish_non_exhaustive()
}
}
impl Drop for RuntimeReporter {
fn drop(&mut self) {
if self.finished {
return;
}
let _ = self.inner.events.send(RenderEvent::Stop {
success: false,
message: "progress reporter dropped before completion".to_owned(),
});
if let Some(renderer) = self.renderer.take() {
let _ = renderer.join();
}
}
}
pub struct TaskProgress {
slot: Arc<ProgressSlot>,
events: SyncSender<RenderEvent>,
cancelled: Arc<AtomicBool>,
active: bool,
}
#[derive(Clone, Debug)]
pub struct CancellationToken(Arc<AtomicBool>);
impl CancellationToken {
pub(crate) fn new() -> Self {
Self(Arc::new(AtomicBool::new(false)))
}
pub(crate) fn shared(&self) -> Arc<AtomicBool> {
Arc::clone(&self.0)
}
pub fn is_cancelled(&self) -> bool {
self.0.load(Ordering::Acquire)
}
pub fn cancel(&self) {
self.0.store(true, Ordering::Release);
}
}
impl TaskProgress {
pub fn is_cancelled(&self) -> bool {
self.cancelled.load(Ordering::Acquire)
}
pub fn identity(&self) -> &TaskIdentity {
&self.slot.identity
}
pub fn current_iteration(&self) -> u64 {
self.slot.current.load(Ordering::Relaxed)
}
pub fn target_iteration(&self) -> Option<u64> {
if self.slot.target_known.load(Ordering::Acquire) {
Some(self.slot.target.load(Ordering::Relaxed))
} else {
None
}
}
pub fn set_target_iteration(&self, target: u64) -> Result<(), ReportingError> {
let current = self.current_iteration();
if current > target {
return Err(ReportingError::InitialIterationBeyondTarget {
identity: self.identity().label().to_owned(),
initial: current,
target,
});
}
self.slot.target.store(target, Ordering::Relaxed);
self.slot.target_known.store(true, Ordering::Release);
Ok(())
}
pub fn status(&self) -> TaskStatus {
TaskStatus::decode(self.slot.status.load(Ordering::Acquire))
}
pub fn set_iteration(&self, iteration: u64) -> Result<(), ReportingError> {
if let Some(target) = self.target_iteration().filter(|target| iteration > *target) {
return Err(ReportingError::IterationBeyondTarget {
identity: self.identity().label().to_owned(),
iteration,
target,
});
}
let previous = self.slot.current.fetch_max(iteration, Ordering::Relaxed);
if iteration < previous {
return Err(ReportingError::IterationRegressed {
identity: self.identity().label().to_owned(),
current: previous,
attempted: iteration,
});
}
Ok(())
}
pub fn should_continue(&self, iteration: u64) -> Result<bool, ReportingError> {
self.set_iteration(iteration)?;
Ok(!self.is_cancelled()
&& self
.target_iteration()
.is_none_or(|target| iteration < target))
}
pub fn set_detail(&self, detail: impl Into<String>) {
*lock(&self.slot.detail) = detail.into().into_boxed_str();
}
pub fn report(&self, message: impl Into<String>) -> Result<(), ReportingError> {
self.events
.send(RenderEvent::TaskMessage {
identity: self.identity().label().to_owned(),
message: message.into(),
})
.map_err(|_| ReportingError::RendererUnavailable)
}
pub fn complete(mut self, reason: Option<String>) -> Result<(), ReportingError> {
if reason.is_none()
&& let Some(target) = self.target_iteration()
{
let current = self.current_iteration();
if current != target {
*lock(&self.slot.detail) = "target not reached".into();
self.slot
.status
.store(TaskStatus::Failed.encode(), Ordering::Release);
self.active = false;
return Err(ReportingError::TargetIterationNotReached {
identity: self.identity().label().to_owned(),
current,
target,
});
}
}
*lock(&self.slot.detail) = reason
.unwrap_or_else(|| "completed".to_owned())
.into_boxed_str();
self.slot
.status
.store(TaskStatus::Completed.encode(), Ordering::Release);
self.active = false;
Ok(())
}
pub fn fail(mut self, reason: impl Into<String>) {
*lock(&self.slot.detail) = reason.into().into_boxed_str();
self.slot
.status
.store(TaskStatus::Failed.encode(), Ordering::Release);
self.active = false;
}
pub(crate) fn cancel(mut self, reason: impl Into<String>) {
*lock(&self.slot.detail) = reason.into().into_boxed_str();
self.slot
.status
.store(TaskStatus::Cancelled.encode(), Ordering::Release);
self.active = false;
}
}
impl fmt::Debug for TaskProgress {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("TaskProgress")
.field("identity", &self.identity().label())
.field("current_iteration", &self.current_iteration())
.field("target_iteration", &self.target_iteration())
.field("active", &self.active)
.finish_non_exhaustive()
}
}
impl Drop for TaskProgress {
fn drop(&mut self) {
if self.active {
*lock(&self.slot.detail) = "interrupted".into();
self.slot
.status
.store(TaskStatus::Failed.encode(), Ordering::Release);
}
}
}
pub struct ActivityTask {
slot: Arc<ProgressSlot>,
events: SyncSender<RenderEvent>,
cancelled: Arc<AtomicBool>,
active: bool,
}
impl ActivityTask {
pub fn is_cancelled(&self) -> bool {
self.cancelled.load(Ordering::Acquire)
}
pub fn identity(&self) -> &TaskIdentity {
&self.slot.identity
}
pub fn status(&self) -> TaskStatus {
TaskStatus::decode(self.slot.status.load(Ordering::Acquire))
}
pub fn set_detail(&self, detail: impl Into<String>) {
*lock(&self.slot.detail) = detail.into().into_boxed_str();
}
pub fn report(&self, message: impl Into<String>) -> Result<(), ReportingError> {
self.events
.send(RenderEvent::TaskMessage {
identity: self.identity().label().to_owned(),
message: message.into(),
})
.map_err(|_| ReportingError::RendererUnavailable)
}
pub fn complete(mut self) {
*lock(&self.slot.detail) = "completed".into();
self.slot
.status
.store(TaskStatus::Completed.encode(), Ordering::Release);
self.active = false;
}
pub fn fail(mut self, reason: impl Into<String>) {
*lock(&self.slot.detail) = reason.into().into_boxed_str();
self.slot
.status
.store(TaskStatus::Failed.encode(), Ordering::Release);
self.active = false;
}
pub(crate) fn cancel(mut self, reason: impl Into<String>) {
*lock(&self.slot.detail) = reason.into().into_boxed_str();
self.slot
.status
.store(TaskStatus::Cancelled.encode(), Ordering::Release);
self.active = false;
}
}
impl fmt::Debug for ActivityTask {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("ActivityTask")
.field("identity", &self.identity().label())
.field("status", &self.status())
.field("active", &self.active)
.finish_non_exhaustive()
}
}
impl Drop for ActivityTask {
fn drop(&mut self) {
if self.active {
*lock(&self.slot.detail) = "interrupted".into();
self.slot
.status
.store(TaskStatus::Failed.encode(), Ordering::Release);
}
}
}
struct ReporterInner {
slots: Arc<[Arc<ProgressSlot>]>,
events: SyncSender<RenderEvent>,
cancelled: Arc<AtomicBool>,
}
struct ProgressSlot {
identity: TaskIdentity,
phase_label: Option<Arc<str>>,
display_kind: TaskDisplayKind,
current: AtomicU64,
target: AtomicU64,
target_known: AtomicBool,
started: AtomicBool,
status: AtomicU8,
detail: Mutex<Box<str>>,
}
enum RenderEvent {
TaskMessage { identity: String, message: String },
Stop { success: bool, message: String },
}
struct TerminalLease;
impl Drop for TerminalLease {
fn drop(&mut self) {
TERMINAL_OWNED.store(false, Ordering::Release);
}
}
fn acquire_terminal() -> Result<(), ReportingError> {
TERMINAL_OWNED
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.map(|_| ())
.map_err(|_| ReportingError::TerminalAlreadyOwned)
}
fn resolve_output(output: OutputMode) -> OutputMode {
match output {
OutputMode::Auto if io::stderr().is_terminal() && io::stdin().is_terminal() => {
OutputMode::Terminal
}
OutputMode::Auto => OutputMode::Plain,
explicit => explicit,
}
}
fn start_reporter(
slots: Arc<[Arc<ProgressSlot>]>,
requested_output: OutputMode,
cancellation: Option<CancellationToken>,
) -> Result<RuntimeReporter, ReportingError> {
acquire_terminal()?;
let lease = TerminalLease;
let cancelled = cancellation
.map(|token| token.shared())
.unwrap_or_else(|| Arc::new(AtomicBool::new(false)));
let mut output = resolve_output(requested_output);
let terminal = if output == OutputMode::Terminal {
match TerminalSession::enter(Arc::clone(&cancelled)) {
Ok(terminal) => Some(terminal),
Err(_) if requested_output == OutputMode::Auto => {
output = OutputMode::Plain;
None
}
Err(error) => return Err(error),
}
} else {
None
};
let (events, receiver) = mpsc::sync_channel(MESSAGE_CAPACITY);
let renderer_slots = Arc::clone(&slots);
let renderer = match thread::Builder::new()
.name("scientific-workflow-progress".to_owned())
.spawn(move || render(receiver, renderer_slots, output, terminal, lease))
{
Ok(renderer) => renderer,
Err(source) => return Err(ReportingError::StartRenderer { source }),
};
Ok(RuntimeReporter {
inner: Arc::new(ReporterInner {
slots,
events,
cancelled,
}),
renderer: Some(renderer),
finished: false,
})
}
fn managed_slot(
slots: &[Arc<ProgressSlot>],
key: &TaskKey,
) -> Result<Arc<ProgressSlot>, ReportingError> {
slots
.iter()
.find(|slot| slot.identity.task_key() == key)
.cloned()
.ok_or_else(|| ReportingError::UnknownManagedTask {
task: key.to_string(),
})
}
fn kind_mismatch(slot: &ProgressSlot, requested: &'static str) -> ReportingError {
ReportingError::ManagedTaskKindMismatch {
task: slot.identity.task_key().to_string(),
requested,
actual: match slot.display_kind {
TaskDisplayKind::Progress => "progress",
TaskDisplayKind::Activity => "activity",
},
}
}
fn mark_slot_reused(slot: &ProgressSlot) -> Result<(), ReportingError> {
slot.status
.compare_exchange(
TaskStatus::Pending.encode(),
TaskStatus::Reused.encode(),
Ordering::AcqRel,
Ordering::Acquire,
)
.map_err(|_| ReportingError::TaskAlreadyStarted {
identity: slot.identity.label().to_owned(),
})?;
*lock(&slot.detail) = "reused".into();
Ok(())
}
fn mark_pending_terminal(
slot: &ProgressSlot,
status: TaskStatus,
detail: &'static str,
) -> Result<(), ReportingError> {
slot.status
.compare_exchange(
TaskStatus::Pending.encode(),
status.encode(),
Ordering::AcqRel,
Ordering::Acquire,
)
.map_err(|_| ReportingError::TaskAlreadyStarted {
identity: slot.identity.label().to_owned(),
})?;
*lock(&slot.detail) = detail.into();
Ok(())
}
fn start_slot(
inner: &ReporterInner,
slot: Arc<ProgressSlot>,
initial_iteration: u64,
target_iteration: Option<u64>,
) -> Result<TaskProgress, ReportingError> {
if let Some(target) = target_iteration.filter(|target| initial_iteration > *target) {
return Err(ReportingError::InitialIterationBeyondTarget {
identity: slot.identity.label().to_owned(),
initial: initial_iteration,
target,
});
}
slot.status
.compare_exchange(
TaskStatus::Pending.encode(),
TaskStatus::Running.encode(),
Ordering::AcqRel,
Ordering::Acquire,
)
.map_err(|_| ReportingError::TaskAlreadyStarted {
identity: slot.identity.label().to_owned(),
})?;
slot.started.store(true, Ordering::Release);
slot.current.store(initial_iteration, Ordering::Relaxed);
if let Some(target) = target_iteration {
slot.target.store(target, Ordering::Relaxed);
slot.target_known.store(true, Ordering::Release);
} else {
slot.target_known.store(false, Ordering::Release);
}
*lock(&slot.detail) = "running".into();
Ok(TaskProgress {
slot,
events: inner.events.clone(),
cancelled: Arc::clone(&inner.cancelled),
active: true,
})
}
fn start_activity_slot(
inner: &ReporterInner,
slot: Arc<ProgressSlot>,
) -> Result<ActivityTask, ReportingError> {
slot.status
.compare_exchange(
TaskStatus::Pending.encode(),
TaskStatus::Running.encode(),
Ordering::AcqRel,
Ordering::Acquire,
)
.map_err(|_| ReportingError::TaskAlreadyStarted {
identity: slot.identity.label().to_owned(),
})?;
slot.started.store(true, Ordering::Release);
slot.target_known.store(false, Ordering::Release);
*lock(&slot.detail) = "running".into();
Ok(ActivityTask {
slot,
events: inner.events.clone(),
cancelled: Arc::clone(&inner.cancelled),
active: true,
})
}
fn build_phase_slots(
phases: &[Phase],
heading: Option<&str>,
) -> Result<Arc<[Arc<ProgressSlot>]>, ReportingError> {
if phases.is_empty() {
return Err(ReportingError::EmptyPhaseSet);
}
let mut phase_ids = HashSet::with_capacity(phases.len());
let capacity = phases.iter().map(|phase| phase.tasks().len()).sum();
let mut slots = Vec::with_capacity(capacity);
for phase in phases {
if !phase_ids.insert(phase.id()) {
return Err(ReportingError::DuplicatePhaseId {
phase: phase.id().get(),
});
}
let phase_label: Arc<str> = heading.unwrap_or_else(|| phase.label()).into();
for task in phase.tasks() {
slots.push(Arc::new(ProgressSlot {
identity: TaskIdentity {
label: task.label().into(),
key: task.key().clone(),
configuration: task.configuration().clone(),
},
phase_label: Some(Arc::clone(&phase_label)),
display_kind: task.display_kind(),
current: AtomicU64::new(0),
target: AtomicU64::new(0),
target_known: AtomicBool::new(false),
started: AtomicBool::new(false),
status: AtomicU8::new(TaskStatus::Pending.encode()),
detail: Mutex::new("pending".into()),
}));
}
}
Ok(slots.into())
}
fn summarize(slots: &[Arc<ProgressSlot>]) -> ProgressSummary {
let mut summary = ProgressSummary {
total: u64::try_from(slots.len()).expect("slot count originated from a u64 task count"),
pending: 0,
running: 0,
completed: 0,
reused: 0,
failed: 0,
cancelled: 0,
skipped: 0,
};
for slot in slots {
match TaskStatus::decode(slot.status.load(Ordering::Acquire)) {
TaskStatus::Pending => summary.pending += 1,
TaskStatus::Running => summary.running += 1,
TaskStatus::Completed => summary.completed += 1,
TaskStatus::Reused => summary.reused += 1,
TaskStatus::Failed => summary.failed += 1,
TaskStatus::Cancelled => summary.cancelled += 1,
TaskStatus::Skipped => summary.skipped += 1,
}
}
summary
}
struct TerminalSession {
stop: Arc<AtomicBool>,
input: Option<JoinHandle<()>>,
}
impl TerminalSession {
fn enter(cancelled: Arc<AtomicBool>) -> Result<Self, ReportingError> {
enable_raw_mode().map_err(|source| ReportingError::TerminalSetup {
operation: "enable raw input mode",
source,
})?;
let mut stderr = io::stderr();
if let Err(source) = execute!(
stderr,
EnterAlternateScreen,
Clear(ClearType::All),
MoveTo(0, 0),
Hide,
EnableMouseCapture
) {
let _ = execute!(stderr, DisableMouseCapture, Show, LeaveAlternateScreen);
let _ = disable_raw_mode();
return Err(ReportingError::TerminalSetup {
operation: "enter the isolated terminal screen",
source,
});
}
let stop = Arc::new(AtomicBool::new(false));
let input_stop = Arc::clone(&stop);
let input = match thread::Builder::new()
.name("scientific-workflow-input".to_owned())
.spawn(move || drain_terminal_input(input_stop, cancelled))
{
Ok(input) => input,
Err(source) => {
let _ = execute!(stderr, DisableMouseCapture, Show, LeaveAlternateScreen);
let _ = disable_raw_mode();
return Err(ReportingError::TerminalSetup {
operation: "start the isolated-screen input drain",
source,
});
}
};
Ok(Self {
stop,
input: Some(input),
})
}
}
impl Drop for TerminalSession {
fn drop(&mut self) {
self.stop.store(true, Ordering::Release);
if let Some(input) = self.input.take() {
let _ = input.join();
}
let mut stderr = io::stderr();
let _ = execute!(stderr, DisableMouseCapture, Show, LeaveAlternateScreen);
let _ = stderr.flush();
let _ = disable_raw_mode();
}
}
fn drain_terminal_input(stop: Arc<AtomicBool>, cancelled: Arc<AtomicBool>) {
while !stop.load(Ordering::Acquire) {
match event::poll(Duration::from_millis(50)) {
Ok(true) => match event::read() {
Ok(Event::Key(key))
if key.kind == KeyEventKind::Press
&& key.code == KeyCode::Char('c')
&& key.modifiers.contains(KeyModifiers::CONTROL) =>
{
cancelled.store(true, Ordering::Release);
}
Ok(_) => {}
Err(_) => {
cancelled.store(true, Ordering::Release);
break;
}
},
Ok(false) => {}
Err(_) => {
cancelled.store(true, Ordering::Release);
break;
}
}
}
}
fn render(
receiver: Receiver<RenderEvent>,
slots: Arc<[Arc<ProgressSlot>]>,
output: OutputMode,
mut terminal_session: Option<TerminalSession>,
_lease: TerminalLease,
) {
let mut terminal = (output == OutputMode::Terminal).then(|| TerminalDisplay::new(&slots));
let mut last_statuses = vec![TaskStatus::Pending; slots.len()];
loop {
if let Some(display) = &mut terminal {
display.refresh(&slots);
}
match receiver.recv_timeout(REFRESH_INTERVAL) {
Ok(RenderEvent::TaskMessage { identity, message }) => {
write_message(output, terminal.as_ref(), &format!("{identity}: {message}"));
}
Ok(RenderEvent::Stop { success, message }) => {
if let Some(display) = &mut terminal {
display.refresh(&slots);
display.finish(&slots);
}
if output == OutputMode::Plain {
write_plain_transitions(&slots, &mut last_statuses);
}
drop(terminal.take());
drop(terminal_session.take());
write_final(output, &slots, success, &message);
break;
}
Err(RecvTimeoutError::Timeout) => {
if output == OutputMode::Plain {
write_plain_transitions(&slots, &mut last_statuses);
}
}
Err(RecvTimeoutError::Disconnected) => break,
}
}
}
struct TerminalDisplay {
multi: MultiProgress,
headings: Vec<ProgressBar>,
bars: Vec<ProgressBar>,
started: Vec<bool>,
idle_style: ProgressStyle,
activity_style: ProgressStyle,
known_style: ProgressStyle,
unknown_style: ProgressStyle,
}
impl TerminalDisplay {
fn new(slots: &[Arc<ProgressSlot>]) -> Self {
let multi = MultiProgress::with_draw_target(ProgressDrawTarget::stderr());
let idle_style = ProgressStyle::with_template("{prefix:.bold} [{msg}]")
.expect("hard-coded idle progress template is valid");
let activity_style =
ProgressStyle::with_template("{prefix:.bold} [{msg}] elapsed {elapsed_precise}")
.expect("hard-coded activity template is valid");
let known_style = ProgressStyle::with_template(
"{prefix:.bold} [{msg}] {wide_bar:.cyan/blue} {pos}/{len} elapsed {elapsed_precise} ETA {eta_precise}",
)
.expect("hard-coded progress template is valid");
let unknown_style = ProgressStyle::with_template(
"{prefix:.bold} [{msg}] {spinner:.cyan} iteration {pos} elapsed {elapsed_precise} ETA unknown",
)
.expect("hard-coded spinner template is valid");
let heading_style =
ProgressStyle::with_template("{prefix:.bold} [{msg}] elapsed {elapsed_precise}")
.expect("hard-coded phase-heading template is valid");
let mut headings = Vec::new();
let mut bars = Vec::with_capacity(slots.len());
let mut previous_phase: Option<&str> = None;
for slot in slots {
if let Some(phase) = slot.phase_label.as_deref()
&& previous_phase != Some(phase)
{
let heading = multi.add(ProgressBar::new(0));
heading.set_style(heading_style.clone());
heading.set_prefix(format!("── {phase} ──"));
headings.push(heading);
previous_phase = Some(phase);
}
let bar = multi.add(ProgressBar::new_spinner());
bar.set_prefix(slot.identity.label().to_owned());
bar.set_style(idle_style.clone());
bar.set_message(terminal_status(TaskStatus::Pending));
bars.push(bar);
}
for bar in &bars {
bar.force_draw();
}
for heading in &headings {
heading.force_draw();
}
Self {
multi,
headings,
bars,
started: vec![false; slots.len()],
idle_style,
activity_style,
known_style,
unknown_style,
}
}
fn refresh(&mut self, slots: &[Arc<ProgressSlot>]) {
let summary = summarize(slots);
let status = if summary.failed > 0 {
"failed"
} else if summary.running > 0 {
"running"
} else if summary.pending > 0 {
"pending"
} else {
"completed"
};
for heading in &self.headings {
heading.set_message(format!(
"{status} · running={} pending={} completed={} reused={} failed={}",
summary.running, summary.pending, summary.completed, summary.reused, summary.failed,
));
}
for ((bar, started), slot) in self.bars.iter().zip(&mut self.started).zip(slots) {
let status = TaskStatus::decode(slot.status.load(Ordering::Acquire));
if !*started && slot.started.load(Ordering::Acquire) {
bar.reset_elapsed();
*started = true;
}
if slot.display_kind == TaskDisplayKind::Activity {
bar.set_style(if *started {
self.activity_style.clone()
} else {
self.idle_style.clone()
});
} else if *started {
if !slot.target_known.load(Ordering::Acquire) {
bar.set_style(self.unknown_style.clone());
} else {
let target = slot.target.load(Ordering::Relaxed);
bar.set_style(self.known_style.clone());
bar.set_length(target);
}
bar.set_position(slot.current.load(Ordering::Relaxed));
} else {
bar.set_style(self.idle_style.clone());
}
let detail = lock(&slot.detail);
if detail.is_empty() || detail.as_ref() == status.label() {
bar.set_message(terminal_status(status));
} else {
bar.set_message(format!("{}: {}", terminal_status(status), detail.as_ref()));
}
if *started && status == TaskStatus::Running {
bar.tick();
}
}
}
fn finish(&self, slots: &[Arc<ProgressSlot>]) {
for (bar, slot) in self.bars.iter().zip(slots) {
let status = TaskStatus::decode(slot.status.load(Ordering::Acquire));
bar.finish_with_message(terminal_status(status));
}
for heading in &self.headings {
heading.finish();
}
let _ = self.multi.clear();
}
}
fn write_message(output: OutputMode, terminal: Option<&TerminalDisplay>, message: &str) {
match output {
OutputMode::Terminal => {
if let Some(display) = terminal {
let _ = display.multi.println(message);
}
}
OutputMode::Plain => eprintln!("[progress] {message}"),
OutputMode::Hidden | OutputMode::Auto => {}
}
}
fn write_plain_transitions(slots: &[Arc<ProgressSlot>], previous: &mut [TaskStatus]) {
for (slot, old) in slots.iter().zip(previous) {
let status = TaskStatus::decode(slot.status.load(Ordering::Acquire));
if status != *old {
let detail = lock(&slot.detail);
eprintln!(
"[task] identity={} status={} detail={} iteration={} target={}",
slot.identity.task_key(),
status.label(),
detail.as_ref(),
slot.current.load(Ordering::Relaxed),
format_target(slot)
);
*old = status;
}
}
}
fn write_final(output: OutputMode, slots: &[Arc<ProgressSlot>], success: bool, message: &str) {
if output == OutputMode::Hidden || (output == OutputMode::Terminal && success) {
return;
}
let summary = summarize(slots);
if !success {
for slot in slots {
let detail = lock(&slot.detail);
eprintln!(
"[task-final] task={} status={} detail={}",
slot.identity.task_key(),
TaskStatus::decode(slot.status.load(Ordering::Acquire)).label(),
detail.as_ref(),
);
}
}
eprintln!(
"[workflow] status={} tasks={} completed={} reused={} failed={} cancelled={} skipped={} pending={} message={}",
if success { "completed" } else { "failed" },
summary.total,
summary.completed,
summary.reused,
summary.failed,
summary.cancelled,
summary.skipped,
summary.pending,
message
);
}
fn format_target(slot: &ProgressSlot) -> String {
if slot.target_known.load(Ordering::Acquire) {
slot.target.load(Ordering::Relaxed).to_string()
} else {
"unknown".to_owned()
}
}
fn terminal_status(status: TaskStatus) -> String {
let style = match status {
TaskStatus::Pending => console::Style::new().dim(),
TaskStatus::Running => console::Style::new().cyan(),
TaskStatus::Completed => console::Style::new().green(),
TaskStatus::Reused => console::Style::new().blue(),
TaskStatus::Failed => console::Style::new().red(),
TaskStatus::Cancelled => console::Style::new().yellow(),
TaskStatus::Skipped => console::Style::new().dim(),
}
.force_styling(true);
style.apply_to(status.label()).to_string()
}
fn lock<T>(mutex: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
mutex
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}