Skip to content

Commit

Permalink
feat: add resendHeadersPrefix in error topic exception filter
Browse files Browse the repository at this point in the history
  • Loading branch information
KillWolfVlad committed Jan 22, 2024
1 parent 30a1cee commit cd7bdea
Show file tree
Hide file tree
Showing 3 changed files with 20 additions and 9 deletions.
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -590,6 +590,7 @@ export class UsersRetryConsumer {
new KafkaConsumerErrorTopicExceptionFilter({
retryTopicPicker: false,
errorTopicPicker: (config: ConfigDto) => config.kafka.errorTopic,
resendHeadersPrefix: "retry",
}),
)
public async onMessage(): Promise<void> {
Expand Down
23 changes: 14 additions & 9 deletions src/consumer/exceptions/kafkaConsumerErrorTopicExceptionFilter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import {
Logger,
RpcExceptionFilter,
} from "@nestjs/common";
import { IHeaders } from "kafkajs";
import { from, Observable, throwError } from "rxjs";

import {
Expand Down Expand Up @@ -103,6 +104,7 @@ export class KafkaConsumerErrorTopicExceptionFilter
host,
this.options.connectionName ?? context.connectionName,
topic,
this.options.resendHeadersPrefix,
),
);
}
Expand Down Expand Up @@ -130,6 +132,7 @@ export class KafkaConsumerErrorTopicExceptionFilter
host: ArgumentsHost,
connectionName: string,
topic: string,
resendHeadersPrefix = "original",
): Promise<void> {
const rpcHost = host.switchToRpc();
const payload: IKafkaConsumerPayload = rpcHost.getData();
Expand All @@ -139,21 +142,23 @@ export class KafkaConsumerErrorTopicExceptionFilter

this.logger.warn(`Send message to ${topic}`);

const headers: IHeaders = {
...payload.rawPayload.message.headers,
[`${resendHeadersPrefix}Topic`]: payload.rawPayload.topic,
[`${resendHeadersPrefix}Partition`]: String(payload.rawPayload.partition),
[`${resendHeadersPrefix}Offset`]: payload.rawPayload.message.offset,
[`${resendHeadersPrefix}Timestamp`]: payload.rawPayload.message.timestamp,
[`${resendHeadersPrefix}TraceId`]: context.traceId,
error: JSON.stringify(serializeError(cause)),
};

await context.kafkaCoreProducer.send(connectionName, {
topic,
messages: [
{
key: payload.rawPayload.message.key,
value: payload.rawPayload.message.value,
headers: {
...payload.rawPayload.message.headers,
originalTopic: payload.rawPayload.topic,
originalPartition: String(payload.rawPayload.partition),
originalOffset: payload.rawPayload.message.offset,
originalTimestamp: payload.rawPayload.message.timestamp,
originalTraceId: context.traceId,
error: JSON.stringify(serializeError(cause)),
},
headers,
},
],
});
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,4 +27,9 @@ export interface IKafkaConsumerErrorTopicExceptionFilterOptions {
* @deprecated Use errorTopicPicker instead
*/
readonly topicPicker?: ((...args: any[]) => string) | false;

/**
* @default "original"
*/
readonly resendHeadersPrefix?: string;
}

0 comments on commit cd7bdea

Please sign in to comment.