diff --git a/deploy/lens/README.md b/deploy/lens/README.md index 6f8c7d3940e..60e9c9991a4 100644 --- a/deploy/lens/README.md +++ b/deploy/lens/README.md @@ -178,16 +178,15 @@ Creation queues the first batch. Posting to `/lens/{id}/runs` queues another, or The default is Next.js dev with no production build (`LENS_DEV_BUILD_UI=0`). Set `LENS_DEV_BUILD_UI=1` when you also want a fresh static dashboard at `http://localhost:4000/ui/`. Build output goes to `.lens-dev/logs/ui-build.log`; a failed build stops startup. Both modes keep the live dashboard on port 3000. Startup checks the live login route before seeding and fails with the UI log path if Next.js exits. `LENS_DEV_STARTUP_TIMEOUT_SECONDS` controls startup readiness retries (default 300; `LENS_DEV_READINESS_REQUEST_TIMEOUT_SECONDS` caps each HTTP probe, default 5) -For local fixture data, run `make lens-dev ARGS=--seed`. Use `make lens-dev ARGS="--seed large"` for 2,000 fixture copies, over one million spans and linked request logs. To seed a running stack without restarting it, use `make lens-dev ARGS="--seed-only --seed large --copies 100"`. The default profile replays one copy of every checked-in capture through authenticated `/v1/traces`, including failures, retries, streaming and multiple agent frameworks. Large seeds use the same parser and compressed ClickHouse writer in batches of four copies, and write matching request logs to PostgreSQL. The first and last batches verify linked spend totals through the proxy +For local fixture data, run `make lens-dev ARGS=--seed`. Use `make lens-dev ARGS="--seed large"` for 2,000 fixture copies spread over the last 24 hours, about 860,000 spans with linked request logs, plus three long sessions of roughly 1,150, 9,200 and 92,000 spans in a single trace for drawer paging and the oversized read path. Their trace IDs are printed at the end. To seed a running stack without restarting it, use `make lens-dev ARGS="--seed-only --seed large --copies 100"`. Every profile replays one copy of every checked-in capture through authenticated `/v1/traces`, including failures, retries, streaming and multiple agent frameworks, and verifies linked spend totals through the proxy. Large seeds then copy that first copy inside ClickHouse and PostgreSQL with `INSERT ... SELECT`, rewriting trace, span and call IDs so each copy keeps its own spend, and verify the last copy through the proxy -Seeds append fresh IDs on every invocation and spread copies over recent timestamps. Restarts without `SEED` do not add data. Lens excludes activity received in the last two minutes, so wait two minutes after seeding before checking investigation previews. `LENS_DEV_SEED_COPIES` overrides total copies, and `LENS_DEV_SEED_BATCH_COPIES` overrides copies per bulk insert (default 4, about 2,000 spans). Start with four or fewer on a constrained machine. Larger batches still respect the existing ClickHouse insert size limit; each capture is decoded separately within the OTLP safety budget. Large seeds test data volume and pagination, rather than concurrent ingestion throughput or review accuracy. They can use substantial disk space; adjust `--copies` for your machine. Seeding expects the generated local tracing configuration. The old `run_tracing_proxy_local.sh --seed` command forwards to Lens dev, using its ports and saved master key +Seeds append fresh IDs on every invocation and spread copies over recent timestamps. Restarts without `SEED` do not add data. Lens excludes activity received in the last two minutes, so wait two minutes after seeding before checking investigation previews. `LENS_DEV_SEED_COPIES` overrides total copies. Large seeds test data volume and pagination, rather than concurrent ingestion throughput or review accuracy. They can use substantial disk space; adjust `--copies` for your machine. Seeding expects the generated local tracing configuration. The old `run_tracing_proxy_local.sh --seed` command forwards to Lens dev, using its ports and saved master key Local ingestion limits are explicit and configurable. Set OTLP and ClickHouse variables before starting the proxy and seeder so both processes use the same settings. Invalid, zero and negative values fail instead of silently falling back. Changing these limits does not require rebuilding Rust | Environment variable | Default | Controls | | --- | --- | --- | | `LENS_DEV_SEED_COPIES` | 1 default, 2000 large | Total fixture copies | -| `LENS_DEV_SEED_BATCH_COPIES` | 4 | Copies per bulk insert | | `LENS_DEV_SEED_TIMEOUT_SECONDS` | 120 | Seeder HTTP timeout | | `OTLP_MAX_BODY_BYTES` | 16777216 | HTTP body and decompressed payload bytes | | `OTLP_MAX_CONCURRENT_INGESTS` | 2 | Concurrent proxy ingestion requests | @@ -202,7 +201,7 @@ Local ingestion limits are explicit and configurable. Set OTLP and ClickHouse va | `CLICKHOUSE_TRACE_MAX_INSERT_BYTES` | 67108864 | Encoded trace or spend insert bytes | | `CLICKHOUSE_INSERT_TIMEOUT_SECONDS` | 30 | ClickHouse insert HTTP timeout | -The wire parsers also enforce their library recursion limits (128 levels for JSON, 100 for protobuf). Raising the configured depth does not remove those parser limits. Bulk seeding parses each capture separately, keeping the per-export limits distinct from the bulk insert limit. Use smaller batches if an insert exceeds its byte budget. For example, `LENS_DEV_SEED_COPIES=100 LENS_DEV_SEED_BATCH_COPIES=2 make lens-dev ARGS="--seed large"` +The wire parsers also enforce their library recursion limits (128 levels for JSON, 100 for protobuf). Raising the configured depth does not remove those parser limits. ## Quality evaluation diff --git a/scripts/seed_tracing_fixtures.py b/scripts/seed_tracing_fixtures.py index 3febe49c451..be11f7cc90b 100644 --- a/scripts/seed_tracing_fixtures.py +++ b/scripts/seed_tracing_fixtures.py @@ -43,6 +43,9 @@ TRACE: Final = TypeAdapter(Trace) NANOSECOND_FIELDS: Final = frozenset({"startTimeUnixNano", "endTimeUnixNano", "timeUnixNano"}) TRACE_ID_FIELDS: Final = frozenset({"traceId", "trace_id", "session_id"}) SPAN_ID_FIELDS: Final = frozenset({"spanId", "parentSpanId", "span_id"}) +COPY_WINDOW_MS: Final = 24 * 60 * 60 * 1000 +LONG_SESSION_SOURCE: Final = "openai_agents_swarm" +LONG_SESSION_REPEATS: Final = (50, 400, 4000) class TenantIdentity(BaseModel): @@ -245,7 +248,6 @@ class SeedOptions(BaseModel): model_config = ConfigDict(frozen=True) profile: Literal["default", "large"] copies: int | None - batch_copies: int = 4 timeout_seconds: float = 120 @@ -258,33 +260,23 @@ def seed_arguments(argv: Sequence[str] | None = None) -> SeedOptions: default=os.environ.get("LENS_DEV_SEED_COPIES"), help="Override fixture copies (default: 1, large: 2000; env: LENS_DEV_SEED_COPIES)", ) - parser.add_argument( - "--batch-copies", - type=int, - default=os.environ.get("LENS_DEV_SEED_BATCH_COPIES", "4"), - help="Copies per bulk insert (default: 4; env: LENS_DEV_SEED_BATCH_COPIES)", - ) parser.add_argument("--timeout-seconds", type=float, default=os.environ.get("LENS_DEV_SEED_TIMEOUT_SECONDS", "120")) arguments: Final = SeedOptions.model_validate(vars(parser.parse_args(argv))) if arguments.copies is not None and arguments.copies < 1: parser.error("--copies must be positive") - if arguments.batch_copies < 1: - parser.error("--batch-copies must be positive") if not math.isfinite(arguments.timeout_seconds) or arguments.timeout_seconds <= 0: parser.error("--timeout-seconds must be finite and positive") return arguments -async def seed_batch( +async def seed_copy( client: httpx.AsyncClient, storage: ClickHouseStorage, database: Prisma, replays: tuple[FixtureReplay, ...], fixtures: tuple[tuple[str, tuple[SpendLogRecord, ...]], ...], pattern: re.Pattern[str], - tenant: TenantIdentity | None, - verify: bool, -) -> TenantIdentity: +) -> tuple[tuple[str, tuple[SpendLogRecord, ...]], ...]: by_name: Final = MappingProxyType(dict(fixtures)) paired: Final = tuple( ( @@ -294,36 +286,36 @@ async def seed_batch( for replay in replays if replay.name in by_name ) - rebased_spends: Final = tuple(chain.from_iterable(rows for _, rows in paired)) - resolved_tenant: Final = await ingest_replays( - client, storage, replays, fixture_capture(*next((name, rows[0]) for name, rows in paired)).trace_id, tenant - ) - stamped_spends: Final[tuple[SpendLogRecord, ...]] = tuple( - {**row, "team_id": resolved_tenant.team_id, "api_key": resolved_tenant.api_key, "user": resolved_tenant.user} - for row in rebased_spends + tenant: Final = await ingest_replays( + client, storage, replays, fixture_capture(*next((name, rows[0]) for name, rows in paired)).trace_id ) + stamped: Final = tuple((name, tuple(stamp(row, tenant) for row in rows)) for name, rows in paired) + stamped_spends: Final = tuple(chain.from_iterable(rows for _, rows in stamped)) await storage.insert_rows("spend_logs", stamped_spends) await database.litellm_spendlogs.create_many(data=[postgres_row(row) for row in stamped_spends]) - if verify: - verified: Final = tuple(await asyncio.gather(*(verify_capture(client, name, rows) for name, rows in paired))) - sys.stdout.write(json.dumps({"spend_rows": len(stamped_spends), "captures": verified}, indent=2) + "\n") - if not all(capture["verified"] for capture in verified): - raise RuntimeError("Seed spend verification failed") - return resolved_tenant + return stamped + + +def stamp(row: SpendLogRecord, tenant: TenantIdentity) -> SpendLogRecord: + return {**row, "team_id": tenant.team_id, "api_key": tenant.api_key, "user": tenant.user} + + +async def verify( + client: httpx.AsyncClient, captures: tuple[tuple[str, tuple[SpendLogRecord, ...]], ...], trace_salt: str +) -> None: + verified: Final = tuple( + await asyncio.gather(*(verify_capture(client, name, rows, trace_salt) for name, rows in captures)) + ) + sys.stdout.write( + json.dumps({"spend_rows": sum(len(rows) for _, rows in captures), "captures": verified}, indent=2) + "\n" + ) + if not all(capture["verified"] for capture in verified): + raise RuntimeError("Seed spend verification failed") async def ingest_replays( - client: httpx.AsyncClient, - storage: ClickHouseStorage, - replays: tuple[FixtureReplay, ...], - trace_id: str, - tenant: TenantIdentity | None, + client: httpx.AsyncClient, storage: ClickHouseStorage, replays: tuple[FixtureReplay, ...], trace_id: str ) -> TenantIdentity: - if tenant is not None: - await storage.insert_rows( - "otel_traces", bulk_span_rows(replays, Tenant(tenant.team_id, tenant.api_key, user_id=tenant.user)) - ) - return tenant for replay in replays: ( await client.post( @@ -347,24 +339,161 @@ def bulk_span_rows(replays: tuple[FixtureReplay, ...], tenant: Tenant) -> tuple[ ) -def replay_batches( - count: int, now_ms: int, namespace: str, pattern: re.Pattern[str], batch_copies: int = 4 -) -> Iterator[tuple[int, tuple[FixtureReplay, ...]]]: - for start, stop in ((start, min(start + batch_copies, count)) for start in range(1, count, batch_copies)): - yield ( - stop, - tuple( - chain.from_iterable( - fixture_replays(TRACE_FIXTURES, now_ms - index * 1000, f"{namespace}-{index}", pattern) - for index in range(start, stop) - ) - ), +@dataclass(frozen=True, slots=True) +class Copies: + """Server-side copies of seeded traces and their spend. + + Each copy `n` hashes trace ids with `session` (or `n` when empty), keeps root span ids, hashes the + other span ids with `n`, rewrites seeded call ids from `source` to `{target}{n}-` and moves `n * step_ms` + earlier. A `session` folds every copy into one trace under a single root that spans all of them. + """ + + trace_ids: tuple[str, ...] + request_ids: tuple[str, ...] + numbers: range + step_ms: int + source: str + target: str + session: str = "" + + +def copied_trace_id(trace_id: str, salt: str) -> str: + return hashlib.sha256(f"{trace_id}:{salt}".encode()).hexdigest()[:32] + + +def clickhouse_call_id(column: str) -> str: + target: Final = "concat({target:String}, toString(c.n), '-')" + return ( + f"if(startsWith({column}, 'resp_'), concat('resp_', base64Encode(replaceAll(" + f"tryBase64Decode(substring({column}, 6)), {{source:String}}, {target}))), " + f"replaceAll({column}, {{source:String}}, {target}))" + ) + + +def clickhouse_hash(column: str, salt: str, length: int) -> str: + return f"if({column} = '', '', substring(lower(hex(SHA256(concat({column}, ':', {salt})))), 1, {length}))" + + +def clickhouse_copy_sql(database: str) -> tuple[str, str]: + trace_salt: Final = "if({session:String} = '', toString(c.n), {session:String})" + roots: Final = ( + f"(SELECT SpanId FROM {database}.otel_traces " + "WHERE TraceId IN {trace_ids:Array(String)} AND ParentSpanId = '')" + ) + folded_root: Final = "{session:String} != '' AND t.ParentSpanId = ''" + shift: Final = f"toIntervalMillisecond(if({folded_root}, {{last:UInt64}}, c.n) * {{step_ms:UInt64}})" + numbers: Final = "CROSS JOIN (SELECT number AS n FROM numbers({first:UInt64}, {count:UInt64})) AS c" + spans: Final = f"""INSERT INTO {database}.otel_traces +SELECT t.* REPLACE ( + t.Timestamp - {shift} AS Timestamp, + {clickhouse_hash("t.TraceId", trace_salt, 32)} AS TraceId, + if(t.ParentSpanId = '', t.SpanId, {clickhouse_hash("t.SpanId", "toString(c.n)", 16)}) AS SpanId, + if(t.ParentSpanId IN {roots}, t.ParentSpanId, {clickhouse_hash("t.ParentSpanId", "toString(c.n)", 16)}) + AS ParentSpanId, + t.Duration + if({folded_root}, {{last:UInt64}} * {{step_ms:UInt64}} * 1000000, 0) AS Duration, + arrayMap(at -> at - {shift}, t.`Events.Timestamp`) AS `Events.Timestamp`, + mapApply((name, value) -> (name, {clickhouse_call_id("value")}), t.SpanAttributes) AS SpanAttributes, + {clickhouse_call_id("t.LiteLLMRequestId")} AS LiteLLMRequestId, + arrayMap(key -> concat(extract(key, '^[^:]*:'), {clickhouse_call_id("replaceRegexpOne(key, '^[^:]*:', '')")}), + t.CallKeys) AS CallKeys +) +FROM {database}.otel_traces AS t {numbers} +WHERE t.TraceId IN {{trace_ids:Array(String)}} + AND ({{session:String}} = '' OR t.ParentSpanId != '' OR c.n = {{first:UInt64}})""" + spend: Final = f"""INSERT INTO {database}.spend_logs +SELECT s.* REPLACE ( + {clickhouse_call_id("s.request_id")} AS request_id, + {clickhouse_call_id("s.response_id")} AS response_id, + {clickhouse_call_id("s.litellm_call_id")} AS litellm_call_id, + {clickhouse_hash("s.trace_id", trace_salt, 32)} AS trace_id, + {clickhouse_hash("s.session_id", trace_salt, 32)} AS session_id, + if(s.span_id IN {roots}, s.span_id, {clickhouse_hash("s.span_id", "toString(c.n)", 16)}) AS span_id, + s.start_time - toIntervalMillisecond(c.n * {{step_ms:UInt64}}) AS start_time, + s.end_time - toIntervalMillisecond(c.n * {{step_ms:UInt64}}) AS end_time, + s.completion_start_time - toIntervalMillisecond(c.n * {{step_ms:UInt64}}) AS completion_start_time +) +FROM {database}.spend_logs AS s {numbers} +WHERE s.request_id IN {{request_ids:Array(String)}}""" + return spans, spend + + +def clickhouse_array(values: tuple[str, ...]) -> str: + return "[" + ",".join("'" + value.replace("\\", "\\\\").replace("'", "\\'") + "'" for value in values) + "]" + + +async def copy_clickhouse(client: httpx.AsyncClient, database: str, copies: Copies) -> None: + parameters: Final = { + "param_trace_ids": clickhouse_array(copies.trace_ids), + "param_request_ids": clickhouse_array(copies.request_ids), + "param_first": str(copies.numbers.start), + "param_count": str(len(copies.numbers)), + "param_last": str(copies.numbers.stop - 1), + "param_step_ms": str(copies.step_ms), + "param_source": copies.source, + "param_target": copies.target, + "param_session": copies.session, + } + for sql in clickhouse_copy_sql(database): + (await client.post("/", params=parameters, content=sql)).raise_for_status() + + +POSTGRES_COPY_TARGET: Final = "($2 || c.n || '-')" +POSTGRES_COPY_SHIFT: Final = "make_interval(secs => c.n * $6::bigint / 1000.0)" +POSTGRES_COPY_SQL: Final = f"""INSERT INTO "LiteLLM_SpendLogs" +SELECT (jsonb_populate_record(s, jsonb_build_object( + 'request_id', CASE WHEN left(s.request_id, 5) = 'resp_' + THEN 'resp_' || translate(encode(convert_to(replace(convert_from(decode(substr(s.request_id, 6), 'base64'), + 'UTF8'), $1, {POSTGRES_COPY_TARGET}), 'UTF8'), 'base64'), E'\\n', '') + ELSE replace(s.request_id, $1, {POSTGRES_COPY_TARGET}) END, + 'session_id', CASE WHEN coalesce(s.session_id, '') = '' THEN s.session_id + ELSE substr(encode(sha256(convert_to( + s.session_id || ':' || CASE WHEN $3 = '' THEN c.n::text ELSE $3 END, 'UTF8')), 'hex'), 1, 32) END, + 'startTime', s."startTime" - {POSTGRES_COPY_SHIFT}, + 'endTime', s."endTime" - {POSTGRES_COPY_SHIFT}, + 'completionStartTime', s."completionStartTime" - {POSTGRES_COPY_SHIFT} +))).* +FROM "LiteLLM_SpendLogs" AS s CROSS JOIN generate_series($4::int, $5::int) AS c(n) +WHERE s.request_id = ANY(string_to_array($7, E'\\n'))""" + + +async def copy_postgres(database: Prisma, copies: Copies) -> None: + await database.execute_raw( + POSTGRES_COPY_SQL, + copies.source, + copies.target, + copies.session, + copies.numbers.start, + copies.numbers.stop - 1, + copies.step_ms, + "\n".join(copies.request_ids), + ) + + +def long_sessions( + replays: tuple[FixtureReplay, ...], + captures: tuple[tuple[str, tuple[SpendLogRecord, ...]], ...], + source: str, + target: str, + repeats: tuple[int, ...] = LONG_SESSION_REPEATS, +) -> tuple[Copies, ...]: + replay: Final = next(replay for replay in replays if replay.name == LONG_SESSION_SOURCE) + rows: Final = dict(captures)[LONG_SESSION_SOURCE] + span_ns: Final = tuple(timestamps(replay.export)) + return tuple( + Copies( + trace_ids=(fixture_capture(LONG_SESSION_SOURCE, rows[0]).trace_id,), + request_ids=tuple(row["request_id"] for row in rows), + numbers=range(count), + step_ms=(max(span_ns) - min(span_ns)) // 1_000_000 + 1000, + source=source, + target=f"{target}s{count}x", + session=f"session{count}", ) + for count in repeats + ) -async def seed( - profile: str = "default", copies: int | None = None, batch_copies: int = 4, timeout_seconds: float = 120 -) -> int: +async def seed(profile: str = "default", copies: int | None = None, timeout_seconds: float = 120) -> int: from prisma import Prisma fixtures: Final = spend_fixtures() @@ -372,29 +501,43 @@ async def seed( count: Final = copies if copies is not None else (2000 if profile == "large" else 1) namespace: Final = uuid4().hex now_ms: Final = time.time_ns() // 1_000_000 - storage: Final = ClickHouseStorage(trace_storage_config({})) + config: Final = trace_storage_config({}) + storage: Final = ClickHouseStorage(config) + replays: Final = fixture_replays(TRACE_FIXTURES, now_ms, namespace + "-0", pattern) + source: Final = f"seed-{namespace}-0-" + target: Final = f"seed-{namespace}-" async with ( httpx.AsyncClient( base_url=os.environ.get("PROXY_BASE_URL", "http://127.0.0.1:4002"), headers={"Authorization": f"Bearer {os.environ['LITELLM_MASTER_KEY']}"}, timeout=timeout_seconds, ) as client, - Prisma() as database, + httpx.AsyncClient(base_url=config.url, params={"database": config.database}, timeout=600) as clickhouse, + Prisma(http={"timeout": httpx.Timeout(600)}) as database, ): - tenant: Final = await seed_batch( - client, - storage, - database, - fixture_replays(TRACE_FIXTURES, now_ms, namespace + "-0", pattern), - fixtures, - pattern, - None, - True, + captures: Final = await seed_copy(client, storage, database, replays, fixtures, pattern) + await verify(client, captures, "") + repeated: Final = Copies( + trace_ids=tuple( + sorted(frozenset(str(span["TraceId"]) for span in bulk_span_rows(replays, Tenant("", "")))) + ), + request_ids=tuple(row["request_id"] for _, rows in captures for row in rows), + numbers=range(1, count), + step_ms=COPY_WINDOW_MS // count, + source=source, + target=target, ) - for stop, replays in replay_batches(count, now_ms, namespace, pattern, batch_copies): - await seed_batch(client, storage, database, replays, fixtures, pattern, tenant, stop == count) - sys.stdout.write(f"Seeded {stop}/{count} fixture copies\n") - sys.stdout.flush() + sessions: Final = long_sessions(replays, captures, source, target) if profile == "large" else () + for plan in (repeated, *sessions) if count > 1 else sessions: + await copy_clickhouse(clickhouse, config.database, plan) + await copy_postgres(database, plan) + if count > 1: + await verify(client, captures, str(count - 1)) + for plan in sessions: + sys.stdout.write( + f"Long session: {len(plan.numbers)} repeats, " + f"trace_id={copied_trace_id(plan.trace_ids[0], plan.session)}\n" + ) sys.stdout.write(f"Seed complete: profile={profile}, copies={count}, namespace={namespace}\n") return 0 @@ -410,17 +553,18 @@ def fixture_capture(name: str, row: SpendLogRecord) -> FixtureCapture: async def verify_capture( - client: httpx.AsyncClient, name: str, rows: tuple[SpendLogRecord, ...] + client: httpx.AsyncClient, name: str, rows: tuple[SpendLogRecord, ...], trace_salt: str = "" ) -> Mapping[str, JsonValue]: capture: Final = fixture_capture(name, rows[0]) - detail: Final = await client.get(f"/v1/traces/{capture.trace_id}") + trace_id: Final = copied_trace_id(capture.trace_id, trace_salt) if trace_salt else capture.trace_id + detail: Final = await client.get(f"/v1/traces/{trace_id}") detail.raise_for_status() trace: Final = TRACE.validate_json(detail.content) expected: Final = sum(row["spend"] or 0 for row in rows) actual: Final = trace["summary"]["spend"] return { "fixture": name, - "trace_id": capture.trace_id, + "trace_id": trace_id, "spend_rows": len(rows), "recorded_spend": expected, "trace_spend": actual, @@ -432,6 +576,4 @@ async def verify_capture( if __name__ == "__main__": arguments: Final = seed_arguments() - raise SystemExit( - asyncio.run(seed(arguments.profile, arguments.copies, arguments.batch_copies, arguments.timeout_seconds)) - ) + raise SystemExit(asyncio.run(seed(arguments.profile, arguments.copies, arguments.timeout_seconds))) diff --git a/tests/test_litellm_rust/test_traces.py b/tests/test_litellm_rust/test_traces.py index c1847a52760..b8af0f0c374 100644 --- a/tests/test_litellm_rust/test_traces.py +++ b/tests/test_litellm_rust/test_traces.py @@ -4,13 +4,15 @@ import json import math import re import time -from collections.abc import Iterator +from collections.abc import Generator, Iterator +from contextlib import closing from dataclasses import dataclass from itertools import chain from types import MappingProxyType from typing import Final from urllib.parse import parse_qs, urlsplit +import httpx import pytest from fastapi import FastAPI from fastapi.testclient import TestClient @@ -19,16 +21,21 @@ from pydantic import BaseModel, ConfigDict, JsonValue, TypeAdapter from litellm.constants import OTLP_MAX_ATTRIBUTE_VALUE_BYTES from litellm.rust_bridge._native import NativeTraceConfig, NativeTraceStorage from litellm.rust_bridge.trace.generated.models import ActivityAvailability, LensAccessParams, TraceQueryHelp -from litellm.rust_bridge.trace.generated.types import TraceScope +from litellm.rust_bridge.trace.generated.types import Trace, TraceScope from litellm.rust_bridge.trace.storage import ClickHouseStorage, TraceStorageConfig, span_rows from litellm.tracing import Tenant, TraceReceiver, TracingPayloadTooLargeError from litellm.tracing.types import SpendLogRecord from scripts.seed_tracing_fixtures import ( TRACE, TRACE_FIXTURES, + Copies, FixtureReplay, + bulk_span_rows, + copied_trace_id, + copy_clickhouse, fixture_capture, fixture_replays, + long_sessions, rebase_spend, response_pattern, spend_fixtures, @@ -503,7 +510,7 @@ def seeded_trace_api(clickhouse_url: str) -> Iterator[SeededTraceAPI]: def _fixture_trace_api( clickhouse_url: str, replays: tuple[FixtureReplay, ...], stamped: tuple[SpendLogRecord, ...] -) -> Iterator[SeededTraceAPI]: +) -> Generator[SeededTraceAPI]: from litellm.proxy._types import LitellmUserRoles, UserAPIKeyAuth from litellm.proxy.auth.user_api_key_auth import user_api_key_auth from litellm.proxy.tracing_endpoints import provide_receiver, provide_trace_query_secret, router @@ -595,24 +602,34 @@ def test_query_correlation_requires_key_or_user_ownership_within_a_team(seeded_t assert all(row["request_id"] != unrelated["request_id"] for row in matches) -@pytest.fixture(scope="module") -def captured_trace_api() -> Iterator[SeededTraceAPI]: +def _captured_replays( + namespace: str, +) -> tuple[tuple[FixtureReplay, ...], tuple[tuple[str, tuple[SpendLogRecord, ...]], ...]]: captures: Final = spend_fixtures() - originals: Final = tuple(chain.from_iterable(rows for _, rows in captures)) - pattern: Final = response_pattern(originals) - replays: Final = fixture_replays(TRACE_FIXTURES, time.time_ns() // 1_000_000, "captured-api", pattern) + pattern: Final = response_pattern(tuple(chain.from_iterable(rows for _, rows in captures))) + replays: Final = fixture_replays(TRACE_FIXTURES, time.time_ns() // 1_000_000, namespace, pattern) by_name: Final = MappingProxyType(dict(captures)) - paired: Final = tuple( - rebase_spend(by_name[replay.name], replay.offset_ms, replay.namespace, pattern) + return replays, tuple( + ( + replay.name, + tuple( + _stamp(row) for row in rebase_spend(by_name[replay.name], replay.offset_ms, replay.namespace, pattern) + ), + ) for replay in replays if replay.name in by_name ) - stamped: Final[tuple[SpendLogRecord, ...]] = tuple( - {**row, "team_id": "team-a", "api_key": "fixture-key", "user": "fixture-user"} - for row in chain.from_iterable(paired) - ) + + +def _stamp(row: SpendLogRecord) -> SpendLogRecord: + return {**row, "team_id": "team-a", "api_key": "fixture-key", "user": "fixture-user"} + + +@pytest.fixture(scope="module") +def captured_trace_api() -> Iterator[SeededTraceAPI]: + replays, paired = _captured_replays("captured-api") with clickhouse_service() as url: - yield from _fixture_trace_api(url, replays, stamped) + yield from _fixture_trace_api(url, replays, tuple(chain.from_iterable(rows for _, rows in paired))) @pytest.mark.parametrize("name", tuple(name for name, _ in spend_fixtures())) @@ -644,3 +661,64 @@ def test_captured_sdk_cost_survives_seeding_and_is_queryable(name: str, captured assert math.isclose(sum(row.spend for row in records), sum(row["spend"] or 0 for row in rows)) assert sum(row.prompt_tokens for row in records) == sum(row["prompt_tokens"] for row in rows) assert sum(row.completion_tokens for row in records) == sum(row["completion_tokens"] for row in rows) + + +def test_server_side_copies_keep_every_capture_linked_to_its_spend() -> None: + replays, paired = _captured_replays("copied-api") + copies: Final = Copies( + trace_ids=tuple(sorted(frozenset(str(span["TraceId"]) for span in bulk_span_rows(replays, Tenant("", ""))))), + request_ids=tuple(row["request_id"] for _, rows in paired for row in rows), + numbers=range(1, 3), + step_ms=60_000, + source="seed-copied-api-", + target="seed-copied-api-c", + ) + (session,) = long_sessions(replays, paired, "seed-copied-api-", "seed-copied-api-c", (3,)) + session_spend: Final = sum(row["spend"] or 0 for row in dict(paired)["openai_agents_swarm"]) + session_spans: Final = len( + span_rows((TRACE_FIXTURES / "openai_agents_swarm.json").read_bytes(), "application/json") + ) + with ( + clickhouse_service() as url, + closing(_fixture_trace_api(url, replays, tuple(chain.from_iterable(rows for _, rows in paired)))) as seeded, + ): + api: Final = next(seeded) + assert api.client.portal is not None + for plan in (copies, session): + api.client.portal.call(_copy_clickhouse, url, plan) + for name, rows in paired: + _assert_capture(api, name, rows, fixture_capture(name, rows[0]).trace_id) + _assert_capture(api, name, rows, copied_trace_id(fixture_capture(name, rows[0]).trace_id, "2")) + trace: Final = _trace(api, copied_trace_id(session.trace_ids[0], session.session)) + assert trace["summary"]["span_count"] == 1 + 3 * (session_spans - 1) + (root,) = (span for span in trace["spans"] if not span["parent_span_id"]) + assert {span["parent_span_id"] for span in trace["spans"] if span["parent_span_id"]} <= { + span["span_id"] for span in trace["spans"] + } + assert root["start_offset_ms"] == min(span["start_offset_ms"] for span in trace["spans"]) + assert root["start_offset_ms"] + root["duration_ms"] >= max( + span["start_offset_ms"] + span["duration_ms"] for span in trace["spans"] + ) + assert trace["summary"]["spend"] == pytest.approx(3 * session_spend) + + +def _trace(api: SeededTraceAPI, trace_id: str) -> Trace: + response: Final = api.client.get(f"/v1/traces/{trace_id}") + assert response.status_code == 200, response.text + return TRACE.validate_json(response.content) + + +def _assert_capture(api: SeededTraceAPI, name: str, rows: tuple[SpendLogRecord, ...], trace_id: str) -> None: + capture: Final = fixture_capture(name, rows[0]) + summary: Final = _trace(api, trace_id)["summary"] + assert summary["span_count"] == len(span_rows((TRACE_FIXTURES / f"{name}.json").read_bytes(), "application/json")) + assert summary["spend"] == ( + pytest.approx(sum(row["spend"] or 0 for row in rows)) + if capture.spend_linked and capture.spend_complete + else None + ) + + +async def _copy_clickhouse(url: str, copies: Copies) -> None: + async with httpx.AsyncClient(base_url=url, params={"database": "trace_test"}) as client: + await copy_clickhouse(client, "trace_test", copies) diff --git a/tests/unit/test_lens_dev.py b/tests/unit/test_lens_dev.py index af61734731e..bc73695f898 100644 --- a/tests/unit/test_lens_dev.py +++ b/tests/unit/test_lens_dev.py @@ -194,12 +194,11 @@ def test_seed_only_with_no_cli_count_preserves_env_controls(tmp_path: Path) -> N proc = _run( tmp_path, "parse_args --seed-only; master_key=sk-local; py() { " - 'printf "%s %s %s\\n" "$LENS_DEV_SEED_COPIES" "$LENS_DEV_SEED_BATCH_COPIES" "$@"; }; py=py; seed_data', + 'printf "%s %s\\n" "$LENS_DEV_SEED_COPIES" "$@"; }; py=py; seed_data', LENS_DEV_SEED_COPIES="3", - LENS_DEV_SEED_BATCH_COPIES="1", ) assert proc.returncode == 0, proc.stderr - assert proc.stdout.startswith("3 1 -m") + assert proc.stdout.startswith("3 -m") def test_proxy_uses_this_checkouts_ui_build(tmp_path: Path) -> None: diff --git a/tests/unit/test_seed_tracing_fixtures.py b/tests/unit/test_seed_tracing_fixtures.py index 30cfde9feec..8dac769b3c1 100644 --- a/tests/unit/test_seed_tracing_fixtures.py +++ b/tests/unit/test_seed_tracing_fixtures.py @@ -11,23 +11,22 @@ import pytest from prisma import Json, Prisma from pydantic import InstanceOf, TypeAdapter +from litellm.rust_bridge.trace.queries import TraceSQLResponse 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_copy, seed_id, spend_fixtures, timestamps, @@ -196,29 +195,39 @@ def test_bulk_export_preserves_all_spans_and_disjoint_copy_ids() -> None: @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 +async def test_first_copy_stamps_the_authenticated_tenant_and_writes_both_stores( + monkeypatch: pytest.MonkeyPatch, +) -> None: + from litellm.rust_bridge.trace.storage import ClickHouseStorage + monkeypatch.setenv("LITELLM_MASTER_KEY", "sk-local") 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) + replays: Final = fixture_replays(TRACE_FIXTURES, 1_800_000_000_000, "first", pattern) storage: Final = AsyncMock(spec=ClickHouseStorage) + storage.query_sql.return_value = TraceSQLResponse.model_validate( + { + "meta": (), + "data": [{"team_id": "local-team", "api_key": "local-hash", "user": "admin"}], + "rows": 1, + "statistics": {"elapsed": 0, "rows_read": 1, "bytes_read": 1}, + } + ) 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" + client: Final = AsyncMock(spec=httpx.AsyncClient) + client.post.return_value = httpx.Response(200, request=httpx.Request("POST", "http://proxy/v1/traces")) + captures: Final = await seed_copy(client, storage, database, replays, fixtures, pattern) + assert tuple(JSON.validate_json(call.kwargs["content"]) for call in client.post.call_args_list) == tuple( + replay.export for replay in replays + ) + rows: Final = tuple(chain.from_iterable(rows for _, rows in captures)) + assert {name for name, _ in captures} == {name for name, _ in fixtures} + assert storage.insert_rows.call_args.args == ("spend_logs", rows) 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: @@ -228,23 +237,6 @@ def test_seed_cli_rejects_nonpositive_copies() -> None: 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: