#1 ingest: normalize the pinned caller id (_norm_cited_id) so a non-bare DOI caller cannot store a URL-form papers.id that defeats DQ-5 idempotency / the startswith(10.) recovery gates. #2 quality.section_label: collapse to a controlled bucket only for an EXACT canonical heading; descriptive titles ('Abstract Nonsense...') keep their real title instead of being mislabelled. #4 ra_grobid_backfill: release the read connection before the slow GROBID network loop, fresh connection for the write (no idle-in-transaction across the loop over the flaky tunnel). #5/#10 tex: flatten_inputs strips unresolved input/include at the depth cap (no literal leak on cycles); _norm_texkey strips only a single leading ./ . #6/#7 arxiv.fetch_source: keep non-.tex members resolvable for input; pick primary on an UN-commented documentclass line. #13 is_arxiv_id: also exclude http:// and arXiv-DOI forms. Tests added/updated for each. Left as deliberate decisions: #3 (pre-section text drop is pre-existing in extract_sections; abstract stored separately), #8 (script normalizer is intentionally self-contained, already documented), #9/#11/#12 (no current trigger / tightening the gzip heuristic would reject valid old LaTeX like documentstyle / title truncation is by design on a write-only column). ruff + mypy clean. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
506 lines
22 KiB
Python
506 lines
22 KiB
Python
"""End-to-end idempotent ingest pipeline for a single paper."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from dataclasses import dataclass
|
|
from pathlib import Path
|
|
|
|
import numpy as np
|
|
import psycopg
|
|
|
|
from codex.config import get_settings
|
|
from codex.db import get_conn
|
|
from codex.embed import get_embedder
|
|
from codex.models import Citation, Paper
|
|
from codex.quality import is_quality_chunk, section_label
|
|
from codex.sources import arxiv, crossref, openalex, semanticscholar
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
@dataclass
|
|
class IngestResult:
|
|
paper_id: str
|
|
chunks_upserted: int
|
|
citations_upserted: int
|
|
formulas_upserted: int = 0
|
|
figures_upserted: int = 0
|
|
|
|
|
|
def _make_bibkey(paper: Paper) -> str | None:
|
|
"""Generate a BibTeX-style key from paper metadata.
|
|
|
|
Format: ConcatSurnames + Year (e.g. "BobenkoPinkallSpringborn2015").
|
|
Surnames are the last whitespace-separated token of each author name.
|
|
>3 authors: FirstSurname + "EtAl" + Year.
|
|
Returns None when authors or year are missing.
|
|
"""
|
|
if not paper.authors or not paper.year:
|
|
return None
|
|
surnames = [name.split()[-1] for name in paper.authors if name.strip()]
|
|
if not surnames:
|
|
return None
|
|
year = str(paper.year)
|
|
if len(surnames) <= 3:
|
|
return "".join(surnames) + year
|
|
return surnames[0] + "EtAl" + year
|
|
|
|
|
|
def _s2_id_for(pid: str) -> str | None:
|
|
"""Format a canonical paper id for the Semantic Scholar paper endpoint.
|
|
|
|
S2 needs a namespaced id (``arXiv:<id>``, ``DOI:<doi>``) or a bare 40-hex
|
|
S2 paperId — a bare arXiv id like ``2305.10988`` 404s. DOIs are detected by
|
|
the ``10.`` prefix; everything else is assumed to be an arXiv id (modern
|
|
``2305.10988`` or legacy ``math/0603097``).
|
|
"""
|
|
p = (pid or "").strip()
|
|
if not p:
|
|
return None
|
|
if p.lower().startswith(("arxiv:", "doi:")):
|
|
return p
|
|
if p.startswith("10."):
|
|
return f"DOI:{p}"
|
|
return f"arXiv:{p}"
|
|
|
|
|
|
def _norm_cited_id(cited_id: str) -> str:
|
|
"""Normalize a DOI cited-id to its canonical bare, lower-cased form.
|
|
|
|
DOIs are case-insensitive and may arrive bare (``10.…``), as a
|
|
``https://doi.org/…`` URL, or ``doi:…``-prefixed: GROBID emits whichever
|
|
form the PDF reference text used, while S2/OpenAlex hand back bare DOIs.
|
|
Stripping the prefix and lower-casing keeps the citation-graph join from
|
|
fragmenting the same reference into distinct case-/URL-variant nodes. arXiv
|
|
ids and S2 paperIds are left untouched. Mirrors ``openalex._normalize_doi``."""
|
|
s = cited_id.strip()
|
|
lowered = s.lower()
|
|
for prefix in ("https://doi.org/", "http://doi.org/", "doi:"):
|
|
if lowered.startswith(prefix):
|
|
return lowered[len(prefix) :]
|
|
return lowered if lowered.startswith("10.") else s
|
|
|
|
|
|
def _recover_abstract(paper: Paper) -> str | None:
|
|
"""Recover a missing abstract from S2, then Crossref (DQ-2 supplement).
|
|
|
|
Tried in order of coverage for this corpus: Semantic Scholar (indexes most
|
|
arXiv + journal abstracts), then Crossref (publisher-deposited, DOI only).
|
|
Network/parse failures degrade to None rather than aborting the ingest.
|
|
"""
|
|
s2_id = _s2_id_for(paper.id)
|
|
if s2_id is not None:
|
|
try:
|
|
abstract = semanticscholar.fetch_abstract(s2_id)
|
|
if abstract and abstract.strip():
|
|
return abstract.strip()
|
|
except Exception:
|
|
logger.warning("S2 abstract recovery failed for %s", paper.id, exc_info=True)
|
|
if paper.id.startswith("10."):
|
|
try:
|
|
abstract = crossref.fetch_abstract(paper.id)
|
|
if abstract and abstract.strip():
|
|
return abstract.strip()
|
|
except Exception:
|
|
logger.warning("Crossref abstract recovery failed for %s", paper.id, exc_info=True)
|
|
return None
|
|
|
|
|
|
def _s2_reference_supplement(paper: Paper) -> list[Citation]:
|
|
"""Fetch a paper's references from Semantic Scholar (DQ-1 supplement).
|
|
|
|
Used when OpenAlex returns an empty ``referenced_works`` list — common for
|
|
arXiv preprints and institutional theses that OpenAlex indexes without a
|
|
parsed bibliography. ``citing_id`` is rewritten to the canonical ``paper.id``
|
|
so the FK holds, and cited DOIs are case-normalised. Network/parse failures
|
|
degrade to an empty list rather than aborting the ingest.
|
|
"""
|
|
s2_id = _s2_id_for(paper.id)
|
|
if s2_id is None:
|
|
return []
|
|
try:
|
|
raw = semanticscholar.fetch_references(s2_id)
|
|
except Exception:
|
|
logger.warning("S2 reference supplement failed for %s", paper.id, exc_info=True)
|
|
return []
|
|
return [
|
|
Citation(citing_id=paper.id, cited_id=_norm_cited_id(c.cited_id), context=c.context)
|
|
for c in raw
|
|
if c.cited_id
|
|
]
|
|
|
|
|
|
def _crossref_reference_supplement(paper: Paper) -> list[Citation]:
|
|
"""Fetch a paper's references from Crossref (roadmap R-B supplement).
|
|
|
|
Third leg of the citation fallback chain (OpenAlex-empty → S2 → Crossref).
|
|
Crossref carries publisher-deposited reference lists keyed on the citing
|
|
work's DOI, so only DOI papers can be looked up (arXiv-only ids → []).
|
|
``citing_id`` is rewritten to the canonical ``paper.id`` and cited DOIs are
|
|
case-normalised. Network/parse failures degrade to [] rather than aborting.
|
|
"""
|
|
if not paper.id.startswith("10."):
|
|
return []
|
|
try:
|
|
raw = crossref.fetch_references(paper.id)
|
|
except Exception:
|
|
logger.warning("Crossref reference supplement failed for %s", paper.id, exc_info=True)
|
|
return []
|
|
return [
|
|
Citation(citing_id=paper.id, cited_id=_norm_cited_id(c.cited_id), context=c.context)
|
|
for c in raw
|
|
if c.cited_id
|
|
]
|
|
|
|
|
|
def ingest_paper(
|
|
paper_id: str,
|
|
source_path: str | None = None,
|
|
rich: bool = False,
|
|
) -> IngestResult:
|
|
"""Idempotent ingest of one paper.
|
|
|
|
Parameters
|
|
----------
|
|
paper_id:
|
|
An arXiv ID (``"2301.07041"``), DOI (``"10.1145/…"``), or
|
|
OpenAlex W-ID (``"W2741809807"``).
|
|
source_path:
|
|
Optional path to a local ``.tex`` or ``.pdf`` file.
|
|
If given, the file is parsed into text chunks and stored.
|
|
Supports:
|
|
- ``.tex`` → :func:`codex.parsing.tex.latex_to_text` + chunk
|
|
- ``.pdf`` → :func:`codex.parsing.nougat.pdf_to_markdown` + chunk
|
|
+ :func:`codex.parsing.grobid.extract_references` for refs
|
|
rich:
|
|
When True and *source_path* is a PDF, also extract formulas via
|
|
:func:`codex.parsing.mathpix.extract_formulas` and figures via
|
|
:func:`codex.parsing.figures.extract_figures` (F-09).
|
|
|
|
Returns
|
|
-------
|
|
IngestResult
|
|
Counts of upserted chunks, citations, formulas, and figures.
|
|
"""
|
|
# ---------------------------------------------------------------
|
|
# 1. Fetch metadata (OpenAlex primary, SemanticScholar fallback)
|
|
# ---------------------------------------------------------------
|
|
paper: Paper | None = openalex.fetch_paper(paper_id)
|
|
|
|
if paper is None:
|
|
# Detect arXiv IDs: numeric pattern like "2301.07041" or "arxiv:..." prefix
|
|
pid_lower = paper_id.lower()
|
|
looks_like_arxiv = pid_lower.startswith("arxiv:") or (
|
|
len(paper_id) > 4 and paper_id[4:5] == "." and paper_id[:4].isdigit()
|
|
)
|
|
if looks_like_arxiv:
|
|
# OpenAlex 404 on an arXiv id → recover authoritative metadata from
|
|
# the arXiv API (DQ-2); fall back to a minimal stub if arXiv also
|
|
# has nothing. Citation recovery from S2 happens in the citation
|
|
# block below (DQ-1 supplement).
|
|
paper = arxiv.fetch_metadata(paper_id) or Paper(id=paper_id, title="")
|
|
else:
|
|
raise ValueError(f"Paper not found: {paper_id}")
|
|
|
|
if paper is None:
|
|
raise ValueError(f"Paper not found: {paper_id}")
|
|
|
|
# Canonical identity is the caller-supplied id — what the CLI / ingest scripts
|
|
# use and what papers.id is contracted to hold (arXiv id or DOI). Pin it to the
|
|
# caller's id (audit C-7), but normalize it the same way cited-ids are (strip a
|
|
# https://doi.org/ or doi: prefix, lower-case DOIs) so a non-bare caller cannot
|
|
# store a URL-form papers.id that would defeat DQ-5's ON CONFLICT (id) idempotency
|
|
# or the `paper.id.startswith("10.")` recovery gates below. arXiv/W ids pass
|
|
# through unchanged. The OpenAlex work id is retained as openalex_id / in
|
|
# paper_identifiers.
|
|
paper.id = _norm_cited_id(paper_id.strip())
|
|
|
|
# DQ-2: OpenAlex often has the paper but no abstract (it cannot redistribute
|
|
# some publishers' abstracts). Supplement from S2, then Crossref (uses the
|
|
# canonical paper.id set above), so the paper gets a real (non-zero) abstract
|
|
# embedding instead of a zero vector.
|
|
if not (paper.abstract or "").strip():
|
|
recovered = _recover_abstract(paper)
|
|
if recovered:
|
|
paper.abstract = recovered
|
|
|
|
# Auto-generate bibkey when source has none (D-04).
|
|
if paper.bibkey is None:
|
|
paper.bibkey = _make_bibkey(paper)
|
|
|
|
# ---------------------------------------------------------------
|
|
# 2. Embed abstract (dense only — schema has one vector column)
|
|
# ---------------------------------------------------------------
|
|
embedder = get_embedder()
|
|
abstract_text = paper.abstract or ""
|
|
if abstract_text:
|
|
abstract_emb_arr: np.ndarray = embedder.encode_dense([abstract_text])
|
|
abstract_emb: list[float] = abstract_emb_arr[0].tolist()
|
|
else:
|
|
dim = embedder.dim
|
|
abstract_emb = [0.0] * dim
|
|
|
|
# ---------------------------------------------------------------
|
|
# 3. Upsert paper (INSERT … ON CONFLICT (id) DO UPDATE SET …)
|
|
# ---------------------------------------------------------------
|
|
upsert_sql = """
|
|
INSERT INTO papers (id, openalex_id, bibkey, title, authors, year, abstract,
|
|
source_path, abstract_emb)
|
|
VALUES (%(id)s, %(openalex_id)s, %(bibkey)s, %(title)s, %(authors)s, %(year)s,
|
|
%(abstract)s, %(source_path)s, %(abstract_emb)s)
|
|
ON CONFLICT (id) DO UPDATE SET
|
|
openalex_id = EXCLUDED.openalex_id,
|
|
bibkey = COALESCE(papers.bibkey, EXCLUDED.bibkey),
|
|
title = EXCLUDED.title,
|
|
authors = EXCLUDED.authors,
|
|
year = EXCLUDED.year,
|
|
abstract = EXCLUDED.abstract,
|
|
source_path = COALESCE(EXCLUDED.source_path, papers.source_path),
|
|
abstract_emb = EXCLUDED.abstract_emb
|
|
"""
|
|
with get_conn() as conn:
|
|
# The id-keyed upsert does not catch a papers.bibkey UNIQUE collision (two
|
|
# same-author/year papers generate the same key). On such a collision retry
|
|
# with an a/b/c… suffix instead of failing the whole ingest (audit C-15).
|
|
base_bibkey = paper.bibkey
|
|
for attempt in range(27):
|
|
try:
|
|
conn.execute(
|
|
upsert_sql,
|
|
{
|
|
"id": paper.id,
|
|
"openalex_id": paper.openalex_id,
|
|
"bibkey": paper.bibkey,
|
|
"title": paper.title,
|
|
"authors": paper.authors,
|
|
"year": paper.year,
|
|
"abstract": paper.abstract,
|
|
"source_path": source_path,
|
|
"abstract_emb": abstract_emb,
|
|
},
|
|
)
|
|
break
|
|
except psycopg.errors.UniqueViolation as exc:
|
|
msg = str(exc).lower()
|
|
conn.rollback()
|
|
if "openalex" in msg:
|
|
# Two different papers.id claim the same OpenAlex work — a real data
|
|
# conflict, not something to auto-rename. Surface it clearly instead
|
|
# of the raw psycopg error (companion gap to C-15's bibkey case).
|
|
raise ValueError(
|
|
f"openalex_id {paper.openalex_id!r} is already attached to a "
|
|
f"different paper; ingesting '{paper.id}' would duplicate it — "
|
|
"resolve the conflicting paper first."
|
|
) from exc
|
|
if base_bibkey is None or "bibkey" not in msg or attempt == 26:
|
|
raise
|
|
paper.bibkey = f"{base_bibkey}{chr(ord('a') + attempt)}"
|
|
logger.warning(
|
|
"bibkey %r already exists; ingesting '%s' as %r instead",
|
|
base_bibkey,
|
|
paper.id,
|
|
paper.bibkey,
|
|
)
|
|
|
|
# Register openalex_id in paper_identifiers (alias table).
|
|
# Every ingest records the primary openalex_id so the graph JOIN
|
|
# resolves it even if papers.openalex_id is later updated.
|
|
if paper.openalex_id:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO paper_identifiers (paper_id, openalex_id)
|
|
VALUES (%s, %s)
|
|
ON CONFLICT DO NOTHING
|
|
""",
|
|
(paper.id, paper.openalex_id),
|
|
)
|
|
|
|
# ---------------------------------------------------------------
|
|
# 4. Parse + chunk + embed source file (if provided)
|
|
# DELETE existing chunks for this paper first, then bulk INSERT
|
|
# ---------------------------------------------------------------
|
|
chunks_upserted = 0
|
|
pdf_citations: list[Citation] = []
|
|
|
|
if source_path is not None:
|
|
from codex.parsing.grobid import extract_references
|
|
from codex.parsing.nougat import pdf_to_markdown
|
|
from codex.parsing.tex import chunk_sections, chunk_text
|
|
|
|
suffix = Path(source_path).suffix.lower()
|
|
# (section_title | None, chunk_text) pairs. `.tex` is section-aware so
|
|
# the stored `section` column reflects the real heading (R-C / R-12);
|
|
# `.txt`/`.pdf` carry no section structure, so the title is None.
|
|
raw_pairs: list[tuple[str | None, str]] = []
|
|
|
|
if suffix == ".tex":
|
|
raw_pairs = chunk_sections(Path(source_path).read_text())
|
|
elif suffix == ".txt":
|
|
text = Path(source_path).read_text(encoding="utf-8", errors="replace")
|
|
raw_pairs = [(None, c) for c in chunk_text(text)]
|
|
elif suffix == ".pdf":
|
|
raw_pairs = [(None, c) for c in chunk_text(pdf_to_markdown(source_path))]
|
|
# Also extract GROBID refs
|
|
grobid_refs = extract_references(source_path)
|
|
for ref in grobid_refs:
|
|
cited_id = ref.get("doi") or ref.get("arxiv_id")
|
|
if cited_id:
|
|
# Normalize the GROBID DOI (bare/URL/``doi:`` → bare lower-case)
|
|
# so PDF-derived edges join the same graph nodes as the
|
|
# S2/OpenAlex paths; arXiv ids pass through untouched.
|
|
pdf_citations.append(
|
|
Citation(citing_id=paper.id, cited_id=_norm_cited_id(cited_id))
|
|
)
|
|
else:
|
|
logger.warning("Unbekannter Dateityp: %s — kein Text-Parsing", source_path)
|
|
|
|
# Quality-gate the pairs (keeps each chunk aligned with its section title).
|
|
settings = get_settings()
|
|
kept_pairs = [
|
|
(title, content)
|
|
for title, content in raw_pairs
|
|
if is_quality_chunk(content, settings=settings)
|
|
]
|
|
|
|
# Delete existing chunks for this paper
|
|
conn.execute("DELETE FROM chunks WHERE paper_id = %s", (paper.id,))
|
|
|
|
if kept_pairs:
|
|
# Embed all chunks in one batch
|
|
chunk_contents = [content for _, content in kept_pairs]
|
|
chunk_embeddings = embedder.encode_dense(chunk_contents)
|
|
chunk_rows = [
|
|
(
|
|
paper.id,
|
|
ord_idx,
|
|
content,
|
|
chunk_embeddings[ord_idx].tolist(),
|
|
# .tex: label from the real \section heading — keeps the
|
|
# controlled bucket when it matches, else stores the cleaned
|
|
# title (R-F); .txt/.pdf (title=None) classify by content.
|
|
# NB: `section` semantics now differ by source (write-only col).
|
|
section_label(title, content),
|
|
)
|
|
for ord_idx, (title, content) in enumerate(kept_pairs)
|
|
]
|
|
with conn.cursor() as cur:
|
|
cur.executemany(
|
|
"INSERT INTO chunks (paper_id, ord, content, embedding, section)"
|
|
" VALUES (%s, %s, %s, %s, %s)",
|
|
chunk_rows,
|
|
)
|
|
chunks_upserted = len(chunk_rows)
|
|
|
|
# ---------------------------------------------------------------
|
|
# 5. Fetch citations (fallback chain: OpenAlex → S2 → Crossref)
|
|
# Insert with ON CONFLICT (citing_id, cited_id) DO NOTHING
|
|
# ---------------------------------------------------------------
|
|
api_citations: list[Citation]
|
|
if paper.openalex_id:
|
|
raw = openalex.fetch_citations(paper.openalex_id)
|
|
# OpenAlex returns citing_id=openalex_id, but papers.id uses the DOI.
|
|
# Rewrite citing_id to paper.id so the FK constraint is satisfied.
|
|
api_citations = [
|
|
Citation(citing_id=paper.id, cited_id=c.cited_id, context=c.context) for c in raw
|
|
]
|
|
# DQ-1: OpenAlex indexes many arXiv preprints / theses WITHOUT a
|
|
# parsed reference list (referenced_works == []). Fall back to the
|
|
# Semantic Scholar references, which often has them.
|
|
if not api_citations:
|
|
api_citations = _s2_reference_supplement(paper)
|
|
else:
|
|
api_citations = _s2_reference_supplement(paper)
|
|
# R-B: if OpenAlex and S2 both came back empty, try Crossref's
|
|
# publisher-deposited references (DOI papers only) as a third leg.
|
|
if not api_citations:
|
|
api_citations = _crossref_reference_supplement(paper)
|
|
|
|
# Merge API citations and GROBID PDF citations; dedup via set
|
|
all_citations_set: set[tuple[str, str]] = set()
|
|
merged_citations: list[Citation] = []
|
|
for cit in api_citations + pdf_citations:
|
|
key = (cit.citing_id, cit.cited_id)
|
|
if key not in all_citations_set:
|
|
all_citations_set.add(key)
|
|
merged_citations.append(cit)
|
|
|
|
if merged_citations:
|
|
with conn.cursor() as cur:
|
|
cur.executemany(
|
|
"""
|
|
INSERT INTO citations (citing_id, cited_id, context)
|
|
VALUES (%s, %s, %s)
|
|
ON CONFLICT (citing_id, cited_id) DO NOTHING
|
|
""",
|
|
[(c.citing_id, c.cited_id, c.context) for c in merged_citations],
|
|
)
|
|
|
|
citations_upserted = len(merged_citations)
|
|
conn.commit()
|
|
|
|
# ---------------------------------------------------------------
|
|
# 6. F-09 Rich Parsing: formulas + figures (PDF only)
|
|
# ---------------------------------------------------------------
|
|
formulas_upserted = 0
|
|
figures_upserted = 0
|
|
|
|
if rich and source_path is not None and Path(source_path).suffix.lower() == ".pdf":
|
|
from codex.parsing.figures import extract_figures
|
|
from codex.parsing.mathpix import extract_formulas
|
|
|
|
settings = get_settings()
|
|
|
|
formulas = extract_formulas(source_path)
|
|
figures = extract_figures(source_path, output_dir=settings.figures_dir)
|
|
|
|
if formulas or figures:
|
|
with get_conn() as conn2:
|
|
with conn2.cursor() as cur:
|
|
# Delete-before-insert keeps re-ingest idempotent (BIGSERIAL has no
|
|
# natural UNIQUE key, so ON CONFLICT DO NOTHING never fires).
|
|
cur.execute("DELETE FROM formulas WHERE paper_id = %s", (paper.id,))
|
|
cur.execute("DELETE FROM figures WHERE paper_id = %s", (paper.id,))
|
|
|
|
if formulas:
|
|
with conn2.cursor() as cur:
|
|
cur.executemany(
|
|
"""
|
|
INSERT INTO formulas (paper_id, page, raw_latex, context, eq_label)
|
|
VALUES (%s, %s, %s, %s, %s)
|
|
""",
|
|
[
|
|
# paper.id, not f.paper_id: the extractor derives its
|
|
# paper_id from the PDF filename stem, which need not match
|
|
# papers.id and would violate the FK (audit C-10).
|
|
(paper.id, f.page, f.raw_latex, f.context, f.eq_label)
|
|
for f in formulas
|
|
],
|
|
)
|
|
formulas_upserted = len(formulas)
|
|
|
|
if figures:
|
|
with conn2.cursor() as cur:
|
|
cur.executemany(
|
|
"""
|
|
INSERT INTO figures (paper_id, page, caption, image_path)
|
|
VALUES (%s, %s, %s, %s)
|
|
""",
|
|
[
|
|
# paper.id, not fig.paper_id (filename stem) — see C-10 above.
|
|
(paper.id, fig.page, fig.caption, fig.image_path)
|
|
for fig in figures
|
|
],
|
|
)
|
|
figures_upserted = len(figures)
|
|
|
|
conn2.commit()
|
|
|
|
return IngestResult(
|
|
paper_id=paper.id,
|
|
chunks_upserted=chunks_upserted,
|
|
citations_upserted=citations_upserted,
|
|
formulas_upserted=formulas_upserted,
|
|
figures_upserted=figures_upserted,
|
|
)
|