diff --git a/apps/fabro-web/app/components/stage-sidebar.tsx b/apps/fabro-web/app/components/stage-sidebar.tsx index c7ae643da..b72ee77bd 100644 --- a/apps/fabro-web/app/components/stage-sidebar.tsx +++ b/apps/fabro-web/app/components/stage-sidebar.tsx @@ -1,6 +1,6 @@ import { useState, useEffect, useRef, type ComponentType } from "react"; import { Link } from "react-router"; -import type { StageStatus } from "@qltysh/fabro-api-client"; +import type { StageState } from "@qltysh/fabro-api-client"; import { ArrowPathIcon, CheckCircleIcon, @@ -16,12 +16,12 @@ import { ACTIVE_STAGE_STATES } from "../lib/stage-sidebar"; export interface Stage { id: string; name: string; - status: StageStatus; + status: StageState; duration: string; dotId?: string; } -export const statusConfig: Record; color: string }> = { +export const statusConfig: Record; color: string }> = { pending: { icon: PauseCircleIcon, color: "text-fg-muted" }, running: { icon: ArrowPathIcon, color: "text-teal-500" }, retrying: { icon: ArrowPathIcon, color: "text-amber" }, diff --git a/apps/fabro-web/app/lib/stage-sidebar.ts b/apps/fabro-web/app/lib/stage-sidebar.ts index c346722c1..747c6a70e 100644 --- a/apps/fabro-web/app/lib/stage-sidebar.ts +++ b/apps/fabro-web/app/lib/stage-sidebar.ts @@ -1,11 +1,11 @@ -import type { PaginatedRunStageList, StageStatus } from "@qltysh/fabro-api-client"; +import type { PaginatedRunStageList, StageState } from "@qltysh/fabro-api-client"; import type { Stage } from "../components/stage-sidebar"; import { isVisibleStage } from "../data/runs"; import { formatDurationSecs } from "./format"; -export const ACTIVE_STAGE_STATES: ReadonlySet = new Set(["running", "retrying"]); -export const SUCCEEDED_STAGE_STATES: ReadonlySet = new Set([ +export const ACTIVE_STAGE_STATES: ReadonlySet = new Set(["running", "retrying"]); +export const SUCCEEDED_STAGE_STATES: ReadonlySet = new Set([ "succeeded", "partially_succeeded", ]); diff --git a/docs/internal/plan-events-as-source-of-truth-follow-ups.md b/docs/internal/plan-events-as-source-of-truth-follow-ups.md index 65dd707ca..650eca25d 100644 --- a/docs/internal/plan-events-as-source-of-truth-follow-ups.md +++ b/docs/internal/plan-events-as-source-of-truth-follow-ups.md @@ -312,7 +312,7 @@ Required new payload: Projection rule: -- `NodeState.parallel_results` projects from `parallel.completed.properties.results` +- `StageProjection.parallel_results` projects from `parallel.completed.properties.results` Why this is the right event: @@ -330,7 +330,7 @@ We already added `checkpoint.completed.diff`, but the memoized plan should not p Decision required: -- either `NodeState.diff` is “latest checkpoint diff for that node visit” +- either `StageProjection.diff` is “latest checkpoint diff for that node visit” - or add a dedicated `node.diff_generated` event This follow-up plan should pick one and update docs/tests accordingly. diff --git a/docs/internal/run-directory-keys.md b/docs/internal/run-directory-keys.md index 7a6e87bb8..74eda074c 100644 --- a/docs/internal/run-directory-keys.md +++ b/docs/internal/run-directory-keys.md @@ -34,7 +34,7 @@ These names are still real, but they are no longer live scratch files by default - Metadata branch files such as `run.json`, `start.json`, `checkpoint.json`, and `retro.json` - `fabro dump` exports such as `run.json`, `start.json`, `status.json`, `checkpoint.json`, `conclusion.json`, `retro.json`, `events.jsonl`, and per-node prompt/response/status/stdout/stderr files -- Retro-agent temp uploads named `events.jsonl`, `run.json`, `graph.fabro`, `checkpoints/{seq:04}.json`, `run.log` when available, and per-stage files under `stages/{node_id}@{visit}/...` inside the retro sandbox +- Retro-agent temp uploads named `events.jsonl`, `run.json`, `graph.fabro`, `checkpoints/{seq:04}.json`, `run.log` when available, and per-stage files under `stages/{rank:03}-{node_id}@{visit}/...` inside the retro sandbox ## Notes diff --git a/docs/public/agents/outputs.mdx b/docs/public/agents/outputs.mdx index 5b6c6d39f..8af17dbf4 100644 --- a/docs/public/agents/outputs.mdx +++ b/docs/public/agents/outputs.mdx @@ -7,7 +7,7 @@ When an agent or prompt node finishes, Fabro captures its response text and prod ## Response capture -After an agent or prompt node completes, Fabro captures the full response text and persists it to `stages/{node_id}@{visit}/response.md` in metadata snapshots and `fabro dump` output. It also writes the final outcome (status, context updates, routing directives) to `stages/{node_id}@{visit}/status.json`. +After an agent or prompt node completes, Fabro captures the full response text and persists it to `stages/{rank:03}-{node_id}@{visit}/response.md` in metadata snapshots and `fabro dump` output. It also writes the final outcome (status, context updates, routing directives) to `stages/{rank:03}-{node_id}@{visit}/status.json`. ## Context updates @@ -92,7 +92,7 @@ review -> approve [label="Approve"] ## Output logging -Fabro writes several files per stage to `stages/{node_id}@{visit}/` in metadata snapshots and `fabro dump` output: +Fabro writes several files per stage to `stages/{rank:03}-{node_id}@{visit}/` in metadata snapshots and `fabro dump` output: | File | Contents | |---|---| diff --git a/docs/public/agents/prompts.mdx b/docs/public/agents/prompts.mdx index 61cc58dc1..f457a8bfc 100644 --- a/docs/public/agents/prompts.mdx +++ b/docs/public/agents/prompts.mdx @@ -295,4 +295,4 @@ Use prompt nodes for analysis, classification, and summarization tasks where too ## Prompt logging -Fabro persists the assembled prompt to `stages/{node_id}@{visit}/prompt.md` in metadata snapshots and `fabro dump` output for every agent and prompt stage. This includes the preamble (if any) and the expanded prompt text. Use these files for debugging when an agent behaves unexpectedly. +Fabro persists the assembled prompt to `stages/{rank:03}-{node_id}@{visit}/prompt.md` in metadata snapshots and `fabro dump` output for every agent and prompt stage. This includes the preamble (if any) and the expanded prompt text. Use these files for debugging when an agent behaves unexpectedly. diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index bae632f8e..fff7cf9c8 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -5103,14 +5103,14 @@ components: - failed - skipped - NodeStatusRecord: - description: Internal node status record. + StageCompletion: + description: Terminal completion metadata for a projected workflow stage. type: object required: - - status + - outcome - timestamp properties: - status: + outcome: $ref: "#/components/schemas/StageOutcome" notes: type: ["string", "null"] @@ -5120,22 +5120,23 @@ components: type: string format: date-time - StageState: - description: Internal stage projection state. + StageProjection: + description: Observable projection data for one workflow stage execution. type: object + required: + - first_event_seq properties: - seq: + first_event_seq: type: integer - format: int32 - minimum: 0 - description: Event-log sequence number of the first event that created this stage. Used to order stages by execution. + format: uint32 + minimum: 1 prompt: type: ["string", "null"] response: type: ["string", "null"] - status: + completion: oneOf: - - $ref: "#/components/schemas/NodeStatusRecord" + - $ref: "#/components/schemas/StageCompletion" - type: "null" provider_used: type: ["object", "null"] @@ -5356,9 +5357,9 @@ components: $ref: "#/components/schemas/PendingInterviewRecord" stages: type: object - description: Map from StageId (`node_id@visit`) to StageState. + description: Map from StageId (`node_id@visit`) to stage projection data. additionalProperties: - $ref: "#/components/schemas/StageState" + $ref: "#/components/schemas/StageProjection" RunSummary: description: Durable run summary derived from the backing store. @@ -6141,7 +6142,7 @@ components: # ── Stage / Turn Schemas ───────────────────────────────────────────── - StageStatus: + StageState: description: Lifecycle projection state of a workflow stage. type: string enum: @@ -6171,7 +6172,7 @@ components: description: Human-readable stage name. example: Propose Changes status: - $ref: "#/components/schemas/StageStatus" + $ref: "#/components/schemas/StageState" duration_secs: type: number description: Time spent in this stage, in seconds. diff --git a/docs/public/execution/checkpoints.mdx b/docs/public/execution/checkpoints.mdx index fb75f6620..2bf801f3d 100644 --- a/docs/public/execution/checkpoints.mdx +++ b/docs/public/execution/checkpoints.mdx @@ -51,7 +51,7 @@ The metadata branch (`fabro/meta/{run_id}`) is an orphan branch that stores stru After each node, the metadata branch is updated with: - **`run.json`** — Refreshed projection snapshot with the new current checkpoint -- **`stages/{node_id}@{visit}/...`** — Per-stage execution trace files (prompts, responses, status, diffs, stdout/stderr, and tool metadata) +- **`stages/{rank:03}-{node_id}@{visit}/...`** — Execution-order-prefixed per-stage trace files (prompts, responses, status, diffs, stdout/stderr, and tool metadata) - **`stages/retro/*.md`** — Retro prompt/response text when present ## What's in a checkpoint diff --git a/docs/public/execution/retros.mdx b/docs/public/execution/retros.mdx index 41c48e9eb..9c8c65019 100644 --- a/docs/public/execution/retros.mdx +++ b/docs/public/execution/retros.mdx @@ -101,7 +101,7 @@ Retro generation happens in two phases after a run completes: 1. **Derive** — Fabro extracts stage durations from durable run events and builds a retro from the checkpoint data. This is deterministic, fast, and produces the quantitative layer. -2. **Narrate** — An LLM agent session analyzes the run data. The agent receives `events.jsonl`, `run.json`, `graph.fabro`, checkpoint snapshots under `checkpoints/{seq:04}.json`, `run.log` when available, and per-stage files under `stages/{node_id}@{visit}/...` inside its sandbox so it can grep and read the event stream, run snapshot, workflow source, checkpoints, logs, and full stage payloads. The narrative fields are merged back into durable retro state. +2. **Narrate** — An LLM agent session analyzes the run data. The agent receives `events.jsonl`, `run.json`, `graph.fabro`, checkpoint snapshots under `checkpoints/{seq:04}.json`, `run.log` when available, and per-stage files under `stages/{rank:03}-{node_id}@{visit}/...` inside its sandbox so it can grep and read the event stream, run snapshot, workflow source, checkpoints, logs, and full stage payloads. The narrative fields are merged back into durable retro state. Both phases run automatically at the end of every CLI run. The API server derives the quantitative layer but does not currently run the narrative agent. diff --git a/docs/public/reference/run-directory.mdx b/docs/public/reference/run-directory.mdx index 1eafe5b55..853652dfc 100644 --- a/docs/public/reference/run-directory.mdx +++ b/docs/public/reference/run-directory.mdx @@ -39,7 +39,7 @@ Metadata branch snapshots and `fabro dump` exports now use the same core layout: - `run.json` for the current projection snapshot, including the current checkpoint - `graph.fabro` for workflow source - `stages/retro/*.md` for retro prompt/response text -- `stages/{node_id}@{visit}/...` for per-stage prompt, response, status, diff, stdout, and stderr files +- `stages/{rank:03}-{node_id}@{visit}/...` for execution-order-prefixed per-stage prompt, response, status, diff, stdout, and stderr files `fabro dump` adds export-only history surfaces on top of that shared layout: diff --git a/lib/crates/fabro-api/build.rs b/lib/crates/fabro-api/build.rs index 004b59741..f95e77e20 100644 --- a/lib/crates/fabro-api/build.rs +++ b/lib/crates/fabro-api/build.rs @@ -312,16 +312,16 @@ fn main() { ("ActorKind", "fabro_types::ActorKind", &[]), ("ActorRef", "fabro_types::ActorRef", &[]), ("QuestionType", "fabro_types::QuestionType", &[]), - ("NodeStatusRecord", "fabro_types::NodeStatusRecord", &[]), + ("StageCompletion", "fabro_types::StageCompletion", &[]), ("StageOutcome", "fabro_types::StageOutcome", &[]), - ("StageStatus", "fabro_types::StageStatus", &[]), + ("StageState", "fabro_types::StageState", &[]), ( "CommandOutputStream", "fabro_types::CommandOutputStream", &[], ), ("CommandTermination", "fabro_types::CommandTermination", &[]), - ("StageState", "fabro_types::StageState", &[]), + ("StageProjection", "fabro_types::StageProjection", &[]), ("SecretMetadata", "fabro_types::SecretMetadata", &[]), ("InterviewOption", "fabro_types::InterviewOption", &[]), ( diff --git a/lib/crates/fabro-api/src/lib.rs b/lib/crates/fabro-api/src/lib.rs index ec3318881..626e8cfb5 100644 --- a/lib/crates/fabro-api/src/lib.rs +++ b/lib/crates/fabro-api/src/lib.rs @@ -31,9 +31,9 @@ pub mod types { pub use fabro_types::{ ActorKind, ActorRef, BilledTokenCounts, CommandOutputStream, CommandTermination, DiffStats, DirtyStatus, EventEnvelope, GitContext, InterviewOption, InterviewQuestionRecord, - NodeStatusRecord, PendingInterviewRecord, PreRunPushOutcome, QuestionType, - RepositoryReference, RunEvent, RunProjection, RunSummary, SecretMetadata, SecretType, - ServerSettings, StageOutcome, StageState, StageStatus, WorkflowSettings, + PendingInterviewRecord, PreRunPushOutcome, QuestionType, RepositoryReference, RunEvent, + RunProjection, RunSummary, SecretMetadata, SecretType, ServerSettings, StageCompletion, + StageOutcome, StageProjection, StageState, WorkflowSettings, }; pub use crate::generated::types::*; diff --git a/lib/crates/fabro-api/tests/run_projection_round_trip.rs b/lib/crates/fabro-api/tests/run_projection_round_trip.rs index 27a4c85d6..2ec82be1c 100644 --- a/lib/crates/fabro-api/tests/run_projection_round_trip.rs +++ b/lib/crates/fabro-api/tests/run_projection_round_trip.rs @@ -58,10 +58,10 @@ fn run_projection_round_trips_populated_projection() { }, "stages": { "build@2": { - "seq": 0, + "first_event_seq": 3, "prompt": null, "response": null, - "status": null, + "completion": null, "provider_used": null, "diff": "diff --git a/file b/file", "script_invocation": null, diff --git a/lib/crates/fabro-api/tests/node_status_record_round_trip.rs b/lib/crates/fabro-api/tests/stage_completion_round_trip.rs similarity index 56% rename from lib/crates/fabro-api/tests/node_status_record_round_trip.rs rename to lib/crates/fabro-api/tests/stage_completion_round_trip.rs index 04c4741a8..dc3171c57 100644 --- a/lib/crates/fabro-api/tests/node_status_record_round_trip.rs +++ b/lib/crates/fabro-api/tests/stage_completion_round_trip.rs @@ -1,24 +1,24 @@ use std::any::{TypeId, type_name}; -use fabro_api::types::NodeStatusRecord as ApiNodeStatusRecord; -use fabro_types::NodeStatusRecord; +use fabro_api::types::StageCompletion as ApiStageCompletion; +use fabro_types::StageCompletion; use serde_json::json; #[test] -fn node_status_record_reuses_canonical_type() { - assert_same_type::(); +fn stage_completion_reuses_canonical_type() { + assert_same_type::(); } #[test] -fn node_status_record_round_trips_representative_json() { +fn stage_completion_round_trips_representative_json() { let value = json!({ - "status": "partially_succeeded", + "outcome": "partially_succeeded", "notes": "continued with warnings", "failure_reason": null, "timestamp": "2026-04-29T12:34:56Z" }); - let record: NodeStatusRecord = serde_json::from_value(value.clone()).unwrap(); + let record: StageCompletion = serde_json::from_value(value.clone()).unwrap(); assert_eq!(serde_json::to_value(record).unwrap(), value); } diff --git a/lib/crates/fabro-api/tests/stage_projection_round_trip.rs b/lib/crates/fabro-api/tests/stage_projection_round_trip.rs new file mode 100644 index 000000000..0523197cf --- /dev/null +++ b/lib/crates/fabro-api/tests/stage_projection_round_trip.rs @@ -0,0 +1,46 @@ +use std::any::{TypeId, type_name}; + +use fabro_api::types::StageProjection as ApiStageProjection; +use fabro_types::StageProjection; +use serde_json::json; + +#[test] +fn stage_projection_reuses_canonical_type() { + assert_same_type::(); +} + +#[test] +fn stage_projection_round_trips_representative_json() { + let value = json!({ + "first_event_seq": 1, + "prompt": "build it", + "response": "done", + "completion": { + "outcome": "succeeded", + "notes": null, + "failure_reason": null, + "timestamp": "2026-04-29T12:34:56Z" + }, + "provider_used": { "provider": "openai", "model": "gpt-5.2" }, + "diff": "diff --git a/file b/file", + "script_invocation": { "command": "cargo test" }, + "script_timing": { "duration_ms": 42 }, + "parallel_results": [{ "branch": 0, "status": "succeeded" }], + "stdout": "ok", + "stderr": "", + "termination": "exited" + }); + + let state: StageProjection = serde_json::from_value(value.clone()).unwrap(); + assert_eq!(serde_json::to_value(state).unwrap(), value); +} + +fn assert_same_type() { + assert_eq!( + TypeId::of::(), + TypeId::of::(), + "{} should be the same type as {}", + type_name::(), + type_name::() + ); +} diff --git a/lib/crates/fabro-api/tests/stage_state_round_trip.rs b/lib/crates/fabro-api/tests/stage_state_round_trip.rs index 29adc6171..f1c4f7740 100644 --- a/lib/crates/fabro-api/tests/stage_state_round_trip.rs +++ b/lib/crates/fabro-api/tests/stage_state_round_trip.rs @@ -10,30 +10,51 @@ fn stage_state_reuses_canonical_type() { } #[test] -fn stage_state_round_trips_representative_json() { - let value = json!({ - "seq": 42, - "prompt": "build it", - "response": "done", - "status": { - "status": "succeeded", - "notes": null, - "failure_reason": null, - "timestamp": "2026-04-29T12:34:56Z" - }, - "provider_used": { "provider": "openai", "model": "gpt-5.2" }, - "diff": "diff --git a/file b/file", - "script_invocation": { "command": "cargo test" }, - "script_timing": { "duration_ms": 42 }, - "parallel_results": [{ "branch": 0, "status": "succeeded" }], - "stdout": "ok", - "stderr": "", - "termination": "exited" - }); +fn stage_state_serializes_as_lifecycle_strings() { + assert_eq!( + serde_json::to_value(StageState::Pending).unwrap(), + json!("pending") + ); + assert_eq!( + serde_json::to_value(StageState::Running).unwrap(), + json!("running") + ); + assert_eq!( + serde_json::to_value(StageState::Retrying).unwrap(), + json!("retrying") + ); + assert_eq!( + serde_json::to_value(StageState::Succeeded).unwrap(), + json!("succeeded") + ); + assert_eq!( + serde_json::to_value(StageState::PartiallySucceeded).unwrap(), + json!("partially_succeeded") + ); + assert_eq!( + serde_json::to_value(StageState::Failed).unwrap(), + json!("failed") + ); + assert_eq!( + serde_json::to_value(StageState::Skipped).unwrap(), + json!("skipped") + ); + assert_eq!( + serde_json::to_value(StageState::Cancelled).unwrap(), + json!("cancelled") + ); +} - let state: StageState = serde_json::from_value(value.clone()).unwrap(); - assert_eq!(state.seq, 42); - assert_eq!(serde_json::to_value(state).unwrap(), value); +#[test] +fn stage_state_deserializes_representative_values() { + assert_eq!( + serde_json::from_value::(json!("retrying")).unwrap(), + StageState::Retrying + ); + assert_eq!( + serde_json::from_value::(json!("partially_succeeded")).unwrap(), + StageState::PartiallySucceeded + ); } fn assert_same_type() { diff --git a/lib/crates/fabro-api/tests/stage_status_round_trip.rs b/lib/crates/fabro-api/tests/stage_status_round_trip.rs deleted file mode 100644 index de78bffcc..000000000 --- a/lib/crates/fabro-api/tests/stage_status_round_trip.rs +++ /dev/null @@ -1,68 +0,0 @@ -use std::any::{TypeId, type_name}; - -use fabro_api::types::StageStatus as ApiStageStatus; -use fabro_types::StageStatus; -use serde_json::json; - -#[test] -fn stage_status_reuses_canonical_type() { - assert_same_type::(); -} - -#[test] -fn stage_status_serializes_as_lifecycle_strings() { - assert_eq!( - serde_json::to_value(StageStatus::Pending).unwrap(), - json!("pending") - ); - assert_eq!( - serde_json::to_value(StageStatus::Running).unwrap(), - json!("running") - ); - assert_eq!( - serde_json::to_value(StageStatus::Retrying).unwrap(), - json!("retrying") - ); - assert_eq!( - serde_json::to_value(StageStatus::Succeeded).unwrap(), - json!("succeeded") - ); - assert_eq!( - serde_json::to_value(StageStatus::PartiallySucceeded).unwrap(), - json!("partially_succeeded") - ); - assert_eq!( - serde_json::to_value(StageStatus::Failed).unwrap(), - json!("failed") - ); - assert_eq!( - serde_json::to_value(StageStatus::Skipped).unwrap(), - json!("skipped") - ); - assert_eq!( - serde_json::to_value(StageStatus::Cancelled).unwrap(), - json!("cancelled") - ); -} - -#[test] -fn stage_status_deserializes_representative_values() { - assert_eq!( - serde_json::from_value::(json!("retrying")).unwrap(), - StageStatus::Retrying - ); - assert_eq!( - serde_json::from_value::(json!("partially_succeeded")).unwrap(), - StageStatus::PartiallySucceeded - ); -} - -fn assert_same_type() { - assert_eq!( - TypeId::of::(), - TypeId::of::(), - "{} should be the same type as {}", - type_name::(), - type_name::() - ); -} diff --git a/lib/crates/fabro-core/src/lib.rs b/lib/crates/fabro-core/src/lib.rs index acf7117f5..853d4eeaf 100644 --- a/lib/crates/fabro-core/src/lib.rs +++ b/lib/crates/fabro-core/src/lib.rs @@ -23,7 +23,7 @@ pub use lifecycle::{ }; pub use outcome::{ FailureCategory, FailureDetail, NodeResult, NodeResultExt, Outcome, OutcomeMeta, StageOutcome, - StageStatus, + StageState, }; pub use retry::{BackoffPolicy, RetryPolicy}; pub use stall::{ActivityMonitor, StallGuard, StallWatchdog}; diff --git a/lib/crates/fabro-core/src/outcome.rs b/lib/crates/fabro-core/src/outcome.rs index bca933718..53d590fbf 100644 --- a/lib/crates/fabro-core/src/outcome.rs +++ b/lib/crates/fabro-core/src/outcome.rs @@ -1,7 +1,7 @@ use std::time::Duration; pub use fabro_types::outcome::{ - FailureCategory, FailureDetail, NodeResult, Outcome, OutcomeMeta, StageOutcome, StageStatus, + FailureCategory, FailureDetail, NodeResult, Outcome, OutcomeMeta, StageOutcome, StageState, }; use crate::error::Error; diff --git a/lib/crates/fabro-dump/src/lib.rs b/lib/crates/fabro-dump/src/lib.rs index d9d00740c..69a670565 100644 --- a/lib/crates/fabro-dump/src/lib.rs +++ b/lib/crates/fabro-dump/src/lib.rs @@ -66,87 +66,81 @@ impl RunDump { entries.push(RunDumpEntry::text("graph.fabro", graph_source.clone())); } - let mut stage_order: Vec<_> = state - .iter_stages() - .map(|(stage_id, stage)| (stage_id.clone(), stage.seq)) - .collect(); - if stage_order.len() > MAX_STAGES_IN_DUMP { + let mut stages: Vec<_> = state.iter_stages().collect(); + if stages.len() > MAX_STAGES_IN_DUMP { bail!( "run dump supports at most {MAX_STAGES_IN_DUMP} stages with the current path prefix width (got {})", - stage_order.len() + stages.len() ); } - stage_order.sort_by(|(left_id, left_seq), (right_id, right_seq)| { - left_seq.cmp(right_seq).then_with(|| left_id.cmp(right_id)) + stages.sort_by(|(left_id, left), (right_id, right)| { + left.first_event_seq + .cmp(&right.first_event_seq) + .then_with(|| left_id.cmp(right_id)) }); + let mut stage_ranks = HashMap::new(); - for (index, (stage_id, _)) in stage_order.iter().enumerate() { + for (index, (stage_id, _)) in stages.iter().enumerate() { let rank = u32::try_from(index + 1).context("stage rank should fit in u32")?; - stage_ranks.insert(stage_id.clone(), rank); + stage_ranks.insert((*stage_id).clone(), rank); } - for (stage_id, _) in stage_order { - let Some(node) = state.stage(&stage_id) else { - continue; - }; - let rank = stage_ranks - .get(&stage_id) - .copied() - .context("stage rank should exist")?; - let base = PathBuf::from("stages").join(stage_dir_name(rank, &stage_id)); + for (index, (stage_id, stage)) in stages.into_iter().enumerate() { + let rank = u32::try_from(index + 1).context("stage rank should fit in u32")?; + let base = PathBuf::from("stages").join(stage_dir_name(rank, stage_id)); - if let Some(prompt) = node.prompt.as_ref() { + if let Some(prompt) = stage.prompt.as_ref() { entries.push(RunDumpEntry::text_path( &base.join("prompt.md"), prompt.clone(), )); } - if let Some(response) = node.response.as_ref() { + if let Some(response) = stage.response.as_ref() { entries.push(RunDumpEntry::text_path( &base.join("response.md"), response.clone(), )); } - if let Some(status) = node.status.as_ref() { - push_json_entry_path(&mut entries, &base.join("status.json"), status)?; + if let Some(completion) = stage.completion.as_ref() { + push_json_entry_path(&mut entries, &base.join("status.json"), completion)?; } - if let Some(provider_used) = node.provider_used.as_ref() { + if let Some(provider_used) = stage.provider_used.as_ref() { entries.push(RunDumpEntry::json_path( &base.join("provider_used.json"), provider_used.clone(), )); } - if let Some(diff) = node.diff.as_ref() { + if let Some(diff) = stage.diff.as_ref() { entries.push(RunDumpEntry::text_path( &base.join("diff.patch"), diff.clone(), )); } - if let Some(script_invocation) = node.script_invocation.as_ref() { + if let Some(script_invocation) = stage.script_invocation.as_ref() { entries.push(RunDumpEntry::json_path( &base.join("script_invocation.json"), script_invocation.clone(), )); } - if let Some(script_timing) = node.script_timing.as_ref() { + if let Some(script_timing) = stage.script_timing.as_ref() { entries.push(RunDumpEntry::json_path( &base.join("script_timing.json"), script_timing.clone(), )); } - if let Some(parallel_results) = node.parallel_results.as_ref() { + if let Some(parallel_results) = stage.parallel_results.as_ref() { entries.push(RunDumpEntry::json_path( &base.join("parallel_results.json"), parallel_results.clone(), )); } - if let Some(stdout) = node.stdout.as_ref() { + if let Some(stdout) = stage.stdout.as_ref() { entries.push(RunDumpEntry::text_path( &base.join("stdout.log"), stdout.clone(), )); } - if let Some(stderr) = node.stderr.as_ref() { + if let Some(stderr) = stage.stderr.as_ref() { entries.push(RunDumpEntry::text_path( &base.join("stderr.log"), stderr.clone(), @@ -211,6 +205,20 @@ impl RunDump { Ok(()) } + fn add_orphan_notice(&mut self, stage_id: &StageId) { + let line = format!("notice: artifact stage {stage_id} was not present in run projection\n"); + if let Some(index) = self.dump_log_index { + if let Some(RunDumpContents::Text(text)) = + self.entries.get_mut(index).map(|entry| &mut entry.contents) + { + text.push_str(&line); + return; + } + } + self.dump_log_index = Some(self.entries.len()); + self.entries.push(RunDumpEntry::text("dump.log", line)); + } + pub fn add_file_bytes(&mut self, path: impl Into, contents: Vec) { self.entries.push(RunDumpEntry::bytes(path, contents)); } @@ -267,20 +275,6 @@ impl RunDump { self.entries.len() } - fn add_orphan_notice(&mut self, stage_id: &StageId) { - let line = format!("notice: artifact stage {stage_id} was not present in run projection\n"); - if let Some(index) = self.dump_log_index { - if let Some(RunDumpContents::Text(text)) = - self.entries.get_mut(index).map(|entry| &mut entry.contents) - { - text.push_str(&line); - return; - } - } - self.dump_log_index = Some(self.entries.len()); - self.entries.push(RunDumpEntry::text("dump.log", line)); - } - pub fn write_to_dir(&self, root: &Path) -> Result { for entry in &self.entries { entry.write_to_dir(root)?; @@ -495,12 +489,12 @@ mod tests { use std::collections::HashMap; use chrono::{TimeZone, Utc}; - use fabro_store::{RunProjection, StageId, StageState}; + use fabro_store::{RunProjection, StageId}; use fabro_types::graph::Graph; use fabro_types::run::RunSpec; use fabro_types::{ - Checkpoint, Conclusion, NodeStatusRecord, RunStatus, SandboxRecord, StageOutcome, - StartRecord, SuccessReason, WorkflowSettings, fixtures, + Checkpoint, Conclusion, RunStatus, SandboxRecord, StageCompletion, StageOutcome, + StartRecord, SuccessReason, WorkflowSettings, first_event_seq, fixtures, }; use futures::executor; @@ -590,32 +584,26 @@ mod tests { }); projection.retro_prompt = Some("retro prompt".to_string()); projection.retro_response = Some("retro response".to_string()); - projection.set_stage(stage_id.clone(), StageState { - seq: 1, - prompt: Some("plan".to_string()), - response: Some("done".to_string()), - status: Some(NodeStatusRecord { - status: StageOutcome::Succeeded, - notes: Some("ok".to_string()), - failure_reason: None, - timestamp: Utc - .with_ymd_and_hms(2026, 4, 20, 12, 1, 0) - .single() - .unwrap(), - }), - provider_used: Some(serde_json::json!({ "provider": "openai" })), - diff: Some("diff --git a/a b/a".to_string()), - script_invocation: Some(serde_json::json!({ "command": "cargo test" })), - script_timing: Some(serde_json::json!({ "duration_ms": 10 })), - parallel_results: Some(serde_json::json!([{ "stage": "fanout@1" }])), - stdout: Some("stdout".to_string()), - stderr: Some("stderr".to_string()), - stdout_bytes: None, - stderr_bytes: None, - streams_separated: None, - live_streaming: None, - termination: None, + let stage = + projection.stage_entry(stage_id.node_id(), stage_id.visit(), first_event_seq(2)); + stage.prompt = Some("plan".to_string()); + stage.response = Some("done".to_string()); + stage.completion = Some(StageCompletion { + outcome: StageOutcome::Succeeded, + notes: Some("ok".to_string()), + failure_reason: None, + timestamp: Utc + .with_ymd_and_hms(2026, 4, 20, 12, 1, 0) + .single() + .unwrap(), }); + stage.provider_used = Some(serde_json::json!({ "provider": "openai" })); + stage.diff = Some("diff --git a/a b/a".to_string()); + stage.script_invocation = Some(serde_json::json!({ "command": "cargo test" })); + stage.script_timing = Some(serde_json::json!({ "duration_ms": 10 })); + stage.parallel_results = Some(serde_json::json!([{ "stage": "fanout@1" }])); + stage.stdout = Some("stdout".to_string()); + stage.stderr = Some("stderr".to_string()); let dump = RunDump::from_projection(&projection).unwrap(); let paths: Vec<&str> = dump @@ -662,7 +650,6 @@ mod tests { assert!(round_tripped.checkpoint.is_some()); assert!(round_tripped.conclusion.is_some()); assert!(round_tripped.sandbox.is_some()); - assert_eq!(node.seq, 1); assert_eq!(node.prompt, None); assert_eq!(node.response, None); assert_eq!(node.diff, None); @@ -675,118 +662,51 @@ mod tests { } #[test] - fn from_projection_orders_stage_dirs_by_seq() { + fn from_projection_prefixes_stage_paths_but_not_artifact_paths() { let mut projection = RunProjection::default(); - projection.set_stage(StageId::new("zebra", 1), StageState { - seq: 2, - prompt: Some("zebra".to_string()), - ..StageState::default() - }); - projection.set_stage(StageId::new("apple", 1), StageState { - seq: 5, - prompt: Some("apple".to_string()), - ..StageState::default() - }); + projection + .stage_entry("zebra", 1, first_event_seq(1)) + .prompt = Some("first".to_string()); + projection + .stage_entry("apple", 1, first_event_seq(2)) + .prompt = Some("second".to_string()); - let dump = RunDump::from_projection(&projection).unwrap(); - let stage_prompt_paths = dump - .entries() - .iter() - .filter_map(|entry| { - entry - .path - .ends_with("prompt.md") - .then_some(entry.path.as_str()) - }) - .collect::>(); - - assert_eq!(stage_prompt_paths, vec![ - "stages/001-zebra@1/prompt.md", - "stages/002-apple@1/prompt.md", - ]); - } - - #[test] - fn from_projection_reads_legacy_nodes_alias_and_ties_seq_zero_by_stage_id() { - let projection: RunProjection = serde_json::from_value(serde_json::json!({ - "nodes": { - "zebra@1": { "prompt": "zebra" }, - "apple@1": { "prompt": "apple" } - } - })) - .unwrap(); - - let dump = RunDump::from_projection(&projection).unwrap(); - let stage_prompt_paths = dump - .entries() - .iter() - .filter_map(|entry| { - entry - .path - .ends_with("prompt.md") - .then_some(entry.path.as_str()) - }) - .collect::>(); - - assert_eq!(stage_prompt_paths, vec![ - "stages/001-apple@1/prompt.md", - "stages/002-zebra@1/prompt.md", - ]); - let serialized = serde_json::to_value(&projection).unwrap(); - assert!(serialized.get("stages").is_some()); - assert!(serialized.get("nodes").is_none()); - } - - #[test] - fn add_artifact_bytes_uses_stage_rank_and_retry() { - let stage_id = StageId::new("build", 1); - let mut projection = RunProjection::default(); - projection.set_stage(stage_id.clone(), StageState { - seq: 7, - ..StageState::default() - }); let mut dump = RunDump::from_projection(&projection).unwrap(); - - dump.add_artifact_bytes(&stage_id, 1, "logs/output.txt", b"first".to_vec()) + dump.add_artifact_bytes(&StageId::new("zebra", 1), 0, "report.txt", b"z".to_vec()) .unwrap(); - dump.add_artifact_bytes(&stage_id, 2, "logs/output.txt", b"second".to_vec()) + dump.add_artifact_bytes(&StageId::new("apple", 1), 0, "report.txt", b"a".to_vec()) .unwrap(); - let paths = dump + let paths: Vec<&str> = dump .entries() .iter() .map(|entry| entry.path.as_str()) - .collect::>(); - assert!(paths.contains(&"artifacts/001-build@1/retry-0001/logs/output.txt")); - assert!(paths.contains(&"artifacts/001-build@1/retry-0002/logs/output.txt")); + .collect(); + + assert!(paths.contains(&"stages/001-zebra@1/prompt.md")); + assert!(paths.contains(&"stages/002-apple@1/prompt.md")); + assert!(paths.contains(&"artifacts/001-zebra@1/retry-0000/report.txt")); + assert!(paths.contains(&"artifacts/002-apple@1/retry-0000/report.txt")); } #[test] fn add_artifact_bytes_places_orphans_under_sentinel() { - let mut dump = RunDump::from_projection(&RunProjection::default()).unwrap(); + let mut projection = RunProjection::default(); + projection + .stage_entry("known", 1, first_event_seq(1)) + .prompt = Some("present".to_string()); - dump.add_artifact_bytes( - &StageId::new("missing", 1), - 1, - "logs/output.txt", - b"orphan".to_vec(), - ) - .unwrap(); + let mut dump = RunDump::from_projection(&projection).unwrap(); + dump.add_artifact_bytes(&StageId::new("missing", 1), 0, "report.txt", b"m".to_vec()) + .unwrap(); - let orphan = dump + let paths: Vec<&str> = dump .entries() .iter() - .find(|entry| entry.path == "artifacts/_orphans/missing@1/retry-0001/logs/output.txt"); - assert!(orphan.is_some()); - let notice = dump - .entries() - .iter() - .find(|entry| entry.path == "dump.log") - .expect("dump log should include orphan notice"); - let RunDumpContents::Text(text) = ¬ice.contents else { - panic!("dump log should be text"); - }; - assert!(text.contains("missing@1")); + .map(|entry| entry.path.as_str()) + .collect(); + assert!(paths.contains(&"artifacts/_orphans/missing@1/retry-0000/report.txt")); + assert!(paths.contains(&"dump.log")); } #[test] diff --git a/lib/crates/fabro-retro/src/retro_agent.rs b/lib/crates/fabro-retro/src/retro_agent.rs index a98275c66..1046f60b8 100644 --- a/lib/crates/fabro-retro/src/retro_agent.rs +++ b/lib/crates/fabro-retro/src/retro_agent.rs @@ -24,7 +24,7 @@ You have access to the run's data files: - `graph.fabro` — the workflow source for the run - `checkpoints/{seq:04}.json` — zero-padded checkpoint snapshots captured during the run - `run.log` — server/worker log output for the run when available -- `stages/{node_id}@{visit}/...` — per-stage prompt, response, status, diff, stdout/stderr, and tool metadata files +- `stages/{rank:03}-{node_id}@{visit}/...` — execution-order-prefixed per-stage prompt, response, status, diff, stdout/stderr, and tool metadata files ## Your task @@ -334,8 +334,8 @@ mod tests { use chrono::{TimeZone, Utc}; use fabro_agent::LocalSandbox; - use fabro_store::{StageId, StageState}; - use fabro_types::{NodeStatusRecord, StageOutcome}; + use fabro_store::StageId; + use fabro_types::{StageCompletion, StageOutcome, first_event_seq}; use tokio::fs; use super::*; @@ -410,32 +410,25 @@ mod tests { let stage_id = StageId::new("build", 2); let mut state = RunProjection::default(); state.graph_source = Some("digraph Ship {}".to_string()); - state.set_stage(stage_id, StageState { - seq: 1, - prompt: Some("plan".to_string()), - response: Some("done".to_string()), - status: Some(NodeStatusRecord { - status: StageOutcome::Succeeded, - notes: Some("ok".to_string()), - failure_reason: None, - timestamp: Utc - .with_ymd_and_hms(2026, 4, 20, 12, 1, 0) - .single() - .unwrap(), - }), - provider_used: Some(serde_json::json!({ "provider": "openai" })), - diff: Some("diff --git a/a b/a".to_string()), - script_invocation: Some(serde_json::json!({ "command": "cargo test" })), - script_timing: Some(serde_json::json!({ "duration_ms": 10 })), - parallel_results: Some(serde_json::json!([{ "stage": "fanout@1" }])), - stdout: Some("stdout".to_string()), - stderr: Some("stderr".to_string()), - stdout_bytes: None, - stderr_bytes: None, - streams_separated: None, - live_streaming: None, - termination: None, + let stage = state.stage_entry(stage_id.node_id(), stage_id.visit(), first_event_seq(2)); + stage.prompt = Some("plan".to_string()); + stage.response = Some("done".to_string()); + stage.completion = Some(StageCompletion { + outcome: StageOutcome::Succeeded, + notes: Some("ok".to_string()), + failure_reason: None, + timestamp: Utc + .with_ymd_and_hms(2026, 4, 20, 12, 1, 0) + .single() + .unwrap(), }); + stage.provider_used = Some(serde_json::json!({ "provider": "openai" })); + stage.diff = Some("diff --git a/a b/a".to_string()); + stage.script_invocation = Some(serde_json::json!({ "command": "cargo test" })); + stage.script_timing = Some(serde_json::json!({ "duration_ms": 10 })); + stage.parallel_results = Some(serde_json::json!([{ "stage": "fanout@1" }])); + stage.stdout = Some("stdout".to_string()); + stage.stderr = Some("stderr".to_string()); upload_data_files( &sandbox, @@ -522,21 +515,19 @@ mod tests { let mut state = RunProjection::default(); let stdout_ref = fabro_types::format_blob_ref(&stdout_id); let stderr_ref = fabro_types::format_blob_ref(&stderr_id); - state.set_stage(stage_id, StageState { - script_invocation: Some(serde_json::json!({ - "command": "cargo test", - "stdout": stdout_ref, - "stderr": stderr_ref, - })), - script_timing: Some(serde_json::json!({ - "exit_code": 0, - "stdout": stdout_ref, - "stderr": stderr_ref, - })), - stdout: Some(stdout_ref), - stderr: Some(stderr_ref), - ..StageState::default() - }); + let stage = state.stage_entry(stage_id.node_id(), stage_id.visit(), first_event_seq(1)); + stage.script_invocation = Some(serde_json::json!({ + "command": "cargo test", + "stdout": stdout_ref, + "stderr": stderr_ref, + })); + stage.script_timing = Some(serde_json::json!({ + "exit_code": 0, + "stdout": stdout_ref, + "stderr": stderr_ref, + })); + stage.stdout = Some(stdout_ref); + stage.stderr = Some(stderr_ref); let reader: BlobReader = Box::new(move |blob_id| { let stdout_blob = stdout_blob.clone(); diff --git a/lib/crates/fabro-server/src/demo/mod.rs b/lib/crates/fabro-server/src/demo/mod.rs index 38cea6acb..b1a56658f 100644 --- a/lib/crates/fabro-server/src/demo/mod.rs +++ b/lib/crates/fabro-server/src/demo/mod.rs @@ -1158,28 +1158,28 @@ mod runs { RunStage { id: "detect-drift".into(), name: "Detect Drift".into(), - status: StageStatus::Succeeded, + status: StageState::Succeeded, duration_secs: Some(72.0), dot_id: Some("detect".into()), }, RunStage { id: "propose-changes".into(), name: "Propose Changes".into(), - status: StageStatus::Succeeded, + status: StageState::Succeeded, duration_secs: Some(154.0), dot_id: Some("propose".into()), }, RunStage { id: "review-changes".into(), name: "Review Changes".into(), - status: StageStatus::Succeeded, + status: StageState::Succeeded, duration_secs: Some(45.0), dot_id: Some("review".into()), }, RunStage { id: "apply-changes".into(), name: "Apply Changes".into(), - status: StageStatus::Running, + status: StageState::Running, duration_secs: Some(118.0), dot_id: Some("apply".into()), }, diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 9db317480..ed9c84d00 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -34,7 +34,7 @@ pub use fabro_api::types::{ PruneRunsRequest, PruneRunsResponse, RenderWorkflowGraphDirection, RenderWorkflowGraphRequest, RewindRequest, RewindResponse, RunArtifactEntry, RunArtifactListResponse, RunBilling, RunBillingStage, RunBillingTotals, RunError, RunManifest, RunStage, RunStatusResponse, - SandboxFileEntry, SandboxFileListResponse, SshAccessRequest, SshAccessResponse, StageStatus, + SandboxFileEntry, SandboxFileListResponse, SshAccessRequest, SshAccessResponse, StageState, StartRunRequest, SubmitAnswerRequest, SystemFeatures, SystemInfoResponse, SystemRunCounts, TimelineEntryResponse, WriteBlobResponse, }; @@ -2173,7 +2173,7 @@ async fn openapi_spec() -> Response { Json(value).into_response() } -fn active_stage_state_from_events(events: &[EventEnvelope], node_id: &str) -> StageStatus { +fn active_stage_state_from_events(events: &[EventEnvelope], node_id: &str) -> StageState { let latest = events.iter().rev().find(|envelope| { envelope.event.node_id.as_deref() == Some(node_id) && matches!( @@ -2183,9 +2183,9 @@ fn active_stage_state_from_events(events: &[EventEnvelope], node_id: &str) -> St }); if latest.is_some_and(|e| e.event.event_name() == "stage.retrying") { - StageStatus::Retrying + StageState::Retrying } else { - StageStatus::Running + StageState::Running } } @@ -2299,8 +2299,8 @@ async fn list_run_stages( for node_id in &checkpoint.completed_nodes { let duration_ms = stage_durations.get(node_id).copied().unwrap_or(0); let status = match checkpoint.node_outcomes.get(node_id) { - Some(outcome) => StageStatus::from(outcome.status), - None => StageStatus::Succeeded, + Some(outcome) => StageState::from(outcome.status), + None => StageState::Succeeded, }; stages.push(RunStage { id: node_id.clone(), @@ -5393,7 +5393,7 @@ async fn get_run_stage_command_log( .map(str::to_string); let live_streaming = node .live_streaming - .unwrap_or_else(|| cas_ref.is_none() && node.status.is_none()); + .unwrap_or_else(|| cas_ref.is_none() && node.completion.is_none()); let run_dir = Storage::new(state.server_storage_dir()) .run_scratch(&id) .root() @@ -5456,7 +5456,7 @@ async fn get_run_stage_command_log( query.offset, limit, LogSource::Full(&[]), - node.status.is_some(), + node.completion.is_some(), None, live_streaming, ) diff --git a/lib/crates/fabro-store/src/lib.rs b/lib/crates/fabro-store/src/lib.rs index f78f499ea..8c0165b24 100644 --- a/lib/crates/fabro-store/src/lib.rs +++ b/lib/crates/fabro-store/src/lib.rs @@ -17,7 +17,7 @@ pub use artifact_store::{ pub use error::{Error, Result}; pub use fabro_types::{ EventEnvelope, PendingInterviewRecord, RunBlobId, RunProjection, RunSummary, StageId, - StageState, + StageProjection, }; pub(crate) use keyed_mutex::KeyedMutex; pub use run_state::RunProjectionReducer; diff --git a/lib/crates/fabro-store/src/run_state.rs b/lib/crates/fabro-store/src/run_state.rs index 2d778ed94..d1f295bd6 100644 --- a/lib/crates/fabro-store/src/run_state.rs +++ b/lib/crates/fabro-store/src/run_state.rs @@ -8,9 +8,9 @@ use fabro_types::run_event::{ }; use fabro_types::{ BilledModelUsage, Checkpoint, Conclusion, EventBody, FailureSignature, InterviewQuestionRecord, - NodeStatusRecord, Outcome, PendingInterviewRecord, PullRequestRecord, RunControlAction, RunId, - RunProjection, RunSpec, RunStatus, RunSummary, SandboxRecord, StageOutcome, StartRecord, - TerminalStatus, + Outcome, PendingInterviewRecord, PullRequestRecord, RunControlAction, RunEvent, RunId, + RunProjection, RunSpec, RunStatus, RunSummary, SandboxRecord, StageCompletion, StageOutcome, + StageProjection, StartRecord, TerminalStatus, first_event_seq, }; use fabro_util::error::render_with_causes; use serde_json::Value; @@ -200,9 +200,28 @@ impl RunProjectionReducer for RunProjection { .and_then(|visit| u32::try_from(*visit).ok()) .unwrap_or(1); if let Some(diff) = props.diff.clone() { - stage_entry_for_event(self, event, node_id, visit).diff = Some(diff); + self.stage_entry(node_id, visit, first_event_seq(event.seq)) + .diff = Some(diff); } } + for (node_id, outcome) in &checkpoint.node_outcomes { + if outcome.status != StageOutcome::Skipped { + continue; + } + let visit = checkpoint + .node_visits + .get(node_id) + .and_then(|visit| u32::try_from(*visit).ok()) + .unwrap_or(1); + if self + .stage(&fabro_types::StageId::new(node_id, visit)) + .is_some() + { + continue; + } + self.stage_entry(node_id, visit, first_event_seq(event.seq)) + .completion = Some(stage_completion_from_outcome(outcome, ts)); + } self.checkpoint = Some(checkpoint.clone()); self.checkpoints.push((event.seq, checkpoint)); } @@ -267,26 +286,28 @@ impl RunProjectionReducer for RunProjection { EventBody::InterviewInterrupted(props) if !props.question_id.is_empty() => { self.pending_interviews.remove(&props.question_id); } - EventBody::StageStarted(_props) => { + EventBody::StageStarted(_) => { let Some(stage_id) = stored.stage_id.as_ref() else { return Ok(()); }; - self.stage_entry_id(stage_id, event.seq); + self.stage_entry( + stage_id.node_id(), + stage_id.visit(), + first_event_seq(event.seq), + ); } EventBody::StagePrompt(props) => { - let Some(node_id) = stored.node_id.as_deref() else { + let Some(stage) = stage_at_visit(self, stored, props.visit, event.seq) else { return Ok(()); }; - let stage = stage_entry_for_event(self, event, node_id, props.visit); stage.prompt = Some(props.text.clone()); stage.provider_used = provider_used_from_prompt(props); } EventBody::PromptCompleted(props) => { - let Some(node_id) = stored.node_id.as_deref() else { + let Some(stage) = stage_at_current_visit(self, stored, event.seq) else { return Ok(()); }; - stage_entry_with_current_visit(self, event, node_id).response = - Some(props.response.clone()); + stage.response = Some(props.response.clone()); } EventBody::StageCompleted(props) => { let Some(node_id) = stored.node_id.as_deref() else { @@ -295,19 +316,18 @@ impl RunProjectionReducer for RunProjection { let visit = stage_visit(node_id, props.node_visits.as_ref(), self).unwrap_or(1); let response = props.response.clone(); let outcome = stage_outcome_from_props(props); - let status = node_status_from_outcome(&outcome, ts); - let node = stage_entry_for_event(self, event, node_id, visit); - node.response = response; - node.status = Some(status); + let completion = stage_completion_from_outcome(&outcome, ts); + let stage = self.stage_entry(node_id, visit, first_event_seq(event.seq)); + stage.response = response; + stage.completion = Some(completion); } EventBody::StageFailed(props) => { - let Some(node_id) = stored.node_id.as_deref() else { + let failure_reason = props.failure.as_ref().map(|detail| detail.message.clone()); + let Some(stage) = stage_at_current_visit(self, stored, event.seq) else { return Ok(()); }; - let failure_reason = props.failure.as_ref().map(|detail| detail.message.clone()); - let node = stage_entry_with_current_visit(self, event, node_id); - node.status = Some(NodeStatusRecord { - status: StageOutcome::Failed { + stage.completion = Some(StageCompletion { + outcome: StageOutcome::Failed { retry_requested: false, }, notes: None, @@ -316,52 +336,50 @@ impl RunProjectionReducer for RunProjection { }); } EventBody::AgentSessionStarted(props) => { - let Some(node_id) = stored.node_id.as_deref() else { + let Some(stage) = stage_at_visit(self, stored, props.visit, event.seq) else { return Ok(()); }; - stage_entry_for_event(self, event, node_id, props.visit).provider_used = - Some(provider_used_from_agent_session_started(props)); + stage.provider_used = Some(provider_used_from_agent_session_started(props)); } EventBody::AgentCliStarted(props) => { - let Some(node_id) = stored.node_id.as_deref() else { + let Some(stage) = stage_at_visit(self, stored, props.visit, event.seq) else { return Ok(()); }; - stage_entry_for_event(self, event, node_id, props.visit).provider_used = - Some(provider_used_from_agent_cli_started(props)); + stage.provider_used = Some(provider_used_from_agent_cli_started(props)); } EventBody::CommandStarted(props) => { - let Some(node_id) = stored.node_id.as_deref() else { + let script_invocation = serde_json::to_value(props).map_err(|err| { + Error::InvalidEvent(format!("invalid command.started payload: {err}")) + })?; + let Some(stage) = stage_at_current_visit(self, stored, event.seq) else { return Ok(()); }; - stage_entry_with_current_visit(self, event, node_id).script_invocation = - Some(serde_json::to_value(props).map_err(|err| { - Error::InvalidEvent(format!("invalid command.started payload: {err}")) - })?); + stage.script_invocation = Some(script_invocation); } EventBody::CommandCompleted(props) => { - let Some(node_id) = stored.node_id.as_deref() else { + let script_timing = serde_json::to_value(props).map_err(|err| { + Error::InvalidEvent(format!("invalid command.completed payload: {err}")) + })?; + let Some(stage) = stage_at_current_visit(self, stored, event.seq) else { return Ok(()); }; - let node = stage_entry_with_current_visit(self, event, node_id); - node.stdout = Some(props.stdout.clone()); - node.stderr = Some(props.stderr.clone()); - node.stdout_bytes = Some(props.stdout_bytes); - node.stderr_bytes = Some(props.stderr_bytes); - node.streams_separated = Some(props.streams_separated); - node.live_streaming = Some(props.live_streaming); - node.termination = Some(props.termination); - node.script_timing = Some(serde_json::to_value(props).map_err(|err| { - Error::InvalidEvent(format!("invalid command.completed payload: {err}")) - })?); + stage.stdout = Some(props.stdout.clone()); + stage.stderr = Some(props.stderr.clone()); + stage.stdout_bytes = Some(props.stdout_bytes); + stage.stderr_bytes = Some(props.stderr_bytes); + stage.streams_separated = Some(props.streams_separated); + stage.live_streaming = Some(props.live_streaming); + stage.termination = Some(props.termination); + stage.script_timing = Some(script_timing); } EventBody::ParallelCompleted(props) => { - let Some(node_id) = stored.node_id.as_deref() else { + let parallel_results = serde_json::to_value(&props.results).map_err(|err| { + Error::InvalidEvent(format!("invalid parallel.completed payload: {err}")) + })?; + let Some(stage) = stage_at_current_visit(self, stored, event.seq) else { return Ok(()); }; - stage_entry_with_current_visit(self, event, node_id).parallel_results = - Some(serde_json::to_value(&props.results).map_err(|err| { - Error::InvalidEvent(format!("invalid parallel.completed payload: {err}")) - })?); + stage.parallel_results = Some(parallel_results); } _ => {} } @@ -370,30 +388,24 @@ impl RunProjectionReducer for RunProjection { } } -fn stage_entry_for_event<'a>( +fn stage_at_visit<'a>( state: &'a mut RunProjection, - event: &EventEnvelope, - node_id: &str, + stored: &RunEvent, visit: u32, -) -> &'a mut fabro_types::StageState { - if let Some(stage_id) = event.event.stage_id.as_ref() { - state.stage_entry_id(stage_id, event.seq) - } else { - state.stage_entry(node_id, visit, event.seq) - } + seq: u32, +) -> Option<&'a mut StageProjection> { + let node_id = stored.node_id.as_deref()?; + Some(state.stage_entry(node_id, visit, first_event_seq(seq))) } -fn stage_entry_with_current_visit<'a>( +fn stage_at_current_visit<'a>( state: &'a mut RunProjection, - event: &EventEnvelope, - node_id: &str, -) -> &'a mut fabro_types::StageState { - if let Some(stage_id) = event.event.stage_id.as_ref() { - state.stage_entry_id(stage_id, event.seq) - } else { - let visit = state.current_visit_for(node_id).unwrap_or(1); - state.stage_entry(node_id, visit, event.seq) - } + stored: &RunEvent, + seq: u32, +) -> Option<&'a mut StageProjection> { + let node_id = stored.node_id.as_deref()?; + let visit = state.current_visit_for(node_id).unwrap_or(1); + Some(state.stage_entry(node_id, visit, first_event_seq(seq))) } pub(crate) fn build_summary(state: &RunProjection, run_id: &RunId) -> RunSummary { @@ -540,12 +552,12 @@ fn stage_outcome_from_props(props: &StageCompletedProps) -> Outcome>, timestamp: DateTime, -) -> NodeStatusRecord { - NodeStatusRecord { - status: outcome.status, +) -> StageCompletion { + StageCompletion { + outcome: outcome.status, notes: outcome.notes.clone(), failure_reason: outcome .failure @@ -595,18 +607,18 @@ fn provider_used_from_agent_cli_started(props: &AgentCliStartedProps) -> Value { #[cfg(test)] mod tests { - use std::collections::HashMap; + use std::collections::{BTreeMap, HashMap}; use chrono::Utc; use fabro_types::run_event::run::RunFailedProps; use fabro_types::run_event::{ - InterviewCompletedProps, InterviewOption, InterviewStartedProps, RunControlEffectProps, - StageStartedProps, + CheckpointCompletedProps, InterviewCompletedProps, InterviewOption, InterviewStartedProps, + RunControlEffectProps, StagePromptProps, StageStartedProps, }; use fabro_types::{ - BlockedReason, Checkpoint, EventBody, FailureReason, QuestionType, RunBlobId, - RunControlAction, RunEvent, RunStatus, StageState, SuccessReason, TerminalStatus, - WorkflowSettings, fixtures, + BlockedReason, Checkpoint, EventBody, FailureReason, Outcome, QuestionType, RunBlobId, + RunControlAction, RunEvent, RunStatus, StageOutcome, SuccessReason, TerminalStatus, + WorkflowSettings, first_event_seq, fixtures, }; use serde_json::json; @@ -633,6 +645,12 @@ mod tests { EventEnvelope { seq, event } } + fn test_stage_event(seq: u32, body: EventBody, stage_id: StageId) -> EventEnvelope { + let mut event = test_event(seq, body, Some(stage_id.node_id())); + event.event.stage_id = Some(stage_id); + event + } + fn test_raw_event( seq: u32, event: &str, @@ -675,7 +693,7 @@ mod tests { } #[test] - fn deserialize_projection_defaults_missing_nodes_and_checkpoints() { + fn deserialize_projection_defaults_missing_stages_and_checkpoints() { let state: RunProjection = serde_json::from_value(serde_json::json!({ "pending_control": "pause" })) @@ -687,7 +705,7 @@ mod tests { } #[test] - fn deserialize_and_round_trip_projection_preserves_stage_ids_and_pending_control() { + fn deserialize_and_round_trip_projection_preserves_stages_and_pending_control() { let state: RunProjection = serde_json::from_value(serde_json::json!({ "spec": { "run_id": "01JW6A7VNFZSFF0SKXJG29Z2M3", @@ -719,7 +737,7 @@ mod tests { ]], "stages": { "build@2": { - "seq": 0, + "first_event_seq": 1, "diff": "diff --git a/file b/file", "stdout": "done" } @@ -729,6 +747,7 @@ mod tests { let stage_id = StageId::new("build", 2); let node = state.stage(&stage_id).unwrap(); + assert_eq!(node.first_event_seq, first_event_seq(1)); assert_eq!(node.diff.as_deref(), Some("diff --git a/file b/file")); assert_eq!(state.list_node_visits("build"), vec![2]); assert_eq!(state.pending_control, Some(RunControlAction::Cancel)); @@ -745,12 +764,10 @@ mod tests { ); assert!(serialized.get("spec").is_some()); assert!(serialized.get("run").is_none()); - assert!(serialized.get("stages").is_some()); - assert!(serialized.get("nodes").is_none()); } #[test] - fn set_stage_round_trips_through_json() { + fn stage_entry_round_trips_through_json() { let mut state = RunProjection::default(); state.pending_control = Some(RunControlAction::Unpause); state.checkpoints = vec![(7, Checkpoint { @@ -766,10 +783,7 @@ mod tests { restart_failure_signatures: HashMap::new(), node_visits: HashMap::from([("build".to_string(), 2usize)]), })]; - state.set_stage(StageId::new("build", 2), StageState { - stdout: Some("done".to_string()), - ..StageState::default() - }); + state.stage_entry("build", 2, first_event_seq(7)).stdout = Some("done".to_string()); let round_tripped: RunProjection = serde_json::from_value(serde_json::to_value(&state).unwrap()).unwrap(); @@ -790,25 +804,97 @@ mod tests { } #[test] - fn stage_started_stamps_seq_from_stage_id() { + fn stage_started_sets_first_event_seq() { let mut state = RunProjection::default(); - let stage_id = StageId::new("build", 2); - let mut event = test_event( - 17, - EventBody::StageStarted(StageStartedProps { - index: 0, - handler_type: "command".to_string(), - attempt: 3, - max_attempts: 3, - }), - Some("build"), - ); - event.event.stage_id = Some(stage_id.clone()); + let stage_id = StageId::new("build", 1); - state.apply_event(&event).unwrap(); + state + .apply_event(&test_stage_event( + 3, + EventBody::StageStarted(StageStartedProps { + index: 0, + handler_type: "agent".to_string(), + attempt: 1, + max_attempts: 1, + }), + stage_id.clone(), + )) + .unwrap(); let stage = state.stage(&stage_id).unwrap(); - assert_eq!(stage.seq, 17); + assert_eq!(stage.first_event_seq, first_event_seq(3)); + } + + #[test] + fn later_stage_events_do_not_overwrite_first_event_seq() { + let mut state = RunProjection::default(); + let stage_id = StageId::new("build", 1); + + state + .apply_event(&test_stage_event( + 3, + EventBody::StageStarted(StageStartedProps { + index: 0, + handler_type: "agent".to_string(), + attempt: 1, + max_attempts: 1, + }), + stage_id.clone(), + )) + .unwrap(); + state + .apply_event(&test_event( + 4, + EventBody::StagePrompt(StagePromptProps { + visit: 1, + text: "prompt".to_string(), + mode: None, + provider: None, + model: None, + }), + Some("build"), + )) + .unwrap(); + + let stage = state.stage(&stage_id).unwrap(); + assert_eq!(stage.first_event_seq, first_event_seq(3)); + assert_eq!(stage.prompt.as_deref(), Some("prompt")); + } + + #[test] + fn checkpoint_completed_creates_projection_entry_for_skipped_stage() { + let mut state = RunProjection::default(); + let stage_id = StageId::new("skip_me", 1); + + state + .apply_event(&test_event( + 5, + EventBody::CheckpointCompleted(CheckpointCompletedProps { + status: "running".to_string(), + current_node: "next".to_string(), + completed_nodes: vec!["skip_me".to_string()], + node_retries: BTreeMap::new(), + context_values: BTreeMap::new(), + node_outcomes: BTreeMap::from([( + "skip_me".to_string(), + Outcome::skipped("condition was false"), + )]), + next_node_id: Some("next".to_string()), + git_commit_sha: None, + loop_failure_signatures: BTreeMap::new(), + restart_failure_signatures: BTreeMap::new(), + node_visits: BTreeMap::from([("skip_me".to_string(), 1usize)]), + diff: None, + }), + None, + )) + .unwrap(); + + let stage = state.stage(&stage_id).unwrap(); + assert_eq!(stage.first_event_seq, first_event_seq(5)); + let completion = stage.completion.as_ref().unwrap(); + assert_eq!(completion.outcome, StageOutcome::Skipped); + assert_eq!(completion.notes.as_deref(), Some("condition was false")); } #[test] diff --git a/lib/crates/fabro-store/src/serializable_projection.rs b/lib/crates/fabro-store/src/serializable_projection.rs index face8b695..7d888b4a3 100644 --- a/lib/crates/fabro-store/src/serializable_projection.rs +++ b/lib/crates/fabro-store/src/serializable_projection.rs @@ -10,23 +10,12 @@ impl Serialize for SerializableProjection<'_> { S: Serializer, { let mut projection = self.0.clone(); - let stage_ids: Vec<_> = projection - .iter_stages() - .map(|(stage_id, _)| stage_id.clone()) - .collect(); - - for stage_id in stage_ids { - let Some(node) = projection.stage(&stage_id).cloned() else { - continue; - }; - projection.set_stage(stage_id, crate::StageState { - prompt: None, - response: None, - diff: None, - stdout: None, - stderr: None, - ..node - }); + for (_, stage) in projection.iter_stages_mut() { + stage.prompt = None; + stage.response = None; + stage.diff = None; + stage.stdout = None; + stage.stderr = None; } projection.serialize(serializer) diff --git a/lib/crates/fabro-store/tests/serializable_projection.rs b/lib/crates/fabro-store/tests/serializable_projection.rs index eeeffe443..6d51065ba 100644 --- a/lib/crates/fabro-store/tests/serializable_projection.rs +++ b/lib/crates/fabro-store/tests/serializable_projection.rs @@ -1,12 +1,12 @@ use std::collections::{BTreeMap, HashMap}; use chrono::{TimeZone, Utc}; -use fabro_store::{RunProjection, SerializableProjection, StageId, StageState}; +use fabro_store::{RunProjection, SerializableProjection, StageId}; use fabro_types::graph::Graph; use fabro_types::run::RunSpec; use fabro_types::{ - Checkpoint, NodeStatusRecord, RunStatus, SandboxRecord, StageOutcome, StartRecord, - TerminalStatus, WorkflowSettings, fixtures, + Checkpoint, RunStatus, SandboxRecord, StageCompletion, StageOutcome, StartRecord, + TerminalStatus, WorkflowSettings, first_event_seq, fixtures, }; use serde_json::json; @@ -77,32 +77,25 @@ fn serializable_projection_round_trips_and_trims_bulky_node_fields() { clone_branch: None, }); projection.pending_interviews = BTreeMap::new(); - projection.set_stage(stage_id.clone(), StageState { - seq: 7, - prompt: Some("plan the work".to_string()), - response: Some("done".to_string()), - status: Some(NodeStatusRecord { - status: StageOutcome::Succeeded, - notes: Some("ok".to_string()), - failure_reason: None, - timestamp: Utc - .with_ymd_and_hms(2026, 4, 20, 12, 1, 0) - .single() - .expect("timestamp should be representable"), - }), - provider_used: Some(json!({ "provider": "openai", "model": "gpt-5.4" })), - diff: Some("diff --git a/a b/a".to_string()), - script_invocation: Some(json!({ "command": "cargo test" })), - script_timing: Some(json!({ "duration_ms": 10 })), - parallel_results: Some(json!([{ "stage": "fanout@1" }])), - stdout: Some("stdout".to_string()), - stderr: Some("stderr".to_string()), - stdout_bytes: None, - stderr_bytes: None, - streams_separated: None, - live_streaming: None, - termination: None, + let stage = projection.stage_entry(stage_id.node_id(), stage_id.visit(), first_event_seq(2)); + stage.prompt = Some("plan the work".to_string()); + stage.response = Some("done".to_string()); + stage.completion = Some(StageCompletion { + outcome: StageOutcome::Succeeded, + notes: Some("ok".to_string()), + failure_reason: None, + timestamp: Utc + .with_ymd_and_hms(2026, 4, 20, 12, 1, 0) + .single() + .expect("timestamp should be representable"), }); + stage.provider_used = Some(json!({ "provider": "openai", "model": "gpt-5.4" })); + stage.diff = Some("diff --git a/a b/a".to_string()); + stage.script_invocation = Some(json!({ "command": "cargo test" })); + stage.script_timing = Some(json!({ "duration_ms": 10 })); + stage.parallel_results = Some(json!([{ "stage": "fanout@1" }])); + stage.stdout = Some("stdout".to_string()); + stage.stderr = Some("stderr".to_string()); let serialized = serde_json::to_value(SerializableProjection(&projection)) .expect("projection should serialize"); @@ -120,12 +113,18 @@ fn serializable_projection_round_trips_and_trims_bulky_node_fields() { ); assert_eq!(round_tripped.status(), Some(RunStatus::Running)); assert!(!round_tripped.is_terminal()); - assert_eq!(node.seq, 7); assert_eq!(node.prompt, None); assert_eq!(node.response, None); assert_eq!(node.diff, None); assert_eq!(node.stdout, None); assert_eq!(node.stderr, None); + assert_eq!(node.first_event_seq, first_event_seq(2)); + assert_eq!( + node.completion + .as_ref() + .map(|completion| completion.outcome), + Some(StageOutcome::Succeeded) + ); assert_eq!( node.provider_used, Some(json!({ "provider": "openai", "model": "gpt-5.4" })) diff --git a/lib/crates/fabro-types/src/lib.rs b/lib/crates/fabro-types/src/lib.rs index 73ae59eee..01f8e08f9 100644 --- a/lib/crates/fabro-types/src/lib.rs +++ b/lib/crates/fabro-types/src/lib.rs @@ -13,7 +13,6 @@ pub mod event_envelope; pub mod failure_signature; pub mod graph; pub mod interview; -pub mod node_status; pub mod outcome; pub mod pull_request; pub mod repository; @@ -27,6 +26,7 @@ pub mod run_summary; pub mod sandbox_record; pub mod secret; pub mod settings; +pub mod stage_completion; pub mod stage_id; pub mod start; pub mod status; @@ -49,9 +49,8 @@ pub use event_envelope::EventEnvelope; pub use failure_signature::FailureSignature; pub use graph::{AttrValue, Edge, Graph, Node, is_llm_handler_type, shape_to_handler_type}; pub use interview::{InterviewQuestionRecord, QuestionType}; -pub use node_status::NodeStatusRecord; pub use outcome::{ - FailureCategory, FailureDetail, NodeResult, Outcome, OutcomeMeta, StageOutcome, StageStatus, + FailureCategory, FailureDetail, NodeResult, Outcome, OutcomeMeta, StageOutcome, StageState, }; pub use pull_request::{ PullRequestDetail, PullRequestGithubDetail, PullRequestRecord, PullRequestRef, PullRequestUser, @@ -71,10 +70,11 @@ pub use run_event::{ MetadataSnapshotPhase, RunEvent, RunNoticeLevel, }; pub use run_id::{RunId, fixtures}; -pub use run_projection::{PendingInterviewRecord, RunProjection, StageState}; +pub use run_projection::{PendingInterviewRecord, RunProjection, StageProjection, first_event_seq}; pub use run_summary::RunSummary; pub use sandbox_record::SandboxRecord; pub use secret::{SecretMetadata, SecretType}; +pub use stage_completion::StageCompletion; pub use stage_id::{ParallelBranchId, StageId}; pub use start::StartRecord; pub use status::{ diff --git a/lib/crates/fabro-types/src/outcome.rs b/lib/crates/fabro-types/src/outcome.rs index 16bbff240..60ad9a811 100644 --- a/lib/crates/fabro-types/src/outcome.rs +++ b/lib/crates/fabro-types/src/outcome.rs @@ -106,7 +106,7 @@ impl<'de> Deserialize<'de> for StageOutcome { )] #[serde(rename_all = "snake_case")] #[strum(serialize_all = "snake_case")] -pub enum StageStatus { +pub enum StageState { Pending, Running, Retrying, @@ -117,7 +117,7 @@ pub enum StageStatus { Cancelled, } -impl StageStatus { +impl StageState { #[must_use] pub fn is_terminal(self) -> bool { matches!( @@ -131,7 +131,7 @@ impl StageStatus { } } -impl From for StageStatus { +impl From for StageState { fn from(outcome: StageOutcome) -> Self { match outcome { StageOutcome::Succeeded => Self::Succeeded, @@ -301,7 +301,7 @@ impl Outcome { mod tests { use serde_json::json; - use super::{StageOutcome, StageStatus}; + use super::{StageOutcome, StageState}; #[test] fn stage_outcome_failed_serde_is_lossy_for_retry_intent() { @@ -321,23 +321,23 @@ mod tests { } #[test] - fn stage_status_projects_terminal_outcomes() { + fn stage_state_projects_terminal_outcomes() { assert_eq!( - StageStatus::from(StageOutcome::Succeeded), - StageStatus::Succeeded + StageState::from(StageOutcome::Succeeded), + StageState::Succeeded ); assert_eq!( - StageStatus::from(StageOutcome::PartiallySucceeded), - StageStatus::PartiallySucceeded + StageState::from(StageOutcome::PartiallySucceeded), + StageState::PartiallySucceeded ); assert_eq!( - StageStatus::from(StageOutcome::Failed { + StageState::from(StageOutcome::Failed { retry_requested: true, }), - StageStatus::Failed + StageState::Failed ); - assert!(StageStatus::Cancelled.is_terminal()); - assert!(!StageStatus::Running.is_terminal()); + assert!(StageState::Cancelled.is_terminal()); + assert!(!StageState::Running.is_terminal()); } } diff --git a/lib/crates/fabro-types/src/run_projection.rs b/lib/crates/fabro-types/src/run_projection.rs index 67cf7de4c..de2e4a1a5 100644 --- a/lib/crates/fabro-types/src/run_projection.rs +++ b/lib/crates/fabro-types/src/run_projection.rs @@ -1,10 +1,11 @@ use std::collections::{BTreeMap, HashMap}; +use std::num::NonZeroU32; use chrono::{DateTime, Utc}; use crate::{ - Checkpoint, Conclusion, InterviewQuestionRecord, InvalidTransition, NodeStatusRecord, - PullRequestRecord, Retro, RunControlAction, RunId, RunSpec, RunStatus, SandboxRecord, StageId, + Checkpoint, Conclusion, InterviewQuestionRecord, InvalidTransition, PullRequestRecord, Retro, + RunControlAction, RunId, RunSpec, RunStatus, SandboxRecord, StageCompletion, StageId, StartRecord, }; @@ -28,8 +29,7 @@ pub struct RunProjection { pub pull_request: Option, pub superseded_by: Option, pub pending_interviews: BTreeMap, - #[serde(alias = "nodes")] - stages: HashMap, + stages: HashMap, } #[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)] @@ -38,13 +38,12 @@ pub struct PendingInterviewRecord { pub started_at: Option>, } -#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)] -pub struct StageState { - #[serde(default)] - pub seq: u32, +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +pub struct StageProjection { + pub first_event_seq: NonZeroU32, pub prompt: Option, pub response: Option, - pub status: Option, + pub completion: Option, pub provider_used: Option, pub diff: Option, pub script_invocation: Option, @@ -64,21 +63,56 @@ pub struct StageState { pub termination: Option, } +/// Convert a 1-based event sequence number into the `NonZeroU32` form used for +/// `StageProjection::first_event_seq`. Run event seqs always start at 1. +#[must_use] +pub fn first_event_seq(seq: u32) -> NonZeroU32 { + NonZeroU32::new(seq).expect("event seq starts at 1") +} + +impl StageProjection { + #[must_use] + pub fn new(first_event_seq: NonZeroU32) -> Self { + Self { + first_event_seq, + prompt: None, + response: None, + completion: None, + provider_used: None, + diff: None, + script_invocation: None, + script_timing: None, + parallel_results: None, + stdout: None, + stderr: None, + stdout_bytes: None, + stderr_bytes: None, + streams_separated: None, + live_streaming: None, + termination: None, + } + } +} + impl RunProjection { - pub fn stage(&self, stage_id: &StageId) -> Option<&StageState> { - self.stages.get(stage_id) + pub fn stage(&self, stage: &StageId) -> Option<&StageProjection> { + self.stages.get(stage) } - pub fn iter_stages(&self) -> impl Iterator { + pub fn iter_stages(&self) -> impl Iterator { self.stages.iter() } + pub fn iter_stages_mut(&mut self) -> impl Iterator { + self.stages.iter_mut() + } + pub fn is_empty(&self) -> bool { self.stages.is_empty() } - pub fn set_stage(&mut self, stage_id: StageId, state: StageState) { - self.stages.insert(stage_id, state); + pub fn stage_mut(&mut self, stage: &StageId) -> Option<&mut StageProjection> { + self.stages.get_mut(stage) } pub fn list_node_visits(&self, node_id: &str) -> Vec { @@ -113,17 +147,15 @@ impl RunProjection { &self.pending_interviews } - pub fn stage_entry_id(&mut self, stage_id: &StageId, seq: u32) -> &mut StageState { + pub fn stage_entry( + &mut self, + node_id: &str, + visit: u32, + first_event_seq: NonZeroU32, + ) -> &mut StageProjection { self.stages - .entry(stage_id.clone()) - .or_insert_with(|| StageState { - seq, - ..Default::default() - }) - } - - pub fn stage_entry(&mut self, node_id: &str, visit: u32, seq: u32) -> &mut StageState { - self.stage_entry_id(&StageId::new(node_id, visit), seq) + .entry(StageId::new(node_id, visit)) + .or_insert_with(|| StageProjection::new(first_event_seq)) } pub fn current_visit_for(&self, node_id: &str) -> Option { diff --git a/lib/crates/fabro-types/src/node_status.rs b/lib/crates/fabro-types/src/stage_completion.rs similarity index 82% rename from lib/crates/fabro-types/src/node_status.rs rename to lib/crates/fabro-types/src/stage_completion.rs index 9b6c8b4d9..12312f2fd 100644 --- a/lib/crates/fabro-types/src/node_status.rs +++ b/lib/crates/fabro-types/src/stage_completion.rs @@ -4,8 +4,8 @@ use serde::{Deserialize, Serialize}; use crate::outcome::StageOutcome; #[derive(Debug, Clone, Serialize, Deserialize)] -pub struct NodeStatusRecord { - pub status: StageOutcome, +pub struct StageCompletion { + pub outcome: StageOutcome, #[serde(default)] pub notes: Option, #[serde(default)] diff --git a/lib/crates/fabro-workflow/src/handler/agent.rs b/lib/crates/fabro-workflow/src/handler/agent.rs index 15c220c5a..b685fabd6 100644 --- a/lib/crates/fabro-workflow/src/handler/agent.rs +++ b/lib/crates/fabro-workflow/src/handler/agent.rs @@ -505,9 +505,9 @@ mod tests { logger.flush().await; let state = run_store.state().await.unwrap(); - let stage_state = state.stage(&StageId::new("plan", 1)).unwrap(); + let node_state = state.stage(&StageId::new("plan", 1)).unwrap(); assert_eq!( - stage_state.prompt.as_deref(), + node_state.prompt.as_deref(), Some("Achieve: Build a feature") ); } @@ -532,8 +532,8 @@ mod tests { logger.flush().await; let state = run_store.state().await.unwrap(); - let stage_state = state.stage(&StageId::new("work", 1)).unwrap(); - assert_eq!(stage_state.prompt.as_deref(), Some("Do work")); + let node_state = state.stage(&StageId::new("work", 1)).unwrap(); + assert_eq!(node_state.prompt.as_deref(), Some("Do work")); } #[tokio::test] @@ -774,9 +774,9 @@ mod tests { logger.flush().await; let state = run_store.state().await.unwrap(); - let stage_state = state.stage(&StageId::new("step", 1)).unwrap(); + let node_state = state.stage(&StageId::new("step", 1)).unwrap(); assert_eq!( - stage_state.provider_used.as_ref().unwrap()["provider"], + node_state.provider_used.as_ref().unwrap()["provider"], "openai" ); } @@ -1262,8 +1262,8 @@ Some text in between. logger.flush().await; let state = run_store.state().await.unwrap(); - let stage_state = state.stage(&StageId::new("report", 1)).unwrap(); - let prompt_content = stage_state.prompt.as_deref().unwrap(); + let node_state = state.stage(&StageId::new("report", 1)).unwrap(); + let prompt_content = node_state.prompt.as_deref().unwrap(); assert!( prompt_content.contains("## Script Output\nAll tests passed"), "prompt.md should contain preamble" diff --git a/lib/crates/fabro-workflow/src/handler/command.rs b/lib/crates/fabro-workflow/src/handler/command.rs index 819e0c3ac..6960ba4d9 100644 --- a/lib/crates/fabro-workflow/src/handler/command.rs +++ b/lib/crates/fabro-workflow/src/handler/command.rs @@ -509,8 +509,8 @@ mod tests { logger.flush().await; let snapshot = run_store.state().await.unwrap(); - let stage_state = snapshot.stage(&StageId::new("script_node", 1)).unwrap(); - let json = stage_state.script_invocation.as_ref().unwrap(); + let node_state = snapshot.stage(&StageId::new("script_node", 1)).unwrap(); + let json = node_state.script_invocation.as_ref().unwrap(); assert_eq!(json["command"], "echo hello"); assert_eq!(json["language"], "shell"); assert_eq!(json["timeout_ms"], serde_json::Value::Null); @@ -540,8 +540,8 @@ mod tests { logger.flush().await; let snapshot = run_store.state().await.unwrap(); - let stage_state = snapshot.stage(&StageId::new("script_node", 1)).unwrap(); - let json = stage_state.script_invocation.as_ref().unwrap(); + let node_state = snapshot.stage(&StageId::new("script_node", 1)).unwrap(); + let json = node_state.script_invocation.as_ref().unwrap(); assert_eq!(json["command"], "echo hello"); assert_eq!(json["language"], "shell"); assert_eq!(json["timeout_ms"], 5000); @@ -567,15 +567,15 @@ mod tests { logger.flush().await; let snapshot = run_store.state().await.unwrap(); - let stage_state = snapshot.stage(&StageId::new("script_node", 1)).unwrap(); - let stdout = stage_state.stdout.as_deref().unwrap(); + let node_state = snapshot.stage(&StageId::new("script_node", 1)).unwrap(); + let stdout = node_state.stdout.as_deref().unwrap(); assert_eq!(command_log_text(&services, stdout).await.trim(), "hello"); - let stderr = stage_state.stderr.as_deref().unwrap(); + let stderr = node_state.stderr.as_deref().unwrap(); assert_eq!(command_log_text(&services, stderr).await, ""); - assert_eq!(stage_state.stdout_bytes, Some(6)); - assert_eq!(stage_state.stderr_bytes, Some(0)); - assert_eq!(stage_state.streams_separated, Some(true)); - assert_eq!(stage_state.live_streaming, Some(true)); + assert_eq!(node_state.stdout_bytes, Some(6)); + assert_eq!(node_state.stderr_bytes, Some(0)); + assert_eq!(node_state.streams_separated, Some(true)); + assert_eq!(node_state.live_streaming, Some(true)); } #[tokio::test] @@ -598,8 +598,8 @@ mod tests { logger.flush().await; let snapshot = run_store.state().await.unwrap(); - let stage_state = snapshot.stage(&StageId::new("script_node", 1)).unwrap(); - let stderr = stage_state.stderr.as_deref().unwrap(); + let node_state = snapshot.stage(&StageId::new("script_node", 1)).unwrap(); + let stderr = node_state.stderr.as_deref().unwrap(); assert_eq!(command_log_text(&services, stderr).await.trim(), "oops"); } @@ -623,8 +623,8 @@ mod tests { logger.flush().await; let snapshot = run_store.state().await.unwrap(); - let stage_state = snapshot.stage(&StageId::new("script_node", 1)).unwrap(); - let json = stage_state.script_timing.as_ref().unwrap(); + let node_state = snapshot.stage(&StageId::new("script_node", 1)).unwrap(); + let json = node_state.script_timing.as_ref().unwrap(); assert!(json["duration_ms"].is_u64()); assert_eq!(json["exit_code"], 0); assert_eq!(json["termination"], "exited"); @@ -648,8 +648,8 @@ mod tests { logger.flush().await; let snapshot = run_store.state().await.unwrap(); - let stage_state = snapshot.stage(&StageId::new("script_node", 1)).unwrap(); - let json = stage_state.script_timing.as_ref().unwrap(); + let node_state = snapshot.stage(&StageId::new("script_node", 1)).unwrap(); + let json = node_state.script_timing.as_ref().unwrap(); assert_eq!(json["exit_code"], 1); assert_eq!(json["termination"], "exited"); } @@ -678,8 +678,8 @@ mod tests { logger.flush().await; let snapshot = run_store.state().await.unwrap(); - let stage_state = snapshot.stage(&StageId::new("script_node", 1)).unwrap(); - let json = stage_state.script_timing.as_ref().unwrap(); + let node_state = snapshot.stage(&StageId::new("script_node", 1)).unwrap(); + let json = node_state.script_timing.as_ref().unwrap(); assert!(json["duration_ms"].is_u64()); assert_eq!(json["exit_code"], serde_json::Value::Null); assert_eq!(json["termination"], "timed_out"); diff --git a/lib/crates/fabro-workflow/src/handler/parallel.rs b/lib/crates/fabro-workflow/src/handler/parallel.rs index 9e703d0e6..92168c79c 100644 --- a/lib/crates/fabro-workflow/src/handler/parallel.rs +++ b/lib/crates/fabro-workflow/src/handler/parallel.rs @@ -736,8 +736,8 @@ mod tests { assert!(results.is_some()); let state = run_store.state().await.unwrap(); - let stage_state = state.stage(&StageId::new("par", 1)).unwrap(); - let parsed = stage_state.parallel_results.as_ref().unwrap(); + let node_state = state.stage(&StageId::new("par", 1)).unwrap(); + let parsed = node_state.parallel_results.as_ref().unwrap(); assert!( parsed.is_array(), "parallel_results.json should be a JSON array" @@ -781,8 +781,8 @@ mod tests { logger.flush().await; let state = run_store.state().await.unwrap(); - let stage_state = state.stage(&fabro_store::StageId::new("par", 1)).unwrap(); - let results = stage_state.parallel_results.as_ref().unwrap(); + let node_state = state.stage(&fabro_store::StageId::new("par", 1)).unwrap(); + let results = node_state.parallel_results.as_ref().unwrap(); assert!(results.is_array()); assert_eq!(results.as_array().unwrap().len(), 2); } diff --git a/lib/crates/fabro-workflow/src/handler/prompt.rs b/lib/crates/fabro-workflow/src/handler/prompt.rs index 37b190220..a76099b80 100644 --- a/lib/crates/fabro-workflow/src/handler/prompt.rs +++ b/lib/crates/fabro-workflow/src/handler/prompt.rs @@ -367,11 +367,8 @@ mod tests { logger.flush().await; let state = run_store.state().await.unwrap(); - let stage_state = state.stage(&StageId::new("classify", 1)).unwrap(); - assert_eq!( - stage_state.provider_used.as_ref().unwrap()["mode"], - "prompt" - ); + let node_state = state.stage(&StageId::new("classify", 1)).unwrap(); + assert_eq!(node_state.provider_used.as_ref().unwrap()["mode"], "prompt"); } struct OneShotCapturingBackend { diff --git a/lib/crates/fabro-workflow/src/outcome.rs b/lib/crates/fabro-workflow/src/outcome.rs index 9092abd86..6fad164aa 100644 --- a/lib/crates/fabro-workflow/src/outcome.rs +++ b/lib/crates/fabro-workflow/src/outcome.rs @@ -1,5 +1,5 @@ pub use fabro_core::outcome::{ - FailureCategory, FailureDetail, OutcomeMeta, StageOutcome, StageStatus, + FailureCategory, FailureDetail, OutcomeMeta, StageOutcome, StageState, }; use fabro_llm::types::TokenCounts as LlmTokenCounts; use fabro_model::{ diff --git a/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs b/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs index 58e1919c0..07dc8d31d 100644 --- a/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs +++ b/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs @@ -772,7 +772,7 @@ async fn execute_persists_start_record_and_node_status() { let node = state.stage(&fabro_store::StageId::new("start", 1)).unwrap(); assert_eq!( - node.status.as_ref().unwrap().status, + node.completion.as_ref().unwrap().outcome, StageOutcome::Succeeded ); } @@ -827,10 +827,10 @@ async fn timeout_causes_fail_status_record() { let status = state .stage(&fabro_store::StageId::new("work", 1)) .unwrap() - .status + .completion .as_ref() .unwrap(); - assert_eq!(status.status, StageOutcome::Failed { + assert_eq!(status.outcome, StageOutcome::Failed { retry_requested: false, }); } diff --git a/lib/crates/fabro-workflow/src/pipeline/finalize.rs b/lib/crates/fabro-workflow/src/pipeline/finalize.rs index 5d1c34979..7e48de0a0 100644 --- a/lib/crates/fabro-workflow/src/pipeline/finalize.rs +++ b/lib/crates/fabro-workflow/src/pipeline/finalize.rs @@ -1,9 +1,10 @@ +use std::collections::HashMap; use std::time::Instant; use fabro_dump::RunDump; use fabro_hooks::{HookContext, HookEvent}; use fabro_types::run_event::{MetadataSnapshotFailureKind, MetadataSnapshotPhase}; -use fabro_types::{BilledTokenCounts, EventBody}; +use fabro_types::{BilledTokenCounts, EventBody, RunProjection}; use fabro_util::error::collect_causes; use fabro_util::time::elapsed_ms; @@ -68,14 +69,22 @@ pub(crate) async fn build_conclusion_from_store( final_git_commit_sha: Option, ) -> Conclusion { let (state_result, events_result) = tokio::join!(run_store.state(), run_store.list_events()); - let checkpoint = state_result.ok().and_then(|state| state.checkpoint); + let projection = state_result.ok(); + let projection_order = projection + .as_ref() + .map(stage_projection_order) + .unwrap_or_default(); + let checkpoint = projection + .as_ref() + .and_then(|state| state.checkpoint.as_ref()); let stage_durations = events_result .map(|events| crate::extract_stage_durations_from_events(&events)) .unwrap_or_default(); build_conclusion_from_parts( - checkpoint.as_ref(), + checkpoint, &stage_durations, + &projection_order, status, failure_reason, run_duration_ms, @@ -85,7 +94,8 @@ pub(crate) async fn build_conclusion_from_store( fn build_conclusion_from_parts( checkpoint: Option<&Checkpoint>, - stage_durations: &std::collections::HashMap, + stage_durations: &HashMap, + projection_order: &HashMap, status: StageOutcome, failure_reason: Option, run_duration_ms: u64, @@ -95,14 +105,31 @@ fn build_conclusion_from_parts( // while the other checkpoint maps are keyed by node_id. Dedupe to one row // per node so the stages table matches the deduped billing total. let (stages, total_retries) = if let Some(cp) = checkpoint { - let mut stages = Vec::new(); + let mut stage_rows = Vec::new(); let mut seen = std::collections::HashSet::new(); let mut retries_sum: u32 = 0; + let mut stage_order = Vec::new(); - for node_id in &cp.completed_nodes { + for (original_checkpoint_order, node_id) in cp.completed_nodes.iter().enumerate() { if !seen.insert(node_id.as_str()) { continue; } + stage_order.push((original_checkpoint_order, node_id.as_str())); + } + let mut extra_node_outcomes = cp + .node_outcomes + .keys() + .filter(|node_id| !seen.contains(node_id.as_str())) + .map(String::as_str) + .collect::>(); + extra_node_outcomes.sort_unstable(); + let extra_offset = stage_order.len(); + for (extra_index, node_id) in extra_node_outcomes.into_iter().enumerate() { + seen.insert(node_id); + stage_order.push((extra_offset + extra_index, node_id)); + } + + for (original_checkpoint_order, node_id) in stage_order { let outcome = cp.node_outcomes.get(node_id); let retries = cp .node_retries @@ -112,16 +139,31 @@ fn build_conclusion_from_parts( .saturating_sub(1); retries_sum += retries; - stages.push(StageSummary { - stage_id: node_id.clone(), - stage_label: node_id.clone(), + let summary = StageSummary { + stage_id: node_id.to_string(), + stage_label: node_id.to_string(), duration_ms: stage_durations.get(node_id).copied().unwrap_or(0), billing_usd_micros: outcome .and_then(|o| o.usage.as_ref()) .and_then(|usage| usage.total_usd_micros), retries, - }); + }; + stage_rows.push(( + projection_order.get(node_id).copied().unwrap_or(u32::MAX), + original_checkpoint_order, + summary, + )); } + stage_rows.sort_by(|left, right| { + left.0 + .cmp(&right.0) + .then_with(|| left.1.cmp(&right.1)) + .then_with(|| left.2.stage_id.cmp(&right.2.stage_id)) + }); + let stages = stage_rows + .into_iter() + .map(|(_, _, summary)| summary) + .collect(); (stages, retries_sum) } else { (vec![], 0) @@ -139,6 +181,19 @@ fn build_conclusion_from_parts( } } +fn stage_projection_order(state: &RunProjection) -> HashMap { + let mut order = HashMap::new(); + for (stage_id, stage) in state.iter_stages() { + order + .entry(stage_id.node_id().to_string()) + .and_modify(|first_seq: &mut u32| { + *first_seq = (*first_seq).min(stage.first_event_seq.get()); + }) + .or_insert_with(|| stage.first_event_seq.get()); + } + order +} + /// `conclusion` is injected because the terminal event hasn't been emitted /// yet — the run store's `projection.conclusion` is still `None` at this point. pub async fn write_finalize_commit( @@ -457,15 +512,18 @@ pub async fn finalize(retroed: Retroed, options: &FinalizeOptions) -> Result, + node_outcomes: HashMap, + ) -> Checkpoint { + Checkpoint { + timestamp: chrono::Utc::now(), + current_node: completed_nodes + .last() + .copied() + .unwrap_or("start") + .to_string(), + completed_nodes: completed_nodes.into_iter().map(str::to_string).collect(), + node_retries: HashMap::new(), + context_values: HashMap::new(), + node_outcomes, + next_node_id: None, + git_commit_sha: None, + loop_failure_signatures: HashMap::new(), + restart_failure_signatures: HashMap::new(), + node_visits: HashMap::new(), + } + } + + #[test] + fn conclusion_stage_order_follows_projection_first_event_order() { + let mut projection = RunProjection::default(); + projection.stage_entry("zebra", 1, first_event_seq(1)); + projection.stage_entry("apple", 1, first_event_seq(2)); + let projection_order = stage_projection_order(&projection); + let checkpoint = checkpoint_with( + vec!["apple", "zebra"], + HashMap::from([ + ("apple".to_string(), Outcome::success()), + ("zebra".to_string(), Outcome::success()), + ]), + ); + + let conclusion = build_conclusion_from_parts( + Some(&checkpoint), + &HashMap::new(), + &projection_order, + StageOutcome::Succeeded, + None, + 10, + None, + ); + + let stage_ids = conclusion + .stages + .iter() + .map(|stage| stage.stage_id.as_str()) + .collect::>(); + assert_eq!(stage_ids, vec!["zebra", "apple"]); + } + + #[test] + fn conclusion_includes_skipped_stage_from_projection_checkpoint_fallback() { + let mut projection = RunProjection::default(); + projection.stage_entry("skipped", 1, first_event_seq(4)); + projection.stage_entry("finished", 1, first_event_seq(5)); + let projection_order = stage_projection_order(&projection); + let checkpoint = checkpoint_with( + vec!["finished"], + HashMap::from([ + ("finished".to_string(), Outcome::success()), + ( + "skipped".to_string(), + Outcome::skipped("condition was false"), + ), + ]), + ); + + let conclusion = build_conclusion_from_parts( + Some(&checkpoint), + &HashMap::new(), + &projection_order, + StageOutcome::Succeeded, + None, + 10, + None, + ); + + let stage_ids = conclusion + .stages + .iter() + .map(|stage| stage.stage_id.as_str()) + .collect::>(); + assert_eq!(stage_ids, vec!["skipped", "finished"]); + } + fn test_services( run_store: RunStoreHandle, emitter: Arc, diff --git a/lib/crates/fabro-workflow/src/pipeline/pull_request.rs b/lib/crates/fabro-workflow/src/pipeline/pull_request.rs index b58740a2a..933f53052 100644 --- a/lib/crates/fabro-workflow/src/pipeline/pull_request.rs +++ b/lib/crates/fabro-workflow/src/pipeline/pull_request.rs @@ -641,7 +641,7 @@ mod tests { AggregateStats, FrictionKind, FrictionPoint, OpenItem, OpenItemKind, StageRetro, }; use fabro_store::Database; - use fabro_types::{BilledTokenCounts, RunSpec, SuccessReason, fixtures}; + use fabro_types::{BilledTokenCounts, RunSpec, SuccessReason, first_event_seq, fixtures}; use fabro_vault::{SecretType, Vault}; use futures::stream; use httpmock::Method::POST; @@ -997,13 +997,8 @@ mod tests { #[test] fn read_plan_text_found() { let mut state = RunProjection::default(); - state.set_stage( - fabro_store::StageId::new("plan", 1), - fabro_store::StageState { - response: Some("This is the plan".to_string()), - ..Default::default() - }, - ); + state.stage_entry("plan", 1, first_event_seq(1)).response = + Some("This is the plan".to_string()); let result = read_plan_text(&state); assert_eq!(result, Some("This is the plan".to_string())); @@ -1012,13 +1007,9 @@ mod tests { #[test] fn read_plan_text_prefix_match() { let mut state = RunProjection::default(); - state.set_stage( - fabro_store::StageId::new("planning", 1), - fabro_store::StageState { - response: Some("Planning content".to_string()), - ..Default::default() - }, - ); + state + .stage_entry("planning", 1, first_event_seq(1)) + .response = Some("Planning content".to_string()); let result = read_plan_text(&state); assert_eq!(result, Some("Planning content".to_string())); @@ -1027,20 +1018,11 @@ mod tests { #[test] fn read_plan_text_prefers_alphabetically_first_plan_node() { let mut state = RunProjection::default(); - state.set_stage( - fabro_store::StageId::new("planning", 1), - fabro_store::StageState { - response: Some("Planning content".to_string()), - ..Default::default() - }, - ); - state.set_stage( - fabro_store::StageId::new("plan", 1), - fabro_store::StageState { - response: Some("Plan content".to_string()), - ..Default::default() - }, - ); + state + .stage_entry("planning", 1, first_event_seq(1)) + .response = Some("Planning content".to_string()); + state.stage_entry("plan", 1, first_event_seq(2)).response = + Some("Plan content".to_string()); let result = read_plan_text(&state); assert_eq!(result, Some("Plan content".to_string())); @@ -1049,10 +1031,7 @@ mod tests { #[test] fn read_plan_text_not_found() { let mut state = RunProjection::default(); - state.set_stage( - fabro_store::StageId::new("implement", 1), - fabro_store::StageState::default(), - ); + state.stage_entry("implement", 1, first_event_seq(1)); let result = read_plan_text(&state); assert_eq!(result, None); diff --git a/lib/crates/fabro-workflow/tests/it/integration.rs b/lib/crates/fabro-workflow/tests/it/integration.rs index 2325be740..29f264dc2 100644 --- a/lib/crates/fabro-workflow/tests/it/integration.rs +++ b/lib/crates/fabro-workflow/tests/it/integration.rs @@ -441,15 +441,18 @@ async fn end_to_end_linear_pipeline() { .contains(&"codergen_step".to_string()) ); - let stage_state = state + let node_state = state .stage(&fabro_types::StageId::new("codergen_step", 1)) .unwrap(); assert!( - stage_state.response.is_some(), + node_state.response.is_some(), "response should be projected" ); - assert!(stage_state.status.is_some(), "status should be projected"); - let prompt_content = stage_state.prompt.as_deref().unwrap(); + assert!( + node_state.completion.is_some(), + "completion should be projected" + ); + let prompt_content = node_state.prompt.as_deref().unwrap(); assert!( prompt_content.ends_with("Implement the feature"), "prompt should end with original prompt, got: {prompt_content}" @@ -9436,14 +9439,14 @@ async fn node_dir_uses_visit_count_on_revisit() { .stage(&fabro_types::StageId::new("gated_work", 2)) .unwrap(); assert_eq!( - first.status.as_ref().unwrap().status, + first.completion.as_ref().unwrap().outcome, StageOutcome::Failed { retry_requested: false, }, "first visit should fail" ); assert_eq!( - second.status.as_ref().unwrap().status, + second.completion.as_ref().unwrap().outcome, StageOutcome::Succeeded, "second visit should succeed" ); @@ -13003,10 +13006,8 @@ async fn asset_collection_local_sandbox_success() { "expected stored artifacts for both files" ); assert_eq!(artifacts[0].node, StageId::new("create_assets", 1)); - assert_eq!(artifacts[0].retry, 1); assert_eq!(artifacts[0].filename, "test-results/output.txt"); assert_eq!(artifacts[1].node, StageId::new("create_assets", 1)); - assert_eq!(artifacts[1].retry, 1); assert_eq!(artifacts[1].filename, "test-results/report.xml"); let report_content = String::from_utf8( artifact_store diff --git a/lib/packages/fabro-api-client/src/.openapi-generator/FILES b/lib/packages/fabro-api-client/src/.openapi-generator/FILES index 3afbe10c3..5120a1a24 100644 --- a/lib/packages/fabro-api-client/src/.openapi-generator/FILES +++ b/lib/packages/fabro-api-client/src/.openapi-generator/FILES @@ -154,7 +154,6 @@ models/model-reference.ts models/model-test-mode.ts models/model-test-result.ts models/model.ts -models/node-status-record.ts models/notification-provider-settings.ts models/notification-route-settings.ts models/object-store-local-settings.ts @@ -286,9 +285,10 @@ models/server-web-settings.ts models/slack-integration-settings.ts models/ssh-access-request.ts models/ssh-access-response.ts +models/stage-completion.ts models/stage-outcome.ts +models/stage-projection.ts models/stage-state.ts -models/stage-status.ts models/stage-turn.ts models/start-run-request.ts models/submit-answer-request.ts diff --git a/lib/packages/fabro-api-client/src/models/index.ts b/lib/packages/fabro-api-client/src/models/index.ts index df9d1f380..986f961a9 100644 --- a/lib/packages/fabro-api-client/src/models/index.ts +++ b/lib/packages/fabro-api-client/src/models/index.ts @@ -133,7 +133,6 @@ export * from './model-limits'; export * from './model-reference'; export * from './model-test-mode'; export * from './model-test-result'; -export * from './node-status-record'; export * from './notification-provider-settings'; export * from './notification-route-settings'; export * from './object-store-local-settings'; @@ -265,9 +264,10 @@ export * from './server-web-settings'; export * from './slack-integration-settings'; export * from './ssh-access-request'; export * from './ssh-access-response'; +export * from './stage-completion'; export * from './stage-outcome'; +export * from './stage-projection'; export * from './stage-state'; -export * from './stage-status'; export * from './stage-turn'; export * from './start-run-request'; export * from './submit-answer-request'; diff --git a/lib/packages/fabro-api-client/src/models/run-projection.ts b/lib/packages/fabro-api-client/src/models/run-projection.ts index 98e758d0b..d3797bddc 100644 --- a/lib/packages/fabro-api-client/src/models/run-projection.ts +++ b/lib/packages/fabro-api-client/src/models/run-projection.ts @@ -33,7 +33,7 @@ import type { RunSpec } from './run-spec'; import type { RunStatus } from './run-status'; // May contain unused imports in some cases // @ts-ignore -import type { StageState } from './stage-state'; +import type { StageProjection } from './stage-projection'; /** * Raw internal run projection derived from the event log. @@ -60,9 +60,9 @@ export interface RunProjection { 'superseded_by'?: string | null; 'pending_interviews'?: { [key: string]: PendingInterviewRecord; }; /** - * Map from StageId (`node_id@visit`) to StageState. + * Map from StageId (`node_id@visit`) to stage projection data. */ - 'stages': { [key: string]: StageState; }; + 'stages': { [key: string]: StageProjection; }; } diff --git a/lib/packages/fabro-api-client/src/models/run-stage.ts b/lib/packages/fabro-api-client/src/models/run-stage.ts index 575b12c96..c98ec9bee 100644 --- a/lib/packages/fabro-api-client/src/models/run-stage.ts +++ b/lib/packages/fabro-api-client/src/models/run-stage.ts @@ -15,7 +15,7 @@ // May contain unused imports in some cases // @ts-ignore -import type { StageStatus } from './stage-status'; +import type { StageState } from './stage-state'; /** * A single stage in a run\'s workflow graph. @@ -29,7 +29,7 @@ export interface RunStage { * Human-readable stage name. */ 'name': string; - 'status': StageStatus; + 'status': StageState; /** * Time spent in this stage, in seconds. */ diff --git a/lib/packages/fabro-api-client/src/models/node-status-record.ts b/lib/packages/fabro-api-client/src/models/stage-completion.ts similarity index 81% rename from lib/packages/fabro-api-client/src/models/node-status-record.ts rename to lib/packages/fabro-api-client/src/models/stage-completion.ts index 6bfb568ac..b0d45cd46 100644 --- a/lib/packages/fabro-api-client/src/models/node-status-record.ts +++ b/lib/packages/fabro-api-client/src/models/stage-completion.ts @@ -18,10 +18,10 @@ import type { StageOutcome } from './stage-outcome'; /** - * Internal node status record. + * Terminal completion metadata for a projected workflow stage. */ -export interface NodeStatusRecord { - 'status': StageOutcome; +export interface StageCompletion { + 'outcome': StageOutcome; 'notes'?: string | null; 'failure_reason'?: string | null; 'timestamp': string; diff --git a/lib/packages/fabro-api-client/src/models/stage-projection.ts b/lib/packages/fabro-api-client/src/models/stage-projection.ts new file mode 100644 index 000000000..41321ec78 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/stage-projection.ts @@ -0,0 +1,58 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.1.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + +// May contain unused imports in some cases +// @ts-ignore +import type { CommandTermination } from './command-termination'; +// May contain unused imports in some cases +// @ts-ignore +import type { StageCompletion } from './stage-completion'; + +/** + * Observable projection data for one workflow stage execution. + */ +export interface StageProjection { + 'first_event_seq': number; + 'prompt'?: string | null; + 'response'?: string | null; + 'completion'?: StageCompletion | null; + /** + * Provider and model metadata recorded for the stage attempt. + */ + 'provider_used'?: object | null; + 'diff'?: string | null; + /** + * Command and environment recorded when the stage script ran. + */ + 'script_invocation'?: object | null; + /** + * Wall-clock and step timing metadata for the stage script. + */ + 'script_timing'?: object | null; + /** + * Per-branch result objects produced by a parallel stage. + */ + 'parallel_results'?: Array | null; + 'stdout'?: string | null; + 'stderr'?: string | null; + 'stdout_bytes'?: number | null; + 'stderr_bytes'?: number | null; + 'streams_separated'?: boolean | null; + 'live_streaming'?: boolean | null; + 'termination'?: CommandTermination | null; +} + + + diff --git a/lib/packages/fabro-api-client/src/models/stage-state.ts b/lib/packages/fabro-api-client/src/models/stage-state.ts index 7838a276c..47c1b494a 100644 --- a/lib/packages/fabro-api-client/src/models/stage-state.ts +++ b/lib/packages/fabro-api-client/src/models/stage-state.ts @@ -13,49 +13,23 @@ */ -// May contain unused imports in some cases -// @ts-ignore -import type { CommandTermination } from './command-termination'; -// May contain unused imports in some cases -// @ts-ignore -import type { NodeStatusRecord } from './node-status-record'; /** - * Internal stage projection state. + * Lifecycle projection state of a workflow stage. */ -export interface StageState { - /** - * Event-log sequence number of the first event that created this stage. Used to order stages by execution. - */ - 'seq'?: number; - 'prompt'?: string | null; - 'response'?: string | null; - 'status'?: NodeStatusRecord | null; - /** - * Provider and model metadata recorded for the stage attempt. - */ - 'provider_used'?: object | null; - 'diff'?: string | null; - /** - * Command and environment recorded when the stage script ran. - */ - 'script_invocation'?: object | null; - /** - * Wall-clock and step timing metadata for the stage script. - */ - 'script_timing'?: object | null; - /** - * Per-branch result objects produced by a parallel stage. - */ - 'parallel_results'?: Array | null; - 'stdout'?: string | null; - 'stderr'?: string | null; - 'stdout_bytes'?: number | null; - 'stderr_bytes'?: number | null; - 'streams_separated'?: boolean | null; - 'live_streaming'?: boolean | null; - 'termination'?: CommandTermination | null; -} + +export const StageState = { + PENDING: 'pending', + RUNNING: 'running', + RETRYING: 'retrying', + SUCCEEDED: 'succeeded', + PARTIALLY_SUCCEEDED: 'partially_succeeded', + FAILED: 'failed', + SKIPPED: 'skipped', + CANCELLED: 'cancelled' +} as const; + +export type StageState = typeof StageState[keyof typeof StageState]; diff --git a/lib/packages/fabro-api-client/src/models/stage-status.ts b/lib/packages/fabro-api-client/src/models/stage-status.ts deleted file mode 100644 index 0756f0150..000000000 --- a/lib/packages/fabro-api-client/src/models/stage-status.ts +++ /dev/null @@ -1,35 +0,0 @@ -/* tslint:disable */ -/* eslint-disable */ -/** - * Fabro Run API - * HTTP API for managing Fabro workflow run executions. - * - * The version of the OpenAPI document: 0.1.0 - * - * - * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). - * https://openapi-generator.tech - * Do not edit the class manually. - */ - - - -/** - * Lifecycle projection state of a workflow stage. - */ - -export const StageStatus = { - PENDING: 'pending', - RUNNING: 'running', - RETRYING: 'retrying', - SUCCEEDED: 'succeeded', - PARTIALLY_SUCCEEDED: 'partially_succeeded', - FAILED: 'failed', - SKIPPED: 'skipped', - CANCELLED: 'cancelled' -} as const; - -export type StageStatus = typeof StageStatus[keyof typeof StageStatus]; - - -