← 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}` },
});
}