From 1fa3cde6a2882ffb0311047be08c044f01561613 Mon Sep 17 00:00:00 2001 From: "devin-ai-integration[bot]" <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Wed, 30 Sep 2026 19:45:11 +0000 Subject: [PATCH] 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 Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- litellm-rust/crates/traces/config/reader.xml | 6 +- .../traces/migrations/0001_otel_traces.sql | 1 - .../traces/migrations/0002_agent_traces.sql | 8 +- .../migrations/0003_agent_traces_mv.sql | 6 +- .../traces/migrations/0004_spend_logs.sql | 3 +- .../migrations/0005_otel_traces_ttl.sql | 1 + .../migrations/0006_agent_traces_ttl.sql | 1 + .../traces/migrations/0007_spend_logs_ttl.sql | 1 + litellm-rust/crates/traces/src/schema.rs | 7 +- litellm-rust/crates/traces/tests/admin_sql.rs | 20 +- .../crates/traces/tests/migrations.rs | 319 +++++++++++++++--- 11 files changed, 310 insertions(+), 63 deletions(-) create mode 100644 litellm-rust/crates/traces/migrations/0005_otel_traces_ttl.sql create mode 100644 litellm-rust/crates/traces/migrations/0006_agent_traces_ttl.sql create mode 100644 litellm-rust/crates/traces/migrations/0007_spend_logs_ttl.sql diff --git a/litellm-rust/crates/traces/config/reader.xml b/litellm-rust/crates/traces/config/reader.xml index 4ff2e9d9d84..73a63035e2e 100644 --- a/litellm-rust/crates/traces/config/reader.xml +++ b/litellm-rust/crates/traces/config/reader.xml @@ -23,9 +23,9 @@ ::/0 litellm_traces_reader - GRANT SELECT ON default.otel_traces - GRANT SELECT ON default.agent_traces - GRANT SELECT ON default.spend_logs + GRANT SELECT ON litellm.otel_traces + GRANT SELECT ON litellm.agent_traces + GRANT SELECT ON litellm.spend_logs diff --git a/litellm-rust/crates/traces/migrations/0001_otel_traces.sql b/litellm-rust/crates/traces/migrations/0001_otel_traces.sql index f228ee8144c..aed869e6ee0 100644 --- a/litellm-rust/crates/traces/migrations/0001_otel_traces.sql +++ b/litellm-rust/crates/traces/migrations/0001_otel_traces.sql @@ -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 diff --git a/litellm-rust/crates/traces/migrations/0002_agent_traces.sql b/litellm-rust/crates/traces/migrations/0002_agent_traces.sql index ea8177fa12f..d994829f9c0 100644 --- a/litellm-rust/crates/traces/migrations/0002_agent_traces.sql +++ b/litellm-rust/crates/traces/migrations/0002_agent_traces.sql @@ -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 diff --git a/litellm-rust/crates/traces/migrations/0003_agent_traces_mv.sql b/litellm-rust/crates/traces/migrations/0003_agent_traces_mv.sql index 3b3c5ed3adc..e162d4da569 100644 --- a/litellm-rust/crates/traces/migrations/0003_agent_traces_mv.sql +++ b/litellm-rust/crates/traces/migrations/0003_agent_traces_mv.sql @@ -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, diff --git a/litellm-rust/crates/traces/migrations/0004_spend_logs.sql b/litellm-rust/crates/traces/migrations/0004_spend_logs.sql index c20c182ba3c..a14930f438f 100644 --- a/litellm-rust/crates/traces/migrations/0004_spend_logs.sql +++ b/litellm-rust/crates/traces/migrations/0004_spend_logs.sql @@ -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) diff --git a/litellm-rust/crates/traces/migrations/0005_otel_traces_ttl.sql b/litellm-rust/crates/traces/migrations/0005_otel_traces_ttl.sql new file mode 100644 index 00000000000..4ac597b8902 --- /dev/null +++ b/litellm-rust/crates/traces/migrations/0005_otel_traces_ttl.sql @@ -0,0 +1 @@ +ALTER TABLE {database}.otel_traces MODIFY TTL toDateTime(Timestamp) + INTERVAL {trace_retention_days} DAY diff --git a/litellm-rust/crates/traces/migrations/0006_agent_traces_ttl.sql b/litellm-rust/crates/traces/migrations/0006_agent_traces_ttl.sql new file mode 100644 index 00000000000..d4c4a113329 --- /dev/null +++ b/litellm-rust/crates/traces/migrations/0006_agent_traces_ttl.sql @@ -0,0 +1 @@ +ALTER TABLE {database}.agent_traces MODIFY TTL toDateTime(StartTs) + INTERVAL {trace_retention_days} DAY diff --git a/litellm-rust/crates/traces/migrations/0007_spend_logs_ttl.sql b/litellm-rust/crates/traces/migrations/0007_spend_logs_ttl.sql new file mode 100644 index 00000000000..131573927ac --- /dev/null +++ b/litellm-rust/crates/traces/migrations/0007_spend_logs_ttl.sql @@ -0,0 +1 @@ +ALTER TABLE {database}.spend_logs MODIFY TTL toDateTime(start_time) + INTERVAL {spend_log_retention_days} DAY diff --git a/litellm-rust/crates/traces/src/schema.rs b/litellm-rust/crates/traces/src/schema.rs index 3959d78b4c6..5258c86ca78 100644 --- a/litellm-rust/crates/traces/src/schema.rs +++ b/litellm-rust/crates/traces/src/schema.rs @@ -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 diff --git a/litellm-rust/crates/traces/tests/admin_sql.rs b/litellm-rust/crates/traces/tests/admin_sql.rs index b2ef6113a15..20a7a229fdc 100644 --- a/litellm-rust/crates/traces/tests/admin_sql.rs +++ b/litellm-rust/crates/traces/tests/admin_sql.rs @@ -37,9 +37,10 @@ async fn database() -> Result> { ); 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> { .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> { 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> { 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> { 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>, ) -> Result<(), Box> { let database = database?; - let connection = Connection::parse(&format!("{}?param_value=wrong", database.url))?; + let connection = Connection::parse(&format!("{}¶m_value=wrong", database.url))?; let values = vec![ "a'b".to_owned(), "back\\slash".to_owned(), diff --git a/litellm-rust/crates/traces/tests/migrations.rs b/litellm-rust/crates/traces/tests/migrations.rs index a77d57cfb0f..48baaee0828 100644 --- a/litellm-rust/crates/traces/tests/migrations.rs +++ b/litellm-rust/crates/traces/tests/migrations.rs @@ -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> { +type TestResult = Result>; + +struct ClickHouseDatabase { + _container: ContainerAsync, + url: String, + client: Client, +} + +#[fixture] +async fn database() -> TestResult { 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>, +) -> 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 { + 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 { + 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 { + 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, +) -> 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, +) -> 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, +) -> 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, +) -> 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)]