From 9b15a9cd3c9f48e62e0956c2ae19656e240b8974 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 4 May 2026 19:02:52 -0400 Subject: [PATCH] fix(server): dedup retried nodes in run billing MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `checkpoint.completed_nodes` records every visit (one entry per re-entry of a looped node), so the billing handler emitted duplicate rows and multiplied `runtime_secs` by visit count. Dedup by node before iterating; duration is now correctly summed across visits exactly once. Per-visit usage isn't reachable in the current data model — `node_outcomes` is keyed by `node_id` and overwritten on each visit — so the billing row's token counts still reflect only the latest visit. Co-Authored-By: Claude Opus 4.7 (1M context) --- .../src/server/handler/billing.rs | 10 ++ lib/crates/fabro-server/src/server/tests.rs | 146 +++++++++++++++++- 2 files changed, 155 insertions(+), 1 deletion(-) diff --git a/lib/crates/fabro-server/src/server/handler/billing.rs b/lib/crates/fabro-server/src/server/handler/billing.rs index 7861ba824..8482bfce0 100644 --- a/lib/crates/fabro-server/src/server/handler/billing.rs +++ b/lib/crates/fabro-server/src/server/handler/billing.rs @@ -1,3 +1,4 @@ +use std::collections::HashSet; use std::num::NonZeroU32; use std::sync::Arc; @@ -165,7 +166,16 @@ async fn get_run_billing( let mut runtime_secs = 0.0_f64; let mut stages = Vec::new(); + // `completed_nodes` records every visit (one entry per re-entry of a + // looped node), but billing is per-node: the duration helper already sums + // across visits, and `node_outcomes` only stores the latest visit's usage. + // Dedup so we emit one row per node and don't multiply the sum by visit + // count. + let mut seen_nodes = HashSet::new(); for node_id in &checkpoint.completed_nodes { + if !seen_nodes.insert(node_id.as_str()) { + continue; + } let duration_ms = stage_durations.get(node_id).copied().unwrap_or(0); runtime_secs += duration_ms as f64 / 1000.0; diff --git a/lib/crates/fabro-server/src/server/tests.rs b/lib/crates/fabro-server/src/server/tests.rs index 129333865..093e0f858 100644 --- a/lib/crates/fabro-server/src/server/tests.rs +++ b/lib/crates/fabro-server/src/server/tests.rs @@ -18,7 +18,8 @@ use fabro_model::Provider; use fabro_types::settings::ServerAuthMethod; use fabro_types::{ AttrValue, AuthMethod, CommandTermination, FailureCategory, FailureDetail, Graph, - InterviewQuestionRecord, QuestionType, RunBlobId, RunId, RunSpec, SystemActorKind, fixtures, + InterviewQuestionRecord, Outcome, QuestionType, RunBlobId, RunId, RunSpec, SystemActorKind, + fixtures, }; use fabro_util::check_report::CheckStatus; use httpmock::Method::{GET, POST}; @@ -2428,6 +2429,149 @@ async fn list_run_stages_distinguishes_visits() { assert!(first.get("dot_id").is_none(), "dot_id should be removed"); } +/// `checkpoint.completed_nodes` records every visit, so a looped node appears +/// once per re-entry. Billing must dedup so a retried node renders as one row +/// and `runtime_secs` is summed across all visits exactly once. +#[tokio::test] +async fn run_billing_dedups_retried_nodes_and_sums_their_durations() { + let state = test_app_state_with_isolated_storage(); + let app = crate::test_support::build_test_router(Arc::clone(&state)); + let run_id = RunId::new(); + + create_durable_run_with_events(&state, run_id, &[ + workflow_event::Event::RunSubmitted { + definition_blob: None, + }, + workflow_event::Event::RunStarting, + workflow_event::Event::RunRunning, + ]) + .await; + + // Visit 1 of `verify` — completed in 1.5s. + append_scoped_stage_event( + &state, + run_id, + "verify", + 1, + &workflow_event::Event::StageCompleted { + node_id: "verify".to_string(), + name: "Verify".to_string(), + index: 1, + duration_ms: 1500, + status: "failed".to_string(), + preferred_label: None, + suggested_next_ids: Vec::new(), + billing: None, + failure: None, + notes: None, + files_touched: Vec::new(), + context_updates: None, + jump_to_node: None, + context_values: None, + node_visits: None, + loop_failure_signatures: None, + restart_failure_signatures: None, + response: None, + attempt: 1, + max_attempts: 1, + }, + ) + .await; + + // Visit 2 of `verify` — completed in 0.8s. + append_scoped_stage_event( + &state, + run_id, + "verify", + 2, + &workflow_event::Event::StageCompleted { + node_id: "verify".to_string(), + name: "Verify".to_string(), + index: 1, + duration_ms: 800, + status: "succeeded".to_string(), + preferred_label: None, + suggested_next_ids: Vec::new(), + billing: None, + failure: None, + notes: None, + files_touched: Vec::new(), + context_updates: None, + jump_to_node: None, + context_values: None, + node_visits: None, + loop_failure_signatures: None, + restart_failure_signatures: None, + response: None, + attempt: 1, + max_attempts: 1, + }, + ) + .await; + + // Checkpoint records `verify` twice (once per visit) — this is what makes + // the dedup necessary. + let run_store = state.store.open_run(&run_id).await.unwrap(); + workflow_event::append_event( + &run_store, + &run_id, + &workflow_event::Event::CheckpointCompleted { + node_id: "verify".to_string(), + status: "running".to_string(), + current_node: "verify".to_string(), + completed_nodes: vec!["verify".to_string(), "verify".to_string()], + node_retries: std::collections::BTreeMap::new(), + context_values: std::collections::BTreeMap::new(), + node_outcomes: std::collections::BTreeMap::from([( + "verify".to_string(), + Outcome::default(), + )]), + next_node_id: Some("done".to_string()), + git_commit_sha: None, + loop_failure_signatures: std::collections::BTreeMap::new(), + restart_failure_signatures: std::collections::BTreeMap::new(), + node_visits: std::collections::BTreeMap::from([("verify".to_string(), 2usize)]), + diff: None, + }, + ) + .await + .unwrap(); + + let response = app + .clone() + .oneshot( + Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}/billing"))) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + let body = response_json!(response, StatusCode::OK).await; + + let stages = body["stages"].as_array().unwrap(); + assert_eq!( + stages.len(), + 1, + "expected one row for the retried verify node" + ); + assert_eq!(stages[0]["stage"]["id"], "verify"); + // Duration on the row is the sum across visits (1.5s + 0.8s = 2.3s). + assert!( + (stages[0]["runtime_secs"].as_f64().unwrap() - 2.3).abs() < f64::EPSILON, + "row runtime_secs should sum visits, got {}", + stages[0]["runtime_secs"] + ); + + // Totals must not double-count: a single 2.3s, not 4.6s. + assert!( + (body["totals"]["runtime_secs"].as_f64().unwrap() - 2.3).abs() < f64::EPSILON, + "totals.runtime_secs should sum visits exactly once, got {}", + body["totals"]["runtime_secs"] + ); +} + #[tokio::test] async fn list_run_stages_shows_retrying_after_failed_event() { let state = test_app_state_with_isolated_storage();