diff --git a/litellm-rust/crates/lens/tests/receiver.rs b/litellm-rust/crates/lens/tests/receiver.rs index bda1cecbbb3..2af3c67955a 100644 --- a/litellm-rust/crates/lens/tests/receiver.rs +++ b/litellm-rust/crates/lens/tests/receiver.rs @@ -10,7 +10,10 @@ use rstest::rstest; use serde_json::json; use sha2::{Digest, Sha256}; use std::{ - sync::{Arc, atomic::Ordering}, + sync::{ + Arc, + atomic::{AtomicBool, Ordering}, + }, time::Duration, }; use wiremock::{ @@ -116,6 +119,107 @@ async fn agent_picker_query_preserves_scope_through_the_internal_read_route() { assert_eq!(response.json::().await.unwrap(), result); } +#[rstest] +#[case::list_paid("list", None, 0.75)] +#[case::list_zero("list", None, 0.0)] +#[case::detail_paid("trace", None, 0.75)] +#[case::paged_zero("trace", Some(1), 0.0)] +#[tokio::test] +async fn internal_reads_refresh_delayed_gateway_amounts( + #[case] operation: &str, + #[case] page_size: Option, + #[case] cost: f64, +) { + let store = MockServer::start().await; + let start_ms = (unix_seconds() as i64 - 600) * 1000; + let available = Arc::new(AtomicBool::new(false)); + Mock::given(body_string_contains("FROM agent_traces_by_key")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({"data": [{ + "trace_id": "trace", "trace_ref": "ref", "team_id": "team", + "api_key_hash": "key", "user_id": "owner", "name": "model call", + "service": "agent", "input_preview": "", "status": "STATUS_CODE_OK", + "start_ms": start_ms, "duration_ms": 1, "span_count": 1, + "agent_count": 0, "agent_invocations": 0, "llm_calls": 1, + "tool_calls": 0, "input_tokens": 1, "output_tokens": 1, + "models": ["test-model"], "error_count": 0, "request_ids": [] + }]}))) + .mount(&store) + .await; + Mock::given(body_string_contains("o.SpanId AS span_id")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({"data": [{ + "trace_id": "trace", "span_id": "span", "parent_span_id": "", + "name": "model call", "type": "llm", "agent": "", + "status": "STATUS_CODE_OK", "status_message": "", "error_truncated": 0, + "start_ns": start_ms * 1_000_000, "duration_ns": 1_000_000, + "service": "agent", "input_preview": "", "model": "test-model", + "input_tokens": 1, "output_tokens": 1, "litellm_request_id": "", + "call_keys": ["provider_response:response"], "call_evidence": "complete", + "team_id": "team", "api_key_hash": "key", "user_id": "owner" + }]}))) + .expect(2) + .mount(&store) + .await; + let spend_available = available.clone(); + Mock::given(body_string_contains("FROM spend_logs FINAL")) + .respond_with(move |_: &wiremock::Request| { + let rows = if spend_available.load(Ordering::Acquire) { + json!([{ + "request_id": "request", "litellm_call_id": "", + "response_id": "response", "upstream_response_id": "", + "trace_id": "", "span_id": "", "team_id": "team", + "api_key": "key", "user": "owner", "spend": cost, + "start_ms": start_ms + }]) + } else { + json!([]) + }; + ResponseTemplate::new(200).set_body_json(json!({"data": rows})) + }) + .expect(2) + .mount(&store) + .await; + let server = serve(&store.uri(), true).await; + let scope = json!({"all_teams": 0, "user_id": "owner", "team_ids": []}); + let (request, summary_path) = if operation == "list" { + ( + json!({ + "operation": operation, "scope": scope, "start_ms": start_ms, + "end_ms": start_ms + 1000, "cursor": null, "limit": 50 + }), + "/data/0", + ) + } else { + ( + json!({ + "operation": operation, "scope": scope, "trace_id": "trace", + "trace_ref": "ref", "cursor": null, "page_size": page_size + }), + "/summary", + ) + }; + let client = http_client().unwrap(); + for expected in [None, Some(cost)] { + if expected.is_some() { + available.store(true, Ordering::Release); + tokio::time::sleep(litellm_traces_cache::LIVE_TTL + Duration::from_millis(200)).await; + } + let response = client + .post(format!("{}/internal/read", server.url)) + .bearer_auth(SERVICE_TOKEN) + .json(&request) + .send() + .await + .unwrap(); + assert_eq!(response.status(), 200); + let body = response.json::().await.unwrap(); + let summary = body.pointer(summary_path).unwrap(); + assert_eq!(summary["spend"], json!(expected)); + assert_eq!(summary["priced_calls"], u64::from(expected.is_some())); + assert_eq!(summary["llm_calls"], 1); + assert!(!body.to_string().contains("gateway_spend_pending")); + } +} + #[rstest] #[tokio::test] async fn feedback_summary_query_preserves_scope_through_the_internal_read_route() { diff --git a/litellm-rust/crates/traces-cache/src/cache.rs b/litellm-rust/crates/traces-cache/src/cache.rs index eda054e6c7d..e326c0c34d4 100644 --- a/litellm-rust/crates/traces-cache/src/cache.rs +++ b/litellm-rust/crates/traces-cache/src/cache.rs @@ -59,7 +59,7 @@ impl SnapshotKey { } /// Native sessions can resume without a terminal record, so their reads retain `LIVE_TTL`. -/// Other traces with known spend settle after `SETTLED_AFTER_MS` of inactivity. +/// Other traces settle after `SETTLED_AFTER_MS` without activity or recoverable missing spend. #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub enum Freshness { Live, @@ -67,7 +67,7 @@ pub enum Freshness { } impl Freshness { - pub fn of(rows: &[TraceSpansRow], spend_known: bool, snapshot_ms: u64) -> Self { + pub fn of(rows: &[TraceSpansRow], trace: &Trace, snapshot_ms: u64) -> Self { if rows .iter() .any(|row| matches!(row.framework.as_str(), "claude-code" | "claude-agent-sdk")) @@ -82,7 +82,7 @@ impl Freshness { let quiet_ms = i64::try_from(snapshot_ms) .unwrap_or(i64::MAX) .saturating_sub(last_end_ms); - if spend_known && quiet_ms >= SETTLED_AFTER_MS as i64 { + if !trace.gateway_spend_pending && quiet_ms >= SETTLED_AFTER_MS as i64 { Self::Settled } else { Self::Live @@ -388,21 +388,54 @@ mod tests { const LAST_END_MS: u64 = 1_790_742_989_010; #[rstest] - #[case::just_ended(LAST_END_MS, true, Freshness::Live)] - #[case::quiet_just_under(LAST_END_MS + SETTLED_AFTER_MS - 1, true, Freshness::Live)] - #[case::quiet_long_enough(LAST_END_MS + SETTLED_AFTER_MS, true, Freshness::Settled)] - #[case::spend_unknown(LAST_END_MS + SETTLED_AFTER_MS, false, Freshness::Live)] - fn freshness_settles_once_spans_stop_and_spend_is_known( + #[case::just_ended(LAST_END_MS, Freshness::Live)] + #[case::quiet_just_under(LAST_END_MS + SETTLED_AFTER_MS - 1, Freshness::Live)] + #[case::quiet_long_enough(LAST_END_MS + SETTLED_AFTER_MS, Freshness::Settled)] + fn freshness_settles_traces_without_model_calls_once_spans_stop( #[case] snapshot_ms: u64, - #[case] spend_known: bool, #[case] expected: Freshness, ) { assert_eq!( - Freshness::of(&[row("root")], spend_known, snapshot_ms), + Freshness::of(&[row("root")], &trace("root"), snapshot_ms), expected ); } + #[rstest] + #[case::missing_log(None)] + #[case::catalog_estimate(Some(0.25))] + fn estimated_amounts_do_not_settle_pending_gateway_spend(#[case] amount: Option) { + let model = TraceSpansRow { + kind: litellm_traces::ObservationType::Llm, + call_keys: vec![litellm_traces::CallKey::ProviderResponse("response".into())], + call_evidence: Some(litellm_traces::CallEvidenceKind::Complete), + ..row("model") + }; + let rows = [model]; + let original = resolve_trace("trace", "ref", &rows, &[]).unwrap(); + let trace = Trace { + summary: TraceSummary { + llm_calls: 1, + spend: amount, + priced_calls: u64::from(amount.is_some()), + ..original.summary + }, + spans: original + .spans + .into_iter() + .map(|span| litellm_traces::Span { + spend: amount, + ..span + }) + .collect(), + ..original + }; + assert_eq!( + Freshness::of(&rows, &trace, LAST_END_MS + SETTLED_AFTER_MS), + Freshness::Live + ); + } + #[rstest] #[case::claude_code("claude-code")] #[case::claude_agent_sdk("claude-agent-sdk")] @@ -412,7 +445,11 @@ mod tests { ..row("native") }; assert_eq!( - Freshness::of(&[row("root"), native], true, LAST_END_MS + SETTLED_AFTER_MS), + Freshness::of( + &[row("root"), native], + &trace("root"), + LAST_END_MS + SETTLED_AFTER_MS, + ), Freshness::Live ); } diff --git a/litellm-rust/crates/traces-cache/src/list.rs b/litellm-rust/crates/traces-cache/src/list.rs index 1b1dcbbefc3..288f6e19c63 100644 --- a/litellm-rust/crates/traces-cache/src/list.rs +++ b/litellm-rust/crates/traces-cache/src/list.rs @@ -162,10 +162,8 @@ async fn resolve_runs( let spend = spend_window(spans).map_or(&[][..], |window| spend_within(&spend_rows, window)); resolve_trace(&row.trace_id, &row.trace_ref, spans, spend).map(|trace| { - ListedRun::Resolved( - Box::new(trace.summary), - Freshness::of(spans, true, snapshot_ms), - ) + let freshness = Freshness::of(spans, &trace, snapshot_ms); + ListedRun::Resolved(Box::new(trace.summary), freshness) }) }) .collect()) diff --git a/litellm-rust/crates/traces-cache/src/reader.rs b/litellm-rust/crates/traces-cache/src/reader.rs index b0f83a89118..00abbd78601 100644 --- a/litellm-rust/crates/traces-cache/src/reader.rs +++ b/litellm-rust/crates/traces-cache/src/reader.rs @@ -211,14 +211,16 @@ impl TraceReader { .await .map_err(|error| Miss::Read(map_store_error(error)))?; let spend_rows = spend(store, access, &rows).await; - let freshness = Freshness::of(&rows, spend_rows.is_some(), snapshot_ms); resolve_trace( trace_id, trace_ref, &rows, spend_rows.as_deref().unwrap_or_default(), ) - .map(|trace| (trace, freshness)) + .map(|trace| { + let freshness = Freshness::of(&rows, &trace, snapshot_ms); + (trace, freshness) + }) .ok_or(Miss::Absent) }) .await @@ -325,6 +327,7 @@ fn page( let create_page = |count: usize| { let end = position.offset.saturating_add(count).min(spans.len()); Trace { + gateway_spend_pending: snapshot.trace().gateway_spend_pending, summary: snapshot.trace().summary.clone(), agents: snapshot.trace().agents.clone(), spans: spans[position.offset..end].to_vec(), diff --git a/litellm-rust/crates/traces-cache/tests/read.rs b/litellm-rust/crates/traces-cache/tests/read.rs index 22293db98e8..b206f63e58e 100644 --- a/litellm-rust/crates/traces-cache/tests/read.rs +++ b/litellm-rust/crates/traces-cache/tests/read.rs @@ -539,6 +539,110 @@ fn spend_row(response_id: &str, cost: f64) -> SpendByResponseIdsRow { } } +#[rstest] +#[case::missing(&[], None, 0, false, 0.5, None)] +#[case::partial(&[Some(0.25)], Some(0.25), 1, false, 0.5, None)] +#[case::delayed_zero(&[Some(0.25)], Some(0.25), 1, false, 0.0, None)] +#[case::null_amount(&[Some(0.25), None], Some(0.25), 1, false, 0.5, None)] +#[case::complete_zero(&[Some(0.25), Some(0.0)], Some(0.25), 2, true, 0.5, None)] +#[case::no_call_id(&[Some(0.25)], Some(0.25), 1, true, 0.5, Some(CallEvidenceKind::Unknown))] +#[case::incomplete_identity(&[Some(0.25)], Some(0.25), 1, true, 0.5, Some(CallEvidenceKind::Partial))] +#[tokio::test] +async fn gateway_cost_refreshes_until_every_model_call_is_priced( + #[case] initial_costs: &[Option], + #[case] initial_total: Option, + #[case] initial_priced: u64, + #[case] settled: bool, + #[case] final_second: f64, + #[case] terminal: Option, +) { + let rows: Vec<_> = std::iter::once(span(0)) + .chain((1..=2).map(|index| TraceSpansRow { + kind: ObservationType::Llm, + call_keys: vec![CallKey::ProviderResponse(format!("response-{index}"))], + call_evidence: Some(if index == 2 { + terminal.unwrap_or(CallEvidenceKind::Complete) + } else { + CallEvidenceKind::Complete + }), + ..span(index) + })) + .collect(); + let store = FakeStore::with_spans("ref", rows.clone()); + store.set_list_runs(vec![run("trace", "ref")]); + store.set_run_spans(rows); + store.state.lock().unwrap().spend = initial_costs + .iter() + .enumerate() + .map(|(index, cost)| SpendByResponseIdsRow { + spend: *cost, + ..spend_row(&format!("response-{}", index + 1), 0.0) + }) + .collect(); + let reader = TraceReader::new(usize::MAX); + let access = access(); + let detail = reader + .get_trace_page(&store, &access, "trace", "ref", None, 1) + .await + .unwrap() + .unwrap(); + let list = reader + .list_traces(&store, &access, 0, i64::MAX, None, 2) + .await + .unwrap(); + assert_eq!(detail.summary.spend, initial_total); + assert_eq!(detail.summary.priced_calls, initial_priced); + assert_eq!(list.data[0].spend, initial_total); + assert_eq!(list.data[0].priced_calls, initial_priced); + + store.state.lock().unwrap().spend = vec![ + spend_row("response-1", 0.25), + spend_row("response-2", final_second), + ]; + tokio::time::sleep(LIVE_TTL + Duration::from_millis(200)).await; + let refreshed = reader + .get_trace(&store, &access, "trace", "ref") + .await + .unwrap() + .unwrap(); + let listed = reader + .list_traces(&store, &access, 0, i64::MAX, None, 2) + .await + .unwrap(); + let expected = if settled { + initial_total + } else { + Some(0.25 + final_second) + }; + assert_eq!(refreshed.summary.spend, expected); + assert_eq!(listed.data[0].spend, expected); + let expected_priced = if settled { initial_priced } else { 2 }; + assert_eq!(refreshed.summary.priced_calls, expected_priced); + assert_eq!(listed.data[0].priced_calls, expected_priced); + assert_eq!( + store.calls(Operation::TraceSpans), + if settled { 1 } else { 2 } + ); + assert_eq!( + store.calls(Operation::RunSpans), + if settled { 1 } else { 2 } + ); + let pinned = reader + .get_trace_page( + &store, + &access, + "trace", + "ref", + detail.next_cursor.as_deref(), + 1, + ) + .await + .unwrap() + .unwrap(); + assert_eq!(pinned.summary.spend, initial_total); + assert_eq!(pinned.summary.priced_calls, initial_priced); +} + #[rstest] #[tokio::test] async fn failed_batch_spend_lookup_falls_back_to_each_run_instead_of_losing_every_cost() { diff --git a/litellm-rust/crates/traces/src/resolve/resolution.rs b/litellm-rust/crates/traces/src/resolve/resolution.rs index e19a98a1ede..007fc10edf8 100644 --- a/litellm-rust/crates/traces/src/resolve/resolution.rs +++ b/litellm-rust/crates/traces/src/resolve/resolution.rs @@ -31,7 +31,11 @@ pub(super) struct Resolution<'a> { call_matches: HashMap>, } -pub(super) type CallMatch<'a> = (Option>, SpendMatch); +pub(super) struct CallMatch<'a> { + pub(super) requests: Option>, + pub(super) state: SpendMatch, + pub(super) spend_pending: bool, +} impl<'a> Resolution<'a> { pub(super) fn new(rows: &'a [TraceSpansRow], spend: &'a [SpendRow]) -> Self { @@ -138,18 +142,16 @@ impl<'a> Resolution<'a> { pub(super) fn call_requests(&self, call: usize) -> Option> { self.call_match(call) - .and_then(|(requests, _)| requests.clone()) + .and_then(|matched| matched.requests.clone()) + } + + pub(super) fn gateway_spend_pending(&self) -> bool { + self.call_matches + .values() + .any(|matched| matched.spend_pending) } fn resolve_call_match(&self, call: usize) -> CallMatch<'a> { - if let Some(requests) = self.resolve_call_requests(call) { - return (Some(requests), SpendMatch::Matched); - } - let evidence = self.requests(call); - (None, evidence.unmatched_reason()) - } - - fn resolve_call_requests(&self, call: usize) -> Option> { let wrappers = self.graph.ancestors(call).into_iter().filter(|ancestor| { self.kind(*ancestor) == ObservationType::Llm && self @@ -178,7 +180,7 @@ impl<'a> Resolution<'a> { .collect() }) .flatten(); - let selected: Requests<'a> = transport_requests + let selected = transport_requests .map(|requests| requests.into_iter().flatten().collect()) .into_iter() .chain(sources.iter().filter_map(SpendEvidence::complete_requests)) @@ -187,15 +189,24 @@ impl<'a> Resolution<'a> { .iter() .chain(&transports) .all(|source| source.agrees_with(selected)) - })?; - Some( - selected - .into_iter() - .map(|request| (request.identity(), request)) - .collect::>() - .into_values() - .collect(), - ) + }); + if let Some(selected) = selected { + let requests = spend::unique(selected); + return CallMatch { + spend_pending: spend::request_cost(&requests).is_none(), + requests: Some(requests), + state: SpendMatch::Matched, + }; + } + // A leaf without complete identifiers may still be priced from a wrapper or transport + // once its spend arrives. Only the absence of every complete source is terminal. + let spend_pending = sources.iter().any(SpendEvidence::has_complete_keys) + || (!transports.is_empty() && transports.iter().all(SpendEvidence::has_complete_keys)); + CallMatch { + requests: None, + state: sources[0].unmatched_reason(), + spend_pending, + } } fn transports(&self, call: usize) -> Vec { diff --git a/litellm-rust/crates/traces/src/resolve/spend.rs b/litellm-rust/crates/traces/src/resolve/spend.rs index a6b2fc0dbc5..e74cc3704af 100644 --- a/litellm-rust/crates/traces/src/resolve/spend.rs +++ b/litellm-rust/crates/traces/src/resolve/spend.rs @@ -162,6 +162,10 @@ pub(super) enum SpendEvidence<'a> { } impl<'a> SpendEvidence<'a> { + pub(super) fn has_complete_keys(&self) -> bool { + matches!(self, Self::Complete(keys) if !keys.is_empty()) + } + pub(super) fn unmatched_reason(&self) -> SpendMatch { match self { Self::Unknown => SpendMatch::NoCallId, diff --git a/litellm-rust/crates/traces/src/resolve/view.rs b/litellm-rust/crates/traces/src/resolve/view.rs index 62c27676605..1b1ae7a783e 100644 --- a/litellm-rust/crates/traces/src/resolve/view.rs +++ b/litellm-rust/crates/traces/src/resolve/view.rs @@ -25,8 +25,8 @@ fn optional(value: &str) -> Option { fn span(resolution: &Resolution<'_>, index: usize, trace_start_ns: i64) -> Span { let row = resolution.row(index); let status = resolution.status_source(index); - let (requests, spend_match) = if let Some((requests, matched)) = resolution.call_match(index) { - (requests.clone(), Some(*matched)) + let (requests, spend_match) = if let Some(matched) = resolution.call_match(index) { + (matched.requests.clone(), Some(matched.state)) } else { (resolution.requests(index).complete_requests(), None) }; @@ -249,6 +249,7 @@ pub fn resolve_trace( }), }; Some(Trace { + gateway_spend_pending: resolution.gateway_spend_pending(), summary, agents, spans, diff --git a/litellm-rust/crates/traces/src/view.rs b/litellm-rust/crates/traces/src/view.rs index 729ab5ef1c3..345e1454855 100644 --- a/litellm-rust/crates/traces/src/view.rs +++ b/litellm-rust/crates/traces/src/view.rs @@ -131,6 +131,9 @@ pub struct TraceSummary { #[macro_rules_attribute::apply(response_type)] #[derive(Clone, Debug, PartialEq)] pub struct Trace { + /// Read-cache metadata from gateway resolution; display estimates must not clear it. + #[serde(skip)] + pub gateway_spend_pending: bool, pub summary: TraceSummary, pub agents: Vec, pub spans: Vec, diff --git a/litellm-rust/crates/traces/tests/resolve.rs b/litellm-rust/crates/traces/tests/resolve.rs index 73f691cb2a6..d81aed59a76 100644 --- a/litellm-rust/crates/traces/tests/resolve.rs +++ b/litellm-rust/crates/traces/tests/resolve.rs @@ -494,21 +494,70 @@ fn model_calls_link_the_spend_log_they_were_priced_from() { } #[rstest] -fn model_call_span_cost_agrees_with_the_run_total_when_priced_from_a_wrapper() { +#[case::wrapper_without_leaf_id(false, litellm_traces::CallEvidenceKind::Unknown)] +#[case::wrapper_with_partial_leaf(false, litellm_traces::CallEvidenceKind::Partial)] +#[case::transport_without_leaf_id(true, litellm_traces::CallEvidenceKind::Unknown)] +#[case::transport_with_partial_leaf(true, litellm_traces::CallEvidenceKind::Partial)] +fn model_call_span_cost_agrees_with_the_run_total_when_priced_from_related_evidence( + #[case] transport: bool, + #[case] leaf_evidence: litellm_traces::CallEvidenceKind, +) { let rows = [ TraceSpansRow { - call_keys: vec![litellm_traces::CallKey::LiteLlmRequest("gateway".into())], + trace_id: "trace".into(), + parent_span_id: if transport { "call" } else { "" }.into(), + kind: if transport { + litellm_traces::ObservationType::Chain + } else { + litellm_traces::ObservationType::Llm + }, + call_keys: vec![if transport { + litellm_traces::CallKey::Transport + } else { + litellm_traces::CallKey::LiteLlmRequest("gateway".into()) + }], call_evidence: Some(litellm_traces::CallEvidenceKind::Complete), - ..llm("wrapper", "", "agent", "") + ..llm("evidence", "", "agent", "") + }, + TraceSpansRow { + trace_id: "trace".into(), + call_evidence: Some(leaf_evidence), + ..llm( + "call", + if transport { "" } else { "evidence" }, + "agent", + "chatcmpl-request", + ) }, - llm("call", "wrapper", "agent", ""), ]; + let pending = resolve_trace("trace", "ref", &rows, &[]).unwrap(); + assert_eq!( + pending.spans[1].spend_match, + Some( + if leaf_evidence == litellm_traces::CallEvidenceKind::Partial { + SpendMatch::IncompleteEvidence + } else { + SpendMatch::NoCallId + } + ) + ); + assert!(pending.gateway_spend_pending); + assert!( + serde_json::to_value(&pending) + .unwrap() + .get("gateway_spend_pending") + .is_none() + ); + assert_eq!(pending.summary.spend, None); let logs = [SpendByResponseIdsRow { + trace_id: "trace".into(), + span_id: "evidence".into(), litellm_call_id: "gateway".into(), ..spend("request", "chatcmpl-request", 0.25) }]; let trace = resolve_trace("trace", "ref", &rows, &logs).unwrap(); assert_eq!(trace.summary.spend, Some(0.25)); + assert!(!trace.gateway_spend_pending); assert_eq!( ( trace.spans[1].spend, @@ -1568,7 +1617,12 @@ fn complete_wrapper_reconciles_ambiguous_response( owned_spend("request-a", "response", "team", "", "key", 0.25), owned_spend("request-b", "response", "team", "", "key", 0.5), owned_spend("request-c", "other-response", "team", "", "key", 0.75), + owned_spend("request-d", "response", "team", "", "key", 0.0), ]; + let pending = resolve_trace("trace", "ref", &rows, &logs[1..]).unwrap(); + assert_eq!(pending.spans[1].spend_match, Some(SpendMatch::Ambiguous)); + assert!(pending.gateway_spend_pending); + assert_eq!(pending.summary.spend, None); let trace = resolve_trace("trace", "ref", &rows, &logs).unwrap(); assert_eq!(trace.summary.spend, expected); assert_eq!(trace.agents[0].spend, expected);