← back to Omega Watches 2

src/app/api/jobs/backfill/route.ts

51 lines

import { NextRequest, NextResponse } from "next/server";
import { queryOne } from "@/lib/db";
import { getSession, requireRole } from "@/lib/auth";

export async function POST(request: NextRequest) {
  const session = await getSession(request);
  if (!session) {
    return NextResponse.json({ success: false, error: { code: "UNAUTHORIZED", message: "Login required" } }, { status: 401 });
  }
  requireRole(session, ["admin"]);

  const { source_name, date_from, date_to } = await request.json();

  if (!source_name || !date_from || !date_to) {
    return NextResponse.json(
      { success: false, error: { code: "MISSING_PARAMS", message: "source_name, date_from, date_to required" } },
      { status: 400 }
    );
  }

  const source = await queryOne("SELECT * FROM data_source WHERE name = $1", [source_name]);
  if (!source) {
    return NextResponse.json(
      { success: false, error: { code: "NOT_FOUND", message: "Source not found" } },
      { status: 404 }
    );
  }

  const parser = await queryOne(
    "SELECT version FROM parser_version WHERE source_id = $1 AND is_current = true",
    [source.id]
  );

  const run = await queryOne(
    `INSERT INTO collector_run (source_id, job_type, status, parser_version, metadata, started_at)
     VALUES ($1, 'backfill', 'pending', $2, $3, NOW())
     RETURNING *`,
    [source.id, parser?.version || "1.0.0", JSON.stringify({ date_from, date_to, requested_by: session.username })]
  );

  return NextResponse.json({
    success: true,
    data: {
      run_id: run.run_id,
      status: "pending",
      message: `Backfill queued for ${source_name} from ${date_from} to ${date_to}. Requires admin approval.`,
      approval_required: true,
    },
  });
}