mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-06 02:48:13 +00:00
fix(lens): paginate trace reads within ClickHouse limits (#44384)
This commit is contained in:
parent
8b1990b4bc
commit
ad8babae33
31 changed files with 1258 additions and 130 deletions
|
|
@ -35,13 +35,14 @@ fn map_error_ref(error: &Error) -> PyErr {
|
|||
use litellm_storage_clickhouse::Error as StorageError;
|
||||
|
||||
match error {
|
||||
Error::Decode(litellm_traces::Error::TooLarge) | Error::InsertTooLarge => {
|
||||
PyOverflowError::new_err(error.to_string())
|
||||
}
|
||||
Error::Decode(litellm_traces::Error::TooLarge)
|
||||
| Error::InsertTooLarge
|
||||
| Error::ReadTooLarge => PyOverflowError::new_err(error.to_string()),
|
||||
Error::InvalidRow
|
||||
| Error::InvalidTable
|
||||
| Error::InvalidCursor(_)
|
||||
| Error::AmbiguousTrace
|
||||
| Error::TraceChanged
|
||||
| Error::Decode(_)
|
||||
| Error::InvalidSchema
|
||||
| Error::InvalidQuery
|
||||
|
|
@ -233,26 +234,44 @@ impl NativeTraceStorage {
|
|||
)
|
||||
}
|
||||
|
||||
#[pyo3(signature = (trace_id, scope, trace_ref, cursor=None, page_size=None))]
|
||||
fn get_trace<'py>(
|
||||
&self,
|
||||
py: Python<'py>,
|
||||
trace_id: String,
|
||||
#[pyo3(from_py_with = litellm_host_python::from_py_argument)] scope: ReadAccessParams,
|
||||
trace_ref: String,
|
||||
cursor: Option<String>,
|
||||
page_size: Option<u32>,
|
||||
) -> PyResult<Bound<'py, PyAny>> {
|
||||
let client = crate::http::host_client(py, ClientVariant::NoRedirect)?;
|
||||
let connection = self.config.storage().reader().clone();
|
||||
crate::execution::run_async(
|
||||
py,
|
||||
async move {
|
||||
litellm_traces_clickhouse::get_trace(
|
||||
&client,
|
||||
&connection,
|
||||
&scope,
|
||||
&trace_id,
|
||||
&trace_ref,
|
||||
)
|
||||
.await
|
||||
if let Some(page_size) = page_size {
|
||||
litellm_traces_clickhouse::get_trace_page(
|
||||
&client,
|
||||
&connection,
|
||||
&scope,
|
||||
&trace_id,
|
||||
&trace_ref,
|
||||
cursor.as_deref(),
|
||||
page_size,
|
||||
)
|
||||
.await
|
||||
} else if cursor.is_some() {
|
||||
Err(Error::InvalidParameters)
|
||||
} else {
|
||||
litellm_traces_clickhouse::get_trace(
|
||||
&client,
|
||||
&connection,
|
||||
&scope,
|
||||
&trace_id,
|
||||
&trace_ref,
|
||||
)
|
||||
.await
|
||||
}
|
||||
},
|
||||
map_error,
|
||||
)
|
||||
|
|
@ -465,6 +484,8 @@ mod tests {
|
|||
#[case::invalid_export(Error::Decode(litellm_traces::Error::InvalidPayload), "ValueError")]
|
||||
#[case::cursor(Error::InvalidCursor("trace"), "ValueError")]
|
||||
#[case::ambiguous(Error::AmbiguousTrace, "ValueError")]
|
||||
#[case::changed_snapshot(Error::TraceChanged, "ValueError")]
|
||||
#[case::read_budget(Error::ReadTooLarge, "OverflowError")]
|
||||
fn trace_read_and_ingest_failures_preserve_public_exception_types(
|
||||
#[case] error: Error,
|
||||
#[case] exception_name: &str,
|
||||
|
|
|
|||
|
|
@ -110,6 +110,13 @@ pub async fn execute_read(
|
|||
.body(sql.to_owned());
|
||||
let mut response = request.send().await.map_err(|_| Error::Transport)?;
|
||||
if !response.status().is_success() {
|
||||
if response
|
||||
.headers()
|
||||
.get("x-clickhouse-exception-code")
|
||||
.is_some_and(|code| code == "396")
|
||||
{
|
||||
return Err(Error::ResponseTooLarge);
|
||||
}
|
||||
return Err(Error::QueryFailed(response.status().as_u16()));
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -111,3 +111,36 @@ async fn typed_fetch_encodes_parameters_and_validates_rows(
|
|||
assert!(matches!(envelope, Err(Error::InvalidResponse)));
|
||||
}
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[case::result_limit("396", true)]
|
||||
#[case::memory_limit("241", false)]
|
||||
#[case::timeout("159", false)]
|
||||
#[case::unknown("", false)]
|
||||
#[tokio::test]
|
||||
async fn server_result_limits_allow_smaller_pages_without_retrying_other_failures(
|
||||
#[case] code: &str,
|
||||
#[case] result_limit: bool,
|
||||
) {
|
||||
use wiremock::{Mock, MockServer, ResponseTemplate, matchers::method};
|
||||
let server = MockServer::start().await;
|
||||
Mock::given(method("POST"))
|
||||
.respond_with(ResponseTemplate::new(500).insert_header("X-ClickHouse-Exception-Code", code))
|
||||
.expect(1)
|
||||
.mount(&server)
|
||||
.await;
|
||||
let connection = Connection::parse(&server.uri()).unwrap();
|
||||
let error = execute_read(
|
||||
&Client::no_redirect_for_test(),
|
||||
&connection,
|
||||
"SELECT 1",
|
||||
&BTreeMap::new(),
|
||||
)
|
||||
.await
|
||||
.unwrap_err();
|
||||
if result_limit {
|
||||
assert!(matches!(error, Error::ResponseTooLarge));
|
||||
} else {
|
||||
assert!(matches!(error, Error::QueryFailed(500)));
|
||||
}
|
||||
}
|
||||
|
|
|
|||
27
litellm-rust/crates/traces-clickhouse/query/spend_batch.sql
Normal file
27
litellm-rust/crates/traces-clickhouse/query/spend_batch.sql
Normal file
|
|
@ -0,0 +1,27 @@
|
|||
SELECT * FROM (
|
||||
SELECT request_id, response_id, upstream_response_id, trace_id, span_id, team_id, api_key, user, spend,
|
||||
toUnixTimestamp64Milli(start_time) AS start_ms
|
||||
FROM (
|
||||
SELECT *,
|
||||
-- A chat request served through the Responses API returns the upstream `resp_` id to the
|
||||
-- client but logs LiteLLM's managed `resp_<base64>` id, which embeds it.
|
||||
if(startsWith(response_id, 'resp_'),
|
||||
extract(tryBase64Decode(substring(response_id, 6)), 'response_id:([^;]+)'),
|
||||
'') AS upstream_response_id
|
||||
FROM spend_logs FINAL
|
||||
WHERE start_time >= fromUnixTimestamp64Milli({start_ms:Int64})
|
||||
AND start_time < fromUnixTimestamp64Milli({end_ms:Int64})
|
||||
AND ({all_teams:UInt8} = 1
|
||||
OR ({user_id:String} != '' AND user = {user_id:String})
|
||||
OR has({team_ids:Array(String)}, team_id))
|
||||
)
|
||||
WHERE response_id IN {response_ids:Array(String)}
|
||||
OR upstream_response_id IN {response_ids:Array(String)}
|
||||
OR request_id IN {request_ids:Array(String)}
|
||||
OR (trace_id != '' AND trace_id IN {trace_ids:Array(String)})
|
||||
ORDER BY start_time DESC
|
||||
)
|
||||
WHERE {has_cursor:UInt8} = 0
|
||||
OR (team_id, start_ms, request_id) > ({after_team:String}, {after_ms:Int64}, {after_id:String})
|
||||
ORDER BY team_id, start_ms, request_id
|
||||
LIMIT {page_size:UInt32}
|
||||
|
|
@ -0,0 +1,30 @@
|
|||
SELECT * FROM (
|
||||
SELECT o.TraceId AS trace_id, o.SpanId AS span_id, o.ParentSpanId AS parent_span_id, o.SpanName AS name,
|
||||
o.ObservationType AS type, toUInt8(o.WrapperCandidate) AS wrapper_candidate, o.AgentName AS agent,
|
||||
o.Framework AS framework, o.StatusCode AS status,
|
||||
substringUTF8(o.StatusMessage, 1, 128) AS status_message,
|
||||
lengthUTF8(o.StatusMessage) > 128 AS error_truncated,
|
||||
toUnixTimestamp64Nano(o.Timestamp) AS start_ns, o.Duration AS duration_ns,
|
||||
o.ServiceName AS service, o.InputPreview AS input_preview, o.Model AS model,
|
||||
o.InputTokens AS input_tokens, o.OutputTokens AS output_tokens,
|
||||
o.LiteLLMRequestId AS litellm_request_id,
|
||||
o.CallKeys AS call_keys, o.CallEvidence AS call_evidence,
|
||||
-- Rows written before ToolCallId keep the call id only in their attributes.
|
||||
if(o.ToolCallId != '' OR o.ObservationType != 'tool', o.ToolCallId,
|
||||
coalesce(nullIf(o.SpanAttributes['gen_ai.tool.call.id'], ''), nullIf(o.SpanAttributes['tool.id'], ''), ''))
|
||||
AS tool_call_id,
|
||||
o.UserId AS user_id, o.TeamId AS team_id, o.ApiKeyHash AS api_key_hash
|
||||
FROM otel_traces AS o
|
||||
WHERE o.TraceId = {trace_id:String}
|
||||
AND ({all_teams:UInt8} = 1
|
||||
OR ({user_id:String} != '' AND o.UserId = {user_id:String})
|
||||
OR has({team_ids:Array(String)}, o.TeamId))
|
||||
AND ({trace_ref:String} = '' OR
|
||||
hex(SHA256(concat(o.TeamId, char(0), o.ApiKeyHash, char(0), o.TraceId))) = {trace_ref:String})
|
||||
AND o.EngineReceivedMs <= {snapshot_ms:UInt64}
|
||||
ORDER BY o.Timestamp, o.EngineReceivedMs, o.StatusMessage
|
||||
LIMIT 1 BY o.SpanId
|
||||
)
|
||||
WHERE span_id > {after_span_id:String}
|
||||
ORDER BY span_id
|
||||
LIMIT {page_size:UInt32}
|
||||
|
|
@ -14,6 +14,8 @@ pub enum Error {
|
|||
InvalidResponse,
|
||||
#[error("ClickHouse insert exceeds the encoded size limit")]
|
||||
InsertTooLarge,
|
||||
#[error("Trace exceeds the interactive read budget; use a filtered trace query")]
|
||||
ReadTooLarge,
|
||||
#[error("ClickHouse schema setup failed with HTTP status {0}")]
|
||||
SchemaFailed(u16),
|
||||
#[error("ClickHouse schema setup transport failed")]
|
||||
|
|
@ -34,6 +36,8 @@ pub enum Error {
|
|||
InvalidCursor(&'static str),
|
||||
#[error("Multiple traces have this ID; provide trace_ref")]
|
||||
AmbiguousTrace,
|
||||
#[error("Trace changed while paging; refresh the trace to continue")]
|
||||
TraceChanged,
|
||||
#[error(transparent)]
|
||||
Decode(#[from] litellm_traces::Error),
|
||||
#[error("trace ingestion task failed")]
|
||||
|
|
|
|||
|
|
@ -17,6 +17,7 @@ pub mod query;
|
|||
mod query_access;
|
||||
mod reads;
|
||||
mod schema;
|
||||
mod span_batches;
|
||||
mod span_row;
|
||||
mod sql;
|
||||
mod table;
|
||||
|
|
@ -30,7 +31,7 @@ pub use litellm_storage_clickhouse::{Connection, Parameter};
|
|||
pub use litellm_traces::{QueryScope, ReadQuery};
|
||||
pub use query::{QueryHelp, execute_read, query_help, query_sql};
|
||||
pub use query_access::QueryReaders;
|
||||
pub use reads::{get_span, get_span_error, get_trace, list_traces};
|
||||
pub use reads::{get_span, get_span_error, get_trace, get_trace_page, list_traces};
|
||||
pub use schema::{
|
||||
NORMALIZED_FIELD_DEFINITIONS, NormalizedFieldDefinition, ensure_schema, schema_statements,
|
||||
};
|
||||
|
|
|
|||
|
|
@ -1,26 +1,54 @@
|
|||
//! Scoped trace reads: the trace list, one trace resolved with its spend, and span payloads.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::sync::{Arc, LazyLock};
|
||||
use std::time::Duration;
|
||||
|
||||
use base64::{Engine, engine::general_purpose::URL_SAFE};
|
||||
use litellm_http::Client;
|
||||
use litellm_storage_clickhouse::fetch;
|
||||
use litellm_storage_clickhouse::{Query, fetch};
|
||||
use litellm_traces::{
|
||||
SpanDetail, SpanErrorPage, SpendLookup, Trace, TracePage, listed_summary,
|
||||
query::named as contracts, resolve_trace, to_ui_content,
|
||||
};
|
||||
use moka::future::Cache;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use sha2::{Digest, Sha256};
|
||||
|
||||
use crate::{
|
||||
Connection, Error,
|
||||
query::named::{
|
||||
ListTraces, ListTracesParams, ReadAccessParams, SpanDetail as SpanDetailQuery,
|
||||
SpanDetailParams, SpanError, SpanErrorParams, SpendByResponseIds, SpendByResponseIdsParams,
|
||||
TraceIdentity, TraceIdentityParams, TracePageSpans, TracePageSpansParams, TraceSpans,
|
||||
TraceSpansParams,
|
||||
ListTracesParams, ListTracesRow, ReadAccessParams, SpanDetail as SpanDetailQuery,
|
||||
SpanDetailParams, SpanError, SpanErrorParams, SpendByResponseIdsParams, TraceIdentity,
|
||||
TraceIdentityParams, TraceSpansParams,
|
||||
},
|
||||
};
|
||||
|
||||
struct RunCandidates;
|
||||
|
||||
impl Query for RunCandidates {
|
||||
type Params = ListTracesParams;
|
||||
type Row = ListTracesRow;
|
||||
const SQL: &'static str = concat!(
|
||||
"SELECT * EXCEPT (request_ids), [] AS request_ids FROM (",
|
||||
include_str!("../query/list_traces.sql"),
|
||||
") ORDER BY start_ms DESC, trace_ref DESC"
|
||||
);
|
||||
}
|
||||
|
||||
// Cursor pages share a bounded snapshot so advancing does not resolve the whole graph again.
|
||||
static TRACE_SNAPSHOTS: LazyLock<Cache<String, Arc<Trace>>> = LazyLock::new(|| {
|
||||
Cache::builder()
|
||||
.max_capacity(64 * 1024 * 1024)
|
||||
.weigher(|_: &String, trace: &Arc<Trace>| {
|
||||
serde_json::to_vec(trace.as_ref())
|
||||
.ok()
|
||||
.and_then(|bytes| u32::try_from(bytes.len().saturating_mul(2)).ok())
|
||||
.unwrap_or(u32::MAX)
|
||||
})
|
||||
.time_to_live(Duration::from_secs(120))
|
||||
.build()
|
||||
});
|
||||
|
||||
const NANOS_PER_MS: i64 = 1_000_000;
|
||||
const SPEND_WINDOW_MS: i64 = 30 * 60 * 1000;
|
||||
|
||||
|
|
@ -121,8 +149,8 @@ async fn spend(
|
|||
start_ms: start_ns.div_euclid(NANOS_PER_MS) - SPEND_WINDOW_MS,
|
||||
end_ms: end_ns.div_euclid(NANOS_PER_MS) + SPEND_WINDOW_MS,
|
||||
});
|
||||
match fetch::<SpendByResponseIds>(client, connection, ¶ms).await {
|
||||
Ok(rows) => rows.into_iter().map(|row| row.0).collect(),
|
||||
match crate::span_batches::read_spend(client, connection, params).await {
|
||||
Ok(rows) => rows,
|
||||
Err(error) => {
|
||||
tracing::warn!(%error, "trace spend lookup unavailable");
|
||||
Vec::new()
|
||||
|
|
@ -139,71 +167,43 @@ pub async fn list_traces(
|
|||
cursor: Option<&str>,
|
||||
limit: u32,
|
||||
) -> Result<TracePage, Error> {
|
||||
if limit == 0 {
|
||||
return Err(Error::InvalidParameters);
|
||||
}
|
||||
let (cursor_ms, cursor_trace_id) = trace_position(cursor)?;
|
||||
let params = ListTracesParams::from(contracts::ListTracesParams {
|
||||
let mut params = ListTracesParams::from(contracts::ListTracesParams {
|
||||
access: access.clone(),
|
||||
start_ms,
|
||||
end_ms,
|
||||
cursor_ms,
|
||||
cursor_trace_id,
|
||||
limit,
|
||||
limit: limit.min(500),
|
||||
});
|
||||
let page: Vec<contracts::ListTracesRow> = fetch::<ListTraces>(client, connection, ¶ms)
|
||||
.await?
|
||||
.into_iter()
|
||||
.map(|row| row.0)
|
||||
.collect();
|
||||
let page: Vec<contracts::ListTracesRow> = loop {
|
||||
match fetch::<RunCandidates>(client, connection, ¶ms).await {
|
||||
Err(litellm_storage_clickhouse::Error::ResponseTooLarge) if params.0.limit > 1 => {
|
||||
params.0.limit /= 2;
|
||||
}
|
||||
Err(litellm_storage_clickhouse::Error::ResponseTooLarge) => {
|
||||
return Err(Error::ReadTooLarge);
|
||||
}
|
||||
result => break result?.into_iter().map(|row| row.0).collect(),
|
||||
}
|
||||
};
|
||||
let next_cursor = page
|
||||
.last()
|
||||
.filter(|_| page.len() == limit as usize)
|
||||
.filter(|_| page.len() == params.0.limit as usize)
|
||||
.map(|last| encode_cursor(&(last.start_ms, &last.trace_ref)));
|
||||
let (Some(page_start), Some(page_end)) = (
|
||||
page.iter().map(|row| row.start_ms).min(),
|
||||
page.iter().map(|row| row.start_ms + row.duration_ms).max(),
|
||||
) else {
|
||||
return Ok(TracePage {
|
||||
data: Vec::new(),
|
||||
next_cursor,
|
||||
});
|
||||
};
|
||||
let span_params = TracePageSpansParams::from(contracts::TracePageSpansParams {
|
||||
access: access.clone(),
|
||||
trace_refs: page.iter().map(|row| row.trace_ref.clone()).collect(),
|
||||
start_ms: page_start,
|
||||
end_ms: page_end + 1,
|
||||
});
|
||||
let span_rows: Vec<contracts::TraceSpansRow> =
|
||||
fetch::<TracePageSpans>(client, connection, &span_params)
|
||||
.await?
|
||||
.into_iter()
|
||||
.map(|row| row.0)
|
||||
.collect();
|
||||
let spend_rows = spend(client, connection, access, &span_rows).await;
|
||||
let mut by_trace: HashMap<(String, String, String), Vec<contracts::TraceSpansRow>> =
|
||||
HashMap::new();
|
||||
for span in span_rows {
|
||||
let key = (
|
||||
span.team_id.clone(),
|
||||
span.api_key_hash.clone(),
|
||||
span.trace_id.clone(),
|
||||
);
|
||||
by_trace.entry(key).or_default().push(span);
|
||||
let mut data = Vec::with_capacity(page.len());
|
||||
for row in &page {
|
||||
let summary =
|
||||
match get_trace(client, connection, access, &row.trace_id, &row.trace_ref).await {
|
||||
Ok(trace) => trace.map_or_else(|| listed_summary(row), |trace| trace.summary),
|
||||
Err(Error::ReadTooLarge) => listed_summary(row),
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
data.push(summary);
|
||||
}
|
||||
let data = page
|
||||
.iter()
|
||||
.map(|row| {
|
||||
let spans = by_trace
|
||||
.get(&(
|
||||
row.team_id.clone(),
|
||||
row.api_key_hash.clone(),
|
||||
row.trace_id.clone(),
|
||||
))
|
||||
.map(Vec::as_slice)
|
||||
.unwrap_or_default();
|
||||
resolve_trace(&row.trace_id, &row.trace_ref, spans, &spend_rows)
|
||||
.map_or_else(|| listed_summary(row), |trace| trace.summary)
|
||||
})
|
||||
.collect();
|
||||
Ok(TracePage { data, next_cursor })
|
||||
}
|
||||
|
||||
|
|
@ -222,11 +222,7 @@ pub async fn get_trace(
|
|||
trace_id: trace_id.to_owned(),
|
||||
trace_ref: trace_ref.clone(),
|
||||
};
|
||||
let rows: Vec<contracts::TraceSpansRow> = fetch::<TraceSpans>(client, connection, ¶ms)
|
||||
.await?
|
||||
.into_iter()
|
||||
.map(|row| row.0)
|
||||
.collect();
|
||||
let rows = crate::span_batches::read_spans(client, connection, params, u64::MAX).await?;
|
||||
if rows.is_empty() {
|
||||
return Ok(None);
|
||||
}
|
||||
|
|
@ -234,6 +230,125 @@ pub async fn get_trace(
|
|||
Ok(resolve_trace(trace_id, &trace_ref, &rows, &spend_rows))
|
||||
}
|
||||
|
||||
#[derive(Deserialize, Serialize)]
|
||||
struct SpanPosition {
|
||||
trace_ref: String,
|
||||
snapshot_ms: u64,
|
||||
offset: usize,
|
||||
version: String,
|
||||
}
|
||||
|
||||
pub async fn get_trace_page(
|
||||
client: &Client,
|
||||
connection: &Connection,
|
||||
access: &ReadAccessParams,
|
||||
trace_id: &str,
|
||||
trace_ref: &str,
|
||||
cursor: Option<&str>,
|
||||
page_size: u32,
|
||||
) -> Result<Option<Trace>, Error> {
|
||||
if !(1..=500).contains(&page_size) {
|
||||
return Err(Error::InvalidParameters);
|
||||
}
|
||||
let Some(trace_ref) = reference(client, connection, access, trace_id, trace_ref).await? else {
|
||||
return Ok(None);
|
||||
};
|
||||
let position = match cursor {
|
||||
Some(cursor) => {
|
||||
let position: SpanPosition = decode_cursor(cursor, "span")?;
|
||||
if position.trace_ref != trace_ref || position.snapshot_ms == 0 {
|
||||
return Err(Error::InvalidCursor("span"));
|
||||
}
|
||||
position
|
||||
}
|
||||
None => SpanPosition {
|
||||
trace_ref: trace_ref.clone(),
|
||||
snapshot_ms: (time::OffsetDateTime::now_utc().unix_timestamp_nanos() / 1_000_000)
|
||||
as u64,
|
||||
offset: 0,
|
||||
version: String::new(),
|
||||
},
|
||||
};
|
||||
let key_bytes = serde_json::to_vec(&(
|
||||
connection.url().as_str(),
|
||||
access,
|
||||
trace_id,
|
||||
&trace_ref,
|
||||
position.snapshot_ms,
|
||||
))
|
||||
.map_err(|_| Error::InvalidParameters)?;
|
||||
let key = format!("{:x}", Sha256::digest(key_bytes));
|
||||
let snapshot = if let Some(trace) = TRACE_SNAPSHOTS.get(&key).await {
|
||||
trace
|
||||
} else {
|
||||
let params = TraceSpansParams {
|
||||
access: access.clone(),
|
||||
trace_id: trace_id.to_owned(),
|
||||
trace_ref: trace_ref.clone(),
|
||||
};
|
||||
let rows =
|
||||
crate::span_batches::read_spans(client, connection, params, position.snapshot_ms)
|
||||
.await?;
|
||||
let spend_rows = spend(client, connection, access, &rows).await;
|
||||
let Some(trace) = resolve_trace(trace_id, &trace_ref, &rows, &spend_rows) else {
|
||||
return Ok(None);
|
||||
};
|
||||
let trace = Arc::new(trace);
|
||||
TRACE_SNAPSHOTS.insert(key, Arc::clone(&trace)).await;
|
||||
trace
|
||||
};
|
||||
let span_ids: Vec<&str> = snapshot
|
||||
.spans
|
||||
.iter()
|
||||
.map(|span| span.span_id.as_str())
|
||||
.collect();
|
||||
let version = format!(
|
||||
"{:x}",
|
||||
Sha256::digest(serde_json::to_vec(&span_ids).map_err(|_| Error::InvalidResponse)?)
|
||||
);
|
||||
if cursor.is_some() && position.version != version {
|
||||
return Err(Error::TraceChanged);
|
||||
}
|
||||
let mut trace = Trace {
|
||||
summary: snapshot.summary.clone(),
|
||||
agents: snapshot.agents.clone(),
|
||||
spans: Vec::new(),
|
||||
next_cursor: None,
|
||||
};
|
||||
if position.offset > snapshot.spans.len() {
|
||||
return Err(Error::InvalidCursor("span"));
|
||||
}
|
||||
let end = position
|
||||
.offset
|
||||
.saturating_add(page_size as usize)
|
||||
.min(snapshot.spans.len());
|
||||
trace.next_cursor = (end < snapshot.spans.len()).then(|| {
|
||||
encode_cursor(&SpanPosition {
|
||||
offset: end,
|
||||
version: version.clone(),
|
||||
..position
|
||||
})
|
||||
});
|
||||
trace.spans = snapshot.spans[position.offset..end].to_vec();
|
||||
while serde_json::to_vec(&trace)
|
||||
.map_err(|_| Error::InvalidResponse)?
|
||||
.len()
|
||||
> litellm_storage_clickhouse::READ_LIMITS.response_bytes
|
||||
{
|
||||
if trace.spans.len() <= 1 {
|
||||
return Err(Error::ReadTooLarge);
|
||||
}
|
||||
trace.spans.truncate(trace.spans.len() / 2);
|
||||
trace.next_cursor = Some(encode_cursor(&SpanPosition {
|
||||
trace_ref: trace_ref.clone(),
|
||||
snapshot_ms: position.snapshot_ms,
|
||||
offset: position.offset + trace.spans.len(),
|
||||
version: version.clone(),
|
||||
}));
|
||||
}
|
||||
Ok(Some(trace))
|
||||
}
|
||||
|
||||
pub async fn get_span(
|
||||
client: &Client,
|
||||
connection: &Connection,
|
||||
|
|
|
|||
181
litellm-rust/crates/traces-clickhouse/src/span_batches.rs
Normal file
181
litellm-rust/crates/traces-clickhouse/src/span_batches.rs
Normal file
|
|
@ -0,0 +1,181 @@
|
|||
use litellm_http::Client;
|
||||
use litellm_storage_clickhouse::{Query, fetch};
|
||||
use litellm_traces::query::named as contracts;
|
||||
use serde::Serialize;
|
||||
|
||||
use crate::{Connection, Error, query::named::TraceSpansRow};
|
||||
|
||||
const PAGE_SIZE: u32 = 256;
|
||||
const MAX_GRAPH_BYTES: usize = 64 * 1024 * 1024;
|
||||
const MAX_GRAPH_SPANS: usize = 100_000;
|
||||
|
||||
#[derive(Default)]
|
||||
struct ReadBudget {
|
||||
bytes: usize,
|
||||
rows: usize,
|
||||
}
|
||||
|
||||
impl ReadBudget {
|
||||
fn reserve(&mut self, bytes: usize) -> Result<(), Error> {
|
||||
self.bytes = self.bytes.saturating_add(bytes);
|
||||
if self.bytes > MAX_GRAPH_BYTES || self.rows == MAX_GRAPH_SPANS {
|
||||
return Err(Error::ReadTooLarge);
|
||||
}
|
||||
self.rows += 1;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn record(&mut self, row: &impl Serialize) -> Result<(), Error> {
|
||||
let bytes = serde_json::to_vec(row).map_err(|_| Error::InvalidResponse)?;
|
||||
self.reserve(bytes.len())
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
struct Parameters {
|
||||
#[serde(flatten)]
|
||||
trace: contracts::TraceSpansParams,
|
||||
after_span_id: String,
|
||||
page_size: u32,
|
||||
snapshot_ms: u64,
|
||||
}
|
||||
|
||||
struct SpanBatch;
|
||||
|
||||
impl Query for SpanBatch {
|
||||
type Params = Parameters;
|
||||
type Row = TraceSpansRow;
|
||||
|
||||
const SQL: &'static str = include_str!("../query/trace_span_batch.sql");
|
||||
}
|
||||
|
||||
pub(crate) async fn read_spans(
|
||||
client: &Client,
|
||||
connection: &Connection,
|
||||
trace: contracts::TraceSpansParams,
|
||||
snapshot_ms: u64,
|
||||
) -> Result<Vec<contracts::TraceSpansRow>, Error> {
|
||||
let mut parameters = Parameters {
|
||||
trace,
|
||||
after_span_id: String::new(),
|
||||
page_size: PAGE_SIZE,
|
||||
snapshot_ms,
|
||||
};
|
||||
let mut spans = Vec::new();
|
||||
let mut budget = ReadBudget::default();
|
||||
loop {
|
||||
let page = match fetch::<SpanBatch>(client, connection, ¶meters).await {
|
||||
Err(litellm_storage_clickhouse::Error::ResponseTooLarge)
|
||||
if parameters.page_size > 1 =>
|
||||
{
|
||||
parameters.page_size /= 2;
|
||||
continue;
|
||||
}
|
||||
Err(litellm_storage_clickhouse::Error::ResponseTooLarge) => {
|
||||
return Err(Error::ReadTooLarge);
|
||||
}
|
||||
result => result?,
|
||||
};
|
||||
let complete = page.len() < parameters.page_size as usize;
|
||||
if let Some(last) = page.last() {
|
||||
parameters.after_span_id.clone_from(&last.0.span_id);
|
||||
}
|
||||
for row in page {
|
||||
budget.record(&row)?;
|
||||
spans.push(row.0);
|
||||
}
|
||||
if complete {
|
||||
spans.sort_by_key(|row| row.start_ns);
|
||||
return Ok(spans);
|
||||
}
|
||||
parameters.page_size = (parameters.page_size * 2).min(PAGE_SIZE);
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
struct SpendParameters {
|
||||
#[serde(flatten)]
|
||||
lookup: crate::query::named::SpendByResponseIdsParams,
|
||||
has_cursor: u8,
|
||||
after_team: String,
|
||||
after_ms: i64,
|
||||
after_id: String,
|
||||
page_size: u32,
|
||||
}
|
||||
|
||||
struct SpendBatch;
|
||||
|
||||
impl Query for SpendBatch {
|
||||
type Params = SpendParameters;
|
||||
type Row = crate::query::named::SpendByResponseIdsRow;
|
||||
|
||||
const SQL: &'static str = include_str!("../query/spend_batch.sql");
|
||||
}
|
||||
|
||||
pub(crate) async fn read_spend(
|
||||
client: &Client,
|
||||
connection: &Connection,
|
||||
lookup: crate::query::named::SpendByResponseIdsParams,
|
||||
) -> Result<Vec<contracts::SpendByResponseIdsRow>, Error> {
|
||||
let mut parameters = SpendParameters {
|
||||
lookup,
|
||||
has_cursor: 0,
|
||||
after_team: String::new(),
|
||||
after_ms: 0,
|
||||
after_id: String::new(),
|
||||
page_size: PAGE_SIZE,
|
||||
};
|
||||
let mut rows = Vec::new();
|
||||
let mut budget = ReadBudget::default();
|
||||
loop {
|
||||
let page = match fetch::<SpendBatch>(client, connection, ¶meters).await {
|
||||
Err(litellm_storage_clickhouse::Error::ResponseTooLarge)
|
||||
if parameters.page_size > 1 =>
|
||||
{
|
||||
parameters.page_size /= 2;
|
||||
continue;
|
||||
}
|
||||
Err(litellm_storage_clickhouse::Error::ResponseTooLarge) => {
|
||||
return Err(Error::ReadTooLarge);
|
||||
}
|
||||
result => result?,
|
||||
};
|
||||
let complete = page.len() < parameters.page_size as usize;
|
||||
if let Some(last) = page.last() {
|
||||
parameters.has_cursor = 1;
|
||||
parameters.after_team.clone_from(&last.0.team_id);
|
||||
parameters.after_ms = last.0.start_ms;
|
||||
parameters.after_id.clone_from(&last.0.request_id);
|
||||
}
|
||||
for row in page {
|
||||
budget.record(&row)?;
|
||||
rows.push(row.0);
|
||||
}
|
||||
if complete {
|
||||
return Ok(rows);
|
||||
}
|
||||
parameters.page_size = (parameters.page_size * 2).min(PAGE_SIZE);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use rstest::rstest;
|
||||
|
||||
#[rstest]
|
||||
#[case::byte_boundary(MAX_GRAPH_BYTES - 1, 0, 1, false)]
|
||||
#[case::byte_overflow(MAX_GRAPH_BYTES - 1, 0, 2, true)]
|
||||
#[case::integer_overflow(MAX_GRAPH_BYTES, 0, usize::MAX, true)]
|
||||
#[case::row_boundary(0, MAX_GRAPH_SPANS - 1, 1, false)]
|
||||
#[case::row_overflow(0, MAX_GRAPH_SPANS, 1, true)]
|
||||
fn accumulation_stops_at_the_graph_budget(
|
||||
#[case] bytes: usize,
|
||||
#[case] rows: usize,
|
||||
#[case] next: usize,
|
||||
#[case] rejected: bool,
|
||||
) {
|
||||
let mut budget = ReadBudget { bytes, rows };
|
||||
assert_eq!(budget.reserve(next).is_err(), rejected);
|
||||
}
|
||||
}
|
||||
436
litellm-rust/crates/traces-clickhouse/tests/reads.rs
Normal file
436
litellm-rust/crates/traces-clickhouse/tests/reads.rs
Normal file
|
|
@ -0,0 +1,436 @@
|
|||
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,
|
||||
};
|
||||
use rstest::rstest;
|
||||
use serde_json::json;
|
||||
|
||||
#[path = "queries/support.rs"]
|
||||
mod fixtures;
|
||||
mod support;
|
||||
|
||||
use fixtures::{DATABASE, SeededDatabase, migrated_database, seeded_database};
|
||||
use support::TestResult;
|
||||
|
||||
#[rstest]
|
||||
#[case::many_runs(50, 21, 0, false)]
|
||||
#[case::one_large_run(1, 1100, 0, false)]
|
||||
#[case::large_rows(1, 280, 20_000, false)]
|
||||
#[case::many_costs(1, 1101, 0, true)]
|
||||
#[tokio::test]
|
||||
async fn large_runs_remain_complete_under_default_reader_limits(
|
||||
#[future(awt)] migrated_database: TestResult<SeededDatabase>,
|
||||
#[case] runs: usize,
|
||||
#[case] steps: usize,
|
||||
#[case] name_bytes: usize,
|
||||
#[case] costed: bool,
|
||||
) -> TestResult {
|
||||
let fixture = migrated_database?;
|
||||
let client = &fixture.database.client;
|
||||
let writer = Connection::writer(&fixture.database.url)?;
|
||||
for run in 0..runs {
|
||||
let rows = (0..steps)
|
||||
.map(|step| {
|
||||
BTreeMap::from([
|
||||
(
|
||||
"Timestamp".into(),
|
||||
json!(1_790_000_000_000_000_000_i64 + step as i64),
|
||||
),
|
||||
("TraceId".into(), json!(format!("trace-{run:04}"))),
|
||||
("SpanId".into(), json!(format!("span-{step:04}"))),
|
||||
(
|
||||
"ParentSpanId".into(),
|
||||
json!(if step == 0 { "" } else { "span-0000" }),
|
||||
),
|
||||
(
|
||||
"SpanName".into(),
|
||||
json!(if name_bytes == 0 {
|
||||
format!("step-{step}")
|
||||
} else {
|
||||
"x".repeat(name_bytes)
|
||||
}),
|
||||
),
|
||||
(
|
||||
"ObservationType".into(),
|
||||
json!(if step == 0 {
|
||||
"agent"
|
||||
} else if costed {
|
||||
"llm"
|
||||
} else {
|
||||
"tool"
|
||||
}),
|
||||
),
|
||||
("TeamId".into(), json!("team-a")),
|
||||
("ApiKeyHash".into(), json!("key-a")),
|
||||
("Duration".into(), json!(1000)),
|
||||
(
|
||||
"LiteLLMRequestId".into(),
|
||||
json!(if costed && step > 0 {
|
||||
format!("response-{step}")
|
||||
} else {
|
||||
String::new()
|
||||
}),
|
||||
),
|
||||
])
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
for chunk in rows.chunks(100) {
|
||||
insert_rows(
|
||||
client,
|
||||
&writer,
|
||||
DATABASE,
|
||||
InsertTable::OtelTraces,
|
||||
chunk.to_vec(),
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
}
|
||||
if costed {
|
||||
let costs = (1..steps)
|
||||
.map(|step| {
|
||||
BTreeMap::from([
|
||||
("request_id".into(), json!(format!("request-{step}"))),
|
||||
("response_id".into(), json!(format!("response-{step}"))),
|
||||
("team_id".into(), json!("team-a")),
|
||||
("api_key".into(), json!("key-a")),
|
||||
("start_time".into(), json!(1_790_000_000_000_i64)),
|
||||
("end_time".into(), json!(1_790_000_000_001_i64)),
|
||||
("spend".into(), json!(0.25)),
|
||||
])
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
insert_rows(client, &writer, DATABASE, InsertTable::SpendLogs, costs).await?;
|
||||
}
|
||||
let reader = fixture
|
||||
.readers
|
||||
.connection(client, &QueryScope::All, "fixture-secret")
|
||||
.await?;
|
||||
let access = ReadAccessParams {
|
||||
all_teams: false,
|
||||
user_id: String::new(),
|
||||
team_ids: vec!["team-a".into()],
|
||||
};
|
||||
let page = list_traces(client, &reader, &access, 0, 2_000_000_000_000, None, 50).await?;
|
||||
assert_eq!(page.data.len(), runs);
|
||||
for summary in &page.data {
|
||||
assert_eq!(summary.span_count, steps as u64);
|
||||
assert_eq!(
|
||||
if costed {
|
||||
summary.llm_calls
|
||||
} else {
|
||||
summary.tool_calls
|
||||
},
|
||||
(steps - 1) as u64
|
||||
);
|
||||
if costed {
|
||||
assert_eq!(summary.spend, Some((steps - 1) as f64 * 0.25));
|
||||
}
|
||||
}
|
||||
let trace_ref = &page
|
||||
.data
|
||||
.iter()
|
||||
.find(|run| run.trace_id == "trace-0000")
|
||||
.ok_or("missing run")?
|
||||
.trace_ref;
|
||||
let detail = get_trace(client, &reader, &access, "trace-0000", trace_ref)
|
||||
.await?
|
||||
.ok_or("missing trace")?;
|
||||
assert_eq!(detail.spans.len(), steps);
|
||||
assert_eq!(detail.spans[0].span_id, "span-0000");
|
||||
assert_eq!(
|
||||
detail.spans[steps - 1].span_id,
|
||||
format!("span-{:04}", steps - 1)
|
||||
);
|
||||
assert_eq!(
|
||||
if costed {
|
||||
detail.summary.llm_calls
|
||||
} else {
|
||||
detail.summary.tool_calls
|
||||
},
|
||||
(steps - 1) as u64
|
||||
);
|
||||
let mut cursor = None;
|
||||
let mut ids = Vec::new();
|
||||
loop {
|
||||
let page = get_trace_page(
|
||||
client,
|
||||
&reader,
|
||||
&access,
|
||||
"trace-0000",
|
||||
trace_ref,
|
||||
cursor.as_deref(),
|
||||
200,
|
||||
)
|
||||
.await?
|
||||
.ok_or("missing page")?;
|
||||
assert_eq!(page.summary, detail.summary);
|
||||
assert!(page.spans.len() <= 200);
|
||||
assert!(
|
||||
serde_json::to_vec(&page)?.len()
|
||||
<= litellm_storage_clickhouse::READ_LIMITS.response_bytes
|
||||
);
|
||||
ids.extend(page.spans.into_iter().map(|span| span.span_id));
|
||||
cursor = page.next_cursor;
|
||||
if cursor.is_none() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
assert_eq!(
|
||||
ids,
|
||||
detail
|
||||
.spans
|
||||
.iter()
|
||||
.map(|span| span.span_id.clone())
|
||||
.collect::<Vec<_>>()
|
||||
);
|
||||
let denied = ReadAccessParams {
|
||||
team_ids: vec!["other-team".into()],
|
||||
..access
|
||||
};
|
||||
assert!(
|
||||
get_trace(client, &reader, &denied, "trace-0000", trace_ref)
|
||||
.await?
|
||||
.is_none()
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[tokio::test]
|
||||
async fn cursor_pages_keep_a_tenant_scoped_snapshot_when_more_spans_arrive(
|
||||
#[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 listed = list_traces(client, &reader, &access, 0, 2_000_000_000_000, None, 10).await?;
|
||||
let summary = listed
|
||||
.data
|
||||
.iter()
|
||||
.find(|summary| summary.span_count == 3)
|
||||
.ok_or("missing fixture")?;
|
||||
let first = get_trace_page(
|
||||
client,
|
||||
&reader,
|
||||
&access,
|
||||
&summary.trace_id,
|
||||
&summary.trace_ref,
|
||||
None,
|
||||
1,
|
||||
)
|
||||
.await?
|
||||
.ok_or("missing first page")?;
|
||||
let original_ids = get_trace(
|
||||
client,
|
||||
&reader,
|
||||
&access,
|
||||
&summary.trace_id,
|
||||
&summary.trace_ref,
|
||||
)
|
||||
.await?
|
||||
.ok_or("missing trace")?
|
||||
.spans
|
||||
.into_iter()
|
||||
.map(|span| span.span_id)
|
||||
.collect::<Vec<_>>();
|
||||
let writer = Connection::writer(&fixture.database.url)?;
|
||||
insert_rows(
|
||||
client,
|
||||
&writer,
|
||||
DATABASE,
|
||||
InsertTable::OtelTraces,
|
||||
vec![BTreeMap::from([
|
||||
("Timestamp".into(), json!(1_790_000_000_000_000_000_i64)),
|
||||
("TraceId".into(), json!(summary.trace_id)),
|
||||
("SpanId".into(), json!("late-span")),
|
||||
("ParentSpanId".into(), json!(first.spans[0].span_id)),
|
||||
("TeamId".into(), json!("team-a")),
|
||||
("ApiKeyHash".into(), json!("key-a")),
|
||||
("EngineReceivedMs".into(), json!(u64::MAX / 2)),
|
||||
])],
|
||||
)
|
||||
.await?;
|
||||
let denied = ReadAccessParams {
|
||||
all_teams: false,
|
||||
user_id: String::new(),
|
||||
team_ids: vec!["not-this-team".into()],
|
||||
};
|
||||
assert!(
|
||||
get_trace_page(
|
||||
client,
|
||||
&reader,
|
||||
&denied,
|
||||
&summary.trace_id,
|
||||
&summary.trace_ref,
|
||||
first.next_cursor.as_deref(),
|
||||
1
|
||||
)
|
||||
.await?
|
||||
.is_none()
|
||||
);
|
||||
let first_cursor = first.next_cursor.clone();
|
||||
let mut cursor = first.next_cursor;
|
||||
let mut ids = first
|
||||
.spans
|
||||
.into_iter()
|
||||
.map(|span| span.span_id)
|
||||
.collect::<Vec<_>>();
|
||||
while let Some(current) = cursor {
|
||||
let next = get_trace_page(
|
||||
client,
|
||||
&reader,
|
||||
&access,
|
||||
&summary.trace_id,
|
||||
&summary.trace_ref,
|
||||
Some(¤t),
|
||||
1,
|
||||
)
|
||||
.await?
|
||||
.ok_or("missing next page")?;
|
||||
assert_eq!(next.summary.span_count, 3);
|
||||
ids.extend(next.spans.into_iter().map(|span| span.span_id));
|
||||
cursor = next.next_cursor;
|
||||
}
|
||||
assert_eq!(ids, original_ids);
|
||||
let refreshed = get_trace(
|
||||
client,
|
||||
&reader,
|
||||
&access,
|
||||
&summary.trace_id,
|
||||
&summary.trace_ref,
|
||||
)
|
||||
.await?
|
||||
.ok_or("missing refreshed trace")?;
|
||||
assert_eq!(refreshed.spans.len(), 4);
|
||||
assert!(matches!(
|
||||
get_trace_page(
|
||||
client,
|
||||
&reader,
|
||||
&access,
|
||||
&summary.trace_id,
|
||||
&summary.trace_ref,
|
||||
Some("invalid"),
|
||||
1
|
||||
)
|
||||
.await,
|
||||
Err(litellm_traces_clickhouse::Error::InvalidCursor("span"))
|
||||
));
|
||||
let backdated = json!({
|
||||
"Timestamp": "2026-09-01 00:00:00.000000000",
|
||||
"TraceId": summary.trace_id,
|
||||
"SpanId": "backdated-span",
|
||||
"EngineReceivedMs": 1,
|
||||
"TeamId": "team-a",
|
||||
"ApiKeyHash": "key-a"
|
||||
});
|
||||
client
|
||||
.post(writer.url().clone())
|
||||
.body(format!(
|
||||
"INSERT INTO {DATABASE}.otel_traces FORMAT JSONEachRow\n{backdated}"
|
||||
))
|
||||
.send()
|
||||
.await?
|
||||
.error_for_status()?;
|
||||
let uncached_reader =
|
||||
Connection::reader(&format!("{}?max_threads=1", fixture.database.url), DATABASE)?;
|
||||
let changed = get_trace_page(
|
||||
client,
|
||||
&uncached_reader,
|
||||
&access,
|
||||
&summary.trace_id,
|
||||
&summary.trace_ref,
|
||||
first_cursor.as_deref(),
|
||||
1,
|
||||
)
|
||||
.await;
|
||||
assert!(
|
||||
matches!(changed, Err(litellm_traces_clickhouse::Error::TraceChanged)),
|
||||
"{changed:?}"
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[tokio::test]
|
||||
async fn an_oversized_span_keeps_the_run_list_available_with_partial_totals(
|
||||
#[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 before = list_traces(client, &reader, &access, 0, 2_000_000_000_000, None, 50).await?;
|
||||
let run = before
|
||||
.data
|
||||
.iter()
|
||||
.find(|run| run.span_count == 3)
|
||||
.ok_or("missing fixture")?;
|
||||
let writer = Connection::writer(&fixture.database.url)?;
|
||||
insert_rows(
|
||||
client,
|
||||
&writer,
|
||||
DATABASE,
|
||||
InsertTable::OtelTraces,
|
||||
vec![BTreeMap::from([
|
||||
("Timestamp".into(), json!(1_790_000_000_000_000_000_i64)),
|
||||
("TraceId".into(), json!(run.trace_id)),
|
||||
("SpanId".into(), json!("oversized-child")),
|
||||
("ParentSpanId".into(), json!("0101010101010101")),
|
||||
(
|
||||
"SpanName".into(),
|
||||
json!("x".repeat(litellm_storage_clickhouse::READ_LIMITS.response_bytes + 1)),
|
||||
),
|
||||
("ObservationType".into(), json!("tool")),
|
||||
("TeamId".into(), json!("team-a")),
|
||||
("ApiKeyHash".into(), json!("key-a")),
|
||||
])],
|
||||
)
|
||||
.await?;
|
||||
let after = list_traces(client, &reader, &access, 0, 2_000_000_000_000, None, 50).await?;
|
||||
assert_eq!(after.data.len(), before.data.len());
|
||||
let limited = after
|
||||
.data
|
||||
.iter()
|
||||
.find(|item| item.trace_ref == run.trace_ref)
|
||||
.ok_or("missing run")?;
|
||||
assert!(limited.resolution_limited);
|
||||
assert_eq!(limited.span_count, 4);
|
||||
assert!(
|
||||
after
|
||||
.data
|
||||
.iter()
|
||||
.filter(|item| item.trace_ref != run.trace_ref)
|
||||
.all(|item| !item.resolution_limited)
|
||||
);
|
||||
assert!(matches!(
|
||||
get_trace_page(
|
||||
client,
|
||||
&reader,
|
||||
&access,
|
||||
&run.trace_id,
|
||||
&run.trace_ref,
|
||||
None,
|
||||
200
|
||||
)
|
||||
.await,
|
||||
Err(litellm_traces_clickhouse::Error::ReadTooLarge)
|
||||
));
|
||||
Ok(())
|
||||
}
|
||||
|
|
@ -171,6 +171,7 @@ pub fn resolve_trace(
|
|||
.map(|(_, (span, _))| span.input_preview.clone())
|
||||
.unwrap_or_default();
|
||||
let summary = TraceSummary {
|
||||
resolution_limited: false,
|
||||
trace_id: trace_id.to_owned(),
|
||||
trace_ref: trace_ref.to_owned(),
|
||||
name: spans[root].name.clone(),
|
||||
|
|
@ -209,11 +210,13 @@ pub fn resolve_trace(
|
|||
summary,
|
||||
agents,
|
||||
spans,
|
||||
next_cursor: None,
|
||||
})
|
||||
}
|
||||
|
||||
pub fn listed_summary(row: &ListTracesRow) -> TraceSummary {
|
||||
TraceSummary {
|
||||
resolution_limited: true,
|
||||
trace_id: row.trace_id.clone(),
|
||||
trace_ref: row.trace_ref.clone(),
|
||||
name: row.name.clone(),
|
||||
|
|
|
|||
|
|
@ -17,7 +17,7 @@ pub enum SpanStatus {
|
|||
}
|
||||
|
||||
#[macro_rules_attribute::apply(response_type)]
|
||||
#[derive(Debug, PartialEq)]
|
||||
#[derive(Clone, Debug, PartialEq)]
|
||||
pub struct Span {
|
||||
pub span_id: String,
|
||||
pub parent_span_id: Option<String>,
|
||||
|
|
@ -41,7 +41,7 @@ pub struct Span {
|
|||
|
||||
/// One distinct agent in a trace: 200 invocations of `researcher` are one node.
|
||||
#[macro_rules_attribute::apply(response_type)]
|
||||
#[derive(Debug, PartialEq)]
|
||||
#[derive(Clone, Debug, PartialEq)]
|
||||
pub struct AgentNode {
|
||||
pub name: String,
|
||||
pub parent_agent: Option<String>,
|
||||
|
|
@ -53,8 +53,10 @@ pub struct AgentNode {
|
|||
}
|
||||
|
||||
#[macro_rules_attribute::apply(response_type)]
|
||||
#[derive(Debug, PartialEq)]
|
||||
#[derive(Clone, Debug, PartialEq)]
|
||||
pub struct TraceSummary {
|
||||
#[cfg_attr(feature = "schema", schemars(extend("x-python-optional" = true)))]
|
||||
pub resolution_limited: bool,
|
||||
pub trace_id: String,
|
||||
#[cfg_attr(feature = "schema", schemars(extend("x-python-optional" = true)))]
|
||||
pub trace_ref: String,
|
||||
|
|
@ -81,11 +83,13 @@ pub struct TraceSummary {
|
|||
}
|
||||
|
||||
#[macro_rules_attribute::apply(response_type)]
|
||||
#[derive(Debug, PartialEq)]
|
||||
#[derive(Clone, Debug, PartialEq)]
|
||||
pub struct Trace {
|
||||
pub summary: TraceSummary,
|
||||
pub agents: Vec<AgentNode>,
|
||||
pub spans: Vec<Span>,
|
||||
#[cfg_attr(feature = "schema", schemars(extend("x-python-optional" = true)))]
|
||||
pub next_cursor: Option<String>,
|
||||
}
|
||||
|
||||
#[macro_rules_attribute::apply(response_type)]
|
||||
|
|
|
|||
|
|
@ -290,6 +290,7 @@ _PRISMA_MODELS: Final[frozenset[str]] = frozenset(
|
|||
"LiteLLM_Config",
|
||||
"LiteLLM_SpendLogs",
|
||||
"LiteLLM_BudgetWindowSpend",
|
||||
"LiteLLM_BackgroundInteractionSettlement",
|
||||
"LiteLLM_ErrorLogs",
|
||||
"LiteLLM_UserNotifications",
|
||||
"LiteLLM_TeamMembership",
|
||||
|
|
|
|||
|
|
@ -140,7 +140,7 @@ async def list_agent_traces(
|
|||
context: Annotated[TraceAccessContext, Depends(provide_trace_access)],
|
||||
start_ms: Annotated[int | None, Query(description="Window start, unix ms. Default: 24h ago")] = None,
|
||||
end_ms: Annotated[int | None, Query(description="Window end, unix ms. Default: now")] = None,
|
||||
cursor: Annotated[str | None, Query()] = None,
|
||||
cursor: Annotated[str | None, Query(max_length=512)] = None,
|
||||
) -> TracePage:
|
||||
now_ms: Final = int(time.time() * 1000)
|
||||
try:
|
||||
|
|
@ -153,6 +153,13 @@ async def list_agent_traces(
|
|||
)
|
||||
except ValueError as error:
|
||||
raise HTTPException(status_code=400, detail=str(error)) from error
|
||||
except OverflowError as error:
|
||||
raise HTTPException(
|
||||
status_code=413, detail="Trace is too large for this view. Use a filtered trace query."
|
||||
) from error
|
||||
except RuntimeError as error:
|
||||
verbose_proxy_logger.warning("Trace read unavailable: %s", error)
|
||||
raise HTTPException(status_code=503, detail="Traces are temporarily unavailable. Please try again.") from error
|
||||
|
||||
|
||||
class TraceQueryRequest(BaseModel):
|
||||
|
|
@ -228,12 +235,21 @@ async def get_agent_trace(
|
|||
trace_id: str,
|
||||
context: Annotated[TraceAccessContext, Depends(provide_trace_access)],
|
||||
trace_ref: Annotated[str, Query()] = "",
|
||||
cursor: Annotated[str | None, Query(max_length=512)] = None,
|
||||
page_size: Annotated[int | None, Query(ge=1, le=500)] = None,
|
||||
) -> Trace:
|
||||
tracing, scope = context.reader()
|
||||
try:
|
||||
trace: Final = await tracing.get_trace(trace_id, scope, trace_ref)
|
||||
trace: Final = await tracing.get_trace(trace_id, scope, trace_ref, cursor, page_size)
|
||||
except ValueError as error:
|
||||
raise HTTPException(status_code=400, detail=str(error)) from error
|
||||
except OverflowError as error:
|
||||
raise HTTPException(
|
||||
status_code=413, detail="Trace is too large for this view. Use a filtered trace query."
|
||||
) from error
|
||||
except RuntimeError as error:
|
||||
verbose_proxy_logger.warning("Trace read unavailable: %s", error)
|
||||
raise HTTPException(status_code=503, detail="Traces are temporarily unavailable. Please try again.") from error
|
||||
if trace is None:
|
||||
raise HTTPException(status_code=404, detail=f"Trace {trace_id} not found")
|
||||
return trace
|
||||
|
|
@ -251,6 +267,13 @@ async def get_agent_trace_span(
|
|||
span: Final = await tracing.get_span(trace_id, span_id, scope, trace_ref)
|
||||
except ValueError as error:
|
||||
raise HTTPException(status_code=400, detail=str(error)) from error
|
||||
except OverflowError as error:
|
||||
raise HTTPException(
|
||||
status_code=413, detail="Trace is too large for this view. Use a filtered trace query."
|
||||
) from error
|
||||
except RuntimeError as error:
|
||||
verbose_proxy_logger.warning("Trace read unavailable: %s", error)
|
||||
raise HTTPException(status_code=503, detail="Traces are temporarily unavailable. Please try again.") from error
|
||||
if span is None:
|
||||
raise HTTPException(status_code=404, detail=f"Span {span_id} not found")
|
||||
return span
|
||||
|
|
@ -269,6 +292,13 @@ async def get_agent_trace_span_error(
|
|||
page: Final = await tracing.get_span_error(trace_id, span_id, scope, trace_ref, cursor)
|
||||
except ValueError as error:
|
||||
raise HTTPException(status_code=400, detail=str(error)) from error
|
||||
except OverflowError as error:
|
||||
raise HTTPException(
|
||||
status_code=413, detail="Trace is too large for this view. Use a filtered trace query."
|
||||
) from error
|
||||
except RuntimeError as error:
|
||||
verbose_proxy_logger.warning("Trace read unavailable: %s", error)
|
||||
raise HTTPException(status_code=503, detail="Traces are temporarily unavailable. Please try again.") from error
|
||||
if page is None:
|
||||
raise HTTPException(status_code=404, detail="Span diagnostic not found or no longer available")
|
||||
return page
|
||||
|
|
|
|||
|
|
@ -45,7 +45,9 @@ class NativeTraceStorage:
|
|||
def list_traces(
|
||||
self, scope: TraceScope, start_ms: int, end_ms: int, cursor: str | None, limit: int
|
||||
) -> Future[JsonValue]: ...
|
||||
def get_trace(self, trace_id: str, scope: TraceScope, trace_ref: str) -> Future[JsonValue]: ...
|
||||
def get_trace(
|
||||
self, trace_id: str, scope: TraceScope, trace_ref: str, cursor: str | None = None, page_size: int | None = None
|
||||
) -> Future[JsonValue]: ...
|
||||
def get_span(self, trace_id: str, span_id: str, scope: TraceScope, trace_ref: str) -> Future[JsonValue]: ...
|
||||
def get_span_error(
|
||||
self, trace_id: str, span_id: str, scope: TraceScope, trace_ref: str, cursor: str | None
|
||||
|
|
|
|||
|
|
@ -99,6 +99,7 @@ class UIMessage(typing_extensions.TypedDict):
|
|||
|
||||
|
||||
class TraceSummary(typing_extensions.TypedDict):
|
||||
resolution_limited: ReadOnly[NotRequired[bool]]
|
||||
trace_id: ReadOnly[str]
|
||||
trace_ref: ReadOnly[NotRequired[str]]
|
||||
name: ReadOnly[str]
|
||||
|
|
@ -145,6 +146,7 @@ class Trace(typing_extensions.TypedDict):
|
|||
summary: ReadOnly[TraceSummary]
|
||||
agents: ReadOnly[tuple[AgentNode, ...]]
|
||||
spans: ReadOnly[tuple[Span, ...]]
|
||||
next_cursor: ReadOnly[NotRequired[str | None]]
|
||||
|
||||
|
||||
class TracePage(typing_extensions.TypedDict):
|
||||
|
|
|
|||
|
|
@ -67,7 +67,9 @@ class NativeStore(Protocol):
|
|||
self, scope: TraceScope, start_ms: int, end_ms: int, cursor: str | None, limit: int
|
||||
) -> Awaitable[JsonValue]: ...
|
||||
|
||||
def get_trace(self, trace_id: str, scope: TraceScope, trace_ref: str) -> Awaitable[JsonValue]: ...
|
||||
def get_trace(
|
||||
self, trace_id: str, scope: TraceScope, trace_ref: str, cursor: str | None = None, page_size: int | None = None
|
||||
) -> Awaitable[JsonValue]: ...
|
||||
|
||||
def get_span(self, trace_id: str, span_id: str, scope: TraceScope, trace_ref: str) -> Awaitable[JsonValue]: ...
|
||||
|
||||
|
|
@ -189,8 +191,15 @@ class ClickHouseStorage:
|
|||
result: Final = await self._native.list_traces(scope, start_ms, end_ms, cursor, limit)
|
||||
return _validate_query_response(_TRACE_PAGE, result)
|
||||
|
||||
async def get_trace(self, trace_id: str, scope: TraceScope, trace_ref: str = "") -> Trace | None:
|
||||
result: Final = await self._native.get_trace(trace_id, scope, trace_ref)
|
||||
async def get_trace(
|
||||
self,
|
||||
trace_id: str,
|
||||
scope: TraceScope,
|
||||
trace_ref: str = "",
|
||||
cursor: str | None = None,
|
||||
page_size: int | None = None,
|
||||
) -> Trace | None:
|
||||
result: Final = await self._native.get_trace(trace_id, scope, trace_ref, cursor, page_size)
|
||||
return _validate_query_response(_TRACE, result)
|
||||
|
||||
async def get_span(self, trace_id: str, span_id: str, scope: TraceScope, trace_ref: str = "") -> SpanDetail | None:
|
||||
|
|
|
|||
|
|
@ -99,8 +99,15 @@ class TraceReceiver:
|
|||
async def list_traces(self, scope: TraceScope, start_ms: int, end_ms: int, cursor: str | None = None) -> TracePage:
|
||||
return await self.storage.list_traces(scope, start_ms, end_ms, cursor, AGENT_TRACING_LIST_PAGE_SIZE)
|
||||
|
||||
async def get_trace(self, trace_id: str, scope: TraceScope, trace_ref: str = "") -> Trace | None:
|
||||
return await self.storage.get_trace(trace_id, scope, trace_ref)
|
||||
async def get_trace(
|
||||
self,
|
||||
trace_id: str,
|
||||
scope: TraceScope,
|
||||
trace_ref: str = "",
|
||||
cursor: str | None = None,
|
||||
page_size: int | None = None,
|
||||
) -> Trace | None:
|
||||
return await self.storage.get_trace(trace_id, scope, trace_ref, cursor, page_size)
|
||||
|
||||
async def get_span(self, trace_id: str, span_id: str, scope: TraceScope, trace_ref: str = "") -> SpanDetail | None:
|
||||
return await self.storage.get_span(trace_id, span_id, scope, trace_ref)
|
||||
|
|
|
|||
|
|
@ -245,6 +245,10 @@
|
|||
"minimum": 0,
|
||||
"type": "integer"
|
||||
},
|
||||
"resolution_limited": {
|
||||
"type": "boolean",
|
||||
"x-python-optional": true
|
||||
},
|
||||
"service": {
|
||||
"type": "string"
|
||||
},
|
||||
|
|
@ -282,6 +286,7 @@
|
|||
}
|
||||
},
|
||||
"required": [
|
||||
"resolution_limited",
|
||||
"trace_id",
|
||||
"trace_ref",
|
||||
"name",
|
||||
|
|
@ -314,6 +319,13 @@
|
|||
},
|
||||
"type": "array"
|
||||
},
|
||||
"next_cursor": {
|
||||
"type": [
|
||||
"string",
|
||||
"null"
|
||||
],
|
||||
"x-python-optional": true
|
||||
},
|
||||
"spans": {
|
||||
"items": {
|
||||
"$ref": "#/$defs/Span"
|
||||
|
|
@ -327,7 +339,8 @@
|
|||
"required": [
|
||||
"summary",
|
||||
"agents",
|
||||
"spans"
|
||||
"spans",
|
||||
"next_cursor"
|
||||
],
|
||||
"title": "Trace",
|
||||
"type": "object"
|
||||
|
|
|
|||
|
|
@ -76,6 +76,10 @@
|
|||
"minimum": 0,
|
||||
"type": "integer"
|
||||
},
|
||||
"resolution_limited": {
|
||||
"type": "boolean",
|
||||
"x-python-optional": true
|
||||
},
|
||||
"service": {
|
||||
"type": "string"
|
||||
},
|
||||
|
|
@ -113,6 +117,7 @@
|
|||
}
|
||||
},
|
||||
"required": [
|
||||
"resolution_limited",
|
||||
"trace_id",
|
||||
"trace_ref",
|
||||
"name",
|
||||
|
|
|
|||
|
|
@ -61,11 +61,12 @@ async def test_ingest_rejects_oversized_body_before_storage() -> None:
|
|||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_reads_delegate_to_storage() -> None:
|
||||
@pytest.mark.parametrize("cursor,page_size", ((None, None), ("next", 200)))
|
||||
async def test_reads_delegate_to_storage(cursor: str | None, page_size: int | None) -> None:
|
||||
storage: Final = _fake_storage()
|
||||
scope: Final[TraceScope] = {"all_teams": 0, "user_id": "", "team_ids": ("team-research",)}
|
||||
assert await TraceReceiver(storage).get_trace("t1", scope) is None
|
||||
storage.get_trace.assert_awaited_once_with("t1", scope, "")
|
||||
assert await TraceReceiver(storage).get_trace("t1", scope, "", cursor, page_size) is None
|
||||
storage.get_trace.assert_awaited_once_with("t1", scope, "", cursor, page_size)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
|
|
|
|||
|
|
@ -266,7 +266,7 @@ def test_get_trace_404_and_200(client, receiver):
|
|||
response = client.get("/v1/traces/t1")
|
||||
assert response.status_code == 200
|
||||
assert response.json() == TRACE_RESPONSE
|
||||
receiver.get_trace.assert_awaited_with("t1", {"all_teams": 0, "user_id": "user", "team_ids": ()}, "")
|
||||
receiver.get_trace.assert_awaited_with("t1", {"all_teams": 0, "user_id": "user", "team_ids": ()}, "", None, None)
|
||||
|
||||
|
||||
def test_get_span_404_and_200(client, receiver):
|
||||
|
|
@ -278,10 +278,45 @@ def test_get_span_404_and_200(client, receiver):
|
|||
receiver.get_span.assert_awaited_with("t1", "s1", {"all_teams": 0, "user_id": "user", "team_ids": ()}, "")
|
||||
|
||||
|
||||
def test_trace_detail_passes_scoped_reference(client, receiver):
|
||||
@pytest.mark.parametrize("suffix,cursor,page_size", [("", None, None), ("&cursor=next&page_size=200", "next", 200)])
|
||||
def test_trace_detail_passes_scoped_reference(client, receiver, suffix, cursor, page_size):
|
||||
receiver.get_trace.return_value = TRACE_RESPONSE
|
||||
assert client.get("/v1/traces/t1?trace_ref=run-one").status_code == 200
|
||||
receiver.get_trace.assert_awaited_with("t1", {"all_teams": 0, "user_id": "user", "team_ids": ()}, "run-one")
|
||||
assert client.get(f"/v1/traces/t1?trace_ref=run-one{suffix}").status_code == 200
|
||||
receiver.get_trace.assert_awaited_with(
|
||||
"t1", {"all_teams": 0, "user_id": "user", "team_ids": ()}, "run-one", cursor, page_size
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"path,method",
|
||||
(
|
||||
("/v1/traces", "list_traces"),
|
||||
("/v1/traces/t1", "get_trace"),
|
||||
("/v1/traces/t1/spans/s1", "get_span"),
|
||||
("/v1/traces/t1/spans/s1/error", "get_span_error"),
|
||||
),
|
||||
)
|
||||
@pytest.mark.parametrize(
|
||||
"error,status,message",
|
||||
(
|
||||
(RuntimeError("private database details"), 503, "Traces are temporarily unavailable. Please try again."),
|
||||
(OverflowError("private query details"), 413, "Trace is too large for this view. Use a filtered trace query."),
|
||||
),
|
||||
)
|
||||
def test_read_failures_are_actionable_without_exposing_database_details(
|
||||
client: TestClient, receiver: MagicMock, path: str, method: str, error: Exception, status: int, message: str
|
||||
) -> None:
|
||||
getattr(receiver, method).side_effect = error
|
||||
response: Final = client.get(path)
|
||||
assert response.status_code == status
|
||||
assert response.json() == {"detail": message}
|
||||
|
||||
|
||||
@pytest.mark.parametrize("query", ("page_size=0", "page_size=501", "cursor=" + "x" * 513))
|
||||
def test_trace_page_rejects_unbounded_parameters(client: TestClient, receiver: MagicMock, query: str) -> None:
|
||||
response: Final = client.get(f"/v1/traces/t1?{query}")
|
||||
assert response.status_code == 422
|
||||
receiver.get_trace.assert_not_awaited()
|
||||
|
||||
|
||||
def test_invalid_export_and_cursor_are_client_errors(client, receiver):
|
||||
|
|
@ -370,7 +405,9 @@ def test_injected_receiver_ingests_with_the_authenticated_tenant(client: TestCli
|
|||
storage: Final = MagicMock(spec=ClickHouseStorage)
|
||||
storage.ingest = AsyncMock(return_value=1)
|
||||
client.app.dependency_overrides[tracing_endpoints.provide_receiver] = lambda: TraceReceiver(storage)
|
||||
response: Final = client.post("/v1/traces", content=b'{"resourceSpans": []}', headers={"content-type": "application/json"})
|
||||
response: Final = client.post(
|
||||
"/v1/traces", content=b'{"resourceSpans": []}', headers={"content-type": "application/json"}
|
||||
)
|
||||
assert response.status_code == 200, response.text
|
||||
assert response.json() == {}
|
||||
storage.ingest.assert_awaited_once_with(
|
||||
|
|
@ -538,9 +575,7 @@ def test_sql_and_help_use_authenticated_scope(
|
|||
result: Final = client.post("/v1/traces/query", json={"sql": "SELECT * FROM otel_traces"})
|
||||
assert result.status_code == 200, result.text
|
||||
assert result.json() == SQL_ENVELOPE
|
||||
receiver.storage.query_sql.assert_awaited_once_with(
|
||||
"SELECT * FROM otel_traces", expected_scope, "test-secret"
|
||||
)
|
||||
receiver.storage.query_sql.assert_awaited_once_with("SELECT * FROM otel_traces", expected_scope, "test-secret")
|
||||
help_result: Final = client.get("/v1/traces/query/help")
|
||||
assert help_result.status_code == 200, help_result.text
|
||||
assert help_result.json() == QUERY_HELP
|
||||
|
|
@ -692,7 +727,7 @@ class _NativeConfig:
|
|||
|
||||
|
||||
class _NativeReturningHelp(ModuleType):
|
||||
def __init__(self, help_payload: Mapping[str, object]) -> None:
|
||||
def __init__(self, help_payload: Mapping[str, object], trace_payload: Mapping[str, object] | None = None) -> None:
|
||||
super().__init__("native_traces")
|
||||
|
||||
class Storage:
|
||||
|
|
@ -702,12 +737,31 @@ class _NativeReturningHelp(ModuleType):
|
|||
async def query_help(self, scope: AllQueryScope, secret: str) -> Mapping[str, object]:
|
||||
return help_payload
|
||||
|
||||
get_trace = AsyncMock(return_value=trace_payload)
|
||||
|
||||
self.trace_read: Final = Storage.get_trace
|
||||
self.NativeTraceConfig: Final = _NativeConfig
|
||||
self.NativeTraceStorage: Final = Storage
|
||||
self.trace_encode_error: Final = bytes
|
||||
self.trace_span_rows: Final = list
|
||||
|
||||
|
||||
@pytest.mark.parametrize("cursor,page_size", ((None, None), ("next", 200)))
|
||||
async def test_storage_preserves_page_cursor_and_normalizes_native_trace_data(
|
||||
monkeypatch: pytest.MonkeyPatch, cursor: str | None, page_size: int | None
|
||||
) -> None:
|
||||
native: Final = _NativeReturningHelp(QUERY_HELP, {**TRACE_RESPONSE, "next_cursor": "more"})
|
||||
monkeypatch.setattr(loader, "_cached_bridge", native)
|
||||
storage: Final = ClickHouseStorage(TraceStorageConfig("http://clickhouse:8123"))
|
||||
scope: Final[TraceScope] = {"all_teams": 0, "user_id": "owner", "team_ids": ()}
|
||||
trace: Final = await storage.get_trace("t1", scope, "run", cursor, page_size)
|
||||
assert trace is not None
|
||||
assert trace["next_cursor"] == "more"
|
||||
assert trace["spans"] == ()
|
||||
assert trace["summary"]["span_count"] == 0
|
||||
native.trace_read.assert_awaited_once_with("t1", scope, "run", cursor, page_size)
|
||||
|
||||
|
||||
async def test_storage_validates_the_native_query_help_value(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.setattr(loader, "_cached_bridge", _NativeReturningHelp(QUERY_HELP))
|
||||
storage: Final = ClickHouseStorage(TraceStorageConfig("http://clickhouse:8123"))
|
||||
|
|
|
|||
|
|
@ -1975,10 +1975,15 @@ export const agentTraceListCall = async ({
|
|||
export const sendOtlpTraceCall = async (accessToken: string, exportRequest: object): Promise<void> =>
|
||||
apiClient.post(`/v1/traces`, { accessToken, body: exportRequest });
|
||||
|
||||
export const agentTraceCall = async (accessToken: string, traceId: string, traceRef?: string): Promise<Trace> =>
|
||||
export const agentTraceCall = async (
|
||||
accessToken: string,
|
||||
traceId: string,
|
||||
traceRef?: string,
|
||||
cursor?: string | null,
|
||||
): Promise<Trace> =>
|
||||
apiClient.get<Trace>(`/v1/traces/${encodeURIComponent(traceId)}`, {
|
||||
accessToken,
|
||||
query: { trace_ref: traceRef || undefined },
|
||||
query: { trace_ref: traceRef || undefined, cursor: cursor ?? undefined, page_size: 200 },
|
||||
});
|
||||
|
||||
export const agentTraceSpanCall = async (
|
||||
|
|
|
|||
|
|
@ -85,6 +85,25 @@ describe("AgentTracesSection", () => {
|
|||
vi.mocked(apiClient.get).mockResolvedValue({ data: [] });
|
||||
});
|
||||
|
||||
it("keeps the original time window and loaded rows when another page fails", async () => {
|
||||
const user = userEvent.setup();
|
||||
const now = vi.spyOn(Date, "now").mockReturnValue(Date.parse("2026-10-01T00:00Z"));
|
||||
vi.mocked(agentTraceListCall).mockResolvedValueOnce({ data: runs.slice(0, 1), next_cursor: "next" });
|
||||
renderSection();
|
||||
expect(await screen.findByTestId("agent-trace-row")).toBeVisible();
|
||||
const first = vi.mocked(agentTraceListCall).mock.calls[0][0];
|
||||
now.mockReturnValue(Date.parse("2026-10-01T01:00Z"));
|
||||
vi.mocked(agentTraceListCall).mockRejectedValue(new ApiError("Please try again", 403, {}));
|
||||
await user.click(screen.getByRole("button", { name: "Load more" }));
|
||||
expect(await screen.findByRole("alert")).toHaveTextContent("Could not load more runs");
|
||||
expect(screen.getAllByTestId("agent-trace-row")).toHaveLength(1);
|
||||
expect(vi.mocked(agentTraceListCall).mock.calls[1][0]).toEqual({ ...first, cursor: "next" });
|
||||
vi.mocked(agentTraceListCall).mockResolvedValueOnce({ data: runs.slice(1, 2), next_cursor: null });
|
||||
await user.click(screen.getByRole("button", { name: "Retry" }));
|
||||
await waitFor(() => expect(screen.getAllByTestId("agent-trace-row")).toHaveLength(2));
|
||||
expect(screen.queryByRole("alert")).not.toBeInTheDocument();
|
||||
});
|
||||
|
||||
it.each([
|
||||
[401, "Your session is no longer valid. Sign out and sign in again."],
|
||||
[403, "Your account does not have access to these traces."],
|
||||
|
|
|
|||
|
|
@ -229,6 +229,8 @@ export function AgentTracesSection({
|
|||
isLoading={traces.isLoading || (checkHistory && history.isLoading)}
|
||||
error={traces.error}
|
||||
hasMore={traces.hasMore}
|
||||
isFetching={traces.isFetching}
|
||||
onRetry={traces.hasMore ? traces.loadMore : traces.refetch}
|
||||
onLoadMore={traces.loadMore}
|
||||
onOpenTrace={toggleRun}
|
||||
selectedKey={openTrace === null ? null : runKey(openTrace)}
|
||||
|
|
|
|||
|
|
@ -16,6 +16,8 @@ interface AgentTracesTableProps {
|
|||
isLoading: boolean;
|
||||
error: Error | null;
|
||||
hasMore: boolean;
|
||||
isFetching?: boolean;
|
||||
onRetry?: () => void;
|
||||
onLoadMore: () => void;
|
||||
onOpenTrace: (trace: TraceSummary) => void;
|
||||
selectedKey?: string | null;
|
||||
|
|
@ -53,6 +55,8 @@ export function AgentTracesTable({
|
|||
isLoading,
|
||||
error,
|
||||
hasMore,
|
||||
isFetching = false,
|
||||
onRetry,
|
||||
onLoadMore,
|
||||
onOpenTrace,
|
||||
selectedKey = null,
|
||||
|
|
@ -105,6 +109,14 @@ export function AgentTracesTable({
|
|||
<span className="truncate text-foreground">
|
||||
{firstLine(previewText(run.input_preview)) || traceDisplayName(run)}
|
||||
</span>
|
||||
{run.resolution_limited && (
|
||||
<span
|
||||
className="shrink-0 text-[10px] text-muted-foreground"
|
||||
title="This run is too large to calculate all totals in this view"
|
||||
>
|
||||
Partial totals
|
||||
</span>
|
||||
)}
|
||||
<span className="hidden shrink-0 font-mono text-[10px] text-muted-foreground 2xl:inline">
|
||||
{run.trace_id}
|
||||
</span>
|
||||
|
|
@ -132,15 +144,24 @@ export function AgentTracesTable({
|
|||
</table>
|
||||
{isLoading && <div className="py-16 text-center text-[12px] text-muted-foreground">Loading runs…</div>}
|
||||
{error && (
|
||||
<div className="py-16 text-center text-[12px] text-muted-foreground">Could not load runs: {error.message}</div>
|
||||
<div role="alert" className="flex items-center justify-center gap-3 py-6 text-[12px] text-muted-foreground">
|
||||
<span>
|
||||
{traces.length ? "Could not load more runs" : "Could not load runs"}: {error.message}
|
||||
</span>
|
||||
{onRetry && (
|
||||
<Button size="xs" variant="outline" disabled={isFetching} onClick={onRetry}>
|
||||
Retry
|
||||
</Button>
|
||||
)}
|
||||
</div>
|
||||
)}
|
||||
{isEmpty && (
|
||||
<div className="py-16 text-center text-[12px] text-muted-foreground">No runs match these filters.</div>
|
||||
)}
|
||||
{hasMore && (
|
||||
<div className="border-t border-border/60 px-3 py-2">
|
||||
<Button size="xs" variant="ghost" onClick={onLoadMore}>
|
||||
Load more
|
||||
<Button size="xs" variant="ghost" disabled={isFetching} onClick={onLoadMore}>
|
||||
{isFetching ? "Loading…" : "Load more"}
|
||||
</Button>
|
||||
</div>
|
||||
)}
|
||||
|
|
|
|||
|
|
@ -155,13 +155,46 @@ describe("RunView", () => {
|
|||
expect(screen.queryByRole("button", { name: "Show details" })).not.toBeInTheDocument();
|
||||
});
|
||||
|
||||
it("keeps loaded steps and totals after a page fails, then retries the same cursor", async () => {
|
||||
const user = userEvent.setup();
|
||||
const summary = { ...research.summary, span_count: 2 };
|
||||
const first: Trace = { ...research, summary, spans: research.spans.slice(0, 1), next_cursor: "next-page" };
|
||||
const second: Trace = {
|
||||
...research,
|
||||
summary,
|
||||
spans: [
|
||||
{ ...research.spans[1], type: "tool", name: "later-page-tool", parent_span_id: research.spans[0].span_id },
|
||||
],
|
||||
next_cursor: null,
|
||||
};
|
||||
vi.mocked(agentTraceCall).mockReset();
|
||||
vi.mocked(agentTraceCall)
|
||||
.mockResolvedValueOnce(first)
|
||||
.mockRejectedValueOnce(new Error("offline"))
|
||||
.mockResolvedValueOnce(second);
|
||||
renderWithProviders(<RunView traceId={research.summary.trace_id} accessToken="sk-test" onBack={vi.fn()} />);
|
||||
|
||||
expect(await screen.findByText("Showing 1 of 2 steps")).toBeVisible();
|
||||
const before = screen.getByRole("banner").textContent;
|
||||
await user.click(screen.getByRole("button", { name: "Load more steps" }));
|
||||
expect(await screen.findByText("Could not load more steps. Your loaded steps are still available.")).toBeVisible();
|
||||
expect(screen.getByRole("tree", { name: "Spans in time order" })).toHaveTextContent(research.spans[0].name);
|
||||
await user.click(screen.getByRole("button", { name: "Retry" }));
|
||||
expect(await screen.findByText("later-page-tool")).toBeVisible();
|
||||
expect(screen.getByRole("tree", { name: "Spans in time order" })).toHaveTextContent(research.spans[0].name);
|
||||
expect(screen.getAllByRole("treeitem")).toHaveLength(2);
|
||||
expect(screen.getByRole("banner")).toHaveTextContent(before ?? "");
|
||||
expect(screen.queryByRole("button", { name: "Load more steps" })).not.toBeInTheDocument();
|
||||
expect(vi.mocked(agentTraceCall).mock.calls.map((call) => call[3])).toEqual([null, "next-page", "next-page"]);
|
||||
});
|
||||
|
||||
it("keeps a way back to the runs table when a run fails to load", async () => {
|
||||
const user = userEvent.setup();
|
||||
const onBack = vi.fn();
|
||||
vi.mocked(agentTraceCall).mockRejectedValue(new Error("trace exceeds the 1000 span read limit"));
|
||||
vi.mocked(agentTraceCall).mockRejectedValue(new Error("Traces are temporarily unavailable"));
|
||||
renderWithProviders(<RunView traceId="big" accessToken="sk-test" onBack={onBack} />);
|
||||
|
||||
expect(await screen.findByText("trace exceeds the 1000 span read limit")).toBeInTheDocument();
|
||||
expect(await screen.findByText("Traces are temporarily unavailable")).toBeInTheDocument();
|
||||
await user.click(screen.getByRole("button", { name: /back to traces/i }));
|
||||
expect(onBack).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@
|
|||
import { useLensDemo } from "@/components/lens/LensDemoContext";
|
||||
import { useTracesApi } from "@/components/lens/services";
|
||||
|
||||
import { useQuery } from "@tanstack/react-query";
|
||||
import { useInfiniteQuery } from "@tanstack/react-query";
|
||||
import { ArrowLeft, Check, Copy } from "lucide-react";
|
||||
import { useCallback, useEffect, useMemo, useState } from "react";
|
||||
|
||||
|
|
@ -368,16 +368,34 @@ interface RunViewProps {
|
|||
embedded?: boolean;
|
||||
}
|
||||
|
||||
/** One agent run: header with totals and "Copy for agent", span tree on the left, span details on the right. */
|
||||
function initialSpanMissing(trace: Trace | undefined, spanId?: string): boolean {
|
||||
return Boolean(spanId && trace && !trace.spans.some((span) => span.span_id === spanId));
|
||||
}
|
||||
|
||||
export function RunView({ traceId, traceRef, initialSpanId, accessToken, onBack, embedded = false }: RunViewProps) {
|
||||
const traces = useTracesApi(accessToken);
|
||||
const [view, setView] = useState<TraceView>("steps");
|
||||
const traceQuery = useQuery({
|
||||
const traceQueryOptions = {
|
||||
queryKey: ["agentTrace", traceId, traceRef, accessToken],
|
||||
queryFn: () => traces.trace(traceId, traceRef),
|
||||
queryFn: ({ pageParam }: { pageParam: string | null }) => traces.trace(traceId, traceRef, pageParam),
|
||||
initialPageParam: null as string | null,
|
||||
getNextPageParam: (lastPage: Trace) => lastPage.next_cursor ?? undefined,
|
||||
staleTime: 30_000,
|
||||
});
|
||||
const trace = traceQuery.data;
|
||||
retry: false,
|
||||
};
|
||||
const traceQuery = useInfiniteQuery(traceQueryOptions);
|
||||
const trace = useMemo(() => {
|
||||
const pages = traceQuery.data?.pages;
|
||||
if (!pages?.length) return undefined;
|
||||
return { ...pages[0], spans: pages.flatMap((page) => page.spans) };
|
||||
}, [traceQuery.data]);
|
||||
const seekingSpan = initialSpanMissing(trace, initialSpanId);
|
||||
const { hasNextPage, isFetching, isError, fetchNextPage } = traceQuery;
|
||||
const canSeek = seekingSpan && hasNextPage;
|
||||
useEffect(() => {
|
||||
if (canSeek && !isFetching && !isError) void fetchNextPage();
|
||||
}, [canSeek, isFetching, isError, fetchNextPage]);
|
||||
const pageAction = isError ? "Retry" : "Load more steps";
|
||||
|
||||
if (traceQuery.isLoading) {
|
||||
return (
|
||||
|
|
@ -396,7 +414,7 @@ export function RunView({ traceId, traceRef, initialSpanId, accessToken, onBack,
|
|||
</div>
|
||||
);
|
||||
}
|
||||
if (traceQuery.isError || !trace) {
|
||||
if (!trace) {
|
||||
return (
|
||||
<div className="p-6 text-[12px]" data-testid="run-view-error">
|
||||
<button
|
||||
|
|
@ -408,6 +426,9 @@ export function RunView({ traceId, traceRef, initialSpanId, accessToken, onBack,
|
|||
</button>
|
||||
<h1 className="mb-2 text-[13px] font-medium">Could not load trace</h1>
|
||||
<span className="text-muted-foreground">{traceQuery.error?.message ?? "Unknown error"}</span>
|
||||
<Button variant="outline" size="sm" className="ml-3" onClick={() => void traceQuery.refetch()}>
|
||||
Retry
|
||||
</Button>
|
||||
</div>
|
||||
);
|
||||
}
|
||||
|
|
@ -422,8 +443,35 @@ export function RunView({ traceId, traceRef, initialSpanId, accessToken, onBack,
|
|||
data-testid="run-view"
|
||||
>
|
||||
<RunHeader trace={trace} onBack={onBack} embedded={embedded} />
|
||||
{(traceQuery.hasNextPage || traceQuery.isError) && (
|
||||
<div className="flex items-center justify-between gap-3 border-b px-3 py-2 text-xs" role="status">
|
||||
<span>
|
||||
{traceQuery.isError
|
||||
? "Could not load more steps. Your loaded steps are still available."
|
||||
: `Showing ${trace.spans.length.toLocaleString()} of ${trace.summary.span_count.toLocaleString()} steps`}
|
||||
</span>
|
||||
{traceQuery.isError && (
|
||||
<Button
|
||||
size="xs"
|
||||
variant="ghost"
|
||||
disabled={traceQuery.isFetching}
|
||||
onClick={() => void traceQuery.refetch()}
|
||||
>
|
||||
Refresh trace
|
||||
</Button>
|
||||
)}
|
||||
<Button
|
||||
size="xs"
|
||||
variant="outline"
|
||||
disabled={traceQuery.isFetching}
|
||||
onClick={() => void (traceQuery.hasNextPage ? traceQuery.fetchNextPage() : traceQuery.refetch())}
|
||||
>
|
||||
{traceQuery.isFetching ? "Loading…" : pageAction}
|
||||
</Button>
|
||||
</div>
|
||||
)}
|
||||
<RunBody
|
||||
key={trace.summary.trace_id}
|
||||
key={`${trace.summary.trace_ref || trace.summary.trace_id}:${initialSpanId && !seekingSpan ? initialSpanId : "root"}`}
|
||||
trace={trace}
|
||||
accessToken={accessToken}
|
||||
initialSpanId={initialSpanId}
|
||||
|
|
|
|||
|
|
@ -16,7 +16,7 @@ export interface TraceWindow {
|
|||
export interface TracesApi {
|
||||
list(window: TraceWindow): Promise<TracePage>;
|
||||
anyRecorded(): Promise<boolean>;
|
||||
trace(traceId: string, traceRef?: string): Promise<Trace>;
|
||||
trace(traceId: string, traceRef?: string, cursor?: string | null): Promise<Trace>;
|
||||
span(traceId: string, spanId: string, traceRef?: string): Promise<SpanDetail>;
|
||||
spanError(
|
||||
traceId: string,
|
||||
|
|
@ -32,7 +32,7 @@ export function liveTracesApi(accessToken: string): TracesApi {
|
|||
const page = await apiClient.get<TracePage>("/v1/traces", { accessToken, query: { start_ms: 0 } });
|
||||
return page.data.length > 0;
|
||||
},
|
||||
trace: (traceId, traceRef) => agentTraceCall(accessToken, traceId, traceRef),
|
||||
trace: (traceId, traceRef, cursor) => agentTraceCall(accessToken, traceId, traceRef, cursor),
|
||||
span: (traceId, spanId, traceRef) => agentTraceSpanCall(accessToken, traceId, spanId, traceRef),
|
||||
spanError: (traceId, spanId, options) => agentTraceSpanErrorCall(accessToken, traceId, spanId, options),
|
||||
};
|
||||
|
|
|
|||
|
|
@ -7,6 +7,11 @@ import { ApiError } from "@/lib/http/client";
|
|||
|
||||
import { LIVE_TAIL_INTERVAL_MS } from "../log_filter_logic";
|
||||
import type { TracePage, TraceSummary } from "./traceTypes";
|
||||
import type { TraceWindow } from "./tracesApi";
|
||||
|
||||
interface LoadedTracePage extends TracePage {
|
||||
window: TraceWindow;
|
||||
}
|
||||
|
||||
export const TRACING_NOT_ENABLED_STATUS = 501;
|
||||
/** A proxy without the tracing routes at all answers 404; treat it like tracing being off. */
|
||||
|
|
@ -55,7 +60,7 @@ export const traceWindowStartMs = (startTime: string, endTime: string, isCustomD
|
|||
|
||||
/**
|
||||
* GET /v1/traces for the Logs page time range, cursor-paginated ("Load more").
|
||||
* Preset ranges re-read "now" on every fetch, moving both bounds so the window keeps its length.
|
||||
* Preset ranges roll on refresh; subsequent pages keep the first page's window.
|
||||
*/
|
||||
export function useAgentTraces({
|
||||
accessToken,
|
||||
|
|
@ -66,19 +71,20 @@ export function useAgentTraces({
|
|||
enabled,
|
||||
}: UseAgentTracesOptions): AgentTracesResult {
|
||||
const traces = useTracesApi(accessToken);
|
||||
const fetchPage = (pageParam: unknown): Promise<TracePage> => {
|
||||
const fetchPage = async (pageParam: unknown): Promise<LoadedTracePage> => {
|
||||
const nowMs = Date.now();
|
||||
return traces.list({
|
||||
const window = (pageParam as TraceWindow | null) ?? {
|
||||
startMs: traceWindowStartMs(startTime, endTime, isCustomDate, nowMs),
|
||||
endMs: isCustomDate ? moment(endTime).valueOf() : nowMs,
|
||||
cursor: pageParam as string | null,
|
||||
});
|
||||
};
|
||||
return { ...(await traces.list(window)), window };
|
||||
};
|
||||
const queryOptions: Parameters<typeof useInfiniteQuery<TracePage, Error>>[0] = {
|
||||
const queryOptions: Parameters<typeof useInfiniteQuery<LoadedTracePage, Error>>[0] = {
|
||||
queryKey: ["agentTraces", accessToken, startTime, endTime, isCustomDate],
|
||||
queryFn: ({ pageParam }) => fetchPage(pageParam),
|
||||
initialPageParam: null,
|
||||
getNextPageParam: (lastPage) => lastPage.next_cursor ?? undefined,
|
||||
getNextPageParam: (lastPage) =>
|
||||
lastPage.next_cursor ? { ...lastPage.window, cursor: lastPage.next_cursor } : undefined,
|
||||
enabled,
|
||||
staleTime: LIVE_TAIL_INTERVAL_MS,
|
||||
retry: (failureCount, error) => !requiresUserAction(error) && failureCount < 1,
|
||||
|
|
@ -87,7 +93,7 @@ export function useAgentTraces({
|
|||
refetchOnReconnect: (q) => !requiresUserAction(q.state.error),
|
||||
refetchIntervalInBackground: false,
|
||||
};
|
||||
const query = useInfiniteQuery<TracePage, Error>(queryOptions);
|
||||
const query = useInfiniteQuery<LoadedTracePage, Error>(queryOptions);
|
||||
|
||||
const loaded = useMemo(() => query.data?.pages.flatMap((page) => page.data) ?? [], [query.data]);
|
||||
const notEnabled = isTracingNotEnabled(query.error);
|
||||
|
|
@ -99,7 +105,9 @@ export function useAgentTraces({
|
|||
notEnabledDetail: notEnabled ? query.error?.message || "Agent tracing is not enabled" : null,
|
||||
error: notEnabled ? null : displayError(query.error),
|
||||
hasMore: query.hasNextPage,
|
||||
loadMore: () => void query.fetchNextPage(),
|
||||
loadMore: () => {
|
||||
if (!query.isFetching) void query.fetchNextPage();
|
||||
},
|
||||
refetch: () => void query.refetch(),
|
||||
};
|
||||
}
|
||||
|
|
|
|||
6
ui/litellm-dashboard/src/lib/http/schema.d.ts
generated
vendored
6
ui/litellm-dashboard/src/lib/http/schema.d.ts
generated
vendored
|
|
@ -47047,6 +47047,8 @@ export interface components {
|
|||
Trace: {
|
||||
/** Agents */
|
||||
agents: components["schemas"]["AgentNode"][];
|
||||
/** Next Cursor */
|
||||
next_cursor?: string | null;
|
||||
/** Spans */
|
||||
spans: components["schemas"]["Span"][];
|
||||
summary: components["schemas"]["TraceSummary"];
|
||||
|
|
@ -47279,6 +47281,8 @@ export interface components {
|
|||
name: string;
|
||||
/** Output Tokens */
|
||||
output_tokens: number;
|
||||
/** Resolution Limited */
|
||||
resolution_limited?: boolean;
|
||||
/** Service */
|
||||
service: string;
|
||||
/** Span Count */
|
||||
|
|
@ -80234,6 +80238,8 @@ export interface operations {
|
|||
parameters: {
|
||||
query?: {
|
||||
trace_ref?: string;
|
||||
cursor?: string | null;
|
||||
page_size?: number | null;
|
||||
};
|
||||
header?: never;
|
||||
path: {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue