diff --git a/Cargo.lock b/Cargo.lock index b4498b453..738b1edbe 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -516,6 +516,16 @@ version = "1.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3a822ea5bc7590f9d40f1ba12c0dc3c2760f3482c6984db1573ad11031420831" +[[package]] +name = "cli-table" +version = "0.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "14da8d951cef7cc4f13ccc9b744d736963d57863c7e6fc33c070ea274546082c" +dependencies = [ + "termcolor", + "unicode-width 0.2.2", +] + [[package]] name = "cmake" version = "0.1.57" @@ -1277,6 +1287,7 @@ dependencies = [ "dirs", "fabro-agent", "fabro-mcp", + "fabro-util", "fabro-workflows", "serde", "tempfile", @@ -1380,6 +1391,7 @@ dependencies = [ "base64", "bytes", "clap", + "cli-table", "dialoguer", "dotenvy", "fabro-util", @@ -1543,6 +1555,7 @@ dependencies = [ "base64", "chrono", "clap", + "cli-table", "console 0.15.11", "daytona-api-client", "daytona-sdk", @@ -4875,6 +4888,15 @@ dependencies = [ "utf-8", ] +[[package]] +name = "termcolor" +version = "1.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "06794f8f6c5c898b3275aebefa6b8a1cb24cd2c6c79397ab15774837a0bc5755" +dependencies = [ + "winapi-util", +] + [[package]] name = "termimad" version = "0.34.1" diff --git a/Cargo.toml b/Cargo.toml index b52b6ead0..8ddf32162 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -30,6 +30,7 @@ jsonschema = "0.42" chrono = { version = "0.4", features = ["clock"] } bollard = "0.18" tar = "0.4" +cli-table = { version = "0.5", default-features = false } console = "0.15" dialoguer = "0.12" git2 = "0.20" @@ -63,6 +64,9 @@ daytona-api-client = { git = "https://github.com/brynary/daytona-sdk-rust", rev lto = "thin" strip = true +[profile.dev.package."*"] +debug = false # Disable debug info for all dependencies + # regex is extremely slow in debug builds (~10s to compile gitleaks patterns) [profile.dev.package.regex] opt-level = 2 diff --git a/docs-internal/events-strategy.md b/docs-internal/events-strategy.md index d21a7a071..b57a592f8 100644 --- a/docs-internal/events-strategy.md +++ b/docs-internal/events-strategy.md @@ -184,9 +184,15 @@ WorkflowRunEvent::MyNewEvent { node_id, duration_ms, .. } => { | Event | JSONL fields | |---|---| -| `CheckpointSaved` | `node_id`, `node_label` | -| `GitCheckpoint` | `run_id`, `node_id`, `node_label`, `status`, `git_commit_sha` | -| `GitCheckpointFailed` | `node_id`, `node_label`, `error` | +| `CheckpointCompleted` | `node_id`, `node_label`, `status`, `git_commit_sha` (optional) | +| `CheckpointFailed` | `node_id`, `node_label`, `error` | +| `GitCommit` | `node_id` (optional), `node_label` (optional), `sha` | +| `GitPush` | `branch`, `success` | +| `GitBranch` | `branch`, `sha` | +| `GitWorktreeAdd` | `path`, `branch` | +| `GitWorktreeRemove` | `path` | +| `GitFetch` | `branch`, `success` | +| `GitReset` | `sha` | ### Human interaction @@ -296,5 +302,5 @@ Error information is stored as plain strings. The `error` field contains the hum | `cli/run.rs` non-verbose listener | `name`, `duration_ms`, `status`, `usage` from `StageCompleted/Failed` | CLI progress output | | `cli/mod.rs` `format_event_summary()` | All events | `-v` verbose output | | `cli/run.rs` cost accumulator | `usage` from `StageCompleted` | Total cost tracking | -| `cli/run.rs` git SHA tracker | `git_commit_sha` from `GitCheckpoint` | Final SHA for `conclusion.json` | +| `cli/run.rs` git SHA tracker | `git_commit_sha` from `CheckpointCompleted` | Final SHA for `conclusion.json` | | External tooling | `progress.jsonl` | Live monitoring, dashboards | diff --git a/docs/execution/observability.mdx b/docs/execution/observability.mdx index 8cf9713d4..edb7679b0 100644 --- a/docs/execution/observability.mdx +++ b/docs/execution/observability.mdx @@ -58,8 +58,14 @@ Events fall into several categories: |---|---|---| | `EdgeSelected` | `from_node`, `to_node`, `label`, `condition` | Transition between nodes | | `LoopRestart` | `from_node`, `to_node` | Loop restart edge taken | -| `CheckpointSaved` | `node_id` | Checkpoint written to disk | -| `GitCheckpoint` | `node_id`, `git_commit_sha` | Checkpoint committed to Git | +| `CheckpointCompleted` | `node_id`, `git_commit_sha` (optional) | Checkpoint saved (with git SHA when git is enabled) | +| `GitCommit` | `node_id`, `sha` | Git commit created | +| `GitPush` | `branch`, `success` | Git push attempted | +| `GitBranch` | `branch`, `sha` | Git branch created | +| `GitWorktreeAdd` | `path`, `branch` | Git worktree added | +| `GitWorktreeRemove` | `path` | Git worktree removed | +| `GitFetch` | `branch`, `success` | Git fetch attempted | +| `GitReset` | `sha` | Git reset executed | | `Failover` | `stage`, `from_provider`, `to_provider`, `error` | LLM provider failover | **Parallel execution:** diff --git a/fabro.toml b/fabro.toml index 577d2b097..3daedcc4e 100644 --- a/fabro.toml +++ b/fabro.toml @@ -14,7 +14,7 @@ auto_stop_interval = 30 repo = "fabro-sh/fabro" [sandbox.daytona.snapshot] -name = "fabro-v3" +name = "fabro-v5" cpu = 4 memory = 8 disk = 20 @@ -25,9 +25,18 @@ RUN apt-get update && apt-get install -y --no-install-recommends \ curl git ca-certificates build-essential pkg-config libssl-dev unzip python3 \ && rm -rf /var/lib/apt/lists/* +# GitHub CLI +RUN curl -fsSL https://cli.github.com/packages/githubcli-archive-keyring.gpg \ + | dd of=/usr/share/keyrings/githubcli-archive-keyring.gpg \ + && echo "deb [arch=$(dpkg --print-architecture) signed-by=/usr/share/keyrings/githubcli-archive-keyring.gpg] https://cli.github.com/packages stable main" \ + | tee /etc/apt/sources.list.d/github-cli.list > /dev/null \ + && apt-get update && apt-get install -y --no-install-recommends gh \ + && rm -rf /var/lib/apt/lists/* + # Rust RUN curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y ENV PATH="/root/.cargo/bin:${PATH}" +ENV CARGO_INCREMENTAL=0 # Bun RUN curl -fsSL https://bun.sh/install | bash diff --git a/lib/crates/fabro-cli/src/main.rs b/lib/crates/fabro-cli/src/main.rs index 2411ead71..5b49cb0e9 100644 --- a/lib/crates/fabro-cli/src/main.rs +++ b/lib/crates/fabro-cli/src/main.rs @@ -110,6 +110,8 @@ enum Command { /// List workflow runs #[command(hide = true)] Ps(fabro_workflows::cli::runs::RunsListArgs), + /// Remove one or more workflow runs + Rm(fabro_workflows::cli::runs::RunsRemoveArgs), /// Pull request operations Pr { #[command(subcommand)] @@ -218,6 +220,13 @@ fn detach_run(args: fabro_workflows::cli::RunArgs) -> Result<()> { )) }); std::fs::create_dir_all(&run_dir)?; + std::fs::write(run_dir.join("id.txt"), &run_id)?; + fabro_workflows::cli::runs::write_run_status( + &run_dir, + fabro_workflows::run_status::RunStatus::Submitted, + None, + ); + std::fs::File::create(run_dir.join("progress.jsonl"))?; let log_file = std::fs::File::create(run_dir.join("detach.log"))?; @@ -391,6 +400,7 @@ async fn main_inner() -> (String, Result<()>) { Command::Init => "init", Command::Install => "install", Command::Ps(_) => "ps", + Command::Rm(_) => "rm", Command::Pr { command } => match command { PrCommand::Create(_) => "pr create", PrCommand::List(_) => "pr list", @@ -685,7 +695,11 @@ async fn main_inner() -> (String, Result<()>) { install::run_install().await?; } Command::Ps(args) => { - fabro_workflows::cli::runs::list_command(&args)?; + let styles = fabro_util::terminal::Styles::detect_stdout(); + fabro_workflows::cli::runs::list_command(&args, &styles)?; + } + Command::Rm(args) => { + fabro_workflows::cli::runs::remove_command(&args).await?; } Command::Pr { command } => { let cli_config = cli_config::load_cli_config(None)?; diff --git a/lib/crates/fabro-cli/tests/cmd/model/bare.trycmd b/lib/crates/fabro-cli/tests/cmd/model/bare.trycmd index 033bd1986..9343f250a 100644 --- a/lib/crates/fabro-cli/tests/cmd/model/bare.trycmd +++ b/lib/crates/fabro-cli/tests/cmd/model/bare.trycmd @@ -1,24 +1,24 @@ ```console $ fabro model -MODEL PROVIDER ALIASES CONTEXT COST SPEED -claude-opus-4-6 anthropic opus, claude-opus 1m $15.0 / $75.0 25 tok/s -claude-sonnet-4-5 anthropic 200k $3.0 / $15.0 50 tok/s -claude-sonnet-4-6 anthropic sonnet, claude-sonnet 200k $3.0 / $15.0 50 tok/s -claude-haiku-4-5 anthropic haiku, claude-haiku 200k $0.8 / $4.0 100 tok/s -gpt-5.2 openai gpt5 1m $1.8 / $14.0 65 tok/s -gpt-5-mini openai gpt5-mini 1m $0.2 / $2.0 70 tok/s -gpt-5.2-codex openai 1m $1.8 / $14.0 100 tok/s -gpt-5.3-codex openai codex 1m $1.8 / $14.0 100 tok/s -gpt-5.3-codex-spark openai codex-spark 131k - / - 1000 tok/s -gpt-5.4 openai gpt54 1m $2.5 / $15.0 70 tok/s -gpt-5.4-pro openai gpt54-pro 1m $30.0 / $180.0 20 tok/s -gemini-3.1-pro-preview gemini gemini-pro 1m $2.0 / $12.0 85 tok/s -gemini-3.1-pro-preview-customtools gemini gemini-customtools 1m $2.0 / $12.0 85 tok/s -gemini-3-flash-preview gemini gemini-flash 1m $0.5 / $3.0 150 tok/s -gemini-3.1-flash-lite-preview gemini gemini-flash-lite 1m $0.2 / $1.5 200 tok/s -kimi-k2.5 kimi kimi 262k $0.6 / $3.0 50 tok/s -glm-4.7 zai glm, glm4 203k $0.6 / $2.2 100 tok/s -minimax-m2.5 minimax minimax 197k $0.3 / $1.2 45 tok/s -mercury-2 inception mercury 131k $0.2 / $0.8 1000 tok/s - + MODEL   PROVIDER   ALIASES   CONTEXT   COST   SPEED  + claude-opus-4-6   anthropic  opus, claude-opus    1m   $15.0 / $75.0   25 tok/s  + claude-sonnet-4-5   anthropic      200k   $3.0 / $15.0   50 tok/s  + claude-sonnet-4-6   anthropic  sonnet, claude-sonnet   200k   $3.0 / $15.0   50 tok/s  + claude-haiku-4-5   anthropic  haiku, claude-haiku    200k   $0.8 / $4.0   100 tok/s  + gpt-5.2   openai   gpt5    1m   $1.8 / $14.0   65 tok/s  + gpt-5-mini   openai   gpt5-mini    1m   $0.2 / $2.0   70 tok/s  + gpt-5.2-codex   openai       1m   $1.8 / $14.0   100 tok/s  + gpt-5.3-codex   openai   codex    1m   $1.8 / $14.0   100 tok/s  + gpt-5.3-codex-spark   openai   codex-spark    131k   - / -  1000 tok/s  + gpt-5.4   openai   gpt54    1m   $2.5 / $15.0   70 tok/s  + gpt-5.4-pro   openai   gpt54-pro    1m  $30.0 / $180.0   20 tok/s  + gemini-3.1-pro-preview   gemini   gemini-pro    1m   $2.0 / $12.0   85 tok/s  + gemini-3.1-pro-preview-customtools  gemini   gemini-customtools    1m   $2.0 / $12.0   85 tok/s  + gemini-3-flash-preview   gemini   gemini-flash    1m   $0.5 / $3.0   150 tok/s  + gemini-3.1-flash-lite-preview   gemini   gemini-flash-lite    1m   $0.2 / $1.5   200 tok/s  + kimi-k2.5   kimi   kimi    262k   $0.6 / $3.0   50 tok/s  + glm-4.7   zai   glm, glm4    203k   $0.6 / $2.2   100 tok/s  + minimax-m2.5   minimax   minimax    197k   $0.3 / $1.2   45 tok/s  + mercury-2   inception  mercury    131k   $0.2 / $0.8  1000 tok/s  + ``` diff --git a/lib/crates/fabro-cli/tests/cmd/model/list-provider.trycmd b/lib/crates/fabro-cli/tests/cmd/model/list-provider.trycmd index 7e884a0bc..83d18371b 100644 --- a/lib/crates/fabro-cli/tests/cmd/model/list-provider.trycmd +++ b/lib/crates/fabro-cli/tests/cmd/model/list-provider.trycmd @@ -1,9 +1,9 @@ ```console $ fabro model list --provider anthropic -MODEL PROVIDER ALIASES CONTEXT COST SPEED -claude-opus-4-6 anthropic opus, claude-opus 1m $15.0 / $75.0 25 tok/s -claude-sonnet-4-5 anthropic 200k $3.0 / $15.0 50 tok/s -claude-sonnet-4-6 anthropic sonnet, claude-sonnet 200k $3.0 / $15.0 50 tok/s -claude-haiku-4-5 anthropic haiku, claude-haiku 200k $0.8 / $4.0 100 tok/s - + MODEL   PROVIDER   ALIASES   CONTEXT   COST   SPEED  + claude-opus-4-6   anthropic  opus, claude-opus    1m  $15.0 / $75.0   25 tok/s  + claude-sonnet-4-5  anthropic      200k   $3.0 / $15.0   50 tok/s  + claude-sonnet-4-6  anthropic  sonnet, claude-sonnet   200k   $3.0 / $15.0   50 tok/s  + claude-haiku-4-5   anthropic  haiku, claude-haiku    200k   $0.8 / $4.0  100 tok/s  + ``` diff --git a/lib/crates/fabro-cli/tests/cmd/model/list-query-aliases.trycmd b/lib/crates/fabro-cli/tests/cmd/model/list-query-aliases.trycmd index 2a5905dd3..6a6ca5c05 100644 --- a/lib/crates/fabro-cli/tests/cmd/model/list-query-aliases.trycmd +++ b/lib/crates/fabro-cli/tests/cmd/model/list-query-aliases.trycmd @@ -1,8 +1,8 @@ ```console $ fabro model list --query codex -MODEL PROVIDER ALIASES CONTEXT COST SPEED -gpt-5.2-codex openai 1m $1.8 / $14.0 100 tok/s -gpt-5.3-codex openai codex 1m $1.8 / $14.0 100 tok/s -gpt-5.3-codex-spark openai codex-spark 131k - / - 1000 tok/s - + MODEL   PROVIDER  ALIASES   CONTEXT   COST   SPEED  + gpt-5.2-codex   openai       1m  $1.8 / $14.0   100 tok/s  + gpt-5.3-codex   openai   codex    1m  $1.8 / $14.0   100 tok/s  + gpt-5.3-codex-spark  openai   codex-spark   131k   - / -  1000 tok/s  + ``` diff --git a/lib/crates/fabro-cli/tests/cmd/model/list-query-case-insensitive.trycmd b/lib/crates/fabro-cli/tests/cmd/model/list-query-case-insensitive.trycmd index b09ce9e71..b7b4381ee 100644 --- a/lib/crates/fabro-cli/tests/cmd/model/list-query-case-insensitive.trycmd +++ b/lib/crates/fabro-cli/tests/cmd/model/list-query-case-insensitive.trycmd @@ -1,6 +1,6 @@ ```console $ fabro model list --query OPUS -MODEL PROVIDER ALIASES CONTEXT COST SPEED -claude-opus-4-6 anthropic opus, claude-opus 1m $15.0 / $75.0 25 tok/s - + MODEL   PROVIDER   ALIASES   CONTEXT   COST   SPEED  + claude-opus-4-6  anthropic  opus, claude-opus   1m  $15.0 / $75.0  25 tok/s  + ``` diff --git a/lib/crates/fabro-cli/tests/cmd/model/list-query.trycmd b/lib/crates/fabro-cli/tests/cmd/model/list-query.trycmd index fd1f45fbb..e3ae96a60 100644 --- a/lib/crates/fabro-cli/tests/cmd/model/list-query.trycmd +++ b/lib/crates/fabro-cli/tests/cmd/model/list-query.trycmd @@ -1,6 +1,6 @@ ```console $ fabro model list --query opus -MODEL PROVIDER ALIASES CONTEXT COST SPEED -claude-opus-4-6 anthropic opus, claude-opus 1m $15.0 / $75.0 25 tok/s - + MODEL   PROVIDER   ALIASES   CONTEXT   COST   SPEED  + claude-opus-4-6  anthropic  opus, claude-opus   1m  $15.0 / $75.0  25 tok/s  + ``` diff --git a/lib/crates/fabro-cli/tests/cmd/model/list.trycmd b/lib/crates/fabro-cli/tests/cmd/model/list.trycmd index 19f79b975..1447601bf 100644 --- a/lib/crates/fabro-cli/tests/cmd/model/list.trycmd +++ b/lib/crates/fabro-cli/tests/cmd/model/list.trycmd @@ -1,24 +1,24 @@ ```console $ fabro model list -MODEL PROVIDER ALIASES CONTEXT COST SPEED -claude-opus-4-6 anthropic opus, claude-opus 1m $15.0 / $75.0 25 tok/s -claude-sonnet-4-5 anthropic 200k $3.0 / $15.0 50 tok/s -claude-sonnet-4-6 anthropic sonnet, claude-sonnet 200k $3.0 / $15.0 50 tok/s -claude-haiku-4-5 anthropic haiku, claude-haiku 200k $0.8 / $4.0 100 tok/s -gpt-5.2 openai gpt5 1m $1.8 / $14.0 65 tok/s -gpt-5-mini openai gpt5-mini 1m $0.2 / $2.0 70 tok/s -gpt-5.2-codex openai 1m $1.8 / $14.0 100 tok/s -gpt-5.3-codex openai codex 1m $1.8 / $14.0 100 tok/s -gpt-5.3-codex-spark openai codex-spark 131k - / - 1000 tok/s -gpt-5.4 openai gpt54 1m $2.5 / $15.0 70 tok/s -gpt-5.4-pro openai gpt54-pro 1m $30.0 / $180.0 20 tok/s -gemini-3.1-pro-preview gemini gemini-pro 1m $2.0 / $12.0 85 tok/s -gemini-3.1-pro-preview-customtools gemini gemini-customtools 1m $2.0 / $12.0 85 tok/s -gemini-3-flash-preview gemini gemini-flash 1m $0.5 / $3.0 150 tok/s -gemini-3.1-flash-lite-preview gemini gemini-flash-lite 1m $0.2 / $1.5 200 tok/s -kimi-k2.5 kimi kimi 262k $0.6 / $3.0 50 tok/s -glm-4.7 zai glm, glm4 203k $0.6 / $2.2 100 tok/s -minimax-m2.5 minimax minimax 197k $0.3 / $1.2 45 tok/s -mercury-2 inception mercury 131k $0.2 / $0.8 1000 tok/s - + MODEL   PROVIDER   ALIASES   CONTEXT   COST   SPEED  + claude-opus-4-6   anthropic  opus, claude-opus    1m   $15.0 / $75.0   25 tok/s  + claude-sonnet-4-5   anthropic      200k   $3.0 / $15.0   50 tok/s  + claude-sonnet-4-6   anthropic  sonnet, claude-sonnet   200k   $3.0 / $15.0   50 tok/s  + claude-haiku-4-5   anthropic  haiku, claude-haiku    200k   $0.8 / $4.0   100 tok/s  + gpt-5.2   openai   gpt5    1m   $1.8 / $14.0   65 tok/s  + gpt-5-mini   openai   gpt5-mini    1m   $0.2 / $2.0   70 tok/s  + gpt-5.2-codex   openai       1m   $1.8 / $14.0   100 tok/s  + gpt-5.3-codex   openai   codex    1m   $1.8 / $14.0   100 tok/s  + gpt-5.3-codex-spark   openai   codex-spark    131k   - / -  1000 tok/s  + gpt-5.4   openai   gpt54    1m   $2.5 / $15.0   70 tok/s  + gpt-5.4-pro   openai   gpt54-pro    1m  $30.0 / $180.0   20 tok/s  + gemini-3.1-pro-preview   gemini   gemini-pro    1m   $2.0 / $12.0   85 tok/s  + gemini-3.1-pro-preview-customtools  gemini   gemini-customtools    1m   $2.0 / $12.0   85 tok/s  + gemini-3-flash-preview   gemini   gemini-flash    1m   $0.5 / $3.0   150 tok/s  + gemini-3.1-flash-lite-preview   gemini   gemini-flash-lite    1m   $0.2 / $1.5   200 tok/s  + kimi-k2.5   kimi   kimi    262k   $0.6 / $3.0   50 tok/s  + glm-4.7   zai   glm, glm4    203k   $0.6 / $2.2   100 tok/s  + minimax-m2.5   minimax   minimax    197k   $0.3 / $1.2   45 tok/s  + mercury-2   inception  mercury    131k   $0.2 / $0.8  1000 tok/s  + ``` diff --git a/lib/crates/fabro-config/Cargo.toml b/lib/crates/fabro-config/Cargo.toml index 2cdd81e53..3debef3f4 100644 --- a/lib/crates/fabro-config/Cargo.toml +++ b/lib/crates/fabro-config/Cargo.toml @@ -17,6 +17,7 @@ anyhow.workspace = true fabro-agent = { path = "../fabro-agent" } fabro-mcp = { path = "../fabro-mcp" } fabro-workflows = { path = "../fabro-workflows" } +fabro-util = { path = "../fabro-util" } dirs.workspace = true serde.workspace = true toml.workspace = true diff --git a/lib/crates/fabro-config/src/lib.rs b/lib/crates/fabro-config/src/lib.rs index 67803db70..e2e08619f 100644 --- a/lib/crates/fabro-config/src/lib.rs +++ b/lib/crates/fabro-config/src/lib.rs @@ -2,31 +2,4 @@ pub mod cli; pub mod project; pub mod server; -use std::path::{Path, PathBuf}; - -/// Expand `~/` prefix to the user's home directory. -pub fn expand_tilde(path: &Path) -> PathBuf { - if let Ok(rest) = path.strip_prefix("~") { - if let Some(home) = dirs::home_dir() { - return home.join(rest); - } - } - path.to_path_buf() -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn expand_tilde_with_home_prefix() { - let result = expand_tilde(Path::new("~/foo/bar")); - assert!(result != Path::new("~/foo/bar")); - assert!(result.ends_with("foo/bar")); - } - - #[test] - fn expand_tilde_without_prefix() { - assert_eq!(expand_tilde(Path::new("/abs/path")), Path::new("/abs/path")); - } -} +pub use fabro_util::path::expand_tilde; diff --git a/lib/crates/fabro-llm/Cargo.toml b/lib/crates/fabro-llm/Cargo.toml index 9d6e14afc..cd8e70bc3 100644 --- a/lib/crates/fabro-llm/Cargo.toml +++ b/lib/crates/fabro-llm/Cargo.toml @@ -27,6 +27,7 @@ reqwest.workspace = true base64.workspace = true bytes.workspace = true tokio-util.workspace = true +cli-table.workspace = true clap.workspace = true dialoguer.workspace = true tracing.workspace = true diff --git a/lib/crates/fabro-llm/src/cli.rs b/lib/crates/fabro-llm/src/cli.rs index d92a62a7d..d9afae73a 100644 --- a/lib/crates/fabro-llm/src/cli.rs +++ b/lib/crates/fabro-llm/src/cli.rs @@ -6,6 +6,8 @@ use std::time::Duration; use anyhow::{bail, Context, Result}; use clap::{Args, Subcommand}; +use cli_table::format::{Border, Justify, Separator}; +use cli_table::{print_stdout, Cell, CellStruct, Color, Style, Table}; use futures::StreamExt; use serde::Deserialize; @@ -107,30 +109,64 @@ fn format_speed(tps: Option) -> String { } } -fn print_models_table(models: &[crate::types::ModelInfo], s: &Styles) { - println!( - "{}", - s.bold_dim.apply_to(format!( - "{:<30} {:<12} {:<24} {:>10} {:>7} {:>7} {:>10}", - "MODEL", "PROVIDER", "ALIASES", "CONTEXT", "COST", "", "SPEED", - )), - ); - for model in models { - let aliases = model.aliases.join(", "); - println!( - "{} {} {} {:>10} {:>7} / {:<7} {}", - s.bold.apply_to(format!("{:<30}", model.id)), - s.dim.apply_to(format!("{:<12}", model.provider)), - s.dim.apply_to(format!("{:<24}", aliases)), - format_context_window(model.limits.context_window), - format_cost(model.costs.input_cost_per_mtok), - format_cost(model.costs.output_cost_per_mtok), - s.cyan - .apply_to(format!("{:>10}", format_speed(model.estimated_output_tps))), - ); +fn color_if(use_color: bool, color: Color) -> Option { + if use_color { + Some(color) + } else { + None } } +fn model_row(model: &crate::types::ModelInfo, use_color: bool) -> Vec { + let aliases = model.aliases.join(", "); + let cost = format!( + "{} / {}", + format_cost(model.costs.input_cost_per_mtok), + format_cost(model.costs.output_cost_per_mtok), + ); + vec![ + model.id.clone().cell().bold(use_color), + model + .provider + .clone() + .cell() + .foreground_color(color_if(use_color, Color::Ansi256(8))), + aliases + .cell() + .foreground_color(color_if(use_color, Color::Ansi256(8))), + format_context_window(model.limits.context_window) + .cell() + .justify(Justify::Right), + cost.cell().justify(Justify::Right), + format_speed(model.estimated_output_tps) + .cell() + .justify(Justify::Right) + .foreground_color(color_if(use_color, Color::Cyan)), + ] +} + +fn models_title() -> Vec { + vec![ + "MODEL".cell().bold(true), + "PROVIDER".cell().bold(true), + "ALIASES".cell().bold(true), + "CONTEXT".cell().bold(true).justify(Justify::Right), + "COST".cell().bold(true).justify(Justify::Right), + "SPEED".cell().bold(true).justify(Justify::Right), + ] +} + +fn print_models_table(models: &[crate::types::ModelInfo], s: &Styles) { + let use_color = s.use_color; + let rows: Vec> = models.iter().map(|m| model_row(m, use_color)).collect(); + let table = rows + .table() + .title(models_title()) + .border(Border::builder().build()) + .separator(Separator::builder().build()); + let _ = print_stdout(table); +} + fn read_stdin_prompt() -> Option { let stdin = io::stdin(); if stdin.is_terminal() { @@ -812,46 +848,48 @@ async fn test_models_via_server( bail!("No models found"); } - println!( - "{}", - s.bold_dim.apply_to(format!( - "{:<30} {:<12} {:>10} {:>7} {:>7} {:>10} RESULT", - "MODEL", "PROVIDER", "CONTEXT", "COST", "", "SPEED", - )), - ); + let use_color = s.use_color; + let mut title = models_title(); + title.push("RESULT".cell().bold(true)); + let mut rows: Vec> = Vec::new(); let mut failures = 0u32; for info in &models_to_test { + eprint!("Testing {}...", info.id); let result = test_model_via_server(&server.client, &server.base_url, &info.id).await; + eprintln!(" done"); - let (status_color, status) = match result { - Ok(resp) if resp.status == "ok" => (&s.green, "ok".to_string()), + let (result_color, status) = match result { + Ok(resp) if resp.status == "ok" => (Color::Green, "ok".to_string()), Ok(resp) => { failures += 1; let msg = resp .error_message .unwrap_or_else(|| "unknown error".to_string()); - (&s.red, format!("error: {msg}")) + (Color::Red, format!("error: {msg}")) } Err(e) => { failures += 1; - (&s.red, format!("error: {e}")) + (Color::Red, format!("error: {e}")) } }; - println!( - "{} {} {:>10} {:>7} / {:<7} {} {}", - s.bold.apply_to(format!("{:<30}", info.id)), - s.dim.apply_to(format!("{:<12}", info.provider)), - format_context_window(info.limits.context_window), - format_cost(info.costs.input_cost_per_mtok), - format_cost(info.costs.output_cost_per_mtok), - s.cyan - .apply_to(format!("{:>10}", format_speed(info.estimated_output_tps))), - status_color.apply_to(&status), + let mut row = model_row(info, use_color); + row.push( + status + .cell() + .foreground_color(color_if(use_color, result_color)), ); + rows.push(row); } + let table = rows + .table() + .title(title) + .border(Border::builder().build()) + .separator(Separator::builder().build()); + print_stdout(table)?; + if failures > 0 { bail!("{failures} model(s) failed"); } @@ -919,16 +957,14 @@ async fn test_models(provider: Option<&str>, model: Option<&str>, s: &Styles) -> bail!("No models found"); } - println!( - "{}", - s.bold_dim.apply_to(format!( - "{:<30} {:<12} {:>10} {:>7} {:>7} {:>10} RESULT", - "MODEL", "PROVIDER", "CONTEXT", "COST", "", "SPEED", - )), - ); + let use_color = s.use_color; + let mut title = models_title(); + title.push("RESULT".cell().bold(true)); + let mut rows: Vec> = Vec::new(); let mut failures = 0u32; for info in &models_to_test { + eprint!("Testing {}...", info.id); let params = GenerateParams::new(&info.id) .provider(&info.provider) .prompt("Say OK") @@ -936,32 +972,36 @@ async fn test_models(provider: Option<&str>, model: Option<&str>, s: &Styles) -> let result = tokio::time::timeout(Duration::from_secs(30), generate::generate(params)).await; + eprintln!(" done"); - let (status_color, status) = match result { - Ok(Ok(_)) => (&s.green, "ok".to_string()), + let (result_color, status) = match result { + Ok(Ok(_)) => (Color::Green, "ok".to_string()), Ok(Err(e)) => { failures += 1; - (&s.red, format!("error: {e}")) + (Color::Red, format!("error: {e}")) } Err(_) => { failures += 1; - (&s.red, "error: timeout (30s)".to_string()) + (Color::Red, "error: timeout (30s)".to_string()) } }; - println!( - "{} {} {:>10} {:>7} / {:<7} {} {}", - s.bold.apply_to(format!("{:<30}", info.id)), - s.dim.apply_to(format!("{:<12}", info.provider)), - format_context_window(info.limits.context_window), - format_cost(info.costs.input_cost_per_mtok), - format_cost(info.costs.output_cost_per_mtok), - s.cyan - .apply_to(format!("{:>10}", format_speed(info.estimated_output_tps))), - status_color.apply_to(&status), + let mut row = model_row(info, use_color); + row.push( + status + .cell() + .foreground_color(color_if(use_color, result_color)), ); + rows.push(row); } + let table = rows + .table() + .title(title) + .border(Border::builder().build()) + .separator(Separator::builder().build()); + print_stdout(table)?; + if failures > 0 { bail!("{failures} model(s) failed"); } diff --git a/lib/crates/fabro-util/src/lib.rs b/lib/crates/fabro-util/src/lib.rs index 194432a50..d8f8e87dd 100644 --- a/lib/crates/fabro-util/src/lib.rs +++ b/lib/crates/fabro-util/src/lib.rs @@ -1,4 +1,5 @@ pub mod check_report; +pub mod path; pub mod redact; pub mod run_log; pub mod telemetry; diff --git a/lib/crates/fabro-util/src/path.rs b/lib/crates/fabro-util/src/path.rs new file mode 100644 index 000000000..992cefb96 --- /dev/null +++ b/lib/crates/fabro-util/src/path.rs @@ -0,0 +1,28 @@ +use std::path::{Path, PathBuf}; + +/// Expand `~/` prefix to the user's home directory. +pub fn expand_tilde(path: &Path) -> PathBuf { + if let Ok(rest) = path.strip_prefix("~") { + if let Some(home) = dirs::home_dir() { + return home.join(rest); + } + } + path.to_path_buf() +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn expand_tilde_with_home_prefix() { + let result = expand_tilde(Path::new("~/foo/bar")); + assert!(result != Path::new("~/foo/bar")); + assert!(result.ends_with("foo/bar")); + } + + #[test] + fn expand_tilde_without_prefix() { + assert_eq!(expand_tilde(Path::new("/abs/path")), Path::new("/abs/path")); + } +} diff --git a/lib/crates/fabro-workflows/Cargo.toml b/lib/crates/fabro-workflows/Cargo.toml index 93f4aea64..42f804caa 100644 --- a/lib/crates/fabro-workflows/Cargo.toml +++ b/lib/crates/fabro-workflows/Cargo.toml @@ -53,6 +53,7 @@ sha2 = { workspace = true } shlex = "1" strsim = "0.11" git2.workspace = true +cli-table.workspace = true console.workspace = true indicatif.workspace = true tokio-util.workspace = true diff --git a/lib/crates/fabro-workflows/src/cli/inspect.rs b/lib/crates/fabro-workflows/src/cli/inspect.rs index 9af2da4af..ecdf20885 100644 --- a/lib/crates/fabro-workflows/src/cli/inspect.rs +++ b/lib/crates/fabro-workflows/src/cli/inspect.rs @@ -93,6 +93,7 @@ mod tests { labels: Default::default(), base_branch: None, workflow_slug: None, + host_repo_path: None, }; manifest.save(&run_dir.join("manifest.json")).unwrap(); @@ -134,12 +135,7 @@ mod tests { }; sandbox.save(&run_dir.join("sandbox.json")).unwrap(); - let output = inspect_run_dir( - "test-run", - &run_dir, - RunStatus::Concluded(StageStatus::Success), - ) - .unwrap(); + let output = inspect_run_dir("test-run", &run_dir, RunStatus::Succeeded).unwrap(); assert_eq!(output.run_id, "test-run"); assert_eq!(output.run_dir, run_dir); @@ -167,6 +163,7 @@ mod tests { labels: Default::default(), base_branch: None, workflow_slug: None, + host_repo_path: None, }; manifest.save(&run_dir.join("manifest.json")).unwrap(); @@ -184,7 +181,7 @@ mod tests { let output = InspectOutput { run_id: "id-1".to_string(), run_dir: PathBuf::from("/tmp/run"), - status: RunStatus::Unknown, + status: RunStatus::Dead, manifest: None, conclusion: None, checkpoint: None, diff --git a/lib/crates/fabro-workflows/src/cli/logs.rs b/lib/crates/fabro-workflows/src/cli/logs.rs index c9d3164e8..084209d06 100644 --- a/lib/crates/fabro-workflows/src/cli/logs.rs +++ b/lib/crates/fabro-workflows/src/cli/logs.rs @@ -22,7 +22,7 @@ pub struct LogsArgs { #[arg(short = 'n', long)] pub tail: Option, /// Formatted colored output with rendered assistant text - #[arg(long)] + #[arg(short = 'p', long)] pub pretty: bool, } @@ -214,15 +214,69 @@ pub fn format_event_pretty(line: &str, styles: &fabro_util::terminal::Styles) -> "WorkflowRunCompleted" => { let duration = format_duration_ms(envelope.get("duration_ms")); + let status_str = match str_field(&envelope, "status") { + Some(s) if !s.is_empty() => s, + _ => "success", + }; + let status_upper = status_str.to_uppercase(); + let status_style = match status_str { + "success" | "partial_success" => &styles.bold_green, + _ => &styles.bold_red, + }; let cost = format_cost(envelope.get("total_cost")); - Some(format!( - "{} {} {} {} {}", + + let mut lines = vec![format!( + "{} {} {} {}", styles.dim.apply_to(&ts), - styles.bold_green.apply_to("\u{2713} Completed"), + status_style.apply_to(format!("\u{2713} {status_upper}")), styles.bold.apply_to(&duration), styles.dim.apply_to(&cost), - "", - )) + )]; + + if let Some(usage) = envelope.get("usage") { + let total = usage + .get("total_tokens") + .and_then(|v| v.as_i64()) + .unwrap_or(0); + let pad = " ".repeat(ts.len() + 1); + if total > 0 { + lines.push(format!( + "{}{}", + pad, + styles + .dim + .apply_to(format!("Tokens: {}", format_tokens(total as u64))) + )); + } + if let Some(cr) = usage.get("cache_read_tokens").and_then(|v| v.as_i64()) { + let cw = usage + .get("cache_write_tokens") + .and_then(|v| v.as_i64()) + .unwrap_or(0); + lines.push(format!( + "{}{}", + pad, + styles.dim.apply_to(format!( + "Cache: {} read, {} write", + format_tokens(cr as u64), + format_tokens(cw as u64) + )) + )); + } + if let Some(r) = usage.get("reasoning_tokens").and_then(|v| v.as_i64()) { + if r > 0 { + lines.push(format!( + "{}{}", + pad, + styles + .dim + .apply_to(format!("Reasoning: {} tokens", format_tokens(r as u64))) + )); + } + } + } + + Some(lines.join("\n")) } "WorkflowRunFailed" => { @@ -447,6 +501,59 @@ pub fn format_event_pretty(line: &str, styles: &fabro_util::terminal::Styles) -> )) } + "PullRequestCreated" => { + let url = str_field(&envelope, "pr_url").unwrap_or("?"); + let draft = envelope + .get("draft") + .and_then(|v| v.as_bool()) + .unwrap_or(false); + let label = if draft { "Draft PR:" } else { "PR:" }; + Some(format!( + "{} {} {}", + styles.dim.apply_to(&ts), + styles.bold.apply_to(label), + url, + )) + } + + "PullRequestFailed" => { + let error = str_field(&envelope, "error").unwrap_or("unknown error"); + Some(format!( + "{} {} {}", + styles.dim.apply_to(&ts), + styles.bold_red.apply_to("PR failed:"), + styles.red.apply_to(error), + )) + } + + "RetroCompleted" => { + let duration = format_duration_ms(envelope.get("duration_ms")); + Some(format!( + "{} {} Retro {}", + styles.dim.apply_to(&ts), + styles.green.apply_to("\u{2713}"), + duration, + )) + } + + "RetroFailed" => { + let error = str_field(&envelope, "error").unwrap_or("unknown error"); + let duration = format_duration_ms(envelope.get("duration_ms")); + Some(format!( + "{} {} Retro {} {}", + styles.dim.apply_to(&ts), + styles.bold_red.apply_to("\u{2717}"), + duration, + styles.red.apply_to(error), + )) + } + + "RetroStarted" => Some(format!( + "{} {} Retro", + styles.dim.apply_to(&ts), + styles.bold_cyan.apply_to("\u{25b6}"), + )), + // Noise events — skip "Agent.SessionStarted" | "Agent.SessionEnded" @@ -460,8 +567,15 @@ pub fn format_event_pretty(line: &str, styles: &fabro_util::terminal::Styles) -> | "SetupStarted" | "SetupCommandStarted" | "SetupCommandCompleted" - | "CheckpointSaved" - | "GitCheckpoint" + | "CheckpointCompleted" + | "CheckpointFailed" + | "GitCommit" + | "GitPush" + | "GitBranch" + | "GitWorktreeAdd" + | "GitWorktreeRemove" + | "GitFetch" + | "GitReset" | "AssetsCaptured" => None, _ => None, @@ -705,11 +819,65 @@ mod tests { #[test] fn pretty_workflow_run_completed() { let styles = no_color_styles(); - let line = r#"{"ts":"2026-01-01T14:23:32Z","run_id":"abc123","event":"WorkflowRunCompleted","duration_ms":25000,"total_cost":0.57}"#; + let line = r#"{"ts":"2026-01-01T14:23:32Z","run_id":"abc123","event":"WorkflowRunCompleted","duration_ms":25000,"status":"success","total_cost":0.57,"usage":{"input_tokens":5000,"output_tokens":2000,"total_tokens":7000,"cache_read_tokens":3000,"cache_write_tokens":500,"reasoning_tokens":800}}"#; let result = format_event_pretty(line, &styles).unwrap(); - assert!(result.contains("Completed"), "got: {result}"); + assert!(result.contains("SUCCESS"), "got: {result}"); assert!(result.contains("25s"), "got: {result}"); assert!(result.contains("$0.57"), "got: {result}"); + assert!(result.contains("7.0k toks"), "got: {result}"); + assert!(result.contains("Cache:"), "got: {result}"); + assert!(result.contains("3.0k toks read"), "got: {result}"); + assert!(result.contains("Reasoning:"), "got: {result}"); + } + + #[test] + fn pretty_workflow_run_completed_backward_compat() { + let styles = no_color_styles(); + // Old JSONL without status/usage still renders + let line = r#"{"ts":"2026-01-01T14:23:32Z","run_id":"abc123","event":"WorkflowRunCompleted","duration_ms":25000,"total_cost":0.57}"#; + let result = format_event_pretty(line, &styles).unwrap(); + assert!(result.contains("SUCCESS"), "got: {result}"); + assert!(result.contains("25s"), "got: {result}"); + assert!(result.contains("$0.57"), "got: {result}"); + // No usage lines when usage is absent + assert!(!result.contains("Tokens:"), "got: {result}"); + } + + #[test] + fn pretty_workflow_run_completed_fail_status() { + let styles = no_color_styles(); + let line = r#"{"ts":"2026-01-01T14:23:32Z","event":"WorkflowRunCompleted","duration_ms":25000,"status":"fail"}"#; + let result = format_event_pretty(line, &styles).unwrap(); + assert!(result.contains("FAIL"), "got: {result}"); + } + + #[test] + fn pretty_pull_request_created() { + let styles = no_color_styles(); + let line = r#"{"ts":"2026-01-01T14:25:00Z","event":"PullRequestCreated","pr_url":"https://github.com/owner/repo/pull/42","pr_number":42,"draft":false}"#; + let result = format_event_pretty(line, &styles).unwrap(); + assert!(result.contains("PR:"), "got: {result}"); + assert!( + result.contains("https://github.com/owner/repo/pull/42"), + "got: {result}" + ); + } + + #[test] + fn pretty_pull_request_created_draft() { + let styles = no_color_styles(); + let line = r#"{"ts":"2026-01-01T14:25:00Z","event":"PullRequestCreated","pr_url":"https://github.com/owner/repo/pull/42","pr_number":42,"draft":true}"#; + let result = format_event_pretty(line, &styles).unwrap(); + assert!(result.contains("Draft PR:"), "got: {result}"); + } + + #[test] + fn pretty_pull_request_failed() { + let styles = no_color_styles(); + let line = r#"{"ts":"2026-01-01T14:25:00Z","event":"PullRequestFailed","error":"auth token expired"}"#; + let result = format_event_pretty(line, &styles).unwrap(); + assert!(result.contains("PR failed:"), "got: {result}"); + assert!(result.contains("auth token expired"), "got: {result}"); } #[test] diff --git a/lib/crates/fabro-workflows/src/cli/mod.rs b/lib/crates/fabro-workflows/src/cli/mod.rs index 8f59dbf95..79b270bd8 100644 --- a/lib/crates/fabro-workflows/src/cli/mod.rs +++ b/lib/crates/fabro-workflows/src/cli/mod.rs @@ -279,6 +279,16 @@ pub fn relative_path(path: &Path) -> String { tilde_path(path) } +/// Return `Some(color)` when color is enabled, `None` otherwise. +/// Used with `cli_table`'s `.foreground_color()` which accepts `Option`. +pub(crate) fn color_if(use_color: bool, color: cli_table::Color) -> Option { + if use_color { + Some(color) + } else { + None + } +} + /// Shorten an absolute path by replacing the home directory prefix with `~`. pub fn tilde_path(path: &Path) -> String { if let Some(home) = dirs::home_dir() { diff --git a/lib/crates/fabro-workflows/src/cli/progress.rs b/lib/crates/fabro-workflows/src/cli/progress.rs index 39c904b3a..833856646 100644 --- a/lib/crates/fabro-workflows/src/cli/progress.rs +++ b/lib/crates/fabro-workflows/src/cli/progress.rs @@ -73,7 +73,7 @@ pub(crate) fn format_duration_short(d: Duration) -> String { } } -fn format_duration_ms(ms: u64) -> String { +pub(crate) fn format_duration_ms(ms: u64) -> String { format_duration_short(Duration::from_millis(ms)) } @@ -598,6 +598,17 @@ impl ProgressUI { } } } + WorkflowRunEvent::RetroStarted => { + self.on_stage_started("retro", "Retro", None); + } + WorkflowRunEvent::RetroCompleted { duration_ms } => { + let dur = format_duration_ms(*duration_ms); + self.finish_stage("retro", "Retro", green_check(), &dur); + } + WorkflowRunEvent::RetroFailed { duration_ms, .. } => { + let dur = format_duration_ms(*duration_ms); + self.finish_stage("retro", "Retro", red_cross(), &dur); + } _ => {} } } diff --git a/lib/crates/fabro-workflows/src/cli/rewind.rs b/lib/crates/fabro-workflows/src/cli/rewind.rs index bc8a37957..44cc027e9 100644 --- a/lib/crates/fabro-workflows/src/cli/rewind.rs +++ b/lib/crates/fabro-workflows/src/cli/rewind.rs @@ -2,6 +2,8 @@ use std::collections::HashMap; use anyhow::{bail, Context, Result}; use clap::Args; +use cli_table::format::{Border, Separator}; +use cli_table::{print_stderr, Cell, CellStruct, Color, Style, Table}; use fabro_git_storage::branchstore::{BranchStore, CommitInfo}; use fabro_git_storage::gitobj::Store; use fabro_util::terminal::Styles; @@ -280,39 +282,53 @@ pub fn print_timeline( return; } - eprintln!( - " {} {} {}", - styles.bold_dim.apply_to(format!("{:<6}", "@")), - styles.bold_dim.apply_to(format!("{:<30}", "Node")), - styles.bold_dim.apply_to("Details"), - ); + let use_color = styles.use_color; - for entry in timeline { - let ordinal_str = format!("@{}", entry.ordinal); - let mut details = Vec::new(); - if entry.visit > 1 { - details.push(format!("visit {}, loop", entry.visit)); - } - if parallel_map.contains_key(&entry.node_name) { - details.push("parallel interior".to_string()); - } - if entry.run_commit_sha.is_none() { - details.push("no run commit".to_string()); - } + let title = vec![ + "@".cell().bold(true), + "Node".cell().bold(true), + "Details".cell().bold(true), + ]; - let detail_str = if details.is_empty() { - String::new() - } else { - format!("({})", details.join(", ")) - }; + let rows: Vec> = timeline + .iter() + .map(|entry| { + let ordinal_str = format!("@{}", entry.ordinal); + let mut details = Vec::new(); + if entry.visit > 1 { + details.push(format!("visit {}, loop", entry.visit)); + } + if parallel_map.contains_key(&entry.node_name) { + details.push("parallel interior".to_string()); + } + if entry.run_commit_sha.is_none() { + details.push("no run commit".to_string()); + } - eprintln!( - " {} {:<30} {}", - styles.cyan.apply_to(format!("{ordinal_str:<6}")), - entry.node_name, - styles.dim.apply_to(detail_str), - ); - } + let detail_str = if details.is_empty() { + String::new() + } else { + format!("({})", details.join(", ")) + }; + + vec![ + ordinal_str + .cell() + .foreground_color(super::color_if(use_color, Color::Cyan)), + entry.node_name.clone().cell(), + detail_str + .cell() + .foreground_color(super::color_if(use_color, Color::Ansi256(8))), + ] + }) + .collect(); + + let table = rows + .table() + .title(title) + .border(Border::builder().build()) + .separator(Separator::builder().build()); + let _ = print_stderr(table); } /// Move both refs backward to the target checkpoint. diff --git a/lib/crates/fabro-workflows/src/cli/run.rs b/lib/crates/fabro-workflows/src/cli/run.rs index 4317266ce..8051c91ea 100644 --- a/lib/crates/fabro-workflows/src/cli/run.rs +++ b/lib/crates/fabro-workflows/src/cli/run.rs @@ -46,7 +46,8 @@ fn resolve_cli_goal( match (goal, goal_file) { (Some(g), _) => Ok(Some(g.clone())), (_, Some(path)) => { - let content = std::fs::read_to_string(path) + let path = fabro_util::path::expand_tilde(path); + let content = std::fs::read_to_string(&path) .with_context(|| format!("failed to read goal file: {}", path.display()))?; debug!(path = %path.display(), "Goal loaded from file"); Ok(Some(content)) @@ -455,17 +456,37 @@ pub async fn run_command( // 3. Create logs directory let run_id = args.run_id.unwrap_or_else(|| ulid::Ulid::new().to_string()); let run_dir = args.run_dir.unwrap_or_else(|| { - let base = dirs::home_dir() - .expect("could not determine home directory") - .join(".fabro") - .join("runs"); - base.join(format!("{}-{}", Local::now().format("%Y%m%d"), run_id)) + if args.dry_run { + std::env::temp_dir().join("fabro-dry-run").join(&run_id) + } else { + let base = dirs::home_dir() + .expect("could not determine home directory") + .join(".fabro") + .join("runs"); + base.join(format!("{}-{}", Local::now().format("%Y%m%d"), run_id)) + } }); tokio::fs::create_dir_all(&run_dir).await?; fabro_util::run_log::activate(&run_dir.join("cli.log")) .context("Failed to activate per-run log")?; tokio::fs::write(run_dir.join("graph.fabro"), &source).await?; tokio::fs::write(run_dir.join("run.pid"), std::process::id().to_string()).await?; + super::runs::write_run_status( + &run_dir, + crate::run_status::RunStatus::Starting, + Some(crate::run_status::StatusReason::SandboxInitializing), + ); + + // Safety net: mark as failed if we exit before engine.run() (e.g. sandbox init failure) + let status_run_dir = run_dir.clone(); + let status_guard = scopeguard::guard((), move |()| { + super::runs::write_run_status( + &status_run_dir, + crate::run_status::RunStatus::Failed, + Some(crate::run_status::StatusReason::SandboxInitFailed), + ); + }); + if workflow_path.extension().is_some_and(|ext| ext == "toml") { if let Ok(toml_contents) = tokio::fs::read(workflow_path).await { tokio::fs::write(run_dir.join("run.toml"), toml_contents).await?; @@ -486,13 +507,17 @@ pub async fn run_command( // 3. Build event emitter let mut emitter = EventEmitter::new(); - // Track the last git commit SHA from GitCheckpoint events + // Track the last git commit SHA from CheckpointCompleted events let last_git_sha: Arc>> = Arc::new(Mutex::new(None)); { let sha_clone = Arc::clone(&last_git_sha); emitter.on_event(move |event| { - if let crate::event::WorkflowRunEvent::GitCheckpoint { git_commit_sha, .. } = event { - *sha_clone.lock().unwrap() = Some(git_commit_sha.clone()); + if let crate::event::WorkflowRunEvent::CheckpointCompleted { + git_commit_sha: Some(sha), + .. + } = event + { + *sha_clone.lock().unwrap() = Some(sha.clone()); } }); } @@ -1197,6 +1222,9 @@ pub async fn run_command( workflow_slug: workflow_slug.clone(), }; + // Defuse the status guard — engine.run() will write "running" and conclusion handles "concluded" + scopeguard::ScopeGuard::into_inner(status_guard); + let run_start = Instant::now(); let engine_result = if let Some(ref checkpoint_path) = args.resume { let checkpoint = Checkpoint::load(checkpoint_path)?; @@ -1217,6 +1245,32 @@ pub async fn run_command( Err(e) => (crate::outcome::StageStatus::Fail, Some(e.to_string())), }; + // Map engine result to RunStatus + StatusReason + let (run_status, status_reason) = match &engine_result { + Ok(o) => match o.status { + StageStatus::Success | StageStatus::Skipped => ( + crate::run_status::RunStatus::Succeeded, + Some(crate::run_status::StatusReason::Completed), + ), + StageStatus::PartialSuccess => ( + crate::run_status::RunStatus::Succeeded, + Some(crate::run_status::StatusReason::PartialSuccess), + ), + StageStatus::Fail | StageStatus::Retry => ( + crate::run_status::RunStatus::Failed, + Some(crate::run_status::StatusReason::WorkflowError), + ), + }, + Err(crate::error::FabroError::Cancelled) => ( + crate::run_status::RunStatus::Failed, + Some(crate::run_status::StatusReason::Cancelled), + ), + Err(_) => ( + crate::run_status::RunStatus::Failed, + Some(crate::run_status::StatusReason::WorkflowError), + ), + }; + // Load checkpoint and stage durations to populate per-stage data let checkpoint = Checkpoint::load(&run_dir.join("checkpoint.json")).ok(); let stage_durations = crate::retro::extract_stage_durations(&run_dir); @@ -1265,6 +1319,7 @@ pub async fn run_command( total_retries, }; let _ = conclusion.save(&run_dir.join("conclusion.json")); + super::runs::write_run_status(&run_dir, run_status, status_reason); } // Auto-derive retro (always, cheap) and optionally run retro agent @@ -1290,7 +1345,7 @@ pub async fn run_command( provider_enum, &model, styles, - Some(&progress_ui), + Some(Arc::clone(&emitter)), ) .await; } @@ -1625,15 +1680,19 @@ async fn run_from_branch( // Set up logs directory let run_dir = args.run_dir.unwrap_or_else(|| { - let base = dirs::home_dir() - .expect("could not determine home directory") - .join(".fabro") - .join("runs"); - base.join(format!( - "{}-{}", - chrono::Local::now().format("%Y%m%d"), - run_id - )) + if args.dry_run { + std::env::temp_dir().join("fabro-dry-run").join(&run_id) + } else { + let base = dirs::home_dir() + .expect("could not determine home directory") + .join(".fabro") + .join("runs"); + base.join(format!( + "{}-{}", + chrono::Local::now().format("%Y%m%d"), + run_id + )) + } }); tokio::fs::create_dir_all(&run_dir).await?; fabro_util::run_log::activate(&run_dir.join("cli.log")) @@ -1847,7 +1906,7 @@ async fn run_from_branch( provider_enum, &model, styles, - None, + Some(Arc::clone(&emitter)), ) .await; } @@ -2276,7 +2335,7 @@ async fn generate_retro( provider_enum: Provider, model: &str, styles: &'static Styles, - progress_ui: Option<&Arc>>, + emitter: Option>, ) { let cp = match Checkpoint::load(&run_dir.join("checkpoint.json")) { Ok(cp) => cp, @@ -2315,27 +2374,14 @@ async fn generate_retro( eprintln!("\n{}", styles.bold.apply_to("=== Retro ===")); let retro_start = std::time::Instant::now(); - let emitter = if let Some(pui) = progress_ui { - let mut em = EventEmitter::new(); - progress::ProgressUI::register(pui, &mut em); - let em = Arc::new(em); - em.emit(&crate::event::WorkflowRunEvent::StageStarted { - node_id: "retro".to_string(), - name: "Retro".to_string(), - index: 0, - handler_type: Some("agent".to_string()), - script: None, - attempt: 1, - max_attempts: 1, - }); - Some(em) + if let Some(ref em) = emitter { + em.emit(&crate::event::WorkflowRunEvent::RetroStarted); } else { eprintln!( "{}", styles.dim.apply_to(format!("Running retro ({model})...")) ); - None - }; + } let narrative_result = if dry_run_mode { Ok(crate::retro_agent::dry_run_narrative()) @@ -2355,26 +2401,19 @@ async fn generate_retro( let retro_dur_elapsed = retro_start.elapsed(); if let Some(ref em) = emitter { - let status = if narrative_result.is_ok() { - "success" - } else { - "fail" - }; - em.emit(&crate::event::WorkflowRunEvent::StageCompleted { - node_id: "retro".to_string(), - name: "Retro".to_string(), - index: 0, - duration_ms: retro_dur_elapsed.as_millis() as u64, - status: status.to_string(), - preferred_label: None, - suggested_next_ids: vec![], - usage: None, - failure: None, - notes: None, - files_touched: vec![], - attempt: 1, - max_attempts: 1, - }); + match &narrative_result { + Ok(_) => { + em.emit(&crate::event::WorkflowRunEvent::RetroCompleted { + duration_ms: retro_dur_elapsed.as_millis() as u64, + }); + } + Err(e) => { + em.emit(&crate::event::WorkflowRunEvent::RetroFailed { + error: e.to_string(), + duration_ms: retro_dur_elapsed.as_millis() as u64, + }); + } + } } let retro_dur = progress::format_duration_short(retro_dur_elapsed); diff --git a/lib/crates/fabro-workflows/src/cli/runs.rs b/lib/crates/fabro-workflows/src/cli/runs.rs index 472684edd..0c55842c8 100644 --- a/lib/crates/fabro-workflows/src/cli/runs.rs +++ b/lib/crates/fabro-workflows/src/cli/runs.rs @@ -1,47 +1,16 @@ use std::collections::HashMap; -use std::fmt; use std::path::{Path, PathBuf}; use anyhow::{bail, Context, Result}; use chrono::{DateTime, Utc}; use clap::Args; +use cli_table::format::{Border, Justify, Separator}; +use cli_table::{print_stdout, Cell, CellStruct, Color, Style, Table}; +use fabro_util::terminal::Styles; use serde::Serialize; use tracing::{debug, info, warn}; -use crate::outcome::StageStatus; - -/// Status of a run directory — either concluded with a `StageStatus`, actively running, or unknown. -#[derive(Debug, Clone, PartialEq, Eq)] -pub enum RunStatus { - Concluded(StageStatus), - Running, - Unknown, -} - -impl Serialize for RunStatus { - fn serialize(&self, serializer: S) -> Result - where - S: serde::Serializer, - { - serializer.serialize_str(&self.to_string()) - } -} - -impl RunStatus { - pub fn is_running(&self) -> bool { - matches!(self, RunStatus::Running) - } -} - -impl fmt::Display for RunStatus { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - match self { - RunStatus::Concluded(s) => write!(f, "{s}"), - RunStatus::Running => write!(f, "running"), - RunStatus::Unknown => write!(f, "unknown"), - } - } -} +pub use crate::run_status::{RunStatus, RunStatusRecord, StatusReason}; #[derive(Args)] pub struct RunFilterArgs { @@ -70,6 +39,10 @@ pub struct RunsListArgs { /// Output as JSON #[arg(long)] pub json: bool, + + /// Show all runs, not just running (like docker ps -a) + #[arg(short = 'a', long)] + pub all: bool, } #[derive(Args)] @@ -86,6 +59,17 @@ pub struct RunsPruneArgs { pub yes: bool, } +#[derive(Args)] +pub struct RunsRemoveArgs { + /// Run IDs or workflow names to remove + #[arg(required = true)] + pub runs: Vec, + + /// Force removal of active runs + #[arg(short, long)] + pub force: bool, +} + #[derive(Debug, Clone, Serialize)] pub struct RunInfo { pub run_id: String, @@ -94,8 +78,17 @@ pub struct RunInfo { #[serde(default, skip_serializing_if = "Option::is_none")] pub workflow_slug: Option, pub status: RunStatus, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub status_reason: Option, pub start_time: String, pub labels: HashMap, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub duration_ms: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub total_cost: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub host_repo_path: Option, + pub goal: String, #[serde(skip)] pub start_time_dt: Option>, #[serde(skip)] @@ -133,27 +126,34 @@ pub fn scan_runs(base: &Path) -> Result> { let run_id = manifest.run_id; let workflow_name = manifest.workflow_name; let workflow_slug = manifest.workflow_slug; + let host_repo_path = manifest.host_repo_path; + let goal = manifest.goal; let start_time_dt = manifest.start_time; let start_time = start_time_dt.to_rfc3339(); let labels = manifest.labels; - let (status, end_time) = read_status(&path); + let si = read_status(&path); runs.push(RunInfo { run_id, dir_name, workflow_name, workflow_slug, - status, + status: si.status, + status_reason: si.reason, start_time, labels, + duration_ms: si.duration_ms, + total_cost: si.total_cost, + host_repo_path, start_time_dt: Some(start_time_dt), - end_time, + end_time: si.end_time, path, + goal, is_orphan: false, }); } else { - // Orphan directory — no manifest.json + // No manifest.json — check for status.json (starting run) vs true orphan let mtime_dt = entry .metadata() .ok() @@ -161,18 +161,34 @@ pub fn scan_runs(base: &Path) -> Result> { .map(|t| -> DateTime { t.into() }); let mtime = mtime_dt.map(|dt| dt.to_rfc3339()).unwrap_or_default(); + let run_id = std::fs::read_to_string(path.join("id.txt")) + .map(|s| s.trim().to_string()) + .unwrap_or_else(|_| dir_name.clone()); + + let si = read_status(&path); + let is_orphan = matches!(si.status, RunStatus::Dead); runs.push(RunInfo { - run_id: dir_name.clone(), + run_id, dir_name, - workflow_name: "[no manifest]".to_string(), + workflow_name: if is_orphan { + "[no manifest]" + } else { + "[starting]" + } + .to_string(), workflow_slug: None, - status: RunStatus::Unknown, + status: si.status, + status_reason: si.reason, start_time: mtime, labels: HashMap::new(), + duration_ms: si.duration_ms, + total_cost: si.total_cost, + host_repo_path: None, start_time_dt: mtime_dt, - end_time: None, + end_time: si.end_time, path, - is_orphan: true, + goal: String::new(), + is_orphan, }); } } @@ -182,17 +198,70 @@ pub fn scan_runs(base: &Path) -> Result> { Ok(runs) } -fn read_status(run_dir: &Path) -> (RunStatus, Option>) { - if let Ok(conclusion) = crate::conclusion::Conclusion::load(&run_dir.join("conclusion.json")) { - return ( - RunStatus::Concluded(conclusion.status), - Some(conclusion.timestamp), - ); +struct StatusInfo { + status: RunStatus, + reason: Option, + end_time: Option>, + duration_ms: Option, + total_cost: Option, +} + +impl StatusInfo { + fn simple(status: RunStatus) -> Self { + Self { + status, + reason: None, + end_time: None, + duration_ms: None, + total_cost: None, + } } - if run_dir.join("run.pid").exists() { - return (RunStatus::Running, None); +} + +/// Write the run status to `status.json` (best-effort). +pub fn write_run_status(run_dir: &Path, status: RunStatus, reason: Option) { + let record = RunStatusRecord::new(status, reason); + if let Err(e) = record.save(&run_dir.join("status.json")) { + warn!("failed to write status.json for {}: {e}", run_dir.display()); } - (RunStatus::Unknown, None) +} + +fn read_status(run_dir: &Path) -> StatusInfo { + // 1. status.json is authoritative + if let Ok(record) = RunStatusRecord::load(&run_dir.join("status.json")) { + // For terminal statuses, pull duration/cost from conclusion.json if available + if record.status.is_terminal() { + if let Ok(conclusion) = + crate::conclusion::Conclusion::load(&run_dir.join("conclusion.json")) + { + return StatusInfo { + status: record.status, + reason: record.reason, + end_time: Some(conclusion.timestamp), + duration_ms: Some(conclusion.duration_ms), + total_cost: conclusion.total_cost, + }; + } + } + return StatusInfo { + status: record.status, + reason: record.reason, + end_time: None, + duration_ms: None, + total_cost: None, + }; + } + // No status.json → Dead (orphan) + StatusInfo::simple(RunStatus::Dead) +} + +/// Which run statuses to include in filtered results. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum StatusFilter { + /// Only include runs that are currently running. + RunningOnly, + /// Include runs of any status. + All, } /// Filter runs by criteria. Orphans are excluded unless `include_orphans` is true. @@ -202,9 +271,13 @@ pub fn filter_runs( workflow: Option<&str>, labels: &[(String, String)], include_orphans: bool, + status_filter: StatusFilter, ) -> Vec { runs.iter() .filter(|r| { + if status_filter == StatusFilter::RunningOnly && !r.status.is_active() { + return false; + } if r.is_orphan && !include_orphans { return false; } @@ -345,7 +418,43 @@ fn collapse_separators(s: &str) -> String { s.chars().filter(|c| *c != '-' && *c != '_').collect() } -pub fn list_command(args: &RunsListArgs) -> Result<()> { +/// Truncate a run ID to a short display form (12 characters). +fn short_run_id(id: &str) -> &str { + if id.len() > 12 { + &id[..12] + } else { + id + } +} + +fn truncate_goal(goal: &str, max_len: usize) -> String { + let line = goal.lines().next().unwrap_or(""); + let chars: Vec = line.chars().collect(); + if chars.len() <= max_len { + return line.to_string(); + } + let truncated: String = chars[..max_len - 3].iter().collect(); + format!("{truncated}...") +} + +use super::color_if; + +fn status_cell(status: &RunStatus, use_color: bool) -> CellStruct { + let text = status.to_string(); + let color = match status { + RunStatus::Succeeded => Some(Color::Green), + RunStatus::Failed => Some(Color::Red), + RunStatus::Running | RunStatus::Starting | RunStatus::Submitted => Some(Color::Cyan), + RunStatus::Removing => Some(Color::Yellow), + RunStatus::Paused => Some(Color::Magenta), + RunStatus::Dead => Some(Color::Ansi256(8)), + }; + text.cell() + .bold(use_color && color != Some(Color::Ansi256(8))) + .foreground_color(color_if(use_color, color.unwrap_or(Color::Ansi256(8)))) +} + +pub fn list_command(args: &RunsListArgs, styles: &Styles) -> Result<()> { let base = default_runs_base(); let runs = scan_runs(&base)?; let label_filters = parse_label_filters(&args.filter.label); @@ -355,6 +464,11 @@ pub fn list_command(args: &RunsListArgs) -> Result<()> { args.filter.workflow.as_deref(), &label_filters, args.filter.orphans, + if args.all { + StatusFilter::All + } else { + StatusFilter::RunningOnly + }, ); if args.json { @@ -363,41 +477,73 @@ pub fn list_command(args: &RunsListArgs) -> Result<()> { } if filtered.is_empty() { - eprintln!("No runs found."); + if args.all { + eprintln!("No runs found."); + } else { + eprintln!("No running processes found. Use -a to show all runs."); + } return Ok(()); } - // Print table header - let header = format!( - "{:<30} {:<25} {:<10} {:<25} LABELS", - "RUN ID", "WORKFLOW", "STATUS", "STARTED" - ); - println!("{header}"); - println!("{}", "-".repeat(100)); + // Reverse to oldest-first for display (scan_runs returns newest-first) + let mut display_runs = filtered; + display_runs.reverse(); - for run in &filtered { - let labels_str = run - .labels + let use_color = styles.use_color; + let title = vec![ + "RUN ID".cell().bold(true), + "WORKFLOW".cell().bold(true), + "STATUS".cell().bold(true), + "DIRECTORY".cell().bold(true), + "DURATION".cell().bold(true), + "GOAL".cell().bold(true), + ]; + + let rows: Vec> = + display_runs .iter() - .map(|(k, v)| format!("{k}={v}")) - .collect::>() - .join(", "); - let run_id_display = if run.run_id.len() > 28 { - format!("{}...", &run.run_id[..25]) - } else { - run.run_id.clone() - }; - let start_display = if run.start_time.len() > 23 { - run.start_time[..23].to_string() - } else { - run.start_time.clone() - }; - println!( - "{:<30} {:<25} {:<10} {:<25} {}", - run_id_display, run.workflow_name, run.status, start_display, labels_str - ); - } - eprintln!("\n{} run(s) listed.", filtered.len()); + .map(|run| { + let run_id_display = short_run_id(&run.run_id); + let duration_display = match run.duration_ms { + Some(ms) => super::progress::format_duration_ms(ms), + None => match run.start_time_dt { + Some(start) => { + let elapsed = Utc::now().signed_duration_since(start); + super::progress::format_duration_ms( + elapsed.num_milliseconds().max(0) as u64 + ) + } + None => "-".to_string(), + }, + }; + let dir_display = run + .host_repo_path + .as_deref() + .map(|p| super::tilde_path(Path::new(p))) + .unwrap_or_else(|| "-".to_string()); + vec![ + run_id_display + .cell() + .foreground_color(color_if(use_color, Color::Ansi256(8))), + run.workflow_name.clone().cell(), + status_cell(&run.status, use_color), + dir_display.cell(), + duration_display.cell(), + truncate_goal(&run.goal, 50) + .cell() + .foreground_color(color_if(use_color, Color::Ansi256(8))), + ] + }) + .collect(); + + let table = rows + .table() + .title(title) + .border(Border::builder().build()) + .separator(Separator::builder().build()); + print_stdout(table)?; + + eprintln!("\n{} run(s) listed.", display_runs.len()); Ok(()) } @@ -461,7 +607,7 @@ pub fn df_from(args: &DfArgs, data_dir: &Path, runs_base: &Path, logs_base: &Pat for run in &runs { let size = dir_size(&run.path); total_run_size += size; - let is_active = run.status.is_running(); + let is_active = run.status.is_active(); if is_active { active_count += 1; } else { @@ -471,7 +617,7 @@ pub fn df_from(args: &DfArgs, data_dir: &Path, runs_base: &Path, logs_base: &Pat run_details.push(RunSizeInfo { run_id: run.run_id.clone(), workflow_name: run.workflow_name.clone(), - status: run.status.clone(), + status: run.status, start_time_dt: run.start_time_dt, size, }); @@ -523,34 +669,48 @@ pub fn df_from(args: &DfArgs, data_dir: &Path, runs_base: &Path, logs_base: &Pat }; let log_reclaim_pct = if total_log_size > 0 { 100 } else { 0 }; - println!( - "{:<14}{:>5}{:>11}{:>12}{:>16}", - "TYPE", "COUNT", "ACTIVE", "SIZE", "RECLAIMABLE" - ); - println!( - "{:<14}{:>5}{:>11}{:>12}{:>12} ({run_reclaim_pct}%)", - "Runs", - runs.len(), - active_count, - format_size(total_run_size), - format_size(reclaimable_run_size), - ); - println!( - "{:<14}{:>5}{:>11}{:>12}{:>12} ({log_reclaim_pct}%)", - "Logs", - log_count, - "-", - format_size(total_log_size), - format_size(total_log_size), - ); - println!( - "{:<14}{:>5}{:>11}{:>12}{:>12} (0%)", - "Databases", - db_count, - "-", - format_size(total_db_size), - format_size(0), - ); + let df_title = vec![ + "TYPE".cell().bold(true), + "COUNT".cell().bold(true).justify(Justify::Right), + "ACTIVE".cell().bold(true).justify(Justify::Right), + "SIZE".cell().bold(true).justify(Justify::Right), + "RECLAIMABLE".cell().bold(true).justify(Justify::Right), + ]; + let df_rows: Vec> = vec![ + vec![ + "Runs".cell(), + runs.len().cell().justify(Justify::Right), + active_count.cell().justify(Justify::Right), + format_size(total_run_size).cell().justify(Justify::Right), + format!("{} ({run_reclaim_pct}%)", format_size(reclaimable_run_size)) + .cell() + .justify(Justify::Right), + ], + vec![ + "Logs".cell(), + log_count.cell().justify(Justify::Right), + "-".cell().justify(Justify::Right), + format_size(total_log_size).cell().justify(Justify::Right), + format!("{} ({log_reclaim_pct}%)", format_size(total_log_size)) + .cell() + .justify(Justify::Right), + ], + vec![ + "Databases".cell(), + db_count.cell().justify(Justify::Right), + "-".cell().justify(Justify::Right), + format_size(total_db_size).cell().justify(Justify::Right), + format!("{} (0%)", format_size(0)) + .cell() + .justify(Justify::Right), + ], + ]; + let df_table = df_rows + .table() + .title(df_title) + .border(Border::builder().build()) + .separator(Separator::builder().build()); + print_stdout(df_table)?; println!(); println!("Data directory: {}", data_dir.display()); @@ -558,50 +718,53 @@ pub fn df_from(args: &DfArgs, data_dir: &Path, runs_base: &Path, logs_base: &Pat // --- Verbose per-run breakdown --- if args.verbose { println!(); - println!( - "{:<30} {:<18} {:<10} {:>5} {:>12}", - "RUN ID", "WORKFLOW", "STATUS", "AGE", "SIZE" - ); + let verbose_title = vec![ + "RUN ID".cell().bold(true), + "WORKFLOW".cell().bold(true), + "STATUS".cell().bold(true), + "AGE".cell().bold(true).justify(Justify::Right), + "SIZE".cell().bold(true).justify(Justify::Right), + ]; let now = chrono::Utc::now(); - for detail in &run_details { - let run_id_display = if detail.run_id.len() > 28 { - format!("{}...", &detail.run_id[..25]) - } else { - detail.run_id.clone() - }; - let workflow_display = if detail.workflow_name.len() > 16 { - format!("{}...", &detail.workflow_name[..13]) - } else { - detail.workflow_name.clone() - }; - let age = if let Some(dt) = detail.start_time_dt { - let dur = now.signed_duration_since(dt); - if dur.num_days() > 0 { - format!("{}d", dur.num_days()) - } else if dur.num_hours() > 0 { - format!("{}h", dur.num_hours()) + let verbose_rows: Vec> = run_details + .iter() + .map(|detail| { + let run_id_display = short_run_id(&detail.run_id); + let workflow_display = truncate_goal(&detail.workflow_name, 16); + let age = if let Some(dt) = detail.start_time_dt { + let dur = now.signed_duration_since(dt); + if dur.num_days() > 0 { + format!("{}d", dur.num_days()) + } else if dur.num_hours() > 0 { + format!("{}h", dur.num_hours()) + } else { + format!("{}m", dur.num_minutes().max(1)) + } } else { - format!("{}m", dur.num_minutes().max(1)) - } - } else { - "-".to_string() - }; - let reclaimable_marker = if !detail.status.is_running() { - " *" - } else { - "" - }; - println!( - "{:<30} {:<18} {:<10} {:>5} {:>10}{}", - run_id_display, - workflow_display, - detail.status, - age, - format_size(detail.size), - reclaimable_marker, - ); - } + "-".to_string() + }; + let size_display = if !detail.status.is_active() { + format!("{} *", format_size(detail.size)) + } else { + format_size(detail.size) + }; + vec![ + run_id_display.cell(), + workflow_display.cell(), + detail.status.to_string().cell(), + age.cell().justify(Justify::Right), + size_display.cell().justify(Justify::Right), + ] + }) + .collect(); + + let verbose_table = verbose_rows + .table() + .title(verbose_title) + .border(Border::builder().build()) + .separator(Separator::builder().build()); + print_stdout(verbose_table)?; println!(); println!("* = reclaimable"); } @@ -640,6 +803,7 @@ pub fn prune_from(args: &RunsPruneArgs, base: &Path) -> Result<()> { args.filter.workflow.as_deref(), &label_filters, args.filter.orphans, + StatusFilter::All, ); // Determine if the user passed any explicit filters @@ -659,8 +823,8 @@ pub fn prune_from(args: &RunsPruneArgs, base: &Path) -> Result<()> { let now = Utc::now(); let cutoff = now - threshold; filtered.retain(|run| { - // Exclude running runs - if run.status.is_running() { + // Exclude active runs + if run.status.is_active() { return false; } // Use end_time if available, fall back to start_time @@ -705,6 +869,66 @@ pub fn prune_from(args: &RunsPruneArgs, base: &Path) -> Result<()> { Ok(()) } +pub async fn remove_command(args: &RunsRemoveArgs) -> Result<()> { + let base = default_runs_base(); + remove_from(args, &base).await +} + +pub async fn remove_from(args: &RunsRemoveArgs, base: &Path) -> Result<()> { + let mut had_errors = false; + + for identifier in &args.runs { + let run = match resolve_run(base, identifier) { + Ok(run) => run, + Err(e) => { + eprintln!("error: {identifier}: {e}"); + had_errors = true; + continue; + } + }; + + if run.status.is_active() && !args.force { + eprintln!( + "cannot remove active run {} (status: {}, use -f to force)", + short_run_id(&run.run_id), + run.status + ); + had_errors = true; + continue; + } + + // Transition status to Removing (best-effort) + write_run_status(&run.path, RunStatus::Removing, None); + + // Best-effort sandbox cleanup + let sandbox_path = run.path.join("sandbox.json"); + if let Ok(record) = crate::sandbox_record::SandboxRecord::load(&sandbox_path) { + if record.provider != "local" { + match super::cp::reconnect(&record).await { + Ok(sandbox) => { + if let Err(e) = sandbox.cleanup().await { + warn!(run_id = %run.run_id, error = %e, "sandbox cleanup failed"); + } + } + Err(e) => { + warn!(run_id = %run.run_id, error = %e, "sandbox reconnect failed"); + } + } + } + } + + // Delete the run directory + std::fs::remove_dir_all(&run.path) + .with_context(|| format!("failed to delete {}", run.path.display()))?; + eprintln!("{}", short_run_id(&run.run_id)); + } + + if had_errors { + bail!("some runs could not be removed"); + } + Ok(()) +} + #[cfg(test)] mod tests { use super::*; @@ -715,7 +939,7 @@ mod tests { dir_name: &str, manifest: Option, conclusion_json: Option, - pid_file: bool, + status: Option<(RunStatus, Option)>, ) -> PathBuf { let dir = base.join(dir_name); fs::create_dir_all(&dir).unwrap(); @@ -733,8 +957,8 @@ mod tests { ) .unwrap(); } - if pid_file { - fs::write(dir.join("run.pid"), "12345").unwrap(); + if let Some((s, r)) = status { + write_run_status(&dir, s, r); } dir } @@ -759,26 +983,23 @@ mod tests { Some( serde_json::json!({ "timestamp": "2026-01-01T12:01:00Z", "status": "success", "duration_ms": 60000 }), ), - false, + Some((RunStatus::Succeeded, Some(StatusReason::Completed))), ); - make_run_dir(base, "fabro-run-orphan", None, None, false); + make_run_dir(base, "fabro-run-orphan", None, None, None); let runs = scan_runs(base).unwrap(); assert_eq!(runs.len(), 2); let completed = runs.iter().find(|r| r.run_id == "abc123").unwrap(); assert_eq!(completed.workflow_name, "my-pipeline"); - assert_eq!( - completed.status, - RunStatus::Concluded(crate::outcome::StageStatus::Success) - ); + assert_eq!(completed.status, RunStatus::Succeeded); assert_eq!(completed.labels.get("env").unwrap(), "prod"); assert!(!completed.is_orphan); let orphan = runs.iter().find(|r| r.is_orphan).unwrap(); assert_eq!(orphan.workflow_name, "[no manifest]"); - assert_eq!(orphan.status, RunStatus::Unknown); + assert_eq!(orphan.status, RunStatus::Dead); } #[test] @@ -798,7 +1019,7 @@ mod tests { "edge_count": 0 })), None, - true, + Some((RunStatus::Running, None)), ); let runs = scan_runs(base).unwrap(); @@ -827,12 +1048,17 @@ mod tests { dir_name: "d1".into(), workflow_name: "p".into(), workflow_slug: None, - status: RunStatus::Concluded(crate::outcome::StageStatus::Success), + status: RunStatus::Succeeded, + status_reason: None, start_time: "2025-06-01T00:00:00Z".into(), labels: HashMap::new(), + duration_ms: None, + total_cost: None, + host_repo_path: None, start_time_dt: None, end_time: None, path: PathBuf::from("/tmp/d1"), + goal: String::new(), is_orphan: false, }, RunInfo { @@ -840,16 +1066,28 @@ mod tests { dir_name: "d2".into(), workflow_name: "p".into(), workflow_slug: None, - status: RunStatus::Concluded(crate::outcome::StageStatus::Success), + status: RunStatus::Succeeded, + status_reason: None, start_time: "2026-03-01T00:00:00Z".into(), labels: HashMap::new(), + duration_ms: None, + total_cost: None, + host_repo_path: None, start_time_dt: None, end_time: None, path: PathBuf::from("/tmp/d2"), + goal: String::new(), is_orphan: false, }, ]; - let filtered = filter_runs(&runs, Some("2026-01-01"), None, &[], false); + let filtered = filter_runs( + &runs, + Some("2026-01-01"), + None, + &[], + false, + StatusFilter::All, + ); assert_eq!(filtered.len(), 1); assert_eq!(filtered[0].run_id, "old"); } @@ -862,12 +1100,17 @@ mod tests { dir_name: "d1".into(), workflow_name: "deploy-prod".into(), workflow_slug: None, - status: RunStatus::Concluded(crate::outcome::StageStatus::Success), + status: RunStatus::Succeeded, + status_reason: None, start_time: "2026-01-01T00:00:00Z".into(), labels: HashMap::new(), + duration_ms: None, + total_cost: None, + host_repo_path: None, start_time_dt: None, end_time: None, path: PathBuf::from("/tmp/d1"), + goal: String::new(), is_orphan: false, }, RunInfo { @@ -875,16 +1118,21 @@ mod tests { dir_name: "d2".into(), workflow_name: "test-suite".into(), workflow_slug: None, - status: RunStatus::Concluded(crate::outcome::StageStatus::Success), + status: RunStatus::Succeeded, + status_reason: None, start_time: "2026-01-01T00:00:00Z".into(), labels: HashMap::new(), + duration_ms: None, + total_cost: None, + host_repo_path: None, start_time_dt: None, end_time: None, path: PathBuf::from("/tmp/d2"), + goal: String::new(), is_orphan: false, }, ]; - let filtered = filter_runs(&runs, None, Some("deploy"), &[], false); + let filtered = filter_runs(&runs, None, Some("deploy"), &[], false, StatusFilter::All); assert_eq!(filtered.len(), 1); assert_eq!(filtered[0].run_id, "a"); } @@ -897,12 +1145,17 @@ mod tests { dir_name: "d1".into(), workflow_name: "p".into(), workflow_slug: None, - status: RunStatus::Concluded(crate::outcome::StageStatus::Success), + status: RunStatus::Succeeded, + status_reason: None, start_time: "2026-01-01T00:00:00Z".into(), labels: HashMap::from([("env".into(), "prod".into())]), + duration_ms: None, + total_cost: None, + host_repo_path: None, start_time_dt: None, end_time: None, path: PathBuf::from("/tmp/d1"), + goal: String::new(), is_orphan: false, }, RunInfo { @@ -910,12 +1163,17 @@ mod tests { dir_name: "d2".into(), workflow_name: "p".into(), workflow_slug: None, - status: RunStatus::Concluded(crate::outcome::StageStatus::Success), + status: RunStatus::Succeeded, + status_reason: None, start_time: "2026-01-01T00:00:00Z".into(), labels: HashMap::from([("env".into(), "staging".into())]), + duration_ms: None, + total_cost: None, + host_repo_path: None, start_time_dt: None, end_time: None, path: PathBuf::from("/tmp/d2"), + goal: String::new(), is_orphan: false, }, ]; @@ -925,6 +1183,7 @@ mod tests { None, &[("env".to_string(), "prod".to_string())], false, + StatusFilter::All, ); assert_eq!(filtered.len(), 1); assert_eq!(filtered[0].run_id, "a"); @@ -937,21 +1196,104 @@ mod tests { dir_name: "d1".into(), workflow_name: "[no manifest]".into(), workflow_slug: None, - status: RunStatus::Unknown, + status: RunStatus::Dead, + status_reason: None, start_time: "".into(), labels: HashMap::new(), + duration_ms: None, + total_cost: None, + host_repo_path: None, start_time_dt: None, end_time: None, path: PathBuf::from("/tmp/d1"), + goal: String::new(), is_orphan: true, }]; - let filtered = filter_runs(&runs, None, None, &[], false); + let filtered = filter_runs(&runs, None, None, &[], false, StatusFilter::All); assert!(filtered.is_empty()); - let filtered = filter_runs(&runs, None, None, &[], true); + let filtered = filter_runs(&runs, None, None, &[], true, StatusFilter::All); assert_eq!(filtered.len(), 1); } + #[test] + fn filter_runs_running_only() { + let runs = vec![ + RunInfo { + run_id: "running-1".into(), + dir_name: "d1".into(), + workflow_name: "p".into(), + workflow_slug: None, + status: RunStatus::Running, + status_reason: None, + start_time: "2026-01-01T00:00:00Z".into(), + labels: HashMap::new(), + duration_ms: None, + total_cost: None, + host_repo_path: None, + start_time_dt: None, + end_time: None, + path: PathBuf::from("/tmp/d1"), + goal: String::new(), + is_orphan: false, + }, + RunInfo { + run_id: "done-1".into(), + dir_name: "d2".into(), + workflow_name: "p".into(), + workflow_slug: None, + status: RunStatus::Succeeded, + status_reason: None, + start_time: "2026-01-01T00:00:00Z".into(), + labels: HashMap::new(), + duration_ms: None, + total_cost: None, + host_repo_path: None, + start_time_dt: None, + end_time: None, + path: PathBuf::from("/tmp/d2"), + goal: String::new(), + is_orphan: false, + }, + ]; + + let filtered = filter_runs(&runs, None, None, &[], false, StatusFilter::RunningOnly); + assert_eq!(filtered.len(), 1); + assert_eq!(filtered[0].run_id, "running-1"); + + let filtered = filter_runs(&runs, None, None, &[], false, StatusFilter::All); + assert_eq!(filtered.len(), 2); + } + + #[test] + fn scan_runs_extracts_host_repo_path() { + let tmp = tempfile::tempdir().unwrap(); + let base = tmp.path(); + + make_run_dir( + base, + "20260101-ABC123", + Some(serde_json::json!({ + "run_id": "abc123", + "workflow_name": "my-pipeline", + "goal": "test goal", + "start_time": "2026-01-01T12:00:00Z", + "node_count": 2, + "edge_count": 1, + "host_repo_path": "/home/user/myproject" + })), + None, + Some((RunStatus::Running, None)), + ); + + let runs = scan_runs(base).unwrap(); + assert_eq!(runs.len(), 1); + assert_eq!( + runs[0].host_repo_path.as_deref(), + Some("/home/user/myproject") + ); + } + #[test] fn prune_dry_run_preserves_dirs() { let tmp = tempfile::tempdir().unwrap(); @@ -971,7 +1313,7 @@ mod tests { Some( serde_json::json!({ "timestamp": "2025-01-01T12:01:00Z", "status": "success", "duration_ms": 60000 }), ), - false, + None, ); let args = RunsPruneArgs { @@ -1008,7 +1350,7 @@ mod tests { Some( serde_json::json!({ "timestamp": "2025-01-01T12:01:00Z", "status": "success", "duration_ms": 60000 }), ), - false, + None, ); // Also add a run that should NOT be pruned (too new) @@ -1026,7 +1368,7 @@ mod tests { Some( serde_json::json!({ "timestamp": "2026-03-01T12:01:00Z", "status": "success", "duration_ms": 60000 }), ), - false, + None, ); let args = RunsPruneArgs { @@ -1053,7 +1395,7 @@ mod tests { let tmp = tempfile::tempdir().unwrap(); let base = tmp.path(); - let orphan_dir = make_run_dir(base, "orphan-dir", None, None, false); + let orphan_dir = make_run_dir(base, "orphan-dir", None, None, None); let args = RunsPruneArgs { filter: RunFilterArgs { @@ -1143,7 +1485,7 @@ mod tests { "edge_count": 0 })), None, - true, + Some((RunStatus::Running, None)), ); // Add a file to give it size fs::write( @@ -1169,7 +1511,7 @@ mod tests { "status": "success", "duration_ms": 60000 })), - false, + None, ); fs::write( runs_base.join("20260307-DONE").join("data.bin"), @@ -1292,7 +1634,7 @@ mod tests { "edge_count": 0 })), None, - false, + None, ); let info = resolve_run(dir.path(), "abc123").unwrap(); @@ -1314,7 +1656,7 @@ mod tests { "edge_count": 0 })), None, - false, + None, ); make_run_dir( dir.path(), @@ -1328,7 +1670,7 @@ mod tests { "edge_count": 0 })), None, - false, + None, ); // "deploy" doesn't match any run ID prefix, so falls back to workflow name @@ -1353,7 +1695,7 @@ mod tests { "edge_count": 0 })), None, - false, + None, ); make_run_dir( dir.path(), @@ -1367,7 +1709,7 @@ mod tests { "edge_count": 0 })), None, - false, + None, ); // "deploy" matches run_id prefix of first run — should prefer that over workflow name match @@ -1399,7 +1741,7 @@ mod tests { "edge_count": 0 })), None, - false, + None, ); make_run_dir( dir.path(), @@ -1413,7 +1755,7 @@ mod tests { "edge_count": 0 })), None, - false, + None, ); let result = resolve_run(dir.path(), "abc"); @@ -1439,7 +1781,7 @@ mod tests { "edge_count": 0 })), None, - false, + None, ); // Slug-style input should match PascalCase workflow name @@ -1462,7 +1804,7 @@ mod tests { "edge_count": 0 })), None, - false, + None, ); let info = resolve_run(dir.path(), "smoke").unwrap(); @@ -1486,7 +1828,7 @@ mod tests { "edge_count": 0 })), None, - false, + None, ); let info = resolve_run(dir.path(), "foo").unwrap(); @@ -1517,7 +1859,7 @@ mod tests { "status": "success", "duration_ms": 300000 })), - false, + Some((RunStatus::Succeeded, Some(StatusReason::Completed))), ); let runs = scan_runs(base).unwrap(); @@ -1547,7 +1889,7 @@ mod tests { "edge_count": 0 })), None, - true, + Some((RunStatus::Running, None)), ); let runs = scan_runs(base).unwrap(); @@ -1582,7 +1924,7 @@ mod tests { "status": "success", "duration_ms": 1000 })), - false, + None, ); let args = RunsPruneArgs { @@ -1618,7 +1960,7 @@ mod tests { "edge_count": 0 })), None, - true, // running + Some((RunStatus::Running, None)), // running ); let args = RunsPruneArgs { @@ -1658,7 +2000,7 @@ mod tests { "status": "success", "duration_ms": 1000 })), - false, + None, ); let args = RunsPruneArgs { @@ -1703,7 +2045,7 @@ mod tests { "status": "success", "duration_ms": 1000 })), - false, + None, ); let args = RunsPruneArgs { @@ -1751,18 +2093,20 @@ mod tests { } #[test] - fn run_status_serializes_as_flat_string() { - let concluded = RunStatus::Concluded(crate::outcome::StageStatus::Success); - assert_eq!(serde_json::to_string(&concluded).unwrap(), "\"success\""); - - let running = RunStatus::Running; - assert_eq!(serde_json::to_string(&running).unwrap(), "\"running\""); - - let unknown = RunStatus::Unknown; - assert_eq!(serde_json::to_string(&unknown).unwrap(), "\"unknown\""); - - let fail = RunStatus::Concluded(crate::outcome::StageStatus::Fail); - assert_eq!(serde_json::to_string(&fail).unwrap(), "\"fail\""); + fn run_status_serializes_as_snake_case_string() { + assert_eq!( + serde_json::to_string(&RunStatus::Succeeded).unwrap(), + "\"succeeded\"" + ); + assert_eq!( + serde_json::to_string(&RunStatus::Running).unwrap(), + "\"running\"" + ); + assert_eq!(serde_json::to_string(&RunStatus::Dead).unwrap(), "\"dead\""); + assert_eq!( + serde_json::to_string(&RunStatus::Failed).unwrap(), + "\"failed\"" + ); } // === Step 3: disk space reporting tests === @@ -1806,7 +2150,7 @@ mod tests { "status": "success", "duration_ms": 1000 })), - false, + None, ); // Run completed 1 day ago — should NOT be pruned with --older-than 2d @@ -1827,7 +2171,7 @@ mod tests { "status": "success", "duration_ms": 1000 })), - false, + None, ); // Run completed 10 days ago — should be pruned with --older-than 7d @@ -1848,7 +2192,7 @@ mod tests { "status": "success", "duration_ms": 1000 })), - false, + None, ); // With --older-than 7d, only the 10-day-old run should be pruned @@ -1877,4 +2221,239 @@ mod tests { "10-day-old run should be pruned with 7d threshold" ); } + + // === status.json tests === + + #[test] + fn read_status_from_status_json() { + let tmp = tempfile::tempdir().unwrap(); + let dir = tmp.path(); + write_run_status(dir, RunStatus::Running, None); + let si = read_status(dir); + assert_eq!(si.status, RunStatus::Running); + } + + #[test] + fn read_status_succeeded_with_conclusion_has_duration() { + let tmp = tempfile::tempdir().unwrap(); + let dir = tmp.path(); + write_run_status(dir, RunStatus::Succeeded, Some(StatusReason::Completed)); + fs::write( + dir.join("conclusion.json"), + serde_json::to_string_pretty(&serde_json::json!({ + "timestamp": "2026-01-01T12:01:00Z", + "status": "success", + "duration_ms": 60000 + })) + .unwrap(), + ) + .unwrap(); + let si = read_status(dir); + assert_eq!(si.status, RunStatus::Succeeded); + assert_eq!(si.duration_ms, Some(60000)); + } + + #[test] + fn read_status_no_status_json_is_dead() { + let tmp = tempfile::tempdir().unwrap(); + let si = read_status(tmp.path()); + assert_eq!(si.status, RunStatus::Dead); + } + + #[test] + fn scan_runs_status_json_without_manifest_is_not_orphan() { + let tmp = tempfile::tempdir().unwrap(); + let base = tmp.path(); + + let dir = base.join("20260301-STARTING"); + fs::create_dir_all(&dir).unwrap(); + fs::write(dir.join("id.txt"), "starting-run-id").unwrap(); + write_run_status( + &dir, + RunStatus::Starting, + Some(StatusReason::SandboxInitializing), + ); + + let runs = scan_runs(base).unwrap(); + assert_eq!(runs.len(), 1); + assert!(!runs[0].is_orphan); + assert_eq!(runs[0].status, RunStatus::Starting); + assert_eq!(runs[0].workflow_name, "[starting]"); + } + + #[test] + fn scan_runs_no_status_json_no_manifest_is_orphan() { + let tmp = tempfile::tempdir().unwrap(); + let base = tmp.path(); + + let dir = base.join("20260301-ORPHAN"); + fs::create_dir_all(&dir).unwrap(); + + let runs = scan_runs(base).unwrap(); + assert_eq!(runs.len(), 1); + assert!(runs[0].is_orphan); + assert_eq!(runs[0].status, RunStatus::Dead); + } + + #[test] + fn filter_runs_running_only_includes_starting_runs() { + let runs = vec![RunInfo { + run_id: "starting-1".into(), + dir_name: "d1".into(), + workflow_name: "[starting]".into(), + workflow_slug: None, + status: RunStatus::Starting, + status_reason: Some(StatusReason::SandboxInitializing), + start_time: "2026-01-01T00:00:00Z".into(), + labels: HashMap::new(), + duration_ms: None, + total_cost: None, + host_repo_path: None, + start_time_dt: None, + end_time: None, + path: PathBuf::from("/tmp/d1"), + goal: String::new(), + is_orphan: false, + }]; + + let filtered = filter_runs(&runs, None, None, &[], false, StatusFilter::RunningOnly); + assert_eq!(filtered.len(), 1); + assert_eq!(filtered[0].run_id, "starting-1"); + } + + #[test] + fn truncate_goal_short_string_unchanged() { + assert_eq!(truncate_goal("short", 50), "short"); + } + + #[test] + fn truncate_goal_long_string_truncated() { + let long = "a]".repeat(30); // 60 chars + let result = truncate_goal(&long, 50); + assert_eq!(result.chars().count(), 50); + assert!(result.ends_with("...")); + } + + #[test] + fn truncate_goal_multibyte_safe() { + let emoji_str = "Hello \u{1F600} world \u{1F600} test \u{1F600} more text here padding"; + // Should not panic + let result = truncate_goal(emoji_str, 15); + assert_eq!(result.chars().count(), 15); + assert!(result.ends_with("...")); + } + + fn make_succeeded_run(base: &Path, run_id: &str) -> PathBuf { + make_run_dir( + base, + &format!("20260101-{run_id}"), + Some(serde_json::json!({ + "run_id": run_id, + "workflow_name": "test-wf", + "goal": "test", + "start_time": "2026-01-01T12:00:00Z", + "node_count": 1, + "edge_count": 0, + "labels": {} + })), + None, + Some((RunStatus::Succeeded, Some(StatusReason::Completed))), + ) + } + + fn make_running_run(base: &Path, run_id: &str) -> PathBuf { + make_run_dir( + base, + &format!("20260101-{run_id}"), + Some(serde_json::json!({ + "run_id": run_id, + "workflow_name": "test-wf", + "goal": "test", + "start_time": "2026-01-01T12:00:00Z", + "node_count": 1, + "edge_count": 0, + "labels": {} + })), + None, + Some((RunStatus::Running, None)), + ) + } + + #[tokio::test] + async fn test_remove_succeeded_run() { + let tmp = tempfile::tempdir().unwrap(); + let base = tmp.path(); + let dir = make_succeeded_run(base, "RUN1"); + assert!(dir.exists()); + + let args = RunsRemoveArgs { + runs: vec!["RUN1".to_string()], + force: false, + }; + remove_from(&args, base).await.unwrap(); + assert!(!dir.exists()); + } + + #[tokio::test] + async fn test_remove_active_run_refused() { + let tmp = tempfile::tempdir().unwrap(); + let base = tmp.path(); + let dir = make_running_run(base, "ACTIVE1"); + assert!(dir.exists()); + + let args = RunsRemoveArgs { + runs: vec!["ACTIVE1".to_string()], + force: false, + }; + let result = remove_from(&args, base).await; + assert!(result.is_err()); + assert!(dir.exists()); + } + + #[tokio::test] + async fn test_remove_active_run_forced() { + let tmp = tempfile::tempdir().unwrap(); + let base = tmp.path(); + let dir = make_running_run(base, "ACTIVE2"); + assert!(dir.exists()); + + let args = RunsRemoveArgs { + runs: vec!["ACTIVE2".to_string()], + force: true, + }; + remove_from(&args, base).await.unwrap(); + assert!(!dir.exists()); + } + + #[tokio::test] + async fn test_remove_unknown_run() { + let tmp = tempfile::tempdir().unwrap(); + let base = tmp.path(); + fs::create_dir_all(base).unwrap(); + + let args = RunsRemoveArgs { + runs: vec!["NONEXISTENT".to_string()], + force: false, + }; + let result = remove_from(&args, base).await; + assert!(result.is_err()); + } + + #[tokio::test] + async fn test_remove_multiple_runs() { + let tmp = tempfile::tempdir().unwrap(); + let base = tmp.path(); + let dir1 = make_succeeded_run(base, "MULTI1"); + let dir2 = make_succeeded_run(base, "MULTI2"); + assert!(dir1.exists()); + assert!(dir2.exists()); + + let args = RunsRemoveArgs { + runs: vec!["MULTI1".to_string(), "MULTI2".to_string()], + force: false, + }; + remove_from(&args, base).await.unwrap(); + assert!(!dir1.exists()); + assert!(!dir2.exists()); + } } diff --git a/lib/crates/fabro-workflows/src/engine.rs b/lib/crates/fabro-workflows/src/engine.rs index a30d281a0..47e7da216 100644 --- a/lib/crates/fabro-workflows/src/engine.rs +++ b/lib/crates/fabro-workflows/src/engine.rs @@ -301,6 +301,10 @@ fn write_manifest(run_dir: &Path, graph: &Graph, config: &RunConfig) -> crate::m labels: config.labels.clone(), base_branch: config.base_branch.clone(), workflow_slug: config.workflow_slug.clone(), + host_repo_path: config + .host_repo_path + .as_ref() + .map(|p| p.to_string_lossy().to_string()), }; let _ = std::fs::create_dir_all(run_dir); let _ = manifest.save(&run_dir.join("manifest.json")); @@ -665,12 +669,12 @@ pub(crate) async fn git_push_host( refspec: &str, github_app: &Option, label: &str, -) { +) -> bool { let (origin_url, _) = match crate::daytona_sandbox::detect_repo_info(repo_path) { Ok(info) => info, Err(e) => { tracing::warn!(error = %e, label, "Cannot detect origin for push"); - return; + return false; } }; @@ -680,12 +684,12 @@ pub(crate) async fn git_push_host( Ok(url) => url, Err(e) => { tracing::warn!(error = %e, label, "Failed to get token for push"); - return; + return false; } }, None => { tracing::warn!(label, "No GitHub App credentials for push"); - return; + return false; } }; @@ -696,13 +700,19 @@ pub(crate) async fn git_push_host( }) .await; match result { - Ok(()) => tracing::info!(label, "Pushed to origin"), - Err(e) => tracing::warn!(error = %e, label, "Failed to push"), + Ok(()) => { + tracing::info!(label, "Pushed to origin"); + true + } + Err(e) => { + tracing::warn!(error = %e, label, "Failed to push"); + false + } } } /// Push the run branch to origin inside a remote sandbox (best-effort). -async fn git_push_remote(sandbox: &dyn Sandbox, branch: &str) { +async fn git_push_remote(sandbox: &dyn Sandbox, branch: &str) -> bool { if let Err(e) = sandbox.refresh_push_credentials().await { tracing::warn!(error = %e, "Failed to refresh push credentials"); } @@ -710,12 +720,15 @@ async fn git_push_remote(sandbox: &dyn Sandbox, branch: &str) { match sandbox.exec_command(&cmd, 60_000, None, None, None).await { Ok(r) if r.exit_code == 0 => { tracing::info!(branch, "Pushed run branch to origin"); + true } Ok(r) => { tracing::warn!(branch, exit_code = r.exit_code, "Failed to push run branch"); + false } Err(e) => { tracing::warn!(branch, error = %e, "Failed to push run branch"); + false } } } @@ -1238,6 +1251,11 @@ impl WorkflowRunEngine { // Write manifest.json (spec 5.6) let manifest = write_manifest(&config.run_dir, graph, config); + crate::cli::runs::write_run_status( + &config.run_dir, + crate::run_status::RunStatus::Running, + None, + ); // Initialize metadata branch for git-native checkpoint storage (best-effort) if let (Some(_), Some(ref repo_path)) = (&config.meta_branch, &config.host_repo_path) { @@ -1867,23 +1885,6 @@ impl WorkflowRunEngine { let checkpoint_path = config.run_dir.join("checkpoint.json"); if let Err(e) = checkpoint.save(&checkpoint_path) { context.append_log(format!("checkpoint save failed: {e}")); - } else { - self.services - .emitter - .emit(&WorkflowRunEvent::CheckpointSaved { - node_id: node.id.clone(), - }); - - // CheckpointSaved hook (non-blocking) - { - let mut hook_ctx = HookContext::new( - HookEvent::CheckpointSaved, - run_id.clone(), - graph.name.clone(), - ); - hook_ctx.node_id = Some(node.id.clone()); - let _ = self.run_hooks(&hook_ctx, hook_work_dir.as_deref()).await; - } } // Step 6b: Write shadow branch first, then run branch commit with trailer @@ -1949,18 +1950,22 @@ impl WorkflowRunEngine { } self.services .emitter - .emit(&WorkflowRunEvent::GitCheckpoint { - run_id: run_id.clone(), + .emit(&WorkflowRunEvent::CheckpointCompleted { node_id: node.id.clone(), status: outcome.status.to_string(), - git_commit_sha: sha.clone(), + git_commit_sha: Some(sha.clone()), }); + self.services.emitter.emit(&WorkflowRunEvent::GitCommit { + node_id: Some(node.id.clone()), + sha: sha.clone(), + }); + // Push run branch (skip in dry-run mode) if !config.dry_run { if let Some(ref branch) = config.run_branch { - if self.services.sandbox.is_remote() { - git_push_remote(&*self.services.sandbox, branch).await; + let push_ok = if self.services.sandbox.is_remote() { + git_push_remote(&*self.services.sandbox, branch).await } else if let Some(ref repo_path) = config.host_repo_path { let refspec = format!("refs/heads/{branch}"); git_push_host( @@ -1969,21 +1974,31 @@ impl WorkflowRunEngine { &config.github_app, "run branch", ) - .await; - } + .await + } else { + false + }; + self.services.emitter.emit(&WorkflowRunEvent::GitPush { + branch: branch.clone(), + success: push_ok, + }); } // Push metadata branch (always from host) if let (Some(ref meta_branch), Some(ref repo_path)) = (&config.meta_branch, &config.host_repo_path) { let refspec = format!("refs/heads/{meta_branch}"); - git_push_host( + let meta_push_ok = git_push_host( repo_path, &refspec, &config.github_app, "metadata branch", ) .await; + self.services.emitter.emit(&WorkflowRunEvent::GitPush { + branch: meta_branch.clone(), + success: meta_push_ok, + }); } } @@ -2010,7 +2025,7 @@ impl WorkflowRunEngine { Err(e) => { self.services .emitter - .emit(&WorkflowRunEvent::GitCheckpointFailed { + .emit(&WorkflowRunEvent::CheckpointFailed { node_id: node.id.clone(), error: e.clone(), }); @@ -2023,6 +2038,27 @@ impl WorkflowRunEngine { }); } } + } else { + // Non-git checkpoint path (start node or git disabled) + self.services + .emitter + .emit(&WorkflowRunEvent::CheckpointCompleted { + node_id: node.id.clone(), + status: outcome.status.to_string(), + git_commit_sha: None, + }); + } + + // CheckpointSaved hook (non-blocking) — fires for both git and non-git paths. + // The Err arm above returns early, so this only runs on success. + { + let mut hook_ctx = HookContext::new( + HookEvent::CheckpointSaved, + run_id.clone(), + graph.name.clone(), + ); + hook_ctx.node_id = Some(node.id.clone()); + let _ = self.run_hooks(&hook_ctx, hook_work_dir.as_deref()).await; } // Step 7: Follow selected edge (or direct jump) @@ -2130,13 +2166,26 @@ impl WorkflowRunEngine { None } }; + + let last_outcome = node_outcomes + .get(completed_nodes.last().unwrap_or(&String::new())) + .cloned() + .unwrap_or_else(Outcome::success); + + let run_usage: Option = node_outcomes + .values() + .filter_map(|o| o.usage.as_ref().map(fabro_llm::types::Usage::from)) + .reduce(|a, b| a + b); + self.services .emitter .emit(&WorkflowRunEvent::WorkflowRunCompleted { duration_ms, artifact_count: artifact_store.list().len(), + status: last_outcome.status.to_string(), total_cost, final_git_commit_sha: last_git_sha.clone(), + usage: run_usage, }); // RunComplete hook (non-blocking) @@ -2158,11 +2207,6 @@ impl WorkflowRunEngine { } } - // Return last outcome, or success if no outcomes recorded - let last_outcome = node_outcomes - .get(completed_nodes.last().unwrap_or(&String::new())) - .cloned() - .unwrap_or_else(Outcome::success); Ok((last_outcome, context)) } } @@ -2913,7 +2957,7 @@ mod tests { let collected = events.lock().unwrap(); // Should have: RunStarted, StageStarted (start), StageCompleted (start), - // CheckpointSaved, RunCompleted + // CheckpointCompleted, RunCompleted assert!(collected.len() >= 4); } @@ -5155,7 +5199,11 @@ mod tests { let git_checkpoint_node_ids: Vec<&str> = collected .iter() .filter_map(|e| match e { - WorkflowRunEvent::GitCheckpoint { node_id, .. } => Some(node_id.as_str()), + WorkflowRunEvent::CheckpointCompleted { + node_id, + git_commit_sha: Some(_), + .. + } => Some(node_id.as_str()), _ => None, }) .collect(); diff --git a/lib/crates/fabro-workflows/src/event.rs b/lib/crates/fabro-workflows/src/event.rs index 13d01a3e3..449216cb0 100644 --- a/lib/crates/fabro-workflows/src/event.rs +++ b/lib/crates/fabro-workflows/src/event.rs @@ -21,10 +21,14 @@ pub enum WorkflowRunEvent { WorkflowRunCompleted { duration_ms: u64, artifact_count: usize, + #[serde(default)] + status: String, #[serde(default, skip_serializing_if = "Option::is_none")] total_cost: Option, #[serde(default, skip_serializing_if = "Option::is_none")] final_git_commit_sha: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + usage: Option, }, WorkflowRunFailed { error: crate::error::FabroError, @@ -108,19 +112,43 @@ pub enum WorkflowRunEvent { stage: String, duration_ms: u64, }, - CheckpointSaved { - node_id: String, - }, - GitCheckpoint { - run_id: String, + CheckpointCompleted { node_id: String, status: String, - git_commit_sha: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + git_commit_sha: Option, }, - GitCheckpointFailed { + CheckpointFailed { node_id: String, error: String, }, + GitCommit { + #[serde(default, skip_serializing_if = "Option::is_none")] + node_id: Option, + sha: String, + }, + GitPush { + branch: String, + success: bool, + }, + GitBranch { + branch: String, + sha: String, + }, + GitWorktreeAdd { + path: String, + branch: String, + }, + GitWorktreeRemove { + path: String, + }, + GitFetch { + branch: String, + success: bool, + }, + GitReset { + sha: String, + }, EdgeSelected { from_node: String, to_node: String, @@ -272,6 +300,14 @@ pub enum WorkflowRunEvent { exit_code: i32, stderr: String, }, + RetroStarted, + RetroCompleted { + duration_ms: u64, + }, + RetroFailed { + error: String, + duration_ms: u64, + }, } impl WorkflowRunEvent { @@ -284,9 +320,13 @@ impl WorkflowRunEvent { Self::WorkflowRunCompleted { duration_ms, artifact_count, + status, .. } => { - info!(duration_ms, artifact_count, "Workflow run completed"); + info!( + duration_ms, + artifact_count, status, "Workflow run completed" + ); } Self::WorkflowRunFailed { error, duration_ms, .. @@ -428,19 +468,45 @@ impl WorkflowRunEvent { } => { warn!(stage, duration_ms, "Interview timeout"); } - Self::CheckpointSaved { node_id } => { - debug!(node_id, "Checkpoint saved"); - } - Self::GitCheckpoint { - run_id, - node_id, - status, - .. + Self::CheckpointCompleted { + node_id, status, .. } => { - debug!(run_id, node_id, status, "Git checkpoint"); + debug!(node_id, status, "Checkpoint completed"); } - Self::GitCheckpointFailed { node_id, error } => { - error!(node_id, error, "Git checkpoint commit failed"); + Self::CheckpointFailed { node_id, error } => { + error!(node_id, error, "Checkpoint failed"); + } + Self::GitCommit { node_id, sha } => { + debug!( + node_id = node_id.as_deref().unwrap_or(""), + sha, "Git commit" + ); + } + Self::GitPush { branch, success } => { + if *success { + debug!(branch, "Git push succeeded"); + } else { + warn!(branch, "Git push failed"); + } + } + Self::GitBranch { branch, sha } => { + debug!(branch, sha, "Git branch created"); + } + Self::GitWorktreeAdd { path, branch } => { + debug!(path, branch, "Git worktree added"); + } + Self::GitWorktreeRemove { path } => { + debug!(path, "Git worktree removed"); + } + Self::GitFetch { branch, success } => { + if *success { + debug!(branch, "Git fetch succeeded"); + } else { + warn!(branch, "Git fetch failed"); + } + } + Self::GitReset { sha } => { + debug!(sha, "Git reset"); } Self::EdgeSelected { from_node, @@ -656,6 +722,15 @@ impl WorkflowRunEvent { command, index, exit_code, "Devcontainer lifecycle command failed" ); } + Self::RetroStarted => { + info!("Retro started"); + } + Self::RetroCompleted { duration_ms } => { + info!(duration_ms, "Retro completed"); + } + Self::RetroFailed { error, duration_ms } => { + error!(error = %error, duration_ms, "Retro failed"); + } } } } @@ -905,11 +980,14 @@ fn rename_fields(event_name: &str, fields: &mut serde_json::Map Some(r.stdout.trim().to_string()), + Ok(r) if r.exit_code == 0 => { + let sha = r.stdout.trim().to_string(); + emitter.emit(&WorkflowRunEvent::GitCommit { + node_id: Some(setup.target_id.clone()), + sha: sha.clone(), + }); + Some(sha) + } _ => None, } } else { @@ -571,6 +589,9 @@ impl Handler for ParallelHandler { if let Some(ref wt_path) = result.worktree_path { let wt_str = wt_path.to_string_lossy().to_string(); crate::engine::git_remove_worktree(&*services.sandbox, &wt_str).await; + services + .emitter + .emit(&WorkflowRunEvent::GitWorktreeRemove { path: wt_str }); } } diff --git a/lib/crates/fabro-workflows/src/lib.rs b/lib/crates/fabro-workflows/src/lib.rs index 5d2f8d0ec..caa84e33d 100644 --- a/lib/crates/fabro-workflows/src/lib.rs +++ b/lib/crates/fabro-workflows/src/lib.rs @@ -49,6 +49,7 @@ pub mod preamble; pub mod pull_request; pub mod retro; pub mod retro_agent; +pub mod run_status; pub mod sandbox_record; pub mod stylesheet; pub mod transform; diff --git a/lib/crates/fabro-workflows/src/manifest.rs b/lib/crates/fabro-workflows/src/manifest.rs index 66df396fb..0f2035ef4 100644 --- a/lib/crates/fabro-workflows/src/manifest.rs +++ b/lib/crates/fabro-workflows/src/manifest.rs @@ -24,6 +24,8 @@ pub struct Manifest { pub base_branch: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub workflow_slug: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub host_repo_path: Option, } impl Manifest { @@ -53,6 +55,7 @@ mod tests { labels: HashMap::from([("env".into(), "test".into())]), base_branch: None, workflow_slug: None, + host_repo_path: None, } } @@ -121,5 +124,6 @@ mod tests { assert!(raw.get("run_branch").is_none()); assert!(raw.get("base_sha").is_none()); assert!(raw.get("workflow_slug").is_none()); + assert!(raw.get("host_repo_path").is_none()); } } diff --git a/lib/crates/fabro-workflows/src/outcome.rs b/lib/crates/fabro-workflows/src/outcome.rs index 3756c2721..7677fe888 100644 --- a/lib/crates/fabro-workflows/src/outcome.rs +++ b/lib/crates/fabro-workflows/src/outcome.rs @@ -61,6 +61,20 @@ pub struct StageUsage { pub cost: Option, } +impl From<&StageUsage> for fabro_llm::types::Usage { + fn from(u: &StageUsage) -> Self { + Self { + input_tokens: u.input_tokens, + output_tokens: u.output_tokens, + total_tokens: u.input_tokens + u.output_tokens, + cache_read_tokens: u.cache_read_tokens, + cache_write_tokens: u.cache_write_tokens, + reasoning_tokens: u.reasoning_tokens, + raw: None, + } + } +} + /// Structured failure information carried through the pipeline. #[derive(Debug, Clone, Serialize, Deserialize)] pub struct FailureDetail { diff --git a/lib/crates/fabro-workflows/src/pull_request.rs b/lib/crates/fabro-workflows/src/pull_request.rs index a8f0edf31..1cfa742b5 100644 --- a/lib/crates/fabro-workflows/src/pull_request.rs +++ b/lib/crates/fabro-workflows/src/pull_request.rs @@ -37,6 +37,10 @@ fn pr_title_from_goal(goal: &str) -> String { .strip_prefix("## ") .or_else(|| first_line.strip_prefix("# ")) .unwrap_or(first_line); + let first_line = first_line + .strip_prefix("Plan:") + .map(|s| s.trim()) + .unwrap_or(first_line); if first_line.chars().count() > 120 { let truncated: String = first_line.chars().take(119).collect(); format!("{truncated}…") @@ -816,6 +820,22 @@ mod tests { ); } + #[test] + fn pr_title_strips_plan_prefix() { + assert_eq!( + pr_title_from_goal("Plan: Add Draft PR Mode"), + "Add Draft PR Mode" + ); + } + + #[test] + fn pr_title_strips_heading_and_plan_prefix() { + assert_eq!( + pr_title_from_goal("## Plan: Add Draft PR Mode"), + "Add Draft PR Mode" + ); + } + #[test] fn pr_title_does_not_strip_h3_prefix() { assert_eq!( diff --git a/lib/crates/fabro-workflows/src/run_status.rs b/lib/crates/fabro-workflows/src/run_status.rs new file mode 100644 index 000000000..fc2a0abb2 --- /dev/null +++ b/lib/crates/fabro-workflows/src/run_status.rs @@ -0,0 +1,311 @@ +use std::fmt; +use std::path::Path; +use std::str::FromStr; + +use chrono::{DateTime, Utc}; +use serde::{Deserialize, Serialize}; + +/// Status of a workflow run in its lifecycle. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum RunStatus { + Submitted, + Starting, + Running, + Paused, + Removing, + Succeeded, + Failed, + Dead, +} + +impl RunStatus { + /// Terminal statuses cannot transition to anything (except Dead via `can_transition_to`). + pub fn is_terminal(self) -> bool { + matches!(self, Self::Succeeded | Self::Failed | Self::Dead) + } + + /// Active statuses represent runs that are in-progress or about to run. + pub fn is_active(self) -> bool { + matches!( + self, + Self::Submitted | Self::Starting | Self::Running | Self::Paused | Self::Removing + ) + } + + /// Check whether a transition from `self` to `to` is valid. + pub fn can_transition_to(self, to: Self) -> bool { + // Any state can transition to Dead + if to == Self::Dead { + return true; + } + // Terminal states cannot transition to anything else + if self.is_terminal() { + return false; + } + matches!( + (self, to), + (Self::Submitted, Self::Starting) + | (Self::Starting, Self::Running) + | (Self::Starting, Self::Failed) + | (Self::Running, Self::Succeeded) + | (Self::Running, Self::Failed) + | (Self::Running, Self::Paused) + | (Self::Running, Self::Removing) + | (Self::Paused, Self::Running) + | (Self::Paused, Self::Failed) + | (Self::Paused, Self::Removing) + | (Self::Removing, Self::Failed) + ) + } + + /// Attempt to transition from `self` to `to`. Returns an error if the transition is invalid. + pub fn transition_to(self, to: Self) -> Result { + if self.can_transition_to(to) { + Ok(to) + } else { + Err(InvalidTransition { from: self, to }) + } + } +} + +impl fmt::Display for RunStatus { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + let s = match self { + Self::Submitted => "submitted", + Self::Starting => "starting", + Self::Running => "running", + Self::Paused => "paused", + Self::Removing => "removing", + Self::Succeeded => "succeeded", + Self::Failed => "failed", + Self::Dead => "dead", + }; + f.write_str(s) + } +} + +impl FromStr for RunStatus { + type Err = ParseRunStatusError; + + fn from_str(s: &str) -> Result { + match s { + "submitted" => Ok(Self::Submitted), + "starting" => Ok(Self::Starting), + "running" => Ok(Self::Running), + "paused" => Ok(Self::Paused), + "removing" => Ok(Self::Removing), + "succeeded" => Ok(Self::Succeeded), + "failed" => Ok(Self::Failed), + "dead" => Ok(Self::Dead), + _ => Err(ParseRunStatusError(s.to_string())), + } + } +} + +#[derive(Debug, Clone)] +pub struct ParseRunStatusError(String); + +impl fmt::Display for ParseRunStatusError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "invalid run status: {:?}", self.0) + } +} + +impl std::error::Error for ParseRunStatusError {} + +#[derive(Debug, Clone, PartialEq)] +pub struct InvalidTransition { + pub from: RunStatus, + pub to: RunStatus, +} + +impl fmt::Display for InvalidTransition { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "invalid status transition: {} -> {}", self.from, self.to) + } +} + +impl std::error::Error for InvalidTransition {} + +/// Reason for the current status. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum StatusReason { + // Succeeded reasons + Completed, + PartialSuccess, + // Failed reasons + WorkflowError, + Cancelled, + Terminated, + TransientInfra, + BudgetExhausted, + SandboxInitFailed, + // Non-terminal reasons + SandboxInitializing, +} + +/// Persisted record of a run's status, written to `status.json`. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct RunStatusRecord { + pub status: RunStatus, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub reason: Option, + pub updated_at: DateTime, +} + +impl RunStatusRecord { + pub fn new(status: RunStatus, reason: Option) -> Self { + Self { + status, + reason, + updated_at: Utc::now(), + } + } + + pub fn save(&self, path: &Path) -> std::io::Result<()> { + let json = serde_json::to_string_pretty(self).map_err(std::io::Error::other)?; + std::fs::write(path, json) + } + + pub fn load(path: &Path) -> std::io::Result { + let data = std::fs::read_to_string(path)?; + serde_json::from_str(&data) + .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e)) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn terminal_states() { + assert!(RunStatus::Succeeded.is_terminal()); + assert!(RunStatus::Failed.is_terminal()); + assert!(RunStatus::Dead.is_terminal()); + assert!(!RunStatus::Running.is_terminal()); + assert!(!RunStatus::Submitted.is_terminal()); + assert!(!RunStatus::Starting.is_terminal()); + assert!(!RunStatus::Paused.is_terminal()); + assert!(!RunStatus::Removing.is_terminal()); + } + + #[test] + fn active_states() { + assert!(RunStatus::Submitted.is_active()); + assert!(RunStatus::Starting.is_active()); + assert!(RunStatus::Running.is_active()); + assert!(RunStatus::Paused.is_active()); + assert!(RunStatus::Removing.is_active()); + assert!(!RunStatus::Succeeded.is_active()); + assert!(!RunStatus::Failed.is_active()); + assert!(!RunStatus::Dead.is_active()); + } + + #[test] + fn valid_transitions() { + assert!(RunStatus::Submitted.can_transition_to(RunStatus::Starting)); + assert!(RunStatus::Starting.can_transition_to(RunStatus::Running)); + assert!(RunStatus::Starting.can_transition_to(RunStatus::Failed)); + assert!(RunStatus::Running.can_transition_to(RunStatus::Succeeded)); + assert!(RunStatus::Running.can_transition_to(RunStatus::Failed)); + assert!(RunStatus::Running.can_transition_to(RunStatus::Paused)); + assert!(RunStatus::Running.can_transition_to(RunStatus::Removing)); + assert!(RunStatus::Paused.can_transition_to(RunStatus::Running)); + assert!(RunStatus::Paused.can_transition_to(RunStatus::Failed)); + assert!(RunStatus::Paused.can_transition_to(RunStatus::Removing)); + assert!(RunStatus::Removing.can_transition_to(RunStatus::Failed)); + } + + #[test] + fn dead_reachable_from_any() { + let all = [ + RunStatus::Submitted, + RunStatus::Starting, + RunStatus::Running, + RunStatus::Paused, + RunStatus::Removing, + RunStatus::Succeeded, + RunStatus::Failed, + RunStatus::Dead, + ]; + for s in all { + assert!( + s.can_transition_to(RunStatus::Dead), + "{s} should be able to transition to Dead" + ); + } + } + + #[test] + fn invalid_transitions() { + assert!(!RunStatus::Submitted.can_transition_to(RunStatus::Running)); + assert!(!RunStatus::Submitted.can_transition_to(RunStatus::Succeeded)); + assert!(!RunStatus::Starting.can_transition_to(RunStatus::Succeeded)); + assert!(!RunStatus::Running.can_transition_to(RunStatus::Starting)); + assert!(!RunStatus::Succeeded.can_transition_to(RunStatus::Running)); + assert!(!RunStatus::Failed.can_transition_to(RunStatus::Running)); + } + + #[test] + fn transition_to_returns_result() { + assert_eq!( + RunStatus::Submitted.transition_to(RunStatus::Starting), + Ok(RunStatus::Starting) + ); + assert!(RunStatus::Failed.transition_to(RunStatus::Running).is_err()); + } + + #[test] + fn display_and_from_str_roundtrip() { + let all = [ + RunStatus::Submitted, + RunStatus::Starting, + RunStatus::Running, + RunStatus::Paused, + RunStatus::Removing, + RunStatus::Succeeded, + RunStatus::Failed, + RunStatus::Dead, + ]; + for s in all { + let text = s.to_string(); + let parsed: RunStatus = text.parse().unwrap(); + assert_eq!(s, parsed); + } + } + + #[test] + fn from_str_invalid() { + assert!("bogus".parse::().is_err()); + } + + #[test] + fn serde_roundtrip() { + let record = + RunStatusRecord::new(RunStatus::Running, Some(StatusReason::SandboxInitializing)); + let json = serde_json::to_string(&record).unwrap(); + let parsed: RunStatusRecord = serde_json::from_str(&json).unwrap(); + assert_eq!(parsed.status, RunStatus::Running); + assert_eq!(parsed.reason, Some(StatusReason::SandboxInitializing)); + } + + #[test] + fn save_and_load() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("status.json"); + let record = RunStatusRecord::new(RunStatus::Succeeded, Some(StatusReason::Completed)); + record.save(&path).unwrap(); + let loaded = RunStatusRecord::load(&path).unwrap(); + assert_eq!(loaded.status, RunStatus::Succeeded); + assert_eq!(loaded.reason, Some(StatusReason::Completed)); + } + + #[test] + fn load_missing_file() { + let result = RunStatusRecord::load(Path::new("/nonexistent/status.json")); + assert!(result.is_err()); + } +} diff --git a/lib/crates/fabro-workflows/tests/daytona_integration.rs b/lib/crates/fabro-workflows/tests/daytona_integration.rs index 15d67cd7d..948181bcc 100644 --- a/lib/crates/fabro-workflows/tests/daytona_integration.rs +++ b/lib/crates/fabro-workflows/tests/daytona_integration.rs @@ -613,19 +613,19 @@ async fn daytona_git_checkpoint_remote_emits_events() { .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); - // Assert GitCheckpoint events were emitted + // Assert CheckpointCompleted events with git SHAs were emitted { let events = events.lock().unwrap(); let git_events: Vec<_> = events .iter() .filter_map(|e| { - if let fabro_workflows::event::WorkflowRunEvent::GitCheckpoint { + if let fabro_workflows::event::WorkflowRunEvent::CheckpointCompleted { node_id, - git_commit_sha, + git_commit_sha: Some(sha), .. } = e { - Some((node_id.clone(), git_commit_sha.clone())) + Some((node_id.clone(), sha.clone())) } else { None } @@ -636,7 +636,7 @@ async fn daytona_git_checkpoint_remote_emits_events() { assert_eq!( git_events.len(), 1, - "expected 1 GitCheckpoint event (work node only), got {}", + "expected 1 CheckpointCompleted event with SHA (work node only), got {}", git_events.len() ); assert!( diff --git a/lib/crates/fabro-workflows/tests/integration.rs b/lib/crates/fabro-workflows/tests/integration.rs index 8f3a2e88d..db7af7329 100644 --- a/lib/crates/fabro-workflows/tests/integration.rs +++ b/lib/crates/fabro-workflows/tests/integration.rs @@ -1910,7 +1910,7 @@ async fn event_streaming_lifecycle() { .any(|e| matches!(e, WorkflowRunEvent::StageCompleted { name, .. } if name == "task"))); assert!(collected .iter() - .any(|e| matches!(e, WorkflowRunEvent::CheckpointSaved { .. }))); + .any(|e| matches!(e, WorkflowRunEvent::CheckpointCompleted { .. }))); assert!(collected .iter() .any(|e| matches!(e, WorkflowRunEvent::WorkflowRunCompleted { .. }))); @@ -10679,7 +10679,7 @@ impl Handler for FileWriterHandler { } } -/// End-to-end test: pipeline with git checkpointing enabled emits `GitCheckpoint` +/// End-to-end test: pipeline with git checkpointing enabled emits `CheckpointCompleted` /// events with valid commit SHAs and writes `diff.patch` per stage. #[tokio::test] async fn git_checkpoint_host_emits_events_and_diff_patch() { @@ -10796,18 +10796,18 @@ async fn git_checkpoint_host_emits_events_and_diff_patch() { .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); - // 6. Assert GitCheckpoint events were emitted + // 6. Assert CheckpointCompleted events with git SHAs were emitted let events = events.lock().unwrap(); let git_events: Vec<_> = events .iter() .filter_map(|e| { - if let WorkflowRunEvent::GitCheckpoint { + if let WorkflowRunEvent::CheckpointCompleted { node_id, - git_commit_sha, + git_commit_sha: Some(sha), .. } = e { - Some((node_id.clone(), git_commit_sha.clone())) + Some((node_id.clone(), sha.clone())) } else { None } @@ -10816,7 +10816,7 @@ async fn git_checkpoint_host_emits_events_and_diff_patch() { // work node gets a checkpoint commit (start is skipped, exit is terminal) assert!( !git_events.is_empty(), - "expected at least 1 GitCheckpoint event, got {}", + "expected at least 1 CheckpointCompleted event with SHA, got {}", git_events.len() ); assert!( diff --git a/test/docs/run_tests.sh b/test/docs/run_tests.sh index 406b3f98d..7e772b7fb 100755 --- a/test/docs/run_tests.sh +++ b/test/docs/run_tests.sh @@ -12,7 +12,8 @@ PARALLEL="${PARALLEL:-1}" [[ "$VERBOSE" == "1" ]] && PARALLEL=1 RESULTS_DIR="$(mktemp -d)" -trap 'rm -rf "$RESULTS_DIR"' EXIT +RUNS_DIR="$(mktemp -d)" +trap 'rm -rf "$RESULTS_DIR" "$RUNS_DIR"' EXIT # Capture command output to log file; when VERBOSE=1, also stream to terminal. capture() { @@ -77,6 +78,7 @@ run_one() { local flags=(--auto-approve) [[ "$PHASE" == "dry-run" ]] && flags+=(--dry-run) [[ "$PHASE" == "haiku" ]] && flags+=(--model claude-haiku-4-5) + [[ "$PHASE" != "dry-run" ]] && flags+=(--run-dir "$RUNS_DIR/$(echo "$rel" | tr '/' '_')") if (cd "$dot_dir" && capture "$result_file.log" "$ARC" run start "$target" "${flags[@]}"); then echo "PASS" > "$result_file"