Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -129,6 +129,7 @@ This library is LibreChat's tracing surface: every agent run it orchestrates is
| -------------------------------------------------------------- | ---------------------------------------------------------------------------- |
| `src/langfuseTraceShaping.ts` | Export-time span rename/retype/drop rules (the shape itself) |
| `src/langfuseToolOutputTracing.ts` | Span processor applying shaping hooks + tool-output redaction |
| `src/langfuseInlineMedia.ts` | Export-time replacement of inline base64 media with type/size descriptors |
| `src/langfuse.ts` | Callback handler, identity/tags/metadata, control-flow + usage normalization |
| `src/instrumentation.ts` | Tracer provider bootstrap, per-tenant routing, deterministic trace ids |
| `src/langfuseConfig.ts` | Config resolution/merging (env vs run vs agent level) |
Expand All @@ -144,6 +145,7 @@ These originated from direct Langfuse-team feedback (PRs #288, #316) and must su
- **Root-observation input/output are the conversation, not the state.** Root input reduces to the user's question, output to the assistant's answer — never full serialized graph state. Tool-dispatch input is scoped to the pending tool calls. Do not emit deprecated trace-level input/output attributes.
- **Control flow is not an error.** `GraphInterrupt` and `ParentCommand` end their traces as successful with `controlFlow` outputs, not as error traces.
- **Usage/cost is accurate per provider.** e.g. Bedrock cache read/write tokens are folded into input tokens so Langfuse cost math is right.
- **Inline media is described, not exported.** Every span that serializes the conversation would otherwise carry its own copy of each attached file, so base64 payloads (raw, `data:` URIs while Langfuse media upload is off, serialized byte buffers) become a type-and-size descriptor unless a host opts in via `inlineMediaTracing`.
- **Redaction is honored everywhere tool output can surface** — tool spans, and any generation input that embeds tool results (e.g. the activity-label prompt).
- **Identity and metadata always propagate**: `userId`, `sessionId`, tags, environment, and trace metadata (`messageId`, `parentMessageId`, `agentId`, `agentName`) — including across LangChain callbacks that fire outside the caller's OTEL context.
- **Trace identity is self-contained.** Root observations never inherit trace ids or parents from foreign ambient OTEL spans (e.g. a host's HTTP auto-instrumentation): the callback handler detaches them so roots stay true roots, deterministic ids apply, and concurrent runs inside one request context (an agent run plus a title run) cannot merge into one trace. Spans created through the Langfuse tracer provider are honored as parents, so hosts can still group runs under their own Langfuse observations deliberately.
Expand Down
6 changes: 4 additions & 2 deletions src/instrumentation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -79,8 +79,9 @@ export function ensureOpenTelemetryContextManager(): void {
}

/** Cache key for processor instances: the destination plus processor-level
* policy (`toolOutputTracing` is baked into each processor's redaction
* behavior), unlike the pure destination identity used for span parenting. */
* policy (`toolOutputTracing` and `inlineMediaTracing` are baked into each
* processor's export behavior), unlike the pure destination identity used
* for span parenting. */
function getLangfuseProcessorCacheKey(
destinationKey: string,
langfuse?: t.LangfuseConfig
Expand All @@ -89,6 +90,7 @@ function getLangfuseProcessorCacheKey(
destinationKey,
mediaUploadEnabled: langfuse?.mediaUploadEnabled,
toolOutputTracing: langfuse?.toolOutputTracing,
inlineMediaTracing: langfuse?.inlineMediaTracing,
});
}

Expand Down
36 changes: 36 additions & 0 deletions src/langfuseConfig.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import type { LangfuseSpanProcessorParams } from '@langfuse/otel';
import type { ResolvedLangfuseToolOutputTracingConfig } from '@/langfuseRuntimeContext';
import type * as t from '@/types';
import { parseBooleanEnv } from '@/utils/misc';
Expand Down Expand Up @@ -175,6 +176,32 @@ export function resolveToolOutputTracingConfig(
};
}

export function resolveInlineMediaTracingEnabled(
runLangfuse?: t.LangfuseConfig,
agentLangfuse?: t.LangfuseConfig
): boolean {
return (
agentLangfuse?.inlineMediaTracing?.enabled ??
runLangfuse?.inlineMediaTracing?.enabled ??
parseBooleanEnv(process.env.LANGFUSE_TRACE_INLINE_MEDIA) ??
false
);
}

/**
* Mirrors `LangfuseSpanProcessor`'s own default, so `data:` URIs are left for
* the SDK exactly when it will upload them and replace them with a reference.
*/
export function resolveLangfuseMediaUploadEnabled(
params?: LangfuseSpanProcessorParams
): boolean {
if (params?.mediaUploadEnabled != null) {
return params.mediaUploadEnabled;
}
const env = process.env.LANGFUSE_MEDIA_UPLOAD_ENABLED;
return !isPresent(env) || !['false', '0'].includes(env.toLowerCase());
}

/**
* Merges header maps case-insensitively, keeping the override's casing.
*
Expand Down Expand Up @@ -239,6 +266,14 @@ export function resolveLangfuseConfig(
...agentLangfuse.toolOutputTracing,
}
: undefined;
const inlineMediaTracing =
runLangfuse.inlineMediaTracing != null ||
agentLangfuse.inlineMediaTracing != null
? {
...runLangfuse.inlineMediaTracing,
...agentLangfuse.inlineMediaTracing,
}
: undefined;
const metadata =
runLangfuse.metadata != null || agentLangfuse.metadata != null
? {
Expand Down Expand Up @@ -277,5 +312,6 @@ export function resolveLangfuseConfig(
...(tags != null ? { tags } : {}),
...(toolNodeTracing != null ? { toolNodeTracing } : {}),
...(toolOutputTracing != null ? { toolOutputTracing } : {}),
...(inlineMediaTracing != null ? { inlineMediaTracing } : {}),
};
}
209 changes: 209 additions & 0 deletions src/langfuseInlineMedia.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,209 @@
import { LangfuseOtelSpanAttributes } from '@langfuse/tracing';
import type { ReadableSpan } from '@opentelemetry/sdk-trace-base';

/** Payloads shorter than this cost less to export than a descriptor saves. */
const MIN_INLINE_MEDIA_LENGTH = 4096;
const MIN_INLINE_MEDIA_BYTES = (MIN_INLINE_MEDIA_LENGTH * 3) / 4;
const BASE64_SAMPLE_LENGTH = 256;
const BASE64_BODY_PATTERN = /^[A-Za-z0-9+/]+$/;
const BASE64_TAIL_PATTERN = /^[A-Za-z0-9+/]*={0,2}$/;
const DATA_URI_PATTERN = /^data:([^;,]{1,127});base64,/;
const MIME_TYPE_KEYS = ['mimeType', 'mime_type', 'media_type', 'mediaType'];
const INLINE_MEDIA_ATTRIBUTES = [
LangfuseOtelSpanAttributes.OBSERVATION_INPUT,
LangfuseOtelSpanAttributes.OBSERVATION_OUTPUT,
];

type JsonValue =
| string
| number
| boolean
| null
| JsonValue[]
| { [key: string]: JsonValue };

type JsonObject = { [key: string]: JsonValue };

type SerializedBuffer = { type: 'Buffer'; data: number[] };

export type LangfuseInlineMediaPolicy = {
/** Also omit `data:` URIs. False while the Langfuse SDK uploads them as media. */
omitDataUris: boolean;
};

function isJsonObject(value: JsonValue): value is JsonObject {
return value != null && typeof value === 'object' && !Array.isArray(value);
}

function isSerializedBuffer(value: JsonValue): value is SerializedBuffer {
return (
isJsonObject(value) &&
value.type === 'Buffer' &&
Array.isArray(value.data) &&
value.data.length >= MIN_INLINE_MEDIA_BYTES &&
typeof value.data[0] === 'number'
);
}

/**
* Samples the head and tail instead of testing the whole string: a full scan
* of a multi-megabyte payload costs ~100ms per span on the export path.
*/
function isBase64Payload(value: string): boolean {
return (
BASE64_BODY_PATTERN.test(value.slice(0, BASE64_SAMPLE_LENGTH)) &&
BASE64_TAIL_PATTERN.test(value.slice(-BASE64_SAMPLE_LENGTH))
);
}

function getBase64ByteLength(value: string, payloadLength: number): number {
let padding = 0;
if (value.endsWith('==')) {
padding = 2;
} else if (value.endsWith('=')) {
padding = 1;
}
return Math.floor((payloadLength * 3) / 4) - padding;
}

function findMimeType(parent?: JsonObject): string | undefined {
if (parent == null) {
return undefined;
}
for (const key of MIME_TYPE_KEYS) {
const mimeType = parent[key];
if (typeof mimeType === 'string' && mimeType !== '') {
return mimeType;
}
}
return undefined;
}

function describeOmittedMedia(bytes: number, mimeType?: string): string {
return `[inline media omitted from trace: ${mimeType ?? 'unknown type'}, ${bytes} bytes]`;
}

function describeInlineString(
value: string,
policy: LangfuseInlineMediaPolicy,
parent?: JsonObject
): string | undefined {
if (value.length < MIN_INLINE_MEDIA_LENGTH) {
return undefined;
}

const dataUri = DATA_URI_PATTERN.exec(value);
if (dataUri != null) {
return policy.omitDataUris
? describeOmittedMedia(
getBase64ByteLength(value, value.length - dataUri[0].length),
dataUri[1]
)
: undefined;
}

return isBase64Payload(value)
? describeOmittedMedia(
getBase64ByteLength(value, value.length),
findMimeType(parent)
)
: undefined;
}

function describeInlineValue(
value: JsonValue,
policy: LangfuseInlineMediaPolicy,
parent?: JsonObject
): string | undefined {
if (typeof value === 'string') {
return describeInlineString(value, policy, parent);
}
if (isSerializedBuffer(value)) {
return describeOmittedMedia(value.data.length, findMimeType(parent));
}
return undefined;
}

/**
* Replaces inline media inside a freshly parsed JSON value. Mutates in place:
* the value is a private parse of the span attribute, and copying a tree that
* holds tens of megabytes would double the cost this exists to remove.
*/
function omitInlineMediaInPlace(
value: JsonValue,
policy: LangfuseInlineMediaPolicy
): boolean {
if (Array.isArray(value)) {
let changed = false;
for (let i = 0; i < value.length; i++) {
const descriptor = describeInlineValue(value[i], policy);
if (descriptor != null) {
value[i] = descriptor;
changed = true;
} else if (omitInlineMediaInPlace(value[i], policy)) {
changed = true;
}
}
return changed;
}

if (!isJsonObject(value)) {
return false;
}

let changed = false;
for (const key in value) {
const descriptor = describeInlineValue(value[key], policy, value);
if (descriptor != null) {
value[key] = descriptor;
changed = true;
} else if (omitInlineMediaInPlace(value[key], policy)) {
changed = true;
}
}
return changed;
}

function omitInlineMediaFromAttribute(
value: string,
policy: LangfuseInlineMediaPolicy
): string | undefined {
if (value.length < MIN_INLINE_MEDIA_LENGTH) {
return undefined;
}
if (value[0] !== '{' && value[0] !== '[') {
return describeInlineString(value, policy);
}

let parsed: JsonValue;
try {
parsed = JSON.parse(value) as JsonValue;
} catch {
return undefined;
}
return omitInlineMediaInPlace(parsed, policy)
? JSON.stringify(parsed)
: undefined;
}

/**
* Replaces inline base64 media in a span's input and output with a short
* type-and-size descriptor. Every span that serializes the conversation (the
* agent, its model node, the prompt, and the generation) otherwise carries its
* own full copy of each attached file.
*/
export function omitLangfuseSpanInlineMedia(
span: ReadableSpan,
policy: LangfuseInlineMediaPolicy
): void {
for (const key of INLINE_MEDIA_ATTRIBUTES) {
const value = span.attributes[key];
if (typeof value !== 'string') {
continue;
}
const next = omitInlineMediaFromAttribute(value, policy);
if (next != null) {
span.attributes[key] = next;
}
}
}
29 changes: 25 additions & 4 deletions src/langfuseToolOutputTracing.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,19 +8,23 @@ import type {
import type { LangfuseSpanProcessorParams } from '@langfuse/otel';
import type { Context } from '@opentelemetry/api';
import type { ResolvedLangfuseToolOutputTracingConfig } from '@/langfuseRuntimeContext';
import type { LangfuseInlineMediaPolicy } from '@/langfuseInlineMedia';
import type * as t from '@/types';
import {
LANGFUSE_TOOL_OUTPUT_REDACTION_TEXT,
hasToolOutputTracingConfig,
normalizeToolName,
resolveInlineMediaTracingEnabled,
resolveLangfuseConfig,
resolveLangfuseMediaUploadEnabled,
resolveToolOutputTracingConfig,
} from '@/langfuseConfig';
import {
shapeLangfuseSpan,
shouldDropLangfuseSpan,
} from '@/langfuseTraceShaping';
import { resolveToolOutputTracingConfigForSpan } from '@/langfuseRuntimeScope';
import { omitLangfuseSpanInlineMedia } from '@/langfuseInlineMedia';

export { LANGFUSE_TOOL_OUTPUT_REDACTION_TEXT, resolveLangfuseConfig };

Expand Down Expand Up @@ -877,29 +881,36 @@ export function redactLangfuseSpanToolOutputs(

export function prepareLangfuseSpanForExport(
span: ReadableSpan,
config?: ResolvedLangfuseToolOutputTracingConfig
config?: ResolvedLangfuseToolOutputTracingConfig,
inlineMediaPolicy?: LangfuseInlineMediaPolicy
): void {
classifyLangfuseToolNodeSpan(span);
if (config != null) {
redactLangfuseSpanToolOutputs(span, config);
}
if (inlineMediaPolicy != null) {
omitLangfuseSpanInlineMedia(span, inlineMediaPolicy);
}
shapeLangfuseSpan(span);
}

class ToolOutputRedactingLangfuseSpanProcessor implements SpanProcessor {
private readonly processor: LangfuseSpanProcessor;
private readonly fallbackConfig?: ResolvedLangfuseToolOutputTracingConfig;
private readonly inlineMediaPolicy?: LangfuseInlineMediaPolicy;
private readonly spanConfigs = new WeakMap<
object,
ResolvedLangfuseToolOutputTracingConfig
>();

constructor(
params?: LangfuseSpanProcessorParams,
fallbackConfig?: ResolvedLangfuseToolOutputTracingConfig
fallbackConfig?: ResolvedLangfuseToolOutputTracingConfig,
inlineMediaPolicy?: LangfuseInlineMediaPolicy
) {
this.processor = new LangfuseSpanProcessor(params);
this.fallbackConfig = fallbackConfig;
this.inlineMediaPolicy = inlineMediaPolicy;
}

onStart(span: Span, parentContext: Context): void {
Expand All @@ -920,7 +931,7 @@ class ToolOutputRedactingLangfuseSpanProcessor implements SpanProcessor {
return;
}
const config = this.spanConfigs.get(span) ?? this.fallbackConfig;
prepareLangfuseSpanForExport(span, config);
prepareLangfuseSpanForExport(span, config, this.inlineMediaPolicy);
this.processor.onEnd(span);
}

Expand All @@ -941,7 +952,17 @@ export function createLangfuseSpanProcessor(
const fallbackConfig = hasToolOutputTracingConfig(runLangfuse, agentLangfuse)
? resolveToolOutputTracingConfig(runLangfuse, agentLangfuse)
: undefined;
return new ToolOutputRedactingLangfuseSpanProcessor(params, fallbackConfig);
const inlineMediaPolicy = resolveInlineMediaTracingEnabled(
runLangfuse,
agentLangfuse
)
? undefined
: { omitDataUris: !resolveLangfuseMediaUploadEnabled(params) };
return new ToolOutputRedactingLangfuseSpanProcessor(
params,
fallbackConfig,
inlineMediaPolicy
);
}

function hasLangfuseEnvKeys(): boolean {
Expand Down
Loading
Loading