Skip to content
Merged
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
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 @@ -3,7 +3,7 @@

import { Run } from "@langchain/core/tracers/base";
import { Span } from "@opentelemetry/api";
import { OpenTelemetryConstants } from "@microsoft/agents-a365-observability";
import { OpenTelemetryConstants, truncateValue } from "@microsoft/agents-a365-observability";

// Type guards
export function isString(value: unknown): value is string {
Expand Down Expand Up @@ -51,8 +51,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 +77,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 +214,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 +263,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,18 +62,32 @@ 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) {
if (isTracingSuppressed(context.active())) {
// Even when suppressed, end any span that was started before suppression kicked in
// to avoid abandoned spans that are never exported or closed.
const suppressedEntry = this.runs.get(run.id);
if (suppressedEntry) {
suppressedEntry.span.end(run.end_time ?? undefined);
}
this.parentByRunId.delete(run.id);
this.runs.delete(run.id);
return;
}
// Skip internal runs
Expand All @@ -80,7 +97,7 @@ export class LangChainTracer extends BaseTracer {
return;
}

const entry = this.runs[run.id];
const entry = this.runs.get(run.id);
if (!entry) {
return;
}
Expand All @@ -91,7 +108,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 +117,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 +146,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