This commit is contained in:
Yujong Lee 2026-10-04 12:26:55 -07:00
parent df0e426692
commit 9c2e0e2853
5 changed files with 88 additions and 8 deletions

View file

@ -121,9 +121,15 @@ async fn resolve_runs<S: TraceStore>(
}
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<S: TraceStore>(
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),
)
})
})

View file

@ -53,6 +53,7 @@ struct State {
span_error: Option<SpanErrorRow>,
list_runs_too_large_above: Option<u32>,
trace_too_large_refs: HashSet<String>,
spend_fails_above_response_ids: Option<usize>,
}
#[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<Vec<SpendByResponseIdsRow>, StoreError<Self::Error>> {
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() {

View file

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

View file

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

View file

@ -87,7 +87,7 @@ export function LensModeSwitch({
</span>
<span className={cn(view === "settings" && "sr-only")}>{label}</span>
{view === "investigations" && setup && (
<span className="absolute -top-3 right-2 rounded-full border border-border bg-card px-1.5 py-0.5 text-[10px] leading-none font-medium text-muted-foreground animate-in fade-in-0 duration-200">
<span className="absolute -top-3 right-2 rounded-full border border-border bg-card px-1.5 py-0.5 text-xs leading-none font-medium text-muted-foreground animate-in fade-in-0 duration-200">
{setup}
</span>
)}