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" ); });