AutonomousAgent / batch_ingest.py
Mathias Heider
Claude Fable 5.1
Never store a credential: redact at every DB write and at boot; keep PDFs when Gemini refuses for budget reasons
00f9c01 unverified
Raw History Blame Contribute Delete
50.2 kB
r"""
batch_ingest.py — preliminary autonomous batch-ingestion prototype for the
AIM Composites Materials Database.
This script takes a folder of PDFs as input and exercises the extraction +
validation + insertion stages of the target autonomous-ingestion architecture
described in Section 6 of the project report:
PDFs -> Gemini extraction -> grounding+validation -> dedup -> SQLite
\-> review_queue.csv (a view)
It intentionally does NOT do source discovery (component 1) and does NOT run
the plot-extraction / image-mapping pipeline. It is a minimal closed loop
sufficient to characterize the cost, throughput, and failure modes of the
extraction-plus-validation stage on a fixed corpus.
As of the extraction-hardening phase, all extraction logic (prompt, schema,
grounding, unit normalization, classification, dedup keys) lives in
``extraction.py`` — this file is just the driver + SQLite mirror. Nothing is
silently dropped: flagged rows are inserted with a ``status`` and
``flag_reason``, and ``review_queue.csv`` is a ``SELECT ... WHERE status != 'ok'``
export rather than a discard pile.
Usage:
export GEMINI_API_KEY=...
python batch_ingest.py --input ./pdfs --db ./materials_mirror.sqlite \
--review review_queue.csv --report run_report.json
# one-off, non-destructive column migration of an existing DB:
python batch_ingest.py --migrate --db ./materials_mirror.sqlite
# write to the shared Postgres (the DB the HF Space reads) instead of
# SQLite — env: DB_HOST/DB_PORT/DB_NAME/DB_USER/DB_PASSWORD or DATABASE_URL.
# Requires a one-time `python pg_migrate.py --apply` first (see pg_mirror.py):
python batch_ingest.py --pg --input ./pdfs
# figure & graph mining (opt-in; SQLite only; see FIGURES.md):
python batch_ingest.py --figures --input ./pdfs --db ./materials_mirror.sqlite
python batch_ingest.py --figures --no-figure-mining --input ./pdfs # harvest+classify only
# figure citation linking (opt-in; local, no API calls; works with --pg;
# ported from the InDeS mapper — see FIGURES.md "Figure citation linking"):
python batch_ingest.py --link-figures --input ./pdfs --links-csv figure_links.csv
python batch_ingest.py --link-figures --embed-figure-images --pg --input ./pdfs
Author: Mathias Heider, ME8930 course project, May 2026.
Extraction prompt and schema now centralized in extraction.py (was adapted from
the live Streamlit app's page_files/categorized/Backend/upload_backend.py,
co-developed with Abhijit on the AIM Composites HF Space).
"""
from __future__ import annotations
import argparse
import dataclasses
import hashlib
import json
import logging
import os
import sqlite3
import sys
import time
from collections.abc import Iterable
from pathlib import Path
from typing import Any, Optional
import requests
import extraction
import migrate as migrate_mod
from extraction import Extraction, PropertyRow, extract_from_pdf, to_rows, verify_against_text
from migrate import (
EXTRA_COLUMNS,
backfill_material_key_grade,
ensure_columns,
ensure_figures_table,
ensure_sources_sha1_unique,
)
# ---------------------------------------------------------------------------
# SQLite mirror of the Postgres schema
# ---------------------------------------------------------------------------
# Base (legacy) columns, kept so the CSV export and page1.py keep working. The
# hardening-phase columns are added on top by migrate.ensure_columns().
SCHEMA_DDL = """
CREATE TABLE IF NOT EXISTS Polymers (
id INTEGER PRIMARY KEY AUTOINCREMENT,
material_name TEXT, material_abbreviation TEXT, section TEXT,
property_name TEXT, value TEXT, unit TEXT, english TEXT,
test_condition TEXT, comments TEXT
);
CREATE TABLE IF NOT EXISTS Fibers (
id INTEGER PRIMARY KEY AUTOINCREMENT,
material_name TEXT, material_abbreviation TEXT, section TEXT,
property_name TEXT, value TEXT, unit TEXT, english TEXT,
test_condition TEXT, comments TEXT
);
CREATE TABLE IF NOT EXISTS Composites_materials (
id INTEGER PRIMARY KEY AUTOINCREMENT,
material_name TEXT, material_abbreviation TEXT, section TEXT,
property_name TEXT, value TEXT, unit TEXT, english TEXT,
test_condition TEXT, comments TEXT
);
CREATE TABLE IF NOT EXISTS sources (
id INTEGER PRIMARY KEY AUTOINCREMENT,
pdf_filename TEXT,
pdf_sha1 TEXT UNIQUE,
ingested_at TEXT,
material_class TEXT,
material_abbreviation TEXT
);
"""
# `sources` identity is the content hash, not the basename: two different PDFs
# that share a filename (vendorA/datasheet.pdf vs vendorB/datasheet.pdf — likely,
# since --input is rglob'd) used to collide on a `pdf_filename UNIQUE` +
# INSERT OR IGNORE, so the second was never recorded and got re-sent to Gemini
# on every run. migrate.ensure_sources_sha1_unique() rebuilds legacy tables.
TABLE_FOR_CLASS = {
"Polymer": "Polymers",
"Fiber": "Fibers",
"Composite": "Composites_materials",
}
ALL_TABLES = tuple(TABLE_FOR_CLASS.values())
# ---------------------------------------------------------------------------
# Run bookkeeping
# ---------------------------------------------------------------------------
@dataclasses.dataclass
class PdfResult:
pdf: str
elapsed_s: float
materials: int
extracted: int
inserted: int
flagged: int # inserted with status != 'ok'
duplicates: int
material_classes: list[str]
error: Optional[str] = None
# figure-mining phase (all zero when --figures is off)
figures_found: int = 0 # harvested PNGs (after junk filters + cap)
figures_mined: int = 0 # plot/table figures that returned a readout
figure_rows: int = 0 # origin='figure' rows inserted
figure_duplicates: int = 0 # figure rows skipped by the dedup grain
vision_calls: int = 0 # classify + mining Gemini calls
figure_error: Optional[str] = None # non-fatal: text rows still inserted
# Gemini refused for budget reasons (spending cap / quota): the PDF was
# not processed and will be once the budget is back; nothing to retry now.
quota: bool = False
figure_filters: dict[str, int] = dataclasses.field(default_factory=dict)
# Gemini usage for this PDF: text call + every vision call (0 when no
# call was made; tokens_out includes thinking tokens, billed as output)
tokens_in: int = 0
tokens_out: int = 0
# figure-linking phase (all zero when --link-figures is off)
rows_linked: int = 0 # text rows that got a figure link (new or backfilled)
link_figures_found: int = 0 # figures harvested for the link pass (local, no API call)
link_stats: dict[str, int] = dataclasses.field(default_factory=dict)
link_error: Optional[str] = None # non-fatal: rows are inserted unlinked
# ---------------------------------------------------------------------------
# Database operations
# ---------------------------------------------------------------------------
def init_db(path: Path) -> sqlite3.Connection:
conn = sqlite3.connect(str(path))
conn.executescript(SCHEMA_DDL)
# Add the hardening-phase columns if they aren't there yet (idempotent).
for table in ALL_TABLES:
ensure_columns(conn, table)
# Rows written before trade_grade joined material_key get re-keyed
# once, so a re-ingest dedups against them instead of doubling them.
backfill_material_key_grade(conn, table)
# Legacy DBs keyed `sources` on pdf_filename; rebuild to pdf_sha1 (idempotent).
ensure_sources_sha1_unique(conn)
# Figure provenance table (figure-mining phase; idempotent).
ensure_figures_table(conn)
conn.commit()
return conn
def run_migrate(db_path: Path) -> None:
"""`--migrate`: the same migration as `python migrate.py --db`, including
the .bak backup of an existing DB. init_db afterwards creates any base
tables the file does not have (an empty/stub file — a touch, an aborted
run — must still end up with the full schema, as the old path guaranteed)."""
if db_path.exists():
migrate_mod.migrate(db_path)
init_db(db_path).close()
# Column order used for inserts: legacy columns first (positional compat), then
# the hardening-phase columns. Mirrors migrate.EXTRA_COLUMNS.
_LEGACY_COLS = [
"material_name", "material_abbreviation", "section", "property_name",
"value", "unit", "english", "test_condition", "comments",
]
_INSERT_COLS = _LEGACY_COLS + [name for name, _type in EXTRA_COLUMNS]
def _row_values(row: PropertyRow) -> tuple:
from datetime import datetime, timezone
# Export guard (figure-mining phase). Every consumer that shows rows to
# the app filters on status='ok', so status is the publish gate. A figure
# row must therefore never be INSERTED as 'ok' — only --promote may set
# that, after a human looked at the PNG. Enforced here, the single point
# both the SQLite and Postgres insert paths go through.
status, flag_reason = row.status, row.flag_reason
if (row.origin or "text") == "figure" and status == "ok":
status = "figure_estimate"
flag_reason = ("figure row inserted with status=ok; downgraded — only "
"--promote may publish a figure reading"
+ (f"; {row.flag_reason}" if row.flag_reason else ""))
mapping = {
"material_name": row.material_name,
"material_abbreviation": row.material_abbreviation,
"section": row.section,
"property_name": row.property_name,
"value": row.value,
"unit": row.unit,
"english": row.english,
"test_condition": row.test_condition,
"comments": row.comments,
"material_key": row.material_key,
"material_class": row.material_class,
"trade_grade": row.trade_grade,
"manufacturer": row.manufacturer,
"matrix": row.matrix,
"fiber": row.fiber,
"fiber_volume_fraction": row.fiber_volume_fraction,
"value_raw": row.value_raw,
"value_num": row.value_num,
"value_min": row.value_min,
"value_max": row.value_max,
"qualifier": row.qualifier,
"unit_canonical": row.unit_canonical,
"value_si": row.value_si,
"source_pdf": row.source_pdf,
"source_sha1": row.source_sha1,
"page": row.page,
"source_quote": row.source_quote,
"status": status,
"flag_reason": flag_reason,
"model": row.model,
"prompt_version": row.prompt_version,
"extracted_at": datetime.now(timezone.utc).isoformat(timespec="seconds"),
"origin": row.origin or "text",
"figure_id": row.figure_id or None,
# figure-linking phase (None when the row links to nothing)
"figure_ref": row.figure_ref or None,
"figure_link_score": row.figure_link_score,
"figure_link_signals": row.figure_link_signals or None,
"image_url": row.image_url or None,
"image": row.image or None,
}
return tuple(mapping[c] for c in _INSERT_COLS)
def seen_sha1(conn: sqlite3.Connection, sha1: str) -> bool:
"""True if this exact PDF (by sha1) was already ingested."""
cur = conn.execute("SELECT 1 FROM sources WHERE pdf_sha1 = ? LIMIT 1", (sha1,))
return cur.fetchone() is not None
def already_inserted(conn: sqlite3.Connection, table: str, row: PropertyRow) -> bool:
"""Source-aware dedup (Task 6).
Grain = (source_sha1, material_key, section, property_name, test_condition,
value_raw, origin). Skip only a true re-ingest of the *same measurement
from the same PDF*; the same property from a different PDF (independent
repeat) is kept. `origin` is part of the grain on purpose: a text row and
a figure row reporting the same number both survive — that agreement is
signal, not duplication (figure-mining phase).
"""
cur = conn.execute(
f"SELECT 1 FROM {table} "
f"WHERE IFNULL(source_sha1,'') = IFNULL(?, '') "
f" AND IFNULL(material_key,'') = IFNULL(?, '') "
f" AND IFNULL(section,'') = IFNULL(?, '') "
f" AND IFNULL(property_name,'') = IFNULL(?, '') "
f" AND IFNULL(test_condition,'') = IFNULL(?, '') "
f" AND IFNULL(value_raw,'') = IFNULL(?, '') "
f" AND IFNULL(origin,'text') = IFNULL(?, 'text') "
f"LIMIT 1",
(row.source_sha1, row.material_key, row.section,
row.property_name, row.test_condition, row.value_raw,
row.origin or "text"),
)
return cur.fetchone() is not None
_FIGURE_COLS = [
"figure_id", "source_pdf", "source_sha1", "page", "bbox", "caption",
"figure_kind", "material_key", "image_path", "image_sha256", "width_px",
"height_px", "route", "mining_status", "n_values", "model",
"figure_prompt_version", "extracted_at",
]
def upsert_figure(conn: sqlite3.Connection, fig: Any, png_bytes: Optional[bytes] = None) -> None:
"""Insert or refresh one harvested figure's provenance row (keyed on
figure_id = sha of the PNG, so a re-harvest is idempotent while a later
classify/mine pass can update kind/status).
``png_bytes``: store the PNG itself (figure-linking phase: only for
figures a text row links to). An upsert without bytes never clears bytes
stored earlier."""
from datetime import datetime, timezone
vals = {
"figure_id": fig.figure_id, "source_pdf": fig.source_pdf,
"source_sha1": fig.source_sha1, "page": fig.page,
"bbox": json.dumps(list(fig.bbox)), "caption": fig.caption,
"figure_kind": fig.figure_kind, "material_key": fig.material_key,
"image_path": fig.image_path, "image_sha256": fig.image_sha256,
"width_px": fig.width_px, "height_px": fig.height_px, "route": fig.route,
"mining_status": fig.mining_status, "n_values": fig.n_values,
"model": fig.model, "figure_prompt_version": fig.figure_prompt_version,
"extracted_at": datetime.now(timezone.utc).isoformat(timespec="seconds"),
"png_bytes": png_bytes or None,
}
all_cols = _FIGURE_COLS + ["png_bytes"]
cols = ", ".join(all_cols)
ph = ", ".join("?" for _ in all_cols)
upd = ", ".join(f"{c}=excluded.{c}" for c in _FIGURE_COLS if c != "figure_id")
upd += ", png_bytes=COALESCE(excluded.png_bytes, figures.png_bytes)"
conn.execute(
f"INSERT INTO figures ({cols}) VALUES ({ph}) "
f"ON CONFLICT(figure_id) DO UPDATE SET {upd}",
tuple(vals[c] for c in all_cols),
)
# --- figure-linking phase: read rows back for a link backfill ----------------
_LINK_BACKFILL_COLS = [
"id", "material_name", "material_abbreviation", "material_key", "material_class",
"section", "property_name", "value_raw", "unit", "test_condition", "comments",
"page", "source_quote",
]
def rows_for_link_backfill(conn: sqlite3.Connection, table: str, sha1: str) -> list[tuple[int, PropertyRow]]:
"""TEXT rows of an already-ingested PDF that carry no figure link yet, as
(row_id, PropertyRow stub) pairs the linker can work on."""
cols = ", ".join(_LINK_BACKFILL_COLS)
cur = conn.execute(
f"SELECT {cols} FROM {table} WHERE source_sha1 = ? "
f"AND IFNULL(origin,'text') = 'text' AND figure_id IS NULL",
(sha1,),
)
return [_stub_row(dict(zip(_LINK_BACKFILL_COLS, rec))) for rec in cur.fetchall()]
def _stub_row(d: dict) -> tuple[int, PropertyRow]:
row = PropertyRow(
material_name=d.get("material_name") or "", material_abbreviation=d.get("material_abbreviation") or "",
material_key=d.get("material_key") or "", material_class=d.get("material_class") or "",
section=d.get("section") or "", property_name=d.get("property_name") or "",
value=d.get("value_raw") or "", unit=d.get("unit") or "", english="",
test_condition=d.get("test_condition") or "", comments=d.get("comments") or "",
value_raw=d.get("value_raw") or "", page=d.get("page"), source_quote=d.get("source_quote") or "",
)
return int(d["id"]), row
def update_row_link(conn: sqlite3.Connection, table: str, row_id: int, row: PropertyRow) -> None:
"""Write a backfilled link onto an existing row (never touches value/status)."""
conn.execute(
f"UPDATE {table} SET figure_id = ?, figure_ref = ?, figure_link_score = ?, "
f"figure_link_signals = ?, image_url = COALESCE(?, image_url), image = COALESCE(?, image) "
f"WHERE id = ?",
(row.figure_id or None, row.figure_ref or None, row.figure_link_score,
row.figure_link_signals or None, row.image_url or None, row.image or None, row_id),
)
def store_linked_figure(conn: sqlite3.Connection, fig: Any) -> None:
"""Figure-linking phase: make sure a figure some text row links to has a
`figures` row and its PNG bytes, without disturbing the kind/status a
--figures run may already have written for it."""
exists = conn.execute("SELECT 1 FROM figures WHERE figure_id = ?", (fig.figure_id,)).fetchone()
if exists:
conn.execute("UPDATE figures SET png_bytes = COALESCE(png_bytes, ?) WHERE figure_id = ?",
(fig.png_bytes or None, fig.figure_id))
else:
upsert_figure(conn, fig, png_bytes=fig.png_bytes or None)
def figures_recorded_for(conn: sqlite3.Connection, sha1: str) -> int:
"""How many figures the `figures` table already holds for this PDF."""
cur = conn.execute("SELECT count(*) FROM figures WHERE source_sha1 = ?", (sha1,))
return int(cur.fetchone()[0])
# Terminal figure states: nothing more to spend on these. Everything else
# (classify_failed / mining_failed / not_mined) is retried on the next --figures
# run — but only those, so an outage costs exactly the failed calls.
_FIGURE_DONE_STATES = ("mined", "skipped_kind")
def figures_done_for(conn: sqlite3.Connection, sha1: str, mine: bool) -> set[str]:
"""figure_ids of this PDF that need no further vision calls."""
states = _FIGURE_DONE_STATES if mine else _FIGURE_DONE_STATES + ("not_mined",)
q = ", ".join("?" for _ in states)
cur = conn.execute(
f"SELECT figure_id FROM figures WHERE source_sha1 = ? AND mining_status IN ({q})",
(sha1, *states),
)
return {r[0] for r in cur.fetchall()}
def source_status(conn: sqlite3.Connection, sha1: str) -> Optional[str]:
"""material_class recorded in `sources` for this PDF ('scanned_no_text' marks
an image-only PDF the text pass refused)."""
cur = conn.execute("SELECT material_class FROM sources WHERE pdf_sha1 = ? LIMIT 1", (sha1,))
r = cur.fetchone()
return r[0] if r else None
def figures_pending_for(conn: sqlite3.Connection, sha1: str, mine: bool) -> int:
"""Figures recorded for this PDF that are NOT done (a rerun should retry them)."""
return figures_recorded_for(conn, sha1) - len(figures_done_for(conn, sha1, mine))
def materials_for_source(conn: sqlite3.Connection, sha1: str) -> list["extraction.Material"]:
"""Rebuild the text-pass material list of an already-ingested PDF from its
rows, so the figure stage can run on a PDF whose text pass happened in an
earlier run (backfill mode)."""
seen: dict[str, extraction.Material] = {}
for table in ALL_TABLES:
cur = conn.execute(
f"SELECT DISTINCT material_name, material_abbreviation, material_class, "
f"trade_grade, manufacturer, matrix, fiber, fiber_volume_fraction "
f"FROM {table} WHERE source_sha1 = ? AND IFNULL(origin,'text') = 'text'",
(sha1,),
)
for r in cur.fetchall():
m = extraction.Material(
material_name=r[0] or "", material_abbreviation=r[1] or "",
material_class=r[2] or "", trade_grade=r[3] or "",
manufacturer=r[4] or "", matrix=r[5] or "", fiber=r[6] or "",
fiber_volume_fraction=r[7] or "",
)
seen.setdefault(extraction.material_key(m), m)
return list(seen.values())
def insert_row(conn: sqlite3.Connection, table: str, row: PropertyRow) -> None:
placeholders = ", ".join("?" for _ in _INSERT_COLS)
cols = ", ".join(_INSERT_COLS)
conn.execute(
f"INSERT INTO {table} ({cols}) VALUES ({placeholders})",
_row_values(row),
)
def record_source(
conn: sqlite3.Connection,
pdf_path: Path,
sha1: str,
material_class: Optional[str],
abbr: Optional[str],
) -> None:
conn.execute(
"INSERT OR IGNORE INTO sources "
"(pdf_filename, pdf_sha1, ingested_at, material_class, material_abbreviation) "
"VALUES (?, ?, datetime('now'), ?, ?)",
(pdf_path.name, sha1, material_class, abbr),
)
# ---------------------------------------------------------------------------
# Driver
# ---------------------------------------------------------------------------
def _empty_result(pdf_path: Path, started: float, error: str) -> PdfResult:
return PdfResult(
pdf=pdf_path.name,
elapsed_s=time.time() - started,
materials=0,
extracted=0,
inserted=0,
flagged=0,
duplicates=0,
material_classes=[],
error=error,
)
@dataclasses.dataclass
class FigureOptions:
"""--figures settings handed to process_pdf (None = figure stage off)."""
out_dir: Path = Path("crawl_out/figures")
max_figures: int = 12
mine: bool = True # False = harvest + classify only (cheap mode)
def _run_figure_stage(
pdf_path: Path, pdf_bytes: bytes, sha1: str, text_materials: list,
api_key: str, conn: Any, db: Any, opts: FigureOptions, result: PdfResult,
) -> None:
"""harvest -> classify -> mine -> insert figure rows. NEVER raises: any
failure lands in result.figure_error and is counted; the PDF's text rows
are already committed by the time this runs."""
try:
import figures as F
done = db.figures_done_for(conn, sha1, opts.mine) if hasattr(db, "figures_done_for") else set()
stage = F.run_figure_stage(
pdf_bytes, pdf_path.name, sha1, text_materials, api_key,
out_dir=opts.out_dir, max_figures=opts.max_figures, mine=opts.mine,
done_figure_ids=done,
)
result.figures_found = len(stage.figures)
result.figures_mined = stage.mined_figures
result.vision_calls = stage.vision.total
result.tokens_in += getattr(stage.vision, "tokens_in", 0)
result.tokens_out += getattr(stage.vision, "tokens_out", 0)
result.figure_filters = {
k: v for k, v in dataclasses.asdict(stage.harvest).items() if v
}
if stage.error:
result.figure_error = stage.error
for fig in stage.figures:
if fig.mining_status == "already_done":
continue # keep the stored kind/status from the earlier run
db.upsert_figure(conn, fig)
frows = fdups = 0
for row in stage.rows:
table = TABLE_FOR_CLASS.get(row.material_class, "Polymers")
if db.already_inserted(conn, table, row):
fdups += 1
continue
db.insert_row(conn, table, row)
frows += 1
result.figure_rows = frows
result.figure_duplicates = fdups
# The top-level rows_* metrics stay TEXT-only (insert_rate = inserted /
# extracted must stay <= 1); figure counts live under result.figure_*
# and run_report.json["figures"].
conn.commit()
if stage.vision.failed_calls and not result.figure_error:
result.figure_error = f"vision_calls_failed:{stage.vision.failed_calls}"
elif stage.vision.incomplete_calls and not result.figure_error:
result.figure_error = f"classify_incomplete:{stage.vision.incomplete_calls}"
except Exception as exc: # pragma: no cover - defensive; run_figure_stage already guards
log = logging.getLogger("batch_ingest")
log.exception("figure stage crashed for %s", pdf_path.name)
result.figure_error = f"figure_stage_error:{type(exc).__name__}:{extraction.redact_secrets(exc)}"
try:
conn.rollback()
except Exception:
pass
@dataclasses.dataclass
class LinkOptions:
"""--link-figures settings handed to process_pdf (None = linking off).
Links TEXT rows to the harvested figure their evidence cites (figure_links.py,
ported from the InDeS mapper). The harvest is local PyMuPDF work; the link
pass makes no Gemini call and never changes a row's value or status.
"""
out_dir: Path = Path("crawl_out/figures")
max_figures: int = 40 # harvest is local and free: cap higher than --figures' 12 so a
# figure-heavy review keeps its later figures linkable
fallback: bool = False # mapper5 page+token fallback links (default: citations only)
embed_images: bool = False # copy the PNG into each linked row's `image` column
upload_s3: bool = True # when S3_BUCKET is set: upload the PNG, fill image_url
def _run_link_stage(
pdf_path: Path, pdf_bytes: bytes, sha1: str, rows: list[PropertyRow],
page_texts: Optional[list[str]], conn: Any, db: Any, opts: LinkOptions,
) -> dict[str, Any]:
"""harvest (local) -> link citations -> image refs -> figure provenance.
NEVER raises: a failure leaves the rows unlinked and lands in link_error.
Returns the PdfResult fields to set."""
info: dict[str, Any] = {"rows_linked": 0, "link_figures_found": 0, "link_stats": {}, "link_error": None}
try:
import figures as F
import figure_links as L
figs = F.harvest_figures(pdf_bytes, pdf_path.name, sha1, opts.out_dir, opts.max_figures,
stats=F.HarvestStats())
info["link_figures_found"] = len(figs)
st = L.link_rows_to_figures(rows, figs, page_texts, allow_fallback=opts.fallback)
L.attach_images(rows, figs, embed=opts.embed_images, upload_s3=opts.upload_s3)
if st.linked_figure_ids and hasattr(db, "store_linked_figure"):
by_id = {f.figure_id: f for f in figs}
for fid in st.linked_figure_ids:
db.store_linked_figure(conn, by_id[fid])
info["rows_linked"] = st.rows_linked
info["link_stats"] = st.as_dict()
except Exception as exc:
logging.getLogger("batch_ingest").exception("link stage failed for %s", pdf_path.name)
info["link_error"] = f"link_stage_error:{type(exc).__name__}:{extraction.redact_secrets(exc)}"
for row in rows: # never insert a half-written link
row.figure_id = row.figure_id if (row.origin or "text") == "figure" else ""
row.figure_ref = ""; row.figure_link_score = None
row.figure_link_signals = ""; row.image_url = ""; row.image = b""
try:
# Pending store_linked_figure writes must not ride the caller's
# next commit as partial figure provenance.
conn.rollback()
except Exception:
pass
return info
def _backfill_links(
pdf_path: Path, pdf_bytes: bytes, sha1: str, conn: Any, db: Any,
opts: LinkOptions, result: PdfResult,
) -> None:
"""Link pass for an already-ingested PDF: rows that carry no figure link
are read back, linked, and updated in place (value/status untouched)."""
if not hasattr(db, "rows_for_link_backfill"):
return
pending: list[tuple[str, int, PropertyRow]] = []
for table in ALL_TABLES:
for row_id, row in db.rows_for_link_backfill(conn, table, sha1):
pending.append((table, row_id, row))
if not pending:
return
rows = [r for _, _, r in pending]
page_texts = extraction.pdf_page_texts(pdf_bytes)
info = _run_link_stage(pdf_path, pdf_bytes, sha1, rows, page_texts, conn, db, opts)
if not info["link_error"]:
for table, row_id, row in pending:
if row.figure_id:
db.update_row_link(conn, table, row_id, row)
conn.commit()
for k, v in info.items():
setattr(result, k, v)
def process_pdf(
pdf_path: Path,
conn: Any,
api_key: str,
db: Any = None,
figure_opts: Optional[FigureOptions] = None,
link_opts: Optional[LinkOptions] = None,
) -> PdfResult:
# `db` = backend module providing seen_sha1 / record_source /
# already_inserted / insert_row. Defaults to this module (SQLite);
# main() passes pg_mirror for --pg. Same logic either way.
db = db or sys.modules[__name__]
started = time.time()
pdf_bytes = pdf_path.read_bytes()
sha1 = hashlib.sha1(pdf_bytes).hexdigest()
# Skip if we already ingested this exact file — unless figures are on and
# this PDF has no figures recorded yet (backfill for PDFs whose text pass
# predates --figures) or has figures still pending (classify/mining
# failed or --no-figure-mining last time): then run ONLY the figure stage,
# and only for the pending figures. Likewise --link-figures backfills
# links onto rows written before linking existed.
if db.seen_sha1(conn, sha1):
result = _empty_result(pdf_path, started, "skipped_seen_sha1")
# Whole-page scans are not figures: a PDF the text pass refused as
# scanned_no_text stays skipped on reruns too (it used to be backfilled
# — spending vision calls on page scans with zero text context).
scanned = hasattr(db, "source_status") and db.source_status(conn, sha1) == "scanned_no_text"
if figure_opts is not None and not scanned and hasattr(db, "figures_recorded_for") and (
db.figures_recorded_for(conn, sha1) == 0
or db.figures_pending_for(conn, sha1, figure_opts.mine) > 0):
mats = db.materials_for_source(conn, sha1)
_run_figure_stage(pdf_path, pdf_bytes, sha1, mats, api_key, conn, db,
figure_opts, result)
result.elapsed_s = time.time() - started
if link_opts is not None and not scanned:
_backfill_links(pdf_path, pdf_bytes, sha1, conn, db, link_opts, result)
result.elapsed_s = time.time() - started
return result
try:
extracted = extract_from_pdf(pdf_bytes, pdf_path.name, api_key)
except requests.RequestException as exc:
# requests quotes the keyed URL in HTTPError/ConnectionError messages.
err = _empty_result(pdf_path, started,
f"gemini_error:{extraction.redact_secrets(exc)}"[:400])
err.quota = bool(getattr(exc, "quota", False))
return err
tokens_in = int(getattr(extracted, "tokens_in", 0) or 0)
tokens_out = int(getattr(extracted, "tokens_out", 0) or 0)
if extracted.doc_status == "scanned_no_text":
# Don't fabricate rows from an image-only PDF (Task 10). Whole-page
# scans are not figures either — the figure stage is skipped too.
db.record_source(conn, pdf_path, sha1, "scanned_no_text", None)
conn.commit()
return _empty_result(pdf_path, started, "scanned_no_text")
if extracted.doc_status != "ok" or not extracted.materials:
# Billed even though nothing came back: keep the usage on the result.
empty = _empty_result(pdf_path, started, "empty_extraction")
empty.tokens_in, empty.tokens_out = tokens_in, tokens_out
return empty
# Ground every value against the PDF text (Task 1).
page_texts = extraction.pdf_page_texts(pdf_bytes)
verify_against_text(extracted, page_texts)
rows = to_rows(extracted, pdf_path.name, sha1)
# Figure-linking phase: attach the cited figure to each text row BEFORE the
# insert, so the row is written once with its link (figure_id, figure_ref,
# score, signals, image_url/image). Local harvest only, no API call.
link_info: dict[str, Any] = {}
if link_opts is not None:
link_info = _run_link_stage(pdf_path, pdf_bytes, sha1, rows, page_texts, conn, db, link_opts)
classes = [extraction.classify_material(m) for m in extracted.materials]
primary_class = classes[0] if classes else None
primary_abbr = rows[0].material_abbreviation if rows else None
db.record_source(conn, pdf_path, sha1, primary_class, primary_abbr)
inserted = flagged = duplicates = linked = 0
for row in rows:
table = TABLE_FOR_CLASS.get(row.material_class, "Polymers")
if db.already_inserted(conn, table, row):
duplicates += 1
continue
db.insert_row(conn, table, row)
inserted += 1
if row.status != "ok":
flagged += 1
if row.figure_id:
linked += 1
conn.commit() # text rows are safe on disk before the figure stage runs
if link_info:
link_info["rows_linked"] = linked # links that actually reached the DB
result = PdfResult(
pdf=pdf_path.name,
elapsed_s=time.time() - started,
materials=len(extracted.materials),
extracted=len(rows),
inserted=inserted,
flagged=flagged,
duplicates=duplicates,
material_classes=sorted(set(classes)),
tokens_in=tokens_in,
tokens_out=tokens_out,
)
for k, v in link_info.items():
setattr(result, k, v)
if figure_opts is not None:
_run_figure_stage(pdf_path, pdf_bytes, sha1, extracted.materials, api_key,
conn, db, figure_opts, result)
result.elapsed_s = time.time() - started
return result
# ---------------------------------------------------------------------------
# review_queue.csv as a view over flagged rows (Task 4)
# ---------------------------------------------------------------------------
_REVIEW_COLUMNS = [
"table_name", "source_pdf", "page", "status", "flag_reason",
"material_name", "material_key", "material_class", "section",
"property_name", "value_raw", "value_num", "unit", "unit_canonical",
"value_si", "test_condition", "source_quote", "comments",
# figure-mining phase: figure rows land here automatically (status is
# never 'ok'); a reviewer opens the PNG behind figure_id and --promote is
# how one gets blessed.
"origin", "figure_id",
]
def export_review_queue(conn: sqlite3.Connection, path: Path) -> int:
"""Write review_queue.csv as SELECT ... WHERE status != 'ok' across tables."""
import csv
select_cols = [c for c in _REVIEW_COLUMNS if c != "table_name"]
rows: list[list[Any]] = []
for table in ALL_TABLES:
cur = conn.execute(
f"SELECT {', '.join(select_cols)} FROM {table} "
f"WHERE IFNULL(status,'ok') != 'ok'"
)
for r in cur.fetchall():
rows.append([table, *r])
with path.open("w", newline="", encoding="utf-8") as fh:
writer = csv.writer(fh)
writer.writerow(_REVIEW_COLUMNS)
writer.writerows(rows)
return len(rows)
def promote_review_queue(conn: sqlite3.Connection, path: Path) -> int:
"""Re-admit corrected rows from a review CSV as status='ok' (Task 4, optional).
Matches on the full dedup grain — (table_name, source_pdf, material_key,
section, property_name, test_condition, value_raw, origin) — and only
touches rows whose status is not already 'ok', so promoting one flagged
row cannot rewrite the flag_reason of an already-ok sibling in another
section, and promoting a text row cannot silently bless the figure row
that reports the same number (or vice versa). A CSV without an `origin`
column (pre-figure-phase export) matches text rows only.
Sets status='ok', flag_reason='promoted'.
"""
import csv
promoted = 0
with path.open("r", newline="", encoding="utf-8") as fh:
for r in csv.DictReader(fh):
table = r.get("table_name")
if table not in ALL_TABLES:
continue
cur = conn.execute(
f"UPDATE {table} SET status='ok', flag_reason='promoted' "
f"WHERE IFNULL(source_pdf,'')=? AND IFNULL(material_key,'')=? "
f" AND IFNULL(section,'')=? "
f" AND IFNULL(property_name,'')=? AND IFNULL(test_condition,'')=? "
f" AND IFNULL(value_raw,'')=? "
f" AND IFNULL(origin,'text')=? "
f" AND IFNULL(status,'ok') != 'ok'",
(r.get("source_pdf") or "", r.get("material_key") or "",
r.get("section") or "",
r.get("property_name") or "", r.get("test_condition") or "",
r.get("value_raw") or "",
(r.get("origin") or "text").strip() or "text"),
)
promoted += cur.rowcount
conn.commit()
return promoted
# ---------------------------------------------------------------------------
# Reporting
# ---------------------------------------------------------------------------
def summarize(results: list[PdfResult]) -> dict[str, Any]:
total_pdfs = len(results)
successes = [r for r in results if not r.error]
extracted = sum(r.extracted for r in successes)
inserted = sum(r.inserted for r in successes)
flagged = sum(r.flagged for r in successes)
duplicates = sum(r.duplicates for r in successes)
materials = sum(r.materials for r in successes)
total_elapsed = sum(r.elapsed_s for r in results)
avg_elapsed = total_elapsed / total_pdfs if total_pdfs else 0.0
error_breakdown: dict[str, int] = {}
for r in results:
if r.error:
key = r.error.split(":", 1)[0]
error_breakdown[key] = error_breakdown.get(key, 0) + 1
# figure-mining phase (all PdfResults, incl. backfill on seen PDFs)
fig_filters: dict[str, int] = {}
fig_errors: dict[str, int] = {}
for r in results:
for k, v in (r.figure_filters or {}).items():
fig_filters[k] = fig_filters.get(k, 0) + v
if r.figure_error:
key = r.figure_error.split(":", 1)[0]
fig_errors[key] = fig_errors.get(key, 0) + 1
figure_rows = sum(r.figure_rows for r in results)
# figure-linking phase
link_stats: dict[str, int] = {}
link_errors: dict[str, int] = {}
for r in results:
for k, v in (r.link_stats or {}).items():
link_stats[k] = link_stats.get(k, 0) + int(v)
if r.link_error:
key = r.link_error.split(":", 1)[0]
link_errors[key] = link_errors.get(key, 0) + 1
rows_linked = sum(r.rows_linked for r in results)
return {
"pdfs_seen": total_pdfs,
"pdfs_ok": len(successes),
"errors_by_kind": error_breakdown,
"materials_extracted": materials,
"rows_extracted": extracted,
"rows_inserted": inserted,
"rows_flagged_in_db": flagged,
"rows_duplicate_skipped": duplicates,
"insert_rate": inserted / extracted if extracted else 0.0,
"flag_rate": flagged / inserted if inserted else 0.0,
"duplicate_rate": duplicates / extracted if extracted else 0.0,
"avg_seconds_per_pdf": round(avg_elapsed, 2),
"total_seconds": round(total_elapsed, 2),
"tokens": {
"tokens_in": sum(r.tokens_in for r in results),
"tokens_out": sum(r.tokens_out for r in results),
},
"figures": {
"figures_found": sum(r.figures_found for r in results),
"figures_mined": sum(r.figures_mined for r in results),
"figure_rows": figure_rows,
"figure_rows_duplicate_skipped": sum(r.figure_duplicates for r in results),
"vision_calls": sum(r.vision_calls for r in results),
"figure_errors_by_kind": fig_errors,
"harvest_filters": fig_filters,
},
"figure_links": {
"rows_linked": rows_linked,
# share of newly inserted text rows that carry a link (backfilled
# links on seen PDFs are counted in rows_linked but not here)
"link_rate": (sum(r.rows_linked for r in successes) / inserted) if inserted else 0.0,
"figures_harvested": sum(r.link_figures_found for r in results),
"link_errors_by_kind": link_errors,
**{k: v for k, v in link_stats.items() if k != "rows_linked"},
},
}
def main() -> int:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--input", type=Path,
help="Folder of PDFs to ingest")
parser.add_argument("--db", default=Path("materials_mirror.sqlite"),
type=Path, help="SQLite mirror path")
parser.add_argument("--review", default=Path("review_queue.csv"),
type=Path, help="CSV view of flagged (status != 'ok') rows")
parser.add_argument("--report", default=Path("run_report.json"),
type=Path, help="JSON run summary")
parser.add_argument("--limit", type=int, default=None,
help="Process at most N PDFs (for testing)")
parser.add_argument("--migrate", action="store_true",
help="Add hardening-phase columns to an existing DB and exit")
parser.add_argument("--promote", type=Path, default=None,
help="Re-admit corrected rows from a review CSV as status='ok'")
parser.add_argument("--pg", action="store_true",
help="Write to the shared Postgres (env DB_HOST/... or "
"DATABASE_URL) instead of the local SQLite mirror")
# --- figure-mining phase (opt-in) ---
parser.add_argument("--figures", action="store_true",
help="Also harvest figures from each PDF, classify them with one "
"vision call per PDF, and mine plots/table-images for "
"property values (origin='figure', status='figure_estimate')")
parser.add_argument("--figures-dir", type=Path, default=Path("crawl_out/figures"),
help="Where harvested figure PNGs go (<dir>/<sha1>/p<page>_<n>.png)")
parser.add_argument("--max-figures-per-pdf", type=int, default=12,
help="Hard cap on harvested figures per PDF (bounds vision calls)")
parser.add_argument("--no-figure-mining", action="store_true",
help="With --figures: harvest + classify only, skip the "
"per-figure mining calls (cheap mode)")
# --- figure-linking phase (opt-in; local, no API calls) ---
parser.add_argument("--link-figures", action="store_true",
help="Link each text row to the harvested figure its evidence "
"cites ('see Fig. 3'): figure_id/figure_ref/score/signals on "
"the row. Also backfills links onto rows of already-ingested "
"PDFs. With --pg, requires the figure migration "
"(pg_migrate.py --apply: figures table + v2 dedup index).")
parser.add_argument("--figure-link-fallback", action="store_true",
help="With --link-figures: also accept the InDeS mapper's "
"page+caption-token fallback links (score < 0.9, lower "
"precision). Default: explicit citations only.")
parser.add_argument("--embed-figure-images", action="store_true",
help="With --link-figures: copy the linked figure's PNG into the "
"row's `image` column so the shared Space shows it without S3")
parser.add_argument("--no-s3", action="store_true",
help="With --link-figures: never upload PNGs even if S3_BUCKET is set")
parser.add_argument("--max-link-figures", type=int, default=40,
help="With --link-figures: harvest cap for the link pass (no API cost; "
"captioned figures are kept first when it bites)")
parser.add_argument("--links-csv", type=Path, default=None,
help="After the run, write every linked row + its figure caption/PNG "
"path to this CSV (labelling sheet for figure-linkage precision)")
args = parser.parse_args()
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s %(message)s",
stream=sys.stderr,
)
log = logging.getLogger("batch_ingest")
# Select the storage backend: this module (SQLite, default) or pg_mirror.
if args.pg:
import pg_mirror as db
log.info("Postgres mode: %s", db.config_summary())
else:
db = sys.modules[__name__]
figure_opts: Optional[FigureOptions] = None
if args.figures:
figure_opts = FigureOptions(out_dir=args.figures_dir,
max_figures=args.max_figures_per_pdf,
mine=not args.no_figure_mining)
link_opts: Optional[LinkOptions] = None
if args.link_figures:
link_opts = LinkOptions(out_dir=args.figures_dir,
max_figures=args.max_link_figures,
fallback=args.figure_link_fallback,
embed_images=args.embed_figure_images,
upload_s3=not args.no_s3)
if args.migrate:
if args.pg:
log.error("Schema changes to the shared Postgres are deliberately "
"kept in one place: run `python pg_migrate.py` (dry-run) "
"then `python pg_migrate.py --apply`.")
return 2
run_migrate(args.db)
log.info("Migration complete: %s", args.db)
return 0
if args.promote:
conn = db.connect_from_env() if args.pg else init_db(args.db)
if args.pg:
db.check_schema(conn)
n = db.promote_review_queue(conn, args.promote)
db.export_review_queue(conn, args.review)
conn.close()
log.info("Promoted %d rows to status='ok' from %s", n, args.promote)
return 0
if not args.input:
log.error("--input is required (folder of PDFs).")
return 2
api_key = os.environ.get("GEMINI_API_KEY") or os.environ.get("GOOGLE_API_KEY")
if not api_key:
log.error("GEMINI_API_KEY (or GOOGLE_API_KEY) is not set.")
return 2
pdfs = sorted(p for p in args.input.rglob("*.pdf"))
if args.limit:
pdfs = pdfs[: args.limit]
if not pdfs:
log.error("No PDFs found under %s", args.input)
return 2
if args.pg:
log.info("Ingesting %d PDFs into Postgres (%s)", len(pdfs), db.config_summary())
conn = db.connect_from_env()
db.check_schema(conn) # refuse to run against an unmigrated schema
if figure_opts is not None or link_opts is not None:
# Not half-supported silently: figure work needs the figures table
# and the origin-aware v2 dedup index (the v1 index would reject a
# figure row matching a text row on the origin-less grain).
problems = db.figures_ready_problems(conn)
if problems:
log.error("--figures/--link-figures with --pg needs the figure "
"migration first — run `python pg_migrate.py` (dry-run) "
"then `python pg_migrate.py --apply`. Problems: %s",
"; ".join(problems))
conn.close()
return 2
else:
log.info("Ingesting %d PDFs into %s", len(pdfs), args.db)
conn = init_db(args.db)
results: list[PdfResult] = []
for i, pdf in enumerate(pdfs, start=1):
log.info("[%d/%d] %s", i, len(pdfs), pdf.name)
result = process_pdf(pdf, conn, api_key, db=db, figure_opts=figure_opts,
link_opts=link_opts)
results.append(result)
log.info(
" -> materials=%d classes=%s extracted=%d inserted=%d "
"flagged=%d duplicates=%d elapsed=%.1fs error=%s",
result.materials,
",".join(result.material_classes) or "-",
result.extracted,
result.inserted,
result.flagged,
result.duplicates,
result.elapsed_s,
result.error,
)
if figure_opts is not None:
log.info(
" -> figures: found=%d mined=%d rows=%d dup=%d vision_calls=%d error=%s",
result.figures_found, result.figures_mined, result.figure_rows,
result.figure_duplicates, result.vision_calls, result.figure_error,
)
if link_opts is not None:
log.info(
" -> figure links: rows_linked=%d figures=%d %s error=%s",
result.rows_linked, result.link_figures_found,
" ".join(f"{k}={v}" for k, v in (result.link_stats or {}).items()
if k in ("via_citation", "via_near_quote", "via_fallback", "cited_but_missing")),
result.link_error,
)
n_flagged = db.export_review_queue(conn, args.review)
summary = summarize(results)
args.report.write_text(json.dumps(summary, indent=2))
log.info("Run summary written to %s", args.report)
log.info("Review queue (%d flagged rows) written to %s", n_flagged, args.review)
if args.links_csv is not None:
import figure_links as L
# A migrated Postgres has the figures table too; only skip the caption
# join on a pg DB from before the figure migration.
n_links = L.export_figure_links(conn, args.links_csv,
with_figures_table=not args.pg or db.figures_table_exists(conn))
log.info("Figure-link labelling sheet (%d linked rows) written to %s", n_links, args.links_csv)
conn.close()
print(json.dumps(summary, indent=2))
return 0
if __name__ == "__main__":
sys.exit(main())