mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-08 03:10:26 +00:00
Re-pin lithos-llm to 55add4596b861a0623d00c3a54aa5c147c8d504b and
pebble to c91810fe51aece80359b9cd8efea971af0c46925, where token usage
and cost travel together as Usage { tokens: TokenCounts, cost:
Option<Cost> }. Fabro now carries that one type everywhere it used to
carry BilledTokenCounts, BilledModelUsage, UsdMicros, or a token count
beside a cost_usd_micros.
fabro-types: billing.rs is usage.rs with ModelRef, ModelUsage { model,
usage }, sum_usage, and usage_is_empty; billing_rollup.rs is
usage_rollup.rs with ProjectionUsageStage, ProjectionUsageByModel,
ProjectionUsageRollup, and usage_rollup_from_projection. Every usage
field is named usage: StageProjection.usage and usage_by_model,
Outcome<Option<ModelUsage>>, stage.completed and stage.failed usage and
usage_by_model, prompt.completed usage, run.completed and run.failed
usage (total_usd_micros is gone), Conclusion.usage, StageSummary.usage,
Run.usage. RunSize buckets by Cost.
fabro-workflow: model_usage_from_llm prices tokens from the catalog with
a Catalog cost source, with_reported_cost keeps a provider cost, and the
pebble handler's stage_usage groups pebble's accounts by model and sums
rows with Usage::saturating_add, so a total has a cost only when every
priced part was priced. The store fold's live usage is the agent's
usage plus its descendants'.
API: the OpenAPI spec deletes BilledTokenCounts, BilledModelUsage,
CompletionUsage, CompletionCost, TokenUsage, and RunBillingSummary,
adds TokenCounts, Cost, Usage, and ModelUsage, and renames every
billing schema, property, tag, path, and operation to usage. fabro-api
reuses lithos-llm's and fabro-types' types through with_replacement,
with a round-trip test per replacement.
Old stored runs get no migration: their pebble events in the old shape
read back with zero usage, and their rebuilt projections lose agent
usage.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
176 lines
5.6 KiB
Rust
176 lines
5.6 KiB
Rust
use axum::body::Body;
|
|
use axum::http::{Request, StatusCode};
|
|
use tokio::time::sleep;
|
|
use tower::ServiceExt;
|
|
|
|
use crate::helpers::{
|
|
MINIMAL_DOT, POLL_ATTEMPTS, POLL_INTERVAL, api, create_and_start_run_from_intent,
|
|
minimal_intent_json, minimal_intent_json_with_dry_run, test_app_state_with_options,
|
|
test_app_with_scheduler, test_settings, wait_for_run_status,
|
|
};
|
|
|
|
const COMMAND_DOT: &str = r#"digraph Test {
|
|
graph [goal="Test"]
|
|
start [shape=Mdiamond]
|
|
echo_task [shape=parallelogram, script="echo command-stage"]
|
|
exit [shape=Msquare]
|
|
start -> echo_task -> exit
|
|
}"#;
|
|
|
|
const WAIT_DOT: &str = r#"digraph Test {
|
|
graph [goal="Test"]
|
|
start [shape=Mdiamond]
|
|
wait_task [shape=insulator, duration="1ms"]
|
|
exit [shape=Msquare]
|
|
start -> wait_task -> exit
|
|
}"#;
|
|
|
|
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
|
async fn aggregate_usage_increments_after_run_completes() {
|
|
let workspace = tempfile::tempdir().unwrap();
|
|
let state = test_app_state_with_options(test_settings(), 5);
|
|
let app = test_app_with_scheduler(state);
|
|
|
|
let run_id = create_and_start_run_from_intent(
|
|
&app,
|
|
minimal_intent_json_with_dry_run(&app, MINIMAL_DOT, workspace.path()).await,
|
|
)
|
|
.await;
|
|
|
|
// Poll until run completes
|
|
let status = wait_for_run_status(&app, &run_id, &["succeeded", "failed"]).await;
|
|
assert_eq!(status, "succeeded");
|
|
|
|
let mut total_runs = 0;
|
|
for _ in 0..POLL_ATTEMPTS {
|
|
let req = Request::builder()
|
|
.method("GET")
|
|
.uri(api("/usage"))
|
|
.body(Body::empty())
|
|
.unwrap();
|
|
|
|
let response = app.clone().oneshot(req).await.unwrap();
|
|
let body = crate::helpers::response_json(
|
|
response,
|
|
StatusCode::OK,
|
|
format!("{}:{}", file!(), line!()),
|
|
)
|
|
.await;
|
|
total_runs = body["totals"]["runs"].as_i64().unwrap();
|
|
if total_runs == 1 {
|
|
break;
|
|
}
|
|
sleep(POLL_INTERVAL).await;
|
|
}
|
|
assert_eq!(total_runs, 1);
|
|
}
|
|
|
|
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
|
async fn run_usage_includes_completed_non_llm_stages() {
|
|
let workspace = tempfile::tempdir().unwrap();
|
|
let state = test_app_state_with_options(test_settings(), 5);
|
|
let app = test_app_with_scheduler(state);
|
|
|
|
let run_id = create_and_start_run_from_intent(
|
|
&app,
|
|
minimal_intent_json(&app, WAIT_DOT, workspace.path()).await,
|
|
)
|
|
.await;
|
|
|
|
let status = wait_for_run_status(&app, &run_id, &["succeeded", "failed"]).await;
|
|
assert_eq!(status, "succeeded");
|
|
|
|
let usage = run_usage(&app, &run_id).await;
|
|
assert_non_llm_usage(&usage, &["wait_task"]);
|
|
}
|
|
|
|
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
|
async fn run_usage_includes_completed_command_stages() {
|
|
let workspace = tempfile::tempdir().unwrap();
|
|
let state = test_app_state_with_options(test_settings(), 5);
|
|
let app = test_app_with_scheduler(state);
|
|
|
|
let run_id = create_and_start_run_from_intent(
|
|
&app,
|
|
minimal_intent_json(&app, COMMAND_DOT, workspace.path()).await,
|
|
)
|
|
.await;
|
|
|
|
let status = wait_for_run_status(&app, &run_id, &["succeeded", "failed"]).await;
|
|
assert_eq!(status, "succeeded");
|
|
|
|
let usage = run_usage(&app, &run_id).await;
|
|
assert_non_llm_usage(&usage, &["echo_task"]);
|
|
}
|
|
|
|
async fn run_usage(app: &axum::Router, run_id: &str) -> serde_json::Value {
|
|
let req = Request::builder()
|
|
.method("GET")
|
|
.uri(api(&format!("/runs/{run_id}/usage")))
|
|
.body(Body::empty())
|
|
.expect("run usage request should build");
|
|
|
|
let response = app.clone().oneshot(req).await.unwrap();
|
|
crate::helpers::response_json(
|
|
response,
|
|
StatusCode::OK,
|
|
format!("GET /api/v1/runs/{run_id}/usage"),
|
|
)
|
|
.await
|
|
}
|
|
|
|
fn assert_non_llm_usage(usage: &serde_json::Value, expected_stage_ids: &[&str]) {
|
|
let stages = usage["stages"]
|
|
.as_array()
|
|
.expect("usage response should include stages");
|
|
let mut stage_ids = stages
|
|
.iter()
|
|
.map(|stage| {
|
|
stage["stage"]["id"]
|
|
.as_str()
|
|
.expect("stage should include an id")
|
|
.to_string()
|
|
})
|
|
.collect::<Vec<_>>();
|
|
stage_ids.sort();
|
|
assert_eq!(stage_ids, expected_stage_ids);
|
|
|
|
assert!(
|
|
stages.iter().all(|stage| {
|
|
stage["model"].is_null()
|
|
&& stage["usage"]["tokens"]["input"] == 0
|
|
&& stage["usage"]["tokens"]["output"] == 0
|
|
&& stage["usage"]["tokens"]["reasoning"] == 0
|
|
&& stage["usage"].get("cost").is_none()
|
|
}),
|
|
"every non-LLM stage should have null model and zero token counts: {stages:?}"
|
|
);
|
|
|
|
let stage_wall_sum: u64 = stages
|
|
.iter()
|
|
.map(|stage| {
|
|
stage["timing"]["wall_time_ms"]
|
|
.as_u64()
|
|
.expect("stage should include timing.wall_time_ms")
|
|
})
|
|
.sum();
|
|
|
|
assert_eq!(
|
|
usage["by_model"]
|
|
.as_array()
|
|
.expect("usage response should include by_model")
|
|
.len(),
|
|
0
|
|
);
|
|
assert_eq!(usage["totals"]["usage"]["tokens"]["input"], 0);
|
|
assert_eq!(usage["totals"]["usage"]["tokens"]["output"], 0);
|
|
assert!(usage["totals"]["usage"].get("cost").is_none());
|
|
|
|
let total_wall_time_ms = usage["totals"]["timing"]["wall_time_ms"]
|
|
.as_u64()
|
|
.expect("totals should include timing.wall_time_ms");
|
|
assert_eq!(
|
|
total_wall_time_ms, stage_wall_sum,
|
|
"total wall_time_ms {total_wall_time_ms} should equal summed stage wall_time_ms {stage_wall_sum}"
|
|
);
|
|
}
|