From 05a14a6c9b15904c07c00f8bad8e250a7e7c7233 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sat, 19 Sep 2026 09:31:09 -0400 Subject: [PATCH] Delete the checkpoint endpoint, fabro parse, and the fabro-workflow shims GET /runs/{id}/checkpoint duplicated what /state serves; the hidden fabro parse command had no user; records, run_status, outcome, and usage_rollup in fabro-workflow only re-exported fabro_types. The importers now name fabro_types directly. format_cost keeps its two callers (the pull request body and the CLI stage display) and moves to fabro_types::usage; the usage rollup tests move beside the function in fabro-types, with test_usage in its test support. Co-Authored-By: Claude Fable 5.1 --- docs/public/api-reference/fabro-api.yaml | 27 -- lib/apps/fabro-cli/src/args.rs | 10 - lib/apps/fabro-cli/src/commands/mod.rs | 1 - lib/apps/fabro-cli/src/commands/parse.rs | 30 -- lib/apps/fabro-cli/src/commands/run/attach.rs | 4 +- lib/apps/fabro-cli/src/commands/run/output.rs | 7 +- .../run/run_progress/stage_display.rs | 3 +- lib/apps/fabro-cli/src/commands/run/wait.rs | 7 +- .../fabro-cli/src/commands/runs/inspect.rs | 3 +- lib/apps/fabro-cli/src/commands/runs/list.rs | 2 +- lib/apps/fabro-cli/src/main.rs | 3 - lib/apps/fabro-cli/src/shared/utilities.rs | 5 - lib/apps/fabro-cli/tests/it/cmd/mod.rs | 1 - lib/apps/fabro-cli/tests/it/cmd/parse.rs | 135 --------- lib/apps/fabro-server/src/demo/mod.rs | 8 - lib/apps/fabro-server/src/server.rs | 13 +- .../src/server/handler/artifacts.rs | 19 -- .../fabro-server/src/server/handler/mod.rs | 1 - .../fabro-server/src/server/handler/runs.rs | 8 +- .../fabro-server/src/server/handler/steer.rs | 3 +- .../fabro-server/src/server/handler/usage.rs | 2 +- .../fabro-server/src/server/petri_runs.rs | 5 +- lib/apps/fabro-server/src/server/tests.rs | 35 +-- lib/components/fabro-workflow/src/lib.rs | 16 +- .../fabro-workflow/src/operations/create.rs | 5 +- lib/components/fabro-workflow/src/outcome.rs | 13 - .../fabro-workflow/src/pipeline/persist.rs | 3 +- .../fabro-workflow/src/pipeline/types.rs | 2 +- .../fabro-workflow/src/pull_request.rs | 12 +- .../fabro-workflow/src/records/conclusion.rs | 1 - .../fabro-workflow/src/records/mod.rs | 8 - .../fabro-workflow/src/records/run.rs | 1 - .../fabro-workflow/src/records/start.rs | 1 - .../fabro-workflow/src/run_lookup.rs | 6 +- .../fabro-workflow/src/run_status.rs | 1 - .../fabro-workflow/src/test_support.rs | 28 -- .../fabro-workflow/src/usage_rollup.rs | 270 ------------------ lib/foundation/fabro-types/src/lib.rs | 2 +- .../fabro-types/src/test_support.rs | 28 +- lib/foundation/fabro-types/src/usage.rs | 6 + .../fabro-types/src/usage_rollup.rs | 266 +++++++++++++++++ .../src/api/run-internals-api.ts | 76 ----- 42 files changed, 340 insertions(+), 737 deletions(-) delete mode 100644 lib/apps/fabro-cli/src/commands/parse.rs delete mode 100644 lib/apps/fabro-cli/tests/it/cmd/parse.rs delete mode 100644 lib/components/fabro-workflow/src/outcome.rs delete mode 100644 lib/components/fabro-workflow/src/records/conclusion.rs delete mode 100644 lib/components/fabro-workflow/src/records/mod.rs delete mode 100644 lib/components/fabro-workflow/src/records/run.rs delete mode 100644 lib/components/fabro-workflow/src/records/start.rs delete mode 100644 lib/components/fabro-workflow/src/run_status.rs delete mode 100644 lib/components/fabro-workflow/src/test_support.rs delete mode 100644 lib/components/fabro-workflow/src/usage_rollup.rs diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 646c1fad7..f01dd5595 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -2648,33 +2648,6 @@ paths: schema: $ref: "#/components/schemas/ErrorResponse" - /api/v1/runs/{id}/checkpoint: - get: - operationId: retrieveRunCheckpoint - tags: [Run Internals] - summary: Retrieve Run Checkpoint - description: Returns the latest checkpoint data for a run, or null if no checkpoint has been recorded yet. - parameters: - - $ref: "#/components/parameters/RunId" - responses: - "200": - description: Checkpoint data (null if not yet available) - content: - application/json: - schema: - oneOf: - - $ref: "#/components/schemas/RunCheckpoint" - - type: "null" - "404": - description: Run not found - headers: - x-request-id: - $ref: "#/components/headers/XRequestId" - content: - application/json: - schema: - $ref: "#/components/schemas/ErrorResponse" - /api/v1/runs/{id}/state: get: operationId: getRunState diff --git a/lib/apps/fabro-cli/src/args.rs b/lib/apps/fabro-cli/src/args.rs index e7dd3e3a1..9db43de8f 100644 --- a/lib/apps/fabro-cli/src/args.rs +++ b/lib/apps/fabro-cli/src/args.rs @@ -549,12 +549,6 @@ pub(crate) struct GraphArgs { pub(crate) allow_invalid: bool, } -#[derive(Args)] -pub(crate) struct ParseArgs { - /// Path to the .fabro workflow file - pub(crate) workflow: PathBuf, -} - #[derive(Args)] pub(crate) struct ArtifactListArgs { #[command(flatten)] @@ -1442,9 +1436,6 @@ pub(crate) enum Commands { Validate(ValidateArgs), /// Render a workflow graph as SVG Graph(GraphArgs), - /// Parse a DOT file and print its AST - #[command(hide = true)] - Parse(ParseArgs), /// Inspect and copy run artifacts (screenshots, reports, traces) Artifact(ArtifactNamespace), /// Export a run's durable state to a directory @@ -1546,7 +1537,6 @@ impl Commands { Self::Preflight(_) => "preflight", Self::Validate(_) => "validate", Self::Graph(_) => "graph", - Self::Parse(_) => "parse", Self::RunsCmd(cmd) => cmd.name(), Self::Model { command } => match command { Some(ModelsCommand::List(_)) => "model list", diff --git a/lib/apps/fabro-cli/src/commands/mod.rs b/lib/apps/fabro-cli/src/commands/mod.rs index 18e22d5f7..a5c357a53 100644 --- a/lib/apps/fabro-cli/src/commands/mod.rs +++ b/lib/apps/fabro-cli/src/commands/mod.rs @@ -10,7 +10,6 @@ pub(crate) mod install; pub(crate) mod mcp; pub(crate) mod model; pub(crate) mod parent; -pub(crate) mod parse; pub(crate) mod pr; pub(crate) mod preflight; pub(crate) mod provider; diff --git a/lib/apps/fabro-cli/src/commands/parse.rs b/lib/apps/fabro-cli/src/commands/parse.rs deleted file mode 100644 index 505671617..000000000 --- a/lib/apps/fabro-cli/src/commands/parse.rs +++ /dev/null @@ -1,30 +0,0 @@ -#![expect( - clippy::disallowed_types, - reason = "sync CLI `parse` command: blocking std::io::Write is the intended output mechanism" -)] -#![expect( - clippy::disallowed_methods, - reason = "sync CLI `parse` command: blocking std::io::stdout is the intended output mechanism" -)] - -use std::io::Write; - -use fabro_config::project::resolve_workflow; -use fabro_graphviz::parser::parse_ast; - -use crate::args::ParseArgs; -use crate::shared::read_workflow_file; - -pub(crate) fn run(args: &ParseArgs) -> anyhow::Result<()> { - let stdout = std::io::stdout(); - run_to(args, stdout.lock()) -} - -fn run_to(args: &ParseArgs, mut out: impl Write) -> anyhow::Result<()> { - let dot_path = resolve_workflow(&args.workflow)?; - let source = read_workflow_file(&dot_path)?; - let ast = parse_ast(&source)?; - serde_json::to_writer_pretty(&mut out, &ast)?; - writeln!(out)?; - Ok(()) -} diff --git a/lib/apps/fabro-cli/src/commands/run/attach.rs b/lib/apps/fabro-cli/src/commands/run/attach.rs index 02b99ece5..dc55ce336 100644 --- a/lib/apps/fabro-cli/src/commands/run/attach.rs +++ b/lib/apps/fabro-cli/src/commands/run/attach.rs @@ -21,11 +21,9 @@ use anyhow::Result; use fabro_api::types; use fabro_interview::{Answer, AnswerValue, Question}; use fabro_types::settings::run::ApprovalMode; -use fabro_types::{InterviewOption, QuestionType, RunId}; +use fabro_types::{InterviewOption, QuestionType, RunId, RunStatus, StageOutcome}; use fabro_util::printer::Printer; use fabro_util::terminal::Styles; -use fabro_workflow::outcome::StageOutcome; -use fabro_workflow::run_status::RunStatus; use tokio::signal::ctrl_c; use tokio::time::{Duration as TokioDuration, sleep}; diff --git a/lib/apps/fabro-cli/src/commands/run/output.rs b/lib/apps/fabro-cli/src/commands/run/output.rs index c1a935aa8..fee4e859f 100644 --- a/lib/apps/fabro-cli/src/commands/run/output.rs +++ b/lib/apps/fabro-cli/src/commands/run/output.rs @@ -6,14 +6,15 @@ use cli_table::format::{Border, Justify, Separator}; use cli_table::{Cell, CellStruct, Style, Table}; use fabro_api::types; use fabro_types::diagnostic::{Diagnostic, RelatedDiagnostic, Severity}; -use fabro_types::{BlobRefEncoding, PullRequestLink, RunId, StageId, parse_blob_ref_encoded}; +use fabro_types::{ + BlobRefEncoding, Conclusion, PullRequestLink, RunId, StageId, StageOutcome, + parse_blob_ref_encoded, +}; use fabro_util::check_report::{CheckDetail, CheckReport, CheckResult, CheckSection, CheckStatus}; use fabro_util::error::render_with_causes; use fabro_util::printer::Printer; use fabro_util::terminal::Styles; use fabro_util::text::strip_goal_decoration; -use fabro_workflow::outcome::StageOutcome; -use fabro_workflow::records::Conclusion; use indicatif::HumanDuration; use crate::server_client; diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/stage_display.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/stage_display.rs index 2a5c863ac..b00a35820 100644 --- a/lib/apps/fabro-cli/src/commands/run/run_progress/stage_display.rs +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/stage_display.rs @@ -3,8 +3,7 @@ use std::convert::TryFrom; use std::time::Duration; use chrono::{DateTime, Utc}; -use fabro_types::{INITIAL_SUBAGENT_GENERATION, LlmOutputKind}; -use fabro_workflow::outcome::{StageOutcome, format_cost}; +use fabro_types::{INITIAL_SUBAGENT_GENERATION, LlmOutputKind, StageOutcome, format_cost}; use indicatif::ProgressBar; use super::event::ProgressUsage; diff --git a/lib/apps/fabro-cli/src/commands/run/wait.rs b/lib/apps/fabro-cli/src/commands/run/wait.rs index feccf3678..5169aae66 100644 --- a/lib/apps/fabro-cli/src/commands/run/wait.rs +++ b/lib/apps/fabro-cli/src/commands/run/wait.rs @@ -10,11 +10,9 @@ use std::io::Write; use anyhow::{Result, bail}; -use fabro_types::RunId; +use fabro_types::{Conclusion, RunId, RunStatus}; use fabro_util::printer::Printer; use fabro_util::terminal::Styles; -use fabro_workflow::records::Conclusion; -use fabro_workflow::run_status::RunStatus; use tokio::time; use tracing::info; @@ -134,10 +132,9 @@ fn print_human_output( #[cfg(test)] mod tests { use fabro_types::{ - FailureCategory, FailureDetail, FailureReason, RunDiff, RunFailure, RunStatus, + Conclusion, FailureCategory, FailureDetail, FailureReason, RunDiff, RunFailure, RunStatus, StageOutcome, SuccessReason, fixtures, }; - use fabro_workflow::records::Conclusion; use lithos_llm::types::{Cost, CostSource, TokenCounts, Usage}; use super::*; diff --git a/lib/apps/fabro-cli/src/commands/runs/inspect.rs b/lib/apps/fabro-cli/src/commands/runs/inspect.rs index de06e8684..ad5789e09 100644 --- a/lib/apps/fabro-cli/src/commands/runs/inspect.rs +++ b/lib/apps/fabro-cli/src/commands/runs/inspect.rs @@ -1,8 +1,7 @@ use std::collections::BTreeMap; use anyhow::Result; -use fabro_types::{StageHandler, StageState}; -use fabro_workflow::run_status::RunStatus; +use fabro_types::{RunStatus, StageHandler, StageState}; use serde::Serialize; use crate::args::InspectArgs; diff --git a/lib/apps/fabro-cli/src/commands/runs/list.rs b/lib/apps/fabro-cli/src/commands/runs/list.rs index e2e964567..412b68cfd 100644 --- a/lib/apps/fabro-cli/src/commands/runs/list.rs +++ b/lib/apps/fabro-cli/src/commands/runs/list.rs @@ -4,9 +4,9 @@ use anyhow::Result; use chrono::Utc; use cli_table::format::{Border, Separator}; use cli_table::{Cell, CellStruct, Color, Style, Table}; +use fabro_types::RunStatus; use fabro_util::terminal::Styles; use fabro_util::text::strip_goal_decoration; -use fabro_workflow::run_status::RunStatus; use super::short_run_id; use crate::args::RunsListArgs; diff --git a/lib/apps/fabro-cli/src/main.rs b/lib/apps/fabro-cli/src/main.rs index befac102b..775dda5ed 100644 --- a/lib/apps/fabro-cli/src/main.rs +++ b/lib/apps/fabro-cli/src/main.rs @@ -295,9 +295,6 @@ async fn main_inner(worker_token: Option) -> (String, Result<()>) { let styles = Styles::detect_stderr(); commands::graph::run(&args, &styles, &base_ctx).await?; } - Commands::Parse(args) => { - commands::parse::run(&args)?; - } Commands::Artifact(ns) => { commands::artifact::dispatch(ns, &base_ctx).await?; } diff --git a/lib/apps/fabro-cli/src/shared/utilities.rs b/lib/apps/fabro-cli/src/shared/utilities.rs index 77a5de01b..29c15e57e 100644 --- a/lib/apps/fabro-cli/src/shared/utilities.rs +++ b/lib/apps/fabro-cli/src/shared/utilities.rs @@ -11,7 +11,6 @@ use std::io::Write; use std::path::{Path, PathBuf}; use std::time::Duration; -use anyhow::Context as _; use cli_table::Color; use fabro_types::RunStatus; use fabro_types::diagnostic::{Diagnostic, Severity}; @@ -32,10 +31,6 @@ pub(crate) fn cyan_spinner(message: impl Into>) - spinner } -pub(crate) fn read_workflow_file(path: &Path) -> anyhow::Result { - std::fs::read_to_string(path).with_context(|| format!("Failed to read {}", path.display())) -} - pub(crate) fn print_json_pretty(value: &T) -> anyhow::Result<()> where T: Serialize + ?Sized, diff --git a/lib/apps/fabro-cli/tests/it/cmd/mod.rs b/lib/apps/fabro-cli/tests/it/cmd/mod.rs index 54f7f1d5a..f3c0b9b4c 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/mod.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/mod.rs @@ -26,7 +26,6 @@ mod model; mod model_list; mod model_test; mod parent; -mod parse; mod pr; mod pr_close; mod pr_create; diff --git a/lib/apps/fabro-cli/tests/it/cmd/parse.rs b/lib/apps/fabro-cli/tests/it/cmd/parse.rs deleted file mode 100644 index ca80cf735..000000000 --- a/lib/apps/fabro-cli/tests/it/cmd/parse.rs +++ /dev/null @@ -1,135 +0,0 @@ -use fabro_test::{fabro_snapshot, test_context}; - -#[test] -fn help() { - let context = test_context!(); - let mut cmd = context.command(); - cmd.args(["parse", "--help"]); - fabro_snapshot!(context.filters(), cmd, @" - success: true - exit_code: 0 - ----- stdout ----- - Parse a DOT file and print its AST - - Usage: fabro parse [OPTIONS] - - Arguments: - Path to the .fabro workflow file - - Options: - --json Output as JSON [env: FABRO_JSON=] - --debug Enable DEBUG-level logging (default is INFO) [env: FABRO_DEBUG=] - --no-upgrade-check Disable automatic upgrade check [env: FABRO_NO_UPGRADE_CHECK=true] - --quiet Suppress non-essential output [env: FABRO_QUIET=] - --verbose Enable verbose output [env: FABRO_VERBOSE=] - -h, --help Print help - ----- stderr ----- - "); -} - -#[test] -fn parse_valid_workflow_prints_ast_json() { - let context = test_context!(); - context.write_temp( - "tiny.fabro", - "digraph Tiny {\n graph [goal=\"Parse a tiny workflow\"]\n start [shape=Mdiamond]\n exit [shape=Msquare]\n main [label=\"Main\", prompt=\"Do the thing\"]\n start -> main -> exit\n}\n", - ); - let mut cmd = context.command(); - cmd.args(["parse", "tiny.fabro"]); - - fabro_snapshot!(context.filters(), cmd, @r###" - success: true - exit_code: 0 - ----- stdout ----- - { - "name": "Tiny", - "statements": [ - { - "GraphAttr": [ - [ - "goal", - { - "Str": "Parse a tiny workflow" - } - ] - ] - }, - { - "Node": { - "id": "start", - "attrs": [ - [ - "shape", - { - "Ident": "Mdiamond" - } - ] - ] - } - }, - { - "Node": { - "id": "exit", - "attrs": [ - [ - "shape", - { - "Ident": "Msquare" - } - ] - ] - } - }, - { - "Node": { - "id": "main", - "attrs": [ - [ - "label", - { - "Str": "Main" - } - ], - [ - "prompt", - { - "Str": "Do the thing" - } - ] - ] - } - }, - { - "Edge": { - "nodes": [ - "start", - "main", - "exit" - ], - "attrs": null - } - } - ] - } - ----- stderr ----- - "###); -} - -#[test] -fn parse_invalid_dot_fails_cleanly() { - let context = test_context!(); - context.write_temp( - "bad.fabro", - "digraph Bad {\n start [shape=Mdiamond]\n exit [shape=Msquare]\n start -> exit\n", - ); - let mut cmd = context.command(); - cmd.args(["parse", "bad.fabro"]); - - fabro_snapshot!(context.filters(), cmd, @r#" - success: false - exit_code: 1 - ----- stdout ----- - ----- stderr ----- - × Parse error: grammar error: Parsing Error: Error { input: "", code: Char } - "#); -} diff --git a/lib/apps/fabro-server/src/demo/mod.rs b/lib/apps/fabro-server/src/demo/mod.rs index d57d59c13..d4ca46711 100644 --- a/lib/apps/fabro-server/src/demo/mod.rs +++ b/lib/apps/fabro-server/src/demo/mod.rs @@ -474,14 +474,6 @@ pub(crate) async fn run_events_stub( Sse::new(tokio_stream::iter(events)).into_response() } -pub(crate) async fn checkpoint_stub( - _auth: RequiredUser, - State(_state): State>, - Path(_id): Path, -) -> Response { - (StatusCode::OK, Json(serde_json::json!(null))).into_response() -} - pub(crate) async fn cancel_stub( _auth: RequiredUser, State(_state): State>, diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index b85d6a34f..c0a40627a 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -94,10 +94,10 @@ use fabro_types::settings::server::{ GithubIntegrationSettings, GithubIntegrationStrategy, LogDestination, }; use fabro_types::{ - AskFabro, AskFabroUnavailableReason, BlobHash, InterviewQuestionRecord, ModelRef, - ModelTestMode, PendingReason, Principal, PullRequestLink, QuestionType, RunControlAction, - RunId, RunRunnableSource, RunStatusKind, RunStreamItem, RunStreamItemKind, SandboxProviderKind, - ServerSettings, + AskFabro, AskFabroUnavailableReason, BlobHash, FailureReason, InterviewQuestionRecord, + ModelRef, ModelTestMode, PendingReason, Principal, PullRequestLink, QuestionType, + RunControlAction, RunId, RunRunnableSource, RunStatus, RunStatusKind, RunStreamItem, + RunStreamItemKind, SandboxProviderKind, ServerSettings, SuccessReason, }; use fabro_util::error::{ SharedError, collect_causes, render_compact_with_causes, render_with_causes, @@ -108,7 +108,6 @@ use fabro_vault::{SecretStore, SecretStoreError, SecretType, Vault}; use fabro_workflow::run_lookup::{ RunInfo, StatusFilter, filter_runs, scan_runs_with_summaries, scratch_base, }; -use fabro_workflow::run_status::{FailureReason, RunStatus, SuccessReason}; use fabro_workflow::{Error as WorkflowError, operations, pull_request}; use futures_util::future::join_all; use lithos_llm::catalog::ProviderId; @@ -1355,13 +1354,13 @@ pub(crate) fn accumulate_concluded_run_usage( .expect("aggregate_usage lock poisoned"); accumulate_usage_rollup( &mut agg, - &fabro_workflow::usage_rollup_from_projection(final_state), + &fabro_types::usage_rollup::usage_rollup_from_projection(final_state), ); } fn accumulate_usage_rollup( accumulator: &mut UsageAccumulator, - rollup: &fabro_workflow::ProjectionUsageRollup, + rollup: &fabro_types::usage_rollup::ProjectionUsageRollup, ) { accumulator.total_runs += 1; accumulator.total_timing = accumulator.total_timing.saturating_add(&rollup.timing); diff --git a/lib/apps/fabro-server/src/server/handler/artifacts.rs b/lib/apps/fabro-server/src/server/handler/artifacts.rs index 3ccc0a8ba..bf505e4c5 100644 --- a/lib/apps/fabro-server/src/server/handler/artifacts.rs +++ b/lib/apps/fabro-server/src/server/handler/artifacts.rs @@ -29,7 +29,6 @@ use super::super::{ pub(super) fn routes() -> Router> { Router::new() - .route("/runs/{id}/checkpoint", get(get_checkpoint)) .route("/runs/{id}/blobs", post(write_run_blob)) .route("/runs/{id}/blobs/{blobHash}", get(read_run_blob)) .route("/runs/{id}/artifacts", get(list_run_artifacts)) @@ -52,24 +51,6 @@ struct ArtifactFilenameParams { retry: Option, } -async fn get_checkpoint( - _auth: RequiredUser, - State(state): State>, - Path(id): Path, -) -> Response { - let id = match parse_run_id_path(&id) { - Ok(id) => id, - Err(response) => return response, - }; - match state.load_run_projection(&id).await { - Ok(projection) => match projection.current_checkpoint() { - Some(cp) => (StatusCode::OK, Json(cp.clone())).into_response(), - None => (StatusCode::OK, Json(serde_json::json!(null))).into_response(), - }, - Err(err) => err.into_response(), - } -} - async fn write_run_blob( RequireRunScoped(id): RequireRunScoped, State(state): State>, diff --git a/lib/apps/fabro-server/src/server/handler/mod.rs b/lib/apps/fabro-server/src/server/handler/mod.rs index a596b049a..886cce4d6 100644 --- a/lib/apps/fabro-server/src/server/handler/mod.rs +++ b/lib/apps/fabro-server/src/server/handler/mod.rs @@ -107,7 +107,6 @@ pub(super) fn demo_routes() -> Router> { "/runs/{id}/stages/{stageId}/logs/output", get(not_implemented), ) - .route("/runs/{id}/checkpoint", get(demo::checkpoint_stub)) .route("/runs/{id}/cancel", post(demo::cancel_stub)) .route("/runs/{id}/start", post(demo::start_run_stub)) .route("/runs/{id}/approve", post(demo::start_run_stub)) diff --git a/lib/apps/fabro-server/src/server/handler/runs.rs b/lib/apps/fabro-server/src/server/handler/runs.rs index 9dec4b39d..8705a43e4 100644 --- a/lib/apps/fabro-server/src/server/handler/runs.rs +++ b/lib/apps/fabro-server/src/server/handler/runs.rs @@ -31,14 +31,14 @@ use fabro_types::diagnostic::Severity; use fabro_types::settings::run::RunMode; use fabro_types::{ AutomationRef, ContextWindowStaleness, ManifestPath, Principal, Run, RunClientProvenance, - RunId, RunProvenance, RunServerProvenance, RunStatusKind, RunTarget, SandboxProviderKind, - StageContextWindow, StageContextWindowUnavailableReason, StageHandler, StageModelUsage, - StageProjection, ValidatedRunTarget, json_scalar_to_toml_value, parse_blob_ref, + RunId, RunProvenance, RunServerProvenance, RunStatus, RunStatusKind, RunTarget, + SandboxProviderKind, StageContextWindow, StageContextWindowUnavailableReason, StageHandler, + StageModelUsage, StageProjection, ValidatedRunTarget, json_scalar_to_toml_value, + parse_blob_ref, }; use fabro_util::error as error_util; use fabro_util::version::FABRO_VERSION; use fabro_workflow::pipeline::Validated; -use fabro_workflow::run_status::RunStatus; use fabro_workflow::{Error as WorkflowError, operations}; use lithos_llm::catalog::ProviderId; use serde::de::IgnoredAny; diff --git a/lib/apps/fabro-server/src/server/handler/steer.rs b/lib/apps/fabro-server/src/server/handler/steer.rs index 281aa58fd..2d2dfe033 100644 --- a/lib/apps/fabro-server/src/server/handler/steer.rs +++ b/lib/apps/fabro-server/src/server/handler/steer.rs @@ -8,8 +8,7 @@ use axum::routing::post; use fabro_api::types::{ InterruptRunRequest, RunControlAcknowledgement, RunControlOutcome, SteerRunRequest, }; -use fabro_types::Principal; -use fabro_workflow::run_status::RunStatus; +use fabro_types::{Principal, RunStatus}; use super::super::{ AnswerTransportError, AppState, RunControlAnswer, durable_run_status, reject_if_archived, diff --git a/lib/apps/fabro-server/src/server/handler/usage.rs b/lib/apps/fabro-server/src/server/handler/usage.rs index e92082e6b..f20a8a99d 100644 --- a/lib/apps/fabro-server/src/server/handler/usage.rs +++ b/lib/apps/fabro-server/src/server/handler/usage.rs @@ -93,7 +93,7 @@ async fn get_run_usage( Err(err) => return err.into_response(), }; - let rollup = fabro_workflow::usage_rollup_from_projection(&projection); + let rollup = fabro_types::usage_rollup::usage_rollup_from_projection(&projection); let by_model = rollup .by_model .iter() diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 227eefcaf..0cb319dd6 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -49,10 +49,11 @@ use fabro_petri::{SqliteRunStore, admission}; use fabro_store::platform_records::{RunLifecycleKind, RunLifecycleRecord}; use fabro_types::settings::McpTransport; use fabro_types::settings::run::{ApprovalMode, McpServerSettings, RunMode}; -use fabro_types::{PetriAdmission, RunId, RunRunnableSource, RunTarget}; +use fabro_types::{ + FailureReason, PetriAdmission, RunId, RunRunnableSource, RunStatus, RunTarget, SuccessReason, +}; use fabro_util::error as error_util; use fabro_workflow::Error as WorkflowError; -use fabro_workflow::run_status::{FailureReason, RunStatus, SuccessReason}; use lithos_llm::catalog::ProviderId; use tokio::task; use tokio_util::sync::CancellationToken; diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index 3aaba9f85..78e2aa5a0 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -7149,34 +7149,6 @@ async fn list_run_events_returns_paginated_json() { assert!(body["meta"]["has_more"].is_boolean()); } -#[tokio::test] -async fn get_checkpoint_returns_null_initially() { - let state = test_app_state(); - let app = crate::test_support::build_test_router(Arc::clone(&state)); - - // Start a run - let req = Request::builder() - .method("POST") - .uri(api("/runs")) - .header("content-type", "application/json") - .body(intent_body(&app, MINIMAL_DOT).await) - .unwrap(); - - let response = app.clone().oneshot(req).await.unwrap(); - let body = body_json(response.into_body()).await; - let run_id = body["id"].as_str().unwrap().parse::().unwrap(); - - // Get checkpoint immediately (before run completes, may be null) - let req = Request::builder() - .method("GET") - .uri(api(&format!("/runs/{run_id}/checkpoint"))) - .body(Body::empty()) - .unwrap(); - - let response = app.oneshot(req).await.unwrap(); - checked_response!(response, StatusCode::OK).await; -} - #[tokio::test] async fn write_and_read_run_blob_accepts_uppercase_hash() { let state = test_app_state(); @@ -7736,7 +7708,6 @@ async fn worker_token_is_rejected_on_user_only_routes() { (Method::GET, "/attach".to_string()), (Method::DELETE, format!("/runs/{run_id}")), (Method::GET, format!("/runs/{run_id}/attach")), - (Method::GET, format!("/runs/{run_id}/checkpoint")), (Method::POST, format!("/runs/{run_id}/pause")), (Method::POST, format!("/runs/{run_id}/unpause")), (Method::GET, format!("/runs/{run_id}/graph")), @@ -9341,11 +9312,11 @@ async fn get_aggregate_usage_saturates_total_cost_across_models() { #[test] fn aggregate_usage_counts_projection_rollup_usage_visits() { let mut accumulator = UsageAccumulator::default(); - let rollup = fabro_workflow::ProjectionUsageRollup { + let rollup = fabro_types::usage_rollup::ProjectionUsageRollup { stages: Vec::new(), totals: test_priced_usage("gpt-5.4", 300, 30).usage, by_model: vec![ - fabro_workflow::ProjectionUsageByModel { + fabro_types::usage_rollup::ProjectionUsageByModel { model: ModelRef::new( lithos_llm::catalog::builtin::openai(), ModelId::new("gpt-5.4"), @@ -9353,7 +9324,7 @@ fn aggregate_usage_counts_projection_rollup_usage_visits() { stages: 1, usage: test_priced_usage("gpt-5.4", 100, 10).usage, }, - fabro_workflow::ProjectionUsageByModel { + fabro_types::usage_rollup::ProjectionUsageByModel { model: ModelRef::new( lithos_llm::catalog::builtin::openai(), ModelId::new("gpt-5.4"), diff --git a/lib/components/fabro-workflow/src/lib.rs b/lib/components/fabro-workflow/src/lib.rs index 5d1907fb1..1b90af1cb 100644 --- a/lib/components/fabro-workflow/src/lib.rs +++ b/lib/components/fabro-workflow/src/lib.rs @@ -3,12 +3,12 @@ //! //! Petri executes every run (`fabro-petri` is the seam). This crate keeps //! what Fabro itself owns: the create-time compile of the Fabro graph the -//! read side displays (`pipeline`, `transforms`, `operations`), the run -//! records and status vocabulary (`records`, `run_status`), the Git +//! read side displays (`pipeline`, `transforms`, `operations`), the Git //! helpers a run's platform effects use (`git`, `sandbox_git`), pull //! request creation (`pull_request`), the run tools an agent session calls //! (`run_tools`, `services`), the built-in web search backend -//! (`web_search`). +//! (`web_search`). The run records and status vocabulary are +//! `fabro_types`'. #![cfg_attr( test, @@ -31,26 +31,16 @@ pub mod error; pub mod file_resolver; pub mod git; pub mod operations; -pub mod outcome; pub mod pipeline; pub mod pull_request; -pub mod records; pub mod run_lookup; -pub mod usage_rollup; pub use error::{Error, Result}; pub use fabro_types::ManifestPath; -pub use usage_rollup::{ - ProjectionUsageByModel, ProjectionUsageRollup, ProjectionUsageStage, - usage_rollup_from_projection, -}; pub mod run_materialization; -pub mod run_status; pub mod run_tools; pub mod sandbox_git; pub mod services; -#[cfg(any(test, feature = "test-support"))] -pub mod test_support; #[doc(hidden)] pub mod transforms; pub mod web_search; diff --git a/lib/components/fabro-workflow/src/operations/create.rs b/lib/components/fabro-workflow/src/operations/create.rs index 7bb58843f..3bac02edf 100644 --- a/lib/components/fabro-workflow/src/operations/create.rs +++ b/lib/components/fabro-workflow/src/operations/create.rs @@ -18,7 +18,7 @@ use fabro_store::{BlobStore, Database}; use fabro_template::TemplateContext; use fabro_types::{ AutomationRef, BlobHash, ForkSourceRef, GitContext, ManifestPath, PetriAdmission, RunId, - RunProvenance, RunStatus, RunTarget, WorkflowSettings, WorkflowVersionId, + RunProvenance, RunSpec, RunStatus, RunTarget, WorkflowSettings, WorkflowVersionId, }; use tokio::task::spawn_blocking; @@ -26,7 +26,6 @@ use super::source::{ResolveWorkflowInput, WorkflowInput, resolve_workflow}; use crate::error::Error; use crate::pipeline::types::PersistOptions; use crate::pipeline::{self, Persisted, TransformOptions, Validated}; -use crate::records::RunSpec; use crate::run_materialization; use crate::transforms::RenderMode; use crate::workflow_bundle::{RunDefinition, WorkflowBundle}; @@ -1786,7 +1785,7 @@ mod tests { ); assert_eq!( last_lifecycle_status(&platform_records(&store, fixtures::RUN_1).await), - Some(crate::run_status::RunStatus::Submitted) + Some(fabro_types::RunStatus::Submitted) ); assert_eq!( created.run_dir, diff --git a/lib/components/fabro-workflow/src/outcome.rs b/lib/components/fabro-workflow/src/outcome.rs deleted file mode 100644 index 289145f13..000000000 --- a/lib/components/fabro-workflow/src/outcome.rs +++ /dev/null @@ -1,13 +0,0 @@ -pub use fabro_types::ModelUsage; -pub use fabro_types::outcome::{ - FailureCategory, FailureDetail, OutcomeMeta, StageOutcome, StageState, -}; - -/// A stage outcome carrying the model usage the stage reported. -pub type Outcome = fabro_types::Outcome>; - -/// Format a USD cost for display, to the cent. -#[must_use] -pub fn format_cost(cost: f64) -> String { - format!("${cost:.2}") -} diff --git a/lib/components/fabro-workflow/src/pipeline/persist.rs b/lib/components/fabro-workflow/src/pipeline/persist.rs index 5b3f3ca7b..3f1cdadb9 100644 --- a/lib/components/fabro-workflow/src/pipeline/persist.rs +++ b/lib/components/fabro-workflow/src/pipeline/persist.rs @@ -32,10 +32,9 @@ mod tests { use std::collections::HashMap; use fabro_graphviz::graph::{AttrValue, Edge, Graph, Node}; - use fabro_types::{PetriAdmission, fixtures, test_support}; + use fabro_types::{PetriAdmission, RunSpec, fixtures, test_support}; use super::*; - use crate::records::RunSpec; fn graph_and_source() -> (Graph, String) { let source = r#"digraph test { diff --git a/lib/components/fabro-workflow/src/pipeline/types.rs b/lib/components/fabro-workflow/src/pipeline/types.rs index e87a60566..6952dbe2f 100644 --- a/lib/components/fabro-workflow/src/pipeline/types.rs +++ b/lib/components/fabro-workflow/src/pipeline/types.rs @@ -3,11 +3,11 @@ use std::sync::Arc; use fabro_graphviz::graph::Graph; use fabro_template::TemplateContext; +use fabro_types::RunSpec; use fabro_types::diagnostic::{Diagnostic, Severity}; use crate::error::Error; use crate::file_resolver::FileResolver; -use crate::records::RunSpec; use crate::transforms::{RenderMode, Transform}; /// Output of the PARSE phase. diff --git a/lib/components/fabro-workflow/src/pull_request.rs b/lib/components/fabro-workflow/src/pull_request.rs index 1a78ac9f5..2bd662c77 100644 --- a/lib/components/fabro-workflow/src/pull_request.rs +++ b/lib/components/fabro-workflow/src/pull_request.rs @@ -8,17 +8,14 @@ use fabro_llm::credentials::CredentialProvider; use fabro_llm::lithos_catalog::Catalog; use fabro_llm::{Client, ClientOptions, Request, selection}; use fabro_store::RunProjection; -use fabro_types::PullRequestLink; use fabro_types::settings::run::MergeStrategy; +use fabro_types::{Conclusion, PullRequestLink, RunSpec, format_cost as outcome_format_cost}; use fabro_util::text::strip_goal_decoration; use lithos_llm::catalog::ProviderId; use lithos_llm::types::{Cost, Message, Role}; use tokio::time::sleep; use tracing::{debug, info, warn}; -use crate::outcome::format_cost as outcome_format_cost; -use crate::records::{Conclusion, RunSpec}; - /// Maximum length of a PR title (Unicode scalar values). const PR_TITLE_MAX_CHARS: usize = 72; @@ -675,8 +672,8 @@ mod tests { use fabro_llm::lithos_catalog::AdapterId; use fabro_llm::{Response, ResponseStream}; use fabro_types::{ - PetriAdmission, RunProjection, RunSpec, WorkflowSettings, first_event_seq, fixtures, - test_support, + PetriAdmission, RunProjection, RunSpec, StageSummary, WorkflowSettings, first_event_seq, + fixtures, test_support, }; use fabro_vault::{SecretType, Vault}; use httpmock::Method::{GET, POST}; @@ -685,7 +682,6 @@ mod tests { use tokio::sync::RwLock as AsyncRwLock; use super::*; - use crate::records::StageSummary; /// Answers every completion with one fixed text, attributed to the route /// that was asked. @@ -852,7 +848,7 @@ capabilities = { text = true, tools = true, response_format = { json_object = tr fn make_test_conclusion() -> Conclusion { Conclusion { timestamp: Utc::now(), - status: crate::outcome::StageOutcome::Succeeded, + status: fabro_types::StageOutcome::Succeeded, timing: fabro_types::RunTiming::wall_only(150_000), failure: None, final_git_commit_sha: None, diff --git a/lib/components/fabro-workflow/src/records/conclusion.rs b/lib/components/fabro-workflow/src/records/conclusion.rs deleted file mode 100644 index 6b98dcf0f..000000000 --- a/lib/components/fabro-workflow/src/records/conclusion.rs +++ /dev/null @@ -1 +0,0 @@ -pub use fabro_types::conclusion::{Conclusion, StageSummary}; diff --git a/lib/components/fabro-workflow/src/records/mod.rs b/lib/components/fabro-workflow/src/records/mod.rs deleted file mode 100644 index 94468f5f9..000000000 --- a/lib/components/fabro-workflow/src/records/mod.rs +++ /dev/null @@ -1,8 +0,0 @@ -mod conclusion; -mod run; -mod start; - -pub use conclusion::{Conclusion, StageSummary}; -pub use fabro_types::checkpoint::Checkpoint; -pub use run::RunSpec; -pub use start::StartRecord; diff --git a/lib/components/fabro-workflow/src/records/run.rs b/lib/components/fabro-workflow/src/records/run.rs deleted file mode 100644 index 6bbe14f21..000000000 --- a/lib/components/fabro-workflow/src/records/run.rs +++ /dev/null @@ -1 +0,0 @@ -pub use fabro_types::run::RunSpec; diff --git a/lib/components/fabro-workflow/src/records/start.rs b/lib/components/fabro-workflow/src/records/start.rs deleted file mode 100644 index 7968868bb..000000000 --- a/lib/components/fabro-workflow/src/records/start.rs +++ /dev/null @@ -1 +0,0 @@ -pub use fabro_types::start::StartRecord; diff --git a/lib/components/fabro-workflow/src/run_lookup.rs b/lib/components/fabro-workflow/src/run_lookup.rs index 8d402977a..a1f3c1b47 100644 --- a/lib/components/fabro-workflow/src/run_lookup.rs +++ b/lib/components/fabro-workflow/src/run_lookup.rs @@ -12,11 +12,10 @@ use chrono::{DateTime, Utc}; use fabro_config::Storage; use fabro_config::user::default_storage_dir; use fabro_store::Database; -use fabro_types::{Run, RunId}; +use fabro_types::{Run, RunId, RunStatus}; use serde::Serialize; use crate::operations::make_run_dir; -use crate::run_status::RunStatus; #[derive(Debug, Clone)] struct RunLocalState { @@ -448,11 +447,10 @@ mod tests { use std::sync::Arc; use fabro_store::RunSummaryStore; - use fabro_types::{RunProjection, RunStatus, fixtures, test_support}; + use fabro_types::{RunProjection, RunSpec, RunStatus, fixtures, test_support}; use super::scan_runs_combined; use crate::operations::make_run_dir; - use crate::records::RunSpec; fn sample_run_spec() -> RunSpec { RunSpec { diff --git a/lib/components/fabro-workflow/src/run_status.rs b/lib/components/fabro-workflow/src/run_status.rs deleted file mode 100644 index 9afd72509..000000000 --- a/lib/components/fabro-workflow/src/run_status.rs +++ /dev/null @@ -1 +0,0 @@ -pub use fabro_types::status::{FailureReason, RunStatus, SuccessReason, TerminalStatus}; diff --git a/lib/components/fabro-workflow/src/test_support.rs b/lib/components/fabro-workflow/src/test_support.rs deleted file mode 100644 index 016bb95a7..000000000 --- a/lib/components/fabro-workflow/src/test_support.rs +++ /dev/null @@ -1,28 +0,0 @@ -use fabro_types::ModelRef; -use lithos_llm::catalog::{ModelId, builtin}; -use lithos_llm::types::{Cost, CostSource, TokenCounts, Usage}; - -/// Construct a fully-populated `ModelUsage` for tests: `input_tokens` and -/// `output_tokens` on an OpenAI model, priced from the catalog at one micro -/// per token. Centralised so callers don't keep rebuilding the same skeleton. -#[must_use] -pub fn test_usage( - model_id: &str, - input_tokens: u64, - output_tokens: u64, -) -> fabro_types::ModelUsage { - fabro_types::ModelUsage::new( - ModelRef::new(builtin::openai(), ModelId::new(model_id)), - Usage { - tokens: TokenCounts { - input: input_tokens, - output: output_tokens, - ..TokenCounts::default() - }, - cost: Some(Cost { - usd_micros: input_tokens.saturating_add(output_tokens), - source: CostSource::Catalog, - }), - }, - ) -} diff --git a/lib/components/fabro-workflow/src/usage_rollup.rs b/lib/components/fabro-workflow/src/usage_rollup.rs deleted file mode 100644 index 63b19f479..000000000 --- a/lib/components/fabro-workflow/src/usage_rollup.rs +++ /dev/null @@ -1,270 +0,0 @@ -pub use fabro_types::usage_rollup::{ - ProjectionUsageByModel, ProjectionUsageRollup, ProjectionUsageStage, - usage_rollup_from_projection, -}; - -#[cfg(test)] -mod tests { - use fabro_types::{ - AttrValue, Graph, ModelRef, Node, RunProjection, RunSpec, StageCompletion, StageOutcome, - first_event_seq, test_support, - }; - use lithos_llm::catalog::{ModelId, builtin}; - use lithos_llm::types::{Cost, CostSource, TokenCounts, Usage}; - - use super::usage_rollup_from_projection; - use crate::test_support::test_usage; - - fn test_projection() -> RunProjection { - RunProjection::new( - "Test run".to_string(), - run_spec_with_boundary_nodes(), - chrono::Utc::now(), - ) - } - - #[test] - fn by_model_splits_a_completed_stage_by_its_usage_rows() { - let mut projection = test_projection(); - let root = test_usage("gpt-root", 100, 10); - let child = test_usage("gpt-child", 7, 1); - let stage = projection.stage_entry("work", 1, first_event_seq(1)); - stage.timing = Some(fabro_types::StageTiming::wall_only(100)); - stage.usage = root.usage.saturating_add(child.usage); - stage.model = Some(root.model().clone()); - stage.usage_by_model = vec![root.clone(), child.clone()]; - stage.completion = Some(StageCompletion { - outcome: StageOutcome::Succeeded, - notes: None, - failure_reason: None, - timestamp: chrono::Utc::now(), - }); - - let rollup = usage_rollup_from_projection(&projection); - - assert_eq!(rollup.totals.tokens.input, 107); - assert_eq!(rollup.stages[0].model.as_ref(), Some(root.model())); - assert_eq!(rollup.by_model.len(), 2, "{:?}", rollup.by_model); - let entry = |model_id: &str| { - rollup - .by_model - .iter() - .find(|entry| entry.model.model_id.as_str() == model_id) - .unwrap_or_else(|| panic!("a row for {model_id}")) - }; - assert_eq!(entry("gpt-root").stages, 1); - assert_eq!(entry("gpt-root").usage.tokens.input, 100); - assert_eq!(entry("gpt-root").usage.cost, root.usage.cost); - assert_eq!(entry("gpt-child").stages, 1); - assert_eq!(entry("gpt-child").usage.tokens.input, 7); - assert_eq!(entry("gpt-child").usage.cost, child.usage.cost); - } - - #[test] - fn rollup_groups_stage_rows_by_node_and_sums_retry_visit_usage() { - let mut projection = test_projection(); - let failed_usage = test_usage("gpt-old", 100, 10); - let success_usage = test_usage("gpt-new", 200, 20); - let first = projection.stage_entry("verify", 1, first_event_seq(1)); - first.timing = Some(fabro_types::StageTiming::wall_only(1200)); - first.usage = failed_usage.usage; - first.model = Some(failed_usage.model().clone()); - first.completion = Some(StageCompletion { - outcome: StageOutcome::Failed { - retry_requested: true, - }, - notes: None, - failure_reason: Some("try again".to_string()), - timestamp: chrono::Utc::now(), - }); - let second = projection.stage_entry("verify", 2, first_event_seq(2)); - second.timing = Some(fabro_types::StageTiming::wall_only(800)); - second.usage = success_usage.usage; - second.model = Some(success_usage.model().clone()); - second.completion = Some(StageCompletion { - outcome: StageOutcome::Succeeded, - notes: None, - failure_reason: None, - timestamp: chrono::Utc::now(), - }); - - let rollup = usage_rollup_from_projection(&projection); - - assert_eq!(rollup.stages.len(), 1); - assert_eq!(rollup.stages[0].node_id, "verify"); - assert_eq!( - rollup.stages[0] - .model - .as_ref() - .map(|model| model.model_id.as_str()), - Some("gpt-new") - ); - assert_eq!(rollup.stages[0].timing.wall_time_ms, 2000); - assert_eq!(rollup.stages[0].usage.tokens.input, 300); - assert_eq!(rollup.stages[0].usage.tokens.output, 30); - assert_eq!( - rollup.stages[0].usage.cost, - Some(Cost { - usd_micros: 330, - source: CostSource::Catalog, - }) - ); - - assert_eq!(rollup.timing.wall_time_ms, 2000); - assert_eq!(rollup.totals.tokens.input, 300); - assert_eq!(rollup.totals.tokens.output, 30); - assert_eq!(rollup.totals.cost.map(|cost| cost.usd_micros), Some(330)); - assert_eq!(rollup.usage_visit_count, 2); - - assert_eq!(rollup.by_model.len(), 2); - assert_eq!(rollup.by_model[0].model.model_id.as_str(), "gpt-new"); - assert_eq!(rollup.by_model[0].stages, 1); - assert_eq!(rollup.by_model[0].usage.tokens.input, 200); - assert_eq!(rollup.by_model[1].model.model_id.as_str(), "gpt-old"); - assert_eq!(rollup.by_model[1].stages, 1); - assert_eq!(rollup.by_model[1].usage.tokens.input, 100); - } - - #[test] - fn rollup_includes_completed_non_llm_stage_rows_with_zero_usage() { - let mut projection = test_projection(); - let stage = projection.stage_entry("build", 1, first_event_seq(1)); - stage.timing = Some(fabro_types::StageTiming::wall_only(25)); - stage.completion = Some(StageCompletion { - outcome: StageOutcome::Succeeded, - notes: None, - failure_reason: None, - timestamp: chrono::Utc::now(), - }); - - let rollup = usage_rollup_from_projection(&projection); - - assert_eq!(rollup.stages.len(), 1); - assert_eq!(rollup.stages[0].node_id, "build"); - assert_eq!(rollup.stages[0].timing.wall_time_ms, 25); - assert!(rollup.stages[0].model.is_none()); - assert_eq!(rollup.stages[0].usage, Usage::default()); - assert_eq!(rollup.timing.wall_time_ms, 25); - assert!(rollup.by_model.is_empty()); - assert!(rollup.usage_if_present().is_none()); - } - - #[test] - fn rollup_excludes_workflow_boundary_stage_rows() { - let mut projection = test_projection(); - projection.spec = run_spec_with_boundary_nodes(); - let start = projection.stage_entry("start", 1, first_event_seq(1)); - start.timing = Some(fabro_types::StageTiming::wall_only(25)); - start.completion = Some(StageCompletion { - outcome: StageOutcome::Succeeded, - notes: None, - failure_reason: None, - timestamp: chrono::Utc::now(), - }); - let exit = projection.stage_entry("exit", 1, first_event_seq(2)); - exit.timing = Some(fabro_types::StageTiming::wall_only(7)); - exit.completion = Some(StageCompletion { - outcome: StageOutcome::Succeeded, - notes: None, - failure_reason: None, - timestamp: chrono::Utc::now(), - }); - - let rollup = usage_rollup_from_projection(&projection); - - assert_eq!(rollup.stages.len(), 0); - assert_eq!(rollup.timing.wall_time_ms, 0); - } - - #[test] - fn rollup_keeps_in_flight_stage_usage_unpriced() { - let mut projection = test_projection(); - let model = ModelRef::new(builtin::openai(), ModelId::new("gpt-5.4")); - let stage = projection.stage_entry("agent", 1, first_event_seq(1)); - stage.started_at = Some(chrono::Utc::now()); - stage.usage = Usage::from(TokenCounts { - input: 500_000, - output: 125_000, - ..TokenCounts::default() - }); - stage.model = Some(model.clone()); - - let rollup = usage_rollup_from_projection(&projection); - - // The rollup keeps the shape of what the events recorded. Costs come - // from the events themselves; an in-flight stage that has recorded no - // cost yet stays unpriced rather than being re-estimated here. - assert_eq!(rollup.stages.len(), 1); - assert_eq!(rollup.stages[0].node_id, "agent"); - assert_eq!(rollup.stages[0].usage.cost, None); - assert_eq!(rollup.stages[0].usage.tokens.input, 500_000); - assert_eq!(rollup.totals.cost, None); - assert_eq!(rollup.by_model.len(), 1); - assert_eq!(rollup.by_model[0].usage.tokens.input, 500_000); - } - - #[test] - fn rollup_totals_lose_their_cost_once_an_unpriced_stage_used_tokens() { - let mut projection = test_projection(); - let priced = test_usage("gpt-priced", 100, 10); - let first = projection.stage_entry("plan", 1, first_event_seq(1)); - first.usage = priced.usage; - first.model = Some(priced.model().clone()); - first.completion = Some(StageCompletion { - outcome: StageOutcome::Succeeded, - notes: None, - failure_reason: None, - timestamp: chrono::Utc::now(), - }); - let second = projection.stage_entry("work", 1, first_event_seq(2)); - second.usage = Usage::from(TokenCounts { - input: 5, - ..TokenCounts::default() - }); - second.model = Some(ModelRef::new(builtin::openai(), ModelId::new("mystery"))); - second.completion = Some(StageCompletion { - outcome: StageOutcome::Succeeded, - notes: None, - failure_reason: None, - timestamp: chrono::Utc::now(), - }); - - let rollup = usage_rollup_from_projection(&projection); - - // A total cost is known only when every part is priced; the per-stage - // rows keep their own. - assert_eq!(rollup.totals.tokens.input, 105); - assert_eq!(rollup.totals.cost, None); - assert_eq!(rollup.stages[0].usage.cost, priced.usage.cost); - assert_eq!(rollup.stages[1].usage.cost, None); - assert_eq!( - rollup.usage_if_present().map(|usage| usage.cost), - Some(None) - ); - } - - fn run_spec_with_boundary_nodes() -> RunSpec { - let mut graph = Graph::new("test"); - graph.nodes.insert("start".to_string(), { - let mut node = Node::new("start"); - node.attrs.insert( - "shape".to_string(), - AttrValue::String("Mdiamond".to_string()), - ); - node - }); - graph.nodes.insert("exit".to_string(), { - let mut node = Node::new("exit"); - node.attrs.insert( - "shape".to_string(), - AttrValue::String("Msquare".to_string()), - ); - node - }); - - RunSpec { - graph, - ..test_support::test_run_spec() - } - } -} diff --git a/lib/foundation/fabro-types/src/lib.rs b/lib/foundation/fabro-types/src/lib.rs index a7f703d44..d2dfd262b 100644 --- a/lib/foundation/fabro-types/src/lib.rs +++ b/lib/foundation/fabro-types/src/lib.rs @@ -195,7 +195,7 @@ pub use transcript::{ MessageId, MessageKind, MessageSource, PairMessageRef, TranscriptMessage, text_of, tool_call_arguments, tool_result_from_json, tool_result_to_json, }; -pub use usage::{ModelRef, ModelUsage, sum_usage, usage_is_empty}; +pub use usage::{ModelRef, ModelUsage, format_cost, sum_usage, usage_is_empty}; pub use variable::{ CreateVariableRequest, UpdateVariableRequest, Variable, VariableListResponse, is_env_style_name, }; diff --git a/lib/foundation/fabro-types/src/test_support.rs b/lib/foundation/fabro-types/src/test_support.rs index 22f637a82..a73b5a579 100644 --- a/lib/foundation/fabro-types/src/test_support.rs +++ b/lib/foundation/fabro-types/src/test_support.rs @@ -1,10 +1,34 @@ use std::collections::HashMap; +use lithos_llm::catalog::{ModelId, builtin}; +use lithos_llm::types::{Cost, CostSource, TokenCounts, Usage}; + use crate::{ - AuthMethod, BlobHash, Graph, IdpIdentity, PetriAdmission, PetriGraphRef, Principal, - RunProvenance, RunSpec, WorkflowSettings, WorkflowVersionId, fixtures, + AuthMethod, BlobHash, Graph, IdpIdentity, ModelRef, ModelUsage, PetriAdmission, PetriGraphRef, + Principal, RunProvenance, RunSpec, WorkflowSettings, WorkflowVersionId, fixtures, }; +/// A fully populated `ModelUsage` for tests: `input_tokens` and +/// `output_tokens` on an OpenAI model, priced from the catalog at one micro +/// per token. +#[must_use] +pub fn test_usage(model_id: &str, input_tokens: u64, output_tokens: u64) -> ModelUsage { + ModelUsage::new( + ModelRef::new(builtin::openai(), ModelId::new(model_id)), + Usage { + tokens: TokenCounts { + input: input_tokens, + output: output_tokens, + ..TokenCounts::default() + }, + cost: Some(Cost { + usd_micros: input_tokens.saturating_add(output_tokens), + source: CostSource::Catalog, + }), + }, + ) +} + #[must_use] pub fn test_principal() -> Principal { Principal::user( diff --git a/lib/foundation/fabro-types/src/usage.rs b/lib/foundation/fabro-types/src/usage.rs index bb86cf521..d43f4478d 100644 --- a/lib/foundation/fabro-types/src/usage.rs +++ b/lib/foundation/fabro-types/src/usage.rs @@ -128,6 +128,12 @@ pub fn usage_is_empty(usage: &Usage) -> bool { *usage == Usage::default() } +/// Format a USD cost for display, to the cent. +#[must_use] +pub fn format_cost(cost: f64) -> String { + format!("${cost:.2}") +} + #[cfg(test)] mod tests { use lithos_llm::types::{Cost, CostSource, TokenCounts}; diff --git a/lib/foundation/fabro-types/src/usage_rollup.rs b/lib/foundation/fabro-types/src/usage_rollup.rs index 1615fe397..14773c941 100644 --- a/lib/foundation/fabro-types/src/usage_rollup.rs +++ b/lib/foundation/fabro-types/src/usage_rollup.rs @@ -196,3 +196,269 @@ fn stage_projection_order(state: &RunProjection) -> HashMap { } order } + +#[cfg(test)] +mod tests { + use lithos_llm::catalog::{ModelId, builtin}; + use lithos_llm::types::{Cost, CostSource, TokenCounts, Usage}; + + use super::usage_rollup_from_projection; + use crate::test_support::{self, test_usage}; + use crate::{ + AttrValue, Graph, ModelRef, Node, RunProjection, RunSpec, StageCompletion, StageOutcome, + first_event_seq, + }; + + fn test_projection() -> RunProjection { + RunProjection::new( + "Test run".to_string(), + run_spec_with_boundary_nodes(), + chrono::Utc::now(), + ) + } + + #[test] + fn by_model_splits_a_completed_stage_by_its_usage_rows() { + let mut projection = test_projection(); + let root = test_usage("gpt-root", 100, 10); + let child = test_usage("gpt-child", 7, 1); + let stage = projection.stage_entry("work", 1, first_event_seq(1)); + stage.timing = Some(crate::StageTiming::wall_only(100)); + stage.usage = root.usage.saturating_add(child.usage); + stage.model = Some(root.model().clone()); + stage.usage_by_model = vec![root.clone(), child.clone()]; + stage.completion = Some(StageCompletion { + outcome: StageOutcome::Succeeded, + notes: None, + failure_reason: None, + timestamp: chrono::Utc::now(), + }); + + let rollup = usage_rollup_from_projection(&projection); + + assert_eq!(rollup.totals.tokens.input, 107); + assert_eq!(rollup.stages[0].model.as_ref(), Some(root.model())); + assert_eq!(rollup.by_model.len(), 2, "{:?}", rollup.by_model); + let entry = |model_id: &str| { + rollup + .by_model + .iter() + .find(|entry| entry.model.model_id.as_str() == model_id) + .unwrap_or_else(|| panic!("a row for {model_id}")) + }; + assert_eq!(entry("gpt-root").stages, 1); + assert_eq!(entry("gpt-root").usage.tokens.input, 100); + assert_eq!(entry("gpt-root").usage.cost, root.usage.cost); + assert_eq!(entry("gpt-child").stages, 1); + assert_eq!(entry("gpt-child").usage.tokens.input, 7); + assert_eq!(entry("gpt-child").usage.cost, child.usage.cost); + } + + #[test] + fn rollup_groups_stage_rows_by_node_and_sums_retry_visit_usage() { + let mut projection = test_projection(); + let failed_usage = test_usage("gpt-old", 100, 10); + let success_usage = test_usage("gpt-new", 200, 20); + let first = projection.stage_entry("verify", 1, first_event_seq(1)); + first.timing = Some(crate::StageTiming::wall_only(1200)); + first.usage = failed_usage.usage; + first.model = Some(failed_usage.model().clone()); + first.completion = Some(StageCompletion { + outcome: StageOutcome::Failed { + retry_requested: true, + }, + notes: None, + failure_reason: Some("try again".to_string()), + timestamp: chrono::Utc::now(), + }); + let second = projection.stage_entry("verify", 2, first_event_seq(2)); + second.timing = Some(crate::StageTiming::wall_only(800)); + second.usage = success_usage.usage; + second.model = Some(success_usage.model().clone()); + second.completion = Some(StageCompletion { + outcome: StageOutcome::Succeeded, + notes: None, + failure_reason: None, + timestamp: chrono::Utc::now(), + }); + + let rollup = usage_rollup_from_projection(&projection); + + assert_eq!(rollup.stages.len(), 1); + assert_eq!(rollup.stages[0].node_id, "verify"); + assert_eq!( + rollup.stages[0] + .model + .as_ref() + .map(|model| model.model_id.as_str()), + Some("gpt-new") + ); + assert_eq!(rollup.stages[0].timing.wall_time_ms, 2000); + assert_eq!(rollup.stages[0].usage.tokens.input, 300); + assert_eq!(rollup.stages[0].usage.tokens.output, 30); + assert_eq!( + rollup.stages[0].usage.cost, + Some(Cost { + usd_micros: 330, + source: CostSource::Catalog, + }) + ); + + assert_eq!(rollup.timing.wall_time_ms, 2000); + assert_eq!(rollup.totals.tokens.input, 300); + assert_eq!(rollup.totals.tokens.output, 30); + assert_eq!(rollup.totals.cost.map(|cost| cost.usd_micros), Some(330)); + assert_eq!(rollup.usage_visit_count, 2); + + assert_eq!(rollup.by_model.len(), 2); + assert_eq!(rollup.by_model[0].model.model_id.as_str(), "gpt-new"); + assert_eq!(rollup.by_model[0].stages, 1); + assert_eq!(rollup.by_model[0].usage.tokens.input, 200); + assert_eq!(rollup.by_model[1].model.model_id.as_str(), "gpt-old"); + assert_eq!(rollup.by_model[1].stages, 1); + assert_eq!(rollup.by_model[1].usage.tokens.input, 100); + } + + #[test] + fn rollup_includes_completed_non_llm_stage_rows_with_zero_usage() { + let mut projection = test_projection(); + let stage = projection.stage_entry("build", 1, first_event_seq(1)); + stage.timing = Some(crate::StageTiming::wall_only(25)); + stage.completion = Some(StageCompletion { + outcome: StageOutcome::Succeeded, + notes: None, + failure_reason: None, + timestamp: chrono::Utc::now(), + }); + + let rollup = usage_rollup_from_projection(&projection); + + assert_eq!(rollup.stages.len(), 1); + assert_eq!(rollup.stages[0].node_id, "build"); + assert_eq!(rollup.stages[0].timing.wall_time_ms, 25); + assert!(rollup.stages[0].model.is_none()); + assert_eq!(rollup.stages[0].usage, Usage::default()); + assert_eq!(rollup.timing.wall_time_ms, 25); + assert!(rollup.by_model.is_empty()); + assert!(rollup.usage_if_present().is_none()); + } + + #[test] + fn rollup_excludes_workflow_boundary_stage_rows() { + let mut projection = test_projection(); + projection.spec = run_spec_with_boundary_nodes(); + let start = projection.stage_entry("start", 1, first_event_seq(1)); + start.timing = Some(crate::StageTiming::wall_only(25)); + start.completion = Some(StageCompletion { + outcome: StageOutcome::Succeeded, + notes: None, + failure_reason: None, + timestamp: chrono::Utc::now(), + }); + let exit = projection.stage_entry("exit", 1, first_event_seq(2)); + exit.timing = Some(crate::StageTiming::wall_only(7)); + exit.completion = Some(StageCompletion { + outcome: StageOutcome::Succeeded, + notes: None, + failure_reason: None, + timestamp: chrono::Utc::now(), + }); + + let rollup = usage_rollup_from_projection(&projection); + + assert_eq!(rollup.stages.len(), 0); + assert_eq!(rollup.timing.wall_time_ms, 0); + } + + #[test] + fn rollup_keeps_in_flight_stage_usage_unpriced() { + let mut projection = test_projection(); + let model = ModelRef::new(builtin::openai(), ModelId::new("gpt-5.4")); + let stage = projection.stage_entry("agent", 1, first_event_seq(1)); + stage.started_at = Some(chrono::Utc::now()); + stage.usage = Usage::from(TokenCounts { + input: 500_000, + output: 125_000, + ..TokenCounts::default() + }); + stage.model = Some(model.clone()); + + let rollup = usage_rollup_from_projection(&projection); + + // The rollup keeps the shape of what the events recorded. Costs come + // from the events themselves; an in-flight stage that has recorded no + // cost yet stays unpriced rather than being re-estimated here. + assert_eq!(rollup.stages.len(), 1); + assert_eq!(rollup.stages[0].node_id, "agent"); + assert_eq!(rollup.stages[0].usage.cost, None); + assert_eq!(rollup.stages[0].usage.tokens.input, 500_000); + assert_eq!(rollup.totals.cost, None); + assert_eq!(rollup.by_model.len(), 1); + assert_eq!(rollup.by_model[0].usage.tokens.input, 500_000); + } + + #[test] + fn rollup_totals_lose_their_cost_once_an_unpriced_stage_used_tokens() { + let mut projection = test_projection(); + let priced = test_usage("gpt-priced", 100, 10); + let first = projection.stage_entry("plan", 1, first_event_seq(1)); + first.usage = priced.usage; + first.model = Some(priced.model().clone()); + first.completion = Some(StageCompletion { + outcome: StageOutcome::Succeeded, + notes: None, + failure_reason: None, + timestamp: chrono::Utc::now(), + }); + let second = projection.stage_entry("work", 1, first_event_seq(2)); + second.usage = Usage::from(TokenCounts { + input: 5, + ..TokenCounts::default() + }); + second.model = Some(ModelRef::new(builtin::openai(), ModelId::new("mystery"))); + second.completion = Some(StageCompletion { + outcome: StageOutcome::Succeeded, + notes: None, + failure_reason: None, + timestamp: chrono::Utc::now(), + }); + + let rollup = usage_rollup_from_projection(&projection); + + // A total cost is known only when every part is priced; the per-stage + // rows keep their own. + assert_eq!(rollup.totals.tokens.input, 105); + assert_eq!(rollup.totals.cost, None); + assert_eq!(rollup.stages[0].usage.cost, priced.usage.cost); + assert_eq!(rollup.stages[1].usage.cost, None); + assert_eq!( + rollup.usage_if_present().map(|usage| usage.cost), + Some(None) + ); + } + + fn run_spec_with_boundary_nodes() -> RunSpec { + let mut graph = Graph::new("test"); + graph.nodes.insert("start".to_string(), { + let mut node = Node::new("start"); + node.attrs.insert( + "shape".to_string(), + AttrValue::String("Mdiamond".to_string()), + ); + node + }); + graph.nodes.insert("exit".to_string(), { + let mut node = Node::new("exit"); + node.attrs.insert( + "shape".to_string(), + AttrValue::String("Msquare".to_string()), + ); + node + }); + + RunSpec { + graph, + ..test_support::test_run_spec() + } + } +} diff --git a/lib/packages/fabro-api-client/src/api/run-internals-api.ts b/lib/packages/fabro-api-client/src/api/run-internals-api.ts index 09e391eaf..6ba34e729 100644 --- a/lib/packages/fabro-api-client/src/api/run-internals-api.ts +++ b/lib/packages/fabro-api-client/src/api/run-internals-api.ts @@ -50,8 +50,6 @@ import type { PetriReleaseRequest } from '../models'; // @ts-ignore import type { RunArtifactListResponse } from '../models'; // @ts-ignore -import type { RunCheckpoint } from '../models'; -// @ts-ignore import type { RunProjection } from '../models'; // @ts-ignore import type { StageContextWindow } from '../models'; @@ -929,46 +927,6 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config options: localVarRequestOptions, }; }, - /** - * Returns the latest checkpoint data for a run, or null if no checkpoint has been recorded yet. - * @summary Retrieve Run Checkpoint - * @param {string} id Unique run identifier (ULID). - * @param {*} [options] Override http request option. - * @throws {RequiredError} - */ - retrieveRunCheckpoint: async (id: string, options: RawAxiosRequestConfig = {}): Promise => { - // verify required parameter 'id' is not null or undefined - assertParamExists('retrieveRunCheckpoint', 'id', id) - const localVarPath = `/api/v1/runs/{id}/checkpoint` - .replace(`{${"id"}}`, encodeURIComponent(String(id))); - // use dummy base URL string because the URL constructor only accepts absolute URLs. - const localVarUrlObj = new URL(localVarPath, DUMMY_BASE_URL); - let baseOptions; - if (configuration) { - baseOptions = configuration.baseOptions; - } - - const localVarRequestOptions = { method: 'GET', ...baseOptions, ...options}; - const localVarHeaderParameter = {} as any; - const localVarQueryParameter = {} as any; - - // authentication SessionCookie required - - // authentication BearerAuth required - // http bearer authentication required - await setBearerAuthToObject(localVarHeaderParameter, configuration) - - localVarHeaderParameter['Accept'] = 'application/json'; - - setSearchParams(localVarUrlObj, localVarQueryParameter); - let headersFromBaseOptions = baseOptions && baseOptions.headers ? baseOptions.headers : {}; - localVarRequestOptions.headers = {...localVarHeaderParameter, ...headersFromBaseOptions, ...options.headers}; - - return { - url: toPathString(localVarUrlObj), - options: localVarRequestOptions, - }; - }, /** * Returns the persisted dense `WorkflowSettings` snapshot used to launch this run. * @summary Retrieve Run Settings @@ -1384,19 +1342,6 @@ export const RunInternalsApiFp = function(configuration?: Configuration) { const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.releasePetriRun']?.[localVarOperationServerIndex]?.url; return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); }, - /** - * Returns the latest checkpoint data for a run, or null if no checkpoint has been recorded yet. - * @summary Retrieve Run Checkpoint - * @param {string} id Unique run identifier (ULID). - * @param {*} [options] Override http request option. - * @throws {RequiredError} - */ - async retrieveRunCheckpoint(id: string, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { - const localVarAxiosArgs = await localVarAxiosParamCreator.retrieveRunCheckpoint(id, options); - const localVarOperationServerIndex = configuration?.serverIndex ?? 0; - const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.retrieveRunCheckpoint']?.[localVarOperationServerIndex]?.url; - return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); - }, /** * Returns the persisted dense `WorkflowSettings` snapshot used to launch this run. * @summary Retrieve Run Settings @@ -1660,16 +1605,6 @@ export const RunInternalsApiFactory = function (configuration?: Configuration, b releasePetriRun(id: string, petriReleaseRequest: PetriReleaseRequest, options?: RawAxiosRequestConfig): AxiosPromise { return localVarFp.releasePetriRun(id, petriReleaseRequest, options).then((request) => request(axios, basePath)); }, - /** - * Returns the latest checkpoint data for a run, or null if no checkpoint has been recorded yet. - * @summary Retrieve Run Checkpoint - * @param {string} id Unique run identifier (ULID). - * @param {*} [options] Override http request option. - * @throws {RequiredError} - */ - retrieveRunCheckpoint(id: string, options?: RawAxiosRequestConfig): AxiosPromise { - return localVarFp.retrieveRunCheckpoint(id, options).then((request) => request(axios, basePath)); - }, /** * Returns the persisted dense `WorkflowSettings` snapshot used to launch this run. * @summary Retrieve Run Settings @@ -1941,17 +1876,6 @@ export class RunInternalsApi extends BaseAPI { return RunInternalsApiFp(this.configuration).releasePetriRun(id, petriReleaseRequest, options).then((request) => request(this.axios, this.basePath)); } - /** - * Returns the latest checkpoint data for a run, or null if no checkpoint has been recorded yet. - * @summary Retrieve Run Checkpoint - * @param {string} id Unique run identifier (ULID). - * @param {*} [options] Override http request option. - * @throws {RequiredError} - */ - public retrieveRunCheckpoint(id: string, options?: RawAxiosRequestConfig) { - return RunInternalsApiFp(this.configuration).retrieveRunCheckpoint(id, options).then((request) => request(this.axios, this.basePath)); - } - /** * Returns the persisted dense `WorkflowSettings` snapshot used to launch this run. * @summary Retrieve Run Settings