rag-service/tests/lifecycle-services.test.ts

552 lines
25 KiB
TypeScript

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<string, unknown>)[name] = value;
}
async function withEnvFlags(updates: Partial<Record<keyof typeof env, unknown>>, handler: () => Promise<void>) {
const previous = new Map<keyof typeof env, unknown>();
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<void>) {
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<VectorStoreClient> = {
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<unknown>) { 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<unknown>) { 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<unknown>) { 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<unknown>) { 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<unknown>) { calls.push("source-lock"); return handler(); },
async withVersionTryLock(_versionId: string, handler: () => Promise<unknown>) { 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<unknown>) { 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<unknown>) { return handler(); },
async withVersionTryLock(_versionId: string, handler: () => Promise<unknown>) { 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<unknown>) { 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<unknown>) { 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<unknown>) {
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<unknown>) { 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<unknown>) { return handler(); },
async withVersionExclusiveLock(_versionId: string, handler: () => Promise<unknown>) { return handler(); },
async markPurging() { calls.push("purging"); },
async markPurged() {
await assert.rejects(access(versionDirectory));
calls.push("purged");
}
};
const vectorStore: Partial<VectorStoreClient> = {
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/);
});