diff --git a/docs/HISTORIAL_SESIONES.md b/docs/HISTORIAL_SESIONES.md index 1a4b2aa..6747beb 100644 --- a/docs/HISTORIAL_SESIONES.md +++ b/docs/HISTORIAL_SESIONES.md @@ -2,14 +2,61 @@ **Proyecto:** Workspace de tools IA para empresas **Modulo:** RAG -**Ultima actualizacion:** 2026-09-15 -**Ultima modificacion por:** Subagente Implement Unit 13 Local E2E +**Ultima actualizacion:** 2026-09-16 +**Ultima modificacion por:** Subagente Actualizacion Operativa OCR **Estado:** Activo --- ## Registro de sesion +### 2026-09-16 - Subagente Actualizacion Operativa OCR +**Agent:** Subagente Actualizacion Operativa OCR · **Model:** ollama/glm-5.3:cloud · **Session:** `ses_f5662ce40ffepxlOg2W0v9zrYv` (subagente de `ses_29bdbd003ffeLrLjUlFgnp08Y7`) +**Responsibility:** Corregir el estado operativo verificado no secreto de EasyPanel/RAG/OCR y fijar el orden seguro de despliegue, sin commit ni push (el agente principal los entrega despues). +**Work:** Actualizado `docs/OPERATIVA.md` con fecha 2026-09-16 y estado verificado: `KNOWLEDGE_LIFECYCLE_ENFORCED=true` (eliminada la redaccion obsoleta de `false`/pendiente de activar), `OCR_INGEST_ENABLED=false`, `OCR_SERVICE_URL` hacia el servicio privado `ocr-service`, token compartido OCR configurado sin revelar su valor, volumen duradero `/data/ingestions`, volumen transitorio `/data/jobs` y 29 variables del RAG verificadas sin duplicados. Sustituido el orden obsoleto por el orden seguro real: Git main, OCR primero con verificacion live/ready + trabajo real autenticado con resultado ligado por integridad y limpieza, RAG segundo con flag en false y verificacion nativa, solo entonces `OCR_INGEST_ENABLED=true` con nuevo deploy, y PDF de FacturaTech como candidato de revision no activado con `sourceRef` logico `Errores Junio 2026 - OCR verificado.md` (aprobacion humana obligatoria). Conservado el rollback de emergencia (OCR en false, deploy RAG, version activa intacta) y documentada la eliminacion transitoria de filas/PDF/resultados del OCR tras la transferencia duradera o automaticamente a las 24 horas. La advertencia de rotacion de credenciales expuestas se mantiene y no se afirma que se hayan rotado. +**Redaction:** Ningun secreto, token, URL con credenciales, huella, hash ni valor leido de `.env*` fue incluido; no se accedio a `.env*`. +**Validation:** Consistencia estructural del documento verificada; busqueda de redaccion obsoleta de lifecycle en `false` o de orden antiguo sin resultados; comprobacion de espacios limitada a los dos documentos editados. +**Files:** `docs/OPERATIVA.md` (corregido), `docs/HISTORIAL_SESIONES.md` (esta entrada). + +--- + +### 2026-09-15 - Subagente Document RAG Environment +**Agent:** Subagente Document RAG Environment · **Model:** ollama/glm-5.3:cloud · **Session:** `ses_f5aa57a88ffeYnnL4d6X1gXTiu` +**Responsibility:** Crear documentacion operativa segura del entorno RAG/EasyPanel, con redaccion estricta de secretos y plan de salida por etapas del OCR derivado del codigo real. +**Work:** Creado `docs/OPERATIVA.md` con lista rapida de EasyPanel, accion de rotacion de credenciales expuestas en chat (sin afirmar que ya se hayan rotado), configuracion actual no secreta (incluye `KNOWLEDGE_LIFECYCLE_ENFORCED=false` e `INGEST_WRITES_ENABLED=true`), nombres de variables secretas con la nota "Configurada en EasyPanel; nunca copiarla en documentacion", y salida por etapas del OCR: `OCR_INGEST_ENABLED=false` antes del deploy, requisitos previos (servicio OCR privado, `OCR_INTERNAL_TOKEN`, `OCR_SERVICE_URL`, volumen `/data/ingestions`), limites opcionales con sus valores por defecto del codigo (`OCR_MAX_UPLOAD_BYTES=52428800`, `OCR_MAX_PAGES=100`, `OCR_PAGE_TIMEOUT_MS=60000`, `OCR_TOTAL_TIMEOUT_MS=900000`) y orden seguro de activacion con desactivacion en caso de fallo conservando el corpus activo. Documentado que el deploy es manual (boton `Deploy` de EasyPanel) y que no existe accion automatizada documentada. +**Redaction:** Ningun valor secreto, parcial, hash ni huella fue copiado, mostrado o persistido. No se accedio a `.env*`, `llaves`, `backups/` ni `docs.zip`. +**Validation:** Verificacion estructural del documento creado y comprobacion de espacios en Markdown. Inspeccion del `git diff` confirmo que no contiene ningun valor de las variables secretas. +**Files:** `docs/OPERATIVA.md` (nuevo), `docs/HISTORIAL_SESIONES.md` (esta entrada). + +--- + +### 2026-09-16 - Subagente Implementacion Runtime OCR +**Agent:** Subagente Implementacion Runtime OCR · **Model:** openai/gpt-5.6-sol · **Session:** `ses_f5692c07cffeBjuFNjshhDRDnL` +**Responsibility:** Implement and prove Unit 14a OCR job execution and authenticated result retrieval only. +**Work:** Persisted accepted PDFs/results, resumed queued jobs safely after restart, executed the existing renderer with an injectable engine, recorded terminal status/errors, and exposed authenticated result retrieval without changing RAG ingestion or cleanup. +**Validation:** Strict-TDD RED 8/10 then recovery RED 9/10; focused GREEN 10/10; OCR 15/15; Node 78/78; localhost Uvicorn upload/status/result `202/200/200`; static, cleanup, and process gates passed. +**Files:** `ocr-service/app/main.py`, `ocr-service/tests/test_api.py`, OpenSpec apply progress, and this history. Task 7.4 remains pending. + +--- + +### 2026-09-16 - Subagente Implementacion Runtime OCR - Unit 14b +**Agent:** Subagente Implementacion Runtime OCR · **Model:** openai/gpt-5.6-sol · **Session:** `ses_f5692c07cffeBjuFNjshhDRDnL` +**Responsibility:** Implement and prove transient OCR cleanup after durable transfer and automatic 24-hour expiry only. +**Work:** Added healthcheck-driven OCR TTL expiry, authenticated client deletion after durable result acceptance, failure retention, and best-effort cleanup fallback without touching durable RAG artifacts or active corpus. +**Validation:** Strict-TDD RED Python 10/11 and Node 15/17; GREEN Python API 11/11, OCR 16/16, focused Node 17/17, canonical Node 80/80; check/build/runtime/cleanup gates passed. +**Files:** OCR API/client/dispatcher tests and code, OpenSpec apply progress, and this history. Task 7.4 remains pending. + +--- + +### 2026-09-16 - Subagente Implementacion Runtime OCR - Unit 14c +**Agent:** Subagente Implementacion Runtime OCR · **Model:** openai/gpt-5.6-sol · **Session:** `ses_f5692c07cffeBjuFNjshhDRDnL` +**Responsibility:** Implement and prove optional logical multipart upload source identity without production or active-content changes. +**Work:** Added trimmed optional `sourceRef`, original-filename fallback, blank validation, OpenAPI contract coverage, and local identity-reuse evidence independent from physical upload names. +**Validation:** Strict-TDD RED 2/4; GREEN focused 4/4, localhost multipart 1/1, canonical Node 81/81; check/build/whitespace/temp cleanup passed. +**Files:** `src/app.ts`, `src/api/openapi.ts`, focused upload/contract tests, OpenSpec apply progress, and this history. Task 7.4 remains pending. + +--- + ### 2026-09-15 - Subagente Implement Unit 13 Local E2E **Agent:** Subagente Implement Unit 13 Local E2E · **Model:** openai/gpt-5.6-sol · **Session:** `ses_f5af65977ffehYI4eK2l177YAL` **Responsibility:** Implement only OCR Unit 13 tasks 7.1–7.3 under strict TDD, preserving Units 8–12 and excluding production acceptance task 7.4 and Git delivery work. diff --git a/docs/OPERATIVA.md b/docs/OPERATIVA.md new file mode 100644 index 0000000..dfad7fa --- /dev/null +++ b/docs/OPERATIVA.md @@ -0,0 +1,157 @@ +# Operativa del servicio RAG + +**Modulo:** RAG +**Ultima actualizacion:** 2026-09-16 +**Version:** 1.1 + +--- + +Este documento registra los hechos operativos del servicio RAG: la configuracion vigente en EasyPanel y la salida por etapas del OCR. No contiene secretos: las credenciales reales viven unicamente en EasyPanel. + +## Estado verificado (2026-09-16) + +- `KNOWLEDGE_LIFECYCLE_ENFORCED=true`: el ciclo de vida del conocimiento esta activo; no queda pendiente ninguna activacion. +- `OCR_INGEST_ENABLED=false`: la ingesta sigue el flujo nativo, sin OCR. +- `OCR_SERVICE_URL`: configurada hacia el servicio privado `ocr-service` de la red interna. +- `OCR_INTERNAL_TOKEN`: secreto compartido configurado en EasyPanel (RAG y servicio OCR); su valor nunca se documenta. +- Volumen duradero del RAG en `/data/ingestions`: conserva originales y artefactos de revision. +- Volumen transitorio del OCR en `/data/jobs`: cola SQLite y trabajos en curso. +- Verificadas exactamente 29 variables de entorno del RAG en EasyPanel, todas unicas y sin nombres duplicados. + +## Lista rapida en EasyPanel + +1. Rotar las credenciales expuestas (ver seccion siguiente); la rotacion sigue pendiente. +2. Comprobar que las variables coinciden con la tabla de configuracion actual. +3. Mantener `OCR_INGEST_ENABLED=false` hasta completar el orden seguro de salida del OCR. +4. Pulsar `Deploy` en EasyPanel despues de cada cambio principal. +5. Verificar `GET /health` tras cada deploy. + +## Accion de seguridad urgente: rotar credenciales + +Credenciales reales quedaron expuestas en una conversacion de chat. Hay que rotarlas: + +- La credencial de OpenRouter (variable `EMBEDDING_API_KEY`; y `ANSWER_API_KEY` si tiene valor propio). +- La contrasena y la URL de PostgreSQL (variable `POSTGRES_URL`). +- El token administrativo del ciclo de vida (variable `LIFECYCLE_ADMIN_TOKEN`). + +Como hacerlo: + +- Generar valores nuevos y actualizarlos unicamente en EasyPanel. +- Este documento no afirma que la rotacion ya se haya realizado. +- Nunca pegar valores de credenciales en chat, documentacion, historial ni ficheros versionados. + +## Como se despliega + +- Tras cada cambio principal (variables de entorno o codigo nuevo subido por Git), el usuario pulsa el boton `Deploy` de EasyPanel. +- No existe ninguna accion automatizada documentada sobre EasyPanel: ni scripts, ni comandos de panel, ni webhooks. El deploy siempre lo dispara una persona desde el panel. +- Verificacion minima tras cada deploy: `GET /health` del RAG. + +## Configuracion actual (valores no secretos) + +Valores vigentes confirmados en EasyPanel (2026-09-16): + +| Variable | Valor actual | +|---|---| +| `NODE_ENV` | `production` | +| `PORT` | `80` | +| `QDRANT_URL` | `http://qdrant:6333` | +| `QDRANT_API_KEY` | vacia (sin valor) | +| `QDRANT_COLLECTION` | `rag_chunks` | +| `EMBEDDING_PROVIDER` | `openrouter` | +| `EMBEDDING_MODEL` | `qwen/qwen3-embedding-8b` | +| `EMBEDDING_BASE_URL` | `https://openrouter.ai/api/v1` | +| `ANSWER_PROVIDER` | `openrouter` | +| `ANSWER_MODEL` | `openai/gpt-4.1-mini` | +| `ANSWER_BASE_URL` | `https://openrouter.ai/api/v1` | +| `POSTGRES_SSL` | `false` | +| `KNOWLEDGE_LIFECYCLE_ENFORCED` | `true` | +| `INGEST_WRITES_ENABLED` | `true` | +| `OCR_INGEST_ENABLED` | `false` | +| `OCR_SERVICE_URL` | configurada hacia el servicio privado `ocr-service` | + +Notas de comportamiento (derivadas del codigo, sin secretos): + +- Si `ANSWER_API_KEY` no se define, el codigo reutiliza `EMBEDDING_API_KEY`. +- Si `POSTGRES_URL` no se define, el codigo prueba `DATABASE_URL`. +- Si `OCR_INGEST_ENABLED` no existe, el codigo tambien la toma como `false`. +- Un cambio en `KNOWLEDGE_LIFECYCLE_ENFORCED`, `INGEST_WRITES_ENABLED` o `OCR_INGEST_ENABLED` se hace siempre en EasyPanel y se registra en este documento. + +## Variables con secretos + +Estas variables contienen credenciales reales. Aqui solo se registran sus nombres: + +| Variable | Estado | +|---|---| +| `EMBEDDING_API_KEY` | Configurada en EasyPanel; nunca copiarla en documentacion. | +| `ANSWER_API_KEY` | Configurada en EasyPanel; nunca copiarla en documentacion. | +| `POSTGRES_URL` | Configurada en EasyPanel; nunca copiarla en documentacion. | +| `LIFECYCLE_ADMIN_TOKEN` | Configurada en EasyPanel; nunca copiarla en documentacion. | +| `OCR_INTERNAL_TOKEN` | Configurada en EasyPanel, compartida con el servicio OCR y enviada como `Authorization: Bearer`; nunca copiarla en documentacion. | + +## Salida por etapas del OCR + +El OCR es un servicio privado e independiente. El RAG solo lo llama si `OCR_INGEST_ENABLED` esta activo; con el flag apagado, la ingesta sigue el flujo nativo actual, sin cambios. Hoy el flag esta en `false` en EasyPanel. + +### Servicio OCR privado (referencia de configuracion) + +- Contenedor `ocr-service`, imagen Python CPU con PaddleOCR `3.4.0` y PaddlePaddle `3.2.2`. +- Una sola replica y un solo worker; limites del servicio: maximo 3 CPU y 5 GiB de RAM. +- Modelos descargados durante el build, nunca en el arranque. +- 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. + +### Volumenes + +| Volumen | Ruta | Papel | +|---|---|---| +| Duradero del RAG | `/data/ingestions` | Originales y artefactos de revision conservados. | +| Transitorio del OCR | `/data/jobs` | Cola SQLite, PDFs y resultados en curso. | + +### 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. +- La limpieza transitoria nunca toca los artefactos duraderos del RAG ni el corpus activo. + +### Orden seguro de salida (cuando cambian ambos servicios) + +1. Publicar el codigo en Git `main` (autorizacion ya concedida de forma separada). +2. Desplegar primero el servicio OCR. +3. Verificar el OCR: `GET /health/live` y `GET /health/ready` responden; un trabajo real autenticado llega a `succeeded`, devuelve un resultado ligado por integridad y admite limpieza transitoria. +4. Desplegar el RAG en segundo lugar manteniendo `OCR_INGEST_ENABLED=false`. +5. Verificar la ruta nativa del RAG: `/health`, ingesta textual, retrieval y corpus activo protegido, sin regresiones. +6. Solo entonces activar el OCR: `OCR_INGEST_ENABLED=true` y volver a pulsar `Deploy` en el RAG. +7. Presentar el PDF de FacturaTech como candidato de revision no activado, con el `sourceRef` logico `Errores Junio 2026 - OCR verificado.md`. La aprobacion humana sigue siendo obligatoria antes de cualquier activacion. + +### Rollback de emergencia + +- Ante cualquier fallo: poner `OCR_INGEST_ENABLED=false`, pulsar `Deploy` en el RAG y conservar la version activa actual del corpus. +- El sistema es fail-closed: un fallo del OCR deja intacta la version activa anterior; no hay activacion parcial. +- Las paginas con OCR no se activan solas; requieren revision humana obligatoria (estado `review_required`). +- Desactivar el OCR no borra el corpus activo ni exige reingesta. + +### Limites opcionales del RAG (valores por defecto del codigo) + +| Variable | Valor por defecto | +|---|---| +| `OCR_MAX_UPLOAD_BYTES` | `52428800` (50 MiB) | +| `OCR_MAX_PAGES` | `100` | +| `OCR_PAGE_TIMEOUT_MS` | `60000` | +| `OCR_TOTAL_TIMEOUT_MS` | `900000` (15 min) | +| `LIFECYCLE_RECONCILE_INTERVAL_MS` | `300000` | +| `LIFECYCLE_INDEXING_STALE_TIMEOUT_MS` | `1800000` | +| `QDRANT_LOGS_COLLECTION` | `rag_eval_logs` | + +## Verificaciones habituales + +- `GET /health` del RAG: estado general con PostgreSQL, Qdrant y reconciliador. +- `GET /health/live` y `GET /health/ready` del servicio OCR: solo accesibles desde la red interna del proyecto. +- Tras cambiar cualquier variable: pulsar `Deploy` y verificar `/health`. + +## Reglas de uso de este documento + +- Consultarlo antes de tareas con servicios externos, credenciales, despliegues o ejecuciones recurrentes. +- Nunca escribir secretos aqui: solo nombres de variables y donde viven (EasyPanel). +- Si un hecho operativo deja de ser valido, actualizarlo o marcarlo como obsoleto en la misma intervencion en que se descubra. \ No newline at end of file diff --git a/ocr-service/app/main.py b/ocr-service/app/main.py index 08500f5..11ede49 100644 --- a/ocr-service/app/main.py +++ b/ocr-service/app/main.py @@ -5,18 +5,21 @@ import os import sqlite3 import threading import uuid -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone from pathlib import Path -from typing import Annotated, Any +from typing import Annotated, Any, Callable -from fastapi import Depends, FastAPI, File, Form, Header, HTTPException, Response, UploadFile +from fastapi import BackgroundTasks, Depends, FastAPI, File, Form, Header, HTTPException, Response, UploadFile +from .engine import OcrEngine from .models import load_runtime_engine +from .render import process_pdf MAX_UPLOAD_BYTES = 50 * 1024 * 1024 MAX_PAGES = 100 QUEUE_CAPACITY = 3 +TRANSIENT_TTL = timedelta(hours=24) ALLOWED_CONFIG = { "languages": ["es", "en"], "dpi": 200, @@ -39,14 +42,20 @@ def page_hash(pages: list[int]) -> str: class JobQueue: - def __init__(self, path: str | Path): + 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)" + "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: @@ -67,45 +76,78 @@ class JobQueue: "createdAt": row["created_at"], } - def submit(self, key: str, request: dict[str, Any]) -> dict[str, Any]: + 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) + 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", datetime.now(timezone.utc).isoformat(), + json.dumps(request["pages"]), "queued", self.now().isoformat(), pdf, None, None, ) - self.connection.execute("INSERT INTO jobs VALUES (?,?,?,?,?,?,?)", values) + 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() - return self.ack(self.connection.execute("SELECT * FROM jobs WHERE job_id=?", (values[0],)).fetchone()) 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 - return {"jobId": job_id, "status": row["status"], "completedPages": 0, - "totalPages": len(json.loads(row["pages"])), "error": 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 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) + def create_app( token: str, db_path: str | Path = ":memory:", max_upload_bytes: int = MAX_UPLOAD_BYTES, engine_ready: bool = False, + 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) + queue = JobQueue(db_path, now) def authorize(authorization: Annotated[str | None, Header()] = None) -> None: scheme, _, supplied = (authorization or "").partition(" ") @@ -118,13 +160,14 @@ def create_app( @application.get("/health/ready") def ready(response: Response) -> dict[str, int | bool]: + queue.sweep_expired() if not engine_ready: response.status_code = 503 return {"ready": engine_ready, "queueDepth": queue.depth(), "queueCapacity": QUEUE_CAPACITY, "concurrency": 1} @application.post("/v1/jobs", status_code=202, dependencies=[Depends(authorize)]) async def create_job( - file: Annotated[UploadFile, File()], request: Annotated[str, Form()], + background_tasks: BackgroundTasks, 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) @@ -153,7 +196,10 @@ 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") - return queue.submit(idempotency_key or "", payload) + 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) + return ack @application.get("/v1/jobs/{job_id}", dependencies=[Depends(authorize)]) def get_job(job_id: str) -> dict[str, Any]: @@ -162,6 +208,15 @@ def create_app( fail(404, "JOB_NOT_FOUND", "OCR job does not exist") return status + @application.get("/v1/jobs/{job_id}/result", dependencies=[Depends(authorize)]) + def get_result(job_id: str) -> dict[str, Any]: + stored = queue.result(job_id) + if stored is None: + fail(404, "JOB_NOT_FOUND", "OCR job does not exist") + if stored[0] != "succeeded" or stored[1] is None: + fail(409, "RESULT_NOT_READY", "OCR job result is not ready") + return stored[1] + @application.delete("/v1/jobs/{job_id}", status_code=204, dependencies=[Depends(authorize)]) def delete_job(job_id: str) -> Response: queue.delete(job_id) @@ -176,4 +231,5 @@ app = create_app( os.getenv("OCR_INTERNAL_TOKEN", ""), os.getenv("OCR_JOBS_DB", ":memory:"), engine_ready=runtime_engine is not None, + engine=runtime_engine, ) diff --git a/ocr-service/tests/test_api.py b/ocr-service/tests/test_api.py index 50ee740..1fe7d7f 100644 --- a/ocr-service/tests/test_api.py +++ b/ocr-service/tests/test_api.py @@ -1,6 +1,8 @@ import hashlib import json +import sqlite3 import sys +from datetime import datetime, timedelta, timezone from pathlib import Path import pytest @@ -9,10 +11,17 @@ from fastapi.testclient import TestClient sys.path.insert(0, str(Path(__file__).parents[1])) from app.main import create_app +from app.engine import EngineLine TOKEN = "unit-5-test-token" PDF = b"%PDF-1.4\nunit five\n%%EOF" +FIXTURE = Path(__file__).parents[2] / "tests" / "fixtures" / "ocr" / "native-three-pages.pdf" + + +class FakeEngine: + def recognize(self, _image: object) -> list[EngineLine]: + return [EngineLine("Factura FAT07", 0.97, (10, 20, 160, 50))] def request_for(pdf: bytes = PDF, pages: list[int] | None = None) -> dict: @@ -129,5 +138,51 @@ def test_auth_queue_pressure_is_retryable_and_status_and_delete_are_authenticate job_id = accepted[0].json()["jobId"] status = client.get(f"/v1/jobs/{job_id}", headers={"Authorization": f"Bearer {TOKEN}"}) assert status.json() == {"jobId": job_id, "status": "queued", "completedPages": 0, "totalPages": 2, "error": None} + assert client.delete(f"/v1/jobs/{job_id}").status_code == 401 assert client.delete(f"/v1/jobs/{job_id}", headers={"Authorization": f"Bearer {TOKEN}"}).status_code == 204 assert client.delete(f"/v1/jobs/{job_id}", headers={"Authorization": f"Bearer {TOKEN}"}).status_code == 204 + + +def test_auth_health_sweeps_only_jobs_older_than_24_hours(tmp_path: Path): + current = [datetime(2026, 9, 16, tzinfo=timezone.utc)] + db_path = tmp_path / "expiry.db" + client = TestClient(create_app(TOKEN, db_path, engine_ready=True, now=lambda: current[0])) + job_id = submit(client, request_for()).json()["jobId"] + headers = {"Authorization": f"Bearer {TOKEN}"} + + current[0] += timedelta(hours=23, minutes=59) + assert client.get("/health/ready").status_code == 200 + 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 + 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 + + +def test_auth_accepted_job_executes_and_exposes_integrity_bound_result(tmp_path: Path): + pdf = FIXTURE.read_bytes() + db_path = tmp_path / "executed.db" + with sqlite3.connect(db_path) as connection: + 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"] + + 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")] + + +def test_auth_result_rejects_unknown_not_ready_and_unauthorized_jobs(client: TestClient): + job_id = submit(client, request_for()).json()["jobId"] + headers = {"Authorization": f"Bearer {TOKEN}"} + assert client.get(f"/v1/jobs/{job_id}/result").status_code == 401 + assert client.get("/v1/jobs/missing/result", headers=headers).json()["detail"]["code"] == "JOB_NOT_FOUND" + pending = client.get(f"/v1/jobs/{job_id}/result", headers=headers) + assert (pending.status_code, pending.json()["detail"]["code"]) == (409, "RESULT_NOT_READY") diff --git a/openspec/changes/ocr-ingest-integration/apply-progress.md b/openspec/changes/ocr-ingest-integration/apply-progress.md index 476b49b..ac19e6e 100644 --- a/openspec/changes/ocr-ingest-integration/apply-progress.md +++ b/openspec/changes/ocr-ingest-integration/apply-progress.md @@ -426,3 +426,48 @@ None — the migration follows the proposal, specifications, design, and closed - Tracked and complete intended-untracked whitespace checks passed. No matching `tsx --test`, E2E Node, or OCR Uvicorn process remained after the corrected self-excluding process check. - Intended untracked inventory: `src/modules/ocr/{dispatcher,indexing,retention,review}.ts`; `tests/ocr/{contracts-deploy,dispatcher,e2e,retention,review}.test.ts`. - Tasks 7.1–7.3 are complete. Task 7.4 remains unchecked and untouched; no production flag, service, database, vector store, embedding provider, OCR endpoint, commit, push, or PR was used. + +## Unit 14a: OCR Job Execution and Result +- Persisted accepted PDF bytes/results in SQLite, resumed queued jobs on idempotent resubmission, executed the existing renderer through an injected engine, recorded terminal status/errors, and added authenticated result retrieval; task 7.4 remains pending. + +### TDD Cycle Evidence +| Task | Test File | Safety Net | RED | GREEN / TRIANGULATE | REFACTOR | +|---|---|---|---|---|---| +| Unit 14a | `ocr-service/tests/test_api.py` | 13/13 | 8/10 missing execution/result; recovery triangulation 9/10 remained queued | 10/10; success, restart recovery, unauthorized, missing, not-ready | Guarded duplicate workers; 10/10 remained green | + +### Work Unit Evidence +| Evidence | Result | +|---|---| +| Focused/regression | API 10/10; OCR 15/15; Node 78/78; `npm run check` and whitespace clean. | +| Runtime harness | Local Uvicorn with deterministic engine returned upload `202`, status `200/succeeded`, result `200`, matching job/document identity; process/temp cleanup verified. | +| Rollback boundary | Revert `ocr-service/app/main.py`, `ocr-service/tests/test_api.py`, and this Unit 14a/history metadata; Units 1–13 remain intact. | + +## Unit 14b: OCR Transient Cleanup +- Added authenticated idempotent client deletion after durable RAG result acceptance and automatic 24-hour OCR row/input/result expiry through the existing 30-second readiness healthcheck; cleanup failures fall back to TTL without failing durable candidates. + +### TDD Cycle Evidence +| Task | Test Files | Safety Net | RED | GREEN / TRIANGULATE | REFACTOR | +|---|---|---|---|---|---| +| Unit 14b | `ocr-service/tests/test_api.py`; `tests/ocr/{client,dispatcher}.test.ts` | Python 10/10; Node 15/15 | Python 10/11 lacked clock/sweep; Node 15/17 lacked delete client/dispatch | Python 11/11; Node 17/17; pre-persistence failure retained remote state and cleanup failure preserved local success | Reused healthcheck and generic retry/auth boundaries; focused suites remained green | + +### Work Unit Evidence +| Evidence | Result | +|---|---| +| Focused/regression | Python API 11/11; OCR 16/16; focused Node 17/17; canonical Node 80/80; check/build/whitespace clean. | +| Runtime harness | Local Uvicorn returned pre-expiry `200`, post-24h `404`, unauthorized delete `401`, idempotent authenticated deletes `204/204`, and post-delete `404`; process/temp cleanup verified. | +| Rollback boundary | Revert only Unit 14b deltas in OCR API/client/dispatcher/tests and this progress/history metadata; durable RAG artifacts, active corpus, and Unit 14a remain intact. | + +## Unit 14c: Logical Upload Source Reference +- Added an optional trimmed multipart `sourceRef` that remains independent from the temporary physical upload path, preserves original-filename fallback, rejects blank values, and documents the contract without activating content. + +### TDD Cycle Evidence +| Task | Test Files | Safety Net | RED | GREEN / TRIANGULATE | REFACTOR | +|---|---|---|---|---|---| +| Unit 14c | `tests/ocr/{upload-source-ref,contracts-deploy}.test.ts` | Existing upload/contracts 13/13 | 2/4; missing contract and blank reference accepted | 4/4; provided/trimmed reuse, physical-name independence, fallback, blank rejection | Shared one parsed value across file/ZIP branches; focused suite remained green | + +### Work Unit Evidence +| Evidence | Result | +|---|---| +| Focused/regression | Focused 4/4; multipart runtime 1/1; canonical Node 81/81; check/build/whitespace clean. | +| Runtime harness | Ephemeral localhost Express multipart uploads reused `Errores Junio 2026 - OCR verificado.md` across two physical PDF names, preserved `fallback.pdf` when omitted, rejected blank input `400`, kept activation false, and removed temporary files. | +| Rollback boundary | Revert only Unit 14c deltas in `src/{app,api/openapi}.ts`, its two tests, and this progress/history metadata; Units 14a–14b, active content, and corpus remain intact. | diff --git a/src/api/openapi.ts b/src/api/openapi.ts index f809cd5..3f8f094 100644 --- a/src/api/openapi.ts +++ b/src/api/openapi.ts @@ -529,6 +529,7 @@ export const openApiDocument = { properties: { file: { type: "string", format: "binary" }, sourceId: { type: "string" }, + sourceRef: { type: "string", minLength: 1, description: "Optional logical source reference; trimmed and defaults to the uploaded filename." }, mode: { type: "string", enum: ["mechanical", "interactive"], default: "mechanical" }, tags: { type: "string", description: "Comma-separated tags." }, isZipFolder: { type: "string", enum: ["true", "false"], default: "false" }, diff --git a/src/app.ts b/src/app.ts index bf71d9b..1bd7778 100644 --- a/src/app.ts +++ b/src/app.ts @@ -486,6 +486,11 @@ export function createApp(options: AppOptions = {}) { const tags = typeof req.body.tags === "string" ? req.body.tags.split(",").map((entry: string) => entry.trim()).filter(Boolean) : []; + const sourceRef = typeof req.body.sourceRef === "string" ? req.body.sourceRef.trim() : undefined; + if (req.body.sourceRef !== undefined && !sourceRef) { + res.status(400).json({ ok: false, error: "sourceRef must be a non-empty string when provided" }); + return; + } let result; @@ -497,7 +502,7 @@ export function createApp(options: AppOptions = {}) { result = await ingestService.ingest({ sourceId: req.body.sourceId ? String(req.body.sourceId) : undefined, sourceType: "folder", - sourceRef: req.file.originalname.replace(/\.zip$/i, ""), + sourceRef: sourceRef ?? req.file.originalname.replace(/\.zip$/i, ""), readPath: extractDirPath, mode: req.body.mode === "interactive" ? "interactive" : "mechanical", tags, @@ -508,7 +513,7 @@ export function createApp(options: AppOptions = {}) { result = await ingestService.ingest({ sourceId: req.body.sourceId ? String(req.body.sourceId) : undefined, sourceType: "file", - sourceRef: req.file.originalname, + sourceRef: sourceRef ?? req.file.originalname, readPath: tempFilePath, mode: req.body.mode === "interactive" ? "interactive" : "mechanical", tags, diff --git a/src/modules/ocr/client.ts b/src/modules/ocr/client.ts index b63c93d..cd6b66b 100644 --- a/src/modules/ocr/client.ts +++ b/src/modules/ocr/client.ts @@ -119,6 +119,10 @@ export class OcrClient { return value; } + async delete(jobId: string): Promise { + await this.requestJson(`/v1/jobs/${encodeURIComponent(jobId)}`, () => ({ method: "DELETE", headers: this.headers() })); + } + private headers(additional: Record = {}): Headers { return new Headers({ Authorization: `Bearer ${this.options.token}`, ...additional }); } @@ -132,7 +136,7 @@ export class OcrClient { if (attempt < 2) { await this.sleep(2_000 * 2 ** attempt); continue; } throw new OcrClientError("OCR_NETWORK_ERROR", undefined, false); } - if (response.ok) return response.json(); + if (response.ok) return response.status === 204 ? null : response.json(); const body = await response.json().catch(() => ({})) as Record; const detail = isObject(body.detail) ? body.detail : body; const code = typeof detail.code === "string" ? detail.code : `OCR_HTTP_${response.status}`; diff --git a/src/modules/ocr/dispatcher.ts b/src/modules/ocr/dispatcher.ts index 01fe743..1643998 100644 --- a/src/modules/ocr/dispatcher.ts +++ b/src/modules/ocr/dispatcher.ts @@ -64,6 +64,7 @@ export class OcrDispatcher { }); const versionComplete = await this.store.completeOcrJob(job.jobId, result); if (versionComplete) await this.store.markReviewRequired(job.versionId); + await this.client.delete(remoteJobId).catch(() => undefined); return "succeeded"; } catch (error) { const detail = error instanceof Error ? error.message : "Unknown OCR dispatch failure"; diff --git a/tests/ocr/client.test.ts b/tests/ocr/client.test.ts index a47785a..37f7a08 100644 --- a/tests/ocr/client.test.ts +++ b/tests/ocr/client.test.ts @@ -155,6 +155,21 @@ test("polling uses bounded exponential backoff until a terminal status", async ( assert.deepEqual(delays, [2_000, 4_000]); }); +test("authenticated delete accepts an empty idempotent response", async () => { + let request: { url: string; method: string; authorization: string | null } | undefined; + const client = new OcrClient({ + baseUrl: "http://ocr.internal:8000", + token: "delete-token", + fetch: (async (input, init) => { + request = { url: String(input), method: init?.method ?? "GET", authorization: new Headers(init?.headers).get("Authorization") }; + return new Response(null, { status: 204 }); + }) as typeof fetch + }); + + await client.delete("job / 1"); + assert.deepEqual(request, { url: "http://ocr.internal:8000/v1/jobs/job%20%2F%201", method: "DELETE", authorization: "Bearer delete-token" }); +}); + 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 e735184..8a5866d 100644 --- a/tests/ocr/contracts-deploy.test.ts +++ b/tests/ocr/contracts-deploy.test.ts @@ -14,6 +14,11 @@ test("OpenAPI describes authenticated OCR ingestion, status, review, approval, a const uploadResponses = api.paths["/ingest/upload"]!.post!.responses as typeof ingestResponses; assert.deepEqual(ingestResponses["202"]!.content!["application/json"]!.schema, { oneOf: [{ $ref: "#/components/schemas/OcrAccepted" }, { $ref: "#/components/schemas/IngestResponse" }] }); assert.deepEqual(uploadResponses["202"]!.content!["application/json"]!.schema, { oneOf: [{ $ref: "#/components/schemas/OcrUploadAccepted" }, { $ref: "#/components/schemas/UploadIngestResponse" }] }); + const uploadRequest = api.components.schemas.UploadIngestRequest!; + assert.deepEqual(uploadRequest.required, ["file"]); + assert.deepEqual((uploadRequest.properties as Record).sourceRef, { + type: "string", minLength: 1, description: "Optional logical source reference; trimmed and defaults to the uploaded filename." + }); for (const [path, method] of [ ["/ingestions/{versionId}", "get"], diff --git a/tests/ocr/dispatcher.test.ts b/tests/ocr/dispatcher.test.ts index 2ce6a54..65594d1 100644 --- a/tests/ocr/dispatcher.test.ts +++ b/tests/ocr/dispatcher.test.ts @@ -260,12 +260,31 @@ test("dispatcher recovers only expired work, reuses its remote job, and complete const client = { async submit() { calls.push("submit"); throw new Error("existing remote work must not be duplicated"); }, async getStatus() { calls.push("status"); return { jobId: "remote-1", status: "succeeded", completedPages: 1, totalPages: 1, error: null }; }, - async getResult() { calls.push("result"); return { pages: [{ page: 1 }] }; } + async getResult() { calls.push("result"); return { pages: [{ page: 1 }] }; }, + async delete() { calls.push("delete"); throw new Error("cleanup unavailable"); } }; const dispatcher = new OcrDispatcher(repository, client as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") })); assert.equal(await dispatcher.recoverExpiredLeases(), 1); - assert.deepEqual(calls, ["recover", "claim-exact", "status", "result", "complete", "review-required"]); + assert.deepEqual(calls, ["recover", "claim-exact", "status", "result", "complete", "review-required", "delete"]); +}); + +test("dispatcher retains remote OCR state when durable result transfer fails", async () => { + const calls: string[] = []; + const job = { jobId: "job-1", versionId: "version-1", documentId: "document-1", remoteJobId: "remote-1", remoteIdempotencyKey: "key", state: "running" as const, requestedPages: [1], completedPages: 0, configVersion: "ocr-v1", attemptCount: 1, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt: null, errorCode: null, errorDetail: null }; + const store = { + async claimNextOcrJob() { return job; }, + async completeOcrJob() { calls.push("persist"); throw new Error("durable write failed"); }, + async failOcrJob() { calls.push("job-failed"); }, async markFailed() { calls.push("version-failed"); } + }; + const client = { + async getStatus() { return { status: "succeeded" }; }, async getResult() { return { pages: [] }; }, + async delete() { calls.push("delete"); } + }; + const dispatcher = new OcrDispatcher(store as never, client as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") })); + + assert.equal(await dispatcher.runOnce(), "failed"); + assert.deepEqual(calls, ["persist", "job-failed", "version-failed"]); }); test("OCR exhaustion fails the leased candidate without activation or engine substitution", async () => { diff --git a/tests/ocr/upload-source-ref.test.ts b/tests/ocr/upload-source-ref.test.ts new file mode 100644 index 0000000..1d494af --- /dev/null +++ b/tests/ocr/upload-source-ref.test.ts @@ -0,0 +1,55 @@ +import assert from "node:assert/strict"; +import { access } from "node:fs/promises"; +import test from "node:test"; +import { createApp } from "../../src/app.js"; + +const logicalSourceRef = "Errores Junio 2026 - OCR verificado.md"; + +test("multipart upload separates logical source identity from physical filenames", async (context) => { + const inputs: Array> = []; + const identities = new Map(); + const activeVersionId = "active-before-upload"; + const ingestService = { + async ingest(input: Record) { + inputs.push(input); + const identity = `${input.sourceId}:${input.sourceRef}`; + const versionId = identities.get(identity) ?? `version-${identities.size + 1}`; + identities.set(identity, versionId); + return { accepted: true, sourceId: input.sourceId, versionId, versionNumber: 2, state: "indexing", phase: "ocr_queued", statusUrl: `/ingestions/${versionId}`, reviewUrl: null, activated: false }; + }, + async cleanup() { return { deleted: 0 }; } + }; + const server = createApp({ ingestService: ingestService as never, startReconciler: false }).listen(0); + context.after(() => server.close()); + const address = server.address(); + assert.ok(address && typeof address === "object"); + + const upload = async (filename: string, sourceRef?: string) => { + const form = new FormData(); + form.set("file", new Blob(["%PDF-1.4\n%%EOF"], { type: "application/pdf" }), filename); + form.set("sourceId", "facturatech:agente-whatsapp-central:errores-junio-2026:v1"); + form.set("activate", "false"); + if (sourceRef !== undefined) form.set("sourceRef", sourceRef); + return fetch(`http://127.0.0.1:${address.port}/ingest/upload`, { method: "POST", body: form }); + }; + + const first = await upload("factura-fisica.pdf", ` ${logicalSourceRef} `); + const second = await upload("otra-copia.pdf", logicalSourceRef); + const fallback = await upload("fallback.pdf"); + const invalid = await upload("invalid.pdf", " "); + const [firstBody, secondBody] = await Promise.all([first.json(), second.json()]) as Array<{ versionId: string; activated: boolean }>; + + assert.equal(first.status, 202); + assert.equal(second.status, 202); + assert.equal(fallback.status, 202); + assert.equal(invalid.status, 400); + assert.equal(firstBody.versionId, secondBody.versionId); + assert.equal(firstBody.activated, false); + assert.deepEqual(inputs.map(({ sourceRef }) => sourceRef), [logicalSourceRef, logicalSourceRef, "fallback.pdf"]); + assert.deepEqual(inputs.map(({ activate }) => activate), [false, false, false]); + assert.equal(activeVersionId, "active-before-upload"); + for (const input of inputs) { + assert.notEqual(input.readPath, input.sourceRef); + await assert.rejects(access(String(input.readPath))); + } +});