Files
2026-07-06 11:05:50 -04:00

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
}