← back to Omega Watches 2

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

54 lines

import { NextRequest, NextResponse } from "next/server";
import { query, 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", "operator"]);

  const { source_name, job_type = "incremental" } = await request.json();

  if (!source_name) {
    return NextResponse.json(
      { success: false, error: { code: "MISSING_SOURCE", message: "source_name is required" } },
      { status: 400 }
    );
  }

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

  // Check ToS compliance
  if (source.tos_reviewed && source.tos_allows_collection === false) {
    return NextResponse.json(
      { success: false, error: { code: "TOS_BLOCKED", message: `Collection blocked by ToS for ${source_name}. Manual collection only.` } },
      { status: 403 }
    );
  }

  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, started_at)
     VALUES ($1, $2, 'pending', $3, NOW())
     RETURNING *`,
    [source.id, job_type, parser?.version || "1.0.0"]
  );

  return NextResponse.json({
    success: true,
    data: { run_id: run.run_id, status: "pending", message: `Job queued for ${source_name}` },
  });
}