fix(lens): seed through isolated ingestion and preserve upstream queue fixes

This commit is contained in:
moe-berri 2026-10-07 12:42:17 -07:00
parent 9a989fb0cb
commit 22e94fb29a
6 changed files with 60 additions and 41 deletions

View file

@ -194,26 +194,7 @@ For local fixture data, run `make lens-dev ARGS=--seed`. Use `make lens-dev ARGS
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_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 |
| `OTLP_MAX_ATTRIBUTE_VALUE_BYTES` | 65536 | Stored attribute/content bytes |
| `OTLP_MAX_DECODE_DEPTH` | 32 | Nested decode depth |
| `OTLP_MAX_DECODE_NODES` | 65536 | JSON values or protobuf fields per export |
| `OTLP_MAX_SPANS` | 4096 | Spans per export |
| `OTLP_MAX_ATTRIBUTES` | 256 | Attributes per resource, scope, span, event or link |
| `OTLP_MAX_EVENTS` | 256 | Events per span |
| `OTLP_MAX_LINKS` | 256 | Links per span |
| `OTLP_MAX_DECODED_SPAN_BYTES` | 16777216 | Decoded span allocation budget |
| `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.
The Rust receiver bounds each upload and its decompressed body to 16 MiB and permits two ingestion requests at once per replica. Exporters should split large batches and retry backpressure. `LENS_DEV_SEED_COPIES` and `LENS_DEV_SEED_TIMEOUT_SECONDS` control the seeder; the receiver's limits are compiled into the service
## Quality evaluation
@ -252,7 +233,7 @@ The hourly development pipeline pins all component images to the same selected c
## Worker dependencies
The worker uses the same digest-pinned Wolfi base and Python version as the component images. Python dependencies and their hashes are locked in `deploy/lens/requirements.lock`. To update them, edit `deploy/lens/requirements.in`, then run `uv pip compile --universal --python-version 3.13 --generate-hashes --no-emit-index-url deploy/lens/requirements.in -o deploy/lens/requirements.lock`. The image installs only the locked wheels with hash verification. CI builds and scans both native architectures
The service builds from the workspace Cargo.lock with a pinned Rust toolchain and a digest-pinned Wolfi runtime. It has no Python package dependencies. CPython and libseccomp support the confined calculation tool. CI builds, runs, and scans native amd64 and arm64 images
## Python analysis boundary
@ -262,14 +243,14 @@ The native worker image builds a syscall policy with libseccomp and includes the
Python execution requires a native Linux worker with Landlock ABI 3 or later and seccomp filtering. Build the image for the host architecture. Missing policy files, an incompatible kernel, or an unsupported host such as a macOS source worker returns a clear tool error. There is no unrestricted execution fallback. Keep the container's non-root user, dropped capabilities, no-new-privileges setting, read-only root and writable temporary mount
The worker permits two Python children at once across all investigations. Set `LENS_PYTHON_CONCURRENCY` to a positive integer to change this worker-wide pool. Queued calls consume no child process or scratch directory; cancelling a queued call does not start it. Model, read and search concurrency are separate
The worker permits two Python children at once across all investigations. Queued calls consume no child process or scratch directory; cancelling a queued call does not start it. Model, read and search concurrency are separate
| Per-call resource | Default |
| --- | --- |
| Elapsed execution time | 60 seconds |
| CPU time | 30 seconds |
| Process address space | 512 MiB |
| Captured stdout or stderr | 8 MiB per stream |
| Captured stdout or stderr | 4 MiB per stream |
| Individual scratch file size | 16 MiB |
| Monitored scratch storage | 64 MiB |
| Monitored scratch entries | 2,048 |
@ -283,12 +264,11 @@ Results include `stdout`, `stderr`, `exit_code`, `error` and `output_complete`.
This is a process boundary sharing the worker's Linux kernel. The checked-in smoke test verifies useful Python operations, filesystem and process restrictions, raw syscall attempts, resource failures, mapping accounting, cleanup and cancellation in the actual image. Run it on the deployment's native architecture and kernel:
```bash
docker build --build-arg LITELLM_RELEASE_TAG=lens-python-test \
-f deploy/lens/Dockerfile -t lens-worker:python-test .
docker run --rm --pull never --read-only --cap-drop ALL \
docker build --target smoke --build-arg LITELLM_RELEASE_TAG=lens-python-test \
-f deploy/lens/Dockerfile -t lens-worker:smoke .
docker run --rm --read-only --cap-drop ALL \
--security-opt no-new-privileges --network none \
--tmpfs /tmp:rw,noexec,nosuid,size=1g --entrypoint python -i \
lens-worker:python-test - < tests/proxy_behavior/lens/worker_python_smoke.py
--tmpfs /tmp:rw,noexec,nosuid,size=1g lens-worker:smoke
```
The same checks can run through pytest by setting `LENS_TEST_WORKER_IMAGE` to an already-built native image. The worker image CI runs the standalone smoke without adding pytest to the production image
The smoke target runs the Rust sandbox integration tests. The production image contains neither Cargo nor the test executable

View file

@ -78,7 +78,10 @@ pub fn response(content_type: Option<&str>, outcome: Result<(), Error>) -> Respo
)
};
let mut response = (status, [(http::header::CONTENT_TYPE, media_type)], body).into_response();
if matches!(status, StatusCode::SERVICE_UNAVAILABLE | StatusCode::TOO_MANY_REQUESTS) {
if matches!(
status,
StatusCode::SERVICE_UNAVAILABLE | StatusCode::TOO_MANY_REQUESTS
) {
response
.headers_mut()
.insert("retry-after", http::HeaderValue::from_static("5"));

View file

@ -1,7 +1,7 @@
import asyncio
import json
import random
from collections.abc import AsyncGenerator, Awaitable, Callable
from collections.abc import AsyncGenerator, AsyncIterator, Awaitable, Callable
from contextlib import AbstractAsyncContextManager, asynccontextmanager
from dataclasses import dataclass
from datetime import datetime, timedelta, timezone

View file

@ -5,10 +5,10 @@ from collections.abc import Iterable, Mapping, Sequence
from contextlib import suppress
from io import BytesIO
from typing import Final
from typing_extensions import TypeIs
import httpx
from pydantic import TypeAdapter, ValidationError
from typing_extensions import TypeIs
from litellm._logging import verbose_proxy_logger
from litellm.integrations.clickhouse.clickhouse_spend_logger import spend_log_row_from_payload
@ -23,11 +23,15 @@ SHUTDOWN_SECONDS: Final = 3.0
_PAYLOAD: Final = TypeAdapter(SpendLogPayload)
def _is_mapping(value: object) -> TypeIs[Mapping[object, object]]:
def _is_mapping(
value: object,
) -> TypeIs[Mapping[object, object]]: # guard-ok: bounds arbitrary callback mappings before validation
return isinstance(value, Mapping)
def _is_sequence(value: object) -> TypeIs[Sequence[object]]:
def _is_sequence(
value: object,
) -> TypeIs[Sequence[object]]: # guard-ok: bounds arbitrary callback sequences before validation
return isinstance(value, (tuple, list))
@ -144,7 +148,7 @@ class LensExporter(CustomLogger):
size = 2 # rebind-ok: count bytes in a bounded batch without copying records
records: Final[deque[bytes]] = deque() # mutable-ok: finite batch drained from the queue
while self.queue and size + len(self.queue[0]) + 1 <= MAX_BATCH_BYTES:
record = self.queue.popleft() # rebind-ok: drain each record into the bounded batch
record: Final = self.queue.popleft()
size += len(record) + 1
records.append(record)
return tuple(records)

View file

@ -176,6 +176,16 @@ wait_for_proxy() {
die "proxy not ready after ${startup_timeout}s; see $log_dir/proxy.log"
}
wait_for_lens() {
local lens_pid="$1"
for _ in $(seq 1 "$startup_timeout"); do
kill -0 "$lens_pid" 2>/dev/null || die "Lens exited; see $log_dir/worker.log"
curl -fsS --max-time "$readiness_request_timeout" "http://127.0.0.1:$lens_port/health/ready" >/dev/null 2>&1 && return
sleep 1
done
die "Lens not ready after ${startup_timeout}s; see $log_dir/worker.log"
}
wait_for_ui() {
local ui_pid="$1"
echo "lens-dev: waiting for the UI (log: $log_dir/ui.log)"
@ -275,7 +285,7 @@ parse_args() {
}
main() {
local config_file exports proxy_pid ui_pid pid key_hint
local config_file exports proxy_pid ui_pid lens_pid pid key_hint
parse_args "$@"
if [ -n "${LENS_DEV_CONFIG:-}" ]; then
[ -f "$LENS_DEV_CONFIG" ] || die "LENS_DEV_CONFIG not found: $LENS_DEV_CONFIG"
@ -353,15 +363,16 @@ main() {
wait_for_ui "$ui_pid"
wait_for_proxy "$proxy_pid"
ensure_worker_token
if [ -n "$seed_profile" ] || [ -n "$seed_logs_profile" ]; then seed_data; fi
LITELLM_RELEASE_TAG="$source_release_tag" \
LITELLM_MODE=PRODUCTION LITELLM_URL="$proxy_url" LENS_WORKER_TOKEN="$(cat "$token_file")" \
LITELLM_LENS_SERVICE_TOKEN="$service_key" LITELLM_LENS_LISTEN="127.0.0.1:$lens_port" \
CLICKHOUSE_URL="$clickhouse_url" CLICKHOUSE_DATABASE=litellm \
"$repo_root/litellm-rust/target/debug/litellm-lens" \
< /dev/null > "$log_dir/worker.log" 2>&1 &
pids+=("$!")
lens_pid=$!
pids+=("$lens_pid")
wait_for_lens "$lens_pid"
if [ -n "$seed_profile" ] || [ -n "$seed_logs_profile" ]; then seed_data; fi
key_hint="password in $key_file"
[ -z "${LENS_DEV_MASTER_KEY:-}" ] || key_hint="password from LENS_DEV_MASTER_KEY"

View file

@ -11,7 +11,8 @@ import os
import re
import sys
import time
from collections.abc import Iterator, Mapping, Sequence
from collections.abc import AsyncIterator, Iterator, Mapping, Sequence
from contextlib import asynccontextmanager
from dataclasses import dataclass
from datetime import datetime, timezone
from functools import cache
@ -24,6 +25,7 @@ from uuid import uuid4
import httpx
from pydantic import BaseModel, ConfigDict, JsonValue, TypeAdapter
from litellm.proxy.lens.ingestion import IngestionKeyCreated
from litellm.rust_bridge.trace.generated.types import AllQueryScope, Trace
from litellm.rust_bridge.trace.storage import ClickHouseStorage, Tenant, span_rows
from litellm.tracing.config import trace_storage_config
@ -528,6 +530,24 @@ def long_sessions(
)
@asynccontextmanager
async def ingestion_client(client: httpx.AsyncClient, timeout_seconds: float) -> AsyncIterator[httpx.AsyncClient]:
response: Final = await client.post("/lens/tracing/keys", json={"name": "Local fixture seed"})
response.raise_for_status()
created: Final = IngestionKeyCreated.model_validate_json(response.content)
try:
if not created.active:
raise RuntimeError("Lens ingestion is not ready; start the Lens service before seeding")
async with httpx.AsyncClient(
base_url=os.environ["LITELLM_LENS_URL"],
headers={"Authorization": f"Bearer {created.key}"},
timeout=timeout_seconds,
) as uploader:
yield uploader
finally:
(await client.delete(f"/lens/tracing/keys/{created.record.id}")).raise_for_status()
async def seed(profile: str = "default", copies: int | None = None, timeout_seconds: float = 120) -> int:
from prisma import Prisma
@ -550,7 +570,8 @@ async def seed(profile: str = "default", copies: int | None = None, timeout_seco
httpx.AsyncClient(base_url=config.url, params={"database": config.database}, timeout=600) as clickhouse,
Prisma(http={"timeout": httpx.Timeout(600)}) as database,
):
captures: Final = await seed_copy(client, storage, database, replays, fixtures, pattern)
async with ingestion_client(client, timeout_seconds) as uploader:
captures: Final = await seed_copy(uploader, storage, database, replays, fixtures, pattern)
await verify(client, captures, "")
repeated: Final = Copies(
trace_ids=tuple(