diff --git a/litellm-rust/crates/pagination/src/cursor.rs b/litellm-rust/crates/pagination/src/cursor.rs index 7ccff2c2455..b997523a42e 100644 --- a/litellm-rust/crates/pagination/src/cursor.rs +++ b/litellm-rust/crates/pagination/src/cursor.rs @@ -41,14 +41,15 @@ impl KeyRing { }) }) .collect::, _>>()?; - if keys.is_empty() { - return Err(Error::InvalidKeys); - } Ok(Self { keys }) } - pub fn signing_key_id(&self) -> &str { - &self.keys[0].id + pub fn signing_key_id(&self) -> Option<&str> { + self.keys.first().map(|key| key.id.as_str()) + } + + fn signing_key(&self) -> Result<&Key, Error> { + self.keys.first().ok_or(Error::MissingKeys) } fn sign(key: &Key, binding: &Binding, payload: &[u8]) -> Vec { @@ -66,15 +67,16 @@ impl KeyRing { binding: &Binding, cursor: &Cursor

, ) -> Result { + let signing_key = self.signing_key()?; let payload = serde_json::to_vec(&Envelope { version: FORMAT_VERSION, - key_id: &self.keys[0].id, + key_id: &signing_key.id, expires_unix: cursor.expires_at.unix_timestamp(), published_ms: cursor.published_ms, revision: &cursor.revision, position: &cursor.position, })?; - let tag = Self::sign(&self.keys[0], binding, &payload); + let tag = Self::sign(signing_key, binding, &payload); Ok(format!( "{}.{}", URL_SAFE_NO_PAD.encode(payload), @@ -89,6 +91,7 @@ impl KeyRing { token: &str, now: OffsetDateTime, ) -> Result, Error> { + self.signing_key()?; let (payload, tag) = token.split_once('.').ok_or(Error::InvalidCursor)?; let payload = URL_SAFE_NO_PAD .decode(payload) diff --git a/litellm-rust/crates/pagination/src/error.rs b/litellm-rust/crates/pagination/src/error.rs index e713563ecf8..ce72a543b30 100644 --- a/litellm-rust/crates/pagination/src/error.rs +++ b/litellm-rust/crates/pagination/src/error.rs @@ -71,6 +71,8 @@ pub enum Error { ResourceTooLarge, #[error("Cursor signing keys must be nonempty")] InvalidKeys, + #[error("Cursor signing keys are not configured")] + MissingKeys, #[error("Cursor encoding failed")] Encode(#[from] serde_json::Error), } @@ -82,7 +84,7 @@ impl Error { Self::TraversalExpired => FailureCode::TraversalExpired, Self::TraversalChanged => FailureCode::TraversalChanged, Self::ResourceTooLarge => FailureCode::ResourceTooLarge, - Self::InvalidKeys | Self::Encode(_) => FailureCode::Unavailable, + Self::InvalidKeys | Self::MissingKeys | Self::Encode(_) => FailureCode::Unavailable, } } } diff --git a/litellm-rust/crates/pagination/tests/cursor.rs b/litellm-rust/crates/pagination/tests/cursor.rs index 541295bd792..e8618826004 100644 --- a/litellm-rust/crates/pagination/tests/cursor.rs +++ b/litellm-rust/crates/pagination/tests/cursor.rs @@ -131,12 +131,31 @@ fn rotated_keys_keep_verifying_earlier_cursors(binding: Binding, now: OffsetDate } #[rstest] -#[case::no_keys(Vec::<&str>::new())] #[case::empty_key(vec![""])] +#[case::empty_key_among_valid(vec!["valid", ""])] fn key_rings_require_nonempty_secrets(#[case] secrets: Vec<&str>) { assert!(matches!(KeyRing::new(secrets), Err(Error::InvalidKeys))); } +#[rstest] +fn a_key_ring_without_keys_refuses_to_sign_or_verify( + keys: KeyRing, + binding: Binding, + now: OffsetDateTime, +) { + let unsigned = KeyRing::new(Vec::<&str>::new()).unwrap(); + let token = keys.encode(&binding, &cursor(now)).unwrap(); + assert_eq!(unsigned.signing_key_id(), None); + assert!(matches!( + unsigned.encode(&binding, &cursor(now)), + Err(Error::MissingKeys) + )); + assert!(matches!( + unsigned.decode::(&binding, &token, now), + Err(Error::MissingKeys) + )); +} + #[rstest] fn traversal_identity_follows_binding_and_publication(binding: Binding, now: OffsetDateTime) { let first = Traversal::new(&binding, 1_790_000_000_000, now); @@ -190,6 +209,7 @@ fn bounded_pages_keep_a_complete_prefix_and_reissue_the_cursor( #[case::changed(Error::TraversalChanged, FailureCode::TraversalChanged, false, true)] #[case::too_large(Error::ResourceTooLarge, FailureCode::ResourceTooLarge, false, false)] #[case::keys(Error::InvalidKeys, FailureCode::Unavailable, true, false)] +#[case::missing_keys(Error::MissingKeys, FailureCode::Unavailable, true, false)] fn errors_map_to_stable_codes( #[case] error: Error, #[case] code: FailureCode, diff --git a/litellm-rust/crates/traces-clickhouse/migrations/0009_trace_rollup_ownership.sql b/litellm-rust/crates/traces-clickhouse/migrations/0009_trace_rollup_ownership.sql index f4445980378..fd696349cf5 100644 --- a/litellm-rust/crates/traces-clickhouse/migrations/0009_trace_rollup_ownership.sql +++ b/litellm-rust/crates/traces-clickhouse/migrations/0009_trace_rollup_ownership.sql @@ -1,4 +1,3 @@ ALTER TABLE {database}.agent_traces_by_key ADD COLUMN IF NOT EXISTS UserIds SimpleAggregateFunction(groupUniqArrayArray, Array(String)) DEFAULT [], - ADD COLUMN IF NOT EXISTS IdentifiedLlmCount SimpleAggregateFunction(sum, UInt64) DEFAULT 0, - ADD COLUMN IF NOT EXISTS ReceivedMs SimpleAggregateFunction(min, UInt64) DEFAULT 0 + ADD COLUMN IF NOT EXISTS IdentifiedLlmCount SimpleAggregateFunction(sum, UInt64) DEFAULT 0 diff --git a/litellm-rust/crates/traces-clickhouse/migrations/0010_trace_cost_completeness.sql b/litellm-rust/crates/traces-clickhouse/migrations/0010_trace_cost_completeness.sql index df70a50b348..af87a6bf40b 100644 --- a/litellm-rust/crates/traces-clickhouse/migrations/0010_trace_cost_completeness.sql +++ b/litellm-rust/crates/traces-clickhouse/migrations/0010_trace_cost_completeness.sql @@ -15,7 +15,6 @@ SELECT countIf(StatusCode = 'STATUS_CODE_ERROR') AS ErrorCount, sum(InputTokens) AS InputTokens, sum(OutputTokens) AS OutputTokens, - min(EngineReceivedMs) AS ReceivedMs, groupUniqArrayIf(toString(Model), Model != '') AS Models, groupUniqArrayIf(SpanName, ObservationType = 'agent') AS AgentNames, groupArrayIf(LiteLLMRequestId, ObservationType = 'llm' OR LiteLLMRequestId != '') AS RequestIds diff --git a/litellm-rust/crates/traces-clickhouse/migrations/0015_agent_traces_received.sql b/litellm-rust/crates/traces-clickhouse/migrations/0015_agent_traces_received.sql new file mode 100644 index 00000000000..4620479afec --- /dev/null +++ b/litellm-rust/crates/traces-clickhouse/migrations/0015_agent_traces_received.sql @@ -0,0 +1,2 @@ +ALTER TABLE {database}.agent_traces_by_key + ADD COLUMN IF NOT EXISTS ReceivedMs SimpleAggregateFunction(min, UInt64) DEFAULT 0 diff --git a/litellm-rust/crates/traces-clickhouse/migrations/0016_agent_traces_mv_received.sql b/litellm-rust/crates/traces-clickhouse/migrations/0016_agent_traces_mv_received.sql new file mode 100644 index 00000000000..df70a50b348 --- /dev/null +++ b/litellm-rust/crates/traces-clickhouse/migrations/0016_agent_traces_mv_received.sql @@ -0,0 +1,23 @@ +ALTER TABLE {database}.agent_traces_by_key_mv MODIFY QUERY +SELECT + TeamId, ApiKeyHash, TraceId, groupUniqArray(UserId) AS UserIds, + 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 = 'llm' AND LiteLLMRequestId != '') AS IdentifiedLlmCount, + countIf(ObservationType = 'tool') AS ToolCount, + countIf(StatusCode = 'STATUS_CODE_ERROR') AS ErrorCount, + sum(InputTokens) AS InputTokens, + sum(OutputTokens) AS OutputTokens, + min(EngineReceivedMs) AS ReceivedMs, + groupUniqArrayIf(toString(Model), Model != '') AS Models, + groupUniqArrayIf(SpanName, ObservationType = 'agent') AS AgentNames, + groupArrayIf(LiteLLMRequestId, ObservationType = 'llm' OR LiteLLMRequestId != '') AS RequestIds +FROM {database}.otel_traces +GROUP BY TeamId, ApiKeyHash, TraceId diff --git a/litellm/rust_bridge/_native.pyi b/litellm/rust_bridge/_native.pyi index c146a6eac92..2affbd56e52 100644 --- a/litellm/rust_bridge/_native.pyi +++ b/litellm/rust_bridge/_native.pyi @@ -26,6 +26,8 @@ def trace_span_rows( body: bytes, content_type: str | None, tenant: Mapping[str, str], max_attribute_value_bytes: int ) -> list[dict[str, JsonValue]]: ... +class TraceReadError(Exception): ... + @final class NativeTraceConfig: def __new__( @@ -34,6 +36,7 @@ class NativeTraceConfig: url: str, retention_days: int, max_attribute_value_bytes: int, + cursor_keys: Sequence[str], ) -> NativeTraceConfig: ... @final @@ -46,7 +49,7 @@ class NativeTraceStorage: self, scope: TraceScope, start_ms: int, end_ms: int, cursor: str | None, limit: int ) -> Future[JsonValue]: ... def get_trace( - self, trace_id: str, scope: TraceScope, trace_ref: str, cursor: str | None = None, page_size: int | None = None + self, trace_id: str, scope: TraceScope, trace_ref: str, cursor: str | None, page_size: int | None ) -> Future[JsonValue]: ... def get_span(self, trace_id: str, span_id: str, scope: TraceScope, trace_ref: str) -> Future[JsonValue]: ... def get_span_error( @@ -358,6 +361,7 @@ __all__ = [ "RustUpstreamError", "TokenCounter", "Tokenizer", + "TraceReadError", "achat_completions", "acompletion", "aembedding", diff --git a/litellm/tracing/config.py b/litellm/tracing/config.py index ef537ba2e67..eeb54411636 100644 --- a/litellm/tracing/config.py +++ b/litellm/tracing/config.py @@ -58,9 +58,7 @@ def _cursor_keys(store: Mapping[str, object], environ: Mapping[str, str]) -> tup supplied: Final = _value(store, "cursor_keys", environ, None) if supplied is None: fallback: Final = environ.get("LITELLM_SALT_KEY") or environ.get("LITELLM_MASTER_KEY") - if not fallback: - raise ValueError("tracing.store.cursor_keys, LITELLM_SALT_KEY or LITELLM_MASTER_KEY is required") - return (fallback,) + return (fallback,) if fallback else () try: candidates: Final = CURSOR_KEYS.validate_python((supplied,) if isinstance(supplied, str) else supplied) except ValidationError as error: @@ -108,3 +106,10 @@ def trace_storage_config(settings: Mapping[str, object], environ: Mapping[str, s ), cursor_keys=_cursor_keys(store, environ), ) + + +def trace_reader_config(settings: Mapping[str, object], environ: Mapping[str, str] = os.environ) -> TraceStorageConfig: + config: Final = trace_storage_config(settings, environ) + if not config.cursor_keys: + raise ValueError("tracing.store.cursor_keys, LITELLM_SALT_KEY or LITELLM_MASTER_KEY is required") + return config diff --git a/litellm/tracing/receiver.py b/litellm/tracing/receiver.py index 0c659aa94e4..36b3892b3a3 100644 --- a/litellm/tracing/receiver.py +++ b/litellm/tracing/receiver.py @@ -26,7 +26,7 @@ from litellm.constants import ( ) from litellm.rust_bridge.trace.generated.types import SpanDetail, SpanErrorPage, Trace, TracePage, TraceScope from litellm.rust_bridge.trace.storage import ClickHouseStorage, Tenant -from litellm.tracing.config import trace_storage_config +from litellm.tracing.config import trace_reader_config from litellm.tracing.otlp_http import InvalidOTLPPayloadError, TracingPayloadTooLargeError, decompress @@ -55,7 +55,7 @@ class TraceReceiver: @classmethod def from_settings(cls, settings: Mapping[str, object]) -> "TraceReceiver": - return cls(storage=ClickHouseStorage(trace_storage_config(settings))) + return cls(storage=ClickHouseStorage(trace_reader_config(settings))) async def start(self) -> None: await self.storage.ensure_schema() diff --git a/tests/unit/tracing/test_config.py b/tests/unit/tracing/test_config.py index 363a6d6db67..465b95c112d 100644 --- a/tests/unit/tracing/test_config.py +++ b/tests/unit/tracing/test_config.py @@ -1,7 +1,7 @@ import pytest from litellm import constants -from litellm.tracing.config import is_clickhouse_tracing_enabled, trace_storage_config +from litellm.tracing.config import is_clickhouse_tracing_enabled, trace_reader_config, trace_storage_config @pytest.mark.parametrize( @@ -124,7 +124,11 @@ def test_legacy_reader_and_split_retention_fields_are_rejected() -> None: @pytest.mark.parametrize( ("store_keys", "environ", "expected"), [ - (["os.environ/NEW_KEY", "literal-old"], {"NEW_KEY": "rotated", "LITELLM_SALT_KEY": "salt"}, ("rotated", "literal-old")), + ( + ["os.environ/NEW_KEY", "literal-old"], + {"NEW_KEY": "rotated", "LITELLM_SALT_KEY": "salt"}, + ("rotated", "literal-old"), + ), ("single-key", {"LITELLM_SALT_KEY": "salt"}, ("single-key",)), (None, {"LITELLM_SALT_KEY": "salt", "LITELLM_MASTER_KEY": "sk-master"}, ("salt",)), (None, {"LITELLM_MASTER_KEY": "sk-master"}, ("sk-master",)), @@ -146,6 +150,10 @@ def test_invalid_cursor_keys_are_rejected(store_keys: object) -> None: trace_storage_config({"store": store}, {"LITELLM_SALT_KEY": "salt"}) -def test_missing_cursor_key_source_is_rejected() -> None: +def test_missing_cursor_key_source_is_rejected_for_readers() -> None: with pytest.raises(ValueError, match="cursor_keys, LITELLM_SALT_KEY or LITELLM_MASTER_KEY is required"): - trace_storage_config({"store": {"type": "clickhouse", "url": "http://localhost:8123"}}, {}) + trace_reader_config({"store": {"type": "clickhouse", "url": "http://localhost:8123"}}, {}) + + +def test_write_only_storage_does_not_need_cursor_keys() -> None: + assert trace_storage_config({"store": {"type": "clickhouse", "url": "http://localhost:8123"}}, {}).cursor_keys == () diff --git a/ui/litellm-dashboard/src/components/lens/LensWorkspace.integration.test.tsx b/ui/litellm-dashboard/src/components/lens/LensWorkspace.integration.test.tsx index 87adbe2adb3..07851eb5d8c 100644 --- a/ui/litellm-dashboard/src/components/lens/LensWorkspace.integration.test.tsx +++ b/ui/litellm-dashboard/src/components/lens/LensWorkspace.integration.test.tsx @@ -95,7 +95,13 @@ describe("Lens interactive demo", () => { const path = new URL(String(input), "http://localhost").pathname; if (path === "/lens") return Response.json({ lenses: [saved], workers: [], tracing_enabled: true }); if (path.endsWith("/runs")) return Response.json(saved.jobs); - if (path === "/v1/traces") return Response.json({ data: data.runs.map((run) => run.trace.summary) }); + if (path === "/v1/traces") { + return Response.json({ + items: data.runs.map((run) => run.trace.summary), + next_cursor: null, + traversal: { id: "demo", published_at: "2026-01-01T00:00:00Z", expires_at: "2026-01-01T00:30:00Z" }, + }); + } return Response.json({ data: [], traces: true, requests: false }); }); renderWithProviders(, {