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" CvAnalyzeExchange = "cv_exchange" CvAnalyzeQueue = "cv.analyze" ) 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() } }