"""Boardroom Map orchestrator web app. Serves the portfolio dashboard and a small JSON API, and starts the background job runner (see jobs.py) that grades the dropped board decks. Most configuration happens through the StartOS *actions* (Configure Sparks / Models / Graders / Grading); this UI is for dropping decks per company, triggering a grading run, watching it, and reading the per-company scorecard ledger. """ from __future__ import annotations import os import threading from fastapi import FastAPI, File, Form, HTTPException, UploadFile from fastapi.responses import HTMLResponse, JSONResponse, PlainTextResponse from fastapi.templating import Jinja2Templates from starlette.requests import Request import bm_config import decks import graders as grader_mod import jobs import ledger as ledger_mod import serving DATA_DIR = os.environ.get("BM_DATA_DIR", "/data") LEDGER_DIR = os.path.join(DATA_DIR, "ledger") runner = jobs.runner INBOX = getattr(jobs, "INBOX", os.path.join(DATA_DIR, "inbox")) REPORTS_DIR = getattr(jobs, "REPORTS_DIR", os.path.join(DATA_DIR, "reports")) PROCESSED = getattr(jobs, "PROCESSED", os.path.join(DATA_DIR, "processed")) app = FastAPI(title="Boardroom Map Orchestrator") templates = Jinja2Templates(directory=os.path.join(os.path.dirname(__file__), "templates")) @app.on_event("startup") def _startup(): runner.start() # ----------------------------------------------------------------------- UI @app.get("/", response_class=HTMLResponse) def index(request: Request): return templates.TemplateResponse("index.html", {"request": request}) @app.get("/healthz") def healthz(): return {"ok": True} # ----------------------------------------------------------------------- status @app.get("/api/status") def status(): cfg = bm_config.load() catalog = {m["alias"] for m in (cfg.get("models") or [])} return { "configured": { "sparks": bool(cfg.get("primarySparkHost")), "models": len(cfg.get("models") or []), "graders": len(cfg.get("graders") or []), }, "networkMode": cfg.get("networkMode"), "adjudicator": bool(cfg.get("adjudicatorEnabled")), "wipeRemoteDocs": bool(cfg.get("wipeRemoteDocs")), "autoRunOnDrop": bool(cfg.get("autoRunOnDrop")), "models": [{"alias": m["alias"], "hfModel": m["hfModel"], "spark": m.get("spark", "primary")} for m in (cfg.get("models") or [])], "panel": [{"name": g.get("name"), "model": g.get("model"), "persona": bool((g.get("persona") or "").strip()), "known": (g.get("model") in catalog)} for g in (cfg.get("graders") or [])], "inbox": _inbox_grouped(), "runtime": runner.snapshot(), } @app.get("/api/events") def events(): return {"events": runner.events()} # ----------------------------------------------------------------------- inbox def _inbox_grouped() -> dict: """Company-grouped inbox view: {companies: {slug: [file dicts]}, skipped: [...]}.""" try: d = decks.discover(INBOX) except Exception: return {"companies": {}, "skipped": []} companies: dict[str, list] = {} for u in d.get("units") or []: lst = companies.setdefault(u["company_slug"], []) for f in u.get("files") or []: try: size = os.path.getsize(f) except OSError: size = 0 lst.append({"name": os.path.basename(f), "bytes": size, "period": u.get("period"), "supported": True}) for fn in u.get("ignored") or []: lst.append({"name": fn, "bytes": 0, "period": None, "supported": False}) # discover() drops units with no supported files, so sweep the company dirs # for anything it didn't list (unsupported strays) and flag them. try: for entry in sorted(os.listdir(INBOX)): cdir = os.path.join(INBOX, entry) if entry.startswith(".") or not os.path.isdir(cdir): continue slug = decks.slugify(entry) seen = {f["name"] for f in companies.get(slug, [])} for fn in sorted(os.listdir(cdir)): if fn.startswith(".") or fn in seen or not os.path.isfile(os.path.join(cdir, fn)): continue try: size = os.path.getsize(os.path.join(cdir, fn)) except OSError: size = 0 companies.setdefault(slug, []).append( {"name": fn, "bytes": size, "period": decks.parse_period_from_name(fn), "supported": os.path.splitext(fn)[1].lower() in decks.SUPPORTED_EXTS}) except OSError: pass return {"companies": companies, "skipped": d.get("skipped") or []} @app.get("/api/inbox") def inbox(): return _inbox_grouped() # ----------------------------------------------------------------------- documents @app.post("/api/upload") async def upload(request: Request, files: list[UploadFile] = File(...), company: str | None = Form(None)): name = (company or request.query_params.get("company") or "").strip() if not name: raise HTTPException(400, "company is required — root-level files are not graded") slug = decks.slugify(name) dest_dir = os.path.join(INBOX, slug) os.makedirs(dest_dir, exist_ok=True) saved = [] for f in files: fn = os.path.basename(f.filename or "deck") dest = os.path.join(dest_dir, fn) with open(dest, "wb") as out: while chunk := await f.read(1 << 20): out.write(chunk) saved.append(fn) return {"ok": True, "company": slug, "saved": saved} @app.post("/api/inbox/clear") def inbox_clear(): import shutil if os.path.isdir(INBOX): for fn in os.listdir(INBOX): p = os.path.join(INBOX, fn) try: if os.path.isfile(p): os.remove(p) elif os.path.isdir(p): shutil.rmtree(p) except OSError: pass return {"ok": True} # ----------------------------------------------------------------------- run @app.post("/api/run") def run_now(): cfg = bm_config.load() if not cfg.get("primarySparkHost"): raise HTTPException(400, "No Spark configured (Configure Sparks).") if not (cfg.get("models") and cfg.get("graders")): raise HTTPException(400, "Configure at least one model and one grader first.") try: units = decks.discover(INBOX).get("units") or [] except Exception: units = [] if not units: raise HTTPException(400, "Inbox is empty — upload decks into a company folder first.") runner.request_run() return {"ok": True, "message": "Grading requested — watch the activity log."} @app.get("/api/serving") def serving_status(): cfg = bm_config.load() if not cfg.get("primarySparkHost"): raise HTTPException(400, "No Spark configured.") return {"serving": serving.health(cfg)} @app.post("/api/grader/build-image") def build_grader_image(): cfg = bm_config.load() if not cfg.get("primarySparkHost"): raise HTTPException(400, "No Spark configured.") threading.Thread(target=lambda: _safe_build(cfg), daemon=True).start() return {"ok": True, "message": "Building grader image on the head Spark — watch the activity log."} def _safe_build(cfg: dict): try: fn = getattr(grader_mod, "ensure_grader_image", None) or grader_mod.ensure_reviewer_image fn(cfg, runner.log) except Exception as e: runner.log(f"[graders] image build failed: {e}") @app.post("/api/stop") def stop(): cfg = bm_config.load() if not cfg.get("primarySparkHost"): raise HTTPException(400, "No Spark configured.") serving.tear_down_all(cfg, runner.log) runner.phase = "idle" runner._persist() return {"ok": True} # ----------------------------------------------------------------------- companies / ledger def _config_companies(cfg: dict) -> list[dict]: """Companies registered in the StartOS config (may have zero graded decks).""" out = [] for c in (cfg.get("companies") or []): if isinstance(c, dict): name = (c.get("name") or c.get("company") or c.get("slug") or "").strip() else: name = str(c).strip() if name: out.append({"slug": decks.slugify(name), "name": name}) return out def _ledger_companies() -> list[dict]: try: led = ledger_mod.Ledger(LEDGER_DIR) return led.all_companies() or [] except Exception: return [] @app.get("/api/companies") def companies(): cfg = bm_config.load() rows, seen = [], set() for c in _ledger_companies(): try: hist = [{"period": h.get("period"), "composite": h.get("composite")} for h in (c.get("history") or [])] latest = None if hist: delta = None cur, prev = hist[-1]["composite"], (hist[-2]["composite"] if len(hist) >= 2 else None) if isinstance(cur, (int, float)) and isinstance(prev, (int, float)): delta = round(cur - prev, 2) latest = {"period": hist[-1]["period"], "composite": cur, "delta": delta} slug = c.get("slug") or decks.slugify(c.get("name") or "") seen.add(slug) rows.append({"slug": slug, "name": c.get("name") or slug, "auto_created": bool(c.get("auto_created")), "deck_count": len(hist), "latest": latest, "history": hist}) except Exception: continue for c in _config_companies(cfg): if c["slug"] not in seen: seen.add(c["slug"]) rows.append({"slug": c["slug"], "name": c["name"], "auto_created": False, "deck_count": 0, "latest": None, "history": []}) rows.sort(key=lambda r: (r["name"] or "").lower()) return {"companies": rows} def _kpi_hit_rate(records: list[dict]) -> dict: """Per canonical KPI across records (oldest→newest): attempts, hits (credit>=0.999), streak of hits ending at the latest attempt, last_credit and profitability tag.""" stats: dict[str, dict] = {} for rec in records: for k in (rec.get("kpi_results") or []): cn = k.get("canonical_name") or k.get("name") if not cn: continue s = stats.setdefault(cn, {"name": k.get("name") or cn, "attempts": 0, "hits": 0, "streak": 0, "last_credit": None, "profitability": False}) if k.get("name"): s["name"] = k["name"] s["profitability"] = bool(k.get("profitability", s["profitability"])) credit = k.get("credit") if not isinstance(credit, (int, float)): continue # no target matched — not an attempt s["attempts"] += 1 s["last_credit"] = credit if credit >= 0.999: s["hits"] += 1 s["streak"] += 1 else: s["streak"] = 0 return stats def _categories_latest(records: list[dict]) -> dict: if not records: return {} def cats(rec): return ((rec.get("qual") or {}).get("categories")) or {} latest = cats(records[-1]) prev = cats(records[-2]) if len(records) >= 2 else {} out = {} for cid in "ABCDEFGH": cur = latest.get(cid) if cur is None and prev.get(cid) is None: continue out[cid] = {"latest_adjusted": (cur or {}).get("adjusted"), "previous_adjusted": (prev.get(cid) or {}).get("adjusted")} return out @app.get("/api/companies/{slug}") def company_detail(slug: str): slug = os.path.basename(slug) company, records = None, [] try: led = ledger_mod.Ledger(LEDGER_DIR) company = led.get_company(slug) if company is not None: records = led.deck_records(slug) or [] except Exception: company, records = None, [] if company is None: cfg = bm_config.load() match = next((c for c in _config_companies(cfg) if c["slug"] == slug), None) if not match: raise HTTPException(404, "no such company") company = {"slug": slug, "name": match["name"], "auto_created": False, "kpi_aliases": {}, "pinned_targets": [], "extracted_targets": {}, "history": []} open_flags = (((records[-1].get("penalties") or {}).get("flags")) or []) if records else [] return { "company": company, "records": records, "kpi_hit_rate": _kpi_hit_rate(records), "categories_latest": _categories_latest(records), "open_flags": open_flags, } @app.delete("/api/companies/{slug}") def delete_company(slug: str, restore_decks: bool = False): """Wipe a company's graded data so it can be re-evaluated from scratch. Removes the ledger (deck records, scorecard, report markdown) plus the per-job report and processed-deck copies. With restore_decks, the original deck files graded earlier are moved from /data/processed back into the company's inbox folder first, ready for a fresh Grade Decks run. A config-registered company reappears as an empty row (name, aliases and pinned targets kept); an auto-created one disappears entirely. """ import shutil slug = os.path.basename(slug) if runner.phase in jobs.RUNNING_PHASES: raise HTTPException(409, "A grading job is running — wait for it to finish.") if not os.path.isdir(os.path.join(LEDGER_DIR, slug)) \ and not os.path.isdir(os.path.join(INBOX, slug)): raise HTTPException(404, "no such company") restored = 0 if restore_decks and os.path.isdir(PROCESSED): dest_dir = os.path.join(INBOX, slug) for job_dir in sorted(os.listdir(PROCESSED)): src_dir = os.path.join(PROCESSED, job_dir, slug) if not os.path.isdir(src_dir): continue os.makedirs(dest_dir, exist_ok=True) for fn in sorted(os.listdir(src_dir)): src = os.path.join(src_dir, fn) dest = os.path.join(dest_dir, fn) if os.path.isfile(src) and not os.path.exists(dest): try: shutil.move(src, dest) restored += 1 except OSError: pass removed = {"ledger": False, "report_dirs": 0, "processed_dirs": 0} led_dir = os.path.join(LEDGER_DIR, slug) if os.path.isdir(led_dir): shutil.rmtree(led_dir, ignore_errors=True) removed["ledger"] = True for base, key in ((REPORTS_DIR, "report_dirs"), (PROCESSED, "processed_dirs")): if not os.path.isdir(base): continue for job_dir in os.listdir(base): d = os.path.join(base, job_dir, slug) if os.path.isdir(d): shutil.rmtree(d, ignore_errors=True) removed[key] += 1 runner.log(f"[api] company '{slug}' deleted " f"(ledger={removed['ledger']}, report dirs={removed['report_dirs']}, " f"processed dirs={removed['processed_dirs']}, " f"decks restored to inbox={restored})") return {"ok": True, "removed": removed, "restored_decks": restored} @app.get("/api/companies/{slug}/scorecard", response_class=PlainTextResponse) def company_scorecard(slug: str): slug = os.path.basename(slug) path = os.path.join(LEDGER_DIR, slug, "SCORECARD.md") if not os.path.exists(path): raise HTTPException(404, "no scorecard yet for this company") return open(path, errors="replace").read() @app.get("/api/companies/{slug}/decks/{deck_id}") def deck_record(slug: str, deck_id: str): slug, deck_id = os.path.basename(slug), os.path.basename(deck_id) rec = None try: led = ledger_mod.Ledger(LEDGER_DIR) rec = led.deck_record(slug, deck_id) except Exception: rec = None if rec is None: raise HTTPException(404, "no such deck record") return JSONResponse(rec) @app.get("/api/companies/{slug}/decks/{deck_id}/report", response_class=PlainTextResponse) def deck_report(slug: str, deck_id: str): slug, deck_id = os.path.basename(slug), os.path.basename(deck_id) if deck_id.endswith(".md"): deck_id = deck_id[:-3] path = os.path.join(LEDGER_DIR, slug, "decks", f"{deck_id}.md") if not os.path.exists(path): raise HTTPException(404, "no report for this deck") return open(path, errors="replace").read() # ----------------------------------------------------------------------- reports (legacy job reports) @app.get("/api/reports") def list_reports(): if not os.path.isdir(REPORTS_DIR): return {"reports": []} jobs_ = sorted((d for d in os.listdir(REPORTS_DIR) if os.path.isdir(os.path.join(REPORTS_DIR, d))), reverse=True) return {"reports": jobs_} @app.get("/api/report", response_class=PlainTextResponse) def latest_report(): path = os.path.join(REPORTS_DIR, "latest.md") if not os.path.exists(path): return "(no report yet — drop decks in a company folder and run a grading job)" return open(path, errors="replace").read().strip() or "(empty report)" @app.get("/api/reports/{job}", response_class=PlainTextResponse) def get_report(job: str): job = os.path.basename(job) path = os.path.join(REPORTS_DIR, job, "report.md") if not os.path.exists(path): raise HTTPException(404, "no such report") return open(path, errors="replace").read()