import "server-only"; import { createHash, randomBytes } from "node:crypto"; import { env } from "@/env"; import { raiseCrowdsecAlert } from "@/lib/crowdsec-alerts"; import type { CrowdsecBlockMeta, CrowdsecConnectionStatus, CrowdsecVerdict, } from "@/lib/crowdsec-api"; import { bumpCrowdsecStat } from "@/lib/crowdsec-stats"; import { logger } from "@/lib/logger"; import { redis } from "@/lib/redis"; import { UNKNOWN_CLIENT_IP } from "./client-ip"; /** * CrowdSec Central API (CAPI) signal push — the "give back" side of the * anti-DDoS pipeline. * * When the gate blocks an IP based on the CTI community reputation, this * module reports that detection back to CrowdSec (POST /v3/signals) so the * community blocklist also protects every other member. Strictly opt-in * (CROWDSEC_REPORT_ENABLED) and always fire-and-forget: a failure here never * blocks the hot path, never throws to the caller, and records the last * outcome for the admin panel. * * A plain CTI API key cannot push signals, so we act as a CAPI "watcher": * 1. generate/load a stable 48-char alnum machine_id + password pair * (persisted in Redis when not provided via env), * 2. register it once (POST /v3/watchers/register), * 3. login to obtain a JWT (POST /v3/watchers/login), cached in Redis and * refreshed against its expiry, * 4. optional console enrollment via attachment key (POST /v3/watchers/enroll), * 5. push the block as a signal (POST /v3/signals), deduped per IP. */ const API_TIMEOUT_MS = 10_000; const TOKEN_CACHE_KEY = "crowdsec:report:token"; const MACHINE_KEY = "crowdsec:report:machine"; const PASSWORD_KEY = "crowdsec:report:pass"; const REGISTERED_KEY = "crowdsec:report:registered"; const ENROLLED_KEY = "crowdsec:report:enrolled"; const LAST_REPORT_KEY = "crowdsec:last-report"; const REPORT_LOCK_PREFIX = "crowdsec:report:"; /** An IP is only reported once per window — the block itself already deters. */ const REPORT_DEDUPE_SECONDS = 6 * 3_600; const SCENARIO = "community/anti-ddos-block"; const SCENARIO_VERSION = "1.0.0"; export class CrowdsecReportError extends Error {} /** True when the reporting channel is switched on AND usable. */ export async function crowdsecReportEnabled(): Promise { // Production parses CROWDSEC_REPORT_ENABLED to a boolean via zod; tests // (SKIP_ENV_VALIDATION) expose the raw env string, so accept both forms. const flag: unknown = env.CROWDSEC_REPORT_ENABLED; if (flag !== true && flag !== "true" && flag !== "1") return false; return (await loadReportCredentials()) !== null; } function isAlnum48(value: string): boolean { return /^[A-Za-z0-9]{48}$/.test(value); } function generateMachineId(): string { // CAPI schema: exactly 48 characters, [A-Za-z0-9]. const alphabet = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789"; const bytes = randomBytes(48); let id = ""; for (let i = 0; i < 48; i += 1) { id += alphabet[bytes[i] % alphabet.length]; } return id; } function generatePassword(): string { // Deliberately generous: each class is present so common password-policy // rules on the CAPI side are satisfied. const upper = "ABCDEFGHIJKLMNOPQRSTUVWXYZ"; const lower = "abcdefghijklmnopqrstuvwxyz"; const digits = "0123456789"; const symbols = "!@#$%^&*()-_=+[]{};:,.?"; const charset = `${upper}${lower}${digits}${symbols}`; const bytes = randomBytes(32); let password = ""; for (let i = 0; i < 8; i += 1) { // Guarantee at least one of each class. const pool = [upper, lower, digits, symbols][i % 4]; password += pool[bytes[i] % pool.length]; } for (let i = 8; i < 32; i += 1) { password += charset[bytes[i] % charset.length]; } return password; } interface ReportCredentials { machineId: string; password: string; } let credentialsCache: ReportCredentials | null = null; let credentialsMissingRedisWarned = false; /** * Load the configured watcher credentials, or generate a stable pair and * persist it in Redis so restarts and other instances reuse the same identity. */ async function loadReportCredentials(): Promise { if (credentialsCache) return credentialsCache; const envMachine = env.CROWDSEC_REPORT_MACHINE_ID?.trim(); const envPassword = env.CROWDSEC_REPORT_PASSWORD; if (envMachine && envPassword) { credentialsCache = { machineId: envMachine, password: envPassword }; return credentialsCache; } if (!redis) { if (!credentialsMissingRedisWarned) { credentialsMissingRedisWarned = true; logger.warn( "[crowdsec-report] Redis is required to persist auto-generated watcher credentials — set CROWDSEC_REPORT_MACHINE_ID and CROWDSEC_REPORT_PASSWORD, or REDIS_URL", ); } return null; } try { let machineId: string | null = null; let password: string | null = null; const storedMachine = await redis.get(MACHINE_KEY); const storedPassword = await redis.get(PASSWORD_KEY); if (storedMachine && isAlnum48(storedMachine)) machineId = storedMachine; if (storedPassword) password = storedPassword; if (!machineId) machineId = generateMachineId(); if (!password) password = generatePassword(); await redis.set(MACHINE_KEY, machineId); await redis.set(PASSWORD_KEY, password); credentialsCache = { machineId, password }; return credentialsCache; } catch { return null; } } async function capiRequest( path: string, init: { method?: "POST" | "GET"; token?: string; body?: unknown } = {}, ): Promise { const controller = new AbortController(); const timer = setTimeout(() => controller.abort(), API_TIMEOUT_MS); const headers: Record = { Accept: "application/json", "Content-Type": "application/json", }; if (init.token) headers.Authorization = `Bearer ${init.token}`; try { return await fetch(`${env.CROWDSEC_CAPI_BASE_URL}${path}`, { method: init.method ?? "POST", headers, body: init.body === undefined ? undefined : JSON.stringify(init.body), signal: controller.signal, cache: "no-store", }); } finally { clearTimeout(timer); } } async function errorDetail(response: Response): Promise { try { const body = (await response.json()) as { message?: string }; return body.message ?? `HTTP ${response.status}`; } catch { return `HTTP ${response.status}`; } } interface CapToken { raw: string; expiresAt: number; } let capToken: CapToken | null = null; let tokenPromise: Promise | null = null; function scenarioHash(): string { return createHash("sha256") .update(`${SCENARIO}:${SCENARIO_VERSION}`) .digest("hex") .slice(0, 16); } async function registerWatcher( machineId: string, password: string, ): Promise { if (!redis) return; try { const already = await redis.get(REGISTERED_KEY); if (already) return; } catch { // proceed anyway — registering is idempotent-ish (400 = already exists) } const response = await capiRequest("/watchers/register", { body: { machine_id: machineId, password }, }); if (response.ok || response.status === 400) { try { await redis.set(REGISTERED_KEY, "1", "EX", 30 * 24 * 3_600); } catch { // fine — will just attempt registration again later } return; } throw new CrowdsecReportError( `watcher registration failed (HTTP ${response.status}): ${await errorDetail(response)}`, ); } /** Login and cache the JWT; guarded against concurrent logins. */ async function acquireCapToken(): Promise { if (capToken && capToken.expiresAt > Date.now() + 60_000) { return capToken.raw; } if (tokenPromise) return tokenPromise; const credentials = await loadReportCredentials(); if (!credentials) { throw new CrowdsecReportError( "CrowdSec reporting is not configured (no watcher credentials)", ); } tokenPromise = (async () => { if (capToken && capToken.expiresAt > Date.now() + 60_000) { return capToken.raw; } if (redis) { try { const cached = await redis.get(TOKEN_CACHE_KEY); if (cached) { const parsed = JSON.parse(cached) as CapToken; if (parsed.expiresAt > Date.now() + 60_000) { capToken = parsed; return parsed.raw; } } } catch { // fall through to a fresh login } } // First signal push needs the watcher to exist on the CAPI. await registerWatcher(credentials.machineId, credentials.password); const response = await capiRequest("/watchers/login", { body: { machine_id: credentials.machineId, password: credentials.password, scenarios: [SCENARIO], }, }); if (!response.ok) { throw new CrowdsecReportError( `CAPI watcher login failed (HTTP ${response.status}): ${await errorDetail(response)}`, ); } const body = (await response.json()) as { token?: string; expire?: string }; if (!body.token) { throw new CrowdsecReportError("CAPI watcher login returned no token"); } let expiresAt = Date.now() + 3_600_000; const expire = body.expire ? Date.parse(body.expire) : NaN; if (Number.isFinite(expire) && expire > Date.now()) { expiresAt = expire; } const token: CapToken = { raw: body.token, expiresAt }; capToken = token; if (redis) { try { await redis.set(TOKEN_CACHE_KEY, JSON.stringify(token), "EX", 3_600); } catch { // in-process view is enough } } return token.raw; })().finally(() => { tokenPromise = null; }); return tokenPromise; } async function enrollWatcher(token: string): Promise { const attachmentKey = env.CROWDSEC_REPORT_ENROLL_KEY?.trim(); if (!attachmentKey) return; if (!redis) return; try { const enrolled = await redis.get(ENROLLED_KEY); if (enrolled) return; } catch { // best effort below } try { const response = await capiRequest("/watchers/enroll", { token, body: { attachment_key: attachmentKey, name: "atomcms-next" }, }); if (response.ok) { try { await redis.set(ENROLLED_KEY, "1", "EX", 30 * 24 * 3_600); } catch { // best effort } } } catch (error) { // Enrollment only affects Console visibility — not worth failing a push. logger.debug("[crowdsec-report] Console enrollment skipped", { err: error instanceof Error ? error.message : String(error), }); } } function formatCapiDuration(seconds: number): string { const safe = Math.max(1, Math.floor(seconds)); const hours = Math.floor(safe / 3_600); const minutes = Math.floor((safe % 3_600) / 60); const rest = safe % 60; if (hours > 0) return `${hours}h${minutes}m${rest}s`; if (minutes > 0) return `${minutes}m${rest}s`; return `${rest}s`; } async function pushSignal(input: { credentials: ReportCredentials; token: string; ip: string; category: string; ttlSeconds: number; verdict: CrowdsecVerdict; meta: CrowdsecBlockMeta; }): Promise { const now = new Date(); try { await enrollWatcher(input.token); const duration = formatCapiDuration(input.ttlSeconds); const response = await capiRequest("/signals", { token: input.token, body: [ { machine_id: input.credentials.machineId, message: input.verdict.reputation ? "atomcms-next anti-DDoS gate blocked a community-flagged IP" : "atomcms-next anti-DDoS gate blocked a repeat rate-limit offender", scenario: SCENARIO, scenario_version: SCENARIO_VERSION, scenario_hash: scenarioHash(), created_at: now.toISOString(), start_at: now.toISOString(), stop_at: now.toISOString(), source: { scope: "ip", value: input.ip, ip: input.ip }, decisions: [ { id: 0, origin: "crowdsec", scenario: SCENARIO, scope: "ip", type: "ban", value: input.ip, duration, }, ], context: [ { key: "crowdsec_category", value: input.category }, { key: "crowdsec_reputation", value: input.verdict.reputation ?? "unknown", }, { key: "crowdsec_score", value: String(input.verdict.score), }, { key: "crowdsec_behaviors", value: input.verdict.behaviors.join(",") || "none", }, ], }, ], }); if (!response.ok) { throw new CrowdsecReportError( `signal push rejected (HTTP ${response.status}): ${await errorDetail(response)}`, ); } // Was the channel down before this? A first success after failures // deserves a recovery notice (separate cooldown key from the failure). const before = await getLastCrowdsecReport(); const status: CrowdsecReportStatus = { ok: true, at: Date.now() }; await setLastCrowdsecReport(status); void bumpCrowdsecStat("reports"); if (before && !before.ok) { void raiseCrowdsecAlert("report-recovered", { type: "ddos", severity: "info", message: "CrowdSec signal push recovered — detections are reaching the community again.", context: { ip: input.ip }, }); } logger.info("[crowdsec-report] Detection shared with the community", { ip: input.ip, category: input.category, ttlSeconds: input.ttlSeconds, score: input.verdict.score, }); return status; } catch (error) { const status: CrowdsecReportStatus = { ok: false, at: Date.now(), message: String(error), }; await setLastCrowdsecReport(status); void bumpCrowdsecStat("report_fail"); // Ops alert, cooldown-gated: a silently broken channel means the // community never learns about the blocks we keep sharing. void raiseCrowdsecAlert("report", { type: "ddos", severity: "warning", message: "CrowdSec signal push failed — detections are not reaching the community.", context: { ip: input.ip, detail: String(error).slice(0, 300), }, }); logger.error("[crowdsec-report] Signal push failed", { ip: input.ip, err: error, }); return status; } } /** * Share a CrowdSec-reputation block back into the community. Safe to call * fire-and-forget: does nothing when reporting is off, never throws, and * dedupes so the same IP is only reported once per window. */ export async function reportCrowdsecSignal(input: { ip: string; category: string; ttlSeconds: number; verdict: CrowdsecVerdict; meta: CrowdsecBlockMeta; }): Promise { const { ip, category, ttlSeconds, verdict, meta } = input; if (!(await crowdsecReportEnabled())) return; if (!ip || ip === UNKNOWN_CLIENT_IP) return; const credentials = await loadReportCredentials(); if (!credentials) return; try { if (redis) { // Cross-instance dedupe: one report per IP per window. const acquired = await redis.set( `${REPORT_LOCK_PREFIX}${ip}`, "1", "EX", REPORT_DEDUPE_SECONDS, "NX", ); if (acquired !== "OK") return; } const token = await acquireCapToken(); await pushSignal({ credentials, token, ip, category, ttlSeconds, verdict, meta, }); } catch (error) { logger.error("[crowdsec-report] Signal push preparation failed", { ip, err: error, }); } } export interface CrowdsecReportStatus { ok: boolean; message?: string; at: number; } let lastReportMemory: CrowdsecReportStatus | null = null; export async function getLastCrowdsecReport(): Promise { if (redis) { try { const raw = await redis.get(LAST_REPORT_KEY); if (raw) return JSON.parse(raw) as CrowdsecReportStatus; } catch { // fall back to the in-process view } } return lastReportMemory; } export async function setLastCrowdsecReport( status: CrowdsecReportStatus, ): Promise { lastReportMemory = status; if (redis) { try { await redis.set( LAST_REPORT_KEY, JSON.stringify(status), "EX", 48 * 3_600, ); } catch { // in-process view is enough } } } /** Probe the full register → login path and report the outcome for the admin. */ export async function verifyCrowdsecReporting(): Promise { if (!env.CROWDSEC_REPORT_ENABLED) { return { ok: false, message: "CROWDSEC_REPORT_ENABLED is not set", at: Date.now(), }; } const credentials = await loadReportCredentials(); if (!credentials) { return { ok: false, message: "CrowdSec reporting is not configured (set CROWDSEC_REPORT_MACHINE_ID + CROWDSEC_REPORT_PASSWORD, or REDIS_URL)", at: Date.now(), }; } try { await acquireCapToken(); return { ok: true, message: "CAPI watcher connected — signal push ready", at: Date.now(), }; } catch (error) { return { ok: false, message: error instanceof Error ? error.message : String(error), at: Date.now(), }; } } /** Test hook only. */ export function resetCrowdsecReportCache(): void { credentialsCache = null; capToken = null; tokenPromise = null; lastReportMemory = null; }