Bill an agent stage's whole session tree from one fold

One usage rule: a stage's usage is its session tree's, the root and every
subagent, live and at completion. The worker's event sink folds pebble's
SessionProjection over the events it records and the stage's billing and
files come from that fold at stage end, so the completed values are what
the run showed live. The store's live usage is the fold's tree usage, and
completion brings the catalog's price for the same tokens instead of
resetting them to the root's.

Fabro keeps catalog pricing: the root at its route, each descendant at its
own route where the catalog knows it and at the root's otherwise, a
provider-reported cost standing in where pebble has one. The rows travel
as billing_by_model on stage.completed and the stage projection, and the
billing rollup splits by_model by them.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-13 07:49:27 -06:00
parent fa6902c448
commit 850d5cca52
No known key found for this signature in database
26 changed files with 812 additions and 144 deletions

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,
@ -7157,6 +7159,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,
@ -7394,6 +7397,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 +7486,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 +7809,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 +7841,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,
@ -8278,6 +8285,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 +8357,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 +15411,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,

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;
}
@ -760,9 +762,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 +780,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 +988,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 +1789,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 +1865,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 +3668,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 +3756,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 +3801,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,
@ -5503,6 +5517,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 +5662,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 +5695,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]

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

@ -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(),
@ -1097,6 +1099,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 +1145,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

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

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;
@ -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,101 @@ impl LiveAgent {
}
}
/// 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 +704,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 +742,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 +752,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 +774,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
@ -986,6 +1110,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 +1151,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 +1166,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 +1184,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 +1203,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 +1229,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,
@ -1171,16 +1296,13 @@ impl CodergenBackend for PebbleBackend {
};
let route = fallback_plan.current().clone();
let stage_usage = billed_model_usage_from_llm(
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);
let root_model = ModelRef::new(
route.target.provider.clone(),
ModelId::new(route.target.model.as_str()),
)
.with_speed(route.controls.speed);
let account = live.account();
let billing = stage_billing(self.catalog.as_ref(), &root_model, &account)?;
live.release_lease();
match reuse_key {
@ -1204,9 +1326,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 +1337,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,
@ -250,6 +251,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,
@ -357,6 +358,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

@ -873,6 +873,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

@ -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")]

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,12 @@ pub struct StageProjection {
pub usage: BilledTokenCounts,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model: Option<ModelRef>,
/// The completed stage's billing split by model, as `stage.completed`
/// 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 +520,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,