jobs-monorepo/infrastructure/rabbitmq/rabbitmq.go
Elshimy Ziad Magdy Taha 23372eb2c4 rabbitmq
2025-10-11 13:55:38 +05:00

244 lines
5.5 KiB
Go

package rabbitmq
import (
"fmt"
"log"
"os"
"time"
amqp "github.com/rabbitmq/amqp091-go"
)
const (
ScraperExchange = "scraper_exchange"
LinkedInQueue = "scraper.linkedin"
IndeedQueue = "scraper.indeed"
BaytQueue = "scraper.bayt"
TokyoDevQueue = "scraper.tokyodev"
JapanDevQueue = "scraper.japandev"
DeadLetterExchange = "scraper_dlx"
)
type RabbitMQClient struct {
Conn *amqp.Connection
Channel *amqp.Channel
}
type MessageHandler func([]byte) error
func NewRabbitMQClient() (*RabbitMQClient, error) {
// Get RabbitMQ URL from environment or use default
rabbitmqURL := os.Getenv("RABBITMQ_URL")
if rabbitmqURL == "" {
rabbitmqURL = "amqp://guest:guest@localhost:5672/"
}
// Connect to RabbitMQ with retry logic
var conn *amqp.Connection
var err error
for i := range 5 {
conn, err = amqp.Dial(rabbitmqURL)
if err == nil {
break
}
log.Printf("Failed to connect to RabbitMQ (attempt %d/5): %v", i+1, err)
time.Sleep(5 * time.Second)
}
if err != nil {
return nil, fmt.Errorf("error connecting to RabbitMQ after 5 attempts: %v", err)
}
// Create a channel
ch, err := conn.Channel()
if err != nil {
conn.Close()
return nil, fmt.Errorf("error creating channel: %v", err)
}
client := &RabbitMQClient{
Conn: conn,
Channel: ch,
}
// Setup exchanges and queues
if err := client.setupInfrastructure(); err != nil {
client.Close()
return nil, fmt.Errorf("error setting up infrastructure: %v", err)
}
log.Println("Successfully connected to RabbitMQ and set up infrastructure")
return client, nil
}
func (r *RabbitMQClient) setupInfrastructure() error {
// Declare dead letter exchange
err := r.Channel.ExchangeDeclare(
DeadLetterExchange,
"direct",
true, // durable
false, // auto-delete
false, // internal
false, // no-wait
nil, // arguments
)
if err != nil {
return fmt.Errorf("error declaring dead letter exchange: %v", err)
}
// Declare main exchange
err = r.Channel.ExchangeDeclare(
ScraperExchange,
"topic",
true, // durable
false, // auto-delete
false, // internal
false, // no-wait
nil, // arguments
)
if err != nil {
return fmt.Errorf("error declaring exchange: %v", err)
}
// Define queues with their routing keys
queues := map[string]string{
LinkedInQueue: "scraper.linkedin",
IndeedQueue: "scraper.indeed",
BaytQueue: "scraper.bayt",
TokyoDevQueue: "scraper.tokyodev",
JapanDevQueue: "scraper.japandev",
}
// Declare queues with dead letter exchange and TTL
for queueName, routingKey := range queues {
dlqName := queueName + "_dlq"
_, err := r.Channel.QueueDeclare(
dlqName,
true, // durable
false, // delete when unused
false, // exclusive
false, // no-wait
nil, // arguments
)
if err != nil {
return fmt.Errorf("error declaring dead letter queue %s: %v", dlqName, err)
}
err = r.Channel.QueueBind(
dlqName,
queueName, // routing key is the original queue name
DeadLetterExchange,
false,
nil,
)
if err != nil {
return fmt.Errorf("error binding dead letter queue %s: %v", dlqName, err)
}
_, err = r.Channel.QueueDeclare(
queueName,
true, // durable
false, // delete when unused
false, // exclusive
false, // no-wait
amqp.Table{
"x-dead-letter-exchange": DeadLetterExchange,
"x-dead-letter-routing-key": queueName,
"x-message-ttl": int64(24 * time.Hour / time.Millisecond), // 24 hours TTL
"x-max-retries": 3,
},
)
if err != nil {
return fmt.Errorf("error declaring queue %s: %v", queueName, err)
}
err = r.Channel.QueueBind(
queueName,
routingKey,
ScraperExchange,
false,
nil,
)
if err != nil {
return fmt.Errorf("error binding queue %s: %v", queueName, err)
}
}
return nil
}
// Publish publishes a message to a specific routing key
func (r *RabbitMQClient) Publish(routingKey, exchange string, data []byte) error {
log.Printf("Publishing to exchange=%s, routingKey=%s, data=%s", ScraperExchange, routingKey, string(data))
err := r.Channel.Publish(
ScraperExchange,
routingKey,
false,
false,
amqp.Publishing{
ContentType: "application/json",
Body: data,
DeliveryMode: amqp.Persistent,
Timestamp: time.Now().UTC(),
},
)
if err != nil {
log.Printf("Error publishing: %v", err)
return fmt.Errorf("error publishing message: %v", err)
}
return nil
}
// Subscribe creates a subscription to a specific queue with production-ready error handling
func (r *RabbitMQClient) Subscribe(queueName, consumerName string, handler MessageHandler) error {
err := r.Channel.Qos(
1, // prefetch count
0, // prefetch size
false, // global
)
if err != nil {
return fmt.Errorf("error setting QoS: %v", err)
}
// Start consuming messages
msgs, err := r.Channel.Consume(
queueName,
consumerName,
false, // auto-ack
false, // exclusive
false, // no-local
false, // no-wait
nil, // args
)
if err != nil {
return fmt.Errorf("error starting consumer: %v", err)
}
go func() {
for msg := range msgs {
// Call the handler
err := handler(msg.Body)
if err != nil {
log.Printf("Error processing message: %v", err)
// Reject and requeue the message (will go to DLX after max retries)
msg.Nack(false, true)
} else {
msg.Ack(false)
}
}
}()
log.Printf("Successfully subscribed to queue: %s ", queueName)
return nil
}
// Close closes the RabbitMQ connection and channel
func (r *RabbitMQClient) Close() {
if r.Channel != nil {
r.Channel.Close()
}
if r.Conn != nil {
r.Conn.Close()
}
}