use std::collections::BTreeMap;
use std::io::{self, IsTerminal};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{OnceLock, mpsc};
use std::thread;
use std::time::{Duration, Instant};
use crate::ui::{RenderOptions, command_progress_line, muted_text, replace_progress_stderr};
const PROGRESS_COALESCE_INTERVAL: Duration = Duration::from_millis(20);
const PROGRESS_FRAMES: [&str; 10] = ["⠋", "⠙", "⠹", "⠸", "⠼", "⠴", "⠦", "⠧", "⠇", "⠏"];
static PROGRESS_SEQUENCE: AtomicU64 = AtomicU64::new(0);
static PROGRESS_COORDINATOR: OnceLock<mpsc::Sender<ProgressMessage>> = OnceLock::new();
pub(crate) struct CommandStatus {
id: u64,
label: String,
step: Option<(usize, usize)>,
active: bool,
frame: usize,
}
pub(crate) struct ActivityStatus {
stop: Option<mpsc::Sender<()>>,
worker: Option<thread::JoinHandle<()>>,
}
impl ActivityStatus {
pub(crate) fn new(label: &str, detail: &str) -> Self {
if !(io::stderr().is_terminal() && RenderOptions::stderr().cursor) {
return Self {
stop: None,
worker: None,
};
}
let (stop, receiver) = mpsc::channel();
let label = label.to_string();
let detail = detail.to_string();
let worker = thread::Builder::new()
.name("bot-forge-activity".into())
.spawn(move || {
let mut status = CommandStatus::new(&label, None, true);
let started = Instant::now();
loop {
if receiver.recv_timeout(Duration::from_millis(120)).is_ok() {
break;
}
status.update(started.elapsed().as_secs(), detail.clone());
}
status.finish();
})
.ok();
Self {
stop: worker.as_ref().map(|_| stop),
worker,
}
}
}
impl Drop for ActivityStatus {
fn drop(&mut self) {
if let Some(stop) = self.stop.take() {
let _ = stop.send(());
}
if let Some(worker) = self.worker.take() {
let _ = worker.join();
}
}
}
impl CommandStatus {
pub(crate) fn new(label: &str, step: Option<(usize, usize)>, enabled: bool) -> Self {
let active = enabled && io::stderr().is_terminal() && RenderOptions::stderr().cursor;
Self {
id: PROGRESS_SEQUENCE.fetch_add(1, Ordering::Relaxed),
label: label.to_string(),
step,
active,
frame: 0,
}
}
pub(crate) fn update(&mut self, elapsed: u64, recent: String) {
if !self.active {
return;
}
progress_sender()
.send(ProgressMessage::Update(ProgressSnapshot {
id: self.id,
label: self.label.clone(),
step: self.step,
elapsed,
recent,
frame: self.frame,
}))
.ok();
self.frame += 1;
}
pub(crate) fn finish(&mut self) {
if self.active {
let (sender, receiver) = mpsc::sync_channel(0);
if progress_sender()
.send(ProgressMessage::Finish(self.id, sender))
.is_ok()
{
let _ = receiver.recv();
}
}
self.active = false;
}
}
impl Drop for CommandStatus {
fn drop(&mut self) {
self.finish();
}
}
#[derive(Clone)]
pub(crate) struct ProgressSnapshot {
pub(crate) id: u64,
pub(crate) label: String,
pub(crate) step: Option<(usize, usize)>,
pub(crate) elapsed: u64,
pub(crate) recent: String,
pub(crate) frame: usize,
}
pub(crate) enum ProgressMessage {
Update(ProgressSnapshot),
Finish(u64, mpsc::SyncSender<()>),
Pause(mpsc::SyncSender<()>),
Resume,
}
pub(crate) struct ProgressPauseGuard {
sender: Option<&'static mpsc::Sender<ProgressMessage>>,
}
pub(crate) fn pause() -> ProgressPauseGuard {
let Some(sender) = PROGRESS_COORDINATOR.get() else {
return ProgressPauseGuard { sender: None };
};
let (acknowledgement, receiver) = mpsc::sync_channel(0);
if sender.send(ProgressMessage::Pause(acknowledgement)).is_ok() {
let _ = receiver.recv();
ProgressPauseGuard {
sender: Some(sender),
}
} else {
ProgressPauseGuard { sender: None }
}
}
impl Drop for ProgressPauseGuard {
fn drop(&mut self) {
if let Some(sender) = self.sender {
let _ = sender.send(ProgressMessage::Resume);
}
}
}
fn progress_sender() -> &'static mpsc::Sender<ProgressMessage> {
PROGRESS_COORDINATOR.get_or_init(|| {
let (sender, receiver) = mpsc::channel();
thread::Builder::new()
.name("bot-forge-progress".into())
.spawn(move || progress_render_loop(receiver))
.expect("progress renderer thread must start");
sender
})
}
fn progress_render_loop(receiver: mpsc::Receiver<ProgressMessage>) {
let mut active = BTreeMap::new();
let mut rendered_lines = 0;
let mut pause_depth = 0_usize;
while let Ok(message) = receiver.recv() {
let mut acknowledgements = Vec::new();
apply_progress_message(
message,
&mut active,
&mut acknowledgements,
&mut pause_depth,
);
loop {
match receiver.recv_timeout(PROGRESS_COALESCE_INTERVAL) {
Ok(message) => apply_progress_message(
message,
&mut active,
&mut acknowledgements,
&mut pause_depth,
),
Err(mpsc::RecvTimeoutError::Timeout) => break,
Err(mpsc::RecvTimeoutError::Disconnected) => {
replace_progress_stderr(&[], rendered_lines);
return;
}
}
}
let lines = if pause_depth > 0 || active.is_empty() {
Vec::new()
} else {
progress_summary_lines(&active)
};
rendered_lines = replace_progress_stderr(&lines, rendered_lines);
for acknowledgement in acknowledgements {
let _ = acknowledgement.send(());
}
}
replace_progress_stderr(&[], rendered_lines);
}
pub(crate) fn apply_progress_message(
message: ProgressMessage,
active: &mut BTreeMap<u64, ProgressSnapshot>,
acknowledgements: &mut Vec<mpsc::SyncSender<()>>,
pause_depth: &mut usize,
) {
match message {
ProgressMessage::Update(snapshot) => {
active.insert(snapshot.id, snapshot);
}
ProgressMessage::Finish(id, acknowledgement) => {
active.remove(&id);
acknowledgements.push(acknowledgement);
}
ProgressMessage::Pause(acknowledgement) => {
*pause_depth = pause_depth.saturating_add(1);
acknowledgements.push(acknowledgement);
}
ProgressMessage::Resume => *pause_depth = pause_depth.saturating_sub(1),
}
}
pub(crate) fn progress_summary_lines(active: &BTreeMap<u64, ProgressSnapshot>) -> Vec<String> {
let frame = active
.values()
.map(|snapshot| snapshot.frame)
.max()
.unwrap_or(0);
let options = RenderOptions::stderr();
let mut lines = vec![command_progress_line(
PROGRESS_FRAMES[frame % PROGRESS_FRAMES.len()],
None,
&format!("{} running", active.len()),
&options,
)];
lines.extend(active.values().map(|snapshot| {
let progress = snapshot
.step
.map(|(current, total)| format!(" {current}/{total}"))
.unwrap_or_default();
muted_text(
&format!(
" {:<20} {:<7} {:>4}s {}",
snapshot.label,
progress.trim(),
snapshot.elapsed,
snapshot.recent
),
&options,
)
}));
lines
}