import "server-only"; import { createHash } from "node:crypto"; import { and, eq, inArray, isNull, lt, sql } from "drizzle-orm"; import { getDb } from "@/db"; import { blocks, characters, dataSources, dataSourceVersions, media, mediaReferences, outboxEvents, pageDataSources, pages, publicSnapshots, revisions, slugAliases, } from "@/db/schema"; import { validatePageBlocks, type SerializedBlock } from "@/lib/blocks"; import { dataSourceSnapshotSchema, evaluateDataSource } from "./data-source"; import { slugSchema, type MutationResult, type SavePageInput } from "./lifecycle"; import { getAdminPage } from "./queries"; import type { DataSourceSnapshot, PageSnapshot } from "./types"; interface SavePageOptions { dependencyPins?: Record; sourceRevisionId?: string | null; } function collectNamedIds( value: unknown, keyName: "dataSourceId" | "mediaId", result = new Set(), ): Set { if (Array.isArray(value)) { for (const item of value) collectNamedIds(item, keyName, result); return result; } if (!value || typeof value !== "object") return result; for (const [key, nested] of Object.entries(value)) { if (key === keyName && typeof nested === "string") result.add(nested); collectNamedIds(nested, keyName, result); } return result; } function jsonObject(value: unknown): value is Record { return Boolean(value) && typeof value === "object" && !Array.isArray(value); } function checkpointHour(date: Date): Date { const hour = new Date(date); hour.setUTCMinutes(0, 0, 0); return hour; } function vectorHash(vector: Record): string { const stableEntries = Object.entries(vector).sort(([left], [right]) => left.localeCompare(right), ); return createHash("sha256").update(JSON.stringify(stableEntries)).digest("hex"); } function pageIdentity(page: typeof pages.$inferSelect) { return { id: page.id, characterId: page.characterId, slug: page.slug, title: page.title, navOrder: page.navOrder, visible: page.visible, publicNote: page.publicNote, version: page.version, }; } function characterIdentity(character: typeof characters.$inferSelect) { return { id: character.id, slug: character.slug, nameTh: character.nameTh, nameEn: character.nameEn, description: character.description, portraitMediaId: character.portraitMediaId, element: character.element, role: character.role, rarity: character.rarity, sortOrder: character.sortOrder, visible: character.visible, }; } export async function savePageDraft( input: SavePageInput, options: SavePageOptions = {}, ): Promise> { const issues: string[] = []; const slugResult = slugSchema.safeParse(input.slug); if (!slugResult.success) { issues.push(...slugResult.error.issues.map((issue) => issue.message)); } if (!input.title.trim()) issues.push("Page title cannot be blank."); const resolutions = validatePageBlocks(input.blocks); for (const resolution of resolutions) { if (resolution.kind === "invalid") { issues.push( ...resolution.issues.map( (issue) => `${resolution.block.type}: ${issue.message}`, ), ); } else if ( resolution.kind === "unsupported" && !jsonObject(resolution.block.config) ) { issues.push(`${resolution.block.type}: configuration must be an object.`); } } if (issues.length > 0 || !slugResult.success) { return { status: "validation-error", issues }; } const normalizedBlocks = resolutions.map((resolution) => resolution.kind === "ready" ? resolution.block : resolution.block, ); const dependencyIds = [ ...collectNamedIds( normalizedBlocks.map((block) => block.config), "dataSourceId", ), ]; const mediaIds = [ ...collectNamedIds( normalizedBlocks.map((block) => block.config), "mediaId", ), ]; const now = input.now ?? new Date(); const db = getDb(); return db.transaction(async (tx): Promise> => { const [current] = await tx .select({ page: pages, character: characters }) .from(pages) .innerJoin(characters, eq(characters.id, pages.characterId)) .where(eq(pages.id, input.pageId)) .limit(1); if (!current) return { status: "not-found" }; if (current.page.version !== input.expectedVersion) { return { status: "conflict", currentVersion: current.page.version }; } if (dependencyIds.length > 0) { const existingSources = await tx .select({ id: dataSources.id }) .from(dataSources) .where(inArray(dataSources.id, dependencyIds)); const existing = new Set(existingSources.map((source) => source.id)); const missing = dependencyIds.filter((id) => !existing.has(id)); if (missing.length > 0) { return { status: "validation-error", issues: missing.map((id) => `Referenced data source does not exist: ${id}`), }; } } const [updated] = await tx .update(pages) .set({ slug: slugResult.data, title: input.title.trim(), visible: input.visible, publicNote: input.publicNote?.trim() || null, version: input.expectedVersion + 1, updatedAt: now, }) .where( and( eq(pages.id, input.pageId), eq(pages.version, input.expectedVersion), ), ) .returning(); if (!updated) { const [latest] = await tx .select({ version: pages.version }) .from(pages) .where(eq(pages.id, input.pageId)) .limit(1); return { status: "conflict", currentVersion: latest?.version ?? input.expectedVersion, }; } if (current.page.slug !== updated.slug) { await tx .insert(slugAliases) .values({ kind: "page", oldCharacterSlug: current.character.slug, oldPageSlug: current.page.slug, targetCharacterId: current.character.id, targetPageId: current.page.id, }) .onConflictDoNothing(); } await tx.delete(blocks).where(eq(blocks.pageId, updated.id)); if (normalizedBlocks.length > 0) { await tx.insert(blocks).values( normalizedBlocks.map((block) => ({ id: block.id, pageId: updated.id, type: block.type, schemaVersion: block.schemaVersion, config: block.config as Record, responsiveLayout: jsonObject(block.responsiveLayout) ? block.responsiveLayout : {}, sortOrder: block.sortOrder, createdAt: now, updatedAt: now, })), ); } const existingDependencies = await tx .select() .from(pageDataSources) .where(eq(pageDataSources.pageId, updated.id)); const previousPins = Object.fromEntries( existingDependencies.map((dependency) => [ dependency.dataSourceId, dependency.pinnedVersion, ]), ); const requestedPins = { ...previousPins, ...options.dependencyPins, }; await tx .delete(pageDataSources) .where(eq(pageDataSources.pageId, updated.id)); if (dependencyIds.length > 0) { await tx.insert(pageDataSources).values( dependencyIds.map((dataSourceId) => ({ pageId: updated.id, dataSourceId, pinnedVersion: requestedPins[dataSourceId] ?? null, })), ); } const sourceRows = dependencyIds.length === 0 ? [] : await tx .select() .from(dataSources) .where(inArray(dataSources.id, dependencyIds)); const versionVector = Object.fromEntries( sourceRows.map((source) => [ source.id, requestedPins[source.id] ?? source.currentVersion, ]), ); const sourceSnapshots: Record = {}; for (const source of sourceRows) { const version = versionVector[source.id]; const [versionRow] = await tx .select() .from(dataSourceVersions) .where( and( eq(dataSourceVersions.dataSourceId, source.id), eq(dataSourceVersions.version, version), ), ) .limit(1); if (versionRow) { sourceSnapshots[source.id] = { id: source.id, name: source.name, version, fields: source.fields as DataSourceSnapshot["fields"], rows: versionRow.rows as DataSourceSnapshot["rows"], }; } } const navigationRows = await tx .select() .from(pages) .where(eq(pages.characterId, updated.characterId)) .orderBy(pages.navOrder, pages.title); const snapshot: PageSnapshot = { character: characterIdentity(current.character), page: pageIdentity(updated), pages: navigationRows.map(pageIdentity), blocks: normalizedBlocks, dataSources: sourceSnapshots, versionVector, }; const revisionSnapshot = { page: snapshot.page, blocks: snapshot.blocks, dependencies: Object.fromEntries( dependencyIds.map((id) => [id, requestedPins[id] ?? null]), ), }; const hour = checkpointHour(now); const [existingCheckpoint] = await tx .select({ id: revisions.id }) .from(revisions) .where( and( eq(revisions.pageId, updated.id), eq(revisions.kind, "checkpoint"), eq(revisions.checkpointHour, hour), ), ) .limit(1); const revisionValues = { pageVersion: updated.version, kind: "checkpoint", label: null, authorId: input.authorId, sourceRevisionId: options.sourceRevisionId ?? null, checkpointHour: hour, snapshot: revisionSnapshot as unknown as Record, dataSourceVersions: versionVector, mediaIds, createdAt: now, expiresAt: new Date(now.getTime() + 30 * 24 * 60 * 60 * 1_000), }; const revisionId = existingCheckpoint ? existingCheckpoint.id : crypto.randomUUID(); if (existingCheckpoint) { await tx .update(revisions) .set(revisionValues) .where(eq(revisions.id, revisionId)); } else { await tx.insert(revisions).values({ id: revisionId, pageId: updated.id, ...revisionValues, }); } const priorCurrentMedia = await tx .select({ mediaId: mediaReferences.mediaId }) .from(mediaReferences) .where( and( eq(mediaReferences.ownerType, "page"), eq(mediaReferences.ownerId, updated.id), eq(mediaReferences.referenceKind, "current"), ), ); const priorRevisionMedia = await tx .select({ mediaId: mediaReferences.mediaId }) .from(mediaReferences) .where( and( eq(mediaReferences.ownerType, "revision"), eq(mediaReferences.ownerId, revisionId), eq(mediaReferences.referenceKind, "revision"), ), ); await tx .delete(mediaReferences) .where( and( eq(mediaReferences.ownerType, "page"), eq(mediaReferences.ownerId, updated.id), eq(mediaReferences.referenceKind, "current"), ), ); await tx .delete(mediaReferences) .where( and( eq(mediaReferences.ownerType, "revision"), eq(mediaReferences.ownerId, revisionId), eq(mediaReferences.referenceKind, "revision"), ), ); if (mediaIds.length > 0) { await tx.insert(mediaReferences).values([ ...mediaIds.map((mediaId) => ({ mediaId, ownerType: "page", ownerId: updated.id, referenceKind: "current", createdAt: now, })), ...mediaIds.map((mediaId) => ({ mediaId, ownerType: "revision", ownerId: revisionId, referenceKind: "revision", createdAt: now, })), ]); } const affectedMediaIds = [ ...new Set([ ...mediaIds, ...priorCurrentMedia.map((item) => item.mediaId), ...priorRevisionMedia.map((item) => item.mediaId), ]), ]; for (const mediaId of affectedMediaIds) { await tx .update(media) .set({ currentReferenceCount: sql`( select count(*)::int from ${mediaReferences} where ${mediaReferences.mediaId} = ${mediaId} and ${mediaReferences.referenceKind} = 'current' )`, revisionReferenceCount: sql`( select count(*)::int from ${mediaReferences} where ${mediaReferences.mediaId} = ${mediaId} and ${mediaReferences.referenceKind} = 'revision' )`, }) .where(eq(media.id, mediaId)); } const publicDto: PageSnapshot = { ...snapshot, pages: snapshot.pages.filter((item) => item.visible), }; await tx.insert(publicSnapshots).values({ pageId: updated.id, pageVersion: updated.version, versionVector, vectorHash: vectorHash(versionVector), dto: publicDto as unknown as Record, visible: current.character.visible && updated.visible, createdAt: now, }); await tx.insert(outboxEvents).values([ { topic: `page:${updated.id}`, aggregateId: updated.id, eventType: "page.updated", payload: { id: updated.id, version: updated.version }, createdAt: now, }, { topic: "directory", aggregateId: updated.id, eventType: "directory.updated", payload: { id: updated.id, version: updated.version }, createdAt: now, }, { topic: "admin", aggregateId: updated.id, eventType: "admin.updated", payload: { id: updated.id, version: updated.version }, createdAt: now, }, ]); return { status: "accepted", value: snapshot }; }); } export async function saveDataSourceVersion( input: Omit & { expectedVersion: number; authorId: string; now?: Date; }, ): Promise> { const next: DataSourceSnapshot = { id: input.id, name: input.name, version: input.expectedVersion + 1, fields: input.fields, rows: input.rows, }; const parsed = dataSourceSnapshotSchema.safeParse(next); if (!parsed.success) { return { status: "validation-error", issues: parsed.error.issues.map((issue) => issue.message), }; } const evaluation = evaluateDataSource(parsed.data); const now = input.now ?? new Date(); const db = getDb(); return db.transaction(async (tx): Promise> => { const [current] = await tx .select() .from(dataSources) .where(eq(dataSources.id, input.id)) .limit(1); if (!current) return { status: "not-found" }; if (current.currentVersion !== input.expectedVersion) { return { status: "conflict", currentVersion: current.currentVersion }; } const [updated] = await tx .update(dataSources) .set({ name: parsed.data.name, fields: parsed.data.fields, currentVersion: parsed.data.version, updatedAt: now, }) .where( and( eq(dataSources.id, input.id), eq(dataSources.currentVersion, input.expectedVersion), ), ) .returning(); if (!updated) { return { status: "conflict", currentVersion: current.currentVersion }; } await tx.insert(dataSourceVersions).values({ dataSourceId: input.id, version: parsed.data.version, rows: parsed.data.rows, evaluation: evaluation as unknown as Record, createdById: input.authorId, createdAt: now, }); const consumers = await tx .select({ pageId: pageDataSources.pageId, pageVersion: pages.version }) .from(pageDataSources) .innerJoin(pages, eq(pages.id, pageDataSources.pageId)) .where( and( eq(pageDataSources.dataSourceId, input.id), isNull(pageDataSources.pinnedVersion), ), ); await tx.insert(outboxEvents).values([ { topic: `data-source:${input.id}`, aggregateId: input.id, eventType: "data-source.updated", payload: { id: input.id, version: parsed.data.version }, createdAt: now, }, ...consumers.map((consumer) => ({ topic: `page:${consumer.pageId}`, aggregateId: consumer.pageId, eventType: "page.updated", payload: { id: consumer.pageId, version: consumer.pageVersion }, createdAt: now, })), { topic: "admin", aggregateId: input.id, eventType: "admin.updated", payload: { id: input.id, version: parsed.data.version }, createdAt: now, }, ]); return { status: "accepted", value: parsed.data }; }); } export async function createNamedRevision( pageId: string, label: string, authorId: string, now = new Date(), ): Promise { const snapshot = await getAdminPageById(pageId); if (!snapshot) return null; const mediaIds = [ ...collectNamedIds( snapshot.blocks.map((block) => block.config), "mediaId", ), ]; const id = crypto.randomUUID(); const dependencies = Object.fromEntries( Object.keys(snapshot.versionVector).map((sourceId) => [sourceId, null]), ); const db = getDb(); await db.transaction(async (tx) => { await tx.insert(revisions).values({ id, pageId, pageVersion: snapshot.page.version, kind: "named", label: label.trim() || "รุ่นที่บันทึก", authorId, snapshot: { page: snapshot.page, blocks: snapshot.blocks, dependencies, } as unknown as Record, dataSourceVersions: snapshot.versionVector, mediaIds, createdAt: now, expiresAt: null, }); if (mediaIds.length > 0) { await tx.insert(mediaReferences).values( mediaIds.map((mediaId) => ({ mediaId, ownerType: "revision", ownerId: id, referenceKind: "revision", createdAt: now, })), ); } }); return id; } export async function restoreRevision( revisionId: string, expectedVersion: number, authorId: string, ): Promise> { const db = getDb(); const [revision] = await db .select() .from(revisions) .where(eq(revisions.id, revisionId)) .limit(1); if (!revision) return { status: "not-found" }; const snapshot = revision.snapshot as { page: PageSnapshot["page"]; blocks: SerializedBlock[]; dependencies: Record; }; return savePageDraft( { pageId: revision.pageId, expectedVersion, authorId, title: snapshot.page.title, slug: snapshot.page.slug, visible: snapshot.page.visible, publicNote: snapshot.page.publicNote, blocks: snapshot.blocks, }, { dependencyPins: revision.dataSourceVersions, sourceRevisionId: revision.id, }, ); } export async function useLatestDataSource( pageId: string, dataSourceId: string, expectedVersion: number, authorId: string, ): Promise> { const snapshot = await getAdminPageById(pageId); if (!snapshot) return { status: "not-found" }; if (snapshot.page.version !== expectedVersion) { return { status: "conflict", currentVersion: snapshot.page.version }; } if (!Object.hasOwn(snapshot.versionVector, dataSourceId)) { return { status: "validation-error", issues: ["The page does not reference this data source."], }; } return savePageDraft( { pageId, expectedVersion, authorId, title: snapshot.page.title, slug: snapshot.page.slug, visible: snapshot.page.visible, publicNote: snapshot.page.publicNote, blocks: snapshot.blocks as SerializedBlock[], }, { dependencyPins: { [dataSourceId]: null } }, ); } export async function expireOldCheckpoints(now = new Date()): Promise { const db = getDb(); return db.transaction(async (tx) => { const expired = await tx .select({ id: revisions.id }) .from(revisions) .where( and( eq(revisions.kind, "checkpoint"), lt(revisions.expiresAt, now), ), ); if (expired.length === 0) return []; const ids = expired.map((revision) => revision.id); const affected = await tx .select({ mediaId: mediaReferences.mediaId }) .from(mediaReferences) .where( and( eq(mediaReferences.ownerType, "revision"), inArray(mediaReferences.ownerId, ids), ), ); await tx .delete(mediaReferences) .where( and( eq(mediaReferences.ownerType, "revision"), inArray(mediaReferences.ownerId, ids), ), ); await tx.delete(revisions).where(inArray(revisions.id, ids)); for (const mediaId of new Set(affected.map((item) => item.mediaId))) { await tx .update(media) .set({ revisionReferenceCount: sql`( select count(*)::int from ${mediaReferences} where ${mediaReferences.mediaId} = ${mediaId} and ${mediaReferences.referenceKind} = 'revision' )`, }) .where(eq(media.id, mediaId)); } return ids; }); } async function getAdminPageById(pageId: string): Promise { const db = getDb(); const [identity] = await db .select({ characterSlug: characters.slug, pageSlug: pages.slug }) .from(pages) .innerJoin(characters, eq(characters.id, pages.characterId)) .where(eq(pages.id, pageId)) .limit(1); return identity ? getAdminPage(identity.characterSlug, identity.pageSlug) : null; }