"""Project what a feedback fold would do to ``proficiency``, without writing. Folding is **irreversible**. ``add_outcome`` accumulates into a running mean and ``recompute_category`` re-derives every row in the category from the new peer rate; neither keeps the pre-fold value anywhere, and the ``verifications.applied_at`` stamp means a second run cannot redo the work either. So the operator gets one shot, and deserves to see the numbers first. The projection does not re-implement the arithmetic. It **copies the database into memory** and runs the real ``add_outcome`` against the copy, then diffs the two ``proficiency`` snapshots. Re-deriving the empirical-Bayes conversion by hand here would be a second implementation to keep in step with ``proficiency_outcome.py``, and the first time it drifted the preview would be confidently wrong about an irreversible action --- the worst available failure. The source connection is opened ``mode=ro``, so a preview cannot write to the database it is previewing even if something below it tried to. Two kinds of movement come out, and the distinction is the non-obvious part: - **direct** --- the row has new client outcomes of its own. - **ripple** --- the row has none, and moves anyway, because ``recompute_category`` recomputes the whole category against a peer rate that the new evidence changed. A model nobody reported on can still shift. """ from __future__ import annotations import csv import sqlite3 import statistics import urllib.request from dataclasses import dataclass from pathlib import Path from typing import Optional # Below this, a move is rounding rather than news. Reported as a count so the # suppression is visible, never silent. MOVE_EPSILON = 0.0005 Key = tuple[str, str, str] def read_only_uri(path: str | Path) -> str: """``file:`` URI for *path*, opened read-only, with characters escaped.""" return ( "file:" + urllib.request.pathname2url(str(Path(path).resolve())) + "?mode=ro" ) def copy_database(path: str | Path) -> sqlite3.Connection: """An in-memory copy of the database at *path*. The source is opened read-only and closed again immediately. Everything the projection does afterwards happens to the copy, so pointing this at the live database is as safe as pointing it at a snapshot --- and pointing it at a snapshot still works, which is what ``--db`` is for. """ src = sqlite3.connect(read_only_uri(path), uri=True) try: dest = sqlite3.connect(":memory:") src.backup(dest) finally: src.close() dest.row_factory = sqlite3.Row return dest def snapshot(conn: sqlite3.Connection) -> dict[Key, dict]: """Every ``proficiency`` row, keyed by (model_id, provider, category).""" conn.row_factory = sqlite3.Row return { (r["model_id"], r["provider"], r["category"]): { "blended_score": r["blended_score"], "source": r["source"], "outcome_score": r["outcome_score"], "outcome_samples": r["outcome_samples"] or 0, } for r in conn.execute( """ SELECT model_id, provider, category, blended_score, source, outcome_score, outcome_samples FROM proficiency """ ) } @dataclass(frozen=True) class Change: """One ``proficiency`` row's projected before/after.""" model_id: str provider: str category: str before_score: Optional[float] after_score: Optional[float] before_source: Optional[str] after_source: Optional[str] before_samples: int after_samples: int successes: int failures: int @property def is_new(self) -> bool: """No row existed before, so there is no prior belief to revise.""" return self.before_score is None @property def kind(self) -> str: return "direct" if (self.successes or self.failures) else "ripple" @property def delta(self) -> Optional[float]: if self.before_score is None or self.after_score is None: return None return self.after_score - self.before_score @property def magnitude(self) -> float: """Sort key: how far this row moves, with new rows sorted last.""" d = self.delta return abs(d) if d is not None else 0.0 def fold_onto(conn: sqlite3.Connection, cfg, grouped: dict) -> None: """Apply *grouped* to *conn* through the real store, marking rows applied. Deliberately calls ``proficiency_store.add_outcome`` rather than ``feedback.apply_failures``: the preview must exercise the storage arithmetic, not the CLI's printing, and reaching through ``feedback`` would also make a dry run indistinguishable from a real one to anything watching that module. """ from proficiency_store import add_outcome for (model_id, provider, category), pairs in sorted(grouped.items()): add_outcome( conn, cfg, model_id, provider, category, [sc for _, sc in pairs] ) conn.executemany( "UPDATE verifications SET applied_at = datetime('now') WHERE id = ?", [(i,) for i, _ in pairs], ) conn.commit() def diff( before: dict[Key, dict], after: dict[Key, dict], grouped: dict ) -> tuple[list[Change], int]: """Rows worth showing, plus the count suppressed as rounding. A row is worth showing when it carries new evidence (even if its score did not budge --- "you spent 42 samples and nothing moved" is a real answer) or when it moved by at least ``MOVE_EPSILON``. """ counts = { key: ( sum(1 for _, s in pairs if s == 1.0), sum(1 for _, s in pairs if s != 1.0), ) for key, pairs in grouped.items() } changes: list[Change] = [] suppressed = 0 for key in sorted(set(before) | set(after)): b = before.get(key) a = after.get(key) successes, failures = counts.get(key, (0, 0)) before_score = b["blended_score"] if b else None after_score = a["blended_score"] if a else None moved = ( before_score is None or after_score is None or abs(after_score - before_score) >= MOVE_EPSILON ) if not moved and not (successes or failures): suppressed += 1 continue changes.append( Change( model_id=key[0], provider=key[1], category=key[2], before_score=before_score, after_score=after_score, before_source=b["source"] if b else None, after_source=a["source"] if a else None, before_samples=b["outcome_samples"] if b else 0, after_samples=a["outcome_samples"] if a else 0, successes=successes, failures=failures, ) ) changes.sort(key=lambda c: (-c.magnitude, c.model_id, c.category)) return changes, suppressed def project(path: str | Path, cfg) -> tuple[list[Change], int, dict]: """Full projection for the database at *path*. Writes nothing to it. Returns ``(changes, suppressed, grouped)``; ``grouped`` is the same shape ``feedback.summarize`` produces, so callers can report sample totals without querying twice. """ from feedback import summarize, unapplied_failures conn = copy_database(path) try: grouped = summarize(unapplied_failures(conn)) before = snapshot(conn) fold_onto(conn, cfg, grouped) after = snapshot(conn) finally: conn.close() changes, suppressed = diff(before, after, grouped) return changes, suppressed, grouped # --- rendering ------------------------------------------------------------- def _score(value: Optional[float]) -> str: return " -- " if value is None else f"{value:.3f} " def _bar(delta: Optional[float], scale: float, width: int = 8) -> str: """A short ASCII bar, scaled to the largest move in this run.""" if delta is None or scale <= 0: return "" n = int(round(abs(delta) / scale * width)) return ("+" if delta > 0 else "-") * max(1, n) if n else "" def format_projection( changes: list[Change], suppressed: int, grouped: dict, sample_count: int ) -> str: """The human summary: a verdict first, then detail, then a digest. The direct rows are listed individually because there are few of them and each one is a model the operator recognises. The ripple rows are NOT --- on the live database a 207-sample fold moved 275 of them, which is a wall, and within a category they all move nearly identically because they are all being pulled toward the same new peer rate. So ripple is digested per category; ``--csv`` still carries every row for anyone who wants them. """ lines: list[str] = [] lines.append( f"would apply {sample_count} sample(s) across {len(grouped)} " f"(model, provider, category) pair(s)" ) lines.append("") if not changes: lines.append("projected proficiency changes: none") return "\n".join(lines) scale = max((c.magnitude for c in changes), default=0.0) direct = [c for c in changes if c.kind == "direct"] ripple = [c for c in changes if c.kind == "ripple"] drops = [c for c in changes if (c.delta or 0) < 0] rises = [c for c in changes if (c.delta or 0) > 0] new_rows = [c for c in changes if c.is_new] lines.append("projected proficiency changes -- DRY RUN, nothing was written") lines.append("") lines.append( f" {len(direct)} row(s) move on their own new evidence, " f"{len(ripple)} move by ripple, {len(new_rows)} are new" ) if drops: worst = min(drops, key=lambda c: c.delta) lines.append( f" largest drop {worst.delta:+.3f} " f"{worst.before_score:.3f} -> {worst.after_score:.3f} " f"{worst.model_id} @{worst.provider} / {worst.category}" ) if rises: best = max(rises, key=lambda c: c.delta) lines.append( f" largest rise {best.delta:+.3f} " f"{best.before_score:.3f} -> {best.after_score:.3f} " f"{best.model_id} @{best.provider} / {best.category}" ) if suppressed: lines.append( f" {suppressed} further row(s) moved by less than {MOVE_EPSILON}" ) lines.append("") if direct: lines.append(" direct -- rows carrying new client outcomes") lines.append( f" {'model':30s} {'provider':10s} {'category':19s} " f"{'outcomes':9s} {'blended_score':16s} {'delta':>7s}" ) for c in direct: outcomes = f"+{c.successes} -{c.failures}" delta = c.delta delta_text = " new " if delta is None else f"{delta:+7.3f}" lines.append( f" {c.model_id[:30]:30s} {c.provider[:10]:10s} " f"{c.category[:19]:19s} {outcomes:9s} " f"{_score(c.before_score)}-> {_score(c.after_score)} " f"{delta_text} {_bar(delta, scale)}" ) lines.append("") if ripple: lines.append( " ripple -- rows with no new outcomes of their own. They move because\n" " recompute_category re-derives every row in a category against\n" " the peer rate the new evidence just changed. Per-row: --csv." ) lines.append( f" {'category':19s} {'rows':>5s} {'delta range':21s} {'median':>7s}" ) for category, rows in _by_category(ripple): deltas = sorted(c.delta for c in rows if c.delta is not None) if not deltas: continue mid = statistics.median(deltas) lines.append( f" {category[:19]:19s} {len(rows):>5d} " f"{deltas[0]:+.3f} .. {deltas[-1]:+.3f} " f"{mid:+7.3f} {_bar(mid, scale)}" ) lines.append("") return "\n".join(lines).rstrip() def _by_category(rows: list[Change]) -> list[tuple[str, list[Change]]]: """Group *rows* by category, widest-moving category first.""" buckets: dict[str, list[Change]] = {} for c in rows: buckets.setdefault(c.category, []).append(c) return sorted( buckets.items(), key=lambda kv: -max(c.magnitude for c in kv[1]), ) CSV_COLUMNS = [ "model_id", "provider", "category", "kind", "is_new", "successes", "failures", "before_score", "after_score", "delta", "before_samples", "after_samples", "before_source", "after_source", ] def write_csv(changes: list[Change], out) -> None: """Emit the projection as CSV, one row per projected change.""" writer = csv.DictWriter(out, fieldnames=CSV_COLUMNS) writer.writeheader() for c in changes: writer.writerow( { "model_id": c.model_id, "provider": c.provider, "category": c.category, "kind": c.kind, "is_new": int(c.is_new), "successes": c.successes, "failures": c.failures, "before_score": "" if c.before_score is None else f"{c.before_score:.6f}", "after_score": "" if c.after_score is None else f"{c.after_score:.6f}", "delta": "" if c.delta is None else f"{c.delta:.6f}", "before_samples": c.before_samples, "after_samples": c.after_samples, "before_source": c.before_source or "", "after_source": c.after_source or "", } )