Merge pull request #862 from fabro-sh/one-usage-rule

One usage rule: bill an agent stage's session tree from one fold
This commit is contained in:
Bryan Helmkamp 2026-09-13 10:36:07 -04:00 • committed by GitHub
commit d5009d976b
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
37 changed files with 1306 additions and 274 deletions

View file

@ -97,6 +97,9 @@ function TokenBreakdown({ billing }: { billing: BilledTokenCounts }) {
</Fragment>
))}
</dl>
<p className="border-line text-fg-3 mt-1.5 border-t pt-1">
Includes subagent tokens, priced at each subagent&apos;s model.
</p>
</div>
);
}

View file

@ -422,6 +422,7 @@ Emitted when a workflow node finishes execution.
| `usage.reasoning_tokens` | number? | Reasoning/thinking tokens |
| `usage.speed` | string? | Speed tier |
| `usage.cost` | number? | Estimated cost in USD |
| `billing_by_model` | array? | For an agent stage, the stage's billing split by model: the root session's route and each subagent's own model, a subagent whose model the catalog does not know billed at the root's. Each row has `model`, `tokens`, and `total_usd_micros`, and the rows sum to the stage's billing. Empty for stages without a coding agent and on events written before it existed |
| `error` | string? | Error message (flattened from failure detail) |
| `failure_class` | string? | `"transient_infra"`, `"deterministic"`, `"budget_exhausted"`, `"compilation_loop"`, `"canceled"`, `"structural"` |
| `failure_signature` | string? | Dedup key for repeated failures |
@ -433,6 +434,12 @@ Emitted when a workflow node finishes execution.
| `restart_failure_signatures` | object? | Restart failure signature counts |
| `response` | string? | Full LLM or agent response text when produced by the stage |
| `notes` | string? | Free-text notes |
An agent stage's usage is its whole session tree's: the root session and
every subagent, live in `StageProjection.usage` and here at completion, both
read from the same fold of the stage's agent events. The root is priced at
its route and each subagent at its own model; where the provider reported a
cost, that cost stands.
| `files_touched` | string[] | File paths modified |
| `attempt` | number | Attempt number (1-based) |
| `max_attempts` | number | Maximum attempts allowed |
@ -466,6 +473,8 @@ Emitted when a stage fails (before retry decision).
| `failure_class` | string | Failure category |
| `failure_signature` | string? | Dedup key for repeated failures |
| `will_retry` | boolean | Whether the stage will be retried |
| `billing` | object? | What the stage spent before it failed, in the shape `stage.completed` uses. An agent stage that fails for good after answering model calls bills its whole session tree, as it would have on completion; a retried attempt and a cancelled stage carry none |
| `billing_by_model` | array? | `billing` split by model, as on `stage.completed` |
### `stage.retrying`

View file

@ -11227,6 +11227,18 @@ components:
agent_control:
$ref: "#/components/schemas/AgentControlState"
description: Whether the agent is executing normally or waiting for steering after an interrupt.
billing_by_model:
type: array
items:
$ref: "#/components/schemas/BilledModelUsage"
default: []
description: >-
The completed stage's `usage` split by model, as `stage.completed`
reported it: the root session's route and each subagent's own
model, a subagent whose model the catalog does not know billed at
the root's. Sums to `usage`. Empty while the stage runs and for
stages without a coding agent; the billing rollup then bills
`usage` to `model`.
agent:
oneOf:
- $ref: "#/components/schemas/AgentSessionProjection"
@ -13355,6 +13367,26 @@ components:
description: Billed USD amount in micros.
example: 720000
BilledModelUsage:
description: >-
Usage and cost billed to one model: one response, or one model's share
of a stage.
type: object
required:
- model
- tokens
properties:
model:
$ref: "#/components/schemas/BillingModelRef"
tokens:
$ref: "#/components/schemas/CompletionUsage"
total_usd_micros:
type: integer
format: int64
description: >-
Cost for `tokens`, when the provider reported one or the catalog
could price them. Absent means no cost data, not zero.
BillingModelRef:
description: Provider-qualified billing model identity used for cost estimates.
type: object

View file

@ -651,6 +651,7 @@ mod tests {
status: "succeeded".into(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,

View file

@ -650,6 +650,7 @@ mod tests {
status: "succeeded".into(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: Some(
billed_model_usage_from_llm(
&fabro_llm::test_support::test_catalog(),

View file

@ -6215,6 +6215,7 @@ fn stage_completed_event(node_id: &str) -> workflow_event::Event {
status: "succeeded".to_string(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -7063,6 +7064,7 @@ async fn list_run_stages_projects_retrying_until_completion() {
status: "succeeded".to_string(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -7102,14 +7104,15 @@ async fn list_run_stages_projects_retrying_until_completion() {
"work",
1,
&workflow_event::Event::StageFailed {
node_id: "work".to_string(),
name: "Work".to_string(),
index: 1,
failure: FailureDetail::new("try again", FailureCategory::TransientInfra),
will_retry: true,
timing: fabro_types::StageTiming::wall_only(10),
billing: None,
actor: None,
node_id: "work".to_string(),
name: "Work".to_string(),
index: 1,
failure: FailureDetail::new("try again", FailureCategory::TransientInfra),
will_retry: true,
timing: fabro_types::StageTiming::wall_only(10),
billing_by_model: Vec::new(),
billing: None,
actor: None,
},
)
.await;
@ -7157,6 +7160,7 @@ async fn list_run_stages_projects_retrying_until_completion() {
status: "partially_succeeded".to_string(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -7370,14 +7374,15 @@ async fn create_billed_retry_run(state: &Arc<AppState>, run_id: RunId) {
"verify",
1,
&workflow_event::Event::StageFailed {
node_id: "verify".to_string(),
name: "Verify".to_string(),
index: 1,
failure: FailureDetail::new("try again", FailureCategory::TransientInfra),
will_retry: true,
timing: fabro_types::StageTiming::wall_only(1200),
billing: Some(test_billed_usage("gpt-old", 100, 10)),
actor: None,
node_id: "verify".to_string(),
name: "Verify".to_string(),
index: 1,
failure: FailureDetail::new("try again", FailureCategory::TransientInfra),
will_retry: true,
timing: fabro_types::StageTiming::wall_only(1200),
billing_by_model: Vec::new(),
billing: Some(test_billed_usage("gpt-old", 100, 10)),
actor: None,
},
)
.await;
@ -7394,6 +7399,7 @@ async fn create_billed_retry_run(state: &Arc<AppState>, run_id: RunId) {
status: "succeeded".to_string(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: Some(test_billed_usage("gpt-new", 200, 20)),
failure: None,
notes: None,
@ -7482,6 +7488,7 @@ async fn list_run_stages_distinguishes_visits() {
status: "failed".to_string(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -7804,6 +7811,7 @@ async fn run_billing_dedups_retried_nodes_and_sums_their_durations() {
status: "failed".to_string(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -7835,6 +7843,7 @@ async fn run_billing_dedups_retried_nodes_and_sums_their_durations() {
status: "succeeded".to_string(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -8111,14 +8120,15 @@ async fn list_run_stages_shows_retrying_after_failed_event() {
"work",
1,
&workflow_event::Event::StageFailed {
node_id: "work".to_string(),
name: "Work".to_string(),
index: 0,
failure: FailureDetail::new("flake", FailureCategory::TransientInfra),
will_retry: true,
timing: fabro_types::StageTiming::wall_only(5),
billing: None,
actor: None,
node_id: "work".to_string(),
name: "Work".to_string(),
index: 0,
failure: FailureDetail::new("flake", FailureCategory::TransientInfra),
will_retry: true,
timing: fabro_types::StageTiming::wall_only(5),
billing_by_model: Vec::new(),
billing: None,
actor: None,
},
)
.await;
@ -8193,14 +8203,15 @@ async fn list_run_stages_shows_retrying_when_failed_will_retry() {
"work",
1,
&workflow_event::Event::StageFailed {
node_id: "work".to_string(),
name: "Work".to_string(),
index: 0,
failure: FailureDetail::new("flake", FailureCategory::TransientInfra),
will_retry: true,
timing: fabro_types::StageTiming::wall_only(5),
billing: None,
actor: None,
node_id: "work".to_string(),
name: "Work".to_string(),
index: 0,
failure: FailureDetail::new("flake", FailureCategory::TransientInfra),
will_retry: true,
timing: fabro_types::StageTiming::wall_only(5),
billing_by_model: Vec::new(),
billing: None,
actor: None,
},
)
.await;
@ -8243,14 +8254,15 @@ async fn run_billing_retried_node_then_succeeded_emits_one_row_with_final_attemp
max_attempts: 3,
},
workflow_event::Event::StageFailed {
node_id: "work".to_string(),
name: "Work".to_string(),
index: 0,
failure: FailureDetail::new("transient", FailureCategory::TransientInfra),
will_retry: true,
timing: fabro_types::StageTiming::wall_only(10),
billing: None,
actor: None,
node_id: "work".to_string(),
name: "Work".to_string(),
index: 0,
failure: FailureDetail::new("transient", FailureCategory::TransientInfra),
will_retry: true,
timing: fabro_types::StageTiming::wall_only(10),
billing_by_model: Vec::new(),
billing: None,
actor: None,
},
workflow_event::Event::StageRetrying {
node_id: "work".to_string(),
@ -8278,6 +8290,7 @@ async fn run_billing_retried_node_then_succeeded_emits_one_row_with_final_attemp
status: "succeeded".to_string(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -8349,6 +8362,7 @@ fn revisit_test_completed_with_visit(
status: "succeeded".to_string(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -15402,6 +15416,7 @@ async fn active_acp_steerable_marker_clears_on_terminal_paths() {
status: "success".to_string(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -15417,14 +15432,15 @@ async fn active_acp_steerable_marker_clears_on_terminal_paths() {
max_attempts: 1,
},
workflow_event::Event::StageFailed {
node_id: "agent".to_string(),
name: "agent".to_string(),
index: 0,
failure: FailureDetail::new("failed", FailureCategory::Deterministic),
will_retry: false,
timing: fabro_types::StageTiming::wall_only(1),
billing: None,
actor: None,
node_id: "agent".to_string(),
name: "agent".to_string(),
index: 0,
failure: FailureDetail::new("failed", FailureCategory::Deterministic),
will_retry: false,
timing: fabro_types::StageTiming::wall_only(1),
billing_by_model: Vec::new(),
billing: None,
actor: None,
},
];

View file

@ -25,7 +25,8 @@ use fabro_types::{
use fabro_util::error::render_compact_with_causes;
use lithos_llm::catalog::{ModelId, ProviderId};
use lithos_llm::types::TokenCounts;
use pebble_coding_agent::events::{CodingEvent, TokenUsage};
use pebble_coding_agent::events::CodingEvent;
use pebble_coding_agent::projection::SessionProjection;
use crate::{Error, EventEnvelope, Result};
@ -524,6 +525,7 @@ impl RunProjectionReducer for RunProjection {
stage.usage.replace_with_billed_usage(billing);
stage.model = Some(billing.model().clone());
}
stage.billing_by_model.clone_from(&props.billing_by_model);
stage.state = StageState::from(outcome.status);
stage.agent_control = AgentControlState::Running;
}
@ -547,6 +549,7 @@ impl RunProjectionReducer for RunProjection {
stage.usage.replace_with_billed_usage(billing);
stage.model = Some(billing.model().clone());
}
stage.billing_by_model.clone_from(&props.billing_by_model);
stage.state =
stage_state_from_failure(props.will_retry, failure_category, stage.termination);
stage.agent_control = AgentControlState::Running;
@ -760,9 +763,16 @@ fn apply_agent_event(
) {
let visit = props.visit;
// Pebble's own fold sees every agent event the stage stored, before the
// fabro-only arms below read the same event.
// fabro-only arms below read the same event. While the stage runs, its
// usage is that fold's: the tree's tokens, the root's and every
// subagent's, with whatever cost the provider reported. The terminal
// billing then brings the catalog's price for the same tokens.
if let Some(stage) = stage_at_stored_or_visit(state, stored, visit, seq) {
stage.agent.get_or_insert_default().apply(&props.event);
let agent = stage.agent.get_or_insert_default();
agent.apply(&props.event);
if stage.completion.is_none() {
stage.usage = live_usage(agent);
}
}
#[expect(
clippy::wildcard_enum_match_arm,
@ -771,17 +781,12 @@ fn apply_agent_event(
match props.coding_event() {
CodingEvent::AssistantMessage {
model,
usage,
cost_usd_micros,
context_window,
..
} => {
let Some(stage) = stage_at_stored_or_visit(state, stored, visit, seq) else {
return;
};
stage
.usage
.add_counts(&billed_counts(*usage, *cost_usd_micros));
if let Some(model) = stage_model_ref(stage, model) {
stage.model = Some(model);
}
@ -984,11 +989,17 @@ fn apply_agent_event(
}
}
/// Token accounting for one assistant message, in fabro's billing shape.
fn billed_counts(usage: TokenUsage, cost_usd_micros: Option<u64>) -> BilledTokenCounts {
/// A running stage's usage, from its agent's fold: the tree's tokens and the
/// cost the provider reported for them, `None` when it reported none.
fn live_usage(agent: &SessionProjection) -> BilledTokenCounts {
let (descendants, descendant_cost) = agent.descendant_usage();
let mut cost = agent.cost_usd_micros;
if let Some(descendant_cost) = descendant_cost {
cost = Some(cost.unwrap_or(0).saturating_add(descendant_cost));
}
BilledTokenCounts::from_token_counts(
TokenCounts::from(usage),
cost_usd_micros.map(|cost| i64::try_from(cost).unwrap_or(i64::MAX)),
TokenCounts::from(agent.usage.saturating_add(descendants)),
cost.map(|cost| i64::try_from(cost).unwrap_or(i64::MAX)),
)
}
@ -1779,6 +1790,7 @@ fn stage_outcome_from_props(props: &StageCompletedProps) -> Outcome<Option<Bille
notes: props.notes.clone(),
failure: props.failure.clone(),
usage: props.billing.clone(),
usage_by_model: props.billing_by_model.clone(),
files_touched: props.files_touched.clone(),
timing: Some(props.timing),
}
@ -1854,12 +1866,12 @@ mod tests {
SandboxProviderKind, StageHandler, StageModelUsage, StageOutcome, StageState, StageTiming,
SubAgentStatus, SuccessReason, WorkflowSettings, first_event_seq, fixtures, test_support,
};
use lithos_llm::types::{ReasoningEffort, Speed};
use lithos_llm::types::{ReasoningEffort, Speed, TokenCounts};
use pebble_coding_agent::events::{
CodingAgentEvent, CodingEvent, ContextWindowBreakdownItem, ContextWindowCategory,
ContextWindowCountMethod, ContextWindowSnapshot, ContextWindowStaleness,
ContextWindowWarning, ErrorData, ErrorKind, SkillActivationSource, SkillSummary,
TokenUsage, ToolCategory, ToolSource, ToolSummary,
CodingAgentEvent, CodingEvent, CompactionReason, ContextWindowBreakdownItem,
ContextWindowCategory, ContextWindowCountMethod, ContextWindowSnapshot,
ContextWindowStaleness, ContextWindowWarning, ErrorData, ErrorKind, SkillActivationSource,
SkillSummary, TokenUsage, ToolCategory, ToolSource, ToolSummary,
};
use pebble_coding_agent::tools::ToolOutputMetadata;
use serde_json::json;
@ -3657,6 +3669,7 @@ mod tests {
status: StageOutcome::Succeeded,
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: Some(usage.clone()),
failure: None,
notes: None,
@ -3744,6 +3757,7 @@ mod tests {
status: StageOutcome::Succeeded,
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: Some(usage),
failure: None,
notes: None,
@ -3788,6 +3802,7 @@ mod tests {
status: StageOutcome::Succeeded,
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: Some(usage.clone()),
failure: None,
notes: None,
@ -3827,14 +3842,15 @@ mod tests {
.apply_event(&test_stage_event(
3,
EventBody::StageFailed(StageFailedProps {
index: 0,
failure: Some(fabro_types::FailureDetail::new(
index: 0,
failure: Some(fabro_types::FailureDetail::new(
"try again",
fabro_types::FailureCategory::TransientInfra,
)),
will_retry: true,
timing: fabro_types::StageTiming::wall_only(444),
billing: Some(usage.clone()),
will_retry: true,
timing: fabro_types::StageTiming::wall_only(444),
billing_by_model: Vec::new(),
billing: Some(usage.clone()),
}),
scoped_stage_id.clone(),
))
@ -5453,6 +5469,7 @@ mod tests {
failure: Some(FailureDetail::new("boom", FailureCategory::TransientInfra)),
will_retry,
timing: fabro_types::StageTiming::wall_only(duration_ms),
billing_by_model: Vec::new(),
billing: None,
}
}
@ -5463,6 +5480,7 @@ mod tests {
failure: Some(FailureDetail::new("cancelled", FailureCategory::Canceled)),
will_retry,
timing: fabro_types::StageTiming::wall_only(duration_ms),
billing_by_model: Vec::new(),
billing: None,
}
}
@ -5503,6 +5521,7 @@ mod tests {
status,
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -5647,11 +5666,28 @@ mod tests {
assert_eq!(stage.model, Some(model));
}
fn child_message_body(input: u64, output: u64) -> EventBody {
EventBody::Agent(AgentEventProps::new(
"code",
1,
CodingAgentEvent::new(
"ses_child",
assistant_message(input, output),
SystemTime::UNIX_EPOCH,
)
.with_parent_session_id("ses_test"),
))
}
/// One usage rule: a stage's usage is its session tree's, live and at
/// completion. The terminal billing carries the tokens the fold already
/// showed plus the catalog's price, so completion changes the cost, not
/// the tokens, and keeps the split by model.
#[test]
fn stage_completed_replaces_live_usage_with_terminal_billing() {
fn stage_completed_keeps_the_trees_live_usage_and_prices_it() {
let mut state = initialized_projection();
let stage_id = StageId::new("build", 1);
let usage = billed_usage();
let model = billed_usage().model().clone();
state
.apply_event(&test_stage_event(
@ -5663,23 +5699,145 @@ mod tests {
state
.apply_event(&test_stage_event(
2,
activated(model.provider.as_str(), model.model_id.as_str()),
stage_id.clone(),
))
.unwrap();
state
.apply_event(&test_stage_event(
3,
agent_message_body(100, 50),
stage_id.clone(),
))
.unwrap();
let mut props = completed_props(42, StageOutcome::Succeeded);
props.billing = Some(usage.clone());
state
.apply_event(&test_stage_event(
3,
4,
child_message_body(7, 1),
stage_id.clone(),
))
.unwrap();
let live = state.stage(&stage_id).unwrap().usage.clone();
assert_eq!(
live,
live_counts(107, 51),
"the subagent's tokens are the stage's too"
);
let tree = BilledModelUsage {
model: model.clone(),
tokens: TokenCounts {
input: 107,
output: 51,
..TokenCounts::default()
},
total_usd_micros: Some(321),
};
let mut props = completed_props(42, StageOutcome::Succeeded);
props.billing = Some(tree.clone());
props.billing_by_model = vec![tree.clone()];
state
.apply_event(&test_stage_event(
5,
EventBody::StageCompleted(props),
stage_id.clone(),
))
.unwrap();
let stage = state.stage(&stage_id).unwrap();
assert_eq!(stage.usage, usage_counts(&usage));
assert_eq!(stage.model.as_ref(), Some(usage.model()));
assert_eq!(
stage.usage.token_counts(),
live.token_counts(),
"completion keeps the tokens the fold showed"
);
assert_eq!(
stage.usage.total_usd_micros,
Some(321),
"and brings the catalog's price"
);
assert_eq!(stage.model.as_ref(), Some(&model));
assert_eq!(stage.billing_by_model, vec![tree]);
}
#[test]
fn live_usage_is_the_trees_with_compactions_and_the_reported_cost() {
let mut state = initialized_projection();
let stage_id = StageId::new("build", 1);
let priced_message = |input: u64, output: u64, cost: u64| {
let CodingEvent::AssistantMessage {
text,
model,
usage,
cost_source,
tool_call_count,
context_window,
reasoning,
..
} = assistant_message(input, output)
else {
unreachable!("assistant_message builds an assistant message")
};
agent_body(CodingEvent::AssistantMessage {
text,
model,
usage,
cost_usd_micros: Some(cost),
cost_source,
tool_call_count,
context_window,
reasoning,
})
};
state
.apply_event(&test_stage_event(
1,
EventBody::StageStarted(started_props()),
stage_id.clone(),
))
.unwrap();
state
.apply_event(&test_stage_event(
2,
priced_message(10, 5, 5),
stage_id.clone(),
))
.unwrap();
state
.apply_event(&test_stage_event(
3,
child_message_body(7, 1),
stage_id.clone(),
))
.unwrap();
state
.apply_event(&test_stage_event(
4,
agent_body(CodingEvent::CompactionCompleted {
original_turn_count: 20,
preserved_turn_count: 6,
summary_token_estimate: 500,
tracked_file_count: 1,
reason: CompactionReason::Threshold,
usage: TokenUsage {
input: 30,
..TokenUsage::default()
},
cost_usd_micros: Some(2),
}),
stage_id.clone(),
))
.unwrap();
let stage = state.stage(&stage_id).unwrap();
assert_eq!(
stage.usage,
BilledTokenCounts {
total_usd_micros: Some(7),
..live_counts(47, 6)
},
"the root's messages and compaction, the child's message, and the provider's cost"
);
}
#[test]
@ -5934,14 +6092,15 @@ mod tests {
.apply_event(&test_event(
3,
EventBody::StageFailed(StageFailedProps {
index: 0,
failure: Some(FailureDetail::new(
index: 0,
failure: Some(FailureDetail::new(
"Script failed with exit code: 100\n\nCancelling due to test failure",
FailureCategory::Canceled,
)),
will_retry: false,
timing: fabro_types::StageTiming::wall_only(10),
billing: None,
will_retry: false,
timing: fabro_types::StageTiming::wall_only(10),
billing_by_model: Vec::new(),
billing: None,
}),
Some("build"),
))

View file

@ -22,6 +22,49 @@ mod tests {
)
}
#[test]
fn by_model_splits_a_completed_stage_by_its_billing_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 = BilledTokenCounts::from_billed_usage(&[root.clone(), child.clone()]);
stage.model = Some(root.model().clone());
stage.billing_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 = billing_rollup_from_projection(&projection);
assert_eq!(rollup.totals.input_tokens, 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").billing.input_tokens, 100);
assert_eq!(
entry("gpt-root").billing.total_usd_micros,
root.total_usd_micros
);
assert_eq!(entry("gpt-child").stages, 1);
assert_eq!(entry("gpt-child").billing.input_tokens, 7);
assert_eq!(
entry("gpt-child").billing.total_usd_micros,
child.total_usd_micros
);
}
#[test]
fn rollup_groups_stage_rows_by_node_and_sums_retry_visit_usage() {
let mut projection = test_projection();

View file

@ -2128,14 +2128,15 @@ mod tests {
// 3. Outcome → StageFailed event
let failure = outcome.failure.clone().unwrap();
let event = Event::StageFailed {
node_id: "code".into(),
name: "code".into(),
index: 0,
failure: failure.clone(),
will_retry: false,
timing: fabro_types::StageTiming::wall_only(0),
billing: None,
actor: None,
node_id: "code".into(),
name: "code".into(),
index: 0,
failure: failure.clone(),
will_retry: false,
timing: fabro_types::StageTiming::wall_only(0),
billing_by_model: Vec::new(),
billing: None,
actor: None,
};
// 4. Verify classification survived all the way through

View file

@ -346,6 +346,7 @@ fn event_body_from_event(event: &Event) -> EventBody {
preferred_label,
suggested_next_ids,
billing,
billing_by_model,
failure,
notes,
files_touched,
@ -366,6 +367,7 @@ fn event_body_from_event(event: &Event) -> EventBody {
preferred_label: preferred_label.clone(),
suggested_next_ids: suggested_next_ids.clone(),
billing: billing.clone(),
billing_by_model: billing_by_model.clone(),
failure: failure.clone(),
notes: notes.clone(),
files_touched: files_touched.clone(),
@ -385,6 +387,7 @@ fn event_body_from_event(event: &Event) -> EventBody {
will_retry,
timing,
billing,
billing_by_model,
..
} => EventBody::StageFailed(fabro_types::StageFailedProps {
index: *index,
@ -392,6 +395,7 @@ fn event_body_from_event(event: &Event) -> EventBody {
will_retry: *will_retry,
timing: *timing,
billing: billing.clone(),
billing_by_model: billing_by_model.clone(),
}),
Event::StageRetrying {
index,
@ -1097,6 +1101,7 @@ mod tests {
status: "succeeded".to_string(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -1142,6 +1147,7 @@ mod tests {
status: "succeeded".to_string(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -1167,17 +1173,18 @@ mod tests {
fn run_event_stage_failure_keeps_failure_detail() {
let usage = test_usage("gpt-5.2", 321, 54);
let stored = to_run_event(&fixtures::RUN_3, &Event::StageFailed {
node_id: "code".to_string(),
name: "Code".to_string(),
index: 1,
failure: FailureDetail::new(
node_id: "code".to_string(),
name: "Code".to_string(),
index: 1,
failure: FailureDetail::new(
"lint failed",
crate::outcome::FailureCategory::Deterministic,
),
will_retry: true,
timing: ::fabro_types::StageTiming::wall_only(5000),
billing: Some(usage.clone()),
actor: None,
will_retry: true,
timing: ::fabro_types::StageTiming::wall_only(5000),
billing_by_model: Vec::new(),
billing: Some(usage.clone()),
actor: None,
});
assert_eq!(stored.event_name(), "stage.failed");

View file

@ -272,6 +272,8 @@ pub enum Event {
preferred_label: Option<String>,
suggested_next_ids: Vec<String>,
billing: Option<BilledModelUsage>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
billing_by_model: Vec<BilledModelUsage>,
#[serde(default, skip_serializing_if = "Option::is_none")]
failure: Option<FailureDetail>,
notes: Option<String>,
@ -294,15 +296,17 @@ pub enum Event {
max_attempts: usize,
},
StageFailed {
node_id: String,
name: String,
index: usize,
failure: FailureDetail,
will_retry: bool,
timing: StageTiming,
billing: Option<BilledModelUsage>,
node_id: String,
name: String,
index: usize,
failure: FailureDetail,
will_retry: bool,
timing: StageTiming,
billing: Option<BilledModelUsage>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
billing_by_model: Vec<BilledModelUsage>,
#[serde(default, skip_serializing_if = "Option::is_none")]
actor: Option<Principal>,
actor: Option<Principal>,
},
StageRetrying {
node_id: String,

View file

@ -585,6 +585,7 @@ mod tests {
status: "succeeded".into(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,

View file

@ -31,7 +31,12 @@ const LAST_FILE_ROUTING_EXTENSIONS: &[&str] = &["json", "md"];
pub enum CodergenResult {
Text {
text: String,
/// The stage's billing: for an agent, the whole session tree's
/// tokens under the root's route.
usage: Option<BilledModelUsage>,
/// `usage` split by model, when the backend billed subagents at
/// their own models. Empty when `usage` is the one row.
usage_by_model: Vec<BilledModelUsage>,
files_touched: Vec<String>,
last_file_touched: Option<String>,
/// Active timing observed by the backend. The wall field is ignored by
@ -302,47 +307,62 @@ impl Handler for AgentHandler {
node_id: node.id.clone(),
}) as Arc<dyn ToolMiddleware>
});
let (response_text, stage_usage, backend_files_touched, last_file_touched, timing) =
if let Some(backend) = &self.backend {
let result = backend
.run(CodergenRunRequest {
node,
prompt: &prompt,
context,
thread_id: thread_id.as_deref(),
emitter: &services.run.emitter,
sandbox: &services.run.sandbox,
tool_middleware,
cancel_token: services.run.cancel_token(),
human_input: Some(human_input),
})
.await;
match result {
Ok(CodergenResult::Full(outcome)) => return Ok(*outcome),
Ok(CodergenResult::Text {
text,
usage,
files_touched,
last_file_touched,
timing,
}) => (text, usage, files_touched, last_file_touched, timing),
Err(Error::Cancelled) => return Err(Error::Cancelled),
Err(e) if e.is_retryable() => {
return Err(e);
}
Err(e) => {
return Ok(e.to_fail_outcome());
}
let (
response_text,
stage_usage,
stage_usage_by_model,
backend_files_touched,
last_file_touched,
timing,
) = if let Some(backend) = &self.backend {
let result = backend
.run(CodergenRunRequest {
node,
prompt: &prompt,
context,
thread_id: thread_id.as_deref(),
emitter: &services.run.emitter,
sandbox: &services.run.sandbox,
tool_middleware,
cancel_token: services.run.cancel_token(),
human_input: Some(human_input),
})
.await;
match result {
Ok(CodergenResult::Full(outcome)) => return Ok(*outcome),
Ok(CodergenResult::Text {
text,
usage,
usage_by_model,
files_touched,
last_file_touched,
timing,
}) => (
text,
usage,
usage_by_model,
files_touched,
last_file_touched,
timing,
),
Err(Error::Cancelled) => return Err(Error::Cancelled),
Err(e) if e.is_retryable() => {
return Err(e);
}
} else {
(
format!("[Simulated] Response for stage: {}", node.id),
None,
Vec::new(),
None,
StageTiming::default(),
)
};
Err(e) => {
return Ok(e.to_fail_outcome());
}
}
} else {
(
format!("[Simulated] Response for stage: {}", node.id),
None,
Vec::new(),
Vec::new(),
None,
StageTiming::default(),
)
};
let response_model = stage_usage
.as_ref()
@ -395,6 +415,7 @@ impl Handler for AgentHandler {
structured_output::exhausted_failure_outcome(node.output_retries());
failed.timing = Some(timing);
failed.usage = stage_usage;
failed.usage_by_model = stage_usage_by_model;
failed.files_touched = backend_files_touched;
return Ok(failed);
}
@ -422,6 +443,7 @@ impl Handler for AgentHandler {
}
}
outcome.usage = stage_usage;
outcome.usage_by_model = stage_usage_by_model;
outcome.files_touched = backend_files_touched;
outcome.timing = Some(timing);
@ -535,6 +557,7 @@ mod tests {
async fn run(&self, _request: CodergenRunRequest<'_>) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: "Done writing results.".to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: vec![self.path.clone()],
last_file_touched: Some(self.path.clone()),
@ -741,6 +764,7 @@ mod tests {
text:
r#"Done. {"outcome": "succeeded", "preferred_next_label": "approve"}"#
.to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -788,6 +812,7 @@ mod tests {
async fn run(&self, _request: CodergenRunRequest<'_>) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: "done".to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -945,6 +970,7 @@ All checks passed.
async fn run(&self, _request: CodergenRunRequest<'_>) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: r#"{"suggested_next_ids": [1]}"#.to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -1003,6 +1029,7 @@ All checks passed.
async fn run(&self, _request: CodergenRunRequest<'_>) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: r#"{"passed": true}"#.to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -1049,6 +1076,7 @@ All checks passed.
*self.captured_prompt.lock().unwrap() = Some(request.prompt.to_string());
Ok(CodergenResult::Text {
text: r#"{"passed": true}"#.to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -1133,6 +1161,7 @@ All checks passed.
);
Ok(CodergenResult::Text {
text: "done".to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -1188,6 +1217,7 @@ All checks passed.
Some(request.thread_id.map(String::from));
Ok(CodergenResult::Text {
text: "ok".to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -1233,6 +1263,7 @@ All checks passed.
Some(request.thread_id.map(String::from));
Ok(CodergenResult::Text {
text: "ok".to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -1441,6 +1472,7 @@ Some text in between.
*self.captured_prompt.lock().unwrap() = Some(request.prompt.to_string());
Ok(CodergenResult::Text {
text: "ok".to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -1502,6 +1534,7 @@ Some text in between.
*self.captured_prompt.lock().unwrap() = Some(request.prompt.to_string());
Ok(CodergenResult::Text {
text: "ok".to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,

View file

@ -208,6 +208,7 @@ mod tests {
assert!(request.prompt.contains("Synthesize every result"));
Ok(CodergenResult::Text {
text: "combined result".to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,

View file

@ -449,6 +449,7 @@ impl AgentAcpBackend {
Ok(CodergenResult::Text {
text: result.text,
usage_by_model: Vec::new(),
usage: None,
files_touched,
last_file_touched,

View file

@ -9,8 +9,8 @@
//! next route to continue it, and this module mirrors each move as the run's
//! `agent.failover` event.
use std::collections::{BTreeSet, HashMap, HashSet};
use std::sync::{Arc, Mutex};
use std::collections::{HashMap, HashSet};
use std::sync::{Arc, Mutex, PoisonError};
use std::time::{Duration, Instant};
use async_trait::async_trait;
@ -24,8 +24,8 @@ use fabro_mcp::pebble::pebble_servers;
use fabro_sandbox::{RunSandbox, SecretRedactor};
use fabro_types::settings::run::RunModelControls;
use fabro_types::{
AgentMcpToolSummary, AgentProfileKind, ModelRef, PermissionLevel, SessionCapability, StageId,
StageTiming, UsdMicros, billing,
AgentMcpToolSummary, AgentProfileKind, BilledModelUsage, ModelRef, PermissionLevel,
SessionCapability, StageId, StageTiming, UsdMicros, billing,
};
use fabro_util::home::Home;
use lithos_llm::catalog::{ModelId, ProviderId};
@ -34,6 +34,7 @@ use pebble_agent::ToolMiddleware;
use pebble_coding_agent::environment::Environment;
use pebble_coding_agent::events::{CodingAgentEvent, CodingEvent, EventSink, EventSinkError};
use pebble_coding_agent::extensions::HumanInputProvider;
use pebble_coding_agent::projection::{DescendantAccount, SessionProjection};
use pebble_coding_agent::state::Message;
use pebble_coding_agent::steering::SteerableSession;
use pebble_coding_agent::subagents::SubagentOptions;
@ -62,7 +63,7 @@ use crate::context::keys::Fidelity;
use crate::error::Error;
use crate::event::{Emitter, Event, StageScope};
use crate::model_fallback::{ModelFallbackNotice, ModelFallbackPolicy};
use crate::outcome::billed_model_usage_from_llm;
use crate::outcome::{Outcome, billed_model_usage_from_llm};
use crate::services::FabroRunToolServices;
use crate::steering_hub::SteeringHub;
use crate::web_search::{self, SearchSecrets};
@ -170,12 +171,27 @@ fn classify_agent_error(error: pebble_coding_agent::Error) -> AgentErrorDisposit
/// `agent.mcp.failed`, and `agent.mcp.disconnected` events, which the store
/// still folds; those mirrors go once every reader is on the projection.
struct WorkflowEventSink {
emitter: Arc<Emitter>,
node_id: String,
scope: StageScope,
emitter: Arc<Emitter>,
node_id: String,
scope: StageScope,
/// The stage's resolved plan, for the controls and origin the mirrored
/// failover event names.
plan: FallbackPlan,
plan: FallbackPlan,
/// Pebble's fold of every event this sink recorded: the stage's one
/// account of what its agent and subagents spent, wrote, and ran. The
/// store folds the same events the same way, so the stage's billing at
/// its end is the usage the run showed live.
projection: Mutex<SessionProjection>,
}
impl WorkflowEventSink {
/// The account as it stands.
fn snapshot(&self) -> SessionProjection {
self.projection
.lock()
.unwrap_or_else(PoisonError::into_inner)
.clone()
}
}
#[async_trait]
@ -272,6 +288,10 @@ impl EventSink for WorkflowEventSink {
if event.event.is_streaming_noise() {
return Ok(());
}
self.projection
.lock()
.unwrap_or_else(PoisonError::into_inner)
.apply(event);
self.emitter
.emit_durable(
&Event::Agent {
@ -291,57 +311,47 @@ impl EventSink for WorkflowEventSink {
// --- Live invocation ------------------------------------------------------
/// One stage invocation's live agent and its accounting.
/// One stage invocation's live agent, its timing, and the sink that
/// accounts for it.
///
/// A stage may run several prompts on one agent (the prompt, output repairs,
/// late steering); the usage, cost, timing, and files of every one of them
/// are summed here, across whatever routes pebble moved through.
/// late steering). What every one of them spent and wrote, subagents
/// included and across whatever routes pebble moved through, is the sink's
/// fold of the events it recorded; the prompt reports here contribute their
/// timing and the route the prompt ended on.
struct LiveAgent {
agent: CodingAgent,
handle: CodingAgentControlHandle,
lease: Option<Arc<ActivationLease>>,
total_usage: TokenCounts,
total_cost: Option<UsdMicros>,
sink: Arc<WorkflowEventSink>,
inference_duration: Duration,
tool_duration: Duration,
/// Every file the stage's prompts wrote or edited, subagents included.
files_touched: BTreeSet<String>,
/// The most recently written path.
last_file_touched: Option<String>,
}
impl LiveAgent {
fn new(agent: CodingAgent, handle: CodingAgentControlHandle) -> Self {
fn new(
agent: CodingAgent,
handle: CodingAgentControlHandle,
sink: Arc<WorkflowEventSink>,
) -> Self {
Self {
agent,
handle,
lease: None,
total_usage: TokenCounts::default(),
total_cost: None,
sink,
inference_duration: Duration::ZERO,
tool_duration: Duration::ZERO,
files_touched: BTreeSet::new(),
last_file_touched: None,
}
}
fn record_report(&mut self, report: &pebble_coding_agent::PromptReport) {
billing::add_usage(&mut self.total_usage, TokenCounts::from(report.usage));
UsdMicros::accumulate(
&mut self.total_cost,
report
.cost_usd_micros
.map(|micros| UsdMicros(i64::try_from(micros).unwrap_or(i64::MAX))),
);
self.inference_duration = self
.inference_duration
.saturating_add(report.timing.inference);
self.tool_duration = self.tool_duration.saturating_add(report.timing.tool);
self.files_touched
.extend(report.files_touched.iter().cloned());
for compaction in &report.compactions {
// The summary call's usage is already in `report.usage`; this is
// the breakdown, for anyone asking why a stage cost what it did.
// The summary call's usage is already in the stage's account; this
// is the breakdown, for anyone asking why a stage cost what it did.
tracing::debug!(
reason = ?compaction.reason,
original_turns = compaction.original_turn_count,
@ -351,9 +361,21 @@ impl LiveAgent {
"agent stage compacted its conversation"
);
}
if report.last_file_touched.is_some() {
self.last_file_touched.clone_from(&report.last_file_touched);
}
}
/// What the stage's prompts have spent and written so far.
fn account(&self) -> SessionProjection {
self.sink.snapshot()
}
/// The path written or edited most recently, when any was.
fn last_file_touched(&self) -> Option<String> {
self.sink
.projection
.lock()
.unwrap_or_else(PoisonError::into_inner)
.last_file_touched
.clone()
}
fn release_lease(&mut self) {
@ -385,6 +407,111 @@ impl LiveAgent {
}
}
/// The route as billing names it: provider, model, and the speed tier the
/// stage asked for.
fn route_model(route: &LlmRoute) -> ModelRef {
ModelRef::new(
route.target.provider.clone(),
ModelId::new(route.target.model.as_str()),
)
.with_speed(route.controls.speed)
}
/// A stage's billing from its account: the whole tree under the root's
/// route, and the rows that split it by model.
struct StageBilling {
total: BilledModelUsage,
by_model: Vec<BilledModelUsage>,
}
/// Bills the stage's account from the catalog: the root session at
/// `root_model`, its route, and each descendant at its own route where the
/// catalog knows it and at the root's otherwise, so a subagent on a cheaper
/// or dearer model is priced as what it ran. A descendant on the root's
/// route joins the root's row. Where pebble carried a provider-reported
/// cost, that cost stands in for the catalog's estimate.
fn stage_billing(
catalog: &Catalog,
root_model: &ModelRef,
account: &SessionProjection,
) -> Result<StageBilling, Error> {
let mut groups: Vec<(ModelRef, TokenCounts, Option<u64>)> = vec![(
root_model.clone(),
TokenCounts::from(account.usage),
account.cost_usd_micros,
)];
for descendant in account.descendants.values() {
let model = descendant_model(catalog, root_model, descendant);
match groups.iter_mut().find(|(grouped, _, _)| *grouped == model) {
Some((_, tokens, cost)) => {
billing::add_usage(tokens, TokenCounts::from(descendant.usage));
add_reported_cost(cost, descendant.cost_usd_micros);
}
None => groups.push((
model,
TokenCounts::from(descendant.usage),
descendant.cost_usd_micros,
)),
}
}
// The root's row first, then the others by model.
groups[1..].sort_by(|left, right| left.0.sort_key().cmp(&right.0.sort_key()));
let mut by_model = Vec::with_capacity(groups.len());
let mut total_tokens = TokenCounts::default();
let mut total_cost = None;
for (model, tokens, reported) in groups {
let row = billed_model_usage_from_llm(catalog, &model, tokens)?
.with_reported_cost(reported.map(usd_micros));
billing::add_usage(&mut total_tokens, row.tokens);
UsdMicros::accumulate(&mut total_cost, row.total_usd_micros.map(UsdMicros));
by_model.push(row);
}
Ok(StageBilling {
total: BilledModelUsage {
model: root_model.clone(),
tokens: total_tokens,
total_usd_micros: total_cost.map(|cost| cost.0),
},
by_model,
})
}
/// The route a descendant is billed at: its own where its start named one
/// the catalog knows, else the root's. A descendant whose start was not seen
/// names only its answers' model, taken to be on the root's provider.
fn descendant_model(
catalog: &Catalog,
root_model: &ModelRef,
account: &DescendantAccount,
) -> ModelRef {
let Some(model) = account.model.as_deref() else {
return root_model.clone();
};
let provider = account
.provider
.as_deref()
.unwrap_or(root_model.provider.as_str());
if provider == root_model.provider.as_str() && model == root_model.model_id.as_str() {
return root_model.clone();
}
if catalog.enabled_provider(provider).is_none() {
return root_model.clone();
}
ModelRef::new(ProviderId::new(provider), ModelId::new(model))
}
/// Folds a reported cost into a total that stays `None` until one is seen.
fn add_reported_cost(total: &mut Option<u64>, cost: Option<u64>) {
if let Some(cost) = cost {
*total = Some(total.unwrap_or(0).saturating_add(cost));
}
}
fn usd_micros(micros: u64) -> UsdMicros {
UsdMicros(i64::try_from(micros).unwrap_or(i64::MAX))
}
/// Everything one stage binds to an agent it builds or resumes.
struct StageBindings<'a> {
node_id: &'a str,
@ -587,21 +714,24 @@ impl PebbleBackend {
plan: &FallbackPlan,
provider: &ProviderContext,
bindings: &StageBindings<'_>,
) -> CodingAgentBuilder {
) -> (CodingAgentBuilder, Arc<WorkflowEventSink>) {
let route = plan.current();
let max_tokens = node_max_output_tokens(node).map(i64::from);
let sink = Arc::new(WorkflowEventSink {
emitter: Arc::clone(bindings.emitter),
node_id: bindings.node_id.to_string(),
scope: bindings.stage_scope.clone(),
plan: plan.clone(),
projection: Mutex::new(SessionProjection::new()),
});
let event_sink = Arc::clone(&sink) as Arc<dyn EventSink>;
builder = builder
.tools(self.stage_tools())
.mcp_servers(pebble_servers(&self.mcp_servers))
.permission_level(PermissionLevel::Full)
.options(self.agent_options(node, route.controls))
.fallback_routes(plan.pebble_routes(max_tokens))
.event_sink(Arc::new(WorkflowEventSink {
emitter: Arc::clone(bindings.emitter),
node_id: bindings.node_id.to_string(),
scope: bindings.stage_scope.clone(),
plan: plan.clone(),
}))
.event_sink(event_sink)
.redactor(Arc::new(SecretRedactor))
.subagents(SubagentOptions::enabled());
if let Some(routes) = bindings.sandbox.port_routes() {
@ -622,7 +752,7 @@ impl PebbleBackend {
if provider.profile_kind == AgentProfileKind::Claude5 {
builder = builder.web_fetch_summarizer(route.selector());
}
builder
(builder, sink)
}
/// A new agent on the plan's current route.
@ -632,15 +762,17 @@ impl PebbleBackend {
plan: &FallbackPlan,
provider: &ProviderContext,
bindings: &StageBindings<'_>,
) -> Result<CodingAgent, Error> {
) -> Result<(CodingAgent, Arc<WorkflowEventSink>), Error> {
let client = self.build_llm_client().await?;
let environment: Arc<dyn Environment> =
Arc::clone(bindings.sandbox) as Arc<dyn Environment>;
let builder = CodingAgent::builder(client, environment).model(plan.current().selector());
self.bind_builder(builder, node, plan, provider, bindings)
let (builder, sink) = self.bind_builder(builder, node, plan, provider, bindings);
let agent = builder
.build()
.await
.map_err(|error| Error::handler_with_source("Failed to start agent session", error))
.map_err(|error| Error::handler_with_source("Failed to start agent session", error))?;
Ok((agent, sink))
}
/// The exported conversation of an earlier stage, continued on the
@ -652,15 +784,17 @@ impl PebbleBackend {
plan: &FallbackPlan,
provider: &ProviderContext,
bindings: &StageBindings<'_>,
) -> Result<CodingAgent, Error> {
) -> Result<(CodingAgent, Arc<WorkflowEventSink>), Error> {
let client = self.build_llm_client().await?;
let environment: Arc<dyn Environment> =
Arc::clone(bindings.sandbox) as Arc<dyn Environment>;
let builder = CodingAgent::resume_from_export(client, environment, export);
self.bind_builder(builder, node, plan, provider, bindings)
let (builder, sink) = self.bind_builder(builder, node, plan, provider, bindings);
let agent = builder
.build()
.await
.map_err(|error| Error::handler_with_source("Failed to resume agent session", error))
.map_err(|error| Error::handler_with_source("Failed to resume agent session", error))?;
Ok((agent, sink))
}
/// Register `live` with the steering hub so steers reach it, and tell
@ -732,6 +866,37 @@ impl PebbleBackend {
}
}
/// The failed outcome of an agent stage that spent before it failed: the
/// failure itself, with the session tree's usage, the files it wrote, and
/// its active time, so the run bills what the stage spent. A billing the
/// catalog cannot price is logged and left off.
fn failed_outcome(&self, error: &Error, live: &LiveAgent, plan: &FallbackPlan) -> Outcome {
let mut outcome = error.to_fail_outcome();
let account = live.account();
match stage_billing(
self.catalog.as_ref(),
&route_model(plan.current()),
&account,
) {
Ok(billing) => {
outcome.usage = Some(billing.total);
outcome.usage_by_model = billing.by_model;
}
Err(billing_error) => {
tracing::debug!(
error = %billing_error,
"failed agent stage could not be billed"
);
}
}
outcome.files_touched = account.files_touched;
outcome.timing = Some(StageTiming::active_only(
crate::millis_u64(live.inference_duration),
crate::millis_u64(live.tool_duration),
));
outcome
}
/// Steers that landed between the answer and the hub's close-the-door
/// check run as further prompts, so the stage never ends with a steer
/// nobody saw.
@ -986,6 +1151,7 @@ impl CodergenBackend for PebbleBackend {
return Ok(CodergenResult::Text {
text: response_text,
usage_by_model: Vec::new(),
usage: Some(stage_usage),
files_touched: Vec::new(),
last_file_touched: None,
@ -1026,13 +1192,13 @@ impl CodergenBackend for PebbleBackend {
let cached = reuse_key.as_ref().and_then(|key| self.take_thread(key));
let is_reused = cached.is_some();
let (agent, mut fallback_plan) = if let Some(thread) = cached {
let ((agent, sink), mut fallback_plan) = if let Some(thread) = cached {
let route = thread.fallback_plan.current().clone();
let provider = self.resolve_provider_context(
route.target.model.as_str(),
Some(route.target.provider.as_str()),
)?;
let agent = self
let session = self
.resume_exported_agent(
thread.export,
node,
@ -1041,7 +1207,7 @@ impl CodergenBackend for PebbleBackend {
&bindings,
)
.await?;
(agent, thread.fallback_plan)
(session, thread.fallback_plan)
} else {
let model = node.model().unwrap_or(&self.model);
let provider = routing::resolve_node_provider_context(
@ -1059,10 +1225,10 @@ impl CodergenBackend for PebbleBackend {
route.target.model.as_str(),
Some(route.target.provider.as_str()),
)?;
let agent = self
let session = self
.build_agent(node, &fallback_plan, &route_provider, &bindings)
.await?;
(agent, fallback_plan)
(session, fallback_plan)
};
if cancel_token.is_cancelled() {
let mut agent = agent;
@ -1078,7 +1244,7 @@ impl CodergenBackend for PebbleBackend {
);
let handle = agent.control_handle();
let mut live = LiveAgent::new(agent, handle);
let mut live = LiveAgent::new(agent, handle, sink);
let route = fallback_plan.current().clone();
if let Err(error) =
self.activate(&mut live, &route, &stage_id, request.thread_id, &bindings)
@ -1104,7 +1270,7 @@ impl CodergenBackend for PebbleBackend {
let mut repair_attempts = 0_i64;
let mut previous_validation_error = None;
loop {
let last_file_touched = live.last_file_touched.clone();
let last_file_touched = live.last_file_touched();
match validate_agent_output_sources(
schema,
&response,
@ -1166,21 +1332,27 @@ impl CodergenBackend for PebbleBackend {
ShutdownReason::Error
};
live.discard(reason).await;
return Err(error);
// Cancellation and a retryable failure go up as the error, so
// the engine cancels or retries as before. A terminal failure
// becomes the stage's failed outcome, carrying what the
// session tree spent and wrote before it failed.
if matches!(error, Error::Cancelled) || error.is_retryable() {
return Err(error);
}
return Ok(CodergenResult::Full(Box::new(self.failed_outcome(
&error,
&live,
&fallback_plan,
))));
}
};
let route = fallback_plan.current().clone();
let stage_usage = billed_model_usage_from_llm(
let account = live.account();
let billing = stage_billing(
self.catalog.as_ref(),
&ModelRef::new(
route.target.provider.clone(),
ModelId::new(route.target.model.as_str()),
)
.with_speed(route.controls.speed),
live.total_usage,
)?
.with_reported_cost(live.total_cost);
&route_model(fallback_plan.current()),
&account,
)?;
live.release_lease();
match reuse_key {
@ -1204,9 +1376,10 @@ impl CodergenBackend for PebbleBackend {
Ok(CodergenResult::Text {
text: response,
usage: Some(stage_usage),
files_touched: live.files_touched.into_iter().collect(),
last_file_touched: live.last_file_touched,
usage: Some(billing.total),
usage_by_model: billing.by_model,
files_touched: account.files_touched,
last_file_touched: account.last_file_touched,
timing: StageTiming::active_only(
crate::millis_u64(live.inference_duration),
crate::millis_u64(live.tool_duration),
@ -1214,3 +1387,174 @@ impl CodergenBackend for PebbleBackend {
})
}
}
#[cfg(test)]
mod tests {
use std::time::SystemTime;
use fabro_llm::test_support::test_catalog;
use lithos_llm::catalog::builtin;
use pebble_coding_agent::events::{CodingAgentEvent, CodingEvent, InputSource, TokenUsage};
use super::*;
fn root(event: CodingEvent) -> CodingAgentEvent {
CodingAgentEvent::new("ses_root".to_string(), event, SystemTime::UNIX_EPOCH)
}
fn child(session_id: &str, event: CodingEvent) -> CodingAgentEvent {
CodingAgentEvent::new(session_id.to_string(), event, SystemTime::UNIX_EPOCH)
.with_parent_session_id("ses_root".to_string())
}
fn started(provider: &str, model: &str) -> CodingEvent {
CodingEvent::SessionStarted {
provider: Some(provider.to_string()),
model: Some(model.to_string()),
}
}
fn message(model: &str, input: u64, output: u64, cost: Option<u64>) -> CodingEvent {
CodingEvent::AssistantMessage {
text: "ok".to_string(),
model: model.to_string(),
usage: TokenUsage {
input,
output,
..TokenUsage::default()
},
cost_usd_micros: cost,
cost_source: None,
tool_call_count: 0,
context_window: None,
reasoning: None,
}
}
fn root_model() -> ModelRef {
ModelRef::new(builtin::openai(), ModelId::new("gpt-5.4"))
}
fn account(events: &[CodingAgentEvent]) -> SessionProjection {
let mut account = SessionProjection::new();
account.apply_all(events);
account
}
#[test]
fn stage_billing_prices_the_root_at_its_route_and_each_descendant_at_its_own() {
let catalog = test_catalog();
let account = account(&[
root(started("openai", "gpt-5.4")),
root(CodingEvent::UserInput {
text: "go".to_string(),
content: None,
source: InputSource::Prompt,
}),
root(message("gpt-5.4", 100_000, 25_000, None)),
// A child on the parent's route joins the parent's row.
child("ses_same", started("openai", "gpt-5.4")),
child("ses_same", message("gpt-5.4", 10_000, 1_000, None)),
// A child on another route is its own row, at that route's rate.
child("ses_other", started("anthropic", "claude-sonnet-5")),
child("ses_other", message("claude-sonnet-5", 20_000, 2_000, None)),
// A child on a route the catalog does not know bills at the root's.
child("ses_unknown", started("nowhere", "mystery")),
child("ses_unknown", message("mystery", 1_000, 100, None)),
root(CodingEvent::ProcessingEnd),
]);
let billing = stage_billing(&catalog, &root_model(), &account).unwrap();
assert_eq!(billing.by_model.len(), 2, "{:?}", billing.by_model);
let root_row = &billing.by_model[0];
assert_eq!(root_row.model, root_model());
assert_eq!(
root_row.tokens.input, 111_000,
"the root, the same-route child, and the unknown-route child"
);
assert_eq!(root_row.tokens.output, 26_100);
let root_priced =
billed_model_usage_from_llm(&catalog, &root_model(), root_row.tokens).unwrap();
assert_eq!(root_row.total_usd_micros, root_priced.total_usd_micros);
let other_model = ModelRef::new(
ProviderId::new("anthropic"),
ModelId::new("claude-sonnet-5"),
);
let other_row = &billing.by_model[1];
assert_eq!(other_row.model, other_model);
assert_eq!(other_row.tokens.input, 20_000);
assert_eq!(other_row.tokens.output, 2_000);
let other_priced =
billed_model_usage_from_llm(&catalog, &other_model, other_row.tokens).unwrap();
assert_eq!(other_row.total_usd_micros, other_priced.total_usd_micros);
assert_ne!(
other_row.total_usd_micros,
billed_model_usage_from_llm(&catalog, &root_model(), other_row.tokens)
.unwrap()
.total_usd_micros,
"priced at its own rate, not the root's"
);
// The total is the tree's tokens under the root's route, at the rows' summed
// cost.
assert_eq!(billing.total.model, root_model());
assert_eq!(billing.total.tokens.input, 131_000);
assert_eq!(billing.total.tokens.output, 28_100);
assert_eq!(
billing.total.total_usd_micros,
Some(root_priced.total_usd_micros.unwrap() + other_priced.total_usd_micros.unwrap())
);
}
#[test]
fn a_provider_reported_cost_stands_in_for_the_catalogs_estimate() {
let catalog = test_catalog();
let account = account(&[
root(started("openai", "gpt-5.4")),
root(message("gpt-5.4", 1_000, 100, Some(4_321))),
child("ses_child", started("anthropic", "claude-sonnet-5")),
child("ses_child", message("claude-sonnet-5", 500, 50, None)),
]);
let billing = stage_billing(&catalog, &root_model(), &account).unwrap();
assert_eq!(billing.by_model[0].total_usd_micros, Some(4_321));
let child_priced = billed_model_usage_from_llm(
&catalog,
&billing.by_model[1].model,
billing.by_model[1].tokens,
)
.unwrap();
assert_eq!(
billing.by_model[1].total_usd_micros,
child_priced.total_usd_micros
);
assert_eq!(
billing.total.total_usd_micros,
Some(4_321 + child_priced.total_usd_micros.unwrap())
);
}
#[test]
fn a_descendant_seen_only_through_its_answers_bills_on_the_roots_provider() {
let catalog = test_catalog();
let mut account = account(&[root(started("openai", "gpt-5.4"))]);
// No `SessionStarted` for the child: only its answer names a model.
account.apply(&child(
"ses_quiet",
message("gpt-5.4-mini", 1_000, 100, None),
));
let billing = stage_billing(&catalog, &root_model(), &account).unwrap();
let child_row = billing
.by_model
.iter()
.find(|row| row.model.model_id.as_str() == "gpt-5.4-mini")
.expect("the child is billed as its answers' model on the root's provider");
assert_eq!(child_row.model.provider, root_model().provider);
assert_eq!(child_row.tokens.input, 1_000);
}
}

View file

@ -160,6 +160,7 @@ mod tests {
async fn run(&self, _request: CodergenRunRequest<'_>) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: "api run".to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -170,6 +171,7 @@ mod tests {
async fn one_shot(&self, _request: OneShotRequest<'_>) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: "api one-shot".to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,

View file

@ -325,6 +325,7 @@ mod tests {
) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: "one-shot response".to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -383,6 +384,7 @@ mod tests {
) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: "one-shot response".to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -421,6 +423,7 @@ mod tests {
) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: r#"{"passed": true}"#.to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -469,6 +472,7 @@ mod tests {
) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: r#"{"outcome": 123}"#.to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -519,6 +523,7 @@ mod tests {
) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: "one-shot response".to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -579,6 +584,7 @@ mod tests {
Some(request.system_prompt.map(String::from));
Ok(CodergenResult::Text {
text: "classified".to_string(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,

View file

@ -125,6 +125,7 @@ mod duration_tests {
status: StageOutcome::Succeeded,
preferred_label: None,
suggested_next_ids: vec![],
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -158,11 +159,12 @@ mod duration_tests {
tool_call_id: None,
actor: None,
body: EventBody::StageFailed(StageFailedProps {
index: 0,
failure: None,
will_retry: true,
timing: StageTiming::wall_only(wall_time_ms),
billing: None,
index: 0,
failure: None,
will_retry: true,
timing: StageTiming::wall_only(wall_time_ms),
billing_by_model: Vec::new(),
billing: None,
}),
};
EventEnvelope { seq, event }
@ -250,6 +252,7 @@ mod duration_tests {
status: StageOutcome::Succeeded,
preferred_label: None,
suggested_next_ids: vec![],
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,

View file

@ -215,6 +215,7 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
status: StageOutcome::Succeeded.to_string(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -290,6 +291,7 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
will_retry: true,
timing,
billing: outcome.usage.clone(),
billing_by_model: outcome.usage_by_model.clone(),
actor,
},
&scope,
@ -342,6 +344,7 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
will_retry: false,
timing,
billing: outcome.usage.clone(),
billing_by_model: outcome.usage_by_model.clone(),
actor,
},
&scope,
@ -357,6 +360,7 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
preferred_label: outcome.preferred_label.clone(),
suggested_next_ids: outcome.suggested_next_ids.clone(),
billing: outcome.usage.clone(),
billing_by_model: outcome.usage_by_model.clone(),
failure: outcome.failure.clone(),
notes: outcome.notes.clone(),
files_touched: outcome.files_touched.clone(),

View file

@ -429,6 +429,7 @@ mod tests {
status: "succeeded".to_string(),
preferred_label: None,
suggested_next_ids: Vec::new(),
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,

View file

@ -2568,6 +2568,7 @@ mod tests {
preferred_label: None,
suggested_next_ids: Vec::new(),
billing,
billing_by_model: Vec::new(),
failure: None,
notes: None,
files_touched: Vec::new(),

View file

@ -1222,6 +1222,7 @@ capabilities = { text = true, tools = true, response_format = { json_object = tr
status: "succeeded".to_string(),
preferred_label: None,
suggested_next_ids: vec![],
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,
@ -1645,6 +1646,7 @@ capabilities = { text = true, tools = true, response_format = { json_object = tr
status: "succeeded".to_string(),
preferred_label: None,
suggested_next_ids: vec![],
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: None,

View file

@ -2202,6 +2202,7 @@ impl CodergenBackend for MockCodergenBackend {
request.node.id,
&request.prompt[..request.prompt.len().min(50)]
),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
@ -7425,6 +7426,7 @@ mod real_llm {
.map_err(|e| Error::handler(e.to_string()))?;
Ok(CodergenResult::Text {
text: response.text(),
usage_by_model: Vec::new(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,

View file

@ -742,6 +742,134 @@ async fn the_stage_timeout_fails_a_slow_agent() {
);
}
/// A stage whose agent fails for good after answering model calls bills
/// those calls: the failed outcome carries the session tree's usage from the
/// same fold the completed outcome would have, and the files it wrote.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_stage_that_fails_after_spending_bills_what_it_spent() {
let stage = Stage::new().await;
let first = stage.file("first.txt");
let second = stage.file("second.txt");
// Two answered calls, each writing a file; the third is refused for good.
stage
.server
.mock_async(|when, then| {
when.method(POST)
.path(CHAT_PATH)
.body_excludes(TOOL_RESULT_MARKER);
sse_headers(
then,
sse_tool_call(
"call-1",
"write_file",
&serde_json::json!({ "file_path": first, "content": "one" }),
),
);
})
.await;
stage
.server
.mock_async(|when, then| {
when.method(POST)
.path(CHAT_PATH)
.body_includes("call-1")
.body_excludes("call-2");
sse_headers(
then,
sse_tool_call(
"call-2",
"write_file",
&serde_json::json!({ "file_path": second, "content": "two" }),
),
);
})
.await;
stage
.server
.mock_async(|when, then| {
when.method(POST).path(CHAT_PATH).body_includes("call-2");
then.status(400)
.header("content-type", "application/json")
.body(r#"{"error":{"message":"the request was rejected","type":"invalid_request_error"}}"#);
})
.await;
let mut graph = agent_graph("Spent", "Write two files");
let work = graph.nodes.get_mut("work").unwrap();
work.attrs
.insert("max_retries".to_string(), AttrValue::Integer(0));
graph.edges.retain(|edge| edge.from != "work");
let mut fail_edge = Edge::new("work", "exit");
fail_edge.attrs.insert(
"condition".to_string(),
AttrValue::String("outcome=failed".to_string()),
);
graph.edges.push(fail_edge);
let backend = stage.backend("openai");
let (_, state) = stage
.run(backend, &graph, CancellationToken::new())
.await
.expect("the fail edge carries the run to exit");
let work = work_stage(&state);
assert_eq!(
work.completion
.as_ref()
.expect("the stage finished")
.outcome,
StageOutcome::Failed {
retry_requested: false,
}
);
assert_eq!(
work.usage.input_tokens,
2 * INPUT_TOKENS_PER_CALL,
"the two answered calls are billed"
);
assert_eq!(work.usage.output_tokens, 2 * OUTPUT_TOKENS_PER_CALL);
assert_eq!(
work.usage.total_usd_micros,
Some(2 * (INPUT_TOKENS_PER_CALL + 2 * OUTPUT_TOKENS_PER_CALL)),
"priced from the catalog like a completed stage"
);
assert_eq!(
work.billing_by_model.len(),
1,
"{:?}",
work.billing_by_model
);
assert_eq!(
work.billing_by_model[0].tokens.input,
u64::try_from(2 * INPUT_TOKENS_PER_CALL).unwrap()
);
assert!(
tokio::fs::try_exists(&second).await.unwrap(),
"the second write landed before the failure"
);
let failed = stage
.events
.lock()
.unwrap()
.iter()
.find(|event| {
event.event_name() == "stage.failed" && event.node_id.as_deref() == Some("work")
})
.cloned()
.expect("the stage failure is emitted");
let EventBody::StageFailed(props) = &failed.body else {
panic!("stage.failed carries its props: {failed:?}");
};
assert!(!props.will_retry);
let billing = props.billing.as_ref().expect("the failed stage is billed");
assert_eq!(
billing.tokens.input,
u64::try_from(2 * INPUT_TOKENS_PER_CALL).unwrap()
);
assert_eq!(props.billing_by_model, vec![billing.clone()]);
}
// --- Questions, subagents, MCP
// --------------------------------------------------
@ -873,6 +1001,57 @@ async fn a_subagent_runs_under_its_parent_session() {
assert_eq!(work_stage(&state).response.as_deref(), Some("Parent done"));
assert_eq!(count(&stage.events, "agent.sub.spawned"), 1);
// One usage rule: the stage bills its whole session tree, live and at
// completion. Four model calls answered: the parent's three and the
// child's one.
let work = work_stage(&state);
assert_eq!(
work.usage.input_tokens,
4 * INPUT_TOKENS_PER_CALL,
"the child's call is the stage's too"
);
assert_eq!(work.usage.output_tokens, 4 * OUTPUT_TOKENS_PER_CALL);
assert_eq!(
work.usage.total_usd_micros,
Some(4 * (INPUT_TOKENS_PER_CALL + 2 * OUTPUT_TOKENS_PER_CALL)),
"priced from the catalog for every call"
);
let agent = work
.agent
.as_ref()
.expect("the stage carries pebble's fold");
let (descendants, _) = agent.descendant_usage();
assert_eq!(
u64::try_from(work.usage.input_tokens).unwrap(),
agent.usage.input + descendants.input,
"the completed usage is what the live fold showed"
);
assert_eq!(
descendants.input,
u64::try_from(INPUT_TOKENS_PER_CALL).unwrap()
);
// The child ran on its parent's model, so the split is one row carrying
// the tree.
assert_eq!(
work.billing_by_model.len(),
1,
"{:?}",
work.billing_by_model
);
assert_eq!(
work.billing_by_model[0].tokens.input,
u64::try_from(4 * INPUT_TOKENS_PER_CALL).unwrap()
);
assert_eq!(
Some(&work.billing_by_model[0].model),
work.model.as_ref(),
"billed under the root's route"
);
assert_eq!(
work.billing_by_model[0].total_usd_micros,
work.usage.total_usd_micros
);
let agent_events = coding_events(&stage.events);
let root_session = agent_events
.iter()

View file

@ -385,6 +385,7 @@ fn main() {
&[],
),
("StageProjection", "fabro_types::StageProjection", &[]),
("BilledModelUsage", "fabro_types::BilledModelUsage", &[]),
(
"StageInferenceProjection",
"fabro_types::StageInferenceProjection",

View file

@ -39,39 +39,39 @@ pub mod types {
};
pub use fabro_types::{
ActivatedSkill, AgentControlState, AgentEventProps, AgentMcpToolSummary,
AgentToolsAvailableProps, AskFabro, AuthMethod, AutomationRef, BilledTokenCounts, BlobHash,
CommandTermination, Conclusion, ContextWindowBreakdownItem, ContextWindowCategory,
ContextWindowCountMethod, ContextWindowSnapshot, ContextWindowStaleness,
ContextWindowWarning, CreateVariableRequest, DiffStats, DiffSummary, DirtyStatus,
EventEnvelope, ExecOutputTail, FailureCategory, FailureDetail, FailureSignature,
GitContext, GitRunTarget, GitRunTarget as AutomationGitWorkflowSource, IdpIdentity,
IntegrationConnectionKind, IntegrationConnectionState, IntegrationConnectionStatus,
IntegrationProvider, IntegrationStatus, InterviewOption, InterviewQuestionRecord,
LlmOutputKind, McpServerDraft as CreateMcpServerRequest, McpServerProjection,
McpServerReplace as ReplaceMcpServerRequest, McpServerStatus, McpServerView as McpServer,
McpTransportView, Model, ModelControls, ModelCosts, ModelFeatures, ModelLimits,
ModelRef as BillingModelRef, ModelTestMode, PairId, PairMessageId, PairMessageRecord,
PairMessageRequest, PairRecord, PairStartRequest, PairStatus, PairTarget,
PairTranscriptEntry, PairTranscriptResponse, ParallelBranchId, ParallelBranchResult,
PendingInterviewRecord, PermissionLevel, Principal, Provider, PullRequest,
PullRequestCreation, PullRequestCreationId, PullRequestCreationStatus, PullRequestDetails,
PullRequestDetailsStatus, PullRequestDetailsUnavailableReason, PullRequestLink,
PullRequestMeta, PullRequestResponse, QuestionType, RepositoryRef, ReviewTarget,
ReviewTargetKind, Run, RunApproval, RunApprovalState, RunClientProvenance, RunEvent,
RunEventDetailContentKind, RunEventDetailResponse, RunFailure, RunIntent, RunIntentArgs,
RunPairStatusResponse, RunProjection, RunProvenance, RunRunnableSource, RunSandbox,
RunSandboxFailure, RunSandboxInstance, RunSandboxKind, RunSandboxPlan, RunSandboxRuntime,
RunServerProvenance, RunSessionMetadata, RunSize, RunTarget, SandboxDetails, SandboxInfo,
SandboxListMeta, SandboxListResponse, SandboxProviderKind, SandboxProviderLookupError,
SandboxService, SandboxServiceListResponse, SecretMetadata, SecretType, ServerSettings,
SessionDetail, SessionId, SessionStatus, SessionSummary, SessionTurn,
SkillActivationSource, SkillSummary, SkillsProjection, StageCompletion, StageContextWindow,
StageContextWindowUnavailableReason, StageHandler, StageId, StageInferenceProjection,
StageModelUsage, StageOutcome, StageProjection, StageState, StageToolBatchProjection,
SubAgentProjection, SubAgentStatus, SystemActorKind, SystemIntegrationStatus,
SystemIntegrationsResponse, TodoListProjection, ToolCategory, ToolSource, ToolSummary,
TurnId, UpdateVariableRequest, UserPrincipal, Variable, VariableListResponse, WorkflowPath,
WorkflowSettings, WorkflowVersion, WorkflowVersionId,
AgentToolsAvailableProps, AskFabro, AuthMethod, AutomationRef, BilledModelUsage,
BilledTokenCounts, BlobHash, CommandTermination, Conclusion, ContextWindowBreakdownItem,
ContextWindowCategory, ContextWindowCountMethod, ContextWindowSnapshot,
ContextWindowStaleness, ContextWindowWarning, CreateVariableRequest, DiffStats,
DiffSummary, DirtyStatus, EventEnvelope, ExecOutputTail, FailureCategory, FailureDetail,
FailureSignature, GitContext, GitRunTarget, GitRunTarget as AutomationGitWorkflowSource,
IdpIdentity, IntegrationConnectionKind, IntegrationConnectionState,
IntegrationConnectionStatus, IntegrationProvider, IntegrationStatus, InterviewOption,
InterviewQuestionRecord, LlmOutputKind, McpServerDraft as CreateMcpServerRequest,
McpServerProjection, McpServerReplace as ReplaceMcpServerRequest, McpServerStatus,
McpServerView as McpServer, McpTransportView, Model, ModelControls, ModelCosts,
ModelFeatures, ModelLimits, ModelRef as BillingModelRef, ModelTestMode, PairId,
PairMessageId, PairMessageRecord, PairMessageRequest, PairRecord, PairStartRequest,
PairStatus, PairTarget, PairTranscriptEntry, PairTranscriptResponse, ParallelBranchId,
ParallelBranchResult, PendingInterviewRecord, PermissionLevel, Principal, Provider,
PullRequest, PullRequestCreation, PullRequestCreationId, PullRequestCreationStatus,
PullRequestDetails, PullRequestDetailsStatus, PullRequestDetailsUnavailableReason,
PullRequestLink, PullRequestMeta, PullRequestResponse, QuestionType, RepositoryRef,
ReviewTarget, ReviewTargetKind, Run, RunApproval, RunApprovalState, RunClientProvenance,
RunEvent, RunEventDetailContentKind, RunEventDetailResponse, RunFailure, RunIntent,
RunIntentArgs, RunPairStatusResponse, RunProjection, RunProvenance, RunRunnableSource,
RunSandbox, RunSandboxFailure, RunSandboxInstance, RunSandboxKind, RunSandboxPlan,
RunSandboxRuntime, RunServerProvenance, RunSessionMetadata, RunSize, RunTarget,
SandboxDetails, SandboxInfo, SandboxListMeta, SandboxListResponse, SandboxProviderKind,
SandboxProviderLookupError, SandboxService, SandboxServiceListResponse, SecretMetadata,
SecretType, ServerSettings, SessionDetail, SessionId, SessionStatus, SessionSummary,
SessionTurn, SkillActivationSource, SkillSummary, SkillsProjection, StageCompletion,
StageContextWindow, StageContextWindowUnavailableReason, StageHandler, StageId,
StageInferenceProjection, StageModelUsage, StageOutcome, StageProjection, StageState,
StageToolBatchProjection, SubAgentProjection, SubAgentStatus, SystemActorKind,
SystemIntegrationStatus, SystemIntegrationsResponse, TodoListProjection, ToolCategory,
ToolSource, ToolSummary, TurnId, UpdateVariableRequest, UserPrincipal, Variable,
VariableListResponse, WorkflowPath, WorkflowSettings, WorkflowVersion, WorkflowVersionId,
};
pub use lithos_llm::catalog::{ModelHandle, ProviderId};
pub use lithos_llm::types::{

View file

@ -4,6 +4,7 @@ use fabro_api::types::{
ActivatedSkill as ApiActivatedSkill, AgentControlState as ApiAgentControlState,
AgentMcpToolSummary as ApiAgentMcpToolSummary,
AgentToolsAvailableProps as ApiAgentToolsAvailableProps,
BilledModelUsage as ApiBilledModelUsage,
ContextWindowBreakdownItem as ApiContextWindowBreakdownItem,
ContextWindowCategory as ApiContextWindowCategory,
ContextWindowCountMethod as ApiContextWindowCountMethod,
@ -23,19 +24,91 @@ use fabro_api::types::{
};
use fabro_types::{
ActivatedSkill, AgentControlState, AgentMcpToolSummary, AgentToolsAvailableProps,
ContextWindowBreakdownItem, ContextWindowCategory, ContextWindowCountMethod,
BilledModelUsage, ContextWindowBreakdownItem, ContextWindowCategory, ContextWindowCountMethod,
ContextWindowSnapshot, ContextWindowStaleness, ContextWindowWarning, LlmOutputKind,
McpServerProjection, McpServerStatus, ParallelBranchId, ParallelBranchResult, PermissionLevel,
SkillActivationSource, SkillSummary, SkillsProjection, StageContextWindow,
McpServerProjection, McpServerStatus, ModelRef, ParallelBranchId, ParallelBranchResult,
PermissionLevel, SkillActivationSource, SkillSummary, SkillsProjection, StageContextWindow,
StageContextWindowUnavailableReason, StageId, StageInferenceProjection, StageProjection,
StageToolBatchProjection, SubAgentProjection, SubAgentStatus, TodoListKind, TodoListProjection,
ToolCategory, ToolSource, ToolSummary,
};
use lithos_llm::catalog::{ModelId, ProviderId};
use lithos_llm::types::TokenCounts;
use serde_json::json;
#[test]
fn stage_projection_reuses_canonical_type() {
assert_same_type::<ApiStageProjection, StageProjection>();
assert_same_type::<ApiBilledModelUsage, BilledModelUsage>();
}
#[test]
fn billing_by_model_rows_match_openapi_json_shape() {
let row = BilledModelUsage {
model: ModelRef::new(ProviderId::new("openai"), ModelId::new("gpt-5.4")),
tokens: TokenCounts {
input: 107,
output: 51,
..TokenCounts::default()
},
total_usd_micros: Some(321),
};
let value = serde_json::to_value(&row).unwrap();
assert_eq!(
value,
json!({
"model": { "provider": "openai", "model_id": "gpt-5.4" },
"tokens": {
"input": 107,
"output": 51,
"reasoning": 0,
"cache_read": 0,
"cache_write": 0
},
"total_usd_micros": 321
})
);
let api_row: ApiBilledModelUsage = serde_json::from_value(value).unwrap();
assert_eq!(api_row, row);
let mut stage = StageProjection::new(std::num::NonZeroU32::new(1).unwrap());
stage.billing_by_model = vec![row.clone()];
let stage_json = serde_json::to_value(&stage).unwrap();
assert_eq!(
stage_json["billing_by_model"],
json!([serde_json::to_value(&row).unwrap()])
);
let without: StageProjection = serde_json::from_value(json!({
"first_event_seq": 1,
"prompt": null,
"response": null,
"completion": null,
"provider_used": null,
"diff": null,
"script_invocation": null,
"script_timing": null,
"parallel_results": null,
"output": null,
"usage": {
"input_tokens": 0,
"output_tokens": 0,
"total_tokens": 0,
"reasoning_tokens": 0,
"cache_read_tokens": 0,
"cache_write_tokens": 0
},
"agent_control": "running",
"state": "running"
}))
.unwrap();
assert!(without.billing_by_model.is_empty());
assert!(
serde_json::to_value(&without)
.unwrap()
.get("billing_by_model")
.is_none(),
"no rows, nothing on the wire"
);
}
#[test]

View file

@ -1,6 +1,9 @@
use std::collections::HashMap;
use crate::{BilledTokenCounts, ModelRef, RunProjection, RunTiming, StageSummary, StageTiming};
use crate::{
BilledTokenCounts, ModelRef, RunProjection, RunTiming, StageProjection, StageSummary,
StageTiming,
};
#[derive(Debug, Clone, PartialEq)]
pub struct ProjectionBillingStage {
@ -159,16 +162,21 @@ pub fn billing_rollup_from_projection(projection: &RunProjection) -> ProjectionB
if let Some(model) = &stage.model {
row.model = Some(model.clone());
}
// A completed agent stage says which model billed which tokens:
// the root's route and each subagent's own. Until then, and for
// a stage without a coding agent, `usage` bills to `model`.
for (model, billing) in model_rows(stage) {
let model_entry =
by_model
.entry(model.clone())
.or_insert_with(|| ProjectionBillingByModel {
model: model.clone(),
stages: 0,
model,
stages: 0,
billing: BilledTokenCounts::default(),
});
model_entry.stages += 1;
model_entry.billing.add_counts(usage);
model_entry.billing.add_counts(&billing);
}
}
}
@ -185,6 +193,28 @@ pub fn billing_rollup_from_projection(projection: &RunProjection) -> ProjectionB
}
}
/// The stage's usage by model: its `billing_by_model` rows when the stage
/// completed with them, else its `usage` under its `model`.
fn model_rows(stage: &StageProjection) -> Vec<(ModelRef, BilledTokenCounts)> {
if stage.billing_by_model.is_empty() {
return stage
.model
.iter()
.map(|model| (model.clone(), stage.usage.clone()))
.collect();
}
stage
.billing_by_model
.iter()
.map(|row| {
(
row.model.clone(),
BilledTokenCounts::from_token_counts(row.tokens, row.total_usd_micros),
)
})
.collect()
}
fn stage_projection_order(state: &RunProjection) -> HashMap<String, u32> {
let mut order = HashMap::new();
for (stage_id, stage) in state.iter_stages() {

View file

@ -9,7 +9,8 @@ use serde_json::Value;
use strum::{Display, EnumString, IntoStaticStr};
use crate::{
ExecOutputTail, FailureSignature, OnFailure, ResolvedOnFailure, StageTiming, SystemActorKind,
BilledModelUsage, ExecOutputTail, FailureSignature, OnFailure, ResolvedOnFailure, StageTiming,
SystemActorKind,
};
pub trait OutcomeMeta:
@ -274,6 +275,11 @@ pub struct Outcome<M: OutcomeMeta = ()> {
pub failure: Option<FailureDetail>,
#[serde(default)]
pub usage: M,
/// The stage's billing split by model, for a stage whose agent ran
/// subagents: the root's route and each subagent's own model. Empty
/// otherwise; `usage` is then the one row.
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub usage_by_model: Vec<BilledModelUsage>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub files_touched: Vec<String>,
/// Stage timing breakdown captured by the workflow engine.
@ -296,6 +302,7 @@ impl<M: OutcomeMeta> Default for Outcome<M> {
notes: None,
failure: None,
usage: M::default(),
usage_by_model: Vec::new(),
files_touched: Vec::new(),
timing: None,
}

View file

@ -1202,6 +1202,7 @@ mod tests {
status: crate::StageOutcome::Succeeded,
preferred_label: None,
suggested_next_ids: vec!["next".to_string()],
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: Some("done".to_string()),
@ -1677,6 +1678,7 @@ mod tests {
status: crate::StageOutcome::Succeeded,
preferred_label: None,
suggested_next_ids: vec!["next".to_string()],
billing_by_model: Vec::new(),
billing: None,
failure: None,
notes: Some("done".to_string()),

View file

@ -36,8 +36,16 @@ pub struct StageCompletedProps {
pub preferred_label: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub suggested_next_ids: Vec<String>,
/// The stage's billing: for an agent stage, the whole session tree's
/// tokens (the root session and every subagent) under the root's route.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub billing: Option<BilledModelUsage>,
/// `billing` split by model: the root session's route and each
/// subagent's own model, a subagent whose model the catalog does not know
/// billed at the root's. Sums to `billing`. Empty for stages without a
/// coding agent and on events written before it existed.
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub billing_by_model: Vec<BilledModelUsage>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub failure: Option<FailureDetail>,
#[serde(default, skip_serializing_if = "Option::is_none")]
@ -64,15 +72,20 @@ pub struct StageCompletedProps {
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct StageFailedProps {
pub index: usize,
pub index: usize,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub failure: Option<FailureDetail>,
pub will_retry: bool,
pub failure: Option<FailureDetail>,
pub will_retry: bool,
/// Per-attempt timing breakdown for this stage visit.
#[serde(default)]
pub timing: StageTiming,
pub timing: StageTiming,
/// The stage's billing: for an agent stage that failed after spending,
/// the whole session tree's tokens under the root's route.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub billing: Option<BilledModelUsage>,
pub billing: Option<BilledModelUsage>,
/// `billing` split by model, as on `stage.completed`.
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub billing_by_model: Vec<BilledModelUsage>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]

View file

@ -14,11 +14,11 @@ use strum::{Display, EnumString, IntoStaticStr};
use crate::run_event::{AgentSessionActivatedProps, StagePromptProps};
use crate::{
AgentBackend, AgentMcpToolSummary, BilledTokenCounts, Checkpoint, Conclusion, GitIdentity,
InterviewQuestionRecord, InvalidTransition, ModelRef, ParallelBranchId, PullRequestCreation,
PullRequestLink, RunApproval, RunControlAction, RunDiff, RunId, RunSandbox, RunSpec, RunStatus,
RunTiming, StageCompletion, StageHandler, StageId, StageState, StageTiming, StartRecord,
timing,
AgentBackend, AgentMcpToolSummary, BilledModelUsage, BilledTokenCounts, Checkpoint, Conclusion,
GitIdentity, InterviewQuestionRecord, InvalidTransition, ModelRef, ParallelBranchId,
PullRequestCreation, PullRequestLink, RunApproval, RunControlAction, RunDiff, RunId,
RunSandbox, RunSpec, RunStatus, RunTiming, StageCompletion, StageHandler, StageId, StageState,
StageTiming, StartRecord, timing,
};
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
@ -297,6 +297,13 @@ pub struct StageProjection {
pub usage: BilledTokenCounts,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model: Option<ModelRef>,
/// The finished stage's billing split by model, as `stage.completed` or
/// `stage.failed` reported it: the root session's route and each
/// subagent's own model. Sums to `usage`. Empty while the stage runs and
/// for stages without a coding agent; the billing rollup then bills
/// `usage` to `model`.
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub billing_by_model: Vec<BilledModelUsage>,
/// Todo/task list owned by the stage's root agent session.
///
/// OpenAI child sessions own separate per-session plans and do not appear
@ -514,6 +521,7 @@ impl StageProjection {
acp_started_at: None,
agent_control: AgentControlState::default(),
agent: None,
billing_by_model: Vec::new(),
provider_used: None,
diff: None,
script_invocation: None,

View file

@ -87,6 +87,7 @@ models/batch-run-lifecycle-request.ts
models/batch-run-lifecycle-response.ts
models/batch-run-lifecycle-result.ts
models/batch-run-lifecycle-summary.ts
models/billed-model-usage.ts
models/billed-token-counts.ts
models/billing-by-model.ts
models/billing-model-ref.ts

View file

@ -0,0 +1,33 @@
/* tslint:disable */
/* eslint-disable */
/**
* Fabro Run API
* HTTP API for managing Fabro workflow run executions.
*
* The version of the OpenAPI document: 0.2.0
*
*
* NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech).
* https://openapi-generator.tech
* Do not edit the class manually.
*/
// May contain unused imports in some cases
// @ts-ignore
import type { BillingModelRef } from './billing-model-ref';
// May contain unused imports in some cases
// @ts-ignore
import type { CompletionUsage } from './completion-usage';
/**
* Usage and cost billed to one model: one response, or one model\'s share of a stage.
*/
export interface BilledModelUsage {
'model': BillingModelRef;
'tokens': CompletionUsage;
/**
* Cost for `tokens`, when the provider reported one or the catalog could price them. Absent means no cost data, not zero.
*/
'total_usd_micros'?: number;
}

View file

@ -58,6 +58,7 @@ export * from './batch-run-lifecycle-request';
export * from './batch-run-lifecycle-response';
export * from './batch-run-lifecycle-result';
export * from './batch-run-lifecycle-summary';
export * from './billed-model-usage';
export * from './billed-token-counts';
export * from './billing-by-model';
export * from './billing-model-ref';

View file

@ -21,6 +21,9 @@ import type { AgentControlState } from './agent-control-state';
import type { AgentSessionProjection } from './agent-session-projection';
// May contain unused imports in some cases
// @ts-ignore
import type { BilledModelUsage } from './billed-model-usage';
// May contain unused imports in some cases
// @ts-ignore
import type { BilledTokenCounts } from './billed-token-counts';
// May contain unused imports in some cases
// @ts-ignore
@ -145,6 +148,10 @@ export interface StageProjection {
* Whether the agent is executing normally or waiting for steering after an interrupt.
*/
'agent_control': AgentControlState;
/**
* The completed stage\'s `usage` split by model, as `stage.completed` reported it: the root session\'s route and each subagent\'s own model, a subagent whose model the catalog does not know billed at the root\'s. Sums to `usage`. Empty while the stage runs and for stages without a coding agent; the billing rollup then bills `usage` to `model`.
*/
'billing_by_model'?: Array<BilledModelUsage>;
'agent'?: AgentSessionProjection | null;
/**
* Lifecycle state of the stage projection.