refactor(workflow): simplify fork rewind cleanup

Extract shared checkpoint CLI helpers and run-store error mapping, and remove redundant rewind/fork state from the workflow operations.
This commit is contained in:
Bryan Helmkamp 2026-04-24 12:13:09 -04:00
parent 4215ed3c16
commit b2f8e0fb75
No known key found for this signature in database
11 changed files with 175 additions and 178 deletions

View file

@ -0,0 +1,118 @@
use anyhow::Result;
use cli_table::format::{Border, Separator};
use cli_table::{Cell, CellStruct, Color, Style, Table};
use fabro_api::types::TimelineEntryResponse;
use fabro_types::RunId;
use fabro_util::printer::Printer;
use fabro_util::terminal::Styles;
use git2::Repository;
use serde::Serialize;
use crate::server_client::Client;
use crate::shared::color_if;
use crate::shared::repo::ensure_matching_repo_origin;
#[derive(Serialize)]
pub(crate) struct TimelineEntryJson {
ordinal: usize,
node_name: String,
visit: usize,
run_commit_sha: Option<String>,
}
pub(crate) async fn ensure_origin_if_local(
client: &Client,
run_id: &RunId,
verb: &str,
) -> Result<()> {
if Repository::discover(".").is_err() {
return Ok(());
}
let state = client.get_run_state(run_id).await?;
if let Some(run_spec) = state.spec {
ensure_matching_repo_origin(run_spec.repo_origin_url.as_deref(), verb)?;
}
Ok(())
}
pub(crate) fn timeline_entries_json(entries: &[TimelineEntryResponse]) -> Vec<TimelineEntryJson> {
entries
.iter()
.map(|entry| TimelineEntryJson {
ordinal: usize::try_from(entry.ordinal.get())
.expect("timeline ordinal should fit in usize"),
node_name: entry.node_name.clone(),
visit: usize::try_from(entry.visit.get())
.expect("timeline visit should fit in usize"),
run_commit_sha: entry.run_commit_sha.clone(),
})
.collect()
}
pub(crate) fn short_id(run_id: &str) -> &str {
&run_id[..8.min(run_id.len())]
}
pub(crate) fn print_timeline(entries: &[TimelineEntryJson], styles: &Styles, printer: Printer) {
if entries.is_empty() {
fabro_util::printerr!(printer, "No checkpoints found.");
return;
}
let use_color = styles.use_color;
let title = vec![
"@".cell().bold(use_color),
"Node".cell().bold(use_color),
"Details".cell().bold(use_color),
];
let rows: Vec<Vec<CellStruct>> = entries
.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 entry.run_commit_sha.is_none() {
details.push("no run commit".to_string());
}
let detail_str = if details.is_empty() {
String::new()
} else {
format!("({})", details.join(", "))
};
vec![
ordinal_str
.cell()
.foreground_color(color_if(use_color, Color::Cyan)),
entry.node_name.clone().cell(),
detail_str
.cell()
.foreground_color(color_if(use_color, Color::Ansi256(8))),
]
})
.collect();
let color_choice = if use_color {
cli_table::ColorChoice::Auto
} else {
cli_table::ColorChoice::Never
};
let table = rows
.table()
.title(title)
.color_choice(color_choice)
.border(Border::builder().build())
.separator(Separator::builder().build());
#[allow(
clippy::print_stderr,
reason = "The checkpoint timeline table is operator feedback, not command output."
)]
if let Ok(display) = table.display() {
eprintln!("{display}");
}
}

View file

@ -11,16 +11,16 @@ pub(crate) async fn run(args: &ForkArgs, styles: &Styles, base_ctx: &CommandCont
let ctx = base_ctx.with_target(&args.server)?;
let client = ctx.server().await?;
let run_id = client.resolve_run(&args.run_id).await?.run_id;
super::rewind::ensure_origin_if_local(client.as_ref(), &run_id, "fork").await?;
super::checkpoints::ensure_origin_if_local(client.as_ref(), &run_id, "fork").await?;
if args.list {
let timeline = client.run_timeline(&run_id).await?;
if ctx.json_output() {
print_json_pretty(&super::rewind::timeline_entries_json(&timeline))?;
print_json_pretty(&super::checkpoints::timeline_entries_json(&timeline))?;
return Ok(());
}
let entries = super::rewind::timeline_entries_json(&timeline);
super::rewind::print_timeline(&entries, styles, printer);
let entries = super::checkpoints::timeline_entries_json(&timeline);
super::checkpoints::print_timeline(&entries, styles, printer);
return Ok(());
}
@ -32,22 +32,18 @@ pub(crate) async fn run(args: &ForkArgs, styles: &Styles, base_ctx: &CommandCont
.await?;
if ctx.json_output() {
print_json_pretty(&serde_json::json!({
"source_run_id": response.source_run_id,
"new_run_id": response.new_run_id,
"target": response.target,
}))?;
print_json_pretty(&response)?;
} else {
fabro_util::printerr!(
printer,
"\nForked run {} -> {}",
super::rewind::short_id(&response.source_run_id),
super::rewind::short_id(&response.new_run_id)
super::checkpoints::short_id(&response.source_run_id),
super::checkpoints::short_id(&response.new_run_id)
);
fabro_util::printerr!(
printer,
"To resume: fabro resume {}",
super::rewind::short_id(&response.new_run_id)
super::checkpoints::short_id(&response.new_run_id)
);
}

View file

@ -8,6 +8,7 @@ use crate::shared::print_json_pretty;
use crate::sleep_inhibitor;
pub(crate) mod attach;
pub(crate) mod checkpoints;
pub(crate) mod command;
pub(crate) mod cp;
pub(crate) mod create;

View file

@ -1,26 +1,11 @@
use anyhow::Result;
use cli_table::format::{Border, Separator};
use cli_table::{Cell, CellStruct, Color, Style, Table};
use fabro_api::types::{RewindRequest, TimelineEntryResponse};
use fabro_types::RunId;
use fabro_util::printer::Printer;
use fabro_api::types::RewindRequest;
use fabro_util::terminal::Styles;
use git2::Repository;
use serde::Serialize;
use super::checkpoints::{ensure_origin_if_local, print_timeline, short_id, timeline_entries_json};
use crate::args::RewindArgs;
use crate::command_context::CommandContext;
use crate::server_client::Client;
use crate::shared::repo::ensure_matching_repo_origin;
use crate::shared::{color_if, print_json_pretty};
#[derive(Serialize)]
pub(crate) struct TimelineEntryJson {
ordinal: usize,
node_name: String,
visit: usize,
run_commit_sha: Option<String>,
}
use crate::shared::print_json_pretty;
pub(crate) async fn run(
args: &RewindArgs,
@ -88,100 +73,3 @@ pub(crate) async fn run(
Ok(())
}
pub(crate) async fn ensure_origin_if_local(
client: &Client,
run_id: &RunId,
verb: &str,
) -> Result<()> {
if Repository::discover(".").is_err() {
return Ok(());
}
let state = client.get_run_state(run_id).await?;
if let Some(run_spec) = state.spec {
ensure_matching_repo_origin(run_spec.repo_origin_url.as_deref(), verb)?;
}
Ok(())
}
pub(crate) fn timeline_entries_json(entries: &[TimelineEntryResponse]) -> Vec<TimelineEntryJson> {
entries
.iter()
.map(|entry| TimelineEntryJson {
ordinal: usize::try_from(entry.ordinal.get())
.expect("timeline ordinal should fit in usize"),
node_name: entry.node_name.clone(),
visit: usize::try_from(entry.visit.get())
.expect("timeline visit should fit in usize"),
run_commit_sha: entry.run_commit_sha.clone(),
})
.collect()
}
pub(crate) fn short_id(run_id: &str) -> &str {
&run_id[..8.min(run_id.len())]
}
pub(crate) fn print_timeline(entries: &[TimelineEntryJson], styles: &Styles, printer: Printer) {
if entries.is_empty() {
fabro_util::printerr!(printer, "No checkpoints found.");
return;
}
let use_color = styles.use_color;
let title = vec![
"@".cell().bold(use_color),
"Node".cell().bold(use_color),
"Details".cell().bold(use_color),
];
let rows: Vec<Vec<CellStruct>> = entries
.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 entry.run_commit_sha.is_none() {
details.push("no run commit".to_string());
}
let detail_str = if details.is_empty() {
String::new()
} else {
format!("({})", details.join(", "))
};
vec![
ordinal_str
.cell()
.foreground_color(color_if(use_color, Color::Cyan)),
entry.node_name.clone().cell(),
detail_str
.cell()
.foreground_color(color_if(use_color, Color::Ansi256(8))),
]
})
.collect();
let color_choice = if use_color {
cli_table::ColorChoice::Auto
} else {
cli_table::ColorChoice::Never
};
let table = rows
.table()
.title(title)
.color_choice(color_choice)
.border(Border::builder().build())
.separator(Separator::builder().build());
#[allow(
clippy::print_stderr,
reason = "The rewind preview table is operator feedback, not command output."
)]
if let Ok(display) = table.display() {
eprintln!("{display}");
}
}

View file

@ -7121,14 +7121,13 @@ async fn rewind_run(
source_run_id,
new_run_id,
target,
archived,
}) => (
StatusCode::OK,
Json(RewindResponse {
source_run_id: source_run_id.to_string(),
new_run_id: new_run_id.to_string(),
target: target.response_target(),
archived,
new_run_id: new_run_id.to_string(),
target: target.response_target(),
archived: true,
archive_error: None,
}),
)

View file

@ -1,16 +1,10 @@
use fabro_store::{Database, Error as StoreError};
use fabro_store::Database;
use fabro_types::{ActorRef, RunId, RunStatus, TerminalStatus};
use super::run_store::map_open_run_error;
use crate::error::Error;
use crate::event::{self, Event};
fn map_open_run_error(run_id: &RunId, err: StoreError) -> Error {
match err {
StoreError::RunNotFound(id) => Error::RunNotFound(id),
other => Error::engine(format!("failed to open run {run_id}: {other}")),
}
}
/// The canonical "run is archived — mutation rejected" error message. Shared
/// by the operations layer, the CLI rewind precheck, and the server HTTP
/// guards so the user sees the same actionable guidance everywhere.

View file

@ -6,7 +6,7 @@ use fabro_types::RunId;
use git2::{Oid, Signature};
use super::run_git;
use super::timeline::{ForkTarget, TimelineEntry, build_timeline};
use super::timeline::{ForkTarget, RunTimeline, TimelineEntry, build_timeline};
use crate::error::Error;
use crate::event::{self, Event};
use crate::git::{MetadataStore, RUN_BRANCH_PREFIX, push_run_branches};
@ -51,14 +51,10 @@ struct ForkedRun {
/// checkpoint.
///
/// Returns the new run ID.
pub fn fork(store: &Store, input: &ForkRunInput) -> Result<RunId> {
#[cfg(test)]
fn fork(store: &Store, input: &ForkRunInput) -> Result<RunId> {
let timeline = build_timeline(store, &input.source_run_id.to_string())?;
let entry = match input.target.as_ref() {
Some(target) => timeline.resolve(target)?,
None => timeline.entries.last().ok_or_else(|| {
anyhow::anyhow!("no checkpoints found for run {}", input.source_run_id)
})?,
};
let entry = resolve_fork_entry(&timeline, &input.source_run_id, input.target.as_ref())?;
Ok(fork_from_entry(store, &input.source_run_id, entry, input.push)?.new_run_id)
}
@ -71,14 +67,8 @@ pub async fn fork_run(store: &Database, input: &ForkRunInput) -> Result<ForkOutc
run_git::with_run_git_store(store, source_run_id, move |git_store| {
let timeline = build_timeline(&git_store, &source_run_id.to_string())
.map_err(|err| Error::engine(err.to_string()))?;
let entry = match target.as_ref() {
Some(target) => timeline
.resolve(target)
.map_err(|err| Error::Validation(err.to_string()))?,
None => timeline.entries.last().ok_or_else(|| {
Error::Validation(format!("no checkpoints found for run {source_run_id}"))
})?,
};
let entry = resolve_fork_entry(&timeline, &source_run_id, target.as_ref())
.map_err(|err| Error::Validation(err.to_string()))?;
let resolved = ResolvedForkTarget {
checkpoint_ordinal: entry.ordinal,
node_id: entry.node_name.clone(),
@ -100,6 +90,20 @@ pub async fn fork_run(store: &Database, input: &ForkRunInput) -> Result<ForkOutc
Ok(outcome)
}
fn resolve_fork_entry<'a>(
timeline: &'a RunTimeline,
source_run_id: &RunId,
target: Option<&ForkTarget>,
) -> Result<&'a TimelineEntry> {
match target {
Some(target) => timeline.resolve(target),
None => timeline
.entries
.last()
.ok_or_else(|| anyhow::anyhow!("no checkpoints found for run {source_run_id}")),
}
}
fn fork_from_entry(
store: &Store,
source_run_id: &RunId,

View file

@ -5,6 +5,7 @@ mod rebuild_meta;
mod resume;
mod rewind;
mod run_git;
mod run_store;
mod source;
mod start;
#[cfg(test)]
@ -17,7 +18,7 @@ pub use archive::{
unarchive,
};
pub use create::{CreateRunInput, CreatedRun, create, make_run_dir};
pub use fork::{ForkOutcome, ForkRunInput, ResolvedForkTarget, fork, fork_run};
pub use fork::{ForkOutcome, ForkRunInput, ResolvedForkTarget, fork_run};
pub use rebuild_meta::{
build_timeline_or_rebuild, find_run_id_by_prefix_or_store, rebuild_metadata_branch,
};

View file

@ -1,5 +1,5 @@
use fabro_store::Database;
use fabro_types::{ActorRef, RunId, RunStatus};
use fabro_types::{ActorRef, RunId};
use tracing::error;
use super::fork::{self, ForkOutcome, ForkRunInput, ResolvedForkTarget};
@ -21,7 +21,6 @@ pub enum RewindOutcome {
source_run_id: RunId,
new_run_id: RunId,
target: ResolvedForkTarget,
archived: bool,
},
Partial {
source_run_id: RunId,
@ -41,15 +40,8 @@ pub async fn rewind(
Error::Precondition(format!("run {} has no status; cannot rewind", input.run_id))
})?;
if matches!(current, RunStatus::Archived { .. }) {
return Err(Error::Precondition(archive::archived_rejection_message(
&input.run_id,
)));
}
if !matches!(
current,
RunStatus::Succeeded { .. } | RunStatus::Failed { .. } | RunStatus::Dead
) {
archive::ensure_not_archived(Some(current), &input.run_id)?;
if current.terminal_status().is_none() {
return Err(Error::Precondition(format!(
"run {} must be terminal (succeeded, failed, or dead) to rewind; current status is {current}",
input.run_id
@ -65,12 +57,11 @@ pub async fn rewind(
match archive::archive(store, &input.run_id, actor).await {
Ok(_) => {
append_superseded_event(store, &forked).await;
append_superseded_event_best_effort(store, &forked).await;
Ok(RewindOutcome::Full {
source_run_id: forked.source_run_id,
new_run_id: forked.new_run_id,
target: forked.target,
archived: true,
})
}
Err(err) => Ok(RewindOutcome::Partial {
@ -82,7 +73,7 @@ pub async fn rewind(
}
}
async fn append_superseded_event(store: &Database, forked: &ForkOutcome) {
async fn append_superseded_event_best_effort(store: &Database, forked: &ForkOutcome) {
let run_store = match store.open_run(&forked.source_run_id).await {
Ok(run_store) => run_store,
Err(err) => {

View file

@ -1,18 +1,12 @@
use fabro_checkpoint::git::Store as GitStore;
use fabro_store::{Database, Error as StoreError};
use fabro_store::Database;
use fabro_types::{RunId, RunProjection};
use git2::Repository;
use tokio::task::spawn_blocking;
use super::run_store::map_open_run_error;
use crate::error::Error;
fn map_open_run_error(run_id: &RunId, err: StoreError) -> Error {
match err {
StoreError::RunNotFound(id) => Error::RunNotFound(id),
other => Error::engine(format!("failed to open run {run_id}: {other}")),
}
}
pub(crate) async fn load_projection(
store: &Database,
run_id: &RunId,

View file

@ -0,0 +1,11 @@
use fabro_store::Error as StoreError;
use fabro_types::RunId;
use crate::error::Error;
pub(super) fn map_open_run_error(run_id: &RunId, err: StoreError) -> Error {
match err {
StoreError::RunNotFound(id) => Error::RunNotFound(id),
other => Error::engine(format!("failed to open run {run_id}: {other}")),
}
}