From 640d3a4a1a2234b899885c130b4d9b6d92e2c95f Mon Sep 17 00:00:00 2001 From: Paco POR-CORREO Date: Mon, 21 Sep 2026 19:20:39 +0200 Subject: [PATCH] feat(ocr): release durable processing safeguards --- Dockerfile | 9 +- docs/API_RAG.md | 8 +- docs/CONTRATO_CICLO_VIDA_Y_OCR.md | 155 ++++++++++- docs/HISTORIAL_SESIONES.md | 101 ++++++- docs/OPERATIVA.md | 26 +- docs/PENDIENTES_RAG.md | 47 +++- ocr-service/Dockerfile | 6 +- ocr-service/README.md | 21 +- ocr-service/app/jobs.py | 419 +++++++++++++++++++++++++++++ ocr-service/app/main.py | 169 +++--------- ocr-service/app/render.py | 58 ++-- ocr-service/tests/test_api.py | 224 +++++++++++++-- ocr-service/tests/test_render.py | 32 +++ package-lock.json | 4 +- package.json | 2 +- src/api/openapi.ts | 5 +- src/app.ts | 12 +- src/config/env.ts | 3 +- src/modules/ocr/artifacts.ts | 53 +++- src/modules/ocr/client.ts | 33 ++- tests/ocr/client.test.ts | 62 ++++- tests/ocr/contracts-deploy.test.ts | 4 + 22 files changed, 1216 insertions(+), 237 deletions(-) create mode 100644 ocr-service/app/jobs.py diff --git a/Dockerfile b/Dockerfile index 963dee2..f896442 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,5 +1,5 @@ FROM node:22-bookworm-slim AS build -ARG RAG_VERSION=0.1.0 +ARG RAG_VERSION=0.2.0 ARG BUILD_REVISION=unknown WORKDIR /app COPY package.json package-lock.json tsconfig.json ./ @@ -10,11 +10,12 @@ COPY public ./public RUN npm run build FROM node:22-bookworm-slim AS runtime -ARG RAG_VERSION=0.1.0 +ARG RAG_VERSION=0.2.0 ARG BUILD_REVISION=unknown WORKDIR /app -ENV NODE_ENV=production OCR_ARTIFACT_ROOT=/data/ingestions RAG_VERSION=${RAG_VERSION} -LABEL org.opencontainers.image.revision=${BUILD_REVISION} +ENV NODE_ENV=production OCR_ARTIFACT_ROOT=/data/ingestions RAG_VERSION=${RAG_VERSION} BUILD_REVISION=${BUILD_REVISION} +LABEL org.opencontainers.image.version=${RAG_VERSION} \ + 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/API_RAG.md b/docs/API_RAG.md index ff1bc68..6e12d49 100644 --- a/docs/API_RAG.md +++ b/docs/API_RAG.md @@ -2,9 +2,9 @@ **Proyecto:** Workspace de tools IA para empresas **Modulo:** RAG -**Ultima actualizacion:** 2026-09-13 -**Ultima modificacion por:** Subagente Implementacion ciclo de vida del conocimiento -**Estado:** Punto 2 desplegado, migrado y validado en produccion +**Ultima actualizacion:** 2026-09-21 +**Ultima modificacion por:** Agente RAG 2 +**Estado:** Produccion en `0.1.0`; version `0.2.0` preparada localmente Esta guia explica como descubrir y consumir la API. El contrato tecnico completo y canonico se publica en OpenAPI 3.1.1; este documento prioriza el camino rapido, las decisiones de uso y los ejemplos habituales. @@ -52,7 +52,7 @@ La API queda formada por las operaciones de descubrimiento, ingesta, retrieval, | `GET` | `/help` | Catalogo resumido derivado de OpenAPI | | `GET` | `/openapi.json` | Contrato OpenAPI 3.1.1 | | `GET` | `/playground` | Interfaz web de prueba | -| `GET` | `/health` | Estado, proveedores, parsers y chunking | +| `GET` | `/health` | Estado, version, commit, proveedores, parsers y chunking | | `GET` | `/sources` | Scopes disponibles en Qdrant | | `GET` | `/sources/{sourceId}` | Fuente registrada en el catalogo | | `GET` | `/sources/{sourceId}/versions` | Versiones de una fuente | diff --git a/docs/CONTRATO_CICLO_VIDA_Y_OCR.md b/docs/CONTRATO_CICLO_VIDA_Y_OCR.md index 249d994..0bed216 100644 --- a/docs/CONTRATO_CICLO_VIDA_Y_OCR.md +++ b/docs/CONTRATO_CICLO_VIDA_Y_OCR.md @@ -2,9 +2,9 @@ **Proyecto:** Workspace de tools IA para empresas **Modulo:** RAG -**Ultima actualizacion:** 2026-09-08 +**Ultima actualizacion:** 2026-09-21 **Ultima modificacion por:** Agente RAG 2 -**Estado:** Punto 2 implementado, migrado y validado en produccion; punto 3 OCR pendiente +**Estado:** Punto 2 validado; pipeline OCR desplegado; aceptacion productiva bloqueada hasta completar el hardening de runtime definido en este contrato ## Resultado esperado @@ -998,7 +998,7 @@ El RAG usa un volumen persistente privado montado en `/data/ingestions`. Estruct candidate-pages.json.gz reviewed-pages.json.gz review-images/ - page-0001.webp + page-0001.png ``` `documentArtifactId` es UUIDv5 del `document_id`, no una ruta ni un nombre aportado por el usuario. El manifiesto relaciona ese ID con `document_id` y evita colisiones en fuentes con varios PDFs. @@ -1144,6 +1144,155 @@ El punto 3 solo se marca completado cuando: - se valida limpieza de temporales y TTL; - se actualizan backlog, API, despliegue e historiales. +# ODD: hardening del runtime OCR antes de la aceptacion productiva + +## Resultado requerido + +La aceptacion de FacturaTech solo puede retomarse cuando OCR procese y sirva las evidencias de revision sin volver a renderizar el PDF bajo carga concurrente. Esta mejora se ejecuta como un unico flujo ODD y sustituye la ampliacion del SDD `ocr-ingest-integration`, que se archiva con su tarea productiva 7.4 incompleta. + +El flujo debe entregar: + +- una cola SQLite durable atendida por un unico worker real; +- recuperacion acotada de trabajos tras reiniciar el contenedor; +- PNG de revision persistidos durante el render inicial de cada pagina; +- lectura de imagenes desde fichero, sin nuevas llamadas a PDFium; +- transferencia RAG secuencial, idempotente y reanudable; +- exclusion mutua defensiva para cualquier operacion residual con PDFium; +- limpieza que no elimine trabajos activos ni evidencias aun no transferidas. + +## Incidente que obliga al hardening + +La candidata FacturaTech v5 `5f2317c6-7a8a-4e08-a614-f8189602ebb8` completo y persistio el OCR de 25 paginas. Despues, RAG solicito las 25 imagenes de revision concurrentemente mediante `Promise.all`. El servicio volvio a renderizar el PDF para cada peticion y PDFium/FreeType termino con `SIGSEGV` en `FPDF_RenderPageBitmap -> FT_Load_Glyph`. + +El contenedor salio con codigo `139`; `OOMKilled=false`. RAG aplico el contrato fail-closed, marco la candidata como fallida y mantuvo intacta la version activa. No fue un fallo de memoria, red, autenticacion ni procesamiento PaddleOCR. + +## Invariantes + +| Area | Contrato obligatorio | +|---|---| +| Admision | La API persiste el trabajo en SQLite antes de responder `202`; no inicia un `BackgroundTask` independiente por solicitud. | +| Concurrencia | Una instancia OCR ejecuta exactamente un worker de procesamiento y como maximo un trabajo `running`. | +| Reinicio | Al arrancar, el worker recupera trabajos `queued` y reclasifica de forma segura los `running` interrumpidos; nunca duplica un resultado `completed`. | +| Render | Cada pagina se renderiza una sola vez por intento de procesamiento. El mismo render alimenta OCR y publica el PNG privado de revision. | +| Imagenes | El endpoint de imagen valida identidad y contencion, y transmite un PNG persistido. No abre ni renderiza el PDF. | +| Transferencia | RAG descarga una pagina cada vez. Una transferencia interrumpida continua desde los artefactos ya publicados y verificados por hash. | +| PDFium | Toda operacion que aun use PDFium queda bajo un lock de proceso e incluye cierre determinista de pagina, documento y bitmap. | +| Limpieza | Un trabajo con lease vigente no se elimina. Un `running` no puede permanecer activo mas de 15 minutos: al expirar su lease se recupera una vez o pasa a `failed`. Tras la transferencia durable se elimina inmediatamente; sin confirmacion, el limite absoluto es 24 horas desde `created_at`, mas como maximo la ejecucion vigente de 15 minutos. | +| Errores | RAG no reintenta HTTP `500`. Solo conserva los reintentos idempotentes ya definidos para conexion, `502` y `503`. | +| Activacion | El hardening, sus pruebas y su despliegue no aprueban, indexan ni activan contenido. La version activa anterior permanece intacta. | + +## Ubicacion y limpieza exacta de los ficheros + +OCR y RAG no comparten filesystem. Cada servicio conserva una copia con una finalidad y una retencion diferentes. + +### Copia transitoria del servicio OCR + +El volumen OCR `/data/jobs` usa esta estructura: + +```text +/data/jobs/ + jobs.db + artifacts/ + / + input.pdf + review-images/ + page-0001.png + page-0002.png +``` + +- `jobs.db` guarda solo metadatos, hashes, estados, intentos, heartbeat, lease y rutas relativas. No guarda el PDF ni los PNG como BLOB. +- Cada fichero se publica con escritura temporal, `fsync` y rename atomico dentro del directorio del trabajo. +- RAG confirma la transferencia solo despues de publicar y releer su copia durable completa. Entonces solicita `DELETE /v1/jobs/:job_id` y OCR elimina recursivamente `/data/jobs/artifacts/` y su fila SQLite. +- El sweeper se ejecuta al arrancar y cada 15 minutos. Un `running` con lease vencido se devuelve a `queued` una sola vez; si vuelve a interrumpirse o supera 15 minutos totales, pasa a `failed`. +- Ningun trabajo queda protegido indefinidamente por llamarse activo. Los trabajos sin confirmacion expiran a las 24 horas desde `created_at`; si en ese instante existe una ejecucion con lease vigente, primero se deja terminar o alcanzar su limite de 15 minutos y despues se purga. +- La purga elimina el directorio completo y la fila SQLite. SQLite se configura con `auto_vacuum=INCREMENTAL` y el mantenimiento ejecuta checkpoint e `incremental_vacuum` para que metadatos eliminados tampoco acumulen espacio muerto. +- Antes de admitir un trabajo y antes de publicar cada PNG, OCR comprueba uso del volumen y espacio libre. El presupuesto transitorio es el menor entre el 10 % de la capacidad del filesystem y 2 GiB; siempre reserva el mayor entre el 10 % y 2 GiB libres. Si no puede respetarlo, rechaza o falla de forma cerrada con `OCR_STORAGE_PRESSURE` y limpia el intento parcial. + +Con cola maxima 3, timeout total de 15 minutos, recuperacion unica, limite de volumen y sweeper periodico, no existe un camino normal que deje PDFs o PNG flotando indefinidamente. + +### Copia durable del RAG + +RAG conserva las imagenes privadas en: + +```text +/data/ingestions//documents//review-images/page-NNNN.png +``` + +- La copia se publica con permisos `0600`, manifiesto y SHA-256 por fichero. +- Mientras la version espera revision, sus evidencias se conservan hasta 30 dias. +- Despues de aprobar o rechazar, las imagenes de revision se conservan 7 dias y despues se eliminan mediante el limpiador de retencion. +- El resto de artefactos de una version activa se conserva mientras siga activa y 30 dias despues de ser reemplazada. +- Una version activa nunca se elimina por TTL; las imagenes no se guardan en PostgreSQL ni Qdrant y no tienen URL publica. + +## Unidades de desarrollo + +Cada unidad se implementa en una rama de feature y termina con un commit de unidad de trabajo Conventional Commit. Pruebas y documentacion propias de la conducta viajan en el mismo commit; la validacion independiente se ejecuta despues y no se mezcla con la implementacion. + +### D1. Cola durable y worker unico + +**Estado:** Implementada y validada de forma independiente el 2026-09-21. + +- Convertir SQLite en la fuente de verdad de trabajos `queued`, `running`, `completed` y `failed`. +- Mantener en SQLite solo metadatos y rutas; persistir PDF y PNG bajo `/data/jobs/artifacts/` para que su borrado libere espacio real. +- Sustituir el lanzamiento por peticion por un loop de worker unico con claim atomico. +- Recuperar trabajos interrumpidos al arrancar, con intento acotado y transiciones auditables. +- Hacer que readiness distinga entre modelo cargado, worker operativo y almacenamiento disponible. + +### D2. Artefactos de revision sin rerender + +**Estado:** Implementada y validada de forma independiente el 2026-09-21. + +- Persistir cada PNG privado durante el render inicial usado por OCR. +- Publicar resultado, manifiesto e imagenes de manera atomica o dejar el trabajo recuperable. +- Servir unicamente ficheros persistidos y rechazar rutas, paginas o hashes no vinculados al trabajo. +- Mantener un lock defensivo alrededor de las operaciones PDFium y liberar recursos en todos los caminos. + +### D3. Transferencia RAG secuencial y reanudable + +**Estado:** Implementada y validada de forma independiente el 2026-09-21. + +- Eliminar la descarga concurrente de imagenes mediante `Promise.all`. +- Transferir paginas en orden, una por una, validando identidad, tipo, tamano y hash. +- Conservar artefactos ya transferidos tras reinicio y reanudar solo los que falten. +- Confirmar limpieza remota unicamente despues de publicar durablemente toda la candidata local. + +### D4. Observabilidad y limites operativos + +**Estado:** Implementada y validada de forma independiente el 2026-09-21. + +- Exponer estados reales de cola, worker y recuperaciones sin filtrar secretos. +- Registrar transiciones, recuperaciones y fallos terminales con `job_id` y categoria segura. +- Mantener un trabajo Paddle concurrente, cola maxima 3 y limites actuales de CPU, RAM y tiempo. +- Aplicar presupuesto de disco, reserva de espacio libre, sweeper cada 15 minutos y vacuum incremental de SQLite. +- Conservar el fallo cerrado y la idempotencia entre RAG y OCR. + +## Validacion independiente + +**Estado:** Superada el 2026-09-21. La primera comprobacion incompleta fue rechazada; la repeticion valida proceso un PDF real de 25 paginas con PaddleOCR, interrumpio y recupero el trabajo sin duplicados y midio la limpieza en completado, fallo y expiracion. + +La validacion comienza solo cuando D1-D4 estan implementadas. Si descubre un defecto, se abre una unidad correctiva separada y se vuelve a ejecutar la validacion afectada; no se corrige silenciosamente durante la comprobacion. + +1. Ejecutar las pruebas focalizadas de cola, reinicio, render, imagenes, transferencia y limpieza. +2. Ejecutar las suites completas Node y Python, `npm run check` y `npm run build`. +3. Procesar localmente un PDF de 25 paginas, reiniciar OCR durante un trabajo y comprobar recuperacion sin duplicados. +4. Solicitar imagenes de revision concurrentemente y demostrar que el endpoint solo lee PNG persistidos y no llama a PDFium. +5. Interrumpir la transferencia RAG y demostrar que continua desde la primera pagina ausente sin redescargar las verificadas. +6. Verificar que un HTTP `500` es terminal para RAG y que conexion, `502` y `503` mantienen el reintento idempotente acotado. +7. Confirmar que los temporales se eliminan solo tras transferencia completa o TTL y nunca mientras el trabajo esta activo. +8. Medir el espacio de `/data/jobs` antes y despues de completar, fallar y expirar trabajos; no deben quedar directorios huerfanos y el espacio debe volver a estar disponible. + +## Despliegue y aceptacion de FacturaTech + +1. Desplegar juntos OCR y RAG `0.2.0` con OCR habilitado desde el inicio. +2. Verificar live, readiness, versiones, commit desplegado, worker unico, volumen SQLite y salud general del RAG. +3. Crear una nueva candidata FacturaTech no activa. Si aparece un fallo cuyo origen no se pueda identificar, usar temporalmente la activacion separada de OCR y RAG solo como procedimiento de diagnostico. +4. Verificar evidencia durable, 25 imagenes de revision y las 34 entradas esperadas. +5. Comprobar exactamente `CBG04a`, `FAT07`, `DSAU08` y `NSAV06`. +6. Presentar la evidencia al usuario y esperar autorizacion explicita. +7. Solo tras esa autorizacion, aprobar, indexar y, si tambien se autoriza, activar. + +La mejora queda cerrada cuando la validacion independiente pasa, FacturaTech completa la revision productiva y la activacion autorizada conserva rollback verificable. Una candidata fallida o una validacion parcial no permite declarar completada la antigua tarea SDD 7.4. + # Verificacion y entrega comun Cada punto se entrega como bloque independiente: diff --git a/docs/HISTORIAL_SESIONES.md b/docs/HISTORIAL_SESIONES.md index e6a1db4..8c28731 100644 --- a/docs/HISTORIAL_SESIONES.md +++ b/docs/HISTORIAL_SESIONES.md @@ -2,7 +2,7 @@ **Proyecto:** Workspace de tools IA para empresas **Modulo:** RAG -**Ultima actualizacion:** 2026-09-17 +**Ultima actualizacion:** 2026-09-21 **Ultima modificacion por:** Agente RAG 2 **Estado:** Activo @@ -10,6 +10,105 @@ ## Registro de sesion +### 2026-09-21 - Agente RAG 2 - Preparacion de la version 0.2.0 +**Agent:** Agente RAG 2 · **Model:** openai/gpt-5.6-sol · **Session:** `ses_29bdbd003ffeLrLjUlFgnp08Y7` +**Work:** Preparada la entrega conjunta RAG/OCR `0.2.0` tras superar D1-D4 la validacion independiente. Ambos servicios muestran version y revision de commit en sus respuestas de salud y etiquetas de imagen. Actualizados el README privado del OCR, la ayuda del RAG, el contrato, la operativa y el backlog. Por decision del usuario, el proximo despliegue instala RAG y OCR juntos con OCR activado; el despliegue separado queda solo como diagnostico si aparece un fallo ambiguo. Produccion continua en `0.1.0` hasta el despliegue manual en EasyPanel. +**Validation:** Node 103/103; Python OCR 24/24; `npm run check`; `npm run build`; `py_compile`; `git diff --check`. La validacion independiente previa proceso 25 paginas con PaddleOCR real y recupero correctamente un reinicio sin duplicados. +**Rollback:** Revertir el commit de hardening y version `0.2.0`; no se ha desplegado ni modificado produccion. +**Files:** codigo, pruebas, Dockerfiles y documentacion de RAG/OCR incluidos en la entrega `0.2.0`. + +--- + +### 2026-09-21 - Subagente Validación Independiente D1-D4 - Veredicto global PASS del hardening OCR (revisado) + +> **Registro de revision honesta (2026-09-21, misma sesion):** La primera version de esta entrada emitio un PASS global que fue **rechazado por el maintainer por incumplir `docs/CONTRATO_CICLO_VIDA_Y_OCR.md:1275`** (criterio 3: procesar localmente un PDF de 25 paginas con reinicio durante el trabajo) y por falta de medicion cuantitativa del criterio 8. La primera validacion uso un fixture de 3 paginas procesado tres veces y un motor fake para el escenario de reinicio, lo cual es evidencia de componente, no cumplimiento literal del criterio. El veredicto fue retirado y la validacion rehecha con runtime real; el texto siguiente refleja la evidencia completa y actual. No se oculta el error original. + +**Agent:** Subagente Validación Independiente D1-D4 · **Model:** openai/glm-5.3 · **Session:** `ses_f3b5e072fffehGGsxlS7r7F1cC` (parent `ses_29bdbd003ffeLrLjUlFgnp08Y7`) +**Role:** Validacion independiente de D1-D4 del hardening OCR; sin implementar correcciones, deploy, commit, push ni acciones sobre FacturaTech. +**Work:** Veredicto global **PASS** tras revalidacion con runtime real. Revision del contrato ODD (`docs/CONTRATO_CICLO_VIDA_Y_OCR.md` lineas 1169-1280) y de los criterios 1-8 (lineas 1273-1279). Se instalaron `paddleocr==3.4.0` y `paddlepaddle==3.2.2` en `ocr-service/.venv` (paso previo necesario: el venv local no los tenia; el motor real PaddleOCR carga en ~9 s con los modelos del contrato `PP-OCRv5_mobile_det`/`latin_PP-OCRv5_mobile_rec`). **Criterio 3 con runtime real completo:** PDF real de exactamente 25 paginas generado fuera del repo (`/tmp/opencode/d1-harness/real25/factura25.pdf`, 28.623 B, sha256 `d4cf427d90c499a48f727491e506c23f162a3bc5a1fe10f3fdc7ac0b41c5faca`) con 39 lineas de texto por pagina incluyendo `SQLSTATE 23505`, `FAT07` y `FE666-xx`. Ruta real de trabajo/worker: `JobQueue` + worker thread + `PaddleOcrEngine` real; render real a 200 DPI (1654x2339, ~871 KB PNG/pagina). Procesamiento iniciado con motor real (~24 s/pagina), `SIGKILL` al proceso en la pagina 5 (4 PNG persistidos de 25; job `running` con lease vigente). Reinicio con motor real: `recover_interrupted` reclasifico `running -> queued` con `recovery_attempts=1` (recuperacion unica) y el worker completo las 25 paginas: `succeeded` en 6:44 min tras el reinicio. **Sin duplicados verificado cuantitativamente:** 1 fila SQLite, 0 claves idempotentes duplicadas, 25 paginas en el resultado, 25 unicas y ordenadas, 25 PNGs en disco con 0 hashes duplicados, `input.pdf` byte-identico al original, `FAT07` reconocido en las 25 paginas (975 lineas OCR reales), bloque `engine` del resultado = `paddleocr|3.4.0|paddlepaddle-3.2.2|cpu|ocr-v1|200`. `attempts=2` es correcto: el invariante de render es "una vez por intento de procesamiento", y el crash anulo el primer intento. **Criterio 8 con medicion cuantitativa en almacenamiento aislado (`/tmp/opencode/d1-harness/real25/storage-c8`):** camino completado: submit 28.623 B/1f -> procesado 21.920.524 B/26f -> DELETE confirmado 0 B/0f, fila SQLite 0; camino fallado por presion: `OCR_STORAGE_PRESSURE` tras publicar 2 paginas, `review-images/` parcial eliminado (input.pdf 28.623 B conservado, fila `failed` conservada); camino expirado por TTL (+25 h): sweep purga fila y directorio, 0 B/0f, 0 filas, 0 directorios huerfanos. Espacio de artefactos vuelve exactamente al baseline en los tres caminos. Metadatos SQLite: verificado con resultados de tamaño real (288 KiB/job de 25 paginas) durante 8 ciclos de purga: crecimiento de main y WAL = 0 bytes (high-water mark reutilizado, sin acumulacion de espacio muerto, conforme a linea 1208). Hallazgo menor (no bloqueante): con payloads artificiales de ~1 MB (37x el maximo real observado), la cadena `wal_checkpoint(PASSIVE)` + `incremental_vacuum` del sweeper deja un high-water mark retenido y el WAL crece hasta ~20 MB durante el sweep; un `wal_checkpoint(TRUNCATE)` tras el vacuum lo liberaria por completo. No es acumulacion (se reutiliza), pero se recomienda como mejora. Resto de evidencia D1-D4 (crash SIGKILL exit-137 en worker con recuperacion unica terminal a la segunda interrupcion, cola maxima 3, concurrencia 1, 25 peticiones de imagen concurrentes sin rerender con PDFium sabotado, exclusion mutua PDFium pico=1, transferencia secuencial con reanudacion desde checkpoints sin redescargar paginas verificadas, DELETE fail-closed `409 JOB_ACTIVE`, readiness sin secretos, politica de reintentos 500-terminal/502-503-red x3) permanece valida de la primera ronda y fue reconfirmada donde aplicaba. +**Validation:** Python OCR 24/24 (`ocr-service/.venv`, ahora con Paddle real instalado), Node 103/103, `npm run check` limpio, `npm run build` limpio, `git diff --check` limpio. Comparacion `git status --short` antes/despues de la revalidacion: sin cambios nuevos en el repo; toda la evidencia vive bajo `/tmp/opencode/d1-harness/real25/`. La instalacion de Paddle en `.venv` es un cambio de entorno local no rastreado por git (`.venv` esta en `.gitignore`) y no altera el codigo del repo; `requirements.txt` ya fijaba las versiones exactas. Sin correcciones aplicadas por protocolo de validacion. +**Files:** `docs/HISTORIAL_SESIONES.md` (esta entrada, corregida). Ningun otro archivo modificado. + +--- + +### 2026-09-21 - Agente RAG 2 - D4 observabilidad y limites operativos OCR +**Agent:** Agente RAG 2 · **Model:** openai/gpt-5.6-sol · **Session:** `ses_29bdbd003ffeLrLjUlFgnp08Y7` +**Work:** Cerrado el hardening D4 con estados reales de cola, worker, sweeper y recuperaciones en readiness, y logs de transicion limitados a `job_id` y categorias seguras. La admision y cada PNG comprueban el presupuesto transitorio y la reserva libre; la presion devuelve `OCR_STORAGE_PRESSURE` y elimina imagenes parciales. El mantenimiento arranca con el servicio, se ejecuta como maximo cada 15 minutos y tambien respeta el vencimiento mas cercano de lease/timeout; recupera una vez, falla cerrado al agotar recuperacion o 15 minutos, protege trabajos activos frente a DELETE, purga el TTL de 24 horas y ejecuta checkpoint WAL y vacuum incremental. Se conservan cola 3, concurrencia 1, un worker Uvicorn y limites documentados de 3 CPU/5 GiB. +**Validation:** Focalizadas API D4 18/18; Python OCR 24/24; Node 103/103; `npm run check`; `npm run build`; `py_compile`; `git diff --check`. El harness TestClient ejercito lifecycle real, trabajo OCR, health preparado, presion en admision/publicacion, limpieza parcial, DELETE activo, recuperacion por lease, timeout terminal, sweeper periodico, TTL y modo incremental-vacuum. Una primera ejecucion revelo y corrigio la anotacion estrecha de respuesta del healthcheck; una segunda detecto y corrigio la expectativa `lastSweepAt` para clientes sin lifespan. La revision final cerro una carrera entre claim y espera del sweeper, y anadio regresion para rechazar en arranque un trabajo cuyo limite total ya expiro. Sin validacion independiente, deploy, aprobacion, indexacion, activacion ni reingesta de FacturaTech. +**Rollback:** Revertir solo la conducta D4 en `ocr-service/app/jobs.py`, `ocr-service/app/main.py` y sus pruebas/documentacion; D1-D3 permanecen funcionalmente separadas. Sin commit ni push. +**Files:** `ocr-service/app/jobs.py`, `ocr-service/app/main.py`, `ocr-service/tests/test_api.py`, `docs/CONTRATO_CICLO_VIDA_Y_OCR.md`, `docs/PENDIENTES_RAG.md`, `docs/OPERATIVA.md`, `docs/HISTORIAL_SESIONES.md`. + +--- + +### 2026-09-21 - Agente RAG 2 - D3 transferencia OCR secuencial y reanudable +**Agent:** Agente RAG 2 · **Model:** openai/gpt-5.6-sol · **Session:** `ses_29bdbd003ffeLrLjUlFgnp08Y7` +**Work:** Sustituida la descarga concurrente con `Promise.all` por transferencia estrictamente secuencial de las paginas OCR solicitadas. Cada PNG validado se publica de forma inmutable junto con un comprobante privado ligado a version, documento, pagina, ruta y SHA-256; tras una interrupcion, RAG verifica y reutiliza paginas completas y reanuda desde la primera ausente. El cliente valida tipo, `Content-Length`, identidad, pagina, firma PNG y hash; conexion, `502` y `503` tienen reintento acotado, mientras `500` es terminal. La limpieza remota conserva su orden posterior a resultado, imagenes, candidata y estado local durables. +**Validation:** Pruebas focalizadas 26/26; Node 103/103; Python OCR 21/21; `npm run check`; `npm run build`; `git diff --check`. Incluye interrupcion en pagina 2, reanudacion sin redescargar pagina 1, ejecucion maxima de una descarga y rechazo fail-closed de checkpoint alterado. Sin validacion independiente, deploy, aprobacion, indexacion, activacion ni reingesta de FacturaTech. +**Files:** `src/app.ts`, `src/modules/ocr/artifacts.ts`, `src/modules/ocr/client.ts`, `tests/ocr/client.test.ts`, `docs/CONTRATO_CICLO_VIDA_Y_OCR.md`, `docs/PENDIENTES_RAG.md`, `docs/HISTORIAL_SESIONES.md`. Sin commit ni push. + +--- + +### 2026-09-21 - Agente RAG 2 - D2 artefactos OCR sin rerender +**Agent:** Agente RAG 2 · **Model:** openai/gpt-5.6-terra · **Session:** `ses_29bdbd003ffeLrLjUlFgnp08Y7` +**Work:** El worker OCR publica cada PNG privado mediante fichero temporal, `fsync` y rename atómico en `artifacts//review-images/` durante el mismo render que usa PaddleOCR. El endpoint de imágenes solo lee ese artefacto asociado a una página solicitada y un trabajo terminado; ausencia o trabajo incompleto falla cerrado. Añadido `RLock` de proceso alrededor de PDFium y preservado el cierre determinista de documento, página y bitmap. +**Validation:** Python OCR 21/21; `npm test` 101/101; `npm run check`; `npm run build`; `git diff --check`. Sin validación independiente, deploy, aprobación, indexación ni activación de contenido. +**Files:** `ocr-service/app/render.py`, `ocr-service/app/jobs.py`, `ocr-service/app/main.py`, `ocr-service/tests/test_api.py`, `ocr-service/tests/test_render.py`, `docs/CONTRATO_CICLO_VIDA_Y_OCR.md`, `docs/PENDIENTES_RAG.md`, `docs/HISTORIAL_SESIONES.md`. Sin commit ni push. + +--- + +### 2026-09-21 - Agente RAG 2 - D1 cola durable OCR +**Agent:** Agente RAG 2 · **Model:** openai/gpt-5.6-terra · **Session:** `ses_29bdbd003ffeLrLjUlFgnp08Y7` +**Work:** Sustituido el PDF BLOB y `BackgroundTasks` por una cola SQLite durable que guarda el PDF privado en `artifacts//input.pdf`. El servicio inicia un unico worker mediante el ciclo de vida de FastAPI, reclama trabajos atomica y secuencialmente, conserva intentos y lease, y recupera una interrupcion al arrancar antes de cerrar el segundo fallo como terminal. Readiness ahora diferencia modelo, worker y almacenamiento. +**Validation:** Python OCR 18/18; `npm test` 101/101; `npm run check`; `npm run build`; `git diff --check`. Sin validacion independiente, deploy, aprobacion, indexacion ni activacion de contenido. +**Files:** `ocr-service/app/jobs.py`, `ocr-service/app/main.py`, `ocr-service/tests/test_api.py`, `docs/CONTRATO_CICLO_VIDA_Y_OCR.md`, `docs/PENDIENTES_RAG.md`, `docs/HISTORIAL_SESIONES.md`. Sin commit ni push. + +--- + +### 2026-09-21 - Agente RAG 2 - Retencion y presupuesto de disco para imagenes OCR +**Agent:** Agente RAG 2 · **Model:** openai/gpt-5.6-sol · **Session:** `ses_29bdbd003ffeLrLjUlFgnp08Y7` +**Work:** Aclarado el ciclo de vida exacto de las imagenes del hardening antes de entregar el plan a otro modelo. Corregida la extension historica `.webp` a `.png`. Definida la copia OCR transitoria bajo `/data/jobs/artifacts/` y la copia RAG durable bajo `/data/ingestions//documents//review-images`. Para evitar espacio no recuperable, el diseño elimina PDF/PNG como BLOB de SQLite: `jobs.db` conserva solo metadatos y rutas, y la purga borra directorio y fila con vacuum incremental. Especificados lease/timeout, una unica recuperacion, sweeper al arranque y cada 15 minutos, TTL absoluto de 24 horas, presupuesto del menor entre 10 % del filesystem y 2 GiB, reserva libre del mayor entre 10 % y 2 GiB, y fallo cerrado `OCR_STORAGE_PRESSURE`. +**Validation:** Contrato releido contra `ocr-service/app/main.py`, `src/modules/ocr/artifacts.ts`, retencion RAG y specs canonicas. El plan ya distingue trabajos realmente activos de trabajos colgados y exige medir recuperacion de espacio en la validacion independiente. No se modifico codigo ni produccion. +**Learned:** Guardar PDFs grandes como BLOB y borrar filas no garantiza devolver espacio al sistema operativo; los binarios transitorios deben vivir en directorios eliminables y SQLite limitarse a metadatos. +**Files:** `docs/CONTRATO_CICLO_VIDA_Y_OCR.md`, `docs/HISTORIAL_SESIONES.md`. Sin commit ni push. + +--- + +### 2026-09-21 - Agente RAG 2 - Transicion del cierre OCR de SDD a ODD +**Agent:** Agente RAG 2 · **Model:** openai/gpt-5.6-sol · **Session:** `ses_29bdbd003ffeLrLjUlFgnp08Y7` +**Work:** Verificado mediante estado nativo que `ocr-ingest-integration` tenia 29/30 tareas completadas y solo 7.4 pendiente. Consolidado el hardening correctivo en la fuente canonica `CONTRATO_CICLO_VIDA_Y_OCR.md` como flujo ODD: cola SQLite durable, worker unico, recuperacion tras reinicio, PNG de revision persistidos sin rerender, transferencia RAG secuencial y reanudable, lock defensivo de PDFium, limpieza segura, unidades de desarrollo separadas de la validacion y aceptacion productiva sin activacion automatica. Actualizada la Fase 3 del backlog con la causa raiz del fallo de FacturaTech v5 y la nueva secuencia 3A/3B. Tras preflight e inicializacion SDD, el agente especializado sincronizo las tres specs canonicas y archivo el cambio preservando 7.4 incompleta y verify-report ausente. +**Validation:** Contrato y backlog releidos; `git diff --check` sin errores; cambio activo ausente; archivo y tres specs canonicas presentes; `tasks.md` archivado conserva 29 tareas marcadas y 7.4 sin marcar. La inicializacion requerida ejecuto `npm test` 101/101, `npm run check` exit 0 y pytest OCR 16/16. No se ejecuto aceptacion productiva ni se aprobo, indexo o activo contenido. +**Learned:** El plugin SDD activo reconoce un preflight canonico de tres grupos, mientras la documentacion instalada describe cuatro; la primera confirmacion de cuatro grupos no genero autoridad y fue necesario usar los marcadores `Gentle AI SDD preflight N/3`. Esta incompatibilidad debe corregirse en Gentle AI. El repo contiene dos proyectos con runners independientes y sin comando workspace-level comun, por lo que `strict_tdd` queda `false` de forma fail-closed. +**Files:** `docs/CONTRATO_CICLO_VIDA_Y_OCR.md`, `docs/PENDIENTES_RAG.md`, `docs/HISTORIAL_SESIONES.md`, `openspec/config.yaml`, `openspec/specs/{ocr-ingest-orchestration,ocr-processing,ocr-review-workflow}/spec.md`, `openspec/changes/archive/2026-09-21-ocr-ingest-integration/`, `/home/pancho/Documentos/Empresa/IA/herramientas/docs/gentle-ai/SEGUIMIENTO_DESCUBRIMIENTOS_MEJORAS_GENTLE_AI.md`. Sin commit ni push. + +--- + +### 2026-09-21 - Subagente Archivo SDD OCR - Archivo honesto de ocr-ingest-integration +**Agent:** Subagente Archivo SDD OCR · **Model:** ollama/glm-5.3:cloud · **Session:** `ses_f3c4dec1dffeJES2Q83hmNsMYB` (subagent of `ses_29bdbd003ffeLrLjUlFgnp08Y7`) +**Responsibility:** Archivo honesto del cambio SDD `ocr-ingest-integration` en modo hybrid: sincronizacion mecanica de las tres delta specs a specs canonicas, movimiento del cambio a archivo con snapshot previo y diff -r vacio obligatorio, y reporte de archivo en Engram. Sin implementar codigo, sin verificacion, sin alterar checkboxes. +**Work:** Estado nativo refrescado sin bloqueos. Creadas las tres specs canonicas nuevas (`openspec/specs/ocr-ingest-orchestration/spec.md`, `openspec/specs/ocr-processing/spec.md`, `openspec/specs/ocr-review-workflow/spec.md`) mediante copia mecanica con shell (cp a mktemp + mv, permisos 0644 alineados al origen) y diff -r vacio frente a cada delta. Movido el cambio completo con `git mv` a `openspec/changes/archive/2026-09-21-ocr-ingest-integration/` con snapshot recursivo previo y diff -r vacio (snapshot vs destino), sin colisiones. Estado preservado: 29/30 tareas completadas, solo 7.4 (aceptacion en produccion) pendiente; verify-report ausente (verificacion no ejecutada). Hecho final registrado: FacturaTech v5 `5f2317c6` completo/persistio OCR de 25 paginas pero fallo con SIGSEGV (`FPDF_RenderPageBitmap -> FT_Load_Glyph`, exit 139, `OOMKilled=false`) al solicitar RAG 25 imagenes de revision concurrentemente; candidata cerro en fallo, version activa intacta, sin aprobacion, indexacion ni activacion. El hardening de runtime y la aceptacion en produccion se movieron a la seccion ODD canonica de `docs/CONTRATO_CICLO_VIDA_Y_OCR.md` y NO estan implementados. Reporte de archivo persistido en Engram `rag-service` bajo `sdd/ocr-ingest-integration/archive-report` (obs #3691, capture_prompt false); 3 veredictos de revision de conflictos registrados (related/compatible/compatible, sin conflictos reales). +**Validation:** diff -r vacio en las 3 copias de specs y en el movimiento a archivo (unico evidencia aceptada); arbol de archivo con 8 ficheros originales intactos; conteo de tasks 29 [x] / 1 [ ] sin alteraciones; cambios visibles en git como renames + directorios nuevos no rastreados. +**Files:** `openspec/specs/{ocr-ingest-orchestration,ocr-processing,ocr-review-workflow}/spec.md` (nuevos), `openspec/changes/archive/2026-09-21-ocr-ingest-integration/` (movido), `docs/HISTORIAL_SESIONES.md` (esta entrada). Sin commits ni pushes. + +--- + +### 2026-09-21 - Subagente Inicializacion SDD RAG - sdd-init +**Agent:** Subagente Inicializacion SDD RAG · **Model:** ollama/glm-5.3:cloud · **Session:** `ses_f3c57954fffeQlFek8JT073tbu` (subagent of `ses_29bdbd003ffeLrLjUlFgnp08Y7`) +**Responsibility:** Guard de inicializacion SDD antes del archivo honesto: detectar proyectos en alcance, persistir contexto/capacidades de testing en modo hybrid y validar el registro de skills. Sin tocar codigo, OCR, estado de produccion, artefactos del cambio activo ni realizar operaciones de archivo. +**Work:** Descubrimiento acotado (raiz + 2 niveles, excluyendo node_modules/dist/.venv/backups/env*) encontro DOS proyectos en alcance: `./` (Node 22 + TS 5.8, `npm test`) y `ocr-service/` (Python 3.11 + FastAPI + PaddleOCR, `ocr-service/.venv/bin/python -m pytest ocr-service/tests`). Ningun comando workspace-level unico cubre ambos (sin Makefile, sin CI, sin script combinado) → el `strict_tdd: true` explicito falla cerrado a `false` segun la puerta de decision del contrato sdd-init. Correccion minima y documentada en `openspec/config.yaml` (contexto de testing, proyecto ocr-service anadido, `strict_tdd: false` con comentario explicativo, `rules.apply`/`rules.verify` ahora listan ambos comandos). Registro `.atl/skill-registry.md` validado como vigente (2026-09-20, sin deriva): no se reescribio. Persistidas 3 observaciones Engram en `rag-service` bajo topic keys canonicos (`sdd-init/rag-service` #3688, `sdd/rag-service/testing-capabilities` #3689, `skill-registry` #3690). +**Validation:** `npm test` 101/101, `npm run check` exit 0, `ocr-service/.venv/bin/python -m pytest ocr-service/tests` 16/16 (1 warning de starlette). +**Files:** `openspec/config.yaml` (modificado), `docs/HISTORIAL_SESIONES.md` (esta entrada). Sin commits ni pushes. + +--- + +### 2026-09-20 - Agente RAG 2 - Resolucion auditada de OCR v4 +**Agent:** Agente RAG 2 · **Model:** openai/gpt-5.6-terra · **Session:** `ses_29bdbd003ffeLrLjUlFgnp08Y7` +**Work:** Consultada v4 en modo autenticado y de solo lectura: la revision devolvio `409 OCR_ARTIFACT_UNAVAILABLE` con accion `use_admin_recovery`; su estado era `review_required`, sin activacion. Tras aprobacion explicita del usuario, se ejecuto la recuperacion administrativa para `ce1b6462-7617-4721-aa3a-8e8216a584ed`. La API respondio `200` con resultado `closed_failed`. +**Validation:** PostgreSQL confirma `state=failed`, `error_code=OCR_ARTIFACT_UNAVAILABLE`, candidata no activa, fuente con su version activa, un registro de auditoria `closed_failed` y cero candidatas OCR bloqueantes para esa fuente. +**Learned:** Una candidata heredada puede informar paginas OCR completas y, aun asi, ser irrecuperable si falta el artefacto durable de revision. El cierre auditado libera la fuente sin indexar ni activar contenido. +**Files:** `docs/PENDIENTES_RAG.md`, `docs/OPERATIVA.md`, `docs/HISTORIAL_SESIONES.md`, `../docs/REGISTRO_SITUACIONES.md`. + +--- + ### 2026-09-17 - Agente RAG 2 - Verificacion de despliegue OCR **Agent:** Agente RAG 2 · **Model:** openai/gpt-5.6-terra · **Session:** `ses_29bdbd003ffeLrLjUlFgnp08Y7` **Work:** Verificado el despliegue de `07e6ed2` tras los deploys manuales de OCR y RAG: `/health` publico devuelve RAG, PostgreSQL, Qdrant y reconciliador sanos; OCR interno devuelve `live=ok`, `ready=true`, cola vacia y version `0.1.0`. La migracion `003_ocr_recovery_audit.sql` y su tabla estan presentes. Con token administrativo, revision y recuperacion de una candidata inexistente devuelven el contrato seguro `404 OCR_CANDIDATE_NOT_FOUND` y no modifican datos. Se agrego la regresion local que representa dos instancias consecutivas del reconciliador y demuestra que trabajo OCR completado sin candidata durable permanece intacto tras reiniciar RAG. Se actualizo la operativa, que aun reflejaba el estado anterior en `false`. diff --git a/docs/OPERATIVA.md b/docs/OPERATIVA.md index 620c9d6..7b50ada 100644 --- a/docs/OPERATIVA.md +++ b/docs/OPERATIVA.md @@ -1,8 +1,8 @@ # Operativa del servicio RAG **Modulo:** RAG -**Ultima actualizacion:** 2026-09-17 -**Version:** 1.2 +**Ultima actualizacion:** 2026-09-21 +**Version:** 1.3 --- @@ -21,8 +21,16 @@ Este documento registra los hechos operativos del servicio RAG: la configuracion - `GET /health` del RAG confirma PostgreSQL, Qdrant y reconciliador sanos. OCR responde internamente `live=ok` y `ready=true`, con cola vacia y version `0.1.0`. - La migracion `003_ocr_recovery_audit.sql` esta aplicada y su tabla de auditoria existe. - Las rutas autenticadas de revision y recuperacion devuelven errores estructurados y seguros para una candidata inexistente: `404`, `OCR_CANDIDATE_NOT_FOUND` y accion `verify_version_id`. +- La candidata OCR heredada de FacturaTech v4 fue cerrada como `failed` el 2026-09-20 mediante recuperacion administrativa auditada tras confirmar `OCR_ARTIFACT_UNAVAILABLE`. No fue indexada ni activada; la fuente conserva su version activa y no tiene candidatas OCR bloqueantes. - Las imagenes en ejecucion tienen digest, pero ambas etiquetas OCI `org.opencontainers.image.revision` valen `unknown`: EasyPanel no esta pasando `BUILD_REVISION` durante el build. La identidad de revision verificable sigue pendiente. +## Proxima version preparada + +- RAG y OCR `0.2.0` estan preparados localmente; produccion sigue en `0.1.0` hasta ejecutar el despliegue. +- Ambas imagenes publican `org.opencontainers.image.version` y `org.opencontainers.image.revision`; las respuestas de salud muestran `version` y `revision`. +- EasyPanel debe construir ambas imagenes con `RAG_VERSION=0.2.0` u `OCR_VERSION=0.2.0` y `BUILD_REVISION=`. +- Por decision del usuario, RAG y OCR se despliegan juntos con `OCR_INGEST_ENABLED=true`. El despliegue por etapas queda reservado para diagnosticar un fallo cuyo origen no sea claro. + ## Lista rapida en EasyPanel 1. Rotar las credenciales expuestas (ver seccion siguiente); la rotacion sigue pendiente. @@ -30,7 +38,7 @@ Este documento registra los hechos operativos del servicio RAG: la configuracion 3. Mantener `OCR_INGEST_ENABLED=true` para el flujo OCR ya habilitado. 4. Pulsar `Deploy` en EasyPanel despues de cada cambio principal. 5. Verificar `GET /health` del RAG y las rutas internas de salud del OCR tras cada deploy. -6. Registrar el digest de cada imagen como identidad exacta de la version desplegada. Configurar `RAG_VERSION`, `OCR_VERSION` y `BUILD_REVISION` en un ciclo posterior si se necesita una revision legible desde health o etiquetas OCI. +6. Registrar el digest de cada imagen y confirmar que `version` y `revision` coinciden con la entrega desplegada. ## Accion de seguridad urgente: rotar credenciales @@ -106,7 +114,8 @@ El OCR es un servicio privado e independiente. El RAG solo lo llama si `OCR_INGE - Solo red interna del proyecto (`easypanel-ia_servicios`), sin dominio publico. - `OCR_SERVICE_URL` apunta al servicio OCR interno por el puerto `8000`; el valor por defecto del codigo es `http://ocr-service:8000`, valido para un servicio llamado `ocr-service` en la red interna. - Limites fijos (no configurables por variables): cola de 3 trabajos, concurrencia 1, render a 200 DPI, maximo 50 MiB por PDF y 100 paginas por trabajo. -- La readiness es falsa hasta que los modelos estan cargados; el healthcheck de arranque concede 90 segundos. +- Cada trabajo dispone de un lease y un limite total de 15 minutos. Solo se recupera una vez; la segunda interrupcion o el timeout son terminales. +- La readiness es falsa hasta que modelo, worker, sweeper y almacenamiento estan operativos; expone conteos por estado, recuperaciones y ultima limpieza sin incluir secretos. El healthcheck de arranque concede 90 segundos. ### Volumenes @@ -118,20 +127,21 @@ El OCR es un servicio privado e independiente. El RAG solo lo llama si `OCR_INGE ### Limpieza transitoria del OCR - Las filas transitorias, el PDF y los resultados del OCR se eliminan tras la transferencia duradera al RAG, cuando el resultado es aceptado. -- Si esa transferencia no llega, la limpieza automatica expira los trabajos a las 24 horas. +- Si esa transferencia no llega, un sweeper ejecutado al arrancar y cada 15 minutos expira los trabajos a las 24 horas, hace checkpoint WAL y vacuum incremental. +- Antes de admitir el PDF y antes de publicar cada PNG, OCR limita su uso al menor entre 10 % del filesystem y 2 GiB, reservando libre el mayor entre 10 % y 2 GiB. La presion falla cerrada como `OCR_STORAGE_PRESSURE`. - La limpieza transitoria nunca toca los artefactos duraderos del RAG ni el corpus activo. ### Verificacion posterior al despliegue (cuando cambian ambos servicios) -1. Publicar el codigo en Git `main` y desplegar OCR y RAG desde EasyPanel. +1. Publicar el codigo en Git `main` y desplegar juntos OCR y RAG `0.2.0` desde EasyPanel con `OCR_INGEST_ENABLED=true`. 2. Verificar el OCR internamente: `GET /health/live` y `GET /health/ready` deben responder con el servicio preparado. 3. Verificar `GET /health` del RAG: PostgreSQL, Qdrant y reconciliador deben estar correctos. 4. Comprobar una ruta OCR protegida sin token. Un `401 Lifecycle admin token is required` confirma que OCR esta habilitado y que la autenticacion administrativa permanece protegida; no es un error que requiera correccion. 5. Contrastar `OCR_INGEST_ENABLED=true` en el entorno persistido de EasyPanel y en el proceso del contenedor si la UI no coincide con el comportamiento efectivo. 6. Consultar con token administrativo una version inexistente en revision y recuperacion: ambas deben devolver un error estructurado `404 OCR_CANDIDATE_NOT_FOUND` sin modificar datos. 7. Confirmar que `rag_schema_migrations` contiene `003_ocr_recovery_audit.sql` y que existe `rag_ocr_recovery_audit`. -8. Revisar los digests y etiquetas OCI de ambas imagenes. El digest es la identidad exacta actual; una revision OCI legible es una mejora operativa posterior si sigue en `unknown`. -9. No crear, recuperar, aprobar, indexar ni activar candidatas hasta la fase aprobada para v4 y FacturaTech. +8. Revisar los digests y confirmar que las etiquetas OCI y las respuestas de salud muestran `0.2.0` y el commit desplegado, no `unknown`. +9. Tras completar estas comprobaciones, crear una candidata FacturaTech no activa. No aprobar, indexar ni activar hasta presentar la evidencia al usuario y recibir autorizacion explicita. ### Rollback de emergencia diff --git a/docs/PENDIENTES_RAG.md b/docs/PENDIENTES_RAG.md index dcb1b07..8650a04 100644 --- a/docs/PENDIENTES_RAG.md +++ b/docs/PENDIENTES_RAG.md @@ -1,6 +1,6 @@ # Pendientes priorizados del RAG -**Ultima actualizacion:** 2026-09-17 +**Ultima actualizacion:** 2026-09-21 **Responsable de la priorizacion:** Usuario **Estado:** Activo @@ -23,23 +23,42 @@ Esta secuencia tiene prioridad sobre la aceptacion productiva pendiente de Factu ### 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. +**Estado:** Completada en produccion el 2026-09-20, con aprobacion explicita del usuario. -**Salida:** v4 queda resuelta de forma auditable, sin crear todavia otra candidata. +1. Completado: v4 respondio `409 OCR_ARTIFACT_UNAVAILABLE` y accion `use_admin_recovery` mediante el contrato autenticado. +2. Completado: confirmada la ausencia de la evidencia durable necesaria para su revision. +3. Completado: recuperacion administrativa auditada aplicada sobre la unica candidata heredada. +4. Completado: v4 quedo en estado `failed`, no activa, sin indexacion ni activacion; la fuente conserva su version activa. +5. Completado: no quedan candidatas `pending`, `indexing` ni `review_required` que bloqueen la fuente. -### Fase 3. Finalizar la aceptacion de FacturaTech +**Salida:** v4 queda resuelta de forma auditable como `failed`, sin crear otra candidata. -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. +### Fase 3. Hardening OCR y aceptacion de FacturaTech -**Salida:** la ingesta OCR de FacturaTech queda aceptada y activada de forma segura. +**Estado:** D1-D4 implementados y validados de forma independiente; queda pendiente desplegar la version `0.2.0` y ejecutar la aceptacion productiva sin activar contenido automaticamente. + +La candidata v5 `5f2317c6-7a8a-4e08-a614-f8189602ebb8` completo 25/25 paginas OCR, pero fallo despues cuando RAG solicito 25 imagenes de revision en paralelo. El rerender concurrente provoco un `SIGSEGV` de PDFium/FreeType, salida `139` y reinicio del contenedor; no hubo OOM. La candidata quedo fallida y la version activa no cambio. + +#### Fase 3A. Hardening del runtime + +1. Completado localmente: cola SQLite durable con un unico worker real y recuperacion tras reinicio. +2. Completado localmente: PNG privados durante el render inicial y servicio desde fichero sin volver a ejecutar PDFium. +3. Completado localmente: transferencia RAG secuencial, idempotente y reanudable mediante comprobantes privados por pagina. +4. Completado localmente: exclusion mutua defensiva de PDFium y cierre determinista de recursos. +5. Completado localmente: observabilidad segura, limites de cola/tiempo/disco, recuperacion por lease, sweeper periodico y vacuum incremental. +6. Completado: validacion independiente con PDF real de 25 paginas, reinicio durante el trabajo, concurrencia, reanudacion, ausencia de duplicados y limpieza cuantitativa. + +#### Fase 3B. Aceptacion productiva + +1. Desplegar juntos OCR y RAG `0.2.0`, con OCR habilitado desde el inicio y sin activar contenido. +2. Crear una nueva candidata no activada para el documento. +3. Verificar la evidencia durable, las 25 imagenes y sus 34 entradas. +4. Comprobar `CBG04a`, `FAT07`, `DSAU08` y `NSAV06`. +5. Requerir aprobacion humana explicita antes de indexar o activar. + +El SDD `ocr-ingest-integration` se archiva con 29/30 tareas completas y 7.4 incompleta. La continuidad del hardening y de la aceptacion se controla mediante el contrato ODD canónico, sin declarar retrospectivamente superada la aceptacion fallida. + +**Salida:** el runtime OCR queda estabilizado y la ingesta OCR de FacturaTech se acepta y activa de forma segura solo con autorizacion explicita. ## 1. Documentacion y descubrimiento de la API diff --git a/ocr-service/Dockerfile b/ocr-service/Dockerfile index cdb210a..21208d0 100644 --- a/ocr-service/Dockerfile +++ b/ocr-service/Dockerfile @@ -1,5 +1,5 @@ FROM python:3.11.13-slim-bookworm -ARG OCR_VERSION=0.1.0 +ARG OCR_VERSION=0.2.0 ARG BUILD_REVISION=unknown LABEL resource.cpu.max="3" \ @@ -10,8 +10,10 @@ ENV PYTHONDONTWRITEBYTECODE=1 \ OCR_JOBS_DB="/data/jobs/jobs.db" \ OCR_LOAD_ENGINE="1" \ OCR_VERSION=${OCR_VERSION} \ + BUILD_REVISION=${BUILD_REVISION} \ HOME="/opt/ocr-home" -LABEL org.opencontainers.image.revision=${BUILD_REVISION} +LABEL org.opencontainers.image.version=${OCR_VERSION} \ + 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/README.md b/ocr-service/README.md index d9acbb2..36c57d1 100644 --- a/ocr-service/README.md +++ b/ocr-service/README.md @@ -1,6 +1,23 @@ -# OCR Service Development Environment +# Private OCR Service -Unit 5 uses the repository-local virtual environment `ocr-service/.venv`. It is intentionally retained between runs and ignored by Git. Do not install these dependencies globally. +The OCR service is an internal RAG dependency. It accepts authenticated PDF jobs, processes them with one PaddleOCR worker, and stores transient artifacts in its private `/data/jobs` volume. + +## Runtime guarantees + +| Area | Behavior | +|---|---| +| Capacity | At most three queued or running jobs and one active OCR execution. | +| Recovery | One recovery after interruption; a second interruption or 15-minute total timeout is terminal. | +| Review images | PNG files are persisted during the OCR render and served without rendering the PDF again. | +| Storage | Admission and PNG publication stop safely when the storage budget or free-space reserve is unavailable. | +| Cleanup | Confirmed jobs are deleted by RAG; unconfirmed jobs expire after 24 hours. | +| Readiness | Requires the model, worker, sweeper, and storage to be operational. | + +The API is private and its interactive documentation is intentionally disabled. Operational limits and deployment checks are documented in `../docs/OPERATIVA.md`; the full lifecycle contract is in `../docs/CONTRATO_CICLO_VIDA_Y_OCR.md`. + +## Development environment + +Use the repository-local virtual environment `ocr-service/.venv`. It is intentionally retained between runs and ignored by Git. Do not install these dependencies globally. ## Install diff --git a/ocr-service/app/jobs.py b/ocr-service/app/jobs.py new file mode 100644 index 0000000..8193de1 --- /dev/null +++ b/ocr-service/app/jobs.py @@ -0,0 +1,419 @@ +import hashlib +import json +import logging +import os +import shutil +import sqlite3 +import threading +import uuid +from datetime import datetime, timedelta, timezone +from pathlib import Path +from typing import Any, Callable + +from .engine import OcrEngine +from .render import process_pdf + + +QUEUE_CAPACITY = 3 +LEASE_DURATION = timedelta(minutes=15) +TRANSIENT_TTL = timedelta(hours=24) +SWEEP_INTERVAL_SECONDS = 15 * 60 +MAX_TRANSIENT_BYTES = 2 * 1024**3 +MIN_FREE_BYTES = 2 * 1024**3 +STORAGE_RATIO = 0.10 + +logger = logging.getLogger("ocr.jobs") + + +class QueueFullError(RuntimeError): + pass + + +class StoragePressureError(RuntimeError): + pass + + +class JobActiveError(RuntimeError): + pass + + +class JobQueue: + def __init__(self, db_path: str | Path, now: Callable[[], datetime], sweep_interval: float = SWEEP_INTERVAL_SECONDS) -> None: + self.now = now + self.sweep_interval = sweep_interval + self.lock = threading.Lock() + self.wake_worker = threading.Event() + self.wake_sweeper = threading.Event() + self.stop_worker = threading.Event() + self.worker: threading.Thread | None = None + self.sweeper: threading.Thread | None = None + self.worker_error: str | None = None + self.sweeper_error: str | None = None + self.last_sweep_at: str | None = None + self.db_path = Path(db_path) + self.artifacts_dir = self.db_path.parent / "artifacts" + if str(db_path) == ":memory:": + self.artifacts_dir = Path.cwd() / ".ocr-test-artifacts" + self.artifacts_dir.mkdir(parents=True, exist_ok=True) + self.connection = sqlite3.connect(str(db_path), check_same_thread=False) + self.connection.row_factory = sqlite3.Row + if self.connection.execute("PRAGMA auto_vacuum").fetchone()[0] != 2: + self.connection.execute("PRAGMA auto_vacuum=INCREMENTAL") + self.connection.execute("VACUUM") + self.connection.execute("PRAGMA journal_mode=WAL") + self.connection.execute( + "CREATE TABLE IF NOT EXISTS jobs (" + "job_id TEXT PRIMARY KEY, idempotency_key TEXT UNIQUE, payload_hash TEXT NOT NULL, " + "document_sha256 TEXT NOT NULL, pages TEXT NOT NULL, status TEXT NOT NULL, created_at TEXT NOT NULL, " + "input_path TEXT, result TEXT, error TEXT, attempts INTEGER NOT NULL DEFAULT 0, " + "recovery_attempts INTEGER NOT NULL DEFAULT 0, started_at TEXT, completed_at TEXT, lease_expires_at TEXT)" + ) + columns = {row[1] for row in self.connection.execute("PRAGMA table_info(jobs)")} + for name, kind in ( + ("input_path", "TEXT"), ("result", "TEXT"), ("error", "TEXT"), + ("attempts", "INTEGER NOT NULL DEFAULT 0"), ("recovery_attempts", "INTEGER NOT NULL DEFAULT 0"), + ("started_at", "TEXT"), ("completed_at", "TEXT"), ("lease_expires_at", "TEXT"), + ): + if name not in columns: + self.connection.execute(f"ALTER TABLE jobs ADD COLUMN {name} {kind}") + self.connection.commit() + self.recover_interrupted() + + @staticmethod + def _iso(value: datetime) -> str: + return value.astimezone(timezone.utc).isoformat() + + def _artifact_path(self, job_id: str) -> Path: + return self.artifacts_dir / job_id / "input.pdf" + + def _write_file(self, target: Path, content: bytes) -> None: + target.parent.mkdir(mode=0o700, parents=True, exist_ok=True) + temporary = target.with_suffix(".tmp") + with temporary.open("wb") as output: + output.write(content) + output.flush() + os.fsync(output.fileno()) + os.replace(temporary, target) + directory_fd = os.open(target.parent, os.O_DIRECTORY) + try: + os.fsync(directory_fd) + finally: + os.close(directory_fd) + os.chmod(target, 0o600) + + def _write_artifact(self, job_id: str, content: bytes) -> str: + target = self._artifact_path(job_id) + self._write_file(target, content) + return str(target.relative_to(self.db_path.parent)) + + def _review_image_path(self, job_id: str, page: int) -> Path: + return self.artifacts_dir / job_id / "review-images" / f"page-{page:04d}.png" + + def _persist_review_image(self, job_id: str, page: int, png: bytes) -> None: + if not self._has_storage_capacity(len(png)): + raise StoragePressureError("OCR_STORAGE_PRESSURE") + self._write_file(self._review_image_path(job_id, page), png) + + def _remove_artifacts(self, job_id: str) -> None: + shutil.rmtree(self.artifacts_dir / job_id, ignore_errors=True) + + def _remove_review_images(self, job_id: str) -> None: + shutil.rmtree(self.artifacts_dir / job_id / "review-images", ignore_errors=True) + + def _owned_storage_bytes(self) -> int: + total = sum(path.stat().st_size for path in self.artifacts_dir.rglob("*") if path.is_file()) + if str(self.db_path) != ":memory:": + for suffix in ("", "-shm", "-wal"): + path = Path(f"{self.db_path}{suffix}") + if path.exists(): + total += path.stat().st_size + return total + + def _has_storage_capacity(self, required_bytes: int = 0) -> bool: + usage = shutil.disk_usage(self.artifacts_dir) + budget = min(int(usage.total * STORAGE_RATIO), MAX_TRANSIENT_BYTES) + reserve = max(int(usage.total * STORAGE_RATIO), MIN_FREE_BYTES) + return self._owned_storage_bytes() + required_bytes <= budget and usage.free - required_bytes >= reserve + + def storage_available(self) -> bool: + try: + if not self._has_storage_capacity(): + return False + probe = self.artifacts_dir / ".write-probe" + with probe.open("wb") as output: + output.write(b"ok") + output.flush() + os.fsync(output.fileno()) + probe.unlink() + return True + except OSError: + return False + + def depth(self) -> int: + row = self.connection.execute("SELECT count(*) AS count FROM jobs WHERE status IN ('queued','running')").fetchone() + return int(row["count"]) + + def operational_state(self) -> dict[str, Any]: + rows = self.connection.execute("SELECT status,count(*) AS count FROM jobs GROUP BY status").fetchall() + states = {status: 0 for status in ("queued", "running", "succeeded", "failed")} + states.update({row["status"]: int(row["count"]) for row in rows}) + worker_alive = self.worker is not None and self.worker.is_alive() + if self.worker_error: + worker_state = "failed" + elif not worker_alive: + worker_state = "stopped" + else: + worker_state = "processing" if states["running"] else "idle" + recovery_attempts = self.connection.execute("SELECT COALESCE(sum(recovery_attempts),0) FROM jobs").fetchone()[0] + return { + "queueDepth": states["queued"] + states["running"], "queueCapacity": QUEUE_CAPACITY, + "queueStates": states, "concurrency": 1, "workerState": worker_state, + "workerOperational": worker_alive and self.worker_error is None, + "sweeperOperational": self.sweeper is not None and self.sweeper.is_alive() and self.sweeper_error is None, + "recoveryAttempts": int(recovery_attempts), "lastSweepAt": self.last_sweep_at, + } + + def contains(self, key: str) -> bool: + return self.connection.execute("SELECT 1 FROM jobs WHERE idempotency_key=?", (key,)).fetchone() is not None + + @staticmethod + def ack(row: sqlite3.Row) -> dict[str, Any]: + return { + "jobId": row["job_id"], "status": row["status"], "documentSha256": row["document_sha256"], + "requestedPages": json.loads(row["pages"]), "configVersion": "ocr-v1", "createdAt": row["created_at"], + } + + def submit(self, key: str, request: dict[str, Any], pdf: bytes) -> tuple[dict[str, Any], bool]: + payload_hash = hashlib.sha256(json.dumps(request, sort_keys=True, separators=(",", ":")).encode()).hexdigest() + with self.lock: + row = self.connection.execute("SELECT * FROM jobs WHERE idempotency_key=?", (key,)).fetchone() + if row: + if row["payload_hash"] != payload_hash: + raise ValueError("IDEMPOTENCY_CONFLICT") + return self.ack(row), False + if self.depth() >= QUEUE_CAPACITY: + raise QueueFullError("QUEUE_FULL") + if not self._has_storage_capacity(len(pdf)): + raise StoragePressureError("OCR_STORAGE_PRESSURE") + job_id = f"ocr_{uuid.uuid4()}" + try: + input_path = self._write_artifact(job_id, pdf) + self.connection.execute( + "INSERT INTO jobs (job_id,idempotency_key,payload_hash,document_sha256,pages,status,created_at,input_path) " + "VALUES (?,?,?,?,?,?,?,?)", + (job_id, key, payload_hash, request["documentSha256"], json.dumps(request["pages"]), "queued", self._iso(self.now()), input_path), + ) + self.connection.commit() + except Exception: + self.connection.rollback() + self._remove_artifacts(job_id) + raise + row = self.connection.execute("SELECT * FROM jobs WHERE job_id=?", (job_id,)).fetchone() + logger.info("ocr_transition job_id=%s from=none to=queued category=accepted", job_id) + self.wake_worker.set() + return self.ack(row), True + + def recover_interrupted(self) -> None: + with self.lock: + interrupted = self.connection.execute( + "SELECT job_id,recovery_attempts,started_at FROM jobs WHERE status='running'" + ).fetchall() + for row in interrupted: + timed_out = row["started_at"] is not None and datetime.fromisoformat(row["started_at"]) <= self.now() - LEASE_DURATION + if timed_out or row["recovery_attempts"] >= 1: + code = "OCR_PROCESSING_TIMEOUT" if timed_out else "OCR_RECOVERY_EXHAUSTED" + self.connection.execute( + "UPDATE jobs SET status='failed', error=?, completed_at=?, lease_expires_at=NULL WHERE job_id=?", + (json.dumps({"code": code, "message": "OCR processing did not complete safely"}), self._iso(self.now()), row["job_id"]), + ) + logger.warning("ocr_transition job_id=%s from=running to=failed category=%s", row["job_id"], code.lower()) + else: + self.connection.execute( + "UPDATE jobs SET status='queued', recovery_attempts=recovery_attempts+1, lease_expires_at=NULL WHERE job_id=?", + (row["job_id"],), + ) + logger.warning("ocr_transition job_id=%s from=running to=queued category=startup_recovery", row["job_id"]) + self.connection.commit() + if interrupted: + self.wake_worker.set() + + def start_worker(self, engine: OcrEngine | None) -> None: + self.sweep_expired() + if self.sweeper is None: + self.sweeper = threading.Thread(target=self._run_sweeper, name="ocr-sweeper", daemon=True) + self.sweeper.start() + if engine is not None and self.worker is None: + self.worker = threading.Thread(target=self._run_worker, args=(engine,), name="ocr-worker", daemon=True) + self.worker.start() + + def stop(self) -> None: + self.stop_worker.set() + self.wake_worker.set() + self.wake_sweeper.set() + if self.worker is not None: + self.worker.join(timeout=5) + if self.sweeper is not None: + self.sweeper.join(timeout=5) + + def _claim(self) -> sqlite3.Row | None: + with self.lock: + row = self.connection.execute("SELECT * FROM jobs WHERE status='queued' ORDER BY created_at LIMIT 1").fetchone() + if row is None: + return None + now = self._iso(self.now()) + lease = self._iso(self.now() + LEASE_DURATION) + changed = self.connection.execute( + "UPDATE jobs SET status='running', attempts=attempts+1, started_at=COALESCE(started_at,?), lease_expires_at=? " + "WHERE job_id=? AND status='queued'", (now, lease, row["job_id"]), + ).rowcount + self.connection.commit() + claimed = self.connection.execute("SELECT * FROM jobs WHERE job_id=?", (row["job_id"],)).fetchone() if changed else None + if claimed: + logger.info("ocr_transition job_id=%s from=queued to=running category=claimed", row["job_id"]) + self.wake_sweeper.set() + return claimed + + def _run_worker(self, engine: OcrEngine) -> None: + try: + self._worker_loop(engine) + except Exception: + self.worker_error = "OCR_WORKER_FAILED" + logger.error("ocr_worker category=worker_failed") + + def _worker_loop(self, engine: OcrEngine) -> None: + while not self.stop_worker.is_set(): + row = self._claim() + if row is None: + self.wake_worker.wait(timeout=1) + self.wake_worker.clear() + continue + try: + input_path = self.db_path.parent / row["input_path"] + result = process_pdf( + row["job_id"], row["document_sha256"], input_path.read_bytes(), json.loads(row["pages"]), engine, + persist_review_image=lambda page, png: self._persist_review_image(row["job_id"], page, png), + ) + update = ("succeeded", json.dumps(result, sort_keys=True, separators=(",", ":")), None) + except StoragePressureError: + self._remove_review_images(row["job_id"]) + update = ("failed", None, json.dumps({"code": "OCR_STORAGE_PRESSURE", "message": "OCR storage capacity is unavailable"})) + except Exception: + update = ("failed", None, json.dumps({"code": "OCR_PROCESSING_FAILED", "message": "OCR processing failed"})) + with self.lock: + changed = self.connection.execute( + "UPDATE jobs SET status=?, result=?, error=?, completed_at=?, lease_expires_at=NULL " + "WHERE job_id=? AND status='running'", + (*update, self._iso(self.now()), row["job_id"]), + ).rowcount + self.connection.commit() + if changed: + category = "completed" if update[0] == "succeeded" else json.loads(update[2])["code"].lower() + logger.info("ocr_transition job_id=%s from=running to=%s category=%s", row["job_id"], update[0], category) + else: + stored = self.connection.execute("SELECT 1 FROM jobs WHERE job_id=?", (row["job_id"],)).fetchone() + (self._remove_review_images if stored else self._remove_artifacts)(row["job_id"]) + + def _next_sweep_delay(self) -> float: + with self.lock: + row = self.connection.execute( + "SELECT min(lease_expires_at),min(started_at) FROM jobs WHERE status='running'" + ).fetchone() + deadlines = [datetime.fromisoformat(row[0])] if row[0] else [] + if row[1]: + deadlines.append(datetime.fromisoformat(row[1]) + LEASE_DURATION) + if not deadlines: + return self.sweep_interval + return max(0.0, min(self.sweep_interval, min((deadline - self.now()).total_seconds() for deadline in deadlines))) + + def _run_sweeper(self) -> None: + while not self.stop_worker.is_set(): + try: + delay = self._next_sweep_delay() + except Exception: + delay = self.sweep_interval + if self.wake_sweeper.wait(delay): + self.wake_sweeper.clear() + if self.stop_worker.is_set(): + return + continue + if self.stop_worker.is_set(): + return + try: + self.sweep_expired() + self.sweeper_error = None + except Exception: + self.sweeper_error = "OCR_SWEEPER_FAILED" + logger.error("ocr_maintenance category=sweeper_failed") + + def status(self, job_id: str) -> dict[str, Any] | None: + row = self.connection.execute("SELECT * FROM jobs WHERE job_id=?", (job_id,)).fetchone() + if not row: + return None + total_pages = len(json.loads(row["pages"])) + return {"jobId": job_id, "status": row["status"], "completedPages": total_pages if row["status"] == "succeeded" else 0, + "totalPages": total_pages, "error": json.loads(row["error"]) if row["error"] else None} + + def result(self, job_id: str) -> tuple[str, dict[str, Any] | None] | None: + 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 review_image(self, job_id: str, page: int) -> tuple[str, bytes] | None: + row = self.connection.execute("SELECT document_sha256,pages,status FROM jobs WHERE job_id=?", (job_id,)).fetchone() + if not row: + return None + if row["status"] != "succeeded": + raise RuntimeError("RESULT_NOT_READY") + if page not in json.loads(row["pages"]): + raise ValueError("INVALID_PAGE") + try: + return row["document_sha256"], self._review_image_path(job_id, page).read_bytes() + except FileNotFoundError as error: + raise RuntimeError("RESULT_NOT_READY") from error + + def delete(self, job_id: str) -> None: + self.sweep_expired() + with self.lock: + row = self.connection.execute("SELECT status FROM jobs WHERE job_id=?", (job_id,)).fetchone() + if row is not None and row["status"] == "running": + raise JobActiveError("JOB_ACTIVE") + self.connection.execute("DELETE FROM jobs WHERE job_id=?", (job_id,)) + self.connection.commit() + self._remove_artifacts(job_id) + if row is not None: + logger.info("ocr_transition job_id=%s from=%s to=deleted category=confirmed_cleanup", job_id, row["status"]) + + def sweep_expired(self) -> None: + current = self.now() + cutoff = self._iso(current - TRANSIENT_TTL) + timeout_cutoff = self._iso(current - LEASE_DURATION) + with self.lock: + interrupted = self.connection.execute( + "SELECT job_id,recovery_attempts,started_at FROM jobs WHERE status='running' AND " + "(lease_expires_at IS NULL OR lease_expires_at <= ? OR started_at <= ?)", + (self._iso(current), timeout_cutoff), + ).fetchall() + for row in interrupted: + timed_out = row["started_at"] is not None and row["started_at"] <= timeout_cutoff + if timed_out or row["recovery_attempts"] >= 1: + code = "OCR_PROCESSING_TIMEOUT" if timed_out else "OCR_RECOVERY_EXHAUSTED" + self.connection.execute( + "UPDATE jobs SET status='failed',error=?,completed_at=?,lease_expires_at=NULL WHERE job_id=?", + (json.dumps({"code": code, "message": "OCR processing did not complete safely"}), self._iso(current), row["job_id"]), + ) + logger.warning("ocr_transition job_id=%s from=running to=failed category=%s", row["job_id"], code.lower()) + else: + self.connection.execute( + "UPDATE jobs SET status='queued',recovery_attempts=recovery_attempts+1,lease_expires_at=NULL WHERE job_id=?", + (row["job_id"],), + ) + logger.warning("ocr_transition job_id=%s from=running to=queued category=lease_recovery", row["job_id"]) + rows = self.connection.execute("SELECT job_id FROM jobs WHERE status != 'running' AND created_at < ?", (cutoff,)).fetchall() + self.connection.executemany("DELETE FROM jobs WHERE job_id=?", [(row["job_id"],) for row in rows]) + self.connection.commit() + self.connection.execute("PRAGMA wal_checkpoint(PASSIVE)") + self.connection.execute("PRAGMA incremental_vacuum") + self.connection.commit() + self.last_sweep_at = self._iso(current) + for row in rows: + self._remove_artifacts(row["job_id"]) + logger.info("ocr_transition job_id=%s from=expired to=deleted category=ttl_cleanup", row["job_id"]) diff --git a/ocr-service/app/main.py b/ocr-service/app/main.py index 7438856..7766b3c 100644 --- a/ocr-service/app/main.py +++ b/ocr-service/app/main.py @@ -2,24 +2,20 @@ import hashlib import hmac import json import os -import sqlite3 -import threading -import uuid -from datetime import datetime, timedelta, timezone +from contextlib import asynccontextmanager +from datetime import datetime, timezone from pathlib import Path from typing import Annotated, Any, Callable -from fastapi import BackgroundTasks, Depends, FastAPI, File, Form, Header, HTTPException, Response, UploadFile +from fastapi import Depends, FastAPI, File, Form, Header, HTTPException, Response, UploadFile from .engine import OcrEngine +from .jobs import JobActiveError, JobQueue, QueueFullError, StoragePressureError from .models import load_runtime_engine -from .render import PdfRenderError, process_pdf, render_pdf_pages MAX_UPLOAD_BYTES = 50 * 1024 * 1024 MAX_PAGES = 100 -QUEUE_CAPACITY = 3 -TRANSIENT_TTL = timedelta(hours=24) ALLOWED_CONFIG = { "languages": ["es", "en"], "dpi": 200, @@ -41,110 +37,6 @@ def page_hash(pages: list[int]) -> str: return hashlib.sha256(value).hexdigest() -class JobQueue: - def __init__(self, path: str | Path, now: Callable[[], datetime]): - self.connection = sqlite3.connect(str(path), check_same_thread=False) - self.connection.row_factory = sqlite3.Row - self.lock = threading.Lock() - self.now = now - self.connection.execute( - "CREATE TABLE IF NOT EXISTS jobs (job_id TEXT PRIMARY KEY, idempotency_key TEXT UNIQUE, " - "payload_hash TEXT, document_sha256 TEXT, pages TEXT, status TEXT, created_at TEXT, " - "pdf BLOB, result TEXT, error TEXT)" - ) - columns = {row[1] for row in self.connection.execute("PRAGMA table_info(jobs)")} - for name, kind in (("pdf", "BLOB"), ("result", "TEXT"), ("error", "TEXT")): - if name not in columns: - self.connection.execute(f"ALTER TABLE jobs ADD COLUMN {name} {kind}") - self.connection.commit() - - def depth(self) -> int: - row = self.connection.execute("SELECT count(*) AS count FROM jobs WHERE status IN ('queued','running')").fetchone() - return int(row["count"]) - - def contains(self, key: str) -> bool: - return self.connection.execute("SELECT 1 FROM jobs WHERE idempotency_key=?", (key,)).fetchone() is not None - - @staticmethod - def ack(row: sqlite3.Row) -> dict[str, Any]: - return { - "jobId": row["job_id"], - "status": "queued", - "documentSha256": row["document_sha256"], - "requestedPages": json.loads(row["pages"]), - "configVersion": "ocr-v1", - "createdAt": row["created_at"], - } - - def submit(self, key: str, request: dict[str, Any], pdf: bytes) -> tuple[dict[str, Any], bool]: - payload_hash = hashlib.sha256(json.dumps(request, sort_keys=True, separators=(",", ":")).encode()).hexdigest() - with self.lock: - row = self.connection.execute("SELECT * FROM jobs WHERE idempotency_key=?", (key,)).fetchone() - if row: - if row["payload_hash"] != payload_hash: - fail(409, "IDEMPOTENCY_CONFLICT", "The idempotency key is already bound to another request") - return self.ack(row), True - if self.depth() >= QUEUE_CAPACITY: - fail(429, "QUEUE_FULL", "The OCR queue is full", True, {"Retry-After": "2"}) - values = ( - f"ocr_{uuid.uuid4()}", key, payload_hash, request["documentSha256"], - json.dumps(request["pages"]), "queued", self.now().isoformat(), pdf, None, None, - ) - self.connection.execute("INSERT INTO jobs VALUES (?,?,?,?,?,?,?,?,?,?)", values) - self.connection.commit() - return self.ack(self.connection.execute("SELECT * FROM jobs WHERE job_id=?", (values[0],)).fetchone()), True - - def execute(self, job_id: str, engine: OcrEngine) -> None: - with self.lock: - row = self.connection.execute("SELECT * FROM jobs WHERE job_id=?", (job_id,)).fetchone() - if row["status"] != "queued": - return - self.connection.execute("UPDATE jobs SET status='running' WHERE job_id=?", (job_id,)) - self.connection.commit() - try: - result = process_pdf(job_id, row["document_sha256"], bytes(row["pdf"]), json.loads(row["pages"]), engine) - update = ("succeeded", json.dumps(result, sort_keys=True, separators=(",", ":")), None, job_id) - except Exception as error: - update = ("failed", None, json.dumps({"code": "OCR_PROCESSING_FAILED", "message": str(error)}), job_id) - with self.lock: - self.connection.execute("UPDATE jobs SET status=?, result=?, error=? WHERE job_id=?", update) - self.connection.commit() - - def status(self, job_id: str) -> dict[str, Any] | None: - row = self.connection.execute("SELECT * FROM jobs WHERE job_id=?", (job_id,)).fetchone() - if not row: - return None - total_pages = len(json.loads(row["pages"])) - return {"jobId": job_id, "status": row["status"], - "completedPages": total_pages if row["status"] == "succeeded" else 0, - "totalPages": total_pages, "error": json.loads(row["error"]) if row["error"] else None} - - def result(self, job_id: str) -> tuple[str, dict[str, Any] | None] | None: - 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,)) - self.connection.commit() - - def sweep_expired(self) -> None: - with self.lock: - self.connection.execute("DELETE FROM jobs WHERE created_at < ?", ((self.now() - TRANSIENT_TTL).isoformat(),)) - self.connection.commit() - - def utc_now() -> datetime: return datetime.now(timezone.utc) @@ -157,9 +49,19 @@ def create_app( engine: OcrEngine | None = None, now: Callable[[], datetime] = utc_now, ) -> FastAPI: - application = FastAPI(title="Private OCR Service", docs_url=None, redoc_url=None) queue = JobQueue(db_path, now) + @asynccontextmanager + async def lifespan(_application: FastAPI): + queue.start_worker(engine) + try: + yield + finally: + queue.stop() + + application = FastAPI(title="Private OCR Service", docs_url=None, redoc_url=None, lifespan=lifespan) + application.state.queue = queue + def authorize(authorization: Annotated[str | None, Header()] = None) -> None: scheme, _, supplied = (authorization or "").partition(" ") if not token or scheme != "Bearer" or not hmac.compare_digest(supplied, token): @@ -167,19 +69,23 @@ def create_app( @application.get("/health/live") def live() -> dict[str, str]: - return {"status": "ok", "service": "ocr", "version": os.getenv("OCR_VERSION", "0.1.0")} + return {"status": "ok", "service": "ocr", "version": os.getenv("OCR_VERSION", "0.2.0"), + "revision": os.getenv("BUILD_REVISION", "unknown")} @application.get("/health/ready") - def ready(response: Response) -> dict[str, int | bool | str]: - queue.sweep_expired() - if not engine_ready: + def ready(response: Response) -> dict[str, Any]: + state = queue.operational_state() + storage_available = queue.storage_available() + ready_state = engine_ready and state["workerOperational"] and state["sweeperOperational"] and storage_available + if not ready_state: response.status_code = 503 - return {"ready": engine_ready, "queueDepth": queue.depth(), "queueCapacity": QUEUE_CAPACITY, "concurrency": 1, - "service": "ocr", "version": os.getenv("OCR_VERSION", "0.1.0")} + return {"ready": ready_state, "engineLoaded": engine_ready, "storageAvailable": storage_available, + **state, "service": "ocr", "version": os.getenv("OCR_VERSION", "0.2.0"), + "revision": os.getenv("BUILD_REVISION", "unknown")} @application.post("/v1/jobs", status_code=202, dependencies=[Depends(authorize)]) async def create_job( - background_tasks: BackgroundTasks, file: Annotated[UploadFile, File()], request: Annotated[str, Form()], + file: Annotated[UploadFile, File()], request: Annotated[str, Form()], idempotency_key: Annotated[str | None, Header(alias="Idempotency-Key")] = None, ) -> dict[str, Any]: content = await file.read(max_upload_bytes + 1) @@ -208,9 +114,14 @@ def create_app( expected_key = f'{digest}:ocr-v1:{page_hash(pages)}' if idempotency_key != expected_key and not queue.contains(idempotency_key or ""): fail(400, "INVALID_IDEMPOTENCY_KEY", "Idempotency-Key does not match request identity") - ack, created = queue.submit(idempotency_key or "", payload, content) - if created and engine is not None: - background_tasks.add_task(queue.execute, ack["jobId"], engine) + try: + ack, _created = queue.submit(idempotency_key or "", payload, content) + except ValueError: + fail(409, "IDEMPOTENCY_CONFLICT", "The idempotency key is already bound to another request") + except QueueFullError: + fail(429, "QUEUE_FULL", "The OCR queue is full", True, {"Retry-After": "2"}) + except StoragePressureError: + fail(503, "OCR_STORAGE_PRESSURE", "OCR storage capacity is unavailable", True, {"Retry-After": "15"}) return ack @application.get("/v1/jobs/{job_id}", dependencies=[Depends(authorize)]) @@ -231,7 +142,12 @@ def create_app( @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) + try: + stored = queue.review_image(job_id, page) + except RuntimeError: + fail(409, "RESULT_NOT_READY", "OCR review image is not ready") + except ValueError: + fail(422, "INVALID_PAGE", "OCR review image does not exist") if stored is None: fail(404, "JOB_NOT_FOUND", "OCR job does not exist") document_sha256, png = stored @@ -243,7 +159,10 @@ def create_app( @application.delete("/v1/jobs/{job_id}", status_code=204, dependencies=[Depends(authorize)]) def delete_job(job_id: str) -> Response: - queue.delete(job_id) + try: + queue.delete(job_id) + except JobActiveError: + fail(409, "JOB_ACTIVE", "An active OCR job cannot be deleted") return Response(status_code=204) return application diff --git a/ocr-service/app/render.py b/ocr-service/app/render.py index 673d0ef..132b953 100644 --- a/ocr-service/app/render.py +++ b/ocr-service/app/render.py @@ -1,6 +1,7 @@ import io import math import statistics +import threading import time from dataclasses import dataclass from typing import Callable @@ -10,6 +11,7 @@ from .engine import EngineLine, OcrEngine RENDER_DPI = 200 MAX_RENDER_PIXELS = 25_000_000 +PDFIUM_LOCK = threading.RLock() class PdfRenderError(ValueError): @@ -30,35 +32,36 @@ def render_pdf_pages(pdf: bytes, pages: list[int], max_pixels: int = MAX_RENDER_ if not pages or pages != sorted(set(pages)) or any(type(page) is not int or page < 1 for page in pages): raise PdfRenderError("PDF pages must be unique, ordered, one-based integers") - try: - document = pypdfium2.PdfDocument(pdf) - except Exception as error: - raise PdfRenderError("PDF cannot be opened for deterministic rendering") from error + with PDFIUM_LOCK: + try: + document = pypdfium2.PdfDocument(pdf) + except Exception as error: + raise PdfRenderError("PDF cannot be opened for deterministic rendering") from error - rendered: list[RenderedPage] = [] - try: - for page_number in pages: - if page_number > len(document): - raise PdfRenderError(f"PDF page {page_number} does not exist") - page = document[page_number - 1] - try: - page_width, page_height = page.get_size() - width = math.ceil(page_width * RENDER_DPI / 72) - height = math.ceil(page_height * RENDER_DPI / 72) - if width * height > max_pixels: - raise PdfRenderError("Rendered page exceeds the 25 megapixels limit") - bitmap = page.render(scale=RENDER_DPI / 72) + rendered: list[RenderedPage] = [] + try: + for page_number in pages: + if page_number > len(document): + raise PdfRenderError(f"PDF page {page_number} does not exist") + page = document[page_number - 1] try: - image = bitmap.to_pil() - output = io.BytesIO() - image.save(output, format="PNG") - rendered.append(RenderedPage(page_number, image.width, image.height, output.getvalue(), image)) + page_width, page_height = page.get_size() + width = math.ceil(page_width * RENDER_DPI / 72) + height = math.ceil(page_height * RENDER_DPI / 72) + if width * height > max_pixels: + raise PdfRenderError("Rendered page exceeds the 25 megapixels limit") + bitmap = page.render(scale=RENDER_DPI / 72) + try: + image = bitmap.to_pil() + output = io.BytesIO() + image.save(output, format="PNG") + rendered.append(RenderedPage(page_number, image.width, image.height, output.getvalue(), image)) + finally: + bitmap.close() finally: - bitmap.close() - finally: - page.close() - finally: - document.close() + page.close() + finally: + document.close() return rendered @@ -84,9 +87,12 @@ def process_pdf( requested_pages: list[int], engine: OcrEngine, processing_ms: Callable[[int], int] | None = None, + persist_review_image: Callable[[int, bytes], None] | None = None, ) -> dict: results = [] for rendered in render_pdf_pages(pdf, requested_pages): + if persist_review_image is not None: + persist_review_image(rendered.page, rendered.png) started = time.perf_counter_ns() indexed = list(enumerate(engine.recognize(rendered.image), start=1)) indexed.sort(key=lambda item: (item[1].bbox[1], item[1].bbox[0], item[0])) diff --git a/ocr-service/tests/test_api.py b/ocr-service/tests/test_api.py index 9a04895..eef67ee 100644 --- a/ocr-service/tests/test_api.py +++ b/ocr-service/tests/test_api.py @@ -2,6 +2,7 @@ import hashlib import json import sqlite3 import sys +import time from datetime import datetime, timedelta, timezone from pathlib import Path @@ -12,6 +13,7 @@ sys.path.insert(0, str(Path(__file__).parents[1])) from app.main import create_app from app.engine import EngineLine +from app.jobs import JobQueue TOKEN = "unit-5-test-token" @@ -59,12 +61,31 @@ def client(tmp_path: Path) -> TestClient: return TestClient(create_app(TOKEN, tmp_path / "jobs.db", engine_ready=True)) +def wait_for_status(client: TestClient, job_id: str, expected: str) -> dict: + headers = {"Authorization": f"Bearer {TOKEN}"} + deadline = time.monotonic() + 5 + while time.monotonic() < deadline: + status = client.get(f"/v1/jobs/{job_id}", headers=headers).json() + if status["status"] == expected: + return status + time.sleep(0.02) + pytest.fail(f"Job {job_id} did not reach {expected}") + + 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", "service": "ocr", "version": "0.1.0"} - assert ready.json() == {"ready": True, "queueDepth": 0, "queueCapacity": 3, "concurrency": 1, - "service": "ocr", "version": "0.1.0"} + assert live.json() == {"status": "ok", "service": "ocr", "version": "0.2.0", "revision": "unknown"} + assert ready.status_code == 503 + assert ready.json() == { + "ready": False, "engineLoaded": True, "storageAvailable": True, + "queueDepth": 0, "queueCapacity": 3, + "queueStates": {"queued": 0, "running": 0, "succeeded": 0, "failed": 0}, + "concurrency": 1, "workerState": "stopped", "workerOperational": False, + "sweeperOperational": False, "recoveryAttempts": 0, + "lastSweepAt": None, + "service": "ocr", "version": "0.2.0", "revision": "unknown", + } for authorization in (None, "Bearer wrong-token"): headers = {"Idempotency-Key": key_for(request_for())} if authorization: @@ -152,10 +173,11 @@ def test_auth_health_sweeps_only_jobs_older_than_24_hours(tmp_path: Path): headers = {"Authorization": f"Bearer {TOKEN}"} current[0] += timedelta(hours=23, minutes=59) - assert client.get("/health/ready").status_code == 200 + assert client.get("/health/ready").status_code == 503 assert client.get(f"/v1/jobs/{job_id}", headers=headers).status_code == 200 current[0] += timedelta(minutes=2) - assert client.get("/health/ready").status_code == 200 + client.app.state.queue.sweep_expired() + assert client.get("/health/ready").status_code == 503 assert client.get(f"/v1/jobs/{job_id}", headers=headers).status_code == 404 with sqlite3.connect(db_path) as connection: assert connection.execute("SELECT count(*) FROM jobs").fetchone()[0] == 0 @@ -168,25 +190,185 @@ def test_auth_accepted_job_executes_and_exposes_integrity_bound_result(tmp_path: connection.execute("CREATE TABLE jobs (job_id TEXT PRIMARY KEY, idempotency_key TEXT UNIQUE, payload_hash TEXT, document_sha256 TEXT, pages TEXT, status TEXT, created_at TEXT)") queued_client = TestClient(create_app(TOKEN, db_path, engine_ready=True)) submit(queued_client, request_for(pdf, [1]), pdf) - client = TestClient(create_app(TOKEN, db_path, engine_ready=True, engine=FakeEngine())) - accepted = submit(client, request_for(pdf, [1]), pdf) - job_id = accepted.json()["jobId"] + with TestClient(create_app(TOKEN, db_path, engine_ready=True, engine=FakeEngine())) as client: + accepted = submit(client, request_for(pdf, [1]), pdf) + job_id = accepted.json()["jobId"] + status = wait_for_status(client, job_id, "succeeded") + result = client.get(f"/v1/jobs/{job_id}/result", headers={"Authorization": f"Bearer {TOKEN}"}) + assert status == {"jobId": job_id, "status": "succeeded", "completedPages": 1, "totalPages": 1, "error": None} + assert result.status_code == 200 + 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")] + ready = client.get("/health/ready") + assert ready.status_code == 200 + assert ready.json()["workerState"] == "idle" + assert ready.json()["sweeperOperational"] is True + assert ready.json()["queueStates"] == {"queued": 0, "running": 0, "succeeded": 1, "failed": 0} - status = client.get(f"/v1/jobs/{job_id}", headers={"Authorization": f"Bearer {TOKEN}"}) - result = client.get(f"/v1/jobs/{job_id}/result", headers={"Authorization": f"Bearer {TOKEN}"}) - assert status.json() == {"jobId": job_id, "status": "succeeded", "completedPages": 1, "totalPages": 1, "error": None} - assert result.status_code == 200 - 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_d2_serves_the_persisted_render_without_rerendering(tmp_path: Path, monkeypatch: pytest.MonkeyPatch): + pdf = FIXTURE.read_bytes() + db_path = tmp_path / "jobs.db" + with TestClient(create_app(TOKEN, db_path, engine_ready=True, engine=FakeEngine())) as client: + job_id = submit(client, request_for(pdf, [1]), pdf).json()["jobId"] + wait_for_status(client, job_id, "succeeded") + image_path = tmp_path / "artifacts" / job_id / "review-images" / "page-0001.png" + expected_png = image_path.read_bytes() + assert image_path.stat().st_mode & 0o777 == 0o600 + + def rerender_is_forbidden(*_args: object, **_kwargs: object) -> object: + raise AssertionError("review endpoint must not invoke PDFium") + + monkeypatch.setattr("app.render.render_pdf_pages", rerender_is_forbidden) + image = client.get(f"/v1/jobs/{job_id}/pages/1/image", headers={"Authorization": f"Bearer {TOKEN}"}) - 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 + assert image.content == expected_png + + +def test_d2_missing_or_unfinished_review_images_fail_closed(tmp_path: Path): + pdf = FIXTURE.read_bytes() + with TestClient(create_app(TOKEN, tmp_path / "jobs.db", engine_ready=True, engine=FakeEngine())) as client: + job_id = submit(client, request_for(pdf, [1]), pdf).json()["jobId"] + headers = {"Authorization": f"Bearer {TOKEN}"} + pending = client.get(f"/v1/jobs/{job_id}/pages/1/image", headers=headers) + assert (pending.status_code, pending.json()["detail"]["code"]) == (409, "RESULT_NOT_READY") + + wait_for_status(client, job_id, "succeeded") + (tmp_path / "artifacts" / job_id / "review-images" / "page-0001.png").unlink() + missing = client.get(f"/v1/jobs/{job_id}/pages/1/image", headers=headers) + assert (missing.status_code, missing.json()["detail"]["code"]) == (409, "RESULT_NOT_READY") + + +def test_d1_persists_input_as_an_artifact_not_a_sqlite_blob(tmp_path: Path): + queue = JobQueue(tmp_path / "jobs.db", lambda: datetime.now(timezone.utc)) + request = request_for() + ack, created = queue.submit(key_for(request), request, PDF) + + assert created is True + artifact = tmp_path / "artifacts" / ack["jobId"] / "input.pdf" + assert artifact.read_bytes() == PDF + assert artifact.stat().st_mode & 0o777 == 0o600 + with sqlite3.connect(tmp_path / "jobs.db") as connection: + columns = {column[1] for column in connection.execute("PRAGMA table_info(jobs)")} + stored_path = connection.execute("SELECT input_path FROM jobs WHERE job_id=?", (ack["jobId"],)).fetchone()[0] + assert "pdf" not in columns + assert stored_path == f"artifacts/{ack['jobId']}/input.pdf" + + +def test_d1_recovers_one_interrupted_job_then_fails_a_second_interruption(tmp_path: Path): + now = datetime(2026, 9, 21, tzinfo=timezone.utc) + db_path = tmp_path / "jobs.db" + queue = JobQueue(db_path, lambda: now) + request = request_for() + ack, _ = queue.submit(key_for(request), request, PDF) + with sqlite3.connect(db_path) as connection: + connection.execute("UPDATE jobs SET status='running' WHERE job_id=?", (ack["jobId"],)) + recovered = JobQueue(db_path, lambda: now) + assert recovered.status(ack["jobId"])["status"] == "queued" + with sqlite3.connect(db_path) as connection: + connection.execute("UPDATE jobs SET status='running' WHERE job_id=?", (ack["jobId"],)) + exhausted = JobQueue(db_path, lambda: now) + status = exhausted.status(ack["jobId"]) + assert status["status"] == "failed" + assert status["error"]["code"] == "OCR_RECOVERY_EXHAUSTED" + + +def test_d4_storage_pressure_rejects_admission_and_cleans_partial_images(tmp_path: Path, monkeypatch: pytest.MonkeyPatch): + application = create_app(TOKEN, tmp_path / "jobs.db", engine_ready=True) + client = TestClient(application) + queue = application.state.queue + monkeypatch.setattr(queue, "_has_storage_capacity", lambda _required=0: False) + + pressure = submit(client, request_for()) + assert pressure.status_code == 503 + assert pressure.json()["detail"] == { + "code": "OCR_STORAGE_PRESSURE", "message": "OCR storage capacity is unavailable", "retryable": True, + } + assert pressure.headers["Retry-After"] == "15" + assert queue.depth() == 0 + + monkeypatch.setattr(queue, "_has_storage_capacity", lambda _required=0: True) + job_id = submit(client, request_for()).json()["jobId"] + + def persist_one_page(*_args, **kwargs): + kwargs["persist_review_image"](1, b"private-png") + return {} + + monkeypatch.setattr("app.jobs.process_pdf", persist_one_page) + monkeypatch.setattr(queue, "_has_storage_capacity", lambda _required=0: False) + queue.start_worker(FakeEngine()) + try: + status = wait_for_status(client, job_id, "failed") + finally: + queue.stop() + assert status["error"]["code"] == "OCR_STORAGE_PRESSURE" + assert not (tmp_path / "artifacts" / job_id / "review-images").exists() + + +def test_d4_lease_recovery_timeout_and_active_delete_are_fail_closed(tmp_path: Path, caplog: pytest.LogCaptureFixture): + current = [datetime(2026, 9, 21, tzinfo=timezone.utc)] + application = create_app(TOKEN, tmp_path / "jobs.db", engine_ready=True, now=lambda: current[0]) + client = TestClient(application) + queue = application.state.queue + job_id = submit(client, request_for()).json()["jobId"] + queue._claim() + headers = {"Authorization": f"Bearer {TOKEN}"} + + active = client.delete(f"/v1/jobs/{job_id}", headers=headers) + assert (active.status_code, active.json()["detail"]["code"]) == (409, "JOB_ACTIVE") + + with queue.connection: + queue.connection.execute( + "UPDATE jobs SET lease_expires_at=? WHERE job_id=?", + (queue._iso(current[0] - timedelta(seconds=1)), job_id), + ) + queue.sweep_expired() + assert queue.status(job_id)["status"] == "queued" + assert queue.operational_state()["recoveryAttempts"] == 1 + + queue._claim() + current[0] += timedelta(minutes=16) + queue.sweep_expired() + assert queue.status(job_id)["error"]["code"] == "OCR_PROCESSING_TIMEOUT" + assert job_id in caplog.text + assert "category=ocr_processing_timeout" in caplog.text + + stale_id = queue.submit(key_for(request_for(PDF + b"stale")), request_for(PDF + b"stale"), PDF + b"stale")[0]["jobId"] + with queue.connection: + queue.connection.execute( + "UPDATE jobs SET status='running',started_at=? WHERE job_id=?", + (queue._iso(current[0] - timedelta(minutes=16)), stale_id), + ) + restarted = JobQueue(tmp_path / "jobs.db", lambda: current[0]) + assert restarted.status(stale_id)["error"]["code"] == "OCR_PROCESSING_TIMEOUT" + + +def test_d4_periodic_sweeper_purges_ttl_and_maintains_incremental_vacuum(tmp_path: Path): + current = [datetime(2026, 9, 21, tzinfo=timezone.utc)] + queue = JobQueue(tmp_path / "jobs.db", lambda: current[0], sweep_interval=0.01) + job_id = queue.submit(key_for(request_for()), request_for(), PDF)[0]["jobId"] + queue.start_worker(None) + current[0] += timedelta(hours=25) + try: + deadline = time.monotonic() + 1 + while queue.status(job_id) is not None and time.monotonic() < deadline: + time.sleep(0.01) + assert queue.status(job_id) is None + assert queue.operational_state()["sweeperOperational"] is True + assert queue.last_sweep_at is not None + assert queue.connection.execute("PRAGMA auto_vacuum").fetchone()[0] == 2 + finally: + queue.stop() def test_auth_result_rejects_unknown_not_ready_and_unauthorized_jobs(client: TestClient): diff --git a/ocr-service/tests/test_render.py b/ocr-service/tests/test_render.py index 56d6542..1644aa1 100644 --- a/ocr-service/tests/test_render.py +++ b/ocr-service/tests/test_render.py @@ -1,4 +1,6 @@ import hashlib +import threading +import time import sys from pathlib import Path @@ -7,6 +9,7 @@ import pytest sys.path.insert(0, str(Path(__file__).parents[1])) from app.engine import EngineLine, PaddleOcrEngine +from app import render from app.render import PdfRenderError, process_pdf, render_pdf_pages @@ -45,6 +48,35 @@ def test_render_rejects_invalid_pages_and_pixel_limit_before_image_creation() -> render_pdf_pages(FIXTURE.read_bytes(), [1], max_pixels=1_000_000) +def test_render_serializes_pdfium_calls_with_the_process_lock(monkeypatch: pytest.MonkeyPatch) -> None: + active = 0 + peak = 0 + + class Document: + def __init__(self, _pdf: bytes) -> None: + nonlocal active, peak + active += 1 + peak = max(peak, active) + time.sleep(0.03) + active -= 1 + raise RuntimeError("test document") + + class Pdfium: + PdfDocument = Document + + monkeypatch.setitem(sys.modules, "pypdfium2", Pdfium()) + def invoke() -> None: + with pytest.raises(PdfRenderError): + render_pdf_pages(FIXTURE.read_bytes(), [1]) + + threads = [threading.Thread(target=invoke) for _ in range(2)] + for thread in threads: + thread.start() + for thread in threads: + thread.join() + assert peak == 1 + + def test_render_builds_repeatable_result_schema_with_deterministic_engine() -> None: def run() -> dict: return process_pdf( diff --git a/package-lock.json b/package-lock.json index 67006b8..1ccf798 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "rag-service", - "version": "0.1.0", + "version": "0.2.0", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "rag-service", - "version": "0.1.0", + "version": "0.2.0", "dependencies": { "@qdrant/js-client-rest": "^1.15.0", "adm-zip": "^0.6.1", diff --git a/package.json b/package.json index c12f6aa..b665f28 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "rag-service", - "version": "0.1.0", + "version": "0.2.0", "private": true, "type": "module", "scripts": { diff --git a/src/api/openapi.ts b/src/api/openapi.ts index b175653..c539fab 100644 --- a/src/api/openapi.ts +++ b/src/api/openapi.ts @@ -18,7 +18,7 @@ export const openApiDocument = { openapi: "3.1.1", info: { title: "RAG Service API", - version: "0.1.0", + version: "0.2.0", description: "HTTP API for ingesting, retrieving, answering from, and evaluating scoped RAG knowledge." }, jsonSchemaDialect: "https://json-schema.org/draft/2020-12/schema", @@ -888,11 +888,12 @@ export const openApiDocument = { }, HealthResponse: { type: "object", - required: ["ok", "service", "version", "environment", "embeddings", "answer", "vectorStore", "postgres", "knowledgeLifecycle", "parsers", "chunking"], + required: ["ok", "service", "version", "revision", "environment", "embeddings", "answer", "vectorStore", "postgres", "knowledgeLifecycle", "parsers", "chunking"], properties: { ok: { type: "boolean" }, service: { type: "string", const: "rag" }, version: { type: "string" }, + revision: { type: "string" }, environment: { type: "string" }, embeddings: ref("ProviderModel"), answer: ref("ProviderModel"), diff --git a/src/app.ts b/src/app.ts index 0388f54..2ba60ee 100644 --- a/src/app.ts +++ b/src/app.ts @@ -22,7 +22,7 @@ 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 { persistComposedCandidateArtifact, persistOcrResultArtifact, persistReviewImageArtifacts, readOcrArtifactPageNumbers, resolveArtifactPath } from "./modules/ocr/artifacts.js"; +import { persistComposedCandidateArtifact, persistOcrResultArtifact, resolveArtifactPath, transferReviewImageArtifacts } from "./modules/ocr/artifacts.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"; @@ -76,9 +76,12 @@ export function createApp(options: AppOptions = {}) { 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 }); + await transferReviewImageArtifacts({ + rootDirectory: ocr.artifactRoot, + versionId: job.versionId, + documentId: job.documentId, + loadImage: async (page) => ocrClient!.getReviewImage(result.jobId, page, result.documentSha256) + }); return durable.result; }, async (versionId) => { const jobs = await catalog.listOcrJobs(versionId); @@ -206,6 +209,7 @@ export function createApp(options: AppOptions = {}) { ok, service: "rag", version: env.ragVersion, + revision: env.buildRevision, environment: env.nodeEnv, embeddings: { provider: embeddingProvider.providerName, diff --git a/src/config/env.ts b/src/config/env.ts index 0e88f71..c3628f6 100644 --- a/src/config/env.ts +++ b/src/config/env.ts @@ -33,7 +33,8 @@ function booleanEnv(name: string, fallback: boolean): boolean { export const env = { nodeEnv: process.env.NODE_ENV ?? "development", - ragVersion: process.env.RAG_VERSION ?? "0.1.0", + ragVersion: process.env.RAG_VERSION ?? "0.2.0", + buildRevision: process.env.BUILD_REVISION ?? "unknown", port: Number(process.env.PORT ?? 3000), qdrantUrl: requireEnv("QDRANT_URL", "http://localhost:6333"), qdrantApiKey: process.env.QDRANT_API_KEY ?? "", diff --git a/src/modules/ocr/artifacts.ts b/src/modules/ocr/artifacts.ts index 9801f93..bb23860 100644 --- a/src/modules/ocr/artifacts.ts +++ b/src/modules/ocr/artifacts.ts @@ -306,27 +306,45 @@ export async function removeReviewedPagesArtifact(rootDirectory: string, version 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); + return native.requestedPages as number[]; } 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 }> }> { + const expectedPages = await readOcrArtifactPageNumbers(input); + if (!sameNumbers(input.images.map(({ page }) => page), expectedPages)) throw new Error("OCR review image identity validation failed"); + const images = new Map(input.images.map((image) => [image.page, image])); + return transferReviewImageArtifacts({ ...input, loadImage: async (page) => images.get(page)! }); +} + +export async function transferReviewImageArtifacts(input: { + rootDirectory: string; versionId: string; documentId: string; + loadImage: (page: number) => Promise<{ 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 expectedPages = native.requestedPages as number[]; 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) { + for (const page of expectedPages) { + const existing = await readReviewImageCheckpoint({ ...input, versionDirectory, documentArtifactId: entry.documentArtifactId, + documentSha256: entry.originalSha256, page }); + if (existing) { images.push(existing); continue; } + const image = { page, ...await input.loadImage(page) }; 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 persisted = { page: image.page, relativePath, artifactPath, sha256: image.sha256, mimeType: "image/png" as const }; + const checkpointPath = path.join(directory, `page-${String(image.page).padStart(4, "0")}.json`); + await publishImmutable(checkpointPath, Buffer.from(canonicalJson({ schemaVersion: "1", versionId: input.versionId, + documentId: input.documentId, documentSha256: entry.originalSha256, page: image.page, relativePath, + sha256: image.sha256, mimeType: "image/png" }))); + images.push(persisted); } const manifestPath = path.join(directory, "manifest.json"); await publishImmutable(manifestPath, Buffer.from(canonicalJson({ schemaVersion: "1", versionId: input.versionId, documentId: input.documentId, @@ -334,6 +352,31 @@ export async function persistReviewImageArtifacts(input: { return { images: images.map(({ page, artifactPath, sha256 }) => ({ page, artifactPath, sha256 })) }; } +async function readReviewImageCheckpoint(input: { + rootDirectory: string; versionId: string; documentId: string; versionDirectory: string; + documentArtifactId: string; documentSha256: string; page: number; +}): Promise<{ page: number; relativePath: string; artifactPath: string; sha256: string; mimeType: "image/png" } | undefined> { + const directory = resolveArtifactPath(input.versionDirectory, path.posix.join("documents", input.documentArtifactId, "review-images")); + const checkpointPath = path.join(directory, `page-${String(input.page).padStart(4, "0")}.json`); + const checkpointFile = await lstat(checkpointPath).catch((error: NodeJS.ErrnoException) => { + if (error.code === "ENOENT") return undefined; + throw error; + }); + if (!checkpointFile) return undefined; + const checkpoint = await readPrivateJson(checkpointPath, undefined, "OCR review image checkpoint integrity validation failed"); + const relativePath = path.posix.join("documents", input.documentArtifactId, "review-images", `page-${String(input.page).padStart(4, "0")}.png`); + if (!isRecord(checkpoint) || checkpoint.schemaVersion !== "1" || checkpoint.versionId !== input.versionId + || checkpoint.documentId !== input.documentId || checkpoint.documentSha256 !== input.documentSha256 + || checkpoint.page !== input.page || checkpoint.relativePath !== relativePath || checkpoint.mimeType !== "image/png" + || typeof checkpoint.sha256 !== "string") throw new Error("OCR review image checkpoint identity validation failed"); + const artifactPath = resolveArtifactPath(input.versionDirectory, relativePath); + const imageFile = await lstat(artifactPath).catch(() => { throw new Error("OCR review image checkpoint integrity validation failed"); }); + const bytes = await readFile(artifactPath); + if (!imageFile.isFile() || (imageFile.mode & 0o777) !== 0o600 || sha256Hex(bytes) !== checkpoint.sha256 + || !bytes.subarray(0, 8).equals(Buffer.from("89504e470d0a1a0a", "hex"))) throw new Error("OCR review image checkpoint integrity validation failed"); + return { page: input.page, relativePath, artifactPath, sha256: checkpoint.sha256, mimeType: "image/png" }; +} + export async function readReviewImageArtifact(input: { rootDirectory: string; versionId: string; documentId: string; page: number; }): Promise<{ bytes: Buffer; sha256: string; mimeType: "image/png" }> { diff --git a/src/modules/ocr/client.ts b/src/modules/ocr/client.ts index 3a6d1ca..d6d49bd 100644 --- a/src/modules/ocr/client.ts +++ b/src/modules/ocr/client.ts @@ -121,16 +121,29 @@ export class OcrClient { 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 }; + for (let attempt = 0; attempt < 3; attempt += 1) { + let response: Response; + try { + response = await this.requestFetch(`${this.baseUrl}/v1/jobs/${encodeURIComponent(jobId)}/pages/${page}/image`, { headers: this.headers() }); + } catch { + if (attempt < 2) { await this.sleep(2_000 * 2 ** attempt); continue; } + throw new OcrClientError("OCR_NETWORK_ERROR", undefined, true); + } + if (!response.ok) { + if (TRANSIENT_STATUSES.has(response.status) && attempt < 2) { await this.sleep(2_000 * 2 ** attempt); continue; } + throw new OcrClientError(`OCR_HTTP_${response.status}`, response.status, TRANSIENT_STATUSES.has(response.status)); + } + const bytes = Buffer.from(await response.arrayBuffer()); + const sha256 = sha256Hex(bytes); + if (response.headers.get("content-type")?.split(";")[0] !== "image/png" + || response.headers.get("content-length") !== String(bytes.length) + || 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 }; + } + throw new OcrClientError("OCR_RETRY_EXHAUSTED", undefined, true); } async delete(jobId: string): Promise { diff --git a/tests/ocr/client.test.ts b/tests/ocr/client.test.ts index 5bb4202..02553a1 100644 --- a/tests/ocr/client.test.ts +++ b/tests/ocr/client.test.ts @@ -11,7 +11,8 @@ import { readOcrResultArtifact, resolveArtifactPath, stageOcrArtifacts, - sweepOrphanArtifacts + sweepOrphanArtifacts, + transferReviewImageArtifacts } from "../../src/modules/ocr/artifacts.js"; import { canonicalJson, hashOrderedPairs, sha256Hex } from "../../src/shared/utils/ids.js"; @@ -181,7 +182,7 @@ test("review image transfer validates authentication, identity, content type, an 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, + "content-type": "image/png", "content-length": String(png.length), "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 @@ -191,6 +192,63 @@ test("review image transfer validates authentication, identity, content type, an await assert.rejects(client.getReviewImage(jobId, 2, documentSha256), /integrity validation failed/); }); +test("review image transfer treats size and HTTP 500 as terminal while retrying transient failures", async () => { + const png = Buffer.from("89504e470d0a1a0a0102", "hex"); + for (const [status, expectedAttempts, retryable] of [[500, 1, false], [503, 3, true]] as const) { + let attempts = 0; + const client = new OcrClient({ baseUrl: "http://ocr.internal:8000", token: "token", + fetch: (async () => { attempts += 1; return new Response(null, { status }); }) as typeof fetch, + sleep: async () => undefined }); + await assert.rejects(client.getReviewImage(jobId, 1, documentSha256), + (error: unknown) => error instanceof OcrClientError && error.status === status && error.retryable === retryable); + assert.equal(attempts, expectedAttempts); + } + const invalidSize = new OcrClient({ baseUrl: "http://ocr.internal:8000", token: "token", + fetch: (async () => new Response(png, { headers: { "content-type": "image/png", "content-length": String(png.length + 1), + "x-document-sha256": documentSha256, "x-content-sha256": sha256Hex(png), "x-page-number": "1" } })) as typeof fetch }); + await assert.rejects(invalidSize.getReviewImage(jobId, 1, documentSha256), /integrity validation failed/); +}); + +test("review image transfer is sequential and resumes from durable per-page checkpoints", async (context) => { + const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-ocr-image-transfer-")); + context.after(() => import("node:fs/promises").then(({ rm }) => rm(rootDirectory, { recursive: true, force: true }))); + const versionId = "12121212-1212-4121-8121-121212121212"; + const documentId = "doc:image-transfer"; + await stageOcrArtifacts({ rootDirectory, versionId, createdAt: "2026-09-21T10:00:00.000Z", documents: [{ + documentId, documentKey: "scan.pdf", bytes: document, requestedPages: [1, 2, 3], pages: [1, 2, 3].map((page) => ({ + page, text: "", rasterCoverage: 1, textSha256: sha256Hex("") + })) + }] }); + const calls: number[] = []; + let active = 0; + let maxActive = 0; + let interrupted = true; + const loadImage = async (page: number) => { + calls.push(page); + active += 1; + maxActive = Math.max(maxActive, active); + await Promise.resolve(); + active -= 1; + if (page === 2 && interrupted) throw new Error("transfer interrupted"); + const bytes = Buffer.from(`89504e470d0a1a0a${String(page).padStart(4, "0")}`, "hex"); + return { bytes, sha256: sha256Hex(bytes) }; + }; + + await assert.rejects(transferReviewImageArtifacts({ rootDirectory, versionId, documentId, loadImage }), /transfer interrupted/); + assert.deepEqual(calls, [1, 2]); + interrupted = false; + const transferred = await transferReviewImageArtifacts({ rootDirectory, versionId, documentId, loadImage }); + + assert.deepEqual(calls, [1, 2, 2, 3]); + assert.equal(maxActive, 1); + assert.deepEqual(transferred.images.map(({ page }) => page), [1, 2, 3]); + assert.deepEqual(await transferReviewImageArtifacts({ rootDirectory, versionId, documentId, loadImage }), transferred); + assert.deepEqual(calls, [1, 2, 2, 3]); + await writeFile(transferred.images[0]!.artifactPath, "corrupt", { mode: 0o600 }); + await assert.rejects(transferReviewImageArtifacts({ rootDirectory, versionId, documentId, loadImage }), /checkpoint integrity validation failed/); + assert.deepEqual(calls, [1, 2, 2, 3]); +}); + 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 })); }); diff --git a/tests/ocr/contracts-deploy.test.ts b/tests/ocr/contracts-deploy.test.ts index 2711423..86664fb 100644 --- a/tests/ocr/contracts-deploy.test.ts +++ b/tests/ocr/contracts-deploy.test.ts @@ -64,6 +64,8 @@ test("OpenAPI OCR schemas preserve progress, review evidence, corrections, and d }); test("OCR deployment defaults and container wiring match the private durable contract", async () => { + assert.equal(env.ragVersion, "0.2.0"); + assert.equal(env.buildRevision, "unknown"); assert.deepEqual({ root: env.ocrArtifactRoot, uploadBytes: (env as unknown as Record).ocrMaxUploadBytes, @@ -76,6 +78,8 @@ test("OCR deployment defaults and container wiring match the private durable con 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, /BUILD_REVISION=\$\{BUILD_REVISION\}/u); + assert.match(dockerfile, /org\.opencontainers\.image\.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);