jobs-monorepo/services/scraper-linkedin/main.go
Elshimy Ziad Magdy Taha 5060643a3f infrastructure fixes
2025-11-03 00:29:33 +05:00

94 lines
2.3 KiB
Go

package main
import (
"context"
"encoding/json"
"log"
"os"
"os/signal"
"syscall"
"time"
"github.com/jobs-scraper/internal/pkg/domain"
"github.com/jobs-scraper/internal/pkg/infrastructure"
"github.com/jobs-scraper/internal/pkg/infrastructure/rabbitmq"
"github.com/jobs-scraper/libs/repo"
"github.com/jobs-scraper/services/scraper/pipeline"
"github.com/joho/godotenv"
)
func main() {
// Try to load .local.env first, then fallback to .env
if err := godotenv.Load(".local.env"); err != nil {
log.Println("No .local.env file found, trying .env")
if err := godotenv.Load(); err != nil {
log.Println("No .env file found, using system environment variables")
}
}
dbConfig := infrastructure.LoadConfigFromEnv()
db, err := infrastructure.NewConnection(dbConfig)
if err != nil {
log.Fatalf("Error connecting to db %v", err)
}
err = db.Ping()
if err != nil {
log.Fatalf("Error pinging db %v", err)
}
log.Println("Successfully connected to db")
rmq, err := rabbitmq.NewRabbitMQClient()
if err != nil {
log.Fatalf("Error connecting to RabbitMQ %v", err)
}
scraper := pipeline.NewScraper(pipeline.Config{
SortBy: "R",
MaxRetries: 3,
BaseDelay: 1 * time.Second,
MaxDelay: 30 * time.Second,
RequestTimeout: 30 * time.Second,
})
jobRepo := repo.NewJobRepository(db)
jobDescriptionRepo := repo.NewJobDescriptionRepository(db)
jobPipeline := pipeline.NewJobPipeline(scraper, 3, 1*time.Second) // 3 workers, 1 second rate limit
err = rmq.Subscribe(rabbitmq.LinkedInQueue, "scraper-consumer", func(data []byte) error {
ctx := context.Background()
var searchParams domain.SearchQuery
err := json.Unmarshal(data, &searchParams)
if err != nil {
log.Printf("Error unmarshaling message: %v", err)
return err
}
log.Printf("Received message: %s", string(data))
err = jobPipeline.ProcessJobsStreaming(ctx, jobRepo, jobDescriptionRepo, searchParams)
if err != nil {
log.Printf("Error processing jobs: %v", err)
return err
}
log.Printf("Successfully processed job search for: %s in %s", searchParams.Keywords, searchParams.Location)
return nil
})
if err != nil {
log.Printf("Error subscribing to LinkedIn queue: %v", err)
}
c := make(chan os.Signal, 1)
signal.Notify(c, os.Interrupt, syscall.SIGTERM)
<-c
rmq.Close()
log.Println("Shutting down scraper...")
}