mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-07 03:00:29 +00:00
Compute LLM cost on-read for in-flight billing stages (#345)
## Summary
The run billing page showed `—` for dollar cost on any active
(in-flight) stage because `total_usd_micros` is only computed at stage
completion. This PR prices stages whose cost is `None` at read time,
using the model and token counts already present in the projection.
## What changed
**`fabro-model/src/billing.rs`** gains two new methods on existing
types:
- `BilledTokenCounts::token_counts()` — extracts the five disjoint token
buckets, dropping the derived sum and optional cost field.
- `Catalog::price_tokens(model, tokens)` — mirrors the cost computation
from `billed_model_usage_from_llm` but callable outside the completion
path. Returns `None` for unknown models or providers with no billing
policy.
**`fabro-workflow/src/billing_rollup.rs`** —
`billing_rollup_from_projection` gains an `Option<&Catalog>` parameter.
A new private helper `stage_usage_with_cost` fills in `total_usd_micros`
on-the-fly for any stage where it is `None` and both a catalog and model
are available. All four accumulator call sites (`is_zero` check,
per-stage row, `totals`, and `by_model`) use the priced copy, keeping
the page internally consistent.
**Call sites** — the read handler (`handler/billing.rs`) passes
`Some(&catalog)` so active stages get priced. The aggregate-billing
sites in `server.rs` and the four finalization sites in `finalize.rs`
pass `None` — completed stages are already priced, and we deliberately
exclude running estimates from org-wide totals to avoid double-counting
when a run later finalizes.
## Design decisions
- **Price any stage with `total_usd_micros == None`**, not just
explicitly in-flight ones. Stages with no billing policy return `None`
again — harmless.
- **No "estimated" label** — cost-so-far is exact for tokens consumed so
far, consistent with the already-unlabeled live token count and runtime.
- In-flight **prompt** stages still show `—` because `stage.model` isn't
set until `PromptCompleted`. This is acceptable; the reported bug
concerns agent stages.
### Fabro Details
<details>
<summary>Ran 9 stages in 44m 13s for $7.17</summary>
| Stage | Duration | Cost | Retries |
|---|---|---|---|
| start | 0s | – | 0 |
| toolchain | 1s | – | 0 |
| preflight_compile | 4m 6s | – | 0 |
| preflight_lint | 4m 58s | – | 0 |
| implement | 13m 24s | $3.84 | 0 |
| simplify_opus | 10m 34s | $1.69 | 0 |
| simplify_gpt | 5m 58s | $1.64 | 0 |
| verify | 3m 40s | – | 0 |
| fmt | 3s | – | 0 |
| **Total** | **44m 13s** | **$7.17** | **0** |
</details>
<details>
<summary>Ran <code>ImplementPlan.fabro</code> (12 nodes and 15
edges)</summary>
```dot
digraph ImplementPlan {
graph [
goal="Implement and simplify",
model_stylesheet="
* { model: claude-opus-4-7; }
"
]
rankdir=LR
start [shape=Mdiamond, label="Start"]
exit [shape=Msquare, label="Exit"]
toolchain [label="Toolchain", shape=parallelogram, script="command -v cargo >/dev/null || { curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y && sudo ln -sf $HOME/.cargo/bin/* /usr/local/bin/; }; cargo --version 2>&1", max_retries=0]
preflight_compile [label="Preflight Compile", shape=parallelogram, script="cargo check -q --workspace 2>&1", max_retries=0]
preflight_lint [label="Preflight Lint", shape=parallelogram, script="cargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1", max_retries=0]
fix_lints [label="Fix Lints", prompt="The preflight lint step failed. Read the build output from context and fix all clippy lint warnings.", max_visits=3]
implement [label="Implement", prompt="Read the plan file referenced in the goal and implement every step. Make all the code changes described in the plan. Use red/green TDD."]
simplify_opus [label="Simplify (Opus)", prompt="@prompts/simplify.md"]
simplify_gpt [label="Simplify (GPT-55)", prompt="@prompts/simplify.md", model="gpt-55"]
verify [label="Verify", shape=parallelogram, script="cargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1 && cargo nextest run --cargo-quiet --workspace --status-level fail 2>&1 && cargo dev docs refresh 2>&1 && cargo dev docs check 2>&1", goal_gate=true, retry_target="fixup"]
fixup [label="Fixup", prompt="The verify step failed. Read the build output from context and fix all clippy lint warnings, test failures, and generated docs errors.", max_visits=3]
fmt [label="Format", shape=parallelogram, script="cargo +nightly-2026-04-14 fmt --all 2>&1", max_retries=0]
start -> toolchain
toolchain -> preflight_compile [condition="outcome=succeeded"]
toolchain -> exit
preflight_compile -> preflight_lint [condition="outcome=succeeded"]
preflight_compile -> exit
preflight_lint -> implement [condition="outcome=succeeded"]
preflight_lint -> fix_lints
fix_lints -> preflight_lint
implement -> simplify_opus -> simplify_gpt -> verify
verify -> fmt [condition="outcome=succeeded"]
verify -> fixup
fixup -> verify
fmt -> exit
}
```
</details>
⚒️ Generated with [Fabro](https://fabro.sh)
---------
Co-authored-by: Fabro <noreply@fabro.sh>
This commit is contained in:
parent
06ee2fea39
commit
2b168b4588
5 changed files with 118 additions and 17 deletions
|
|
@ -358,6 +358,19 @@ impl BilledTokenCounts {
|
|||
}
|
||||
}
|
||||
|
||||
/// Returns the five disjoint per-call token buckets, dropping the derived
|
||||
/// `total_tokens` sum and the optional `total_usd_micros` cost.
|
||||
#[must_use]
|
||||
pub fn token_counts(&self) -> TokenCounts {
|
||||
TokenCounts {
|
||||
input_tokens: self.input_tokens,
|
||||
output_tokens: self.output_tokens,
|
||||
reasoning_tokens: self.reasoning_tokens,
|
||||
cache_read_tokens: self.cache_read_tokens,
|
||||
cache_write_tokens: self.cache_write_tokens,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn add_counts(&mut self, source: &Self) {
|
||||
self.input_tokens += source.input_tokens;
|
||||
self.output_tokens += source.output_tokens;
|
||||
|
|
@ -435,6 +448,27 @@ impl Catalog {
|
|||
self.provider(&model_ref.provider)
|
||||
.and_then(|provider| ModelBillingFacts::for_policy(provider.billing_policy, tokens))
|
||||
}
|
||||
|
||||
/// Price a partial token sample for `model` using catalog pricing.
|
||||
///
|
||||
/// Returns `None` when the provider has no billing policy, the model is
|
||||
/// unknown, or the pricing algorithm cannot produce a result for the given
|
||||
/// tokens. Used by read-side rollups so in-flight stages can show an
|
||||
/// exact cost for the tokens consumed so far.
|
||||
#[must_use]
|
||||
pub fn price_tokens(&self, model: &ModelRef, tokens: &TokenCounts) -> Option<i64> {
|
||||
let facts = self.billing_facts_for(model, tokens)?;
|
||||
let input = ModelBillingInput {
|
||||
usage: ModelUsage {
|
||||
model: model.clone(),
|
||||
tokens: tokens.clone(),
|
||||
},
|
||||
facts,
|
||||
};
|
||||
self.pricing_for(model)
|
||||
.and_then(|pricing| pricing.bill(&input))
|
||||
.map(|amount| amount.0)
|
||||
}
|
||||
}
|
||||
|
||||
fn costs_for_speed(
|
||||
|
|
|
|||
|
|
@ -3238,7 +3238,7 @@ async fn execute_run_in_process(state: Arc<AppState>, run_id: RunId) {
|
|||
.expect("aggregate_billing lock poisoned");
|
||||
accumulate_billing_rollup(
|
||||
&mut agg,
|
||||
&fabro_workflow::billing_rollup_from_projection(projection),
|
||||
&fabro_workflow::billing_rollup_from_projection(projection, None),
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
@ -3523,7 +3523,7 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
|
|||
.expect("aggregate_billing lock poisoned");
|
||||
accumulate_billing_rollup(
|
||||
&mut agg,
|
||||
&fabro_workflow::billing_rollup_from_projection(&final_state),
|
||||
&fabro_workflow::billing_rollup_from_projection(&final_state, None),
|
||||
);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -79,7 +79,8 @@ async fn get_run_billing(
|
|||
};
|
||||
let projection = cached.projection;
|
||||
|
||||
let rollup = fabro_workflow::billing_rollup_from_projection(&projection);
|
||||
let catalog = state.catalog();
|
||||
let rollup = fabro_workflow::billing_rollup_from_projection(&projection, Some(&catalog));
|
||||
let by_model = rollup
|
||||
.by_model
|
||||
.iter()
|
||||
|
|
|
|||
|
|
@ -1,6 +1,30 @@
|
|||
use std::borrow::Cow;
|
||||
use std::collections::HashMap;
|
||||
|
||||
use fabro_types::{BilledTokenCounts, ModelRef, RunProjection};
|
||||
use fabro_model::Catalog;
|
||||
use fabro_types::{BilledTokenCounts, ModelRef, RunProjection, StageProjection};
|
||||
|
||||
fn stage_usage_with_cost<'a>(
|
||||
catalog: Option<&Catalog>,
|
||||
stage: &'a StageProjection,
|
||||
) -> Cow<'a, BilledTokenCounts> {
|
||||
let Some(catalog) = catalog else {
|
||||
return Cow::Borrowed(&stage.usage);
|
||||
};
|
||||
let Some(model) = stage.model.as_ref() else {
|
||||
return Cow::Borrowed(&stage.usage);
|
||||
};
|
||||
if stage.usage.total_usd_micros.is_some() {
|
||||
return Cow::Borrowed(&stage.usage);
|
||||
}
|
||||
|
||||
let Some(total_usd_micros) = catalog.price_tokens(model, &stage.usage.token_counts()) else {
|
||||
return Cow::Borrowed(&stage.usage);
|
||||
};
|
||||
let mut usage = stage.usage.clone();
|
||||
usage.total_usd_micros = Some(total_usd_micros);
|
||||
Cow::Owned(usage)
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct ProjectionBillingStage {
|
||||
|
|
@ -34,7 +58,10 @@ impl ProjectionBillingRollup {
|
|||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn billing_rollup_from_projection(projection: &RunProjection) -> ProjectionBillingRollup {
|
||||
pub fn billing_rollup_from_projection(
|
||||
projection: &RunProjection,
|
||||
catalog: Option<&Catalog>,
|
||||
) -> ProjectionBillingRollup {
|
||||
let mut stage_indices = HashMap::<String, usize>::new();
|
||||
let mut stages = Vec::<ProjectionBillingStage>::new();
|
||||
let mut by_model = HashMap::<ModelRef, ProjectionBillingByModel>::new();
|
||||
|
|
@ -46,7 +73,9 @@ pub fn billing_rollup_from_projection(projection: &RunProjection) -> ProjectionB
|
|||
if is_boundary_stage(projection, stage_id.node_id()) {
|
||||
continue;
|
||||
}
|
||||
if stage.completion.is_none() && stage.duration_ms.is_none() && stage.usage.is_zero() {
|
||||
let usage = stage_usage_with_cost(catalog, stage);
|
||||
let usage = usage.as_ref();
|
||||
if stage.completion.is_none() && stage.duration_ms.is_none() && usage.is_zero() {
|
||||
continue;
|
||||
}
|
||||
|
||||
|
|
@ -68,10 +97,10 @@ pub fn billing_rollup_from_projection(projection: &RunProjection) -> ProjectionB
|
|||
runtime_ms = runtime_ms.saturating_add(duration_ms);
|
||||
}
|
||||
|
||||
if !stage.usage.is_zero() {
|
||||
if !usage.is_zero() {
|
||||
billed_visit_count += 1;
|
||||
row.billing.add_counts(&stage.usage);
|
||||
totals.add_counts(&stage.usage);
|
||||
row.billing.add_counts(usage);
|
||||
totals.add_counts(usage);
|
||||
|
||||
if let Some(model) = &stage.model {
|
||||
row.model = Some(model.clone());
|
||||
|
|
@ -84,7 +113,7 @@ pub fn billing_rollup_from_projection(projection: &RunProjection) -> ProjectionB
|
|||
billing: BilledTokenCounts::default(),
|
||||
});
|
||||
model_entry.stages += 1;
|
||||
model_entry.billing.add_counts(&stage.usage);
|
||||
model_entry.billing.add_counts(usage);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -126,6 +155,7 @@ fn is_boundary_stage(projection: &RunProjection, node_id: &str) -> bool {
|
|||
mod tests {
|
||||
use std::collections::HashMap;
|
||||
|
||||
use fabro_model::{Catalog, ModelRef, ProviderId};
|
||||
use fabro_types::{
|
||||
AttrValue, BilledModelUsage, BilledTokenCounts, Graph, Node, RunProjection, RunSpec,
|
||||
StageCompletion, StageOutcome, WorkflowSettings, first_event_seq, fixtures,
|
||||
|
|
@ -190,7 +220,7 @@ mod tests {
|
|||
timestamp: chrono::Utc::now(),
|
||||
});
|
||||
|
||||
let rollup = billing_rollup_from_projection(&projection);
|
||||
let rollup = billing_rollup_from_projection(&projection, None);
|
||||
|
||||
assert_eq!(rollup.stages.len(), 1);
|
||||
assert_eq!(rollup.stages[0].node_id, "verify");
|
||||
|
|
@ -233,7 +263,7 @@ mod tests {
|
|||
timestamp: chrono::Utc::now(),
|
||||
});
|
||||
|
||||
let rollup = billing_rollup_from_projection(&projection);
|
||||
let rollup = billing_rollup_from_projection(&projection, None);
|
||||
|
||||
assert_eq!(rollup.stages.len(), 1);
|
||||
assert_eq!(rollup.stages[0].node_id, "build");
|
||||
|
|
@ -266,12 +296,48 @@ mod tests {
|
|||
timestamp: chrono::Utc::now(),
|
||||
});
|
||||
|
||||
let rollup = billing_rollup_from_projection(&projection);
|
||||
let rollup = billing_rollup_from_projection(&projection, None);
|
||||
|
||||
assert_eq!(rollup.stages.len(), 0);
|
||||
assert_eq!(rollup.runtime_ms, 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rollup_prices_in_flight_stage_usage_using_catalog() {
|
||||
let mut projection = test_projection();
|
||||
let model = ModelRef {
|
||||
provider: ProviderId::openai(),
|
||||
model_id: "gpt-5.4".to_string(),
|
||||
speed: None,
|
||||
};
|
||||
let stage = projection.stage_entry("agent", 1, first_event_seq(1));
|
||||
stage.started_at = Some(chrono::Utc::now());
|
||||
stage.usage = BilledTokenCounts {
|
||||
input_tokens: 500_000,
|
||||
output_tokens: 125_000,
|
||||
total_tokens: 625_000,
|
||||
..BilledTokenCounts::default()
|
||||
};
|
||||
stage.model = Some(model.clone());
|
||||
|
||||
let priced = billing_rollup_from_projection(&projection, Some(Catalog::builtin()));
|
||||
let unpriced = billing_rollup_from_projection(&projection, None);
|
||||
|
||||
assert_eq!(priced.stages.len(), 1);
|
||||
assert_eq!(priced.stages[0].node_id, "agent");
|
||||
let stage_cost = priced.stages[0].billing.total_usd_micros;
|
||||
assert!(
|
||||
stage_cost.is_some_and(|cost| cost > 0),
|
||||
"expected priced stage cost, got {stage_cost:?}"
|
||||
);
|
||||
assert_eq!(priced.totals.total_usd_micros, stage_cost);
|
||||
assert_eq!(priced.by_model.len(), 1);
|
||||
assert_eq!(priced.by_model[0].billing.total_usd_micros, stage_cost);
|
||||
assert_eq!(unpriced.stages.len(), 1);
|
||||
assert_eq!(unpriced.stages[0].billing.total_usd_micros, None);
|
||||
assert_eq!(unpriced.totals.total_usd_micros, None);
|
||||
}
|
||||
|
||||
fn run_spec_with_boundary_nodes() -> RunSpec {
|
||||
let mut graph = Graph::new("test");
|
||||
graph.nodes.insert("start".to_string(), {
|
||||
|
|
|
|||
|
|
@ -82,7 +82,7 @@ pub(crate) async fn build_conclusion_from_store(
|
|||
.unwrap_or_default();
|
||||
let projection_billing = projection
|
||||
.as_ref()
|
||||
.map(billing_rollup_from_projection)
|
||||
.map(|projection| billing_rollup_from_projection(projection, None))
|
||||
.unwrap_or_default();
|
||||
let checkpoint = projection
|
||||
.as_ref()
|
||||
|
|
@ -438,7 +438,7 @@ async fn compute_final_patch(
|
|||
}
|
||||
|
||||
pub(crate) fn billing_from_projection(projection: &RunProjection) -> Option<BilledTokenCounts> {
|
||||
billing_rollup_from_projection(projection).billing_if_present()
|
||||
billing_rollup_from_projection(projection, None).billing_if_present()
|
||||
}
|
||||
|
||||
pub(crate) fn build_terminal_event(
|
||||
|
|
@ -554,7 +554,7 @@ pub async fn finalize(executed: Executed, options: &FinalizeOptions) -> Result<C
|
|||
.unwrap_or_default();
|
||||
let projection_billing = projection
|
||||
.as_ref()
|
||||
.map(billing_rollup_from_projection)
|
||||
.map(|projection| billing_rollup_from_projection(projection, None))
|
||||
.unwrap_or_default();
|
||||
let checkpoint = projection
|
||||
.as_ref()
|
||||
|
|
@ -979,7 +979,7 @@ mod tests {
|
|||
});
|
||||
|
||||
let projection_order = stage_projection_order(&projection);
|
||||
let projection_billing = billing_rollup_from_projection(&projection);
|
||||
let projection_billing = billing_rollup_from_projection(&projection, None);
|
||||
let mut latest_outcome = Outcome::success();
|
||||
latest_outcome.usage = Some(success_usage);
|
||||
latest_outcome.duration_ms = Some(800);
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue