import assert from "node:assert/strict"; import { mkdtemp, rm, writeFile } from "node:fs/promises"; import os from "node:os"; import path from "node:path"; import test from "node:test"; import { CatalogError } from "../../src/modules/catalog/errors.js"; import type { EmbeddingProvider } from "../../src/modules/embeddings/provider.js"; import { PostgresOcrIndexingStore, OcrReadyIndexingService, type ApprovedOcrCandidate } from "../../src/modules/ocr/indexing.js"; import { persistComposedCandidateArtifact, persistOcrResultArtifact, persistReviewedPagesArtifact, persistReviewImageArtifacts, stageOcrArtifacts } from "../../src/modules/ocr/artifacts.js"; import { buildOcrIdentity } from "../../src/modules/ocr/client.js"; import { computeNativeMetrics } from "../../src/modules/ocr/detection.js"; import { chunkDocument, documentalChunkingPolicy } from "../../src/modules/process/chunking.js"; import type { VectorStoreClient } from "../../src/modules/vectorstore/client.js"; import { buildChunkId, buildVersionedQdrantPointId, normalizeContentForHash, sha256Hex } from "../../src/shared/utils/ids.js"; const versionId = "eeeeeeee-eeee-4eee-8eee-eeeeeeeeeeee"; const documentId = "doc:indexed-review"; const sourceId = "src:reviewed"; const jobId = "indexed-review"; const png = Buffer.from("89504e470d0a1a0a01020304", "hex"); function ocrPage(page: number, text: string) { return { page, width: 100, height: 100, processingMs: 1, text, metrics: { lineCount: 1, nonWhitespaceCharacters: computeNativeMetrics(text).nonWhitespaceCharacters, inkCoverage: 0.4, medianConfidence: 0.97, p10Confidence: 0.97, lowConfidenceLineRatio: 0 }, lines: [{ lineId: `p${page}-l1`, text, confidence: 0.97, bbox: [1, 2, 30, 10] as [number, number, number, number] }] }; } async function fixture(rootDirectory: string) { const first = "Reviewed alpha content. ".repeat(110); const second = "Reviewed beta content. ".repeat(3).trim(); const original = Buffer.from("%PDF-reviewed-index"); const identity = buildOcrIdentity({ documentSha256: sha256Hex(original), pages: [1, 2], idempotencyKey: "indexed-review-key" }); await stageOcrArtifacts({ rootDirectory, versionId, createdAt: "2026-09-22T12:00:00.000Z", documents: [{ documentId, documentKey: "review.pdf", bytes: original, requestedPages: [1, 2], pages: [1, 2].map((page) => ({ page, text: "", rasterCoverage: 1, textSha256: sha256Hex("") })) }] }); await persistOcrResultArtifact({ rootDirectory, versionId, documentId, result: { schemaVersion: "1", jobId, ...identity, engine: { name: "paddleocr", version: "3.4.0", runtime: "paddlepaddle-3.2.2", device: "cpu", configVersion: "ocr-v2", dpi: 200 }, pages: [ocrPage(1, first), ocrPage(2, second)] } }); await persistReviewImageArtifacts({ rootDirectory, versionId, documentId, images: [1, 2].map((page) => ({ page, bytes: png, sha256: sha256Hex(png) })) }); const { candidate } = await persistComposedCandidateArtifact({ rootDirectory, versionId, jobs: [{ documentId, remoteJobId: jobId, requestedPages: [1, 2], state: "succeeded" }] }); const pages = candidate.documents[0]!.pages; const reviewedText = pages.map(({ candidateText }) => candidateText).join("\n\n"); await persistReviewedPagesArtifact({ rootDirectory, versionId, sourceId, candidateSha256: candidate.candidateSha256, reviewedTextSha256: sha256Hex(reviewedText), reviewedBy: "reviewer", documents: [{ documentId, pages: pages.map(({ page, candidateText, lines }) => ({ page, candidateText, ocr: { lines } })) }] }); return { reviewedText, pages }; } function candidate(reviewedText: string): ApprovedOcrCandidate { return { versionId, sourceId, state: "indexing", activateRequested: true, expectedActiveVersionId: "active-1", reviewedText, reviewedTextSha256: sha256Hex(reviewedText), processingFingerprint: "fingerprint", metadataHash: "metadata" }; } function harness(pageCount: number, options: { embeddingFailure?: boolean; countOffset?: number } = {}) { const sql: string[] = []; const points: Array<{ id: string; vector: number[]; payload: Record }> = []; let embedCalls = 0; const rows = Array.from({ length: pageCount }, (_, index) => ({ version_id: versionId, source_id: sourceId, version_number: 4, state: "indexing", base_active_version_id: "active-1", current_active_version_id: "active-1", processing_fingerprint: "fingerprint", metadata_hash: "metadata", tags: ["reviewed"], embedding_provider: "test-provider", embedding_model: "test-model", embedding_dimensions: 3, qdrant_collection: "rag_documents", expected_document_count: 1, document_id: documentId, document_key: "review.pdf", title: "review.pdf", mime_type: "application/pdf", page_number: index + 1, reviewed_text_hash: "pending" })); const query = async (statement: string, params?: unknown[]) => { sql.push(statement); if (/^(BEGIN|COMMIT|ROLLBACK|SET CONSTRAINTS)/u.test(statement.trim())) return { rowCount: 0, rows: [] }; if (statement.includes("FROM rag_source_versions v")) return { rowCount: rows.length, rows }; if (statement.includes("UPDATE rag_version_documents")) return { rowCount: 1, rows: [] }; if (statement.includes("UPDATE rag_source_versions")) return { rowCount: 1, rows: [{ version_id: versionId }] }; throw new Error(`Unexpected SQL: ${statement} ${String(params)}`); }; const pool = { query, async connect() { return { query, release() {} }; } }; const embeddings: EmbeddingProvider = { providerName: "test-provider", modelName: "test-model", dimensions: 3, async embed(input) { embedCalls += 1; if (options.embeddingFailure) throw new Error("embedding unavailable"); return input.map(() => [0.1, 0.2, 0.3]); } }; const vectors: Partial = { kind: "fake", async upsert(chunks) { points.push(...chunks); }, async countVersionPoints() { return points.length + (options.countOffset ?? 0); } }; return { pool: pool as never, embeddings, vectors: vectors as VectorStoreClient, points, sql, embedCalls: () => embedCalls }; } test("reviewed artifact indexing survives restart and writes canonical versioned chunks before ready", async (context) => { const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-reviewed-index-")); context.after(() => rm(rootDirectory, { recursive: true, force: true })); const { reviewedText, pages } = await fixture(rootDirectory); const runtime = harness(pages.length); const originalQuery = (runtime.pool as { query: (sql: string, params?: unknown[]) => Promise<{ rowCount: number; rows: Array> }> }).query; for (let index = 0; index < pages.length; index += 1) { const row = (await originalQuery("SELECT * FROM rag_source_versions v", [])).rows[index]!; row.reviewed_text_hash = sha256Hex(pages[index]!.candidateText); } const restarted = new PostgresOcrIndexingStore(runtime.pool, rootDirectory, runtime.embeddings, runtime.vectors); const result = await new OcrReadyIndexingService(restarted).index(candidate(reviewedText)); const expected = chunkDocument("review.pdf", normalizeContentForHash(reviewedText), documentalChunkingPolicy); assert.deepEqual(result, { versionId, state: "ready", activated: false }); assert.equal(runtime.points.length, expected.length); assert.deepEqual(runtime.points.map(({ id }) => id), expected.map((chunk) => buildVersionedQdrantPointId(versionId, buildChunkId(documentId, "documental", chunk.index)))); assert.ok(runtime.points.every(({ payload }) => payload.processing_fingerprint === "fingerprint" && payload.source_version_id === versionId)); assert.match(runtime.sql.join("\n"), /SET state = 'ready'/u); assert.doesNotMatch(runtime.sql.join("\n"), /active_version_id\s*=/u); }); test("identity, corrupt artifact, embedding, and partial-write failures remain fail closed", async (context) => { const rootDirectory = await mkdtemp(path.join(os.tmpdir(), "rag-reviewed-index-fail-")); context.after(() => rm(rootDirectory, { recursive: true, force: true })); const { reviewedText, pages } = await fixture(rootDirectory); for (const scenario of [ { name: "identity", runtime: harness(pages.length), value: { ...candidate(reviewedText), processingFingerprint: "wrong" } }, { name: "embedding", runtime: harness(pages.length, { embeddingFailure: true }), value: candidate(reviewedText) }, { name: "count", runtime: harness(pages.length, { countOffset: -1 }), value: candidate(reviewedText) } ]) { await assert.rejects(new OcrReadyIndexingService(new PostgresOcrIndexingStore(scenario.runtime.pool, rootDirectory, scenario.runtime.embeddings, scenario.runtime.vectors)).index(scenario.value), scenario.name === "count" ? (error) => error instanceof CatalogError && error.code === "SOURCE_VERSION_INCONSISTENT" : Error); assert.doesNotMatch(scenario.runtime.sql.join("\n"), /SET state = 'ready'/u); if (scenario.name === "identity") assert.equal(scenario.runtime.embedCalls(), 0); } await writeFile(path.join(rootDirectory, versionId, "reviewed-pages.json"), "corrupt", { mode: 0o600 }); const corrupt = harness(pages.length); await assert.rejects(new OcrReadyIndexingService(new PostgresOcrIndexingStore(corrupt.pool, rootDirectory, corrupt.embeddings, corrupt.vectors)).index(candidate(reviewedText)), /integrity validation failed/u); assert.equal(corrupt.embedCalls(), 0); assert.equal(corrupt.points.length, 0); });