diff --git a/apps/fabro-web/app/components/stage-sidebar.tsx b/apps/fabro-web/app/components/stage-sidebar.tsx index 409fdd628..2970d0745 100644 --- a/apps/fabro-web/app/components/stage-sidebar.tsx +++ b/apps/fabro-web/app/components/stage-sidebar.tsx @@ -1,5 +1,5 @@ import { useState, useEffect, useRef } from "react"; -import { Link } from "react-router"; +import { Link, useRevalidator } from "react-router"; import { CheckCircleIcon, ArrowPathIcon, PauseCircleIcon, XCircleIcon } from "@heroicons/react/24/solid"; import { DocumentTextIcon, MapIcon } from "@heroicons/react/24/outline"; import { formatDurationSecs } from "../lib/format"; @@ -29,11 +29,41 @@ interface StageSidebarProps { activeLink?: "settings" | "graph"; } +const STAGE_EVENTS = new Set([ + "stage.started", "stage.completed", "stage.failed", + "run.completed", "run.failed", +]); + export function StageSidebar({ stages, runId, selectedStageId, activeLink }: StageSidebarProps) { + const revalidator = useRevalidator(); + // Track when we first observed each running stage (for ticking timer) const runningStartRef = useRef>(new Map()); const [, setTick] = useState(0); + // Subscribe to run-specific SSE for live stage updates + useEffect(() => { + const source = new EventSource(`/api/v1/runs/${runId}/attach?since_seq=1`); + let debounceTimer: ReturnType | undefined; + + source.onmessage = (msg) => { + try { + const payload = JSON.parse(msg.data); + if (STAGE_EVENTS.has(payload.event)) { + clearTimeout(debounceTimer); + debounceTimer = setTimeout(() => revalidator.revalidate(), 300); + } + } catch { + // ignore malformed events + } + }; + + return () => { + clearTimeout(debounceTimer); + source.close(); + }; + }, [runId]); + // Track start times for running stages useEffect(() => { const running = new Set( diff --git a/apps/fabro-web/app/routes/run-detail.tsx b/apps/fabro-web/app/routes/run-detail.tsx index 16093b9c1..d6ad109e5 100644 --- a/apps/fabro-web/app/routes/run-detail.tsx +++ b/apps/fabro-web/app/routes/run-detail.tsx @@ -1,7 +1,7 @@ import { useEffect } from "react"; import { ChevronDownIcon, ChevronRightIcon } from "@heroicons/react/20/solid"; import { Menu, MenuButton, MenuItem, MenuItems } from "@headlessui/react"; -import { Link, Outlet, useFetcher, useLocation, useRevalidator } from "react-router"; +import { Link, Outlet, useFetcher, useLocation } from "react-router"; import { mapRunSummaryToRunItem, runStatusDisplay, isRunStatus } from "../data/runs"; import type { RunSummaryResponse } from "../data/runs"; import { apiJson } from "../api"; @@ -59,11 +59,6 @@ export function meta({ data }: any) { return [{ title: run ? `${run.title} — Fabro` : "Run — Fabro" }]; } -const RUN_EVENTS = new Set([ - "stage.started", "stage.completed", "stage.failed", - "run.completed", "run.failed", "run.running", "run.paused", -]); - export default function RunDetail({ loaderData, params }: any) { const { run } = loaderData; const { pathname } = useLocation(); @@ -71,30 +66,6 @@ export default function RunDetail({ loaderData, params }: any) { const previewFetcher = useFetcher(); const demoMode = useDemoMode(); const tabs = allTabs.filter((t) => !t.broken && (!t.demoOnly || demoMode)); - const revalidator = useRevalidator(); - - // Subscribe to SSE for live updates (stages + run status) - useEffect(() => { - const source = new EventSource("/api/v1/attach"); - let debounceTimer: ReturnType | undefined; - - source.onmessage = (msg) => { - try { - const payload = JSON.parse(msg.data); - if (payload.run_id === params.id && RUN_EVENTS.has(payload.event)) { - clearTimeout(debounceTimer); - debounceTimer = setTimeout(() => revalidator.revalidate(), 300); - } - } catch { - // ignore malformed events - } - }; - - return () => { - clearTimeout(debounceTimer); - source.close(); - }; - }, [params.id]); useEffect(() => { if (previewFetcher.data?.url) { diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 80c57e77a..c8e2f5175 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -2144,15 +2144,22 @@ async fn list_run_stages( }); } - // Add current node as running if the run is still active. - if run_is_active && !checkpoint.completed_nodes.contains(&checkpoint.current_node) { - stages.push(RunStage { - id: checkpoint.current_node.clone(), - name: checkpoint.current_node.clone(), - status: ApiStageStatus::Running, - duration_secs: None, - dot_id: Some(checkpoint.current_node.clone()), - }); + // Add next node as running if the run is still active. + // The checkpoint's current_node is the last *completed* stage; next_node_id + // is the stage that is currently executing. + if let Some(next_id) = &checkpoint.next_node_id { + if run_is_active + && next_id != "exit" + && !checkpoint.completed_nodes.contains(next_id) + { + stages.push(RunStage { + id: next_id.clone(), + name: next_id.clone(), + status: ApiStageStatus::Running, + duration_secs: None, + dot_id: Some(next_id.clone()), + }); + } } (StatusCode::OK, Json(ListResponse::new(stages))).into_response()