MCPcopy Create free account
hub / github.com/codealong-dev/ride-sharing / ConsumeMessages

Method ConsumeMessages

shared/messaging/rabbitmq.go:53–121  ·  view source on GitHub ↗
(queueName string, handler MessageHandler)

Source from the content-addressed store, hash-verified

51type MessageHandler func(context.Context, amqp.Delivery) error
52
53func (r *RabbitMQ) ConsumeMessages(queueName string, handler MessageHandler) error {
54 // Set prefetch count to 1 for fair dispatch
55 // This tells RabbitMQ not to give more than one message to a service at a time.
56 // The worker will only get the next message after it has acknowledged the previous one.
57 err := r.Channel.Qos(
58 1, // prefetchCount: Limit to 1 unacknowledged message per consumer
59 0, // prefetchSize: No specific limit on message size
60 false, // global: Apply prefetchCount to each consumer individually
61 )
62 if err != nil {
63 return fmt.Errorf("failed to set QoS: %v", err)
64 }
65
66 msgs, err := r.Channel.Consume(
67 queueName, // queue
68 "", // consumer
69 false, // auto-ack
70 false, // exclusive
71 false, // no-local
72 false, // no-wait
73 nil, // args
74 )
75 if err != nil {
76 return err
77 }
78
79 go func() {
80 for msg := range msgs {
81 if err := tracing.TracedConsumer(msg, func(ctx context.Context, d amqp.Delivery) error {
82 log.Printf("Received a message: %s", msg.Body)
83
84 cfg := retry.DefaultConfig()
85 err := retry.WithBackoff(ctx, cfg, func() error {
86 return handler(ctx, d)
87 })
88 if err != nil {
89 log.Printf("Message processing failed after %d retries for message ID: %s, err: %v", cfg.MaxRetries, d.MessageId, err)
90
91 // Add failure context before sending to the DLQ
92 headers := amqp.Table{}
93 if d.Headers != nil {
94 headers = d.Headers
95 }
96
97 headers["x-death-reason"] = err.Error()
98 headers["x-origin-exchange"] = d.Exchange
99 headers["x-original-routing-key"] = d.RoutingKey
100 headers["x-retry-count"] = cfg.MaxRetries
101 d.Headers = headers
102
103 // Reject without requeue - message will go to the DLQ
104 _ = d.Reject(false)
105 return err
106 }
107
108 // Only Ack if the handler succeeds
109 if ackErr := msg.Ack(false); ackErr != nil {
110 log.Printf("ERROR: Failed to Ack message: %v. Message body: %s", ackErr, msg.Body)

Callers 4

ListenMethod · 0.80
ListenMethod · 0.80
ListenMethod · 0.80
ListenMethod · 0.80

Calls 3

TracedConsumerFunction · 0.92
DefaultConfigFunction · 0.92
WithBackoffFunction · 0.92

Tested by

no test coverage detected