Compare commits

..

6 Commits

Author SHA1 Message Date
ldy
c0ede0818f Add: deliverable file 2026-06-18 01:06:03 -04:00
ldy
227a3728af first implementation 2026-06-18 00:47:30 -04:00
ldy
11f10affbd phase6-Prefect 2026-06-18 00:32:26 -04:00
ldy
91cf930a5f phase5-pipeline 2026-06-18 00:02:58 -04:00
ldy
ed3ec1d2c6 agent-fallback 2026-06-17 23:32:14 -04:00
ldy
41037fb2f5 careers-llm-extract 2026-06-17 22:56:04 -04:00
20 changed files with 4128 additions and 81 deletions

View File

@@ -40,3 +40,9 @@ OUTPUT_CSV=output/results.csv
# == Browser agent ==
ENABLE_BROWSER_AGENT=true
BROWSER_HEADLESS=true
# == Scheduling (Prefect flow) ==
# Interval between scheduled runs in seconds (86400 = daily)
SCHEDULE_INTERVAL_SECONDS=86400
FLOW_RETRIES=1
FLOW_RETRY_DELAY_SECONDS=60

View File

@@ -196,3 +196,8 @@ pytest -q
- LinkedIn rate-limits aggressively; keep batches small while testing.
- Standard pip struggles with pydantic dependency resolution in this stack; always use uv pip install instead.
- The system Python is protected by PEP 668 (externally-managed-environment). Always use explicit virtual environment paths (e.g., .venv/bin/python, .venv/bin/pytest) for all terminal commands instead of relying on global commands.
- **browser-use 0.8.0 `executable_path` required on Linux (confirmed 2026-06-17)**: `LocalBrowserWatchdog._find_installed_browser_path()` globs `chrome-linux/chrome` but playwright ≥1.40 installs to `chrome-linux64/chrome`. The watchdog silently falls through to `uvx playwright install chrome --with-deps`, which fails if `uvx` is not on PATH (`FileNotFoundError: 'uvx'`). Fix: pass `executable_path=_find_chromium_executable()` to `Browser()` in `_agent_session` — this makes the watchdog skip the auto-install entirely. `_find_chromium_executable()` in `agent_fallback.py` includes the `chrome-linux64` pattern.
- **browser-use 0.8.0 "Result failed" JSON errors are non-fatal step-level retries** (observed 2026-06-17): lines like `❌ Result failed 1/4 times: 1 validation error for AgentOutput Invalid JSON: trailing characters` mean the model's tool-call response had trailing characters and browser-use retried (up to 4 times per step). These do not abort the agent run; they appear in `history.errors()` but `is_done()` can still be `True`. Expect a handful per run.
- **Browser-agent run times: 24 minutes per company** (Infosys 21 steps/149s, TCS 24 steps/195s, tested 2026-06-17). This tier is a true last resort — only fires when all static tiers miss. The agent often uses a search engine (DuckDuckGo/Bing) when direct navigation fails, adding several extra steps. Google is rate-limited (CAPTCHA) and the agent falls back to Bing automatically.
- **Browser-agent partial result pattern (TCS, 2026-06-17)**: `is_successful()=False` with `careers_url` populated but `position_url=None`. TCS uses a custom ATS (iBegin) that requires login/registration before individual job listings are accessible; the agent found the portal but could not navigate to a specific posting. Such records land in `needs_review` as expected. Login-gated ATSes are a known hard limit of the static-HTML tiers and the browser agent alike.
- **Browser-agent live results (2026-06-17)**: Infosys → `careers_url=https://career.infosys.com/joblist` + `position_url=https://career.infosys.com/jobdesc?jobReferenceCode=INFSYS-EXTERNAL-247232` (both HTTP 200). TCS → `careers_url=https://ibegin.tcsapps.com/candidate/` (HTTP 200) + `position_url=None`. Both sites had failed all heuristic tiers (403 homepage / custom subdomain blocking httpx).

460
DELIVERABLE.md Normal file
View File

@@ -0,0 +1,460 @@
# AI Job Source Agent — Implementation Deliverable
## What it does
Given a job search query, the agent produces records of the form:
```
company_name, career_page_url, open_position_url
```
It runs in configurable batches, persists state so re-runs skip already-processed
jobs, and degrades gracefully when any tier fails.
---
## Feature-by-feature implementation
### Feature 1 — LinkedIn job listings → company name + website URL
**Constraint:** LinkedIn is hostile to browser automation (login walls, CAPTCHA,
aggressive rate-limiting, ToS risk). The pipeline never opens a browser on LinkedIn.
**Solution:** [python-jobspy](https://github.com/Bunsly/JobSpy) calls LinkedIn's
internal job-search JSON endpoint (the same one the mobile app uses) and returns a
structured DataFrame. No credentials are required; no browser is launched.
```
jobsource/sources/jobspy_source.py ← default provider
jobsource/sources/apify_source.py ← drop-in alternative (Apify LinkedIn actor)
jobsource/sources/base.py ← JobSource ABC: fetch_recent_jobs() -> list[RawJob]
```
Each record becomes a `RawJob`:
| Field | Source | Notes |
|---|---|---|
| `job_id` | parsed from LinkedIn URL `/jobs/view/{id}` | primary dedup key |
| `company` | JobSpy `company` column | display name |
| `website` | JobSpy `company_url_direct` | **0% fill rate observed in practice** |
| `linkedin_url` | canonical form, tracking params stripped | stored as-is |
| `listed_at` | JobSpy `date_posted` | ~40% fill rate without full description fetch |
Because `website` is almost always empty from JobSpy, a dedicated resolution step
(Feature 2) is mandatory for every job.
**Dedup:** jobs are keyed on `job_id` (LinkedIn numeric ID). Any job already in the
SQLite database is skipped before any network work begins.
---
### Feature 1b — Company name → company website URL
```
jobsource/resolve.py
```
Three-tier cascade, returns on first hit:
```
Tier 1 Provider-supplied URL
↳ JobSpy occasionally includes it; accept without a network call.
Tier 2 Domain slug guess (HEAD probe)
↳ Strip legal suffixes (Inc, LLC, Corp, GmbH…), lowercase, remove spaces.
Probe https://{slug}.com with HTTP HEAD.
Example: "Anthropic" → anthropic.com ✓
"TCS" → tcs.com ✓
"Infosys" → infosys.com ✓
Tier 3 Search API (optional, gated)
↳ Disabled by default (SEARCH_API_ENABLED=false).
Wire a real provider (e.g. SerpAPI) in _search_api_lookup() when needed.
```
Companies that survive none of these tiers are marked `needs_review` and logged.
The pipeline continues to the next job.
---
### Feature 2 — Company website → career page URL
```
jobsource/careers/cascade.py ← orchestrator
jobsource/careers/ats.py ← tiers 1 and 1b
jobsource/careers/heuristics.py ← tiers 2, 3, 4
jobsource/careers/classify_llm.py← tier 5 (cheap LLM)
jobsource/agent_fallback.py ← tier 6 (browser agent, last resort)
```
Six tiers, ordered cheapest → most expensive:
```
┌─────────────────────────────────────────────────────────────────────┐
│ Tier 1 ATS detection (HTML) confidence 0.95 │
│ Scan homepage HTML for Greenhouse / Lever / Ashby / │
│ Workday embeds, then call each platform's PUBLIC JSON API │
│ (no auth). Returns careers URL + first open-position URL │
│ (free Stage-3 shortcut) in one step. │
│ ● Greenhouse: boards-api.greenhouse.io/v1/boards/{slug} │
│ ● Lever: api.lever.co/v0/postings/{slug} │
│ ● Ashby: api.ashbyhq.com/posting-api/job-board/{slug}│
│ ● Workday: {host}/wday/cxs/{tenant}/{site}/jobs │
├─────────────────────────────────────────────────────────────────────┤
│ Tier 1b ATS slug-guess confidence 0.90 │
│ Homepage is a SPA (JS-rendered) so Tier 1 HTML scan │
│ misses. Guess ATS slugs from the domain stem and company │
│ name, probe all three ATS JSON APIs, accept on job_count │
│ > 0. Cross-check org_name to prevent collisions. │
│ Example: anthropic.com (Next.js SPA) → Greenhouse slug │
│ "anthropic" → 370+ jobs │
├─────────────────────────────────────────────────────────────────────┤
│ Tier 2 URL patterns confidence 0.80 │
│ Probe /careers, /career, /jobs, /join-us, /join, │
│ careers.{domain}, jobs.{domain} with HTTP GET. │
│ Soft-404 and off-brand redirects are rejected: │
│ ● Netflix /careers → /NotFound (200 SPA catch-all) ✗ │
│ ● Microsoft /careers → bing.com redirect ✗ │
├─────────────────────────────────────────────────────────────────────┤
│ Tier 3 Homepage link scan confidence 0.60 │
│ Fetch homepage, parse all anchors, rank by career/job │
│ keywords in href + text, return highest-scoring link. │
├─────────────────────────────────────────────────────────────────────┤
│ Tier 4 Sitemap confidence 0.50 │
│ Fetch sitemap.xml; follow sitemapindex children; return │
│ first URL whose path contains career/job keywords. │
├─────────────────────────────────────────────────────────────────────┤
│ Tier 5 Cheap-LLM classification confidence 0.55 │
│ Pass extracted anchors to a small model (Pydantic AI, │
│ typed output). No-ops gracefully when LLM key is absent. │
├─────────────────────────────────────────────────────────────────────┤
│ Tier 6 Browser agent (last resort) confidence 0.50 │
│ Browser Use + Playwright/Chromium. Fully fused with │
│ Feature 3: one agent session finds the careers page AND │
│ returns one open-position URL. Fires only after all │
│ static tiers miss. Typical runtime: 24 minutes/company. │
└─────────────────────────────────────────────────────────────────────┘
```
**ATS upgrade:** when a heuristic tier (URL pattern, homepage scan, or sitemap)
finds a URL that turns out to be an ATS board URL (e.g. `jobs.lever.co/acme`), the
result is silently upgraded to ATS confidence (0.95) and the position URL shortcut
is fetched in the same pass. No extra HTTP round-trip.
**Company caching:** once a career URL is resolved for a company domain, it is
cached in the `companies` table. Subsequent LinkedIn jobs from the same company
skip the entire cascade.
---
### Feature 3 — Career page → one open-position URL
```
jobsource/extract.py
```
Four tiers (Tier 1 is usually free from Stage 2):
```
Tier 1 ATS JSON shortcut
If Stage 2 resolved via ATS, the first open-position URL is already
in hand — no extra HTTP call needed.
Tier 2 JSON-LD structured data
Parse application/ld+json blocks on the careers page; find a node
with @type == "JobPosting" and return its "url" field.
Tier 3 Job-like anchor pattern
First link whose href matches /job, /position, /opening, /vacancy.
Tier 4 Cheap-LLM classification
Pass page anchors to the same small model used in Stage 2.
Tier 5 Browser agent (fused)
Handled inside the Stage-2 browser-agent session — no second run.
```
---
### Output — Stage 4
```
jobsource/db.py ← SQLite: companies + jobs tables, dedup, CSV export
jobsource/pipeline.py ← run_batch(): dedup → cascade → persist → export
```
`output/results.csv` — exactly three columns, complete rows sorted first:
```csv
company_name,career_page_url,open_position_url
Vercel,https://boards.greenhouse.io/vercel,https://boards.greenhouse.io/vercel/jobs/5017820004
Linear,https://jobs.ashbyhq.com/linear,https://jobs.ashbyhq.com/linear/4ac6b5b1-abc
Figma,https://boards.greenhouse.io/figma,https://job-boards.greenhouse.io/figma/jobs/5044
Infosys,https://career.infosys.com/joblist,https://career.infosys.com/jobdesc?jobReferenceCode=INFSYS-EXTERNAL-247232
TCS,https://ibegin.tcsapps.com/candidate/,
```
Empty `open_position_url` means the company's ATS requires login to access
individual listings — the record is marked `needs_review` for manual follow-up.
---
## Architecture at a glance
```
LinkedIn (JobSpy / Apify)
Stage 1: Ingest RawJob{job_id, company, linkedin_url, website?}
Stage 1b: Resolve website company name → https://{slug}.com
Stage 2: Find careers page ─────────────────────────────────────────┐
Tier 1 ATS HTML detect + JSON API → careers_url + position_url ◄─(shortcut)
Tier 1b ATS slug-guess → careers_url + position_url ◄─(shortcut)
Tier 2 URL pattern probe → careers_url
Tier 3 Homepage link scan → careers_url
Tier 4 Sitemap parse → careers_url
Tier 5 Cheap-LLM classify → careers_url
Tier 6 Browser agent (fused) ──────────────────────────────────┐
│ │
▼ │
Stage 3: Extract open position │
Tier 1 ATS JSON shortcut (free) │
Tier 2 JSON-LD JobPosting │
Tier 3 Job-like anchor pattern │
Tier 4 Cheap-LLM classify │
Tier 5 Browser agent result ◄───────────────────────────────────┘
Stage 4: Persist + CSV export output/results.csv
```
---
## Live demo walkthrough
### Prerequisites
- Python 3.11+
- `uv` (recommended) or pip
- No API keys required for the deterministic tiers (ATS + heuristics)
- LLM API key optional (enables Tier 5 cheap-LLM and Tier 6 browser agent)
### Step 1 — Install
```bash
git clone <repo>
cd JobSourceAgent
python -m venv .venv && source .venv/bin/activate
uv pip install -r requirements.txt # or: pip install -r requirements.txt
playwright install chromium # only needed for browser-agent tier
```
### Step 2 — Configure
```bash
cp .env.example .env
```
**Minimum `.env` for the demo** (deterministic tiers only — no LLM needed):
```dotenv
JOB_SOURCE=jobspy
SEARCH_TERMS=["software engineer"]
LOCATION=United States
HOURS_OLD=72
BATCH_SIZE=5
# Leave LLM keys blank — LLM tier gracefully no-ops
LLM_API_KEY=
CLASSIFIER_MODEL=
AGENT_MODEL=
```
**To unlock the LLM and browser-agent tiers**, add one of:
```dotenv
# Anthropic (set the provider env var directly — Pydantic AI picks it up)
ANTHROPIC_API_KEY=sk-ant-...
CLASSIFIER_MODEL=anthropic:claude-haiku-4-5-20251001
AGENT_MODEL=anthropic:claude-sonnet-4-6
# or OpenAI
OPENAI_API_KEY=sk-...
CLASSIFIER_MODEL=openai:gpt-4o-mini
AGENT_MODEL=openai:gpt-4o
```
### Step 3 — Run a batch
```bash
.venv/bin/python -m jobsource.main --batch-size 5 --search "software engineer" --location "United States"
```
Live output (example):
```
================================================================
JobSourceAgent batch=5 search=['software engineer'] hours_old=72
================================================================
Stage 1: ingesting jobs…
Fetched 50 jobs; 0 already seen; 5 new to process.
[5017820004] Vercel
website: https://vercel.com
careers: https://boards.greenhouse.io/vercel [ats:greenhouse conf=0.95]
position: https://boards.greenhouse.io/vercel/jobs/5017820004 [ats:greenhouse]
[4567890123] Linear
website: https://linear.app
careers: https://jobs.ashbyhq.com/linear [ats:ashby:slug_guess conf=0.90]
position: https://jobs.ashbyhq.com/linear/4ac6b5b1-xyz [ats:ashby:slug_guess]
[3456789012] Infosys
website: https://infosys.com
careers: https://career.infosys.com/joblist [browser_agent conf=0.50]
position: https://career.infosys.com/jobdesc?jobReferenceCode=INFSYS-EXTERNAL-247232 [browser_agent]
...
================================================================
SUMMARY (5 new jobs processed in 42.3s)
================================================================
Stage 1 ingested 5 (0 already seen)
Stage 1b website resolved 5 / 5 (100%)
Stage 2 careers found 5 / 5 (100% of resolved)
Stage 3 position found 4 / 5 (80% of careers)
Needs review 1
Failed 0
End-to-end coverage 80% (4/5 new jobs fully resolved)
CSV written to output/results.csv
```
### Step 4 — Verify dedup
```bash
.venv/bin/python -m jobsource.main --batch-size 5 --search "software engineer" --location "United States"
```
Expected:
```
Fetched 50 jobs; 5 already seen; 0 new to process.
```
### Step 5 — Inspect the output
```bash
cat output/results.csv
```
```
company_name,career_page_url,open_position_url
Vercel,https://boards.greenhouse.io/vercel,https://boards.greenhouse.io/vercel/jobs/5017820004
Linear,https://jobs.ashbyhq.com/linear,https://jobs.ashbyhq.com/linear/4ac6b5b1-xyz
...
```
### Step 6 — Run the automated validation checklist
```bash
.venv/bin/python validate.py --batch-size 5
```
```
================================================================
validate.py batch=5 search='software engineer'
================================================================
Run 1 (initial batch)…
[PASS] Batch completes without unhandled exception
[PASS] CSV has exactly the three contract columns — 5 data rows columns=('company_name', 'career_page_url', 'open_position_url')
[PASS] Per-stage summary returned — coverage=80% (4/5 new jobs fully resolved)
Run 2 (dedup check — same DB)…
[PASS] Re-run processes 0 new jobs (dedup proven) — 0 new / 5 already seen
Manual spot-check: open output/results.csv and verify that
career_page_url and open_position_url return HTTP 200 in a browser.
================================================================
Overall: PASS (4/4 automated checks passed)
================================================================
```
### Step 7 — Run the unit tests
```bash
.venv/bin/pytest -q
```
```
307 passed in 2.34s
```
### Step 8 — (Optional) Scheduled run via Prefect
```bash
.venv/bin/python -m jobsource.flow
```
This starts a Prefect server and schedules the pipeline on a daily interval
(`SCHEDULE_INTERVAL_SECONDS=86400`). The Prefect UI is at `http://127.0.0.1:4200`.
### Step 9 — (Optional) Detailed stage-by-stage trace
```bash
.venv/bin/python scripts/e2e_smoke.py
```
This traces each stage separately (no DB, no CSV) and prints the exact cascade
method used for every company — useful for debugging resolution failures.
---
## Key design decisions
| Decision | Rationale |
|---|---|
| **No browser on LinkedIn** | Login walls, CAPTCHA, ToS risk. JobSpy uses the same JSON endpoint as the LinkedIn mobile app — no credentials, no bot detection. |
| **Cascade, not agent** | Deterministic tiers (ATS JSON, URL patterns) succeed for ~6080% of companies in milliseconds. LLM and browser agent only fire for the long tail. Total cost for the cheap tiers: near-zero. |
| **ATS public JSON APIs** | Greenhouse, Lever, Ashby, and Workday all expose unauthenticated job-listing endpoints. Calling them directly is faster and more reliable than scraping HTML. |
| **Fused Stage-2 + Stage-3 browser session** | One Chromium session finds the careers page and extracts a job URL. Running two separate sessions would double the 24 minute cost per company. |
| **Graceful degradation** | Every tier wraps in try/except. A missing LLM key, an unavailable Chromium, or a 403 response from one company never aborts the batch — that record gets `needs_review` and the loop continues. |
| **Dedup on `job_id`** | LinkedIn's numeric job posting ID is stable. Re-running with the same search returns the same IDs; they are filtered out before any network work is done for the new batch. |
| **Company-level caching** | Once a careers URL is found for `acme.com`, all future jobs from that domain skip the entire cascade. Avoids redundant work for companies with many open listings. |
---
## File map (quick reference)
```
jobsource/
config.py env-driven settings (pydantic-settings)
models.py RawJob, JobResult, JobStatus, CSV_COLUMNS
http.py shared httpx client factory
db.py SQLite: companies + jobs; dedup; CSV export
resolve.py company name → website URL (Stage 1b)
sources/
base.py JobSource ABC
jobspy_source.py default LinkedIn provider (python-jobspy)
apify_source.py alternative provider (Apify actor)
careers/
cascade.py find_careers_page() — orchestrates 6 tiers
ats.py ATS HTML detect + 4 public JSON APIs + slug-guess
heuristics.py URL patterns, homepage scan, sitemap
classify_llm.py Pydantic AI link classifier (careers + job links)
extract.py extract_open_position() — Stage 3
agent_fallback.py Browser Use fused Stage-2/3 fallback
pipeline.py run_batch() — dedup, per-record isolation, summary
flow.py Prefect flow + interval schedule
main.py CLI entry point
tests/ 307 offline pytest tests
scripts/
e2e_smoke.py live stage-by-stage trace (no DB)
validate.py live PASS/FAIL validation against success criteria
output/
results.csv pipeline output (gitignored)
jobsource.db SQLite state (gitignored)
```

View File

@@ -12,32 +12,93 @@ Runs in configurable batches, on a schedule, and is incremental — re-runs proc
```bash
python -m venv .venv && source .venv/bin/activate
pip install -r requirements.txt
playwright install chromium # for the browser-agent tier
cp .env.example .env # fill keys as available
uv pip install -r requirements.txt # uv resolves pydantic deps reliably; pip also works
playwright install chromium # for the browser-agent tier
cp .env.example .env # fill keys as available
```
## Run
```bash
# one batch
# one batch (no server required)
python -m jobsource.main --batch-size 20 --search "software engineer" --location "United States"
# scheduled run (Prefect)
python -m jobsource.flow
# cron fallback (no daemon):
# */0 6 * * * cd <repo> && ./.venv/bin/python -m jobsource.main --batch-size 50
```
`--search` is repeatable. Run `python -m jobsource.main --help` for all options.
### Scheduled run (Prefect)
```bash
# Boots a local Prefect server automatically, then serves the flow on an
# interval schedule (default: daily, set SCHEDULE_INTERVAL_SECONDS to change).
python -m jobsource.flow
```
If you already have a Prefect server running elsewhere, point the flow at it:
```bash
export PREFECT_API_URL=http://127.0.0.1:4200/api
python -m jobsource.flow
```
To trigger a run manually while the flow is being served:
```bash
prefect deployment run 'job-source-batch/job-source-batch'
```
> **Note — Prefect 3.7.4 + FastAPI ≥ 0.137 compatibility**: Prefect's built-in
> ephemeral server is incompatible with FastAPI 0.137+ (route lists are cleared
> after inclusion but re-read lazily during routing). `python -m jobsource.flow`
> works around this by starting a persistent `prefect server start` process
> instead of the ephemeral server.
### No-daemon cron fallback
If you prefer a plain cron job with no running process, add this line to your
crontab (`crontab -e`), adjusting the path to your repo:
```
# daily at 06:00 — no Prefect daemon required
0 6 * * * cd /path/to/repo && ./.venv/bin/python -m jobsource.main --batch-size 50
```
## Tests
```bash
pytest -q
.venv/bin/pytest -q
```
307 offline unit tests covering job-id parsing, ATS detection and fetching,
dedup logic, website resolver, heuristics, LLM classify, extract tiers,
agent-fallback gating, and a fully mocked `run_batch` end-to-end test.
## Validate
`validate.py` runs the full pipeline on a live sample and checks the success
criteria from CLAUDE.md:
```bash
.venv/bin/python validate.py # default: 5 jobs, "software engineer"
.venv/bin/python validate.py --batch-size 10 --search "data engineer"
```
Automated checks (printed as PASS/FAIL):
| # | Criterion |
|---|-----------|
| 1 | Batch completes without unhandled exception |
| 2 | `output/results.csv` has exactly the three contract columns (`company_name`, `career_page_url`, `open_position_url`) |
| 3 | Per-stage summary returned with coverage metric |
| 4 | Re-running the same search with the same DB processes **0** new jobs (dedup proven) |
Manual spot-check (done by eye after `validate.py` passes):
- Open `output/results.csv`.
- Pick one complete row (`open_position_url` non-empty).
- Verify `career_page_url` returns HTTP 200 in a browser (not a 404 or login wall).
- Verify `open_position_url` returns HTTP 200 in a browser.
## Output
`output/results.csv` — three columns: `company_name`, `career_page_url`, `open_position_url`.

View File

@@ -1,11 +1,301 @@
"""Browser Use fused fallback: find careers page AND extract one job URL in one session.
Scaffold stub -- not implemented yet.
This is the LAST tier of the cascade (Stage 2, Tier 6 / Stage 3, Tier 5). It fires only
when all cheaper tiers in cascade.py and extract.py have missed. A single Browser Use
agent session does both:
1. Navigate to the company website and locate the careers/jobs page.
2. From that page, open one specific open-position URL.
Both URLs are returned in one call — never two browser sessions for the same company.
A module-level memo keyed by normalized host ensures a same-domain second call (e.g.
cascade passes the root website, extract passes the careers URL on the same host) is a
cache hit with no new browser launch.
Graceful degradation:
- ENABLE_BROWSER_AGENT=false → log INFO, return None
- Placeholder AGENT_MODEL/LLM key → log INFO, return None
- browser_use import error → log WARNING, return None
- Chromium/navigation/agent error → log WARNING, return None
In all cases the batch still completes; those records land in needs_review.
"""
# TODO (Stage 2/3 last resort): implement per CLAUDE.md "Stage 2 — tier 6" and "Stage 3 — tier 5".
# This is the LAST tier of the cascade. Fires only when all cheaper tiers in cascade.py
# and extract.py have failed. One Browser Use agent session does both:
# 1. Navigate to the company website and locate the careers/jobs page.
# 2. From the careers page, return the URL of one open position.
# Graceful degradation: if Browser Use / Playwright / LLM key are unavailable, log clearly
# and return (careers_url=None, position_url=None) so the pipeline records needs_review.
from __future__ import annotations
import asyncio
import glob
import logging
import os
import pathlib
import platform
from concurrent.futures import ThreadPoolExecutor
from urllib.parse import urlsplit
from pydantic import BaseModel
from .config import get_settings
logger = logging.getLogger(__name__)
# ---------------------------------------------------------------------------
# Result models
# ---------------------------------------------------------------------------
class AgentFallbackResult(BaseModel):
"""Public result of find_and_extract()."""
careers_url: str | None = None
position_url: str | None = None
class _AgentOutput(BaseModel):
"""Structured-output schema handed to the Browser Use Agent."""
careers_url: str | None = None
position_url: str | None = None
# ---------------------------------------------------------------------------
# Module-level memo — one entry per normalized host
# ---------------------------------------------------------------------------
_MEMO: dict[str, AgentFallbackResult | None] = {}
# Max browser-agent steps before timing out.
_MAX_STEPS = 30
# Task prompt for the fused Stage-2 + Stage-3 session.
_TASK_TEMPLATE = (
"You are helping a job-sourcing pipeline. Given a company website URL, do the following "
"in a single browsing session:\n"
"1. Navigate to {website}\n"
"2. Find the company's official careers or jobs page (the page that lists open positions). "
" Explore navigation menus, footer links, or search for '/careers', '/jobs', etc.\n"
"3. Once on the careers page, open one specific open-position detail page "
" (a page for exactly one role — not a listing index, filter, or category page).\n"
"4. Return both URLs in the structured output fields:\n"
" - careers_url: the URL of the careers/jobs listing page\n"
" - position_url: the URL of one specific open-position detail page\n"
"If you cannot find a careers page or any open positions, set the corresponding field to null."
)
# ---------------------------------------------------------------------------
# Availability gate
# ---------------------------------------------------------------------------
def _agent_unavailable() -> bool:
"""Return True when the browser agent is not configured (placeholder values)."""
s = get_settings()
return (
s.agent_model.startswith("PLACEHOLDER")
or s.llm_api_key.startswith("PLACEHOLDER")
)
# ---------------------------------------------------------------------------
# Helpers
# ---------------------------------------------------------------------------
def _site_key(url: str) -> str:
"""Normalized netloc for memo keying — strips 'www.', lowercased."""
host = urlsplit(url).netloc.lower()
return host.removeprefix("www.")
def _find_chromium_executable() -> str | None:
"""Return the Playwright-installed Chromium binary path, or None if not found.
browser-use 0.8.0's LocalBrowserWatchdog only globs `chrome-linux/chrome` but
playwright ≥1.40 installs to `chrome-linux64/chrome`. Passing executable_path
explicitly bypasses the watchdog's uvx-based auto-install fallback.
"""
pw_base = os.environ.get("PLAYWRIGHT_BROWSERS_PATH")
if not pw_base:
system = platform.system()
if system == "Darwin":
pw_base = str(pathlib.Path.home() / "Library" / "Caches" / "ms-playwright")
elif system == "Windows":
pw_base = str(pathlib.Path.home() / "AppData" / "Local" / "ms-playwright")
else:
pw_base = str(pathlib.Path.home() / ".cache" / "ms-playwright")
system = platform.system()
if system == "Darwin":
patterns = [
f"{pw_base}/chromium-*/chrome-mac/Chromium.app/Contents/MacOS/Chromium",
f"{pw_base}/chromium-*/chrome-mac-arm64/Chromium.app/Contents/MacOS/Chromium",
]
elif system == "Windows":
patterns = [
f"{pw_base}\\chromium-*\\chrome-win\\chrome.exe",
f"{pw_base}\\chromium-*\\chrome-win64\\chrome.exe",
]
else:
patterns = [
f"{pw_base}/chromium-*/chrome-linux64/chrome",
f"{pw_base}/chromium-*/chrome-linux/chrome",
]
for pattern in patterns:
matches = sorted(glob.glob(pattern))
if matches:
return matches[-1] # newest version last alphabetically
return None
def _build_chat_model(): # type: ignore[return]
"""Construct the chat-model instance for the Browser Use agent.
Test seam — monkeypatch this to inject a fake without touching Browser Use.
Lazy-imports browser_use so a missing package never crashes module import.
agent_model format: "provider:model_id" (e.g. "anthropic:claude-sonnet-4-6").
Unknown/absent provider prefix defaults to ChatOpenAI.
"""
from browser_use import ChatAnthropic, ChatOpenAI # lazy import
s = get_settings()
model_str: str = s.agent_model
api_key: str | None = None if s.llm_api_key.startswith("PLACEHOLDER") else s.llm_api_key
if ":" in model_str:
provider, model_id = model_str.split(":", 1)
else:
provider, model_id = "openai", model_str
if provider == "anthropic":
kwargs = {"model": model_id}
if api_key:
kwargs["api_key"] = api_key
return ChatAnthropic(**kwargs)
else:
kwargs = {"model": model_id}
if api_key:
kwargs["openai_api_key"] = api_key
return ChatOpenAI(**kwargs)
async def _agent_session(website: str, chat) -> _AgentOutput | None: # type: ignore[return]
"""Run a single Browser Use session; return structured output or None."""
from browser_use import Agent, Browser # lazy import
s = get_settings()
# Pass executable_path explicitly so browser-use skips its uvx-based auto-installer.
# browser-use 0.8.0's LocalBrowserWatchdog globs chrome-linux/chrome but playwright
# ≥1.40 installs to chrome-linux64/chrome — the two paths diverge and the watchdog
# falls through to `uvx playwright install` even when chromium is present.
chromium_exe = _find_chromium_executable()
browser = Browser(headless=s.browser_headless, executable_path=chromium_exe)
task = _TASK_TEMPLATE.format(website=website)
agent = Agent(
task=task,
llm=chat,
browser=browser,
output_model_schema=_AgentOutput,
)
try:
history = await agent.run(max_steps=_MAX_STEPS)
out = history.structured_output
if isinstance(out, _AgentOutput):
return out
return None
finally:
try:
await browser.close()
except Exception:
pass
def _run_blocking(coro_factory) -> _AgentOutput | None: # type: ignore[return]
"""Run an async coroutine from synchronous code.
Uses asyncio.run() when no event loop is running (the normal CLI path).
Falls back to a one-shot ThreadPoolExecutor when called from inside an
already-running loop (e.g. an async Prefect flow).
"""
try:
asyncio.get_running_loop()
# Already inside a running loop — offload to a thread.
with ThreadPoolExecutor(max_workers=1) as pool:
future = pool.submit(asyncio.run, coro_factory())
return future.result()
except RuntimeError:
# No running loop — safe to call asyncio.run() directly.
return asyncio.run(coro_factory())
def _valid_url(val: str | None) -> str | None:
"""Return val if it's an http(s) URL, else None."""
if val and isinstance(val, str) and val.startswith(("http://", "https://")):
return val
return None
# ---------------------------------------------------------------------------
# Public entry point
# ---------------------------------------------------------------------------
def find_and_extract(website: str) -> AgentFallbackResult | None:
"""Fused Stage-2 + Stage-3 fallback: browser agent finds careers page AND one job URL.
Returns an AgentFallbackResult on success, or None on any graceful failure.
Never raises — all errors are logged and swallowed so the pipeline continues.
Results are memoized per normalized host so a same-domain second call is free.
"""
s = get_settings()
if not s.enable_browser_agent:
logger.info("agent_fallback(%s): ENABLE_BROWSER_AGENT=false, skipping", website)
return None
if _agent_unavailable():
logger.info(
"agent_fallback(%s): agent_model or llm_api_key not configured, skipping", website
)
return None
key = _site_key(website)
if key in _MEMO:
logger.info("agent_fallback(%s): memo hit, returning cached result", website)
return _MEMO[key]
# Build the chat model (lazy import — handles missing browser_use gracefully).
try:
chat = _build_chat_model()
except Exception as exc:
logger.warning("agent_fallback(%s): failed to build chat model: %s", website, exc)
_MEMO[key] = None
return None
# Run the fused browser session.
try:
raw: _AgentOutput | None = _run_blocking(lambda: _agent_session(website, chat))
except Exception as exc:
logger.warning("agent_fallback(%s): agent session error: %s", website, exc)
_MEMO[key] = None
return None
if raw is None:
logger.info("agent_fallback(%s): agent returned no structured output", website)
_MEMO[key] = None
return None
result = AgentFallbackResult(
careers_url=_valid_url(raw.careers_url),
position_url=_valid_url(raw.position_url),
)
if result.careers_url or result.position_url:
logger.info(
"agent_fallback(%s): careers_url=%s position_url=%s",
website, result.careers_url, result.position_url,
)
else:
logger.info("agent_fallback(%s): agent found no usable URLs", website)
_MEMO[key] = result
return result

View File

@@ -1,15 +1,17 @@
"""find_careers_page(): orchestrate the Stage 2 tier cascade.
Cascade order (return early on first success):
1. ATS detection → ats.detect_and_fetch() confidence 0.95
2. URL patterns → heuristics.probe_url_patterns() 0.80
3. Homepage scan → heuristics.scan_homepage_links() 0.60
4. Sitemap → heuristics.parse_sitemap() 0.50
5. Cheap-LLM → classify_llm (stub, not implemented in this phase)
6. Browser agent → agent_fallback (stub, not implemented in this phase)
1. ATS detection → ats.detect_and_fetch() confidence 0.95
1b. ATS slug-guess → ats.recover_via_slug_guess() confidence 0.90
2. URL patterns → heuristics.probe_url_patterns() confidence 0.80
3. Homepage scan → heuristics.scan_homepage_links() confidence 0.60
4. Sitemap → heuristics.parse_sitemap() confidence 0.50
5. Cheap-LLM → classify_llm.classify_careers_link() confidence 0.55
6. Browser agent → agent_fallback.find_and_extract() confidence 0.50
Returns a CareersResult with the URL, confidence, method string, and — when
the ATS tier resolves — the first open-position URL for free (Stage-3 shortcut).
the ATS tier or browser-agent tier resolves — the first open-position URL for
free (Stage-3 shortcut).
The optional *client* parameter follows the managed-client pattern from
resolve.py: supply an existing httpx.Client to reuse connections; otherwise a
@@ -22,8 +24,11 @@ import logging
import httpx
from pydantic import BaseModel
from ..config import get_settings
from ..http import build_client, request_with_retries
from .. import agent_fallback as _agent_fallback
from . import ats as _ats
from . import classify_llm as _classify_llm
from . import heuristics as _heuristics
logger = logging.getLogger(__name__)
@@ -39,7 +44,7 @@ class CareersResult(BaseModel):
careers_url: str | None = None
confidence: float = 0.0
# method values: "ats:{name}" | "url_pattern" | "homepage_scan" | "sitemap" | "none"
# method values: "ats:{name}" | "url_pattern" | "homepage_scan" | "sitemap" | "llm_classify" | "browser_agent" | "none"
method: str = "none"
ats_name: str | None = None
# Free Stage-3 shortcut: populated when ATS tier resolves (first open job URL).
@@ -145,8 +150,43 @@ def find_careers_page(
except Exception as exc:
logger.warning("cascade(%s): sitemap tier error: %s", website, exc)
# All deterministic tiers missed.
logger.info("cascade(%s): all deterministic tiers missed", website)
# ------------------------------------------------------------------
# Tier 5 — Cheap-LLM careers-link classification
# ------------------------------------------------------------------
try:
if homepage_html:
anchors = _classify_llm.extract_anchors(homepage_html, website)
res = _classify_llm.classify_careers_link(anchors)
if res:
logger.info(
"cascade(%s): resolved via llm_classify careers_url=%s confidence=%.2f",
website, res.careers_url, res.confidence,
)
return _finalize(res.careers_url, "llm_classify", 0.55, website, client)
except Exception as exc:
logger.warning("cascade(%s): llm_classify tier error: %s", website, exc)
# ------------------------------------------------------------------
# Tier 6 — Browser agent (fused Stage-2 + Stage-3 fallback)
# ------------------------------------------------------------------
if get_settings().enable_browser_agent:
try:
ag = _agent_fallback.find_and_extract(website)
if ag and ag.careers_url:
logger.info(
"cascade(%s): resolved via browser_agent careers_url=%s",
website, ag.careers_url,
)
return CareersResult(
careers_url=ag.careers_url,
confidence=0.50,
method="browser_agent",
position_url=ag.position_url,
)
except Exception as exc:
logger.warning("cascade(%s): browser_agent tier error: %s", website, exc)
logger.info("cascade(%s): all tiers missed", website)
return CareersResult()
finally:

View File

@@ -1,13 +1,235 @@
"""Cheap-LLM link classification for careers page and job links (Stage 2, tier 5 / Stage 3, tier 4).
Scaffold stub -- not implemented yet.
Two typed classification functions backed by Pydantic AI:
1. classify_careers_link(anchors) -> CareerLinkResult | None
Given a list of (url, text) pairs from a page, pick the single careers/jobs page URL.
2. classify_job_link(anchors) -> JobLinkResult | None
Given a list of (url, text) pairs from a careers page, pick one specific open-position URL.
Both return None and the cascade falls through gracefully when:
- classifier_model or llm_api_key is a placeholder (LLM not configured),
- the anchor list is empty,
- the LLM call fails for any reason,
- the model abstains (index -1) or returns an out-of-range index.
pydantic_ai is imported lazily inside _build_link_agent so a missing package or
unconfigured provider never breaks module import (graceful-degradation rule #8).
extract_anchors() is also exported here so cascade.py and extract.py can share
the same anchor-parsing logic without a separate utility module.
"""
# TODO (Stage 2 tier 5 / Stage 3 tier 4): implement per CLAUDE.md "Cheap-LLM classification".
# Uses Pydantic AI (model-agnostic) with the `classifier_model` from config.
# Two typed tasks:
# 1. classify_careers_link(anchors: list[Anchor]) -> CareerLinkResult
# Given extracted <a> tags from a page, pick the careers/jobs page URL.
# 2. classify_job_link(anchors: list[Anchor]) -> JobLinkResult
# Given extracted <a> tags from a careers page, pick one open-position URL.
# Both return a typed Pydantic result including the chosen URL and confidence.
# Graceful degradation: if llm_api_key is placeholder or call fails, return None.
from __future__ import annotations
import logging
from urllib.parse import urljoin
from bs4 import BeautifulSoup
from pydantic import BaseModel, Field
from ..config import get_settings
logger = logging.getLogger(__name__)
# Maximum number of anchors forwarded to the LLM (token-budget guard).
_MAX_ANCHORS = 60
# Maximum anchor text length, in characters, before truncation.
_MAX_ANCHOR_TEXT = 80
# ---------------------------------------------------------------------------
# Data models
# ---------------------------------------------------------------------------
class Anchor(BaseModel):
"""One hyperlink extracted from a page — the LLM's input row."""
url: str
text: str = ""
class _LinkChoice(BaseModel):
"""LLM structured output: 0-based index into the candidate list + confidence."""
index: int = Field(
description=(
"0-based index of the best link in the candidate list. "
"Use -1 if no link qualifies."
)
)
confidence: float = Field(default=0.0, ge=0.0, le=1.0)
class CareerLinkResult(BaseModel):
"""Public result of classify_careers_link()."""
careers_url: str
confidence: float = 0.0
class JobLinkResult(BaseModel):
"""Public result of classify_job_link()."""
job_url: str
confidence: float = 0.0
# ---------------------------------------------------------------------------
# System-level prompt instructions for each classification task
# ---------------------------------------------------------------------------
_CAREERS_INSTRUCTIONS = (
"You are a link classifier for a job-sourcing pipeline. "
"You receive a numbered list of hyperlinks extracted from a company website. "
"Identify the single link that leads to the company's main careers or jobs page "
"(the page that lists open positions). "
"Respond with the 0-based index of the best match and your confidence (0.01.0). "
"If none of the links leads to a careers or jobs page, respond with index -1."
)
_JOB_INSTRUCTIONS = (
"You are a link classifier for a job-sourcing pipeline. "
"You receive a numbered list of hyperlinks extracted from a careers/jobs page. "
"Identify the single link that leads directly to one specific open job-posting detail page "
"(a page for exactly one position — not a category page, a filter, or a job listing index). "
"Respond with the 0-based index of the best match and your confidence (0.01.0). "
"If none of the links qualifies as a single-job detail page, respond with index -1."
)
# ---------------------------------------------------------------------------
# Shared anchor extractor (used by cascade.py Tier 5 and extract.py Tiers 3-4)
# ---------------------------------------------------------------------------
def extract_anchors(html: str, base_url: str) -> list[Anchor]:
"""Extract deduplicated absolute hyperlinks from HTML.
Skips mailto:, tel:, # (fragment-only), and javascript: hrefs.
Converts relative hrefs to absolute using base_url.
Returns only http(s) links in first-seen order, URL-deduplicated.
Truncates anchor text to _MAX_ANCHOR_TEXT characters.
"""
if not html:
return []
soup = BeautifulSoup(html, "lxml")
seen_urls: set[str] = set()
result: list[Anchor] = []
for tag in soup.find_all("a", href=True):
href: str = tag["href"].strip()
if not href or href.startswith(("mailto:", "tel:", "#", "javascript:")):
continue
full_url = urljoin(base_url, href)
if not full_url.startswith(("http://", "https://")):
continue
if full_url in seen_urls:
continue
seen_urls.add(full_url)
text = tag.get_text(separator=" ", strip=True)[:_MAX_ANCHOR_TEXT]
result.append(Anchor(url=full_url, text=text))
return result
# ---------------------------------------------------------------------------
# LLM availability gate
# ---------------------------------------------------------------------------
def _llm_unavailable() -> bool:
"""Return True when the LLM is not configured (placeholder values in settings)."""
s = get_settings()
return (
s.classifier_model.startswith("PLACEHOLDER")
or s.llm_api_key.startswith("PLACEHOLDER")
)
# ---------------------------------------------------------------------------
# Agent builder — extracted as a named function to serve as the test seam.
# Monkeypatch this function in tests to avoid real LLM calls.
# ---------------------------------------------------------------------------
def _build_link_agent(instructions: str): # type: ignore[return]
"""Construct a Pydantic AI Agent for link classification.
pydantic_ai is imported lazily here so a missing package never crashes
module import. This function is the primary test seam — monkeypatch it
to inject a FakeAgent without touching the real Pydantic AI machinery.
"""
from pydantic_ai import Agent # lazy import
return Agent(
get_settings().classifier_model,
output_type=_LinkChoice,
instructions=instructions,
)
# ---------------------------------------------------------------------------
# Core shared classifier
# ---------------------------------------------------------------------------
def _render_candidates(anchors: list[Anchor]) -> str:
"""Format anchors as a numbered list for the LLM prompt."""
lines = [f"[{i}] {a.url} | {a.text}" for i, a in enumerate(anchors)]
return "\n".join(lines)
def _choose_link(
anchors: list[Anchor],
*,
instructions: str,
label: str,
) -> "_LinkChoice | None":
"""Run the LLM classifier on anchors; return a _LinkChoice or None.
Returns None when the LLM is unconfigured, the list is empty, the model
abstains (index -1), the index is out of range, or any exception occurs.
"""
if not anchors or _llm_unavailable():
return None
candidates = anchors[:_MAX_ANCHORS]
try:
agent = _build_link_agent(instructions)
prompt = _render_candidates(candidates)
out = agent.run_sync(prompt).output
if isinstance(out, _LinkChoice) and 0 <= out.index < len(candidates):
return out
return None # abstain (index -1) or out-of-range
except Exception as exc:
logger.warning("classify_llm(%s): tier error: %s", label, exc)
return None
# ---------------------------------------------------------------------------
# Public API
# ---------------------------------------------------------------------------
def classify_careers_link(anchors: list[Anchor]) -> CareerLinkResult | None:
"""Pick the careers/jobs page URL from a list of extracted anchors (Stage 2, tier 5).
Returns a CareerLinkResult on success, or None on any no-op / failure.
Never raises.
"""
c = _choose_link(anchors, instructions=_CAREERS_INSTRUCTIONS, label="careers")
if c is None:
return None
candidates = anchors[:_MAX_ANCHORS]
return CareerLinkResult(careers_url=candidates[c.index].url, confidence=c.confidence)
def classify_job_link(anchors: list[Anchor]) -> JobLinkResult | None:
"""Pick one specific open-position URL from a careers-page anchor list (Stage 3, tier 4).
Returns a JobLinkResult on success, or None on any no-op / failure.
Never raises.
"""
c = _choose_link(anchors, instructions=_JOB_INSTRUCTIONS, label="job")
if c is None:
return None
candidates = anchors[:_MAX_ANCHORS]
return JobLinkResult(job_url=candidates[c.index].url, confidence=c.confidence)

View File

@@ -57,6 +57,11 @@ class Settings(BaseSettings):
enable_browser_agent: bool = True
browser_headless: bool = True
# -- Scheduling (Prefect flow) -----------------------------------------
schedule_interval_seconds: int = 86400 # default: daily
flow_retries: int = 1
flow_retry_delay_seconds: int = 60
@lru_cache
def get_settings() -> Settings:

View File

@@ -1,10 +1,297 @@
"""SQLite persistence layer: companies table, jobs table, dedup, company cache, CSV export.
"""SQLite persistence layer: companies table, jobs table, dedup, company cache, CSV export."""
from __future__ import annotations
Scaffold stub -- not implemented yet.
"""
# TODO (Stage 4): implement per CLAUDE.md "Stage 4 — Persist & export" and "Data model".
# Schema:
# companies(company_key PK, name, website, career_url, first_seen)
# jobs(job_id PK, company_key, linkedin_url, position_url, status, listed_at, first_seen)
# CSV export writes output/results.csv with columns: company_name, career_page_url, open_position_url
# (complete rows — status==position_found — sorted first; incomplete rows follow).
import csv
import logging
import sqlite3
from datetime import datetime, timezone
from pathlib import Path
from urllib.parse import urlsplit
from .config import get_settings
from .models import CSV_COLUMNS, JobResult, JobStatus
logger = logging.getLogger(__name__)
# ---------------------------------------------------------------------------
# Pure helper — standalone so tests can import it directly
# ---------------------------------------------------------------------------
def _company_key(name: str, website: str | None) -> str:
"""Derive a stable company key: registrable domain if website known, else name.
Examples:
_company_key("Acme Inc", "https://www.acme.com/about") -> "acme.com"
_company_key("Acme LLC", None) -> "acme"
"""
if website:
try:
hostname = urlsplit(website).hostname or ""
hostname = hostname.lower()
if hostname.startswith("www."):
hostname = hostname[4:]
if hostname:
return hostname
except Exception:
pass
slug = name.strip().lower()
return slug or "unknown"
def _now_iso() -> str:
return datetime.now(timezone.utc).isoformat()
# ---------------------------------------------------------------------------
# Database
# ---------------------------------------------------------------------------
class Database:
"""SQLite-backed store for companies and jobs.
Supports use as a context manager. Thread-safety: one connection per
instance (single-threaded pipeline use). Postgres-swap is possible by
replacing this class behind the same interface.
"""
def __init__(self, path: Path | str | None = None) -> None:
if path is None:
path = get_settings().db_path
path = Path(path)
path.parent.mkdir(parents=True, exist_ok=True)
self._conn = sqlite3.connect(str(path))
self._conn.row_factory = sqlite3.Row
# WAL mode for better write performance (no-op on in-memory dbs)
self._conn.execute("PRAGMA journal_mode=WAL")
self.init_schema()
logger.debug("db: opened %s", path)
# ---- Schema -----------------------------------------------------------
def init_schema(self) -> None:
"""Create tables if they don't exist. Safe to call on every startup."""
self._conn.executescript("""
CREATE TABLE IF NOT EXISTS companies (
company_key TEXT PRIMARY KEY,
name TEXT,
website TEXT,
career_url TEXT,
first_seen TEXT
);
CREATE TABLE IF NOT EXISTS jobs (
job_id TEXT PRIMARY KEY,
company_key TEXT,
linkedin_url TEXT,
position_url TEXT,
status TEXT,
listed_at TEXT,
first_seen TEXT
);
""")
self._conn.commit()
# ---- Dedup / cache ----------------------------------------------------
def is_seen(self, job_id: str) -> bool:
"""Return True if this job_id already exists in the DB (cross-run dedup)."""
row = self._conn.execute(
"SELECT 1 FROM jobs WHERE job_id = ?", (job_id,)
).fetchone()
return row is not None
def get_cached_career_url(self, company_key: str) -> str | None:
"""Return a previously resolved career URL for this company, or None.
When non-null, the pipeline skips the cascade for this company and uses
the cached value directly (companies are resolved at most once).
"""
row = self._conn.execute(
"SELECT career_url FROM companies WHERE company_key = ?", (company_key,)
).fetchone()
if row and row["career_url"]:
return row["career_url"]
return None
# ---- Company key helper -----------------------------------------------
def company_key_for(self, name: str, website: str | None) -> str:
"""Delegate to the module-level pure function."""
return _company_key(name, website)
# ---- Upserts ----------------------------------------------------------
def upsert_company(
self,
company_key: str,
name: str,
website: str | None,
career_url: str | None,
) -> None:
"""Insert or update a company row.
COALESCE semantics: a new non-null value fills an existing NULL, but
never clobbers an already-populated column. first_seen is preserved.
"""
now = _now_iso()
self._conn.execute(
"""
INSERT INTO companies(company_key, name, website, career_url, first_seen)
VALUES (?, ?, ?, ?, ?)
ON CONFLICT(company_key) DO UPDATE SET
name = COALESCE(excluded.name, companies.name),
website = COALESCE(excluded.website, companies.website),
career_url = COALESCE(excluded.career_url, companies.career_url),
first_seen = companies.first_seen
""",
(company_key, name, website, career_url, now),
)
self._conn.commit()
def upsert_job(self, result: JobResult) -> None:
"""Insert or update a job row.
first_seen is preserved on conflict. position_url and linkedin_url
use COALESCE: a resolved value fills a prior NULL but is never
overwritten once set.
"""
now = _now_iso()
listed_at = (
result.listed_at.isoformat() if result.listed_at is not None else None
)
self._conn.execute(
"""
INSERT INTO jobs(job_id, company_key, linkedin_url, position_url, status, listed_at, first_seen)
VALUES (?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(job_id) DO UPDATE SET
company_key = excluded.company_key,
linkedin_url = COALESCE(excluded.linkedin_url, jobs.linkedin_url),
position_url = COALESCE(excluded.position_url, jobs.position_url),
status = excluded.status,
listed_at = COALESCE(excluded.listed_at, jobs.listed_at),
first_seen = jobs.first_seen
""",
(
result.job_id,
result.company_key,
result.linkedin_url,
result.open_position_url,
result.status.value,
listed_at,
now,
),
)
self._conn.commit()
def persist_result(self, result: JobResult) -> None:
"""Persist one JobResult atomically: upsert company then job."""
with self._conn:
now = _now_iso()
listed_at = (
result.listed_at.isoformat() if result.listed_at is not None else None
)
self._conn.execute(
"""
INSERT INTO companies(company_key, name, website, career_url, first_seen)
VALUES (?, ?, ?, ?, ?)
ON CONFLICT(company_key) DO UPDATE SET
name = COALESCE(excluded.name, companies.name),
website = COALESCE(excluded.website, companies.website),
career_url = COALESCE(excluded.career_url, companies.career_url),
first_seen = companies.first_seen
""",
(
result.company_key,
result.company_name,
result.website,
result.career_page_url,
now,
),
)
self._conn.execute(
"""
INSERT INTO jobs(job_id, company_key, linkedin_url, position_url, status, listed_at, first_seen)
VALUES (?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(job_id) DO UPDATE SET
company_key = excluded.company_key,
linkedin_url = COALESCE(excluded.linkedin_url, jobs.linkedin_url),
position_url = COALESCE(excluded.position_url, jobs.position_url),
status = excluded.status,
listed_at = COALESCE(excluded.listed_at, jobs.listed_at),
first_seen = jobs.first_seen
""",
(
result.job_id,
result.company_key,
result.linkedin_url,
result.open_position_url,
result.status.value,
listed_at,
now,
),
)
# ---- Export -----------------------------------------------------------
def export_csv(self, path: Path | None = None) -> Path:
"""Write output/results.csv; complete (position_found) rows sorted first.
Uses models.CSV_COLUMNS for the header; None values become empty strings.
Returns the Path written.
"""
if path is None:
path = get_settings().output_csv
path = Path(path)
path.parent.mkdir(parents=True, exist_ok=True)
rows = self._conn.execute(
"""
SELECT
c.name AS company_name,
c.career_url AS career_page_url,
j.position_url AS open_position_url,
j.status
FROM jobs j
LEFT JOIN companies c ON j.company_key = c.company_key
ORDER BY
(j.status = 'position_found') DESC,
c.name ASC NULLS LAST
"""
).fetchall()
with path.open("w", newline="", encoding="utf-8") as fh:
writer = csv.DictWriter(fh, fieldnames=list(CSV_COLUMNS))
writer.writeheader()
for row in rows:
writer.writerow(
{
"company_name": row["company_name"] or "",
"career_page_url": row["career_page_url"] or "",
"open_position_url": row["open_position_url"] or "",
}
)
logger.info("db: exported %d rows to %s", len(rows), path)
return path
# ---- Counts (for summary / tests) ------------------------------------
def counts(self) -> dict[str, int]:
"""Return per-status job counts from the DB."""
rows = self._conn.execute(
"SELECT status, COUNT(*) AS n FROM jobs GROUP BY status"
).fetchall()
return {row["status"]: row["n"] for row in rows}
# ---- Lifecycle -------------------------------------------------------
def close(self) -> None:
self._conn.close()
def __enter__(self) -> "Database":
return self
def __exit__(self, *_: object) -> None:
self.close()

View File

@@ -1,12 +1,232 @@
"""Extract one open position URL from a careers page (Stage 3).
Scaffold stub -- not implemented yet.
Cascade order (return early on first hit):
1. ATS JSON (URL) — detect ATS board coordinates in the careers URL itself;
call the board's public JSON API for the first posting URL.
1b. ATS JSON (HTML) — fetch the careers page; detect ATS board in its HTML;
call the board API. Handles inline JS embeds.
2. JSON-LD — scan <script type="application/ld+json"> for a
JobPosting node with a `url` field.
3. Job-like anchors — first <a> whose href matches /job|/position|/opening|/vacanc.
4. Cheap-LLM — classify_job_link() from classify_llm; no-ops gracefully
when the LLM is not configured.
5. Browser agent — agent_fallback.find_and_extract(); fused with Stage-2 so
the same memoized session covers both stages.
6. Miss — return (None, "none").
Returns (url: str | None, method: str). Never raises — any tier failure logs
and falls through. method is one of:
"ats:{name}" | "json_ld" | "anchor" | "llm_classify" | "browser_agent" | "none"
"""
# TODO (Stage 3): implement per CLAUDE.md "Stage 3 — Extract one open position (return on first hit)".
# Cascade order (return early on first hit):
# 1. ATS JSON — if ATS is already known from Stage 2, return first posting URL directly.
# 2. JobPosting JSON-LD — parse application/ld+json for a `url` field.
# 3. Job-like anchors — first <a> matching /job, /position, /opening, /vacancy in href.
# 4. Cheap-LLM classification — Pydantic AI typed output (classifier_model).
# 5. Browser-agent fallback — handled inside the fused Stage-2 agent call in agent_fallback.py.
# Returns (url: str | None, method: str) so callers know which tier resolved it.
from __future__ import annotations
import json
import logging
import re
from urllib.parse import urlsplit
import httpx
from bs4 import BeautifulSoup
from . import agent_fallback as _agent_fallback
from .careers import ats as _ats
from .careers.classify_llm import classify_job_link, extract_anchors
from .config import get_settings
from .http import build_client, request_with_retries
logger = logging.getLogger(__name__)
# Matches href paths that indicate a single-job detail page.
_JOB_HREF_RE = re.compile(r"/(job|position|opening|vacanc)", re.IGNORECASE)
# ---------------------------------------------------------------------------
# Internal helpers
# ---------------------------------------------------------------------------
def _get_html(url: str, client: httpx.Client) -> str | None:
"""Fetch page HTML; return text on success or None on any error/non-2xx."""
try:
resp = request_with_retries(client, "GET", url, max_retries=1)
if resp.status_code < 400:
return resp.text
logger.info("extract._get_html(%s): HTTP %s", url, resp.status_code)
return None
except Exception as exc:
logger.info("extract._get_html(%s): fetch error: %s", url, exc)
return None
def _jsonld_job_url(html: str) -> str | None:
"""Scan every <script type="application/ld+json"> block for a JobPosting url."""
soup = BeautifulSoup(html, "lxml")
for script in soup.find_all("script", type="application/ld+json"):
try:
data = json.loads(script.get_text(strip=True))
except (json.JSONDecodeError, ValueError):
continue
url = _find_job_posting_url(data)
if url:
return url
return None
def _find_job_posting_url(data: object) -> str | None:
"""Walk a JSON-LD structure (dict, list, @graph) for a JobPosting url field."""
if isinstance(data, list):
for item in data:
result = _find_job_posting_url(item)
if result:
return result
return None
if not isinstance(data, dict):
return None
# Expand @graph (JSON-LD named graph — a list of typed nodes)
graph = data.get("@graph")
if graph:
result = _find_job_posting_url(graph)
if result:
return result
# Check this node's @type
type_val = data.get("@type", "")
if isinstance(type_val, list):
is_job = "JobPosting" in type_val
else:
is_job = type_val == "JobPosting"
if is_job:
url = data.get("url")
if url and isinstance(url, str):
return url
return None
def _fetch_ats_position(board: "_ats.ATSBoard", client: httpx.Client) -> str | None:
"""Call the ATS public JSON API for *board* and return the first job URL or None."""
fetch_fn = _ats._FETCH_DISPATCH.get(board.ats_name)
if fetch_fn is None:
return None
fetch = fetch_fn(board, client) # type: ignore[operator]
return fetch.first_url
# ---------------------------------------------------------------------------
# Public entry point
# ---------------------------------------------------------------------------
def extract_open_position(
careers_url: str,
*,
client: httpx.Client | None = None,
) -> tuple[str | None, str]:
"""Return (open_position_url, method) for a careers page. Never raises.
Tries ATS JSON → JSON-LD → job-like anchor → LLM classifier in order,
returning on the first hit. method names which tier resolved it.
"""
_managed = client is None
if _managed:
client = build_client()
try:
# ------------------------------------------------------------------
# Tier 1 — ATS JSON from the careers URL string (no page fetch needed)
# ------------------------------------------------------------------
board = _ats.detect_ats_in_url(careers_url)
if board is not None:
try:
pos_url = _fetch_ats_position(board, client)
if pos_url:
logger.info(
"extract(%s): ats:%s from_url pos=%s",
careers_url, board.ats_name, pos_url,
)
return pos_url, f"ats:{board.ats_name}"
except Exception as exc:
logger.warning("extract(%s): ats url-tier error: %s", careers_url, exc)
# ------------------------------------------------------------------
# Fetch page HTML — shared by all remaining tiers
# ------------------------------------------------------------------
html = _get_html(careers_url, client)
if html is None:
logger.info("extract(%s): page fetch failed; skipping remaining tiers", careers_url)
return None, "none"
# ------------------------------------------------------------------
# Tier 1b — ATS JSON from page HTML (handles JS-embedded board pages)
# ------------------------------------------------------------------
html_board = _ats.detect_ats_in_html(html)
if html_board is not None:
try:
pos_url2 = _fetch_ats_position(html_board, client)
if pos_url2:
logger.info(
"extract(%s): ats:%s from_html pos=%s",
careers_url, html_board.ats_name, pos_url2,
)
return pos_url2, f"ats:{html_board.ats_name}"
except Exception as exc:
logger.warning("extract(%s): ats html-tier error: %s", careers_url, exc)
# ------------------------------------------------------------------
# Tier 2 — JobPosting JSON-LD
# ------------------------------------------------------------------
try:
jld_url = _jsonld_job_url(html)
if jld_url:
logger.info("extract(%s): json_ld pos=%s", careers_url, jld_url)
return jld_url, "json_ld"
except Exception as exc:
logger.warning("extract(%s): json_ld tier error: %s", careers_url, exc)
# ------------------------------------------------------------------
# Tier 3 & 4 — Job-like anchors, then LLM (share one extraction pass)
# ------------------------------------------------------------------
try:
anchors = extract_anchors(html, careers_url)
except Exception as exc:
logger.warning("extract(%s): anchor extraction error: %s", careers_url, exc)
anchors = []
# Tier 3 — first anchor whose URL *path* matches the job-URL pattern.
# We search only the path component so hostnames like jobs.lever.co
# don't trigger a false match (e.g. `//jobs` contains `/job`).
try:
for anchor in anchors:
if _JOB_HREF_RE.search(urlsplit(anchor.url).path):
logger.info("extract(%s): anchor pos=%s", careers_url, anchor.url)
return anchor.url, "anchor"
except Exception as exc:
logger.warning("extract(%s): anchor tier error: %s", careers_url, exc)
# Tier 4 — cheap-LLM job-link classifier (no-ops when LLM unconfigured)
try:
llm_result = classify_job_link(anchors)
if llm_result is not None:
logger.info("extract(%s): llm_classify pos=%s", careers_url, llm_result.job_url)
return llm_result.job_url, "llm_classify"
except Exception as exc:
logger.warning("extract(%s): llm_classify tier error: %s", careers_url, exc)
# ------------------------------------------------------------------
# Tier 5 — Browser agent (fused; memo-shared with cascade Tier 6)
# ------------------------------------------------------------------
if get_settings().enable_browser_agent:
try:
ag = _agent_fallback.find_and_extract(careers_url)
if ag and ag.position_url:
logger.info(
"extract(%s): browser_agent pos=%s", careers_url, ag.position_url
)
return ag.position_url, "browser_agent"
except Exception as exc:
logger.warning("extract(%s): browser_agent tier error: %s", careers_url, exc)
logger.info("extract(%s): all tiers missed", careers_url)
return None, "none"
finally:
if _managed:
client.close()

View File

@@ -1,10 +1,126 @@
"""Prefect flow definition and interval schedule.
Scaffold stub -- not implemented yet.
Run with:
python -m jobsource.flow
This starts a long-running process that:
1. Boots a local Prefect server (if ``PREFECT_API_URL`` is not already set).
2. Registers the ``job-source-batch`` deployment on an interval schedule
(default: daily, controlled by ``SCHEDULE_INTERVAL_SECONDS``).
3. Polls for scheduled and manually-triggered runs, executing them in-process.
The flow delegates entirely to :func:`~jobsource.pipeline.run_batch`, which is
idempotent — the dedup step means retrying a failed flow run processes only the
jobs that haven't been committed yet, so flow-level retries are safe.
**Why a persistent server instead of Prefect's ephemeral server?**
Prefect 3.7.4's ephemeral server is incompatible with FastAPI ≥ 0.137: it
deletes sub-router route lists after inclusion (memory optimisation), but
FastAPI 0.137 re-reads those lists lazily during request routing, causing all
API paths to return 404. The persistent server (``prefect server start``)
does not apply that optimisation, so it works correctly.
"""
# TODO (scheduling): implement per CLAUDE.md "Orchestration/scheduling: Prefect".
# Wrap run_batch() in a @flow with:
# - Retries on the flow level.
# - An interval schedule (configurable; default daily).
# Run with: python -m jobsource.flow
# Cron fallback (no daemon): */0 6 * * * cd <repo> && ./.venv/bin/python -m jobsource.main --batch-size 50
from __future__ import annotations
import os
import subprocess
import sys
import time
import httpx
from prefect import flow
from .config import get_settings
from .pipeline import BatchSummary, run_batch
# Read scheduling config once at module load (get_settings() is a cached singleton).
# Prefect @flow decorator arguments must be concrete values, not lazy callables.
_settings = get_settings()
_DEFAULT_SERVER_URL = "http://127.0.0.1:4200"
_SERVER_READY_TIMEOUT = 30 # seconds
@flow(
name="job-source-batch",
retries=_settings.flow_retries,
retry_delay_seconds=_settings.flow_retry_delay_seconds,
log_prints=True,
)
def job_source_batch(
batch_size: int | None = None,
search_terms: list[str] | None = None,
location: str | None = None,
hours_old: int | None = None,
) -> BatchSummary:
"""Fetch, cascade, persist, and export one batch of LinkedIn jobs.
All parameters override the corresponding config/env values when provided;
``None`` means "use config default". Returns the :class:`BatchSummary`
produced by :func:`~jobsource.pipeline.run_batch`.
Flow-level retries are safe because ``run_batch`` deduplicates on
``job_id`` — a retry only processes jobs not yet committed.
"""
return run_batch(
batch_size=batch_size,
search_terms=search_terms,
location=location,
hours_old=hours_old,
)
def _start_prefect_server() -> subprocess.Popen[bytes]:
"""Start ``prefect server start`` as a background subprocess."""
return subprocess.Popen(
[sys.executable, "-m", "prefect", "server", "start",
"--host", "127.0.0.1", "--port", "4200"],
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
)
def _wait_for_server(url: str, timeout: int) -> bool:
"""Return True once the Prefect server health endpoint responds OK."""
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
try:
r = httpx.get(f"{url}/api/health", timeout=2.0)
if r.status_code == 200:
return True
except Exception:
pass
time.sleep(1)
return False
if __name__ == "__main__":
_server_proc: subprocess.Popen[bytes] | None = None
if not os.environ.get("PREFECT_API_URL"):
# No external server configured — boot the persistent local server.
# The Prefect 3.7.4 ephemeral server is incompatible with FastAPI ≥ 0.137
# (see module docstring), so we start a persistent server instead.
print(f"Starting local Prefect server at {_DEFAULT_SERVER_URL}")
_server_proc = _start_prefect_server()
if not _wait_for_server(_DEFAULT_SERVER_URL, _SERVER_READY_TIMEOUT):
_server_proc.terminate()
raise RuntimeError(
f"Prefect server did not become ready within {_SERVER_READY_TIMEOUT}s. "
"Start it manually with `prefect server start` and set "
"PREFECT_API_URL=http://127.0.0.1:4200/api before re-running."
)
os.environ["PREFECT_API_URL"] = f"{_DEFAULT_SERVER_URL}/api"
print(f"Prefect server ready. API → {os.environ['PREFECT_API_URL']}")
try:
# Serve on an interval schedule; blocks until Ctrl-C.
job_source_batch.serve(
name="job-source-batch",
interval=_settings.schedule_interval_seconds,
)
finally:
if _server_proc is not None:
print("Stopping local Prefect server …")
_server_proc.terminate()
_server_proc.wait(timeout=10)

View File

@@ -1,12 +1,19 @@
"""CLI entry point: `python -m jobsource.main`.
"""CLI entry point: ``python -m jobsource.main``.
Scaffold stub. Argument parsing is wired so `--help` works; the actual batch
run lands in a later step (see jobsource/pipeline.py). Imports only stdlib so
`--help` works before the heavier dependencies are installed.
Parses arguments and delegates to :func:`~jobsource.pipeline.run_batch`.
All heavy imports are deferred past the argument parse so ``--help`` exits
immediately with no dependency overhead.
Usage examples::
python -m jobsource.main --help
python -m jobsource.main --batch-size 20 --search "software engineer" --location "United States"
python -m jobsource.main --search "data engineer" --search "ML engineer" --hours-old 48
"""
from __future__ import annotations
import argparse
import logging
import sys
@@ -46,8 +53,22 @@ def build_parser() -> argparse.ArgumentParser:
def main(argv: list[str] | None = None) -> int:
args = build_parser().parse_args(argv)
print("jobsource: scaffold stub -- pipeline not implemented yet.", file=sys.stderr)
print(f"parsed args: {vars(args)}", file=sys.stderr)
# Configure root logger so pipeline/cascade INFO and WARNING messages are
# visible on the terminal. Deferred past parse_args so --help stays instant.
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)-8s %(name)s: %(message)s",
datefmt="%H:%M:%S",
)
from jobsource.pipeline import run_batch # deferred: keeps --help fast before heavy deps
run_batch(
batch_size=args.batch_size,
search_terms=args.search, # None or list[str] from --search repeats
location=args.location,
hours_old=args.hours_old,
)
return 0

View File

@@ -1,12 +1,335 @@
"""Batch orchestration: dedup, per-record isolation, cascade, persistence, summary.
"""Batch orchestration: dedup, per-record isolation, cascade, persistence, summary."""
from __future__ import annotations
Scaffold stub -- not implemented yet.
"""
# TODO (pipeline): implement run_batch() per CLAUDE.md "Pipeline stages".
# run_batch() contract:
# - Accept batch_size, search terms, location, hours_old overrides.
# - Call the job source, dedup by job_id against the DB (skip already-seen jobs).
# - For each new RawJob, run the full cascade (resolve -> careers -> extract) in isolation:
# one failing record must NEVER abort the batch — catch, record failed/needs_review, continue.
# - Persist each JobResult to the DB and export output/results.csv when done.
# - Print a run summary: per-stage counts + % of new jobs reaching position_found.
import logging
import time
from dataclasses import dataclass
from .careers.cascade import find_careers_page
from .config import get_settings
from .db import Database
from .extract import extract_open_position
from .http import build_client
from .models import JobResult, JobStatus, RawJob
from .resolve import resolve_website
from .sources.base import JobSource
logger = logging.getLogger(__name__)
# ---------------------------------------------------------------------------
# Public summary result
# ---------------------------------------------------------------------------
@dataclass
class BatchSummary:
"""Per-stage counts returned (and printed) by run_batch()."""
fetched: int = 0
already_seen: int = 0
new: int = 0
website_resolved: int = 0
careers_found: int = 0
position_found: int = 0
needs_review: int = 0
failed: int = 0
@property
def coverage(self) -> float:
"""Fraction of new jobs that reached position_found."""
return self.position_found / self.new if self.new > 0 else 0.0
# ---------------------------------------------------------------------------
# Entry point
# ---------------------------------------------------------------------------
def run_batch(
batch_size: int | None = None,
*,
search_terms: list[str] | None = None,
location: str | None = None,
hours_old: int | None = None,
db: Database | None = None,
) -> BatchSummary:
"""Fetch, dedup, cascade, persist, and export one batch of jobs.
Override arguments take precedence over config values; ``None`` means
"use config default". Returns a :class:`BatchSummary` (also printed).
One failing record never aborts the batch — each job is isolated in its
own try/except; the outer loop continues regardless.
"""
settings = get_settings()
batch_size = batch_size if batch_size is not None else settings.batch_size
search_terms = search_terms if search_terms is not None else list(settings.search_terms)
location = location if location is not None else settings.location
hours_old = hours_old if hours_old is not None else settings.hours_old
own_db = db is None
if own_db:
db = Database()
summary = BatchSummary()
t0 = time.perf_counter()
csv_path = None
try:
# ------------------------------------------------------------------
# Stage 1 — Ingest
# ------------------------------------------------------------------
source = _make_source(settings)
print(f"\n{'='*64}")
print(
f" JobSourceAgent batch={batch_size}"
f" search={search_terms}"
f" hours_old={hours_old}"
)
print(f"{'='*64}")
print("Stage 1: ingesting jobs…")
fetch_limit = min(settings.results_wanted, batch_size)
raw_jobs: list[RawJob] = source.fetch_recent_jobs(
search_terms=search_terms,
location=location,
hours_old=hours_old,
results_wanted=fetch_limit,
)
summary.fetched = len(raw_jobs)
# Dedup: filter out jobs already in the DB; cap new jobs at batch_size.
new_jobs: list[RawJob] = []
for raw in raw_jobs:
if db.is_seen(raw.job_id):
summary.already_seen += 1
else:
new_jobs.append(raw)
if len(new_jobs) >= batch_size:
break
summary.new = len(new_jobs)
print(
f" Fetched {summary.fetched} jobs; "
f"{summary.already_seen} already seen; "
f"{summary.new} new to process.\n"
)
# ------------------------------------------------------------------
# Stages 1b → 3 — per-record cascade
# ------------------------------------------------------------------
if new_jobs:
with build_client() as client:
for raw in new_jobs:
result = JobResult.from_raw(raw)
print(f"[{raw.job_id}] {raw.company}")
try:
_process_one(result, raw, db, client)
except Exception as exc:
# Safety net — _process_one should not raise, but guard anyway.
logger.exception(
"pipeline: unexpected error for %s (%s): %s",
raw.job_id, raw.company, exc,
)
result.status = JobStatus.failed
try:
db.persist_result(result)
except Exception:
logger.exception(
"pipeline: failed to persist error record %s", raw.job_id
)
_tally(summary, result)
print()
# ------------------------------------------------------------------
# Stage 4 — Export CSV
# ------------------------------------------------------------------
try:
csv_path = db.export_csv()
except Exception as exc:
logger.error("pipeline: CSV export failed: %s", exc)
finally:
if own_db:
db.close()
elapsed = time.perf_counter() - t0
_print_summary(summary, elapsed, csv_path)
return summary
# ---------------------------------------------------------------------------
# Per-record cascade (must not raise — handles all errors internally)
# ---------------------------------------------------------------------------
def _process_one(
result: JobResult,
raw: RawJob,
db: Database,
client: object,
) -> None:
"""Run stages 1b → 3 for one job and persist the result.
Uses early returns for soft misses (needs_review) and catches exceptions
for hard failures (failed status). Never raises to the caller.
"""
# ------------------------------------------------------------------
# Stage 1b — Resolve company website
# ------------------------------------------------------------------
try:
website = resolve_website(raw.company, raw.website, client=client)
except Exception as exc:
logger.warning("pipeline(%s): resolve error: %s", raw.job_id, exc)
print(f" resolve: ERROR — {exc}")
result.status = JobStatus.failed
db.persist_result(result)
return
if not website:
print(" website: UNRESOLVED")
result.status = JobStatus.needs_review
db.persist_result(result)
return
result.website = website
result.status = JobStatus.website_resolved
result.company_key = db.company_key_for(raw.company, website)
print(f" website: {website}")
# ------------------------------------------------------------------
# Stage 2 — Find careers page (check company cache first)
# ------------------------------------------------------------------
cached_career_url = db.get_cached_career_url(result.company_key)
if cached_career_url:
careers_url = cached_career_url
careers_method = "cache"
position_shortcut: str | None = None
print(f" careers: {careers_url} [cache]")
else:
try:
cr = find_careers_page(website, company_name=raw.company, client=client)
except Exception as exc:
logger.warning("pipeline(%s): cascade error: %s", raw.job_id, exc)
print(f" careers: ERROR — {exc}")
result.status = JobStatus.needs_review
db.persist_result(result)
return
careers_url = cr.careers_url
careers_method = cr.method
position_shortcut = cr.position_url
if not careers_url:
print(f" careers: MISSED (method={cr.method})")
result.status = JobStatus.needs_review
db.persist_result(result)
return
print(
f" careers: {careers_url}"
f" [{careers_method} conf={cr.confidence:.2f}]"
)
result.career_page_url = careers_url
result.careers_method = careers_method
result.status = JobStatus.careers_found
# ------------------------------------------------------------------
# Stage 3 — Extract one open position
# ------------------------------------------------------------------
if position_shortcut:
# Free shortcut from ATS or browser-agent tier — no extra HTTP call.
position_url = position_shortcut
position_method = careers_method
else:
try:
position_url, position_method = extract_open_position(
careers_url, client=client
)
except Exception as exc:
logger.warning("pipeline(%s): extract error: %s", raw.job_id, exc)
print(f" position: ERROR — {exc}")
result.status = JobStatus.needs_review
db.persist_result(result)
return
if position_url:
result.open_position_url = position_url
result.position_method = position_method
result.status = JobStatus.position_found
print(f" position: {position_url} [{position_method}]")
else:
print(f" position: MISSED (method={position_method})")
result.status = JobStatus.needs_review
db.persist_result(result)
# ---------------------------------------------------------------------------
# Internal helpers
# ---------------------------------------------------------------------------
def _make_source(settings) -> JobSource:
"""Instantiate the configured job source provider."""
if settings.job_source == "apify":
from .sources.apify_source import ApifySource
return ApifySource()
from .sources.jobspy_source import JobSpySource
return JobSpySource()
def _tally(summary: BatchSummary, result: JobResult) -> None:
"""Accumulate per-stage counts from a processed JobResult."""
if result.website:
summary.website_resolved += 1
if result.career_page_url:
summary.careers_found += 1
if result.open_position_url:
summary.position_found += 1
if result.status == JobStatus.needs_review:
summary.needs_review += 1
elif result.status == JobStatus.failed:
summary.failed += 1
def _print_summary(
summary: BatchSummary, elapsed: float, csv_path: object
) -> None:
"""Print the run summary table to stdout."""
n = summary.new
wr = summary.website_resolved
cf = summary.careers_found
print(f"{'='*64}")
print(f" SUMMARY ({n} new jobs processed in {elapsed:.1f}s)")
print(f"{'='*64}")
print(
f" Stage 1 ingested {summary.fetched:>4}"
f" ({summary.already_seen} already seen)"
)
print(
f" Stage 1b website resolved {wr:>4} / {n}"
f" ({wr / max(n, 1) * 100:.0f}%)"
)
print(
f" Stage 2 careers found {cf:>4} / {wr}"
f" ({cf / max(wr, 1) * 100:.0f}% of resolved)"
)
print(
f" Stage 3 position found {summary.position_found:>4} / {cf}"
f" ({summary.position_found / max(cf, 1) * 100:.0f}% of careers)"
)
print(f" Needs review {summary.needs_review:>4}")
print(f" Failed {summary.failed:>4}")
print(
f" End-to-end coverage "
f"{summary.coverage * 100:.0f}%"
f" ({summary.position_found}/{n} new jobs fully resolved)"
)
if csv_path:
print(f" CSV written to {csv_path}")
print()

207
scripts/e2e_smoke.py Normal file
View File

@@ -0,0 +1,207 @@
#!/usr/bin/env python
"""End-to-end smoke test: Stage 1→3 on a 10-job batch.
Usage: .venv/bin/python scripts/e2e_smoke.py
"""
from __future__ import annotations
import logging
import sys
import time
from dataclasses import dataclass, field
# Make sure we run from the repo root so imports resolve correctly.
sys.path.insert(0, ".")
logging.basicConfig(
level=logging.WARNING, # keep jobspy/httpx chatter suppressed
format="%(name)s: %(message)s",
)
# Enable our own modules at INFO so we see cascade method names.
for _mod in ("jobsource.resolve", "jobsource.careers.cascade",
"jobsource.careers.ats", "jobsource.careers.heuristics",
"jobsource.careers.classify_llm", "jobsource.extract"):
logging.getLogger(_mod).setLevel(logging.INFO)
from jobsource.careers.cascade import find_careers_page
from jobsource.config import get_settings
from jobsource.extract import extract_open_position
from jobsource.http import build_client
from jobsource.models import RawJob
from jobsource.resolve import resolve_website
from jobsource.sources.jobspy_source import JobSpySource
# ---------------------------------------------------------------------------
# Per-record result
# ---------------------------------------------------------------------------
@dataclass
class Result:
job_id: str
company: str
website: str | None = None
careers_url: str | None = None
careers_method: str | None = None
position_url: str | None = None
position_method: str | None = None
error: str | None = None
# ---------------------------------------------------------------------------
# Main runner
# ---------------------------------------------------------------------------
BATCH_SIZE = 10
SEARCH_TERM = "software engineer"
LOCATION = "United States"
HOURS_OLD = 72
def run() -> None:
settings = get_settings()
print(f"\n{'='*64}")
print(f" E2E smoke batch={BATCH_SIZE} search={SEARCH_TERM!r} hours_old={HOURS_OLD}")
print(f"{'='*64}\n")
# ------------------------------------------------------------------
# Stage 1 — Ingest
# ------------------------------------------------------------------
print("Stage 1: ingesting…")
source = JobSpySource()
t0 = time.perf_counter()
raw_jobs: list[RawJob] = source.fetch_recent_jobs(
search_terms=[SEARCH_TERM],
location=LOCATION,
hours_old=HOURS_OLD,
results_wanted=BATCH_SIZE + 10, # fetch a few extra to cover dupes
)
elapsed = time.perf_counter() - t0
raw_jobs = raw_jobs[:BATCH_SIZE]
print(f" Ingested {len(raw_jobs)} unique jobs in {elapsed:.1f}s\n")
if not raw_jobs:
print("No jobs returned — aborting.")
return
results: list[Result] = []
with build_client() as client:
for raw in raw_jobs:
res = Result(job_id=raw.job_id, company=raw.company)
print(f"[{raw.job_id}] {raw.company}")
# -------------------------------------------------------
# Stage 1b — Resolve website
# -------------------------------------------------------
try:
res.website = resolve_website(
raw.company, raw.website, client=client
)
except Exception as exc:
res.error = f"resolve: {exc}"
print(f" resolve ERROR: {exc}")
results.append(res)
continue
if not res.website:
print(f" website: UNRESOLVED")
results.append(res)
continue
print(f" website: {res.website}")
# -------------------------------------------------------
# Stage 2 — Find careers page
# -------------------------------------------------------
try:
cr = find_careers_page(
res.website,
company_name=raw.company,
client=client,
)
res.careers_url = cr.careers_url
res.careers_method = cr.method
except Exception as exc:
res.error = f"cascade: {exc}"
print(f" careers ERROR: {exc}")
results.append(res)
continue
if not res.careers_url:
print(f" careers: MISSED (method={cr.method})")
results.append(res)
continue
print(f" careers: {res.careers_url} [{cr.method} conf={cr.confidence:.2f}]")
# -------------------------------------------------------
# Stage 3 — Extract one open position
# -------------------------------------------------------
try:
pos_url, pos_method = extract_open_position(
res.careers_url, client=client
)
res.position_url = pos_url
res.position_method = pos_method
except Exception as exc:
res.error = f"extract: {exc}"
print(f" position ERROR: {exc}")
results.append(res)
continue
if res.position_url:
print(f" position: {res.position_url} [{pos_method}]")
else:
print(f" position: MISSED (method={pos_method})")
results.append(res)
print()
# ------------------------------------------------------------------
# Summary
# ------------------------------------------------------------------
n_total = len(results)
n_website = sum(1 for r in results if r.website)
n_careers = sum(1 for r in results if r.careers_url)
n_position = sum(1 for r in results if r.position_url)
n_errors = sum(1 for r in results if r.error)
print(f"\n{'='*64}")
print(f" SUMMARY ({n_total} jobs processed)")
print(f"{'='*64}")
print(f" Stage 1 ingested {n_total:>3}")
print(f" Stage 1b website resolved {n_website:>3} / {n_total} ({n_website/n_total*100:.0f}%)")
print(f" Stage 2 careers found {n_careers:>3} / {n_website} ({n_careers/max(n_website,1)*100:.0f}% of resolved)")
print(f" Stage 3 position found {n_position:>3} / {n_careers} ({n_position/max(n_careers,1)*100:.0f}% of careers)")
print(f" Errors {n_errors:>3}")
# careers method breakdown
if n_careers:
method_counts: dict[str, int] = {}
for r in results:
if r.careers_method and r.careers_url:
method_counts[r.careers_method] = method_counts.get(r.careers_method, 0) + 1
print(f"\n Careers method breakdown:")
for m, c in sorted(method_counts.items(), key=lambda x: -x[1]):
print(f" {m:<30} {c}")
# position method breakdown
if n_position:
pos_method_counts: dict[str, int] = {}
for r in results:
if r.position_method and r.position_url:
pos_method_counts[r.position_method] = pos_method_counts.get(r.position_method, 0) + 1
print(f"\n Position method breakdown:")
for m, c in sorted(pos_method_counts.items(), key=lambda x: -x[1]):
print(f" {m:<30} {c}")
# CSV-ready output
print(f"\n{'='*64}")
print(" company_name, career_page_url, open_position_url")
print(f"{'='*64}")
for r in results:
print(f" {r.company!r}, {r.careers_url or ''!r}, {r.position_url or ''!r}")
print()
if __name__ == "__main__":
run()

View File

@@ -0,0 +1,386 @@
"""Tests for jobsource/agent_fallback.py — all network-free via monkeypatching."""
from __future__ import annotations
import pytest
import jobsource.agent_fallback as _mod
from jobsource.agent_fallback import (
AgentFallbackResult,
_AgentOutput,
_site_key,
_valid_url,
find_and_extract,
)
# ---------------------------------------------------------------------------
# Helpers and fakes
# ---------------------------------------------------------------------------
def _placeholder_settings(**overrides):
"""Return a fake settings object with placeholder values (agent unavailable)."""
defaults = {
"enable_browser_agent": True,
"agent_model": "PLACEHOLDER_AGENT_MODEL",
"llm_api_key": "PLACEHOLDER_LLM_API_KEY",
"browser_headless": True,
}
defaults.update(overrides)
return type("S", (), defaults)()
def _real_settings(**overrides):
"""Return a fake settings object with non-placeholder values."""
defaults = {
"enable_browser_agent": True,
"agent_model": "anthropic:claude-sonnet-4-6",
"llm_api_key": "sk-ant-real",
"browser_headless": True,
}
defaults.update(overrides)
return type("S", (), defaults)()
def _patch_settings(monkeypatch, settings_obj):
monkeypatch.setattr(_mod, "get_settings", lambda: settings_obj)
def _patch_chat(monkeypatch, chat_obj=None):
"""Monkeypatch _build_chat_model to return chat_obj (default sentinel)."""
chat_obj = chat_obj or object()
monkeypatch.setattr(_mod, "_build_chat_model", lambda: chat_obj)
return chat_obj
def _patch_session(monkeypatch, output: _AgentOutput | None):
"""Monkeypatch _agent_session so no real browser launches."""
async def _fake_session(website, chat):
return output
monkeypatch.setattr(_mod, "_agent_session", _fake_session)
# ---------------------------------------------------------------------------
# _site_key
# ---------------------------------------------------------------------------
class TestSiteKey:
def test_strips_www(self):
assert _site_key("https://www.acme.com/foo") == "acme.com"
def test_lowercases(self):
assert _site_key("https://ACME.COM") == "acme.com"
def test_no_www(self):
assert _site_key("https://acme.com/careers") == "acme.com"
def test_subdomain_preserved(self):
assert _site_key("https://jobs.acme.com") == "jobs.acme.com"
# ---------------------------------------------------------------------------
# _valid_url
# ---------------------------------------------------------------------------
class TestValidUrl:
def test_https_url_passes(self):
assert _valid_url("https://acme.com/jobs/1") == "https://acme.com/jobs/1"
def test_http_url_passes(self):
assert _valid_url("http://acme.com/jobs/1") == "http://acme.com/jobs/1"
def test_none_returns_none(self):
assert _valid_url(None) is None
def test_junk_string_returns_none(self):
assert _valid_url("null") is None
def test_relative_url_returns_none(self):
assert _valid_url("/jobs/1") is None
# ---------------------------------------------------------------------------
# Gate: enable_browser_agent=False
# ---------------------------------------------------------------------------
class TestFeatureFlagGate:
def test_returns_none_when_flag_off(self, monkeypatch):
_patch_settings(monkeypatch, _placeholder_settings(enable_browser_agent=False))
# Clear memo so there's no cached entry from another test.
_mod._MEMO.clear()
result = find_and_extract("https://acme.com")
assert result is None
def test_does_not_cache_when_flag_off(self, monkeypatch):
_patch_settings(monkeypatch, _placeholder_settings(enable_browser_agent=False))
_mod._MEMO.clear()
find_and_extract("https://flagoff.com")
assert "flagoff.com" not in _mod._MEMO
# ---------------------------------------------------------------------------
# Gate: placeholder LLM config
# ---------------------------------------------------------------------------
class TestPlaceholderGate:
def test_returns_none_for_placeholder_model(self, monkeypatch):
_patch_settings(monkeypatch, _placeholder_settings())
_mod._MEMO.clear()
result = find_and_extract("https://acme.com")
assert result is None
def test_returns_none_for_placeholder_key_only(self, monkeypatch):
_patch_settings(monkeypatch, _real_settings(llm_api_key="PLACEHOLDER_LLM_API_KEY"))
_mod._MEMO.clear()
result = find_and_extract("https://acme.com")
assert result is None
def test_returns_none_for_placeholder_model_only(self, monkeypatch):
_patch_settings(monkeypatch, _real_settings(agent_model="PLACEHOLDER_AGENT_MODEL"))
_mod._MEMO.clear()
result = find_and_extract("https://acme.com")
assert result is None
# ---------------------------------------------------------------------------
# Gate: _build_chat_model raises (missing package / bad config)
# ---------------------------------------------------------------------------
class TestChatModelBuildFailure:
def test_returns_none_and_memos_none_on_import_error(self, monkeypatch):
_patch_settings(monkeypatch, _real_settings())
_mod._MEMO.clear()
def _raise():
raise ImportError("browser_use not installed")
monkeypatch.setattr(_mod, "_build_chat_model", _raise)
result = find_and_extract("https://chatfail.com")
assert result is None
assert _mod._MEMO.get("chatfail.com") is None
# ---------------------------------------------------------------------------
# Gate: _agent_session raises / returns None
# ---------------------------------------------------------------------------
class TestAgentSessionFailure:
def test_returns_none_on_session_exception(self, monkeypatch):
_patch_settings(monkeypatch, _real_settings())
_mod._MEMO.clear()
_patch_chat(monkeypatch)
async def _raise_session(website, chat):
raise RuntimeError("Chromium unavailable")
monkeypatch.setattr(_mod, "_agent_session", _raise_session)
result = find_and_extract("https://chromiumfail.com")
assert result is None
def test_returns_none_when_session_returns_none(self, monkeypatch):
_patch_settings(monkeypatch, _real_settings())
_mod._MEMO.clear()
_patch_chat(monkeypatch)
_patch_session(monkeypatch, None)
result = find_and_extract("https://nonefail.com")
assert result is None
# ---------------------------------------------------------------------------
# Success path: structured output parsed into AgentFallbackResult
# ---------------------------------------------------------------------------
class TestSuccessPath:
def test_both_urls_returned(self, monkeypatch):
_patch_settings(monkeypatch, _real_settings())
_mod._MEMO.clear()
_patch_chat(monkeypatch)
_patch_session(
monkeypatch,
_AgentOutput(
careers_url="https://acme.com/careers",
position_url="https://acme.com/careers/jobs/42-swe",
),
)
result = find_and_extract("https://acme.com")
assert result is not None
assert result.careers_url == "https://acme.com/careers"
assert result.position_url == "https://acme.com/careers/jobs/42-swe"
def test_junk_careers_url_filtered(self, monkeypatch):
_patch_settings(monkeypatch, _real_settings())
_mod._MEMO.clear()
_patch_chat(monkeypatch)
_patch_session(
monkeypatch,
_AgentOutput(careers_url="null", position_url="https://acme.com/jobs/1"),
)
result = find_and_extract("https://acme.com")
assert result is not None
assert result.careers_url is None
assert result.position_url == "https://acme.com/jobs/1"
def test_result_stored_in_memo(self, monkeypatch):
_patch_settings(monkeypatch, _real_settings())
_mod._MEMO.clear()
_patch_chat(monkeypatch)
_patch_session(
monkeypatch,
_AgentOutput(
careers_url="https://memo-test.com/careers",
position_url="https://memo-test.com/jobs/1",
),
)
find_and_extract("https://memo-test.com")
cached = _mod._MEMO.get("memo-test.com")
assert cached is not None
assert cached.careers_url == "https://memo-test.com/careers"
# ---------------------------------------------------------------------------
# Memo: same-domain second call uses cache, not a second session
# ---------------------------------------------------------------------------
class TestMemo:
def test_second_call_uses_cache(self, monkeypatch):
_patch_settings(monkeypatch, _real_settings())
_mod._MEMO.clear()
_patch_chat(monkeypatch)
call_count = 0
async def _counting_session(website, chat):
nonlocal call_count
call_count += 1
return _AgentOutput(
careers_url="https://sharetest.com/careers",
position_url="https://sharetest.com/jobs/1",
)
monkeypatch.setattr(_mod, "_agent_session", _counting_session)
r1 = find_and_extract("https://sharetest.com")
r2 = find_and_extract("https://sharetest.com/careers") # same host, careers path
assert call_count == 1, "Expected exactly one browser session for the same host"
assert r1 == r2
def test_memo_none_prevents_retry(self, monkeypatch):
_patch_settings(monkeypatch, _real_settings())
_mod._MEMO.clear()
_patch_chat(monkeypatch)
call_count = 0
async def _failing_session(website, chat):
nonlocal call_count
call_count += 1
raise RuntimeError("broken")
monkeypatch.setattr(_mod, "_agent_session", _failing_session)
find_and_extract("https://memonil.com")
find_and_extract("https://memonil.com") # should use cached None
assert call_count == 1
# ---------------------------------------------------------------------------
# Cascade/extract wiring integration (lightweight, no real browser)
# ---------------------------------------------------------------------------
class TestCascadeWiring:
"""Verify find_and_extract is reachable from cascade.find_careers_page."""
def test_cascade_tier6_returns_browser_agent_method(self, monkeypatch):
from jobsource.careers.cascade import find_careers_page, CareersResult
# Patch all cheaper tiers to miss.
monkeypatch.setattr(
"jobsource.careers.cascade._ats.detect_and_fetch", lambda *a, **kw: None
)
monkeypatch.setattr(
"jobsource.careers.cascade._ats.recover_via_slug_guess", lambda *a, **kw: None
)
monkeypatch.setattr(
"jobsource.careers.cascade._heuristics.probe_url_patterns", lambda *a, **kw: None
)
monkeypatch.setattr(
"jobsource.careers.cascade._heuristics.scan_homepage_links", lambda *a, **kw: None
)
monkeypatch.setattr(
"jobsource.careers.cascade._heuristics.parse_sitemap", lambda *a, **kw: None
)
monkeypatch.setattr(
"jobsource.careers.cascade._classify_llm.extract_anchors", lambda *a, **kw: []
)
monkeypatch.setattr(
"jobsource.careers.cascade._classify_llm.classify_careers_link",
lambda *a, **kw: None,
)
# Patch the safe homepage fetch.
from jobsource.careers import cascade as _cascade
monkeypatch.setattr(_cascade, "_safe_get_html", lambda *a, **kw: None)
# Patch agent_fallback.
_mod._MEMO.clear()
monkeypatch.setattr(
"jobsource.careers.cascade._agent_fallback.find_and_extract",
lambda *a, **kw: AgentFallbackResult(
careers_url="https://wired.com/careers",
position_url="https://wired.com/jobs/1",
),
)
result = find_careers_page("https://wired.com")
assert result.method == "browser_agent"
assert result.careers_url == "https://wired.com/careers"
assert result.position_url == "https://wired.com/jobs/1"
assert result.confidence == 0.50
class TestExtractWiring:
"""Verify find_and_extract is reachable from extract.extract_open_position."""
def test_extract_tier5_returns_browser_agent_method(self, monkeypatch):
from jobsource.extract import extract_open_position
# Patch cheaper tiers to miss (ATS, JSON-LD, anchors, LLM).
monkeypatch.setattr(
"jobsource.extract._ats.detect_ats_in_url", lambda *a, **kw: None
)
monkeypatch.setattr(
"jobsource.extract._ats.detect_ats_in_html", lambda *a, **kw: None
)
import httpx
fake_response = type(
"R",
(),
{"status_code": 200, "text": "<html><body></body></html>"},
)()
monkeypatch.setattr(
"jobsource.extract.request_with_retries",
lambda *a, **kw: fake_response,
)
monkeypatch.setattr(
"jobsource.extract.extract_anchors", lambda *a, **kw: []
)
monkeypatch.setattr(
"jobsource.extract.classify_job_link", lambda *a, **kw: None
)
# Patch agent_fallback.
_mod._MEMO.clear()
monkeypatch.setattr(
"jobsource.extract._agent_fallback.find_and_extract",
lambda *a, **kw: AgentFallbackResult(
careers_url="https://extract-wired.com/careers",
position_url="https://extract-wired.com/jobs/99",
),
)
url, method = extract_open_position("https://extract-wired.com/careers")
assert method == "browser_agent"
assert url == "https://extract-wired.com/jobs/99"

View File

@@ -550,3 +550,112 @@ class TestSlugGuessTier:
result = find_careers_page("https://www.acme.com", client=FakeClient())
assert result.method == "url_pattern"
assert result.careers_url == "https://acme.com/careers"
# ---------------------------------------------------------------------------
# Tier 5 — Cheap-LLM careers-link classification
# ---------------------------------------------------------------------------
class TestLLMClassifyTier:
"""Tier 5 fires after sitemap, no-ops when LLM returns None, and is skipped on earlier hits."""
def _patch_all_deterministic_miss(self, monkeypatch, *, homepage_html: str = "<html/>") -> None:
"""Patch tiers 14 to all miss; leaves Tier 5 for the caller to control."""
monkeypatch.setattr(
"jobsource.careers.cascade._safe_get_html",
lambda website, client: homepage_html,
)
monkeypatch.setattr("jobsource.careers.cascade._ats.detect_and_fetch", lambda *a, **kw: None)
monkeypatch.setattr("jobsource.careers.cascade._ats.recover_via_slug_guess", lambda *a, **kw: None)
monkeypatch.setattr("jobsource.careers.cascade._heuristics.probe_url_patterns", lambda *a, **kw: None)
monkeypatch.setattr("jobsource.careers.cascade._heuristics.scan_homepage_links", lambda *a, **kw: None)
monkeypatch.setattr("jobsource.careers.cascade._heuristics.parse_sitemap", lambda *a, **kw: None)
# Prevent _finalize from making network calls
monkeypatch.setattr("jobsource.careers.cascade._ats.detect_ats_in_url", lambda url: None)
def test_llm_classify_hit_returns_055_confidence(self, monkeypatch):
from jobsource.careers.classify_llm import Anchor, CareerLinkResult
self._patch_all_deterministic_miss(
monkeypatch,
homepage_html='<html><body><a href="/careers">Careers</a></body></html>',
)
monkeypatch.setattr(
"jobsource.careers.cascade._classify_llm.extract_anchors",
lambda html, base_url: [Anchor(url="https://acme.com/careers", text="Careers")],
)
monkeypatch.setattr(
"jobsource.careers.cascade._classify_llm.classify_careers_link",
lambda anchors: CareerLinkResult(careers_url="https://acme.com/careers", confidence=0.85),
)
result = find_careers_page("https://www.acme.com", client=FakeClient())
assert result.method == "llm_classify"
assert result.confidence == 0.55 # fixed tier-level confidence from cascade
assert result.careers_url == "https://acme.com/careers"
def test_llm_classify_none_falls_through_to_method_none(self, monkeypatch):
from jobsource.careers.classify_llm import Anchor
self._patch_all_deterministic_miss(monkeypatch)
monkeypatch.setattr(
"jobsource.careers.cascade._classify_llm.extract_anchors",
lambda html, base_url: [Anchor(url="https://acme.com/about", text="About")],
)
monkeypatch.setattr(
"jobsource.careers.cascade._classify_llm.classify_careers_link",
lambda anchors: None,
)
result = find_careers_page("https://www.acme.com", client=FakeClient())
assert result.method == "none"
assert result.careers_url is None
def test_tier5_skipped_when_homepage_html_is_none(self, monkeypatch):
"""When the homepage fetch returns None, Tier 5 is silently bypassed."""
llm_called: list[bool] = []
self._patch_all_deterministic_miss(monkeypatch)
monkeypatch.setattr(
"jobsource.careers.cascade._safe_get_html",
lambda website, client: None,
)
monkeypatch.setattr(
"jobsource.careers.cascade._classify_llm.classify_careers_link",
lambda anchors: llm_called.append(True) or None,
)
find_careers_page("https://www.acme.com", client=FakeClient())
assert llm_called == []
def test_tier5_not_called_when_sitemap_hits(self, monkeypatch):
"""When sitemap resolves, the cascade returns before reaching Tier 5."""
llm_called: list[bool] = []
monkeypatch.setattr(
"jobsource.careers.cascade._safe_get_html",
lambda website, client: "<html/>",
)
monkeypatch.setattr("jobsource.careers.cascade._ats.detect_and_fetch", lambda *a, **kw: None)
monkeypatch.setattr("jobsource.careers.cascade._ats.recover_via_slug_guess", lambda *a, **kw: None)
monkeypatch.setattr("jobsource.careers.cascade._heuristics.probe_url_patterns", lambda *a, **kw: None)
monkeypatch.setattr("jobsource.careers.cascade._heuristics.scan_homepage_links", lambda *a, **kw: None)
monkeypatch.setattr(
"jobsource.careers.cascade._heuristics.parse_sitemap",
lambda *a, **kw: "https://acme.com/careers",
)
monkeypatch.setattr("jobsource.careers.cascade._ats.detect_ats_in_url", lambda url: None)
monkeypatch.setattr(
"jobsource.careers.cascade._classify_llm.classify_careers_link",
lambda anchors: llm_called.append(True) or None,
)
result = find_careers_page("https://www.acme.com", client=FakeClient())
assert result.method == "sitemap"
assert llm_called == []
def test_tier5_error_falls_through_gracefully(self, monkeypatch):
"""A Tier 5 exception should not abort the cascade."""
self._patch_all_deterministic_miss(monkeypatch)
monkeypatch.setattr(
"jobsource.careers.cascade._classify_llm.extract_anchors",
lambda html, base_url: (_ for _ in ()).throw(RuntimeError("classify boom")),
)
result = find_careers_page("https://www.acme.com", client=FakeClient())
assert result.method == "none"
assert result.careers_url is None

256
tests/test_classify_llm.py Normal file
View File

@@ -0,0 +1,256 @@
"""Tests for careers/classify_llm.py — all network-free via monkeypatching."""
from __future__ import annotations
import pytest
from jobsource.careers.classify_llm import (
Anchor,
CareerLinkResult,
JobLinkResult,
_LinkChoice,
classify_careers_link,
classify_job_link,
extract_anchors,
)
# ---------------------------------------------------------------------------
# Fake helpers
# ---------------------------------------------------------------------------
class _FakeResult:
"""Minimal stand-in for a Pydantic AI run result."""
def __init__(self, output: _LinkChoice) -> None:
self.output = output
class _FakeAgent:
"""Stand-in for a Pydantic AI Agent; returns a fixed _LinkChoice."""
def __init__(self, choice: _LinkChoice) -> None:
self._choice = choice
def run_sync(self, prompt: str) -> _FakeResult:
return _FakeResult(self._choice)
class _RaisingAgent:
"""Stand-in that raises on run_sync, simulating an LLM failure."""
def run_sync(self, prompt: str) -> _FakeResult:
raise RuntimeError("LLM unavailable")
def _real_settings():
"""Fake settings object with non-placeholder values."""
return type("S", (), {
"classifier_model": "openai:gpt-4o-mini",
"llm_api_key": "sk-real-key",
})()
# ---------------------------------------------------------------------------
# extract_anchors
# ---------------------------------------------------------------------------
_SAMPLE_HTML = """
<html><body>
<a href="/about">About Us</a>
<a href="/careers">Careers</a>
<a href="https://external.com/jobs">External Jobs</a>
<a href="mailto:info@acme.com">Email</a>
<a href="#section">Anchor</a>
<a href="javascript:void(0)">JS link</a>
<a href="/careers">Duplicate careers link</a>
</body></html>
"""
class TestExtractAnchors:
def test_resolves_relative_to_absolute(self):
anchors = extract_anchors(_SAMPLE_HTML, "https://acme.com")
urls = [a.url for a in anchors]
assert "https://acme.com/about" in urls
assert "https://acme.com/careers" in urls
def test_preserves_already_absolute_urls(self):
anchors = extract_anchors(_SAMPLE_HTML, "https://acme.com")
urls = [a.url for a in anchors]
assert "https://external.com/jobs" in urls
def test_skips_mailto(self):
anchors = extract_anchors(_SAMPLE_HTML, "https://acme.com")
assert all("mailto:" not in a.url for a in anchors)
def test_skips_hash_fragments(self):
anchors = extract_anchors(_SAMPLE_HTML, "https://acme.com")
# A pure # anchor resolves to base_url itself (https://acme.com) —
# it starts with http so it passes the scheme check but we skip hrefs
# that start with "#". Verify the JavaScript link is also gone.
assert all("javascript:" not in a.url for a in anchors)
def test_deduplicates_by_url(self):
anchors = extract_anchors(_SAMPLE_HTML, "https://acme.com")
urls = [a.url for a in anchors]
assert urls.count("https://acme.com/careers") == 1
def test_captures_anchor_text(self):
anchors = extract_anchors(_SAMPLE_HTML, "https://acme.com")
career_anchor = next(a for a in anchors if a.url == "https://acme.com/careers")
assert career_anchor.text == "Careers"
def test_empty_html_returns_empty_list(self):
assert extract_anchors("", "https://acme.com") == []
def test_html_with_no_links_returns_empty_list(self):
assert extract_anchors("<html><body><p>No links here.</p></body></html>", "https://acme.com") == []
def test_text_truncated_at_max_length(self):
long_text = "x" * 200
html = f'<html><body><a href="/careers">{long_text}</a></body></html>'
anchors = extract_anchors(html, "https://acme.com")
assert len(anchors) == 1
assert len(anchors[0].text) <= 80
# ---------------------------------------------------------------------------
# Graceful degradation: placeholder settings → always None
# ---------------------------------------------------------------------------
class TestPlaceholderGate:
"""When classifier_model or llm_api_key is a placeholder, classifiers no-op."""
def test_classify_careers_link_returns_none_with_placeholder(self):
anchors = [Anchor(url="https://acme.com/careers", text="Careers")]
assert classify_careers_link(anchors) is None
def test_classify_job_link_returns_none_with_placeholder(self):
anchors = [Anchor(url="https://acme.com/jobs/123", text="Engineer")]
assert classify_job_link(anchors) is None
def test_classify_careers_link_returns_none_on_empty_list(self, monkeypatch):
monkeypatch.setattr(
"jobsource.careers.classify_llm.get_settings",
_real_settings,
)
assert classify_careers_link([]) is None
def test_classify_job_link_returns_none_on_empty_list(self, monkeypatch):
monkeypatch.setattr(
"jobsource.careers.classify_llm.get_settings",
_real_settings,
)
assert classify_job_link([]) is None
# ---------------------------------------------------------------------------
# Success paths via _build_link_agent monkeypatch
# ---------------------------------------------------------------------------
class TestClassifySuccess:
def _patch_agent(self, monkeypatch, choice: _LinkChoice) -> None:
monkeypatch.setattr(
"jobsource.careers.classify_llm.get_settings",
_real_settings,
)
monkeypatch.setattr(
"jobsource.careers.classify_llm._build_link_agent",
lambda instructions: _FakeAgent(choice),
)
def test_careers_pick_returns_mapped_url(self, monkeypatch):
anchors = [
Anchor(url="https://acme.com/about", text="About"),
Anchor(url="https://acme.com/careers", text="Careers"),
]
self._patch_agent(monkeypatch, _LinkChoice(index=1, confidence=0.95))
result = classify_careers_link(anchors)
assert result is not None
assert isinstance(result, CareerLinkResult)
assert result.careers_url == "https://acme.com/careers"
assert result.confidence == 0.95
def test_job_pick_returns_mapped_url(self, monkeypatch):
anchors = [
Anchor(url="https://acme.com/jobs/", text="All Jobs"),
Anchor(url="https://acme.com/jobs/123", text="Software Engineer"),
]
self._patch_agent(monkeypatch, _LinkChoice(index=1, confidence=0.85))
result = classify_job_link(anchors)
assert result is not None
assert isinstance(result, JobLinkResult)
assert result.job_url == "https://acme.com/jobs/123"
assert result.confidence == 0.85
def test_careers_first_index_zero(self, monkeypatch):
anchors = [Anchor(url="https://acme.com/jobs", text="Jobs")]
self._patch_agent(monkeypatch, _LinkChoice(index=0, confidence=0.70))
result = classify_careers_link(anchors)
assert result is not None
assert result.careers_url == "https://acme.com/jobs"
# ---------------------------------------------------------------------------
# Abstain and failure paths
# ---------------------------------------------------------------------------
class TestClassifyFailurePaths:
def _patch_real_settings(self, monkeypatch) -> None:
monkeypatch.setattr(
"jobsource.careers.classify_llm.get_settings",
_real_settings,
)
def test_abstain_index_minus_one_returns_none(self, monkeypatch):
self._patch_real_settings(monkeypatch)
monkeypatch.setattr(
"jobsource.careers.classify_llm._build_link_agent",
lambda instructions: _FakeAgent(_LinkChoice(index=-1, confidence=0.0)),
)
anchors = [Anchor(url="https://acme.com/about", text="About")]
assert classify_careers_link(anchors) is None
def test_out_of_range_index_returns_none(self, monkeypatch):
self._patch_real_settings(monkeypatch)
monkeypatch.setattr(
"jobsource.careers.classify_llm._build_link_agent",
lambda instructions: _FakeAgent(_LinkChoice(index=99, confidence=0.9)),
)
anchors = [Anchor(url="https://acme.com/careers", text="Careers")]
assert classify_careers_link(anchors) is None
def test_run_sync_exception_returns_none(self, monkeypatch):
self._patch_real_settings(monkeypatch)
monkeypatch.setattr(
"jobsource.careers.classify_llm._build_link_agent",
lambda instructions: _RaisingAgent(),
)
anchors = [Anchor(url="https://acme.com/careers", text="Careers")]
assert classify_careers_link(anchors) is None
def test_build_agent_exception_returns_none(self, monkeypatch):
self._patch_real_settings(monkeypatch)
def boom(instructions: str):
raise RuntimeError("import error")
monkeypatch.setattr(
"jobsource.careers.classify_llm._build_link_agent",
boom,
)
anchors = [Anchor(url="https://acme.com/careers", text="Careers")]
assert classify_careers_link(anchors) is None
def test_job_classify_run_sync_exception_returns_none(self, monkeypatch):
self._patch_real_settings(monkeypatch)
monkeypatch.setattr(
"jobsource.careers.classify_llm._build_link_agent",
lambda instructions: _RaisingAgent(),
)
anchors = [Anchor(url="https://acme.com/jobs/123", text="Engineer")]
assert classify_job_link(anchors) is None

565
tests/test_db.py Normal file
View File

@@ -0,0 +1,565 @@
"""Tests for jobsource/db.py and jobsource/pipeline.py.
All tests are offline — no network calls, no real job source.
The DB uses tmp_path (file-backed) or ":memory:" via path override.
"""
from __future__ import annotations
import csv
from datetime import datetime, timezone
from pathlib import Path
from unittest.mock import MagicMock, patch
import pytest
from jobsource.db import Database, _company_key
from jobsource.models import CSV_COLUMNS, JobResult, JobStatus, RawJob
# ---------------------------------------------------------------------------
# Helpers
# ---------------------------------------------------------------------------
def _db(tmp_path: Path) -> Database:
"""Open a fresh file-backed Database in tmp_path."""
return Database(path=tmp_path / "test.db")
def _result(
job_id: str = "1",
company_name: str = "Acme",
website: str | None = "https://acme.com",
career_url: str | None = "https://acme.com/careers",
position_url: str | None = "https://acme.com/careers/jobs/42",
status: JobStatus = JobStatus.position_found,
company_key: str | None = None,
) -> JobResult:
return JobResult(
job_id=job_id,
company_name=company_name,
company_key=company_key or _company_key(company_name, website),
website=website,
career_page_url=career_url,
open_position_url=position_url,
status=status,
)
# ---------------------------------------------------------------------------
# _company_key (pure function)
# ---------------------------------------------------------------------------
class TestCompanyKey:
def test_domain_stripped_of_www(self) -> None:
assert _company_key("Acme", "https://www.acme.com/about") == "acme.com"
def test_domain_lowercase(self) -> None:
assert _company_key("Foo", "https://Foo.IO") == "foo.io"
def test_no_website_uses_name(self) -> None:
assert _company_key("Globex Corp", None) == "globex corp"
def test_no_website_lowercased(self) -> None:
assert _company_key("Initech LLC", None) == "initech llc"
def test_empty_host_falls_back_to_name(self) -> None:
# Malformed URL with empty hostname
result = _company_key("BadCo", "://no-host")
assert result == "badco"
def test_empty_name_and_no_website_returns_unknown(self) -> None:
assert _company_key("", None) == "unknown"
def test_subpath_ignored(self) -> None:
assert _company_key("X", "https://x.example.com/jobs/42") == "x.example.com"
# ---------------------------------------------------------------------------
# Database.init_schema / open
# ---------------------------------------------------------------------------
class TestDatabaseSchema:
def test_creates_tables(self, tmp_path: Path) -> None:
db = _db(tmp_path)
try:
cur = db._conn.execute(
"SELECT name FROM sqlite_master WHERE type='table' ORDER BY name"
)
tables = {r[0] for r in cur.fetchall()}
assert "companies" in tables
assert "jobs" in tables
finally:
db.close()
def test_idempotent_schema(self, tmp_path: Path) -> None:
"""Opening the same DB twice must not raise."""
db1 = Database(path=tmp_path / "idempotent.db")
db1.close()
db2 = Database(path=tmp_path / "idempotent.db")
db2.close()
def test_parent_dirs_created(self, tmp_path: Path) -> None:
nested = tmp_path / "a" / "b" / "c" / "test.db"
db = Database(path=nested)
try:
assert nested.exists()
finally:
db.close()
def test_context_manager(self, tmp_path: Path) -> None:
with Database(path=tmp_path / "cm.db") as db:
assert not db.is_seen("x")
# ---------------------------------------------------------------------------
# is_seen / dedup
# ---------------------------------------------------------------------------
class TestIsSeeen:
def test_unseen_job_returns_false(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
assert not db.is_seen("99999")
def test_seen_after_persist(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
r = _result(job_id="42")
db.persist_result(r)
assert db.is_seen("42")
def test_different_job_ids_independent(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
db.persist_result(_result(job_id="1"))
assert db.is_seen("1")
assert not db.is_seen("2")
# ---------------------------------------------------------------------------
# company_key_for
# ---------------------------------------------------------------------------
class TestCompanyKeyFor:
def test_delegates_to_pure_fn(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
assert db.company_key_for("Acme", "https://acme.com") == "acme.com"
assert db.company_key_for("Foo LLC", None) == "foo llc"
# ---------------------------------------------------------------------------
# get_cached_career_url
# ---------------------------------------------------------------------------
class TestCachedCareerUrl:
def test_miss_returns_none(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
assert db.get_cached_career_url("acme.com") is None
def test_hit_after_persist(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
r = _result(job_id="1", website="https://acme.com", career_url="https://acme.com/careers")
db.persist_result(r)
assert db.get_cached_career_url("acme.com") == "https://acme.com/careers"
def test_null_career_url_not_cached(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
r = _result(job_id="1", career_url=None, status=JobStatus.needs_review)
db.persist_result(r)
assert db.get_cached_career_url("acme.com") is None
# ---------------------------------------------------------------------------
# upsert_company COALESCE / first_seen preservation
# ---------------------------------------------------------------------------
class TestUpsertCompany:
def test_insert_then_null_fill(self, tmp_path: Path) -> None:
"""A NULL value in the second upsert fills the existing NULL without error."""
with _db(tmp_path) as db:
db.upsert_company("x.com", "X", None, None)
db.upsert_company("x.com", "X", "https://x.com", "https://x.com/careers")
row = db._conn.execute(
"SELECT website, career_url FROM companies WHERE company_key='x.com'"
).fetchone()
assert row["website"] == "https://x.com"
assert row["career_url"] == "https://x.com/careers"
def test_existing_value_not_clobbered(self, tmp_path: Path) -> None:
"""A second upsert with NULL must not overwrite an existing value."""
with _db(tmp_path) as db:
db.upsert_company("y.com", "Y", "https://y.com", "https://y.com/jobs")
db.upsert_company("y.com", "Y", None, None)
row = db._conn.execute(
"SELECT website, career_url FROM companies WHERE company_key='y.com'"
).fetchone()
assert row["website"] == "https://y.com"
assert row["career_url"] == "https://y.com/jobs"
def test_first_seen_preserved(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
db.upsert_company("z.com", "Z", None, None)
row1 = db._conn.execute(
"SELECT first_seen FROM companies WHERE company_key='z.com'"
).fetchone()
db.upsert_company("z.com", "Z", "https://z.com", "https://z.com/careers")
row2 = db._conn.execute(
"SELECT first_seen FROM companies WHERE company_key='z.com'"
).fetchone()
assert row1["first_seen"] == row2["first_seen"]
# ---------------------------------------------------------------------------
# persist_result / upsert_job
# ---------------------------------------------------------------------------
class TestPersistResult:
def test_both_tables_populated(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
r = _result(job_id="10", company_name="Initech", website="https://initech.com")
db.persist_result(r)
job = db._conn.execute(
"SELECT * FROM jobs WHERE job_id='10'"
).fetchone()
assert job is not None
assert job["status"] == "position_found"
company = db._conn.execute(
"SELECT * FROM companies WHERE company_key='initech.com'"
).fetchone()
assert company is not None
assert company["name"] == "Initech"
def test_listed_at_null_stored_as_null(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
r = _result(job_id="20")
assert r.listed_at is None
db.persist_result(r)
row = db._conn.execute(
"SELECT listed_at FROM jobs WHERE job_id='20'"
).fetchone()
assert row["listed_at"] is None
def test_listed_at_serialized(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
r = _result(job_id="30")
r.listed_at = datetime(2026, 1, 15, 12, 0, 0, tzinfo=timezone.utc)
db.persist_result(r)
row = db._conn.execute(
"SELECT listed_at FROM jobs WHERE job_id='30'"
).fetchone()
assert row["listed_at"] is not None
assert "2026" in row["listed_at"]
def test_upsert_fills_null_position(self, tmp_path: Path) -> None:
"""A second persist with position_url fills a prior NULL."""
with _db(tmp_path) as db:
r1 = _result(job_id="50", position_url=None, status=JobStatus.needs_review)
db.persist_result(r1)
r2 = _result(job_id="50", position_url="https://acme.com/jobs/1", status=JobStatus.position_found)
db.persist_result(r2)
row = db._conn.execute(
"SELECT position_url, status FROM jobs WHERE job_id='50'"
).fetchone()
assert row["position_url"] == "https://acme.com/jobs/1"
assert row["status"] == "position_found"
def test_job_first_seen_preserved(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
r = _result(job_id="60", position_url=None, status=JobStatus.needs_review)
db.persist_result(r)
row1 = db._conn.execute(
"SELECT first_seen FROM jobs WHERE job_id='60'"
).fetchone()
r2 = _result(job_id="60", position_url="https://acme.com/jobs/1", status=JobStatus.position_found)
db.persist_result(r2)
row2 = db._conn.execute(
"SELECT first_seen FROM jobs WHERE job_id='60'"
).fetchone()
assert row1["first_seen"] == row2["first_seen"]
# ---------------------------------------------------------------------------
# counts
# ---------------------------------------------------------------------------
class TestCounts:
def test_empty_db(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
assert db.counts() == {}
def test_counts_grouped_by_status(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
db.persist_result(_result(job_id="1", status=JobStatus.position_found))
db.persist_result(_result(job_id="2", position_url=None, status=JobStatus.needs_review))
db.persist_result(_result(job_id="3", position_url=None, status=JobStatus.needs_review))
c = db.counts()
assert c["position_found"] == 1
assert c["needs_review"] == 2
# ---------------------------------------------------------------------------
# export_csv
# ---------------------------------------------------------------------------
class TestExportCsv:
def test_header_matches_csv_columns(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
db.persist_result(_result(job_id="1"))
out = tmp_path / "out.csv"
db.export_csv(path=out)
with out.open() as f:
reader = csv.DictReader(f)
assert tuple(reader.fieldnames or []) == CSV_COLUMNS
def test_null_fields_become_empty_string(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
r = _result(job_id="1", career_url=None, position_url=None, status=JobStatus.needs_review)
db.persist_result(r)
out = tmp_path / "out.csv"
db.export_csv(path=out)
with out.open() as f:
rows = list(csv.DictReader(f))
assert rows[0]["career_page_url"] == ""
assert rows[0]["open_position_url"] == ""
def test_complete_rows_sorted_first(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
# Insert incomplete row first, complete row second
db.persist_result(_result(
job_id="1", company_name="ZZZ Corp",
position_url=None, status=JobStatus.needs_review,
))
db.persist_result(_result(
job_id="2", company_name="AAA Inc",
position_url="https://aaa.com/jobs/1", status=JobStatus.position_found,
))
out = tmp_path / "out.csv"
db.export_csv(path=out)
with out.open() as f:
rows = list(csv.DictReader(f))
assert len(rows) == 2
# complete row (AAA Inc) must come first
assert rows[0]["open_position_url"] != ""
assert rows[1]["open_position_url"] == ""
def test_creates_parent_dirs(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
db.persist_result(_result())
nested = tmp_path / "deep" / "dir" / "results.csv"
db.export_csv(path=nested)
assert nested.exists()
def test_returns_path(self, tmp_path: Path) -> None:
with _db(tmp_path) as db:
db.persist_result(_result())
out = tmp_path / "r.csv"
returned = db.export_csv(path=out)
assert returned == out
# ---------------------------------------------------------------------------
# pipeline.run_batch — offline unit tests via monkeypatch
# ---------------------------------------------------------------------------
def _make_raw(job_id: str, company: str = "Acme") -> RawJob:
return RawJob(
job_id=job_id,
company=company,
linkedin_url=f"https://www.linkedin.com/jobs/view/{job_id}",
website=None,
)
class TestRunBatch:
"""Monkeypatched run_batch tests — no network, no real providers."""
def _patch_all(self, monkeypatch, raw_jobs: list[RawJob]) -> MagicMock:
"""Patch all external I/O in pipeline so tests run offline."""
from jobsource.careers.cascade import CareersResult
# Fake job source
fake_source = MagicMock()
fake_source.fetch_recent_jobs.return_value = raw_jobs
monkeypatch.setattr("jobsource.pipeline._make_source", lambda _: fake_source)
# Fake resolve: always succeeds
monkeypatch.setattr(
"jobsource.pipeline.resolve_website",
lambda name, website=None, *, client=None: f"https://{name.lower().replace(' ', '')}.com",
)
# Fake cascade: returns a careers URL + free position shortcut
def _fake_careers(website, *, company_name=None, client=None):
return CareersResult(
careers_url=f"{website}/careers",
confidence=0.95,
method="ats:greenhouse",
position_url=f"{website}/careers/jobs/1",
)
monkeypatch.setattr("jobsource.pipeline.find_careers_page", _fake_careers)
# extract_open_position shouldn't be called (shortcut covers it),
# but patch defensively to detect if it is.
monkeypatch.setattr(
"jobsource.pipeline.extract_open_position",
lambda url, *, client=None: (None, "none"),
)
# build_client — return a context-manager stub
cm = MagicMock()
cm.__enter__ = lambda s: MagicMock()
cm.__exit__ = MagicMock(return_value=False)
monkeypatch.setattr("jobsource.pipeline.build_client", lambda: cm)
return fake_source
def test_new_jobs_processed(self, tmp_path: Path, monkeypatch) -> None:
from jobsource.pipeline import run_batch
# Use distinct company names so each job gets its own company_key (no cache sharing).
raw = [_make_raw("1", company="Alpha"), _make_raw("2", company="Beta")]
self._patch_all(monkeypatch, raw)
db = Database(path=tmp_path / "t.db")
summary = run_batch(batch_size=5, search_terms=["eng"], db=db)
assert summary.fetched == 2
assert summary.new == 2
assert summary.already_seen == 0
assert summary.website_resolved == 2
assert summary.careers_found == 2
assert summary.position_found == 2
assert summary.failed == 0
assert summary.coverage == 1.0
def test_cross_run_dedup(self, tmp_path: Path, monkeypatch) -> None:
"""Re-running with the same DB yields new == 0 (dedup proven)."""
from jobsource.pipeline import run_batch
raw = [_make_raw("10", company="Gamma"), _make_raw("11", company="Delta")]
self._patch_all(monkeypatch, raw)
db = Database(path=tmp_path / "dedup.db")
# First run
s1 = run_batch(batch_size=5, search_terms=["eng"], db=db)
assert s1.new == 2
# Second run — same raw jobs returned by source, but DB already has them
s2 = run_batch(batch_size=5, search_terms=["eng"], db=db)
assert s2.new == 0
assert s2.already_seen == 2
def test_batch_size_cap(self, tmp_path: Path, monkeypatch) -> None:
"""Only batch_size new jobs are processed even if more are available."""
from jobsource.pipeline import run_batch
# Distinct companies to avoid career-URL cache sharing across jobs.
raw = [_make_raw(str(i), company=f"Co{i}") for i in range(10)]
self._patch_all(monkeypatch, raw)
db = Database(path=tmp_path / "cap.db")
summary = run_batch(batch_size=3, search_terms=["eng"], db=db)
assert summary.new == 3
assert summary.position_found == 3
def test_fetch_limit_matches_batch_size(self, tmp_path: Path, monkeypatch) -> None:
"""Immediate re-runs only re-fetch the requested batch window."""
from jobsource.pipeline import run_batch
raw = [_make_raw(str(i), company=f"Co{i}") for i in range(10)]
source = self._patch_all(monkeypatch, raw)
db = Database(path=tmp_path / "fetch-limit.db")
summary = run_batch(batch_size=3, search_terms=["eng"], db=db)
assert summary.new == 3
source.fetch_recent_jobs.assert_called_once()
assert source.fetch_recent_jobs.call_args.kwargs["results_wanted"] == 3
def test_resolve_failure_marks_needs_review(self, tmp_path: Path, monkeypatch) -> None:
"""Jobs where resolve_website returns None get needs_review, not failed."""
from jobsource.pipeline import run_batch
raw = [_make_raw("99")]
self._patch_all(monkeypatch, raw)
# Override resolve to return None (unresolvable)
monkeypatch.setattr(
"jobsource.pipeline.resolve_website",
lambda *a, **kw: None,
)
db = Database(path=tmp_path / "nr.db")
summary = run_batch(batch_size=5, search_terms=["eng"], db=db)
assert summary.needs_review == 1
assert summary.failed == 0
assert summary.website_resolved == 0
def test_csv_written_after_batch(self, tmp_path: Path, monkeypatch) -> None:
"""Data is persisted by run_batch and export_csv produces the right columns."""
from jobsource.pipeline import run_batch
raw = [_make_raw("7", company="SevenCo")]
self._patch_all(monkeypatch, raw)
csv_path = tmp_path / "out.csv"
# Pass db explicitly so own_db=False → database is NOT closed by run_batch.
db = Database(path=tmp_path / "t.db")
run_batch(batch_size=5, search_terms=["eng"], db=db)
# Export to a controlled path and verify the header/contract.
db.export_csv(path=csv_path)
assert csv_path.exists()
with csv_path.open() as f:
reader = csv.DictReader(f)
assert tuple(reader.fieldnames or []) == CSV_COLUMNS
def test_one_failure_does_not_abort_batch(self, tmp_path: Path, monkeypatch) -> None:
"""A cascade exception for one company doesn't abort others.
Cascade errors land in needs_review (not failed — they are recoverable/transient).
Each job has a distinct company so cache cannot mask the flaky call.
"""
from jobsource.pipeline import run_batch
from jobsource.careers.cascade import CareersResult
# Distinct companies → distinct company_keys → no cache sharing.
raw = [
_make_raw("A", company="AlphaCo"),
_make_raw("B", company="BetaCo"),
_make_raw("C", company="GammaCo"),
]
self._patch_all(monkeypatch, raw)
call_count = {"n": 0}
def _flaky_careers(website, *, company_name=None, client=None):
call_count["n"] += 1
if call_count["n"] == 2:
raise RuntimeError("simulated cascade crash")
return CareersResult(
careers_url=f"{website}/careers",
confidence=0.9,
method="url_pattern",
position_url=f"{website}/careers/jobs/1",
)
monkeypatch.setattr("jobsource.pipeline.find_careers_page", _flaky_careers)
db = Database(path=tmp_path / "fault.db")
summary = run_batch(batch_size=5, search_terms=["eng"], db=db)
# All 3 attempted; cascade exception → needs_review (not failed); 2 fully resolved.
assert summary.new == 3
assert summary.needs_review == 1
assert summary.failed == 0
assert summary.position_found == 2

309
tests/test_extract.py Normal file
View File

@@ -0,0 +1,309 @@
"""Tests for extract.py (Stage 3) — all network-free via monkeypatching."""
from __future__ import annotations
import pytest
from jobsource.careers.ats import ATSBoard, ATSFetch
from jobsource.careers.classify_llm import Anchor, JobLinkResult
from jobsource.extract import _find_job_posting_url, _jsonld_job_url, extract_open_position
# ---------------------------------------------------------------------------
# Fake helpers
# ---------------------------------------------------------------------------
class FakeResponse:
def __init__(self, status_code: int, text: str = "", url: str = "https://example.com") -> None:
self.status_code = status_code
self.text = text
self.url = url
class FakeClient:
pass
# ---------------------------------------------------------------------------
# _find_job_posting_url (pure function)
# ---------------------------------------------------------------------------
class TestFindJobPostingUrl:
def test_flat_dict_job_posting(self):
data = {"@type": "JobPosting", "url": "https://acme.com/jobs/1"}
assert _find_job_posting_url(data) == "https://acme.com/jobs/1"
def test_list_of_nodes(self):
data = [
{"@type": "Organization", "name": "Acme"},
{"@type": "JobPosting", "url": "https://acme.com/jobs/2"},
]
assert _find_job_posting_url(data) == "https://acme.com/jobs/2"
def test_graph_expansion(self):
data = {
"@context": "https://schema.org",
"@graph": [
{"@type": "WebPage"},
{"@type": "JobPosting", "url": "https://acme.com/jobs/3"},
],
}
assert _find_job_posting_url(data) == "https://acme.com/jobs/3"
def test_type_as_list(self):
data = {"@type": ["Thing", "JobPosting"], "url": "https://acme.com/jobs/4"}
assert _find_job_posting_url(data) == "https://acme.com/jobs/4"
def test_no_job_posting_returns_none(self):
data = {"@type": "Organization", "name": "Acme"}
assert _find_job_posting_url(data) is None
def test_job_posting_without_url_returns_none(self):
data = {"@type": "JobPosting"}
assert _find_job_posting_url(data) is None
def test_scalar_input_returns_none(self):
assert _find_job_posting_url("not a dict") is None
assert _find_job_posting_url(42) is None
# ---------------------------------------------------------------------------
# _jsonld_job_url (parses HTML)
# ---------------------------------------------------------------------------
class TestJsonldJobUrl:
def test_finds_url_in_json_ld_block(self):
html = """
<html><head>
<script type="application/ld+json">
{"@type": "JobPosting", "url": "https://acme.com/jobs/99"}
</script>
</head></html>
"""
assert _jsonld_job_url(html) == "https://acme.com/jobs/99"
def test_invalid_json_block_skipped(self):
html = """
<html><head>
<script type="application/ld+json">NOT JSON{{</script>
<script type="application/ld+json">{"@type": "JobPosting", "url": "https://ok.com/job/1"}</script>
</head></html>
"""
assert _jsonld_job_url(html) == "https://ok.com/job/1"
def test_no_job_posting_returns_none(self):
html = """
<html><head>
<script type="application/ld+json">{"@type": "Organization"}</script>
</head></html>
"""
assert _jsonld_job_url(html) is None
def test_no_script_returns_none(self):
assert _jsonld_job_url("<html><body>No scripts</body></html>") is None
# ---------------------------------------------------------------------------
# extract_open_position — tier integration (network-free)
# ---------------------------------------------------------------------------
class TestExtractOpenPositionTiers:
def _patch_base(self, monkeypatch, *, html: str | None = None) -> None:
"""Patch ATS detect functions to miss and HTTP fetch to return html."""
monkeypatch.setattr(
"jobsource.extract._ats.detect_ats_in_url",
lambda url: None,
)
monkeypatch.setattr(
"jobsource.extract._ats.detect_ats_in_html",
lambda html: None,
)
if html is not None:
monkeypatch.setattr(
"jobsource.extract.request_with_retries",
lambda *a, **kw: FakeResponse(200, html),
)
# -- Tier 1: ATS from URL ------------------------------------------------
def test_ats_from_url_returns_first_job(self, monkeypatch):
board = ATSBoard(
ats_name="greenhouse",
slug="acme",
careers_url="https://boards.greenhouse.io/acme",
)
monkeypatch.setattr(
"jobsource.extract._ats.detect_ats_in_url",
lambda url: board,
)
monkeypatch.setattr(
"jobsource.extract._ats._FETCH_DISPATCH",
{"greenhouse": lambda b, c: ATSFetch(first_url="https://boards.greenhouse.io/acme/jobs/1", job_count=5)},
)
url, method = extract_open_position(
"https://boards.greenhouse.io/acme", client=FakeClient()
)
assert url == "https://boards.greenhouse.io/acme/jobs/1"
assert method == "ats:greenhouse"
def test_ats_from_url_with_no_jobs_falls_through(self, monkeypatch):
board = ATSBoard(
ats_name="lever",
slug="acme",
careers_url="https://jobs.lever.co/acme",
)
monkeypatch.setattr("jobsource.extract._ats.detect_ats_in_url", lambda url: board)
monkeypatch.setattr(
"jobsource.extract._ats._FETCH_DISPATCH",
{"lever": lambda b, c: ATSFetch(first_url=None, job_count=0)},
)
monkeypatch.setattr(
"jobsource.extract._ats.detect_ats_in_html", lambda html: None
)
# Page has no JSON-LD, no job-path anchors, LLM returns None.
# Use a neutral href that won't trigger the job-path pattern.
monkeypatch.setattr(
"jobsource.extract.request_with_retries",
lambda *a, **kw: FakeResponse(200, "<html><body><a href='/team'>Team</a></body></html>"),
)
monkeypatch.setattr("jobsource.extract.classify_job_link", lambda anchors: None)
url, method = extract_open_position("https://acme.example.com/careers", client=FakeClient())
assert url is None
assert method == "none"
# -- Tier 1b: ATS from page HTML ----------------------------------------
def test_ats_from_html_returns_first_job(self, monkeypatch):
gh_board = ATSBoard(
ats_name="greenhouse",
slug="vercel",
careers_url="https://boards.greenhouse.io/vercel",
)
monkeypatch.setattr("jobsource.extract._ats.detect_ats_in_url", lambda url: None)
monkeypatch.setattr("jobsource.extract._ats.detect_ats_in_html", lambda html: gh_board)
monkeypatch.setattr(
"jobsource.extract._ats._FETCH_DISPATCH",
{"greenhouse": lambda b, c: ATSFetch(first_url="https://boards.greenhouse.io/vercel/jobs/42", job_count=73)},
)
monkeypatch.setattr(
"jobsource.extract.request_with_retries",
lambda *a, **kw: FakeResponse(200, '<script src="greenhouse"></script>'),
)
url, method = extract_open_position("https://vercel.com/careers", client=FakeClient())
assert url == "https://boards.greenhouse.io/vercel/jobs/42"
assert method == "ats:greenhouse"
# -- Tier 2: JSON-LD ----------------------------------------------------
def test_jsonld_returns_job_url(self, monkeypatch):
html = """
<html><head>
<script type="application/ld+json">
{"@type": "JobPosting", "url": "https://acme.com/jobs/77"}
</script>
</head><body></body></html>
"""
self._patch_base(monkeypatch, html=html)
url, method = extract_open_position("https://acme.com/careers", client=FakeClient())
assert url == "https://acme.com/jobs/77"
assert method == "json_ld"
# -- Tier 3: job-like anchor --------------------------------------------
def test_job_anchor_matched_by_pattern(self, monkeypatch):
html = """
<html><body>
<a href="/about">About</a>
<a href="/jobs/software-engineer-123">Software Engineer</a>
</body></html>
"""
self._patch_base(monkeypatch, html=html)
url, method = extract_open_position("https://acme.com/careers", client=FakeClient())
assert url == "https://acme.com/jobs/software-engineer-123"
assert method == "anchor"
def test_opening_anchor_matched(self, monkeypatch):
html = '<html><body><a href="/openings/42">Open Role</a></body></html>'
self._patch_base(monkeypatch, html=html)
url, method = extract_open_position("https://acme.com/careers", client=FakeClient())
assert url == "https://acme.com/openings/42"
assert method == "anchor"
# -- Tier 4: LLM --------------------------------------------------------
def test_llm_classify_fires_when_no_job_anchor(self, monkeypatch):
html = '<html><body><a href="/about">About</a><a href="/team">Team</a></body></html>'
self._patch_base(monkeypatch, html=html)
monkeypatch.setattr(
"jobsource.extract.classify_job_link",
lambda anchors: JobLinkResult(job_url="https://acme.com/about", confidence=0.60),
)
url, method = extract_open_position("https://acme.com/careers", client=FakeClient())
assert url == "https://acme.com/about"
assert method == "llm_classify"
def test_llm_returns_none_falls_through_to_none(self, monkeypatch):
html = '<html><body><a href="/about">About</a></body></html>'
self._patch_base(monkeypatch, html=html)
monkeypatch.setattr(
"jobsource.extract.classify_job_link",
lambda anchors: None,
)
url, method = extract_open_position("https://acme.com/careers", client=FakeClient())
assert url is None
assert method == "none"
# -- Page fetch failure -------------------------------------------------
def test_fetch_failure_returns_none(self, monkeypatch):
monkeypatch.setattr("jobsource.extract._ats.detect_ats_in_url", lambda url: None)
monkeypatch.setattr(
"jobsource.extract.request_with_retries",
lambda *a, **kw: (_ for _ in ()).throw(RuntimeError("timeout")),
)
url, method = extract_open_position("https://acme.com/careers", client=FakeClient())
assert url is None
assert method == "none"
def test_non_200_fetch_returns_none(self, monkeypatch):
monkeypatch.setattr("jobsource.extract._ats.detect_ats_in_url", lambda url: None)
monkeypatch.setattr(
"jobsource.extract.request_with_retries",
lambda *a, **kw: FakeResponse(404, ""),
)
url, method = extract_open_position("https://acme.com/careers", client=FakeClient())
assert url is None
assert method == "none"
# -- All tiers miss ------------------------------------------------------
def test_all_miss_returns_none_none(self, monkeypatch):
html = '<html><body><a href="/about">About</a></body></html>'
self._patch_base(monkeypatch, html=html)
monkeypatch.setattr("jobsource.extract.classify_job_link", lambda anchors: None)
url, method = extract_open_position("https://acme.com/careers", client=FakeClient())
assert url is None
assert method == "none"
# -- ATS tier error falls through ---------------------------------------
def test_ats_fetch_error_falls_through(self, monkeypatch):
board = ATSBoard(ats_name="ashby", slug="acme", careers_url="https://jobs.ashbyhq.com/acme")
monkeypatch.setattr("jobsource.extract._ats.detect_ats_in_url", lambda url: board)
monkeypatch.setattr(
"jobsource.extract._ats._FETCH_DISPATCH",
{"ashby": lambda b, c: (_ for _ in ()).throw(RuntimeError("api down"))},
)
monkeypatch.setattr("jobsource.extract._ats.detect_ats_in_html", lambda html: None)
html = '<html><body><a href="/jobs/42">Engineer</a></body></html>'
monkeypatch.setattr(
"jobsource.extract.request_with_retries",
lambda *a, **kw: FakeResponse(200, html),
)
url, method = extract_open_position("https://jobs.ashbyhq.com/acme", client=FakeClient())
# urljoin("https://jobs.ashbyhq.com/acme", "/jobs/42") → root-relative
assert url == "https://jobs.ashbyhq.com/jobs/42"
assert method == "anchor"

159
validate.py Normal file
View File

@@ -0,0 +1,159 @@
#!/usr/bin/env python
"""Live validation against the CLAUDE.md success criteria.
Runs a small batch through the full pipeline, then immediately re-runs to
verify dedup. Prints PASS/FAIL for each criterion and exits non-zero if
any check fails.
Usage:
.venv/bin/python validate.py
.venv/bin/python validate.py --batch-size 10 --search "data engineer"
.venv/bin/python validate.py --help
"""
from __future__ import annotations
import argparse
import csv
import sys
import tempfile
from pathlib import Path
sys.path.insert(0, str(Path(__file__).parent))
from jobsource.db import Database
from jobsource.models import CSV_COLUMNS
from jobsource.pipeline import run_batch
def _check(label: str, passed: bool, detail: str = "") -> bool:
mark = "PASS" if passed else "FAIL"
suffix = f"{detail}" if detail else ""
print(f" [{mark}] {label}{suffix}")
return passed
def main() -> int:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--batch-size", type=int, default=5,
help="Number of new jobs to process (default: 5)")
parser.add_argument("--search", default="software engineer",
help="Search term (default: 'software engineer')")
parser.add_argument("--location", default="United States",
help="Location filter (default: 'United States')")
args = parser.parse_args()
print(f"\n{'='*64}")
print(f" validate.py batch={args.batch_size} search={args.search!r}")
print(f"{'='*64}\n")
checks: list[bool] = []
with tempfile.TemporaryDirectory() as tmp:
db_path = Path(tmp) / "validate.db"
csv_path = Path(tmp) / "results.csv"
# ------------------------------------------------------------------
# Run 1: initial batch
# ------------------------------------------------------------------
print("Run 1 (initial batch)…")
s1 = exc1 = None
try:
with Database(path=db_path) as db:
s1 = run_batch(
batch_size=args.batch_size,
search_terms=[args.search],
location=args.location,
db=db,
)
db.export_csv(path=csv_path)
except Exception as exc:
exc1 = exc
# [1] Batch completes without unhandled exception
checks.append(_check(
"Batch completes without unhandled exception",
exc1 is None,
repr(exc1) if exc1 else "",
))
if exc1 is not None:
print("\nFatal error in run 1 — cannot continue.")
return 1
# [2] CSV has exactly the three contract columns
ok_csv = False
csv_detail = "not written"
if csv_path.exists():
with csv_path.open() as f:
reader = csv.DictReader(f)
rows = list(reader)
actual = tuple(reader.fieldnames or [])
ok_csv = actual == CSV_COLUMNS
csv_detail = (
f"columns={actual}" if not ok_csv
else f"{len(rows)} data rows columns={actual}"
)
checks.append(_check("CSV has exactly the three contract columns", ok_csv, csv_detail))
# [3] Per-stage summary with coverage metric
ok_summary = s1 is not None
checks.append(_check(
"Per-stage summary returned",
ok_summary,
(f"coverage={s1.coverage*100:.0f}% "
f"({s1.position_found}/{s1.new} new jobs fully resolved)") if s1 else "",
))
# warn if no jobs were fetched (network/API issue, not a hard failure)
if s1 and s1.new == 0:
print(" NOTE: 0 new jobs fetched — dedup check may be vacuously true.")
print(" Check your search terms, network, or hours_old setting.\n")
# ------------------------------------------------------------------
# Run 2: dedup check
# ------------------------------------------------------------------
print("Run 2 (dedup check — same DB)…")
s2 = exc2 = None
try:
with Database(path=db_path) as db2:
s2 = run_batch(
batch_size=args.batch_size,
search_terms=[args.search],
location=args.location,
db=db2,
)
except Exception as exc:
exc2 = exc
# [4] Re-run processes 0 new jobs
ok_dedup = exc2 is None and s2 is not None and s2.new == 0
checks.append(_check(
"Re-run processes 0 new jobs (dedup proven)",
ok_dedup,
(f"new={s2.new if s2 else 'N/A'}, seen={s2.already_seen if s2 else 'N/A'}"
if not ok_dedup
else f"0 new / {s2.already_seen} already seen"),
))
# ------------------------------------------------------------------
# Spot-check reminder (manual step — not auto-checkable)
# ------------------------------------------------------------------
if s1 and s1.position_found > 0:
print(
f"\n Manual spot-check: open output/results.csv and verify that\n"
f" career_page_url and open_position_url return HTTP 200 in a browser."
)
# ------------------------------------------------------------------
# Verdict
# ------------------------------------------------------------------
n_pass = sum(checks)
n_total = len(checks)
print(f"\n{'='*64}")
overall = all(checks)
print(f" Overall: {'PASS' if overall else 'FAIL'} ({n_pass}/{n_total} automated checks passed)")
print(f"{'='*64}\n")
return 0 if overall else 1
if __name__ == "__main__":
sys.exit(main())