diff --git a/CHANGELOG.md b/CHANGELOG.md index 872e875d..4a1a9a8a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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. +- **`truncateValue`** / **`MAX_ATTRIBUTE_LENGTH`** — Exported utilities for attribute value truncation (8192 char limit). + ### Breaking Changes (`@microsoft/agents-a365-observability-hosting`) - **`ScopeUtils.deriveAgentDetails(turnContext, authToken)`** — New required `authToken: string` parameter. diff --git a/packages/agents-a365-observability-extensions-langchain/src/LangChainTraceInstrumentor.ts b/packages/agents-a365-observability-extensions-langchain/src/LangChainTraceInstrumentor.ts index eac853f7..2fedc96b 100644 --- a/packages/agents-a365-observability-extensions-langchain/src/LangChainTraceInstrumentor.ts +++ b/packages/agents-a365-observability-extensions-langchain/src/LangChainTraceInstrumentor.ts @@ -20,8 +20,9 @@ class LangChainTraceInstrumentorImpl extends InstrumentationBase ) { - 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); }; @@ -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); } /** @@ -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; } diff --git a/packages/agents-a365-observability-extensions-langchain/src/Utils.ts b/packages/agents-a365-observability-extensions-langchain/src/Utils.ts index b96247e0..1f14db1a 100644 --- a/packages/agents-a365-observability-extensions-langchain/src/Utils.ts +++ b/packages/agents-a365-observability-extensions-langchain/src/Utils.ts @@ -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 { @@ -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); } @@ -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))); } } @@ -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))); } } @@ -263,7 +263,7 @@ 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 @@ -271,7 +271,7 @@ export function setSystemInstructionsAttribute(run: Run, span: Span) { .map((m: Record) => String((m.lc_kwargs as Record | 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) diff --git a/packages/agents-a365-observability-extensions-langchain/src/tracer.ts b/packages/agents-a365-observability-extensions-langchain/src/tracer.ts index 2fdb64d8..3d6a1ca2 100644 --- a/packages/agents-a365-observability-extensions-langchain/src/tracer.ts +++ b/packages/agents-a365-observability-extensions-langchain/src/tracer.ts @@ -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"; 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 = {}; - private parentByRunId: Record = {}; + private isContentRecordingEnabled: boolean; + private runs = new Map(); + private parentByRunId = new Map(); - constructor(tracer: Tracer) { + constructor(tracer: Tracer, options?: { isContentRecordingEnabled?: boolean }) { super(); this.tracer = tracer; + this.isContentRecordingEnabled = options?.isContentRecordingEnabled ?? false; } name = "OpenTelemetryLangChainTracer"; @@ -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); } @@ -59,6 +62,12 @@ 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; + } + const startTime = run.start_time ?? Date.now(); const span = this.tracer.startSpan(spanName, { kind: SpanKind.INTERNAL, @@ -66,11 +75,19 @@ export class LangChainTracer extends BaseTracer { 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 @@ -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; } @@ -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))); } else { span.setStatus({ code: SpanStatusCode.OK }); @@ -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); + } + } 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); } } @@ -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; } diff --git a/packages/agents-a365-observability-extensions-openai/package.json b/packages/agents-a365-observability-extensions-openai/package.json index 41c7d05a..c72aee19 100644 --- a/packages/agents-a365-observability-extensions-openai/package.json +++ b/packages/agents-a365-observability-extensions-openai/package.json @@ -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:", @@ -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" diff --git a/packages/agents-a365-observability-extensions-openai/src/OpenAIAgentsTraceInstrumentor.ts b/packages/agents-a365-observability-extensions-openai/src/OpenAIAgentsTraceInstrumentor.ts index 27a7693b..0fe4e2bf 100644 --- a/packages/agents-a365-observability-extensions-openai/src/OpenAIAgentsTraceInstrumentor.ts +++ b/packages/agents-a365-observability-extensions-openai/src/OpenAIAgentsTraceInstrumentor.ts @@ -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; } /** @@ -100,7 +105,8 @@ export class OpenAIAgentsTraceInstrumentor extends InstrumentationBase = new Map(); private readonly otelSpans: Map = new Map(); private readonly tokens: Map = new Map(); @@ -48,9 +50,23 @@ export class OpenAIAgentsTraceProcessor implements TracingProcessor { ['generation' + Constants.GEN_AI_REQUEST_CONTENT_KEY, OpenTelemetryConstants.GEN_AI_INPUT_MESSAGES_KEY], ]); - constructor(tracer: OtelTracer, options?: { suppressInvokeAgentInput?: boolean }) { + constructor(tracer: OtelTracer, options?: { suppressInvokeAgentInput?: boolean; isContentRecordingEnabled?: boolean }) { this.tracer = tracer; this.suppressInvokeAgentInput = options?.suppressInvokeAgentInput ?? false; + this.isContentRecordingEnabled = options?.isContentRecordingEnabled ?? false; + } + + private static readonly CONTENT_KEYS = new Set([ + OpenTelemetryConstants.GEN_AI_INPUT_MESSAGES_KEY, + OpenTelemetryConstants.GEN_AI_OUTPUT_MESSAGES_KEY, + OpenTelemetryConstants.GEN_AI_EVENT_CONTENT, + OpenTelemetryConstants.GEN_AI_TOOL_ARGS_KEY, + Constants.GEN_AI_REQUEST_CONTENT_KEY, + Constants.GEN_AI_RESPONSE_CONTENT_KEY, + ]); + + private static isContentKey(key: string): boolean { + return OpenAIAgentsTraceProcessor.CONTENT_KEYS.has(key); } private getNewKey(spanType: string, key: string): string | null { @@ -96,6 +112,11 @@ export class OpenAIAgentsTraceProcessor implements TracingProcessor { return; } + if (this.otelSpans.size >= OpenAIAgentsTraceProcessor.MAX_SPANS_IN_FLIGHT) { + logger.warn(`[OpenAIAgentsTraceProcessor] Max spans in flight (${OpenAIAgentsTraceProcessor.MAX_SPANS_IN_FLIGHT}) reached, skipping span`); + return; + } + const startTime = new Date(startedAt).getTime(); // Find parent span @@ -187,18 +208,19 @@ export class OpenAIAgentsTraceProcessor implements TracingProcessor { */ private processSpanData(otelSpan: OtelSpan, data: SpanData, traceId: string): void { const type = data.type; + const contentRecording = this.isContentRecordingEnabled; switch (type) { case 'response': - this.processResponseSpanData(otelSpan, data); + this.processResponseSpanData(otelSpan, data, contentRecording); break; case 'generation': - this.processGenerationSpanData(otelSpan, data, traceId); + this.processGenerationSpanData(otelSpan, data, traceId, contentRecording); break; case 'function': - this.processFunctionSpanData(otelSpan, data, traceId); + this.processFunctionSpanData(otelSpan, data, traceId, contentRecording); break; case 'mcp_tools': @@ -218,7 +240,7 @@ export class OpenAIAgentsTraceProcessor implements TracingProcessor { /** * Process response span data */ - private processResponseSpanData(otelSpan: OtelSpan, data: SpanData): void { + private processResponseSpanData(otelSpan: OtelSpan, data: SpanData, contentRecording: boolean): void { const responseData = data as Record; // Handle both formats: _response/_input (actual format) and response/input (legacy format) const responseObj = responseData._response || responseData.response; @@ -227,13 +249,13 @@ export class OpenAIAgentsTraceProcessor implements TracingProcessor { const resp = responseObj as Record; // Store the output field for GEN_AI_RESPONSE_CONTENT_KEY - if (resp.output) { + if (resp.output && contentRecording) { if (typeof resp.output === 'string') { - otelSpan.setAttribute(OpenTelemetryConstants.GEN_AI_OUTPUT_MESSAGES_KEY, resp.output); + otelSpan.setAttribute(OpenTelemetryConstants.GEN_AI_OUTPUT_MESSAGES_KEY, truncateValue(resp.output)); } else { otelSpan.setAttribute( OpenTelemetryConstants.GEN_AI_OUTPUT_MESSAGES_KEY, - this.buildOutputMessages(resp.output as Array<{ role: string; content: Array<{ type: string; text: string }> }>) + truncateValue(this.buildOutputMessages(resp.output as Array<{ role: string; content: Array<{ type: string; text: string }> }>)) ); } } @@ -251,26 +273,26 @@ export class OpenAIAgentsTraceProcessor implements TracingProcessor { otelSpan.updateName(`${InferenceOperationType.CHAT} ${modelName}`); } - if (inputObj && !this.suppressInvokeAgentInput) { + if (inputObj && !this.suppressInvokeAgentInput && contentRecording) { if (typeof inputObj === 'string') { try { const parsed = JSON.parse(inputObj as string); if (Array.isArray(parsed)) { otelSpan.setAttribute( OpenTelemetryConstants.GEN_AI_INPUT_MESSAGES_KEY, - this.buildInputMessages(parsed) + truncateValue(this.buildInputMessages(parsed)) ); return; } } catch { // If parsing fails, fall back to raw string behavior } - otelSpan.setAttribute(OpenTelemetryConstants.GEN_AI_INPUT_MESSAGES_KEY, inputObj); + otelSpan.setAttribute(OpenTelemetryConstants.GEN_AI_INPUT_MESSAGES_KEY, truncateValue(inputObj as string)); } else if (Array.isArray(inputObj)) { // build the input messages from array otelSpan.setAttribute( OpenTelemetryConstants.GEN_AI_INPUT_MESSAGES_KEY, - this.buildInputMessages(inputObj) + truncateValue(this.buildInputMessages(inputObj)) ); } } @@ -305,15 +327,18 @@ export class OpenAIAgentsTraceProcessor implements TracingProcessor { /** * Process generation span data */ - private processGenerationSpanData(otelSpan: OtelSpan, data: SpanData, traceId: string): void { + private processGenerationSpanData(otelSpan: OtelSpan, data: SpanData, traceId: string, contentRecording: boolean): void { const attrs = Utils.getAttributesFromGenerationSpanData(data); Object.entries(attrs).forEach(([key, value]) => { const shouldExcludeKey = key === OpenTelemetryConstants.GEN_AI_EXECUTION_TYPE_KEY || key === Constants.GEN_AI_EXECUTION_PAYLOAD_KEY; if (value !== null && value !== undefined && !shouldExcludeKey) { const newKey = this.getNewKey(data.type, key); - if (newKey !== OpenTelemetryConstants.GEN_AI_INPUT_MESSAGES_KEY || !this.suppressInvokeAgentInput) { - otelSpan.setAttribute(newKey || key, value as string | number | boolean); + const resolvedKey = newKey || key; + if (resolvedKey !== OpenTelemetryConstants.GEN_AI_INPUT_MESSAGES_KEY || !this.suppressInvokeAgentInput) { + if (!OpenAIAgentsTraceProcessor.isContentKey(resolvedKey) || contentRecording) { + otelSpan.setAttribute(resolvedKey, value as string | number | boolean); + } } } }); @@ -331,13 +356,16 @@ export class OpenAIAgentsTraceProcessor implements TracingProcessor { /** * Process function/tool span data */ - private processFunctionSpanData(otelSpan: OtelSpan, data: SpanData, traceId: string): void { + private processFunctionSpanData(otelSpan: OtelSpan, data: SpanData, traceId: string, contentRecording: boolean): void { const functionData = data as Record; const attrs = Utils.getAttributesFromFunctionSpanData(data); Object.entries(attrs).forEach(([key, value]) => { if (value !== null && value !== undefined && key !== OpenTelemetryConstants.GEN_AI_EXECUTION_TYPE_KEY) { const newKey = this.getNewKey(data.type, key); - otelSpan.setAttribute(newKey || key, value as string | number | boolean); + const resolvedKey = newKey || key; + if (!OpenAIAgentsTraceProcessor.isContentKey(resolvedKey) || contentRecording) { + otelSpan.setAttribute(resolvedKey, value as string | number | boolean); + } } otelSpan.setAttribute(OpenTelemetryConstants.GEN_AI_TOOL_TYPE_KEY, 'function'); }); diff --git a/packages/agents-a365-observability-extensions-openai/src/Utils.ts b/packages/agents-a365-observability-extensions-openai/src/Utils.ts index 54331ea7..761b711b 100644 --- a/packages/agents-a365-observability-extensions-openai/src/Utils.ts +++ b/packages/agents-a365-observability-extensions-openai/src/Utils.ts @@ -3,7 +3,7 @@ // ------------------------------------------------------------------------------ import { SpanStatusCode } from '@opentelemetry/api'; -import { OpenTelemetryConstants } from '@microsoft/agents-a365-observability'; +import { OpenTelemetryConstants, truncateValue } from '@microsoft/agents-a365-observability'; import * as Constants from './Constants'; import { Span as AgentsSpan, SpanData } from '@openai/agents-core/dist/tracing/spans'; @@ -14,9 +14,9 @@ import { Span as AgentsSpan, SpanData } from '@openai/agents-core/dist/tracing/s */ export function safeJsonDumps(obj: unknown): string { try { - return JSON.stringify(obj); + return truncateValue(JSON.stringify(obj)); } catch { - return String(obj); + return truncateValue(String(obj)); } } @@ -127,12 +127,12 @@ export function getAttributesFromFunctionSpanData(data: SpanData): Record(); private readonly _defaultRefreshSkewMs = 60_000; private readonly _defaultMaxTokenAgeMs = 3_600_000; + private readonly _maxCacheSize = 10_000; + private readonly _maxExpSeconds = 86_400; // 24 hours private readonly _keyLocks = new Map>(); private readonly _configProvider: IConfigurationProvider; @@ -85,6 +87,12 @@ export class AgenticTokenCache { return; } entry = { scopes: effectiveScopes }; + if (this._map.size >= this._maxCacheSize) { + const oldest = this._map.keys().next().value; + if (oldest !== undefined) { + this._map.delete(oldest); + } + } this._map.set(key, entry); } if (!Array.isArray(entry.scopes) || entry.scopes.length === 0) { @@ -157,7 +165,9 @@ export class AgenticTokenCache { const payloadSegment = parts[1]; const padded = payloadSegment + '='.repeat((4 - (payloadSegment.length % 4)) % 4); const json = JSON.parse(Buffer.from(padded, 'base64').toString('utf8')) as { exp?: unknown }; - return typeof json.exp === 'number' ? json.exp : undefined; + if (typeof json.exp !== 'number') return undefined; + const maxExp = Math.floor(Date.now() / 1000) + this._maxExpSeconds; + return Math.min(json.exp, maxExp); } catch { return undefined; } diff --git a/packages/agents-a365-observability-hosting/src/utils/ScopeUtils.ts b/packages/agents-a365-observability-hosting/src/utils/ScopeUtils.ts index ff6520d3..3e4c3835 100644 --- a/packages/agents-a365-observability-hosting/src/utils/ScopeUtils.ts +++ b/packages/agents-a365-observability-hosting/src/utils/ScopeUtils.ts @@ -25,7 +25,10 @@ import { resolveEmbodiedAgentIds } from './TurnContextUtils'; export class ScopeUtils { - private static setInputMessageTags(scope: InvokeAgentScope | InferenceScope, turnContext: TurnContext): InvokeAgentScope | InferenceScope { + private static setInputMessageTags( + scope: InvokeAgentScope | InferenceScope, + turnContext: TurnContext, + ): InvokeAgentScope | InferenceScope { if (turnContext?.activity?.text) { scope.recordInputMessages([turnContext.activity.text]); } diff --git a/packages/agents-a365-observability/src/configuration/ObservabilityConfiguration.ts b/packages/agents-a365-observability/src/configuration/ObservabilityConfiguration.ts index 2e6f143e..95756448 100644 --- a/packages/agents-a365-observability/src/configuration/ObservabilityConfiguration.ts +++ b/packages/agents-a365-observability/src/configuration/ObservabilityConfiguration.ts @@ -61,4 +61,5 @@ export class ObservabilityConfiguration extends RuntimeConfiguration { ?? process.env.A365_OBSERVABILITY_LOG_LEVEL ?? 'none'; } + } diff --git a/packages/agents-a365-observability/src/index.ts b/packages/agents-a365-observability/src/index.ts index b45f86d7..3ff38e27 100644 --- a/packages/agents-a365-observability/src/index.ts +++ b/packages/agents-a365-observability/src/index.ts @@ -55,6 +55,7 @@ export { InferenceScope } from './tracing/scopes/InferenceScope'; export { OutputScope } from './tracing/scopes/OutputScope'; export { logger, setLogger, getLogger, resetLogger, formatError } from './utils/logging'; export type { ILogger } from './utils/logging'; +export { truncateValue, MAX_ATTRIBUTE_LENGTH } from './tracing/util'; // Exporter utilities export { isPerRequestExportEnabled } from './tracing/exporter/utils'; diff --git a/packages/agents-a365-observability/src/tracing/exporter/Agent365Exporter.ts b/packages/agents-a365-observability/src/tracing/exporter/Agent365Exporter.ts index b5b8df75..937468e0 100644 --- a/packages/agents-a365-observability/src/tracing/exporter/Agent365Exporter.ts +++ b/packages/agents-a365-observability/src/tracing/exporter/Agent365Exporter.ts @@ -171,8 +171,8 @@ export class Agent365Exporter implements SpanExporter { // Select endpoint path based on S2S flag (includes tenantId in path) const endpointRelativePath = this.options.useS2SEndpoint - ? `/observabilityService/tenants/${tenantId}/agents/${agentId}/traces` - : `/observability/tenants/${tenantId}/agents/${agentId}/traces`; + ? `/observabilityService/tenants/${encodeURIComponent(tenantId)}/agents/${encodeURIComponent(agentId)}/traces` + : `/observability/tenants/${encodeURIComponent(tenantId)}/agents/${encodeURIComponent(agentId)}/traces`; let url: string; const domainOverride = getAgent365ObservabilityDomainOverride(this.configProvider); @@ -260,7 +260,7 @@ export class Agent365Exporter implements SpanExporter { // Retry transient errors if ([408, 429].includes(response.status) || (response.status >= 500 && response.status < 600)) { if (attempt < DEFAULT_MAX_RETRIES) { - const sleepMs = 200 * (attempt + 1); + const sleepMs = 200 * (attempt + 1) + Math.floor(Math.random() * 100); logger.warn(`[Agent365Exporter] Transient error ${response.status}, correlation ID: ${correlationId}, retrying after ${sleepMs}ms`); await this.sleep(sleepMs); continue; diff --git a/packages/agents-a365-observability/src/tracing/util.ts b/packages/agents-a365-observability/src/tracing/util.ts index 2c051f2b..c12c2efe 100644 --- a/packages/agents-a365-observability/src/tracing/util.ts +++ b/packages/agents-a365-observability/src/tracing/util.ts @@ -21,3 +21,25 @@ export const isAgent365ExporterEnabled = ( const provider = configProvider ?? defaultObservabilityConfigurationProvider; return provider.getConfiguration().isObservabilityExporterEnabled; }; + +/** + * Maximum length for span attribute values. + * Values exceeding this limit will be truncated with a suffix. + */ +export const MAX_ATTRIBUTE_LENGTH = 8_192; + +const TRUNCATION_SUFFIX = '...[truncated]'; + +/** + * Truncate a string value to {@link MAX_ATTRIBUTE_LENGTH} characters. + * If the value exceeds the limit, it is trimmed and a truncation suffix is appended, + * with the total length capped at exactly {@link MAX_ATTRIBUTE_LENGTH}. + * @param value The string to truncate + * @returns The original string if within limits, otherwise the truncated string + */ +export function truncateValue(value: string): string { + if (value.length > MAX_ATTRIBUTE_LENGTH) { + return value.substring(0, MAX_ATTRIBUTE_LENGTH - TRUNCATION_SUFFIX.length) + TRUNCATION_SUFFIX; + } + return value; +} diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 7cf4f654..f194ae85 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -855,10 +855,10 @@ importers: specifier: workspace:* version: link:../../packages/agents-a365-tooling-extensions-langchain '@microsoft/agents-activity': - specifier: ^1.3.1 + specifier: 'catalog:' version: 1.3.1 '@microsoft/agents-hosting': - specifier: ^1.3.1 + specifier: 'catalog:' version: 1.3.1 dotenv: specifier: ^17.2.3 diff --git a/tests/observability/extension/hosting/agentic-token-cache.test.ts b/tests/observability/extension/hosting/agentic-token-cache.test.ts index fee7dbd1..06f457cc 100644 --- a/tests/observability/extension/hosting/agentic-token-cache.test.ts +++ b/tests/observability/extension/hosting/agentic-token-cache.test.ts @@ -156,4 +156,67 @@ describe('AgenticTokenCacheInstance', () => { const tokenAfter = AgenticTokenCacheInstance.getObservabilityToken('agentE', 'tenantE'); expect(tokenAfter).toBeNull(); }); + + it('evicts oldest entry when cache exceeds max size', async () => { + const { AgenticTokenCache } = require('@microsoft/agents-a365-observability-hosting'); + const cache = new AgenticTokenCache(); + const map = (cache as any)._map as Map; + + // Pre-fill the map to capacity + const MAX = (cache as any)._maxCacheSize as number; + for (let i = 0; i < MAX; i++) { + map.set(`agent-${i}:tenant-${i}`, { scopes: ['s'], token: `t-${i}`, acquiredOn: Date.now() }); + } + expect(map.size).toBe(MAX); + + // Insert one more via RefreshObservabilityToken + const token = makeJwtWithExp(300); + const auth = makeAuthorizationMock([{ token }]); + await cache.RefreshObservabilityToken( + 'agent-new', + 'tenant-new', + asTurnContext(makeTurnContext()), + auth as any, + ['scope.read'] + ); + + // Size should still be at MAX (oldest evicted, new one added) + expect(map.size).toBe(MAX); + // First entry should have been evicted + expect(map.has('agent-0:tenant-0')).toBe(false); + // New entry should exist + expect(map.has('agent-new:tenant-new')).toBe(true); + }); + + it('caps JWT exp claim to 24 hours', async () => { + const { AgenticTokenCache } = require('@microsoft/agents-a365-observability-hosting'); + const cache = new AgenticTokenCache(); + + // Create JWT with exp 48 hours from now + const farFutureExp = Math.floor(Date.now() / 1000) + (48 * 60 * 60); + const header = Buffer.from(JSON.stringify({ alg: 'none', typ: 'JWT' })).toString('base64url'); + const payload = Buffer.from(JSON.stringify({ exp: farFutureExp })).toString('base64url'); + const farFutureToken = `${header}.${payload}.sig`; + + const auth = makeAuthorizationMock([{ token: farFutureToken }]); + await cache.RefreshObservabilityToken( + 'agent-exp', + 'tenant-exp', + asTurnContext(makeTurnContext()), + auth as any, + ['scope.read'] + ); + + const map = (cache as any)._map as Map; + const entry = map.get('agent-exp:tenant-exp'); + expect(entry).toBeDefined(); + expect(entry.expiresOn).toBeDefined(); + + // The expiresOn should be capped to ~24 hours from now (not 48 hours) + const maxAllowed = Date.now() + (24 * 60 * 60 * 1000) + 5000; // 24h + small tolerance + expect(entry.expiresOn).toBeLessThanOrEqual(maxAllowed); + // And should be well below the 48-hour uncapped value + const uncapped = farFutureExp * 1000; + expect(entry.expiresOn).toBeLessThan(uncapped); + }); }); diff --git a/tests/observability/extension/openai/OpenAIAgentsTraceProcessor.test.ts b/tests/observability/extension/openai/OpenAIAgentsTraceProcessor.test.ts index 9fd453cb..4890d1cd 100644 --- a/tests/observability/extension/openai/OpenAIAgentsTraceProcessor.test.ts +++ b/tests/observability/extension/openai/OpenAIAgentsTraceProcessor.test.ts @@ -34,7 +34,7 @@ describe('OpenAIAgentsTraceProcessor', () => { let processor: OpenAIAgentsTraceProcessor; beforeEach(() => { - processor = new OpenAIAgentsTraceProcessor(tracer); + processor = new OpenAIAgentsTraceProcessor(tracer, { isContentRecordingEnabled: true }); }); afterEach(async () => { @@ -480,7 +480,7 @@ describe('OpenAIAgentsTraceProcessor', () => { }); it('does not record GEN_AI_INPUT_MESSAGES when disabled', async () => { - const processor = new OpenAIAgentsTraceProcessor(tracer, { suppressInvokeAgentInput: true }); + const processor = new OpenAIAgentsTraceProcessor(tracer, { suppressInvokeAgentInput: true, isContentRecordingEnabled: true }); const traceData = { traceId: 'trace-suppress', name: 'Agent' } as any; await processor.onTraceStart(traceData); @@ -503,8 +503,8 @@ describe('OpenAIAgentsTraceProcessor', () => { expect(keys).not.toContain(OpenTelemetryConstants.GEN_AI_INPUT_MESSAGES_KEY); }); - it('records GEN_AI_INPUT_MESSAGES when enabled (default)', async () => { - const processor = new OpenAIAgentsTraceProcessor(tracer); + it('records GEN_AI_INPUT_MESSAGES when content recording is enabled', async () => { + const processor = new OpenAIAgentsTraceProcessor(tracer, { isContentRecordingEnabled: true }); const traceData = { traceId: 'trace-allow', name: 'Agent' } as any; await processor.onTraceStart(traceData); @@ -528,7 +528,7 @@ describe('OpenAIAgentsTraceProcessor', () => { }); it('suppresses input on response spans when disabled', async () => { - const processor = new OpenAIAgentsTraceProcessor(tracer, { suppressInvokeAgentInput: true }); + const processor = new OpenAIAgentsTraceProcessor(tracer, { suppressInvokeAgentInput: true, isContentRecordingEnabled: true }); const traceData = { traceId: 'trace-resp', name: 'Agent' } as any; await processor.onTraceStart(traceData); @@ -552,7 +552,7 @@ describe('OpenAIAgentsTraceProcessor', () => { }); it('records full array JSON when only assistant messages are present', async () => { - const processor = new OpenAIAgentsTraceProcessor(tracer); + const processor = new OpenAIAgentsTraceProcessor(tracer, { isContentRecordingEnabled: true }); const traceData = { traceId: 'trace-assistant-only', name: 'Agent' } as any; await processor.onTraceStart(traceData); @@ -588,7 +588,7 @@ describe('OpenAIAgentsTraceProcessor', () => { expect(parsed).toEqual(inputArray); }); it('records user text content for array _input on response spans', async () => { - const processor = new OpenAIAgentsTraceProcessor(tracer); + const processor = new OpenAIAgentsTraceProcessor(tracer, { isContentRecordingEnabled: true }); const traceData = { traceId: 'trace-array-input', name: 'Agent' } as any; await processor.onTraceStart(traceData); @@ -621,7 +621,7 @@ describe('OpenAIAgentsTraceProcessor', () => { }); it('parses stringified array _input and records only user text content', async () => { - const processor = new OpenAIAgentsTraceProcessor(tracer); + const processor = new OpenAIAgentsTraceProcessor(tracer, { isContentRecordingEnabled: true }); const traceData = { traceId: 'trace-array-input-string', name: 'Agent' } as any; await processor.onTraceStart(traceData); @@ -657,7 +657,7 @@ describe('OpenAIAgentsTraceProcessor', () => { }); it('records [gen_ai.input.messages] attribute for array input with non standard schema on response spans', async () => { - const processor = new OpenAIAgentsTraceProcessor(tracer); + const processor = new OpenAIAgentsTraceProcessor(tracer, { isContentRecordingEnabled: true }); const traceData = { traceId: 'trace-array-input', name: 'Agent' } as any; await processor.onTraceStart(traceData); const inputArray = [ @@ -690,7 +690,7 @@ describe('OpenAIAgentsTraceProcessor', () => { }); it('records GEN_AI_OUTPUT_MESSAGES as plain string when output is a string', async () => { - const processor = new OpenAIAgentsTraceProcessor(tracer); + const processor = new OpenAIAgentsTraceProcessor(tracer, { isContentRecordingEnabled: true }); const traceData = { traceId: 'trace-output-string', name: 'Agent' } as any; await processor.onTraceStart(traceData); @@ -717,7 +717,7 @@ describe('OpenAIAgentsTraceProcessor', () => { }); it('records GEN_AI_OUTPUT_MESSAGES as aggregated texts when output is structured', async () => { - const processor = new OpenAIAgentsTraceProcessor(tracer); + const processor = new OpenAIAgentsTraceProcessor(tracer, { isContentRecordingEnabled: true }); const traceData = { traceId: 'trace-output-structured', name: 'Agent' } as any; await processor.onTraceStart(traceData); @@ -755,5 +755,40 @@ describe('OpenAIAgentsTraceProcessor', () => { const parsed = JSON.parse(value); expect(parsed).toEqual(['Hello user 1', 'Hello user 2']); }); + + it('suppresses all content attributes when isContentRecordingEnabled is false', async () => { + const processor = new OpenAIAgentsTraceProcessor(tracer); + const traceData = { traceId: 'trace-no-content', name: 'Agent' } as any; + await processor.onTraceStart(traceData); + + // Generation span with input/output + const genSpan = { + spanId: 'gen-no-content', + traceId: 'trace-no-content', + startedAt: new Date().toISOString(), + spanData: { + type: 'generation' as const, + model: 'gpt-4', + input: [{ role: 'user', content: 'secret prompt' }], + output: { id: 'resp-1', choices: [{ text: 'secret response' }] }, + }, + } as any; + + await processor.onSpanStart(genSpan); + await processor.onSpanEnd(genSpan); + + const mock = spansByName['generation']; + const attrs = mock._attrs as Array<[string, unknown]>; + const contentKeys = attrs.filter(([k]) => + k === OpenTelemetryConstants.GEN_AI_INPUT_MESSAGES_KEY || + k === OpenTelemetryConstants.GEN_AI_OUTPUT_MESSAGES_KEY + ); + expect(contentKeys).toHaveLength(0); + + // Model attribute should still be present (non-content) + const modelAttr = attrs.find(([k]) => k === OpenTelemetryConstants.GEN_AI_REQUEST_MODEL_KEY); + expect(modelAttr).toBeDefined(); + }); + }); }); diff --git a/tests/observability/integration/openai-agent-instrument.test.ts b/tests/observability/integration/openai-agent-instrument.test.ts index 30b70366..9f02fb9b 100644 --- a/tests/observability/integration/openai-agent-instrument.test.ts +++ b/tests/observability/integration/openai-agent-instrument.test.ts @@ -43,6 +43,7 @@ describe("OpenAI Trace Processor Integration Tests", () => { enabled: true, tracerName: TEST_INSTRUMENTATION_NAME, tracerVersion: TEST_INSTRUMENTATION_VERSION, + isContentRecordingEnabled: true, }); // Start observability diff --git a/tests/observability/tracing/truncation.test.ts b/tests/observability/tracing/truncation.test.ts new file mode 100644 index 00000000..9001abe4 --- /dev/null +++ b/tests/observability/tracing/truncation.test.ts @@ -0,0 +1,51 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT License. + +import { describe, it, expect } from '@jest/globals'; +import { truncateValue, MAX_ATTRIBUTE_LENGTH } from '../../../packages/agents-a365-observability/src/tracing/util'; + +describe('truncateValue', () => { + const SUFFIX = '...[truncated]'; + + it('should return the original string when within limit', () => { + const value = 'hello world'; + expect(truncateValue(value)).toBe(value); + }); + + it('should return the original string when exactly at limit', () => { + const value = 'x'.repeat(MAX_ATTRIBUTE_LENGTH); + expect(truncateValue(value)).toBe(value); + expect(truncateValue(value).length).toBe(MAX_ATTRIBUTE_LENGTH); + }); + + it('should truncate when 1 character over limit', () => { + const value = 'x'.repeat(MAX_ATTRIBUTE_LENGTH + 1); + const result = truncateValue(value); + expect(result.length).toBe(MAX_ATTRIBUTE_LENGTH); + expect(result.endsWith(SUFFIX)).toBe(true); + }); + + it('should truncate long strings to exactly MAX_ATTRIBUTE_LENGTH', () => { + const value = 'a'.repeat(MAX_ATTRIBUTE_LENGTH * 2); + const result = truncateValue(value); + expect(result.length).toBe(MAX_ATTRIBUTE_LENGTH); + expect(result.endsWith(SUFFIX)).toBe(true); + }); + + it('should preserve the beginning of the string when truncating', () => { + const prefix = 'PREFIX_'; + const value = prefix + 'x'.repeat(MAX_ATTRIBUTE_LENGTH); + const result = truncateValue(value); + expect(result.startsWith(prefix)).toBe(true); + }); + + it('should return empty string unchanged', () => { + expect(truncateValue('')).toBe(''); + }); +}); + +describe('MAX_ATTRIBUTE_LENGTH', () => { + it('should be 8192', () => { + expect(MAX_ATTRIBUTE_LENGTH).toBe(8_192); + }); +});