'use server'; import { cookies } from 'next/headers'; import { redirect } from 'next/navigation'; import { revalidatePath } from 'next/cache'; import { z } from 'zod'; import { getDb, connectors, connectorRuns, connectorBackfills, sales, normalizedRecords, taxonomyProposals, categories, auditLog, events, eq, and, inArray, sql } from '@rareindex/database'; import { newId } from '@rareindex/shared'; import { ADMIN_COOKIE, adminCookieValue, isAdmin, tokenMatches } from './auth'; import { enqueue, QUEUES } from './queue'; async function guard(): Promise { if (!(await isAdmin())) throw new Error('Admin authentication required'); } async function audit(entityType: string, entityId: string, action: string, reason: string, details: Record = {}) { await getDb().insert(auditLog).values({ id: newId('event'), entityType, entityId, action, reason, actor: 'admin', details }); } export async function adminLogin(formData: FormData): Promise { const token = String(formData.get('token') ?? ''); const next = String(formData.get('next') ?? '/admin'); if (!process.env.ADMIN_TOKEN) redirect('/admin/login?error=ADMIN_TOKEN+is+not+set+on+the+server'); if (!tokenMatches(token)) redirect('/admin/login?error=Invalid+token'); const jar = await cookies(); jar.set(ADMIN_COOKIE, adminCookieValue()!, { httpOnly: true, sameSite: 'strict', secure: process.env.NODE_ENV === 'production', path: '/', maxAge: 60 * 60 * 12 }); redirect(next.startsWith('/admin') ? next : '/admin'); } export async function adminLogout(): Promise { const jar = await cookies(); jar.delete(ADMIN_COOKIE); redirect('/admin/login'); } const ConnectorAction = z.enum(['run', 'probe', 'recrawl', 'backfill', 'backfill_pause', 'backfill_reset', 'retry', 'pause', 'resume', 'maintenance']); /** * Connector control-center actions (§144, SPEC §29): run/probe/backfill enqueue `crawl.run`; * pause/resume/maintenance flip status; backfill_pause/backfill_reset manage the resumable campaign. */ export async function connectorAction(formData: FormData): Promise { await guard(); const id = String(formData.get('id') ?? ''); const action = ConnectorAction.parse(formData.get('action')); const db = getDb(); const [c] = await db.select().from(connectors).where(eq(connectors.id, id)).limit(1); if (!c) throw new Error(`unknown connector ${id}`); if (action === 'pause' || action === 'resume' || action === 'maintenance') { await db.update(connectors).set({ status: action === 'pause' ? 'paused' : action === 'maintenance' ? 'maintenance' : 'active', updatedAt: new Date() }).where(eq(connectors.id, id)); await audit('connector', id, action, `admin ${action}`); } else if (action === 'backfill_pause') { await db.update(connectorBackfills).set({ status: 'paused', updatedAt: new Date() }).where(and(eq(connectorBackfills.connectorId, id), eq(connectorBackfills.status, 'running'))); await audit('connector', id, action, 'admin paused backfill campaign'); } else if (action === 'backfill_reset') { await db.update(connectorBackfills).set({ status: 'failed', lastError: 'reset by admin', finishedAt: new Date(), updatedAt: new Date() }).where(and(eq(connectorBackfills.connectorId, id), inArray(connectorBackfills.status, ['running', 'paused']))); await audit('connector', id, action, 'admin reset backfill campaign'); } else { const mode = action === 'recrawl' || action === 'backfill' ? 'backfill' : action === 'probe' ? 'probe' : 'incremental'; const payload: Record = { connectorId: id, mode, trigger: action === 'retry' ? 'retry' : 'manual', requestedAt: new Date().toISOString() }; if (action === 'probe') payload.limit = 25; if (action === 'retry') { const [last] = await db.select().from(connectorRuns).where(eq(connectorRuns.connectorId, id)).orderBy(sql`started_at desc`).limit(1); if (last) payload.cursor = last.cursor; } const jobId = await enqueue(QUEUES.crawlRun, payload, { singletonKey: `${id}:${mode}`, priority: 10 }); await db.insert(events).values({ id: newId('event'), type: 'crawl_requested', entityType: 'connector', entityId: id, payload: { ...payload, jobId } }); await audit('connector', id, action, `admin enqueued crawl.run (${mode})`, { jobId }); } revalidatePath('/admin/connectors'); revalidatePath(`/admin/connectors/${id}`); } export async function updateConnectorConfig(formData: FormData): Promise { await guard(); const id = String(formData.get('id') ?? ''); const raw = String(formData.get('config') ?? '{}'); let config: Record; try { config = z.record(z.string(), z.unknown()).parse(JSON.parse(raw)); } catch (err) { throw new Error(`config must be a JSON object: ${err instanceof Error ? err.message : String(err)}`); } const refresh = Number(formData.get('refresh') ?? NaN); const priority = String(formData.get('priority') ?? 'medium'); await getDb() .update(connectors) .set({ config, ...(Number.isFinite(refresh) && refresh > 0 ? { refreshFrequencyMinutes: Math.round(refresh) } : {}), ...(['high', 'medium', 'low'].includes(priority) ? { priority } : {}), updatedAt: new Date() }) .where(eq(connectors.id, id)); await audit('connector', id, 'edited', 'admin updated config/schedule', { config, refresh, priority }); revalidatePath(`/admin/connectors/${id}`); } export async function saleStatusAction(formData: FormData): Promise { await guard(); const id = String(formData.get('id') ?? ''); const status = z.enum(['valid', 'flagged', 'excluded']).parse(formData.get('status')); const reason = String(formData.get('reason') ?? '').trim() || `admin set status ${status}`; const db = getDb(); const [s] = await db.select({ status: sales.status, flags: sales.flags }).from(sales).where(eq(sales.id, id)).limit(1); if (!s) throw new Error('sale not found'); const flags = status === 'valid' ? s.flags.filter((f) => !f.startsWith('admin:')) : [...new Set([...s.flags, `admin:${status}`])]; await db.update(sales).set({ status, flags }).where(eq(sales.id, id)); await audit('sale', id, status === 'valid' ? 'restored' : status, reason, { from: s.status, to: status }); revalidatePath('/admin/data-quality'); } export async function manualMatchAction(formData: FormData): Promise { await guard(); const recordId = String(formData.get('recordId') ?? ''); const assetId = String(formData.get('assetId') ?? ''); const db = getDb(); if (!assetId) { await db.update(normalizedRecords).set({ status: 'rejected', rejectReason: 'admin: rejected', processedAt: new Date() }).where(eq(normalizedRecords.id, recordId)); await audit('normalized_record', recordId, 'rejected', 'admin rejected unmatched record'); } else { await db.update(normalizedRecords).set({ status: 'matched', assetId, matchMethod: 'manual', matchConfidence: 1, rejectReason: null, processedAt: null }).where(eq(normalizedRecords.id, recordId)); await audit('normalized_record', recordId, 'merged', 'admin manual match', { assetId }); await enqueue(QUEUES.normalize, { normalizedRecordId: recordId, reason: 'manual_match' }).catch(() => null); } revalidatePath('/admin/data-quality'); } export async function taxonomyDecision(formData: FormData): Promise { await guard(); const id = String(formData.get('id') ?? ''); const decision = z.enum(['approved', 'rejected']).parse(formData.get('decision')); const db = getDb(); const [p] = await db.select().from(taxonomyProposals).where(eq(taxonomyProposals.id, id)).limit(1); if (!p) throw new Error('proposal not found'); if (decision === 'approved') { const parent = String(formData.get('parent') ?? p.parentSlug ?? '') || null; const [parentRow] = parent ? await db.select().from(categories).where(eq(categories.slug, parent)).limit(1) : []; const [{ n }] = (await db.execute(sql`select count(*)::int as n from categories`)) as unknown as [{ n: number }]; await db .insert(categories) .values({ slug: p.proposedSlug, parentSlug: parentRow?.slug ?? null, familySlug: parentRow?.familySlug ?? p.proposedSlug, name: p.name, level: parentRow ? parentRow.level + 1 : 0, phase: 3, conditionScale: parentRow?.conditionScale ?? 'general', graders: parentRow?.graders ?? [], indexTicker: parentRow?.indexTicker ?? null, sortOrder: n + 1, description: `Added from taxonomy proposal ${p.id}`, }) .onConflictDoNothing(); } await db.update(taxonomyProposals).set({ status: decision, decidedBy: 'admin', decidedAt: new Date().toISOString() }).where(eq(taxonomyProposals.id, id)); await audit('taxonomy_proposal', id, decision, `admin ${decision} ${p.proposedSlug}`); revalidatePath('/admin/taxonomy'); } export async function proposeTaxonomyNode(formData: FormData): Promise { await guard(); const slug = String(formData.get('slug') ?? '').trim().toLowerCase().replace(/[^a-z0-9_]+/g, '_'); const name = String(formData.get('name') ?? '').trim(); const parent = String(formData.get('parent') ?? '').trim() || null; if (!slug || !name) throw new Error('slug and name required'); await getDb().insert(taxonomyProposals).values({ id: newId('event').replace('evt_', 'tp_'), proposedSlug: slug, name, parentSlug: parent, evidence: { source: 'admin' } }); revalidatePath('/admin/taxonomy'); }