jobs-monorepo/internal/pipeline/job_pipeline_workers.go
Elshimy Ziad Magdy Taha 7bf35ccb6b refactor
2025-09-29 22:44:30 +05:00

65 lines
1.4 KiB
Go

package pipeline
import (
"context"
"log"
"sync"
"github.com/jobs-scraper/internal/domain"
)
func GetJobs(context context.Context, scraperService *Scraper, searchQuery domain.SearchQuery) <-chan domain.Job {
jobChan := make(chan domain.Job, 100)
go func() {
defer close(jobChan)
if err := scraperService.ScrapeLinkedInJobsStreaming(context, 10, jobChan, searchQuery); err != nil {
log.Printf("Error scraping jobs: %v", err)
}
}()
return jobChan
}
func GetJobDescription(context context.Context, scraperService *Scraper, jobChan <-chan domain.Job, numWorkers int) <-chan domain.JobWithDescription {
jobDescriptionChan := make(chan domain.JobWithDescription, 100)
var wg sync.WaitGroup
// Start worker goroutines
for range numWorkers {
wg.Add(1)
go func() {
defer wg.Done()
for job := range jobChan {
select {
case <-context.Done():
return
default:
}
// Scrape job description
if jd, jc, err := scraperService.ScrapeJobDescriptionWithContext(context, job); err != nil {
log.Printf("Error scraping jobs: %v", err)
} else {
jobDescriptionChan <- domain.JobWithDescription{
Job: job,
JobDescription: domain.JobDescription{
JobID: job.ID,
Description: jd,
Criteria: jc,
},
}
}
}
}()
}
// Close the output channel when all workers are done
go func() {
wg.Wait()
close(jobDescriptionChan)
}()
return jobDescriptionChan
}