Download ingest.py from graziul/differance-engine: direct link, hf CLI and curl.
- Browser
- Download file 14.8 kB
-
https://huggingface.co/graziul/differance-engine/resolve/main/ingest.py
- Command line
-
hf download hf://graziul/differance-engine/ingest.py
-
curl -L -o ingest.py https://huggingface.co/graziul/differance-engine/resolve/main/ingest.py
14.8 kB
| """ | |
| arXiv ingestion module. | |
| Polls the arXiv API for new preprints in ML/AI categories, deduplicates by | |
| arXiv ID, and applies a triage gate: only papers whose abstracts contain | |
| novelty-claim language proceed to extraction. | |
| Categories polled: | |
| cs.LG — Machine Learning | |
| cs.AI — Artificial Intelligence | |
| stat.ML — Machine Learning (Statistics) | |
| cs.CL — Computation and Language (NLP) | |
| Rate limit: arXiv asks for polite delays (one call per 3 seconds). | |
| We use 5s between calls and limit to 50 results per call by default. | |
| """ | |
| from __future__ import annotations | |
| import re | |
| import time | |
| import urllib.parse | |
| import urllib.request | |
| import xml.etree.ElementTree as ET | |
| from dataclasses import dataclass, field | |
| from datetime import datetime, timezone | |
| from typing import Any | |
| from .db import Database, get_db | |
| # --------------------------------------------------------------------------- | |
| # arXiv API constants | |
| # --------------------------------------------------------------------------- | |
| ARXIV_API_BASE = "https://export.arxiv.org/api/query" | |
| DEFAULT_CATEGORIES = ["cs.LG", "cs.AI", "stat.ML", "cs.CL"] | |
| # arXiv API namespaces | |
| _NS = { | |
| "atom": "http://www.w3.org/2005/Atom", | |
| "arxiv": "http://arxiv.org/schemas/atom", | |
| } | |
| # --------------------------------------------------------------------------- | |
| # Triage gate | |
| # --------------------------------------------------------------------------- | |
| # Phrases that suggest the paper claims novelty (pass triage) | |
| _NOVELTY_PATTERNS: list[re.Pattern] = [ | |
| re.compile(p, re.IGNORECASE) | |
| for p in [ | |
| r"\bnovel\b", | |
| r"\bwe\s+(propose|introduce|present)\b", | |
| r"\bnew\s+(method|architecture|framework|approach|technique|algorithm|model|paradigm)\b", | |
| r"\bstate.of.the.art\b", | |
| r"\boutperforms?\b", | |
| r"\b(unlike|differs?\s+from|in\s+contrast\s+to)\s+(prior|previous|existing|traditional)\b", | |
| r"\badvances?\s+(the\s+)?(state|field)\b", | |
| r"\bfirst\s+(method|approach|architecture|time)\b", | |
| r"\b(breakthrough|groundbreaking|pioneering)\b", | |
| r"\bcontribution\b", | |
| r"\bwe\s+(achieve|obtain|demonstrate)\b", | |
| ] | |
| ] | |
| # Phrases that suggest the paper claims NO novelty (skip or flag) | |
| _SKIP_PATTERNS: list[re.Pattern] = [ | |
| re.compile(p, re.IGNORECASE) | |
| for p in [ | |
| r"\b(survey|review|tutorial)\s+(of|on)\b", | |
| r"\bcomprehensive\s+(survey|review)\b", | |
| r"\b(literature\s+review|related\s+work\b)", | |
| r"\b(benchmark|benchmarking)\b", | |
| r"\b(reproduce|replicate|reproduction)\b", | |
| r"\b(dataset|corpus|collection)\s+(release|introduction|description)\b", | |
| r"\b(position\s+paper|opinion|commentary)\b", | |
| r"\b(workshop|competition|challenge)\s+(report|summary|overview)\b", | |
| r"\b(extended\s+abstract|demo|poster)\b", | |
| ] | |
| ] | |
| def triage(abstract: str) -> tuple[bool, str]: | |
| """Determine whether an abstract passes the triage gate. | |
| Returns (passed, reason). | |
| """ | |
| # Check skip patterns first (hard no) | |
| for pat in _SKIP_PATTERNS: | |
| if pat.search(abstract): | |
| return False, f"skip_pattern_match: {pat.pattern[:60]}" | |
| # Check novelty patterns (soft yes) | |
| matches: list[str] = [] | |
| for pat in _NOVELTY_PATTERNS: | |
| m = pat.search(abstract) | |
| if m: | |
| matches.append(m.group(0)) | |
| if matches: | |
| return True, f"novelty_signals: {', '.join(matches[:3])}" | |
| return False, "no_novelty_signals_detected" | |
| # --------------------------------------------------------------------------- | |
| # arXiv API client | |
| # --------------------------------------------------------------------------- | |
| def _fetch_arxiv( | |
| categories: list[str] | None = None, | |
| max_results: int = 50, | |
| start: int = 0, | |
| sort_by: str = "submittedDate", | |
| sort_order: str = "descending", | |
| ) -> str: | |
| """Fetch raw XML from arXiv API. Returns the XML as a string. | |
| Tries urllib first, falls back to requests if urllib fails (some | |
| container environments have DNS/config issues with urllib). | |
| """ | |
| if categories is None: | |
| categories = DEFAULT_CATEGORIES | |
| cat_query = "+OR+".join(f"cat:{c}" for c in categories) | |
| params = { | |
| "search_query": cat_query, | |
| "start": str(start), | |
| "max_results": str(max_results), | |
| "sortBy": sort_by, | |
| "sortOrder": sort_order, | |
| } | |
| url = f"{ARXIV_API_BASE}?{urllib.parse.urlencode(params)}" | |
| headers = {"User-Agent": "DifferanceEngine/0.1 (mailto:chris@graziul.io)"} | |
| # Strategy 1: urllib | |
| try: | |
| req = urllib.request.Request(url, headers=headers) | |
| with urllib.request.urlopen(req, timeout=30) as resp: | |
| return resp.read().decode("utf-8") | |
| except Exception as e: | |
| print(f" [ingest] urllib fetch failed ({type(e).__name__}: {str(e)[:100]}), trying requests...") | |
| # Strategy 2: requests (more robust in containerized environments) | |
| import requests as _requests | |
| resp = _requests.get(url, headers=headers, timeout=30) | |
| resp.raise_for_status() | |
| return resp.text | |
| def parse_arxiv_xml(xml_str: str) -> list[dict]: | |
| """Parse arXiv API Atom XML into a list of paper dicts.""" | |
| root = ET.fromstring(xml_str) | |
| papers: list[dict] = [] | |
| for entry in root.findall("atom:entry", _NS): | |
| arxiv_id_full = entry.find("atom:id", _NS).text or "" | |
| # Strip the "http://arxiv.org/abs/" prefix to get the canonical ID | |
| arxiv_id = arxiv_id_full.split("/abs/")[-1] if "/abs/" in arxiv_id_full else arxiv_id_full | |
| title = " ".join((entry.find("atom:title", _NS).text or "").split()) | |
| abstract = " ".join((entry.find("atom:summary", _NS).text or "").split()) | |
| # Authors | |
| authors: list[str] = [] | |
| for author_elem in entry.findall("atom:author", _NS): | |
| name_elem = author_elem.find("atom:name", _NS) | |
| if name_elem is not None and name_elem.text: | |
| authors.append(name_elem.text.strip()) | |
| # Categories | |
| categories: list[str] = [] | |
| for cat_elem in entry.findall("atom:category", _NS): | |
| term = cat_elem.get("term", "") | |
| if term: | |
| categories.append(term) | |
| # Dates | |
| published = entry.find("atom:published", _NS) | |
| published_str = published.text if published is not None else "" | |
| updated = entry.find("atom:updated", _NS) | |
| updated_str = updated.text if updated is not None else "" | |
| # PDF link | |
| pdf_url = "" | |
| for link in entry.findall("atom:link", _NS): | |
| if link.get("title") == "pdf": | |
| pdf_url = link.get("href", "") | |
| break | |
| papers.append({ | |
| "arxiv_id": arxiv_id, | |
| "title": title, | |
| "abstract": abstract, | |
| "authors": authors, | |
| "categories": categories, | |
| "published": published_str, | |
| "updated": updated_str, | |
| "pdf_url": pdf_url, | |
| }) | |
| return papers | |
| # --------------------------------------------------------------------------- | |
| # Ingestion runner | |
| # --------------------------------------------------------------------------- | |
| class IngestResult: | |
| ingested: int = 0 | |
| triaged_in: int = 0 | |
| triaged_out: int = 0 | |
| skipped_existing: int = 0 | |
| def ingest_daily( | |
| db: Database | None = None, | |
| categories: list[str] | None = None, | |
| max_results: int = 50, | |
| max_pages: int = 2, | |
| ) -> IngestResult: | |
| """Run daily ingestion: fetch new papers, deduplicate, triage, store. | |
| Pages through results up to max_pages * max_results papers. | |
| """ | |
| if db is None: | |
| db = get_db() | |
| db.connect() | |
| if categories is None: | |
| categories = DEFAULT_CATEGORIES | |
| result = IngestResult() | |
| for page in range(max_pages): | |
| start = page * max_results | |
| try: | |
| xml_str = _fetch_arxiv( | |
| categories=categories, | |
| max_results=max_results, | |
| start=start, | |
| ) | |
| except Exception as e: | |
| print(f" [ingest] arXiv API error (page {page}, start={start}): {e}") | |
| if page == 0: | |
| raise # Fail hard on first page error; tolerate subsequent pages | |
| break | |
| papers = parse_arxiv_xml(xml_str) | |
| if not papers: | |
| break # No more results | |
| for paper in papers: | |
| # Deduplicate | |
| if db.paper_exists(paper["arxiv_id"]): | |
| result.skipped_existing += 1 | |
| continue | |
| # Triage | |
| passed, reason = triage(paper["abstract"]) | |
| # Store | |
| db.insert_paper(paper) | |
| db.update_triage(paper["arxiv_id"], passed, reason) | |
| result.ingested += 1 | |
| if passed: | |
| result.triaged_in += 1 | |
| else: | |
| result.triaged_out += 1 | |
| # Respect arXiv rate limit | |
| if page < max_pages - 1: | |
| time.sleep(5) | |
| return result | |
| def _fetch_paper_via_hf_hub(arxiv_id: str) -> dict | None: | |
| """Fallback: try to fetch paper metadata via huggingface_hub papers API. | |
| HF Hub mirrors arXiv metadata and may be reachable when arXiv is not. | |
| Returns a paper dict matching the arXiv parse format, or None. | |
| Uses the HF_TOKEN from the environment (automatically available in Spaces). | |
| The list_papers API was added in huggingface_hub 0.26+; if unavailable, | |
| falls back to searching daily papers. | |
| """ | |
| try: | |
| from huggingface_hub import HfApi | |
| api = HfApi() | |
| # Try the dedicated list_papers API (huggingface_hub >= 0.26) | |
| if hasattr(api, "list_papers"): | |
| papers = api.list_papers(query=arxiv_id, limit=1) | |
| paper_list = list(papers) | |
| if paper_list: | |
| p = paper_list[0] | |
| return { | |
| "arxiv_id": p.id or arxiv_id, | |
| "title": p.title or "", | |
| "abstract": p.summary or "", | |
| "authors": p.authors or [], | |
| "categories": p.tags or [], | |
| "published": p.published_at.isoformat() if p.published_at else "", | |
| "updated": p.updated_at.isoformat() if getattr(p, "updated_at", None) else "", | |
| "pdf_url": p.url_pdf or "", | |
| } | |
| # Fallback: search daily papers | |
| if hasattr(api, "search_papers"): | |
| papers = api.search_papers(query=arxiv_id, limit=1) | |
| paper_list = list(papers) | |
| if paper_list: | |
| p = paper_list[0] | |
| return { | |
| "arxiv_id": p.id or arxiv_id, | |
| "title": p.title or "", | |
| "abstract": p.summary or "", | |
| "authors": p.authors or [], | |
| "categories": p.tags or [], | |
| "published": p.published_at.isoformat() if p.published_at else "", | |
| "updated": p.updated_at.isoformat() if getattr(p, "updated_at", None) else "", | |
| "pdf_url": p.url_pdf or "", | |
| } | |
| print(f" [ingest] HF Hub paper API not available in this huggingface_hub version") | |
| return None | |
| except Exception as e: | |
| print(f" [ingest] HF Hub paper fallback also failed: {type(e).__name__}: {str(e)[:120]}") | |
| return None | |
| def _try_fetch_citations(arxiv_id: str, db: Database | None = None): | |
| """Fetch citation count from Semantic Scholar and update the DB. | |
| Runs synchronously but catches all errors — this is best-effort enrichment, | |
| not mission-critical. If it fails, the paper is still ingested. | |
| """ | |
| import json as _json | |
| try: | |
| url = f"https://api.semanticscholar.org/graph/v1/paper/ArXiv:{arxiv_id}?fields=citationCount" | |
| req = urllib.request.Request(url, headers={"User-Agent": "DifferanceEngine/1.0"}) | |
| with urllib.request.urlopen(req, timeout=10) as resp: | |
| data = _json.loads(resp.read()) | |
| count = data.get("citationCount", 0) | |
| if count and db: | |
| db.update_citation(arxiv_id, count) | |
| return count | |
| except Exception: | |
| return 0 | |
| def ingest_single( | |
| arxiv_id: str, | |
| db: Database | None = None, | |
| ) -> dict | None: | |
| """Ingest a single paper by arXiv ID. | |
| Returns the paper dict if found and ingested, None if not found. | |
| Tries arXiv API first, falls back to huggingface_hub papers API. | |
| """ | |
| if db is None: | |
| db = get_db() | |
| db.connect() | |
| # Check if already ingested (version-aware) — if so, return it | |
| existing = db.find_paper(arxiv_id) | |
| if existing: | |
| return existing | |
| paper = None | |
| errors: list[str] = [] | |
| # Strategy 1: arXiv API (primary) | |
| try: | |
| params = { | |
| "id_list": arxiv_id, | |
| "max_results": "1", | |
| } | |
| url = f"{ARXIV_API_BASE}?{urllib.parse.urlencode(params)}" | |
| req = urllib.request.Request(url) | |
| req.add_header("User-Agent", "DifferanceEngine/0.1 (mailto:chris@graziul.io)") | |
| with urllib.request.urlopen(req, timeout=30) as resp: | |
| xml_str = resp.read().decode("utf-8") | |
| papers = parse_arxiv_xml(xml_str) | |
| if papers: | |
| paper = papers[0] | |
| except Exception as e: | |
| err_msg = f"arXiv API: {type(e).__name__}: {str(e)[:120]}" | |
| errors.append(err_msg) | |
| print(f" [ingest] {err_msg}") | |
| # Strategy 2: Also try with requests (sometimes urllib fails on weird network configs) | |
| if paper is None: | |
| try: | |
| import requests as _requests | |
| params = { | |
| "id_list": arxiv_id, | |
| "max_results": "1", | |
| } | |
| url = f"{ARXIV_API_BASE}?{urllib.parse.urlencode(params)}" | |
| resp = _requests.get( | |
| url, | |
| headers={"User-Agent": "DifferanceEngine/0.1 (mailto:chris@graziul.io)"}, | |
| timeout=30, | |
| ) | |
| resp.raise_for_status() | |
| papers = parse_arxiv_xml(resp.text) | |
| if papers: | |
| paper = papers[0] | |
| except Exception as e: | |
| err_msg = f"arXiv via requests: {type(e).__name__}: {str(e)[:120]}" | |
| errors.append(err_msg) | |
| print(f" [ingest] {err_msg}") | |
| # Strategy 3: huggingface_hub papers API (last resort) | |
| if paper is None: | |
| print(f" [ingest] arXiv fetch failed, trying HF Hub papers API...") | |
| paper = _fetch_paper_via_hf_hub(arxiv_id) | |
| if paper is None: | |
| print(f" [ingest] All strategies failed for {arxiv_id}: {'; '.join(errors)}") | |
| return None | |
| passed, reason = triage(paper["abstract"]) | |
| db.insert_paper(paper) | |
| db.update_triage(paper["arxiv_id"], passed, reason) | |
| # Fetch citation count from Semantic Scholar in background | |
| _try_fetch_citations(paper["arxiv_id"], db) | |
| return paper | |