From 2f8ee2ff3670f4eabf188095c5a8ec2647227060 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sun, 20 Sep 2026 15:01:14 -0400 Subject: [PATCH] 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 --- .../fabro-petri/src/projection/coordinator.rs | 24 +++--------- .../fabro-petri/src/projection/engine.rs | 31 +++++++-------- .../fabro-petri/src/projection/mod.rs | 12 +++++- .../fabro-petri/src/projection/platform.rs | 39 ++++--------------- .../fabro-petri/src/projection/progress.rs | 36 +++++++---------- lib/foundation/fabro-types/src/conclusion.rs | 24 ++++++++++++ .../fabro-types/src/run_projection.rs | 13 +++++++ lib/foundation/fabro-types/src/status.rs | 21 ++++++++++ 8 files changed, 111 insertions(+), 89 deletions(-) diff --git a/lib/components/fabro-petri/src/projection/coordinator.rs b/lib/components/fabro-petri/src/projection/coordinator.rs index ffa5093ce..79d784e4a 100644 --- a/lib/components/fabro-petri/src/projection/coordinator.rs +++ b/lib/components/fabro-petri/src/projection/coordinator.rs @@ -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 } => { diff --git a/lib/components/fabro-petri/src/projection/engine.rs b/lib/components/fabro-petri/src/projection/engine.rs index a936ab494..e542b816b 100644 --- a/lib/components/fabro-petri/src/projection/engine.rs +++ b/lib/components/fabro-petri/src/projection/engine.rs @@ -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 { } } +/// A stage's completion from an attempt's status: its outcome and, for a +/// failure, the message. +fn completion(status: &Status, at: DateTime) -> 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) { diff --git a/lib/components/fabro-petri/src/projection/mod.rs b/lib/components/fabro-petri/src/projection/mod.rs index 4191bdd62..6cc9f2cbf 100644 --- a/lib/components/fabro-petri/src/projection/mod.rs +++ b/lib/components/fabro-petri/src/projection/mod.rs @@ -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) { } } +/// 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 { diff --git a/lib/components/fabro-petri/src/projection/platform.rs b/lib/components/fabro-petri/src/projection/platform.rs index d5c592768..ebff0cd24 100644 --- a/lib/components/fabro-petri/src/projection/platform.rs +++ b/lib/components/fabro-petri/src/projection/platform.rs @@ -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), diff --git a/lib/components/fabro-petri/src/projection/progress.rs b/lib/components/fabro-petri/src/projection/progress.rs index 7bfe03428..95f8996e5 100644 --- a/lib/components/fabro-petri/src/projection/progress.rs +++ b/lib/components/fabro-petri/src/projection/progress.rs @@ -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); } diff --git a/lib/foundation/fabro-types/src/conclusion.rs b/lib/foundation/fabro-types/src/conclusion.rs index a06697070..65ec19c23 100644 --- a/lib/foundation/fabro-types/src/conclusion.rs +++ b/lib/foundation/fabro-types/src/conclusion.rs @@ -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, + status: StageOutcome, + failure: Option, + ) -> Self { + Self { + timestamp, + status, + timing: RunTiming::default(), + failure, + final_git_commit_sha: None, + stages: Vec::new(), + usage: None, + total_retries: 0, + diff: RunDiff::default(), + } + } +} diff --git a/lib/foundation/fabro-types/src/run_projection.rs b/lib/foundation/fabro-types/src/run_projection.rs index 339aed897..89a884ff5 100644 --- a/lib/foundation/fabro-types/src/run_projection.rs +++ b/lib/foundation/fabro-types/src/run_projection.rs @@ -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, model: Option) -> 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] diff --git a/lib/foundation/fabro-types/src/status.rs b/lib/foundation/fabro-types/src/status.rs index 2b2c63e75..1a7db7061 100644 --- a/lib/foundation/fabro-types/src/status.rs +++ b/lib/foundation/fabro-types/src/status.rs @@ -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 { match self { Self::Succeeded { reason } => Some(TerminalStatus::Succeeded { reason }),