diff --git a/apps/fabro-web/app/components/stage-sidebar.tsx b/apps/fabro-web/app/components/stage-sidebar.tsx
index 4936eb0f8..46138ac5e 100644
--- a/apps/fabro-web/app/components/stage-sidebar.tsx
+++ b/apps/fabro-web/app/components/stage-sidebar.tsx
@@ -11,7 +11,7 @@ import {
} from "@heroicons/react/24/solid";
import { Bars3BottomLeftIcon, DocumentTextIcon, MapIcon } from "@heroicons/react/24/outline";
import { formatDurationSecs } from "../lib/format";
-import { ACTIVE_STAGE_STATES } from "../lib/stage-sidebar";
+import { ACTIVE_STAGE_STATES, formatStageLabel } from "../lib/stage-sidebar";
export interface Stage {
id: string;
@@ -101,7 +101,7 @@ export function StageSidebar({ stages, runId, selectedStageId, activeLink }: Sta
}`}
>
- {stage.visit > 1 ? `${stage.name} (${stage.visit})` : stage.name}
+ {formatStageLabel(stage)}
{stageDuration(stage)}
diff --git a/apps/fabro-web/app/lib/run-events.ts b/apps/fabro-web/app/lib/run-events.ts
index d964cc4d6..28115899b 100644
--- a/apps/fabro-web/app/lib/run-events.ts
+++ b/apps/fabro-web/app/lib/run-events.ts
@@ -136,10 +136,10 @@ export function subscribeToRunEvents(
}
function stageIdFromPayload(payload: RunEventPayload): string | undefined {
- if (typeof payload.stage_id === "string") return payload.stage_id;
- if (typeof payload.node_id === "string") return payload.node_id;
- const nodeId = payload.properties?.node_id;
- return typeof nodeId === "string" ? nodeId : undefined;
+ // Only return a true `node_id@visit` StageId. A bare `node_id` would not
+ // match the suffixed `stageTurns(runId, "verify@1")` cache key, so falling
+ // back to it would silently no-op the invalidation.
+ return typeof payload.stage_id === "string" ? payload.stage_id : undefined;
}
export function useRunEvents(runId: string | undefined) {
diff --git a/apps/fabro-web/app/lib/stage-sidebar.ts b/apps/fabro-web/app/lib/stage-sidebar.ts
index 73fdc1127..a820f5498 100644
--- a/apps/fabro-web/app/lib/stage-sidebar.ts
+++ b/apps/fabro-web/app/lib/stage-sidebar.ts
@@ -10,6 +10,15 @@ export const SUCCEEDED_STAGE_STATES: ReadonlySet = new Set([
"partially_succeeded",
]);
+/**
+ * Display label for a stage. Suffixes `(N)` for visits > 1 so a looped node
+ * (e.g. `verify`) renders as `verify`, `verify (2)`, `verify (3)` in the
+ * sidebar and stage header.
+ */
+export function formatStageLabel(stage: { name: string; visit: number }): string {
+ return stage.visit > 1 ? `${stage.name} (${stage.visit})` : stage.name;
+}
+
export function mapRunStagesToSidebarStages(
stagesResult: PaginatedRunStageList | null | undefined,
): Stage[] {
@@ -39,21 +48,27 @@ export function aggregateGraphNodeStatus(stages: readonly Stage[]): Map<
string,
{ displayStatus: StageState; latestStageId: string }
> {
- const grouped = new Map();
+ // Single pass per nodeId: track the visit with the highest `visit` overall
+ // (drives click target + terminal status) and the highest-visit *active*
+ // stage (drives display when any visit is in flight).
+ const latest = new Map();
+ const latestActive = new Map();
for (const stage of stages) {
- const list = grouped.get(stage.nodeId) ?? [];
- list.push(stage);
- grouped.set(stage.nodeId, list);
+ const prevLatest = latest.get(stage.nodeId);
+ if (!prevLatest || stage.visit > prevLatest.visit) {
+ latest.set(stage.nodeId, stage);
+ }
+ if (ACTIVE_STAGE_STATES.has(stage.status)) {
+ const prevActive = latestActive.get(stage.nodeId);
+ if (!prevActive || stage.visit > prevActive.visit) {
+ latestActive.set(stage.nodeId, stage);
+ }
+ }
}
const result = new Map();
- for (const [nodeId, list] of grouped) {
- list.sort((a, b) => a.visit - b.visit);
- const latest = list[list.length - 1];
- const activeVisit = [...list]
- .reverse()
- .find((s) => ACTIVE_STAGE_STATES.has(s.status));
- const display = activeVisit ?? latest;
- result.set(nodeId, { displayStatus: display.status, latestStageId: latest.id });
+ for (const [nodeId, latestStage] of latest) {
+ const display = latestActive.get(nodeId) ?? latestStage;
+ result.set(nodeId, { displayStatus: display.status, latestStageId: latestStage.id });
}
return result;
}
\ No newline at end of file
diff --git a/apps/fabro-web/app/routes/run-stages.tsx b/apps/fabro-web/app/routes/run-stages.tsx
index bb66bfd72..90ba4ed0c 100644
--- a/apps/fabro-web/app/routes/run-stages.tsx
+++ b/apps/fabro-web/app/routes/run-stages.tsx
@@ -41,7 +41,7 @@ import { EmptyState } from "../components/state";
import { CopyButton } from "../components/ui";
import { formatDurationSecs } from "../lib/format";
import { fetchRunCommandLog, useRunEventsList, useRunStageTurns, useRunStages } from "../lib/queries";
-import { mapRunStagesToSidebarStages } from "../lib/stage-sidebar";
+import { formatStageLabel, mapRunStagesToSidebarStages } from "../lib/stage-sidebar";
import { getNumber, getString, type UnknownRecord } from "../lib/unknown";
import {
CommandOutputStream,
@@ -638,7 +638,7 @@ export default function RunStages() {
- {selectedStage.visit > 1 ? `${selectedStage.name} (${selectedStage.visit})` : selectedStage.name}
+ {formatStageLabel(selectedStage)}
Router> {
.route("/runs/{id}/billing", get(get_run_billing))
}
-/// Pick the stage state from the latest lifecycle event for `stage_id`,
-/// falling back to the projection's stored completion when no lifecycle
-/// events have landed yet (e.g. an empty event log for a completed run
-/// recovered from snapshot only).
-fn stage_status_from_events(
- events: &[EventEnvelope],
- stage_id: &StageId,
- projection: &RunProjection,
-) -> StageState {
- let latest = events.iter().rev().find(|envelope| {
- envelope.event.stage_id.as_ref() == Some(stage_id)
- && matches!(
- &envelope.event.body,
- EventBody::StageStarted(_)
- | EventBody::StageRetrying(_)
- | EventBody::StageCompleted(_)
- | EventBody::StageFailed(_)
- )
- });
-
- if let Some(envelope) = latest {
- return match &envelope.event.body {
- EventBody::StageStarted(_) => StageState::Running,
- EventBody::StageRetrying(_) => StageState::Retrying,
- EventBody::StageFailed(props) => {
- if props.will_retry {
- StageState::Retrying
- } else {
- StageState::Failed
- }
- }
- EventBody::StageCompleted(props) => StageState::from(props.status),
- _ => StageState::Pending,
- };
+/// Map a `stage.*` lifecycle event body to the [`StageState`] it implies.
+/// Returns `None` for any other variant.
+fn stage_state_from_lifecycle(body: &EventBody) -> Option {
+ match body {
+ EventBody::StageStarted(_) => Some(StageState::Running),
+ EventBody::StageRetrying(_) => Some(StageState::Retrying),
+ EventBody::StageFailed(props) => Some(if props.will_retry {
+ StageState::Retrying
+ } else {
+ StageState::Failed
+ }),
+ EventBody::StageCompleted(props) => Some(StageState::from(props.status)),
+ _ => None,
}
+}
- projection
- .stage(stage_id)
- .and_then(|stage| stage.completion.as_ref())
- .map_or(StageState::Pending, |c| StageState::from(c.outcome))
+/// Single-pass scan over `events` building the latest [`StageState`] for each
+/// [`StageId`] from lifecycle events (started/retrying/completed/failed). Each
+/// later lifecycle event overwrites earlier ones, leaving the latest as the
+/// stored value — equivalent to "scan in reverse, take first match" but in O(E)
+/// for the whole list rather than O(stages × events).
+fn latest_stage_states(events: &[EventEnvelope]) -> HashMap {
+ let mut states = HashMap::new();
+ for envelope in events {
+ let Some(stage_id) = envelope.event.stage_id.as_ref() else {
+ continue;
+ };
+ let Some(state) = stage_state_from_lifecycle(&envelope.event.body) else {
+ continue;
+ };
+ states.insert(stage_id.clone(), state);
+ }
+ states
}
async fn list_run_stages(
@@ -77,21 +70,30 @@ async fn list_run_stages(
let projection = RunProjection::apply_events(&events).unwrap_or_default();
let stage_durations = fabro_workflow::extract_stage_durations_by_stage_id(&events);
+ let lifecycle_states = latest_stage_states(&events);
let mut entries: Vec<(&StageId, &fabro_types::StageProjection)> =
projection.iter_stages().collect();
- entries.sort_by_key(|(_, projection)| projection.first_event_seq);
+ entries.sort_by_key(|(_, stage)| stage.first_event_seq);
let mut stages = Vec::with_capacity(entries.len());
- for (stage_id, _projection_stage) in entries {
- let duration_ms = stage_durations.get(stage_id).copied();
+ for (stage_id, stage_projection) in entries {
+ let node_id = stage_id.node_id().to_string();
let visit = NonZeroU32::new(stage_id.visit()).expect("StageId.visit is 1-based");
+ // Prefer the latest lifecycle event; fall back to the projection's
+ // stored completion (e.g. for runs recovered from snapshot only).
+ let status = lifecycle_states.get(stage_id).copied().unwrap_or_else(|| {
+ stage_projection
+ .completion
+ .as_ref()
+ .map_or(StageState::Pending, |c| StageState::from(c.outcome))
+ });
stages.push(RunStage {
id: stage_id.to_string(),
- name: stage_id.node_id().to_string(),
- status: stage_status_from_events(&events, stage_id, &projection),
- duration_secs: duration_ms.map(|ms| ms as f64 / 1000.0),
- node_id: stage_id.node_id().to_string(),
+ name: node_id.clone(),
+ status,
+ duration_secs: stage_durations.get(stage_id).map(|ms| *ms as f64 / 1000.0),
+ node_id,
visit,
});
}
@@ -215,4 +217,4 @@ async fn get_run_billing(
};
(StatusCode::OK, Json(response)).into_response()
-}
+}
\ No newline at end of file
diff --git a/lib/crates/fabro-workflow/src/lib.rs b/lib/crates/fabro-workflow/src/lib.rs
index fd6398bec..40a487330 100644
--- a/lib/crates/fabro-workflow/src/lib.rs
+++ b/lib/crates/fabro-workflow/src/lib.rs
@@ -20,6 +20,7 @@ use std::sync::Arc;
use fabro_retro::retro::CompletedStage;
use fabro_store::EventEnvelope;
+use fabro_types::EventBody;
/// Callback invoked when a workflow node starts executing.
pub type OnNodeCallback = Option>;
@@ -86,23 +87,23 @@ pub fn build_completed_stages(cp: &records::Checkpoint, run_failed: bool) -> Vec
stages
}
+/// Extract the `duration_ms` from a `stage.completed` / `stage.failed`
+/// event body, or `None` for any other variant.
+fn stage_completion_duration_ms(body: &EventBody) -> Option {
+ match body {
+ EventBody::StageCompleted(props) => Some(props.duration_ms),
+ EventBody::StageFailed(props) => Some(props.duration_ms),
+ _ => None,
+ }
+}
+
pub fn extract_stage_durations_from_events(events: &[EventEnvelope]) -> HashMap {
let mut durations = HashMap::new();
for envelope in events {
- let event = &envelope.event;
- let event_name = event.event_name();
- if event_name != "stage.completed" && event_name != "stage.failed" {
- continue;
- }
- let Some(node_id) = event.node_id.as_deref() else {
+ let Some(duration_ms) = stage_completion_duration_ms(&envelope.event.body) else {
continue;
};
- let Some(duration_ms) = event
- .properties()
- .ok()
- .and_then(|properties| properties.get("duration_ms").cloned())
- .and_then(|duration| duration.as_u64())
- else {
+ let Some(node_id) = envelope.event.node_id.as_deref() else {
continue;
};
durations.insert(node_id.to_string(), duration_ms);
@@ -120,20 +121,10 @@ pub fn extract_stage_durations_by_stage_id(
) -> HashMap {
let mut durations = HashMap::new();
for envelope in events {
- let event = &envelope.event;
- let event_name = event.event_name();
- if event_name != "stage.completed" && event_name != "stage.failed" {
- continue;
- }
- let Some(stage_id) = event.stage_id.as_ref() else {
+ let Some(duration_ms) = stage_completion_duration_ms(&envelope.event.body) else {
continue;
};
- let Some(duration_ms) = event
- .properties()
- .ok()
- .and_then(|properties| properties.get("duration_ms").cloned())
- .and_then(|duration| duration.as_u64())
- else {
+ let Some(stage_id) = envelope.event.stage_id.as_ref() else {
continue;
};
durations.insert(stage_id.clone(), duration_ms);
@@ -188,4 +179,4 @@ mod stage_scope;
pub mod test_support;
#[doc(hidden)]
pub mod transforms;
-pub mod workflow_bundle;
+pub mod workflow_bundle;
\ No newline at end of file