import test from "node:test"; import assert from "node:assert/strict"; import { access, mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; import os from "node:os"; import path from "node:path"; import { createApp } from "../src/app.js"; import { openApiDocument } from "../src/api/openapi.js"; import { env } from "../src/config/env.js"; import { CatalogError } from "../src/modules/catalog/errors.js"; import { buildLegacyMigrationPlan } from "../src/modules/catalog/legacy-migration.js"; import { KnowledgeLifecycleReconciler } from "../src/modules/catalog/reconciler.js"; import { CatalogRepository } from "../src/modules/catalog/repository.js"; import { activateMigratedVersion, createCatalogVersion, rollbackMigratedBatch } from "../src/scripts/migrate-legacy-lifecycle.js"; import { IngestService } from "../src/modules/ingest/service.js"; import { RetrieveService } from "../src/modules/retrieve/service.js"; import type { EmbeddingProvider } from "../src/modules/embeddings/provider.js"; import type { VectorStoreClient } from "../src/modules/vectorstore/client.js"; function setEnvFlag(name: keyof typeof env, value: unknown) { (env as unknown as Record)[name] = value; } async function withEnvFlags(updates: Partial>, handler: () => Promise) { const previous = new Map(); for (const [name, value] of Object.entries(updates) as Array<[keyof typeof env, unknown]>) { previous.set(name, env[name]); setEnvFlag(name, value); } try { await handler(); } finally { for (const [name, value] of previous) { setEnvFlag(name, value); } } } async function withTempMarkdown(content: string, handler: (filePath: string) => Promise) { const dir = await mkdtemp(path.join(os.tmpdir(), "rag-lifecycle-test-")); try { const filePath = path.join(dir, "source.md"); await writeFile(filePath, content, "utf8"); await handler(filePath); } finally { await rm(dir, { recursive: true, force: true }); } } function fakeEmbeddingProvider(): EmbeddingProvider { return { providerName: "test-provider", modelName: "test-model", dimensions: 3, async embed(input: string[]) { return input.map(() => [0.1, 0.2, 0.3]); } }; } function fakeVectorStore() { const chunks: unknown[] = []; const vectorStore: Partial = { kind: "fake", async upsert(input) { chunks.push(...input); }, async countVersionPoints() { return chunks.length; }, async validateCollectionDimensions() { return undefined; } }; return { vectorStore: vectorStore as VectorStoreClient, chunks }; } function fakeTransactionPool(handler: (sql: string, params?: unknown[]) => Promise<{ rowCount: number; rows: unknown[] }>) { const client = { async query(sql: string, params?: unknown[]) { if (/^(BEGIN|COMMIT|ROLLBACK|SET CONSTRAINTS)/.test(sql.trim())) { return { rowCount: 0, rows: [] }; } return handler(sql, params); }, release() { return undefined; } }; return { async connect() { return client; }, async query(sql: string, params?: unknown[]) { return handler(sql, params); } }; } test("activation failure leaves a newly indexed version ready instead of failed", async () => { await withEnvFlags({ knowledgeLifecycleEnforced: true }, async () => withTempMarkdown("# Title\n\nUseful content", async (filePath) => { const { vectorStore } = fakeVectorStore(); const calls: string[] = []; const catalog = { async beginAttempt() { return "attempt-1"; }, async updateAttempt() { return undefined; }, async withSourceLock(_sourceId: string, handler: () => Promise) { return handler(); }, async findReusableVersion() { return undefined; }, async createPendingVersion() { return { versionId: "version-1", versionNumber: 1, state: "pending", previousVersionId: null }; }, async markIndexing() { calls.push("indexing"); }, async markReady() { calls.push("ready"); return { versionId: "version-1", versionNumber: 1, state: "ready", previousVersionId: null }; }, async activateVersion() { throw new CatalogError("precondition failed", 409, "ACTIVE_VERSION_PRECONDITION_FAILED"); }, async markFailed() { calls.push("failed"); } }; const service = new IngestService(fakeEmbeddingProvider(), vectorStore, catalog as never); await assert.rejects( service.ingest({ sourceType: "file", sourceRef: "source.md", readPath: filePath, activate: true, expectedActiveVersionId: "other" }), /precondition failed/ ); assert.deepEqual(calls, ["indexing", "ready"]); })); }); test("active no-op validates the expected active version precondition", async () => { await withEnvFlags({ knowledgeLifecycleEnforced: true }, async () => withTempMarkdown("# Title\n\nUseful content", async (filePath) => { const { vectorStore, chunks } = fakeVectorStore(); let preconditionChecked = false; const catalog = { async beginAttempt() { return "attempt-1"; }, async updateAttempt() { return undefined; }, async withSourceLock(_sourceId: string, handler: () => Promise) { return handler(); }, async findReusableVersion() { return { versionId: "active-1", versionNumber: 1, state: "active", previousVersionId: null, expectedDocumentCount: 1 }; }, async assertActiveVersionPrecondition() { preconditionChecked = true; throw new CatalogError("Active version precondition failed", 409, "ACTIVE_VERSION_PRECONDITION_FAILED"); }, async markFailed() { throw new Error("markFailed must not run for no-op precondition errors"); } }; const service = new IngestService(fakeEmbeddingProvider(), vectorStore, catalog as never); await assert.rejects( service.ingest({ sourceType: "file", sourceRef: "source.md", readPath: filePath, activate: true, expectedActiveVersionId: "wrong" }), /precondition failed/i ); assert.equal(preconditionChecked, true); assert.equal(chunks.length, 0); })); }); test("/sources fails closed and /cleanup refuses deletion when lifecycle enforcement is enabled", async () => { setEnvFlag("knowledgeLifecycleEnforced", true); setEnvFlag("lifecycleAdminToken", "test-token"); const server = createApp().listen(0); try { const address = server.address(); assert.ok(address && typeof address === "object"); const baseUrl = `http://127.0.0.1:${address.port}`; const sources = await fetch(`${baseUrl}/sources`); assert.equal(sources.status, 503); const cleanupWithoutAdmin = await fetch(`${baseUrl}/cleanup`, { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify({ scope: { sourceId: "src:default:folder:docs" } }) }); assert.equal(cleanupWithoutAdmin.status, 401); const cleanup = await fetch(`${baseUrl}/cleanup`, { method: "POST", headers: { "Content-Type": "application/json", Authorization: "Bearer test-token" }, body: JSON.stringify({ scope: { sourceId: "src:default:folder:docs" } }) }); assert.equal(cleanup.status, 410); } finally { server.close(); setEnvFlag("knowledgeLifecycleEnforced", false); setEnvFlag("lifecycleAdminToken", ""); } }); test("retrieval fails closed when lifecycle is enforced without a catalog", async () => { await withEnvFlags({ knowledgeLifecycleEnforced: true }, async () => { const service = new RetrieveService(fakeEmbeddingProvider(), fakeVectorStore().vectorStore); await assert.rejects( service.retrieve("documental", "specific", "hello"), (error) => error instanceof CatalogError && error.statusCode === 503 && error.code === "CATALOG_UNAVAILABLE" ); }); }); test("repository state guards do not mark non-pending/indexing versions as failed", async () => { const sqlCalls: string[] = []; const pool = { async query(sql: string) { sqlCalls.push(sql); return sql.includes("UPDATE rag_source_versions") ? { rowCount: 0, rows: [] } : { rowCount: 1, rows: [] }; } }; const repository = new CatalogRepository(pool as never); await repository.markFailed("ready-version", "TEST", "must not change ready"); assert.equal(sqlCalls.some((sql) => sql.includes("UPDATE rag_version_documents")), false); assert.match(sqlCalls[0], /state IN \('pending', 'indexing'\)/); }); test("reconciler reports invariant violations and recovers orphaned versions", async () => { const calls: string[] = []; const catalog = { async withGlobalTryLock(_lockName: string, handler: () => Promise) { return handler(); }, async resolveActiveVersions() { return [{ sourceId: "src:one", sourceVersionId: "v1", expectedPointCount: 2, sourceVersionNumber: 1, embeddingDimensions: 3, qdrantCollection: env.qdrantCollection }]; }, async validateActiveInvariant() { return ["src:bad"]; }, async listOrphanedIndexingCandidates() { return []; }, async markNeedsReingest(sourceId: string) { calls.push(sourceId); } }; const vectorStore = { async validateCollectionDimensions() { return undefined; }, async countVersionPoints() { return 1; } }; const reconciler = new KnowledgeLifecycleReconciler(catalog as never, vectorStore as never); const status = await reconciler.runOnce(); assert.equal(status.ok, false); assert.deepEqual(status.inconsistentSources, ["src:one"]); assert.deepEqual(status.invariantViolations, ["src:bad"]); assert.equal(status.orphanedVersionsRecovered, 0); assert.deepEqual(calls, ["src:one"]); }); test("reconciler recovers stale complete indexing versions to ready", async () => { const calls: string[] = []; const catalog = { async withGlobalTryLock(_lockName: string, handler: () => Promise) { return handler(); }, async resolveActiveVersions() { return []; }, async validateActiveInvariant() { return []; }, async listOrphanedIndexingCandidates() { return [{ sourceId: "src:one", versionId: "v1", expectedDocumentCount: 1, expectedPointCount: 2, embeddingDimensions: 3, qdrantCollection: env.qdrantCollection }]; }, async withSourceTryLock(_sourceId: string, handler: () => Promise) { calls.push("source-lock"); return handler(); }, async withVersionTryLock(_versionId: string, handler: () => Promise) { calls.push("version-lock"); return handler(); }, async hasCompleteVersionDocuments() { return true; }, async recoverVersionReady() { calls.push("ready"); return true; }, async markFailed() { calls.push("failed"); } }; const vectorStore = { async validateCollectionDimensions() { calls.push("dimensions"); }, async countVersionPoints() { calls.push("count"); return 2; } }; const reconciler = new KnowledgeLifecycleReconciler(catalog as never, vectorStore as never); const status = await reconciler.runOnce(); assert.equal(status.ok, true); assert.equal(status.orphanedVersionsRecovered, 1); assert.equal(status.orphanedVersionsFailed, 0); assert.deepEqual(calls, ["source-lock", "version-lock", "dimensions", "count", "ready"]); }); test("reconciler marks stale partial indexing versions as failed", async () => { const calls: string[] = []; const catalog = { async withGlobalTryLock(_lockName: string, handler: () => Promise) { return handler(); }, async resolveActiveVersions() { return []; }, async validateActiveInvariant() { return []; }, async listOrphanedIndexingCandidates() { return [{ sourceId: "src:one", versionId: "v1", expectedDocumentCount: 1, expectedPointCount: 2, embeddingDimensions: 3, qdrantCollection: env.qdrantCollection }]; }, async withSourceTryLock(_sourceId: string, handler: () => Promise) { return handler(); }, async withVersionTryLock(_versionId: string, handler: () => Promise) { return handler(); }, async hasCompleteVersionDocuments() { return true; }, async recoverVersionReady() { calls.push("ready"); return true; }, async markFailed() { calls.push("failed"); } }; const vectorStore = { async validateCollectionDimensions() { return undefined; }, async countVersionPoints() { return 1; } }; const reconciler = new KnowledgeLifecycleReconciler(catalog as never, vectorStore as never); const status = await reconciler.runOnce(); assert.equal(status.orphanedVersionsRecovered, 0); assert.equal(status.orphanedVersionsFailed, 1); assert.deepEqual(calls, ["failed"]); }); test("reconciler does not touch stale versions while a source lock is held", async () => { const calls: string[] = []; const catalog = { async withGlobalTryLock(_lockName: string, handler: () => Promise) { return handler(); }, async resolveActiveVersions() { return []; }, async validateActiveInvariant() { return []; }, async listOrphanedIndexingCandidates() { return [{ sourceId: "src:one", versionId: "v1", expectedDocumentCount: 1, expectedPointCount: 2, embeddingDimensions: 3, qdrantCollection: env.qdrantCollection }]; }, async withSourceTryLock() { return undefined; }, async markFailed() { calls.push("failed"); } }; const reconciler = new KnowledgeLifecycleReconciler(catalog as never, fakeVectorStore().vectorStore); const status = await reconciler.runOnce(); assert.equal(status.orphanedVersionsRecovered, 0); assert.equal(status.orphanedVersionsFailed, 0); assert.deepEqual(calls, []); }); test("reconciler skips when another reconciler holds the advisory lock", async () => { const catalog = { async withGlobalTryLock() { return undefined; }, async resolveActiveVersions() { throw new Error("must not run without the lock"); } }; const reconciler = new KnowledgeLifecycleReconciler(catalog as never, fakeVectorStore().vectorStore); const status = await reconciler.runOnce(); assert.equal(status.ok, true); assert.equal(status.inconsistentSources.length, 0); }); test("retrieval rejects active versions from a different Qdrant collection", async () => { await withEnvFlags({ knowledgeLifecycleEnforced: true }, async () => { const catalog = { async resolveActiveVersions() { return [{ sourceId: "src:one", sourceVersionId: "v1", sourceVersionNumber: 1, expectedPointCount: 1, embeddingDimensions: 3, qdrantCollection: "other_collection" }]; }, async markNeedsReingest() { return undefined; }, async withVersionSharedLocks(_versionIds: string[], handler: () => Promise) { return handler(); } }; const service = new RetrieveService(fakeEmbeddingProvider(), fakeVectorStore().vectorStore, catalog as never); await assert.rejects( service.retrieve("documental", "specific", "hello"), (error) => error instanceof CatalogError && error.code === "SOURCE_VERSION_INCONSISTENT" ); }); }); test("retrieval holds shared locks while validating counts and searching Qdrant", async () => { await withEnvFlags({ knowledgeLifecycleEnforced: true }, async () => { const calls: string[] = []; const catalog = { async resolveActiveVersions() { calls.push("resolve"); return [{ sourceId: "src:one", sourceVersionId: "v1", sourceVersionNumber: 1, expectedPointCount: 1, embeddingDimensions: 3, qdrantCollection: env.qdrantCollection }]; }, async withVersionSharedLocks(_versionIds: string[], handler: () => Promise) { calls.push("lock:start"); const result = await handler(); calls.push("lock:end"); return result; }, async markNeedsReingest() { calls.push("needs-reingest"); } }; const vectorStore = { kind: "fake", async validateCollectionDimensions() { calls.push("dimensions"); }, async countVersionPoints() { calls.push("count"); return 1; }, async search() { calls.push("search"); return []; }, async browseScope() { return []; } }; const service = new RetrieveService(fakeEmbeddingProvider(), vectorStore as never, catalog as never); await service.retrieve("documental", "specific", "hello"); assert.deepEqual(calls, ["resolve", "lock:start", "dimensions", "count", "search", "lock:end"]); }); }); test("lifecycle ingest rejects vectors that do not match configured embedding dimensions before upsert", async () => { await withEnvFlags({ knowledgeLifecycleEnforced: true }, async () => withTempMarkdown("# Title\n\nUseful content", async (filePath) => { const calls: string[] = []; const provider: EmbeddingProvider = { providerName: "test-provider", modelName: "test-model", dimensions: 3, async embed(input: string[]) { return input.map(() => [0.1, 0.2]); } }; const vectorStore = { kind: "fake", async upsert() { calls.push("upsert"); }, async countVersionPoints() { return 0; } }; const catalog = { async beginAttempt() { return "attempt-1"; }, async updateAttempt() { return undefined; }, async withSourceLock(_sourceId: string, handler: () => Promise) { return handler(); }, async findReusableVersion() { return undefined; }, async createPendingVersion() { return { versionId: "version-1", versionNumber: 1, state: "pending", previousVersionId: null }; }, async markIndexing() { calls.push("indexing"); }, async markReady() { calls.push("ready"); return { versionId: "version-1", versionNumber: 1, state: "ready", previousVersionId: null }; }, async markFailed() { calls.push("failed"); } }; const service = new IngestService(provider, vectorStore as never, catalog as never); await assert.rejects( service.ingest({ sourceType: "file", sourceRef: "source.md", readPath: filePath, activate: true, expectedActiveVersionId: null }), (error) => error instanceof CatalogError && error.code === "EMBEDDING_DIMENSIONS_INVALID" ); assert.deepEqual(calls, ["indexing", "failed"]); })); }); test("cleanup respects disabled ingest writes even on the legacy path", async () => { await withEnvFlags({ ingestWritesEnabled: false, knowledgeLifecycleEnforced: false }, async () => { const service = new IngestService(fakeEmbeddingProvider(), fakeVectorStore().vectorStore); await assert.rejects( service.cleanup({ sourceId: "src:default:folder:docs" }), (error) => error instanceof CatalogError && error.code === "INGEST_WRITES_DISABLED" ); }); }); test("markPurged fails when the version is no longer in purging state", async () => { const pool = fakeTransactionPool(async (sql) => { if (sql.includes("UPDATE rag_source_versions")) { return { rowCount: 0, rows: [] }; } return { rowCount: 1, rows: [] }; }); const repository = new CatalogRepository(pool as never); await assert.rejects( repository.markPurged("version-1"), (error) => error instanceof CatalogError && error.code === "INVALID_VERSION_STATE" ); }); test("version purge removes durable OCR artifacts before marking the version purged", async (context) => { const artifactRoot = await mkdtemp(path.join(os.tmpdir(), "rag-purge-artifacts-")); const versionDirectory = path.join(artifactRoot, "version-1"); await mkdir(versionDirectory); await writeFile(path.join(versionDirectory, "manifest.json"), "{}"); context.after(() => rm(artifactRoot, { recursive: true, force: true })); await withEnvFlags({ lifecycleAdminToken: "purge-token", ocrArtifactRoot: artifactRoot }, async () => { const calls: string[] = []; const catalog = { async withSourceLock(_sourceId: string, handler: () => Promise) { return handler(); }, async withVersionExclusiveLock(_versionId: string, handler: () => Promise) { return handler(); }, async markPurging() { calls.push("purging"); }, async markPurged() { await assert.rejects(access(versionDirectory)); calls.push("purged"); } }; const vectorStore: Partial = { kind: "fake", async deleteVersionPoints() { calls.push("delete-points"); return 1; }, async countVersionPoints() { calls.push("count-points"); return 0; } }; const server = createApp({ catalog: catalog as never, vectorStore: vectorStore as VectorStoreClient, startReconciler: false }).listen(0); context.after(() => server.close()); const address = server.address(); assert.ok(address && typeof address === "object"); const response = await fetch(`http://127.0.0.1:${address.port}/sources/source-1/versions/version-1`, { method: "DELETE", headers: { authorization: "Bearer purge-token" } }); assert.equal(response.status, 200); assert.deepEqual(calls, ["purging", "delete-points", "count-points", "purged"]); }); }); test("legacy migration apply reuses a catalog version created before payload update", async () => { const source = buildLegacyMigrationPlan([{ id: "p1", payload: { source_id: "src:default:file:source.md", source_type: "file", source_ref: "source.md", document_id: "doc:src:default:file:source.md:source.md", document_key: "source.md", chunk_id: "chk:doc:src:default:file:source.md:source.md:documental:0001", content: "content", embedding_provider: "openrouter", embedding_model: "qwen", embedding_dimensions: 3 } }]).sources[0]; const sqlCalls: string[] = []; const pool = fakeTransactionPool(async (sql) => { sqlCalls.push(sql); if (sql.includes("FROM rag_legacy_migration_items")) { return { rowCount: 1, rows: [{ version_id: "version-1", expected_point_count: 1 }] }; } if (sql.includes("FROM rag_source_versions v")) { return { rowCount: 1, rows: [{ version_id: "version-1", version_number: "1", expected_point_count: "1", source_type: "file", source_ref: "source.md", qdrant_collection: env.qdrantCollection }] }; } return { rowCount: 1, rows: [] }; }); const result = await createCatalogVersion(pool as never, "batch-1", source, env.qdrantCollection); assert.deepEqual(result, { versionId: "version-1", versionNumber: 1 }); assert.equal(sqlCalls.some((sql) => sql.includes("INSERT INTO rag_sources")), false); }); test("legacy migration activation is idempotent after payload update and prior activation", async () => { const sqlCalls: string[] = []; const pool = fakeTransactionPool(async (sql) => { sqlCalls.push(sql); if (sql.includes("SELECT v.state")) { return { rowCount: 1, rows: [{ state: "active", active_version_id: "version-1" }] }; } return { rowCount: 1, rows: [] }; }); await activateMigratedVersion(pool as never, "src:default:file:source.md", "version-1", 1); assert.equal(sqlCalls.some((sql) => sql.includes("SET state = 'active'")), false); }); test("legacy rollback prepares the whole batch before restoring and validates restored legacy counts before catalog cleanup", async () => { const calls: string[] = []; const pool = fakeTransactionPool(async (sql) => { if (sql.includes("SELECT state, snapshot_ref")) { calls.push("batch-lock"); return { rowCount: 1, rows: [{ state: "verified", snapshot_ref: "snapshot-url" }] }; } if (sql.includes("SELECT source_id, version_id, expected_point_count")) { calls.push("items-lock"); return { rowCount: 1, rows: [{ source_id: "src:default:file:source.md", version_id: "version-1", expected_point_count: 1 }] }; } if (sql.includes("SET active_version_id = NULL")) calls.push("active-null"); if (sql.includes("SET state = 'purging'")) calls.push("purging"); if (sql.includes("COUNT(*) AS count")) { calls.push("zero-active-check"); return { rowCount: 1, rows: [{ count: "0" }] }; } if (sql.includes("DELETE FROM rag_version_documents")) calls.push("delete-docs"); if (sql.includes("UPDATE rag_legacy_migration_items SET version_id = NULL")) calls.push("items-rolled-back"); return { rowCount: 1, rows: [] }; }); const qdrant = { async recoverSnapshot() { calls.push("restore-snapshot"); return true; }, async scroll() { calls.push("verify-legacy-count"); return { points: [{ id: "p1", payload: { source_id: "src:default:file:source.md" } }], next_page_offset: undefined }; } }; const result = await rollbackMigratedBatch(pool as never, qdrant as never, "batch-1", "snapshot-url"); assert.equal(result.ok, true); assert.ok(calls.indexOf("zero-active-check") < calls.indexOf("restore-snapshot")); assert.ok(calls.indexOf("verify-legacy-count") < calls.indexOf("delete-docs")); assert.ok(calls.includes("items-rolled-back")); }); test("OpenAPI documents /sources 503 responses", () => { const sources = openApiDocument.paths["/sources"].get.responses; assert.ok("503" in sources); }); test("legacy migration dry-run blocks weak document identity", () => { const plan = buildLegacyMigrationPlan([ { id: "p1", payload: { source_id: "src:default:folder:docs", source_ref: "docs", source_type: "folder", document_id: "doc:src:default:folder:docs:readme.md", chunk_id: "chk:doc:src:default:folder:docs:readme.md:documental:0001", content: "content", embedding_provider: "openrouter", embedding_model: "qwen", embedding_dimensions: 3 } } ]); assert.equal(plan.ok, false); assert.match(plan.blockedReasons.join("\n"), /weak document identity/); });