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.
341 lines
14 KiB
Python
341 lines
14 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 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.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
|
|
|
|
|
|
@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
|
|
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)
|
|
|
|
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)
|
|
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
|
|
|
|
# ── 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
|
|
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)
|
|
|
|
result.results_count = len(discovered)
|
|
logger.info(
|
|
f"[proactive] Discovery: {result.queries_run} queries, "
|
|
f"{result.results_count} results"
|
|
)
|
|
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 and keep one result per canonical URL."""
|
|
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 kept
|
|
|
|
# ── 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."""
|
|
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,
|
|
"snippet": normalize_whitespace(cand.snippet)[:2000] or None,
|
|
"posted_text": cand.posted_text,
|
|
"salary_text": cand.salary_text,
|
|
"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"],
|
|
source_engine=lead["source_engine"],
|
|
source_query=lead["source_query"],
|
|
snippet=lead.get("snippet"),
|
|
salary_text=lead.get("salary_text"),
|
|
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.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)
|
|
logger.info(f"[proactive] Report saved: {path}")
|
|
if 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}") |