From aa615a1141f340ca94ba49023128950eda22eec4 Mon Sep 17 00:00:00 2001 From: loki5512344 Date: Tue, 15 Sep 2026 23:55:17 +0200 Subject: [PATCH] fix: XDP unsafe loader + Redis sync; enforce 250-line/4-file layout limits MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - xdp: CString for if_nametoindex (was UB), real detach via prog fd (fabricated borrow_raw(-1) silently never detached), SAFETY comments, saturating expiry math; +5 unit tests - redis: real pubsub reconnect with exponential backoff (was sleep+return); KEYS -> SCAN in heartbeat sweep - ci: cargo test --all-features, repo-gates job — module size gate (scripts/check_module_size.sh, ratchet baseline) + default-secrets grep - refactor src/ to <=250 LOC/file, <=4 .rs/dir without behavior change; thin bins (rampart.rs 330 -> 6 LOC), new app/, subnet/, intel/, profile/, prefix/, challenge/, filter/, probe/, inventory/, metrics/, node/ - docs: TODO v5.0 (status refresh, new rules, findings backlog), README quickstart now matches real binaries - verify: fmt/clippy -D warnings/test --all-features (164 tests)/clang XDP green --- .github/workflows/ci.yml | 12 +- .module_size_baseline | 3 + Makefile | 13 +- README.md | 48 +-- TODO.md | 126 ++++-- scripts/check_default_secrets.sh | 63 +++ scripts/check_module_size.sh | 161 +++++++ src/app/mod.rs | 6 + src/app/runtime.rs | 229 ++++++++++ src/app/services.rs | 119 ++++++ src/bin/rampart-cli.rs | 8 +- src/bin/rampart-manager.rs | 6 +- src/bin/rampart.rs | 328 +------------- src/cli/commands/mod.rs | 4 +- src/cli/commands/{ => node}/drain.rs | 0 src/cli/commands/{ => node}/emergency.rs | 0 src/cli/commands/node/mod.rs | 4 + src/cli/commands/{ => node}/status.rs | 0 src/config/sections.rs | 404 ------------------ src/config/sections/detect.rs | 121 ++++++ src/config/sections/edge.rs | 135 ++++++ src/config/sections/mod.rs | 11 + src/config/sections/platform.rs | 148 +++++++ src/engine/challenge/difficulty.rs | 88 ++++ src/engine/challenge/mod.rs | 8 + src/engine/{challenge.rs => challenge/pow.rs} | 85 ---- src/engine/mod.rs | 7 +- src/engine/subnet/mod.rs | 3 + .../{subnet_monitor.rs => subnet/monitor.rs} | 0 .../{subnet_tracker.rs => subnet/tracker.rs} | 0 src/lib.rs | 1 + src/manager/api/{ => inventory}/health.rs | 0 src/manager/api/inventory/mod.rs | 4 + src/manager/api/{ => inventory}/nodes.rs | 0 src/manager/api/{ => inventory}/servers.rs | 0 src/manager/api/mod.rs | 4 +- src/manager/sync/heartbeat.rs | 116 ++++- src/store/redis.rs | 147 +++++-- src/traffic/{ => intel}/detector.rs | 0 src/traffic/{ => intel}/ewma.rs | 0 src/traffic/intel/mod.rs | 4 + src/traffic/{ => intel}/reputation.rs | 0 src/traffic/mod.rs | 12 +- src/traffic/prefix/key.rs | 91 ++++ src/traffic/prefix/mod.rs | 14 + src/traffic/{prefix.rs => prefix/stats.rs} | 99 +---- src/traffic/{ => profile}/alert.rs | 0 src/traffic/profile/mod.rs | 3 + src/traffic/{ => profile}/profiler.rs | 0 src/tui/{ => metrics}/fetch.rs | 0 src/tui/metrics/mod.rs | 3 + src/tui/{ => metrics}/prometheus.rs | 0 src/tui/mod.rs | 5 +- src/xdp/diagnostics.rs | 285 ------------ src/xdp/filter/attach.rs | 58 +++ src/xdp/filter/maps.rs | 43 ++ src/xdp/{filter.rs => filter/mod.rs} | 105 +++-- src/xdp/mod.rs | 13 +- src/xdp/probe/diagnostics/mod.rs | 13 + src/xdp/probe/diagnostics/report.rs | 130 ++++++ src/xdp/probe/diagnostics/types.rs | 154 +++++++ src/xdp/{ => probe}/metrics.rs | 0 src/xdp/probe/mod.rs | 16 + src/xdp/{ => probe}/stats.rs | 0 64 files changed, 2052 insertions(+), 1408 deletions(-) create mode 100644 .module_size_baseline create mode 100755 scripts/check_default_secrets.sh create mode 100755 scripts/check_module_size.sh create mode 100644 src/app/mod.rs create mode 100644 src/app/runtime.rs create mode 100644 src/app/services.rs rename src/cli/commands/{ => node}/drain.rs (100%) rename src/cli/commands/{ => node}/emergency.rs (100%) create mode 100644 src/cli/commands/node/mod.rs rename src/cli/commands/{ => node}/status.rs (100%) delete mode 100644 src/config/sections.rs create mode 100644 src/config/sections/detect.rs create mode 100644 src/config/sections/edge.rs create mode 100644 src/config/sections/mod.rs create mode 100644 src/config/sections/platform.rs create mode 100644 src/engine/challenge/difficulty.rs create mode 100644 src/engine/challenge/mod.rs rename src/engine/{challenge.rs => challenge/pow.rs} (70%) create mode 100644 src/engine/subnet/mod.rs rename src/engine/{subnet_monitor.rs => subnet/monitor.rs} (100%) rename src/engine/{subnet_tracker.rs => subnet/tracker.rs} (100%) rename src/manager/api/{ => inventory}/health.rs (100%) create mode 100644 src/manager/api/inventory/mod.rs rename src/manager/api/{ => inventory}/nodes.rs (100%) rename src/manager/api/{ => inventory}/servers.rs (100%) rename src/traffic/{ => intel}/detector.rs (100%) rename src/traffic/{ => intel}/ewma.rs (100%) create mode 100644 src/traffic/intel/mod.rs rename src/traffic/{ => intel}/reputation.rs (100%) create mode 100644 src/traffic/prefix/key.rs create mode 100644 src/traffic/prefix/mod.rs rename src/traffic/{prefix.rs => prefix/stats.rs} (67%) rename src/traffic/{ => profile}/alert.rs (100%) create mode 100644 src/traffic/profile/mod.rs rename src/traffic/{ => profile}/profiler.rs (100%) rename src/tui/{ => metrics}/fetch.rs (100%) create mode 100644 src/tui/metrics/mod.rs rename src/tui/{ => metrics}/prometheus.rs (100%) delete mode 100644 src/xdp/diagnostics.rs create mode 100644 src/xdp/filter/attach.rs create mode 100644 src/xdp/filter/maps.rs rename src/xdp/{filter.rs => filter/mod.rs} (66%) create mode 100644 src/xdp/probe/diagnostics/mod.rs create mode 100644 src/xdp/probe/diagnostics/report.rs create mode 100644 src/xdp/probe/diagnostics/types.rs rename src/xdp/{ => probe}/metrics.rs (100%) create mode 100644 src/xdp/probe/mod.rs rename src/xdp/{ => probe}/stats.rs (100%) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 362c17a..e795e1f 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -40,7 +40,17 @@ jobs: - uses: actions/checkout@v4 - uses: dtolnay/rust-toolchain@stable - uses: Swatinem/rust-cache@v2 - - run: cargo test + - run: cargo test --all-features + + repo-gates: + name: Repo — module size & default secrets + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - name: Module size gate + run: bash scripts/check_module_size.sh + - name: Default secrets gate + run: bash scripts/check_default_secrets.sh rust-deny: name: Rust — cargo-deny diff --git a/.module_size_baseline b/.module_size_baseline new file mode 100644 index 0000000..d6b61de --- /dev/null +++ b/.module_size_baseline @@ -0,0 +1,3 @@ +# Module size ratchet baseline — read by scripts/check_module_size.sh +# Oversized files: plain path. Overcrowded directories: path with trailing /. +# Regenerate after a refactor: bash scripts/check_module_size.sh --update-baseline diff --git a/Makefile b/Makefile index 82d4306..3796aeb 100644 --- a/Makefile +++ b/Makefile @@ -3,10 +3,10 @@ CARGO = cargo TARGET_DIR = target -.PHONY: all build release test check check-all fmt fmt-check clippy clean deny audit +.PHONY: all build release test test-all check check-all fmt fmt-check clippy clean deny audit .PHONY: ebpf ebpf-clean .PHONY: docker docker-build docker-up docker-down docker-logs -.PHONY: ci ci-full +.PHONY: ci ci-full repo-gates all: check test build @@ -27,6 +27,9 @@ check-all: test: $(CARGO) test +test-all: + $(CARGO) test --all-features + fmt: $(CARGO) fmt --all @@ -75,7 +78,11 @@ docker-logs: checkstyle: fmt-check clippy @echo "✓ Checkstyle passed (rustfmt + clippy)" -ci: checkstyle test build deny +repo-gates: + bash scripts/check_module_size.sh + bash scripts/check_default_secrets.sh + +ci: checkstyle test-all build deny repo-gates @echo "✓ CI passed" ci-full: ci ebpf diff --git a/README.md b/README.md index 0a5c334..21666d8 100644 --- a/README.md +++ b/README.md @@ -54,7 +54,7 @@ Protocol-specific logic lives in **modular protocol plugins**, so the same platf │ L7 handshake analysis · rate limit · HMAC · death-code patterns │ ├────────────────────────────────────────────────────────────────────┤ │ Layer 4: Protocol Plugins (feature crates) │ -│ minecraft (first plugin) · http (planned) · grpc (planned) │ +│ http (feature protocol-http) · grpc (planned) │ └────────────────────────────────────────────────────────────────────┘ Traffic Intel (EWMA thresholds, profiling, reputation) runs across all layers @@ -69,12 +69,13 @@ Attacker → [XDP/eBPF] → [PoW] → [Userspace Core] → [Plugin] → Your Ser | Component | Role | Stack | |-----------|------|-------| -| **rampart-core** | Edge engine: XDP loader, PoW challenge, L7 filtering, traffic intel | Rust (tokio, libbpf) + C (XDP) | +| **rampart** | Edge engine: XDP loader, PoW challenge, L7 filtering, traffic intel | Rust (tokio, libbpf) + C (XDP) | | **rampart-manager** | Management API + Redis sync | Rust (axum, redis) | | **rampart-cli** | CLI tool for operators | Rust (clap) | +| **rampart-tui** | Live metrics terminal dashboard, polls the Prometheus `/metrics` endpoint | Rust (ratatui) | | **protocol plugins** | Protocol-aware filtering as feature crates | Rust | -| ↳ `minecraft` | First plugin (MC handshake analysis) | Rust | -| ↳ `http`, `grpc` | Planned | Rust | +| ↳ `http` | HTTP/1.1 handler, compiled with the `protocol-http` feature | Rust | +| ↳ `grpc` | Planned | Rust | | **docs/kb** | Bilingual knowledge base: attack anatomy, defense levels, practice guides | Markdown | ### Performance @@ -92,14 +93,15 @@ Full benchmark suite in progress. ### Quick Start ```bash -# Build -cargo build --release +# Build (the HTTP protocol handler is a feature) +cargo build --release --features protocol-http -# Create config -rampart config init > /etc/rampart/config.toml +# Install the default config +sudo mkdir -p /etc/rampart +sudo cp deploy/config/edge.toml /etc/rampart/config.toml -# Run edge node -./target/release/rampart-core --config /etc/rampart/config.toml +# Run edge node (config path comes from $RAMPART_CONFIG, default /etc/rampart/config.toml) +RAMPART_CONFIG=/etc/rampart/config.toml ./target/release/rampart ``` ### Documentation @@ -118,8 +120,6 @@ rampart config init > /etc/rampart/config.toml - Stabilize the protocol plugin API - BPF hook modules for deep protocol parsing in XDP -- Terminal UI (ratatui TUI) -- HTTP protocol plugin --- @@ -153,7 +153,7 @@ Rampart фильтрует трафик на трёх уровнях до тог │ L7 handshake analysis · rate limit · HMAC · death-code паттерны │ ├────────────────────────────────────────────────────────────────────┤ │ Слой 4: Протокол-плагины (feature crates) │ -│ minecraft (первый плагин) · http (в планах) · grpc (в планах) │ +│ http (feature protocol-http) · grpc (в планах) │ └────────────────────────────────────────────────────────────────────┘ Traffic Intel (EWMA thresholds, профилирование, репутация) работает поперёк всех слоёв @@ -168,12 +168,13 @@ Rampart фильтрует трафик на трёх уровнях до тог | Компонент | Роль | Технологии | |-----------|------|------------| -| **rampart-core** | Edge-движок: XDP loader, PoW challenge, L7-фильтрация, traffic intel | Rust (tokio, libbpf) + C (XDP) | +| **rampart** | Edge-движок: XDP loader, PoW challenge, L7-фильтрация, traffic intel | Rust (tokio, libbpf) + C (XDP) | | **rampart-manager** | Management API + Redis sync | Rust (axum, redis) | | **rampart-cli** | CLI для операторов | Rust (clap) | +| **rampart-tui** | Терминальный дашборд live-метрик, опрашивает Prometheus `/metrics` | Rust (ratatui) | | **Протокол-плагины** | Протоколозависимая фильтрация в виде feature crates | Rust | -| ↳ `minecraft` | Первый плагин (анализ MC-handshake) | Rust | -| ↳ `http`, `grpc` | В планах | Rust | +| ↳ `http` | HTTP/1.1-обработчик, собирается с фичей `protocol-http` | Rust | +| ↳ `grpc` | В планах | Rust | | **docs/kb** | Двуязычная база знаний: анатомия атак, уровни защиты, практические руководства | Markdown | ### Производительность @@ -191,14 +192,15 @@ Rampart фильтрует трафик на трёх уровнях до тог ### Быстрый старт ```bash -# Сборка -cargo build --release +# Сборка (HTTP-обработчик собирается фичей) +cargo build --release --features protocol-http -# Создание конфига -rampart config init > /etc/rampart/config.toml +# Установка дефолтного конфига +sudo mkdir -p /etc/rampart +sudo cp deploy/config/edge.toml /etc/rampart/config.toml -# Запуск edge ноды -./target/release/rampart-core --config /etc/rampart/config.toml +# Запуск edge ноды (путь конфига берётся из $RAMPART_CONFIG, по умолчанию /etc/rampart/config.toml) +RAMPART_CONFIG=/etc/rampart/config.toml ./target/release/rampart ``` ### Документация @@ -217,8 +219,6 @@ rampart config init > /etc/rampart/config.toml - Стабилизация API протокол-плагинов - BPF hook модули для глубокого парсинга протоколов в XDP -- Терминальный интерфейс (ratatui TUI) -- HTTP протокол-плагин --- diff --git a/TODO.md b/TODO.md index fcf0f35..7ac7874 100644 --- a/TODO.md +++ b/TODO.md @@ -12,7 +12,10 @@ ### KISS - Не добавляй абстракцию до третьего повторения. -- **Функция ≤ 60 строк, модуль ≤ 300 строк** (жёсткий лимит; больше — декомпозиция). +- **Функция ≤ 60 строк, файл ≤ 250 строк, ≤ 4 .rs-файлов на папку** (жёсткий лимит; + больше — декомпозиция или группировка доменами с ре-экспортами в mod.rs). + Гейт: `scripts/check_module_size.sh` + ratchet-базлайн `.module_size_baseline` + (легаси — WARN, новые нарушения — FAIL; после рефакторинга `--update-baseline`). - Не используй generics где хватит `&str` и `Vec`. ### DRY @@ -44,19 +47,24 @@ ``` guard/ ├── Cargo.toml # ОДИН пакет rampart, features = ["protocol-http", ...] +├── scripts/ # check_module_size.sh (гейт 250/4), check_default_secrets.sh ├── src/ -│ ├── bin/{rampart, rampart-manager, rampart-cli}.rs -│ ├── engine/ # listener, tunnel (generic TCP proxy), challenge (PoW) +│ ├── bin/{rampart, rampart-manager, rampart-cli, rampart-tui}.rs # тонкие, логика в lib +│ ├── app/ # thin-main: runtime wiring, services (spawn-циклы) +│ ├── engine/ # listener, tunnel; challenge/ (pow+difficulty), subnet/ (tracker+monitor) │ ├── filter/ # blacklist, rate_limit, geo — trait Filter -│ ├── traffic/ # EWMA, detector, profiler, reputation, alert +│ ├── traffic/ # hook + prefix/ (key+stats) + intel/ (ewma,detector,reputation) + profile/ (profiler,alert) │ ├── store/ # redis (+ trait StateStore) -│ ├── manager/ # api/, auth/, sync/ -│ ├── cli/ # команды CLI -│ └── protocol/ # trait ProtocolHandler + registry (реализаций пока 0) +│ ├── manager/ # api/ (+ api/inventory/), auth/, sync/ +│ ├── cli/ # команды; commands/node/ (status,drain,emergency) +│ ├── config/ # config + sections/{edge,platform,detect} +│ ├── tui/ # app/state/ui + metrics/ (prometheus,fetch) +│ ├── protocol/ # trait ProtocolHandler + registry; http/ под фичей +│ └── xdp/ # filter/ (attach,maps), probe/ (diagnostics/,stats,metrics), globals, noop ├── xdp/ -│ ├── core/ # universal_filter.c + maps/stats/config/common.h +│ ├── core/ # universal_filter.c + maps/stats/config/common.h + prefix_stats/syn_challenge │ └── hooks/hook_api.h # контракт подключаемых BPF-протокол-хуков -├── tests/ # интеграционные +├── tests/ # интеграционные (вне гейта 4-файлов: cargo требует 1 файл = 1 бинарь) └── docs/ # kb/ (knowledge base) + research/ + ops-доки ``` @@ -78,46 +86,71 @@ guard/ ## 2. Ближайшие задачи (v0.3) -### Subnet-level detection (ботнет с ротацией IP) -- [ ] **XDP**: карта `prefix_stats` (LRU_HASH, ключ /24 v4 | /64 v6) — счётчики SYN/pps - per-префикс рядом с per-IP (референс: caddy-mitigator CIDR promotion, lnvps_fw carpet-bomb). -- [ ] **Detector**: префикс превышает порог при том что отдельные IP под лимитом - → распределённая атака → флаг подсети. -- [ ] **Мягкая эскалация для подсетей**: monitor → strict limits → challenge → блок. - Хард-бан /24 только через challenge (CGNAT: за одним /24 легитимно живут сотни людей). -- [ ] Блок самой подсети — уже умеем: `blacklist_map` это LPM trie (CIDR из коробки). +### Subnet-level detection (ботнет с ротацией IP) — ✅ готово (2026-08/09) +- [x] **XDP**: карта `prefix_stats` (LRU_HASH, ключ /24 v4 | /64 v6) — счётчики SYN/pps + per-префикс (xdp/core/prefix_stats.h, трафик-слой: src/traffic/prefix/). +- [x] **Detector**: превышение порога префиксом при IP под лимитом → флаг подсети + (src/traffic/intel/detector.rs, src/engine/subnet/). +- [x] **Мягкая эскалация**: monitor → strict_limit → challenge → block + (SubnetVerdict-лестница; блок /24 идёт через ban_cidr, не слепой hard-ban). +- [x] Блок подсети через `blacklist_map` (LPM trie) — `XdpFilter::ban_cidr`. -### Движок без протоколов — сделать полезным -- [ ] **Первый протокол-плагин**: `protocol-http` (feature) — минимальный HTTP/1.1 - handshake-анализ (request line, заголовки, размер), чтобы edge-нода заработала - для веб-сервисов. -- [ ] **TCP-proxy режим**: generic upstream forwarding за ProtocolHandler - (tunnel.rs уже generic — проверить интеграцию). -- [ ] **Fail-fast сообщение** при пустом registry — улучшить текст подсказки сборки. +### Движок без протоколов — сделать полезным — ✅ готово +- [x] **Первый протокол-плагин**: `protocol-http` — HTTP/1.1 request-head анализ + (src/protocol/http/, tests/http_protocol.rs). +- [x] **TCP-proxy режим**: tunnel.rs + ProtocolHandler интегрированы в listener/app wiring. +- [x] **Fail-fast сообщение** при пустом registry (`ProtocolRegistry::primary()`). -### Подключение мёртвого интеллекта (правило: «мёртвый код = баг») -- [ ] Layer Traffic Intel подключить в hot path: AttackDetector/IpReputation → - метрики + auto-ban (сейчас не вызывается). -- [ ] Blacklist: `clear_expired()` по таймеру. -- [ ] RateLimiter: TTL-эвикция idle bucket'ов + cap карты. +### Подключение мёртвого интеллекта (правило: «мёртвый код = баг») — ✅ готово +- [x] Traffic Intel в hot path: AttackDetector/IpReputation → метрики + auto-ban + (src/traffic/hook.rs, src/engine/tunnel.rs, src/app/services.rs). +- [x] Blacklist: `clear_expired()` по таймеру (src/app/). +- [x] RateLimiter: TTL-эвикция idle bucket'ов. -### Безопасность (перенос из аудита v0.3, актуальное) -- [ ] Rate limiter на login endpoint manager API (5/60с). -- [ ] JWT: валидация ролей/audience, secret ≥ 32 байт. -- [ ] Redis: `KEYS` → `SCAN`, reconnect pubsub-подписчика. +### Безопасность — ✅ кроме ролей +- [x] Rate limiter на login endpoint manager API (per-IP, tests в api/auth.rs). +- [x] JWT: audience-валидация, secret ≥ 32 байт fail-fast (rampart-manager.rs). +- [ ] JWT-роли: сейчас единственный hardcoded `admin`; RBAC-ролей нет — либо убрать поле, + либо делать роли (решение отложить до второго потребителя). +- [x] Redis: `KEYS` → `SCAN` (scan_options), pubsub-подписчик — честный reconnect + с экспоненциальным backoff (src/store/redis.rs). ### XDP -- [ ] Verifier-проверка на реальном ядре (в контейнере нет CAP_BPF — компиляция OK, - загрузка не проверялась). -- [ ] Rust loader (`src/xdp/`): пути к xdp/core/universal_filter.c, patch глобалов - G_* из config.toml, ringbuf events → blacklist. -- [ ] Smoke-test attach в CI (VM runner с CAP_BPF). +- [x] Verifier-проверка на реальном ядре (live-тест 2026-09: два бага RST-challenge найдены и + исправлены, d6bcae5). +- [x] Rust loader (`src/xdp/filter/`): open/load/attach, patch глобалов G_* из config.toml + (src/xdp/globals.rs), CString-safe if_nametoindex, реальный detach по fd прогрессы. +- [ ] Ringbuf events → blacklist: `drain_events()` сейчас только логирует debug; + нужен матчинг события → `ban_ip` (перенести из «loader», боится ошибок в data path). +- [ ] Smoke-test attach в CI (VM runner с CAP_BPF) + netns+veth харнесс — вывод №7 из + конкурентного анализа, приоритет v0.4. ### Документация -- [ ] docs/deployment.md, configuration.md, runbook.md — переписать под новую структуру - (сейчас упоминают старые крейты/MC). -- [ ] docs/kb/README.md — индекс KB со ссылками на все статьи. -- [ ] TUI (ratatui): live-метрики из Prometheus endpoint (planned, v0.4). +- [x] docs/deployment.md, configuration.md, runbook.md — старых крейтов/MC не упоминают (grep чисто). +- [x] docs/kb/README.md — индекс KB со ссылками. +- [x] TUI (ratatui): live-метрики из Prometheus endpoint — готово (src/tui/, 564613b). + +## 2a. Раунд 2026-09-15 — исправлено по итогам ревью + +- **unsafe-баги XDP-лоадера** (src/xdp/filter/): `&str.as_ptr()` без NUL в `if_nametoindex` → `CString`; + фиктивный `borrow_raw(-1)`-fd в `unload()` → реальный detach по fd прогрессы с propagatable ошибкой; + `unsafe impl Send` получил `// SAFETY:`; `saturating_mul` в expiry. +5 юнит-тестов. +- **Redis**: fake-reconnect pubsub (sleep+return) → настоящий reconnect-loop с backoff; + `KEYS` → `SCAN` в heartbeat. +- **CI**: `cargo test --all-features` (раньше тесты фич не гонялись); job `repo-gates`: + гейт 250/4 + grep дефолтных секретов; Makefile `ci` обновлён. +- **README**: несуществующие `rampart-core`/`rampart config init` → реальные команды; + таблица компонентов/плагинов приведена к коду (http готов, minecraft удалён, добавлен rampart-tui). +- **Рефакторинг под новые лимиты 250/4** без изменения поведения: app/, challenge/, subnet/, + prefix/, intel/, profile/, commands/node/, api/inventory/, tui/metrics/, filter/, probe/; + bin-файлы тонкие (rampart.rs: 330 → 6 строк). + +### Новые найденные проблемы (чинить в v0.4) +- [ ] `cargo check --no-default-features` сломан исторически: manager/sync и store ссылаются + на redis без `#[cfg(feature = "store-redis")]` — нарушение «features additive». +- [ ] `src/app/runtime.rs`: узкий `#[allow(clippy::exit)]` — паллиатив после переноса main-логики + в lib; правильно — возвращать exit-код из `run()` вместо `process::exit`. +- [ ] IPv6-банов в XDP-putи нет (`ban_cidr` bail'ит на v6) — IPv6-паритет (вывод №8). ## 3. Backlog @@ -190,7 +223,8 @@ guard/ 5. **По умолчанию безопасно**: нет дефолтных секретов; отсутствие обязательного env = fail-fast. 6. **Интеграционный тест на слой**: config parse, filter logic, registry fail-fast (есть); новый слой = новый тест. -7. **Модуль ≤ 300 строк**: CI-гейт через grep/wc скрипт или ревью. +7. **Файл ≤ 250 строк, папка ≤ 4 .rs**: CI-гейт `scripts/check_module_size.sh` + (ratchet-базлайн `.module_size_baseline`: легаси — WARN, новые нарушения — FAIL). 8. **CI guardrails**: `cargo clippy --all-targets -- -D warnings`, `cargo test`, clang-build xdp/core/universal_filter.c, grep на `changeme`. 9. **README/TODO не врут**: каждое число имеет ссылку на тест или отчёт. @@ -201,7 +235,7 @@ guard/ ☐ cargo check / cargo test проходят ☐ cargo clippy --all-targets -- -D warnings — 0 warnings ☐ cargo fmt --check проходит -☐ Ни один модуль не превышает 300 строк +☐ Ни один файл не превышает 250 строк; ни в одной папке src/ больше 4 .rs-файлов ☐ Unit тесты покрывают happy path + 2+ error cases ☐ Нет мёртвого кода: pub без вызовов, конфиг-поле без потребителя, метрика без writer ☐ Нет дефолтных секретов @@ -214,7 +248,7 @@ guard/ ``` ❌ Тесты после кода. Пиши вместе. -❌ Модуль > 300 строк — сигнал декомпозировать немедленно. +❌ Файл > 250 строк или > 4 .rs в папке — сигнал декомпозировать/сгруппировать немедленно. ❌ TODO в коде без issue. ❌ Мёртвый код: pub без вызовов, конфиг-поле без потребителя, метрика без writer. ❌ «Бумажный слой»: фича описана, но не вызывается. @@ -225,4 +259,4 @@ guard/ --- -*Версия: 4.0 | Обновлён: 2026-08-24 (universal redesign)* +*Версия: 5.0 | Обновлён: 2026-09-15 (bug-fix раунд: XDP unsafe, Redis reconnect/SCAN, CI-гейты, лимиты 250/4)* diff --git a/scripts/check_default_secrets.sh b/scripts/check_default_secrets.sh new file mode 100755 index 0000000..fb54e13 --- /dev/null +++ b/scripts/check_default_secrets.sh @@ -0,0 +1,63 @@ +#!/usr/bin/env bash +# +# check_default_secrets.sh — reject the placeholder secret "changeme" +# committed as a real value in code, deploy configs, or docs. +# +# Intentional occurrences are whitelisted in two layers: +# 1. File level — the validation code itself: +# src/bin/rampart-manager.rs (startup guard rejecting the default) +# src/manager/api/auth.rs (test constant) +# 2. Line level — only lines that *use* "changeme" as a value are reported. +# A matching line is skipped when it contains "must not be" or "test" +# (case-insensitive), states the prohibition ("запр" root: запрет / +# запрещено / запрещён — the project docs are partly Russian), or is a +# comment / markdown table row (starts with '#', '//', '*', '|', '--'). +# +# Any remaining occurrence is an offender -> exit 1, offending lines listed. + +set -u + +PATTERN="changeme" +SCAN_DIRS="src deploy docs" +FILE_WHITELIST="src/bin/rampart-manager.rs src/manager/api/auth.rs" + +cd "$(CDPATH= cd -- "$(dirname -- "$0")/.." && pwd)" || exit 1 + +found=0 +for dir in $SCAN_DIRS; do + [ -d "$dir" ] || continue + while IFS= read -r match; do + file="${match%%:*}" + rest="${match#*:}" + lineno="${rest%%:*}" + content="${rest#*:}" + + skip=0 + for whitelisted in $FILE_WHITELIST; do + [ "$file" = "$whitelisted" ] && skip=1 + done + if printf '%s' "$content" | grep -qiE 'must not be|test|запр'; then + skip=1 + fi + trimmed="${content#"${content%%[![:space:]]*}"}" + if [ "${trimmed:0:2}" = "//" ] || [ "${trimmed:0:2}" = "--" ]; then + skip=1 + fi + case "${trimmed:0:1}" in + '#'|'*'|'|') skip=1 ;; + esac + + if [ "$skip" -eq 0 ]; then + found=1 + printf 'FAIL %s:%s: %s\n' "$file" "$lineno" "$trimmed" >&2 + fi + done < <(grep -rn --binary-files=without-match "$PATTERN" "$dir" 2>/dev/null) +done + +if [ "$found" -ne 0 ]; then + echo "secrets gate: FAIL — 'changeme' used as a value (see lines above)" >&2 + exit 1 +fi + +echo "secrets gate: PASS — no unexpected '$PATTERN' occurrences in: $SCAN_DIRS" +exit 0 diff --git a/scripts/check_module_size.sh b/scripts/check_module_size.sh new file mode 100755 index 0000000..2581499 --- /dev/null +++ b/scripts/check_module_size.sh @@ -0,0 +1,161 @@ +#!/usr/bin/env bash +# +# check_module_size.sh — module size ratchet. +# +# Project standard (TODO.md §4): max 250 lines per .rs file; max 4 .rs files +# per source directory (top level, non-recursive, mod.rs included). +# +# Rationale — ratchet pattern: the gate blocks regressions immediately while +# legacy debt drains away file by file. Paths that currently violate the +# standard are pinned one-per-line in .module_size_baseline and are reported +# as WARN ("legacy, pending refactor"); only NEW violations fail the build. +# Baseline entries that no longer violate are stale and also fail, so the +# baseline never hides a shrink that should be locked in — regenerate it with +# --update-baseline and commit the result. +# +# Scope note: we scan src/ only and deliberately do NOT scan tests/. +# Cargo treats every top-level tests/*.rs as its own integration-test +# binary, so those files must stay "one file per binary" at tests/ root; +# splitting an oversized test into a module directory would break cargo's +# test discovery. + +set -u + +MAX_LINES=250 +MAX_FILES=4 +SRC_DIR="src" +BASELINE_FILE=".module_size_baseline" + +export LC_ALL=C +cd "$(CDPATH= cd -- "$(dirname -- "$0")/.." && pwd)" || exit 1 + +UPDATE_BASELINE=0 +for arg in "$@"; do + case "$arg" in + --update-baseline) UPDATE_BASELINE=1 ;; + *) echo "usage: $0 [--update-baseline]" >&2; exit 2 ;; + esac +done + +count_lines() { + wc -l < "$1" | tr -d '[:space:]' +} + +count_dir_rs() { + find "$1" -maxdepth 1 -type f -name '*.rs' | wc -l | tr -d '[:space:]' +} + +# Print current violations, one path per line (dirs carry a trailing /), sorted. +current_violations() { + { + find "$SRC_DIR" -type f -name '*.rs' -not -path '*/target/*' 2>/dev/null | + while IFS= read -r file; do + [ "$(count_lines "$file")" -gt "$MAX_LINES" ] && printf '%s\n' "$file" + done + + find "$SRC_DIR" -type d -not -path '*/target/*' 2>/dev/null | + while IFS= read -r dir; do + [ "$(count_dir_rs "$dir")" -gt "$MAX_FILES" ] && printf '%s/\n' "$dir" + done + } | sort +} + +describe_entry() { + entry="$1" + case "$entry" in + */) + dir="${entry%/}" + if [ -d "$dir" ]; then + printf '%s — %d top-level .rs files (max %d)' "$dir" "$(count_dir_rs "$dir")" "$MAX_FILES" + else + printf '%s — directory no longer exists' "$dir" + fi + ;; + *) + if [ -f "$entry" ]; then + printf '%s — %d lines (max %d)' "$entry" "$(count_lines "$entry")" "$MAX_LINES" + else + printf '%s — file no longer exists' "$entry" + fi + ;; + esac +} + +list_with_prefix() { + # $1 = prefix, $2 = file with one path per line + while IFS= read -r entry; do + [ -n "$entry" ] && printf ' %s %s\n' "$1" "$(describe_entry "$entry")" + done < "$2" +} + +current=$(current_violations) + +if [ "$UPDATE_BASELINE" -eq 1 ]; then + { + echo "# Module size ratchet baseline — read by scripts/check_module_size.sh" + echo "# Oversized files: plain path. Overcrowded directories: path with trailing /." + echo "# Regenerate after a refactor: bash scripts/check_module_size.sh --update-baseline" + printf '%s\n' "$current" | sed '/^$/d' + } > "$BASELINE_FILE" + entries=$(printf '%s\n' "$current" | sed '/^$/d' | wc -l | tr -d '[:space:]') + echo "module-size: baseline regenerated — $entries entries written to $BASELINE_FILE" + exit 0 +fi + +cur_f=$(mktemp) base_f=$(mktemp) new_f=$(mktemp) stale_f=$(mktemp) warn_f=$(mktemp) +trap 'rm -f "$cur_f" "$base_f" "$new_f" "$stale_f" "$warn_f"' EXIT + +printf '%s\n' "$current" | sed '/^$/d' > "$cur_f" +have_base=0 +if [ -f "$BASELINE_FILE" ]; then + have_base=1 + sed -e '/^[[:space:]]*$/d' -e '/^[[:space:]]*#/d' -e 's/[[:space:]]*$//' "$BASELINE_FILE" | sort > "$base_f" +else + : > "$base_f" +fi + +comm -23 "$cur_f" "$base_f" > "$new_f" # violations not baselined -> FAIL +comm -13 "$cur_f" "$base_f" > "$stale_f" # baselined but compliant -> FAIL (stale) +comm -12 "$cur_f" "$base_f" > "$warn_f" # baselined and violating -> WARN + +echo "module-size gate: max $MAX_LINES lines per .rs file, max $MAX_FILES .rs files per directory (scope: $SRC_DIR/)" +echo + +if [ "$have_base" -eq 0 ]; then + echo "note: no $BASELINE_FILE found — all current violations count as new" + echo +fi + +if [ -s "$warn_f" ]; then + echo "WARN (legacy, pending refactor):" + list_with_prefix "WARN " "$warn_f" + echo +fi + +fail=0 +if [ -s "$stale_f" ]; then + fail=1 + echo "STALE baseline entries — these no longer violate the limits:" + list_with_prefix "STALE" "$stale_f" + echo " fix: run 'bash scripts/check_module_size.sh --update-baseline' and commit $BASELINE_FILE" + echo +fi +if [ -s "$new_f" ]; then + fail=1 + echo "NEW violations:" + list_with_prefix "FAIL " "$new_f" + echo +fi + +if [ "$fail" -ne 0 ]; then + echo "FAIL — module size gate" + exit 1 +fi + +warn_n=$(sed -n '$=' "$warn_f") +if [ "${warn_n:-0}" -gt 0 ]; then + echo "PASS — module size gate ($warn_n legacy entries baselined as WARN, no new violations)" +else + echo "PASS — module size gate (no violations)" +fi +exit 0 diff --git a/src/app/mod.rs b/src/app/mod.rs new file mode 100644 index 0000000..ba696c2 --- /dev/null +++ b/src/app/mod.rs @@ -0,0 +1,6 @@ +//! Wiring-логика edge-демона: сборка сервисов и главный цикл запуска. + +pub mod runtime; +pub mod services; + +pub use runtime::run; diff --git a/src/app/runtime.rs b/src/app/runtime.rs new file mode 100644 index 0000000..23a5b1a --- /dev/null +++ b/src/app/runtime.rs @@ -0,0 +1,229 @@ +//! Runtime edge-демона: сборка компонентов, фоновые задачи и запуск listener. + +use crate::config::Config; +use crate::engine::challenge::DifficultyAdjuster; +use crate::engine::listener; +use crate::engine::subnet::monitor; +use crate::engine::tunnel::Gateway; +use crate::filter::blacklist::Blacklist; +use crate::filter::rate_limit::RateLimiter; +use crate::metrics; +use crate::store::clickhouse::{ClickHouseEvent, ClickHouseWriter}; +use crate::traffic::hook::TrafficHook; +use crate::traffic::intel::detector::{AttackDetector, AttackStatus}; +use crate::traffic::intel::reputation::IpReputation; +use crate::traffic::profile::alert::{AlertDispatcher, send_webhook}; +use std::sync::Arc; +use std::sync::Mutex; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Duration; +use tokio::sync::watch; +use tracing_subscriber::EnvFilter; + +use super::services::{build_registry, build_whitelist, start_subnet_tracker, start_xdp}; + +fn attack_status_value(status: AttackStatus) -> i64 { + match status { + AttackStatus::Normal => 0, + AttackStatus::Suspicious => 1, + AttackStatus::UnderAttack => 2, + } +} + +/// Точка входа edge-демона (`rampart`): то, что делал `main` бинаря. +/// +/// # Errors +/// Ошибки старта: нечитаемый конфиг, пустой реестр протоколов, сбой listener. +pub async fn run() -> anyhow::Result<()> { + tracing_subscriber::fmt() + .with_env_filter(EnvFilter::from_default_env().add_directive("rampart=info".parse()?)) + .init(); + + let config_path = std::env::var("RAMPART_CONFIG").unwrap_or_else(|_| "/etc/rampart/config.toml".to_string()); + let config = Arc::new(Config::from_file(&config_path)?); + + let registry = Arc::new(build_registry(&config)); + if registry.is_empty() { + anyhow::bail!( + "no protocol plugins compiled; available feature flags: protocol-http, \ + store-redis, geoip, xdp, io-uring. Build with --features protocol-http \ + or link an external ProtocolHandler implementation" + ); + } + let whitelist = build_whitelist(&config)?; + + let rate_limiter = Arc::new(RateLimiter::new( + config.limits.rate_limit_pps, + config.limits.rate_limit_burst, + )); + let blacklist = Arc::new(Blacklist::new()); + let reputation = Arc::new(IpReputation::new()); + let detector = Arc::new(Mutex::new(AttackDetector::new())); + let hook = Arc::new(TrafficHook::new( + reputation.clone(), + blacklist.clone(), + config.detect.autoban.clone(), + config.ban.ban_duration_secs, + )); + let alert_dispatcher = AlertDispatcher::new(); + let webhook_url = config.detect.alert.webhook_url.clone(); + let allowed_1s = Arc::new(AtomicU64::new(0)); + + let (shutdown_tx, shutdown_rx) = watch::channel(false); + + let sig_tx = shutdown_tx.clone(); + // shutdown-таймер обязан завершить процесс; в исходном коде вызов жил + // внутри `fn main` и был освобождён от clippy::exit. Чистый перенос в + // lib сохранил бы поведение, поэтому узкий allow — минимальная цена. + #[allow(clippy::exit)] + tokio::spawn(async move { + wait_for_signal().await; + tracing::info!("shutdown signal received, draining connections..."); + let _ = sig_tx.send(true); + tokio::time::sleep(Duration::from_secs(5)).await; + tracing::info!("shutdown timeout reached, exiting"); + std::process::exit(0); + }); + + #[cfg(feature = "store-redis")] + if let Some(redis_url) = &config.store.redis_url + && !redis_url.is_empty() + { + let bl = blacklist.clone(); + let sd = shutdown_rx.clone(); + let url = redis_url.clone(); + tokio::spawn(async move { + let client = match redis::Client::open(url.as_str()) { + Ok(c) => c, + Err(e) => { + tracing::warn!("Invalid redis_url: {e}, blacklist sync disabled"); + return; + }, + }; + crate::store::start_blacklist_sync(&client, bl, sd).await; + }); + } + + if config.metrics.enabled { + let metrics_addr = format!("0.0.0.0:{}", config.metrics.port); + tracing::info!("Metrics server listening on {metrics_addr}"); + tokio::spawn(async move { + metrics::run_metrics_server(&metrics_addr).await; + }); + } + + let clickhouse: Option>> = match &config.store.clickhouse_url { + Some(url) if !url.is_empty() => { + let writer = Arc::new(tokio::sync::Mutex::new(ClickHouseWriter::new(url))); + crate::store::clickhouse::start_flush_task(writer.clone(), shutdown_rx.clone()); + Some(writer) + }, + _ => None, + }; + + let xdp_filter = start_xdp(&config, &shutdown_rx)?; + let subnet_tracker = start_subnet_tracker(&config); + monitor::spawn( + &config.detect.prefix, + subnet_tracker.clone(), + xdp_filter.clone(), + shutdown_rx.clone(), + ); + + let rl = rate_limiter.clone(); + let bl = blacklist.clone(); + let det = detector.clone(); + let intel = hook.clone(); + let a1s = allowed_1s.clone(); + let ch = clickhouse.clone(); + let mut sd = shutdown_rx.clone(); + tokio::spawn(async move { + let mut sec_tick = tokio::time::interval(Duration::from_secs(1)); + let mut min_tick = tokio::time::interval(Duration::from_secs(60)); + sec_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + min_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + loop { + tokio::select! { + biased; + _ = sd.changed() => { + if *sd.borrow() { + return; + } + } + _ = sec_tick.tick() => { + let pps = a1s.swap(0, Ordering::Relaxed) as f64; + let cps = intel.take_cps() as f64; + metrics::INTEL_CPS.set(cps as i64); + let status = det.lock().expect("detector lock poisoned").analyze(pps, cps); + metrics::ATTACK_STATUS.set(attack_status_value(status)); + if let Some(alert) = alert_dispatcher.on_status(status) { + // Дедупликация: алерт только на переходе состояния атаки. + metrics::INTEL_ALERTS_TOTAL.inc(); + tracing::warn!(pps, cps, "{alert}"); + push_attack_event(&ch, pps).await; + if let Some(url) = webhook_url.clone() { + tokio::spawn(send_webhook(url, alert.clone())); + } + } + if status == AttackStatus::UnderAttack { + let banned = intel.ban_offenders(); + if banned > 0 { + tracing::info!(banned, "escalation: offenders auto-banned during attack"); + } + } + } + _ = min_tick.tick() => { + rl.sweep(); + bl.clear_expired(); + } + } + } + }); + + tracing::info!("Rampart edge starting on {}:{}", config.bind.address, config.bind.port); + tracing::info!("Upstreams: {:?}", config.backend.upstreams); + + let adjuster = Arc::new(Mutex::new(DifficultyAdjuster::default())); + let gateway = Arc::new(Gateway::new( + config.clone(), + rate_limiter, + blacklist, + adjuster, + whitelist, + reputation, + xdp_filter, + clickhouse, + allowed_1s, + registry, + subnet_tracker, + )); + listener::run(config, gateway, hook, shutdown_rx).await +} + +async fn push_attack_event(writer: &Option>>, pps: f64) { + let Some(writer) = writer else { + return; + }; + let event = ClickHouseEvent { + timestamp: chrono::Utc::now(), + event_type: "attack".to_string(), + ip: String::new(), + data_float: pps, + data_int: 0, + data_string: "under_attack".to_string(), + }; + if let Err(e) = writer.lock().await.push(event).await { + tracing::debug!("clickhouse push error: {e}"); + } +} + +async fn wait_for_signal() { + let ctrl_c = tokio::signal::ctrl_c(); + let mut term = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) + .expect("failed to install SIGTERM handler"); + + tokio::select! { + _ = ctrl_c => {} + _ = term.recv() => {} + } +} diff --git a/src/app/services.rs b/src/app/services.rs new file mode 100644 index 0000000..c2b40bb --- /dev/null +++ b/src/app/services.rs @@ -0,0 +1,119 @@ +//! Сервисы edge-демона: реестр протоколов, whitelist, XDP и subnet-трекер. + +use crate::config::Config; +use crate::engine::subnet::tracker::SubnetTracker; +use crate::protocol::ProtocolRegistry; +use crate::xdp::XdpFilter; +#[cfg(feature = "xdp")] +use crate::xdp::XdpGlobals; +use std::collections::HashSet; +use std::net::IpAddr; +use std::sync::Arc; +use std::sync::Mutex; +use std::time::Duration; +use tokio::sync::watch; + +/// Собирает реестр протоколов из скомпилированных реализаций. +/// +/// Обработчики поставляются фичами: `protocol-http` регистрирует +/// HTTP/1.1-обработчик с политикой из секции `[protocol.http]`. +/// Внешние крейты регистрируют свои обработчики в этом же месте. +pub(crate) fn build_registry(config: &Config) -> ProtocolRegistry { + #[cfg(feature = "protocol-http")] + { + let mut registry = ProtocolRegistry::new(); + register_http(&mut registry, config); + registry + } + #[cfg(not(feature = "protocol-http"))] + { + let _ = config; + ProtocolRegistry::new() + } +} + +#[cfg(feature = "protocol-http")] +fn register_http(registry: &mut ProtocolRegistry, config: &Config) { + registry.register(Box::new(crate::protocol::http::HttpProtocolHandler::new( + &config.protocol.http, + &config.backend.upstreams, + ))); +} + +#[cfg(feature = "xdp")] +pub(crate) fn start_xdp( + config: &Arc, + shutdown_rx: &watch::Receiver, +) -> anyhow::Result>>> { + use crate::xdp::XdpMetrics; + + if !config.xdp.enabled { + return Ok(None); + } + + let filter = XdpFilter::new(&config.xdp.interface); + let shared = Arc::new(Mutex::new(filter)); + { + let mut guard = shared.lock().expect("xdp lock poisoned"); + let globals = XdpGlobals::from_config(&config.xdp)?; + guard.set_globals(globals); + guard.load()?; + } + let xdp_metrics = XdpMetrics::register()?; + + let sd = shutdown_rx.clone(); + let shared_thread = shared.clone(); + std::thread::spawn(move || { + while !*sd.borrow() { + let guard = match shared_thread.lock() { + Ok(g) => g, + Err(_) => break, + }; + guard.drain_events(); + if let Ok(stats) = guard.get_stats() { + xdp_metrics.update(&stats); + } + drop(guard); + std::thread::sleep(Duration::from_secs(5)); + } + if let Ok(mut guard) = shared_thread.lock() { + guard.unload().ok(); + } + }); + Ok(Some(shared)) +} + +#[cfg(not(feature = "xdp"))] +pub(crate) fn start_xdp( + _config: &Arc, + _shutdown_rx: &watch::Receiver, +) -> anyhow::Result>>> { + Ok(None) +} + +/// Юзерспейс-агрегатор префиксов: активен, когда detect.prefix включён, +/// а XDP-путь не работает (не собран или выключен в конфиге). +pub(crate) fn start_subnet_tracker(config: &Config) -> Option> { + let xdp_active = cfg!(feature = "xdp") && config.xdp.enabled; + if config.detect.prefix.enabled && !xdp_active { + tracing::info!( + window_secs = config.detect.prefix.window_secs, + syn_threshold = config.detect.prefix.syn_threshold, + "subnet tracker enabled (userspace prefix aggregation)" + ); + Some(Arc::new(SubnetTracker::new())) + } else { + None + } +} + +pub(crate) fn build_whitelist(config: &Config) -> anyhow::Result>> { + let mut set = HashSet::with_capacity(config.whitelist.len()); + for entry in &config.whitelist { + let ip: IpAddr = entry + .parse() + .map_err(|_| anyhow::anyhow!("invalid whitelist entry: {entry}"))?; + set.insert(ip); + } + Ok(Arc::new(set)) +} diff --git a/src/bin/rampart-cli.rs b/src/bin/rampart-cli.rs index a508672..a9528a1 100644 --- a/src/bin/rampart-cli.rs +++ b/src/bin/rampart-cli.rs @@ -55,7 +55,7 @@ async fn main() -> anyhow::Result<()> { let cli = Cli::parse(); match cli.command { - Commands::Status => commands::status::run().await, + Commands::Status => commands::node::status::run().await, Commands::Doctor => commands::doctor::run().await, Commands::Config { key, value } => commands::config::run(key, value).await, Commands::Blacklist { action } => match action { @@ -64,9 +64,9 @@ async fn main() -> anyhow::Result<()> { BlacklistAction::List => commands::blacklist::list().await, }, Commands::Emergency { mode } => match mode { - EmergencyMode::Enable => commands::emergency::enable().await, - EmergencyMode::Disable => commands::emergency::disable().await, + EmergencyMode::Enable => commands::node::emergency::enable().await, + EmergencyMode::Disable => commands::node::emergency::disable().await, }, - Commands::Drain { node } => commands::drain::run(&node).await, + Commands::Drain { node } => commands::node::drain::run(&node).await, } } diff --git a/src/bin/rampart-manager.rs b/src/bin/rampart-manager.rs index f48f235..b7768fa 100644 --- a/src/bin/rampart-manager.rs +++ b/src/bin/rampart-manager.rs @@ -46,16 +46,16 @@ async fn main() -> anyhow::Result<()> { tokio::spawn(sync::heartbeat::start_heartbeat_check(state.clone())); let public = Router::new() - .route("/api/v1/health", get(api::health::health_check)) + .route("/api/v1/health", get(api::inventory::health::health_check)) .route("/api/v1/auth/login", post(api::auth::login)); let protected = Router::new() - .route("/api/v1/servers", get(api::servers::list_servers)) + .route("/api/v1/servers", get(api::inventory::servers::list_servers)) .route( "/api/v1/blacklist", get(api::blacklist::list_blacklist).post(api::blacklist::add_blacklist), ) - .route("/api/v1/nodes", get(api::nodes::list_nodes)) + .route("/api/v1/nodes", get(api::inventory::nodes::list_nodes)) .route_layer(middleware::from_fn(auth::auth_middleware)); let cors = match std::env::var("CORS_ORIGIN") { diff --git a/src/bin/rampart.rs b/src/bin/rampart.rs index 08154ef..d1c1011 100644 --- a/src/bin/rampart.rs +++ b/src/bin/rampart.rs @@ -1,330 +1,6 @@ -use rampart::config::Config; -use rampart::engine::challenge::DifficultyAdjuster; -use rampart::engine::listener; -use rampart::engine::subnet_monitor; -use rampart::engine::subnet_tracker::SubnetTracker; -use rampart::engine::tunnel::Gateway; -use rampart::filter::blacklist::Blacklist; -use rampart::filter::rate_limit::RateLimiter; -use rampart::metrics; -use rampart::protocol::ProtocolRegistry; -use rampart::store::clickhouse::{ClickHouseEvent, ClickHouseWriter}; -use rampart::traffic::alert::{AlertDispatcher, send_webhook}; -use rampart::traffic::detector::{AttackDetector, AttackStatus}; -use rampart::traffic::hook::TrafficHook; -use rampart::traffic::reputation::IpReputation; -use rampart::xdp::XdpFilter; -#[cfg(feature = "xdp")] -use rampart::xdp::XdpGlobals; -use std::collections::HashSet; -use std::net::IpAddr; -use std::sync::Arc; -use std::sync::Mutex; -use std::sync::atomic::{AtomicU64, Ordering}; -use std::time::Duration; -use tokio::sync::watch; -use tracing_subscriber::EnvFilter; - -fn attack_status_value(status: AttackStatus) -> i64 { - match status { - AttackStatus::Normal => 0, - AttackStatus::Suspicious => 1, - AttackStatus::UnderAttack => 2, - } -} - -/// Собирает реестр протоколов из скомпилированных реализаций. -/// -/// Обработчики поставляются фичами: `protocol-http` регистрирует -/// HTTP/1.1-обработчик с политикой из секции `[protocol.http]`. -/// Внешние крейты регистрируют свои обработчики в этом же месте. -fn build_registry(config: &Config) -> ProtocolRegistry { - #[cfg(feature = "protocol-http")] - { - let mut registry = ProtocolRegistry::new(); - register_http(&mut registry, config); - registry - } - #[cfg(not(feature = "protocol-http"))] - { - let _ = config; - ProtocolRegistry::new() - } -} - -#[cfg(feature = "protocol-http")] -fn register_http(registry: &mut ProtocolRegistry, config: &Config) { - registry.register(Box::new(rampart::protocol::http::HttpProtocolHandler::new( - &config.protocol.http, - &config.backend.upstreams, - ))); -} +use rampart::app; #[tokio::main] async fn main() -> anyhow::Result<()> { - tracing_subscriber::fmt() - .with_env_filter(EnvFilter::from_default_env().add_directive("rampart=info".parse()?)) - .init(); - - let config_path = std::env::var("RAMPART_CONFIG").unwrap_or_else(|_| "/etc/rampart/config.toml".to_string()); - let config = Arc::new(Config::from_file(&config_path)?); - - let registry = Arc::new(build_registry(&config)); - if registry.is_empty() { - anyhow::bail!( - "no protocol plugins compiled; available feature flags: protocol-http, \ - store-redis, geoip, xdp, io-uring. Build with --features protocol-http \ - or link an external ProtocolHandler implementation" - ); - } - let whitelist = build_whitelist(&config)?; - - let rate_limiter = Arc::new(RateLimiter::new( - config.limits.rate_limit_pps, - config.limits.rate_limit_burst, - )); - let blacklist = Arc::new(Blacklist::new()); - let reputation = Arc::new(IpReputation::new()); - let detector = Arc::new(Mutex::new(AttackDetector::new())); - let hook = Arc::new(TrafficHook::new( - reputation.clone(), - blacklist.clone(), - config.detect.autoban.clone(), - config.ban.ban_duration_secs, - )); - let alert_dispatcher = AlertDispatcher::new(); - let webhook_url = config.detect.alert.webhook_url.clone(); - let allowed_1s = Arc::new(AtomicU64::new(0)); - - let (shutdown_tx, shutdown_rx) = watch::channel(false); - - let sig_tx = shutdown_tx.clone(); - tokio::spawn(async move { - wait_for_signal().await; - tracing::info!("shutdown signal received, draining connections..."); - let _ = sig_tx.send(true); - tokio::time::sleep(Duration::from_secs(5)).await; - tracing::info!("shutdown timeout reached, exiting"); - std::process::exit(0); - }); - - #[cfg(feature = "store-redis")] - if let Some(redis_url) = &config.store.redis_url - && !redis_url.is_empty() - { - let bl = blacklist.clone(); - let sd = shutdown_rx.clone(); - let url = redis_url.clone(); - tokio::spawn(async move { - let client = match redis::Client::open(url.as_str()) { - Ok(c) => c, - Err(e) => { - tracing::warn!("Invalid redis_url: {e}, blacklist sync disabled"); - return; - }, - }; - rampart::store::start_blacklist_sync(&client, bl, sd).await; - }); - } - - if config.metrics.enabled { - let metrics_addr = format!("0.0.0.0:{}", config.metrics.port); - tracing::info!("Metrics server listening on {metrics_addr}"); - tokio::spawn(async move { - metrics::run_metrics_server(&metrics_addr).await; - }); - } - - let clickhouse: Option>> = match &config.store.clickhouse_url { - Some(url) if !url.is_empty() => { - let writer = Arc::new(tokio::sync::Mutex::new(ClickHouseWriter::new(url))); - rampart::store::clickhouse::start_flush_task(writer.clone(), shutdown_rx.clone()); - Some(writer) - }, - _ => None, - }; - - let xdp_filter = start_xdp(&config, &shutdown_rx)?; - let subnet_tracker = start_subnet_tracker(&config); - subnet_monitor::spawn( - &config.detect.prefix, - subnet_tracker.clone(), - xdp_filter.clone(), - shutdown_rx.clone(), - ); - - let rl = rate_limiter.clone(); - let bl = blacklist.clone(); - let det = detector.clone(); - let intel = hook.clone(); - let a1s = allowed_1s.clone(); - let ch = clickhouse.clone(); - let mut sd = shutdown_rx.clone(); - tokio::spawn(async move { - let mut sec_tick = tokio::time::interval(Duration::from_secs(1)); - let mut min_tick = tokio::time::interval(Duration::from_secs(60)); - sec_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); - min_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); - loop { - tokio::select! { - biased; - _ = sd.changed() => { - if *sd.borrow() { - return; - } - } - _ = sec_tick.tick() => { - let pps = a1s.swap(0, Ordering::Relaxed) as f64; - let cps = intel.take_cps() as f64; - metrics::INTEL_CPS.set(cps as i64); - let status = det.lock().expect("detector lock poisoned").analyze(pps, cps); - metrics::ATTACK_STATUS.set(attack_status_value(status)); - if let Some(alert) = alert_dispatcher.on_status(status) { - // Дедупликация: алерт только на переходе состояния атаки. - metrics::INTEL_ALERTS_TOTAL.inc(); - tracing::warn!(pps, cps, "{alert}"); - push_attack_event(&ch, pps).await; - if let Some(url) = webhook_url.clone() { - tokio::spawn(send_webhook(url, alert.clone())); - } - } - if status == AttackStatus::UnderAttack { - let banned = intel.ban_offenders(); - if banned > 0 { - tracing::info!(banned, "escalation: offenders auto-banned during attack"); - } - } - } - _ = min_tick.tick() => { - rl.sweep(); - bl.clear_expired(); - } - } - } - }); - - tracing::info!("Rampart edge starting on {}:{}", config.bind.address, config.bind.port); - tracing::info!("Upstreams: {:?}", config.backend.upstreams); - - let adjuster = Arc::new(Mutex::new(DifficultyAdjuster::default())); - let gateway = Arc::new(Gateway::new( - config.clone(), - rate_limiter, - blacklist, - adjuster, - whitelist, - reputation, - xdp_filter, - clickhouse, - allowed_1s, - registry, - subnet_tracker, - )); - listener::run(config, gateway, hook, shutdown_rx).await -} - -async fn push_attack_event(writer: &Option>>, pps: f64) { - let Some(writer) = writer else { - return; - }; - let event = ClickHouseEvent { - timestamp: chrono::Utc::now(), - event_type: "attack".to_string(), - ip: String::new(), - data_float: pps, - data_int: 0, - data_string: "under_attack".to_string(), - }; - if let Err(e) = writer.lock().await.push(event).await { - tracing::debug!("clickhouse push error: {e}"); - } -} - -#[cfg(feature = "xdp")] -fn start_xdp( - config: &Arc, - shutdown_rx: &watch::Receiver, -) -> anyhow::Result>>> { - use rampart::xdp::XdpMetrics; - - if !config.xdp.enabled { - return Ok(None); - } - - let filter = XdpFilter::new(&config.xdp.interface); - let shared = Arc::new(Mutex::new(filter)); - { - let mut guard = shared.lock().expect("xdp lock poisoned"); - let globals = XdpGlobals::from_config(&config.xdp)?; - guard.set_globals(globals); - guard.load()?; - } - let xdp_metrics = XdpMetrics::register()?; - - let sd = shutdown_rx.clone(); - let shared_thread = shared.clone(); - std::thread::spawn(move || { - while !*sd.borrow() { - let guard = match shared_thread.lock() { - Ok(g) => g, - Err(_) => break, - }; - guard.drain_events(); - if let Ok(stats) = guard.get_stats() { - xdp_metrics.update(&stats); - } - drop(guard); - std::thread::sleep(Duration::from_secs(5)); - } - if let Ok(mut guard) = shared_thread.lock() { - guard.unload().ok(); - } - }); - Ok(Some(shared)) -} - -#[cfg(not(feature = "xdp"))] -fn start_xdp( - _config: &Arc, - _shutdown_rx: &watch::Receiver, -) -> anyhow::Result>>> { - Ok(None) -} - -/// Юзерспейс-агрегатор префиксов: активен, когда detect.prefix включён, -/// а XDP-путь не работает (не собран или выключен в конфиге). -fn start_subnet_tracker(config: &Config) -> Option> { - let xdp_active = cfg!(feature = "xdp") && config.xdp.enabled; - if config.detect.prefix.enabled && !xdp_active { - tracing::info!( - window_secs = config.detect.prefix.window_secs, - syn_threshold = config.detect.prefix.syn_threshold, - "subnet tracker enabled (userspace prefix aggregation)" - ); - Some(Arc::new(SubnetTracker::new())) - } else { - None - } -} - -fn build_whitelist(config: &Config) -> anyhow::Result>> { - let mut set = HashSet::with_capacity(config.whitelist.len()); - for entry in &config.whitelist { - let ip: IpAddr = entry - .parse() - .map_err(|_| anyhow::anyhow!("invalid whitelist entry: {entry}"))?; - set.insert(ip); - } - Ok(Arc::new(set)) -} - -async fn wait_for_signal() { - let ctrl_c = tokio::signal::ctrl_c(); - let mut term = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) - .expect("failed to install SIGTERM handler"); - - tokio::select! { - _ = ctrl_c => {} - _ = term.recv() => {} - } + app::run().await } diff --git a/src/cli/commands/mod.rs b/src/cli/commands/mod.rs index b5d001a..9a7d149 100644 --- a/src/cli/commands/mod.rs +++ b/src/cli/commands/mod.rs @@ -1,6 +1,4 @@ pub mod blacklist; pub mod config; pub mod doctor; -pub mod drain; -pub mod emergency; -pub mod status; +pub mod node; diff --git a/src/cli/commands/drain.rs b/src/cli/commands/node/drain.rs similarity index 100% rename from src/cli/commands/drain.rs rename to src/cli/commands/node/drain.rs diff --git a/src/cli/commands/emergency.rs b/src/cli/commands/node/emergency.rs similarity index 100% rename from src/cli/commands/emergency.rs rename to src/cli/commands/node/emergency.rs diff --git a/src/cli/commands/node/mod.rs b/src/cli/commands/node/mod.rs new file mode 100644 index 0000000..fd7613d --- /dev/null +++ b/src/cli/commands/node/mod.rs @@ -0,0 +1,4 @@ +//! CLI-команды управления edge-узлами: статус, drain, аварийный режим. +pub mod drain; +pub mod emergency; +pub mod status; diff --git a/src/cli/commands/status.rs b/src/cli/commands/node/status.rs similarity index 100% rename from src/cli/commands/status.rs rename to src/cli/commands/node/status.rs diff --git a/src/config/sections.rs b/src/config/sections.rs deleted file mode 100644 index d558660..0000000 --- a/src/config/sections.rs +++ /dev/null @@ -1,404 +0,0 @@ -//! Секции единого конфига платформы. - -use serde::Deserialize; - -#[derive(Debug, Clone, Deserialize)] -pub struct BindConfig { - #[serde(default = "default_bind_address")] - pub address: String, - #[serde(default = "default_bind_port")] - pub port: u16, -} - -impl Default for BindConfig { - fn default() -> Self { - Self { - address: default_bind_address(), - port: default_bind_port(), - } - } -} - -fn default_bind_address() -> String { - "0.0.0.0".to_string() -} -fn default_bind_port() -> u16 { - 25565 -} - -/// Generic upstream-бэкенды (список addr:port) за edge-нодой. -#[derive(Debug, Clone, Deserialize)] -pub struct BackendConfig { - #[serde(default = "default_upstreams")] - pub upstreams: Vec, -} - -impl Default for BackendConfig { - fn default() -> Self { - Self { - upstreams: default_upstreams(), - } - } -} - -fn default_upstreams() -> Vec { - vec!["127.0.0.1:25566".to_string()] -} - -#[derive(Debug, Clone, Deserialize)] -pub struct WorkerConfig { - #[serde(default = "default_worker_count")] - pub count: usize, -} - -impl Default for WorkerConfig { - fn default() -> Self { - Self { count: 4 } - } -} - -fn default_worker_count() -> usize { - 4 -} - -#[derive(Debug, Clone, Deserialize)] -pub struct LimitsConfig { - #[serde(default = "default_handshake_timeout")] - pub handshake_timeout_secs: u64, - #[serde(default = "default_max_connections_per_ip")] - pub max_connections_per_ip: u32, - #[serde(default = "default_rate_limit_pps")] - pub rate_limit_pps: f64, - #[serde(default = "default_rate_limit_burst")] - pub rate_limit_burst: f64, -} - -impl Default for LimitsConfig { - fn default() -> Self { - Self { - handshake_timeout_secs: default_handshake_timeout(), - max_connections_per_ip: default_max_connections_per_ip(), - rate_limit_pps: default_rate_limit_pps(), - rate_limit_burst: default_rate_limit_burst(), - } - } -} - -fn default_handshake_timeout() -> u64 { - 5 -} -fn default_max_connections_per_ip() -> u32 { - 10 -} -fn default_rate_limit_pps() -> f64 { - 5.0 -} -fn default_rate_limit_burst() -> f64 { - 10.0 -} - -#[derive(Debug, Clone, Deserialize)] -pub struct BanConfig { - #[serde(default = "default_ban_duration")] - pub ban_duration_secs: u64, -} - -impl Default for BanConfig { - fn default() -> Self { - Self { - ban_duration_secs: default_ban_duration(), - } - } -} - -fn default_ban_duration() -> u64 { - 3600 -} - -#[derive(Debug, Clone, Deserialize)] -pub struct StoreConfig { - pub redis_url: Option, - #[serde(default = "default_blacklist_cache_ttl")] - pub blacklist_cache_ttl_secs: u64, - pub clickhouse_url: Option, -} - -impl Default for StoreConfig { - fn default() -> Self { - Self { - redis_url: None, - blacklist_cache_ttl_secs: default_blacklist_cache_ttl(), - clickhouse_url: None, - } - } -} - -fn default_blacklist_cache_ttl() -> u64 { - 300 -} - -#[derive(Debug, Clone, Deserialize)] -pub struct XdpConfig { - #[serde(default)] - pub enabled: bool, - #[serde(default = "default_xdp_interface")] - pub interface: String, - #[serde(default = "default_xdp_port_start")] - pub protected_port_start: u16, - #[serde(default = "default_xdp_port_end")] - pub protected_port_end: u16, - /// "pass" | "drop" | "rate-limit" → G_UDP_POLICY. - #[serde(default = "default_xdp_udp_policy")] - pub udp_policy: String, - #[serde(default = "default_xdp_udp_hit_count")] - pub udp_rate_hit_count: u32, - #[serde(default = "default_xdp_udp_window_ms")] - pub udp_rate_window_ms: u64, - #[serde(default)] - pub syn_challenge_enabled: bool, - /// 128-битный hex (16 символов) для G_CHALLENGE_SECRET; None = плейсхолдер из config.h. - #[serde(default)] - pub challenge_secret_hex: Option, - #[serde(default = "default_xdp_challenge_timeout_ms")] - pub challenge_timeout_ms: u32, - #[serde(default = "default_xdp_throttle_enabled")] - pub throttle_enabled: bool, - #[serde(default = "default_xdp_events_enabled")] - pub events_enabled: bool, -} - -impl Default for XdpConfig { - fn default() -> Self { - Self { - enabled: false, - interface: default_xdp_interface(), - protected_port_start: default_xdp_port_start(), - protected_port_end: default_xdp_port_end(), - udp_policy: default_xdp_udp_policy(), - udp_rate_hit_count: default_xdp_udp_hit_count(), - udp_rate_window_ms: default_xdp_udp_window_ms(), - syn_challenge_enabled: false, - challenge_secret_hex: None, - challenge_timeout_ms: default_xdp_challenge_timeout_ms(), - throttle_enabled: default_xdp_throttle_enabled(), - events_enabled: default_xdp_events_enabled(), - } - } -} - -fn default_xdp_interface() -> String { - "eth0".to_string() -} -fn default_xdp_port_start() -> u16 { - 1 -} -fn default_xdp_port_end() -> u16 { - 65535 -} -fn default_xdp_udp_policy() -> String { - "pass".to_string() -} -fn default_xdp_udp_hit_count() -> u32 { - 100 -} -fn default_xdp_udp_window_ms() -> u64 { - 1000 -} -fn default_xdp_challenge_timeout_ms() -> u32 { - 3000 -} -fn default_xdp_throttle_enabled() -> bool { - true -} -fn default_xdp_events_enabled() -> bool { - true -} - -#[derive(Debug, Clone, Deserialize)] -pub struct LoggingConfig { - #[serde(default = "default_log_level")] - pub level: String, - #[serde(default = "default_log_format")] - pub format: String, -} - -impl Default for LoggingConfig { - fn default() -> Self { - Self { - level: default_log_level(), - format: default_log_format(), - } - } -} - -fn default_log_level() -> String { - "info".to_string() -} -fn default_log_format() -> String { - "text".to_string() -} - -#[derive(Debug, Clone, Deserialize)] -pub struct MetricsConfig { - #[serde(default = "default_metrics_enabled")] - pub enabled: bool, - #[serde(default = "default_metrics_port")] - pub port: u16, -} - -impl Default for MetricsConfig { - fn default() -> Self { - Self { - enabled: default_metrics_enabled(), - port: default_metrics_port(), - } - } -} - -fn default_metrics_enabled() -> bool { - true -} -fn default_metrics_port() -> u16 { - 9090 -} - -#[derive(Debug, Clone, Deserialize)] -pub struct PowConfig { - #[serde(default)] - pub enabled: bool, - #[serde(default = "default_pow_difficulty")] - pub difficulty: u8, -} - -impl Default for PowConfig { - fn default() -> Self { - Self { - enabled: false, - difficulty: default_pow_difficulty(), - } - } -} - -fn default_pow_difficulty() -> u8 { - 4 -} - -/// Секция `[detect]`: детекторы распределённых атак. -#[derive(Debug, Clone, Default, Deserialize)] -pub struct DetectConfig { - #[serde(default)] - pub prefix: DetectPrefixConfig, - #[serde(default)] - pub autoban: DetectAutobanConfig, - #[serde(default)] - pub alert: DetectAlertConfig, -} - -/// Секция `[detect.autoban]`: авто-бан источников по репутации и сигналам -/// детектора атак. Порог сравнивается со скором [`crate::traffic::reputation::IpReputation`], -/// TTL берётся из `ban.ban_duration_secs`. -#[derive(Debug, Clone, Deserialize)] -pub struct DetectAutobanConfig { - #[serde(default)] - pub enabled: bool, - #[serde(default = "default_autoban_reputation_threshold")] - pub reputation_threshold: i32, -} - -impl Default for DetectAutobanConfig { - fn default() -> Self { - Self { - enabled: false, - reputation_threshold: default_autoban_reputation_threshold(), - } - } -} - -fn default_autoban_reputation_threshold() -> i32 { - -50 -} - -/// Секция `[detect.alert]`: webhook-алерты на переходах состояния атаки -/// («атака началась» / «атака закончилась»). -#[derive(Debug, Clone, Default, Deserialize)] -pub struct DetectAlertConfig { - #[serde(default)] - pub webhook_url: Option, -} - -/// Секция `[detect.prefix]`: subnet-level детектор распределённых атак. -#[derive(Debug, Clone, Deserialize)] -pub struct DetectPrefixConfig { - #[serde(default)] - pub enabled: bool, - #[serde(default = "default_prefix_syn_threshold")] - pub syn_threshold: u64, - #[serde(default = "default_prefix_window_secs")] - pub window_secs: u64, - #[serde(default = "default_prefix_min_unique_sources")] - pub min_unique_sources: u64, -} - -impl Default for DetectPrefixConfig { - fn default() -> Self { - Self { - enabled: false, - syn_threshold: default_prefix_syn_threshold(), - window_secs: default_prefix_window_secs(), - min_unique_sources: default_prefix_min_unique_sources(), - } - } -} - -fn default_prefix_syn_threshold() -> u64 { - 500 -} -fn default_prefix_window_secs() -> u64 { - 10 -} -fn default_prefix_min_unique_sources() -> u64 { - 16 -} - -/// Секция `[protocol]`: протокольные обработчики (plugin-by-feature). -#[cfg(feature = "protocol-http")] -#[derive(Debug, Clone, Default, Deserialize)] -pub struct ProtocolConfig { - #[serde(default)] - pub http: HttpProtocolConfig, -} - -/// Секция `[protocol.http]`: политика HTTP/1.1-обработчика на edge-ноде. -#[cfg(feature = "protocol-http")] -#[derive(Debug, Clone, Deserialize)] -pub struct HttpProtocolConfig { - #[serde(default = "default_http_max_header_bytes")] - pub max_header_bytes: usize, - #[serde(default = "default_http_timeout_secs")] - pub timeout_secs: u64, - #[serde(default)] - pub blocked_paths: Vec, - #[serde(default)] - pub require_user_agent: bool, -} - -#[cfg(feature = "protocol-http")] -impl Default for HttpProtocolConfig { - fn default() -> Self { - Self { - max_header_bytes: default_http_max_header_bytes(), - timeout_secs: default_http_timeout_secs(), - blocked_paths: Vec::new(), - require_user_agent: false, - } - } -} - -#[cfg(feature = "protocol-http")] -fn default_http_max_header_bytes() -> usize { - 8192 -} -#[cfg(feature = "protocol-http")] -fn default_http_timeout_secs() -> u64 { - 5 -} diff --git a/src/config/sections/detect.rs b/src/config/sections/detect.rs new file mode 100644 index 0000000..271f935 --- /dev/null +++ b/src/config/sections/detect.rs @@ -0,0 +1,121 @@ +use serde::Deserialize; + +/// Секция `[detect]`: детекторы распределённых атак. +#[derive(Debug, Clone, Default, Deserialize)] +pub struct DetectConfig { + #[serde(default)] + pub prefix: DetectPrefixConfig, + #[serde(default)] + pub autoban: DetectAutobanConfig, + #[serde(default)] + pub alert: DetectAlertConfig, +} + +/// Секция `[detect.autoban]`: авто-бан источников по репутации и сигналам +/// детектора атак. Порог сравнивается со скором [`crate::traffic::reputation::IpReputation`], +/// TTL берётся из `ban.ban_duration_secs`. +#[derive(Debug, Clone, Deserialize)] +pub struct DetectAutobanConfig { + #[serde(default)] + pub enabled: bool, + #[serde(default = "default_autoban_reputation_threshold")] + pub reputation_threshold: i32, +} + +impl Default for DetectAutobanConfig { + fn default() -> Self { + Self { + enabled: false, + reputation_threshold: default_autoban_reputation_threshold(), + } + } +} + +fn default_autoban_reputation_threshold() -> i32 { + -50 +} + +/// Секция `[detect.alert]`: webhook-алерты на переходах состояния атаки +/// («атака началась» / «атака закончилась»). +#[derive(Debug, Clone, Default, Deserialize)] +pub struct DetectAlertConfig { + #[serde(default)] + pub webhook_url: Option, +} + +/// Секция `[detect.prefix]`: subnet-level детектор распределённых атак. +#[derive(Debug, Clone, Deserialize)] +pub struct DetectPrefixConfig { + #[serde(default)] + pub enabled: bool, + #[serde(default = "default_prefix_syn_threshold")] + pub syn_threshold: u64, + #[serde(default = "default_prefix_window_secs")] + pub window_secs: u64, + #[serde(default = "default_prefix_min_unique_sources")] + pub min_unique_sources: u64, +} + +impl Default for DetectPrefixConfig { + fn default() -> Self { + Self { + enabled: false, + syn_threshold: default_prefix_syn_threshold(), + window_secs: default_prefix_window_secs(), + min_unique_sources: default_prefix_min_unique_sources(), + } + } +} + +fn default_prefix_syn_threshold() -> u64 { + 500 +} +fn default_prefix_window_secs() -> u64 { + 10 +} +fn default_prefix_min_unique_sources() -> u64 { + 16 +} + +/// Секция `[protocol]`: протокольные обработчики (plugin-by-feature). +#[cfg(feature = "protocol-http")] +#[derive(Debug, Clone, Default, Deserialize)] +pub struct ProtocolConfig { + #[serde(default)] + pub http: HttpProtocolConfig, +} + +/// Секция `[protocol.http]`: политика HTTP/1.1-обработчика на edge-ноде. +#[cfg(feature = "protocol-http")] +#[derive(Debug, Clone, Deserialize)] +pub struct HttpProtocolConfig { + #[serde(default = "default_http_max_header_bytes")] + pub max_header_bytes: usize, + #[serde(default = "default_http_timeout_secs")] + pub timeout_secs: u64, + #[serde(default)] + pub blocked_paths: Vec, + #[serde(default)] + pub require_user_agent: bool, +} + +#[cfg(feature = "protocol-http")] +impl Default for HttpProtocolConfig { + fn default() -> Self { + Self { + max_header_bytes: default_http_max_header_bytes(), + timeout_secs: default_http_timeout_secs(), + blocked_paths: Vec::new(), + require_user_agent: false, + } + } +} + +#[cfg(feature = "protocol-http")] +fn default_http_max_header_bytes() -> usize { + 8192 +} +#[cfg(feature = "protocol-http")] +fn default_http_timeout_secs() -> u64 { + 5 +} diff --git a/src/config/sections/edge.rs b/src/config/sections/edge.rs new file mode 100644 index 0000000..457ccf8 --- /dev/null +++ b/src/config/sections/edge.rs @@ -0,0 +1,135 @@ +use serde::Deserialize; + +#[derive(Debug, Clone, Deserialize)] +pub struct BindConfig { + #[serde(default = "default_bind_address")] + pub address: String, + #[serde(default = "default_bind_port")] + pub port: u16, +} + +impl Default for BindConfig { + fn default() -> Self { + Self { + address: default_bind_address(), + port: default_bind_port(), + } + } +} + +fn default_bind_address() -> String { + "0.0.0.0".to_string() +} +fn default_bind_port() -> u16 { + 25565 +} + +/// Generic upstream-бэкенды (список addr:port) за edge-нодой. +#[derive(Debug, Clone, Deserialize)] +pub struct BackendConfig { + #[serde(default = "default_upstreams")] + pub upstreams: Vec, +} + +impl Default for BackendConfig { + fn default() -> Self { + Self { + upstreams: default_upstreams(), + } + } +} + +fn default_upstreams() -> Vec { + vec!["127.0.0.1:25566".to_string()] +} + +#[derive(Debug, Clone, Deserialize)] +pub struct WorkerConfig { + #[serde(default = "default_worker_count")] + pub count: usize, +} + +impl Default for WorkerConfig { + fn default() -> Self { + Self { count: 4 } + } +} + +fn default_worker_count() -> usize { + 4 +} + +#[derive(Debug, Clone, Deserialize)] +pub struct LimitsConfig { + #[serde(default = "default_handshake_timeout")] + pub handshake_timeout_secs: u64, + #[serde(default = "default_max_connections_per_ip")] + pub max_connections_per_ip: u32, + #[serde(default = "default_rate_limit_pps")] + pub rate_limit_pps: f64, + #[serde(default = "default_rate_limit_burst")] + pub rate_limit_burst: f64, +} + +impl Default for LimitsConfig { + fn default() -> Self { + Self { + handshake_timeout_secs: default_handshake_timeout(), + max_connections_per_ip: default_max_connections_per_ip(), + rate_limit_pps: default_rate_limit_pps(), + rate_limit_burst: default_rate_limit_burst(), + } + } +} + +fn default_handshake_timeout() -> u64 { + 5 +} +fn default_max_connections_per_ip() -> u32 { + 10 +} +fn default_rate_limit_pps() -> f64 { + 5.0 +} +fn default_rate_limit_burst() -> f64 { + 10.0 +} + +#[derive(Debug, Clone, Deserialize)] +pub struct BanConfig { + #[serde(default = "default_ban_duration")] + pub ban_duration_secs: u64, +} + +impl Default for BanConfig { + fn default() -> Self { + Self { + ban_duration_secs: default_ban_duration(), + } + } +} + +fn default_ban_duration() -> u64 { + 3600 +} + +#[derive(Debug, Clone, Deserialize)] +pub struct PowConfig { + #[serde(default)] + pub enabled: bool, + #[serde(default = "default_pow_difficulty")] + pub difficulty: u8, +} + +impl Default for PowConfig { + fn default() -> Self { + Self { + enabled: false, + difficulty: default_pow_difficulty(), + } + } +} + +fn default_pow_difficulty() -> u8 { + 4 +} diff --git a/src/config/sections/mod.rs b/src/config/sections/mod.rs new file mode 100644 index 0000000..0febf42 --- /dev/null +++ b/src/config/sections/mod.rs @@ -0,0 +1,11 @@ +//! Секции единого конфига платформы. + +mod detect; +mod edge; +mod platform; + +pub use detect::{DetectAlertConfig, DetectAutobanConfig, DetectConfig, DetectPrefixConfig}; +#[cfg(feature = "protocol-http")] +pub use detect::{HttpProtocolConfig, ProtocolConfig}; +pub use edge::{BackendConfig, BanConfig, BindConfig, LimitsConfig, PowConfig, WorkerConfig}; +pub use platform::{LoggingConfig, MetricsConfig, StoreConfig, XdpConfig}; diff --git a/src/config/sections/platform.rs b/src/config/sections/platform.rs new file mode 100644 index 0000000..eb66198 --- /dev/null +++ b/src/config/sections/platform.rs @@ -0,0 +1,148 @@ +use serde::Deserialize; + +#[derive(Debug, Clone, Deserialize)] +pub struct StoreConfig { + pub redis_url: Option, + #[serde(default = "default_blacklist_cache_ttl")] + pub blacklist_cache_ttl_secs: u64, + pub clickhouse_url: Option, +} + +impl Default for StoreConfig { + fn default() -> Self { + Self { + redis_url: None, + blacklist_cache_ttl_secs: default_blacklist_cache_ttl(), + clickhouse_url: None, + } + } +} + +fn default_blacklist_cache_ttl() -> u64 { + 300 +} + +#[derive(Debug, Clone, Deserialize)] +pub struct XdpConfig { + #[serde(default)] + pub enabled: bool, + #[serde(default = "default_xdp_interface")] + pub interface: String, + #[serde(default = "default_xdp_port_start")] + pub protected_port_start: u16, + #[serde(default = "default_xdp_port_end")] + pub protected_port_end: u16, + /// "pass" | "drop" | "rate-limit" → G_UDP_POLICY. + #[serde(default = "default_xdp_udp_policy")] + pub udp_policy: String, + #[serde(default = "default_xdp_udp_hit_count")] + pub udp_rate_hit_count: u32, + #[serde(default = "default_xdp_udp_window_ms")] + pub udp_rate_window_ms: u64, + #[serde(default)] + pub syn_challenge_enabled: bool, + /// 128-битный hex (16 символов) для G_CHALLENGE_SECRET; None = плейсхолдер из config.h. + #[serde(default)] + pub challenge_secret_hex: Option, + #[serde(default = "default_xdp_challenge_timeout_ms")] + pub challenge_timeout_ms: u32, + #[serde(default = "default_xdp_throttle_enabled")] + pub throttle_enabled: bool, + #[serde(default = "default_xdp_events_enabled")] + pub events_enabled: bool, +} + +impl Default for XdpConfig { + fn default() -> Self { + Self { + enabled: false, + interface: default_xdp_interface(), + protected_port_start: default_xdp_port_start(), + protected_port_end: default_xdp_port_end(), + udp_policy: default_xdp_udp_policy(), + udp_rate_hit_count: default_xdp_udp_hit_count(), + udp_rate_window_ms: default_xdp_udp_window_ms(), + syn_challenge_enabled: false, + challenge_secret_hex: None, + challenge_timeout_ms: default_xdp_challenge_timeout_ms(), + throttle_enabled: default_xdp_throttle_enabled(), + events_enabled: default_xdp_events_enabled(), + } + } +} + +fn default_xdp_interface() -> String { + "eth0".to_string() +} +fn default_xdp_port_start() -> u16 { + 1 +} +fn default_xdp_port_end() -> u16 { + 65535 +} +fn default_xdp_udp_policy() -> String { + "pass".to_string() +} +fn default_xdp_udp_hit_count() -> u32 { + 100 +} +fn default_xdp_udp_window_ms() -> u64 { + 1000 +} +fn default_xdp_challenge_timeout_ms() -> u32 { + 3000 +} +fn default_xdp_throttle_enabled() -> bool { + true +} +fn default_xdp_events_enabled() -> bool { + true +} + +#[derive(Debug, Clone, Deserialize)] +pub struct LoggingConfig { + #[serde(default = "default_log_level")] + pub level: String, + #[serde(default = "default_log_format")] + pub format: String, +} + +impl Default for LoggingConfig { + fn default() -> Self { + Self { + level: default_log_level(), + format: default_log_format(), + } + } +} + +fn default_log_level() -> String { + "info".to_string() +} +fn default_log_format() -> String { + "text".to_string() +} + +#[derive(Debug, Clone, Deserialize)] +pub struct MetricsConfig { + #[serde(default = "default_metrics_enabled")] + pub enabled: bool, + #[serde(default = "default_metrics_port")] + pub port: u16, +} + +impl Default for MetricsConfig { + fn default() -> Self { + Self { + enabled: default_metrics_enabled(), + port: default_metrics_port(), + } + } +} + +fn default_metrics_enabled() -> bool { + true +} +fn default_metrics_port() -> u16 { + 9090 +} diff --git a/src/engine/challenge/difficulty.rs b/src/engine/challenge/difficulty.rs new file mode 100644 index 0000000..8e9105b --- /dev/null +++ b/src/engine/challenge/difficulty.rs @@ -0,0 +1,88 @@ +use crate::metrics; +use std::collections::VecDeque; +use std::time::Instant; + +/// Адаптирует сложность PoW к текущему темпу подключений. +pub struct DifficultyAdjuster { + window: VecDeque, + min: u8, + max: u8, + current: u8, +} + +impl Default for DifficultyAdjuster { + fn default() -> Self { + Self::new(4, 10) + } +} + +impl DifficultyAdjuster { + #[must_use] + pub fn new(min: u8, max: u8) -> Self { + Self { + window: VecDeque::new(), + min: min.max(4), + max: max.min(10), + current: min.max(4), + } + } + + pub fn record_connection(&mut self) { + let now = Instant::now(); + self.window.push_back(now); + while let Some(&t) = self.window.front() { + if now.duration_since(t).as_secs() >= 1 { + self.window.pop_front(); + } else { + break; + } + } + let new_diff = self.compute_difficulty(); + if self.current != new_diff { + tracing::info!( + old = self.current, + new = new_diff, + window = self.window.len(), + "pow: difficulty adjusted" + ); + self.current = new_diff; + metrics::POW_CURRENT_DIFFICULTY.set(self.current as i64); + } + } + + #[must_use] + pub fn current_difficulty(&self) -> u8 { + metrics::POW_CURRENT_DIFFICULTY.set(self.current as i64); + self.current + } + + fn compute_difficulty(&self) -> u8 { + match self.window.len() { + cps if cps > 500 => self.max.max(self.min), + cps if cps > 200 => 8, + cps if cps > 50 => 6, + _ => self.min, + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn adjuster_raises_difficulty_under_load() { + let mut adjuster = DifficultyAdjuster::new(4, 10); + for _ in 0..600 { + adjuster.record_connection(); + } + assert_eq!(adjuster.current_difficulty(), 10); + } + + #[test] + fn adjuster_stays_minimal_when_idle() { + let mut adjuster = DifficultyAdjuster::new(4, 10); + adjuster.record_connection(); + assert_eq!(adjuster.current_difficulty(), 4); + } +} diff --git a/src/engine/challenge/mod.rs b/src/engine/challenge/mod.rs new file mode 100644 index 0000000..c75559d --- /dev/null +++ b/src/engine/challenge/mod.rs @@ -0,0 +1,8 @@ +//! Универсальный SHA-256 hashcash: генерация challenge, решатель, +//! верификатор и адаптивная сложность. + +mod difficulty; +mod pow; + +pub use difficulty::DifficultyAdjuster; +pub use pow::{Challenge, enforce, solve}; diff --git a/src/engine/challenge.rs b/src/engine/challenge/pow.rs similarity index 70% rename from src/engine/challenge.rs rename to src/engine/challenge/pow.rs index 318d633..0646a47 100644 --- a/src/engine/challenge.rs +++ b/src/engine/challenge/pow.rs @@ -1,10 +1,5 @@ -//! Универсальный SHA-256 hashcash: генерация challenge, решатель, -//! верификатор и адаптивная сложность. - -use crate::metrics; use rand::RngCore; use sha2::{Digest, Sha256}; -use std::collections::VecDeque; use std::net::IpAddr; use std::time::{Duration, Instant}; use subtle::ConstantTimeEq; @@ -92,70 +87,6 @@ pub fn solve(challenge: &str, difficulty: u8) -> Option { None } -/// Адаптирует сложность PoW к текущему темпу подключений. -pub struct DifficultyAdjuster { - window: VecDeque, - min: u8, - max: u8, - current: u8, -} - -impl Default for DifficultyAdjuster { - fn default() -> Self { - Self::new(4, 10) - } -} - -impl DifficultyAdjuster { - #[must_use] - pub fn new(min: u8, max: u8) -> Self { - Self { - window: VecDeque::new(), - min: min.max(4), - max: max.min(10), - current: min.max(4), - } - } - - pub fn record_connection(&mut self) { - let now = Instant::now(); - self.window.push_back(now); - while let Some(&t) = self.window.front() { - if now.duration_since(t).as_secs() >= 1 { - self.window.pop_front(); - } else { - break; - } - } - let new_diff = self.compute_difficulty(); - if self.current != new_diff { - tracing::info!( - old = self.current, - new = new_diff, - window = self.window.len(), - "pow: difficulty adjusted" - ); - self.current = new_diff; - metrics::POW_CURRENT_DIFFICULTY.set(self.current as i64); - } - } - - #[must_use] - pub fn current_difficulty(&self) -> u8 { - metrics::POW_CURRENT_DIFFICULTY.set(self.current as i64); - self.current - } - - fn compute_difficulty(&self) -> u8 { - match self.window.len() { - cps if cps > 500 => self.max.max(self.min), - cps if cps > 200 => 8, - cps if cps > 50 => 6, - _ => self.min, - } - } -} - /// Проводит текстовый PoW-gate в потоке: выдаёт challenge и проверяет ответ. /// /// # Errors @@ -241,20 +172,4 @@ mod tests { let big_nonce = "0".repeat(MAX_NONCE_LEN + 1); assert!(!challenge.verify(&big_nonce)); } - - #[test] - fn adjuster_raises_difficulty_under_load() { - let mut adjuster = DifficultyAdjuster::new(4, 10); - for _ in 0..600 { - adjuster.record_connection(); - } - assert_eq!(adjuster.current_difficulty(), 10); - } - - #[test] - fn adjuster_stays_minimal_when_idle() { - let mut adjuster = DifficultyAdjuster::new(4, 10); - adjuster.record_connection(); - assert_eq!(adjuster.current_difficulty(), 4); - } } diff --git a/src/engine/mod.rs b/src/engine/mod.rs index 36626ed..c2708f9 100644 --- a/src/engine/mod.rs +++ b/src/engine/mod.rs @@ -2,6 +2,9 @@ pub mod challenge; pub mod listener; -pub mod subnet_monitor; -pub mod subnet_tracker; +pub mod subnet; pub mod tunnel; + +// Совместимость путей после переезда в subnet/: `engine::subnet_tracker` / +// `engine::subnet_monitor` остаются валидными алиасами модулей. +pub use subnet::{monitor as subnet_monitor, tracker as subnet_tracker}; diff --git a/src/engine/subnet/mod.rs b/src/engine/subnet/mod.rs new file mode 100644 index 0000000..1462cec --- /dev/null +++ b/src/engine/subnet/mod.rs @@ -0,0 +1,3 @@ +//! Subnet-level детектор: юзерспейс-агрегатор префиксов и монитор вердиктов. +pub mod monitor; +pub mod tracker; diff --git a/src/engine/subnet_monitor.rs b/src/engine/subnet/monitor.rs similarity index 100% rename from src/engine/subnet_monitor.rs rename to src/engine/subnet/monitor.rs diff --git a/src/engine/subnet_tracker.rs b/src/engine/subnet/tracker.rs similarity index 100% rename from src/engine/subnet_tracker.rs rename to src/engine/subnet/tracker.rs diff --git a/src/lib.rs b/src/lib.rs index 1526efe..cdd7320 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,4 +1,5 @@ //! Rampart — универсальная платформа сетевой защиты (L3/L4/L7). +pub mod app; pub mod cli; pub mod config; pub mod engine; diff --git a/src/manager/api/health.rs b/src/manager/api/inventory/health.rs similarity index 100% rename from src/manager/api/health.rs rename to src/manager/api/inventory/health.rs diff --git a/src/manager/api/inventory/mod.rs b/src/manager/api/inventory/mod.rs new file mode 100644 index 0000000..70d55b6 --- /dev/null +++ b/src/manager/api/inventory/mod.rs @@ -0,0 +1,4 @@ +//! Инвентарь manager-API: health, узлы и серверы. +pub mod health; +pub mod nodes; +pub mod servers; diff --git a/src/manager/api/nodes.rs b/src/manager/api/inventory/nodes.rs similarity index 100% rename from src/manager/api/nodes.rs rename to src/manager/api/inventory/nodes.rs diff --git a/src/manager/api/servers.rs b/src/manager/api/inventory/servers.rs similarity index 100% rename from src/manager/api/servers.rs rename to src/manager/api/inventory/servers.rs diff --git a/src/manager/api/mod.rs b/src/manager/api/mod.rs index 990010f..9181329 100644 --- a/src/manager/api/mod.rs +++ b/src/manager/api/mod.rs @@ -1,5 +1,3 @@ pub mod auth; pub mod blacklist; -pub mod health; -pub mod nodes; -pub mod servers; +pub mod inventory; diff --git a/src/manager/sync/heartbeat.rs b/src/manager/sync/heartbeat.rs index d8413c4..5ce0e54 100644 --- a/src/manager/sync/heartbeat.rs +++ b/src/manager/sync/heartbeat.rs @@ -1,8 +1,17 @@ use crate::manager::AppState; +use futures::StreamExt; use redis::AsyncCommands; +use redis::ScanOptions; use std::sync::Arc; use tokio::time::{Duration, interval}; +/// SCAN pattern for node records (blocking KEYS is forbidden under load). +const NODES_KEY_PATTERN: &str = "rampart:nodes:*"; +/// SCAN COUNT hint per round. +const SCAN_BATCH: usize = 100; +/// Node is offline when the last heartbeat is older than this. +const NODE_TTL_SECS: i64 = 60; + pub async fn start_heartbeat_check(state: Arc) { let mut ticker = interval(Duration::from_secs(30)); loop { @@ -15,30 +24,99 @@ pub async fn start_heartbeat_check(state: Arc) { async fn check_nodes(state: &AppState) -> anyhow::Result<()> { let mut conn = state.redis_client.get_multiplexed_async_connection().await?; - let keys: Vec = redis::cmd("KEYS").arg("rampart:nodes:*").query_async(&mut conn).await?; + let keys = collect_node_keys(&mut conn).await?; let now = chrono::Utc::now().timestamp(); for key in &keys { let raw: Option = conn.get(key).await?; - if let Some(json) = raw - && let Ok(mut node) = serde_json::from_str::(&json) - { - let hb = node["last_heartbeat"] - .as_str() - .and_then(|s| chrono::DateTime::parse_from_rfc3339(s).ok()) - .map(|t| t.timestamp()) - .unwrap_or(0); - if now - hb > 60 { - if let Some(obj) = node.as_object_mut() { - obj.insert("status".to_string(), serde_json::Value::String("offline".to_string())); - if let Ok(updated) = serde_json::to_string(&node) { - let _: () = conn.set(key.as_str(), updated).await.unwrap_or_default(); - } - } - tracing::warn!("Node {key} is offline (heartbeat expired)"); - } - } + let Some(updated) = raw.and_then(|json| offline_update(&json, now)) else { + continue; + }; + tracing::warn!("Node {key} is offline (heartbeat expired)"); + let _: () = conn.set(key.as_str(), updated).await.unwrap_or_default(); } Ok(()) } + +/// Incrementally iterates the keyspace via SCAN (cursor-based, non-blocking). +async fn collect_node_keys(conn: &mut redis::aio::MultiplexedConnection) -> anyhow::Result> { + let opts = ScanOptions::default() + .with_pattern(NODES_KEY_PATTERN) + .with_count(SCAN_BATCH); + let mut iter = conn.scan_options::(opts).await?; + let mut keys = Vec::new(); + while let Some(key) = iter.next().await { + keys.push(key); + } + Ok(keys) +} + +/// Pure decision: returns the updated node JSON with `status: "offline"` when +/// the RFC3339 `last_heartbeat` is older than `NODE_TTL_SECS` (a missing or +/// unparseable timestamp counts as expired, matching the previous behavior); +/// `None` when the node is fresh, unparsable, or not a JSON object. +fn offline_update(raw: &str, now: i64) -> Option { + let mut node = serde_json::from_str::(raw).ok()?; + let hb = node["last_heartbeat"] + .as_str() + .and_then(|s| chrono::DateTime::parse_from_rfc3339(s).ok()) + .map(|t| t.timestamp()) + .unwrap_or(0); + if now - hb <= NODE_TTL_SECS { + return None; + } + node.as_object_mut()? + .insert("status".to_string(), serde_json::Value::String("offline".to_string())); + serde_json::to_string(&node).ok() +} + +#[cfg(test)] +mod tests { + use super::*; + + fn rfc3339(ts: i64) -> String { + chrono::DateTime::from_timestamp(ts, 0) + .expect("in-range timestamp") + .to_rfc3339() + } + + #[test] + fn expired_heartbeat_is_marked_offline() { + let now = 2_000_000_i64; + let raw = format!( + r#"{{"name":"edge1","last_heartbeat":"{}","status":"online"}}"#, + rfc3339(now - 120) + ); + let updated = offline_update(&raw, now).expect("expired node should be marked offline"); + assert!(updated.contains("offline")); + assert!(updated.contains("edge1")); + } + + #[test] + fn fresh_heartbeat_is_left_untouched() { + let now = 2_000_000_i64; + let raw = format!(r#"{{"last_heartbeat":"{}","status":"online"}}"#, rfc3339(now - 5)); + assert!(offline_update(&raw, now).is_none()); + } + + #[test] + fn missing_or_broken_timestamp_counts_as_expired() { + let now = 2_000_000_i64; + let updated = offline_update(r#"{"status":"online"}"#, now).expect("missing heartbeat counts as expired"); + assert!(updated.contains("offline")); + let updated = offline_update(r#"{"last_heartbeat":"garbage","status":"online"}"#, now) + .expect("unparseable heartbeat counts as expired"); + assert!(updated.contains("offline")); + } + + #[test] + fn invalid_json_is_ignored() { + assert!(offline_update("not json", 2_000_000).is_none()); + } + + #[test] + fn node_key_pattern_targets_rampart_nodes() { + assert_eq!(NODES_KEY_PATTERN, "rampart:nodes:*"); + } +} diff --git a/src/store/redis.rs b/src/store/redis.rs index 58eb760..6b4abc4 100644 --- a/src/store/redis.rs +++ b/src/store/redis.rs @@ -31,53 +31,116 @@ impl RedisStore { } } +const BLACKLIST_CHANNEL: &str = "rampart:blacklist:events"; +const BACKOFF_MIN: Duration = Duration::from_secs(1); +const BACKOFF_MAX: Duration = Duration::from_secs(30); + +/// Capped exponential backoff step: `prev` doubled, clamped to `max`. +fn next_backoff(prev: Duration, max: Duration) -> Duration { + prev.saturating_mul(2).min(max) +} + +/// What the outer reconnect loop should do after one pubsub session. +enum Flow { + /// Shutdown requested — exit the task. + Stop, + /// Connection dropped or session failed — reconnect with backoff. + Retry, +} + +/// Resolves only when shutdown is truly requested (value `true`) or every +/// sender has been dropped; spurious non-shutdown notifications keep it pending +/// so it never cancels in-flight connect/subscribe work. +async fn wait_for_shutdown(shutdown: &mut watch::Receiver) { + loop { + if shutdown.changed().await.is_err() || *shutdown.borrow_and_update() { + return; + } + } +} + +/// Blacklist pubsub subscriber with real reconnect: on connect/subscribe +/// failure or stream end it retries with capped exponential backoff (reset to +/// `BACKOFF_MIN` after each successfully applied message). Every await is +/// raced against the shutdown watch so the task stays cancellable. pub async fn start_blacklist_sync( client: &redis::Client, blacklist: Arc, mut shutdown: watch::Receiver, ) { + let mut backoff = BACKOFF_MIN; + loop { + match subscribe_once(client, &blacklist, &mut shutdown, &mut backoff).await { + Flow::Stop => { + tracing::info!("shutting down blacklist subscriber"); + return; + }, + Flow::Retry => {}, + } + + tokio::select! { + _ = tokio::time::sleep(backoff) => {}, + () = wait_for_shutdown(&mut shutdown) => { + tracing::info!("shutting down blacklist subscriber"); + return; + }, + } + backoff = next_backoff(backoff, BACKOFF_MAX); + } +} + +/// One connect+subscribe+consume session. Sets `*backoff` to `BACKOFF_MIN` +/// whenever a message is applied, proving the link is healthy. Each blocking +/// await is raced against shutdown. (Connect+subscribe stay inline: the +/// `PubSub` type produced by `into_pubsub` is awkward to name in a split +/// helper's return type.) +async fn subscribe_once( + client: &redis::Client, + blacklist: &Blacklist, + shutdown: &mut watch::Receiver, + backoff: &mut Duration, +) -> Flow { #[allow(deprecated)] - let conn = match client.get_async_connection().await { - Ok(c) => c, - Err(e) => { - tracing::error!("failed to connect to Redis for blacklist sync: {e}"); - return; + let conn = tokio::select! { + connected = client.get_async_connection() => match connected { + Ok(c) => c, + Err(e) => { + tracing::error!("failed to connect to Redis for blacklist sync: {e}"); + return Flow::Retry; + }, }, + () = wait_for_shutdown(shutdown) => return Flow::Stop, }; let mut pubsub = conn.into_pubsub(); - if let Err(e) = pubsub.subscribe("rampart:blacklist:events").await { - tracing::error!("failed to subscribe to blacklist events: {e}"); - return; + tokio::select! { + subscribed = pubsub.subscribe(BLACKLIST_CHANNEL) => match subscribed { + Ok(()) => tracing::info!("subscribed to {BLACKLIST_CHANNEL}"), + Err(e) => { + tracing::error!("failed to subscribe to blacklist events: {e}"); + return Flow::Retry; + }, + }, + () = wait_for_shutdown(shutdown) => return Flow::Stop, } - tracing::info!("subscribed to rampart:blacklist:events"); + let mut stream = pubsub.on_message(); loop { - let mut stream = pubsub.on_message(); - let msg_fut = stream.next(); - tokio::pin!(msg_fut); - tokio::select! { - _ = shutdown.changed() => { - if *shutdown.borrow() { - tracing::info!("shutting down blacklist subscriber"); - return; - } - } - result = &mut msg_fut => { - match result { - Some(msg) => { - if let Err(e) = handle_event(&msg, &blacklist) { - tracing::error!("blacklist event error: {e}"); - } + maybe_msg = stream.next() => match maybe_msg { + Some(msg) => { + if let Err(e) = handle_event(&msg, blacklist) { + tracing::error!("blacklist event error: {e}"); + } else { + *backoff = BACKOFF_MIN; } - None => { - tracing::error!("pubsub stream ended"); - tokio::time::sleep(Duration::from_secs(1)).await; - return; - } - } - } + }, + None => { + tracing::error!("blacklist pubsub stream ended; resubscribing"); + return Flow::Retry; + }, + }, + () = wait_for_shutdown(shutdown) => return Flow::Stop, } } } @@ -126,3 +189,23 @@ impl StateStore for RedisStore { Ok(()) } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn backoff_doubles_until_cap() { + let mut b = BACKOFF_MIN; + for expected_secs in [2, 4, 8, 16, 30, 30] { + b = next_backoff(b, BACKOFF_MAX); + assert_eq!(b, Duration::from_secs(expected_secs)); + } + } + + #[test] + fn backoff_never_exceeds_max_even_with_huge_prev() { + let b = next_backoff(Duration::from_secs(u64::MAX / 2), BACKOFF_MAX); + assert!(b <= BACKOFF_MAX); + } +} diff --git a/src/traffic/detector.rs b/src/traffic/intel/detector.rs similarity index 100% rename from src/traffic/detector.rs rename to src/traffic/intel/detector.rs diff --git a/src/traffic/ewma.rs b/src/traffic/intel/ewma.rs similarity index 100% rename from src/traffic/ewma.rs rename to src/traffic/intel/ewma.rs diff --git a/src/traffic/intel/mod.rs b/src/traffic/intel/mod.rs new file mode 100644 index 0000000..493cb59 --- /dev/null +++ b/src/traffic/intel/mod.rs @@ -0,0 +1,4 @@ +//! Интеллект трафика: EWMA-сглаживание, детектор атак и репутация IP. +pub mod detector; +pub mod ewma; +pub mod reputation; diff --git a/src/traffic/reputation.rs b/src/traffic/intel/reputation.rs similarity index 100% rename from src/traffic/reputation.rs rename to src/traffic/intel/reputation.rs diff --git a/src/traffic/mod.rs b/src/traffic/mod.rs index 9c4bda7..bf7eb56 100644 --- a/src/traffic/mod.rs +++ b/src/traffic/mod.rs @@ -1,9 +1,11 @@ //! Профилирование трафика, детектор атак и репутация IP. -pub mod alert; -pub mod detector; -pub mod ewma; pub mod hook; +pub mod intel; pub mod prefix; -pub mod profiler; -pub mod reputation; +pub mod profile; + +// Совместимость путей после группировки: домены `intel/` и `profile/` +// остаются доступны по прежним путям `traffic::{detector, ewma, ...}`. +pub use intel::{detector, ewma, reputation}; +pub use profile::{alert, profiler}; diff --git a/src/traffic/prefix/key.rs b/src/traffic/prefix/key.rs new file mode 100644 index 0000000..b75c789 --- /dev/null +++ b/src/traffic/prefix/key.rs @@ -0,0 +1,91 @@ +use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; + +/// `family` для IPv4-префиксов (/24). +pub const PREFIX_FAMILY_V4: u8 = 4; +/// `family` для IPv6-префиксов (/64). +pub const PREFIX_FAMILY_V6: u8 = 6; + +/// Ключ карты `prefix_stats` (`BPF_MAP_TYPE_LRU_HASH`, max_entries 65536). +#[repr(C)] +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct PrefixKey { + pub family: u8, + pub pad: [u8; 7], + pub addr: [u8; 16], +} + +impl PrefixKey { + /// Строит ключ префикса для адреса (v4 → /24 в первых 4 байтах, v6 → /64). + #[must_use] + pub fn new(ip: IpAddr) -> Self { + let mut key = Self::default(); + match ip { + IpAddr::V4(v4) => { + key.family = PREFIX_FAMILY_V4; + key.addr[..4].copy_from_slice(&masked_v4(v4).octets()); + }, + IpAddr::V6(v6) => { + key.family = PREFIX_FAMILY_V6; + key.addr.copy_from_slice(&masked_v6_segments(v6)); + }, + } + key + } + + /// Префикс как IP-адрес (маскированный). + #[must_use] + pub fn to_prefix(self) -> Option { + match self.family { + PREFIX_FAMILY_V4 => Some(IpAddr::V4(Ipv4Addr::new( + self.addr[0], + self.addr[1], + self.addr[2], + self.addr[3], + ))), + PREFIX_FAMILY_V6 => Some(IpAddr::V6(Ipv6Addr::from(self.addr))), + _ => None, + } + } + + #[must_use] + pub fn from_bytes(bytes: [u8; 24]) -> Self { + Self { + family: bytes[0], + pad: bytes[1..8].try_into().expect("7 pad bytes"), + addr: bytes[8..24].try_into().expect("16 addr bytes"), + } + } + + #[must_use] + pub fn to_bytes(self) -> [u8; 24] { + let mut out = [0u8; 24]; + out[0] = self.family; + out[1..8].copy_from_slice(&self.pad); + out[8..24].copy_from_slice(&self.addr); + out + } +} + +/// Возвращает префикс адреса: IPv4 и IPv4-mapped → /24, IPv6 → /64. +#[must_use] +pub fn prefix_of(ip: IpAddr) -> IpAddr { + match ip { + IpAddr::V4(v4) => IpAddr::V4(masked_v4(v4)), + IpAddr::V6(v6) => match v6.to_ipv4_mapped() { + // ::ffff:a.b.c.d живёт в том же /24-пространстве, что и a.b.c.d. + Some(v4) => IpAddr::V4(masked_v4(v4)), + None => IpAddr::V6(Ipv6Addr::from(masked_v6_segments(v6))), + }, + } +} + +fn masked_v4(ip: Ipv4Addr) -> Ipv4Addr { + let o = ip.octets(); + Ipv4Addr::new(o[0], o[1], o[2], 0) +} + +fn masked_v6_segments(ip: Ipv6Addr) -> [u8; 16] { + let mut segs = ip.octets(); + segs[8..].fill(0); + segs +} diff --git a/src/traffic/prefix/mod.rs b/src/traffic/prefix/mod.rs new file mode 100644 index 0000000..1d3d641 --- /dev/null +++ b/src/traffic/prefix/mod.rs @@ -0,0 +1,14 @@ +//! Subnet-level (prefix) детектор распределённых атак. +//! +//! Источник данных — трейт [`PrefixStatsSource`]: карта XDP `prefix_stats` +//! (feature `xdp`) либо юзерспейс-агрегатор [`crate::engine::subnet_tracker`]. +//! Контракт с XDP-агентом: [`PrefixKey`] / [`PrefixStatsVal`] должны совпадать +//! побайтово со структурами ядра. + +mod key; +mod stats; + +pub use key::{PREFIX_FAMILY_V4, PREFIX_FAMILY_V6, PrefixKey, prefix_of}; +#[cfg(feature = "xdp")] +pub use stats::XdpPrefixStats; +pub use stats::{PrefixSnapshot, PrefixStatsSource, PrefixStatsVal, PrefixVerdict, SubnetDetector, SubnetVerdict}; diff --git a/src/traffic/prefix.rs b/src/traffic/prefix/stats.rs similarity index 67% rename from src/traffic/prefix.rs rename to src/traffic/prefix/stats.rs index 0a4a81b..1c63844 100644 --- a/src/traffic/prefix.rs +++ b/src/traffic/prefix/stats.rs @@ -1,27 +1,6 @@ -//! Subnet-level (prefix) детектор распределённых атак. -//! -//! Источник данных — трейт [`PrefixStatsSource`]: карта XDP `prefix_stats` -//! (feature `xdp`) либо юзерспейс-агрегатор [`crate::engine::subnet_tracker`]. -//! Контракт с XDP-агентом: [`PrefixKey`] / [`PrefixStatsVal`] должны совпадать -//! побайтово со структурами ядра. - use anyhow::Result; use std::future::Future; -use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; - -/// `family` для IPv4-префиксов (/24). -pub const PREFIX_FAMILY_V4: u8 = 4; -/// `family` для IPv6-префиксов (/64). -pub const PREFIX_FAMILY_V6: u8 = 6; - -/// Ключ карты `prefix_stats` (`BPF_MAP_TYPE_LRU_HASH`, max_entries 65536). -#[repr(C)] -#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] -pub struct PrefixKey { - pub family: u8, - pub pad: [u8; 7], - pub addr: [u8; 16], -} +use std::net::IpAddr; /// Значение карты `prefix_stats`. #[repr(C)] @@ -56,58 +35,6 @@ impl PrefixStatsVal { } } -impl PrefixKey { - /// Строит ключ префикса для адреса (v4 → /24 в первых 4 байтах, v6 → /64). - #[must_use] - pub fn new(ip: IpAddr) -> Self { - let mut key = Self::default(); - match ip { - IpAddr::V4(v4) => { - key.family = PREFIX_FAMILY_V4; - key.addr[..4].copy_from_slice(&masked_v4(v4).octets()); - }, - IpAddr::V6(v6) => { - key.family = PREFIX_FAMILY_V6; - key.addr.copy_from_slice(&masked_v6_segments(v6)); - }, - } - key - } - - /// Префикс как IP-адрес (маскированный). - #[must_use] - pub fn to_prefix(self) -> Option { - match self.family { - PREFIX_FAMILY_V4 => Some(IpAddr::V4(Ipv4Addr::new( - self.addr[0], - self.addr[1], - self.addr[2], - self.addr[3], - ))), - PREFIX_FAMILY_V6 => Some(IpAddr::V6(Ipv6Addr::from(self.addr))), - _ => None, - } - } - - #[must_use] - pub fn from_bytes(bytes: [u8; 24]) -> Self { - Self { - family: bytes[0], - pad: bytes[1..8].try_into().expect("7 pad bytes"), - addr: bytes[8..24].try_into().expect("16 addr bytes"), - } - } - - #[must_use] - pub fn to_bytes(self) -> [u8; 24] { - let mut out = [0u8; 24]; - out[0] = self.family; - out[1..8].copy_from_slice(&self.pad); - out[8..24].copy_from_slice(&self.addr); - out - } -} - /// Агрегат по одному префиксу за окно детекции. /// /// `unique_sources == 0` означает, что источник не считает уникальные адреса @@ -129,30 +56,6 @@ pub trait PrefixStatsSource { fn snapshot(&self) -> impl Future>> + Send; } -/// Возвращает префикс адреса: IPv4 и IPv4-mapped → /24, IPv6 → /64. -#[must_use] -pub fn prefix_of(ip: IpAddr) -> IpAddr { - match ip { - IpAddr::V4(v4) => IpAddr::V4(masked_v4(v4)), - IpAddr::V6(v6) => match v6.to_ipv4_mapped() { - // ::ffff:a.b.c.d живёт в том же /24-пространстве, что и a.b.c.d. - Some(v4) => IpAddr::V4(masked_v4(v4)), - None => IpAddr::V6(Ipv6Addr::from(masked_v6_segments(v6))), - }, - } -} - -fn masked_v4(ip: Ipv4Addr) -> Ipv4Addr { - let o = ip.octets(); - Ipv4Addr::new(o[0], o[1], o[2], 0) -} - -fn masked_v6_segments(ip: Ipv6Addr) -> [u8; 16] { - let mut segs = ip.octets(); - segs[8..].fill(0); - segs -} - /// Эскалация вердиктов subnet-level детектора. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum SubnetVerdict { diff --git a/src/traffic/alert.rs b/src/traffic/profile/alert.rs similarity index 100% rename from src/traffic/alert.rs rename to src/traffic/profile/alert.rs diff --git a/src/traffic/profile/mod.rs b/src/traffic/profile/mod.rs new file mode 100644 index 0000000..82773ea --- /dev/null +++ b/src/traffic/profile/mod.rs @@ -0,0 +1,3 @@ +//! Профилирование трафика edge-ноды и алерты детектора атак. +pub mod alert; +pub mod profiler; diff --git a/src/traffic/profiler.rs b/src/traffic/profile/profiler.rs similarity index 100% rename from src/traffic/profiler.rs rename to src/traffic/profile/profiler.rs diff --git a/src/tui/fetch.rs b/src/tui/metrics/fetch.rs similarity index 100% rename from src/tui/fetch.rs rename to src/tui/metrics/fetch.rs diff --git a/src/tui/metrics/mod.rs b/src/tui/metrics/mod.rs new file mode 100644 index 0000000..4265f9e --- /dev/null +++ b/src/tui/metrics/mod.rs @@ -0,0 +1,3 @@ +//! Метрики TUI: парсер Prometheus-exposition и fetcher эндпоинтов. +pub mod fetch; +pub mod prometheus; diff --git a/src/tui/prometheus.rs b/src/tui/metrics/prometheus.rs similarity index 100% rename from src/tui/prometheus.rs rename to src/tui/metrics/prometheus.rs diff --git a/src/tui/mod.rs b/src/tui/mod.rs index 23417a9..8cf9d58 100644 --- a/src/tui/mod.rs +++ b/src/tui/mod.rs @@ -4,7 +4,8 @@ //! exposition endpoint (`/metrics`). pub mod app; -pub mod fetch; -pub mod prometheus; +pub mod metrics; pub mod state; pub mod ui; + +pub use metrics::{fetch, prometheus}; diff --git a/src/xdp/diagnostics.rs b/src/xdp/diagnostics.rs deleted file mode 100644 index bd06c46..0000000 --- a/src/xdp/diagnostics.rs +++ /dev/null @@ -1,285 +0,0 @@ -//! Диагностика окружения перед XDP-attach: ядро, BTF, драйвер NIC, привилегии. -//! -//! Основной источник боли при XDP — непонятные ошибки attach на неподдерживаемых -//! ядрах и драйверах. Модуль собирает [`EnvironmentReport`] ДО загрузки BPF-программы -//! и даёт человекочитаемый вердикт с предупреждениями. - -use anyhow::{Result, bail}; -use std::fmt; - -/// Минимально поддерживаемая версия ядра: 5.15 LTS. -pub const MIN_KERNEL: KernelVersion = KernelVersion::new(5, 15, 0); - -const OSRELEASE_PATH: &str = "/proc/sys/kernel/osrelease"; -const BTF_PATH: &str = "/sys/kernel/btf/vmlinux"; -const STATUS_PATH: &str = "/proc/self/status"; - -/// Драйверы с известной поддержкой native XDP (по данным xdpgeneric/native -/// матриц upstream-ядра). Отсутствие в списке не означает отсутствие поддержки, -/// но для таких драйверов предсказываем generic mode. -const NATIVE_XDP_DRIVERS: &[&str] = &[ - "virtio_net", - "ixgbe", - "ixgbevf", - "i40e", - "iavf", - "ice", - "mlx4_core", - "mlx5_core", - "igb", - "igc", - "e1000e", - "vmxnet3", - "bnxt_en", - "nfp", - "sfc", - "ena", - "hv_netvsc", - "macvlan", -]; - -/// Версия ядра Linux, пригодная для сравнения (`major.minor.patch`). -#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] -pub struct KernelVersion { - pub major: u16, - pub minor: u16, - pub patch: u16, -} - -impl KernelVersion { - #[must_use] - pub const fn new(major: u16, minor: u16, patch: u16) -> Self { - Self { major, minor, patch } - } - - /// Парсит строку формата `uname -r`. Берёт первые три числовых компонента, - /// дистрибутивные суффиксы (`6.8.0-45-generic`, `5.15.0-rc2`) отбрасываются. - /// Строка без цифр не парсится. - #[must_use] - pub fn parse(release: &str) -> Option { - let mut nums = release - .split(|c: char| !c.is_ascii_digit()) - .filter_map(|part| part.parse::().ok()); - let major = nums.next()?; - let minor = nums.next().unwrap_or(0); - let patch = nums.next().unwrap_or(0); - Some(Self::new(major, minor, patch)) - } -} - -impl fmt::Display for KernelVersion { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - write!(f, "{}.{}.{}", self.major, self.minor, self.patch) - } -} - -/// Режим присоединения XDP-программы к интерфейсу. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub enum AttachMode { - /// Программа исполняется в контексте драйвера NIC — минимальные накладные расходы. - Native, - /// Программа вызывается из сетевого стека после ingress — работает везде, дороже по CPU. - Generic, - /// Аппаратная разгрузка в NIC — требует явного включения и поддержки железа, - /// автоматически не выбирается. - Offload, -} - -impl fmt::Display for AttachMode { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - f.write_str(match self { - Self::Native => "native", - Self::Generic => "generic", - Self::Offload => "offload", - }) - } -} - -/// Абстракция над procfs/sysfs: реальная реализация — [`FilesystemProbe`], -/// в тестах — мок. Изолирует диагностику от реального ядра. -pub trait SystemProbe { - /// Читает файл целиком. - /// - /// # Errors - /// Проксирует ошибку чтения файла. - fn read_file(&self, path: &str) -> std::io::Result; - - #[must_use] - fn path_exists(&self, path: &str) -> bool; - - #[must_use] - fn symlink_target(&self, path: &str) -> Option; - - #[must_use] - fn effective_uid(&self) -> u32; -} - -/// Реальная реализация [`SystemProbe`] поверх `std::fs`. -pub struct FilesystemProbe; - -impl SystemProbe for FilesystemProbe { - fn read_file(&self, path: &str) -> std::io::Result { - std::fs::read_to_string(path) - } - - fn path_exists(&self, path: &str) -> bool { - std::path::Path::new(path).exists() - } - - fn symlink_target(&self, path: &str) -> Option { - Some(std::fs::read_link(path).ok()?.to_string_lossy().into_owned()) - } - - /// Эвристика: эффективный uid берётся из `/proc/self/status`. - /// root почти всегда имеет CAP_BPF/CAP_NET_ADMIN (если явно не урезаны); - /// для не-root подтверждение возможно только попыткой attach. - /// Нечитаемый статус трактуем как непривилегированный процесс. - fn effective_uid(&self) -> u32 { - std::fs::read_to_string(STATUS_PATH) - .ok() - .and_then(|status| parse_euid(&status)) - .unwrap_or(1) - } -} - -fn parse_euid(status: &str) -> Option { - status.lines().find_map(|line| { - let mut fields = line.strip_prefix("Uid:")?.split_whitespace(); - fields.nth(1)?.parse().ok() - }) -} - -fn driver_from_link(target: &str) -> Option { - let name = target.rsplit('/').next()?; - (!name.is_empty()).then(|| name.to_owned()) -} - -/// Поддерживает ли драйвер native XDP (по таблице известных драйверов). -#[must_use] -pub fn driver_supports_native_xdp(driver: &str) -> bool { - NATIVE_XDP_DRIVERS.contains(&driver) -} - -/// Структурированный отчёт о пригодности окружения для XDP. -pub struct EnvironmentReport { - pub kernel_version: Option, - pub btf_available: bool, - pub privileged: bool, - pub driver_name: Option, - pub warnings: Vec, -} - -impl EnvironmentReport { - /// Собирает отчёт через произвольную реализацию [`SystemProbe`]. - #[must_use] - pub fn collect(probe: &dyn SystemProbe, interface: &str) -> Self { - let kernel_version = probe - .read_file(OSRELEASE_PATH) - .ok() - .and_then(|r| KernelVersion::parse(r.trim())); - let btf_available = probe.path_exists(BTF_PATH); - let privileged = probe.effective_uid() == 0; - let driver_link = probe.symlink_target(&format!("/sys/class/net/{interface}/device/driver")); - let driver_name = driver_link.as_deref().and_then(driver_from_link); - - let mut warnings = Vec::new(); - if !privileged { - warnings.push( - "процесс не от root: CAP_BPF/CAP_NET_ADMIN не подтверждены, \ - attach скорее всего завершится EPERM" - .to_owned(), - ); - } - if !btf_available { - warnings.push(format!( - "{BTF_PATH} недоступен: CO-RE релокации невозможны, переносимость BPF-программы ограничена" - )); - } - if kernel_version.is_some_and(|v| v < KernelVersion::new(5, 11, 0)) { - warnings.push( - "ядро < 5.11: память BPF-карт ограничена RLIMIT_MEMLOCK — увеличьте `ulimit -l` или обновите ядро" - .to_owned(), - ); - } - if driver_name.is_none() { - warnings.push(format!( - "драйвер интерфейса '{interface}' не определён (нет device/driver symlink): ожидается generic mode" - )); - } - - Self { - kernel_version, - btf_available, - privileged, - driver_name, - warnings, - } - } - - /// Ожидаемый режим attach по таблице драйверов. - #[must_use] - pub fn attach_mode(&self) -> AttachMode { - match self.driver_name.as_deref() { - Some(driver) if driver_supports_native_xdp(driver) => AttachMode::Native, - _ => AttachMode::Generic, - } - } - - /// Человекочитаемое объяснение вердикта. - #[must_use] - pub fn verdict(&self) -> String { - let kernel = self - .kernel_version - .map_or_else(|| "неизвестна".to_owned(), |v| v.to_string()); - match (self.attach_mode(), self.driver_name.as_deref()) { - (AttachMode::Native, Some(driver)) => format!( - "ядро {kernel}: драйвер '{driver}' поддерживает native XDP — программа работает в драйвере, минимальные накладные расходы" - ), - (_, driver) => format!( - "ядро {kernel}: драйвер '{}' не поддерживает native XDP → будет generic mode, CPU дороже", - driver.unwrap_or("неизвестный") - ), - } - } - - /// Fail-fast проверка минимальной версии ядра. - /// - /// # Errors - /// Версия ядра не определена или ниже [`MIN_KERNEL`]. - pub fn validate(&self) -> Result<()> { - let Some(version) = self.kernel_version else { - bail!("не удалось определить версию ядра ({OSRELEASE_PATH}) — XDP attach отклонён"); - }; - if version < MIN_KERNEL { - bail!("ядро {version} ниже минимально поддерживаемой {MIN_KERNEL}: XDP attach отклонён, обновите ядро"); - } - Ok(()) - } -} - -/// Предстартовая диагностика перед загрузкой XDP-программы: -/// структурный отчёт в лог, предупреждения, fail-fast на старом ядре. -/// -/// # Errors -/// См. [`EnvironmentReport::validate`]. -#[cfg(feature = "xdp")] -pub(crate) fn preflight(interface: &str) -> Result<()> { - let report = EnvironmentReport::collect(&FilesystemProbe, interface); - let kernel = report - .kernel_version - .map_or_else(|| "unknown".to_owned(), |v| v.to_string()); - tracing::info!( - interface, - kernel = %kernel, - btf = report.btf_available, - privileged = report.privileged, - driver = report.driver_name.as_deref().unwrap_or("unknown"), - mode = %report.attach_mode(), - "XDP environment: {}", - report.verdict() - ); - for warning in &report.warnings { - tracing::warn!(interface, "{warning}"); - } - report.validate() -} diff --git a/src/xdp/filter/attach.rs b/src/xdp/filter/attach.rs new file mode 100644 index 0000000..3b076f8 --- /dev/null +++ b/src/xdp/filter/attach.rs @@ -0,0 +1,58 @@ +use anyhow::{Context, Result}; +use libbpf_rs::{MapCore, Object, OpenObject, RingBuffer, RingBufferBuilder}; +use std::ffi::CString; + +use crate::xdp::globals::{RODATA_MAP_NAME, XdpGlobals}; + +/// Builds the NUL-terminated C string libc's `if_nametoindex` requires. +/// Rejects interface names containing an interior NUL byte instead of +/// silently truncating at it (the `&str`-as-`*const c_char` UB we removed). +pub(super) fn interface_cstr(name: &str) -> Result { + CString::new(name).with_context(|| format!("interface name {name:?} contains an interior NUL byte")) +} + +/// Патчит карту `rampart_.rodata` (volatile const глобалы) ДО загрузки объекта. +pub(super) fn patch_rodata(open_obj: &mut OpenObject, globals: &XdpGlobals) -> Result<()> { + let image = globals.build_rodata_image(); + let mut map = open_obj + .maps_mut() + .find(|m| m.name() == RODATA_MAP_NAME) + .with_context(|| format!("map '{RODATA_MAP_NAME}' not found in XDP object"))?; + map.set_initial_value(&image) + .with_context(|| format!("failed to set initial value of '{RODATA_MAP_NAME}'")) +} + +pub(super) fn build_ringbuf(obj: &Object) -> Result> { + let map = obj + .maps() + .find(|m| m.name() == "events_map") + .context("events_map not found")?; + let mut builder = RingBufferBuilder::new(); + builder.add(&map, |data: &[u8]| { + if data.len() >= 16 { + let ty = u32::from_ne_bytes(data[0..4].try_into().expect("4 bytes for type")); + let ip4 = u32::from_ne_bytes(data[4..8].try_into().expect("4 bytes for ip")); + let val = u64::from_ne_bytes(data[8..16].try_into().expect("8 bytes for val")); + tracing::debug!(event = ty, src_ip = ip4, data = val, "xdp event"); + } + 0 + })?; + Ok(builder.build()?) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn interface_cstr_accepts_plain_names() { + let c = interface_cstr("eth0").expect("valid interface name"); + assert_eq!(c.as_bytes(), b"eth0"); + assert_eq!(c.as_bytes_with_nul(), b"eth0\0"); + } + + #[test] + fn interface_cstr_rejects_interior_nul() { + assert!(interface_cstr("et\0h0").is_err()); + } +} diff --git a/src/xdp/filter/maps.rs b/src/xdp/filter/maps.rs new file mode 100644 index 0000000..330ae43 --- /dev/null +++ b/src/xdp/filter/maps.rs @@ -0,0 +1,43 @@ +pub(super) fn unix_ns() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_nanos() as u64 +} + +/// Encodes an LPM-trie key for `blacklist_map`: prefix length in byte 0 and the +/// four IPv4 octets in bytes 4..8 (bytes 1..3 are padding the kernel ignores). +pub(super) fn blacklist_key(prefix_len: u8, v4: &[u8; 4]) -> [u8; 8] { + let mut key = [0u8; 8]; + key[0] = prefix_len; + key[4..8].copy_from_slice(v4); + key +} + +/// Absolute ban expiry in ns-since-epoch, saturating so an absurd duration +/// clamps to `u64::MAX` rather than wrapping into the past. +pub(super) fn blacklist_expiry_ns(now_ns: u64, duration_secs: u64) -> u64 { + now_ns.saturating_add(duration_secs.saturating_mul(1_000_000_000)) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn blacklist_key_puts_prefix_len_first_and_octets_last() { + assert_eq!(blacklist_key(24, &[10, 0, 1, 200]), [24, 0, 0, 0, 10, 0, 1, 200]); + assert_eq!(blacklist_key(32, &[127, 0, 0, 1]), [32, 0, 0, 0, 127, 0, 0, 1]); + } + + #[test] + fn blacklist_expiry_converts_secs_to_ns() { + assert_eq!(blacklist_expiry_ns(1_000, 2), 1_000 + 2_000_000_000); + } + + #[test] + fn blacklist_expiry_saturates_instead_of_wrapping() { + assert_eq!(blacklist_expiry_ns(u64::MAX, 1), u64::MAX); + assert_eq!(blacklist_expiry_ns(0, u64::MAX), u64::MAX); + } +} diff --git a/src/xdp/filter.rs b/src/xdp/filter/mod.rs similarity index 66% rename from src/xdp/filter.rs rename to src/xdp/filter/mod.rs index fcaff5e..21de55f 100644 --- a/src/xdp/filter.rs +++ b/src/xdp/filter/mod.rs @@ -1,13 +1,19 @@ use anyhow::{Context, Result, bail}; -use libbpf_rs::{MapCore, MapFlags, Object, ObjectBuilder, OpenObject, RingBuffer, RingBufferBuilder, Xdp, XdpFlags}; +use libbpf_rs::{MapCore, MapFlags, Object, ObjectBuilder, RingBuffer, Xdp, XdpFlags}; use std::net::IpAddr; use std::net::Ipv4Addr; use std::os::unix::io::AsFd; use super::XdpStats; -use super::globals::{RODATA_MAP_NAME, XdpGlobals}; +use super::globals::XdpGlobals; use crate::traffic::prefix::{PrefixKey, PrefixStatsVal}; +use attach::{build_ringbuf, interface_cstr, patch_rodata}; +use maps::{blacklist_expiry_ns, blacklist_key, unix_ns}; + +mod attach; +mod maps; + pub struct XdpFilter { obj: Option, ringbuf: Option>, @@ -34,7 +40,7 @@ impl XdpFilter { pub fn load(&mut self) -> Result<()> { // Диагностика ДО загрузки: fail-fast на неподдерживаемом ядре/драйвере. - super::diagnostics::preflight(&self.interface)?; + super::probe::preflight(&self.interface)?; let bpf_obj = include_bytes!(concat!(env!("OUT_DIR"), "/universal_filter.o")); let mut open_obj = ObjectBuilder::default() @@ -43,9 +49,15 @@ impl XdpFilter { patch_rodata(&mut open_obj, &self.globals)?; let obj = open_obj.load().context("Failed to load XDP object (verifier error?)")?; - let ifindex = unsafe { libc::if_nametoindex(self.interface.as_ptr() as *const libc::c_char) }; + let c_interface = interface_cstr(&self.interface)?; + // SAFETY: `c_interface` is a valid, NUL-terminated C string and is kept + // alive until the end of this scope (past the `if_nametoindex` call), so + // the raw pointer satisfies the `const char *` contract of the libc FFI. + // Rust `&str` is NOT NUL-terminated, hence the `CString` conversion above. + let ifindex = unsafe { libc::if_nametoindex(c_interface.as_ptr()) }; if ifindex == 0 { - bail!("interface '{}' not found", self.interface); + let err = std::io::Error::last_os_error(); + bail!("interface '{}' not found: {err}", self.interface); } let prog = obj @@ -64,9 +76,8 @@ impl XdpFilter { } pub fn unload(&mut self) -> Result<()> { - if self.ifindex != 0 { - let fd = unsafe { std::os::unix::io::BorrowedFd::borrow_raw(std::os::unix::io::RawFd::from(-1)) }; - let _ = Xdp::new(fd).detach(self.ifindex, XdpFlags::NONE); + if let Some(obj) = self.obj.as_ref() { + self.detach_program(obj)?; } self.ringbuf = None; self.obj = None; @@ -75,6 +86,28 @@ impl XdpFilter { Ok(()) } + /// Detaches the XDP program using the real program fd owned by the loaded + /// `Object` (mirroring `load`'s `Xdp::new(prog.as_fd()).attach(..)`). A + /// detach failure while a filter is active is surfaced loudly and returned, + /// never swallowed — otherwise the kernel keeps dropping packets we no + /// longer intend to filter. + fn detach_program(&self, obj: &Object) -> Result<()> { + let prog = obj + .progs() + .find(|p| p.name() == "rampart_universal_filter") + .context("XDP program 'rampart_universal_filter' not found while detaching")?; + if let Err(e) = Xdp::new(prog.as_fd()).detach(self.ifindex, XdpFlags::NONE) { + tracing::error!( + interface = %self.interface, + ifindex = self.ifindex, + error = %e, + "XDP detach failed while filter was active" + ); + return Err(e).with_context(|| format!("failed to detach XDP from '{}'", self.interface)); + } + Ok(()) + } + pub fn drain_events(&self) { if let Some(rb) = &self.ringbuf { let _ = rb.consume(); @@ -107,22 +140,15 @@ impl XdpFilter { fn blacklist_update(&self, prefix_len: u8, v4: &[u8; 4], duration_secs: u64) -> Result<()> { let map = self.find_map("blacklist_map")?; - let mut key = [0u8; 8]; - key[0] = prefix_len; - key[4..8].copy_from_slice(v4); - map.update( - &key, - &(unix_ns() + duration_secs * 1_000_000_000).to_le_bytes(), - MapFlags::ANY, - )?; + let key = blacklist_key(prefix_len, v4); + let expiry_ns = blacklist_expiry_ns(unix_ns(), duration_secs); + map.update(&key, &expiry_ns.to_le_bytes(), MapFlags::ANY)?; Ok(()) } pub fn unban_ip(&self, ip: Ipv4Addr) -> Result<()> { let map = self.find_map("blacklist_map")?; - let mut key = [0u8; 8]; - key[0] = 32; - key[4..8].copy_from_slice(&ip.octets()); + let key = blacklist_key(32, &ip.octets()); map.delete(&key)?; Ok(()) } @@ -175,6 +201,11 @@ impl XdpFilter { } } +// SAFETY: `XdpFilter` is only ever accessed from the single thread that owns it, +// held behind `Arc>`; that Mutex serialization is the invariant +// that makes moving it across threads sound. Auto-`Send` is not derived because +// `libbpf_rs::Object` and the `RingBuffer<'static>` wrap non-`Send` libbpf +// handles, so the marker is asserted manually under the ownership discipline above. unsafe impl Send for XdpFilter {} impl Drop for XdpFilter { @@ -182,39 +213,3 @@ impl Drop for XdpFilter { let _ = self.unload(); } } - -fn unix_ns() -> u64 { - std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .unwrap_or_default() - .as_nanos() as u64 -} - -/// Патчит карту `rampart_.rodata` (volatile const глобалы) ДО загрузки объекта. -fn patch_rodata(open_obj: &mut OpenObject, globals: &XdpGlobals) -> Result<()> { - let image = globals.build_rodata_image(); - let mut map = open_obj - .maps_mut() - .find(|m| m.name() == RODATA_MAP_NAME) - .with_context(|| format!("map '{RODATA_MAP_NAME}' not found in XDP object"))?; - map.set_initial_value(&image) - .with_context(|| format!("failed to set initial value of '{RODATA_MAP_NAME}'")) -} - -fn build_ringbuf(obj: &Object) -> Result> { - let map = obj - .maps() - .find(|m| m.name() == "events_map") - .context("events_map not found")?; - let mut builder = RingBufferBuilder::new(); - builder.add(&map, |data: &[u8]| { - if data.len() >= 16 { - let ty = u32::from_ne_bytes(data[0..4].try_into().expect("4 bytes for type")); - let ip4 = u32::from_ne_bytes(data[4..8].try_into().expect("4 bytes for ip")); - let val = u64::from_ne_bytes(data[8..16].try_into().expect("8 bytes for val")); - tracing::debug!(event = ty, src_ip = ip4, data = val, "xdp event"); - } - 0 - })?; - Ok(builder.build()?) -} diff --git a/src/xdp/mod.rs b/src/xdp/mod.rs index 0afe229..30ca8a8 100644 --- a/src/xdp/mod.rs +++ b/src/xdp/mod.rs @@ -1,10 +1,9 @@ //! Загрузчик XDP-программы (L3/L4-уровень защиты). -mod stats; -pub use stats::XdpStats; -mod diagnostics; -pub use diagnostics::{ - AttachMode, EnvironmentReport, FilesystemProbe, KernelVersion, MIN_KERNEL, SystemProbe, driver_supports_native_xdp, +mod probe; +pub use probe::{ + AttachMode, EnvironmentReport, FilesystemProbe, KernelVersion, MIN_KERNEL, SystemProbe, XdpStats, + driver_supports_native_xdp, }; #[cfg(feature = "xdp")] @@ -18,9 +17,7 @@ mod filter; pub use filter::XdpFilter; #[cfg(feature = "xdp")] -mod metrics; -#[cfg(feature = "xdp")] -pub use metrics::XdpMetrics; +pub use probe::XdpMetrics; #[cfg(not(feature = "xdp"))] mod noop; diff --git a/src/xdp/probe/diagnostics/mod.rs b/src/xdp/probe/diagnostics/mod.rs new file mode 100644 index 0000000..6e5ff8e --- /dev/null +++ b/src/xdp/probe/diagnostics/mod.rs @@ -0,0 +1,13 @@ +//! Диагностика окружения перед XDP-attach: ядро, BTF, драйвер NIC, привилегии. +//! +//! Основной источник боли при XDP — непонятные ошибки attach на неподдерживаемых +//! ядрах и драйверах. Модуль собирает [`EnvironmentReport`] ДО загрузки BPF-программы +//! и даёт человекочитаемый вердикт с предупреждениями. + +mod report; +mod types; + +pub use report::EnvironmentReport; +#[cfg(feature = "xdp")] +pub(crate) use report::preflight; +pub use types::{AttachMode, FilesystemProbe, KernelVersion, MIN_KERNEL, SystemProbe, driver_supports_native_xdp}; diff --git a/src/xdp/probe/diagnostics/report.rs b/src/xdp/probe/diagnostics/report.rs new file mode 100644 index 0000000..24f6f03 --- /dev/null +++ b/src/xdp/probe/diagnostics/report.rs @@ -0,0 +1,130 @@ +use anyhow::{Result, bail}; + +use super::types::{ + AttachMode, BTF_PATH, FilesystemProbe, KernelVersion, MIN_KERNEL, OSRELEASE_PATH, SystemProbe, driver_from_link, + driver_supports_native_xdp, +}; + +/// Структурированный отчёт о пригодности окружения для XDP. +pub struct EnvironmentReport { + pub kernel_version: Option, + pub btf_available: bool, + pub privileged: bool, + pub driver_name: Option, + pub warnings: Vec, +} + +impl EnvironmentReport { + /// Собирает отчёт через произвольную реализацию [`SystemProbe`]. + #[must_use] + pub fn collect(probe: &dyn SystemProbe, interface: &str) -> Self { + let kernel_version = probe + .read_file(OSRELEASE_PATH) + .ok() + .and_then(|r| KernelVersion::parse(r.trim())); + let btf_available = probe.path_exists(BTF_PATH); + let privileged = probe.effective_uid() == 0; + let driver_link = probe.symlink_target(&format!("/sys/class/net/{interface}/device/driver")); + let driver_name = driver_link.as_deref().and_then(driver_from_link); + + let mut warnings = Vec::new(); + if !privileged { + warnings.push( + "процесс не от root: CAP_BPF/CAP_NET_ADMIN не подтверждены, \ + attach скорее всего завершится EPERM" + .to_owned(), + ); + } + if !btf_available { + warnings.push(format!( + "{BTF_PATH} недоступен: CO-RE релокации невозможны, переносимость BPF-программы ограничена" + )); + } + if kernel_version.is_some_and(|v| v < KernelVersion::new(5, 11, 0)) { + warnings.push( + "ядро < 5.11: память BPF-карт ограничена RLIMIT_MEMLOCK — увеличьте `ulimit -l` или обновите ядро" + .to_owned(), + ); + } + if driver_name.is_none() { + warnings.push(format!( + "драйвер интерфейса '{interface}' не определён (нет device/driver symlink): ожидается generic mode" + )); + } + + Self { + kernel_version, + btf_available, + privileged, + driver_name, + warnings, + } + } + + /// Ожидаемый режим attach по таблице драйверов. + #[must_use] + pub fn attach_mode(&self) -> AttachMode { + match self.driver_name.as_deref() { + Some(driver) if driver_supports_native_xdp(driver) => AttachMode::Native, + _ => AttachMode::Generic, + } + } + + /// Человекочитаемое объяснение вердикта. + #[must_use] + pub fn verdict(&self) -> String { + let kernel = self + .kernel_version + .map_or_else(|| "неизвестна".to_owned(), |v| v.to_string()); + match (self.attach_mode(), self.driver_name.as_deref()) { + (AttachMode::Native, Some(driver)) => format!( + "ядро {kernel}: драйвер '{driver}' поддерживает native XDP — программа работает в драйвере, минимальные накладные расходы" + ), + (_, driver) => format!( + "ядро {kernel}: драйвер '{}' не поддерживает native XDP → будет generic mode, CPU дороже", + driver.unwrap_or("неизвестный") + ), + } + } + + /// Fail-fast проверка минимальной версии ядра. + /// + /// # Errors + /// Версия ядра не определена или ниже [`MIN_KERNEL`]. + pub fn validate(&self) -> Result<()> { + let Some(version) = self.kernel_version else { + bail!("не удалось определить версию ядра ({OSRELEASE_PATH}) — XDP attach отклонён"); + }; + if version < MIN_KERNEL { + bail!("ядро {version} ниже минимально поддерживаемой {MIN_KERNEL}: XDP attach отклонён, обновите ядро"); + } + Ok(()) + } +} + +/// Предстартовая диагностика перед загрузкой XDP-программы: +/// структурный отчёт в лог, предупреждения, fail-fast на старом ядре. +/// +/// # Errors +/// См. [`EnvironmentReport::validate`]. +#[cfg(feature = "xdp")] +pub(crate) fn preflight(interface: &str) -> Result<()> { + let report = EnvironmentReport::collect(&FilesystemProbe, interface); + let kernel = report + .kernel_version + .map_or_else(|| "unknown".to_owned(), |v| v.to_string()); + tracing::info!( + interface, + kernel = %kernel, + btf = report.btf_available, + privileged = report.privileged, + driver = report.driver_name.as_deref().unwrap_or("unknown"), + mode = %report.attach_mode(), + "XDP environment: {}", + report.verdict() + ); + for warning in &report.warnings { + tracing::warn!(interface, "{warning}"); + } + report.validate() +} diff --git a/src/xdp/probe/diagnostics/types.rs b/src/xdp/probe/diagnostics/types.rs new file mode 100644 index 0000000..64663fd --- /dev/null +++ b/src/xdp/probe/diagnostics/types.rs @@ -0,0 +1,154 @@ +use std::fmt; + +/// Минимально поддерживаемая версия ядра: 5.15 LTS. +pub const MIN_KERNEL: KernelVersion = KernelVersion::new(5, 15, 0); + +pub(super) const OSRELEASE_PATH: &str = "/proc/sys/kernel/osrelease"; +pub(super) const BTF_PATH: &str = "/sys/kernel/btf/vmlinux"; +const STATUS_PATH: &str = "/proc/self/status"; + +/// Драйверы с известной поддержкой native XDP (по данным xdpgeneric/native +/// матриц upstream-ядра). Отсутствие в списке не означает отсутствие поддержки, +/// но для таких драйверов предсказываем generic mode. +const NATIVE_XDP_DRIVERS: &[&str] = &[ + "virtio_net", + "ixgbe", + "ixgbevf", + "i40e", + "iavf", + "ice", + "mlx4_core", + "mlx5_core", + "igb", + "igc", + "e1000e", + "vmxnet3", + "bnxt_en", + "nfp", + "sfc", + "ena", + "hv_netvsc", + "macvlan", +]; + +/// Версия ядра Linux, пригодная для сравнения (`major.minor.patch`). +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +pub struct KernelVersion { + pub major: u16, + pub minor: u16, + pub patch: u16, +} + +impl KernelVersion { + #[must_use] + pub const fn new(major: u16, minor: u16, patch: u16) -> Self { + Self { major, minor, patch } + } + + /// Парсит строку формата `uname -r`. Берёт первые три числовых компонента, + /// дистрибутивные суффиксы (`6.8.0-45-generic`, `5.15.0-rc2`) отбрасываются. + /// Строка без цифр не парсится. + #[must_use] + pub fn parse(release: &str) -> Option { + let mut nums = release + .split(|c: char| !c.is_ascii_digit()) + .filter_map(|part| part.parse::().ok()); + let major = nums.next()?; + let minor = nums.next().unwrap_or(0); + let patch = nums.next().unwrap_or(0); + Some(Self::new(major, minor, patch)) + } +} + +impl fmt::Display for KernelVersion { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "{}.{}.{}", self.major, self.minor, self.patch) + } +} + +/// Режим присоединения XDP-программы к интерфейсу. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum AttachMode { + /// Программа исполняется в контексте драйвера NIC — минимальные накладные расходы. + Native, + /// Программа вызывается из сетевого стека после ingress — работает везде, дороже по CPU. + Generic, + /// Аппаратная разгрузка в NIC — требует явного включения и поддержки железа, + /// автоматически не выбирается. + Offload, +} + +impl fmt::Display for AttachMode { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str(match self { + Self::Native => "native", + Self::Generic => "generic", + Self::Offload => "offload", + }) + } +} + +/// Абстракция над procfs/sysfs: реальная реализация — [`FilesystemProbe`], +/// в тестах — мок. Изолирует диагностику от реального ядра. +pub trait SystemProbe { + /// Читает файл целиком. + /// + /// # Errors + /// Проксирует ошибку чтения файла. + fn read_file(&self, path: &str) -> std::io::Result; + + #[must_use] + fn path_exists(&self, path: &str) -> bool; + + #[must_use] + fn symlink_target(&self, path: &str) -> Option; + + #[must_use] + fn effective_uid(&self) -> u32; +} + +/// Реальная реализация [`SystemProbe`] поверх `std::fs`. +pub struct FilesystemProbe; + +impl SystemProbe for FilesystemProbe { + fn read_file(&self, path: &str) -> std::io::Result { + std::fs::read_to_string(path) + } + + fn path_exists(&self, path: &str) -> bool { + std::path::Path::new(path).exists() + } + + fn symlink_target(&self, path: &str) -> Option { + Some(std::fs::read_link(path).ok()?.to_string_lossy().into_owned()) + } + + /// Эвристика: эффективный uid берётся из `/proc/self/status`. + /// root почти всегда имеет CAP_BPF/CAP_NET_ADMIN (если явно не урезаны); + /// для не-root подтверждение возможно только попыткой attach. + /// Нечитаемый статус трактуем как непривилегированный процесс. + fn effective_uid(&self) -> u32 { + std::fs::read_to_string(STATUS_PATH) + .ok() + .and_then(|status| parse_euid(&status)) + .unwrap_or(1) + } +} + +fn parse_euid(status: &str) -> Option { + status.lines().find_map(|line| { + let mut fields = line.strip_prefix("Uid:")?.split_whitespace(); + fields.nth(1)?.parse().ok() + }) +} + +pub(super) fn driver_from_link(target: &str) -> Option { + let name = target.rsplit('/').next()?; + (!name.is_empty()).then(|| name.to_owned()) +} + +/// Поддерживает ли драйвер native XDP (по таблице известных драйверов). +#[must_use] +pub fn driver_supports_native_xdp(driver: &str) -> bool { + NATIVE_XDP_DRIVERS.contains(&driver) +} diff --git a/src/xdp/metrics.rs b/src/xdp/probe/metrics.rs similarity index 100% rename from src/xdp/metrics.rs rename to src/xdp/probe/metrics.rs diff --git a/src/xdp/probe/mod.rs b/src/xdp/probe/mod.rs new file mode 100644 index 0000000..2513829 --- /dev/null +++ b/src/xdp/probe/mod.rs @@ -0,0 +1,16 @@ +//! Зондирование XDP-окружения: диагностика, статистика и метрики. + +mod diagnostics; +#[cfg(feature = "xdp")] +pub(crate) use diagnostics::preflight; +pub use diagnostics::{ + AttachMode, EnvironmentReport, FilesystemProbe, KernelVersion, MIN_KERNEL, SystemProbe, driver_supports_native_xdp, +}; + +mod stats; +pub use stats::XdpStats; + +#[cfg(feature = "xdp")] +mod metrics; +#[cfg(feature = "xdp")] +pub use metrics::XdpMetrics; diff --git a/src/xdp/stats.rs b/src/xdp/probe/stats.rs similarity index 100% rename from src/xdp/stats.rs rename to src/xdp/probe/stats.rs