From be536557d200824503fe82ed3ea8a461263d490d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Josef=20=C5=A0im=C3=A1nek?= Date: Fri, 2 Oct 2026 22:17:15 +0200 Subject: [PATCH] Add standalone `rb-progress` reporting --- .github/workflows/ci.yml | 12 + Cargo.lock | 9 + crates/rb-progress/Cargo.toml | 23 + crates/rb-progress/README.md | 39 ++ crates/rb-progress/examples/failure.rs | 18 + crates/rb-progress/examples/parallel.rs | 34 ++ .../rb-progress/examples/progress_ui_probe.rs | 138 +++++ crates/rb-progress/examples/standalone.rs | 12 + crates/rb-progress/examples/tasks.rs | 48 ++ crates/rb-progress/src/events.rs | 11 + crates/rb-progress/src/lib.rs | 14 + crates/rb-progress/src/plain.rs | 46 ++ crates/rb-progress/src/render.rs | 363 +++++++++++++ crates/rb-progress/src/reporter.rs | 237 +++++++++ crates/rb-progress/src/reporter/tests.rs | 499 ++++++++++++++++++ crates/rb-progress/src/state.rs | 163 ++++++ crates/rb-progress/src/task.rs | 163 ++++++ crates/rb-progress/src/terminal.rs | 148 ++++++ crates/rb-progress/src/terminal/tests.rs | 167 ++++++ crates/rb-progress/tests/pty_ui_test.py | 298 +++++++++++ 20 files changed, 2442 insertions(+) create mode 100644 crates/rb-progress/Cargo.toml create mode 100644 crates/rb-progress/README.md create mode 100644 crates/rb-progress/examples/failure.rs create mode 100644 crates/rb-progress/examples/parallel.rs create mode 100644 crates/rb-progress/examples/progress_ui_probe.rs create mode 100644 crates/rb-progress/examples/standalone.rs create mode 100644 crates/rb-progress/examples/tasks.rs create mode 100644 crates/rb-progress/src/events.rs create mode 100644 crates/rb-progress/src/lib.rs create mode 100644 crates/rb-progress/src/plain.rs create mode 100644 crates/rb-progress/src/render.rs create mode 100644 crates/rb-progress/src/reporter.rs create mode 100644 crates/rb-progress/src/reporter/tests.rs create mode 100644 crates/rb-progress/src/state.rs create mode 100644 crates/rb-progress/src/task.rs create mode 100644 crates/rb-progress/src/terminal.rs create mode 100644 crates/rb-progress/src/terminal/tests.rs create mode 100644 crates/rb-progress/tests/pty_ui_test.py diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 4faa843..595858d 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -122,6 +122,18 @@ jobs: - name: Run integration tests run: cargo test --verbose --test '*' + - name: Test standalone progress reporting + run: | + cargo test -p rb-progress --no-default-features + cargo build -p rb-progress --no-default-features --examples + + - name: Test progress UI in a pseudo-terminal + if: matrix.os == 'ubuntu-latest' + run: | + cargo test -p rb-progress --features rb-task + cargo build -p rb-progress --features rb-task --example progress_ui_probe --example tasks + python3 crates/rb-progress/tests/pty_ui_test.py target/debug/examples/progress_ui_probe target/debug/examples/parallel + - name: Build release run: cargo build --release diff --git a/Cargo.lock b/Cargo.lock index fab1d45..3942645 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -582,6 +582,15 @@ dependencies = [ "which", ] +[[package]] +name = "rb-progress" +version = "0.3.1" +dependencies = [ + "rb-task", + "terminal_size", + "unicode-width 0.2.2", +] + [[package]] name = "rb-task" version = "0.3.1" diff --git a/crates/rb-progress/Cargo.toml b/crates/rb-progress/Cargo.toml new file mode 100644 index 0000000..720666c --- /dev/null +++ b/crates/rb-progress/Cargo.toml @@ -0,0 +1,23 @@ +[package] +name = "rb-progress" +version.workspace = true +edition.workspace = true +license.workspace = true +description = "Terminal progress reporting for task execution" + +[features] +default = [] +rb-task = ["dep:rb-task"] + +[dependencies] +terminal_size = "0.4.3" +unicode-width = "0.2.2" +rb-task = { path = "../rb-task", optional = true } + +[[example]] +name = "progress_ui_probe" +required-features = ["rb-task"] + +[[example]] +name = "tasks" +required-features = ["rb-task"] diff --git a/crates/rb-progress/README.md b/crates/rb-progress/README.md new file mode 100644 index 0000000..8566a98 --- /dev/null +++ b/crates/rb-progress/README.md @@ -0,0 +1,39 @@ +# rb-progress + +Progress reporting with a live terminal tree and plain, append-only output when +stderr is redirected. Supports parallel workers, elapsed times, and errors. + +```rust,no_run +use rb_progress::Reporter; + +let reporter = Reporter::start("Polishing the silver").unwrap(); +let progress = reporter.handle(); +progress.event("polishing", 0, 3, None); +// Do the work, then report completion. +progress.event("polishing", 3, 3, None); +reporter.finish_with_summary("The silver is ready, sir"); +``` + +Terminal output (timing varies): + +```text +┌─ ✓ Polishing the silver +│ ├─ ✓ polishing 3/3 (0.0s) +└─ ✓ The silver is ready, sir (0.0s) +``` + +Clone a `ProgressHandle` to report from multiple threads. Use `worker_started`, +`worker_event`, and `worker_finished` or `worker_failed` for each worker. +Use `fail` to report an overall error. + +`rb-task` is optional and disabled by default. Enable the `rb-task` feature to +connect `Reporter::task_sink()` to `TaskEvents`. + +Run the examples: + +```sh +cargo run -p rb-progress --example standalone +cargo run -p rb-progress --example parallel +cargo run -p rb-progress --example failure +cargo run -p rb-progress --features rb-task --example tasks +``` diff --git a/crates/rb-progress/examples/failure.rs b/crates/rb-progress/examples/failure.rs new file mode 100644 index 0000000..dbece7e --- /dev/null +++ b/crates/rb-progress/examples/failure.rs @@ -0,0 +1,18 @@ +use rb_progress::Reporter; +use std::{thread, time::Duration}; + +fn inspect_teacup() -> Result<(), String> { + thread::sleep(Duration::from_millis(800)); + Err("Replace the teacup".into()) +} + +fn main() { + let reporter = Reporter::start("Inspecting the china").unwrap(); + let progress = reporter.handle(); + progress.event("inspecting", 0, 1, None); + match inspect_teacup() { + Ok(()) => progress.event("inspecting", 1, 1, None), + Err(error) => progress.fail(error), + } + reporter.finish(); +} diff --git a/crates/rb-progress/examples/parallel.rs b/crates/rb-progress/examples/parallel.rs new file mode 100644 index 0000000..0c442f2 --- /dev/null +++ b/crates/rb-progress/examples/parallel.rs @@ -0,0 +1,34 @@ +use rb_progress::Reporter; +use std::{thread, time::Duration}; + +fn main() { + let reporter = Reporter::start("Preparing the dining room").unwrap(); + let progress = reporter.handle(); + progress.event("preparing", 0, 3, None); + + thread::scope(|scope| { + for (worker, job) in ["polishing silver", "folding napkins", "arranging flowers"] + .into_iter() + .enumerate() + { + let progress = progress.clone(); + scope.spawn(move || { + progress.worker_started(worker, job); + for done in 1..=3 { + thread::sleep(Duration::from_millis(300)); + progress.worker_event( + Some(worker), + job, + done, + 3, + Some(format!("{done}/3 pieces")), + ); + } + progress.worker_finished(worker); + }); + } + }); + + progress.event("preparing", 3, 3, None); + reporter.finish_with_summary("The dining room is ready, sir"); +} diff --git a/crates/rb-progress/examples/progress_ui_probe.rs b/crates/rb-progress/examples/progress_ui_probe.rs new file mode 100644 index 0000000..1c2b7da --- /dev/null +++ b/crates/rb-progress/examples/progress_ui_probe.rs @@ -0,0 +1,138 @@ +use rb_progress::Reporter; +use std::{thread, time::Duration}; + +fn main() { + if std::env::args().any(|arg| arg == "--aggregate") { + let reporter = Reporter::start("processing documents").unwrap(); + reporter.handle().event("reading", 1, 1, None); + reporter.handle().event("writing", 1, 1, None); + reporter.finish_with_summary("documents ready"); + return; + } + if std::env::args().any(|arg| matches!(arg.as_str(), "--tasks" | "--stress" | "--failure")) { + task_probe(); + return; + } + let reporter = Reporter::start("preparing pty-probe bundle").unwrap(); + let progress = reporter.handle(); + thread::sleep(Duration::from_millis(350)); + progress.event("resolved", 1, 1, None); + for worker in 0..12 { + progress.event("downloading", worker, 12, None); + progress.worker_event( + Some(worker), + "downloading", + worker, + 12, + Some(format!("gem-{worker}")), + ); + progress.event("unpacking", worker + 1, 12, None); + } + thread::sleep(Duration::from_millis(300)); + progress.event("packing LMDB", 0, 1, Some("records".into())); + progress.worker_event(Some(9), "packing LMDB", 1, 1, Some("records".into())); + progress.event("packing LMDB", 1, 1, Some("records".into())); + progress.event("preparing native extensions", 0, 0, None); + if std::env::args().any(|arg| arg == "--handoff-output") { + let mut graph = rb_task::TaskGraph::new(); + let task = graph.add("preparing content", []); + for line in ["preview one", "preview two", "preview three"] { + reporter.task_sink().event(rb_task::TaskEvent::Output { + task, + worker: 11, + line: line.into(), + elapsed: Duration::ZERO, + }); + } + } + progress.suspend(); + let native_lines = 10; + eprint!("\x1b[1B\r\x1b[{}L\x1b[1A\r", native_lines - 1); + let rows = [ + "├─ ● installing pty-probe bundle | compiling native extensions", + "│ ├─ ● worker 0 (compiling) demo-1.0", + "│ checking for ruby.h... yes", + "│ └─ ● worker 1 (compiling) demo-1.0", + "│ compiling demo.c", + "└─ overall: 1/2 native tasks | pending 1 | elapsed 0.3s", + ]; + for index in 0..native_lines { + let line = rows.get(index).copied().unwrap_or(""); + eprint!("\x1b[2K{line}"); + if index + 1 < native_lines { + eprintln!(); + } + } + eprint!( + "\x1b[{}A\r\x1b[0J\x1b[{}M", + native_lines - 1, + native_lines - 1 + ); + progress.event("preparing", 1, 1, None); + thread::sleep(Duration::from_millis(50)); + reporter.finish_with_summary("your meticulously prepared bundle is ready, sir (.rb)"); + println!("probe complete"); +} + +fn task_probe() { + use rb_task::{Executor, TaskAction, TaskEvents, TaskGraph}; + use std::{collections::BTreeMap, sync::Arc}; + + let stress = std::env::args().any(|arg| arg == "--stress"); + let fail = std::env::args().any(|arg| arg == "--failure"); + let reporter = Reporter::start("processing documents").unwrap(); + let events = TaskEvents::default(); + events.subscribe(reporter.task_sink()); + let mut graph = TaskGraph::new(); + let prepare = graph.add("preparing content", []); + let publish = graph.add("exporting documents", [prepare]); + let mut actions = BTreeMap::from([ + ( + prepare, + Arc::new(|context: rb_task::TaskContext| { + context.progress("preparing content", 0, 1, None); + context.output("reading document"); + thread::sleep(Duration::from_millis(300)); + Ok(()) + }) as TaskAction, + ), + ( + publish, + Arc::new(|context: rb_task::TaskContext| { + context.progress("exporting documents", 0, 1, None); + context.output("writing document"); + thread::sleep(Duration::from_millis(300)); + Ok(()) + }) as TaskAction, + ), + ]); + if stress || fail { + for index in 0..16 { + let task = graph.add(format!("document {index}"), []); + actions.insert( + task, + Arc::new(move |context: rb_task::TaskContext| { + context.progress("failed", 0, 1, None); + context.output(format!( + "\x1b[2J\x1b[31m{}\x1b[0m\nworking on document {index}\npreview e\u{301}", + "資料".repeat(100) + )); + thread::sleep(Duration::from_millis(30)); + if fail && index == 0 { + Err("document failure".into()) + } else { + Ok(()) + } + }) as TaskAction, + ); + } + } + let result = Executor::new(if stress || fail { 6 } else { 2 }, events).run(graph, actions); + assert_eq!(result.is_err(), fail); + reporter.finish_with_summary(if fail { + "processing failed" + } else { + "documents ready" + }); + println!("probe complete"); +} diff --git a/crates/rb-progress/examples/standalone.rs b/crates/rb-progress/examples/standalone.rs new file mode 100644 index 0000000..569214f --- /dev/null +++ b/crates/rb-progress/examples/standalone.rs @@ -0,0 +1,12 @@ +use rb_progress::Reporter; +use std::{thread, time::Duration}; + +fn main() { + let reporter = Reporter::start("Polishing the silver").unwrap(); + let handle = reporter.handle(); + for done in 0..=3 { + handle.event("polishing", done, 3, None); + thread::sleep(Duration::from_millis(350)); + } + reporter.finish_with_summary("The silver is ready, sir"); +} diff --git a/crates/rb-progress/examples/tasks.rs b/crates/rb-progress/examples/tasks.rs new file mode 100644 index 0000000..86de6c2 --- /dev/null +++ b/crates/rb-progress/examples/tasks.rs @@ -0,0 +1,48 @@ +use rb_progress::Reporter; +use rb_task::{Executor, TaskAction, TaskContext, TaskEvents, TaskGraph}; +use std::{collections::BTreeMap, sync::Arc, thread, time::Duration}; + +fn main() -> Result<(), Box> { + let reporter = Reporter::start("Preparing afternoon tea").unwrap(); + let events = TaskEvents::default(); + events.subscribe(reporter.task_sink()); + let mut graph = TaskGraph::new(); + let kettle = graph.add("Boil water", []); + let table = graph.add("Lay the table", []); + let serve = graph.add("Serve tea", [kettle, table]); + let actions = BTreeMap::from([ + ( + kettle, + Arc::new(|ctx: TaskContext| { + ctx.output("The kettle is warming."); + thread::sleep(Duration::from_millis(1200)); + Ok(()) + }) as TaskAction, + ), + ( + table, + Arc::new(|ctx: TaskContext| { + ctx.progress("arranging", 0, 1, Some("cups and saucers".into())); + ctx.output("Two places, neatly arranged."); + thread::sleep(Duration::from_millis(600)); + Ok(()) + }) as TaskAction, + ), + ( + serve, + Arc::new(|ctx: TaskContext| { + ctx.output("Your tea awaits, sir."); + thread::sleep(Duration::from_millis(400)); + Ok(()) + }) as TaskAction, + ), + ]); + let result = Executor::new(2, events).run(graph, actions); + reporter.finish_with_summary(if result.is_ok() { + "Tea is served" + } else { + "Service interrupted" + }); + result?; + Ok(()) +} diff --git a/crates/rb-progress/src/events.rs b/crates/rb-progress/src/events.rs new file mode 100644 index 0000000..f14a62b --- /dev/null +++ b/crates/rb-progress/src/events.rs @@ -0,0 +1,11 @@ +use std::time::Duration; + +#[derive(Clone, Debug)] +pub struct ProgressEvent { + pub worker: Option, + pub phase: &'static str, + pub done: usize, + pub total: usize, + pub detail: Option, + pub elapsed: Option, +} diff --git a/crates/rb-progress/src/lib.rs b/crates/rb-progress/src/lib.rs new file mode 100644 index 0000000..160564b --- /dev/null +++ b/crates/rb-progress/src/lib.rs @@ -0,0 +1,14 @@ +#![doc = include_str!("../README.md")] + +mod events; +mod plain; +mod render; +mod reporter; +mod state; +#[cfg(feature = "rb-task")] +mod task; +mod terminal; + +pub use events::ProgressEvent; +pub use reporter::{ProgressHandle, Reporter}; +pub use terminal::format_duration; diff --git a/crates/rb-progress/src/plain.rs b/crates/rb-progress/src/plain.rs new file mode 100644 index 0000000..c744c8f --- /dev/null +++ b/crates/rb-progress/src/plain.rs @@ -0,0 +1,46 @@ +use crate::{ + events::ProgressEvent, + state::ProgressState, + terminal::{format_duration, plain_text}, +}; + +pub(crate) fn line(state: &ProgressState, prefix: &str, text: &str) { + if let Some(output) = &state.output { + let mut output = output.lock().unwrap(); + for text in text.split('\n') { + let _ = writeln!(output, "[{prefix}] {}", plain_text(text)); + } + let _ = output.flush(); + } +} + +pub(crate) fn progress(state: &ProgressState, event: &ProgressEvent) { + let prefix = event + .worker + .map_or_else(|| "overall".into(), |worker| format!("worker: {worker}")); + let detail = event.detail.as_deref().unwrap_or(""); + line( + state, + &prefix, + &format!("{} {}/{} {detail}", event.phase, event.done, event.total), + ); +} + +pub(crate) fn finish(state: &ProgressState, summary: Option<&str>) { + let status = if state.failed_tasks > 0 { + "failed" + } else { + "complete" + }; + let detail = summary.map_or(String::new(), |summary| format!(": {summary}")); + line( + state, + "overall", + &format!( + "{status}{detail} | completed {} | failed {} ({})", + state.completed_tasks, + state.failed_tasks, + format_duration(state.elapsed()) + ), + ); +} diff --git a/crates/rb-progress/src/render.rs b/crates/rb-progress/src/render.rs new file mode 100644 index 0000000..0b8ee5c --- /dev/null +++ b/crates/rb-progress/src/render.rs @@ -0,0 +1,363 @@ +use crate::state::{ProgressState, WorkerStatus}; +use crate::terminal::{ + activity_icon, clear_sequence, fit_terminal, format_duration, limit_tree_viewport, + viewport_height, +}; + +fn print(state: &ProgressState, arguments: std::fmt::Arguments<'_>) { + if let Some(output) = &state.output { + let _ = output.lock().unwrap().write_fmt(arguments); + } +} + +fn println(state: &ProgressState, arguments: std::fmt::Arguments<'_>) { + print(state, arguments); + print(state, format_args!("\n")); +} + +fn flush(state: &ProgressState) { + if let Some(output) = &state.output { + let _ = output.lock().unwrap().flush(); + } +} + +pub(crate) fn draw(state: &mut ProgressState) { + if state.plain { + return; + } + if state.phase.is_empty() && state.history.is_empty() && state.workers.is_empty() { + return; + } + if state.external_output { + if !state.handoff_drawn { + clear_render(state); + let lines = tree_frame(state); + repaint(state, &lines); + let active_row = lines + .iter() + .position(|line| line.starts_with("│ ├─ ●")) + .unwrap_or(1); + let active_offset = lines.len() - active_row - 1; + if active_offset > 0 { + print(state, format_args!("\x1b[{}A\r", active_offset)); + } + state.handoff_drawn = true; + state.handoff_phase = Some(state.phase); + state.handoff_tail_rows = lines.len() - active_row; + state.rendered_lines = 0; + state.frame_capacity = 0; + flush(state); + } + return; + } + let lines = tree_frame(state); + repaint(state, &lines); + flush(state); +} + +pub(crate) fn tree_frame(state: &ProgressState) -> Vec { + let elapsed = state.elapsed(); + let icon = activity_icon(elapsed); + let header = fit_terminal(&format!("┌─ ● {}", state.title)); + let mut body = body_lines(state); + if body.is_empty() { + body.push("│ · working…".into()); + } + let footer = if state.phase.is_empty() { + let active = state + .workers + .values() + .filter(|worker| worker.status == WorkerStatus::Running) + .count(); + fit_terminal(&format!( + "└─ {icon} completed {} | active {} | failed {} | elapsed {}", + state.completed_tasks, + active, + state.failed_tasks, + format_duration(elapsed) + )) + } else { + fit_terminal(&format!( + "└─ {icon} overall: {} {}/{} | workers {} | elapsed {}", + state.phase, + state.done, + state.total, + state.workers.len(), + format_duration(elapsed) + )) + }; + limit_tree_viewport(header, body, footer, viewport_height()) +} + +pub(crate) fn body_lines(state: &ProgressState) -> Vec { + let mut lines = state + .history + .iter() + .map(|record| { + let stats = if record.total > 0 { + format!("{}/{}", record.done, record.total) + } else { + "done".into() + }; + format!( + "│ ├─ ✓ {} {} ({})", + record.phase, + stats, + format_duration(record.elapsed) + ) + }) + .collect::>(); + if !state.phase.is_empty() { + lines.push(format!( + "│ ├─ ● {} ({})", + state.phase, + format_duration(state.phase_elapsed()) + )); + } + for (position, (worker, worker_state)) in state.workers.iter().enumerate() { + let branch = if position + 1 == state.workers.len() { + "└─" + } else { + "├─" + }; + let detail = worker_state + .detail + .as_deref() + .map_or(String::new(), |detail| format!(" | {detail}")); + let marker = match worker_state.status { + WorkerStatus::Running => "●", + WorkerStatus::Succeeded => "✓", + WorkerStatus::Failed => "×", + }; + lines.push(format!( + "│ {} {marker} worker {} ({}){} ({})", + branch, + worker, + worker_state.phase, + detail, + format_duration(worker_state.elapsed.map_or_else( + || worker_state.started.elapsed(), + |elapsed| if worker_state.status != WorkerStatus::Running { + elapsed + } else { + elapsed + worker_state.started.elapsed() + } + ),) + )); + for line in &worker_state.output { + lines.push(format!("│ {line}")); + } + } + lines.into_iter().map(|line| fit_terminal(&line)).collect() +} + +#[cfg(test)] +pub(crate) fn history_lines(state: &ProgressState) -> Vec { + state + .history + .iter() + .enumerate() + .map(|(index, record)| { + let stats = if record.total > 0 { + format!("{}/{}", record.done, record.total) + } else { + "done".into() + }; + let detail = record + .detail + .as_deref() + .map_or(String::new(), |detail| format!(" | {detail}")); + fit_terminal(&format!( + "{} ✓ {} {}{} ({})", + if index == 0 { "┌─" } else { "├─" }, + record.phase, + stats, + detail, + format_duration(record.elapsed) + )) + }) + .collect() +} + +fn repaint(state: &mut ProgressState, lines: &[String]) { + let first_frame = state.rendered_lines == 0 && !state.anchor_started; + if state.rendered_lines > 0 { + print( + state, + format_args!("\x1b[{}A\r", state.rendered_lines.saturating_sub(1)), + ); + } else if first_frame { + print(state, format_args!("\r")); + state.anchor_started = true; + } else { + print(state, format_args!("\r")); + } + if state.frame_capacity > 0 { + print(state, format_args!("\x1b[{}M", state.frame_capacity)); + } + if first_frame { + if let Some(first) = lines.first() { + print(state, format_args!("\x1b[2K{}", first)); + } + if lines.len() > 1 { + print(state, format_args!("\n")); + print(state, format_args!("\x1b[{}L", lines.len() - 1)); + for (index, line) in lines.iter().skip(1).enumerate() { + print(state, format_args!("\x1b[2K{}", line)); + if index + 1 < lines.len() - 1 { + print(state, format_args!("\n")); + } + } + } + } else { + print(state, format_args!("\x1b[{}L", lines.len())); + for (index, line) in lines.iter().enumerate() { + print(state, format_args!("\x1b[2K{}", line)); + if index + 1 < lines.len() { + print(state, format_args!("\n")); + } + } + } + state.rendered_lines = lines.len(); + state.frame_capacity = lines.len(); +} + +fn clear_render(state: &mut ProgressState) { + if state.rendered_lines == 0 { + return; + } + print( + state, + format_args!("{}", clear_sequence(state.rendered_lines)), + ); + if state.frame_capacity > 0 { + print(state, format_args!("\x1b[{}M", state.frame_capacity)); + } + state.rendered_lines = 0; + state.frame_capacity = 0; +} + +pub(crate) fn finish(state: &mut ProgressState, summary: Option) { + if state.plain { + crate::plain::finish(state, summary.as_deref()); + return; + } + if state.started.is_none() { + return; + } + if state.handoff_drawn { + let line = fit_terminal(&format!( + "├─ ✓ {} | {} done ({})", + state.title, + state.phase, + format_duration(state.phase_elapsed()) + )); + print(state, format_args!("\r")); + for index in 0..state.handoff_tail_rows { + print(state, format_args!("\x1b[2K")); + if index + 1 < state.handoff_tail_rows { + println(state, format_args!("")); + } + } + if state.handoff_tail_rows > 1 { + print( + state, + format_args!("\x1b[{}A\r", state.handoff_tail_rows - 1), + ); + } + print(state, format_args!("\x1b[{}M", state.handoff_tail_rows)); + if let Some(phase) = state.handoff_phase + && let Some(record) = state + .history + .iter() + .rev() + .find(|record| record.phase == phase) + { + let stats = if record.total > 0 { + format!("{}/{}", record.done, record.total) + } else { + "done".into() + }; + println( + state, + format_args!( + "├─ ✓ {} {} ({})", + phase, + stats, + format_duration(record.elapsed) + ), + ); + } + println(state, format_args!("{}", line)); + if let Some(summary) = summary { + println(state, format_args!("└─ ✓ {}", fit_terminal(&summary))); + } + state.handoff_drawn = false; + state.handoff_phase = None; + state.handoff_tail_rows = 0; + state.rendered_lines = 0; + state.frame_capacity = 0; + } else { + let failed = state.failed_tasks > 0; + let marker = if failed { "×" } else { "✓" }; + let header = fit_terminal(&format!("┌─ {marker} {}", state.title)); + let mut body = state + .history + .iter() + .map(|record| { + fit_terminal(&format!( + "│ ├─ ✓ {} {}/{} ({})", + record.phase, + record.done, + record.total, + format_duration(record.elapsed) + )) + }) + .collect::>(); + if !state.phase.is_empty() { + body.push(fit_terminal(&format!( + "│ ├─ {marker} {} {}/{} ({})", + state.phase, + state.done, + state.total, + format_duration(state.phase_elapsed()) + ))); + } + let failure_lines = state + .failures + .iter() + .flat_map(|error| { + error + .lines() + .map(|line| fit_terminal(&format!("│ × {line}"))) + }) + .collect::>(); + body.extend(failure_lines); + let compact_error = failed && viewport_height() <= 3; + let footer = fit_terminal(&format!( + "└─ {marker} {} ({})", + if compact_error { + state + .failures + .last() + .map(String::as_str) + .unwrap_or("failed") + } else { + summary + .as_deref() + .unwrap_or(if failed { "failed" } else { "complete" }) + }, + format_duration(state.elapsed()) + )); + let limit = viewport_height().saturating_sub(1).max(2); + if body.len() > limit.saturating_sub(2) { + body.drain(..body.len() - limit.saturating_sub(2)); + } + let lines = limit_tree_viewport(header, body, footer, limit); + repaint(state, &lines); + println(state, format_args!("")); + state.rendered_lines = 0; + state.frame_capacity = 0; + } + flush(state); +} diff --git a/crates/rb-progress/src/reporter.rs b/crates/rb-progress/src/reporter.rs new file mode 100644 index 0000000..5aeeba0 --- /dev/null +++ b/crates/rb-progress/src/reporter.rs @@ -0,0 +1,237 @@ +use std::{ + io::{self, IsTerminal, Write}, + sync::{ + Arc, Mutex, + atomic::{AtomicBool, Ordering}, + }, + thread, + time::{Duration, Instant}, +}; + +#[cfg(feature = "rb-task")] +use crate::task::TaskBridge; +#[cfg(feature = "rb-task")] +use rb_task::TaskEventSink; + +use crate::{events::ProgressEvent, render, state::ProgressState}; + +#[derive(Clone)] +pub struct ProgressHandle { + progress: Arc>, +} + +impl ProgressHandle { + pub fn emit(&self, event: ProgressEvent) { + let mut progress = self.progress.lock().unwrap(); + if progress.closed { + return; + } + if progress.plain { + crate::plain::progress(&progress, &event); + } + progress.show(event); + render::draw(&mut progress); + } + + pub fn event(&self, phase: &'static str, done: usize, total: usize, detail: Option) { + self.worker_event(None, phase, done, total, detail); + } + + pub fn worker_event( + &self, + worker: Option, + phase: &'static str, + done: usize, + total: usize, + detail: Option, + ) { + self.emit(ProgressEvent { + worker, + phase, + done, + total, + detail, + elapsed: None, + }); + } + + pub fn worker_started(&self, worker: usize, label: impl Into) { + let mut progress = self.progress.lock().unwrap(); + if progress.closed { + return; + } + let label = label.into(); + if progress.plain { + crate::plain::line( + &progress, + &format!("worker: {worker}"), + &format!("started: {label}"), + ); + } + progress.show(ProgressEvent { + worker: Some(worker), + phase: "running", + done: 0, + total: 1, + detail: Some(label), + elapsed: None, + }); + progress.set_worker_status(worker, crate::state::WorkerStatus::Running); + render::draw(&mut progress); + } + + pub fn worker_finished(&self, worker: usize) { + self.end_worker(worker, None); + } + + pub fn worker_failed(&self, worker: usize, error: impl Into) { + self.end_worker(worker, Some(error.into())); + } + + fn end_worker(&self, worker: usize, error: Option) { + let mut progress = self.progress.lock().unwrap(); + let status = if error.is_some() { + crate::state::WorkerStatus::Failed + } else { + crate::state::WorkerStatus::Succeeded + }; + if progress.closed || !progress.set_worker_status(worker, status) { + return; + } + let message = if let Some(error) = error { + progress.workers.get_mut(&worker).unwrap().detail = Some(error.clone()); + progress.record_failure(format!("{error} (worker {worker})")); + format!("failed: {error}") + } else { + "complete".into() + }; + if progress.plain { + crate::plain::line(&progress, &format!("worker: {worker}"), &message); + } + render::draw(&mut progress); + } + + pub fn fail(&self, error: impl Into) { + let mut progress = self.progress.lock().unwrap(); + if progress.closed { + return; + } + let error = error.into(); + if progress.plain { + crate::plain::line(&progress, "overall", &format!("failed: {error}")); + } + progress.record_failure(error); + render::draw(&mut progress); + } + + pub fn suspend(&self) { + let mut progress = self.progress.lock().unwrap(); + if progress.closed { + return; + } + progress.external_output = true; + render::draw(&mut progress); + } +} + +pub struct Reporter { + progress: Arc>, + running: Arc, + ticker: Option>, + finished: bool, +} + +impl Reporter { + pub fn start(title: impl Into) -> Option { + if !io::stderr().is_terminal() { + return Some(Self::with_plain_writer(title, io::stderr())); + } + Some(Self::with_writer(title, io::stderr())) + } + + pub fn with_writer(title: impl Into, writer: impl Write + Send + 'static) -> Self { + Self::build(title.into(), writer, false) + } + + pub fn with_plain_writer( + title: impl Into, + writer: impl Write + Send + 'static, + ) -> Self { + Self::build(title.into(), writer, true) + } + + fn build(title: String, writer: impl Write + Send + 'static, plain: bool) -> Self { + let progress = Arc::new(Mutex::new(ProgressState { + plain, + title, + output: Some(Mutex::new(Box::new(io::BufWriter::new(writer)))), + started: Some(Instant::now()), + ..Default::default() + })); + let running = Arc::new(AtomicBool::new(true)); + let ticker_progress = progress.clone(); + let ticker_running = running.clone(); + if plain { + let state = progress.lock().unwrap(); + crate::plain::line(&state, "overall", &state.title); + } + let ticker = (!plain).then(|| { + thread::spawn(move || { + while ticker_running.load(Ordering::Relaxed) { + render::draw(&mut ticker_progress.lock().unwrap()); + thread::sleep(Duration::from_millis(250)); + } + }) + }); + Self { + progress, + running, + ticker, + finished: false, + } + } + + pub fn finish(mut self) { + self.stop(None); + } + + pub fn finish_with_summary(mut self, summary: impl Into) { + self.stop(Some(summary.into())); + } + + pub fn handle(&self) -> ProgressHandle { + ProgressHandle { + progress: self.progress.clone(), + } + } + + #[cfg(feature = "rb-task")] + pub fn task_sink(&self) -> Arc { + Arc::new(TaskBridge { + progress: self.progress.clone(), + }) + } + + fn stop(&mut self, summary: Option) { + if self.finished { + return; + } + self.finished = true; + self.running.store(false, Ordering::Relaxed); + if let Some(ticker) = self.ticker.take() { + let _ = ticker.join(); + } + let mut progress = self.progress.lock().unwrap(); + render::finish(&mut progress, summary); + progress.closed = true; + } +} + +impl Drop for Reporter { + fn drop(&mut self) { + self.stop(None); + } +} + +#[cfg(test)] +mod tests; diff --git a/crates/rb-progress/src/reporter/tests.rs b/crates/rb-progress/src/reporter/tests.rs new file mode 100644 index 0000000..879ff6f --- /dev/null +++ b/crates/rb-progress/src/reporter/tests.rs @@ -0,0 +1,499 @@ +use super::*; +#[cfg(feature = "rb-task")] +use rb_task::TaskEvent; + +#[test] +fn completed_phases_remain_as_checked_history_rows() { + let mut progress = ProgressState { + title: "demo".into(), + ..Default::default() + }; + progress.show(ProgressEvent { + worker: None, + phase: "resolving", + done: 1, + total: 1, + detail: None, + elapsed: None, + }); + progress.show(ProgressEvent { + worker: None, + phase: "downloading", + done: 0, + total: 2, + detail: None, + elapsed: None, + }); + let rows = render::history_lines(&progress); + assert_eq!(rows.len(), 1); + assert!(rows[0].contains("✓ resolving")); + assert!(rows[0].contains("1/1")); +} + +#[test] +fn native_handoff_is_rendered_once_until_completion() { + let mut progress = ProgressState { + title: "demo".into(), + phase: "preparing native extensions", + ..Default::default() + }; + progress.external_output = true; + render::draw(&mut progress); + assert!(progress.handoff_drawn); + let lines = progress.rendered_lines; + render::draw(&mut progress); + assert_eq!(progress.rendered_lines, lines); +} + +#[test] +fn fetch_progress_updates_keep_separate_phase_rows() { + let mut progress = ProgressState { + title: "demo".into(), + ..Default::default() + }; + for phase in ["downloading", "unpacking", "downloading", "unpacking"] { + progress.show(ProgressEvent { + worker: None, + phase, + done: 1, + total: 2, + detail: None, + elapsed: None, + }); + } + assert_eq!(progress.history.len(), 2); + assert_eq!(progress.history[0].phase, "downloading"); + assert_eq!(progress.history[1].phase, "unpacking"); + assert_eq!(progress.phase, "unpacking"); + assert_eq!(render::history_lines(&progress).len(), 2); + assert!(render::history_lines(&progress)[0].contains("downloading")); +} + +#[test] +fn empty_reporter_does_not_claim_prompt_rows() { + let mut progress = ProgressState::default(); + render::draw(&mut progress); + assert_eq!(progress.rendered_lines, 0); + assert!(!progress.anchor_started); +} + +#[derive(Clone, Default)] +struct Buffer(Arc>>); +impl Write for Buffer { + fn write(&mut self, bytes: &[u8]) -> io::Result { + self.0.lock().unwrap().extend_from_slice(bytes); + Ok(bytes.len()) + } + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } +} +impl Buffer { + fn text(&self) -> String { + String::from_utf8(self.0.lock().unwrap().clone()).unwrap() + } +} + +#[test] +fn reporters_have_independent_sinks_and_ignore_events_after_finish() { + let first = Buffer::default(); + let second = Buffer::default(); + let a = Reporter::with_writer("Building documents", first.clone()); + let b = Reporter::with_writer("Copying images", second.clone()); + let old = a.handle(); + old.event("preparing content", 0, 3, None); + assert!(!a.progress.lock().unwrap().external_output); + b.handle().event("copying", 1, 2, None); + a.finish_with_summary("Documents ready"); + let finished = first.text(); + old.event("late event", 1, 1, None); + assert_eq!(first.text(), finished); + b.handle().event("copying", 2, 2, None); + b.finish_with_summary("Images ready"); + assert!(finished.contains("Building documents") && finished.contains("Documents ready")); + assert!(!finished.contains("Copying images") && !finished.contains("bundle")); + assert!( + second.text().contains("Images ready") && !second.text().contains("Building documents") + ); +} + +#[cfg(feature = "rb-task")] +#[test] +fn task_output_stays_on_its_worker_and_completion_freezes_the_clock() { + let reporter = Reporter::with_writer("Processing files", Buffer::default()); + let sink = reporter.task_sink(); + let mut graph = rb_task::TaskGraph::new(); + let task = graph.add("document", []); + sink.event(TaskEvent::Started { + task, + worker: 0, + label: "document".into(), + elapsed: Duration::ZERO, + }); + for line in ["one", "two", "three", "four"] { + sink.event(TaskEvent::Output { + task, + worker: 0, + line: line.into(), + elapsed: Duration::ZERO, + }); + } + sink.event(TaskEvent::Finished { + task, + worker: 0, + elapsed: Duration::from_secs(12), + }); + let state = reporter.progress.lock().unwrap(); + assert!(state.phase.is_empty()); + assert_eq!( + state.workers[&0] + .output + .iter() + .map(String::as_str) + .collect::>(), + ["two", "three", "four"] + ); + let before = render::body_lines(&state) + .into_iter() + .find(|line| line.contains("worker 0")); + drop(state); + reporter + .progress + .lock() + .unwrap() + .workers + .get_mut(&0) + .unwrap() + .started -= Duration::from_secs(2); + assert_eq!( + render::body_lines(&reporter.progress.lock().unwrap()) + .into_iter() + .find(|line| line.contains("worker 0")), + before + ); + reporter.finish(); +} + +#[cfg(feature = "rb-task")] +#[test] +fn failed_task_is_not_presented_as_successful_completion() { + let buffer = Buffer::default(); + let reporter = Reporter::with_writer("processing documents", buffer.clone()); + let mut graph = rb_task::TaskGraph::new(); + let task = graph.add("export", []); + reporter.task_sink().event(TaskEvent::Failed { + task, + worker: 0, + error: "fixture".into(), + elapsed: Duration::from_secs(1), + }); + reporter.finish(); + let output = buffer.text(); + assert!(output.contains("┌─ × processing documents")); + assert!(output.contains("└─ × failed")); + assert!(!output.contains("└─ ✓ complete")); +} + +#[cfg(feature = "rb-task")] +#[test] +fn worker_reuse_keeps_failures_and_clears_only_the_previous_task_output() { + let buffer = Buffer::default(); + let reporter = Reporter::with_writer("documents", buffer.clone()); + let sink = reporter.task_sink(); + let mut graph = rb_task::TaskGraph::new(); + let first = graph.add("first", []); + let second = graph.add("second", []); + sink.event(TaskEvent::Started { + task: first, + worker: 0, + label: "first".into(), + elapsed: Duration::ZERO, + }); + sink.event(TaskEvent::Output { + task: first, + worker: 0, + line: "one\ntwo\nthree\nfour".into(), + elapsed: Duration::ZERO, + }); + sink.event(TaskEvent::Progress { + task: first, + worker: 0, + phase: "running", + done: 0, + total: 1, + detail: None, + elapsed: Duration::ZERO, + }); + let output = reporter.progress.lock().unwrap().workers[&0].output.clone(); + assert_eq!( + output.into_iter().collect::>(), + ["two", "three", "four"] + ); + sink.event(TaskEvent::Failed { + task: first, + worker: 0, + error: "fixture".into(), + elapsed: Duration::from_secs(1), + }); + sink.event(TaskEvent::Started { + task: second, + worker: 0, + label: "second".into(), + elapsed: Duration::ZERO, + }); + let empty = reporter.progress.lock().unwrap().workers[&0] + .output + .is_empty(); + assert!(empty); + sink.event(TaskEvent::Finished { + task: second, + worker: 0, + elapsed: Duration::from_secs(2), + }); + let frame = render::tree_frame(&reporter.progress.lock().unwrap()); + assert!( + frame + .last() + .unwrap() + .contains("completed 1 | active 0 | failed 1") + ); + reporter.finish(); + assert!(buffer.text().contains("└─ × failed")); +} + +#[cfg(feature = "rb-task")] +#[test] +fn phase_names_do_not_stop_running_timers() { + let reporter = Reporter::with_writer("documents", Vec::new()); + let mut graph = rb_task::TaskGraph::new(); + let task = graph.add("work", []); + for phase in ["done", "failed", "running"] { + reporter.task_sink().event(TaskEvent::Progress { + task, + worker: 0, + phase, + done: 0, + total: 1, + detail: None, + elapsed: Duration::from_secs(1), + }); + let mut state = reporter.progress.lock().unwrap(); + state.workers.get_mut(&0).unwrap().started -= Duration::from_secs(2); + let lines = render::body_lines(&state); + drop(state); + assert!(lines.iter().any(|line| line.contains("(3.0s)"))); + } + reporter.finish(); +} + +#[cfg(feature = "rb-task")] +#[test] +fn plain_output_prefixes_every_line_and_keeps_errors_after_worker_reuse() { + let buffer = Buffer::default(); + let reporter = Reporter::with_plain_writer("china", buffer.clone()); + assert!(reporter.ticker.is_none()); + let sink = reporter.task_sink(); + let mut graph = rb_task::TaskGraph::new(); + let first = graph.add("inspect", []); + let second = graph.add("replace", []); + sink.event(TaskEvent::Started { + task: first, + worker: 0, + label: "inspect".into(), + elapsed: Duration::ZERO, + }); + sink.event(TaskEvent::Output { + task: first, + worker: 0, + line: "one\ntwo\x1b[2J".into(), + elapsed: Duration::ZERO, + }); + sink.event(TaskEvent::Failed { + task: first, + worker: 0, + error: "Replace the teacup".into(), + elapsed: Duration::from_secs(1), + }); + sink.event(TaskEvent::Started { + task: second, + worker: 0, + label: "replace".into(), + elapsed: Duration::ZERO, + }); + sink.event(TaskEvent::Finished { + task: second, + worker: 0, + elapsed: Duration::from_secs(2), + }); + reporter.finish(); + let text = buffer.text(); + assert!(!text.contains('\x1b')); + assert!(!text.contains('\r')); + assert!(text.contains("[worker: 0 | task: 0] one\n[worker: 0 | task: 0] two\n")); + assert!(text.contains("[worker: 0 | task: 0] failed: Replace the teacup (1.0s)")); + assert!(text.contains("[overall] failed | completed 1 | failed 1")); + sink.event(TaskEvent::Output { + task: second, + worker: 0, + line: "late".into(), + elapsed: Duration::ZERO, + }); + assert_eq!(buffer.text(), text); +} + +#[cfg(feature = "rb-task")] +#[test] +fn terminal_finish_retains_the_error_after_worker_reuse() { + let buffer = Buffer::default(); + let reporter = Reporter::with_writer("china", buffer.clone()); + let sink = reporter.task_sink(); + let mut graph = rb_task::TaskGraph::new(); + let first = graph.add("inspect", []); + let second = graph.add("replace", []); + sink.event(TaskEvent::Failed { + task: first, + worker: 0, + error: "Replace the teacup".into(), + elapsed: Duration::from_secs(1), + }); + sink.event(TaskEvent::Started { + task: second, + worker: 0, + label: "replace".into(), + elapsed: Duration::ZERO, + }); + sink.event(TaskEvent::Finished { + task: second, + worker: 0, + elapsed: Duration::from_secs(2), + }); + buffer.0.lock().unwrap().clear(); + reporter.finish(); + assert!(buffer.text().contains("Replace the teacup")); +} + +#[test] +fn live_summaries_include_an_activity_icon_for_both_reporting_styles() { + let mut state = ProgressState::default(); + for phase in ["", "polishing"] { + state.phase = phase; + let frame = render::tree_frame(&state); + assert!(frame.last().unwrap().starts_with("└─ ⠋ ")); + } +} + +#[test] +fn standalone_plain_progress_needs_no_executor() { + let buffer = Buffer::default(); + let reporter = Reporter::with_plain_writer("silver", buffer.clone()); + let handle = reporter.handle(); + handle.event("polishing", 1, 3, None); + handle.worker_event(Some(2), "polishing", 2, 3, Some("spoons".into())); + reporter.finish_with_summary("silver ready"); + let output = buffer.text(); + assert!(output.contains("[overall] polishing 1/3")); + assert!(output.contains("[worker: 2] polishing 2/3 spoons")); + assert!(output.contains("[overall] complete: silver ready")); + assert!(!output.contains('\x1b')); + handle.event("late", 3, 3, None); + assert_eq!(buffer.text(), output); +} + +#[test] +fn standalone_failure_survives_later_progress_and_finish() { + for plain in [false, true] { + let buffer = Buffer::default(); + let reporter = if plain { + Reporter::with_plain_writer("china", buffer.clone()) + } else { + Reporter::with_writer("china", buffer.clone()) + }; + let handle = reporter.handle(); + handle.fail("Replace the teacup"); + handle.event("cleanup", 1, 1, None); + if !plain { + buffer.0.lock().unwrap().clear(); + } + reporter.finish(); + let output = buffer.text(); + assert!(output.contains("Replace the teacup")); + assert!(output.contains(if plain { + "[overall] failed" + } else { + "└─ × failed" + })); + handle.fail("late failure"); + assert_eq!(buffer.text(), output); + } +} + +#[test] +fn standalone_workers_complete_once_freeze_timers_and_can_be_reused() { + let buffer = Buffer::default(); + let reporter = Reporter::with_plain_writer("dining room", buffer.clone()); + let handle = reporter.handle(); + handle.worker_started(0, "silver"); + handle.worker_event(Some(0), "polishing", 3, 3, None); + assert_eq!(reporter.progress.lock().unwrap().completed_tasks, 0); + handle.worker_finished(0); + handle.worker_finished(0); + handle.worker_failed(0, "late failure"); + let mut state = reporter.progress.lock().unwrap(); + assert_eq!(state.completed_tasks, 1); + assert_eq!(state.failed_tasks, 0); + let before = render::body_lines(&state); + state.workers.get_mut(&0).unwrap().started -= Duration::from_secs(10); + assert_eq!(render::body_lines(&state), before); + drop(state); + handle.worker_started(0, "china"); + handle.worker_failed(0, "Replace the teacup"); + handle.worker_failed(0, "duplicate"); + handle.worker_started(0, "napkins"); + handle.worker_finished(0); + reporter.finish(); + let text = buffer.text(); + assert!(text.contains("[worker: 0] failed: Replace the teacup")); + assert!(text.contains("[overall] failed | completed 2 | failed 1")); + handle.worker_started(1, "late"); + handle.worker_finished(1); + assert_eq!(buffer.text(), text); +} + +#[test] +fn standalone_threads_share_workers_and_completion_counts() { + use std::sync::Barrier; + let buffer = Buffer::default(); + let reporter = Reporter::with_writer("dining room", buffer.clone()); + let handle = reporter.handle(); + let ready = Barrier::new(4); + let release = Barrier::new(4); + thread::scope(|scope| { + for worker in 0..3 { + let handle = handle.clone(); + let ready = &ready; + let release = &release; + scope.spawn(move || { + handle.worker_started(worker, "polishing"); + ready.wait(); + release.wait(); + handle.worker_finished(worker); + }); + } + ready.wait(); + let state = reporter.progress.lock().unwrap(); + let worker_count = state.workers.len(); + let frame = render::tree_frame(&state).join("\n"); + drop(state); + release.wait(); + assert_eq!(worker_count, 3); + assert!(frame.contains("active 3")); + for worker in 0..3 { + assert!(frame.contains(&format!("worker {worker}"))); + } + }); + let state = reporter.progress.lock().unwrap(); + let frame = render::tree_frame(&state).join("\n"); + assert!(frame.contains("completed 3 | active 0 | failed 0")); + drop(state); + reporter.finish(); +} diff --git a/crates/rb-progress/src/state.rs b/crates/rb-progress/src/state.rs new file mode 100644 index 0000000..a944771 --- /dev/null +++ b/crates/rb-progress/src/state.rs @@ -0,0 +1,163 @@ +use crate::events::ProgressEvent; +use std::{ + collections::BTreeMap, + io::Write, + sync::Mutex, + time::{Duration, Instant}, +}; + +#[derive(Default)] +pub(crate) struct ProgressState { + pub(crate) output: Option>>, + pub(crate) closed: bool, + pub(crate) plain: bool, + pub(crate) failures: Vec, + pub(crate) completed_tasks: usize, + pub(crate) failed_tasks: usize, + pub(crate) title: String, + pub(crate) started: Option, + pub(crate) phase: &'static str, + pub(crate) phase_started: Option, + pub(crate) phase_detail: Option, + pub(crate) done: usize, + pub(crate) total: usize, + pub(crate) workers: BTreeMap, + pub(crate) history: Vec, + pub(crate) rendered_lines: usize, + pub(crate) frame_capacity: usize, + pub(crate) anchor_started: bool, + pub(crate) external_output: bool, + pub(crate) handoff_drawn: bool, + pub(crate) handoff_phase: Option<&'static str>, + pub(crate) handoff_tail_rows: usize, +} + +pub(crate) struct PhaseRecord { + pub(crate) phase: &'static str, + pub(crate) elapsed: Duration, + pub(crate) done: usize, + pub(crate) total: usize, + pub(crate) detail: Option, +} + +#[derive(Clone, Copy, Default, Eq, PartialEq)] +pub(crate) enum WorkerStatus { + #[default] + Running, + Succeeded, + Failed, +} + +pub(crate) struct WorkerProgress { + pub(crate) status: WorkerStatus, + pub(crate) phase: &'static str, + pub(crate) detail: Option, + pub(crate) started: Instant, + pub(crate) elapsed: Option, + pub(crate) output: std::collections::VecDeque, +} + +impl ProgressState { + pub(crate) fn set_worker_status(&mut self, worker: usize, status: WorkerStatus) -> bool { + let Some(state) = self.workers.get_mut(&worker) else { + return false; + }; + if status != WorkerStatus::Running && state.status != WorkerStatus::Running { + return false; + } + if status == WorkerStatus::Running { + state.started = Instant::now(); + state.elapsed = None; + state.output.clear(); + } else { + state.elapsed = Some(state.elapsed.unwrap_or_default() + state.started.elapsed()); + } + state.status = status; + if status == WorkerStatus::Succeeded { + self.completed_tasks += 1; + } + true + } + + pub(crate) fn record_failure(&mut self, error: String) { + self.failed_tasks += 1; + self.failures.push(error); + } + + pub(crate) fn show(&mut self, event: ProgressEvent) { + self.started.get_or_insert_with(Instant::now); + if event.worker.is_none() { + let phase_changed = self.phase != event.phase; + if !self.phase.is_empty() && phase_changed { + let elapsed = self + .phase_started + .map_or(Duration::ZERO, |started| started.elapsed()); + if let Some(record) = self + .history + .iter_mut() + .find(|record| record.phase == self.phase) + { + record.elapsed += elapsed; + record.done = self.done; + record.total = self.total; + record.detail = self.phase_detail.clone(); + } else { + self.history.push(PhaseRecord { + phase: self.phase, + elapsed, + done: self.done, + total: self.total, + detail: self.phase_detail.clone(), + }); + } + } + self.phase = event.phase; + if phase_changed || self.phase_started.is_none() { + self.phase_started = Some(Instant::now()); + } + self.done = event.done; + self.total = event.total; + self.phase_detail = event.detail.clone(); + if !self.external_output { + self.handoff_drawn = false; + } + } + if let Some(worker) = event.worker { + let started = if event.elapsed.is_none() { + self.workers + .get(&worker) + .map_or_else(Instant::now, |state| state.started) + } else { + Instant::now() + }; + self.workers.insert( + worker, + WorkerProgress { + status: self + .workers + .get(&worker) + .map_or(WorkerStatus::Running, |state| state.status), + phase: event.phase, + detail: event.detail, + started, + elapsed: event.elapsed, + output: self + .workers + .get(&worker) + .map(|state| state.output.clone()) + .unwrap_or_default(), + }, + ); + } + } + + pub(crate) fn elapsed(&self) -> Duration { + self.started + .map_or(Duration::ZERO, |started| started.elapsed()) + } + + pub(crate) fn phase_elapsed(&self) -> Duration { + self.phase_started + .map_or(Duration::ZERO, |started| started.elapsed()) + } +} diff --git a/crates/rb-progress/src/task.rs b/crates/rb-progress/src/task.rs new file mode 100644 index 0000000..c1a1aa1 --- /dev/null +++ b/crates/rb-progress/src/task.rs @@ -0,0 +1,163 @@ +use crate::{ + events::ProgressEvent, + state::{ProgressState, WorkerStatus}, + terminal::format_duration, +}; +use rb_task::{TaskEvent, TaskEventSink}; +use std::sync::{Arc, Mutex}; + +pub(crate) struct TaskBridge { + pub(crate) progress: Arc>, +} + +impl TaskEventSink for TaskBridge { + fn event(&self, event: TaskEvent) { + let mut progress = self.progress.lock().unwrap(); + if progress.closed { + return; + } + if progress.plain { + write_plain_event(&progress, &event); + } + if let TaskEvent::Failed { + task, + worker, + error, + .. + } = &event + { + progress.record_failure(format!("{error} (worker {worker}, task {})", task.index())); + } + let transition = match &event { + TaskEvent::Started { worker, .. } => Some((*worker, WorkerStatus::Running)), + TaskEvent::Finished { worker, .. } => Some((*worker, WorkerStatus::Succeeded)), + TaskEvent::Failed { worker, .. } => Some((*worker, WorkerStatus::Failed)), + _ => None, + }; + let event = match event { + TaskEvent::Started { + task, + worker, + label, + elapsed, + } => ProgressEvent { + worker: Some(worker), + phase: "running", + done: 0, + total: 1, + detail: Some(format!("{}: {}", task.index(), label)), + elapsed: Some(elapsed), + }, + TaskEvent::Output { worker, line, .. } => { + if let Some(state) = progress.workers.get_mut(&worker) { + for line in line.lines() { + state.output.push_back(line.to_owned()); + while state.output.len() > 3 { + state.output.pop_front(); + } + } + } + crate::render::draw(&mut progress); + return; + } + TaskEvent::Progress { + worker, + phase, + done, + total, + detail, + elapsed, + .. + } => ProgressEvent { + worker: Some(worker), + phase, + done, + total, + detail, + elapsed: Some(elapsed), + }, + TaskEvent::Finished { + task, + worker, + elapsed, + } => ProgressEvent { + worker: Some(worker), + phase: "done", + done: 1, + total: 1, + detail: Some(format!("{}: complete", task.index())), + elapsed: Some(elapsed), + }, + TaskEvent::Failed { + task, + worker, + error, + elapsed, + } => ProgressEvent { + worker: Some(worker), + phase: "failed", + done: 1, + total: 1, + detail: Some(format!("{}: {error}", task.index())), + elapsed: Some(elapsed), + }, + }; + let elapsed = event.elapsed; + progress.show(event); + if let Some((worker, status)) = transition { + progress.set_worker_status(worker, status); + progress.workers.get_mut(&worker).unwrap().elapsed = elapsed; + } + crate::render::draw(&mut progress); + } +} +fn write_plain_event(state: &ProgressState, event: &TaskEvent) { + let (worker, task, message) = match event { + TaskEvent::Started { + worker, + task, + label, + .. + } => (worker, task, format!("started: {label}")), + TaskEvent::Output { + worker, task, line, .. + } => (worker, task, line.clone()), + TaskEvent::Progress { + worker, + task, + phase, + done, + total, + detail, + .. + } => ( + worker, + task, + format!("{phase} {done}/{total} {}", detail.as_deref().unwrap_or("")), + ), + TaskEvent::Finished { + worker, + task, + elapsed, + } => ( + worker, + task, + format!("complete ({})", format_duration(*elapsed)), + ), + TaskEvent::Failed { + worker, + task, + elapsed, + error, + } => ( + worker, + task, + format!("failed: {error} ({})", format_duration(*elapsed)), + ), + }; + crate::plain::line( + state, + &format!("worker: {worker} | task: {}", task.index()), + &message, + ); +} diff --git a/crates/rb-progress/src/terminal.rs b/crates/rb-progress/src/terminal.rs new file mode 100644 index 0000000..d89d6fc --- /dev/null +++ b/crates/rb-progress/src/terminal.rs @@ -0,0 +1,148 @@ +use std::time::Duration; +use terminal_size::{Height, Width, terminal_size}; +use unicode_width::{UnicodeWidthChar, UnicodeWidthStr}; + +pub(crate) fn activity_icon(elapsed: Duration) -> char { + const FRAMES: [char; 10] = ['⠋', '⠙', '⠹', '⠸', '⠼', '⠴', '⠦', '⠧', '⠇', '⠏']; + FRAMES[((elapsed.as_millis() / 250) % FRAMES.len() as u128) as usize] +} + +pub fn format_duration(duration: Duration) -> String { + let seconds = duration.as_secs(); + if seconds < 10 { + return format!("{:.1}s", duration.as_secs_f64()); + } + if seconds < 60 { + return format!("{}s", seconds); + } + format!("{}m{:02}s", seconds / 60, seconds % 60) +} + +pub(crate) fn fit_terminal(text: &str) -> String { + let width = terminal_size() + .map(|(Width(width), Height(_))| usize::from(width)) + .or_else(|| { + std::env::var("COLUMNS") + .ok() + .and_then(|value| value.parse().ok()) + }) + .unwrap_or(80) + .saturating_sub(1); + fit_width(text, width) +} + +fn fit_width(text: &str, width: usize) -> String { + if width == 0 { + return String::new(); + } + let text = plain_text(text); + if UnicodeWidthStr::width(text.as_str()) <= width { + return text; + } + let mut result = String::new(); + let mut used = 0; + for character in text.chars() { + let columns = character.width().unwrap_or(0); + if used + columns >= width { + break; + } + used += columns; + result.push(character); + } + result.push('…'); + result +} + +pub(crate) fn plain_text(text: &str) -> String { + let mut result = String::new(); + let mut chars = text.chars().peekable(); + while let Some(character) = chars.next() { + match character { + '\x1b' => match chars.next() { + Some('[') => { + for character in chars.by_ref() { + if ('@'..='~').contains(&character) { + break; + } + } + } + Some(']') => { + while let Some(character) = chars.next() { + if character == '\x07' { + break; + } + if character == '\x1b' && chars.peek() == Some(&'\\') { + chars.next(); + break; + } + } + } + _ => {} + }, + '\n' | '\r' => result.push(' '), + '\t' => result.push_str(" "), + character if !character.is_control() => result.push(character), + _ => {} + } + } + result +} + +pub(crate) fn viewport_height() -> usize { + terminal_size() + .map(|(_, Height(height))| usize::from(height)) + .or_else(|| { + std::env::var("LINES") + .ok() + .and_then(|value| value.parse().ok()) + }) + .unwrap_or(24) + .saturating_sub(1) + .max(1) +} + +pub(crate) fn limit_tree_viewport( + header: String, + mut body: Vec, + footer: String, + limit: usize, +) -> Vec { + let limit = limit.max(2); + let body_limit = limit.saturating_sub(2); + if body.len() > body_limit { + if body_limit == 0 { + body.clear(); + } else if body_limit == 1 { + body = vec!["│ · working…".into()]; + } else { + let start = body.len().saturating_sub(body_limit); + body.drain(..start); + } + } + let mut lines = Vec::with_capacity(body.len() + 2); + lines.push(header); + lines.extend(body); + lines.push(footer); + lines +} + +pub(crate) fn clear_sequence(lines: usize) -> String { + if lines == 0 { + return String::new(); + } + let mut sequence = format!("\x1b[{}A\r", lines.saturating_sub(1)); + for index in 0..lines { + sequence.push_str("\x1b[2K"); + if index + 1 < lines { + sequence.push('\n'); + } + } + sequence.push('\r'); + if lines > 1 { + sequence.push_str(&format!("\x1b[{}A\r", lines - 1)); + } + sequence +} + +#[cfg(test)] +mod tests; diff --git a/crates/rb-progress/src/terminal/tests.rs b/crates/rb-progress/src/terminal/tests.rs new file mode 100644 index 0000000..fdbd9b5 --- /dev/null +++ b/crates/rb-progress/src/terminal/tests.rs @@ -0,0 +1,167 @@ +use super::*; + +#[test] +fn ticker_line_is_bounded() { + assert!(fit_terminal(&"x".repeat(500)).ends_with('…')); +} + +#[test] +fn clearing_returns_to_top_of_rendered_block() { + let sequence = clear_sequence(4); + assert!(sequence.ends_with("\x1b[3A\r")); + assert_eq!(sequence.matches("\x1b[2K").count(), 4); +} + +#[test] +fn short_durations_show_progress_before_one_second() { + assert_eq!(format_duration(Duration::from_millis(500)), "0.5s"); + assert_eq!(format_duration(Duration::from_secs(2)), "2.0s"); +} + +#[test] +fn viewport_keeps_active_root_and_footer_visible() { + let visible = limit_tree_viewport( + "┌─ processing documents".into(), + vec![ + "│ ✓ resolved".into(), + "│ ● worker 0".into(), + "│ ● worker 1".into(), + ], + "└─ overall: working".into(), + 4, + ); + assert_eq!( + visible, + [ + "┌─ processing documents", + "│ ● worker 0", + "│ ● worker 1", + "└─ overall: working" + ] + ); +} + +#[test] +fn viewport_drops_old_history_before_active_workers() { + let visible = limit_tree_viewport( + "┌─ processing documents".into(), + vec![ + "│ ✓ resolved".into(), + "│ ✓ downloaded".into(), + "│ ● worker 0".into(), + "│ output".into(), + ], + "└─ overall: working".into(), + 4, + ); + assert_eq!( + visible, + [ + "┌─ processing documents", + "│ ● worker 0", + "│ output", + "└─ overall: working" + ] + ); +} + +#[test] +fn viewport_always_leaves_one_line_for_footer() { + let visible = limit_tree_viewport( + "┌─ processing documents".into(), + vec!["│ ● worker 0".into()], + "└─ overall: working".into(), + 2, + ); + assert_eq!(visible, ["┌─ processing documents", "└─ overall: working"]); +} + +#[test] +fn viewport_preserves_all_rows_when_they_fit() { + let visible = limit_tree_viewport( + "┌─ processing documents".into(), + vec!["│ ✓ resolved".into(), "│ ● worker 0".into()], + "└─ overall: working".into(), + 4, + ); + assert_eq!( + visible, + [ + "┌─ processing documents", + "│ ✓ resolved", + "│ ● worker 0", + "└─ overall: working" + ] + ); +} + +#[test] +fn tree_viewport_keeps_header_and_footer() { + let lines = limit_tree_viewport( + "┌─ ● preparing demo bundle".into(), + (0..8) + .map(|index| format!("│ ├─ detail {index}")) + .collect(), + "└─ overall: working".into(), + 5, + ); + assert_eq!(lines.len(), 5); + assert!(lines[0].starts_with("┌─")); + assert!(lines[4].starts_with("└─")); +} + +#[test] +fn tree_viewport_collapses_body_to_working_marker() { + let lines = limit_tree_viewport( + "┌─ ● preparing demo bundle".into(), + (0..8) + .map(|index| format!("│ ├─ detail {index}")) + .collect(), + "└─ overall: working".into(), + 3, + ); + assert_eq!( + lines, + vec![ + "┌─ ● preparing demo bundle", + "│ · working…", + "└─ overall: working", + ] + ); +} + +#[test] +fn truncation_counts_terminal_columns() { + assert_eq!(fit_width("資料資料資料", 7), "資料資…"); + assert_eq!( + fit_width("e\u{301}e\u{301}e\u{301}", 3), + "e\u{301}e\u{301}e\u{301}" + ); + assert_eq!(fit_width("data", 0), ""); +} + +#[test] +fn task_text_cannot_move_the_cursor_or_create_extra_rows() { + assert_eq!( + fit_width("\x1b[2Jhello\nthere\r\x1b]0;title\x07!", 80), + "hello there !" + ); +} + +#[test] +fn duration_formatting_is_owned_by_the_renderer() { + assert_eq!(format_duration(Duration::from_millis(2260)), "2.3s"); + assert_eq!(format_duration(Duration::from_millis(12250)), "12s"); + assert_eq!(format_duration(Duration::from_millis(62250)), "1m02s"); +} + +#[test] +fn activity_icon_cycles_at_the_refresh_interval() { + let frames = (0..10) + .map(|tick| activity_icon(Duration::from_millis(tick * 250))) + .collect::(); + assert_eq!(frames, "⠋⠙⠹⠸⠼⠴⠦⠧⠇⠏"); + assert_eq!(activity_icon(Duration::from_millis(249)), '⠋'); + assert_eq!(activity_icon(Duration::from_millis(2500)), '⠋'); + assert!(frames.chars().all(|frame| frame.width() == Some(1))); +} diff --git a/crates/rb-progress/tests/pty_ui_test.py b/crates/rb-progress/tests/pty_ui_test.py new file mode 100644 index 0000000..f821d70 --- /dev/null +++ b/crates/rb-progress/tests/pty_ui_test.py @@ -0,0 +1,298 @@ +#!/usr/bin/env python3 +import codecs +import fcntl +import os +import pty +import select +import struct +import subprocess +import termios +import unicodedata +import sys +import time + + +class Screen: + def __init__(self, rows, columns): + self.rows = rows + self.columns = columns + self.cells = [[' '] * columns for _ in range(rows)] + self.row = 0 + self.column = 0 + self.pending = b'' + self.decoder = codecs.getincrementaldecoder('utf-8')() + + def newline(self): + self.row += 1 + self.column = 0 + if self.row >= self.rows: + self.cells.pop(0) + self.cells.append([' '] * self.columns) + self.row = self.rows - 1 + + def write(self, value): + for char in self.decoder.decode(value): + if char == '\r': + self.column = 0 + elif char == '\n': + self.newline() + elif char.isprintable(): + width = 0 if unicodedata.combining(char) else 2 if unicodedata.east_asian_width(char) in ('W', 'F') else 1 + if width == 0: + previous = self.column - 1 + while previous > 0 and self.cells[self.row][previous] == '': + previous -= 1 + if previous >= 0: + self.cells[self.row][previous] += char + continue + if self.column + width > self.columns: + self.newline() + self.cells[self.row][self.column] = char + if width == 2: + self.cells[self.row][self.column + 1] = '' + self.column += width + + def csi(self, command, params): + values = [int(value) if value else 0 for value in params.split(';')] if params else [] + if command == 'A': + self.row = max(0, self.row - (values[0] if values and values[0] else 1)) + elif command == 'B': + self.row = min(self.rows - 1, self.row + (values[0] if values and values[0] else 1)) + elif command == 'K': + self.cells[self.row] = [' '] * self.columns + elif command == 'L': + count = values[0] if values and values[0] else 1 + for _ in range(count): + self.cells.insert(self.row, [' '] * self.columns) + self.cells.pop() + elif command == 'M': + count = values[0] if values and values[0] else 1 + for _ in range(count): + self.cells.pop(self.row) + self.cells.append([' '] * self.columns) + elif command == 'J' and values and values[0] == 2: + self.cells = [[' '] * self.columns for _ in range(self.rows)] + + def feed(self, data): + data = self.pending + data + self.pending = b'' + index = 0 + plain = bytearray() + while index < len(data): + if data[index:index + 2] == b'\x1b[': + if plain: + self.write(bytes(plain)) + plain.clear() + end = index + 2 + while end < len(data) and not (0x40 <= data[end] <= 0x7e): + end += 1 + if end == len(data): + self.pending = data[index:] + break + params = data[index + 2:end].decode('ascii', errors='ignore') + self.csi(chr(data[end]), params) + index = end + 1 + elif data[index] == 0x1b: + if index + 1 == len(data): + self.pending = data[index:] + break + index += 2 + else: + plain.append(data[index]) + index += 1 + if plain: + self.write(bytes(plain)) + + def line(self, row): + return ''.join(self.cells[row]).rstrip() + + def has(self, text): + return any(self.line(row).startswith(text) for row in range(self.rows)) + + def find(self, text): + for row in range(self.rows): + if self.line(row).startswith(text): + return row + return None + + +def capture(binary, args, rows, columns, prompt_row, observe=None): + master, slave = pty.openpty() + fcntl.ioctl(slave, termios.TIOCSWINSZ, struct.pack('HHHH', rows, columns, 0, 0)) + env = os.environ.copy() + env['COLUMNS'] = str(columns) + env['LINES'] = str(rows) + process = subprocess.Popen([binary] + args, stdin=slave, stdout=slave, stderr=slave, env=env) + os.close(slave) + screen = Screen(rows, columns) + for _ in range(prompt_row): + screen.newline() + screen.write(b'[shell]$ rb --db sync') + screen.newline() + output = bytearray() + progress_started = False + deadline = time.monotonic() + 20 + try: + while True: + remaining = deadline - time.monotonic() + if remaining <= 0 or not select.select([master], [], [], remaining)[0]: + raise AssertionError('progress probe timed out') + try: + chunk = os.read(master, 4096) + except OSError: + break + if not chunk: + break + output.extend(chunk) + screen.feed(chunk) + if observe is not None: + observe(screen) + if not screen.has('[shell]$ rb --db sync'): + raise AssertionError('prompt was erased or scrolled during rendering') + if not progress_started and b'resolved' in output: + prompt = screen.find('[shell]$ rb --db sync') + if not screen.line(prompt + 1).strip(): + raise AssertionError('renderer left an empty row below the prompt') + progress_started = True + if process.wait(timeout=5) != 0: + raise AssertionError('progress probe exited unsuccessfully') + return screen, output + except Exception as error: + raise AssertionError(f'{args} at {columns}x{rows}: {error}\n' + '\n'.join(screen.line(row) for row in range(rows))) from error + finally: + if process.poll() is None: + process.kill() + process.wait(timeout=5) + os.close(master) + + +def run(binary, tasks=False, handoff_output=False, stress=False, failure=False, rows=24, columns=80, prompt_row=8): + args = ['--failure'] if failure else ['--stress'] if stress else ['--tasks'] if tasks else ['--handoff-output'] if handoff_output else [] + screen, output = capture(binary, args, rows, columns, prompt_row) + if not screen.has('[shell]$ rb --db sync'): + raise SystemExit(f'prompt was erased or scrolled\n{screen.line(1)!r}') + if b'\x1b[2J' in output: + raise SystemExit('renderer cleared the entire terminal') + if b'probe complete' not in output: + raise SystemExit('probe did not complete') + if tasks or stress or failure: + title = screen.find('┌─ × processing documents' if failure else '┌─ ✓ processing documents') + summary = screen.find('└─ × processing failed' if failure else '└─ ✓ documents ready') + prompt = screen.find('[shell]$ rb --db sync') + if title != prompt + 1 or summary != title + (2 if failure else 1): + raise SystemExit('task-only reporter left duplicate or blank rows') + if failure and not any('document failure' in screen.line(row) for row in range(rows)): + raise SystemExit('final report lost the task error') + if screen.find('probe complete') != summary + 1: + raise SystemExit('task-only final output did not start on the next line') + titles = sum('processing documents' in screen.line(row) for row in range(screen.rows)) + if titles != 1 or any(screen.line(row) for row in range(summary + 2, screen.rows)): + raise SystemExit('task-only completion left duplicated or stale rows') + required = [] if stress or failure else [b'preparing content', b'exporting documents', b'reading document', b'writing document'] + for text in required: + if text not in output: + raise SystemExit(f'task-only reporter lost {text!r}') + if stress and columns >= 80 and '資料'.encode() not in output: + raise SystemExit('stress probe lost Unicode output') + print(f'PASS PTY generic tasks {columns}x{rows} prompt={prompt_row} failure={failure}') + return + completed_row = screen.find('├─ ✓ preparing pty-probe bundle') + if completed_row is None: + raise SystemExit('final bundle row was not completed') + if sum(line.startswith('├─ ✓ preparing pty-probe bundle') for line in (screen.line(row) for row in range(screen.rows))) != 1: + raise SystemExit('final bundle row was duplicated') + if screen.has('├─ ● preparing pty-probe bundle'): + raise SystemExit('final bundle row remained active') + if screen.find('├─ ✓ preparing native extensions') is None: + raise SystemExit('native handoff row was lost') + resolved_row = screen.find('┌─ ● preparing pty-probe bundle') + if resolved_row is None: + raise SystemExit('top tree entry was not preserved') + summary_row = screen.find('└─ ✓ your meticulously prepared bundle is ready, sir') + if summary_row != completed_row + 1: + raise SystemExit('final bundle summary was not placed below the tree') + probe_row = screen.find('probe complete') + if probe_row != summary_row + 1: + raise SystemExit('final output was not placed directly below the completed tree') + if handoff_output and any('● preparing native extensions' in screen.line(row) or 'preview ' in screen.line(row) for row in range(screen.rows)): + raise SystemExit('handoff left active rows or worker previews in the completed tree') + print('PASS PTY prompt and viewport preservation') + + +def verify_aggregate_completion(binary): + for rows in [4, 8]: + screen, output = capture(binary, ['--aggregate'], rows, 40, 0) + assert screen.find('[shell]$ rb --db sync') == 0 + assert screen.find('┌─ ✓ processing documents') == 1 + summary = screen.find('└─ ✓ documents ready') + assert summary is not None and summary < rows - 1 + assert screen.row == summary + 1 and screen.column == 0 + assert not any(screen.line(row) for row in range(summary + 1, rows)) + assert b'\x1b[2J' not in output + print(f'PASS PTY aggregate completion 40x{rows}') + + +def verify_standalone_parallel(binary): + plain = subprocess.run([binary], capture_output=True, timeout=20) + assert plain.returncode == 0 + assert b'completed 3 | failed 0' in plain.stderr + assert b'\x1b' not in plain.stderr and b'\r' not in plain.stderr + for worker in range(3): + assert f'[worker: {worker}] complete'.encode() in plain.stderr + for rows, columns in [(24, 80), (8, 40), (4, 40)]: + simultaneous = [] + def observe(screen): + frame = '\n'.join(screen.line(row) for row in range(rows)) + if all(f'worker {worker}' in frame for worker in range(3)): + simultaneous.append(frame) + screen, output = capture(binary, [], rows, columns, 0, observe) + if rows == 24: + assert simultaneous, 'parallel workers were never visible together' + assert screen.find('[shell]$ rb --db sync') == 0 + assert screen.find('┌─ ✓ Preparing the dining room') == 1 + summary = screen.find('└─ ✓ The dining room is ready, sir') + assert summary is not None + assert screen.row == summary + 1 and screen.column == 0 + assert all(screen.line(row) for row in range(summary + 1)) + assert not any(screen.line(row) for row in range(summary + 1, rows)) + assert sum('Preparing the dining room' in screen.line(row) for row in range(rows)) == 1 + assert b'\x1b[2J' not in output + print(f'PASS PTY standalone parallel {columns}x{rows}') + + +def verify_parser(): + data = 'prompt\r\n資料 e\u0301\x1b[1A\r\x1b[2Ktitle'.encode() + expected = Screen(8, 40) + expected.feed(data) + assert expected.line(0) == 'title' + assert expected.line(1) == '資料 e\u0301' + assert expected.cells[1][2] == '料' + assert expected.cells[1][5] == 'e\u0301' + for size in [1, 2, 3, 5, 7]: + fragmented = Screen(8, 40) + for start in range(0, len(data), size): + fragmented.feed(data[start:start + size]) + assert fragmented.cells == expected.cells, f'parser lost bytes at chunk size {size}' + assert (fragmented.row, fragmented.column) == (expected.row, expected.column) + + +if __name__ == '__main__': + verify_parser() + if len(sys.argv) != 3: + raise SystemExit('usage: pty_ui_test.py PATH_TO_PROBE PATH_TO_PARALLEL') + plain = subprocess.run([sys.argv[1], '--failure'], capture_output=True, timeout=20) + assert plain.returncode == 0 + assert b'[worker: ' in plain.stderr and b'failed: document failure' in plain.stderr + assert b'\x1b' not in plain.stderr and b'\r' not in plain.stderr + print('PASS redirected output uses worker-prefixed plain reporting') + verify_standalone_parallel(sys.argv[2]) + verify_aggregate_completion(sys.argv[1]) + run(sys.argv[1]) + run(sys.argv[1], tasks=True) + run(sys.argv[1], handoff_output=True) + for rows in [8, 12, 24, 40]: + for columns in [40, 80, 120]: + for prompt in [0, rows // 2]: + run(sys.argv[1], stress=True, rows=rows, columns=columns, prompt_row=prompt) + run(sys.argv[1], failure=True) + run(sys.argv[1], failure=True, rows=8, columns=40, prompt_row=0)