From 9c2e0e28533ed347bc06e38df766281be81bf864 Mon Sep 17 00:00:00 2001 From: Yujong Lee Date: Sun, 4 Oct 2026 12:26:55 -0700 Subject: [PATCH] wip --- litellm-rust/crates/traces-cache/src/list.rs | 14 +++- .../crates/traces-cache/tests/read.rs | 77 ++++++++++++++++++- .../traces-clickhouse/src/span_batches.rs | 1 - litellm/proxy/tracing_endpoints.py | 2 +- .../src/components/lens/LensModeSwitch.tsx | 2 +- 5 files changed, 88 insertions(+), 8 deletions(-) diff --git a/litellm-rust/crates/traces-cache/src/list.rs b/litellm-rust/crates/traces-cache/src/list.rs index 7093411bc8c..1b1dcbbefc3 100644 --- a/litellm-rust/crates/traces-cache/src/list.rs +++ b/litellm-rust/crates/traces-cache/src/list.rs @@ -121,9 +121,15 @@ async fn resolve_runs( } Err(error) => return Err(map_store_error(error)), }; - let spend_rows = spend(store, access, &spans).await; - let spend_known = spend_rows.is_some(); - let spend_rows = spend_rows.unwrap_or_default(); + let Some(spend_rows) = spend(store, access, &spans).await else { + // The batch's combined spend read failed; a run's own narrower window may still + // resolve, so fall back per run instead of leaving every run in the batch costless. + let mut resolved = Vec::with_capacity(runs.len()); + for row in runs { + resolved.push(resolve_run(reader, store, access, row).await?); + } + return Ok(resolved); + }; let mut spans = spans; spans.sort_by(|left, right| { run_key(&left.team_id, &left.api_key_hash, &left.trace_id) @@ -158,7 +164,7 @@ async fn resolve_runs( resolve_trace(&row.trace_id, &row.trace_ref, spans, spend).map(|trace| { ListedRun::Resolved( Box::new(trace.summary), - Freshness::of(spans, spend_known, snapshot_ms), + Freshness::of(spans, true, snapshot_ms), ) }) }) diff --git a/litellm-rust/crates/traces-cache/tests/read.rs b/litellm-rust/crates/traces-cache/tests/read.rs index b5d2442582c..2c7098b3d37 100644 --- a/litellm-rust/crates/traces-cache/tests/read.rs +++ b/litellm-rust/crates/traces-cache/tests/read.rs @@ -53,6 +53,7 @@ struct State { span_error: Option, list_runs_too_large_above: Option, trace_too_large_refs: HashSet, + spend_fails_above_response_ids: Option, } #[derive(Default)] @@ -115,6 +116,12 @@ impl FakeStore { .insert(trace_ref.to_owned()); } + /// Fails `spend` only when the lookup covers more than `limit` response ids, so a batch + /// covering several runs fails while each run's own narrower lookup still succeeds. + fn set_spend_fails_above_response_ids(&self, limit: usize) { + self.state.lock().unwrap().spend_fails_above_response_ids = Some(limit); + } + fn calls(&self, operation: Operation) -> usize { match operation { Operation::TraceRefs => self.calls.trace_refs.load(Ordering::SeqCst), @@ -207,11 +214,17 @@ impl TraceStore for FakeStore { async fn spend( &self, - _: &SpendByResponseIdsParams, + params: &SpendByResponseIdsParams, ) -> Result, StoreError> { self.calls.spend.fetch_add(1, Ordering::SeqCst); let state = self.state.lock().unwrap(); Self::failure(&state, Operation::Spend)?; + if state + .spend_fails_above_response_ids + .is_some_and(|limit| params.response_ids.len() > limit) + { + return Err(StoreError::Failed(FakeError)); + } Ok(state.spend.clone()) } @@ -504,6 +517,68 @@ async fn oversized_run_batch_falls_back_to_each_run_and_keeps_listed_summaries() assert_eq!(store.calls(Operation::TraceSpans), 2); } +fn spend_row(response_id: &str, cost: f64) -> SpendByResponseIdsRow { + SpendByResponseIdsRow { + request_id: response_id.into(), + litellm_call_id: String::new(), + response_id: response_id.into(), + upstream_response_id: String::new(), + trace_id: String::new(), + span_id: String::new(), + team_id: "team".into(), + api_key: "key".into(), + user: "user".into(), + spend: Some(cost), + start_ms: START_NS / 1_000_000, + } +} + +#[rstest] +#[tokio::test] +async fn failed_batch_spend_lookup_falls_back_to_each_run_instead_of_losing_every_cost() { + let mut first = span(0); + first.trace_id = "trace-a".into(); + first.kind = ObservationType::Llm; + first.litellm_request_id = "response-a".into(); + first.call_keys = vec![CallKey::ProviderResponse("response-a".into())]; + first.call_evidence = Some(CallEvidenceKind::Complete); + let mut second = span(0); + second.trace_id = "trace-b".into(); + second.kind = ObservationType::Llm; + second.litellm_request_id = "response-b".into(); + second.call_keys = vec![CallKey::ProviderResponse("response-b".into())]; + second.call_evidence = Some(CallEvidenceKind::Complete); + + let store = FakeStore::default(); + store.set_list_runs(vec![run("trace-a", "ref-a"), run("trace-b", "ref-b")]); + store.set_run_spans(vec![first.clone(), second.clone()]); + { + let mut state = store.state.lock().unwrap(); + state.trace_spans.insert("ref-a".to_owned(), vec![first]); + state.trace_spans.insert("ref-b".to_owned(), vec![second]); + state.spend = vec![spend_row("response-a", 1.5), spend_row("response-b", 2.5)]; + } + // The batch covers both runs' response ids (2); each run resolved on its own only ever + // asks for its own (1), so this fails only the combined read, not the per-run fallback. + store.set_spend_fails_above_response_ids(1); + + let page = TraceReader::new(usize::MAX) + .list_traces(&store, &access(), 0, i64::MAX, None, 8) + .await + .unwrap(); + + assert_eq!(page.data.len(), 2); + let by_ref: HashMap<&str, f64> = page + .data + .iter() + .map(|run| (run.trace_ref.as_str(), run.spend.expect("run's own spend read should have succeeded"))) + .collect(); + assert_eq!(by_ref["ref-a"], 1.5); + assert_eq!(by_ref["ref-b"], 2.5); + assert_eq!(store.calls(Operation::RunSpans), 1); + assert_eq!(store.calls(Operation::TraceSpans), 2); +} + #[rstest] #[tokio::test] async fn failed_spend_lookup_preserves_the_trace_with_unknown_spend() { diff --git a/litellm-rust/crates/traces-clickhouse/src/span_batches.rs b/litellm-rust/crates/traces-clickhouse/src/span_batches.rs index 7d77feadeee..5e406d67cfa 100644 --- a/litellm-rust/crates/traces-clickhouse/src/span_batches.rs +++ b/litellm-rust/crates/traces-clickhouse/src/span_batches.rs @@ -45,7 +45,6 @@ trait Keyset: Serialize + Sized + Send + Sync { type Row: Serialize + DeserializeOwned + Send; const SQL: &'static str; - /// The position just after `last`. fn after(self, last: &Self::Row) -> Self; } diff --git a/litellm/proxy/tracing_endpoints.py b/litellm/proxy/tracing_endpoints.py index 503280ea7e8..0d3723c45b8 100644 --- a/litellm/proxy/tracing_endpoints.py +++ b/litellm/proxy/tracing_endpoints.py @@ -177,7 +177,7 @@ def read_failure(error: TraceChanged | ValueError | OverflowError | RuntimeError headers={"Retry-After": str(TRACE_READ_RETRY_AFTER_SECONDS)}, ) case _: - assert_never(error) + return assert_never(error) @router.get("/v1/traces", response_model=TracePage) diff --git a/ui/litellm-dashboard/src/components/lens/LensModeSwitch.tsx b/ui/litellm-dashboard/src/components/lens/LensModeSwitch.tsx index dcca4abb268..0cfa9e47bce 100644 --- a/ui/litellm-dashboard/src/components/lens/LensModeSwitch.tsx +++ b/ui/litellm-dashboard/src/components/lens/LensModeSwitch.tsx @@ -87,7 +87,7 @@ export function LensModeSwitch({ {label} {view === "investigations" && setup && ( - + {setup} )}