File size: 12,203 Bytes
11dde75
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
de5014e
 
 
 
 
 
e5bfacd
 
 
 
 
 
 
 
 
 
03f160f
 
 
 
 
 
 
 
 
 
 
 
11dde75
 
 
 
de5014e
 
e5bfacd
de5014e
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
e5bfacd
 
de5014e
 
 
 
e5bfacd
 
 
 
 
 
 
 
de5014e
 
e5bfacd
 
de5014e
 
 
 
e5bfacd
 
 
 
de5014e
 
11dde75
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
de5014e
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
11dde75
 
 
 
 
 
 
 
 
 
 
de5014e
 
 
11dde75
de5014e
 
 
 
 
11dde75
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
de5014e
11dde75
de5014e
 
 
 
 
 
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
r"""
migrate.py — non-destructive column migration for the materials mirror.

Adds the extraction-hardening columns to the existing ``Polymers`` / ``Fibers``
/ ``Composites_materials`` tables via ``ALTER TABLE ... ADD COLUMN``. The
original columns (``value``, ``unit``, ``english`` ...) are left untouched and
keep being populated, so the CSV export and the Streamlit ``page1.py`` continue
to work.

SQLite ``ADD COLUMN`` is cheap and never rewrites existing rows; existing rows
get NULL for the new columns. New rows written by ``batch_ingest.py`` fill them.

Usage:
    python migrate.py --db ./materials_mirror.sqlite          # back up + migrate
    python migrate.py --db ./materials_mirror.sqlite --no-backup
    # equivalently:  python batch_ingest.py --migrate --db ...
"""

from __future__ import annotations

import argparse
import shutil
import sqlite3
import sys
from pathlib import Path

# (column_name, sqlite_type) — appended after the legacy columns. Order is the
# canonical insert order used by batch_ingest._INSERT_COLS, so keep it stable.
EXTRA_COLUMNS: list[tuple[str, str]] = [
    ("material_key", "TEXT"),
    ("material_class", "TEXT"),
    ("trade_grade", "TEXT"),
    ("manufacturer", "TEXT"),
    ("matrix", "TEXT"),
    ("fiber", "TEXT"),
    ("fiber_volume_fraction", "TEXT"),
    ("value_raw", "TEXT"),
    ("value_num", "REAL"),
    ("value_min", "REAL"),
    ("value_max", "REAL"),
    ("qualifier", "TEXT"),
    ("unit_canonical", "TEXT"),
    ("value_si", "REAL"),
    ("source_pdf", "TEXT"),
    ("source_sha1", "TEXT"),
    ("page", "INTEGER"),
    ("source_quote", "TEXT"),
    ("status", "TEXT DEFAULT 'ok'"),
    ("flag_reason", "TEXT"),
    ("model", "TEXT"),
    ("prompt_version", "TEXT"),
    # Populated explicitly by batch_ingest at insert time. SQLite forbids a
    # non-constant DEFAULT (datetime('now')) on ALTER TABLE ADD COLUMN, so this
    # stays a plain column rather than carrying a SQL default.
    ("extracted_at", "TEXT"),
    # Figure-mining phase (additive): where the value came from. 'text' =
    # grounded in the PDF text (every pre-existing row reads 'text' via the
    # default); 'figure' = read off a plot/table image by figures.py — an
    # estimate, status 'figure_estimate', never exported unless promoted.
    ("origin", "TEXT DEFAULT 'text'"),
    ("figure_id", "TEXT"),
    # Figure-linking phase (additive): a TEXT row's citation of a harvested
    # figure (figure_links.py). image_url / image are the two columns the
    # shared Space reads to show a row's plot; both already exist on the
    # Postgres category tables (InDeS schema), so only the three link columns
    # are new there.
    ("figure_ref", "TEXT"),
    ("figure_link_score", "REAL"),
    ("figure_link_signals", "TEXT"),
    ("image_url", "TEXT"),
    ("image", "BLOB"),
    # Processing route per property (prompt 2.1, additive): how the specimen
    # behind a value was made. process_type is normalized to
    # extraction.PROCESS_TYPE_ENUM; name / conditions are as printed;
    # process_quote + process_page are the sentence that names the route and
    # process_status says how that sentence grounded in the PDF text. All
    # NULL when the document reports no route (and on every pre-2.1 row).
    ("process_type", "TEXT"),
    ("process_name", "TEXT"),
    ("process_conditions", "TEXT"),
    ("process_quote", "TEXT"),
    ("process_page", "INTEGER"),
    ("process_status", "TEXT"),
]

TARGET_TABLES = ("Polymers", "Fibers", "Composites_materials")

# Provenance of harvested figures (figure-mining phase). One row per PNG on
# disk; property rows point at it through figure_id. Additive. The Postgres
# twin lives in pg_mirror.FIGURES_TABLE_DDL (its bytes column is image_bytes).
FIGURES_DDL = """
CREATE TABLE IF NOT EXISTS figures (
    figure_id TEXT PRIMARY KEY,
    source_pdf TEXT,
    source_sha1 TEXT,
    page INTEGER,
    bbox TEXT,
    caption TEXT,
    figure_kind TEXT,
    material_key TEXT,
    image_path TEXT,
    image_sha256 TEXT,
    width_px INTEGER,
    height_px INTEGER,
    route TEXT,
    mining_status TEXT,
    n_values INTEGER,
    model TEXT,
    figure_prompt_version TEXT,
    extracted_at TEXT,
    png_bytes BLOB
);
CREATE INDEX IF NOT EXISTS ix_figures_source_sha1 ON figures (source_sha1);
"""

# Columns added to `figures` after its first version (additive, like
# EXTRA_COLUMNS for the material tables). png_bytes: the PNG itself, stored
# for figures a text row links to, so the link survives the ephemeral work
# dir on the Space (figure-linking phase).
FIGURE_EXTRA_COLUMNS: list[tuple[str, str]] = [
    ("png_bytes", "BLOB"),
]


def ensure_figures_table(conn: sqlite3.Connection) -> bool:
    """Create the `figures` table + index if missing, and add any missing
    FIGURE_EXTRA_COLUMNS to an existing one. Idempotent; True if created."""
    existed = conn.execute(
        "SELECT 1 FROM sqlite_master WHERE type='table' AND name='figures'"
    ).fetchone() is not None
    conn.executescript(FIGURES_DDL)
    have = _existing_columns(conn, "figures")
    for name, coltype in FIGURE_EXTRA_COLUMNS:
        if name not in have:
            conn.execute(f"ALTER TABLE figures ADD COLUMN {name} {coltype}")
    return not existed


def _existing_columns(conn: sqlite3.Connection, table: str) -> set[str]:
    cur = conn.execute(f"PRAGMA table_info({table})")
    return {row[1] for row in cur.fetchall()}


def ensure_columns(conn: sqlite3.Connection, table: str) -> list[str]:
    """Add any missing EXTRA_COLUMNS to `table`. Returns the columns added.

    Idempotent: safe to call on every run. Only touches a table that exists.
    """
    info = conn.execute(
        "SELECT name FROM sqlite_master WHERE type='table' AND name=?", (table,)
    ).fetchone()
    if not info:
        return []
    have = _existing_columns(conn, table)
    added: list[str] = []
    for name, coltype in EXTRA_COLUMNS:
        if name not in have:
            conn.execute(f"ALTER TABLE {table} ADD COLUMN {name} {coltype}")
            added.append(name)
    return added


def _sources_unique_column(conn: sqlite3.Connection) -> str | None:
    """Which column carries the UNIQUE constraint on `sources` (None if no table)."""
    row = conn.execute(
        "SELECT sql FROM sqlite_master WHERE type='table' AND name='sources'"
    ).fetchone()
    if not row or not row[0]:
        return None
    ddl = row[0].lower()
    if "pdf_sha1 text unique" in ddl:
        return "pdf_sha1"
    if "pdf_filename text unique" in ddl:
        return "pdf_filename"
    return ""


def ensure_sources_sha1_unique(conn: sqlite3.Connection) -> bool:
    """Rebuild `sources` so its UNIQUE key is `pdf_sha1`, not `pdf_filename`.

    Legacy DBs keyed the doc logbook on the basename. Two different PDFs that
    share a filename (likely — --input is rglob'd across vendor folders) then
    collided: `INSERT OR IGNORE` dropped the second, `seen_sha1()` never saw
    it, and it was re-sent to Gemini on every run. Content hash is the real
    identity. SQLite cannot alter a UNIQUE constraint in place, so this is a
    copy-rebuild; it keeps the earliest row per sha1 if duplicates exist.
    Idempotent — returns True only when a rebuild happened.
    """
    key = _sources_unique_column(conn)
    if key in (None, "pdf_sha1"):
        return False
    # One transaction: Python's sqlite3 executescript() autocommits each
    # statement otherwise, and a crash between the CREATE and the RENAME left
    # every later init_db() failing on "sources__new already exists" (or, after
    # the DROP, silently orphaned the whole logbook). DDL is transactional in
    # SQLite, so BEGIN/COMMIT makes the rebuild all-or-nothing; the leading
    # DROP IF EXISTS makes a retry after a crash self-healing.
    conn.executescript(
        """
        BEGIN;
        DROP TABLE IF EXISTS sources__new;
        CREATE TABLE sources__new (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            pdf_filename TEXT,
            pdf_sha1 TEXT UNIQUE,
            ingested_at TEXT,
            material_class TEXT,
            material_abbreviation TEXT
        );
        INSERT INTO sources__new
            (id, pdf_filename, pdf_sha1, ingested_at, material_class, material_abbreviation)
        SELECT id, pdf_filename, pdf_sha1, ingested_at, material_class, material_abbreviation
        FROM sources
        WHERE id IN (SELECT MIN(id) FROM sources GROUP BY IFNULL(pdf_sha1, '__null__' || id));
        DROP TABLE sources;
        ALTER TABLE sources__new RENAME TO sources;
        COMMIT;
        """
    )
    return True


def backfill_material_key_grade(conn: sqlite3.Connection, table: str) -> int:
    """Re-key pipeline rows written before trade_grade became part of material_key.

    Rows ingested before 2026-08 carry ``material_key = <name>`` while their
    ``trade_grade`` column is populated; the current rule is
    ``<name>|<grade>`` (see extraction.material_key). Without this backfill a
    re-ingest of the same PDF would not dedup against the old rows and would
    insert every graded row a second time. Only touches pipeline rows
    (``source_sha1`` set) whose key has no ``|`` yet. Idempotent. Returns the
    number of rows updated. Computed in Python so it is byte-identical to
    what extraction.material_key() produces.
    """
    import re as _re
    info = conn.execute(
        "SELECT name FROM sqlite_master WHERE type='table' AND name=?", (table,)
    ).fetchone()
    if not info:
        return 0
    have = _existing_columns(conn, table)
    if not {"material_key", "trade_grade", "source_sha1"} <= have:
        return 0
    rows = conn.execute(
        f"SELECT id, material_key, trade_grade FROM {table} "
        f"WHERE source_sha1 IS NOT NULL AND IFNULL(trade_grade,'') <> '' "
        f"  AND IFNULL(material_key,'') <> '' AND material_key NOT LIKE '%|%'"
    ).fetchall()
    n = 0
    for rid, key, grade in rows:
        g = _re.sub(r"\s+", " ", (grade or "").strip().lower())
        if g and g != key and g not in key:
            conn.execute(f"UPDATE {table} SET material_key=? WHERE id=?", (f"{key}|{g}", rid))
            n += 1
    return n


def migrate(db_path: Path, backup: bool = True) -> dict[str, list[str]]:
    if backup and db_path.exists():
        bak = db_path.with_suffix(db_path.suffix + ".bak")
        shutil.copy2(db_path, bak)
        print(f"Backed up {db_path} -> {bak}", file=sys.stderr)

    conn = sqlite3.connect(str(db_path))
    try:
        result: dict[str, list[str]] = {}
        for table in TARGET_TABLES:
            added = ensure_columns(conn, table)
            n = backfill_material_key_grade(conn, table)
            if n:
                added = added + [f"(re-keyed {n} rows: material_key += '|trade_grade')"]
            result[table] = added
        if ensure_sources_sha1_unique(conn):
            result["sources"] = ["UNIQUE(pdf_sha1) (rebuilt from UNIQUE(pdf_filename))"]
        else:
            result["sources"] = []
        result["figures"] = ["created"] if ensure_figures_table(conn) else []
        conn.commit()
    finally:
        conn.close()
    return result


def main() -> int:
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument("--db", required=True, type=Path, help="SQLite mirror path")
    parser.add_argument("--no-backup", action="store_true",
                        help="Skip the .bak backup before altering")
    args = parser.parse_args()

    if not args.db.exists():
        print(f"DB {args.db} does not exist; nothing to migrate "
              f"(a fresh DB is created with all columns by batch_ingest).",
              file=sys.stderr)
        return 1

    result = migrate(args.db, backup=not args.no_backup)
    for table, added in result.items():
        if not added:
            print(f"{table}: already up to date")
        elif table == "sources":
            print(f"{table}: rebuilt -> {added[0]}")
        elif table == "figures":
            print(f"{table}: table created")
        else:
            print(f"{table}: added {len(added)} columns: {', '.join(added)}")
    return 0


if __name__ == "__main__":
    sys.exit(main())