internal/ and libs/ were both shared-code roots with no rule for which
one a package belonged in, so the split carried no information. Adopt a
rule the language answers by itself:
shared Go -> internal/
shared TypeScript -> libs/ (npm workspace packages)
libs/{ports,repo,server} move to internal/, leaving libs/ holding only
the two TypeScript packages (@jobs-scraper/rabbitmq-ts and
@jobs-scraper/browser-automation), which matches npm workspace
convention. internal/ is also Go's marker for code not importable from
outside the repo, which is accurate here since none of it is published.
Import paths are rewritten mechanically (jobs-scraper/libs/ ->
jobs-scraper/internal/) across 15 lines in 10 files. The 6 moved files
are pure renames with no content change. Doing this after the module
collapse in the previous commit meant no go.mod or replace-directive
edits were needed.
gofmt is applied to services/api/pkg/http/job.go, whose import group the
rewrite left out of order. Three files were already unformatted before
this refactor (internal/openai/interface.go, internal/ports/job-queries.go,
internal/repo/job-analysis-result.go) and are deliberately left alone to
keep this diff limited to the move.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
94 lines
2.3 KiB
Go
94 lines
2.3 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"log"
|
|
"os"
|
|
"os/signal"
|
|
"syscall"
|
|
|
|
"time"
|
|
|
|
"github.com/jobs-scraper/internal/domain"
|
|
"github.com/jobs-scraper/internal/infrastructure"
|
|
"github.com/jobs-scraper/internal/infrastructure/rabbitmq"
|
|
"github.com/jobs-scraper/internal/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...")
|
|
}
|