From 936881b38ff3629738cfc85ca8a3467f1c1d84c2 Mon Sep 17 00:00:00 2001 From: Paco POR-CORREO Date: Mon, 14 Sep 2026 14:20:52 +0200 Subject: [PATCH] feat(ocr): add durable review schema --- migrations/002_ocr_review.sql | 61 ++++++++++++++++++++++++++ tests/catalog/migration-002.test.ts | 67 +++++++++++++++++++++++++++++ 2 files changed, 128 insertions(+) create mode 100644 migrations/002_ocr_review.sql create mode 100644 tests/catalog/migration-002.test.ts diff --git a/migrations/002_ocr_review.sql b/migrations/002_ocr_review.sql new file mode 100644 index 0000000..9bd9eda --- /dev/null +++ b/migrations/002_ocr_review.sql @@ -0,0 +1,61 @@ +CREATE TABLE IF NOT EXISTS rag_ocr_jobs ( + job_id uuid PRIMARY KEY DEFAULT gen_random_uuid(), + version_id uuid NOT NULL, + document_id text NOT NULL, + remote_job_id text NULL, + remote_idempotency_key text NOT NULL, + state text NOT NULL CHECK (state IN ('queued', 'running', 'succeeded', 'failed')), + requested_pages integer[] NOT NULL, + completed_pages integer NOT NULL DEFAULT 0 CHECK (completed_pages >= 0), + config_version text NOT NULL, + attempt_count integer NOT NULL DEFAULT 0 CHECK (attempt_count >= 0), + heartbeat_at timestamptz NULL, + lease_expires_at timestamptz NULL, + next_attempt_at timestamptz NULL, + error_code text NULL, + error_detail text NULL, + created_at timestamptz NOT NULL DEFAULT now(), + started_at timestamptz NULL, + completed_at timestamptz NULL, + UNIQUE (version_id, document_id), + FOREIGN KEY (version_id, document_id) + REFERENCES rag_version_documents(version_id, document_id) ON DELETE RESTRICT +); + +CREATE TABLE IF NOT EXISTS rag_document_pages ( + version_id uuid NOT NULL, + document_id text NOT NULL, + page_number integer NOT NULL CHECK (page_number > 0), + extraction_method text NOT NULL CHECK (extraction_method IN ('native', 'ocr', 'blank')), + native_text_hash char(64) NULL, + ocr_text_hash char(64) NULL, + candidate_text_hash char(64) NULL, + reviewed_text_hash char(64) NULL, + metrics jsonb NOT NULL DEFAULT '{}', + risk_tokens jsonb NOT NULL DEFAULT '[]', + blocked_reason text NULL, + created_at timestamptz NOT NULL DEFAULT now(), + PRIMARY KEY (version_id, document_id, page_number), + FOREIGN KEY (version_id, document_id) + REFERENCES rag_version_documents(version_id, document_id) ON DELETE RESTRICT +); + +CREATE TABLE IF NOT EXISTS rag_review_corrections ( + correction_id uuid PRIMARY KEY DEFAULT gen_random_uuid(), + version_id uuid NOT NULL, + document_id text NOT NULL, + page_number integer NOT NULL, + line_id text NOT NULL, + expected_line_hash char(64) NOT NULL, + replacement_text text NOT NULL, + reviewed_by text NOT NULL, + created_at timestamptz NOT NULL DEFAULT now(), + FOREIGN KEY (version_id, document_id, page_number) + REFERENCES rag_document_pages(version_id, document_id, page_number) ON DELETE RESTRICT, + UNIQUE (version_id, document_id, page_number, line_id) +); + +CREATE UNIQUE INDEX IF NOT EXISTS rag_one_pending_ocr_identity + ON rag_source_versions(source_id, original_manifest_hash, processing_fingerprint, metadata_hash) + WHERE source_content_hash IS NULL + AND state IN ('pending', 'indexing', 'review_required'); diff --git a/tests/catalog/migration-002.test.ts b/tests/catalog/migration-002.test.ts new file mode 100644 index 0000000..3447224 --- /dev/null +++ b/tests/catalog/migration-002.test.ts @@ -0,0 +1,67 @@ +import assert from "node:assert/strict"; +import { readFile } from "node:fs/promises"; +import test from "node:test"; + +const migrationUrl = new URL("../../migrations/002_ocr_review.sql", import.meta.url); + +async function readMigration(): Promise { + return readFile(migrationUrl, "utf8"); +} + +test("migration 002 persists durable OCR jobs with idempotency and lease fields", async () => { + const sql = await readMigration(); + + assert.match(sql, /CREATE TABLE IF NOT EXISTS rag_ocr_jobs/i); + assert.match(sql, /remote_job_id text NULL/i); + assert.match(sql, /remote_idempotency_key text NOT NULL/i); + assert.match(sql, /requested_pages integer\[\] NOT NULL/i); + assert.match(sql, /state text NOT NULL CHECK \(state IN \('queued', 'running', 'succeeded', 'failed'\)\)/i); + assert.match(sql, /attempt_count integer NOT NULL DEFAULT 0/i); + assert.match(sql, /heartbeat_at timestamptz NULL/i); + assert.match(sql, /lease_expires_at timestamptz NULL/i); + assert.match(sql, /next_attempt_at timestamptz NULL/i); + assert.match(sql, /UNIQUE \(version_id, document_id\)/i); + assert.match( + sql, + /FOREIGN KEY \(version_id, document_id\)\s+REFERENCES rag_version_documents\(version_id, document_id\) ON DELETE RESTRICT/i + ); +}); + +test("migration 002 stores auditable page extraction hashes, metrics, and risks", async () => { + const sql = await readMigration(); + + assert.match(sql, /CREATE TABLE IF NOT EXISTS rag_document_pages/i); + assert.match(sql, /page_number integer NOT NULL CHECK \(page_number > 0\)/i); + assert.match(sql, /extraction_method text NOT NULL CHECK \(extraction_method IN \('native', 'ocr', 'blank'\)\)/i); + assert.match(sql, /native_text_hash char\(64\) NULL/i); + assert.match(sql, /ocr_text_hash char\(64\) NULL/i); + assert.match(sql, /candidate_text_hash char\(64\) NULL/i); + assert.match(sql, /reviewed_text_hash char\(64\) NULL/i); + assert.match(sql, /metrics jsonb NOT NULL DEFAULT '\{\}'/i); + assert.match(sql, /risk_tokens jsonb NOT NULL DEFAULT '\[\]'/i); + assert.match(sql, /blocked_reason text NULL/i); + assert.match(sql, /PRIMARY KEY \(version_id, document_id, page_number\)/i); +}); + +test("migration 002 makes review correction targets unique and page-bound", async () => { + const sql = await readMigration(); + + assert.match(sql, /CREATE TABLE IF NOT EXISTS rag_review_corrections/i); + assert.match(sql, /expected_line_hash char\(64\) NOT NULL/i); + assert.match(sql, /replacement_text text NOT NULL/i); + assert.match(sql, /reviewed_by text NOT NULL/i); + assert.match( + sql, + /FOREIGN KEY \(version_id, document_id, page_number\)\s+REFERENCES rag_document_pages\(version_id, document_id, page_number\) ON DELETE RESTRICT/i + ); + assert.match(sql, /UNIQUE \(version_id, document_id, page_number, line_id\)/i); +}); + +test("migration 002 prevents duplicate non-terminal OCR identities", async () => { + const sql = await readMigration(); + + assert.match( + sql, + /CREATE UNIQUE INDEX IF NOT EXISTS rag_one_pending_ocr_identity\s+ON rag_source_versions\(source_id, original_manifest_hash, processing_fingerprint, metadata_hash\)\s+WHERE source_content_hash IS NULL\s+AND state IN \('pending', 'indexing', 'review_required'\)/i + ); +});