From 0c0e589a78124fbb4ea85d7b6125664e617aadd9 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sun, 13 Sep 2026 07:58:34 -0600 Subject: [PATCH] Bill a failed agent stage what it spent An agent stage that failed billed nothing: the backend returned a bare error and the outcome built from it carried no usage. A terminal failure now becomes the stage's failed outcome from the same fold that bills a completed stage, with the tree's usage, the rows by model, the files it wrote, and its active time; stage.failed carries billing and billing_by_model and the store keeps both. Cancellation and retryable failures still go up as the error. Co-Authored-By: Claude Fable 5.1 --- docs/internal/events.md | 2 + lib/apps/fabro-server/src/server/tests.rs | 102 +++++++------- lib/components/fabro-store/src/run_state.rs | 25 ++-- lib/components/fabro-workflow/src/error.rs | 17 +-- .../fabro-workflow/src/event/convert.rs | 19 +-- .../fabro-workflow/src/event/events.rs | 18 +-- .../fabro-workflow/src/handler/llm/pebble.rs | 68 ++++++++-- lib/components/fabro-workflow/src/lib.rs | 11 +- .../fabro-workflow/src/lifecycle/event.rs | 2 + .../fabro-workflow/tests/it/pebble_agent.rs | 128 ++++++++++++++++++ .../fabro-types/src/run_event/stage.rs | 15 +- .../fabro-types/src/run_projection.rs | 9 +- 12 files changed, 311 insertions(+), 105 deletions(-) diff --git a/docs/internal/events.md b/docs/internal/events.md index b419d9396..add8b9411 100644 --- a/docs/internal/events.md +++ b/docs/internal/events.md @@ -473,6 +473,8 @@ Emitted when a stage fails (before retry decision). | `failure_class` | string | Failure category | | `failure_signature` | string? | Dedup key for repeated failures | | `will_retry` | boolean | Whether the stage will be retried | +| `billing` | object? | What the stage spent before it failed, in the shape `stage.completed` uses. An agent stage that fails for good after answering model calls bills its whole session tree, as it would have on completion; a retried attempt and a cancelled stage carry none | +| `billing_by_model` | array? | `billing` split by model, as on `stage.completed` | ### `stage.retrying` diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index 8f9ef998b..2fd594239 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -7104,14 +7104,15 @@ async fn list_run_stages_projects_retrying_until_completion() { "work", 1, &workflow_event::Event::StageFailed { - node_id: "work".to_string(), - name: "Work".to_string(), - index: 1, - failure: FailureDetail::new("try again", FailureCategory::TransientInfra), - will_retry: true, - timing: fabro_types::StageTiming::wall_only(10), - billing: None, - actor: None, + node_id: "work".to_string(), + name: "Work".to_string(), + index: 1, + failure: FailureDetail::new("try again", FailureCategory::TransientInfra), + will_retry: true, + timing: fabro_types::StageTiming::wall_only(10), + billing_by_model: Vec::new(), + billing: None, + actor: None, }, ) .await; @@ -7373,14 +7374,15 @@ async fn create_billed_retry_run(state: &Arc, run_id: RunId) { "verify", 1, &workflow_event::Event::StageFailed { - node_id: "verify".to_string(), - name: "Verify".to_string(), - index: 1, - failure: FailureDetail::new("try again", FailureCategory::TransientInfra), - will_retry: true, - timing: fabro_types::StageTiming::wall_only(1200), - billing: Some(test_billed_usage("gpt-old", 100, 10)), - actor: None, + node_id: "verify".to_string(), + name: "Verify".to_string(), + index: 1, + failure: FailureDetail::new("try again", FailureCategory::TransientInfra), + will_retry: true, + timing: fabro_types::StageTiming::wall_only(1200), + billing_by_model: Vec::new(), + billing: Some(test_billed_usage("gpt-old", 100, 10)), + actor: None, }, ) .await; @@ -8118,14 +8120,15 @@ async fn list_run_stages_shows_retrying_after_failed_event() { "work", 1, &workflow_event::Event::StageFailed { - node_id: "work".to_string(), - name: "Work".to_string(), - index: 0, - failure: FailureDetail::new("flake", FailureCategory::TransientInfra), - will_retry: true, - timing: fabro_types::StageTiming::wall_only(5), - billing: None, - actor: None, + node_id: "work".to_string(), + name: "Work".to_string(), + index: 0, + failure: FailureDetail::new("flake", FailureCategory::TransientInfra), + will_retry: true, + timing: fabro_types::StageTiming::wall_only(5), + billing_by_model: Vec::new(), + billing: None, + actor: None, }, ) .await; @@ -8200,14 +8203,15 @@ async fn list_run_stages_shows_retrying_when_failed_will_retry() { "work", 1, &workflow_event::Event::StageFailed { - node_id: "work".to_string(), - name: "Work".to_string(), - index: 0, - failure: FailureDetail::new("flake", FailureCategory::TransientInfra), - will_retry: true, - timing: fabro_types::StageTiming::wall_only(5), - billing: None, - actor: None, + node_id: "work".to_string(), + name: "Work".to_string(), + index: 0, + failure: FailureDetail::new("flake", FailureCategory::TransientInfra), + will_retry: true, + timing: fabro_types::StageTiming::wall_only(5), + billing_by_model: Vec::new(), + billing: None, + actor: None, }, ) .await; @@ -8250,14 +8254,15 @@ async fn run_billing_retried_node_then_succeeded_emits_one_row_with_final_attemp max_attempts: 3, }, workflow_event::Event::StageFailed { - node_id: "work".to_string(), - name: "Work".to_string(), - index: 0, - failure: FailureDetail::new("transient", FailureCategory::TransientInfra), - will_retry: true, - timing: fabro_types::StageTiming::wall_only(10), - billing: None, - actor: None, + node_id: "work".to_string(), + name: "Work".to_string(), + index: 0, + failure: FailureDetail::new("transient", FailureCategory::TransientInfra), + will_retry: true, + timing: fabro_types::StageTiming::wall_only(10), + billing_by_model: Vec::new(), + billing: None, + actor: None, }, workflow_event::Event::StageRetrying { node_id: "work".to_string(), @@ -15427,14 +15432,15 @@ async fn active_acp_steerable_marker_clears_on_terminal_paths() { max_attempts: 1, }, workflow_event::Event::StageFailed { - node_id: "agent".to_string(), - name: "agent".to_string(), - index: 0, - failure: FailureDetail::new("failed", FailureCategory::Deterministic), - will_retry: false, - timing: fabro_types::StageTiming::wall_only(1), - billing: None, - actor: None, + node_id: "agent".to_string(), + name: "agent".to_string(), + index: 0, + failure: FailureDetail::new("failed", FailureCategory::Deterministic), + will_retry: false, + timing: fabro_types::StageTiming::wall_only(1), + billing_by_model: Vec::new(), + billing: None, + actor: None, }, ]; diff --git a/lib/components/fabro-store/src/run_state.rs b/lib/components/fabro-store/src/run_state.rs index afe270b99..df6d0184b 100644 --- a/lib/components/fabro-store/src/run_state.rs +++ b/lib/components/fabro-store/src/run_state.rs @@ -549,6 +549,7 @@ impl RunProjectionReducer for RunProjection { stage.usage.replace_with_billed_usage(billing); stage.model = Some(billing.model().clone()); } + stage.billing_by_model.clone_from(&props.billing_by_model); stage.state = stage_state_from_failure(props.will_retry, failure_category, stage.termination); stage.agent_control = AgentControlState::Running; @@ -3841,14 +3842,15 @@ mod tests { .apply_event(&test_stage_event( 3, EventBody::StageFailed(StageFailedProps { - index: 0, - failure: Some(fabro_types::FailureDetail::new( + index: 0, + failure: Some(fabro_types::FailureDetail::new( "try again", fabro_types::FailureCategory::TransientInfra, )), - will_retry: true, - timing: fabro_types::StageTiming::wall_only(444), - billing: Some(usage.clone()), + will_retry: true, + timing: fabro_types::StageTiming::wall_only(444), + billing_by_model: Vec::new(), + billing: Some(usage.clone()), }), scoped_stage_id.clone(), )) @@ -5467,6 +5469,7 @@ mod tests { failure: Some(FailureDetail::new("boom", FailureCategory::TransientInfra)), will_retry, timing: fabro_types::StageTiming::wall_only(duration_ms), + billing_by_model: Vec::new(), billing: None, } } @@ -5477,6 +5480,7 @@ mod tests { failure: Some(FailureDetail::new("cancelled", FailureCategory::Canceled)), will_retry, timing: fabro_types::StageTiming::wall_only(duration_ms), + billing_by_model: Vec::new(), billing: None, } } @@ -6088,14 +6092,15 @@ mod tests { .apply_event(&test_event( 3, EventBody::StageFailed(StageFailedProps { - index: 0, - failure: Some(FailureDetail::new( + index: 0, + failure: Some(FailureDetail::new( "Script failed with exit code: 100\n\nCancelling due to test failure", FailureCategory::Canceled, )), - will_retry: false, - timing: fabro_types::StageTiming::wall_only(10), - billing: None, + will_retry: false, + timing: fabro_types::StageTiming::wall_only(10), + billing_by_model: Vec::new(), + billing: None, }), Some("build"), )) diff --git a/lib/components/fabro-workflow/src/error.rs b/lib/components/fabro-workflow/src/error.rs index 40f271ad1..a0e85035b 100644 --- a/lib/components/fabro-workflow/src/error.rs +++ b/lib/components/fabro-workflow/src/error.rs @@ -2128,14 +2128,15 @@ mod tests { // 3. Outcome → StageFailed event let failure = outcome.failure.clone().unwrap(); let event = Event::StageFailed { - node_id: "code".into(), - name: "code".into(), - index: 0, - failure: failure.clone(), - will_retry: false, - timing: fabro_types::StageTiming::wall_only(0), - billing: None, - actor: None, + node_id: "code".into(), + name: "code".into(), + index: 0, + failure: failure.clone(), + will_retry: false, + timing: fabro_types::StageTiming::wall_only(0), + billing_by_model: Vec::new(), + billing: None, + actor: None, }; // 4. Verify classification survived all the way through diff --git a/lib/components/fabro-workflow/src/event/convert.rs b/lib/components/fabro-workflow/src/event/convert.rs index fa0bb6d61..d2f733f2c 100644 --- a/lib/components/fabro-workflow/src/event/convert.rs +++ b/lib/components/fabro-workflow/src/event/convert.rs @@ -387,6 +387,7 @@ fn event_body_from_event(event: &Event) -> EventBody { will_retry, timing, billing, + billing_by_model, .. } => EventBody::StageFailed(fabro_types::StageFailedProps { index: *index, @@ -394,6 +395,7 @@ fn event_body_from_event(event: &Event) -> EventBody { will_retry: *will_retry, timing: *timing, billing: billing.clone(), + billing_by_model: billing_by_model.clone(), }), Event::StageRetrying { index, @@ -1171,17 +1173,18 @@ mod tests { fn run_event_stage_failure_keeps_failure_detail() { let usage = test_usage("gpt-5.2", 321, 54); let stored = to_run_event(&fixtures::RUN_3, &Event::StageFailed { - node_id: "code".to_string(), - name: "Code".to_string(), - index: 1, - failure: FailureDetail::new( + node_id: "code".to_string(), + name: "Code".to_string(), + index: 1, + failure: FailureDetail::new( "lint failed", crate::outcome::FailureCategory::Deterministic, ), - will_retry: true, - timing: ::fabro_types::StageTiming::wall_only(5000), - billing: Some(usage.clone()), - actor: None, + will_retry: true, + timing: ::fabro_types::StageTiming::wall_only(5000), + billing_by_model: Vec::new(), + billing: Some(usage.clone()), + actor: None, }); assert_eq!(stored.event_name(), "stage.failed"); diff --git a/lib/components/fabro-workflow/src/event/events.rs b/lib/components/fabro-workflow/src/event/events.rs index 57a639a38..68618d764 100644 --- a/lib/components/fabro-workflow/src/event/events.rs +++ b/lib/components/fabro-workflow/src/event/events.rs @@ -296,15 +296,17 @@ pub enum Event { max_attempts: usize, }, StageFailed { - node_id: String, - name: String, - index: usize, - failure: FailureDetail, - will_retry: bool, - timing: StageTiming, - billing: Option, + node_id: String, + name: String, + index: usize, + failure: FailureDetail, + will_retry: bool, + timing: StageTiming, + billing: Option, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + billing_by_model: Vec, #[serde(default, skip_serializing_if = "Option::is_none")] - actor: Option, + actor: Option, }, StageRetrying { node_id: String, diff --git a/lib/components/fabro-workflow/src/handler/llm/pebble.rs b/lib/components/fabro-workflow/src/handler/llm/pebble.rs index fffb259db..6d7686fba 100644 --- a/lib/components/fabro-workflow/src/handler/llm/pebble.rs +++ b/lib/components/fabro-workflow/src/handler/llm/pebble.rs @@ -63,7 +63,7 @@ use crate::context::keys::Fidelity; use crate::error::Error; use crate::event::{Emitter, Event, StageScope}; use crate::model_fallback::{ModelFallbackNotice, ModelFallbackPolicy}; -use crate::outcome::billed_model_usage_from_llm; +use crate::outcome::{Outcome, billed_model_usage_from_llm}; use crate::services::FabroRunToolServices; use crate::steering_hub::SteeringHub; use crate::web_search::{self, SearchSecrets}; @@ -407,6 +407,16 @@ impl LiveAgent { } } +/// The route as billing names it: provider, model, and the speed tier the +/// stage asked for. +fn route_model(route: &LlmRoute) -> ModelRef { + ModelRef::new( + route.target.provider.clone(), + ModelId::new(route.target.model.as_str()), + ) + .with_speed(route.controls.speed) +} + /// A stage's billing from its account: the whole tree under the root's /// route, and the rows that split it by model. struct StageBilling { @@ -856,6 +866,37 @@ impl PebbleBackend { } } + /// The failed outcome of an agent stage that spent before it failed: the + /// failure itself, with the session tree's usage, the files it wrote, and + /// its active time, so the run bills what the stage spent. A billing the + /// catalog cannot price is logged and left off. + fn failed_outcome(&self, error: &Error, live: &LiveAgent, plan: &FallbackPlan) -> Outcome { + let mut outcome = error.to_fail_outcome(); + let account = live.account(); + match stage_billing( + self.catalog.as_ref(), + &route_model(plan.current()), + &account, + ) { + Ok(billing) => { + outcome.usage = Some(billing.total); + outcome.usage_by_model = billing.by_model; + } + Err(billing_error) => { + tracing::debug!( + error = %billing_error, + "failed agent stage could not be billed" + ); + } + } + outcome.files_touched = account.files_touched; + outcome.timing = Some(StageTiming::active_only( + crate::millis_u64(live.inference_duration), + crate::millis_u64(live.tool_duration), + )); + outcome + } + /// Steers that landed between the answer and the hub's close-the-door /// check run as further prompts, so the stage never ends with a steer /// nobody saw. @@ -1291,18 +1332,27 @@ impl CodergenBackend for PebbleBackend { ShutdownReason::Error }; live.discard(reason).await; - return Err(error); + // Cancellation and a retryable failure go up as the error, so + // the engine cancels or retries as before. A terminal failure + // becomes the stage's failed outcome, carrying what the + // session tree spent and wrote before it failed. + if matches!(error, Error::Cancelled) || error.is_retryable() { + return Err(error); + } + return Ok(CodergenResult::Full(Box::new(self.failed_outcome( + &error, + &live, + &fallback_plan, + )))); } }; - let route = fallback_plan.current().clone(); - let 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)?; + let billing = stage_billing( + self.catalog.as_ref(), + &route_model(fallback_plan.current()), + &account, + )?; live.release_lease(); match reuse_key { diff --git a/lib/components/fabro-workflow/src/lib.rs b/lib/components/fabro-workflow/src/lib.rs index 4b66dea59..94f822240 100644 --- a/lib/components/fabro-workflow/src/lib.rs +++ b/lib/components/fabro-workflow/src/lib.rs @@ -159,11 +159,12 @@ mod duration_tests { tool_call_id: None, actor: None, body: EventBody::StageFailed(StageFailedProps { - index: 0, - failure: None, - will_retry: true, - timing: StageTiming::wall_only(wall_time_ms), - billing: None, + index: 0, + failure: None, + will_retry: true, + timing: StageTiming::wall_only(wall_time_ms), + billing_by_model: Vec::new(), + billing: None, }), }; EventEnvelope { seq, event } diff --git a/lib/components/fabro-workflow/src/lifecycle/event.rs b/lib/components/fabro-workflow/src/lifecycle/event.rs index 278479006..cf0bda9c3 100644 --- a/lib/components/fabro-workflow/src/lifecycle/event.rs +++ b/lib/components/fabro-workflow/src/lifecycle/event.rs @@ -291,6 +291,7 @@ impl RunLifecycle for EventLifecycle { will_retry: true, timing, billing: outcome.usage.clone(), + billing_by_model: outcome.usage_by_model.clone(), actor, }, &scope, @@ -343,6 +344,7 @@ impl RunLifecycle for EventLifecycle { will_retry: false, timing, billing: outcome.usage.clone(), + billing_by_model: outcome.usage_by_model.clone(), actor, }, &scope, diff --git a/lib/components/fabro-workflow/tests/it/pebble_agent.rs b/lib/components/fabro-workflow/tests/it/pebble_agent.rs index b749ad120..7eae5bd9b 100644 --- a/lib/components/fabro-workflow/tests/it/pebble_agent.rs +++ b/lib/components/fabro-workflow/tests/it/pebble_agent.rs @@ -742,6 +742,134 @@ async fn the_stage_timeout_fails_a_slow_agent() { ); } +/// A stage whose agent fails for good after answering model calls bills +/// those calls: the failed outcome carries the session tree's usage from the +/// same fold the completed outcome would have, and the files it wrote. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_stage_that_fails_after_spending_bills_what_it_spent() { + let stage = Stage::new().await; + let first = stage.file("first.txt"); + let second = stage.file("second.txt"); + // Two answered calls, each writing a file; the third is refused for good. + stage + .server + .mock_async(|when, then| { + when.method(POST) + .path(CHAT_PATH) + .body_excludes(TOOL_RESULT_MARKER); + sse_headers( + then, + sse_tool_call( + "call-1", + "write_file", + &serde_json::json!({ "file_path": first, "content": "one" }), + ), + ); + }) + .await; + stage + .server + .mock_async(|when, then| { + when.method(POST) + .path(CHAT_PATH) + .body_includes("call-1") + .body_excludes("call-2"); + sse_headers( + then, + sse_tool_call( + "call-2", + "write_file", + &serde_json::json!({ "file_path": second, "content": "two" }), + ), + ); + }) + .await; + stage + .server + .mock_async(|when, then| { + when.method(POST).path(CHAT_PATH).body_includes("call-2"); + then.status(400) + .header("content-type", "application/json") + .body(r#"{"error":{"message":"the request was rejected","type":"invalid_request_error"}}"#); + }) + .await; + + let mut graph = agent_graph("Spent", "Write two files"); + let work = graph.nodes.get_mut("work").unwrap(); + work.attrs + .insert("max_retries".to_string(), AttrValue::Integer(0)); + graph.edges.retain(|edge| edge.from != "work"); + let mut fail_edge = Edge::new("work", "exit"); + fail_edge.attrs.insert( + "condition".to_string(), + AttrValue::String("outcome=failed".to_string()), + ); + graph.edges.push(fail_edge); + + let backend = stage.backend("openai"); + let (_, state) = stage + .run(backend, &graph, CancellationToken::new()) + .await + .expect("the fail edge carries the run to exit"); + + let work = work_stage(&state); + assert_eq!( + work.completion + .as_ref() + .expect("the stage finished") + .outcome, + StageOutcome::Failed { + retry_requested: false, + } + ); + assert_eq!( + work.usage.input_tokens, + 2 * INPUT_TOKENS_PER_CALL, + "the two answered calls are billed" + ); + assert_eq!(work.usage.output_tokens, 2 * OUTPUT_TOKENS_PER_CALL); + assert_eq!( + work.usage.total_usd_micros, + Some(2 * (INPUT_TOKENS_PER_CALL + 2 * OUTPUT_TOKENS_PER_CALL)), + "priced from the catalog like a completed stage" + ); + assert_eq!( + work.billing_by_model.len(), + 1, + "{:?}", + work.billing_by_model + ); + assert_eq!( + work.billing_by_model[0].tokens.input, + u64::try_from(2 * INPUT_TOKENS_PER_CALL).unwrap() + ); + assert!( + tokio::fs::try_exists(&second).await.unwrap(), + "the second write landed before the failure" + ); + + let failed = stage + .events + .lock() + .unwrap() + .iter() + .find(|event| { + event.event_name() == "stage.failed" && event.node_id.as_deref() == Some("work") + }) + .cloned() + .expect("the stage failure is emitted"); + let EventBody::StageFailed(props) = &failed.body else { + panic!("stage.failed carries its props: {failed:?}"); + }; + assert!(!props.will_retry); + let billing = props.billing.as_ref().expect("the failed stage is billed"); + assert_eq!( + billing.tokens.input, + u64::try_from(2 * INPUT_TOKENS_PER_CALL).unwrap() + ); + assert_eq!(props.billing_by_model, vec![billing.clone()]); +} + // --- Questions, subagents, MCP // -------------------------------------------------- diff --git a/lib/foundation/fabro-types/src/run_event/stage.rs b/lib/foundation/fabro-types/src/run_event/stage.rs index 524fff69f..af97b9299 100644 --- a/lib/foundation/fabro-types/src/run_event/stage.rs +++ b/lib/foundation/fabro-types/src/run_event/stage.rs @@ -72,15 +72,20 @@ pub struct StageCompletedProps { #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct StageFailedProps { - pub index: usize, + pub index: usize, #[serde(default, skip_serializing_if = "Option::is_none")] - pub failure: Option, - pub will_retry: bool, + pub failure: Option, + pub will_retry: bool, /// Per-attempt timing breakdown for this stage visit. #[serde(default)] - pub timing: StageTiming, + pub timing: StageTiming, + /// The stage's billing: for an agent stage that failed after spending, + /// the whole session tree's tokens under the root's route. #[serde(default, skip_serializing_if = "Option::is_none")] - pub billing: Option, + pub billing: Option, + /// `billing` split by model, as on `stage.completed`. + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub billing_by_model: Vec, } #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] diff --git a/lib/foundation/fabro-types/src/run_projection.rs b/lib/foundation/fabro-types/src/run_projection.rs index 5e0540d78..ed7ccb9ad 100644 --- a/lib/foundation/fabro-types/src/run_projection.rs +++ b/lib/foundation/fabro-types/src/run_projection.rs @@ -297,10 +297,11 @@ 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`. + /// The finished stage's billing split by model, as `stage.completed` or + /// `stage.failed` reported it: the root session's route and each + /// subagent's own model. Sums to `usage`. Empty while the stage runs and + /// for stages without a coding agent; the billing rollup then bills + /// `usage` to `model`. #[serde(default, skip_serializing_if = "Vec::is_empty")] pub billing_by_model: Vec, /// Todo/task list owned by the stage's root agent session.