diff --git a/apps/fabro-web/app/app.css b/apps/fabro-web/app/app.css index cd7443e35..eccab3275 100644 --- a/apps/fabro-web/app/app.css +++ b/apps/fabro-web/app/app.css @@ -151,6 +151,21 @@ --diffs-fg-number-deletion-override: rgba(232, 107, 107, 0.5); } +/* ── Server-rendered SVG graph overrides ── */ +/* The server injects @media (prefers-color-scheme: dark) but the app uses + .dark/.light classes, so we apply dark-mode colors via the class. */ + +.dark .graph-svg text { fill: #c6d4e0; } +.dark .graph-svg [stroke="#357f9e"] { stroke: #67B2D7; } +.dark .graph-svg [stroke="#666666"] { stroke: #5a7a94; } +.dark .graph-svg polygon[fill="#357f9e"] { fill: #67B2D7; } +.dark .graph-svg polygon[fill="#666666"] { fill: #5a7a94; } +.dark .graph-svg ellipse[fill="#ffffff"], +.dark .graph-svg polygon[fill="#ffffff"] { fill: #1a2b3c; } +.dark .graph-svg ellipse[stroke="#357f9e"], +.dark .graph-svg polygon[stroke="#357f9e"] { stroke: #67B2D7; } +.dark .graph-svg [fill="#1a1a1a"] { fill: #c6d4e0; } + .light { --diffs-bg-buffer-override: #ffffff; --diffs-bg-context-override: #f8fafc; diff --git a/apps/fabro-web/app/data/runs.ts b/apps/fabro-web/app/data/runs.ts index b6629cbf2..8fc2bc09b 100644 --- a/apps/fabro-web/app/data/runs.ts +++ b/apps/fabro-web/app/data/runs.ts @@ -144,6 +144,13 @@ export function isRunStatus(s: string): s is RunStatus { return knownRunStatuses.has(s); } +/** Graph control nodes hidden from stage lists in the UI. */ +const hiddenStageIds = new Set(["start", "exit"]); + +export function isVisibleStage(id: string): boolean { + return !hiddenStageIds.has(id); +} + export const ciConfig: Record = { passing: { label: "Passing", dot: "bg-mint", text: "text-mint" }, failing: { label: "Changes needed", dot: "bg-coral", text: "text-coral" }, diff --git a/apps/fabro-web/app/routes/run-graph.tsx b/apps/fabro-web/app/routes/run-graph.tsx index ae7a494b0..e7ce63cbc 100644 --- a/apps/fabro-web/app/routes/run-graph.tsx +++ b/apps/fabro-web/app/routes/run-graph.tsx @@ -6,6 +6,7 @@ import { DocumentTextIcon, MapIcon } from "@heroicons/react/24/outline"; import { useTheme } from "../lib/theme"; import { getGraphTheme } from "../lib/graph-theme"; import { apiFetch, apiJsonOrNull } from "../api"; +import { isVisibleStage } from "../data/runs"; import { formatDurationSecs } from "../lib/format"; import type { PaginatedRunStageList } from "@qltysh/fabro-api-client"; @@ -26,7 +27,7 @@ export async function loader({ request, params }: any) { apiJsonOrNull(`/runs/${params.id}/stages`, { request }), apiFetch(`/runs/${params.id}/graph`, { request }), ]); - const stages: Stage[] = (stagesResult?.data ?? []).map((s) => ({ + const stages: Stage[] = (stagesResult?.data ?? []).filter((s) => isVisibleStage(s.id)).map((s) => ({ id: s.id, name: s.name, dotId: s.dot_id ?? s.id, @@ -208,12 +209,27 @@ export default function RunGraph({ loaderData }: any) { let cancelled = false; async function render() { - const { instance } = await import("@viz-js/viz"); - const viz = await instance(); - if (cancelled) return; - try { - const svg = viz.renderSVGElement(buildDot(direction, graphTheme)); + let svg: SVGSVGElement; + + if (graphSvg) { + // Use server-rendered SVG from the real graph endpoint. + const parser = new DOMParser(); + const doc = parser.parseFromString(graphSvg, "image/svg+xml"); + const parsed = doc.documentElement; + if (!(parsed instanceof SVGSVGElement)) { + setError("Invalid SVG from server"); + return; + } + svg = parsed; + } else { + // Fall back to hardcoded demo graph rendered client-side. + const { instance } = await import("@viz-js/viz"); + const viz = await instance(); + if (cancelled) return; + svg = viz.renderSVGElement(buildDot(direction, graphTheme)); + } + stripGraphTitle(svg); annotateRunningNodes(svg, graphTheme, stages); @@ -229,7 +245,7 @@ export default function RunGraph({ loaderData }: any) { setPan({ x: 0, y: 0 }); render(); return () => { cancelled = true; }; - }, [direction, graphTheme]); + }, [direction, graphTheme, graphSvg]); const onPointerDown = useCallback((e: React.PointerEvent) => { if ((e.target as HTMLElement).closest("button")) return; @@ -325,7 +341,7 @@ export default function RunGraph({ loaderData }: any) {
-
+
- -
- -
- -
- -
- - -
-
- -
-
-

Loading diagram...

-
-
-
- ); -} +const STAGE_EVENTS = new Set([ + "stage.started", "stage.completed", "stage.failed", + "run.completed", "run.failed", +]); export default function RunOverview({ loaderData }: any) { const { id } = useParams(); - const { stages, graphDot } = loaderData; + const { stages, graphSvg } = loaderData; + 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 SSE for live stage updates + 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 === id && STAGE_EVENTS.has(payload.event)) { + clearTimeout(debounceTimer); + debounceTimer = setTimeout(() => revalidator.revalidate(), 300); + } + } catch { + // ignore malformed events + } + }; + + return () => { + clearTimeout(debounceTimer); + source.close(); + }; + }, [id]); + + // Track start times for running stages + useEffect(() => { + const running = new Set( + stages.filter((s: Stage) => s.status === "running").map((s: Stage) => s.id), + ); + for (const stageId of running) { + if (!runningStartRef.current.has(stageId)) { + runningStartRef.current.set(stageId, Date.now()); + } + } + for (const stageId of runningStartRef.current.keys()) { + if (!running.has(stageId)) { + runningStartRef.current.delete(stageId); + } + } + }, [stages]); + + // Tick every second while any stage is running + useEffect(() => { + if (!stages.some((s: Stage) => s.status === "running")) return; + const interval = setInterval(() => setTick((t) => t + 1), 1000); + return () => clearInterval(interval); + }, [stages]); + + function stageDuration(stage: Stage): string { + if (stage.status === "running") { + const start = runningStartRef.current.get(stage.id); + if (start) return formatDurationSecs(Math.floor((Date.now() - start) / 1000)); + return "0s"; + } + return stage.duration; + } return (
- {graphDot ? ( -
- -
+ {graphSvg ? ( +
) : (

No workflow graph available.

)} diff --git a/apps/fabro-web/app/routes/run-settings.tsx b/apps/fabro-web/app/routes/run-settings.tsx index 59d5d0f7b..0cfdc89ed 100644 --- a/apps/fabro-web/app/routes/run-settings.tsx +++ b/apps/fabro-web/app/routes/run-settings.tsx @@ -3,6 +3,7 @@ import { CheckCircleIcon, ArrowPathIcon, PauseCircleIcon, XCircleIcon } from "@h import { DocumentTextIcon, MapIcon } from "@heroicons/react/24/outline"; import { CollapsibleFile } from "../components/collapsible-file"; import { apiJson } from "../api"; +import { isVisibleStage } from "../data/runs"; import { formatDurationSecs } from "../lib/format"; import type { PaginatedRunStageList } from "@qltysh/fabro-api-client"; import type { RunSettings } from "../lib/workflow-api"; @@ -31,7 +32,7 @@ export async function loader({ request, params }: any) { apiJson(`/runs/${params.id}/stages`, { request }), apiJson(`/runs/${params.id}/settings`, { request }), ]); - const stages: Stage[] = apiStages.map((s) => ({ + const stages: Stage[] = apiStages.filter((s) => isVisibleStage(s.id)).map((s) => ({ id: s.id, name: s.name, status: s.status as StageStatus, diff --git a/apps/fabro-web/app/routes/run-stages.tsx b/apps/fabro-web/app/routes/run-stages.tsx index 8edf3af2a..2f3d04ff2 100644 --- a/apps/fabro-web/app/routes/run-stages.tsx +++ b/apps/fabro-web/app/routes/run-stages.tsx @@ -1,13 +1,12 @@ -import { useState } from "react"; import { Link, useParams } from "react-router"; -import { ChevronRightIcon } from "@heroicons/react/20/solid"; import { CheckCircleIcon, ArrowPathIcon, PauseCircleIcon, XCircleIcon } from "@heroicons/react/24/solid"; import { DocumentTextIcon, MapIcon, CommandLineIcon, ChatBubbleLeftIcon } from "@heroicons/react/24/outline"; -import { ToolRow, ToolBlock } from "../components/tool-use"; +import { ToolBlock } from "../components/tool-use"; import type { ToolUse } from "../components/tool-use"; -import { apiJson } from "../api"; +import { apiJson, apiJsonOrNull } from "../api"; +import { isVisibleStage } from "../data/runs"; import { formatDurationSecs } from "../lib/format"; -import type { PaginatedRunStageList, StageTurn as ApiStageTurn, PaginatedStageTurnList } from "@qltysh/fabro-api-client"; +import type { PaginatedRunStageList, StageTurn as ApiStageTurn, PaginatedStageTurnList, PaginatedEventList } from "@qltysh/fabro-api-client"; export const handle = { wide: true }; @@ -20,21 +19,113 @@ interface Stage { duration: string; } +type TurnType = + | { kind: "system"; content: string } + | { kind: "assistant"; content: string } + | { kind: "tool"; tools: ToolUse[] }; + +interface RawEvent { + node_id?: string; + event: string; + properties?: Record; + text?: string; + tool_name?: string; + tool_call_id?: string; + arguments?: unknown; + output?: unknown; + is_error?: boolean; +} + +function turnsFromEvents(events: RawEvent[], stageId: string): TurnType[] { + const stageEvents = events.filter((e) => e.node_id === stageId); + const turns: TurnType[] = []; + // Collect tool pairs: started → completed + const pendingTools = new Map(); + + for (const e of stageEvents) { + switch (e.event) { + case "stage.prompt": + turns.push({ kind: "system", content: e.properties?.text as string ?? e.text ?? "" }); + break; + case "agent.message": + turns.push({ kind: "assistant", content: e.properties?.text as string ?? e.text ?? "" }); + break; + case "agent.tool.started": { + const callId = e.properties?.tool_call_id as string ?? e.tool_call_id ?? ""; + pendingTools.set(callId, { + toolName: e.properties?.tool_name as string ?? e.tool_name ?? "", + input: typeof (e.properties?.arguments ?? e.arguments) === "string" + ? (e.properties?.arguments ?? e.arguments) as string + : JSON.stringify(e.properties?.arguments ?? e.arguments ?? ""), + }); + break; + } + case "agent.tool.completed": { + const callId = e.properties?.tool_call_id as string ?? e.tool_call_id ?? ""; + const started = pendingTools.get(callId); + const output = e.properties?.output ?? e.output ?? ""; + const result = typeof output === "string" ? output : JSON.stringify(output); + const tool: ToolUse = { + id: callId, + toolName: started?.toolName ?? e.properties?.tool_name as string ?? e.tool_name ?? "", + input: started?.input ?? "", + result, + isError: (e.properties?.is_error ?? e.is_error) === true, + }; + pendingTools.delete(callId); + turns.push({ kind: "tool", tools: [tool] }); + break; + } + } + } + return turns; +} + export async function loader({ request, params }: any) { const { data: apiStages } = await apiJson(`/runs/${params.id}/stages`, { request }); - const stages: Stage[] = apiStages.map((s) => ({ + const stages: Stage[] = apiStages.filter((s) => isVisibleStage(s.id)).map((s) => ({ id: s.id, name: s.name, status: s.status as StageStatus, duration: s.duration_secs != null ? formatDurationSecs(s.duration_secs) : "--", })); - // Fetch turns for the selected stage (first stage if none specified) const selectedStageId = params.stageId ?? stages[0]?.id; - let turns: ApiStageTurn[] = []; + + // Try demo turns endpoint first, fall back to events. + let turns: TurnType[] = []; if (selectedStageId) { - const { data } = await apiJson(`/runs/${params.id}/stages/${selectedStageId}/turns`, { request }); - turns = data; + const turnsResult = await apiJsonOrNull( + `/runs/${params.id}/stages/${selectedStageId}/turns`, + { request }, + ); + if (turnsResult?.data?.length) { + turns = turnsResult.data.map((t: ApiStageTurn): TurnType => { + if (t.kind === "tool" && t.tools) { + return { + kind: "tool", + tools: t.tools.map((tu) => ({ + id: tu.id, + toolName: tu.tool_name, + input: tu.input, + result: tu.result, + isError: tu.is_error, + durationMs: tu.duration_ms, + })), + }; + } + return { kind: t.kind as "system" | "assistant", content: t.content ?? "" }; + }); + } else { + // Fetch events and build turns from them. + const eventsResult = await apiJsonOrNull( + `/runs/${params.id}/events?limit=1000`, + { request }, + ); + if (eventsResult?.data) { + turns = turnsFromEvents(eventsResult.data as unknown as RawEvent[], selectedStageId); + } + } } return { stages, turns }; @@ -48,13 +139,6 @@ const statusConfig: Record @@ -85,27 +169,10 @@ function AssistantBlock({ content }: { content: string }) { export default function RunStages({ loaderData }: any) { const { id, stageId } = useParams(); - const { stages, turns: apiTurns } = loaderData; + const { stages, turns } = loaderData; - const mappedTurns: TurnType[] = apiTurns.map((t) => { - if (t.kind === "tool" && t.tools) { - return { - kind: "tool" as const, - tools: t.tools.map((tu) => ({ - id: tu.id, - toolName: tu.tool_name, - input: tu.input, - result: tu.result, - isError: tu.is_error, - durationMs: tu.duration_ms, - })), - }; - } - return { kind: t.kind as "system" | "assistant", content: t.content ?? "" }; - }); - - const selectedStage = stages.find((s) => s.id === stageId) ?? stages[0]; - const selectedConfig = statusConfig[selectedStage.status]; + const selectedStage = stages.find((s: Stage) => s.id === stageId) ?? stages[0]; + const selectedConfig = selectedStage ? statusConfig[selectedStage.status] : statusConfig.pending; const SelectedIcon = selectedConfig.icon; return ( @@ -170,7 +237,7 @@ export default function RunStages({ loaderData }: any) { {selectedStage.duration}
- {mappedTurns.map((turn, i) => { + {turns.map((turn: TurnType, i: number) => { switch (turn.kind) { case "system": return ; diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 5de6af9ca..80c57e77a 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -32,10 +32,11 @@ pub use fabro_api::types::{ PruneRunEntry, PruneRunsRequest, PruneRunsResponse, QuestionType as ApiQuestionType, RenderWorkflowGraphDirection, RenderWorkflowGraphRequest, RunArtifactEntry, RunArtifactListResponse, RunBilling, RunBillingStage, RunBillingTotals, - RunControlAction as ApiRunControlAction, RunError, RunManifest, RunStatus, RunStatusResponse, - SandboxFileEntry, SandboxFileListResponse, SecretType as ApiSecretType, ServerSettings, - SshAccessRequest, SshAccessResponse, StartRunRequest, StatusReason as ApiStatusReason, - SubmitAnswerRequest, SystemFeatures, SystemInfoResponse, SystemRunCounts, WriteBlobResponse, + RunControlAction as ApiRunControlAction, RunError, RunManifest, RunStage, RunStatus, + RunStatusResponse, SandboxFileEntry, SandboxFileListResponse, SecretType as ApiSecretType, + ServerSettings, SshAccessRequest, SshAccessResponse, StageStatus as ApiStageStatus, + StartRunRequest, StatusReason as ApiStatusReason, SubmitAnswerRequest, SystemFeatures, + SystemInfoResponse, SystemRunCounts, WriteBlobResponse, }; use fabro_auth::parse_credential_secret; use fabro_config::{Storage, resolve_server_from_file}; @@ -1080,7 +1081,7 @@ fn real_routes() -> Router> { .route("/runs/{id}/pause", post(pause_run)) .route("/runs/{id}/unpause", post(unpause_run)) .route("/runs/{id}/graph", get(get_graph)) - .route("/runs/{id}/stages", get(not_implemented)) + .route("/runs/{id}/stages", get(list_run_stages)) .route("/runs/{id}/artifacts", get(list_run_artifacts)) .route("/runs/{id}/stages/{stageId}/turns", get(not_implemented)) .route( @@ -2059,6 +2060,104 @@ async fn get_aggregate_billing( (StatusCode::OK, Json(response)).into_response() } +async fn list_run_stages( + _auth: AuthenticatedService, + State(state): State>, + Path(id): Path, + Query(_pagination): Query, +) -> Response { + let id = match parse_run_id_path(&id) { + Ok(id) => id, + Err(response) => return response, + }; + + // Try live run first. + let (checkpoint, run_is_active) = { + let runs = state.runs.lock().expect("runs lock poisoned"); + match runs.get(&id) { + Some(managed_run) => { + let active = !matches!( + managed_run.status, + RunStatus::Completed | RunStatus::Failed | RunStatus::Cancelled + ); + (managed_run.checkpoint.clone(), active) + } + None => (None, false), + } + }; + + // Fall back to stored run. + let (checkpoint, run_is_active) = if checkpoint.is_some() { + (checkpoint, run_is_active) + } else { + match state.store.open_run_reader(&id).await { + Ok(run_store) => match run_store.state().await { + Ok(run_state) => { + let active = run_state + .status + .as_ref() + .map_or(false, |s| !s.status.is_terminal()); + (run_state.checkpoint, active) + } + Err(_) => (None, false), + }, + Err(_) => return ApiError::not_found("Run not found.").into_response(), + } + }; + + let Some(checkpoint) = checkpoint else { + return ( + StatusCode::OK, + Json(ListResponse::new(Vec::::new())), + ) + .into_response(); + }; + + // Get durations from events. + let stage_durations = match state.store.open_run_reader(&id).await { + Ok(run_store) => match run_store.list_events().await { + Ok(events) => fabro_workflow::extract_stage_durations_from_events(&events), + Err(_) => HashMap::new(), + }, + Err(_) => HashMap::new(), + }; + + let mut stages = Vec::new(); + for node_id in &checkpoint.completed_nodes { + let duration_ms = stage_durations.get(node_id).copied().unwrap_or(0); + let status = match checkpoint.node_outcomes.get(node_id) { + Some(outcome) => match outcome.status { + fabro_types::outcome::StageStatus::Success + | fabro_types::outcome::StageStatus::PartialSuccess => ApiStageStatus::Completed, + fabro_types::outcome::StageStatus::Fail => ApiStageStatus::Failed, + fabro_types::outcome::StageStatus::Skipped => ApiStageStatus::Cancelled, + fabro_types::outcome::StageStatus::Retry => ApiStageStatus::Pending, + }, + None => ApiStageStatus::Completed, + }; + stages.push(RunStage { + id: node_id.clone(), + name: node_id.clone(), + status, + duration_secs: Some(duration_ms as f64 / 1000.0), + dot_id: Some(node_id.clone()), + }); + } + + // 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()), + }); + } + + (StatusCode::OK, Json(ListResponse::new(stages))).into_response() +} + async fn get_run_billing( _auth: AuthenticatedService, State(state): State>, @@ -6411,13 +6510,12 @@ async fn get_graph( }; let live_dot_source = { let runs = state.runs.lock().expect("runs lock poisoned"); - match runs.get(&id) { - Some(managed_run) => managed_run.dot_source.clone(), - None => return ApiError::not_found("Run not found.").into_response(), - } + runs.get(&id).map(|managed_run| managed_run.dot_source.clone()) }; - if !live_dot_source.is_empty() { - return render_graph_bytes(&live_dot_source).await; + if let Some(dot) = &live_dot_source { + if !dot.is_empty() { + return render_graph_bytes(dot).await; + } } match state.store.open_run_reader(&id).await {