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