Spaces:
Running
Cost "f": price the agent at the model it runs; thinking/vision token breakdown; model switch; near-duplicate gate
Browse filesThe cost column priced gemini-3.5-flash calls at Gemini 2.5 Flash rates
($0.30 / $2.50 per M) while the model is $1.50 / $9.00: it showed about a
quarter of the bill (runs 544-593: $12.98 shown, ~$48.50 real, ~$0.41 per
PDF; AI Studio billed $4.27 on 28 Sep for cycles the Space showed at <$1).
- config: MODEL_PRICES per model (3.5 flash, 3.5 flash-lite, 3 flash preview,
2.5 flash, 2.5 flash-lite), BATCH_DISCOUNT, model_price(), cost_usd(model,
batch), run_cost_usd() pricing text and vision tokens at their own models;
AGENT_PRICE_*_PER_M only as an explicit override of every model.
- extraction: usage_detail() -> (in, out, thinking); Extraction.tokens_thinking;
text_model()/vision_model() read GEMINI_MODEL / GEMINI_VISION_MODEL at call
time (dataclass model defaults were bound at import); extraction_from_response()
split out (accepts dicts, for the Batch API).
- figures: calls go to vision_model(); VisionStats.tokens_thinking + model.
- batch_ingest: PdfResult tokens_thinking, vision_tokens_in/out, text_model,
vision_model, duplicate_of; process_pdf(pre_extracted=...) skips the text
call; summarize() reports the breakdown.
- agent_runs: tokens_thinking, vision_tokens_in/out, duplicate_text, text_model,
vision_model, cost_usd (stored at finish, priced per model).
- Run Control "Gemini cost": text model, vision model (default: same), near-dup
threshold; Logs: Thinking, Near-dup, Model, Cost (stored; older runs re-priced
at gemini-3.5-flash), $/PDF, and a spend line over the listed cycles.
- agent/textdedupe.py: MinHash (128 x word 5-shingles) of the PDF text; a PDF
>= dedupe_text_similarity (default 0.9) to an ingested document is skipped
before any Gemini call (registry 'duplicate_text'); claims are registered at
the gate under a lock (parallel workers) and released when a PDF fails for good.
- test_cost_e2e.py (38 checks).
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018Db9kKLweSLwRrQm63Uot3
- agent/agdb.py +20 -2
- agent/config.py +78 -8
- agent/orchestrator.py +78 -10
- agent/textdedupe.py +182 -0
- batch_ingest.py +30 -1
- extraction.py +55 -11
- figures.py +16 -7
- page_files/Logs.py +38 -7
- page_files/Run_Control.py +47 -5
- test_cost_e2e.py +326 -0
- test_tokens.py +1 -1
|
@@ -157,6 +157,15 @@ _RUNS_EXTRA_COLS = (
|
|
| 157 |
("heartbeat_at", "timestamptz"), # last sign of progress (events)
|
| 158 |
("kind", "text DEFAULT 'cycle'"), # 'cycle' | 'discovery' (28 Sep 2026)
|
| 159 |
("queued", "integer DEFAULT 0"), # discovery pass: candidates queued
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 160 |
)
|
| 161 |
|
| 162 |
# Integer counters persisted by finish_run, in column order. Metric keys equal
|
|
@@ -167,13 +176,16 @@ _RUN_COUNTER_COLS = (
|
|
| 167 |
"figures_found", "figures_mined", "figure_rows", "vision_calls",
|
| 168 |
"relevance_rejected", "url_seen_skipped", "download_failed",
|
| 169 |
"tokens_in", "tokens_out", "rows_linked", "queued",
|
|
|
|
| 170 |
)
|
| 171 |
_METRIC_ALIASES = {"url_seen_skipped": ("url_seen_skipped", "skipped_url")}
|
| 172 |
|
| 173 |
|
| 174 |
def ensure_agent_schema(conn) -> None:
|
|
|
|
| 175 |
with conn.cursor() as cur:
|
| 176 |
cur.execute(DDL)
|
|
|
|
| 177 |
for name, typ in _RUNS_EXTRA_COLS:
|
| 178 |
cur.execute(f"ALTER TABLE agent_runs ADD COLUMN IF NOT EXISTS {name} {typ}")
|
| 179 |
conn.commit()
|
|
@@ -289,14 +301,20 @@ def finish_run(conn, run_id: int, status: str, metrics: dict, report: dict) -> N
|
|
| 289 |
return 0
|
| 290 |
|
| 291 |
sources = metrics.get("sources")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 292 |
with conn.cursor() as cur:
|
| 293 |
cur.execute(
|
| 294 |
f"UPDATE agent_runs SET finished_at = now(), status = %s, {sets}, "
|
| 295 |
-
f"sources = %s::jsonb, report = %s "
|
|
|
|
| 296 |
f"WHERE id = %s AND status = 'running'",
|
| 297 |
(status, *[_val(c) for c in cols],
|
| 298 |
json.dumps(sources) if sources else None,
|
| 299 |
-
redact(json.dumps(report))[:100_000], run_id),
|
| 300 |
)
|
| 301 |
conn.commit()
|
| 302 |
|
|
|
|
| 157 |
("heartbeat_at", "timestamptz"), # last sign of progress (events)
|
| 158 |
("kind", "text DEFAULT 'cycle'"), # 'cycle' | 'discovery' (28 Sep 2026)
|
| 159 |
("queued", "integer DEFAULT 0"), # discovery pass: candidates queued
|
| 160 |
+
# 30 Sep 2026: cost breakdown and per-run pricing (the old derived cost
|
| 161 |
+
# used Gemini 2.5 Flash prices for a gemini-3.5-flash agent).
|
| 162 |
+
("tokens_thinking", "bigint DEFAULT 0"), # thinking share of tokens_out
|
| 163 |
+
("vision_tokens_in", "bigint DEFAULT 0"), # vision share of tokens_in
|
| 164 |
+
("vision_tokens_out", "bigint DEFAULT 0"), # vision share of tokens_out
|
| 165 |
+
("duplicate_text", "integer DEFAULT 0"), # PDFs skipped as near-duplicates
|
| 166 |
+
("text_model", "text"), # Gemini model of the text calls
|
| 167 |
+
("vision_model", "text"), # Gemini model of the vision calls
|
| 168 |
+
("cost_usd", "double precision"), # priced per model at finish
|
| 169 |
)
|
| 170 |
|
| 171 |
# Integer counters persisted by finish_run, in column order. Metric keys equal
|
|
|
|
| 176 |
"figures_found", "figures_mined", "figure_rows", "vision_calls",
|
| 177 |
"relevance_rejected", "url_seen_skipped", "download_failed",
|
| 178 |
"tokens_in", "tokens_out", "rows_linked", "queued",
|
| 179 |
+
"tokens_thinking", "vision_tokens_in", "vision_tokens_out", "duplicate_text",
|
| 180 |
)
|
| 181 |
_METRIC_ALIASES = {"url_seen_skipped": ("url_seen_skipped", "skipped_url")}
|
| 182 |
|
| 183 |
|
| 184 |
def ensure_agent_schema(conn) -> None:
|
| 185 |
+
from agent import textdedupe
|
| 186 |
with conn.cursor() as cur:
|
| 187 |
cur.execute(DDL)
|
| 188 |
+
cur.execute(textdedupe.DDL)
|
| 189 |
for name, typ in _RUNS_EXTRA_COLS:
|
| 190 |
cur.execute(f"ALTER TABLE agent_runs ADD COLUMN IF NOT EXISTS {name} {typ}")
|
| 191 |
conn.commit()
|
|
|
|
| 301 |
return 0
|
| 302 |
|
| 303 |
sources = metrics.get("sources")
|
| 304 |
+
from agent import config as C
|
| 305 |
+
spent = int(metrics.get("tokens_in", 0) or 0) + int(metrics.get("tokens_out", 0) or 0)
|
| 306 |
+
text_model = (metrics.get("text_model") or None) if spent else None
|
| 307 |
+
vision_model = (metrics.get("vision_model") or None) if int(metrics.get("vision_tokens_in", 0) or 0) else None
|
| 308 |
+
cost = round(C.run_cost_usd(metrics), 6) if spent else 0.0
|
| 309 |
with conn.cursor() as cur:
|
| 310 |
cur.execute(
|
| 311 |
f"UPDATE agent_runs SET finished_at = now(), status = %s, {sets}, "
|
| 312 |
+
f"sources = %s::jsonb, report = %s, text_model = %s, vision_model = %s, "
|
| 313 |
+
f"cost_usd = %s "
|
| 314 |
f"WHERE id = %s AND status = 'running'",
|
| 315 |
(status, *[_val(c) for c in cols],
|
| 316 |
json.dumps(sources) if sources else None,
|
| 317 |
+
redact(json.dumps(report))[:100_000], text_model, vision_model, cost, run_id),
|
| 318 |
)
|
| 319 |
conn.commit()
|
| 320 |
|
|
@@ -132,6 +132,18 @@ DEFAULTS = {
|
|
| 132 |
# Thinking tokens are billed as output (8x input); "low" cut a cycle's
|
| 133 |
# bill by roughly two thirds. extraction.THINKING is set from this each cycle.
|
| 134 |
"gemini_thinking": os.environ.get("GEMINI_THINKING", "low").strip().lower() or "dynamic",
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 135 |
"max_figures_per_pdf": _i("AGENT_MAX_FIGURES", 12),
|
| 136 |
# Figure citation linking (figure_links.py): a text row that cites
|
| 137 |
# "Fig. 3" gets figure_id/figure_ref/score/signals + the PNG embedded in
|
|
@@ -172,16 +184,74 @@ DISCOVERY_TIMEOUT_MINUTES = _f("AGENT_DISCOVERY_TIMEOUT_MINUTES", 45.0)
|
|
| 172 |
# query). 0 disables; the hard budget still applies.
|
| 173 |
CYCLE_STALL_MINUTES = _f("AGENT_CYCLE_STALL_MINUTES", 15.0)
|
| 174 |
|
| 175 |
-
#
|
| 176 |
-
#
|
| 177 |
-
#
|
| 178 |
-
|
| 179 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 180 |
|
| 181 |
|
| 182 |
-
def
|
| 183 |
-
|
| 184 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 185 |
|
| 186 |
# Keep-alive: the container requests its own public URL every N minutes so
|
| 187 |
# the hosting platform never sees the Space idle (free Spaces sleep after 48 h
|
|
|
|
| 132 |
# Thinking tokens are billed as output (8x input); "low" cut a cycle's
|
| 133 |
# bill by roughly two thirds. extraction.THINKING is set from this each cycle.
|
| 134 |
"gemini_thinking": os.environ.get("GEMINI_THINKING", "low").strip().lower() or "dynamic",
|
| 135 |
+
# Gemini models (30 Sep 2026). gemini-3.5-flash is $1.50 / $9.00 per M
|
| 136 |
+
# in/out; gemini-3.5-flash-lite $0.30 / $2.50. The text model is what the
|
| 137 |
+
# paper's numbers are stamped with (rows carry it in `model`), so it only
|
| 138 |
+
# changes on purpose. Vision model "" = same as the text model; the one
|
| 139 |
+
# classify call per PDF is an easy task for the Lite model.
|
| 140 |
+
"gemini_model": os.environ.get("GEMINI_MODEL", "gemini-3.5-flash").strip() or "gemini-3.5-flash",
|
| 141 |
+
"gemini_vision_model": os.environ.get("GEMINI_VISION_MODEL", "").strip(),
|
| 142 |
+
# Near-duplicate gate: a PDF whose text is >= this similar to a document
|
| 143 |
+
# already ingested (regional datasheet variants, a preprint and its
|
| 144 |
+
# published version, the same PDF under another URL) is skipped before
|
| 145 |
+
# any Gemini call. 0 disables.
|
| 146 |
+
"dedupe_text_similarity": _f("AGENT_DEDUPE_TEXT_SIMILARITY", 0.9),
|
| 147 |
"max_figures_per_pdf": _i("AGENT_MAX_FIGURES", 12),
|
| 148 |
# Figure citation linking (figure_links.py): a text row that cites
|
| 149 |
# "Fig. 3" gets figure_id/figure_ref/score/signals + the PNG embedded in
|
|
|
|
| 184 |
# query). 0 disables; the hard budget still applies.
|
| 185 |
CYCLE_STALL_MINUTES = _f("AGENT_CYCLE_STALL_MINUTES", 15.0)
|
| 186 |
|
| 187 |
+
# Gemini list prices, USD per million tokens (input, output; thinking tokens
|
| 188 |
+
# bill as output). From ai.google.dev/gemini-api/docs/pricing, 30 Sep 2026.
|
| 189 |
+
# Until 30 Sep the cost column used one pair, 0.30 / 2.50 (Gemini 2.5 Flash),
|
| 190 |
+
# while the agent ran gemini-3.5-flash at 1.50 / 9.00: it showed ~1/4 of the
|
| 191 |
+
# real bill. Runs now store their cost (priced per model at the end of the
|
| 192 |
+
# cycle); older runs are re-priced at read time with LEGACY_RUN_MODEL.
|
| 193 |
+
MODEL_PRICES = {
|
| 194 |
+
"gemini-3.5-flash": (1.50, 9.00),
|
| 195 |
+
"gemini-3.5-flash-lite": (0.30, 2.50),
|
| 196 |
+
"gemini-3-flash-preview": (0.50, 3.00),
|
| 197 |
+
"gemini-2.5-flash": (0.30, 2.50),
|
| 198 |
+
"gemini-2.5-flash-lite": (0.10, 0.40),
|
| 199 |
+
}
|
| 200 |
+
BATCH_DISCOUNT = 0.5 # Batch API: half the standard price
|
| 201 |
+
# Every run before per-run pricing (tokens were first logged on 22 Sep 2026)
|
| 202 |
+
# used this model.
|
| 203 |
+
LEGACY_RUN_MODEL = "gemini-3.5-flash"
|
| 204 |
+
# Explicit override for a model not in the table (or a price change): both
|
| 205 |
+
# must be set; they then apply to every model.
|
| 206 |
+
_PRICE_IN_ENV = os.environ.get("AGENT_PRICE_IN_PER_M", "").strip()
|
| 207 |
+
_PRICE_OUT_ENV = os.environ.get("AGENT_PRICE_OUT_PER_M", "").strip()
|
| 208 |
+
|
| 209 |
+
|
| 210 |
+
def model_price(model: str | None) -> tuple[float, float]:
|
| 211 |
+
"""(input, output) USD per million tokens for `model`."""
|
| 212 |
+
if _PRICE_IN_ENV and _PRICE_OUT_ENV:
|
| 213 |
+
try:
|
| 214 |
+
return float(_PRICE_IN_ENV), float(_PRICE_OUT_ENV)
|
| 215 |
+
except ValueError:
|
| 216 |
+
pass
|
| 217 |
+
m = (model or LEGACY_RUN_MODEL).strip().lower()
|
| 218 |
+
if m.startswith("models/"):
|
| 219 |
+
m = m[len("models/"):]
|
| 220 |
+
if m in MODEL_PRICES:
|
| 221 |
+
return MODEL_PRICES[m]
|
| 222 |
+
# dated / suffixed variants ("gemini-3.5-flash-001"): longest known prefix
|
| 223 |
+
for k in sorted(MODEL_PRICES, key=len, reverse=True):
|
| 224 |
+
if m.startswith(k):
|
| 225 |
+
return MODEL_PRICES[k]
|
| 226 |
+
return MODEL_PRICES[LEGACY_RUN_MODEL]
|
| 227 |
+
|
| 228 |
+
|
| 229 |
+
# Kept for callers that show "the" price (the Logs column help).
|
| 230 |
+
PRICE_IN_PER_M, PRICE_OUT_PER_M = model_price(os.environ.get("GEMINI_MODEL", LEGACY_RUN_MODEL))
|
| 231 |
+
|
| 232 |
+
|
| 233 |
+
def cost_usd(tokens_in: int, tokens_out: int, model: str | None = None,
|
| 234 |
+
batch: bool = False) -> float:
|
| 235 |
+
pin, pout = model_price(model)
|
| 236 |
+
c = (float(tokens_in or 0) * pin + float(tokens_out or 0) * pout) / 1e6
|
| 237 |
+
return c * (BATCH_DISCOUNT if batch else 1.0)
|
| 238 |
|
| 239 |
|
| 240 |
+
def run_cost_usd(metrics: dict) -> float:
|
| 241 |
+
"""A cycle's cost from its metrics: text and vision tokens priced at their
|
| 242 |
+
own models, batch-mode text calls at the batch discount."""
|
| 243 |
+
tin = int(metrics.get("tokens_in", 0) or 0)
|
| 244 |
+
tout = int(metrics.get("tokens_out", 0) or 0)
|
| 245 |
+
vin = int(metrics.get("vision_tokens_in", 0) or 0)
|
| 246 |
+
vout = int(metrics.get("vision_tokens_out", 0) or 0)
|
| 247 |
+
bin_ = int(metrics.get("batch_tokens_in", 0) or 0)
|
| 248 |
+
bout = int(metrics.get("batch_tokens_out", 0) or 0)
|
| 249 |
+
tm = metrics.get("text_model") or LEGACY_RUN_MODEL
|
| 250 |
+
vm = metrics.get("vision_model") or tm
|
| 251 |
+
text_in, text_out = max(0, tin - vin - bin_), max(0, tout - vout - bout)
|
| 252 |
+
return (cost_usd(text_in, text_out, tm)
|
| 253 |
+
+ cost_usd(bin_, bout, tm, batch=True)
|
| 254 |
+
+ cost_usd(vin, vout, vm))
|
| 255 |
|
| 256 |
# Keep-alive: the container requests its own public URL every N minutes so
|
| 257 |
# the hosting platform never sees the Space idle (free Spaces sleep after 48 h
|
|
@@ -88,6 +88,9 @@ def node_plan(state: AgentState) -> AgentState:
|
|
| 88 |
# Cost knob: how much every Gemini call in this cycle may think.
|
| 89 |
import extraction as _ex
|
| 90 |
_ex.THINKING = str(cfg.get("gemini_thinking", _ex.THINKING) or "dynamic").strip().lower()
|
|
|
|
|
|
|
|
|
|
| 91 |
# Budget hold: Gemini refused the previous cycle for budget reasons
|
| 92 |
# (prepaid credits depleted, spending cap). Downloading more would
|
| 93 |
# only pile PDFs up on the ephemeral disk, so this cycle downloads
|
|
@@ -651,6 +654,10 @@ def node_ingest(state: AgentState) -> AgentState:
|
|
| 651 |
m.setdefault("vision_calls", 0)
|
| 652 |
m.setdefault("tokens_in", 0)
|
| 653 |
m.setdefault("tokens_out", 0)
|
|
|
|
|
|
|
|
|
|
|
|
|
| 654 |
m.setdefault("rows_linked", 0)
|
| 655 |
results = []
|
| 656 |
_adopt_leftover_pdfs(state)
|
|
@@ -721,6 +728,16 @@ def node_ingest(state: AgentState) -> AgentState:
|
|
| 721 |
+ "; ".join(problems)
|
| 722 |
+ " -- run `python pg_migrate.py --apply`",
|
| 723 |
level="warn", node="ingest", run_id=state["run_id"])
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 724 |
rows_by_file: dict[str, int] = {}
|
| 725 |
workers = max(1, min(int(cfg.get("ingest_workers", C.INGEST_WORKERS_DEFAULT)),
|
| 726 |
C.HARD_MAX_INGEST_WORKERS, len(state["downloads"])))
|
|
@@ -736,7 +753,7 @@ def node_ingest(state: AgentState) -> AgentState:
|
|
| 736 |
kept: list[str] = [] # PDFs left on disk for the next cycle
|
| 737 |
budget_reason = ""
|
| 738 |
for rec, result, exc in _ingest_all(state["downloads"], figure_opts, link_opts,
|
| 739 |
-
workers, stop):
|
| 740 |
pdf_path = C.PDF_DIR / rec["filename"]
|
| 741 |
if isinstance(exc, _BudgetStop) or (result is not None and getattr(result, "quota", False)):
|
| 742 |
# Gemini refused for budget reasons (monthly spending cap,
|
|
@@ -754,6 +771,8 @@ def node_ingest(state: AgentState) -> AgentState:
|
|
| 754 |
if exc is not None:
|
| 755 |
m["errors"] += 1
|
| 756 |
agdb.update_ingest_status(meta_conn, rec["filename"], f"error")
|
|
|
|
|
|
|
| 757 |
if rec.get("_cand_id"):
|
| 758 |
agdb.mark_candidate(meta_conn, rec["_cand_id"], "failed", str(exc))
|
| 759 |
tb = getattr(exc, "__traceback_text__", "")
|
|
@@ -762,10 +781,23 @@ def node_ingest(state: AgentState) -> AgentState:
|
|
| 762 |
+ (f"\n{tb}" if tb else ""),
|
| 763 |
level="error", node="ingest", run_id=state["run_id"])
|
| 764 |
continue
|
| 765 |
-
|
| 766 |
-
m["tokens_in"] += int(getattr(result, "tokens_in", 0) or 0)
|
| 767 |
-
m["tokens_out"] += int(getattr(result, "tokens_out", 0) or 0)
|
| 768 |
m["rows_linked"] += int(getattr(result, "rows_linked", 0) or 0)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 769 |
if result.error:
|
| 770 |
if str(result.error).startswith("gemini_error"):
|
| 771 |
# Transient (rate limit, 5xx, network): keep the PDF for
|
|
@@ -777,6 +809,8 @@ def node_ingest(state: AgentState) -> AgentState:
|
|
| 777 |
level="warn", node="ingest", run_id=state["run_id"])
|
| 778 |
continue
|
| 779 |
agdb.update_ingest_status(meta_conn, rec["filename"], result.error)
|
|
|
|
|
|
|
| 780 |
if rec.get("_cand_id"):
|
| 781 |
agdb.mark_candidate(meta_conn, rec["_cand_id"], "failed", result.error)
|
| 782 |
agdb.log_event(meta_conn, f"{rec['filename']}: {result.error}",
|
|
@@ -864,12 +898,24 @@ def node_ingest(state: AgentState) -> AgentState:
|
|
| 864 |
return state
|
| 865 |
|
| 866 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 867 |
class _BudgetStop(Exception):
|
| 868 |
"""Marker: this PDF was not attempted because Gemini already refused
|
| 869 |
this cycle for budget reasons (spending cap / quota)."""
|
| 870 |
|
| 871 |
|
| 872 |
-
def _ingest_one(rec: dict, figure_opts, link_opts, stop: Optional[threading.Event] = None
|
|
|
|
| 873 |
"""Worker: extract one PDF on a private Postgres connection.
|
| 874 |
Returns (rec, PdfResult or None, exception or None)."""
|
| 875 |
pdf_path = C.PDF_DIR / rec["filename"]
|
|
@@ -877,6 +923,12 @@ def _ingest_one(rec: dict, figure_opts, link_opts, stop: Optional[threading.Even
|
|
| 877 |
return rec, None, _BudgetStop()
|
| 878 |
pg = None
|
| 879 |
try:
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 880 |
pg = pg_mirror.connect_from_env()
|
| 881 |
result = batch_ingest.process_pdf(pdf_path, pg, C.GEMINI_API_KEY,
|
| 882 |
db=pg_mirror,
|
|
@@ -902,8 +954,20 @@ def _ingest_one(rec: dict, figure_opts, link_opts, stop: Optional[threading.Even
|
|
| 902 |
pass
|
| 903 |
|
| 904 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 905 |
def _ingest_all(downloads: list, figure_opts, link_opts, workers: int,
|
| 906 |
-
stop: Optional[threading.Event] = None):
|
| 907 |
"""Yield (rec, result, exc) for every download, in completion order.
|
| 908 |
workers == 1 keeps the old one-at-a-time behaviour (and the old
|
| 909 |
connection count) — the only difference is the private connection.
|
|
@@ -911,11 +975,12 @@ def _ingest_all(downloads: list, figure_opts, link_opts, workers: int,
|
|
| 911 |
started come back as _BudgetStop without a request."""
|
| 912 |
if workers <= 1:
|
| 913 |
for rec in downloads:
|
| 914 |
-
yield _ingest_one(rec, figure_opts, link_opts, stop)
|
| 915 |
return
|
| 916 |
from concurrent.futures import ThreadPoolExecutor, as_completed
|
| 917 |
with ThreadPoolExecutor(max_workers=workers, thread_name_prefix="ingest") as pool:
|
| 918 |
-
futures = [pool.submit(_ingest_one, rec, figure_opts, link_opts, stop
|
|
|
|
| 919 |
for fut in as_completed(futures):
|
| 920 |
yield fut.result()
|
| 921 |
|
|
@@ -1008,8 +1073,11 @@ def node_report(state: AgentState) -> AgentState:
|
|
| 1008 |
f"{m.get('figures_found', 0)} harvested, "
|
| 1009 |
f"{m.get('figure_rows', 0)} estimate rows, "
|
| 1010 |
f"{m.get('rows_linked', 0)} rows linked; "
|
| 1011 |
-
f"tokens {m.get('tokens_in', 0)} in / {m.get('tokens_out', 0)} out "
|
| 1012 |
-
f"
|
|
|
|
|
|
|
|
|
|
| 1013 |
node="report", run_id=state["run_id"])
|
| 1014 |
finally:
|
| 1015 |
conn.close()
|
|
|
|
| 88 |
# Cost knob: how much every Gemini call in this cycle may think.
|
| 89 |
import extraction as _ex
|
| 90 |
_ex.THINKING = str(cfg.get("gemini_thinking", _ex.THINKING) or "dynamic").strip().lower()
|
| 91 |
+
# Cost knob: which models the text and figure calls go to.
|
| 92 |
+
_ex.GEMINI_MODEL = str(cfg.get("gemini_model") or _ex.GEMINI_MODEL).strip()
|
| 93 |
+
_ex.GEMINI_VISION_MODEL = str(cfg.get("gemini_vision_model") or "").strip()
|
| 94 |
# Budget hold: Gemini refused the previous cycle for budget reasons
|
| 95 |
# (prepaid credits depleted, spending cap). Downloading more would
|
| 96 |
# only pile PDFs up on the ephemeral disk, so this cycle downloads
|
|
|
|
| 654 |
m.setdefault("vision_calls", 0)
|
| 655 |
m.setdefault("tokens_in", 0)
|
| 656 |
m.setdefault("tokens_out", 0)
|
| 657 |
+
m.setdefault("tokens_thinking", 0)
|
| 658 |
+
m.setdefault("vision_tokens_in", 0)
|
| 659 |
+
m.setdefault("vision_tokens_out", 0)
|
| 660 |
+
m.setdefault("duplicate_text", 0)
|
| 661 |
m.setdefault("rows_linked", 0)
|
| 662 |
results = []
|
| 663 |
_adopt_leftover_pdfs(state)
|
|
|
|
| 728 |
+ "; ".join(problems)
|
| 729 |
+ " -- run `python pg_migrate.py --apply`",
|
| 730 |
level="warn", node="ingest", run_id=state["run_id"])
|
| 731 |
+
# Near-duplicate gate (agent/textdedupe.py): skip a PDF whose text
|
| 732 |
+
# repeats a document already ingested, before any Gemini call.
|
| 733 |
+
text_gate = None
|
| 734 |
+
try:
|
| 735 |
+
sim = float(cfg.get("dedupe_text_similarity", 0.0) or 0.0)
|
| 736 |
+
except (TypeError, ValueError):
|
| 737 |
+
sim = 0.0
|
| 738 |
+
if 0.0 < sim <= 1.0:
|
| 739 |
+
from agent import textdedupe
|
| 740 |
+
text_gate = textdedupe.Gate(agdb.connect, sim)
|
| 741 |
rows_by_file: dict[str, int] = {}
|
| 742 |
workers = max(1, min(int(cfg.get("ingest_workers", C.INGEST_WORKERS_DEFAULT)),
|
| 743 |
C.HARD_MAX_INGEST_WORKERS, len(state["downloads"])))
|
|
|
|
| 753 |
kept: list[str] = [] # PDFs left on disk for the next cycle
|
| 754 |
budget_reason = ""
|
| 755 |
for rec, result, exc in _ingest_all(state["downloads"], figure_opts, link_opts,
|
| 756 |
+
workers, stop, text_gate=text_gate):
|
| 757 |
pdf_path = C.PDF_DIR / rec["filename"]
|
| 758 |
if isinstance(exc, _BudgetStop) or (result is not None and getattr(result, "quota", False)):
|
| 759 |
# Gemini refused for budget reasons (monthly spending cap,
|
|
|
|
| 771 |
if exc is not None:
|
| 772 |
m["errors"] += 1
|
| 773 |
agdb.update_ingest_status(meta_conn, rec["filename"], f"error")
|
| 774 |
+
if text_gate is not None:
|
| 775 |
+
text_gate.release_file(rec["filename"])
|
| 776 |
if rec.get("_cand_id"):
|
| 777 |
agdb.mark_candidate(meta_conn, rec["_cand_id"], "failed", str(exc))
|
| 778 |
tb = getattr(exc, "__traceback_text__", "")
|
|
|
|
| 781 |
+ (f"\n{tb}" if tb else ""),
|
| 782 |
level="error", node="ingest", run_id=state["run_id"])
|
| 783 |
continue
|
| 784 |
+
_add_usage(m, result)
|
|
|
|
|
|
|
| 785 |
m["rows_linked"] += int(getattr(result, "rows_linked", 0) or 0)
|
| 786 |
+
if result.error == "duplicate_text":
|
| 787 |
+
# Near-duplicate of a document already ingested: nothing was
|
| 788 |
+
# sent to Gemini. Final, like a DOI-level duplicate.
|
| 789 |
+
m["duplicate_text"] += 1
|
| 790 |
+
agdb.update_ingest_status(meta_conn, rec["filename"], "duplicate_text")
|
| 791 |
+
if rec.get("_cand_id"):
|
| 792 |
+
agdb.mark_candidate(meta_conn, rec["_cand_id"], "duplicate",
|
| 793 |
+
f"text repeats {result.duplicate_of}")
|
| 794 |
+
agdb.log_event(meta_conn, f"{rec['filename']}: skipped before extraction — "
|
| 795 |
+
f"its text repeats an ingested document "
|
| 796 |
+
f"({str(result.duplicate_of)[:12]})",
|
| 797 |
+
node="ingest", run_id=state["run_id"])
|
| 798 |
+
pdf_path.unlink(missing_ok=True)
|
| 799 |
+
continue
|
| 800 |
+
results.append(result)
|
| 801 |
if result.error:
|
| 802 |
if str(result.error).startswith("gemini_error"):
|
| 803 |
# Transient (rate limit, 5xx, network): keep the PDF for
|
|
|
|
| 809 |
level="warn", node="ingest", run_id=state["run_id"])
|
| 810 |
continue
|
| 811 |
agdb.update_ingest_status(meta_conn, rec["filename"], result.error)
|
| 812 |
+
if text_gate is not None:
|
| 813 |
+
text_gate.release_file(rec["filename"]) # a variant may stand in later
|
| 814 |
if rec.get("_cand_id"):
|
| 815 |
agdb.mark_candidate(meta_conn, rec["_cand_id"], "failed", result.error)
|
| 816 |
agdb.log_event(meta_conn, f"{rec['filename']}: {result.error}",
|
|
|
|
| 898 |
return state
|
| 899 |
|
| 900 |
|
| 901 |
+
def _add_usage(m: dict, result) -> None:
|
| 902 |
+
"""Fold one PdfResult's Gemini usage into the cycle metrics."""
|
| 903 |
+
for k in ("tokens_in", "tokens_out", "tokens_thinking",
|
| 904 |
+
"vision_tokens_in", "vision_tokens_out"):
|
| 905 |
+
m[k] = int(m.get(k, 0) or 0) + int(getattr(result, k, 0) or 0)
|
| 906 |
+
if getattr(result, "text_model", ""):
|
| 907 |
+
m["text_model"] = result.text_model
|
| 908 |
+
if getattr(result, "vision_model", ""):
|
| 909 |
+
m["vision_model"] = result.vision_model
|
| 910 |
+
|
| 911 |
+
|
| 912 |
class _BudgetStop(Exception):
|
| 913 |
"""Marker: this PDF was not attempted because Gemini already refused
|
| 914 |
this cycle for budget reasons (spending cap / quota)."""
|
| 915 |
|
| 916 |
|
| 917 |
+
def _ingest_one(rec: dict, figure_opts, link_opts, stop: Optional[threading.Event] = None,
|
| 918 |
+
text_gate=None):
|
| 919 |
"""Worker: extract one PDF on a private Postgres connection.
|
| 920 |
Returns (rec, PdfResult or None, exception or None)."""
|
| 921 |
pdf_path = C.PDF_DIR / rec["filename"]
|
|
|
|
| 923 |
return rec, None, _BudgetStop()
|
| 924 |
pg = None
|
| 925 |
try:
|
| 926 |
+
if text_gate is not None:
|
| 927 |
+
dup = _near_duplicate(text_gate, pdf_path)
|
| 928 |
+
if dup:
|
| 929 |
+
res = batch_ingest._empty_result(pdf_path, time.time(), "duplicate_text")
|
| 930 |
+
res.duplicate_of = dup
|
| 931 |
+
return rec, res, None
|
| 932 |
pg = pg_mirror.connect_from_env()
|
| 933 |
result = batch_ingest.process_pdf(pdf_path, pg, C.GEMINI_API_KEY,
|
| 934 |
db=pg_mirror,
|
|
|
|
| 954 |
pass
|
| 955 |
|
| 956 |
|
| 957 |
+
def _near_duplicate(text_gate, pdf_path) -> Optional[str]:
|
| 958 |
+
"""sha1 of the ingested document this PDF's text repeats, or None. The
|
| 959 |
+
gate is an optimisation: any failure means 'not a duplicate'."""
|
| 960 |
+
try:
|
| 961 |
+
data = pdf_path.read_bytes()
|
| 962 |
+
import hashlib
|
| 963 |
+
return text_gate(data, hashlib.sha1(data).hexdigest(), pdf_path.name)
|
| 964 |
+
except Exception as exc:
|
| 965 |
+
log.warning("near-duplicate gate failed for %s: %s", pdf_path.name, exc)
|
| 966 |
+
return None
|
| 967 |
+
|
| 968 |
+
|
| 969 |
def _ingest_all(downloads: list, figure_opts, link_opts, workers: int,
|
| 970 |
+
stop: Optional[threading.Event] = None, text_gate=None):
|
| 971 |
"""Yield (rec, result, exc) for every download, in completion order.
|
| 972 |
workers == 1 keeps the old one-at-a-time behaviour (and the old
|
| 973 |
connection count) — the only difference is the private connection.
|
|
|
|
| 975 |
started come back as _BudgetStop without a request."""
|
| 976 |
if workers <= 1:
|
| 977 |
for rec in downloads:
|
| 978 |
+
yield _ingest_one(rec, figure_opts, link_opts, stop, text_gate)
|
| 979 |
return
|
| 980 |
from concurrent.futures import ThreadPoolExecutor, as_completed
|
| 981 |
with ThreadPoolExecutor(max_workers=workers, thread_name_prefix="ingest") as pool:
|
| 982 |
+
futures = [pool.submit(_ingest_one, rec, figure_opts, link_opts, stop, text_gate)
|
| 983 |
+
for rec in downloads]
|
| 984 |
for fut in as_completed(futures):
|
| 985 |
yield fut.result()
|
| 986 |
|
|
|
|
| 1073 |
f"{m.get('figures_found', 0)} harvested, "
|
| 1074 |
f"{m.get('figure_rows', 0)} estimate rows, "
|
| 1075 |
f"{m.get('rows_linked', 0)} rows linked; "
|
| 1076 |
+
f"tokens {m.get('tokens_in', 0)} in / {m.get('tokens_out', 0)} out, "
|
| 1077 |
+
f"{m.get('tokens_thinking', 0)} of it thinking "
|
| 1078 |
+
f"(~${C.run_cost_usd(m):.4f} at {m.get('text_model') or C.LEGACY_RUN_MODEL} prices)"
|
| 1079 |
+
+ (f"; {m['duplicate_text']} near-duplicate PDF(s) skipped"
|
| 1080 |
+
if m.get("duplicate_text") else ""),
|
| 1081 |
node="report", run_id=state["run_id"])
|
| 1082 |
finally:
|
| 1083 |
conn.close()
|
|
@@ -0,0 +1,182 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
"""Near-duplicate gate: skip a PDF whose text repeats a document already
|
| 2 |
+
ingested, before any Gemini call is made (30 Sep 2026).
|
| 3 |
+
|
| 4 |
+
Why: the same content arrives under different bytes — a manufacturer's
|
| 5 |
+
regional datasheet variants (…global_version… / …eu_version…), a preprint
|
| 6 |
+
and its published version, the same paper from a repository and a publisher.
|
| 7 |
+
The sha1 check only catches byte-identical files, so each variant was sent
|
| 8 |
+
to Gemini (~$0.10–0.40 each) for rows the dedup grain then threw away.
|
| 9 |
+
|
| 10 |
+
How: MinHash over word 5-shingles of the PDF's extracted text (PyMuPDF, the
|
| 11 |
+
same text the grounding check uses). 128 slots, stored as 1 KB per document
|
| 12 |
+
in agent_text_fingerprints. The estimated Jaccard similarity of two
|
| 13 |
+
documents is the share of equal slots; at >= the threshold (default 0.9,
|
| 14 |
+
Run Control) the new PDF is skipped. Different datasheets from one template
|
| 15 |
+
share boilerplate but not their numbers, and every changed number breaks
|
| 16 |
+
five shingles, so a family of sheets stays well below 0.9.
|
| 17 |
+
|
| 18 |
+
A fingerprint is registered when its PDF passes the gate (a claim), so two
|
| 19 |
+
variants in the same cycle — even on parallel workers — cannot both pass.
|
| 20 |
+
If the claimed PDF then fails for good, release() drops the claim.
|
| 21 |
+
Documents ingested before this existed have no fingerprint (their PDFs are
|
| 22 |
+
gone), so the gate only knows documents from its deploy on.
|
| 23 |
+
"""
|
| 24 |
+
from __future__ import annotations
|
| 25 |
+
|
| 26 |
+
import hashlib
|
| 27 |
+
import logging
|
| 28 |
+
import re
|
| 29 |
+
import threading
|
| 30 |
+
from typing import Callable, Optional
|
| 31 |
+
|
| 32 |
+
import numpy as np
|
| 33 |
+
|
| 34 |
+
log = logging.getLogger("agent.textdedupe")
|
| 35 |
+
|
| 36 |
+
NUM_PERM = 128
|
| 37 |
+
SHINGLE = 5
|
| 38 |
+
MIN_SHINGLES = 40 # below this there is too little text to judge
|
| 39 |
+
_PRIME = np.uint64((1 << 61) - 1)
|
| 40 |
+
_rng = np.random.RandomState(20260930)
|
| 41 |
+
_A = _rng.randint(1, 1 << 31, size=NUM_PERM).astype(np.uint64)
|
| 42 |
+
_B = _rng.randint(0, 1 << 31, size=NUM_PERM).astype(np.uint64)
|
| 43 |
+
_WORD = re.compile(r"[a-z0-9]+(?:[.,][0-9]+)*")
|
| 44 |
+
|
| 45 |
+
DDL = """
|
| 46 |
+
CREATE TABLE IF NOT EXISTS agent_text_fingerprints (
|
| 47 |
+
sha1 text PRIMARY KEY,
|
| 48 |
+
filename text,
|
| 49 |
+
minhash bytea NOT NULL,
|
| 50 |
+
n_shingles integer,
|
| 51 |
+
created_at timestamptz DEFAULT now()
|
| 52 |
+
);
|
| 53 |
+
"""
|
| 54 |
+
|
| 55 |
+
|
| 56 |
+
def _words(text: str) -> list[str]:
|
| 57 |
+
return _WORD.findall((text or "").lower())
|
| 58 |
+
|
| 59 |
+
|
| 60 |
+
def fingerprint(text: str) -> Optional[tuple[bytes, int]]:
|
| 61 |
+
"""(minhash bytes, number of distinct shingles), or None for too little text."""
|
| 62 |
+
w = _words(text)
|
| 63 |
+
if len(w) < SHINGLE:
|
| 64 |
+
return None
|
| 65 |
+
sh = {" ".join(w[i:i + SHINGLE]) for i in range(len(w) - SHINGLE + 1)}
|
| 66 |
+
if len(sh) < MIN_SHINGLES:
|
| 67 |
+
return None
|
| 68 |
+
# one 31-bit hash per shingle, then NUM_PERM universal hashes of it
|
| 69 |
+
h = np.fromiter((int.from_bytes(hashlib.blake2b(s.encode(), digest_size=4).digest(), "little")
|
| 70 |
+
for s in sh), dtype=np.uint64, count=len(sh))
|
| 71 |
+
perm = (np.outer(h, _A) + _B) % _PRIME # (n_shingles, NUM_PERM), < 2^61: no overflow
|
| 72 |
+
return perm.min(axis=0).astype(np.uint64).tobytes(), len(sh)
|
| 73 |
+
|
| 74 |
+
|
| 75 |
+
def similarity(a: bytes, b: bytes) -> float:
|
| 76 |
+
x = np.frombuffer(a, dtype=np.uint64)
|
| 77 |
+
y = np.frombuffer(b, dtype=np.uint64)
|
| 78 |
+
if x.shape != y.shape or not len(x):
|
| 79 |
+
return 0.0
|
| 80 |
+
return float(np.mean(x == y))
|
| 81 |
+
|
| 82 |
+
|
| 83 |
+
def ensure_table(conn) -> None:
|
| 84 |
+
with conn.cursor() as cur:
|
| 85 |
+
cur.execute(DDL)
|
| 86 |
+
conn.commit()
|
| 87 |
+
|
| 88 |
+
|
| 89 |
+
class Gate:
|
| 90 |
+
"""Callable for batch_ingest.process_pdf(text_gate=...):
|
| 91 |
+
gate(pdf_bytes, sha1) -> sha1 of the ingested document this PDF repeats,
|
| 92 |
+
or None (and the PDF's fingerprint is claimed)."""
|
| 93 |
+
|
| 94 |
+
_lock = threading.Lock() # parallel ingest workers share one process
|
| 95 |
+
|
| 96 |
+
def __init__(self, connect: Callable, threshold: float,
|
| 97 |
+
text_of: Optional[Callable[[bytes], str]] = None):
|
| 98 |
+
self.connect = connect
|
| 99 |
+
self.threshold = float(threshold or 0.0)
|
| 100 |
+
self.text_of = text_of or _pdf_text
|
| 101 |
+
self._cache: Optional[list[tuple[str, bytes]]] = None
|
| 102 |
+
self.claimed: dict[str, str] = {} # sha1 -> filename, claimed by this gate
|
| 103 |
+
|
| 104 |
+
@property
|
| 105 |
+
def enabled(self) -> bool:
|
| 106 |
+
return 0.0 < self.threshold <= 1.0
|
| 107 |
+
|
| 108 |
+
def _load(self, conn) -> list[tuple[str, bytes]]:
|
| 109 |
+
if self._cache is None:
|
| 110 |
+
with conn.cursor() as cur:
|
| 111 |
+
cur.execute("SELECT sha1, minhash FROM agent_text_fingerprints")
|
| 112 |
+
self._cache = [(r[0], bytes(r[1])) for r in cur.fetchall()]
|
| 113 |
+
return self._cache
|
| 114 |
+
|
| 115 |
+
def __call__(self, pdf_bytes: bytes, sha1: str, filename: str = "") -> Optional[str]:
|
| 116 |
+
if not self.enabled:
|
| 117 |
+
return None
|
| 118 |
+
fp = fingerprint(self.text_of(pdf_bytes))
|
| 119 |
+
if fp is None:
|
| 120 |
+
return None
|
| 121 |
+
mh, n = fp
|
| 122 |
+
with self._lock:
|
| 123 |
+
conn = self.connect()
|
| 124 |
+
try:
|
| 125 |
+
best, best_sha = 0.0, None
|
| 126 |
+
for other_sha, other in self._load(conn):
|
| 127 |
+
if other_sha == sha1:
|
| 128 |
+
continue # itself (a kept PDF coming back)
|
| 129 |
+
s = similarity(mh, other)
|
| 130 |
+
if s > best:
|
| 131 |
+
best, best_sha = s, other_sha
|
| 132 |
+
if best_sha is not None and best >= self.threshold:
|
| 133 |
+
log.info("near-duplicate: %s ~ %s (%.2f)", sha1[:12], best_sha[:12], best)
|
| 134 |
+
return best_sha
|
| 135 |
+
with conn.cursor() as cur:
|
| 136 |
+
cur.execute(
|
| 137 |
+
"INSERT INTO agent_text_fingerprints (sha1, filename, minhash, n_shingles) "
|
| 138 |
+
"VALUES (%s, %s, %s, %s) ON CONFLICT (sha1) DO NOTHING",
|
| 139 |
+
(sha1, filename or None, mh, n))
|
| 140 |
+
conn.commit()
|
| 141 |
+
if self._cache is not None and not any(s == sha1 for s, _ in self._cache):
|
| 142 |
+
self._cache.append((sha1, mh))
|
| 143 |
+
self.claimed[sha1] = filename
|
| 144 |
+
return None
|
| 145 |
+
finally:
|
| 146 |
+
try:
|
| 147 |
+
conn.close()
|
| 148 |
+
except Exception:
|
| 149 |
+
pass
|
| 150 |
+
|
| 151 |
+
def release(self, sha1: str) -> None:
|
| 152 |
+
"""Drop a claim whose PDF failed for good (so a variant can stand in)."""
|
| 153 |
+
if not sha1:
|
| 154 |
+
return
|
| 155 |
+
with self._lock:
|
| 156 |
+
conn = self.connect()
|
| 157 |
+
try:
|
| 158 |
+
with conn.cursor() as cur:
|
| 159 |
+
cur.execute("DELETE FROM agent_text_fingerprints WHERE sha1 = %s", (sha1,))
|
| 160 |
+
conn.commit()
|
| 161 |
+
if self._cache is not None:
|
| 162 |
+
self._cache = [(s, m) for s, m in self._cache if s != sha1]
|
| 163 |
+
self.claimed.pop(sha1, None)
|
| 164 |
+
finally:
|
| 165 |
+
try:
|
| 166 |
+
conn.close()
|
| 167 |
+
except Exception:
|
| 168 |
+
pass
|
| 169 |
+
|
| 170 |
+
|
| 171 |
+
def release_file(self, filename: str) -> None:
|
| 172 |
+
for sha, fn in list(self.claimed.items()):
|
| 173 |
+
if fn == filename:
|
| 174 |
+
self.release(sha)
|
| 175 |
+
|
| 176 |
+
|
| 177 |
+
def _pdf_text(pdf_bytes: bytes) -> str:
|
| 178 |
+
import extraction
|
| 179 |
+
try:
|
| 180 |
+
return "\n".join(extraction.pdf_page_texts(pdf_bytes))
|
| 181 |
+
except Exception:
|
| 182 |
+
return ""
|
|
@@ -157,6 +157,16 @@ class PdfResult:
|
|
| 157 |
# call was made; tokens_out includes thinking tokens, billed as output)
|
| 158 |
tokens_in: int = 0
|
| 159 |
tokens_out: int = 0
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 160 |
# figure-linking phase (all zero when --link-figures is off)
|
| 161 |
rows_linked: int = 0 # text rows that got a figure link (new or backfilled)
|
| 162 |
link_figures_found: int = 0 # figures harvested for the link pass (local, no API call)
|
|
@@ -526,6 +536,11 @@ def _run_figure_stage(
|
|
| 526 |
result.vision_calls = stage.vision.total
|
| 527 |
result.tokens_in += getattr(stage.vision, "tokens_in", 0)
|
| 528 |
result.tokens_out += getattr(stage.vision, "tokens_out", 0)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 529 |
result.figure_filters = {
|
| 530 |
k: v for k, v in dataclasses.asdict(stage.harvest).items() if v
|
| 531 |
}
|
|
@@ -650,10 +665,15 @@ def process_pdf(
|
|
| 650 |
db: Any = None,
|
| 651 |
figure_opts: Optional[FigureOptions] = None,
|
| 652 |
link_opts: Optional[LinkOptions] = None,
|
|
|
|
|
|
|
| 653 |
) -> PdfResult:
|
| 654 |
# `db` = backend module providing seen_sha1 / record_source /
|
| 655 |
# already_inserted / insert_row. Defaults to this module (SQLite);
|
| 656 |
# main() passes pg_mirror for --pg. Same logic either way.
|
|
|
|
|
|
|
|
|
|
| 657 |
db = db or sys.modules[__name__]
|
| 658 |
started = time.time()
|
| 659 |
pdf_bytes = pdf_path.read_bytes()
|
|
@@ -684,7 +704,8 @@ def process_pdf(
|
|
| 684 |
return result
|
| 685 |
|
| 686 |
try:
|
| 687 |
-
extracted =
|
|
|
|
| 688 |
except requests.RequestException as exc:
|
| 689 |
# requests quotes the keyed URL in HTTPError/ConnectionError messages.
|
| 690 |
err = _empty_result(pdf_path, started,
|
|
@@ -693,6 +714,8 @@ def process_pdf(
|
|
| 693 |
return err
|
| 694 |
tokens_in = int(getattr(extracted, "tokens_in", 0) or 0)
|
| 695 |
tokens_out = int(getattr(extracted, "tokens_out", 0) or 0)
|
|
|
|
|
|
|
| 696 |
|
| 697 |
if extracted.doc_status == "scanned_no_text":
|
| 698 |
# Don't fabricate rows from an image-only PDF (Task 10). Whole-page
|
|
@@ -704,6 +727,7 @@ def process_pdf(
|
|
| 704 |
# Billed even though nothing came back: keep the usage on the result.
|
| 705 |
empty = _empty_result(pdf_path, started, "empty_extraction")
|
| 706 |
empty.tokens_in, empty.tokens_out = tokens_in, tokens_out
|
|
|
|
| 707 |
return empty
|
| 708 |
|
| 709 |
# Ground every value against the PDF text (Task 1).
|
|
@@ -751,6 +775,8 @@ def process_pdf(
|
|
| 751 |
material_classes=sorted(set(classes)),
|
| 752 |
tokens_in=tokens_in,
|
| 753 |
tokens_out=tokens_out,
|
|
|
|
|
|
|
| 754 |
)
|
| 755 |
for k, v in link_info.items():
|
| 756 |
setattr(result, k, v)
|
|
@@ -895,6 +921,9 @@ def summarize(results: list[PdfResult]) -> dict[str, Any]:
|
|
| 895 |
"tokens": {
|
| 896 |
"tokens_in": sum(r.tokens_in for r in results),
|
| 897 |
"tokens_out": sum(r.tokens_out for r in results),
|
|
|
|
|
|
|
|
|
|
| 898 |
},
|
| 899 |
"figures": {
|
| 900 |
"figures_found": sum(r.figures_found for r in results),
|
|
|
|
| 157 |
# call was made; tokens_out includes thinking tokens, billed as output)
|
| 158 |
tokens_in: int = 0
|
| 159 |
tokens_out: int = 0
|
| 160 |
+
# Breakdown for cost analysis (all included in tokens_in / tokens_out):
|
| 161 |
+
# thinking tokens of every call, and the vision calls' share.
|
| 162 |
+
tokens_thinking: int = 0
|
| 163 |
+
vision_tokens_in: int = 0
|
| 164 |
+
vision_tokens_out: int = 0
|
| 165 |
+
text_model: str = "" # model of the text call ("" = no call made)
|
| 166 |
+
vision_model: str = "" # model of the vision calls ("" = none made)
|
| 167 |
+
# near-duplicate gate: sha1 of the already-ingested document this one
|
| 168 |
+
# repeats (then nothing was sent to Gemini)
|
| 169 |
+
duplicate_of: Optional[str] = None
|
| 170 |
# figure-linking phase (all zero when --link-figures is off)
|
| 171 |
rows_linked: int = 0 # text rows that got a figure link (new or backfilled)
|
| 172 |
link_figures_found: int = 0 # figures harvested for the link pass (local, no API call)
|
|
|
|
| 536 |
result.vision_calls = stage.vision.total
|
| 537 |
result.tokens_in += getattr(stage.vision, "tokens_in", 0)
|
| 538 |
result.tokens_out += getattr(stage.vision, "tokens_out", 0)
|
| 539 |
+
result.tokens_thinking += getattr(stage.vision, "tokens_thinking", 0)
|
| 540 |
+
result.vision_tokens_in += getattr(stage.vision, "tokens_in", 0)
|
| 541 |
+
result.vision_tokens_out += getattr(stage.vision, "tokens_out", 0)
|
| 542 |
+
if stage.vision.total:
|
| 543 |
+
result.vision_model = getattr(stage.vision, "model", "") or extraction.vision_model()
|
| 544 |
result.figure_filters = {
|
| 545 |
k: v for k, v in dataclasses.asdict(stage.harvest).items() if v
|
| 546 |
}
|
|
|
|
| 665 |
db: Any = None,
|
| 666 |
figure_opts: Optional[FigureOptions] = None,
|
| 667 |
link_opts: Optional[LinkOptions] = None,
|
| 668 |
+
*,
|
| 669 |
+
pre_extracted: Optional[Extraction] = None,
|
| 670 |
) -> PdfResult:
|
| 671 |
# `db` = backend module providing seen_sha1 / record_source /
|
| 672 |
# already_inserted / insert_row. Defaults to this module (SQLite);
|
| 673 |
# main() passes pg_mirror for --pg. Same logic either way.
|
| 674 |
+
#
|
| 675 |
+
# pre_extracted: the text call's result obtained elsewhere (the Gemini
|
| 676 |
+
# Batch API); no text call is made here then.
|
| 677 |
db = db or sys.modules[__name__]
|
| 678 |
started = time.time()
|
| 679 |
pdf_bytes = pdf_path.read_bytes()
|
|
|
|
| 704 |
return result
|
| 705 |
|
| 706 |
try:
|
| 707 |
+
extracted = (pre_extracted if pre_extracted is not None
|
| 708 |
+
else extract_from_pdf(pdf_bytes, pdf_path.name, api_key))
|
| 709 |
except requests.RequestException as exc:
|
| 710 |
# requests quotes the keyed URL in HTTPError/ConnectionError messages.
|
| 711 |
err = _empty_result(pdf_path, started,
|
|
|
|
| 714 |
return err
|
| 715 |
tokens_in = int(getattr(extracted, "tokens_in", 0) or 0)
|
| 716 |
tokens_out = int(getattr(extracted, "tokens_out", 0) or 0)
|
| 717 |
+
tokens_thinking = int(getattr(extracted, "tokens_thinking", 0) or 0)
|
| 718 |
+
text_model_used = (getattr(extracted, "model", "") or "") if (tokens_in or tokens_out) else ""
|
| 719 |
|
| 720 |
if extracted.doc_status == "scanned_no_text":
|
| 721 |
# Don't fabricate rows from an image-only PDF (Task 10). Whole-page
|
|
|
|
| 727 |
# Billed even though nothing came back: keep the usage on the result.
|
| 728 |
empty = _empty_result(pdf_path, started, "empty_extraction")
|
| 729 |
empty.tokens_in, empty.tokens_out = tokens_in, tokens_out
|
| 730 |
+
empty.tokens_thinking, empty.text_model = tokens_thinking, text_model_used
|
| 731 |
return empty
|
| 732 |
|
| 733 |
# Ground every value against the PDF text (Task 1).
|
|
|
|
| 775 |
material_classes=sorted(set(classes)),
|
| 776 |
tokens_in=tokens_in,
|
| 777 |
tokens_out=tokens_out,
|
| 778 |
+
tokens_thinking=tokens_thinking,
|
| 779 |
+
text_model=text_model_used,
|
| 780 |
)
|
| 781 |
for k, v in link_info.items():
|
| 782 |
setattr(result, k, v)
|
|
|
|
| 921 |
"tokens": {
|
| 922 |
"tokens_in": sum(r.tokens_in for r in results),
|
| 923 |
"tokens_out": sum(r.tokens_out for r in results),
|
| 924 |
+
"tokens_thinking": sum(r.tokens_thinking for r in results),
|
| 925 |
+
"vision_tokens_in": sum(r.vision_tokens_in for r in results),
|
| 926 |
+
"vision_tokens_out": sum(r.vision_tokens_out for r in results),
|
| 927 |
},
|
| 928 |
"figures": {
|
| 929 |
"figures_found": sum(r.figures_found for r in results),
|
|
@@ -64,8 +64,24 @@ log = logging.getLogger("extraction")
|
|
| 64 |
# extraction call. Previews get retired with ~2 weeks notice; stable models
|
| 65 |
# don't. Overridable via env so a future swap needs no code change.
|
| 66 |
GEMINI_MODEL = os.environ.get("GEMINI_MODEL", "gemini-3.5-flash")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 67 |
PROMPT_VERSION = "2.0"
|
| 68 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 69 |
# How much the model may "think" per call. Thinking tokens are billed as
|
| 70 |
# output (8x the input price): on 28 Sep 2026 a 25-paper cycle cost $2.02,
|
| 71 |
# $1.86 of it output tokens, ~30k per PDF where the extraction JSON itself
|
|
@@ -248,13 +264,31 @@ class Material:
|
|
| 248 |
@dataclasses.dataclass
|
| 249 |
class Extraction:
|
| 250 |
materials: list[Material] = dataclasses.field(default_factory=list)
|
| 251 |
-
model: str =
|
| 252 |
prompt_version: str = PROMPT_VERSION
|
| 253 |
doc_status: str = "ok" # "ok" | "scanned_no_text" | "empty_extraction"
|
| 254 |
# Gemini usage for the text call (0 when no call was made). tokens_out
|
| 255 |
-
# includes the model's thinking tokens, which are billed as output
|
|
|
|
| 256 |
tokens_in: int = 0
|
| 257 |
tokens_out: int = 0
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 258 |
|
| 259 |
|
| 260 |
def usage_from_response(resp: Any) -> tuple[int, int]:
|
|
@@ -318,7 +352,7 @@ class PropertyRow:
|
|
| 318 |
status: str = "ok"
|
| 319 |
flag_reason: str = ""
|
| 320 |
# bookkeeping
|
| 321 |
-
model: str =
|
| 322 |
prompt_version: str = PROMPT_VERSION
|
| 323 |
# provenance kind (figure-mining phase): 'text' (grounded in the PDF
|
| 324 |
# text) or 'figure' (read off a plot/table image — an estimate; see
|
|
@@ -609,8 +643,8 @@ def _upload_pdf_file(
|
|
| 609 |
uri = info.get("uri", uri)
|
| 610 |
|
| 611 |
|
| 612 |
-
def _parse_extraction_json(resp:
|
| 613 |
-
data = resp.json()
|
| 614 |
candidates = data.get("candidates", [])
|
| 615 |
if not candidates:
|
| 616 |
return None
|
|
@@ -697,18 +731,28 @@ def extract_from_pdf(pdf_bytes: bytes, filename: str, api_key: str) -> Extractio
|
|
| 697 |
"responseSchema": EXTRACTION_SCHEMA,
|
| 698 |
},
|
| 699 |
}
|
| 700 |
-
|
|
|
|
| 701 |
resp = gemini_request(url, payload)
|
| 702 |
if resp is None:
|
| 703 |
-
return Extraction(doc_status="empty_extraction")
|
| 704 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 705 |
raw = _parse_extraction_json(resp)
|
| 706 |
if not raw:
|
| 707 |
# The call was made and billed even though nothing usable came back.
|
| 708 |
-
return Extraction(doc_status="empty_extraction",
|
| 709 |
-
tokens_in=tokens_in, tokens_out=tokens_out
|
|
|
|
| 710 |
out = _extraction_from_json(raw)
|
| 711 |
-
out.
|
|
|
|
| 712 |
return out
|
| 713 |
|
| 714 |
|
|
|
|
| 64 |
# extraction call. Previews get retired with ~2 weeks notice; stable models
|
| 65 |
# don't. Overridable via env so a future swap needs no code change.
|
| 66 |
GEMINI_MODEL = os.environ.get("GEMINI_MODEL", "gemini-3.5-flash")
|
| 67 |
+
# Model for the figure (vision) calls: classify + mining. Empty = same as
|
| 68 |
+
# GEMINI_MODEL. Classification is an easy task (one call per PDF picking a
|
| 69 |
+
# figure kind), so a cheaper model here costs little quality. The agent sets
|
| 70 |
+
# both from Run Control each cycle; read them at call time via text_model() /
|
| 71 |
+
# vision_model(), never through a name imported at module load.
|
| 72 |
+
GEMINI_VISION_MODEL = os.environ.get("GEMINI_VISION_MODEL", "").strip()
|
| 73 |
PROMPT_VERSION = "2.0"
|
| 74 |
|
| 75 |
+
|
| 76 |
+
def text_model() -> str:
|
| 77 |
+
"""The model the text extraction call uses right now."""
|
| 78 |
+
return (GEMINI_MODEL or "gemini-3.5-flash").strip()
|
| 79 |
+
|
| 80 |
+
|
| 81 |
+
def vision_model() -> str:
|
| 82 |
+
"""The model the figure classify / mining calls use right now."""
|
| 83 |
+
return (GEMINI_VISION_MODEL or "").strip() or text_model()
|
| 84 |
+
|
| 85 |
# How much the model may "think" per call. Thinking tokens are billed as
|
| 86 |
# output (8x the input price): on 28 Sep 2026 a 25-paper cycle cost $2.02,
|
| 87 |
# $1.86 of it output tokens, ~30k per PDF where the extraction JSON itself
|
|
|
|
| 264 |
@dataclasses.dataclass
|
| 265 |
class Extraction:
|
| 266 |
materials: list[Material] = dataclasses.field(default_factory=list)
|
| 267 |
+
model: str = dataclasses.field(default_factory=text_model)
|
| 268 |
prompt_version: str = PROMPT_VERSION
|
| 269 |
doc_status: str = "ok" # "ok" | "scanned_no_text" | "empty_extraction"
|
| 270 |
# Gemini usage for the text call (0 when no call was made). tokens_out
|
| 271 |
+
# includes the model's thinking tokens, which are billed as output;
|
| 272 |
+
# tokens_thinking is that thinking share on its own (for cost analysis).
|
| 273 |
tokens_in: int = 0
|
| 274 |
tokens_out: int = 0
|
| 275 |
+
tokens_thinking: int = 0
|
| 276 |
+
|
| 277 |
+
|
| 278 |
+
def usage_detail(resp: Any) -> tuple[int, int, int]:
|
| 279 |
+
"""(tokens_in, tokens_out, tokens_thinking) from a generateContent
|
| 280 |
+
response's ``usageMetadata``. tokens_out is billed output (candidates +
|
| 281 |
+
thinking); tokens_thinking is the thinking part alone. (0, 0, 0) when
|
| 282 |
+
absent or unparseable; never raises."""
|
| 283 |
+
try:
|
| 284 |
+
data = resp.json() if hasattr(resp, "json") else (resp or {})
|
| 285 |
+
u = (data or {}).get("usageMetadata") or {}
|
| 286 |
+
tin = int(u.get("promptTokenCount") or 0)
|
| 287 |
+
think = int(u.get("thoughtsTokenCount") or 0)
|
| 288 |
+
tout = int(u.get("candidatesTokenCount") or 0) + think
|
| 289 |
+
return max(0, tin), max(0, tout), max(0, think)
|
| 290 |
+
except Exception:
|
| 291 |
+
return 0, 0, 0
|
| 292 |
|
| 293 |
|
| 294 |
def usage_from_response(resp: Any) -> tuple[int, int]:
|
|
|
|
| 352 |
status: str = "ok"
|
| 353 |
flag_reason: str = ""
|
| 354 |
# bookkeeping
|
| 355 |
+
model: str = dataclasses.field(default_factory=text_model)
|
| 356 |
prompt_version: str = PROMPT_VERSION
|
| 357 |
# provenance kind (figure-mining phase): 'text' (grounded in the PDF
|
| 358 |
# text) or 'figure' (read off a plot/table image — an estimate; see
|
|
|
|
| 643 |
uri = info.get("uri", uri)
|
| 644 |
|
| 645 |
|
| 646 |
+
def _parse_extraction_json(resp: Any) -> Optional[dict[str, Any]]:
|
| 647 |
+
data = resp.json() if hasattr(resp, "json") else (resp or {})
|
| 648 |
candidates = data.get("candidates", [])
|
| 649 |
if not candidates:
|
| 650 |
return None
|
|
|
|
| 731 |
"responseSchema": EXTRACTION_SCHEMA,
|
| 732 |
},
|
| 733 |
}
|
| 734 |
+
model = text_model()
|
| 735 |
+
url = GEMINI_URL_TEMPLATE.format(model=model, key=api_key)
|
| 736 |
resp = gemini_request(url, payload)
|
| 737 |
if resp is None:
|
| 738 |
+
return Extraction(doc_status="empty_extraction", model=model)
|
| 739 |
+
return extraction_from_response(resp, model)
|
| 740 |
+
|
| 741 |
+
|
| 742 |
+
def extraction_from_response(resp: Any, model: Optional[str] = None) -> Extraction:
|
| 743 |
+
"""Extraction (with usage) from a generateContent response object or its
|
| 744 |
+
parsed JSON dict (the Batch API hands back dicts)."""
|
| 745 |
+
model = model or text_model()
|
| 746 |
+
tokens_in, tokens_out, tokens_thinking = usage_detail(resp)
|
| 747 |
raw = _parse_extraction_json(resp)
|
| 748 |
if not raw:
|
| 749 |
# The call was made and billed even though nothing usable came back.
|
| 750 |
+
return Extraction(doc_status="empty_extraction", model=model,
|
| 751 |
+
tokens_in=tokens_in, tokens_out=tokens_out,
|
| 752 |
+
tokens_thinking=tokens_thinking)
|
| 753 |
out = _extraction_from_json(raw)
|
| 754 |
+
out.model = model
|
| 755 |
+
out.tokens_in, out.tokens_out, out.tokens_thinking = tokens_in, tokens_out, tokens_thinking
|
| 756 |
return out
|
| 757 |
|
| 758 |
|
|
@@ -152,7 +152,7 @@ class Figure:
|
|
| 152 |
material_key: str = ""
|
| 153 |
mining_status: str = "" # "" | not_mined | skipped_kind | mined | mining_failed | classify_failed
|
| 154 |
n_values: int = 0
|
| 155 |
-
model: str =
|
| 156 |
figure_prompt_version: str = FIGURE_PROMPT_VERSION
|
| 157 |
# not persisted: raw PNG bytes (on-disk artifact) + the JPEG copy sent to Gemini
|
| 158 |
png_bytes: bytes = dataclasses.field(default=b"", repr=False, compare=False)
|
|
@@ -200,6 +200,8 @@ class VisionStats:
|
|
| 200 |
incomplete_calls: int = 0 # classify responses that skipped some indices
|
| 201 |
tokens_in: int = 0 # Gemini usageMetadata over all vision calls
|
| 202 |
tokens_out: int = 0 # (candidates + thinking tokens)
|
|
|
|
|
|
|
| 203 |
|
| 204 |
@property
|
| 205 |
def total(self) -> int:
|
|
@@ -212,14 +214,17 @@ class VisionStats:
|
|
| 212 |
self.incomplete_calls += other.incomplete_calls
|
| 213 |
self.tokens_in += other.tokens_in
|
| 214 |
self.tokens_out += other.tokens_out
|
|
|
|
|
|
|
| 215 |
|
| 216 |
def record_usage(self, resp: Any) -> None:
|
| 217 |
"""Add a response's usageMetadata (no-op for None / unparseable)."""
|
| 218 |
if resp is None:
|
| 219 |
return
|
| 220 |
-
tin, tout = extraction.
|
| 221 |
self.tokens_in += tin
|
| 222 |
self.tokens_out += tout
|
|
|
|
| 223 |
|
| 224 |
|
| 225 |
@dataclasses.dataclass
|
|
@@ -724,7 +729,7 @@ def classify_figures(
|
|
| 724 |
text_materials: list[str],
|
| 725 |
api_key: str,
|
| 726 |
*,
|
| 727 |
-
model: str =
|
| 728 |
stats: Optional[VisionStats] = None,
|
| 729 |
) -> list[Figure]:
|
| 730 |
"""One Gemini vision call classifies all of a PDF's figures.
|
|
@@ -734,6 +739,7 @@ def classify_figures(
|
|
| 734 |
and ``figure_kind='other'`` (nothing raises; the caller counts it).
|
| 735 |
"""
|
| 736 |
vs = stats if stats is not None else VisionStats()
|
|
|
|
| 737 |
if not figures:
|
| 738 |
return figures
|
| 739 |
for batch in _batches_by_size(figures, CLASSIFY_MAX_BYTES):
|
|
@@ -883,7 +889,7 @@ def mine_figure(
|
|
| 883 |
text_materials: list[str],
|
| 884 |
api_key: str,
|
| 885 |
*,
|
| 886 |
-
model: str =
|
| 887 |
stats: Optional[VisionStats] = None,
|
| 888 |
) -> Optional[MinedFigure]:
|
| 889 |
"""One vision call: structured readout of a property_plot / table_image.
|
|
@@ -892,6 +898,7 @@ def mine_figure(
|
|
| 892 |
sets ``'mined'`` and ``fig.n_values`` on success.
|
| 893 |
"""
|
| 894 |
vs = stats if stats is not None else VisionStats()
|
|
|
|
| 895 |
mats = ", ".join(f"'{m}'" for m in dict.fromkeys(text_materials) if m) or "(none)"
|
| 896 |
prompt = MINING_PROMPT.format(
|
| 897 |
caption=fig.caption or "(no caption)", materials=mats,
|
|
@@ -1093,7 +1100,7 @@ def figure_properties_to_rows(
|
|
| 1093 |
mat.properties.append(prop)
|
| 1094 |
fig_of_prop[id(prop)] = fig
|
| 1095 |
|
| 1096 |
-
ext = Extraction(materials=list(buckets.values()), model=
|
| 1097 |
prompt_version=FIGURE_PROMPT_VERSION)
|
| 1098 |
rows = to_rows(ext, source_pdf, source_sha1)
|
| 1099 |
# to_rows iterates materials x properties in order; walk the same order to
|
|
@@ -1104,7 +1111,7 @@ def figure_properties_to_rows(
|
|
| 1104 |
fig = fig_of_prop[id(prop)]
|
| 1105 |
row.origin = "figure"
|
| 1106 |
row.figure_id = fig.figure_id
|
| 1107 |
-
row.model = fig.model or
|
| 1108 |
row.prompt_version = fig.figure_prompt_version
|
| 1109 |
row.source_quote = fig.caption or f"[{fig.label}, no caption]"
|
| 1110 |
row.page = fig.page
|
|
@@ -1137,7 +1144,7 @@ def run_figure_stage(
|
|
| 1137 |
out_dir: Path | str = DEFAULT_FIGURES_DIR,
|
| 1138 |
max_figures: int = DEFAULT_MAX_FIGURES,
|
| 1139 |
mine: bool = True,
|
| 1140 |
-
model: str =
|
| 1141 |
done_figure_ids: Optional[set[str]] = None,
|
| 1142 |
) -> FigureStageResult:
|
| 1143 |
"""harvest -> classify -> mine -> rows. Never raises; errors are recorded
|
|
@@ -1152,6 +1159,8 @@ def run_figure_stage(
|
|
| 1152 |
to overwrite their stored status.
|
| 1153 |
"""
|
| 1154 |
hs, vs = HarvestStats(), VisionStats()
|
|
|
|
|
|
|
| 1155 |
res = FigureStageResult(figures=[], rows=[], harvest=hs, vision=vs)
|
| 1156 |
done = done_figure_ids or set()
|
| 1157 |
try:
|
|
|
|
| 152 |
material_key: str = ""
|
| 153 |
mining_status: str = "" # "" | not_mined | skipped_kind | mined | mining_failed | classify_failed
|
| 154 |
n_values: int = 0
|
| 155 |
+
model: str = dataclasses.field(default_factory=extraction.vision_model)
|
| 156 |
figure_prompt_version: str = FIGURE_PROMPT_VERSION
|
| 157 |
# not persisted: raw PNG bytes (on-disk artifact) + the JPEG copy sent to Gemini
|
| 158 |
png_bytes: bytes = dataclasses.field(default=b"", repr=False, compare=False)
|
|
|
|
| 200 |
incomplete_calls: int = 0 # classify responses that skipped some indices
|
| 201 |
tokens_in: int = 0 # Gemini usageMetadata over all vision calls
|
| 202 |
tokens_out: int = 0 # (candidates + thinking tokens)
|
| 203 |
+
tokens_thinking: int = 0 # the thinking share of tokens_out
|
| 204 |
+
model: str = "" # the model the calls went to (for pricing)
|
| 205 |
|
| 206 |
@property
|
| 207 |
def total(self) -> int:
|
|
|
|
| 214 |
self.incomplete_calls += other.incomplete_calls
|
| 215 |
self.tokens_in += other.tokens_in
|
| 216 |
self.tokens_out += other.tokens_out
|
| 217 |
+
self.tokens_thinking += other.tokens_thinking
|
| 218 |
+
self.model = self.model or other.model
|
| 219 |
|
| 220 |
def record_usage(self, resp: Any) -> None:
|
| 221 |
"""Add a response's usageMetadata (no-op for None / unparseable)."""
|
| 222 |
if resp is None:
|
| 223 |
return
|
| 224 |
+
tin, tout, think = extraction.usage_detail(resp)
|
| 225 |
self.tokens_in += tin
|
| 226 |
self.tokens_out += tout
|
| 227 |
+
self.tokens_thinking += think
|
| 228 |
|
| 229 |
|
| 230 |
@dataclasses.dataclass
|
|
|
|
| 729 |
text_materials: list[str],
|
| 730 |
api_key: str,
|
| 731 |
*,
|
| 732 |
+
model: Optional[str] = None,
|
| 733 |
stats: Optional[VisionStats] = None,
|
| 734 |
) -> list[Figure]:
|
| 735 |
"""One Gemini vision call classifies all of a PDF's figures.
|
|
|
|
| 739 |
and ``figure_kind='other'`` (nothing raises; the caller counts it).
|
| 740 |
"""
|
| 741 |
vs = stats if stats is not None else VisionStats()
|
| 742 |
+
model = model or extraction.vision_model()
|
| 743 |
if not figures:
|
| 744 |
return figures
|
| 745 |
for batch in _batches_by_size(figures, CLASSIFY_MAX_BYTES):
|
|
|
|
| 889 |
text_materials: list[str],
|
| 890 |
api_key: str,
|
| 891 |
*,
|
| 892 |
+
model: Optional[str] = None,
|
| 893 |
stats: Optional[VisionStats] = None,
|
| 894 |
) -> Optional[MinedFigure]:
|
| 895 |
"""One vision call: structured readout of a property_plot / table_image.
|
|
|
|
| 898 |
sets ``'mined'`` and ``fig.n_values`` on success.
|
| 899 |
"""
|
| 900 |
vs = stats if stats is not None else VisionStats()
|
| 901 |
+
model = model or extraction.vision_model()
|
| 902 |
mats = ", ".join(f"'{m}'" for m in dict.fromkeys(text_materials) if m) or "(none)"
|
| 903 |
prompt = MINING_PROMPT.format(
|
| 904 |
caption=fig.caption or "(no caption)", materials=mats,
|
|
|
|
| 1100 |
mat.properties.append(prop)
|
| 1101 |
fig_of_prop[id(prop)] = fig
|
| 1102 |
|
| 1103 |
+
ext = Extraction(materials=list(buckets.values()), model=extraction.vision_model(),
|
| 1104 |
prompt_version=FIGURE_PROMPT_VERSION)
|
| 1105 |
rows = to_rows(ext, source_pdf, source_sha1)
|
| 1106 |
# to_rows iterates materials x properties in order; walk the same order to
|
|
|
|
| 1111 |
fig = fig_of_prop[id(prop)]
|
| 1112 |
row.origin = "figure"
|
| 1113 |
row.figure_id = fig.figure_id
|
| 1114 |
+
row.model = fig.model or extraction.vision_model()
|
| 1115 |
row.prompt_version = fig.figure_prompt_version
|
| 1116 |
row.source_quote = fig.caption or f"[{fig.label}, no caption]"
|
| 1117 |
row.page = fig.page
|
|
|
|
| 1144 |
out_dir: Path | str = DEFAULT_FIGURES_DIR,
|
| 1145 |
max_figures: int = DEFAULT_MAX_FIGURES,
|
| 1146 |
mine: bool = True,
|
| 1147 |
+
model: Optional[str] = None,
|
| 1148 |
done_figure_ids: Optional[set[str]] = None,
|
| 1149 |
) -> FigureStageResult:
|
| 1150 |
"""harvest -> classify -> mine -> rows. Never raises; errors are recorded
|
|
|
|
| 1159 |
to overwrite their stored status.
|
| 1160 |
"""
|
| 1161 |
hs, vs = HarvestStats(), VisionStats()
|
| 1162 |
+
model = model or extraction.vision_model()
|
| 1163 |
+
vs.model = model
|
| 1164 |
res = FigureStageResult(figures=[], rows=[], harvest=hs, vision=vs)
|
| 1165 |
done = done_figure_ids or set()
|
| 1166 |
try:
|
|
@@ -134,7 +134,8 @@ def _history_block() -> pd.DataFrame:
|
|
| 134 |
"candidates, relevance_rejected, skipped_doi, url_seen_skipped, "
|
| 135 |
"download_failed, queued, downloaded, pdfs_ingested, rows_inserted, "
|
| 136 |
"rows_flagged, duplicates, errors, figures_found, figures_mined, "
|
| 137 |
-
"figure_rows, vision_calls, rows_linked, tokens_in, tokens_out, "
|
|
|
|
| 138 |
"heartbeat_at, report "
|
| 139 |
"FROM agent_runs ORDER BY id DESC LIMIT 50")
|
| 140 |
|
|
@@ -189,13 +190,31 @@ def _history_block() -> pd.DataFrame:
|
|
| 189 |
"skipped_doi", "url_seen_skipped", "download_failed", "queued", "downloaded",
|
| 190 |
"pdfs_ingested", "rows_inserted", "rows_flagged", "duplicates",
|
| 191 |
"errors", "figures_found", "figures_mined", "figure_rows",
|
| 192 |
-
"vision_calls", "rows_linked", "tokens_in", "tokens_out"
|
|
|
|
| 193 |
if c in disp.columns:
|
| 194 |
disp[c] = _to_int64(disp[c])
|
| 195 |
-
#
|
|
|
|
|
|
|
|
|
|
| 196 |
disp["cost_usd"] = [
|
| 197 |
-
|
| 198 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 199 |
disp["trigger"] = disp["trigger"].map(_s)
|
| 200 |
disp["status"] = disp["status"].map(_s)
|
| 201 |
|
|
@@ -264,10 +283,22 @@ def _history_block() -> pd.DataFrame:
|
|
| 264 |
"tokens_out": st.column_config.NumberColumn(
|
| 265 |
"Tokens out", format="localized",
|
| 266 |
help="Gemini output tokens incl. thinking (text + vision)"),
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 267 |
"cost_usd": st.column_config.NumberColumn(
|
| 268 |
"Cost", format="$%.4f",
|
| 269 |
-
help=
|
| 270 |
-
|
|
|
|
|
|
|
|
|
|
| 271 |
},
|
| 272 |
hide_index=True, width='stretch',
|
| 273 |
height=min(35 * 11 + 3, 35 * (len(disp) + 1) + 3))
|
|
|
|
| 134 |
"candidates, relevance_rejected, skipped_doi, url_seen_skipped, "
|
| 135 |
"download_failed, queued, downloaded, pdfs_ingested, rows_inserted, "
|
| 136 |
"rows_flagged, duplicates, errors, figures_found, figures_mined, "
|
| 137 |
+
"figure_rows, vision_calls, rows_linked, duplicate_text, tokens_in, tokens_out, "
|
| 138 |
+
"tokens_thinking, vision_tokens_in, vision_tokens_out, text_model, cost_usd, "
|
| 139 |
"heartbeat_at, report "
|
| 140 |
"FROM agent_runs ORDER BY id DESC LIMIT 50")
|
| 141 |
|
|
|
|
| 190 |
"skipped_doi", "url_seen_skipped", "download_failed", "queued", "downloaded",
|
| 191 |
"pdfs_ingested", "rows_inserted", "rows_flagged", "duplicates",
|
| 192 |
"errors", "figures_found", "figures_mined", "figure_rows",
|
| 193 |
+
"vision_calls", "rows_linked", "duplicate_text", "tokens_in", "tokens_out",
|
| 194 |
+
"tokens_thinking"):
|
| 195 |
if c in disp.columns:
|
| 196 |
disp[c] = _to_int64(disp[c])
|
| 197 |
+
# Cost: stored per run since 30 Sep 2026 (each model at its own
|
| 198 |
+
# price); older runs re-priced from their tokens at the model they
|
| 199 |
+
# ran on (gemini-3.5-flash) — the old derived column used Gemini
|
| 200 |
+
# 2.5 Flash prices and showed about a quarter of the real bill.
|
| 201 |
disp["cost_usd"] = [
|
| 202 |
+
(float(c) if pd.notna(c) else
|
| 203 |
+
(C.cost_usd(i, o, C.LEGACY_RUN_MODEL) if (pd.notna(i) and pd.notna(o)) else None))
|
| 204 |
+
for c, i, o in zip(disp["cost_usd"], disp["tokens_in"], disp["tokens_out"])]
|
| 205 |
+
disp["cost_per_pdf"] = [
|
| 206 |
+
(c / n) if (c is not None and pd.notna(n) and int(n) > 0) else None
|
| 207 |
+
for c, n in zip(disp["cost_usd"], disp["pdfs_ingested"])]
|
| 208 |
+
disp["text_model"] = [
|
| 209 |
+
_s(mm) or (C.LEGACY_RUN_MODEL if (pd.notna(i) and int(i or 0) > 0) else "")
|
| 210 |
+
for mm, i in zip(disp["text_model"], disp["tokens_in"])]
|
| 211 |
+
disp = disp.drop(columns=["vision_tokens_in", "vision_tokens_out"])
|
| 212 |
+
_spent = sum(c for c in disp["cost_usd"] if c is not None)
|
| 213 |
+
_pdfs = int(pd.to_numeric(disp["pdfs_ingested"], errors="coerce").fillna(0).sum())
|
| 214 |
+
if _spent > 0:
|
| 215 |
+
st.caption(f"Gemini spend over these cycles: **${_spent:,.2f}** for {_pdfs:,} "
|
| 216 |
+
f"PDFs ingested" + (f" — **${_spent / _pdfs:.3f} per PDF**" if _pdfs else "")
|
| 217 |
+
+ ". Near-duplicates skipped before extraction cost nothing.")
|
| 218 |
disp["trigger"] = disp["trigger"].map(_s)
|
| 219 |
disp["status"] = disp["status"].map(_s)
|
| 220 |
|
|
|
|
| 283 |
"tokens_out": st.column_config.NumberColumn(
|
| 284 |
"Tokens out", format="localized",
|
| 285 |
help="Gemini output tokens incl. thinking (text + vision)"),
|
| 286 |
+
"tokens_thinking": st.column_config.NumberColumn(
|
| 287 |
+
"Thinking", format="localized",
|
| 288 |
+
help="The thinking share of Tokens out (billed as output)"),
|
| 289 |
+
"duplicate_text": st.column_config.NumberColumn(
|
| 290 |
+
"Near-dup", format="localized",
|
| 291 |
+
help="PDFs skipped before extraction: their text repeats a document "
|
| 292 |
+
"already ingested"),
|
| 293 |
+
"text_model": st.column_config.TextColumn(
|
| 294 |
+
"Model", width="small", help="Gemini model of the text extraction calls"),
|
| 295 |
"cost_usd": st.column_config.NumberColumn(
|
| 296 |
"Cost", format="$%.4f",
|
| 297 |
+
help="Gemini list price per model (text and vision calls each at their "
|
| 298 |
+
"own model; batch calls at half price). Runs before 30 Sep are "
|
| 299 |
+
"re-priced at gemini-3.5-flash ($1.50 / $9.00 per M in/out)."),
|
| 300 |
+
"cost_per_pdf": st.column_config.NumberColumn(
|
| 301 |
+
"$/PDF", format="$%.3f", help="Cost / PDFs ingested"),
|
| 302 |
},
|
| 303 |
hide_index=True, width='stretch',
|
| 304 |
height=min(35 * 11 + 3, 35 * (len(disp) + 1) + 3))
|
|
@@ -272,9 +272,25 @@ _THINK_OPTS = ["low", "minimal", "off", "dynamic", "high"]
|
|
| 272 |
_think_now = str(cfg.get("gemini_thinking", "low") or "dynamic").lower()
|
| 273 |
if _think_now not in _THINK_OPTS:
|
| 274 |
_THINK_OPTS.append(_think_now) # an explicit token budget set by env
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 275 |
with st.container(border=True):
|
| 276 |
-
T.card_title("Gemini cost", sub="thinking tokens are billed as output
|
| 277 |
-
|
|
|
|
| 278 |
with k1:
|
| 279 |
gemini_thinking = st.selectbox("Thinking per call", _THINK_OPTS,
|
| 280 |
index=_THINK_OPTS.index(_think_now),
|
|
@@ -287,9 +303,31 @@ with st.container(border=True):
|
|
| 287 |
"If the model rejects the option the call is repeated "
|
| 288 |
"without it.")
|
| 289 |
with k2:
|
| 290 |
-
st.
|
| 291 |
-
|
| 292 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 293 |
|
| 294 |
# --- figure mining ----------------------------------------------------------
|
| 295 |
with st.container(border=True):
|
|
@@ -349,6 +387,10 @@ with st.container(horizontal=True, vertical_alignment="center", gap="medium"):
|
|
| 349 |
"discovery_per_query": int(disc_per_query),
|
| 350 |
"figures_enabled": figures_on, "figure_mining": figure_mine,
|
| 351 |
"gemini_thinking": str(gemini_thinking),
|
|
|
|
|
|
|
|
|
|
|
|
|
| 352 |
"max_figures_per_pdf": int(max_figs),
|
| 353 |
"link_figures": link_on,
|
| 354 |
})
|
|
|
|
| 272 |
_think_now = str(cfg.get("gemini_thinking", "low") or "dynamic").lower()
|
| 273 |
if _think_now not in _THINK_OPTS:
|
| 274 |
_THINK_OPTS.append(_think_now) # an explicit token budget set by env
|
| 275 |
+
_MODEL_OPTS = list(C.MODEL_PRICES)
|
| 276 |
+
_model_now = str(cfg.get("gemini_model") or "gemini-3.5-flash").strip()
|
| 277 |
+
if _model_now not in _MODEL_OPTS:
|
| 278 |
+
_MODEL_OPTS.insert(0, _model_now)
|
| 279 |
+
_VMODEL_OPTS = ["(same as text model)"] + list(C.MODEL_PRICES)
|
| 280 |
+
_vmodel_now = str(cfg.get("gemini_vision_model") or "").strip() or "(same as text model)"
|
| 281 |
+
if _vmodel_now not in _VMODEL_OPTS:
|
| 282 |
+
_VMODEL_OPTS.append(_vmodel_now)
|
| 283 |
+
|
| 284 |
+
|
| 285 |
+
def _price_label(m: str) -> str:
|
| 286 |
+
pin, pout = C.model_price(m)
|
| 287 |
+
return f"{m} (${pin:g} / ${pout:g} per M)"
|
| 288 |
+
|
| 289 |
+
|
| 290 |
with st.container(border=True):
|
| 291 |
+
T.card_title("Gemini cost", sub="thinking tokens are billed as output (6× the input price "
|
| 292 |
+
"on gemini-3.5-flash)")
|
| 293 |
+
k1, k2, k3 = st.columns(3, vertical_alignment="bottom")
|
| 294 |
with k1:
|
| 295 |
gemini_thinking = st.selectbox("Thinking per call", _THINK_OPTS,
|
| 296 |
index=_THINK_OPTS.index(_think_now),
|
|
|
|
| 303 |
"If the model rejects the option the call is repeated "
|
| 304 |
"without it.")
|
| 305 |
with k2:
|
| 306 |
+
gemini_model = st.selectbox("Text extraction model", _MODEL_OPTS,
|
| 307 |
+
index=_MODEL_OPTS.index(_model_now),
|
| 308 |
+
format_func=_price_label,
|
| 309 |
+
disabled=not admin, key="rc_model",
|
| 310 |
+
help="Every row is stamped with this model. Change it only "
|
| 311 |
+
"after comparing on the gold set (cost_ab.py) — the "
|
| 312 |
+
"paper's numbers are per model.")
|
| 313 |
+
with k3:
|
| 314 |
+
gemini_vision_model = st.selectbox("Figure (vision) model", _VMODEL_OPTS,
|
| 315 |
+
index=_VMODEL_OPTS.index(_vmodel_now),
|
| 316 |
+
format_func=lambda m: m if m.startswith("(")
|
| 317 |
+
else _price_label(m),
|
| 318 |
+
disabled=not admin, key="rc_vmodel",
|
| 319 |
+
help="The one classify call per PDF (and mining, if "
|
| 320 |
+
"on). An easy task: the Lite model is ~5x cheaper.")
|
| 321 |
+
dedupe_sim = st.slider("Skip near-duplicate PDFs (text similarity ≥)", 0.0, 1.0,
|
| 322 |
+
float(cfg.get("dedupe_text_similarity", 0.9) or 0.0), 0.01,
|
| 323 |
+
disabled=not admin, key="rc_dedupe",
|
| 324 |
+
help="Before any Gemini call, a PDF whose text is this similar to one "
|
| 325 |
+
"already ingested (regional datasheet variants, preprint vs "
|
| 326 |
+
"published version) is skipped. 0 = off.")
|
| 327 |
+
st.caption("Real cost of the 28–29 Sep cycles on gemini-3.5-flash with dynamic thinking and "
|
| 328 |
+
"mining on: ~$0.41 per PDF (36k in / 40k out tokens). Watch Cost and $/PDF in "
|
| 329 |
+
"Runs & Logs after a change; rows per paper and the flagged share tell you if "
|
| 330 |
+
"recall moved.")
|
| 331 |
|
| 332 |
# --- figure mining ----------------------------------------------------------
|
| 333 |
with st.container(border=True):
|
|
|
|
| 387 |
"discovery_per_query": int(disc_per_query),
|
| 388 |
"figures_enabled": figures_on, "figure_mining": figure_mine,
|
| 389 |
"gemini_thinking": str(gemini_thinking),
|
| 390 |
+
"gemini_model": str(gemini_model),
|
| 391 |
+
"gemini_vision_model": ("" if str(gemini_vision_model).startswith("(")
|
| 392 |
+
else str(gemini_vision_model)),
|
| 393 |
+
"dedupe_text_similarity": float(dedupe_sim),
|
| 394 |
"max_figures_per_pdf": int(max_figs),
|
| 395 |
"link_figures": link_on,
|
| 396 |
})
|
|
@@ -0,0 +1,326 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
"""Cost follow-up "f" (30 Sep 2026), scratch Postgres, no network:
|
| 2 |
+
|
| 3 |
+
* per-model pricing — the old cost column priced a gemini-3.5-flash agent
|
| 4 |
+
at Gemini 2.5 Flash rates (~1/4 of the bill); runs now store cost_usd
|
| 5 |
+
priced per model, text and vision separately;
|
| 6 |
+
* thinking / vision token breakdown on PdfResult and agent_runs;
|
| 7 |
+
* the text and vision models follow Run Control at call time;
|
| 8 |
+
* the near-duplicate gate skips a PDF whose text repeats an ingested one
|
| 9 |
+
before any Gemini call, and releases its claim when a PDF fails for good.
|
| 10 |
+
|
| 11 |
+
python test_cost_e2e.py (same DB_* env as test_e2e.py)
|
| 12 |
+
"""
|
| 13 |
+
import hashlib
|
| 14 |
+
import json
|
| 15 |
+
import os
|
| 16 |
+
import sys
|
| 17 |
+
|
| 18 |
+
os.environ.setdefault("AGENT_WORK_DIR", "/tmp/agent_cost_test")
|
| 19 |
+
os.environ["GEMINI_API_KEY"] = "test-key-not-used"
|
| 20 |
+
os.environ["AGENT_USE_NTRS"] = "0"
|
| 21 |
+
os.environ["AGENT_USE_S2"] = "0"
|
| 22 |
+
os.environ["AGENT_USE_OPENALEX_TOPICS"] = "0"
|
| 23 |
+
os.environ["AGENT_USE_DATASHEETS"] = "0"
|
| 24 |
+
os.environ["AGENT_SEED_QUERY_GRID"] = "0"
|
| 25 |
+
os.environ.pop("AGENT_PRICE_IN_PER_M", None)
|
| 26 |
+
os.environ.pop("AGENT_PRICE_OUT_PER_M", None)
|
| 27 |
+
|
| 28 |
+
import fitz # noqa: E402
|
| 29 |
+
|
| 30 |
+
import batch_ingest # noqa: E402
|
| 31 |
+
import extraction # noqa: E402
|
| 32 |
+
import figures # noqa: E402
|
| 33 |
+
import pdf_crawler # noqa: E402
|
| 34 |
+
from agent import agdb, config as C, orchestrator, textdedupe # noqa: E402
|
| 35 |
+
|
| 36 |
+
C.GEMINI_API_KEY = "test-key-not-used"
|
| 37 |
+
fails = []
|
| 38 |
+
|
| 39 |
+
|
| 40 |
+
def check(name, cond):
|
| 41 |
+
print(("PASS " if cond else "FAIL ") + name)
|
| 42 |
+
if not cond:
|
| 43 |
+
fails.append(name)
|
| 44 |
+
|
| 45 |
+
|
| 46 |
+
def close(a, b, tol=1e-6):
|
| 47 |
+
return abs(float(a) - float(b)) <= tol
|
| 48 |
+
|
| 49 |
+
|
| 50 |
+
class FakeResp:
|
| 51 |
+
def __init__(self, payload):
|
| 52 |
+
self._p = payload
|
| 53 |
+
self.status_code = 200
|
| 54 |
+
self.headers = {}
|
| 55 |
+
|
| 56 |
+
def json(self):
|
| 57 |
+
return self._p
|
| 58 |
+
|
| 59 |
+
|
| 60 |
+
# --- pricing ------------------------------------------------------------------
|
| 61 |
+
check("price: gemini-3.5-flash is $1.50 / $9.00", C.model_price("gemini-3.5-flash") == (1.50, 9.00))
|
| 62 |
+
check("price: gemini-3.5-flash-lite is $0.30 / $2.50",
|
| 63 |
+
C.model_price("gemini-3.5-flash-lite") == (0.30, 2.50))
|
| 64 |
+
check("price: 'models/' prefix and dated suffix resolve",
|
| 65 |
+
C.model_price("models/gemini-3.5-flash-001") == (1.50, 9.00)
|
| 66 |
+
and C.model_price("gemini-3.5-flash-lite-preview-06") == (0.30, 2.50))
|
| 67 |
+
check("price: unknown model falls back to the legacy run model",
|
| 68 |
+
C.model_price("some-future-model") == C.model_price(C.LEGACY_RUN_MODEL))
|
| 69 |
+
# runs 544-593 (29 Sep): 4.25M in / 4.68M out -> $48.5 real, not the $12.98 shown
|
| 70 |
+
check("cost: runs 544-593 re-priced at gemini-3.5-flash = $48.50",
|
| 71 |
+
close(C.cost_usd(4_250_000, 4_680_000), 48.495))
|
| 72 |
+
check("cost: batch halves the price", close(C.cost_usd(1e6, 1e6, "gemini-3.5-flash", batch=True), 5.25))
|
| 73 |
+
rc = C.run_cost_usd({"tokens_in": 1_000_000, "tokens_out": 200_000,
|
| 74 |
+
"vision_tokens_in": 400_000, "vision_tokens_out": 50_000,
|
| 75 |
+
"text_model": "gemini-3.5-flash", "vision_model": "gemini-3.5-flash-lite"})
|
| 76 |
+
# text: 600k in, 150k out at 3.5 flash = 0.9 + 1.35; vision 400k/50k at lite = 0.12 + 0.125
|
| 77 |
+
check("run cost: text and vision priced at their own models", close(rc, 0.9 + 1.35 + 0.12 + 0.125))
|
| 78 |
+
|
| 79 |
+
# --- usage breakdown ------------------------------------------------------------
|
| 80 |
+
resp = {"candidates": [{"content": {"parts": [{"text": "{}"}]}}],
|
| 81 |
+
"usageMetadata": {"promptTokenCount": 9000, "candidatesTokenCount": 3000,
|
| 82 |
+
"thoughtsTokenCount": 25000}}
|
| 83 |
+
check("usage_detail: (in, out incl. thinking, thinking)",
|
| 84 |
+
extraction.usage_detail(FakeResp(resp)) == (9000, 28000, 25000))
|
| 85 |
+
check("usage_detail: works on a plain dict (Batch API responses)",
|
| 86 |
+
extraction.usage_detail(resp) == (9000, 28000, 25000))
|
| 87 |
+
check("usage_detail: garbage -> zeros", extraction.usage_detail(None) == (0, 0, 0))
|
| 88 |
+
|
| 89 |
+
# --- models follow the runtime setting --------------------------------------------
|
| 90 |
+
doc = fitz.open()
|
| 91 |
+
page = doc.new_page()
|
| 92 |
+
page.insert_text((72, 100), "The tensile modulus of PEEK-CF30 was measured as 24.5 GPa. " * 3)
|
| 93 |
+
TINY = doc.tobytes()
|
| 94 |
+
doc.close()
|
| 95 |
+
GOOD = json.dumps({"materials": [{
|
| 96 |
+
"material_name": "PEEK CF30", "material_abbreviation": "PEEK-CF30",
|
| 97 |
+
"material_class": "Composite", "matrix": "PEEK", "fiber": "Carbon",
|
| 98 |
+
"properties": [{"section": "Mechanical", "property_name": "Tensile Modulus",
|
| 99 |
+
"value_raw": "24.5", "value_num": 24.5, "unit": "GPa",
|
| 100 |
+
"source_quote": "measured as 24.5 GPa", "page": 1}]}]})
|
| 101 |
+
URLS = []
|
| 102 |
+
_orig_req = extraction.gemini_request
|
| 103 |
+
|
| 104 |
+
|
| 105 |
+
def fake_req(url, payload, **kw):
|
| 106 |
+
URLS.append(url)
|
| 107 |
+
return FakeResp({"candidates": [{"content": {"parts": [{"text": GOOD}]}}],
|
| 108 |
+
"usageMetadata": {"promptTokenCount": 100, "candidatesTokenCount": 10,
|
| 109 |
+
"thoughtsTokenCount": 40}})
|
| 110 |
+
|
| 111 |
+
|
| 112 |
+
try:
|
| 113 |
+
extraction.gemini_request = fake_req
|
| 114 |
+
figures.gemini_request = fake_req
|
| 115 |
+
extraction.GEMINI_MODEL, extraction.GEMINI_VISION_MODEL = "gemini-3.5-flash-lite", ""
|
| 116 |
+
ex = extraction.extract_from_pdf(TINY, "t.pdf", "k")
|
| 117 |
+
check("text call goes to the runtime model", "models/gemini-3.5-flash-lite:" in URLS[-1])
|
| 118 |
+
check("Extraction stamped with the runtime model and thinking tokens",
|
| 119 |
+
ex.model == "gemini-3.5-flash-lite" and ex.tokens_thinking == 40 and ex.tokens_out == 50)
|
| 120 |
+
rows = extraction.to_rows(ex, "t.pdf", "abc")
|
| 121 |
+
check("rows carry the runtime model", rows and all(r.model == "gemini-3.5-flash-lite" for r in rows))
|
| 122 |
+
extraction.GEMINI_MODEL, extraction.GEMINI_VISION_MODEL = "gemini-3.5-flash", "gemini-3.5-flash-lite"
|
| 123 |
+
check("vision_model(): explicit vision model", extraction.vision_model() == "gemini-3.5-flash-lite")
|
| 124 |
+
vs = figures.VisionStats()
|
| 125 |
+
vs.record_usage(fake_req("x", {}))
|
| 126 |
+
check("VisionStats keeps thinking separately", (vs.tokens_out, vs.tokens_thinking) == (50, 40))
|
| 127 |
+
extraction.GEMINI_VISION_MODEL = ""
|
| 128 |
+
check("vision_model(): empty = the text model", extraction.vision_model() == "gemini-3.5-flash")
|
| 129 |
+
finally:
|
| 130 |
+
extraction.gemini_request = _orig_req
|
| 131 |
+
figures.gemini_request = _orig_req
|
| 132 |
+
extraction.GEMINI_MODEL, extraction.GEMINI_VISION_MODEL = "gemini-3.5-flash", ""
|
| 133 |
+
|
| 134 |
+
# figure classification URL: run classify on a stub figure list via the batch helper
|
| 135 |
+
URLS.clear()
|
| 136 |
+
try:
|
| 137 |
+
figures.gemini_request = fake_req
|
| 138 |
+
extraction.GEMINI_VISION_MODEL = "gemini-3.5-flash-lite"
|
| 139 |
+
f1 = figures.Figure(figure_id="f1", source_pdf="t.pdf", source_sha1="abc", page=1,
|
| 140 |
+
bbox=(0, 0, 100, 100), caption="Fig. 1 Tensile modulus", image_path="",
|
| 141 |
+
image_sha256="", width_px=200, height_px=200, route="raster",
|
| 142 |
+
jpeg_bytes=b"\xff\xd8fakejpeg")
|
| 143 |
+
check("Figure.model defaults to the runtime vision model", f1.model == "gemini-3.5-flash-lite")
|
| 144 |
+
figures.classify_figures([f1], ["PEEK"], "k", stats=figures.VisionStats())
|
| 145 |
+
check("figure classify call goes to the vision model",
|
| 146 |
+
bool(URLS) and "models/gemini-3.5-flash-lite:" in URLS[-1])
|
| 147 |
+
finally:
|
| 148 |
+
figures.gemini_request = _orig_req
|
| 149 |
+
extraction.GEMINI_VISION_MODEL = ""
|
| 150 |
+
|
| 151 |
+
# --- MinHash fingerprints -----------------------------------------------------------
|
| 152 |
+
BASE = (" ".join(f"Hexcel HexPly 8552 epoxy matrix product data sheet section {i}: "
|
| 153 |
+
f"tensile strength {1500 + i} MPa at 23 C, compression modulus {60 + i} GPa, "
|
| 154 |
+
f"interlaminar shear {90 + i} MPa after cure 180 C for 120 minutes."
|
| 155 |
+
for i in range(25)))
|
| 156 |
+
VARIANT = BASE.replace("section 0:", "section 0 (EU version, Hexcel Composites GmbH):") + \
|
| 157 |
+
" Copyright Hexcel Corporation. EU edition."
|
| 158 |
+
OTHER = BASE.replace("8552", "8551")
|
| 159 |
+
for n in range(25): # a sibling datasheet: same template, different numbers
|
| 160 |
+
OTHER = OTHER.replace(f"{1500 + n} MPa", f"{1300 + 3 * n} MPa").replace(f"{60 + n} GPa", f"{50 + 2 * n} GPa")
|
| 161 |
+
fa, fb, fc = textdedupe.fingerprint(BASE), textdedupe.fingerprint(VARIANT), textdedupe.fingerprint(OTHER)
|
| 162 |
+
check("fingerprint: identical text -> similarity 1.0", textdedupe.similarity(fa[0], fa[0]) == 1.0)
|
| 163 |
+
sab = textdedupe.similarity(fa[0], fb[0])
|
| 164 |
+
sac = textdedupe.similarity(fa[0], fc[0])
|
| 165 |
+
print(f" regional variant ~ {sab:.2f}, sibling datasheet ~ {sac:.2f}")
|
| 166 |
+
check("fingerprint: regional variant >= 0.9", sab >= 0.9)
|
| 167 |
+
check("fingerprint: sibling datasheet (same template, other numbers) < 0.9", sac < 0.9)
|
| 168 |
+
check("fingerprint: too little text -> None", textdedupe.fingerprint("tensile modulus 24 GPa") is None)
|
| 169 |
+
|
| 170 |
+
|
| 171 |
+
# --- E2E cycle: models from config, breakdown + cost stored, duplicate skipped ------
|
| 172 |
+
def make_pdf(text: str) -> bytes:
|
| 173 |
+
d = fitz.open()
|
| 174 |
+
words = text.split()
|
| 175 |
+
for start in range(0, len(words), 350):
|
| 176 |
+
p = d.new_page()
|
| 177 |
+
chunk = words[start:start + 350]
|
| 178 |
+
y = 60
|
| 179 |
+
for i in range(0, len(chunk), 12):
|
| 180 |
+
p.insert_text((40, y), " ".join(chunk[i:i + 12]), fontsize=8)
|
| 181 |
+
y += 11
|
| 182 |
+
data = d.tobytes()
|
| 183 |
+
d.close()
|
| 184 |
+
return data
|
| 185 |
+
|
| 186 |
+
|
| 187 |
+
PDFS = {"base": make_pdf(BASE), "variant": make_pdf(VARIANT), "other": make_pdf(OTHER)}
|
| 188 |
+
ORDER = ["base", "variant", "other"]
|
| 189 |
+
EXTRACTED = []
|
| 190 |
+
|
| 191 |
+
|
| 192 |
+
def fake_search_openalex(query, limit):
|
| 193 |
+
for k in ORDER:
|
| 194 |
+
yield pdf_crawler.Candidate(title=f"{k} hexply 8552 datasheet", pdf_url=f"https://example.org/{k}.pdf",
|
| 195 |
+
source="openalex", query=query, doi=f"10.9999/aim.cost.{k}", year="2026",
|
| 196 |
+
abstract="thermoplastic composite tensile modulus carbon fiber datasheet")
|
| 197 |
+
|
| 198 |
+
|
| 199 |
+
def fake_download(cand, pdf_dir, state):
|
| 200 |
+
if cand.pdf_url in state.seen_urls:
|
| 201 |
+
return None
|
| 202 |
+
state.seen_urls.add(cand.pdf_url)
|
| 203 |
+
key = cand.pdf_url.rsplit("/", 1)[1][:-4]
|
| 204 |
+
data = PDFS[key]
|
| 205 |
+
sha = hashlib.sha256(data).hexdigest()
|
| 206 |
+
fname = f"openalex_{key}_{sha[:8]}.pdf"
|
| 207 |
+
(pdf_dir / fname).write_bytes(data)
|
| 208 |
+
return {"filename": fname, "title": cand.title, "doi": cand.doi, "url": cand.pdf_url,
|
| 209 |
+
"year": cand.year, "source": cand.source, "sha256": sha, "query": cand.query, "bytes": len(data)}
|
| 210 |
+
|
| 211 |
+
|
| 212 |
+
FAIL_KEYS = set()
|
| 213 |
+
|
| 214 |
+
|
| 215 |
+
def fake_extract(pdf_bytes, filename, api_key):
|
| 216 |
+
key = next(k for k in ORDER if f"_{k}_" in filename)
|
| 217 |
+
EXTRACTED.append(key)
|
| 218 |
+
if key in FAIL_KEYS:
|
| 219 |
+
return extraction.Extraction(doc_status="empty_extraction", tokens_in=500, tokens_out=100,
|
| 220 |
+
tokens_thinking=60)
|
| 221 |
+
val = "1500" if key != "other" else "1300"
|
| 222 |
+
mat = extraction.Material(
|
| 223 |
+
material_name=f"HexPly {key}", material_abbreviation=f"HP-{key}", material_class="Composite",
|
| 224 |
+
matrix="Epoxy", fiber="Carbon", fiber_volume_fraction="",
|
| 225 |
+
properties=[extraction.Property(section="Mechanical", property_name="Tensile Strength",
|
| 226 |
+
value_raw=val, value_num=float(val), unit="MPa",
|
| 227 |
+
test_condition="23 C", source_quote=f"tensile strength {val} MPa",
|
| 228 |
+
page=1)])
|
| 229 |
+
ext = extraction.Extraction(materials=[mat], doc_status="ok")
|
| 230 |
+
ext.tokens_in, ext.tokens_out, ext.tokens_thinking = 30000, 8000, 5000
|
| 231 |
+
return ext
|
| 232 |
+
|
| 233 |
+
|
| 234 |
+
pdf_crawler.search_openalex = fake_search_openalex
|
| 235 |
+
pdf_crawler.search_arxiv = lambda q, n: iter(())
|
| 236 |
+
pdf_crawler.download_pdf = fake_download
|
| 237 |
+
batch_ingest.extract_from_pdf = fake_extract
|
| 238 |
+
|
| 239 |
+
|
| 240 |
+
def q(sql, params=()):
|
| 241 |
+
conn = agdb.connect()
|
| 242 |
+
try:
|
| 243 |
+
with conn.cursor() as cur:
|
| 244 |
+
cur.execute(sql, params) if params else cur.execute(sql)
|
| 245 |
+
try:
|
| 246 |
+
return cur.fetchall()
|
| 247 |
+
except Exception:
|
| 248 |
+
return None
|
| 249 |
+
finally:
|
| 250 |
+
conn.commit()
|
| 251 |
+
conn.close()
|
| 252 |
+
|
| 253 |
+
|
| 254 |
+
orchestrator.bootstrap()
|
| 255 |
+
cfg_now = agdb.get_config(agdb.connect())
|
| 256 |
+
check("defaults: text model gemini-3.5-flash, vision = same, dedupe 0.9",
|
| 257 |
+
cfg_now["gemini_model"] == "gemini-3.5-flash" and cfg_now["gemini_vision_model"] == ""
|
| 258 |
+
and close(cfg_now["dedupe_text_similarity"], 0.9))
|
| 259 |
+
conn = agdb.connect()
|
| 260 |
+
agdb.set_config(conn, {"queries_per_cycle": 1, "max_rounds_per_cycle": 1, "max_new_pdfs_per_cycle": 5,
|
| 261 |
+
"min_new_pdfs_per_cycle": 5, "figures_enabled": False, "link_figures": False,
|
| 262 |
+
"expand_queries": False, "ingest_workers": 1, "use_candidate_queue": False,
|
| 263 |
+
"gemini_model": "gemini-3.5-flash-lite",
|
| 264 |
+
"gemini_vision_model": "gemini-3.5-flash-lite"})
|
| 265 |
+
conn.close()
|
| 266 |
+
|
| 267 |
+
orchestrator.run_cycle(trigger="test")
|
| 268 |
+
check("node_plan applied the configured models",
|
| 269 |
+
extraction.GEMINI_MODEL == "gemini-3.5-flash-lite"
|
| 270 |
+
and extraction.GEMINI_VISION_MODEL == "gemini-3.5-flash-lite")
|
| 271 |
+
row = q("SELECT status, downloaded, pdfs_ingested, duplicate_text, tokens_in, tokens_out, "
|
| 272 |
+
"tokens_thinking, text_model, cost_usd FROM agent_runs WHERE kind = 'cycle' "
|
| 273 |
+
"ORDER BY id DESC LIMIT 1")[0]
|
| 274 |
+
print(" run:", row)
|
| 275 |
+
check("cycle: 3 downloaded, 2 ingested, 1 near-duplicate skipped", row[:4] == ("ok", 3, 2, 1))
|
| 276 |
+
check("the duplicate never reached Gemini", sorted(EXTRACTED) == ["base", "other"])
|
| 277 |
+
check("tokens and thinking summed over the 2 extracted PDFs", row[4:7] == (60000, 16000, 10000))
|
| 278 |
+
check("run stamped with the text model", row[7] == "gemini-3.5-flash-lite")
|
| 279 |
+
check("cost stored, priced at the lite model", close(row[8], C.cost_usd(60000, 16000, "gemini-3.5-flash-lite")))
|
| 280 |
+
reg = q("SELECT ingest_status FROM agent_doi_seen WHERE filename LIKE 'openalex_variant_%'")
|
| 281 |
+
check("registry: the variant marked duplicate_text", reg == [("duplicate_text",)])
|
| 282 |
+
check("the variant PDF was removed from disk", not list(C.PDF_DIR.glob("openalex_variant_*.pdf")))
|
| 283 |
+
fps = q("SELECT count(*) FROM agent_text_fingerprints")[0][0]
|
| 284 |
+
check("fingerprints: the 2 ingested documents are claimed", fps == 2)
|
| 285 |
+
ev = q("SELECT count(*) FROM agent_events WHERE message LIKE '%skipped before extraction%'")[0][0]
|
| 286 |
+
check("event explains the skip", ev == 1)
|
| 287 |
+
done = q("SELECT message FROM agent_events WHERE message LIKE 'cycle done:%' ORDER BY id DESC LIMIT 1")[0][0]
|
| 288 |
+
check("cycle-done event names thinking tokens and the model price",
|
| 289 |
+
"10000 of it thinking" in done and "gemini-3.5-flash-lite prices" in done
|
| 290 |
+
and "1 near-duplicate" in done)
|
| 291 |
+
|
| 292 |
+
# A PDF that fails for good releases its claim, so a variant may stand in.
|
| 293 |
+
EXTRACTED.clear()
|
| 294 |
+
q("DELETE FROM agent_doi_seen"); q('DELETE FROM "Composites_materials"'); q("DELETE FROM sources")
|
| 295 |
+
q("DELETE FROM agent_text_fingerprints")
|
| 296 |
+
conn = agdb.connect()
|
| 297 |
+
agdb.set_state(conn, "crawler_state", {"seen_urls": [], "seen_hashes": [], "failed": {}})
|
| 298 |
+
agdb.set_config(conn, {"gemini_model": "gemini-3.5-flash", "gemini_vision_model": ""})
|
| 299 |
+
conn.close()
|
| 300 |
+
FAIL_KEYS.add("base")
|
| 301 |
+
ORDER[:] = ["base", "other"]
|
| 302 |
+
orchestrator.run_cycle(trigger="test")
|
| 303 |
+
fps = q("SELECT count(*) FROM agent_text_fingerprints")[0][0]
|
| 304 |
+
check("failed PDF's fingerprint released (only 'other' is claimed)", fps == 1)
|
| 305 |
+
row = q("SELECT pdfs_ingested, tokens_in, text_model, cost_usd FROM agent_runs WHERE kind = 'cycle' "
|
| 306 |
+
"ORDER BY id DESC LIMIT 1")[0]
|
| 307 |
+
check("billed-but-empty extraction still counted in cost, at 3.5 flash",
|
| 308 |
+
row[1] == 30500 and row[2] == "gemini-3.5-flash"
|
| 309 |
+
and close(row[3], C.cost_usd(30500, 8100, "gemini-3.5-flash")))
|
| 310 |
+
FAIL_KEYS.clear()
|
| 311 |
+
|
| 312 |
+
# Dedupe off -> nothing claimed, nothing skipped
|
| 313 |
+
EXTRACTED.clear()
|
| 314 |
+
q("DELETE FROM agent_doi_seen"); q('DELETE FROM "Composites_materials"'); q("DELETE FROM sources")
|
| 315 |
+
q("DELETE FROM agent_text_fingerprints")
|
| 316 |
+
conn = agdb.connect()
|
| 317 |
+
agdb.set_state(conn, "crawler_state", {"seen_urls": [], "seen_hashes": [], "failed": {}})
|
| 318 |
+
agdb.set_config(conn, {"dedupe_text_similarity": 0.0})
|
| 319 |
+
conn.close()
|
| 320 |
+
ORDER[:] = ["base", "variant"]
|
| 321 |
+
orchestrator.run_cycle(trigger="test")
|
| 322 |
+
check("dedupe off: both variants extracted", sorted(EXTRACTED) == ["base", "variant"])
|
| 323 |
+
check("dedupe off: no fingerprints written", q("SELECT count(*) FROM agent_text_fingerprints")[0][0] == 0)
|
| 324 |
+
|
| 325 |
+
print("\n%d checks failed" % len(fails))
|
| 326 |
+
sys.exit(1 if fails else 0)
|
|
@@ -113,7 +113,7 @@ r = batch_ingest.PdfResult(pdf="x", elapsed_s=0, materials=0, extracted=0, inser
|
|
| 113 |
flagged=0, duplicates=0, material_classes=[])
|
| 114 |
check("PdfResult has tokens_in/tokens_out defaulting to 0", (r.tokens_in, r.tokens_out) == (0, 0))
|
| 115 |
check("summarize reports token totals",
|
| 116 |
-
|
| 117 |
|
| 118 |
print("\n%d checks failed" % len(fails))
|
| 119 |
sys.exit(1 if fails else 0)
|
|
|
|
| 113 |
flagged=0, duplicates=0, material_classes=[])
|
| 114 |
check("PdfResult has tokens_in/tokens_out defaulting to 0", (r.tokens_in, r.tokens_out) == (0, 0))
|
| 115 |
check("summarize reports token totals",
|
| 116 |
+
{"tokens_in": 0, "tokens_out": 0}.items() <= batch_ingest.summarize([r])["tokens"].items())
|
| 117 |
|
| 118 |
print("\n%d checks failed" % len(fails))
|
| 119 |
sys.exit(1 if fails else 0)
|