package main import ( "context" "encoding/json" "log" "os" "os/signal" "syscall" "time" "github.com/jobs-scraper/services/scraper/pipeline" "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/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...") }