From ac046e37c18e5691b7811e0fc91d217a5af5a355 Mon Sep 17 00:00:00 2001 From: Paco POR-CORREO Date: Mon, 14 Sep 2026 23:24:29 +0200 Subject: [PATCH] feat(ocr): add client and artifacts --- docs/HISTORIAL_SESIONES.md | 30 ++- .../ocr-ingest-integration/apply-progress.md | 53 ++++- .../changes/ocr-ingest-integration/tasks.md | 4 +- src/modules/ocr/artifacts.ts | 129 +++++++++++ src/modules/ocr/client.ts | 174 +++++++++++++++ tests/ocr/client.test.ts | 203 ++++++++++++++++++ 6 files changed, 587 insertions(+), 6 deletions(-) create mode 100644 src/modules/ocr/artifacts.ts create mode 100644 src/modules/ocr/client.ts create mode 100644 tests/ocr/client.test.ts diff --git a/docs/HISTORIAL_SESIONES.md b/docs/HISTORIAL_SESIONES.md index c0f7379..68c3a86 100644 --- a/docs/HISTORIAL_SESIONES.md +++ b/docs/HISTORIAL_SESIONES.md @@ -3,7 +3,7 @@ **Proyecto:** Workspace de tools IA para empresas **Modulo:** RAG **Ultima actualizacion:** 2026-09-14 -**Ultima modificacion por:** Subagente OCR Unit 6 Render +**Ultima modificacion por:** Subagente OCR Unit 7 Client **Estado:** Activo --- @@ -636,3 +636,31 @@ Continuidad operativa y evolutiva del modulo RAG. - `ocr-service/tests/test_render.py` - `openspec/changes/ocr-ingest-integration/{tasks.md,apply-progress.md}` - `docs/HISTORIAL_SESIONES.md` + +--- + +### 2026-09-14 - Subagente OCR Unit 7 Client - Cliente y artefactos OCR + +**Agente:** **Subagente OCR Unit 7 Client** +**Rol/responsabilidad:** Implementar exclusivamente Unit 7, tareas 4.2 y 4.3 y la parte de reintentos/integridad de 4.1, mediante TDD estricto, sin iniciar dispatcher ni routing de Unit 8. +**Modelo:** openai/gpt-5.6-sol +**Session ID OpenCode:** `ses_f5e4fa3e6ffeabETlcJxsURUKb` +**Directorio:** `/home/pancho/Documentos/Empresa/Desarrollo/IA/RAG` + +**Trabajo realizado:** +- Añadido el cliente OCR privado con clave idempotente estable, dos reenvios transitorios como maximo, clasificacion terminal de fallos deterministas, presion `429` reintentable, polling acotado y validacion estricta de identidad/esquema/resultados. +- Añadido el gestor de artefactos con originales y manifiesto canonico durables, permisos `0600`, IDs UUIDv5, hashes verificables, contencion de rutas y barrido de huerfanos antiguo y conservador. +- Mantenida 4.1 pendiente: Unit 7 solo cubre reintentos e integridad; los casos HTTP de ingesta, estado y fallo cerrado pertenecen a Unit 8. + +**Validacion y estado final:** +- RED valido por ausencia de los dos modulos; GREEN focalizado 6/6 y harness determinista 1/1 con tres intentos, misma clave y backoff exacto de 2/4 segundos. +- `npm test` 51/51, `npm run check`, `npm run build` y comprobaciones de espacios correctos; temporales eliminados y ningun proceso o llamada OCR externa iniciados. +- El slice funcional minimo es de 506 lineas autoradas y la contabilidad nativa total con metadatos obligatorios es de 585 lineas. El maintainer aprobo explicitamente `size:exception` porque cliente, artefactos y pruebas forman una unidad cohesiva; no se comprimio codigo ni se inicio un segundo particionado. +- El intento nativo se reinicio para validar la excepcion sin modificar produccion ni pruebas. El nuevo work unit es `unit-7-size-exception-validation` y remedia la evidencia fallida `sha256:dd224c3864adc8380dc7d9fd147cb448f5da8138277f3f346a35d7eeb3aca158`. +- No se hizo commit, push, PR, review nativa ni trabajo de Unit 8; no se accedio a rutas restringidas y se preservo `ocr-service/.venv`. + +**Archivos modificados:** +- `src/modules/ocr/{client,artifacts}.ts` +- `tests/ocr/client.test.ts` +- `openspec/changes/ocr-ingest-integration/{tasks.md,apply-progress.md}` +- `docs/HISTORIAL_SESIONES.md` diff --git a/openspec/changes/ocr-ingest-integration/apply-progress.md b/openspec/changes/ocr-ingest-integration/apply-progress.md index cd57eda..570f0ee 100644 --- a/openspec/changes/ocr-ingest-integration/apply-progress.md +++ b/openspec/changes/ocr-ingest-integration/apply-progress.md @@ -3,9 +3,9 @@ ## Current State - **Mode:** Strict TDD -- **Delivery:** Feature-branch-chain, Units 1–6 complete; maintainer-approved `size:exception` for Unit 2 only -- **Completed tasks:** 1.1, 1.2, 1.3, 1.4, 2.1, 2.2, 2.3, 2.4, 2.5, 3.1, 3.2, 3.3, 3.4 -- **Overall task progress:** 13/30 complete +- **Delivery:** Feature-branch-chain; Units 1–7 implemented; the maintainer approved `size:exception` for Unit 7's cohesive 506 functional authored lines and 585 total native-accounting lines +- **Completed tasks:** 1.1, 1.2, 1.3, 1.4, 2.1, 2.2, 2.3, 2.4, 2.5, 3.1, 3.2, 3.3, 3.4, 4.2, 4.3 +- **Overall task progress:** 15/30 complete ## Unit 1: Migration @@ -224,3 +224,50 @@ None — the migration follows the proposal, specifications, design, and closed - Start: Unit 5 private queue/API is complete. End: Unit 6 render/engine/schema/container behavior is complete. Unit 7 client work was not started. - Native Unit 6 implementation/tests/docs are 397 authored additions plus deletions, within the 400-line budget; required SDD progress and workspace history metadata are administrative evidence outside that native slice. - No specification or design deviation. The Docker image is private by deployment contract; CPU/RAM labels document limits that the platform must enforce. + +## Unit 7: OCR Client and Durable Artifacts + +### Implementation Summary + +- Added a private OCR HTTP client that preserves the exact idempotency key across one initial submission and at most two transient resubmissions, treats `429` as retryable queue pressure without immediate resubmission, and treats deterministic failures as terminal. +- Added 2–15 second polling backoff plus strict acknowledgement, status, engine, job, document, ordered-page, metrics, line, and schema validation before a result can be persisted. +- Added durable `0600` original and canonical manifest writes, UUIDv5 document artifact identities, lexical path containment, verifiable hashes, and age-gated orphan sweeping that preserves retained, recent, non-UUID, and symlink entries. +- Left task 4.1 pending because Unit 7 proves only its retries/integrity clause; Unit 8 still owns flag-off `201`, OCR `202`/status, OCR-down fail-closed, and catalog-down `503` HTTP coverage. + +### TDD Cycle Evidence + +| Task | Test File | Layer | Safety Net | RED | GREEN | TRIANGULATE | REFACTOR | +|---|---|---|---|---|---|---|---| +| 4.2 | `tests/ocr/client.test.ts` | Unit/runtime fetch stub | Existing OCR detection suite passed 5/5; production file was new. | Focused suite failed with `ERR_MODULE_NOT_FOUND` for `src/modules/ocr/client.js` before production code existed. | Final focused suite passed 6/6; runtime retry scenario passed 1/1. | Network then `503` produced exactly three same-key attempts with 2s/4s backoff; `429` and `422` each stopped after one attempt; malformed acknowledgement/schema/pages failed integrity; polling used 2s/4s. | Adapted multipart bytes to a type-safe `Uint8Array`; focused suite and type check remained green. | +| 4.3 | `tests/ocr/client.test.ts` | Filesystem integration | Existing OCR detection suite passed 5/5; production file was new. | The same missing-module RED preceded artifact production code. | Private originals, canonical manifest, safe resolution, and sweep passed in the final 6/6 suite. | Two out-of-order documents proved canonical ordering/hashes; traversal was rejected; sweep removed only an old orphan while preserving retained/recent directories and a symlink. | A test-side assertion was corrected to inspect the retained dangling symlink with `lstat`; no production behavior changed, and the suite passed 6/6. | + +### Test Summary + +- **Tests written and passing:** 6 focused tests; unit/fetch-stub and filesystem integration layers. +- **Approval tests:** None — both production modules are new. +- **RED:** exit 1, module-not-found before either production file existed. +- **GREEN/REFACTOR:** exit 0, 6 passed, 0 failed. + +### Work Unit Evidence + +| Evidence | Result | +|---|---| +| Focused test command and exact result | `NODE_ENV=test npx --no-install tsx --test tests/ocr/client.test.ts` — exit 0; 6 passed, 0 failed. | +| Runtime harness command/scenario and exact result | `NODE_ENV=test npx --no-install tsx --test --test-name-pattern "runtime retry stub" tests/ocr/client.test.ts` — exit 0; 1 passed, 0 failed. Deterministic stub observed network failure, `503`, then `202`: exactly 3 submission attempts, one unchanged idempotency key, 2,000/4,000 ms backoffs, followed by a strictly validated two-page result. | +| Rollback boundary | Remove `src/modules/ocr/client.ts`, `src/modules/ocr/artifacts.ts`, and `tests/ocr/client.test.ts`; revert tasks 4.2/4.3 and this Unit 7 progress/history metadata. Units 1–6 remain intact. | + +### Validation and Boundary + +- `npm test` passed 51/51; `npm run check`, `npm run build`, tracked whitespace, and no-index whitespace checks for all three new files passed. +- Test-created `rag-ocr-artifacts-*` and `rag-ocr-sweep-*` directories were removed; no server, external OCR call, or persistent process was started; `ocr-service/.venv` was preserved. +- Start: Unit 6 OCR rendering and container runtime are complete. End: Unit 7 client retries/integrity and durable artifacts/orphan sweep are implemented and green. +- Out of scope and untouched: dispatcher, reconciler, ingest branch, app upload/status routes, Unit 8 tests, commits, pushes, PRs, native review, deployment, and restricted paths. +- The implementation/test slice is 506 authored additions (client 174, artifacts 129, tests 203). This exceeds the 400-line default after one honest cohesive assessment; no code-golf or second slicing pass was attempted. The maintainer explicitly approved `size:exception`, including 585 total native-accounting lines with required metadata. +- No specification or design deviation. + +### Approved Size-Exception Corrective Context + +- The native attempt was reset after failed evidence revision `sha256:dd224c3864adc8380dc7d9fd147cb448f5da8138277f3f346a35d7eeb3aca158`; the failure was budget-only, not a functional defect. +- Corrective work unit `unit-7-size-exception-validation` uses attempt token `sha256:0109cfd8d3b37a5672c65a37700269b89cb80e4e357a39f68a920779a1715c39` and must settle against the failed revision through the parent. +- Production and test files remain unchanged for the corrective rerun; only approval/reset evidence and fresh validation metadata may change. +- Task 4.1 remains pending, while tasks 4.2 and 4.3 remain complete. Unit 8 was not started. diff --git a/openspec/changes/ocr-ingest-integration/tasks.md b/openspec/changes/ocr-ingest-integration/tasks.md index 844da40..c7fb89d 100644 --- a/openspec/changes/ocr-ingest-integration/tasks.md +++ b/openspec/changes/ocr-ingest-integration/tasks.md @@ -52,8 +52,8 @@ Tracker feature/ocr-ingest-integration is draft/no-merge and sole main target. P ## 4 Orchestration (Units 7–8; O1–O4) - [ ] 4.1 RED HTTP: flag-off 201; OCR 202/status; OCR-down fail-closed; catalog-down 503; retries/integrity -- [ ] 4.2 GREEN: `src/modules/ocr/client.ts` retries and integrity -- [ ] 4.3 GREEN: `src/modules/ocr/artifacts.ts` originals, manifest, sweep +- [x] 4.2 GREEN: `src/modules/ocr/client.ts` retries and integrity +- [x] 4.3 GREEN: `src/modules/ocr/artifacts.ts` originals, manifest, sweep - [ ] 4.4 GREEN: `src/modules/ocr/dispatcher.ts` leases; `src/modules/catalog/reconciler.ts` - [ ] 4.5 GREEN: `src/modules/ingest/service.ts` branch; `src/app.ts` upload/status diff --git a/src/modules/ocr/artifacts.ts b/src/modules/ocr/artifacts.ts new file mode 100644 index 0000000..1e22d60 --- /dev/null +++ b/src/modules/ocr/artifacts.ts @@ -0,0 +1,129 @@ +import { createHash, randomUUID } from "node:crypto"; +import { link, lstat, mkdir, open, readdir, rm, stat, unlink } from "node:fs/promises"; +import path from "node:path"; +import { canonicalJson, hashOrderedPairs, sha256Hex } from "../../shared/utils/ids.js"; + +const UUID = /^[0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/iu; +const UUID_NAMESPACE_URL = Buffer.from("6ba7b8119dad11d180b400c04fd430c8", "hex"); + +interface StageInput { + rootDirectory: string; + versionId: string; + createdAt: string; + documents: Array<{ documentId: string; documentKey: string; bytes: Buffer }>; +} + +export function resolveArtifactPath(versionDirectory: string, relativePath: string): string { + if (path.isAbsolute(relativePath)) throw new Error("Artifact path escapes version directory"); + const root = path.resolve(versionDirectory); + const resolved = path.resolve(root, relativePath); + if (resolved === root || !resolved.startsWith(`${root}${path.sep}`)) throw new Error("Artifact path escapes version directory"); + return resolved; +} + +export async function stageOcrArtifacts(input: StageInput): Promise<{ + versionDirectory: string; + manifestPath: string; + manifestSha256: string; + originalManifestHash: string; +}> { + assertUuid(input.versionId); + if (input.documents.length === 0 || Number.isNaN(Date.parse(input.createdAt))) throw new TypeError("Artifact manifest input is invalid"); + const documentIds = new Set(input.documents.map(({ documentId }) => documentId)); + const documentKeys = new Set(input.documents.map(({ documentKey }) => documentKey)); + if (documentIds.size !== input.documents.length || documentKeys.size !== input.documents.length) throw new TypeError("Artifact documents must be unique"); + + await mkdir(input.rootDirectory, { recursive: true, mode: 0o700 }); + const versionDirectory = path.join(path.resolve(input.rootDirectory), input.versionId); + await mkdir(versionDirectory, { mode: 0o700 }); + try { + const documents = []; + for (const document of [...input.documents].sort((left, right) => Buffer.compare(Buffer.from(left.documentKey), Buffer.from(right.documentKey)))) { + const documentArtifactId = uuidV5(document.documentId); + const relativeDirectory = path.posix.join("documents", documentArtifactId); + const directory = resolveArtifactPath(versionDirectory, relativeDirectory); + await mkdir(directory, { recursive: true, mode: 0o700 }); + const originalPath = path.posix.join(relativeDirectory, "original.pdf"); + await durableWrite(resolveArtifactPath(versionDirectory, originalPath), document.bytes); + documents.push({ + documentId: document.documentId, + documentKey: document.documentKey, + documentArtifactId, + originalPath, + originalSha256: sha256Hex(document.bytes) + }); + } + const originalManifestHash = hashOrderedPairs(documents.map(({ documentKey, originalSha256 }) => [documentKey, originalSha256]))!; + const manifest = { schemaVersion: "1", versionId: input.versionId, createdAt: input.createdAt, originalManifestHash, documents }; + const serialized = canonicalJson(manifest); + const manifestPath = path.join(versionDirectory, "manifest.json"); + await durableWrite(manifestPath, Buffer.from(serialized)); + await syncDirectory(versionDirectory); + await syncDirectory(path.resolve(input.rootDirectory)); + return { versionDirectory, manifestPath, manifestSha256: sha256Hex(serialized), originalManifestHash }; + } catch (error) { + await rm(versionDirectory, { recursive: true, force: true }); + throw error; + } +} + +export async function sweepOrphanArtifacts(input: { + rootDirectory: string; + retainedVersionIds: ReadonlySet; + olderThan: Date; +}): Promise { + for (const versionId of input.retainedVersionIds) assertUuid(versionId); + let entries; + try { + entries = await readdir(input.rootDirectory, { withFileTypes: true }); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return []; + throw error; + } + const removed: string[] = []; + for (const entry of entries.sort((left, right) => left.name.localeCompare(right.name))) { + if (!entry.isDirectory() || !UUID.test(entry.name) || input.retainedVersionIds.has(entry.name)) continue; + const candidate = path.join(path.resolve(input.rootDirectory), entry.name); + const current = await lstat(candidate); + if (!current.isDirectory() || current.mtimeMs >= input.olderThan.getTime()) continue; + await rm(candidate, { recursive: true, force: true }); + removed.push(entry.name); + } + if (removed.length > 0) await syncDirectory(path.resolve(input.rootDirectory)); + return removed; +} + +async function durableWrite(target: string, bytes: Buffer): Promise { + const temporary = `${target}.${randomUUID()}.tmp`; + const handle = await open(temporary, "wx", 0o600); + try { + await handle.writeFile(bytes); + await handle.chmod(0o600); + await handle.sync(); + } finally { + await handle.close(); + } + try { + await link(temporary, target); + } finally { + await unlink(temporary).catch(() => undefined); + } + await syncDirectory(path.dirname(target)); +} + +async function syncDirectory(directory: string): Promise { + const handle = await open(directory, "r"); + try { await handle.sync(); } finally { await handle.close(); } +} + +function uuidV5(value: string): string { + const bytes = createHash("sha1").update(Buffer.concat([UUID_NAMESPACE_URL, Buffer.from(value)])).digest().subarray(0, 16); + bytes[6] = (bytes[6]! & 0x0f) | 0x50; + bytes[8] = (bytes[8]! & 0x3f) | 0x80; + const hex = bytes.toString("hex"); + return `${hex.slice(0, 8)}-${hex.slice(8, 12)}-${hex.slice(12, 16)}-${hex.slice(16, 20)}-${hex.slice(20)}`; +} + +function assertUuid(value: string): void { + if (!UUID.test(value)) throw new TypeError("Version identity must be a UUID"); +} diff --git a/src/modules/ocr/client.ts b/src/modules/ocr/client.ts new file mode 100644 index 0000000..dc4ec86 --- /dev/null +++ b/src/modules/ocr/client.ts @@ -0,0 +1,174 @@ +import { sha256Hex } from "../../shared/utils/ids.js"; + +const OCR_CONFIG = { + languages: ["es", "en"], + dpi: 200, + engine: "paddleocr", + engineVersion: "3.4.0", + runtimeVersion: "3.2.2", + configVersion: "ocr-v1", + returnLayout: true +} as const; +const TRANSIENT_STATUSES = new Set([502, 503]); +const JOB_STATUSES = new Set(["queued", "running", "succeeded", "failed"]); + +export interface OcrAck { + jobId: string; + status: "queued"; + documentSha256: string; + requestedPages: number[]; + configVersion: "ocr-v1"; + createdAt: string; +} + +export interface OcrJobStatus { + jobId: string; + status: "queued" | "running" | "succeeded" | "failed"; + completedPages: number; + totalPages: number; + error: { code: string; message: string } | null; +} + +export interface OcrResult { + schemaVersion: "1"; + jobId: string; + documentSha256: string; + engine: { name: "paddleocr"; version: "3.4.0"; runtime: "paddlepaddle-3.2.2"; device: "cpu"; configVersion: "ocr-v1"; dpi: 200 }; + pages: Array<{ + page: number; + width: number; + height: number; + processingMs: number; + text: string; + metrics: { lineCount: number; nonWhitespaceCharacters: number; medianConfidence: number; p10Confidence: number; lowConfidenceLineRatio: number }; + lines: Array<{ lineId: string; text: string; confidence: number; bbox: [number, number, number, number] }>; + }>; +} + +export class OcrClientError extends Error { + constructor(public readonly code: string, public readonly status: number | undefined, public readonly retryable: boolean) { + super(`OCR request failed: ${code}`); + } +} + +interface OcrClientOptions { + baseUrl: string; + token: string; + fetch?: typeof globalThis.fetch; + sleep?: (milliseconds: number) => Promise; +} + +export class OcrClient { + private readonly baseUrl: string; + private readonly requestFetch: typeof globalThis.fetch; + private readonly sleep: (milliseconds: number) => Promise; + + constructor(private readonly options: OcrClientOptions) { + this.baseUrl = options.baseUrl.replace(/\/+$/u, ""); + this.requestFetch = options.fetch ?? globalThis.fetch; + this.sleep = options.sleep ?? ((milliseconds) => new Promise((resolve) => setTimeout(resolve, milliseconds))); + } + + async submit(file: Buffer, expected: { documentSha256: string; pages: number[] }): Promise { + assertExpected(expected); + if (sha256Hex(file) !== expected.documentSha256) throw integrityError(); + const payload = { documentSha256: expected.documentSha256, pages: expected.pages, ...OCR_CONFIG }; + const idempotencyKey = `${expected.documentSha256}:ocr-v1:${sha256Hex(JSON.stringify(expected.pages))}`; + const value = await this.requestJson("/v1/jobs", () => { + const form = new FormData(); + form.set("file", new Blob([new Uint8Array(file)], { type: "application/pdf" }), "original.pdf"); + form.set("request", JSON.stringify(payload)); + return { method: "POST", headers: this.headers({ "Idempotency-Key": idempotencyKey }), body: form }; + }); + if (!isObject(value) + || value.jobId === undefined + || value.status !== "queued" + || value.documentSha256 !== expected.documentSha256 + || value.configVersion !== "ocr-v1" + || !sameNumbers(value.requestedPages, expected.pages) + || typeof value.createdAt !== "string" + || Number.isNaN(Date.parse(value.createdAt))) throw integrityError(); + return value as unknown as OcrAck; + } + + async getStatus(jobId: string): Promise { + const value = await this.requestJson(`/v1/jobs/${encodeURIComponent(jobId)}`, () => ({ headers: this.headers() })); + if (!isObject(value) || value.jobId !== jobId || typeof value.status !== "string" || !JOB_STATUSES.has(value.status) + || !isCount(value.completedPages) || !isCount(value.totalPages) || value.completedPages > value.totalPages + || !(value.error === null || (isObject(value.error) && typeof value.error.code === "string" && typeof value.error.message === "string"))) { + throw integrityError(); + } + return value as unknown as OcrJobStatus; + } + + async pollUntilTerminal(jobId: string): Promise { + let delay = 2_000; + while (true) { + const status = await this.getStatus(jobId); + if (status.status === "succeeded" || status.status === "failed") return status; + await this.sleep(delay); + delay = Math.min(delay * 2, 15_000); + } + } + + async getResult(jobId: string, expected: { documentSha256: string; pages: number[] }): Promise { + assertExpected(expected); + const value = await this.requestJson(`/v1/jobs/${encodeURIComponent(jobId)}/result`, () => ({ headers: this.headers() })); + if (!validResult(value, jobId, expected)) throw integrityError(); + return value; + } + + private headers(additional: Record = {}): Headers { + return new Headers({ Authorization: `Bearer ${this.options.token}`, ...additional }); + } + + private async requestJson(pathname: string, buildInit: () => RequestInit): Promise { + for (let attempt = 0; attempt < 3; attempt += 1) { + let response: Response; + try { + response = await this.requestFetch(`${this.baseUrl}${pathname}`, buildInit()); + } catch { + 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(); + 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}`; + if (TRANSIENT_STATUSES.has(response.status) && attempt < 2) { await this.sleep(2_000 * 2 ** attempt); continue; } + throw new OcrClientError(code, response.status, response.status === 429); + } + throw new OcrClientError("OCR_RETRY_EXHAUSTED", undefined, false); + } +} + +function assertExpected(expected: { documentSha256: string; pages: number[] }): void { + if (!/^[a-f0-9]{64}$/u.test(expected.documentSha256) + || expected.pages.length === 0 + || expected.pages.some((page, index) => !Number.isInteger(page) || page < 1 || (index > 0 && page <= expected.pages[index - 1]!))) { + throw new TypeError("OCR request identity is invalid"); + } +} + +function validResult(value: unknown, jobId: string, expected: { documentSha256: string; pages: number[] }): value is OcrResult { + if (!isObject(value) || value.schemaVersion !== "1" || value.jobId !== jobId || value.documentSha256 !== expected.documentSha256 + || !isObject(value.engine) || canonicalEngine(value.engine) !== "paddleocr|3.4.0|paddlepaddle-3.2.2|cpu|ocr-v1|200" + || !Array.isArray(value.pages) || !sameNumbers(value.pages.map((page) => isObject(page) ? page.page : undefined), expected.pages)) return false; + return value.pages.every((page) => isObject(page) && isCount(page.page) && positiveCount(page.width) && positiveCount(page.height) + && isCount(page.processingMs) && typeof page.text === "string" && isObject(page.metrics) && Array.isArray(page.lines) + && page.lines.every((line) => isObject(line) && typeof line.lineId === "string" && line.lineId.length > 0 && typeof line.text === "string" + && validRatio(line.confidence) && Array.isArray(line.bbox) && line.bbox.length === 4 && line.bbox.every(Number.isFinite)) + && page.text === page.lines.map((line) => (line as Record).text).join("\n") + && page.metrics.lineCount === page.lines.length && isCount(page.metrics.nonWhitespaceCharacters) + && validRatio(page.metrics.medianConfidence) && validRatio(page.metrics.p10Confidence) && validRatio(page.metrics.lowConfidenceLineRatio)); +} + +function canonicalEngine(engine: Record): string { + return [engine.name, engine.version, engine.runtime, engine.device, engine.configVersion, engine.dpi].join("|"); +} +function isObject(value: unknown): value is Record { return typeof value === "object" && value !== null && !Array.isArray(value); } +function isCount(value: unknown): value is number { return Number.isInteger(value) && Number(value) >= 0; } +function positiveCount(value: unknown): value is number { return isCount(value) && value > 0; } +function validRatio(value: unknown): value is number { return typeof value === "number" && Number.isFinite(value) && value >= 0 && value <= 1; } +function sameNumbers(value: unknown, expected: number[]): boolean { return Array.isArray(value) && value.length === expected.length && value.every((entry, index) => entry === expected[index]); } +function integrityError(): Error { return new Error("OCR response integrity validation failed"); } diff --git a/tests/ocr/client.test.ts b/tests/ocr/client.test.ts new file mode 100644 index 0000000..a47785a --- /dev/null +++ b/tests/ocr/client.test.ts @@ -0,0 +1,203 @@ +import assert from "node:assert/strict"; +import { access, lstat, mkdir, mkdtemp, readFile, stat, symlink, utimes } from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; +import test from "node:test"; +import { OcrClient, OcrClientError } from "../../src/modules/ocr/client.js"; +import { + resolveArtifactPath, + stageOcrArtifacts, + sweepOrphanArtifacts +} from "../../src/modules/ocr/artifacts.js"; +import { canonicalJson, hashOrderedPairs, sha256Hex } from "../../src/shared/utils/ids.js"; + +const document = Buffer.from("%PDF-1.4\nunit-7\n%%EOF\n"); +const documentSha256 = sha256Hex(document); +const jobId = "ocr_job-7"; + +function jsonResponse(status: number, body: unknown): Response { + return new Response(JSON.stringify(body), { status, headers: { "content-type": "application/json" } }); +} + +function ack(overrides: Record = {}): Record { + return { + jobId, + status: "queued", + documentSha256, + requestedPages: [1, 3], + configVersion: "ocr-v1", + createdAt: "2026-09-14T10:00:00.000Z", + ...overrides + }; +} + +function result(overrides: Record = {}): Record { + return { + schemaVersion: "1", + jobId, + documentSha256, + engine: { + name: "paddleocr", + version: "3.4.0", + runtime: "paddlepaddle-3.2.2", + device: "cpu", + configVersion: "ocr-v1", + dpi: 200 + }, + pages: [1, 3].map((page) => ({ + page, + width: 1700, + height: 2200, + processingMs: 25, + text: `page ${page}`, + metrics: { + lineCount: 1, + nonWhitespaceCharacters: 5, + medianConfidence: 0.95, + p10Confidence: 0.95, + lowConfidenceLineRatio: 0 + }, + lines: [{ lineId: `p${page}-l1`, text: `page ${page}`, confidence: 0.95, bbox: [1, 2, 3, 4] }] + })), + ...overrides + }; +} + +test("runtime retry stub preserves identity, exact attempts, backoff, and result integrity", async () => { + const calls: Array<{ url: string; key: string | null }> = []; + const delays: number[] = []; + const responses: Array = [ + new TypeError("connection reset"), + jsonResponse(503, { detail: { code: "ENGINE_UNAVAILABLE" } }), + jsonResponse(202, ack()), + jsonResponse(200, result()) + ]; + const client = new OcrClient({ + baseUrl: "http://ocr.internal:8000/", + token: "internal-test-token", + fetch: (async (input, init) => { + calls.push({ url: String(input), key: new Headers(init?.headers).get("Idempotency-Key") }); + const response = responses.shift(); + if (response instanceof Error) throw response; + return response as Response; + }) as typeof fetch, + sleep: async (milliseconds) => { delays.push(milliseconds); } + }); + + const submitted = await client.submit(document, { documentSha256, pages: [1, 3] }); + const completed = await client.getResult(jobId, { documentSha256, pages: [1, 3] }); + + assert.equal(submitted.jobId, jobId); + assert.deepEqual(delays, [2_000, 4_000]); + assert.equal(calls.filter(({ url }) => url.endsWith("/v1/jobs")).length, 3); + assert.equal(new Set(calls.slice(0, 3).map(({ key }) => key)).size, 1); + assert.equal(calls[0]?.key, `${documentSha256}:ocr-v1:${sha256Hex("[1,3]")}`); + assert.deepEqual(completed.pages.map(({ page }) => page), [1, 3]); +}); + +test("queue pressure and deterministic failures are surfaced without retries", async () => { + for (const [status, retryable] of [[429, true], [422, false]] as const) { + let attempts = 0; + const client = new OcrClient({ + baseUrl: "http://ocr.internal:8000", + token: "token", + fetch: (async () => { + attempts += 1; + return jsonResponse(status, { detail: { code: status === 429 ? "QUEUE_FULL" : "UNSUPPORTED_PDF" } }); + }) as typeof fetch, + sleep: async () => { throw new Error("must not sleep"); } + }); + + await assert.rejects( + client.submit(document, { documentSha256, pages: [1, 3] }), + (error: unknown) => error instanceof OcrClientError && error.status === status && error.retryable === retryable + ); + assert.equal(attempts, 1); + } +}); + +test("strict validation rejects mismatched acknowledgements and result contracts", async () => { + for (const response of [ + jsonResponse(202, ack({ requestedPages: [3, 1] })), + jsonResponse(200, result({ schemaVersion: "2" })), + jsonResponse(200, result({ pages: [result().pages as unknown] })) + ]) { + const client = new OcrClient({ + baseUrl: "http://ocr.internal:8000", + token: "token", + fetch: (async () => response) as typeof fetch, + sleep: async () => undefined + }); + const operation = response.status === 202 + ? client.submit(document, { documentSha256, pages: [1, 3] }) + : client.getResult(jobId, { documentSha256, pages: [1, 3] }); + await assert.rejects(operation, /OCR response integrity validation failed/); + } +}); + +test("polling uses bounded exponential backoff until a terminal status", async () => { + const delays: number[] = []; + const states = ["queued", "running", "succeeded"] as const; + const client = new OcrClient({ + baseUrl: "http://ocr.internal:8000", + token: "token", + fetch: (async () => jsonResponse(200, { + jobId, + status: states.shift(), + completedPages: states.length === 0 ? 2 : 0, + totalPages: 2, + error: null + })) as typeof fetch, + sleep: async (milliseconds) => { delays.push(milliseconds); } + }); + + assert.equal((await client.pollUntilTerminal(jobId)).status, "succeeded"); + assert.deepEqual(delays, [2_000, 4_000]); +}); + +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 })); }); + const versionId = "11111111-1111-4111-8111-111111111111"; + const staged = await stageOcrArtifacts({ + rootDirectory, + versionId, + createdAt: "2026-09-14T10:00:00.000Z", + documents: [ + { documentId: "doc:source:b", documentKey: "b.pdf", bytes: Buffer.from("%PDF-b") }, + { documentId: "doc:source:a", documentKey: "a.pdf", bytes: Buffer.from("%PDF-a") } + ] + }); + const persisted = JSON.parse(await readFile(staged.manifestPath, "utf8")); + + assert.equal(staged.originalManifestHash, hashOrderedPairs([["b.pdf", sha256Hex("%PDF-b")], ["a.pdf", sha256Hex("%PDF-a")]])); + assert.equal(staged.manifestSha256, sha256Hex(canonicalJson(persisted))); + assert.deepEqual(persisted.documents.map((entry: { documentKey: string }) => entry.documentKey), ["a.pdf", "b.pdf"]); + for (const entry of persisted.documents) { + const originalPath = resolveArtifactPath(staged.versionDirectory, entry.originalPath); + assert.equal((await stat(originalPath)).mode & 0o777, 0o600); + assert.equal(sha256Hex(await readFile(originalPath)), entry.originalSha256); + } + assert.equal((await stat(staged.manifestPath)).mode & 0o777, 0o600); +}); + +test("safe resolution blocks traversal and orphan sweep preserves retained, recent, and non-directory entries", async (context) => { + const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-ocr-sweep-")); + context.after(async () => { await import("node:fs/promises").then(({ rm }) => rm(rootDirectory, { recursive: true, force: true })); }); + const retained = "22222222-2222-4222-8222-222222222222"; + const orphan = "33333333-3333-4333-8333-333333333333"; + const recent = "44444444-4444-4444-8444-444444444444"; + await Promise.all([retained, orphan, recent].map((id) => mkdir(path.join(rootDirectory, id)))); + await utimes(path.join(rootDirectory, orphan), new Date(0), new Date(0)); + await symlink(path.join(rootDirectory, orphan), path.join(rootDirectory, "55555555-5555-4555-8555-555555555555")); + + assert.throws(() => resolveArtifactPath(path.join(rootDirectory, retained), "../manifest.json"), /escapes version directory/); + assert.deepEqual(await sweepOrphanArtifacts({ + rootDirectory, + retainedVersionIds: new Set([retained]), + olderThan: new Date("2026-09-14T09:00:00.000Z") + }), [orphan]); + await assert.rejects(access(path.join(rootDirectory, orphan))); + await Promise.all([retained, recent].map((id) => access(path.join(rootDirectory, id)))); + assert.equal((await lstat(path.join(rootDirectory, "55555555-5555-4555-8555-555555555555"))).isSymbolicLink(), true); +});