fix(lens): complete service routing and reject stale investigation results

This commit is contained in:
moe-berri 2026-10-07 12:19:52 -07:00
parent ac39101ba7
commit fc7707f168
44 changed files with 666 additions and 378 deletions

View file

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

View file

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

40
deploy/lens/smoke.sh Normal file
View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

@ -668,6 +668,7 @@ lensWorker:
publicUrl: ""
service:
port: 4318
annotations: {}
ingress:
enabled: false
className: ""

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

@ -650,6 +650,7 @@ lensWorker:
publicUrl: ""
service:
port: 4318
annotations: {}
ingress:
enabled: false
className: ""

View file

@ -4206,6 +4206,7 @@ dependencies = [
"typify",
"unicode-casefold",
"url",
"uuid",
"wiremock",
]

View file

@ -43,3 +43,4 @@ prettyplease = "0.2"
[dev-dependencies]
rstest.workspace = true
wiremock.workspace = true
uuid.workspace = true

View file

@ -98,6 +98,8 @@ pub fn router(state: Arc<State>) -> Router {
]),
);
public
.clone()
.nest("/lens-ingest", public)
.merge(
Router::new()
.route("/internal/read", post(read))

View file

@ -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::<serde_json::Value>().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"))

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

@ -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<typeof userEvent.setup
expect(settings.getByRole("heading", { name: "Connect a worker" })).toBeVisible();
await user.click(settings.getByRole("combobox", { name: "Analysis model" }));
await user.click(await screen.findByRole("option", { name: "analysis" }));
await user.click(settings.getByRole("button", { name: "Get install command" }));
await user.click(settings.getByRole("button", { name: "Enable investigations" }));
expect(await settings.findByRole("heading", { name: "Worker connected" })).toBeVisible();
await user.click(settings.getByRole("button", { name: "New investigation" }));
expect(await screen.findByRole("region", { name: "New investigation" })).toBeVisible();

View file

@ -328,7 +328,7 @@ describe("Lens interactive demo", () => {
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();

View file

@ -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<typeof fetch>();
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");

View file

@ -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 (
<div className="space-y-2">
{pendingActivation && (
<p role="status" className="text-sm text-muted-foreground">
Key saved. Lens has not confirmed it yet. Once the service is connected, keys sync within 30 seconds.
</p>
)}
<CodeBlock code={tracingKey} display={maskSecret(tracingKey)} tabs={<FileLabel>Your tracing key</FileLabel>} />
<p className="text-sm text-muted-foreground">
Hidden for safety. Copy copies the full key, and the environment step below includes it. This key can only

View file

@ -16,7 +16,6 @@ function AnalysisKeyPickerForm() {
useExisting: true,
analysisKey: null,
access: { model: null, budget: "100" },
address: "http://localhost:4000",
},
});
return (

View file

@ -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<WorkerFormInput>();
const { control } = useFormContext<WorkerFormInput>();
const useExisting = useWatch({ control, name: "useExisting" });
return (
<div className="min-w-0 space-y-5">
@ -29,16 +24,6 @@ export function WorkerForm({ editingWorker }: { editingWorker: string | null })
render={({ field }) => <Switch checked={field.value} onCheckedChange={field.onChange} />}
/>
</label>
{!editingWorker && (
<div className="space-y-2">
<label htmlFor="worker-proxy-address" className="block text-sm font-medium">
LiteLLM proxy URL
</label>
<Input id="worker-proxy-address" {...register("address")} />
<p className="text-xs text-muted-foreground">Your server must be able to reach this address.</p>
{errors.address?.message && <p className="text-sm text-destructive">{errors.address.message}</p>}
</div>
)}
</div>
</details>
</div>

View file

@ -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 (
<p role="status" className="text-sm text-muted-foreground">
Connecting your Lens service… This page updates automatically. Check the service logs if it does not connect.
</p>
);
return (
<>
<Button variant="default" className="w-full gap-2" onClick={() => copyCommand.mutate(command)}>
{copyCommand.isSuccess ? <CheckCircle2 className="size-4" /> : <Copy className="size-4" />}
{copyCommand.isSuccess ? "Copied" : "Copy Docker command"}
</Button>
<details className="text-sm">
<summary className="cursor-pointer text-muted-foreground">View command</summary>
<p className="mt-3 text-xs text-muted-foreground">Contains a private worker token.</p>
<pre
aria-label="Docker command preview"
className="mt-3 max-h-48 overflow-auto rounded-md bg-muted/40 p-3 text-xs leading-5"
>
{command}
</pre>
</details>
<details className="text-sm">
<summary className="cursor-pointer text-muted-foreground">Using Docker Compose or Helm?</summary>
<p className="mt-3 text-muted-foreground">
Save this private token as LENS_WORKER_TOKEN in Compose or in your Helm worker token secret. Keep it for
future upgrades.
</p>
<Button variant="outline" className="mt-3" onClick={() => copyToken.mutate(created.token)}>
{copyToken.isSuccess ? "Token copied" : "Copy worker token"}
</Button>
</details>
<div className="space-y-3 border-t pt-5">
<div role="status" className="flex items-center gap-2.5 text-sm">
<Loader2 className="size-4 shrink-0 animate-spin text-muted-foreground" /> Waiting for your worker to connect…
</div>
<details className="pl-6.5 text-sm text-muted-foreground">
<summary className="cursor-pointer">Not connecting?</summary>
<p className="mt-2 break-words leading-6">
Check that Docker is running and can reach {address}. Inspect the container logs for connection or
authentication errors. This page updates automatically.
</p>
</details>
</div>
{(copyCommand.isError || copyToken.isError) && (
<p role="alert" className="text-sm text-destructive">
{CLIPBOARD_FAILED}
</p>
)}
</>
);
}
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 (
<SettingsCard {...props} data-slot="worker-install" className={cn("flex flex-col gap-5", className)}>
<header className="space-y-1">
<h3 className="text-base font-semibold">Run the worker</h3>
<p className="text-sm text-muted-foreground">
Your Lens service connects automatically. A separately hosted worker can use the Docker command.
</p>
<h3 className="text-base font-semibold">Connecting Lens</h3>
<p className="text-sm text-muted-foreground">Your Lens service connects automatically.</p>
</header>
<InstallSteps address={address} created={created} />
<p role="status" className="text-sm text-muted-foreground">
Connecting your Lens service… This page updates automatically. Check the service logs if it does not connect.
</p>
</SettingsCard>
);
return (

View file

@ -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(<WorkerSettings workers={[]} />, { 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(<WorkerSettings workers={[{ ...created.worker, last_seen: new Date().toISOString() }]} />);
expect(screen.getByRole("heading", { name: "Worker connected" })).toBeVisible();
expect(screen.queryByRole("status")).not.toBeInTheDocument();
@ -135,10 +125,10 @@ describe("Worker setup", () => {
renderWithLens(<WorkerSettingsHost />, { 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(<WorkerSettings workers={[]} />, { 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();
});
});

View file

@ -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({
<header className="space-y-1">
<h3 className="text-base font-semibold">{editing ? "Analysis access" : "Connect a worker"}</h3>
<p className="text-sm text-muted-foreground">
{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."}
</p>
</header>
<WorkerForm editingWorker={editing?.id ?? null} />
@ -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 (
<WorkerInstall address={form.getValues("address")} created={screen.created} connected={screen.connected}>
<WorkerInstall connected={screen.connected}>
{readyAction ?? (
<Button className="w-full" onClick={() => prepareWorker.reset()}>
Done

View file

@ -5,10 +5,9 @@ import { lensKeys } from "../../data/queries";
import type { LensApi } from "../../data/service";
import { useLensApi } from "../../data/LensServices";
import type { WorkerCreated } from "../../model/types";
import { validateWorkerAddress, analysisAccessSchema, type AnalysisAccess } from "./workerSchema";
import { analysisAccessSchema, type AnalysisAccess } from "./workerSchema";
export interface WorkerRegistration {
readonly address: string;
readonly useExisting: boolean;
readonly analysisKey: string | null;
readonly access: AnalysisAccess;
@ -30,8 +29,7 @@ async function releaseUnusedKey(api: LensApi, keyId: string): Promise<void> {
}
async function prepareWorker(api: LensApi, registration: WorkerRegistration): Promise<WorkerCreated | null> {
const { address, useExisting, analysisKey, access, editingWorker } = registration;
validateWorkerAddress(address);
const { useExisting, analysisKey, access, editingWorker } = registration;
if (useExisting && !analysisKey) throw new Error("Choose an existing key");
const keyId = useExisting && analysisKey ? analysisKey : await createAnalysisKey(api, access);
const newKey = useExisting ? null : keyId;

View file

@ -1,22 +0,0 @@
import { proxyBaseUrl } from "@/components/networking";
import { serverRootPath } from "@/lib/serverRootPath";
export function initialProxyAddress(): string {
const url = new URL(proxyBaseUrl || serverRootPath, window.location.origin);
if (["localhost", "127.0.0.1", "[::1]"].includes(url.hostname)) url.hostname = "host.docker.internal";
return url.toString().replace(/\/$/, "");
}
export function workerSetupCommand(address: string, token: string, image: string): string {
const quote = (value: string) => "'" + value.replaceAll("'", "'\\''") + "'";
return [
"docker run -d --restart unless-stopped --read-only --cap-drop ALL",
" --tmpfs /tmp:rw,noexec,nosuid,size=1g",
" --security-opt no-new-privileges --add-host host.docker.internal:host-gateway",
` -e ${quote("LITELLM_URL=" + address)}`,
` -e ${quote("LENS_WORKER_TOKEN=" + token)}`,
" -p 127.0.0.1:4318:4318 --memory 2g --cpus 2",
" -e LITELLM_LENS_SERVICE_TOKEN -e CLICKHOUSE_URL",
` ${quote(image)}`,
].join(" \\\n");
}

View file

@ -1,41 +1,12 @@
import { describe, expect, it } from "vitest";
import { validateWorkerAddress, analysisAccessSchema, workerFormSchema } from "./workerSchema";
import { workerSetupCommand } from "./workerCommand";
import { expect, it } from "vitest";
import { analysisAccessSchema, workerFormSchema } from "./workerSchema";
const workerDefaults = {
useExisting: false,
analysisKey: null,
access: { model: null, budget: "100" },
address: "http://localhost:4000",
};
describe("worker setup", () => {
it.each(["https://gateway.example/proxy", "http://host.docker.internal:4000"])("accepts %s", (address) => {
expect(() => validateWorkerAddress(address)).not.toThrow();
});
it.each(["ftp://gateway.example", "https://user:pass@gateway.example", "https://user@gateway.example"])(
"rejects %s",
(address) => {
expect(() => validateWorkerAddress(address)).toThrow("Enter an HTTP or HTTPS proxy URL without credentials");
},
);
it("rejects malformed addresses", () => {
expect(() => validateWorkerAddress("not a URL")).toThrow("Invalid URL");
});
it("quotes apostrophes literally and retains the pinned image and runtime restrictions", () => {
const image = "registry.example/lens-worker:v1.2.3-rc.4";
const command = workerSetupCommand("https://gateway.example/proxy?name=it's", "token'quoted", image);
expect(command).toContain("'LITELLM_URL=https://gateway.example/proxy?name=it'\\''s'");
expect(command).toContain("'LENS_WORKER_TOKEN=token'\\''quoted'");
expect(command).toContain("--read-only --cap-drop ALL");
expect(command).toContain("--security-opt no-new-privileges --add-host host.docker.internal:host-gateway");
expect(command.split("\n").at(-1)?.trim()).toBe(`'${image}'`);
});
});
it.each(["", "0", "-1", "NaN", "Infinity"])("rejects an invalid analysis budget: %s", (budget) => {
const result = analysisAccessSchema.safeParse({ model: "analysis", budget });
expect(result.success).toBe(false);
@ -51,15 +22,14 @@ it("requires a model and converts an accepted analysis budget to a number", () =
});
});
it("validates the selected-key branch and non-empty proxy address with field paths", () => {
it("requires a billing key when using existing access", () => {
const result = workerFormSchema.safeParse({
...workerDefaults,
useExisting: true,
address: " ",
});
expect(result.success).toBe(false);
if (result.success) return;
expect(result.error.issues.map(({ path }) => path)).toEqual([["analysisKey"], ["address"]]);
expect(result.error.issues.map(({ path }) => path)).toEqual([["analysisKey"]]);
});
it("validates analysis access only when creating a new virtual key", () => {

View file

@ -1,27 +1,10 @@
import { z } from "zod";
export const workerAddressSchema = z.string().superRefine((address, ctx) => {
try {
const parsed = new URL(address);
if (!["http:", "https:"].includes(parsed.protocol) || parsed.username || parsed.password) {
ctx.addIssue({ code: "custom", message: "Enter an HTTP or HTTPS proxy URL without credentials" });
}
} catch (error) {
ctx.addIssue({ code: "custom", message: error instanceof Error ? error.message : "Invalid URL" });
}
});
export function validateWorkerAddress(address: string) {
const result = workerAddressSchema.safeParse(address);
if (!result.success) throw new Error(result.error.issues[0].message);
}
const analysisAccessFields = { model: z.string().nullable(), budget: z.string() };
const workerFormFields = {
useExisting: z.boolean(),
analysisKey: z.string().nullable(),
access: z.object(analysisAccessFields),
address: z.string(),
};
export const analysisAccessSchema = z
.object(analysisAccessFields)
@ -52,13 +35,6 @@ export const workerFormSchema = z.object(workerFormFields).superRefine((values,
});
}
}
if (!values.address.trim()) {
ctx.addIssue({
code: "custom",
message: "Enter a proxy URL",
path: ["address"],
});
}
});
export type WorkerFormInput = z.input<typeof workerFormSchema>;

View file

@ -27,6 +27,7 @@ export interface RequestOptions {
signal?: AbortSignal;
/** Send browser cookies with the request; needed for cookie-authenticated proxy routes. */
credentials?: RequestCredentials;
redirect?: RequestRedirect;
}
export class ApiError extends Error {
@ -152,7 +153,7 @@ export function createApiClient(config: ApiClientConfig): ApiClient {
const doFetch: typeof fetch = (input, init) => (fetchImpl ?? fetch)(input, init);
async function fetchChecked(method: HttpMethod, path: string, options: RequestOptions = {}): Promise<Response> {
const { accessToken, body, rawBody, query, headers: extraHeaders, signal, credentials } = options;
const { accessToken, body, rawBody, query, headers: extraHeaders, signal, credentials, redirect } = options;
const url = appendQuery(`${getBaseUrl()}${path}`, query);
@ -168,7 +169,7 @@ export function createApiClient(config: ApiClientConfig): ApiClient {
Object.assign(headers, extraHeaders);
}
const init: RequestInit = { method, headers, signal, credentials };
const init: RequestInit = { method, headers, signal, credentials, redirect };
if (rawBody !== undefined) {
init.body = rawBody;
} else if (body !== undefined) {

View file

@ -13079,7 +13079,7 @@ export interface paths {
* Example:
* ```bash
* curl --location --request DELETE 'http://0.0.0.0:4000/project/delete' \
* --header 'Authorization: Bearer sk-1234' \
* --header "Authorization: Bearer $LITELLM_MASTER_KEY" \
* --header 'Content-Type: application/json' \
* --data '{
* "project_ids": ["project-123", "project-456"]
@ -13109,7 +13109,7 @@ export interface paths {
* Example:
* ```bash
* curl --location 'http://0.0.0.0:4000/project/info?project_id=project-123' \
* --header 'Authorization: Bearer sk-1234'
* --header "Authorization: Bearer $LITELLM_MASTER_KEY"
* ```
*/
get: operations["project_info_project_info_get"];
@ -13135,7 +13135,7 @@ export interface paths {
* Example:
* ```bash
* curl --location 'http://0.0.0.0:4000/project/list' \
* --header 'Authorization: Bearer sk-1234'
* --header "Authorization: Bearer $LITELLM_MASTER_KEY"
* ```
*/
get: operations["list_projects_project_list_get"];
@ -13188,7 +13188,7 @@ export interface paths {
*
* ```bash
* curl --location 'http://0.0.0.0:4000/project/new' \
* --header 'Authorization: Bearer sk-1234' \
* --header "Authorization: Bearer $LITELLM_MASTER_KEY" \
* --header 'Content-Type: application/json' \
* --data '{
* "project_alias": "flight-search-assistant",
@ -13215,7 +13215,7 @@ export interface paths {
*
* ```bash
* curl --location 'http://0.0.0.0:4000/project/new' \
* --header 'Authorization: Bearer sk-1234' \
* --header "Authorization: Bearer $LITELLM_MASTER_KEY" \
* --header 'Content-Type: application/json' \
* --data '{
* "project_alias": "hotel-recommendations",
@ -13269,7 +13269,7 @@ export interface paths {
* Example:
* ```bash
* curl --location 'http://0.0.0.0:4000/project/update' \
* --header 'Authorization: Bearer sk-1234' \
* --header "Authorization: Bearer $LITELLM_MASTER_KEY" \
* --header 'Content-Type: application/json' \
* --data '{
* "project_id": "project-123",