| """ |
| 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_BASE = "https://export.arxiv.org/api/query" |
|
|
| DEFAULT_CATEGORIES = ["cs.LG", "cs.AI", "stat.ML", "cs.CL"] |
|
|
| |
| _NS = { |
| "atom": "http://www.w3.org/2005/Atom", |
| "arxiv": "http://arxiv.org/schemas/atom", |
| } |
|
|
| |
| |
| |
|
|
| |
| _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", |
| ] |
| ] |
|
|
| |
| _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). |
| """ |
| |
| for pat in _SKIP_PATTERNS: |
| if pat.search(abstract): |
| return False, f"skip_pattern_match: {pat.pattern[:60]}" |
|
|
| |
| 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" |
|
|
|
|
| |
| |
| |
|
|
| 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)"} |
|
|
| |
| 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...") |
|
|
| |
| 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 "" |
| |
| 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: 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: list[str] = [] |
| for cat_elem in entry.findall("atom:category", _NS): |
| term = cat_elem.get("term", "") |
| if term: |
| categories.append(term) |
|
|
| |
| 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_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 |
|
|
|
|
| |
| |
| |
|
|
| @dataclass |
| 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 |
| break |
|
|
| papers = parse_arxiv_xml(xml_str) |
| if not papers: |
| break |
|
|
| for paper in papers: |
| |
| if db.paper_exists(paper["arxiv_id"]): |
| result.skipped_existing += 1 |
| continue |
|
|
| |
| passed, reason = triage(paper["abstract"]) |
|
|
| |
| 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 |
|
|
| |
| 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() |
| |
| 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 "", |
| } |
|
|
| |
| 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() |
|
|
| |
| existing = db.find_paper(arxiv_id) |
| if existing: |
| return existing |
|
|
| paper = None |
| errors: list[str] = [] |
|
|
| |
| 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}") |
|
|
| |
| 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}") |
|
|
| |
| 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) |
| |
| _try_fetch_citations(paper["arxiv_id"], db) |
|
|
| return paper |
|
|