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())