import type { AvailableScope, ChunkMode, IngestSourceInput, SourceType, SourceVersionState, RetrieveScope } from "../../shared/types/rag.js"; import { withCatalogAdvisoryLock, withCatalogAdvisorySharedLocks, withCatalogTryAdvisoryLock, withTransaction, type PgPool, type PgPoolClient } from "./client.js"; import { CatalogError } from "./errors.js"; import { normalizeExpectedActiveVersion } from "./lifecycle.js"; export interface CatalogDocumentInput { documentId: string; documentKey: string; originalHash: string; originalHashKind: "bytes" | "legacy-derived"; contentHash: string | null; mimeType: string; title: string; chunkCount: number; } export interface CreateVersionInput { sourceId: string; sourceType: SourceType; sourceRef: string; originalManifestHash: string; sourceContentHash: string | null; processingFingerprint: string; metadataHash: string; tags: string[]; activateRequested: boolean; embeddingProvider: string; embeddingModel: string; embeddingDimensions: number; qdrantCollection: string; expectedDocumentCount: number; expectedPointCount: number; documents: CatalogDocumentInput[]; } export interface CatalogVersionRow { versionId: string; sourceId: string; versionNumber: number; previousVersionId: string | null; state: SourceVersionState; tags: string[]; sourceContentHash: string | null; processingFingerprint: string; metadataHash: string; embeddingProvider: string; embeddingModel: string; embeddingDimensions: number; expectedDocumentCount: number; expectedPointCount: number; verifiedPointCount: number; qdrantCollection: string; } export interface ActiveVersionScope { sourceId: string; sourceVersionId: string; sourceVersionNumber: number; expectedPointCount: number; embeddingDimensions: number; qdrantCollection: string; } export interface OrphanedVersionCandidate { sourceId: string; versionId: string; expectedDocumentCount: number; expectedPointCount: number; embeddingDimensions: number; qdrantCollection: string; } type VersionDbRow = { version_id: string; source_id: string; version_number: string | number; previous_version_id: string | null; state: SourceVersionState; tags: string[]; source_content_hash: string | null; processing_fingerprint: string; metadata_hash: string; embedding_provider: string; embedding_model: string; embedding_dimensions: string | number; expected_document_count: string | number; expected_point_count: string | number; verified_point_count: string | number; qdrant_collection: string; }; function toVersion(row: VersionDbRow): CatalogVersionRow { return { versionId: row.version_id, sourceId: row.source_id, versionNumber: Number(row.version_number), previousVersionId: row.previous_version_id, state: row.state, tags: row.tags ?? [], sourceContentHash: row.source_content_hash, processingFingerprint: row.processing_fingerprint, metadataHash: row.metadata_hash, embeddingProvider: row.embedding_provider, embeddingModel: row.embedding_model, embeddingDimensions: Number(row.embedding_dimensions), expectedDocumentCount: Number(row.expected_document_count), expectedPointCount: Number(row.expected_point_count), verifiedPointCount: Number(row.verified_point_count), qdrantCollection: row.qdrant_collection }; } export class CatalogRepository { constructor(private readonly pool: PgPool) {} async healthcheck(): Promise<{ ok: boolean; kind: "postgres" }> { try { await this.pool.query("SELECT 1"); return { ok: true, kind: "postgres" }; } catch { return { ok: false, kind: "postgres" }; } } async withSourceLock(sourceId: string, handler: () => Promise): Promise { return withCatalogAdvisoryLock(this.pool, `rag:source:${sourceId}`, handler); } async withSourceTryLock(sourceId: string, handler: () => Promise): Promise { return withCatalogTryAdvisoryLock(this.pool, `rag:source:${sourceId}`, handler); } async withVersionSharedLocks(versionIds: string[], handler: () => Promise): Promise { return withCatalogAdvisorySharedLocks(this.pool, versionIds.map((versionId) => `rag:version:${versionId}`), handler); } async withVersionExclusiveLock(versionId: string, handler: () => Promise): Promise { return withCatalogAdvisoryLock(this.pool, `rag:version:${versionId}`, handler); } async withVersionTryLock(versionId: string, handler: () => Promise): Promise { return withCatalogTryAdvisoryLock(this.pool, `rag:version:${versionId}`, handler); } async withGlobalTryLock(lockName: string, handler: () => Promise): Promise { return withCatalogTryAdvisoryLock(this.pool, lockName, handler); } async isMaintenanceEnabled(client?: PgPoolClient, lockRow = false): Promise { const executor = client ?? this.pool; const result = await executor.query<{ maintenance: boolean }>( `SELECT maintenance FROM rag_system_state WHERE singleton = true${lockRow ? " FOR UPDATE" : ""}` ); return Boolean(result.rows[0]?.maintenance); } async beginAttempt(source: IngestSourceInput & { sourceId: string }): Promise { return withTransaction(this.pool, async (client) => { if (await this.isMaintenanceEnabled(client, true)) { throw new CatalogError("Knowledge catalog is in maintenance mode", 503, "CATALOG_MAINTENANCE"); } const sourceResult = await client.query<{ source_type: SourceType; source_ref: string }>( `INSERT INTO rag_sources(source_id, source_type, source_ref) VALUES ($1, $2, $3) ON CONFLICT (source_id) DO NOTHING RETURNING source_type, source_ref`, [source.sourceId, source.sourceType, source.sourceRef] ); const existing = sourceResult.rowCount ? sourceResult.rows[0] : (await client.query<{ source_type: SourceType; source_ref: string }>( "SELECT source_type, source_ref FROM rag_sources WHERE source_id = $1", [source.sourceId] )).rows[0]; if (!existing || existing.source_type !== source.sourceType || existing.source_ref !== source.sourceRef) { throw new CatalogError("sourceId already exists with a different sourceType or sourceRef", 409, "SOURCE_IDENTITY_CONFLICT"); } const result = await client.query<{ attempt_id: string }>( `INSERT INTO rag_ingestion_attempts(source_id, state, input_locator) VALUES ($1, 'received', $2) RETURNING attempt_id`, [source.sourceId, source.readPath ?? source.sourceRef] ); return result.rows[0].attempt_id; }); } async updateAttempt(attemptId: string, values: { state: "processing" | "completed" | "failed"; versionId?: string | null; inputHash?: string | null; errorCode?: string; errorDetail?: string }): Promise { await this.pool.query( `UPDATE rag_ingestion_attempts SET state = $2, version_id = COALESCE($3, version_id), input_hash = COALESCE($4, input_hash), error_code = $5, error_detail = $6, completed_at = CASE WHEN $2 IN ('completed', 'failed') THEN now() ELSE completed_at END WHERE attempt_id = $1`, [attemptId, values.state, values.versionId ?? null, values.inputHash ?? null, values.errorCode ?? null, values.errorDetail ?? null] ); } async findReusableVersion(input: Pick): Promise { if (!input.sourceContentHash) { return undefined; } const result = await this.pool.query( `SELECT * FROM rag_source_versions WHERE source_id = $1 AND source_content_hash = $2 AND processing_fingerprint = $3 AND metadata_hash = $4 AND state NOT IN ('failed', 'rejected', 'purged') ORDER BY version_number DESC LIMIT 1`, [input.sourceId, input.sourceContentHash, input.processingFingerprint, input.metadataHash] ); return result.rows[0] ? toVersion(result.rows[0]) : undefined; } async createPendingVersion(input: CreateVersionInput): Promise { return withTransaction(this.pool, async (client) => { if (await this.isMaintenanceEnabled(client, true)) { throw new CatalogError("Knowledge catalog is in maintenance mode", 503, "CATALOG_MAINTENANCE"); } await client.query( `INSERT INTO rag_sources(source_id, source_type, source_ref) VALUES ($1, $2, $3) ON CONFLICT (source_id) DO NOTHING`, [input.sourceId, input.sourceType, input.sourceRef] ); const source = await client.query<{ active_version_id: string | null; source_type: SourceType; source_ref: string }>( "SELECT active_version_id, source_type, source_ref FROM rag_sources WHERE source_id = $1 FOR UPDATE", [input.sourceId] ); if (!source.rowCount) { throw new CatalogError("Source not found after creation", 500, "SOURCE_CREATE_FAILED"); } if (source.rows[0].source_type !== input.sourceType || source.rows[0].source_ref !== input.sourceRef) { throw new CatalogError("sourceId already exists with a different sourceType or sourceRef", 409, "SOURCE_IDENTITY_CONFLICT"); } const baseActiveVersionId = source.rows[0].active_version_id; const versionNumberResult = await client.query<{ next_version_number: string }>( "SELECT COALESCE(MAX(version_number), 0) + 1 AS next_version_number FROM rag_source_versions WHERE source_id = $1", [input.sourceId] ); const versionNumber = Number(versionNumberResult.rows[0].next_version_number); const versionResult = await client.query( `INSERT INTO rag_source_versions( source_id, version_number, previous_version_id, state, original_manifest_hash, source_content_hash, processing_fingerprint, metadata_hash, tags, activate_requested, base_active_version_id, embedding_provider, embedding_model, embedding_dimensions, qdrant_collection, expected_document_count, expected_point_count ) VALUES ($1, $2, $3, 'pending', $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16) RETURNING *`, [ input.sourceId, versionNumber, baseActiveVersionId, input.originalManifestHash, input.sourceContentHash, input.processingFingerprint, input.metadataHash, input.tags, input.activateRequested, baseActiveVersionId, input.embeddingProvider, input.embeddingModel, input.embeddingDimensions, input.qdrantCollection, input.expectedDocumentCount, input.expectedPointCount ] ); const version = toVersion(versionResult.rows[0]); for (const document of input.documents) { await client.query( `INSERT INTO rag_version_documents( version_id, document_id, document_key, original_hash, original_hash_kind, content_hash, mime_type, title, index_state, chunk_count ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, 'pending', $9)`, [ version.versionId, document.documentId, document.documentKey, document.originalHash, document.originalHashKind, document.contentHash, document.mimeType, document.title, document.chunkCount ] ); } return version; }); } async markIndexing(versionId: string): Promise { const result = await this.pool.query( `UPDATE rag_source_versions SET state = 'indexing', indexing_started_at = now() WHERE version_id = $1 AND state = 'pending'`, [versionId] ); if (result.rowCount !== 1) { throw new CatalogError("Version could not transition from pending to indexing", 409, "INVALID_VERSION_STATE"); } await this.pool.query("UPDATE rag_version_documents SET index_state = 'indexing' WHERE version_id = $1", [versionId]); } async markReady(versionId: string, verifiedPointCount: number): Promise { return withTransaction(this.pool, async (client) => { await client.query("UPDATE rag_version_documents SET index_state = 'ready' WHERE version_id = $1", [versionId]); const result = await client.query( `UPDATE rag_source_versions SET state = 'ready', verified_point_count = $2, ready_at = now() WHERE version_id = $1 AND state = 'indexing' AND expected_point_count = $2 RETURNING *`, [versionId, verifiedPointCount] ); if (result.rowCount !== 1) { throw new CatalogError("Version could not transition from indexing to ready", 409, "INVALID_VERSION_STATE"); } return toVersion(result.rows[0]); }); } async markFailed(versionId: string | undefined, code: string, detail: string): Promise { if (!versionId) { return; } const result = await this.pool.query( `UPDATE rag_source_versions SET state = 'failed', error_code = $2, error_detail = $3 WHERE version_id = $1 AND state IN ('pending', 'indexing')`, [versionId, code, detail.slice(0, 2000)] ); if (result.rowCount === 1) { await this.pool.query("UPDATE rag_version_documents SET index_state = 'failed', error_detail = $2 WHERE version_id = $1", [versionId, detail.slice(0, 2000)]); } } async activateVersion(sourceId: string, versionId: string, expectedActiveVersionId: string | null | undefined): Promise { const expected = normalizeExpectedActiveVersion(expectedActiveVersionId, true); return withTransaction(this.pool, async (client) => { if (await this.isMaintenanceEnabled(client, true)) { throw new CatalogError("Knowledge catalog is in maintenance mode", 503, "CATALOG_MAINTENANCE"); } const source = await client.query<{ active_version_id: string | null }>( "SELECT active_version_id FROM rag_sources WHERE source_id = $1 FOR UPDATE", [sourceId] ); if (!source.rowCount) { throw new CatalogError("Source not found", 404, "SOURCE_NOT_FOUND"); } const currentActive = source.rows[0].active_version_id; if (currentActive !== expected) { throw new CatalogError("Active version precondition failed", 409, "ACTIVE_VERSION_PRECONDITION_FAILED"); } const version = await client.query( "SELECT * FROM rag_source_versions WHERE source_id = $1 AND version_id = $2 FOR UPDATE", [sourceId, versionId] ); if (!version.rowCount) { throw new CatalogError("Version not found", 404, "VERSION_NOT_FOUND"); } if (!["ready", "superseded", "active"].includes(version.rows[0].state)) { throw new CatalogError("Only ready or superseded versions can be activated", 409, "VERSION_NOT_ACTIVATABLE"); } if (version.rows[0].state === "active") { if (currentActive !== versionId) { throw new CatalogError("Active version invariant is inconsistent", 503, "ACTIVE_VERSION_INVARIANT_FAILED"); } return toVersion(version.rows[0]); } if (currentActive) { await client.query( "UPDATE rag_source_versions SET state = 'superseded', superseded_at = now() WHERE source_id = $1 AND version_id = $2", [sourceId, currentActive] ); } const updated = await client.query( `UPDATE rag_source_versions SET state = 'active', activated_at = now(), error_code = NULL, error_detail = NULL WHERE source_id = $1 AND version_id = $2 RETURNING *`, [sourceId, versionId] ); if (updated.rowCount !== 1) { throw new CatalogError("Version activation failed", 500, "VERSION_ACTIVATION_FAILED"); } await client.query("UPDATE rag_sources SET active_version_id = $2, needs_reingest = false WHERE source_id = $1", [sourceId, versionId]); return toVersion(updated.rows[0]); }); } async assertActiveVersionPrecondition(sourceId: string, expectedActiveVersionId: string | null | undefined): Promise { const expected = normalizeExpectedActiveVersion(expectedActiveVersionId, true); const result = await this.pool.query<{ active_version_id: string | null }>( "SELECT active_version_id FROM rag_sources WHERE source_id = $1", [sourceId] ); if (!result.rowCount) { throw new CatalogError("Source not found", 404, "SOURCE_NOT_FOUND"); } if (result.rows[0].active_version_id !== expected) { throw new CatalogError("Active version precondition failed", 409, "ACTIVE_VERSION_PRECONDITION_FAILED"); } } async listSources(): Promise { const result = await this.pool.query<{ source_id: string; source_ref: string; tags: string[] | null; active_version_id: string | null; version_number: string | number | null; state: SourceVersionState | null; needs_reingest: boolean; updated_at: Date; }>( `SELECT s.source_id, s.source_ref, v.tags, s.active_version_id, v.version_number, v.state, s.needs_reingest, s.updated_at FROM rag_sources s LEFT JOIN rag_source_versions v ON v.version_id = s.active_version_id WHERE s.disabled_at IS NULL ORDER BY s.source_ref ASC` ); return result.rows.map((row) => ({ sourceId: row.source_id, sourceRef: row.source_ref, chunkModes: ["documental", "codigo"] as ChunkMode[], tags: row.tags ?? [], activeVersionId: row.active_version_id, activeVersionNumber: row.version_number === null ? null : Number(row.version_number), state: row.state, needsReingest: row.needs_reingest, updatedAt: row.updated_at.toISOString() })); } async getSource(sourceId: string): Promise { return (await this.listSources()).find((source) => source.sourceId === sourceId); } async listVersions(sourceId: string): Promise { const result = await this.pool.query( "SELECT * FROM rag_source_versions WHERE source_id = $1 ORDER BY version_number DESC", [sourceId] ); return result.rows.map(toVersion); } async getVersion(sourceId: string, versionId: string): Promise { const result = await this.pool.query( "SELECT * FROM rag_source_versions WHERE source_id = $1 AND version_id = $2", [sourceId, versionId] ); return result.rows[0] ? toVersion(result.rows[0]) : undefined; } async resolveActiveVersions(scope?: RetrieveScope): Promise { const mustTags = scope?.tags ?? []; const result = await this.pool.query<{ source_id: string; version_id: string; version_number: string | number; expected_point_count: string | number; embedding_dimensions: string | number; qdrant_collection: string; }>( `SELECT s.source_id, v.version_id, v.version_number, v.expected_point_count, v.embedding_dimensions, v.qdrant_collection FROM rag_sources s JOIN rag_source_versions v ON v.version_id = s.active_version_id AND v.state = 'active' WHERE s.disabled_at IS NULL AND ($1::text IS NULL OR s.source_id = $1) AND ($2::text IS NULL OR s.source_ref = $2) AND ($3::text[] IS NULL OR v.tags @> $3::text[])`, [scope?.sourceId ?? null, scope?.sourceRef ?? null, mustTags.length ? mustTags : null] ); return result.rows.map((row) => ({ sourceId: row.source_id, sourceVersionId: row.version_id, sourceVersionNumber: Number(row.version_number), expectedPointCount: Number(row.expected_point_count), embeddingDimensions: Number(row.embedding_dimensions), qdrantCollection: row.qdrant_collection })); } async markNeedsReingest(sourceId: string): Promise { await this.pool.query("UPDATE rag_sources SET needs_reingest = true WHERE source_id = $1", [sourceId]); } async rollback(sourceId: string, targetVersionId: string, expectedActiveVersionId: string): Promise { return this.withSourceLock(sourceId, async () => this.withVersionExclusiveLock(targetVersionId, async () => { const version = await this.activateVersion(sourceId, targetVersionId, expectedActiveVersionId); const source = await this.getSource(sourceId); if (source?.activeVersionId !== targetVersionId) { throw new CatalogError("Rollback verification failed", 503, "ROLLBACK_VERIFICATION_FAILED"); } return version; })); } async markPurging(sourceId: string, versionId: string): Promise { return withTransaction(this.pool, async (client) => { const version = await client.query( "SELECT * FROM rag_source_versions WHERE source_id = $1 AND version_id = $2 FOR UPDATE", [sourceId, versionId] ); if (!version.rowCount) { throw new CatalogError("Version not found", 404, "VERSION_NOT_FOUND"); } if (version.rows[0].state === "active") { throw new CatalogError("Active versions cannot be purged", 409, "ACTIVE_VERSION_PURGE_FORBIDDEN"); } if (!["ready", "superseded", "failed", "rejected", "purging"].includes(version.rows[0].state)) { throw new CatalogError("Version state cannot be purged", 409, "VERSION_NOT_PURGEABLE"); } const updated = await client.query( "UPDATE rag_source_versions SET state = 'purging' WHERE source_id = $1 AND version_id = $2 RETURNING *", [sourceId, versionId] ); return toVersion(updated.rows[0]); }); } async markPurged(versionId: string): Promise { await withTransaction(this.pool, async (client) => { await client.query("DELETE FROM rag_version_documents WHERE version_id = $1", [versionId]); const result = await client.query( `UPDATE rag_source_versions SET state = 'purged', verified_point_count = 0, artifact_state = 'retention_deleted', purged_at = now() WHERE version_id = $1 AND state = 'purging'`, [versionId] ); if (result.rowCount !== 1) { throw new CatalogError("Version could not transition from purging to purged", 409, "INVALID_VERSION_STATE"); } }); } async validateActiveInvariant(): Promise { const result = await this.pool.query<{ source_id: string }>( `SELECT s.source_id FROM rag_sources s LEFT JOIN rag_source_versions v ON v.version_id = s.active_version_id AND v.source_id = s.source_id AND v.state = 'active' WHERE (s.active_version_id IS NOT NULL AND v.version_id IS NULL) UNION SELECT v.source_id FROM rag_source_versions v JOIN rag_sources s ON s.source_id = v.source_id WHERE v.state = 'active' AND s.active_version_id IS DISTINCT FROM v.version_id` ); return result.rows.map((row) => row.source_id); } async listOrphanedIndexingCandidates(maxAgeMs = 1800000): Promise { const result = await this.pool.query<{ version_id: string; source_id: string; expected_document_count: string | number; expected_point_count: string | number; embedding_dimensions: string | number; qdrant_collection: string; }>( `SELECT version_id, source_id, expected_document_count, expected_point_count, embedding_dimensions, qdrant_collection FROM rag_source_versions WHERE ( state = 'pending' AND created_at < now() - ($1::text || ' milliseconds')::interval ) OR ( state = 'indexing' AND indexing_started_at IS NOT NULL AND indexing_started_at < now() - ($1::text || ' milliseconds')::interval )`, [String(maxAgeMs)] ); return result.rows.map((row) => ({ sourceId: row.source_id, versionId: row.version_id, expectedDocumentCount: Number(row.expected_document_count), expectedPointCount: Number(row.expected_point_count), embeddingDimensions: Number(row.embedding_dimensions), qdrantCollection: row.qdrant_collection })); } async hasCompleteVersionDocuments(versionId: string, expectedDocumentCount: number, expectedPointCount: number): Promise { const result = await this.pool.query<{ document_count: string | number; point_count: string | number; missing_content_count: string | number; }>( `SELECT COUNT(*) AS document_count, COALESCE(SUM(chunk_count), 0) AS point_count, COUNT(*) FILTER (WHERE content_hash IS NULL) AS missing_content_count FROM rag_version_documents WHERE version_id = $1`, [versionId] ); const row = result.rows[0]; return Number(row?.document_count ?? 0) === expectedDocumentCount && Number(row?.point_count ?? 0) === expectedPointCount && Number(row?.missing_content_count ?? 0) === 0; } async recoverVersionReady(versionId: string, verifiedPointCount: number): Promise { return withTransaction(this.pool, async (client) => { await client.query("UPDATE rag_version_documents SET index_state = 'ready' WHERE version_id = $1", [versionId]); const result = await client.query( `UPDATE rag_source_versions SET state = 'ready', verified_point_count = $2, ready_at = now(), error_code = NULL, error_detail = NULL WHERE version_id = $1 AND state IN ('pending', 'indexing') AND expected_point_count = $2`, [versionId, verifiedPointCount] ); return result.rowCount === 1; }); } async recoverOrphanedIndexingVersions(maxAgeMs = 1800000): Promise> { const candidates = await this.listOrphanedIndexingCandidates(maxAgeMs); const recovered: Array<{ versionId: string; sourceId: string; action: "ready" | "failed" }> = []; for (const candidate of candidates) { await this.markFailed(candidate.versionId, "ORPHANED_INDEXING", "Recovered by reconciler after stale pending/indexing state."); recovered.push({ versionId: candidate.versionId, sourceId: candidate.sourceId, action: "failed" }); } return recovered; } } export { CatalogError };