From cea1fa739de8c2247010eb21b5917db9639e6f16 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 1 May 2026 19:56:22 -0400 Subject: [PATCH 1/3] refactor(run-projection): use stage vocabulary --- ...an-events-as-source-of-truth-follow-ups.md | 4 +- docs/internal/run-directory-keys.md | 2 +- docs/public/agents/outputs.mdx | 4 +- docs/public/agents/prompts.mdx | 2 +- docs/public/api-reference/fabro-api.yaml | 30 ++- docs/public/execution/checkpoints.mdx | 2 +- docs/public/execution/retros.mdx | 2 +- docs/public/reference/run-directory.mdx | 2 +- lib/crates/fabro-api/build.rs | 4 +- lib/crates/fabro-api/src/lib.rs | 6 +- .../tests/run_projection_round_trip.rs | 7 +- ...trip.rs => stage_completion_round_trip.rs} | 14 +- ...trip.rs => stage_projection_round_trip.rs} | 17 +- .../fabro-cli/src/commands/run/attach.rs | 2 +- lib/crates/fabro-cli/src/commands/run/diff.rs | 2 +- lib/crates/fabro-cli/tests/it/cmd/dump.rs | 12 +- lib/crates/fabro-cli/tests/it/cmd/inspect.rs | 2 +- lib/crates/fabro-cli/tests/it/cmd/run.rs | 2 +- lib/crates/fabro-cli/tests/it/cmd/runner.rs | 4 +- .../fabro-cli/tests/it/scenario/smoke.rs | 2 +- lib/crates/fabro-dump/src/lib.rs | 137 ++++++---- lib/crates/fabro-retro/src/retro_agent.rs | 111 ++++---- lib/crates/fabro-server/src/server.rs | 8 +- lib/crates/fabro-store/src/lib.rs | 3 +- lib/crates/fabro-store/src/run_state.rs | 254 ++++++++++++++---- .../src/serializable_projection.rs | 17 +- .../tests/serializable_projection.rs | 60 +++-- lib/crates/fabro-types/src/lib.rs | 6 +- lib/crates/fabro-types/src/run_projection.rs | 67 +++-- .../{node_status.rs => stage_completion.rs} | 4 +- lib/crates/fabro-workflow/src/git.rs | 14 +- .../fabro-workflow/src/handler/agent.rs | 8 +- .../fabro-workflow/src/handler/command.rs | 16 +- .../fabro-workflow/src/handler/parallel.rs | 4 +- .../fabro-workflow/src/handler/prompt.rs | 2 +- .../fabro-workflow/src/operations/fork.rs | 2 +- .../src/pipeline/execute/tests.rs | 10 +- .../fabro-workflow/src/pipeline/finalize.rs | 191 +++++++++++-- .../src/pipeline/pull_request.rs | 46 +--- .../fabro-workflow/tests/it/integration.rs | 45 ++-- .../src/.openapi-generator/FILES | 4 +- .../fabro-api-client/src/models/index.ts | 4 +- .../fabro-api-client/src/models/model.ts | 2 +- .../src/models/run-projection.ts | 10 +- ...e-status-record.ts => stage-completion.ts} | 6 +- .../{node-state.ts => stage-projection.ts} | 9 +- 46 files changed, 755 insertions(+), 407 deletions(-) rename lib/crates/fabro-api/tests/{node_status_record_round_trip.rs => stage_completion_round_trip.rs} (56%) rename lib/crates/fabro-api/tests/{node_state_round_trip.rs => stage_projection_round_trip.rs} (68%) rename lib/crates/fabro-types/src/{node_status.rs => stage_completion.rs} (82%) rename lib/packages/fabro-api-client/src/models/{node-status-record.ts => stage-completion.ts} (81%) rename lib/packages/fabro-api-client/src/models/{node-state.ts => stage-projection.ts} (81%) 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 bd1858b5a..37f975acf 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -5074,14 +5074,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"] @@ -5091,17 +5091,23 @@ components: type: string format: date-time - NodeState: - description: Internal node projection state. + StageProjection: + description: Observable projection data for one workflow stage execution. type: object + required: + - first_event_seq properties: + first_event_seq: + type: integer + 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: {} diff: @@ -5247,7 +5253,7 @@ components: description: Raw internal run projection derived from the event log. type: object required: - - nodes + - stages properties: spec: oneOf: @@ -5310,11 +5316,11 @@ components: type: object additionalProperties: $ref: "#/components/schemas/PendingInterviewRecord" - nodes: + stages: type: object - description: Map from StageId (`node_id@visit`) to NodeState. + description: Map from StageId (`node_id@visit`) to stage projection data. additionalProperties: - $ref: "#/components/schemas/NodeState" + $ref: "#/components/schemas/StageProjection" RunSummary: description: Durable run summary derived from the backing store. 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 7f8a4c41e..f95e77e20 100644 --- a/lib/crates/fabro-api/build.rs +++ b/lib/crates/fabro-api/build.rs @@ -312,7 +312,7 @@ 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", &[]), ("StageState", "fabro_types::StageState", &[]), ( @@ -321,7 +321,7 @@ fn main() { &[], ), ("CommandTermination", "fabro_types::CommandTermination", &[]), - ("NodeState", "fabro_types::NodeState", &[]), + ("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 78e3f7ab1..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, - NodeState, NodeStatusRecord, PendingInterviewRecord, PreRunPushOutcome, QuestionType, - RepositoryReference, RunEvent, RunProjection, RunSummary, SecretMetadata, SecretType, - ServerSettings, StageOutcome, StageState, 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 20c6bac2d..2ec82be1c 100644 --- a/lib/crates/fabro-api/tests/run_projection_round_trip.rs +++ b/lib/crates/fabro-api/tests/run_projection_round_trip.rs @@ -56,11 +56,12 @@ fn run_projection_round_trips_populated_projection() { "started_at": "2026-04-29T12:35:00Z" } }, - "nodes": { + "stages": { "build@2": { + "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, @@ -96,7 +97,7 @@ fn run_projection_round_trips_with_pending_control_unset() { "pull_request": null, "superseded_by": null, "pending_interviews": {}, - "nodes": {} + "stages": {} }); let projection: RunProjection = serde_json::from_value(value.clone()).unwrap(); 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/node_state_round_trip.rs b/lib/crates/fabro-api/tests/stage_projection_round_trip.rs similarity index 68% rename from lib/crates/fabro-api/tests/node_state_round_trip.rs rename to lib/crates/fabro-api/tests/stage_projection_round_trip.rs index 7bc1ea598..0523197cf 100644 --- a/lib/crates/fabro-api/tests/node_state_round_trip.rs +++ b/lib/crates/fabro-api/tests/stage_projection_round_trip.rs @@ -1,21 +1,22 @@ use std::any::{TypeId, type_name}; -use fabro_api::types::NodeState as ApiNodeState; -use fabro_types::NodeState; +use fabro_api::types::StageProjection as ApiStageProjection; +use fabro_types::StageProjection; use serde_json::json; #[test] -fn node_state_reuses_canonical_type() { - assert_same_type::(); +fn stage_projection_reuses_canonical_type() { + assert_same_type::(); } #[test] -fn node_state_round_trips_representative_json() { +fn stage_projection_round_trips_representative_json() { let value = json!({ + "first_event_seq": 1, "prompt": "build it", "response": "done", - "status": { - "status": "succeeded", + "completion": { + "outcome": "succeeded", "notes": null, "failure_reason": null, "timestamp": "2026-04-29T12:34:56Z" @@ -30,7 +31,7 @@ fn node_state_round_trips_representative_json() { "termination": "exited" }); - let state: NodeState = serde_json::from_value(value.clone()).unwrap(); + let state: StageProjection = serde_json::from_value(value.clone()).unwrap(); assert_eq!(serde_json::to_value(state).unwrap(), value); } diff --git a/lib/crates/fabro-cli/src/commands/run/attach.rs b/lib/crates/fabro-cli/src/commands/run/attach.rs index 4e71d6a28..58e377bf6 100644 --- a/lib/crates/fabro-cli/src/commands/run/attach.rs +++ b/lib/crates/fabro-cli/src/commands/run/attach.rs @@ -506,7 +506,7 @@ mod tests { "sandbox": null, "final_patch": null, "pull_request": null, - "nodes": {} + "stages": {} }) } diff --git a/lib/crates/fabro-cli/src/commands/run/diff.rs b/lib/crates/fabro-cli/src/commands/run/diff.rs index e4ee319fb..1c9394a71 100644 --- a/lib/crates/fabro-cli/src/commands/run/diff.rs +++ b/lib/crates/fabro-cli/src/commands/run/diff.rs @@ -51,7 +51,7 @@ pub(crate) async fn run(args: DiffArgs, base_ctx: &CommandContext) -> Result<()> fn resolve_diff(state: &RunProjection, args: &DiffArgs) -> Result { if let Some(ref node_id) = args.node { if let Some(visit) = state.list_node_visits(node_id).into_iter().max() { - if let Some(node) = state.node(&fabro_store::StageId::new(node_id, visit)) { + if let Some(node) = state.stage(&fabro_store::StageId::new(node_id, visit)) { if let Some(patch) = node.diff.clone() { debug!(node_id, visit, "Reading per-node diff from projected state"); return Ok(patch); diff --git a/lib/crates/fabro-cli/tests/it/cmd/dump.rs b/lib/crates/fabro-cli/tests/it/cmd/dump.rs index 11b92c8ed..9b60f7453 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/dump.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/dump.rs @@ -282,12 +282,12 @@ fn dump_exports_completed_run_snapshot() { graph.fabro run.json run.log - stages/exit@1/status.json - stages/report@1/response.md - stages/report@1/status.json - stages/run_tests@1/response.md - stages/run_tests@1/status.json - stages/start@1/status.json + stages/001-start@1/status.json + stages/002-run_tests@1/response.md + stages/002-run_tests@1/status.json + stages/003-report@1/response.md + stages/003-report@1/status.json + stages/004-exit@1/status.json "); } diff --git a/lib/crates/fabro-cli/tests/it/cmd/inspect.rs b/lib/crates/fabro-cli/tests/it/cmd/inspect.rs index c9dad593b..26f6be124 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/inspect.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/inspect.rs @@ -83,7 +83,7 @@ fn inspect_resolves_selector_via_server_endpoint() { .path(format!("/api/v1/runs/{}/state", run_id.as_str())); then.status(200) .header("content-type", "application/json") - .body(r#"{"nodes": {}}"#); + .body(r#"{"stages": {}}"#); }); let mut cmd = context.command(); diff --git a/lib/crates/fabro-cli/tests/it/cmd/run.rs b/lib/crates/fabro-cli/tests/it/cmd/run.rs index 43bff5713..0b4df300e 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/run.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/run.rs @@ -63,7 +63,7 @@ fn remote_run_state_response() -> serde_json::Value { "sandbox": null, "final_patch": null, "pull_request": null, - "nodes": {} + "stages": {} }) } diff --git a/lib/crates/fabro-cli/tests/it/cmd/runner.rs b/lib/crates/fabro-cli/tests/it/cmd/runner.rs index d93c7c1cf..fafa1681c 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/runner.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/runner.rs @@ -553,7 +553,7 @@ methods = ["dev-token"] let state = run_state(&run_dir); let probe_stage_id = StageId::new("probe", 1); let _probe = state - .node(&probe_stage_id) + .stage(&probe_stage_id) .expect("probe node state should exist"); let stdout = command_log_text(&run_dir, &probe_stage_id, CommandOutputStream::Stdout); assert!( @@ -723,7 +723,7 @@ fn runner_reports_missing_run_spec_without_prefetching_events() { "sandbox": null, "final_patch": null, "pull_request": null, - "nodes": {} + "stages": {} }) .to_string(), ); diff --git a/lib/crates/fabro-cli/tests/it/scenario/smoke.rs b/lib/crates/fabro-cli/tests/it/scenario/smoke.rs index 9cb223a33..6820875fe 100644 --- a/lib/crates/fabro-cli/tests/it/scenario/smoke.rs +++ b/lib/crates/fabro-cli/tests/it/scenario/smoke.rs @@ -21,7 +21,7 @@ fn live_run_state_response() -> serde_json::Value { "sandbox": null, "final_patch": null, "pull_request": null, - "nodes": {} + "stages": {} }) } diff --git a/lib/crates/fabro-dump/src/lib.rs b/lib/crates/fabro-dump/src/lib.rs index d12df2210..262669469 100644 --- a/lib/crates/fabro-dump/src/lib.rs +++ b/lib/crates/fabro-dump/src/lib.rs @@ -48,69 +48,75 @@ impl RunDump { } let mut stage_ids: Vec<_> = state - .iter_nodes() - .map(|(stage_id, _)| stage_id.clone()) + .iter_stages() + .map(|(stage_id, stage)| (stage_id.clone(), stage.first_event_seq)) .collect(); - stage_ids.sort(); + stage_ids.sort_by(|(left_id, left_seq), (right_id, right_seq)| { + left_seq + .get() + .cmp(&right_seq.get()) + .then_with(|| left_id.cmp(right_id)) + }); - for stage_id in stage_ids { - let Some(node) = state.node(&stage_id) else { + for (index, (stage_id, _)) in stage_ids.into_iter().enumerate() { + let Some(stage) = state.stage(&stage_id) else { continue; }; - let base = PathBuf::from("stages").join(stage_id.to_string()); + let rank = index + 1; + let base = PathBuf::from("stages").join(format!("{rank:03}-{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(), @@ -425,19 +431,24 @@ fn ensure_parent_dir(path: &Path) -> Result<()> { #[cfg(test)] mod tests { use std::collections::HashMap; + use std::num::NonZeroU32; use chrono::{TimeZone, Utc}; - use fabro_store::{NodeState, RunProjection, StageId}; + use fabro_store::{RunProjection, StageId}; use fabro_types::graph::Graph; use fabro_types::run::RunSpec; use fabro_types::{ - Checkpoint, Conclusion, NodeStatusRecord, RunStatus, SandboxRecord, StageOutcome, + Checkpoint, Conclusion, RunStatus, SandboxRecord, StageCompletion, StageOutcome, StartRecord, SuccessReason, WorkflowSettings, fixtures, }; use futures::executor; use super::{RunDump, RunDumpContents, RunDumpEntry}; + fn nonzero(value: u32) -> NonZeroU32 { + NonZeroU32::new(value).expect("test sequence must be non-zero") + } + fn sample_run_spec() -> RunSpec { RunSpec { run_id: fixtures::RUN_1, @@ -522,31 +533,25 @@ mod tests { }); projection.retro_prompt = Some("retro prompt".to_string()); projection.retro_response = Some("retro response".to_string()); - projection.set_node(stage_id.clone(), NodeState { - 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(), nonzero(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 @@ -559,16 +564,16 @@ mod tests { assert!(paths.contains(&"graph.fabro")); assert!(paths.contains(&"stages/retro/prompt.md")); assert!(paths.contains(&"stages/retro/response.md")); - assert!(paths.contains(&"stages/build@2/prompt.md")); - assert!(paths.contains(&"stages/build@2/response.md")); - assert!(paths.contains(&"stages/build@2/status.json")); - assert!(paths.contains(&"stages/build@2/provider_used.json")); - assert!(paths.contains(&"stages/build@2/diff.patch")); - assert!(paths.contains(&"stages/build@2/script_invocation.json")); - assert!(paths.contains(&"stages/build@2/script_timing.json")); - assert!(paths.contains(&"stages/build@2/parallel_results.json")); - assert!(paths.contains(&"stages/build@2/stdout.log")); - assert!(paths.contains(&"stages/build@2/stderr.log")); + assert!(paths.contains(&"stages/001-build@2/prompt.md")); + assert!(paths.contains(&"stages/001-build@2/response.md")); + assert!(paths.contains(&"stages/001-build@2/status.json")); + assert!(paths.contains(&"stages/001-build@2/provider_used.json")); + assert!(paths.contains(&"stages/001-build@2/diff.patch")); + assert!(paths.contains(&"stages/001-build@2/script_invocation.json")); + assert!(paths.contains(&"stages/001-build@2/script_timing.json")); + assert!(paths.contains(&"stages/001-build@2/parallel_results.json")); + assert!(paths.contains(&"stages/001-build@2/stdout.log")); + assert!(paths.contains(&"stages/001-build@2/stderr.log")); assert!(!paths.contains(&"start.json")); assert!(!paths.contains(&"status.json")); assert!(!paths.contains(&"checkpoint.json")); @@ -585,7 +590,7 @@ mod tests { panic!("run.json should be json"); }; let round_tripped: RunProjection = serde_json::from_value(value.clone()).unwrap(); - let node = round_tripped.node(&stage_id).expect("node should exist"); + let node = round_tripped.stage(&stage_id).expect("node should exist"); assert!(round_tripped.spec.is_some()); assert!(round_tripped.start.is_some()); @@ -604,6 +609,30 @@ mod tests { ); } + #[test] + fn from_projection_prefixes_stage_paths_but_not_artifact_paths() { + let mut projection = RunProjection::default(); + projection.stage_entry("zebra", 1, nonzero(1)).prompt = Some("first".to_string()); + projection.stage_entry("apple", 1, nonzero(2)).prompt = Some("second".to_string()); + + let mut dump = RunDump::from_projection(&projection).unwrap(); + dump.add_artifact_bytes(&StageId::new("zebra", 1), "report.txt", b"z".to_vec()) + .unwrap(); + dump.add_artifact_bytes(&StageId::new("apple", 1), "report.txt", b"a".to_vec()) + .unwrap(); + + let paths: Vec<&str> = dump + .entries() + .iter() + .map(|entry| entry.path.as_str()) + .collect(); + + assert!(paths.contains(&"stages/001-zebra@1/prompt.md")); + assert!(paths.contains(&"stages/002-apple@1/prompt.md")); + assert!(paths.contains(&"artifacts/zebra@1/report.txt")); + assert!(paths.contains(&"artifacts/apple@1/report.txt")); + } + #[test] fn hydrate_referenced_blobs_ignores_legacy_artifact_file_refs() { let blob = serde_json::to_vec("hydrated legacy text").unwrap(); diff --git a/lib/crates/fabro-retro/src/retro_agent.rs b/lib/crates/fabro-retro/src/retro_agent.rs index bc72467f4..fb412f371 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 @@ -330,16 +330,21 @@ async fn upload_data_files( #[cfg(test)] mod tests { + use std::num::NonZeroU32; use std::sync::Arc; use chrono::{TimeZone, Utc}; use fabro_agent::LocalSandbox; - use fabro_store::{NodeState, StageId}; - use fabro_types::{NodeStatusRecord, StageOutcome}; + use fabro_store::StageId; + use fabro_types::{StageCompletion, StageOutcome}; use tokio::fs; use super::*; + fn nonzero(value: u32) -> NonZeroU32 { + NonZeroU32::new(value).expect("test sequence must be non-zero") + } + #[test] fn submit_retro_schema_is_valid_json() { let schema: serde_json::Value = serde_json::from_str(SUBMIT_RETRO_SCHEMA).unwrap(); @@ -410,31 +415,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_node(stage_id, NodeState { - 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(), nonzero(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, @@ -455,8 +454,8 @@ mod tests { .expect("run.json should parse"); assert!(run_json.get("spec").is_some()); assert!(run_json.get("run").is_none()); - assert!(run_json["nodes"]["build@2"]["prompt"].is_null()); - assert!(run_json["nodes"]["build@2"]["diff"].is_null()); + assert!(run_json["stages"]["build@2"]["prompt"].is_null()); + assert!(run_json["stages"]["build@2"]["diff"].is_null()); assert_eq!( fs::read_to_string(target_dir.join("graph.fabro")) .await @@ -464,19 +463,19 @@ mod tests { "digraph Ship {}" ); assert_eq!( - fs::read_to_string(target_dir.join("stages/build@2/prompt.md")) + fs::read_to_string(target_dir.join("stages/001-build@2/prompt.md")) .await .expect("prompt file should exist"), "plan" ); assert_eq!( - fs::read_to_string(target_dir.join("stages/build@2/response.md")) + fs::read_to_string(target_dir.join("stages/001-build@2/response.md")) .await .expect("response file should exist"), "done" ); assert_eq!( - fs::read_to_string(target_dir.join("stages/build@2/stdout.log")) + fs::read_to_string(target_dir.join("stages/001-build@2/stdout.log")) .await .expect("stdout file should exist"), "stdout" @@ -494,7 +493,7 @@ mod tests { "server log\n" ); assert!( - target_dir.join("stages/build@2/status.json").exists(), + target_dir.join("stages/001-build@2/status.json").exists(), "status file should exist" ); assert!( @@ -521,21 +520,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_node(stage_id, NodeState { - 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), - ..NodeState::default() - }); + let stage = state.stage_entry(stage_id.node_id(), stage_id.visit(), nonzero(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(); @@ -556,20 +553,20 @@ mod tests { .expect("retro files should upload"); assert_eq!( - fs::read_to_string(target_dir.join("stages/build@1/stdout.log")) + fs::read_to_string(target_dir.join("stages/001-build@1/stdout.log")) .await .expect("stdout file should exist"), "resolved stdout" ); assert_eq!( - fs::read_to_string(target_dir.join("stages/build@1/stderr.log")) + fs::read_to_string(target_dir.join("stages/001-build@1/stderr.log")) .await .expect("stderr file should exist"), "resolved stderr" ); let script_timing: serde_json::Value = serde_json::from_str( - &fs::read_to_string(target_dir.join("stages/build@1/script_timing.json")) + &fs::read_to_string(target_dir.join("stages/001-build@1/script_timing.json")) .await .expect("script timing should exist"), ) @@ -578,7 +575,7 @@ mod tests { assert_eq!(script_timing["stderr"], "resolved stderr"); let script_invocation: serde_json::Value = serde_json::from_str( - &fs::read_to_string(target_dir.join("stages/build@1/script_invocation.json")) + &fs::read_to_string(target_dir.join("stages/001-build@1/script_invocation.json")) .await .expect("script invocation should exist"), ) @@ -593,18 +590,18 @@ mod tests { ) .expect("run.json should parse"); assert_eq!( - run_json["nodes"]["build@1"]["script_timing"]["stdout"], + run_json["stages"]["build@1"]["script_timing"]["stdout"], "resolved stdout" ); assert_eq!( - run_json["nodes"]["build@1"]["script_timing"]["stderr"], + run_json["stages"]["build@1"]["script_timing"]["stderr"], "resolved stderr" ); assert_eq!( - run_json["nodes"]["build@1"]["script_invocation"]["stdout"], + run_json["stages"]["build@1"]["script_invocation"]["stdout"], "resolved stdout" ); - assert!(run_json["nodes"]["build@1"]["stdout"].is_null()); - assert!(run_json["nodes"]["build@1"]["stderr"].is_null()); + assert!(run_json["stages"]["build@1"]["stdout"].is_null()); + assert!(run_json["stages"]["build@1"]["stderr"].is_null()); } } diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 72bdc8cb2..c39f5c200 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -5366,7 +5366,7 @@ async fn get_run_stage_command_log( .into_response(); } }; - let Some(node) = run_state.node(&stage_id) else { + let Some(node) = run_state.stage(&stage_id) else { return ApiError::not_found("Stage not found.").into_response(); }; @@ -5379,7 +5379,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() @@ -5442,7 +5442,7 @@ async fn get_run_stage_command_log( query.offset, limit, LogSource::Full(&[]), - node.status.is_some(), + node.completion.is_some(), None, live_streaming, ) @@ -10890,7 +10890,7 @@ slug = "fabro" let response = app.oneshot(req).await.unwrap(); let body = response_json!(response, StatusCode::OK).await; - assert!(body["nodes"].is_object()); + assert!(body["stages"].is_object()); } #[tokio::test] diff --git a/lib/crates/fabro-store/src/lib.rs b/lib/crates/fabro-store/src/lib.rs index a88f9c77e..5a97b4b9c 100644 --- a/lib/crates/fabro-store/src/lib.rs +++ b/lib/crates/fabro-store/src/lib.rs @@ -13,7 +13,8 @@ mod types; pub use artifact_store::{ArtifactStore, NodeArtifact, stage_storage_segment}; pub use error::{Error, Result}; pub use fabro_types::{ - EventEnvelope, NodeState, PendingInterviewRecord, RunBlobId, RunProjection, RunSummary, StageId, + EventEnvelope, PendingInterviewRecord, RunBlobId, RunProjection, RunSummary, StageId, + 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 ada4202cb..98309385e 100644 --- a/lib/crates/fabro-store/src/run_state.rs +++ b/lib/crates/fabro-store/src/run_state.rs @@ -1,4 +1,5 @@ use std::collections::{BTreeMap, HashMap}; +use std::num::NonZeroU32; use std::str::FromStr; use chrono::{DateTime, Utc}; @@ -8,8 +9,8 @@ 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, + Outcome, PendingInterviewRecord, PullRequestRecord, RunControlAction, RunId, RunProjection, + RunSpec, RunStatus, RunSummary, SandboxRecord, StageCompletion, StageOutcome, StartRecord, TerminalStatus, }; use fabro_util::error::render_with_causes; @@ -200,9 +201,28 @@ impl RunProjectionReducer for RunProjection { .and_then(|visit| u32::try_from(*visit).ok()) .unwrap_or(1); if let Some(diff) = props.diff.clone() { - self.node_mut(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,20 +287,33 @@ impl RunProjectionReducer for RunProjection { EventBody::InterviewInterrupted(props) if !props.question_id.is_empty() => { self.pending_interviews.remove(&props.question_id); } + EventBody::StageStarted(_) => { + let Some(stage_id) = stored.stage_id.as_ref() else { + return Ok(()); + }; + 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 { return Ok(()); }; let visit = props.visit; - self.node_mut(node_id, visit).prompt = Some(props.text.clone()); - self.node_mut(node_id, visit).provider_used = provider_used_from_prompt(props); + let provider_used = provider_used_from_prompt(props); + let stage = self.stage_entry(node_id, visit, first_event_seq(event.seq)); + stage.prompt = Some(props.text.clone()); + stage.provider_used = provider_used; } EventBody::PromptCompleted(props) => { let Some(node_id) = stored.node_id.as_deref() else { return Ok(()); }; let visit = self.current_visit_for(node_id).unwrap_or(1); - self.node_mut(node_id, visit).response = Some(props.response.clone()); + self.stage_entry(node_id, visit, first_event_seq(event.seq)) + .response = Some(props.response.clone()); } EventBody::StageCompleted(props) => { let Some(node_id) = stored.node_id.as_deref() else { @@ -289,10 +322,10 @@ 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 = self.node_mut(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 { @@ -300,9 +333,9 @@ impl RunProjectionReducer for RunProjection { }; let visit = self.current_visit_for(node_id).unwrap_or(1); let failure_reason = props.failure.as_ref().map(|detail| detail.message.clone()); - let node = self.node_mut(node_id, visit); - node.status = Some(NodeStatusRecord { - status: StageOutcome::Failed { + let stage = self.stage_entry(node_id, visit, first_event_seq(event.seq)); + stage.completion = Some(StageCompletion { + outcome: StageOutcome::Failed { retry_requested: false, }, notes: None, @@ -314,52 +347,54 @@ impl RunProjectionReducer for RunProjection { let Some(node_id) = stored.node_id.as_deref() else { return Ok(()); }; - self.node_mut(node_id, props.visit).provider_used = - Some(provider_used_from_agent_session_started(props)); + self.stage_entry(node_id, props.visit, first_event_seq(event.seq)) + .provider_used = Some(provider_used_from_agent_session_started(props)); } EventBody::AgentCliStarted(props) => { let Some(node_id) = stored.node_id.as_deref() else { return Ok(()); }; - self.node_mut(node_id, props.visit).provider_used = - Some(provider_used_from_agent_cli_started(props)); + self.stage_entry(node_id, props.visit, first_event_seq(event.seq)) + .provider_used = Some(provider_used_from_agent_cli_started(props)); } EventBody::CommandStarted(props) => { let Some(node_id) = stored.node_id.as_deref() else { return Ok(()); }; let visit = self.current_visit_for(node_id).unwrap_or(1); - self.node_mut(node_id, visit).script_invocation = - Some(serde_json::to_value(props).map_err(|err| { - Error::InvalidEvent(format!("invalid command.started payload: {err}")) - })?); + self.stage_entry(node_id, visit, first_event_seq(event.seq)) + .script_invocation = Some(serde_json::to_value(props).map_err(|err| { + Error::InvalidEvent(format!("invalid command.started payload: {err}")) + })?); } EventBody::CommandCompleted(props) => { let Some(node_id) = stored.node_id.as_deref() else { return Ok(()); }; let visit = self.current_visit_for(node_id).unwrap_or(1); - let node = self.node_mut(node_id, visit); - 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| { + let script_timing = serde_json::to_value(props).map_err(|err| { Error::InvalidEvent(format!("invalid command.completed payload: {err}")) - })?); + })?; + let stage = self.stage_entry(node_id, visit, first_event_seq(event.seq)); + 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 { return Ok(()); }; let visit = self.current_visit_for(node_id).unwrap_or(1); - self.node_mut(node_id, visit).parallel_results = - Some(serde_json::to_value(&props.results).map_err(|err| { - Error::InvalidEvent(format!("invalid parallel.completed payload: {err}")) - })?); + let parallel_results = serde_json::to_value(&props.results).map_err(|err| { + Error::InvalidEvent(format!("invalid parallel.completed payload: {err}")) + })?; + self.stage_entry(node_id, visit, first_event_seq(event.seq)) + .parallel_results = Some(parallel_results); } _ => {} } @@ -368,6 +403,10 @@ impl RunProjectionReducer for RunProjection { } } +fn first_event_seq(seq: u32) -> NonZeroU32 { + NonZeroU32::new(seq).expect("event seq starts at 1") +} + pub(crate) fn build_summary(state: &RunProjection, run_id: &RunId) -> RunSummary { let workflow_name = state.spec.as_ref().map(|spec| { if spec.graph.name.is_empty() { @@ -512,12 +551,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 @@ -567,23 +606,29 @@ fn provider_used_from_agent_cli_started(props: &AgentCliStartedProps) -> Value { #[cfg(test)] mod tests { - use std::collections::HashMap; + use std::collections::{BTreeMap, HashMap}; + use std::num::NonZeroU32; use chrono::Utc; use fabro_types::run_event::run::RunFailedProps; use fabro_types::run_event::{ - InterviewCompletedProps, InterviewOption, InterviewStartedProps, RunControlEffectProps, + CheckpointCompletedProps, InterviewCompletedProps, InterviewOption, InterviewStartedProps, + RunControlEffectProps, StagePromptProps, StageStartedProps, }; use fabro_types::{ - BlockedReason, Checkpoint, EventBody, FailureReason, NodeState, QuestionType, RunBlobId, - RunControlAction, RunEvent, RunStatus, SuccessReason, TerminalStatus, WorkflowSettings, - fixtures, + BlockedReason, Checkpoint, EventBody, FailureReason, Outcome, QuestionType, RunBlobId, + RunControlAction, RunEvent, RunStatus, StageOutcome, SuccessReason, TerminalStatus, + WorkflowSettings, fixtures, }; use serde_json::json; use super::{RunProjection, RunProjectionReducer, build_summary}; use crate::{Error, EventEnvelope, StageId}; + fn nonzero(value: u32) -> NonZeroU32 { + NonZeroU32::new(value).expect("test sequence must be non-zero") + } + fn test_event(seq: u32, body: EventBody, node_id: Option<&str>) -> EventEnvelope { let event = RunEvent { id: format!("evt-{seq}"), @@ -604,6 +649,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, @@ -646,7 +697,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" })) @@ -658,7 +709,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", @@ -688,8 +739,9 @@ mod tests { "node_visits": { "build": 2 } } ]], - "nodes": { + "stages": { "build@2": { + "first_event_seq": 1, "diff": "diff --git a/file b/file", "stdout": "done" } @@ -698,7 +750,8 @@ mod tests { .unwrap(); let stage_id = StageId::new("build", 2); - let node = state.node(&stage_id).unwrap(); + let node = state.stage(&stage_id).unwrap(); + assert_eq!(node.first_event_seq, nonzero(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)); @@ -706,7 +759,7 @@ mod tests { let round_tripped: RunProjection = serde_json::from_value(serde_json::to_value(&state).unwrap()).unwrap(); let serialized = serde_json::to_value(&state).unwrap(); - let round_tripped_node = round_tripped.node(&stage_id).unwrap(); + let round_tripped_node = round_tripped.stage(&stage_id).unwrap(); assert_eq!(round_tripped_node.stdout.as_deref(), Some("done")); assert_eq!(round_tripped.list_node_visits("build"), vec![2]); assert_eq!( @@ -718,7 +771,7 @@ mod tests { } #[test] - fn set_node_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 { @@ -734,17 +787,14 @@ mod tests { restart_failure_signatures: HashMap::new(), node_visits: HashMap::from([("build".to_string(), 2usize)]), })]; - state.set_node(StageId::new("build", 2), NodeState { - stdout: Some("done".to_string()), - ..NodeState::default() - }); + state.stage_entry("build", 2, nonzero(7)).stdout = Some("done".to_string()); let round_tripped: RunProjection = serde_json::from_value(serde_json::to_value(&state).unwrap()).unwrap(); assert_eq!( round_tripped - .node(&StageId::new("build", 2)) + .stage(&StageId::new("build", 2)) .unwrap() .stdout .as_deref(), @@ -757,6 +807,100 @@ mod tests { ); } + #[test] + fn stage_started_sets_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(); + + let stage = state.stage(&stage_id).unwrap(); + assert_eq!(stage.first_event_seq, nonzero(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, nonzero(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, nonzero(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] fn interview_events_populate_and_clear_pending_interviews() { let mut state = RunProjection::default(); diff --git a/lib/crates/fabro-store/src/serializable_projection.rs b/lib/crates/fabro-store/src/serializable_projection.rs index 073a21790..6bcf0c837 100644 --- a/lib/crates/fabro-store/src/serializable_projection.rs +++ b/lib/crates/fabro-store/src/serializable_projection.rs @@ -11,22 +11,19 @@ impl Serialize for SerializableProjection<'_> { { let mut projection = self.0.clone(); let stage_ids: Vec<_> = projection - .iter_nodes() + .iter_stages() .map(|(stage_id, _)| stage_id.clone()) .collect(); for stage_id in stage_ids { - let Some(node) = projection.node(&stage_id).cloned() else { + let Some(stage) = projection.stage_mut(&stage_id) else { continue; }; - projection.set_node(stage_id, crate::NodeState { - prompt: None, - response: None, - diff: None, - stdout: None, - stderr: None, - ..node - }); + 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 b7feb76d3..1560b5c7c 100644 --- a/lib/crates/fabro-store/tests/serializable_projection.rs +++ b/lib/crates/fabro-store/tests/serializable_projection.rs @@ -1,15 +1,20 @@ use std::collections::{BTreeMap, HashMap}; +use std::num::NonZeroU32; use chrono::{TimeZone, Utc}; -use fabro_store::{NodeState, RunProjection, SerializableProjection, StageId}; +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, + Checkpoint, RunStatus, SandboxRecord, StageCompletion, StageOutcome, StartRecord, TerminalStatus, WorkflowSettings, fixtures, }; use serde_json::json; +fn nonzero(value: u32) -> NonZeroU32 { + NonZeroU32::new(value).expect("test sequence must be non-zero") +} + fn sample_run_spec() -> RunSpec { RunSpec { run_id: fixtures::RUN_1, @@ -77,37 +82,31 @@ fn serializable_projection_round_trips_and_trims_bulky_node_fields() { clone_branch: None, }); projection.pending_interviews = BTreeMap::new(); - projection.set_node(stage_id.clone(), NodeState { - 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(), nonzero(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"); let round_tripped: RunProjection = serde_json::from_value(serialized).expect("serialized projection should deserialize"); - let node = round_tripped.node(&stage_id).expect("node should remain"); + let node = round_tripped.stage(&stage_id).expect("node should remain"); assert_eq!(round_tripped.spec().map(RunSpec::id), Some(fixtures::RUN_1)); assert_eq!( @@ -124,6 +123,13 @@ fn serializable_projection_round_trips_and_trims_bulky_node_fields() { assert_eq!(node.diff, None); assert_eq!(node.stdout, None); assert_eq!(node.stderr, None); + assert_eq!(node.first_event_seq, nonzero(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 27de2185d..6181ff184 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,7 +49,6 @@ 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, StageState, }; @@ -71,10 +70,11 @@ pub use run_event::{ MetadataSnapshotPhase, RunEvent, RunNoticeLevel, }; pub use run_id::{RunId, fixtures}; -pub use run_projection::{NodeState, PendingInterviewRecord, RunProjection}; +pub use run_projection::{PendingInterviewRecord, RunProjection, StageProjection}; 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/run_projection.rs b/lib/crates/fabro-types/src/run_projection.rs index f3d404b73..971068a7e 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,7 +29,7 @@ pub struct RunProjection { pub pull_request: Option, pub superseded_by: Option, pub pending_interviews: BTreeMap, - nodes: HashMap, + stages: HashMap, } #[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)] @@ -37,11 +38,12 @@ pub struct PendingInterviewRecord { pub started_at: Option>, } -#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)] -pub struct NodeState { +#[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, @@ -61,26 +63,50 @@ pub struct NodeState { pub termination: Option, } +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 node(&self, node: &StageId) -> Option<&NodeState> { - self.nodes.get(node) + pub fn stage(&self, stage: &StageId) -> Option<&StageProjection> { + self.stages.get(stage) } - pub fn iter_nodes(&self) -> impl Iterator { - self.nodes.iter() + pub fn iter_stages(&self) -> impl Iterator { + self.stages.iter() } pub fn is_empty(&self) -> bool { - self.nodes.is_empty() + self.stages.is_empty() } - pub fn set_node(&mut self, node: StageId, state: NodeState) { - self.nodes.insert(node, 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 { let mut visits = self - .nodes + .stages .keys() .filter(|node| node.node_id() == node_id) .map(StageId::visit) @@ -110,12 +136,19 @@ impl RunProjection { &self.pending_interviews } - pub fn node_mut(&mut self, node_id: &str, visit: u32) -> &mut NodeState { - self.nodes.entry(StageId::new(node_id, visit)).or_default() + pub fn stage_entry( + &mut self, + node_id: &str, + visit: u32, + first_event_seq: NonZeroU32, + ) -> &mut StageProjection { + self.stages + .entry(StageId::new(node_id, visit)) + .or_insert_with(|| StageProjection::new(first_event_seq)) } pub fn current_visit_for(&self, node_id: &str) -> Option { - self.nodes + self.stages .keys() .filter(|node| node.node_id() == node_id) .map(StageId::visit) 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/git.rs b/lib/crates/fabro-workflow/src/git.rs index e8c5ead70..e6f7d08eb 100644 --- a/lib/crates/fabro-workflow/src/git.rs +++ b/lib/crates/fabro-workflow/src/git.rs @@ -556,13 +556,13 @@ mod tests { .git_entries() .unwrap(); let paths: Vec<&str> = files.iter().map(|(path, _)| path.as_str()).collect(); - assert!(paths.contains(&"stages/work@2/prompt.md")); - assert!(paths.contains(&"stages/work@2/response.md")); - assert!(paths.contains(&"stages/work@2/status.json")); - assert!(paths.contains(&"stages/work@2/provider_used.json")); - assert!(paths.contains(&"stages/work@2/script_invocation.json")); - assert!(paths.contains(&"stages/work@2/script_timing.json")); - assert!(paths.contains(&"stages/work@2/parallel_results.json")); + assert!(paths.contains(&"stages/001-work@2/prompt.md")); + assert!(paths.contains(&"stages/001-work@2/response.md")); + assert!(paths.contains(&"stages/001-work@2/status.json")); + assert!(paths.contains(&"stages/001-work@2/provider_used.json")); + assert!(paths.contains(&"stages/001-work@2/script_invocation.json")); + assert!(paths.contains(&"stages/001-work@2/script_timing.json")); + assert!(paths.contains(&"stages/001-work@2/parallel_results.json")); } #[test] diff --git a/lib/crates/fabro-workflow/src/handler/agent.rs b/lib/crates/fabro-workflow/src/handler/agent.rs index 32d25ddf0..b685fabd6 100644 --- a/lib/crates/fabro-workflow/src/handler/agent.rs +++ b/lib/crates/fabro-workflow/src/handler/agent.rs @@ -505,7 +505,7 @@ mod tests { logger.flush().await; let state = run_store.state().await.unwrap(); - let node_state = state.node(&StageId::new("plan", 1)).unwrap(); + let node_state = state.stage(&StageId::new("plan", 1)).unwrap(); assert_eq!( node_state.prompt.as_deref(), Some("Achieve: Build a feature") @@ -532,7 +532,7 @@ mod tests { logger.flush().await; let state = run_store.state().await.unwrap(); - let node_state = state.node(&StageId::new("work", 1)).unwrap(); + let node_state = state.stage(&StageId::new("work", 1)).unwrap(); assert_eq!(node_state.prompt.as_deref(), Some("Do work")); } @@ -774,7 +774,7 @@ mod tests { logger.flush().await; let state = run_store.state().await.unwrap(); - let node_state = state.node(&StageId::new("step", 1)).unwrap(); + let node_state = state.stage(&StageId::new("step", 1)).unwrap(); assert_eq!( node_state.provider_used.as_ref().unwrap()["provider"], "openai" @@ -1262,7 +1262,7 @@ Some text in between. logger.flush().await; let state = run_store.state().await.unwrap(); - let node_state = state.node(&StageId::new("report", 1)).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"), diff --git a/lib/crates/fabro-workflow/src/handler/command.rs b/lib/crates/fabro-workflow/src/handler/command.rs index b123727b6..6960ba4d9 100644 --- a/lib/crates/fabro-workflow/src/handler/command.rs +++ b/lib/crates/fabro-workflow/src/handler/command.rs @@ -509,7 +509,7 @@ mod tests { logger.flush().await; let snapshot = run_store.state().await.unwrap(); - let node_state = snapshot.node(&StageId::new("script_node", 1)).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"); @@ -540,7 +540,7 @@ mod tests { logger.flush().await; let snapshot = run_store.state().await.unwrap(); - let node_state = snapshot.node(&StageId::new("script_node", 1)).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"); @@ -567,7 +567,7 @@ mod tests { logger.flush().await; let snapshot = run_store.state().await.unwrap(); - let node_state = snapshot.node(&StageId::new("script_node", 1)).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 = node_state.stderr.as_deref().unwrap(); @@ -598,7 +598,7 @@ mod tests { logger.flush().await; let snapshot = run_store.state().await.unwrap(); - let node_state = snapshot.node(&StageId::new("script_node", 1)).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,7 +623,7 @@ mod tests { logger.flush().await; let snapshot = run_store.state().await.unwrap(); - let node_state = snapshot.node(&StageId::new("script_node", 1)).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); @@ -648,7 +648,7 @@ mod tests { logger.flush().await; let snapshot = run_store.state().await.unwrap(); - let node_state = snapshot.node(&StageId::new("script_node", 1)).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,7 +678,7 @@ mod tests { logger.flush().await; let snapshot = run_store.state().await.unwrap(); - let node_state = snapshot.node(&StageId::new("script_node", 1)).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); @@ -706,7 +706,7 @@ mod tests { let snapshot = run_store.state().await.unwrap(); let node = snapshot - .node(&StageId::new("script_node", 1)) + .stage(&StageId::new("script_node", 1)) .cloned() .unwrap(); diff --git a/lib/crates/fabro-workflow/src/handler/parallel.rs b/lib/crates/fabro-workflow/src/handler/parallel.rs index 383cb005f..92168c79c 100644 --- a/lib/crates/fabro-workflow/src/handler/parallel.rs +++ b/lib/crates/fabro-workflow/src/handler/parallel.rs @@ -736,7 +736,7 @@ mod tests { assert!(results.is_some()); let state = run_store.state().await.unwrap(); - let node_state = state.node(&StageId::new("par", 1)).unwrap(); + let node_state = state.stage(&StageId::new("par", 1)).unwrap(); let parsed = node_state.parallel_results.as_ref().unwrap(); assert!( parsed.is_array(), @@ -781,7 +781,7 @@ mod tests { logger.flush().await; let state = run_store.state().await.unwrap(); - let node_state = state.node(&fabro_store::StageId::new("par", 1)).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 3881faf4f..a76099b80 100644 --- a/lib/crates/fabro-workflow/src/handler/prompt.rs +++ b/lib/crates/fabro-workflow/src/handler/prompt.rs @@ -367,7 +367,7 @@ mod tests { logger.flush().await; let state = run_store.state().await.unwrap(); - let node_state = state.node(&StageId::new("classify", 1)).unwrap(); + let node_state = state.stage(&StageId::new("classify", 1)).unwrap(); assert_eq!(node_state.provider_used.as_ref().unwrap()["mode"], "prompt"); } diff --git a/lib/crates/fabro-workflow/src/operations/fork.rs b/lib/crates/fabro-workflow/src/operations/fork.rs index 3c119527f..9c7c48d12 100644 --- a/lib/crates/fabro-workflow/src/operations/fork.rs +++ b/lib/crates/fabro-workflow/src/operations/fork.rs @@ -401,7 +401,7 @@ mod tests { let forked_events = forked.list_events().await.unwrap(); let forked_state = fabro_store::RunProjection::apply_events(&forked_events).unwrap(); let node = forked_state - .node(&StageId::new("work", 1)) + .stage(&StageId::new("work", 1)) .expect("forked state should retain historical node projection"); assert_eq!(node.response.as_deref(), Some("historical response")); diff --git a/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs b/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs index 8377ca3ec..07dc8d31d 100644 --- a/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs +++ b/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs @@ -770,9 +770,9 @@ async fn execute_persists_start_record_and_node_status() { ); assert_eq!(start.base_sha.as_deref(), Some("abc123")); - let node = state.node(&fabro_store::StageId::new("start", 1)).unwrap(); + 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 ); } @@ -825,12 +825,12 @@ async fn timeout_causes_fail_status_record() { .await; let state = executed.engine.run.run_store.state().await.unwrap(); let status = state - .node(&fabro_store::StageId::new("work", 1)) + .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 1e509ddea..f158bf5bc 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; @@ -67,11 +68,14 @@ pub(crate) async fn build_conclusion_from_store( run_duration_ms: u64, final_git_commit_sha: Option, ) -> Conclusion { - let checkpoint = run_store - .state() - .await - .ok() - .and_then(|state| state.checkpoint); + let projection = run_store.state().await.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 = run_store .list_events() .await @@ -79,8 +83,9 @@ pub(crate) async fn build_conclusion_from_store( .unwrap_or_default(); build_conclusion_from_parts( - checkpoint.as_ref(), + checkpoint, &stage_durations, + &projection_order, status, failure_reason, run_duration_ms, @@ -90,7 +95,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, @@ -100,14 +106,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 @@ -117,16 +140,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) @@ -144,6 +182,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( @@ -462,15 +513,18 @@ pub async fn finalize(retroed: Retroed, options: &FinalizeOptions) -> Result Result NonZeroU32 { + NonZeroU32::new(value).expect("test sequence must be non-zero") + } + + fn checkpoint_with( + completed_nodes: Vec<&str>, + 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, nonzero(1)); + projection.stage_entry("apple", 1, nonzero(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, nonzero(4)); + projection.stage_entry("finished", 1, nonzero(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 a9b4cba5c..7067e5039 100644 --- a/lib/crates/fabro-workflow/src/pipeline/pull_request.rs +++ b/lib/crates/fabro-workflow/src/pipeline/pull_request.rs @@ -213,7 +213,7 @@ fn parse_dot_summary(dot: &str) -> (String, usize, usize) { /// directory scan behavior. fn read_plan_text(state: &RunProjection) -> Option { let mut plan_nodes = state - .iter_nodes() + .iter_stages() .filter_map(|(stage_id, node)| { stage_id.node_id().starts_with("plan").then_some(( stage_id.node_id(), @@ -625,6 +625,7 @@ pub async fn pull_request(concluded: Concluded, options: &PullRequestOptions) -> #[cfg(test)] mod tests { use std::collections::HashMap; + use std::num::NonZeroU32; use std::sync::Arc; use std::time::Duration; @@ -653,6 +654,10 @@ mod tests { use crate::event::{Event, append_event}; use crate::records::StageSummary; + fn nonzero(value: u32) -> NonZeroU32 { + NonZeroU32::new(value).expect("test sequence must be non-zero") + } + struct MockProvider { name: String, response_text: String, @@ -997,13 +1002,7 @@ mod tests { #[test] fn read_plan_text_found() { let mut state = RunProjection::default(); - state.set_node( - fabro_store::StageId::new("plan", 1), - fabro_store::NodeState { - response: Some("This is the plan".to_string()), - ..Default::default() - }, - ); + state.stage_entry("plan", 1, nonzero(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 +1011,8 @@ mod tests { #[test] fn read_plan_text_prefix_match() { let mut state = RunProjection::default(); - state.set_node( - fabro_store::StageId::new("planning", 1), - fabro_store::NodeState { - response: Some("Planning content".to_string()), - ..Default::default() - }, - ); + state.stage_entry("planning", 1, nonzero(1)).response = + Some("Planning content".to_string()); let result = read_plan_text(&state); assert_eq!(result, Some("Planning content".to_string())); @@ -1027,20 +1021,9 @@ mod tests { #[test] fn read_plan_text_prefers_alphabetically_first_plan_node() { let mut state = RunProjection::default(); - state.set_node( - fabro_store::StageId::new("planning", 1), - fabro_store::NodeState { - response: Some("Planning content".to_string()), - ..Default::default() - }, - ); - state.set_node( - fabro_store::StageId::new("plan", 1), - fabro_store::NodeState { - response: Some("Plan content".to_string()), - ..Default::default() - }, - ); + state.stage_entry("planning", 1, nonzero(1)).response = + Some("Planning content".to_string()); + state.stage_entry("plan", 1, nonzero(2)).response = Some("Plan content".to_string()); let result = read_plan_text(&state); assert_eq!(result, Some("Plan content".to_string())); @@ -1049,10 +1032,7 @@ mod tests { #[test] fn read_plan_text_not_found() { let mut state = RunProjection::default(); - state.set_node( - fabro_store::StageId::new("implement", 1), - fabro_store::NodeState::default(), - ); + state.stage_entry("implement", 1, nonzero(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 1c522911e..8a9c891c7 100644 --- a/lib/crates/fabro-workflow/tests/it/integration.rs +++ b/lib/crates/fabro-workflow/tests/it/integration.rs @@ -442,13 +442,16 @@ async fn end_to_end_linear_pipeline() { ); let node_state = state - .node(&fabro_types::StageId::new("codergen_step", 1)) + .stage(&fabro_types::StageId::new("codergen_step", 1)) .unwrap(); assert!( node_state.response.is_some(), "response should be projected" ); - assert!(node_state.status.is_some(), "status should be projected"); + 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"), @@ -1882,7 +1885,7 @@ async fn smoke_test_with_mock_codergen_backend() { "should NOT have traversed fix path" ); - let plan_state = state.node(&fabro_types::StageId::new("plan", 1)).unwrap(); + let plan_state = state.stage(&fabro_types::StageId::new("plan", 1)).unwrap(); let plan_response = plan_state .response .as_deref() @@ -3863,7 +3866,7 @@ async fn integration_smoke_plan_implement_review_done() { assert!(cp.completed_nodes.contains(&"implement".to_string())); assert!(cp.completed_nodes.contains(&"review".to_string())); - let plan_state = state.node(&fabro_types::StageId::new("plan", 1)).unwrap(); + let plan_state = state.stage(&fabro_types::StageId::new("plan", 1)).unwrap(); assert!(plan_state.prompt.is_some()); assert!(plan_state.response.is_some()); @@ -6410,7 +6413,7 @@ mod real_llm { // Verify actual LLM responses were written let plan_response = state - .node(&fabro_types::StageId::new("plan", 1)) + .stage(&fabro_types::StageId::new("plan", 1)) .and_then(|node| node.response.as_deref()) .unwrap(); assert!( @@ -6742,7 +6745,7 @@ mod real_llm { assert_eq!(outcome.status, StageOutcome::Succeeded); let response = state - .node(&fabro_types::StageId::new("classify", 1)) + .stage(&fabro_types::StageId::new("classify", 1)) .and_then(|node| node.response.as_deref()) .unwrap(); assert!(!response.is_empty(), "response.md should be non-empty"); @@ -7906,7 +7909,7 @@ async fn hook_stage_start_proceed_allows_execution() { assert!( state - .node(&fabro_types::StageId::new("work", 1)) + .stage(&fabro_types::StageId::new("work", 1)) .and_then(|node| node.response.as_ref()) .is_some(), "response should exist when StageStart hook proceeds" @@ -7933,7 +7936,7 @@ async fn hook_stage_start_skip_bypasses_node() { assert!( state - .node(&fabro_types::StageId::new("work", 1)) + .stage(&fabro_types::StageId::new("work", 1)) .and_then(|node| node.response.as_ref()) .is_none(), "response should not exist when StageStart hook skips node" @@ -7997,7 +8000,7 @@ async fn hook_stage_start_matcher_filters_by_node_id() { assert!( state - .node(&fabro_types::StageId::new("step1", 1)) + .stage(&fabro_types::StageId::new("step1", 1)) .and_then(|node| node.response.as_ref()) .is_some(), "step1 should execute because matcher doesn't match it" @@ -8005,7 +8008,7 @@ async fn hook_stage_start_matcher_filters_by_node_id() { assert!( state - .node(&fabro_types::StageId::new("step2", 1)) + .stage(&fabro_types::StageId::new("step2", 1)) .and_then(|node| node.response.as_ref()) .is_none(), "step2 should be skipped because matcher matches it" @@ -8498,14 +8501,14 @@ async fn hook_matcher_regex_pattern() { assert!( state - .node(&fabro_types::StageId::new("step1", 1)) + .stage(&fabro_types::StageId::new("step1", 1)) .and_then(|node| node.response.as_ref()) .is_none(), "step1 should be skipped by regex ^step" ); assert!( state - .node(&fabro_types::StageId::new("step2", 1)) + .stage(&fabro_types::StageId::new("step2", 1)) .and_then(|node| node.response.as_ref()) .is_none(), "step2 should be skipped by regex ^step" @@ -8715,7 +8718,7 @@ async fn run_fidelity_prompt_pipeline(fidelity: &str) -> String { .expect("pipeline should succeed"); state - .node(&fabro_types::StageId::new("report", 1)) + .stage(&fabro_types::StageId::new("report", 1)) .and_then(|node| node.prompt.clone()) .expect("report prompt should exist") } @@ -9430,20 +9433,20 @@ async fn node_dir_uses_visit_count_on_revisit() { assert_eq!(outcome.status, StageOutcome::Succeeded); let first = state - .node(&fabro_types::StageId::new("gated_work", 1)) + .stage(&fabro_types::StageId::new("gated_work", 1)) .unwrap(); let second = state - .node(&fabro_types::StageId::new("gated_work", 2)) + .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" ); @@ -10335,7 +10338,7 @@ async fn full_pipeline_with_cli_backend_node() { assert_eq!(outcome.status, StageOutcome::Succeeded); let api_response = state - .node(&fabro_types::StageId::new("api_work", 1)) + .stage(&fabro_types::StageId::new("api_work", 1)) .and_then(|node| node.response.as_deref()) .unwrap(); assert!( @@ -10344,7 +10347,7 @@ async fn full_pipeline_with_cli_backend_node() { ); let cli_response = state - .node(&fabro_types::StageId::new("cli_work", 1)) + .stage(&fabro_types::StageId::new("cli_work", 1)) .and_then(|node| node.response.as_deref()) .unwrap(); assert_eq!( @@ -10353,7 +10356,7 @@ async fn full_pipeline_with_cli_backend_node() { ); let provider_json = state - .node(&fabro_types::StageId::new("cli_work", 1)) + .stage(&fabro_types::StageId::new("cli_work", 1)) .unwrap() .provider_used .as_ref() @@ -10454,7 +10457,7 @@ async fn stylesheet_backend_property_routes_to_cli() { assert_eq!(outcome.status, StageOutcome::Succeeded); let response = state - .node(&fabro_types::StageId::new("work", 1)) + .stage(&fabro_types::StageId::new("work", 1)) .and_then(|node| node.response.as_deref()) .unwrap(); assert_eq!( diff --git a/lib/packages/fabro-api-client/src/.openapi-generator/FILES b/lib/packages/fabro-api-client/src/.openapi-generator/FILES index b84e0dcea..5120a1a24 100644 --- a/lib/packages/fabro-api-client/src/.openapi-generator/FILES +++ b/lib/packages/fabro-api-client/src/.openapi-generator/FILES @@ -154,8 +154,6 @@ models/model-reference.ts models/model-test-mode.ts models/model-test-result.ts models/model.ts -models/node-state.ts -models/node-status-record.ts models/notification-provider-settings.ts models/notification-route-settings.ts models/object-store-local-settings.ts @@ -287,7 +285,9 @@ 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-turn.ts models/start-run-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 85bc80c70..986f961a9 100644 --- a/lib/packages/fabro-api-client/src/models/index.ts +++ b/lib/packages/fabro-api-client/src/models/index.ts @@ -133,8 +133,6 @@ export * from './model-limits'; export * from './model-reference'; export * from './model-test-mode'; export * from './model-test-result'; -export * from './node-state'; -export * from './node-status-record'; export * from './notification-provider-settings'; export * from './notification-route-settings'; export * from './object-store-local-settings'; @@ -266,7 +264,9 @@ 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-turn'; export * from './start-run-request'; diff --git a/lib/packages/fabro-api-client/src/models/model.ts b/lib/packages/fabro-api-client/src/models/model.ts index 834661c41..779786680 100644 --- a/lib/packages/fabro-api-client/src/models/model.ts +++ b/lib/packages/fabro-api-client/src/models/model.ts @@ -67,7 +67,7 @@ export interface Model { */ 'default': boolean; /** - * Whether credential material is present for this model\'s provider on the server (vault entry or environment variable). Does NOT imply the credential is valid or that requests will succeed; call `POST /models/{id}/test` to verify usability. + * Whether credential material is present for this model\'s provider on the server (vault entry or environment variable). Does NOT imply the credential is valid or that requests will succeed; call `POST /models/{id}/test` to verify usability. */ 'configured': boolean; } 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 fed006842..d3797bddc 100644 --- a/lib/packages/fabro-api-client/src/models/run-projection.ts +++ b/lib/packages/fabro-api-client/src/models/run-projection.ts @@ -13,9 +13,6 @@ */ -// May contain unused imports in some cases -// @ts-ignore -import type { NodeState } from './node-state'; // May contain unused imports in some cases // @ts-ignore import type { PendingInterviewRecord } from './pending-interview-record'; @@ -34,6 +31,9 @@ import type { RunSpec } from './run-spec'; // May contain unused imports in some cases // @ts-ignore import type { RunStatus } from './run-status'; +// May contain unused imports in some cases +// @ts-ignore +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 NodeState. + * Map from StageId (`node_id@visit`) to stage projection data. */ - 'nodes': { [key: string]: NodeState; }; + 'stages': { [key: string]: StageProjection; }; } 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/node-state.ts b/lib/packages/fabro-api-client/src/models/stage-projection.ts similarity index 81% rename from lib/packages/fabro-api-client/src/models/node-state.ts rename to lib/packages/fabro-api-client/src/models/stage-projection.ts index df5eb16a1..336c9b34b 100644 --- a/lib/packages/fabro-api-client/src/models/node-state.ts +++ b/lib/packages/fabro-api-client/src/models/stage-projection.ts @@ -18,15 +18,16 @@ import type { CommandTermination } from './command-termination'; // May contain unused imports in some cases // @ts-ignore -import type { NodeStatusRecord } from './node-status-record'; +import type { StageCompletion } from './stage-completion'; /** - * Internal node projection state. + * Observable projection data for one workflow stage execution. */ -export interface NodeState { +export interface StageProjection { + 'first_event_seq': number; 'prompt'?: string | null; 'response'?: string | null; - 'status'?: NodeStatusRecord | null; + 'completion'?: StageCompletion | null; 'provider_used'?: any; 'diff'?: string | null; 'script_invocation'?: any; From 840dc42d3c6bf4d34bcd6687dba1fb78083b90cc Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 1 May 2026 20:41:57 -0400 Subject: [PATCH 2/3] refactor(run-projection): dedupe first_event_seq helper and tighten dump Expose `fabro_types::first_event_seq` next to `StageProjection`, replacing six identical `nonzero` test helpers and the private one in `run_state`. Add `RunProjection::iter_stages_mut` so `SerializableProjection` can clear bulky fields without the collect-then-lookup dance, and let the `fabro-dump` loop iterate `(&StageId, &StageProjection)` borrows directly to drop the per-stage `StageId::clone()` and redundant HashMap lookup. Co-Authored-By: Claude Opus 4.7 (1M context) --- lib/crates/fabro-dump/src/lib.rs | 35 ++++++++----------- lib/crates/fabro-retro/src/retro_agent.rs | 11 ++---- lib/crates/fabro-store/src/run_state.rs | 24 ++++--------- .../src/serializable_projection.rs | 10 +----- .../tests/serializable_projection.rs | 11 ++---- lib/crates/fabro-types/src/lib.rs | 2 +- lib/crates/fabro-types/src/run_projection.rs | 11 ++++++ .../fabro-workflow/src/pipeline/finalize.rs | 17 ++++----- .../src/pipeline/pull_request.rs | 25 +++++++------ 9 files changed, 59 insertions(+), 87 deletions(-) diff --git a/lib/crates/fabro-dump/src/lib.rs b/lib/crates/fabro-dump/src/lib.rs index 262669469..57c51cbf7 100644 --- a/lib/crates/fabro-dump/src/lib.rs +++ b/lib/crates/fabro-dump/src/lib.rs @@ -47,21 +47,14 @@ impl RunDump { entries.push(RunDumpEntry::text("graph.fabro", graph_source.clone())); } - let mut stage_ids: Vec<_> = state - .iter_stages() - .map(|(stage_id, stage)| (stage_id.clone(), stage.first_event_seq)) - .collect(); - stage_ids.sort_by(|(left_id, left_seq), (right_id, right_seq)| { - left_seq - .get() - .cmp(&right_seq.get()) + let mut stages: Vec<_> = state.iter_stages().collect(); + 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)) }); - for (index, (stage_id, _)) in stage_ids.into_iter().enumerate() { - let Some(stage) = state.stage(&stage_id) else { - continue; - }; + for (index, (stage_id, stage)) in stages.into_iter().enumerate() { let rank = index + 1; let base = PathBuf::from("stages").join(format!("{rank:03}-{stage_id}")); @@ -431,7 +424,6 @@ fn ensure_parent_dir(path: &Path) -> Result<()> { #[cfg(test)] mod tests { use std::collections::HashMap; - use std::num::NonZeroU32; use chrono::{TimeZone, Utc}; use fabro_store::{RunProjection, StageId}; @@ -439,16 +431,12 @@ mod tests { use fabro_types::run::RunSpec; use fabro_types::{ Checkpoint, Conclusion, RunStatus, SandboxRecord, StageCompletion, StageOutcome, - StartRecord, SuccessReason, WorkflowSettings, fixtures, + StartRecord, SuccessReason, WorkflowSettings, first_event_seq, fixtures, }; use futures::executor; use super::{RunDump, RunDumpContents, RunDumpEntry}; - fn nonzero(value: u32) -> NonZeroU32 { - NonZeroU32::new(value).expect("test sequence must be non-zero") - } - fn sample_run_spec() -> RunSpec { RunSpec { run_id: fixtures::RUN_1, @@ -533,7 +521,8 @@ mod tests { }); projection.retro_prompt = Some("retro prompt".to_string()); projection.retro_response = Some("retro response".to_string()); - let stage = projection.stage_entry(stage_id.node_id(), stage_id.visit(), nonzero(2)); + 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 { @@ -612,8 +601,12 @@ mod tests { #[test] fn from_projection_prefixes_stage_paths_but_not_artifact_paths() { let mut projection = RunProjection::default(); - projection.stage_entry("zebra", 1, nonzero(1)).prompt = Some("first".to_string()); - projection.stage_entry("apple", 1, nonzero(2)).prompt = Some("second".to_string()); + 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 mut dump = RunDump::from_projection(&projection).unwrap(); dump.add_artifact_bytes(&StageId::new("zebra", 1), "report.txt", b"z".to_vec()) diff --git a/lib/crates/fabro-retro/src/retro_agent.rs b/lib/crates/fabro-retro/src/retro_agent.rs index fb412f371..1046f60b8 100644 --- a/lib/crates/fabro-retro/src/retro_agent.rs +++ b/lib/crates/fabro-retro/src/retro_agent.rs @@ -330,21 +330,16 @@ async fn upload_data_files( #[cfg(test)] mod tests { - use std::num::NonZeroU32; use std::sync::Arc; use chrono::{TimeZone, Utc}; use fabro_agent::LocalSandbox; use fabro_store::StageId; - use fabro_types::{StageCompletion, StageOutcome}; + use fabro_types::{StageCompletion, StageOutcome, first_event_seq}; use tokio::fs; use super::*; - fn nonzero(value: u32) -> NonZeroU32 { - NonZeroU32::new(value).expect("test sequence must be non-zero") - } - #[test] fn submit_retro_schema_is_valid_json() { let schema: serde_json::Value = serde_json::from_str(SUBMIT_RETRO_SCHEMA).unwrap(); @@ -415,7 +410,7 @@ mod tests { let stage_id = StageId::new("build", 2); let mut state = RunProjection::default(); state.graph_source = Some("digraph Ship {}".to_string()); - let stage = state.stage_entry(stage_id.node_id(), stage_id.visit(), nonzero(2)); + 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 { @@ -520,7 +515,7 @@ 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); - let stage = state.stage_entry(stage_id.node_id(), stage_id.visit(), nonzero(1)); + 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, diff --git a/lib/crates/fabro-store/src/run_state.rs b/lib/crates/fabro-store/src/run_state.rs index 98309385e..8527ab9d0 100644 --- a/lib/crates/fabro-store/src/run_state.rs +++ b/lib/crates/fabro-store/src/run_state.rs @@ -1,5 +1,4 @@ use std::collections::{BTreeMap, HashMap}; -use std::num::NonZeroU32; use std::str::FromStr; use chrono::{DateTime, Utc}; @@ -11,7 +10,7 @@ use fabro_types::{ BilledModelUsage, Checkpoint, Conclusion, EventBody, FailureSignature, InterviewQuestionRecord, Outcome, PendingInterviewRecord, PullRequestRecord, RunControlAction, RunId, RunProjection, RunSpec, RunStatus, RunSummary, SandboxRecord, StageCompletion, StageOutcome, StartRecord, - TerminalStatus, + TerminalStatus, first_event_seq, }; use fabro_util::error::render_with_causes; use serde_json::Value; @@ -403,10 +402,6 @@ impl RunProjectionReducer for RunProjection { } } -fn first_event_seq(seq: u32) -> NonZeroU32 { - NonZeroU32::new(seq).expect("event seq starts at 1") -} - pub(crate) fn build_summary(state: &RunProjection, run_id: &RunId) -> RunSummary { let workflow_name = state.spec.as_ref().map(|spec| { if spec.graph.name.is_empty() { @@ -607,7 +602,6 @@ fn provider_used_from_agent_cli_started(props: &AgentCliStartedProps) -> Value { #[cfg(test)] mod tests { use std::collections::{BTreeMap, HashMap}; - use std::num::NonZeroU32; use chrono::Utc; use fabro_types::run_event::run::RunFailedProps; @@ -618,17 +612,13 @@ mod tests { use fabro_types::{ BlockedReason, Checkpoint, EventBody, FailureReason, Outcome, QuestionType, RunBlobId, RunControlAction, RunEvent, RunStatus, StageOutcome, SuccessReason, TerminalStatus, - WorkflowSettings, fixtures, + WorkflowSettings, first_event_seq, fixtures, }; use serde_json::json; use super::{RunProjection, RunProjectionReducer, build_summary}; use crate::{Error, EventEnvelope, StageId}; - fn nonzero(value: u32) -> NonZeroU32 { - NonZeroU32::new(value).expect("test sequence must be non-zero") - } - fn test_event(seq: u32, body: EventBody, node_id: Option<&str>) -> EventEnvelope { let event = RunEvent { id: format!("evt-{seq}"), @@ -751,7 +741,7 @@ mod tests { let stage_id = StageId::new("build", 2); let node = state.stage(&stage_id).unwrap(); - assert_eq!(node.first_event_seq, nonzero(1)); + 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)); @@ -787,7 +777,7 @@ mod tests { restart_failure_signatures: HashMap::new(), node_visits: HashMap::from([("build".to_string(), 2usize)]), })]; - state.stage_entry("build", 2, nonzero(7)).stdout = Some("done".to_string()); + 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(); @@ -826,7 +816,7 @@ mod tests { .unwrap(); let stage = state.stage(&stage_id).unwrap(); - assert_eq!(stage.first_event_seq, nonzero(3)); + assert_eq!(stage.first_event_seq, first_event_seq(3)); } #[test] @@ -861,7 +851,7 @@ mod tests { .unwrap(); let stage = state.stage(&stage_id).unwrap(); - assert_eq!(stage.first_event_seq, nonzero(3)); + assert_eq!(stage.first_event_seq, first_event_seq(3)); assert_eq!(stage.prompt.as_deref(), Some("prompt")); } @@ -895,7 +885,7 @@ mod tests { .unwrap(); let stage = state.stage(&stage_id).unwrap(); - assert_eq!(stage.first_event_seq, nonzero(5)); + 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")); diff --git a/lib/crates/fabro-store/src/serializable_projection.rs b/lib/crates/fabro-store/src/serializable_projection.rs index 6bcf0c837..7d888b4a3 100644 --- a/lib/crates/fabro-store/src/serializable_projection.rs +++ b/lib/crates/fabro-store/src/serializable_projection.rs @@ -10,15 +10,7 @@ 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(stage) = projection.stage_mut(&stage_id) else { - continue; - }; + for (_, stage) in projection.iter_stages_mut() { stage.prompt = None; stage.response = None; stage.diff = None; diff --git a/lib/crates/fabro-store/tests/serializable_projection.rs b/lib/crates/fabro-store/tests/serializable_projection.rs index 1560b5c7c..6d51065ba 100644 --- a/lib/crates/fabro-store/tests/serializable_projection.rs +++ b/lib/crates/fabro-store/tests/serializable_projection.rs @@ -1,5 +1,4 @@ use std::collections::{BTreeMap, HashMap}; -use std::num::NonZeroU32; use chrono::{TimeZone, Utc}; use fabro_store::{RunProjection, SerializableProjection, StageId}; @@ -7,14 +6,10 @@ use fabro_types::graph::Graph; use fabro_types::run::RunSpec; use fabro_types::{ Checkpoint, RunStatus, SandboxRecord, StageCompletion, StageOutcome, StartRecord, - TerminalStatus, WorkflowSettings, fixtures, + TerminalStatus, WorkflowSettings, first_event_seq, fixtures, }; use serde_json::json; -fn nonzero(value: u32) -> NonZeroU32 { - NonZeroU32::new(value).expect("test sequence must be non-zero") -} - fn sample_run_spec() -> RunSpec { RunSpec { run_id: fixtures::RUN_1, @@ -82,7 +77,7 @@ fn serializable_projection_round_trips_and_trims_bulky_node_fields() { clone_branch: None, }); projection.pending_interviews = BTreeMap::new(); - let stage = projection.stage_entry(stage_id.node_id(), stage_id.visit(), nonzero(2)); + 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 { @@ -123,7 +118,7 @@ fn serializable_projection_round_trips_and_trims_bulky_node_fields() { assert_eq!(node.diff, None); assert_eq!(node.stdout, None); assert_eq!(node.stderr, None); - assert_eq!(node.first_event_seq, nonzero(2)); + assert_eq!(node.first_event_seq, first_event_seq(2)); assert_eq!( node.completion .as_ref() diff --git a/lib/crates/fabro-types/src/lib.rs b/lib/crates/fabro-types/src/lib.rs index 6181ff184..01f8e08f9 100644 --- a/lib/crates/fabro-types/src/lib.rs +++ b/lib/crates/fabro-types/src/lib.rs @@ -70,7 +70,7 @@ pub use run_event::{ MetadataSnapshotPhase, RunEvent, RunNoticeLevel, }; pub use run_id::{RunId, fixtures}; -pub use run_projection::{PendingInterviewRecord, RunProjection, StageProjection}; +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}; diff --git a/lib/crates/fabro-types/src/run_projection.rs b/lib/crates/fabro-types/src/run_projection.rs index 971068a7e..de2e4a1a5 100644 --- a/lib/crates/fabro-types/src/run_projection.rs +++ b/lib/crates/fabro-types/src/run_projection.rs @@ -63,6 +63,13 @@ pub struct StageProjection { 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 { @@ -96,6 +103,10 @@ impl RunProjection { 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() } diff --git a/lib/crates/fabro-workflow/src/pipeline/finalize.rs b/lib/crates/fabro-workflow/src/pipeline/finalize.rs index f158bf5bc..c3fe3c215 100644 --- a/lib/crates/fabro-workflow/src/pipeline/finalize.rs +++ b/lib/crates/fabro-workflow/src/pipeline/finalize.rs @@ -593,7 +593,6 @@ pub async fn finalize(retroed: Retroed, options: &FinalizeOptions) -> Result NonZeroU32 { - NonZeroU32::new(value).expect("test sequence must be non-zero") - } - fn checkpoint_with( completed_nodes: Vec<&str>, node_outcomes: HashMap, @@ -745,8 +742,8 @@ mod tests { #[test] fn conclusion_stage_order_follows_projection_first_event_order() { let mut projection = RunProjection::default(); - projection.stage_entry("zebra", 1, nonzero(1)); - projection.stage_entry("apple", 1, nonzero(2)); + 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"], @@ -777,8 +774,8 @@ mod tests { #[test] fn conclusion_includes_skipped_stage_from_projection_checkpoint_fallback() { let mut projection = RunProjection::default(); - projection.stage_entry("skipped", 1, nonzero(4)); - projection.stage_entry("finished", 1, nonzero(5)); + 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"], diff --git a/lib/crates/fabro-workflow/src/pipeline/pull_request.rs b/lib/crates/fabro-workflow/src/pipeline/pull_request.rs index 7067e5039..933f53052 100644 --- a/lib/crates/fabro-workflow/src/pipeline/pull_request.rs +++ b/lib/crates/fabro-workflow/src/pipeline/pull_request.rs @@ -625,7 +625,6 @@ pub async fn pull_request(concluded: Concluded, options: &PullRequestOptions) -> #[cfg(test)] mod tests { use std::collections::HashMap; - use std::num::NonZeroU32; use std::sync::Arc; use std::time::Duration; @@ -642,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; @@ -654,10 +653,6 @@ mod tests { use crate::event::{Event, append_event}; use crate::records::StageSummary; - fn nonzero(value: u32) -> NonZeroU32 { - NonZeroU32::new(value).expect("test sequence must be non-zero") - } - struct MockProvider { name: String, response_text: String, @@ -1002,7 +997,8 @@ mod tests { #[test] fn read_plan_text_found() { let mut state = RunProjection::default(); - state.stage_entry("plan", 1, nonzero(1)).response = Some("This is the plan".to_string()); + 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())); @@ -1011,8 +1007,9 @@ mod tests { #[test] fn read_plan_text_prefix_match() { let mut state = RunProjection::default(); - state.stage_entry("planning", 1, nonzero(1)).response = - Some("Planning content".to_string()); + 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())); @@ -1021,9 +1018,11 @@ mod tests { #[test] fn read_plan_text_prefers_alphabetically_first_plan_node() { let mut state = RunProjection::default(); - state.stage_entry("planning", 1, nonzero(1)).response = - Some("Planning content".to_string()); - state.stage_entry("plan", 1, nonzero(2)).response = Some("Plan content".to_string()); + 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())); @@ -1032,7 +1031,7 @@ mod tests { #[test] fn read_plan_text_not_found() { let mut state = RunProjection::default(); - state.stage_entry("implement", 1, nonzero(1)); + state.stage_entry("implement", 1, first_event_seq(1)); let result = read_plan_text(&state); assert_eq!(result, None); From 56a2257d8a52b508b648089e2380e01222ce8754 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 1 May 2026 20:47:40 -0400 Subject: [PATCH 3/3] refactor(run-projection): extract stage_at_visit helpers in reducer Eight reducer arms repeated the same node_id-presence check followed by a visit-derivation step (either explicit from props, or `current_visit_for(...).unwrap_or(1)`) and a `stage_entry` call. Pull those into `stage_at_visit` and `stage_at_current_visit` so each arm just binds the projection entry and writes its fields. The visit-derivation strategy is now legible from the helper name instead of buried in a free-floating `let visit = ...` line. Co-Authored-By: Claude Opus 4.7 (1M context) --- lib/crates/fabro-store/src/run_state.rs | 84 +++++++++++++------------ 1 file changed, 45 insertions(+), 39 deletions(-) diff --git a/lib/crates/fabro-store/src/run_state.rs b/lib/crates/fabro-store/src/run_state.rs index 8527ab9d0..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, - Outcome, PendingInterviewRecord, PullRequestRecord, RunControlAction, RunId, RunProjection, - RunSpec, RunStatus, RunSummary, SandboxRecord, StageCompletion, StageOutcome, StartRecord, - TerminalStatus, first_event_seq, + 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; @@ -297,22 +297,17 @@ impl RunProjectionReducer for RunProjection { ); } 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 visit = props.visit; - let provider_used = provider_used_from_prompt(props); - let stage = self.stage_entry(node_id, visit, first_event_seq(event.seq)); stage.prompt = Some(props.text.clone()); - stage.provider_used = provider_used; + 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(()); }; - let visit = self.current_visit_for(node_id).unwrap_or(1); - self.stage_entry(node_id, visit, first_event_seq(event.seq)) - .response = Some(props.response.clone()); + stage.response = Some(props.response.clone()); } EventBody::StageCompleted(props) => { let Some(node_id) = stored.node_id.as_deref() else { @@ -327,12 +322,10 @@ impl RunProjectionReducer for RunProjection { 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 visit = self.current_visit_for(node_id).unwrap_or(1); - let failure_reason = props.failure.as_ref().map(|detail| detail.message.clone()); - let stage = self.stage_entry(node_id, visit, first_event_seq(event.seq)); stage.completion = Some(StageCompletion { outcome: StageOutcome::Failed { retry_requested: false, @@ -343,38 +336,33 @@ 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(()); }; - self.stage_entry(node_id, props.visit, first_event_seq(event.seq)) - .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(()); }; - self.stage_entry(node_id, props.visit, first_event_seq(event.seq)) - .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(()); }; - let visit = self.current_visit_for(node_id).unwrap_or(1); - self.stage_entry(node_id, visit, first_event_seq(event.seq)) - .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 { - return Ok(()); - }; - let visit = self.current_visit_for(node_id).unwrap_or(1); let script_timing = serde_json::to_value(props).map_err(|err| { Error::InvalidEvent(format!("invalid command.completed payload: {err}")) })?; - let stage = self.stage_entry(node_id, visit, first_event_seq(event.seq)); + let Some(stage) = stage_at_current_visit(self, stored, event.seq) else { + return Ok(()); + }; stage.stdout = Some(props.stdout.clone()); stage.stderr = Some(props.stderr.clone()); stage.stdout_bytes = Some(props.stdout_bytes); @@ -385,15 +373,13 @@ impl RunProjectionReducer for RunProjection { stage.script_timing = Some(script_timing); } EventBody::ParallelCompleted(props) => { - let Some(node_id) = stored.node_id.as_deref() else { - return Ok(()); - }; - let visit = self.current_visit_for(node_id).unwrap_or(1); let parallel_results = serde_json::to_value(&props.results).map_err(|err| { Error::InvalidEvent(format!("invalid parallel.completed payload: {err}")) })?; - self.stage_entry(node_id, visit, first_event_seq(event.seq)) - .parallel_results = Some(parallel_results); + let Some(stage) = stage_at_current_visit(self, stored, event.seq) else { + return Ok(()); + }; + stage.parallel_results = Some(parallel_results); } _ => {} } @@ -402,6 +388,26 @@ impl RunProjectionReducer for RunProjection { } } +fn stage_at_visit<'a>( + state: &'a mut RunProjection, + stored: &RunEvent, + visit: u32, + 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_at_current_visit<'a>( + state: &'a mut RunProjection, + 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 { let workflow_name = state.spec.as_ref().map(|spec| { if spec.graph.name.is_empty() {