diff --git a/litellm-rust/Cargo.lock b/litellm-rust/Cargo.lock index daa0cf631bd..f7b667c8ab2 100644 --- a/litellm-rust/Cargo.lock +++ b/litellm-rust/Cargo.lock @@ -3932,22 +3932,6 @@ dependencies = [ "uuid", ] -[[package]] -name = "litellm-gateway-traces" -version = "0.1.0" -dependencies = [ - "axum", - "litellm-traces", - "litellm-traces-cache", - "rstest", - "serde", - "serde_json", - "thiserror 2.0.19", - "tokio", - "tower", - "tracing", -] - [[package]] name = "litellm-gateway-ui" version = "0.1.0" diff --git a/litellm-rust/crates/gateway-traces/Cargo.toml b/litellm-rust/crates/gateway-traces/Cargo.toml deleted file mode 100644 index a8abedd3675..00000000000 --- a/litellm-rust/crates/gateway-traces/Cargo.toml +++ /dev/null @@ -1,20 +0,0 @@ -[package] -name = "litellm-gateway-traces" -version = "0.1.0" -edition.workspace = true -license.workspace = true -repository.workspace = true - -[dependencies] -axum = { workspace = true, features = ["json", "query"] } -litellm-traces.workspace = true -litellm-traces-cache.workspace = true -serde.workspace = true -tracing.workspace = true - -[dev-dependencies] -rstest.workspace = true -serde_json.workspace = true -thiserror.workspace = true -tokio.workspace = true -tower = { version = "0.5", features = ["util"] } diff --git a/litellm-rust/crates/gateway-traces/src/error.rs b/litellm-rust/crates/gateway-traces/src/error.rs deleted file mode 100644 index b557c120652..00000000000 --- a/litellm-rust/crates/gateway-traces/src/error.rs +++ /dev/null @@ -1,87 +0,0 @@ -use axum::{ - Json, - http::{HeaderValue, StatusCode, header}, - response::{IntoResponse, Response}, -}; -use litellm_traces_cache::ReadError; -use serde::Serialize; - -const RETRY_AFTER_SECONDS: &str = "2"; - -#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)] -#[serde(rename_all = "snake_case")] -pub enum ReadFailureCode { - InvalidRequest, - TraceChanged, - TooLarge, - Unavailable, -} - -#[derive(Debug)] -pub struct ReadFailure { - code: ReadFailureCode, - message: String, -} - -impl From> for ReadFailure { - fn from(error: ReadError) -> Self { - let (code, message) = match &error { - ReadError::InvalidParameters - | ReadError::InvalidCursor(_) - | ReadError::AmbiguousTrace => (ReadFailureCode::InvalidRequest, error.to_string()), - ReadError::TraceChanged => (ReadFailureCode::TraceChanged, error.to_string()), - ReadError::TooLarge => ( - ReadFailureCode::TooLarge, - "Trace is too large for this view. Use a filtered trace query.".to_owned(), - ), - ReadError::Encode(_) | ReadError::Store(_) => { - tracing::warn!(%error, "trace read unavailable"); - ( - ReadFailureCode::Unavailable, - "Traces are temporarily unavailable. Please try again.".to_owned(), - ) - } - }; - Self { code, message } - } -} - -#[derive(Serialize)] -struct Body<'a> { - detail: Detail<'a>, -} - -#[derive(Serialize)] -struct Detail<'a> { - code: ReadFailureCode, - message: &'a str, -} - -impl IntoResponse for ReadFailure { - fn into_response(self) -> Response { - let status = match self.code { - ReadFailureCode::InvalidRequest => StatusCode::BAD_REQUEST, - ReadFailureCode::TraceChanged => StatusCode::CONFLICT, - ReadFailureCode::TooLarge => StatusCode::PAYLOAD_TOO_LARGE, - ReadFailureCode::Unavailable => StatusCode::SERVICE_UNAVAILABLE, - }; - let body = Json(Body { - detail: Detail { - code: self.code, - message: &self.message, - }, - }); - match self.code { - ReadFailureCode::Unavailable => ( - status, - [( - header::RETRY_AFTER, - HeaderValue::from_static(RETRY_AFTER_SECONDS), - )], - body, - ) - .into_response(), - _ => (status, body).into_response(), - } - } -} diff --git a/litellm-rust/crates/gateway-traces/src/lib.rs b/litellm-rust/crates/gateway-traces/src/lib.rs deleted file mode 100644 index 2b31f22c0a5..00000000000 --- a/litellm-rust/crates/gateway-traces/src/lib.rs +++ /dev/null @@ -1,24 +0,0 @@ -mod error; -mod runs; - -use std::sync::Arc; - -use axum::{Router, routing::get}; -pub use error::ReadFailure; -use litellm_traces_cache::{TraceReader, TraceStore}; - -pub struct Traces { - pub reader: TraceReader, - pub store: S, -} - -pub fn router(traces: Arc>) -> Router -where - S: TraceStore + Send + 'static, -{ - Router::new() - .route("/v1/traces", get(runs::list::)) - .route("/v1/traces/histogram", get(runs::histogram::)) - .route("/v1/traces/values/{field}", get(runs::values::)) - .with_state(traces) -} diff --git a/litellm-rust/crates/gateway-traces/src/runs.rs b/litellm-rust/crates/gateway-traces/src/runs.rs deleted file mode 100644 index 030a7c2a5a5..00000000000 --- a/litellm-rust/crates/gateway-traces/src/runs.rs +++ /dev/null @@ -1,135 +0,0 @@ -use std::{ - sync::Arc, - time::{SystemTime, UNIX_EPOCH}, -}; - -use axum::{ - Extension, Json, - extract::{Path, Query, State}, -}; -use litellm_traces::{ - QueryScope, TracePage, - search::{RunField, RunFilter, RunSearch, RunValues, TraceHistogram}, - store::RunOrder, -}; -use litellm_traces_cache::{PageRequest, TraceStore}; -use serde::Deserialize; - -use crate::{ReadFailure, Traces}; - -const DAY_MS: i64 = 24 * 60 * 60 * 1000; -const PAGE_SIZE: u32 = 50; - -#[derive(Deserialize)] -pub(crate) struct Runs { - start_ms: Option, - end_ms: Option, - #[serde(default)] - q: String, -} - -impl Runs { - fn filter(&self) -> RunFilter { - let now_ms = now_ms(); - RunFilter { - start_ms: self.start_ms.unwrap_or(now_ms - DAY_MS), - end_ms: self.end_ms.unwrap_or(now_ms), - search: RunSearch::parse(&self.q), - trace_refs: Vec::new(), - } - } -} - -#[derive(Deserialize)] -pub(crate) struct Page { - cursor: Option, -} - -pub(crate) async fn list( - State(traces): State>>, - Extension(access): Extension, - Query(runs): Query, - Query(Page { cursor }): Query, -) -> Result, ReadFailure> { - let page = PageRequest { - cursor, - limit: PAGE_SIZE, - ..PageRequest::default() - }; - Ok(Json( - traces - .reader - .list_traces( - &traces.store, - &access, - &runs.filter(), - RunOrder::NEWEST, - &page, - ) - .await?, - )) -} - -#[derive(Deserialize)] -pub(crate) struct Buckets { - #[serde(default = "default_buckets")] - buckets: u32, -} - -fn default_buckets() -> u32 { - 60 -} - -pub(crate) async fn histogram( - State(traces): State>>, - Extension(access): Extension, - Query(runs): Query, - Query(Buckets { buckets }): Query, -) -> Result, ReadFailure> { - Ok(Json( - traces - .reader - .histogram(&traces.store, &access, &runs.filter(), buckets) - .await?, - )) -} - -#[derive(Deserialize)] -pub(crate) struct Values { - #[serde(default)] - contains: String, - #[serde(default = "default_values")] - limit: u32, -} - -fn default_values() -> u32 { - 20 -} - -pub(crate) async fn values( - State(traces): State>>, - Extension(access): Extension, - Path(field): Path, - Query(runs): Query, - Query(values): Query, -) -> Result, ReadFailure> { - Ok(Json( - traces - .reader - .values( - &traces.store, - &access, - &runs.filter(), - field, - &values.contains, - values.limit, - ) - .await?, - )) -} - -fn now_ms() -> i64 { - SystemTime::now() - .duration_since(UNIX_EPOCH) - .map_or(0, |elapsed| elapsed.as_millis() as i64) -} diff --git a/litellm-rust/crates/gateway-traces/tests/routes.rs b/litellm-rust/crates/gateway-traces/tests/routes.rs deleted file mode 100644 index de4b00449d3..00000000000 --- a/litellm-rust/crates/gateway-traces/tests/routes.rs +++ /dev/null @@ -1,283 +0,0 @@ -use std::{ - sync::{Arc, Mutex}, - time::{SystemTime, UNIX_EPOCH}, -}; - -use axum::{ - Extension, Router, - body::{Body, to_bytes}, - http::{Request, StatusCode, header}, - response::Response, -}; -use litellm_gateway_traces::{Traces, router}; -use litellm_traces::{ - QueryScope, - search::{RunField, RunFilter, RunSearch}, - store::{ - CallQuery, CallRow, CountValue, RunCount, RunCountQuery, RunQuery, RunRow, RunSelection, - SpanQuery, SpanRow, SpanText, SpanTextQuery, - }, -}; -use litellm_traces_cache::{StoreError, StoreResult, TraceReader, TraceStore}; -use rstest::rstest; -use serde_json::{Value, json}; -use tower::ServiceExt; - -#[derive(Debug, thiserror::Error)] -#[error("fake trace store failed")] -struct FakeError; - -#[derive(Clone, Copy)] -enum Outcome { - Empty, - TooLarge, - Failed, -} - -struct FakeStore { - outcome: Outcome, - lists: Mutex>, - counts: Mutex>, -} - -impl FakeStore { - fn listed_filter(&self) -> RunFilter { - let lists = self.lists.lock().unwrap(); - let [(_, query)] = lists.as_slice() else { - panic!("expected one list read, got {}", lists.len()); - }; - let RunSelection::Matching(filter) = &query.selection else { - panic!("expected a run search"); - }; - filter.clone() - } - - fn counted(&self) -> RunCountQuery { - self.counts.lock().unwrap()[0].clone() - } -} - -impl TraceStore for FakeStore { - type Error = FakeError; - - fn source(&self) -> &str { - "fake" - } - - async fn runs( - &self, - access: &QueryScope, - query: &RunQuery, - ) -> StoreResult, FakeError> { - self.lists - .lock() - .unwrap() - .push((access.clone(), query.clone())); - match self.outcome { - Outcome::Empty => Ok(Vec::new()), - Outcome::TooLarge => Err(StoreError::TooLarge), - Outcome::Failed => Err(StoreError::Failed(FakeError)), - } - } - - async fn run_counts( - &self, - _: &QueryScope, - query: &RunCountQuery, - ) -> StoreResult, FakeError> { - self.counts.lock().unwrap().push(query.clone()); - Ok(vec![RunCount { - bucket: 0, - failed: false, - value: "researcher".into(), - runs: 2, - }]) - } - - async fn spans(&self, _: &QueryScope, _: &SpanQuery) -> StoreResult, FakeError> { - Ok(Vec::new()) - } - - async fn span_text( - &self, - _: &QueryScope, - _: &SpanTextQuery, - ) -> StoreResult, FakeError> { - Ok(Vec::new()) - } - - async fn calls(&self, _: &QueryScope, _: &CallQuery) -> StoreResult, FakeError> { - Ok(Vec::new()) - } -} - -fn access() -> QueryScope { - QueryScope::Owned { - user_id: "user".into(), - team_ids: vec!["team".into()], - } -} - -fn app(outcome: Outcome) -> (Router, Arc>) { - let traces = Arc::new(Traces { - reader: TraceReader::new(usize::MAX), - store: FakeStore { - outcome, - lists: Mutex::new(Vec::new()), - counts: Mutex::new(Vec::new()), - }, - }); - let app = router(Arc::clone(&traces)).layer(Extension(access())); - (app, traces) -} - -async fn get(app: Router, uri: &str) -> Response { - app.oneshot(Request::get(uri).body(Body::empty()).unwrap()) - .await - .unwrap() -} - -async fn json_body(response: Response) -> Value { - serde_json::from_slice(&to_bytes(response.into_body(), 65536).await.unwrap()).unwrap() -} - -#[tokio::test] -async fn list_reads_the_window_and_parsed_search_in_the_callers_scope() { - let (app, traces) = app(Outcome::Empty); - let response = get( - app, - "/v1/traces?start_ms=10&end_ms=20&q=plan%20-agent:res*%20model:%22gpt%20x%22", - ) - .await; - - assert_eq!(response.status(), StatusCode::OK); - assert_eq!( - json_body(response).await, - json!({"data": [], "next_cursor": null}) - ); - let filter = traces.store.listed_filter(); - assert_eq!((filter.start_ms, filter.end_ms), (10, 20)); - let lists = traces.store.lists.lock().unwrap(); - assert_eq!((&lists[0].0, lists[0].1.limit), (&access(), 50)); - assert_eq!( - filter.search, - RunSearch::parse(r#"plan -agent:res* model:"gpt x""#) - ); -} - -#[tokio::test] -async fn list_defaults_to_the_last_day() { - let (app, traces) = app(Outcome::Empty); - let before = now_ms(); - assert_eq!(get(app, "/v1/traces").await.status(), StatusCode::OK); - let after = now_ms(); - - let filter = traces.store.listed_filter(); - assert!((before..=after).contains(&filter.end_ms)); - assert_eq!(filter.end_ms - filter.start_ms, 24 * 60 * 60 * 1000); - assert_eq!(filter.search, RunSearch::default()); -} - -#[tokio::test] -async fn histogram_buckets_the_matching_runs_of_the_window() { - let (app, traces) = app(Outcome::Empty); - let response = get( - app, - "/v1/traces/histogram?start_ms=0&end_ms=40&q=status:ok&buckets=4", - ) - .await; - - assert_eq!(response.status(), StatusCode::OK); - let body = json_body(response).await; - assert_eq!(body["buckets"].as_array().map(Vec::len), Some(4)); - assert_eq!( - body["buckets"][0], - json!({"start_ms": 0, "end_ms": 10, "total": 2, "failed": 0, "agents": [{"agent": "researcher", "runs": 2}]}) - ); - let query = traces.store.counted(); - assert_eq!(query.by.buckets, Some(4)); - assert_eq!(query.filter.search, RunSearch::parse("status:ok")); -} - -#[tokio::test] -async fn values_suggest_a_field_narrowed_by_the_search() { - let (app, traces) = app(Outcome::Empty); - let response = get( - app, - "/v1/traces/values/agent?start_ms=0&end_ms=40&q=model:gpt*&contains=res&limit=5", - ) - .await; - - assert_eq!(response.status(), StatusCode::OK); - assert_eq!(json_body(response).await, json!({"values": ["researcher"]})); - let query = traces.store.counted(); - assert_eq!(query.by.value, Some(CountValue::Field(RunField::Agent))); - assert_eq!((query.contains.as_str(), query.limit), ("res", Some(5))); - assert_eq!(query.filter.search, RunSearch::parse("model:gpt*")); -} - -#[rstest] -#[case::bad_cursor( - Outcome::Empty, - "/v1/traces?cursor=not-a-cursor", - StatusCode::BAD_REQUEST, - "invalid_request" -)] -#[case::bad_window(Outcome::Empty, "/v1/traces?start_ms=x", StatusCode::BAD_REQUEST, "")] -#[case::reversed_window( - Outcome::Empty, - "/v1/traces?start_ms=20&end_ms=10", - StatusCode::BAD_REQUEST, - "invalid_request" -)] -#[case::too_many_buckets( - Outcome::Empty, - "/v1/traces/histogram?buckets=241", - StatusCode::BAD_REQUEST, - "invalid_request" -)] -#[case::too_many_values( - Outcome::Empty, - "/v1/traces/values/agent?limit=101", - StatusCode::BAD_REQUEST, - "invalid_request" -)] -#[case::unknown_field(Outcome::Empty, "/v1/traces/values/color", StatusCode::BAD_REQUEST, "")] -#[case::too_large( - Outcome::TooLarge, - "/v1/traces", - StatusCode::PAYLOAD_TOO_LARGE, - "too_large" -)] -#[case::store_down( - Outcome::Failed, - "/v1/traces", - StatusCode::SERVICE_UNAVAILABLE, - "unavailable" -)] -#[tokio::test] -async fn failures_keep_the_python_status_and_code( - #[case] outcome: Outcome, - #[case] uri: &str, - #[case] status: StatusCode, - #[case] code: &str, -) { - let (app, _) = app(outcome); - let response = get(app, uri).await; - - assert_eq!(response.status(), status); - assert_eq!( - response.headers().get(header::RETRY_AFTER).is_some(), - status == StatusCode::SERVICE_UNAVAILABLE - ); - if !code.is_empty() { - assert_eq!(json_body(response).await["detail"]["code"], code); - } -} - -fn now_ms() -> i64 { - SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap() - .as_millis() as i64 -} diff --git a/litellm-rust/crates/python-bridge/src/routes/traces.rs b/litellm-rust/crates/python-bridge/src/routes/traces.rs index 1892bda37bc..7816fd4ab59 100644 --- a/litellm-rust/crates/python-bridge/src/routes/traces.rs +++ b/litellm-rust/crates/python-bridge/src/routes/traces.rs @@ -1,4 +1,8 @@ -use std::{collections::BTreeMap, sync::Arc}; +use std::{ + collections::BTreeMap, + sync::Arc, + time::{SystemTime, UNIX_EPOCH}, +}; use litellm_http::ClientVariant; use litellm_traces::{ @@ -6,7 +10,7 @@ use litellm_traces::{ search::{RunField, RunFilter, RunSearch}, store::{RunOrder, SpanPart, TextRange}, }; -use litellm_traces_cache::{PageRequest, ReadError, TraceReader}; +use litellm_traces_cache::{PageRequest, ReadError, TraceReader, resolve_run_window}; use litellm_traces_clickhouse::{ClickHouseTraces, Config, Error, InsertTable, QueryReaders}; use prost::Message; use pyo3::{ @@ -16,6 +20,7 @@ use pyo3::{ }; pyo3::import_exception!(litellm.rust_bridge.trace.errors, TraceChanged); +pyo3::import_exception!(litellm.rust_bridge.trace.errors, TraceQueryError); #[derive(Message)] struct OtlpErrorStatus { @@ -109,11 +114,44 @@ fn parsed(kind: &str, value: &str) -> PyResult { } fn map_sql_error(error: Error) -> PyErr { + map_sql_error_ref(&error) +} + +fn map_sql_error_ref(error: &Error) -> PyErr { match error { - Error::Storage(litellm_storage_clickhouse::Error::QueryFailed(400 | 404)) => { - PyValueError::new_err(error.to_string()) + Error::Storage(litellm_storage_clickhouse::Error::QueryFailed(failure)) => { + TraceQueryError::new_err(( + sql_failure_kind(failure), + failure.code, + failure.message.clone(), + )) } - error => map_error(error), + Error::Storage(litellm_storage_clickhouse::Error::ResponseTooLarge) => { + TraceQueryError::new_err(( + "limited", + None::, + "Query exceeded the response size limit", + )) + } + Error::Cached(source) => map_sql_error_ref(source), + error => map_error_ref(error), + } +} + +fn sql_failure_kind(failure: &litellm_storage_clickhouse::QueryFailure) -> &'static str { + // https://github.com/ClickHouse/ClickHouse/blob/v26.9.6.6-stable/src/Common/ErrorCodes.cpp + match failure.code { + Some(158 | 159 | 160 | 167 | 168 | 191 | 202 | 229 | 241 | 290 | 307 | 396 | 776) => { + "limited" + } + Some( + 6 | 34 | 35 | 36 | 42 | 43 | 44 | 46 | 47 | 48 | 50 | 53 | 60 | 62 | 63 | 69 | 70 | 72 + | 73 | 78 | 80 | 81 | 115 | 164 | 291 | 344 | 392 | 452 | 472 | 497, + ) => "rejected", + Some(192 | 193 | 194 | 516) => "unavailable", + _ if matches!(failure.status, 408 | 413 | 429) => "limited", + _ if matches!(failure.status, 400 | 404 | 405 | 406 | 411 | 415 | 422) => "rejected", + _ => "unavailable", } } @@ -249,15 +287,26 @@ impl NativeTraceStorage { &self, py: Python<'py>, #[pyo3(from_py_with = litellm_host_python::from_py_argument)] scope: QueryScope, - start_ms: i64, - end_ms: i64, + start_ms: Option, + end_ms: Option, q: &str, cursor: Option, limit: u32, #[pyo3(from_py_with = litellm_host_python::from_py_argument)] order: RunOrder, trace_refs: Vec, ) -> PyResult> { - let filter = run_filter(start_ms, end_ms, q, trace_refs); + let now_ms = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_millis() as i64; + let window = resolve_run_window( + start_ms, + end_ms, + cursor.as_deref(), + (now_ms - 86_400_000, now_ms), + ) + .map_err(map_read_error)?; + let filter = run_filter(window.0, window.1, q, trace_refs); let page = PageRequest { cursor, limit }; let client = crate::http::host_client(py, ClientVariant::NoRedirect)?; let connection = self.config.storage().reader().clone(); @@ -533,7 +582,7 @@ impl NativeTraceStorage { let connection = readers.connection(&client, &scope, &secret).await?; litellm_traces_clickhouse::query_help(&client, &connection).await }, - map_sql_error, + map_error, ) } } @@ -563,6 +612,21 @@ mod tests { use super::*; + fn initialize_python_path(py: Python<'_>) { + let repository = std::path::Path::new(env!("CARGO_MANIFEST_DIR")) + .ancestors() + .nth(3) + .unwrap() + .to_str() + .unwrap(); + pyo3::types::PyModule::import(py, "sys") + .unwrap() + .getattr("path") + .unwrap() + .call_method1("insert", (0, repository)) + .unwrap(); + } + #[rstest] #[case::row(Error::InvalidRow, "ValueError")] #[case::insert_limit(Error::InvalidLimit("CLICKHOUSE_TRACE_MAX_INSERT_BYTES"), "ValueError")] @@ -598,26 +662,71 @@ mod tests { } #[rstest] - #[case::invalid_sql(400, "ValueError")] - #[case::missing_table(404, "ValueError")] - #[case::unavailable(503, "RuntimeError")] - fn wrapped_query_status_preserves_public_exception_type( + #[case::syntax(500, Some(62), "rejected")] + #[case::unknown_column(500, Some(47), "rejected")] + #[case::readonly(500, Some(164), "rejected")] + #[case::denied_table(403, Some(497), "rejected")] + #[case::settings_constraint(500, Some(452), "rejected")] + #[case::memory_limit(500, Some(241), "limited")] + #[case::timeout(408, Some(159), "limited")] + #[case::result_limit(500, Some(396), "limited")] + #[case::invalid_sql_status(400, None, "rejected")] + #[case::unavailable(503, None, "unavailable")] + #[case::invalid_reader_credentials(403, Some(516), "unavailable")] + fn raw_query_failures_preserve_diagnostics_and_curated_exception_type( #[case] status: u16, - #[case] exception_name: &str, + #[case] code: Option, + #[case] kind: &str, ) { Python::initialize(); Python::attach(|py| { - let error = Error::Storage(litellm_storage_clickhouse::Error::QueryFailed(status)); - let message = error.to_string(); + initialize_python_path(py); + let error = Error::Storage(litellm_storage_clickhouse::Error::QueryFailed( + litellm_storage_clickhouse::QueryFailure { + status, + code, + message: "engine diagnostic".into(), + }, + )); + let curated = map_error_ref(&error); + assert_eq!(curated.get_type(py).name().unwrap(), "RuntimeError"); let exception = map_sql_error(error); - assert_eq!(exception.get_type(py).name().unwrap(), exception_name); + assert_eq!(exception.get_type(py).name().unwrap(), "TraceQueryError"); + assert_eq!( + exception + .value(py) + .getattr("kind") + .unwrap() + .extract::() + .unwrap(), + kind + ); + assert_eq!( + exception + .value(py) + .getattr("database_code") + .unwrap() + .extract::>() + .unwrap(), + code + ); assert_eq!( exception.value(py).str().unwrap().to_str().unwrap(), - message + "engine diagnostic" ); }); } + #[rstest] + fn raw_query_transport_failures_preserve_generic_exception_type() { + Python::initialize(); + Python::attach(|py| { + let exception = + map_sql_error(Error::Storage(litellm_storage_clickhouse::Error::Transport)); + assert_eq!(exception.get_type(py).name().unwrap(), "RuntimeError"); + }); + } + #[rstest] #[case::decode_budget(Error::Decode(litellm_traces::Error::TooLarge), "OverflowError")] #[case::invalid_export(Error::Decode(litellm_traces::Error::InvalidPayload), "ValueError")] @@ -655,18 +764,7 @@ mod tests { ) { Python::initialize(); Python::attach(|py| { - let repository = std::path::Path::new(env!("CARGO_MANIFEST_DIR")) - .ancestors() - .nth(3) - .unwrap() - .to_str() - .unwrap(); - pyo3::types::PyModule::import(py, "sys") - .unwrap() - .getattr("path") - .unwrap() - .call_method1("insert", (0, repository)) - .unwrap(); + initialize_python_path(py); let message = error.to_string(); let exception = map_read_error(error); assert_eq!(exception.get_type(py).name().unwrap(), exception_name); diff --git a/litellm-rust/crates/storage-clickhouse/src/error.rs b/litellm-rust/crates/storage-clickhouse/src/error.rs index acec9b91675..b143cbe3532 100644 --- a/litellm-rust/crates/storage-clickhouse/src/error.rs +++ b/litellm-rust/crates/storage-clickhouse/src/error.rs @@ -1,3 +1,11 @@ +#[derive(Debug, thiserror::Error)] +#[error("ClickHouse query failed with HTTP status {status}: {message}")] +pub struct QueryFailure { + pub status: u16, + pub code: Option, + pub message: String, +} + #[derive(Debug, thiserror::Error)] pub enum Error { #[error("invalid ClickHouse insert row")] @@ -16,8 +24,8 @@ pub enum Error { InvalidParameters, #[error("unknown ClickHouse read query")] InvalidQuery, - #[error("ClickHouse query failed with HTTP status {0}")] - QueryFailed(u16), + #[error(transparent)] + QueryFailed(QueryFailure), #[error("ClickHouse insert failed with HTTP status {0}")] InsertFailed(u16), #[error("ClickHouse insert exceeds the encoded size limit")] diff --git a/litellm-rust/crates/storage-clickhouse/src/lib.rs b/litellm-rust/crates/storage-clickhouse/src/lib.rs index 7ab2aa9bc0a..9c87e253821 100644 --- a/litellm-rust/crates/storage-clickhouse/src/lib.rs +++ b/litellm-rust/crates/storage-clickhouse/src/lib.rs @@ -2,7 +2,7 @@ mod error; mod insert; mod read; -pub use error::Error; +pub use error::{Error, QueryFailure}; pub use insert::{insert_compressed_rows, insert_encoded_rows}; pub use read::{Parameter, Query, READ_LIMITS, ReadLimits, execute_read, fetch, fetch_json}; use url::Url; diff --git a/litellm-rust/crates/storage-clickhouse/src/read.rs b/litellm-rust/crates/storage-clickhouse/src/read.rs index c4bfdef393a..85b5c758e94 100644 --- a/litellm-rust/crates/storage-clickhouse/src/read.rs +++ b/litellm-rust/crates/storage-clickhouse/src/read.rs @@ -3,7 +3,7 @@ use std::{collections::BTreeMap, time::Duration}; use litellm_http::Client; use serde::{Deserialize, Serialize, de::DeserializeOwned}; -use crate::{Connection, Error}; +use crate::{Connection, Error, QueryFailure}; #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub struct ReadLimits { @@ -107,17 +107,38 @@ pub async fn execute_read( let request = client .post(url) .timeout(Duration::from_secs(15)) + .header("X-ClickHouse-Format", "JSON") .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); + let status = response.status().as_u16(); + let code = response + .headers() + .get("x-clickhouse-exception-code") + .and_then(|value| value.to_str().ok()) + .and_then(|value| value.parse().ok()); + if !response.status().is_success() || code.is_some() { + let mut body = Vec::new(); + while let Some(chunk) = response.chunk().await.map_err(|_| Error::Transport)? { + let remaining = DIAGNOSTIC_BYTES - body.len(); + body.extend_from_slice(&chunk[..chunk.len().min(remaining)]); + if body.len() == DIAGNOSTIC_BYTES { + break; + } } - return Err(Error::QueryFailed(response.status().as_u16())); + let message = serde_json::from_slice::(&body) + .ok() + .and_then(|json| { + json.get("exception") + .and_then(|value| value.as_str()) + .map(str::to_owned) + }) + .unwrap_or_else(|| { + let valid = std::str::from_utf8(&body) + .err() + .map_or(body.len(), |error| error.valid_up_to()); + String::from_utf8_lossy(&body[..valid]).into_owned() + }); + return Err(query_failure(status, code, &message)); } let mut body = Vec::new(); @@ -130,13 +151,42 @@ pub async fn execute_read( let json: serde_json::Value = serde_json::from_slice(&body).map_err(|_| Error::InvalidResponse)?; - if json.get("exception").is_some() || !json.get("data").is_some_and(serde_json::Value::is_array) - { + if let Some(message) = json.get("exception").and_then(serde_json::Value::as_str) { + return Err(query_failure(status, code, message)); + } + if !json.get("data").is_some_and(serde_json::Value::is_array) { return Err(Error::InvalidResponse); } String::from_utf8(body).map_err(|_| Error::InvalidResponse) } +const DIAGNOSTIC_BYTES: usize = 4096; + +fn query_failure(status: u16, code: Option, message: &str) -> Error { + let code = code.or_else(|| { + message + .strip_prefix("Code: ")? + .split_once('.')? + .0 + .parse() + .ok() + }); + let message = if message.is_empty() { + "ClickHouse query failed".to_owned() + } else { + let end = (0..=DIAGNOSTIC_BYTES.min(message.len())) + .rev() + .find(|&end| message.is_char_boundary(end)) + .unwrap_or(0); + message[..end].to_owned() + }; + Error::QueryFailed(QueryFailure { + status, + code, + message, + }) +} + pub trait Query { type Params: Serialize; type Row: DeserializeOwned; diff --git a/litellm-rust/crates/storage-clickhouse/tests/transport.rs b/litellm-rust/crates/storage-clickhouse/tests/transport.rs index f70fe52e1de..b6a1dff4e11 100644 --- a/litellm-rust/crates/storage-clickhouse/tests/transport.rs +++ b/litellm-rust/crates/storage-clickhouse/tests/transport.rs @@ -107,25 +107,42 @@ async fn typed_fetch_encodes_parameters_and_validates_rows( ); assert_eq!(envelope.unwrap(), body); } else { - assert!(matches!(rows, Err(Error::InvalidResponse))); - assert!(matches!(envelope, Err(Error::InvalidResponse))); + assert!(rows.is_err()); + assert!(envelope.is_err()); } } #[rstest] -#[case::result_limit("396", true)] -#[case::memory_limit("241", false)] -#[case::timeout("159", false)] -#[case::unknown("", false)] +#[case::syntax(500, Some("62"), "Code: 62. Invalid SQL", Some(62))] +#[case::readonly(500, Some("164"), "Code: 164. Writes are denied", Some(164))] +#[case::memory_limit(500, Some("241"), "Memory limit exceeded", Some(241))] +#[case::timeout(408, Some("159"), "Time limit exceeded", Some(159))] +#[case::result_limit(500, Some("396"), "Result limit exceeded", Some(396))] +#[case::status_only(503, None, "Server unavailable", None)] +#[case::body_code(500, None, "Code: 47. Unknown identifier", Some(47))] +#[case::late_exception( + 200, + None, + r#"{"data":[],"exception":"Code: 62. Invalid SQL"}"#, + Some(62) +)] #[tokio::test] -async fn server_result_limits_allow_smaller_pages_without_retrying_other_failures( - #[case] code: &str, - #[case] result_limit: bool, +async fn bounded_read_preserves_database_diagnostics( + #[case] status: u16, + #[case] header: Option<&str>, + #[case] body: &str, + #[case] expected_code: Option, ) { use wiremock::{Mock, MockServer, ResponseTemplate, matchers::method}; let server = MockServer::start().await; + let response = match header { + Some(code) => { + ResponseTemplate::new(status).insert_header("X-ClickHouse-Exception-Code", code) + } + None => ResponseTemplate::new(status), + }; Mock::given(method("POST")) - .respond_with(ResponseTemplate::new(500).insert_header("X-ClickHouse-Exception-Code", code)) + .respond_with(response.set_body_string(body)) .expect(1) .mount(&server) .await; @@ -138,11 +155,77 @@ async fn server_result_limits_allow_smaller_pages_without_retrying_other_failure ) .await .unwrap_err(); - if result_limit { - assert!(matches!(error, Error::ResponseTooLarge)); + let Error::QueryFailed(failure) = error else { + panic!("expected database diagnostic: {error:?}"); + }; + assert_eq!(failure.status, status); + assert_eq!(failure.code, expected_code); + let expected_message = serde_json::from_str::(body) + .ok() + .and_then(|json| { + json.get("exception") + .and_then(|value| value.as_str()) + .map(str::to_owned) + }) + .unwrap_or_else(|| body.to_owned()); + assert_eq!(failure.message, expected_message); +} + +#[rstest] +#[case::http_failure(500)] +#[case::late_json_failure(200)] +#[tokio::test] +async fn database_diagnostics_are_bounded_and_preserve_unicode(#[case] status: u16) { + use wiremock::{Mock, MockServer, ResponseTemplate, matchers::method}; + let server = MockServer::start().await; + let message = format!("Code: 62. x{}", "雪".repeat(10_000)); + let body = if status == 200 { + serde_json::json!({"data": [], "exception": message}).to_string() } else { - assert!(matches!(error, Error::QueryFailed(500))); - } + message.clone() + }; + Mock::given(method("POST")) + .respond_with(ResponseTemplate::new(status).set_body_string(body)) + .expect(1) + .mount(&server) + .await; + let error = execute_read( + &Client::no_redirect_for_test(), + &Connection::parse(&server.uri()).unwrap(), + "SELECT 1", + &BTreeMap::new(), + ) + .await + .unwrap_err(); + let Error::QueryFailed(failure) = error else { + panic!("expected database diagnostic: {error:?}"); + }; + assert_eq!(failure.code, Some(62)); + assert!(failure.message.len() <= 4096); + assert!(message.starts_with(&failure.message)); +} + +#[rstest] +#[tokio::test] +async fn bounded_read_preserves_transport_failures() { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let connection = + Connection::parse(&format!("http://{}", listener.local_addr().unwrap())).unwrap(); + let peer = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + drop(stream); + }); + assert!(matches!( + execute_read( + &Client::no_redirect_for_test(), + &connection, + "SELECT 1", + &BTreeMap::new() + ) + .await, + Err(Error::Transport) + )); + peer.await.unwrap(); } #[test] diff --git a/litellm-rust/crates/traces-cache/src/cache.rs b/litellm-rust/crates/traces-cache/src/cache.rs index 367526ea613..681e04bf18e 100644 --- a/litellm-rust/crates/traces-cache/src/cache.rs +++ b/litellm-rust/crates/traces-cache/src/cache.rs @@ -1,6 +1,6 @@ use std::{future::Future, sync::Arc, time::Duration}; -use litellm_traces::{QueryScope, Trace, TraceSummary, store::SpanRow}; +use litellm_traces::{QueryScope, Trace, TraceSummary, search::RunFilter, store::SpanRow}; use moka::{Expiry, future::Cache}; use serde::Serialize; use sha2::{Digest, Sha256}; @@ -50,6 +50,14 @@ impl SnapshotKey { Self::digest(&("run", source, access, run)) } + pub(crate) fn run_page_scope( + source: &str, + access: &QueryScope, + filter: &RunFilter, + ) -> Result { + Self::digest(&("run_page_v1", source, access, filter)).map(|key| key.0) + } + pub(crate) fn scope(source: &str, access: &QueryScope) -> Result { Self::digest(&("scope", source, access)) } diff --git a/litellm-rust/crates/traces-cache/src/cursor.rs b/litellm-rust/crates/traces-cache/src/cursor.rs index 82d5a61e282..d5e5acdc7a1 100644 --- a/litellm-rust/crates/traces-cache/src/cursor.rs +++ b/litellm-rust/crates/traces-cache/src/cursor.rs @@ -44,15 +44,24 @@ impl Cursor { #[serde(deny_unknown_fields)] pub(super) struct RunPosition { order: RunOrder, + query_scope: String, + window: (i64, i64), value: i64, trace_ref: String, } impl RunPosition { - pub(super) fn after(order: RunOrder, row: &RunRow) -> Self { + pub(super) fn after( + order: RunOrder, + row: &RunRow, + query_scope: &str, + window: (i64, i64), + ) -> Self { let RunCursor { value, trace_ref } = order.cursor(row); Self { order, + query_scope: query_scope.to_owned(), + window, value, trace_ref, } @@ -79,12 +88,19 @@ pub(super) struct TextPosition { pub(super) fn run_position( cursor: Option<&str>, order: RunOrder, + query_scope: &str, + window: (i64, i64), ) -> Result, ReadError> { let Some(cursor) = cursor.filter(|cursor| !cursor.is_empty()) else { return Ok(None); }; match Cursor::decode(cursor, "trace")? { - Cursor::Run(position) if position.order == order && !position.trace_ref.is_empty() => { + Cursor::Run(position) + if position.order == order + && !position.trace_ref.is_empty() + && position.query_scope == query_scope + && position.window == window => + { Ok(Some(RunCursor { value: position.value, trace_ref: position.trace_ref, @@ -94,6 +110,35 @@ pub(super) fn run_position( } } +pub fn resolve_run_window( + start_ms: Option, + end_ms: Option, + cursor: Option<&str>, + default_window: (i64, i64), +) -> Result<(i64, i64), ReadError> { + let Some(cursor) = cursor.filter(|cursor| !cursor.is_empty()) else { + let window = ( + start_ms.unwrap_or(default_window.0), + end_ms.unwrap_or(default_window.1), + ); + return if window.0 < window.1 { + Ok(window) + } else { + Err(ReadError::InvalidParameters) + }; + }; + let Cursor::Run(position) = Cursor::decode(cursor, "trace")? else { + return Err(ReadError::InvalidCursor("trace")); + }; + if position.window.0 >= position.window.1 + || start_ms.is_some_and(|start| start != position.window.0) + || end_ms.is_some_and(|end| end != position.window.1) + { + return Err(ReadError::InvalidCursor("trace")); + } + Ok(position.window) +} + pub(super) fn span_position(cursor: &str) -> Result> { match Cursor::decode(cursor, "span")? { Cursor::Span(position) => Ok(position), @@ -129,9 +174,13 @@ mod tests { use super::*; + const WINDOW: (i64, i64) = (10, 100); + fn run(order: RunOrder, value: i64, trace_ref: &str) -> String { Cursor::Run(RunPosition { order, + query_scope: "query".into(), + window: WINDOW, value, trace_ref: trace_ref.into(), }) @@ -169,22 +218,48 @@ mod tests { #[rstest] #[case::newest(RunOrder::NEWEST, 1_790_742_989_377)] #[case::zero_value(BY_ERRORS, 0)] - fn run_cursor_round_trips_under_its_own_order(#[case] order: RunOrder, #[case] value: i64) { - let position = run_position::(Some(&run(order, value, "4BAD")), order) - .unwrap() - .unwrap(); + fn run_cursor_round_trips_under_its_query(#[case] order: RunOrder, #[case] value: i64) { + let position = run_position::( + Some(&run(order, value, "4BAD")), + order, + "query", + WINDOW, + ) + .unwrap() + .unwrap(); assert_eq!( (position.value, position.trace_ref.as_str()), (value, "4BAD") ); } + #[rstest] + #[case::order(BY_ERRORS, "query", WINDOW)] + #[case::direction(RunOrder { descending: false, ..RunOrder::NEWEST }, "query", WINDOW)] + #[case::scope(RunOrder::NEWEST, "other-query", WINDOW)] + #[case::window(RunOrder::NEWEST, "query", (11, 100))] + fn run_cursor_rejects_a_changed_query( + #[case] order: RunOrder, + #[case] scope: &str, + #[case] window: (i64, i64), + ) { + assert!(matches!( + run_position::( + Some(&run(RunOrder::NEWEST, 1, "ref")), + order, + scope, + window + ), + Err(ReadError::InvalidCursor("trace")) + )); + } + #[rstest] #[case::absent(None)] #[case::empty(Some(""))] fn missing_run_cursor_starts_from_the_first_page(#[case] cursor: Option<&str>) { assert!( - run_position::(cursor, RunOrder::NEWEST) + run_position::(cursor, RunOrder::NEWEST, "query", WINDOW) .unwrap() .is_none() ); @@ -193,17 +268,59 @@ mod tests { #[rstest] #[case::not_base64("abc".into())] #[case::not_json(URL_SAFE.encode("not-json"))] - #[case::untagged_tuple(json(serde_json::json!([1, "ref"])))] - #[case::other_key(run(BY_ERRORS, 1, "ref"))] - #[case::other_direction(run(RunOrder { descending: false, ..RunOrder::NEWEST }, 1, "ref"))] #[case::empty_ref(run(RunOrder::NEWEST, 1, ""))] #[case::span_cursor(span())] #[case::text_cursor(text(SpanPart::Error, 0, "A".repeat(64)))] - #[case::without_order(json(serde_json::json!({"kind": "run", "position": {"value": 1, "trace_ref": "r"}})))] - #[case::unknown_field(json(serde_json::json!({"kind": "run", "position": {"order": {"key": "start_ms", "descending": true}, "value": 1, "trace_ref": "r", "extra": 1}})))] - fn run_cursors_not_minted_under_the_requested_order_are_rejected(#[case] cursor: String) { + #[case::missing_fields(json(serde_json::json!({"kind": "run", "position": {"value": 1, "trace_ref": "r"}})))] + fn malformed_run_cursors_are_rejected(#[case] cursor: String) { assert!(matches!( - run_position::(Some(&cursor), RunOrder::NEWEST), + run_position::(Some(&cursor), RunOrder::NEWEST, "query", WINDOW), + Err(ReadError::InvalidCursor("trace")) + )); + } + + #[rstest] + #[case::default(None, None, (0, 50))] + #[case::start(Some(10), None, (10, 50))] + #[case::end(None, Some(40), (0, 40))] + #[case::explicit(Some(20), Some(40), (20, 40))] + fn first_page_resolves_only_missing_window_bounds( + #[case] start: Option, + #[case] end: Option, + #[case] expected: (i64, i64), + ) { + assert_eq!( + resolve_run_window::(start, end, None, (0, 50)).unwrap(), + expected + ); + } + + #[rstest] + #[case::omitted(None, None)] + #[case::start(Some(WINDOW.0), None)] + #[case::end(None, Some(WINDOW.1))] + #[case::explicit(Some(WINDOW.0), Some(WINDOW.1))] + fn cursor_keeps_its_window_when_the_default_clock_advances( + #[case] start: Option, + #[case] end: Option, + ) { + let cursor = run(RunOrder::NEWEST, 1, "ref"); + assert_eq!( + resolve_run_window::(start, end, Some(&cursor), (200, 300)).unwrap(), + WINDOW + ); + } + + #[rstest] + #[case::start(Some(WINDOW.0 + 1), None)] + #[case::end(None, Some(WINDOW.1 + 1))] + fn cursor_rejects_explicit_window_changes( + #[case] start: Option, + #[case] end: Option, + ) { + let cursor = run(RunOrder::NEWEST, 1, "ref"); + assert!(matches!( + resolve_run_window::(start, end, Some(&cursor), (200, 300)), Err(ReadError::InvalidCursor("trace")) )); } diff --git a/litellm-rust/crates/traces-cache/src/lib.rs b/litellm-rust/crates/traces-cache/src/lib.rs index ece59fec654..295081389ee 100644 --- a/litellm-rust/crates/traces-cache/src/lib.rs +++ b/litellm-rust/crates/traces-cache/src/lib.rs @@ -8,6 +8,7 @@ mod spend; mod store; pub use cache::{Freshness, LIVE_TTL, SETTLED_TTL, Snapshot, SnapshotCache, SnapshotKey}; +pub use cursor::resolve_run_window; pub use error::{Error, ReadError}; pub use reader::{MAX_GRAPH_BYTES, MAX_GRAPH_SPANS, MAX_TEXT_SPANS, PageRequest, TraceReader}; pub use store::{StoreError, StoreResult, TraceStore}; diff --git a/litellm-rust/crates/traces-cache/src/list.rs b/litellm-rust/crates/traces-cache/src/list.rs index d70a6339721..6d5eddf88d4 100644 --- a/litellm-rust/crates/traces-cache/src/list.rs +++ b/litellm-rust/crates/traces-cache/src/list.rs @@ -41,9 +41,16 @@ fn cache_key( } fn summary(row: &RunRow, listed: Option<&ListedRun>) -> TraceSummary { - match listed { + let summary = match listed { Some(ListedRun::Resolved(summary, _)) => (**summary).clone(), Some(ListedRun::Limited) | None => listed_summary(row), + }; + TraceSummary { + start_time: litellm_traces::iso_time(row.start_ms), + duration_ms: row.duration_ns as f64 / 1_000_000.0, + span_count: row.span_count, + error_count: row.error_count, + ..summary } } @@ -98,7 +105,11 @@ async fn resolve_runs( let (Some(start_ms), Some(end_ms)) = ( runs.iter().map(|row| row.start_ms).min(), runs.iter() - .map(|row| row.start_ms.saturating_add(row.duration_ms)) + .map(|row| { + row.start_ms.saturating_add( + i64::try_from(row.duration_ns.div_ceil(1_000_000)).unwrap_or(i64::MAX), + ) + }) .max(), ) else { return Ok(Vec::new()); diff --git a/litellm-rust/crates/traces-cache/src/reader.rs b/litellm-rust/crates/traces-cache/src/reader.rs index 9c94c63b60c..0b679c5cbdc 100644 --- a/litellm-rust/crates/traces-cache/src/reader.rs +++ b/litellm-rust/crates/traces-cache/src/reader.rs @@ -86,7 +86,13 @@ impl TraceReader { if page.limit == 0 || filter.start_ms >= filter.end_ms { return Err(ReadError::InvalidParameters); } - let after = run_position(page.cursor.as_deref(), order)?; + let query_scope = SnapshotKey::run_page_scope(store.source(), access, filter)?; + let after = run_position( + page.cursor.as_deref(), + order, + &query_scope, + (filter.start_ms, filter.end_ms), + )?; let scope = SnapshotKey::scope(store.source(), access)?; let accepted = self.lists.limits.get(&scope).await.unwrap_or(u32::MAX); let mut page_size = page.limit.min(500).min(accepted); @@ -108,10 +114,15 @@ impl TraceReader { }; let more = rows.len() > page_size as usize; rows.truncate(page_size as usize); - let next_cursor = rows - .last() - .filter(|_| more) - .map(|last| Cursor::Run(RunPosition::after(order, last)).encode()); + let next_cursor = rows.last().filter(|_| more).map(|last| { + Cursor::Run(RunPosition::after( + order, + last, + &query_scope, + (filter.start_ms, filter.end_ms), + )) + .encode() + }); let data = { let mut summaries = Vec::with_capacity(rows.len()); for batch in run_batches(&rows) { diff --git a/litellm-rust/crates/traces-cache/tests/read.rs b/litellm-rust/crates/traces-cache/tests/read.rs index 39c0a103993..f32d5dd025e 100644 --- a/litellm-rust/crates/traces-cache/tests/read.rs +++ b/litellm-rust/crates/traces-cache/tests/read.rs @@ -8,18 +8,18 @@ use std::{ }; use litellm_traces::{ - CallEvidenceKind, CallKey, ObservationType, QueryScope, SpanStatus, + CallEvidenceKind, CallKey, ObservationType, QueryScope, SpanStatus, TraceSummary, search::{AgentRuns, HistogramBucket, RunField, RunFilter, RunSearch}, store::{ CallQuery, CallRow, CountBy, CountValue, RunCount, RunCountQuery, RunOrder, RunQuery, - RunRow, RunSelection, SpanPart, SpanQuery, SpanRow, SpanSelection, SpanText, SpanTextQuery, - TextRange, + RunRow, RunSelection, RunSortKey, SpanPart, SpanQuery, SpanRow, SpanSelection, SpanText, + SpanTextQuery, TextRange, }, }; use litellm_traces_cache::{ LIVE_TTL, PageRequest, ReadError, StoreError, StoreResult, TraceReader, TraceStore, }; -use rstest::rstest; +use rstest::{fixture, rstest}; const START_NS: i64 = 1_790_742_989_000_000_000; @@ -61,6 +61,7 @@ struct State { #[derive(Default)] struct FakeStore { + source: Option<&'static str>, state: Mutex, calls: Mutex>, } @@ -68,6 +69,7 @@ struct FakeStore { impl FakeStore { fn with_spans(trace_ref: &str, spans: Vec) -> Self { Self { + source: None, state: Mutex::new(State { trace_spans: HashMap::from([(trace_ref.to_owned(), spans)]), ..State::default() @@ -166,7 +168,7 @@ impl TraceStore for FakeStore { type Error = FakeError; fn source(&self) -> &str { - "fake" + self.source.unwrap_or("fake") } async fn runs(&self, _: &QueryScope, query: &RunQuery) -> StoreResult, FakeError> { @@ -189,12 +191,25 @@ impl TraceStore for FakeStore { { return Err(StoreError::TooLarge); } - Ok(state + let mut rows: Vec<_> = state .list_runs .iter() - .take(query.limit as usize) + .filter(|row| { + query.after.as_ref().is_none_or(|after| { + let value = (query.order.value(row), &row.trace_ref); + let cursor = (after.value, &after.trace_ref); + if query.order.descending { + value < cursor + } else { + value > cursor + } + }) + }) .cloned() - .collect()) + .collect(); + rows.sort_by(|left, right| query.order.compare(left, right)); + rows.truncate(query.limit as usize); + Ok(rows) } async fn run_counts( @@ -227,10 +242,18 @@ impl TraceStore for FakeStore { .cloned() .unwrap_or_default() } - SpanSelection::Runs { .. } => { + SpanSelection::Runs { window, .. } => { self.record(Operation::RunSpans); Self::failure(&state, Operation::RunSpans)?; - state.run_spans.clone() + state + .run_spans + .iter() + .filter(|row| { + i128::from(row.start_ns) >= i128::from(window.start) * 1_000_000 + && i128::from(row.start_ns) < i128::from(window.end) * 1_000_000 + }) + .cloned() + .collect() } }; Ok(keyset( @@ -328,7 +351,6 @@ fn newest(limit: u32) -> PageRequest { PageRequest { cursor: None, limit, - ..Default::default() } } @@ -378,7 +400,7 @@ fn run(trace_id: &str, trace_ref: &str) -> RunRow { input_preview: String::new(), status: SpanStatus::Ok, start_ms: 1_790_742_989_000, - duration_ms: 10, + duration_ns: 10_000_000, span_count: 1, agent_count: 1, agent_invocations: 1, @@ -526,6 +548,55 @@ async fn response_size_splits_pages_and_rejects_a_single_oversized_span() { )); } +#[rstest] +#[tokio::test] +async fn listed_run_resolution_keeps_spans_crossing_a_fractional_millisecond() { + let store = FakeStore::default(); + let root = SpanRow { + start_ns: START_NS + 900_000, + duration_ns: 900_000, + ..span(0) + }; + let child = SpanRow { + start_ns: START_NS + 1_200_000, + duration_ns: 100_000, + status: SpanStatus::Error, + agent: "child".into(), + ..span(1) + }; + let listed = RunRow { + duration_ns: root.duration_ns, + span_count: 2, + error_count: 1, + agent_count: 2, + agent_invocations: 2, + agent_names: vec!["agent".into(), "child".into()], + ..run("trace", "ref") + }; + store.set_list_runs(vec![listed.clone()]); + store.set_run_spans(vec![root, child]); + let page = TraceReader::new(usize::MAX) + .list_traces( + &store, + &access(), + &everything(), + RunOrder::NEWEST, + &newest(1), + ) + .await + .unwrap(); + assert_eq!(page.data.len(), 1); + assert_eq!(page.data[0].span_count, listed.span_count); + assert_eq!(page.data[0].error_count, listed.error_count); + assert_eq!(page.data[0].agent_count, listed.agent_count); + assert_eq!(page.data[0].agent_names, listed.agent_names); + assert_eq!( + page.data[0].duration_ms, + listed.duration_ns as f64 / 1_000_000.0 + ); + assert!(!page.data[0].resolution_limited); +} + #[rstest] #[tokio::test] async fn list_run_budget_halves_the_limit_and_cursor_requires_a_run_past_the_page() { @@ -561,6 +632,120 @@ async fn list_run_budget_halves_the_limit_and_cursor_requires_a_run_past_the_pag } } +#[fixture] +fn paging_scope() -> QueryScope { + QueryScope::Owned { + user_id: "user".into(), + team_ids: vec!["team".into()], + } +} + +#[rstest] +#[case::start(window(1, i64::MAX, ""), paging_scope(), "fake", RunOrder::NEWEST)] +#[case::end(window(0, i64::MAX - 1, ""), paging_scope(), "fake", RunOrder::NEWEST)] +#[case::text(window(0, i64::MAX, "find"), paging_scope(), "fake", RunOrder::NEWEST)] +#[case::field( + window(0, i64::MAX, "agent:worker"), + paging_scope(), + "fake", + RunOrder::NEWEST +)] +#[case::attribute( + window(0, i64::MAX, "attr.stage:production"), + paging_scope(), + "fake", + RunOrder::NEWEST +)] +#[case::trace_refs(RunFilter { trace_refs: vec!["ref".into()], ..everything() }, paging_scope(), "fake", RunOrder::NEWEST)] +#[case::user(everything(), QueryScope::Owned { user_id: "another-user".into(), team_ids: vec!["team".into()] }, "fake", RunOrder::NEWEST)] +#[case::teams(everything(), QueryScope::Owned { user_id: "user".into(), team_ids: vec!["another-team".into()] }, "fake", RunOrder::NEWEST)] +#[case::source(everything(), paging_scope(), "other", RunOrder::NEWEST)] +#[case::sort(everything(), paging_scope(), "fake", RunOrder { descending: false, ..RunOrder::NEWEST })] +#[tokio::test] +async fn run_cursors_reject_a_changed_query_before_reading_storage( + #[case] filter: RunFilter, + #[case] scope: QueryScope, + #[case] source: &'static str, + #[case] order: RunOrder, +) { + let original = FakeStore::default(); + original.set_list_runs(vec![run("trace-a", "ref-a"), run("trace-b", "ref-b")]); + let reader = TraceReader::new(usize::MAX); + let first = reader + .list_traces( + &original, + &paging_scope(), + &everything(), + RunOrder::NEWEST, + &newest(1), + ) + .await + .unwrap(); + let changed = FakeStore { + source: Some(source), + ..Default::default() + }; + let result = reader + .list_traces( + &changed, + &scope, + &filter, + order, + &PageRequest { + cursor: first.next_cursor, + limit: 1, + }, + ) + .await; + assert!(matches!(result, Err(ReadError::InvalidCursor("trace")))); + assert_eq!(changed.calls(Operation::ListRuns), 0); +} + +#[rstest] +#[tokio::test] +async fn run_cursors_continue_without_gaps_when_page_size_changes() { + let store = FakeStore::default(); + store.set_list_runs(vec![ + run("trace-a", "ref-a"), + run("trace-b", "ref-b"), + run("trace-c", "ref-c"), + ]); + let reader = TraceReader::new(usize::MAX); + let first = reader + .list_traces( + &store, + &access(), + &everything(), + RunOrder::NEWEST, + &newest(1), + ) + .await + .unwrap(); + let cursor = first.next_cursor.unwrap(); + let second = reader + .list_traces( + &store, + &access(), + &everything(), + RunOrder::NEWEST, + &PageRequest { + cursor: Some(cursor), + limit: 2, + }, + ) + .await + .unwrap(); + let refs: Vec<_> = first + .data + .iter() + .chain(&second.data) + .map(|run| run.trace_ref.as_str()) + .collect(); + assert_eq!(refs, ["ref-c", "ref-b", "ref-a"]); + assert!(second.next_cursor.is_none()); + assert_eq!(store.calls(Operation::ListRuns), 2); +} + #[rstest] #[tokio::test] async fn oversized_run_batch_falls_back_to_each_run_and_keeps_listed_summaries() { @@ -955,6 +1140,71 @@ async fn failed_reads_are_not_cached() { assert_eq!(store.calls(Operation::TraceSpans), 2); } +#[rstest] +#[case::start(RunSortKey::StartMs)] +#[case::duration(RunSortKey::DurationMs)] +#[case::spans(RunSortKey::SpanCount)] +#[case::errors(RunSortKey::ErrorCount)] +#[case::reference(RunSortKey::TraceRef)] +#[tokio::test] +async fn cached_run_keeps_all_canonical_metrics_current_for_every_sort(#[case] key: RunSortKey) { + let store = FakeStore::default(); + store.set_list_runs(vec![run("trace", "ref")]); + store.set_run_spans(vec![SpanRow { + kind: ObservationType::Llm, + litellm_request_id: "response".into(), + call_keys: vec![CallKey::ProviderResponse("response".into())], + call_evidence: Some(CallEvidenceKind::Complete), + ..span(0) + }]); + store.state.lock().unwrap().spend = vec![spend_row("response", 1.5)]; + let reader = TraceReader::new(usize::MAX); + let first = reader + .list_traces( + &store, + &access(), + &everything(), + RunOrder::NEWEST, + &newest(1), + ) + .await + .unwrap(); + let changed = RunRow { + start_ms: START_NS / 1_000_000 + 20, + duration_ns: 23_000_000, + span_count: 7, + error_count: 3, + agent_names: vec!["new-agent".into()], + frameworks: vec!["new-framework".into()], + ..run("trace", "ref") + }; + store.set_list_runs(vec![changed.clone()]); + let second = reader + .list_traces( + &store, + &access(), + &everything(), + RunOrder { + key, + descending: false, + }, + &newest(1), + ) + .await + .unwrap(); + let expected = TraceSummary { + start_time: litellm_traces::iso_time(changed.start_ms), + duration_ms: changed.duration_ns as f64 / 1_000_000.0, + span_count: changed.span_count, + error_count: changed.error_count, + ..first.data[0].clone() + }; + assert_eq!(expected.spend, Some(1.5)); + assert_eq!(second.data, vec![expected]); + assert_eq!(store.calls(Operation::RunSpans), 1); + assert_eq!(store.calls(Operation::Spend), 1); +} + #[rstest] #[tokio::test] async fn listed_runs_are_read_once_until_a_live_run_expires() { @@ -969,7 +1219,10 @@ async fn listed_runs_are_read_once_until_a_live_run_expires() { }; let store = FakeStore::default(); store.set_list_runs(vec![ - run("trace-live", "ref-live"), + RunRow { + start_ms: live.start_ns / 1_000_000, + ..run("trace-live", "ref-live") + }, run("trace-settled", "ref-settled"), ]); store.set_run_spans(vec![live, settled]); diff --git a/litellm-rust/crates/traces-clickhouse/migrations/0001_otel_traces.sql b/litellm-rust/crates/traces-clickhouse/migrations/0001_otel_traces.sql index fb5eaa367d7..0ae650f43db 100644 --- a/litellm-rust/crates/traces-clickhouse/migrations/0001_otel_traces.sql +++ b/litellm-rust/crates/traces-clickhouse/migrations/0001_otel_traces.sql @@ -31,6 +31,7 @@ CREATE TABLE IF NOT EXISTS {database}.otel_traces SpanAttributes['gen_ai.operation.name'] = 'execute_tool', 'tool', 'chain'), AgentName LowCardinality(String) DEFAULT SpanAttributes['gen_ai.agent.name'], + Framework LowCardinality(String) DEFAULT '', LiteLLMRequestId String DEFAULT SpanAttributes['gen_ai.response.id'], Model LowCardinality(String) DEFAULT SpanAttributes['gen_ai.request.model'], InputTokens UInt32 DEFAULT toUInt32OrZero(SpanAttributes['gen_ai.usage.input_tokens']), diff --git a/litellm-rust/crates/traces-clickhouse/migrations/0003_agent_traces.sql b/litellm-rust/crates/traces-clickhouse/migrations/0003_agent_traces.sql index 821cc2f3723..2ac58f9bb8c 100644 --- a/litellm-rust/crates/traces-clickhouse/migrations/0003_agent_traces.sql +++ b/litellm-rust/crates/traces-clickhouse/migrations/0003_agent_traces.sql @@ -18,6 +18,8 @@ CREATE TABLE IF NOT EXISTS {database}.agent_traces_by_key OutputTokens SimpleAggregateFunction(sum, UInt64), Models SimpleAggregateFunction(groupUniqArrayArray, Array(String)), AgentNames SimpleAggregateFunction(groupUniqArrayArray, Array(String)), + AgentIdentities SimpleAggregateFunction(groupUniqArrayArray, Array(String)), + Frameworks SimpleAggregateFunction(groupUniqArrayArray, Array(String)), RequestIds SimpleAggregateFunction(groupArrayArray, Array(String)) ) ENGINE = AggregatingMergeTree diff --git a/litellm-rust/crates/traces-clickhouse/migrations/0005_agent_traces_mv.sql b/litellm-rust/crates/traces-clickhouse/migrations/0005_agent_traces_mv.sql index 94dad81f998..082e1410ed6 100644 --- a/litellm-rust/crates/traces-clickhouse/migrations/0005_agent_traces_mv.sql +++ b/litellm-rust/crates/traces-clickhouse/migrations/0005_agent_traces_mv.sql @@ -16,7 +16,9 @@ SELECT sum(InputTokens) AS InputTokens, sum(OutputTokens) AS OutputTokens, groupUniqArrayIf(toString(Model), Model != '') AS Models, - groupUniqArrayIf(SpanName, ObservationType = 'agent') AS AgentNames, + groupUniqArrayIf(if(AgentName = '', SpanName, AgentName), AgentName != '' OR ObservationType = 'agent') AS AgentNames, + groupUniqArrayIf(if(AgentName = '', SpanName, AgentName), ObservationType = 'agent') AS AgentIdentities, + groupUniqArrayIf(toString(Framework), Framework != '') AS Frameworks, groupArrayIf(LiteLLMRequestId, LiteLLMRequestId != '') AS RequestIds FROM {database}.otel_traces GROUP BY TeamId, ApiKeyHash, TraceId diff --git a/litellm-rust/crates/traces-clickhouse/migrations/0010_trace_cost_completeness.sql b/litellm-rust/crates/traces-clickhouse/migrations/0010_trace_cost_completeness.sql index af87a6bf40b..bbcd4fd87d2 100644 --- a/litellm-rust/crates/traces-clickhouse/migrations/0010_trace_cost_completeness.sql +++ b/litellm-rust/crates/traces-clickhouse/migrations/0010_trace_cost_completeness.sql @@ -16,7 +16,9 @@ SELECT sum(InputTokens) AS InputTokens, sum(OutputTokens) AS OutputTokens, groupUniqArrayIf(toString(Model), Model != '') AS Models, - groupUniqArrayIf(SpanName, ObservationType = 'agent') AS AgentNames, + groupUniqArrayIf(if(AgentName = '', SpanName, AgentName), AgentName != '' OR ObservationType = 'agent') AS AgentNames, + groupUniqArrayIf(if(AgentName = '', SpanName, AgentName), ObservationType = 'agent') AS AgentIdentities, + groupUniqArrayIf(toString(Framework), Framework != '') AS Frameworks, groupArrayIf(LiteLLMRequestId, ObservationType = 'llm' OR LiteLLMRequestId != '') AS RequestIds FROM {database}.otel_traces GROUP BY TeamId, ApiKeyHash, TraceId diff --git a/litellm-rust/crates/traces-clickhouse/migrations/0016_trace_rollup_agent_labels.sql b/litellm-rust/crates/traces-clickhouse/migrations/0016_trace_rollup_agent_labels.sql deleted file mode 100644 index 25eacc740ac..00000000000 --- a/litellm-rust/crates/traces-clickhouse/migrations/0016_trace_rollup_agent_labels.sql +++ /dev/null @@ -1,4 +0,0 @@ -ALTER TABLE {database}.agent_traces_by_key - ADD COLUMN IF NOT EXISTS AgentLabels SimpleAggregateFunction(groupUniqArrayArray, Array(String)) DEFAULT [], - ADD COLUMN IF NOT EXISTS AgentIdentities SimpleAggregateFunction(groupUniqArrayArray, Array(String)) DEFAULT [], - ADD COLUMN IF NOT EXISTS Frameworks SimpleAggregateFunction(groupUniqArrayArray, Array(String)) DEFAULT [] diff --git a/litellm-rust/crates/traces-clickhouse/migrations/0017_trace_rollup_agent_labels_mv.sql b/litellm-rust/crates/traces-clickhouse/migrations/0017_trace_rollup_agent_labels_mv.sql deleted file mode 100644 index 1fd7b1295a4..00000000000 --- a/litellm-rust/crates/traces-clickhouse/migrations/0017_trace_rollup_agent_labels_mv.sql +++ /dev/null @@ -1,25 +0,0 @@ -ALTER TABLE {database}.agent_traces_by_key_mv MODIFY QUERY -SELECT - TeamId, ApiKeyHash, TraceId, groupUniqArray(UserId) AS UserIds, - min(Timestamp) AS StartTs, - max(Timestamp + toIntervalNanosecond(Duration)) AS EndTs, - any(ServiceName) AS ServiceName, - anyLastIf(toNullable(SpanName), ParentSpanId = '') AS RootName, - anyLastIf(toNullable(InputPreview), ParentSpanId = '') AS RootInput, - anyLastIf(toNullable(StatusCode), ParentSpanId = '') AS RootStatus, - count() AS SpanCount, - countIf(ObservationType = 'agent') AS AgentCount, - countIf(ObservationType = 'llm') AS LlmCount, - countIf(ObservationType = 'llm' AND LiteLLMRequestId != '') AS IdentifiedLlmCount, - countIf(ObservationType = 'tool') AS ToolCount, - countIf(StatusCode = 'STATUS_CODE_ERROR') AS ErrorCount, - sum(InputTokens) AS InputTokens, - sum(OutputTokens) AS OutputTokens, - groupUniqArrayIf(toString(Model), Model != '') AS Models, - groupUniqArrayIf(SpanName, ObservationType = 'agent') AS AgentNames, - groupUniqArrayIf(AgentName, AgentName != '') AS AgentLabels, - groupUniqArrayIf(if(AgentName = '', SpanName, AgentName), ObservationType = 'agent') AS AgentIdentities, - groupUniqArrayIf(toString(Framework), Framework != '') AS Frameworks, - groupArrayIf(LiteLLMRequestId, ObservationType = 'llm' OR LiteLLMRequestId != '') AS RequestIds -FROM {database}.otel_traces -GROUP BY TeamId, ApiKeyHash, TraceId diff --git a/litellm-rust/crates/traces-clickhouse/query/help/failed_spans.sql b/litellm-rust/crates/traces-clickhouse/query/help/failed_spans.sql index b0f3cc413c1..4d989acbece 100644 --- a/litellm-rust/crates/traces-clickhouse/query/help/failed_spans.sql +++ b/litellm-rust/crates/traces-clickhouse/query/help/failed_spans.sql @@ -1,7 +1,12 @@ SELECT TeamId AS team, ApiKeyHash AS api_key, TraceId AS trace_id, SpanId AS span_id, StatusMessage AS message -FROM otel_traces -WHERE Timestamp >= now() - INTERVAL 1 DAY - AND StatusCode = 'STATUS_CODE_ERROR' +FROM ( + SELECT * + FROM otel_traces + WHERE Timestamp >= now() - INTERVAL 1 DAY + ORDER BY Timestamp, EngineReceivedMs, StatusMessage, Duration, StatusCode + LIMIT 1 BY TeamId, ApiKeyHash, TraceId, SpanId +) +WHERE StatusCode = 'STATUS_CODE_ERROR' ORDER BY Timestamp DESC, team, api_key, trace_id, span_id LIMIT 100 diff --git a/litellm-rust/crates/traces-clickhouse/query/matching_runs.sql b/litellm-rust/crates/traces-clickhouse/query/matching_runs.sql index 14964d09495..6249cfd4111 100644 --- a/litellm-rust/crates/traces-clickhouse/query/matching_runs.sql +++ b/litellm-rust/crates/traces-clickhouse/query/matching_runs.sql @@ -4,23 +4,50 @@ SELECT TraceId AS trace_id, ifNull(any(RootName), '') AS name, any(ServiceName) AS service, ifNull(any(RootInput), '') AS input_preview, ifNull(any(RootStatus), '') AS status, toUnixTimestamp64Milli(min(StartTs)) AS start_ms, - dateDiff('millisecond', min(StartTs), max(EndTs)) AS duration_ms, - sum(SpanCount) AS span_count, + if(any(metrics.span_count) > 0, + toUInt64(greatest(toUnixTimestamp64Nano(any(metrics.end_ts)) - toUnixTimestamp64Nano(any(metrics.start_ts)), 0)), + toUInt64(greatest(toUnixTimestamp64Nano(max(EndTs)) - toUnixTimestamp64Nano(min(StartTs)), 0))) AS duration_ns, + if(any(metrics.span_count) > 0, any(metrics.span_count), sum(SpanCount)) AS span_count, sum(AgentCount) AS agent_invocations, sum(LlmCount) AS llm_calls, sum(ToolCount) AS tool_calls, sum(InputTokens) AS input_tokens, sum(OutputTokens) AS output_tokens, - groupUniqArrayArray(Models) AS models, sum(ErrorCount) AS error_count, - arraySort(if(empty(groupUniqArrayArray(AgentLabels)), - groupUniqArrayArray(AgentNames), - groupUniqArrayArray(AgentLabels))) AS search_agents, + groupUniqArrayArray(Models) AS models, + if(any(metrics.span_count) > 0, any(metrics.error_count), sum(ErrorCount)) AS error_count, + arraySort(groupUniqArrayArray(AgentNames)) AS search_agents, if(error_count > 0, 'error', 'ok') AS search_status, length(groupUniqArrayArray(AgentIdentities)) AS agent_count, arraySort(groupUniqArrayArray(Frameworks)) AS frameworks FROM owned_runs LEFT JOIN ( SELECT TeamId, ApiKeyHash, TraceId, - groupUniqArrayArray(arrayFilter(i -> ResourceAttributes[{attribute_keys:Array(String)}[i]] ILIKE {attribute_patterns:Array(String)}[i] - OR SpanAttributes[{attribute_keys:Array(String)}[i]] ILIKE {attribute_patterns:Array(String)}[i], + min(Timestamp) AS start_ts, + max(Timestamp + toIntervalNanosecond(Duration)) AS end_ts, + count() AS span_count, + countIf(StatusCode = 'STATUS_CODE_ERROR') AS error_count + FROM ( + SELECT TeamId, ApiKeyHash, TraceId, SpanId, Timestamp, Duration, StatusCode + FROM owned_spans + WHERE (TeamId, ApiKeyHash, TraceId) IN ( + SELECT TeamId, ApiKeyHash, TraceId + FROM owned_runs + WHERE {trace_id:String} = '' OR TraceId = {trace_id:String} + GROUP BY TeamId, ApiKeyHash, TraceId + HAVING {trace_id:String} != '' OR ( + min(StartTs) >= fromUnixTimestamp64Milli({start_ms:Int64}) + AND min(StartTs) < fromUnixTimestamp64Milli({end_ms:Int64})) + ) + AND ({trace_id:String} != '' OR Timestamp >= fromUnixTimestamp64Milli({start_ms:Int64})) + ORDER BY Timestamp, EngineReceivedMs, StatusMessage, Duration, StatusCode + LIMIT 1 BY TeamId, ApiKeyHash, TraceId, SpanId + ) + GROUP BY TeamId, ApiKeyHash, TraceId +) AS metrics USING (TeamId, ApiKeyHash, TraceId) +LEFT JOIN ( + SELECT TeamId, ApiKeyHash, TraceId, + groupUniqArrayArray(arrayFilter(i -> (mapContains(ResourceAttributes, {attribute_keys:Array(String)}[i]) + AND ResourceAttributes[{attribute_keys:Array(String)}[i]] ILIKE {attribute_patterns:Array(String)}[i]) + OR (mapContains(SpanAttributes, {attribute_keys:Array(String)}[i]) + AND SpanAttributes[{attribute_keys:Array(String)}[i]] ILIKE {attribute_patterns:Array(String)}[i]), arrayEnumerate({attribute_keys:Array(String)}))) AS matched_attributes FROM owned_spans WHERE notEmpty({attribute_keys:Array(String)}) diff --git a/litellm-rust/crates/traces-clickhouse/query/run_attribute_counts.sql b/litellm-rust/crates/traces-clickhouse/query/run_attribute_counts.sql index f2c3d7969ae..32886c877dd 100644 --- a/litellm-rust/crates/traces-clickhouse/query/run_attribute_counts.sql +++ b/litellm-rust/crates/traces-clickhouse/query/run_attribute_counts.sql @@ -2,7 +2,7 @@ SELECT bucket, failed, value, uniqExact(team_id, api_key_hash, trace_id) AS runs FROM ( SELECT runs.team_id AS team_id, runs.api_key_hash AS api_key_hash, runs.trace_id AS trace_id, if({buckets:UInt32} = 0, toUInt32(0), - toUInt32(intDiv((runs.start_ms - {start_ms:Int64}) * {buckets:UInt32}, {end_ms:Int64} - {start_ms:Int64}))) AS bucket, + toUInt32(intDiv((runs.start_ms - {start_ms:Int64} + 1) * {buckets:UInt32} - 1, {end_ms:Int64} - {start_ms:Int64}))) AS bucket, toUInt8({by_failed:UInt8} = 1 AND runs.error_count > 0) AS failed, value FROM owned_spans AS spans diff --git a/litellm-rust/crates/traces-clickhouse/query/run_counts.sql b/litellm-rust/crates/traces-clickhouse/query/run_counts.sql index 00ad7b25d45..2bf3c8891de 100644 --- a/litellm-rust/crates/traces-clickhouse/query/run_counts.sql +++ b/litellm-rust/crates/traces-clickhouse/query/run_counts.sql @@ -1,5 +1,5 @@ SELECT if({buckets:UInt32} = 0, toUInt32(0), - toUInt32(intDiv((start_ms - {start_ms:Int64}) * {buckets:UInt32}, {end_ms:Int64} - {start_ms:Int64}))) AS bucket, + toUInt32(intDiv((start_ms - {start_ms:Int64} + 1) * {buckets:UInt32} - 1, {end_ms:Int64} - {start_ms:Int64}))) AS bucket, toUInt8({by_failed:UInt8} = 1 AND error_count > 0) AS failed, value, count() AS runs @@ -16,7 +16,7 @@ ARRAY JOIN multiIf( {value:String} = 'service', [service], {value:String} = 'team', [team_id], []) AS value -WHERE ({value:String} = '' OR value != '') AND value ILIKE {contains:String} +WHERE ({value:String} IN ('', 'primary_agent') OR value != '') AND value ILIKE {contains:String} GROUP BY bucket, failed, value ORDER BY runs DESC, bucket, failed, value LIMIT {limit:UInt64} diff --git a/litellm-rust/crates/traces-clickhouse/query/runs_page.sql b/litellm-rust/crates/traces-clickhouse/query/runs_page.sql index cd6e9d46192..135edca3966 100644 --- a/litellm-rust/crates/traces-clickhouse/query/runs_page.sql +++ b/litellm-rust/crates/traces-clickhouse/query/runs_page.sql @@ -1,8 +1,8 @@ page AS ( SELECT * EXCEPT (search_status), - multiIf({sort_key:String} = 'duration_ms', duration_ms, - {sort_key:String} = 'span_count', toInt64(span_count), - {sort_key:String} = 'error_count', toInt64(error_count), + multiIf({sort_key:String} = 'duration_ms', toInt64(least(duration_ns, toUInt64(9223372036854775807))), + {sort_key:String} = 'span_count', toInt64(least(span_count, toUInt64(9223372036854775807))), + {sort_key:String} = 'error_count', toInt64(least(error_count, toUInt64(9223372036854775807))), {sort_key:String} = 'trace_ref', toInt64(0), start_ms) AS sort_value FROM runs diff --git a/litellm-rust/crates/traces-clickhouse/query/span_page.sql b/litellm-rust/crates/traces-clickhouse/query/span_page.sql index d51c278adc0..bd95d2d10c4 100644 --- a/litellm-rust/crates/traces-clickhouse/query/span_page.sql +++ b/litellm-rust/crates/traces-clickhouse/query/span_page.sql @@ -1,5 +1,5 @@ AND EngineReceivedMs <= {as_of_ms:UInt64} -ORDER BY Timestamp, EngineReceivedMs, StatusMessage +ORDER BY Timestamp, EngineReceivedMs, StatusMessage, Duration, StatusCode LIMIT 1 BY TeamId, ApiKeyHash, TraceId, SpanId ) WHERE (team_id, api_key_hash, trace_id, span_id) > ({after_team:String}, {after_key:String}, {after_trace:String}, {after_span:String}) diff --git a/litellm-rust/crates/traces-clickhouse/query/span_text.sql b/litellm-rust/crates/traces-clickhouse/query/span_text.sql index 824c4fbdf91..3b724411d1e 100644 --- a/litellm-rust/crates/traces-clickhouse/query/span_text.sql +++ b/litellm-rust/crates/traces-clickhouse/query/span_text.sql @@ -14,6 +14,6 @@ FROM ( FROM owned_spans WHERE TraceId = {trace_id:String} AND SpanId IN {span_ids:Array(String)} AND hex(SHA256(concat(TeamId, char(0), ApiKeyHash, char(0), TraceId))) = {trace_ref:String} - ORDER BY Timestamp, EngineReceivedMs, StatusMessage +ORDER BY Timestamp, EngineReceivedMs, StatusMessage, Duration, StatusCode LIMIT 1 BY SpanId ) diff --git a/litellm-rust/crates/traces-clickhouse/src/query/named.rs b/litellm-rust/crates/traces-clickhouse/src/query/named.rs index 4b94365f181..acee1f0a965 100644 --- a/litellm-rust/crates/traces-clickhouse/src/query/named.rs +++ b/litellm-rust/crates/traces-clickhouse/src/query/named.rs @@ -159,7 +159,7 @@ struct RunRowEncoding { #[serde(deserialize_with = "super::number::deserialize")] pub start_ms: i64, #[serde(deserialize_with = "super::number::deserialize")] - pub duration_ms: i64, + pub duration_ns: u64, #[serde(deserialize_with = "super::number::deserialize")] pub span_count: u64, #[serde(deserialize_with = "super::number::deserialize")] @@ -667,7 +667,7 @@ mod tests { #[case::unquoted(false)] #[case::quoted(true)] fn rows_decode_clickhouse_numbers(#[case] quoted: bool) { - let run = 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}); + let run = 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_ns": 20_000_000, "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}); assert_eq!(decoded::(run.clone(), quoted), run); let span = json!({"trace_id": "trace", "span_id": "span", "parent_span_id": "parent", "name": "agent", "type": "agent", "wrapper_candidate": 1, "agent": "agent", "framework": "claude-agent-sdk", "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"}); assert_eq!(decoded::(span.clone(), quoted), span); diff --git a/litellm-rust/crates/traces-clickhouse/src/query_access.rs b/litellm-rust/crates/traces-clickhouse/src/query_access.rs index cb6b1d25f7a..51436713867 100644 --- a/litellm-rust/crates/traces-clickhouse/src/query_access.rs +++ b/litellm-rust/crates/traces-clickhouse/src/query_access.rs @@ -77,11 +77,9 @@ impl QueryReaders { .map_err(|_| Error::InvalidScope)?; let user = format!("litellm_traces_{:x}", Sha256::digest(&identity)); let password = credential(secret, b"password", &identity)?; + let cache_key = format!("{user}:{:x}", Sha256::digest(&password)); self.readers - .try_get_with( - user.clone(), - self.provision(client, scope, &user, &password), - ) + .try_get_with(cache_key, self.provision(client, scope, &user, &password)) .await .map_err(Error::Cached) } diff --git a/litellm-rust/crates/traces-clickhouse/src/reads.rs b/litellm-rust/crates/traces-clickhouse/src/reads.rs index e71253ddf9d..257a1ec56ae 100644 --- a/litellm-rust/crates/traces-clickhouse/src/reads.rs +++ b/litellm-rust/crates/traces-clickhouse/src/reads.rs @@ -32,6 +32,9 @@ impl ClickHouseTraces { .await .map_err(|error| match error { StorageError::ResponseTooLarge => StoreError::TooLarge, + StorageError::QueryFailed(failure) if failure.code == Some(396) => { + StoreError::TooLarge + } error => StoreError::Failed(Error::Storage(error)), }) } diff --git a/litellm-rust/crates/traces-clickhouse/templates/query_help.jinja b/litellm-rust/crates/traces-clickhouse/templates/query_help.jinja index a879d3be755..ead9022183b 100644 --- a/litellm-rust/crates/traces-clickhouse/templates/query_help.jinja +++ b/litellm-rust/crates/traces-clickhouse/templates/query_help.jinja @@ -153,6 +153,7 @@ Filter calls by nested metadata {% block time_window -%} Always bound Timestamp or start_time and use LIMIT; add TeamId/ApiKeyHash or team_id/api_key filters when investigating one tenant +Raw SQL returns one bounded response without a cursor. Callers own ORDER BY, LIMIT and keyset predicates for pagination {%- endblock %} {% block reader_limits -%} @@ -164,7 +165,7 @@ LiteLLM provisions SELECT-only readers from the configured ClickHouse connection {%- endblock %} {% block output_format -%} -Do not add FORMAT clauses; the endpoint requires ClickHouse JSON output +The endpoint fixes output to ClickHouse JSON, overriding any SQL FORMAT clause. Trace SQL queries require ClickHouse 26.8 or later {%- endblock %} {% block json_values -%} @@ -197,6 +198,8 @@ Use spend_logs FINAL to collapse replacement rows before totals. Shared response {% block trace_rollups -%} agent_traces_by_key uses SimpleAggregateFunction columns; group by TeamId, ApiKeyHash and TraceId, using min(StartTs), max(EndTs), sum(SpanCount) and groupUniqArrayArray(Models). Do not use Merge combinators +Rollup counts and raw span reads can include repeated exports. Curated trace metrics select one copy per TeamId, ApiKeyHash, TraceId and SpanId, ordered by Timestamp, EngineReceivedMs, StatusMessage, Duration and StatusCode, before filtering status or aggregating +AgentNames contains searchable agent labels, including SpanName for unnamed agents. AgentIdentities contains distinct labels from agent spans, and Frameworks contains observed framework names {%- endblock %} {% block sampling -%} diff --git a/litellm-rust/crates/traces-clickhouse/tests/admin_sql.rs b/litellm-rust/crates/traces-clickhouse/tests/admin_sql.rs index 5c3303e9e8a..8361b3cbdcd 100644 --- a/litellm-rust/crates/traces-clickhouse/tests/admin_sql.rs +++ b/litellm-rust/crates/traces-clickhouse/tests/admin_sql.rs @@ -54,9 +54,12 @@ async fn database() -> Result> { } #[rstest] +#[case::default_format("")] +#[case::explicit_csv(" FORMAT CSV")] #[tokio::test] async fn admin_sql_reads_rows_with_enforced_settings( #[future(awt)] database: Result>, + #[case] format: &str, ) -> Result<(), Box> { let database = database?; let connection = Connection::parse(&format!( @@ -67,7 +70,7 @@ async fn admin_sql_reads_rows_with_enforced_settings( let result = read( &database.client, &connection, - "SELECT n AS answer FROM otel_traces", + &format!("SELECT n AS answer FROM otel_traces{format}"), ) .await?; let json: Value = serde_json::from_str(&result)?; @@ -76,7 +79,7 @@ async fn admin_sql_reads_rows_with_enforced_settings( let result = read( &database.client, &connection, - "SELECT n AS answer FROM agent_traces_by_key", + &format!("SELECT n AS answer FROM agent_traces_by_key{format}"), ) .await?; let json: Value = serde_json::from_str(&result)?; @@ -145,7 +148,7 @@ async fn admin_sql_rejects_errors_after_output_starts( matches!( result, Err(Error::Storage( - litellm_storage_clickhouse::Error::InvalidResponse + litellm_storage_clickhouse::Error::QueryFailed(_) )) ), "expected an error embedded in a successful HTTP response: {result:?}" @@ -175,7 +178,7 @@ async fn admin_sql_enforces_result_row_limit( matches!( result, Err(Error::Storage( - litellm_storage_clickhouse::Error::ResponseTooLarge + litellm_storage_clickhouse::Error::QueryFailed(_) )) ), "{result:?}" diff --git a/litellm-rust/crates/traces-clickhouse/tests/migrations.rs b/litellm-rust/crates/traces-clickhouse/tests/migrations.rs index 6f154ec32f3..8d71dc9303d 100644 --- a/litellm-rust/crates/traces-clickhouse/tests/migrations.rs +++ b/litellm-rust/crates/traces-clickhouse/tests/migrations.rs @@ -517,7 +517,7 @@ async fn listed_agent_names_preserve_scope_and_cursor( .collect::>(); assert_eq!( names["shared"], - serde_json::json!(["research_agent", "reviewer"]) + serde_json::json!(["research_agent", "reviewer", "unnamed"]) ); assert_eq!(names["second"], serde_json::json!(["support_agent"])); let frameworks = [&first["data"][0], &second["data"][0]] diff --git a/litellm-rust/crates/traces-clickhouse/tests/queries.rs b/litellm-rust/crates/traces-clickhouse/tests/queries.rs index 78df4b30d85..758c7d5d54c 100644 --- a/litellm-rust/crates/traces-clickhouse/tests/queries.rs +++ b/litellm-rust/crates/traces-clickhouse/tests/queries.rs @@ -8,7 +8,7 @@ use litellm_traces_cache::{StoreResult, TraceStore}; use litellm_traces_clickhouse::{ClickHouseTraces, Error, QueryScope, query_help, query_sql}; use rstest::{fixture, rstest}; use serde::Deserialize; -use serde_json::Value; +use serde_json::{Value, json}; #[path = "queries/support.rs"] mod fixtures; @@ -123,6 +123,59 @@ async fn documented_queries_render_and_return_expected_rows( Ok(()) } +#[rstest] +#[tokio::test] +async fn documented_failed_spans_filters_status_after_selecting_the_canonical_copy( + #[future(awt)] migrated_database: TestResult, + fixture_clock: TestResult, +) -> TestResult { + let fixture = migrated_database?; + let clock = fixture_clock?; + let writer = litellm_traces_clickhouse::Connection::writer(&fixture.database.url)?; + let rows = [ + ("changed-status", "STATUS_CODE_OK", "", 1), + ("changed-status", "STATUS_CODE_ERROR", "", 2), + ("stable-error", "STATUS_CODE_ERROR", "first", 1), + ("stable-error", "STATUS_CODE_ERROR", "later", 1), + ] + .map(|(span_id, status, message, duration)| { + BTreeMap::from([ + ("Timestamp".into(), json!(clock * 1_000_000_000)), + ("TeamId".into(), json!("team-a")), + ("ApiKeyHash".into(), json!("key-a")), + ("TraceId".into(), json!("duplicates")), + ("SpanId".into(), json!(span_id)), + ("StatusCode".into(), json!(status)), + ("StatusMessage".into(), json!(message)), + ("Duration".into(), json!(duration)), + ]) + }); + litellm_traces_clickhouse::insert_rows( + &fixture.database.client, + &writer, + fixtures::DATABASE, + litellm_traces_clickhouse::InsertTable::OtelTraces, + rows.to_vec(), + ) + .await?; + let reader = fixture + .readers + .connection(&fixture.database.client, &QueryScope::All, "fixture-secret") + .await?; + let sql = include_str!("../query/help/failed_spans.sql") + .replace("now()", &format!("toDateTime({clock})")); + let result: QueryResult = + serde_json::from_str(&query_sql(&fixture.database.client, &reader, &sql).await?)?; + assert_eq!( + result.data, + [json!({ + "team": "team-a", "api_key": "key-a", "trace_id": "duplicates", + "span_id": "stable-error", "message": "first" + })] + ); + Ok(()) +} + #[fixture] fn fixture_clock() -> TestResult { let spans = litellm_traces::decode_otlp( diff --git a/litellm-rust/crates/traces-clickhouse/tests/query_access.rs b/litellm-rust/crates/traces-clickhouse/tests/query_access.rs index ef5b76c2097..362661ec85c 100644 --- a/litellm-rust/crates/traces-clickhouse/tests/query_access.rs +++ b/litellm-rust/crates/traces-clickhouse/tests/query_access.rs @@ -6,6 +6,7 @@ use litellm_traces_clickhouse::{ }; use rstest::{fixture, rstest}; use serde_json::{Value, json}; +use wiremock::{Mock, MockServer, ResponseTemplate, matchers::method}; mod support; use support::{ClickHouseDatabase, database as start_database}; @@ -102,6 +103,45 @@ async fn queries_and_help_are_scoped_by_the_database( Ok(()) } +#[rstest] +#[tokio::test] +async fn reader_cache_reprovisions_after_credential_rotation() +-> Result<(), Box> { + let server = MockServer::start().await; + Mock::given(method("POST")) + .respond_with(ResponseTemplate::new(200)) + .mount(&server) + .await; + let client = Client::no_redirect_for_test(); + let readers = QueryReaders::new(Connection::writer(&server.uri())?, "trace_test".into()); + let old_reader = readers + .connection(&client, &QueryScope::All, "old-master-secret") + .await?; + let provisioned = server.received_requests().await.unwrap().len(); + assert!(provisioned > 0); + + let cached_old = readers + .connection(&client, &QueryScope::All, "old-master-secret") + .await?; + assert!(cached_old.url() == old_reader.url()); + assert_eq!(server.received_requests().await.unwrap().len(), provisioned); + + let new_reader = readers + .connection(&client, &QueryScope::All, "new-master-secret") + .await?; + let rotated = server.received_requests().await.unwrap().len(); + assert!(rotated > provisioned); + assert_eq!(old_reader.url().username(), new_reader.url().username()); + assert!(old_reader.url().password() != new_reader.url().password()); + + let cached_new = readers + .connection(&client, &QueryScope::All, "new-master-secret") + .await?; + assert!(cached_new.url() == new_reader.url()); + assert_eq!(server.received_requests().await.unwrap().len(), rotated); + Ok(()) +} + #[rstest] #[tokio::test] async fn rotating_master_secret_revokes_previous_reader_credentials( @@ -125,8 +165,8 @@ async fn rotating_master_secret_revokes_previous_reader_credentials( let old_rows: Value = serde_json::from_str(&old_result)?; assert_eq!(old_rows["data"], json!([{ "id": "a1" }, { "id": "a2" }])); - let rotated_readers = QueryReaders::new(database.writer.clone(), "trace_test".into()); - let new_reader = rotated_readers + let new_reader = database + .readers .connection(&database.client, &scope, "new-master-secret") .await?; assert!( diff --git a/litellm-rust/crates/traces-clickhouse/tests/reads.rs b/litellm-rust/crates/traces-clickhouse/tests/reads.rs index a1a73637564..9866e8bc096 100644 --- a/litellm-rust/crates/traces-clickhouse/tests/reads.rs +++ b/litellm-rust/crates/traces-clickhouse/tests/reads.rs @@ -110,7 +110,6 @@ async fn list_costs_match_each_run_when_response_ids_are_reused( &PageRequest { cursor: None, limit: 50, - ..Default::default() }, ) .await?; @@ -248,7 +247,6 @@ async fn large_runs_remain_complete_under_default_reader_limits( &PageRequest { cursor: None, limit: 500, - ..Default::default() }, ) .await?; @@ -410,7 +408,6 @@ async fn cursor_pages_keep_a_tenant_scoped_snapshot_when_more_spans_arrive( &PageRequest { cursor: None, limit: 10, - ..Default::default() }, ) .await?; @@ -578,7 +575,6 @@ async fn an_oversized_span_keeps_the_run_list_available_with_partial_totals( &PageRequest { cursor: None, limit: 50, - ..Default::default() }, ) .await?; @@ -617,7 +613,6 @@ async fn an_oversized_span_keeps_the_run_list_available_with_partial_totals( &PageRequest { cursor: None, limit: 50, - ..Default::default() }, ) .await?; @@ -632,7 +627,6 @@ async fn an_oversized_span_keeps_the_run_list_available_with_partial_totals( &PageRequest { cursor: None, limit: 50, - ..Default::default() }, ) .await?; @@ -760,7 +754,6 @@ async fn gateway_ids_resolve_through_detail_and_batch_reads_with_legacy_fallback &PageRequest { cursor: None, limit: 50, - ..Default::default() }, ) .await?; @@ -827,7 +820,6 @@ async fn a_run_shared_with_another_user_stays_hidden_before_its_rollup_rows_merg let page = PageRequest { cursor: None, limit: 50, - ..Default::default() }; let team = QueryScope::Owned { user_id: String::new(), diff --git a/litellm-rust/crates/traces-clickhouse/tests/search.rs b/litellm-rust/crates/traces-clickhouse/tests/search.rs index 3d9b228c7a3..0d24f2e3534 100644 --- a/litellm-rust/crates/traces-clickhouse/tests/search.rs +++ b/litellm-rust/crates/traces-clickhouse/tests/search.rs @@ -9,7 +9,7 @@ use litellm_traces::{ }; use litellm_traces_cache::{PageRequest, TraceReader, TraceStore}; use litellm_traces_clickhouse::{ - ClickHouseTraces, Connection, InsertTable, QueryScope, insert_rows, + ClickHouseTraces, Connection, InsertTable, QueryScope, encode_rows, insert_rows, }; use rstest::rstest; use serde_json::{Value, json}; @@ -38,11 +38,7 @@ fn filter(start_ms: i64, q: &str) -> RunFilter { } fn page(cursor: Option, limit: u32) -> PageRequest { - PageRequest { - cursor, - limit, - ..Default::default() - } + PageRequest { cursor, limit } } fn reader() -> TraceReader { @@ -294,6 +290,56 @@ async fn list_q_pages_through_matches_only( Ok(()) } +#[rstest] +#[tokio::test] +async fn named_and_unnamed_agents_in_one_run_remain_searchable( + #[future(awt)] migrated_database: TestResult, +) -> TestResult { + let fixture = migrated_database?; + let store = seed(&fixture).await?; + let writer = Connection::writer(&fixture.database.url)?; + insert_rows( + &fixture.database.client, + &writer, + DATABASE, + InsertTable::OtelTraces, + vec![span_row( + &runs()[0], + 5, + &step("alpha-unnamed", "helper", "agent"), + "alpha-root", + )], + ) + .await?; + let reader = reader(); + let values = reader + .values( + &store, + &team_a(), + &filter(T0_MS, "trace_id:alpha"), + RunField::Agent, + "", + 10, + ) + .await?; + assert_eq!(values.values, ["helper", "researcher"]); + for agent in values.values { + let listed = reader + .list_traces( + &store, + &team_a(), + &filter(T0_MS, &format!("agent:{agent}")), + RunOrder::NEWEST, + &page(None, 10), + ) + .await?; + assert_eq!(listed.data.len(), 1); + assert_eq!(listed.data[0].trace_id, "alpha"); + assert!(listed.data[0].agent_names.contains(&agent)); + } + Ok(()) +} + #[rstest] #[case::all("", [ (1, 0, vec![("researcher", 1)]), @@ -337,6 +383,123 @@ async fn histogram_counts_matching_runs_per_bucket( Ok(()) } +#[rstest] +#[case::successful(false)] +#[case::failed(true)] +#[tokio::test] +async fn histogram_counts_runs_without_an_agent_or_service( + #[future(awt)] migrated_database: TestResult, + #[case] failed: bool, +) -> TestResult { + let fixture = migrated_database?; + let store = seed(&fixture).await?; + let writer = Connection::writer(&fixture.database.url)?; + insert_rows( + &fixture.database.client, + &writer, + DATABASE, + InsertTable::OtelTraces, + vec![BTreeMap::from([ + ("Timestamp".into(), json!(T0_MS * 1_000_000)), + ("TraceId".into(), json!("unlabelled")), + ("SpanId".into(), json!("root")), + ("ParentSpanId".into(), json!("")), + ("SpanName".into(), json!("run")), + ("ServiceName".into(), json!("")), + ("ObservationType".into(), json!("chain")), + ("TeamId".into(), json!("team-a")), + ("ApiKeyHash".into(), json!("key-a")), + ( + "StatusCode".into(), + json!(if failed { + "STATUS_CODE_ERROR" + } else { + "STATUS_CODE_OK" + }), + ), + ])], + ) + .await?; + let reader = reader(); + let filter = filter(T0_MS, "trace_id:unlabelled"); + let histogram = reader.histogram(&store, &team_a(), &filter, 3).await?; + let count = reader.count_traces(&store, &team_a(), &filter).await?; + let bucket = &histogram.buckets[0]; + assert_eq!(count, 1); + assert_eq!(bucket.total, count); + assert_eq!(bucket.failed, u64::from(failed)); + assert_eq!( + bucket.total, + bucket.failed + bucket.agents.iter().map(|agent| agent.runs).sum::() + ); + Ok(()) +} + +#[rstest] +#[case::agent(CountValue::PrimaryAgent)] +#[case::attribute(CountValue::Attribute("bucket.tag".into()))] +#[tokio::test] +async fn histogram_counts_use_the_displayed_bucket_boundaries( + #[future(awt)] migrated_database: TestResult, + #[case] value: CountValue, +) -> TestResult { + let fixture = migrated_database?; + let store = seed(&fixture).await?; + let writer = Connection::writer(&fixture.database.url)?; + insert_rows( + &fixture.database.client, + &writer, + DATABASE, + InsertTable::OtelTraces, + (0..10) + .map(|offset| { + BTreeMap::from([ + ("Timestamp".into(), json!((T0_MS + offset) * 1_000_000)), + ("TraceId".into(), json!(format!("bucket-{offset}"))), + ("SpanId".into(), json!("root")), + ("ParentSpanId".into(), json!("")), + ("SpanName".into(), json!("run")), + ("ServiceName".into(), json!("svc")), + ("ObservationType".into(), json!("chain")), + ("TeamId".into(), json!("team-a")), + ("ApiKeyHash".into(), json!("key-a")), + ("SpanAttributes".into(), json!({"bucket.tag": "tag"})), + ]) + }) + .collect(), + ) + .await?; + let filter = RunFilter { + start_ms: T0_MS, + end_ms: T0_MS + 10, + search: RunSearch::parse("trace_id:bucket*"), + ..Default::default() + }; + let counts = store + .run_counts( + &team_a(), + &RunCountQuery { + filter: filter.clone(), + by: CountBy { + buckets: Some(3), + value: Some(value), + ..Default::default() + }, + contains: String::new(), + limit: None, + }, + ) + .await?; + let histogram = litellm_traces::search::histogram(&counts, filter.start_ms, filter.end_ms, 3); + for bucket in histogram.buckets { + let expected = (T0_MS..T0_MS + 10) + .filter(|start| (bucket.start_ms..bucket.end_ms).contains(start)) + .count() as u64; + assert_eq!(bucket.total, expected); + } + Ok(()) +} + #[rstest] #[case::agents_in_scope(RunField::Agent, "", &["researcher", "writer"])] #[case::names_by_frequency(RunField::Name, "", &["plan trip", "write report"])] @@ -551,25 +714,50 @@ async fn listed_runs_name_the_agent_they_matched_even_when_resolution_is_limited #[rstest] #[tokio::test] -async fn listing_runs_scans_the_rollup_once( +async fn listing_runs_skips_out_of_window_span_rows( #[future(awt)] migrated_database: TestResult, ) -> TestResult { let fixture = migrated_database?; let store = seed(&fixture).await?; + let writer = Connection::writer(&fixture.database.url)?; + insert_rows( + &fixture.database.client, + &writer, + DATABASE, + InsertTable::OtelTraces, + (0..5000) + .map(|index| { + BTreeMap::from([ + ( + "Timestamp".into(), + json!((T0_MS - 24 * HOUR_MS) * 1_000_000), + ), + ("TraceId".into(), json!("historical")), + ("SpanId".into(), json!(format!("historical-{index}"))), + ("ParentSpanId".into(), json!("")), + ("SpanName".into(), json!("run")), + ("ObservationType".into(), json!("chain")), + ("TeamId".into(), json!("team-a")), + ("ApiKeyHash".into(), json!("key-a")), + ]) + }) + .collect(), + ) + .await?; let query = RunQuery { - selection: RunSelection::Matching(filter(0, "")), + selection: RunSelection::Matching(filter(T0_MS, "")), order: RunOrder::NEWEST, after: None, limit: 50, }; let listed = store.runs(&QueryScope::All, &query).await?; assert_eq!(listed.len(), 4); - let budget = table_rows(&fixture, "agent_traces_by_key").await? - + table_rows(&fixture, "otel_traces").await?; + let budget = 2 * table_rows(&fixture, "agent_traces_by_key").await? + + runs().iter().flat_map(rows).count() as u64; let read = rows_read_by(&fixture, "FROM owned_runs").await?; assert!( read <= budget, - "listing 4 runs read {read} rows, more than the {budget} rollup and span rows that exist" + "listing current runs read {read} rows, exceeding their {budget} rollup and in-window span rows" ); Ok(()) } @@ -701,6 +889,9 @@ async fn add_attributes(fixture: &SeededDatabase) -> TestResult { #[case::span_attribute("attr.tenant.tier:gold", &["alpha"])] #[case::resource_attribute("attr.tenant.tier:SILV*", &["beta"])] #[case::excluded_attribute("-attr.tenant.tier:gold", &["gamma", "beta"])] +#[case::any_attribute("attr.tenant.tier:*", &["beta", "alpha"])] +#[case::missing_attribute("-attr.tenant.tier:*", &["gamma"])] +#[case::missing_wildcard("attr.missing:*", &[])] #[case::two_attributes("attr.tenant.tier:gold attr.tenant.tier:silver", &[])] #[case::attribute_and_field("attr.tenant.tier:* status:error", &["beta"])] #[case::unknown_attribute("attr.missing:gold", &[])] @@ -727,6 +918,176 @@ async fn service_team_and_attribute_filters_select_runs( Ok(()) } +#[rstest] +#[case::duration_ascending(RunSortKey::DurationMs, false)] +#[case::duration_descending(RunSortKey::DurationMs, true)] +#[case::span_count_ascending(RunSortKey::SpanCount, false)] +#[case::span_count_descending(RunSortKey::SpanCount, true)] +#[case::error_count_ascending(RunSortKey::ErrorCount, false)] +#[case::error_count_descending(RunSortKey::ErrorCount, true)] +#[case::start_ascending(RunSortKey::StartMs, false)] +#[case::start_descending(RunSortKey::StartMs, true)] +#[case::reference_ascending(RunSortKey::TraceRef, false)] +#[case::reference_descending(RunSortKey::TraceRef, true)] +#[tokio::test] +async fn metric_sorting_pages_by_the_deduplicated_displayed_values( + #[future(awt)] migrated_database: TestResult, + #[case] key: RunSortKey, + #[case] descending: bool, +) -> TestResult { + let fixture = migrated_database?; + let store = seed(&fixture).await?; + let writer = Connection::writer(&fixture.database.url)?; + let span = |trace: &str, id: &str, duration: u64, failed: bool| { + BTreeMap::from([ + ("Timestamp".into(), json!(T0_MS * 1_000_000)), + ("TraceId".into(), json!(trace)), + ("SpanId".into(), json!(id)), + ( + "ParentSpanId".into(), + json!(if id == "root" { "" } else { "root" }), + ), + ("SpanName".into(), json!(id)), + ("ServiceName".into(), json!("svc")), + ( + "ObservationType".into(), + json!(if id == "root" { "chain" } else { "tool" }), + ), + ("TeamId".into(), json!("team-a")), + ("ApiKeyHash".into(), json!("key-a")), + ("Duration".into(), json!(duration)), + ("EngineReceivedMs".into(), json!(1)), + ( + "StatusCode".into(), + json!(if failed { + "STATUS_CODE_ERROR" + } else { + "STATUS_CODE_OK" + }), + ), + ]) + }; + let duplicate = |changes: [(&str, serde_json::Value); 1]| { + BTreeMap::from_iter( + span("metric-a", "root", 50_000_000, true) + .into_iter() + .chain(changes.into_iter().map(|(key, value)| (key.into(), value))), + ) + }; + fixture + .database + .client + .post(writer.url().clone()) + .query(&[( + "query", + format!("INSERT INTO {DATABASE}.otel_traces FORMAT JSONEachRow"), + )]) + .body(encode_rows(vec![ + span("metric-a", "root", 400_000, false), + span("metric-b", "root", 900_000, false), + span("metric-b", "child", 200_000, false), + span("metric-c", "root", 1_100_000, false), + span("metric-c", "child-one", 200_000, true), + span("metric-c", "child-two", 300_000, true), + duplicate([("EngineReceivedMs", json!(2))]), + duplicate([("Timestamp", json!((T0_MS + 1) * 1_000_000))]), + duplicate([("StatusMessage", json!("later"))]), + span("metric-b", "root", 800_000, false), + span("metric-b", "root", 800_000, true), + ])?) + .send() + .await? + .error_for_status()?; + let reader = reader(); + let order = RunOrder { key, descending }; + let filter = filter(T0_MS, "trace_id:metric*"); + let mut listed = Vec::new(); + let mut cursor = None; + loop { + let result = reader + .list_traces(&store, &team_a(), &filter, order, &page(cursor.take(), 1)) + .await?; + listed.extend(result.data); + let Some(next) = result.next_cursor else { + break; + }; + cursor = Some(next); + } + let expected = if descending { + ["metric-c", "metric-b", "metric-a"] + } else { + ["metric-a", "metric-b", "metric-c"] + }; + if matches!(key, RunSortKey::StartMs | RunSortKey::TraceRef) { + assert_eq!( + listed + .iter() + .map(|run| run.trace_id.as_str()) + .collect::>(), + expected.into_iter().collect() + ); + } else { + assert_eq!( + listed + .iter() + .map(|run| run.trace_id.as_str()) + .collect::>(), + expected + ); + } + for run in listed { + let (duration_ns, span_count, error_count) = match run.trace_id.as_str() { + "metric-a" => (400_000, 1, 0), + "metric-b" => (800_000, 2, 1), + "metric-c" => (1_100_000, 3, 2), + other => return Err(format!("unexpected run {other}").into()), + }; + assert_eq!(run.duration_ms, f64::from(duration_ns) / 1_000_000.0); + assert_eq!(run.span_count, span_count); + assert_eq!(run.error_count, error_count); + } + let errors = RunFilter { + search: RunSearch::parse("trace_id:metric* status:error"), + ..filter.clone() + }; + let filtered = reader + .list_traces(&store, &team_a(), &errors, order, &page(None, 10)) + .await?; + assert_eq!( + filtered + .data + .iter() + .map(|run| run.trace_id.as_str()) + .collect::>(), + BTreeSet::from(["metric-b", "metric-c"]) + ); + let count = reader.count_traces(&store, &team_a(), &errors).await?; + let histogram = reader.histogram(&store, &team_a(), &errors, 3).await?; + assert_eq!(count, filtered.data.len() as u64); + assert_eq!( + histogram + .buckets + .iter() + .map(|bucket| bucket.total) + .sum::(), + count + ); + assert_eq!( + histogram + .buckets + .iter() + .map(|bucket| bucket.failed) + .sum::(), + count + ); + let all = reader.histogram(&store, &team_a(), &filter, 3).await?; + assert_eq!( + all.buckets.iter().map(|bucket| bucket.failed).sum::(), + count + ); + Ok(()) +} + #[rstest] #[case::newest(RunOrder::NEWEST)] #[case::oldest(RunOrder { descending: false, ..RunOrder::NEWEST })] diff --git a/litellm-rust/crates/traces/src/resolve/view.rs b/litellm-rust/crates/traces/src/resolve/view.rs index 2b3a49d4e12..07f97b60c49 100644 --- a/litellm-rust/crates/traces/src/resolve/view.rs +++ b/litellm-rust/crates/traces/src/resolve/view.rs @@ -224,7 +224,7 @@ pub fn listed_summary(row: &RunRow) -> TraceSummary { frameworks: row.frameworks.clone(), input_preview: row.input_preview.clone(), start_time: iso_time(row.start_ms), - duration_ms: row.duration_ms as f64, + duration_ms: row.duration_ns as f64 / NANOS_PER_MS, status: row.status, span_count: row.span_count, agent_count: row.agent_count, diff --git a/litellm-rust/crates/traces/src/search.rs b/litellm-rust/crates/traces/src/search.rs index 7be72d3c689..b6e437e7289 100644 --- a/litellm-rust/crates/traces/src/search.rs +++ b/litellm-rust/crates/traces/src/search.rs @@ -30,7 +30,7 @@ pub enum RunField { /// What a `key:value` filter matches: a run field, or `attr.`, a span or resource /// attribute that any span of the run carries. -#[derive(Clone, Debug, Eq, PartialEq)] +#[derive(Clone, Debug, Eq, PartialEq, Serialize)] pub enum SearchKey { Field(RunField), Attribute(String), @@ -48,7 +48,7 @@ impl SearchKey { } } -#[derive(Clone, Debug, Default, Eq, PartialEq)] +#[derive(Clone, Debug, Default, Eq, PartialEq, Serialize)] pub struct RunFilter { pub start_ms: i64, pub end_ms: i64, @@ -57,7 +57,7 @@ pub struct RunFilter { pub trace_refs: Vec, } -#[derive(Clone, Debug, Eq, PartialEq)] +#[derive(Clone, Debug, Eq, PartialEq, Serialize)] pub struct FieldFilter { pub key: SearchKey, /// Matched against the whole value, ignoring case; `*` matches any run of characters. @@ -66,7 +66,7 @@ pub struct FieldFilter { } /// The parsed `q` of the runs list. Every text term and every filter must hold. -#[derive(Clone, Debug, Default, Eq, PartialEq)] +#[derive(Clone, Debug, Default, Eq, PartialEq, Serialize)] pub struct RunSearch { /// Each must appear in the trace id, input or name, ignoring case. pub text: Vec, diff --git a/litellm-rust/crates/traces/src/store.rs b/litellm-rust/crates/traces/src/store.rs index 34d7ec75b6e..b26eb57edb8 100644 --- a/litellm-rust/crates/traces/src/store.rs +++ b/litellm-rust/crates/traces/src/store.rs @@ -56,7 +56,7 @@ impl RunOrder { let count = |count: u64| i64::try_from(count).unwrap_or(i64::MAX); match self.key { RunSortKey::StartMs => row.start_ms, - RunSortKey::DurationMs => row.duration_ms, + RunSortKey::DurationMs => count(row.duration_ns), RunSortKey::SpanCount => count(row.span_count), RunSortKey::ErrorCount => count(row.error_count), RunSortKey::TraceRef => 0, @@ -108,7 +108,7 @@ pub struct RunRow { #[serde(serialize_with = "crate::wire::serialize_status")] pub status: crate::SpanStatus, pub start_ms: i64, - pub duration_ms: i64, + pub duration_ns: u64, pub span_count: u64, pub agent_count: u64, pub agent_invocations: u64, @@ -139,15 +139,14 @@ pub enum CountValue { /// Each dimension left unset collapses to one group: bucket 0, not failed, or an empty value. #[derive(Clone, Debug, Default, Eq, PartialEq)] pub struct CountBy { - /// Equal-width slices of the filter window; run `i` lands in - /// `(start_ms - window.start) * buckets / window.len()`. + /// Equal-width slices of the filter window, using the histogram's integer boundaries. pub buckets: Option, pub failed: bool, pub value: Option, } /// Matching runs per group, most runs first. A run with several values for a field, such as -/// several models, counts once under each; empty values are not counted. +/// several models, counts once under each; totals and primary-agent groups retain empty values. #[derive(Clone, Debug, PartialEq)] pub struct RunCountQuery { pub filter: RunFilter, diff --git a/litellm-rust/crates/traces/tests/resolve.rs b/litellm-rust/crates/traces/tests/resolve.rs index 2c62dce27a1..6e001856c4e 100644 --- a/litellm-rust/crates/traces/tests/resolve.rs +++ b/litellm-rust/crates/traces/tests/resolve.rs @@ -779,7 +779,9 @@ fn gateway_id_miss_only_vetoes_rows_that_carry_a_call_id( } #[rstest] -fn listed_summary_keeps_rollup_counts_with_unknown_cost() { +#[case::whole_milliseconds(51_385_000_000)] +#[case::fractional_milliseconds(51_385_123_456)] +fn listed_summary_keeps_rollup_counts_with_unknown_cost(#[case] duration_ns: u64) { let summary = listed_summary(&RunRow { trace_id: "t1".into(), trace_ref: "ref".into(), @@ -791,7 +793,7 @@ fn listed_summary_keeps_rollup_counts_with_unknown_cost() { input_preview: "hi".into(), status: SpanStatus::Ok, start_ms: 1_790_742_989_377, - duration_ms: 51_385, + duration_ns, span_count: 126, agent_count: 2, agent_invocations: 0, @@ -805,6 +807,7 @@ fn listed_summary_keeps_rollup_counts_with_unknown_cost() { error_count: 1, }); assert_eq!(summary.spend, None); + assert_eq!(summary.duration_ms, duration_ns as f64 / 1_000_000.0); assert_eq!(summary.status, SpanStatus::Ok); assert_eq!( ( diff --git a/litellm/proxy/lens/models.py b/litellm/proxy/lens/models.py index 85fa7a4e039..926227509f9 100644 --- a/litellm/proxy/lens/models.py +++ b/litellm/proxy/lens/models.py @@ -1,9 +1,7 @@ -import base64 from datetime import datetime, timedelta, timezone from typing import Annotated, Final, Literal, TypeAlias -from pydantic import AfterValidator, BaseModel, ConfigDict, Field, TypeAdapter, ValidationError, model_validator -from typing_extensions import ReadOnly, TypedDict +from pydantic import AfterValidator, BaseModel, ConfigDict, Field, model_validator from litellm.rust_bridge.trace.generated.types import TraceSummary @@ -53,56 +51,6 @@ def parse_execution(value: str) -> tuple[str, str]: return trace_ref, trace_id -_LEGACY_ID: Final[TypeAdapter[tuple[str, str, str] | tuple[str, str, str, str]]] = TypeAdapter( - tuple[str, str, str] | tuple[str, str, str, str] -) - - -def _legacy_execution_id(value: str) -> str | None: - """Selections saved before executions were runs named them by source, team, trace id and reference.""" - try: - parts: Final = _LEGACY_ID.validate_json(base64.urlsafe_b64decode(value)) - except (ValueError, ValidationError): - return value - trace_ref: Final = parts[3] if len(parts) == 4 else "" - return execution_id(trace_ref, parts[2]) if parts[0] == "traces" and trace_ref else None - - -def _term(key: str, value: str) -> str: - return f'{key}:"{value}"' if any(c.isspace() for c in value) else f"{key}:{value}" - - -class _LegacyFilter(BaseModel): - key: str - value: str - - -class _LegacySelection(BaseModel): - model_config = ConfigDict(extra="allow") - q: str = "" - source: str = "" - service: str = "" - agent_name: str = "" - filters: tuple[_LegacyFilter, ...] = () - team_id: str = "" - execution_ids: tuple[str, ...] = () - - def current(self) -> dict[str, object]: - terms: Final = ( - self.q, - _term("agent", self.agent_name) if self.agent_name else "", - _term("service", self.service) if self.service else "", - _term("team", self.team_id) if self.team_id else "", - *(_term(f"attr.{f.key}", f.value) for f in self.filters), - ) - ids: Final = tuple(i for i in map(_legacy_execution_id, self.execution_ids) if i is not None) - return {**(self.model_extra or {}), "q": " ".join(t for t in terms if t), "execution_ids": ids} - - -_LEGACY_KEYS: Final = frozenset({"source", "service", "agent_name", "filters", "team_id"}) -_FIELDS: Final = TypeAdapter(dict[str, object]) - - class Check(Record): id: str = Field(min_length=1) instruction: str = Field(min_length=3) @@ -115,18 +63,6 @@ class ActivitySelection(Record): sample_percent: float = Field(default=100, gt=0, le=100, allow_inf_nan=False) execution_ids: tuple[str, ...] = () - @model_validator(mode="before") - @classmethod - def from_saved_filters(cls, data: object) -> object: - """Selections saved with separate agent, service, team and attribute filters load as `q`.""" - try: - fields: Final = _FIELDS.validate_python(data) - except ValidationError: - return data - if not _LEGACY_KEYS & fields.keys(): - return data - return _LegacySelection.model_validate(fields).current() - class LensSettings(ActivitySelection): name: str = Field(min_length=1) @@ -254,17 +190,6 @@ class ActivityAvailability(Record): traces: bool = False -class _SavedSample(TypedDict, total=False): - executions: ReadOnly[list[dict[str, object]]] - - -class _SavedJob(TypedDict, total=False): - sample: ReadOnly[_SavedSample | None] - - -_SAVED_JOB: Final = TypeAdapter(_SavedJob) - - class RunAssessment(Record): execution_id: str issue_checks: tuple[str, ...] = () @@ -287,20 +212,6 @@ class Step(Record): class Job(Record): - @model_validator(mode="before") - @classmethod - def without_legacy_sample(cls, data: object) -> object: - """Samples saved before executions were runs no longer resolve, so they load as absent.""" - try: - saved: Final = _SAVED_JOB.validate_python(data) - fields: Final = _FIELDS.validate_python(data) - except ValidationError: - return data - sample: Final = saved.get("sample") or {} - if not any("source" in e for e in sample.get("executions", [])): - return data - return {**fields, "sample": None} - id: str status: Literal["queued", "running", "completed", "failed", "cancelled"] = "queued" stage: str = "Queued" diff --git a/litellm/proxy/lens/sources.py b/litellm/proxy/lens/sources.py index 2cb535714c0..11afe38cee7 100644 --- a/litellm/proxy/lens/sources.py +++ b/litellm/proxy/lens/sources.py @@ -188,21 +188,23 @@ class SourceReader: ) -> Mapping[tuple[str, SpanPart], SpanText]: if not span_ids: return {} - reads: Final[list[tuple[SpanPart, tuple[SpanText, ...]]]] = [ - ( - part, - await self.storage.span_text( - execution.trace_id, - execution.trace_ref, - span_ids, + reads: Final[tuple[tuple[SpanPart, tuple[SpanText, ...]], ...]] = tuple( + [ + ( part, - access, - max_chars=max_chars if not tail else budget - budget // 3, - tail=tail, - ), - ) - for part, _, budget in PARTS - ] + await self.storage.span_text( + execution.trace_id, + execution.trace_ref, + span_ids, + part, + access, + max_chars=max_chars if not tail else budget - budget // 3, + tail=tail, + ), + ) + for part, _, budget in PARTS + ] + ) return { (text["span_id"], part): text for part, texts in reads for text in texts } # comprehension-ok: flatten one read per part @@ -222,12 +224,7 @@ class SourceReader: heads: Final = await self._texts(access, execution, page_ids, BUDGET) long: Final = tuple(span["span_id"] for span in page if offset == 0 and _total(_pieces(span, heads)) > BUDGET) tails: Final = await self._texts(access, execution, long, 0, tail=True) - parts: Final = tuple( - [ - await self._part(access, execution, span, heads, tails, offset) - for span in page # comprehension-ok: sequential reads keep storage load bounded - ] - ) + parts: Final = tuple([await self._part(access, execution, span, heads, tails, offset) for span in page]) root_seen: Final = any(span.get("parent_span_id") is None for span in spans) return ExecutionContent( execution=execution, diff --git a/litellm/proxy/tracing_endpoints.py b/litellm/proxy/tracing_endpoints.py index e6392c9b2c9..2dc768f4d2f 100644 --- a/litellm/proxy/tracing_endpoints.py +++ b/litellm/proxy/tracing_endpoints.py @@ -29,7 +29,7 @@ from litellm.proxy.auth.authorization_dependencies import LogTeamLookupDependenc from litellm.proxy.auth.user_api_key_auth import user_api_key_auth from litellm.proxy.common_utils.http_parsing_utils import is_otlp_trace_request from litellm.proxy.tracing_runtime import provide_receiver, require_receiver -from litellm.rust_bridge.trace.errors import TraceChanged +from litellm.rust_bridge.trace.errors import TraceChanged, TraceQueryError from litellm.rust_bridge.trace.generated.models import TraceQueryHelp from litellm.rust_bridge.trace.generated.types import ( AllQueryScope, @@ -182,7 +182,8 @@ RunQuery = Annotated[ Query( max_length=1000, description='Free text and key:value filters, e.g. `agent:research* -status:ok "book a flight"`. ' - "Keys: name, agent, status, model, input, trace_id. `*` globs and a leading `-` negates", + "Keys: name, agent, status, model, input, trace_id, service, team and attr.. " + "`*` globs and a leading `-` negates", ), ] @@ -204,7 +205,8 @@ def trace_window(start_ms: StartMs = None, end_ms: EndMs = None) -> TraceWindow: @router.get("/v1/traces", response_model=TracePage) async def list_agent_traces( context: Annotated[TraceAccessContext, Depends(provide_trace_access)], - window: Annotated[TraceWindow, Depends(trace_window)], + start_ms: StartMs = None, + end_ms: EndMs = None, q: RunQuery = "", cursor: Annotated[str | None, Query(max_length=512)] = None, sort_by: Literal["start_ms", "duration_ms", "span_count", "error_count"] = "start_ms", @@ -213,9 +215,7 @@ async def list_agent_traces( order: Final = RunOrder(key=sort_by, descending=sort_dir == "desc") try: tracing, scope = context.reader() - return await tracing.list_traces( - scope=scope, start_ms=window.start_ms, end_ms=window.end_ms, q=q, cursor=cursor, order=order - ) + return await tracing.list_traces(scope=scope, start_ms=start_ms, end_ms=end_ms, q=q, cursor=cursor, order=order) except (TraceChanged, ValueError, OverflowError, RuntimeError) as error: raise read_failure(error) from error @@ -293,6 +293,16 @@ async def provide_trace_query_access( return TraceQueryAccess(storage, scope, secret) +def sql_failure_status(error: TraceQueryError) -> tuple[int, str]: + match error.kind: + case "rejected": + return 400, "query_rejected" + case "limited": + return 422, "query_limit_exceeded" + case "unavailable": + return 503, "query_unavailable" + + @router.post("/v1/traces/query", response_model=TraceSQLResponse, response_model_exclude_unset=True) async def query_agent_traces( body: TraceQueryRequest, @@ -300,11 +310,27 @@ async def query_agent_traces( ) -> TraceSQLResponse: try: return await access.storage.query_sql(body.sql, read_access(access.scope), access.secret) + except TraceQueryError as error: + status, code = sql_failure_status(error) + raise HTTPException( + status_code=status, + detail={"code": code, "database_code": error.database_code, "message": error.message}, + ) from error except ValueError as error: - raise HTTPException(status_code=400, detail=str(error)) from error + raise HTTPException( + status_code=400, + detail={"code": "query_rejected", "database_code": None, "message": str(error)}, + ) from error except RuntimeError as error: verbose_proxy_logger.warning("Trace SQL query unavailable: %s", error) - raise HTTPException(status_code=503, detail="Trace SQL query failed or exceeded reader limits") from error + raise HTTPException( + status_code=503, + detail={ + "code": "query_unavailable", + "database_code": None, + "message": "Trace SQL is temporarily unavailable", + }, + ) from error @router.get("/v1/traces/query/help", response_model=TraceQueryHelp, response_model_exclude_unset=True) diff --git a/litellm/rust_bridge/_native.pyi b/litellm/rust_bridge/_native.pyi index 4f9e3b8f115..e74f9fcb25d 100644 --- a/litellm/rust_bridge/_native.pyi +++ b/litellm/rust_bridge/_native.pyi @@ -11,7 +11,7 @@ from litellm.rust_bridge.embeddings.entrypoints import LiteLLMEmbeddingRequest from litellm.rust_bridge.messages.entrypoints import LiteLLMMessagesRequest from litellm.rust_bridge.ocr.entrypoints import LiteLLMOcrRequest from litellm.rust_bridge.responses.entrypoints import LiteLLMResponsesRequest -from litellm.rust_bridge.trace.generated.types import QueryScope +from litellm.rust_bridge.trace.generated.types import QueryScope, RunOrder from litellm.types.llms.anthropic_messages.anthropic_response import AnthropicMessagesResponse from litellm.types.llms.openai import ResponsesAPIResponse from litellm.types.utils import EmbeddingResponse, ModelResponse @@ -45,12 +45,12 @@ class NativeTraceStorage: def list_traces( self, scope: QueryScope, - start_ms: int, - end_ms: int, + start_ms: int | None, + end_ms: int | None, q: str, cursor: str | None, limit: int, - order: str = "newest", + order: RunOrder, trace_refs: Sequence[str] = (), ) -> Future[JsonValue]: ... def count_traces( diff --git a/litellm/rust_bridge/trace/errors.py b/litellm/rust_bridge/trace/errors.py index 3c84d233be4..5e08c62820b 100644 --- a/litellm/rust_bridge/trace/errors.py +++ b/litellm/rust_bridge/trace/errors.py @@ -1,5 +1,21 @@ +from typing import Final, Literal + + class TraceChanged(Exception): """The paging snapshot no longer matches the stored trace, so the client must start a new traversal. Raised by the Rust trace reader when a cursor's snapshot version differs from the graph it rebuilt. """ + + +class TraceQueryError(Exception): + def __init__( + self, + kind: Literal["rejected", "limited", "unavailable"], + database_code: int | None, + message: str, + ) -> None: + self.kind: Final = kind + self.database_code: Final = database_code + self.message: Final = message + super().__init__(message) diff --git a/litellm/rust_bridge/trace/storage.py b/litellm/rust_bridge/trace/storage.py index dfe6c377ad6..8071d2399b4 100644 --- a/litellm/rust_bridge/trace/storage.py +++ b/litellm/rust_bridge/trace/storage.py @@ -52,8 +52,8 @@ class NativeStore(Protocol): def list_traces( self, scope: QueryScope, - start_ms: int, - end_ms: int, + start_ms: int | None, + end_ms: int | None, q: str, cursor: str | None, limit: int, @@ -201,8 +201,8 @@ class ClickHouseStorage: async def list_traces( self, scope: QueryScope, - start_ms: int, - end_ms: int, + start_ms: int | None, + end_ms: int | None, q: str = "", cursor: str | None = None, limit: int = AGENT_TRACING_LIST_PAGE_SIZE, diff --git a/litellm/tracing/receiver.py b/litellm/tracing/receiver.py index 1957e1d8f87..99c9e84670e 100644 --- a/litellm/tracing/receiver.py +++ b/litellm/tracing/receiver.py @@ -109,8 +109,8 @@ class TraceReceiver: async def list_traces( self, scope: QueryScope, - start_ms: int, - end_ms: int, + start_ms: int | None, + end_ms: int | None, q: str = "", cursor: str | None = None, order: RunOrder = NEWEST, diff --git a/tests/test_litellm_rust/test_traces.py b/tests/test_litellm_rust/test_traces.py index b8af0f0c374..822f11a893d 100644 --- a/tests/test_litellm_rust/test_traces.py +++ b/tests/test_litellm_rust/test_traces.py @@ -18,10 +18,10 @@ from fastapi import FastAPI from fastapi.testclient import TestClient from pydantic import BaseModel, ConfigDict, JsonValue, TypeAdapter -from litellm.constants import OTLP_MAX_ATTRIBUTE_VALUE_BYTES +from litellm.constants import AGENT_TRACING_LIST_PAGE_SIZE, OTLP_MAX_ATTRIBUTE_VALUE_BYTES from litellm.rust_bridge._native import NativeTraceConfig, NativeTraceStorage -from litellm.rust_bridge.trace.generated.models import ActivityAvailability, LensAccessParams, TraceQueryHelp -from litellm.rust_bridge.trace.generated.types import Trace, TraceScope +from litellm.rust_bridge.trace.generated.models import TraceQueryHelp +from litellm.rust_bridge.trace.generated.types import AllQueryScope, Trace, TracePage from litellm.rust_bridge.trace.storage import ClickHouseStorage, TraceStorageConfig, span_rows from litellm.tracing import Tenant, TraceReceiver, TracingPayloadTooLargeError from litellm.tracing.types import SpendLogRecord @@ -45,6 +45,7 @@ from tests.test_litellm_rust.support.recording_server import RecordingServer, Re pytestmark = pytest.mark.requires_rust_extension QUERY_ROWS: Final = TypeAdapter(tuple[dict[str, JsonValue], ...]) +_TRACE_PAGE: Final = TypeAdapter(TracePage) class CapturedSpendRow(BaseModel): @@ -90,23 +91,20 @@ def span_row() -> dict[str, JsonValue]: } -@pytest.fixture -def span_params() -> dict[str, str | int | list[str]]: - return {"trace_id": "trace-1", "trace_ref": "", "all_teams": 1, "user_id": "", "team_ids": []} - - @pytest.mark.asyncio async def test_trace_reader_projects_connection_and_parameters( - recording_server: RecordingServer, span_row: dict[str, JsonValue], span_params: dict[str, str | int | list[str]] + recording_server: RecordingServer, + span_row: dict[str, JsonValue], ) -> None: recording_server.enqueue(ResponseSpec(body={"data": [span_row]})) url: Final = recording_server.base_url.replace("http://", "http://reader:p%40ss%2Fword%25@") storage: Final = _native_storage("trace_test", url + "?database=wrong") - rows: Final = json.loads(await storage.query("trace_spans", span_params)) + trace: Final = TypeAdapter(Trace).validate_python( + await storage.get_trace("trace-1", AllQueryScope(kind="all"), "ref") + ) request: Final = recording_server.requests[0] parameters: Final = parse_qs(urlsplit(request.path).query) - assert rows == {"data": [span_row]} - assert b"o.TraceId = {trace_id:String}" in request.raw_body + assert trace["spans"][0]["span_id"] == span_row["span_id"] assert parameters["database"] == ["trace_test"] assert parameters["param_trace_id"] == ["trace-1"] assert parameters["readonly"] == ["1"] @@ -116,21 +114,11 @@ async def test_trace_reader_projects_connection_and_parameters( @pytest.mark.asyncio -async def test_trace_reader_rejects_success_status_with_embedded_error( - recording_server: RecordingServer, span_params: dict[str, str | int | list[str]] -) -> None: +async def test_trace_reader_rejects_success_status_with_embedded_error(recording_server: RecordingServer) -> None: recording_server.enqueue(ResponseSpec(body={"data": [], "exception": "query failed"})) storage: Final = _native_storage("trace_test", recording_server.base_url) - with pytest.raises(RuntimeError, match="invalid or failed JSON"): - await storage.query("trace_spans", span_params) - - -@pytest.mark.asyncio -async def test_reader_rejects_arbitrary_sql_before_sending(recording_server: RecordingServer) -> None: - recording_server.expected_requests = 0 - storage: Final = _native_storage("trace_test", recording_server.base_url) - with pytest.raises(ValueError, match="unknown ClickHouse read query"): - await storage.query("SELECT 1", {}) + with pytest.raises(RuntimeError, match="query failed"): + await storage.get_trace("trace-1", AllQueryScope(kind="all"), "ref") @pytest.mark.asyncio @@ -158,12 +146,97 @@ async def test_from_env_reads_with_clickhouse_url( recording_server.enqueue(ResponseSpec(body={"data": []})) monkeypatch.setenv("CLICKHOUSE_URL", recording_server.base_url) monkeypatch.delenv("CLICKHOUSE_READER_URL", raising=False) - scope: Final[TraceScope] = {"all_teams": 1, "user_id": "", "team_ids": ()} + scope: Final = AllQueryScope(kind="all") page: Final = await TraceReceiver.from_env().list_traces(scope, 0, 1) assert page == {"data": (), "next_cursor": None} assert len(recording_server.requests) == 1 +@pytest.fixture +def run_row() -> dict[str, JsonValue]: + return { + "trace_id": "trace", + "trace_ref": "ref", + "team_id": "", + "api_key_hash": "", + "user_id": "", + "name": "run", + "service": "test", + "input_preview": "", + "status": "STATUS_CODE_OK", + "start_ms": 1000, + "duration_ns": 1_000_000, + "span_count": 0, + "agent_count": 0, + "agent_invocations": 0, + "agent_names": [], + "frameworks": [], + "llm_calls": 0, + "tool_calls": 0, + "input_tokens": 0, + "output_tokens": 0, + "models": [], + "error_count": 0, + } + + +@pytest.mark.parametrize("params", ({}, {"start_ms": 1})) +def test_trace_pages_keep_the_effective_window_through_the_http_native_boundary( + recording_server: RecordingServer, + run_row: dict[str, JsonValue], + params: dict[str, int], +) -> None: + from litellm.proxy._types import LitellmUserRoles, UserAPIKeyAuth + from litellm.proxy.auth.user_api_key_auth import user_api_key_auth + from litellm.proxy.tracing_endpoints import provide_receiver, router + + recording_server.expected_requests = None + recording_server.default_response = ResponseSpec(body={"data": []}) + rows: Final = tuple( + {**run_row, "trace_id": f"trace-{index:03d}", "trace_ref": f"ref-{index:03d}"} + for index in range(AGENT_TRACING_LIST_PAGE_SIZE + 1) + ) + recording_server.enqueue(ResponseSpec(body={"data": rows})) + storage: Final = ClickHouseStorage(TraceStorageConfig(recording_server.base_url, "trace_test")) + app: Final = FastAPI() + app.include_router(router) + app.dependency_overrides[user_api_key_auth] = lambda: UserAPIKeyAuth( + user_role=LitellmUserRoles.PROXY_ADMIN, token="test" + ) + app.dependency_overrides[provide_receiver] = lambda: TraceReceiver(storage) + with TestClient(app) as client: + first: Final = client.get("/v1/traces", params=params) + assert first.status_code == 200, first.text + first_body: Final = _TRACE_PAGE.validate_python(first.json()) + first_parameters: Final = parse_qs(urlsplit(recording_server.requests[0].path).query) + cursor: Final = first_body["next_cursor"] + assert cursor + assert len(first_body["data"]) == AGENT_TRACING_LIST_PAGE_SIZE + recording_server.enqueue(ResponseSpec(body={"data": [rows[0]]})) + second: Final = client.get("/v1/traces", params={**params, "cursor": cursor}) + assert second.status_code == 200, second.text + second_body: Final = _TRACE_PAGE.validate_python(second.json()) + refs: Final = tuple(row["trace_ref"] for row in (*first_body["data"], *second_body["data"])) + assert refs == tuple(row["trace_ref"] for row in reversed(rows)) + assert second_body["next_cursor"] is None + run_parameters: Final = tuple( + parse_qs(urlsplit(request.path).query) + for request in recording_server.requests + if "param_sort_key" in parse_qs(urlsplit(request.path).query) + ) + assert len(run_parameters) == 2 + assert (run_parameters[1]["param_start_ms"], run_parameters[1]["param_end_ms"]) == ( + first_parameters["param_start_ms"], + first_parameters["param_end_ms"], + ) + sent: Final = len(recording_server.requests) + changed: Final = client.get( + "/v1/traces", params={"cursor": cursor, "end_ms": int(first_parameters["param_end_ms"][0]) + 1} + ) + assert changed.status_code == 400, changed.text + assert len(recording_server.requests) == sent + + @pytest.mark.asyncio async def test_schema_setup_uses_configured_retention(recording_server: RecordingServer) -> None: recording_server.expected_requests = None @@ -395,17 +468,24 @@ def test_trace_help_endpoint_runs_native_schema_and_metadata_discovery( @pytest.mark.parametrize( - ("clickhouse_status", "body", "expected_status"), + ("clickhouse_status", "body", "expected_status", "database_code"), ( - (400, b"ClickHouse rejected the query", 400), - (404, b"ClickHouse rejected the query", 400), - (500, b"ClickHouse rejected the query", 503), - (503, b"ClickHouse rejected the query", 503), - (200, b'{"data":[]}', 503), + (400, b"ClickHouse rejected the query", 400, None), + (404, b"ClickHouse rejected the query", 400, None), + (500, b"Code: 62. Invalid syntax", 400, 62), + (403, b"Code: 497. Access denied", 400, 497), + (500, b"Code: 241. Memory limit exceeded", 422, 241), + (500, b"ClickHouse rejected the query", 503, None), + (503, b"ClickHouse rejected the query", 503, None), + (200, b'{"data":[]}', 503, None), ), ) def test_trace_sql_endpoint_distinguishes_query_errors_from_reader_failures( - recording_server: RecordingServer, clickhouse_status: int, body: bytes, expected_status: int + recording_server: RecordingServer, + clickhouse_status: int, + body: bytes, + expected_status: int, + database_code: int | None, ) -> None: from fastapi import FastAPI from fastapi.testclient import TestClient @@ -434,6 +514,13 @@ def test_trace_sql_endpoint_distinguishes_query_errors_from_reader_failures( with TestClient(app) as client: failed: Final = client.post("/v1/traces/query", json={"sql": "SELEC 42"}) assert failed.status_code == expected_status, failed.text + assert failed.json()["detail"]["database_code"] == database_code + assert ( + failed.json()["detail"]["code"] + == {400: "query_rejected", 422: "query_limit_exceeded", 503: "query_unavailable"}[expected_status] + ) + if database_code is not None: + assert failed.json()["detail"]["message"] == body.decode() recovered: Final = client.post("/v1/traces/query", json={"sql": "SELECT 42 AS answer"}) assert recovered.status_code == 200, recovered.text assert recovered.json() == envelope @@ -445,14 +532,13 @@ async def test_trace_receiver_reads_with_only_one_clickhouse_url( recording_server: RecordingServer, monkeypatch: pytest.MonkeyPatch, span_row: dict[str, JsonValue], - span_params: dict[str, str | int | list[str]], ) -> None: monkeypatch.setenv("CLICKHOUSE_URL", recording_server.base_url) monkeypatch.setenv("CLICKHOUSE_DATABASE", "trace_test") monkeypatch.delenv("CLICKHOUSE_READER_URL", raising=False) recording_server.enqueue(ResponseSpec(body={"data": [span_row]})) receiver: Final = TraceReceiver.from_env() - trace: Final = await receiver.get_trace("trace-1", {"all_teams": 1, "user_id": "", "team_ids": ()}, "ref") + trace: Final = await receiver.get_trace("trace-1", AllQueryScope(kind="all"), "ref") assert trace is not None assert trace["spans"][0]["span_id"] == span_row["span_id"] assert trace["spans"][0]["duration_ms"] == int(str(span_row["duration_ns"])) / 1_000_000 @@ -461,20 +547,6 @@ async def test_trace_receiver_reads_with_only_one_clickhouse_url( assert parameters["readonly"] == ["1"] -@pytest.mark.asyncio -async def test_lens_read_uses_the_shared_native_query_and_returns_typed_rows( - recording_server: RecordingServer, -) -> None: - recording_server.enqueue(ResponseSpec(body={"data": [{"traces": 0, "requests": 1}]})) - storage: Final = ClickHouseStorage(TraceStorageConfig(recording_server.base_url, "trace_test")) - rows: Final = await storage.lens_availability(LensAccessParams(all_teams=0, team="team-a", key_hash="key-a")) - assert rows == (ActivityAvailability(traces=False, requests=True),) - parameters: Final = parse_qs(urlsplit(recording_server.requests[0].path).query) - assert parameters["param_all_teams"] == ["0"] - assert parameters["param_team"] == ["team-a"] - assert parameters["param_key_hash"] == ["key-a"] - - @dataclass(frozen=True, slots=True) class SeededTraceAPI: client: TestClient diff --git a/tests/unit/proxy/lens/test_sources.py b/tests/unit/proxy/lens/test_sources.py index cb79bea7ed5..833a477924e 100644 --- a/tests/unit/proxy/lens/test_sources.py +++ b/tests/unit/proxy/lens/test_sources.py @@ -1,10 +1,9 @@ -import base64 -import json -from collections.abc import Sequence +from collections.abc import Mapping, Sequence from dataclasses import dataclass, field from typing import Final import pytest +from pydantic import ValidationError from litellm.proxy.lens.models import ( ActivitySelection, @@ -300,27 +299,15 @@ def test_lens_scope_reads_with_the_same_access_as_traces(scope: Scope, access: Q def test_execution_ids_carry_reference_and_trace_id() -> None: assert parse_execution(execution_id(REF, "trace:with:colons")) == (REF, "trace:with:colons") - with pytest.raises(ValueError): + with pytest.raises(ValueError, match=r"^Not an execution ID$"): parse_execution("short:trace") -def legacy_id(*parts: str) -> str: - return base64.urlsafe_b64encode(json.dumps(parts).encode()).decode() - - -def test_saved_filters_load_as_one_search() -> None: +def test_current_selection_preserves_search_and_execution_ids() -> None: selection: Final = ActivitySelection.model_validate( { - "source": "both", - "agent_name": "research agent", - "service": "billing", - "team_id": "alpha", - "filters": [{"key": "tenant.tier", "value": "gold"}], - "execution_ids": [ - legacy_id("traces", "alpha", "trace", REF), - legacy_id("requests", "alpha", "request", ""), - legacy_id("traces", "alpha", "old"), - ], + "q": 'agent:"research agent" service:billing team:alpha attr.tenant.tier:gold', + "execution_ids": [execution_id(REF, "trace")], "sample_percent": 50, } ) @@ -329,17 +316,53 @@ def test_saved_filters_load_as_one_search() -> None: assert selection.sample_percent == 50 -def test_saved_job_samples_from_before_runs_load_as_absent() -> None: +@pytest.mark.parametrize( + ("field", "value"), + ( + ("source", "both"), + ("agent_name", "research agent"), + ("service", "billing"), + ("team_id", "alpha"), + ("filters", ({"key": "tenant.tier", "value": "gold"},)), + ), +) +def test_selection_rejects_obsolete_fields(field: str, value: str | tuple[Mapping[str, str], ...]) -> None: + with pytest.raises(ValidationError, match=field): + ActivitySelection.model_validate({"q": "agent:research", field: value}) + + +def test_saved_job_preserves_its_current_sample() -> None: job: Final = Job.model_validate( { "id": "job", "created_at": NOW, "start": NOW, "end": NOW, - "settings": {"name": "Research", "model": "analysis", "context": "Find failures", "agent_name": "a"}, + "settings": {"name": "Research", "model": "analysis", "context": "Find failures", "q": "agent:a"}, "revision": 1, - "sample": {"executions": [{"id": "x", "source": "traces"}], "eligible": 1}, + "sample": { + "executions": [{"id": execution_id(REF, "trace"), "trace_id": "trace", "trace_ref": REF}], + "eligible": 1, + }, } ) - assert job.sample is None + assert job.sample is not None + assert job.sample.executions == (Execution(id=execution_id(REF, "trace"), trace_id="trace", trace_ref=REF),) + assert job.sample.eligible == 1 assert job.settings.q == "agent:a" + assert Job.model_validate_json(job.model_dump_json()) == job + + +def test_saved_job_rejects_an_obsolete_execution_shape() -> None: + with pytest.raises(ValidationError, match="source"): + Job.model_validate( + { + "id": "job", + "created_at": NOW, + "start": NOW, + "end": NOW, + "settings": {"name": "Research", "model": "analysis", "context": "Find failures"}, + "revision": 1, + "sample": {"executions": [{"id": "x", "source": "traces"}], "eligible": 1}, + } + ) diff --git a/tests/unit/proxy/test_tracing_endpoints.py b/tests/unit/proxy/test_tracing_endpoints.py index e358ae52331..8bff4c38400 100644 --- a/tests/unit/proxy/test_tracing_endpoints.py +++ b/tests/unit/proxy/test_tracing_endpoints.py @@ -20,7 +20,7 @@ from litellm.proxy.auth.authorization_dependencies import get_log_team_lookup from litellm.proxy.auth.user_api_key_auth import user_api_key_auth from litellm.proxy.tracing_runtime import manage_tracing, provide_storage from litellm.rust_bridge import loader -from litellm.rust_bridge.trace.errors import TraceChanged +from litellm.rust_bridge.trace.errors import TraceChanged, TraceQueryError from litellm.rust_bridge.trace.generated.models import TraceQueryHelp from litellm.rust_bridge.trace.generated.types import AllQueryScope, OwnedQueryScope, QueryScope from litellm.rust_bridge.trace.queries import TraceSQLResponse @@ -340,11 +340,27 @@ def test_histogram_and_values_failures_map_to_read_statuses( assert client.get("/v1/traces/values/name").status_code == status -def test_list_traces_defaults_to_last_24h(client, receiver): - client.get("/v1/traces") - kwargs = receiver.list_traces.call_args.kwargs - assert kwargs["end_ms"] - kwargs["start_ms"] == tracing_endpoints.MS_PER_DAY - assert kwargs["cursor"] is None +@pytest.mark.parametrize( + ("params", "start_ms", "end_ms", "cursor"), + ( + ({}, None, None, None), + ({"cursor": "next"}, None, None, "next"), + ({"start_ms": 10, "cursor": "next"}, 10, None, "next"), + ({"end_ms": 20, "cursor": "next"}, None, 20, "next"), + ), +) +def test_list_traces_preserves_omitted_bounds_for_cursor_window( + client: TestClient, + receiver: MagicMock, + params: Mapping[str, str | int], + start_ms: int | None, + end_ms: int | None, + cursor: str | None, +) -> None: + response: Final = client.get("/v1/traces", params=params) + assert response.status_code == 200 + kwargs: Final = receiver.list_traces.call_args.kwargs + assert (kwargs["start_ms"], kwargs["end_ms"], kwargs["cursor"]) == (start_ms, end_ms, cursor) assert kwargs["q"] == "" @@ -723,15 +739,57 @@ def test_sql_rejects_missing_identity_without_querying( @pytest.mark.parametrize( - ("error", "status"), ((ValueError("invalid SQL"), 400), (RuntimeError("reader unavailable"), 503)) + ("error", "status", "detail"), + ( + ( + TraceQueryError("rejected", 10, "Syntax error at SELECT"), + 400, + {"code": "query_rejected", "database_code": 10, "message": "Syntax error at SELECT"}, + ), + ( + TraceQueryError("rejected", 20, "INSERT is denied by readonly mode"), + 400, + {"code": "query_rejected", "database_code": 20, "message": "INSERT is denied by readonly mode"}, + ), + ( + TraceQueryError("limited", 30, "Memory limit exceeded"), + 422, + {"code": "query_limit_exceeded", "database_code": 30, "message": "Memory limit exceeded"}, + ), + ( + TraceQueryError("limited", None, "Query exceeded the response size limit"), + 422, + { + "code": "query_limit_exceeded", + "database_code": None, + "message": "Query exceeded the response size limit", + }, + ), + ( + TraceQueryError("unavailable", 40, "Database temporarily unavailable"), + 503, + {"code": "query_unavailable", "database_code": 40, "message": "Database temporarily unavailable"}, + ), + ( + RuntimeError("private transport or credential details"), + 503, + {"code": "query_unavailable", "database_code": None, "message": "Trace SQL is temporarily unavailable"}, + ), + ( + ValueError("SQL query must not be empty"), + 400, + {"code": "query_rejected", "database_code": None, "message": "SQL query must not be empty"}, + ), + ), ) def test_sql_reports_rejected_queries_and_unavailable_readers( - client: TestClient, receiver: MagicMock, error: Exception, status: int + client: TestClient, receiver: MagicMock, error: Exception, status: int, detail: str | Mapping[str, object] ) -> None: client.app.dependency_overrides[tracing_endpoints.provide_trace_query_secret] = lambda: "test-secret" receiver.storage.query_sql = AsyncMock(side_effect=error) result: Final = client.post("/v1/traces/query", json={"sql": "SELECT 1"}) assert result.status_code == status, result.text + assert result.json()["detail"] == detail receiver.storage.query_sql.assert_awaited_once_with( "SELECT 1", {"kind": "owned", "user_id": "user", "team_ids": ()}, "test-secret" ) diff --git a/ui/litellm-dashboard/src/components/lens/traces/list/AgentTracesTable.tsx b/ui/litellm-dashboard/src/components/lens/traces/list/AgentTracesTable.tsx index a7c18c34286..8b9469b21c4 100644 --- a/ui/litellm-dashboard/src/components/lens/traces/list/AgentTracesTable.tsx +++ b/ui/litellm-dashboard/src/components/lens/traces/list/AgentTracesTable.tsx @@ -290,7 +290,8 @@ export function AgentTracesTable({ }: AgentTracesTableProps) { const settled = !isLoading && !error; const isEmpty = settled && !hasMore && traces.length === 0; - const autoContinue = settled && hasMore && traces.length > 0 && !isPlaceholder; + const canContinue = settled && hasMore && !isPlaceholder; + const autoContinue = canContinue && traces.length > 0; const { columnVisibility, onColumnVisibilityChange } = usePersistedColumnVisibility("lens-traces"); const sorting = useMemo(() => toSorting(order), [order]); const columns = useMemo(() => (picks ? [pickColumn(picks), ...RUN_COLUMNS] : RUN_COLUMNS), [picks]); diff --git a/ui/litellm-dashboard/src/components/lens/traces/list/runOrder.test.ts b/ui/litellm-dashboard/src/components/lens/traces/list/runOrder.test.ts index 0c827309a19..1ad9b8da139 100644 --- a/ui/litellm-dashboard/src/components/lens/traces/list/runOrder.test.ts +++ b/ui/litellm-dashboard/src/components/lens/traces/list/runOrder.test.ts @@ -7,32 +7,33 @@ import type { TracePage, TraceSummary } from "../types"; const template = (traceList as TracePage).data[0] as TraceSummary; const run = (overrides: Partial): TraceSummary => ({ ...template, ...overrides }); -const RUNS: TraceSummary[] = [ - run({ +const RUN_OVERRIDES: readonly Partial[] = [ + { trace_id: "a", trace_ref: "ref-a", start_time: "2026-09-30T06:00:00Z", duration_ms: 500, span_count: 3, error_count: 0, - }), - run({ + }, + { trace_id: "b", trace_ref: "ref-b", start_time: "2026-09-30T07:00:00Z", duration_ms: 500, span_count: 9, error_count: 2, - }), - run({ + }, + { trace_id: "c", trace_ref: "ref-c", start_time: "2026-09-30T05:00:00Z", duration_ms: 50, span_count: 1, error_count: 1, - }), + }, ]; +const RUNS: TraceSummary[] = RUN_OVERRIDES.map(run); const ids = (runs: readonly TraceSummary[]) => runs.map((item) => item.trace_id); diff --git a/ui/litellm-dashboard/src/components/lens/traces/list/useTraceHistogram.ts b/ui/litellm-dashboard/src/components/lens/traces/list/useTraceHistogram.ts index cc052b96b49..87d5dee0eae 100644 --- a/ui/litellm-dashboard/src/components/lens/traces/list/useTraceHistogram.ts +++ b/ui/litellm-dashboard/src/components/lens/traces/list/useTraceHistogram.ts @@ -1,4 +1,4 @@ -import { keepPreviousData, useQuery } from "@tanstack/react-query"; +import { keepPreviousData, useQuery, type UseQueryOptions } from "@tanstack/react-query"; import type { TimeWindow } from "@/components/shared/timeRange/timeRange"; import { type RunSelection, useTracesApi } from "../api"; @@ -40,12 +40,13 @@ export function useTraceHistogram( ): TraceHistogramResult { const traces = useTracesApi(accessToken); const { window, q } = selection; - const histogram = useQuery({ + const histogramOptions = { queryKey: ["agentTraceHistogram", traces.scope, window.startMs, window.endMs, q], queryFn: () => traces.histogram(selection, BUCKETS), enabled, placeholderData: keepPreviousData, select: toBuckets, - }); + } satisfies UseQueryOptions; + const histogram = useQuery(histogramOptions); return { buckets: histogram.data ?? emptyBuckets(window), isLoading: histogram.isLoading }; } diff --git a/ui/litellm-dashboard/src/lib/http/schema.d.ts b/ui/litellm-dashboard/src/lib/http/schema.d.ts index 55f56853ce9..7e541e11798 100644 --- a/ui/litellm-dashboard/src/lib/http/schema.d.ts +++ b/ui/litellm-dashboard/src/lib/http/schema.d.ts @@ -81305,13 +81305,15 @@ export interface operations { list_agent_traces_v1_traces_get: { parameters: { query?: { - /** @description Free text and key:value filters, e.g. `agent:research* -status:ok "book a flight"`. Keys: name, agent, status, model, input, trace_id. `*` globs and a leading `-` negates */ - q?: string; - cursor?: string | null; /** @description Window start, unix ms. Default: 24h ago */ start_ms?: number | null; /** @description Window end, unix ms. Default: now */ end_ms?: number | null; + /** @description Free text and key:value filters, e.g. `agent:research* -status:ok "book a flight"`. Keys: name, agent, status, model, input, trace_id, service, team and attr.. `*` globs and a leading `-` negates */ + q?: string; + cursor?: string | null; + sort_by?: "start_ms" | "duration_ms" | "span_count" | "error_count"; + sort_dir?: "asc" | "desc"; }; header?: never; path?: never; @@ -81362,7 +81364,7 @@ export interface operations { agent_trace_histogram_v1_traces_histogram_get: { parameters: { query?: { - /** @description Free text and key:value filters, e.g. `agent:research* -status:ok "book a flight"`. Keys: name, agent, status, model, input, trace_id. `*` globs and a leading `-` negates */ + /** @description Free text and key:value filters, e.g. `agent:research* -status:ok "book a flight"`. Keys: name, agent, status, model, input, trace_id, service, team and attr.. `*` globs and a leading `-` negates */ q?: string; buckets?: number; /** @description Window start, unix ms. Default: 24h ago */ @@ -81452,7 +81454,7 @@ export interface operations { agent_trace_values_v1_traces_values__field__get: { parameters: { query?: { - /** @description Free text and key:value filters, e.g. `agent:research* -status:ok "book a flight"`. Keys: name, agent, status, model, input, trace_id. `*` globs and a leading `-` negates */ + /** @description Free text and key:value filters, e.g. `agent:research* -status:ok "book a flight"`. Keys: name, agent, status, model, input, trace_id, service, team and attr.. `*` globs and a leading `-` negates */ q?: string; contains?: string; limit?: number;