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("
");
+ });
+
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