diff --git a/deploy/lens/compose.yaml b/deploy/lens/compose.yaml index 1d507fc7b74..a95684a45c3 100644 --- a/deploy/lens/compose.yaml +++ b/deploy/lens/compose.yaml @@ -3,8 +3,15 @@ services: image: ${LENS_WORKER_IMAGE:-${LITELLM_VERSION:+ghcr.io/berriai/litellm-lens-worker:v}${LITELLM_VERSION:-}} environment: LITELLM_URL: ${LITELLM_URL:?Set the URL reachable from this container} - LENS_WORKER_TOKEN: ${LENS_WORKER_TOKEN:?Create a worker credential in the Lens UI} - LENS_PYTHON_CONCURRENCY: ${LENS_PYTHON_CONCURRENCY:-2} + LENS_WORKER_TOKEN: ${LENS_WORKER_TOKEN:-} + LITELLM_LENS_SERVICE_TOKEN: ${LITELLM_LENS_SERVICE_TOKEN:?Set the same secret on LiteLLM and Lens} + CLICKHOUSE_URL: ${CLICKHOUSE_URL:?Set the ClickHouse URL reachable from Lens} + CLICKHOUSE_DATABASE: ${CLICKHOUSE_DATABASE:-litellm} + AGENT_TRACING_RETENTION_DAYS: ${AGENT_TRACING_RETENTION_DAYS:-14} + ports: + - "127.0.0.1:${LENS_PORT:-4318}:4318" + mem_limit: 2g + cpus: 2 restart: unless-stopped read_only: true tmpfs: diff --git a/deploy/lens/config.yaml b/deploy/lens/config.yaml index cb12a2b0919..43bfe32ec26 100644 --- a/deploy/lens/config.yaml +++ b/deploy/lens/config.yaml @@ -2,6 +2,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/deploy/lens/stack.yaml b/deploy/lens/stack.yaml index 852aa9d8488..85d1b38a3d5 100644 --- a/deploy/lens/stack.yaml +++ b/deploy/lens/stack.yaml @@ -10,9 +10,7 @@ services: import os, sys from urllib.parse import quote postgres_password = quote(os.environ["POSTGRES_PASSWORD"], safe="") - clickhouse_password = quote(os.environ["CLICKHOUSE_PASSWORD"], safe="") os.environ["DATABASE_URL"] = f"postgresql://litellm:{postgres_password}@db:5432/litellm" - os.environ["CLICKHOUSE_URL"] = f"http://default:{clickhouse_password}@clickhouse:8123" os.execv("docker/prod_entrypoint.sh", ["docker/prod_entrypoint.sh", *sys.argv[1:]]) command: ["--config", "/app/lens-config.yaml", "--port", "4000"] environment: @@ -20,7 +18,9 @@ services: LITELLM_SALT_KEY: ${LITELLM_SALT_KEY:?Set a permanent encryption key and keep it across upgrades} POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:?Set a permanent database password} STORE_MODEL_IN_DB: "True" - CLICKHOUSE_PASSWORD: ${CLICKHOUSE_PASSWORD:?Set a permanent ClickHouse password} + LITELLM_LENS_URL: http://lens-worker:4318 + LITELLM_LENS_PUBLIC_URL: ${LITELLM_LENS_PUBLIC_URL:-http://localhost:4318} + LITELLM_LENS_SERVICE_TOKEN: ${LITELLM_LENS_SERVICE_TOKEN:?Set the shared Lens service secret} LENS_WORKER_IMAGE: ghcr.io/berriai/litellm-lens-worker:v${LITELLM_VERSION} volumes: - ./config.yaml:/app/lens-config.yaml:ro @@ -30,19 +30,24 @@ services: depends_on: db: condition: service_healthy - clickhouse: - condition: service_healthy restart: unless-stopped lens-worker: - profiles: [lens] image: ghcr.io/berriai/litellm-lens-worker:v${LITELLM_VERSION} environment: LITELLM_URL: http://litellm:4000 LENS_WORKER_TOKEN: ${LENS_WORKER_TOKEN:-} - LENS_PYTHON_CONCURRENCY: ${LENS_PYTHON_CONCURRENCY:-2} + LITELLM_LENS_SERVICE_TOKEN: ${LITELLM_LENS_SERVICE_TOKEN} + CLICKHOUSE_HOST: clickhouse + CLICKHOUSE_PASSWORD: ${CLICKHOUSE_PASSWORD:?Set a permanent ClickHouse password} + CLICKHOUSE_DATABASE: ${CLICKHOUSE_DATABASE:-litellm} + AGENT_TRACING_RETENTION_DAYS: ${AGENT_TRACING_RETENTION_DAYS:-14} depends_on: [litellm] - networks: [proxy] + networks: [proxy, storage] + ports: + - "127.0.0.1:${LENS_PORT:-4318}:4318" + mem_limit: 2g + cpus: 2 restart: unless-stopped read_only: true tmpfs: diff --git a/helm/litellm-helm/templates/_helpers.tpl b/helm/litellm-helm/templates/_helpers.tpl index 9630633912e..7ba8331bdcd 100644 --- a/helm/litellm-helm/templates/_helpers.tpl +++ b/helm/litellm-helm/templates/_helpers.tpl @@ -321,3 +321,47 @@ through an emptyDir. Empty when the sidecar is off or uses 127.0.0.1 TCP. - name: LITELLM_COLLECTOR_DRAIN_TIMEOUT_SECONDS value: {{ .Values.collector.drainTimeoutSeconds | quote }} {{- end -}} + +{{- define "litellm.lensWorker.image" -}} +{{- if .Values.lensWorker.image.digest -}} +{{- if not (regexMatch "^sha256:[0-9a-f]{64}$" .Values.lensWorker.image.digest) -}} +{{- fail "lensWorker.image.digest must be sha256 followed by 64 lowercase hex characters" -}} +{{- end -}} +{{- printf "%s@%s" .Values.lensWorker.image.repository .Values.lensWorker.image.digest -}} +{{- else -}} +{{- $backendTag := .Values.image.tag | default .Chart.AppVersion -}} +{{- $releaseTag := ternary (printf "v%s" $backendTag) $backendTag (regexMatch "^[0-9]" $backendTag) -}} +{{- $tag := .Values.lensWorker.image.tag | default $releaseTag -}} +{{- $repository := .Values.lensWorker.image.repository -}} +{{- if and (hasPrefix "sha-" $tag) (eq $repository "ghcr.io/berriai/litellm-lens-worker") -}} +{{- $repository = "ghcr.io/berriai/litellm-lens-worker-dev" -}} +{{- end -}} +{{- printf "%s:%s" $repository $tag -}} +{{- end -}} +{{- end -}} + +{{- define "litellm.gateway.collectorSocketDir" -}} +{{- if and .Values.gateway.collector.enabled (hasPrefix "unix://" .Values.gateway.collector.address) -}} +{{- dir (trimPrefix "unix://" .Values.gateway.collector.address) -}} +{{- end -}} +{{- end -}} + +{{/* +LITELLM_COLLECTOR_* env shared by the producer (gateway container) and the +consumer (collector container), so both agree on the transport and the +shutdown drain window. +*/}} +{{- define "litellm.gateway.collectorEnv" -}} +{{- with .Values.gateway.collector }} +- name: LITELLM_COLLECTOR_ENABLED + value: "true" +- name: LITELLM_COLLECTOR_ADDRESS + value: {{ .address | quote }} +- name: LITELLM_COLLECTOR_BUFFER_SIZE + value: {{ .bufferSize | quote }} +- name: LITELLM_COLLECTOR_ON_UNAVAILABLE + value: {{ .onUnavailable | quote }} +- name: LITELLM_COLLECTOR_DRAIN_TIMEOUT_SECONDS + value: {{ .drainTimeoutSeconds | quote }} +{{- end }} +{{- end -}} diff --git a/helm/litellm-helm/templates/deployment.yaml b/helm/litellm-helm/templates/deployment.yaml index 299d41e2019..4aac75fba8c 100644 --- a/helm/litellm-helm/templates/deployment.yaml +++ b/helm/litellm-helm/templates/deployment.yaml @@ -56,6 +56,17 @@ spec: image: "{{ .Values.image.repository }}:{{ .Values.image.tag | default .Chart.AppVersion }}" imagePullPolicy: {{ .Values.image.pullPolicy }} 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.proxyEnv" . | nindent 12 }} {{- if .Values.liteadmin.enabled }} - name: LITELLM_ADMIN_AGENT_URL diff --git a/helm/litellm-helm/templates/lens/deployment.yaml b/helm/litellm-helm/templates/lens/deployment.yaml new file mode 100644 index 00000000000..aa7686ad16b --- /dev/null +++ b/helm/litellm-helm/templates/lens/deployment.yaml @@ -0,0 +1,95 @@ +{{- if .Values.lensWorker.enabled }} +apiVersion: apps/v1 +kind: Deployment +metadata: + name: {{ include "litellm.fullname" . }}-lens-worker + labels: + {{- include "litellm.labels" . | nindent 4 }} + app.kubernetes.io/component: lens-worker +spec: + replicas: {{ .Values.lensWorker.replicaCount }} + selector: + matchLabels: + app.kubernetes.io/instance: {{ .Release.Name }} + app.kubernetes.io/component: lens-worker + template: + metadata: + labels: + {{- include "litellm.labels" . | nindent 8 }} + app.kubernetes.io/component: lens-worker + spec: + automountServiceAccountToken: false + {{- with .Values.imagePullSecrets }} + imagePullSecrets: + {{- toYaml . | nindent 8 }} + {{- end }} + securityContext: + runAsNonRoot: true + runAsUser: 65532 + runAsGroup: 65532 + fsGroup: 65532 + seccompProfile: + type: RuntimeDefault + containers: + - name: lens-worker + image: {{ include "litellm.lensWorker.image" . | quote }} + imagePullPolicy: {{ .Values.lensWorker.image.pullPolicy }} + securityContext: + allowPrivilegeEscalation: false + readOnlyRootFilesystem: true + capabilities: + drop: [ALL] + env: + - name: LITELLM_URL + value: {{ .Values.lensWorker.url | default (printf "http://%s:%v" (include "litellm.fullname" .) .Values.service.port) | 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 }} + - name: CLICKHOUSE_URL + valueFrom: + secretKeyRef: + name: {{ required "lensWorker.clickhouseSecret.name is required" .Values.lensWorker.clickhouseSecret.name | quote }} + key: {{ .Values.lensWorker.clickhouseSecret.key | quote }} + {{- if .Values.lensWorker.tokenSecret.name }} + - name: LENS_WORKER_TOKEN + valueFrom: + secretKeyRef: + name: {{ .Values.lensWorker.tokenSecret.name | quote }} + key: {{ .Values.lensWorker.tokenSecret.key | quote }} + {{- end }} + ports: + - name: otlp + containerPort: 4318 + livenessProbe: + httpGet: + path: /health/live + port: otlp + readinessProbe: + httpGet: + path: /health/ready + port: otlp + resources: + {{- toYaml .Values.lensWorker.resources | nindent 12 }} + volumeMounts: + - name: tmp + mountPath: /tmp + volumes: + - name: tmp + emptyDir: + medium: Memory + sizeLimit: {{ .Values.lensWorker.tmpSizeLimit }} + {{- with .Values.lensWorker.nodeSelector }} + nodeSelector: + {{- toYaml . | nindent 8 }} + {{- end }} + {{- with .Values.lensWorker.tolerations }} + tolerations: + {{- toYaml . | nindent 8 }} + {{- end }} + {{- with .Values.lensWorker.affinity }} + affinity: + {{- toYaml . | nindent 8 }} + {{- end }} +{{- end }} diff --git a/helm/litellm-helm/templates/lens/ingress.yaml b/helm/litellm-helm/templates/lens/ingress.yaml new file mode 100644 index 00000000000..d2b73390cd2 --- /dev/null +++ b/helm/litellm-helm/templates/lens/ingress.yaml @@ -0,0 +1,29 @@ +{{- if and .Values.lensWorker.enabled .Values.lensWorker.ingress.enabled }} +apiVersion: networking.k8s.io/v1 +kind: Ingress +metadata: + name: {{ include "litellm.fullname" . }}-lens-worker + {{- with .Values.lensWorker.ingress.annotations }} + annotations: + {{- toYaml . | nindent 4 }} + {{- end }} +spec: + {{- with .Values.lensWorker.ingress.className }} + ingressClassName: {{ . | quote }} + {{- end }} + {{- with .Values.lensWorker.ingress.tls }} + tls: + {{- toYaml . | nindent 4 }} + {{- end }} + rules: + - host: {{ required "lensWorker.ingress.host is required" .Values.lensWorker.ingress.host | quote }} + http: + paths: + - path: /v1/ + pathType: Prefix + backend: + service: + name: {{ include "litellm.fullname" . }}-lens-worker + port: + name: otlp +{{- end }} diff --git a/helm/litellm-helm/templates/lens/service.yaml b/helm/litellm-helm/templates/lens/service.yaml new file mode 100644 index 00000000000..a09dde1a4fd --- /dev/null +++ b/helm/litellm-helm/templates/lens/service.yaml @@ -0,0 +1,14 @@ +{{- if .Values.lensWorker.enabled }} +apiVersion: v1 +kind: Service +metadata: + name: {{ include "litellm.fullname" . }}-lens-worker +spec: + selector: + app.kubernetes.io/instance: {{ .Release.Name }} + app.kubernetes.io/component: lens-worker + ports: + - name: otlp + port: {{ .Values.lensWorker.service.port }} + targetPort: otlp +{{- end }} diff --git a/helm/litellm-helm/values.yaml b/helm/litellm-helm/values.yaml index 83dbb3c5aa0..8da4cd52325 100644 --- a/helm/litellm-helm/values.yaml +++ b/helm/litellm-helm/values.yaml @@ -647,3 +647,41 @@ serviceMonitor: namespaceSelector: matchNames: [] # - test-namespace + +lensWorker: + enabled: false + replicaCount: 1 + image: + repository: ghcr.io/berriai/litellm-lens-worker + tag: "" + digest: "" + pullPolicy: IfNotPresent + tokenSecret: + name: "" + key: token + serviceTokenSecret: + name: "" + key: service-token + clickhouseSecret: + name: "" + key: url + publicUrl: "" + service: + port: 4318 + ingress: + enabled: false + className: "" + host: "" + annotations: {} + tls: [] + url: "" + tmpSizeLimit: 1Gi + resources: + requests: + cpu: 100m + memory: 256Mi + limits: + memory: 2Gi + nodeSelector: {} + tolerations: [] + affinity: {} diff --git a/helm/litellm/templates/backend/deployment.yaml b/helm/litellm/templates/backend/deployment.yaml index 5d3be1439bd..ebb9a1dd838 100644 --- a/helm/litellm/templates/backend/deployment.yaml +++ b/helm/litellm/templates/backend/deployment.yaml @@ -57,6 +57,17 @@ 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 }} - 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/lens/deployment.yaml b/helm/litellm/templates/lens/deployment.yaml index 787581b9ad1..f3ba3780231 100644 --- a/helm/litellm/templates/lens/deployment.yaml +++ b/helm/litellm/templates/lens/deployment.yaml @@ -42,11 +42,34 @@ spec: env: - name: LITELLM_URL value: {{ .Values.lensWorker.url | default (printf "http://%s:%v" (include "litellm.backend.fullname" .) .Values.backend.service.port) | 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 }} + - name: CLICKHOUSE_URL + valueFrom: + secretKeyRef: + name: {{ required "lensWorker.clickhouseSecret.name is required" .Values.lensWorker.clickhouseSecret.name | quote }} + key: {{ .Values.lensWorker.clickhouseSecret.key | quote }} + {{- if .Values.lensWorker.tokenSecret.name }} - name: LENS_WORKER_TOKEN valueFrom: secretKeyRef: - name: {{ required "lensWorker.tokenSecret.name must reference a Lens worker token" .Values.lensWorker.tokenSecret.name | quote }} + name: {{ .Values.lensWorker.tokenSecret.name | quote }} key: {{ .Values.lensWorker.tokenSecret.key | quote }} + {{- end }} + ports: + - name: otlp + containerPort: 4318 + livenessProbe: + httpGet: + path: /health/live + port: otlp + readinessProbe: + httpGet: + path: /health/ready + port: otlp resources: {{- toYaml .Values.lensWorker.resources | nindent 12 }} volumeMounts: diff --git a/helm/litellm/templates/lens/ingress.yaml b/helm/litellm/templates/lens/ingress.yaml new file mode 100644 index 00000000000..d2b73390cd2 --- /dev/null +++ b/helm/litellm/templates/lens/ingress.yaml @@ -0,0 +1,29 @@ +{{- if and .Values.lensWorker.enabled .Values.lensWorker.ingress.enabled }} +apiVersion: networking.k8s.io/v1 +kind: Ingress +metadata: + name: {{ include "litellm.fullname" . }}-lens-worker + {{- with .Values.lensWorker.ingress.annotations }} + annotations: + {{- toYaml . | nindent 4 }} + {{- end }} +spec: + {{- with .Values.lensWorker.ingress.className }} + ingressClassName: {{ . | quote }} + {{- end }} + {{- with .Values.lensWorker.ingress.tls }} + tls: + {{- toYaml . | nindent 4 }} + {{- end }} + rules: + - host: {{ required "lensWorker.ingress.host is required" .Values.lensWorker.ingress.host | quote }} + http: + paths: + - path: /v1/ + pathType: Prefix + backend: + service: + name: {{ include "litellm.fullname" . }}-lens-worker + port: + name: otlp +{{- end }} diff --git a/helm/litellm/templates/lens/service.yaml b/helm/litellm/templates/lens/service.yaml new file mode 100644 index 00000000000..a09dde1a4fd --- /dev/null +++ b/helm/litellm/templates/lens/service.yaml @@ -0,0 +1,14 @@ +{{- if .Values.lensWorker.enabled }} +apiVersion: v1 +kind: Service +metadata: + name: {{ include "litellm.fullname" . }}-lens-worker +spec: + selector: + app.kubernetes.io/instance: {{ .Release.Name }} + app.kubernetes.io/component: lens-worker + ports: + - name: otlp + port: {{ .Values.lensWorker.service.port }} + targetPort: otlp +{{- end }} diff --git a/helm/litellm/values.yaml b/helm/litellm/values.yaml index cf3334f6156..7b1d591babe 100644 --- a/helm/litellm/values.yaml +++ b/helm/litellm/values.yaml @@ -641,6 +641,21 @@ lensWorker: tokenSecret: name: "" key: token + serviceTokenSecret: + name: "" + key: service-token + clickhouseSecret: + name: "" + key: url + publicUrl: "" + service: + port: 4318 + ingress: + enabled: false + className: "" + host: "" + annotations: {} + tls: [] url: "" tmpSizeLimit: 1Gi resources: diff --git a/litellm-rust/Cargo.lock b/litellm-rust/Cargo.lock index f394dc29eac..c7a8f8fba6f 100644 --- a/litellm-rust/Cargo.lock +++ b/litellm-rust/Cargo.lock @@ -4201,6 +4201,7 @@ dependencies = [ "tempfile", "thiserror 2.0.19", "tokio", + "tower-http", "tracing", "typify", "unicode-casefold", diff --git a/litellm-rust/crates/lens/Cargo.toml b/litellm-rust/crates/lens/Cargo.toml index 6c68a5b06da..4b1ac05aca5 100644 --- a/litellm-rust/crates/lens/Cargo.toml +++ b/litellm-rust/crates/lens/Cargo.toml @@ -30,6 +30,7 @@ tempfile.workspace = true thiserror.workspace = true tokio = { workspace = true, features = ["signal", "sync", "process", "io-util"] } tracing.workspace = true +tower-http = { version = "0.6.11", features = ["cors"] } url.workspace = true unicode-casefold = "0.2" diff --git a/litellm-rust/crates/lens/src/config.rs b/litellm-rust/crates/lens/src/config.rs index 5f8b3e81bd6..347800efeea 100644 --- a/litellm-rust/crates/lens/src/config.rs +++ b/litellm-rust/crates/lens/src/config.rs @@ -33,8 +33,11 @@ impl Config { { return Err(Error::Configuration("LITELLM_URL")); } - let worker_token = required("LENS_WORKER_TOKEN")?; let service_token = required("LITELLM_LENS_SERVICE_TOKEN")?; + let worker_token = std::env::var("LENS_WORKER_TOKEN") + .ok() + .filter(|value| !value.is_empty()) + .unwrap_or_else(|| service_token.clone()); if service_token.len() < 32 { return Err(Error::Configuration( "LITELLM_LENS_SERVICE_TOKEN must contain at least 32 characters", @@ -51,7 +54,7 @@ impl Config { release: required("LITELLM_RELEASE_TAG")?, storage: StorageConfig::new( std::env::var("CLICKHOUSE_DATABASE").unwrap_or_else(|_| "litellm".into()), - &required("CLICKHOUSE_URL")?, + &clickhouse_url()?, std::env::var("AGENT_TRACING_RETENTION_DAYS") .unwrap_or_else(|_| "14".into()) .parse() @@ -62,6 +65,21 @@ impl Config { } } +fn clickhouse_url() -> Result { + if let Ok(url) = required("CLICKHOUSE_URL") { + return Ok(url); + } + let mut url = url::Url::parse("http://localhost:8123") + .map_err(|_| Error::Configuration("CLICKHOUSE_HOST"))?; + url.set_host(Some(&required("CLICKHOUSE_HOST")?)) + .map_err(|_| Error::Configuration("CLICKHOUSE_HOST"))?; + url.set_username(&std::env::var("CLICKHOUSE_USER").unwrap_or_else(|_| "default".into())) + .map_err(|_| Error::Configuration("CLICKHOUSE_USER"))?; + url.set_password(Some(&required("CLICKHOUSE_PASSWORD")?)) + .map_err(|_| Error::Configuration("CLICKHOUSE_PASSWORD"))?; + Ok(url.into()) +} + pub fn http_client() -> Result { let settings = HttpSettings { connect_timeout: Duration::from_secs(5), diff --git a/litellm-rust/crates/lens/src/control.rs b/litellm-rust/crates/lens/src/control.rs index 88b639a6556..be90edb1690 100644 --- a/litellm-rust/crates/lens/src/control.rs +++ b/litellm-rust/crates/lens/src/control.rs @@ -14,6 +14,7 @@ pub struct Control { base: Url, token: Arc, model_slots: Arc, + attempt: Option, } impl Control { @@ -26,6 +27,7 @@ impl Control { base, token: token.into(), model_slots: Arc::new(Semaphore::new(16)), + attempt: None, } } @@ -52,6 +54,10 @@ impl Control { Some(body) => request.json(body), None => request, }; + let request = match self.attempt { + Some(attempt) => request.header("x-litellm-lens-attempt", attempt), + None => request, + }; let mut response = request.send().await?; let status = response.status(); if !status.is_success() { @@ -141,6 +147,11 @@ pub struct JobClient { } impl JobClient { + pub fn with_attempt(mut self, attempt: u64) -> Self { + self.control.attempt = Some(attempt); + self + } + pub fn new( control: Control, lens_id: &str, diff --git a/litellm-rust/crates/lens/src/lib.rs b/litellm-rust/crates/lens/src/lib.rs index 64680205b32..ed7e35eecb4 100644 --- a/litellm-rust/crates/lens/src/lib.rs +++ b/litellm-rust/crates/lens/src/lib.rs @@ -81,16 +81,98 @@ impl State { } pub fn router(state: Arc) -> Router { - Router::new() + let public = Router::new() .route("/health/live", get(|| async { StatusCode::OK })) .route("/health/ready", get(ready)) .route("/v1/traces", post(traces)) .route("/v1/logs", post(logs)) - .route("/internal/read", post(read)) - .route("/internal/spend", post(spend)) + .route("/v1/traces/receipt", post(receipt)) + .layer( + tower_http::cors::CorsLayer::new() + .allow_origin(tower_http::cors::Any) + .allow_methods([http::Method::POST, http::Method::GET]) + .allow_headers([ + http::header::AUTHORIZATION, + http::header::CONTENT_TYPE, + http::header::CONTENT_ENCODING, + ]), + ); + public + .merge( + Router::new() + .route("/internal/read", post(read)) + .route("/internal/spend", post(spend)) + .route("/internal/credentials", post(credentials)) + .route("/internal/status", get(status)), + ) .with_state(state) } +#[derive(serde::Deserialize)] +#[serde(deny_unknown_fields)] +struct ReceiptRequest { + trace_id: String, + #[serde(default)] + span_ids: Vec, +} + +async fn receipt( + AppState(state): AppState>, + headers: HeaderMap, + body: Body, +) -> Result, Error> { + let tenant = state.credentials.tenant(&headers)?; + state.require_storage()?; + let _permit = state + .read_slots + .try_acquire() + .map_err(|_| Error::Unavailable)?; + let body = tokio::time::timeout(Duration::from_secs(5), to_bytes(body, 64 * 1024)) + .await + .map_err(|_| Error::Unavailable)? + .map_err(|_| Error::TooLarge)?; + let request: ReceiptRequest = + serde_json::from_slice(&body).map_err(|_| Error::InvalidRequest)?; + let received = litellm_traces_clickhouse::trace_received( + &state.storage.client, + state.storage.config.storage().reader(), + &tenant, + &request.trace_id, + &request.span_ids, + ) + .await?; + Ok(Json(serde_json::json!({"received": received}))) +} + +async fn status( + AppState(state): AppState>, + headers: HeaderMap, +) -> Result, Error> { + auth::authorize_service(&headers, &state.service_token)?; + Ok(Json(serde_json::json!({ + "storage_ready": state.schema_ready.load(Ordering::Acquire), + "credentials_ready": state.credentials.ready(), + "release": std::env::var("LITELLM_RELEASE_TAG").unwrap_or_default(), + "protocol_version": wire::PROTOCOL_VERSION, + }))) +} + +async fn credentials( + AppState(state): AppState>, + headers: HeaderMap, + body: Body, +) -> Result { + auth::authorize_service(&headers, &state.service_token)?; + let body = tokio::time::timeout(Duration::from_secs(5), to_bytes(body, 8 * 1024 * 1024)) + .await + .map_err(|_| Error::Unavailable)? + .map_err(|_| Error::TooLarge)?; + state + .credentials + .replace(serde_json::from_slice(&body).map_err(|_| Error::InvalidRequest)?)?; + Ok(StatusCode::NO_CONTENT) +} + async fn ready(AppState(state): AppState>) -> StatusCode { if state.schema_ready.load(Ordering::Acquire) && state.credentials.ready() { StatusCode::OK diff --git a/litellm-rust/crates/lens/src/main.rs b/litellm-rust/crates/lens/src/main.rs index cb8c6dda019..135ed69cf38 100644 --- a/litellm-rust/crates/lens/src/main.rs +++ b/litellm-rust/crates/lens/src/main.rs @@ -51,13 +51,13 @@ async fn run() -> Result<(), litellm_lens::Error> { config.worker_token.clone(), ); let storage = Storage::new(config.storage, client.clone(), config.service_token.clone()); - let state = Arc::new(State::new(storage, config.service_token)); + let state = Arc::new(State::new(storage, config.service_token.clone())); let listener = tokio::net::TcpListener::bind(config.address).await?; let auth_task = tokio::spawn(auth::refresh_loop( state.credentials.clone(), client, - control.url("lens/worker/ingestion-credentials")?, - config.worker_token, + control.url("lens/internal/ingestion-credentials")?, + config.service_token, )); let provision_task = tokio::spawn(provision(state.clone())); let mut worker = tokio::spawn(Worker::new(control, config.release).serve()); diff --git a/litellm-rust/crates/lens/src/worker.rs b/litellm-rust/crates/lens/src/worker.rs index 9f0ef4f5170..851f1e8a679 100644 --- a/litellm-rust/crates/lens/src/worker.rs +++ b/litellm-rust/crates/lens/src/worker.rs @@ -23,6 +23,7 @@ struct Identity { #[derive(Deserialize)] struct JobIdentity { id: String, + attempts: u64, } impl Worker { @@ -48,7 +49,8 @@ impl Worker { if claim.is_err() || !validator.is_valid(&payload) { let identity: Identity = serde_json::from_value(payload)?; let client = - JobClient::new(self.control.clone(), &identity.lens_id, &identity.job.id, 1)?; + JobClient::new(self.control.clone(), &identity.lens_id, &identity.job.id, 1)? + .with_attempt(identity.job.attempts); self.failure(&client, "The worker could not read this investigation. Update the worker to match the gateway, then retry.").await?; return Ok(true); } @@ -58,7 +60,8 @@ impl Worker { &claim.lens_id, &claim.job.id, claim.job.settings.concurrency.get() as usize, - )?; + )? + .with_attempt(u64::try_from(claim.job.attempts).map_err(|_| Error::InvalidRequest)?); let work = async { let sample: wire::Sample = client.get("sample").await?; claim.reviews = Some(client.get("reviews").await?); diff --git a/litellm-rust/crates/lens/tests/sandbox.rs b/litellm-rust/crates/lens/tests/sandbox.rs index ed0267fba94..f612ceec5a5 100644 --- a/litellm-rust/crates/lens/tests/sandbox.rs +++ b/litellm-rust/crates/lens/tests/sandbox.rs @@ -29,7 +29,7 @@ fn workspace() -> Workspace { } fn request(code: &str) -> wire::PythonRequest { - serde_json::from_value(json!({"code": code})).unwrap() + serde_json::from_value(json!({"action": "python", "code": code})).unwrap() } fn succeeded(reply: &Value) { diff --git a/litellm-rust/crates/traces-clickhouse/src/lib.rs b/litellm-rust/crates/traces-clickhouse/src/lib.rs index aed87691958..9328ee1419e 100644 --- a/litellm-rust/crates/traces-clickhouse/src/lib.rs +++ b/litellm-rust/crates/traces-clickhouse/src/lib.rs @@ -16,6 +16,7 @@ mod insert; pub mod query; mod query_access; mod reads; +mod receipt; mod schema; mod span_batches; mod span_row; @@ -32,6 +33,7 @@ pub use litellm_traces::{QueryScope, ReadQuery}; pub use query::{QueryHelp, execute_read, query_help, query_sql}; pub use query_access::QueryReaders; pub use reads::ClickHouseTraces; +pub use receipt::trace_received; pub use schema::{ NORMALIZED_FIELD_DEFINITIONS, NormalizedFieldDefinition, apply_migrations, ensure_schema, reconcile_retention, schema_statements, diff --git a/litellm-rust/crates/traces-clickhouse/src/receipt.rs b/litellm-rust/crates/traces-clickhouse/src/receipt.rs new file mode 100644 index 00000000000..5ae262d2369 --- /dev/null +++ b/litellm-rust/crates/traces-clickhouse/src/receipt.rs @@ -0,0 +1,57 @@ +use crate::{Connection, Error, Parameter}; +use litellm_http::Client; +use litellm_traces::Tenant; +use serde::Deserialize; +use std::collections::{BTreeMap, BTreeSet}; + +#[derive(Deserialize)] +struct Receipt { + received: u32, +} + +#[derive(Deserialize)] +struct Rows { + data: Vec, +} + +pub async fn trace_received( + client: &Client, + connection: &Connection, + tenant: &Tenant, + trace_id: &str, + span_ids: &[String], +) -> Result { + let valid_id = + |value: &str, length| value.len() == length && value.bytes().all(|b| b.is_ascii_hexdigit()); + if !valid_id(trace_id, 32) + || span_ids.len() > 1000 + || span_ids.iter().any(|id| !valid_id(id, 16)) + { + return Err(Error::InvalidParameters); + } + let spans: BTreeSet<_> = span_ids.iter().map(|id| id.to_ascii_lowercase()).collect(); + let expected = spans.len(); + let parameters = BTreeMap::from([ + ( + "trace_id".into(), + Parameter::Text(trace_id.to_ascii_lowercase()), + ), + ( + "api_key_hash".into(), + Parameter::Text(tenant.api_key_hash.clone()), + ), + ( + "span_ids".into(), + Parameter::Strings(spans.into_iter().collect()), + ), + ]); + let response = litellm_storage_clickhouse::execute_read(client, connection, + "SELECT toUInt32(uniqExact(SpanId)) AS received FROM otel_traces WHERE TraceId={trace_id:String} AND ApiKeyHash={api_key_hash:String} AND (empty({span_ids:Array(String)}) OR has({span_ids:Array(String)}, SpanId))", ¶meters).await?; + let rows: Rows = serde_json::from_str(&response).map_err(|_| Error::InvalidResponse)?; + let row = rows.data.first().ok_or(Error::InvalidResponse)?; + Ok(if expected == 0 { + row.received > 0 + } else { + row.received as usize == expected + }) +} diff --git a/litellm/proxy/lens/endpoints.py b/litellm/proxy/lens/endpoints.py index 88e561a0c0e..d3451d6ace0 100644 --- a/litellm/proxy/lens/endpoints.py +++ b/litellm/proxy/lens/endpoints.py @@ -7,7 +7,7 @@ from types import MappingProxyType from typing import Annotated, Final, TypeAlias from uuid import uuid4 -from fastapi import APIRouter, Depends, HTTPException, Query, Request, Response +from fastapi import APIRouter, Depends, Header, HTTPException, Query, Request, Response from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer from pydantic import AwareDatetime, Field @@ -18,15 +18,17 @@ from litellm.proxy.auth.resolvers.exceptions import KeyNotFoundError from litellm.proxy.auth.user_api_key_auth import user_api_key_auth from litellm.proxy.db.routing_prisma_wrapper import writer_wrapper from litellm.proxy.lens.billing import validate_key +from litellm.proxy.lens.inference import Deployment, deployment_prices from litellm.proxy.lens.ingestion import ( IngestionCredential, IngestionKey, IngestionKeyCreated, IngestionKeyRequest, IngestionSnapshot, + ServiceConnection, + ServiceStatus, new_key, ) -from litellm.proxy.lens.inference import Deployment, deployment_prices from litellm.proxy.lens.models import ( ActivitySelection, Claim, @@ -76,6 +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.types.llms.base import LiteLLMBaseModel router: Final = APIRouter(prefix="/lens", tags=["Lens"]) @@ -124,37 +127,42 @@ async def worker_auth(credentials: Annotated[HTTPAuthorizationCredentials, Depen WorkerAuth: TypeAlias = Annotated[Worker, Depends(worker_auth)] +Attempt: TypeAlias = Annotated[int, Header(alias="X-LiteLLM-Lens-Attempt", ge=1)] -@router.post("/tracing/keys", response_model=IngestionKeyCreated) -async def create_ingestion_key(body: IngestionKeyRequest, auth: Auth) -> IngestionKeyCreated: - user_scope(auth, write=True) +async def service_auth(credentials: Annotated[HTTPAuthorizationCredentials, Depends(_bearer)]) -> None: try: - created: Final = new_key(body, auth.user_id or "") + connection: Final = LensConnection.from_env() except ValueError as error: - raise HTTPException(422, str(error)) from error - await repository().save_ingestion_key(created.record) - return created + raise HTTPException(503, "Configure the Lens service connection") from error + if not secrets.compare_digest(credentials.credentials, connection.token): + raise HTTPException(401, "Invalid Lens service credential") -@router.get("/tracing/keys", response_model=tuple[IngestionKey, ...]) -async def list_ingestion_keys(auth: Auth) -> tuple[IngestionKey, ...]: +ServiceAuth: TypeAlias = Annotated[None, Depends(service_auth)] + + +@router.get("/service", response_model=ServiceConnection) +async def service_connection(auth: Auth) -> ServiceConnection: + import os + + import httpx + user_scope(auth) - return await repository().ingestion_keys() + 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) + if response.status_code == 200: + status: Final = ServiceStatus.model_validate_json(response.content) + return ServiceConnection(url=public_url, connected=True, status=status) + except (ValueError, httpx.HTTPError): + pass + return ServiceConnection(url=public_url, connected=False, status=ServiceStatus()) -@router.delete("/tracing/keys/{key_id}") -async def revoke_ingestion_key(key_id: str, auth: Auth) -> bool: - user_scope(auth, write=True) - await repository().revoke_ingestion_key(key_id) - return True - - -@router.get("/worker/ingestion-credentials", response_model=IngestionSnapshot) -async def ingestion_credentials(worker: WorkerAuth, response: Response) -> IngestionSnapshot: - if not worker.scope.all_teams: - raise HTTPException(403, "Ingestion requires an administrator-managed Lens service") - response.headers["Cache-Control"] = "no-store" +async def credential_snapshot() -> IngestionSnapshot: keys: Final = await repository().ingestion_keys() now: Final = int(datetime.now(timezone.utc).timestamp()) return IngestionSnapshot( @@ -167,7 +175,51 @@ async def ingestion_credentials(worker: WorkerAuth, response: Response) -> Inges ) -async def assigned(lens_id: str, job_id: str, worker: Worker) -> tuple[Lens, Job]: +async def publish_credentials() -> bool: + import httpx + + try: + 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) + return response.status_code == 204 + except (ValueError, httpx.HTTPError): + return False + + +@router.post("/tracing/keys", response_model=IngestionKeyCreated) +async def create_ingestion_key(body: IngestionKeyRequest, auth: Auth) -> IngestionKeyCreated: + user_scope(auth, write=True) + try: + created: Final = new_key(body, auth.user_id or "") + except ValueError as error: + raise HTTPException(422, str(error)) from error + await repository().save_ingestion_key(created.record) + return created.model_copy(update={"active": await publish_credentials()}) + + +@router.get("/tracing/keys", response_model=tuple[IngestionKey, ...]) +async def list_ingestion_keys(auth: Auth) -> tuple[IngestionKey, ...]: + user_scope(auth) + return await repository().ingestion_keys() + + +@router.delete("/tracing/keys/{key_id}") +async def revoke_ingestion_key(key_id: str, auth: Auth) -> bool: + user_scope(auth, write=True) + await repository().revoke_ingestion_key(key_id) + await publish_credentials() + return True + + +@router.get("/internal/ingestion-credentials", response_model=IngestionSnapshot) +async def ingestion_credentials(service: ServiceAuth, response: Response) -> IngestionSnapshot: + response.headers["Cache-Control"] = "no-store" + return await credential_snapshot() + + +async def assigned(lens_id: str, job_id: str, worker: Worker, attempt: int = 1) -> tuple[Lens, Job]: lens: Final = await get_lens(lens_id, worker.scope) job: Final = current_job(lens) if ( @@ -175,6 +227,7 @@ async def assigned(lens_id: str, job_id: str, worker: Worker) -> tuple[Lens, Job or job.id != job_id or job.status != "running" or job.worker_id != worker.id + or job.attempts != attempt or job.lease_until is None or job.lease_until <= datetime.now(timezone.utc) ): @@ -484,6 +537,7 @@ class WorkerBilling(LiteLLMBaseModel): class WorkerName(WorkerBilling): name: str = Field(default="Lens worker", min_length=1) + managed: bool = False def configured_worker_image() -> str: @@ -501,7 +555,11 @@ async def register_worker(body: WorkerName, auth: Auth) -> WorkerCreated: scope: Final = user_scope(auth, write=True) image: Final = configured_worker_image() await validate_key(body.analysis_key_id) - token: Final = "lens-" + secrets.token_urlsafe(40) + try: + token: Final = LensConnection.from_env().token if body.managed else "lens-" + secrets.token_urlsafe(40) + except ValueError as error: + raise HTTPException(503, "Configure the Lens service before enabling investigations") from error + token_hash: Final = hashlib.sha256(token.encode()).hexdigest() worker: Final = Worker( id=str(uuid4()), name=body.name, @@ -509,7 +567,10 @@ async def register_worker(body: WorkerName, auth: Auth) -> WorkerCreated: analysis_key_id=body.analysis_key_id, last_seen=datetime(1970, 1, 1, tzinfo=timezone.utc), ) - await repository().save_worker(worker, hashlib.sha256(token.encode()).hexdigest()) + if body.managed: + managed: Final = await repository().configure_service_worker(worker, token_hash) + return WorkerCreated(worker=managed, token="", image=image, managed=True) + await repository().save_worker(worker, token_hash) return WorkerCreated(worker=worker, token=token, image=image) @@ -560,8 +621,8 @@ async def claim(worker: WorkerAuth, protocol_version: int = 1, worker_release: s @router.post("/worker/{lens_id}/{job_id}/progress", response_model=bool) -async def progress(lens_id: str, job_id: str, body: Progress, worker: WorkerAuth) -> bool: - _, assigned_job = await assigned(lens_id, job_id, worker) +async def progress(lens_id: str, job_id: str, body: Progress, worker: WorkerAuth, attempt: Attempt = 1) -> bool: + _, assigned_job = await assigned(lens_id, job_id, worker, attempt) if body.review is not None: if assigned_job.sample is None or body.review.execution_id not in frozenset( execution.id for execution in assigned_job.sample.executions @@ -579,14 +640,14 @@ async def progress(lens_id: str, job_id: str, body: Progress, worker: WorkerAuth @router.get("/worker/{lens_id}/{job_id}/reviews", response_model=tuple[Review, ...]) -async def cached_reviews(lens_id: str, job_id: str, worker: WorkerAuth) -> tuple[Review, ...]: - _, job = await assigned(lens_id, job_id, worker) +async def cached_reviews(lens_id: str, job_id: str, worker: WorkerAuth, attempt: Attempt = 1) -> tuple[Review, ...]: + _, job = await assigned(lens_id, job_id, worker, attempt) return await repository().reviews(lens_id, job) @router.get("/worker/{lens_id}/{job_id}/sample", response_model=Sample) -async def sample(lens_id: str, job_id: str, worker: WorkerAuth, storage: StorageDep) -> Sample: - lens, job = await assigned(lens_id, job_id, worker) +async def sample(lens_id: str, job_id: str, worker: WorkerAuth, storage: StorageDep, attempt: Attempt = 1) -> Sample: + lens, job = await assigned(lens_id, job_id, worker, attempt) if job.sample is not None: return job.sample pages: list[Sample] = [] # mutable-ok: freeze selection after stable cursor traversal @@ -610,7 +671,7 @@ 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: + if active is None or active.id != job_id or active.worker_id != worker.id or active.attempts != attempt: raise HTTPException(409, "Job was cancelled or reassigned") return ( replace_job(e, active.model_copy(update=MappingProxyType({"sample": selected}))) @@ -634,8 +695,9 @@ async def content( storage: StorageDep, cursor: str = "", offset: int = Query(default=0, ge=0), + attempt: Attempt = 1, ) -> ExecutionContent: - lens, job = await assigned(lens_id, job_id, worker) + lens, job = await assigned(lens_id, job_id, worker, attempt) selected: Final = job.sample or Sample(executions=(), eligible=0) execution: Final = next((e for e in selected.executions if e.id == execution_id), None) if execution is None: @@ -656,11 +718,11 @@ 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 + 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 - lens, job = await assigned(lens_id, job_id, worker) + lens, job = await assigned(lens_id, job_id, worker, attempt) try: completion: Final = await analyze(repository(), lens, job, worker, body, request) except (ProxyException, HTTPException) as error: @@ -671,14 +733,14 @@ 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) -> 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: + if old and old.status in ("completed", "failed") and old.worker_id == worker.id and old.attempts == attempt: if old.review_versions and old.status == "completed": await repository().complete_reviews(lens_id, old, old.review_versions) return lens - _, job = await assigned(lens_id, job_id, worker) + _, job = await assigned(lens_id, job_id, worker, attempt) now: Final = datetime.now(timezone.utc) selected: Final = job.sample or Sample(executions=(), eligible=0) allowed: Final = frozenset(e.id for e in selected.executions) @@ -707,7 +769,7 @@ 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: + if active is None or active.id != job_id or active.worker_id != worker.id or active.attempts != attempt: return e restored: Final = e.model_copy( update=MappingProxyType( @@ -797,8 +859,8 @@ def merge_results(lens: Lens, result: Result, revision: int, now: datetime, job_ @router.post("/worker/{lens_id}/{job_id}/heartbeat", response_model=bool) -async def heartbeat(lens_id: str, job_id: str, worker: WorkerAuth) -> bool: - return await progress(lens_id, job_id, Progress(), worker) +async def heartbeat(lens_id: str, job_id: str, worker: WorkerAuth, attempt: Attempt = 1) -> bool: + return await progress(lens_id, job_id, Progress(), worker, attempt) async def claim_candidate(candidate: Lens, worker: Worker, now: datetime) -> Claim | None: diff --git a/litellm/proxy/lens/ingestion.py b/litellm/proxy/lens/ingestion.py index 0297224951b..912d819dc68 100644 --- a/litellm/proxy/lens/ingestion.py +++ b/litellm/proxy/lens/ingestion.py @@ -44,6 +44,20 @@ class IngestionSnapshot(Record): class IngestionKeyCreated(Record): key: str record: IngestionKey + active: bool = False + + +class ServiceStatus(Record): + storage_ready: bool = False + credentials_ready: bool = False + release: str = "" + protocol_version: int = 0 + + +class ServiceConnection(Record): + url: str + connected: bool + status: ServiceStatus def new_key(request: IngestionKeyRequest, user_id: str) -> IngestionKeyCreated: diff --git a/litellm/proxy/lens/models.py b/litellm/proxy/lens/models.py index 7d3f791e55a..8bb9d611537 100644 --- a/litellm/proxy/lens/models.py +++ b/litellm/proxy/lens/models.py @@ -402,6 +402,7 @@ class WorkerCreated(Record): image: str worker: Worker token: str + managed: bool = False class LensList(Record): diff --git a/litellm/proxy/lens/repository.py b/litellm/proxy/lens/repository.py index 3f3e70236d8..7b57ef1beb9 100644 --- a/litellm/proxy/lens/repository.py +++ b/litellm/proxy/lens/repository.py @@ -391,6 +391,19 @@ class LensRepository: 'UPDATE "LiteLLM_LensWorker" SET data=$1::jsonb WHERE id=$2', worker.model_dump_json(), worker.id ) + async def configure_service_worker(self, worker: Worker, token_hash: str) -> Worker: + 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 ' + "SET data=jsonb_set(EXCLUDED.data, '{id}', to_jsonb(existing.id)) RETURNING data", + worker.id, + token_hash, + worker.model_dump_json(), + ) + ) + return Worker.model_validate(rows[0].data) + async def set_worker_billing(self, worker_id: str, key_id: str) -> Worker | None: rows: Final = _ROWS.validate_python( await self.db.query_raw( diff --git a/ui/litellm-dashboard/src/components/lens/data/service.ts b/ui/litellm-dashboard/src/components/lens/data/service.ts index 9bd7919308b..9de75206902 100644 --- a/ui/litellm-dashboard/src/components/lens/data/service.ts +++ b/ui/litellm-dashboard/src/components/lens/data/service.ts @@ -181,7 +181,7 @@ export function liveLensApi(client: LensClient, apiClient: ApiClient, accessToke required( client.POST("/lens/workers/register", { headers, - body: { name: "Lens worker", analysis_key_id: analysisKeyId }, + body: { name: "Lens worker", analysis_key_id: analysisKeyId, managed: true }, }), ), setWorkerBillingKey: (workerId, analysisKeyId) => 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 52d5597cc5a..f8c4cd96bcd 100644 --- a/ui/litellm-dashboard/src/components/lens/onboarding/tracing/TracingSetupCard.tsx +++ b/ui/litellm-dashboard/src/components/lens/onboarding/tracing/TracingSetupCard.tsx @@ -2,6 +2,8 @@ import { ArrowRight, ArrowUpRight, Check, Copy, KeyRound, Loader2, Send } from "lucide-react"; import { useState } from "react"; +import { useQuery } from "@tanstack/react-query"; +import type { components } from "@/lib/http/schema"; import { useTimeout } from "usehooks-ts"; import { cn } from "@/lib/cva.config"; @@ -12,7 +14,7 @@ import { copyToClipboard } from "@/utils/dataUtils"; import anthropicLogo from "../../../../../public/assets/logos/anthropic.svg"; import openaiLogo from "../../../../../public/assets/logos/openai_small.svg"; import otelLogo from "../../../../../public/assets/logos/opentelemetry.svg"; -import { agentTraceCall, apiClient, getProxyBaseUrl, sendOtlpTraceCall } from "../../../networking"; +import { agentTraceCall, apiClient, getProxyBaseUrl } from "../../../networking"; import { ActiveDot } from "../../traces/ui/ActiveDot"; import { sampleTraceExport } from "./sampleTrace"; import { FRAMEWORKS, frameworkSnippet, type FrameworkGuide } from "./tracingSetupGuides"; @@ -20,13 +22,9 @@ import type { TraceSummary } from "../../traces/types"; const COPIED_RESET_MS = 1500; const DOCS_URL = "https://docs.litellm.ai/docs/proxy/lens"; -const EXAMPLE_MODEL = "openai/gpt-6-sol"; +const EXAMPLE_MODEL = "openai/gpt-6.1-sol"; const SAMPLE_TRACE_POLL_MS = 1000; -export const TRACING_KEY_REQUEST = { - key_alias: "Agent tracing", - allowed_routes: ["/v1/traces"], - metadata: { purpose: "agent_tracing" }, -} as const; +export const TRACING_KEY_REQUEST = { name: "Agent tracing" } as const; const SAMPLE_TRACE_POLL_ATTEMPTS = 15; type Installer = "pip" | "uv"; @@ -39,25 +37,25 @@ const PY_INSTALL: Record string> = { export const tracingEnvSnippet = (proxyUrl: string, tracingKey: string | null = null): string => [ - ...(tracingKey ? [`export LITELLM_TRACING_KEY=${tracingKey}`] : []), + `export LITELLM_TRACING_KEY="${tracingKey ?? ""}"`, `export OTEL_EXPORTER_OTLP_TRACES_ENDPOINT="${proxyUrl}/v1/traces"`, - `export OTEL_EXPORTER_OTLP_TRACES_HEADERS="Authorization=Bearer $${tracingKey ? "LITELLM_TRACING_KEY" : "LITELLM_API_KEY"}"`, + `export OTEL_EXPORTER_OTLP_TRACES_HEADERS="Authorization=Bearer $LITELLM_TRACING_KEY"`, 'export OTEL_EXPORTER_OTLP_PROTOCOL="http/protobuf"', 'export OTEL_METRICS_EXPORTER="none"', 'export OTEL_LOGS_EXPORTER="none"', ].join("\n"); -export const codingAgentPrompt = (proxyUrl: string, guide: FrameworkGuide, model: string): string => +export const codingAgentPrompt = (proxyUrl: string, traceUrl: string, guide: FrameworkGuide, model: string): string => [ `Send this ${guide.label} project's OpenTelemetry traces to LiteLLM.`, - "Keep the existing model configuration, authentication, and application behavior. Never hardcode a key; read it from LITELLM_API_KEY.", + "Keep the existing model configuration, authentication, and application behavior. Never hardcode keys. Read the model key from LITELLM_API_KEY and the dedicated tracing key from LITELLM_TRACING_KEY.", "Set the trace destination wherever this project loads environment variables:", - tracingEnvSnippet(proxyUrl), + tracingEnvSnippet(traceUrl), guide.install ?? `Install and enable the ${guide.plugin?.label}: ${guide.plugin?.url}`, guide.plugin?.instruction ?? "Initialize OpenTelemetry before creating the agent. If the app already configures a tracer provider, keep it and point its exporter at the destination above instead.", "Adapt this example to the existing application, replacing research_agent with the agent's name:", - frameworkSnippet(guide, proxyUrl, model), + frameworkSnippet(guide, proxyUrl, model, traceUrl), guide.note ?? "", "Run the agent once and confirm its named run appears in Lens > Traces.", ] @@ -74,17 +72,14 @@ export const maskSecret = (secret: string): string => export const otlpEndpoints = (proxyUrl: string): readonly (readonly [string, string, boolean])[] => [ ["Traces endpoint", `${proxyUrl}/v1/traces`, true], - ["Auth header", "Authorization: Bearer ", true], + ["Auth header", "Authorization: Bearer ", true], ["Protocol", "OTLP/HTTP (protobuf or JSON)", false], ]; export const PROXY_CONFIG_SNIPPET = [ - "general_settings:", - " tracing:", - " store:", - " type: clickhouse", - " url: os.environ/CLICKHOUSE_URL", - " retention_days: 14", + 'export LITELLM_LENS_URL="http://lens-worker:4318"', + 'export LITELLM_LENS_PUBLIC_URL="https://traces.example.com"', + 'export LITELLM_LENS_SERVICE_TOKEN=""', ].join("\n"); function CodeBlock({ @@ -203,9 +198,13 @@ async function waitForTrace(accessToken: string, traceId: string): Promise void; }) { const [state, setState] = useState({ kind: "idle" }); @@ -213,7 +212,16 @@ function SendTestTrace({ setState({ kind: "sending" }); const sample = sampleTraceExport(Date.now()); try { - await sendOtlpTraceCall(accessToken, sample.body); + if (!tracingKey) throw new Error("Generate a tracing key first"); + const response = await fetch(`${traceUrl}/v1/traces`, { + method: "POST", + credentials: "omit", + redirect: "error", + headers: { "Content-Type": "application/json", Authorization: `Bearer ${tracingKey}` }, + body: JSON.stringify(sample.body), + signal: AbortSignal.timeout(15000), + }); + if (!response.ok) throw new Error(`Trace upload failed (HTTP ${response.status})`); } catch { setState({ kind: "failed", message: "Could not send the test trace." }); return; @@ -241,7 +249,7 @@ function SendTestTrace({ const busy = state.kind === "sending" || state.kind === "waiting"; return (
-
{missingAfterCheck && (

- No traces received yet. Check the exporter URL and LiteLLM key in your agent’s environment, then check its - logs for export errors. + No traces received yet. Check the Lens URL and tracing key in your agent’s environment, then check its logs + for export errors.

)} @@ -305,7 +313,7 @@ function TracingKey({ setCreating(true); setError(""); try { - const result = await apiClient.post<{ key?: string }>("/key/generate", { + const result = await apiClient.post("/lens/tracing/keys", { accessToken, body: TRACING_KEY_REQUEST, }); @@ -323,8 +331,8 @@ function TracingKey({ Your tracing key} />

Hidden for safety. Copy copies the full key, and the environment step below includes it. This key can only - send traces, so your agent still needs its own key for model calls. Manage it under Virtual Keys as - "Agent tracing". + send traces and check delivery. Your agent still needs its own key for model calls. Save this key before + leaving the page.

); @@ -339,7 +347,7 @@ function TracingKey({ )} Generate tracing key - Or use any existing LiteLLM virtual key. + Use a dedicated Lens key for tracing. {error &&

{error}

} ); @@ -416,17 +424,17 @@ function EnableTracing({ checked, checking, onCheck }: { checked: boolean; check <>

- Set your ClickHouse URL, add this to config.yaml, then restart the proxy. Ask your proxy administrator if you - don’t manage this deployment. + Run the Lens service with ClickHouse access, then set these variables on LiteLLM and restart it. Use the same + service secret on both services.

- config.yaml} /> + LiteLLM environment} /> - ClickHouse and proxy setup
{checked && !checking && ( @@ -444,10 +452,20 @@ function EnableTracing({ checked, checking, onCheck }: { checked: boolean; check ); } -function CodingAgentSetup({ proxyUrl, guide, model }: { proxyUrl: string; guide: FrameworkGuide; model: string }) { +function CodingAgentSetup({ + proxyUrl, + traceUrl, + guide, + model, +}: { + proxyUrl: string; + traceUrl: string; + guide: FrameworkGuide; + model: string; +}) { const [codingAgent, setCodingAgent] = useState("Claude Code"); const [copied, setCopied] = useState(null); - const command = codingAgentCommand(codingAgent, codingAgentPrompt(proxyUrl, guide, model)); + const command = codingAgentCommand(codingAgent, codingAgentPrompt(proxyUrl, traceUrl, guide, model)); useTimeout(() => setCopied(null), copied === null ? null : COPIED_RESET_MS); const copy = async () => { if (await copyToClipboard(command)) setCopied(command); @@ -471,8 +489,8 @@ function CodingAgentSetup({ proxyUrl, guide, model }: { proxyUrl: string; guide:

- Run the setup command in your agent’s project. It uses your LITELLM_API_KEY - . + Run the setup command in your agent’s project. It uses LITELLM_TRACING_KEY{" "} + for traces and keeps your model key separate .