diff --git a/extralit-frontend/v1/domain/entities/document/DocumentLayout.ts b/extralit-frontend/v1/domain/entities/document/DocumentLayout.ts index 075c430e3..b10fc86d0 100644 --- a/extralit-frontend/v1/domain/entities/document/DocumentLayout.ts +++ b/extralit-frontend/v1/domain/entities/document/DocumentLayout.ts @@ -44,12 +44,7 @@ export class BoundingBox { /** Fractions of the page, for overlays that position with percentages. */ toRelativeRect(pageWidth: number, pageHeight: number): Rect { - return { - left: this.l / pageWidth, - top: this.t / pageHeight, - width: this.width / pageWidth, - height: this.height / pageHeight, - }; + return this.toRect(pageWidth, pageHeight, 1, 1); } } @@ -62,19 +57,41 @@ export class Provenance { ) {} } +export interface LayoutItemFields { + /** Citation anchor, e.g. `#/texts/12`. Stable for the lifetime of the stored layout. */ + selfRef: string; + label: string; + readingOrder: number; + prov?: Provenance[]; + parentRef?: string | null; + contentLayer?: string | null; + level?: number | null; + text?: string | null; + html?: string | null; +} + export class LayoutItem { - constructor( - /** Citation anchor, e.g. `#/texts/12`. Stable for the lifetime of the stored layout. */ - public readonly selfRef: string, - public readonly label: string, - public readonly readingOrder: number, - public readonly prov: Provenance[] = [], - public readonly parentRef: string | null = null, - public readonly contentLayer: string | null = null, - public readonly level: number | null = null, - public readonly text: string | null = null, - public readonly html: string | null = null - ) {} + readonly selfRef: string; + readonly label: string; + readonly readingOrder: number; + readonly prov: Provenance[]; + readonly parentRef: string | null; + readonly contentLayer: string | null; + readonly level: number | null; + readonly text: string | null; + readonly html: string | null; + + constructor(fields: LayoutItemFields) { + this.selfRef = fields.selfRef; + this.label = fields.label; + this.readingOrder = fields.readingOrder; + this.prov = fields.prov ?? []; + this.parentRef = fields.parentRef ?? null; + this.contentLayer = fields.contentLayer ?? null; + this.level = fields.level ?? null; + this.text = fields.text ?? null; + this.html = fields.html ?? null; + } /** Every page this item touches; more than one when it spans a page break. */ get pageNumbers(): number[] { diff --git a/extralit-frontend/v1/infrastructure/repositories/DocumentRepository.test.ts b/extralit-frontend/v1/infrastructure/repositories/DocumentRepository.test.ts index 6dc0d230d..f20ee8c23 100644 --- a/extralit-frontend/v1/infrastructure/repositories/DocumentRepository.test.ts +++ b/extralit-frontend/v1/infrastructure/repositories/DocumentRepository.test.ts @@ -93,10 +93,20 @@ describe("DocumentRepository", () => { expect(heading.readingOrder).toBe(0); expect(heading.contentLayer).toBe("body"); expect(heading.level).toBe(2); + expect(heading.text).toBe("Methods"); + expect(heading.html).toBeNull(); expect(heading.prov[0].pageNo).toBe(1); expect(heading.prov[0].charspan).toEqual([0, 7]); }); + it("keeps text and html on the field each belongs to", async () => { + const layout = await new DocumentRepository(axiosMock(() => BACKEND_LAYOUT)).getDocumentLayout("d-1"); + const table = layout.itemByRef("#/tables/0"); + + expect(table.text).toBeNull(); + expect(table.html).toBe("
a
"); + }); + it("passes page and label filters as query params", async () => { const axios = axiosMock(() => BACKEND_LAYOUT); const repository = new DocumentRepository(axios); diff --git a/extralit-frontend/v1/infrastructure/repositories/DocumentRepository.ts b/extralit-frontend/v1/infrastructure/repositories/DocumentRepository.ts index cf02ec7b5..f2c2ff4b0 100644 --- a/extralit-frontend/v1/infrastructure/repositories/DocumentRepository.ts +++ b/extralit-frontend/v1/infrastructure/repositories/DocumentRepository.ts @@ -63,17 +63,17 @@ const toProvenance = (prov: BackendProvenance): Provenance => new Provenance(prov.page_no, toBoundingBox(prov.bbox), prov.charspan); const toLayoutItem = (item: BackendLayoutItem): LayoutItem => - new LayoutItem( - item.self_ref, - item.label, - item.reading_order, - (item.prov ?? []).map(toProvenance), - item.parent_ref ?? null, - item.content_layer ?? null, - item.level ?? null, - item.text ?? null, - item.html ?? null - ); + new LayoutItem({ + selfRef: item.self_ref, + label: item.label, + readingOrder: item.reading_order, + prov: (item.prov ?? []).map(toProvenance), + parentRef: item.parent_ref, + contentLayer: item.content_layer, + level: item.level, + text: item.text, + html: item.html, + }); const toLayoutPage = (page: BackendLayoutPage): LayoutPage => new LayoutPage(page.page_no, page.width, page.height); diff --git a/extralit-server/scripts/bench_layout_store.py b/extralit-server/scripts/bench_layout_store.py index 19f3f59d9..32359bc68 100644 --- a/extralit-server/scripts/bench_layout_store.py +++ b/extralit-server/scripts/bench_layout_store.py @@ -125,7 +125,7 @@ def write_lance(): ) timed("query: one document (lance)", lambda: store.load_items(one).num_rows) - print(f"\nfragments {store.fragment_count(ITEMS_DATASET)}") + print(f"\nfragments {len(store.open(ITEMS_DATASET).get_fragments())}") print(f"bytes: parquet {du(parquet_dir):,} over {len(list(parquet_dir.iterdir())):,} objects") if not args.workspace: lance_dir = root / "lance" diff --git a/extralit-server/src/extralit_server/api/schemas/v1/workflows.py b/extralit-server/src/extralit_server/api/schemas/v1/workflows.py index 9ceff77c3..fe675d8d1 100644 --- a/extralit-server/src/extralit_server/api/schemas/v1/workflows.py +++ b/extralit-server/src/extralit_server/api/schemas/v1/workflows.py @@ -4,7 +4,7 @@ from pydantic import BaseModel, Field, field_validator -from extralit_server.contexts.ocr.parsers import list_parsers +from extralit_server.contexts.ocr.parsers.registry import list_parsers class StartWorkflowRequest(BaseModel): diff --git a/extralit-server/src/extralit_server/contexts/ocr/layout_store.py b/extralit-server/src/extralit_server/contexts/ocr/layout_store.py index 5dcfa1dbc..886f5559f 100644 --- a/extralit-server/src/extralit_server/contexts/ocr/layout_store.py +++ b/extralit-server/src/extralit_server/contexts/ocr/layout_store.py @@ -221,19 +221,16 @@ def source(self, name: str) -> Any: # --- maintenance --------------------------------------------------------------------------- - def fragment_count(self, name: str) -> int: - dataset = self.open(name) - return 0 if dataset is None else len(dataset.get_fragments()) - def maybe_compact(self) -> None: """Best effort: a failed compaction must never fail the extraction that triggered it.""" for name in (ITEMS_DATASET, PAGES_DATASET): try: - if self.fragment_count(name) <= COMPACT_FRAGMENT_THRESHOLD: - continue dataset = self.open(name) + if dataset is None or len(dataset.get_fragments()) <= COMPACT_FRAGMENT_THRESHOLD: + continue + # `compact_files` advances the handle in place, so cleanup sees the new version. dataset.optimize.compact_files(target_rows_per_fragment=TARGET_ROWS_PER_FRAGMENT) - self.open(name).cleanup_old_versions(older_than=CLEANUP_OLDER_THAN) + dataset.cleanup_old_versions(older_than=CLEANUP_OLDER_THAN) except Exception as error: _LOGGER.warning(f"Layout compaction of {name} at {self.root_uri} failed: {error}") diff --git a/extralit-server/src/extralit_server/contexts/ocr/parsers/__init__.py b/extralit-server/src/extralit_server/contexts/ocr/parsers/__init__.py index 72fb50181..e69de29bb 100644 --- a/extralit-server/src/extralit_server/contexts/ocr/parsers/__init__.py +++ b/extralit-server/src/extralit_server/contexts/ocr/parsers/__init__.py @@ -1,65 +0,0 @@ -"""Swappable PDF→`DoclingDocument` parsers. - -Each parser normalizes its backend into `LayoutBlock`s and hands them to the shared builder, -so the document that comes out is the same shape regardless of which one ran. -""" - -from __future__ import annotations - -import logging -from collections.abc import Sequence -from typing import Optional, Protocol - -from docling_core.types.doc import DoclingDocument - -_LOGGER = logging.getLogger(__name__) - - -class LayoutParser(Protocol): - """Parse PDF bytes into a `DoclingDocument`.""" - - def __call__( - self, - pdf_bytes: bytes, - *, - name: str, - pages: Optional[Sequence[int]] = None, - filename: Optional[str] = None, - ) -> DoclingDocument: ... - - -_PARSERS: dict[str, LayoutParser] = {} - -from extralit_server.contexts.ocr.parsers.pdf_inspector import parse as _parse_pdf_inspector - -_PARSERS["pdf_inspector"] = _parse_pdf_inspector - -_PYMUPDF_AVAILABLE = False -try: - from extralit_server.contexts.ocr.parsers.pymupdf import parse as _parse_pymupdf - - _PARSERS["pymupdf"] = _parse_pymupdf - _PYMUPDF_AVAILABLE = True -except ImportError as e: # AGPL extra, deliberately optional - _LOGGER.debug(f"pymupdf layout parser unavailable: {e}") - - -def list_parsers() -> list[str]: - """Names of every parser installed in this environment.""" - return sorted(_PARSERS) - - -def get_parser(name: str) -> LayoutParser: - """Look up a parser by name.""" - try: - return _PARSERS[name] - except KeyError: - raise ValueError(f"unknown layout parser {name!r}; available: {list_parsers()}") from None - - -def default_parser_name() -> str: - """Prefer pymupdf's higher-fidelity geometry when the extra is installed.""" - return "pymupdf" if _PYMUPDF_AVAILABLE else "pdf_inspector" - - -__all__ = ["LayoutParser", "default_parser_name", "get_parser", "list_parsers"] diff --git a/extralit-server/src/extralit_server/contexts/ocr/parsers/registry.py b/extralit-server/src/extralit_server/contexts/ocr/parsers/registry.py new file mode 100644 index 000000000..c228161f9 --- /dev/null +++ b/extralit-server/src/extralit_server/contexts/ocr/parsers/registry.py @@ -0,0 +1,60 @@ +"""Swappable PDF→`DoclingDocument` parsers. + +Each parser normalizes its backend into `LayoutBlock`s and hands them to the shared builder, +so the document that comes out is the same shape regardless of which one ran. +""" + +from __future__ import annotations + +import logging +from collections.abc import Sequence +from typing import Optional, Protocol + +from docling_core.types.doc import DoclingDocument + +from extralit_server.contexts.ocr.parsers.pdf_inspector import parse as _parse_pdf_inspector + +try: + from extralit_server.contexts.ocr.parsers.pymupdf import parse as _parse_pymupdf +except ImportError as e: # AGPL extra, deliberately optional + _parse_pymupdf = None + logging.getLogger(__name__).debug(f"pymupdf layout parser unavailable: {e}") + + +class LayoutParser(Protocol): + """Parse PDF bytes into a `DoclingDocument`.""" + + def __call__( + self, + pdf_bytes: bytes, + *, + name: str, + pages: Optional[Sequence[int]] = None, + filename: Optional[str] = None, + ) -> DoclingDocument: ... + + +_PARSERS: dict[str, LayoutParser] = {"pdf_inspector": _parse_pdf_inspector} +if _parse_pymupdf is not None: + _PARSERS["pymupdf"] = _parse_pymupdf + + +def list_parsers() -> list[str]: + """Names of every parser installed in this environment.""" + return sorted(_PARSERS) + + +def get_parser(name: str) -> LayoutParser: + """Look up a parser by name.""" + try: + return _PARSERS[name] + except KeyError: + raise ValueError(f"unknown layout parser {name!r}; available: {list_parsers()}") from None + + +def default_parser_name() -> str: + """Prefer pymupdf's higher-fidelity geometry when the extra is installed.""" + return "pymupdf" if "pymupdf" in _PARSERS else "pdf_inspector" + + +__all__ = ["LayoutParser", "default_parser_name", "get_parser", "list_parsers"] diff --git a/extralit-server/src/extralit_server/contexts/workflows.py b/extralit-server/src/extralit_server/contexts/workflows.py index d57cff7a8..dc2cc94eb 100644 --- a/extralit-server/src/extralit_server/contexts/workflows.py +++ b/extralit-server/src/extralit_server/contexts/workflows.py @@ -8,10 +8,11 @@ from rq.exceptions import NoSuchJobError from rq.group import Group from rq.job import Job, JobStatus +from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from extralit_server.jobs.queues import REDIS_CONNECTION -from extralit_server.models.database import DocumentWorkflow +from extralit_server.models.database import Document, DocumentWorkflow _LOGGER = logging.getLogger(__name__) @@ -511,20 +512,26 @@ def get_failed_jobs_in_group(group_id: str) -> list[dict[str, Any]]: return [] -async def is_current_workflow_run(db: AsyncSession, document_id: UUID, workflow_id: Optional[str]) -> bool: - """Whether this job still belongs to the document's newest workflow run. +async def writer_skip_reason(db: AsyncSession, document_id: UUID, workflow_id: Optional[str]) -> Optional[str]: + """Why this job must not write its artifacts, or None when it may. Checked before every commit. - `send_stop_job_command` only *asks* a worker to stop, so a forced restart can leave the previous - run alive long enough to overwrite the new one's PDF, layout or metadata. Every writer checks - this generation token before it writes; the workflow row is the token. - - A job with no workflow in its meta (direct call, test, ad-hoc enqueue) is always current. + Two ways a job outlives what it was started for. The document can be deleted mid-run, and + writing then resurrects artifacts for a row that no longer exists. Or a forced restart can + supersede it — `send_stop_job_command` only *asks* a worker to stop, so the previous run can + stay alive long enough to overwrite the new one's PDF, layout or metadata. The workflow row is + the generation token; a job with no workflow in its meta (direct call, test, ad-hoc enqueue) + is always current. """ + if await db.scalar(select(Document.id).where(Document.id == document_id)) is None: + return "document deleted" + if not workflow_id: - return True + return None workflow = await DocumentWorkflow.get_by_document_id(db, document_id) - return workflow is None or str(workflow.id) == str(workflow_id) + if workflow is not None and str(workflow.id) != str(workflow_id): + return "workflow superseded" + return None def stop_workflow_jobs(group_id: str) -> list[str]: diff --git a/extralit-server/src/extralit_server/jobs/document_jobs.py b/extralit-server/src/extralit_server/jobs/document_jobs.py index 89295e666..7f92c6add 100644 --- a/extralit-server/src/extralit_server/jobs/document_jobs.py +++ b/extralit-server/src/extralit_server/jobs/document_jobs.py @@ -6,7 +6,6 @@ from rq import Retry, get_current_job from rq.decorators import job -from sqlalchemy import select from extralit_server.api.schemas.v1.document.metadata import DocumentProcessingMetadata from extralit_server.contexts import files @@ -14,10 +13,9 @@ from extralit_server.contexts.document.metadata import update_processing_metadata from extralit_server.contexts.document.preprocessing import PDFPreprocessor from extralit_server.contexts.ocr.triage import triage_pdf -from extralit_server.contexts.workflows import is_current_workflow_run +from extralit_server.contexts.workflows import writer_skip_reason from extralit_server.database import AsyncSessionLocal from extralit_server.jobs.queues import DEFAULT_QUEUE, REDIS_CONNECTION -from extralit_server.models.database import Document _LOGGER = logging.getLogger(__name__) @@ -86,12 +84,10 @@ async def analysis_and_preprocess_job( # A forced restart may already be running; its rotation and metadata must not lose to this # one's, because stopping a started job is only a request. async with AsyncSessionLocal() as db: - if await db.scalar(select(Document.id).where(Document.id == document_id)) is None: - _LOGGER.info(f"Document {document_id} was deleted before its analysis could store") - return {"document_id": str(document_id), "skipped": "document deleted"} - if not await is_current_workflow_run(db, document_id, current_job.meta.get("workflow_id")): - _LOGGER.info(f"Analysis run for document {document_id} was superseded before it could store") - return {"document_id": str(document_id), "skipped": "workflow superseded"} + skip = await writer_skip_reason(db, document_id, current_job.meta.get("workflow_id")) + if skip is not None: + _LOGGER.info(f"Analysis for document {document_id} was not stored: {skip}") + return {"document_id": str(document_id), "skipped": skip} # The PDF rewrite is the last S3 write of this job — dependents key on it. object_path = s3_url.replace(f"/api/v1/file/{workspace_name}/", "") @@ -121,21 +117,22 @@ async def analysis_and_preprocess_job( metadata={"processing_applied": "ocrmypdf_rotation", "original_filename": filename}, ) + preprocessing_result = { + "processing_time": processing_response.metadata.processing_time, + "ocr_applied": False, + "rotation_ran": processing_response.metadata.rotation_ran, + "error": processing_response.metadata.error, + } combined_result = { "document_id": str(document_id), "analysis_result": analysis_result, - "preprocessing_result": { - "processing_time": processing_response.metadata.processing_time, - "ocr_applied": False, - "rotation_ran": processing_response.metadata.rotation_ran, - "error": processing_response.metadata.error, - }, + "preprocessing_result": preprocessing_result, } # The layout job writes the same JSON column concurrently; both go through the row lock. def apply(metadata: DocumentProcessingMetadata) -> None: metadata.update_analysis_results(analysis_result) - metadata.update_preprocessing_results(combined_result["preprocessing_result"]) + metadata.update_preprocessing_results(preprocessing_result) async with AsyncSessionLocal() as db: if await update_processing_metadata(db, document_id, apply) is None: diff --git a/extralit-server/src/extralit_server/jobs/ocr_jobs.py b/extralit-server/src/extralit_server/jobs/ocr_jobs.py index 8671dbb6e..456470163 100644 --- a/extralit-server/src/extralit_server/jobs/ocr_jobs.py +++ b/extralit-server/src/extralit_server/jobs/ocr_jobs.py @@ -7,19 +7,17 @@ from rq import Retry, get_current_job from rq.decorators import job -from sqlalchemy import select from extralit_server.api.schemas.v1.document.metadata import LayoutMetadata from extralit_server.contexts import files from extralit_server.contexts.document.metadata import update_processing_metadata from extralit_server.contexts.ocr import storage from extralit_server.contexts.ocr.layout_store import LayoutStore -from extralit_server.contexts.ocr.parsers import default_parser_name, get_parser from extralit_server.contexts.ocr.parsers.pdf_inspector import classify -from extralit_server.contexts.workflows import is_current_workflow_run +from extralit_server.contexts.ocr.parsers.registry import default_parser_name, get_parser +from extralit_server.contexts.workflows import writer_skip_reason from extralit_server.database import AsyncSessionLocal from extralit_server.jobs.queues import OCR_QUEUE, REDIS_CONNECTION -from extralit_server.models.database import Document _LOGGER = logging.getLogger(__name__) @@ -97,15 +95,10 @@ async def async_document_layout_job( store = LayoutStore.for_workspace(workspace_name) async with store.locked(): async with AsyncSessionLocal() as db: - still_exists = await db.scalar(select(Document.id).where(Document.id == document_id)) - superseded = not await is_current_workflow_run(db, document_id, workflow_id) - if still_exists is None: - _LOGGER.info(f"Document {document_id} was deleted before its layout was stored") - return {"document_id": str(document_id), "parser": parser_name, "skipped": "document deleted"} - if superseded: - # A forced restart already began; its layout must not lose to this one's. - _LOGGER.info(f"Layout run for document {document_id} was superseded before it could store") - return {"document_id": str(document_id), "parser": parser_name, "skipped": "workflow superseded"} + skip = await writer_skip_reason(db, document_id, workflow_id) + if skip is not None: + _LOGGER.info(f"Layout for document {document_id} was not stored: {skip}") + return {"document_id": str(document_id), "parser": parser_name, "skipped": skip} paths = await storage.store_layout(s3_client, workspace_name, document_id, doc, store=store) diff --git a/extralit-server/src/extralit_server/workflows/documents.py b/extralit-server/src/extralit_server/workflows/documents.py index 43bd496c4..77794c6d4 100644 --- a/extralit-server/src/extralit_server/workflows/documents.py +++ b/extralit-server/src/extralit_server/workflows/documents.py @@ -5,7 +5,7 @@ from rq.group import Group from rq.job import Dependency -from extralit_server.contexts.ocr.parsers import default_parser_name +from extralit_server.contexts.ocr.parsers.registry import default_parser_name from extralit_server.database import AsyncSessionLocal from extralit_server.jobs.document_jobs import analysis_and_preprocess_job from extralit_server.jobs.ocr_jobs import async_document_layout_job @@ -62,6 +62,14 @@ async def create_document_workflow( await db.commit() await db.refresh(workflow) + def meta(step: str) -> dict[str, str]: + return { + "document_id": str(document_id), + "reference": reference, + "workflow_step": step, + "workflow_id": str(workflow.id), + } + # Step 4: Prepare jobs using Queue.prepare_data(); the @job decorator kwargs are inert here. analysis_job_data = DEFAULT_QUEUE.prepare_data( analysis_and_preprocess_job, @@ -70,12 +78,7 @@ async def create_document_workflow( job_id=f"analysis_preprocess_{document_id}_{run_suffix}", retry=Retry(max=3, interval=[10, 30, 60]), result_ttl=JOB_RESULT_TTL, - meta={ - "document_id": str(document_id), - "reference": reference, - "workflow_step": "analysis_and_preprocess", - "workflow_id": str(workflow.id), - }, + meta=meta("analysis_and_preprocess"), ) analysis_jobs = group.enqueue_many(queue=DEFAULT_QUEUE, job_datas=[analysis_job_data]) @@ -93,12 +96,7 @@ async def create_document_workflow( depends_on=on_analysis, retry=Retry(max=2, interval=[30, 60]), result_ttl=JOB_RESULT_TTL, - meta={ - "document_id": str(document_id), - "reference": reference, - "workflow_step": "text_extraction", - "workflow_id": str(workflow.id), - }, + meta=meta("text_extraction"), ) group.enqueue_many(queue=OCR_QUEUE, job_datas=[text_extraction_job_data]) @@ -114,12 +112,7 @@ async def create_document_workflow( depends_on=on_analysis, retry=Retry(max=2, interval=[30, 60]), result_ttl=JOB_RESULT_TTL, - meta={ - "document_id": str(document_id), - "reference": reference, - "workflow_step": "document_layout", - "workflow_id": str(workflow.id), - }, + meta=meta("document_layout"), ) group.enqueue_many(queue=OCR_QUEUE, job_datas=[layout_job_data]) diff --git a/extralit-server/tests/unit/api/schemas/v1/test_workflows.py b/extralit-server/tests/unit/api/schemas/v1/test_workflows.py index 7a6bf817f..4d2f898bb 100644 --- a/extralit-server/tests/unit/api/schemas/v1/test_workflows.py +++ b/extralit-server/tests/unit/api/schemas/v1/test_workflows.py @@ -4,7 +4,7 @@ from pydantic import ValidationError from extralit_server.api.schemas.v1.workflows import StartWorkflowRequest -from extralit_server.contexts.ocr.parsers import list_parsers +from extralit_server.contexts.ocr.parsers.registry import list_parsers def _request(**overrides) -> StartWorkflowRequest: diff --git a/extralit-server/tests/unit/contexts/ocr/test_layout_store.py b/extralit-server/tests/unit/contexts/ocr/test_layout_store.py index 8ed72e182..e19d8a5a4 100644 --- a/extralit-server/tests/unit/contexts/ocr/test_layout_store.py +++ b/extralit-server/tests/unit/contexts/ocr/test_layout_store.py @@ -22,6 +22,11 @@ WORKSPACE = "ws-layout" +def fragment_count(store: LayoutStore, name: str) -> int: + dataset = store.open(name) + return 0 if dataset is None else len(dataset.get_fragments()) + + def items(document_id: str, count: int = 3, label: str = "text") -> pa.Table: return pa.Table.from_pylist( [ @@ -242,19 +247,17 @@ def test_compaction_fires_past_the_threshold_and_preserves_rows(self, local_stor for document_id in document_ids: local_store.replace_document(document_id, items(document_id, 2), pages(document_id)) - before = local_store.fragment_count(ITEMS_DATASET) + before = fragment_count(local_store, ITEMS_DATASET) local_store.maybe_compact() assert before > 3 - assert local_store.fragment_count(ITEMS_DATASET) < before + assert fragment_count(local_store, ITEMS_DATASET) < before for document_id in document_ids: assert local_store.load_items(document_id).num_rows == 2 def test_a_broken_compaction_never_reaches_the_caller(self, local_store, monkeypatch): document_id = str(uuid4()) local_store.replace_document(document_id, items(document_id), pages(document_id)) - monkeypatch.setattr( - LayoutStore, "fragment_count", lambda *args, **kwargs: (_ for _ in ()).throw(RuntimeError("boom")) - ) + monkeypatch.setattr(LayoutStore, "open", lambda *args, **kwargs: (_ for _ in ()).throw(RuntimeError("boom"))) local_store.maybe_compact() diff --git a/extralit-server/tests/unit/contexts/ocr/test_pdf_inspector_parser.py b/extralit-server/tests/unit/contexts/ocr/test_pdf_inspector_parser.py index a79092a23..d1919a676 100644 --- a/extralit-server/tests/unit/contexts/ocr/test_pdf_inspector_parser.py +++ b/extralit-server/tests/unit/contexts/ocr/test_pdf_inspector_parser.py @@ -5,13 +5,13 @@ import pytest from docling_core.types.doc import CoordOrigin, DocItemLabel -from extralit_server.contexts.ocr.parsers import get_parser, list_parsers from extralit_server.contexts.ocr.parsers.pdf_inspector import ( classify, page_sizes, parse, role_to_label, ) +from extralit_server.contexts.ocr.parsers.registry import get_parser, list_parsers FIXTURES = Path(__file__).parents[3] / "fixtures" / "pdf" PAGE_WIDTH, PAGE_HEIGHT = 612.0, 792.0 diff --git a/extralit-server/tests/unit/contexts/ocr/test_pymupdf_parser.py b/extralit-server/tests/unit/contexts/ocr/test_pymupdf_parser.py index 60f71599d..ece1d814c 100644 --- a/extralit-server/tests/unit/contexts/ocr/test_pymupdf_parser.py +++ b/extralit-server/tests/unit/contexts/ocr/test_pymupdf_parser.py @@ -7,8 +7,8 @@ pytest.importorskip("pymupdf4llm") -from extralit_server.contexts.ocr.parsers import get_parser, list_parsers from extralit_server.contexts.ocr.parsers.pymupdf import parse +from extralit_server.contexts.ocr.parsers.registry import get_parser, list_parsers FIXTURES = Path(__file__).parents[3] / "fixtures" / "pdf" PAGE_WIDTH, PAGE_HEIGHT = 612.0, 792.0 diff --git a/extralit-server/tests/unit/jobs/test_document_jobs.py b/extralit-server/tests/unit/jobs/test_document_jobs.py index 1ace996eb..5c06c46bd 100644 --- a/extralit-server/tests/unit/jobs/test_document_jobs.py +++ b/extralit-server/tests/unit/jobs/test_document_jobs.py @@ -57,7 +57,7 @@ async def put_object(client, workspace, key, data, **kwargs): patch(f"{MODULE}.PDFPreprocessor", return_value=preprocessor), patch(f"{MODULE}.triage_pdf", return_value=triage()) as triage_pdf, patch(f"{MODULE}.update_processing_metadata", AsyncMock()) as update_metadata, - patch(f"{MODULE}.is_current_workflow_run", AsyncMock(return_value=True)) as is_current, + patch(f"{MODULE}.writer_skip_reason", AsyncMock(return_value=None)) as skip, patch(f"{MODULE}.get_current_job", return_value=current_job), patch(f"{MODULE}.AsyncSessionLocal") as session, ): @@ -76,7 +76,7 @@ async def put_object(client, workspace, key, data, **kwargs): "triage_pdf": triage_pdf, "update_metadata": update_metadata, "written": written, - "is_current": is_current, + "skip": skip, } @@ -170,7 +170,7 @@ async def test_a_storage_failure_surfaces_on_the_job(self, job_context): @pytest.mark.asyncio class TestSupersededRuns: async def test_a_superseded_run_rewrites_nothing(self, job_context): - job_context["is_current"].return_value = False + job_context["skip"].return_value = "workflow superseded" result = await run_job() diff --git a/extralit-server/tests/unit/jobs/test_ocr_jobs.py b/extralit-server/tests/unit/jobs/test_ocr_jobs.py index 0f6d01ca8..f4df568dd 100644 --- a/extralit-server/tests/unit/jobs/test_ocr_jobs.py +++ b/extralit-server/tests/unit/jobs/test_ocr_jobs.py @@ -29,7 +29,6 @@ async def update_metadata(db, document_id, mutate): return None db = MagicMock() - db.scalar = AsyncMock(return_value=uuid4()) with ( patch(f"{MODULE}.files.get_s3_client", AsyncMock(return_value=AsyncMock())), @@ -40,12 +39,12 @@ async def update_metadata(db, document_id, mutate): patch(f"{MODULE}.storage.store_layout", store_layout), patch(f"{MODULE}.update_processing_metadata", update_metadata), patch(f"{MODULE}.get_current_job", return_value=MagicMock(meta={"workflow_id": "wf-1"})), - patch(f"{MODULE}.is_current_workflow_run", AsyncMock(return_value=True)) as is_current, + patch(f"{MODULE}.writer_skip_reason", AsyncMock(return_value=None)) as skip, patch(f"{MODULE}.AsyncSessionLocal") as session, ): session.return_value.__aenter__ = AsyncMock(return_value=db) session.return_value.__aexit__ = AsyncMock(return_value=False) - yield {"calls": calls, "db": db, "store": store, "is_current": is_current} + yield {"calls": calls, "store": store, "skip": skip} async def run_job(document_id=None): @@ -73,7 +72,7 @@ async def test_metadata_is_updated_after_the_lock_is_released(self, job_context) assert ("update_metadata", 0) in job_context["calls"] async def test_a_document_deleted_mid_parse_is_not_resurrected(self, job_context): - job_context["db"].scalar = AsyncMock(return_value=None) + job_context["skip"].return_value = "document deleted" result = await run_job() @@ -89,10 +88,10 @@ async def test_pages_needing_ocr_are_surfaced(self, job_context): class TestSupersededRuns: async def test_a_superseded_run_does_not_write(self, job_context): # Stopping a started job is only a request, so a forced restart can overlap this run. - job_context["is_current"].return_value = False + job_context["skip"].return_value = "workflow superseded" result = await run_job() - assert job_context["is_current"].await_args.args[2] == "wf-1" + assert job_context["skip"].await_args.args[2] == "wf-1" assert result["skipped"] == "workflow superseded" assert job_context["calls"] == [] diff --git a/extralit-server/tests/unit/workflows/test_document_workflow_layout.py b/extralit-server/tests/unit/workflows/test_document_workflow_layout.py index 567fe496b..9417752fa 100644 --- a/extralit-server/tests/unit/workflows/test_document_workflow_layout.py +++ b/extralit-server/tests/unit/workflows/test_document_workflow_layout.py @@ -63,7 +63,7 @@ def prepared_for(calls, step): @pytest.mark.asyncio class TestLayoutJobSequencing: async def test_layout_runs_by_default_with_the_default_parser(self, enqueued): - from extralit_server.contexts.ocr.parsers import default_parser_name + from extralit_server.contexts.ocr.parsers.registry import default_parser_name await run_workflow(layout_parser=None) diff --git a/extralit-server/tests/unit/workflows/test_workflow_generation.py b/extralit-server/tests/unit/workflows/test_workflow_generation.py index bc1978b31..6caf0a418 100644 --- a/extralit-server/tests/unit/workflows/test_workflow_generation.py +++ b/extralit-server/tests/unit/workflows/test_workflow_generation.py @@ -1,35 +1,48 @@ -"""Tests for the workflow-generation guard every artifact writer checks.""" +"""Tests for the guard every artifact writer checks before it commits.""" from unittest.mock import AsyncMock, MagicMock, patch from uuid import uuid4 import pytest -from extralit_server.contexts.workflows import is_current_workflow_run +from extralit_server.contexts.workflows import writer_skip_reason MODULE = "extralit_server.contexts.workflows" pytestmark = pytest.mark.asyncio -class TestIsCurrentWorkflowRun: - async def test_the_newest_run_is_current(self): +def db(document_exists: bool = True) -> AsyncMock: + session = AsyncMock() + session.scalar = AsyncMock(return_value=uuid4() if document_exists else None) + return session + + +class TestWriterSkipReason: + async def test_the_newest_run_may_write(self): workflow = MagicMock(id=uuid4()) with patch(f"{MODULE}.DocumentWorkflow.get_by_document_id", AsyncMock(return_value=workflow)): - assert await is_current_workflow_run(AsyncMock(), uuid4(), str(workflow.id)) is True + assert await writer_skip_reason(db(), uuid4(), str(workflow.id)) is None - async def test_a_superseded_run_is_not_current(self): + async def test_a_superseded_run_is_told_why(self): with patch(f"{MODULE}.DocumentWorkflow.get_by_document_id", AsyncMock(return_value=MagicMock(id=uuid4()))): - assert await is_current_workflow_run(AsyncMock(), uuid4(), str(uuid4())) is False + assert await writer_skip_reason(db(), uuid4(), str(uuid4())) == "workflow superseded" + + async def test_a_deleted_document_is_told_why(self): + with patch(f"{MODULE}.DocumentWorkflow.get_by_document_id", AsyncMock()) as lookup: + assert await writer_skip_reason(db(document_exists=False), uuid4(), str(uuid4())) == "document deleted" + + # Deletion wins outright; a superseded run of a deleted document is still just deleted. + lookup.assert_not_awaited() - async def test_a_job_without_a_workflow_is_always_current(self): + async def test_a_job_without_a_workflow_may_always_write(self): # Direct calls and ad-hoc enqueues carry no workflow; they must not be gated on one. with patch(f"{MODULE}.DocumentWorkflow.get_by_document_id", AsyncMock()) as lookup: - assert await is_current_workflow_run(AsyncMock(), uuid4(), None) is True + assert await writer_skip_reason(db(), uuid4(), None) is None lookup.assert_not_awaited() - async def test_a_document_without_a_workflow_row_is_current(self): + async def test_a_document_without_a_workflow_row_may_write(self): with patch(f"{MODULE}.DocumentWorkflow.get_by_document_id", AsyncMock(return_value=None)): - assert await is_current_workflow_run(AsyncMock(), uuid4(), str(uuid4())) is True + assert await writer_skip_reason(db(), uuid4(), str(uuid4())) is None