fix(tracing): append ReceivedMs migrations, make cursor keys reader-only, sync native stub

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
Yujong Lee 2026-10-03 23:14:58 +00:00
parent 97ded17960
commit 1c12e4cf73
12 changed files with 94 additions and 23 deletions

View file

@ -41,14 +41,15 @@ impl KeyRing {
})
})
.collect::<Result<Vec<_>, _>>()?;
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<u8> {
@ -66,15 +67,16 @@ impl KeyRing {
binding: &Binding,
cursor: &Cursor<P>,
) -> Result<String, Error> {
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<Cursor<P>, Error> {
self.signing_key()?;
let (payload, tag) = token.split_once('.').ok_or(Error::InvalidCursor)?;
let payload = URL_SAFE_NO_PAD
.decode(payload)

View file

@ -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,
}
}
}

View file

@ -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::<Position>(&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,

View file

@ -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

View file

@ -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

View file

@ -0,0 +1,2 @@
ALTER TABLE {database}.agent_traces_by_key
ADD COLUMN IF NOT EXISTS ReceivedMs SimpleAggregateFunction(min, UInt64) DEFAULT 0

View file

@ -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

View file

@ -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",

View file

@ -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

View file

@ -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()

View file

@ -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 == ()

View file

@ -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(<LensWorkspace accessToken="live-token" userRole="Admin" readOnly={false} />, {