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>
This commit is contained in:
devin-ai-integration[bot] 2026-10-07 12:50:35 -07:00 • committed by GitHub
parent c877e0f055
commit a9b9700790
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
12 changed files with 148 additions and 4 deletions

View file

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

View file

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

View file

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

View file

@ -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::<crate::query::lens::LensContentParams>(parameters).is_ok(),

View file

@ -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<String> {
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<String, Parameter> {
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<String, Parameter> {
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<SeededDatabase>,
) -> 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(())
}

View file

@ -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<ClickHouseDatabase>,
#[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)),

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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