From 5e4fc9ab5957b0720a08f5548e8d5e0482947bcd Mon Sep 17 00:00:00 2001 From: openhands Date: Wed, 23 Sep 2026 14:24:44 +0200 Subject: [PATCH] feat(security): give back to CrowdSec and harden the CTI budget - Bound the in-process verdict cache (FIFO eviction at 2000 entries) so a flood of distinct bucket-tripping IPs cannot grow it without limit. - Record block metadata (reputation, score, behaviors, category, TTL) in antiddos:block:meta:{ip}, surfaced as the reason in the admin block list; unban now also clears the metadata and report locks. - Track daily CTI enrichment usage in Redis (crowdsec:usage:{date}); warn once at 80% and pause lookups until tomorrow at CROWDSEC_CTI_DAILY_QUOTA (default 10000, 0 = unlimited) so a via-spread DDoS cannot burn the plan. - Add opt-in signal push to the CrowdSec community (CAPI watcher): stable auto-generated 48-char machine_id/password pair persisted in Redis (or via env), one-time registration, cached JWT login, optional Console enrollment, and POST /v3/signals with a ban decision, deduped per IP. Never throws and reports last status to the admin panel with a verify action. - Admin page: quota usage bar, reporting status/verify channel, and CrowdSec block reasons in the active-blocks list. --- .env.example | 19 + src/actions/admin-antiddos.ts | 19 + src/app/admin/devops/antiddos/page.tsx | 185 +++++++-- src/env.ts | 25 ++ src/lib/crowdsec-api.test.ts | 106 +++++ src/lib/crowdsec-api.ts | 185 ++++++++- src/lib/crowdsec-report.test.ts | 280 +++++++++++++ src/lib/crowdsec-report.ts | 541 +++++++++++++++++++++++++ 8 files changed, 1323 insertions(+), 37 deletions(-) create mode 100644 src/lib/crowdsec-report.test.ts create mode 100644 src/lib/crowdsec-report.ts diff --git a/.env.example b/.env.example index af2d5a41..d2855327 100644 --- a/.env.example +++ b/.env.example @@ -89,6 +89,25 @@ CROWDSEC_BLOCK_SCORE=4 CROWDSEC_BLOCK_TTL_SECONDS=86400 # Endpoint — override only for tests/staging. CROWDSEC_CTI_BASE_URL=https://cti.api.crowdsec.net/v2 +# Daily enrichment-call ceiling (freemium plan ≈ 10k/day). Once today's +# counter reaches it, reputation lookups pause until tomorrow so a spread +# DDoS cannot silently burn the whole quota. 0 = unlimited. +CROWDSEC_CTI_DAILY_QUOTA=10000 + +# --- CROWDSEC SIGNAL PUSH (share our blocks back, optional) --- +# Opt-in: pushes blocked IPs + behaviors to the CrowdSec Central API (CAPI) so +# the community blocklist protects other members too. Set to "true" to enable. +# Requires watcher credentials — either set both CROWDSEC_REPORT_MACHINE_ID +# (48 chars, [A-Za-z0-9]) and CROWDSEC_REPORT_PASSWORD now, or leave them +# unset and let the app generate a stable pair persisted in Redis automatically. +CROWDSEC_REPORT_ENABLED=false +CROWDSEC_REPORT_MACHINE_ID= +CROWDSEC_REPORT_PASSWORD= +# Optional: attachment key from https://app.crowdsec.net → Console settings — +# links our watcher to your account so pushed signals show up there. +CROWDSEC_REPORT_ENROLL_KEY= +# Central API base — override only for tests/staging. +CROWDSEC_CAPI_BASE_URL=https://api.crowdsec.net/v3 # --- PATHS --- BADGE_UPLOAD_DIR=./public/assets/images/badges diff --git a/src/actions/admin-antiddos.ts b/src/actions/admin-antiddos.ts index 3ab25ff0..48cb13a9 100644 --- a/src/actions/admin-antiddos.ts +++ b/src/actions/admin-antiddos.ts @@ -19,6 +19,10 @@ import { setLastCrowdsecVerify, verifyCrowdsecConnection, } from "@/lib/crowdsec-api"; +import { + setLastCrowdsecReport, + verifyCrowdsecReporting, +} from "@/lib/crowdsec-report"; import { db, WebsiteSetting } from "@/lib/db"; import { logger } from "@/lib/logger"; import { PERMS } from "@/lib/permissions"; @@ -223,6 +227,8 @@ export async function unbanAntiddosIp(formData: FormData): Promise { if (redis) { await Promise.all([ redis.del(`antiddos:block:${ip}`), + redis.del(`antiddos:block:meta:${ip}`), + redis.del(`crowdsec:report:${ip}`), redis.del(`antiddos:v:${ip}`), ]); } @@ -272,6 +278,19 @@ export async function verifyCrowdsecConfiguration(): Promise { revalidatePath("/admin/devops/antiddos"); } +/** Test the CrowdSec signal-push (CAPI watcher) channel. */ +export async function verifyCrowdsecReportingConfiguration(): Promise { + const staff = await requirePermission(PERMS.SETTINGS_VIEW); + const status = await verifyCrowdsecReporting(); + await setLastCrowdsecReport(status); + logger.info("CrowdSec reporting configuration verified", { + staff: staff.username, + ok: status.ok, + message: status.message, + }); + revalidatePath("/admin/devops/antiddos"); +} + /** Test the configured Cloudflare API credentials against the zone. */ export async function verifyCloudflareConfiguration(): Promise { const staff = await requirePermission(PERMS.SETTINGS_VIEW); diff --git a/src/app/admin/devops/antiddos/page.tsx b/src/app/admin/devops/antiddos/page.tsx index 4f5ca068..7ddf1cfb 100644 --- a/src/app/admin/devops/antiddos/page.tsx +++ b/src/app/admin/devops/antiddos/page.tsx @@ -15,6 +15,7 @@ import { unbanAntiddosIp, verifyCloudflareConfiguration, verifyCrowdsecConfiguration, + verifyCrowdsecReportingConfiguration, } from "@/actions/admin-antiddos"; import { Badge } from "@/components/ui/badge"; import { Button } from "@/components/ui/button"; @@ -34,9 +35,16 @@ import { } from "@/lib/cloudflare-api"; import { CROWDSEC_BLOCK_SOURCE, + type CrowdsecBlockMeta, crowdsecEnabled, + getCrowdsecBlockMeta, + getCrowdsecQuotaUsage, getLastCrowdsecVerify, } from "@/lib/crowdsec-api"; +import { + crowdsecReportEnabled, + getLastCrowdsecReport, +} from "@/lib/crowdsec-report"; import { db, WebsiteSetting } from "@/lib/db"; import { canAccess, getAdminContext, PERMS } from "@/lib/permissions"; import { redis } from "@/lib/redis"; @@ -76,6 +84,7 @@ export default async function AdminAntiDdosPage() { ttlMs: number; count: number; source: "gate" | "crowdsec"; + meta: CrowdsecBlockMeta | null; }[] = []; let redisOk = false; const rateStore = redis; @@ -100,15 +109,19 @@ export default async function AdminAntiDdosPage() { rateStore.pttl(key), rateStore.get(key), ]); + const ip = key.replace("antiddos:block:", ""); + const source = + value === CROWDSEC_BLOCK_SOURCE + ? ("crowdsec" as const) + : ("gate" as const); return { - ip: key.replace("antiddos:block:", ""), + ip, ttlMs: ttlMs > 0 ? ttlMs : 0, - count: violationCounts.get(key.replace("antiddos:block:", "")) ?? 0, + count: violationCounts.get(ip) ?? 0, // The gate writes "1"; "crowdsec" marks a community-reputation block. - source: - value === CROWDSEC_BLOCK_SOURCE - ? ("crowdsec" as const) - : ("gate" as const), + source, + // Why CrowdSec blocked this IP, when the meta was recorded. + meta: source === "crowdsec" ? await getCrowdsecBlockMeta(ip) : null, }; }), ); @@ -131,6 +144,9 @@ export default async function AdminAntiDdosPage() { const lastVerify = await getLastCloudflareVerify(); const crowdsecConfigured = crowdsecEnabled(); const lastCrowdsecVerify = await getLastCrowdsecVerify(); + const crowdsecUsage = redisOk ? await getCrowdsecQuotaUsage() : null; + const reportingEnabled = await crowdsecReportEnabled(); + const lastReport = await getLastCrowdsecReport(); return (
@@ -195,7 +211,9 @@ export default async function AdminAntiDdosPage() { {crowdsecConfigured ? "Connected" : "Not configured"}

- Community reputation auto-block + {crowdsecUsage && crowdsecUsage.quota > 0 + ? `${crowdsecUsage.used.toLocaleString()} / ${crowdsecUsage.quota.toLocaleString()} CTI calls today${crowdsecUsage.exhausted ? " (paused)" : ""}` + : "Community reputation auto-block"}

@@ -485,28 +503,44 @@ export default async function AdminAntiDdosPage() {

) : (
- {blocks.map((b) => ( -
- {b.ip} - - {b.source === "crowdsec" ? ( - CrowdSec - ) : ( - Gate - )} - TTL {seconds(b.ttlMs)} · violations {b.count} - -
- - -
-
- ))} + {blocks.map((b) => { + const behaviorLabel = + b.meta && b.meta.behaviors.length > 0 + ? b.meta.behaviors.join(", ") + : null; + return ( +
+ {b.ip} + + {b.source === "crowdsec" ? ( + CrowdSec + ) : ( + Gate + )} + TTL {seconds(b.ttlMs)} · violations {b.count} + {b.meta && ( + + {b.meta.reputation ?? "unknown"} · score{" "} + {b.meta.score} · {b.meta.category} + {behaviorLabel ? ` · ${behaviorLabel}` : ""} + + )} + +
+ + +
+
+ ); + })}
)} @@ -642,9 +676,100 @@ export default async function AdminAntiDdosPage() { Verdicts are looked up lazily for IPs that already triggered a rate bucket (never on the per-request hot path), cached for an hour, and blocked IPs show a{" "} - CrowdSec badge in the list above. + CrowdSec badge in the list above + with the community reasoning (reputation, score, behaviors).

)} + + {crowdsecUsage && ( +
+

+ Reputation lookups today +

+ {crowdsecUsage.quota > 0 ? ( + <> +
+
+
= crowdsecUsage.quota * 0.8 + ? "var(--admin-accent)" + : "var(--color-primary)", + }} + /> +
+ + {crowdsecUsage.used.toLocaleString()} /{" "} + {crowdsecUsage.quota.toLocaleString()} + +
+

+ {crowdsecUsage.exhausted + ? "Quota spent for today — reputation lookups are paused until tomorrow (admin via CROWDSEC_CTI_DAILY_QUOTA)." + : "Visible in the env via CROWDSEC_CTI_DAILY_QUOTA (0 = unlimited). Lookups pause at the ceiling to protect the plan."} +

+ + ) : ( +

+ Tracking disabled (CROWDSEC_CTI_DAILY_QUOTA = 0 / unlimited). +

+ )} +
+ )} + +
+

Community signal push

+
+ + {reportingEnabled ? "Enabled" : "Off"} + + {!reportingEnabled && ( +

+ Set{" "} + + CROWDSEC_REPORT_ENABLED=true + {" "} + to share blocked IPs back into the CrowdSec community + blocklist. Watcher credentials are auto-generated and + persisted in Redis. +

+ )} + {reportingEnabled && ( +

+ Blocked IPs are pushed to the Central API (deduped per IP) so + the community blocklist protects other members too. +

+ )} +
+ +
+
+ {lastReport && ( +

+ + {lastReport.ok ? "Push healthy" : "Push failed"} + + + {lastReport.ok + ? `Last signal accepted ${new Date(lastReport.at).toLocaleString()}` + : `${lastReport.message ?? "unknown"} (${new Date(lastReport.at).toLocaleString()})`} + +

+ )} +
diff --git a/src/env.ts b/src/env.ts index 3cce5ae1..234cf2cc 100644 --- a/src/env.ts +++ b/src/env.ts @@ -177,6 +177,31 @@ const schema = z .int() .positive() .default(86_400), + // Daily CTI enrichment quota guard (freemium plan ≈ 10k lookups/day). + // The gate stops consulting the API once the counter for today exceeds + // it, so a spread DDoS can never silently burn the whole quota; 0 + // disables the guard. + CROWDSEC_CTI_DAILY_QUOTA: z.coerce.number().int().min(0).default(10_000), + // Share our own detections back into the CrowdSec community blocklist + // (signal push over the Central API). Opt-in: flipping this on publicly + // shares blocked IPs + behaviors, so it defaults to off. + CROWDSEC_REPORT_ENABLED: z + .string() + .optional() + .transform((value) => value === "true" || value === "1"), + // Central API (CAPI) base endpoint; overridden for tests/staging. + CROWDSEC_CAPI_BASE_URL: z + .string() + .url() + .default("https://api.crowdsec.net/v3"), + // Watcher credentials for signal push. When omitted, a stable pair is + // generated once and persisted in Redis (48-char alnum machine id, + // per the CAPI schema). + CROWDSEC_REPORT_MACHINE_ID: z.string().optional(), + CROWDSEC_REPORT_PASSWORD: z.string().optional(), + // Optional attachment key from the CrowdSec Console — links our + // watcher to your account so pushed signals show up there. + CROWDSEC_REPORT_ENROLL_KEY: z.string().optional(), }) .superRefine((data, ctx) => { if (data.NODE_ENV !== "production") return; diff --git a/src/lib/crowdsec-api.test.ts b/src/lib/crowdsec-api.test.ts index 44c79a6d..dae6e106 100644 --- a/src/lib/crowdsec-api.test.ts +++ b/src/lib/crowdsec-api.test.ts @@ -3,7 +3,11 @@ import { type CrowdsecVerdict, crowdsecEnabled, getCrowdsecApiConfig, + getCrowdsecBlockMeta, + getCrowdsecQuotaUsage, getLastCrowdsecVerify, + getMemoryVerdictCacheSize, + lookupCrowdsecVerdict, maybeAutoBlockCrowdsec, resetCrowdsecCache, setLastCrowdsecVerify, @@ -34,6 +38,12 @@ vi.mock("@/lib/redis", () => ({ for (const key of keys) state.map.delete(key); return keys.length; }, + incr: async (key: string) => { + const next = (Number(state.map.get(key)) || 0) + 1; + state.map.set(key, String(next)); + return next; + }, + expire: async () => 1, pttl: async () => 60_000, }, __esModule: true, @@ -44,6 +54,7 @@ vi.mock("@/lib/logger", () => ({ info: vi.fn(), warn: vi.fn(), error: vi.fn(), + debug: vi.fn(), }, })); @@ -419,4 +430,99 @@ describe("crowdsec-api", () => { status, ); }); + + it("records why it blocked an IP next to the gate key", async () => { + vi.stubEnv("CROWDSEC_API_KEY", "cs_key"); + fetchMock.mockResolvedValue(jsonResponse(maliciousItem(blockIp()))); + + await maybeAutoBlockCrowdsec({ + ip: blockIp(), + category: "api", + ttlSeconds: 86_400, + scoreThreshold: 4, + enabled: true, + }); + + const meta = await getCrowdsecBlockMeta(blockIp()); + expect(meta).not.toBeNull(); + expect(meta?.source).toBe("crowdsec"); + expect(meta?.reputation).toBe("malicious"); + expect(meta?.score).toBe(5); + expect(meta?.behaviors).toEqual(["http:bruteforce", "http:scan"]); + expect(meta?.category).toBe("api"); + expect(meta?.ttlSeconds).toBe(86_400); + expect(meta?.blockedAt).toBeGreaterThan(0); + }); + + it("stops consulting the API once today's quota is spent", async () => { + vi.stubEnv("CROWDSEC_API_KEY", "cs_key"); + vi.stubEnv("CROWDSEC_CTI_DAILY_QUOTA", "2"); + // Fresh Response per call — a consumed body must never be re-parsed. + fetchMock.mockImplementation(() => + Promise.resolve(jsonResponse(maliciousItem(blockIp()))), + ); + + await maybeAutoBlockCrowdsec({ + ip: blockIp(), + category: "api", + ttlSeconds: 600, + scoreThreshold: 4, + enabled: true, + }); + await maybeAutoBlockCrowdsec({ + ip: "198.51.100.2", + category: "api", + ttlSeconds: 600, + scoreThreshold: 4, + enabled: true, + }); + expect(fetchMock).toHaveBeenCalledTimes(2); + expect(state.map.get(`antiddos:block:198.51.100.2`)).toBe("crowdsec"); + + // Third bucket-tripping IP arrives after the quota counter hit 2. + await maybeAutoBlockCrowdsec({ + ip: "198.51.100.3", + category: "api", + ttlSeconds: 600, + scoreThreshold: 4, + enabled: true, + }); + expect(fetchMock).toHaveBeenCalledTimes(2); + expect(state.map.has(`antiddos:block:198.51.100.3`)).toBe(false); + + const usage = await getCrowdsecQuotaUsage(); + expect(usage.quota).toBe(2); + expect(usage.used).toBe(2); + expect(usage.exhausted).toBe(true); + }); + + it("exposes today's quota usage for the admin panel", async () => { + vi.stubEnv("CROWDSEC_CTI_DAILY_QUOTA", "10000"); + const before = await getCrowdsecQuotaUsage(); + expect(before.quota).toBe(10000); + expect(before.used).toBe(0); + expect(before.exhausted).toBe(false); + expect(before.date).toMatch(/^\d{4}-\d{2}-\d{2}$/); + + state.map.set(`crowdsec:usage:${before.date}`, "9876"); + const after = await getCrowdsecQuotaUsage(); + expect(after.used).toBe(9876); + }); + + it("caps the in-process verdict cache so it cannot grow forever", async () => { + vi.stubEnv("CROWDSEC_API_KEY", "cs_key"); + vi.stubEnv("CROWDSEC_CTI_DAILY_QUOTA", "0"); + fetchMock.mockImplementation((url: string | URL) => + Promise.resolve( + jsonResponse(maliciousItem(String(url).split("/").pop() ?? "ip")), + ), + ); + + // One lookup per distinct IP (never cached before), exceeding the cap — + // the oldest entries are evicted first, so the cache stays bounded. + for (let i = 0; i < 2100; i += 1) { + await lookupCrowdsecVerdict(`198.51.100.${i}`); + } + expect(getMemoryVerdictCacheSize()).toBe(2000); + }); }); diff --git a/src/lib/crowdsec-api.ts b/src/lib/crowdsec-api.ts index 03bff31b..aaf60ff6 100644 --- a/src/lib/crowdsec-api.ts +++ b/src/lib/crowdsec-api.ts @@ -1,6 +1,7 @@ import "server-only"; import { env } from "@/env"; +import { reportCrowdsecSignal } from "@/lib/crowdsec-report"; import { logger } from "@/lib/logger"; import { redis } from "@/lib/redis"; import { UNKNOWN_CLIENT_IP } from "./client-ip"; @@ -28,9 +29,15 @@ import { UNKNOWN_CLIENT_IP } from "./client-ip"; * Credentials come from env only (`CROWDSEC_API_KEY`) and are never written * into the admin-visible config — same contract as the Cloudflare token. * - * Note: sharing our own blocks back into the community (signal push) is not - * part of this module — CrowdSec's report channel requires a full Security - * Engine / CAPI machine enrollment, not a CTI API key. + * Quota guard: every enrichment call counts against a per-day Redis counter so + * a spread DDoS (many distinct IPs tripping buckets) can exhaust the day's + * freemium quota only until the configured ceiling, after which lookups pause + * until tomorrow instead of hammering a 429 wall. + * + * Sharing detections back: after a block is created this module fires the + * signal push in `@/lib/crowdsec-report` (Central API watcher login + POST + * /signals), strictly opt-in via CROWDSEC_REPORT_ENABLED and always + * fire-and-forget. */ export class CrowdsecApiError extends Error {} @@ -70,6 +77,27 @@ export interface CrowdsecConnectionStatus { at: number; } +/** Why a CrowdSec-sourced block exists — persisted next to the block key. */ +export interface CrowdsecBlockMeta { + source: typeof CROWDSEC_BLOCK_SOURCE; + category: string; + reputation: CrowdsecReputation | null; + score: number; + behaviors: string[]; + ttlSeconds: number; + blockedAt: number; +} + +/** Daily CTI usage counter as shown in the admin panel. */ +export interface CrowdsecQuotaUsage { + /** UTC calendar day the counter belongs to (YYYY-MM-DD). */ + date: string; + used: number; + /** 0 = unlimited. */ + quota: number; + exhausted: boolean; +} + /** Value written into the shared block key so the admin UI can label the source. */ export const CROWDSEC_BLOCK_SOURCE = "crowdsec"; @@ -82,6 +110,13 @@ const AUTH_BACKOFF_MS = 300_000; const VERDICT_PREFIX = "crowdsec:cti:"; const LOOKUP_LOCK_PREFIX = "crowdsec:lock:"; const LAST_VERIFY_KEY = "crowdsec:last-verify"; +const BLOCK_META_PREFIX = "antiddos:block:meta:"; +const QUOTA_PREFIX = "crowdsec:usage:"; +const QUOTA_KEY_TTL_SECONDS = 48 * 3_600; +/** In-process verdict cache cap so a flood of distinct IPs cannot grow it forever. */ +const MEMORY_VERDICT_CACHE_MAX = 2_000; +/** Warn at this fraction of the daily quota, once per day. */ +const QUOTA_WARN_RATIO = 0.8; /** Well-known, community-safe address used by the admin "verify" button. */ const PROBE_IP = "1.1.1.1"; @@ -197,6 +232,26 @@ export function verdictIsMalicious( const memoryVerdicts = new Map(); +/** + * Insert/refresh an in-process verdict while keeping the cache bounded: Map + * iteration order is insertion order, so the oldest (leftmost) entry is + * dropped first and re-inserted entries are refreshed to the back. + */ +function rememberVerdict(verdict: CrowdsecVerdict): void { + memoryVerdicts.delete(verdict.ip); + memoryVerdicts.set(verdict.ip, verdict); + while (memoryVerdicts.size > MEMORY_VERDICT_CACHE_MAX) { + const oldest = memoryVerdicts.keys().next(); + if (oldest.done) break; + memoryVerdicts.delete(oldest.value); + } +} + +/** Test hook only — reports the bounded in-process cache size. */ +export function getMemoryVerdictCacheSize(): number { + return memoryVerdicts.size; +} + function verdictKey(ip: string): string { return `${VERDICT_PREFIX}${ip}`; } @@ -211,7 +266,7 @@ async function readVerdictCache(ip: string): Promise { const raw = await redis.get(verdictKey(ip)); if (raw) { const parsed = JSON.parse(raw) as CrowdsecVerdict; - memoryVerdicts.set(ip, parsed); + rememberVerdict(parsed); return parsed; } } catch { @@ -222,7 +277,7 @@ async function readVerdictCache(ip: string): Promise { } async function writeVerdictCache(verdict: CrowdsecVerdict): Promise { - memoryVerdicts.set(verdict.ip, verdict); + rememberVerdict(verdict); if (redis) { try { await redis.set( @@ -256,11 +311,79 @@ async function acquireLookupLock(ip: string): Promise { } let backoffUntil = 0; +let quotaWarnedDate: string | null = null; + +function quotaDate(): string { + return new Date().toISOString().slice(0, 10); +} + +function quotaKey(date: string): string { + return `${QUOTA_PREFIX}${date}`; +} + +/** + * Today's configured ceiling. Coerced because tests run with + * SKIP_ENV_VALIDATION (raw process.env strings, no zod defaults) while + * production gets a parsed number. 0 (or unset) = unlimited. + */ +function dailyQuota(): number { + const raw = Number(env.CROWDSEC_CTI_DAILY_QUOTA ?? 0); + return Number.isFinite(raw) && raw > 0 ? raw : 0; +} + +/** + * Today's CTI usage against the configured daily ceiling. Counted in Redis so + * every instance shares one budget. + */ +export async function getCrowdsecQuotaUsage(): Promise { + const date = quotaDate(); + const quota = dailyQuota(); + let used = 0; + if (redis) { + try { + used = Number((await redis.get(quotaKey(date))) ?? 0); + if (!Number.isFinite(used)) used = 0; + } catch { + // counter unavailable — report zero rather than blocking the admin + } + } + return { date, used, quota, exhausted: quota > 0 && used >= quota }; +} + +async function reserveQuota(): Promise { + const quota = dailyQuota(); + if (quota <= 0) return true; + const date = quotaDate(); + const usage = await getCrowdsecQuotaUsage(); + if (usage.used >= quota) { + // Stop consulting the API for the rest of the day: a flood of distinct + // bucket-tripping IPs would otherwise burn every remaining call and + // then sit in a 429 storm anyway. + return false; + } + if (usage.used >= quota * QUOTA_WARN_RATIO && quotaWarnedDate !== date) { + quotaWarnedDate = date; + logger.warn("[crowdsec-api] CTI daily quota nearing its limit", { + used: usage.used, + quota, + }); + } + if (redis) { + try { + await redis.incr(quotaKey(date)); + await redis.expire(quotaKey(date), QUOTA_KEY_TTL_SECONDS); + } catch { + // best effort — an uncounted call is better than a failed lookup + } + } + return true; +} /** * Community reputation verdict for an IP, from cache when possible. Returns - * null when the API is not configured, the lookup failed, or the API is in - * backoff — never throws, so it is safe on the gate's hot path. + * null when the API is not configured, the lookup failed, the API is in + * backoff, or today's quota is spent — never throws, so it is safe on the + * gate's hot path. */ export async function lookupCrowdsecVerdict( ip: string, @@ -284,6 +407,15 @@ export async function lookupCrowdsecVerdict( const raced = await readVerdictCache(ip); if (raced) return raced; + // Cache miss costs a paid call — reserve against today's quota first. + if (!(await reserveQuota())) { + logger.warn( + "[crowdsec-api] CTI daily quota exhausted — pausing lookups until tomorrow", + { quota: dailyQuota() }, + ); + return null; + } + const response = await crowdsecRequest(`/smoke/${encodeURIComponent(ip)}`); if (response.status === 404) { @@ -358,6 +490,27 @@ export async function maybeAutoBlockCrowdsec(input: { if (existingTtl >= ttlSeconds * 1000) return; await redis.set(blockKey, CROWDSEC_BLOCK_SOURCE, "EX", ttlSeconds); + // Record why this block exists so the admin panel can surface the + // community reasoning (reputation, score, behaviors) for the IP. + const meta: CrowdsecBlockMeta = { + source: CROWDSEC_BLOCK_SOURCE, + category, + reputation: verdict.reputation, + score: verdict.score, + behaviors: verdict.behaviors, + ttlSeconds, + blockedAt: Date.now(), + }; + try { + await redis.set( + `${BLOCK_META_PREFIX}${ip}`, + JSON.stringify(meta), + "EX", + ttlSeconds, + ); + } catch { + // metadata is display sugar only — the block itself is already set. + } // CrowdSec only records the block in the gate's own key. It never // creates Cloudflare edge rules — the gate's own escalation logic is // the only place that may mirror a host-level block to the edge. @@ -372,6 +525,9 @@ export async function maybeAutoBlockCrowdsec(input: { behaviors: verdict.behaviors, }, ); + // Opt-in community signal push (CAPI), fire-and-forget: never awaited, + // never throws, and internally deduped per IP. + void reportCrowdsecSignal({ ip, category, ttlSeconds, verdict, meta }); } catch (error) { logger.error("[crowdsec-api] Automatic IP block failed", { ip, @@ -461,9 +617,24 @@ export async function setLastCrowdsecVerify( } } +/** Why an IP is blocked by CrowdSec, when known. */ +export async function getCrowdsecBlockMeta( + ip: string, +): Promise { + if (!redis) return null; + try { + const raw = await redis.get(`${BLOCK_META_PREFIX}${ip}`); + if (!raw) return null; + return JSON.parse(raw) as CrowdsecBlockMeta; + } catch { + return null; + } +} + /** Test hook only — drop in-memory state between unit runs. */ export function resetCrowdsecCache(): void { memoryVerdicts.clear(); backoffUntil = 0; + quotaWarnedDate = null; lastVerifyMemory = null; } diff --git a/src/lib/crowdsec-report.test.ts b/src/lib/crowdsec-report.test.ts new file mode 100644 index 00000000..8776b089 --- /dev/null +++ b/src/lib/crowdsec-report.test.ts @@ -0,0 +1,280 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import type { CrowdsecVerdict } from "./crowdsec-api"; +import { + type CrowdsecReportStatus, + crowdsecReportEnabled, + getLastCrowdsecReport, + reportCrowdsecSignal, + resetCrowdsecReportCache, + verifyCrowdsecReporting, +} from "./crowdsec-report"; + +// The signal-push watcher is tested against a deterministic in-memory Redis +// fake (NX lock + token cache) and a mocked fetch that routes the CAPI paths. +const state = vi.hoisted(() => ({ map: new Map() })); + +vi.mock("@/lib/redis", () => ({ + redis: { + get: async (key: string) => state.map.get(key) ?? null, + set: async ( + key: string, + value: string, + _mode?: string, + _seconds?: number, + nx?: string, + ) => { + if (nx === "NX" && state.map.has(key)) return null; + state.map.set(key, value); + return "OK"; + }, + del: async (...keys: string[]) => { + for (const key of keys) state.map.delete(key); + return keys.length; + }, + }, + __esModule: true, +})); + +vi.mock("@/lib/logger", () => ({ + logger: { + info: vi.fn(), + warn: vi.fn(), + error: vi.fn(), + debug: vi.fn(), + }, +})); + +const CAPI = "https://capi.example.test/v3"; +const MACHINE = "m".repeat(48); +const PASSWORD = "Strong!1P@ssw0rdStrong!1P@ssw0rd"; + +function jsonResponse(body: unknown, status = 200): Response { + return new Response(JSON.stringify(body), { + status, + headers: { "content-type": "application/json" }, + }); +} + +function signalInput(ip = "198.51.100.9") { + const verdict: CrowdsecVerdict = { + ip, + reputation: "malicious", + score: 5, + aggressiveness: 4, + confidence: "0.95", + behaviors: ["http:bruteforce", "http:scan"], + falsePositive: false, + checkedAt: Date.now(), + }; + return { + ip, + category: "api", + ttlSeconds: 86_400, + verdict, + meta: { + source: "crowdsec" as const, + category: "api", + reputation: verdict.reputation, + score: verdict.score, + behaviors: verdict.behaviors, + ttlSeconds: 86_400, + blockedAt: Date.now(), + }, + }; +} + +describe("crowdsec-report", () => { + let fetchMock: ReturnType; + + function routeCapi(overrides: Record = {}) { + const statusFor = (path: string) => + overrides[path] ?? (path === "/signals" ? 200 : 200); + fetchMock.mockImplementation((url: string) => { + const path = String(url).replace(CAPI, ""); + const status = statusFor(path); + if (status !== 200) { + return Promise.resolve(jsonResponse({ message: "boom" }, status)); + } + if (path === "/watchers/login") { + return Promise.resolve( + jsonResponse({ + token: "jwt-xyz", + expire: new Date(Date.now() + 3_600_000).toISOString(), + }), + ); + } + return Promise.resolve(jsonResponse({})); + }); + } + + beforeEach(() => { + vi.unstubAllGlobals(); + vi.unstubAllEnvs(); + state.map.clear(); + resetCrowdsecReportCache(); + fetchMock = vi.fn(); + vi.stubGlobal("fetch", fetchMock); + vi.stubEnv("CROWDSEC_REPORT_ENABLED", "true"); + vi.stubEnv("CROWDSEC_REPORT_MACHINE_ID", MACHINE); + vi.stubEnv("CROWDSEC_REPORT_PASSWORD", PASSWORD); + vi.stubEnv("CROWDSEC_CAPI_BASE_URL", CAPI); + }); + + afterEach(() => { + vi.unstubAllGlobals(); + vi.unstubAllEnvs(); + state.map.clear(); + resetCrowdsecReportCache(); + vi.restoreAllMocks(); + }); + + it("is enabled only when the toggle and credentials are present", async () => { + expect(await crowdsecReportEnabled()).toBe(true); + + vi.stubEnv("CROWDSEC_REPORT_ENABLED", ""); + resetCrowdsecReportCache(); + expect(await crowdsecReportEnabled()).toBe(false); + + // Machine id supplied but no password: falls back to generating a + // stable credential pair persisted in Redis. + vi.stubEnv("CROWDSEC_REPORT_ENABLED", "true"); + vi.stubEnv("CROWDSEC_REPORT_PASSWORD", ""); + resetCrowdsecReportCache(); + expect(await crowdsecReportEnabled()).toBe(true); + const storedMachine = state.map.get("crowdsec:report:machine"); + expect(storedMachine).toMatch(/^[A-Za-z0-9]{48}$/); + expect(state.map.get("crowdsec:report:pass")).toBeTruthy(); + }); + + it("does nothing when the channel is disabled", async () => { + vi.stubEnv("CROWDSEC_REPORT_ENABLED", ""); + resetCrowdsecReportCache(); + + await reportCrowdsecSignal(signalInput()); + expect(fetchMock).not.toHaveBeenCalled(); + }); + + it("registers once, caches the token and pushes one signal per IP", async () => { + routeCapi(); + + await reportCrowdsecSignal(signalInput("198.51.100.10")); + await reportCrowdsecSignal(signalInput("198.51.100.11")); + await reportCrowdsecSignal(signalInput("198.51.100.10")); + // Let the fire-and-forget network body land. + await new Promise((resolve) => setTimeout(resolve, 20)); + + const urls = fetchMock.mock.calls.map((call) => String(call[0])); + expect(urls.filter((u) => u.endsWith("/watchers/register"))).toHaveLength( + 1, + ); + expect(urls.filter((u) => u.endsWith("/watchers/login"))).toHaveLength(1); + expect(urls.filter((u) => u.endsWith("/signals"))).toHaveLength(2); + // No enrollment requested without an attachment key. + expect(urls.some((u) => u.endsWith("/watchers/enroll"))).toBe(false); + }); + + it("builds a well-formed CrowdSec signal with a ban decision", async () => { + routeCapi(); + + await reportCrowdsecSignal(signalInput()); + await new Promise((resolve) => setTimeout(resolve, 20)); + + const signalsCall = fetchMock.mock.calls.find((call) => + String(call[0]).endsWith("/signals"), + ); + expect(signalsCall).toBeDefined(); + if (!signalsCall) throw new Error("expected a /signals call"); + const init = signalsCall[1] as { + body: string; + headers: Record; + }; + const body = JSON.parse(init.body) as Record[]; + expect(body).toHaveLength(1); + const signal = body[0] as { + machine_id: string; + scenario: string; + scenario_version: string; + source: { scope: string; value: string; ip: string }; + decisions: { + scope: string; + type: string; + value: string; + duration: string; + }[]; + context: { key: string; value: string }[]; + created_at: string; + start_at: string; + stop_at: string; + }; + expect(signal.machine_id).toBe(MACHINE); + expect(signal.scenario).toBe("community/anti-ddos-block"); + expect(signal.scenario_version).toBe("1.0.0"); + expect(signal.source).toEqual({ + scope: "ip", + value: "198.51.100.9", + ip: "198.51.100.9", + }); + expect(signal.decisions).toHaveLength(1); + expect(signal.decisions[0]).toMatchObject({ + origin: "crowdsec", + scope: "ip", + type: "ban", + value: "198.51.100.9", + }); + expect(String(signal.decisions[0].duration)).toMatch(/^24h0m0s$/); + for (const key of ["created_at", "start_at", "stop_at"] as const) { + expect(typeof signal[key]).toBe("string"); + } + expect( + signal.context.find((c) => c.key === "crowdsec_reputation")?.value, + ).toBe("malicious"); + }); + + it("records a healthy last-report state after a successful push", async () => { + routeCapi(); + await reportCrowdsecSignal(signalInput()); + await new Promise((resolve) => setTimeout(resolve, 20)); + const last = await getLastCrowdsecReport(); + expect(last?.ok).toBe(true); + }); + + it("never throws and logs the failure when the CAPI rejects the signal", async () => { + fetchMock.mockImplementation((url: string) => { + const path = String(url).replace(CAPI, ""); + if (path === "/watchers/login") { + return Promise.resolve( + jsonResponse({ + token: "jwt-xyz", + expire: new Date(Date.now() + 3_600_000).toISOString(), + }), + ); + } + if (path === "/signals") { + return Promise.resolve(jsonResponse({ message: "boom" }, 500)); + } + return Promise.resolve(jsonResponse({})); + }); + + await expect(reportCrowdsecSignal(signalInput())).resolves.toBeUndefined(); + await new Promise((resolve) => setTimeout(resolve, 20)); + const last: CrowdsecReportStatus | null = await getLastCrowdsecReport(); + expect(last?.ok).toBe(false); + expect(last?.message).toContain("signal push rejected"); + }); + + it("verifies the watcher channel end to end", async () => { + routeCapi(); + const status = await verifyCrowdsecReporting(); + expect(status.ok).toBe(true); + expect(String(fetchMock.mock.calls[0][0])).toContain("/watchers/register"); + }); + + it("reports a clear reason when verification is impossible", async () => { + vi.stubEnv("CROWDSEC_REPORT_ENABLED", ""); + resetCrowdsecReportCache(); + const status = await verifyCrowdsecReporting(); + expect(status.ok).toBe(false); + expect(status.message).toContain("CROWDSEC_REPORT_ENABLED"); + expect(fetchMock).not.toHaveBeenCalled(); + }); +}); diff --git a/src/lib/crowdsec-report.ts b/src/lib/crowdsec-report.ts new file mode 100644 index 00000000..7a6a852d --- /dev/null +++ b/src/lib/crowdsec-report.ts @@ -0,0 +1,541 @@ +import "server-only"; + +import { createHash, randomBytes } from "node:crypto"; +import { env } from "@/env"; +import type { + CrowdsecBlockMeta, + CrowdsecConnectionStatus, + CrowdsecVerdict, +} from "@/lib/crowdsec-api"; +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: "atomcms-next anti-DDoS gate blocked a community-flagged IP", + 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)}`, + ); + } + const status: CrowdsecReportStatus = { ok: true, at: Date.now() }; + await setLastCrowdsecReport(status); + 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); + 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; +}