- Reorganized project structure into services/ and apps/ directories for better separation of concerns - Added comprehensive CI/CD pipeline with GitHub Actions for testing and Docker builds - Created .dockerignore file to optimize container builds - Updated Makefile with new targets for each service and application - Added detailed README with architecture overview, setup instructions and development guidelines - Moved cron-analyzer to dedicate
248 lines
5.7 KiB
Go
248 lines
5.7 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"
|
|
GlassDoorQueue = "scraper.glassdoor"
|
|
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",
|
|
GlassDoorQueue: "scraper.glassdoor",
|
|
}
|
|
|
|
// 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()
|
|
}
|
|
}
|