Spaces:
Running
Running
File size: 51,889 Bytes
11dde75 e5bfacd de5014e e5bfacd 11dde75 e5bfacd 11dde75 de5014e 11dde75 de5014e 11dde75 de5014e 11dde75 de5014e 00f9c01 de5014e 5ea82ce 2cf2ccd e5bfacd 11dde75 de5014e 11dde75 e5bfacd 11dde75 de5014e 11dde75 de5014e 11dde75 de5014e e5bfacd 11dde75 de5014e 11dde75 de5014e 11dde75 de5014e 11dde75 de5014e e5bfacd de5014e e5bfacd de5014e e5bfacd de5014e e5bfacd de5014e e5bfacd de5014e e5bfacd de5014e e5bfacd de5014e 11dde75 de5014e 5ea82ce 2cf2ccd de5014e 00f9c01 de5014e e5bfacd 00f9c01 e5bfacd 11dde75 de5014e e5bfacd 2cf2ccd 11dde75 2cf2ccd 11dde75 de5014e e5bfacd 11dde75 de5014e e5bfacd de5014e 11dde75 2cf2ccd 11dde75 00f9c01 5ea82ce 2cf2ccd 11dde75 de5014e 11dde75 5ea82ce 2cf2ccd 5ea82ce 11dde75 e5bfacd 11dde75 e5bfacd 11dde75 e5bfacd de5014e e5bfacd 11dde75 de5014e 11dde75 5ea82ce 2cf2ccd 11dde75 e5bfacd 11dde75 de5014e 11dde75 de5014e 11dde75 de5014e 11dde75 de5014e 11dde75 de5014e 11dde75 de5014e e5bfacd 11dde75 5ea82ce 2cf2ccd 5ea82ce de5014e e5bfacd 11dde75 de5014e e5bfacd 11dde75 de5014e e5bfacd 11dde75 e5bfacd 11dde75 e5bfacd 11dde75 e5bfacd 11dde75 de5014e e5bfacd 11dde75 e5bfacd 11dde75 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708 709 710 711 712 713 714 715 716 717 718 719 720 721 722 723 724 725 726 727 728 729 730 731 732 733 734 735 736 737 738 739 740 741 742 743 744 745 746 747 748 749 750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788 789 790 791 792 793 794 795 796 797 798 799 800 801 802 803 804 805 806 807 808 809 810 811 812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 845 846 847 848 849 850 851 852 853 854 855 856 857 858 859 860 861 862 863 864 865 866 867 868 869 870 871 872 873 874 875 876 877 878 879 880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 900 901 902 903 904 905 906 907 908 909 910 911 912 913 914 915 916 917 918 919 920 921 922 923 924 925 926 927 928 929 930 931 932 933 934 935 936 937 938 939 940 941 942 943 944 945 946 947 948 949 950 951 952 953 954 955 956 957 958 959 960 961 962 963 964 965 966 967 968 969 970 971 972 973 974 975 976 977 978 979 980 981 982 983 984 985 986 987 988 989 990 991 992 993 994 995 996 997 998 999 1000 1001 1002 1003 1004 1005 1006 1007 1008 1009 1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022 1023 1024 1025 1026 1027 1028 1029 1030 1031 1032 1033 1034 1035 1036 1037 1038 1039 1040 1041 1042 1043 1044 1045 1046 1047 1048 1049 1050 1051 1052 1053 1054 1055 1056 1057 1058 1059 1060 1061 1062 1063 1064 1065 1066 1067 1068 1069 1070 1071 1072 1073 1074 1075 1076 1077 1078 1079 1080 1081 1082 1083 1084 1085 1086 1087 1088 1089 1090 1091 1092 1093 1094 1095 1096 1097 1098 1099 1100 1101 1102 1103 1104 1105 1106 1107 1108 1109 1110 1111 1112 1113 1114 1115 1116 1117 1118 1119 1120 1121 1122 1123 1124 1125 1126 1127 1128 1129 1130 1131 1132 1133 1134 1135 1136 1137 1138 1139 1140 | 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
# Breakdown for cost analysis (all included in tokens_in / tokens_out):
# thinking tokens of every call, and the vision calls' share.
tokens_thinking: int = 0
vision_tokens_in: int = 0
vision_tokens_out: int = 0
text_model: str = "" # model of the text call ("" = no call made)
vision_model: str = "" # model of the vision calls ("" = none made)
# near-duplicate gate: sha1 of the already-ingested document this one
# repeats (then nothing was sent to Gemini)
duplicate_of: Optional[str] = None
# 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.tokens_thinking += getattr(stage.vision, "tokens_thinking", 0)
result.vision_tokens_in += getattr(stage.vision, "tokens_in", 0)
result.vision_tokens_out += getattr(stage.vision, "tokens_out", 0)
if stage.vision.total:
result.vision_model = getattr(stage.vision, "model", "") or extraction.vision_model()
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,
*,
pre_extracted: Optional[Extraction] = 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.
#
# pre_extracted: the text call's result obtained elsewhere (the Gemini
# Batch API); no text call is made here then.
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 = (pre_extracted if pre_extracted is not None
else 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)
tokens_thinking = int(getattr(extracted, "tokens_thinking", 0) or 0)
text_model_used = (getattr(extracted, "model", "") or "") if (tokens_in or tokens_out) else ""
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
empty.tokens_thinking, empty.text_model = tokens_thinking, text_model_used
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,
tokens_thinking=tokens_thinking,
text_model=text_model_used,
)
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),
"tokens_thinking": sum(r.tokens_thinking for r in results),
"vision_tokens_in": sum(r.vision_tokens_in for r in results),
"vision_tokens_out": sum(r.vision_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())
|