rag-service/src/modules/catalog/reconciler.ts

135 lines
5.3 KiB
TypeScript

import { env } from "../../config/env.js";
import type { VectorStoreClient } from "../vectorstore/client.js";
import type { CatalogRepository } from "./repository.js";
import type { OcrDispatcher } from "../ocr/dispatcher.js";
import type { OcrRetentionService } from "../ocr/retention.js";
export interface ReconcilerStatus {
ok: boolean;
lastRunAt?: string;
lastError?: string;
inconsistentSources: string[];
orphanedVersionsRecovered?: number;
orphanedVersionsFailed?: number;
ocrLeasesRecovered?: number;
ocrCandidatesRecovered?: number;
ocrRetentionDeleted?: number;
invariantViolations?: string[];
}
export class KnowledgeLifecycleReconciler {
private status: ReconcilerStatus = { ok: true, inconsistentSources: [] };
private timer: NodeJS.Timeout | undefined;
constructor(
private readonly catalog: CatalogRepository | undefined,
private readonly vectorStore: VectorStoreClient,
private readonly ocrDispatcher?: Pick<OcrDispatcher, "recoverExpiredLeases" | "dispatchAvailable">,
private readonly ocrRetention?: Pick<OcrRetentionService, "runOnce">
) {}
getStatus(): ReconcilerStatus {
return this.status;
}
start(): void {
if (!env.knowledgeLifecycleEnforced || !this.catalog || this.timer) {
return;
}
void this.runOnce();
this.timer = setInterval(() => {
void this.runOnce();
}, env.lifecycleReconcileIntervalMs);
this.timer.unref();
}
async runOnce(): Promise<ReconcilerStatus> {
if (!this.catalog) {
this.status = { ok: false, lastRunAt: new Date().toISOString(), lastError: "Knowledge catalog is not configured", inconsistentSources: [] };
return this.status;
}
const catalog = this.catalog;
try {
const lockedStatus = await catalog.withGlobalTryLock("rag:knowledge-lifecycle:reconciler", async () => {
const activeVersions = await catalog.resolveActiveVersions();
const inconsistentSources: string[] = [];
const invariantViolations = await catalog.validateActiveInvariant();
const ocrLeasesRecovered = await this.ocrDispatcher?.recoverExpiredLeases() ?? 0;
await this.ocrDispatcher?.dispatchAvailable();
const retention = await this.ocrRetention?.runOnce();
const orphanRecovery = await this.recoverOrphanedVersions(catalog);
for (const version of activeVersions) {
if (version.qdrantCollection !== env.qdrantCollection) {
inconsistentSources.push(version.sourceId);
await catalog.markNeedsReingest(version.sourceId);
continue;
}
await this.vectorStore.validateCollectionDimensions(version.embeddingDimensions);
const count = await this.vectorStore.countVersionPoints(version.sourceVersionId);
if (count !== version.expectedPointCount) {
inconsistentSources.push(version.sourceId);
await catalog.markNeedsReingest(version.sourceId);
}
}
return {
ok: inconsistentSources.length === 0 && invariantViolations.length === 0,
lastRunAt: new Date().toISOString(),
inconsistentSources,
orphanedVersionsRecovered: orphanRecovery.ready,
orphanedVersionsFailed: orphanRecovery.failed,
ocrLeasesRecovered,
ocrRetentionDeleted: (retention?.deleted ?? 0) + (retention?.resumed ?? 0),
invariantViolations
};
});
this.status = lockedStatus ?? { ...this.status, lastRunAt: new Date().toISOString() };
return this.status;
} catch (error) {
this.status = {
ok: false,
lastRunAt: new Date().toISOString(),
lastError: error instanceof Error ? error.message : "Unknown reconciler error",
inconsistentSources: []
};
return this.status;
}
}
private async recoverOrphanedVersions(catalog: CatalogRepository): Promise<{ ready: number; failed: number }> {
const candidates = await catalog.listOrphanedIndexingCandidates(env.lifecycleIndexingStaleTimeoutMs);
let ready = 0;
let failed = 0;
for (const candidate of candidates) {
const recovered = await catalog.withSourceTryLock(candidate.sourceId, async () => catalog.withVersionTryLock(candidate.versionId, async () => {
const documentsComplete = await catalog.hasCompleteVersionDocuments(candidate.versionId, candidate.expectedDocumentCount, candidate.expectedPointCount);
const collectionMatches = candidate.qdrantCollection === env.qdrantCollection;
if (collectionMatches) {
await this.vectorStore.validateCollectionDimensions(candidate.embeddingDimensions);
}
const pointCount = collectionMatches ? await this.vectorStore.countVersionPoints(candidate.versionId) : 0;
if (collectionMatches && documentsComplete && pointCount === candidate.expectedPointCount) {
const changed = await catalog.recoverVersionReady(candidate.versionId, pointCount);
if (changed) {
return "ready" as const;
}
}
await catalog.markFailed(candidate.versionId, "ORPHANED_INDEXING_INCOMPLETE", "Stale pending/indexing version has incomplete documents, dimensions, collection, or point count.");
return "failed" as const;
}));
if (recovered === "ready") {
ready += 1;
} else if (recovered === "failed") {
failed += 1;
}
}
return { ready, failed };
}
}