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,