fix(traces): correct ClickHouse rollup partitioning, dedupe keys, and retention changes (#43901)

* fix(traces): correct ClickHouse rollup partitioning, dedupe keys, and retention changes

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* test(traces): pin spend dedupe timestamps within one second

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

---------

Co-authored-by: Yujong Lee <yujong@berri.ai>
Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
devin-ai-integration[bot] 2026-09-30 19:45:11 +00:00 • committed by GitHub
parent 50f5cc9bbb
commit 1fa3cde6a2
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
11 changed files with 310 additions and 63 deletions

View file

@ -23,9 +23,9 @@
<networks><ip>::/0</ip></networks>
<profile>litellm_traces_reader</profile>
<grants>
<query>GRANT SELECT ON default.otel_traces</query>
<query>GRANT SELECT ON default.agent_traces</query>
<query>GRANT SELECT ON default.spend_logs</query>
<query>GRANT SELECT ON litellm.otel_traces</query>
<query>GRANT SELECT ON litellm.agent_traces</query>
<query>GRANT SELECT ON litellm.spend_logs</query>
</grants>
</litellm_traces_reader>
</users>

View file

@ -44,5 +44,4 @@ CREATE TABLE IF NOT EXISTS {database}.otel_traces
ENGINE = MergeTree
PARTITION BY toDate(Timestamp)
ORDER BY (TeamId, ServiceName, toDateTime(Timestamp), TraceId)
TTL toDateTime(Timestamp) + INTERVAL {trace_retention_days} DAY
SETTINGS ttl_only_drop_parts = 1

View file

@ -5,9 +5,9 @@ CREATE TABLE IF NOT EXISTS {database}.agent_traces
StartTs SimpleAggregateFunction(min, DateTime64(9)),
EndTs SimpleAggregateFunction(max, DateTime64(9)),
ServiceName SimpleAggregateFunction(any, LowCardinality(String)),
RootName SimpleAggregateFunction(anyLast, String),
RootInput SimpleAggregateFunction(anyLast, String),
RootStatus SimpleAggregateFunction(anyLast, String),
RootName SimpleAggregateFunction(anyLast, Nullable(String)),
RootInput SimpleAggregateFunction(anyLast, Nullable(String)),
RootStatus SimpleAggregateFunction(anyLast, Nullable(String)),
SpanCount SimpleAggregateFunction(sum, UInt64),
AgentCount SimpleAggregateFunction(sum, UInt64),
LlmCount SimpleAggregateFunction(sum, UInt64),
@ -20,6 +20,4 @@ CREATE TABLE IF NOT EXISTS {database}.agent_traces
RequestIds SimpleAggregateFunction(groupArrayArray, Array(String))
)
ENGINE = AggregatingMergeTree
PARTITION BY toDate(StartTs)
ORDER BY (TeamId, TraceId)
TTL toDateTime(StartTs) + INTERVAL {trace_retention_days} DAY

View file

@ -4,9 +4,9 @@ SELECT
min(Timestamp) AS StartTs,
max(Timestamp + toIntervalNanosecond(Duration)) AS EndTs,
any(ServiceName) AS ServiceName,
anyLastIf(SpanName, ParentSpanId = '') AS RootName,
anyLastIf(InputPreview, ParentSpanId = '') AS RootInput,
anyLastIf(StatusCode, ParentSpanId = '') AS RootStatus,
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,

View file

@ -39,5 +39,4 @@ CREATE TABLE IF NOT EXISTS {database}.spend_logs
)
ENGINE = ReplacingMergeTree(end_time)
PARTITION BY toYYYYMM(start_time)
ORDER BY (team_id, toDateTime(start_time), request_id)
TTL toDateTime(start_time) + INTERVAL {spend_log_retention_days} DAY
ORDER BY (team_id, start_time, request_id)

View file

@ -0,0 +1 @@
ALTER TABLE {database}.otel_traces MODIFY TTL toDateTime(Timestamp) + INTERVAL {trace_retention_days} DAY

View file

@ -0,0 +1 @@
ALTER TABLE {database}.agent_traces MODIFY TTL toDateTime(StartTs) + INTERVAL {trace_retention_days} DAY

View file

@ -0,0 +1 @@
ALTER TABLE {database}.spend_logs MODIFY TTL toDateTime(start_time) + INTERVAL {spend_log_retention_days} DAY

View file

@ -1,13 +1,17 @@
use litellm_http::Client;
use std::time::Duration;
use crate::Connection;
use crate::Error;
const MIGRATIONS: [&str; 4] = [
const MIGRATIONS: [&str; 7] = [
include_str!("../migrations/0001_otel_traces.sql"),
include_str!("../migrations/0002_agent_traces.sql"),
include_str!("../migrations/0003_agent_traces_mv.sql"),
include_str!("../migrations/0004_spend_logs.sql"),
include_str!("../migrations/0005_otel_traces_ttl.sql"),
include_str!("../migrations/0006_agent_traces_ttl.sql"),
include_str!("../migrations/0007_spend_logs_ttl.sql"),
];
pub fn schema_statements(
@ -49,6 +53,7 @@ pub async fn ensure_schema(
for statement in schema_statements(database, trace_retention_days, spend_log_retention_days)? {
let response = client
.post(connection.url().clone())
.timeout(Duration::from_secs(10))
.body(statement)
.send()
.await

View file

@ -37,9 +37,10 @@ async fn database() -> Result<Database, Box<dyn std::error::Error>> {
);
let client = Client::no_redirect_for_test();
for sql in [
"CREATE TABLE otel_traces (n UInt8) ENGINE = Memory",
"INSERT INTO otel_traces VALUES (1)",
"CREATE TABLE private_traces (n UInt8) ENGINE = Memory",
"CREATE DATABASE litellm",
"CREATE TABLE litellm.otel_traces (n UInt8) ENGINE = Memory",
"INSERT INTO litellm.otel_traces VALUES (1)",
"CREATE TABLE litellm.private_traces (n UInt8) ENGINE = Memory",
] {
client
.post(&admin_url)
@ -48,7 +49,10 @@ async fn database() -> Result<Database, Box<dyn std::error::Error>> {
.await?
.error_for_status()?;
}
let url = admin_url.replacen("http://", "http://litellm_traces_reader:test_password@", 1);
let url = format!(
"{}?database=litellm",
admin_url.replacen("http://", "http://litellm_traces_reader:test_password@", 1)
);
Ok(Database {
_container: container,
url,
@ -64,7 +68,7 @@ async fn admin_sql_reads_rows_with_enforced_settings(
) -> Result<(), Box<dyn std::error::Error>> {
let database = database?;
let connection = Connection::parse(&format!(
"{}?readonly=0&default_format=TabSeparated&query=SELECT+2",
"{}&readonly=0&default_format=TabSeparated&query=SELECT+2",
database.url,
))?;
@ -98,7 +102,7 @@ async fn reader_rejects_writes_and_privilege_escalation(
#[case] sql: &str,
) -> Result<(), Box<dyn std::error::Error>> {
let database = database?;
let connection = Connection::parse(&format!("{}?readonly=0", database.url))?;
let connection = Connection::parse(&format!("{}&readonly=0", database.url))?;
let result = read(&database.client, &connection, sql).await;
@ -142,7 +146,7 @@ async fn admin_sql_enforces_result_row_limit(
) -> Result<(), Box<dyn std::error::Error>> {
let database = database?;
let connection = Connection::parse(&format!(
"{}?max_result_rows=0&result_overflow_mode=throw&wait_end_of_query=1",
"{}&max_result_rows=0&result_overflow_mode=throw&wait_end_of_query=1",
database.url,
))?;
@ -227,7 +231,7 @@ async fn query_parameters_preserve_values_and_replace_url_parameters(
#[future(awt)] database: Result<Database, Box<dyn std::error::Error>>,
) -> Result<(), Box<dyn std::error::Error>> {
let database = database?;
let connection = Connection::parse(&format!("{}?param_value=wrong", database.url))?;
let connection = Connection::parse(&format!("{}&param_value=wrong", database.url))?;
let values = vec![
"a'b".to_owned(),
"back\\slash".to_owned(),

View file

@ -1,19 +1,28 @@
use std::collections::BTreeMap;
use std::{collections::BTreeMap, time::Duration};
use litellm_http::Client;
use litellm_traces::{Connection, encode_rows, ensure_schema, execute_read, schema_statements};
use rstest::rstest;
use litellm_traces::{
Connection, Error, encode_rows, ensure_schema, execute_read, schema_statements,
};
use rstest::{fixture, rstest};
use testcontainers_modules::{
clickhouse::ClickHouse,
testcontainers::{ImageExt, runners::AsyncRunner},
testcontainers::{ContainerAsync, ImageExt, runners::AsyncRunner},
};
const CLICKHOUSE_TAG: &str =
"26.9.6.6@sha256:eb4870e7ca7ed70c259eebfcfbee6cf797017f6b5436c2926bbbfe3d4d28486e";
#[rstest]
#[tokio::test]
async fn schema_supports_span_rollups_and_spend_joins() -> Result<(), Box<dyn std::error::Error>> {
type TestResult<T = ()> = Result<T, Box<dyn std::error::Error>>;
struct ClickHouseDatabase {
_container: ContainerAsync<ClickHouse>,
url: String,
client: Client,
}
#[fixture]
async fn database() -> TestResult<ClickHouseDatabase> {
let container = ClickHouse::default()
.with_tag(CLICKHOUSE_TAG)
.with_env_var("CLICKHOUSE_SKIP_USER_SETUP", "1")
@ -24,10 +33,83 @@ async fn schema_supports_span_rollups_and_spend_joins() -> Result<(), Box<dyn st
container.get_host().await?,
container.get_host_port_ipv4(8123).await?
);
let client = Client::no_redirect_for_test();
let writer = Connection::writer(&url, "default", "")?;
ensure_schema(&client, &writer, "trace_test", 7, 14).await?;
ensure_schema(&client, &writer, "trace_test", 7, 14).await?;
Ok(ClickHouseDatabase {
_container: container,
url,
client: Client::no_redirect_for_test(),
})
}
async fn insert_rows(
database: &ClickHouseDatabase,
table: &str,
rows: Vec<BTreeMap<String, serde_json::Value>>,
) -> TestResult {
database
.client
.post(&database.url)
.query(&[
(
"query",
format!("INSERT INTO trace_test.{table} FORMAT JSONEachRow"),
),
("date_time_input_format", "best_effort".into()),
])
.body(encode_rows(rows)?)
.send()
.await?
.error_for_status()?;
Ok(())
}
async fn execute_write(database: &ClickHouseDatabase, sql: &str) -> TestResult {
database
.client
.post(&database.url)
.body(sql.to_owned())
.send()
.await?
.error_for_status()?;
Ok(())
}
async fn read_json(database: &ClickHouseDatabase, sql: &str) -> TestResult<serde_json::Value> {
let connection = Connection::configured(&database.url, "trace_test", "default", "")?;
let body = execute_read(&database.client, &connection, sql, &BTreeMap::new()).await?;
Ok(serde_json::from_str(&body)?)
}
async fn table_rows(database: &ClickHouseDatabase, table: &str) -> TestResult<u64> {
let response = read_json(
database,
&format!("SELECT count() AS rows FROM trace_test.{table}"),
)
.await?;
Ok(response["data"][0]["rows"]
.as_u64()
.expect("ClickHouse returns row counts as unsigned integers"))
}
async fn mutation_rows(database: &ClickHouseDatabase) -> TestResult<u64> {
let response = read_json(
database,
"SELECT count() AS rows FROM system.mutations WHERE database = 'trace_test'",
)
.await?;
Ok(response["data"][0]["rows"]
.as_u64()
.expect("ClickHouse returns mutation counts as unsigned integers"))
}
#[rstest]
#[tokio::test]
async fn schema_supports_span_rollups_and_spend_joins(
#[future(awt)] database: TestResult<ClickHouseDatabase>,
) -> TestResult {
let database = database?;
let writer = Connection::writer(&database.url, "default", "")?;
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
let timestamp = time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64;
let span = serde_json::from_value(serde_json::json!({
"Timestamp": timestamp, "TraceId": "trace-1", "SpanId": "span-1", "ParentSpanId": "",
@ -40,53 +122,210 @@ async fn schema_supports_span_rollups_and_spend_joins() -> Result<(), Box<dyn st
"start_time": timestamp / 1_000_000, "end_time": timestamp / 1_000_000 + 100,
"completion_start_time": null
}))?;
for (table, row) in [("otel_traces", span), ("spend_logs", spend)] {
client
.post(&url)
.query(&[
(
"query",
format!("INSERT INTO trace_test.{table} FORMAT JSONEachRow"),
),
("date_time_input_format", "best_effort".into()),
])
.body(encode_rows(vec![row])?)
.send()
.await?
.error_for_status()?;
}
let connection = Connection::configured(&url, "trace_test", "default", "")?;
let body = execute_read(&client, &connection,
insert_rows(&database, "otel_traces", vec![span]).await?;
insert_rows(&database, "spend_logs", vec![spend]).await?;
let body = read_json(
&database,
"SELECT o.TeamId, o.ApiKeyHash, o.ObservationType, o.InputPreview, s.spend, \
toString(toUnixTimestamp64Nano(o.Timestamp)) AS timestamp_ns, \
toString(toUnixTimestamp64Milli(s.start_time)) AS start_ms \
FROM otel_traces o JOIN spend_logs s ON o.LiteLLMRequestId = s.response_id AND o.TeamId = s.team_id",
&BTreeMap::new()).await?;
let response: serde_json::Value = serde_json::from_str(&body)?;
FROM trace_test.otel_traces o JOIN trace_test.spend_logs s \
ON o.LiteLLMRequestId = s.response_id AND o.TeamId = s.team_id",
)
.await?;
assert_eq!(
response["data"],
body["data"],
serde_json::json!([{
"TeamId": "team-1", "ApiKeyHash": "hash-1", "ObservationType": "agent",
"InputPreview": "hello world", "spend": 0.125,
"timestamp_ns": timestamp.to_string(), "start_ms": (timestamp / 1_000_000).to_string()
"timestamp_ns": timestamp.to_string(), "start_ms": (timestamp / 1_000_000).to_string()
}])
);
let body = execute_read(
&client,
&connection,
let body = read_json(
&database,
"SELECT toUInt32(sum(SpanCount)) AS spans, toUInt32(sum(InputTokens)) AS tokens \
FROM agent_traces WHERE TeamId = 'team-1' AND TraceId = 'trace-1'",
&BTreeMap::new(),
FROM trace_test.agent_traces WHERE TeamId = 'team-1' AND TraceId = 'trace-1'",
)
.await?;
let response: serde_json::Value = serde_json::from_str(&body)?;
assert_eq!(
response["data"],
body["data"],
serde_json::json!([{"spans": 1, "tokens": 12}])
);
Ok(())
}
#[rstest]
#[tokio::test]
async fn rollup_merges_spans_across_days_without_losing_root_fields(
#[future(awt)] database: TestResult<ClickHouseDatabase>,
) -> TestResult {
let database = database?;
let writer = Connection::writer(&database.url, "default", "")?;
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
let day_start = time::OffsetDateTime::now_utc()
.replace_time(time::Time::MIDNIGHT)
.unix_timestamp_nanos() as i64;
let root = serde_json::from_value(serde_json::json!({
"Timestamp": day_start - 1_000_000_000, "TraceId": "cross-day", "SpanId": "span-root",
"ParentSpanId": "", "ServiceName": "proxy", "SpanName": "root", "Input": "root input",
"StatusCode": "STATUS_CODE_ERROR",
"ResourceAttributes": {"litellm.team_id": "team-1"}
}))?;
insert_rows(&database, "otel_traces", vec![root]).await?;
let child = serde_json::from_value(serde_json::json!({
"Timestamp": day_start + 1_000_000_000, "TraceId": "cross-day", "SpanId": "span-child",
"ParentSpanId": "span-root", "ServiceName": "proxy", "SpanName": "child",
"StatusCode": "STATUS_CODE_UNSET",
"ResourceAttributes": {"litellm.team_id": "team-1"}
}))?;
insert_rows(&database, "otel_traces", vec![child]).await?;
execute_write(&database, "OPTIMIZE TABLE trace_test.agent_traces FINAL").await?;
let response = read_json(
&database,
"SELECT count() AS rows, any(RootName) AS RootName, any(RootInput) AS RootInput, \
any(RootStatus) AS RootStatus, sum(SpanCount) AS SpanCount \
FROM trace_test.agent_traces",
)
.await?;
assert_eq!(
response["data"],
serde_json::json!([{
"rows": 1, "RootName": "root", "RootInput": "root input",
"RootStatus": "STATUS_CODE_ERROR", "SpanCount": 2
}])
);
Ok(())
}
#[rstest]
#[tokio::test]
async fn spend_deduplication_preserves_subsecond_requests_and_retries(
#[future(awt)] database: TestResult<ClickHouseDatabase>,
) -> TestResult {
let database = database?;
let writer = Connection::writer(&database.url, "default", "")?;
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
let now_ms = time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64 / 1_000_000;
let base_start_time = now_ms / 1000 * 1000;
let first_start_time = base_start_time + 100;
let second_start_time = base_start_time + 200;
let first = serde_json::from_value(serde_json::json!({
"request_id": "same-request", "team_id": "team-1", "spend": 1.0,
"start_time": first_start_time, "end_time": first_start_time + 1000
}))?;
let second = serde_json::from_value(serde_json::json!({
"request_id": "same-request", "team_id": "team-1", "spend": 2.0,
"start_time": second_start_time, "end_time": second_start_time + 1200
}))?;
let retry = serde_json::from_value(serde_json::json!({
"request_id": "same-request", "team_id": "team-1", "spend": 1.0,
"start_time": first_start_time, "end_time": first_start_time + 2000
}))?;
insert_rows(&database, "spend_logs", vec![first]).await?;
insert_rows(&database, "spend_logs", vec![second]).await?;
insert_rows(&database, "spend_logs", vec![retry]).await?;
execute_write(&database, "OPTIMIZE TABLE trace_test.spend_logs FINAL").await?;
let rows = read_json(
&database,
"SELECT toString(toUnixTimestamp64Milli(start_time)) AS start_time, \
toString(toUnixTimestamp64Milli(end_time)) AS end_time \
FROM trace_test.spend_logs ORDER BY start_time",
)
.await?;
assert_eq!(
rows["data"],
serde_json::json!([
{
"start_time": first_start_time.to_string(),
"end_time": (first_start_time + 2000).to_string()
},
{
"start_time": second_start_time.to_string(),
"end_time": (second_start_time + 1200).to_string()
}
])
);
assert_eq!(table_rows(&database, "spend_logs").await?, 2);
Ok(())
}
#[rstest]
#[tokio::test]
async fn retention_changes_materialize_existing_rows_and_remain_idempotent(
#[future(awt)] database: TestResult<ClickHouseDatabase>,
) -> TestResult {
let database = database?;
let writer = Connection::writer(&database.url, "default", "")?;
ensure_schema(&database.client, &writer, "trace_test", 30, 30).await?;
let old_time = time::OffsetDateTime::now_utc() - time::Duration::days(20);
let old_timestamp_ns = old_time.unix_timestamp_nanos() as i64;
let old_timestamp_ms = old_timestamp_ns / 1_000_000;
let span = serde_json::from_value(serde_json::json!({
"Timestamp": old_timestamp_ns, "TraceId": "expired", "SpanId": "span-old",
"ParentSpanId": "", "ServiceName": "proxy", "SpanName": "old-root", "Input": "old input",
"ResourceAttributes": {"litellm.team_id": "team-1"}
}))?;
let spend = serde_json::from_value(serde_json::json!({
"request_id": "old-request", "team_id": "team-1", "spend": 1.0,
"start_time": old_timestamp_ms, "end_time": old_timestamp_ms + 1000
}))?;
insert_rows(&database, "otel_traces", vec![span]).await?;
insert_rows(&database, "spend_logs", vec![spend]).await?;
ensure_schema(&database.client, &writer, "trace_test", 14, 14).await?;
let deadline = tokio::time::Instant::now() + Duration::from_secs(60);
loop {
let response = read_json(
&database,
"SELECT countIf(is_done = 0) AS pending \
FROM system.mutations WHERE database = 'trace_test'",
)
.await?;
let pending = response["data"][0]["pending"]
.as_u64()
.expect("ClickHouse returns pending mutation counts as unsigned integers");
if pending == 0 {
break;
}
assert!(
tokio::time::Instant::now() < deadline,
"ClickHouse TTL mutations did not finish before the deadline"
);
tokio::time::sleep(Duration::from_millis(100)).await;
}
execute_write(&database, "OPTIMIZE TABLE trace_test.otel_traces FINAL").await?;
execute_write(&database, "OPTIMIZE TABLE trace_test.agent_traces FINAL").await?;
execute_write(&database, "OPTIMIZE TABLE trace_test.spend_logs FINAL").await?;
assert_eq!(table_rows(&database, "otel_traces").await?, 0);
assert_eq!(table_rows(&database, "agent_traces").await?, 0);
assert_eq!(table_rows(&database, "spend_logs").await?, 0);
let mutation_count = mutation_rows(&database).await?;
ensure_schema(&database.client, &writer, "trace_test", 14, 14).await?;
assert_eq!(mutation_rows(&database).await?, mutation_count);
Ok(())
}
#[rstest]
#[tokio::test]
async fn schema_statement_timeout_maps_to_transport_error() -> TestResult {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
let address = listener.local_addr()?;
let server = tokio::spawn(async move {
let (_connection, _) = listener.accept().await.expect("accept schema request");
std::future::pending::<()>().await;
});
let client = Client::no_redirect_for_test();
let url = format!("http://{address}");
let writer = Connection::writer(&url, "default", "")?;
let result = tokio::time::timeout(
Duration::from_secs(12),
ensure_schema(&client, &writer, "trace_test", 7, 14),
)
.await;
server.abort();
assert!(matches!(result, Ok(Err(Error::Transport))), "{result:?}");
Ok(())
}
#[rstest]
#[case::empty("", 7, 14)]
#[case::sql("db; DROP DATABASE default", 7, 14)]