Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

fix:consumer processes deliveries only once #101

Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
118 changes: 60 additions & 58 deletions internal/pkg/rabbitmq/consumer.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,11 @@ package rabbitmq
import (
"context"
"fmt"
"reflect"
"runtime"
"strings"
"time"

"github.com/ahmetb/go-linq/v3"
"github.com/iancoleman/strcase"
jsoniter "github.com/json-iterator/go"
Expand All @@ -11,10 +16,6 @@ import (
"github.com/streadway/amqp"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/trace"
"reflect"
"runtime"
"strings"
"time"
)

//go:generate mockery --name IConsumer
Expand Down Expand Up @@ -104,60 +105,61 @@ func (c Consumer[T]) ConsumeMessage(msg interface{}, dependencies T) error {
}

go func() {

select {
case <-c.ctx.Done():
defer func(ch *amqp.Channel) {
err := ch.Close()
if err != nil {
c.log.Errorf("failed to close channel closed for for queue: %s", q.Name)
}
}(ch)
c.log.Infof("channel closed for for queue: %s", q.Name)
return

case delivery, ok := <-deliveries:
{
if !ok {
c.log.Errorf("NOT OK deliveries channel closed for queue: %s", q.Name)
return
}

// Extract headers
c.ctx = otel.ExtractAMQPHeaders(c.ctx, delivery.Headers)

err := c.handler(q.Name, delivery, dependencies)
if err != nil {
c.log.Error(err.Error())
}

consumedMessages = append(consumedMessages, snakeTypeName)

_, span := c.jaegerTracer.Start(c.ctx, consumerHandlerName)

h, err := jsoniter.Marshal(delivery.Headers)

if err != nil {
c.log.Errorf("Error in marshalling headers in consumer: %v", string(h))
}

span.SetAttributes(attribute.Key("message-id").String(delivery.MessageId))
span.SetAttributes(attribute.Key("correlation-id").String(delivery.CorrelationId))
span.SetAttributes(attribute.Key("queue").String(q.Name))
span.SetAttributes(attribute.Key("exchange").String(delivery.Exchange))
span.SetAttributes(attribute.Key("routing-key").String(delivery.RoutingKey))
span.SetAttributes(attribute.Key("ack").Bool(true))
span.SetAttributes(attribute.Key("timestamp").String(delivery.Timestamp.String()))
span.SetAttributes(attribute.Key("body").String(string(delivery.Body)))
span.SetAttributes(attribute.Key("headers").String(string(h)))

// Cannot use defer inside a for loop
time.Sleep(1 * time.Millisecond)
span.End()

err = delivery.Ack(false)
if err != nil {
c.log.Errorf("We didn't get a ack for delivery: %v", string(delivery.Body))
for {
select {
case <-c.ctx.Done():
defer func(ch *amqp.Channel) {
err := ch.Close()
if err != nil {
c.log.Errorf("failed to close channel closed for for queue: %s", q.Name)
}
}(ch)
c.log.Infof("channel closed for for queue: %s", q.Name)
return

case delivery, ok := <-deliveries:
{
if !ok {
c.log.Errorf("NOT OK deliveries channel closed for queue: %s", q.Name)
return
}

// Extract headers
c.ctx = otel.ExtractAMQPHeaders(c.ctx, delivery.Headers)

err := c.handler(q.Name, delivery, dependencies)
if err != nil {
c.log.Error(err.Error())
}

consumedMessages = append(consumedMessages, snakeTypeName)

_, span := c.jaegerTracer.Start(c.ctx, consumerHandlerName)

h, err := jsoniter.Marshal(delivery.Headers)

if err != nil {
c.log.Errorf("Error in marshalling headers in consumer: %v", string(h))
}

span.SetAttributes(attribute.Key("message-id").String(delivery.MessageId))
span.SetAttributes(attribute.Key("correlation-id").String(delivery.CorrelationId))
span.SetAttributes(attribute.Key("queue").String(q.Name))
span.SetAttributes(attribute.Key("exchange").String(delivery.Exchange))
span.SetAttributes(attribute.Key("routing-key").String(delivery.RoutingKey))
span.SetAttributes(attribute.Key("ack").Bool(true))
span.SetAttributes(attribute.Key("timestamp").String(delivery.Timestamp.String()))
span.SetAttributes(attribute.Key("body").String(string(delivery.Body)))
span.SetAttributes(attribute.Key("headers").String(string(h)))

// Cannot use defer inside a for loop
time.Sleep(1 * time.Millisecond)
span.End()

err = delivery.Ack(false)
if err != nil {
c.log.Errorf("We didn't get a ack for delivery: %v", string(delivery.Body))
}
}
}
}
Expand Down
Loading