From 91f00d04105bfbf42b57ca7df8755e8c6a4ad927 Mon Sep 17 00:00:00 2001 From: Yujong Lee Date: Sat, 3 Oct 2026 23:32:32 +0000 Subject: [PATCH] fix(tracing): keep listed runs unique under backdated spans, fence oversized fallback, report changed diagnostics Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- ...ved.sql => 0016_agent_traces_received.sql} | 0 ....sql => 0017_agent_traces_mv_received.sql} | 0 .../traces-clickhouse/query/list_traces.sql | 4 +- .../traces-clickhouse/src/query/named.rs | 4 +- .../crates/traces-clickhouse/src/reads.rs | 150 +++++++++--- .../crates/traces-clickhouse/tests/reads.rs | 216 +++++++++++++++++- litellm-rust/crates/traces/src/query/named.rs | 2 + .../crates/traces/tests/query/named.rs | 2 +- litellm-rust/crates/traces/tests/resolve.rs | 1 + 9 files changed, 349 insertions(+), 30 deletions(-) rename litellm-rust/crates/traces-clickhouse/migrations/{0015_agent_traces_received.sql => 0016_agent_traces_received.sql} (100%) rename litellm-rust/crates/traces-clickhouse/migrations/{0016_agent_traces_mv_received.sql => 0017_agent_traces_mv_received.sql} (100%) diff --git a/litellm-rust/crates/traces-clickhouse/migrations/0015_agent_traces_received.sql b/litellm-rust/crates/traces-clickhouse/migrations/0016_agent_traces_received.sql similarity index 100% rename from litellm-rust/crates/traces-clickhouse/migrations/0015_agent_traces_received.sql rename to litellm-rust/crates/traces-clickhouse/migrations/0016_agent_traces_received.sql diff --git a/litellm-rust/crates/traces-clickhouse/migrations/0016_agent_traces_mv_received.sql b/litellm-rust/crates/traces-clickhouse/migrations/0017_agent_traces_mv_received.sql similarity index 100% rename from litellm-rust/crates/traces-clickhouse/migrations/0016_agent_traces_mv_received.sql rename to litellm-rust/crates/traces-clickhouse/migrations/0017_agent_traces_mv_received.sql diff --git a/litellm-rust/crates/traces-clickhouse/query/list_traces.sql b/litellm-rust/crates/traces-clickhouse/query/list_traces.sql index 918f39e8f0c..746a7cb15b8 100644 --- a/litellm-rust/crates/traces-clickhouse/query/list_traces.sql +++ b/litellm-rust/crates/traces-clickhouse/query/list_traces.sql @@ -19,6 +19,7 @@ FROM agent_traces_by_key WHERE ({all_teams:UInt8} = 1 OR ({user_id:String} != '' AND UserIds = [{user_id:String}]) OR has({team_ids:Array(String)}, TeamId)) + AND ({snapshot_ms:UInt64} = 0 OR ReceivedMs <= {snapshot_ms:UInt64}) GROUP BY TeamId, ApiKeyHash, TraceId HAVING min(StartTs) >= fromUnixTimestamp64Milli({start_ms:Int64}) AND min(StartTs) < fromUnixTimestamp64Milli({end_ms:Int64}) @@ -30,10 +31,11 @@ LIMIT {limit:UInt32} ) SELECT page.* EXCEPT (trace_start, trace_end), identities.agent_names AS agent_names, identities.agent_count AS agent_count, - identities.frameworks AS frameworks + identities.frameworks AS frameworks, identities.fenced_start_ms AS fenced_start_ms FROM page LEFT JOIN ( SELECT TeamId, ApiKeyHash, TraceId, + toUnixTimestamp64Milli(min(Timestamp)) AS fenced_start_ms, arraySort(groupUniqArrayIf(AgentName, AgentName != '')) AS agent_names, arraySort(groupUniqArrayIf(toString(Framework), Framework != '')) AS frameworks, uniqExactIf(if(AgentName = '', SpanName, AgentName), ObservationType = 'agent') AS agent_count diff --git a/litellm-rust/crates/traces-clickhouse/src/query/named.rs b/litellm-rust/crates/traces-clickhouse/src/query/named.rs index 4f54f427897..91624411de3 100644 --- a/litellm-rust/crates/traces-clickhouse/src/query/named.rs +++ b/litellm-rust/crates/traces-clickhouse/src/query/named.rs @@ -48,6 +48,8 @@ struct ListTracesRowEncoding { pub status: litellm_traces::SpanStatus, #[serde(deserialize_with = "super::number::deserialize")] pub start_ms: i64, + #[serde(default, deserialize_with = "super::number::deserialize")] + pub fenced_start_ms: i64, #[serde(deserialize_with = "super::number::deserialize")] pub duration_ms: i64, #[serde(deserialize_with = "super::number::deserialize")] @@ -341,7 +343,7 @@ mod tests { #[case::quoted(true)] fn rows_decode_into_neutral_contracts(#[case] quoted: bool) { round_trip::( - json!({"trace_id": "trace", "trace_ref": "ref", "team_id": "team", "api_key_hash": "key", "user_id": "user", "name": "agent", "service": "service", "input_preview": "input", "status": "STATUS_CODE_OK", "start_ms": -1, "duration_ms": 20, "span_count": u64::MAX, "agent_count": 1, "agent_invocations": 2, "agent_names": ["agent"], "frameworks": ["claude-agent-sdk"], "llm_calls": 3, "tool_calls": 4, "input_tokens": 5, "output_tokens": 6, "models": ["model"], "error_count": 0, "request_ids": ["request"]}), + json!({"trace_id": "trace", "trace_ref": "ref", "team_id": "team", "api_key_hash": "key", "user_id": "user", "name": "agent", "service": "service", "input_preview": "input", "status": "STATUS_CODE_OK", "start_ms": -1, "fenced_start_ms": -1, "duration_ms": 20, "span_count": u64::MAX, "agent_count": 1, "agent_invocations": 2, "agent_names": ["agent"], "frameworks": ["claude-agent-sdk"], "llm_calls": 3, "tool_calls": 4, "input_tokens": 5, "output_tokens": 6, "models": ["model"], "error_count": 0, "request_ids": ["request"]}), quoted, ); round_trip::( diff --git a/litellm-rust/crates/traces-clickhouse/src/reads.rs b/litellm-rust/crates/traces-clickhouse/src/reads.rs index 27c4e09b6dd..6ee2112155d 100644 --- a/litellm-rust/crates/traces-clickhouse/src/reads.rs +++ b/litellm-rust/crates/traces-clickhouse/src/reads.rs @@ -133,6 +133,16 @@ impl ListTraversal { ) } + /// A backdated span received after publication can lower a listed run's rollup start below + /// the cursor; its fenced start still sits at or above the cursor, so the run is not repeated. + fn unseen(&self, row: &contracts::ListTracesRow) -> bool { + let cursor = &self.cursor.position; + cursor.start_ms == 0 + || row.fenced_start_ms == 0 + || (row.fenced_start_ms, row.trace_ref.as_str()) + < (cursor.start_ms, cursor.trace_ref.as_str()) + } + fn continue_after( &self, keys: &KeyRing, @@ -275,17 +285,24 @@ pub async fn list_traces( .iter() .map(|row| (row.trace_ref.clone(), row.start_ms)) .collect(); - let items: Vec = stream::iter(page.chunks(16)) + let last_candidate = page.last().map(|row| ListPosition { + start_ms: row.start_ms, + trace_ref: row.trace_ref.clone(), + }); + let unseen: Vec = page + .into_iter() + .filter(|row| traversal.unseen(row)) + .collect(); + let items: Vec = stream::iter(unseen.chunks(16)) .then(|batch| list_summaries(client, connection, access, batch, snapshot_ms)) .try_collect::>() .await? .into_iter() .flatten() .collect(); - let next_cursor = items - .last() + let next_cursor = last_candidate .filter(|_| !exhausted) - .map(|last| traversal.continue_after(keys, &starts, last)) + .map(|last| keys.encode(&traversal.binding, &traversal.cursor.advance(last))) .transpose()?; let page = Page { items, @@ -318,27 +335,34 @@ async fn list_summaries( start_ms, end_ms: end_ms.saturating_add(1), }); - let spans = match crate::span_batches::read_list_spans(client, connection, params, snapshot_ms) - .await - { - Ok(spans) => spans, - Err(Error::ReadTooLarge) => { - return stream::iter(runs) - .then(|row| async move { - match get_trace(client, connection, access, &row.trace_id, &row.trace_ref).await - { - Ok(trace) => { - Ok(trace.map_or_else(|| listed_summary(row), |trace| trace.summary)) + let spans = + match crate::span_batches::read_list_spans(client, connection, params, snapshot_ms).await { + Ok(spans) => spans, + Err(Error::ReadTooLarge) => { + return stream::iter(runs) + .then(|row| async move { + match read_trace( + client, + connection, + access, + &row.trace_id, + &row.trace_ref, + snapshot_ms, + ) + .await + { + Ok(trace) => Ok( + trace.map_or_else(|| listed_summary(row), |trace| trace.summary) + ), + Err(Error::ReadTooLarge) => Ok(listed_summary(row)), + Err(error) => Err(error), } - Err(Error::ReadTooLarge) => Ok(listed_summary(row)), - Err(error) => Err(error), - } - }) - .try_collect() - .await; - } - Err(error) => return Err(error), - }; + }) + .try_collect() + .await; + } + Err(error) => return Err(error), + }; let by_trace = spans.into_iter().into_group_map_by(|span| { ( span.team_id.clone(), @@ -376,6 +400,17 @@ pub async fn get_trace( access: &ReadAccessParams, trace_id: &str, trace_ref: &str, +) -> Result, Error> { + read_trace(client, connection, access, trace_id, trace_ref, u64::MAX).await +} + +async fn read_trace( + client: &Client, + connection: &Connection, + access: &ReadAccessParams, + trace_id: &str, + trace_ref: &str, + snapshot_ms: u64, ) -> Result, Error> { let Some(trace_ref) = reference(client, connection, access, trace_id, trace_ref).await? else { return Ok(None); @@ -385,7 +420,7 @@ pub async fn get_trace( trace_id: trace_id.to_owned(), trace_ref: trace_ref.clone(), }; - let rows = crate::span_batches::read_spans(client, connection, params, u64::MAX).await?; + let rows = crate::span_batches::read_spans(client, connection, params, snapshot_ms).await?; if rows.is_empty() { return Ok(None); } @@ -566,7 +601,7 @@ pub async fn get_span_error( trace_ref, span_id: span_id.to_owned(), error_offset: offset, - error_version: position.revision.clone(), + error_version: String::new(), }); let Some(row) = fetch::(client, connection, ¶ms) .await? @@ -625,6 +660,68 @@ mod tests { OffsetDateTime::from_unix_timestamp(1_790_000_000).unwrap() } + fn listed_row( + trace_ref: &str, + start_ms: i64, + fenced_start_ms: i64, + ) -> contracts::ListTracesRow { + contracts::ListTracesRow { + trace_id: "t".into(), + trace_ref: trace_ref.into(), + team_id: "team-a".into(), + api_key_hash: String::new(), + user_id: "user-a".into(), + name: String::new(), + service: String::new(), + input_preview: String::new(), + status: litellm_traces::SpanStatus::Ok, + start_ms, + fenced_start_ms, + duration_ms: 1, + span_count: 1, + agent_count: 0, + agent_invocations: 0, + llm_calls: 0, + tool_calls: 0, + input_tokens: 0, + output_tokens: 0, + models: Vec::new(), + agent_names: Vec::new(), + frameworks: Vec::new(), + error_count: 0, + request_ids: Vec::new(), + } + } + + #[rstest] + #[case::first_page_keeps_everything(None, 50, 900, true)] + #[case::earlier_fenced_start(Some((500, "b")), 50, 400, true)] + #[case::same_start_lower_ref(Some((500, "b")), 500, 500, true)] + #[case::unknown_fenced_start(Some((500, "b")), 50, 0, true)] + #[case::backdated_after_publication(Some((500, "b")), 50, 900, false)] + #[case::same_start_same_ref(Some((500, "a")), 50, 500, false)] + fn unseen_drops_runs_whose_fenced_start_was_already_listed( + #[case] cursor: Option<(i64, &str)>, + #[case] start_ms: i64, + #[case] fenced_start_ms: i64, + #[case] expected: bool, + ) { + let keys = keys(); + let token = cursor.map(|(start_ms, trace_ref)| { + let first = ListTraversal::open(&keys, &access(), (0, 1000), None, now()).unwrap(); + let position = ListPosition { + start_ms, + trace_ref: trace_ref.into(), + }; + keys.encode(&first.binding, &first.cursor.advance(position)) + .unwrap() + }); + let traversal = + ListTraversal::open(&keys, &access(), (0, 1000), token.as_deref(), now()).unwrap(); + let row = listed_row("a", start_ms, fenced_start_ms); + assert_eq!(traversal.unseen(&row), expected); + } + #[rstest] fn list_traversal_pins_publication_and_continues_after_the_last_run() { let keys = keys(); @@ -641,6 +738,7 @@ mod tests { input_preview: String::new(), status: litellm_traces::SpanStatus::Ok, start_ms: 1_790_742_989_377, + fenced_start_ms: 1_790_742_989_377, duration_ms: 1, span_count: 1, agent_count: 0, diff --git a/litellm-rust/crates/traces-clickhouse/tests/reads.rs b/litellm-rust/crates/traces-clickhouse/tests/reads.rs index e8aa9386087..e05ae483c6c 100644 --- a/litellm-rust/crates/traces-clickhouse/tests/reads.rs +++ b/litellm-rust/crates/traces-clickhouse/tests/reads.rs @@ -2,7 +2,8 @@ use std::collections::BTreeMap; use litellm_traces::query::named::ReadAccessParams; use litellm_traces_clickhouse::{ - Connection, InsertTable, QueryScope, get_trace, get_trace_page, insert_rows, list_traces, + Connection, InsertTable, QueryScope, get_span_error, get_trace, get_trace_page, insert_rows, + list_traces, }; use rstest::rstest; use serde_json::json; @@ -18,6 +19,26 @@ fn keys() -> litellm_pagination::KeyRing { litellm_pagination::KeyRing::new(["fixture-cursor-secret"]).unwrap() } +/// `insert_rows` stamps `EngineReceivedMs` with the wall clock; tests that pin a receipt time +/// write the row directly. +async fn insert_received( + client: &litellm_http::Client, + writer: &Connection, + row: serde_json::Value, +) -> TestResult { + let response = client + .post(writer.url().clone()) + .body(format!( + "INSERT INTO {DATABASE}.otel_traces FORMAT JSONEachRow\n{row}" + )) + .send() + .await?; + if !response.status().is_success() { + return Err(response.text().await?.into()); + } + Ok(()) +} + #[rstest] #[case::api_key("key-a", "")] #[case::user("", "user-a")] @@ -628,6 +649,21 @@ async fn an_oversized_span_keeps_the_run_list_available_with_partial_totals( ])], ) .await?; + insert_received( + client, + &writer, + json!({ + "Timestamp": "2026-09-21 00:00:00.000000000", + "TraceId": run.trace_id, + "SpanId": "late-child", + "ParentSpanId": "0101010101010101", + "ObservationType": "tool", + "TeamId": "team-a", + "ApiKeyHash": "key-a", + "EngineReceivedMs": u64::MAX / 2 + }), + ) + .await?; let after = list_traces( client, &reader, @@ -670,3 +706,181 @@ async fn an_oversized_span_keeps_the_run_list_available_with_partial_totals( )); Ok(()) } + +#[rstest] +#[tokio::test] +async fn a_backdated_span_received_after_publication_does_not_repeat_a_listed_run( + #[future(awt)] seeded_database: TestResult, +) -> TestResult { + let fixture = seeded_database?; + let client = &fixture.database.client; + let reader = fixture + .readers + .connection(client, &QueryScope::All, "fixture-secret") + .await?; + let access = ReadAccessParams { + all_teams: true, + user_id: String::new(), + team_ids: Vec::new(), + }; + let mut cursor = None; + let mut listed = Vec::new(); + let fixture_run = loop { + let page = list_traces( + client, + &reader, + &keys(), + &access, + 0, + 2_000_000_000_000, + cursor.as_deref(), + 1, + ) + .await?; + let run = page + .items + .into_iter() + .next() + .ok_or("fixture run not listed")?; + cursor = page.next_cursor; + listed.push(run.trace_ref.clone()); + if run.span_count == 3 { + break run; + } + }; + let continuation = cursor.clone().ok_or("fixture run was the last run")?; + let writer = Connection::writer(&fixture.database.url)?; + insert_received( + client, + &writer, + json!({ + "Timestamp": "2020-01-01 00:00:00.000000000", + "TraceId": fixture_run.trace_id, + "SpanId": "backdated-root", + "TeamId": "team-a", + "ApiKeyHash": "key-a", + "EngineReceivedMs": u64::MAX / 2 + }), + ) + .await?; + for statement in [ + format!("SYSTEM START MERGES {DATABASE}.agent_traces_by_key"), + format!("OPTIMIZE TABLE {DATABASE}.agent_traces_by_key FINAL"), + ] { + let response = client + .post(fixture.database.url.clone()) + .body(statement) + .send() + .await?; + if !response.status().is_success() { + return Err(response.text().await?.into()); + } + } + let mut cursor = Some(continuation); + while let Some(current) = cursor { + let page = list_traces( + client, + &reader, + &keys(), + &access, + 0, + 2_000_000_000_000, + Some(¤t), + 1, + ) + .await?; + listed.extend(page.items.into_iter().map(|run| run.trace_ref)); + cursor = page.next_cursor; + } + let mut unique = listed.clone(); + unique.sort(); + unique.dedup(); + assert_eq!(unique.len(), listed.len(), "{listed:?}"); + let fresh = list_traces( + client, + &reader, + &keys(), + &access, + 0, + 2_000_000_000_000, + None, + 50, + ) + .await?; + let fresh_listings = fresh + .items + .iter() + .filter(|run| run.trace_ref == fixture_run.trace_ref) + .count(); + assert_eq!(fresh_listings, 1); + Ok(()) +} + +#[rstest] +#[tokio::test] +async fn a_changed_diagnostic_reports_traversal_changed_instead_of_not_found( + #[future(awt)] seeded_database: TestResult, +) -> TestResult { + let fixture = seeded_database?; + let client = &fixture.database.client; + let reader = fixture + .readers + .connection(client, &QueryScope::All, "fixture-secret") + .await?; + let access = ReadAccessParams { + all_teams: true, + user_id: String::new(), + team_ids: Vec::new(), + }; + let writer = Connection::writer(&fixture.database.url)?; + let diagnostic = |received: u64, message: String| { + json!({ + "Timestamp": "2026-09-21 00:00:00.000000000", + "TraceId": "diagnostic-trace", + "SpanId": "diagnostic-span", + "StatusCode": "STATUS_CODE_ERROR", + "StatusMessage": message, + "TeamId": "team-a", + "ApiKeyHash": "key-a", + "EngineReceivedMs": received + }) + }; + insert_received(client, &writer, diagnostic(100, "a".repeat(40_000))).await?; + let first = get_span_error( + client, + &reader, + &keys(), + &access, + "diagnostic-trace", + "diagnostic-span", + "", + None, + ) + .await? + .ok_or("missing diagnostic")?; + let continuation = first.next_cursor.ok_or("diagnostic fit in one page")?; + insert_received(client, &writer, diagnostic(50, "b".repeat(40_000))).await?; + let uncached_reader = + Connection::reader(&format!("{}?max_threads=1", fixture.database.url), DATABASE)?; + let changed = get_span_error( + client, + &uncached_reader, + &keys(), + &access, + "diagnostic-trace", + "diagnostic-span", + "", + Some(&continuation), + ) + .await; + assert!( + matches!( + changed, + Err(litellm_traces_clickhouse::Error::Pagination( + litellm_pagination::Error::TraversalChanged + )) + ), + "{changed:?}" + ); + Ok(()) +} diff --git a/litellm-rust/crates/traces/src/query/named.rs b/litellm-rust/crates/traces/src/query/named.rs index e9df28ee9c8..7f08693965c 100644 --- a/litellm-rust/crates/traces/src/query/named.rs +++ b/litellm-rust/crates/traces/src/query/named.rs @@ -41,6 +41,8 @@ pub struct ListTracesRow { #[serde(serialize_with = "crate::wire::serialize_status")] pub status: crate::SpanStatus, pub start_ms: i64, + #[serde(default)] + pub fenced_start_ms: i64, pub duration_ms: i64, pub span_count: u64, pub agent_count: u64, diff --git a/litellm-rust/crates/traces/tests/query/named.rs b/litellm-rust/crates/traces/tests/query/named.rs index 7db1a5895f7..4fca4281674 100644 --- a/litellm-rust/crates/traces/tests/query/named.rs +++ b/litellm-rust/crates/traces/tests/query/named.rs @@ -55,7 +55,7 @@ fn named_requests_preserve_all_access_cases( #[rstest] fn result_contracts_preserve_public_field_names() { round_trip::( - json!({"trace_id": "trace", "trace_ref": "ref", "team_id": "team", "api_key_hash": "key", "user_id": "user", "name": "agent", "service": "service", "input_preview": "input", "status": "STATUS_CODE_OK", "start_ms": -1, "duration_ms": 20, "span_count": u64::MAX, "agent_count": 1, "agent_invocations": 2, "agent_names": ["agent"], "frameworks": ["framework"], "llm_calls": 3, "tool_calls": 4, "input_tokens": 5, "output_tokens": 6, "models": ["model"], "error_count": 0, "request_ids": ["request"]}), + json!({"trace_id": "trace", "trace_ref": "ref", "team_id": "team", "api_key_hash": "key", "user_id": "user", "name": "agent", "service": "service", "input_preview": "input", "status": "STATUS_CODE_OK", "start_ms": -1, "fenced_start_ms": -1, "duration_ms": 20, "span_count": u64::MAX, "agent_count": 1, "agent_invocations": 2, "agent_names": ["agent"], "frameworks": ["framework"], "llm_calls": 3, "tool_calls": 4, "input_tokens": 5, "output_tokens": 6, "models": ["model"], "error_count": 0, "request_ids": ["request"]}), ); round_trip::( json!({"trace_id": "trace", "span_id": "span", "parent_span_id": "parent", "name": "agent", "type": "agent", "wrapper_candidate": 1, "agent": "agent", "framework": "framework", "status": "STATUS_CODE_ERROR", "status_message": "error", "error_truncated": 1, "start_ns": -1, "duration_ns": u64::MAX, "service": "service", "input_preview": "input", "model": "model", "input_tokens": u32::MAX, "output_tokens": 6, "litellm_request_id": "request", "call_keys": ["provider_response:request"], "call_evidence": "complete", "tool_call_id": "call", "team_id": "team", "api_key_hash": "key", "user_id": "user"}), diff --git a/litellm-rust/crates/traces/tests/resolve.rs b/litellm-rust/crates/traces/tests/resolve.rs index 367f1d146ee..891f03c62fe 100644 --- a/litellm-rust/crates/traces/tests/resolve.rs +++ b/litellm-rust/crates/traces/tests/resolve.rs @@ -577,6 +577,7 @@ fn listed_summary_keeps_rollup_counts_with_unknown_cost() { input_preview: "hi".into(), status: SpanStatus::Ok, start_ms: 1_790_742_989_377, + fenced_start_ms: 1_790_742_989_377, duration_ms: 51_385, span_count: 126, agent_count: 2,