import assert from "node:assert/strict"; import { access, mkdir, mkdtemp, rm } 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 { KnowledgeLifecycleReconciler } from "../../src/modules/catalog/reconciler.js"; import { CatalogRepository } from "../../src/modules/catalog/repository.js"; import { OcrRetentionService, type OcrRetentionCandidate, type OcrRetentionStore } from "../../src/modules/ocr/retention.js"; const activeId = "11111111-1111-4111-8111-111111111111"; const failedId = "22222222-2222-4222-8222-222222222222"; const deletingId = "33333333-3333-4333-8333-333333333333"; function candidate(versionId: string, state: OcrRetentionCandidate["state"], artifactState: OcrRetentionCandidate["artifactState"] = "present"): OcrRetentionCandidate { return { versionId, sourceId: `source:${versionId}`, state, artifactState }; } test("repository retention selection applies exact TTLs and CAS-protects active versions", async () => { const queries: Array<{ sql: string; params?: unknown[] }> = []; const pool = { async query(sql: string, params?: unknown[]) { queries.push({ sql, params }); if (sql.includes("SELECT v.version_id")) return { rowCount: 1, rows: [{ version_id: failedId, source_id: "source:failed", state: "failed", artifact_state: "present" }] }; return { rowCount: 1, rows: [] }; } }; const repository = new CatalogRepository(pool as never); const now = new Date("2026-09-15T12:00:00.000Z"); assert.deepEqual(await repository.listOcrRetentionCandidates(now), [{ ...candidate(failedId, "failed"), sourceId: "source:failed" }]); assert.equal(await repository.expireOcrReview("source:review", activeId, now), true); assert.equal(await repository.claimOcrRetentionDeletion("source:failed", failedId, "failed", "present"), true); assert.equal(await repository.completeOcrRetentionDeletion("source:failed", failedId, "failed"), true); assert.match(queries[0]!.sql, /review_required.*30 days.*failed.*rejected.*7 days.*superseded.*30 days/s); assert.match(queries[0]!.sql, /active_version_id IS DISTINCT FROM v\.version_id/); assert.deepEqual(queries[0]!.params, [now]); assert.match(queries[1]!.sql, /state = 'rejected'.*REVIEW_EXPIRED.*retention_due_at/s); assert.match(queries[2]!.sql, /SET artifact_state = 'retention_deleting'.*state = \$3.*artifact_state = \$4/s); assert.match(queries[2]!.sql, /active_version_id IS DISTINCT FROM v\.version_id/); assert.match(queries[3]!.sql, /artifact_state = 'retention_deleted'.*artifact_state = 'retention_deleting'/s); }); test("activation cannot race an artifact deletion already claimed by retention", async () => { const pool = { async query(sql: string) { if (sql.includes("maintenance")) return { rowCount: 1, rows: [{ maintenance: false }] }; if (sql.includes("FROM rag_sources")) return { rowCount: 1, rows: [{ active_version_id: null }] }; return { rowCount: 1, rows: [{ version_id: failedId, source_id: "source:failed", version_number: 2, previous_version_id: null, state: "ready", tags: [], source_content_hash: "hash", processing_fingerprint: "fingerprint", metadata_hash: "metadata", embedding_provider: "test", embedding_model: "test", embedding_dimensions: 3, expected_document_count: 1, expected_point_count: 1, verified_point_count: 1, qdrant_collection: "rag", artifact_state: "retention_deleting" }] }; }, async connect() { return { query: this.query, release() {} }; } }; const repository = new CatalogRepository(pool as never); await assert.rejects(repository.activateVersion("source:failed", failedId, null), (error) => error instanceof Error && (error as { code?: string }).code === "VERSION_ARTIFACTS_UNAVAILABLE"); }); test("retention expires review before removal and preserves active artifacts", async (context) => { const root = await mkdtemp(path.join(os.tmpdir(), "rag-retention-")); context.after(() => rm(root, { recursive: true, force: true })); await Promise.all([activeId, failedId].map((id) => mkdir(path.join(root, id)))); const calls: string[] = []; const store: OcrRetentionStore = { async listOcrRetentionCandidates() { return [candidate(activeId, "active"), candidate(failedId, "failed"), candidate(deletingId, "review_required")]; }, async withVersionTryLock(versionId, handler) { calls.push(`lock:${versionId}`); return handler(); }, async expireOcrReview(_sourceId, versionId) { calls.push(`expire:${versionId}`); return true; }, async claimOcrRetentionDeletion(_sourceId, versionId) { calls.push(`claim:${versionId}`); return true; }, async completeOcrRetentionDeletion(_sourceId, versionId) { calls.push(`complete:${versionId}`); return true; } }; assert.deepEqual(await new OcrRetentionService(store, root).runOnce(), { expired: 1, deleted: 1, resumed: 0, skipped: 1 }); await access(path.join(root, activeId)); await assert.rejects(access(path.join(root, failedId))); assert.equal(calls.some((call) => call.includes(activeId)), false); assert.equal(calls.includes(`expire:${deletingId}`), true); }); test("interrupted deletion resumes idempotently without affecting another version", async (context) => { const root = await mkdtemp(path.join(os.tmpdir(), "rag-retention-resume-")); context.after(() => rm(root, { recursive: true, force: true })); await Promise.all([deletingId, activeId].map((id) => mkdir(path.join(root, id)))); let completed = false; let failCompletion = true; const store: OcrRetentionStore = { async listOcrRetentionCandidates() { return completed ? [] : [candidate(deletingId, "rejected", "retention_deleting")]; }, async withVersionTryLock(_versionId, handler) { return handler(); }, async expireOcrReview() { return false; }, async claimOcrRetentionDeletion() { return true; }, async completeOcrRetentionDeletion() { if (failCompletion) { failCompletion = false; throw new Error("simulated restart"); } completed = true; return true; } }; const retention = new OcrRetentionService(store, root); await assert.rejects(retention.runOnce(), /simulated restart/); await assert.rejects(access(path.join(root, deletingId))); assert.deepEqual(await retention.runOnce(), { expired: 0, deleted: 0, resumed: 1, skipped: 0 }); assert.deepEqual(await retention.runOnce(), { expired: 0, deleted: 0, resumed: 0, skipped: 0 }); await access(path.join(root, activeId)); }); test("OCR flag off keeps native HTTP synchronous and hides candidate surfaces", async (context) => { const previous = { enabled: env.ocrIngestEnabled, lifecycle: env.knowledgeLifecycleEnforced, token: env.lifecycleAdminToken }; Object.assign(env, { ocrIngestEnabled: false, knowledgeLifecycleEnforced: true, lifecycleAdminToken: "retention-token" }); context.after(() => Object.assign(env, previous)); let candidateReads = 0; const native = { accepted: true, sourceId: "source:native", versionId: "native-1", state: "active" }; const app = createApp({ ingestService: { async ingest() { return native; }, async cleanup() { return { deleted: 0 }; } } as never, catalog: { async getIngestionStatus() { candidateReads += 1; return {}; } } as never, reviewService: { async view() { candidateReads += 1; return {}; }, async approve() { candidateReads += 1; return {} as never; }, async reject() { candidateReads += 1; return {} as never; } }, indexingService: { async index() { candidateReads += 1; return {} as never; } }, startReconciler: false }); const server = app.listen(0); context.after(() => server.close()); const address = server.address(); assert.ok(address && typeof address === "object"); const base = `http://127.0.0.1:${address.port}`; assert.equal((await fetch(`${base}/ingest`, { method: "POST", headers: { "content-type": "application/json" }, body: "{}" })).status, 201); for (const [suffix, method] of [["", "GET"], ["/review", "GET"], ["/approve", "POST"], ["/reject", "POST"]] as const) { assert.equal((await fetch(`${base}/ingestions/${failedId}${suffix}`, { method, headers: { authorization: "Bearer retention-token", "content-type": "application/json" }, body: method === "POST" ? "{}" : undefined })).status, 404); } assert.equal(candidateReads, 0); }); test("reconciler runs retention inside its exclusive global lock", async () => { let locked = false; const catalog = { async withGlobalTryLock(_name: string, handler: () => Promise) { locked = true; try { return await handler(); } finally { locked = false; } }, async resolveActiveVersions() { return []; }, async validateActiveInvariant() { return []; }, async listOrphanedIndexingCandidates() { return []; } }; const retention = { async runOnce() { assert.equal(locked, true); return { expired: 0, deleted: 0, resumed: 0, skipped: 0 }; } }; const reconciler = new KnowledgeLifecycleReconciler(catalog as never, {} as never, undefined, retention); assert.equal((await reconciler.runOnce()).ocrRetentionDeleted, 0); });