use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use trusty_progress::relay::{relay_requested, StageEvent, StageState, ENV_RELAY};
use super::{Outcome, ProgressBus, ProgressEvent};
const DRAIN_INTERVAL: std::time::Duration = std::time::Duration::from_millis(100);
pub fn line_for(event: &ProgressEvent) -> String {
let state = match &event.outcome {
None if event.done == 0 => StageState::Started,
None => StageState::Advanced,
Some(Outcome::Completed) => StageState::Completed,
Some(Outcome::Failed { .. }) => StageState::Failed,
Some(Outcome::Skipped { .. }) => StageState::Skipped,
};
let mut relayed = StageEvent::new(event.stage.label(), &event.target, state)
.with_counts(event.done, event.total);
if let Some(text) = event
.outcome
.as_ref()
.and_then(Outcome::reason)
.or(event.detail.as_deref())
{
relayed = relayed.with_detail(text);
}
relayed.encode()
}
#[derive(Debug)]
pub struct StageRelay {
bus: Option<ProgressBus>,
stop: Arc<AtomicBool>,
task: Option<tokio::task::JoinHandle<()>>,
}
impl StageRelay {
pub fn from_env() -> Self {
if relay_requested(std::env::var(ENV_RELAY).ok().as_deref()) {
Self::started()
} else {
Self::off()
}
}
pub fn off() -> Self {
Self {
bus: None,
stop: Arc::new(AtomicBool::new(true)),
task: None,
}
}
pub fn started() -> Self {
let bus = ProgressBus::new();
let stop = Arc::new(AtomicBool::new(false));
let task = tokio::spawn(drain_to_stderr(bus.clone(), Arc::clone(&stop)));
Self {
bus: Some(bus),
stop,
task: Some(task),
}
}
pub fn bus(&self) -> Option<&ProgressBus> {
self.bus.as_ref()
}
pub async fn finish(mut self) {
self.stop.store(true, Ordering::Relaxed);
if let Some(task) = self.task.take() {
let _ = task.await;
}
}
}
impl Drop for StageRelay {
fn drop(&mut self) {
self.stop.store(true, Ordering::Relaxed);
}
}
async fn drain_to_stderr(bus: ProgressBus, stop: Arc<AtomicBool>) {
loop {
let stopping = stop.load(Ordering::Relaxed);
write_lines(&bus.drain());
if stopping {
return;
}
tokio::time::sleep(DRAIN_INTERVAL).await;
}
}
fn write_lines(events: &[ProgressEvent]) {
use std::io::Write as _;
if events.is_empty() {
return;
}
let mut rendered = String::new();
for event in events {
rendered.push_str(&line_for(event));
rendered.push('\n');
}
let mut err = std::io::stderr().lock();
let _ = err.write_all(rendered.as_bytes());
let _ = err.flush();
}