From 3b54f996ca4a448f6357ab8de9900d1902a1c826 Mon Sep 17 00:00:00 2001 From: skillfactor-pipeline Date: Tue, 7 Jul 2026 11:33:43 +0200 Subject: [PATCH] feat(batch): budget/resume module, QA harness, runner skeleton + architecture design Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01RCcxND1mMWu6Lt2c2PxexN --- data/progress.json | 4 + docs/batch-architecture.md | 69 ++++++++++++++ docs/qa/testset_fable-issues.log | 0 docs/qa/testset_fable-sample.md | 154 +++++++++++++++++++++++++++++++ pipeline/batch_run.py | 106 +++++++++++++++++++++ pipeline/progress.py | 66 +++++++++++++ pipeline/qa_sample.py | 102 ++++++++++++++++++++ 7 files changed, 501 insertions(+) create mode 100644 data/progress.json create mode 100644 docs/batch-architecture.md create mode 100644 docs/qa/testset_fable-issues.log create mode 100644 docs/qa/testset_fable-sample.md create mode 100644 pipeline/batch_run.py create mode 100644 pipeline/progress.py create mode 100644 pipeline/qa_sample.py diff --git a/data/progress.json b/data/progress.json new file mode 100644 index 00000000..344eeb68 --- /dev/null +++ b/data/progress.json @@ -0,0 +1,4 @@ +{ + "jsearch_requests_total": 0, + "occupations": {} +} \ No newline at end of file diff --git a/docs/batch-architecture.md b/docs/batch-architecture.md new file mode 100644 index 00000000..dc52547f --- /dev/null +++ b/docs/batch-architecture.md @@ -0,0 +1,69 @@ +# Batch architecture — full-catalog build-out (design, Fable 5, 2026-07-07) + +Design is final; Sonnet implements the TODOs in `pipeline/batch_run.py` +without redesigning. Contracts below are what the skeleton, `progress.py`, +`extract_local.py` and `qa_sample.py` already assume. + +## Staging + +| Stage | Scope | LLM | Where | +|---|---|---|---| +| 1 packages | all 3 039 ESCO occupations | none | local + throttled Gitea push | +| 2 evidence | all occupations, ~60 ads each | Ollama (extraction only) | local + RTX-3090 box | +| 3 depth tier 1 | top-200 by evidence volume | Claude (per templates/) | Claude session | +| 3 depth tier 3 | rest, night batch | Claude, resumable | Claude sessions | + +Order within stage 2 follows `docs/batch-plan.md` (HR/IT -> finance -> +management), then the remaining catalog by ISCO group. + +## Budgets & throttles + +- **JSearch: hard cap 33 000 requests lifetime**, enforced in + `progress.spend_request()` BEFORE each HTTP call — not advisory. + ~6 requests/occupation x 3 039 = ~18 k expected, cap leaves headroom. +- Rate: <= 2 req/s (sleep 0.5 s), cache first — a cache hit costs nothing. +- Thin market (< 20 usable ads): ESCO-synonym fallback -> Adzuna (free) -> + mark `evidence: failed / thin market`, move on. No endless retries. +- Gitea: >= 1 s between repo operations, check-before-create (resumable). +- Ollama: sequential, one ad per call (~6.5 s measured on gemma3:27b); + ~180 k ads ≈ 14 days GPU time — run as detached script on the webserver + (survives Claude sessions), NOT inside the Claude loop. + +## Data contracts + +- `data/raw/jobs//jsearch___p.json` — raw cache + (existing recruiter files stay where they are; new fetches use subdirs) +- `data/raw/jobs/_ads.json` — deduped `[{job_id, title, employer, + country, description}]` (same shape as recruiter_ads.json) +- `data/evidence/.jsonl` — one validated extraction per line + (schema: `pipeline/prompts/extract_schema.json` + `job_id`) +- `data/progress.json` — see `pipeline/progress.py` docstring; the ONLY + place run state lives. Deleting it = full re-scan, but no re-spend + (caches + file existence still short-circuit). + +## MSSQL change (prerequisite for stage 2) + +```sql +ALTER TABLE evidence_job ADD occupation_slug NVARCHAR(200) NULL; +UPDATE evidence_job SET occupation_slug = 'recruitment-consultant' + WHERE occupation_slug IS NULL; -- backfill: all 350 existing rows are recruiter +CREATE INDEX ix_evidence_job_occ ON evidence_job(occupation_slug); +``` + +`p3b_store_evidence.py` becomes incremental: MERGE on +`(occupation_slug, job_id)` instead of drop-and-rebuild; `evidence_entity` +keyed by job_id as today. `p3c_aggregate.py` parameterized by slug. + +## QA cadence + +After each occupation's extraction: `qa_sample.py --rate 0.02` (min 1 +record). Claude reviews the generated `docs/qa/-sample.md` in the next +session pass; two consecutive harness failures (exit 1) = stop the batch, +diagnose (systemic prompt/model drift, not noise). Depth files are checked +against `pipeline/templates/QUALITY_BAR.md` (mechanical part greppable). + +## Failure policy + +Any per-occupation failure is recorded (`error` field) and skipped — the +batch never blocks on a single occupation. `--slugs` re-runs individual +failures after a fix. BudgetExhausted stops the run with exit 2. diff --git a/docs/qa/testset_fable-issues.log b/docs/qa/testset_fable-issues.log new file mode 100644 index 00000000..e69de29b diff --git a/docs/qa/testset_fable-sample.md b/docs/qa/testset_fable-sample.md new file mode 100644 index 00000000..bd41adab --- /dev/null +++ b/docs/qa/testset_fable-sample.md @@ -0,0 +1,154 @@ +# QA sample — testset_fable + +2 of 20 records (10%). Judge each extraction against QUALITY_BAR.md 'Extraction quality': faithful, no invented items, plausible seniority. + +## Recruiter - Government Contracting (us) + +**Ad (first 2000 chars):** + +> Join Titania as a Full-Cycle Recruiter supporting technical and professional services hiring. Manage a focused workload (2–8 requisitions) while partnering closely with hiring managers and growing your expertise in a collaborative, structured environment. +> +> Position Overview +> +> Titania is seeking a Recruiter to manage full-cycle recruiting for a consistent workload of approximately 2–8 open requisitions at a time. This role supports hiring across technical and professional services positions and partners closely with hiring managers to deliver effective hiring outcomes. +> +> This opportunity is well-suited for a recruiting professional who is interested in developing expertise in full-cycle recruiting, technical recruiting, and applicant tracking systems, while working in a structured and relationship-driven environment. +> +> Key Responsibilities +> +> Full-Cycle Recruiting +> • Manage the full recruiting lifecycle, including intake, sourcing, screening, interviewing, and offer support +> • Partner with hiring managers to define job requirements, qualifications, and hiring timelines +> • Create and post job descriptions aligned with business needs and compliance standards +> • Source candidates using job boards, LinkedIn, Clearancejobs.com, and other recruiting tools +> • Conduct phone screens and coordinate structured interview processes +> • Maintain candidate pipelines for current and future hiring needs +> +> Applicant Tracking System (ATS) & Reporting +> • Utilize Greenhouse for job postings, candidate tracking, interview scheduling, and documentation +> • Ensure data accuracy, workflow consistency, and audit readiness +> • Track recruiting metrics such as time-to-fill, pipeline activity, and candidate flow +> +> Stakeholder Management +> • Build strong working relationships with hiring managers and internal stakeholders +> • Provide guidance on recruiting best practices, interview techniques, and candidate evaluation +> • Communicate regularly on hiring progress and market conditions +> +> Interview & Offer Coordination +> • Sch + +**Extraction:** + +```json +{ + "hard_skills": [ + "full-cycle recruiting", + "sourcing", + "screening", + "interviewing", + "job description creation", + "candidate pipeline management", + "recruiting metrics tracking" + ], + "tools": [ + "Greenhouse", + "LinkedIn / LinkedIn Recruiter", + "Clearancejobs.com" + ], + "methods": [ + "structured interviewing", + "applicant tracking system (ATS)", + "compliance standards", + "process improvement" + ], + "responsibilities": [ + "intake process management", + "hiring manager partnership", + "job posting", + "candidate sourcing", + "interview scheduling", + "offer support", + "recruiting record maintenance", + "stakeholder communication" + ], + "qualifications": [ + "2+ years recruiting experience", + "experience with applicant tracking system (ATS)" + ], + "soft_skills": [ + "communication skills", + "interpersonal skills", + "organizational skills", + "relationship building" + ], + "seniority": "mid" +} +``` + +## Senior Recruiter (us) + +**Ad (first 2000 chars):** + +> About Us +> +> When military bases, national security agencies, federal buildings, healthcare facilities, and leading commercial clients in Washington D.C and across the nation need superior facility support services, we answer the call. We offer a unique combination of government experience and facilities expertise that has helped us earn a position as one of the country's leading integrated service providers. +> Job Summary +> +> EMCOR Government Services is seeking a Senior Recruiter in Alexandria, VA. This is a hybrid role requiring the successful candidate to work from the Alexandria VA office a minimum of 3 days a week. +> +> The Senior Recruiter is a vital part of a dynamic, interactive, and engaged Talent Acquisition and Human Resources team. The Senior Recruiter leads full-cycle recruiting strategy and execution to support a fast-paced government contracting environment. This role is responsible for building scalable recruiting processes, managing high-volume and cleared hiring pipelines, and partnering with business leaders to meet mission-critical workforce needs. The ideal candidate brings deep experience recruiting for federal contracts, including cleared personnel, and understands compliance requirements. The Senior Recruiter plays a critical role in not only filling current operational talent needs but also works very closely with business development to source candidates to compete and win new work. +> +> #EGS +> +> #LI-NS1 +> +> Essential Duties & Responsibilities +> +> Works collaboratively to develop and execute enterprise-wide recruiting and sourcing strategies to address business needs, strategic priorities, and contract requirements. Maintains a broad pipeline of talented candidates. Builds a network of talent through varied sources. Supports executive leadership, program managers, and capture teams forecasted hiring needs and workforce plans Source and hire high-quality candidates for leadership, managerial, professional, and technical roles across a diverse array of organizationa + +**Extraction:** + +```json +{ + "hard_skills": [ + "full lifecycle recruiting", + "sourcing", + "screening", + "interviewing", + "talent pipeline management", + "compliance", + "workforce planning" + ], + "tools": [ + "iCIMS", + "LinkedIn Recruiter", + "Indeed", + "Cleared Jobs", + "Microsoft Excel", + "Microsoft Word" + ], + "methods": [ + "social media sourcing", + "internet sourcing", + "progressive sourcing methods", + "performance metrics optimization" + ], + "responsibilities": [ + "recruiting strategy execution", + "candidate pipeline building", + "talent pool development", + "offer negotiation", + "hiring manager consultation", + "compliance adherence" + ], + "qualifications": [ + "bachelor's degree or equivalent", + "5+ years recruiting experience", + "experience with government contracts", + "experience sourcing cleared personnel", + "experience with government proposals" + ], + "soft_skills": [], + "seniority": "senior" +} +``` diff --git a/pipeline/batch_run.py b/pipeline/batch_run.py new file mode 100644 index 00000000..8fafe6c4 --- /dev/null +++ b/pipeline/batch_run.py @@ -0,0 +1,106 @@ +"""Full-catalog batch runner — SKELETON (Fable 5, 2026-07-07). + +Sonnet: implement the TODO functions, do not redesign the flow. The staging +logic, resume mechanics and budget enforcement are final; see +docs/batch-architecture.md for the design and data contracts. + +Stages (per skillfactor_finalize.md / auftrag_sonnet.md): + 1. packages: deterministic generator for ALL ESCO occupations (no LLM) + 2. evidence: per occupation fetch (JSearch, budget-capped) -> extract + (Ollama, extract_local.py) -> aggregate (market.md) + 3. depth: tier 1 = top-200 by evidence volume, tier 3 = night batch + +Every step is idempotent and resumable: progress.json + file existence +decide what still needs work; an abort at any point loses at most the +current occupation's in-flight step. + +Usage: + python pipeline/batch_run.py --stage packages [--limit N] + python pipeline/batch_run.py --stage evidence [--slugs a,b] [--limit N] + python pipeline/batch_run.py --stage tier # (re)compute tiers only +""" +import argparse +import os +import subprocess +import sys + +sys.path.insert(0, os.path.dirname(__file__)) +import progress +from db import connect + +BASE = os.path.join(os.path.dirname(__file__), "..") +ADS_PER_OCC_TARGET = 60 # ~6 requests/occupation (1 req ~ 10 ads) +COUNTRIES = ("us", "gb") +REQ_RATE_SLEEP = 0.5 # ~2 req/s max (finalize order) +TOP_TIER_SIZE = 200 + + +def esco_occupations(cur): + """All ESCO occupations with slug + altLabels for query building.""" + # TODO(sonnet): return [(slug, preferred_label, [alt_labels...]), ...] + # from esco_occupation; slug generation must match p2_generate.py + raise NotImplementedError + + +def stage_packages(state, limit=0): + """Run the existing generator for every occupation without a package. + TODO(sonnet): refactor p2_generate.py so it can be called per-occupation + (it currently generates a fixed balanced set of 100). Reuse, don't fork. + Gitea pushes go through p4_publish.py, throttled (sleep >= 1s/repo, + resumable by checking repo existence via API before create).""" + raise NotImplementedError + + +def stage_evidence(state, slugs=None, limit=0): + """Per occupation: fetch (cached) -> extract (Ollama) -> aggregate. + + fetch: generalize p3a (queries = preferredLabel + top-3 altLabels x + COUNTRIES, cache data/raw/jobs//, budget via + progress.spend_request BEFORE each HTTP call). + < 20 usable ads -> synonym fallback, then Adzuna (free), then + mark evidence="failed", error="thin market" and move on. + extract: subprocess extract_local.py --ads data/raw/jobs/_ads.json + --out data/evidence/.jsonl (already implemented + tested) + store: extend p3b: incremental MERGE keyed by (occupation_slug, job_id) + -- needs ALTER TABLE evidence_job ADD occupation_slug (see + docs/batch-architecture.md, backfill 'recruitment-consultant') + aggregate: generalize p3c per occupation (market.md + market-evidence + sections); QA: run qa_sample.py per batch, stop the run if the + harness exits 1 twice in a row (systemic problem, not noise). + """ + raise NotImplementedError + + +def stage_tier(state): + """Rank occupations for depth staging: tier 1 = TOP_TIER_SIZE by + (ads collected, crosswalk strength, marketplace demand), rest tier 3. + Pure bookkeeping in progress.json — no API calls.""" + raise NotImplementedError + + +def main(): + ap = argparse.ArgumentParser() + ap.add_argument("--stage", required=True, + choices=["packages", "evidence", "tier"]) + ap.add_argument("--slugs", default="") + ap.add_argument("--limit", type=int, default=0) + args = ap.parse_args() + + state = progress.load() + try: + if args.stage == "packages": + stage_packages(state, args.limit) + elif args.stage == "evidence": + slugs = [s for s in args.slugs.split(",") if s] or None + stage_evidence(state, slugs, args.limit) + else: + stage_tier(state) + except progress.BudgetExhausted as exc: + print(f"STOP: {exc}") + sys.exit(2) + finally: + progress.save(state) + + +if __name__ == "__main__": + main() diff --git a/pipeline/progress.py b/pipeline/progress.py new file mode 100644 index 00000000..76d84a93 --- /dev/null +++ b/pipeline/progress.py @@ -0,0 +1,66 @@ +"""Resume + budget bookkeeping for the full-catalog batch runs. + +Single source of truth: data/progress.json +{ + "jsearch_requests_total": 123, # lifetime counter, hard cap below + "occupations": { + "": { + "package": "done|pending", # phase 1 (deterministic generator) + "evidence": "done|fetching|extracting|aggregating|failed|pending", + "depth": "done|in_progress|failed|pending", # phases 2b/3 + "tier": 1|2|3, # 1 = top-200, 3 = night batch rest + "requests": 7, # JSearch requests spent on this slug + "ads": 61, # unique ads with full text + "error": "last error message" # only present after a failure + }, ... + } +} + +Writes are atomic (tmp + os.replace) so an aborted run never corrupts the +file. Every JSearch request MUST go through spend_request() — it enforces +the hard budget. Cache hits cost nothing and are not counted. +""" +import json +import os + +BASE = os.path.join(os.path.dirname(__file__), "..") +PROGRESS = os.path.join(BASE, "data", "progress.json") + +JSEARCH_BUDGET_TOTAL = 33000 # hard cap per finalize order (skillfactor_finalize.md) + + +class BudgetExhausted(RuntimeError): + pass + + +def load(): + if os.path.exists(PROGRESS): + with open(PROGRESS, encoding="utf-8") as f: + return json.load(f) + return {"jsearch_requests_total": 0, "occupations": {}} + + +def save(state): + tmp = PROGRESS + ".tmp" + with open(tmp, "w", encoding="utf-8") as f: + json.dump(state, f, ensure_ascii=False, indent=1) + os.replace(tmp, PROGRESS) + + +def occ(state, slug): + return state["occupations"].setdefault(slug, { + "package": "pending", "evidence": "pending", "depth": "pending", + "tier": 3, "requests": 0, "ads": 0, + }) + + +def spend_request(state, slug, n=1): + """Count n JSearch requests against the global and per-slug budget. + Call BEFORE the HTTP request; raises BudgetExhausted at the cap.""" + if state["jsearch_requests_total"] + n > JSEARCH_BUDGET_TOTAL: + raise BudgetExhausted( + f"JSearch budget {JSEARCH_BUDGET_TOTAL} reached " + f"({state['jsearch_requests_total']} spent)") + state["jsearch_requests_total"] += n + occ(state, slug)["requests"] += n + save(state) diff --git a/pipeline/qa_sample.py b/pipeline/qa_sample.py new file mode 100644 index 00000000..32708d61 --- /dev/null +++ b/pipeline/qa_sample.py @@ -0,0 +1,102 @@ +"""QA harness — 2 % spot-check sample per extraction batch. + +Mechanical checks run here; faithfulness (does the extraction match the ad?) +is judged by Claude reading the generated review file. Deviations from +QUALITY_BAR.md ("Extraction quality" section) are logged, not auto-fixed. + +Usage: + python pipeline/qa_sample.py --extractions data/evidence/.jsonl \ + --ads data/raw/jobs/_ads.json [--rate 0.02] [--seed ] + +Output: + docs/qa/-sample.md ad text + extraction side-by-side for review + docs/qa/-issues.log mechanical deviations (one line each) + +Deterministic sampling (seeded by batch name) so re-runs pick the same +records and reviews stay comparable. +""" +import argparse +import json +import os +import random +import sys + +sys.path.insert(0, os.path.dirname(__file__)) +from extract_local import ARRAY_FIELDS, SENIORITY + +BASE = os.path.join(os.path.dirname(__file__), "..") +QA_DIR = os.path.join(BASE, "docs", "qa") + + +def mechanical_issues(rec): + issues = [] + hard = {x.lower() for x in rec.get("hard_skills", [])} + soft = {x.lower() for x in rec.get("soft_skills", [])} + tools = {x.lower() for x in rec.get("tools", [])} + if hard & soft: + issues.append(f"cross-field dupe hard/soft: {sorted(hard & soft)}") + if hard & tools: + issues.append(f"cross-field dupe hard/tools: {sorted(hard & tools)}") + if rec.get("seniority") not in SENIORITY: + issues.append(f"invalid seniority {rec.get('seniority')!r}") + for field, (max_items, max_len) in ARRAY_FIELDS.items(): + vals = rec.get(field, []) + if len(vals) > max_items: + issues.append(f"{field}: {len(vals)} items > {max_items}") + for v in vals: + if "�" in v or len(v) > max_len: + issues.append(f"{field}: bad item {v!r}") + # lowercase rule: only flag all-caps sentences, products are fine + if v.isupper() and len(v) > 6: + issues.append(f"{field}: shouting {v!r}") + return issues + + +def main(): + ap = argparse.ArgumentParser() + ap.add_argument("--extractions", required=True) + ap.add_argument("--ads", required=True) + ap.add_argument("--rate", type=float, default=0.02) + ap.add_argument("--seed", default=None, help="defaults to extractions filename") + args = ap.parse_args() + + batch = os.path.splitext(os.path.basename(args.extractions))[0] + rng = random.Random(args.seed or batch) + + recs = [json.loads(l) for l in open(args.extractions, encoding="utf-8") if l.strip()] + ads = {a["job_id"]: a for a in json.load(open(args.ads, encoding="utf-8"))} + n = max(1, round(len(recs) * args.rate)) + sample = rng.sample(recs, min(n, len(recs))) + + os.makedirs(QA_DIR, exist_ok=True) + issues_total = 0 + review = [f"# QA sample — {batch}", "", + f"{len(sample)} of {len(recs)} records ({args.rate:.0%}). " + "Judge each extraction against QUALITY_BAR.md 'Extraction quality': " + "faithful, no invented items, plausible seniority.", ""] + with open(os.path.join(QA_DIR, f"{batch}-issues.log"), "w", encoding="utf-8") as log: + for rec in sample: + ad = ads.get(rec["job_id"], {}) + probs = mechanical_issues(rec) + issues_total += len(probs) + for p in probs: + log.write(f"{rec['job_id']}\t{p}\n") + review += [f"## {ad.get('title', rec['job_id'])} ({ad.get('country','?')})", "", + "**Ad (first 2000 chars):**", "", + "> " + (ad.get("description", "AD TEXT MISSING")[:2000] + ).replace("\n", "\n> "), "", + "**Extraction:**", "", "```json", + json.dumps({k: v for k, v in rec.items() if k != "job_id"}, + ensure_ascii=False, indent=1), + "```", ""] + if probs: + review += ["**Mechanical issues:** " + "; ".join(probs), ""] + out = os.path.join(QA_DIR, f"{batch}-sample.md") + open(out, "w", encoding="utf-8").write("\n".join(review)) + print(f"{len(sample)} sampled, {issues_total} mechanical issues -> {out}") + if issues_total > len(sample): # more than 1 issue per record on average + sys.exit(1) + + +if __name__ == "__main__": + main()