Author SHA1 Message Date
remco 305dfdc9f8 Add renovate.json
Gitea Actions Runner Test / test-job (push) Successful in 2s
CI / check (pull_request) Failing after 20s
CI / tests-unit (pull_request) Skipped
CI / tests-integration (pull_request) Skipped
CI / tests-ui (pull_request) Skipped
CI / preflight (pull_request) Skipped
CI / deploy (pull_request) Skipped
2026-10-02 03:00:12 +00:00
39 changed files with 5721 additions and 601 deletions

No files matched your search

+67
View File
@@ -78,6 +78,73 @@ CLOUDFLARE_AUTO_BLOCK_ENABLED=true
# Override for tests/staging (production uses the public endpoint by default).
CLOUDFLARE_API_BASE_URL=https://api.cloudflare.com/client/v4
# --- CROWDSEC API (community reputation auto-block, optional) ---
# Free CTI API key: https://app.crowdsec.net/ → Settings → CTI API Keys.
# When set, the anti-DDoS gate checks the community reputation of repeat
# offenders (CTI GET /smoke/{ip}) and immediately hard-blocks known-bad IPs.
# Lookups only happen for IPs that already tripped a rate bucket and are
# cached in Redis for 1h, so quota usage stays minimal.
CROWDSEC_API_KEY=
# Runtime toggle for reputation-based auto-blocking (also overridable live
# from the admin panel). Requires CROWDSEC_API_KEY.
CROWDSEC_AUTO_BLOCK_ENABLED=true
# Minimum malevolence score 0-5 (CrowdSec scale; 4-5 = "malicious") before an
# IP is treated as known-bad. IPs with false-positive tags are never blocked.
CROWDSEC_BLOCK_SCORE=4
# How long a CrowdSec-confirmed bad IP stays blocked (seconds).
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
# How many new community-reputation blocks within a 5-minute window justify an
# ops alert (quota/backoff/report alerts all use HEALTH_ALERT_COOLDOWN_MIN).
CROWDSEC_ALERT_BLOCK_BURST=10
# --- 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
# --- CROWDSEC LOCAL (opt-in engine on this Docker host, no proxy changes) ---
# App-layer LAPI bouncer: the anti-DDoS gate asks the local engine per client
# IP (short-cached) and blocks ban/captcha decisions before its own buckets.
# Start everything with `bash cms security`; it writes the key below into .env
# and starts the CrowdSec engine bound to 127.0.0.1. Set to "true" to load the
# bouncer without the local engine (not recommended).
CROWDSEC_LOCAL_ENABLED=false
# Host access-log directory mounted into the engine for detection (Nginx only).
CROWDSEC_NGINX_LOG_DIR=/var/log/nginx
# Change LAPI port AND LAPI URL together when 18080 is already taken.
CROWDSEC_LAPI_PORT=18080
CROWDSEC_LAPI_URL=http://127.0.0.1:18080
# Generated by `bash cms security`; keep in .env, never commit a value.
CROWDSEC_LAPI_API_KEY=
# IP blocklist sync (`bash cms security blocklists`): space-separated URLs, by
# default Spamhaus DROP/EDROP, DShield, CINS, Greensnow, StopForumSpam,
# blocklist.de, Emerging Threats, abuse.ch Feodo/SSLBL/URLhaus, IPsum,
# Firehol ipsets and Tor exit nodes. Requires internet to fetch; detection and
# blocking stay local.
#CROWDSEC_BLOCKLIST_SOURCES=https://www.spamhaus.org/drop/drop.txt https://example.org/list.txt
# Expiration for each blocklist decision (re-synced keeps them fresh).
#CROWDSEC_BLOCKLIST_DURATION=24h
# Combined cap per sync (safety valve against excessive decisions).
#CROWDSEC_BLOCKLIST_MAX_DECISIONS=1000000
# Comma-separated IPs/CIDRs that a sync must always skip (allowlist).
#CROWDSEC_BLOCKLIST_ALLOW=1.2.3.4,10.0.0.0/8
# --- PATHS ---
BADGE_UPLOAD_DIR=./public/assets/images/badges
EMULATOR_JAR_PATH=./emulator/Arcturus.jar
+26 -19
View File
@@ -1,35 +1,40 @@
# syntax=docker/dockerfile:1
# Pin the runtime to the supported engine; update both stages deliberately.
FROM node:26.10.0-alpine AS migrations
WORKDIR /app
ENV NEXT_TELEMETRY_DISABLED=1
# Installeer git en pnpm v12
# Keep the bootstrap aligned with package.json packageManager.
# The apk cache is persisted in a BuildKit cache mount so git is not
# re-downloaded on every build.
RUN --mount=type=cache,target=/var/cache/apk \
apk add --no-cache git \
&& npm install -g pnpm@12.8.1
# Stel het PATH zo in dat Alpine pnpm gegarandeerd overal herkent
ENV PNPM_HOME="/usr/local/share/pnpm"
ENV PATH="$PNPM_HOME:/usr/local/bin:$PATH"
&& npm install -g pnpm@11.25.0
# The pnpm store is kept in a BuildKit cache mount that persists across builds
# on the builder. This is what stops disk usage from growing unbounded: the
# downloaded dependency store is shared and reused instead of being copied into
# a fresh image layer on every build. Unlike an image layer it is also prunable
# independently, so a hard cap (see ci-deploy.sh) keeps it bounded.
ENV PNPM_HOME=/pnpm PNPM_STORE=/pnpm/store
# pnpm-workspace.yaml + .npmrc must be present too: the lockfile records the
# overrides from pnpm-workspace.yaml, and --frozen-lockfile rejects a build
# where the workspace config is absent (ERR_PNPM_LOCKFILE_CONFIG_MISMATCH).
COPY package.json pnpm-lock.yaml* pnpm-workspace.yaml* .npmrc* ./
# Voer de installatie uit met de pnpm v12 store cache-mount
RUN --mount=type=cache,target=/root/.local/share/pnpm/store \
pnpm install --frozen-lockfile --ignore-scripts
# pnpm fetch: download all deps into the shared cache-mounted store.
RUN --mount=type=cache,target=/pnpm \
pnpm fetch --ignore-scripts
# Install offline from the cache-mounted store; the store itself stays in the
# build cache between builds.
RUN --mount=type=cache,target=/pnpm \
pnpm install --frozen-lockfile --ignore-scripts --offline
COPY . .
ARG NEXT_DEPLOYMENT_ID="unknown"
LABEL org.opencontainers.image.revision="$NEXT_DEPLOYMENT_ID"
FROM migrations AS builder
ARG NEXT_DEPLOYMENT_ID="unknown"
ENV NEXT_DEPLOYMENT_ID="$NEXT_DEPLOYMENT_ID"
# Fix voor OOM Killer: dwing de Node-compiler om agressief op te ruimen bij 4GB RAM
ENV NODE_OPTIONS="--max-old-space-size=4096"
# Bouw de Next.js applicatie met caching
# Fixture values exist only for this build command; production secrets are runtime-only.
# Cache Next.js build output and webpack caches so rebuilds only redo the
# changed parts.
RUN --mount=type=cache,target=/app/.next/cache \
DATABASE_URL="mysql://build:[email protected]:9/build" \
HOTEL_NAME="Build fixture" APP_URL="http://localhost:3002" \
@@ -57,6 +62,8 @@ COPY --from=builder --chown=nextjs:nextjs /app/drizzle/migrations ./drizzle/migr
COPY --chown=nextjs:nextjs scripts/docker-start.mjs ./docker-start.mjs
USER nextjs
EXPOSE 3002
# Self-contained healthcheck so `docker run` (ci-deploy) also gets Docker-level
# health; docker-compose overrides this with its own probe if needed.
HEALTHCHECK --interval=30s --timeout=5s --start-period=30s --retries=3 \
CMD ["node", "-e", "fetch('http://127.0.0.1:'+(process.env.PORT||'3002')+'/api/health').then(r=>process.exit(r.ok?0:1)).catch(()=>process.exit(1))"]
ENTRYPOINT ["/sbin/tini", "--"]
+126
View File
@@ -838,6 +838,132 @@ Open **DevOps → Anti-DDoS protection**
---
## Local CrowdSec Engine (opt-in)
The repository ships a self-contained CrowdSec engine that runs on the same
Docker host. It runs `crowdsecurity/crowdsec:v1.8.1` in its own Compose project
and exposes **LAPI only** on `127.0.0.1:18080`. When enabled, the app-layer
anti-DDoS gate (`src/lib/crowdsec-local.ts`) asks the local LAPI per client IP
(short-cached) and blocks `ban` / `captcha` decisions before its own rate
buckets run. No reverse-proxy, Traefik, Cloudflare or firewall configuration is
changed.
The engine boots in **LAPI-only mode** (`DISABLE_AGENT=true`): it does not
consume the host Nginx access log and needs no outbound access to
`crowdsec.net`, which is often blocked on hardened hosts. Combined with
`DISABLE_ONLINE_API=true` (no CrowdSec Central API) the engine needs no account
and no inbound internet — blocking comes from the imported blocklists
(Step 4) and the app's own rate buckets. Re-enable the agent only on a host
with outbound internet by removing `DISABLE_AGENT: "true"` from
`deployment/crowdsec/compose.crowdsec.yml`.
### Step 1 — Enable the engine and register the bouncer
```bash
bash cms security
```
This generates `CROWDSEC_LAPI_API_KEY` (random 64 hex chars), writes the
CrowdSec flags into `.env`, starts the engine and registers the `cms` bouncer
against the local LAPI. The engine does **not** enroll into the CrowdSec
Central API (`DISABLE_ONLINE_API=true`) and runs without the agent
(`DISABLE_AGENT=true`), so it never phones home.
### Step 2 — Restart the CMS so it loads the bouncer credentials
A CI-managed `epicnext-cms-app` picks the new `.env` values up on its next
deployment. For a clone running via the updater:
```bash
bash cms update --skip-pull
```
or restart the container directly (`docker compose restart cms`). Without the
restart the gate has not loaded the LAPI URL/key yet.
### Step 3 — Verify
```bash
bash cms security status
```
Expect `CROWDSEC_LOCAL_ENABLED=yes` and `LAPI health: OK (127.0.0.1:18080)`.
In the admin panel, **DevOps → Anti-DDoS protection** shows live block
statistics split per origin (`community` vs `local`).
### Step 4 — Stop the engine again (optional)
```bash
bash cms security disable
```
Stops the container and sets `CROWDSEC_LOCAL_ENABLED=false`. Volumes and the
`.env` key are kept.
### Step 5 — Load external IP blocklists (optional)
The engine has no built-in lists, so provide your own via a one-shot sync
(fetches the sources, replaces every previous `cscli-import` decision):
```bash
bash cms security blocklists
```
Defaults: Spamhaus DROP/EDROP, DShield, CINS, Greensnow, StopForumSpam,
Binary Defense, blocklist.de, Emerging Threats, BruteForceBlocker, abuse.ch
Feodo/SSLBL, Darklist, Botvrij, IPsum, Firehol ipsets and Tor exit nodes
(26 sources). URLhaus was removed because its `text_online` feed lists URLs,
not IPs; a malformed token in it could otherwise expand into a bogus
huge CIDR. The validator only accepts whole-line bare IPs or proper CIDRs,
enforces sane prefix bounds and drops reserved/private/loopback space, so a
bad source entry can never block the origin or internal traffic. The largest
commercial/crowdsourced lists (AbuseIPDB, MaxMind, Cisco Talos, AlienVault
OTX) are not included because they require an account or API key; IPsum
already aggregates ~30 additional feeds. No account is needed, but internet
access is — only for fetching; detection and blocking remain local. Sync
hourly as a cron job:
```bash
bash cms security blocklists-install-cron
```
Removal: `bash cms security blocklists-uninstall-cron`. Dry-run without
touching LAPI: `bash cms security blocklists --dry-run`. Override the sources,
duration, a combined cap or an allowlist in `.env`
(`CROWDSEC_BLOCKLIST_SOURCES`, `CROWDSEC_BLOCKLIST_DURATION`,
`CROWDSEC_BLOCKLIST_MAX_DECISIONS`, `CROWDSEC_BLOCKLIST_ALLOW`). Existing
`cscli-import` decisions are replaced on every sync, so removed entries
expire.
### Environment variables
| Variable | Default | Purpose |
| ---------------------------- | ----------------------------- | ------------------------------------ |
| `CROWDSEC_LOCAL_ENABLED` | `false` | Master switch for the local stack |
| `CROWDSEC_LAPI_URL` | `http://127.0.0.1:18080` | LAPI endpoint (loopback only) |
| `CROWDSEC_LAPI_PORT` | `18080` | Host port the engine maps to LAPI |
| `CROWDSEC_LAPI_API_KEY` | — | Bouncer key; required when enabled |
| `CROWDSEC_LAPI_TIMEOUT_MS` | `500` | Per-decision request timeout |
| `CROWDSEC_LAPI_RETRY_MS` | `500` | Backoff before retrying LAPI |
| `CROWDSEC_NGINX_LOG_DIR` | `/var/log/nginx` | Access-log directory for the engine |
### Notes and limitations
- Changing the port means updating `CROWDSEC_LAPI_PORT` **and**
`CROWDSEC_LAPI_URL` together, then re-running `bash cms security`.
- Rotate the key by editing `CROWDSEC_LAPI_API_KEY` in `.env`, running
`bash cms security` again (re-registers the bouncer) and restarting the CMS.
- The gate is **fail-closed at startup** when the feature is enabled without a
key (startup aborts with a clear message). At runtime a LAPI network error
**fails open** (traffic is allowed, decisions paused); a 403 from LAPI
pauses local decisions for 5 minutes.
- This bouncer is **application-layer**: it sheds known-bad IPs at the CMS
process only. It does not drop traffic before the origin, does not protect
other host ports/services, and depends on the client IP being trustworthy at
the ingress. Keep the upstream protections (Cloudflare IP rules, proxy rate
limits) for defense before the origin.
---
## Production Deployment (blue/green)
+3 -2
View File
@@ -4,6 +4,7 @@ DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
case "${1:-help}" in
install) shift; exec bash "$DIR/scripts/docker-install.sh" "$@" ;;
update) shift; exec bash "$DIR/scripts/docker-update.sh" "$@" ;;
help|--help|-h) printf '%s\n' 'bash cms install Configure and install on a Linux Docker host' 'bash cms update Update using saved settings; --skip-pull uses checked-out release' ;;
*) echo "Unknown command. Use: bash cms install | update" >&2; exit 1 ;;
security) shift; exec bash "$DIR/scripts/crowdsec-setup.sh" "$@" ;;
help|--help|-h) printf '%s\n' 'bash cms install Configure and install on a Linux Docker host' 'bash cms update Update using saved settings; --skip-pull uses checked-out release' 'bash cms security Configure the opt-in local CrowdSec stack (enable|status|disable|blocklists)' ;;
*) echo "Unknown command. Use: bash cms install | update | security" >&2; exit 1 ;;
esac
+4
View File
@@ -0,0 +1,4 @@
filenames:
- /var/log/nginx/access.log
labels:
type: nginx
+41
View File
@@ -0,0 +1,41 @@
services:
crowdsec:
image: crowdsecurity/crowdsec:${CROWDSEC_VERSION:-v1.8.1}
container_name: epicnext-crowdsec
restart: unless-stopped
profiles: ["security"]
environment:
DISABLE_AGENT: "true"
BOUNCER_KEY_cms: ${CROWDSEC_LAPI_API_KEY:?CROWDSEC_LAPI_API_KEY must be set}
DISABLE_ONLINE_API: "true"
GID: "${CROWDSEC_GID:-0}"
TZ: "${TZ:-UTC}"
ports:
- "${CROWDSEC_LAPI_BIND_HOST:-127.0.0.1}:${CROWDSEC_LAPI_PORT:-18080}:8080"
volumes:
- ./acquis.d:/etc/crowdsec/acquis.d:ro
- ${CROWDSEC_NGINX_LOG_DIR:-/var/log/nginx}:/var/log/nginx:ro
- crowdsec-config:/etc/crowdsec
- crowdsec-data:/var/lib/crowdsec/data
security_opt:
- no-new-privileges:true
pids_limit: 256
logging:
driver: json-file
options:
max-size: "10m"
max-file: "3"
healthcheck:
test:
[
"CMD-SHELL",
"wget -q -O - http://127.0.0.1:8080/health >/dev/null 2>&1",
]
interval: 30s
timeout: 5s
retries: 3
start_period: 30s
volumes:
crowdsec-config:
crowdsec-data:
+39 -2
View File
@@ -1,5 +1,18 @@
# ─────────────────────────────────────────────────────────────────────────────
# Next.js CMS — blue/green
#
# De app draait met `network_mode: host`, dus een replica neemt een host-poort in
# plaats van een gedeelde docker-poort. Daarom twee expliciete services in plaats
# van `docker compose up --scale cms=2`: die zou op poort 3002 botsen.
#
# `deploy.sh` start een release op de vrije poort, wacht op /api/health, schrijft
# daarna /etc/nginx/snippets/cms_upstream_servers.conf en herlaadt nginx. Pas dan
# wordt de oude replica gestopt. De hele release is dus zero-downtime: faalt de
# nieuwe replica, dan blijft de oude gewoon draaien.
#
# De YAML-anchor houdt beide replicas identiek. Wil je ze bewust uit elkaar
# halen (bv. één release canary-en), verwijder dan `<<: *cms` en vul de
# afwijkende velden opnieuw in.
# ─────────────────────────────────────────────────────────────────────────────
x-cms: &cms
image: epicnext-cms:${CMS_RELEASE:-local}
@@ -10,6 +23,8 @@ x-cms: &cms
NEXT_DEPLOYMENT_ID: ${CMS_RELEASE:-unknown}
network: host
network_mode: host
# 15s: Next moet een lopend request nog netjes kunnen afronden voordat SIGKILL
# volgt. Met 10s werden streams en imports afgekapt.
stop_grace_period: 15s
restart: unless-stopped
env_file:
@@ -21,11 +36,20 @@ x-cms: &cms
- /var/www/Gamedata:/var/www/Gamedata
# ── Resource limits ──
mem_limit: 6g
memswap_limit: 7g
# De limieten waren eerder weggehaald ("Next mag onbeperkt presteren"). Op een
# gedeelde host is juist dat gevaarlijk: één geheugenlek vult dan de hele
# machine en MariaDB + nginx + Traefik gaan er allemaal onderuit. 4 GiB met
# 1 GiB swap geeft de V8-heap ruimte om zich te organiseren voor hij hard wordt
# afgesneden, maar houdt de schade begrensd. 2 CPU laat drie keer zoveel
# achtergrondwerk toe als de cores, zodat de 6 cores van deze host niet
# volledig door twee replicas worden opgeëist.
mem_limit: 4g
memswap_limit: 5g
cpus: 2.0
pids_limit: 512
# Leest de poort uit de eigen omgeving, dus dezelfde healthcheck werkt voor
# 3002 én 3003 zonder dat deze tweemaal in de compose hoeft te staan.
healthcheck:
test: ["CMD", "node", "-e", "fetch('http://127.0.0.1:'+(process.env.PORT||'3002')+'/api/health').then(r=>{process.exit(r.ok?0:1)}).catch(()=>process.exit(1))"]
interval: 15s
@@ -34,6 +58,7 @@ x-cms: &cms
start_period: 40s
services:
# Blauwe replica: host-poort 3002.
cms:
<<: *cms
container_name: epicnext-cms
@@ -41,6 +66,8 @@ services:
- HOSTNAME=0.0.0.0
- PORT=3002
# Groene replica: host-poort 3003. Meestal uitgeschakeld; alleen tijdens een
# release gestart, totdat nginx hem in de upstream-lijst heeft overgenomen.
cms-green:
<<: *cms
container_name: epicnext-cms-green
@@ -49,6 +76,7 @@ services:
- HOSTNAME=0.0.0.0
- PORT=3003
# ── Byparr (Cloudflare bypass for clone sources) ──
byparr:
image: ghcr.io/thephaseless/byparr:latest
container_name: byparr
@@ -56,10 +84,19 @@ services:
restart: unless-stopped
environment:
- LOG_LEVEL=INFO
# Resource limits verwijderd: Headless Chrome heeft bij zware pagina-scrapes
# soms tijdelijk meer dan 1 GB RAM nodig. Nu krijgt hij alle ruimte.
pids_limit: 256
healthcheck:
test: ["CMD", "curl", "http://localhost:8191/health"]
interval: 30s
timeout: 10s
retries: 3
start_period: 30s
# De database draait niet meer in Docker. `mariadb-turbo` is verwijderd: de
# service is nooit gestart, de volume bestond niet, en de echte MariaDB draait
# al als host-proces op 127.0.0.1:3306. De optimalisatie-vlaggen daar stonden
# dus al langer niets meer in beheer.
+18
View File
@@ -85,6 +85,24 @@ The direct template intentionally records the CDN/edge socket address when place
`pnpm test:integration` additionally starts disposable Nginx containers from the actual templates, supplies a temporary test certificate, and sends real HTTPS requests with forged identity headers. It checks direct-mode replacement even with an inherited real-IP rule, rejection of untrusted peers, and acceptance through an explicitly trusted peer. This requires Docker Engine and the OpenSSL CLI and does not read deployment credentials. The templates must still pass `nginx -t` on the intended host after its hostname/certificate substitution, then the listener and trusted-header checks above; the disposable fixture cannot certify that host or its firewall.
## Opt-in: CrowdSec on the same Docker host
A self-contained CrowdSec engine ships in `deployment/crowdsec`. It runs `crowdsecurity/crowdsec:v1.8.1` in its own Compose project in **LAPI-only mode** (`DISABLE_AGENT=true`) and exposes LAPI only on `127.0.0.1:18080`. No reverse-proxy, Traefik, Cloudflare or firewall configuration is changed.
```sh
bash cms security
```
The command generates `CROWDSEC_LAPI_API_KEY`, writes the CrowdSec flags into `.env`, starts the engine and registers the `cms` bouncer. The anti-DDoS gate then consults the local LAPI per client IP (short-cached) and blocks `ban`/`captcha` decisions before its own rate buckets. `bash cms security status` reports engine state and `bash cms security disable` stops the engine and flips the toggle off.
The engine does not enroll into the CrowdSec Central API (`DISABLE_ONLINE_API=true`) and runs without the agent (`DISABLE_AGENT=true`), so it needs no account and no outbound access to `crowdsec.net` (often blocked on hardened hosts). Blocking comes from the imported blocklists plus the app's own rate buckets; it does not parse the host Nginx log. The app still has its separate opt-in traffic-sharing channel via `CROWDSEC_REPORT_ENABLED`. Change `CROWDSEC_LAPI_PORT` and `CROWDSEC_LAPI_URL` together when `18080` is already in use. `CROWDSEC_NGINX_LOG_DIR` is honored for when the agent is re-enabled.
This bouncer is application-layer: it sheds known-bad IPs at the CMS process and only for traffic that reaches the Next.js proxy. It does not drop traffic before the origin, does not protect other host ports/services, and depends on the client IP being trustworthy at the ingress. Keep the upstream protections (Cloudflare IP rules, proxy rate limits) for defense before the origin.
The running CMS loads the new env values on its next restart or deployment. For a CI-managed `epicnext-cms-app`, the next deploy (which sources `.env`) applies them; for a clone, `bash cms update --skip-pull` restarts it. `.env` now holds the LAPI key — keep its permissions restrictive.
External IP blocklists (Spamhaus, DShield, CINS, blocklist.de, abuse.ch, IPsum, Firehol, Tor exit nodes, …) can be synced into the local LAPI with `bash cms security blocklists`, and hourly with `bash cms security blocklists-install-cron` (no account, but internet to fetch). Configure via `CROWDSEC_BLOCKLIST_*`.
## Routine and selected-release updates
```sh
+3 -5
View File
@@ -380,9 +380,7 @@ describe("Redis application cache", () => {
let fetches = 0;
const fetch = async () => ({ revision: ++fetches });
expect(await cache.cached(key, 60_000, fetch)).toEqual({ revision: 1 });
expect(
await cache.cached(key, 60000, async () => ({ revision: 1 })),
).resolves.toMatchObject({ revision: 1 });
expect(await appRedis?.get(key)).toBe('{"revision":1}');
expect(await appRedis?.ttl(key)).toBeGreaterThan(0);
cache.invalidateMemory(key);
expect(await cache.cached(key, 60_000, fetch)).toEqual({ revision: 1 });
@@ -622,7 +620,7 @@ describe("real news publication, scheduling and cache delivery", () => {
expect(await publicNews.getPublishedArticle(existing.slug)).toBeNull();
const negativeRevision = await appRedis?.get(NEWS_REVISION_KEY);
const negativeKey = `news:${negativeRevision}:article:v2:slug:${existing.slug}`;
expect(await appRedis?.get(negativeKey)).toBeNull();
expect(await appRedis?.get(negativeKey)).toBe("null");
expect(await appRedis?.ttl(negativeKey)).toBeGreaterThan(0);
const body = `<p>${"Contenuto completo è 📰 ".repeat(4000)}</p>`;
@@ -839,7 +837,7 @@ describe("real news publication, scheduling and cache delivery", () => {
expect(await publicNews.getPublishedArticle(existing.slug)).toBeNull();
const negativeRevision = await redis.get(NEWS_REVISION_KEY);
const negativeKey = `news:${negativeRevision}:article:v2:slug:${existing.slug}`;
expect(await redis.get(negativeKey)).toBeNull();
expect(await redis.get(negativeKey)).toBe("null");
const publish = articleForm({
id: String(existing.id),
baseToken: articleEditToken(existing),
+19 -19
View File
@@ -5,14 +5,14 @@
"engines": {
"node": ">=26.10.0 <27"
},
"packageManager": "pnpm@12.8.1",
"packageManager": "pnpm@12.6.0+sha512.3ef68f951cb111ac204b4a5a16f0b2ddf0da56a96e0413e81d855d9f0b55ef926714709028e1cd00c405c2c5fb7b9e8ec4dc46777c805d0373c2f2ff00fd20ec",
"scripts": {
"dev": "pnpm assets:editor && next dev",
"build": "pnpm assets:editor && NODE_OPTIONS='--max-old-space-size=4096' next build",
"build": "pnpm assets:editor && next build",
"start": "next start",
"toolchain:check": "node scripts/check-node-toolchain.mjs",
"lint": "biome check . || true",
"biome:lint": "biome check . || true",
"lint": "biome check .",
"biome:lint": "biome check .",
"format": "biome format --write .",
"diag:permissions": "tsx scripts/diagnose-permission-page.ts",
"jobs:worker": "node --conditions=react-server --import tsx scripts/jobs-worker.ts",
@@ -50,7 +50,7 @@
"@dnd-kit/utilities": "3.2.2",
"@formatjs/icu-messageformat-parser": "3.5.20",
"@hookform/resolvers": "5.9.1",
"@tanstack/react-query": "5.104.1",
"@tanstack/react-query": "5.104.0",
"@tanstack/react-virtual": "3.14.13",
"class-variance-authority": "0.7.1",
"clsx": "2.1.1",
@@ -63,20 +63,20 @@
"jpeg-js": "0.4.4",
"jsonc-parser": "3.3.1",
"jszip": "3.10.2",
"lucide-react": "1.50.0",
"lucide-react": "1.48.0",
"lzma-wasm": "1.0.7",
"motion": "14.0.0",
"motion": "13.4.4",
"music-metadata": "11.16.1",
"mysql2": "3.24.5",
"next": "16.3.8",
"mysql2": "3.24.4",
"next": "16.3.6",
"next-auth": "5.0.0-beta.32",
"next-intl": "4.14.9",
"next-intl": "4.14.7",
"otplib": "13.5.0",
"pino": "10.4.0",
"pino": "10.3.1",
"react": "19.3.0",
"react-dom": "19.3.0",
"react-hook-form": "7.89.0",
"resend": "6.32.0",
"resend": "6.30.0",
"server-only": "0.0.1",
"sharp": "^0.35.5",
"sonner": "2.0.8",
@@ -87,25 +87,25 @@
"devDependencies": {
"@axe-core/playwright": "4.13.0",
"@babel/parser": "7.29.9",
"@biomejs/biome": "2.5.15",
"@biomejs/biome": "2.5.14",
"@playwright/test": "1.63.0",
"@tailwindcss/forms": "0.5.11",
"@tailwindcss/postcss": "4.3.3",
"@tailwindcss/typography": "0.5.20",
"@types/node": "26.6.4",
"@types/node": "26.6.3",
"@types/react": "19.3.0",
"@types/react-dom": "19.3.0",
"@vitest/coverage-v8": "5.0.3",
"@vitest/coverage-v8": "5.0.2",
"drizzle-kit": "0.31.11",
"esbuild": "0.28.2",
"msw": "3.0.1",
"msw": "2.15.0",
"pino-pretty": "13.1.3",
"postcss": "8.5.28",
"tailwindcss": "4.3.3",
"testcontainers": "12.2.0",
"testcontainers": "12.1.0",
"tsx": "4.23.15",
"typescript": "7.0.2",
"vite": "8.3.2",
"vitest": "5.0.3"
"vite": "8.3.1",
"vitest": "5.0.2"
}
}
+359 -444
View File
File diff suppressed because it is too large. Load diff
+258
View File
@@ -0,0 +1,258 @@
#!/usr/bin/env bash
set -Eeuo pipefail
DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
cd "$DIR"
ENV_FILE="$DIR/.env"
COMPOSE_FILE="deployment/crowdsec/compose.crowdsec.yml"
PROJECT_NAME="epicnext-crowdsec"
CONTAINER_NAME="epicnext-crowdsec"
DEFAULT_DURATION="24h"
DEFAULT_MAX_DECISIONS="250000"
DEFAULT_SOURCES=(
# DDoS / abuse stoplists
"https://www.spamhaus.org/drop/drop.txt"
"https://www.spamhaus.org/drop/edrop.txt"
"https://www.dshield.org/block.txt"
"https://cinsscore.com/list/ci-badguys.txt"
"https://blocklist.greensnow.co/greensnow.txt"
"https://www.stopforumspam.com/downloads/toxic_ip_cidr.txt"
"https://www.binarydefense.com/banlist.txt"
# Brute force / credential stuffing
"https://lists.blocklist.de/lists/all.txt"
"https://lists.blocklist.de/lists/ssh.txt"
"https://lists.blocklist.de/lists/apache.txt"
"https://rules.emergingthreats.net/blockrules/compromised-ips.txt"
"https://danger.rulez.sk/projects/bruteforceblocker/blist.php"
# Malware C2 / botnets
"https://feodotracker.abuse.ch/downloads/ipblocklist.txt"
"https://sslbl.abuse.ch/blacklist/sslipblacklist.txt"
"https://www.botvrij.eu/data/ioclist.ip-dst.raw"
# Live SSH/spam attackers (last 48h), bare IPs only
"https://www.darklist.de/raw.php"
# Aggregated threat intel
"https://raw.githubusercontent.com/stamparm/ipsum/master/levels/3.txt"
"https://raw.githubusercontent.com/stamparm/ipsum/master/levels/2.txt"
# Firehol ipsets (security scanners, abusers, proxies, anonymous)
"https://raw.githubusercontent.com/firehol/blocklist-ipsets/master/firehol_level1.netset"
"https://raw.githubusercontent.com/firehol/blocklist-ipsets/master/firehol_level2.netset"
"https://raw.githubusercontent.com/firehol/blocklist-ipsets/master/firehol_abusers_1d.netset"
"https://raw.githubusercontent.com/firehol/blocklist-ipsets/master/firehol_abusers_30d.netset"
"https://raw.githubusercontent.com/firehol/blocklist-ipsets/master/firehol_proxies.netset"
"https://raw.githubusercontent.com/firehol/blocklist-ipsets/master/firehol_anonymous.netset"
"https://raw.githubusercontent.com/firehol/blocklist-ipsets/master/firehol_level3.netset"
# Tor exit nodes
"https://check.torproject.org/torbulkexitlist"
)
mode="${1:-sync}"
dry_run=false
case "$mode" in
sync) ;;
install-cron|uninstall-cron) ;;
*) printf 'ERROR: unknown mode "%s". Modes: sync [--dry-run] | install-cron | uninstall-cron\n' "$mode" >&2; exit 1 ;;
esac
[[ "${2:-}" = --dry-run ]] && dry_run=true
umask 077
fail() { printf 'ERROR: %s\n' "$*" >&2; exit 1; }
for command in docker curl flock; do command -v "$command" >/dev/null || fail "Required command: $command"; done
docker compose version >/dev/null 2>&1 || fail "Docker Compose plugin required."
exec 9>"$DIR/.deploy.lock"
flock -w 30 9 || fail "Another installation, update or sync is running."
[[ -f "$ENV_FILE" ]] || fail "Create .env first (bash cms install)."
env_get() {
local key="$1" line
while IFS= read -r line || [[ -n "$line" ]]; do
case "$line" in
"$key="*) line="${line#*=}"; line="${line%\"}"; line="${line#\"}"; printf '%s' "$line"; return 0 ;;
esac
done < "$ENV_FILE"
return 1
}
compose_cmd() {
docker compose --project-name "$PROJECT_NAME" --env-file "$ENV_FILE" -f "$COMPOSE_FILE" --profile security "$@"
}
container_running() {
[[ "$(docker inspect -f '{{.State.Running}}' "$CONTAINER_NAME" 2>/dev/null || true)" = true ]]
}
cscli_exec() {
compose_cmd exec -T crowdsec cscli "$@"
}
fetch_sources() {
local target="$1"
local sources=( "${DEFAULT_SOURCES[@]}" )
IFS=' ' read -r -a parsed <<< "${CROWDSEC_BLOCKLIST_SOURCES:-}"
[[ "${#parsed[@]}" -gt 0 ]] && sources=( "${parsed[@]}" )
local index=0 url
for url in "${sources[@]}"; do
[[ -n "$url" ]] || continue
index=$((index + 1))
if ! curl -fsSL -A "EpicNext-CMS blocklist sync" --retry 2 --max-time 90 -o "$target/source-$index.txt" "$url"; then
printf 'Warning: failed to fetch %s — continuing with the remaining sources.\n' "$url"
else
printf 'Fetched %s\n' "$url"
fi
done
}
install_cron() {
mkdir -p "$DIR/logs"
local cron_line="0 * * * * /usr/bin/env bash $DIR/scripts/blocklists-sync.sh >> $DIR/logs/blocklists-sync.log 2>&1"
if crontab -l 2>/dev/null | grep -Fq "$DIR/scripts/blocklists-sync.sh"; then
printf 'Cron entry already present:\n%s\n' "$cron_line"
else
( crontab -l 2>/dev/null; printf '%s\n' "$cron_line" ) | crontab -
printf 'Installed hourly cron entry:\n%s\n' "$cron_line"
fi
}
uninstall_cron() {
if crontab -l 2>/dev/null | grep -Fq "$DIR/scripts/blocklists-sync.sh"; then
crontab -l 2>/dev/null | grep -Fv "$DIR/scripts/blocklists-sync.sh" | crontab -
printf 'Removed cron entry matching %s.\n' "$DIR/scripts/blocklists-sync.sh"
else
printf 'No cron entry to remove.\n'
fi
}
if [[ "$mode" = install-cron ]]; then
install_cron
exit 0
fi
if [[ "$mode" = uninstall-cron ]]; then
uninstall_cron
exit 0
fi
duration="$(env_get CROWDSEC_BLOCKLIST_DURATION 2>/dev/null || true)"
[[ -n "$duration" ]] || duration="$DEFAULT_DURATION"
max_decisions="$(env_get CROWDSEC_BLOCKLIST_MAX_DECISIONS 2>/dev/null || true)"
[[ -n "$max_decisions" ]] || max_decisions="$DEFAULT_MAX_DECISIONS"
[[ "$max_decisions" =~ ^[0-9]+$ ]] || fail "CROWDSEC_BLOCKLIST_MAX_DECISIONS must be a number."
allowlist="${CROWDSEC_BLOCKLIST_ALLOW:-$(env_get CROWDSEC_BLOCKLIST_ALLOW 2>/dev/null || true)}"
work="$(mktemp -d "$DIR/.blocklists.XXXXXX")"
trap 'rm -rf -- "$work"' EXIT
build_lists() {
fetch_sources "$work"
cat "$work"/source-*.txt 2>/dev/null | awk '{print $1}' \
| grep -E '^([0-9]{1,3}\.){3}[0-9]{1,3}(/[0-9]{1,2})?$|^([0-9a-fA-F]{1,4}:){2,}[0-9a-fA-F:]*[0-9a-fA-F](/[0-9]{1,3})?$' \
| awk '
# Only globally routable attacker space may become a decision. Reserved,
# private, loopback, link-local, CGNAT, test and multicast ranges never
# represent an external attacker and must not be imported (they could
# otherwise block the origin itself or internal traffic).
function isReserved4(prefix, a, b, c) {
if (a == 0 || a == 127 || a >= 224) return 1
if (a == 10) return 1
if (a == 100 && (prefix < 10 || (prefix >= 10 && b >= 64 && b <= 127))) return 1
if (a == 169 && (prefix < 16 || (prefix >= 16 && b == 254))) return 1
if (a == 172 && (prefix < 12 || (prefix >= 12 && b >= 16 && b <= 31))) return 1
if (a == 192 && b == 168) return 1
if (a == 192 && b == 0) return 1
if ((a == 198 && (b == 18 || b == 19)) || (a == 198 && b == 51 && c == 100)) return 1
if (a == 203 && b == 0 && c == 113) return 1
return 0
}
function isReserved6(line, prefix, first, h) {
if (prefix < 32) return 1
if (line ~ /^::/) return 1
first = tolower(line); sub(/^::?/, "", first); sub(/:.*/, "", first)
h = "0x" substr(first, 1, 2)
if (h >= 252) return 1 # ULA fc00::/7, link-local fe80::/10, multicast ff00::/8
return 0
}
{
if (index($0, "/") > 0) {
n = split($0, seg, "/")
if (n != 2 || seg[2] !~ /^[0-9]+$/) next
if (index(seg[1], ":") > 0) {
pref = seg[2] + 0
if (pref < 32 || pref > 128) next
if (isReserved6(seg[1], pref)) next
print
next
}
pref = seg[2] + 0
if (pref < 8 || pref > 32) next
split(seg[1], oct, ".")
ok = 1
for (i = 1; i <= 4; i++) {
if (oct[i] !~ /^[0-9]+$/ || oct[i] + 0 > 255) { ok = 0; break }
if (length(oct[i]) > 1 && oct[i] ~ /^0/) { ok = 0; break }
}
if (!ok) next
if (isReserved4(pref, oct[1] + 0, oct[2] + 0, oct[3] + 0)) next
print
next
}
if (index($0, ":") > 0) {
if (isReserved6($0, 128)) next
print
next
}
n = split($0, part, ".")
if (n != 4) next
ok = 1
for (i = 1; i <= 4; i++) {
if (part[i] !~ /^[0-9]+$/ || part[i] + 0 > 255) { ok = 0; break }
if (length(part[i]) > 1 && part[i] ~ /^0/) { ok = 0; break }
}
if (!ok) next
if (isReserved4(32, part[1] + 0, part[2] + 0, part[3] + 0)) next
print
}' \
| sort -u > "$work/candidates.txt"
if [[ -n "$allowlist" ]]; then
printf '%s\n' "$allowlist" | tr ',' '\n' | while IFS= read -r line; do printf '%s\n' "$line"; done | sort -u > "$work/allow.txt"
comm -23 "$work/candidates.txt" "$work/allow.txt" > "$work/final.txt"
else
cp "$work/candidates.txt" "$work/final.txt"
fi
if [[ "$(wc -l < "$work/final.txt" | tr -d ' ')" -gt "$max_decisions" ]]; then
sort -u "$work/final.txt" | head -n "$max_decisions" > "$work/final.limited.txt" || true
mv "$work/final.limited.txt" "$work/final.txt"
printf 'Note: capped the combined list at %s decisions (CROWDSEC_BLOCKLIST_MAX_DECISIONS).\n' "$max_decisions"
fi
count_total=$(wc -l < "$work/final.txt" | tr -d ' ')
count_ip=$(grep -cv '/' "$work/final.txt" || true)
count_range=$(grep -c '/' "$work/final.txt" || true)
if [[ "$count_total" -lt 1 ]]; then
fail "No valid addresses could be parsed from the configured sources. Configure CROWDSEC_BLOCKLIST_SOURCES."
fi
}
build_lists
printf 'Parsed %s targets (%s IPs, %s ranges).\n' "$count_total" "$count_ip" "$count_range"
if $dry_run; then
printf 'Dry run: would replace the cscli-import decisions with these %s targets.\n' "$count_total"
exit 0
fi
container_running || fail "The CrowdSec engine is not running. Start it first with: bash cms security"
printf 'Removing previous cscli-import decisions...\n'
cscli_exec decisions delete --origin cscli-import >/dev/null 2>&1 || true
{
printf 'duration,scope,value\n'
awk -v d="$duration" '{ if (index($0, "/") > 0) printf "%s,range,%s\n", d, $0; else printf "%s,ip,%s\n", d, $0 }' "$work/final.txt"
} > "$work/import.csv"
printf 'Importing %s decisions into the local LAPI (duration %s)...\n' "$count_total" "$duration"
cscli_exec decisions import -i - --format csv --batch 1000 < "$work/import.csv"
printf 'Done. The app bouncer picks these up within a few seconds.\n'
+154
View File
@@ -0,0 +1,154 @@
#!/usr/bin/env bash
set -Eeuo pipefail
DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
cd "$DIR"
ENV_FILE="$DIR/.env"
COMPOSE_FILE="deployment/crowdsec/compose.crowdsec.yml"
PROJECT_NAME="epicnext-crowdsec"
CONTAINER_NAME="epicnext-crowdsec"
DEFAULT_PORT="18080"
mode="${1:-enable}"
case "$mode" in
enable|--enable) ;;
status|--status) ;;
disable|--disable) ;;
blocklists|blocklists-install-cron|blocklists-uninstall-cron) ;;
*) echo "Usage: bash cms security [enable|status|disable|blocklists|blocklists-install-cron|blocklists-uninstall-cron]" >&2; exit 1 ;;
esac
umask 077
fail() { printf 'ERROR: %s\n' "$*" >&2; exit 1; }
for command in docker flock; do command -v "$command" >/dev/null || fail "Required command: $command"; done
docker info >/dev/null 2>&1 || fail "Docker is not reachable."
docker compose version >/dev/null 2>&1 || fail "Docker Compose plugin required."
exec 9>"$DIR/.deploy.lock"
flock -w 30 9 || fail "Another installation or update is running."
[[ -f "$ENV_FILE" ]] || fail "Create .env first (bash cms install)."
case "$mode" in
blocklists) exec bash "$DIR/scripts/blocklists-sync.sh" sync "${2:-}" ;;
blocklists-install-cron) exec bash "$DIR/scripts/blocklists-sync.sh" install-cron ;;
blocklists-uninstall-cron) exec bash "$DIR/scripts/blocklists-sync.sh" uninstall-cron ;;
esac
env_get() {
local key="$1" line
while IFS= read -r line || [[ -n "$line" ]]; do
case "$line" in
"$key="*) line="${line#*=}"; line="${line%\"}"; line="${line#\"}"; printf '%s' "$line"; return 0 ;;
esac
done < "$ENV_FILE"
return 1
}
env_set() {
local key="$1" value="$2" tmp
tmp="$(mktemp "$DIR/.env.crowdsec.XXXXXX")"
if awk -v k="$key" -v v="$value" 'BEGIN{FS=OFS="=";done=0} { if ($1==k) { print k "=" v; done=1 } else print } END { if (!done) print k "=" v }' "$ENV_FILE" > "$tmp"; then
chmod 600 "$tmp"
mv -f -- "$tmp" "$ENV_FILE"
else
rm -f -- "$tmp"
fail "Could not update .env"
fi
}
compose_cmd() {
docker compose --project-name "$PROJECT_NAME" --env-file "$ENV_FILE" -f "$COMPOSE_FILE" --profile security "$@"
}
health_probe() {
local url="$1"
if command -v curl >/dev/null 2>&1; then
curl -fsS --max-time 3 "$url" >/dev/null 2>&1
else
compose_cmd exec -T crowdsec wget -q -O - "$url" >/dev/null 2>&1
fi
}
container_running() {
[[ "$(docker inspect -f '{{.State.Running}}' "$CONTAINER_NAME" 2>/dev/null || true)" = true ]]
}
if [[ "$mode" = disable || "$mode" = --disable ]]; then
set +e
compose_cmd stop crowdsec
rc=$?
set -e
[[ $rc -eq 0 ]] || printf 'CrowdSec engine was not running or could not be stopped.\n'
env_set CROWDSEC_LOCAL_ENABLED false
printf 'CrowdSec local stack disabled. The engine container is stopped; volumes and .env key were kept.\n'
exit 0
fi
if [[ "$mode" = status || "$mode" = --status ]]; then
enabled=no
[[ "$(env_get CROWDSEC_LOCAL_ENABLED 2>/dev/null || true)" = true ]] && enabled=yes
port="$(env_get CROWDSEC_LAPI_PORT 2>/dev/null || true)"
[[ -z "$port" ]] && port="$DEFAULT_PORT"
printf 'CROWDSEC_LOCAL_ENABLED=%s\n' "$enabled"
if container_running; then
printf 'Engine: running\n'
if health_probe "http://127.0.0.1:$port/health"; then
printf 'LAPI health: OK (127.0.0.1:%s)\n' "$port"
else
printf 'LAPI health: UNREACHABLE (127.0.0.1:%s)\n' "$port"
fi
else
printf 'Engine: not running\n'
printf 'Start with: bash cms security\n'
fi
exit 0
fi
port="$(env_get CROWDSEC_LAPI_PORT 2>/dev/null || true)"
[[ -n "$port" ]] || port="$DEFAULT_PORT"
[[ "$port" =~ ^[0-9]{1,5}$ ]] || fail "CROWDSEC_LAPI_PORT must be a port number."
if (( port < 1024 || port > 65535 )); then
fail "CROWDSEC_LAPI_PORT must be within 1024-65535."
fi
key="$(env_get CROWDSEC_LAPI_API_KEY 2>/dev/null || true)"
[[ -n "$key" ]] || key="$(od -An -N32 -tx1 /dev/urandom | tr -d ' \n')"
url="$(env_get CROWDSEC_LAPI_URL 2>/dev/null || true)"
[[ -n "$url" ]] || url="http://127.0.0.1:$port"
log_dir="${CROWDSEC_NGINX_LOG_DIR:-$(env_get CROWDSEC_NGINX_LOG_DIR 2>/dev/null || true)}"
[[ -n "$log_dir" ]] || log_dir="/var/log/nginx"
if ! container_running && command -v ss >/dev/null 2>&1; then
if ss -ltn "( sport = :$port )" 2>/dev/null | grep -q LISTEN; then
fail "Port $port is already in use. Set CROWDSEC_LAPI_PORT (and CROWDSEC_LAPI_URL) in .env to a free port and re-run."
fi
fi
if [[ ! -r "$log_dir/access.log" ]]; then
printf 'Warning: %s/access.log is not readable. The engine will run but has no detections until an access log is available.\n' "$log_dir"
fi
env_set CROWDSEC_LOCAL_ENABLED true
env_set CROWDSEC_LAPI_URL "$url"
env_set CROWDSEC_LAPI_PORT "$port"
env_set CROWDSEC_LAPI_API_KEY "$key"
env_set CROWDSEC_NGINX_LOG_DIR "$log_dir"
compose_cmd config --quiet || fail "CrowdSec Compose configuration is invalid; fix CROWDSEC_* settings in .env."
set +e
compose_cmd up -d --wait crowdsec
rc=$?
set -e
if [[ $rc -ne 0 ]]; then
compose_cmd up -d crowdsec
fi
attempt=0
while ! health_probe "http://127.0.0.1:$port/health"; do
attempt=$((attempt + 1))
[[ $attempt -lt 30 ]] || fail "CrowdSec LAPI did not become healthy on port $port."
sleep 2
done
printf 'CrowdSec engine running in LAPI-only mode on 127.0.0.1:%s (container %s).\n' "$port" "$CONTAINER_NAME"
compose_cmd exec -T crowdsec cscli bouncers list >/dev/null 2>&1 \
&& printf 'Bouncer "cms" was registered against the local LAPI.\n' \
|| printf 'Warning: could not list bouncers. Diagnose with: docker compose exec -T %s cscli bouncers list\n' "$CONTAINER_NAME"
printf 'Restart the CMS container (or run your next deployment) so it loads the new bouncer env. For a clone: bash cms update --skip-pull\n'
+26
View File
@@ -15,3 +15,29 @@ it("imports runtime validation without starting the CMS", () => {
);
assert.equal(result.status, 0, result.stderr);
});
it("rejects the local CrowdSec bouncer without a key", () => {
const result = spawnSync(
process.execPath,
[
"--input-type=module",
"-e",
"import {validateRuntime} from './scripts/docker-start.mjs';try{validateRuntime({HOTEL_NAME:'x',AUTH_SECRET:'01234567890123456789012345678901',DATABASE_URL:'mysql://u:p@h/db',APP_URL:'http://h',CROWDSEC_LOCAL_ENABLED:'true'});process.exit(1)}catch(error){if(!String(error.message).includes('CROWDSEC_LAPI_API_KEY'))throw error}",
],
{ encoding: "utf8" },
);
assert.equal(result.status, 0, result.stderr);
});
it("accepts a complete local CrowdSec configuration", () => {
const result = spawnSync(
process.execPath,
[
"--input-type=module",
"-e",
"import {validateRuntime} from './scripts/docker-start.mjs';validateRuntime({HOTEL_NAME:'x',AUTH_SECRET:'01234567890123456789012345678901',DATABASE_URL:'mysql://u:p@h/db',APP_URL:'http://h',CROWDSEC_LOCAL_ENABLED:'true',CROWDSEC_LAPI_API_KEY:'fixture-key',CROWDSEC_LAPI_URL:'http://127.0.0.1:18080'})",
],
{ encoding: "utf8" },
);
assert.equal(result.status, 0, result.stderr);
});
+16
View File
@@ -23,6 +23,22 @@ export function validateRuntime(settings) {
invalid.push(key);
}
}
const localEnabled = ["true", "1"].includes(
String(settings.CROWDSEC_LOCAL_ENABLED ?? "")
.trim()
.toLowerCase(),
);
if (localEnabled) {
if (!settings.CROWDSEC_LAPI_API_KEY?.trim())
invalid.push("CROWDSEC_LAPI_API_KEY");
if (settings.CROWDSEC_LAPI_URL) {
try {
new URL(settings.CROWDSEC_LAPI_URL);
} catch {
invalid.push("CROWDSEC_LAPI_URL");
}
}
}
if (invalid.length)
throw new Error(`Invalid runtime configuration: ${invalid.join(", ")}`);
}
+91
View File
@@ -3,6 +3,7 @@ import {
copyFileSync,
mkdirSync,
mkdtempSync,
readFileSync,
rmSync,
writeFileSync,
} from "node:fs";
@@ -84,3 +85,93 @@ it.skipIf(!hasCompose)(
},
30_000,
);
it("documents the local CrowdSec switches in .env.example", () => {
const examples = readFileSync(path.join(root, ".env.example"), "utf8");
for (const key of [
"CROWDSEC_LOCAL_ENABLED",
"CROWDSEC_LAPI_URL",
"CROWDSEC_LAPI_PORT",
"CROWDSEC_LAPI_API_KEY",
"CROWDSEC_NGINX_LOG_DIR",
]) {
expect(examples).toContain(key);
}
});
it.skipIf(!hasCompose)(
"renders the standalone CrowdSec stack with a loopback-only LAPI",
() => {
const directory = mkdtempSync(path.join(tmpdir(), "cms-crowdsec-"));
try {
mkdirSync(path.join(directory, "deployment/crowdsec/acquis.d"), {
recursive: true,
});
copyFileSync(
path.join(root, "deployment/crowdsec/compose.crowdsec.yml"),
path.join(directory, "deployment/crowdsec/compose.crowdsec.yml"),
);
copyFileSync(
path.join(root, "deployment/crowdsec/acquis.d/nginx.yaml"),
path.join(directory, "deployment/crowdsec/acquis.d/nginx.yaml"),
);
writeFileSync(
path.join(directory, ".env"),
[
"CROWDSEC_LAPI_API_KEY=fixture-key",
"CROWDSEC_LAPI_PORT=18080",
"CROWDSEC_LAPI_URL=http://127.0.0.1:18080",
"CROWDSEC_NGINX_LOG_DIR=/var/log/nginx",
].join("\n"),
);
const environment = { ...process.env };
for (const key of Object.keys(environment))
if (
key.startsWith("COMPOSE_") ||
key.startsWith("CROWDSEC_") ||
key.startsWith("TZ")
)
delete environment[key];
const result = spawnSync(
"docker",
[
"compose",
"--project-name",
"crowdsec-fixture",
"--env-file",
".env",
"-f",
"deployment/crowdsec/compose.crowdsec.yml",
"--profile",
"security",
"config",
"--format",
"json",
],
{ cwd: directory, env: environment, encoding: "utf8", timeout: 15_000 },
);
expect(result.status, result.stderr).toBe(0);
const config = JSON.parse(result.stdout);
const service = config.services.crowdsec;
expect(service).toBeDefined();
expect(service.image).toContain("crowdsecurity/crowdsec:");
expect(service.environment.BOUNCER_KEY_cms).toBe("fixture-key");
expect(service.environment.DISABLE_ONLINE_API).toBe("true");
expect(
service.ports.some(
(published) =>
published.host_ip === "127.0.0.1" &&
published.published === "18080" &&
published.target === 8080,
),
).toBe(true);
const targets = service.volumes.map((volume) => volume.target);
expect(targets).toContain("/var/log/nginx");
expect(targets).toContain("/etc/crowdsec/acquis.d");
expect(service.healthcheck.test.join(" ")).toContain("wget");
} finally {
rmSync(directory, { recursive: true, force: true });
}
},
30_000,
);
+60
View File
@@ -15,6 +15,14 @@ import {
setLastCloudflareVerify,
verifyCloudflareConnection,
} from "@/lib/cloudflare-api";
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";
@@ -33,6 +41,17 @@ function positiveInt(raw: FormDataEntryValue | null, fallback: number): number {
return Math.floor(n);
}
function clampInt(
raw: FormDataEntryValue | null,
fallback: number,
min: number,
max: number,
): number {
const n = Number(str(raw));
if (!Number.isFinite(n)) return fallback;
return Math.min(max, Math.max(min, Math.floor(n)));
}
function parseTiers(raw: FormDataEntryValue | null): AntiddosBlockTier[] {
const tiers: AntiddosBlockTier[] = [];
for (const part of str(raw).split(",")) {
@@ -99,6 +118,17 @@ function configFromForm(formData: FormData): AntiddosConfig {
defaults.globalHaltMs,
),
cloudflareAutoBlock: str(formData.get("cfa_auto_block")) === "1",
crowdsecAutoBlock: str(formData.get("cs_auto_block")) === "1",
crowdsecBlockScore: clampInt(
formData.get("cs_block_score"),
defaults.crowdsecBlockScore,
0,
5,
),
crowdsecBlockTtlSeconds: positiveInt(
formData.get("cs_block_ttl_sec"),
defaults.crowdsecBlockTtlSeconds,
),
};
}
@@ -123,6 +153,9 @@ async function persistSettings(config: AntiddosConfig): Promise<void> {
],
["antiddos_global_halt_ms", String(config.globalHaltMs)],
["antiddos_cfa_auto_block", config.cloudflareAutoBlock ? "1" : "0"],
["antiddos_cs_auto_block", config.crowdsecAutoBlock ? "1" : "0"],
["antiddos_cs_block_score", String(config.crowdsecBlockScore)],
["antiddos_cs_block_ttl", String(config.crowdsecBlockTtlSeconds)],
];
await Promise.all(
entries.map(([key, value]) =>
@@ -195,6 +228,7 @@ export async function unbanAntiddosIp(formData: FormData): Promise<void> {
await Promise.all([
redis.del(`antiddos:block:${ip}`),
redis.del(`antiddos:block:meta:${ip}`),
redis.del(`crowdsec:report:${ip}`),
redis.del(`antiddos:v:${ip}`),
]);
}
@@ -231,6 +265,32 @@ export async function removeCloudflareRule(formData: FormData): Promise<void> {
revalidatePath("/admin/devops/antiddos");
}
/** Test the configured CrowdSec API credentials against the CTI endpoint. */
export async function verifyCrowdsecConfiguration(): Promise<void> {
const staff = await requirePermission(PERMS.SETTINGS_VIEW);
const status = await verifyCrowdsecConnection();
await setLastCrowdsecVerify(status);
logger.info("CrowdSec API configuration verified", {
staff: staff.username,
ok: status.ok,
message: status.message,
});
revalidatePath("/admin/devops/antiddos");
}
/** Test the CrowdSec signal-push (CAPI watcher) channel. */
export async function verifyCrowdsecReportingConfiguration(): Promise<void> {
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<void> {
const staff = await requirePermission(PERMS.SETTINGS_VIEW);
+396 -4
View File
@@ -1,4 +1,11 @@
import { BadgeCheck, Cloud, Lock, Server, ShieldAlert } from "lucide-react";
import {
BadgeCheck,
Cloud,
Lock,
Radar,
Server,
ShieldAlert,
} from "lucide-react";
import { headers } from "next/headers";
import { redirect } from "next/navigation";
import {
@@ -7,6 +14,8 @@ import {
saveAntiddosSettings,
unbanAntiddosIp,
verifyCloudflareConfiguration,
verifyCrowdsecConfiguration,
verifyCrowdsecReportingConfiguration,
} from "@/actions/admin-antiddos";
import { Badge } from "@/components/ui/badge";
import { Button } from "@/components/ui/button";
@@ -24,6 +33,19 @@ import {
listCloudflareBlocks,
sweepExpiredCloudflareBlocks,
} 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 { type CrowdsecDailyStat, getCrowdsecStats } from "@/lib/crowdsec-stats";
import { db, WebsiteSetting } from "@/lib/db";
import { canAccess, getAdminContext, PERMS } from "@/lib/permissions";
import { redis } from "@/lib/redis";
@@ -36,6 +58,40 @@ function seconds(ttlMs: number): string {
return `${Math.floor(s / 3600)}h ${Math.floor((s % 3600) / 60)}m`;
}
function BarSparkline({ values }: { values: number[] }) {
if (values.length === 0) return null;
const max = Math.max(...values, 1);
return (
<div className="flex items-end gap-[3px] h-10" aria-hidden="true">
{values.map((v, i) => (
<div
// biome-ignore lint/suspicious/noArrayIndexKey: static timeline position is the bar's identity
key={i}
className="w-full rounded-sm bg-primary/60"
style={{
height: `${Math.max(v > 0 ? 6 : 2, (v / max) * 100)}%`,
opacity: v === 0 ? 0.15 : 0.6 + (v / max) * 0.4,
}}
/>
))}
</div>
);
}
/** Merge per-day breakdown maps (categories / reputations) into range totals. */
function mergeBreakdowns(
rows: CrowdsecDailyStat[],
kind: keyof Pick<CrowdsecDailyStat, "categories" | "reputations">,
): Record<string, number> {
const totals: Record<string, number> = {};
for (const row of rows) {
for (const [k, v] of Object.entries(row[kind])) {
totals[k] = (totals[k] ?? 0) + v;
}
}
return totals;
}
export default async function AdminAntiDdosPage() {
const { session, permissions } = await getAdminContext();
if (!canAccess(permissions, PERMS.SETTINGS_VIEW, session.user.rank)) {
@@ -62,7 +118,8 @@ export default async function AdminAntiDdosPage() {
ip: string;
ttlMs: number;
count: number;
source: "gate";
source: "gate" | "crowdsec";
meta: CrowdsecBlockMeta | null;
}[] = [];
let redisOk = false;
const rateStore = redis;
@@ -83,13 +140,23 @@ export default async function AdminAntiDdosPage() {
}
const withTtl = await Promise.all(
blockKeys.slice(0, 100).map(async (key) => {
const ttlMs = await rateStore.pttl(key);
const [ttlMs, value] = await Promise.all([
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,
ttlMs: ttlMs > 0 ? ttlMs : 0,
count: violationCounts.get(ip) ?? 0,
source: "gate" as const,
// The gate writes "1"; "crowdsec" marks a community-reputation block.
source,
// Why CrowdSec blocked this IP, when the meta was recorded.
meta: source === "crowdsec" ? await getCrowdsecBlockMeta(ip) : null,
};
}),
);
@@ -110,6 +177,12 @@ export default async function AdminAntiDdosPage() {
cloudflareBlocks.push(...(await listCloudflareBlocks()));
}
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();
const crowdsecStats = redisOk ? await getCrowdsecStats(14) : [];
return (
<div className="space-y-6">
@@ -164,6 +237,23 @@ export default async function AdminAntiDdosPage() {
</CardContent>
</Card>
<Card>
<CardHeader className="flex flex-row items-center justify-between space-y-0 pb-2">
<CardTitle className="text-sm font-medium">CrowdSec</CardTitle>
<Radar className="h-4 w-4 text-muted-foreground" />
</CardHeader>
<CardContent>
<Badge variant={crowdsecConfigured ? "default" : "secondary"}>
{crowdsecConfigured ? "Connected" : "Not configured"}
</Badge>
<p className="text-xs text-muted-foreground mt-1">
{crowdsecUsage && crowdsecUsage.quota > 0
? `${crowdsecUsage.used.toLocaleString()} / ${crowdsecUsage.quota.toLocaleString()} CTI calls today${crowdsecUsage.exhausted ? " (paused)" : ""}`
: "Community reputation auto-block"}
</p>
</CardContent>
</Card>
<Card>
<CardHeader className="flex flex-row items-center justify-between space-y-0 pb-2">
<CardTitle className="text-sm font-medium">Active blocks</CardTitle>
@@ -261,6 +351,53 @@ export default async function AdminAntiDdosPage() {
block.
</p>
<label className="flex items-center gap-2 text-sm">
<input
type="checkbox"
name="cs_auto_block"
value="1"
defaultChecked={effective.crowdsecAutoBlock}
/>
Automatically block IPs flagged as malicious by the CrowdSec
community
</label>
<p className="text-xs text-muted-foreground -mt-2">
Requires <span className="font-mono">CROWDSEC_API_KEY</span> in
the environment. When a repeat offender has a bad community
reputation it is hard-blocked immediately (no need to cross the
local violation threshold). IPs carrying CrowdSec false-positive
tags are never blocked.
</p>
<div className="flex flex-wrap items-center gap-4">
<label className="block">
<span className="text-xs font-medium">
Minimum reputation score (0–5)
</span>
<input
name="cs_block_score"
type="number"
min={0}
max={5}
defaultValue={effective.crowdsecBlockScore}
className="w-24 mt-1"
/>
<span className="text-xs text-muted-foreground ml-2">
4–5 = malicious (CrowdSec scale)
</span>
</label>
<label className="block">
<span className="text-xs font-medium">
Block duration (sec)
</span>
<input
name="cs_block_ttl_sec"
type="number"
defaultValue={effective.crowdsecBlockTtlSeconds}
className="w-32 mt-1"
/>
</label>
</div>
<div className="grid grid-cols-1 gap-4 md:grid-cols-3">
{(
[
@@ -403,6 +540,10 @@ export default async function AdminAntiDdosPage() {
) : (
<div className="space-y-2">
{blocks.map((b) => {
const behaviorLabel =
b.meta && b.meta.behaviors.length > 0
? b.meta.behaviors.join(", ")
: null;
return (
<div
key={b.ip}
@@ -410,8 +551,22 @@ export default async function AdminAntiDdosPage() {
>
<span className="font-mono">{b.ip}</span>
<span className="flex items-center gap-2 text-xs text-muted-foreground">
{b.source === "crowdsec" ? (
<Badge variant="default">CrowdSec</Badge>
) : (
<Badge variant="secondary">Gate</Badge>
)}
TTL {seconds(b.ttlMs)} · violations {b.count}
{b.meta && (
<span
className="max-w-xs truncate"
title={`${b.meta.reputation ?? "unknown"} · score ${b.meta.score} · ${b.meta.category}${behaviorLabel ? ` · ${behaviorLabel}` : ""}`}
>
{b.meta.reputation ?? "unknown"} · score{" "}
{b.meta.score} · {b.meta.category}
{behaviorLabel ? ` · ${behaviorLabel}` : ""}
</span>
)}
</span>
<form action={unbanAntiddosIp}>
<input type="hidden" name="ip" value={b.ip} />
@@ -505,6 +660,243 @@ export default async function AdminAntiDdosPage() {
</CardContent>
</Card>
<Card>
<CardHeader>
<CardTitle className="flex items-center gap-2">
<Radar className="h-4 w-4" /> CrowdSec reputation API
</CardTitle>
</CardHeader>
<CardContent className="space-y-4">
<div className="flex flex-wrap items-center gap-3">
<Badge variant={crowdsecConfigured ? "default" : "secondary"}>
{crowdsecConfigured ? "API configured" : "API not configured"}
</Badge>
{!crowdsecConfigured && (
<p className="text-xs text-muted-foreground">
Set <span className="font-mono">CROWDSEC_API_KEY</span> to
enable community-reputation auto-blocks. When a repeat offender
is flagged as malicious by the CrowdSec community it is
hard-blocked immediately without waiting for the local violation
threshold.
</p>
)}
<form action={verifyCrowdsecConfiguration}>
<Button
type="submit"
size="sm"
variant="outline"
disabled={!crowdsecConfigured}
>
Verify connection
</Button>
</form>
</div>
{lastCrowdsecVerify && crowdsecConfigured && (
<p className="text-xs">
<Badge
variant={lastCrowdsecVerify.ok ? "default" : "destructive"}
>
{lastCrowdsecVerify.ok ? "Reachable" : "Failed"}
</Badge>
<span className="ml-2 text-muted-foreground">
{lastCrowdsecVerify.ok
? `CTI endpoint verified ${new Date(lastCrowdsecVerify.at).toLocaleString()}`
: lastCrowdsecVerify.message}
</span>
</p>
)}
{crowdsecConfigured && (
<p className="text-sm text-muted-foreground">
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{" "}
<Badge variant="default">CrowdSec</Badge> badge in the list above
with the community reasoning (reputation, score, behaviors).
</p>
)}
{crowdsecUsage && (
<div className="rounded-md border p-3">
<p className="text-xs font-medium mb-1">
Reputation lookups today
</p>
{crowdsecUsage.quota > 0 ? (
<>
<div className="flex items-center gap-2">
<div className="h-2 flex-1 overflow-hidden rounded-full bg-muted">
<div
className="h-full rounded-full"
style={{
width: `${Math.min(100, (crowdsecUsage.used / crowdsecUsage.quota) * 100)}%`,
background: crowdsecUsage.exhausted
? "var(--color-destructive)"
: crowdsecUsage.used >= crowdsecUsage.quota * 0.8
? "var(--admin-accent)"
: "var(--color-primary)",
}}
/>
</div>
<span
className={`text-xs ${crowdsecUsage.exhausted ? "text-destructive" : "text-muted-foreground"}`}
>
{crowdsecUsage.used.toLocaleString()} /{" "}
{crowdsecUsage.quota.toLocaleString()}
</span>
</div>
<p className="text-xs text-muted-foreground mt-1">
{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."}
</p>
</>
) : (
<p className="text-xs text-muted-foreground">
Tracking disabled (CROWDSEC_CTI_DAILY_QUOTA = 0 / unlimited).
</p>
)}
</div>
)}
{crowdsecStats.length > 0 && (
<div className="rounded-md border p-3">
<div className="flex flex-wrap items-center justify-between gap-2">
<p className="text-xs font-medium">
Daily activity (last {crowdsecStats.length} days)
</p>
<a
href="/admin/alerts"
className="text-xs text-muted-foreground underline-offset-2 hover:underline"
>
Ops alert history →
</a>
</div>
<div className="mt-2 grid gap-4 sm:grid-cols-3">
<BarSparkline values={crowdsecStats.map((row) => row.blocks)} />
<div className="col-span-2 flex flex-wrap items-center gap-1.5">
{Object.entries(
mergeBreakdowns(crowdsecStats, "categories"),
).map(([category, count]) => (
<Badge key={category} variant="secondary">
{category} · {count.toLocaleString()}
</Badge>
))}
{Object.entries(
mergeBreakdowns(crowdsecStats, "reputations"),
).map(([reputation, count]) => (
<Badge
key={reputation}
variant={
reputation === "malicious" ? "destructive" : "secondary"
}
>
{reputation} · {count.toLocaleString()}
</Badge>
))}
</div>
</div>
<div className="mt-3 max-h-40 overflow-y-auto">
<table className="w-full text-xs">
<thead>
<tr className="text-left text-muted-foreground">
<th className="pb-1 pr-2 font-medium">Date</th>
<th className="pb-1 pr-2 font-medium text-right">
Lookups
</th>
<th className="pb-1 pr-2 font-medium text-right">
Blocks
</th>
<th className="pb-1 pr-2 font-medium text-right">
Reports
</th>
<th className="pb-1 font-medium text-right">Failures</th>
</tr>
</thead>
<tbody>
{crowdsecStats.map((row) => (
<tr key={row.date} className="border-t">
<td className="py-1 pr-2 text-muted-foreground">
{row.date === new Date().toISOString().slice(0, 10)
? "Today"
: row.date.slice(5)}
</td>
<td className="py-1 pr-2 text-right">
{row.lookups.toLocaleString()}
</td>
<td className="py-1 pr-2 text-right">
{row.blocks.toLocaleString()}
</td>
<td className="py-1 pr-2 text-right">
{row.reports.toLocaleString()}
</td>
<td className="py-1 text-right">
{row.reportFailures > 0 ? (
<span className="text-destructive">
{row.reportFailures.toLocaleString()}
</span>
) : (
"–"
)}
</td>
</tr>
))}
</tbody>
</table>
</div>
</div>
)}
<div className="rounded-md border p-3">
<p className="text-xs font-medium mb-1">Community signal push</p>
<div className="flex flex-wrap items-center gap-3">
<Badge variant={reportingEnabled ? "default" : "secondary"}>
{reportingEnabled ? "Enabled" : "Off"}
</Badge>
{!reportingEnabled && (
<p className="text-xs text-muted-foreground">
Set{" "}
<span className="font-mono">
CROWDSEC_REPORT_ENABLED=true
</span>{" "}
to share blocked IPs back into the CrowdSec community
blocklist. Watcher credentials are auto-generated and
persisted in Redis.
</p>
)}
{reportingEnabled && (
<p className="text-xs text-muted-foreground">
Blocked IPs are pushed to the Central API (deduped per IP) so
the community blocklist protects other members too.
</p>
)}
<form action={verifyCrowdsecReportingConfiguration}>
<Button
type="submit"
size="sm"
variant="outline"
disabled={!reportingEnabled}
>
Verify channel
</Button>
</form>
</div>
{lastReport && (
<p className="text-xs mt-2">
<Badge variant={lastReport.ok ? "default" : "destructive"}>
{lastReport.ok ? "Push healthy" : "Push failed"}
</Badge>
<span className="ml-2 text-muted-foreground">
{lastReport.ok
? `Last signal accepted ${new Date(lastReport.at).toLocaleString()}`
: `${lastReport.message ?? "unknown"} (${new Date(lastReport.at).toLocaleString()})`}
</span>
</p>
)}
</div>
</CardContent>
</Card>
{stored.size === 0 && (
<p className="text-xs text-muted-foreground">
Persisted site settings: none yet — the form values above reflect the
@@ -1,47 +1,266 @@
import { describe, it, expect, vi, beforeAll, afterAll, afterEach } from "vitest";
import { onlineManager, QueryObserver } from "@tanstack/react-query";
import { delay, HttpResponse, http } from "msw";
import { setupServer } from "msw/node";
import { http, HttpResponse } from "msw";
import { QueryClient } from "@tanstack/react-query";
import { furnitureJobsQueryOptions } from "./furniture-jobs-query";
import {
afterAll,
afterEach,
beforeAll,
describe,
expect,
it,
vi,
} from "vitest";
import {
furnitureCursorFixture,
furnitureJobFixture,
} from "@/test/furniture-jobs-fixtures";
import {
createFurnitureJobsClient,
furnitureJobsQueryOptions,
readFurnitureJobResponse,
} from "./furniture-jobs-query";
let calls = 0;
let release: (() => void) | null = null;
// Cast to any to bypass strict MSW v3 config types in the test environment
const server = setupServer(
http.get("https://localhost/api/admin/studio/furniture/jobs", async () => {
calls++;
if (release) await new Promise<void>((resolve) => { release = resolve; });
return HttpResponse.json({ jobs: [], nextCursor: null });
})
);
beforeAll(() => server.listen({ onUnhandledRequest: "bypass" } as any));
afterEach(() => {
calls = 0;
release = null;
const endpoint = "http://localhost/api/admin/studio/import-jobs";
const server = setupServer();
const clients: ReturnType<typeof createFurnitureJobsClient>[] = [];
function client() {
const value = createFurnitureJobsClient();
clients.push(value);
return value;
}
beforeAll(() => server.listen({ onUnhandledRequest: "error" }));
afterEach(async () => {
await Promise.all(
clients.map(async (value) => {
await value.cancelQueries();
value.clear();
}),
);
clients.length = 0;
server.resetHandlers();
});
afterAll(() => {
try { server.close(); } catch (e) {}
});
afterAll(() => server.close());
describe("furniture history HTTP query", () => {
const client = () => new QueryClient({ defaultOptions: { queries: { retry: false } } });
it("deduplicates simultaneous refreshes and caches their validated result", async () => {
let calls = 0;
let release: (() => void) | undefined;
server.use(
http.get(endpoint, async () => {
calls++;
await new Promise<void>((resolve) => {
release = resolve;
});
return HttpResponse.json({
ok: true,
jobs: [furnitureJobFixture],
nextCursor: furnitureCursorFixture,
});
}),
);
const cache = client();
// Use type assertions to allow the test to manipulate options freely
const options = furnitureJobsQueryOptions(false, null) as any;
options.queryKey = ["furniture-import-history", false, null];
options.queryFn = () => fetch("https://localhost/api/admin/studio/furniture/jobs").then(r => r.json());
const options = furnitureJobsQueryOptions(true, null);
const first = cache.fetchQuery(options);
const second = cache.fetchQuery(options);
await vi.waitFor(() => expect(calls).toBe(1));
if (release) release();
await Promise.all([first, second]);
release?.();
const [a, b] = await Promise.all([first, second]);
expect(a).toEqual(b);
expect(a.jobs[0].id).toBe(furnitureJobFixture.id);
expect(cache.getQueryData(options.queryKey)).toEqual(a);
});
it("keeps different pages and mounted history clients isolated", async () => {
const urls: string[] = [];
server.use(
http.get(endpoint, ({ request }) => {
urls.push(request.url);
return HttpResponse.json({ ok: true, jobs: [], nextCursor: null });
}),
);
const first = client(),
second = client();
await first.fetchQuery(furnitureJobsQueryOptions(true, null));
await first.fetchQuery(
furnitureJobsQueryOptions(true, furnitureCursorFixture),
);
await second.fetchQuery(furnitureJobsQueryOptions(true, null));
expect(urls).toHaveLength(3);
expect(new URL(urls[1]).searchParams.get("before")).toBe(
furnitureCursorFixture,
);
first.clear();
expect(
second.getQueryData(furnitureJobsQueryOptions(true, null).queryKey),
).toEqual({ jobs: [], nextCursor: null });
});
it.each([403, 500])(
"surfaces HTTP %i without automatic retries or a false empty result",
async (status) => {
let calls = 0;
server.use(
http.get(endpoint, () => {
calls++;
return HttpResponse.json(
{ error: "History unavailable" },
{ status },
);
}),
);
const cache = client(),
options = furnitureJobsQueryOptions(false, null);
await expect(cache.fetchQuery(options)).rejects.toMatchObject({
status,
message: "History unavailable",
});
expect(calls).toBe(1);
expect(cache.getQueryData(options.queryKey)).toBeUndefined();
},
);
it.each([
{ ok: true, jobs: null },
{ ok: true, jobs: [{ ...furnitureJobFixture, state: "invented" }] },
{
ok: true,
jobs: [
{
...furnitureJobFixture,
items: [{ ...furnitureJobFixture.items[0], state: "unknown" }],
},
],
},
{ ok: true, jobs: [{ ...furnitureJobFixture, createdAt: "not a date" }] },
{ ok: true, jobs: [], nextCursor: "../other-account" },
{
ok: true,
jobs: [],
nextCursor:
"2030-99-99T25:61:61.000Z|aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa",
},
])("rejects malformed successful payloads", async (payload) => {
server.use(http.get(endpoint, () => HttpResponse.json(payload)));
const cache = client(),
options = furnitureJobsQueryOptions(false, null);
await expect(cache.fetchQuery(options)).rejects.toThrow(
"Invalid import history response",
);
expect(cache.getQueryData(options.queryKey)).toBeUndefined();
});
it("reports invalid JSON as a transport failure", async () => {
server.use(
http.get(
endpoint,
() => new HttpResponse("<html>unexpected</html>", { status: 200 }),
),
);
await expect(
client().fetchQuery(furnitureJobsQueryOptions(false, null)),
).rejects.toThrow("Invalid import history response");
});
it("times out a stalled HTTP request without retries", async () => {
let calls = 0;
server.use(
http.get(endpoint, async () => {
calls++;
await delay(100);
return HttpResponse.json({ ok: true, jobs: [] });
}),
);
await expect(
client().fetchQuery(furnitureJobsQueryOptions(false, null, 10)),
).rejects.toThrow("Import history request timed out");
expect(calls).toBe(1);
});
it("aborts a cancelled request without replacing cache data with an empty result", async () => {
let entered = false;
let transportAborted = false;
server.use(
http.get(endpoint, async ({ request }) => {
request.signal.addEventListener("abort", () => {
transportAborted = true;
});
entered = true;
await delay(100);
return HttpResponse.json({ ok: true, jobs: [] });
}),
);
const cache = client(),
options = furnitureJobsQueryOptions(false, null);
const pending = cache.fetchQuery(options);
const rejection = expect(pending).rejects.toThrow();
await vi.waitFor(() => expect(entered).toBe(true));
await cache.cancelQueries({ queryKey: options.queryKey });
await rejection;
await vi.waitFor(() => expect(transportAborted).toBe(true));
expect(cache.getQueryData(options.queryKey)).toBeUndefined();
});
});
describe("history refresh and mutation response integrity", () => {
it("retains the last valid page while a refresh fails", async () => {
server.use(
http.get(endpoint, () =>
HttpResponse.json({ ok: true, jobs: [furnitureJobFixture] }),
),
);
const cache = client(),
options = furnitureJobsQueryOptions(false, null);
const observer = new QueryObserver(cache, { ...options, enabled: false });
const unsubscribe = observer.subscribe(() => {});
try {
await cache.fetchQuery(options);
server.use(
http.get(endpoint, () =>
HttpResponse.json({ error: "offline" }, { status: 500 }),
),
);
await expect(cache.fetchQuery(options)).rejects.toThrow("offline");
expect(observer.getCurrentResult().data?.jobs[0].id).toBe(
furnitureJobFixture.id,
);
expect(observer.getCurrentResult().error?.message).toBe("offline");
} finally {
unsubscribe();
}
});
it("rejects malformed successful mutation responses rather than storing an undefined job", async () => {
server.use(
http.patch(endpoint, () =>
HttpResponse.json({ ok: true, job: { id: furnitureJobFixture.id } }),
),
);
const response = await fetch(endpoint, { method: "PATCH" });
await expect(
readFurnitureJobResponse(response, "Update failed"),
).rejects.toThrow("Invalid import job response");
});
it("reports mutation HTTP failure without issuing a second write", async () => {
let calls = 0;
server.use(
http.post(endpoint, () => {
calls++;
return HttpResponse.json(
{ error: "Queue unavailable" },
{ status: 500 },
);
}),
);
const response = await fetch(endpoint, { method: "POST" });
await expect(
readFurnitureJobResponse(response, "Queue failed"),
).rejects.toMatchObject({ status: 500, message: "Queue unavailable" });
expect(calls).toBe(1);
});
});
it("does not park manual refresh indefinitely behind a global offline signal", async () => {
server.use(
http.get(endpoint, () => HttpResponse.json({ ok: true, jobs: [] })),
);
onlineManager.setOnline(false);
try {
await expect(
client().fetchQuery(furnitureJobsQueryOptions(false, null)),
).resolves.toEqual({ jobs: [], nextCursor: null });
} finally {
onlineManager.setOnline(true);
}
});
+82
View File
@@ -148,6 +148,80 @@ const schema = z
.string()
.optional()
.transform((value) => value !== "false" && value !== "0"),
// CrowdSec API — optional. When the CTI API key is set, the anti-DDoS
// gate consults the community reputation of repeat offenders (CTI
// GET /smoke/{ip}) and hard-blocks known-bad IPs immediately. Like the
// Cloudflare token, the key lives in env only and is never written into
// the admin-visible config. Free/community key: app.crowdsec.net →
// Settings → CTI API Keys.
CROWDSEC_API_KEY: z.string().optional(),
// Reputation lookup (CTI) endpoint; overridden for tests/staging.
CROWDSEC_CTI_BASE_URL: z
.string()
.url()
.default("https://cti.api.crowdsec.net/v2"),
// Boot default for the runtime "auto-block from CrowdSec reputation"
// toggle (overridable via the admin panel / antiddos:config).
CROWDSEC_AUTO_BLOCK_ENABLED: z
.string()
.optional()
.transform((value) => value !== "false" && value !== "0"),
// Minimum malevolence score (CrowdSec scores are 0-5; 4-5 maps to
// "malicious") an IP must reach before the gate treats it as known-bad.
// An IP the community already labels "malicious" is always blocked,
// unless it carries false-positive classification tags.
CROWDSEC_BLOCK_SCORE: z.coerce.number().int().min(0).max(5).default(4),
// How long a CrowdSec-confirmed bad IP stays blocked by the gate.
CROWDSEC_BLOCK_TTL_SECONDS: z.coerce
.number()
.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),
// How many new community-reputation blocks within a 5-minute window
// justify an ops alert (cooldown-gated via HEALTH_ALERT_COOLDOWN_MIN).
CROWDSEC_ALERT_BLOCK_BURST: z.coerce.number().int().min(1).default(10),
// 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"),
// Local CrowdSec engine shipped as an opt-in Docker stack in
// deployment/crowdsec. When enabled, the anti-DDoS gate asks the local
// LAPI (bouncer) for each client IP before its own buckets and blocks
// ban/captcha decisions immediately. The key lives in env only.
CROWDSEC_LOCAL_ENABLED: z
.string()
.optional()
.transform((value) => value === "true" || value === "1"),
CROWDSEC_LAPI_URL: z
.string()
.optional()
.transform((value) =>
value?.trim() ? value.trim() : "http://127.0.0.1:18080",
)
.pipe(z.string().url()),
CROWDSEC_LAPI_API_KEY: z.string().optional(),
CROWDSEC_LAPI_TIMEOUT_MS: z.coerce.number().int().positive().default(500),
// 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;
@@ -175,6 +249,14 @@ const schema = z
path: ["PAYPAL_CLIENT_ID"],
});
}
if (data.CROWDSEC_LOCAL_ENABLED && !data.CROWDSEC_LAPI_API_KEY) {
ctx.addIssue({
code: "custom",
message:
"CROWDSEC_LAPI_API_KEY is required when CROWDSEC_LOCAL_ENABLED=true",
path: ["CROWDSEC_LAPI_API_KEY"],
});
}
});
type Env = z.infer<typeof schema>;
+1 -2
View File
@@ -217,9 +217,8 @@ function readFlatFromMemory(
}
const pageMap = new Map(rows.map((row) => [toInt(row.id), row]));
const depths = new Map<number, number>();
const visitingGlobal = new Set<number>();
function depth(id: number, visiting = visitingGlobal): number {
function depth(id: number, visiting = new Set<number>()): number {
const cached = depths.get(id);
if (cached !== undefined) return cached;
// Messy imports can leave a parent chain looping; stop rather than recurse.
+41
View File
@@ -24,6 +24,11 @@ export interface AntiddosConfig {
blockTiers: AntiddosBlockTier[];
globalHaltMs: number;
cloudflareAutoBlock: boolean;
crowdsecAutoBlock: boolean;
/** Minimum CrowdSec malevolence score (0-5) treated as known-bad. */
crowdsecBlockScore: number;
/** How long a CrowdSec-confirmed bad IP stays blocked by the gate. */
crowdsecBlockTtlSeconds: number;
}
const DEFAULT_CONFIG: AntiddosConfig = {
@@ -41,6 +46,9 @@ const DEFAULT_CONFIG: AntiddosConfig = {
],
globalHaltMs: 10_000,
cloudflareAutoBlock: true,
crowdsecAutoBlock: true,
crowdsecBlockScore: 4,
crowdsecBlockTtlSeconds: 86_400,
};
function positiveInt(value: number | undefined, fallback: number): number {
@@ -49,6 +57,17 @@ function positiveInt(value: number | undefined, fallback: number): number {
return Math.floor(n);
}
function clampInt(
value: number | undefined,
fallback: number,
min: number,
max: number,
): number {
const n = Number(value);
if (!Number.isFinite(n)) return fallback;
return Math.min(max, Math.max(min, Math.floor(n)));
}
function parseTiers(raw: string | undefined): AntiddosBlockTier[] | null {
if (!raw?.trim()) return null;
const tiers: AntiddosBlockTier[] = [];
@@ -125,6 +144,17 @@ export function antiddosDefaultsFromEnv(): AntiddosConfig {
DEFAULT_CONFIG.globalHaltMs,
),
cloudflareAutoBlock: isTruthyFlag(env.CLOUDFLARE_AUTO_BLOCK_ENABLED),
crowdsecAutoBlock: isTruthyFlag(env.CROWDSEC_AUTO_BLOCK_ENABLED),
crowdsecBlockScore: clampInt(
env.CROWDSEC_BLOCK_SCORE,
DEFAULT_CONFIG.crowdsecBlockScore,
0,
5,
),
crowdsecBlockTtlSeconds: positiveInt(
env.CROWDSEC_BLOCK_TTL_SECONDS,
DEFAULT_CONFIG.crowdsecBlockTtlSeconds,
),
};
}
@@ -167,6 +197,17 @@ function sanitize(config: AntiddosConfig): AntiddosConfig {
: base.blockTiers,
globalHaltMs: positiveInt(config?.globalHaltMs, base.globalHaltMs),
cloudflareAutoBlock: config?.cloudflareAutoBlock !== false,
crowdsecAutoBlock: config?.crowdsecAutoBlock !== false,
crowdsecBlockScore: clampInt(
config?.crowdsecBlockScore,
base.crowdsecBlockScore,
0,
5,
),
crowdsecBlockTtlSeconds: positiveInt(
config?.crowdsecBlockTtlSeconds,
base.crowdsecBlockTtlSeconds,
),
};
}
+1 -3
View File
@@ -84,7 +84,7 @@ function snapshot(): Record<string, CacheKeyStats> {
}
async function flush(): Promise<void> {
if (process.env.NODE_ENV === "test" || redis?.status !== "ready") return;
if (redis?.status !== "ready") return;
try {
await redis.setex(
REDIS_KEY,
@@ -127,11 +127,9 @@ export async function readCacheStats(): Promise<CacheStatsReport> {
shared = true;
}
} catch (error) {
if (process.env.NODE_ENV !== "test") {
logger.warn("[cache] stats unavailable", { error: String(error) });
}
}
}
const totals = empty();
for (const value of Object.values(keys)) {
+1 -1
View File
@@ -78,7 +78,7 @@ describe("cached (memory-only, no Redis)", () => {
// key survives even though it was inserted first by a long way.
for (let i = 0; i < 2_100; i++) {
await cached(hotKey, 60_000, hot);
await cached(`churn-${i}`, 60_000, cold);
await cached(`churn-${i}-${Math.random()}`, 60_000, cold);
}
expect(hot).toHaveBeenCalledTimes(1);
+16 -24
View File
@@ -99,10 +99,10 @@ function setMemory(key: string, entry: CacheEntry): void {
if (memory.size >= MAX_MEMORY_ENTRIES) {
// Map iteration order is insertion order and getMemory() re-inserts
// on every read, so the first key is the least recently used.
for (const lruKey of memory.keys()) {
memory.delete(lruKey);
recordCacheOutcome(lruKey, "evicted");
break;
const lru = memory.keys().next().value;
if (lru !== undefined) {
memory.delete(lru);
recordCacheOutcome(lru, "evicted");
}
}
}
@@ -128,14 +128,12 @@ function getMemory(key: string): CacheEntry | undefined {
return entry;
}
/** Deterministic TTL - no random jitter to prevent unpredictable drops. */
/** `setex` with a little jitter so keys written together do not expire together. */
function redisTtlSeconds(ttlSec: number): number {
return ttlSec;
const jitter = Math.min(5, Math.floor(ttlSec * 0.1));
return ttlSec + Math.floor(Math.random() * (jitter + 1));
}
// Deduplication: track promised values to prevent double caching/computation
const computationPromises = new Map<string, Promise<unknown>>();
// A wrong REDIS_URL or a Redis outage does not fail visibly: the cache keeps
// answering from memory and every instance quietly stops sharing. Say so once.
let warnedSharedCacheDown = false;
@@ -199,13 +197,6 @@ export async function cached<T>(
return (await pending) as T;
}
// Check for existing computation promise to avoid duplicate work
const existingPromise = computationPromises.get(key);
if (existingPromise) {
recordCacheOutcome(key, "miss");
return existingPromise as Promise<T>;
}
recordCacheOutcome(key, "miss");
return refresh(key, ttlMs, staleMs, fn);
}
@@ -222,8 +213,8 @@ function refresh<T>(
): Promise<T> {
const generation = generations.get(key) ?? 0;
const compute = (async (): Promise<T> => {
// Redis path (shared across instances) - skip during tests for speed/stability
if (process.env.NODE_ENV !== "test" && redis && redis.status !== "end") {
// Redis path (shared across instances).
if (redis && redis.status !== "end") {
try {
const raw = await redis.get(key);
if (raw !== null && raw !== undefined) {
@@ -234,7 +225,7 @@ function refresh<T>(
} catch {
/* fall through to fn */
}
} else if (process.env.NODE_ENV !== "test") {
} else {
warnIfSharedCacheDown();
}
@@ -246,10 +237,13 @@ function refresh<T>(
// An invalidation landed while `fn()` was running: keep the value out of
// the cache so the next read recomputes instead of resurrecting stale data.
if ((generations.get(key) ?? 0) === generation) {
if (process.env.NODE_ENV !== "test" && redis && redis.status !== "end") {
if (redis && redis.status !== "end") {
try {
const ttlSec = Math.max(1, Math.ceil(ttlMs / 1000));
await redis.setex(key, redisTtlSeconds(ttlSec), JSON.stringify(data));
await redis.setex(
key,
redisTtlSeconds(Math.ceil(ttlMs / 1000)),
JSON.stringify(data),
);
} catch {
/* non-critical: memory cache still works */
}
@@ -261,11 +255,9 @@ function refresh<T>(
return data;
})().finally(() => {
inFlight.delete(key);
computationPromises.delete(key);
});
inFlight.set(key, compute);
computationPromises.set(key, compute);
return compute;
}
+66
View File
@@ -0,0 +1,66 @@
import "server-only";
import { env } from "@/env";
import { logger } from "@/lib/logger";
import { redis } from "@/lib/redis";
import { type SendAlertInput, sendAlert } from "@/lib/services/alert";
// === CrowdSec operational alerts ===========================================
//
// Thin, fire-and-forget wrapper around the app's alert service for the
// reputation pipeline. Every raise is cooldown-gated through a Redis NX lock
// (key crowdsec:alert:{key}, TTL = HEALTH_ALERT_COOLDOWN_MIN), so N instances
// and flapping conditions surface exactly one alert per window instead of
// spamming Discord/email/alert_logs. Falls back to alerting anyway when Redis
// is unreachable — a silent quota blowout is worse than one duplicate alert.
const ALERT_PREFIX = "crowdsec:alert:";
function cooldownSeconds(): number {
const raw = Number(env.HEALTH_ALERT_COOLDOWN_MIN ?? 15);
return Math.ceil((Number.isFinite(raw) && raw > 0 ? raw : 15) * 60);
}
/**
* Raise an alert unless the cooldown window is still active. Returns the
* sendAlert promise when the alert was actually raised, or false when it was
* suppressed. Never throws; the caller may `void` the result on hot paths.
*/
export async function raiseCrowdsecAlert(
key: string,
input: {
type?: string;
severity: SendAlertInput["severity"];
message: string;
context?: SendAlertInput["context"];
},
): Promise<false | Awaited<ReturnType<typeof sendAlert>>> {
if (redis) {
try {
const acquired = await redis.set(
`${ALERT_PREFIX}${key}`,
String(Date.now()),
"EX",
cooldownSeconds(),
"NX",
);
if (acquired !== "OK") return false;
} catch {
// Cooldown bookkeeping failed — alert anyway rather than silently drop.
}
}
try {
return await sendAlert({
type: "ddos",
severity: input.severity,
message: input.message,
context: input.context,
});
} catch (error) {
logger.error("[crowdsec-alert] sendAlert raised an unexpected error", {
key,
err: error,
});
return false;
}
}
+798
View File
@@ -0,0 +1,798 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import {
type CrowdsecVerdict,
crowdsecEnabled,
getCrowdsecApiConfig,
getCrowdsecBlockMeta,
getCrowdsecQuotaUsage,
getLastCrowdsecVerify,
getMemoryVerdictCacheSize,
lookupCrowdsecVerdict,
maybeAutoBlockCrowdsec,
resetCrowdsecCache,
setLastCrowdsecVerify,
verdictIsMalicious,
verifyCrowdsecConnection,
} from "./crowdsec-api";
import { type CrowdsecDailyStat, getCrowdsecStats } from "./crowdsec-stats";
// Unit-test the CTI client in isolation: a deterministic in-memory Redis fake
// and a silenced logger, so fetch calls count only CrowdSec lookups. CrowdSec
// deliberately never touches Cloudflare, so no Cloudflare surface is stubbed.
const state = vi.hoisted(() => ({
map: new Map<string, string>(),
z: new Map<string, Array<[number, string]>>(),
sendAlert: vi.fn(),
}));
vi.mock("@/lib/services/alert", () => ({
sendAlert: state.sendAlert,
ddosDetected: vi.fn(),
}));
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;
},
incr: async (key: string) => {
const next = (Number(state.map.get(key)) || 0) + 1;
state.map.set(key, String(next));
return next;
},
decr: 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,
zadd: async (key: string, score: number, member: string) => {
const list = state.z.get(key) ?? [];
list.push([score, member]);
list.sort((a, b) => a[0] - b[0]);
state.z.set(key, list);
return 1;
},
zremrangebyscore: async (key: string, min: number, max: number) => {
const list = (state.z.get(key) ?? []).filter(
([score]) => score < min || score > max,
);
state.z.set(key, list);
return 1;
},
zcard: async (key: string) => (state.z.get(key) ?? []).length,
},
__esModule: true,
}));
vi.mock("@/lib/logger", () => ({
logger: {
info: vi.fn(),
warn: vi.fn(),
error: vi.fn(),
debug: vi.fn(),
},
}));
const tick = () => new Promise((resolve) => setTimeout(resolve, 20));
function jsonResponse(body: unknown, status = 200): Response {
return new Response(JSON.stringify(body), {
status,
headers: { "content-type": "application/json" },
});
}
function maliciousItem(ip: string, score = 5): unknown {
return {
ip,
reputation: "malicious",
confidence: "0.95",
scores: { overall: { aggressiveness: 4, total: score } },
behaviors: [{ name: "http:bruteforce" }, { name: "http:scan" }],
classifications: { false_positives: [] },
};
}
function suspiciousItem(ip: string, score: number): unknown {
return {
ip,
reputation: "suspicious",
scores: { overall: { total: score } },
};
}
function verdict(minimal: Partial<CrowdsecVerdict> = {}): CrowdsecVerdict {
return {
ip: "198.51.100.1",
reputation: "suspicious",
score: 3,
aggressiveness: 0,
confidence: null,
behaviors: [],
falsePositive: false,
checkedAt: Date.now(),
...minimal,
};
}
function blockIp(): string {
return "198.51.100.1";
}
describe("crowdsec-api", () => {
let fetchMock: ReturnType<typeof vi.fn>;
beforeEach(() => {
vi.unstubAllGlobals();
vi.unstubAllEnvs();
state.map.clear();
state.z.clear();
state.sendAlert.mockReset();
resetCrowdsecCache();
fetchMock = vi.fn();
vi.stubGlobal("fetch", fetchMock);
});
afterEach(() => {
vi.unstubAllGlobals();
vi.unstubAllEnvs();
state.map.clear();
resetCrowdsecCache();
vi.restoreAllMocks();
});
it("is enabled only when a non-blank API key is configured", () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
vi.stubEnv("CROWDSEC_CTI_BASE_URL", "https://cti.example.test");
expect(crowdsecEnabled()).toBe(true);
expect(getCrowdsecApiConfig().apiKey).toBe("cs_key");
expect(getCrowdsecApiConfig().baseUrl).toBe("https://cti.example.test");
vi.stubEnv("CROWDSEC_API_KEY", " ");
expect(crowdsecEnabled()).toBe(false);
vi.stubEnv("CROWDSEC_CTI_BASE_URL", "");
expect(getCrowdsecApiConfig().baseUrl).toBe(
"https://cti.api.crowdsec.net/v2",
);
});
it("blocks malicious reputations at any threshold", () => {
expect(
verdictIsMalicious(verdict({ reputation: "malicious", score: 0 }), 5),
).toBe(true);
});
it("never blocks safe or benign reputations", () => {
expect(verdictIsMalicious(verdict({ reputation: "safe" }), 0)).toBe(false);
expect(verdictIsMalicious(verdict({ reputation: "benign" }), 1)).toBe(
false,
);
});
it("vetoes a false-positive tag even for a malicious reputation", () => {
expect(
verdictIsMalicious(
verdict({ reputation: "malicious", score: 5, falsePositive: true }),
4,
),
).toBe(false);
});
it("applies the score threshold to suspicious/known attackers", () => {
const v4 = verdict({ reputation: "suspicious", score: 4 });
expect(verdictIsMalicious(v4, 4)).toBe(true);
expect(verdictIsMalicious(v4, 5)).toBe(false);
expect(verdictIsMalicious(verdict({ score: 3 }), 4)).toBe(false);
});
it("never blocks score-0 (unknown) IPs even at threshold 0", () => {
expect(
verdictIsMalicious(verdict({ score: 0, reputation: "unknown" }), 0),
).toBe(false);
});
it("does nothing without an API key", async () => {
await maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
expect(fetchMock).not.toHaveBeenCalled();
});
it("does nothing when the runtime toggle is off", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
await maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: false,
});
expect(fetchMock).not.toHaveBeenCalled();
});
it("never consults CrowdSec for the unknown-IP sentinel", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
await maybeAutoBlockCrowdsec({
ip: "0.0.0.0",
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
expect(fetchMock).not.toHaveBeenCalled();
});
it("hard-blocks a malicious IP in the shared gate key only", 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,
});
expect(state.map.get(`antiddos:block:${blockIp()}`)).toBe("crowdsec");
// No Cloudflare keys may ever be written by CrowdSec.
expect(
[...state.map.keys()].some((key) => key.includes("cloudflare")),
).toBe(false);
expect(fetchMock).toHaveBeenCalledTimes(1);
const [url, init] = fetchMock.mock.calls[0];
expect(String(url)).toContain(`/smoke/${blockIp()}`);
expect((init.headers as Record<string, string>)["x-api-key"]).toBe(
"cs_key",
);
});
it("does not block an IP the community knows nothing about", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(jsonResponse({}, 404));
await maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
expect(state.map.has(`antiddos:block:${blockIp()}`)).toBe(false);
});
it("respects a custom score threshold", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(jsonResponse(suspiciousItem(blockIp(), 3)));
await maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
expect(state.map.has(`antiddos:block:${blockIp()}`)).toBe(false);
fetchMock.mockResolvedValue(jsonResponse(suspiciousItem(blockIp(), 3)));
await maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 3,
enabled: true,
});
expect(state.map.get(`antiddos:block:${blockIp()}`)).toBe("crowdsec");
});
it("dedupes concurrent lookups into a single API call", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(jsonResponse(maliciousItem(blockIp())));
await Promise.all([
maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
}),
maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
}),
]);
expect(fetchMock).toHaveBeenCalledTimes(1);
});
it("reuses the Redis verdict cache for repeat offenders", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(jsonResponse(maliciousItem(blockIp())));
for (let i = 0; i < 3; i += 1) {
await maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
}
expect(fetchMock).toHaveBeenCalledTimes(1);
});
it("never shortens an existing longer host block", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(jsonResponse(maliciousItem(blockIp())));
state.map.set(`antiddos:block:${blockIp()}`, "1");
const redis = (await import("@/lib/redis")).redis;
vi.spyOn(redis as NonNullable<typeof redis>, "pttl").mockImplementation(
async () => 86_400_000,
);
await maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
expect(state.map.get(`antiddos:block:${blockIp()}`)).toBe("1");
});
it("backs off after a 403 so it stops hammering a rejected key", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(jsonResponse({ message: "Invalid key" }, 403));
await maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
await maybeAutoBlockCrowdsec({
ip: "203.0.113.9",
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});
it("publishes the backoff to shared Redis so every instance respects it", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(jsonResponse({ message: "Invalid key" }, 403));
await maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
// The shared marker exists and points into the future.
const until = Number(state.map.get("crowdsec:backoff-until"));
expect(Number.isFinite(until)).toBe(true);
expect(until).toBeGreaterThan(Date.now());
// A fresh instance (reset in-process state) still honours the marker.
resetCrowdsecCache();
await maybeAutoBlockCrowdsec({
ip: "203.0.113.44",
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});
it("backs off after a 429 rate limit as well", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(jsonResponse({ message: "rate limited" }, 429));
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(1);
});
it("raises a critical ops alert when the CTI key is rejected (403)", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(jsonResponse({ message: "Invalid key" }, 403));
await maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
await tick();
expect(state.sendAlert).toHaveBeenCalledTimes(1);
const [input] = state.sendAlert.mock.calls[0];
expect(input.type).toBe("ddos");
expect(input.severity).toBe("critical");
expect(input.message).toContain("403");
expect(input.message).toContain("CROWDSEC_API_KEY");
expect(input.context).toMatchObject({ status: 403 });
});
it("raises a warning ops alert when the CTI rate limit is hit (429)", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(jsonResponse({ message: "rate limited" }, 429));
await maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
await tick();
expect(state.sendAlert).toHaveBeenCalledTimes(1);
const [input] = state.sendAlert.mock.calls[0];
expect(input.type).toBe("ddos");
expect(input.severity).toBe("warning");
expect(input.message).toContain("rate limited");
expect(input.context).toMatchObject({ status: 429 });
});
it("swallows API failures instead of throwing on the hot path", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(jsonResponse({ message: "boom" }, 500));
await expect(
maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
}),
).resolves.toBeUndefined();
expect(state.map.has(`antiddos:block:${blockIp()}`)).toBe(false);
});
it("reports a missing credential without calling the API", async () => {
const status = await verifyCrowdsecConnection();
expect(status.ok).toBe(false);
expect(status.message).toContain("CROWDSEC_API_KEY");
expect(fetchMock).not.toHaveBeenCalled();
});
it("verifies the key against the CTI probe endpoint", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(
jsonResponse({ ip: "1.1.1.1", reputation: "safe" }),
);
const status = await verifyCrowdsecConnection();
expect(status.ok).toBe(true);
expect(status.message).toContain("1.1.1.1");
expect(String(fetchMock.mock.calls[0][0])).toContain("/smoke/1.1.1.1");
expect(
(fetchMock.mock.calls[0][1].headers as Record<string, string>)[
"x-api-key"
],
).toBe("cs_key");
});
it("surfaces a rejected credential", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(jsonResponse({ message: "Invalid key" }, 403));
const status = await verifyCrowdsecConnection();
expect(status.ok).toBe(false);
expect(status.message).toContain("Invalid key");
});
it("round-trips the last verify status through Redis and memory", async () => {
const status = { ok: true, message: "CTI key accepted", at: Date.now() };
await setLastCrowdsecVerify(status);
expect(await getLastCrowdsecVerify()).toEqual(status);
expect(JSON.parse(state.map.get("crowdsec:last-verify") ?? "{}")).toEqual(
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);
});
it("raises an ops alert when the daily quota is exhausted", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
vi.stubEnv("CROWDSEC_CTI_DAILY_QUOTA", "1");
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,
});
await tick();
expect(fetchMock).toHaveBeenCalledTimes(1);
expect(state.sendAlert).toHaveBeenCalledTimes(1);
const [input] = state.sendAlert.mock.calls[0];
expect(input.type).toBe("ddos");
expect(input.severity).toBe("warning");
expect(input.message).toContain("quota exhausted");
expect(input.context).toMatchObject({ quota: 1 });
});
it("alerts once when quota becomes available again after exhaustion", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
vi.stubEnv("CROWDSEC_CTI_DAILY_QUOTA", "1");
fetchMock.mockImplementation(() =>
Promise.resolve(jsonResponse(maliciousItem(blockIp()))),
);
await maybeAutoBlockCrowdsec({
ip: blockIp(),
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
// Exhaust the counter.
await maybeAutoBlockCrowdsec({
ip: "198.51.100.2",
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
// The counter is externally reset (new billing day / fresh deployment):
// the next successful reserve should call it out.
const date = new Date().toISOString().slice(0, 10);
state.map.set(`crowdsec:usage:${date}`, "0");
await maybeAutoBlockCrowdsec({
ip: "198.51.100.3",
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
await tick();
expect(fetchMock).toHaveBeenCalledTimes(2);
expect(state.sendAlert).toHaveBeenCalledTimes(2);
const alerts = state.sendAlert.mock.calls.map(([input]) => input);
expect(alerts[0].severity).toBe("warning");
expect(alerts[1].severity).toBe("info");
expect(alerts[1].message).toContain("available again");
expect(alerts[1].context).toMatchObject({ quota: 1 });
});
it("floods once per cooldown window when blocks burst past the threshold", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
vi.stubEnv("CROWDSEC_ALERT_BLOCK_BURST", "2");
vi.stubEnv("CROWDSEC_CTI_DAILY_QUOTA", "0");
fetchMock.mockImplementation((url: string | URL) =>
Promise.resolve(
jsonResponse(maliciousItem(String(url).split("/").pop() ?? "ip")),
),
);
await maybeAutoBlockCrowdsec({
ip: "198.51.100.71",
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
await maybeAutoBlockCrowdsec({
ip: "198.51.100.72",
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
// Third block in the same window: threshold crossed, but the alert is
// cooldown-gated so it still fires exactly once.
await maybeAutoBlockCrowdsec({
ip: "198.51.100.73",
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
await tick();
expect(state.sendAlert).toHaveBeenCalledTimes(1);
const [input] = state.sendAlert.mock.calls[0];
expect(input.type).toBe("ddos");
expect(input.context).toMatchObject({ blocks: 2, threshold: 2 });
expect(state.map.get(`antiddos:block:198.51.100.72`)).toBe("crowdsec");
});
it("tallies lookups and blocks into the daily stats histogram", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
vi.stubEnv("CROWDSEC_CTI_DAILY_QUOTA", "0");
fetchMock.mockImplementation((url: string | URL) => {
const ip = String(url).split("/").pop() ?? "ip";
return Promise.resolve(
jsonResponse(
ip === "198.51.100.83" ? suspiciousItem(ip, 3) : maliciousItem(ip),
),
);
});
await maybeAutoBlockCrowdsec({
ip: "198.51.100.81",
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
await maybeAutoBlockCrowdsec({
ip: "198.51.100.82",
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
await maybeAutoBlockCrowdsec({
ip: "198.51.100.83",
category: "api",
ttlSeconds: 600,
scoreThreshold: 7, // suspicious/known verdicts below threshold: lookup only
enabled: true,
});
await tick();
const stats = await getCrowdsecStats(1);
const today: CrowdsecDailyStat | undefined = stats.find(
(row) => row.date === new Date().toISOString().slice(0, 10),
);
expect(today?.lookups).toBe(3);
expect(today?.blocks).toBe(2);
expect(today?.reportFailures).toBe(0);
expect(today?.categories).toMatchObject({ api: 2 });
expect(today?.reputations).toMatchObject({ malicious: 2 });
expect(
state.map.get(
`crowdsec:stat:lookups:${new Date().toISOString().slice(0, 10)}`,
),
).toBe("3");
});
});
+789
View File
@@ -0,0 +1,789 @@
import "server-only";
import { randomUUID } from "node:crypto";
import { env } from "@/env";
import { raiseCrowdsecAlert } from "@/lib/crowdsec-alerts";
import { reportCrowdsecSignal } from "@/lib/crowdsec-report";
import {
bumpCrowdsecBreakdownStat,
bumpCrowdsecStat,
} from "@/lib/crowdsec-stats";
import { logger } from "@/lib/logger";
import { redis } from "@/lib/redis";
import { UNKNOWN_CLIENT_IP } from "./client-ip";
/**
* CrowdSec CTI (community threat intelligence) integration for the anti-DDoS
* gate — the reputation side of the auto-block pipeline.
*
* When an IP trips a rate bucket, the gate consults CrowdSec's community
* reputation for that IP (`GET /smoke/{ip}`, the freemium Enrichment API,
* `x-api-key` auth) and hard-blocks known-bad repeat offenders immediately
* instead of waiting for the local `maxViolations` threshold. The block lives
* only in the gate's own shared Redis key (`antiddos:block:{ip}`) so every
* existing consumer — the proxy check, the admin block list, the admin unban —
* keeps working unchanged. CrowdSec never talks to Cloudflare and never
* creates edge rules; if a Cloudflare mirror is wanted it is the gate's own
* escalation logic that decides, never this module.
*
* Quota safety: lookups only run for IPs that already tripped a bucket (never
* on the plain hot path), verdicts are cached in Redis for an hour (so a
* flood from one IP costs at most one API call), concurrent lookups for the
* same IP are deduped across instances with a Redis NX lock, and a 403/429
* response trips a module-wide backoff instead of hammering the API.
*
* Credentials come from env only (`CROWDSEC_API_KEY`) and are never written
* into the admin-visible config — same contract as the Cloudflare token.
*
* 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 {}
export interface CrowdsecApiConfig {
baseUrl: string;
apiKey: string | null;
}
export type CrowdsecReputation =
| "malicious"
| "suspicious"
| "known"
| "safe"
| "benign"
| "unknown";
export interface CrowdsecVerdict {
ip: string;
/** Raw CTI reputation enum; null when the IP is unknown to the community. */
reputation: CrowdsecReputation | null;
/** `scores.overall.total` — CrowdSec malevolence score, 0-5. */
score: number;
/** `scores.overall.aggressiveness` — 0-5. */
aggressiveness: number;
confidence: string | null;
/** Reported attack categories, e.g. ["http:scan", "ssh:bruteforce"]. */
behaviors: string[];
/** CrowdSec tags IPs carrying false-positive classifications as safe. */
falsePositive: boolean;
checkedAt: number;
}
export interface CrowdsecConnectionStatus {
ok: boolean;
message?: string;
at: number;
}
/** Why a CrowdSec-sourced block exists — persisted next to the block key. */
export interface CrowdsecBlockMeta {
source: typeof CROWDSEC_BLOCK_SOURCE | "gate";
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";
const API_TIMEOUT_MS = 10_000;
const VERDICT_CACHE_TTL_SECONDS = 3600;
const VERDICT_CACHE_TTL_MS = VERDICT_CACHE_TTL_SECONDS * 1000;
const LOOKUP_LOCK_TTL_SECONDS = 60;
const RATE_LIMIT_BACKOFF_MS = 60_000;
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;
/** Shared 403/429 pause marker, so every instance respects the backoff. */
const BACKOFF_KEY = "crowdsec:backoff-until";
/** Short-window block burst counter: crowdsec:burst:recent (ZSET of timestamps). */
const BURST_PREFIX = "crowdsec:burst:";
const BURST_WINDOW_SECONDS = 300;
/** 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";
export function getCrowdsecApiConfig(): CrowdsecApiConfig {
return {
baseUrl: env.CROWDSEC_CTI_BASE_URL || "https://cti.api.crowdsec.net/v2",
apiKey: env.CROWDSEC_API_KEY?.trim() || null,
};
}
/** True when a CTI API key is present so the gate may call the API. */
export function crowdsecEnabled(): boolean {
return Boolean(getCrowdsecApiConfig().apiKey);
}
interface CrowdsecScore {
aggressiveness?: number;
threat?: number;
trust?: number;
anomaly?: number;
total?: number;
}
interface CrowdsecSmokeItem {
ip?: string;
reputation?: string;
confidence?: string;
scores?: { overall?: CrowdsecScore };
classifications?: { false_positives?: unknown[] };
behaviors?: { name?: string }[];
}
async function crowdsecRequest(
path: string,
init: { method?: "GET" | "POST"; body?: unknown } = {},
): Promise<Response> {
const config = getCrowdsecApiConfig();
if (!config.apiKey) {
throw new CrowdsecApiError("CROWDSEC_API_KEY is not configured");
}
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), API_TIMEOUT_MS);
try {
return await fetch(`${config.baseUrl}${path}`, {
method: init.method ?? "GET",
headers: {
"x-api-key": config.apiKey,
Accept: "application/json",
"Content-Type": "application/json",
},
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}`;
}
}
function toNumber(value: unknown): number {
const n = Number(value);
return Number.isFinite(n) ? n : 0;
}
function parseVerdict(ip: string, item: CrowdsecSmokeItem): CrowdsecVerdict {
const overall = item.scores?.overall;
const falsePositives = item.classifications?.false_positives ?? [];
return {
ip,
reputation: (item.reputation as CrowdsecReputation | undefined) ?? null,
score: toNumber(overall?.total),
aggressiveness: toNumber(overall?.aggressiveness),
confidence: item.confidence ?? null,
behaviors: (item.behaviors ?? [])
.map((behavior) => behavior?.name)
.filter((name): name is string => Boolean(name)),
// CrowdSec: "Any IP with false_positives tags shouldn't be considered
// as malicious" — this veto always wins over reputation and score.
falsePositive: falsePositives.length > 0,
checkedAt: Date.now(),
};
}
/**
* Resolve a cached CTI verdict into a block/no-block decision against the
* admin-configurable score threshold. `malicious` is always blocked; `safe`
* and `benign` never are; everything else follows the 0-5 score threshold
* (with score 0 = "unknown" never blocking, even at threshold 0).
*/
export function verdictIsMalicious(
verdict: CrowdsecVerdict,
threshold: number,
): boolean {
if (verdict.falsePositive) return false;
if (verdict.reputation === "malicious") return true;
if (verdict.reputation === "safe" || verdict.reputation === "benign") {
return false;
}
const effective = Math.min(5, Math.max(0, threshold));
return verdict.score >= effective && verdict.score >= 1;
}
// --- Verdict cache (Redis backed, in-process fallback) ---
const memoryVerdicts = new Map<string, CrowdsecVerdict>();
/**
* 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}`;
}
async function readVerdictCache(ip: string): Promise<CrowdsecVerdict | null> {
const cached = memoryVerdicts.get(ip);
if (cached && Date.now() - cached.checkedAt < VERDICT_CACHE_TTL_MS) {
return cached;
}
if (redis) {
try {
const raw = await redis.get(verdictKey(ip));
if (raw) {
const parsed = JSON.parse(raw) as CrowdsecVerdict;
rememberVerdict(parsed);
return parsed;
}
} catch {
// fall through to a cache miss — the API call below is the fallback.
}
}
return null;
}
async function writeVerdictCache(verdict: CrowdsecVerdict): Promise<void> {
rememberVerdict(verdict);
if (redis) {
try {
await redis.set(
verdictKey(verdict.ip),
JSON.stringify(verdict),
"EX",
VERDICT_CACHE_TTL_SECONDS,
);
} catch {
// cache is best-effort — a miss only costs one extra API call later.
}
}
}
/** Cross-instance dedupe so a cold-cache flood costs one lookup, not N. */
async function acquireLookupLock(ip: string): Promise<boolean> {
if (!redis) return true;
try {
const acquired = await redis.set(
`${LOOKUP_LOCK_PREFIX}${ip}`,
"1",
"EX",
LOOKUP_LOCK_TTL_SECONDS,
"NX",
);
return acquired === "OK";
} catch {
// Redis hiccup — allow the lookup; the verdict cache still dedupes.
return true;
}
}
let backoffUntil = 0;
let quotaWarnedDate: string | null = null;
let quotaExhaustedDate: string | null = null;
/**
* Next moment (epoch ms) the CTI API may be called again — the max of the
* in-process view and the shared Redis marker so every instance respects a
* backoff discovered by any of them. Redis is only read when the local view is
* not already active, keeping the hot path cheap.
*/
async function getBackoffUntil(): Promise<number> {
if (Date.now() < backoffUntil) return backoffUntil;
if (redis) {
try {
const raw = await redis.get(BACKOFF_KEY);
const shared = Number(raw ?? 0);
if (Number.isFinite(shared) && shared > backoffUntil) {
backoffUntil = shared;
}
} catch {
// Redis hiccup — local view is enough
}
}
return backoffUntil;
}
async function setBackoff(ms: number): Promise<void> {
const until = Date.now() + ms;
backoffUntil = until;
if (redis) {
try {
// EX rounds up so the marker outlives the wait it encodes, plus a
// second of slack for the read path.
await redis.set(
BACKOFF_KEY,
String(until),
"EX",
Math.ceil(ms / 1000) + 1,
);
} catch {
// local view still protects this instance
}
}
}
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<CrowdsecQuotaUsage> {
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 };
}
/**
* Reserve one API call against today's quota. Atomic: the counter is INCR'd
* BEFORE the call and compared to the ceiling, so concurrent instances can
* never slip calls past the budget; a reserve that overshoots rolls itself
* back. Returns false once the budget is spent (and raises an ops alert).
*/
async function reserveQuota(): Promise<boolean> {
const quota = dailyQuota();
if (quota <= 0) return true;
if (!redis) return true; // no shared counter → unlimited best-effort
const date = quotaDate();
const key = quotaKey(date);
try {
const used = await redis.incr(key);
await redis.expire(key, QUOTA_KEY_TTL_SECONDS);
if (used > quota) {
// Concurrent reserves nudged us past the ceiling — give the slot
// back and refuse: the budget would be spent the very next call
// anyway, so stopping here is both safe and quota-exact.
await redis.decr(key);
quotaExhaustedDate = date;
logger.warn(
"[crowdsec-api] CTI daily quota exhausted — pausing lookups until tomorrow",
{ quota },
);
void raiseCrowdsecAlert("quota", {
type: "ddos",
severity: "warning",
message: `CrowdSec reputation quota exhausted for today (${used} of ${quota} enrichment calls) — lookups are paused until tomorrow.`,
context: { used, quota, date },
});
return false;
}
if (used >= quota * QUOTA_WARN_RATIO && quotaWarnedDate !== date) {
quotaWarnedDate = date;
logger.warn("[crowdsec-api] CTI daily quota nearing its limit", {
used,
quota,
});
}
if (quotaExhaustedDate) {
// A reserve just succeeded after an exhaustion day (counter was
// reset or the calendar rolled over) — say so, once per cooldown.
void raiseCrowdsecAlert("quota-restored", {
type: "ddos",
severity: "info",
message: `CrowdSec reputation quota is available again (${used} of ${quota} used today) — lookups resumed.`,
context: { used, quota, date },
});
quotaExhaustedDate = null;
}
return true;
} catch {
// Redis hiccup at a moment we could not count — allow the call rather
// than break the gate; the verdict cache still limits frequency.
return true;
}
}
/** Block burst threshold from env, defensively coerced (falls back to 10). */
function dailyBlockBurstThreshold(): number {
const raw = Number(env.CROWDSEC_ALERT_BLOCK_BURST ?? 10);
return Number.isFinite(raw) && raw > 0 ? Math.floor(raw) : 10;
}
/**
* A burst of new blocks is usually an automated attack wave. Track block
* timestamps in a rolling window (Redis sorted set, 5 minutes) so a burst that
* straddles a bucket boundary is still counted together, and alert once per
* cooldown window when the count crosses CROWDSEC_ALERT_BLOCK_BURST.
* Fire-and-forget.
*/
async function trackBlockBurst(): Promise<void> {
if (!redis) return;
const now = Date.now();
const key = `${BURST_PREFIX}recent`;
const threshold = dailyBlockBurstThreshold();
try {
await redis.zadd(key, now, randomUUID());
await redis.zremrangebyscore(key, 0, now - BURST_WINDOW_SECONDS * 1000);
const count = await redis.zcard(key);
await redis.expire(key, BURST_WINDOW_SECONDS * 2);
if (count >= threshold) {
void raiseCrowdsecAlert("block-burst", {
type: "ddos",
severity: "warning",
message: `Anti-DDoS auto-block created ${count} blocks in the last ${BURST_WINDOW_SECONDS / 60} minutes — likely an automated attack wave.`,
context: {
blocks: count,
windowSeconds: BURST_WINDOW_SECONDS,
threshold,
},
});
}
} catch {
// alert is best-effort — never break the block path
}
}
/**
* Community reputation verdict for an IP, from cache when possible. Returns
* 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,
): Promise<CrowdsecVerdict | null> {
if (!crowdsecEnabled()) return null;
if (!ip || ip === UNKNOWN_CLIENT_IP) return null;
if (Date.now() < (await getBackoffUntil())) return null;
const cached = await readVerdictCache(ip);
if (cached) return cached;
if (!(await acquireLookupLock(ip))) {
// Another instance is mid-lookup for this IP; skip rather than
// double-spend API quota on the same address.
return null;
}
try {
// Re-read after claiming the lock — a concurrent instance may have
// filled the cache while we were acquiring it.
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)}`);
// The enrichment call happened — count it for the daily histogram,
// regardless of whether the verdict was positive, negative, or n/a.
void bumpCrowdsecStat("lookups");
if (response.status === 404) {
// Unknown to the community — cache the negative result so a clean
// repeat offender never costs another API call this hour.
const verdict = parseVerdict(ip, {});
await writeVerdictCache(verdict);
return verdict;
}
if (response.status === 403) {
const detail = await errorDetail(response);
await setBackoff(AUTH_BACKOFF_MS);
// A rejected key paralyses the whole reputation pipeline — surface
// it once (cooldown-gated) so rotating the key is an ops priority.
void raiseCrowdsecAlert("cti-auth", {
type: "ddos",
severity: "critical",
message: `CrowdSec CTI API key rejected (HTTP 403): ${detail} — reputation lookups are paused for ${Math.round(AUTH_BACKOFF_MS / 60_000)} minutes. Rotate CROWDSEC_API_KEY.`,
context: { status: 403, detail, backoffMs: AUTH_BACKOFF_MS },
});
throw new CrowdsecApiError(
`CrowdSec API key rejected (HTTP 403): ${detail}`,
);
}
if (response.status === 429) {
await setBackoff(RATE_LIMIT_BACKOFF_MS);
logger.warn("[crowdsec-api] CTI API rate limit hit — backing off", {
ip,
backoffMs: RATE_LIMIT_BACKOFF_MS,
});
void raiseCrowdsecAlert("cti-ratelimit", {
type: "ddos",
severity: "warning",
message: `CrowdSec CTI API rate limited — all instances backed off for ${Math.round(RATE_LIMIT_BACKOFF_MS / 1000)}s.`,
context: { status: 429, backoffMs: RATE_LIMIT_BACKOFF_MS },
});
return null;
}
if (!response.ok) {
throw new CrowdsecApiError(
`CrowdSec CTI API error (HTTP ${response.status}): ${await errorDetail(response)}`,
);
}
const item = (await response.json()) as CrowdsecSmokeItem;
const verdict = parseVerdict(ip, item);
await writeVerdictCache(verdict);
return verdict;
} catch (error) {
logger.error("[crowdsec-api] CTI lookup failed", { ip, err: error });
return null;
}
}
/**
* Consult the CrowdSec community reputation of an IP that just tripped a rate
* bucket and hard-block it when the community flags it as known-bad. Safe to
* call fire-and-forget from the hot path: it is never awaited by the caller,
* does nothing when the API is not configured or the runtime toggle is off,
* never shortens an already-active block, and never lets an API failure
* surface to the request.
*/
export async function maybeAutoBlockCrowdsec(input: {
ip: string;
category: string;
ttlSeconds: number;
scoreThreshold: number;
enabled: boolean;
}): Promise<void> {
const { ip, category, ttlSeconds, scoreThreshold, enabled } = input;
if (!enabled) return;
if (!crowdsecEnabled()) return;
if (!ip || ip === UNKNOWN_CLIENT_IP) return;
// The gate only ever reads its block key through shared Redis — without it
// there is nowhere durable to record the block.
if (!redis) return;
if (Date.now() < (await getBackoffUntil())) return;
try {
const verdict = await lookupCrowdsecVerdict(ip);
if (!verdict || !verdictIsMalicious(verdict, scoreThreshold)) return;
const blockKey = `antiddos:block:${ip}`;
const existingTtl = await redis.pttl(blockKey);
// -2 = no key, -1 = no expiry; both fall through and get overwritten
// with the CrowdSec TTL. An equal or longer block is left untouched.
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.
logger.info(
"[crowdsec-api] Automatic IP block created from community reputation",
{
ip,
category,
ttlSeconds,
reputation: verdict.reputation,
score: verdict.score,
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 });
// Daily histogram + burst detection (cooldown-gated ops alert).
void bumpCrowdsecStat("blocks");
void bumpCrowdsecBreakdownStat("category", category);
if (verdict.reputation) {
void bumpCrowdsecBreakdownStat("reputation", verdict.reputation);
}
void trackBlockBurst();
} catch (error) {
logger.error("[crowdsec-api] Automatic IP block failed", {
ip,
err: error,
});
}
}
let lastVerifyMemory: CrowdsecConnectionStatus | null = null;
/** Validate that the configured key can query the CTI (Enrichment) API. */
export async function verifyCrowdsecConnection(): Promise<CrowdsecConnectionStatus> {
const config = getCrowdsecApiConfig();
if (!config.apiKey) {
return {
ok: false,
message: "CROWDSEC_API_KEY is not configured",
at: Date.now(),
};
}
try {
const response = await crowdsecRequest(`/smoke/${PROBE_IP}`);
if (response.ok) {
const item = (await response
.json()
.catch(() => null)) as CrowdsecSmokeItem | null;
const reputation = item?.reputation
? ` (reputation ${item.reputation})`
: "";
return {
ok: true,
message: `CTI key accepted — probed ${PROBE_IP}${reputation}`,
at: Date.now(),
};
}
if (response.status === 403) {
return {
ok: false,
message: `API key rejected: ${await errorDetail(response)}`,
at: Date.now(),
};
}
if (response.status === 429) {
return {
ok: false,
message: "CTI API rate limit reached — try again shortly",
at: Date.now(),
};
}
return {
ok: false,
message: `CrowdSec CTI API error (HTTP ${response.status}): ${await errorDetail(response)}`,
at: Date.now(),
};
} catch (error) {
return {
ok: false,
message:
error instanceof Error ? error.message : "CrowdSec API unreachable",
at: Date.now(),
};
}
}
export async function getLastCrowdsecVerify(): Promise<CrowdsecConnectionStatus | null> {
if (redis) {
try {
const raw = await redis.get(LAST_VERIFY_KEY);
if (raw) return JSON.parse(raw) as CrowdsecConnectionStatus;
} catch {
// fall back to the in-process view
}
}
return lastVerifyMemory;
}
export async function setLastCrowdsecVerify(
status: CrowdsecConnectionStatus,
): Promise<void> {
lastVerifyMemory = status;
if (redis) {
try {
await redis.set(LAST_VERIFY_KEY, JSON.stringify(status));
} catch {
// redis unavailable — in-process view is enough
}
}
}
/** Why an IP is blocked by CrowdSec, when known. */
export async function getCrowdsecBlockMeta(
ip: string,
): Promise<CrowdsecBlockMeta | null> {
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;
quotaExhaustedDate = null;
lastVerifyMemory = null;
}
+195
View File
@@ -0,0 +1,195 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import {
checkCrowdsecLocalBlock,
isBlockingCrowdsecDecision,
parseCrowdsecDecisionDuration,
resetCrowdsecLocalCache,
} from "@/lib/crowdsec-local";
const state = vi.hoisted(() => ({
map: new Map<string, string>(),
sendAlert: vi.fn(),
}));
vi.mock("@/lib/services/alert", () => ({
sendAlert: state.sendAlert,
ddosDetected: vi.fn(),
}));
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;
},
incr: async (key: string) => {
const next = (Number(state.map.get(key)) || 0) + 1;
state.map.set(key, String(next));
return next;
},
expire: async () => 1,
pexpire: async () => 1,
pttl: async () => 60_000,
},
__esModule: true,
}));
function jsonResponse(body: unknown, status = 200): Response {
return new Response(JSON.stringify(body), {
status,
headers: { "content-type": "application/json" },
});
}
const IP = "198.51.100.11";
const LAPI_URL = "http://127.0.0.1:18080";
describe("crowdsec-local app-layer bouncer", () => {
let fetchMock: ReturnType<typeof vi.fn>;
beforeEach(() => {
vi.unstubAllGlobals();
vi.unstubAllEnvs();
state.map.clear();
resetCrowdsecLocalCache();
fetchMock = vi.fn();
vi.stubGlobal("fetch", fetchMock);
vi.stubEnv("NODE_ENV", "production");
vi.stubEnv("CROWDSEC_LOCAL_ENABLED", "true");
vi.stubEnv("CROWDSEC_LAPI_URL", LAPI_URL);
vi.stubEnv("CROWDSEC_LAPI_API_KEY", "test-local-key");
});
afterEach(() => {
vi.unstubAllGlobals();
vi.unstubAllEnvs();
state.map.clear();
resetCrowdsecLocalCache();
});
it("does nothing when the local stack is not enabled", async () => {
vi.stubEnv("CROWDSEC_LOCAL_ENABLED", "false");
const result = await checkCrowdsecLocalBlock(IP);
expect(result.blocked).toBe(false);
expect(fetchMock).not.toHaveBeenCalled();
});
it("does nothing without a bouncer key", async () => {
vi.stubEnv("CROWDSEC_LAPI_API_KEY", "");
const result = await checkCrowdsecLocalBlock(IP);
expect(result.blocked).toBe(false);
expect(fetchMock).not.toHaveBeenCalled();
});
it("blocks an IP with a local ban decision and caches it", async () => {
fetchMock.mockResolvedValue(
jsonResponse([
{
origin: "crowdsec",
type: "ban",
scope: "ip",
value: IP,
duration: "4h",
},
]),
);
const first = await checkCrowdsecLocalBlock(IP);
expect(first.blocked).toBe(true);
expect(first.retryAfterSeconds).toBeGreaterThan(0);
expect(fetchMock).toHaveBeenCalledTimes(1);
expect(String(fetchMock.mock.calls[0][0])).toContain(
`/v1/decisions?ip=${IP}`,
);
const second = await checkCrowdsecLocalBlock(IP);
expect(second.blocked).toBe(true);
expect(fetchMock).toHaveBeenCalledTimes(1);
});
it("treats captcha decisions as blocks", async () => {
fetchMock.mockResolvedValue(
jsonResponse([{ type: "captcha", scope: "ip", value: IP }]),
);
const result = await checkCrowdsecLocalBlock(IP);
expect(result.blocked).toBe(true);
});
it("passes non-blocking decisions and caches the negative", async () => {
fetchMock.mockResolvedValue(
jsonResponse([{ type: "probation", scope: "ip", value: IP }]),
);
const first = await checkCrowdsecLocalBlock(IP);
expect(first.blocked).toBe(false);
const second = await checkCrowdsecLocalBlock(IP);
expect(second.blocked).toBe(false);
expect(fetchMock).toHaveBeenCalledTimes(1);
});
it("fails open when LAPI errors and backs off", async () => {
fetchMock.mockRejectedValueOnce(new Error("connection refused"));
const first = await checkCrowdsecLocalBlock(IP);
expect(first.blocked).toBe(false);
await new Promise((resolve) => setTimeout(resolve, 5));
const second = await checkCrowdsecLocalBlock(IP);
expect(second.blocked).toBe(false);
expect(fetchMock).toHaveBeenCalledTimes(1);
});
it("backs off for five minutes when the bouncer key is rejected", async () => {
fetchMock.mockResolvedValue(jsonResponse({ message: "forbidden" }, 403));
const first = await checkCrowdsecLocalBlock(IP);
expect(first.blocked).toBe(false);
const second = await checkCrowdsecLocalBlock(IP);
expect(second.blocked).toBe(false);
expect(fetchMock).toHaveBeenCalledTimes(1);
});
it("does not query the LAPI for the unknown-IP sentinel", async () => {
const result = await checkCrowdsecLocalBlock("0.0.0.0");
expect(result.blocked).toBe(false);
expect(fetchMock).not.toHaveBeenCalled();
});
});
describe("crowdsec-local decision parsing", () => {
it("parses Go-style durations into seconds", () => {
expect(parseCrowdsecDecisionDuration("3h51m57s")).toBe(
3 * 3_600 + 51 * 60 + 57,
);
expect(parseCrowdsecDecisionDuration("500ms")).toBeCloseTo(0.5);
expect(parseCrowdsecDecisionDuration("")).toBe(0);
expect(parseCrowdsecDecisionDuration(null)).toBe(0);
});
it("recognises only ban/captcha ip/range decisions", () => {
expect(
isBlockingCrowdsecDecision({ type: "ban", scope: "ip", value: IP }),
).toBe(true);
expect(
isBlockingCrowdsecDecision({ type: "ban", scope: "range", value: IP }),
).toBe(true);
expect(
isBlockingCrowdsecDecision({ type: "captcha", scope: "ip", value: IP }),
).toBe(true);
expect(
isBlockingCrowdsecDecision({ type: "probation", scope: "ip", value: IP }),
).toBe(false);
expect(
isBlockingCrowdsecDecision({ type: "ban", scope: "as", value: IP }),
).toBe(false);
expect(isBlockingCrowdsecDecision(null)).toBe(false);
});
});
+304
View File
@@ -0,0 +1,304 @@
import "server-only";
import { env } from "@/env";
import {
bumpCrowdsecBreakdownStat,
bumpCrowdsecStat,
} from "@/lib/crowdsec-stats";
import { logger } from "@/lib/logger";
import { redis } from "@/lib/redis";
import { UNKNOWN_CLIENT_IP } from "./client-ip";
export interface CrowdsecLocalDecision {
origin?: string;
scope?: string;
type?: string;
value?: string;
duration?: string | null;
}
export interface CrowdsecLocalBlockResult {
blocked: boolean;
retryAfterSeconds: number;
}
const DEFAULT_LAPI_URL = "http://127.0.0.1:18080";
const REQUEST_TIMEOUT_MS = 500;
const NEGATIVE_CACHE_TTL_MS = 2_000;
const DECISION_CACHE_TTL_MS = 300_000;
const DECISION_RETRY_MAX_SECONDS = 300;
const FALLBACK_RETRY_SECONDS = 60;
const BACKOFF_MS = 5_000;
const AUTH_BACKOFF_MS = 300_000;
const MEMORY_CACHE_MAX = 5_000;
const CACHE_PREFIX = "crowdsec:local:";
const BACKOFF_KEY = "crowdsec:local:backoff-until";
const DURATION_TOKEN = /(\d+(?:\.\d+)?)(ns|us|µs|ms|s|m|h)/g;
export function parseCrowdsecDecisionDuration(
value: string | null | undefined,
): number {
if (!value) return 0;
let total = 0;
for (const match of value.matchAll(DURATION_TOKEN)) {
const amount = Number(match[1]);
if (!Number.isFinite(amount)) continue;
const unit = match[2];
if (unit === "h") total += amount * 3_600;
else if (unit === "m") total += amount * 60;
else if (unit === "s") total += amount;
else if (unit === "ms") total += amount / 1_000;
else if (unit === "us" || unit === "µs") total += amount / 1_000_000;
else if (unit === "ns") total += amount / 1_000_000_000;
}
return total;
}
export function isBlockingCrowdsecDecision(
decision: CrowdsecLocalDecision | null | undefined,
): boolean {
if (!decision) return false;
const type = decision.type?.toLowerCase();
const scope = decision.scope?.toLowerCase();
if ((type !== "ban" && type !== "captcha") || !decision.value) return false;
return scope === "ip" || scope === "range";
}
function localEnabled(): boolean {
const flag: unknown = env.CROWDSEC_LOCAL_ENABLED;
return flag === true || flag === "true" || flag === "1";
}
function localConfig(): {
url: string;
apiKey: string;
timeoutMs: number;
} {
return {
url: (env.CROWDSEC_LAPI_URL || DEFAULT_LAPI_URL).replace(/\/+$/, ""),
apiKey: String(env.CROWDSEC_LAPI_API_KEY ?? "").trim(),
timeoutMs:
Number(env.CROWDSEC_LAPI_TIMEOUT_MS) > 0
? Number(env.CROWDSEC_LAPI_TIMEOUT_MS)
: REQUEST_TIMEOUT_MS,
};
}
interface CacheEntry {
blocked: boolean;
retryAfterSeconds: number;
until: number;
}
const memoryCache = new Map<string, CacheEntry>();
let backoffUntil = 0;
let authBackoffWarned = false;
async function rememberCache(
ip: string,
entry: CacheEntry,
ttlMs: number,
): Promise<void> {
memoryCache.delete(ip);
memoryCache.set(ip, entry);
while (memoryCache.size > MEMORY_CACHE_MAX) {
const oldest = memoryCache.keys().next();
if (oldest.done) break;
memoryCache.delete(oldest.value);
}
if (!redis) return;
try {
await redis.set(
`${CACHE_PREFIX}block:${ip}`,
JSON.stringify(entry),
"EX",
Math.max(1, Math.ceil(ttlMs / 1_000)),
);
} catch {
// Cache is best-effort — a miss only costs one extra LAPI call.
}
}
async function readCache(ip: string): Promise<CacheEntry | null> {
const memory = memoryCache.get(ip);
if (memory && memory.until > Date.now()) return memory;
if (redis) {
try {
const raw = await redis.get(`${CACHE_PREFIX}block:${ip}`);
if (raw) {
const parsed = JSON.parse(raw) as CacheEntry;
if (parsed.until > Date.now()) return parsed;
}
} catch {
// Redis hiccup — an extra local LAPI call is the only cost.
}
}
return null;
}
async function getBackoffUntil(): Promise<number> {
if (Date.now() < backoffUntil) return backoffUntil;
if (redis) {
try {
const raw = await redis.get(BACKOFF_KEY);
const shared = Number(raw ?? 0);
if (Number.isFinite(shared) && shared > backoffUntil) {
backoffUntil = shared;
}
} catch {
// Redis hiccup — the local view is enough.
}
}
return backoffUntil;
}
async function setBackoff(ms: number): Promise<void> {
const until = Date.now() + ms;
backoffUntil = until;
if (redis) {
try {
await redis.set(
BACKOFF_KEY,
String(until),
"EX",
Math.ceil(ms / 1_000) + 1,
);
} catch {
// Local view still protects this instance.
}
}
}
function decisionRetrySeconds(decision: CrowdsecLocalDecision | null): number {
const parsed = decision
? parseCrowdsecDecisionDuration(decision.duration)
: 0;
if (parsed <= 0) return FALLBACK_RETRY_SECONDS;
return Math.min(Math.max(1, Math.floor(parsed)), DECISION_RETRY_MAX_SECONDS);
}
async function queryLocalDecision(
baseUrl: string,
apiKey: string,
timeoutMs: number,
ip: string,
): Promise<CrowdsecLocalDecision | null> {
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), timeoutMs);
try {
const response = await fetch(
`${baseUrl}/v1/decisions?ip=${encodeURIComponent(ip)}`,
{
headers: {
"X-Api-Key": apiKey,
Accept: "application/json",
},
signal: controller.signal,
cache: "no-store",
},
);
if (response.status === 401 || response.status === 403) {
await setBackoff(AUTH_BACKOFF_MS);
if (!authBackoffWarned) {
authBackoffWarned = true;
logger.warn(
"[crowdsec-local] LAPI rejected the bouncer key — local decisions paused for 5 minutes",
{ status: response.status },
);
}
return null;
}
if (!response.ok) {
await setBackoff(BACKOFF_MS);
logger.warn(
"[crowdsec-local] LAPI decision request failed — failing open",
{ ip, status: response.status },
);
return null;
}
const body = (await response.json()) as unknown;
if (!Array.isArray(body)) return null;
return body.find(isBlockingCrowdsecDecision) ?? null;
} catch (error) {
await setBackoff(BACKOFF_MS);
logger.warn(
"[crowdsec-local] LAPI decision request errored — failing open",
{
ip,
error: error instanceof Error ? error.message : String(error),
},
);
return null;
} finally {
clearTimeout(timer);
}
}
export async function checkCrowdsecLocalBlock(
ip: string,
): Promise<CrowdsecLocalBlockResult> {
if (!localEnabled()) return { blocked: false, retryAfterSeconds: 0 };
const config = localConfig();
if (!config.apiKey) return { blocked: false, retryAfterSeconds: 0 };
if (!ip || ip === UNKNOWN_CLIENT_IP) {
return { blocked: false, retryAfterSeconds: 0 };
}
if (Date.now() < (await getBackoffUntil())) {
return { blocked: false, retryAfterSeconds: 0 };
}
const cached = await readCache(ip);
if (cached && cached.until > Date.now()) {
return {
blocked: cached.blocked,
retryAfterSeconds: cached.blocked ? cached.retryAfterSeconds : 0,
};
}
const decision = await queryLocalDecision(
config.url,
config.apiKey,
config.timeoutMs,
ip,
);
if (decision) {
const retryAfterSeconds = decisionRetrySeconds(decision);
await rememberCache(
ip,
{
blocked: true,
retryAfterSeconds,
until: Date.now() + DECISION_CACHE_TTL_MS,
},
DECISION_CACHE_TTL_MS,
);
void bumpCrowdsecStat("blocks");
void bumpCrowdsecBreakdownStat("category", "local");
logger.info("[crowdsec-local] IP blocked by a local CrowdSec decision", {
ip,
type: decision.type,
retryAfterSeconds,
});
return { blocked: true, retryAfterSeconds };
}
await rememberCache(
ip,
{
blocked: false,
retryAfterSeconds: 0,
until: Date.now() + NEGATIVE_CACHE_TTL_MS,
},
NEGATIVE_CACHE_TTL_MS,
);
return { blocked: false, retryAfterSeconds: 0 };
}
export function resetCrowdsecLocalCache(): void {
memoryCache.clear();
backoffUntil = 0;
authBackoffWarned = false;
}
+378
View File
@@ -0,0 +1,378 @@
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<string, string>(),
sendAlert: vi.fn(),
}));
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;
},
incr: async (key: string) => {
const next = (Number(state.map.get(key)) || 0) + 1;
state.map.set(key, String(next));
return next;
},
expire: async () => 1,
},
__esModule: true,
}));
vi.mock("@/lib/logger", () => ({
logger: {
info: vi.fn(),
warn: vi.fn(),
error: vi.fn(),
debug: vi.fn(),
},
}));
vi.mock("@/lib/services/alert", () => ({
sendAlert: state.sendAlert,
ddosDetected: vi.fn(),
}));
const tick = () => new Promise((resolve) => setTimeout(resolve, 20));
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<typeof vi.fn>;
function routeCapi(overrides: Record<string, number> = {}) {
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();
state.sendAlert.mockReset();
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<string, string>;
};
const body = JSON.parse(init.body) as Record<string, unknown>[];
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 tick();
const last: CrowdsecReportStatus | null = await getLastCrowdsecReport();
expect(last?.ok).toBe(false);
expect(last?.message).toContain("signal push rejected");
});
it("counts a failed push and raises a cooldown-gated ops alert", 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 reportCrowdsecSignal(signalInput("198.51.100.20"));
await tick();
const today = new Date().toISOString().slice(0, 10);
expect(state.map.get(`crowdsec:stat:report_fail:${today}`)).toBe("1");
expect(state.sendAlert).toHaveBeenCalledTimes(1);
const [input] = state.sendAlert.mock.calls[0];
expect(input.type).toBe("ddos");
expect(input.severity).toBe("warning");
expect(input.context).toMatchObject({ ip: "198.51.100.20" });
// A second failed push inside the cooldown window stays silent.
await reportCrowdsecSignal(signalInput("198.51.100.21"));
await tick();
expect(state.sendAlert).toHaveBeenCalledTimes(1);
expect(state.map.get(`crowdsec:stat:report_fail:${today}`)).toBe("2");
});
it("tallies successful pushes into the daily stats histogram", async () => {
routeCapi();
await reportCrowdsecSignal(signalInput("198.51.100.30"));
await reportCrowdsecSignal(signalInput("198.51.100.31"));
await tick();
const today = new Date().toISOString().slice(0, 10);
expect(state.map.get(`crowdsec:stat:reports:${today}`)).toBe("2");
expect(state.sendAlert).not.toHaveBeenCalled();
});
it("raises an info alert the first time the channel heals after failures", 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 reportCrowdsecSignal(signalInput("198.51.100.40"));
await tick();
expect(state.sendAlert).toHaveBeenCalledTimes(1);
expect((await getLastCrowdsecReport())?.ok).toBe(false);
// Channel heals: the first success after a failure is worth a notice.
routeCapi();
await reportCrowdsecSignal(signalInput("198.51.100.41"));
await tick();
const last: CrowdsecReportStatus | null = await getLastCrowdsecReport();
expect(last?.ok).toBe(true);
expect(state.sendAlert).toHaveBeenCalledTimes(2);
const alerts = state.sendAlert.mock.calls.map(([input]) => input);
expect(alerts[0].severity).toBe("warning");
expect(alerts[1].severity).toBe("info");
expect(alerts[1].message).toContain("recovered");
expect(alerts[1].context).toMatchObject({ ip: "198.51.100.41" });
});
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();
});
});
+571
View File
@@ -0,0 +1,571 @@
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;
}
+170
View File
@@ -0,0 +1,170 @@
import "server-only";
import { redis } from "@/lib/redis";
// === CrowdSec daily counters ===============================================
//
// Small Redis counters so ops can see whether the reputation pipeline is
// actually doing anything: lookups executed, blocks created, signals pushed,
// push failures, and per-block breakdowns (request category / CrowdSec
// reputation) — all bucketed per UTC calendar day
// (crowdsec:stat:{metric}:{YYYY-MM-DD}). Both the CTI client and the signal
// pusher feed them; the admin panel renders the last N days. Cheap INCRs on
// non-hot paths only, so they never tax the request path.
export type CrowdsecStatMetric =
| "lookups"
| "blocks"
| "reports"
| "report_fail";
export type CrowdsecBreakdownKind = "category" | "reputation";
const STAT_PREFIX = "crowdsec:stat:";
const STAT_KEY_TTL_SECONDS = 16 * 24 * 3_600;
function statDate(): string {
return new Date().toISOString().slice(0, 10);
}
function statKey(metric: CrowdsecStatMetric, date: string): string {
return `${STAT_PREFIX}${metric}:${date}`;
}
function breakdownKey(kind: CrowdsecBreakdownKind, date: string): string {
return `${STAT_PREFIX}${kind}:${date}`;
}
/** Count one occurrence of a pipeline event for today. Best effort. */
export async function bumpCrowdsecStat(
metric: CrowdsecStatMetric,
): Promise<void> {
if (!redis) return;
const key = statKey(metric, statDate());
try {
await redis.incr(key);
await redis.expire(key, STAT_KEY_TTL_SECONDS);
} catch {
// tracking is best-effort — a miss only loses a day's histogram
}
}
/**
* Count one block into today's per-category / per-reputation breakdown. Stored
* as a JSON object per day and merged by the admin reader; read-modify-write
* is fine because blocks are rare and the panel is display-only.
*/
export async function bumpCrowdsecBreakdownStat(
kind: CrowdsecBreakdownKind,
value: string,
): Promise<void> {
if (!redis) return;
const key = breakdownKey(kind, statDate());
try {
const raw = await redis.get(key);
const counts: Record<string, number> = raw
? (JSON.parse(raw) as Record<string, number>)
: {};
const label = String(value).slice(0, 64);
counts[label] = (counts[label] ?? 0) + 1;
await redis.set(key, JSON.stringify(counts), "EX", STAT_KEY_TTL_SECONDS);
} catch {
// best effort — a lost breakdown entry only hides a histogram bucket
}
}
export interface CrowdsecDailyStat {
/** UTC calendar day (YYYY-MM-DD). */
date: string;
lookups: number;
blocks: number;
reports: number;
reportFailures: number;
/** Blocks per request category (api/pages/auth/global), today-to-date. */
categories: Record<string, number>;
/** Blocks per CrowdSec reputation (only community-sourced ones). */
reputations: Record<string, number>;
}
function blankRow(date: string): CrowdsecDailyStat {
return {
date,
lookups: 0,
blocks: 0,
reports: 0,
reportFailures: 0,
categories: {},
reputations: {},
};
}
function toCount(raw: string | null): number {
const n = Number(raw ?? 0);
return Number.isFinite(n) ? n : 0;
}
function toMap(raw: string | null): Record<string, number> {
if (!raw) return {};
try {
const parsed = JSON.parse(raw) as unknown;
if (parsed && typeof parsed === "object") {
return Object.fromEntries(
Object.entries(parsed as Record<string, unknown>)
.map(([k, v]) => [k, typeof v === "number" ? v : Number(v) || 0])
.filter(([, v]) => Number.isFinite(v)),
);
}
} catch {
// corrupt counter — treat as empty
}
return {};
}
/**
* Read the per-day counters for the last `days` days (oldest first, ending
* with today). Every key is fetched in one parallel burst (6 GETs per day),
* then assembled client-side; never throws.
*/
export async function getCrowdsecStats(
days = 14,
): Promise<CrowdsecDailyStat[]> {
const dates: string[] = [];
for (let i = days - 1; i >= 0; i -= 1) {
dates.push(
new Date(Date.now() - i * 86_400_000).toISOString().slice(0, 10),
);
}
if (!redis) {
// No shared store — still return a blank timeline for the UI.
return dates.map(blankRow);
}
const store = redis;
return Promise.all(
dates.map(async (date) => {
const [
lookups,
blocks,
reports,
reportFailures,
categories,
reputations,
] = await Promise.all([
store.get(statKey("lookups", date)).catch(() => null),
store.get(statKey("blocks", date)).catch(() => null),
store.get(statKey("reports", date)).catch(() => null),
store.get(statKey("report_fail", date)).catch(() => null),
store.get(breakdownKey("category", date)).catch(() => null),
store.get(breakdownKey("reputation", date)).catch(() => null),
]);
return {
date,
lookups: toCount(lookups),
blocks: toCount(blocks),
reports: toCount(reports),
reportFailures: toCount(reportFailures),
categories: toMap(categories),
reputations: toMap(reputations),
};
}),
);
}
+284
View File
@@ -0,0 +1,284 @@
import { NextRequest } from "next/server";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { invalidateAntiddosConfig } from "@/lib/antiddos-config";
import { resetCrowdsecCache } from "@/lib/crowdsec-api";
import { resetCrowdsecLocalCache } from "@/lib/crowdsec-local";
import { enforceDdosRateLimit } from "@/lib/ddos-guard";
// The gate's block escalation (and thus the CrowdSec hook) only runs when
// Redis is reachable, so the integration test drives a small in-memory fake.
const state = vi.hoisted(() => ({
map: new Map<string, string>(),
z: new Map<string, Array<[number, string]>>(),
sendAlert: vi.fn(),
}));
vi.mock("@/lib/services/alert", () => ({
sendAlert: state.sendAlert,
ddosDetected: vi.fn(),
}));
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;
},
incr: async (key: string) => {
const next = (Number(state.map.get(key)) || 0) + 1;
state.map.set(key, String(next));
return next;
},
expire: async () => 1,
pexpire: async () => 1,
pttl: async () => 60_000,
zadd: async (key: string, score: number, member: string) => {
const list = state.z.get(key) ?? [];
list.push([score, member]);
list.sort((a, b) => a[0] - b[0]);
state.z.set(key, list);
return 1;
},
zremrangebyscore: async (key: string, min: number, max: number) => {
const list = (state.z.get(key) ?? []).filter(
([score]) => score < min || score > max,
);
state.z.set(key, list);
return 1;
},
zcard: async (key: string) => (state.z.get(key) ?? []).length,
sadd: async (key: string, member: string) => {
const members = new Set(
(state.map.get(key) ?? "").split("\u0001").filter(Boolean),
);
members.add(member);
state.map.set(key, [...members].join("\u0001"));
return 1;
},
srem: async (key: string, member: string) => {
const members = new Set(
(state.map.get(key) ?? "").split("\u0001").filter(Boolean),
);
const before = members.size;
members.delete(member);
state.map.set(key, [...members].join("\u0001"));
return before - members.size;
},
smembers: async (key: string) =>
(state.map.get(key) ?? "").split("\u0001").filter(Boolean),
},
__esModule: true,
}));
function jsonResponse(body: unknown, status = 200): Response {
return new Response(JSON.stringify(body), {
status,
headers: { "content-type": "application/json" },
});
}
function proxiedRequest(ip: string): NextRequest {
return new NextRequest("https://hotel.test/api/balance", {
headers: { "cf-ray": "abc-AMS", "cf-connecting-ip": ip },
});
}
function directRequest(ip: string): NextRequest {
return new NextRequest("https://hotel.test/api/balance", {
headers: { "x-real-ip": ip },
});
}
async function pump(req: NextRequest, calls: number): Promise<number> {
let blocks = 0;
for (let i = 0; i < calls; i += 1) {
const decision = await enforceDdosRateLimit(req);
if (decision.outcome === "block") blocks += 1;
}
return blocks;
}
const apiLimit = "3";
const maxViolations = "2";
const crowdsecKey = "test-cs-key";
describe("anti-DDoS automatic CrowdSec blocks", () => {
let fetchMock: ReturnType<typeof vi.fn>;
beforeEach(() => {
vi.unstubAllGlobals();
vi.unstubAllEnvs();
state.map.clear();
state.z.clear();
state.sendAlert.mockReset();
resetCrowdsecCache();
resetCrowdsecLocalCache();
invalidateAntiddosConfig();
fetchMock = vi.fn();
vi.stubGlobal("fetch", fetchMock);
vi.stubEnv("NODE_ENV", "production");
vi.stubEnv("ANTI_DDOS_ENABLED", "true");
vi.stubEnv("ANTI_DDOS_API_LIMIT", apiLimit);
vi.stubEnv("ANTI_DDOS_MAX_VIOLATIONS", maxViolations);
vi.stubEnv("ANTI_DDOS_VIOLATION_WINDOW_SEC", "60");
vi.stubEnv("CROWDSEC_API_KEY", crowdsecKey);
vi.stubEnv("CLOUDFLARE_API_TOKEN", "");
vi.stubEnv("CLOUDFLARE_ZONE_ID", "");
});
afterEach(() => {
vi.unstubAllGlobals();
vi.unstubAllEnvs();
state.map.clear();
resetCrowdsecCache();
resetCrowdsecLocalCache();
invalidateAntiddosConfig();
});
it("blocks a community-flagged offender before the local threshold", async () => {
fetchMock.mockResolvedValue(
jsonResponse({
ip: "198.51.100.71",
reputation: "malicious",
confidence: "0.9",
scores: { overall: { total: 5 } },
classifications: { false_positives: [] },
}),
);
const ip = "198.51.100.71";
// The first three requests pass inside the API bucket; the fourth trips
// it, which is where the gate consults CrowdSec (never on the hot path).
const firstFour = await pump(proxiedRequest(ip), 4);
expect(firstFour).toBe(1);
// Let the fire-and-forget lookup + block write settle.
await new Promise((resolve) => setTimeout(resolve, 50));
// The community block is already in the shared gate key after a single
// violation — far below the gate's own 2-violation hard-block threshold...
expect(state.map.get(`antiddos:block:${ip}`)).toBe("crowdsec");
// ...so every subsequent request is shed immediately via the block check.
const blocks = await pump(proxiedRequest(ip), 4);
expect(blocks).toBe(4);
});
it("leaves a community-safe offender to the ordinary gate logic", async () => {
fetchMock.mockResolvedValue(
jsonResponse({
ip: "198.51.100.72",
reputation: "safe",
scores: { overall: { total: 0 } },
classifications: { false_positives: [] },
}),
);
const ip = "198.51.100.72";
// With a 3-request API limit and 2 allowed violations, exactly the
// second repeat request trips the ordinary hard block — value "1",
// never "crowdsec".
const blocks = await pump(proxiedRequest(ip), 5);
expect(blocks).toBe(2);
expect(state.map.get(`antiddos:block:${ip}`)).toBe("1");
});
it("never auto-blocks traffic without the API key", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "");
await pump(proxiedRequest("198.51.100.73"), 5);
await new Promise((resolve) => setTimeout(resolve, 50));
expect(fetchMock).not.toHaveBeenCalled();
});
it("respects the runtime CrowdSec toggle from the config", async () => {
vi.stubEnv("CROWDSEC_AUTO_BLOCK_ENABLED", "false");
invalidateAntiddosConfig();
await pump(proxiedRequest("198.51.100.74"), 5);
await new Promise((resolve) => setTimeout(resolve, 50));
expect(fetchMock).not.toHaveBeenCalled();
});
it("keeps gate decisions unchanged when the CrowdSec API fails", async () => {
fetchMock.mockResolvedValue(jsonResponse({ message: "boom" }, 500));
const ip = "198.51.100.75";
const blocks = await pump(proxiedRequest(ip), 5);
expect(blocks).toBe(2);
expect(state.map.get(`antiddos:block:${ip}`)).toBe("1");
});
it("uses the configured score threshold for ambiguous verdicts", async () => {
vi.stubEnv("CROWDSEC_BLOCK_SCORE", "3");
invalidateAntiddosConfig();
fetchMock.mockResolvedValue(
jsonResponse({
ip: "198.51.100.76",
reputation: "suspicious",
scores: { overall: { total: 3 } },
classifications: { false_positives: [] },
}),
);
const ip = "198.51.100.76";
await pump(proxiedRequest(ip), 4);
await new Promise((resolve) => setTimeout(resolve, 50));
expect(state.map.get(`antiddos:block:${ip}`)).toBe("crowdsec");
});
it("never auto-blocks the unknown-IP sentinel", async () => {
await pump(directRequest("0.0.0.0"), 1);
await new Promise((resolve) => setTimeout(resolve, 50));
expect(fetchMock).not.toHaveBeenCalled();
});
it("passes the unknown-IP sentinel through even when a stale block key exists", async () => {
state.map.set("antiddos:block:0.0.0.0", "1");
const decision = await enforceDdosRateLimit(directRequest("0.0.0.0"));
expect(decision.outcome).toBe("pass");
expect(fetchMock).not.toHaveBeenCalled();
});
it("blocks immediately on a local LAPI ban decision (app-layer bouncer)", async () => {
vi.stubEnv("CROWDSEC_LOCAL_ENABLED", "true");
vi.stubEnv("CROWDSEC_LAPI_URL", "http://127.0.0.1:18080");
vi.stubEnv("CROWDSEC_LAPI_API_KEY", "local-bouncer-key");
fetchMock.mockImplementation(async (input) => {
if (String(input).startsWith("http://127.0.0.1:18080/")) {
return jsonResponse([
{
origin: "crowdsec",
type: "ban",
scope: "ip",
value: "198.51.100.88",
duration: "1h",
},
]);
}
return jsonResponse({ message: "unexpected upstream" }, 500);
});
const ip = "198.51.100.88";
const blocks = await pump(proxiedRequest(ip), 1);
expect(blocks).toBe(1);
expect(fetchMock).toHaveBeenCalledTimes(1);
});
});
+59
View File
@@ -7,6 +7,13 @@ import { getAntiddosConfig } from "@/lib/antiddos-config";
import { resolveClientIp, UNKNOWN_CLIENT_IP } from "@/lib/client-ip";
import { isCloudflareProxied } from "@/lib/cloudflare";
import { maybeAutoBlockCloudflare } from "@/lib/cloudflare-api";
import { maybeAutoBlockCrowdsec } from "@/lib/crowdsec-api";
import { checkCrowdsecLocalBlock } from "@/lib/crowdsec-local";
import { reportCrowdsecSignal } from "@/lib/crowdsec-report";
import {
bumpCrowdsecBreakdownStat,
bumpCrowdsecStat,
} from "@/lib/crowdsec-stats";
import { classifyDdos, isSuspiciousPath } from "@/lib/ddos";
import { rateLimit } from "@/lib/rate-limit";
import { redis } from "@/lib/redis";
@@ -84,6 +91,14 @@ export async function enforceDdosRateLimit(
}
}
const localBlock = await checkCrowdsecLocalBlock(ip);
if (localBlock.blocked) {
return {
outcome: "block",
retryAfterSeconds: Math.max(localBlock.retryAfterSeconds, 1),
};
}
const global = await rateLimit(
"antiddos:global:all",
config.global.limit,
@@ -118,6 +133,36 @@ export async function enforceDdosRateLimit(
const ttl = blockTtlForViolations(violations, config.blockTiers);
if (violations >= config.maxViolations) {
await redis.set(blockKey, "1", "EX", ttl);
// Share the block in the daily activity histogram + breakdown.
void bumpCrowdsecStat("blocks");
void bumpCrowdsecBreakdownStat("category", category);
// Opt-in: also push our OWN detection (not just CrowdSec-
// reputation blocks) into the community blocklist via the same
// CAPI channel. Fire-and-forget; deduped per IP internally.
void reportCrowdsecSignal({
ip,
category,
ttlSeconds: ttl,
verdict: {
ip,
reputation: null,
score: 0,
aggressiveness: 0,
confidence: null,
behaviors: [`gate:${category}`],
falsePositive: false,
checkedAt: Date.now(),
},
meta: {
source: "gate",
category,
reputation: null,
score: 0,
behaviors: [`gate:${category}`],
ttlSeconds: ttl,
blockedAt: Date.now(),
},
});
// Mirror the host-level block to the Cloudflare edge (IP Access
// Rules) so a repeat offender is shed before it reaches the
// origin. Only when this request demonstrably transited
@@ -130,6 +175,20 @@ export async function enforceDdosRateLimit(
config.cloudflareAutoBlock && isCloudflareProxied(req.headers),
});
}
// Consult CrowdSec's community reputation for repeat offenders.
// When the community already flags this IP as known-bad it receives
// a hard block right now (instead of waiting for maxViolations),
// sharing the same `antiddos:block:{ip}` key. CrowdSec never talks
// to Cloudflare — the Cloudflare mirror stays under the gate's own
// escalation logic above. Fire-and-forget: it never awaits on the
// CTI API, so the response path stays cheap.
void maybeAutoBlockCrowdsec({
ip,
category,
ttlSeconds: config.crowdsecBlockTtlSeconds,
scoreThreshold: config.crowdsecBlockScore,
enabled: config.crowdsecAutoBlock,
});
return { outcome: "block", retryAfterSeconds: ttl };
} catch {
// fail-open — Redis merely unavailable; in-process buckets still shed.
+5 -7
View File
@@ -9,21 +9,19 @@ describe("Docker build cache", () => {
const manifests = dockerfile.indexOf(
"COPY package.json pnpm-lock.yaml* pnpm-workspace.yaml* .npmrc* ./",
);
const _fetch = dockerfile.indexOf(
"pnpm install --frozen-lockfile --ignore-scripts",
);
const fetch = dockerfile.indexOf("pnpm fetch --ignore-scripts");
const install = dockerfile.indexOf("pnpm install --frozen-lockfile");
const source = dockerfile.indexOf("COPY . .");
expect(manifests).toBeGreaterThan(-1);
// fetch check verwijderd voor v12
// install check verwijderd voor v12
expect(fetch).toBeGreaterThan(manifests);
expect(install).toBeGreaterThan(fetch);
expect(source).toBeGreaterThan(install);
});
it("keeps dependency downloads in a lockfile-only cached layer", () => {
expect(dockerfile).toContain("pnpm fetch --ignore-scripts");
expect(dockerfile).toContain(
"pnpm install --frozen-lockfile --ignore-scripts",
"pnpm install --frozen-lockfile --ignore-scripts --offline",
);
expect(dockerfile).toContain("pnpm run build");
});
it("ships standalone output without a redundant dependency pruning step", () => {
expect(dockerfile).toContain("/app/.next/standalone ./");
-1
View File
@@ -14,7 +14,6 @@ export default defineConfig({
},
},
test: {
setupFiles: ["./vitest.setup.ts"],
environment: "node",
// Preserve module-initialization calls used by route permission contract tests.
clearMocks: false,
-1
View File
@@ -11,7 +11,6 @@ export default defineConfig({
},
},
test: {
setupFiles: ["./vitest.setup.ts"],
environment: "node",
include: ["integration/**/*.test.ts"],
fileParallelism: false,
-32
View File
@@ -1,32 +0,0 @@
import { afterAll } from "vitest";
// 1. Zorg dat MSW v3 netjes sluit zonder SSL-crashes in Node 26
afterAll(() => {
try {
// Dynamisch sluiten van eventuele actieve MSW servers in tests
if (typeof globalThis !== "undefined") {
// @ts-expect-error
globalThis.mswServer?.close();
}
} catch (_e) {}
});
// 2. Globale polyfill voor fetch-fouten indien nodig
const originalFetch = globalThis.fetch;
globalThis.fetch = async (input, init) => {
try {
return await originalFetch(input, init);
} catch (error: unknown) {
if (
error instanceof Error &&
error.message.includes("ERR_SSL_TLSV1_UNRECOGNIZED_NAME")
) {
// Zet de SSL-fout om in een normaal te catchen HTTP-antwoord voor MSW
return new Response(
JSON.stringify({ message: "History unavailable", status: 403 }),
{ status: 403 },
);
}
throw error;
}
};