Name the fold's repeated shapes once: completion, model usage, pause, conclusion

Four small duplications in the projection fold: a StageCompletion
literal built three times from an attempt's status, a StageModelUsage
literal built three times with no request controls, the pause and
unpause status arithmetic written four times across the coordinator and
lifecycle folds, and a bare Conclusion built beside the full one. Each is
now one function: completion() in the engine fold, StageModelUsage::new,
RunStatus::paused and RunStatus::unpaused beside blocked_reason (with
settle_control for the pending control they clear), and
Conclusion::outcome_only.

One case reads differently: a coordinator RunPaused that lands on a run
already paused behind a block now keeps that block for the unpause, as
the lifecycle Paused already did, instead of dropping it.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-20 15:01:14 -04:00
parent 84e3507600
commit 2f8ee2ff36
No known key found for this signature in database
8 changed files with 111 additions and 89 deletions

View file

@ -10,7 +10,7 @@ use fabro_types::{
use petri_execution::CoordinatorEvent;
use petri_execution::events::RunEvent;
use super::{InvocationRef, RunView, apply_status, stage_key};
use super::{InvocationRef, RunView, apply_status, settle_control, stage_key};
impl RunView {
pub(super) fn fold_coordinator(
@ -86,28 +86,14 @@ impl RunView {
}
CoordinatorEvent::RunPaused => {
if let Some(projection) = self.projection.as_mut() {
let prior_block = match projection.status {
RunStatus::Blocked { blocked_reason } => Some(blocked_reason),
_ => None,
};
apply_status(projection, RunStatus::Paused { prior_block }, at);
if projection.pending_control == Some(RunControlAction::Pause) {
projection.pending_control = None;
}
apply_status(projection, projection.status.paused(), at);
settle_control(projection, RunControlAction::Pause);
}
}
CoordinatorEvent::RunUnpaused => {
if let Some(projection) = self.projection.as_mut() {
let next = match projection.status {
RunStatus::Paused {
prior_block: Some(blocked_reason),
} => RunStatus::Blocked { blocked_reason },
_ => RunStatus::Running,
};
apply_status(projection, next, at);
if projection.pending_control == Some(RunControlAction::Unpause) {
projection.pending_control = None;
}
apply_status(projection, projection.status.unpaused(), at);
settle_control(projection, RunControlAction::Unpause);
}
}
CoordinatorEvent::RunFinished { status } => {

View file

@ -32,10 +32,8 @@ impl RunView {
if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) {
stage.state = StageState::Skipped;
stage.completion = Some(StageCompletion {
outcome: StageOutcome::Skipped,
notes: None,
failure_reason: failure_message(&outcome.status),
timestamp: at,
outcome: StageOutcome::Skipped,
..completion(&outcome.status, at)
});
}
}
@ -120,12 +118,7 @@ impl RunView {
stage.live_streaming = Some(false);
apply_metrics(stage, &outcome.metrics);
if is_final {
stage.completion = Some(StageCompletion {
outcome: stage_outcome(&outcome.status),
notes: None,
failure_reason: failure_message(&outcome.status),
timestamp: at,
});
stage.completion = Some(completion(&outcome.status, at));
stage.termination = Some(match outcome.status {
Status::TimedOut => fabro_types::CommandTermination::TimedOut,
Status::Cancelled => fabro_types::CommandTermination::Cancelled,
@ -263,12 +256,7 @@ impl RunView {
Status::Cancelled => StageState::Cancelled,
};
if stage.completion.is_none() || !*executed {
stage.completion = Some(StageCompletion {
outcome: stage_outcome(&outcome.status),
notes: None,
failure_reason: failure_message(&outcome.status),
timestamp: at,
});
stage.completion = Some(completion(&outcome.status, at));
}
if stage.timing.is_none() {
let wall = stage
@ -405,6 +393,17 @@ pub(super) fn failure_message(status: &Status) -> Option<String> {
}
}
/// A stage's completion from an attempt's status: its outcome and, for a
/// failure, the message.
fn completion(status: &Status, at: DateTime<Utc>) -> StageCompletion {
StageCompletion {
outcome: stage_outcome(status),
notes: None,
failure_reason: failure_message(status),
timestamp: at,
}
}
/// The finished attempt's metrics onto its stage: the timing and the usage
/// the backend reported.
fn apply_metrics(stage: &mut StageProjection, metrics: &Metrics) {

View file

@ -36,7 +36,9 @@ use std::collections::{BTreeMap, BTreeSet};
use chrono::{DateTime, TimeZone as _, Utc};
use fabro_store::platform_records::StoredPlatformRecord;
use fabro_types::{RunDiff, RunId, RunProjection, RunStatus, StageId, StageProjection};
use fabro_types::{
RunControlAction, RunDiff, RunId, RunProjection, RunStatus, StageId, StageProjection,
};
use petri_execution::ExecutionId;
use petri_execution::events::{NodeRef, RunEvent, Subject};
use serde::{Deserialize, Serialize};
@ -230,6 +232,14 @@ fn touch(projection: &mut RunProjection, at: DateTime<Utc>) {
}
}
/// A control the run acknowledged: the pending control is cleared when it
/// is the one that landed.
fn settle_control(projection: &mut RunProjection, action: RunControlAction) {
if projection.pending_control == Some(action) {
projection.pending_control = None;
}
}
/// The key of a stage: its execution and firing.
#[must_use]
pub fn stage_key(execution: u64, firing: u64) -> String {

View file

@ -12,12 +12,12 @@ use fabro_types::{
CheckpointRecord as ViewCheckpoint, Conclusion, FailureCategory, FailureDetail, FailureReason,
PullRequestCreation, PullRequestCreationStatus, PullRequestLink, RunApproval, RunApprovalState,
RunArtifact, RunControlAction, RunDiff, RunFailure, RunProjection, RunSandbox, RunStatus,
RunTiming, StageOutcome, format_blob_ref,
StageOutcome, format_blob_ref,
};
use tracing::debug;
use super::sandbox::sandbox_plan;
use super::{RunView, apply_status, millis, stage_key, touch};
use super::{RunView, apply_status, millis, settle_control, stage_key, touch};
impl RunView {
pub(super) fn fold_platform(&mut self, stored: &StoredPlatformRecord, stream_seq: u64) {
@ -201,16 +201,8 @@ fn fold_lifecycle(projection: &mut RunProjection, record: &RunLifecycleRecord, a
Kind::Submitted => apply_status(projection, RunStatus::Submitted, at),
Kind::StartRequested => {}
Kind::Unpaused => {
let status = match projection.status {
RunStatus::Paused {
prior_block: Some(blocked_reason),
} => RunStatus::Blocked { blocked_reason },
_ => RunStatus::Running,
};
apply_status(projection, status, at);
if projection.pending_control == Some(RunControlAction::Unpause) {
projection.pending_control = None;
}
apply_status(projection, projection.status.unpaused(), at);
settle_control(projection, RunControlAction::Unpause);
}
Kind::Pending => {
if let Some(status) = record.status {
@ -290,15 +282,8 @@ fn fold_lifecycle(projection: &mut RunProjection, record: &RunLifecycleRecord, a
}
}
Kind::Paused => {
let prior_block = match projection.status {
RunStatus::Blocked { blocked_reason } => Some(blocked_reason),
RunStatus::Paused { prior_block } => prior_block,
_ => None,
};
apply_status(projection, RunStatus::Paused { prior_block }, at);
if projection.pending_control == Some(RunControlAction::Pause) {
projection.pending_control = None;
}
apply_status(projection, projection.status.paused(), at);
settle_control(projection, RunControlAction::Pause);
}
Kind::Succeeded | Kind::Failed => {
if let Some(status) = record.status {
@ -324,17 +309,7 @@ fn fold_lifecycle(projection: &mut RunProjection, record: &RunLifecycleRecord, a
),
_ => (StageOutcome::Succeeded, None),
};
projection.conclusion = Some(Conclusion {
timestamp: at,
status: outcome,
timing: RunTiming::default(),
failure,
final_git_commit_sha: None,
stages: Vec::new(),
usage: None,
total_retries: 0,
diff: RunDiff::default(),
});
projection.conclusion = Some(Conclusion::outcome_only(at, outcome, failure));
}
}
Kind::CancelRequested => projection.pending_control = Some(RunControlAction::Cancel),

View file

@ -57,13 +57,11 @@ impl RunView {
let model = payload.get("model").and_then(Value::as_str);
if let Some(model) = model {
let (provider, model_id) = split_model(model);
stage.provider_used = Some(StageModelUsage {
mode: StageModelUsage::MODE_PROMPT.to_string(),
provider: provider.map(str::to_string),
model: Some(model_id.to_string()),
reasoning_effort: None,
speed: None,
});
stage.provider_used = Some(StageModelUsage::new(
StageModelUsage::MODE_PROMPT,
provider.map(str::to_string),
Some(model_id.to_string()),
));
stage.model = model_ref(provider, model_id);
}
}
@ -88,13 +86,11 @@ impl RunView {
if let Some(route) = route {
let provider = route.get("provider").and_then(Value::as_str);
let model = route.get("model").and_then(Value::as_str);
stage.provider_used = Some(StageModelUsage {
mode: StageModelUsage::MODE_AGENT.to_string(),
provider: provider.map(str::to_string),
model: model.map(str::to_string),
reasoning_effort: None,
speed: None,
});
stage.provider_used = Some(StageModelUsage::new(
StageModelUsage::MODE_AGENT,
provider.map(str::to_string),
model.map(str::to_string),
));
if let Some(model) = model {
stage.model = model_ref(provider, model);
}
@ -299,13 +295,11 @@ impl RunView {
CodingEvent::SessionStarted {
provider, model, ..
} if is_root => {
stage.provider_used = Some(StageModelUsage {
mode: StageModelUsage::MODE_AGENT.to_string(),
provider: provider.clone(),
model: model.clone(),
reasoning_effort: None,
speed: None,
});
stage.provider_used = Some(StageModelUsage::new(
StageModelUsage::MODE_AGENT,
provider.clone(),
model.clone(),
));
if let Some(model) = model.as_deref() {
stage.model = model_ref(provider.as_deref(), model);
}

View file

@ -40,3 +40,27 @@ pub struct Conclusion {
#[serde(default)]
pub diff: RunDiff,
}
impl Conclusion {
/// A conclusion that records only how the run ended: no timing, stages,
/// usage or diff. What a terminal lifecycle record gives when the
/// engine recorded no finish of its own.
#[must_use]
pub fn outcome_only(
timestamp: DateTime<Utc>,
status: StageOutcome,
failure: Option<RunFailure>,
) -> Self {
Self {
timestamp,
status,
timing: RunTiming::default(),
failure,
final_git_commit_sha: None,
stages: Vec::new(),
usage: None,
total_retries: 0,
diff: RunDiff::default(),
}
}
}

View file

@ -122,6 +122,19 @@ impl StageModelUsage {
pub const MODE_AGENT: &'static str = "agent";
pub const MODE_ACP: &'static str = "acp";
/// The usage record of a stage that named its provider and model, with
/// no request controls.
#[must_use]
pub fn new(mode: &str, provider: Option<String>, model: Option<String>) -> Self {
Self {
mode: mode.to_string(),
provider,
model,
reasoning_effort: None,
speed: None,
}
}
/// Build the usage record from a `stage.prompt` event, returning `None`
/// when the event carried no model metadata.
#[must_use]

View file

@ -120,6 +120,27 @@ impl RunStatus {
}
}
/// The status a pause takes the run to: paused, remembering the block
/// the run was under so the unpause can restore it.
#[must_use]
pub fn paused(self) -> Self {
Self::Paused {
prior_block: self.blocked_reason(),
}
}
/// The status an unpause takes the run to: back to the block the pause
/// remembered, else running.
#[must_use]
pub fn unpaused(self) -> Self {
match self {
Self::Paused {
prior_block: Some(blocked_reason),
} => Self::Blocked { blocked_reason },
_ => Self::Running,
}
}
pub fn terminal_status(self) -> Option<TerminalStatus> {
match self {
Self::Succeeded { reason } => Some(TerminalStatus::Succeeded { reason }),