mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-11 03:38:38 +00:00
* feat: seed Lens dev with configurable load profiles * feat: run Lens UI live through the dev launcher * fix: verify Lens UI startup before seeding * fix(ui): render run timestamps on one compact line The agent runs table printed the long locale form with timezone, which wrapped to two lines per row. Use a fixed-width 24h form with milliseconds that matches the timeline axis, keep the long form in the hover title, and show the timezone once in the column header Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * refactor(ui): use the root query client for the Lens demo Drop the demo's nested QueryClient. Cache keys are already partitioned by scope, and the root client now skips retries on 4xx ApiErrors, which covers the demo's not-in-demo and read-only rejections and live 4xx alike. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * chore: drop unused synthetic_spend reference from query help Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * chore(traces): drop the deeplite fixtures and map every SDK to a logo The deeplite captures predate the example repository and carried synthetic spend rows, which leaked a fixture-only column into the query help SQL and pinned tests to its shape. Replace them with the SDK captures in the ClickHouse round trip and query API tests, and derive the seed tenant lookup from the capture metadata The runs table only knew the two Anthropic framework slugs. Register the slugs the normalizer emits for LangChain, LangGraph, Deep Agents, CrewAI, Google ADK, LlamaIndex, OpenAI Agents, Pydantic AI, Strands, Vercel AI SDK, Codex, Cursor and Copilot, and fall back to the generic agent glyph when a run has no known framework Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * refactor(ui): inject the Lens sample through data sources instead of demo checks Components no longer ask whether they are in the demo. The traces source carries live and handoff, navigation state comes from a URL or memory route, and the preview action comes from context instead of onDemo props Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * feat(ui): infinite scroll for the agent runs list Replace the Load more button with a sentinel that fetches the next cursor page as the list nears its end. Placeholder rows hold the tail while more runs exist, and a failed page stops auto-loading until Retry. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * feat(ui): keep the whole Lens view in the URL Lens navigation now lives entirely in query params through nuqs: the sample session (demo=true, with a Demo data switch in the header), the open run (trace, trace_ref), the selected step, view and detail section (span, view, span_tab) and the list filters and range (q, agent, status, hours). Any Lens view is a shareable link and the back button walks runs RunView takes its selection injected: the drawer feeds it URL state and the investigations evidence sheet keeps a local one, so a finding's original run never writes step ids into the URL. Leaving the sample session clears every Lens key except the tab so sample ids never point at live data Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * style(traces): cargo fmt captures tests Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: Yujong Lee <yujong@berri.ai> Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
252 lines
12 KiB
Python
252 lines
12 KiB
Python
import json
|
|
import re
|
|
from datetime import datetime
|
|
from itertools import chain
|
|
from pathlib import Path
|
|
from typing import Final
|
|
from unittest.mock import AsyncMock
|
|
|
|
import httpx
|
|
import pytest
|
|
from prisma import Json, Prisma
|
|
from pydantic import InstanceOf, TypeAdapter
|
|
|
|
from litellm.rust_bridge.trace.storage import Tenant, span_rows
|
|
from litellm.tracing.types import SpendLogRecord
|
|
from scripts.seed_tracing_fixtures import (
|
|
JSON,
|
|
TRACE_FIXTURES,
|
|
TenantIdentity,
|
|
bulk_span_rows,
|
|
fixture_capture,
|
|
fixture_replays,
|
|
postgres_row,
|
|
rebase,
|
|
rebase_spend,
|
|
replay_batches,
|
|
response_ids,
|
|
response_pattern,
|
|
seed_arguments,
|
|
seed_batch,
|
|
seed_id,
|
|
spend_fixtures,
|
|
timestamps,
|
|
)
|
|
|
|
CALL_KEYS: Final = TypeAdapter(tuple[str, ...])
|
|
DATETIMES: Final = TypeAdapter(tuple[datetime, datetime])
|
|
SPAN_IDENTITY: Final = TypeAdapter(tuple[str, str, str, int])
|
|
JSON_FIELDS: Final[TypeAdapter[tuple[Json, Json, Json]]] = TypeAdapter(
|
|
tuple[InstanceOf[Json], InstanceOf[Json], InstanceOf[Json]]
|
|
)
|
|
|
|
|
|
@pytest.mark.requires_rust_extension
|
|
@pytest.mark.parametrize(
|
|
"path",
|
|
sorted(TRACE_FIXTURES.glob("*.json")),
|
|
ids=tuple(path.stem for path in sorted(TRACE_FIXTURES.glob("*.json"))),
|
|
)
|
|
def test_all_fixture_replays_are_recent_and_preserve_spans(path: Path) -> None:
|
|
export: Final = JSON.validate_json(path.read_bytes())
|
|
now_ms: Final = max(timestamps(export)) // 1_000_000 + 86_400_000
|
|
replays: Final = fixture_replays(TRACE_FIXTURES, now_ms, "all-fixtures", re.compile(r"(?!)"))
|
|
replay: Final = next(item for item in replays if item.name == path.stem)
|
|
original: Final = span_rows(path.read_bytes(), "application/json")
|
|
replayed: Final = span_rows(json.dumps(replay.export).encode(), "application/json")
|
|
group: Final = tuple(item for item in replays if item.namespace == replay.namespace)
|
|
|
|
assert max(max(timestamps(item.export)) for item in group) // 1_000_000 == now_ms - 1000
|
|
assert len(frozenset(item.offset_ms for item in group)) == 1
|
|
assert tuple(timestamps(replay.export)) == tuple(
|
|
timestamp + replay.offset_ms * 1_000_000 for timestamp in timestamps(export)
|
|
)
|
|
for before, after in zip(original, replayed, strict=True):
|
|
trace_id, span_id, parent_id, timestamp = SPAN_IDENTITY.validate_python(
|
|
(before["TraceId"], before["SpanId"], before["ParentSpanId"], before["Timestamp"])
|
|
)
|
|
assert after["TraceId"] == seed_id(trace_id, replay.namespace, 32)
|
|
assert after["SpanId"] == seed_id(span_id, replay.namespace, 16)
|
|
assert after["ParentSpanId"] == seed_id(parent_id, replay.namespace, 16)
|
|
assert after["Timestamp"] == timestamp + replay.offset_ms * 1_000_000
|
|
assert (after["Duration"], after["InputTokens"], after["OutputTokens"], after["StatusCode"]) == (
|
|
before["Duration"],
|
|
before["InputTokens"],
|
|
before["OutputTokens"],
|
|
before["StatusCode"],
|
|
)
|
|
if path.stem.startswith("query_"):
|
|
assert all(item.namespace == replay.namespace for item in replays if item.name.startswith("query_"))
|
|
else:
|
|
assert all(item.namespace != replay.namespace for item in replays if item.name != path.stem)
|
|
|
|
|
|
@pytest.mark.requires_rust_extension
|
|
def test_replay_preserves_trace_topology_usage_and_event_timing() -> None:
|
|
export: Final = JSON.validate_json((TRACE_FIXTURES / "deepagents_swarm.json").read_bytes())
|
|
original: Final = span_rows(json.dumps(export).encode(), "application/json")
|
|
spend_rows: Final = dict(spend_fixtures())["deepagents_swarm"]
|
|
pattern: Final = re.compile("|".join(re.escape(row["response_id"]) for row in spend_rows))
|
|
shifted: Final = rebase(export, 123_000_000, "first-run", pattern)
|
|
replayed: Final = span_rows(json.dumps(shifted).encode(), "application/json")
|
|
other_run: Final = span_rows(
|
|
json.dumps(rebase(export, 123_000_000, "second-run", pattern)).encode(), "application/json"
|
|
)
|
|
span_ids: Final = {before["SpanId"]: after["SpanId"] for before, after in zip(original, replayed, strict=True)}
|
|
|
|
assert tuple(timestamps(shifted)) == tuple(timestamp + 123_000_000 for timestamp in timestamps(export))
|
|
assert {span["TraceId"] for span in original}.isdisjoint(span["TraceId"] for span in replayed)
|
|
assert {span["TraceId"] for span in replayed}.isdisjoint(span["TraceId"] for span in other_run)
|
|
for before, after in zip(original, replayed, strict=True):
|
|
assert after["ParentSpanId"] == span_ids.get(before["ParentSpanId"], "")
|
|
assert after["Timestamp"] == before["Timestamp"] + 123_000_000
|
|
assert after["Duration"] == before["Duration"]
|
|
assert after["InputTokens"] == before["InputTokens"]
|
|
assert after["OutputTokens"] == before["OutputTokens"]
|
|
assert after["StatusCode"] == before["StatusCode"]
|
|
assert after["LiteLLMRequestId"] == (
|
|
f"seed-first-run-{before['LiteLLMRequestId']}" if before["LiteLLMRequestId"] else ""
|
|
)
|
|
|
|
|
|
def test_postgres_rows_preserve_clickhouse_cost_identity_and_payloads() -> None:
|
|
spends: Final = dict(spend_fixtures())["deepagents_swarm"]
|
|
|
|
for spend, postgres in ((spend, postgres_row(spend)) for spend in spends):
|
|
start_time, end_time = DATETIMES.validate_python((postgres["startTime"], postgres["endTime"]))
|
|
messages, response, proxy_request = JSON_FIELDS.validate_python(
|
|
(postgres["messages"], postgres["response"], postgres["proxy_server_request"])
|
|
)
|
|
assert postgres["request_id"] == spend["response_id"]
|
|
assert (postgres["api_key"], postgres["team_id"], postgres["user"], postgres["session_id"]) == (
|
|
spend["api_key"],
|
|
spend["team_id"],
|
|
spend["user"],
|
|
spend["session_id"],
|
|
)
|
|
assert postgres["spend"] == spend["spend"]
|
|
assert postgres["total_tokens"] == spend["prompt_tokens"] + spend["completion_tokens"]
|
|
assert round(start_time.timestamp() * 1000) == spend["start_time"]
|
|
assert round(end_time.timestamp() * 1000) == spend["end_time"]
|
|
assert postgres["request_duration_ms"] == spend["end_time"] - spend["start_time"]
|
|
assert JSON.validate_python(getattr(messages, "data")) == JSON.validate_json(spend["messages"])
|
|
assert JSON.validate_python(getattr(response, "data")) == JSON.validate_json(spend["response"])
|
|
assert JSON.validate_python(getattr(proxy_request, "data")) is None
|
|
|
|
|
|
@pytest.mark.requires_rust_extension
|
|
@pytest.mark.parametrize("name,spends", spend_fixtures())
|
|
def test_captured_spend_replay_preserves_real_cost_and_call_identity(
|
|
name: str, spends: tuple[SpendLogRecord, ...]
|
|
) -> None:
|
|
export: Final = JSON.validate_json((TRACE_FIXTURES / f"{name}.json").read_bytes())
|
|
pattern: Final = response_pattern(spends)
|
|
offset_ms: Final = 1123
|
|
namespace: Final = f"captured-{name}"
|
|
shifted: Final = rebase(export, offset_ms * 1_000_000, namespace, pattern)
|
|
spans: Final = span_rows(json.dumps(shifted).encode(), "application/json")
|
|
replayed: Final = rebase_spend(spends, offset_ms, namespace, pattern)
|
|
keys: Final = frozenset(chain.from_iterable(CALL_KEYS.validate_python(span["CallKeys"]) for span in spans))
|
|
capture: Final = fixture_capture(name, replayed[0])
|
|
|
|
assert capture.trace_id in frozenset(span["TraceId"] for span in spans)
|
|
for before, after in zip(spends, replayed, strict=True):
|
|
assert after["spend"] == before["spend"]
|
|
assert (after["prompt_tokens"], after["completion_tokens"], after["total_tokens"]) == (
|
|
before["prompt_tokens"],
|
|
before["completion_tokens"],
|
|
before["total_tokens"],
|
|
)
|
|
assert after["request_id"] != before["request_id"]
|
|
assert after["start_time"] == before["start_time"] + offset_ms
|
|
assert after["end_time"] == before["end_time"] + offset_ms
|
|
if before["litellm_call_id"]:
|
|
assert after["litellm_call_id"] != before["litellm_call_id"]
|
|
identities: Final = frozenset(f"provider_response:{identity}" for identity in response_ids((after,))) | {
|
|
f"litellm_request:{after['litellm_call_id']}"
|
|
}
|
|
assert bool(identities & keys) is capture.spend_linked
|
|
|
|
|
|
@pytest.mark.parametrize("call_id", (None, "gateway"))
|
|
def test_spend_fixture_loading_preserves_gateway_ids_and_defaults_legacy_rows(
|
|
tmp_path: Path, call_id: str | None
|
|
) -> None:
|
|
original: Final = dict(spend_fixtures())["deepagents_swarm"][0]
|
|
fields: Final = {key: value for key, value in original.items() if key != "litellm_call_id"}
|
|
supplied: Final = fields if call_id is None else {**fields, "litellm_call_id": call_id}
|
|
(tmp_path / "example_spend_logs.jsonl").write_text(json.dumps(supplied) + "\n")
|
|
loaded: Final = spend_fixtures(tmp_path)
|
|
assert loaded == (("example", ({**original, "litellm_call_id": call_id or ""},)),)
|
|
|
|
|
|
@pytest.mark.requires_rust_extension
|
|
def test_bulk_export_preserves_all_spans_and_disjoint_copy_ids() -> None:
|
|
first: Final = fixture_replays(TRACE_FIXTURES, 1_800_000_000_000, "copy-1", re.compile(r"(?!)"))
|
|
second: Final = fixture_replays(TRACE_FIXTURES, 1_800_000_001_000, "copy-2", re.compile(r"(?!)"))
|
|
merged: Final = bulk_span_rows(first + second, Tenant("", ""))
|
|
separate: Final = tuple(
|
|
span for replay in first + second for span in span_rows(json.dumps(replay.export).encode(), "application/json")
|
|
)
|
|
assert tuple(merged) == separate
|
|
first_ids: Final = frozenset(span["TraceId"] for span in bulk_span_rows(first, Tenant("", "")))
|
|
second_ids: Final = frozenset(span["TraceId"] for span in bulk_span_rows(second, Tenant("", "")))
|
|
assert first_ids.isdisjoint(second_ids)
|
|
|
|
|
|
@pytest.mark.requires_rust_extension
|
|
@pytest.mark.asyncio
|
|
async def test_bulk_seed_stamps_authenticated_tenant_and_writes_both_stores() -> None:
|
|
from litellm.rust_bridge.trace.storage import ClickHouseStorage, Tenant
|
|
|
|
fixtures: Final = spend_fixtures()
|
|
pattern: Final = response_pattern(tuple(chain.from_iterable(rows for _, rows in fixtures)))
|
|
replays: Final = fixture_replays(TRACE_FIXTURES, 1_800_000_000_000, "bulk", pattern)
|
|
storage: Final = AsyncMock(spec=ClickHouseStorage)
|
|
database: Final = AsyncMock(spec=Prisma, litellm_spendlogs=AsyncMock())
|
|
tenant: Final = TenantIdentity(team_id="local-team", api_key="local-hash", user="admin")
|
|
async with httpx.AsyncClient() as client:
|
|
result: Final = await seed_batch(client, storage, database, replays, fixtures, pattern, tenant, False)
|
|
assert result == tenant
|
|
trace_table, trace_rows = storage.insert_rows.call_args_list[0].args
|
|
assert trace_table == "otel_traces"
|
|
assert trace_rows == bulk_span_rows(replays, Tenant("local-team", "local-hash", user_id="admin"))
|
|
table, rows = storage.insert_rows.call_args_list[1].args
|
|
assert table == "spend_logs"
|
|
assert len(rows) == sum(len(original) for _, original in fixtures)
|
|
assert all((row["team_id"], row["api_key"], row["user"]) == ("local-team", "local-hash", "admin") for row in rows)
|
|
saved: Final = database.litellm_spendlogs.create_many.call_args.kwargs["data"]
|
|
assert tuple(row["request_id"] for row in saved) == tuple(row["request_id"] for row in rows)
|
|
assert tuple(row["spend"] for row in saved) == tuple(row["spend"] for row in rows)
|
|
storage.query_sql.assert_not_called()
|
|
|
|
|
|
def test_seed_cli_rejects_nonpositive_copies() -> None:
|
|
with pytest.raises(SystemExit) as error:
|
|
seed_arguments(["--copies", "0"])
|
|
assert error.value.code == 2
|
|
assert seed_arguments(["--profile", "large", "--copies", "5"]).copies == 5
|
|
|
|
|
|
def test_bulk_batches_cover_every_copy_including_partial_tail() -> None:
|
|
batches: Final = tuple(replay_batches(8, 1_800_000_000_000, "batch", re.compile(r"(?!)"), 3))
|
|
assert tuple(stop for stop, _ in batches) == (4, 7, 8)
|
|
expected: Final = tuple(
|
|
tuple(fixture_replays(TRACE_FIXTURES, 1_800_000_000_000 - index * 1000, f"batch-{index}", re.compile(r"(?!)")))
|
|
for index in range(1, 8)
|
|
)
|
|
assert tuple(chain.from_iterable(replays for _, replays in batches)) == tuple(chain.from_iterable(expected))
|
|
|
|
|
|
def test_seed_cli_rejects_nonpositive_batch_size() -> None:
|
|
with pytest.raises(SystemExit) as error:
|
|
seed_arguments(["--batch-copies", "0"])
|
|
assert error.value.code == 2
|
|
assert seed_arguments(["--batch-copies", "2"]).batch_copies == 2
|
|
|
|
|
|
@pytest.mark.parametrize("timeout", ("0", "-1", "inf", "nan"))
|
|
def test_seed_cli_rejects_invalid_http_timeouts(timeout: str) -> None:
|
|
with pytest.raises(SystemExit) as error:
|
|
seed_arguments(["--timeout-seconds", timeout])
|
|
assert error.value.code == 2
|