diff --git a/litellm-rust/crates/traces/config/reader.xml b/litellm-rust/crates/traces/config/reader.xml index 73a63035e2e..3ab337a13fc 100644 --- a/litellm-rust/crates/traces/config/reader.xml +++ b/litellm-rust/crates/traces/config/reader.xml @@ -24,7 +24,7 @@ litellm_traces_reader GRANT SELECT ON litellm.otel_traces - GRANT SELECT ON litellm.agent_traces + GRANT SELECT ON litellm.agent_traces_by_key GRANT SELECT ON litellm.spend_logs diff --git a/litellm-rust/crates/traces/migrations/0002_agent_traces.sql b/litellm-rust/crates/traces/migrations/0002_agent_traces.sql index d994829f9c0..bb5cf67794d 100644 --- a/litellm-rust/crates/traces/migrations/0002_agent_traces.sql +++ b/litellm-rust/crates/traces/migrations/0002_agent_traces.sql @@ -1,6 +1,7 @@ -CREATE TABLE IF NOT EXISTS {database}.agent_traces +CREATE TABLE IF NOT EXISTS {database}.agent_traces_by_key ( TeamId LowCardinality(String), + ApiKeyHash String, TraceId String, StartTs SimpleAggregateFunction(min, DateTime64(9)), EndTs SimpleAggregateFunction(max, DateTime64(9)), @@ -20,4 +21,4 @@ CREATE TABLE IF NOT EXISTS {database}.agent_traces RequestIds SimpleAggregateFunction(groupArrayArray, Array(String)) ) ENGINE = AggregatingMergeTree -ORDER BY (TeamId, TraceId) +ORDER BY (TeamId, ApiKeyHash, TraceId) 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 e162d4da569..94dad81f998 100644 --- a/litellm-rust/crates/traces/migrations/0003_agent_traces_mv.sql +++ b/litellm-rust/crates/traces/migrations/0003_agent_traces_mv.sql @@ -1,6 +1,7 @@ -CREATE MATERIALIZED VIEW IF NOT EXISTS {database}.agent_traces_mv TO {database}.agent_traces AS +CREATE MATERIALIZED VIEW IF NOT EXISTS {database}.agent_traces_by_key_mv +TO {database}.agent_traces_by_key AS SELECT - TeamId, TraceId, + TeamId, ApiKeyHash, TraceId, min(Timestamp) AS StartTs, max(Timestamp + toIntervalNanosecond(Duration)) AS EndTs, any(ServiceName) AS ServiceName, @@ -18,4 +19,4 @@ SELECT groupUniqArrayIf(SpanName, ObservationType = 'agent') AS AgentNames, groupArrayIf(LiteLLMRequestId, LiteLLMRequestId != '') AS RequestIds FROM {database}.otel_traces -GROUP BY TeamId, TraceId +GROUP BY TeamId, ApiKeyHash, TraceId diff --git a/litellm-rust/crates/traces/migrations/0006_agent_traces_ttl.sql b/litellm-rust/crates/traces/migrations/0006_agent_traces_ttl.sql index d4c4a113329..8681f0622a4 100644 --- a/litellm-rust/crates/traces/migrations/0006_agent_traces_ttl.sql +++ b/litellm-rust/crates/traces/migrations/0006_agent_traces_ttl.sql @@ -1 +1 @@ -ALTER TABLE {database}.agent_traces MODIFY TTL toDateTime(StartTs) + INTERVAL {trace_retention_days} DAY +ALTER TABLE {database}.agent_traces_by_key MODIFY TTL toDateTime(StartTs) + INTERVAL {trace_retention_days} DAY diff --git a/litellm-rust/crates/traces/migrations/0008_agent_traces_by_key.sql b/litellm-rust/crates/traces/migrations/0008_agent_traces_by_key.sql deleted file mode 100644 index bb5cf67794d..00000000000 --- a/litellm-rust/crates/traces/migrations/0008_agent_traces_by_key.sql +++ /dev/null @@ -1,24 +0,0 @@ -CREATE TABLE IF NOT EXISTS {database}.agent_traces_by_key -( - TeamId LowCardinality(String), - ApiKeyHash String, - TraceId String, - StartTs SimpleAggregateFunction(min, DateTime64(9)), - EndTs SimpleAggregateFunction(max, DateTime64(9)), - ServiceName SimpleAggregateFunction(any, LowCardinality(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), - ToolCount SimpleAggregateFunction(sum, UInt64), - ErrorCount SimpleAggregateFunction(sum, UInt64), - InputTokens SimpleAggregateFunction(sum, UInt64), - OutputTokens SimpleAggregateFunction(sum, UInt64), - Models SimpleAggregateFunction(groupUniqArrayArray, Array(String)), - AgentNames SimpleAggregateFunction(groupUniqArrayArray, Array(String)), - RequestIds SimpleAggregateFunction(groupArrayArray, Array(String)) -) -ENGINE = AggregatingMergeTree -ORDER BY (TeamId, ApiKeyHash, TraceId) diff --git a/litellm-rust/crates/traces/migrations/0009_agent_traces_by_key_mv.sql b/litellm-rust/crates/traces/migrations/0009_agent_traces_by_key_mv.sql deleted file mode 100644 index 94dad81f998..00000000000 --- a/litellm-rust/crates/traces/migrations/0009_agent_traces_by_key_mv.sql +++ /dev/null @@ -1,22 +0,0 @@ -CREATE MATERIALIZED VIEW IF NOT EXISTS {database}.agent_traces_by_key_mv -TO {database}.agent_traces_by_key AS -SELECT - TeamId, ApiKeyHash, TraceId, - min(Timestamp) AS StartTs, - max(Timestamp + toIntervalNanosecond(Duration)) AS EndTs, - any(ServiceName) AS ServiceName, - 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, - countIf(ObservationType = 'tool') AS ToolCount, - countIf(StatusCode = 'STATUS_CODE_ERROR') AS ErrorCount, - sum(InputTokens) AS InputTokens, - sum(OutputTokens) AS OutputTokens, - groupUniqArrayIf(toString(Model), Model != '') AS Models, - groupUniqArrayIf(SpanName, ObservationType = 'agent') AS AgentNames, - groupArrayIf(LiteLLMRequestId, LiteLLMRequestId != '') AS RequestIds -FROM {database}.otel_traces -GROUP BY TeamId, ApiKeyHash, TraceId diff --git a/litellm-rust/crates/traces/src/schema.rs b/litellm-rust/crates/traces/src/schema.rs index 3e3ba7efe07..5a154eb87c3 100644 --- a/litellm-rust/crates/traces/src/schema.rs +++ b/litellm-rust/crates/traces/src/schema.rs @@ -6,7 +6,7 @@ use crate::Error; const SCHEMA_REQUEST_TIMEOUT: Duration = Duration::from_secs(30); -const MIGRATIONS: [&str; 9] = [ +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"), @@ -14,8 +14,6 @@ const MIGRATIONS: [&str; 9] = [ include_str!("../migrations/0005_otel_traces_ttl.sql"), include_str!("../migrations/0006_agent_traces_ttl.sql"), include_str!("../migrations/0007_spend_logs_ttl.sql"), - include_str!("../migrations/0008_agent_traces_by_key.sql"), - include_str!("../migrations/0009_agent_traces_by_key_mv.sql"), ]; pub fn schema_statements( diff --git a/litellm-rust/crates/traces/tests/admin_sql.rs b/litellm-rust/crates/traces/tests/admin_sql.rs index 255740f9634..ab0eb873a28 100644 --- a/litellm-rust/crates/traces/tests/admin_sql.rs +++ b/litellm-rust/crates/traces/tests/admin_sql.rs @@ -40,8 +40,8 @@ async fn database() -> Result> { "CREATE DATABASE litellm", "CREATE TABLE litellm.otel_traces (n UInt8) ENGINE = Memory", "INSERT INTO litellm.otel_traces VALUES (1)", - "CREATE TABLE litellm.agent_traces (n UInt8) ENGINE = Memory", - "INSERT INTO litellm.agent_traces VALUES (2)", + "CREATE TABLE litellm.agent_traces_by_key (n UInt8) ENGINE = Memory", + "INSERT INTO litellm.agent_traces_by_key VALUES (4)", "CREATE TABLE litellm.spend_logs (n UInt8) ENGINE = Memory", "INSERT INTO litellm.spend_logs VALUES (3)", "CREATE TABLE litellm.private_traces (n UInt8) ENGINE = Memory", @@ -86,6 +86,15 @@ async fn admin_sql_reads_rows_with_enforced_settings( let json: Value = serde_json::from_str(&result)?; assert_eq!(json["data"][0]["answer"], 1); + let result = read( + &database.client, + &connection, + "SELECT n AS answer FROM agent_traces_by_key", + ) + .await?; + let json: Value = serde_json::from_str(&result)?; + assert_eq!(json["data"][0]["answer"], 4); + Ok(()) } diff --git a/litellm-rust/crates/traces/tests/migrations.rs b/litellm-rust/crates/traces/tests/migrations.rs index 1fff72fc26d..bf14710cf48 100644 --- a/litellm-rust/crates/traces/tests/migrations.rs +++ b/litellm-rust/crates/traces/tests/migrations.rs @@ -107,7 +107,7 @@ async fn schema_supports_span_rollups_and_spend_joins( #[future(awt)] database: TestResult, ) -> TestResult { let database = database?; - let writer = Connection::writer(&database.url, "default", "")?; + let writer = Connection::writer(&database.url)?; 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; @@ -144,7 +144,7 @@ async fn schema_supports_span_rollups_and_spend_joins( let body = read_json( &database, "SELECT toUInt32(sum(SpanCount)) AS spans, toUInt32(sum(InputTokens)) AS tokens \ - FROM trace_test.agent_traces WHERE TeamId = 'team-1' AND TraceId = 'trace-1'", + FROM trace_test.agent_traces_by_key WHERE TeamId = 'team-1' AND TraceId = 'trace-1'", ) .await?; assert_eq!( @@ -160,7 +160,7 @@ async fn keyed_rollup_keeps_same_trace_ids_separate_by_api_key( #[future(awt)] database: TestResult, ) -> TestResult { let database = database?; - let writer = Connection::writer(&database.url, "default", "")?; + let writer = Connection::writer(&database.url)?; ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?; let timestamp = time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64; let rows = vec![ @@ -204,7 +204,7 @@ async fn rollup_merges_spans_across_days_without_losing_root_fields( #[future(awt)] database: TestResult, ) -> TestResult { let database = database?; - let writer = Connection::writer(&database.url, "default", "")?; + let writer = Connection::writer(&database.url)?; ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?; let day_start = time::OffsetDateTime::now_utc() .replace_time(time::Time::MIDNIGHT) @@ -223,12 +223,16 @@ async fn rollup_merges_spans_across_days_without_losing_root_fields( "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?; + execute_write( + &database, + "OPTIMIZE TABLE trace_test.agent_traces_by_key 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", + FROM trace_test.agent_traces_by_key", ) .await?; assert_eq!( @@ -247,7 +251,7 @@ async fn spend_deduplication_preserves_subsecond_requests_and_retries( #[future(awt)] database: TestResult, ) -> TestResult { let database = database?; - let writer = Connection::writer(&database.url, "default", "")?; + let writer = Connection::writer(&database.url)?; 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; @@ -299,7 +303,7 @@ 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", "")?; + let writer = Connection::writer(&database.url)?; 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; @@ -315,6 +319,7 @@ async fn retention_changes_materialize_existing_rows_and_remain_idempotent( }))?; insert_rows(&database, "otel_traces", vec![span]).await?; insert_rows(&database, "spend_logs", vec![spend]).await?; + assert_eq!(table_rows(&database, "agent_traces_by_key").await?, 1); ensure_schema(&database.client, &writer, "trace_test", 14, 14).await?; let deadline = tokio::time::Instant::now() + Duration::from_secs(60); loop { @@ -337,10 +342,14 @@ async fn retention_changes_materialize_existing_rows_and_remain_idempotent( 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.agent_traces_by_key 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, "agent_traces_by_key").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?; @@ -359,7 +368,7 @@ async fn schema_statement_timeout_maps_to_transport_error() -> TestResult { }); let client = Client::no_redirect_for_test(); let url = format!("http://{address}"); - let writer = Connection::writer(&url, "default", "")?; + let writer = Connection::writer(&url)?; let result = tokio::time::timeout( Duration::from_secs(12), ensure_schema(&client, &writer, "trace_test", 7, 14), diff --git a/litellm/integrations/clickhouse/schema.py b/litellm/integrations/clickhouse/schema.py index 8352d46bb13..6bec35c5630 100644 --- a/litellm/integrations/clickhouse/schema.py +++ b/litellm/integrations/clickhouse/schema.py @@ -3,7 +3,6 @@ from typing import Final from litellm.rust_bridge.traces import TraceStorage OTEL_TRACES_TABLE: Final = "otel_traces" -AGENT_TRACES_TABLE: Final = "agent_traces" AGENT_TRACES_BY_KEY_TABLE: Final = "agent_traces_by_key" SPEND_LOGS_TABLE: Final = "spend_logs"