From a9b97007900f241994eeaf8582e7fde6bc557d1f Mon Sep 17 00:00:00 2001 From: "devin-ai-integration[bot]" <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Wed, 7 Oct 2026 12:50:35 -0700 Subject: [PATCH] perf(lens): bound single trace reads by the sampled start time (#45088) * perf(lens): bound single trace reads by the sampled start time Co-Authored-By: Ishaan Jaffer <155045088+ishaan-berri@users.noreply.github.com> * test(lens): allow unused query fixture field in load tests Co-Authored-By: Ishaan Jaffer <155045088+ishaan-berri@users.noreply.github.com> * test(lens): pass start_time in every lens content and evidence test 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> Co-authored-by: Ishaan Jaffer <155045088+ishaan-berri@users.noreply.github.com> --- .../traces-clickhouse/query/lens_content.sql | 2 + .../traces-clickhouse/query/lens_evidence.sql | 2 + .../traces-clickhouse/src/query/lens.rs | 2 + .../traces-clickhouse/src/query/number.rs | 2 +- .../crates/traces-clickhouse/tests/load.rs | 83 ++++++++++++++++++- .../traces-clickhouse/tests/migrations.rs | 30 ++++++- litellm/proxy/lens/sources.py | 2 + litellm/rust_bridge/trace/generated/models.py | 2 + .../traces-clickhouse/LensContentParams.json | 4 + .../traces-clickhouse/LensEvidenceParams.json | 4 + tests/unit/proxy/lens/test_sources.py | 17 +++- tests/unit/rust_bridge/trace/test_queries.py | 2 + 12 files changed, 148 insertions(+), 4 deletions(-) diff --git a/litellm-rust/crates/traces-clickhouse/query/lens_content.sql b/litellm-rust/crates/traces-clickhouse/query/lens_content.sql index 99fb56a5f48..52355a11061 100644 --- a/litellm-rust/crates/traces-clickhouse/query/lens_content.sql +++ b/litellm-rust/crates/traces-clickhouse/query/lens_content.sql @@ -17,6 +17,7 @@ SELECT * FROM ( FROM otel_traces WHERE {source:String}='traces' AND ({all_teams:UInt8}=1 OR TeamId={team:String}) AND ({key_hash:String}='' OR ApiKeyHash={key_hash:String}) + AND Timestamp >= parseDateTime64BestEffortOrZero({start_time:String}, 9) - INTERVAL 7 DAY AND ({trace_ref:String}='' OR hex(SHA256(concat(TeamId, char(0), ApiKeyHash, char(0), TraceId)))={trace_ref:String}) AND TraceId={id:String} AND TeamId={record_team:String} AND SpanId > {cursor:String} ORDER BY SpanId LIMIT 1 BY SpanId LIMIT 40 @@ -35,5 +36,6 @@ SELECT * FROM ( FROM spend_logs FINAL WHERE {source:String}='requests' AND ({all_teams:UInt8}=1 OR team_id={team:String}) AND ({key_hash:String}='' OR api_key={key_hash:String}) + AND spend_logs.start_time >= parseDateTime64BestEffortOrZero({start_time:String}, 3) - INTERVAL 7 DAY AND request_id={id:String} AND team_id={record_team:String} LIMIT 1 ) diff --git a/litellm-rust/crates/traces-clickhouse/query/lens_evidence.sql b/litellm-rust/crates/traces-clickhouse/query/lens_evidence.sql index a0d600cdfde..b53617364cf 100644 --- a/litellm-rust/crates/traces-clickhouse/query/lens_evidence.sql +++ b/litellm-rust/crates/traces-clickhouse/query/lens_evidence.sql @@ -2,6 +2,7 @@ SELECT sum(matches) AS count FROM ( SELECT count() AS matches FROM otel_traces WHERE {source:String}='traces' AND ({all_teams:UInt8}=1 OR TeamId={team:String}) AND ({key_hash:String}='' OR ApiKeyHash={key_hash:String}) + AND Timestamp >= parseDateTime64BestEffortOrZero({start_time:String}, 9) - INTERVAL 7 DAY AND ({trace_ref:String}='' OR hex(SHA256(concat(TeamId, char(0), ApiKeyHash, char(0), TraceId)))={trace_ref:String}) AND TraceId={id:String} AND TeamId={record_team:String} AND SpanId={span:String} AND position(concat('Input: ',Input,'\nOutput: ',Output,'\nStatus: ',StatusCode,' ',StatusMessage),{quote:String})>0 @@ -9,6 +10,7 @@ SELECT sum(matches) AS count FROM ( SELECT count() AS matches FROM spend_logs FINAL WHERE {source:String}='requests' AND ({all_teams:UInt8}=1 OR team_id={team:String}) AND ({key_hash:String}='' OR api_key={key_hash:String}) + AND spend_logs.start_time >= parseDateTime64BestEffortOrZero({start_time:String}, 3) - INTERVAL 7 DAY AND request_id={id:String} AND team_id={record_team:String} AND request_id={span:String} AND position(concat('Input: ',messages,'\nOutput: ',response,'\nError: ',error_str),{quote:String})>0 ) diff --git a/litellm-rust/crates/traces-clickhouse/src/query/lens.rs b/litellm-rust/crates/traces-clickhouse/src/query/lens.rs index 6439e696dad..77474c7143d 100644 --- a/litellm-rust/crates/traces-clickhouse/src/query/lens.rs +++ b/litellm-rust/crates/traces-clickhouse/src/query/lens.rs @@ -209,6 +209,7 @@ pub struct LensContentParams { pub source: ContentSource, pub id: String, pub record_team: String, + pub start_time: String, pub trace_ref: String, pub cursor: String, #[serde(deserialize_with = "super::number::deserialize")] @@ -252,6 +253,7 @@ pub struct LensEvidenceParams { pub source: ContentSource, pub id: String, pub record_team: String, + pub start_time: String, pub trace_ref: String, pub span: String, pub quote: String, diff --git a/litellm-rust/crates/traces-clickhouse/src/query/number.rs b/litellm-rust/crates/traces-clickhouse/src/query/number.rs index 9283903fee1..7a07845d0cd 100644 --- a/litellm-rust/crates/traces-clickhouse/src/query/number.rs +++ b/litellm-rust/crates/traces-clickhouse/src/query/number.rs @@ -99,7 +99,7 @@ mod tests { fn content_rejects_unsupported_sources(#[case] source: &str, #[case] valid: bool) { let parameters = serde_json::json!({ "all_teams": 0, "team": "team", "key_hash": "", "source": source, "id": "id", - "record_team": "team", "trace_ref": "", "cursor": "", "offset": 0 + "record_team": "team", "start_time": "", "trace_ref": "", "cursor": "", "offset": 0 }); assert_eq!( serde_json::from_value::(parameters).is_ok(), diff --git a/litellm-rust/crates/traces-clickhouse/tests/load.rs b/litellm-rust/crates/traces-clickhouse/tests/load.rs index 07c1095dfc3..aaec2c17e54 100644 --- a/litellm-rust/crates/traces-clickhouse/tests/load.rs +++ b/litellm-rust/crates/traces-clickhouse/tests/load.rs @@ -25,7 +25,8 @@ async fn seed_days(fixture: &SeededDatabase, first_day: u64, days: u64) -> TestR "INSERT INTO {DATABASE}.otel_traces \ (Timestamp, TraceId, SpanId, ParentSpanId, SpanName, ServiceName, ObservationType, TeamId, ApiKeyHash, Duration, SpanAttributes) \ SELECT now64(9) - toIntervalHour(intDiv(number, {SPANS_PER_DAY}) * 24 + 12 + {first_day} * 24), \ - concat('load-', toString(number + {first_row})), concat('span-', toString(number + {first_row})), \ + if({first_day} = 0, concat('load-', toString(number + {first_row})), 'load-0'), \ + concat('span-', toString(number + {first_row})), \ '', 'span', 'service', 'agent', 'load-team', '', 0, \ if({first_day}=0 AND number < {SPANS_PER_DAY}, map('payload', repeat('x', 3000)), map()) \ FROM numbers({count})" @@ -41,6 +42,62 @@ async fn seed_days(fixture: &SeededDatabase, first_day: u64, days: u64) -> TestR Ok(()) } +async fn trace_start_time(fixture: &SeededDatabase) -> TestResult { + let query = format!( + "SELECT toString(Timestamp, 'UTC') AS start_time FROM {DATABASE}.otel_traces \ + WHERE TraceId = 'load-0' LIMIT 1 FORMAT JSON" + ); + let response = fixture + .database + .client + .post(&fixture.database.url) + .body(query) + .send() + .await? + .error_for_status()? + .text() + .await?; + let result: Value = serde_json::from_str(&response)?; + result["data"][0]["start_time"] + .as_str() + .map(str::to_owned) + .ok_or_else(|| "trace start time missing".into()) +} + +fn content_parameters(start_time: &str) -> BTreeMap { + BTreeMap::from([ + ("source".into(), Parameter::Text("traces".into())), + ("all_teams".into(), Parameter::Integer(0)), + ("team".into(), Parameter::Text("load-team".into())), + ("key_hash".into(), Parameter::Text(String::new())), + ("id".into(), Parameter::Text("load-0".into())), + ("record_team".into(), Parameter::Text("load-team".into())), + ("start_time".into(), Parameter::Text(start_time.into())), + ("trace_ref".into(), Parameter::Text(String::new())), + ("cursor".into(), Parameter::Text(String::new())), + ("offset".into(), Parameter::Integer(1)), + ]) +} + +async fn content(fixture: &SeededDatabase, start_time: &str, query_id: &str) -> TestResult { + let connection = Connection::configured( + &format!("{}?query_id={query_id}", fixture.database.url), + DATABASE, + "default", + "", + )?; + let response = execute_named_read( + &fixture.database.client, + &connection, + ReadQuery::Content, + &content_parameters(start_time), + ) + .await?; + let result: Value = serde_json::from_str(&response)?; + assert!(!result["data"].as_array().ok_or("content rows")?.is_empty()); + Ok(()) +} + fn sample_parameters(start: u64, end: u64) -> BTreeMap { BTreeMap::from([ ("source".into(), Parameter::Text("traces".into())), @@ -144,3 +201,27 @@ async fn lens_sample_reads_scale_with_window_not_retention( ); Ok(()) } + +#[rstest] +#[tokio::test] +async fn lens_content_reads_scale_with_trace_not_retention( + #[future(awt)] migrated_database: TestResult, +) -> TestResult { + let fixture = migrated_database?; + seed_days(&fixture, 0, 8).await?; + let start_time = trace_start_time(&fixture).await?; + let before_id = format!("lens_content_before_{}", std::process::id()); + content(&fixture, &start_time, &before_id).await?; + let before = query_read_rows(&fixture, &before_id).await?; + + seed_days(&fixture, 8, 24).await?; + let after_id = format!("lens_content_after_{}", std::process::id()); + content(&fixture, &start_time, &after_id).await?; + let after = query_read_rows(&fixture, &after_id).await?; + println!("lens_content read_rows: before={before}, after={after}"); + assert!( + after * 100 <= before * 105, + "read_rows grew from {before} to {after}" + ); + Ok(()) +} diff --git a/litellm-rust/crates/traces-clickhouse/tests/migrations.rs b/litellm-rust/crates/traces-clickhouse/tests/migrations.rs index e51d9083c59..7630b3033a7 100644 --- a/litellm-rust/crates/traces-clickhouse/tests/migrations.rs +++ b/litellm-rust/crates/traces-clickhouse/tests/migrations.rs @@ -1150,6 +1150,7 @@ async fn lens_filters_reads_and_evidence_keep_reused_trace_ids_separate( ("source".into(), Parameter::Text("traces".into())), ("id".into(), Parameter::Text("shared".into())), ("record_team".into(), Parameter::Text("team".into())), + ("start_time".into(), Parameter::Text(String::new())), ("trace_ref".into(), Parameter::Text(first_ref.into())), ("cursor".into(), Parameter::Text(String::new())), ("offset".into(), Parameter::Integer(1)), @@ -1177,6 +1178,7 @@ async fn lens_filters_reads_and_evidence_keep_reused_trace_ids_separate( ("source".into(), Parameter::Text("traces".into())), ("id".into(), Parameter::Text("shared".into())), ("record_team".into(), Parameter::Text("team".into())), + ("start_time".into(), Parameter::Text(String::new())), ("trace_ref".into(), Parameter::Text(first_ref.into())), ("span".into(), Parameter::Text("root".into())), ("quote".into(), Parameter::Text(opposite.into())), @@ -1339,7 +1341,7 @@ async fn lens_selection_pages_without_losing_or_repeating_runs( #[case::traces("traces", 9)] #[case::requests("requests", 3)] #[tokio::test] -async fn lens_content_keeps_original_span_and_request_timestamps( +async fn lens_content_keeps_original_timestamps_with_start_time_slack( #[future(awt)] database: TestResult, #[case] source: &str, #[case] precision: usize, @@ -1379,6 +1381,30 @@ async fn lens_content_keeps_original_span_and_request_timestamps( ) .await?; let connection = Connection::configured(&database.url, "trace_test", "default", "")?; + let start_time_body = execute_read( + &database.client, + &connection, + "SELECT toString(fromUnixTimestamp64Nano({timestamp:Int64})) AS start_time FORMAT JSON", + &BTreeMap::from([( + "timestamp".into(), + Parameter::Integer(root_start + 86_400_000_000_000), + )]), + ) + .await?; + let start_time: serde_json::Value = serde_json::from_str(&start_time_body)?; + let start_time = start_time["data"][0]["start_time"] + .as_str() + .ok_or("start time missing")? + .to_owned(); + let parsed_time_body = execute_read( + &database.client, + &connection, + "SELECT toString(parseDateTime64BestEffortOrZero({start_time:String}, 9)) AS start_time FORMAT JSON", + &BTreeMap::from([("start_time".into(), Parameter::Text(start_time.clone()))]), + ) + .await?; + let parsed_time: serde_json::Value = serde_json::from_str(&parsed_time_body)?; + assert_eq!(parsed_time["data"][0]["start_time"], start_time); let parameters = BTreeMap::from([ ("source".into(), Parameter::Text(source.into())), ("all_teams".into(), Parameter::Integer(0)), @@ -1386,6 +1412,7 @@ async fn lens_content_keeps_original_span_and_request_timestamps( ("record_team".into(), Parameter::Text("team".into())), ("key_hash".into(), Parameter::Text(String::new())), ("trace_ref".into(), Parameter::Text(String::new())), + ("start_time".into(), Parameter::Text(start_time)), ("id".into(), Parameter::Text("run".into())), ("cursor".into(), Parameter::Text(String::new())), ("offset".into(), Parameter::Integer(1)), @@ -1453,6 +1480,7 @@ async fn lens_content_keeps_output_visible_after_long_input( ("record_team".into(), Parameter::Text("team".into())), ("key_hash".into(), Parameter::Text(String::new())), ("trace_ref".into(), Parameter::Text(String::new())), + ("start_time".into(), Parameter::Text(String::new())), ("id".into(), Parameter::Text("request".into())), ("cursor".into(), Parameter::Text(String::new())), ("offset".into(), Parameter::Integer(1)), diff --git a/litellm/proxy/lens/sources.py b/litellm/proxy/lens/sources.py index 015be69ecca..b1d142a6748 100644 --- a/litellm/proxy/lens/sources.py +++ b/litellm/proxy/lens/sources.py @@ -142,6 +142,7 @@ class SourceReader: id=execution.trace_id, trace_ref=execution.trace_ref, record_team=execution.team_id, + start_time=execution.start_time, cursor=cursor, offset=offset + 1, ) @@ -175,6 +176,7 @@ class SourceReader: id=execution.trace_id, trace_ref=execution.trace_ref, record_team=execution.team_id, + start_time=execution.start_time, span=evidence.span_id, quote=evidence.quote, ) diff --git a/litellm/rust_bridge/trace/generated/models.py b/litellm/rust_bridge/trace/generated/models.py index ea2c8bda648..5d84003aba2 100644 --- a/litellm/rust_bridge/trace/generated/models.py +++ b/litellm/rust_bridge/trace/generated/models.py @@ -215,6 +215,7 @@ class LensContentParams(LiteLLMBaseModel): source: ContentSource id: str record_team: str + start_time: str trace_ref: str cursor: str offset: int = Field(..., ge=0, le=4294967295) @@ -232,6 +233,7 @@ class LensEvidenceParams(LiteLLMBaseModel): source: ContentSource id: str record_team: str + start_time: str trace_ref: str span: str quote: str diff --git a/scripts/trace_codegen/schemas/traces-clickhouse/LensContentParams.json b/scripts/trace_codegen/schemas/traces-clickhouse/LensContentParams.json index 5ee5ab558ce..6026ccd26e1 100644 --- a/scripts/trace_codegen/schemas/traces-clickhouse/LensContentParams.json +++ b/scripts/trace_codegen/schemas/traces-clickhouse/LensContentParams.json @@ -39,6 +39,9 @@ "source": { "$ref": "#/$defs/ContentSource" }, + "start_time": { + "type": "string" + }, "team": { "type": "string" }, @@ -53,6 +56,7 @@ "source", "id", "record_team", + "start_time", "trace_ref", "cursor", "offset" diff --git a/scripts/trace_codegen/schemas/traces-clickhouse/LensEvidenceParams.json b/scripts/trace_codegen/schemas/traces-clickhouse/LensEvidenceParams.json index 07b9c216083..dbe9b32fdd6 100644 --- a/scripts/trace_codegen/schemas/traces-clickhouse/LensEvidenceParams.json +++ b/scripts/trace_codegen/schemas/traces-clickhouse/LensEvidenceParams.json @@ -36,6 +36,9 @@ "span": { "type": "string" }, + "start_time": { + "type": "string" + }, "team": { "type": "string" }, @@ -50,6 +53,7 @@ "source", "id", "record_team", + "start_time", "trace_ref", "span", "quote" diff --git a/tests/unit/proxy/lens/test_sources.py b/tests/unit/proxy/lens/test_sources.py index 8512db04eff..ad1f90f97ce 100644 --- a/tests/unit/proxy/lens/test_sources.py +++ b/tests/unit/proxy/lens/test_sources.py @@ -10,8 +10,10 @@ from litellm.proxy.lens.sources import SourceReader, execution_id, parse_executi from litellm.rust_bridge.trace.generated.models import ( ActivityAvailability, AgentRow, + CountRow, ExecutionRow, LensContentParams, + LensEvidenceParams, PartRow, ) from tests.unit.proxy.lens.test_agent_workspace import python_data @@ -183,9 +185,17 @@ async def test_recorded_times_survive_source_catalog_reads_search_and_python( class ContentStorage: async def lens_content(self, parameters: LensContentParams) -> tuple[PartRow, ...]: - assert parameters.source == source and parameters.record_team == "team" + assert ( + parameters.source == source + and parameters.record_team == "team" + and parameters.start_time == run.start_time + ) return rows + async def lens_evidence(self, parameters: LensEvidenceParams) -> tuple[CountRow, ...]: + assert parameters.start_time == run.start_time + return (CountRow(count=1),) + reader: Final = SourceReader(ContentStorage()) async def read(identity: str, cursor: str, offset: int) -> ExecutionContent: @@ -217,6 +227,11 @@ async def test_recorded_times_survive_source_catalog_reads_search_and_python( assert computed.sessions[0].parts == expected assert min(computed.sessions[0].parts, key=lambda part: part.start_time).span_id == rows[-1].span_id assert await workspace.valid(Evidence(execution_id=run.id, span_id=rows[0].span_id, quote=rows[0].content)) + assert await reader.verify_evidence( + Scope(team_id="team"), + run, + Evidence(execution_id=run.id, span_id=rows[0].span_id, quote=rows[0].content), + ) @pytest.mark.asyncio diff --git a/tests/unit/rust_bridge/trace/test_queries.py b/tests/unit/rust_bridge/trace/test_queries.py index 3c0556b1d97..9fd2964af20 100644 --- a/tests/unit/rust_bridge/trace/test_queries.py +++ b/tests/unit/rust_bridge/trace/test_queries.py @@ -18,6 +18,7 @@ def test_named_query_rejects_offsets_outside_the_native_integer_range(offset: in "source": "traces", "id": "trace", "record_team": "team", + "start_time": "", "trace_ref": "ref", "cursor": "", "offset": offset, @@ -34,6 +35,7 @@ def test_named_query_rejects_parameters_for_a_different_query() -> None: source="traces", id="trace", record_team="team", + start_time="", trace_ref="ref", cursor="", offset=0,