Skip to content

CorrelationFrom

const CorrelationFrom: object

Defined in: decorators/with-correlation.decorator.ts:277

Pre-built extraction presets for WithCorrelation.

Instead of writing a custom extract function for each transport, use a preset to extract the correlation ID from the transport’s native metadata (AMQP properties, Kafka headers, NATS headers, gRPC metadata).

For transports that have no headers (Bull/BullMQ, PostgreSQL NOTIFY), use the bare @WithCorrelation() which reads from the data payload.

readonly amqp: () => CorrelationDecoratorOptions

RabbitMQ (AMQP) — extracts from ctx.getMessage().properties.correlationId.

AMQP has a first-class correlationId property on the message. On the producer side, set it with:

channel.publish(exchange, key, buffer, { correlationId: getCorrelationId() });

CorrelationDecoratorOptions

readonly grpc: (key) => CorrelationDecoratorOptions

gRPC — extracts from gRPC Metadata (second argument).

In NestJS gRPC handlers the signature is (data, metadata, call).

string = DEFAULT_CORRELATION_HEADER

Metadata key. Defaults to 'x-correlation-id'.

CorrelationDecoratorOptions

readonly kafka: (header) => CorrelationDecoratorOptions

Kafka — extracts from ctx.getMessage().headers[key].

Kafka headers are Buffer | string | undefined. The value is .toString()-ed. On the producer side, set it with:

{ headers: correlationHeaders() }

string = DEFAULT_CORRELATION_HEADER

Header key. Defaults to 'x-correlation-id'.

CorrelationDecoratorOptions

readonly nats: (header) => CorrelationDecoratorOptions

NATS — extracts from ctx.getHeaders().get(key).

string = DEFAULT_CORRELATION_HEADER

Header key. Defaults to 'x-correlation-id'.

CorrelationDecoratorOptions

// RabbitMQ
@MessagePattern('user.created')
@WithCorrelation(CorrelationFrom.amqp())
async handle(@Payload() data: any, @Ctx() ctx: RmqContext) { }
// Kafka
@EventPattern('order.placed')
@WithCorrelation(CorrelationFrom.kafka())
async handle(@Payload() data: any, @Ctx() ctx: KafkaContext) { }
// NATS
@MessagePattern('user.created')
@WithCorrelation(CorrelationFrom.nats())
async handle(@Payload() data: any, @Ctx() ctx: NatsContext) { }
// gRPC
@GrpcMethod('UsersService', 'FindOne')
@WithCorrelation(CorrelationFrom.grpc())
async findOne(data: any, metadata: Metadata) { }