import { sql, type Kysely } from "kysely"; import { ulid } from "ulidx"; import { MediaUsageWorkRepository } from "../../database/repositories/media-usage-work.js"; import { MediaUsageRepository, type MediaUsageSource, } from "../../database/repositories/media-usage.js"; import type { Database } from "../../database/types.js"; import { validateIdentifier } from "../../database/validate.js"; import { chunks, SQL_BATCH_SIZE } from "../../utils/chunks.js"; import { loadContentMediaUsageFields, MediaUsageFieldDiscoveryError, type ContentMediaUsageFieldDiscovery, } from "./content-fields.js"; import { CONTENT_MEDIA_USAGE_ADAPTER_ID, CONTENT_MEDIA_USAGE_COLLECTION_SCOPE, withContentUsageCollectionLock, } from "./content-refresh.js"; import { CONTENT_SOURCE_SCHEMA_VERSION, type ContentMediaUsageSnapshot, loadContentMediaUsageSnapshots, } from "./content-snapshots.js"; import { buildContentMediaUsageSourceKey, MEDIA_USAGE_CONTENT_SOURCE_VARIANTS, } from "./source-key.js"; export type ContentMediaUsageRepairStatus = "complete" | "partial" | "failed" | "stale"; export const CONTENT_MEDIA_USAGE_REPAIR_ERROR = { COLLECTION_NOT_FOUND: "COLLECTION_NOT_FOUND", CONTENT_NOT_FOUND: "CONTENT_NOT_FOUND", DRAFT_REVISION_NOT_FOUND: "DRAFT_REVISION_NOT_FOUND", DRAFT_REVISION_MISMATCH: "DRAFT_REVISION_MISMATCH", DRAFT_REVISION_INVALID: "DRAFT_REVISION_INVALID", CONTENT_USAGE_REPAIR_ERROR: "CONTENT_USAGE_REPAIR_ERROR", CONTENT_USAGE_REPAIR_CONFLICT: "CONTENT_USAGE_REPAIR_CONFLICT", INVALID_REPEATER_VALIDATION: "INVALID_REPEATER_VALIDATION", } as const; export type ContentMediaUsageRepairErrorCode = (typeof CONTENT_MEDIA_USAGE_REPAIR_ERROR)[keyof typeof CONTENT_MEDIA_USAGE_REPAIR_ERROR]; export interface ContentMediaUsageRepairCollectionInput { collectionSlug: string; } export interface ContentMediaUsageRepairScope { adapterId: typeof CONTENT_MEDIA_USAGE_ADAPTER_ID; scopeType: typeof CONTENT_MEDIA_USAGE_COLLECTION_SCOPE; scopeKey: string; } export interface ContentMediaUsageRepairCollectionResult { scope: ContentMediaUsageRepairScope; status: ContentMediaUsageRepairStatus; indexedSourceCount: number; failedSourceCount: number; skippedSourceCount: number; deletedSourceCount: number; lastErrorCode: string | null; startedAt: string; completedAt: string | null; } export interface ContentMediaUsageRepairAllResult { status: ContentMediaUsageRepairStatus; collections: ContentMediaUsageRepairCollectionResult[]; indexedSourceCount: number; failedSourceCount: number; skippedSourceCount: number; deletedSourceCount: number; } export interface ContentMediaUsageCollectionScan { collectionSlug: string; contentIds: string[]; } interface ContentMediaUsageCollectionRecord { id: string; slug: string; } interface ContentMediaUsageInitialCollectionResult { collection: ContentMediaUsageCollectionRecord; result: ContentMediaUsageRepairCollectionResult; } interface ContentMediaUsageRepairExecution { collectionId: string; startingEpoch: number | string; } export async function repairContentMediaUsageAll( db: Kysely, ): Promise { const collections = await loadContentMediaUsageCollectionRecords(db); const results: ContentMediaUsageInitialCollectionResult[] = []; for (const collection of collections) { results.push({ collection, result: await repairContentMediaUsageCollectionSafely(db, collection), }); } return aggregateContentMediaUsageRepairAll( await filterExistingContentMediaUsageCollectionResults(db, results), ); } export async function scanContentMediaUsageCollection( db: Kysely, collectionSlug: string, expectedCollectionId?: string, ): Promise { validateIdentifier(collectionSlug, "collection slug"); let collectionQuery = db .selectFrom("_emdash_collections") .select("id") .where("slug", "=", collectionSlug); if (expectedCollectionId !== undefined) { collectionQuery = collectionQuery.where("id", "=", expectedCollectionId); } const collection = await collectionQuery.executeTakeFirst(); if (!collection) return null; const tableName = getContentTableName(collectionSlug); const rows = await sql<{ id: string }>` SELECT content.id FROM ${sql.ref(tableName)} AS content WHERE EXISTS ( SELECT 1 FROM _emdash_collections AS collection WHERE collection.id = ${collection.id} AND collection.slug = ${collectionSlug} ) ORDER BY content.id ASC `.execute(db); return { collectionSlug, contentIds: rows.rows.map((row) => row.id), }; } export async function repairContentMediaUsageCollection( db: Kysely, input: ContentMediaUsageRepairCollectionInput, ): Promise { validateIdentifier(input.collectionSlug, "collection slug"); return withContentUsageCollectionLock(input.collectionSlug, () => repairContentMediaUsageCollectionUnlocked(db, input.collectionSlug), ); } async function loadContentMediaUsageCollectionRecords( db: Kysely, ): Promise { return db .selectFrom("_emdash_collections") .select(["id", "slug"]) .orderBy("slug", "asc") .execute(); } async function repairContentMediaUsageCollectionSafely( db: Kysely, collection: ContentMediaUsageCollectionRecord, ): Promise { try { return await withContentUsageCollectionLock(collection.slug, () => repairContentMediaUsageCollectionUnlocked(db, collection.slug, collection.id), ); } catch (error) { console.error(`[media-usage] Failed to repair collection ${collection.slug}:`, error); const now = new Date().toISOString(); return { scope: contentMediaUsageCollectionScope(collection.slug), status: "failed", indexedSourceCount: 0, failedSourceCount: 0, skippedSourceCount: 0, deletedSourceCount: 0, lastErrorCode: CONTENT_MEDIA_USAGE_REPAIR_ERROR.CONTENT_USAGE_REPAIR_ERROR, startedAt: now, completedAt: now, }; } } async function filterExistingContentMediaUsageCollectionResults( db: Kysely, results: readonly ContentMediaUsageInitialCollectionResult[], ): Promise { const identityBound = await isIncrementalCaptureActive(db); const currentCollections = await loadContentMediaUsageCollectionRecordsSafely(db); const currentIdsBySlug = new Map( currentCollections.map((collection) => [collection.slug, collection.id]), ); const includedResults: ContentMediaUsageRepairCollectionResult[] = []; const excludedResults: ContentMediaUsageInitialCollectionResult[] = []; for (const { collection, result } of results) { if (currentIdsBySlug.get(collection.slug) === collection.id) { includedResults.push(result); } else { excludedResults.push({ collection, result }); } } if (excludedResults.length > 0) { const repo = new MediaUsageRepository(db); for (const excluded of excludedResults) { await repo.deleteIndexStatus( excluded.result.scope, identityBound ? excluded.collection.id : undefined, ); } } return includedResults; } async function loadContentMediaUsageCollectionRecordsSafely( db: Kysely, ): Promise { try { return await loadContentMediaUsageCollectionRecords(db); } catch { // Retry once before failing; returning unpruned results can over-report deleted collections. } try { return await loadContentMediaUsageCollectionRecords(db); } catch (error) { console.error("[media-usage] Failed to reconcile all-content repair collections:", error); throw error; } } function aggregateContentMediaUsageRepairAll( collections: readonly ContentMediaUsageRepairCollectionResult[], ): ContentMediaUsageRepairAllResult { return { status: determineRepairAllStatus(collections), collections: [...collections], indexedSourceCount: sumCollectionRepairCount(collections, "indexedSourceCount"), failedSourceCount: sumCollectionRepairCount(collections, "failedSourceCount"), skippedSourceCount: sumCollectionRepairCount(collections, "skippedSourceCount"), deletedSourceCount: sumCollectionRepairCount(collections, "deletedSourceCount"), }; } function determineRepairAllStatus( collections: readonly ContentMediaUsageRepairCollectionResult[], ): ContentMediaUsageRepairStatus { if (collections.length === 0) return "complete"; if (collections.every((collection) => collection.status === "complete")) return "complete"; if (collections.some((collection) => collection.status === "stale")) return "stale"; if (collections.some((collection) => collection.status === "partial")) return "partial"; if (collections.every((collection) => collection.status === "failed")) return "failed"; return "partial"; } function sumCollectionRepairCount( collections: readonly ContentMediaUsageRepairCollectionResult[], key: "indexedSourceCount" | "failedSourceCount" | "skippedSourceCount" | "deletedSourceCount", ): number { return collections.reduce((sum, collection) => sum + collection[key], 0); } async function repairContentMediaUsageCollectionUnlocked( db: Kysely, collectionSlug: string, expectedCollectionId?: string, ): Promise { const requestedAt = new Date().toISOString(); const scope = contentMediaUsageCollectionScope(collectionSlug); const identityBound = await isIncrementalCaptureActive(db); const collection = await loadContentMediaUsageCollectionRecord(db, collectionSlug); if ( !collection || (identityBound && expectedCollectionId && collection.id !== expectedCollectionId) ) { return { scope, status: "failed", indexedSourceCount: 0, failedSourceCount: 0, skippedSourceCount: 0, deletedSourceCount: 0, lastErrorCode: CONTENT_MEDIA_USAGE_REPAIR_ERROR.COLLECTION_NOT_FOUND, startedAt: requestedAt, completedAt: requestedAt, }; } const repo = new MediaUsageRepository(db); const runToken = ulid(); let startedAt = requestedAt; let execution: ContentMediaUsageRepairExecution | undefined; if (identityBound) { const run = await repo.beginIndexStatusRepairAtCurrentEpoch({ ...scope, collectionId: collection.id, runToken, schemaVersion: CONTENT_SOURCE_SCHEMA_VERSION, }); if (!run) { return { scope, status: "failed", indexedSourceCount: 0, failedSourceCount: 0, skippedSourceCount: 0, deletedSourceCount: 0, lastErrorCode: CONTENT_MEDIA_USAGE_REPAIR_ERROR.CONTENT_USAGE_REPAIR_ERROR, startedAt, completedAt: requestedAt, }; } startedAt = run.startedAt; execution = { collectionId: collection.id, startingEpoch: run.changeEpoch }; } else { await repo.beginIndexStatusRepair({ ...scope, runToken, schemaVersion: CONTENT_SOURCE_SCHEMA_VERSION, startedAt, }); } try { const scan = await scanContentMediaUsageCollection(db, collectionSlug, execution?.collectionId); if (!scan) { const completedAt = new Date().toISOString(); return await finalizeRepairStatus(db, repo, { ...scope, runToken, counts: { indexedSourceCount: 0, failedSourceCount: 0, skippedSourceCount: 0, deletedSourceCount: 0, lastErrorCode: CONTENT_MEDIA_USAGE_REPAIR_ERROR.COLLECTION_NOT_FOUND, missingContentIds: new Set(), }, status: "failed", startedAt, completedAt, execution, }); } const counts = await repairScannedContentSources(db, repo, scan, execution?.collectionId); const finalScan = await scanContentMediaUsageCollection( db, collectionSlug, execution?.collectionId, ); if (!finalScan) { counts.failedSourceCount++; counts.lastErrorCode = CONTENT_MEDIA_USAGE_REPAIR_ERROR.COLLECTION_NOT_FOUND; } else if (!sameContentIds(repairedContentIds(scan.contentIds, counts), finalScan.contentIds)) { markRepairConflict(counts); } const completedAt = new Date().toISOString(); const status = determineRepairStatus(counts); return await finalizeRepairStatus(db, repo, { ...scope, runToken, counts, status, startedAt, completedAt, execution, }); } catch (error) { if (!(error instanceof MediaUsageFieldDiscoveryError)) { console.error(`[media-usage] Failed to repair collection ${collectionSlug}:`, error); } const completedAt = new Date().toISOString(); const lastErrorCode = error instanceof MediaUsageFieldDiscoveryError ? error.code : CONTENT_MEDIA_USAGE_REPAIR_ERROR.CONTENT_USAGE_REPAIR_ERROR; return finalizeRepairStatus(db, repo, { ...scope, runToken, counts: { indexedSourceCount: 0, failedSourceCount: 0, skippedSourceCount: 0, deletedSourceCount: 0, lastErrorCode, missingContentIds: new Set(), }, status: "failed", startedAt, completedAt, execution, }); } } interface RepairCounts { indexedSourceCount: number; failedSourceCount: number; skippedSourceCount: number; deletedSourceCount: number; lastErrorCode: ContentMediaUsageRepairErrorCode | null; missingContentIds: Set; } async function repairScannedContentSources( db: Kysely, repo: MediaUsageRepository, scan: ContentMediaUsageCollectionScan, collectionId?: string, ): Promise { const counts: RepairCounts = { indexedSourceCount: 0, failedSourceCount: 0, skippedSourceCount: 0, deletedSourceCount: 0, lastErrorCode: null, missingContentIds: new Set(), }; const fieldDiscovery = await loadContentMediaUsageFields(db, scan.collectionSlug, collectionId); const observedSources = await repo.findSources(buildContentSourceKeysForScan(scan, collectionId)); for (const contentId of scan.contentIds) { await repairContentSource( db, repo, scan.collectionSlug, contentId, fieldDiscovery, observedSources, counts, collectionId, ); } await reconcileOrphanedContentSources(db, repo, scan.collectionSlug, counts, collectionId); return counts; } async function repairContentSource( db: Kysely, repo: MediaUsageRepository, collectionSlug: string, contentId: string, fieldDiscovery: ContentMediaUsageFieldDiscovery, observedSources: Map, counts: RepairCounts, collectionId?: string, ): Promise { const sourceKeys = buildContentSourceKeys(collectionSlug, contentId, collectionId); const snapshotsResult = await loadContentMediaUsageSnapshots( db, collectionSlug, contentId, fieldDiscovery, collectionId ? { collectionId, identityVersion: 1 } : undefined, ); if (!snapshotsResult.success) { if (snapshotsResult.error === CONTENT_MEDIA_USAGE_REPAIR_ERROR.CONTENT_NOT_FOUND) { markRepairConflict(counts); counts.missingContentIds.add(contentId); return; } counts.lastErrorCode = snapshotsResult.error; if (snapshotsResult.snapshots) { // Partial snapshot failures keep absent variants in place; the attempted source is marked below. await repairSnapshotSources(repo, snapshotsResult.snapshots, observedSources, counts); } if (snapshotsResult.source) { const result = await repo.markSourceAttemptedIfMatching( { ...snapshotsResult.source, sourceCompleteness: "failed", lastErrorCode: snapshotsResult.error, }, observedSources.get(snapshotsResult.source.sourceKey) ?? null, ); if (result.attempted) { counts.failedSourceCount++; } else { markRepairConflict(counts); } return; } counts.failedSourceCount++; return; } const expectedSourceKeys = await repairSnapshotSources( repo, snapshotsResult.snapshots, observedSources, counts, ); for (const sourceKey of sourceKeys) { if (expectedSourceKeys.has(sourceKey)) continue; const observedSource = observedSources.get(sourceKey); if (!observedSource) continue; await deleteObservedSource(repo, sourceKey, observedSource, counts); } } async function repairSnapshotSources( repo: MediaUsageRepository, snapshots: readonly ContentMediaUsageSnapshot[], observedSources: Map, counts: RepairCounts, ): Promise> { const expectedSourceKeys = new Set(); for (const snapshot of snapshots) { expectedSourceKeys.add(snapshot.source.sourceKey); const result = await repo.replaceSourceIfMatching( snapshot.source, snapshot.occurrences, observedSources.get(snapshot.source.sourceKey) ?? null, ); if (result.replaced || result.unchanged) { counts.indexedSourceCount++; } else { markRepairConflict(counts); } } return expectedSourceKeys; } function markRepairConflict(counts: RepairCounts): void { counts.skippedSourceCount++; counts.lastErrorCode ??= CONTENT_MEDIA_USAGE_REPAIR_ERROR.CONTENT_USAGE_REPAIR_CONFLICT; } async function reconcileOrphanedContentSources( db: Kysely, repo: MediaUsageRepository, collectionSlug: string, counts: RepairCounts, collectionId?: string, ): Promise { const sources = await repo.findCollectionContentSources(collectionSlug, collectionId); const existingContentIds = await findExistingContentIds( db, collectionSlug, sources.flatMap((source) => (source.contentId ? [source.contentId] : [])), collectionId, ); for (const source of sources) { if (!source.contentId) { await deleteObservedSource(repo, source.sourceKey, source, counts); continue; } if (existingContentIds.has(source.contentId)) continue; await deleteObservedSourceIfContentAbsent(repo, collectionSlug, source, counts); } } async function findExistingContentIds( db: Kysely, collectionSlug: string, contentIds: readonly string[], collectionId?: string, ): Promise> { validateIdentifier(collectionSlug, "collection slug"); const existingContentIds = new Set(); const uniqueContentIds = [...new Set(contentIds)]; if (uniqueContentIds.length === 0) return existingContentIds; const tableName = getContentTableName(collectionSlug); for (const contentIdBatch of chunks(uniqueContentIds, SQL_BATCH_SIZE)) { const result = await sql<{ id: string }>` SELECT content.id FROM ${sql.ref(tableName)} AS content WHERE content.id IN (${sql.join(contentIdBatch)}) ${ collectionId ? sql`AND EXISTS ( SELECT 1 FROM _emdash_collections AS collection WHERE collection.id = ${collectionId} AND collection.slug = ${collectionSlug} )` : sql`` } `.execute(db); for (const row of result.rows) { existingContentIds.add(row.id); } } return existingContentIds; } async function deleteObservedSourceIfContentAbsent( repo: MediaUsageRepository, collectionSlug: string, observedSource: MediaUsageSource, counts: RepairCounts, ): Promise { if (!observedSource.contentId) return; const result = await repo.deleteSourceIfMatchingContentAbsent( observedSource.sourceKey, observedSource, collectionSlug, observedSource.contentId, ); if (result.deleted) { counts.deletedSourceCount++; return; } if (result.contentPresent) return; if (result.source) { markRepairConflict(counts); } } async function deleteObservedSource( repo: MediaUsageRepository, sourceKey: string, observedSource: MediaUsageSource, counts: RepairCounts, ): Promise { const result = await repo.deleteSourceIfMatching(sourceKey, observedSource); if (result.deleted) { counts.deletedSourceCount++; return; } if (result.source) { markRepairConflict(counts); } } interface FinalizeInput extends ContentMediaUsageRepairScope { runToken: string; counts: RepairCounts; status: Exclude; startedAt: string; completedAt: string; execution?: ContentMediaUsageRepairExecution; } async function finalizeRepairStatus( db: Kysely, repo: MediaUsageRepository, input: FinalizeInput, ): Promise { let result; if (input.execution) { if (input.status === "complete") { await new MediaUsageWorkRepository(db).deleteWorkThroughEpoch( input.execution.collectionId, input.execution.startingEpoch, ); } result = await repo.finalizeIndexStatusRepairAtEpoch({ adapterId: input.adapterId, scopeType: input.scopeType, scopeKey: input.scopeKey, collectionId: input.execution.collectionId, runToken: input.runToken, startingEpoch: input.execution.startingEpoch, status: input.status, schemaVersion: CONTENT_SOURCE_SCHEMA_VERSION, indexedSourceCount: input.counts.indexedSourceCount, failedSourceCount: input.counts.failedSourceCount, lastErrorCode: input.counts.lastErrorCode, }); } else { result = await repo.finalizeIndexStatusRepairIfRunning({ adapterId: input.adapterId, scopeType: input.scopeType, scopeKey: input.scopeKey, runToken: input.runToken, status: input.status, schemaVersion: CONTENT_SOURCE_SCHEMA_VERSION, completedAt: input.completedAt, indexedSourceCount: input.counts.indexedSourceCount, failedSourceCount: input.counts.failedSourceCount, lastErrorCode: input.counts.lastErrorCode, }); } return { scope: { adapterId: input.adapterId, scopeType: input.scopeType, scopeKey: input.scopeKey, }, status: result.finalized ? input.status : "stale", indexedSourceCount: input.counts.indexedSourceCount, failedSourceCount: input.counts.failedSourceCount, skippedSourceCount: input.counts.skippedSourceCount, deletedSourceCount: input.counts.deletedSourceCount, lastErrorCode: result.finalized ? input.counts.lastErrorCode : (result.status?.lastErrorCode ?? CONTENT_MEDIA_USAGE_REPAIR_ERROR.CONTENT_USAGE_REPAIR_CONFLICT), startedAt: input.startedAt, completedAt: result.finalized ? (result.status?.completedAt ?? input.completedAt) : null, }; } function determineRepairStatus( counts: RepairCounts, ): Exclude { if (counts.failedSourceCount === 0 && counts.skippedSourceCount === 0) return "complete"; const trustedProgress = counts.indexedSourceCount + counts.deletedSourceCount; if (counts.failedSourceCount > 0 && trustedProgress === 0) { return "failed"; } return "partial"; } function buildContentSourceKeys( collectionSlug: string, contentId: string, collectionId?: string, ): string[] { return MEDIA_USAGE_CONTENT_SOURCE_VARIANTS.map((sourceVariant) => buildContentMediaUsageSourceKey({ collectionId, collectionSlug, contentId, sourceVariant }), ); } function buildContentSourceKeysForScan( scan: ContentMediaUsageCollectionScan, collectionId?: string, ): string[] { return scan.contentIds.flatMap((contentId) => buildContentSourceKeys(scan.collectionSlug, contentId, collectionId), ); } async function loadContentMediaUsageCollectionRecord( db: Kysely, collectionSlug: string, ): Promise { const row = await db .selectFrom("_emdash_collections") .select(["id", "slug"]) .where("slug", "=", collectionSlug) .executeTakeFirst(); return row ?? null; } async function isIncrementalCaptureActive(db: Kysely): Promise { const row = await db .selectFrom("_emdash_media_usage_activation") .select("state") .where("task_key", "=", "incremental_capture") .executeTakeFirst(); return row?.state === "active"; } function sameContentIds(left: readonly string[], right: readonly string[]): boolean { if (left.length !== right.length) return false; const rightIds = new Set(right); return left.every((id) => rightIds.has(id)); } function repairedContentIds( contentIds: readonly string[], counts: Pick, ): string[] { if (counts.missingContentIds.size === 0) return [...contentIds]; return contentIds.filter((contentId) => !counts.missingContentIds.has(contentId)); } function contentMediaUsageCollectionScope(collectionSlug: string): ContentMediaUsageRepairScope { return { adapterId: CONTENT_MEDIA_USAGE_ADAPTER_ID, scopeType: CONTENT_MEDIA_USAGE_COLLECTION_SCOPE, scopeKey: collectionSlug, }; } function getContentTableName(collectionSlug: string): string { validateIdentifier(collectionSlug, "collection slug"); return `ec_${collectionSlug}`; }