Skip to content
Merged
Show file tree
Hide file tree
Changes from 9 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
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Added

- **`OpenAIAgentsInstrumentationConfig.isContentRecordingEnabled`** — Optional `boolean` to enable content recording in OpenAI trace processor.
- **`LangChainTraceInstrumentor.instrument(module, options?)`** — New optional `{ isContentRecordingEnabled?: boolean }` parameter to enable content recording in LangChain tracer.
Comment thread
fpfp100 marked this conversation as resolved.
- **`truncateValue`** / **`MAX_ATTRIBUTE_LENGTH`** — Exported utilities for attribute value truncation (8192 char limit).

Comment thread
fpfp100 marked this conversation as resolved.
### Breaking Changes (`@microsoft/agents-a365-observability-hosting`)

- **`ScopeUtils.deriveAgentDetails(turnContext, authToken)`** — New required `authToken: string` parameter.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,8 +20,9 @@ class LangChainTraceInstrumentorImpl extends InstrumentationBase<Instrumentation
private _hasBeenEnabled = false;
private _isPatched = false;
protected otelTracer: Tracer;
private isContentRecordingEnabled: boolean;

private constructor() {
private constructor(options?: { isContentRecordingEnabled?: boolean }) {
if (LangChainTraceInstrumentorImpl._instance !== null) {
throw new Error("LangChainTraceInstrumentor can only be instantiated once.");
}
Expand All @@ -41,14 +42,15 @@ class LangChainTraceInstrumentorImpl extends InstrumentationBase<Instrumentation
"agent365-langchain",
"1.0.0"
);
this.isContentRecordingEnabled = options?.isContentRecordingEnabled ?? false;

LangChainTraceInstrumentorImpl._instance = this;
logger.info("[LangChainTraceInstrumentor] Initialized and automatically enabled");
}

static getInstance(): LangChainTraceInstrumentorImpl {
static getInstance(options?: { isContentRecordingEnabled?: boolean }): LangChainTraceInstrumentorImpl {
if (!LangChainTraceInstrumentorImpl._instance) {
LangChainTraceInstrumentorImpl._instance = new LangChainTraceInstrumentorImpl();
LangChainTraceInstrumentorImpl._instance = new LangChainTraceInstrumentorImpl(options);
Comment thread
fpfp100 marked this conversation as resolved.
}
return LangChainTraceInstrumentorImpl._instance;
}
Expand Down Expand Up @@ -102,7 +104,7 @@ class LangChainTraceInstrumentorImpl extends InstrumentationBase<Instrumentation
this: CallbackManagerModuleType,
...args: Parameters<typeof CallbackManager["_configureSync"]>
) {
args[0] = addTracerToHandlers(instrumentor.otelTracer, args[0]);
args[0] = addTracerToHandlers(instrumentor.otelTracer, args[0], { isContentRecordingEnabled: instrumentor.isContentRecordingEnabled });
logger.info("[LangChainTraceInstrumentor] _configureSync wrapped to add LangChainTracer");
return original.apply(this, args);
};
Expand Down Expand Up @@ -150,10 +152,12 @@ export class LangChainTraceInstrumentor {
}

/**
* Initialize and auto-instrument for LangChain
* Initialize and auto-instrument for LangChain
* @param module The CallbackManager module to instrument
* @param options Optional configuration options
*/
static instrument(module: CallbackManagerModuleType): void {
LangChainTraceInstrumentorImpl.getInstance().manuallyInstrumentImpl(module);
static instrument(module: CallbackManagerModuleType, options?: { isContentRecordingEnabled?: boolean }): void {
LangChainTraceInstrumentorImpl.getInstance(options).manuallyInstrumentImpl(module);
}

/**
Expand Down Expand Up @@ -186,21 +190,22 @@ export class LangChainTraceInstrumentor {

export function addTracerToHandlers(
tracer: Tracer,
handlers: CallbackManagerModule.Callbacks | undefined
handlers: CallbackManagerModule.Callbacks | undefined,
options?: { isContentRecordingEnabled?: boolean }
): CallbackManagerModule.Callbacks {
if (handlers == null) {
return [new LangChainTracer(tracer)];
return [new LangChainTracer(tracer, options)];
}

if (Array.isArray(handlers)) {
if (!handlers.some((h) => h instanceof LangChainTracer)) {
handlers.push(new LangChainTracer(tracer));
handlers.push(new LangChainTracer(tracer, options));
}
return handlers;
}

if (!handlers.inheritableHandlers.some((h) => h instanceof LangChainTracer)) {
handlers.addHandler(new LangChainTracer(tracer), true);
handlers.addHandler(new LangChainTracer(tracer, options), true);
}
return handlers;
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,15 @@ import { Run } from "@langchain/core/tracers/base";
import { Span } from "@opentelemetry/api";
import { OpenTelemetryConstants } from "@microsoft/agents-a365-observability";

const MAX_ATTRIBUTE_LENGTH = 8_192;

function truncateValue(value: string): string {
if (value.length > MAX_ATTRIBUTE_LENGTH) {
return value.substring(0, MAX_ATTRIBUTE_LENGTH) + '...[truncated]';
Comment thread
fpfp100 marked this conversation as resolved.
Outdated
}
Comment thread
fpfp100 marked this conversation as resolved.
Outdated
return value;
}

// Type guards
export function isString(value: unknown): value is string {
return typeof value === "string";
Expand Down Expand Up @@ -51,8 +60,8 @@ export function setToolAttributes(run: Run, span: Span) {
if (isString(run.name)) {
span.setAttribute(OpenTelemetryConstants.GEN_AI_TOOL_NAME_KEY, run.name);
}
if (run.inputs) span.setAttribute(OpenTelemetryConstants.GEN_AI_TOOL_ARGS_KEY, JSON.stringify(run.inputs?.input ?? run.inputs));
if (run.outputs?.output?.kwargs?.content) span.setAttribute(OpenTelemetryConstants.GEN_AI_TOOL_CALL_RESULT_KEY, JSON.stringify(run.outputs?.output?.kwargs?.content));
if (run.inputs) span.setAttribute(OpenTelemetryConstants.GEN_AI_TOOL_ARGS_KEY, truncateValue(JSON.stringify(run.inputs?.input ?? run.inputs)));
if (run.outputs?.output?.kwargs?.content) span.setAttribute(OpenTelemetryConstants.GEN_AI_TOOL_CALL_RESULT_KEY, truncateValue(JSON.stringify(run.outputs?.output?.kwargs?.content)));
span.setAttribute(OpenTelemetryConstants.GEN_AI_TOOL_TYPE_KEY, "extension");
if (run.outputs?.output?.tool_call_id) span.setAttribute(OpenTelemetryConstants.GEN_AI_TOOL_CALL_ID_KEY, run.outputs?.output?.tool_call_id);
}
Expand All @@ -77,7 +86,7 @@ export function setInputMessagesAttribute(run: Run, span: Span) {
.filter(Boolean);

if (processed.length > 0) {
span.setAttribute(OpenTelemetryConstants.GEN_AI_INPUT_MESSAGES_KEY, JSON.stringify(processed));
span.setAttribute(OpenTelemetryConstants.GEN_AI_INPUT_MESSAGES_KEY, truncateValue(JSON.stringify(processed)));
}
}

Expand Down Expand Up @@ -214,7 +223,7 @@ export function setOutputMessagesAttribute(run: Run, span: Span) {
}

if (messages.length > 0) {
span.setAttribute(OpenTelemetryConstants.GEN_AI_OUTPUT_MESSAGES_KEY, JSON.stringify(messages));
span.setAttribute(OpenTelemetryConstants.GEN_AI_OUTPUT_MESSAGES_KEY, truncateValue(JSON.stringify(messages)));
}
}

Expand Down Expand Up @@ -263,15 +272,15 @@ export function setSystemInstructionsAttribute(run: Run, span: Span) {
}

const prompts = Array.isArray(inputs.prompts) ? inputs.prompts.map(p => String(p ?? "").trim()).filter(Boolean).join("\n") : "";
if (prompts) return span.setAttribute(OpenTelemetryConstants.GEN_AI_SYSTEM_INSTRUCTIONS_KEY, prompts);
if (prompts) return span.setAttribute(OpenTelemetryConstants.GEN_AI_SYSTEM_INSTRUCTIONS_KEY, truncateValue(prompts));

const messages = Array.isArray(inputs.messages) ? inputs.messages : [];
const systemText = messages
.filter((m: Record<string, unknown>) => m.lc_type === "system")
.map((m: Record<string, unknown>) => String((m.lc_kwargs as Record<string, unknown> | undefined)?.content ?? "").trim())
.filter(Boolean)
.join("\n");
if (systemText) span.setAttribute(OpenTelemetryConstants.GEN_AI_SYSTEM_INSTRUCTIONS_KEY, systemText);
if (systemText) span.setAttribute(OpenTelemetryConstants.GEN_AI_SYSTEM_INSTRUCTIONS_KEY, truncateValue(systemText));
}

// Tokens (input and output)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,20 +4,23 @@
import { context, trace, Span, SpanKind, SpanStatusCode, Tracer } from "@opentelemetry/api";
import { BaseTracer, Run } from "@langchain/core/tracers/base";
import { isTracingSuppressed } from "@opentelemetry/core";
import { logger, OpenTelemetryConstants } from "@microsoft/agents-a365-observability";
import { logger, OpenTelemetryConstants, truncateValue } from "@microsoft/agents-a365-observability";
Comment thread
fpfp100 marked this conversation as resolved.
import * as Utils from "./Utils";

type RunWithSpan = { run: Run; span: Span; startTime: number; lastAccessTime: number };

export class LangChainTracer extends BaseTracer {
private static readonly MAX_RUNS = 10_000;
private tracer: Tracer;
private runs: Record<string, RunWithSpan> = {};
private parentByRunId: Record<string, string | undefined> = {};
private isContentRecordingEnabled: boolean;
private runs = new Map<string, RunWithSpan>();
private parentByRunId = new Map<string, string | undefined>();


constructor(tracer: Tracer) {
constructor(tracer: Tracer, options?: { isContentRecordingEnabled?: boolean }) {
super();
this.tracer = tracer;
this.isContentRecordingEnabled = options?.isContentRecordingEnabled ?? false;
}

name = "OpenTelemetryLangChainTracer";
Expand All @@ -27,7 +30,7 @@ export class LangChainTracer extends BaseTracer {
}

async onRunCreate(run: Run) {
this.parentByRunId[run.id] = run.parent_run_id;
this.parentByRunId.set(run.id, run.parent_run_id);
if (super.onRunCreate) await super.onRunCreate(run);
this.startTracing(run);
}
Comment thread
fpfp100 marked this conversation as resolved.
Expand Down Expand Up @@ -59,14 +62,20 @@ export class LangChainTracer extends BaseTracer {
spanName = `${operation} ${Utils.getModel(run) || run.name}`.trim();
}

if (this.runs.size >= LangChainTracer.MAX_RUNS) {
logger.warn(`[LangChainTracer] Max runs (${LangChainTracer.MAX_RUNS}) reached, skipping span`);
this.parentByRunId.delete(run.id);
return;
}
Comment thread
fpfp100 marked this conversation as resolved.

const startTime = run.start_time ?? Date.now();
const span = this.tracer.startSpan(spanName, {
kind: SpanKind.INTERNAL,
startTime,
attributes: { [OpenTelemetryConstants.GEN_AI_SYSTEM_KEY]: "langchain" },
}, activeContext);

this.runs[run.id] = { run, span, startTime, lastAccessTime: startTime };
this.runs.set(run.id, { run, span, startTime, lastAccessTime: startTime });
}

protected async _endTrace(run: Run) {
Expand All @@ -77,11 +86,13 @@ export class LangChainTracer extends BaseTracer {
const operation = Utils.getOperationType(run);
Comment thread
fpfp100 marked this conversation as resolved.
if (run.tags?.includes("langsmith:hidden") || run.name?.startsWith("Branch") || operation === "unknown") {
logger.info(`Skipping internal run: ${run.name} (parent: ${run.parent_run_id})`);
this.parentByRunId.delete(run.id);
return;
}

const entry = this.runs[run.id];
const entry = this.runs.get(run.id);
if (!entry) {
this.parentByRunId.delete(run.id);
return;
}

Expand All @@ -91,7 +102,7 @@ export class LangChainTracer extends BaseTracer {

if (run.error) {
span.setStatus({ code: SpanStatusCode.ERROR });
span.setAttribute(OpenTelemetryConstants.ERROR_MESSAGE_KEY, String(run.error));
span.setAttribute(OpenTelemetryConstants.ERROR_MESSAGE_KEY, truncateValue(String(run.error)));

Comment thread
fpfp100 marked this conversation as resolved.
} else {
span.setStatus({ code: SpanStatusCode.OK });
Expand All @@ -100,22 +111,27 @@ export class LangChainTracer extends BaseTracer {
// Set all attributes
Utils.setOperationTypeAttribute(operation, span);
Utils.setAgentAttributes(run, span);
Utils.setToolAttributes(run, span);
Utils.setInputMessagesAttribute(run, span);
Utils.setOutputMessagesAttribute(run, span);
Utils.setSystemInstructionsAttribute(run, span);
Utils.setModelAttribute(run, span);
Utils.setProviderNameAttribute(run, span);
Utils.setSessionIdAttribute(run, span);
Utils.setTokenAttributes(run, span);

// Content attributes gated by content recording setting
const contentRecording = this.isContentRecordingEnabled;
if (contentRecording) {
Utils.setToolAttributes(run, span);
Utils.setInputMessagesAttribute(run, span);
Utils.setOutputMessagesAttribute(run, span);
Utils.setSystemInstructionsAttribute(run, span);
}
Comment thread
fpfp100 marked this conversation as resolved.
Comment thread
fpfp100 marked this conversation as resolved.

} catch (error) {
logger.error(`[LangChainTracer] Error setting span attributes for run ${run.name}: ${error instanceof Error ? error.message : String(error)}`);
span.setStatus({ code: SpanStatusCode.ERROR });
} finally {
span.end(run.end_time ?? undefined);
delete this.runs[run.id];
delete this.parentByRunId[run.id];
this.runs.delete(run.id);
this.parentByRunId.delete(run.id);
await super._endTrace(run);
}
}
Expand All @@ -124,9 +140,9 @@ export class LangChainTracer extends BaseTracer {
let pid = run.parent_run_id;

while (pid) {
const entry = this.runs[pid];
const entry = this.runs.get(pid);
if (entry) return entry.span.spanContext();
pid = this.parentByRunId[pid];
pid = this.parentByRunId.get(pid);
}
return undefined;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,8 +49,7 @@
"@microsoft/agents-a365-runtime": "workspace:*",
"@openai/agents": "catalog:",
"@opentelemetry/api": "catalog:",
"@opentelemetry/instrumentation": "catalog:",
"hono": "catalog:"
"@opentelemetry/instrumentation": "catalog:"
},
"devDependencies": {
"@eslint/js": "catalog:",
Expand All @@ -66,7 +65,13 @@
"typescript-eslint": "catalog:"
},
"peerDependencies": {
"@openai/agents": "catalog:"
"@openai/agents": "catalog:",
"hono": "catalog:"
},
"peerDependenciesMeta": {
"hono": {
"optional": true
}
},
"engines": {
"node": ">=18.0.0"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,11 @@ export interface OpenAIAgentsInstrumentationConfig extends InstrumentationConfig
* Defaults to false.
*/
suppressInvokeAgentInput?: boolean;
/**
* Whether to enable content recording (input/output messages, tool args, etc.).
* @default false
*/
isContentRecordingEnabled?: boolean;
}

/**
Expand Down Expand Up @@ -100,7 +105,8 @@ export class OpenAIAgentsTraceInstrumentor extends InstrumentationBase<OpenAIAge
trace.getTracerProvider();

this.processor = new OpenAIAgentsTraceProcessor(agent365Tracer, {
suppressInvokeAgentInput: this._config.suppressInvokeAgentInput ?? false
suppressInvokeAgentInput: this._config.suppressInvokeAgentInput ?? false,
isContentRecordingEnabled: this._config.isContentRecordingEnabled ?? false,
});

// Register the processor directly using the imported setTraceProcessors function
Expand Down
Loading
Loading