bot-forge 1.0.2

Rust CLI for installing agent skills and developer tools from configurable forms.
Documentation
//! Process-wide coordination for concurrent command progress rendering.
//!
//! Each command publishes snapshots independently. A single renderer coalesces pending updates,
//! keeps rows in stable command order, and temporarily clears them while other output owns stderr.

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,
}

/// Background activity indicator for launcher actions before command progress is available.
pub(crate) struct ActivityStatus {
    stop: Option<mpsc::Sender<()>>,
    worker: Option<thread::JoinHandle<()>>,
}

impl ActivityStatus {
    /// Start a spinner when the current terminal can render coordinated progress.
    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 {
    /// Register a command with the process-wide progress renderer when status display is enabled.
    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,
        }
    }

    /// Publish the latest elapsed time and sanitized output summary.
    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;
    }

    /// Remove this command's row and wait until the renderer has applied the removal.
    ///
    /// Waiting for acknowledgement prevents subsequent stderr output from racing with a stale
    /// progress row. Calling `finish` more than once is harmless.
    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)]
/// Latest renderable state for one command, keyed and ordered by its monotonic ID.
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>>,
}

/// Clear coordinated progress rows until the returned guard is dropped.
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,
        );

        // Commands report independently. Coalesce queued snapshots so one refresh represents the
        // complete running set instead of rendering one row at a time.
        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);
}

/// Apply one coordinator message while preserving nested pause depth and finish acknowledgements.
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),
    }
}

/// Render one aggregate row followed by command rows in stable ID order.
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
}