From b970e412d9449fdf25887772e65de3b759f5e430 Mon Sep 17 00:00:00 2001 From: ishaan-berri <155045088+ishaan-berri@users.noreply.github.com> Date: Wed, 7 Oct 2026 18:44:19 -0700 Subject: [PATCH] perf(lens): faster trace opens and list pages at scale (#45228) * perf(lens): read trace spans in larger batches Co-Authored-By: Ishaan Jaffer <155045088+ishaan-berri@users.noreply.github.com> * perf(lens): evaluate the trace list page once per identity lookup Co-Authored-By: Ishaan Jaffer <155045088+ishaan-berri@users.noreply.github.com> * fix(lens): queue trace reads instead of rejecting past 8 in flight Co-Authored-By: Ishaan Jaffer <155045088+ishaan-berri@users.noreply.github.com> --------- Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- litellm-rust/crates/lens/Cargo.toml | 1 + litellm-rust/crates/lens/src/lib.rs | 61 ++++++++++++++++--- .../traces-clickhouse/query/list_traces.sql | 7 +-- .../traces-clickhouse/src/span_batches.rs | 25 ++++++-- .../tests/queries/support.rs | 1 + .../crates/traces-clickhouse/tests/reads.rs | 39 +++--------- 6 files changed, 83 insertions(+), 51 deletions(-) diff --git a/litellm-rust/crates/lens/Cargo.toml b/litellm-rust/crates/lens/Cargo.toml index d513e979f97..620d37cef7c 100644 --- a/litellm-rust/crates/lens/Cargo.toml +++ b/litellm-rust/crates/lens/Cargo.toml @@ -42,5 +42,6 @@ prettyplease = "0.2" [dev-dependencies] rstest.workspace = true +tokio = { workspace = true, features = ["test-util"] } wiremock.workspace = true uuid.workspace = true diff --git a/litellm-rust/crates/lens/src/lib.rs b/litellm-rust/crates/lens/src/lib.rs index 25830fda896..887f7f571fb 100644 --- a/litellm-rust/crates/lens/src/lib.rs +++ b/litellm-rust/crates/lens/src/lib.rs @@ -26,6 +26,7 @@ use litellm_traces_clickhouse::InsertTable; use serde_json::Value; use std::{ collections::BTreeMap, + future::Future, sync::{ Arc, atomic::{AtomicBool, Ordering}, @@ -33,6 +34,7 @@ use std::{ time::Duration, }; pub use storage::Storage; +use tokio::sync::Semaphore; #[allow( dead_code, @@ -46,7 +48,8 @@ pub use storage::Storage; pub mod wire { include!(concat!(env!("OUT_DIR"), "/wire.rs")); } -use tokio::sync::Semaphore; + +const READ_QUEUE_WAIT: Duration = Duration::from_secs(10); pub struct State { pub credentials: Arc, @@ -80,6 +83,15 @@ impl State { } } +async fn wait_for_read_slot

( + acquire: impl Future>, +) -> Result { + tokio::time::timeout(READ_QUEUE_WAIT, acquire) + .await + .map_err(|_| Error::Unavailable)? + .map_err(|_| Error::Unavailable) +} + pub fn router(state: Arc) -> Router { let public = Router::new() .route("/health/live", get(|| async { StatusCode::OK })) @@ -125,10 +137,7 @@ async fn receipt( ) -> Result, Error> { let tenant = state.credentials.tenant(&headers)?; state.require_storage()?; - let _permit = state - .read_slots - .try_acquire() - .map_err(|_| Error::Unavailable)?; + let _permit = wait_for_read_slot(state.read_slots.acquire()).await?; let body = tokio::time::timeout(Duration::from_secs(5), to_bytes(body, 64 * 1024)) .await .map_err(|_| Error::Unavailable)? @@ -206,11 +215,7 @@ async fn read( ) -> Result, Error> { auth::authorize_service(&headers, &state.service_token)?; state.require_storage()?; - let permit = state - .read_slots - .clone() - .try_acquire_owned() - .map_err(|_| Error::Unavailable)?; + let permit = wait_for_read_slot(state.read_slots.clone().acquire_owned()).await?; let body = tokio::time::timeout(Duration::from_secs(10), to_bytes(body, 1024 * 1024)) .await .map_err(|_| Error::Unavailable)? @@ -279,3 +284,39 @@ pub async fn provision(state: Arc) { tokio::time::sleep(Duration::from_secs(10)).await; } } + +#[cfg(test)] +mod tests { + use super::{Error, READ_QUEUE_WAIT, wait_for_read_slot}; + use std::sync::Arc; + use tokio::sync::Semaphore; + + #[tokio::test] + async fn ninth_read_waits_for_a_permit_and_succeeds() { + let slots = Arc::new(Semaphore::new(8)); + let permits = (0..8) + .map(|_| slots.clone().try_acquire_owned().expect("available permit")) + .collect::>(); + let waiting_slots = slots.clone(); + let waiting = + tokio::spawn(async move { wait_for_read_slot(waiting_slots.acquire_owned()).await }); + + tokio::task::yield_now().await; + assert!(!waiting.is_finished()); + drop(permits); + assert!(waiting.await.expect("joined read").is_ok()); + } + + #[tokio::test(start_paused = true)] + async fn read_queue_timeout_returns_unavailable() { + let slots = Arc::new(Semaphore::new(0)); + let waiting = tokio::spawn(wait_for_read_slot(slots.acquire_owned())); + + tokio::task::yield_now().await; + tokio::time::advance(READ_QUEUE_WAIT).await; + assert!(matches!( + waiting.await.expect("joined read"), + Err(Error::Unavailable) + )); + } +} diff --git a/litellm-rust/crates/traces-clickhouse/query/list_traces.sql b/litellm-rust/crates/traces-clickhouse/query/list_traces.sql index c52adf7ef49..fa9a67c1e20 100644 --- a/litellm-rust/crates/traces-clickhouse/query/list_traces.sql +++ b/litellm-rust/crates/traces-clickhouse/query/list_traces.sql @@ -5,7 +5,6 @@ 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, - min(StartTs) AS trace_start, max(EndTs) AS trace_end, dateDiff('millisecond', min(StartTs), max(EndTs)) AS duration_ms, sum(SpanCount) AS span_count, sum(AgentCount) AS agent_invocations, @@ -27,7 +26,7 @@ HAVING min(StartTs) >= fromUnixTimestamp64Milli({start_ms:Int64}) ORDER BY start_ms DESC, trace_ref DESC LIMIT {limit:UInt32} ) -SELECT page.* EXCEPT (trace_start, trace_end), +SELECT page.*, identities.agent_names AS agent_names, identities.agent_count AS agent_count, identities.frameworks AS frameworks FROM page @@ -37,10 +36,8 @@ LEFT JOIN ( arraySort(groupUniqArrayIf(toString(Framework), Framework != '')) AS frameworks, uniqExactIf(if(AgentName = '', SpanName, AgentName), ObservationType = 'agent') AS agent_count FROM otel_traces - WHERE Timestamp >= (SELECT min(trace_start) FROM page) - AND Timestamp <= (SELECT max(trace_end) FROM page) + WHERE Timestamp >= fromUnixTimestamp64Milli({start_ms:Int64}) AND TraceId IN (SELECT trace_id FROM page) - AND (TeamId, ApiKeyHash, TraceId) IN (SELECT team_id, api_key_hash, trace_id FROM page) GROUP BY TeamId, ApiKeyHash, TraceId ) AS identities ON page.team_id = identities.TeamId AND page.api_key_hash = identities.ApiKeyHash diff --git a/litellm-rust/crates/traces-clickhouse/src/span_batches.rs b/litellm-rust/crates/traces-clickhouse/src/span_batches.rs index 5e406d67cfa..a08ede7a43f 100644 --- a/litellm-rust/crates/traces-clickhouse/src/span_batches.rs +++ b/litellm-rust/crates/traces-clickhouse/src/span_batches.rs @@ -4,7 +4,7 @@ use std::{future::Future, marker::PhantomData}; use litellm_http::Client; -use litellm_storage_clickhouse::{Query, fetch}; +use litellm_storage_clickhouse::{Query, ReadLimits, fetch}; use litellm_traces::query::named as contracts; use litellm_traces_cache::{MAX_GRAPH_BYTES, MAX_GRAPH_SPANS, StoreError}; use serde::{Serialize, de::DeserializeOwned}; @@ -14,7 +14,12 @@ use crate::{ query::named::{SpendByResponseIdsParams, SpendByResponseIdsRow, TraceSpansRow}, }; -const PAGE_SIZE: u32 = 256; +const PAGE_SIZE: u32 = 8192; +const SPAN_BATCH_READ_LIMITS: ReadLimits = ReadLimits { + result_rows: PAGE_SIZE as u64, + response_bytes: 16 * 1024 * 1024, + ..litellm_storage_clickhouse::READ_LIMITS +}; #[derive(Default)] struct ReadBudget { @@ -67,6 +72,7 @@ struct Paged(PhantomData); impl Query for Paged { type Params = Batch; type Row = K::Row; + const READ_LIMITS: ReadLimits = SPAN_BATCH_READ_LIMITS; const SQL: &'static str = K::SQL; } @@ -316,8 +322,8 @@ mod tests { } #[rstest] - #[case::fits(1000, PAGE_SIZE, &[256, 256, 256, 256])] - #[case::uniform_large_rows(1000, 100, &[256, 128, 64, 64, 64, 64, 64, 64, 64, 64, 64, 64, 64, 64, 64, 64, 64, 64])] + #[case::fits(1000, PAGE_SIZE, &[8192])] + #[case::uniform_large_rows(200, 100, &[8192, 4096, 2048, 1024, 512, 256, 128, 64, 64, 64, 64])] #[tokio::test] async fn a_rejected_page_size_is_not_retried( #[case] total: u32, @@ -334,6 +340,13 @@ mod tests { assert_eq!(table.requests.lock().unwrap().as_slice(), requests); } + #[test] + fn span_batches_use_larger_read_limits() { + let limits = as Query>::READ_LIMITS; + assert_eq!(limits.result_rows, 8192); + assert_eq!(limits.response_bytes, 16 * 1024 * 1024); + } + #[rstest] #[tokio::test] async fn a_single_oversized_row_fails_the_read() { @@ -346,7 +359,9 @@ mod tests { assert!(matches!(result, Err(StoreError::TooLarge)), "{result:?}"); assert_eq!( table.requests.lock().unwrap().as_slice(), - &[256, 128, 64, 32, 16, 8, 4, 2, 1] + &[ + 8192, 4096, 2048, 1024, 512, 256, 128, 64, 32, 16, 8, 4, 2, 1 + ] ); } } diff --git a/litellm-rust/crates/traces-clickhouse/tests/queries/support.rs b/litellm-rust/crates/traces-clickhouse/tests/queries/support.rs index 7532d9aecc9..2b425bcfc1a 100644 --- a/litellm-rust/crates/traces-clickhouse/tests/queries/support.rs +++ b/litellm-rust/crates/traces-clickhouse/tests/queries/support.rs @@ -13,6 +13,7 @@ pub const DATABASE: &str = "trace_test"; pub struct SeededDatabase { pub database: ClickHouseDatabase, + #[allow(dead_code)] // dead_code: also shared with reads.rs, which uses a direct storage reader pub readers: QueryReaders, } diff --git a/litellm-rust/crates/traces-clickhouse/tests/reads.rs b/litellm-rust/crates/traces-clickhouse/tests/reads.rs index c938b6f451b..35dfafafd1d 100644 --- a/litellm-rust/crates/traces-clickhouse/tests/reads.rs +++ b/litellm-rust/crates/traces-clickhouse/tests/reads.rs @@ -3,9 +3,7 @@ use std::collections::BTreeMap; use litellm_http::Client; use litellm_traces::query::named::ReadAccessParams; use litellm_traces_cache::{ReadError, TraceReader}; -use litellm_traces_clickhouse::{ - ClickHouseTraces, Connection, InsertTable, QueryScope, insert_rows, -}; +use litellm_traces_clickhouse::{ClickHouseTraces, Connection, InsertTable, insert_rows}; use rstest::rstest; use serde_json::json; @@ -83,10 +81,7 @@ async fn list_costs_match_each_run_when_response_ids_are_reused( .collect(), ) .await?; - let connection = fixture - .readers - .connection(client, &QueryScope::All, "fixture-secret") - .await?; + let connection = Connection::reader(&fixture.database.url, DATABASE)?; let (reader, store) = make_reader(client, connection); let access = ReadAccessParams { all_teams: false, @@ -212,10 +207,7 @@ async fn large_runs_remain_complete_under_default_reader_limits( .collect::>(); insert_rows(client, &writer, DATABASE, InsertTable::SpendLogs, costs).await?; } - let connection = fixture - .readers - .connection(client, &QueryScope::All, "fixture-secret") - .await?; + let connection = Connection::reader(&fixture.database.url, DATABASE)?; let (reader, store) = make_reader(client, connection); let access = ReadAccessParams { all_teams: false, @@ -365,10 +357,7 @@ async fn cursor_pages_keep_a_tenant_scoped_snapshot_when_more_spans_arrive( ) -> TestResult { let fixture = seeded_database?; let client = &fixture.database.client; - let connection = fixture - .readers - .connection(client, &QueryScope::All, "fixture-secret") - .await?; + let connection = Connection::reader(&fixture.database.url, DATABASE)?; let (reader, store) = make_reader(client, connection.clone()); let access = ReadAccessParams { all_teams: true, @@ -528,10 +517,7 @@ async fn an_oversized_span_keeps_the_run_list_available_with_partial_totals( ) -> TestResult { let fixture = seeded_database?; let client = &fixture.database.client; - let connection = fixture - .readers - .connection(client, &QueryScope::All, "fixture-secret") - .await?; + let connection = Connection::reader(&fixture.database.url, DATABASE)?; let (reader, store) = make_reader(client, connection.clone()); let access = ReadAccessParams { all_teams: true, @@ -557,10 +543,7 @@ async fn an_oversized_span_keeps_the_run_list_available_with_partial_totals( ("TraceId".into(), json!(run.trace_id)), ("SpanId".into(), json!("oversized-child")), ("ParentSpanId".into(), json!("0101010101010101")), - ( - "SpanName".into(), - json!("x".repeat(litellm_storage_clickhouse::READ_LIMITS.response_bytes + 1)), - ), + ("SpanName".into(), json!("x".repeat(16 * 1024 * 1024 + 1))), ("ObservationType".into(), json!("tool")), ("TeamId".into(), json!("team-a")), ("ApiKeyHash".into(), json!("key-a")), @@ -699,10 +682,7 @@ async fn assigned_call_ids_require_shared_ownership_through_detail_and_batch_rea .collect(), ) .await?; - let connection = fixture - .readers - .connection(client, &QueryScope::All, "fixture-secret") - .await?; + let connection = Connection::reader(&fixture.database.url, DATABASE)?; let (reader, store) = make_reader(client, connection); let access = ReadAccessParams { all_teams: false, @@ -815,10 +795,7 @@ async fn native_cost_correlation_survives_session_grouping_and_excludes_other_ow .collect(), ) .await?; - let connection = fixture - .readers - .connection(client, &QueryScope::All, "fixture-secret") - .await?; + let connection = Connection::reader(&fixture.database.url, DATABASE)?; let (reader, store) = make_reader(client, connection); let access = ReadAccessParams { all_teams: true,