use std::io::{IsTerminal, Write};
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::mpsc::{Receiver, Sender, channel};
use std::sync::{Arc, Mutex, PoisonError, TryLockError};
#[cfg(not(target_arch = "wasm32"))]
use std::time::Instant;
#[cfg(target_arch = "wasm32")]
use web_time::Instant;
pub const NO_PROGRESS_ENV: &str = "FOCR_NO_PROGRESS";
const BAR_CELLS: usize = 24;
static ACTIVE_LINE_WIDTH: AtomicUsize = AtomicUsize::new(0);
static ACTIVE_GENERATION: AtomicUsize = AtomicUsize::new(0);
static NEXT_GENERATION: AtomicUsize = AtomicUsize::new(1);
static PROGRESS_SUPPRESSED: AtomicBool = AtomicBool::new(false);
static OUTPUT_LOCK: Mutex<()> = Mutex::new(());
struct Inner {
label: &'static str,
total: usize,
done: usize,
current: String,
started: Instant,
}
struct RenderFrame {
label: &'static str,
total: usize,
done: usize,
current: String,
elapsed_secs: f64,
}
impl From<&Inner> for RenderFrame {
fn from(inner: &Inner) -> Self {
Self {
label: inner.label,
total: inner.total,
done: inner.done,
current: inner.current.clone(),
elapsed_secs: inner.started.elapsed().as_secs_f64(),
}
}
}
struct Shared {
state: Mutex<Inner>,
finished: AtomicBool,
suspended: AtomicBool,
generation: usize,
}
enum RenderCommand {
Draw,
Finish(Option<Sender<()>>),
}
pub struct Progress {
inner: Option<Arc<Shared>>,
renderer: Option<Sender<RenderCommand>>,
}
impl Progress {
pub fn new(label: &'static str, total: usize, human_mode: bool) -> Self {
if !human_mode
|| total == 0
|| PROGRESS_SUPPRESSED.load(Ordering::Acquire)
|| !stderr_supports_progress()
{
return Self {
inner: None,
renderer: None,
};
}
let inner = Arc::new(Shared {
state: Mutex::new(Inner {
label,
total,
done: 0,
current: String::new(),
started: Instant::now(),
}),
finished: AtomicBool::new(false),
suspended: AtomicBool::new(false),
generation: NEXT_GENERATION.fetch_add(1, Ordering::Relaxed),
});
let (renderer, commands) = channel();
let renderer_state = Arc::clone(&inner);
if std::thread::Builder::new()
.name("focr-progress".into())
.spawn(move || render_commands(renderer_state, commands))
.is_err()
{
return Self {
inner: None,
renderer: None,
};
}
Self {
inner: Some(inner),
renderer: Some(renderer),
}
}
pub fn is_enabled(&self) -> bool {
self.inner.is_some() && self.renderer.is_some()
}
pub fn start_item(&self, current: impl Into<String>) {
let Some(inner) = &self.inner else { return };
if inner.finished.load(Ordering::Acquire) {
return;
}
inner.suspended.store(false, Ordering::Release);
let mut g = lock(&inner.state);
if inner.finished.load(Ordering::Acquire) {
return;
}
g.current = current.into();
drop(g);
self.send(RenderCommand::Draw);
}
pub fn complete_item(&self) {
let Some(inner) = &self.inner else { return };
if inner.finished.load(Ordering::Acquire) {
return;
}
let mut g = lock(&inner.state);
if inner.finished.load(Ordering::Acquire) {
return;
}
g.done = (g.done + 1).min(g.total);
}
pub fn suspend(&self) {
let Some(inner) = &self.inner else { return };
if inner.finished.load(Ordering::Acquire) {
return;
}
inner.suspended.store(true, Ordering::Release);
clear_bar_line(inner.generation);
}
pub fn note(&self, msg: &str) {
let Some(inner) = &self.inner else {
stderr_message(format_args!("{msg}"));
return;
};
if inner.finished.load(Ordering::Acquire) || PROGRESS_SUPPRESSED.load(Ordering::Acquire) {
return;
}
inner.suspended.store(true, Ordering::Release);
{
let _output = output_lock();
clear_generation_locked(inner.generation);
let mut err = std::io::stderr().lock();
let _ = writeln!(err, "{msg}");
let _ = err.flush();
}
inner.suspended.store(false, Ordering::Release);
self.send(RenderCommand::Draw);
}
pub fn finish(&self) {
let Some(inner) = &self.inner else { return };
if inner.finished.swap(true, Ordering::AcqRel) {
return;
}
inner.suspended.store(true, Ordering::Release);
let (ack_tx, ack_rx) = channel();
self.send(RenderCommand::Finish(Some(ack_tx)));
if ack_rx.recv().is_err() {
clear_bar_line(inner.generation);
}
}
pub fn retire(&self) {
let Some(inner) = &self.inner else { return };
if inner.finished.swap(true, Ordering::AcqRel) {
return;
}
inner.suspended.store(true, Ordering::Release);
self.send(RenderCommand::Finish(None));
try_clear_bar_line(inner.generation);
}
pub fn page_sink(&self) -> crate::PageSink {
let Some(inner) = &self.inner else {
return Box::new(|_page, _body| {});
};
let Some(renderer) = self.renderer.as_ref().cloned() else {
return Box::new(|_page, _body| {});
};
let inner = Arc::clone(inner);
Box::new(move |page, _body| {
if inner.finished.load(Ordering::Acquire) || PROGRESS_SUPPRESSED.load(Ordering::Acquire)
{
return;
}
let mut g = lock(&inner.state);
if inner.finished.load(Ordering::Acquire) {
return;
}
g.done = page.min(g.total);
g.current = format!("page {page}/{}", g.total);
drop(g);
let _ = renderer.send(RenderCommand::Draw);
})
}
fn send(&self, command: RenderCommand) {
if let Some(renderer) = &self.renderer {
let _ = renderer.send(command);
}
}
}
impl Drop for Progress {
fn drop(&mut self) {
self.retire();
}
}
fn lock(inner: &Mutex<Inner>) -> std::sync::MutexGuard<'_, Inner> {
inner.lock().unwrap_or_else(PoisonError::into_inner)
}
fn output_lock() -> std::sync::MutexGuard<'static, ()> {
OUTPUT_LOCK.lock().unwrap_or_else(PoisonError::into_inner)
}
fn try_output_lock() -> Option<std::sync::MutexGuard<'static, ()>> {
match OUTPUT_LOCK.try_lock() {
Ok(output) => Some(output),
Err(TryLockError::Poisoned(e)) => Some(e.into_inner()),
Err(TryLockError::WouldBlock) => None,
}
}
fn render_commands(inner: Arc<Shared>, commands: Receiver<RenderCommand>) {
while let Ok(command) = commands.recv() {
match command {
RenderCommand::Draw => draw_latest(&inner),
RenderCommand::Finish(ack) => {
clear_bar_line(inner.generation);
if let Some(ack) = ack {
let _ = ack.send(());
}
return;
}
}
}
clear_bar_line(inner.generation);
}
fn draw_latest(inner: &Shared) {
let frame = {
let g = lock(&inner.state);
RenderFrame::from(&*g)
};
let _ = draw_frame(inner, &frame);
}
fn draw_frame(inner: &Shared, frame: &RenderFrame) -> usize {
let _output = output_lock();
if inner.finished.load(Ordering::Acquire)
|| inner.suspended.load(Ordering::Acquire)
|| PROGRESS_SUPPRESSED.load(Ordering::Acquire)
{
clear_generation_locked(inner.generation);
return 0;
}
let line = render_line(
frame.label,
frame.done,
frame.total,
&frame.current,
frame.elapsed_secs,
);
let width = line.chars().count();
let previous_width = ACTIVE_LINE_WIDTH.load(Ordering::Acquire);
let pad = previous_width.saturating_sub(width);
ACTIVE_GENERATION.store(inner.generation, Ordering::Release);
ACTIVE_LINE_WIDTH.store(width.max(previous_width), Ordering::Release);
let mut err = std::io::stderr().lock();
let _ = write!(err, "\r{line}{:pad$}", "");
let _ = err.flush();
drop(err);
if inner.finished.load(Ordering::Acquire)
|| inner.suspended.load(Ordering::Acquire)
|| PROGRESS_SUPPRESSED.load(Ordering::Acquire)
{
clear_generation_locked(inner.generation);
0
} else {
width
}
}
fn clear_generation_locked(generation: usize) {
if ACTIVE_GENERATION.load(Ordering::Acquire) != generation {
return;
}
let width = ACTIVE_LINE_WIDTH.swap(0, Ordering::AcqRel);
ACTIVE_GENERATION.store(0, Ordering::Release);
if width == 0 {
return;
}
let mut err = std::io::stderr().lock();
let _ = write!(err, "\r{:width$}\r", "");
let _ = err.flush();
ACTIVE_LINE_WIDTH.store(0, Ordering::Release);
}
fn render_line(label: &str, done: usize, total: usize, current: &str, elapsed_secs: f64) -> String {
let total = total.max(1);
let done = done.min(total);
let filled = done * BAR_CELLS / total;
let pct = done * 100 / total;
let mut line = format!(
"[focr] {label} [{}{}] {pct:>3}%",
"#".repeat(filled),
"-".repeat(BAR_CELLS - filled)
);
if !current.is_empty() {
line.push(' ');
line.extend(current.chars().map(|ch| {
if ch.is_ascii() && !ch.is_ascii_control() {
ch
} else {
'?'
}
}));
}
line.push_str(&format!(", {} elapsed", fmt_duration(elapsed_secs)));
if done > 0 && done < total {
let eta = elapsed_secs / done as f64 * (total - done) as f64;
line.push_str(&format!(", ~{} left", fmt_duration(eta)));
}
line
}
pub fn clear_active_line() {
let _output = output_lock();
clear_active_line_locked();
}
pub fn stderr_message(args: std::fmt::Arguments<'_>) {
let _output = output_lock();
clear_active_line_locked();
let mut err = std::io::stderr().lock();
let _ = writeln!(err, "{args}");
let _ = err.flush();
}
pub fn try_stderr_message(args: std::fmt::Arguments<'_>) -> bool {
let Some(_output) = try_output_lock() else {
return false;
};
clear_active_line_locked();
let mut err = std::io::stderr().lock();
let _ = writeln!(err, "{args}");
let _ = err.flush();
true
}
fn clear_bar_line(generation: usize) {
let _output = output_lock();
clear_generation_locked(generation);
}
pub fn suppress_for_interrupt() {
PROGRESS_SUPPRESSED.store(true, Ordering::Release);
try_clear_active_line();
}
fn try_clear_bar_line(generation: usize) {
if let Some(_output) = try_output_lock() {
clear_generation_locked(generation);
}
}
fn try_clear_active_line() {
if let Some(_output) = try_output_lock() {
clear_active_line_locked();
}
}
fn clear_active_line_locked() {
let width = ACTIVE_LINE_WIDTH.swap(0, Ordering::AcqRel);
ACTIVE_GENERATION.store(0, Ordering::Release);
if width == 0 {
return;
}
let mut err = std::io::stderr().lock();
let _ = write!(err, "\r{:width$}\r", "");
let _ = err.flush();
}
fn fmt_duration(secs: f64) -> String {
let s = secs.max(0.0).round() as u64;
if s < 60 {
format!("{s}s")
} else if s < 3600 {
format!("{}m{:02}s", s / 60, s % 60)
} else {
format!("{}h{:02}m", s / 3600, (s % 3600) / 60)
}
}
fn stderr_supports_progress() -> bool {
if std::env::var_os("FOCR_TIMING").is_some() {
return false;
}
progress_allowed(
std::io::stderr().is_terminal(),
std::env::var("TERM").ok().as_deref(),
std::env::var(NO_PROGRESS_ENV).ok().as_deref(),
)
}
pub fn progress_allowed(
stderr_is_tty: bool,
term: Option<&str>,
no_progress: Option<&str>,
) -> bool {
if !stderr_is_tty {
return false;
}
if term == Some("dumb") {
return false;
}
match no_progress {
None => true,
Some(v) => {
let v = v.trim();
v.is_empty() || v == "0"
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::mpsc;
use std::time::Duration;
static TEST_GLOBAL_LOCK: Mutex<()> = Mutex::new(());
fn test_bar(total: usize) -> (Progress, Arc<Shared>, std::thread::JoinHandle<()>) {
let shared = Arc::new(Shared {
state: Mutex::new(Inner {
label: "ocr",
total,
done: 0,
current: String::new(),
started: Instant::now(),
}),
finished: AtomicBool::new(false),
suspended: AtomicBool::new(false),
generation: NEXT_GENERATION.fetch_add(1, Ordering::Relaxed),
});
let (renderer, commands) = channel();
let renderer_state = Arc::clone(&shared);
let renderer_thread = std::thread::spawn(move || render_commands(renderer_state, commands));
(
Progress {
inner: Some(Arc::clone(&shared)),
renderer: Some(renderer),
},
shared,
renderer_thread,
)
}
#[test]
fn not_a_tty_never_draws() {
assert!(!progress_allowed(false, None, None));
assert!(!progress_allowed(false, Some("xterm-256color"), None));
assert!(!progress_allowed(false, Some("xterm"), Some("0")));
}
#[test]
fn tty_gate_respects_term_and_kill_switch() {
assert!(progress_allowed(true, None, None));
assert!(progress_allowed(true, Some("xterm-256color"), None));
assert!(!progress_allowed(true, Some("dumb"), None));
assert!(!progress_allowed(true, Some("xterm"), Some("1")));
assert!(!progress_allowed(true, Some("xterm"), Some("true")));
assert!(progress_allowed(true, Some("xterm"), Some("0")));
assert!(progress_allowed(true, Some("xterm"), Some("")));
}
#[test]
fn robot_and_json_modes_construct_disabled() {
let bar = Progress::new("ocr", 10, false);
assert!(!bar.is_enabled());
bar.start_item("page 1/10");
bar.complete_item();
bar.suspend();
bar.retire();
let mut sink = bar.page_sink();
sink(3, "body");
}
#[test]
fn zero_total_constructs_disabled() {
assert!(!Progress::new("ocr", 0, true).is_enabled());
}
#[test]
fn render_line_shape() {
let line = render_line("ocr", 4, 12, "page 5/12", 42.0);
assert!(line.starts_with("[focr] ocr ["));
assert!(line.contains(" 33%"), "{line}");
assert!(line.contains("page 5/12"), "{line}");
assert!(line.contains("42s elapsed"), "{line}");
assert!(line.contains("~1m24s left"), "{line}");
assert!(!line.contains('\r') && !line.contains('\n'));
assert!(line.is_ascii());
}
#[test]
fn render_line_sanitizes_terminal_control_and_wide_characters() {
let line = render_line(
"batch",
0,
1,
"image 1/1: bad\r\n\u{1b}[31m-wide-\u{754c}.png",
0.0,
);
assert!(line.is_ascii(), "{line:?}");
assert!(!line.contains('\r') && !line.contains('\n'));
assert!(!line.contains('\u{1b}'));
assert!(line.contains("bad???[31m-wide-?.png"), "{line:?}");
}
#[test]
fn finish_does_not_wait_for_page_sink_state_lock() {
let _serial = TEST_GLOBAL_LOCK.lock().unwrap();
let (bar, shared, renderer) = test_bar(2);
let guard = lock(&shared.state);
let (tx, rx) = mpsc::channel();
let worker = std::thread::spawn(move || {
bar.finish();
let _ = tx.send(());
});
assert!(
rx.recv_timeout(Duration::from_secs(1)).is_ok(),
"finish blocked behind a busy page-sink mutex"
);
drop(guard);
worker.join().expect("finish worker");
renderer.join().expect("progress renderer");
assert!(shared.finished.load(Ordering::Acquire));
}
#[test]
fn renderer_cannot_publish_after_retirement() {
let _serial = TEST_GLOBAL_LOCK.lock().unwrap();
ACTIVE_LINE_WIDTH.store(0, Ordering::Release);
ACTIVE_GENERATION.store(0, Ordering::Release);
let (bar, shared, progress_renderer) = test_bar(2);
let frame = {
let mut g = lock(&shared.state);
g.current = "page 1/2".into();
RenderFrame::from(&*g)
};
let output_guard = output_lock();
let (started_tx, started_rx) = mpsc::channel();
let (done_tx, done_rx) = mpsc::channel();
let renderer_state = Arc::clone(&shared);
let renderer = std::thread::spawn(move || {
let _ = started_tx.send(());
let rendered_width = draw_frame(&renderer_state, &frame);
let _ = done_tx.send(rendered_width);
});
started_rx.recv().expect("renderer starts");
bar.retire();
assert!(
done_rx.try_recv().is_err(),
"renderer must be at output gate"
);
drop(output_guard);
assert_eq!(
done_rx
.recv_timeout(Duration::from_secs(1))
.expect("renderer retires"),
0
);
renderer.join().expect("renderer worker");
progress_renderer.join().expect("progress renderer");
assert_eq!(ACTIVE_LINE_WIDTH.load(Ordering::Acquire), 0);
}
#[test]
fn stale_renderer_cannot_clear_a_newer_generation() {
let _serial = TEST_GLOBAL_LOCK.lock().unwrap();
let old_generation = NEXT_GENERATION.fetch_add(1, Ordering::Relaxed);
let newer_generation = NEXT_GENERATION.fetch_add(1, Ordering::Relaxed);
ACTIVE_GENERATION.store(newer_generation, Ordering::Release);
ACTIVE_LINE_WIDTH.store(37, Ordering::Release);
let shared = Arc::new(Shared {
state: Mutex::new(Inner {
label: "old",
total: 1,
done: 0,
current: String::new(),
started: Instant::now(),
}),
finished: AtomicBool::new(true),
suspended: AtomicBool::new(true),
generation: old_generation,
});
let (commands, receiver) = channel();
let renderer = std::thread::spawn(move || render_commands(shared, receiver));
let (ack_tx, ack_rx) = channel();
commands
.send(RenderCommand::Finish(Some(ack_tx)))
.expect("finish stale renderer");
ack_rx
.recv_timeout(Duration::from_secs(1))
.expect("stale renderer acknowledges retirement");
renderer.join().expect("stale renderer");
assert_eq!(ACTIVE_GENERATION.load(Ordering::Acquire), newer_generation);
assert_eq!(ACTIVE_LINE_WIDTH.load(Ordering::Acquire), 37);
ACTIVE_GENERATION.store(0, Ordering::Release);
ACTIVE_LINE_WIDTH.store(0, Ordering::Release);
}
#[test]
fn teardown_stderr_never_waits_for_terminal_output() {
let _serial = TEST_GLOBAL_LOCK.lock().unwrap();
let output_guard = output_lock();
let (tx, rx) = mpsc::channel();
let worker = std::thread::spawn(move || {
let wrote = try_stderr_message(format_args!("teardown"));
let _ = tx.send(wrote);
});
assert!(
!rx.recv_timeout(Duration::from_secs(1))
.expect("teardown stderr returns"),
"teardown stderr unexpectedly acquired the busy output lock"
);
worker.join().expect("teardown stderr worker");
drop(output_guard);
}
#[test]
fn page_sink_never_waits_for_terminal_output() {
let _serial = TEST_GLOBAL_LOCK.lock().unwrap();
let (bar, _shared, renderer) = test_bar(2);
let mut sink = bar.page_sink();
let output_guard = output_lock();
let (tx, rx) = mpsc::channel();
let callback = std::thread::spawn(move || {
sink(1, "body");
let _ = tx.send(());
});
assert!(
rx.recv_timeout(Duration::from_secs(1)).is_ok(),
"model-runtime callback waited for terminal I/O"
);
callback.join().expect("page callback");
drop(output_guard);
bar.finish();
renderer.join().expect("progress renderer");
}
#[test]
fn render_line_eta_bounds() {
assert!(!render_line("ocr", 0, 5, "page 1/5", 3.0).contains("left"));
let done = render_line("ocr", 5, 5, "", 30.0);
assert!(!done.contains("left"));
assert!(done.contains("100%"));
}
#[test]
fn duration_formatting() {
assert_eq!(fmt_duration(3.2), "3s");
assert_eq!(fmt_duration(62.4), "1m02s");
assert_eq!(fmt_duration(3725.0), "1h02m");
assert_eq!(fmt_duration(-1.0), "0s");
}
}