import { after, type NextRequest, NextResponse } from "next/server"; import { logAuthorizationEvent } from "@/lib/admin/authorization-events"; import { validateCsrfToken } from "@/lib/foundation/security"; import { canAccess, getApiAdminContext } from "@/lib/permissions"; import { logServerError } from "@/lib/server-log"; import { beginCatalogExport, catalogExportEnabled, isCatalogMutation, } from "@/lib/services/catalog-git-queue"; const MUTATING_METHODS = new Set(["POST", "PUT", "PATCH", "DELETE"]); const MAX_BODY_BYTES = 10 * 1024 * 1024; // 10 MB type AdminContext = NonNullable>>; type RouteContext = { params?: Promise> }; type AdminHandler = ( request: NextRequest, context: AdminContext, routeContext: RouteContext, ) => Promise | Response; export function withAdmin( options: { permission?: string; requireCsrf?: boolean; maxBodyBytes?: number; }, handler: AdminHandler, ) { return async (request: NextRequest, routeContext: RouteContext = {}) => { // CSRF required for mutating admin APIs unless explicitly opted out. const csrfRequired = options.requireCsrf !== false && MUTATING_METHODS.has(request.method); if (csrfRequired) { const csrfToken = request.headers.get("x-csrf-token") ?? request.headers.get("csrf-token") ?? ""; const valid = await validateCsrfToken(csrfToken); if (!valid) { return NextResponse.json( { ok: false, error: "Invalid or missing CSRF token" }, { status: 403 }, ); } } if (MUTATING_METHODS.has(request.method)) { const contentLength = request.headers.get("content-length"); const maxBytes = options.maxBodyBytes ?? MAX_BODY_BYTES; if (contentLength && Number(contentLength) > maxBytes) { return NextResponse.json( { ok: false, error: `Request body exceeds ${maxBytes} bytes` }, { status: 413 }, ); } } const context = await getApiAdminContext(); if (!context) return NextResponse.json( { ok: false, error: "Unauthorized" }, { status: 401 }, ); if ( options.permission && !canAccess( context.permissions, options.permission, context.session.user.rank, ) ) { await logAuthorizationEvent({ kind: "permission.denied", userId: context.session.user.id, username: context.session.user.username, rank: context.session.user.rank, permission: options.permission, source: request.nextUrl.pathname, reason: "API permission check denied", }); return NextResponse.json( { ok: false, error: "Forbidden" }, { status: 403 }, ); } let finishExport: (() => Promise) | undefined; try { if ( catalogExportEnabled() && isCatalogMutation(request.method, request.nextUrl.pathname) ) { finishExport = await beginCatalogExport(); } const response = await handler(request, context, routeContext); if ( finishExport && response.body && response.headers.get("content-type")?.includes("text/event-stream") ) { const [client, completion] = response.body.tee(); const finish = finishExport; after(async () => { const reader = completion.getReader(); try { while (!(await reader.read()).done) { /* Wait for all import results. */ } } finally { reader.releaseLock(); await finish(); } }); return new Response(client, { status: response.status, headers: response.headers, }); } await finishExport?.(); return response; } catch (error) { await finishExport?.().catch(() => undefined); const errorId = logServerError("admin.api_failed", error, { path: request.nextUrl.pathname, userId: context.session.user.id, }); return NextResponse.json( { ok: false, error: "Internal server error", errorId }, { status: 500 }, ); } }; }