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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-19 09:31:09 -04:00
parent 0752d4c7c4
commit 05a14a6c9b
No known key found for this signature in database
42 changed files with 340 additions and 737 deletions

View file

@ -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

View file

@ -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",

View file

@ -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;

View file

@ -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(())
}

View file

@ -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};

View file

@ -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;

View file

@ -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;

View file

@ -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::*;

View file

@ -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;

View file

@ -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;

View file

@ -295,9 +295,6 @@ async fn main_inner(worker_token: Option<String>) -> (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?;
}

View file

@ -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<std::borrow::Cow<'static, str>>) -
spinner
}
pub(crate) fn read_workflow_file(path: &Path) -> anyhow::Result<String> {
std::fs::read_to_string(path).with_context(|| format!("Failed to read {}", path.display()))
}
pub(crate) fn print_json_pretty<T>(value: &T) -> anyhow::Result<()>
where
T: Serialize + ?Sized,

View file

@ -26,7 +26,6 @@ mod model;
mod model_list;
mod model_test;
mod parent;
mod parse;
mod pr;
mod pr_close;
mod pr_create;

View file

@ -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] <WORKFLOW>
Arguments:
<WORKFLOW> 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 }
"#);
}

View file

@ -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<Arc<AppState>>,
Path(_id): Path<String>,
) -> Response {
(StatusCode::OK, Json(serde_json::json!(null))).into_response()
}
pub(crate) async fn cancel_stub(
_auth: RequiredUser,
State(_state): State<Arc<AppState>>,

View file

@ -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);

View file

@ -29,7 +29,6 @@ use super::super::{
pub(super) fn routes() -> Router<Arc<AppState>> {
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<u32>,
}
async fn get_checkpoint(
_auth: RequiredUser,
State(state): State<Arc<AppState>>,
Path(id): Path<String>,
) -> 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<Arc<AppState>>,

View file

@ -107,7 +107,6 @@ pub(super) fn demo_routes() -> Router<Arc<AppState>> {
"/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))

View file

@ -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;

View file

@ -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,

View file

@ -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()

View file

@ -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;

View file

@ -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::<RunId>().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"),

View file

@ -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;

View file

@ -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,

View file

@ -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<Option<ModelUsage>>;
/// Format a USD cost for display, to the cent.
#[must_use]
pub fn format_cost(cost: f64) -> String {
format!("${cost:.2}")
}

View file

@ -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 {

View file

@ -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.

View file

@ -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,

View file

@ -1 +0,0 @@
pub use fabro_types::conclusion::{Conclusion, StageSummary};

View file

@ -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;

View file

@ -1 +0,0 @@
pub use fabro_types::run::RunSpec;

View file

@ -1 +0,0 @@
pub use fabro_types::start::StartRecord;

View file

@ -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 {

View file

@ -1 +0,0 @@
pub use fabro_types::status::{FailureReason, RunStatus, SuccessReason, TerminalStatus};

View file

@ -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,
}),
},
)
}

View file

@ -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()
}
}
}

View file

@ -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,
};

View file

@ -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(

View file

@ -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};

View file

@ -196,3 +196,269 @@ fn stage_projection_order(state: &RunProjection) -> HashMap<String, u32> {
}
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()
}
}
}

View file

@ -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<RequestArgs> => {
// 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<RunCheckpoint>> {
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<void> {
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<RunCheckpoint> {
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