From fc7707f16837a5e14ee59ce299702bdfe2e89021 Mon Sep 17 00:00:00 2001 From: moe-berri Date: Wed, 7 Oct 2026 12:19:52 -0700 Subject: [PATCH] fix(lens): complete service routing and reject stale investigation results --- .github/workflows/lens-worker.yml | 5 +- deploy/lens/compose.yaml | 1 + deploy/lens/smoke.sh | 40 +++++++ deploy/lens/stack.yaml | 7 +- docker/docker-compose.tracing.yml | 30 ++++- docker/tracing-config.yaml | 4 +- helm/litellm-helm/templates/ingress.yaml | 14 +++ helm/litellm-helm/templates/lens/service.yaml | 4 + helm/litellm-helm/values.yaml | 1 + helm/litellm/templates/_helpers.tpl | 14 +++ .../litellm/templates/backend/deployment.yaml | 12 +- .../litellm/templates/gateway/deployment.yaml | 1 + helm/litellm/templates/ingress.yaml | 10 ++ helm/litellm/templates/lens/service.yaml | 4 + helm/litellm/tests/lens_worker_tests.yaml | 40 +++++-- helm/litellm/values.yaml | 1 + litellm-rust/Cargo.lock | 1 + litellm-rust/crates/lens/Cargo.toml | 1 + litellm-rust/crates/lens/src/lib.rs | 2 + litellm-rust/crates/lens/tests/clickhouse.rs | 47 ++++++-- litellm-rust/crates/lens/tests/receiver.rs | 80 +++++++++++++ litellm/proxy/_lazy_openapi_snapshot.json | 10 ++ litellm/proxy/lens/endpoints.py | 53 +++++++-- litellm/proxy/lens/repository.py | 12 +- litellm/tracing/exporter.py | 52 +++++---- litellm/tracing/remote.py | 14 +-- scripts/lens_dev.sh | 37 ++++-- .../database/test_lens_repository.py | 67 +++++++++++ tests/unit/proxy/lens/test_endpoints.py | 51 ++++++++ .../lens/LensSetup.integration.test.tsx | 11 +- .../lens/LensWorkspace.integration.test.tsx | 29 +++-- .../TracingSetupCard.integration.test.tsx | 109 +++++++++++------- .../onboarding/tracing/TracingSetupCard.tsx | 20 +++- .../AnalysisKeyPicker.integration.test.tsx | 1 - .../lens/settings/worker/WorkerForm.tsx | 17 +-- .../lens/settings/worker/WorkerInstall.tsx | 83 ++----------- .../WorkerSettings.integration.test.tsx | 41 +++---- .../lens/settings/worker/WorkerSettings.tsx | 11 +- .../lens/settings/worker/usePrepareWorker.ts | 6 +- .../lens/settings/worker/workerCommand.ts | 22 ---- .../lens/settings/worker/workerSchema.test.ts | 38 +----- .../lens/settings/worker/workerSchema.ts | 24 ---- ui/litellm-dashboard/src/lib/http/client.ts | 5 +- ui/litellm-dashboard/src/lib/http/schema.d.ts | 12 +- 44 files changed, 666 insertions(+), 378 deletions(-) create mode 100644 deploy/lens/smoke.sh delete mode 100644 ui/litellm-dashboard/src/components/lens/settings/worker/workerCommand.ts diff --git a/.github/workflows/lens-worker.yml b/.github/workflows/lens-worker.yml index 4dc7bcda02f..d440dd6e5a7 100644 --- a/.github/workflows/lens-worker.yml +++ b/.github/workflows/lens-worker.yml @@ -52,10 +52,7 @@ jobs: - name: Verify version and unprivileged runtime env: RELEASE_TAG: sha-${{ github.sha }} - run: | - version=$(docker run --rm --network none --read-only --cap-drop ALL --security-opt no-new-privileges lens-worker --version) - test "$version" = "litellm-lens $RELEASE_TAG protocol=7" - test "$(docker run --rm --network none --read-only --entrypoint id lens-worker -u)" = 65532 + run: bash deploy/lens/smoke.sh lens-worker "$RELEASE_TAG" - name: Verify confined Python on the native architecture env: RELEASE_TAG: sha-${{ github.sha }} diff --git a/deploy/lens/compose.yaml b/deploy/lens/compose.yaml index a95684a45c3..663747baa62 100644 --- a/deploy/lens/compose.yaml +++ b/deploy/lens/compose.yaml @@ -12,6 +12,7 @@ services: - "127.0.0.1:${LENS_PORT:-4318}:4318" mem_limit: 2g cpus: 2 + pids_limit: 64 restart: unless-stopped read_only: true tmpfs: diff --git a/deploy/lens/smoke.sh b/deploy/lens/smoke.sh new file mode 100644 index 00000000000..a09060dd635 --- /dev/null +++ b/deploy/lens/smoke.sh @@ -0,0 +1,40 @@ +#!/usr/bin/env bash +set -euo pipefail + +image="${1:?pass the built image reference}" +release="${2:?pass the expected release tag}" +version="$(docker run --rm --network none --read-only --cap-drop ALL --security-opt no-new-privileges "$image" --version)" +test "$version" = "litellm-lens $release protocol=7" +container="$(docker run -d --network none --read-only --cap-drop ALL \ + --security-opt no-new-privileges --pids-limit 64 --memory 2g --cpus 2 \ + --tmpfs /tmp:rw,noexec,nosuid,size=256m \ + -e LITELLM_URL=http://127.0.0.1:1 \ + -e CLICKHOUSE_URL=http://127.0.0.1:1 \ + -e LITELLM_LENS_SERVICE_TOKEN=isolated-runtime-smoke-secret-32-characters \ + "$image")" +trap 'docker rm -f "$container" >/dev/null' EXIT +test "$(docker exec "$container" id -u)" = 65532 +docker exec -i "$container" python3.13 -I -S - <<'PY' +import time +import urllib.error +import urllib.request + +for attempt in range(50): + try: + with urllib.request.urlopen("http://127.0.0.1:4318/health/live", timeout=1) as response: + assert response.status == 200 + break + except urllib.error.URLError: + if attempt == 49: + raise + time.sleep(0.1) + +for path, expected in (("health/ready", 503), ("internal/status", 401)): + try: + urllib.request.urlopen(f"http://127.0.0.1:4318/{path}", timeout=1) + except urllib.error.HTTPError as error: + assert error.code == expected, (path, error.code) + else: + raise AssertionError(f"{path} should return {expected}") +print("Unprivileged Lens service remains live with unavailable dependencies") +PY diff --git a/deploy/lens/stack.yaml b/deploy/lens/stack.yaml index 85d1b38a3d5..ba764129d36 100644 --- a/deploy/lens/stack.yaml +++ b/deploy/lens/stack.yaml @@ -26,7 +26,7 @@ services: - ./config.yaml:/app/lens-config.yaml:ro ports: - "127.0.0.1:${LITELLM_PORT:-4000}:4000" - networks: [proxy, storage] + networks: [proxy, database] depends_on: db: condition: service_healthy @@ -48,6 +48,7 @@ services: - "127.0.0.1:${LENS_PORT:-4318}:4318" mem_limit: 2g cpus: 2 + pids_limit: 64 restart: unless-stopped read_only: true tmpfs: @@ -61,7 +62,7 @@ services: POSTGRES_DB: litellm POSTGRES_USER: litellm POSTGRES_PASSWORD: ${POSTGRES_PASSWORD} - networks: [storage] + networks: [database] volumes: - postgres_data:/var/lib/postgresql/data healthcheck: @@ -89,6 +90,8 @@ services: networks: proxy: + database: + internal: true storage: internal: true diff --git a/docker/docker-compose.tracing.yml b/docker/docker-compose.tracing.yml index 87d7197725e..85708028c6b 100644 --- a/docker/docker-compose.tracing.yml +++ b/docker/docker-compose.tracing.yml @@ -13,8 +13,9 @@ services: LITELLM_SALT_KEY: sk-local-tracing-salt-key DATABASE_URL: postgresql://litellm:litellm@db:5432/litellm STORE_MODEL_IN_DB: "True" - CLICKHOUSE_URL: http://default:local-tracing@clickhouse:8123 - CLICKHOUSE_DATABASE: litellm + LITELLM_LENS_URL: http://lens-worker:4318 + LITELLM_LENS_PUBLIC_URL: http://localhost:4318 + LITELLM_LENS_SERVICE_TOKEN: ${LITELLM_LENS_SERVICE_TOKEN:?set LITELLM_LENS_SERVICE_TOKEN} OPENAI_API_KEY: ${OPENAI_API_KEY:-} LENS_WORKER_IMAGE: ${LENS_WORKER_IMAGE:-} volumes: @@ -24,8 +25,29 @@ services: depends_on: db: condition: service_healthy - clickhouse: - condition: service_healthy + + lens-worker: + build: + context: .. + dockerfile: deploy/lens/Dockerfile + args: + LITELLM_RELEASE_TAG: ${LITELLM_RELEASE_TAG:?set LITELLM_RELEASE_TAG to the source commit} + environment: + LITELLM_URL: http://litellm:4000 + LITELLM_LENS_SERVICE_TOKEN: ${LITELLM_LENS_SERVICE_TOKEN:?set LITELLM_LENS_SERVICE_TOKEN} + CLICKHOUSE_URL: http://default:local-tracing@clickhouse:8123 + CLICKHOUSE_DATABASE: litellm + ports: + - "127.0.0.1:4318:4318" + read_only: true + cap_drop: [ALL] + security_opt: [no-new-privileges:true] + tmpfs: + - /tmp:rw,nosuid,nodev,size=256m,mode=1777 + mem_limit: 2g + cpus: 2 + pids_limit: 64 + restart: unless-stopped db: image: postgres:16 diff --git a/docker/tracing-config.yaml b/docker/tracing-config.yaml index d8e3759641f..7cdd10b4e35 100644 --- a/docker/tracing-config.yaml +++ b/docker/tracing-config.yaml @@ -8,6 +8,4 @@ general_settings: master_key: os.environ/LITELLM_MASTER_KEY tracing: store: - type: clickhouse - url: os.environ/CLICKHOUSE_URL - retention_days: 14 + type: lens diff --git a/helm/litellm-helm/templates/ingress.yaml b/helm/litellm-helm/templates/ingress.yaml index ea9ffcbb54c..5a6bbe9d1cd 100644 --- a/helm/litellm-helm/templates/ingress.yaml +++ b/helm/litellm-helm/templates/ingress.yaml @@ -44,6 +44,20 @@ spec: - host: {{ .host | quote }} http: paths: + {{- if $.Values.lensWorker.enabled }} + - path: /lens-ingest + pathType: Prefix + backend: + {{- if semverCompare ">=1.19-0" $.Capabilities.KubeVersion.GitVersion }} + service: + name: {{ $fullName }}-lens-worker + port: + number: {{ $.Values.lensWorker.service.port }} + {{- else }} + serviceName: {{ $fullName }}-lens-worker + servicePort: {{ $.Values.lensWorker.service.port }} + {{- end }} + {{- end }} {{- range .paths }} - path: {{ .path }} {{- if and .pathType (semverCompare ">=1.18-0" $.Capabilities.KubeVersion.GitVersion) }} diff --git a/helm/litellm-helm/templates/lens/service.yaml b/helm/litellm-helm/templates/lens/service.yaml index a09dde1a4fd..ef063b38cd8 100644 --- a/helm/litellm-helm/templates/lens/service.yaml +++ b/helm/litellm-helm/templates/lens/service.yaml @@ -3,6 +3,10 @@ apiVersion: v1 kind: Service metadata: name: {{ include "litellm.fullname" . }}-lens-worker + {{- with .Values.lensWorker.service.annotations }} + annotations: + {{- toYaml . | nindent 4 }} + {{- end }} spec: selector: app.kubernetes.io/instance: {{ .Release.Name }} diff --git a/helm/litellm-helm/values.yaml b/helm/litellm-helm/values.yaml index 8da4cd52325..c53fd5b0315 100644 --- a/helm/litellm-helm/values.yaml +++ b/helm/litellm-helm/values.yaml @@ -668,6 +668,7 @@ lensWorker: publicUrl: "" service: port: 4318 + annotations: {} ingress: enabled: false className: "" diff --git a/helm/litellm/templates/_helpers.tpl b/helm/litellm/templates/_helpers.tpl index eb7433c279a..0b5ccc14ddb 100644 --- a/helm/litellm/templates/_helpers.tpl +++ b/helm/litellm/templates/_helpers.tpl @@ -514,3 +514,17 @@ shutdown drain window. value: {{ .drainTimeoutSeconds | quote }} {{- end }} {{- end -}} + +{{- define "litellm.lensConnectionEnv" -}} +{{- if .Values.lensWorker.enabled }} +- name: LITELLM_LENS_URL + value: {{ printf "http://%s-lens-worker:%v" (include "litellm.fullname" .) .Values.lensWorker.service.port | quote }} +- name: LITELLM_LENS_PUBLIC_URL + value: {{ required "lensWorker.publicUrl is required" .Values.lensWorker.publicUrl | quote }} +- name: LITELLM_LENS_SERVICE_TOKEN + valueFrom: + secretKeyRef: + name: {{ required "lensWorker.serviceTokenSecret.name is required" .Values.lensWorker.serviceTokenSecret.name | quote }} + key: {{ .Values.lensWorker.serviceTokenSecret.key | quote }} +{{- end }} +{{- end -}} diff --git a/helm/litellm/templates/backend/deployment.yaml b/helm/litellm/templates/backend/deployment.yaml index ebb9a1dd838..b752a7c10ff 100644 --- a/helm/litellm/templates/backend/deployment.yaml +++ b/helm/litellm/templates/backend/deployment.yaml @@ -57,17 +57,7 @@ spec: containerPort: 4001 protocol: TCP env: - {{- if .Values.lensWorker.enabled }} - - name: LITELLM_LENS_URL - value: {{ printf "http://%s-lens-worker:%v" (include "litellm.fullname" .) .Values.lensWorker.service.port | quote }} - - name: LITELLM_LENS_PUBLIC_URL - value: {{ required "lensWorker.publicUrl is required" .Values.lensWorker.publicUrl | quote }} - - name: LITELLM_LENS_SERVICE_TOKEN - valueFrom: - secretKeyRef: - name: {{ required "lensWorker.serviceTokenSecret.name is required" .Values.lensWorker.serviceTokenSecret.name | quote }} - key: {{ .Values.lensWorker.serviceTokenSecret.key | quote }} - {{- end }} + {{- include "litellm.lensConnectionEnv" . | nindent 12 }} - name: LENS_WORKER_IMAGE value: {{ include "litellm.lensWorker.image" . | quote }} {{- include "litellm.serverEnv" (dict "root" $ "component" .Values.backend) | nindent 12 }} diff --git a/helm/litellm/templates/gateway/deployment.yaml b/helm/litellm/templates/gateway/deployment.yaml index 49b452b3053..9b3b58c97ed 100644 --- a/helm/litellm/templates/gateway/deployment.yaml +++ b/helm/litellm/templates/gateway/deployment.yaml @@ -55,6 +55,7 @@ spec: containerPort: 4000 protocol: TCP env: + {{- include "litellm.lensConnectionEnv" . | nindent 12 }} {{- include "litellm.serverEnv" (dict "root" $ "component" .Values.gateway) | nindent 12 }} {{- if .Values.gateway.config.create }} - name: CONFIG_FILE_PATH diff --git a/helm/litellm/templates/ingress.yaml b/helm/litellm/templates/ingress.yaml index e9f7ed4ec3f..7bfbe85db79 100644 --- a/helm/litellm/templates/ingress.yaml +++ b/helm/litellm/templates/ingress.yaml @@ -156,6 +156,16 @@ spec: port: number: {{ $gatewayPort }} {{- end }} + {{- if .Values.lensWorker.enabled }} + {{- $builtinPathKeys = append $builtinPathKeys "/lens-ingest|Prefix" }} + - path: /lens-ingest + pathType: Prefix + backend: + service: + name: {{ include "litellm.fullname" . }}-lens-worker + port: + number: {{ .Values.lensWorker.service.port }} + {{- end }} {{- /* --- Operator-supplied extra paths (ingress.extraPaths) --- Rendered after every built-in path so an entry can never take diff --git a/helm/litellm/templates/lens/service.yaml b/helm/litellm/templates/lens/service.yaml index a09dde1a4fd..ef063b38cd8 100644 --- a/helm/litellm/templates/lens/service.yaml +++ b/helm/litellm/templates/lens/service.yaml @@ -3,6 +3,10 @@ apiVersion: v1 kind: Service metadata: name: {{ include "litellm.fullname" . }}-lens-worker + {{- with .Values.lensWorker.service.annotations }} + annotations: + {{- toYaml . | nindent 4 }} + {{- end }} spec: selector: app.kubernetes.io/instance: {{ .Release.Name }} diff --git a/helm/litellm/tests/lens_worker_tests.yaml b/helm/litellm/tests/lens_worker_tests.yaml index 9a83a4af6a4..b93797dae18 100644 --- a/helm/litellm/tests/lens_worker_tests.yaml +++ b/helm/litellm/tests/lens_worker_tests.yaml @@ -11,7 +11,9 @@ tests: set: backend.image.tag: sha-0123456789abcdef lensWorker.enabled: true - lensWorker.tokenSecret.name: lens-credential + lensWorker.publicUrl: https://traces.example + lensWorker.serviceTokenSecret.name: lens-service + lensWorker.clickhouseSecret.name: lens-storage asserts: - equal: path: spec.template.spec.containers[0].image @@ -31,7 +33,9 @@ tests: set: backend.image.tag: sha-0123456789abcdef lensWorker.enabled: true - lensWorker.tokenSecret.name: lens-credential + lensWorker.publicUrl: https://traces.example + lensWorker.serviceTokenSecret.name: lens-service + lensWorker.clickhouseSecret.name: lens-storage lensWorker.image.repository: registry.example/lens-worker asserts: - equal: @@ -41,7 +45,9 @@ tests: template: lens/deployment.yaml set: lensWorker.enabled: true - lensWorker.tokenSecret.name: lens-credential + lensWorker.publicUrl: https://traces.example + lensWorker.serviceTokenSecret.name: lens-service + lensWorker.clickhouseSecret.name: lens-storage lensWorker.image.tag: replaced-release lensWorker.image.digest: sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa asserts: @@ -71,20 +77,22 @@ tests: asserts: - hasDocuments: count: 0 - - it: requires a limited worker credential when enabled + - it: requires a shared service secret when enabled template: lens/deployment.yaml set: lensWorker.enabled: true asserts: - failedTemplate: - errorMessage: lensWorker.tokenSecret.name must reference a Lens worker token + errorMessage: lensWorker.serviceTokenSecret.name is required - it: uses the chart release and a secret without granting Kubernetes access template: lens/deployment.yaml chart: appVersion: v1.2.3 set: lensWorker.enabled: true - lensWorker.tokenSecret.name: lens-credential + lensWorker.publicUrl: https://traces.example + lensWorker.serviceTokenSecret.name: lens-service + lensWorker.clickhouseSecret.name: lens-storage asserts: - equal: path: spec.template.spec.containers[0].image @@ -92,8 +100,8 @@ tests: - equal: path: spec.template.spec.containers[0].env[1].valueFrom.secretKeyRef value: - name: lens-credential - key: token + name: lens-service + key: service-token - equal: path: spec.template.spec.automountServiceAccountToken value: false @@ -120,7 +128,9 @@ tests: template: lens/deployment.yaml set: lensWorker.enabled: true - lensWorker.tokenSecret.name: lens-credential + lensWorker.publicUrl: https://traces.example + lensWorker.serviceTokenSecret.name: lens-service + lensWorker.clickhouseSecret.name: lens-storage lensWorker.url: https://gateway.example/proxy lensWorker.image.repository: registry.example/lens-worker lensWorker.image.tag: branch-main-1234567 @@ -137,7 +147,9 @@ tests: appVersion: 1.2.3-rc.4 set: lensWorker.enabled: true - lensWorker.tokenSecret.name: lens-credential + lensWorker.publicUrl: https://traces.example + lensWorker.serviceTokenSecret.name: lens-service + lensWorker.clickhouseSecret.name: lens-storage asserts: - equal: path: spec.template.spec.containers[0].image @@ -147,7 +159,9 @@ tests: set: backend.image.tag: branch-main-1234567 lensWorker.enabled: true - lensWorker.tokenSecret.name: lens-credential + lensWorker.publicUrl: https://traces.example + lensWorker.serviceTokenSecret.name: lens-service + lensWorker.clickhouseSecret.name: lens-storage asserts: - equal: path: spec.template.spec.containers[0].image @@ -167,7 +181,9 @@ tests: set: backend.image.tag: 1.2.3-dev.4 lensWorker.enabled: true - lensWorker.tokenSecret.name: lens-credential + lensWorker.publicUrl: https://traces.example + lensWorker.serviceTokenSecret.name: lens-service + lensWorker.clickhouseSecret.name: lens-storage asserts: - equal: path: spec.template.spec.containers[0].image diff --git a/helm/litellm/values.yaml b/helm/litellm/values.yaml index 7b1d591babe..e8820102098 100644 --- a/helm/litellm/values.yaml +++ b/helm/litellm/values.yaml @@ -650,6 +650,7 @@ lensWorker: publicUrl: "" service: port: 4318 + annotations: {} ingress: enabled: false className: "" diff --git a/litellm-rust/Cargo.lock b/litellm-rust/Cargo.lock index c7a8f8fba6f..987491d514b 100644 --- a/litellm-rust/Cargo.lock +++ b/litellm-rust/Cargo.lock @@ -4206,6 +4206,7 @@ dependencies = [ "typify", "unicode-casefold", "url", + "uuid", "wiremock", ] diff --git a/litellm-rust/crates/lens/Cargo.toml b/litellm-rust/crates/lens/Cargo.toml index 4b1ac05aca5..d513e979f97 100644 --- a/litellm-rust/crates/lens/Cargo.toml +++ b/litellm-rust/crates/lens/Cargo.toml @@ -43,3 +43,4 @@ prettyplease = "0.2" [dev-dependencies] rstest.workspace = true wiremock.workspace = true +uuid.workspace = true diff --git a/litellm-rust/crates/lens/src/lib.rs b/litellm-rust/crates/lens/src/lib.rs index ed7e35eecb4..25830fda896 100644 --- a/litellm-rust/crates/lens/src/lib.rs +++ b/litellm-rust/crates/lens/src/lib.rs @@ -98,6 +98,8 @@ pub fn router(state: Arc) -> Router { ]), ); public + .clone() + .nest("/lens-ingest", public) .merge( Router::new() .route("/internal/read", post(read)) diff --git a/litellm-rust/crates/lens/tests/clickhouse.rs b/litellm-rust/crates/lens/tests/clickhouse.rs index 14a7aba7d63..770e2434d61 100644 --- a/litellm-rust/crates/lens/tests/clickhouse.rs +++ b/litellm-rust/crates/lens/tests/clickhouse.rs @@ -15,12 +15,20 @@ use std::{ }; #[rstest] +#[case::own_trace("isolated-ingestion-key", vec![], true)] +#[case::own_span("isolated-ingestion-key", vec!["aabbccdd00112233"], true)] +#[case::missing_span("isolated-ingestion-key", vec!["ffffffffffffffff"], false)] +#[case::other_key("other-ingestion-key", vec![], false)] #[tokio::test] #[ignore = "requires an isolated ClickHouse instance in LENS_TEST_CLICKHOUSE_URL"] -async fn traces_round_trip_through_real_clickhouse_with_scoped_reads() { +async fn traces_round_trip_through_real_clickhouse_with_scoped_reads( + #[case] key: &str, + #[case] spans: Vec<&str>, + #[case] expected: bool, +) { let url = std::env::var("LENS_TEST_CLICKHOUSE_URL").expect("set LENS_TEST_CLICKHOUSE_URL"); let client = http_client().unwrap(); - let database = format!("lens_test_{}_{}", std::process::id(), unix_seconds()); + let database = format!("lens_test_{}", uuid::Uuid::new_v4().simple()); let config = Config::new(database.clone(), &url, 14, 65_536).unwrap(); let storage = Storage::new( config.clone(), @@ -37,16 +45,19 @@ async fn traces_round_trip_through_real_clickhouse_with_scoped_reads() { .credentials .replace(Snapshot { issued_at: unix_seconds(), - keys: vec![Credential { - token_hash: format!("{:x}", Sha256::digest("isolated-ingestion-key")), - tenant: Tenant { - team_id: "team-a".into(), - user_id: "user-a".into(), - api_key_hash: "key-a".into(), - ..Tenant::default() - }, - expires_at: None, - }], + keys: ["isolated-ingestion-key", "other-ingestion-key"] + .into_iter() + .map(|key| Credential { + token_hash: format!("{:x}", Sha256::digest(key)), + tenant: Tenant { + team_id: "team-a".into(), + user_id: "user-a".into(), + api_key_hash: format!("{:x}", Sha256::digest(key)), + ..Tenant::default() + }, + expires_at: None, + }) + .collect(), }) .unwrap(); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); @@ -70,6 +81,18 @@ async fn traces_round_trip_through_real_clickhouse_with_scoped_reads() { .await .unwrap(); assert_eq!(written.status(), 200, "{}", written.text().await.unwrap()); + let receipt = client + .post(format!("{endpoint}/v1/traces/receipt")) + .bearer_auth(key) + .json(&json!({"trace_id": trace_id, "span_ids": spans})) + .send() + .await + .unwrap(); + assert_eq!(receipt.status(), 200); + assert_eq!( + receipt.json::().await.unwrap(), + json!({"received": expected}) + ); let read = json!({"operation":"list","scope":{"all_teams":0,"user_id":"user-a","team_ids":[]},"start_ms":now/1_000_000-1000,"end_ms":now/1_000_000+1000,"cursor":null,"limit":50}); let found = client .post(format!("{endpoint}/internal/read")) diff --git a/litellm-rust/crates/lens/tests/receiver.rs b/litellm-rust/crates/lens/tests/receiver.rs index 1ea01c241b4..3383d29aab3 100644 --- a/litellm-rust/crates/lens/tests/receiver.rs +++ b/litellm-rust/crates/lens/tests/receiver.rs @@ -108,6 +108,86 @@ async fn ingestion_confirms_storage_and_overwrites_exporter_tenant() { assert_eq!(row["ApiKeyHash"], "authenticated-key"); } +#[rstest] +#[tokio::test] +async fn shared_ingress_prefix_exposes_uploads_without_internal_control_routes() { + let store = MockServer::start().await; + Mock::given(method("POST")) + .respond_with(ResponseTemplate::new(200)) + .expect(1) + .mount(&store) + .await; + let server = serve(&store.uri(), true).await; + let client = http_client().unwrap(); + let upload = client + .post(format!("{}/lens-ingest/v1/traces", server.url)) + .bearer_auth(KEY) + .json(&export()) + .send() + .await + .unwrap(); + assert_eq!(upload.status(), 200); + let internal = client + .get(format!("{}/lens-ingest/internal/status", server.url)) + .bearer_auth(SERVICE_TOKEN) + .send() + .await + .unwrap(); + assert_eq!(internal.status(), 404); + let preflight = client + .request( + http::Method::OPTIONS, + format!("{}/lens-ingest/v1/traces", server.url), + ) + .header("origin", "https://dashboard.example") + .header("access-control-request-method", "POST") + .header( + "access-control-request-headers", + "authorization,content-type", + ) + .send() + .await + .unwrap(); + assert_eq!(preflight.headers()["access-control-allow-origin"], "*"); + assert!( + !preflight + .headers() + .contains_key("access-control-allow-credentials") + ); +} + +#[rstest] +#[tokio::test] +async fn only_the_service_secret_can_replace_ingestion_credentials() { + let server = serve("http://127.0.0.1:1", true).await; + let client = http_client().unwrap(); + let snapshot = json!({"issued_at": unix_seconds(), "keys": []}); + let denied = client + .post(format!("{}/internal/credentials", server.url)) + .bearer_auth(KEY) + .json(&snapshot) + .send() + .await + .unwrap(); + assert_eq!(denied.status(), 401); + let accepted = client + .post(format!("{}/internal/credentials", server.url)) + .bearer_auth(SERVICE_TOKEN) + .json(&snapshot) + .send() + .await + .unwrap(); + assert_eq!(accepted.status(), 204); + let revoked = client + .post(format!("{}/v1/traces", server.url)) + .bearer_auth(KEY) + .json(&export()) + .send() + .await + .unwrap(); + assert_eq!(revoked.status(), 401); +} + #[rstest] #[case::refused(503)] #[case::disk_full(507)] diff --git a/litellm/proxy/_lazy_openapi_snapshot.json b/litellm/proxy/_lazy_openapi_snapshot.json index d864ec04a8d..fd470de1422 100644 --- a/litellm/proxy/_lazy_openapi_snapshot.json +++ b/litellm/proxy/_lazy_openapi_snapshot.json @@ -34933,6 +34933,11 @@ "title": "Image", "type": "string" }, + "managed": { + "default": false, + "title": "Managed", + "type": "boolean" + }, "token": { "title": "Token", "type": "string" @@ -34956,6 +34961,11 @@ "title": "Analysis Key Id", "type": "string" }, + "managed": { + "default": false, + "title": "Managed", + "type": "boolean" + }, "name": { "default": "Lens worker", "minLength": 1, diff --git a/litellm/proxy/lens/endpoints.py b/litellm/proxy/lens/endpoints.py index d3451d6ace0..eae8747306a 100644 --- a/litellm/proxy/lens/endpoints.py +++ b/litellm/proxy/lens/endpoints.py @@ -78,7 +78,7 @@ from litellm.proxy.lens.state import ( summarized, ) from litellm.proxy.tracing_runtime import provide_storage -from litellm.tracing.remote import LensConnection +from litellm.tracing.remote import LensConnection, bounded_response from litellm.types.llms.base import LiteLLMBaseModel router: Final = APIRouter(prefix="/lens", tags=["Lens"]) @@ -152,12 +152,14 @@ async def service_connection(auth: Auth) -> ServiceConnection: public_url: Final = os.environ.get("LITELLM_LENS_PUBLIC_URL", "").rstrip("/") try: connection: Final = LensConnection.from_env() - async with connection.client() as client: - response: Final = await client.get("/internal/status", timeout=2) + async with ( + connection.client() as client, + client.stream("GET", "/internal/status", timeout=2) as response, + ): if response.status_code == 200: - status: Final = ServiceStatus.model_validate_json(response.content) + status: Final = ServiceStatus.model_validate_json(await bounded_response(response, 16 * 1024)) return ServiceConnection(url=public_url, connected=True, status=status) - except (ValueError, httpx.HTTPError): + except (ValueError, RuntimeError, httpx.HTTPError): pass return ServiceConnection(url=public_url, connected=False, status=ServiceStatus()) @@ -182,7 +184,9 @@ async def publish_credentials() -> bool: connection: Final = LensConnection.from_env() snapshot: Final = await credential_snapshot() async with connection.client() as client: - response: Final = await client.post("/internal/credentials", json=snapshot.model_dump(mode="json"), timeout=2) + response: Final = await client.post( + "/internal/credentials", json=snapshot.model_dump(mode="json"), timeout=2 + ) return response.status_code == 204 except (ValueError, httpx.HTTPError): return False @@ -671,7 +675,15 @@ async def sample(lens_id: str, job_id: str, worker: WorkerAuth, storage: Storage def freeze(e: Lens) -> Lens: active: Final = current_job(e) - if active is None or active.id != job_id or active.worker_id != worker.id or active.attempts != attempt: + if ( + active is None + or active.id != job_id + or active.worker_id != worker.id + or active.attempts != attempt + or active.status != "running" + or active.lease_until is None + or active.lease_until <= datetime.now(timezone.utc) + ): raise HTTPException(409, "Job was cancelled or reassigned") return ( replace_job(e, active.model_copy(update=MappingProxyType({"sample": selected}))) @@ -718,7 +730,13 @@ def model_failure(error: HTTPException | ProxyException) -> HTTPException: @router.post("/worker/{lens_id}/{job_id}/model", response_model=ModelResult) async def model( - lens_id: str, job_id: str, body: ModelRequest, worker: WorkerAuth, request: Request, response: Response, attempt: Attempt = 1 + lens_id: str, + job_id: str, + body: ModelRequest, + worker: WorkerAuth, + request: Request, + response: Response, + attempt: Attempt = 1, ) -> ModelResult: from litellm.proxy.lens.inference import analyze @@ -733,7 +751,9 @@ async def model( @router.post("/worker/{lens_id}/{job_id}/result", response_model=Lens) -async def result(lens_id: str, job_id: str, body: Result, worker: WorkerAuth, storage: StorageDep, attempt: Attempt = 1) -> Lens: +async def result( + lens_id: str, job_id: str, body: Result, worker: WorkerAuth, storage: StorageDep, attempt: Attempt = 1 +) -> Lens: lens: Final = await get_lens(lens_id, worker.scope) old: Final = next((j for j in lens.jobs if j.id == job_id), None) if old and old.status in ("completed", "failed") and old.worker_id == worker.id and old.attempts == attempt: @@ -769,7 +789,15 @@ async def result(lens_id: str, job_id: str, body: Result, worker: WorkerAuth, st def finish(e: Lens) -> Lens: active: Final = current_job(e) - if active is None or active.id != job_id or active.worker_id != worker.id or active.attempts != attempt: + if ( + active is None + or active.id != job_id + or active.worker_id != worker.id + or active.attempts != attempt + or active.status != "running" + or active.lease_until is None + or active.lease_until <= datetime.now(timezone.utc) + ): return e restored: Final = e.model_copy( update=MappingProxyType( @@ -836,7 +864,10 @@ async def result(lens_id: str, job_id: str, body: Result, worker: WorkerAuth, st ) finished: Final = required(await repository().update(lens_id, finish)) - if body.review_versions and any(j.id == job_id and j.status == "completed" for j in finished.jobs): + if body.review_versions and any( + j.id == job_id and j.status == "completed" and j.attempts == attempt and j.worker_id == worker.id + for j in finished.jobs + ): await repository().complete_reviews(lens_id, job, body.review_versions) return finished diff --git a/litellm/proxy/lens/repository.py b/litellm/proxy/lens/repository.py index 7b57ef1beb9..05099014091 100644 --- a/litellm/proxy/lens/repository.py +++ b/litellm/proxy/lens/repository.py @@ -173,7 +173,7 @@ class LensRepository: async def claim_candidates(self, scope: Scope, now: datetime) -> AsyncIterator[Lens]: cursor = "" # rebind-ok: advance a bounded keyset page while True: - rows: Final = _ROWS.validate_python( + rows = _ROWS.validate_python( # rebind-ok: fetch the next bounded keyset page await self.db.query_raw( """SELECT data FROM "LiteLLM_Lens" WHERE id > $1 AND ($2::boolean OR ( @@ -183,9 +183,9 @@ class LensRepository: AND ( EXISTS (SELECT 1 FROM jsonb_array_elements(data->'jobs') AS job WHERE job->>'status'='queued' OR (job->>'status'='running' - AND (job->>'lease_until' IS NULL OR (job->>'lease_until')::timestamptz<=$5))) + AND (job->>'lease_until' IS NULL OR (job->>'lease_until')::timestamptz<=$5::timestamptz))) OR ((data->'settings'->>'enabled')::boolean - AND (data->>'next_run_at')::timestamptz<=$5 + AND (data->>'next_run_at')::timestamptz<=$5::timestamptz AND NOT EXISTS (SELECT 1 FROM jsonb_array_elements(data->'jobs') AS job WHERE job->>'status' IN ('queued', 'running')))) ORDER BY id LIMIT 50""", @@ -193,10 +193,10 @@ class LensRepository: scope.all_teams, scope.team_id, scope.api_key_hash, - now, + now.isoformat(), ) ) - candidates: Final = tuple(Lens.model_validate(row.data) for row in rows) + candidates = tuple(Lens.model_validate(row.data) for row in rows) # rebind-ok: decode this page for candidate in candidates: yield candidate if len(candidates) < 50: @@ -395,7 +395,7 @@ class LensRepository: rows: Final = _ROWS.validate_python( await self.db.query_raw( 'INSERT INTO "LiteLLM_LensWorker" AS existing (id,token_hash,data) VALUES ($1,$2,$3::jsonb) ' - 'ON CONFLICT (token_hash) DO UPDATE ' + "ON CONFLICT (token_hash) DO UPDATE " "SET data=jsonb_set(EXCLUDED.data, '{id}', to_jsonb(existing.id)) RETURNING data", worker.id, token_hash, diff --git a/litellm/tracing/exporter.py b/litellm/tracing/exporter.py index a277b7cd5a1..a6ce1e7e6ec 100644 --- a/litellm/tracing/exporter.py +++ b/litellm/tracing/exporter.py @@ -1,10 +1,11 @@ import asyncio import json from collections import deque -from collections.abc import Iterable, Mapping +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 @@ -22,6 +23,14 @@ SHUTDOWN_SECONDS: Final = 3.0 _PAYLOAD: Final = TypeAdapter(SpendLogPayload) +def _is_mapping(value: object) -> TypeIs[Mapping[object, object]]: + return isinstance(value, Mapping) + + +def _is_sequence(value: object) -> TypeIs[Sequence[object]]: + return isinstance(value, (tuple, list)) + + def _check_size(value: object, remaining: int, depth: int = 0) -> int: if remaining <= 0 or depth > 32: raise OverflowError("Trace record exceeds the export budget") @@ -29,9 +38,9 @@ def _check_size(value: object, remaining: int, depth: int = 0) -> int: if len(value) > remaining: raise OverflowError("Trace record exceeds the export budget") return remaining - len(value.encode()) - if isinstance(value, Mapping): + if _is_mapping(value): return _check_sequence(value.items(), remaining, depth) - if isinstance(value, (tuple, list)): + if _is_sequence(value): return _check_sequence(value, remaining, depth) return remaining - 32 @@ -39,7 +48,7 @@ def _check_size(value: object, remaining: int, depth: int = 0) -> int: def _check_sequence(values: Iterable[object], remaining: int, depth: int) -> int: budget = remaining # rebind-ok: consumes a finite serialization budget for value in values: - budget = _check_size(value, budget - 8, depth + 1) # rebind-ok: consumes a finite serialization budget + budget = _check_size(value, budget - 8, depth + 1) if budget < 0: raise OverflowError("Trace record exceeds the export budget") return budget @@ -48,10 +57,10 @@ def _check_sequence(values: Iterable[object], remaining: int, depth: int) -> int def encode_record(value: Mapping[str, object]) -> bytes: _check_size(value, MAX_EVENT_BYTES) with BytesIO() as output: - for part in json.JSONEncoder(ensure_ascii=False, allow_nan=False, separators=(",", ":")).iterencode( + parts: Final = json.JSONEncoder(ensure_ascii=False, allow_nan=False, separators=(",", ":")).iterencode( dict(value) - ): - encoded: Final = part.encode() + ) + for encoded in (part.encode() for part in parts): if output.tell() + len(encoded) > MAX_EVENT_BYTES: raise OverflowError("Trace record exceeds the export budget") output.write(encoded) @@ -135,7 +144,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: Final = self.queue.popleft() + record = self.queue.popleft() # rebind-ok: drain each record into the bounded batch size += len(record) + 1 records.append(record) return tuple(records) @@ -160,7 +169,7 @@ class LensExporter(CustomLogger): except httpx.HTTPError: pass if attempt < 2: - await asyncio.sleep(2**attempt) + await asyncio.sleep(float(1 << attempt)) self._warn("retry limit reached") return False @@ -170,18 +179,21 @@ class LensExporter(CustomLogger): self.wake.clear() await self.wake.wait() continue - batch: Final = self._batch() - try: - if await self._send(batch): - self.rows_written += len(batch) - else: - self.rows_dropped += len(batch) - except asyncio.CancelledError: + await self._drain_batch() + + async def _drain_batch(self) -> None: + batch: Final = self._batch() + try: + if await self._send(batch): + self.rows_written += len(batch) + else: self.rows_dropped += len(batch) - raise - finally: - self.buffered_events -= len(batch) - self.buffered_bytes -= sum(len(record) for record in batch) + except asyncio.CancelledError: + self.rows_dropped += len(batch) + raise + finally: + self.buffered_events -= len(batch) + self.buffered_bytes -= sum(len(record) for record in batch) async def aclose(self) -> None: self.closed = True diff --git a/litellm/tracing/remote.py b/litellm/tracing/remote.py index bd5badf77bc..8612dbc82a1 100644 --- a/litellm/tracing/remote.py +++ b/litellm/tracing/remote.py @@ -12,7 +12,7 @@ from litellm.rust_bridge.trace.errors import TraceChanged from litellm.rust_bridge.trace.generated.types import QueryScope, ReadQueryName, TraceScope MAX_RESPONSE_BYTES: Final = 64 * 1024 * 1024 -_JSON: Final = TypeAdapter(JsonValue) +_JSON: Final[TypeAdapter[JsonValue]] = TypeAdapter(JsonValue) @dataclass(frozen=True, slots=True, repr=False) @@ -86,7 +86,7 @@ class RemoteTraceStore: return await self._read( { "operation": "list", - "scope": scope.model_dump(), + "scope": scope, "start_ms": start_ms, "end_ms": end_ms, "cursor": cursor, @@ -100,7 +100,7 @@ class RemoteTraceStore: return await self._read( { "operation": "trace", - "scope": scope.model_dump(), + "scope": scope, "trace_id": trace_id, "trace_ref": trace_ref, "cursor": cursor, @@ -112,7 +112,7 @@ class RemoteTraceStore: return await self._read( { "operation": "span", - "scope": scope.model_dump(), + "scope": scope, "trace_id": trace_id, "trace_ref": trace_ref, "span_id": span_id, @@ -125,7 +125,7 @@ class RemoteTraceStore: return await self._read( { "operation": "span_error", - "scope": scope.model_dump(), + "scope": scope, "trace_id": trace_id, "trace_ref": trace_ref, "span_id": span_id, @@ -134,10 +134,10 @@ class RemoteTraceStore: ) async def query_sql(self, sql: str, scope: QueryScope, secret: str) -> str: - return json.dumps(await self._read({"operation": "sql", "sql": sql, "scope": scope.model_dump()})) + return json.dumps(await self._read({"operation": "sql", "sql": sql, "scope": scope})) async def query_help(self, scope: QueryScope, secret: str) -> JsonValue: - return await self._read({"operation": "help", "scope": scope.model_dump()}) + return await self._read({"operation": "help", "scope": scope}) async def query(self, name: ReadQueryName, parameters: Mapping[str, str | int | float | Sequence[str]]) -> str: return json.dumps(await self._read({"operation": "query", "name": name, "parameters": dict(parameters)})) diff --git a/scripts/lens_dev.sh b/scripts/lens_dev.sh index 69ef8c0f8bd..d1d6aa1cfda 100755 --- a/scripts/lens_dev.sh +++ b/scripts/lens_dev.sh @@ -17,10 +17,12 @@ repo_root="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" source_release_tag="sha-$(git -C "$repo_root" rev-parse HEAD)" proxy_port="${LENS_DEV_PROXY_PORT:-4000}" ui_port="${LENS_DEV_UI_PORT:-3000}" +lens_port="${LENS_DEV_SERVICE_PORT:-4318}" state_dir="${LENS_DEV_STATE_DIR:-$repo_root/.lens-dev}" log_dir="$state_dir/logs" token_file="$state_dir/worker_token" key_file="$state_dir/master_key" +service_key_file="$state_dir/service_key" proxy_url="http://localhost:$proxy_port" py="${LENS_DEV_PYTHON:-$repo_root/.venv/bin/python}" database_url="${LENS_DEV_DATABASE_URL:-postgresql://litellm:litellm@127.0.0.1:15432/litellm}" @@ -64,7 +66,9 @@ ensure_services() { services+=(clickhouse) fi if [ "${#services[@]}" -gt 0 ]; then - docker compose -f docker/docker-compose.tracing.yml up -d --wait "${services[@]}" + LITELLM_MASTER_KEY="$master_key" LITELLM_LENS_SERVICE_TOKEN="${service_key:-}" \ + LITELLM_RELEASE_TAG="$source_release_tag" \ + docker compose -f docker/docker-compose.tracing.yml up -d --wait "${services[@]}" else echo "lens-dev: reusing running Postgres and ClickHouse" fi @@ -73,18 +77,16 @@ ensure_services() { write_default_config() { cat > "$1" <<'EOF' model_list: - - model_name: gpt-4.1-mini + - model_name: gpt-6.1-sol litellm_params: - model: openai/gpt-4.1-mini + model: openai/gpt-6.1-sol api_key: os.environ/OPENAI_API_KEY general_settings: master_key: os.environ/LITELLM_MASTER_KEY store_prompts_in_spend_logs: true tracing: store: - type: clickhouse - url: os.environ/CLICKHOUSE_URL - retention_days: 14 + type: lens EOF } @@ -119,8 +121,10 @@ proxy_env() { export LITELLM_SALT_KEY=sk-local-tracing-salt-key export DATABASE_URL="$database_url" export STORE_MODEL_IN_DB=True - export CLICKHOUSE_URL="$clickhouse_url" - export CLICKHOUSE_DATABASE=litellm + unset CLICKHOUSE_URL CLICKHOUSE_DATABASE + export LITELLM_LENS_URL="http://127.0.0.1:$lens_port" + export LITELLM_LENS_PUBLIC_URL="http://localhost:$lens_port" + export LITELLM_LENS_SERVICE_TOKEN="${service_key:-}" export LITELLM_LOCAL_MODEL_COST_MAP=True export PROXY_BASE_URL="$proxy_url" export LITELLM_UI_PATH="$repo_root/ui/litellm-dashboard/out" @@ -221,6 +225,8 @@ build_dashboard() { seed_data() { ( proxy_env "" + export CLICKHOUSE_URL="$clickhouse_url" + export CLICKHOUSE_DATABASE=litellm export LENS_DEV_UI_URL="http://localhost:$ui_port" if [ -n "$seed_profile" ]; then "$py" -m scripts.seed_tracing_fixtures --profile "$seed_profile" ${seed_options[@]+"${seed_options[@]}"} @@ -288,9 +294,15 @@ main() { [[ "$readiness_request_timeout" =~ ^[1-9][0-9]*$ ]] || die "LENS_DEV_READINESS_REQUEST_TIMEOUT_SECONDS must be a positive integer" listening "$proxy_port" && die "port $proxy_port is in use; set LENS_DEV_PROXY_PORT" listening "$ui_port" && die "port $ui_port is in use; set LENS_DEV_UI_PORT" + listening "$lens_port" && die "port $lens_port is in use; set LENS_DEV_SERVICE_PORT" [ "$proxy_port" != "$ui_port" ] || die "proxy and UI ports must differ" + [ "$lens_port" != "$proxy_port" ] && [ "$lens_port" != "$ui_port" ] || die "Lens service port must differ from proxy and UI ports" mkdir -p "$log_dir" load_master_key + if [ ! -s "$service_key_file" ]; then + (umask 077 && openssl rand -hex 32 > "$service_key_file") + fi + service_key="$(cat "$service_key_file")" uv sync --inexact --frozen --extra proxy --group proxy-dev --no-install-project ensure_services @@ -304,6 +316,7 @@ main() { echo "lens-dev: checking the Rust bridge (litellm.rust_bridge._native) is current; the ClickHouse trace store uses it" PYO3_PYTHON="$py" VIRTUAL_ENV="$repo_root/.venv" uvx --from maturin==1.15.0 maturin develop \ --release --manifest-path litellm-rust/crates/python-bridge/Cargo.toml --features extension-module + cargo build --locked --manifest-path litellm-rust/Cargo.toml -p litellm-lens if [ ! -x ui/litellm-dashboard/node_modules/.bin/next ]; then (cd ui/litellm-dashboard && "$repo_root/scripts/with_dashboard_node.sh" npm ci) @@ -331,7 +344,8 @@ main() { ( cd ui/litellm-dashboard - NEXT_PUBLIC_BASE_URL="" LENS_DEV_PROXY_URL="$proxy_url" exec "$repo_root/scripts/with_dashboard_node.sh" npx next dev -p "$ui_port" + NEXT_PUBLIC_BASE_URL="" NEXT_PUBLIC_USE_REWRITES=true LENS_DEV_PROXY_URL="$proxy_url" \ + exec "$repo_root/scripts/with_dashboard_node.sh" npx next dev -p "$ui_port" ) < /dev/null > "$log_dir/ui.log" 2>&1 & ui_pid=$! pids+=("$ui_pid") @@ -343,7 +357,9 @@ main() { LITELLM_RELEASE_TAG="$source_release_tag" \ LITELLM_MODE=PRODUCTION LITELLM_URL="$proxy_url" LENS_WORKER_TOKEN="$(cat "$token_file")" \ - "$py" -c "import asyncio, logging; from litellm.proxy.lens.worker import main; logging.basicConfig(level=logging.INFO); asyncio.run(main())" \ + 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+=("$!") @@ -356,6 +372,7 @@ Lens dev is up. Ctrl-C stops everything. Lens: http://localhost:$ui_port/ui/lens/ (hot-reloads) Logs: http://localhost:$ui_port/ui/?page=logs API: $proxy_url + Traces: http://localhost:$lens_port/v1/traces Logs: $log_dir/proxy.log $log_dir/worker.log $log_dir/ui.log diff --git a/tests/integration/database/test_lens_repository.py b/tests/integration/database/test_lens_repository.py index e29e92a2505..dd9f467fd76 100644 --- a/tests/integration/database/test_lens_repository.py +++ b/tests/integration/database/test_lens_repository.py @@ -88,6 +88,73 @@ async def test_heartbeat_never_restores_revoked_access(lens_db: Prisma) -> None: await lens_db.execute_raw('DELETE FROM "LiteLLM_LensWorker" WHERE id=$1', worker.id) +@pytest.mark.asyncio +async def test_managed_registration_is_atomic_and_keeps_the_original_worker_id(lens_db: Prisma) -> None: + now: Final = datetime.now(timezone.utc) + repo: Final = LensRepository(WriterDatabase(PrismaWrapper(lens_db))) + token_hash: Final = uuid4().hex + workers: Final = tuple( + Worker(id=uuid4().hex, name="Managed Lens", scope=Scope(all_teams=True), last_seen=now) for _ in range(8) + ) + try: + registered: Final = await asyncio.gather(*(repo.configure_service_worker(w, token_hash) for w in workers)) + assert len(frozenset(w.id for w in registered)) == 1 + assert await repo.worker(token_hash) == registered[0] + await repo.revoke_worker(registered[0].id) + restored: Final = await repo.configure_service_worker(workers[-1], token_hash) + assert restored.id == registered[0].id + assert restored.revoked is False + finally: + await lens_db.execute_raw('DELETE FROM "LiteLLM_LensWorker" WHERE token_hash=$1', token_hash) + + +@pytest.mark.asyncio +async def test_claim_pages_only_yield_work_the_worker_can_claim(lens_db: Prisma) -> None: + now: Final = datetime.now(timezone.utc) + scope: Final = Scope(team_id=uuid4().hex) + repo: Final = LensRepository(WriterDatabase(PrismaWrapper(lens_db))) + prefix: Final = uuid4().hex + base: Final = Lens( + id=prefix, + scope=scope, + settings=LensSettings( + name="Candidate pagination", model="test", enabled=False, context="Find repeated failures" + ), + created_at=now, + next_run_at=now + timedelta(days=1), + budget_month=now.strftime("%Y-%m"), + ) + queued: Final = tuple( + queue_job(base.model_copy(update={"id": f"{prefix}-{i:03d}"}), now, uuid4().hex) for i in range(52) + ) + other_scope: Final = queued[0].model_copy(update={"id": f"{prefix}-other", "scope": Scope(team_id=uuid4().hex)}) + due: Final = base.model_copy( + update={ + "id": f"{prefix}-due", + "settings": base.settings.model_copy(update={"enabled": True}), + "next_run_at": now, + } + ) + live: Final = claim_job(queued[0], Worker(id=prefix, name="worker", scope=scope, last_seen=now), now) + expired: Final = live.model_copy( + update={ + "id": f"{prefix}-expired", + "jobs": (live.jobs[0].model_copy(update={"lease_until": now - timedelta(seconds=1)}),), + } + ) + rows: Final = (*queued[1:], live, base, due, expired, other_scope) + try: + for row in rows: + await repo.create(row) + found: Final = tuple([candidate async for candidate in repo.claim_candidates(scope, now)]) + assert frozenset(candidate.id for candidate in found) == frozenset( + candidate.id for candidate in (*queued[1:], due, expired) + ) + assert len(found) == 53 + finally: + await lens_db.execute_raw('DELETE FROM "LiteLLM_Lens" WHERE id LIKE $1', prefix + "%") + + @pytest.mark.asyncio async def test_trace_findings_include_archived_assessments_without_counting_retries_or_counterexamples( lens_db: Prisma, diff --git a/tests/unit/proxy/lens/test_endpoints.py b/tests/unit/proxy/lens/test_endpoints.py index 34d51110761..2238da7e00f 100644 --- a/tests/unit/proxy/lens/test_endpoints.py +++ b/tests/unit/proxy/lens/test_endpoints.py @@ -2,6 +2,7 @@ from datetime import datetime, timedelta, timezone from types import SimpleNamespace from typing import Final +import httpx import pytest from fastapi import HTTPException from pydantic import ValidationError @@ -38,6 +39,8 @@ from litellm.proxy.lens.models import ( ) from litellm.proxy.lens.repository import Row from litellm.proxy.lens.state import claim_job, queue_job, replace_job +from litellm.rust_bridge.trace.storage import ClickHouseStorage +from litellm.tracing.remote import RemoteTraceStore from tests.unit.proxy.lens.test_agent_workspace import execution from tests.unit.proxy.lens.test_state import NOW, lens, worker @@ -64,6 +67,54 @@ class ResultDatabase: return len(self.completed) +@pytest.mark.asyncio +@pytest.mark.parametrize("change", ("cancelled", "expired", "reclaimed", "reassigned")) +async def test_result_cannot_commit_after_losing_ownership_during_evidence_validation( + monkeypatch: pytest.MonkeyPatch, change: str +) -> None: + from litellm.proxy import proxy_server + from tests.unit.proxy.lens.test_state import finding + + claimed: Final = claim_job(queue_job(lens(), NOW, "job"), worker(), NOW) + active: Final = claimed.jobs[0].model_copy( + update={ + "lease_until": datetime.max.replace(tzinfo=timezone.utc), + "sample": Sample(executions=(execution("run"),), eligible=1), + } + ) + competing: Final = active.model_copy( + update={ + "status": "cancelled" if change == "cancelled" else "running", + "lease_until": NOW if change == "expired" else active.lease_until, + "attempts": 2 if change == "reclaimed" else 1, + "worker_id": "other-worker" if change == "reassigned" else active.worker_id, + } + ) + db: Final = ResultDatabase(replace_job(claimed, active)) + monkeypatch.setattr(proxy_server, "prisma_client", SimpleNamespace(db=db)) + + async def evidence(request: httpx.Request) -> httpx.Response: + db.stored = replace_job(db.stored, competing) + return httpx.Response(200, json={"data": [{"count": 1}]}) + + async with httpx.AsyncClient(base_url="http://lens.test", transport=httpx.MockTransport(evidence)) as client: + saved: Final = await result( + "lens", + "job", + Result( + coverage=Coverage(screened=1, investigated=1), + findings=(finding("run"),), + assessments=(RunAssessment(execution_id="run"),), + review_versions=(ReviewVersion(execution_id="run", content_version="v1"),), + ), + worker(), + ClickHouseStorage(RemoteTraceStore(client)), + ) + assert saved.jobs[0] == competing + assert saved.findings == () + assert db.completed == () + + @pytest.mark.asyncio @pytest.mark.parametrize( "selected,check_id,quoted", diff --git a/ui/litellm-dashboard/src/components/lens/LensSetup.integration.test.tsx b/ui/litellm-dashboard/src/components/lens/LensSetup.integration.test.tsx index 7d41939112b..205de9207fb 100644 --- a/ui/litellm-dashboard/src/components/lens/LensSetup.integration.test.tsx +++ b/ui/litellm-dashboard/src/components/lens/LensSetup.integration.test.tsx @@ -22,6 +22,12 @@ function serve({ enabled = false, traces = false, requests = false, connected = list.mockResolvedValue({ lenses: [], workers: connected ? [worker()] : [], tracing_enabled: enabled }); network.mockImplementation(async (input, init) => { const { path, method, body, query } = await readRequest(input, init); + if (path === "/lens/service") + return Response.json({ + url: "https://traces.test", + connected: true, + status: { storage_ready: true, credentials_ready: true }, + }); if (path === "/v1/traces") return enabled ? Response.json({ data: traces ? [data.runs[0].trace.summary] : [] }) @@ -37,7 +43,8 @@ function serve({ enabled = false, traces = false, requests = false, connected = if (path === "/key/generate") return Response.json({ token_id: worker().analysis_key_id }); if (path === "/lens/workers/register") { list.mockResolvedValue({ lenses: [], workers: [worker()], tracing_enabled: true }); - return Response.json({ worker: worker(), token: "test-worker-token", image: "test-worker-image" }); + const created = { worker: worker(), token: "", image: "test-worker-image", managed: true }; + return Response.json(created); } if (path === "/models") return Response.json({ data: [{ id: "analysis" }] }); if (path === "/model_group/info") @@ -77,7 +84,7 @@ async function connectWorkerFromSettings(user: ReturnType { await expectUrl(onUrlUpdate, (url) => expect(url.get("tab")).toBe("settings")); const panel = within(await screen.findByRole("region", { name: "Settings" })); expect(panel.getByRole("heading", { name: "Connect a worker" })).toBeVisible(); - expect(panel.getByRole("button", { name: "Get install command" })).toBeVisible(); + expect(panel.getByRole("button", { name: "Enable investigations" })).toBeVisible(); }); it("keeps a pending worker install across tab switches and offers the first investigation once it connects", async () => { @@ -349,7 +349,8 @@ describe("Lens interactive demo", () => { if (path === "/lens") return Response.json({ lenses: [], workers: workers(), tracing_enabled: true }); if (path === "/lens/workers/register" && method === "POST") { workers.mockReturnValue([worker]); - return Response.json({ token: "lens-test-token", image: "lens-worker:v1", worker }); + const created = { token: "", managed: true, image: "lens-worker:v1", worker }; + return Response.json(created); } if (path === "/key/list") return Response.json({ keys: [{ token, key_alias: "Analysis" }], total_pages: 1 }); if (path === "/key/info") return Response.json({ info: { models: [], max_budget: null } }); @@ -368,14 +369,28 @@ describe("Lens interactive demo", () => { await user.click(panel.getByRole("switch", { name: "Use an existing virtual key" })); await user.click(panel.getByRole("combobox", { name: "Charge analysis to" })); await user.click(await screen.findByRole("option", { name: "Analysis" })); - await user.click(panel.getByRole("button", { name: "Get install command" })); - expect(await panel.findByText("Waiting for your worker to connect…")).toBeInTheDocument(); + await user.click(panel.getByRole("button", { name: "Enable investigations" })); + expect( + await panel.findByText( + "Connecting your Lens service… This page updates automatically. Check the service logs if it does not connect.", + ), + ).toBeInTheDocument(); const tabs = within(screen.getByRole("tablist", { name: "Lens" })); await user.click(tabs.getByRole("tab", { name: "Traces" })); - await waitFor(() => expect(panel.getByText("Waiting for your worker to connect…")).not.toBeVisible()); + await waitFor(() => + expect( + panel.getByText( + "Connecting your Lens service… This page updates automatically. Check the service logs if it does not connect.", + ), + ).not.toBeVisible(), + ); await user.click(tabs.getByRole("tab", { name: "Settings" })); - expect(panel.getByText("Waiting for your worker to connect…")).toBeVisible(); - expect(panel.getByLabelText("Docker command preview")).toHaveTextContent("LENS_WORKER_TOKEN=lens-test-token"); + expect( + panel.getByText( + "Connecting your Lens service… This page updates automatically. Check the service logs if it does not connect.", + ), + ).toBeVisible(); + expect(panel.queryByLabelText("Docker command preview")).not.toBeInTheDocument(); workers.mockReturnValue([{ ...worker, last_seen: new Date().toISOString() }]); await testQueryClient.refetchQueries({ queryKey: lensKeys.lists() }); expect(await panel.findByRole("heading", { name: "Worker connected" })).toBeVisible(); diff --git a/ui/litellm-dashboard/src/components/lens/onboarding/tracing/TracingSetupCard.integration.test.tsx b/ui/litellm-dashboard/src/components/lens/onboarding/tracing/TracingSetupCard.integration.test.tsx index f2860e40c22..16a0a4fe946 100644 --- a/ui/litellm-dashboard/src/components/lens/onboarding/tracing/TracingSetupCard.integration.test.tsx +++ b/ui/litellm-dashboard/src/components/lens/onboarding/tracing/TracingSetupCard.integration.test.tsx @@ -3,7 +3,7 @@ import userEvent from "@testing-library/user-event"; import { beforeEach, describe, expect, it, vi } from "vitest"; import { chooseSelectOption, renderWithProviders } from "@/../tests/test-utils"; import { copyToClipboard } from "@/utils/dataUtils"; -import { agentTraceCall, apiClient, sendOtlpTraceCall } from "../../../networking"; +import { agentTraceCall, apiClient } from "../../../networking"; import { codingAgentCommand, codingAgentPrompt, @@ -17,15 +17,14 @@ import type { Trace } from "../../traces/types"; vi.mock("../../../networking", () => ({ getProxyBaseUrl: () => "http://proxy.test/", - sendOtlpTraceCall: vi.fn(), agentTraceCall: vi.fn(), - apiClient: { post: vi.fn() }, + apiClient: { post: vi.fn(), get: vi.fn() }, })); vi.mock("@/utils/dataUtils", () => ({ copyToClipboard: vi.fn().mockResolvedValue(true) })); const SECRET = "sk-abcdefghijklmnopWXYZ"; -const renderCard = ( +const renderCard = async ( props: { detail?: string | null; connected?: boolean; @@ -46,49 +45,62 @@ const renderCard = ( onOpenTrace={onOpenTrace} />, ); + if (!props.detail) await screen.findByRole("combobox", { name: "Your agent framework" }); return { onOpenTrace, card: screen.getByTestId("tracing-setup-card") }; }; -beforeEach(() => vi.clearAllMocks()); +const network = vi.fn(); +beforeEach(() => { + vi.clearAllMocks(); + vi.stubGlobal("fetch", network); + network.mockResolvedValue(Response.json({})); + vi.mocked(apiClient.get).mockResolvedValue({ + url: "https://traces.test", + connected: true, + status: { storage_ready: true, credentials_ready: true }, + }); + vi.mocked(apiClient.post).mockResolvedValue({ key: SECRET, active: true }); +}); describe("TracingSetupCard", () => { it("guides agent connection while waiting for the first trace", async () => { const user = userEvent.setup(); - const { card } = renderCard(); + const { card } = await renderCard(); expect(screen.getByRole("heading", { name: "Connect your agent" })).toBeVisible(); expect(screen.getByText("Tracing enabled")).toBeVisible(); expect(screen.getByText("Waiting for your first trace")).toBeVisible(); expect(screen.queryByRole("button", { name: "Preview sample" })).not.toBeInTheDocument(); - expect(sendOtlpTraceCall).not.toHaveBeenCalled(); + expect(network).not.toHaveBeenCalled(); expect(card).not.toHaveTextContent("store: clickhouse"); await user.click(screen.getByText("Set up manually")); - expect(screen.getByText(/^export OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=/)).toBeVisible(); + expect(screen.getByText(/^export LITELLM_TRACING_KEY=/)).toBeVisible(); expect(card).not.toHaveTextContent(/langsmith/i); }); it("keeps connection details visible and copies the full trace endpoint", async () => { const user = userEvent.setup(); - renderCard(); + await renderCard(); expect(screen.getByRole("combobox", { name: "Your agent framework" })).toBeVisible(); - expect(screen.getByRole("button", { name: "Copy http://proxy.test/v1/traces" })).toBeVisible(); - await user.click(screen.getByRole("button", { name: "Copy http://proxy.test/v1/traces" })); - expect(copyToClipboard).toHaveBeenLastCalledWith("http://proxy.test/v1/traces"); + expect(screen.getByRole("button", { name: "Copy https://traces.test/v1/traces" })).toBeVisible(); + await user.click(screen.getByRole("button", { name: "Copy https://traces.test/v1/traces" })); + expect(copyToClipboard).toHaveBeenLastCalledWith("https://traces.test/v1/traces"); }); - it("shows connection guidance for another agent without a demo", () => { - renderCard({ connected: true }); + it("shows connection guidance for another agent without a demo", async () => { + await renderCard({ connected: true }); expect(screen.getByRole("heading", { name: "Connect another agent" })).toBeVisible(); expect(screen.queryByRole("button", { name: "Preview sample" })).not.toBeInTheDocument(); }); it("builds the coding agent command for the selected framework and keeps both manual installers", async () => { const user = userEvent.setup(); - renderCard(); + await renderCard(); await chooseSelectOption(user, screen.getByRole("combobox", { name: "Your agent framework" }), "CrewAI"); const prompt = codingAgentPrompt( "http://proxy.test", + "https://traces.test", FRAMEWORKS.find((guide) => guide.id === "crewai")!, - "openai/gpt-6-sol", + "openai/gpt-6.1-sol", ); expect(screen.getByText(/^claude /)).not.toBeVisible(); await user.click(screen.getByRole("button", { name: "Copy setup command" })); @@ -114,7 +126,7 @@ describe("TracingSetupCard", () => { it("uses the selected framework's tracing and agent name without asking for a model", async () => { const user = userEvent.setup(); - const { card } = renderCard(); + const { card } = await renderCard(); await chooseSelectOption(user, screen.getByRole("combobox", { name: "Your agent framework" }), "Vercel AI SDK"); await user.click(screen.getByText("Set up manually")); expect(screen.queryByRole("combobox", { name: "Model" })).not.toBeInTheDocument(); @@ -123,7 +135,7 @@ describe("TracingSetupCard", () => { expect(card).toHaveTextContent("functionId: AGENT_NAME"); expect(card).toHaveTextContent("Use a model configured on this proxy."); expect(screen.getByText(/^import \{ createOpenAICompatible/)).toHaveTextContent( - 'const model = litellm("openai/gpt-6-sol")', + 'const model = litellm("openai/gpt-6.1-sol")', ); expect(card).toHaveTextContent('baseURL: "http://proxy.test/v1"'); }); @@ -131,7 +143,7 @@ describe("TracingSetupCard", () => { it("keeps plugin model settings and uses a generated tracing key only for tracing", async () => { const user = userEvent.setup(); vi.mocked(apiClient.post).mockResolvedValue({ key: SECRET }); - const { card } = renderCard(); + const { card } = await renderCard(); await chooseSelectOption(user, screen.getByRole("combobox", { name: "Your agent framework" }), "Hermes"); await user.click(screen.getByText("Set up manually")); expect(screen.queryByRole("combobox", { name: "Model" })).not.toBeInTheDocument(); @@ -139,39 +151,38 @@ describe("TracingSetupCard", () => { await user.click(screen.getByRole("button", { name: "Generate tracing key" })); await screen.findByText("Your tracing key"); expect(card).toHaveTextContent("gen_ai.agent.name: research_agent"); - expect(card).toHaveTextContent("endpoint: http://proxy.test/v1/traces"); + expect(card).toHaveTextContent("endpoint: https://traces.test/v1/traces"); expect(card).toHaveTextContent('Authorization: "Bearer ${LITELLM_TRACING_KEY}"'); expect(card).not.toHaveTextContent(SECRET); }); - it("hides the actions a read-only viewer cannot perform", () => { - const { card } = renderCard({ readOnly: true }); + it("hides the actions a read-only viewer cannot perform", async () => { + const { card } = await renderCard({ readOnly: true }); expect(screen.queryByRole("button", { name: "Send a test trace" })).not.toBeInTheDocument(); expect(screen.queryByRole("button", { name: "Generate tracing key" })).not.toBeInTheDocument(); - expect(card).toHaveTextContent("ask a proxy admin for one"); + expect(card).toHaveTextContent("Ask your proxy admin for a dedicated Lens tracing key."); expect(card).toHaveTextContent("Connection details"); }); - it("offers a scoped tracing key only to callers allowed to set key routes", () => { - const { card } = renderCard({ canMintTracingKey: false }); + it("offers a scoped tracing key only to callers allowed to create tracing keys", async () => { + const { card } = await renderCard({ canMintTracingKey: false }); expect(screen.queryByRole("button", { name: "Generate tracing key" })).not.toBeInTheDocument(); expect(screen.getByRole("button", { name: "Send a test trace" })).toBeVisible(); - expect(card).toHaveTextContent("Use any LiteLLM virtual key you already have"); + expect(card).toHaveTextContent("Ask your proxy admin for a dedicated Lens tracing key."); }); it("generates a tracing key that stays masked on screen but copies in full", async () => { const user = userEvent.setup(); vi.mocked(apiClient.post).mockResolvedValue({ key: SECRET }); - const { card } = renderCard(); + const { card } = await renderCard(); await user.click(screen.getByText("Set up manually")); await user.click(screen.getByRole("button", { name: "Generate tracing key" })); expect(await screen.findByText("Your tracing key")).toBeVisible(); - expect(apiClient.post).toHaveBeenCalledWith("/key/generate", { + expect(apiClient.post).toHaveBeenCalledWith("/lens/tracing/keys", { accessToken: "sk-admin", body: TRACING_KEY_REQUEST, }); - expect(TRACING_KEY_REQUEST.allowed_routes).toEqual(["/v1/traces"]); expect(card).not.toHaveTextContent(SECRET); expect(card).toHaveTextContent(maskSecret(SECRET)); await user.click(screen.getAllByRole("button", { name: "Copy" })[0]); @@ -181,23 +192,33 @@ describe("TracingSetupCard", () => { it("sends a test trace, waits for it to land, then opens it", async () => { const user = userEvent.setup(); const summary = { trace_id: "abc", name: "weather_agent" } as Trace["summary"]; - vi.mocked(sendOtlpTraceCall).mockResolvedValue(undefined); + network.mockResolvedValue(Response.json({})); vi.mocked(agentTraceCall).mockResolvedValue({ summary, agents: [], spans: [] } as unknown as Trace); - const { onOpenTrace } = renderCard(); + const { onOpenTrace } = await renderCard(); + await user.click(screen.getByRole("button", { name: "Generate tracing key" })); + await screen.findByText("Your tracing key"); await user.click(screen.getByRole("button", { name: "Send a test trace" })); await user.click(await screen.findByRole("button", { name: /View trace/ })); - expect(sendOtlpTraceCall).toHaveBeenCalledOnce(); + const uploadOptions = { + method: "POST", + credentials: "omit", + redirect: "error", + headers: { Accept: "application/json", "Content-Type": "application/json", Authorization: `Bearer ${SECRET}` }, + }; + expect(network).toHaveBeenCalledWith("https://traces.test/v1/traces", expect.objectContaining(uploadOptions)); expect(vi.mocked(agentTraceCall).mock.calls[0][1]).toMatch(/^[0-9a-f]{32}$/); expect(onOpenTrace).toHaveBeenCalledWith(summary); }); it("reports a failed send instead of claiming success", async () => { const user = userEvent.setup(); - vi.mocked(sendOtlpTraceCall).mockRejectedValue(new Error("boom")); - renderCard(); + network.mockRejectedValue(new Error("boom")); + await renderCard(); + await user.click(screen.getByRole("button", { name: "Generate tracing key" })); + await screen.findByText("Your tracing key"); await user.click(screen.getByRole("button", { name: "Send a test trace" })); expect(await screen.findByText("Could not send the test trace.")).toBeVisible(); @@ -208,10 +229,10 @@ describe("TracingSetupCard", () => { it("guides proxy setup before agent setup and allows checking readiness", async () => { const user = userEvent.setup(); const onCheck = vi.fn(); - const { card } = renderCard({ detail: "Agent tracing is not enabled", onCheck }); + const { card } = await renderCard({ detail: "Agent tracing is not enabled", onCheck }); expect(screen.getByRole("heading", { name: "Enable tracing" })).toBeVisible(); - expect(card).toHaveTextContent("type: clickhouse"); - expect(card).toHaveTextContent("url: os.environ/CLICKHOUSE_URL"); + expect(card).toHaveTextContent("LITELLM_LENS_URL"); + expect(card).toHaveTextContent("LITELLM_LENS_SERVICE_TOKEN"); expect(screen.queryByRole("combobox", { name: "Your agent framework" })).not.toBeInTheDocument(); expect(screen.queryByRole("button", { name: "Send a test trace" })).not.toBeInTheDocument(); await user.click(screen.getByRole("button", { name: "Check setup" })); @@ -222,18 +243,18 @@ describe("TracingSetupCard", () => { describe("setup snippets", () => { it("uses the instance trace endpoint and keeps tracing and inference keys separate", () => { - const env = tracingEnvSnippet("http://proxy.test"); - expect(env).toContain('OTEL_EXPORTER_OTLP_TRACES_ENDPOINT="http://proxy.test/v1/traces"'); + const env = tracingEnvSnippet("https://traces.test"); + expect(env).toContain('OTEL_EXPORTER_OTLP_TRACES_ENDPOINT="https://traces.test/v1/traces"'); expect(env).toContain('OTEL_EXPORTER_OTLP_PROTOCOL="http/protobuf"'); expect(env).not.toContain("export LITELLM_API_KEY="); - expect(env).toContain("Bearer $LITELLM_API_KEY"); - const withKey = tracingEnvSnippet("http://proxy.test", SECRET); - expect(withKey).toContain(`export LITELLM_TRACING_KEY=${SECRET}\n`); + expect(env).toContain("Bearer $LITELLM_TRACING_KEY"); + const withKey = tracingEnvSnippet("https://traces.test", SECRET); + expect(withKey).toContain(`export LITELLM_TRACING_KEY="${SECRET}"\n`); expect(withKey).toContain('OTEL_EXPORTER_OTLP_TRACES_HEADERS="Authorization=Bearer $LITELLM_TRACING_KEY"'); expect(withKey).not.toContain("export LITELLM_API_KEY="); - const prompt = codingAgentPrompt("http://proxy.test", FRAMEWORKS[0], "openai/gpt-6-sol"); - expect(prompt).toContain('OTEL_EXPORTER_OTLP_TRACES_ENDPOINT="http://proxy.test/v1/traces"'); + const prompt = codingAgentPrompt("http://proxy.test", "https://traces.test", FRAMEWORKS[0], "openai/gpt-6.1-sol"); + expect(prompt).toContain('OTEL_EXPORTER_OTLP_TRACES_ENDPOINT="https://traces.test/v1/traces"'); expect(prompt).toContain("Keep the existing model configuration"); expect(prompt).toContain('AGENT_NAME = "research_agent"'); expect(prompt).toContain("name=AGENT_NAME"); diff --git a/ui/litellm-dashboard/src/components/lens/onboarding/tracing/TracingSetupCard.tsx b/ui/litellm-dashboard/src/components/lens/onboarding/tracing/TracingSetupCard.tsx index f8c4cd96bcd..33009ea77f7 100644 --- a/ui/litellm-dashboard/src/components/lens/onboarding/tracing/TracingSetupCard.tsx +++ b/ui/litellm-dashboard/src/components/lens/onboarding/tracing/TracingSetupCard.tsx @@ -4,6 +4,7 @@ import { ArrowRight, ArrowUpRight, Check, Copy, KeyRound, Loader2, Send } from " import { useState } from "react"; import { useQuery } from "@tanstack/react-query"; import type { components } from "@/lib/http/schema"; +import { createApiClient, type RequestOptions } from "@/lib/http/client"; import { useTimeout } from "usehooks-ts"; import { cn } from "@/lib/cva.config"; @@ -213,15 +214,15 @@ function SendTestTrace({ const sample = sampleTraceExport(Date.now()); try { if (!tracingKey) throw new Error("Generate a tracing key first"); - const response = await fetch(`${traceUrl}/v1/traces`, { - method: "POST", + const client = createApiClient({ getBaseUrl: () => traceUrl }); + const options: RequestOptions = { credentials: "omit", redirect: "error", - headers: { "Content-Type": "application/json", Authorization: `Bearer ${tracingKey}` }, - body: JSON.stringify(sample.body), + accessToken: tracingKey, + body: sample.body, signal: AbortSignal.timeout(15000), - }); - if (!response.ok) throw new Error(`Trace upload failed (HTTP ${response.status})`); + }; + await client.post("/v1/traces", options); } catch { setState({ kind: "failed", message: "Could not send the test trace." }); return; @@ -309,6 +310,7 @@ function TracingKey({ }) { const [creating, setCreating] = useState(false); const [error, setError] = useState(""); + const [pendingActivation, setPendingActivation] = useState(false); const create = async () => { setCreating(true); setError(""); @@ -318,6 +320,7 @@ function TracingKey({ body: TRACING_KEY_REQUEST, }); if (!result.key) throw new Error("The proxy did not return the new key"); + setPendingActivation(!result.active); onCreated(result.key); } catch (cause) { setError(cause instanceof Error ? cause.message : "Could not create a key"); @@ -328,6 +331,11 @@ function TracingKey({ if (tracingKey) { return (
+ {pendingActivation && ( +

+ Key saved. Lens has not confirmed it yet. Once the service is connected, keys sync within 30 seconds. +

+ )} Your tracing key} />

Hidden for safety. Copy copies the full key, and the environment step below includes it. This key can only diff --git a/ui/litellm-dashboard/src/components/lens/settings/worker/AnalysisKeyPicker.integration.test.tsx b/ui/litellm-dashboard/src/components/lens/settings/worker/AnalysisKeyPicker.integration.test.tsx index 2873c6f18e6..93b330153dd 100644 --- a/ui/litellm-dashboard/src/components/lens/settings/worker/AnalysisKeyPicker.integration.test.tsx +++ b/ui/litellm-dashboard/src/components/lens/settings/worker/AnalysisKeyPicker.integration.test.tsx @@ -16,7 +16,6 @@ function AnalysisKeyPickerForm() { useExisting: true, analysisKey: null, access: { model: null, budget: "100" }, - address: "http://localhost:4000", }, }); return ( diff --git a/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerForm.tsx b/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerForm.tsx index fb1e7a57ee8..b379c77862c 100644 --- a/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerForm.tsx +++ b/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerForm.tsx @@ -1,7 +1,6 @@ "use client"; import { Controller, useFormContext, useWatch } from "react-hook-form"; -import { Input } from "@/components/ui/input"; import { Switch } from "@/components/ui/switch"; import { AnalysisKeyPicker } from "./AnalysisKeyPicker"; @@ -9,11 +8,7 @@ import { AnalysisAccessFields } from "./AnalysisAccessFields"; import type { WorkerFormInput } from "./workerSchema"; export function WorkerForm({ editingWorker }: { editingWorker: string | null }) { - const { - control, - register, - formState: { errors }, - } = useFormContext(); + const { control } = useFormContext(); const useExisting = useWatch({ control, name: "useExisting" }); return (

@@ -29,16 +24,6 @@ export function WorkerForm({ editingWorker }: { editingWorker: string | null }) render={({ field }) => } /> - {!editingWorker && ( -
- - -

Your server must be able to reach this address.

- {errors.address?.message &&

{errors.address.message}

} -
- )}
diff --git a/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerInstall.tsx b/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerInstall.tsx index 16346c882e7..d166bc5f132 100644 --- a/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerInstall.tsx +++ b/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerInstall.tsx @@ -1,96 +1,27 @@ "use client"; import type { ComponentProps, ReactNode } from "react"; -import { useMutation } from "@tanstack/react-query"; -import { CheckCircle2, Copy, Loader2 } from "lucide-react"; -import { Button } from "@/components/ui/button"; +import { CheckCircle2 } from "lucide-react"; import { cn } from "@/lib/cva.config"; -import type { WorkerCreated } from "../../model/types"; import { SettingsCard } from "../SettingsSection"; -import { workerSetupCommand } from "./workerCommand"; - -const CLIPBOARD_FAILED = "Clipboard access failed. Allow clipboard access and try again."; - -function useCopy() { - return useMutation({ retry: false, mutationFn: (text: string) => navigator.clipboard.writeText(text) }); -} - -function InstallSteps({ address, created }: { address: string; created: WorkerCreated }) { - const command = workerSetupCommand(address, created.token, created.image); - const copyCommand = useCopy(); - const copyToken = useCopy(); - if (created.managed) - return ( -

- Connecting your Lens service… This page updates automatically. Check the service logs if it does not connect. -

- ); - return ( - <> - -
- View command -

Contains a private worker token.

-
-          {command}
-        
-
-
- Using Docker Compose or Helm? -

- Save this private token as LENS_WORKER_TOKEN in Compose or in your Helm worker token secret. Keep it for - future upgrades. -

- -
-
-
- Waiting for your worker to connect… -
-
- Not connecting? -

- Check that Docker is running and can reach {address}. Inspect the container logs for connection or - authentication errors. This page updates automatically. -

-
-
- {(copyCommand.isError || copyToken.isError) && ( -

- {CLIPBOARD_FAILED} -

- )} - - ); -} export type WorkerInstallProps = ComponentProps<"div"> & { - address: string; - created: WorkerCreated; connected: boolean; /** Rendered once the worker connects, in place of the install steps. */ children: ReactNode; }; -export function WorkerInstall({ address, created, connected, children, className, ...props }: WorkerInstallProps) { +export function WorkerInstall({ connected, children, className, ...props }: WorkerInstallProps) { if (!connected) return (
-

Run the worker

-

- Your Lens service connects automatically. A separately hosted worker can use the Docker command. -

+

Connecting Lens

+

Your Lens service connects automatically.

- +

+ Connecting your Lens service… This page updates automatically. Check the service logs if it does not connect. +

); return ( diff --git a/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerSettings.integration.test.tsx b/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerSettings.integration.test.tsx index 16de3a834a8..7f25bef85d2 100644 --- a/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerSettings.integration.test.tsx +++ b/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerSettings.integration.test.tsx @@ -30,7 +30,8 @@ const calls = (method: string, path: string) => const writes = () => sent.filter((request) => request.method !== "GET"); const created = { - token: "lens-test-token", + token: "", + managed: true, image: "ghcr.io/berriai/litellm-lens-worker:v1.2.3", worker: { id: "worker", @@ -63,32 +64,21 @@ describe("Worker setup", () => { vi.stubGlobal("fetch", network); serve(keyRoute); }); - it("generates a complete command using one worker credential and the configured proxy address", async () => { + it("enables the installed Lens service without exposing a worker credential", async () => { serve((request) => (request.path === "/lens/workers/register" ? created : keyRoute(request))); const user = userEvent.setup(); const { rerender } = renderWithLens(, { accessToken: "admin" }); await user.click(screen.getByText("Advanced options")); await user.click(screen.getByRole("switch", { name: "Use an existing virtual key" })); - expect(screen.getByRole("textbox", { name: "LiteLLM proxy URL" })).toHaveValue("https://gateway.example/proxy"); - expect(screen.getByRole("button", { name: "Get install command" })).toBeDisabled(); + expect(screen.getByRole("button", { name: "Enable investigations" })).toBeDisabled(); await user.click(screen.getByRole("combobox", { name: "Charge analysis to" })); await user.click(await screen.findByRole("option", { name: "Analysis" })); - await user.click(screen.getByRole("button", { name: "Get install command" })); + await user.click(screen.getByRole("button", { name: "Enable investigations" })); expect(calls("POST", "/lens/workers/register").map(({ body }) => body)).toEqual([ - { name: "Lens worker", analysis_key_id: "b".repeat(64) }, + { name: "Lens worker", analysis_key_id: "b".repeat(64), managed: true }, ]); - expect(screen.getByRole("status")).toHaveTextContent("Waiting for your worker to connect"); - expect(screen.getByLabelText("Docker command preview")).not.toBeVisible(); - await user.click(screen.getByRole("button", { name: "Copy Docker command" })); - const command = await navigator.clipboard.readText(); - expect(command).toContain("LITELLM_URL=https://gateway.example/proxy"); - expect(command).toContain("LENS_WORKER_TOKEN=lens-test-token"); - expect(command).toContain("--add-host host.docker.internal:host-gateway"); - expect(command).toContain(created.image); - await user.click(screen.getByText("Using Docker Compose or Helm?")); - await user.click(screen.getByRole("button", { name: "Copy worker token" })); - expect(await navigator.clipboard.readText()).toBe(created.token); - expect(await screen.findByRole("button", { name: "Token copied" })).toBeVisible(); + expect(screen.getByRole("status")).toHaveTextContent("Connecting your Lens service"); + expect(screen.queryByRole("button", { name: "Copy Docker command" })).not.toBeInTheDocument(); rerender(); expect(screen.getByRole("heading", { name: "Worker connected" })).toBeVisible(); expect(screen.queryByRole("status")).not.toBeInTheDocument(); @@ -135,10 +125,10 @@ describe("Worker setup", () => { renderWithLens(, { accessToken: "admin" }); const revoke = await screen.findByRole("button", { name: "Revoke access" }); expect(screen.queryByRole("button", { name: "Add worker" })).not.toBeInTheDocument(); - expect(screen.queryByRole("button", { name: "Get install command" })).not.toBeInTheDocument(); + expect(screen.queryByRole("button", { name: "Enable investigations" })).not.toBeInTheDocument(); await user.click(revoke); expect(calls("DELETE", "/lens/workers/worker")).toHaveLength(1); - expect(await screen.findByRole("button", { name: "Get install command" })).toBeDisabled(); + expect(await screen.findByRole("button", { name: "Enable investigations" })).toBeDisabled(); expect(screen.getByRole("combobox", { name: "Analysis model" })).toBeVisible(); expect(listCalls()).toBe(2); }); @@ -158,13 +148,12 @@ describe("Worker setup", () => { return { keys: [], total_pages: 0 }; }); renderWithLens(, { accessToken: "admin" }); - expect(screen.getByRole("button", { name: "Get install command" })).toBeDisabled(); - expect(screen.getByRole("textbox", { name: "LiteLLM proxy URL", hidden: true })).not.toBeVisible(); + expect(screen.getByRole("button", { name: "Enable investigations" })).toBeDisabled(); await user.click(screen.getByRole("combobox", { name: "Analysis model" })); await user.click(await screen.findByRole("option", { name: "analysis-model" })); await user.clear(screen.getByLabelText("Monthly limit (USD)")); await user.type(screen.getByLabelText("Monthly limit (USD)"), "12"); - await user.click(screen.getByRole("button", { name: "Get install command" })); + await user.click(screen.getByRole("button", { name: "Enable investigations" })); expect(await screen.findByRole("alert")).toHaveTextContent("Registration unavailable"); expect(writes()[0]).toMatchObject({ path: "/key/generate", @@ -176,8 +165,8 @@ describe("Worker setup", () => { metadata: { purpose: "lens" }, }, }); - await user.click(screen.getByRole("button", { name: "Get install command" })); - expect(await screen.findByRole("status")).toHaveTextContent("Waiting for your worker"); + await user.click(screen.getByRole("button", { name: "Enable investigations" })); + expect(await screen.findByRole("status")).toHaveTextContent("Connecting your Lens service"); expect(calls("POST", "/key/delete").map(({ body }) => body)).toEqual([{ keys: ["limited-key-id"] }]); expect(writes().map(({ path }) => path)).toEqual([ "/key/generate", @@ -186,7 +175,7 @@ describe("Worker setup", () => { "/key/generate", "/lens/workers/register", ]); - expect(writes().at(-1)?.body).toEqual({ name: "Lens worker", analysis_key_id: "retry-key-id" }); + expect(writes().at(-1)?.body).toEqual({ name: "Lens worker", analysis_key_id: "retry-key-id", managed: true }); expect(screen.queryByText("sk-secret-not-displayed")).not.toBeInTheDocument(); }); }); diff --git a/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerSettings.tsx b/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerSettings.tsx index 8a30908e51c..fb151548444 100644 --- a/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerSettings.tsx +++ b/ui/litellm-dashboard/src/components/lens/settings/worker/WorkerSettings.tsx @@ -9,7 +9,6 @@ import type { LensList, Worker } from "../../model/types"; import { useWorkerConnected } from "../../hooks/useWorkerConnected"; import { SettingsCard } from "../SettingsSection"; import { usePrepareWorker } from "./usePrepareWorker"; -import { initialProxyAddress } from "./workerCommand"; import { WorkerForm } from "./WorkerForm"; import { WorkerInstall } from "./WorkerInstall"; import { WorkerList } from "./WorkerList"; @@ -21,7 +20,6 @@ function defaultWorkerFormValues(): WorkerFormInput { useExisting: false, analysisKey: null, access: { model: null, budget: "100" }, - address: typeof window === "undefined" ? "" : initialProxyAddress(), }; } @@ -36,7 +34,7 @@ function ErrorText({ message }: { message: string | undefined }) { function submitLabel(editing: Worker | null, busy: boolean): string { if (busy) return "Preparing…"; - return editing ? "Save analysis access" : "Get install command"; + return editing ? "Save analysis access" : "Enable investigations"; } function WorkerFormCard({ @@ -59,7 +57,9 @@ function WorkerFormCard({

{editing ? "Analysis access" : "Connect a worker"}

- {editing ? "Choose which key pays for analysis." : "Deploy the worker on your server to run investigations."} + {editing + ? "Choose which key pays for analysis." + : "Choose a model and spending limit. Your Lens service runs investigations automatically."}

@@ -111,7 +111,6 @@ export function WorkerSettings({ const submit = (editing: Worker | null) => form.handleSubmit((values) => { const registration = { - address: values.address, useExisting: values.useExisting, analysisKey: values.analysisKey, access: values.access, @@ -156,7 +155,7 @@ export function WorkerSettings({ ); case "install": return ( - + {readyAction ?? (