From 8e6626875baeb8087c364c8b000d233785072cf5 Mon Sep 17 00:00:00 2001 From: Thiago Oliveira Santos Date: Fri, 4 Sep 2026 15:38:42 -0300 Subject: [PATCH 1/4] feat(otel-nestjs-instrumentation): add RPC messaging spans and metrics Instrument Nest microservice handlers with messaging semantic conventions and messaging.process.duration for SQS, Kafka, and RabbitMQ transports, without affecting HTTP handlers or duplicating http.server.duration. Co-authored-by: Cursor --- .../src/internal/index.ts | 2 + .../src/internal/otel-instrumentation.ts | 50 ++++++- .../record-messaging-process-duration.ts | 34 +++++ .../resolve-rpc-messaging-metadata.ts | 135 ++++++++++++++++++ .../src/otel.interceptor.ts | 49 +++++-- .../test/internal/internal-utilities.spec.ts | 12 +- .../internal/otel-instrumentation.spec.ts | 68 ++++++++- .../record-messaging-process-duration.spec.ts | 52 +++++++ .../resolve-rpc-messaging-metadata.spec.ts | 107 ++++++++++++++ .../test/otel-integration.spec.ts | 14 +- .../test/otel-interceptor.spec.ts | 16 +++ 11 files changed, 508 insertions(+), 31 deletions(-) create mode 100644 libs/otel-nestjs-instrumentation/src/internal/record-messaging-process-duration.ts create mode 100644 libs/otel-nestjs-instrumentation/src/internal/resolve-rpc-messaging-metadata.ts create mode 100644 libs/otel-nestjs-instrumentation/test/internal/record-messaging-process-duration.spec.ts create mode 100644 libs/otel-nestjs-instrumentation/test/internal/resolve-rpc-messaging-metadata.spec.ts 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..649580c 100644 --- a/libs/otel-nestjs-instrumentation/src/internal/otel-instrumentation.ts +++ b/libs/otel-nestjs-instrumentation/src/internal/otel-instrumentation.ts @@ -1,5 +1,6 @@ import { ExecutionContext } from '@nestjs/common'; import * as otel from '@opentelemetry/api'; +import { resolveRpcMessagingMetadata } from './resolve-rpc-messaging-metadata'; import { tracerName } from './tracer-name'; /** @@ -39,7 +40,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; @@ -64,7 +65,7 @@ export const otelInstrumentation = { let spanKind = otel.SpanKind.INTERNAL; // default for unknown contexts const attributes: Record = {}; - effectiveType ??= context.getType() as 'http' | 'rpc'; + effectiveType ??= context.getType(); if (effectiveType === 'http') { spanKind = otel.SpanKind.SERVER; @@ -83,7 +84,25 @@ export const otelInstrumentation = { // Ignore request extraction errors } } else if (effectiveType === 'rpc') { - spanKind = otel.SpanKind.SERVER; + const messaging = resolveRpcMessagingMetadata(context); + if (messaging) { + spanKind = messaging.spanKind; + attributes['messaging.system'] = messaging.system; + attributes['messaging.destination.name'] = messaging.destination; + attributes['messaging.operation.type'] = 'process'; + attributes['messaging.operation.name'] = messaging.operationName; + if (messaging.messageId) { + attributes['messaging.message.id'] = messaging.messageId; + } + if (Object.keys(messaging.propagationCarrier).length > 0) { + spanContext = otel.propagation.extract( + spanContext, + messaging.propagationCarrier, + ); + } + } else { + spanKind = otel.SpanKind.SERVER; + } try { attributes['rpc.method'] = context.getHandler()?.name ?? 'Call'; } catch { @@ -101,8 +120,17 @@ export const otelInstrumentation = { // Ignore NestJS context extraction errors } + const messaging = + effectiveType === 'rpc' + ? resolveRpcMessagingMetadata(context) + : undefined; + const spanName = + messaging?.recordMetric === true + ? `process ${messaging.operationName}` + : transactionName; + const span = tracer.startSpan( - transactionName, + spanName, { kind: spanKind, attributes, @@ -110,8 +138,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..733a17d --- /dev/null +++ b/libs/otel-nestjs-instrumentation/src/internal/record-messaging-process-duration.ts @@ -0,0 +1,34 @@ +import { metrics, type Histogram } from '@opentelemetry/api'; +import { tracerName } from './tracer-name'; +import type { RpcMessagingMetadata } 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, { + 'messaging.system': metadata.system, + 'messaging.destination.name': metadata.destination, + 'messaging.operation.type': 'process', + 'messaging.operation.name': metadata.operationName, + }); +} 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..3c52718 --- /dev/null +++ b/libs/otel-nestjs-instrumentation/src/internal/resolve-rpc-messaging-metadata.ts @@ -0,0 +1,135 @@ +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; +}; + +function extractSqsCarrier(message: SqsMessageLike): Record { + const carrier: Record = {}; + for (const [key, value] of Object.entries(message.MessageAttributes ?? {})) { + if (value.StringValue) { + carrier[key.toLowerCase()] = value.StringValue; + } + } + return carrier; +} + +function extractKafkaCarrier( + message: KafkaMessageLike | undefined, +): Record { + const carrier: Record = {}; + for (const [key, value] of Object.entries(message?.headers ?? {})) { + if (value === undefined) continue; + carrier[key.toLowerCase()] = Buffer.isBuffer(value) + ? value.toString('utf8') + : String(value); + } + return carrier; +} + +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 { + system: 'kafka', + destination: rpcContext.getTopic(), + operationName, + messageId: + message?.offset === undefined ? undefined : String(message.offset), + propagationCarrier: extractKafkaCarrier(message), + spanKind: otel.SpanKind.CONSUMER, + recordMetric: true, + }; + } + + const message = rpcContext.getMessage?.(); + if (isSqsMessage(message)) { + return { + system: 'aws_sqs', + destination: operationName, + operationName, + messageId: message.MessageId, + propagationCarrier: extractSqsCarrier(message), + spanKind: otel.SpanKind.CONSUMER, + recordMetric: true, + }; + } + + if (typeof rpcContext.getPattern === 'function') { + const pattern = rpcContext.getPattern(); + const routingKey = rpcContext.getChannelRef?.()?.fields?.routingKey; + return { + system: 'rabbitmq', + destination: routingKey ?? pattern, + operationName, + propagationCarrier: {}, + spanKind: otel.SpanKind.CONSUMER, + recordMetric: true, + }; + } + + 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..bb7d733 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 })), @@ -202,6 +204,7 @@ describe('OtelInstrumentation', () => { const result = otelInstrumentation.create( 'RpcController.rpcHandler', rpcContext, + 'rpc', ); expect(mockTracer.startSpan).toHaveBeenCalledWith( @@ -209,6 +212,10 @@ describe('OtelInstrumentation', () => { { kind: mockOtelApi.SpanKind.SERVER, attributes: { + 'messaging.system': 'nestjs', + 'messaging.destination.name': 'RpcController.rpcHandler', + 'messaging.operation.type': 'process', + 'messaging.operation.name': 'RpcController.rpcHandler', 'rpc.method': 'rpcHandler', 'nestjs.controller': 'RpcController', 'nestjs.handler': 'rpcHandler', @@ -216,9 +223,52 @@ 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: { + 'messaging.system': 'kafka', + 'messaging.destination.name': 'orders.created', + 'messaging.operation.type': 'process', + 'messaging.operation.name': 'OrdersController.handleOrder', + 'messaging.message.id': '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 +301,20 @@ 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) + 'messaging.system': 'nestjs', + 'messaging.destination.name': 'RpcController.unknownMethod', + 'messaging.operation.type': 'process', + 'messaging.operation.name': 'RpcController.unknownMethod', + 'rpc.method': '', 'nestjs.controller': 'RpcController', - 'nestjs.handler': 'unknown', // Falls back to 'unknown' when name is empty + 'nestjs.handler': 'unknown', }, }, {}, @@ -311,15 +365,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 + 'messaging.system': 'nestjs', + 'messaging.destination.name': 'UnknownTransaction', + 'messaging.operation.type': 'process', + 'messaging.operation.name': '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..3dade45 --- /dev/null +++ b/libs/otel-nestjs-instrumentation/test/internal/record-messaging-process-duration.spec.ts @@ -0,0 +1,52 @@ +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'; + +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, { + 'messaging.system': 'aws_sqs', + 'messaging.destination.name': 'announcements', + 'messaging.operation.type': 'process', + 'messaging.operation.name': 'announcement-dispatch', + }); + }); + + 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..9ad3466 --- /dev/null +++ b/libs/otel-nestjs-instrumentation/test/internal/resolve-rpc-messaging-metadata.spec.ts @@ -0,0 +1,107 @@ +import * as otel from '@opentelemetry/api'; +import { resolveRpcMessagingMetadata } from '../../src/internal/resolve-rpc-messaging-metadata'; +import { createMockExecutionContext } from '../test-utils'; + +describe('resolveRpcMessagingMetadata', () => { + it('returns undefined for HTTP contexts', () => { + expect( + resolveRpcMessagingMetadata(createMockExecutionContext('http')), + ).toBeUndefined(); + }); + + it('resolves Kafka metadata from KafkaContext-like RPC contexts', () => { + const context = createMockExecutionContext('rpc', { + controller: 'OrdersController', + handler: 'handleOrder', + }); + jest.spyOn(context.switchToRpc(), 'getContext').mockReturnValue({ + getTopic: () => 'orders.created', + getMessage: () => ({ + offset: '42', + headers: { + traceparent: Buffer.from( + '00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01', + ), + }, + }), + }); + + expect(resolveRpcMessagingMetadata(context)).toEqual({ + system: 'kafka', + destination: 'orders.created', + operationName: 'OrdersController.handleOrder', + messageId: '42', + propagationCarrier: { + traceparent: '00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01', + }, + spanKind: otel.SpanKind.CONSUMER, + recordMetric: true, + }); + }); + + it('resolves SQS metadata from message-shaped RPC contexts', () => { + const context = createMockExecutionContext('rpc', { + controller: 'Announcements', + handler: 'handleAnnouncement', + }); + jest.spyOn(context.switchToRpc(), 'getContext').mockReturnValue({ + getMessage: () => ({ + MessageId: 'msg-123', + MessageAttributes: { + traceparent: { + StringValue: + '00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01', + }, + }, + }), + }); + + expect(resolveRpcMessagingMetadata(context)).toEqual({ + system: 'aws_sqs', + destination: 'Announcements.handleAnnouncement', + operationName: 'Announcements.handleAnnouncement', + messageId: 'msg-123', + propagationCarrier: { + traceparent: '00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01', + }, + spanKind: otel.SpanKind.CONSUMER, + recordMetric: true, + }); + }); + + it('resolves RabbitMQ metadata from pattern-based RPC contexts', () => { + const context = createMockExecutionContext('rpc', { + controller: 'EventsController', + handler: 'handleEvent', + }); + jest.spyOn(context.switchToRpc(), 'getContext').mockReturnValue({ + getPattern: () => 'events.created', + getChannelRef: () => ({ fields: { routingKey: 'events.created' } }), + }); + + expect(resolveRpcMessagingMetadata(context)).toEqual({ + system: 'rabbitmq', + destination: 'events.created', + operationName: 'EventsController.handleEvent', + propagationCarrier: {}, + spanKind: otel.SpanKind.CONSUMER, + recordMetric: true, + }); + }); + + 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, + }); + }); +}); 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, From 2dc4cb663d17f29f9bd5551092540e679db67543 Mon Sep 17 00:00:00 2001 From: Thiago Oliveira Santos Date: Fri, 4 Sep 2026 15:59:10 -0300 Subject: [PATCH 2/4] fix(otel-nestjs-instrumentation): narrow context type for effectiveType context.getType() returns string; only assign when it is http or rpc so nest build passes strict TypeScript checking in CI. Co-authored-by: Cursor --- .../src/internal/otel-instrumentation.ts | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/libs/otel-nestjs-instrumentation/src/internal/otel-instrumentation.ts b/libs/otel-nestjs-instrumentation/src/internal/otel-instrumentation.ts index 649580c..3650bd6 100644 --- a/libs/otel-nestjs-instrumentation/src/internal/otel-instrumentation.ts +++ b/libs/otel-nestjs-instrumentation/src/internal/otel-instrumentation.ts @@ -65,7 +65,12 @@ export const otelInstrumentation = { let spanKind = otel.SpanKind.INTERNAL; // default for unknown contexts const attributes: Record = {}; - effectiveType ??= context.getType(); + if (!effectiveType) { + const contextType = context.getType(); + if (contextType === 'http' || contextType === 'rpc') { + effectiveType = contextType; + } + } if (effectiveType === 'http') { spanKind = otel.SpanKind.SERVER; From b476ee81bcb216c712eedebd8e6c4c9a9182d409 Mon Sep 17 00:00:00 2001 From: Thiago Oliveira Santos Date: Fri, 4 Sep 2026 16:06:02 -0300 Subject: [PATCH 3/4] refactor(otel-nestjs-instrumentation): reduce duplicated messaging code Extract shared consumer metadata and span-attribute helpers, and resolve RPC messaging metadata once per span to satisfy Sonar duplication limits. Co-authored-by: Cursor --- .../src/internal/otel-instrumentation.ts | 32 ++++----- .../resolve-rpc-messaging-metadata.ts | 70 +++++++++++++------ 2 files changed, 60 insertions(+), 42 deletions(-) diff --git a/libs/otel-nestjs-instrumentation/src/internal/otel-instrumentation.ts b/libs/otel-nestjs-instrumentation/src/internal/otel-instrumentation.ts index 3650bd6..64248bd 100644 --- a/libs/otel-nestjs-instrumentation/src/internal/otel-instrumentation.ts +++ b/libs/otel-nestjs-instrumentation/src/internal/otel-instrumentation.ts @@ -1,6 +1,9 @@ import { ExecutionContext } from '@nestjs/common'; import * as otel from '@opentelemetry/api'; -import { resolveRpcMessagingMetadata } from './resolve-rpc-messaging-metadata'; +import { + resolveRpcMessagingMetadata, + toMessagingSpanAttributes, +} from './resolve-rpc-messaging-metadata'; import { tracerName } from './tracer-name'; /** @@ -64,6 +67,7 @@ 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; if (!effectiveType) { const contextType = context.getType(); @@ -89,20 +93,14 @@ export const otelInstrumentation = { // Ignore request extraction errors } } else if (effectiveType === 'rpc') { - const messaging = resolveRpcMessagingMetadata(context); - if (messaging) { - spanKind = messaging.spanKind; - attributes['messaging.system'] = messaging.system; - attributes['messaging.destination.name'] = messaging.destination; - attributes['messaging.operation.type'] = 'process'; - attributes['messaging.operation.name'] = messaging.operationName; - if (messaging.messageId) { - attributes['messaging.message.id'] = messaging.messageId; - } - if (Object.keys(messaging.propagationCarrier).length > 0) { + 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, - messaging.propagationCarrier, + messagingMetadata.propagationCarrier, ); } } else { @@ -125,13 +123,9 @@ export const otelInstrumentation = { // Ignore NestJS context extraction errors } - const messaging = - effectiveType === 'rpc' - ? resolveRpcMessagingMetadata(context) - : undefined; const spanName = - messaging?.recordMetric === true - ? `process ${messaging.operationName}` + messagingMetadata?.recordMetric === true + ? `process ${messagingMetadata.operationName}` : transactionName; const span = tracer.startSpan( 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 index 3c52718..6feba55 100644 --- a/libs/otel-nestjs-instrumentation/src/internal/resolve-rpc-messaging-metadata.ts +++ b/libs/otel-nestjs-instrumentation/src/internal/resolve-rpc-messaging-metadata.ts @@ -29,6 +29,38 @@ export type RpcMessagingMetadata = { 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 extractSqsCarrier(message: SqsMessageLike): Record { const carrier: Record = {}; for (const [key, value] of Object.entries(message.MessageAttributes ?? {})) { @@ -93,42 +125,34 @@ export function resolveRpcMessagingMetadata( if (typeof rpcContext.getTopic === 'function') { const message = rpcContext.getMessage?.() as KafkaMessageLike | undefined; - return { - system: 'kafka', - destination: rpcContext.getTopic(), + return buildConsumerMetadata( + 'kafka', + rpcContext.getTopic(), operationName, - messageId: - message?.offset === undefined ? undefined : String(message.offset), - propagationCarrier: extractKafkaCarrier(message), - spanKind: otel.SpanKind.CONSUMER, - recordMetric: true, - }; + { + messageId: + message?.offset === undefined ? undefined : String(message.offset), + propagationCarrier: extractKafkaCarrier(message), + }, + ); } const message = rpcContext.getMessage?.(); if (isSqsMessage(message)) { - return { - system: 'aws_sqs', - destination: operationName, - operationName, + return buildConsumerMetadata('aws_sqs', operationName, operationName, { messageId: message.MessageId, propagationCarrier: extractSqsCarrier(message), - spanKind: otel.SpanKind.CONSUMER, - recordMetric: true, - }; + }); } if (typeof rpcContext.getPattern === 'function') { const pattern = rpcContext.getPattern(); const routingKey = rpcContext.getChannelRef?.()?.fields?.routingKey; - return { - system: 'rabbitmq', - destination: routingKey ?? pattern, + return buildConsumerMetadata( + 'rabbitmq', + routingKey ?? pattern, operationName, - propagationCarrier: {}, - spanKind: otel.SpanKind.CONSUMER, - recordMetric: true, - }; + ); } return buildGenericRpcMetadata(operationName); From 14b4d2fccdaf081500922e8f5aa36a7acce64e8d Mon Sep 17 00:00:00 2001 From: Thiago Oliveira Santos Date: Fri, 4 Sep 2026 16:12:25 -0300 Subject: [PATCH 4/4] refactor(otel-nestjs-instrumentation): dedupe messaging helpers and tests Reuse toMessagingSpanAttributes for metrics recording, consolidate carrier extraction, and table-drive transport metadata tests to satisfy Sonar. Co-authored-by: Cursor --- .../record-messaging-process-duration.ts | 15 +- .../resolve-rpc-messaging-metadata.ts | 38 ++-- .../internal/otel-instrumentation.spec.ts | 64 +++++-- .../record-messaging-process-duration.spec.ts | 18 +- .../resolve-rpc-messaging-metadata.spec.ts | 170 +++++++++--------- 5 files changed, 180 insertions(+), 125 deletions(-) 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 index 733a17d..266fa92 100644 --- a/libs/otel-nestjs-instrumentation/src/internal/record-messaging-process-duration.ts +++ b/libs/otel-nestjs-instrumentation/src/internal/record-messaging-process-duration.ts @@ -1,6 +1,9 @@ import { metrics, type Histogram } from '@opentelemetry/api'; import { tracerName } from './tracer-name'; -import type { RpcMessagingMetadata } from './resolve-rpc-messaging-metadata'; +import { + type RpcMessagingMetadata, + toMessagingSpanAttributes, +} from './resolve-rpc-messaging-metadata'; let processDurationHistogram: Histogram | undefined; @@ -25,10 +28,8 @@ export function recordMessagingProcessDuration( ): void { if (!metadata.recordMetric) return; - getProcessDurationHistogram().record(durationMs, { - 'messaging.system': metadata.system, - 'messaging.destination.name': metadata.destination, - 'messaging.operation.type': 'process', - 'messaging.operation.name': metadata.operationName, - }); + 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 index 6feba55..85bbce7 100644 --- a/libs/otel-nestjs-instrumentation/src/internal/resolve-rpc-messaging-metadata.ts +++ b/libs/otel-nestjs-instrumentation/src/internal/resolve-rpc-messaging-metadata.ts @@ -61,27 +61,41 @@ function buildConsumerMetadata( }; } -function extractSqsCarrier(message: SqsMessageLike): Record { +function normalizeCarrier( + entries: Record, + resolveValue: (value: unknown) => string | undefined, +): Record { const carrier: Record = {}; - for (const [key, value] of Object.entries(message.MessageAttributes ?? {})) { - if (value.StringValue) { - carrier[key.toLowerCase()] = value.StringValue; + 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 { - const carrier: Record = {}; - for (const [key, value] of Object.entries(message?.headers ?? {})) { - if (value === undefined) continue; - carrier[key.toLowerCase()] = Buffer.isBuffer(value) - ? value.toString('utf8') - : String(value); - } - return carrier; + 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 { 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 bb7d733..35bdebc 100644 --- a/libs/otel-nestjs-instrumentation/test/internal/otel-instrumentation.spec.ts +++ b/libs/otel-nestjs-instrumentation/test/internal/otel-instrumentation.spec.ts @@ -24,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; @@ -212,10 +236,11 @@ describe('OtelInstrumentation', () => { { kind: mockOtelApi.SpanKind.SERVER, attributes: { - 'messaging.system': 'nestjs', - 'messaging.destination.name': 'RpcController.rpcHandler', - 'messaging.operation.type': 'process', - 'messaging.operation.name': 'RpcController.rpcHandler', + ...messagingAttributes({ + system: 'nestjs', + destination: 'RpcController.rpcHandler', + operationName: 'RpcController.rpcHandler', + }), 'rpc.method': 'rpcHandler', 'nestjs.controller': 'RpcController', 'nestjs.handler': 'rpcHandler', @@ -255,11 +280,12 @@ describe('OtelInstrumentation', () => { { kind: mockOtelApi.SpanKind.CONSUMER, attributes: { - 'messaging.system': 'kafka', - 'messaging.destination.name': 'orders.created', - 'messaging.operation.type': 'process', - 'messaging.operation.name': 'OrdersController.handleOrder', - 'messaging.message.id': '7', + ...messagingAttributes({ + system: 'kafka', + destination: 'orders.created', + operationName: 'OrdersController.handleOrder', + messageId: '7', + }), 'rpc.method': 'handleOrder', 'nestjs.controller': 'OrdersController', 'nestjs.handler': 'handleOrder', @@ -308,10 +334,11 @@ describe('OtelInstrumentation', () => { { kind: mockOtelApi.SpanKind.SERVER, attributes: { - 'messaging.system': 'nestjs', - 'messaging.destination.name': 'RpcController.unknownMethod', - 'messaging.operation.type': 'process', - 'messaging.operation.name': 'RpcController.unknownMethod', + ...messagingAttributes({ + system: 'nestjs', + destination: 'RpcController.unknownMethod', + operationName: 'RpcController.unknownMethod', + }), 'rpc.method': '', 'nestjs.controller': 'RpcController', 'nestjs.handler': 'unknown', @@ -371,12 +398,11 @@ describe('OtelInstrumentation', () => { 'RpcController.rpcHandler', { kind: mockOtelApi.SpanKind.SERVER, - attributes: { - 'messaging.system': 'nestjs', - 'messaging.destination.name': 'UnknownTransaction', - 'messaging.operation.type': 'process', - 'messaging.operation.name': 'UnknownTransaction', - }, + 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 index 3dade45..bbef3eb 100644 --- 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 @@ -12,6 +12,7 @@ import { __resetMessagingProcessDurationForTests, recordMessagingProcessDuration, } from '../../src/internal/record-messaging-process-duration'; +import { toMessagingSpanAttributes } from '../../src/internal/resolve-rpc-messaging-metadata'; describe('recordMessagingProcessDuration', () => { beforeEach(() => { @@ -29,12 +30,17 @@ describe('recordMessagingProcessDuration', () => { recordMetric: true, }); - expect(mockRecord).toHaveBeenCalledWith(12.5, { - 'messaging.system': 'aws_sqs', - 'messaging.destination.name': 'announcements', - 'messaging.operation.type': 'process', - 'messaging.operation.name': 'announcement-dispatch', - }); + 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', () => { 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 index 9ad3466..a44f61d 100644 --- 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 @@ -1,7 +1,21 @@ import * as otel from '@opentelemetry/api'; -import { resolveRpcMessagingMetadata } from '../../src/internal/resolve-rpc-messaging-metadata'; +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( @@ -9,86 +23,6 @@ describe('resolveRpcMessagingMetadata', () => { ).toBeUndefined(); }); - it('resolves Kafka metadata from KafkaContext-like RPC contexts', () => { - const context = createMockExecutionContext('rpc', { - controller: 'OrdersController', - handler: 'handleOrder', - }); - jest.spyOn(context.switchToRpc(), 'getContext').mockReturnValue({ - getTopic: () => 'orders.created', - getMessage: () => ({ - offset: '42', - headers: { - traceparent: Buffer.from( - '00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01', - ), - }, - }), - }); - - expect(resolveRpcMessagingMetadata(context)).toEqual({ - system: 'kafka', - destination: 'orders.created', - operationName: 'OrdersController.handleOrder', - messageId: '42', - propagationCarrier: { - traceparent: '00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01', - }, - spanKind: otel.SpanKind.CONSUMER, - recordMetric: true, - }); - }); - - it('resolves SQS metadata from message-shaped RPC contexts', () => { - const context = createMockExecutionContext('rpc', { - controller: 'Announcements', - handler: 'handleAnnouncement', - }); - jest.spyOn(context.switchToRpc(), 'getContext').mockReturnValue({ - getMessage: () => ({ - MessageId: 'msg-123', - MessageAttributes: { - traceparent: { - StringValue: - '00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01', - }, - }, - }), - }); - - expect(resolveRpcMessagingMetadata(context)).toEqual({ - system: 'aws_sqs', - destination: 'Announcements.handleAnnouncement', - operationName: 'Announcements.handleAnnouncement', - messageId: 'msg-123', - propagationCarrier: { - traceparent: '00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01', - }, - spanKind: otel.SpanKind.CONSUMER, - recordMetric: true, - }); - }); - - it('resolves RabbitMQ metadata from pattern-based RPC contexts', () => { - const context = createMockExecutionContext('rpc', { - controller: 'EventsController', - handler: 'handleEvent', - }); - jest.spyOn(context.switchToRpc(), 'getContext').mockReturnValue({ - getPattern: () => 'events.created', - getChannelRef: () => ({ fields: { routingKey: 'events.created' } }), - }); - - expect(resolveRpcMessagingMetadata(context)).toEqual({ - system: 'rabbitmq', - destination: 'events.created', - operationName: 'EventsController.handleEvent', - propagationCarrier: {}, - spanKind: otel.SpanKind.CONSUMER, - recordMetric: true, - }); - }); - it('falls back to generic Nest RPC metadata without metrics', () => { const context = createMockExecutionContext('rpc', { controller: 'RpcController', @@ -104,4 +38,78 @@ describe('resolveRpcMessagingMetadata', () => { 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, + }); + }, + ); });