diff --git a/docs/API.md b/docs/API.md index b2f0c5b..69a963d 100644 --- a/docs/API.md +++ b/docs/API.md @@ -851,7 +851,7 @@ def load_equation_spec(path: str | Path) -> EquationSpec: | Implementación | Estado | Notas | |----------------|--------|-------| -| `OpenAlexSource` | **v1** | **Referencia/backbone**, sobre `httpx`. Entrega mínimo + enriquecimiento (refs inline + afiliaciones per-autor + instituciones; `cited_by_id` lo puebla el chaining/Enricher, no el seed). Traducción **passthrough**: envuelve la ecuación en `title_and_abstract.search:(...)` y **reporta** los límites WoS (NEAR/comodín/tags) sin traducirlos. Flag `native=True` (query cruda). **Negaciones (`exclude`):** cada `AND NOT ""` se inyecta **dentro** de la única expresión `search:((query) AND NOT "")` (campo no repetido; el filtro de año queda como predicado separado por coma, fuera del `search`) y se reporta en el `translation_report`; ignorado con `native`. Credenciales inyectadas (arg → `OPENALEX_API_KEY` → `~/.openalex/credentials` → polite pool; ADR 0012). Cursor paging con tope `max_results` (default 200). Puebla `Manifest.openalex_version` (ADR 0017). `transport` inyectable (tests sin red). | +| `OpenAlexSource` | **v1** | **Referencia/backbone**, sobre `httpx`. Entrega mínimo + enriquecimiento (refs inline + afiliaciones per-autor + instituciones; `cited_by_id` lo puebla el chaining/Enricher, no el seed). Traducción **passthrough**: envuelve la ecuación en `title_and_abstract.search:(...)` y **reporta** los límites WoS (NEAR/comodín/tags) sin traducirlos. Flag `native=True` (query cruda). **Negaciones (`exclude`):** cada `AND NOT ""` se inyecta **dentro** de la única expresión `search:((query) AND NOT "")` (campo no repetido; el filtro de año queda como predicado separado por coma, fuera del `search`) y se reporta en el `translation_report`; ignorado con `native`. Credenciales inyectadas (arg → `OPENALEX_API_KEY` → `~/.openalex/credentials` → polite pool; ADR 0012). Cursor paging con tope `max_results` (default 200). Puebla `Manifest.openalex_version` (ADR 0017). `transport` inyectable (tests sin red). Un **429** (rate limit del pool anónimo) en `seed()` → `NetworkError` (exit 4) con mensaje **accionable**: declarar `--email` mueve la petición al polite pool (remedio primario); api_key opcional (ADR 0012, #210). | | `BibtexSource` | **v1, secundaria** | Sembrar desde *pearls* vía `load()`. Extra **`[bibtex]`** (import perezoso de `bibtexparser`); acceso defensivo (campos faltantes sin `KeyError`). Mínimo universal. `seed()` lanza `NotImplementedError`. `.bib` con error grave → `ValueError`; sin entradas válidas → `UserWarning` (no no-op silencioso). Carga bulk con `from_arrow`. | | `ScieloSource` / `RedalycSource` / `LaReferenciaSource` | futuro | Fuentes regionales, mínimo universal. Declaradas, no implementadas (ADR 0018). | | `RisSource` / `CsvSource` | futuro | No implementados. | @@ -860,7 +860,9 @@ def load_equation_spec(path: str | Path) -> EquationSpec: el `Forager` y el `Enricher`): - **`fetch_citing(openalex_id) -> list[dict]`** (singular, forward chaining): `GET works?filter=cites:`, - con retry/backoff ante 429/5xx. + con retry/backoff ante 429/5xx. Al **agotar** los reintentos con **429** → `NetworkError` (exit 4) + **accionable** (polite pool/`--email`; ADR 0012, #210); los 5xx agotados conservan + `httpx.HTTPStatusError`. Asimetría deliberada: solo el 429 tiene remedio del lado del usuario. - **`fetch_citing_batch(ids, *, max_per_paper, since=None) -> dict[seed_id, list[citer_id]]`**: trae los citantes **batcheando por OR** (`cites:W1|W2|...`, lotes ≤50), pagina por cursor y **atribuye página a página** con **presupuesto por semilla** (corta cuando todas alcanzan `max_per_paper`; sin starvation). diff --git a/docs/decisiones/0012-openalex-credenciales.md b/docs/decisiones/0012-openalex-credenciales.md index 7451a8a..bde4b7a 100644 --- a/docs/decisiones/0012-openalex-credenciales.md +++ b/docs/decisiones/0012-openalex-credenciales.md @@ -1,6 +1,7 @@ # 0012 — Credenciales de OpenAlex: email del pool cortés + API key opcional, inyectados -- **Estado:** Aceptada +- **Estado:** Aceptada · **enmendada 2026-06-29** (#210: un 429 se traduce a `NetworkError` + accionable que apunta al polite pool, en `seed` y en el chaining — ver "Seguimiento" al final) - **Fecha:** 2026-06-15 - **Decidido por:** IA (Claude Opus 4.8), validado por el Product Owner humano (ver [`registro-ia.md`](registro-ia.md)) @@ -44,3 +45,24 @@ secreto en código**, ningún `os.environ.get("...", "literal")` para la key. en el README cuando llegue el Hito 4, y testear ambos caminos (con y sin key) contra API mockeada — sin red en CI. - No cambia el núcleo puro: las credenciales viven en la costura `Source`, inyectadas. + +## Seguimiento — 2026-06-29 (#210: el 429 aflora como error accionable hacia el polite pool) + +> Realización de la consecuencia "se avisa" (lección 7): un **429 (Too Many Requests)** del pool +> anónimo ya **no** aflora como `httpx.HTTPStatusError` pelado, sino como **`NetworkError`** (exit 4, +> `service/errors.py`) con un mensaje que nombra el **remedio primario — declarar el email mueve la +> petición al polite pool** (límite más generoso), la **api_key como opcional** y referencia este ADR +> (`_MSG_RATE_LIMIT_429`). Aplica en **ambos caminos**: `seed()` (al recibir 429) y +> `fetch_citing`/chaining (`_fetch_all_with_retry`, al **agotar** los reintentos con 429). +> +> No cambia la **política** de credenciales (la decisión de este ADR): la implementa. Statuses no-429 +> y los retryables 5xx (500/502/503/504) conservan su conducta anterior — asimetría deliberada: solo +> el 429 tiene remedio del lado del usuario (el polite pool). +> +> Extender el fix más allá del `seed` del DoD original —que el chaining levante `NetworkError` en vez +> de `HTTPStatusError` al agotar reintentos— fue **decisión consciente del PO**: dejar el footgun del +> 429 cubierto en todos los caminos de cara a la promoción de la librería. +> +> **Decidido por:** Product Owner humano (2026-06-29, #210); implementado y verificado (gate verde). +> Ver `src/bib2graph/sources/openalex.py` (`_MSG_RATE_LIMIT_429`, `seed`, `_fetch_all_with_retry`) y +> `docs/API.md` §2. diff --git a/pyproject.toml b/pyproject.toml index dc1043b..f40a6fc 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -204,12 +204,13 @@ addopts = [ "--cov=bib2graph", "--cov-report=term-missing", "-m", - "not network", + "not network and not slow", ] markers = [ "unit: tests puros, sin red ni I/O (default)", "integration: tests con I/O local (duckdb/store); corren en el gate", "network: tests que requieren red real (p.ej. OpenAlex); excluidos del gate por defecto (correr con -m network)", + "slow: benchmarks de escala (50 K+ filas); excluidos del gate por defecto (correr con -m slow)", ] [tool.coverage.run] diff --git a/src/bib2graph/backends/duckdb.py b/src/bib2graph/backends/duckdb.py index 4fe0583..461dea9 100644 --- a/src/bib2graph/backends/duckdb.py +++ b/src/bib2graph/backends/duckdb.py @@ -40,6 +40,7 @@ _apply_curation_to_rows, _merge_curation_status, _merge_provenance, + _merge_rows, compute_corpus_hash, ) from bib2graph.constants import LIST_COLUMNS, CurationStatus @@ -254,6 +255,118 @@ def _build_upsert_sql() -> str: _UPSERT_SQL = _build_upsert_sql() + +def _build_bulk_update_existing_sql() -> str: + """SQL para actualizar filas que ya existen en ``corpus`` (conflictos reales). + + Lee el lote desde la vista registrada ``_incoming_upsert`` y actualiza solo + las filas cuyo ``id`` ya existe en ``corpus``. Las UDFs Python + ``_merge_curation_status_udf`` y ``_merge_provenance_udf`` se invocan + únicamente para las filas que efectivamente están en conflicto (no para + inserciones nuevas), eliminando el overhead de UDF para el caso de primer + persist sin conflictos. + + La columna ``_seq`` NO se actualiza: las filas existentes conservan su + orden de primera aparición (ADR 0024, D3). + """ + scalar_cols = [ + "source_id", + "doi", + "title", + "year", + "abstract", + "source", + "language", + "publisher", + "is_seed", + ] + list_cols = list(LIST_COLUMNS) + + scalar_sets = [f" {c} = COALESCE(src.{c}, corpus.{c})" for c in scalar_cols] + + # D3 para listas: unión ordenada deduplicada; NULL si ambos son NULL + list_sets = [ + f" {c} = CASE\n" + f" WHEN corpus.{c} IS NULL AND src.{c} IS NULL THEN NULL\n" + f" ELSE list_sort(list_distinct(list_concat(\n" + f" COALESCE(corpus.{c}, []),\n" + f" COALESCE(src.{c}, [])\n" + f" )))\n" + f" END" + for c in list_cols + ] + + # UDFs para los campos especiales de merge + special_sets = [ + " curation_status = _merge_curation_status_udf(\n" + " corpus.curation_status, corpus.provenance,\n" + " src.curation_status, src.provenance\n" + " )", + " provenance = _merge_provenance_udf(corpus.provenance, src.provenance)", + ] + + all_sets = scalar_sets + list_sets + special_sets + sets_sql = ",\n".join(all_sets) + + return ( + f"UPDATE corpus SET\n" + f"{sets_sql}\n" + f"FROM _incoming_upsert AS src\n" + f"WHERE corpus.id = src.id" + ) + + +def _build_bulk_insert_new_sql() -> str: + """SQL para insertar filas nuevas (sin conflicto con ``corpus``). + + Lee el lote desde la vista registrada ``_incoming_upsert`` y hace INSERT + solo de las filas cuyo ``id`` NO existe todavía en ``corpus``. No invoca + ninguna UDF Python (las filas son nuevas, no hay merge que hacer). + + ADR 0024: ``_seq = CAST($1 AS BIGINT) + ROW_NUMBER() OVER (ORDER BY _row_idx)`` + asigna valores únicos y crecientes en el orden de aparición del lote Arrow. + ``_row_idx`` es una columna auxiliar (índice 0-based) agregada al registrar + la vista, que garantiza que el ROW_NUMBER() respete el orden de la tabla + Arrow entrante sin depender del scan de DuckDB (no garantizado por SQL estándar + sin ORDER BY). ``$1`` es el ``MAX(_seq)`` calculado antes de la llamada. + """ + schema_cols = ", ".join(f.name for f in CORPUS_SCHEMA) + insert_cols = f"{schema_cols}, _seq" + return ( + f"INSERT INTO corpus ({insert_cols})\n" + f"SELECT {schema_cols},\n" + f" CAST($1 AS BIGINT) + ROW_NUMBER() OVER (ORDER BY _row_idx) AS _seq\n" + f"FROM _incoming_upsert\n" + f"WHERE id NOT IN (SELECT id FROM corpus)" + ) + + +def _build_simple_insert_sql() -> str: + """SQL de INSERT directo sin filtro de conflicto (usado tras DELETE en overwrite). + + Después de un ``DELETE FROM corpus`` la tabla está vacía y cualquier + fila de ``_incoming_upsert`` es nueva: el filtro ``NOT IN`` sería un + no-op costoso. Este SQL hace el INSERT sin ninguna condición. + + ADR 0024: ``_seq`` empieza en 1. ``ROW_NUMBER() OVER (ORDER BY _row_idx)`` + garantiza que el orden de filas en la tabla Arrow se preserve fielmente en + ``_seq`` (la columna ``_row_idx`` es el índice 0-based agregado al registrar + la vista; sin ORDER BY el scan de DuckDB no tiene orden garantizado por SQL). + """ + schema_cols = ", ".join(f.name for f in CORPUS_SCHEMA) + insert_cols = f"{schema_cols}, _seq" + return ( + f"INSERT INTO corpus ({insert_cols})\n" + f"SELECT {schema_cols},\n" + f" ROW_NUMBER() OVER (ORDER BY _row_idx) AS _seq\n" + f"FROM _incoming_upsert" + ) + + +_BULK_UPDATE_EXISTING_SQL = _build_bulk_update_existing_sql() +_BULK_INSERT_NEW_SQL = _build_bulk_insert_new_sql() +_SIMPLE_INSERT_SQL = _build_simple_insert_sql() + # --------------------------------------------------------------------------- # Helpers de conversión Python ↔ DuckDB # --------------------------------------------------------------------------- @@ -287,6 +400,45 @@ def _arrow_table_from_con(con: duckdb.DuckDBPyConnection) -> pa.Table: return result.cast(CORPUS_SCHEMA) +def _dedup_merge_table(table: pa.Table) -> pa.Table: + """Deduplica las filas de ``table`` por ``id``, mergeando duplicados con semántica D3. + + Si el lote entrante contiene múltiples filas con el mismo ``id``, las fusiona + progresivamente usando la misma lógica que ``_merge_rows`` de ``backends.memory`` + (la misma semántica que las UDFs ``_merge_curation_status_udf`` y + ``_merge_provenance_udf``). El orden de primera aparición se preserva. + + Garantiza que el lote no viole el PRIMARY KEY de ``corpus`` (``id UNIQUE``), + que antes causaba ``ConstraintException`` en el upsert bulk cuando el batch + entrante ya contenía ids repetidos. + + Args: + table: Tabla Arrow de entrada (puede tener ``id`` duplicados). + + Returns: + Tabla Arrow sin ids duplicados, con el mismo schema canónico. + Si la tabla no tiene duplicados, la devuelve intacta. + """ + if len(table) == 0: + return table + rows = table.to_pylist() + seen: dict[str, dict[str, object]] = {} + order: list[str] = [] + for row in rows: + id_ = str(row["id"]) + if id_ not in seen: + seen[id_] = row + order.append(id_) + else: + # Merge progresivo: preserva provenance acumulada y curación más reciente + seen[id_] = _merge_rows(seen[id_], row) + if len(seen) == len(rows): + # Sin duplicados: devolver la tabla original sin coste de reconstrucción + return table + deduped = [seen[id_] for id_ in order] + return pa.Table.from_pylist(deduped, schema=CORPUS_SCHEMA) + + # --------------------------------------------------------------------------- # DuckDBBackend # --------------------------------------------------------------------------- @@ -421,24 +573,53 @@ def merge_provenance_udf( def _upsert_table(self, table: pa.Table) -> None: """Hace upsert de todas las filas de ``table`` en la base de datos. - ADR 0024: calcula ``start = COALESCE(MAX(_seq), 0)`` una sola vez y - asigna ``_seq = start + 1 + i`` (0-based) a cada fila. Filas ya - existentes (ON CONFLICT) ignoran el ``_seq`` provisto porque el - DO UPDATE no lo actualiza, preservando así el orden de primera - aparición (D3). + Ingesta vectorizada de dos pasos — registra la tabla Arrow en DuckDB + (zero-copy) y ejecuta exactamente 2 + 1 sentencias SQL por lote: + + 1. **UPDATE** filas ya existentes (las UDFs Python se invocan solo para + los conflictos reales; para un primer persist sin conflictos = 0 UDFs). + 2. **INSERT** filas nuevas con ``WHERE id NOT IN (SELECT id FROM corpus)`` + (vectorizado, cero UDFs). + + Elimina el round-trip por fila previo (O(n) → O(1) viajes al motor SQL) + y evita además llamar a las UDFs Python para filas sin conflicto, + logrando la mayor ganancia de velocidad en el caso de primer persist + (corpus vacío → todas las filas son nuevas → 0 UDFs). + + Pre-procesamiento del lote: + - **Dedup-MERGE** (``_dedup_merge_table``): si el lote contiene ids + duplicados, los fusiona con la misma semántica D3 que las UDFs + (provenance append-only, curación más reciente gana) antes de registrar + la vista, evitando violaciones de PRIMARY KEY. + - **Columna ``_row_idx``**: se agrega al registrar la vista para que + ``ROW_NUMBER() OVER (ORDER BY _row_idx)`` en ``_BULK_INSERT_NEW_SQL`` + produzca un ``_seq`` determinista y estable que respeta el orden de + aparición de la tabla Arrow entrante (ADR 0024, D3). + + Las filas existentes conservan su ``_seq`` original (D3). """ - rows = table.to_pylist() - if not rows: + if len(table) == 0: return - # ADR 0024: base para el _seq monótono de esta inserción por lote + # Dedup-MERGE: fusiona ids duplicados en el lote ANTES del upsert SQL + table = _dedup_merge_table(table) + # ADR 0024: base para el _seq monótono de las filas nuevas en este lote result = self._con.execute( "SELECT COALESCE(MAX(_seq), 0) FROM corpus" ).fetchone() start: int = int(result[0]) if result else 0 - for i, row in enumerate(rows): - # _seq = start + 1 + i garantiza valores únicos y crecientes en este lote - params = [*_row_to_params(row), start + 1 + i] - self._con.execute(_UPSERT_SQL, params) + # _row_idx: columna auxiliar de orden para ROW_NUMBER() OVER (ORDER BY _row_idx) + table_with_idx = table.append_column( + "_row_idx", + pa.array(range(len(table)), type=pa.int64()), + ) + self._con.register("_incoming_upsert", table_with_idx) + try: + # Paso 1: actualizar filas existentes (UDFs solo para conflictos reales) + self._con.execute(_BULK_UPDATE_EXISTING_SQL) + # Paso 2: insertar filas nuevas (vectorizado, sin UDF) + self._con.execute(_BULK_INSERT_NEW_SQL, [start]) + finally: + self._con.unregister("_incoming_upsert") def _clone(self) -> DuckDBBackend: """Crea una nueva instancia apuntando al mismo archivo (path compartido). @@ -566,9 +747,12 @@ def merge(self, other_table: pa.Table) -> DuckDBBackend: luego las nuevas filas de ``other_table`` en su orden de aparición) se garantiza por la columna interna ``_seq``: - Filas de ``self`` ya tienen ``_seq`` asignado (se preservan en el - ON CONFLICT DO UPDATE, que no actualiza ``_seq``). - - Filas nuevas de ``other_table`` reciben ``_seq`` mayor (calculado - como ``MAX(_seq) + 1 + i`` en ``_upsert_table``). + UPDATE, que no actualiza ``_seq``). + - Filas nuevas de ``other_table`` reciben ``_seq`` mayor, calculado + con ``ROW_NUMBER() OVER (ORDER BY _row_idx)`` desde ``MAX(_seq)`` + en ``_upsert_table``; la columna auxiliar ``_row_idx`` (índice + 0-based del lote Arrow) garantiza que el orden de primera aparición + se preserve de forma determinista, sin depender del scan de DuckDB. - ``_arrow_table_from_con`` lee con ``ORDER BY _seq``. No se requiere DELETE+reinsert. @@ -1005,9 +1189,10 @@ def close(self) -> None: def overwrite_corpus(self, table: pa.Table) -> None: """Reemplaza TODA la tabla ``corpus`` con el contenido de ``table``. - Hace TRUNCATE + INSERT (no upsert) para que el estado en disco sea - exactamente ``table``, sin residuos de filas previas. Preserva las - tablas hermanas (``loop_state_log``, ``referenced_but_not_fetched``). + Hace TRUNCATE + INSERT masivo (no upsert fila-a-fila) para que el estado + en disco sea exactamente ``table``, sin residuos de filas previas. + Preserva las tablas hermanas (``loop_state_log``, + ``referenced_but_not_fetched``). Úsalo solo en la ruta de ingesta (``seed``, ``restore``, ``chain``, ``thesaurus``) donde ya tenés el corpus completo y correcto en @@ -1024,10 +1209,22 @@ def overwrite_corpus(self, table: pa.Table) -> None: Debe cumplir ``CORPUS_SCHEMA``. """ self._con.execute("DELETE FROM corpus") - rows = table.to_pylist() - for i, row in enumerate(rows): - params = [*_row_to_params(row), i + 1] - self._con.execute(_UPSERT_SQL, params) + if len(table) == 0: + return + # Dedup-MERGE: fusiona ids duplicados en el lote ANTES del INSERT masivo + table = _dedup_merge_table(table) + # Tras DELETE la tabla está vacía: INSERT directo sin filtro NOT IN + # (más eficiente que _BULK_INSERT_NEW_SQL para el caso de tabla limpia) + # _row_idx: columna auxiliar de orden para ROW_NUMBER() OVER (ORDER BY _row_idx) + table_with_idx = table.append_column( + "_row_idx", + pa.array(range(len(table)), type=pa.int64()), + ) + self._con.register("_incoming_upsert", table_with_idx) + try: + self._con.execute(_SIMPLE_INSERT_SQL) + finally: + self._con.unregister("_incoming_upsert") def query(self, sql: str) -> pa.Table: """Ejecuta una consulta SQL sobre el backend y devuelve tabla Arrow. diff --git a/src/bib2graph/sources/openalex.py b/src/bib2graph/sources/openalex.py index f15ee56..29bc311 100644 --- a/src/bib2graph/sources/openalex.py +++ b/src/bib2graph/sources/openalex.py @@ -35,6 +35,7 @@ from bib2graph.constants import Col, CurationStatus from bib2graph.corpus import Corpus, EquationRef, _rows_with_ids from bib2graph.schemas import CORPUS_SCHEMA, ProvenanceEvent +from bib2graph.service.errors import NetworkError from .base import SeedResult @@ -72,6 +73,18 @@ _RETRY_MAX_ATTEMPTS: int = 3 _RETRY_BACKOFF_BASE: float = 1.0 # segundos; duplica por cada intento +# Mensaje accionable para 429 (pool anónimo → polite pool). +# ADR 0012: el email mueve al polite pool (límite más generoso); +# la api_key mejora aún más (opcional). +_MSG_RATE_LIMIT_429 = ( + "OpenAlex respondió 429 (Too Many Requests): límite de tasa del pool anónimo" + " alcanzado. Remedio primario — declarar tu email mueve la petición al polite" + " pool, que tiene un límite más generoso. En el CLI: b2g seed --email" + " tu@email.com … En código: OpenAlexSource(email='tu@email.com')." + " Opcional: una api_key mejora el límite aún más" + " (variable de entorno OPENALEX_API_KEY). Ver ADR 0012." +) + # --------------------------------------------------------------------------- # Traducción de ecuación (función pura, sin I/O) @@ -552,7 +565,12 @@ def seed( equation_id = f"eq-{datetime.now(UTC).strftime('%Y%m%dT%H%M%S')}" fetched_at = datetime.now(UTC).isoformat() - works, openalex_version = self._fetch_all(executed_query) + try: + works, openalex_version = self._fetch_all(executed_query) + except httpx.HTTPStatusError as exc: + if exc.response.status_code == 429: + raise NetworkError(_MSG_RATE_LIMIT_429) from exc + raise # R5: bulk-load — construir tabla Arrow de una vez en vez de N add_paper/clone. rows = [ @@ -612,9 +630,10 @@ def _fetch_all_with_retry( Tupla ``(works, openalex_version)`` con retry/backoff. Raises: - httpx.HTTPStatusError: Si se agotan los reintentos. + NetworkError: Si se agotan los reintentos con 429 (con mensaje accionable). + httpx.HTTPStatusError: Si se agotan los reintentos con otro código retryable. """ - last_exc: Exception | None = None + last_exc: httpx.HTTPStatusError | None = None for attempt in range(_RETRY_MAX_ATTEMPTS): try: return self._fetch_all(filter_str) @@ -627,6 +646,8 @@ def _fetch_all_with_retry( raise # error no retryable: propagar inmediatamente # Se agotaron los reintentos assert last_exc is not None + if last_exc.response.status_code == 429: + raise NetworkError(_MSG_RATE_LIMIT_429) from last_exc raise last_exc def fetch_citing(self, openalex_id: str) -> list[dict[str, Any]]: diff --git a/tests/unit/test_idempotencia_pipeline_bitabit.py b/tests/unit/test_idempotencia_pipeline_bitabit.py index 08113b3..a067c38 100644 --- a/tests/unit/test_idempotencia_pipeline_bitabit.py +++ b/tests/unit/test_idempotencia_pipeline_bitabit.py @@ -33,6 +33,7 @@ from pathlib import Path from typing import Any +import httpx import pyarrow as pa import pyarrow.parquet as pq import pytest @@ -85,6 +86,26 @@ def _communities_of_run(corpus: Any) -> dict[str, dict[str, int]]: } +def _make_offline_transport() -> httpx.MockTransport: + """MockTransport vacío para mantener ``run_build`` SIN red. + + Cuando el corpus tiene seeds aceptadas, ``run_build`` corre una pasada + ``cited_by`` contra OpenAlex (ADR 0038 §enrich) ANTES de proyectar las + redes. El parquet del ejemplo tiene seeds aceptadas, así que sin inyectar + transport esa pasada pegaría a la red real → flake (httpx.ReadTimeout en CI). + Este transport devuelve 0 citantes: la pasada es no-op determinista y el + test queda hermético, como declara su módulo ("Sin red").""" + + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response( + 200, + json={"results": [], "meta": {"count": 0, "next_cursor": None}}, + headers={"x-openalex-api-version": "2026-05-01"}, + ) + + return httpx.MockTransport(handler) + + # --------------------------------------------------------------------------- # 1. corpus_hash idéntico entre dos cargas puras del mismo parquet (sin DuckDB) # --------------------------------------------------------------------------- @@ -203,12 +224,16 @@ def test_comunidades_identicas_entre_dos_runs_pipeline(tmp_path: Path) -> None: from bib2graph.cli.commands.restore import run_restore from bib2graph.stores.duckdb import DuckDBStore - # Run A: pipeline completo (incluye run_build para ejercer el camino de escritura) + # Run A: pipeline completo (incluye run_build para ejercer el camino de escritura). + # transport inyectado → la pasada cited_by (ADR 0038 §enrich) NO toca la red: + # devuelve 0 citantes, manteniendo el test hermético y determinista. store_a = tmp_path / "runA.duckdb" run_restore(store_a, slice_parquet) corpus_a = DuckDBStore(store_a).load() communities_a = _communities_of_run(corpus_a) - run_build(store_a, out_dir=tmp_path / "networks_a") + run_build( + store_a, out_dir=tmp_path / "networks_a", transport=_make_offline_transport() + ) # Run B (store distinto, mismo parquet): solo restore + comunidades. # No se reconstruye run_build: su determinismo de output no se asevera acá diff --git a/tests/unit/test_r5_robustness.py b/tests/unit/test_r5_robustness.py index 2557f2f..5a28ce2 100644 --- a/tests/unit/test_r5_robustness.py +++ b/tests/unit/test_r5_robustness.py @@ -375,10 +375,13 @@ def handler_429_luego_200(request: httpx.Request) -> httpx.Response: @pytest.mark.unit def test_fetch_citing_retry_agotado_lanza_error() -> None: - """fetch_citing lanza HTTPStatusError si se agotan los reintentos. + """fetch_citing lanza NetworkError si se agotan los reintentos con 429. - R5: con 3 intentos de 429 consecutivos, la excepción se propaga. + R5 / #210: con 3 intentos de 429 consecutivos, la excepción se convierte + a NetworkError accionable (no HTTPStatusError pelado). """ + from bib2graph.service.errors import NetworkError + calls: list[int] = [0] def handler_siempre_429(request: httpx.Request) -> httpx.Response: @@ -390,11 +393,13 @@ def handler_siempre_429(request: httpx.Request) -> httpx.Response: with ( patch("bib2graph.sources.openalex.time.sleep"), - pytest.raises(httpx.HTTPStatusError) as exc_info, + pytest.raises(NetworkError) as exc_info, ): source.fetch_citing("W12345") - assert exc_info.value.response.status_code == 429 + msg = str(exc_info.value) + assert "429" in msg + assert "polite" in msg.lower() # Se intentó _RETRY_MAX_ATTEMPTS veces (cada intento = al menos 1 HTTP) assert calls[0] >= 1 @@ -425,6 +430,35 @@ def handler_404(request: httpx.Request) -> httpx.Response: assert not mock_sleep.called +@pytest.mark.unit +def test_fetch_citing_429_levanta_network_error_accionable() -> None: + """fetch_citing con 429 agotado lanza NetworkError con mensaje completo. + + #210: el path de chaining (_fetch_all_with_retry) convierte 429 + agotado en NetworkError con mensaje que menciona email, polite-pool, + ADR 0012 y api_key como opcional. Sin red real (MockTransport). + """ + from bib2graph.service.errors import NetworkError + + def handler_siempre_429(request: httpx.Request) -> httpx.Response: + return httpx.Response(429, text="Rate limit exceeded") + + transport = httpx.MockTransport(handler_siempre_429) + source = OpenAlexSource(transport=transport) + + with ( + patch("bib2graph.sources.openalex.time.sleep"), + pytest.raises(NetworkError) as exc_info, + ): + source.fetch_citing("W99999") + + msg = str(exc_info.value) + assert "email" in msg.lower() + assert "polite" in msg.lower() + assert "ADR 0012" in msg + assert "api_key" in msg # la api_key se nombra como opcional en el mensaje + + # =========================================================================== # 6. Bulk-load: test de no-regresión # =========================================================================== diff --git a/tests/unit/test_sources.py b/tests/unit/test_sources.py index f741665..8f2a546 100644 --- a/tests/unit/test_sources.py +++ b/tests/unit/test_sources.py @@ -782,6 +782,39 @@ def test_translate_exclude_strip_comillas_internas() -> None: assert after_not.count('"') == 1 # solo la comilla de cierre +# --------------------------------------------------------------------------- +# #210 — 429 → NetworkError accionable (polite pool / email) +# --------------------------------------------------------------------------- + + +@pytest.mark.unit +def test_seed_429_levanta_network_error_accionable() -> None: + """Un 429 de OpenAlex en seed() lanza NetworkError con mensaje accionable. + + Sin email declarado → pool anónimo → 429. El error debe mencionar email + y polite pool como remedio primario (no un HTTPStatusError pelado). + Referencia: ADR 0012. + """ + from bib2graph.service.errors import NetworkError + + def handler_429(request: httpx.Request) -> httpx.Response: + return httpx.Response(429, text="Too Many Requests") + + transport = httpx.MockTransport(handler_429) + source = OpenAlexSource(transport=transport) + + with pytest.raises(NetworkError) as exc_info: + source.seed("ecological exchange") + + msg = str(exc_info.value) + assert "429" in msg + assert "email" in msg.lower() + assert "polite" in msg.lower() + assert "OPENALEX_API_KEY" in msg + assert "ADR 0012" in msg + assert "api_key" in msg # la api_key se nombra como opcional en el mensaje + + @pytest.mark.unit def test_translate_exclude_con_anio_forma_completa() -> None: """Con exclude + year, los NOT van dentro de search y el año va fuera con coma. diff --git a/tests/unit/test_stores.py b/tests/unit/test_stores.py index d75a40d..b897031 100644 --- a/tests/unit/test_stores.py +++ b/tests/unit/test_stores.py @@ -290,3 +290,204 @@ def test_duckdb_backend_close_es_idempotente(tmp_path: Path) -> None: backend2 = DuckDBBackend(path=db_path) assert len(backend2) == 0 backend2.close() + + +# --------------------------------------------------------------------------- +# Correctitud: lote con ids duplicados → dedup-MERGE (regresión #211) +# --------------------------------------------------------------------------- + + +@pytest.mark.integration +def test_persist_lote_con_duplicados_dedup_merge(tmp_path: Path) -> None: + """persist de lote con id duplicado: una fila final, provenance acumulada. + + Regresión #211: el upsert bulk anterior lanzaba ConstraintException + (PRIMARY KEY) cuando el lote tenía dos filas con el mismo id. Ahora se + deduplica-MERGE antes del upsert, igual que el loop fila-a-fila previo: + provenance append-only y curación "más reciente gana". + """ + db_path = tmp_path / "lib_dup.duckdb" + paper_id = "oa:test0000aaaa0000" + + prov_seed = json.dumps( + [ + { + "action": "seeded", + "equation_id": "eq1", + "chaining_hop": None, + "source": "openalex", + "fetched_at": None, + "decided_by": None, + "decided_at": None, + } + ] + ) + prov_accept = json.dumps( + [ + { + "action": "accepted", + "equation_id": None, + "chaining_hop": None, + "source": None, + "fetched_at": None, + "decided_by": "revisor", + "decided_at": "2026-01-01T00:00:00+00:00", + } + ] + ) + + # Lote con el mismo id dos veces: candidate→seeded primero, accepted→prov_accept después + row_a = _make_row(id=paper_id, curation_status="candidate", provenance=prov_seed) + row_b = _make_row(id=paper_id, curation_status="accepted", provenance=prov_accept) + corpus = _make_corpus([row_a, row_b]) + + store = DuckDBStore(db_path) + store.persist(corpus) # No debe lanzar ConstraintException + + loaded = store.load() + assert len(loaded) == 1, "Debe haber solo 1 fila tras dedup" + + row = loaded.to_arrow().to_pylist()[0] + assert row["curation_status"] == "accepted", "Curación más reciente gana" + events = json.loads(row["provenance"]) + assert len(events) == 2, "Provenance acumulada de ambas apariciones" + actions = {e["action"] for e in events} + assert "seeded" in actions + assert "accepted" in actions + + +@pytest.mark.integration +def test_persist_overwrite_corpus_lote_con_duplicados(tmp_path: Path) -> None: + """overwrite_corpus con lote duplicado: una fila final, provenance acumulada. + + Complementa el test anterior para la ruta de ``overwrite_corpus`` + (DELETE + INSERT masivo), que también recibe el lote antes del upsert + y debe deduplicar con la misma semántica. + """ + import pyarrow as pa + + from bib2graph.backends.duckdb import DuckDBBackend + from bib2graph.schemas import CORPUS_SCHEMA + + db_path = tmp_path / "lib_ow_dup.duckdb" + paper_id = "oa:test1111bbbb1111" + + prov_a = json.dumps( + [ + { + "action": "seeded", + "equation_id": "eq2", + "chaining_hop": None, + "source": "openalex", + "fetched_at": None, + "decided_by": None, + "decided_at": None, + } + ] + ) + prov_b = json.dumps( + [ + { + "action": "accepted", + "equation_id": None, + "chaining_hop": None, + "source": None, + "fetched_at": None, + "decided_by": "curador", + "decided_at": "2026-02-01T00:00:00+00:00", + } + ] + ) + + row_a = _make_row(id=paper_id, curation_status="candidate", provenance=prov_a) + row_b = _make_row(id=paper_id, curation_status="accepted", provenance=prov_b) + table = pa.Table.from_pylist([row_a, row_b], schema=CORPUS_SCHEMA) + + backend = DuckDBBackend(path=db_path) + backend.overwrite_corpus(table) # No debe lanzar ConstraintException + + assert len(backend) == 1, "Debe haber solo 1 fila tras dedup en overwrite" + + row = backend.to_arrow().to_pylist()[0] + assert row["curation_status"] == "accepted" + events = json.loads(row["provenance"]) + assert len(events) == 2 + actions = {e["action"] for e in events} + assert "seeded" in actions + assert "accepted" in actions + + +# --------------------------------------------------------------------------- +# Correctitud: orden determinista de _seq para múltiples filas nuevas en un lote +# --------------------------------------------------------------------------- + + +@pytest.mark.integration +def test_persist_multifila_nueva_seq_determinista(tmp_path: Path) -> None: + """persist de lote con N filas nuevas → orden D3 estable en to_arrow(). + + Regresión de la observación MEDIO del verifier (#211): ROW_NUMBER() OVER () + sin ORDER BY no garantiza el orden en SQL estándar. Con la columna auxiliar + ``_row_idx`` (ORDER BY _row_idx), el _seq respeta el orden de aparición de + la tabla Arrow entrante. + + El test verifica el ORDER (no solo el corpus_hash, que es order-independent) + con MÚLTIPLES filas nuevas en el lote, pues el bug solo se manifiesta con ≥2 + filas nuevas y dependencia del scan interno de DuckDB. + """ + db_path = tmp_path / "lib_ord.duckdb" + # 6 ids en orden deliberado (mezcla de hex para evitar colusión con ascii sort) + ids = [ + "oa:ffff000000000001", + "oa:aaaa000000000002", + "oa:9999000000000003", + "oa:bbbb000000000004", + "oa:1111000000000005", + "oa:cccc000000000006", + ] + rows = [_make_row(id=id_, title=f"Paper {i}") for i, id_ in enumerate(ids)] + corpus = _make_corpus(rows) + + store = DuckDBStore(db_path) + store.persist(corpus) + + loaded_ids = store.load().to_arrow().column("id").to_pylist() + assert loaded_ids == ids, ( + f"El orden D3 debe coincidir con el de la tabla Arrow entrante.\n" + f"Esperado: {ids}\nObtenido: {loaded_ids}" + ) + + +# --------------------------------------------------------------------------- +# Benchmark de escala — upsert masivo (issue #211) +# Marcado @slow: NO corre en el gate por defecto (-m "not network and not slow"). +# La idempotencia a escala pequeña ya está cubierta por test_persist_idempotente. +# --------------------------------------------------------------------------- + + +@pytest.mark.integration +@pytest.mark.slow +def test_persist_escala_masiva(tmp_path: Path) -> None: + """persist de 50.000 filas es correcto e idempotente (benchmark de escala). + + Consolida los dos tests de 50 K anteriores en uno: verifica que todas las + filas persisten y que una segunda pasada no duplica (idempotencia D3). + + No hay assert de wall-clock (era la fuente de flake en runners CI cargados). + La señal de rendimiento queda implícita en que el test termine en tiempo + razonable (no bloquea el gate por estar marcado @slow). + """ + N = 50_000 + rows = [_make_row(id=f"oa:{i:016x}", title=f"Paper {i}") for i in range(N)] + corpus = _make_corpus(rows) + + db_path = tmp_path / "lib_bench.duckdb" + store = DuckDBStore(db_path) + store.persist(corpus) + + assert len(store.load()) == N, "No se persistieron todas las filas" + + # Segunda pasada: idempotencia (sin duplicados, mismo hash) + store.persist(corpus) + loaded = store.load() + assert len(loaded) == N, "La segunda persist duplicó filas" diff --git a/uv.lock b/uv.lock index 12632f5..5d2f786 100644 --- a/uv.lock +++ b/uv.lock @@ -101,7 +101,7 @@ wheels = [ [[package]] name = "bib2graph" -version = "0.10.0" +version = "0.10.1" source = { editable = "." } dependencies = [ { name = "click" },