diff --git a/app/ingestion/pipeline.py b/app/ingestion/pipeline.py index 9ce433b..e7717af 100644 --- a/app/ingestion/pipeline.py +++ b/app/ingestion/pipeline.py @@ -20,7 +20,6 @@ from app.db.models import ( ArtifactType, Chunk, Document, - DocumentArtifact, DocumentStatus, Page, ProcessingEvent, @@ -35,7 +34,7 @@ from app.ingestion.knowledge_ingest import load_asset_metadata_for_document from app.ingestion.ocr import run_ocr from app.ingestion.table_processor import persist_tables from app.logging_config import get_logger -from app.storage.artifacts import ensure_artifact +from app.storage.artifacts import ensure_artifact, latest_artifact from app.storage.local_paths import ( key_docling_json, key_markdown, @@ -65,12 +64,9 @@ def process_document_id( # noqa: PLR0911, PLR0912, PLR0915 source_path = Path(doc.source_path) sha = doc.sha256 - original_artifact = db.execute( - select(DocumentArtifact).where( - DocumentArtifact.document_id == doc.id, - DocumentArtifact.artifact_type == ArtifactType.ORIGINAL_PDF, - ) - ).scalar_one_or_none() + original_artifact = latest_artifact( + db, document_id=doc.id, artifact_type=ArtifactType.ORIGINAL_PDF + ) work_dir = work_dir_for(document_id) local_pdf = work_dir / f"{sha}.pdf" diff --git a/app/storage/artifacts.py b/app/storage/artifacts.py index 13d11f4..3578705 100644 --- a/app/storage/artifacts.py +++ b/app/storage/artifacts.py @@ -51,3 +51,32 @@ def ensure_artifact( ) db.add(artifact) return artifact + + +def latest_artifact( + db: Session, + *, + document_id: uuid.UUID, + artifact_type: str, +) -> DocumentArtifact | None: + """Newest artifact of the given type, or ``None``. + + Duplicates of one type are legal: repeated knowledge-ingest of identical + content (same sha256, new asset_id → new canonical object key) appends a + second ORIGINAL_PDF row because ensure_artifact identity is + (document_id, storage_key). All such rows reference byte-identical objects + (sha256 verified at ingest), so the newest reference is always a safe, + deterministic pick — never scalar_one_or_none here (G6 MultipleResultsFound). + """ + return ( + db.execute( + select(DocumentArtifact) + .where( + DocumentArtifact.document_id == document_id, + DocumentArtifact.artifact_type == artifact_type, + ) + .order_by(DocumentArtifact.created_at.desc(), DocumentArtifact.id.desc()) + ) + .scalars() + .first() + ) diff --git a/tests/test_artifact_lookup.py b/tests/test_artifact_lookup.py new file mode 100644 index 0000000..aee9e27 --- /dev/null +++ b/tests/test_artifact_lookup.py @@ -0,0 +1,93 @@ +"""G6: ORIGINAL_PDF artifact lookup must tolerate duplicate rows. + +Repeated knowledge-ingest of identical content (same sha256, new asset_id → +new canonical object key) reuses the Document row but appends a second +ORIGINAL_PDF artifact (ensure_artifact identity is document_id+storage_key). +The old ``scalar_one_or_none()`` in ``process_document_id`` then raised +``MultipleResultsFound`` → Celery retry storm. The lookup must instead pick +the newest artifact deterministically — any row is byte-identical (sha256 +verified at ingest), so the latest reference is always safe. +""" + +from __future__ import annotations + +import uuid + +from app.db.models import ArtifactType, DocumentArtifact +from app.storage.artifacts import latest_artifact + + +class _FakeScalars: + def __init__(self, rows): + self._rows = rows + + def first(self): + return self._rows[0] if self._rows else None + + +class _FakeResult: + def __init__(self, rows): + self._rows = rows + + def scalars(self): + return _FakeScalars(self._rows) + + +class _FakeSession: + def __init__(self, rows): + self.rows = rows + self.statements = [] + + def execute(self, stmt): + self.statements.append(stmt) + return _FakeResult(self.rows) + + +def _artifact(key: str) -> DocumentArtifact: + return DocumentArtifact( + id=uuid.uuid4(), + document_id=uuid.uuid4(), + artifact_type=ArtifactType.ORIGINAL_PDF, + storage_bucket="teamhub-mailhub-originals", + storage_key=key, + ) + + +def test_latest_artifact_returns_first_row_without_raising(): + rows = [_artifact("mailhub/2026/07/11/b/original/x.pdf"), + _artifact("mailhub/2026/07/11/a/original/x.pdf")] + db = _FakeSession(rows) + got = latest_artifact(db, document_id=rows[0].document_id, + artifact_type=ArtifactType.ORIGINAL_PDF) + assert got is rows[0] + + +def test_latest_artifact_none_when_absent(): + db = _FakeSession([]) + assert latest_artifact(db, document_id=uuid.uuid4(), + artifact_type=ArtifactType.ORIGINAL_PDF) is None + + +def test_latest_artifact_query_is_deterministic_newest_first(): + """The emitted SQL must filter by document_id+artifact_type and order + newest-first with a stable tiebreaker — that is what makes the pick + deterministic when duplicates exist.""" + db = _FakeSession([]) + latest_artifact(db, document_id=uuid.uuid4(), artifact_type=ArtifactType.ORIGINAL_PDF) + (stmt,) = db.statements + sql = str(stmt.compile(compile_kwargs={"literal_binds": False})) + assert "document_artifacts.document_id" in sql + assert "document_artifacts.artifact_type" in sql + assert "ORDER BY document_artifacts.created_at DESC, document_artifacts.id DESC" in sql + + +def test_pipeline_uses_tolerant_lookup(): + """process_document_id must not use scalar_one_or_none for the ORIGINAL_PDF + lookup (the G6 crash); it goes through latest_artifact instead.""" + import inspect + + from app.ingestion import pipeline + + src = inspect.getsource(pipeline.process_document_id) + assert "latest_artifact" in src + assert "scalar_one_or_none" not in src