From 5fef85cfb23b186c116b8ff52a387717e797ef35 Mon Sep 17 00:00:00 2001 From: Paco POR-CORREO Date: Wed, 16 Sep 2026 18:35:18 +0200 Subject: [PATCH] feat(ocr): durable review indexing production wiring Wire the complete OCR review and indexing production pipeline: durable OCR result handoff before remote deletion (17a), native-page evidence and quality-gated candidate composition (17b), review images and restart-safe candidate loading (17c), transactional approval and rejection decisions (17d), reviewed-artifact indexing store with exact count verification (17e), and production approve-to-ready wiring with fail-closed OCR_INDEXING_UNAVAILABLE (17f). Activation remains a separately authorized operation; task 7.4 stays pending. --- docs/HISTORIAL_SESIONES.md | 47 ++- ocr-service/app/main.py | 25 +- ocr-service/app/render.py | 8 +- ocr-service/tests/test_api.py | 9 + ocr-service/tests/test_render.py | 2 + .../ocr-ingest-integration/apply-progress.md | 100 +++++ src/api/openapi.ts | 18 +- src/app.ts | 52 ++- src/modules/catalog/reconciler.ts | 5 +- src/modules/catalog/repository.ts | 69 +++- src/modules/ingest/service.ts | 8 +- src/modules/ocr/artifacts.ts | 357 +++++++++++++++++- src/modules/ocr/client.ts | 22 +- src/modules/ocr/dispatcher.ts | 29 +- src/modules/ocr/indexing.ts | 134 +++++++ src/modules/ocr/review.ts | 201 ++++++++-- tests/catalog/repository-ocr.test.ts | 39 ++ tests/ocr/approval-indexing-wiring.test.ts | 81 ++++ tests/ocr/client.test.ts | 133 ++++++- tests/ocr/contracts-deploy.test.ts | 4 +- tests/ocr/dispatcher.test.ts | 92 ++++- tests/ocr/e2e.test.ts | 20 +- tests/ocr/indexing-store.test.ts | 110 ++++++ tests/ocr/review.test.ts | 155 +++++++- 24 files changed, 1633 insertions(+), 87 deletions(-) create mode 100644 tests/ocr/approval-indexing-wiring.test.ts create mode 100644 tests/ocr/indexing-store.test.ts diff --git a/docs/HISTORIAL_SESIONES.md b/docs/HISTORIAL_SESIONES.md index 40f7ecb..e577d74 100644 --- a/docs/HISTORIAL_SESIONES.md +++ b/docs/HISTORIAL_SESIONES.md @@ -3,13 +3,58 @@ **Proyecto:** Workspace de tools IA para empresas **Modulo:** RAG **Ultima actualizacion:** 2026-09-16 -**Ultima modificacion por:** Subagente Deteccion PDF Mixto +**Ultima modificacion por:** Subagent Transactional Review Decisions **Estado:** Activo --- ## Registro de sesion +### 2026-09-16 - Orchestrator Inline Unit 17f (OpenAI executor exhausted) +**Agent:** gentle-orchestrator (GLM inline, parent `ses_29bdbd003ffeLrLjUlFgnp08Y7`) +**Work:** Unit 17f: production approve→index wiring — `PostgresOcrIndexingStore`+`OcrReadyIndexingService` constructed with PostgreSQL, approve route indexes after durable approval and returns `ready`; 503 `OCR_INDEXING_UNAVAILABLE` when unconfigured. RED 2/2 → GREEN 2/2; E2E updated to approve→index→ready; canonical 99/99; check/build/whitespace green. No production/migration/commit; task 7.4 pending. + +--- + +### 2026-09-16 - Orchestrator Inline Unit 17e (OpenAI executor exhausted) +**Agent:** gentle-orchestrator (GLM inline, parent `ses_29bdbd003ffeLrLjUlFgnp08Y7`) +**Work:** Unit 17e: `PostgresOcrIndexingStore` + `OcrReadyIndexingService` — durable reviewed-artifact indexing with canonical chunking, embeddings, versioned Qdrant points, exact count verification, fail-closed identity/embedding/corruption; `indexing→ready` only, no activation/active-pointer. RED module-not-found → GREEN 2/2; canonical 97/97; check/build/whitespace green. No production/migration/commit; task 7.4 pending. + +--- + +### 2026-09-16 - Subagent Transactional Review Decisions +**Agent:** Subagent Transactional Review Decisions · **Session:** `ses_f5534595effeAPCqGLe5exfvl3` (sub of `ses_29bdbd003ffeLrLjUlFgnp08Y7`) +**Work:** Unit 17d: locked PostgreSQL approval/rejection, immutable reviewed-page publication/readback, durable corrections/rejection reasons, production wiring stopping approval at `indexing`. Strict-TDD: RED store/route, focused 16/16, runtime 2/2, canonical 95/95, check/build green. No indexing/Qdrant/activation/production/migration/commit; task 7.4 pending. + +--- + +### 2026-09-16 - Subagente Review Candidate Loader OCR +**Agent:** Subagente Review Candidate Loader OCR · **Model:** openai/gpt-5.6-sol · **Session:** `ses_f554bd29cffekII6j1rSUxMkW3` (subagent of `ses_29bdbd003ffeLrLjUlFgnp08Y7`) +**Responsibility:** Implement only Unit 17c durable review images and authenticated restart-safe production candidate loading without decisions, indexing, activation, or production access. +**Work:** Added authenticated OCR page-image transfer, immutable private review-image artifacts, lifecycle-bound durable candidate reconstruction, production review GET and OpenAPI wiring, and containment-safe authenticated image serving. Production approval/rejection remain unavailable. +**Validation:** Strict-TDD safety 37/37 Node and 16/16 Python; genuine Node/Python RED; focused 40/40, relevant 48/48, canonical Node 92/92, offline Python 16/16, check/build/whitespace/cleanup green. +**Files:** OCR API/client/artifacts/review, catalog repository, app wiring, focused tests, OpenSpec apply progress, and this history. Task 7.4 remains pending; no production or decision mutation occurred. + +--- + +### 2026-09-16 - Subagente Candidate Composition OCR +**Agent:** Subagente Candidate Composition OCR · **Model:** openai/gpt-5.6-sol · **Session:** `ses_f55bdc181ffeitgdlkjunLiCCt` (subagent of `ses_29bdbd003ffeLrLjUlFgnp08Y7`) +**Responsibility:** Implement only Unit 17b native-page evidence and quality-gated restart-safe candidate composition before `review_required`. +**Work:** Added private native evidence, OCR ink coverage, immutable composed-candidate artifacts, exact lifecycle hash/metrics persistence, fail-closed boundaries, and reconciler recovery after completed OCR jobs. +**Validation:** Strict-TDD safety 25/25; genuine Node/Python RED; focused 36/36, runtime 2/2, regressions 15/15, canonical Node 89/89, offline Python 16/16, check/build/whitespace/cleanup green. +**Files:** OCR artifacts/client/dispatcher, catalog repository/reconciler, ingest/app wiring, OCR renderer, focused tests, OpenSpec apply progress, and this history. Task 7.4 and production review wiring remain pending. + +--- + +### 2026-09-16 - Subagente Handoff Durable OCR +**Agent:** Subagente Handoff Durable OCR · **Model:** openai/gpt-5.6-sol · **Session:** `ses_f55dc1c2effeK3U6mfC2NNOmGY` (subagent of `ses_29bdbd003ffeLrLjUlFgnp08Y7`) +**Responsibility:** Implement only Unit 17a durable validated OCR-result persistence and fail-closed dispatcher ordering, without production access or task 7.4 completion. +**Work:** Added canonical private OCR-result artifacts with immutable atomic publication, exact restart readback, schema/content/identity/permission validation, and dispatcher/application ordering before database completion, review transition, and remote deletion. +**Validation:** Strict-TDD safety net 17/17; genuine RED 8/11; focused GREEN/refactor 19/19; relevant regression 21/21; canonical Node 84/84; check/build/whitespace and process cleanup passed. +**Files:** `src/app.ts`, `src/modules/ocr/{artifacts,client,dispatcher}.ts`, `tests/ocr/{client,dispatcher,e2e}.test.ts`, OpenSpec apply progress, and this history. Task 7.4 remains pending; candidate/active versions and production were untouched. + +--- + ### 2026-09-16 - Subagente Deteccion PDF Mixto **Agent:** Subagente Deteccion PDF Mixto · **Model:** openai/gpt-5.6-sol · **Session:** `ses_f5640de36ffekplsg06rgaNmX9` (subagent of `ses_29bdbd003ffeLrLjUlFgnp08Y7`) **Responsibility:** Implement only Unit 16 raster-aware PDF detection and routing without production calls, candidate mutation, or task 7.4 completion. diff --git a/ocr-service/app/main.py b/ocr-service/app/main.py index 11ede49..2fabdd8 100644 --- a/ocr-service/app/main.py +++ b/ocr-service/app/main.py @@ -13,7 +13,7 @@ from fastapi import BackgroundTasks, Depends, FastAPI, File, Form, Header, HTTPE from .engine import OcrEngine from .models import load_runtime_engine -from .render import process_pdf +from .render import PdfRenderError, process_pdf, render_pdf_pages MAX_UPLOAD_BYTES = 50 * 1024 * 1024 @@ -123,6 +123,17 @@ class JobQueue: row = self.connection.execute("SELECT status,result FROM jobs WHERE job_id=?", (job_id,)).fetchone() return (row["status"], json.loads(row["result"]) if row["result"] else None) if row else None + def image(self, job_id: str, page: int) -> tuple[str, bytes] | None: + row = self.connection.execute("SELECT status,document_sha256,pdf FROM jobs WHERE job_id=?", (job_id,)).fetchone() + if not row: + return None + if row["status"] != "succeeded" or row["pdf"] is None: + fail(409, "RESULT_NOT_READY", "OCR review image is not ready") + try: + return row["document_sha256"], render_pdf_pages(bytes(row["pdf"]), [page])[0].png + except PdfRenderError as error: + fail(422, "INVALID_PAGE", str(error)) + def delete(self, job_id: str) -> None: with self.lock: self.connection.execute("DELETE FROM jobs WHERE job_id=?", (job_id,)) @@ -217,6 +228,18 @@ def create_app( fail(409, "RESULT_NOT_READY", "OCR job result is not ready") return stored[1] + @application.get("/v1/jobs/{job_id}/pages/{page}/image", dependencies=[Depends(authorize)]) + def get_review_image(job_id: str, page: int) -> Response: + stored = queue.image(job_id, page) + if stored is None: + fail(404, "JOB_NOT_FOUND", "OCR job does not exist") + document_sha256, png = stored + return Response(png, media_type="image/png", headers={ + "X-Document-Sha256": document_sha256, + "X-Content-Sha256": hashlib.sha256(png).hexdigest(), + "X-Page-Number": str(page), + }) + @application.delete("/v1/jobs/{job_id}", status_code=204, dependencies=[Depends(authorize)]) def delete_job(job_id: str) -> Response: queue.delete(job_id) diff --git a/ocr-service/app/render.py b/ocr-service/app/render.py index 4848f52..673d0ef 100644 --- a/ocr-service/app/render.py +++ b/ocr-service/app/render.py @@ -62,11 +62,15 @@ def render_pdf_pages(pdf: bytes, pages: list[int], max_pixels: int = MAX_RENDER_ return rendered -def _metrics(lines: list[EngineLine], text: str) -> dict[str, int | float]: +def _metrics(lines: list[EngineLine], text: str, image: object) -> dict[str, int | float]: confidences = sorted(line.confidence for line in lines) + grayscale = image.convert("L") + grayscale.thumbnail((256, 256)) + pixels = list(grayscale.getdata()) return { "lineCount": len(lines), "nonWhitespaceCharacters": sum(not character.isspace() for character in text), + "inkCoverage": sum(value < 250 for value in pixels) / len(pixels), "medianConfidence": statistics.median(confidences) if confidences else 0.0, "p10Confidence": confidences[math.floor((len(confidences) - 1) * 0.1)] if confidences else 0.0, "lowConfidenceLineRatio": sum(value < 0.5 for value in confidences) / len(confidences) if confidences else 0.0, @@ -104,7 +108,7 @@ def process_pdf( "height": rendered.height, "processingMs": processing_ms(rendered.page) if processing_ms else elapsed, "text": text, - "metrics": _metrics(lines, text), + "metrics": _metrics(lines, text, rendered.image), "lines": serialized_lines, }) return { diff --git a/ocr-service/tests/test_api.py b/ocr-service/tests/test_api.py index 1fe7d7f..689199e 100644 --- a/ocr-service/tests/test_api.py +++ b/ocr-service/tests/test_api.py @@ -178,6 +178,15 @@ def test_auth_accepted_job_executes_and_exposes_integrity_bound_result(tmp_path: assert (result.json()["jobId"], result.json()["documentSha256"]) == (job_id, hashlib.sha256(pdf).hexdigest()) assert [(page["page"], page["text"]) for page in result.json()["pages"]] == [(1, "Factura FAT07")] + image = client.get(f"/v1/jobs/{job_id}/pages/1/image", headers={"Authorization": f"Bearer {TOKEN}"}) + assert image.status_code == 200 + assert image.headers["content-type"] == "image/png" + assert image.headers["x-document-sha256"] == hashlib.sha256(pdf).hexdigest() + assert image.headers["x-content-sha256"] == hashlib.sha256(image.content).hexdigest() + assert image.content.startswith(b"\x89PNG\r\n\x1a\n") + assert client.get(f"/v1/jobs/{job_id}/pages/1/image").status_code == 401 + assert client.get(f"/v1/jobs/{job_id}/pages/4/image", headers={"Authorization": f"Bearer {TOKEN}"}).status_code == 422 + def test_auth_result_rejects_unknown_not_ready_and_unauthorized_jobs(client: TestClient): job_id = submit(client, request_for()).json()["jobId"] diff --git a/ocr-service/tests/test_render.py b/ocr-service/tests/test_render.py index 628e622..56d6542 100644 --- a/ocr-service/tests/test_render.py +++ b/ocr-service/tests/test_render.py @@ -60,6 +60,7 @@ def test_render_builds_repeatable_result_schema_with_deterministic_engine() -> N assert first == run() assert first["schemaVersion"] == "1" assert first["documentSha256"] == DOCUMENT_SHA256 + assert 0 < first["pages"][0]["metrics"]["inkCoverage"] < 1 assert first["engine"] == { "name": "paddleocr", "version": "3.4.0", @@ -77,6 +78,7 @@ def test_render_builds_repeatable_result_schema_with_deterministic_engine() -> N "metrics": { "lineCount": 2, "nonWhitespaceCharacters": 15, + "inkCoverage": first["pages"][0]["metrics"]["inkCoverage"], "medianConfidence": 0.86, "p10Confidence": 0.74, "lowConfidenceLineRatio": 0.0, diff --git a/openspec/changes/ocr-ingest-integration/apply-progress.md b/openspec/changes/ocr-ingest-integration/apply-progress.md index ea3c2d7..ddb6b74 100644 --- a/openspec/changes/ocr-ingest-integration/apply-progress.md +++ b/openspec/changes/ocr-ingest-integration/apply-progress.md @@ -486,3 +486,103 @@ None — the migration follows the proposal, specifications, design, and closed | Focused/canonical | Focused 11/11; canonical Node 82/82; check/build/whitespace clean. | | Runtime harness | Real 25-page PDF reported raster coverage 1 on every page, selected pages 1–25, and native text still lacked CBG04a/FAT07/DSAU08/NSAV06. | | Rollback boundary | Revert Unit 16 parser/detection/ingest/test/contract deltas and this metadata; preserve Units 1–15, candidates, and active content. | + +## Unit 17a: Durable OCR Result Handoff +- Persisted each fully validated OCR result as canonical private content before PostgreSQL job completion, review transition, or best-effort remote deletion. The envelope binds version, document, remote job, source hash, engine, pages, full page text, metrics, and line text/confidence/boxes. +- Added immediate and restart-safe exact readback with schema, content hash, expected identity, permission, and optional artifact-byte hash validation. Immutable retries accept only byte-identical content; corruption and identity mismatches fail closed. +- Task 7.4 remains pending; no candidate composition, review endpoint restoration, activation, production wiring execution, migration, deployment, commit, or push occurred. + +### TDD Cycle Evidence +| Task | Test Files | Safety Net | RED | GREEN / TRIANGULATE | REFACTOR | +|---|---|---|---|---|---| +| Unit 17a | `tests/ocr/{client,dispatcher}.test.ts` | Existing focused suites 17/17 | 8/11: missing artifact exports and dispatcher skipped the durable-write boundary | Focused 19/19; success ordering, simulated remote deletion/restart, repeat write, corruption, identity mismatch, and write failure | Centralized artifact integrity errors; focused 19/19 remained green | + +### Work Unit Evidence +| Evidence | Result | +|---|---| +| Focused/regression | Focused client/dispatcher 19/19; client/dispatcher/E2E regression 21/21; canonical Node 84/84; check/build/whitespace clean. | +| Runtime harness | Real temporary private artifacts survived simulated remote deletion and exact readback after a restart boundary; write failure occurred before job completion/review/deletion and left the active version untouched. | +| Rollback boundary | Revert Unit 17a deltas in `src/app.ts`, `src/modules/ocr/{artifacts,client,dispatcher}.ts`, `tests/ocr/{client,dispatcher,e2e}.test.ts`, and this progress/history metadata; preserve Units 1–16 and task 7.4. | + +### Boundary and Accounting +- Forecast: 250 changed lines against the authorized 360-line cap. Actual: 227 additions plus deletions across source, tests, progress, and history. +- Unit 17a does not restore the production review endpoint and cannot by itself satisfy the mandatory failed-evidence remediation binding. + +## Unit 17b: Native Evidence and Candidate Composition +- Persisted private canonical native-page evidence beside each original and bound it to version, document, original hash, selected OCR pages, page text hashes, and raster metrics. +- Added OCR `inkCoverage`, exact restart readback of native/OCR evidence, quality classification, existing-domain candidate composition, immutable candidate publication/readback, and lifecycle hash/metrics agreement before `review_required`. +- Added reconciler recovery for the crash window after all OCR jobs succeed but before candidate finalization. Any missing, corrupt, mismatched, blocked, conflicting, or failed evidence leaves the version fail-closed and the remote OCR result undeleted. + +### TDD Cycle Evidence +| Task | Safety Net | RED | GREEN / TRIANGULATE | REFACTOR | +|---|---|---|---|---| +| Unit 17b | Focused OCR/catalog suites 25/25 | Node failed on missing candidate exports/repository/recovery and skipped finalization; Python failed on missing `inkCoverage` | Focused 36/36 plus runtime 2/2; native/OCR/blank, corruption, missing input, quality block, immutable conflict, readback mismatch, and restart recovery | Low-resolution ink sampling and shared private JSON validation; all gates remained green | + +### Work Unit Evidence +| Evidence | Result | +|---|---| +| Focused/regression | Focused 36/36; review/E2E/retention 15/15; canonical Node 89/89; offline Python 16/16. | +| Runtime harness | Bounded filesystem/restart scenarios passed 2/2: mixed native/OCR/blank evidence survived exact readback, and a completed version resumed candidate finalization after restart. | +| Gates and cleanup | `npm run check`, `npm run build`, tracked/untracked whitespace passed; temporary artifact directories absent; final process check found no test, Uvicorn, or pytest process. | +| Rollback boundary | Revert only Unit 17b deltas in OCR render/client/artifacts/orchestration/repository/reconciler/ingest code, focused tests, and this metadata; preserve Unit 17a durable results and Units 1–16. | + +### Boundary and Accounting +- Forecast: 370 changed lines against the authorized 400-line cap. Actual: 400 additions plus deletions against native begin tree `787b4b55c4d86e401a8b1f195173e3ed0e0e4d94`. +- Task 7.4 remains pending; no candidate or active version, production service, secret, deployment, commit, or push was touched. Production review construction remains unwired, so Unit 17b does not claim the broader production review `503` is fixed. + +## Unit 17c: Review Artifacts and Production Candidate Loading +- Added authenticated OCR-service page-image transfer and immutable private `0600` review-image artifacts with version/document/page/hash metadata and containment-safe resolution. +- Added a concrete PostgreSQL lifecycle-context read path and restart-safe durable candidate reader that validates state, identity, page hashes, metrics, risks, image permissions, and image hashes before returning review data. +- Wired and documented production review GET and protected image GET routes while intentionally leaving production approval/rejection/indexing unavailable. Task 7.4 remains pending. + +### TDD Cycle Evidence +| Task | Safety Net | RED | GREEN / TRIANGULATE | REFACTOR | +|---|---|---|---|---| +| Unit 17c | Node OCR/catalog 37/37; Python 16/16 | Node failed on missing image client/artifact/reader exports and absent OpenAPI image path; Python image route returned 404 | Focused 40/40; relevant 48/48; contract 3/3; authenticated image transfer, restart loading, lifecycle mismatch, corruption, traversal-shaped identity, disabled/unauthorized and decision isolation | Consolidated identifier imports and native-document validation; focused suites remained green | + +### Work Unit Evidence +| Evidence | Result | +|---|---| +| Focused/regression | Focused 40/40; relevant OCR/catalog 48/48; image contract 3/3; canonical Node 92/92; offline Python 16/16. | +| Runtime harness | Bounded localhost Express dispatch persisted result and all review images before completion, composed the candidate, entered review, restarted the reader, served authenticated review/image bytes, and kept approval/rejection at `503`; FastAPI integration returned authenticated PNG identity/hash headers and rejected unauthorized/invalid pages. | +| Gates and cleanup | `npm run check`, `npm run build`, canonical tests, offline Python, tracked/untracked whitespace, temporary-artifact cleanup, and final process inspection passed. | +| Rollback boundary | Revert only Unit 17c deltas in OCR API/client/artifacts/review/catalog/app code, focused tests, and this metadata; preserve Units 17a–17b and task 7.4. | + +### Boundary and Accounting +- Unit 17c is 377 additions plus deletions against native begin tree `b64f46c0bb4edd8f9878b3920bf4d3372a29073a`, below the 400-line cap. +- No approval, rejection, correction, embedding, Qdrant, activation, production, SSH, secret, migration, commit, or push operation occurred. + +## Unit 17d: Transactional Review Decisions +- Added a concrete PostgreSQL review store that locks the candidate/source/page lifecycle rows, reloads durable evidence, and revalidates version, active version, candidate, page, line, and correction identities before any decision write. +- Approval publishes and exactly rereads canonical private immutable `reviewed-pages.json`, records all corrections and reviewed page hashes transactionally, then transitions only `review_required → indexing`. Rejection records its bounded reason transactionally and creates no reviewed artifact. +- Production approval/rejection routes now use the durable store when PostgreSQL is configured. Approval stops at `indexing`; neither decision route invokes indexing, embeddings, Qdrant, activation, or the active pointer. Task 7.4 remains pending. + +### TDD Cycle Evidence +| Task | Safety Net | RED | GREEN / TRIANGULATE | REFACTOR | +|---|---|---|---|---| +| Unit 17d | Review/client 20/20 | Missing store/artifact exports failed module loading; route/contract RED failed 2/5 | Focused review/E2E/contracts 16/16; transactional runtime 2/2; approval, rejection, replay, stale state, corrupt path, auth/off gates, and zero-index boundary | Shared correction preparation and immutable publication cleanup; focused suites remained green | + +### Work Unit Evidence +- Review/E2E/contracts 16/16; canonical Node 95/95; check/build passed; fake-PostgreSQL transactional runtime 2/2; rejection produced no reviewed artifact; approval stopped at `indexing` with zero indexing calls. Python unchanged. +- Rollback: revert only Unit 17d deltas in `src/modules/ocr/{artifacts,review}.ts`, `src/{app,api/openapi}.ts`, the three focused tests, and this metadata; preserve Units 17a–17c and task 7.4. + +### Boundary and Accounting +- Direct delta against native begin tree `5ab7b1eb5fe07cf2025e36b35fcd132f9428eb01`: within the 400-line cap after metadata condensation. Stale/mismatched/replayed/partial/corrupt decisions fail closed; transaction rollback removes newly published reviewed artifacts on downstream failure. +- Migrations 001/002 unchanged. No indexing, embedding, Qdrant, activation, active-pointer, production, SSH, secret, re-OCR, deployment, commit, push, or task 7.4 operation occurred. + +## Unit 17e: Reviewed Indexing Store +- Added `PostgresOcrIndexingStore` (`indexing.ts`): loads composed candidate + reviewed-pages artifacts with full identity/integrity validation, chunks reviewed text with the canonical `documentalChunkingPolicy`, embeds via the existing provider, writes canonical versioned Qdrant points, and verifies exact point count before `indexing → ready`. +- Added `OcrReadyIndexingService`: indexes approved reviewed candidates to `ready` only; activation/settlement methods fail closed as reserved for a separately authorized unit. No active-pointer changes. +- Strict TDD: RED from missing exports (module-not-found); GREEN focused `tests/ocr/indexing-store.test.ts` 2/2 (restart-safe indexing, canonical point IDs, count verification, fail-closed identity/embedding/count/corruption). Canonical Node 97/97, check/build/whitespace green. Python unchanged. +- Rollback: revert `src/modules/ocr/indexing.ts`, remove `tests/ocr/indexing-store.test.ts`, and this metadata; preserve Units 17a–17d and task 7.4. + +### Boundary and Accounting +- Direct delta against begin tree `f351a3d4514d89c9208c72a93be3f9411b41bedd`: within the 400-line cap. No migration, Qdrant production writes, activation, active-pointer, production/SSH/secret access, deploy, commit, or push. Task 7.4 remains pending. + +## Unit 17f: Production Approval→Indexing Wiring +- Production `createApp` now constructs `PostgresOcrIndexingStore` + `OcrReadyIndexingService` when PostgreSQL is configured (overridable via `options.indexingService`), and the authenticated approve route invokes indexing after a successful durable approval, returning the resulting `ready` state; without a configured indexing service it fails closed with 503 `OCR_INDEXING_UNAVAILABLE`. Rejection stays decision-only; activation remains reserved. +- Strict TDD: genuine RED 2/2 (approve returned `indexing` without indexing; missing-service returned 500), GREEN focused wiring 2/2; E2E contract updated to approve→index→ready with exactly one indexing call. Canonical Node 99/99, check/build/whitespace green. Python unchanged. +- Rollback: revert `src/app.ts` wiring, `tests/ocr/approval-indexing-wiring.test.ts`, E2E delta, and this metadata; preserve Units 17a–17e and task 7.4. + +### Boundary and Accounting +- No migration, activation, active-pointer, production/SSH/secret access, deploy, commit, or push. Task 7.4 remains pending. diff --git a/src/api/openapi.ts b/src/api/openapi.ts index 3f8f094..3808269 100644 --- a/src/api/openapi.ts +++ b/src/api/openapi.ts @@ -280,6 +280,22 @@ export const openApiDocument = { } } }, + "/ingestions/{versionId}/documents/{documentId}/pages/{page}/image": { + get: { + tags: ["OCR Review"], summary: "Get a private OCR review page image", security: [{ bearerAuth: [] }], + parameters: [ + { name: "versionId", in: "path", required: true, schema: { type: "string", format: "uuid" } }, + { name: "documentId", in: "path", required: true, schema: { type: "string" } }, + { name: "page", in: "path", required: true, schema: { type: "integer", minimum: 1 } } + ], + responses: { + "200": { description: "Integrity-validated private PNG review image.", content: { "image/png": { schema: { type: "string", format: "binary" } } } }, + "401": jsonResponse("Missing or invalid lifecycle admin token.", ref("Error")), + "404": jsonResponse("OCR is disabled or the review image does not exist.", ref("Error")), + "503": serverError, "500": serverError + } + } + }, "/ingestions/{versionId}/approve": { post: { tags: ["OCR Review"], summary: "Correct and approve an OCR candidate", security: [{ bearerAuth: [] }], @@ -616,7 +632,7 @@ export const openApiDocument = { }, OcrDecisionResponse: { type: "object", required: ["versionId", "state", "activated"], - properties: { versionId: { type: "string", format: "uuid" }, state: { type: "string", enum: ["ready", "active", "rejected"] }, activated: { type: "boolean" }, activatedVersionId: { type: "string", format: "uuid" }, errorCode: { type: "string", enum: ["DUPLICATE_REUSABLE_VERSION"] } } + properties: { versionId: { type: "string", format: "uuid" }, state: { type: "string", enum: ["indexing", "rejected"] }, activated: { type: "boolean" } } }, CleanupRequest: { type: "object", diff --git a/src/app.ts b/src/app.ts index 1bd7778..a79c3bc 100644 --- a/src/app.ts +++ b/src/app.ts @@ -22,9 +22,10 @@ import { KnowledgeLifecycleReconciler } from "./modules/catalog/reconciler.js"; import { CatalogError, CatalogRepository } from "./modules/catalog/repository.js"; import { OcrClient } from "./modules/ocr/client.js"; import { OcrDispatcher } from "./modules/ocr/dispatcher.js"; -import { resolveArtifactPath } from "./modules/ocr/artifacts.js"; -import type { OcrReviewService } from "./modules/ocr/review.js"; +import { persistComposedCandidateArtifact, persistOcrResultArtifact, persistReviewImageArtifacts, readOcrArtifactPageNumbers, resolveArtifactPath } from "./modules/ocr/artifacts.js"; +import { DurableOcrReviewReader, OcrReviewService, PostgresOcrReviewStore } from "./modules/ocr/review.js"; import type { OcrIndexingService } from "./modules/ocr/indexing.js"; +import { OcrReadyIndexingService, PostgresOcrIndexingStore } from "./modules/ocr/indexing.js"; import { OcrRetentionService } from "./modules/ocr/retention.js"; import type { ChatMessage, ChunkMode, RetrieveIntent, RetrieveScope } from "./shared/types/rag.js"; import { sha256Hex } from "./shared/utils/ids.js"; @@ -72,13 +73,30 @@ export function createApp(options: AppOptions = {}) { const bytes = await readFile(resolveArtifactPath(versionDirectory, document.originalPath)); if (sha256Hex(bytes) !== document.originalSha256) throw new Error("OCR artifact integrity validation failed"); return { bytes, documentSha256: document.originalSha256 }; + }, async (job, result) => { + const durable = await persistOcrResultArtifact({ rootDirectory: ocr.artifactRoot, versionId: job.versionId, documentId: job.documentId, result }); + const pages = await readOcrArtifactPageNumbers({ rootDirectory: ocr.artifactRoot, versionId: job.versionId, documentId: job.documentId }); + const images = await Promise.all(pages.map(async (page) => ({ page, ...await ocrClient!.getReviewImage(result.jobId, page, result.documentSha256) }))); + await persistReviewImageArtifacts({ rootDirectory: ocr.artifactRoot, versionId: job.versionId, documentId: job.documentId, images }); + return durable.result; + }, async (versionId) => { + const jobs = await catalog.listOcrJobs(versionId); + const { candidate } = await persistComposedCandidateArtifact({ rootDirectory: ocr.artifactRoot, versionId, jobs }); + await catalog.persistOcrCandidate(versionId, candidate.documents.flatMap(({ documentId, pages }) => pages.map((page) => ({ documentId, ...page })))); }) : undefined; const retention = catalog ? new OcrRetentionService(catalog, ocr.artifactRoot) : undefined; const reconciler = new KnowledgeLifecycleReconciler(catalog, vectorStore, ocrDispatcher, retention); const evaluationLogs = new EvaluationLogService(embeddingProvider); const ingestService = options.ingestService ?? new IngestService(embeddingProvider, vectorStore, catalog, ocr); + const durableReview = catalog ? new DurableOcrReviewReader(catalog, ocr.artifactRoot) : undefined; + const reviewService = options.reviewService ?? (catalogPool && durableReview + ? new OcrReviewService(new PostgresOcrReviewStore(catalogPool, durableReview, ocr.artifactRoot)) + : undefined); const retrieveService = new RetrieveService(embeddingProvider, vectorStore, catalog); const answerService = new AnswerService(retrieveService); + const indexingService = options.indexingService ?? (catalogPool + ? new OcrReadyIndexingService(new PostgresOcrIndexingStore(catalogPool, ocr.artifactRoot, embeddingProvider, vectorStore)) + : undefined); if (options.startReconciler !== false) reconciler.start(); function sendError(res: express.Response, error: unknown, fallback: string) { @@ -269,27 +287,43 @@ export function createApp(options: AppOptions = {}) { app.get("/ingestions/:versionId/review", async (req, res) => { if (!requireOcrEnabled(res)) return; if (!requireLifecycleAdmin(req, res)) return; - if (!options.reviewService) { + const reader = reviewService ?? durableReview; + if (!reader) { res.status(503).json({ ok: false, error: "OCR review service is not configured", code: "OCR_REVIEW_UNAVAILABLE" }); return; } try { - res.json(await options.reviewService.view(String(req.params.versionId))); + res.json(await reader.view(String(req.params.versionId))); } catch (error) { sendError(res, error, "Unknown OCR review error"); } }); + app.get("/ingestions/:versionId/documents/:documentId/pages/:page/image", async (req, res) => { + if (!requireOcrEnabled(res)) return; + if (!requireLifecycleAdmin(req, res)) return; + if (!durableReview) return void res.status(503).json({ ok: false, error: "OCR review service is not configured", code: "OCR_REVIEW_UNAVAILABLE" }); + try { + const image = await durableReview.image(String(req.params.versionId), String(req.params.documentId), Number(req.params.page)); + res.set({ "Content-Type": image.mimeType, "Cache-Control": "private, no-store", "X-Content-Sha256": image.sha256 }).send(image.bytes); + } catch (error) { sendError(res, error, "Unknown OCR review image error"); } + }); + app.post("/ingestions/:versionId/approve", async (req, res) => { if (!requireOcrEnabled(res)) return; if (!requireLifecycleAdmin(req, res)) return; - if (!options.reviewService || !options.indexingService) { + if (!reviewService) { res.status(503).json({ ok: false, error: "OCR approval service is not configured", code: "OCR_REVIEW_UNAVAILABLE" }); return; } try { - const approved = await options.reviewService.approve(String(req.params.versionId), req.body); - res.json(await options.indexingService.index(approved)); + const approved = await reviewService.approve(String(req.params.versionId), req.body); + if (!indexingService) { + res.status(503).json({ ok: false, error: "OCR indexing service is not configured", code: "OCR_INDEXING_UNAVAILABLE" }); + return; + } + const indexed = await indexingService.index(approved); + res.json({ versionId: approved.versionId, state: indexed.state, activated: indexed.activated }); } catch (error) { sendError(res, error, "Unknown OCR approval error"); } @@ -298,12 +332,12 @@ export function createApp(options: AppOptions = {}) { app.post("/ingestions/:versionId/reject", async (req, res) => { if (!requireOcrEnabled(res)) return; if (!requireLifecycleAdmin(req, res)) return; - if (!options.reviewService) { + if (!reviewService) { res.status(503).json({ ok: false, error: "OCR review service is not configured", code: "OCR_REVIEW_UNAVAILABLE" }); return; } try { - res.json(await options.reviewService.reject(String(req.params.versionId), req.body)); + res.json(await reviewService.reject(String(req.params.versionId), req.body)); } catch (error) { sendError(res, error, "Unknown OCR rejection error"); } diff --git a/src/modules/catalog/reconciler.ts b/src/modules/catalog/reconciler.ts index 4b6cc34..dfc5e71 100644 --- a/src/modules/catalog/reconciler.ts +++ b/src/modules/catalog/reconciler.ts @@ -12,6 +12,7 @@ export interface ReconcilerStatus { orphanedVersionsRecovered?: number; orphanedVersionsFailed?: number; ocrLeasesRecovered?: number; + ocrCandidatesRecovered?: number; ocrRetentionDeleted?: number; invariantViolations?: string[]; } @@ -23,7 +24,7 @@ export class KnowledgeLifecycleReconciler { constructor( private readonly catalog: CatalogRepository | undefined, private readonly vectorStore: VectorStoreClient, - private readonly ocrDispatcher?: Pick, + private readonly ocrDispatcher?: Pick, private readonly ocrRetention?: Pick ) {} @@ -57,6 +58,7 @@ export class KnowledgeLifecycleReconciler { const invariantViolations = await catalog.validateActiveInvariant(); const ocrLeasesRecovered = await this.ocrDispatcher?.recoverExpiredLeases() ?? 0; await this.ocrDispatcher?.dispatchAvailable(); + const ocrCandidatesRecovered = await this.ocrDispatcher?.recoverCompletedCandidates() ?? 0; const retention = await this.ocrRetention?.runOnce(); const orphanRecovery = await this.recoverOrphanedVersions(catalog); @@ -81,6 +83,7 @@ export class KnowledgeLifecycleReconciler { orphanedVersionsRecovered: orphanRecovery.ready, orphanedVersionsFailed: orphanRecovery.failed, ocrLeasesRecovered, + ocrCandidatesRecovered, ocrRetentionDeleted: (retention?.deleted ?? 0) + (retention?.resumed ?? 0), invariantViolations }; diff --git a/src/modules/catalog/repository.ts b/src/modules/catalog/repository.ts index 811ad46..3dfb11c 100644 --- a/src/modules/catalog/repository.ts +++ b/src/modules/catalog/repository.ts @@ -5,6 +5,7 @@ import { CatalogError } from "./errors.js"; import { normalizeExpectedActiveVersion } from "./lifecycle.js"; import type { OcrResult } from "../ocr/client.js"; import type { OcrRetentionCandidate } from "../ocr/retention.js"; +import type { OcrReviewContext } from "../ocr/review.js"; export interface CatalogDocumentInput { documentId: string; @@ -395,6 +396,24 @@ export class CatalogRepository { return result.rows.map(toOcrJob); } + async listOcrJobs(versionId: string): Promise { + const result = await this.pool.query( + `SELECT ${ocrJobColumns()} FROM rag_ocr_jobs WHERE version_id = $1 ORDER BY document_id`, [versionId] + ); + return result.rows.map(toOcrJob); + } + + async listOcrVersionsAwaitingCandidate(): Promise { + const result = await this.pool.query<{ version_id: string }>( + `SELECT DISTINCT v.version_id FROM rag_source_versions v + WHERE v.state = 'indexing' AND EXISTS (SELECT 1 FROM rag_ocr_jobs j WHERE j.version_id = v.version_id) + AND NOT EXISTS (SELECT 1 FROM rag_ocr_jobs j WHERE j.version_id = v.version_id AND j.state <> 'succeeded') + AND EXISTS (SELECT 1 FROM rag_document_pages p WHERE p.version_id = v.version_id AND p.candidate_text_hash IS NULL) + ORDER BY v.version_id` + ); + return result.rows.map(({ version_id }) => version_id); + } + async setOcrRemoteJob(jobId: string, remoteJobId: string, leaseMs: number): Promise { const result = await this.pool.query( `UPDATE rag_ocr_jobs SET remote_job_id = $2, heartbeat_at = now(), @@ -424,7 +443,7 @@ export class CatalogRepository { const { version_id: versionId, document_id: documentId } = identity.rows[0]; for (const page of result.pages) { await client.query( - `UPDATE rag_document_pages SET ocr_text_hash = $4, candidate_text_hash = $4, metrics = $5 + `UPDATE rag_document_pages SET ocr_text_hash = $4, metrics = $5 WHERE version_id = $1 AND document_id = $2 AND page_number = $3 AND extraction_method = 'ocr'`, [versionId, documentId, page.page, sha256Hex(page.text), page.metrics] ); @@ -444,6 +463,25 @@ export class CatalogRepository { }); } + async persistOcrCandidate(versionId: string, pages: Array<{ + documentId: string; page: number; method: "native" | "ocr" | "blank"; nativeTextSha256: string; + ocrTextSha256: string | null; candidateTextSha256: string; metrics: Record; risks: string[]; + }>): Promise { + await withTransaction(this.pool, async (client) => { + for (const page of pages) { + const result = await client.query( + `UPDATE rag_document_pages SET extraction_method = $4, candidate_text_hash = $5, risk_tokens = $6::jsonb + WHERE version_id = $1 AND document_id = $2 AND page_number = $3 AND extraction_method = $7 + AND native_text_hash = $8 AND ocr_text_hash IS NOT DISTINCT FROM $9::char(64) + AND metrics = $10::jsonb AND blocked_reason IS NULL`, + [versionId, page.documentId, page.page, page.method, page.candidateTextSha256, JSON.stringify(page.risks), + page.ocrTextSha256 === null ? "native" : "ocr", page.nativeTextSha256, page.ocrTextSha256, JSON.stringify(page.metrics)] + ); + if (result.rowCount !== 1) throw new CatalogError("OCR candidate lifecycle evidence does not match durable artifacts", 409, "OCR_ARTIFACT_INTEGRITY_FAILED"); + } + }); + } + async failOcrJob(jobId: string, code: string, detail: string): Promise { await this.pool.query( `UPDATE rag_ocr_jobs SET state = 'failed', lease_expires_at = NULL, next_attempt_at = NULL, @@ -531,20 +569,43 @@ export class CatalogRepository { }; } + async loadOcrReviewContext(versionId: string): Promise { + const version = await this.pool.query<{ + version_id: string; source_id: string; state: string; base_active_version_id: string | null; current_active_version_id: string | null; + activate_requested: boolean; processing_fingerprint: string; metadata_hash: string; + }>(`SELECT v.version_id, v.source_id, v.state, v.base_active_version_id, s.active_version_id AS current_active_version_id, + v.activate_requested, v.processing_fingerprint, v.metadata_hash + FROM rag_source_versions v JOIN rag_sources s ON s.source_id = v.source_id WHERE v.version_id = $1`, [versionId]); + if (!version.rowCount) return undefined; + const pages = await this.pool.query<{ + document_id: string; page_number: string | number; native_text_hash: string; ocr_text_hash: string | null; + candidate_text_hash: string | null; metrics: Record; risk_tokens: string[]; + }>(`SELECT document_id, page_number, native_text_hash, ocr_text_hash, candidate_text_hash, metrics, risk_tokens + FROM rag_document_pages WHERE version_id = $1 ORDER BY document_id, page_number`, [versionId]); + if (pages.rows.some(({ candidate_text_hash }) => !candidate_text_hash)) throw new CatalogError("OCR candidate lifecycle evidence is incomplete", 409, "OCR_ARTIFACT_INTEGRITY_FAILED"); + const row = version.rows[0]; + return { versionId: row.version_id, sourceId: row.source_id, state: row.state, baseActiveVersionId: row.base_active_version_id, + currentActiveVersionId: row.current_active_version_id, activateRequested: row.activate_requested, + processingFingerprint: row.processing_fingerprint, metadataHash: row.metadata_hash, + pages: pages.rows.map((page) => ({ documentId: page.document_id, page: Number(page.page_number), nativeTextSha256: page.native_text_hash, + ocrTextSha256: page.ocr_text_hash, candidateTextSha256: page.candidate_text_hash!, metrics: page.metrics, risks: page.risk_tokens })) }; + } + async markReviewRequired(versionId: string): Promise { const result = await this.pool.query( `UPDATE rag_source_versions SET state = 'review_required' WHERE version_id = $1 AND state = 'indexing' AND EXISTS (SELECT 1 FROM rag_ocr_jobs WHERE version_id = $1) AND NOT EXISTS (SELECT 1 FROM rag_ocr_jobs WHERE version_id = $1 AND state <> 'succeeded') - AND NOT EXISTS ( + AND NOT EXISTS ( SELECT 1 FROM rag_version_documents d WHERE d.version_id = $1 AND d.content_hash IS NULL AND NOT EXISTS ( SELECT 1 FROM rag_ocr_jobs j WHERE j.version_id = d.version_id AND j.document_id = d.document_id AND j.state = 'succeeded' - ) - )`, + ) + ) + AND NOT EXISTS (SELECT 1 FROM rag_document_pages WHERE version_id = $1 AND (candidate_text_hash IS NULL OR blocked_reason IS NOT NULL))`, [versionId] ); if (result.rowCount !== 1) throw new CatalogError("Version is not ready for OCR review", 409, "INVALID_VERSION_STATE"); diff --git a/src/modules/ingest/service.ts b/src/modules/ingest/service.ts index 28cc71b..4e23abf 100644 --- a/src/modules/ingest/service.ts +++ b/src/modules/ingest/service.ts @@ -365,7 +365,7 @@ export class IngestService { const parsed = await parseDocument(original.filePath); plans.push({ original, - pages: [], + pages: [{ page: 1, text: parsed.content, rasterCoverage: 0, textSha256: sha256Hex(parsed.content) }], requestedPages: [], contentHash: sha256Hex(normalizeContentForHash(parsed.content)), mimeType: parsed.mimeType, @@ -407,10 +407,12 @@ export class IngestService { rootDirectory: this.ocr!.artifactRoot, versionId, createdAt: new Date().toISOString(), - documents: plans.map(({ original }) => ({ + documents: plans.map(({ original, pages, requestedPages }) => ({ documentId: original.documentId, documentKey: original.documentKey, - bytes: original.bytes + bytes: original.bytes, + pages: pages.map((page) => ({ ...page, rasterCoverage: page.rasterCoverage ?? 0 })), + requestedPages })) }); if (staged.originalManifestHash !== originalManifestHash) { diff --git a/src/modules/ocr/artifacts.ts b/src/modules/ocr/artifacts.ts index 1e22d60..d003b55 100644 --- a/src/modules/ocr/artifacts.ts +++ b/src/modules/ocr/artifacts.ts @@ -1,7 +1,10 @@ import { createHash, randomUUID } from "node:crypto"; -import { link, lstat, mkdir, open, readdir, rm, stat, unlink } from "node:fs/promises"; +import { link, lstat, mkdir, open, readFile, readdir, rm, stat, unlink } from "node:fs/promises"; import path from "node:path"; import { canonicalJson, hashOrderedPairs, sha256Hex } from "../../shared/utils/ids.js"; +import { isValidOcrResult, type OcrResult } from "./client.js"; +import { composeCandidate, type CandidatePage } from "./composition.js"; +import { classifyOcrPage } from "./detection.js"; const UUID = /^[0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/iu; const UUID_NAMESPACE_URL = Buffer.from("6ba7b8119dad11d180b400c04fd430c8", "hex"); @@ -10,7 +13,13 @@ interface StageInput { rootDirectory: string; versionId: string; createdAt: string; - documents: Array<{ documentId: string; documentKey: string; bytes: Buffer }>; + documents: Array<{ + documentId: string; + documentKey: string; + bytes: Buffer; + pages?: Array<{ page: number; text: string; rasterCoverage?: number; textSha256: string }>; + requestedPages?: number[]; + }>; } export function resolveArtifactPath(versionDirectory: string, relativePath: string): string { @@ -45,12 +54,23 @@ export async function stageOcrArtifacts(input: StageInput): Promise<{ await mkdir(directory, { recursive: true, mode: 0o700 }); const originalPath = path.posix.join(relativeDirectory, "original.pdf"); await durableWrite(resolveArtifactPath(versionDirectory, originalPath), document.bytes); + const nativePagesPath = path.posix.join(relativeDirectory, "native-pages.json"); + const nativePages = { + schemaVersion: "1", versionId: input.versionId, documentId: document.documentId, + documentSha256: sha256Hex(document.bytes), requestedPages: document.requestedPages ?? [], + pages: document.pages?.map((page) => ({ ...page, rasterCoverage: page.rasterCoverage ?? 0 })) ?? [] + }; + validateNativePages(nativePages); + const nativeBytes = Buffer.from(canonicalJson(nativePages)); + await durableWrite(resolveArtifactPath(versionDirectory, nativePagesPath), nativeBytes); documents.push({ documentId: document.documentId, documentKey: document.documentKey, documentArtifactId, originalPath, - originalSha256: sha256Hex(document.bytes) + originalSha256: sha256Hex(document.bytes), + nativePagesPath, + nativePagesSha256: sha256Hex(nativeBytes) }); } const originalManifestHash = hashOrderedPairs(documents.map(({ documentKey, originalSha256 }) => [documentKey, originalSha256]))!; @@ -67,6 +87,253 @@ export async function stageOcrArtifacts(input: StageInput): Promise<{ } } +interface OcrResultArtifactInput { + rootDirectory: string; + versionId: string; + documentId: string; + result: OcrResult; +} + +interface OcrResultReadInput extends Omit { + jobId: string; + documentSha256: string; + pages: number[]; + artifactSha256?: string; +} + +export async function persistOcrResultArtifact(input: OcrResultArtifactInput): Promise<{ + artifactPath: string; + artifactSha256: string; + result: OcrResult; +}> { + assertUuid(input.versionId); + const pages = input.result.pages.map(({ page }) => page); + if (!isValidOcrResult(input.result, input.result.jobId, { documentSha256: input.result.documentSha256, pages })) { + throw resultIntegrityError(); + } + const resultSha256 = sha256Hex(canonicalJson(input.result)); + const serialized = canonicalJson({ + schemaVersion: "1", + versionId: input.versionId, + documentId: input.documentId, + resultSha256, + result: input.result + }); + const artifactPath = resultArtifactPath(input.rootDirectory, input.versionId, input.documentId); + await publishImmutable(artifactPath, Buffer.from(serialized)); + const artifactSha256 = sha256Hex(serialized); + const result = await readOcrResultArtifact({ + rootDirectory: input.rootDirectory, + versionId: input.versionId, + documentId: input.documentId, + jobId: input.result.jobId, + documentSha256: input.result.documentSha256, + pages, + artifactSha256 + }); + return { artifactPath, artifactSha256, result }; +} + +export async function readOcrResultArtifact(input: OcrResultReadInput): Promise { + assertUuid(input.versionId); + const artifactPath = resultArtifactPath(input.rootDirectory, input.versionId, input.documentId); + const file = await lstat(artifactPath).catch(() => { throw resultIntegrityError(); }); + if (!file.isFile() || (file.mode & 0o777) !== 0o600) throw resultIntegrityError(); + const bytes = await readFile(artifactPath); + if (input.artifactSha256 && sha256Hex(bytes) !== input.artifactSha256) throw resultIntegrityError(); + let envelope: unknown; + try { envelope = JSON.parse(bytes.toString("utf8")); } catch { throw resultIntegrityError(); } + if (!isRecord(envelope) || envelope.schemaVersion !== "1" || envelope.versionId !== input.versionId || envelope.documentId !== input.documentId) { + throw new Error("OCR result artifact identity validation failed"); + } + if (typeof envelope.resultSha256 !== "string" || sha256Hex(canonicalJson(envelope.result)) !== envelope.resultSha256 + || !isValidOcrResult(envelope.result, input.jobId, { documentSha256: input.documentSha256, pages: input.pages })) { + throw resultIntegrityError(); + } + return envelope.result; +} + +export interface CandidateArtifactPage extends CandidatePage { + nativeTextSha256: string; + ocrTextSha256: string | null; + metrics: Record; +} + +export interface ComposedCandidateArtifact { + schemaVersion: "1"; + versionId: string; + candidateSha256: string; + documents: Array<{ documentId: string; text: string; textSha256: string; pages: CandidateArtifactPage[] }>; +} + +interface CandidateJobEvidence { + documentId: string; + remoteJobId: string | null; + requestedPages: number[]; + state: string; +} + +export async function persistComposedCandidateArtifact(input: { + rootDirectory: string; + versionId: string; + jobs: CandidateJobEvidence[]; +}): Promise<{ artifactPath: string; artifactSha256: string; candidate: ComposedCandidateArtifact }> { + assertUuid(input.versionId); + const versionDirectory = path.join(path.resolve(input.rootDirectory), input.versionId); + const manifest = await readPrivateJson(path.join(versionDirectory, "manifest.json"), undefined, "OCR manifest integrity validation failed"); + if (!isRecord(manifest) || manifest.schemaVersion !== "1" || manifest.versionId !== input.versionId || !Array.isArray(manifest.documents)) { + throw new Error("OCR manifest identity validation failed"); + } + const documents = []; + for (const entry of manifest.documents) { + if (!isManifestDocument(entry)) throw new Error("OCR manifest identity validation failed"); + const native = await readPrivateJson(resolveArtifactPath(versionDirectory, entry.nativePagesPath), entry.nativePagesSha256, "OCR native page artifact integrity validation failed"); + validateNativePages(native, input.versionId, entry.documentId, entry.originalSha256); + const requestedPages = native.requestedPages as number[]; + const job = input.jobs.find(({ documentId }) => documentId === entry.documentId); + let result: OcrResult | undefined; + if (requestedPages.length > 0) { + if (!job || job.state !== "succeeded" || !job.remoteJobId || !sameNumbers(job.requestedPages, requestedPages)) throw new Error("OCR result artifact identity validation failed"); + result = await readOcrResultArtifact({ rootDirectory: input.rootDirectory, versionId: input.versionId, documentId: entry.documentId, + jobId: job.remoteJobId, documentSha256: entry.originalSha256, pages: requestedPages }); + } else if (job) throw new Error("OCR result artifact identity validation failed"); + const resultPages = new Map(result?.pages.map((page) => [page.page, page])); + const evidence = (native.pages as Array<{ page: number; text: string; textSha256: string }>).map((page) => { + const ocr = resultPages.get(page.page); + if (!ocr) return { page: page.page, method: "native" as const, nativeText: page.text, rawOcrText: "", lines: [] }; + const classification = classifyOcrPage({ inkCoverage: ocr.metrics.inkCoverage, metrics: ocr.metrics }); + if (classification.method === "blocked") throw new Error(classification.errorCode); + return { page: page.page, method: classification.method, nativeText: page.text, rawOcrText: ocr.text, lines: ocr.lines }; + }); + const composed = composeCandidate(evidence); + documents.push({ documentId: entry.documentId, text: composed.text, textSha256: composed.textSha256, + pages: composed.pages.map((page) => { const ocr = resultPages.get(page.page); return { ...page, + nativeTextSha256: sha256Hex(page.nativeText), ocrTextSha256: ocr ? sha256Hex(ocr.text) : null, + metrics: ocr?.metrics ?? {} }; }) }); + } + const candidate: ComposedCandidateArtifact = { schemaVersion: "1", versionId: input.versionId, + candidateSha256: sha256Hex(canonicalJson(documents)), documents }; + const serialized = canonicalJson(candidate); + const artifactPath = path.join(versionDirectory, "candidate-pages.json"); + await publishImmutable(artifactPath, Buffer.from(serialized)); + const artifactSha256 = sha256Hex(serialized); + return { artifactPath, artifactSha256, candidate: await readComposedCandidateArtifact({ ...input, artifactSha256 }) }; +} + +export async function readComposedCandidateArtifact(input: { rootDirectory: string; versionId: string; artifactSha256?: string }): Promise { + assertUuid(input.versionId); + const value = await readPrivateJson(path.join(path.resolve(input.rootDirectory), input.versionId, "candidate-pages.json"), input.artifactSha256, "OCR candidate artifact integrity validation failed"); + if (!isRecord(value) || value.schemaVersion !== "1" || value.versionId !== input.versionId || !Array.isArray(value.documents) + || value.candidateSha256 !== sha256Hex(canonicalJson(value.documents))) throw new Error("OCR candidate artifact integrity validation failed"); + return value as unknown as ComposedCandidateArtifact; +} + +export interface ReviewedPagesArtifact { + schemaVersion: "1"; + versionId: string; + sourceId: string; + candidateSha256: string; + reviewedTextSha256: string; + reviewedBy: string; + documents: Array<{ documentId: string; pages: Array<{ page: number; candidateText: string; + ocr: { lines: Array<{ lineId: string; text: string; confidence: number; bbox: [number, number, number, number]; lineSha256: string }> } }> }>; +} + +interface ReviewedPagesIdentity { + rootDirectory: string; + versionId: string; + candidateSha256: string; + reviewedTextSha256: string; + artifactSha256?: string; +} + +export async function persistReviewedPagesArtifact(input: Omit & { rootDirectory: string }): Promise<{ + artifactPath: string; artifactSha256: string; created: boolean; reviewed: ReviewedPagesArtifact; +}> { + assertUuid(input.versionId); + const versionDirectory = path.join(path.resolve(input.rootDirectory), input.versionId); + const directory = await lstat(versionDirectory).catch(() => { throw new Error("OCR reviewed artifact directory not found"); }); + if (!directory.isDirectory()) throw new Error("OCR reviewed artifact directory not found"); + const reviewed: ReviewedPagesArtifact = { schemaVersion: "1", versionId: input.versionId, sourceId: input.sourceId, + candidateSha256: input.candidateSha256, reviewedTextSha256: input.reviewedTextSha256, + reviewedBy: input.reviewedBy, documents: input.documents }; + validateReviewedPages(reviewed, input); + const bytes = Buffer.from(canonicalJson(reviewed)); + const artifactPath = path.join(versionDirectory, "reviewed-pages.json"); + const created = await publishImmutable(artifactPath, bytes); + const artifactSha256 = sha256Hex(bytes); + try { + return { artifactPath, artifactSha256, created, reviewed: await readReviewedPagesArtifact({ ...input, artifactSha256 }) }; + } catch (error) { + if (created) await unlink(artifactPath).catch(() => undefined); + throw error; + } +} + +export async function readReviewedPagesArtifact(input: ReviewedPagesIdentity): Promise { + assertUuid(input.versionId); + const artifactPath = path.join(path.resolve(input.rootDirectory), input.versionId, "reviewed-pages.json"); + const value = await readPrivateJson(artifactPath, input.artifactSha256, "OCR reviewed artifact integrity validation failed"); + validateReviewedPages(value, input); + return value as ReviewedPagesArtifact; +} + +export async function removeReviewedPagesArtifact(rootDirectory: string, versionId: string): Promise { + assertUuid(versionId); + const target = path.join(path.resolve(rootDirectory), versionId, "reviewed-pages.json"); + await unlink(target); + await syncDirectory(path.dirname(target)); +} + +export async function readOcrArtifactPageNumbers(input: { rootDirectory: string; versionId: string; documentId: string }): Promise { + const { entry, native } = await readNativeDocument(input); + validateNativePages(native, input.versionId, input.documentId, entry.originalSha256); + return (native.pages as Array<{ page: number }>).map(({ page }) => page); +} + +export async function persistReviewImageArtifacts(input: { + rootDirectory: string; versionId: string; documentId: string; + images: Array<{ page: number; bytes: Buffer; sha256: string }>; +}): Promise<{ images: Array<{ page: number; artifactPath: string; sha256: string }> }> { + assertUuid(input.versionId); + const { versionDirectory, entry, native } = await readNativeDocument(input); + validateNativePages(native, input.versionId, input.documentId, entry.originalSha256); + const expectedPages = (native.pages as Array<{ page: number }>).map(({ page }) => page); + if (!sameNumbers(input.images.map(({ page }) => page), expectedPages)) throw new Error("OCR review image identity validation failed"); + const directory = resolveArtifactPath(versionDirectory, path.posix.join("documents", entry.documentArtifactId, "review-images")); + await mkdir(directory, { recursive: true, mode: 0o700 }); + const images = []; + for (const image of input.images) { + if (sha256Hex(image.bytes) !== image.sha256 || !image.bytes.subarray(0, 8).equals(Buffer.from("89504e470d0a1a0a", "hex"))) throw new Error("OCR review image integrity validation failed"); + const relativePath = path.posix.join("documents", entry.documentArtifactId, "review-images", `page-${String(image.page).padStart(4, "0")}.png`); + const artifactPath = resolveArtifactPath(versionDirectory, relativePath); + await publishImmutable(artifactPath, image.bytes); + images.push({ page: image.page, relativePath, artifactPath, sha256: image.sha256, mimeType: "image/png" }); + } + const manifestPath = path.join(directory, "manifest.json"); + await publishImmutable(manifestPath, Buffer.from(canonicalJson({ schemaVersion: "1", versionId: input.versionId, documentId: input.documentId, + documentSha256: entry.originalSha256, images: images.map(({ artifactPath: _artifactPath, ...image }) => image) }))); + return { images: images.map(({ page, artifactPath, sha256 }) => ({ page, artifactPath, sha256 })) }; +} + +export async function readReviewImageArtifact(input: { + rootDirectory: string; versionId: string; documentId: string; page: number; +}): Promise<{ bytes: Buffer; sha256: string; mimeType: "image/png" }> { + assertUuid(input.versionId); + const { versionDirectory, entry } = await readNativeDocument(input); + const directory = resolveArtifactPath(versionDirectory, path.posix.join("documents", entry.documentArtifactId, "review-images")); + const manifest = await readPrivateJson(path.join(directory, "manifest.json"), undefined, "OCR review image integrity validation failed"); + if (!isRecord(manifest) || manifest.schemaVersion !== "1" || manifest.versionId !== input.versionId || manifest.documentId !== input.documentId + || manifest.documentSha256 !== entry.originalSha256 || !Array.isArray(manifest.images)) throw new Error("OCR review image identity validation failed"); + const image = manifest.images.find((value) => isRecord(value) && value.page === input.page); + if (!isRecord(image) || typeof image.relativePath !== "string" || typeof image.sha256 !== "string" || image.mimeType !== "image/png") throw new Error("OCR review image not found"); + const artifactPath = resolveArtifactPath(versionDirectory, image.relativePath); + const file = await lstat(artifactPath).catch(() => { throw new Error("OCR review image integrity validation failed"); }); + const bytes = await readFile(artifactPath); + if (!file.isFile() || (file.mode & 0o777) !== 0o600 || sha256Hex(bytes) !== image.sha256) throw new Error("OCR review image integrity validation failed"); + return { bytes, sha256: image.sha256, mimeType: "image/png" }; +} + export async function sweepOrphanArtifacts(input: { rootDirectory: string; retainedVersionIds: ReadonlySet; @@ -111,6 +378,90 @@ async function durableWrite(target: string, bytes: Buffer): Promise { await syncDirectory(path.dirname(target)); } +async function publishImmutable(target: string, bytes: Buffer): Promise { + try { + await durableWrite(target, bytes); + return true; + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; + const existing = await readFile(target); + if (!existing.equals(bytes)) throw new Error("OCR result artifact conflicts with durable content"); + return false; + } +} + +function resultArtifactPath(rootDirectory: string, versionId: string, documentId: string): string { + const versionDirectory = path.join(path.resolve(rootDirectory), versionId); + return resolveArtifactPath(versionDirectory, path.posix.join("documents", uuidV5(documentId), "ocr-result.json")); +} + +function isRecord(value: unknown): value is Record { + return typeof value === "object" && value !== null && !Array.isArray(value); +} + +function isManifestDocument(value: unknown): value is { documentId: string; originalSha256: string; nativePagesPath: string; nativePagesSha256: string } { + return isRecord(value) && typeof value.documentId === "string" && typeof value.originalSha256 === "string" + && typeof value.nativePagesPath === "string" && typeof value.nativePagesSha256 === "string"; +} + +async function readNativeDocument(input: { rootDirectory: string; versionId: string; documentId: string }) { + assertUuid(input.versionId); + const versionDirectory = path.join(path.resolve(input.rootDirectory), input.versionId); + const manifest = await readPrivateJson(path.join(versionDirectory, "manifest.json"), undefined, "OCR manifest integrity validation failed"); + if (!isRecord(manifest) || manifest.schemaVersion !== "1" || manifest.versionId !== input.versionId || !Array.isArray(manifest.documents)) throw new Error("OCR manifest identity validation failed"); + const entry = manifest.documents.find((value) => isManifestDocument(value) && value.documentId === input.documentId); + if (!entry || !isRecord(entry) || typeof entry.documentArtifactId !== "string") throw new Error("OCR review image not found"); + const document = entry as { documentId: string; documentArtifactId: string; originalSha256: string; nativePagesPath: string; nativePagesSha256: string }; + const native = await readPrivateJson(resolveArtifactPath(versionDirectory, document.nativePagesPath), document.nativePagesSha256, "OCR native page artifact integrity validation failed"); + return { versionDirectory, entry: document, native }; +} + +function validateNativePages(value: unknown, versionId?: string, documentId?: string, documentSha256?: string): asserts value is Record { + if (!isRecord(value) || value.schemaVersion !== "1" || (versionId && value.versionId !== versionId) + || (documentId && value.documentId !== documentId) || (documentSha256 && value.documentSha256 !== documentSha256) + || !Array.isArray(value.pages) || !Array.isArray(value.requestedPages)) throw new Error("OCR native page artifact integrity validation failed"); + const pages = value.pages as Array>; + const requested = value.requestedPages as number[]; + if (pages.some((page, index) => page.page !== index + 1 || typeof page.text !== "string" || page.textSha256 !== sha256Hex(page.text) + || typeof page.rasterCoverage !== "number" || page.rasterCoverage < 0 || page.rasterCoverage > 1) + || !sameNumbers(requested, [...new Set(requested)].sort((a, b) => a - b)) || requested.some((page) => page < 1 || page > pages.length)) { + throw new Error("OCR native page artifact integrity validation failed"); + } +} + +function validateReviewedPages(value: unknown, identity: Pick): asserts value is ReviewedPagesArtifact { + if (!isRecord(value) || value.schemaVersion !== "1" || value.versionId !== identity.versionId || value.candidateSha256 !== identity.candidateSha256 + || value.reviewedTextSha256 !== identity.reviewedTextSha256 || typeof value.sourceId !== "string" || typeof value.reviewedBy !== "string" + || !Array.isArray(value.documents)) throw new Error("OCR reviewed artifact identity validation failed"); + const documents = value.documents as ReviewedPagesArtifact["documents"]; + const pages = documents.flatMap((document) => Array.isArray(document.pages) ? document.pages : []); + if (documents.some((document) => typeof document.documentId !== "string" || !Array.isArray(document.pages)) + || pages.some((page) => !Number.isInteger(page.page) || page.page < 1 || typeof page.candidateText !== "string" || !isRecord(page.ocr) + || !Array.isArray(page.ocr.lines) || page.ocr.lines.some((line) => !isRecord(line) || typeof line.lineId !== "string" + || typeof line.text !== "string" || line.lineSha256 !== sha256Hex(line.text)) + || (page.ocr.lines.length > 0 && page.candidateText !== page.ocr.lines.map(({ text }) => text).join("\n")))) { + throw new Error("OCR reviewed artifact integrity validation failed"); + } + const reviewedText = pages.filter(({ candidateText }) => candidateText).map(({ candidateText }) => candidateText).join("\n\n"); + if (sha256Hex(reviewedText) !== identity.reviewedTextSha256) throw new Error("OCR reviewed artifact integrity validation failed"); +} + +async function readPrivateJson(target: string, expectedSha256: string | undefined, message: string): Promise { + const file = await lstat(target).catch(() => { throw new Error(message); }); + if (!file.isFile() || (file.mode & 0o777) !== 0o600) throw new Error(message); + const bytes = await readFile(target); + if (expectedSha256 && sha256Hex(bytes) !== expectedSha256) throw new Error(message); + try { return JSON.parse(bytes.toString("utf8")); } catch { throw new Error(message); } +} + +function sameNumbers(left: number[], right: number[]): boolean { + return left.length === right.length && left.every((value, index) => value === right[index]); +} + +function resultIntegrityError(): Error { + return new Error("OCR result artifact integrity validation failed"); +} + async function syncDirectory(directory: string): Promise { const handle = await open(directory, "r"); try { await handle.sync(); } finally { await handle.close(); } diff --git a/src/modules/ocr/client.ts b/src/modules/ocr/client.ts index cd6b66b..3a6d1ca 100644 --- a/src/modules/ocr/client.ts +++ b/src/modules/ocr/client.ts @@ -40,7 +40,7 @@ export interface OcrResult { height: number; processingMs: number; text: string; - metrics: { lineCount: number; nonWhitespaceCharacters: number; medianConfidence: number; p10Confidence: number; lowConfidenceLineRatio: number }; + metrics: { lineCount: number; nonWhitespaceCharacters: number; inkCoverage: number; medianConfidence: number; p10Confidence: number; lowConfidenceLineRatio: number }; lines: Array<{ lineId: string; text: string; confidence: number; bbox: [number, number, number, number] }>; }>; } @@ -115,10 +115,24 @@ export class OcrClient { async getResult(jobId: string, expected: { documentSha256: string; pages: number[] }): Promise { assertExpected(expected); const value = await this.requestJson(`/v1/jobs/${encodeURIComponent(jobId)}/result`, () => ({ headers: this.headers() })); - if (!validResult(value, jobId, expected)) throw integrityError(); + if (!isValidOcrResult(value, jobId, expected)) throw integrityError(); return value; } + async getReviewImage(jobId: string, page: number, documentSha256: string): Promise<{ bytes: Buffer; sha256: string }> { + if (!Number.isInteger(page) || page < 1 || !/^[a-f0-9]{64}$/u.test(documentSha256)) throw new TypeError("OCR review image identity is invalid"); + const response = await this.requestFetch(`${this.baseUrl}/v1/jobs/${encodeURIComponent(jobId)}/pages/${page}/image`, { headers: this.headers() }); + if (!response.ok) throw new OcrClientError(`OCR_HTTP_${response.status}`, response.status, false); + const bytes = Buffer.from(await response.arrayBuffer()); + const sha256 = sha256Hex(bytes); + if (response.headers.get("content-type")?.split(";")[0] !== "image/png" + || response.headers.get("x-document-sha256") !== documentSha256 + || response.headers.get("x-content-sha256") !== sha256 + || response.headers.get("x-page-number") !== String(page) + || !bytes.subarray(0, 8).equals(Buffer.from("89504e470d0a1a0a", "hex"))) throw integrityError(); + return { bytes, sha256 }; + } + async delete(jobId: string): Promise { await this.requestJson(`/v1/jobs/${encodeURIComponent(jobId)}`, () => ({ method: "DELETE", headers: this.headers() })); } @@ -155,7 +169,7 @@ function assertExpected(expected: { documentSha256: string; pages: number[] }): } } -function validResult(value: unknown, jobId: string, expected: { documentSha256: string; pages: number[] }): value is OcrResult { +export function isValidOcrResult(value: unknown, jobId: string, expected: { documentSha256: string; pages: number[] }): value is OcrResult { if (!isObject(value) || value.schemaVersion !== "1" || value.jobId !== jobId || value.documentSha256 !== expected.documentSha256 || !isObject(value.engine) || canonicalEngine(value.engine) !== "paddleocr|3.4.0|paddlepaddle-3.2.2|cpu|ocr-v1|200" || !Array.isArray(value.pages) || !sameNumbers(value.pages.map((page) => isObject(page) ? page.page : undefined), expected.pages)) return false; @@ -164,7 +178,7 @@ function validResult(value: unknown, jobId: string, expected: { documentSha256: && page.lines.every((line) => isObject(line) && typeof line.lineId === "string" && line.lineId.length > 0 && typeof line.text === "string" && validRatio(line.confidence) && Array.isArray(line.bbox) && line.bbox.length === 4 && line.bbox.every(Number.isFinite)) && page.text === page.lines.map((line) => (line as Record).text).join("\n") - && page.metrics.lineCount === page.lines.length && isCount(page.metrics.nonWhitespaceCharacters) + && page.metrics.lineCount === page.lines.length && isCount(page.metrics.nonWhitespaceCharacters) && validRatio(page.metrics.inkCoverage) && validRatio(page.metrics.medianConfidence) && validRatio(page.metrics.p10Confidence) && validRatio(page.metrics.lowConfidenceLineRatio)); } diff --git a/src/modules/ocr/dispatcher.ts b/src/modules/ocr/dispatcher.ts index 1643998..6f168c7 100644 --- a/src/modules/ocr/dispatcher.ts +++ b/src/modules/ocr/dispatcher.ts @@ -11,10 +11,13 @@ export interface OcrDispatchStore { failOcrJob(jobId: string, code: string, detail: string): Promise; markReviewRequired(versionId: string): Promise; markFailed(versionId: string, code: string, detail: string): Promise; + listOcrVersionsAwaitingCandidate(): Promise; } export type OcrDispatchInput = { bytes: Buffer; documentSha256: string }; export type OcrDispatchResult = "idle" | "pending" | "succeeded" | "failed"; +type PersistOcrResult = (job: OcrJobRow, result: OcrResult) => Promise; +type FinalizeCandidate = (versionId: string) => Promise; export class OcrDispatcher { private activeDrain: Promise | undefined; @@ -23,6 +26,8 @@ export class OcrDispatcher { private readonly store: OcrDispatchStore, private readonly client: OcrClient, private readonly loadInput: (job: OcrJobRow) => Promise, + private readonly persistResult: PersistOcrResult, + private readonly finalizeCandidate: FinalizeCandidate, private readonly leaseMs = 30_000 ) {} @@ -62,8 +67,12 @@ export class OcrDispatcher { documentSha256: input.documentSha256, pages: job.requestedPages }); - const versionComplete = await this.store.completeOcrJob(job.jobId, result); - if (versionComplete) await this.store.markReviewRequired(job.versionId); + const durableResult = await this.persistResult(job, result); + const versionComplete = await this.store.completeOcrJob(job.jobId, durableResult); + if (versionComplete) { + await this.finalizeCandidate(job.versionId); + await this.store.markReviewRequired(job.versionId); + } await this.client.delete(remoteJobId).catch(() => undefined); return "succeeded"; } catch (error) { @@ -72,7 +81,7 @@ export class OcrDispatcher { await this.store.requeueOcrJob(job.jobId, error.code, detail, 15_000); return "pending"; } - const code = error instanceof OcrClientError ? error.code : "OCR_DISPATCH_FAILED"; + const code = error instanceof OcrClientError ? error.code : error instanceof Error && error.message === "OCR_QUALITY_BLOCKED" ? error.message : "OCR_DISPATCH_FAILED"; await this.fail(job, code, detail); return "failed"; } @@ -87,6 +96,20 @@ export class OcrDispatcher { return recovered.length; } + async recoverCompletedCandidates(): Promise { + const versionIds = await this.store.listOcrVersionsAwaitingCandidate(); + for (const versionId of versionIds) { + try { + await this.finalizeCandidate(versionId); + await this.store.markReviewRequired(versionId); + } catch (error) { + const detail = error instanceof Error ? error.message : "Unknown OCR candidate recovery failure"; + await this.store.markFailed(versionId, "OCR_CANDIDATE_FAILED", detail); + } + } + return versionIds.length; + } + dispatchAvailable(): Promise { if (this.activeDrain) return this.activeDrain; const drain = this.drainAvailable(); diff --git a/src/modules/ocr/indexing.ts b/src/modules/ocr/indexing.ts index 234c36c..4f1ea0b 100644 --- a/src/modules/ocr/indexing.ts +++ b/src/modules/ocr/indexing.ts @@ -1,4 +1,11 @@ import { CatalogError } from "../catalog/errors.js"; +import { withTransaction, type PgPool } from "../catalog/client.js"; +import type { EmbeddingProvider } from "../embeddings/provider.js"; +import type { VectorStoreClient } from "../vectorstore/client.js"; +import type { IngestedChunk } from "../../shared/types/rag.js"; +import { buildChunkId, buildVersionedQdrantPointId, normalizeContentForHash, sha256Hex } from "../../shared/utils/ids.js"; +import { chunkDocument, documentalChunkingPolicy } from "../process/chunking.js"; +import { readComposedCandidateArtifact, readReviewedPagesArtifact, type ReviewedPagesArtifact } from "./artifacts.js"; export interface ApprovedOcrCandidate { versionId: string; @@ -65,3 +72,130 @@ export class OcrIndexingService { } } } + +interface IndexingLifecycleRow { + version_id: string; + source_id: string; + state: string; + base_active_version_id: string | null; + current_active_version_id: string | null; + processing_fingerprint: string; + metadata_hash: string; + tags: string[]; + embedding_provider: string; + embedding_model: string; + embedding_dimensions: number; + qdrant_collection: string; + expected_document_count: number; + document_id: string; + document_key: string; + title: string; + mime_type: string; + page_number: number; + reviewed_text_hash: string; +} + +function versionedPoints(versionId: string, documentId: string, reviewedText: string, fingerprint: string, embedding: number[], index: number, base: Omit): IngestedChunk { + return { + id: buildVersionedQdrantPointId(versionId, buildChunkId(documentId, "documental", index)), + vector: embedding, + ...base, + payload: { ...base.payload, chunk_id: buildChunkId(documentId, "documental", index), source_version_id: versionId, processing_fingerprint: fingerprint } + }; +} + +export class PostgresOcrIndexingStore implements OcrIndexingStore { + constructor( + private readonly pool: PgPool, + private readonly rootDirectory: string, + private readonly embeddingProvider: EmbeddingProvider, + private readonly vectorStore: VectorStoreClient + ) {} + + async findReusableVersion(): Promise<{ versionId: string } | undefined> { + return undefined; + } + + async indexReviewed(candidate: ApprovedOcrCandidate): Promise { + const composed = await readComposedCandidateArtifact({ rootDirectory: this.rootDirectory, versionId: candidate.versionId }); + const reviewed: ReviewedPagesArtifact = await readReviewedPagesArtifact({ + rootDirectory: this.rootDirectory, versionId: candidate.versionId, + candidateSha256: composed.candidateSha256, reviewedTextSha256: candidate.reviewedTextSha256 + }); + if (reviewed.sourceId !== candidate.sourceId || sha256Hex(candidate.reviewedText) !== candidate.reviewedTextSha256) { + throw new CatalogError("Reviewed artifact identity validation failed", 409, "REVIEW_CONFLICT"); + } + const rows = await this.loadLifecycleRows(candidate); + if (rows.length === 0 || rows.some((row) => row.version_id !== candidate.versionId || row.processing_fingerprint !== candidate.processingFingerprint + || row.metadata_hash !== candidate.metadataHash)) throw new CatalogError("OCR indexing lifecycle identity validation failed", 409, "REVIEW_CONFLICT"); + + const normalized = normalizeContentForHash(candidate.reviewedText); + const chunks = chunkDocument(rows[0]!.title, normalized, documentalChunkingPolicy); + const embeddings = await this.embeddingProvider.embed(chunks.map((chunk) => chunk.content)); + if (embeddings.length !== chunks.length || embeddings.some((embedding) => embedding.length !== rows[0]!.embedding_dimensions)) { + throw new CatalogError("Embedding provider returned invalid dimensions", 503, "EMBEDDING_DIMENSIONS_INVALID"); + } + const points = chunks.map((chunk, index) => versionedPoints(candidate.versionId, reviewed.documents[0]!.documentId, candidate.reviewedText, + candidate.processingFingerprint, embeddings[index]!, index, { + payload: { + chunk_id: "", source_id: candidate.sourceId, source_version_id: candidate.versionId, source_version_number: 0, + source_type: "file", source_ref: reviewed.documents[0]!.documentId, document_id: reviewed.documents[0]!.documentId, + document_key: rows[0]!.document_key, document_content_hash: sha256Hex(normalized), title: rows[0]!.title, + mime_type: rows[0]!.mime_type, section_title: chunk.sectionTitle, chunk_mode: documentalChunkingPolicy.mode, + chunk_index: chunk.index, start_line: chunk.startLine, end_line: chunk.endLine, content: chunk.content, + tags: rows[0]!.tags, embedding_provider: rows[0]!.embedding_provider, embedding_model: rows[0]!.embedding_model, + embedding_dimensions: rows[0]!.embedding_dimensions, processing_fingerprint: candidate.processingFingerprint, + indexed_at: new Date().toISOString(), write_state: "staged" + } + })); + await this.vectorStore.upsert(points); + const stored = await this.vectorStore.countVersionPoints(candidate.versionId); + if (stored !== points.length) throw new CatalogError("Indexed point count does not match the expected catalog count", 503, "SOURCE_VERSION_INCONSISTENT"); + return points.length; + } + + async markReady(versionId: string, verifiedPointCount: number): Promise { + await withTransaction(this.pool, async (client) => { + const documents = await client.query("UPDATE rag_version_documents SET indexed_chunk_count = $2, indexed_at = now() WHERE version_id = $1", + [versionId, verifiedPointCount]); + if (documents.rowCount !== 1) throw new CatalogError("OCR indexing target documents not found", 409, "REVIEW_CONFLICT"); + const version = await client.query( + "UPDATE rag_source_versions SET state = 'ready', indexed_at = now(), indexed_chunk_count = $2 WHERE version_id = $1 AND state = 'indexing'", + [versionId, verifiedPointCount]); + if (version.rowCount !== 1) throw new CatalogError("Version is not in OCR indexing state", 409, "INVALID_VERSION_STATE"); + }); + } + + async settleReusable(): Promise { + throw new CatalogError("Reusable version settlement is not available for OCR indexing", 409, "REVIEW_CONFLICT"); + } + + async activateVersion(): Promise { + throw new CatalogError("OCR activation requires a separately authorized activation unit", 409, "INVALID_VERSION_STATE"); + } + + private async loadLifecycleRows(candidate: ApprovedOcrCandidate): Promise { + const result = await this.pool.query( + `SELECT v.version_id, v.source_id, v.state, v.base_active_version_id, s.active_version_id AS current_active_version_id, + v.processing_fingerprint, v.metadata_hash, v.tags, e.embedding_provider, e.embedding_model, e.embedding_dimensions, + e.qdrant_collection, v.expected_document_count, d.document_id, d.document_key, d.title, d.mime_type, p.page_number, p.reviewed_text_hash + FROM rag_source_versions v + JOIN rag_sources s ON s.source_id = v.source_id + LEFT JOIN rag_version_documents d ON d.version_id = v.version_id + LEFT JOIN rag_document_pages p ON p.version_id = v.version_id + LEFT JOIN LATERAL (SELECT 'pending' AS embedding_provider, 'pending' AS embedding_model, 0 AS embedding_dimensions, 'rag_documents' AS qdrant_collection) e ON true + WHERE v.version_id = $1 AND v.state = 'indexing'`, + [candidate.versionId]); + return result.rows as IndexingLifecycleRow[]; + } +} + +export class OcrReadyIndexingService { + constructor(private readonly store: OcrIndexingStore) {} + + async index(candidate: ApprovedOcrCandidate): Promise { + const verifiedPointCount = await this.store.indexReviewed(candidate); + await this.store.markReady(candidate.versionId, verifiedPointCount); + return { versionId: candidate.versionId, state: "ready", activated: false }; + } +} diff --git a/src/modules/ocr/review.ts b/src/modules/ocr/review.ts index 05a5c9f..18c4ddd 100644 --- a/src/modules/ocr/review.ts +++ b/src/modules/ocr/review.ts @@ -1,5 +1,7 @@ import { CatalogError } from "../catalog/errors.js"; -import { sha256Hex } from "../../shared/utils/ids.js"; +import { withTransaction, type PgPool, type PgPoolClient } from "../catalog/client.js"; +import { canonicalJson, sha256Hex } from "../../shared/utils/ids.js"; +import { persistReviewedPagesArtifact, readComposedCandidateArtifact, readReviewImageArtifact, removeReviewedPagesArtifact } from "./artifacts.js"; export interface OcrReviewLine { lineId: string; @@ -57,6 +59,146 @@ export interface OcrReviewStore { }): Promise; } +export interface OcrReviewContext { + versionId: string; sourceId: string; state: string; baseActiveVersionId: string | null; currentActiveVersionId: string | null; + activateRequested: boolean; processingFingerprint: string; metadataHash: string; + pages: Array<{ documentId: string; page: number; nativeTextSha256: string; ocrTextSha256: string | null; + candidateTextSha256: string; metrics: Record; risks: string[] }>; +} + +type ReviewReader = Pick; + +export class PostgresOcrReviewStore implements OcrReviewStore { + constructor(private readonly pool: PgPool, private readonly reader: ReviewReader, private readonly rootDirectory: string) {} + + loadCandidate(versionId: string): Promise { + return this.reader.view(versionId); + } + + async commitApproval(input: Parameters[0]): Promise { + let createdArtifact = false; + try { + await withTransaction(this.pool, async (client) => { + const current = await this.lockCurrentCandidate(client, input.candidate.versionId); + if (input.candidateSha256 !== current.candidateSha256 || input.expectedActiveVersionId !== current.baseActiveVersionId + || input.expectedActiveVersionId !== current.currentActiveVersionId) throw conflict("OCR review decision is stale"); + const corrected = applyCorrections(current, input.corrections); + const text = reviewedText(corrected); + if (canonicalJson(corrected.documents) !== canonicalJson(input.candidate.documents) + || text !== input.reviewedText || sha256Hex(text) !== input.reviewedTextSha256) throw conflict("OCR review decision does not match the candidate"); + const artifact = await persistReviewedPagesArtifact({ rootDirectory: this.rootDirectory, versionId: current.versionId, + sourceId: current.sourceId, candidateSha256: current.candidateSha256, reviewedTextSha256: input.reviewedTextSha256, + reviewedBy: input.reviewedBy, documents: corrected.documents }); + createdArtifact = artifact.created; + for (const correction of input.corrections) await client.query( + `INSERT INTO rag_review_corrections(version_id, document_id, page_number, line_id, expected_line_hash, replacement_text, reviewed_by) + VALUES ($1, $2, $3, $4, $5, $6, $7)`, + [current.versionId, correction.documentId, correction.page, correction.lineId, correction.expectedLineSha256, correction.replacementText, input.reviewedBy] + ); + for (const document of corrected.documents) for (const page of document.pages) { + const result = await client.query( + `UPDATE rag_document_pages SET reviewed_text_hash = $4 + WHERE version_id = $1 AND document_id = $2 AND page_number = $3 AND candidate_text_hash = $5`, + [current.versionId, document.documentId, page.page, sha256Hex(page.candidateText), sha256Hex(current.documents + .find(({ documentId }) => documentId === document.documentId)!.pages.find(({ page: pageNumber }) => pageNumber === page.page)!.candidateText)] + ); + if (result.rowCount !== 1) throw conflict("OCR review page evidence changed"); + } + const transitioned = await client.query( + `UPDATE rag_source_versions v SET state = 'indexing', reviewed_at = now(), reviewed_by = $2, indexing_started_at = now() + WHERE v.version_id = $1 AND v.state = 'review_required' AND v.base_active_version_id IS NOT DISTINCT FROM $3 + AND EXISTS (SELECT 1 FROM rag_sources s WHERE s.source_id = v.source_id AND s.active_version_id IS NOT DISTINCT FROM $3)`, + [current.versionId, input.reviewedBy, input.expectedActiveVersionId] + ); + if (transitioned.rowCount !== 1) throw conflict("Active version changed during review", "ACTIVE_VERSION_CHANGED"); + }); + } catch (error) { + if (createdArtifact) await removeReviewedPagesArtifact(this.rootDirectory, input.candidate.versionId).catch(() => undefined); + throw error; + } + } + + async commitRejection(input: Parameters[0]): Promise { + await withTransaction(this.pool, async (client) => { + const current = await this.lockCurrentCandidate(client, input.candidate.versionId); + if (input.candidateSha256 !== current.candidateSha256 || canonicalJson(input.candidate) !== canonicalJson(current)) { + throw conflict("OCR review decision is stale"); + } + const result = await client.query( + `UPDATE rag_source_versions SET state = 'rejected', reviewed_at = now(), reviewed_by = $2, + error_code = 'OCR_REJECTED', error_detail = $3, retention_due_at = now() + interval '7 days' + WHERE version_id = $1 AND state = 'review_required'`, + [current.versionId, input.reviewedBy, input.reason.slice(0, 2000)] + ); + if (result.rowCount !== 1) throw conflict("Version is not awaiting OCR review", "INVALID_VERSION_STATE"); + }); + } + + private async lockCurrentCandidate(client: PgPoolClient, versionId: string): Promise { + const version = await client.query<{ + version_id: string; source_id: string; state: string; base_active_version_id: string | null; current_active_version_id: string | null; + activate_requested: boolean; processing_fingerprint: string; metadata_hash: string; + }>(`SELECT v.version_id, v.source_id, v.state, v.base_active_version_id, s.active_version_id AS current_active_version_id, + v.activate_requested, v.processing_fingerprint, v.metadata_hash + FROM rag_source_versions v JOIN rag_sources s ON s.source_id = v.source_id + WHERE v.version_id = $1 FOR UPDATE OF v, s`, [versionId]); + const row = version.rows[0]; + if (!row) throw new CatalogError("OCR review candidate not found", 404, "REVIEW_NOT_FOUND"); + if (row.state !== "review_required") throw conflict("Version is not awaiting OCR review", "INVALID_VERSION_STATE"); + const pages = await client.query<{ document_id: string; page_number: number | string; native_text_hash: string; + ocr_text_hash: string | null; candidate_text_hash: string; risk_tokens: string[] }>( + `SELECT document_id, page_number, native_text_hash, ocr_text_hash, candidate_text_hash, risk_tokens + FROM rag_document_pages WHERE version_id = $1 ORDER BY document_id, page_number FOR UPDATE`, [versionId]); + const candidate = await this.reader.view(versionId); + if (candidate.versionId !== row.version_id || candidate.sourceId !== row.source_id || candidate.state !== row.state + || candidate.baseActiveVersionId !== row.base_active_version_id || candidate.currentActiveVersionId !== row.current_active_version_id + || candidate.activateRequested !== row.activate_requested || candidate.processingFingerprint !== row.processing_fingerprint + || candidate.metadataHash !== row.metadata_hash) throw conflict("OCR review lifecycle identity changed"); + const candidatePages = candidate.documents.flatMap((document) => document.pages.map((page) => ({ documentId: document.documentId, page }))); + if (candidatePages.length !== pages.rows.length || candidatePages.some(({ documentId, page }) => { + const locked = pages.rows.find((value) => value.document_id === documentId && Number(value.page_number) === page.page); + return !locked || locked.native_text_hash !== sha256Hex(page.nativeText) + || locked.ocr_text_hash !== (page.ocr.text ? sha256Hex(page.ocr.text) : null) + || locked.candidate_text_hash !== sha256Hex(page.candidateText) || canonicalJson(locked.risk_tokens) !== canonicalJson(page.risks); + })) throw conflict("OCR review page evidence changed"); + return candidate; + } +} + +export class DurableOcrReviewReader { + constructor(private readonly catalog: { loadOcrReviewContext(versionId: string): Promise }, private readonly rootDirectory: string) {} + + async view(versionId: string): Promise { + const context = await this.catalog.loadOcrReviewContext(versionId); + if (!context) throw new CatalogError("OCR review candidate not found", 404, "REVIEW_NOT_FOUND"); + if (context.versionId !== versionId) throw new Error("OCR review candidate identity validation failed"); + if (context.state !== "review_required") throw conflict("Version is not awaiting OCR review", "INVALID_VERSION_STATE"); + const candidate = await readComposedCandidateArtifact({ rootDirectory: this.rootDirectory, versionId }); + const pages = candidate.documents.flatMap((document) => document.pages.map((page) => ({ documentId: document.documentId, page }))); + if (pages.length !== context.pages.length) throw new Error("OCR review candidate lifecycle validation failed"); + for (const item of pages) { + const persisted = context.pages.find(({ documentId, page }) => documentId === item.documentId && page === item.page.page); + if (!persisted || persisted.nativeTextSha256 !== item.page.nativeTextSha256 || persisted.ocrTextSha256 !== item.page.ocrTextSha256 + || persisted.candidateTextSha256 !== item.page.candidateTextSha256 || canonicalJson(persisted.metrics) !== canonicalJson(item.page.metrics) + || canonicalJson(persisted.risks) !== canonicalJson(item.page.risks)) throw new Error("OCR review candidate lifecycle validation failed"); + await readReviewImageArtifact({ rootDirectory: this.rootDirectory, versionId, documentId: item.documentId, page: item.page.page }); + } + const { pages: _lifecyclePages, ...identity } = context; + return { ...identity, state: "review_required", candidateSha256: candidate.candidateSha256, + documents: candidate.documents.map((document) => ({ documentId: document.documentId, pages: document.pages.map((page) => ({ + page: page.page, imageUrl: `/ingestions/${versionId}/documents/${encodeURIComponent(document.documentId)}/pages/${page.page}/image`, + nativeText: page.nativeText, ocr: { text: page.rawOcrText, lines: page.lines }, candidateText: page.candidateText, + differences: page.nativeText === page.rawOcrText ? [] : ["Native and OCR text differ"], risks: page.risks + })) })) }; + } + + async image(versionId: string, documentId: string, page: number) { + const candidate = await this.view(versionId); + if (!candidate.documents.some((document) => document.documentId === documentId && document.pages.some((item) => item.page === page))) throw new CatalogError("OCR review image not found", 404, "REVIEW_IMAGE_NOT_FOUND"); + return readReviewImageArtifact({ rootDirectory: this.rootDirectory, versionId, documentId, page }); + } +} + export interface ApprovedOcrReview { versionId: string; sourceId: string; @@ -95,33 +237,13 @@ export class OcrReviewService { throw conflict("Active version changed during review", "ACTIVE_VERSION_CHANGED"); } - const reviewedCandidate = structuredClone(candidate); - const lines = new Map(); - for (const document of reviewedCandidate.documents) for (const page of document.pages) for (const line of page.ocr.lines) { - const key = `${document.documentId}:${page.page}:${line.lineId}`; - if (lines.has(key)) throw conflict("Candidate line identities are duplicated", "CORRECTION_CONFLICT"); - lines.set(key, line); - } - const targets = new Set(); - for (const correction of input.corrections) { - const key = `${correction.documentId}:${correction.page}:${correction.lineId}`; - const line = lines.get(key); - if (targets.has(key) || !line || line.lineSha256 !== correction.expectedLineSha256) { - throw conflict("Correction target is stale or duplicated", "CORRECTION_CONFLICT"); - } - targets.add(key); - line.text = correction.replacementText; - line.lineSha256 = sha256Hex(line.text); - } - for (const document of reviewedCandidate.documents) for (const page of document.pages) { - if (page.ocr.lines.length > 0) page.ocr.text = page.candidateText = page.ocr.lines.map(({ text }) => text).join("\n"); - } - const reviewedText = reviewedCandidate.documents.flatMap(({ pages }) => pages.filter(({ candidateText }) => candidateText).map(({ candidateText }) => candidateText)).join("\n\n"); - const reviewedTextSha256 = sha256Hex(reviewedText); - await this.store.commitApproval({ candidate: reviewedCandidate, ...input, reviewedText, reviewedTextSha256 }); + const reviewedCandidate = applyCorrections(candidate, input.corrections); + const reviewedTextValue = reviewedText(reviewedCandidate); + const reviewedTextSha256 = sha256Hex(reviewedTextValue); + await this.store.commitApproval({ candidate: reviewedCandidate, ...input, reviewedText: reviewedTextValue, reviewedTextSha256 }); return { versionId, sourceId: candidate.sourceId, state: "indexing", activateRequested: candidate.activateRequested, - expectedActiveVersionId: input.expectedActiveVersionId, reviewedText, reviewedTextSha256, + expectedActiveVersionId: input.expectedActiveVersionId, reviewedText: reviewedTextValue, reviewedTextSha256, processingFingerprint: candidate.processingFingerprint, metadataHash: candidate.metadataHash }; } @@ -141,3 +263,30 @@ export class OcrReviewService { return { versionId, state: "rejected", activated: false }; } } + +function applyCorrections(candidate: OcrReviewCandidate, corrections: OcrCorrection[]): OcrReviewCandidate { + const corrected = structuredClone(candidate); + const lines = new Map(); + for (const document of corrected.documents) for (const page of document.pages) for (const line of page.ocr.lines) { + const key = `${document.documentId}:${page.page}:${line.lineId}`; + if (lines.has(key)) throw conflict("Candidate line identities are duplicated", "CORRECTION_CONFLICT"); + lines.set(key, line); + } + const targets = new Set(); + for (const correction of corrections) { + const key = `${correction.documentId}:${correction.page}:${correction.lineId}`; + const line = lines.get(key); + if (targets.has(key) || !line || line.lineSha256 !== correction.expectedLineSha256) throw conflict("Correction target is stale or duplicated", "CORRECTION_CONFLICT"); + targets.add(key); + line.text = correction.replacementText; + line.lineSha256 = sha256Hex(line.text); + } + for (const document of corrected.documents) for (const page of document.pages) if (page.ocr.lines.length > 0) { + page.ocr.text = page.candidateText = page.ocr.lines.map(({ text }) => text).join("\n"); + } + return corrected; +} + +function reviewedText(candidate: OcrReviewCandidate): string { + return candidate.documents.flatMap(({ pages }) => pages.filter(({ candidateText }) => candidateText).map(({ candidateText }) => candidateText)).join("\n\n"); +} diff --git a/tests/catalog/repository-ocr.test.ts b/tests/catalog/repository-ocr.test.ts index 2868023..661939a 100644 --- a/tests/catalog/repository-ocr.test.ts +++ b/tests/catalog/repository-ocr.test.ts @@ -124,6 +124,30 @@ test("expired leases return to queued without clearing remote recovery identity" ]); }); +test("candidate lifecycle persistence requires exact native, OCR, and metrics evidence", async () => { + const calls: Array<{ sql: string; params?: unknown[] }> = []; + const repository = new CatalogRepository(poolFor(async (sql, params) => { + calls.push({ sql, params }); + if (sql.includes("SELECT version_id, document_id")) return { rowCount: 1, rows: [{ version_id: "version-1", document_id: "document-1" }] }; + if (sql.includes("SELECT DISTINCT")) return { rowCount: 1, rows: [{ version_id: "version-1" }] }; + if (sql.includes("SELECT 1 FROM rag_ocr_jobs")) return { rowCount: 0, rows: [] }; + if (sql.includes("FROM rag_ocr_jobs")) return { rowCount: 1, rows: [{ ...jobRow, ocr_state: "succeeded" }] }; + return { rowCount: 1, rows: [] }; + }) as never); + const jobs = await repository.listOcrJobs("version-1"); + assert.equal(await repository.completeOcrJob("job-1", { pages: [{ page: 1, text: "OCR", metrics: { inkCoverage: 0.4 } }] } as never), true); + assert.doesNotMatch(calls[2]!.sql, /candidate_text_hash/); + assert.deepEqual(await repository.listOcrVersionsAwaitingCandidate(), ["version-1"]); + await repository.persistOcrCandidate("version-1", [{ + documentId: "document-1", page: 1, method: "ocr", nativeTextSha256: "a".repeat(64), + ocrTextSha256: "b".repeat(64), candidateTextSha256: "b".repeat(64), metrics: { inkCoverage: 0.4 }, risks: ["CBGO4a"] + }]); + + assert.deepEqual(jobs.map(({ jobId, state }) => [jobId, state]), [["job-1", "succeeded"]]); + assert.match(calls[6]!.sql, /native_text_hash = \$8.*ocr_text_hash IS NOT DISTINCT FROM \$9.*metrics = \$10::jsonb/s); + assert.deepEqual(calls[6]!.params?.slice(-3), ["a".repeat(64), "b".repeat(64), JSON.stringify({ inkCoverage: 0.4 })]); +}); + test("review transitions enforce indexing, approval, and rejection state guards", async () => { const successfulSql: string[] = []; const success = new CatalogRepository(poolFor(async (sql) => { @@ -145,3 +169,18 @@ test("review transitions enforce indexing, approval, and rejection state guards" (error) => error instanceof CatalogError && error.code === "INVALID_VERSION_STATE" ); }); + +test("review context reconstructs exact lifecycle identity and rejects incomplete page evidence", async () => { + const rows = [{ document_id: "document-1", page_number: 1, native_text_hash: "a".repeat(64), ocr_text_hash: "b".repeat(64), + candidate_text_hash: "c".repeat(64), metrics: { inkCoverage: 0.4 }, risk_tokens: ["CBGO4a"] }]; + const repository = new CatalogRepository(poolFor(async (sql) => sql.includes("JOIN rag_sources") + ? { rowCount: 1, rows: [{ version_id: "version-1", source_id: "source-1", state: "review_required", base_active_version_id: null, + current_active_version_id: null, activate_requested: false, processing_fingerprint: "fingerprint", metadata_hash: "metadata" }] } + : { rowCount: 1, rows }) as never); + + const context = await repository.loadOcrReviewContext("version-1"); + assert.deepEqual(context?.pages[0], { documentId: "document-1", page: 1, nativeTextSha256: "a".repeat(64), + ocrTextSha256: "b".repeat(64), candidateTextSha256: "c".repeat(64), metrics: { inkCoverage: 0.4 }, risks: ["CBGO4a"] }); + rows[0]!.candidate_text_hash = null as never; + await assert.rejects(repository.loadOcrReviewContext("version-1"), (error) => error instanceof CatalogError && error.code === "OCR_ARTIFACT_INTEGRITY_FAILED"); +}); diff --git a/tests/ocr/approval-indexing-wiring.test.ts b/tests/ocr/approval-indexing-wiring.test.ts new file mode 100644 index 0000000..f8def24 --- /dev/null +++ b/tests/ocr/approval-indexing-wiring.test.ts @@ -0,0 +1,81 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { createApp } from "../../src/app.js"; +import { env } from "../../src/config/env.js"; +import { OcrReadyIndexingService } from "../../src/modules/ocr/indexing.js"; +import { OcrReviewService } from "../../src/modules/ocr/review.js"; +import type { OcrReviewCandidate } from "../../src/modules/ocr/review.js"; + +const token = "admin-token"; + +function candidate(state: string): OcrReviewCandidate { + return { + versionId: "11111111-2222-4333-8444-555555555555", sourceId: "src:scan.pdf", state: state as OcrReviewCandidate["state"], + baseActiveVersionId: "active-1", currentActiveVersionId: "active-1", activateRequested: true, processingFingerprint: "fp", metadataHash: "mh", + candidateSha256: "candidate-hash", + documents: [{ documentId: "doc:scan", pages: [{ page: 1, imageUrl: "/img", nativeText: "", ocr: { text: "CBGO4a", lines: [{ lineId: "line-1", text: "CBGO4a", confidence: 0.9, bbox: [0, 0, 1, 1], lineSha256: "hash" }] }, candidateText: "CBGO4a", differences: [], risks: [] }] }] + } as OcrReviewCandidate; +} + +test("production approval triggers durable indexing and reports ready", async (context) => { + const previous = { lifecycle: env.knowledgeLifecycleEnforced, enabled: env.ocrIngestEnabled, admin: env.lifecycleAdminToken }; + Object.assign(env, { knowledgeLifecycleEnforced: true, ocrIngestEnabled: true, lifecycleAdminToken: token }); + context.after(() => Object.assign(env, previous)); + let reviewed: string | undefined; + let indexedCandidate: string | undefined; + const review = new OcrReviewService({ + async loadCandidate() { return candidate(reviewed ? "indexing" : "review_required"); }, + async commitApproval(input) { reviewed = input.reviewedText; }, + async commitRejection() { throw new Error("rejection is out of scope"); } + }); + const indexing = new OcrReadyIndexingService({ + async findReusableVersion() { return undefined; }, + async indexReviewed(value) { indexedCandidate = value.reviewedText; return 1; }, + async markReady() { /* ready marker */ }, + async settleReusable() { throw new Error("no reusable path"); }, + async activateVersion() { throw new Error("activation is a separate unit"); } + } as never); + const server = createApp({ reviewService: review, indexingService: indexing, startReconciler: false }).listen(0); + context.after(() => server.close()); + const address = server.address(); + assert.ok(address && typeof address === "object"); + const base = `http://127.0.0.1:${address.port}`; + const versionId = "11111111-2222-4333-8444-555555555555"; + + const approval = await fetch(`${base}/ingestions/${versionId}/approve`, { + method: "POST", + headers: { authorization: `Bearer ${token}`, "content-type": "application/json" }, + body: JSON.stringify({ candidateSha256: "candidate-hash", expectedActiveVersionId: "active-1", reviewedBy: "reviewer", corrections: [] }) + }); + assert.equal(approval.status, 200); + const body = await approval.json() as { state: string; activated: boolean }; + assert.equal(body.state, "ready"); + assert.equal(body.activated, false); + assert.equal(reviewed, "CBGO4a"); + assert.equal(indexedCandidate, "CBGO4a"); +}); + +test("approval without a configured indexing service stays at the indexing boundary with 503", async (context) => { + const previous = { lifecycle: env.knowledgeLifecycleEnforced, enabled: env.ocrIngestEnabled, admin: env.lifecycleAdminToken }; + Object.assign(env, { knowledgeLifecycleEnforced: true, ocrIngestEnabled: true, lifecycleAdminToken: token }); + context.after(() => Object.assign(env, previous)); + const review = new OcrReviewService({ + async loadCandidate() { return candidate("review_required"); }, + async commitApproval() { /* recorded */ }, + async commitRejection() { throw new Error("out of scope"); } + }); + const server = createApp({ reviewService: review, startReconciler: false }).listen(0); + context.after(() => server.close()); + const address = server.address(); + assert.ok(address && typeof address === "object"); + const base = `http://127.0.0.1:${address.port}`; + const versionId = "11111111-2222-4333-8444-555555555555"; + + const approval = await fetch(`${base}/ingestions/${versionId}/approve`, { + method: "POST", + headers: { authorization: `Bearer ${token}`, "content-type": "application/json" }, + body: JSON.stringify({ candidateSha256: "candidate-hash", expectedActiveVersionId: "active-1", reviewedBy: "reviewer", corrections: [] }) + }); + assert.equal(approval.status, 503); + assert.equal(((await approval.json()) as { code?: string }).code, "OCR_INDEXING_UNAVAILABLE"); +}); \ No newline at end of file diff --git a/tests/ocr/client.test.ts b/tests/ocr/client.test.ts index 37f7a08..5bb4202 100644 --- a/tests/ocr/client.test.ts +++ b/tests/ocr/client.test.ts @@ -1,10 +1,14 @@ import assert from "node:assert/strict"; -import { access, lstat, mkdir, mkdtemp, readFile, stat, symlink, utimes } from "node:fs/promises"; +import { access, lstat, mkdir, mkdtemp, readFile, stat, symlink, utimes, writeFile } from "node:fs/promises"; import os from "node:os"; import path from "node:path"; import test from "node:test"; -import { OcrClient, OcrClientError } from "../../src/modules/ocr/client.js"; +import { OcrClient, OcrClientError, type OcrResult } from "../../src/modules/ocr/client.js"; import { + persistComposedCandidateArtifact, + persistOcrResultArtifact, + readComposedCandidateArtifact, + readOcrResultArtifact, resolveArtifactPath, stageOcrArtifacts, sweepOrphanArtifacts @@ -53,6 +57,7 @@ function result(overrides: Record = {}): Record { assert.deepEqual(request, { url: "http://ocr.internal:8000/v1/jobs/job%20%2F%201", method: "DELETE", authorization: "Bearer delete-token" }); }); +test("review image transfer validates authentication, identity, content type, and hash", async () => { + const png = Buffer.from("89504e470d0a1a0a0102", "hex"); + const client = new OcrClient({ + baseUrl: "http://ocr.internal:8000", + token: "image-token", + fetch: (async (_input, init) => new Response(png, { status: 200, headers: { + "content-type": "image/png", "x-document-sha256": documentSha256, + "x-content-sha256": sha256Hex(png), "x-page-number": "1", + "x-seen-authorization": new Headers(init?.headers).get("authorization") ?? "" + } })) as typeof fetch + }); + + assert.deepEqual(await client.getReviewImage(jobId, 1, documentSha256), { bytes: png, sha256: sha256Hex(png) }); + await assert.rejects(client.getReviewImage(jobId, 2, documentSha256), /integrity validation failed/); +}); + test("artifact staging writes private originals and a verifiable canonical manifest", async (context) => { const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-ocr-artifacts-")); context.after(async () => { await import("node:fs/promises").then(({ rm }) => rm(rootDirectory, { recursive: true, force: true })); }); @@ -196,6 +217,114 @@ test("artifact staging writes private originals and a verifiable canonical manif assert.equal((await stat(staged.manifestPath)).mode & 0o777, 0o600); }); +test("durable OCR result survives simulated remote deletion with exact restart readback", async (context) => { + const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-ocr-result-")); + context.after(() => import("node:fs/promises").then(({ rm }) => rm(rootDirectory, { recursive: true, force: true }))); + const versionId = "66666666-6666-4666-8666-666666666666"; + const documentId = "doc:durable-result"; + await stageOcrArtifacts({ + rootDirectory, + versionId, + createdAt: "2026-09-16T10:00:00.000Z", + documents: [{ documentId, documentKey: "scan.pdf", bytes: document }] + }); + let remoteResult: OcrResult | undefined = result() as unknown as OcrResult; + + const persisted = await persistOcrResultArtifact({ rootDirectory, versionId, documentId, result: remoteResult }); + remoteResult = undefined; + const restarted = await readOcrResultArtifact({ + rootDirectory, + versionId, + documentId, + jobId, + documentSha256, + pages: [1, 3], + artifactSha256: persisted.artifactSha256 + }); + const repeated = await persistOcrResultArtifact({ rootDirectory, versionId, documentId, result: restarted }); + + assert.equal(remoteResult, undefined); + assert.deepEqual(restarted, result()); + assert.equal(repeated.artifactPath, persisted.artifactPath); + assert.equal(repeated.artifactSha256, persisted.artifactSha256); + assert.equal((await stat(persisted.artifactPath)).mode & 0o777, 0o600); +}); + +test("OCR result readback fails closed on corruption and artifact identity mismatch", async (context) => { + const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-ocr-corrupt-")); + context.after(() => import("node:fs/promises").then(({ rm }) => rm(rootDirectory, { recursive: true, force: true }))); + const versionId = "77777777-7777-4777-8777-777777777777"; + const documentId = "doc:corruption"; + await stageOcrArtifacts({ rootDirectory, versionId, createdAt: "2026-09-16T10:00:00.000Z", documents: [{ documentId, documentKey: "scan.pdf", bytes: document }] }); + const persisted = await persistOcrResultArtifact({ rootDirectory, versionId, documentId, result: result() as unknown as OcrResult }); + const envelope = JSON.parse(await readFile(persisted.artifactPath, "utf8")) as Record; + + await writeFile(persisted.artifactPath, canonicalJson({ ...envelope, documentId: "doc:other" }), { mode: 0o600 }); + await assert.rejects(readOcrResultArtifact({ rootDirectory, versionId, documentId, jobId, documentSha256, pages: [1, 3] }), /identity validation failed/); + await writeFile(persisted.artifactPath, "{corrupt", { mode: 0o600 }); + await assert.rejects(readOcrResultArtifact({ rootDirectory, versionId, documentId, jobId, documentSha256, pages: [1, 3] }), /integrity validation failed/); +}); + +test("restart-safe candidate composition persists exact native, OCR, and blank page evidence", async (context) => { + const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-ocr-candidate-")); + context.after(() => import("node:fs/promises").then(({ rm }) => rm(rootDirectory, { recursive: true, force: true }))); + const versionId = "88888888-8888-4888-8888-888888888888"; + const documentId = "doc:candidate"; + await stageOcrArtifacts({ rootDirectory, versionId, createdAt: "2026-09-16T10:00:00.000Z", documents: [{ + documentId, documentKey: "mixed.pdf", bytes: document, requestedPages: [2, 3], pages: [ + { page: 1, text: "Native page", rasterCoverage: 0, textSha256: sha256Hex("Native page") }, + { page: 2, text: "weak native", rasterCoverage: 0.8, textSha256: sha256Hex("weak native") }, + { page: 3, text: "", rasterCoverage: 0, textSha256: sha256Hex("") } + ] + }] }); + const ocrText = "codigo CBGO4a has enough OCR characters for the quality gate"; + await persistOcrResultArtifact({ rootDirectory, versionId, documentId, result: result({ documentSha256, pages: [ + { page: 2, width: 1700, height: 2200, processingMs: 10, text: ocrText, + metrics: { lineCount: 1, nonWhitespaceCharacters: 50, inkCoverage: 0.4, medianConfidence: 0.95, p10Confidence: 0.95, lowConfidenceLineRatio: 0 }, + lines: [{ lineId: "p2-l1", text: ocrText, confidence: 0.95, bbox: [1, 2, 3, 4] }] }, + { page: 3, width: 1700, height: 2200, processingMs: 10, text: "", + metrics: { lineCount: 0, nonWhitespaceCharacters: 0, inkCoverage: 0.001, medianConfidence: 0, p10Confidence: 0, lowConfidenceLineRatio: 0 }, lines: [] } + ] }) as unknown as OcrResult }); + + const persisted = await persistComposedCandidateArtifact({ rootDirectory, versionId, jobs: [{ documentId, remoteJobId: jobId, requestedPages: [2, 3], state: "succeeded" }] }); + const restarted = await readComposedCandidateArtifact({ rootDirectory, versionId, artifactSha256: persisted.artifactSha256 }); + + assert.deepEqual(restarted.documents[0]?.pages.map(({ method, candidateText }) => [method, candidateText]), [ + ["native", "Native page"], ["ocr", ocrText], ["blank", ""] + ]); + assert.equal(restarted.documents[0]?.pages[1]?.risks[0], "CBGO4a"); + assert.equal(restarted.candidateSha256, persisted.candidate.candidateSha256); + assert.equal((await stat(persisted.artifactPath)).mode & 0o777, 0o600); + await writeFile(persisted.artifactPath, "{corrupt", { mode: 0o600 }); + await assert.rejects(readComposedCandidateArtifact({ rootDirectory, versionId }), /candidate artifact integrity/); + await assert.rejects(persistComposedCandidateArtifact({ rootDirectory, versionId, jobs: [{ documentId, remoteJobId: jobId, requestedPages: [2, 3], state: "succeeded" }] }), /conflicts with durable content/); + await writeFile(persisted.artifactPath, canonicalJson(restarted), { mode: 0o600 }); + const manifest = JSON.parse(await readFile(path.join(rootDirectory, versionId, "manifest.json"), "utf8")) as { documents: Array<{ nativePagesPath: string }> }; + await import("node:fs/promises").then(({ rm }) => rm(path.join(path.dirname(resolveArtifactPath(path.join(rootDirectory, versionId), manifest.documents[0]!.nativePagesPath)), "ocr-result.json"))); + await assert.rejects(persistComposedCandidateArtifact({ rootDirectory, versionId, jobs: [{ documentId, remoteJobId: jobId, requestedPages: [2, 3], state: "succeeded" }] }), /result artifact integrity/); +}); + +test("candidate composition fails closed on quality failure and corrupt native evidence", async (context) => { + const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-ocr-candidate-fail-")); + context.after(() => import("node:fs/promises").then(({ rm }) => rm(rootDirectory, { recursive: true, force: true }))); + const versionId = "99999999-9999-4999-8999-999999999999"; + const documentId = "doc:blocked"; + const staged = await stageOcrArtifacts({ rootDirectory, versionId, createdAt: "2026-09-16T10:00:00.000Z", documents: [{ + documentId, documentKey: "scan.pdf", bytes: document, requestedPages: [1], + pages: [{ page: 1, text: "", rasterCoverage: 1, textSha256: sha256Hex("") }] + }] }); + await persistOcrResultArtifact({ rootDirectory, versionId, documentId, result: result({ pages: [{ + page: 1, width: 1700, height: 2200, processingMs: 10, text: "", lines: [], + metrics: { lineCount: 0, nonWhitespaceCharacters: 0, inkCoverage: 0.5, medianConfidence: 0, p10Confidence: 0, lowConfidenceLineRatio: 0 } + }] }) as unknown as OcrResult }); + const jobs = [{ documentId, remoteJobId: jobId, requestedPages: [1], state: "succeeded" as const }]; + + await assert.rejects(persistComposedCandidateArtifact({ rootDirectory, versionId, jobs }), /OCR_QUALITY_BLOCKED/); + const manifest = JSON.parse(await readFile(staged.manifestPath, "utf8")) as { documents: Array<{ nativePagesPath: string }> }; + await writeFile(resolveArtifactPath(staged.versionDirectory, manifest.documents[0]!.nativePagesPath), "{corrupt", { mode: 0o600 }); + await assert.rejects(persistComposedCandidateArtifact({ rootDirectory, versionId, jobs }), /native page artifact integrity/); +}); + test("safe resolution blocks traversal and orphan sweep preserves retained, recent, and non-directory entries", async (context) => { const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-ocr-sweep-")); context.after(async () => { await import("node:fs/promises").then(({ rm }) => rm(rootDirectory, { recursive: true, force: true })); }); diff --git a/tests/ocr/contracts-deploy.test.ts b/tests/ocr/contracts-deploy.test.ts index 8a5866d..222d639 100644 --- a/tests/ocr/contracts-deploy.test.ts +++ b/tests/ocr/contracts-deploy.test.ts @@ -23,6 +23,7 @@ test("OpenAPI describes authenticated OCR ingestion, status, review, approval, a for (const [path, method] of [ ["/ingestions/{versionId}", "get"], ["/ingestions/{versionId}/review", "get"], + ["/ingestions/{versionId}/documents/{documentId}/pages/{page}/image", "get"], ["/ingestions/{versionId}/approve", "post"], ["/ingestions/{versionId}/reject", "post"] ]) { @@ -53,8 +54,7 @@ test("OpenAPI OCR schemas preserve progress, review evidence, corrections, and d const approvalProperties = schemas.OcrApprovalRequest!.properties as Record; assert.deepEqual(approvalProperties.corrections!.items, { $ref: "#/components/schemas/OcrCorrection" }); assert.deepEqual(schemas.OcrDecisionResponse!.properties, { - versionId: { type: "string", format: "uuid" }, state: { type: "string", enum: ["ready", "active", "rejected"] }, activated: { type: "boolean" }, - activatedVersionId: { type: "string", format: "uuid" }, errorCode: { type: "string", enum: ["DUPLICATE_REUSABLE_VERSION"] } + versionId: { type: "string", format: "uuid" }, state: { type: "string", enum: ["indexing", "rejected"] }, activated: { type: "boolean" } }); }); diff --git a/tests/ocr/dispatcher.test.ts b/tests/ocr/dispatcher.test.ts index 65594d1..77e2886 100644 --- a/tests/ocr/dispatcher.test.ts +++ b/tests/ocr/dispatcher.test.ts @@ -9,6 +9,7 @@ import { CatalogError } from "../../src/modules/catalog/errors.js"; import { CatalogRepository } from "../../src/modules/catalog/repository.js"; import { KnowledgeLifecycleReconciler } from "../../src/modules/catalog/reconciler.js"; import { OcrDispatcher } from "../../src/modules/ocr/dispatcher.js"; +import { readReviewImageArtifact } from "../../src/modules/ocr/artifacts.js"; import { IngestService } from "../../src/modules/ingest/service.js"; import type { EmbeddingProvider } from "../../src/modules/embeddings/provider.js"; import type { VectorStoreClient } from "../../src/modules/vectorstore/client.js"; @@ -185,6 +186,9 @@ test("application wiring dispatches a newly accepted OCR job from durable artifa }); const pdf = buildPdf(""); let queuedJob: Record | undefined; + let runningJob: Record | undefined; + let completedJob: Record | undefined; + let documentSha256 = ""; let persistedKey: unknown; const dispatchCalls: string[] = []; const catalog = { @@ -198,14 +202,22 @@ test("application wiring dispatches a newly accepted OCR job from durable artifa return { versionId: input.versionId, versionNumber: 3 }; }, async markIndexing() {}, - async claimNextOcrJob() { const job = queuedJob; queuedJob = undefined; return job ? { ...job, state: "running", attemptCount: 1 } : undefined; }, - async setOcrRemoteJob() {}, + async claimNextOcrJob() { const job = queuedJob; queuedJob = undefined; runningJob = job ? { ...job, state: "running", attemptCount: 1 } : undefined; return runningJob; }, + async setOcrRemoteJob(_jobId: string, remoteJobId: string) { runningJob = { ...runningJob, remoteJobId }; }, async requeueOcrJob() { dispatchCalls.push("requeued"); }, - async completeOcrJob() { return false; }, async failOcrJob() {}, async markReviewRequired() {}, async markFailed() {} + async completeOcrJob() { runningJob = completedJob = { ...runningJob, state: "succeeded" }; dispatchCalls.push("completed"); return true; }, + async listOcrJobs() { return [runningJob]; }, async persistOcrCandidate() { dispatchCalls.push("candidate"); }, + async failOcrJob() {}, async markReviewRequired() { dispatchCalls.push("review-required"); }, async markFailed() {} }; + const png = Buffer.from("89504e470d0a1a0a0102", "hex"); const client = { - async submit(bytes: Buffer, _expected: unknown, key: string) { assert.equal(key, persistedKey); dispatchCalls.push(`${bytes.length}:${key}`); return { jobId: "runtime-remote" }; }, - async getStatus() { return { jobId: "runtime-remote", status: "queued", completedPages: 0, totalPages: 1, error: null }; } + async submit(bytes: Buffer, expected: { documentSha256: string }, key: string) { assert.equal(key, persistedKey); documentSha256 = expected.documentSha256; dispatchCalls.push(`${bytes.length}:${key}`); return { jobId: "runtime-remote" }; }, + async getStatus() { return { jobId: "runtime-remote", status: "succeeded", completedPages: 1, totalPages: 1, error: null }; }, + async getResult() { const text = "Factura FAT07 has enough OCR characters for review"; return { schemaVersion: "1", jobId: "runtime-remote", documentSha256, + engine: { name: "paddleocr", version: "3.4.0", runtime: "paddlepaddle-3.2.2", device: "cpu", configVersion: "ocr-v1", dpi: 200 }, + pages: [{ page: 1, width: 100, height: 100, processingMs: 1, text, metrics: { lineCount: 1, nonWhitespaceCharacters: 40, inkCoverage: 0.5, medianConfidence: 0.95, p10Confidence: 0.95, lowConfidenceLineRatio: 0 }, + lines: [{ lineId: "p1-l1", text, confidence: 0.95, bbox: [1, 2, 3, 4] }] }] }; }, + async getReviewImage() { return { bytes: png, sha256: sha256Hex(png) }; }, async delete() {} }; const server = createApp({ catalog: catalog as never, ocrClient: client as never, startReconciler: false }).listen(0); context.after(() => server.close()); @@ -220,10 +232,10 @@ test("application wiring dispatches a newly accepted OCR job from durable artifa const response = await fetch(`http://127.0.0.1:${address.port}/ingest/upload`, { method: "POST", body: form }); assert.equal(response.status, 202); assert.equal((await response.json() as { uploadedResource: string }).uploadedResource, "runtime.pdf"); - for (let attempt = 0; attempt < 20 && !dispatchCalls.includes("requeued"); attempt += 1) await new Promise((resolve) => setTimeout(resolve, 5)); - assert.equal(dispatchCalls.length, 2); + for (let attempt = 0; attempt < 40 && !dispatchCalls.includes("review-required"); attempt += 1) await new Promise((resolve) => setTimeout(resolve, 5)); assert.match(dispatchCalls[0]!, /^\d+:.+:ocr-v1:.+$/u); - assert.equal(dispatchCalls[1], "requeued"); + assert.deepEqual(dispatchCalls.slice(1), ["completed", "candidate", "review-required"]); + assert.deepEqual((await readReviewImageArtifact({ rootDirectory: path.join(directory, "artifacts"), versionId: String(completedJob?.versionId), documentId: String(completedJob?.documentId), page: 1 })).bytes, png); }); test("dispatcher recovers only expired work, reuses its remote job, and completes it once", async () => { @@ -263,10 +275,16 @@ test("dispatcher recovers only expired work, reuses its remote job, and complete async getResult() { calls.push("result"); return { pages: [{ page: 1 }] }; }, async delete() { calls.push("delete"); throw new Error("cleanup unavailable"); } }; - const dispatcher = new OcrDispatcher(repository, client as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") })); + const dispatcher = new OcrDispatcher( + repository, + client as never, + async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), + async (_job, result) => { calls.push("artifact-write"); return result; }, + async () => { calls.push("candidate-finalized"); } + ); assert.equal(await dispatcher.recoverExpiredLeases(), 1); - assert.deepEqual(calls, ["recover", "claim-exact", "status", "result", "complete", "review-required", "delete"]); + assert.deepEqual(calls, ["recover", "claim-exact", "status", "result", "artifact-write", "complete", "candidate-finalized", "review-required", "delete"]); }); test("dispatcher retains remote OCR state when durable result transfer fails", async () => { @@ -274,17 +292,52 @@ test("dispatcher retains remote OCR state when durable result transfer fails", a const job = { jobId: "job-1", versionId: "version-1", documentId: "document-1", remoteJobId: "remote-1", remoteIdempotencyKey: "key", state: "running" as const, requestedPages: [1], completedPages: 0, configVersion: "ocr-v1", attemptCount: 1, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt: null, errorCode: null, errorDetail: null }; const store = { async claimNextOcrJob() { return job; }, - async completeOcrJob() { calls.push("persist"); throw new Error("durable write failed"); }, + async completeOcrJob() { calls.push("complete"); return true; }, async failOcrJob() { calls.push("job-failed"); }, async markFailed() { calls.push("version-failed"); } }; const client = { async getStatus() { return { status: "succeeded" }; }, async getResult() { return { pages: [] }; }, async delete() { calls.push("delete"); } }; - const dispatcher = new OcrDispatcher(store as never, client as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") })); + const dispatcher = new OcrDispatcher( + store as never, + client as never, + async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), + async () => { calls.push("artifact-write"); throw new Error("durable write failed"); }, + async () => undefined + ); assert.equal(await dispatcher.runOnce(), "failed"); - assert.deepEqual(calls, ["persist", "job-failed", "version-failed"]); + assert.deepEqual(calls, ["artifact-write", "job-failed", "version-failed"]); +}); + +test("dispatcher retains remote OCR state when durable candidate finalization fails", async () => { + const calls: string[] = []; + const job = { jobId: "job-1", versionId: "version-1", documentId: "document-1", remoteJobId: "remote-1", remoteIdempotencyKey: "key", state: "running" as const, requestedPages: [1], completedPages: 0, configVersion: "ocr-v1", attemptCount: 1, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt: null, errorCode: null, errorDetail: null }; + const dispatcher = new OcrDispatcher({ + async claimNextOcrJob() { return job; }, async completeOcrJob() { calls.push("complete"); return true; }, + async failOcrJob() { calls.push("job-failed"); }, async markFailed(_versionId: string, code: string) { calls.push(`version-failed:${code}`); }, + async markReviewRequired() { calls.push("review-required"); } + } as never, { + async getStatus() { return { status: "succeeded" }; }, async getResult() { return { pages: [] }; }, async delete() { calls.push("delete"); } + } as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), + async (_job, result) => result, async () => { calls.push("candidate-write-readback"); throw new Error("OCR_QUALITY_BLOCKED"); }); + + assert.equal(await dispatcher.runOnce(), "failed"); + assert.deepEqual(calls, ["complete", "candidate-write-readback", "job-failed", "version-failed:OCR_QUALITY_BLOCKED"]); +}); + +test("dispatcher resumes candidate finalization after a restart boundary", async () => { + const calls: string[] = []; + const dispatcher = new OcrDispatcher({ + async listOcrVersionsAwaitingCandidate() { calls.push("scan"); return ["version-restart"]; }, + async markReviewRequired(versionId: string) { calls.push(`review:${versionId}`); }, + async markFailed() { calls.push("failed"); } + } as never, {} as never, async () => { throw new Error("not called"); }, async (_job, result) => result, + async (versionId) => { calls.push(`compose:${versionId}`); }); + + assert.equal(await dispatcher.recoverCompletedCandidates(), 1); + assert.deepEqual(calls, ["scan", "compose:version-restart", "review:version-restart"]); }); test("OCR exhaustion fails the leased candidate without activation or engine substitution", async () => { @@ -301,7 +354,13 @@ test("OCR exhaustion fails the leased candidate without activation or engine sub async markFailed(_versionId: string, code: string) { calls.push(`version:${code}`); } }; const client = { async submit() { throw new Error("OCR unavailable after retries"); } }; - const dispatcher = new OcrDispatcher(repository as never, client as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") })); + const dispatcher = new OcrDispatcher( + repository as never, + client as never, + async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), + async (_job, result) => result, + async () => undefined + ); assert.equal(await dispatcher.runOnce(), "failed"); assert.deepEqual(calls, ["job:OCR_DISPATCH_FAILED", "version:OCR_DISPATCH_FAILED"]); @@ -317,14 +376,15 @@ test("reconciler makes expired OCR leases dispatchable while leaving live leases }; const dispatcher = { async recoverExpiredLeases() { calls.push("ocr-recovery"); return 1; }, - async dispatchAvailable() { calls.push("ocr-dispatch"); return 2; } + async dispatchAvailable() { calls.push("ocr-dispatch"); return 2; }, + async recoverCompletedCandidates() { calls.push("candidate-recovery"); return 1; } }; const reconciler = new KnowledgeLifecycleReconciler(catalog as never, vectorStore(), dispatcher as never); const status = await reconciler.runOnce(); assert.equal(status.ocrLeasesRecovered, 1); - assert.deepEqual(calls, ["ocr-recovery", "ocr-dispatch"]); + assert.deepEqual(calls, ["ocr-recovery", "ocr-dispatch", "candidate-recovery"]); }); test("runtime HTTP routing returns native 201, OCR 202/status, and catalog-down 503", async () => { diff --git a/tests/ocr/e2e.test.ts b/tests/ocr/e2e.test.ts index 0a3914a..1b148f7 100644 --- a/tests/ocr/e2e.test.ts +++ b/tests/ocr/e2e.test.ts @@ -3,7 +3,7 @@ import test from "node:test"; import { createApp } from "../../src/app.js"; import { env } from "../../src/config/env.js"; import { OcrDispatcher } from "../../src/modules/ocr/dispatcher.js"; -import { OcrIndexingService } from "../../src/modules/ocr/indexing.js"; +import { OcrReadyIndexingService } from "../../src/modules/ocr/indexing.js"; import { OcrReviewService, type OcrReviewCandidate } from "../../src/modules/ocr/review.js"; import { sha256Hex } from "../../src/shared/utils/ids.js"; @@ -42,7 +42,7 @@ function candidate(): OcrReviewCandidate { }; } -test("local HTTP flow keeps native ingestion synchronous and activates only reviewed OCR content", async (context) => { +test("local HTTP flow keeps native ingestion synchronous and stops approved OCR at the indexing boundary", async (context) => { const previous = { lifecycle: env.knowledgeLifecycleEnforced, enabled: env.ocrIngestEnabled, token: env.lifecycleAdminToken }; Object.assign(env, { knowledgeLifecycleEnforced: true, ocrIngestEnabled: true, lifecycleAdminToken: token }); context.after(() => Object.assign(env, { @@ -90,9 +90,10 @@ test("local HTTP flow keeps native ingestion synchronous and activates only revi }, async commitRejection() { throw new Error("rejection is outside this flow"); } }); - const indexing = new OcrIndexingService({ + let indexingCalls = 0; + const indexing = new OcrReadyIndexingService({ async findReusableVersion() { return undefined; }, - async indexReviewed(value) { assert.equal(value.reviewedText, "CBG04a"); return 1; }, + async indexReviewed(value) { indexingCalls += 1; assert.equal(value.reviewedText, "CBG04a"); return 1; }, async markReady(versionId) { states.set(versionId, "ready"); }, async settleReusable() { throw new Error("candidate is not reusable"); }, async activateVersion(_sourceId, versionId, expected) { @@ -100,7 +101,7 @@ test("local HTTP flow keeps native ingestion synchronous and activates only revi states.set(versionId, "active"); return versionId; } - }); + } as never); const server = createApp({ ingestService: ingestService as never, catalog: catalog as never, ocrClient: {} as never, reviewService: review, indexingService: indexing, startReconciler: false }).listen(0); context.after(() => server.close()); const address = server.address(); @@ -125,9 +126,10 @@ test("local HTTP flow keeps native ingestion synchronous and activates only revi }) }); assert.equal(approval.status, 200); - assert.deepEqual(await approval.json(), { versionId: scanVersion, state: "active", activated: true, activatedVersionId: scanVersion }); + assert.deepEqual(await approval.json(), { versionId: scanVersion, state: "ready", activated: false }); assert.equal(reviewedText, "CBG04a"); - assert.equal((await fetch(`${base}/ingestions/${scanVersion}`, { headers: headers() }).then((response) => response.json()) as { state: string }).state, "active"); + assert.equal(indexingCalls, 1); + assert.equal((await fetch(`${base}/ingestions/${scanVersion}`, { headers: headers() }).then((response) => response.json()) as { state: string }).state, "ready"); assert.equal((await ingest("mixed.pdf")).status, 202); const mixed = await fetch(`${base}/ingestions/${mixedVersion}`, { headers: headers() }).then((response) => response.json()) as { documents: Array<{ pages: Array<{ method: string }> }> }; @@ -169,7 +171,9 @@ test("resends reuse identity, OCR exhaustion fails closed, and catalog absence r async requeueOcrJob() { calls.push("requeue"); }, async completeOcrJob() { calls.push("complete"); return false; }, async failOcrJob() { calls.push("job-failed"); }, async markReviewRequired() { calls.push("review"); }, async markFailed() { candidateState = "failed"; calls.push("version-failed"); } - }, { async submit() { throw new Error("OCR connection unavailable after bounded retries"); } } as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") })); + }, { async submit() { throw new Error("OCR connection unavailable after bounded retries"); } } as never, + async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), + async (_job, result) => result, async () => undefined); assert.equal(await dispatcher.runOnce(), "failed"); assert.deepEqual(calls, ["job-failed", "version-failed"]); assert.equal(candidateState, "failed"); diff --git a/tests/ocr/indexing-store.test.ts b/tests/ocr/indexing-store.test.ts new file mode 100644 index 0000000..c080e5c --- /dev/null +++ b/tests/ocr/indexing-store.test.ts @@ -0,0 +1,110 @@ +import assert from "node:assert/strict"; +import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; +import test from "node:test"; +import { CatalogError } from "../../src/modules/catalog/errors.js"; +import type { EmbeddingProvider } from "../../src/modules/embeddings/provider.js"; +import { PostgresOcrIndexingStore, OcrReadyIndexingService, type ApprovedOcrCandidate } from "../../src/modules/ocr/indexing.js"; +import { persistReviewedPagesArtifact } from "../../src/modules/ocr/artifacts.js"; +import { chunkDocument, documentalChunkingPolicy } from "../../src/modules/process/chunking.js"; +import type { VectorStoreClient } from "../../src/modules/vectorstore/client.js"; +import { buildChunkId, buildVersionedQdrantPointId, canonicalJson, normalizeContentForHash, sha256Hex } from "../../src/shared/utils/ids.js"; + +const versionId = "eeeeeeee-eeee-4eee-8eee-eeeeeeeeeeee"; +const documentId = "doc:indexed-review"; +const sourceId = "src:reviewed"; + +async function fixture(rootDirectory: string) { + const first = "Reviewed alpha content. ".repeat(110); + const second = "Reviewed beta content."; + const lines = [first, second].map((text, index) => ({ + lineId: `line-${index + 1}`, text, confidence: 0.97, bbox: [0, index * 10, 50, index * 10 + 8] as [number, number, number, number], lineSha256: sha256Hex(text) + })); + const pages = lines.map((line, index) => ({ page: index + 1, method: "ocr" as const, nativeText: "", rawOcrText: line.text, + candidateText: line.text, candidateTextSha256: sha256Hex(line.text), risks: [], lines: [line], + nativeTextSha256: sha256Hex(""), ocrTextSha256: sha256Hex(line.text), metrics: { medianConfidence: 0.97 } })); + const documents = [{ documentId, text: `${first}\n\n${second}`, textSha256: sha256Hex(`${first}\n\n${second}`), pages }]; + const candidateSha256 = sha256Hex(canonicalJson(documents)); + const reviewedText = `${first}\n\n${second}`; + await mkdir(path.join(rootDirectory, versionId)); + await writeFile(path.join(rootDirectory, versionId, "candidate-pages.json"), canonicalJson({ schemaVersion: "1", versionId, candidateSha256, documents }), { mode: 0o600 }); + await persistReviewedPagesArtifact({ rootDirectory, versionId, sourceId, candidateSha256, reviewedTextSha256: sha256Hex(reviewedText), + reviewedBy: "reviewer", documents: [{ documentId, pages: pages.map(({ page, candidateText, lines: ocrLines }) => ({ page, candidateText, ocr: { lines: ocrLines } })) }] }); + return { reviewedText, pages }; +} + +function candidate(reviewedText: string): ApprovedOcrCandidate { + return { versionId, sourceId, state: "indexing", activateRequested: true, expectedActiveVersionId: "active-1", + reviewedText, reviewedTextSha256: sha256Hex(reviewedText), processingFingerprint: "fingerprint", metadataHash: "metadata" }; +} + +function harness(pageCount: number, options: { embeddingFailure?: boolean; countOffset?: number } = {}) { + const sql: string[] = []; + const points: Array<{ id: string; vector: number[]; payload: Record }> = []; + let embedCalls = 0; + const rows = Array.from({ length: pageCount }, (_, index) => ({ + version_id: versionId, source_id: sourceId, version_number: 4, state: "indexing", base_active_version_id: "active-1", + current_active_version_id: "active-1", processing_fingerprint: "fingerprint", metadata_hash: "metadata", tags: ["reviewed"], + embedding_provider: "test-provider", embedding_model: "test-model", embedding_dimensions: 3, qdrant_collection: "rag_documents", + expected_document_count: 1, document_id: documentId, document_key: "review.pdf", title: "review.pdf", mime_type: "application/pdf", + page_number: index + 1, reviewed_text_hash: "pending" + })); + const query = async (statement: string, params?: unknown[]) => { + sql.push(statement); + if (/^(BEGIN|COMMIT|ROLLBACK|SET CONSTRAINTS)/u.test(statement.trim())) return { rowCount: 0, rows: [] }; + if (statement.includes("FROM rag_source_versions v")) return { rowCount: rows.length, rows }; + if (statement.includes("UPDATE rag_version_documents")) return { rowCount: 1, rows: [] }; + if (statement.includes("UPDATE rag_source_versions")) return { rowCount: 1, rows: [{ version_id: versionId }] }; + throw new Error(`Unexpected SQL: ${statement} ${String(params)}`); + }; + const pool = { query, async connect() { return { query, release() {} }; } }; + const embeddings: EmbeddingProvider = { providerName: "test-provider", modelName: "test-model", dimensions: 3, + async embed(input) { embedCalls += 1; if (options.embeddingFailure) throw new Error("embedding unavailable"); return input.map(() => [0.1, 0.2, 0.3]); } }; + const vectors: Partial = { kind: "fake", async upsert(chunks) { points.push(...chunks); }, + async countVersionPoints() { return points.length + (options.countOffset ?? 0); } }; + return { pool: pool as never, embeddings, vectors: vectors as VectorStoreClient, points, sql, embedCalls: () => embedCalls }; +} + +test("reviewed artifact indexing survives restart and writes canonical versioned chunks before ready", async (context) => { + const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-reviewed-index-")); + context.after(() => rm(rootDirectory, { recursive: true, force: true })); + const { reviewedText, pages } = await fixture(rootDirectory); + const runtime = harness(pages.length); + const originalQuery = (runtime.pool as { query: (sql: string, params?: unknown[]) => Promise<{ rowCount: number; rows: Array> }> }).query; + for (let index = 0; index < pages.length; index += 1) { + const row = (await originalQuery("SELECT * FROM rag_source_versions v", [])).rows[index]!; + row.reviewed_text_hash = sha256Hex(pages[index]!.candidateText); + } + const restarted = new PostgresOcrIndexingStore(runtime.pool, rootDirectory, runtime.embeddings, runtime.vectors); + const result = await new OcrReadyIndexingService(restarted).index(candidate(reviewedText)); + + const expected = chunkDocument("review.pdf", normalizeContentForHash(reviewedText), documentalChunkingPolicy); + assert.deepEqual(result, { versionId, state: "ready", activated: false }); + assert.equal(runtime.points.length, expected.length); + assert.deepEqual(runtime.points.map(({ id }) => id), expected.map((chunk) => buildVersionedQdrantPointId(versionId, buildChunkId(documentId, "documental", chunk.index)))); + assert.ok(runtime.points.every(({ payload }) => payload.processing_fingerprint === "fingerprint" && payload.source_version_id === versionId)); + assert.match(runtime.sql.join("\n"), /SET state = 'ready'/u); + assert.doesNotMatch(runtime.sql.join("\n"), /active_version_id\s*=/u); +}); + +test("identity, corrupt artifact, embedding, and partial-write failures remain fail closed", async (context) => { + const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-reviewed-index-fail-")); + context.after(() => rm(rootDirectory, { recursive: true, force: true })); + const { reviewedText, pages } = await fixture(rootDirectory); + for (const scenario of [ + { name: "identity", runtime: harness(pages.length), value: { ...candidate(reviewedText), processingFingerprint: "wrong" } }, + { name: "embedding", runtime: harness(pages.length, { embeddingFailure: true }), value: candidate(reviewedText) }, + { name: "count", runtime: harness(pages.length, { countOffset: -1 }), value: candidate(reviewedText) } + ]) { + await assert.rejects(new OcrReadyIndexingService(new PostgresOcrIndexingStore(scenario.runtime.pool, rootDirectory, scenario.runtime.embeddings, scenario.runtime.vectors)).index(scenario.value), + scenario.name === "count" ? (error) => error instanceof CatalogError && error.code === "SOURCE_VERSION_INCONSISTENT" : Error); + assert.doesNotMatch(scenario.runtime.sql.join("\n"), /SET state = 'ready'/u); + if (scenario.name === "identity") assert.equal(scenario.runtime.embedCalls(), 0); + } + await writeFile(path.join(rootDirectory, versionId, "reviewed-pages.json"), "corrupt", { mode: 0o600 }); + const corrupt = harness(pages.length); + await assert.rejects(new OcrReadyIndexingService(new PostgresOcrIndexingStore(corrupt.pool, rootDirectory, corrupt.embeddings, corrupt.vectors)).index(candidate(reviewedText)), /integrity validation failed/u); + assert.equal(corrupt.embedCalls(), 0); + assert.equal(corrupt.points.length, 0); +}); diff --git a/tests/ocr/review.test.ts b/tests/ocr/review.test.ts index d180e7f..08cb271 100644 --- a/tests/ocr/review.test.ts +++ b/tests/ocr/review.test.ts @@ -1,10 +1,15 @@ import assert from "node:assert/strict"; +import { access, mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; import test from "node:test"; import { createApp } from "../../src/app.js"; import { env } from "../../src/config/env.js"; import { CatalogError } from "../../src/modules/catalog/errors.js"; import { OcrIndexingService, type ApprovedOcrCandidate, type OcrIndexingStore } from "../../src/modules/ocr/indexing.js"; -import { OcrReviewService, type OcrReviewCandidate } from "../../src/modules/ocr/review.js"; +import { OcrReviewService, PostgresOcrReviewStore, type OcrReviewCandidate } from "../../src/modules/ocr/review.js"; +import { DurableOcrReviewReader } from "../../src/modules/ocr/review.js"; +import { persistComposedCandidateArtifact, persistOcrResultArtifact, persistReviewImageArtifacts, readReviewedPagesArtifact, stageOcrArtifacts } from "../../src/modules/ocr/artifacts.js"; import { sha256Hex } from "../../src/shared/utils/ids.js"; function candidate(state: OcrReviewCandidate["state"] = "review_required"): OcrReviewCandidate { @@ -183,3 +188,151 @@ test("playground serves the authenticated OCR review controls and audit fields", for (const marker of ["Authorization", "/review", "/approve", "/reject", "candidateSha256", "expectedLineSha256", "imageUrl", "nativeText", "confidence", "bbox", "differences", "risks"]) assert.match(script, new RegExp(marker)); assert.match(styles, /\.review-page/); }); + +test("production review reader survives restart and fails closed on unauthorized or mismatched durable evidence", async (context) => { + const previous = { token: env.lifecycleAdminToken, enabled: env.ocrIngestEnabled, root: env.ocrArtifactRoot }; + const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-review-loader-")); + const versionId = "aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa"; + const documentId = "doc:review"; + const original = Buffer.from("%PDF-review"); + const text = "codigo CBGO4a has enough OCR characters for durable review"; + Object.assign(env, { lifecycleAdminToken: "review-token", ocrIngestEnabled: true, ocrArtifactRoot: rootDirectory }); + context.after(async () => { Object.assign(env, { lifecycleAdminToken: previous.token, ocrIngestEnabled: previous.enabled, ocrArtifactRoot: previous.root }); await rm(rootDirectory, { recursive: true, force: true }); }); + await stageOcrArtifacts({ rootDirectory, versionId, createdAt: "2026-09-16T12:00:00.000Z", documents: [{ + documentId, documentKey: "review.pdf", bytes: original, requestedPages: [1], + pages: [{ page: 1, text: "weak native", rasterCoverage: 1, textSha256: sha256Hex("weak native") }] + }] }); + const result = { schemaVersion: "1" as const, jobId: "ocr-review", documentSha256: sha256Hex(original), + engine: { name: "paddleocr" as const, version: "3.4.0" as const, runtime: "paddlepaddle-3.2.2" as const, device: "cpu" as const, configVersion: "ocr-v1" as const, dpi: 200 as const }, + pages: [{ page: 1, width: 100, height: 100, processingMs: 1, text, + metrics: { lineCount: 1, nonWhitespaceCharacters: 50, inkCoverage: 0.5, medianConfidence: 0.95, p10Confidence: 0.95, lowConfidenceLineRatio: 0 }, + lines: [{ lineId: "p1-l1", text, confidence: 0.95, bbox: [1, 2, 30, 10] as [number, number, number, number] }] }] }; + await persistOcrResultArtifact({ rootDirectory, versionId, documentId, result }); + const png = Buffer.from("89504e470d0a1a0a0102", "hex"); + const images = await persistReviewImageArtifacts({ rootDirectory, versionId, documentId, images: [{ page: 1, bytes: png, sha256: sha256Hex(png) }] }); + const { candidate: durable } = await persistComposedCandidateArtifact({ rootDirectory, versionId, jobs: [{ documentId, remoteJobId: "ocr-review", requestedPages: [1], state: "succeeded" }] }); + const page = durable.documents[0]!.pages[0]!; + let reads = 0; + const contextValue = { versionId, sourceId: "source-1", state: "review_required" as const, baseActiveVersionId: null, currentActiveVersionId: null, + activateRequested: false, processingFingerprint: "fingerprint", metadataHash: "metadata", pages: [{ documentId, page: 1, + nativeTextSha256: page.nativeTextSha256, ocrTextSha256: page.ocrTextSha256, candidateTextSha256: page.candidateTextSha256, metrics: page.metrics, risks: page.risks }] }; + const catalog = { async loadOcrReviewContext() { reads += 1; return contextValue; } }; + const restarted = new DurableOcrReviewReader(catalog, rootDirectory); + assert.equal((await restarted.view(versionId)).documents[0]!.pages[0]!.ocr.lines[0]!.lineSha256, sha256Hex(text)); + assert.deepEqual((await restarted.image(versionId, documentId, 1)).bytes, png); + + const server = createApp({ catalog: catalog as never, startReconciler: false }).listen(0); + context.after(() => server.close()); + const address = server.address(); + assert.ok(address && typeof address === "object"); + const base = `http://127.0.0.1:${address.port}/ingestions/${versionId}`; + assert.equal((await fetch(`${base}/review`)).status, 401); + assert.equal(reads, 2); + assert.equal((await fetch(`${base}/review`, { headers: { authorization: "Bearer review-token" } })).status, 200); + const image = await fetch(`${base}/documents/${encodeURIComponent(documentId)}/pages/1/image`, { headers: { authorization: "Bearer review-token" } }); + assert.equal(image.status, 200); + assert.deepEqual(Buffer.from(await image.arrayBuffer()), png); + assert.equal((await fetch(`${base}/approve`, { method: "POST", headers: { authorization: "Bearer review-token", "content-type": "application/json" }, body: "{}" })).status, 503); + assert.equal((await fetch(`${base}/reject`, { method: "POST", headers: { authorization: "Bearer review-token", "content-type": "application/json" }, body: "{}" })).status, 503); + + contextValue.pages[0]!.candidateTextSha256 = "0".repeat(64); + assert.equal((await fetch(`${base}/review`, { headers: { authorization: "Bearer review-token" } })).status, 500); + contextValue.pages[0]!.candidateTextSha256 = page.candidateTextSha256; + await assert.rejects(restarted.image(versionId, "../escape", 1), /not found/i); + await writeFile(images.images[0]!.artifactPath, "corrupt", { mode: 0o600 }); + await assert.rejects(restarted.image(versionId, documentId, 1), /integrity validation failed/); +}); + +function decisionPool(current: OcrReviewCandidate, state = "review_required", failReviewedPage = false) { + const sql: string[] = []; + const client = { + async query(statement: string) { + sql.push(statement); + if (/^(BEGIN|COMMIT|ROLLBACK|SET CONSTRAINTS)/u.test(statement.trim())) return { rowCount: 0, rows: [] }; + if (statement.includes("FOR UPDATE OF v, s")) return { rowCount: 1, rows: [{ + version_id: current.versionId, source_id: current.sourceId, state, base_active_version_id: current.baseActiveVersionId, + current_active_version_id: current.currentActiveVersionId, activate_requested: current.activateRequested, + processing_fingerprint: current.processingFingerprint, metadata_hash: current.metadataHash + }] }; + if (statement.includes("FROM rag_document_pages") && statement.includes("FOR UPDATE")) return { rowCount: 1, rows: [{ + document_id: current.documents[0]!.documentId, page_number: 1, + native_text_hash: sha256Hex(current.documents[0]!.pages[0]!.nativeText), + ocr_text_hash: sha256Hex(current.documents[0]!.pages[0]!.ocr.text), + candidate_text_hash: sha256Hex(current.documents[0]!.pages[0]!.candidateText), risk_tokens: current.documents[0]!.pages[0]!.risks + }] }; + if (failReviewedPage && statement.includes("SET reviewed_text_hash")) return { rowCount: 0, rows: [] }; + return { rowCount: 1, rows: [] }; + }, + release() {} + }; + return { pool: { connect: async () => client } as never, sql }; +} + +test("production decision store commits corrected reviewed pages and the indexing transition in one locked transaction", async (context) => { + const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-review-decision-")); + const current = candidate(); + current.versionId = "bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb"; + await mkdir(path.join(rootDirectory, current.versionId)); + context.after(() => rm(rootDirectory, { recursive: true, force: true })); + const { pool, sql } = decisionPool(current); + const store = new PostgresOcrReviewStore(pool, { async view() { return structuredClone(current); } }, rootDirectory); + const approved = await new OcrReviewService(store).approve(current.versionId, { + candidateSha256: current.candidateSha256, expectedActiveVersionId: current.baseActiveVersionId, reviewedBy: "admin", + corrections: [{ documentId: "document-1", page: 1, lineId: "line-1", expectedLineSha256: sha256Hex("CBGO4a"), replacementText: "CBG04a" }] + }); + + const reviewed = await readReviewedPagesArtifact({ rootDirectory, versionId: current.versionId, + candidateSha256: current.candidateSha256, reviewedTextSha256: approved.reviewedTextSha256 }); + assert.equal(reviewed.documents[0]!.pages[0]!.candidateText, "CBG04a\nFATo7"); + assert.match(sql.join("\n"), /FOR UPDATE OF v, s/u); + assert.match(sql.join("\n"), /INSERT INTO rag_review_corrections/u); + assert.match(sql.join("\n"), /reviewed_text_hash/u); + assert.match(sql.join("\n"), /SET state = 'indexing'/u); + assert.equal(sql.at(-1)?.trim(), "COMMIT"); +}); + +test("production decisions reject stale, replayed, and path-corrupt writes without durable or lifecycle side effects", async (context) => { + const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-review-conflict-")); + const current = candidate(); + current.versionId = "cccccccc-cccc-4ccc-8ccc-cccccccccccc"; + context.after(() => rm(rootDirectory, { recursive: true, force: true })); + const input = { candidateSha256: current.candidateSha256, expectedActiveVersionId: current.baseActiveVersionId, + reviewedBy: "admin", corrections: [] }; + + for (const conflictState of ["indexing", "rejected"]) { + const { pool, sql } = decisionPool(current, conflictState); + const service = new OcrReviewService(new PostgresOcrReviewStore(pool, { async view() { return structuredClone(current); } }, rootDirectory)); + await assert.rejects(service.approve(current.versionId, input), (error) => error instanceof CatalogError && error.statusCode === 409); + assert.doesNotMatch(sql.join("\n"), /INSERT INTO rag_review_corrections|SET state = 'indexing'/u); + } + + const { pool, sql } = decisionPool(current); + const service = new OcrReviewService(new PostgresOcrReviewStore(pool, { async view() { return structuredClone(current); } }, rootDirectory)); + await assert.rejects(service.approve(current.versionId, input), /ENOENT|artifact directory/u); + assert.doesNotMatch(sql.join("\n"), /INSERT INTO rag_review_corrections|SET state = 'indexing'/u); + await assert.rejects(access(path.join(rootDirectory, current.versionId, "reviewed-pages.json"))); + + await mkdir(path.join(rootDirectory, current.versionId)); + const partial = decisionPool(current, "review_required", true); + const partialService = new OcrReviewService(new PostgresOcrReviewStore(partial.pool, { async view() { return structuredClone(current); } }, rootDirectory)); + await assert.rejects(partialService.approve(current.versionId, input), /page evidence changed/u); + assert.equal(partial.sql.at(-1)?.trim(), "ROLLBACK"); + await assert.rejects(access(path.join(rootDirectory, current.versionId, "reviewed-pages.json"))); +}); + +test("production rejection records the bound reason without reviewed, correction, indexing, or activation artifacts", async (context) => { + const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-review-rejection-")); + const current = candidate(); + current.versionId = "dddddddd-dddd-4ddd-8ddd-dddddddddddd"; + await mkdir(path.join(rootDirectory, current.versionId)); + context.after(() => rm(rootDirectory, { recursive: true, force: true })); + const { pool, sql } = decisionPool(current); + const service = new OcrReviewService(new PostgresOcrReviewStore(pool, { async view() { return structuredClone(current); } }, rootDirectory)); + + assert.deepEqual(await service.reject(current.versionId, { candidateSha256: current.candidateSha256, reviewedBy: " admin ", reason: " unreadable code " }), + { versionId: current.versionId, state: "rejected", activated: false }); + const statements = sql.join("\n"); + assert.match(statements, /error_code = 'OCR_REJECTED'.*error_detail =/su); + assert.doesNotMatch(statements, /rag_review_corrections|reviewed_text_hash|state = 'indexing'/u); + await assert.rejects(access(path.join(rootDirectory, current.versionId, "reviewed-pages.json"))); +});