From 07e6ed23fb925e7f74421573f5337b26dc95e392 Mon Sep 17 00:00:00 2001 From: Paco POR-CORREO Date: Thu, 17 Sep 2026 14:27:48 +0200 Subject: [PATCH] fix(ocr): harden review recovery --- Dockerfile | 7 +++- docs/PENDIENTES_RAG.md | 46 ++++++++++++++++++++- docs/SISTEMA_RAG_BASE.md | 17 +++++++- migrations/003_ocr_recovery_audit.sql | 9 +++++ ocr-service/Dockerfile | 4 ++ ocr-service/app/main.py | 7 ++-- ocr-service/tests/test_api.py | 5 ++- src/api/openapi.ts | 30 +++++++++++++- src/app.ts | 48 ++++++++++++++++------ src/config/env.ts | 1 + src/modules/catalog/reconciler.ts | 4 +- src/modules/catalog/repository.ts | 2 +- src/modules/ocr/artifacts.ts | 22 +++++++++- src/modules/ocr/dispatcher.ts | 15 ------- src/modules/ocr/review.ts | 58 ++++++++++++++++++++++++++- tests/catalog/migration-003.test.ts | 11 +++++ tests/catalog/repository-ocr.test.ts | 2 +- tests/ocr/contracts-deploy.test.ts | 7 ++++ tests/ocr/dispatcher.test.ts | 18 +-------- tests/ocr/review.test.ts | 48 +++++++++++++++++++--- 20 files changed, 294 insertions(+), 67 deletions(-) create mode 100644 migrations/003_ocr_recovery_audit.sql create mode 100644 tests/catalog/migration-003.test.ts diff --git a/Dockerfile b/Dockerfile index ea3dc14..963dee2 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,4 +1,6 @@ FROM node:22-bookworm-slim AS build +ARG RAG_VERSION=0.1.0 +ARG BUILD_REVISION=unknown WORKDIR /app COPY package.json package-lock.json tsconfig.json ./ RUN npm ci @@ -8,8 +10,11 @@ COPY public ./public RUN npm run build FROM node:22-bookworm-slim AS runtime +ARG RAG_VERSION=0.1.0 +ARG BUILD_REVISION=unknown WORKDIR /app -ENV NODE_ENV=production OCR_ARTIFACT_ROOT=/data/ingestions +ENV NODE_ENV=production OCR_ARTIFACT_ROOT=/data/ingestions RAG_VERSION=${RAG_VERSION} +LABEL org.opencontainers.image.revision=${BUILD_REVISION} COPY package.json package-lock.json ./ RUN npm ci --omit=dev COPY --from=build /app/dist ./dist diff --git a/docs/PENDIENTES_RAG.md b/docs/PENDIENTES_RAG.md index 3ef132f..f002786 100644 --- a/docs/PENDIENTES_RAG.md +++ b/docs/PENDIENTES_RAG.md @@ -1,11 +1,44 @@ # Pendientes priorizados del RAG -**Ultima actualizacion:** 2026-09-13 +**Ultima actualizacion:** 2026-09-17 **Responsable de la priorizacion:** Usuario **Estado:** Activo Este documento es la fuente canonica del orden de trabajo pendiente del modulo RAG. La numeracion ya refleja la prioridad final indicada por el usuario. +## Hoja de ruta OCR prioritaria + +Esta secuencia tiene prioridad sobre la aceptacion productiva pendiente de FacturaTech. No cerrar v4, crear una candidata nueva, reingestar FacturaTech, aprobar, indexar ni activar contenido fuera del orden indicado. + +### Fase 1. Corregir y verificar el codigo + +1. Sustituir errores OCR previsibles por respuestas estructuradas, seguras y accionables. +2. Incorporar una recuperacion administrativa controlada para candidatas con evidencia que no cumple el contrato. +3. Exponer versiones de producto y metadatos de build verificables en RAG y OCR. +4. Cubrir errores, recuperacion, reinicios y versiones con pruebas; desplegar y verificar en produccion. + +**Salida:** el comportamiento nuevo esta validado y v4 permanece intacta. + +### Fase 2. Resolver la candidata v4 + +1. Consultar v4 mediante el contrato de errores ya corregido. +2. Confirmar que su evidencia no cumple los requisitos. +3. Aplicar una operacion administrativa solo con aprobacion explicita del usuario. +4. Verificar que la version activa no cambia ni se indexa contenido. +5. Confirmar que la fuente ya no queda bloqueada para una nueva candidata. + +**Salida:** v4 queda resuelta de forma auditable, sin crear todavia otra candidata. + +### Fase 3. Finalizar la aceptacion de FacturaTech + +1. Crear una nueva candidata no activada para el documento. +2. Verificar la evidencia durable y revisar sus 34 entradas. +3. Comprobar `CBG04a`, `FAT07`, `DSAU08` y `NSAV06`. +4. Requerir aprobacion humana antes de indexar o activar. +5. Completar la tarea SDD 7.4 con su evidencia de validacion. + +**Salida:** la ingesta OCR de FacturaTech queda aceptada y activada de forma segura. + ## 1. Documentacion y descubrimiento de la API **Estado:** Completado y validado definitivamente en produccion el 2026-09-08. @@ -68,6 +101,17 @@ El punto 2 queda cerrado. El punto 3, OCR integrado en la ingesta, puede comenza - Verificar el resultado antes de sustituir una fuente vigente. - Evitar que se repita el problema detectado con el PDF de FacturaTech. +### Paquete de mejora posterior: OCR reutilizable para entradas visuales + +**Prioridad:** Diferida; no bloquea la correccion actual ni la aceptacion de FacturaTech. + +- Mantener el contrato actual de trabajos PDF. +- Aceptar directamente `image/jpeg` y `image/png`, ademas de PDF. +- Normalizar cada entrada visual a un resultado por pagina con texto, lineas, posiciones y metricas. +- Emitir resultados OCR o errores estructurados sin depender del RAG. +- Mantener RAG como consumidor opcional del OCR y responsable de versionado, revision e indexacion de conocimiento. +- Definir limites de tamano, dimensiones y paginas, autenticacion, retencion y pruebas para cada formato. + ## 4. Mejora del retrieval - Anadir busqueda hibrida semantica y textual para codigos exactos como `FAT07`, `504` o `SQLSTATE[23505]`. diff --git a/docs/SISTEMA_RAG_BASE.md b/docs/SISTEMA_RAG_BASE.md index 851e14f..691ab5b 100644 --- a/docs/SISTEMA_RAG_BASE.md +++ b/docs/SISTEMA_RAG_BASE.md @@ -2,9 +2,9 @@ **Proyecto:** Workspace de tools IA para empresas **Modulo:** RAG -**Ultima actualizacion:** 2026-09-11 +**Ultima actualizacion:** 2026-09-17 **Ultima modificacion por:** Agente RAG 2 -**Estado:** V1 operativa; ciclo de vida desplegado y pendiente de preparar PostgreSQL y migrar el corpus +**Estado:** RAG y OCR desplegados; pendientes correcciones de revision OCR y versionado verificable de servicios --- @@ -115,6 +115,19 @@ Esta mejora esta implementada y desplegada. Falta preparar PostgreSQL, registrar Los pendientes vigentes, su prioridad y su estado se mantienen en [`PENDIENTES_RAG.md`](./PENDIENTES_RAG.md). +### Flujo RAG y OCR + +RAG y OCR son servicios separados con responsabilidades distintas: + +- RAG recibe documentos compatibles, extrae texto nativo, conserva sus versiones y convierte el texto aprobado en conocimiento consultable. +- OCR recibe un PDF y una lista cerrada de paginas, procesa solo esas paginas y devuelve texto, lineas, posiciones y metricas de calidad. + +Para un PDF, RAG revisa primero todas sus paginas. Usa OCR cuando el texto nativo de una pagina es insuficiente o cuando su contenido rasterizado supera el umbral configurado. Reune todas las paginas seleccionadas y las envia en un solo trabajo OCR; no procesa una pagina, reinicia el documento y continua en ciclos. + +El resultado OCR queda como evidencia revisable. RAG conserva la responsabilidad de revision humana, indexacion, versionado y activacion. OCR no escribe en PostgreSQL ni en Qdrant, ni decide que contenido se activa. + +Actualmente OCR procesa PDFs. La admision directa de imagenes JPG y PNG es una mejora futura registrada en [`PENDIENTES_RAG.md`](./PENDIENTES_RAG.md). El detalle de contratos, umbrales y estados esta en [`CONTRATO_CICLO_VIDA_Y_OCR.md`](./CONTRATO_CICLO_VIDA_Y_OCR.md). + --- ## Alcance de este documento diff --git a/migrations/003_ocr_recovery_audit.sql b/migrations/003_ocr_recovery_audit.sql new file mode 100644 index 0000000..dc3ee86 --- /dev/null +++ b/migrations/003_ocr_recovery_audit.sql @@ -0,0 +1,9 @@ +CREATE TABLE IF NOT EXISTS rag_ocr_recovery_audit ( + recovery_id uuid PRIMARY KEY DEFAULT gen_random_uuid(), + version_id uuid NOT NULL REFERENCES rag_source_versions(version_id) ON DELETE RESTRICT, + recovered_by text NOT NULL, + reason text NOT NULL, + outcome text NOT NULL CHECK (outcome IN ('closed_failed')), + error_code text NOT NULL, + created_at timestamptz NOT NULL DEFAULT now() +); diff --git a/ocr-service/Dockerfile b/ocr-service/Dockerfile index 336d3f5..cdb210a 100644 --- a/ocr-service/Dockerfile +++ b/ocr-service/Dockerfile @@ -1,4 +1,6 @@ FROM python:3.11.13-slim-bookworm +ARG OCR_VERSION=0.1.0 +ARG BUILD_REVISION=unknown LABEL resource.cpu.max="3" \ resource.memory.max="5GiB" @@ -7,7 +9,9 @@ ENV PYTHONDONTWRITEBYTECODE=1 \ PYTHONUNBUFFERED=1 \ OCR_JOBS_DB="/data/jobs/jobs.db" \ OCR_LOAD_ENGINE="1" \ + OCR_VERSION=${OCR_VERSION} \ HOME="/opt/ocr-home" +LABEL org.opencontainers.image.revision=${BUILD_REVISION} RUN apt-get update \ && apt-get install --yes --no-install-recommends libgl1 libglib2.0-0 libgomp1 \ diff --git a/ocr-service/app/main.py b/ocr-service/app/main.py index 2fabdd8..7438856 100644 --- a/ocr-service/app/main.py +++ b/ocr-service/app/main.py @@ -167,14 +167,15 @@ def create_app( @application.get("/health/live") def live() -> dict[str, str]: - return {"status": "ok"} + return {"status": "ok", "service": "ocr", "version": os.getenv("OCR_VERSION", "0.1.0")} @application.get("/health/ready") - def ready(response: Response) -> dict[str, int | bool]: + def ready(response: Response) -> dict[str, int | bool | str]: queue.sweep_expired() if not engine_ready: response.status_code = 503 - return {"ready": engine_ready, "queueDepth": queue.depth(), "queueCapacity": QUEUE_CAPACITY, "concurrency": 1} + return {"ready": engine_ready, "queueDepth": queue.depth(), "queueCapacity": QUEUE_CAPACITY, "concurrency": 1, + "service": "ocr", "version": os.getenv("OCR_VERSION", "0.1.0")} @application.post("/v1/jobs", status_code=202, dependencies=[Depends(authorize)]) async def create_job( diff --git a/ocr-service/tests/test_api.py b/ocr-service/tests/test_api.py index 689199e..9a04895 100644 --- a/ocr-service/tests/test_api.py +++ b/ocr-service/tests/test_api.py @@ -62,8 +62,9 @@ def client(tmp_path: Path) -> TestClient: def test_auth_rejects_missing_and_wrong_bearer_and_health_exposes_no_secret(client: TestClient): live = client.get("/health/live") ready = client.get("/health/ready") - assert live.json() == {"status": "ok"} - assert ready.json() == {"ready": True, "queueDepth": 0, "queueCapacity": 3, "concurrency": 1} + assert live.json() == {"status": "ok", "service": "ocr", "version": "0.1.0"} + assert ready.json() == {"ready": True, "queueDepth": 0, "queueCapacity": 3, "concurrency": 1, + "service": "ocr", "version": "0.1.0"} for authorization in (None, "Bearer wrong-token"): headers = {"Idempotency-Key": key_for(request_for())} if authorization: diff --git a/src/api/openapi.ts b/src/api/openapi.ts index 3808269..b175653 100644 --- a/src/api/openapi.ts +++ b/src/api/openapi.ts @@ -280,6 +280,22 @@ export const openApiDocument = { } } }, + "/ingestions/{versionId}/recover": { + post: { + tags: ["OCR Review"], summary: "Close an unrecoverable OCR review candidate", security: [{ bearerAuth: [] }], + parameters: [{ name: "versionId", in: "path", required: true, schema: { type: "string", format: "uuid" } }], + requestBody: { required: true, content: jsonContent(ref("OcrRecoveryRequest")) }, + responses: { + "200": jsonResponse("Candidate was closed without approval, indexing, or activation.", ref("OcrRecoveryResponse")), + "400": jsonResponse("Recovery actor and reason are required.", ref("Error")), + "401": jsonResponse("Missing or invalid lifecycle admin token.", ref("Error")), + "404": jsonResponse("OCR candidate does not exist.", ref("Error")), + "409": jsonResponse("Candidate state or evidence precondition does not permit recovery.", ref("Error")), + "422": jsonResponse("OCR artifact format or integrity is invalid.", ref("Error")), + "503": serverError, "500": serverError + } + } + }, "/ingestions/{versionId}/documents/{documentId}/pages/{page}/image": { get: { tags: ["OCR Review"], summary: "Get a private OCR review page image", security: [{ bearerAuth: [] }], @@ -474,7 +490,8 @@ export const openApiDocument = { properties: { ok: { type: "boolean", const: false }, error: { type: "string" }, - code: { type: "string" } + code: { type: "string" }, + action: { type: "string" } } }, Scope: { @@ -630,6 +647,14 @@ export const openApiDocument = { type: "object", required: ["candidateSha256", "reviewedBy", "reason"], properties: { candidateSha256: { type: "string", pattern: "^[a-f0-9]{64}$" }, reviewedBy: { type: "string" }, reason: { type: "string" } } }, + OcrRecoveryRequest: { + type: "object", required: ["recoveredBy", "reason"], + properties: { recoveredBy: { type: "string" }, reason: { type: "string" } } + }, + OcrRecoveryResponse: { + type: "object", required: ["versionId", "state", "outcome"], + properties: { versionId: { type: "string", format: "uuid" }, state: { type: "string", const: "failed" }, outcome: { type: "string", const: "closed_failed" } } + }, OcrDecisionResponse: { type: "object", required: ["versionId", "state", "activated"], properties: { versionId: { type: "string", format: "uuid" }, state: { type: "string", enum: ["indexing", "rejected"] }, activated: { type: "boolean" } } @@ -863,10 +888,11 @@ export const openApiDocument = { }, HealthResponse: { type: "object", - required: ["ok", "service", "environment", "embeddings", "answer", "vectorStore", "postgres", "knowledgeLifecycle", "parsers", "chunking"], + required: ["ok", "service", "version", "environment", "embeddings", "answer", "vectorStore", "postgres", "knowledgeLifecycle", "parsers", "chunking"], properties: { ok: { type: "boolean" }, service: { type: "string", const: "rag" }, + version: { type: "string" }, environment: { type: "string" }, embeddings: ref("ProviderModel"), answer: ref("ProviderModel"), diff --git a/src/app.ts b/src/app.ts index a79c3bc..0388f54 100644 --- a/src/app.ts +++ b/src/app.ts @@ -23,7 +23,7 @@ import { CatalogError, CatalogRepository } from "./modules/catalog/repository.js import { OcrClient } from "./modules/ocr/client.js"; import { OcrDispatcher } from "./modules/ocr/dispatcher.js"; import { persistComposedCandidateArtifact, persistOcrResultArtifact, persistReviewImageArtifacts, readOcrArtifactPageNumbers, resolveArtifactPath } from "./modules/ocr/artifacts.js"; -import { DurableOcrReviewReader, OcrReviewService, PostgresOcrReviewStore } from "./modules/ocr/review.js"; +import { classifyOcrReviewError, DurableOcrReviewReader, OcrReviewRecoveryService, 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"; @@ -39,6 +39,7 @@ interface AppOptions { catalog?: CatalogRepository; ocrClient?: OcrClient; reviewService?: Pick; + reviewRecoveryService?: Pick; indexingService?: Pick; startReconciler?: boolean; } @@ -89,8 +90,10 @@ export function createApp(options: AppOptions = {}) { 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)) + const reviewStore = catalogPool && durableReview ? new PostgresOcrReviewStore(catalogPool, durableReview, ocr.artifactRoot) : undefined; + const reviewService = options.reviewService ?? (reviewStore ? new OcrReviewService(reviewStore) : undefined); + const reviewRecoveryService = options.reviewRecoveryService ?? (reviewStore && durableReview + ? new OcrReviewRecoveryService(durableReview, reviewStore) : undefined); const retrieveService = new RetrieveService(embeddingProvider, vectorStore, catalog); const answerService = new AnswerService(retrieveService); @@ -110,6 +113,16 @@ export function createApp(options: AppOptions = {}) { res.status(statusCode).json(body); } + function sendOcrReviewError(res: express.Response, error: unknown) { + const classified = classifyOcrReviewError(error); + if (classified.statusCode === 500) console.error("OCR review request failed", error); + const action = classified.code === "OCR_CANDIDATE_NOT_FOUND" ? "verify_version_id" + : classified.code === "OCR_REVIEW_STATE_INVALID" ? "inspect_ingestion_status" + : classified.statusCode === 409 || classified.statusCode === 422 ? "use_admin_recovery" + : "retry_or_contact_support"; + res.status(classified.statusCode).json({ ok: false, error: classified.message, code: classified.code, action }); + } + function dispatchAcceptedOcr(result: Awaited>): void { if ("phase" in result) void ocrDispatcher?.dispatchAvailable(); } @@ -192,6 +205,7 @@ export function createApp(options: AppOptions = {}) { res.status(ok ? 200 : 503).json({ ok, service: "rag", + version: env.ragVersion, environment: env.nodeEnv, embeddings: { provider: embeddingProvider.providerName, @@ -280,7 +294,7 @@ export function createApp(options: AppOptions = {}) { } res.json(status); } catch (error) { - sendError(res, error, "Unknown ingestion status error"); + sendOcrReviewError(res, error); } }); @@ -295,7 +309,21 @@ export function createApp(options: AppOptions = {}) { try { res.json(await reader.view(String(req.params.versionId))); } catch (error) { - sendError(res, error, "Unknown OCR review error"); + sendOcrReviewError(res, error); + } + }); + + app.post("/ingestions/:versionId/recover", async (req, res) => { + if (!requireOcrEnabled(res)) return; + if (!requireLifecycleAdmin(req, res)) return; + if (!reviewRecoveryService) { + res.status(503).json({ ok: false, error: "OCR recovery service is not configured", code: "OCR_RECOVERY_UNAVAILABLE" }); + return; + } + try { + res.json(await reviewRecoveryService.recover(String(req.params.versionId), req.body)); + } catch (error) { + sendOcrReviewError(res, error); } }); @@ -306,7 +334,7 @@ export function createApp(options: AppOptions = {}) { 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"); } + } catch (error) { sendOcrReviewError(res, error); } }); app.post("/ingestions/:versionId/approve", async (req, res) => { @@ -324,9 +352,7 @@ export function createApp(options: AppOptions = {}) { } 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"); - } + } catch (error) { sendOcrReviewError(res, error); } }); app.post("/ingestions/:versionId/reject", async (req, res) => { @@ -338,9 +364,7 @@ export function createApp(options: AppOptions = {}) { } try { res.json(await reviewService.reject(String(req.params.versionId), req.body)); - } catch (error) { - sendError(res, error, "Unknown OCR rejection error"); - } + } catch (error) { sendOcrReviewError(res, error); } }); app.get("/models/answer", async (_req, res) => { diff --git a/src/config/env.ts b/src/config/env.ts index b2d3beb..0e88f71 100644 --- a/src/config/env.ts +++ b/src/config/env.ts @@ -33,6 +33,7 @@ function booleanEnv(name: string, fallback: boolean): boolean { export const env = { nodeEnv: process.env.NODE_ENV ?? "development", + ragVersion: process.env.RAG_VERSION ?? "0.1.0", port: Number(process.env.PORT ?? 3000), qdrantUrl: requireEnv("QDRANT_URL", "http://localhost:6333"), qdrantApiKey: process.env.QDRANT_API_KEY ?? "", diff --git a/src/modules/catalog/reconciler.ts b/src/modules/catalog/reconciler.ts index dfc5e71..bcc7f14 100644 --- a/src/modules/catalog/reconciler.ts +++ b/src/modules/catalog/reconciler.ts @@ -24,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 ) {} @@ -58,7 +58,6 @@ 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); @@ -83,7 +82,6 @@ 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 3dfb11c..831c334 100644 --- a/src/modules/catalog/repository.ts +++ b/src/modules/catalog/repository.ts @@ -582,7 +582,7 @@ export class CatalogRepository { 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"); + if (pages.rows.some(({ candidate_text_hash }) => !candidate_text_hash)) throw new CatalogError("OCR candidate lifecycle evidence is incomplete", 409, "OCR_REVIEW_EVIDENCE_INCOMPLETE"); 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, diff --git a/src/modules/ocr/artifacts.ts b/src/modules/ocr/artifacts.ts index d003b55..9801f93 100644 --- a/src/modules/ocr/artifacts.ts +++ b/src/modules/ocr/artifacts.ts @@ -6,6 +6,16 @@ import { isValidOcrResult, type OcrResult } from "./client.js"; import { composeCandidate, type CandidatePage } from "./composition.js"; import { classifyOcrPage } from "./detection.js"; +export class OcrArtifactError extends Error { + constructor( + message: string, + readonly statusCode: 409 | 422, + readonly code: "OCR_ARTIFACT_UNAVAILABLE" | "OCR_ARTIFACT_FORMAT_INVALID" | "OCR_ARTIFACT_INTEGRITY_FAILED" + ) { + super(message); + } +} + 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"); @@ -222,9 +232,17 @@ export async function persistComposedCandidateArtifact(input: { 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"); + const artifactPath = path.join(path.resolve(input.rootDirectory), input.versionId, "candidate-pages.json"); + const file = await lstat(artifactPath).catch(() => { throw new OcrArtifactError("OCR candidate artifact is unavailable", 409, "OCR_ARTIFACT_UNAVAILABLE"); }); + if (!file.isFile()) throw new OcrArtifactError("OCR candidate artifact has an invalid format", 422, "OCR_ARTIFACT_FORMAT_INVALID"); + let value: unknown; + try { + value = await readPrivateJson(artifactPath, input.artifactSha256, "OCR candidate artifact integrity validation failed"); + } catch { + throw new OcrArtifactError("OCR candidate artifact integrity validation failed", 422, "OCR_ARTIFACT_INTEGRITY_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"); + || value.candidateSha256 !== sha256Hex(canonicalJson(value.documents))) throw new OcrArtifactError("OCR candidate artifact integrity validation failed", 422, "OCR_ARTIFACT_INTEGRITY_FAILED"); return value as unknown as ComposedCandidateArtifact; } diff --git a/src/modules/ocr/dispatcher.ts b/src/modules/ocr/dispatcher.ts index 6f168c7..cf1d95b 100644 --- a/src/modules/ocr/dispatcher.ts +++ b/src/modules/ocr/dispatcher.ts @@ -11,7 +11,6 @@ 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 }; @@ -96,20 +95,6 @@ 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/review.ts b/src/modules/ocr/review.ts index 18c4ddd..0c09687 100644 --- a/src/modules/ocr/review.ts +++ b/src/modules/ocr/review.ts @@ -1,7 +1,7 @@ import { CatalogError } from "../catalog/errors.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"; +import { OcrArtifactError, persistReviewedPagesArtifact, readComposedCandidateArtifact, readReviewImageArtifact, removeReviewedPagesArtifact } from "./artifacts.js"; export interface OcrReviewLine { lineId: string; @@ -59,6 +59,10 @@ export interface OcrReviewStore { }): Promise; } +export interface OcrReviewRecoveryStore { + closeUnrecoverableCandidate(input: { versionId: string; recoveredBy: string; reason: string; errorCode: string }): Promise; +} + export interface OcrReviewContext { versionId: string; sourceId: string; state: string; baseActiveVersionId: string | null; currentActiveVersionId: string | null; activateRequested: boolean; processingFingerprint: string; metadataHash: string; @@ -134,6 +138,25 @@ export class PostgresOcrReviewStore implements OcrReviewStore { }); } + async closeUnrecoverableCandidate(input: Parameters[0]): Promise { + await withTransaction(this.pool, async (client) => { + const version = await client.query<{ state: string }>("SELECT state FROM rag_source_versions WHERE version_id = $1 FOR UPDATE", [input.versionId]); + if (!version.rowCount) throw new CatalogError("OCR review candidate not found", 404, "REVIEW_NOT_FOUND"); + if (version.rows[0]!.state !== "review_required") throw conflict("Version is not awaiting OCR review", "INVALID_VERSION_STATE"); + const result = await client.query( + `UPDATE rag_source_versions SET state = 'failed', error_code = $2, error_detail = $3, retention_due_at = now() + interval '7 days' + WHERE version_id = $1 AND state = 'review_required'`, + [input.versionId, input.errorCode, input.reason.slice(0, 2000)] + ); + if (result.rowCount !== 1) throw conflict("Version is not awaiting OCR review", "INVALID_VERSION_STATE"); + await client.query( + `INSERT INTO rag_ocr_recovery_audit(version_id, recovered_by, reason, outcome, error_code) + VALUES ($1, $2, $3, 'closed_failed', $4)`, + [input.versionId, input.recoveredBy, input.reason.slice(0, 2000), input.errorCode] + ); + }); + } + 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; @@ -264,6 +287,39 @@ export class OcrReviewService { } } +export class OcrReviewRecoveryService { + constructor(private readonly reader: Pick, private readonly store: OcrReviewRecoveryStore) {} + + async recover(versionId: string, input: { recoveredBy: string; reason: string }): Promise<{ versionId: string; state: "failed"; outcome: "closed_failed" }> { + const recoveredBy = input?.recoveredBy?.trim(); + const reason = input?.reason?.trim(); + if (!recoveredBy || !reason) throw new CatalogError("Recovery actor and reason are required", 400, "OCR_RECOVERY_INVALID"); + try { + await this.reader.view(versionId); + } catch (error) { + const classified = classifyOcrReviewError(error); + if (classified.statusCode !== 409 && classified.statusCode !== 422) throw classified; + await this.store.closeUnrecoverableCandidate({ versionId, recoveredBy, reason, errorCode: classified.code }); + return { versionId, state: "failed", outcome: "closed_failed" }; + } + throw conflict("OCR review candidate has valid durable evidence", "OCR_RECOVERY_NOT_REQUIRED"); + } +} + +export function classifyOcrReviewError(error: unknown): CatalogError { + if (error instanceof CatalogError) { + if (error.code === "REVIEW_NOT_FOUND") return new CatalogError("OCR review candidate was not found", 404, "OCR_CANDIDATE_NOT_FOUND"); + if (error.code === "INVALID_VERSION_STATE") return new CatalogError("OCR review is not available in the current lifecycle state", 409, "OCR_REVIEW_STATE_INVALID"); + if (error.code === "OCR_ARTIFACT_INTEGRITY_FAILED") return new CatalogError("OCR review evidence is incomplete", 409, "OCR_REVIEW_EVIDENCE_INCOMPLETE"); + return error; + } + if (error instanceof OcrArtifactError) return new CatalogError(error.message, error.statusCode, error.code); + if (error instanceof Error && /OCR (?:candidate|review).*?(?:artifact|evidence|validation)|OCR review image/i.test(error.message)) { + return new CatalogError("OCR review artifact integrity validation failed", 422, "OCR_ARTIFACT_INTEGRITY_FAILED"); + } + return new CatalogError("Unexpected OCR review failure", 500, "OCR_REVIEW_UNEXPECTED"); +} + function applyCorrections(candidate: OcrReviewCandidate, corrections: OcrCorrection[]): OcrReviewCandidate { const corrected = structuredClone(candidate); const lines = new Map(); diff --git a/tests/catalog/migration-003.test.ts b/tests/catalog/migration-003.test.ts new file mode 100644 index 0000000..1995631 --- /dev/null +++ b/tests/catalog/migration-003.test.ts @@ -0,0 +1,11 @@ +import assert from "node:assert/strict"; +import { readFile } from "node:fs/promises"; +import test from "node:test"; + +test("migration 003 records explicit OCR recovery decisions", async () => { + const sql = await readFile(new URL("../../migrations/003_ocr_recovery_audit.sql", import.meta.url), "utf8"); + assert.match(sql, /CREATE TABLE IF NOT EXISTS rag_ocr_recovery_audit/i); + assert.match(sql, /version_id uuid NOT NULL REFERENCES rag_source_versions\(version_id\)/i); + assert.match(sql, /recovered_by text NOT NULL/i); + assert.match(sql, /outcome text NOT NULL CHECK \(outcome IN \('closed_failed'\)\)/i); +}); diff --git a/tests/catalog/repository-ocr.test.ts b/tests/catalog/repository-ocr.test.ts index 661939a..8f1e1c9 100644 --- a/tests/catalog/repository-ocr.test.ts +++ b/tests/catalog/repository-ocr.test.ts @@ -182,5 +182,5 @@ test("review context reconstructs exact lifecycle identity and rejects incomplet 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"); + await assert.rejects(repository.loadOcrReviewContext("version-1"), (error) => error instanceof CatalogError && error.code === "OCR_REVIEW_EVIDENCE_INCOMPLETE"); }); diff --git a/tests/ocr/contracts-deploy.test.ts b/tests/ocr/contracts-deploy.test.ts index 222d639..2711423 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}/recover", "post"], ["/ingestions/{versionId}/documents/{documentId}/pages/{page}/image", "get"], ["/ingestions/{versionId}/approve", "post"], ["/ingestions/{versionId}/reject", "post"] @@ -35,13 +36,17 @@ test("OpenAPI describes authenticated OCR ingestion, status, review, approval, a const approve = api.paths["/ingestions/{versionId}/approve"]!.post!; const reject = api.paths["/ingestions/{versionId}/reject"]!.post!; + const recover = api.paths["/ingestions/{versionId}/recover"]!.post!; assert.ok("409" in (approve.responses as Record)); assert.ok("400" in (reject.responses as Record)); assert.ok("409" in (reject.responses as Record)); + assert.ok("422" in (recover.responses as Record)); assert.deepEqual((approve.requestBody as { content: Record }).content["application/json"]!.schema, { $ref: "#/components/schemas/OcrApprovalRequest" }); assert.deepEqual((reject.requestBody as { content: Record }).content["application/json"]!.schema, { $ref: "#/components/schemas/OcrRejectionRequest" }); + assert.deepEqual((recover.requestBody as { content: Record }).content["application/json"]!.schema, + { $ref: "#/components/schemas/OcrRecoveryRequest" }); assert.match(buildApiHelp().authentication, /Bearer.*OCR/u); }); @@ -70,6 +75,8 @@ test("OCR deployment defaults and container wiring match the private durable con const dockerfile = await readFile(new URL("../../Dockerfile", import.meta.url), "utf8"); assert.match(dockerfile, /VOLUME \["\/data\/ingestions"\]/u); assert.match(dockerfile, /ENV NODE_ENV=production OCR_ARTIFACT_ROOT=\/data\/ingestions/u); + assert.match(dockerfile, /RAG_VERSION=\$\{RAG_VERSION\}/u); + assert.match(dockerfile, /org\.opencontainers\.image\.revision/u); assert.match(dockerfile, /USER node/u); assert.match(dockerfile, /dist\/modules\/catalog\/migrations\.js.*dist\/server\.js/u); assert.doesNotMatch(dockerfile, /(?:COPY|ADD)\s+\.env|ENV\s+(?:OCR_INTERNAL_TOKEN|LIFECYCLE_ADMIN_TOKEN)=/u); diff --git a/tests/ocr/dispatcher.test.ts b/tests/ocr/dispatcher.test.ts index 77e2886..672dca9 100644 --- a/tests/ocr/dispatcher.test.ts +++ b/tests/ocr/dispatcher.test.ts @@ -327,19 +327,6 @@ test("dispatcher retains remote OCR state when durable candidate finalization fa 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 () => { const calls: string[] = []; const repository = { @@ -376,15 +363,14 @@ 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 recoverCompletedCandidates() { calls.push("candidate-recovery"); return 1; } + async dispatchAvailable() { calls.push("ocr-dispatch"); return 2; } }; 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", "candidate-recovery"]); + assert.deepEqual(calls, ["ocr-recovery", "ocr-dispatch"]); }); test("runtime HTTP routing returns native 201, OCR 202/status, and catalog-down 503", async () => { diff --git a/tests/ocr/review.test.ts b/tests/ocr/review.test.ts index 08cb271..eb139d1 100644 --- a/tests/ocr/review.test.ts +++ b/tests/ocr/review.test.ts @@ -7,9 +7,9 @@ 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, PostgresOcrReviewStore, type OcrReviewCandidate } from "../../src/modules/ocr/review.js"; +import { classifyOcrReviewError, 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 { OcrArtifactError, 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 { @@ -46,7 +46,7 @@ test("review HTTP rejects unauthorized and premature access without exposing art value = candidate("indexing"); const premature = await fetch(url, { headers: { authorization: "Bearer review-token" } }); assert.equal(premature.status, 409); - assert.deepEqual(await premature.json(), { ok: false, error: "Version is not awaiting OCR review", code: "INVALID_VERSION_STATE" }); + assert.deepEqual(await premature.json(), { ok: false, error: "OCR review is not available in the current lifecycle state", code: "OCR_REVIEW_STATE_INVALID", action: "inspect_ingestion_status" }); }); test("review exposes audit detail and commits current corrections as one immutable set", async () => { @@ -155,7 +155,7 @@ test("authenticated rejection is durable, conflict-safe, and never indexes or ac assert.equal(reads, 0); const invalid = await fetch(`${url}/reject`, { method: "POST", headers: { authorization: "Bearer review-token", "content-type": "application/json" }, body: "{}" }); assert.equal(invalid.status, 400); - assert.deepEqual(await invalid.json(), { ok: false, error: "Candidate hash, reviewer, and rejection reason are required", code: "INVALID_REJECTION" }); + assert.deepEqual(await invalid.json(), { ok: false, error: "Candidate hash, reviewer, and rejection reason are required", code: "INVALID_REJECTION", action: "retry_or_contact_support" }); assert.equal(reads, 0); const stale = await fetch(`${url}/reject`, { method: "POST", headers: { authorization: "Bearer review-token", "content-type": "application/json" }, body: JSON.stringify({ candidateSha256: "stale", reviewedBy: "admin", reason: "Unreadable code" }) }); assert.equal(stale.status, 409); @@ -173,6 +173,44 @@ test("authenticated rejection is durable, conflict-safe, and never indexes or ac assert.equal(audits.length, 1); }); +test("review routes expose safe actionable errors and recovery stays explicit", async (context) => { + const previous = { token: env.lifecycleAdminToken, enabled: env.ocrIngestEnabled }; + Object.assign(env, { lifecycleAdminToken: "review-token", ocrIngestEnabled: true }); + context.after(() => Object.assign(env, previous)); + const recoveries: Array<{ recoveredBy: string; reason: string }> = []; + const review = { + async view() { throw new Error("unexpected database connection failure"); }, + async approve() { throw new Error("not called"); }, + async reject() { throw new Error("not called"); } + }; + const recovery = { + async recover(_versionId: string, input: { recoveredBy: string; reason: string }) { + recoveries.push(input); + return { versionId: "version-2", state: "failed" as const, outcome: "closed_failed" as const }; + } + }; + const server = createApp({ reviewService: review, reviewRecoveryService: recovery, 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/version-2`; + + const unsafe = await fetch(`${base}/review`, { headers: { authorization: "Bearer review-token" } }); + assert.equal(unsafe.status, 500); + assert.deepEqual(await unsafe.json(), { ok: false, error: "Unexpected OCR review failure", code: "OCR_REVIEW_UNEXPECTED", action: "retry_or_contact_support" }); + const unauthorized = await fetch(`${base}/recover`, { method: "POST" }); + assert.equal(unauthorized.status, 401); + const recovered = await fetch(`${base}/recover`, { + method: "POST", headers: { authorization: "Bearer review-token", "content-type": "application/json" }, + body: JSON.stringify({ recoveredBy: "admin", reason: "Missing durable candidate" }) + }); + assert.equal(recovered.status, 200); + assert.deepEqual(await recovered.json(), { versionId: "version-2", state: "failed", outcome: "closed_failed" }); + assert.deepEqual(recoveries, [{ recoveredBy: "admin", reason: "Missing durable candidate" }]); + assert.deepEqual(classifyOcrReviewError(new OcrArtifactError("missing", 409, "OCR_ARTIFACT_UNAVAILABLE")), + new CatalogError("missing", 409, "OCR_ARTIFACT_UNAVAILABLE")); +}); + test("playground serves the authenticated OCR review controls and audit fields", async (context) => { const server = createApp({ startReconciler: false }).listen(0); context.after(() => server.close()); @@ -236,7 +274,7 @@ test("production review reader survives restart and fails closed on unauthorized 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); + assert.equal((await fetch(`${base}/review`, { headers: { authorization: "Bearer review-token" } })).status, 422); 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 });