Files
EpicNext-Cms/src/lib/crowdsec-report.ts
T
openhands 3e1a3f92c8
Gitea Actions Runner Test / test-job (push) Successful in 1s
CI / check (push) Successful in 30s
CI / tests-integration (push) Successful in 1m42s
CI / tests-unit (push) Successful in 1m50s
CI / tests-ui (push) Successful in 2m42s
CI / preflight (push) Skipped
CI / deploy (push) Successful in 2m3s
feat(security): recovery alerts, gate-block sharing, rolling-window burst and admin breakdown for CrowdSec
2026-09-23 15:06:16 +02:00

572 lines
16 KiB
TypeScript

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<boolean> {
// 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<ReportCredentials | null> {
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<Response> {
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), API_TIMEOUT_MS);
const headers: Record<string, string> = {
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<string> {
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<string> | 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<void> {
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<string> {
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<void> {
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<CrowdsecReportStatus> {
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<void> {
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<CrowdsecReportStatus | null> {
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<void> {
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<CrowdsecConnectionStatus> {
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;
}