diff --git a/libs/otel-nestjs-instrumentation/src/internal/index.ts b/libs/otel-nestjs-instrumentation/src/internal/index.ts index edf7b0c..9f7540e 100644 --- a/libs/otel-nestjs-instrumentation/src/internal/index.ts +++ b/libs/otel-nestjs-instrumentation/src/internal/index.ts @@ -2,5 +2,7 @@ export * from './emitter-symbol'; export * from './get-transaction-name'; export * from './internal-context'; export * from './otel-instrumentation'; +export * from './record-messaging-process-duration'; +export * from './resolve-rpc-messaging-metadata'; export * from './tracer-name'; export * from './transaction-symbol'; diff --git a/libs/otel-nestjs-instrumentation/src/internal/otel-instrumentation.ts b/libs/otel-nestjs-instrumentation/src/internal/otel-instrumentation.ts index 5fe074a..64248bd 100644 --- a/libs/otel-nestjs-instrumentation/src/internal/otel-instrumentation.ts +++ b/libs/otel-nestjs-instrumentation/src/internal/otel-instrumentation.ts @@ -1,5 +1,9 @@ import { ExecutionContext } from '@nestjs/common'; import * as otel from '@opentelemetry/api'; +import { + resolveRpcMessagingMetadata, + toMessagingSpanAttributes, +} from './resolve-rpc-messaging-metadata'; import { tracerName } from './tracer-name'; /** @@ -39,7 +43,7 @@ export const otelInstrumentation = { transactionName: string, context: ExecutionContext, effectiveType?: 'http' | 'rpc', - ) { + ): string | undefined { // Create a new span since none exists const tracer = otel.trace.getTracer(tracerName); if (!tracer) return undefined; @@ -63,8 +67,14 @@ export const otelInstrumentation = { // Determine span kind and attributes based on effective type let spanKind = otel.SpanKind.INTERNAL; // default for unknown contexts const attributes: Record = {}; + let messagingMetadata: ReturnType; - effectiveType ??= context.getType() as 'http' | 'rpc'; + if (!effectiveType) { + const contextType = context.getType(); + if (contextType === 'http' || contextType === 'rpc') { + effectiveType = contextType; + } + } if (effectiveType === 'http') { spanKind = otel.SpanKind.SERVER; @@ -83,7 +93,19 @@ export const otelInstrumentation = { // Ignore request extraction errors } } else if (effectiveType === 'rpc') { - spanKind = otel.SpanKind.SERVER; + messagingMetadata = resolveRpcMessagingMetadata(context); + if (messagingMetadata) { + spanKind = messagingMetadata.spanKind; + Object.assign(attributes, toMessagingSpanAttributes(messagingMetadata)); + if (Object.keys(messagingMetadata.propagationCarrier).length > 0) { + spanContext = otel.propagation.extract( + spanContext, + messagingMetadata.propagationCarrier, + ); + } + } else { + spanKind = otel.SpanKind.SERVER; + } try { attributes['rpc.method'] = context.getHandler()?.name ?? 'Call'; } catch { @@ -101,8 +123,13 @@ export const otelInstrumentation = { // Ignore NestJS context extraction errors } + const spanName = + messagingMetadata?.recordMetric === true + ? `process ${messagingMetadata.operationName}` + : transactionName; + const span = tracer.startSpan( - transactionName, + spanName, { kind: spanKind, attributes, @@ -110,8 +137,18 @@ export const otelInstrumentation = { spanContext, ); - const spanContextData = span.spanContext(); - const traceId = spanContextData?.traceId; + const traceId = span.spanContext()?.traceId; + if (!traceId) { + span.end(); + return undefined; + } + + const activeContext = otel.trace.setSpan(spanContext, span); + ( + otel.context as typeof otel.context & { + enterWith: (ctx: otel.Context) => void; + } + ).enterWith(activeContext); return traceId; }, diff --git a/libs/otel-nestjs-instrumentation/src/internal/record-messaging-process-duration.ts b/libs/otel-nestjs-instrumentation/src/internal/record-messaging-process-duration.ts new file mode 100644 index 0000000..266fa92 --- /dev/null +++ b/libs/otel-nestjs-instrumentation/src/internal/record-messaging-process-duration.ts @@ -0,0 +1,35 @@ +import { metrics, type Histogram } from '@opentelemetry/api'; +import { tracerName } from './tracer-name'; +import { + type RpcMessagingMetadata, + toMessagingSpanAttributes, +} from './resolve-rpc-messaging-metadata'; + +let processDurationHistogram: Histogram | undefined; + +function getProcessDurationHistogram(): Histogram { + processDurationHistogram ??= metrics + .getMeter(tracerName) + .createHistogram('messaging.process.duration', { + description: 'Measures the duration of inbound messaging operations.', + unit: 'ms', + }); + return processDurationHistogram; +} + +/** @internal Test-only reset for module-level meter instrument cache. */ +export function __resetMessagingProcessDurationForTests(): void { + processDurationHistogram = undefined; +} + +export function recordMessagingProcessDuration( + durationMs: number, + metadata: RpcMessagingMetadata, +): void { + if (!metadata.recordMetric) return; + + getProcessDurationHistogram().record( + durationMs, + toMessagingSpanAttributes(metadata), + ); +} diff --git a/libs/otel-nestjs-instrumentation/src/internal/resolve-rpc-messaging-metadata.ts b/libs/otel-nestjs-instrumentation/src/internal/resolve-rpc-messaging-metadata.ts new file mode 100644 index 0000000..85bbce7 --- /dev/null +++ b/libs/otel-nestjs-instrumentation/src/internal/resolve-rpc-messaging-metadata.ts @@ -0,0 +1,173 @@ +import { ExecutionContext } from '@nestjs/common'; +import * as otel from '@opentelemetry/api'; +import { getTransactionName } from './get-transaction-name'; + +type SqsMessageLike = { + MessageId?: string; + MessageAttributes?: Record; +}; + +type KafkaMessageLike = { + headers?: Record; + offset?: string | number; +}; + +type RpcContextLike = { + getTopic?: () => string; + getPattern?: () => string; + getMessage?: () => unknown; + getChannelRef?: () => { fields?: { routingKey?: string } }; +}; + +export type RpcMessagingMetadata = { + system: string; + destination: string; + operationName: string; + messageId?: string; + propagationCarrier: Record; + spanKind: otel.SpanKind; + recordMetric: boolean; +}; + +export function toMessagingSpanAttributes( + metadata: RpcMessagingMetadata, +): Record { + const attributes: Record = { + 'messaging.system': metadata.system, + 'messaging.destination.name': metadata.destination, + 'messaging.operation.type': 'process', + 'messaging.operation.name': metadata.operationName, + }; + if (metadata.messageId) { + attributes['messaging.message.id'] = metadata.messageId; + } + return attributes; +} + +function buildConsumerMetadata( + system: string, + destination: string, + operationName: string, + extras?: Pick, +): RpcMessagingMetadata { + return { + system, + destination, + operationName, + messageId: extras?.messageId, + propagationCarrier: extras?.propagationCarrier ?? {}, + spanKind: otel.SpanKind.CONSUMER, + recordMetric: true, + }; +} + +function normalizeCarrier( + entries: Record, + resolveValue: (value: unknown) => string | undefined, +): Record { + const carrier: Record = {}; + for (const [key, value] of Object.entries(entries)) { + const normalized = resolveValue(value); + if (normalized) { + carrier[key.toLowerCase()] = normalized; + } + } + return carrier; +} + +function extractSqsCarrier(message: SqsMessageLike): Record { + return normalizeCarrier(message.MessageAttributes ?? {}, (value) => { + if ( + typeof value === 'object' && + value !== null && + 'StringValue' in value && + typeof (value as { StringValue?: string }).StringValue === 'string' + ) { + return (value as { StringValue: string }).StringValue; + } + return undefined; + }); +} + +function extractKafkaCarrier( + message: KafkaMessageLike | undefined, +): Record { + return normalizeCarrier(message?.headers ?? {}, (value) => { + if (value === undefined) return undefined; + return Buffer.isBuffer(value) ? value.toString('utf8') : String(value); + }); +} + +function isSqsMessage(message: unknown): message is SqsMessageLike { + return ( + typeof message === 'object' && + message !== null && + 'MessageId' in message && + typeof (message as SqsMessageLike).MessageId === 'string' + ); +} + +function buildGenericRpcMetadata(operationName: string): RpcMessagingMetadata { + return { + system: 'nestjs', + destination: operationName, + operationName, + propagationCarrier: {}, + spanKind: otel.SpanKind.SERVER, + recordMetric: false, + }; +} + +/** + * Resolves messaging semantic attributes for NestJS RPC/microservice handlers. + * + * Works across custom and built-in transports (SQS, Kafka, RabbitMQ, etc.) + * by duck-typing `context.switchToRpc().getContext()`. + */ +export function resolveRpcMessagingMetadata( + context: ExecutionContext, +): RpcMessagingMetadata | undefined { + if (context.getType() !== 'rpc') return undefined; + + const operationName = getTransactionName(context); + let rpcContext: RpcContextLike; + try { + rpcContext = context.switchToRpc().getContext(); + } catch { + return buildGenericRpcMetadata(operationName); + } + + if (typeof rpcContext.getTopic === 'function') { + const message = rpcContext.getMessage?.() as KafkaMessageLike | undefined; + return buildConsumerMetadata( + 'kafka', + rpcContext.getTopic(), + operationName, + { + messageId: + message?.offset === undefined ? undefined : String(message.offset), + propagationCarrier: extractKafkaCarrier(message), + }, + ); + } + + const message = rpcContext.getMessage?.(); + if (isSqsMessage(message)) { + return buildConsumerMetadata('aws_sqs', operationName, operationName, { + messageId: message.MessageId, + propagationCarrier: extractSqsCarrier(message), + }); + } + + if (typeof rpcContext.getPattern === 'function') { + const pattern = rpcContext.getPattern(); + const routingKey = rpcContext.getChannelRef?.()?.fields?.routingKey; + return buildConsumerMetadata( + 'rabbitmq', + routingKey ?? pattern, + operationName, + ); + } + + return buildGenericRpcMetadata(operationName); +} diff --git a/libs/otel-nestjs-instrumentation/src/otel.interceptor.ts b/libs/otel-nestjs-instrumentation/src/otel.interceptor.ts index 9f362e6..30afad2 100644 --- a/libs/otel-nestjs-instrumentation/src/otel.interceptor.ts +++ b/libs/otel-nestjs-instrumentation/src/otel.interceptor.ts @@ -10,11 +10,15 @@ import { emitterSymbol, InternalContext, otelInstrumentation, + resolveRpcMessagingMetadata, + recordMessagingProcessDuration, } from './internal'; import EventEmitter from 'events'; import otel, { Span } from '@opentelemetry/api'; import { startOtelInstrumentationIfAbsent } from './start-otel-instrumentation-if-absent'; +const NANOSECONDS_PER_MILLISECOND = 1_000_000; + /** * NestJS interceptor that manages OpenTelemetry span lifecycle. * @@ -106,22 +110,38 @@ export class OtelInterceptor implements NestInterceptor { * @returns Observable that completes when the request is finished */ intercept(context: ExecutionContext, next: CallHandler) { - // If no guard ran (e.g. gRPC/microservice), force RPC since HTTP spans - // would have been started by the guard already. - startOtelInstrumentationIfAbsent( - context, - this.context, - this.emitter, - 'rpc', - ); + const isRpc = context.getType() === 'rpc'; + const messagingMetadata = isRpc + ? resolveRpcMessagingMetadata(context) + : undefined; + const startedAt = + messagingMetadata?.recordMetric === true + ? process.hrtime.bigint() + : undefined; + + if (isRpc) { + startOtelInstrumentationIfAbsent( + context, + this.context, + this.emitter, + 'rpc', + ); + } else { + startOtelInstrumentationIfAbsent(context, this.context, this.emitter); + } + const span = otel.trace.getActiveSpan(); if (!span) return next.handle(); const traceId = span.spanContext().traceId; return next.handle().pipe( tap({ - next: () => this.finishSpan(traceId, span), + next: () => { + this.recordMessagingMetrics(startedAt, messagingMetadata); + this.finishSpan(traceId, span); + }, error: (error) => { + this.recordMessagingMetrics(startedAt, messagingMetadata); this.recordError(error); this.finishSpan(traceId, span); }, @@ -129,6 +149,17 @@ export class OtelInterceptor implements NestInterceptor { ); } + private recordMessagingMetrics( + startedAt: bigint | undefined, + messagingMetadata: ReturnType, + ): void { + if (startedAt === undefined || !messagingMetadata?.recordMetric) return; + + const durationMs = + Number(process.hrtime.bigint() - startedAt) / NANOSECONDS_PER_MILLISECOND; + recordMessagingProcessDuration(durationMs, messagingMetadata); + } + /** * Complete the span and emit the appropriate events. * diff --git a/libs/otel-nestjs-instrumentation/test/internal/internal-utilities.spec.ts b/libs/otel-nestjs-instrumentation/test/internal/internal-utilities.spec.ts index 10815f9..424b2fd 100644 --- a/libs/otel-nestjs-instrumentation/test/internal/internal-utilities.spec.ts +++ b/libs/otel-nestjs-instrumentation/test/internal/internal-utilities.spec.ts @@ -8,19 +8,21 @@ const mockOtelApi = { end: jest.fn(), })), })), + setSpan: jest.fn((context, span) => ({ ...context, span })), }, context: { active: jest.fn(() => ({})), + enterWith: jest.fn(), }, propagation: { extract: jest.fn((context, _headers) => context), }, SpanKind: { - INTERNAL: 1, - SERVER: 2, - CLIENT: 3, - PRODUCER: 4, - CONSUMER: 5, + INTERNAL: 0, + SERVER: 1, + CLIENT: 2, + PRODUCER: 3, + CONSUMER: 4, }, SpanStatusCode: { UNSET: 0, diff --git a/libs/otel-nestjs-instrumentation/test/internal/otel-instrumentation.spec.ts b/libs/otel-nestjs-instrumentation/test/internal/otel-instrumentation.spec.ts index 5c09740..35bdebc 100644 --- a/libs/otel-nestjs-instrumentation/test/internal/otel-instrumentation.spec.ts +++ b/libs/otel-nestjs-instrumentation/test/internal/otel-instrumentation.spec.ts @@ -2,9 +2,11 @@ const mockOtelApi = { trace: { getTracer: jest.fn(), getActiveSpan: jest.fn(), + setSpan: jest.fn((context, span) => ({ ...context, span })), }, context: { active: jest.fn(() => ({})), + enterWith: jest.fn(), }, propagation: { extract: jest.fn((context, headers) => ({ extracted: true, headers })), @@ -22,8 +24,32 @@ jest.mock('@opentelemetry/api', () => mockOtelApi); import { ExecutionContext } from '@nestjs/common'; import { otelInstrumentation } from '../../src/internal/otel-instrumentation'; +import { toMessagingSpanAttributes } from '../../src/internal/resolve-rpc-messaging-metadata'; import { createMockExecutionContext } from '../test-utils'; +function messagingAttributes( + metadata: { + system: string; + destination: string; + operationName: string; + messageId?: string; + }, + extra: Record = {}, +): Record { + return { + ...toMessagingSpanAttributes({ + system: metadata.system, + destination: metadata.destination, + operationName: metadata.operationName, + messageId: metadata.messageId, + propagationCarrier: {}, + spanKind: 0, + recordMetric: true, + }), + ...extra, + }; +} + describe('OtelInstrumentation', () => { let mockExecutionContext: ExecutionContext; let mockTracer: any; @@ -202,6 +228,7 @@ describe('OtelInstrumentation', () => { const result = otelInstrumentation.create( 'RpcController.rpcHandler', rpcContext, + 'rpc', ); expect(mockTracer.startSpan).toHaveBeenCalledWith( @@ -209,6 +236,11 @@ describe('OtelInstrumentation', () => { { kind: mockOtelApi.SpanKind.SERVER, attributes: { + ...messagingAttributes({ + system: 'nestjs', + destination: 'RpcController.rpcHandler', + operationName: 'RpcController.rpcHandler', + }), 'rpc.method': 'rpcHandler', 'nestjs.controller': 'RpcController', 'nestjs.handler': 'rpcHandler', @@ -216,9 +248,53 @@ describe('OtelInstrumentation', () => { }, {}, ); + expect(mockOtelApi.context.enterWith).toHaveBeenCalled(); expect(result).toBe('test-trace-id'); }); + it('should create consumer span with messaging attributes for Kafka RPC contexts', () => { + const rpcContext = createMockExecutionContext('rpc', { + controller: 'OrdersController', + handler: 'handleOrder', + }); + jest.spyOn(rpcContext.switchToRpc(), 'getContext').mockReturnValue({ + getTopic: () => 'orders.created', + getMessage: () => ({ + offset: '7', + headers: { traceparent: '00-abc' }, + }), + }); + + otelInstrumentation.create( + 'OrdersController.handleOrder', + rpcContext, + 'rpc', + ); + + expect(mockOtelApi.propagation.extract).toHaveBeenCalledWith( + {}, + { traceparent: '00-abc' }, + ); + expect(mockTracer.startSpan).toHaveBeenCalledWith( + 'process OrdersController.handleOrder', + { + kind: mockOtelApi.SpanKind.CONSUMER, + attributes: { + ...messagingAttributes({ + system: 'kafka', + destination: 'orders.created', + operationName: 'OrdersController.handleOrder', + messageId: '7', + }), + 'rpc.method': 'handleOrder', + 'nestjs.controller': 'OrdersController', + 'nestjs.handler': 'handleOrder', + }, + }, + { extracted: true, headers: { traceparent: '00-abc' } }, + ); + }); + it('should handle unknown context type', () => { const unknownContext = createMockExecutionContext('ws', { controller: 'WsController', @@ -251,16 +327,21 @@ describe('OtelInstrumentation', () => { Object.defineProperty(mockHandler, 'name', { value: '' }); jest.spyOn(rpcContext, 'getHandler').mockReturnValue(mockHandler); - otelInstrumentation.create('RpcController.rpcHandler', rpcContext); + otelInstrumentation.create('RpcController.rpcHandler', rpcContext, 'rpc'); expect(mockTracer.startSpan).toHaveBeenCalledWith( 'RpcController.rpcHandler', { kind: mockOtelApi.SpanKind.SERVER, attributes: { - 'rpc.method': '', // Empty string preserved by ?? (not null/undefined) + ...messagingAttributes({ + system: 'nestjs', + destination: 'RpcController.unknownMethod', + operationName: 'RpcController.unknownMethod', + }), + 'rpc.method': '', 'nestjs.controller': 'RpcController', - 'nestjs.handler': 'unknown', // Falls back to 'unknown' when name is empty + 'nestjs.handler': 'unknown', }, }, {}, @@ -311,16 +392,17 @@ describe('OtelInstrumentation', () => { throw new Error('Handler extraction failed'); }); - otelInstrumentation.create('RpcController.rpcHandler', rpcContext); + otelInstrumentation.create('RpcController.rpcHandler', rpcContext, 'rpc'); expect(mockTracer.startSpan).toHaveBeenCalledWith( 'RpcController.rpcHandler', { kind: mockOtelApi.SpanKind.SERVER, - attributes: { - // Both RPC and NestJS context extraction should fail since getHandler throws - // No attributes should be set when the entire block fails - }, + attributes: messagingAttributes({ + system: 'nestjs', + destination: 'UnknownTransaction', + operationName: 'UnknownTransaction', + }), }, {}, ); diff --git a/libs/otel-nestjs-instrumentation/test/internal/record-messaging-process-duration.spec.ts b/libs/otel-nestjs-instrumentation/test/internal/record-messaging-process-duration.spec.ts new file mode 100644 index 0000000..bbef3eb --- /dev/null +++ b/libs/otel-nestjs-instrumentation/test/internal/record-messaging-process-duration.spec.ts @@ -0,0 +1,58 @@ +const mockRecord = jest.fn(); + +jest.mock('@opentelemetry/api', () => ({ + metrics: { + getMeter: jest.fn(() => ({ + createHistogram: jest.fn(() => ({ record: mockRecord })), + })), + }, +})); + +import { + __resetMessagingProcessDurationForTests, + recordMessagingProcessDuration, +} from '../../src/internal/record-messaging-process-duration'; +import { toMessagingSpanAttributes } from '../../src/internal/resolve-rpc-messaging-metadata'; + +describe('recordMessagingProcessDuration', () => { + beforeEach(() => { + mockRecord.mockClear(); + __resetMessagingProcessDurationForTests(); + }); + + it('records messaging.process.duration for messaging transports', () => { + recordMessagingProcessDuration(12.5, { + system: 'aws_sqs', + destination: 'announcements', + operationName: 'announcement-dispatch', + propagationCarrier: {}, + spanKind: 4, + recordMetric: true, + }); + + expect(mockRecord).toHaveBeenCalledWith( + 12.5, + toMessagingSpanAttributes({ + system: 'aws_sqs', + destination: 'announcements', + operationName: 'announcement-dispatch', + propagationCarrier: {}, + spanKind: 4, + recordMetric: true, + }), + ); + }); + + it('skips metrics for generic Nest RPC handlers', () => { + recordMessagingProcessDuration(12.5, { + system: 'nestjs', + destination: 'RpcController.rpcHandler', + operationName: 'RpcController.rpcHandler', + propagationCarrier: {}, + spanKind: 1, + recordMetric: false, + }); + + expect(mockRecord).not.toHaveBeenCalled(); + }); +}); diff --git a/libs/otel-nestjs-instrumentation/test/internal/resolve-rpc-messaging-metadata.spec.ts b/libs/otel-nestjs-instrumentation/test/internal/resolve-rpc-messaging-metadata.spec.ts new file mode 100644 index 0000000..a44f61d --- /dev/null +++ b/libs/otel-nestjs-instrumentation/test/internal/resolve-rpc-messaging-metadata.spec.ts @@ -0,0 +1,115 @@ +import * as otel from '@opentelemetry/api'; +import { + resolveRpcMessagingMetadata, + type RpcMessagingMetadata, +} from '../../src/internal/resolve-rpc-messaging-metadata'; +import { createMockExecutionContext } from '../test-utils'; + +const TRACE_PARENT = '00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01'; + +const consumerDefaults: Pick< + RpcMessagingMetadata, + 'spanKind' | 'recordMetric' | 'propagationCarrier' +> = { + propagationCarrier: {}, + spanKind: otel.SpanKind.CONSUMER, + recordMetric: true, +}; + +describe('resolveRpcMessagingMetadata', () => { + it('returns undefined for HTTP contexts', () => { + expect( + resolveRpcMessagingMetadata(createMockExecutionContext('http')), + ).toBeUndefined(); + }); + + it('falls back to generic Nest RPC metadata without metrics', () => { + const context = createMockExecutionContext('rpc', { + controller: 'RpcController', + handler: 'rpcHandler', + }); + + expect(resolveRpcMessagingMetadata(context)).toEqual({ + system: 'nestjs', + destination: 'RpcController.rpcHandler', + operationName: 'RpcController.rpcHandler', + propagationCarrier: {}, + spanKind: otel.SpanKind.SERVER, + recordMetric: false, + }); + }); + + it.each([ + [ + 'Kafka', + 'OrdersController', + 'handleOrder', + { + getTopic: () => 'orders.created', + getMessage: () => ({ + offset: '42', + headers: { + traceparent: Buffer.from(TRACE_PARENT), + }, + }), + }, + { + system: 'kafka', + destination: 'orders.created', + operationName: 'OrdersController.handleOrder', + messageId: '42', + propagationCarrier: { traceparent: TRACE_PARENT }, + }, + ], + [ + 'SQS', + 'Announcements', + 'handleAnnouncement', + { + getMessage: () => ({ + MessageId: 'msg-123', + MessageAttributes: { + traceparent: { StringValue: TRACE_PARENT }, + }, + }), + }, + { + system: 'aws_sqs', + destination: 'Announcements.handleAnnouncement', + operationName: 'Announcements.handleAnnouncement', + messageId: 'msg-123', + propagationCarrier: { traceparent: TRACE_PARENT }, + }, + ], + [ + 'RabbitMQ', + 'EventsController', + 'handleEvent', + { + getPattern: () => 'events.created', + getChannelRef: () => ({ fields: { routingKey: 'events.created' } }), + }, + { + system: 'rabbitmq', + destination: 'events.created', + operationName: 'EventsController.handleEvent', + }, + ], + ] as const)( + 'resolves %s metadata from RPC contexts', + (_label, controller, handler, rpcContext, expected) => { + const context = createMockExecutionContext('rpc', { + controller, + handler, + }); + jest + .spyOn(context.switchToRpc(), 'getContext') + .mockReturnValue(rpcContext); + + expect(resolveRpcMessagingMetadata(context)).toEqual({ + ...consumerDefaults, + ...expected, + }); + }, + ); +}); diff --git a/libs/otel-nestjs-instrumentation/test/otel-integration.spec.ts b/libs/otel-nestjs-instrumentation/test/otel-integration.spec.ts index 1e608b7..504c9cf 100644 --- a/libs/otel-nestjs-instrumentation/test/otel-integration.spec.ts +++ b/libs/otel-nestjs-instrumentation/test/otel-integration.spec.ts @@ -36,19 +36,23 @@ const mockOtelApi = { }; }), })), + setSpan: jest.fn((context, span) => ({ ...context, span })), }, context: { active: jest.fn(() => ({})), + enterWith: jest.fn(() => { + activeSpanCount++; + }), }, propagation: { extract: jest.fn((context, _headers) => context), }, SpanKind: { - INTERNAL: 1, - SERVER: 2, - CLIENT: 3, - PRODUCER: 4, - CONSUMER: 5, + INTERNAL: 0, + SERVER: 1, + CLIENT: 2, + PRODUCER: 3, + CONSUMER: 4, }, SpanStatusCode: { UNSET: 0, diff --git a/libs/otel-nestjs-instrumentation/test/otel-interceptor.spec.ts b/libs/otel-nestjs-instrumentation/test/otel-interceptor.spec.ts index f8bf804..ce9fa98 100644 --- a/libs/otel-nestjs-instrumentation/test/otel-interceptor.spec.ts +++ b/libs/otel-nestjs-instrumentation/test/otel-interceptor.spec.ts @@ -11,6 +11,18 @@ const mockOtelApi = { end: jest.fn(), })), })), + setSpan: jest.fn((context, span) => ({ ...context, span })), + }, + context: { + active: jest.fn(() => ({})), + enterWith: jest.fn(), + }, + SpanKind: { + INTERNAL: 0, + SERVER: 1, + CLIENT: 2, + PRODUCER: 3, + CONSUMER: 4, }, SpanStatusCode: { UNSET: 0, @@ -434,6 +446,10 @@ describe('OtelInterceptor', () => { }); it('should call startOtelInstrumentationIfAbsent as fallback for RPC contexts', async () => { + mockExecutionContext = createMockExecutionContext('rpc', { + controller: 'TestController', + handler: 'testMethod', + }); mockOtelApi.trace.getActiveSpan.mockReturnValue(null); mockOtelInstrumentation.getCurrentTransactionId.mockReturnValue( undefined,