From d2ecf1910208258ad2c3d7e46f069e27f6a425f2 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 repository job transitions --- package.json | 2 +- src/modules/catalog/repository.ts | 166 +++++++++++++++++++++++++++ tests/catalog/repository-ocr.test.ts | 147 ++++++++++++++++++++++++ 3 files changed, 314 insertions(+), 1 deletion(-) create mode 100644 tests/catalog/repository-ocr.test.ts diff --git a/package.json b/package.json index 84284b9..c12f6aa 100644 --- a/package.json +++ b/package.json @@ -9,7 +9,7 @@ "start": "node dist/server.js", "migrate:lifecycle": "node dist/modules/catalog/migrations.js", "migrate:legacy:lifecycle": "node dist/scripts/migrate-legacy-lifecycle.js", - "test": "NODE_ENV=test tsx --test tests/**/*.test.ts", + "test": "NODE_ENV=test tsx --test tests/*.test.ts tests/**/*.test.ts", "check": "tsc --noEmit -p tsconfig.json" }, "dependencies": { diff --git a/src/modules/catalog/repository.ts b/src/modules/catalog/repository.ts index 1b6be21..7653836 100644 --- a/src/modules/catalog/repository.ts +++ b/src/modules/catalog/repository.ts @@ -70,6 +70,35 @@ export interface OrphanedVersionCandidate { qdrantCollection: string; } +export interface PendingOcrIdentity { + sourceId: string; + originalManifestHash: string; + processingFingerprint: string; + metadataHash: string; +} + +export interface OcrJobRow { + jobId: string; + documentId: string; + remoteJobId: string | null; + remoteIdempotencyKey: string; + state: "queued" | "running" | "succeeded" | "failed"; + requestedPages: number[]; + completedPages: number; + configVersion: string; + attemptCount: number; + heartbeatAt: Date | null; + leaseExpiresAt: Date | null; + nextAttemptAt: Date | null; + errorCode: string | null; + errorDetail: string | null; +} + +export interface PendingOcrVersion { + version: CatalogVersionRow; + jobs: OcrJobRow[]; +} + type VersionDbRow = { version_id: string; source_id: string; @@ -89,6 +118,37 @@ type VersionDbRow = { qdrant_collection: string; }; +type OcrJobDbRow = { + ocr_job_id: string; + ocr_document_id: string; + ocr_remote_job_id: string | null; + ocr_remote_idempotency_key: string; + ocr_state: OcrJobRow["state"]; + ocr_requested_pages: number[]; + ocr_completed_pages: string | number; + ocr_config_version: string; + ocr_attempt_count: string | number; + ocr_heartbeat_at: Date | null; + ocr_lease_expires_at: Date | null; + ocr_next_attempt_at: Date | null; + ocr_error_code: string | null; + ocr_error_detail: string | null; +}; + +const OCR_JOB_COLUMNS = [ + ["job_id", "ocr_job_id"], ["document_id", "ocr_document_id"], + ["remote_job_id", "ocr_remote_job_id"], ["remote_idempotency_key", "ocr_remote_idempotency_key"], + ["state", "ocr_state"], ["requested_pages", "ocr_requested_pages"], + ["completed_pages", "ocr_completed_pages"], ["config_version", "ocr_config_version"], + ["attempt_count", "ocr_attempt_count"], ["heartbeat_at", "ocr_heartbeat_at"], + ["lease_expires_at", "ocr_lease_expires_at"], ["next_attempt_at", "ocr_next_attempt_at"], + ["error_code", "ocr_error_code"], ["error_detail", "ocr_error_detail"] +] as const; + +function ocrJobColumns(prefix = ""): string { + return OCR_JOB_COLUMNS.map(([column, alias]) => `${prefix}${column} AS ${alias}`).join(", "); +} + function toVersion(row: VersionDbRow): CatalogVersionRow { return { versionId: row.version_id, @@ -110,6 +170,25 @@ function toVersion(row: VersionDbRow): CatalogVersionRow { }; } +function toOcrJob(row: OcrJobDbRow): OcrJobRow { + return { + jobId: row.ocr_job_id, + documentId: row.ocr_document_id, + remoteJobId: row.ocr_remote_job_id, + remoteIdempotencyKey: row.ocr_remote_idempotency_key, + state: row.ocr_state, + requestedPages: row.ocr_requested_pages, + completedPages: Number(row.ocr_completed_pages), + configVersion: row.ocr_config_version, + attemptCount: Number(row.ocr_attempt_count), + heartbeatAt: row.ocr_heartbeat_at, + leaseExpiresAt: row.ocr_lease_expires_at, + nextAttemptAt: row.ocr_next_attempt_at, + errorCode: row.ocr_error_code, + errorDetail: row.ocr_error_detail + }; +} + export class CatalogRepository { constructor(private readonly pool: PgPool) {} @@ -220,6 +299,93 @@ export class CatalogRepository { return result.rows[0] ? toVersion(result.rows[0]) : undefined; } + async findPendingOcrVersion(input: PendingOcrIdentity): Promise { + const result = await this.pool.query( + `SELECT v.*, ${ocrJobColumns("j.")} + FROM rag_source_versions v + JOIN rag_version_documents d ON d.version_id = v.version_id + JOIN rag_ocr_jobs j ON j.version_id = d.version_id AND j.document_id = d.document_id + WHERE v.source_id = $1 + AND v.original_manifest_hash = $2 + AND v.processing_fingerprint = $3 + AND v.metadata_hash = $4 + AND v.source_content_hash IS NULL + AND v.state IN ('pending', 'indexing', 'review_required') + ORDER BY j.document_id`, + [input.sourceId, input.originalManifestHash, input.processingFingerprint, input.metadataHash] + ); + return result.rows[0] ? { version: toVersion(result.rows[0]), jobs: result.rows.map(toOcrJob) } : undefined; + } + + async claimNextOcrJob(leaseMs: number): Promise { + return withTransaction(this.pool, async (client) => { + const candidate = await client.query<{ job_id: string }>( + `SELECT job_id FROM rag_ocr_jobs + WHERE state = 'queued' AND (next_attempt_at IS NULL OR next_attempt_at <= now()) + ORDER BY created_at FOR UPDATE SKIP LOCKED LIMIT 1` + ); + if (!candidate.rowCount) return undefined; + const result = await client.query( + `UPDATE rag_ocr_jobs + SET state = 'running', attempt_count = attempt_count + 1, + started_at = COALESCE(started_at, now()), heartbeat_at = now(), + lease_expires_at = now() + ($2::text || ' milliseconds')::interval + WHERE job_id = $1 + RETURNING ${ocrJobColumns()}`, + [candidate.rows[0].job_id, String(leaseMs)] + ); + return result.rows[0] ? toOcrJob(result.rows[0]) : undefined; + }); + } + + async recoverExpiredOcrLeases(): Promise { + const result = await this.pool.query( + `UPDATE rag_ocr_jobs + SET state = 'queued', heartbeat_at = NULL, lease_expires_at = NULL, next_attempt_at = now() + WHERE state = 'running' AND lease_expires_at <= now() + RETURNING ${ocrJobColumns()}` + ); + return result.rows.map(toOcrJob); + } + + async markReviewRequired(versionId: string): Promise { + const result = await this.pool.query( + `UPDATE rag_source_versions SET state = 'review_required' + WHERE version_id = $1 AND state = 'indexing' + AND EXISTS (SELECT 1 FROM rag_ocr_jobs WHERE version_id = $1) + AND NOT EXISTS (SELECT 1 FROM rag_ocr_jobs WHERE version_id = $1 AND state <> 'succeeded') + AND NOT EXISTS ( + SELECT 1 FROM rag_version_documents d + WHERE d.version_id = $1 AND d.content_hash IS NULL + AND NOT EXISTS ( + SELECT 1 FROM rag_ocr_jobs j + WHERE j.version_id = d.version_id AND j.document_id = d.document_id AND j.state = 'succeeded' + ) + )`, + [versionId] + ); + if (result.rowCount !== 1) throw new CatalogError("Version is not ready for OCR review", 409, "INVALID_VERSION_STATE"); + } + + async markOcrReviewIndexing(versionId: string, reviewedBy: string): Promise { + const result = await this.pool.query( + `UPDATE rag_source_versions SET state = 'indexing', reviewed_at = now(), reviewed_by = $2, indexing_started_at = now() + WHERE version_id = $1 AND state = 'review_required'`, + [versionId, reviewedBy] + ); + if (result.rowCount !== 1) throw new CatalogError("Version is not awaiting OCR review", 409, "INVALID_VERSION_STATE"); + } + + async rejectOcrVersion(versionId: string, reviewedBy: string, reason: string): Promise { + const result = await this.pool.query( + `UPDATE rag_source_versions + SET state = 'rejected', reviewed_at = now(), reviewed_by = $2, error_code = 'OCR_REJECTED', error_detail = $3 + WHERE version_id = $1 AND state = 'review_required'`, + [versionId, reviewedBy, reason.slice(0, 2000)] + ); + if (result.rowCount !== 1) throw new CatalogError("Version is not awaiting OCR review", 409, "INVALID_VERSION_STATE"); + } + async createPendingVersion(input: CreateVersionInput): Promise { return withTransaction(this.pool, async (client) => { if (await this.isMaintenanceEnabled(client, true)) { diff --git a/tests/catalog/repository-ocr.test.ts b/tests/catalog/repository-ocr.test.ts new file mode 100644 index 0000000..2868023 --- /dev/null +++ b/tests/catalog/repository-ocr.test.ts @@ -0,0 +1,147 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { CatalogError, CatalogRepository } from "../../src/modules/catalog/repository.js"; + +const versionRow = { + version_id: "version-1", + source_id: "source-1", + version_number: 2, + previous_version_id: "version-0", + state: "indexing", + tags: [], + source_content_hash: null, + processing_fingerprint: "fingerprint", + metadata_hash: "metadata", + embedding_provider: "test", + embedding_model: "test", + embedding_dimensions: 3, + expected_document_count: 1, + expected_point_count: 0, + verified_point_count: 0, + qdrant_collection: "rag" +}; + +const jobRow = { + ocr_job_id: "job-1", + ocr_document_id: "document-1", + ocr_remote_job_id: "remote-1", + ocr_remote_idempotency_key: "document-hash:ocr-v1:pages-hash", + ocr_state: "running", + ocr_requested_pages: [1, 3], + ocr_completed_pages: 1, + ocr_config_version: "ocr-v1", + ocr_attempt_count: 2, + ocr_heartbeat_at: new Date("2026-09-14T10:00:00Z"), + ocr_lease_expires_at: new Date("2026-09-14T10:01:00Z"), + ocr_next_attempt_at: null, + ocr_error_code: null, + ocr_error_detail: null +}; + +function poolFor(handler: (sql: string, params?: unknown[]) => Promise<{ rowCount: number; rows: unknown[] }>) { + const client = { + query: async (sql: string, params?: unknown[]) => /^(BEGIN|COMMIT|ROLLBACK|SET CONSTRAINTS)/.test(sql.trim()) + ? { rowCount: 0, rows: [] } + : handler(sql, params), + release() {} + }; + return { query: handler, connect: async () => client }; +} + +test("pending OCR lookup returns the persisted version and jobs for a duplicate ingestion", async () => { + const pool = poolFor(async (sql, params) => { + assert.match(sql, /JOIN rag_version_documents/); + assert.match(sql, /JOIN rag_ocr_jobs/); + assert.match(sql, /source_content_hash IS NULL/); + assert.match(sql, /state IN \('pending', 'indexing', 'review_required'\)/); + assert.deepEqual(params, ["source-1", "manifest", "fingerprint", "metadata"]); + return { rowCount: 1, rows: [{ ...versionRow, ...jobRow }] }; + }); + const repository = new CatalogRepository(pool as never); + + const candidate = await repository.findPendingOcrVersion({ + sourceId: "source-1", + originalManifestHash: "manifest", + processingFingerprint: "fingerprint", + metadataHash: "metadata" + }); + + assert.equal(candidate?.version.versionId, "version-1"); + assert.deepEqual(candidate?.jobs.map((job) => [job.jobId, job.remoteJobId, job.remoteIdempotencyKey]), [ + ["job-1", "remote-1", "document-hash:ocr-v1:pages-hash"] + ]); +}); + +test("pending OCR lookup returns undefined when only terminal candidates exist", async () => { + const repository = new CatalogRepository(poolFor(async () => ({ rowCount: 0, rows: [] })) as never); + const candidate = await repository.findPendingOcrVersion({ + sourceId: "source-1", + originalManifestHash: "manifest", + processingFingerprint: "fingerprint", + metadataHash: "metadata" + }); + assert.equal(candidate, undefined); +}); + +test("lease claiming uses a locked queue row and keeps its remote identity", async () => { + const calls: string[] = []; + const pool = poolFor(async (sql, params) => { + calls.push(sql); + if (sql.includes("FOR UPDATE SKIP LOCKED")) { + assert.match(sql, /state = 'queued'/); + return { rowCount: 1, rows: [{ job_id: "job-1" }] }; + } + assert.match(sql, /attempt_count = attempt_count \+ 1/); + assert.deepEqual(params, ["job-1", "30000"]); + return { rowCount: 1, rows: [jobRow] }; + }); + const repository = new CatalogRepository(pool as never); + + const claimed = await repository.claimNextOcrJob(30_000); + + assert.equal(calls.length, 2); + assert.equal(claimed?.remoteJobId, "remote-1"); + assert.equal(claimed?.remoteIdempotencyKey, "document-hash:ocr-v1:pages-hash"); +}); + +test("lease claiming returns undefined when no queued job is eligible", async () => { + const repository = new CatalogRepository(poolFor(async () => ({ rowCount: 0, rows: [] })) as never); + assert.equal(await repository.claimNextOcrJob(30_000), undefined); +}); + +test("expired leases return to queued without clearing remote recovery identity", async () => { + const repository = new CatalogRepository(poolFor(async (sql) => { + assert.match(sql, /state = 'running'/); + assert.match(sql, /lease_expires_at <= now\(\)/); + assert.doesNotMatch(sql, /remote_(job_id|idempotency_key)\s*=/); + return { rowCount: 1, rows: [{ ...jobRow, ocr_state: "queued" }] }; + }) as never); + + const recovered = await repository.recoverExpiredOcrLeases(); + + assert.deepEqual(recovered.map((job) => [job.state, job.remoteJobId, job.remoteIdempotencyKey]), [ + ["queued", "remote-1", "document-hash:ocr-v1:pages-hash"] + ]); +}); + +test("review transitions enforce indexing, approval, and rejection state guards", async () => { + const successfulSql: string[] = []; + const success = new CatalogRepository(poolFor(async (sql) => { + successfulSql.push(sql); + return { rowCount: 1, rows: [] }; + }) as never); + await success.markReviewRequired("version-1"); + await success.markOcrReviewIndexing("version-1", "reviewer-1"); + await success.rejectOcrVersion("version-1", "reviewer-1", "Unreadable code"); + assert.match(successfulSql[0], /state = 'indexing'.*NOT EXISTS/s); + assert.match(successfulSql[0], /FROM rag_version_documents/); + assert.match(successfulSql[0], /rag_ocr_jobs WHERE version_id = \$1 AND state <> 'succeeded'/); + assert.match(successfulSql[1], /state = 'review_required'/); + assert.match(successfulSql[2], /state = 'review_required'/); + + const blocked = new CatalogRepository(poolFor(async () => ({ rowCount: 0, rows: [] })) as never); + await assert.rejects( + blocked.markReviewRequired("version-2"), + (error) => error instanceof CatalogError && error.code === "INVALID_VERSION_STATE" + ); +});