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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-13 07:58:34 -06:00 • committed by Bryan Helmkamp
parent 9101a90471
commit 0c0e589a78
12 changed files with 311 additions and 105 deletions

View file

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

View file

@ -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<AppState>, run_id: RunId) {
"verify",
1,
&workflow_event::Event::StageFailed {
node_id: "verify".to_string(),
name: "Verify".to_string(),
index: 1,
failure: FailureDetail::new("try again", FailureCategory::TransientInfra),
will_retry: true,
timing: fabro_types::StageTiming::wall_only(1200),
billing: Some(test_billed_usage("gpt-old", 100, 10)),
actor: None,
node_id: "verify".to_string(),
name: "Verify".to_string(),
index: 1,
failure: FailureDetail::new("try again", FailureCategory::TransientInfra),
will_retry: true,
timing: fabro_types::StageTiming::wall_only(1200),
billing_by_model: Vec::new(),
billing: Some(test_billed_usage("gpt-old", 100, 10)),
actor: None,
},
)
.await;
@ -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,
},
];

View file

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

View file

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

View file

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

View file

@ -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<BilledModelUsage>,
node_id: String,
name: String,
index: usize,
failure: FailureDetail,
will_retry: bool,
timing: StageTiming,
billing: Option<BilledModelUsage>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
billing_by_model: Vec<BilledModelUsage>,
#[serde(default, skip_serializing_if = "Option::is_none")]
actor: Option<Principal>,
actor: Option<Principal>,
},
StageRetrying {
node_id: String,

View file

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

View file

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

View file

@ -291,6 +291,7 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
will_retry: true,
timing,
billing: outcome.usage.clone(),
billing_by_model: outcome.usage_by_model.clone(),
actor,
},
&scope,
@ -343,6 +344,7 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
will_retry: false,
timing,
billing: outcome.usage.clone(),
billing_by_model: outcome.usage_by_model.clone(),
actor,
},
&scope,

View file

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

View file

@ -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<FailureDetail>,
pub will_retry: bool,
pub failure: Option<FailureDetail>,
pub will_retry: bool,
/// Per-attempt timing breakdown for this stage visit.
#[serde(default)]
pub timing: StageTiming,
pub timing: StageTiming,
/// The stage's billing: for an agent stage that failed after spending,
/// the whole session tree's tokens under the root's route.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub billing: Option<BilledModelUsage>,
pub billing: Option<BilledModelUsage>,
/// `billing` split by model, as on `stage.completed`.
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub billing_by_model: Vec<BilledModelUsage>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]

View file

@ -297,10 +297,11 @@ pub struct StageProjection {
pub usage: BilledTokenCounts,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model: Option<ModelRef>,
/// The completed stage's billing split by model, as `stage.completed`
/// reported it: the root session's route and each subagent's own model.
/// Sums to `usage`. Empty while the stage runs and for stages without a
/// coding agent; the billing rollup then bills `usage` to `model`.
/// The finished stage's billing split by model, as `stage.completed` or
/// `stage.failed` reported it: the root session's route and each
/// subagent's own model. Sums to `usage`. Empty while the stage runs and
/// for stages without a coding agent; the billing rollup then bills
/// `usage` to `model`.
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub billing_by_model: Vec<BilledModelUsage>,
/// Todo/task list owned by the stage's root agent session.