From 975286c2ace46a0d8f3af459f75378de99d5b9ba Mon Sep 17 00:00:00 2001 From: ankitjain235 Date: Mon, 7 Sep 2020 19:13:36 +0530 Subject: [PATCH] Add headers to Kafka MQT error topics (#1701) Co-authored-by: Vishal --- pkg/mqtrigger/messageQueue/kafka/kafka.go | 32 +++++++++++++++++++---- 1 file changed, 27 insertions(+), 5 deletions(-) diff --git a/pkg/mqtrigger/messageQueue/kafka/kafka.go b/pkg/mqtrigger/messageQueue/kafka/kafka.go index 8541a821..0214927f 100644 --- a/pkg/mqtrigger/messageQueue/kafka/kafka.go +++ b/pkg/mqtrigger/messageQueue/kafka/kafka.go @@ -47,6 +47,9 @@ func init() { var ( // Need to use raw string to support escape sequence for - & . chars validKafkaTopicName = regexp.MustCompile(`^[a-zA-Z0-9][a-zA-Z0-9\-\._]*[a-zA-Z0-9]$`) + + // Map for ErrorTopic messages to maintain recycle counter + errorMessageMap = make(map[string]int) ) type ( @@ -293,14 +296,32 @@ func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *fv1.Me zap.String("trigger", trigger.ObjectMeta.Name), zap.String("body", string(body))) + generateErrorHeaders := func(errString string) []sarama.RecordHeader { + var errorHeaders []sarama.RecordHeader + if kafka.version.IsAtLeast(sarama.V0_11_0_0) { + if count, ok := errorMessageMap[errString]; ok { + errorMessageMap[errString] = count + 1 + } else { + errorMessageMap[errString] = 1 + } + errorHeaders = append(errorHeaders, sarama.RecordHeader{Key: []byte("MessageSource"), Value: []byte(trigger.Spec.Topic)}) + errorHeaders = append(errorHeaders, sarama.RecordHeader{Key: []byte("RecycleCounter"), Value: []byte(strconv.Itoa(errorMessageMap[errString]))}) + } + return errorHeaders + } + if err != nil { + errorString := string("request body error: " + string(body)) + errorHeaders := generateErrorHeaders(errorString) errorHandler(kafka.logger, trigger, producer, url, - errors.Wrapf(err, "request body error: %v", string(body))) + errors.Wrapf(err, errorString), errorHeaders) return } if resp.StatusCode != 200 { + errorString := string("request returned failure: " + string(resp.StatusCode)) + errorHeaders := generateErrorHeaders(errorString) errorHandler(kafka.logger, trigger, producer, url, - fmt.Errorf("request returned failure: %v", resp.StatusCode)) + fmt.Errorf("request returned failure: %v", resp.StatusCode), errorHeaders) return } if len(trigger.Spec.ResponseTopic) > 0 { @@ -334,11 +355,12 @@ func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *fv1.Me consumer.MarkOffset(msg, "") // mark message as processed } -func errorHandler(logger *zap.Logger, trigger *fv1.MessageQueueTrigger, producer sarama.SyncProducer, funcUrl string, err error) { +func errorHandler(logger *zap.Logger, trigger *fv1.MessageQueueTrigger, producer sarama.SyncProducer, funcUrl string, err error, errorTopicHeaders []sarama.RecordHeader) { if len(trigger.Spec.ErrorTopic) > 0 { _, _, e := producer.SendMessage(&sarama.ProducerMessage{ - Topic: trigger.Spec.ErrorTopic, - Value: sarama.StringEncoder(err.Error()), + Topic: trigger.Spec.ErrorTopic, + Value: sarama.StringEncoder(err.Error()), + Headers: errorTopicHeaders, }) if e != nil { logger.Error("failed to publish message to error topic",