191 lines
4.9 KiB
Go
191 lines
4.9 KiB
Go
package scan
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"net/http"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// Runner scans all enabled companies with an adapter, filtered by title.
|
|
type Runner struct {
|
|
DB *sql.DB
|
|
Client *http.Client
|
|
Matcher TitleMatcher
|
|
Workers int
|
|
}
|
|
|
|
// NewRunner returns a Runner with sensible defaults.
|
|
func NewRunner(db *sql.DB, cfg *Config) *Runner {
|
|
return &Runner{
|
|
DB: db,
|
|
Client: HTTPClient(15 * time.Second),
|
|
Matcher: NewTitleMatcher(cfg.TitleFilter),
|
|
Workers: 6,
|
|
}
|
|
}
|
|
|
|
// Run scans every enabled, adapter-supported company. Results stream on the
|
|
// returned channel; the channel closes when all companies are done.
|
|
func (r *Runner) Run(ctx context.Context, companies []Company) <-chan ScanResult {
|
|
return r.RunAll(ctx, companies, nil)
|
|
}
|
|
|
|
// RunAll scans companies AND aggregators concurrently. Aggregators are
|
|
// keyword/tag/feed sources (Remotive, RemoteOK, USAJobs, generic RSS) that
|
|
// aren't tied to one company; each aggregator appears as its own ScanResult
|
|
// row named after the aggregator label.
|
|
func (r *Runner) RunAll(ctx context.Context, companies []Company, aggregators []Aggregator) <-chan ScanResult {
|
|
out := make(chan ScanResult, len(companies)+len(aggregators))
|
|
workers := r.Workers
|
|
if workers < 1 {
|
|
workers = 1
|
|
}
|
|
sem := make(chan struct{}, workers)
|
|
var wg sync.WaitGroup
|
|
for _, c := range companies {
|
|
if !c.Enabled {
|
|
continue
|
|
}
|
|
info, ok := DetectAdapter(c)
|
|
if !ok {
|
|
out <- ScanResult{Company: c.Name, Err: errUnsupported}
|
|
continue
|
|
}
|
|
wg.Add(1)
|
|
sem <- struct{}{}
|
|
go func(c Company, info AdapterInfo) {
|
|
defer wg.Done()
|
|
defer func() { <-sem }()
|
|
start := time.Now()
|
|
res := r.scanOne(ctx, c, info)
|
|
res.DurationMs = time.Since(start).Milliseconds()
|
|
out <- res
|
|
}(c, info)
|
|
}
|
|
for _, a := range aggregators {
|
|
if !a.Enabled {
|
|
continue
|
|
}
|
|
info, ok := DetectAggregator(a)
|
|
if !ok {
|
|
out <- ScanResult{Company: a.Name, Err: errUnsupported}
|
|
continue
|
|
}
|
|
wg.Add(1)
|
|
sem <- struct{}{}
|
|
go func(a Aggregator, info AdapterInfo) {
|
|
defer wg.Done()
|
|
defer func() { <-sem }()
|
|
start := time.Now()
|
|
// Aggregators reuse scanOne via a synthetic Company so the
|
|
// title-filter and record path stay shared.
|
|
res := r.scanOne(ctx, Company{Name: a.Name}, info)
|
|
res.DurationMs = time.Since(start).Milliseconds()
|
|
out <- res
|
|
}(a, info)
|
|
}
|
|
go func() {
|
|
wg.Wait()
|
|
close(out)
|
|
}()
|
|
return out
|
|
}
|
|
|
|
func (r *Runner) scanOne(ctx context.Context, c Company, info AdapterInfo) ScanResult {
|
|
res := ScanResult{Company: c.Name}
|
|
jobs, err := Fetch(ctx, r.Client, c, info)
|
|
if err != nil {
|
|
res.Err = err
|
|
return res
|
|
}
|
|
res.Found = len(jobs)
|
|
|
|
// Title filter.
|
|
kept := make([]Job, 0, len(jobs))
|
|
for _, j := range jobs {
|
|
if r.Matcher(j.Title) {
|
|
kept = append(kept, j)
|
|
}
|
|
}
|
|
res.Filtered = res.Found - len(kept)
|
|
|
|
// Single transaction per company — batches all kept jobs in one commit
|
|
// to eliminate the N+1 tx pattern under parallel workers.
|
|
newJobs, err := r.recordJobs(ctx, kept)
|
|
if err != nil {
|
|
res.Err = err
|
|
return res
|
|
}
|
|
res.New = len(newJobs)
|
|
res.Duplicate = len(kept) - res.New
|
|
res.NewJobs = newJobs
|
|
return res
|
|
}
|
|
|
|
// recordJobs inserts a batch of jobs under a single transaction and returns
|
|
// the subset that were newly inserted (not duplicates).
|
|
func (r *Runner) recordJobs(ctx context.Context, jobs []Job) ([]Job, error) {
|
|
if len(jobs) == 0 {
|
|
return nil, nil
|
|
}
|
|
tx, err := r.DB.BeginTx(ctx, nil)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer tx.Rollback() // safe: no-op after explicit commit
|
|
|
|
insertJob, err := tx.PrepareContext(ctx,
|
|
`INSERT OR IGNORE INTO jobs (url, company, title, source) VALUES (?, ?, ?, ?)`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer insertJob.Close()
|
|
|
|
upsertScan, err := tx.PrepareContext(ctx,
|
|
`INSERT INTO scanned_urls (url, company, title, source, status)
|
|
VALUES (?, ?, ?, ?, 'added')
|
|
ON CONFLICT(url) DO UPDATE SET last_seen = CURRENT_TIMESTAMP`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer upsertScan.Close()
|
|
|
|
var newJobs []Job
|
|
for _, j := range jobs {
|
|
result, err := insertJob.ExecContext(ctx, j.URL, j.Company, j.Title, j.Source)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
rows, err := result.RowsAffected()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if rows > 0 {
|
|
newJobs = append(newJobs, j)
|
|
}
|
|
if _, err := upsertScan.ExecContext(ctx, j.URL, j.Company, j.Title, j.Source); err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
if err := tx.Commit(); err != nil {
|
|
return nil, err
|
|
}
|
|
return newJobs, nil
|
|
}
|
|
|
|
// errUnsupported is a sentinel error surfaced to the user when a company lacks
|
|
// an adapter-compatible API. It carries through to Scan UI for visibility.
|
|
var errUnsupported = &unsupportedErr{}
|
|
|
|
type unsupportedErr struct{}
|
|
|
|
func (u *unsupportedErr) Error() string { return "no tier-1 adapter available" }
|
|
|
|
// IsUnsupported lets UI code show a different row style for unsupported companies.
|
|
func IsUnsupported(err error) bool {
|
|
_, ok := err.(*unsupportedErr)
|
|
return ok
|
|
}
|