import assert from "node:assert/strict"; import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; import os from "node:os"; import path from "node:path"; import test from "node:test"; import { createApp } from "../../src/app.js"; import { env } from "../../src/config/env.js"; import { CatalogError } from "../../src/modules/catalog/errors.js"; import { CatalogRepository } from "../../src/modules/catalog/repository.js"; import { KnowledgeLifecycleReconciler } from "../../src/modules/catalog/reconciler.js"; import { OcrDispatcher, type OcrSchedulerClock } from "../../src/modules/ocr/dispatcher.js"; import { readReviewImageArtifact } from "../../src/modules/ocr/artifacts.js"; import { IngestService } from "../../src/modules/ingest/service.js"; import type { EmbeddingProvider } from "../../src/modules/embeddings/provider.js"; import type { VectorStoreClient } from "../../src/modules/vectorstore/client.js"; import { sha256Hex } from "../../src/shared/utils/ids.js"; import { buildOcrIdentity, OcrClientError } from "../../src/modules/ocr/client.js"; function setEnvFlag(name: keyof typeof env, value: unknown): void { (env as unknown as Record)[name] = value; } function buildPdf(text: string): Buffer { const stream = text ? `BT /F1 10 Tf 40 760 Td (${text.replace(/([\\()])/g, "\\$1")}) Tj ET` : ""; const objects = [ "<< /Type /Catalog /Pages 2 0 R >>", "<< /Type /Pages /Kids [3 0 R] /Count 1 >>", "<< /Type /Page /Parent 2 0 R /MediaBox [0 0 612 792] /Resources << /Font << /F1 5 0 R >> >> /Contents 4 0 R >>", `<< /Length ${Buffer.byteLength(stream)} >>\nstream\n${stream}\nendstream`, "<< /Type /Font /Subtype /Type1 /BaseFont /Helvetica >>" ]; let pdf = "%PDF-1.4\n"; const offsets = [0]; objects.forEach((object, index) => { offsets.push(Buffer.byteLength(pdf)); pdf += `${index + 1} 0 obj\n${object}\nendobj\n`; }); const xref = Buffer.byteLength(pdf); pdf += `xref\n0 ${objects.length + 1}\n0000000000 65535 f \n`; pdf += offsets.slice(1).map((offset) => `${String(offset).padStart(10, "0")} 00000 n \n`).join(""); pdf += `trailer\n<< /Size ${objects.length + 1} /Root 1 0 R >>\nstartxref\n${xref}\n%%EOF\n`; return Buffer.from(pdf); } function embeddingProvider(): EmbeddingProvider { return { providerName: "test-provider", modelName: "test-model", dimensions: 3, async embed(input) { return input.map(() => [0.1, 0.2, 0.3]); } }; } function vectorStore(): VectorStoreClient { return { kind: "fake", async upsert() {}, async countVersionPoints() { return 0; } } as VectorStoreClient; } async function flushPromises(): Promise { await new Promise((resolve) => setImmediate(resolve)); } test("OCR-disabled textual PDFs retain the native synchronous path", async (context) => { const previousLifecycle = env.knowledgeLifecycleEnforced; setEnvFlag("knowledgeLifecycleEnforced", true); context.after(() => setEnvFlag("knowledgeLifecycleEnforced", previousLifecycle)); const directory = await mkdtemp(path.join(os.tmpdir(), "rag-ocr-disabled-")); context.after(() => rm(directory, { recursive: true, force: true })); const filePath = path.join(directory, "native.pdf"); const sufficient = Array.from({ length: 24 }, (_, index) => `Alpha${index} beta${index}`).join(" "); await writeFile(filePath, buildPdf(sufficient)); const calls: string[] = []; const catalog = { async beginAttempt() { return "attempt-1"; }, async updateAttempt() {}, async withSourceLock(_sourceId: string, handler: () => Promise) { return handler(); }, async findReusableVersion() { return { versionId: "native-1", versionNumber: 1, state: "active", previousVersionId: null, expectedDocumentCount: 1 }; }, async assertActiveVersionPrecondition() { calls.push("native"); } }; const service = new IngestService(embeddingProvider(), vectorStore(), catalog as never, { enabled: false, artifactRoot: directory }); const result = await service.ingest({ sourceType: "file", sourceRef: "native.pdf", readPath: filePath, activate: true, expectedActiveVersionId: null }); assert.equal(result.state, "active"); assert.equal("phase" in result, false); assert.deepEqual(calls, ["native"]); }); test("OCR-required PDFs are durably accepted once and duplicate pending ingestion reuses status identity", async (context) => { const previousLifecycle = env.knowledgeLifecycleEnforced; setEnvFlag("knowledgeLifecycleEnforced", true); context.after(() => setEnvFlag("knowledgeLifecycleEnforced", previousLifecycle)); const directory = await mkdtemp(path.join(os.tmpdir(), "rag-ocr-accept-")); context.after(() => rm(directory, { recursive: true, force: true })); const filePath = path.join(directory, "scanned.pdf"); await writeFile(filePath, buildPdf("")); let createdInput: Record | undefined; let pending: Record | undefined; const catalog = { async beginAttempt() { return "attempt-1"; }, async updateAttempt() {}, async withSourceLock(_sourceId: string, handler: () => Promise) { return handler(); }, async findPendingOcrVersion() { return pending; }, async createPendingVersion(input: Record) { createdInput = input; pending = { version: { versionId: input.versionId, versionNumber: 4, state: "indexing" }, jobs: [{ jobId: "job-1" }] }; return (pending as { version: unknown }).version; }, async markIndexing() {} }; const service = new IngestService(embeddingProvider(), vectorStore(), catalog as never, { enabled: true, artifactRoot: path.join(directory, "artifacts") }); const source = { sourceId: "src:scan", sourceType: "file" as const, sourceRef: "scanned.pdf", readPath: filePath, activate: true, expectedActiveVersionId: null }; const accepted = await service.ingest(source); const duplicate = await service.ingest(source); assert.deepEqual(accepted, { accepted: true, sourceId: "src:scan", versionId: accepted.versionId, versionNumber: 4, state: "indexing", phase: "ocr_queued", statusUrl: `/ingestions/${accepted.versionId}`, reviewUrl: null, activated: false }); assert.equal(duplicate.versionId, accepted.versionId); assert.equal((createdInput?.sourceContentHash), null); assert.deepEqual((createdInput?.ocrJobs as Array<{ requestedPages: number[] }>)[0]?.requestedPages, [1]); const originalPath = path.join(directory, "artifacts", String(accepted.versionId), "documents"); assert.equal((await readFile(path.join(directory, "artifacts", String(accepted.versionId), "manifest.json"), "utf8")).includes("scanned.pdf"), true); assert.equal(path.isAbsolute(originalPath), true); }); test("OCR acceptance retains every supported native document and reloads a concurrent winner", async (context) => { const previousLifecycle = env.knowledgeLifecycleEnforced; setEnvFlag("knowledgeLifecycleEnforced", true); context.after(() => setEnvFlag("knowledgeLifecycleEnforced", previousLifecycle)); const directory = await mkdtemp(path.join(os.tmpdir(), "rag-ocr-multi-")); context.after(() => rm(directory, { recursive: true, force: true })); await writeFile(path.join(directory, "scanned.pdf"), buildPdf("")); await writeFile(path.join(directory, "notes.txt"), "Native companion document\n"); const winner = { version: { versionId: "22222222-2222-4222-8222-222222222222", versionNumber: 7, state: "indexing" }, jobs: [{ jobId: "winner-job" }] }; let lookupCount = 0; let candidateInput: Record | undefined; const catalog = { async beginAttempt() { return "attempt-multi"; }, async updateAttempt() {}, async withSourceLock(_sourceId: string, handler: () => Promise) { return handler(); }, async findPendingOcrVersion() { lookupCount += 1; return lookupCount === 1 ? undefined : winner; }, async createPendingVersion(input: Record) { candidateInput = input; throw Object.assign(new Error("duplicate pending identity"), { code: "23505" }); } }; const service = new IngestService(embeddingProvider(), vectorStore(), catalog as never, { enabled: true, artifactRoot: path.join(directory, "artifacts") }); const accepted = await service.ingest({ sourceId: "src:multi", sourceType: "folder", sourceRef: "bundle", readPath: directory, activate: true, expectedActiveVersionId: null }); assert.equal(accepted.versionId, winner.version.versionId); assert.equal(candidateInput?.expectedDocumentCount, 2); assert.deepEqual((candidateInput?.documents as Array<{ documentKey: string }>).map(({ documentKey }) => documentKey), ["notes.txt", "scanned.pdf"]); assert.equal((candidateInput?.ocrJobs as unknown[]).length, 1); }); test("application wiring dispatches a newly accepted OCR job from durable artifacts", async (context) => { const previous = { lifecycle: env.knowledgeLifecycleEnforced, enabled: (env as unknown as Record).ocrIngestEnabled, root: (env as unknown as Record).ocrArtifactRoot }; const directory = await mkdtemp(path.join(os.tmpdir(), "rag-ocr-runtime-")); context.after(() => rm(directory, { recursive: true, force: true })); setEnvFlag("knowledgeLifecycleEnforced", true); setEnvFlag("ocrIngestEnabled" as keyof typeof env, true); setEnvFlag("ocrArtifactRoot" as keyof typeof env, path.join(directory, "artifacts")); context.after(() => { setEnvFlag("knowledgeLifecycleEnforced", previous.lifecycle); setEnvFlag("ocrIngestEnabled" as keyof typeof env, previous.enabled); setEnvFlag("ocrArtifactRoot" as keyof typeof env, previous.root); }); const pdf = buildPdf(""); let queuedJob: Record | undefined; let runningJob: Record | undefined; let completedJob: Record | undefined; let documentSha256 = ""; let persistedKey: unknown; const dispatchCalls: string[] = []; const catalog = { async beginAttempt() { return "runtime-attempt"; }, async updateAttempt() {}, async withSourceLock(_sourceId: string, handler: () => Promise) { return handler(); }, async findPendingOcrVersion() { return undefined; }, async createPendingVersion(input: Record) { const job = (input.ocrJobs as Array>)[0]!; persistedKey = job.remoteIdempotencyKey; queuedJob = { ...job, jobId: "runtime-job", versionId: input.versionId, remoteJobId: null, state: "queued", completedPages: 0, attemptCount: 0, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt: null, errorCode: null, errorDetail: null }; return { versionId: input.versionId, versionNumber: 3 }; }, async markIndexing() {}, async claimNextOcrJob() { const job = queuedJob; queuedJob = undefined; runningJob = job ? { ...job, state: "running", attemptCount: 1 } : undefined; return runningJob; }, async setOcrRemoteJob(_jobId: string, remoteJobId: string) { runningJob = { ...runningJob, remoteJobId }; }, async requeueOcrJob() { dispatchCalls.push("requeued"); }, async completeOcrJob() { runningJob = completedJob = { ...runningJob, state: "succeeded" }; dispatchCalls.push("completed"); return true; }, async listOcrJobs() { return [runningJob]; }, async persistOcrCandidate() { dispatchCalls.push("candidate"); }, async failOcrJob() {}, async markReviewRequired() { dispatchCalls.push("review-required"); }, async markFailed() {} }; const png = Buffer.from("89504e470d0a1a0a0102", "hex"); const client = { async submit(bytes: Buffer, expected: { documentSha256: string }, key: string) { assert.equal(key, persistedKey); documentSha256 = expected.documentSha256; dispatchCalls.push(`${bytes.length}:${key}`); return { jobId: "runtime-remote" }; }, async getStatus() { return { jobId: "runtime-remote", status: "succeeded", ...buildOcrIdentity({ documentSha256, pages: [1], idempotencyKey: String(persistedKey) }), completedPages: 1, totalPages: 1, error: null }; }, async getResult() { const text = "Factura FAT07 has enough OCR characters for review"; return { schemaVersion: "1", jobId: "runtime-remote", documentSha256, ...buildOcrIdentity({ documentSha256, pages: [1], idempotencyKey: String(persistedKey) }), engine: { name: "paddleocr", version: "3.4.0", runtime: "paddlepaddle-3.2.2", device: "cpu", configVersion: "ocr-v2", dpi: 200 }, pages: [{ page: 1, width: 100, height: 100, processingMs: 1, text, metrics: { lineCount: 1, nonWhitespaceCharacters: 40, inkCoverage: 0.5, medianConfidence: 0.95, p10Confidence: 0.95, lowConfidenceLineRatio: 0 }, lines: [{ lineId: "p1-l1", text, confidence: 0.95, bbox: [1, 2, 3, 4] }] }] }; }, async getReviewImage() { return { bytes: png, sha256: sha256Hex(png) }; }, async delete() {} }; const server = createApp({ catalog: catalog as never, ocrClient: client as never, startReconciler: false }).listen(0); context.after(() => server.close()); const address = server.address(); assert.ok(address && typeof address === "object"); const form = new FormData(); form.set("file", new Blob([new Uint8Array(pdf)], { type: "application/pdf" }), "runtime.pdf"); form.set("sourceId", "src:runtime"); form.set("activate", "true"); form.set("expectedActiveVersionId", "null"); const response = await fetch(`http://127.0.0.1:${address.port}/ingest/upload`, { method: "POST", body: form }); assert.equal(response.status, 202); assert.equal((await response.json() as { uploadedResource: string }).uploadedResource, "runtime.pdf"); for (let attempt = 0; attempt < 200 && !dispatchCalls.includes("review-required"); attempt += 1) await new Promise((resolve) => setTimeout(resolve, 10)); assert.match(dispatchCalls[0]!, /^\d+:.+:ocr-v2:.+$/u); assert.deepEqual(dispatchCalls.slice(1), ["completed", "candidate", "review-required"]); assert.deepEqual((await readReviewImageArtifact({ rootDirectory: path.join(directory, "artifacts"), versionId: String(completedJob?.versionId), documentId: String(completedJob?.documentId), page: 1 })).bytes, png); }); test("dispatcher recovers only expired work, reuses its remote job, and completes it once", async () => { const calls: string[] = []; const job = { jobId: "job-expired", versionId: "version-1", documentId: "document-1", remoteJobId: "remote-1", remoteIdempotencyKey: "persisted-key", state: "queued" as const, requestedPages: [1], completedPages: 0, configVersion: "ocr-v2", attemptCount: 1, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt: null, errorCode: null, errorDetail: null }; let available = true; const repository = { async recoverExpiredOcrLeases() { calls.push("recover"); return [job]; }, async claimNextOcrJob() { throw new Error("recovery must not claim unrelated queued work"); }, async getNextOcrAttemptAt() { return undefined; }, async claimOcrJob(jobId: string) { assert.equal(jobId, job.jobId); if (!available) return undefined; available = false; calls.push("claim-exact"); return { ...job, state: "running" as const }; }, async setOcrRemoteJob() { calls.push("submit-persist"); }, async requeueOcrJob() { calls.push("requeue"); }, async completeOcrJob() { calls.push("complete"); return true; }, async failOcrJob() { calls.push("job-failed"); }, async markReviewRequired() { calls.push("review-required"); }, async markFailed() { calls.push("version-failed"); } }; const client = { async submit() { calls.push("submit"); throw new Error("existing remote work must not be duplicated"); }, async getStatus() { calls.push("status"); return { jobId: "remote-1", status: "succeeded", completedPages: 1, totalPages: 1, error: null }; }, async getResult() { calls.push("result"); return { pages: [{ page: 1 }] }; }, async delete() { calls.push("delete"); throw new Error("cleanup unavailable"); } }; const dispatcher = new OcrDispatcher( repository, client as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), async (_job, result) => { calls.push("artifact-write"); return result; }, async () => { calls.push("candidate-finalized"); } ); assert.equal(await dispatcher.recoverExpiredLeases(), 1); assert.deepEqual(calls, ["recover", "claim-exact", "status", "result", "artifact-write", "complete", "candidate-finalized", "review-required", "delete"]); }); test("dispatcher retains remote OCR state when durable result transfer fails", async () => { const calls: string[] = []; const job = { jobId: "job-1", versionId: "version-1", documentId: "document-1", remoteJobId: "remote-1", remoteIdempotencyKey: "key", state: "running" as const, requestedPages: [1], completedPages: 0, configVersion: "ocr-v2", attemptCount: 1, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt: null, errorCode: null, errorDetail: null }; const store = { async claimNextOcrJob() { return job; }, async completeOcrJob() { calls.push("complete"); return true; }, async requeueOcrJob() { calls.push("requeued"); }, async failOcrJob() { calls.push("job-failed"); }, async markFailed() { calls.push("version-failed"); } }; const client = { async getStatus() { return { status: "succeeded" }; }, async getResult() { return { pages: [] }; }, async delete() { calls.push("delete"); } }; const dispatcher = new OcrDispatcher( store as never, client as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), async () => { calls.push("artifact-write"); throw new Error("durable write failed"); }, async () => undefined ); assert.equal(await dispatcher.runOnce(), "pending"); assert.deepEqual(calls, ["artifact-write", "requeued"]); }); test("dispatcher fails closed on an unrecoverable remote OCR response", async () => { const calls: string[] = []; const job = { jobId: "job-integrity", versionId: "version-1", documentId: "document-1", remoteJobId: "remote-1", remoteIdempotencyKey: "key", state: "running" as const, requestedPages: [1], completedPages: 0, configVersion: "ocr-v2", attemptCount: 1, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt: null, errorCode: null, errorDetail: null }; const dispatcher = new OcrDispatcher({ async claimNextOcrJob() { return job; }, async requeueOcrJob() { calls.push("requeued"); }, async failOcrJob(_jobId: string, code: string) { calls.push(`job:${code}`); }, async markFailed(_versionId: string, code: string) { calls.push(`version:${code}`); } } as never, { async getStatus() { throw new OcrClientError("OCR_RESPONSE_INTEGRITY_FAILED", undefined, false); } } as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), async (_job, result) => result, async () => undefined); assert.equal(await dispatcher.runOnce(), "failed"); assert.deepEqual(calls, ["job:OCR_RESPONSE_INTEGRITY_FAILED", "version:OCR_RESPONSE_INTEGRITY_FAILED"]); }); test("dispatcher propagates the remote terminal failure code", async () => { const calls: string[] = []; const job = { jobId: "job-failed", versionId: "version-1", documentId: "document-1", remoteJobId: "remote-1", remoteIdempotencyKey: "key", state: "running" as const, requestedPages: [1], completedPages: 0, configVersion: "ocr-v2", attemptCount: 1, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt: null, errorCode: null, errorDetail: null }; const dispatcher = new OcrDispatcher({ async claimNextOcrJob() { return job; }, async failOcrJob(_jobId: string, code: string) { calls.push(`job:${code}`); }, async markFailed(_versionId: string, code: string) { calls.push(`version:${code}`); } } as never, { async getStatus() { return { status: "failed", error: { code: "OCR_ENGINE_FAILED", message: "engine failed" } }; } } as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), async (_job, result) => result, async () => undefined); assert.equal(await dispatcher.runOnce(), "failed"); assert.deepEqual(calls, ["job:OCR_ENGINE_FAILED", "version:OCR_ENGINE_FAILED"]); }); test("durable scheduler follows persisted backoff and stops after terminal success", async () => { let now = 1_000; let timer: { callback: () => void; delayMs: number; cleared: boolean; unref(): void } | undefined; const clock: OcrSchedulerClock = { now: () => now, setTimeout(callback, delayMs) { timer = { callback: () => { timer = undefined; callback(); }, delayMs, cleared: false, unref() {} }; return timer as unknown as NodeJS.Timeout; }, clearTimeout(handle) { (handle as unknown as { cleared: boolean }).cleared = true; } }; const job = { jobId: "job-scheduled", versionId: "version-1", documentId: "document-1", remoteJobId: "remote-1", remoteIdempotencyKey: "key", state: "queued" as "queued" | "running" | "succeeded", requestedPages: [1], completedPages: 0, configVersion: "ocr-v2", attemptCount: 0, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt: null as Date | null, errorCode: null, errorDetail: null }; const statuses = ["queued", "running", "succeeded"] as const; let statusCalls = 0; let completed = 0; const store = { async claimNextOcrJob() { if (job.state !== "queued" || (job.nextAttemptAt && job.nextAttemptAt.getTime() > now)) return undefined; job.state = "running"; job.attemptCount += 1; return { ...job, state: "running" as const }; }, async getNextOcrAttemptAt() { return job.state === "queued" ? job.nextAttemptAt ?? new Date(now) : undefined; }, async requeueOcrJob(_jobId: string, _code: string, _detail: string, delayMs: number) { job.state = "queued"; job.nextAttemptAt = new Date(now + delayMs); }, async completeOcrJob() { job.state = "succeeded"; completed += 1; return true; }, async failOcrJob() {}, async markReviewRequired() {}, async markFailed() {} }; const dispatcher = new OcrDispatcher(store as never, { async getStatus() { const status = statuses[statusCalls++]!; return { status, error: null }; }, async getResult() { return { pages: [] }; }, async delete() {} } as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), async (_job, result) => result, async () => undefined, 30_000, clock); dispatcher.start(); await flushPromises(); assert.equal(timer?.delayMs, 2_000); now += 2_000; timer!.callback(); await flushPromises(); assert.equal(timer?.delayMs, 4_000); now += 4_000; timer!.callback(); await flushPromises(); assert.equal(statusCalls, 3); assert.equal(completed, 1); assert.equal(job.state, "succeeded"); assert.equal(timer, undefined); dispatcher.stop(); }); test("durable scheduler leaves no timer after a terminal remote failure", async () => { let state: "queued" | "running" | "failed" = "queued"; let timersCreated = 0; const clock: OcrSchedulerClock = { now: () => 1_000, setTimeout() { timersCreated += 1; return { unref() {} } as unknown as NodeJS.Timeout; }, clearTimeout() {} }; const store = { async claimNextOcrJob() { if (state !== "queued") return undefined; state = "running"; return { jobId: "job-failed", versionId: "version-1", documentId: "document-1", remoteJobId: "remote-1", remoteIdempotencyKey: "key", state: "running" as const, requestedPages: [1], completedPages: 0, configVersion: "ocr-v2", attemptCount: 1, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt: null, errorCode: null, errorDetail: null }; }, async getNextOcrAttemptAt() { return undefined; }, async failOcrJob() { state = "failed"; }, async markFailed() {} }; const dispatcher = new OcrDispatcher(store as never, { async getStatus() { return { status: "failed", error: { code: "OCR_ENGINE_FAILED", message: "engine failed" } }; } } as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), async (_job, result) => result, async () => undefined, 30_000, clock); dispatcher.start(); await flushPromises(); assert.equal(state, "failed"); assert.equal(timersCreated, 0); dispatcher.stop(); }); test("scheduler rebuilds a future persisted attempt after restart without resubmitting OCR", async () => { let now = 10_000; let scheduled: { callback: () => void; delayMs: number; unref(): void } | undefined; const clock: OcrSchedulerClock = { now: () => now, setTimeout(callback, delayMs) { scheduled = { callback: () => { scheduled = undefined; callback(); }, delayMs, unref() {} }; return scheduled as unknown as NodeJS.Timeout; }, clearTimeout() {} }; let state: "queued" | "running" | "succeeded" = "queued"; const nextAttemptAt = new Date(now + 8_000); let submitCalls = 0; const store = { async claimNextOcrJob() { if (state !== "queued" || now < nextAttemptAt.getTime()) return undefined; state = "running"; return { jobId: "job-restart", versionId: "version-1", documentId: "document-1", remoteJobId: "remote-existing", remoteIdempotencyKey: "key", state: "running" as const, requestedPages: [1], completedPages: 0, configVersion: "ocr-v2", attemptCount: 2, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt, errorCode: null, errorDetail: null }; }, async getNextOcrAttemptAt() { return state === "queued" ? nextAttemptAt : undefined; }, async completeOcrJob() { state = "succeeded"; return true; }, async failOcrJob() {}, async markReviewRequired() {}, async markFailed() {} }; const dispatcher = new OcrDispatcher(store as never, { async submit() { submitCalls += 1; throw new Error("must reuse persisted remote job"); }, async getStatus() { return { status: "succeeded" }; }, async getResult() { return { pages: [] }; }, async delete() {} } as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), async (_job, result) => result, async () => undefined, 30_000, clock); dispatcher.start(); await flushPromises(); assert.equal(scheduled?.delayMs, 8_000); now += 8_000; scheduled!.callback(); await flushPromises(); assert.equal(state, "succeeded"); assert.equal(submitCalls, 0); dispatcher.stop(); }); test("concurrent dispatchers and duplicate drains process a leased job only once", async () => { let state: "queued" | "running" | "succeeded" = "queued"; let claimCalls = 0; let statusCalls = 0; let completeCalls = 0; const job = { jobId: "job-concurrent", versionId: "version-1", documentId: "document-1", remoteJobId: "remote-1", remoteIdempotencyKey: "key", state: "running" as const, requestedPages: [1], completedPages: 0, configVersion: "ocr-v2", attemptCount: 1, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt: null, errorCode: null, errorDetail: null }; const store = { async claimNextOcrJob() { claimCalls += 1; await Promise.resolve(); if (state !== "queued") return undefined; state = "running"; return job; }, async getNextOcrAttemptAt() { return undefined; }, async completeOcrJob() { state = "succeeded"; completeCalls += 1; return true; }, async failOcrJob() {}, async markReviewRequired() {}, async markFailed() {} }; const client = { async getStatus() { statusCalls += 1; return { status: "succeeded" }; }, async getResult() { return { pages: [] }; }, async delete() {} }; const createDispatcher = () => new OcrDispatcher(store as never, client as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), async (_job, result) => result, async () => undefined); const first = createDispatcher(); const second = createDispatcher(); const firstDrain = first.dispatchAvailable(); assert.equal(first.dispatchAvailable(), firstDrain); await Promise.all([firstDrain, second.dispatchAvailable()]); assert.equal(state, "succeeded"); assert.equal(statusCalls, 1); assert.equal(completeCalls, 1); assert.ok(claimCalls >= 2); }); test("dispatcher retains remote OCR state when durable candidate finalization fails", async () => { const calls: string[] = []; const job = { jobId: "job-1", versionId: "version-1", documentId: "document-1", remoteJobId: "remote-1", remoteIdempotencyKey: "key", state: "running" as const, requestedPages: [1], completedPages: 0, configVersion: "ocr-v2", attemptCount: 1, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt: null, errorCode: null, errorDetail: null }; const dispatcher = new OcrDispatcher({ async claimNextOcrJob() { return job; }, async completeOcrJob() { calls.push("complete"); return true; }, async failOcrJob() { calls.push("job-failed"); }, async markFailed(_versionId: string, code: string) { calls.push(`version-failed:${code}`); }, async markReviewRequired() { calls.push("review-required"); } } as never, { async getStatus() { return { status: "succeeded" }; }, async getResult() { return { pages: [] }; }, async delete() { calls.push("delete"); } } as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), async (_job, result) => result, async () => { calls.push("candidate-write-readback"); throw new Error("OCR_QUALITY_BLOCKED"); }); assert.equal(await dispatcher.runOnce(), "failed"); assert.deepEqual(calls, ["complete", "candidate-write-readback", "version-failed:OCR_QUALITY_BLOCKED"]); }); test("local finalization resumes after every durable boundary exactly once without resubmitting OCR", async () => { for (const failedBoundary of ["result", "images", "report", "diagnostics"] as const) { const completedBoundaries = new Set(); let boundaryFailed = false; let jobState: "queued" | "running" | "succeeded" = "queued"; let versionState: "indexing" | "review_required" = "indexing"; let reviewTransitions = 0; let submitCalls = 0; const job = { jobId: `job-${failedBoundary}`, versionId: `version-${failedBoundary}`, documentId: "document-1", remoteJobId: "remote-1", remoteIdempotencyKey: "persisted-key", state: "running" as const, requestedPages: [1], completedPages: 0, configVersion: "ocr-v2", attemptCount: 1, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt: null, errorCode: null, errorDetail: null }; const persistBoundary = (boundary: typeof failedBoundary) => { completedBoundaries.add(boundary); if (!boundaryFailed && boundary === failedBoundary) { boundaryFailed = true; throw new Error(`simulated crash after ${boundary}`); } }; const dispatcher = new OcrDispatcher({ async claimNextOcrJob() { if (jobState !== "queued") return undefined; jobState = "running"; return job; }, async claimOcrJob() { return undefined; }, async recoverExpiredOcrLeases() { return []; }, async setOcrRemoteJob() { throw new Error("a known remote job must never be resubmitted"); }, async requeueOcrJob() { jobState = "queued"; }, async completeOcrJob() { assert.equal(jobState, "running"); jobState = "succeeded"; return true; }, async failOcrJob() { throw new Error("a recoverable local crash must not fail OCR"); }, async markReviewRequired() { assert.equal(versionState, "indexing"); versionState = "review_required"; reviewTransitions += 1; }, async markFailed() { throw new Error("a recoverable local crash must not fail the version"); }, async listOcrVersionsAwaitingCandidate() { return jobState === "succeeded" && versionState === "indexing" ? [job.versionId] : []; }, async listOcrJobs() { return [{ ...job, state: "succeeded" as const }]; } } as never, { async submit() { submitCalls += 1; throw new Error("a known remote job must never be submitted"); }, async getStatus() { return { status: "succeeded" }; }, async getResult() { return { pages: [{ page: 1 }] }; }, async delete() {} } as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), async (_job, result) => { persistBoundary("result"); persistBoundary("images"); return result; }, async () => { persistBoundary("report"); persistBoundary("diagnostics"); }); const firstAttempt = await dispatcher.runOnce(); if (failedBoundary === "result" || failedBoundary === "images") { assert.equal(firstAttempt, "pending", failedBoundary); assert.equal(await dispatcher.runOnce(), "succeeded", failedBoundary); } else { assert.equal(firstAttempt, "succeeded", failedBoundary); } await dispatcher.recoverCompletedCandidates(); await dispatcher.recoverCompletedCandidates(); assert.deepEqual([...completedBoundaries].sort(), ["diagnostics", "images", "report", "result"]); assert.equal(reviewTransitions, 1, failedBoundary); assert.equal(submitCalls, 0, failedBoundary); assert.equal(versionState, "review_required", failedBoundary); } }); test("administrative routes expose only classified errors", async (context) => { const previousToken = env.lifecycleAdminToken; setEnvFlag("lifecycleAdminToken", "admin-token"); context.after(() => setEnvFlag("lifecycleAdminToken", previousToken)); let failure: Error = new Error("unexpected failure at /private/path with token=secret"); const server = createApp({ catalog: { async activateVersion() { throw failure; } } as never, startReconciler: false }).listen(0); context.after(() => server.close()); const address = server.address(); assert.ok(address && typeof address === "object"); const activate = () => fetch(`http://127.0.0.1:${address.port}/sources/source-1/versions/version-1/activate`, { method: "POST", headers: { authorization: "Bearer admin-token", "content-type": "application/json" }, body: JSON.stringify({ expectedActiveVersionId: "active-1" }) }); const unexpected = await activate(); assert.equal(unexpected.status, 500); const unexpectedBody = await unexpected.json() as { error: string; code?: string }; assert.equal(unexpectedBody.error, "Unknown activation error"); assert.equal(unexpectedBody.code, undefined); assert.doesNotMatch(JSON.stringify(unexpectedBody), /private\/path|token=secret/u); failure = new CatalogError("Active version changed during activation", 409, "ACTIVE_VERSION_CHANGED"); const classified = await activate(); assert.equal(classified.status, 409); assert.deepEqual(await classified.json(), { ok: false, error: "Active version changed during activation", code: "ACTIVE_VERSION_CHANGED" }); }); test("OCR exhaustion fails the leased candidate without activation or engine substitution", async () => { const calls: string[] = []; const repository = { async claimNextOcrJob() { return { jobId: "job-1", versionId: "version-1", documentId: "document-1", remoteJobId: null, remoteIdempotencyKey: "same-key", state: "running", requestedPages: [1], completedPages: 0, configVersion: "ocr-v2", attemptCount: 1, heartbeatAt: null, leaseExpiresAt: null, nextAttemptAt: null, errorCode: null, errorDetail: null }; }, async setOcrRemoteJob() { calls.push("remote"); }, async requeueOcrJob() { calls.push("requeue"); }, async completeOcrJob() { calls.push("complete"); return false; }, async failOcrJob(_jobId: string, code: string) { calls.push(`job:${code}`); }, async markReviewRequired() { calls.push("review"); }, async markFailed(_versionId: string, code: string) { calls.push(`version:${code}`); } }; const client = { async submit() { throw new Error("OCR unavailable after retries"); } }; const dispatcher = new OcrDispatcher( repository as never, client as never, async () => ({ bytes: Buffer.from("pdf"), documentSha256: sha256Hex("pdf") }), async (_job, result) => result, async () => undefined ); assert.equal(await dispatcher.runOnce(), "failed"); assert.deepEqual(calls, ["job:OCR_DISPATCH_FAILED", "version:OCR_DISPATCH_FAILED"]); }); test("reconciler makes expired OCR leases dispatchable while leaving live leases to PostgreSQL", async () => { const calls: string[] = []; const catalog = { async withGlobalTryLock(_name: string, handler: () => Promise) { return handler(); }, async resolveActiveVersions() { return []; }, async validateActiveInvariant() { return []; }, async listOrphanedIndexingCandidates() { return []; } }; const dispatcher = { async recoverExpiredLeases() { calls.push("ocr-recovery"); return 1; }, async dispatchAvailable() { calls.push("ocr-dispatch"); return 2; } }; const reconciler = new KnowledgeLifecycleReconciler(catalog as never, vectorStore(), dispatcher as never); const status = await reconciler.runOnce(); assert.equal(status.ocrLeasesRecovered, 1); assert.deepEqual(calls, ["ocr-recovery", "ocr-dispatch"]); }); test("reconciler restart completes durable OCR work without resubmitting OCR", async () => { const candidate = { state: "indexing", durableArtifact: false }; const calls: string[] = []; const catalog = { async withGlobalTryLock(_name: string, handler: () => Promise) { return handler(); }, async resolveActiveVersions() { return []; }, async validateActiveInvariant() { return []; }, async listOrphanedIndexingCandidates() { return []; } }; const dispatcher = { async recoverExpiredLeases() { calls.push("ocr-recovery"); return 0; }, async dispatchAvailable() { calls.push("ocr-dispatch"); return 0; }, async recoverCompletedCandidates() { candidate.state = "review_required"; candidate.durableArtifact = true; calls.push("candidate-recovery"); return 1; } }; for (const reconciler of [ new KnowledgeLifecycleReconciler(catalog as never, vectorStore(), dispatcher as never), new KnowledgeLifecycleReconciler(catalog as never, vectorStore(), dispatcher as never) ]) { assert.equal((await reconciler.runOnce()).ok, true); } assert.deepEqual(candidate, { state: "review_required", durableArtifact: true }); assert.deepEqual(calls, ["ocr-recovery", "ocr-dispatch", "candidate-recovery", "ocr-recovery", "ocr-dispatch", "candidate-recovery"]); }); test("runtime HTTP routing returns native 201, OCR 202/status, and catalog-down 503", async () => { const previous = { knowledgeLifecycleEnforced: env.knowledgeLifecycleEnforced, lifecycleAdminToken: env.lifecycleAdminToken, ocrIngestEnabled: env.ocrIngestEnabled }; setEnvFlag("knowledgeLifecycleEnforced", true); setEnvFlag("lifecycleAdminToken", "unit-8-token"); setEnvFlag("ocrIngestEnabled", true); const versionId = "11111111-1111-4111-8111-111111111111"; const ingestService = { async ingest(input: Record) { return input.sourceRef === "native.pdf" ? { accepted: true, sourceId: "src:native", versionId: "native-1", versionNumber: 1, state: "active", filesDiscovered: 1, documentsProcessed: 1, chunksStored: 1, activated: true, noOp: false, collectionName: "rag" } : { accepted: true, sourceId: "src:scan", versionId, versionNumber: 2, state: "indexing", phase: "ocr_queued", statusUrl: `/ingestions/${versionId}`, reviewUrl: null, activated: false }; }, async cleanup() { return { deleted: 0 }; } }; const catalog = { async claimNextOcrJob() { return undefined; }, async getIngestionStatus(id: string) { assert.equal(id, versionId); return { sourceId: "src:scan", versionId, state: "indexing", phase: "ocr_running", activated: false, documents: [{ documentId: "doc:scan", state: "ocr_running", completedPages: 0, totalPages: 1, pages: [{ page: 1, method: "ocr", state: "ocr_running" }] }], error: null, statusUrl: `/ingestions/${versionId}`, reviewUrl: null }; } }; const app = createApp({ ingestService: ingestService as never, catalog: catalog as never, startReconciler: false }); const server = app.listen(0); try { const address = server.address(); assert.ok(address && typeof address === "object"); const baseUrl = `http://127.0.0.1:${address.port}`; const request = (sourceRef: string) => fetch(`${baseUrl}/ingest`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ sourceType: "file", sourceRef }) }); assert.equal((await request("native.pdf")).status, 201); const accepted = await request("scan.pdf"); assert.equal(accepted.status, 202); assert.equal((await accepted.json() as { statusUrl: string }).statusUrl, `/ingestions/${versionId}`); assert.equal((await fetch(`${baseUrl}/ingestions/${versionId}`)).status, 401); const status = await fetch(`${baseUrl}/ingestions/${versionId}`, { headers: { authorization: "Bearer unit-8-token" } }); assert.equal(status.status, 200); assert.equal((await status.json() as { phase: string }).phase, "ocr_running"); const unavailable = createApp({ ingestService: ingestService as never, catalog: undefined, startReconciler: false }).listen(0); try { const unavailableAddress = unavailable.address(); assert.ok(unavailableAddress && typeof unavailableAddress === "object"); const response = await fetch(`http://127.0.0.1:${unavailableAddress.port}/ingestions/${versionId}`, { headers: { authorization: "Bearer unit-8-token" } }); assert.equal(response.status, 503); assert.equal((await response.json() as { code: string }).code, "CATALOG_UNAVAILABLE"); const retrieval = await fetch(`http://127.0.0.1:${unavailableAddress.port}/retrieve`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ query: "previous active content", mode: "documental", intent: "specific" }) }); assert.equal(retrieval.status, 503); assert.equal((await retrieval.json() as { code: string }).code, "CATALOG_UNAVAILABLE"); } finally { unavailable.close(); } } finally { server.close(); setEnvFlag("knowledgeLifecycleEnforced", previous.knowledgeLifecycleEnforced); setEnvFlag("lifecycleAdminToken", previous.lifecycleAdminToken); setEnvFlag("ocrIngestEnabled", previous.ocrIngestEnabled); } }); test("status reports all native and OCR documents and fails the version when one document fails", async () => { let query = 0; const pool = { async query() { query += 1; if (query === 1) return { rowCount: 1, rows: [{ source_id: "src:mixed", state: "failed", error_code: "OCR_QUALITY_BLOCKED", error_detail: "Visible ink failed OCR quality" }] }; if (query === 2) return { rowCount: 2, rows: [ { document_id: "doc:native", index_state: "indexing", job_state: null, completed_pages: null, requested_pages: null }, { document_id: "doc:ocr", index_state: "failed", job_state: "failed", completed_pages: 0, requested_pages: [2] } ] }; return { rowCount: 3, rows: [ { document_id: "doc:native", page_number: 1, extraction_method: "native", blocked_reason: null }, { document_id: "doc:ocr", page_number: 1, extraction_method: "native", blocked_reason: null }, { document_id: "doc:ocr", page_number: 2, extraction_method: "ocr", blocked_reason: "OCR_QUALITY_BLOCKED" } ] }; } }; const status = await new CatalogRepository(pool as never).getIngestionStatus("version-mixed") as { phase: string; documents: Array<{ documentId: string; completedPages: number; totalPages: number }>; error: { code: string } }; assert.equal(status.phase, "failed"); assert.deepEqual(status.documents.map(({ documentId, completedPages, totalPages }) => [documentId, completedPages, totalPages]), [ ["doc:native", 1, 1], ["doc:ocr", 1, 2] ]); assert.equal(status.error.code, "OCR_QUALITY_BLOCKED"); }); test("status counts an OCR-requested blank page only once", async () => { let query = 0; const pool = { async query() { query += 1; if (query === 1) return { rowCount: 1, rows: [{ source_id: "src:mixed", state: "review_required", error_code: null, error_detail: null }] }; if (query === 2) return { rowCount: 1, rows: [{ document_id: "doc:mixed", index_state: "indexing", job_state: "succeeded", completed_pages: 1, requested_pages: [2] }] }; return { rowCount: 2, rows: [ { document_id: "doc:mixed", page_number: 1, extraction_method: "native", blocked_reason: null }, { document_id: "doc:mixed", page_number: 2, extraction_method: "blank", blocked_reason: null } ] }; } }; const status = await new CatalogRepository(pool as never).getIngestionStatus("version-mixed") as { documents: Array<{ completedPages: number; totalPages: number }> }; assert.deepEqual(status.documents.map(({ completedPages, totalPages }) => [completedPages, totalPages]), [[2, 2]]); });