package pipeline import ( "context" "log" "github.com/jobs-scraper/shared/domain" ) func GetJobs(context context.Context, scraperService *Scraper, searchQuery domain.SearchQuery) <-chan domain.Job { jobChan := make(chan domain.Job) go func() { defer close(jobChan) if err := scraperService.ScrapeLinkedInJobsStreaming(context, searchQuery.NumPages, 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) go func() { defer close(jobDescriptionChan) for job := range jobChan { select { case <-context.Done(): return default: } if jd, jc, err := scraperService.ScrapeJobDescriptionWithContext(context, job); err != nil { log.Printf("Error scraping job description: %v", err) } else { jobDescriptionChan <- domain.JobWithDescription{ Job: job, JobDescription: domain.JobDescription{ JobID: job.ID, Description: jd, Criteria: jc, }, } } } }() return jobDescriptionChan // unless you have proxies, you will get rate limited by linkedin, // but if you do have you can use this code below, // it will work faster (if u pass more than one worker) // + dont forget to use buffered channels // 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 job description: %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) // }() }