Merge main into fabro/run/01KKS6PG9929P116A2RRKXC738

Resolve conflict in engine.rs: keep PR's simplified refspec
(meta_branch is already `fabro/meta/{run_id}`) while incorporating
main's git_push_host return value and GitPush event emission.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-03-15 17:32:57 -04:00
commit 05f73e1804
37 changed files with 2257 additions and 564 deletions

22
Cargo.lock generated
View file

@ -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"

View file

@ -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

View file

@ -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 |

View file

@ -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:**

View file

@ -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

View file

@ -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)?;

View file

@ -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 

```

View file

@ -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 

```

View file

@ -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 

```

View file

@ -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 

```

View file

@ -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 

```

View file

@ -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 

```

View file

@ -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

View file

@ -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;

View file

@ -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

View file

@ -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<f64>) -> 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<Color> {
if use_color {
Some(color)
} else {
None
}
}
fn model_row(model: &crate::types::ModelInfo, use_color: bool) -> Vec<CellStruct> {
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<CellStruct> {
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<Vec<CellStruct>> = 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<String> {
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<CellStruct>> = 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<CellStruct>> = 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");
}

View file

@ -1,4 +1,5 @@
pub mod check_report;
pub mod path;
pub mod redact;
pub mod run_log;
pub mod telemetry;

View file

@ -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"));
}
}

View file

@ -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

View file

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

View file

@ -22,7 +22,7 @@ pub struct LogsArgs {
#[arg(short = 'n', long)]
pub tail: Option<usize>,
/// 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]

View file

@ -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<Color>`.
pub(crate) fn color_if(use_color: bool, color: cli_table::Color) -> Option<cli_table::Color> {
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() {

View file

@ -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);
}
_ => {}
}
}

View file

@ -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<Vec<CellStruct>> = 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.

View file

@ -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<Mutex<Option<String>>> = 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<Mutex<progress::ProgressUI>>>,
emitter: Option<Arc<EventEmitter>>,
) {
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);

File diff suppressed because it is too large Load diff

View file

@ -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<fabro_github::GitHubAppCredentials>,
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<fabro_llm::types::Usage> = 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();

View file

@ -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<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
final_git_commit_sha: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
usage: Option<fabro_llm::types::Usage>,
},
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<String>,
},
GitCheckpointFailed {
CheckpointFailed {
node_id: String,
error: String,
},
GitCommit {
#[serde(default, skip_serializing_if = "Option::is_none")]
node_id: Option<String>,
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<String, serde_js
default_node_label(fields);
rename(fields, "start_node", "start_node_id");
} else if event_name == "SubgraphCompleted"
|| event_name == "CheckpointSaved"
|| event_name == "GitCheckpoint"
|| event_name == "GitCheckpointFailed"
|| event_name == "CheckpointCompleted"
|| event_name == "CheckpointFailed"
{
default_node_label(fields);
} else if event_name == "GitCommit" {
if fields.contains_key("node_id") {
default_node_label(fields);
}
} else if event_name.starts_with("DevcontainerLifecycleCommand")
|| event_name == "DevcontainerLifecycleFailed"
{
@ -1832,24 +1910,43 @@ mod tests {
}
#[test]
fn rename_fields_checkpoint_saved() {
let event = WorkflowRunEvent::CheckpointSaved {
node_id: "plan".to_string(),
fn rename_fields_checkpoint_completed() {
let event = WorkflowRunEvent::CheckpointCompleted {
node_id: "work".to_string(),
status: "success".to_string(),
git_commit_sha: Some("abc123".to_string()),
};
let (name, fields) = flatten_event(&event);
assert_eq!(name, "CheckpointSaved");
assert_eq!(fields["node_id"], "plan");
assert_eq!(fields["node_label"], "plan");
assert_eq!(name, "CheckpointCompleted");
assert_eq!(fields["node_id"], "work");
assert_eq!(fields["node_label"], "work");
// Without git_commit_sha
let event_no_git = WorkflowRunEvent::CheckpointCompleted {
node_id: "plan".to_string(),
status: "success".to_string(),
git_commit_sha: None,
};
let json = serde_json::to_string(&event_no_git).unwrap();
assert!(!json.contains("git_commit_sha"));
let deserialized: WorkflowRunEvent = serde_json::from_str(&json).unwrap();
assert!(matches!(
deserialized,
WorkflowRunEvent::CheckpointCompleted {
git_commit_sha: None,
..
}
));
}
#[test]
fn rename_fields_git_checkpoint_failed() {
let event = WorkflowRunEvent::GitCheckpointFailed {
fn rename_fields_checkpoint_failed() {
let event = WorkflowRunEvent::CheckpointFailed {
node_id: "fix_lints".to_string(),
error: "git add failed (exit 1): fatal: not a git repository".to_string(),
};
let (name, fields) = flatten_event(&event);
assert_eq!(name, "GitCheckpointFailed");
assert_eq!(name, "CheckpointFailed");
assert_eq!(fields["node_id"], "fix_lints");
assert_eq!(fields["node_label"], "fix_lints");
assert_eq!(
@ -1858,6 +1955,136 @@ mod tests {
);
}
#[test]
fn rename_fields_git_commit_with_node_id() {
let event = WorkflowRunEvent::GitCommit {
node_id: Some("work".to_string()),
sha: "abc123".to_string(),
};
let (name, fields) = flatten_event(&event);
assert_eq!(name, "GitCommit");
assert_eq!(fields["node_id"], "work");
assert_eq!(fields["node_label"], "work");
assert_eq!(fields["sha"], "abc123");
}
#[test]
fn rename_fields_git_commit_without_node_id() {
let event = WorkflowRunEvent::GitCommit {
node_id: None,
sha: "abc123".to_string(),
};
let (name, fields) = flatten_event(&event);
assert_eq!(name, "GitCommit");
assert!(!fields.contains_key("node_label"));
assert_eq!(fields["sha"], "abc123");
}
#[test]
fn git_commit_serialization() {
let event = WorkflowRunEvent::GitCommit {
node_id: Some("work".to_string()),
sha: "abc123".to_string(),
};
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("GitCommit"));
assert!(json.contains("\"sha\":\"abc123\""));
let deserialized: WorkflowRunEvent = serde_json::from_str(&json).unwrap();
assert!(matches!(deserialized, WorkflowRunEvent::GitCommit { sha, .. } if sha == "abc123"));
// node_id None is omitted
let event_none = WorkflowRunEvent::GitCommit {
node_id: None,
sha: "def456".to_string(),
};
let json_none = serde_json::to_string(&event_none).unwrap();
assert!(!json_none.contains("node_id"));
}
#[test]
fn git_push_serialization() {
let event = WorkflowRunEvent::GitPush {
branch: "fabro/run/123".to_string(),
success: true,
};
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("GitPush"));
assert!(json.contains("\"success\":true"));
let deserialized: WorkflowRunEvent = serde_json::from_str(&json).unwrap();
assert!(matches!(
deserialized,
WorkflowRunEvent::GitPush { success: true, .. }
));
}
#[test]
fn git_branch_serialization() {
let event = WorkflowRunEvent::GitBranch {
branch: "fabro/run/123/work".to_string(),
sha: "abc123".to_string(),
};
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("GitBranch"));
let deserialized: WorkflowRunEvent = serde_json::from_str(&json).unwrap();
assert!(
matches!(deserialized, WorkflowRunEvent::GitBranch { branch, .. } if branch == "fabro/run/123/work")
);
}
#[test]
fn git_worktree_add_serialization() {
let event = WorkflowRunEvent::GitWorktreeAdd {
path: "/tmp/wt".to_string(),
branch: "work".to_string(),
};
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("GitWorktreeAdd"));
let deserialized: WorkflowRunEvent = serde_json::from_str(&json).unwrap();
assert!(
matches!(deserialized, WorkflowRunEvent::GitWorktreeAdd { path, .. } if path == "/tmp/wt")
);
}
#[test]
fn git_worktree_remove_serialization() {
let event = WorkflowRunEvent::GitWorktreeRemove {
path: "/tmp/wt".to_string(),
};
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("GitWorktreeRemove"));
let deserialized: WorkflowRunEvent = serde_json::from_str(&json).unwrap();
assert!(
matches!(deserialized, WorkflowRunEvent::GitWorktreeRemove { path } if path == "/tmp/wt")
);
}
#[test]
fn git_fetch_serialization() {
let event = WorkflowRunEvent::GitFetch {
branch: "main".to_string(),
success: false,
};
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("GitFetch"));
assert!(json.contains("\"success\":false"));
let deserialized: WorkflowRunEvent = serde_json::from_str(&json).unwrap();
assert!(matches!(
deserialized,
WorkflowRunEvent::GitFetch { success: false, .. }
));
}
#[test]
fn git_reset_serialization() {
let event = WorkflowRunEvent::GitReset {
sha: "abc123".to_string(),
};
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("GitReset"));
let deserialized: WorkflowRunEvent = serde_json::from_str(&json).unwrap();
assert!(matches!(deserialized, WorkflowRunEvent::GitReset { sha } if sha == "abc123"));
}
#[test]
fn rename_fields_sandbox_snapshot_pulling() {
let event = WorkflowRunEvent::Sandbox {
@ -1970,6 +2197,51 @@ mod tests {
);
}
#[test]
fn workflow_run_completed_serialization_with_status_and_usage() {
let event = WorkflowRunEvent::WorkflowRunCompleted {
duration_ms: 30000,
artifact_count: 2,
status: "success".to_string(),
total_cost: Some(1.23),
final_git_commit_sha: Some("abc123".to_string()),
usage: Some(Usage {
input_tokens: 5000,
output_tokens: 2000,
total_tokens: 7000,
cache_read_tokens: Some(3000),
cache_write_tokens: Some(500),
reasoning_tokens: Some(800),
raw: None,
}),
};
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"status\":\"success\""));
assert!(json.contains("\"total_tokens\":7000"));
assert!(json.contains("\"cache_read_tokens\":3000"));
assert!(json.contains("\"reasoning_tokens\":800"));
let deserialized: WorkflowRunEvent = serde_json::from_str(&json).unwrap();
assert!(matches!(
deserialized,
WorkflowRunEvent::WorkflowRunCompleted { status, usage: Some(u), .. }
if status == "success" && u.total_tokens == 7000
));
}
#[test]
fn workflow_run_completed_backward_compat_without_new_fields() {
// Old JSONL without status/usage should deserialize with defaults
let json =
r#"{"WorkflowRunCompleted":{"duration_ms":5000,"artifact_count":1,"total_cost":0.25}}"#;
let deserialized: WorkflowRunEvent = serde_json::from_str(json).unwrap();
assert!(matches!(
deserialized,
WorkflowRunEvent::WorkflowRunCompleted { status, usage, .. }
if status.is_empty() && usage.is_none()
));
}
#[test]
fn devcontainer_resolved_serializes() {
let event = WorkflowRunEvent::DevcontainerResolved {
@ -2060,4 +2332,78 @@ mod tests {
assert_eq!(fields["command_index"], 1);
assert!(!fields.contains_key("index"));
}
#[test]
fn retro_started_event_serialization() {
let event = WorkflowRunEvent::RetroStarted;
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("RetroStarted"));
let deserialized: WorkflowRunEvent = serde_json::from_str(&json).unwrap();
assert!(matches!(deserialized, WorkflowRunEvent::RetroStarted));
}
#[test]
fn retro_completed_event_serialization() {
let event = WorkflowRunEvent::RetroCompleted { duration_ms: 5000 };
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"duration_ms\":5000"));
let deserialized: WorkflowRunEvent = serde_json::from_str(&json).unwrap();
assert!(matches!(
deserialized,
WorkflowRunEvent::RetroCompleted { duration_ms: 5000 }
));
}
#[test]
fn retro_failed_event_serialization() {
let event = WorkflowRunEvent::RetroFailed {
error: "LLM timeout".to_string(),
duration_ms: 3000,
};
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("LLM timeout"));
assert!(json.contains("\"duration_ms\":3000"));
let deserialized: WorkflowRunEvent = serde_json::from_str(&json).unwrap();
assert!(matches!(deserialized, WorkflowRunEvent::RetroFailed { .. }));
}
#[test]
fn flatten_retro_started() {
let event = WorkflowRunEvent::RetroStarted;
let (name, _fields) = flatten_event(&event);
assert_eq!(name, "RetroStarted");
}
#[test]
fn flatten_retro_failed() {
let event = WorkflowRunEvent::RetroFailed {
error: "timeout".to_string(),
duration_ms: 1000,
};
let (name, fields) = flatten_event(&event);
assert_eq!(name, "RetroFailed");
assert_eq!(fields["error"], "timeout");
assert_eq!(fields["duration_ms"], 1000);
}
#[test]
fn emitter_captures_retro_events() {
let mut emitter = EventEmitter::new();
let received = Arc::new(Mutex::new(Vec::new()));
let r = Arc::clone(&received);
emitter.on_event(move |event| {
if let WorkflowRunEvent::RetroStarted = event {
r.lock().unwrap().push("started".to_string());
}
if let WorkflowRunEvent::RetroCompleted { .. } = event {
r.lock().unwrap().push("completed".to_string());
}
});
emitter.emit(&WorkflowRunEvent::RetroStarted);
emitter.emit(&WorkflowRunEvent::RetroCompleted { duration_ms: 100 });
let events = received.lock().unwrap();
assert_eq!(events.len(), 2);
assert_eq!(events[0], "started");
assert_eq!(events[1], "completed");
}
}

View file

@ -329,6 +329,10 @@ impl Handler for ParallelHandler {
"failed to create branch {branch_name}"
)));
}
services.emitter.emit(&WorkflowRunEvent::GitBranch {
branch: branch_name.clone(),
sha: bsha.clone(),
});
if !crate::engine::git_replace_worktree(
&*services.sandbox,
&wt_path_str,
@ -340,6 +344,10 @@ impl Handler for ParallelHandler {
"failed to add worktree {wt_path_str}"
)));
}
services.emitter.emit(&WorkflowRunEvent::GitWorktreeAdd {
path: wt_path_str.clone(),
branch: branch_name.clone(),
});
let reset_cmd = format!("{} reset --hard {bsha}", crate::engine::GIT_REMOTE);
let reset_result = services
.sandbox
@ -350,6 +358,9 @@ impl Handler for ParallelHandler {
"failed to reset worktree {wt_path_str}"
)));
}
services
.emitter
.emit(&WorkflowRunEvent::GitReset { sha: bsha.clone() });
branch_context.set(keys::INTERNAL_WORK_DIR, serde_json::json!(&wt_path_str));
@ -476,7 +487,14 @@ impl Handler for ParallelHandler {
.exec_command(&sha_cmd, 10_000, None, None, None)
.await;
match sha_result {
Ok(r) if r.exit_code == 0 => 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 });
}
}

View file

@ -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;

View file

@ -24,6 +24,8 @@ pub struct Manifest {
pub base_branch: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub workflow_slug: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub host_repo_path: Option<String>,
}
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());
}
}

View file

@ -61,6 +61,20 @@ pub struct StageUsage {
pub cost: Option<f64>,
}
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 {

View file

@ -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!(

View file

@ -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<Self, InvalidTransition> {
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<Self, Self::Err> {
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<StatusReason>,
pub updated_at: DateTime<Utc>,
}
impl RunStatusRecord {
pub fn new(status: RunStatus, reason: Option<StatusReason>) -> 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<Self> {
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::<RunStatus>().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());
}
}

View file

@ -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!(

View file

@ -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!(

View file

@ -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"