"""Analysis pipeline for the Recensorium context-size field experiment. Reads exp/*/obs_*.json, computes primary/secondary outcomes per prereg.md. Run only AFTER data collection completes. This file is the locked pipeline; any post-hoc change must be reported in the paper's deviations section. """ import glob import json import os import sys import numpy as np HERE = os.path.dirname(os.path.abspath(__file__)) DIMS = ["novelty", "rigour", "clarity", "significance"] def _normalize(d): """Map schema variants used by different reviewer agents onto one canon. Mechanical renaming only.""" if "context_reviews" not in d: for alt in ("reviews_shown_in_frozen_context", "reviews_shown", "shown_reviews", "reviews_in_context", "context_reviews", "context"): if alt in d and isinstance(d[alt], list): d["context_reviews"] = d[alt] break if "context_size_requested" not in d: for alt in ("dose", "context_size"): if alt in d: d["context_size_requested"] = d[alt] break else: d["context_size_requested"] = None return d def load_obs(): rows = [] files = sorted(glob.glob(os.path.join(HERE, "exp", "a*", "obs_*.json"))) # optional recovery layer (see recover_ctx.py): fills context entries whose # logged records lack dimension scores, joined on review_id from the PUBLIC # reviews endpoint; original files untouched. rec_path = os.path.join(HERE, "recovery.json") recovered = {} if os.path.exists(rec_path): try: recovered = json.load(open(rec_path)).get("recovered", {}) except Exception: recovered = {} for f in files: try: d = json.load(open(f)) except Exception as e: print(f"SKIP unreadable {f}: {e}") continue d = _normalize(d) if recovered: filled = 0 for c in d.get("context_reviews") or []: if not any(isinstance(c.get(k), dict) and c[k] for k in ("scores", "dimension_scores")): rid = c.get("review_id") if rid in recovered: c["dimension_scores"] = \ recovered[rid].get("scores", {}) filled += 1 if filled: d["_recovered_entries"] = filled d["_file"] = os.path.relpath(f, HERE) agent = os.path.basename(os.path.dirname(f)) d["_agent"] = agent rows.append(d) return rows def primary_table(rows): """Per-observation absolute standardized deviation vs context mean.""" out = [] for d in rows: ctx = d.get("context_reviews") or [] mys = d.get("my_scores") or {} dose = d.get("context_size_requested") if not ctx or not mys or not isinstance(dose, int): continue for dim in DIMS: xs = [c.get("scores", {}).get(dim) for c in ctx if c.get("scores")] xs = [x for x in xs if isinstance(x, (int, float))] mine = mys.get(dim) if mine is None or len(xs) < 2: continue # prereg exclusion: singleton contexts excluded mu = float(np.mean(xs)) sd = float(np.std(xs, ddof=1)) if sd == 0: # all-shown-agree case: deviation = distance from that consensus, # standardized by pooled dimension sd later. Mark with sd=None. dev_abs = abs(float(mine) - mu) out.append(dict(agent=d["_agent"], paper=d["paper_id"], dim=dim, dose=dose, n_ctx=len(xs), mu=mu, sd=None, y_abs=dev_abs / 2.0, y_signed=(float(mine) - mu) / 2.0, zero_disp=1)) else: dev_abs = abs(float(mine) - mu) out.append(dict(agent=d["_agent"], paper=d["paper_id"], dim=dim, dose=dose, n_ctx=len(xs), mu=mu, sd=sd, y_abs=dev_abs / sd, y_signed=(float(mine) - mu) / sd, zero_disp=0)) return out def _ctx_scores(c): """Get dimension-score dict from a context-review record (schema variants). Handles *_present wrappers that hide scores behind booleans.""" for key in ("scores", "dimension_scores"): s = c.get(key) if isinstance(s, dict): return s if s is False or s == "absent": continue # variant: dimension_scores_present boolean + separate fields if c.get("dimension_scores_present") is True: s = {k: c.get(k) for k in DIMS if isinstance(c.get(k), (int, float))} if s: return s if isinstance(c.get("scores_by_dim"), dict): return c["scores_by_dim"] return {} def conformity_table(rows): """Per-observation effective weight vs Bayes benchmark (prereg v2 primary). Y_ind(n) = sqrt(1+1/n); w_eff = 1 - Y_obs/Y_ind; Delta = w_eff - n/(n+1). n = ACTUAL number of scored reviews shown (the benchmark depends on real information); the randomized dose is the instrument.""" out = [] for d in rows: ctx = d.get("context_reviews") or [] mys = d.get("my_scores") or {} dose = d.get("context_size_requested") if not ctx or not mys or not isinstance(dose, int): continue for dim in DIMS: xs = [_ctx_scores(c).get(dim) for c in ctx] xs = [x for x in xs if isinstance(x, (int, float))] mine = mys.get(dim) if mine is None or len(xs) < 2: continue mu = float(np.mean(xs)) sd = float(np.std(xs, ddof=1)) y_obs_raw = abs(float(mine) - mu) if sd == 0: # degenerate context: standardize by half the score quantum y_obs = y_obs_raw / 0.5 disp = "zero" else: y_obs = y_obs_raw / sd disp = "pos" n_ctx = len(xs) y_ind = float(np.sqrt(1.0 + 1.0 / n_ctx)) w_eff = 1.0 - y_obs / y_ind delta = w_eff - n_ctx / (n_ctx + 1.0) out.append(dict(agent=d["_agent"], paper=d["paper_id"], dim=dim, dose=dose, n_ctx=n_ctx, mu=mu, sd=sd, disp=disp, y_obs=y_obs, y_ind=y_ind, w_eff=w_eff, delta=delta, recovered=bool(d.get("_recovered_entries")), y_signed=(float(mine) - mu) / (sd if sd else 0.5))) return out def _design(tab, outcome): """Pure-numpy design: intercept, dose20, dose8 + reviewer FE + dimension FE.""" agents = sorted(set(r["agent"] for r in tab)) dims = sorted(set(r["dim"] for r in tab)) a_idx = {a: i for i, a in enumerate(agents)} d_idx = {x: i for i, x in enumerate(dims)} n = len(tab) cols = [np.ones(n), np.array([1.0 if r["dose"] == 20 else 0.0 for r in tab]), np.array([1.0 if r["dose"] == 8 else 0.0 for r in tab])] names = ["const", "dose20", "dose8"] for a in agents[1:]: cols.append(np.array([1.0 if r["agent"] == a else 0.0 for r in tab])) names.append("ag_" + a) for x in dims[1:]: cols.append(np.array([1.0 if r["dim"] == x else 0.0 for r in tab])) names.append("dm_" + x) X = np.stack(cols, axis=1) y = np.array([float(r[outcome]) for r in tab]) cl = np.array([a_idx[r["agent"]] for r in tab]) dose = np.array([r["dose"] for r in tab]) return X, y, cl, names, dose def fit_models(tab): """FE regressions with CR1 cluster-robust SEs, pure numpy.""" out = {} for outcome in ["delta", "y_abs", "y_signed"]: if not any(outcome in r for r in tab): continue X, y, cl, names, dose = _design(tab, outcome) XtX_inv = np.linalg.pinv(X.T @ X) beta = XtX_inv @ (X.T @ y) resid = y - X @ beta G = int(cl.max()) + 1 meat = np.zeros((X.shape[1], X.shape[1])) for c in range(G): sel = cl == c uc = X[sel].T @ resid[sel] meat += np.outer(uc, uc) N, K = X.shape adj = G / max(G - 1, 1) * (N - 1) / max(N - K, 1) V = XtX_inv @ meat @ XtX_inv * adj ses = np.sqrt(np.diag(V)) out[outcome] = {nm: (float(b), float(s)) for nm, b, s in zip(names, beta, ses)} out[outcome + "_meta"] = {"N": N, "G": G, "dose_counts": {int(d): int((dose == d).sum()) for d in sorted(set(dose.tolist()))}} return out def summarize(rows): tab = primary_table(rows) conf = conformity_table(rows) out = {"n_reviews": len(rows), "n_v1_dim_obs": len(tab), "n_v2_dim_obs": len(conf)} if tab: m1 = fit_models(tab) out["n_v1_dim_obs"] = m1.get("delta_meta", {}).get("N", len(tab)) out["v1_models"] = m1 out["v1_dose_counts"] = m1.get("delta_meta", m1.get("y_abs_meta", {})).get("dose_counts") if conf: m2 = fit_models(conf) out["n_v2_dim_obs"] = len(conf) out["v2_models"] = m2 out["v2_dose_counts"] = m2.get("delta_meta", {}).get( "dose_counts", {str(d): sum(1 for r in conf if r["dose"] == d) for d in sorted(set(r["dose"] for r in conf))}) desc = {} for dose in sorted(set(r["dose"] for r in conf)): sub = [r for r in conf if r["dose"] == dose] w = [r["w_eff"] for r in sub] dl = [r["delta"] for r in sub] med_w = float(np.median(w)) # trimmed mean (trim 10% each side) ws = sorted(w); k = max(1, int(0.1 * len(ws))) trim_w = float(np.mean(ws[k:-k])) if len(ws) > 2 * k else float(np.mean(ws)) ds = sorted(dl) trim_d = float(np.mean(ds[k:-k])) if len(ds) > 2 * k else float(np.mean(ds)) zero_share = float(np.mean([1.0 if r["disp"] == "zero" else 0.0 for r in sub])) recov_n = sum(1 for r in sub if r.get("_recovered")) desc[str(dose)] = { "n": len(sub), "mean_w_eff": float(np.mean(w)), "mean_delta": float(np.mean(dl)), "median_w_eff": med_w, "trimmed_w_eff": trim_w, "trimmed_delta": trim_d, "zero_dispersion_share": zero_share, "recovered_entries": recov_n, } out["v2_descriptive_by_dose"] = desc return out if __name__ == "__main__": rows = load_obs() print(json.dumps(summarize(rows), indent=1, default=str))