From df8762663b6ac2b2799661a7804c75295b74fbfb Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sun, 13 Sep 2026 07:49:27 -0600 Subject: [PATCH] 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 --- .../src/commands/run/run_progress/event.rs | 1 + .../src/commands/run/run_progress/mod.rs | 1 + lib/apps/fabro-server/src/server/tests.rs | 10 + lib/components/fabro-store/src/run_state.rs | 202 +++++++- .../fabro-workflow/src/billing_rollup.rs | 43 ++ .../fabro-workflow/src/event/convert.rs | 4 + .../fabro-workflow/src/event/events.rs | 2 + lib/components/fabro-workflow/src/git.rs | 1 + .../fabro-workflow/src/handler/agent.rs | 113 +++-- .../fabro-workflow/src/handler/fan_in.rs | 1 + .../fabro-workflow/src/handler/llm/acp.rs | 1 + .../fabro-workflow/src/handler/llm/pebble.rs | 434 +++++++++++++++--- .../fabro-workflow/src/handler/llm/router.rs | 2 + .../fabro-workflow/src/handler/prompt.rs | 6 + lib/components/fabro-workflow/src/lib.rs | 2 + .../fabro-workflow/src/lifecycle/event.rs | 2 + .../fabro-workflow/src/operations/fork.rs | 1 + .../fabro-workflow/src/operations/start.rs | 1 + .../src/pipeline/pull_request.rs | 2 + .../fabro-workflow/tests/it/integration.rs | 2 + .../fabro-workflow/tests/it/pebble_agent.rs | 51 ++ .../fabro-types/src/billing_rollup.rs | 38 +- lib/foundation/fabro-types/src/outcome.rs | 9 +- .../fabro-types/src/run_event/mod.rs | 2 + .../fabro-types/src/run_event/stage.rs | 8 + .../fabro-types/src/run_projection.rs | 17 +- 26 files changed, 812 insertions(+), 144 deletions(-) diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs index 895494d78..7d9377dde 100644 --- a/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs @@ -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, diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs index 94c6ce88c..ca3f5a020 100644 --- a/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs @@ -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(), diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index 1b86f0c90..8f9ef998b 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -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, 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, diff --git a/lib/components/fabro-store/src/run_state.rs b/lib/components/fabro-store/src/run_state.rs index e0cedb00f..afe270b99 100644 --- a/lib/components/fabro-store/src/run_state.rs +++ b/lib/components/fabro-store/src/run_state.rs @@ -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) -> 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 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] diff --git a/lib/components/fabro-workflow/src/billing_rollup.rs b/lib/components/fabro-workflow/src/billing_rollup.rs index 5dc6ca164..2c3a56b34 100644 --- a/lib/components/fabro-workflow/src/billing_rollup.rs +++ b/lib/components/fabro-workflow/src/billing_rollup.rs @@ -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(); diff --git a/lib/components/fabro-workflow/src/event/convert.rs b/lib/components/fabro-workflow/src/event/convert.rs index dbbf259c3..fa0bb6d61 100644 --- a/lib/components/fabro-workflow/src/event/convert.rs +++ b/lib/components/fabro-workflow/src/event/convert.rs @@ -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, diff --git a/lib/components/fabro-workflow/src/event/events.rs b/lib/components/fabro-workflow/src/event/events.rs index 01ee3eba7..57a639a38 100644 --- a/lib/components/fabro-workflow/src/event/events.rs +++ b/lib/components/fabro-workflow/src/event/events.rs @@ -272,6 +272,8 @@ pub enum Event { preferred_label: Option, suggested_next_ids: Vec, billing: Option, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + billing_by_model: Vec, #[serde(default, skip_serializing_if = "Option::is_none")] failure: Option, notes: Option, diff --git a/lib/components/fabro-workflow/src/git.rs b/lib/components/fabro-workflow/src/git.rs index 5bc514a19..c169215f2 100644 --- a/lib/components/fabro-workflow/src/git.rs +++ b/lib/components/fabro-workflow/src/git.rs @@ -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, diff --git a/lib/components/fabro-workflow/src/handler/agent.rs b/lib/components/fabro-workflow/src/handler/agent.rs index da9c912b5..e2b0c4a09 100644 --- a/lib/components/fabro-workflow/src/handler/agent.rs +++ b/lib/components/fabro-workflow/src/handler/agent.rs @@ -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, + /// `usage` split by model, when the backend billed subagents at + /// their own models. Empty when `usage` is the one row. + usage_by_model: Vec, files_touched: Vec, last_file_touched: Option, /// 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 }); - 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 { 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 { 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 { 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 { 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, diff --git a/lib/components/fabro-workflow/src/handler/fan_in.rs b/lib/components/fabro-workflow/src/handler/fan_in.rs index 2665744fc..410409b10 100644 --- a/lib/components/fabro-workflow/src/handler/fan_in.rs +++ b/lib/components/fabro-workflow/src/handler/fan_in.rs @@ -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, diff --git a/lib/components/fabro-workflow/src/handler/llm/acp.rs b/lib/components/fabro-workflow/src/handler/llm/acp.rs index 04983a443..444e754ae 100644 --- a/lib/components/fabro-workflow/src/handler/llm/acp.rs +++ b/lib/components/fabro-workflow/src/handler/llm/acp.rs @@ -449,6 +449,7 @@ impl AgentAcpBackend { Ok(CodergenResult::Text { text: result.text, + usage_by_model: Vec::new(), usage: None, files_touched, last_file_touched, diff --git a/lib/components/fabro-workflow/src/handler/llm/pebble.rs b/lib/components/fabro-workflow/src/handler/llm/pebble.rs index 07f4231f0..fffb259db 100644 --- a/lib/components/fabro-workflow/src/handler/llm/pebble.rs +++ b/lib/components/fabro-workflow/src/handler/llm/pebble.rs @@ -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, - node_id: String, - scope: StageScope, + emitter: Arc, + 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, +} + +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>, - total_usage: TokenCounts, - total_cost: Option, + sink: Arc, inference_duration: Duration, tool_duration: Duration, - /// Every file the stage's prompts wrote or edited, subagents included. - files_touched: BTreeSet, - /// The most recently written path. - last_file_touched: Option, } impl LiveAgent { - fn new(agent: CodingAgent, handle: CodingAgentControlHandle) -> Self { + fn new( + agent: CodingAgent, + handle: CodingAgentControlHandle, + sink: Arc, + ) -> 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 { + 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, +} + +/// 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 { + let mut groups: Vec<(ModelRef, TokenCounts, Option)> = 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, cost: Option) { + 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) { 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; 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 { + ) -> Result<(CodingAgent, Arc), Error> { let client = self.build_llm_client().await?; let environment: Arc = Arc::clone(bindings.sandbox) as Arc; 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 { + ) -> Result<(CodingAgent, Arc), Error> { let client = self.build_llm_client().await?; let environment: Arc = Arc::clone(bindings.sandbox) as Arc; 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) -> 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); + } +} diff --git a/lib/components/fabro-workflow/src/handler/llm/router.rs b/lib/components/fabro-workflow/src/handler/llm/router.rs index 45e2250d8..6ecbd7a01 100644 --- a/lib/components/fabro-workflow/src/handler/llm/router.rs +++ b/lib/components/fabro-workflow/src/handler/llm/router.rs @@ -160,6 +160,7 @@ mod tests { async fn run(&self, _request: CodergenRunRequest<'_>) -> Result { 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 { Ok(CodergenResult::Text { text: "api one-shot".to_string(), + usage_by_model: Vec::new(), usage: None, files_touched: Vec::new(), last_file_touched: None, diff --git a/lib/components/fabro-workflow/src/handler/prompt.rs b/lib/components/fabro-workflow/src/handler/prompt.rs index 85dd74c45..1f5137c3e 100644 --- a/lib/components/fabro-workflow/src/handler/prompt.rs +++ b/lib/components/fabro-workflow/src/handler/prompt.rs @@ -325,6 +325,7 @@ mod tests { ) -> Result { 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 { 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 { 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 { 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 { 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, diff --git a/lib/components/fabro-workflow/src/lib.rs b/lib/components/fabro-workflow/src/lib.rs index 918d08973..4b66dea59 100644 --- a/lib/components/fabro-workflow/src/lib.rs +++ b/lib/components/fabro-workflow/src/lib.rs @@ -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, diff --git a/lib/components/fabro-workflow/src/lifecycle/event.rs b/lib/components/fabro-workflow/src/lifecycle/event.rs index 6303ea1f3..278479006 100644 --- a/lib/components/fabro-workflow/src/lifecycle/event.rs +++ b/lib/components/fabro-workflow/src/lifecycle/event.rs @@ -215,6 +215,7 @@ impl RunLifecycle 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 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(), diff --git a/lib/components/fabro-workflow/src/operations/fork.rs b/lib/components/fabro-workflow/src/operations/fork.rs index 00c2e4af4..d601cc1bc 100644 --- a/lib/components/fabro-workflow/src/operations/fork.rs +++ b/lib/components/fabro-workflow/src/operations/fork.rs @@ -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, diff --git a/lib/components/fabro-workflow/src/operations/start.rs b/lib/components/fabro-workflow/src/operations/start.rs index 83b5fbe1e..76e1f9d4c 100644 --- a/lib/components/fabro-workflow/src/operations/start.rs +++ b/lib/components/fabro-workflow/src/operations/start.rs @@ -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(), diff --git a/lib/components/fabro-workflow/src/pipeline/pull_request.rs b/lib/components/fabro-workflow/src/pipeline/pull_request.rs index 879d91eaa..ffcf4e88e 100644 --- a/lib/components/fabro-workflow/src/pipeline/pull_request.rs +++ b/lib/components/fabro-workflow/src/pipeline/pull_request.rs @@ -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, diff --git a/lib/components/fabro-workflow/tests/it/integration.rs b/lib/components/fabro-workflow/tests/it/integration.rs index fb2df9eaf..118fc20ff 100644 --- a/lib/components/fabro-workflow/tests/it/integration.rs +++ b/lib/components/fabro-workflow/tests/it/integration.rs @@ -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, diff --git a/lib/components/fabro-workflow/tests/it/pebble_agent.rs b/lib/components/fabro-workflow/tests/it/pebble_agent.rs index d8d88a466..b749ad120 100644 --- a/lib/components/fabro-workflow/tests/it/pebble_agent.rs +++ b/lib/components/fabro-workflow/tests/it/pebble_agent.rs @@ -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() diff --git a/lib/foundation/fabro-types/src/billing_rollup.rs b/lib/foundation/fabro-types/src/billing_rollup.rs index 9ec5aee23..01b4a31e2 100644 --- a/lib/foundation/fabro-types/src/billing_rollup.rs +++ b/lib/foundation/fabro-types/src/billing_rollup.rs @@ -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 { let mut order = HashMap::new(); for (stage_id, stage) in state.iter_stages() { diff --git a/lib/foundation/fabro-types/src/outcome.rs b/lib/foundation/fabro-types/src/outcome.rs index 4aadc9aac..02810714a 100644 --- a/lib/foundation/fabro-types/src/outcome.rs +++ b/lib/foundation/fabro-types/src/outcome.rs @@ -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 { pub failure: Option, #[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, #[serde(default, skip_serializing_if = "Vec::is_empty")] pub files_touched: Vec, /// Stage timing breakdown captured by the workflow engine. @@ -296,6 +302,7 @@ impl Default for Outcome { notes: None, failure: None, usage: M::default(), + usage_by_model: Vec::new(), files_touched: Vec::new(), timing: None, } diff --git a/lib/foundation/fabro-types/src/run_event/mod.rs b/lib/foundation/fabro-types/src/run_event/mod.rs index 048823c6a..642d43d12 100644 --- a/lib/foundation/fabro-types/src/run_event/mod.rs +++ b/lib/foundation/fabro-types/src/run_event/mod.rs @@ -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()), diff --git a/lib/foundation/fabro-types/src/run_event/stage.rs b/lib/foundation/fabro-types/src/run_event/stage.rs index 422780c4d..524fff69f 100644 --- a/lib/foundation/fabro-types/src/run_event/stage.rs +++ b/lib/foundation/fabro-types/src/run_event/stage.rs @@ -36,8 +36,16 @@ pub struct StageCompletedProps { pub preferred_label: Option, #[serde(default, skip_serializing_if = "Vec::is_empty")] pub suggested_next_ids: Vec, + /// 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, + /// `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, #[serde(default, skip_serializing_if = "Option::is_none")] pub failure: Option, #[serde(default, skip_serializing_if = "Option::is_none")] diff --git a/lib/foundation/fabro-types/src/run_projection.rs b/lib/foundation/fabro-types/src/run_projection.rs index 9f8e446b8..5e0540d78 100644 --- a/lib/foundation/fabro-types/src/run_projection.rs +++ b/lib/foundation/fabro-types/src/run_projection.rs @@ -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, + /// 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, /// 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,