import { Router } from "express"; import { z } from "zod"; import { db } from "../db/client.js"; import { hashToken } from "../services/tokens.js"; import { syncServerTasks } from "../services/taskSync.js"; import { asyncHandler } from "../utils/asyncHandler.js"; export const agentReportRouter = Router(); const listeningPortSchema = z.object({ protocol: z.enum(["tcp", "udp"]), port: z.number().int().min(1).max(65535), address: z.string().max(100), process: z.string().max(100).optional(), }); const systemSchema = z.object({ ip_addresses: z.array(z.string()).optional(), cpu: z.object({ model: z.string().optional(), cores: z.number().optional(), load_percent: z.number().nullable().optional() }).optional(), memory: z.object({ total_bytes: z.number().optional(), used_bytes: z.number().optional() }).optional(), disks: z.array(z.object({ mount: z.string(), size_bytes: z.number(), used_bytes: z.number() })).optional(), // Deliberately lenient: one odd line from `ss` must never cost the agent its whole report (tasks included), // so entries are validated one by one and bad ones dropped rather than failing the request. listening_ports: z .array(z.unknown()) .max(5000) .optional() .transform((entries) => entries?.flatMap((e) => { const parsed = listeningPortSchema.safeParse(e); return parsed.success ? [parsed.data] : []; })), }); const reportSchema = z.object({ hostname: z.string().max(255).optional(), os_type: z.string().optional(), reported_at: z.string().optional(), system: systemSchema.nullable().optional(), tasks: z.array( z.object({ schedule_type: z.enum(["cron", "systemd_timer", "windows_task"]), name: z.string().min(1), command: z.string().optional(), schedule_expression: z.string().optional(), source: z.string().optional(), enabled: z.boolean().optional(), next_run_at: z.string().optional(), metadata: z.unknown().optional(), }), ), }); agentReportRouter.post("/", asyncHandler(async (req, res) => { const authHeader = req.header("authorization") ?? ""; const match = authHeader.match(/^Bearer\s+(.+)$/i); if (!match) { return res.status(401).json({ error: "missing_token" }); } const tokenHash = hashToken(match[1]); const server = await db.query.servers.findFirst({ where: (s, { eq }) => eq(s.apiTokenHash, tokenHash), }); if (!server) { return res.status(401).json({ error: "invalid_token" }); } const parsed = reportSchema.safeParse(req.body); if (!parsed.success) { return res.status(400).json({ error: "invalid_body", details: parsed.error.flatten() }); } await syncServerTasks(server.id, { hostname: parsed.data.hostname, osType: parsed.data.os_type, system: parsed.data.system, tasks: parsed.data.tasks.map((t) => ({ scheduleType: t.schedule_type, name: t.name, command: t.command, scheduleExpression: t.schedule_expression, source: t.source, enabled: t.enabled, nextRunAt: t.next_run_at, metadata: t.metadata, })), }); res.status(202).json({ ok: true, taskCount: parsed.data.tasks.length }); }));