Compare commits

...
16 Commits
Author SHA1 Message Date
Nimar 4a13377e35 chore: release v3.167.2 2026-04-10 15:40:12 +02:00
NimarandGitHub 30af822ac9 chore(deps): bump defu (#13100) 2026-04-10 13:32:54 +00:00
NimarandGitHub c2c0b661e7 chore(deps): bump hono to 4.12.12 (#13099) 2026-04-10 13:22:12 +00:00
Hassieb PakzadandGitHub 2e94ebfe4b fix(evals): prevent llm-as-a-judge queue stalls (#13037) 2026-04-10 14:11:59 +02:00
NimarandGitHub b8544b3423 chore(deps): bump next to 16.2.3 (#13092) 2026-04-10 10:20:59 +00:00
NimarandGitHub 24cc309fb8 chore(deps): bump lodash 4.18.1 (#13090) 2026-04-10 09:44:35 +00:00
NimarandGitHub 1ca70d7033 chore(deps): bump langchain 1.1.39 and related (#13089) 2026-04-10 09:33:32 +00:00
NimarandGitHub ba980c302e chore(deps): bump slack and thus axios 1.15.0 (#13088) 2026-04-10 09:26:48 +00:00
Hassieb PakzadandGitHub ea197e4287 fix(llm-execution-tracing): imperatively set internal tracing environment on events (#13085) 2026-04-10 11:28:48 +02:00
NimarandGitHub 0b20e4d366 chore(deps): build go migrate with clickhouse only (#13082)
* chore(deps): build go migrate with clickhouse only

* add comment
2026-04-10 09:21:55 +00:00
NimarandGitHub 31a1a34616 chore(deps): bump node mocks to 1.17.2 (#13087) 2026-04-10 09:16:59 +00:00
Valery MeleshkinandGitHub 3c3d4bf129 chore: add scripts to provision and run local cloud dependencies (Postgres, Redis, ClickHouse, MinIO) and setup/maintenance helpers (#13054) 2026-04-10 11:19:24 +02:00
07cae52cc7 fix: validate Azure blob storage container names (#13080)
* fix: validate Azure blob storage container names

Azure requires container names to be 3-63 chars, lowercase alphanumeric
and hyphens only. Add Zod superRefine validation to the form schema,
tRPC router, and public API schema so invalid names like "Feedback N8N Bot"
are rejected at submission time with a clear error message.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>

* fix: address PR review feedback for Azure container name validation

- Add empty-string guard in validateAzureContainerName to avoid double
  error when bucketName is blank
- Add .min(1) to public API bucketName schema to match tRPC form schema
- Add Fern docs note describing Azure container naming constraints
- Add server test for invalid Azure container name rejection
- Add client test for empty-string guard behavior

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-04-10 09:11:04 +00:00
Ben BachemandGitHub a81edec0be fix(scores-table): Unused omittedFilter prop (#13079) 2026-04-10 09:07:16 +00:00
Ben BachemandGitHub 497179934d fix(web): Table padding issues (#13060)
* fix(web): Table padding issues

* Increase cell padding in `SelectDashboardDialog` and `SelectWidgetDialog`

* Fix memoization comparison for cellPadding in DataTable

* Set cellPadding="comfortable" for `MembersTable` in org settings
2026-04-10 09:07:08 +00:00
NimarandGitHub ad9dfc41a2 chore(deps): bump vitest to 4.1.4 (#13086) 2026-04-10 09:03:48 +00:00
55 changed files with 3186 additions and 894 deletions
+12 -3
View File
@@ -1,10 +1,19 @@
# Dev container Dockerfile
FROM --platform=${BUILDPLATFORM} golang:1.24 AS migrate-builder
ARG TARGETOS
ARG TARGETARCH
ENV CGO_ENABLED=0 \
GOBIN=/out \
GOOS=${TARGETOS} \
GOARCH=${TARGETARCH}
# Build only the ClickHouse migrate CLI used in this repo.
RUN /usr/local/go/bin/go install -trimpath -tags 'clickhouse' -ldflags='-s -w' \
github.com/golang-migrate/migrate/v4/cmd/migrate@v4.19.1
FROM mcr.microsoft.com/devcontainers/universal:2
# Install golang-migrate for database migrations
RUN curl -L https://github.com/golang-migrate/migrate/releases/download/v4.19.1/migrate.linux-amd64.tar.gz | tar xvz && \
chmod +x migrate && \
mv migrate /usr/local/bin/migrate
COPY --from=migrate-builder /out/migrate /usr/local/bin/migrate
# Activate the repo's pinned pnpm via Corepack
RUN corepack enable && corepack prepare pnpm@10.33.0 --activate
+1 -1
View File
@@ -29,7 +29,7 @@
"@langfuse/shared": "workspace:*",
"@opentelemetry/api": ">=1.0.0 <1.10.0",
"https-proxy-agent": "^7.0.6",
"next": "16.2.2",
"next": "16.2.3",
"next-auth": "^4.24.13",
"zod": "^4.3.6"
},
@@ -68,7 +68,7 @@ types:
type: BlobStorageIntegrationType
bucketName:
type: string
docs: Name of the storage bucket
docs: Name of the storage bucket. For AZURE_BLOB_STORAGE, must be a valid Azure container name (3-63 chars, lowercase letters, numbers, and hyphens only, must start and end with a letter or number, no consecutive hyphens).
endpoint:
type: optional<string>
docs: Custom endpoint URL (required for S3_COMPATIBLE type)
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "langfuse",
"version": "3.167.1",
"version": "3.167.2",
"author": "engineering@langfuse.com",
"license": "MIT",
"private": true,
+4 -10
View File
@@ -1,10 +1,9 @@
import tseslint from "typescript-eslint";
import nextCoreWebVitals from "eslint-config-next/core-web-vitals";
import eslintPluginPrettierRecommended from "eslint-plugin-prettier/recommended";
import turboConfig from "eslint-config-turbo/flat";
import "eslint-plugin-only-warn";
export default tseslint.config(
export default [
// Global ignores - include config files
{
name: "langfuse/ignores",
@@ -57,17 +56,12 @@ export default tseslint.config(
// Prettier (last)
eslintPluginPrettierRecommended,
// TypeScript config for TS files
// Note: The old config had a bug (duplicate extends) that prevented TS rules from applying
// Only adding parser + plugin + custom rules to match old behavior
// Layer repo-specific TS rules on top of Next's built-in flat TS config.
// Next already provides the parser and @typescript-eslint plugin here.
{
name: "langfuse/next/typescript",
files: ["**/*.ts", "**/*.tsx"],
plugins: {
"@typescript-eslint": tseslint.plugin,
},
languageOptions: {
parser: tseslint.parser,
globals: {
React: "readonly",
JSX: "readonly",
@@ -103,4 +97,4 @@ export default tseslint.config(
"react/jsx-key": ["error", { warnOnDuplicates: true }],
},
},
);
];
+1 -1
View File
@@ -16,7 +16,7 @@
],
"dependencies": {
"@eslint/js": "^9.39.2",
"eslint-config-next": "16.2.2",
"eslint-config-next": "16.2.3",
"eslint-config-prettier": "^10.1.8",
"eslint-config-turbo": "2.9.5",
"eslint-plugin-only-warn": "^1.1.0",
+9 -9
View File
@@ -84,17 +84,17 @@
"@azure/storage-blob": "^12.26.0",
"@clickhouse/client": "^1.13.0",
"@google-cloud/storage": "^7.19.0",
"@langchain/anthropic": "^1.3.12",
"@langchain/aws": "^1.3.3",
"@langchain/core": "^1.1.34",
"@langchain/google-genai": "^2.1.13",
"@langchain/google-vertexai": "^2.1.13",
"@langchain/openai": "^1.2.3",
"@langchain/anthropic": "^1.3.26",
"@langchain/aws": "^1.3.4",
"@langchain/core": "^1.1.39",
"@langchain/google-genai": "^2.1.26",
"@langchain/google-vertexai": "^2.1.26",
"@langchain/openai": "^1.4.2",
"@opentelemetry/api": ">=1.0.0 <1.10.0",
"@prisma/client": "^6.19.3",
"@react-email/components": "^0.5.1",
"@react-email/render": "^1.2.1",
"@slack/oauth": "^3.0.4",
"@slack/oauth": "3.0.5",
"@slack/web-api": "^7.15.0",
"@types/bcryptjs": "^2.4.6",
"ajv": "^8.18.0",
@@ -108,9 +108,9 @@
"ioredis": "^5.8.2",
"ipaddr.js": "^2.2.0",
"jsonpath-plus": "10.3.0",
"langchain": "^1.2.15",
"langchain": "^1.3.0",
"langfuse-langchain": "3.38.20",
"lodash": "^4.17.23",
"lodash": "^4.18.1",
"lossless-json": "^4.1.1",
"next-auth": "^4.24.13",
"nodemailer": "^7.0.11",
+1 -1
View File
@@ -1 +1 @@
export const VERSION = "v3.167.1";
export const VERSION = "v3.167.2";
@@ -43,7 +43,11 @@ import type { BaseCallbackHandler } from "@langchain/core/callbacks/base";
import { ProxyAgent } from "undici";
import { getInternalTracingHandler } from "./getInternalTracingHandler";
import { decrypt } from "../../encryption";
import { decryptAndParseExtraHeaders } from "./utils";
import {
decryptAndParseExtraHeaders,
executeWithRuntimeTimeout,
RUNTIME_TIMEOUT_ADAPTERS,
} from "./utils";
import { logger } from "../logger";
import { LLMCompletionError } from "./errors";
@@ -453,6 +457,19 @@ export async function fetchLLMCompletion(
metadata: traceSinkParams?.metadata,
};
const runtimeTimeoutEnabled = RUNTIME_TIMEOUT_ADAPTERS.has(
modelParams.adapter,
);
const runtimeTimeoutController = runtimeTimeoutEnabled
? new AbortController()
: undefined;
const runConfigWithTimeout = runtimeTimeoutController
? {
...runConfig,
signal: runtimeTimeoutController.signal,
}
: runConfig;
const thinkingTypes = getThinkingBlockTypes(modelParams.adapter);
try {
@@ -460,17 +477,24 @@ export async function fetchLLMCompletion(
if (params.structuredOutputSchema) {
// Thinking-capable adapters may produce reasoning blocks that corrupt JSON schema
// parsing. Force function calling so the parser reads from tool_calls instead.
const structuredOutputSchema = params.structuredOutputSchema;
const structuredOutputConfig =
thinkingTypes != null
? { method: "functionCalling" as const }
: undefined;
const structuredOutput = await (chatModel as ChatOpenAI)
.withStructuredOutput(
params.structuredOutputSchema,
structuredOutputConfig,
)
.invoke(finalMessages, runConfig);
const structuredOutput = await executeWithRuntimeTimeout({
enabled: runtimeTimeoutEnabled,
timeoutMs,
abortController: runtimeTimeoutController,
operation: () =>
(chatModel as ChatOpenAI)
.withStructuredOutput(
structuredOutputSchema,
structuredOutputConfig,
)
.invoke(finalMessages, runConfigWithTimeout),
});
return structuredOutput;
}
@@ -481,9 +505,15 @@ export async function fetchLLMCompletion(
function: tool,
}));
const result = await chatModel
.bindTools(langchainTools)
.invoke(finalMessages, runConfig);
const result = await executeWithRuntimeTimeout({
enabled: runtimeTimeoutEnabled,
timeoutMs,
abortController: runtimeTimeoutController,
operation: () =>
chatModel
.bindTools(langchainTools)
.invoke(finalMessages, runConfigWithTimeout),
});
// For thinking adapters, strip reasoning blocks from content before parsing
// so ToolCallResponseSchema can validate. Extract reasoning separately.
@@ -512,20 +542,37 @@ export async function fetchLLMCompletion(
}
if (streaming)
return chatModel
.pipe(new BytesOutputParser())
.stream(finalMessages, runConfig);
return await executeWithRuntimeTimeout({
enabled: runtimeTimeoutEnabled,
timeoutMs,
abortController: runtimeTimeoutController,
operation: () =>
chatModel
.pipe(new BytesOutputParser())
.stream(finalMessages, runConfigWithTimeout),
});
// content with thinking blocks can't be handled by StringOutputParser
// Invoke model directly and extract text + reasoning separately.
if (thinkingTypes != null) {
const aiMessage = await chatModel.invoke(finalMessages, runConfig);
const aiMessage = await executeWithRuntimeTimeout({
enabled: runtimeTimeoutEnabled,
timeoutMs,
abortController: runtimeTimeoutController,
operation: () => chatModel.invoke(finalMessages, runConfigWithTimeout),
});
return extractCompletionWithReasoning(aiMessage, thinkingTypes);
}
const completion = await chatModel
.pipe(new StringOutputParser())
.invoke(finalMessages, runConfig);
const completion = await executeWithRuntimeTimeout({
enabled: runtimeTimeoutEnabled,
timeoutMs,
abortController: runtimeTimeoutController,
operation: () =>
chatModel
.pipe(new StringOutputParser())
.invoke(finalMessages, runConfigWithTimeout),
});
return completion;
} catch (e) {
@@ -4,6 +4,11 @@ import { processEventBatch } from "../ingestion/processEventBatch";
import { logger } from "../logger";
import { traceException } from "../instrumentation";
type TracedEvent = {
type: string;
body: Record<string, unknown>;
};
/**
* Extracts and merges generation details from a list of processed events.
* Handles multiple generation-create and generation-update events with the same id.
@@ -73,6 +78,66 @@ export function extractGenerationDetails(
};
}
export function prepareTracedEventsForIngestion(
events: TracedEvent[],
{ environment, prompt }: Pick<TraceSinkParams, "environment" | "prompt">,
): TracedEvent[] {
const blockedSpanIds = new Set<string>();
const blockedSpanNames = [
"RunnableLambda",
"StructuredOutputParser",
"StrOutputParser",
"JsonOutputParser",
];
for (const event of events) {
const eventName = event.body.name;
if (typeof eventName !== "string" || eventName.length === 0) {
continue;
}
if (
blockedSpanNames.includes(eventName) &&
typeof event.body.id === "string"
) {
blockedSpanIds.add(event.body.id);
}
}
return events
.filter((event) => {
if (typeof event.body.id === "string") {
return !blockedSpanIds.has(event.body.id);
}
return true;
})
.map((event) => {
return {
...event,
body: {
...event.body,
environment,
},
};
})
.map((event) => {
if (event.type === "generation-create" && prompt) {
return {
...event,
body: {
...event.body,
promptName: prompt.name,
promptVersion: prompt.version,
},
};
}
return event;
});
}
export function getInternalTracingHandler(traceSinkParams: TraceSinkParams): {
handler: CallbackHandler;
processTracedEvents: () => Promise<void>;
@@ -91,46 +156,13 @@ export function getInternalTracingHandler(traceSinkParams: TraceSinkParams): {
traceSinkParams.targetProjectId,
);
// Filter out unnecessary Langchain spans
const blockedSpanIds = new Set();
const blockedSpanNames = [
"RunnableLambda",
"StructuredOutputParser",
"StrOutputParser",
"JsonOutputParser",
];
for (const event of events) {
const eventName = "name" in event.body ? event.body.name : "";
if (!eventName) continue;
if (blockedSpanNames.includes(eventName) && "id" in event.body) {
blockedSpanIds.add(event.body.id);
}
}
const processedEvents = events
.filter((event) => {
if ("id" in event.body) {
return !blockedSpanIds.has(event.body.id);
}
return true;
})
.map((event: any) => {
// to add the prompt name and version to only generation-type observations
if (event.type === "generation-create" && prompt) {
return {
...event,
body: {
...event.body,
...{ promptName: prompt.name, promptVersion: prompt.version },
},
};
}
return event;
});
const processedEvents = prepareTracedEventsForIngestion(
events as TracedEvent[],
{
environment,
prompt,
},
);
await processEventBatch(
JSON.parse(JSON.stringify(processedEvents)), // stringify to emulate network event batch from network call
+40
View File
@@ -1,9 +1,49 @@
import { z } from "zod";
import { decrypt } from "../../encryption";
import { LLMAdapter } from "./types";
const ExtraHeaderSchema = z.record(z.string(), z.string());
export const RUNTIME_TIMEOUT_ADAPTERS = new Set([
LLMAdapter.VertexAI,
LLMAdapter.GoogleAIStudio,
]);
export async function executeWithRuntimeTimeout<T>({
enabled,
timeoutMs,
abortController,
operation,
}: {
enabled: boolean;
timeoutMs: number;
abortController?: AbortController;
operation: () => Promise<T>;
}): Promise<T> {
if (!enabled) {
return operation();
}
const timeoutError = new Error(`Request timed out after ${timeoutMs}ms`);
let timeoutId: ReturnType<typeof setTimeout> | undefined;
try {
return await Promise.race([
operation(),
new Promise<never>((_, reject) => {
timeoutId = setTimeout(() => {
abortController?.abort(timeoutError);
reject(timeoutError);
}, timeoutMs);
}),
]);
} finally {
if (timeoutId) clearTimeout(timeoutId);
}
}
export function decryptAndParseExtraHeaders(
extraHeaders: string | null | undefined,
) {
+1314 -654
View File
File diff suppressed because it is too large Load Diff
+24
View File
@@ -19,6 +19,30 @@ minimumReleaseAgeExclude:
- "@turbo/windows-arm64@2.9.5"
- "@turbo/linux-64@2.9.5"
- "@turbo/linux-arm64@2.9.5"
- "vitest@4.1.4"
- "@vitest/mocker@4.1.4"
- "@vitest/pretty-format@4.1.4"
- "@vitest/snapshot@4.1.4"
- "@vitest/spy@4.1.4"
- "@vitest/runner@4.1.4"
- "@vitest/expect@4.1.4"
- "@vitest/utils@4.1.4"
- "@vitest/coverage-v8@4.1.4"
- "axios@1.15.0"
- "next@16.2.3"
- "@next/env@16.2.3"
- "eslint-config-next@16.2.3"
- "@next/eslint-plugin-next@16.2.3"
- "@next/swc-darwin-arm64@16.2.3"
- "@next/swc-darwin-x64@16.2.3"
- "@next/swc-linux-arm64-gnu@16.2.3"
- "@next/swc-linux-arm64-musl@16.2.3"
- "@next/swc-linux-x64-gnu@16.2.3"
- "@next/swc-linux-x64-musl@16.2.3"
- "@next/swc-win32-arm64-msvc@16.2.3"
- "@next/swc-win32-x64-msvc@16.2.3"
- "defu@6.1.7"
- "hono@4.12.12"
allowBuilds:
"@prisma/client": true
"@prisma/engines": true
+461
View File
@@ -0,0 +1,461 @@
#!/usr/bin/env bash
set -euo pipefail
CODEX_SERVICES_ROOT="${CODEX_SERVICES_ROOT:-$PWD/.codex/services}"
# NOTE: POSTGRES_PORT and POSTGRES_USER are effectively immutable once
# `$CODEX_SERVICES_ROOT/postgres/data` is initialized. Changing either value on
# reruns requires deleting the initialized Postgres data directory and allowing
# `initdb` to recreate the cluster with the new settings.
POSTGRES_PORT="${POSTGRES_PORT:-5432}"
REDIS_PORT="${REDIS_PORT:-6379}"
CLICKHOUSE_HTTP_PORT="${CLICKHOUSE_HTTP_PORT:-8123}"
CLICKHOUSE_NATIVE_PORT="${CLICKHOUSE_NATIVE_PORT:-9000}"
MINIO_API_PORT="${MINIO_API_PORT:-9090}"
MINIO_CONSOLE_PORT="${MINIO_CONSOLE_PORT:-9091}"
POSTGRES_USER="${POSTGRES_USER:-postgres}"
POSTGRES_PASSWORD="${POSTGRES_PASSWORD:-postgres}"
POSTGRES_DB="${POSTGRES_DB:-postgres}"
REDIS_AUTH="${REDIS_AUTH:-myredissecret}"
CLICKHOUSE_USER="${CLICKHOUSE_USER:-clickhouse}"
CLICKHOUSE_PASSWORD="${CLICKHOUSE_PASSWORD:-clickhouse}"
MINIO_ROOT_USER="${MINIO_ROOT_USER:-minio}"
MINIO_ROOT_PASSWORD="${MINIO_ROOT_PASSWORD:-miniosecret}"
MINIO_RELEASE_TAG="${MINIO_RELEASE_TAG:-RELEASE.2025-09-07T16-13-09Z}"
MC_RELEASE_TAG="${MC_RELEASE_TAG:-RELEASE.2025-08-13T08-35-41Z}"
MINIO_SHA256_AMD64="${MINIO_SHA256_AMD64:-7c5bd8512c6e966455b1d198209358b2d191c77a83ab377c4073281065fb855f}"
MINIO_SHA256_ARM64="${MINIO_SHA256_ARM64:-5c83cd2cf151717ba0243f73e1c7802ff36e272b67144bdd7f1f7d684fd6f03d}"
MC_SHA256_AMD64="${MC_SHA256_AMD64:-01f866e9c5f9b87c2b09116fa5d7c06695b106242d829a8bb32990c00312e891}"
MC_SHA256_ARM64="${MC_SHA256_ARM64:-14c8c9616cfce4636add161304353244e8de383b2e2752c0e9dad01d4c27c12c}"
export DEBIAN_FRONTEND=noninteractive
ensure_apt_package() {
local package="$1"
if dpkg -s "$package" >/dev/null 2>&1; then
return 0
fi
if [ -z "${CODEX_APT_UPDATED:-}" ]; then
apt-get update
CODEX_APT_UPDATED=1
fi
apt-get install -y "$package"
}
stop_service_if_running() {
local service_name="$1"
if command -v systemctl >/dev/null 2>&1; then
systemctl stop "$service_name" >/dev/null 2>&1 || true
fi
if command -v service >/dev/null 2>&1; then
service "$service_name" stop >/dev/null 2>&1 || true
fi
}
stop_system_postgres_clusters() {
if command -v pg_lsclusters >/dev/null 2>&1 && command -v pg_ctlcluster >/dev/null 2>&1; then
while read -r version cluster_name _ status _; do
if [ "$status" = "online" ]; then
pg_ctlcluster "$version" "$cluster_name" stop >/dev/null 2>&1 || true
fi
done < <(pg_lsclusters --no-header 2>/dev/null || true)
fi
stop_service_if_running postgresql
}
ensure_clickhouse_repo() {
ensure_apt_package ca-certificates
ensure_apt_package curl
ensure_apt_package gnupg
local keyring="/etc/apt/keyrings/clickhouse.gpg"
local source_file="/etc/apt/sources.list.d/clickhouse.list"
mkdir -p /etc/apt/keyrings
if [ ! -f "$keyring" ]; then
curl -fsSL https://packages.clickhouse.com/rpm/lts/repodata/repomd.xml.key \
| gpg --dearmor -o "$keyring"
fi
if [ ! -f "$source_file" ]; then
echo "deb [signed-by=$keyring] https://packages.clickhouse.com/deb stable main" > "$source_file"
apt-get update
fi
}
ensure_postgres_binaries() {
ensure_apt_package postgresql
ensure_apt_package postgresql-client
stop_system_postgres_clusters
}
ensure_redis_binary() {
ensure_apt_package redis-server
stop_service_if_running redis-server
}
ensure_clickhouse_binaries() {
if command -v clickhouse-server >/dev/null 2>&1 && command -v clickhouse-client >/dev/null 2>&1; then
stop_service_if_running clickhouse-server
return 0
fi
ensure_clickhouse_repo
apt-get install -y clickhouse-server clickhouse-client
stop_service_if_running clickhouse-server
}
detect_minio_arch() {
local machine_arch
machine_arch="$(uname -m)"
case "$machine_arch" in
x86_64|amd64)
echo "amd64"
;;
aarch64|arm64)
echo "arm64"
;;
*)
echo "Unsupported architecture for MinIO binaries: $machine_arch" >&2
exit 1
;;
esac
}
download_and_verify_sha256() {
local url="$1"
local output_path="$2"
local expected_sha256="$3"
local tmp_download
tmp_download="$(mktemp)"
trap 'rm -f "$tmp_download"' RETURN
curl -fsSL "$url" -o "$tmp_download"
local actual_sha256
actual_sha256="$(sha256sum "$tmp_download" | awk '{print $1}')"
if [ "$actual_sha256" != "$expected_sha256" ]; then
echo "SHA256 mismatch for $url" >&2
echo "expected: $expected_sha256" >&2
echo "actual: $actual_sha256" >&2
return 1
fi
mv "$tmp_download" "$output_path"
trap - RETURN
}
ensure_minio_binaries() {
local bin_dir="$CODEX_SERVICES_ROOT/bin"
local minio_arch
local minio_sha256
local mc_sha256
mkdir -p "$bin_dir"
minio_arch="$(detect_minio_arch)"
case "$minio_arch" in
amd64)
minio_sha256="$MINIO_SHA256_AMD64"
mc_sha256="$MC_SHA256_AMD64"
;;
arm64)
minio_sha256="$MINIO_SHA256_ARM64"
mc_sha256="$MC_SHA256_ARM64"
;;
esac
if [ ! -x "$bin_dir/minio" ]; then
download_and_verify_sha256 \
"https://dl.min.io/server/minio/release/linux-${minio_arch}/archive/minio.${MINIO_RELEASE_TAG}" \
"$bin_dir/minio" \
"$minio_sha256"
chmod +x "$bin_dir/minio"
fi
if [ ! -x "$bin_dir/mc" ]; then
download_and_verify_sha256 \
"https://dl.min.io/client/mc/release/linux-${minio_arch}/archive/mc.${MC_RELEASE_TAG}" \
"$bin_dir/mc" \
"$mc_sha256"
chmod +x "$bin_dir/mc"
fi
export PATH="$bin_dir:$PATH"
}
find_postgres_bin() {
local name="$1"
if command -v "$name" >/dev/null 2>&1; then
command -v "$name"
return 0
fi
find /usr/lib/postgresql -type f -name "$name" 2>/dev/null | sort -V | tail -n 1
}
wait_for_port() {
local host="$1"
local port="$2"
local timeout_seconds="${3:-45}"
local deadline=$((SECONDS + timeout_seconds))
until (echo >"/dev/tcp/$host/$port") >/dev/null 2>&1; do
if [ "$SECONDS" -ge "$deadline" ]; then
return 1
fi
sleep 1
done
}
wait_for_http() {
local url="$1"
local timeout_seconds="${2:-45}"
local deadline=$((SECONDS + timeout_seconds))
until curl -fsS "$url" >/dev/null 2>&1; do
if [ "$SECONDS" -ge "$deadline" ]; then
return 1
fi
sleep 1
done
}
escape_sql_literal() {
local value="$1"
value="${value//\\/\\\\}"
printf "%s" "${value//\'/\'\'}"
}
escape_clickhouse_identifier() {
local value="$1"
printf '`%s`' "${value//\`/\`\`}"
}
escape_redis_config_string() {
local value="$1"
# Redis treats backslashes and double-quotes as escape delimiters inside
# quoted config strings, so both must be escaped before writing requirepass.
value="${value//\\/\\\\}"
value="${value//$'\n'/\\n}"
value="${value//\"/\\\"}"
printf "%s" "$value"
}
ensure_postgres_running() {
ensure_postgres_binaries
local initdb
local pg_ctl
local psql
local pg_isready
initdb="$(find_postgres_bin initdb)"
pg_ctl="$(find_postgres_bin pg_ctl)"
psql="$(find_postgres_bin psql)"
pg_isready="$(find_postgres_bin pg_isready)"
if [ -z "$initdb" ] || [ -z "$pg_ctl" ] || [ -z "$psql" ] || [ -z "$pg_isready" ]; then
echo "Unable to find required PostgreSQL binaries (initdb, pg_ctl, psql, pg_isready)."
exit 1
fi
local pg_root="$CODEX_SERVICES_ROOT/postgres"
local pg_data="$pg_root/data"
local pg_log="$pg_root/postgres.log"
local pg_socket_dir="$pg_root"
local -a pg_runner
mkdir -p "$pg_root"
if [ "${EUID:-$(id -u)}" -eq 0 ] && id -u postgres >/dev/null 2>&1; then
chown -R postgres:postgres "$pg_root"
pg_runner=(runuser -u postgres --)
else
pg_runner=()
fi
if [ ! -f "$pg_data/PG_VERSION" ]; then
"${pg_runner[@]}" "$initdb" -D "$pg_data" -U "$POSTGRES_USER" --auth-host=md5 --auth-local=trust >/dev/null
{
echo "listen_addresses = '127.0.0.1'"
echo "port = $POSTGRES_PORT"
echo "log_statement = 'all'"
echo "timezone = 'UTC'"
echo "unix_socket_directories = '$pg_socket_dir'"
} >> "$pg_data/postgresql.conf"
fi
if ! "${pg_runner[@]}" "$pg_ctl" -D "$pg_data" status >/dev/null 2>&1; then
"${pg_runner[@]}" "$pg_ctl" -D "$pg_data" -l "$pg_log" -w start
fi
if ! "$pg_isready" -h "$pg_socket_dir" -p "$POSTGRES_PORT" -U "$POSTGRES_USER" >/dev/null 2>&1; then
echo "PostgreSQL did not become ready on socket $pg_socket_dir (port $POSTGRES_PORT)"
exit 1
fi
PGPASSWORD="${POSTGRES_PASSWORD}" "${pg_runner[@]}" "$psql" -h "$pg_socket_dir" -p "$POSTGRES_PORT" -U "$POSTGRES_USER" -d postgres -v postgres_user="$POSTGRES_USER" -v postgres_db="$POSTGRES_DB" -v postgres_password="$POSTGRES_PASSWORD" <<SQL >/dev/null
SELECT format('ALTER USER %I WITH PASSWORD %L', :'postgres_user', :'postgres_password')\gexec
SELECT format('CREATE DATABASE %I', :'postgres_db')
WHERE NOT EXISTS (SELECT FROM pg_database WHERE datname = :'postgres_db')\gexec
SQL
}
ensure_redis_running() {
ensure_redis_binary
local redis_root="$CODEX_SERVICES_ROOT/redis"
local redis_conf="$redis_root/redis.conf"
local redis_log="$redis_root/redis.log"
local redis_pid="$redis_root/redis.pid"
local redis_auth_escaped
mkdir -p "$redis_root"
if wait_for_port 127.0.0.1 "$REDIS_PORT" 1; then
echo "Redis already running on 127.0.0.1:$REDIS_PORT; keeping existing runtime config."
return 0
fi
redis_auth_escaped="$(escape_redis_config_string "$REDIS_AUTH")"
cat > "$redis_conf" <<CONF
bind 127.0.0.1
port $REDIS_PORT
requirepass "$redis_auth_escaped"
maxmemory-policy noeviction
daemonize yes
pidfile "$redis_pid"
logfile "$redis_log"
dir "$redis_root"
CONF
redis-server "$redis_conf"
if ! wait_for_port 127.0.0.1 "$REDIS_PORT" 30; then
echo "Redis did not start on 127.0.0.1:$REDIS_PORT"
exit 1
fi
}
ensure_clickhouse_running() {
ensure_clickhouse_binaries
local clickhouse_root="$CODEX_SERVICES_ROOT/clickhouse"
local clickhouse_data="$clickhouse_root/data"
local clickhouse_log="$clickhouse_root/clickhouse.log"
local clickhouse_err="$clickhouse_root/clickhouse.err.log"
local clickhouse_pid="$clickhouse_root/clickhouse.pid"
mkdir -p "$clickhouse_data"
if [ "${EUID:-$(id -u)}" -eq 0 ] && id -u clickhouse >/dev/null 2>&1; then
chown -R clickhouse:clickhouse "$clickhouse_root"
fi
if ! wait_for_http "http://127.0.0.1:$CLICKHOUSE_HTTP_PORT/ping" 1; then
clickhouse-server \
--daemon \
--config-file=/etc/clickhouse-server/config.xml \
--pid-file="$clickhouse_pid" \
--log-file="$clickhouse_log" \
--errorlog-file="$clickhouse_err" \
-- \
--path="$clickhouse_data" \
--http_port="$CLICKHOUSE_HTTP_PORT" \
--tcp_port="$CLICKHOUSE_NATIVE_PORT"
fi
if ! wait_for_http "http://127.0.0.1:$CLICKHOUSE_HTTP_PORT/ping" 45; then
echo "ClickHouse did not start on 127.0.0.1:$CLICKHOUSE_HTTP_PORT"
exit 1
fi
local clickhouse_password_sql
local clickhouse_user_identifier
clickhouse_password_sql="$(escape_sql_literal "$CLICKHOUSE_PASSWORD")"
clickhouse_user_identifier="$(escape_clickhouse_identifier "$CLICKHOUSE_USER")"
clickhouse-client --host 127.0.0.1 --port "$CLICKHOUSE_NATIVE_PORT" -q "CREATE USER IF NOT EXISTS $clickhouse_user_identifier IDENTIFIED WITH plaintext_password BY '$clickhouse_password_sql'"
clickhouse-client --host 127.0.0.1 --port "$CLICKHOUSE_NATIVE_PORT" -q "GRANT ALL ON *.* TO $clickhouse_user_identifier WITH GRANT OPTION"
}
ensure_minio_running() {
ensure_minio_binaries
local minio_root="$CODEX_SERVICES_ROOT/minio"
local minio_data="$minio_root/data"
local minio_log="$minio_root/minio.log"
local minio_pid="$minio_root/minio.pid"
local minio_already_running="false"
mkdir -p "$minio_data"
if wait_for_port 127.0.0.1 "$MINIO_API_PORT" 1; then
echo "MinIO already running on 127.0.0.1:$MINIO_API_PORT; skipping server start."
minio_already_running="true"
fi
if [ "$minio_already_running" != "true" ]; then
(
export MINIO_ROOT_USER MINIO_ROOT_PASSWORD
nohup minio server \
--address "127.0.0.1:$MINIO_API_PORT" \
--console-address "127.0.0.1:$MINIO_CONSOLE_PORT" \
"$minio_data" >"$minio_log" 2>&1 &
echo $! > "$minio_pid"
)
fi
if ! wait_for_port 127.0.0.1 "$MINIO_API_PORT" 45; then
echo "MinIO did not start on 127.0.0.1:$MINIO_API_PORT"
exit 1
fi
if ! mc alias set local "http://127.0.0.1:$MINIO_API_PORT" "$MINIO_ROOT_USER" "$MINIO_ROOT_PASSWORD" >/dev/null 2>&1; then
if [ "$minio_already_running" = "true" ]; then
echo "MinIO is running but credentials do not match MINIO_ROOT_USER/MINIO_ROOT_PASSWORD; skipping bucket reconciliation."
return 0
fi
echo "Failed to configure MinIO client alias for fresh MinIO startup."
exit 1
fi
if ! mc mb --ignore-existing local/langfuse >/dev/null 2>&1; then
if [ "$minio_already_running" = "true" ]; then
echo "Failed to reconcile MinIO bucket 'langfuse'; will retry on next run."
return 0
fi
echo "Failed to create MinIO bucket 'langfuse' after fresh startup."
exit 1
fi
}
ensure_cloud_dependencies() {
mkdir -p "$CODEX_SERVICES_ROOT"
ensure_postgres_running
ensure_redis_running
ensure_clickhouse_running
ensure_minio_running
echo "Cloud dependencies are installed and running:"
echo "- PostgreSQL on 127.0.0.1:$POSTGRES_PORT"
echo "- Redis on 127.0.0.1:$REDIS_PORT"
echo "- ClickHouse HTTP on 127.0.0.1:$CLICKHOUSE_HTTP_PORT, native on 127.0.0.1:$CLICKHOUSE_NATIVE_PORT"
echo "- MinIO API on 127.0.0.1:$MINIO_API_PORT, console on 127.0.0.1:$MINIO_CONSOLE_PORT"
}
+20
View File
@@ -0,0 +1,20 @@
#!/usr/bin/env bash
set -euo pipefail
if ! command -v corepack >/dev/null 2>&1; then
echo "corepack is required. Use a Codex base environment with Node.js 24 support."
exit 1
fi
corepack enable
corepack prepare pnpm@10.33.0 --activate
# shellcheck source=/dev/null
source "$(dirname "${BASH_SOURCE[0]}")/cloud_services.sh"
ensure_cloud_dependencies
pnpm install --frozen-lockfile
# Keep generated Prisma artifacts aligned after dependency or schema updates.
pnpm run db:generate
+42
View File
@@ -0,0 +1,42 @@
#!/usr/bin/env bash
set -euo pipefail
ensure_env_file() {
local target_path="$1"
local fallback_path="$2"
if [ -f "$target_path" ]; then
return 0
fi
cp "$fallback_path" "$target_path"
}
if ! command -v corepack >/dev/null 2>&1; then
echo "corepack is required. Use a Codex base environment with Node.js 24 support."
exit 1
fi
corepack enable
corepack prepare pnpm@10.33.0 --activate
ensure_env_file .env .env.dev.example
ensure_env_file .env.test .env.test.example
# shellcheck source=/dev/null
source "$(dirname "${BASH_SOURCE[0]}")/cloud_services.sh"
ensure_cloud_dependencies
pnpm install --frozen-lockfile
# Install Chromium into the default user-level Playwright cache so frontend
# browser review works on first bootstrap.
pnpm run playwright:install
# Generate the shared Prisma client explicitly in the current worktree before
# the workspace-wide db:generate task, which may be satisfied by Turbo cache.
pnpm --filter=shared run db:generate
# Prisma client generation is needed for typecheck/build tasks in Codex.
pnpm run db:generate
+15 -3
View File
@@ -19,6 +19,20 @@ FROM --platform=${TARGETPLATFORM:-linux/amd64} alpine AS runtime-base
RUN rm -rf /usr/local/lib/node_modules/corepack && \
rm -f /usr/local/bin/corepack /usr/local/bin/yarn /usr/local/bin/yarnpkg
FROM --platform=${BUILDPLATFORM} golang:1.24 AS migrate-builder
ARG TARGETOS
ARG TARGETARCH
ENV CGO_ENABLED=0 \
GOBIN=/out \
GOOS=${TARGETOS} \
GOARCH=${TARGETARCH}
# Build only the ClickHouse migrate CLI Langfuse uses at runtime.
# compile this ourselves instead of downloading the upstream release
# because prebuilt bins bundle many unused drivers and thus inherit CVEs
# eg.: https://github.com/golang-migrate/migrate/issues/1357
RUN /usr/local/go/bin/go install -trimpath -tags 'clickhouse' -ldflags='-s -w' \
github.com/golang-migrate/migrate/v4/cmd/migrate@v4.19.1
FROM --platform=${TARGETPLATFORM:-linux/amd64} build-base AS pruner
WORKDIR /app
@@ -134,9 +148,7 @@ RUN if [ -n "$NEXT_PUBLIC_LANGFUSE_CLOUD_REGION" ]; then \
RUN rm -rf /usr/local/lib/node_modules/npm && \
rm -f /usr/local/bin/npm /usr/local/bin/npx
RUN MIGRATE_TARGET_ARCH=$(echo ${TARGETPLATFORM:-linux/amd64} | sed 's/\//-/g') && \
wget -q -O- https://github.com/golang-migrate/migrate/releases/download/v4.19.1/migrate.$MIGRATE_TARGET_ARCH.tar.gz | tar xvz && \
mv migrate /usr/bin/migrate
COPY --from=migrate-builder /out/migrate /usr/bin/migrate
COPY --from=builder --chown=nextjs:nodejs /app/web/next.config.mjs .
COPY --from=builder --chown=nextjs:nodejs /app/web/package.json .
+7 -7
View File
@@ -1,6 +1,6 @@
{
"name": "web",
"version": "3.167.1",
"version": "3.167.2",
"private": true,
"license": "MIT",
"engines": {
@@ -43,7 +43,7 @@
"@headlessui/tailwindcss": "0.2.2",
"@heroicons/react": "^2.2.0",
"@hookform/resolvers": "^5.2.2",
"@langchain/core": "^1.1.34",
"@langchain/core": "^1.1.39",
"@langfuse/ee": "workspace:*",
"@langfuse/shared": "workspace:*",
"@lezer/highlight": "^1.2.3",
@@ -122,12 +122,12 @@
"ioredis": "^5.8.2",
"ip-address": "^9.0.5",
"json-schema-faker": "^0.5.9",
"langchain": "^1.2.15",
"langchain": "^1.3.0",
"langfuse": "3.38.4",
"lodash": "^4.17.23",
"lodash": "^4.18.1",
"lucide-react": "^0.552.0",
"nanoid": "^3.3.11",
"next": "16.2.2",
"next": "16.2.3",
"next-auth": "^4.24.13",
"next-query-params": "^5.1.0",
"next-themes": "^0.4.6",
@@ -192,10 +192,10 @@
"@typescript/native-preview": "7.0.0-dev.20260122.3",
"dotenv-cli": "^7.4.2",
"eslint": "^9.39.2",
"eslint-config-next": "16.2.2",
"eslint-config-next": "16.2.3",
"jest": "^30.2.0",
"jest-environment-jsdom": "^30.2.0",
"node-mocks-http": "^1.14.1",
"node-mocks-http": "^1.17.2",
"prettier": "^3.8.1",
"prettier-plugin-tailwindcss": "^0.7.2",
"tailwindcss": "^4.2.2",
+5 -1
View File
@@ -6325,7 +6325,11 @@ components:
$ref: '#/components/schemas/BlobStorageIntegrationType'
bucketName:
type: string
description: Name of the storage bucket
description: >-
Name of the storage bucket. For AZURE_BLOB_STORAGE, must be a valid
Azure container name (3-63 chars, lowercase letters, numbers, and
hyphens only, must start and end with a letter or number, no
consecutive hyphens).
endpoint:
type: string
nullable: true
@@ -445,6 +445,24 @@ describe("Blob Storage Integrations API", () => {
);
});
it("should reject invalid Azure container names", async () => {
const azureConfig = {
...validBlobStorageConfig,
projectId: testProject1Id,
type: "AZURE_BLOB_STORAGE" as const,
endpoint: "https://myaccount.blob.core.windows.net",
bucketName: "Feedback N8N Bot",
};
const result = await makeAPICall(
"PUT",
"/api/public/integrations/blob-storage",
azureConfig,
createBasicAuthHeader(testApiKey, testApiSecretKey),
);
expect(result.status).toBe(400);
});
it("should handle export modes with dates", async () => {
const customDateConfig = {
...validBlobStorageConfig,
+9 -1
View File
@@ -77,6 +77,7 @@ interface DataTableProps<TData, TValue> {
tableName: string;
getRowClassName?: (row: TData) => string;
topAlignCells?: boolean;
cellPadding?: "compact" | "comfortable";
}
export interface AsyncTableData<T> {
@@ -165,6 +166,7 @@ export function DataTable<TData extends object, TValue>({
tableName,
getRowClassName,
topAlignCells = false,
cellPadding = "compact",
}: DataTableProps<TData, TValue>) {
const [columnFilters, setColumnFilters] = useState<ColumnFiltersState>([]);
const rowheighttw = getRowHeightTailwindClass(rowHeight, customRowHeights);
@@ -416,6 +418,7 @@ export function DataTable<TData extends object, TValue>({
onRowClick={hasRowClickAction ? handleOnRowClick : undefined}
getRowClassName={getRowClassName}
topAlignCells={topAlignCells}
cellPadding={cellPadding}
tableSnapshot={{
columnVisibility,
columnOrder,
@@ -434,6 +437,7 @@ export function DataTable<TData extends object, TValue>({
onRowClick={hasRowClickAction ? handleOnRowClick : undefined}
getRowClassName={getRowClassName}
topAlignCells={topAlignCells}
cellPadding={cellPadding}
/>
)}
</Table>
@@ -480,6 +484,7 @@ interface TableBodyComponentProps<TData> {
onRowClick?: (row: TData, event?: React.MouseEvent) => void;
getRowClassName?: (row: TData) => string;
topAlignCells?: boolean;
cellPadding?: "compact" | "comfortable";
tableSnapshot?: {
columnVisibility?: VisibilityState;
columnOrder?: ColumnOrderState;
@@ -533,6 +538,7 @@ function TableBodyComponent<TData>({
onRowClick,
getRowClassName,
topAlignCells = false,
cellPadding = "compact",
}: TableBodyComponentProps<TData>) {
return (
<TableBody>
@@ -562,7 +568,8 @@ function TableBodyComponent<TData>({
<TableCell
key={cell.id}
className={cn(
"overflow-hidden border-b px-1 text-xs first:pl-2",
"overflow-hidden border-b text-xs first:pl-2",
cellPadding === "comfortable" ? "p-1" : "px-1",
isSmallRowHeight && "whitespace-nowrap",
getPinningClasses(cell.column),
)}
@@ -654,6 +661,7 @@ const MemoizedTableBody = React.memo(TableBodyComponent, (prev, next) => {
if (prev.data.isLoading !== next.data.isLoading) return false;
if (prev.rowheighttw !== next.rowheighttw) return false;
if (prev.rowHeight !== next.rowHeight) return false;
if (prev.cellPadding !== next.cellPadding) return false;
// Then do more expensive deep equality checks
if (
@@ -358,6 +358,7 @@ export default function ModelTable({ projectId }: { projectId: string }) {
columnOrder={columnOrder}
onColumnOrderChange={setColumnOrder}
rowHeight={rowHeight}
cellPadding="comfortable"
onRowClick={(row) => {
router.push(`/project/${projectId}/settings/models/${row.modelId}`);
}}
@@ -284,6 +284,7 @@ export function ScoreConfigsTable({ projectId }: { projectId: string }) {
columnOrder={columnOrder}
onColumnOrderChange={setColumnOrder}
rowHeight={rowHeight}
cellPadding="comfortable"
className="gap-2"
/>
</SettingsTableCard>
+27 -17
View File
@@ -13,8 +13,9 @@ import { Avatar, AvatarImage } from "@/src/components/ui/avatar";
import useColumnVisibility from "@/src/features/column-visibility/hooks/useColumnVisibility";
import { useSidebarFilterState } from "@/src/features/filters/hooks/useSidebarFilterState";
import {
scoreFilterConfig,
getScoreFilterConfig,
SCORE_COLUMN_TO_BACKEND_KEY,
type ScoresTableHiddenColumn,
} from "@/src/features/filters/config/scores-config";
import { DEFAULT_SIDEBAR_IMPLICIT_ENVIRONMENT_CONFIG } from "@/src/features/filters/constants/internal-environments";
import { transformFiltersForBackend } from "@/src/features/filters/lib/filter-transform";
@@ -80,6 +81,16 @@ export type ScoresTableRow = {
executionTraceId?: string;
};
export type ScoresTableProps = {
projectId: string;
userId?: string;
traceId?: string;
observationId?: string;
hiddenColumns?: ScoresTableHiddenColumn[];
localStorageSuffix?: string;
disableUrlPersistence?: boolean;
};
function createFilterState(
userFilterState: FilterState,
omittedFilters: Record<string, string>[],
@@ -104,16 +115,15 @@ export default function ScoresTable({
hiddenColumns = [],
localStorageSuffix = "",
disableUrlPersistence = false,
}: {
projectId: string;
userId?: string;
traceId?: string;
observationId?: string;
omittedFilter?: string[];
hiddenColumns?: string[];
localStorageSuffix?: string;
disableUrlPersistence?: boolean;
}) {
}: ScoresTableProps) {
const scoresFilterConfig = useMemo(
() => getScoreFilterConfig(hiddenColumns),
[hiddenColumns],
);
const hiddenColumnSet = useMemo(
() => new Set<string>(hiddenColumns),
[hiddenColumns],
);
const { isBetaEnabled } = useV4Beta();
// In v4beta, scores must exclusively use events-backed endpoints (no traces-table route).
const useEventsBackedScores = isBetaEnabled;
@@ -288,7 +298,7 @@ export default function ScoresTable({
);
const queryFilter = useSidebarFilterState(
scoreFilterConfig,
scoresFilterConfig,
newFilterOptions,
{
loading: filterOptions.isPending || environmentFilterOptions.isPending,
@@ -322,7 +332,7 @@ export default function ScoresTable({
const backendFilterState = transformFiltersForBackend(
filterState,
SCORE_COLUMN_TO_BACKEND_KEY,
scoreFilterConfig.columnDefinitions,
scoresFilterConfig.columnDefinitions,
);
const getCountPayload = {
@@ -714,7 +724,7 @@ export default function ScoresTable({
];
const columns = rawColumns.filter(
(c) => !!c.id && !hiddenColumns.includes(c.id),
(c) => !!c.id && !hiddenColumnSet.has(c.id),
);
const [columnVisibility, setColumnVisibility] =
@@ -819,15 +829,15 @@ export default function ScoresTable({
},
validationContext: {
columns,
filterColumnDefinition: scoreFilterConfig.columnDefinitions,
filterColumnDefinition: scoresFilterConfig.columnDefinitions,
},
currentFilterState: queryFilter.explicitFilterState,
});
return (
<DataTableControlsProvider
tableName={scoreFilterConfig.tableName}
defaultSidebarCollapsed={scoreFilterConfig.defaultSidebarCollapsed}
tableName={scoresFilterConfig.tableName}
defaultSidebarCollapsed={scoresFilterConfig.defaultSidebarCollapsed}
>
<div className="flex h-full w-full flex-col">
{/* Toolbar spanning full width */}
@@ -594,12 +594,12 @@ export const ObservationPreview = ({
<ScoresTable
projectId={projectId}
traceId={traceId}
omittedFilter={["Observation ID"]}
observationId={preloadedObservation.id}
hiddenColumns={[
"traceId",
"observationId",
"traceName",
"traceTags",
"jobConfigurationId",
"userId",
]}
+7 -2
View File
@@ -583,9 +583,14 @@ export const TracePreview = ({
<div className="flex h-full min-h-0 w-full flex-col overflow-hidden pr-3 md:flex-1">
<ScoresTable
projectId={trace.projectId}
omittedFilter={["Trace ID"]}
traceId={trace.id}
hiddenColumns={["traceName", "jobConfigurationId", "userId"]}
hiddenColumns={[
"traceId",
"traceName",
"traceTags",
"jobConfigurationId",
"userId",
]}
localStorageSuffix="TracePreview"
disableUrlPersistence
/>
@@ -477,6 +477,7 @@ export function ObservationDetailView({
"traceId",
"observationId",
"traceName",
"traceTags",
"jobConfigurationId",
"userId",
]}
@@ -418,9 +418,14 @@ export function TraceDetailView({
<div className="flex h-full min-h-0 w-full flex-col overflow-hidden pr-3">
<ScoresTable
projectId={projectId}
omittedFilter={["Trace ID"]}
traceId={trace.id}
hiddenColumns={["traceName", "jobConfigurationId", "userId"]}
hiddenColumns={[
"traceId",
"traceName",
"traceTags",
"jobConfigurationId",
"userId",
]}
localStorageSuffix="TracePreview"
disableUrlPersistence={isPeekMode}
/>
+7 -4
View File
@@ -2,6 +2,8 @@ import * as React from "react";
import { cn } from "@/src/utils/tailwind";
type TableDensity = "compact" | "comfortable";
const Table = React.forwardRef<
HTMLTableElement,
React.HTMLAttributes<HTMLTableElement>
@@ -74,7 +76,7 @@ const TableHead = React.forwardRef<
<th
ref={ref}
className={cn(
"bg-background text-muted-foreground relative h-10 border-b px-4 text-left align-middle font-medium [&:has([role=checkbox])]:pr-0",
"bg-background text-muted-foreground relative h-10 border-b px-2 text-left align-middle font-medium [&:has([role=checkbox])]:pr-0",
className,
)}
{...props}
@@ -84,12 +86,13 @@ TableHead.displayName = "TableHead";
const TableCell = React.forwardRef<
HTMLTableCellElement,
React.TdHTMLAttributes<HTMLTableCellElement>
>(({ className, ...props }, ref) => (
React.TdHTMLAttributes<HTMLTableCellElement> & { density?: TableDensity }
>(({ className, density = "compact", ...props }, ref) => (
<td
ref={ref}
className={cn(
"h-full px-2 py-0 align-middle [&:has([role=checkbox])]:pr-0",
"h-full align-middle [&:has([role=checkbox])]:pr-0",
density === "comfortable" ? "p-2" : "px-2 py-0",
"border-b [:last-child_>_&]:border-b-0",
className,
)}
+1 -1
View File
@@ -1 +1 @@
export const VERSION = "v3.167.1";
export const VERSION = "v3.167.2";
@@ -187,6 +187,7 @@ export function AuditLogsTable(props: AuditLogsTableProps) {
state: paginationState,
}}
rowHeight={rowHeight}
cellPadding="comfortable"
/>
</SettingsTableCard>
</>
@@ -6,7 +6,8 @@ import {
createTRPCRouter,
protectedProjectProcedure,
} from "@/src/server/api/trpc";
import { blobStorageIntegrationFormSchema } from "@/src/features/blobstorage-integration/types";
import { blobStorageIntegrationFormSchemaBase } from "@/src/features/blobstorage-integration/types";
import { validateAzureContainerName } from "@/src/features/blobstorage-integration/validation";
import { upsertBlobStorageIntegration } from "@/src/features/blobstorage-integration/service";
import { TRPCError } from "@trpc/server";
import {
@@ -56,7 +57,11 @@ export const blobStorageIntegrationRouter = createTRPCRouter({
}),
update: protectedProjectProcedure
.input(blobStorageIntegrationFormSchema.extend({ projectId: z.string() }))
.input(
blobStorageIntegrationFormSchemaBase
.extend({ projectId: z.string() })
.superRefine(validateAzureContainerName),
)
.mutation(async ({ input, ctx }) => {
try {
throwIfNoProjectAccess({
@@ -5,8 +5,9 @@ import {
BlobStorageExportMode,
AnalyticsIntegrationExportSource,
} from "@langfuse/shared";
import { validateAzureContainerName } from "@/src/features/blobstorage-integration/validation";
export const blobStorageIntegrationFormSchema = z.object({
export const blobStorageIntegrationFormSchemaBase = z.object({
type: z.enum(BlobStorageIntegrationType),
bucketName: z.string().min(1, { message: "Bucket name is required" }),
endpoint: z.string().url().optional().nullable(),
@@ -36,6 +37,9 @@ export const blobStorageIntegrationFormSchema = z.object({
compressed: z.boolean().default(true),
});
export const blobStorageIntegrationFormSchema =
blobStorageIntegrationFormSchemaBase.superRefine(validateAzureContainerName);
export type BlobStorageIntegrationFormSchema = z.infer<
typeof blobStorageIntegrationFormSchema
>;
@@ -0,0 +1,83 @@
import {
AZURE_CONTAINER_NAME_REGEX,
validateAzureContainerName,
} from "./validation";
import { z } from "zod";
describe("AZURE_CONTAINER_NAME_REGEX", () => {
const valid = [
"abc",
"my-container",
"a1b2c3",
"123",
"a-b",
"a".repeat(63),
"container-name-1",
];
const invalid = [
"ab", // too short
"a", // too short
"a".repeat(64), // too long
"ABC", // uppercase
"My-Container", // mixed case
"-abc", // starts with hyphen
"abc-", // ends with hyphen
"my--container", // consecutive hyphens
"has space", // spaces
"has.dot", // dots
"has/slash", // slashes
"Feedback N8N Bot", // the original issue
"", // empty
];
it.each(valid)("accepts valid name: %s", (name) => {
expect(AZURE_CONTAINER_NAME_REGEX.test(name)).toBe(true);
});
it.each(invalid)("rejects invalid name: %s", (name) => {
expect(AZURE_CONTAINER_NAME_REGEX.test(name)).toBe(false);
});
});
describe("validateAzureContainerName via schema", () => {
const schema = z
.object({ type: z.string(), bucketName: z.string() })
.superRefine(validateAzureContainerName);
it("rejects invalid Azure container name", () => {
const result = schema.safeParse({
type: "AZURE_BLOB_STORAGE",
bucketName: "Feedback N8N Bot",
});
expect(result.success).toBe(false);
if (!result.success) {
expect(result.error.issues[0].path).toEqual(["bucketName"]);
}
});
it("allows invalid container name for S3 type", () => {
const result = schema.safeParse({
type: "S3",
bucketName: "Feedback N8N Bot",
});
expect(result.success).toBe(true);
});
it("allows valid Azure container name", () => {
const result = schema.safeParse({
type: "AZURE_BLOB_STORAGE",
bucketName: "valid-container",
});
expect(result.success).toBe(true);
});
it("skips Azure validation when bucketName is empty", () => {
const result = schema.safeParse({
type: "AZURE_BLOB_STORAGE",
bucketName: "",
});
// Should pass superRefine (empty guard), letting .min(1) handle it upstream
expect(result.success).toBe(true);
});
});
@@ -0,0 +1,31 @@
import { z } from "zod";
/**
* Azure container names must be 3-63 characters, lowercase letters, numbers,
* and hyphens only. Must start and end with a letter or number. No consecutive
* hyphens.
*
* @see https://learn.microsoft.com/en-us/rest/api/storageservices/naming-and-referencing-containers--blobs--and-metadata#container-names
*/
export const AZURE_CONTAINER_NAME_REGEX =
/^[a-z0-9](?!.*--)[a-z0-9-]{1,61}[a-z0-9]$/;
export const AZURE_CONTAINER_NAME_ERROR =
"Azure container names must be 3-63 characters, lowercase letters, numbers, and hyphens only. Must start and end with a letter or number, no consecutive hyphens.";
export function validateAzureContainerName(
data: { type: string; bucketName: string },
ctx: z.RefinementCtx,
) {
if (!data.bucketName) return;
if (
data.type === "AZURE_BLOB_STORAGE" &&
!AZURE_CONTAINER_NAME_REGEX.test(data.bucketName)
) {
ctx.addIssue({
code: z.ZodIssueCode.custom,
message: AZURE_CONTAINER_NAME_ERROR,
path: ["bucketName"],
});
}
}
@@ -102,11 +102,20 @@ export function SelectDashboardDialog({
selectedDashboardId === d.id ? "bg-muted" : ""
}`}
>
<TableCell className="font-medium">{d.name}</TableCell>
<TableCell className="truncate" title={d.description}>
<TableCell
density="comfortable"
className="font-medium"
>
{d.name}
</TableCell>
<TableCell
density="comfortable"
className="truncate"
title={d.description}
>
{d.description}
</TableCell>
<TableCell>
<TableCell density="comfortable">
{new Date(d.updatedAt).toLocaleString()}
</TableCell>
</TableRow>
@@ -0,0 +1,22 @@
import { getScoreFilterConfig } from "./scores-config";
describe("getScoreFilterConfig", () => {
it("omits sidebar facets for hidden score columns", () => {
const config = getScoreFilterConfig([
"traceId",
"traceName",
"observationId",
"traceTags",
]);
expect(config.facets.map((facet) => facet.column)).not.toContain("traceId");
expect(config.facets.map((facet) => facet.column)).not.toContain(
"traceName",
);
expect(config.facets.map((facet) => facet.column)).not.toContain(
"observationId",
);
expect(config.facets.map((facet) => facet.column)).not.toContain("tags");
expect(config.facets.map((facet) => facet.column)).toContain("userId");
});
});
@@ -8,6 +8,20 @@ export const SCORE_COLUMN_TO_BACKEND_KEY: ColumnToBackendKeyMap = {
tags: "trace_tags",
};
export type ScoresTableHiddenColumn =
| "traceId"
| "traceName"
| "observationId"
| "jobConfigurationId"
| "userId"
| "traceTags";
const SCORES_HIDDEN_COLUMN_TO_FILTER_COLUMN: Partial<
Record<ScoresTableHiddenColumn, string>
> = {
traceTags: "tags",
};
export const scoreFilterConfig: FilterConfig = {
tableName: "scores",
@@ -83,3 +97,26 @@ export const scoreFilterConfig: FilterConfig = {
},
],
};
export function getScoreFilterConfig(
hiddenColumns: ScoresTableHiddenColumn[] = [],
): FilterConfig {
if (hiddenColumns.length === 0) {
return scoreFilterConfig;
}
const hiddenColumnSet = new Set<string>(
hiddenColumns.map(
(column) => SCORES_HIDDEN_COLUMN_TO_FILTER_COLUMN[column] ?? column,
),
);
return {
...scoreFilterConfig,
defaultExpanded: scoreFilterConfig.defaultExpanded?.filter(
(column) => !hiddenColumnSet.has(column),
),
facets: scoreFilterConfig.facets.filter(
(facet) => !hiddenColumnSet.has(facet.column),
),
};
}
@@ -117,7 +117,11 @@ export function ApiKeyList(props: { entityId: string; scope: ApiKeyScope }) {
<TableBody className="text-muted-foreground">
{apiKeysQuery.data?.length === 0 ? (
<TableRow>
<TableCell colSpan={5} className="text-center">
<TableCell
density="comfortable"
colSpan={5}
className="text-center"
>
None
</TableCell>
</TableRow>
@@ -127,29 +131,32 @@ export function ApiKeyList(props: { entityId: string; scope: ApiKeyScope }) {
key={apiKey.id}
className="hover:bg-primary-foreground"
>
<TableCell className="hidden md:table-cell">
<TableCell
density="comfortable"
className="hidden md:table-cell"
>
{apiKey.createdAt.toLocaleDateString()}
</TableCell>
<TableCell>
<TableCell density="comfortable">
<ApiKeyNote
apiKey={apiKey}
entityId={entityId}
scope={scope}
/>
</TableCell>
<TableCell className="font-mono">
<TableCell density="comfortable" className="font-mono">
<CodeView
className="inline-block text-xs"
content={apiKey.publicKey}
/>
</TableCell>
<TableCell className="font-mono">
<TableCell density="comfortable" className="font-mono">
{apiKey.displaySecretKey}
</TableCell>
{/* <TableCell>
{apiKey.lastUsedAt?.toLocaleDateString() ?? "Never"}
</TableCell> */}
<TableCell>
<TableCell density="comfortable">
<DeleteApiKeyButton
entityId={entityId}
apiKeyId={apiKey.id}
@@ -93,7 +93,11 @@ export function LlmApiKeyList(props: { projectId: string }) {
<TableBody className="text-muted-foreground">
{apiKeys.data?.data.length === 0 ? (
<TableRow>
<TableCell colSpan={6} className="text-center">
<TableCell
density="comfortable"
colSpan={6}
className="text-center"
>
None
</TableCell>
</TableRow>
@@ -104,18 +108,28 @@ export function LlmApiKeyList(props: { projectId: string }) {
className="hover:bg-primary-foreground cursor-default"
onClick={() => setEditingKeyId(apiKey.id)}
>
<TableCell className="font-mono">{apiKey.provider}</TableCell>
<TableCell className="font-mono">{apiKey.adapter}</TableCell>
<TableCell className="max-w-md overflow-auto font-mono">
<TableCell density="comfortable" className="font-mono">
{apiKey.provider}
</TableCell>
<TableCell density="comfortable" className="font-mono">
{apiKey.adapter}
</TableCell>
<TableCell
density="comfortable"
className="max-w-md overflow-auto font-mono"
>
{apiKey.baseURL ?? "default"}
</TableCell>
<TableCell className="font-mono">
<TableCell density="comfortable" className="font-mono">
{apiKey.displaySecretKey}
</TableCell>
{hasExtraHeaderKeys ? (
<TableCell> {apiKey.extraHeaderKeys.join(", ")} </TableCell>
<TableCell density="comfortable">
{" "}
{apiKey.extraHeaderKeys.join(", ")}{" "}
</TableCell>
) : null}
<TableCell className="text-right">
<TableCell density="comfortable" className="text-right">
<div
className="flex justify-end space-x-2"
onClick={(e) => e.stopPropagation()}
@@ -1,4 +1,5 @@
import { z } from "zod";
import { validateAzureContainerName } from "@/src/features/blobstorage-integration/validation";
/**
* Enums
@@ -26,7 +27,7 @@ export const CreateBlobStorageIntegrationRequest = z
.object({
projectId: z.string(),
type: BlobStorageIntegrationType,
bucketName: z.string(),
bucketName: z.string().min(1),
endpoint: z.string().nullable().optional(),
region: z.string(),
accessKeyId: z.string().nullable().optional(),
@@ -57,7 +58,8 @@ export const CreateBlobStorageIntegrationRequest = z
"exportStartDate is required when exportMode is FROM_CUSTOM_DATE",
path: ["exportStartDate"],
},
);
)
.superRefine(validateAzureContainerName);
export const BlobStorageIntegrationResponse = z
.object({
@@ -432,6 +432,7 @@ export function MembersTable({
onColumnVisibilityChange={setColumnVisibility}
columnOrder={columnOrder}
onColumnOrderChange={setColumnOrder}
cellPadding="comfortable"
/>
</SettingsTableCard>
) : (
@@ -464,6 +465,7 @@ export function MembersTable({
onColumnVisibilityChange={setColumnVisibility}
columnOrder={columnOrder}
onColumnOrderChange={setColumnOrder}
cellPadding="comfortable"
/>
)}
</>
@@ -119,19 +119,20 @@ export function SelectWidgetDialog({
selectedWidgetId === widget.id ? "bg-muted" : ""
}`}
>
<TableCell className="font-medium">
<TableCell density="comfortable" className="font-medium">
{widget.name}
</TableCell>
<TableCell
density="comfortable"
className="truncate"
title={widget.description}
>
{widget.description}
</TableCell>
<TableCell>
<TableCell density="comfortable">
{startCase(widget.view.toLowerCase())}
</TableCell>
<TableCell>
<TableCell density="comfortable">
{getChartTypeDisplayName(
widget.chartType as DashboardWidgetChartType,
)}
@@ -392,7 +392,7 @@ const BlobStorageIntegrationSettingsForm = ({
</FormControl>
<FormDescription>
{integrationType === "AZURE_BLOB_STORAGE"
? "The Azure storage container name"
? "Azure container name (3-63 chars, lowercase letters, numbers, and hyphens only)"
: "The S3 bucket name"}
</FormDescription>
<FormMessage />
@@ -180,7 +180,7 @@ function ScoresTab({ userId, projectId }: TabProps) {
<ScoresTable
projectId={projectId}
userId={userId}
omittedFilter={["User ID"]}
hiddenColumns={["userId"]}
/>
);
}
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "worker",
"version": "3.167.1",
"version": "3.167.2",
"description": "",
"license": "MIT",
"private": true,
@@ -78,7 +78,7 @@
"@types/pg": "^8.11.10",
"@types/uuid": "^9.0.8",
"@typescript/native-preview": "7.0.0-dev.20260122.3",
"@vitest/coverage-v8": "^4.1.2",
"@vitest/coverage-v8": "^4.1.4",
"eslint": "^9.39.2",
"msw": "^2.6.5",
"prettier": "^3.8.1",
@@ -86,7 +86,7 @@
"tsc-watch": "^6.2.0",
"tsx": "^4.20.5",
"typescript": "^5.7.2",
"vitest": "^4.1.2",
"vitest": "^4.1.4",
"wait-for-expect": "^3.0.2"
}
}
@@ -0,0 +1,163 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
const invokeMock = vi.fn();
const streamMock = vi.fn();
const chatVertexAIConstructorMock = vi.fn().mockImplementation(() => ({
invoke: invokeMock,
pipe: vi.fn().mockReturnValue({
stream: streamMock,
}),
}));
const VERTEXAI_USE_DEFAULT_CREDENTIALS = "__VERTEXAI_DEFAULT_CREDENTIALS__";
process.env.CLICKHOUSE_URL ??= "http://localhost:8123";
process.env.CLICKHOUSE_USER ??= "default";
process.env.CLICKHOUSE_PASSWORD ??= "password";
process.env.LANGFUSE_S3_EVENT_UPLOAD_BUCKET ??= "test-bucket";
process.env.ENCRYPTION_KEY ??=
"0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef";
class MockLLMCompletionError extends Error {
responseStatusCode: number;
isRetryable: boolean;
blockReason: null;
constructor(params: {
message: string;
responseStatusCode?: number;
isRetryable?: boolean;
}) {
super(params.message);
this.name = "LLMCompletionError";
this.responseStatusCode = params.responseStatusCode ?? 500;
this.isRetryable = params.isRetryable ?? false;
this.blockReason = null;
}
shouldBlockConfig() {
return false;
}
getEvaluatorBlockReason() {
return null;
}
}
describe("fetchLLMCompletion runtime timeouts", () => {
let originalTimeout: number;
let originalCloudRegion: string | undefined;
let env: typeof import("../../../packages/shared/src/env").env;
let encrypt: typeof import("../../../packages/shared/src/encryption").encrypt;
let fetchLLMCompletion: typeof import("../../../packages/shared/src/server/llm/fetchLLMCompletion").fetchLLMCompletion;
beforeEach(async () => {
vi.useFakeTimers();
invokeMock.mockReset();
streamMock.mockReset();
chatVertexAIConstructorMock.mockClear();
vi.resetModules();
originalCloudRegion = process.env.NEXT_PUBLIC_LANGFUSE_CLOUD_REGION;
delete process.env.NEXT_PUBLIC_LANGFUSE_CLOUD_REGION;
vi.doMock("@langchain/google-vertexai", () => ({
ChatVertexAI: chatVertexAIConstructorMock,
}));
vi.doMock("../../../packages/shared/src/server/llm/errors", () => ({
LLMCompletionError: MockLLMCompletionError,
}));
({ env } = await import("../../../packages/shared/src/env"));
({ encrypt } = await import("../../../packages/shared/src/encryption"));
({ fetchLLMCompletion } =
await import("../../../packages/shared/src/server/llm/fetchLLMCompletion"));
originalTimeout = env.LANGFUSE_FETCH_LLM_COMPLETION_TIMEOUT_MS;
});
afterEach(() => {
env.LANGFUSE_FETCH_LLM_COMPLETION_TIMEOUT_MS = originalTimeout;
if (originalCloudRegion === undefined) {
delete process.env.NEXT_PUBLIC_LANGFUSE_CLOUD_REGION;
} else {
process.env.NEXT_PUBLIC_LANGFUSE_CLOUD_REGION = originalCloudRegion;
}
vi.useRealTimers();
});
it("wraps non-streaming VertexAI timeouts as non-retryable LLMCompletionError", async () => {
env.LANGFUSE_FETCH_LLM_COMPLETION_TIMEOUT_MS = 25;
invokeMock.mockImplementation(() => new Promise(() => {}));
const completionPromise = fetchLLMCompletion({
streaming: false,
messages: [
{
role: "user",
content: "What is 2+2? Answer only with the number.",
type: "public-api-created",
},
],
modelParams: {
provider: "google-vertex-ai",
adapter: "google-vertex-ai",
model: "gemini-2.0-flash",
temperature: 0,
max_tokens: 10,
},
llmConnection: {
secretKey: encrypt(VERTEXAI_USE_DEFAULT_CREDENTIALS),
config: null,
},
});
const completionRejection = expect(completionPromise).rejects.toMatchObject(
{
name: "LLMCompletionError",
message: "Request timed out after 25ms",
isRetryable: false,
},
);
await vi.runOnlyPendingTimersAsync();
await completionRejection;
});
it("wraps streaming VertexAI timeouts as non-retryable LLMCompletionError", async () => {
env.LANGFUSE_FETCH_LLM_COMPLETION_TIMEOUT_MS = 25;
streamMock.mockImplementation(() => new Promise(() => {}));
const completionPromise = fetchLLMCompletion({
streaming: true,
messages: [
{
role: "user",
content: "Stream the answer.",
type: "public-api-created",
},
],
modelParams: {
provider: "google-vertex-ai",
adapter: "google-vertex-ai",
model: "gemini-2.0-flash",
temperature: 0,
max_tokens: 10,
},
llmConnection: {
secretKey: encrypt(VERTEXAI_USE_DEFAULT_CREDENTIALS),
config: null,
},
});
const completionRejection = expect(completionPromise).rejects.toMatchObject(
{
name: "LLMCompletionError",
message: "Request timed out after 25ms",
isRetryable: false,
},
);
await vi.runOnlyPendingTimersAsync();
await completionRejection;
});
});
@@ -0,0 +1,153 @@
import { describe, expect, it } from "vitest";
import {
LangfuseInternalTraceEnvironment,
prepareTracedEventsForIngestion,
} from "@langfuse/shared/src/server";
describe("prepareTracedEventsForIngestion", () => {
it("sets the configured environment on every forwarded event", () => {
const preparedEvents = prepareTracedEventsForIngestion(
[
{
type: "trace-create",
body: {
id: "trace-123",
name: "internal-trace",
environment: "wrong-environment",
},
},
{
type: "span-create",
body: {
id: "blocked-span",
traceId: "trace-123",
name: "RunnableLambda",
environment: "wrong-environment",
},
},
{
type: "span-update",
body: {
id: "blocked-span",
traceId: "trace-123",
},
},
{
type: "generation-create",
body: {
id: "generation-123",
traceId: "trace-123",
name: "ChatOpenAI",
},
},
{
type: "generation-update",
body: {
id: "generation-123",
traceId: "trace-123",
output: {
role: "assistant",
content: "Berlin",
},
},
},
{
type: "span-create",
body: {
id: "allowed-span",
traceId: "trace-123",
name: "Retriever",
environment: "another-environment",
},
},
{
type: "span-update",
body: {
id: "allowed-span",
traceId: "trace-123",
},
},
],
{
environment: LangfuseInternalTraceEnvironment.PromptExperiments,
},
);
expect(preparedEvents.map((event) => event.body.id)).not.toContain(
"blocked-span",
);
expect(preparedEvents).toHaveLength(5);
expect(
preparedEvents.every(
(event) =>
event.body.environment ===
LangfuseInternalTraceEnvironment.PromptExperiments,
),
).toBe(true);
});
it("adds prompt metadata only to generation-create events", () => {
const preparedEvents = prepareTracedEventsForIngestion(
[
{
type: "trace-create",
body: {
id: "trace-123",
name: "internal-trace",
},
},
{
type: "generation-create",
body: {
id: "generation-123",
traceId: "trace-123",
name: "ChatOpenAI",
},
},
{
type: "generation-update",
body: {
id: "generation-123",
traceId: "trace-123",
},
},
],
{
environment: LangfuseInternalTraceEnvironment.PromptExperiments,
prompt: {
name: "internal-prompt",
version: 3,
},
},
);
expect(preparedEvents).toEqual(
expect.arrayContaining([
expect.objectContaining({
type: "generation-create",
body: expect.objectContaining({
id: "generation-123",
environment: LangfuseInternalTraceEnvironment.PromptExperiments,
promptName: "internal-prompt",
promptVersion: 3,
}),
}),
expect.objectContaining({
type: "generation-update",
body: expect.not.objectContaining({
promptName: "internal-prompt",
promptVersion: 3,
}),
}),
expect.objectContaining({
type: "trace-create",
body: expect.not.objectContaining({
promptName: "internal-prompt",
promptVersion: 3,
}),
}),
]),
);
});
});
+1 -1
View File
@@ -1 +1 @@
export const VERSION = "v3.167.1";
export const VERSION = "v3.167.2";
@@ -0,0 +1,108 @@
import { beforeEach, describe, expect, it, vi, type Mock } from "vitest";
vi.mock("@langfuse/shared/src/server", () => ({
convertQueueNameToMetricName: vi.fn().mockImplementation((name) => name),
logger: {
info: vi.fn(),
warn: vi.fn(),
error: vi.fn(),
},
recordDistribution: vi.fn(),
}));
vi.mock("@langfuse/shared/src/db", () => ({
prisma: {
datasetRuns: {
findFirstOrThrow: vi.fn(),
},
jobExecution: {
findFirstOrThrow: vi.fn(),
},
},
}));
vi.mock("crypto", () => ({
randomUUID: vi.fn().mockReturnValue("retry-job-id"),
}));
import { prisma } from "@langfuse/shared/src/db";
import { logger, recordDistribution } from "@langfuse/shared/src/server";
import { retryLLMRateLimitError } from "./retry-handler";
describe("retryLLMRateLimitError", () => {
beforeEach(() => {
vi.clearAllMocks();
(prisma.jobExecution.findFirstOrThrow as Mock).mockResolvedValue({
createdAt: new Date(),
});
});
it("returns queue_unavailable instead of throwing when queue.add fails", async () => {
const add = vi.fn().mockRejectedValue(new Error("redis unavailable"));
const result = await retryLLMRateLimitError(
{
data: {
timestamp: new Date(),
payload: {
projectId: "project-id",
jobExecutionId: "job-execution-id",
},
},
},
{
table: "job_executions",
idField: "jobExecutionId",
queue: { add },
queueName: "llm-as-a-judge-execution-queue-1",
jobName: "llm-as-a-judge-execution-job",
delayFn: () => 30_000,
},
);
expect(result).toEqual({
outcome: "queue_unavailable",
});
expect(add).toHaveBeenCalledOnce();
expect(logger.warn).toHaveBeenCalledWith(
expect.stringContaining("Failed to enqueue retry job"),
expect.any(Error),
);
expect(recordDistribution).toHaveBeenCalledTimes(2);
});
it("returns queue_unavailable instead of throwing when age lookup fails", async () => {
(prisma.jobExecution.findFirstOrThrow as Mock).mockRejectedValue(
new Error("database unavailable"),
);
const result = await retryLLMRateLimitError(
{
data: {
timestamp: new Date(),
payload: {
projectId: "project-id",
jobExecutionId: "job-execution-id",
},
},
},
{
table: "job_executions",
idField: "jobExecutionId",
queue: { add: vi.fn() },
queueName: "llm-as-a-judge-execution-queue-1",
jobName: "llm-as-a-judge-execution-job",
delayFn: () => 30_000,
},
);
expect(result).toEqual({
outcome: "queue_unavailable",
});
expect(logger.error).toHaveBeenCalledWith(
expect.stringContaining("Failed to handle 429 retry"),
expect.any(Error),
);
expect(recordDistribution).not.toHaveBeenCalled();
});
});
+83 -41
View File
@@ -28,14 +28,23 @@ interface RetryConfig {
delayFn: (attempt: number) => number;
}
export type RetryScheduleResult =
| {
outcome: "scheduled";
delay: number;
retryBaggage: RetryBaggage;
}
| {
outcome: "skipped";
reason: "too_old";
}
| {
outcome: "queue_unavailable";
};
/**
* Handles rate limiting and retry logic for queue jobs
* Automatically retries jobs that fail with 429/5xx errors unless they're older than 24h
*
* @param error - The error that occurred
* @param job - The job that failed
* @param config - Retry configuration
* @returns true if retry was handled and job was added to the queue, false if regular processing should continue
*/
export async function retryLLMRateLimitError(
job: {
@@ -46,7 +55,7 @@ export async function retryLLMRateLimitError(
};
},
config: RetryConfig,
): Promise<void> {
): Promise<RetryScheduleResult> {
try {
const jobId = job.data.payload[config.idField];
@@ -67,65 +76,98 @@ export async function retryLLMRateLimitError(
`Job ${jobId} is rate limited for more than 24h. Stop retrying.`,
);
return; // Don't retry
return {
outcome: "skipped",
reason: "too_old",
};
}
// Retry the job with delay
const delay = config.delayFn((job.data.retryBaggage?.attempt ?? 0) + 1);
const retryBaggage: RetryBaggage | undefined = job.data.retryBaggage
const retryBaggage: RetryBaggage = job.data.retryBaggage
? {
originalJobTimestamp: new Date(
job.data.retryBaggage.originalJobTimestamp,
),
attempt: job.data.retryBaggage.attempt + 1,
}
: undefined;
: {
originalJobTimestamp: new Date(job.data.timestamp),
attempt: 1,
};
if (!config.queue) {
logger.warn(
`Retry queue ${config.queueName} is not available for job ${jobId}. Falling back to normal error handling.`,
);
return {
outcome: "queue_unavailable",
};
}
// Record retry attempt distribution per queue
if (retryBaggage) {
recordDistribution(
`${convertQueueNameToMetricName(config.queueName)}.retries`,
retryBaggage.attempt,
{
queue: config.queueName,
},
);
recordDistribution(
`${convertQueueNameToMetricName(config.queueName)}.retries`,
retryBaggage.attempt,
{
queue: config.queueName,
},
);
// Record delay distribution per queue
recordDistribution(
`${convertQueueNameToMetricName(config.queueName)}.total_retry_delay_ms`,
new Date().getTime() -
new Date(retryBaggage.originalJobTimestamp).getTime(), // this is the total delay
{
queue: config.queueName,
unit: "milliseconds",
},
);
}
// Record delay distribution per queue
recordDistribution(
`${convertQueueNameToMetricName(config.queueName)}.total_retry_delay_ms`,
new Date().getTime() -
new Date(retryBaggage.originalJobTimestamp).getTime(), // this is the total delay
{
queue: config.queueName,
unit: "milliseconds",
},
);
logger.info(
`Job ${jobId} is rate limited. Retrying in ${delay}ms. Attempt: ${retryBaggage?.attempt}. Total delay: ${retryBaggage ? new Date().getTime() - new Date(retryBaggage?.originalJobTimestamp).getTime() : "unavailable"}ms.`,
);
await config.queue?.add(
config.queueName,
{
name: config.jobName,
id: randomUUID(),
timestamp: new Date(),
payload: job.data.payload,
retryBaggage: retryBaggage,
},
{ delay },
);
try {
await config.queue.add(
config.queueName,
{
name: config.jobName,
id: randomUUID(),
timestamp: new Date(),
payload: job.data.payload,
retryBaggage: retryBaggage,
},
{ delay },
);
} catch (addErr) {
logger.warn(
`Failed to enqueue retry job for ${jobId}. Falling back to normal error handling.`,
addErr,
);
return {
outcome: "queue_unavailable",
};
}
return {
outcome: "scheduled",
delay,
retryBaggage,
};
} catch (innerErr) {
const jobId = job.data.payload[config.idField];
logger.error(
`Failed to handle 429 retry for ${jobId}. Continuing regular processing.`,
`Failed to handle 429 retry for ${jobId}. Falling back to caller error handling.`,
innerErr,
);
throw innerErr;
return {
outcome: "queue_unavailable",
};
}
}
@@ -0,0 +1,90 @@
import { beforeEach, describe, expect, it, vi, type Mock } from "vitest";
import { Job } from "bullmq";
vi.mock("@langfuse/shared/src/server", () => ({
ExperimentCreateQueue: {
getInstance: vi.fn().mockReturnValue({
add: vi.fn(),
}),
},
QueueJobs: {
ExperimentCreateJob: "experiment-create-job",
},
QueueName: {
ExperimentCreate: "experiment-create-queue",
},
isLLMCompletionError: vi.fn(),
logger: {
error: vi.fn(),
},
traceException: vi.fn(),
}));
vi.mock("../../features/utils", () => ({
retryLLMRateLimitError: vi.fn(),
}));
vi.mock("../../features/experiments/experimentServiceClickhouse", () => ({
createExperimentJobClickhouse: vi.fn(),
}));
vi.mock("../../errors/UnrecoverableError", async () => {
const actual = await vi.importActual("../../errors/UnrecoverableError");
return {
...actual,
isUnrecoverableError: vi.fn(),
};
});
import { isLLMCompletionError } from "@langfuse/shared/src/server";
import { createExperimentJobClickhouse } from "../../features/experiments/experimentServiceClickhouse";
import { retryLLMRateLimitError } from "../../features/utils";
import { experimentCreateQueueProcessor } from "../experimentQueue";
import { isUnrecoverableError } from "../../errors/UnrecoverableError";
describe("experimentCreateQueueProcessor", () => {
const createMockJob = (): Job<any> =>
({
data: {
payload: {
projectId: "project-id",
runId: "run-id",
},
},
}) as Job<any>;
beforeEach(() => {
vi.clearAllMocks();
(isLLMCompletionError as Mock).mockReturnValue(false);
(isUnrecoverableError as Mock).mockReturnValue(false);
});
it("rethrows retryable LLM errors when the retry queue is unavailable", async () => {
const llmError = new Error("Rate limit exceeded");
(llmError as Error & { isRetryable: boolean }).isRetryable = true;
(createExperimentJobClickhouse as Mock).mockRejectedValue(llmError);
(isLLMCompletionError as Mock).mockReturnValue(true);
(retryLLMRateLimitError as Mock).mockResolvedValue({
outcome: "queue_unavailable",
});
await expect(
experimentCreateQueueProcessor(createMockJob()),
).rejects.toThrow("Rate limit exceeded");
expect(retryLLMRateLimitError).toHaveBeenCalledWith(
expect.objectContaining({
data: expect.objectContaining({
payload: expect.objectContaining({
projectId: "project-id",
runId: "run-id",
}),
}),
}),
expect.objectContaining({
table: "dataset_runs",
idField: "runId",
}),
);
});
});
@@ -1,10 +1,38 @@
import { describe, it, expect, vi, beforeEach, type Mock } from "vitest";
import {
beforeAll,
beforeEach,
describe,
expect,
it,
vi,
type Mock,
} from "vitest";
import { Job } from "bullmq";
import { JobExecutionStatus } from "@prisma/client";
import { llmAsJudgeExecutionQueueProcessorBuilder } from "../evalQueue";
import { QueueName, type TQueueJobTypes } from "@langfuse/shared/src/server";
import { UnrecoverableError } from "../../errors/UnrecoverableError";
const QueueName = {
LLMAsJudgeExecution: "llm-as-a-judge-execution-queue",
EvaluationExecutionSecondaryQueue: "evaluation-execution-secondary-queue",
} as const;
const QueueJobs = {
LLMAsJudgeExecution: "llm-as-a-judge-execution-job",
EvaluationExecution: "evaluation-execution-job",
} as const;
const JobExecutionStatus = {
DELAYED: "DELAYED",
ERROR: "ERROR",
} as const;
vi.mock("@langfuse/shared", () => ({
removeEmptyEnvVariables: <T>(value: T) => value,
JobExecutionStatus: {
DELAYED: "DELAYED",
ERROR: "ERROR",
},
}));
// Mock prisma
vi.mock("@langfuse/shared/src/db", () => ({
prisma: {
@@ -20,10 +48,20 @@ vi.mock("../../features/evaluation/observationEval", () => ({
}));
// Mock logger and span
vi.mock("@langfuse/shared/src/server", async () => {
const actual = await vi.importActual("@langfuse/shared/src/server");
vi.mock("@langfuse/shared/src/server", () => {
const getQueueInstance = vi.fn().mockReturnValue({
add: vi.fn(),
});
return {
...actual,
QueueName: {
LLMAsJudgeExecution: "llm-as-a-judge-execution-queue",
EvaluationExecutionSecondaryQueue: "evaluation-execution-secondary-queue",
},
QueueJobs: {
LLMAsJudgeExecution: "llm-as-a-judge-execution-job",
EvaluationExecution: "evaluation-execution-job",
},
logger: {
debug: vi.fn(),
info: vi.fn(),
@@ -34,10 +72,14 @@ vi.mock("@langfuse/shared/src/server", async () => {
getCurrentSpan: vi.fn().mockReturnValue({
setAttribute: vi.fn(),
}),
EvalExecutionQueue: {
getInstance: vi.fn(),
},
SecondaryEvalExecutionQueue: {
getInstance: vi.fn(),
},
LLMAsJudgeExecutionQueue: {
getInstance: vi.fn().mockReturnValue({
add: vi.fn(),
}),
getInstance: getQueueInstance,
},
isLLMCompletionError: vi.fn(),
};
@@ -73,12 +115,18 @@ describe("llmAsJudgeExecutionQueueProcessor", () => {
const jobExecutionId = "job-exec-456";
const observationS3Path = "evals/test/observation.json";
const queueName = `${QueueName.LLMAsJudgeExecution}-1`;
const llmAsJudgeExecutionQueueProcessor =
llmAsJudgeExecutionQueueProcessorBuilder(queueName);
let llmAsJudgeExecutionQueueProcessor: (
job: Job<{
payload: {
projectId: string;
jobExecutionId: string;
observationS3Path: string;
};
retryBaggage?: { attempt: number };
}>,
) => Promise<unknown>;
const createMockJob = (
overrides: Partial<TQueueJobTypes[QueueName.LLMAsJudgeExecution]> = {},
): Job<TQueueJobTypes[QueueName.LLMAsJudgeExecution]> => {
const createMockJob = (overrides: Record<string, unknown> = {}): Job<any> => {
return {
data: {
id: "queue-job-123",
@@ -92,9 +140,16 @@ describe("llmAsJudgeExecutionQueueProcessor", () => {
retryBaggage: { attempt: 0 },
...overrides,
},
} as unknown as Job<TQueueJobTypes[QueueName.LLMAsJudgeExecution]>;
} as Job<any>;
};
beforeAll(async () => {
const { llmAsJudgeExecutionQueueProcessorBuilder } =
await import("../evalQueue");
llmAsJudgeExecutionQueueProcessor =
llmAsJudgeExecutionQueueProcessorBuilder(queueName);
});
beforeEach(() => {
vi.clearAllMocks();
(isLLMCompletionError as Mock).mockReturnValue(false);
@@ -146,6 +201,9 @@ describe("llmAsJudgeExecutionQueueProcessor", () => {
// Mark as retryable
(rateLimitError as unknown as { isRetryable: boolean }).isRetryable =
true;
(retryLLMRateLimitError as Mock).mockResolvedValue({
outcome: "scheduled",
});
const job = createMockJob();
await llmAsJudgeExecutionQueueProcessor(job);
@@ -180,6 +238,9 @@ describe("llmAsJudgeExecutionQueueProcessor", () => {
true;
(processObservationEval as Mock).mockRejectedValue(rateLimitError);
(isLLMCompletionError as Mock).mockReturnValue(true);
(retryLLMRateLimitError as Mock).mockResolvedValue({
outcome: "scheduled",
});
const job = createMockJob();
@@ -188,6 +249,61 @@ describe("llmAsJudgeExecutionQueueProcessor", () => {
llmAsJudgeExecutionQueueProcessor(job),
).resolves.not.toThrow();
});
it("should set ERROR when retryable LLM errors are not re-enqueued", async () => {
const rateLimitError = new Error("Rate limit exceeded");
(processObservationEval as Mock).mockRejectedValue(rateLimitError);
(isLLMCompletionError as Mock).mockReturnValue(true);
(rateLimitError as unknown as { isRetryable: boolean }).isRetryable =
true;
(retryLLMRateLimitError as Mock).mockResolvedValue({
outcome: "skipped",
reason: "too_old",
});
const job = createMockJob();
await llmAsJudgeExecutionQueueProcessor(job);
expect(prisma.jobExecution.update).toHaveBeenCalledWith({
where: {
id: jobExecutionId,
projectId,
},
data: expect.objectContaining({
status: JobExecutionStatus.ERROR,
endTime: expect.any(Date),
error: "Rate limit exceeded",
executionTraceId: "test-trace-id",
}),
});
});
it("should set ERROR when the retry queue is unavailable", async () => {
const rateLimitError = new Error("Rate limit exceeded");
(processObservationEval as Mock).mockRejectedValue(rateLimitError);
(isLLMCompletionError as Mock).mockReturnValue(true);
(rateLimitError as unknown as { isRetryable: boolean }).isRetryable =
true;
(retryLLMRateLimitError as Mock).mockResolvedValue({
outcome: "queue_unavailable",
});
const job = createMockJob();
await llmAsJudgeExecutionQueueProcessor(job);
expect(prisma.jobExecution.update).toHaveBeenCalledWith({
where: {
id: jobExecutionId,
projectId,
},
data: expect.objectContaining({
status: JobExecutionStatus.ERROR,
endTime: expect.any(Date),
error: "Rate limit exceeded",
executionTraceId: "test-trace-id",
}),
});
});
});
describe("LLM completion errors (non-retryable)", () => {
@@ -356,7 +472,7 @@ describe("llmAsJudgeExecutionQueueProcessor", () => {
const job = createMockJob({
retryBaggage: { attempt: 3 },
});
await llmAsJudgeExecutionQueueProcessor(job);
await llmAsJudgeExecutionQueueProcessor(job as Job<any>);
expect(mockSpan.setAttribute).toHaveBeenCalledWith(
"messaging.bullmq.job.input.retryBaggage.attempt",
@@ -373,7 +489,7 @@ describe("llmAsJudgeExecutionQueueProcessor", () => {
const job = createMockJob();
delete (job.data as { retryBaggage?: unknown }).retryBaggage;
await llmAsJudgeExecutionQueueProcessor(job);
await llmAsJudgeExecutionQueueProcessor(job as Job<any>);
expect(mockSpan.setAttribute).toHaveBeenCalledWith(
"messaging.bullmq.job.input.retryBaggage.attempt",
+30 -26
View File
@@ -211,7 +211,7 @@ export const evalJobExecutorQueueProcessorBuilder = (
? SecondaryEvalExecutionQueue.getInstance({ shardName: queueName })
: EvalExecutionQueue.getInstance({ shardName: queueName });
await retryLLMRateLimitError(job, {
const retryResult = await retryLLMRateLimitError(job, {
table: "job_executions",
idField: "jobExecutionId",
queue,
@@ -220,20 +220,22 @@ export const evalJobExecutorQueueProcessorBuilder = (
delayFn: delayInMs,
});
// Use the deterministic execution trace ID to update the job execution
await prisma.jobExecution.update({
where: {
id: job.data.payload.jobExecutionId,
projectId: job.data.payload.projectId,
},
data: {
status: JobExecutionStatus.DELAYED,
executionTraceId,
},
});
if (retryResult.outcome === "scheduled") {
// Use the deterministic execution trace ID to update the job execution
await prisma.jobExecution.update({
where: {
id: job.data.payload.jobExecutionId,
projectId: job.data.payload.projectId,
},
data: {
status: JobExecutionStatus.DELAYED,
executionTraceId,
},
});
// Return early as we have already scheduled a delayed retry
return;
// Return early as we have already scheduled a delayed retry
return;
}
}
// At this point there will be only 4xx LLMCompletionErrors that are not retryable and application errors
@@ -305,7 +307,7 @@ export const llmAsJudgeExecutionQueueProcessorBuilder =
const queue = LLMAsJudgeExecutionQueue.getInstance({
shardName: queueName,
});
await retryLLMRateLimitError(job, {
const retryResult = await retryLLMRateLimitError(job, {
table: "job_executions",
idField: "jobExecutionId",
queue,
@@ -314,18 +316,20 @@ export const llmAsJudgeExecutionQueueProcessorBuilder =
delayFn: delayInMs,
});
await prisma.jobExecution.update({
where: {
id: job.data.payload.jobExecutionId,
projectId: job.data.payload.projectId,
},
data: {
status: JobExecutionStatus.DELAYED,
executionTraceId,
},
});
if (retryResult.outcome === "scheduled") {
await prisma.jobExecution.update({
where: {
id: job.data.payload.jobExecutionId,
projectId: job.data.payload.projectId,
},
data: {
status: JobExecutionStatus.DELAYED,
executionTraceId,
},
});
return;
return;
}
}
await prisma.jobExecution.update({
+3 -2
View File
@@ -23,7 +23,7 @@ export const experimentCreateQueueProcessor = async (
return true;
} catch (e) {
if (isLLMCompletionError(e) && e.isRetryable) {
await retryLLMRateLimitError(job, {
const retryResult = await retryLLMRateLimitError(job, {
table: "dataset_runs",
idField: "runId",
queue: ExperimentCreateQueue.getInstance(),
@@ -32,7 +32,8 @@ export const experimentCreateQueueProcessor = async (
delayFn: delayInMs,
});
return;
if (retryResult.outcome === "scheduled") return;
if (retryResult.outcome === "queue_unavailable") throw e;
}
if (isLLMCompletionError(e) || isUnrecoverableError(e)) return;