Skip to content

WithCorrelation

WithCorrelation(): DualMethodDecorator

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

Queue-agnostic MethodDecorator that propagates a correlation ID into the pipeline’s async-local correlation context.

Place it under any transport or scheduling decorator a framework adds. It does not replace or interfere with them.

It extracts the correlation ID from the method arguments (via path or extract) and runs the method inside runWithCorrelationId, which keeps the current id, or makes one with correlationSource.create(), when no ID is found.

⚠️ Array payloads: When using the default dot-path extraction, the first argument must be an object (e.g. a Bull Job). If it is an array, the dot-path cannot resolve and the decorator logs a warning at runtime. For transports that deliver an array as the first argument, use the extract option:

@WithCorrelation({ extract: (items) => items?.[0]?.correlationId })

Inside the method body, read the active ID with:

import { getCorrelationId } from '@cqrs-ddd/pipeline-correlation';
const id = getCorrelationId();

DualMethodDecorator

class Consumers {
// BullMQ job (default path: data.correlationId)
@WithCorrelation()
async sendEmail(job: Job) {
await mailer.send(job.data);
}
// Another key in the job data
@WithCorrelation('data.x-request-id')
async sendSms(job: Job) {}
// RabbitMQ message (amqplib)
@WithCorrelation({ extract: (message) => message.properties.correlationId })
async onUserCreated(message: ConsumeMessage) {}
// Kafka message (kafkajs)
@WithCorrelation({
extract: ({ message }) => message.headers?.['x-correlation-id']?.toString(),
})
async onOrderPlaced(payload: EachMessagePayload) {}
// PostgreSQL LISTEN/NOTIFY
@WithCorrelation({ path: 'correlationId' })
async onNotification(notification: { correlationId?: string }) {}
// Scheduled job: no id in the arguments, so a new one is made
@WithCorrelation()
async hourlySync() {}
}

WithCorrelation(path): DualMethodDecorator

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

Queue-agnostic MethodDecorator that propagates a correlation ID into the pipeline’s async-local correlation context.

Place it under any transport or scheduling decorator a framework adds. It does not replace or interfere with them.

It extracts the correlation ID from the method arguments (via path or extract) and runs the method inside runWithCorrelationId, which keeps the current id, or makes one with correlationSource.create(), when no ID is found.

⚠️ Array payloads: When using the default dot-path extraction, the first argument must be an object (e.g. a Bull Job). If it is an array, the dot-path cannot resolve and the decorator logs a warning at runtime. For transports that deliver an array as the first argument, use the extract option:

@WithCorrelation({ extract: (items) => items?.[0]?.correlationId })

Inside the method body, read the active ID with:

import { getCorrelationId } from '@cqrs-ddd/pipeline-correlation';
const id = getCorrelationId();

string

DualMethodDecorator

class Consumers {
// BullMQ job (default path: data.correlationId)
@WithCorrelation()
async sendEmail(job: Job) {
await mailer.send(job.data);
}
// Another key in the job data
@WithCorrelation('data.x-request-id')
async sendSms(job: Job) {}
// RabbitMQ message (amqplib)
@WithCorrelation({ extract: (message) => message.properties.correlationId })
async onUserCreated(message: ConsumeMessage) {}
// Kafka message (kafkajs)
@WithCorrelation({
extract: ({ message }) => message.headers?.['x-correlation-id']?.toString(),
})
async onOrderPlaced(payload: EachMessagePayload) {}
// PostgreSQL LISTEN/NOTIFY
@WithCorrelation({ path: 'correlationId' })
async onNotification(notification: { correlationId?: string }) {}
// Scheduled job: no id in the arguments, so a new one is made
@WithCorrelation()
async hourlySync() {}
}

WithCorrelation(options): DualMethodDecorator

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

Queue-agnostic MethodDecorator that propagates a correlation ID into the pipeline’s async-local correlation context.

Place it under any transport or scheduling decorator a framework adds. It does not replace or interfere with them.

It extracts the correlation ID from the method arguments (via path or extract) and runs the method inside runWithCorrelationId, which keeps the current id, or makes one with correlationSource.create(), when no ID is found.

⚠️ Array payloads: When using the default dot-path extraction, the first argument must be an object (e.g. a Bull Job). If it is an array, the dot-path cannot resolve and the decorator logs a warning at runtime. For transports that deliver an array as the first argument, use the extract option:

@WithCorrelation({ extract: (items) => items?.[0]?.correlationId })

Inside the method body, read the active ID with:

import { getCorrelationId } from '@cqrs-ddd/pipeline-correlation';
const id = getCorrelationId();

CorrelationDecoratorOptions

DualMethodDecorator

class Consumers {
// BullMQ job (default path: data.correlationId)
@WithCorrelation()
async sendEmail(job: Job) {
await mailer.send(job.data);
}
// Another key in the job data
@WithCorrelation('data.x-request-id')
async sendSms(job: Job) {}
// RabbitMQ message (amqplib)
@WithCorrelation({ extract: (message) => message.properties.correlationId })
async onUserCreated(message: ConsumeMessage) {}
// Kafka message (kafkajs)
@WithCorrelation({
extract: ({ message }) => message.headers?.['x-correlation-id']?.toString(),
})
async onOrderPlaced(payload: EachMessagePayload) {}
// PostgreSQL LISTEN/NOTIFY
@WithCorrelation({ path: 'correlationId' })
async onNotification(notification: { correlationId?: string }) {}
// Scheduled job: no id in the arguments, so a new one is made
@WithCorrelation()
async hourlySync() {}
}