Implement BDEF v1.1 grading: scoring core, per-deck pipeline, ledger, dashboard, StartOS layer
- Deterministic scoring.py (quant 60 / qual 40 / flags -15, profitability heaviest) - Per-company JSON ledger with forecast-target chaining deck N-1 -> N - Single-shot sandbox agent with guided-JSON fallback ladder (no tool loop) - Portfolio dashboard with sparklines, KPI hit rates, BDEF category bars - 48 unit tests green; endpoints smoke-tested; npm check+build green Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Fable 5
parent
1dde915540
commit
b1d7aed9f4
+425
-153
@@ -1,40 +1,61 @@
|
||||
"""The Boardroom Map job runner — convenes the review panel over dropped documents.
|
||||
"""The Boardroom Map job runner — grades dropped board decks against the BDEF.
|
||||
|
||||
Runs as a background thread inside the FastAPI app. It does NOT run on a clock
|
||||
like Nightshift; it reacts to triggers:
|
||||
Runs as a background thread inside the FastAPI app. It does NOT run on a clock;
|
||||
it reacts to triggers:
|
||||
|
||||
* an explicit "Run Review" (drops /data/state/run_request), or
|
||||
* autoRunOnDrop: files landing in /data/inbox, once the inbox is stable.
|
||||
* an explicit "Grade Decks" run request (drops /data/state/run_request), or
|
||||
* autoRunOnDrop: files landing under /data/inbox/<company-slug>/, once the
|
||||
inbox is stable across two ticks.
|
||||
|
||||
One job at a time. A job:
|
||||
1. extract text from the inbox locally (CPU) — only text crosses to the Sparks
|
||||
2. rsync the text to a per-job dir on the head Spark
|
||||
3. serve the needed models in WAVES; run the reviewers for each wave
|
||||
4. optionally run the local lead-reviewer synthesis
|
||||
5. pull the reports back to /data/reports/<job>, assemble latest.md
|
||||
6. wipe the documents from the Sparks (unless disabled) and tear serving down
|
||||
One job at a time. A job iterates the discovered deck units OLDEST FIRST per
|
||||
company, and for each deck:
|
||||
|
||||
All state (phase, current job, per-reviewer status, last report) is mirrored to
|
||||
1. extract text locally (CPU) — only text crosses to the Sparks
|
||||
2. rsync the per-deck bundle (docs/, BDEF.md, personas/, schemas/, out/,
|
||||
adjudicator-out/) to {remoteWorkDir}/jobs/<job>/<company>/<deck>/
|
||||
3. serve the needed models in WAVES (graders' models ∪ the extractor model);
|
||||
the extractor runs when its model's wave is up, graders in their waves
|
||||
4. pull out/, validate: extraction.json invalid => the DECK fails (the job
|
||||
continues); grader JSONs are validated individually, invalid ones dropped;
|
||||
fewer than 2 valid grade reports => the deck fails
|
||||
5. adjudicate (optional, non-fatal): a local model weighs the panel's evidence
|
||||
6. score deterministically (scoring.score_deck) against the company's pinned
|
||||
targets + the prior deck's forward targets, and record it in the ledger
|
||||
7. render DECK_REPORT.md + refresh the company SCORECARD.md and the
|
||||
/data/reports copies
|
||||
|
||||
Then it wipes the remote job dir (unless disabled), tears serving down, and
|
||||
moves the graded originals to /data/processed/<job>/<slug>/ (the company folder
|
||||
stays in the inbox for reuse). One deck's failure never kills the job: the job
|
||||
ends "done" if at least one deck was graded.
|
||||
|
||||
All state (phase, per-deck status, panel status, last report) is mirrored to
|
||||
/data/state/runtime.json so the Web UI can render it.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import shutil
|
||||
import threading
|
||||
import time
|
||||
import traceback
|
||||
from collections import deque
|
||||
from datetime import datetime
|
||||
from datetime import datetime, timezone
|
||||
|
||||
import adjudicator as adj_mod
|
||||
import bm_config
|
||||
import decks
|
||||
import extraction
|
||||
import graders as gr_mod
|
||||
import ledger as ledger_mod
|
||||
import preflight
|
||||
import reviewers as rev_mod
|
||||
import scorecard
|
||||
import scoring
|
||||
import serving
|
||||
import spark_client as sc
|
||||
import synthesis as synth_mod
|
||||
import validate
|
||||
|
||||
DATA_DIR = os.environ.get("BM_DATA_DIR", "/data")
|
||||
INBOX = os.path.join(DATA_DIR, "inbox")
|
||||
@@ -42,40 +63,77 @@ PROCESSED = os.path.join(DATA_DIR, "processed")
|
||||
STATE_DIR = os.path.join(DATA_DIR, "state")
|
||||
JOBS_DIR = os.path.join(STATE_DIR, "jobs")
|
||||
REPORTS_DIR = os.path.join(DATA_DIR, "reports")
|
||||
LEDGER_DIR = os.path.join(DATA_DIR, "ledger")
|
||||
RUNTIME_PATH = os.path.join(STATE_DIR, "runtime.json")
|
||||
REQUEST_PATH = os.path.join(STATE_DIR, "run_request")
|
||||
|
||||
TICK_SECONDS = 10
|
||||
RUNNING_PHASES = ("extracting", "grading", "adjudicating", "scoring", "collecting")
|
||||
|
||||
# Canonical-ish reporting periods: 2026-Q2, 2026-H1, FY2026, 2026-05, 2026.
|
||||
_PERIOD_RE = re.compile(
|
||||
r"^(?:FY\s?-?\d{4}|\d{4}(?:[-/ ]?(?:Q[1-4]|H[12]|0[1-9]|1[0-2]))?)$", re.IGNORECASE)
|
||||
|
||||
|
||||
def _inbox_signature() -> tuple[int, str]:
|
||||
"""(count, signature) of supported files in the inbox, for stability checks."""
|
||||
"""(count, signature) of supported files anywhere in the inbox tree, for
|
||||
autoRunOnDrop stability checks (decks live in per-company subfolders)."""
|
||||
if not os.path.isdir(INBOX):
|
||||
return (0, "")
|
||||
items = []
|
||||
for fn in sorted(os.listdir(INBOX)):
|
||||
p = os.path.join(INBOX, fn)
|
||||
if os.path.isfile(p) and os.path.splitext(fn)[1].lower() in extraction.SUPPORTED:
|
||||
items.append(f"{fn}:{os.path.getsize(p)}:{int(os.path.getmtime(p))}")
|
||||
for root, _dirs, files in os.walk(INBOX):
|
||||
for fn in files:
|
||||
p = os.path.join(root, fn)
|
||||
if os.path.splitext(fn)[1].lower() in extraction.SUPPORTED:
|
||||
try:
|
||||
items.append(f"{os.path.relpath(p, INBOX)}:{os.path.getsize(p)}:"
|
||||
f"{int(os.path.getmtime(p))}")
|
||||
except OSError:
|
||||
pass
|
||||
items.sort()
|
||||
return (len(items), "|".join(items))
|
||||
|
||||
|
||||
def _token(s: str) -> str:
|
||||
t = re.sub(r"[^a-z0-9]+", "-", (s or "").lower()).strip("-")
|
||||
return t or "deck"
|
||||
|
||||
|
||||
def _composite(record) -> float | None:
|
||||
"""Best-effort composite lookup on the scoring record (shape owned by scoring.py)."""
|
||||
if not isinstance(record, dict):
|
||||
return None
|
||||
for k in ("composite", "composite_score"):
|
||||
v = record.get(k)
|
||||
if isinstance(v, (int, float)):
|
||||
return v
|
||||
for parent in ("scores", "totals", "score"):
|
||||
d = record.get(parent)
|
||||
if isinstance(d, dict) and isinstance(d.get("composite"), (int, float)):
|
||||
return d["composite"]
|
||||
return None
|
||||
|
||||
|
||||
class JobRunner:
|
||||
def __init__(self):
|
||||
self._events = deque(maxlen=500)
|
||||
self._lock = threading.Lock()
|
||||
self.phase = "idle" # idle | extracting | reviewing | synthesizing | collecting | done | error
|
||||
self.phase = "idle" # idle | extracting | grading | adjudicating | scoring | collecting | done | error
|
||||
self.job_id = None
|
||||
self.message = ""
|
||||
self.panel: list[dict] = []
|
||||
self.waves_total = 0
|
||||
self.wave_index = 0
|
||||
self.decks_total = 0
|
||||
self.deck_index = 0
|
||||
self.company = None
|
||||
self.period = None
|
||||
self.decks: list[dict] = []
|
||||
self.last_report_path = None
|
||||
self._thread = None
|
||||
self._last_sig = None
|
||||
self._stable_sig = None
|
||||
self._last_done_sig = None
|
||||
for d in (STATE_DIR, JOBS_DIR, REPORTS_DIR, INBOX, PROCESSED):
|
||||
for d in (STATE_DIR, JOBS_DIR, REPORTS_DIR, LEDGER_DIR, INBOX, PROCESSED):
|
||||
os.makedirs(d, exist_ok=True)
|
||||
self._restore()
|
||||
|
||||
@@ -108,11 +166,12 @@ class JobRunner:
|
||||
self.job_id = d.get("job_id")
|
||||
self.message = d.get("message", "")
|
||||
self.panel = d.get("panel", [])
|
||||
self.decks = d.get("decks", [])
|
||||
self.last_report_path = d.get("last_report_path")
|
||||
for e in d.get("events", []):
|
||||
self._events.append(e)
|
||||
# A job can't survive a restart; reset a stuck running phase.
|
||||
if self.phase in ("extracting", "reviewing", "synthesizing", "collecting"):
|
||||
if self.phase in RUNNING_PHASES:
|
||||
self.phase = "idle"
|
||||
except Exception:
|
||||
pass
|
||||
@@ -125,6 +184,11 @@ class JobRunner:
|
||||
"panel": self.panel,
|
||||
"waves_total": self.waves_total,
|
||||
"wave_index": self.wave_index,
|
||||
"decks_total": self.decks_total,
|
||||
"deck_index": self.deck_index,
|
||||
"company": self.company,
|
||||
"period": self.period,
|
||||
"decks": self.decks,
|
||||
"last_report_path": self.last_report_path,
|
||||
}
|
||||
|
||||
@@ -136,7 +200,7 @@ class JobRunner:
|
||||
self._thread.start()
|
||||
|
||||
def request_run(self):
|
||||
"""Public hook (used by the API) to request a review immediately."""
|
||||
"""Public hook (used by the API) to request a grading run immediately."""
|
||||
try:
|
||||
with open(REQUEST_PATH, "w") as f:
|
||||
f.write(str(time.time()))
|
||||
@@ -160,23 +224,17 @@ class JobRunner:
|
||||
if os.path.exists(REQUEST_PATH):
|
||||
os.remove(REQUEST_PATH)
|
||||
triggered = True
|
||||
self.log("[runner] review requested")
|
||||
self.log("[runner] grading run requested")
|
||||
elif cfg.get("autoRunOnDrop"):
|
||||
count, sig = _inbox_signature()
|
||||
if count and sig == self._last_sig and sig != self._last_done_sig:
|
||||
# stable across two ticks and not the batch we last processed
|
||||
triggered = True
|
||||
self.log("[runner] inbox stable — auto-running review")
|
||||
self.log("[runner] inbox stable — auto-running grading")
|
||||
self._last_sig = sig
|
||||
|
||||
if not triggered:
|
||||
return
|
||||
count, _ = _inbox_signature()
|
||||
if not count:
|
||||
self.log("[runner] nothing to review (inbox empty of supported files)")
|
||||
self.phase = "idle"
|
||||
self._persist()
|
||||
return
|
||||
self._run_job(cfg)
|
||||
|
||||
# ------------------------------------------------------------- the job
|
||||
@@ -186,114 +244,124 @@ class JobRunner:
|
||||
self.message = ""
|
||||
self.waves_total = 0
|
||||
self.wave_index = 0
|
||||
self.decks_total = 0
|
||||
self.deck_index = 0
|
||||
self.company = None
|
||||
self.period = None
|
||||
self.decks = []
|
||||
self.panel = []
|
||||
local_job = os.path.join(JOBS_DIR, job_id)
|
||||
local_docs = os.path.join(local_job, "docs")
|
||||
remote_job = f"{cfg['remoteWorkDir'].rstrip('/')}/jobs/{job_id}"
|
||||
rubric = cfg.get("reviewInstructions") or bm_config.DEFAULT_RUBRIC
|
||||
self.log(f"=== Review job {job_id} begins ===")
|
||||
remote_root = f"{cfg['remoteWorkDir'].rstrip('/')}/jobs/{job_id}"
|
||||
self.log(f"=== Grading job {job_id} begins ===")
|
||||
|
||||
try:
|
||||
# 1. Extract locally (only text crosses to the Sparks).
|
||||
self.phase = "extracting"; self._persist()
|
||||
manifest = extraction.extract_inbox(INBOX, local_docs, self.log)
|
||||
ok_docs = [m for m in manifest if m["ok"]]
|
||||
if not ok_docs:
|
||||
raise RuntimeError("no documents could be extracted (unsupported or empty inbox)")
|
||||
self.log(f"[runner] extracted {len(ok_docs)} document(s)")
|
||||
# 1. Discover deck units (per-company folders; oldest first).
|
||||
disc = decks.discover(INBOX)
|
||||
units = disc.get("units") or []
|
||||
for fn in disc.get("skipped") or []:
|
||||
self.log(f"[runner] WARNING: skipping root-level inbox file '{fn}' — "
|
||||
"decks belong in /data/inbox/<company-slug>/")
|
||||
if not units:
|
||||
raise RuntimeError("nothing to grade — drop decks into /data/inbox/<company-slug>/")
|
||||
|
||||
# 2. Resolve the panel against the model catalog.
|
||||
catalog = {m["alias"] for m in (cfg.get("models") or [])}
|
||||
panel = rev_mod.roster(cfg)
|
||||
# 2. Ledger, merged with the configured portfolio companies.
|
||||
led = ledger_mod.Ledger(LEDGER_DIR)
|
||||
led.merge_config_companies(cfg.get("companies") or [])
|
||||
|
||||
# Resolve the grader panel + extractor model against the catalog.
|
||||
models = cfg.get("models") or []
|
||||
catalog = {m["alias"] for m in models}
|
||||
panel = gr_mod.roster(cfg)
|
||||
valid = [r for r in panel if r["model"] in catalog]
|
||||
invalid = [r for r in panel if r["model"] not in catalog]
|
||||
for r in invalid:
|
||||
self.log(f"[runner] WARNING: reviewer '{r['name']}' uses unknown model '{r['model']}' — skipped")
|
||||
if not valid:
|
||||
raise RuntimeError("no reviewers reference a configured model (see Configure Models/Reviewers)")
|
||||
self.panel = [{"name": r["name"], "model": r["model"], "status": "pending"} for r in valid]
|
||||
for r in panel:
|
||||
if r["model"] not in catalog:
|
||||
self.log(f"[runner] WARNING: grader '{r['name']}' uses unknown model "
|
||||
f"'{r['model']}' — skipped")
|
||||
if len(valid) < 2:
|
||||
raise RuntimeError(
|
||||
"need at least 2 graders referencing configured models — every deck "
|
||||
"requires >= 2 valid grade reports (see Configure Models/Graders)")
|
||||
extractor_model = (cfg.get("extractorModel") or "").strip() or \
|
||||
(models[0]["alias"] if models else "")
|
||||
if extractor_model not in catalog:
|
||||
raise RuntimeError(f"extractor model '{extractor_model}' is not in the model catalog")
|
||||
|
||||
needed = {r["model"] for r in valid}
|
||||
if cfg.get("synthesisEnabled"):
|
||||
sm = synth_mod.pick_model(cfg)
|
||||
if sm:
|
||||
needed.add(sm)
|
||||
needed = {r["model"] for r in valid} | {extractor_model}
|
||||
adjudicate = bool(cfg.get("adjudicatorEnabled"))
|
||||
adj_model = adj_mod.pick_model(cfg) if adjudicate else ""
|
||||
needed_all = needed | ({adj_model} if adjudicate and adj_model else set())
|
||||
|
||||
# Air-gapped mode can't route to second-Spark models (internal net).
|
||||
if cfg.get("networkMode") == "airgapped":
|
||||
cat = {m["alias"]: m for m in (cfg.get("models") or [])}
|
||||
offenders = [a for a in needed if cat.get(a, {}).get("spark") == "secondary"]
|
||||
cat = {m["alias"]: m for m in models}
|
||||
offenders = [a for a in needed_all if cat.get(a, {}).get("spark") == "secondary"]
|
||||
if offenders:
|
||||
raise RuntimeError(
|
||||
"air-gapped mode requires all models on the head Spark, but these are "
|
||||
f"on the secondary: {', '.join(sorted(offenders))}. Move them to the "
|
||||
"primary Spark or switch to local-services mode.")
|
||||
|
||||
# 3. Ship text to the Spark + ensure infra.
|
||||
self.phase = "reviewing"; self._persist()
|
||||
push = sc.push_dir(sc.head(cfg), local_docs, f"{remote_job}/docs")
|
||||
if push.returncode != 0:
|
||||
raise RuntimeError(f"shipping documents to the Spark failed: {push.stderr}")
|
||||
rev_mod.ensure_reviewer_image(cfg, self.log)
|
||||
# 3. Infra once per job.
|
||||
gr_mod.ensure_grader_image(cfg, self.log)
|
||||
serving.ensure_network(cfg, self.log)
|
||||
# In local_services mode, warn early if the optional web_search backend
|
||||
# (SearXNG, self-signed HTTPS) is unreachable — non-fatal.
|
||||
preflight.check_searxng(cfg, self.log)
|
||||
|
||||
# 4. Run the panel in waves (synthesis is handled separately, below).
|
||||
review_aliases = {r["model"] for r in valid}
|
||||
waves = serving.plan_waves(cfg, review_aliases)
|
||||
self.waves_total = len(waves)
|
||||
hf = bm_config.hf_token()
|
||||
collected = []
|
||||
for i, wave in enumerate(waves, 1):
|
||||
self.wave_index = i
|
||||
wave_aliases = {m["alias"] for m in wave}
|
||||
wpanel = [r for r in valid if r["model"] in wave_aliases]
|
||||
self.log(f"[runner] wave {i}/{len(waves)}: models={sorted(wave_aliases)} "
|
||||
f"reviewers={[r['name'] for r in wpanel]}")
|
||||
serving.bring_up_wave(cfg, wave, hf, self.log)
|
||||
self._await_serving(cfg, wave)
|
||||
preflight.check_wave(cfg, wave, self.log)
|
||||
res = rev_mod.run_wave_reviewers(cfg, remote_job, wpanel, rubric, self.log)
|
||||
collected.extend(res)
|
||||
self._mark_panel(res)
|
||||
serving.tear_down_wave(cfg, wave, self.log)
|
||||
|
||||
# 5. Synthesis (its own single-model wave).
|
||||
synth_ok = False
|
||||
if cfg.get("synthesisEnabled"):
|
||||
self.phase = "synthesizing"; self._persist()
|
||||
sm = synth_mod.pick_model(cfg)
|
||||
swave = serving.plan_waves(cfg, {sm})
|
||||
for wave in swave:
|
||||
serving.bring_up_wave(cfg, wave, hf, self.log)
|
||||
self._await_serving(cfg, wave)
|
||||
preflight.check_wave(cfg, wave, self.log)
|
||||
sres = synth_mod.run_synthesis(cfg, remote_job, rubric, self.log)
|
||||
synth_ok = bool(sres.get("report"))
|
||||
serving.tear_down_wave(cfg, wave, self.log)
|
||||
# 4. Grade each deck unit, oldest first. One deck's failure never
|
||||
# kills the job.
|
||||
self.decks_total = len(units)
|
||||
succeeded = 0
|
||||
for idx, unit in enumerate(units, 1):
|
||||
self.deck_index = idx
|
||||
self.company = unit["company_slug"]
|
||||
self.period = unit.get("period")
|
||||
entry = {"company": unit["company_slug"], "period": unit.get("period"),
|
||||
"status": "running"}
|
||||
self.decks.append(entry)
|
||||
self.panel = [{"name": r["name"], "model": r["model"], "status": "pending"}
|
||||
for r in valid]
|
||||
self._persist()
|
||||
try:
|
||||
result = self._grade_deck(cfg, led, job_id, remote_root, unit, idx,
|
||||
valid, extractor_model, adjudicate, hf)
|
||||
entry.update({"status": "done", "period": result["period"],
|
||||
"composite": result["composite"]})
|
||||
succeeded += 1
|
||||
comp = result["composite"]
|
||||
self.log(f"[runner] deck done: {unit['company_slug']} {result['period']}"
|
||||
f" composite={comp if comp is not None else '?'}")
|
||||
except Exception as e:
|
||||
entry.update({"status": "failed", "error": str(e)[:300]})
|
||||
self.log(f"[runner] DECK FAILED ({unit['company_slug']} "
|
||||
f"{unit.get('period') or '?'}): {e}")
|
||||
self.log(traceback.format_exc().splitlines()[-1])
|
||||
self._persist()
|
||||
|
||||
# 6. Collect reports + assemble.
|
||||
# 5. Job-level report.
|
||||
self.phase = "collecting"; self._persist()
|
||||
self._collect(cfg, job_id, remote_job, local_job, valid, manifest, synth_ok)
|
||||
self._write_job_report(job_id)
|
||||
|
||||
# 7. Confidentiality: wipe the documents from the Spark.
|
||||
# 6. Confidentiality: wipe the deck text from the Spark + teardown.
|
||||
if cfg.get("wipeRemoteDocs", True):
|
||||
sc.run(sc.head(cfg), f"rm -rf {remote_job}", timeout=60)
|
||||
self.log("[runner] wiped document text from the Spark")
|
||||
sc.run(sc.head(cfg), f"rm -rf {remote_root}", timeout=120)
|
||||
self.log("[runner] wiped deck text from the Spark")
|
||||
serving.tear_down_all(cfg, self.log)
|
||||
|
||||
# 8. Clear the inbox (move originals aside so they aren't re-reviewed).
|
||||
self._drain_inbox(job_id)
|
||||
# 7. Move the graded originals aside (the company folders stay).
|
||||
self._drain_inbox(job_id, units)
|
||||
self._last_done_sig = _inbox_signature()[1]
|
||||
self.phase = "done"
|
||||
self.message = f"Reviewed {len(ok_docs)} document(s) with {len(valid)} reviewer(s)."
|
||||
self.log(f"=== Review job {job_id} complete ===")
|
||||
|
||||
if succeeded:
|
||||
self.phase = "done"
|
||||
self.message = (f"Graded {succeeded}/{len(units)} deck(s) with "
|
||||
f"{len(valid)} grader(s).")
|
||||
else:
|
||||
self.phase = "error"
|
||||
self.message = f"All {len(units)} deck(s) failed — see the activity log."
|
||||
self.log(f"=== Grading job {job_id} complete ({succeeded}/{len(units)} decks) ===")
|
||||
self._persist()
|
||||
except Exception as e:
|
||||
self.phase = "error"
|
||||
self.message = f"Review failed: {e}"
|
||||
self.message = f"Grading failed: {e}"
|
||||
self.log(f"[runner] JOB FAILED — {e}")
|
||||
self.log(traceback.format_exc().splitlines()[-1])
|
||||
try:
|
||||
@@ -302,6 +370,205 @@ class JobRunner:
|
||||
pass
|
||||
self._persist()
|
||||
|
||||
# ------------------------------------------------------------- one deck
|
||||
def _grade_deck(self, cfg: dict, led, job_id: str, remote_root: str, unit: dict,
|
||||
idx: int, panel: list[dict], extractor_model: str,
|
||||
adjudicate: bool, hf) -> dict:
|
||||
"""Grade one deck unit end to end. Raises on deck failure (caller continues)."""
|
||||
slug_c = unit["company_slug"]
|
||||
# The staging/remote dir token; the final deck_id is the resolved period.
|
||||
token = _token(unit["period"]) if unit.get("period") else f"deck-{idx:02d}"
|
||||
local_deck = os.path.join(JOBS_DIR, job_id, slug_c, token)
|
||||
remote_deck = f"{remote_root}/{slug_c}/{token}"
|
||||
head = sc.head(cfg)
|
||||
|
||||
# --- extract locally + stage + push -----------------------------------
|
||||
self.phase = "extracting"; self._persist()
|
||||
manifest = extraction.extract_files(unit["files"], os.path.join(local_deck, "docs"),
|
||||
self.log)
|
||||
ok_docs = [m for m in manifest if m["ok"]]
|
||||
for m in manifest:
|
||||
if not m["ok"]:
|
||||
self.log(f"[runner] WARNING: {m['source']}: {m['error']}")
|
||||
if not ok_docs:
|
||||
raise RuntimeError("no document text could be extracted from this deck")
|
||||
gr_mod.stage_deck_files(cfg, local_deck, panel)
|
||||
push = sc.push_dir(head, local_deck, remote_deck)
|
||||
if push.returncode != 0:
|
||||
raise RuntimeError(f"shipping deck text to the Spark failed: {push.stderr}")
|
||||
|
||||
# --- serve in waves; extractor + graders run in their model's wave ----
|
||||
self.phase = "grading"; self._persist()
|
||||
waves = serving.plan_waves(cfg, {r["model"] for r in panel} | {extractor_model})
|
||||
self.waves_total = len(waves)
|
||||
for i, wave in enumerate(waves, 1):
|
||||
self.wave_index = i
|
||||
aliases = {m["alias"] for m in wave}
|
||||
wpanel = [r for r in panel if r["model"] in aliases]
|
||||
self.log(f"[runner] wave {i}/{len(waves)}: models={sorted(aliases)} "
|
||||
f"graders={[r['name'] for r in wpanel]}"
|
||||
f"{' +extractor' if extractor_model in aliases else ''}")
|
||||
serving.bring_up_wave(cfg, wave, hf, self.log)
|
||||
try:
|
||||
self._await_serving(cfg, wave)
|
||||
preflight.check_wave(cfg, wave, self.log)
|
||||
if extractor_model in aliases:
|
||||
er = gr_mod.run_extractor(cfg, remote_deck, extractor_model, self.log)
|
||||
if not er.get("report"):
|
||||
raise RuntimeError("extractor produced no extraction.json")
|
||||
if wpanel:
|
||||
res = gr_mod.run_wave_graders(cfg, remote_deck, wpanel, self.log)
|
||||
self._mark_panel(res)
|
||||
finally:
|
||||
serving.tear_down_wave(cfg, wave, self.log)
|
||||
|
||||
# --- pull the panel outputs + validate --------------------------------
|
||||
local_out = os.path.join(local_deck, "out")
|
||||
pull = sc.pull_dir(head, f"{remote_deck}/out", local_out)
|
||||
if pull.returncode != 0:
|
||||
raise RuntimeError(f"pulling panel outputs from the Spark failed: {pull.stderr}")
|
||||
|
||||
ext_path = os.path.join(local_out, "extraction.json")
|
||||
ext_obj, ext_err = validate.validate_file(ext_path, "extraction")
|
||||
if ext_obj is None or ext_err:
|
||||
raise RuntimeError(f"extraction.json invalid: {ext_err or 'missing'}")
|
||||
|
||||
period, deck_id = self._resolve_period(unit, ext_obj, ext_path, token)
|
||||
self.period = period; self._persist()
|
||||
|
||||
grades, panel_meta = [], []
|
||||
for r in panel:
|
||||
gpath = os.path.join(local_out, f"{r['rid']}.json")
|
||||
gobj, gerr = (None, "no output file")
|
||||
if os.path.exists(gpath) and not os.path.exists(gpath + ".invalid"):
|
||||
gobj, gerr = validate.validate_file(gpath, "grades")
|
||||
ok = gobj is not None and not gerr
|
||||
if ok:
|
||||
grades.append(gobj)
|
||||
else:
|
||||
self.log(f"[runner] grader {r['rid']} report dropped: {gerr}")
|
||||
panel_meta.append({"rid": r["rid"], "model": r["model"], "valid": ok})
|
||||
if len(grades) < 2:
|
||||
raise RuntimeError(f"only {len(grades)} valid grade report(s) (need >= 2)")
|
||||
|
||||
# --- adjudication (non-fatal) ------------------------------------------
|
||||
adjudication_md = None
|
||||
if adjudicate:
|
||||
self.phase = "adjudicating"; self._persist()
|
||||
try:
|
||||
adj_model = adj_mod.pick_model(cfg)
|
||||
for wave in serving.plan_waves(cfg, {adj_model}):
|
||||
serving.bring_up_wave(cfg, wave, hf, self.log)
|
||||
try:
|
||||
self._await_serving(cfg, wave)
|
||||
preflight.check_wave(cfg, wave, self.log)
|
||||
adj_mod.run_adjudication(cfg, remote_deck, self.log)
|
||||
finally:
|
||||
serving.tear_down_wave(cfg, wave, self.log)
|
||||
local_adj = os.path.join(local_deck, "adjudicator-out")
|
||||
sc.pull_dir(head, f"{remote_deck}/adjudicator-out", local_adj)
|
||||
apath = os.path.join(local_adj, "ADJUDICATION.md")
|
||||
if os.path.exists(apath):
|
||||
adjudication_md = open(apath, errors="replace").read().strip() or None
|
||||
if not adjudication_md:
|
||||
self.log("[runner] WARNING: no adjudication produced (continuing without)")
|
||||
except Exception as e:
|
||||
self.log(f"[runner] WARNING: adjudication failed (non-fatal): {e}")
|
||||
|
||||
# --- deterministic scoring + ledger ------------------------------------
|
||||
self.phase = "scoring"; self._persist()
|
||||
company = led.ensure_company(slug_c)
|
||||
pinned = company.get("pinned_targets") or []
|
||||
aliases_map = company.get("kpi_aliases") or {}
|
||||
prior = led.prior_targets(slug_c, period)
|
||||
report_deck_dir = os.path.join(REPORTS_DIR, job_id, slug_c, deck_id)
|
||||
meta = {
|
||||
"company": slug_c,
|
||||
"period": period,
|
||||
"deck_id": deck_id,
|
||||
"job_id": job_id,
|
||||
"graded_at": datetime.now(timezone.utc).isoformat(),
|
||||
"panel": panel_meta,
|
||||
"artifacts": {
|
||||
"report_dir": report_deck_dir,
|
||||
"extraction": os.path.join(report_deck_dir, "extraction.json"),
|
||||
"grades": [os.path.join(report_deck_dir, f"{p['rid']}.json")
|
||||
for p in panel_meta if p["valid"]],
|
||||
"adjudication": (os.path.join(report_deck_dir, "ADJUDICATION.md")
|
||||
if adjudication_md else None),
|
||||
},
|
||||
}
|
||||
record = scoring.score_deck(ext_obj, grades, pinned, prior, aliases_map,
|
||||
cfg["weights"], meta)
|
||||
rec_path = led.record_deck(slug_c, record, ext_obj.get("forward_targets") or [])
|
||||
self.log(f"[runner] ledger updated: {rec_path}")
|
||||
|
||||
# --- reports ------------------------------------------------------------
|
||||
self.phase = "collecting"; self._persist()
|
||||
deck_md = scorecard.render_deck_report(record, ext_obj, adjudication_md)
|
||||
if rec_path and str(rec_path).endswith(".json"):
|
||||
ledger_md = os.path.splitext(str(rec_path))[0] + ".md"
|
||||
else:
|
||||
ledger_md = os.path.join(LEDGER_DIR, slug_c, "decks", f"{deck_id}.md")
|
||||
os.makedirs(os.path.dirname(ledger_md), exist_ok=True)
|
||||
with open(ledger_md, "w") as f:
|
||||
f.write(deck_md)
|
||||
|
||||
os.makedirs(report_deck_dir, exist_ok=True)
|
||||
with open(os.path.join(report_deck_dir, "DECK_REPORT.md"), "w") as f:
|
||||
f.write(deck_md)
|
||||
for fn in sorted(os.listdir(local_out)): # extraction + raw grader jsons (+ .invalid)
|
||||
src = os.path.join(local_out, fn)
|
||||
if os.path.isfile(src):
|
||||
shutil.copyfile(src, os.path.join(report_deck_dir, fn))
|
||||
if adjudication_md:
|
||||
with open(os.path.join(report_deck_dir, "ADJUDICATION.md"), "w") as f:
|
||||
f.write(adjudication_md + "\n")
|
||||
|
||||
# Refresh the company scorecard + the /data/reports latest copy.
|
||||
sc_md = scorecard.render_scorecard(led.get_company(slug_c), led.deck_records(slug_c))
|
||||
sc_path = os.path.join(LEDGER_DIR, slug_c, "SCORECARD.md")
|
||||
os.makedirs(os.path.dirname(sc_path), exist_ok=True)
|
||||
with open(sc_path, "w") as f:
|
||||
f.write(sc_md)
|
||||
with open(os.path.join(REPORTS_DIR, "latest-scorecard.md"), "w") as f:
|
||||
f.write(sc_md)
|
||||
self.log(f"[runner] reports saved to {report_deck_dir}")
|
||||
|
||||
return {"period": period, "deck_id": deck_id, "composite": _composite(record),
|
||||
"report_dir": report_deck_dir}
|
||||
|
||||
def _resolve_period(self, unit: dict, ext_obj: dict, ext_path: str,
|
||||
token: str) -> tuple[str, str]:
|
||||
"""(period, deck_id) for this deck. Filename-derived period wins; else the
|
||||
extractor's deck.period if canonical-ish; else the file's mtime month
|
||||
(flagged as period_inferred in the extraction's red-flag candidates)."""
|
||||
if unit.get("period"):
|
||||
return unit["period"], _token(unit["period"])
|
||||
p = ((ext_obj.get("deck") or {}).get("period") or "").strip()
|
||||
if p and _PERIOD_RE.match(p):
|
||||
self.log(f"[runner] period '{p}' taken from the deck text")
|
||||
return p, _token(p)
|
||||
try:
|
||||
mtime = os.path.getmtime(unit["files"][0])
|
||||
except OSError:
|
||||
mtime = time.time()
|
||||
period = time.strftime("%Y-%m", time.localtime(mtime))
|
||||
ext_obj.setdefault("red_flag_candidates", []).append({
|
||||
"code": "period_inferred",
|
||||
"description": ("Reporting period was not stated in the filename or the deck "
|
||||
f"text; inferred from the file's modification time as {period}."),
|
||||
"severity": 2,
|
||||
"evidence": "",
|
||||
})
|
||||
try:
|
||||
with open(ext_path, "w") as f:
|
||||
json.dump(ext_obj, f, indent=2)
|
||||
except Exception:
|
||||
pass
|
||||
self.log(f"[runner] WARNING: period inferred from file mtime: {period}")
|
||||
return period, (_token(period) or token)
|
||||
|
||||
# ------------------------------------------------------------- helpers
|
||||
def _await_serving(self, cfg, wave, timeout=900):
|
||||
self.log("[runner] waiting for wave serving to come online…")
|
||||
@@ -329,50 +596,55 @@ class JobRunner:
|
||||
p["status"] = "no-report"
|
||||
self._persist()
|
||||
|
||||
def _collect(self, cfg, job_id, remote_job, local_job, valid, manifest, synth_ok):
|
||||
out_local = os.path.join(REPORTS_DIR, job_id)
|
||||
os.makedirs(out_local, exist_ok=True)
|
||||
sc.pull_dir(sc.head(cfg), f"{remote_job}/out", os.path.join(out_local, "reviewers"))
|
||||
if synth_ok:
|
||||
sc.pull_dir(sc.head(cfg), f"{remote_job}/synth-out", os.path.join(out_local, "synthesis"))
|
||||
|
||||
# Assemble a single latest.md: the consolidated report if present, else a
|
||||
# concatenation of the individual reports.
|
||||
parts = [f"# Boardroom Map review — {job_id}\n",
|
||||
"Documents reviewed: " + ", ".join(m["source"] for m in manifest if m["ok"]) + "\n",
|
||||
"Panel: " + ", ".join(f"{r['name']} ({r['model']})" for r in valid) + "\n"]
|
||||
consolidated = os.path.join(out_local, "synthesis", "CONSOLIDATED_REPORT.md")
|
||||
if synth_ok and os.path.exists(consolidated):
|
||||
parts.append("\n---\n\n## Consolidated report (lead reviewer)\n\n")
|
||||
parts.append(open(consolidated, errors="replace").read())
|
||||
parts.append("\n\n---\n")
|
||||
parts.append("\n## Individual reviewer reports\n")
|
||||
rev_local = os.path.join(out_local, "reviewers")
|
||||
if os.path.isdir(rev_local):
|
||||
for fn in sorted(os.listdir(rev_local)):
|
||||
if fn.endswith(".md"):
|
||||
parts.append(f"\n### {fn[:-3]}\n\n")
|
||||
parts.append(open(os.path.join(rev_local, fn), errors="replace").read())
|
||||
parts.append("\n")
|
||||
assembled = "".join(parts)
|
||||
with open(os.path.join(out_local, "report.md"), "w") as f:
|
||||
def _write_job_report(self, job_id: str):
|
||||
"""latest.md — one job summary across every deck graded (or failed)."""
|
||||
lines = [f"# Boardroom Map grading job — {job_id}\n"]
|
||||
done = [d for d in self.decks if d.get("status") == "done"]
|
||||
failed = [d for d in self.decks if d.get("status") == "failed"]
|
||||
lines.append(f"Decks graded: {len(done)}/{len(self.decks)}\n")
|
||||
lines.append("\n## Results\n")
|
||||
for d in self.decks:
|
||||
if d.get("status") == "done":
|
||||
comp = d.get("composite")
|
||||
comp_s = f"{comp:.1f}" if isinstance(comp, (int, float)) else "?"
|
||||
lines.append(f"- **{d['company']}** — {d.get('period') or '?'}: "
|
||||
f"composite **{comp_s}** / 100\n")
|
||||
else:
|
||||
lines.append(f"- **{d['company']}** — {d.get('period') or '?'}: "
|
||||
f"FAILED — {d.get('error', 'unknown error')}\n")
|
||||
if done:
|
||||
lines.append("\nPer-deck reports (DECK_REPORT.md, extraction, raw grades, "
|
||||
f"adjudication): `/data/reports/{job_id}/<company>/<deck>/`.\n")
|
||||
lines.append("Company scorecards: `/data/ledger/<company>/SCORECARD.md` "
|
||||
"(latest copy at `/data/reports/latest-scorecard.md`).\n")
|
||||
if failed:
|
||||
lines.append("\nFailed decks were still moved to "
|
||||
f"`/data/processed/{job_id}/` — re-drop them to regrade.\n")
|
||||
assembled = "".join(lines)
|
||||
out_dir = os.path.join(REPORTS_DIR, job_id)
|
||||
os.makedirs(out_dir, exist_ok=True)
|
||||
with open(os.path.join(out_dir, "report.md"), "w") as f:
|
||||
f.write(assembled)
|
||||
with open(os.path.join(REPORTS_DIR, "latest.md"), "w") as f:
|
||||
f.write(assembled)
|
||||
self.last_report_path = os.path.join(out_local, "report.md")
|
||||
self.log(f"[runner] reports saved to {out_local}")
|
||||
self.last_report_path = os.path.join(out_dir, "report.md")
|
||||
|
||||
def _drain_inbox(self, job_id):
|
||||
dest = os.path.join(PROCESSED, job_id)
|
||||
os.makedirs(dest, exist_ok=True)
|
||||
for fn in os.listdir(INBOX):
|
||||
src = os.path.join(INBOX, fn)
|
||||
if os.path.isfile(src):
|
||||
try:
|
||||
shutil.move(src, os.path.join(dest, fn))
|
||||
except Exception:
|
||||
pass
|
||||
self.log(f"[runner] inbox cleared (originals moved to processed/{job_id})")
|
||||
def _drain_inbox(self, job_id: str, units: list[dict]):
|
||||
"""Move each unit's ORIGINAL files to /data/processed/<job>/<slug>/. The
|
||||
per-company inbox folders are kept — the user reuses them next quarter."""
|
||||
moved = 0
|
||||
for unit in units:
|
||||
dest = os.path.join(PROCESSED, job_id, unit["company_slug"])
|
||||
os.makedirs(dest, exist_ok=True)
|
||||
for src in unit["files"]:
|
||||
if os.path.isfile(src):
|
||||
try:
|
||||
shutil.move(src, os.path.join(dest, os.path.basename(src)))
|
||||
moved += 1
|
||||
except Exception:
|
||||
pass
|
||||
self.log(f"[runner] inbox cleared ({moved} file(s) moved to processed/{job_id}; "
|
||||
"company folders kept)")
|
||||
|
||||
|
||||
# Module-level singleton used by app.py
|
||||
|
||||
Reference in New Issue
Block a user