You cannot select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

741 lines
30 KiB
Python

"""Proactive search orchestrator: discovery → dedup → verification → ledger."""
from __future__ import annotations
import random
import time
from dataclasses import dataclass, field
from datetime import date, datetime
from itertools import zip_longest
from typing import Optional
from loguru import logger
from gimme_job.config import GlobalConfig
from gimme_job.proactive.engines import EngineBlockedError, SearchResult, build_engine, resolve_url
from gimme_job.proactive.plan import ProactivePlan
from gimme_job.proactive.roles import (
NON_CLINICAL as NON_CLINICAL_ROLE,
)
from gimme_job.proactive.roles import (
SUPPORT as SUPPORT_ROLE,
)
from gimme_job.proactive.roles import (
classify_role,
record_role_candidate,
)
from gimme_job.proactive.verifier import (
ACTIVE,
CLOSED,
FILTERED,
STALE,
VERIFY,
evidence_to_dict,
verify,
)
from gimme_job.utils.hashing import compute_fingerprint
from gimme_job.utils.text import normalize_whitespace
from gimme_job.utils.urls import canonical_lead_url
# Host substrings that indicate an official ATS posting (preferred primary)
_ATS_HOST_MARKERS = (
"myworkdayjobs",
"icims.com",
"greenhouse.io",
"lever.co",
"smartrecruiters",
"workdayjobs",
"jobvite.com",
)
# SearchResult.source values emitted by the direct ATS board fetchers
_BOARD_SOURCES = (
"greenhouse",
"lever",
"smartrecruiters",
"workday",
"icims",
"usajobs",
)
@dataclass
class ProactiveRunResult:
run_date: date = field(default_factory=date.today)
mode: str = "daily"
started_at: datetime = field(default_factory=datetime.utcnow)
ended_at: Optional[datetime] = None
queries_run: int = 0
results_count: int = 0
board_results_count: int = 0
boards_discovered: int = 0
verified_count: int = 0
new_count: int = 0
reopened_count: int = 0
still_open_count: int = 0
verify_count: int = 0
closed_count: int = 0
stale_count: int = 0
filtered_count: int = 0
engine_blocks: int = 0
errors: list[str] = field(default_factory=list)
report_path: Optional[str] = None
role_audit_path: Optional[str] = None
def finish(self) -> None:
self.ended_at = datetime.utcnow()
class ProactiveEngine:
def __init__(self, global_config: GlobalConfig, session_factory=None):
self.cfg = global_config
self.session_factory = session_factory
def run(
self,
mode: str = "daily",
dry_run: bool = False,
limit_queries: Optional[int] = None,
max_verify: Optional[int] = None,
) -> ProactiveRunResult:
"""Run the proactive discovery + verification pipeline once."""
from gimme_job.db.engine import init_db
init_db() # idempotent — ensures new tables exist on existing DBs
plan = ProactivePlan.load()
result = ProactiveRunResult(mode=mode)
if mode == "weekly":
queries = plan.deepscan.generate_queries()
verify_cap = plan.deepscan.budget.max_verify_pages
else:
queries = plan.generate_queries(result.run_date)
verify_cap = plan.budget.max_verify_pages
if limit_queries:
queries = queries[:limit_queries]
if not queries:
logger.warning("No queries generated — check sites/proactive.yaml")
result.finish()
return result
if max_verify is not None:
verify_cap = max_verify
logger.info(f"Proactive run: {len(queries)} queries, mode={mode}, dry_run={dry_run}")
from gimme_job.runtime.browser import BrowserManager
bm = BrowserManager(
profile_name=self.cfg.runtime.profile_name,
headless=self.cfg.runtime.headless,
slow_mo=self.cfg.runtime.slow_mo_ms,
)
try:
bm.open_context()
page = bm.new_page()
page.set_default_timeout(self.cfg.runtime.default_timeout_ms)
page.set_default_navigation_timeout(self.cfg.runtime.navigation_timeout_ms)
discovered = self._discover(page, plan, queries, result)
candidates = self._dedupe_discovered(discovered, plan)
logger.info(f"[proactive] {len(candidates)} unique candidate URLs")
verified = self._verify_candidates(page, plan, candidates, result, max_pages=verify_cap)
if not dry_run:
self._persist(verified, result)
self._record_run(result)
self._deliver_report(result, plan)
self._run_role_audit(plan, result)
finally:
bm.close()
result.finish()
logger.info(
f"[proactive] Done: {result.queries_run} queries, {result.verified_count} verified, "
f"NEW={result.new_count} REOPENED={result.reopened_count} "
f"STILL_OPEN={result.still_open_count} VERIFY={result.verify_count} "
f"CLOSED={result.closed_count} STALE={result.stale_count} "
f"filtered={result.filtered_count} blocks={result.engine_blocks}"
)
return result
def reverify(
self,
dry_run: bool = False,
limit: Optional[int] = None,
) -> ProactiveRunResult:
"""Re-open previously CLOSED/STALE leads to detect reopenings (doc §15.3).
A lead that went inactive can come back (requisition reopened, reposted).
This pass re-verifies its URL and, if it reads as ACTIVE again, transitions
it to REOPENED through the normal upsert path.
"""
from gimme_job.db.engine import db_session, init_db
from gimme_job.models.dto import ProactiveLeadCandidate
init_db()
plan = ProactivePlan.load()
result = ProactiveRunResult(mode="reverify")
from gimme_job.runtime.browser import BrowserManager
bm = BrowserManager(
profile_name=self.cfg.runtime.profile_name,
headless=self.cfg.runtime.headless,
slow_mo=self.cfg.runtime.slow_mo_ms,
)
try:
bm.open_context()
page = bm.new_page()
page.set_default_timeout(self.cfg.runtime.default_timeout_ms)
page.set_default_navigation_timeout(self.cfg.runtime.navigation_timeout_ms)
from gimme_job.proactive.verifier import ACTIVE, verify
with db_session() as session:
from gimme_job.db.proactive_repo import ProactiveLeadRepo
recs = ProactiveLeadRepo().get_inactive(session, limit=limit)
leads = [
{
"id": rec.id,
"fingerprint": rec.fingerprint,
"official_url": rec.official_url,
"employer": rec.employer,
"title": rec.title,
"state": rec.state,
"source_engine": rec.source_engine,
"source_query": rec.source_query,
"snippet": rec.snippet,
"salary_text": rec.salary_text,
"fte": rec.fte,
"license_requirement": rec.license_requirement,
}
for rec in recs
]
logger.info(f"[proactive] reverify: {len(leads)} inactive leads")
reopened = 0
for lead in leads:
if not lead["official_url"]:
continue
result.verified_count += 1
logger.info(f"[proactive] reverify: {lead['official_url']}")
try:
vr = verify(page, lead["official_url"], plan)
except Exception as e:
logger.warning(f"[proactive] reverify error {lead['official_url']}: {e}")
result.errors.append(str(e)[:200])
continue
if vr.status not in (ACTIVE,):
# still closed/stale — no transition
result.closed_count += 1
continue
candidate = ProactiveLeadCandidate(
status=vr.status,
employer=lead["employer"],
title=vr.evidence.title or lead["title"],
state=lead["state"],
official_url=lead["official_url"],
source_engine=lead["source_engine"],
source_query=lead["source_query"],
snippet=lead["snippet"],
salary_text=vr.evidence.salary_text or lead["salary_text"],
fte=vr.evidence.fte or lead["fte"],
license_requirement=vr.evidence.license_requirement
or lead["license_requirement"],
evidence={**evidence_to_dict(vr.evidence), "reason": vr.reason},
fingerprint=lead["fingerprint"],
)
if dry_run:
logger.info(f"[proactive] reverify would reopen: {lead['title']}")
reopened += 1
continue
with db_session() as session:
changed, final_status = ProactiveLeadRepo().upsert(
session, candidate, result.run_date
)
if changed and final_status == "REOPENED":
result.reopened_count += 1
reopened += 1
else:
result.still_open_count += 1
if not dry_run:
self._record_run(result)
logger.info(
f"[proactive] reverify done: {len(leads)} checked, "
f"{result.reopened_count} reopened, {result.closed_count} still closed"
)
return result
finally:
bm.close()
# ── Discovery ──────────────────────────────────────────────────────────
def _discover(
self,
page,
plan: ProactivePlan,
queries: list[str],
result: ProactiveRunResult,
) -> list[SearchResult]:
discovered: list[SearchResult] = []
seen_urls: set[str] = set()
blocked_engines: set[str] = set()
for query in queries:
for engine_name in plan.enabled_engines():
if engine_name in blocked_engines:
continue # circuit breaker: skip engines blocked this run
engine = build_engine(engine_name)
cfg = plan.engine_config(engine_name)
try:
results = engine.search(page, query, cfg.max_results_per_query)
result.queries_run += 1
for r in results:
r.url = resolve_url("https://www.google.com", r.url)
canon = canonical_lead_url(r.url)
if not canon or canon in seen_urls:
continue
role, role_reason = classify_role(r.title, plan)
if role in (SUPPORT_ROLE, NON_CLINICAL_ROLE):
record_role_candidate(
r.title, r.url, r.source, role, role_reason
)
continue
# Unknown SERP pages are logged by the board adapters
# only — page titles from search results are usually
# navigational, not job titles (audit noise).
seen_urls.add(canon)
discovered.append(r)
except EngineBlockedError as e:
# Engine anti-bot wall — not a real error; skip engine for rest of run
blocked_engines.add(engine_name)
result.engine_blocks += 1
logger.warning(
f"[proactive] {engine_name} blocked ({e}) — "
f"skipping it for the rest of this run"
)
except Exception as e:
logger.warning(f"[proactive] {engine_name} query failed: {e}")
result.errors.append(str(e)[:200])
self._politeness_delay(cfg.delay_min_ms, cfg.delay_max_ms)
self._politeness_delay(plan.engine_config("duckduckgo").delay_min_ms, 3000)
# Direct ATS boards — hidden jobs not indexed by search engines
if plan.boards.enabled:
discovered.extend(self._discover_boards(plan, result, seen_urls))
if result.mode == "weekly" and plan.boards.discovery.enabled:
discovered.extend(self._discover_new_boards(page, plan, result, seen_urls))
result.results_count = len(discovered)
logger.info(
f"[proactive] Discovery: {result.queries_run} queries, "
f"{result.results_count} results, {result.board_results_count} board matches"
)
return discovered
def _discover_boards(
self,
plan: ProactivePlan,
result: ProactiveRunResult,
seen_urls: set[str],
) -> list[SearchResult]:
"""Fetch employer ATS boards directly — catches postings hidden from SERPs.
Board results are appended to the same dedup pool so they flow through
the normal verify → ledger pipeline.
"""
from gimme_job.proactive.boards import fetch_board
discovered: list[SearchResult] = []
boards = plan.employers.all_boards()
if not boards:
return discovered
logger.info(f"[proactive] Fetching {len(boards)} ATS boards...")
for cfg in boards:
try:
results = fetch_board(cfg)
except Exception as e:
logger.warning(f"[proactive] Board {cfg.source}:{cfg.name} failed: {e}")
result.errors.append(str(e)[:200])
continue
for r in results:
r.url = resolve_url("https://www.google.com", r.url)
canon = canonical_lead_url(r.url)
if not canon or canon in seen_urls:
continue
seen_urls.add(canon)
discovered.append(r)
result.board_results_count += len(results)
self._politeness_delay(
plan.boards.delay_min_ms, plan.boards.delay_max_ms
)
if discovered:
logger.info(f"[proactive] Boards → {len(discovered)} new candidate URLs")
return discovered
def _discover_new_boards(
self,
page,
plan: ProactivePlan,
result: ProactiveRunResult,
seen_urls: set[str],
) -> list[SearchResult]:
"""Weekly: grow the ATS catalog via site: queries, then fetch new boards.
Discovery probes each candidate first; verified boards are written to
sites/employers.auto.yaml (machine-owned) and fetched in the same run
so newly found hidden jobs flow straight into verification.
"""
from gimme_job.proactive.boards import fetch_board
from gimme_job.proactive.discovery import run_discovery, write_auto_boards
report = run_discovery(
page,
plan,
engine=build_engine(plan.boards.discovery.engine),
keywords=plan.boards.discovery.keywords,
max_queries=plan.boards.discovery.max_queries,
min_keyword_hits=plan.boards.discovery.min_keyword_hits,
)
result.queries_run += report.queries_run
result.errors.extend(report.errors)
verified = report.verified()
result.boards_discovered = len(verified)
if not verified:
return []
if plan.boards.discovery.write:
path, added = write_auto_boards(verified)
logger.info(
f"[proactive] Board discovery: {len(verified)} verified, "
f"{added} new → {path}"
)
discovered: list[SearchResult] = []
for cand in verified:
try:
results = fetch_board(cand.to_config())
except Exception as e:
logger.warning(f"[proactive] New board {cand.label()} failed: {e}")
result.errors.append(str(e)[:200])
continue
for r in results:
r.url = resolve_url("https://www.google.com", r.url)
canon = canonical_lead_url(r.url)
if not canon or canon in seen_urls:
continue
seen_urls.add(canon)
discovered.append(r)
result.board_results_count += len(results)
self._politeness_delay(plan.boards.delay_min_ms, plan.boards.delay_max_ms)
if discovered:
logger.info(
f"[proactive] New boards → {len(discovered)} candidate URLs"
)
return discovered
@staticmethod
def _politeness_delay(min_ms: int, max_ms: int) -> None:
delay = random.randint(min_ms, max_ms) / 1000.0
time.sleep(delay)
@staticmethod
def _dedupe_discovered(
discovered: list[SearchResult], plan: ProactivePlan
) -> list[SearchResult]:
"""Drop denylisted employers, dedup by canonical URL, then merge
cross-source duplicates of the same posting by (title, employer).
When the same job is found on both an official ATS and an aggregator
mirror, the official URL is kept as the primary and the alternates are
recorded in secondary_urls (doc §19).
"""
seen: set[str] = set()
kept: list[SearchResult] = []
for r in discovered:
canon = canonical_lead_url(r.url)
if not canon or canon in seen:
continue
seen.add(canon)
if plan.is_denylisted(r.title):
logger.debug(f"[proactive] Denylisted: {r.title}")
continue
kept.append(r)
return ProactiveEngine._merge_duplicates(kept, plan)
@staticmethod
def _merge_duplicates(
results: list[SearchResult], plan: ProactivePlan
) -> list[SearchResult]:
"""Merge postings that appear on multiple sources (doc §19)."""
aggregator_domains = plan.verification.aggregator_domains or []
buckets: dict[tuple[str, str], list[SearchResult]] = {}
for r in results:
title_key = normalize_whitespace(r.title).lower() or r.url
employer_key = normalize_whitespace(r.employer_hint or "").lower()
buckets.setdefault((employer_key, title_key), []).append(r)
merged: list[SearchResult] = []
for (employer_key, _title), group in buckets.items():
if len(group) == 1:
merged.append(group[0])
continue
# representative: official ATS domain > more complete URL
def score(r: SearchResult) -> int:
host = r.url.lower()
if any(agg in host for agg in aggregator_domains):
return 0 # aggregator mirror
if any(plat in host for plat in _ATS_HOST_MARKERS):
return 2 # official ATS
return 1 # generic employer site
group.sort(key=score, reverse=True)
primary = group[0]
alts = [r.url for r in group[1:]]
for alt in alts:
if alt not in primary.secondary_urls:
primary.secondary_urls.append(alt)
merged.append(primary)
logger.debug(
f"[proactive] merged {len(group)} sources for '{primary.title}' "
f"→ primary {primary.url}"
)
return merged
@staticmethod
def _prioritize_for_verification(
candidates: list[SearchResult],
) -> list[SearchResult]:
"""Interleave board results with SERP results (boards first).
Board results are appended after every SERP result at discovery time;
with a fixed verify budget they were never reached (verified_count was
consumed entirely by SERP candidates). Alternating keeps both sources
represented while giving hidden board jobs the first slots.
"""
boards = [c for c in candidates if c.source in _BOARD_SOURCES]
serp = [c for c in candidates if c.source not in _BOARD_SOURCES]
ordered: list[SearchResult] = []
for board, search in zip_longest(boards, serp):
if board is not None:
ordered.append(board)
if search is not None:
ordered.append(search)
return ordered
# ── Verification ───────────────────────────────────────────────────────
def _verify_candidates(
self,
page,
plan: ProactivePlan,
candidates: list[SearchResult],
result: ProactiveRunResult,
max_pages: int,
) -> list[dict]:
"""Verify each candidate URL. Returns list of lead dicts to persist."""
candidates = self._prioritize_for_verification(candidates)
board_n = sum(1 for c in candidates if c.source in _BOARD_SOURCES)
logger.info(
f"[proactive] Verify queue: {board_n} board + "
f"{len(candidates) - board_n} SERP candidates, cap {max_pages}"
)
leads: list[dict] = []
for i, cand in enumerate(candidates[:max_pages]):
total = min(len(candidates), max_pages)
logger.info(f"[proactive] Verifying {i + 1}/{total}: {cand.url}")
try:
vr = verify(page, cand.url, plan)
except Exception as e:
logger.warning(f"[proactive] Verify error for {cand.url}: {e}")
result.errors.append(str(e)[:200])
continue
result.verified_count += 1
status = vr.status
if status == VERIFY:
result.verify_count += 1
elif status == CLOSED:
result.closed_count += 1
elif status == STALE:
result.stale_count += 1
elif status == FILTERED:
result.filtered_count += 1
if status == FILTERED:
continue
leads.append(
{
"status": status,
"title": normalize_whitespace(vr.evidence.title) or cand.title,
"url": canonical_lead_url(vr.evidence.final_url or cand.url),
"state": plan.state_for_query(cand.query),
"employer_hint": cand.employer_hint,
"source_engine": cand.source,
"source_query": cand.query,
"secondary_urls": cand.secondary_urls,
"snippet": normalize_whitespace(cand.snippet)[:2000] or None,
"posted_text": cand.posted_text,
"salary_text": vr.evidence.salary_text or cand.salary_text,
"fte": vr.evidence.fte,
"license_requirement": vr.evidence.license_requirement,
"evidence": {**evidence_to_dict(vr.evidence), "reason": vr.reason},
"screenshot_path": vr.screenshot_path,
"dom_path": vr.dom_path,
"reason": vr.reason,
}
)
self._politeness_delay(1500, 3000)
return leads
def _persist(self, verified: list[dict], result: ProactiveRunResult) -> None:
"""Persist leads. Only ACTIVE leads are transitioned here — VERIFY/CLOSED/
STALE counts were already recorded at classification time."""
from gimme_job.db.engine import db_session
from gimme_job.db.proactive_repo import ProactiveLeadRepo
from gimme_job.models.dto import ProactiveLeadCandidate
repo = ProactiveLeadRepo()
run_date = result.run_date
with db_session() as session:
for lead in verified:
if lead["status"] != ACTIVE:
continue
candidate = ProactiveLeadCandidate(
status=lead["status"],
employer=lead.get("employer_hint"),
title=lead["title"],
state=lead.get("state"),
official_url=lead["url"],
secondary_urls=lead.get("secondary_urls") or [],
source_engine=lead["source_engine"],
source_query=lead["source_query"],
snippet=lead.get("snippet"),
salary_text=lead.get("salary_text"),
fte=lead.get("fte"),
license_requirement=lead.get("license_requirement"),
evidence=lead["evidence"],
fingerprint=self._fingerprint(lead["url"]),
)
changed, final_status = repo.upsert(session, candidate, run_date)
if not changed:
continue
if final_status == "NEW":
result.new_count += 1
elif final_status == "REOPENED":
result.reopened_count += 1
elif final_status == "STILL_OPEN":
result.still_open_count += 1
@staticmethod
def _fingerprint(url: str) -> str:
return compute_fingerprint("proactive", None, None, None, url, url)
def _record_run(self, result: ProactiveRunResult) -> None:
try:
from gimme_job.db.engine import db_session
from gimme_job.db.proactive_repo import ProactiveRunRepo
with db_session() as session:
ProactiveRunRepo().record_run(
session,
started_at=result.started_at,
ended_at=result.ended_at,
mode=result.mode,
status="success" if not result.errors else "partial",
queries_run=result.queries_run,
results_count=result.results_count,
verified_count=result.verified_count,
new_count=result.new_count,
reopened_count=result.reopened_count,
still_open_count=result.still_open_count,
verify_count=result.verify_count,
closed_count=result.closed_count,
stale_count=result.stale_count,
filtered_count=result.filtered_count,
error_summary="; ".join(result.errors[:5]) or None,
)
except Exception as e:
logger.warning(f"[proactive] Failed to record run: {e}")
def _deliver_report(self, result: ProactiveRunResult, plan: ProactivePlan) -> None:
"""Build and deliver the daily report. Never finishes silently (doc §23)."""
try:
from gimme_job.proactive.html_report import terminal_link
from gimme_job.proactive.report import build_report, deliver_report
report_text = build_report(result.run_date, plan)
path = deliver_report(
report_text, result.run_date, self.session_factory, plan=plan
)
if path:
result.report_path = str(path)
logger.info(f"[proactive] Report saved: {terminal_link(path, path)}")
try:
from gimme_job.db.engine import db_session
from gimme_job.db.proactive_repo import ProactiveRunRepo
with db_session() as session:
latest = ProactiveRunRepo().get_recent(session, limit=1)
if latest:
latest[0].report_path = str(path)
except Exception as e:
logger.debug(f"[proactive] Report path update failed: {e}")
except Exception as e:
logger.error(f"[proactive] Report delivery failed: {e}")
def _run_role_audit(self, plan: ProactivePlan, result: ProactiveRunResult) -> None:
"""Daily offline LLM pass over borderline titles (never blocks the run).
gemma4 reviews unknown/support/non_clinical candidates logged during
discovery. Support/non-clinical proposals are auto-applied to the
machine-owned role_terms.auto.yaml; target proposals stay in the audit
report for human review (they would loosen the gate).
"""
cfg = plan.role_audit
if not cfg.enabled:
return
try:
from gimme_job.proactive.html_report import terminal_link
from gimme_job.proactive.role_audit import append_auto_terms, run_audit
path, parsed, entries = run_audit(
result.run_date,
model=cfg.model,
base_url=cfg.base_url,
max_titles=cfg.max_titles,
)
result.role_audit_path = str(path)
logger.info(
f"[proactive] Role audit: {len(entries)} candidates → "
f"{terminal_link(path, path)}"
)
proposals = parsed.get("proposals") or []
applicable = [p for p in proposals if p.get("role") in cfg.apply_roles]
if cfg.write and applicable:
auto_path, added = append_auto_terms(applicable)
logger.info(
f"[proactive] Role audit: {added} term(s) applied → {auto_path}"
)
skipped = len(proposals) - len(applicable)
if skipped:
logger.info(
f"[proactive] Role audit: {skipped} proposal(s) report-only "
"(target/unknown — review the report)"
)
except ValueError as e:
logger.debug(f"[proactive] Role audit skipped: {e}")
except Exception as e:
logger.warning(f"[proactive] Role audit failed (non-fatal): {e}")