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>
This commit is contained in:
Yujong Lee 2026-10-03 23:32:32 +00:00
parent 1c12e4cf73
commit 91f00d0410
9 changed files with 349 additions and 30 deletions

View file

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

View file

@ -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::<ListTracesRow>(
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::<TraceSpansRow>(

View file

@ -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<litellm_traces::TraceSummary> = 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<contracts::ListTracesRow> = page
.into_iter()
.filter(|row| traversal.unseen(row))
.collect();
let items: Vec<litellm_traces::TraceSummary> = stream::iter(unseen.chunks(16))
.then(|batch| list_summaries(client, connection, access, batch, snapshot_ms))
.try_collect::<Vec<_>>()
.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<Option<Trace>, 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<Option<Trace>, 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::<SpanError>(client, connection, &params)
.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,

View file

@ -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<SeededDatabase>,
) -> 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(&current),
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<SeededDatabase>,
) -> 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(())
}

View file

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

View file

@ -55,7 +55,7 @@ fn named_requests_preserve_all_access_cases(
#[rstest]
fn result_contracts_preserve_public_field_names() {
round_trip::<ListTracesRow>(
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::<TraceSpansRow>(
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"}),

View file

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