fix(lens): refresh runs until gateway costs are complete (#45105)

This commit is contained in:
tin-berri 2026-10-08 00:24:41 -07:00 • committed by GitHub
parent 8683f1418e
commit 057034d0d3
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
10 changed files with 363 additions and 44 deletions

View file

@ -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::<serde_json::Value>().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<u32>,
#[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::<serde_json::Value>().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() {

View file

@ -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<f64>) {
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
);
}

View file

@ -162,10 +162,8 @@ async fn resolve_runs<S: TraceStore>(
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())

View file

@ -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<E>(
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(),

View file

@ -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<f64>],
#[case] initial_total: Option<f64>,
#[case] initial_priced: u64,
#[case] settled: bool,
#[case] final_second: f64,
#[case] terminal: Option<CallEvidenceKind>,
) {
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() {

View file

@ -31,7 +31,11 @@ pub(super) struct Resolution<'a> {
call_matches: HashMap<usize, CallMatch<'a>>,
}
pub(super) type CallMatch<'a> = (Option<Requests<'a>>, SpendMatch);
pub(super) struct CallMatch<'a> {
pub(super) requests: Option<Requests<'a>>,
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<Requests<'a>> {
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<Requests<'a>> {
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::<IndexMap<_, _>>()
.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<usize> {

View file

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

View file

@ -25,8 +25,8 @@ fn optional(value: &str) -> Option<String> {
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,

View file

@ -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<AgentNode>,
pub spans: Vec<Span>,

View file

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