fix: XDP unsafe loader + Redis sync; enforce 250-line/4-file layout limits

- 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
This commit is contained in:
loki5512344 2026-09-15 23:55:17 +02:00
parent d6bcae54c8
commit aa615a1141
Signed by: boba
GPG key ID: 253067914055423B
64 changed files with 2052 additions and 1408 deletions

View file

@ -40,7 +40,17 @@ jobs:
- uses: actions/checkout@v4 - uses: actions/checkout@v4
- uses: dtolnay/rust-toolchain@stable - uses: dtolnay/rust-toolchain@stable
- uses: Swatinem/rust-cache@v2 - 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: rust-deny:
name: Rust — cargo-deny name: Rust — cargo-deny

3
.module_size_baseline Normal file
View file

@ -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

View file

@ -3,10 +3,10 @@
CARGO = cargo CARGO = cargo
TARGET_DIR = target 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: ebpf ebpf-clean
.PHONY: docker docker-build docker-up docker-down docker-logs .PHONY: docker docker-build docker-up docker-down docker-logs
.PHONY: ci ci-full .PHONY: ci ci-full repo-gates
all: check test build all: check test build
@ -27,6 +27,9 @@ check-all:
test: test:
$(CARGO) test $(CARGO) test
test-all:
$(CARGO) test --all-features
fmt: fmt:
$(CARGO) fmt --all $(CARGO) fmt --all
@ -75,7 +78,11 @@ docker-logs:
checkstyle: fmt-check clippy checkstyle: fmt-check clippy
@echo "✓ Checkstyle passed (rustfmt + 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" @echo "✓ CI passed"
ci-full: ci ebpf ci-full: ci ebpf

View file

@ -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 │ │ L7 handshake analysis · rate limit · HMAC · death-code patterns │
├────────────────────────────────────────────────────────────────────┤ ├────────────────────────────────────────────────────────────────────┤
│ Layer 4: Protocol Plugins (feature crates) │ │ 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) Traffic Intel (EWMA thresholds, profiling, reputation)
runs across all layers runs across all layers
@ -69,12 +69,13 @@ Attacker → [XDP/eBPF] → [PoW] → [Userspace Core] → [Plugin] → Your Ser
| Component | Role | Stack | | 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-manager** | Management API + Redis sync | Rust (axum, redis) |
| **rampart-cli** | CLI tool for operators | Rust (clap) | | **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 | | **protocol plugins** | Protocol-aware filtering as feature crates | Rust |
| ↳ `minecraft` | First plugin (MC handshake analysis) | Rust | | ↳ `http` | HTTP/1.1 handler, compiled with the `protocol-http` feature | Rust |
| ↳ `http`, `grpc` | Planned | Rust | | ↳ `grpc` | Planned | Rust |
| **docs/kb** | Bilingual knowledge base: attack anatomy, defense levels, practice guides | Markdown | | **docs/kb** | Bilingual knowledge base: attack anatomy, defense levels, practice guides | Markdown |
### Performance ### Performance
@ -92,14 +93,15 @@ Full benchmark suite in progress.
### Quick Start ### Quick Start
```bash ```bash
# Build # Build (the HTTP protocol handler is a feature)
cargo build --release cargo build --release --features protocol-http
# Create config # Install the default config
rampart config init > /etc/rampart/config.toml sudo mkdir -p /etc/rampart
sudo cp deploy/config/edge.toml /etc/rampart/config.toml
# Run edge node # Run edge node (config path comes from $RAMPART_CONFIG, default /etc/rampart/config.toml)
./target/release/rampart-core --config /etc/rampart/config.toml RAMPART_CONFIG=/etc/rampart/config.toml ./target/release/rampart
``` ```
### Documentation ### Documentation
@ -118,8 +120,6 @@ rampart config init > /etc/rampart/config.toml
- Stabilize the protocol plugin API - Stabilize the protocol plugin API
- BPF hook modules for deep protocol parsing in XDP - 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 паттерны │ │ L7 handshake analysis · rate limit · HMAC · death-code паттерны │
├────────────────────────────────────────────────────────────────────┤ ├────────────────────────────────────────────────────────────────────┤
│ Слой 4: Протокол-плагины (feature crates) │ │ Слой 4: Протокол-плагины (feature crates) │
│ minecraft (первый плагин) · http (в планах) · grpc (в планах) │ │ http (feature protocol-http) · grpc (в планах) │
└────────────────────────────────────────────────────────────────────┘ └────────────────────────────────────────────────────────────────────┘
Traffic Intel (EWMA thresholds, профилирование, репутация) 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-manager** | Management API + Redis sync | Rust (axum, redis) |
| **rampart-cli** | CLI для операторов | Rust (clap) | | **rampart-cli** | CLI для операторов | Rust (clap) |
| **rampart-tui** | Терминальный дашборд live-метрик, опрашивает Prometheus `/metrics` | Rust (ratatui) |
| **Протокол-плагины** | Протоколозависимая фильтрация в виде feature crates | Rust | | **Протокол-плагины** | Протоколозависимая фильтрация в виде feature crates | Rust |
| ↳ `minecraft` | Первый плагин (анализ MC-handshake) | Rust | | ↳ `http` | HTTP/1.1-обработчик, собирается с фичей `protocol-http` | Rust |
| ↳ `http`, `grpc` | В планах | Rust | | ↳ `grpc` | В планах | Rust |
| **docs/kb** | Двуязычная база знаний: анатомия атак, уровни защиты, практические руководства | Markdown | | **docs/kb** | Двуязычная база знаний: анатомия атак, уровни защиты, практические руководства | Markdown |
### Производительность ### Производительность
@ -191,14 +192,15 @@ Rampart фильтрует трафик на трёх уровнях до тог
### Быстрый старт ### Быстрый старт
```bash ```bash
# Сборка # Сборка (HTTP-обработчик собирается фичей)
cargo build --release 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 ноды # Запуск edge ноды (путь конфига берётся из $RAMPART_CONFIG, по умолчанию /etc/rampart/config.toml)
./target/release/rampart-core --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 протокол-плагинов - Стабилизация API протокол-плагинов
- BPF hook модули для глубокого парсинга протоколов в XDP - BPF hook модули для глубокого парсинга протоколов в XDP
- Терминальный интерфейс (ratatui TUI)
- HTTP протокол-плагин
--- ---

126
TODO.md
View file

@ -12,7 +12,10 @@
### KISS ### 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<u8>`. - Не используй generics где хватит `&str` и `Vec<u8>`.
### DRY ### DRY
@ -44,19 +47,24 @@
``` ```
guard/ guard/
├── Cargo.toml # ОДИН пакет rampart, features = ["protocol-http", ...] ├── Cargo.toml # ОДИН пакет rampart, features = ["protocol-http", ...]
├── scripts/ # check_module_size.sh (гейт 250/4), check_default_secrets.sh
├── src/ ├── src/
│ ├── bin/{rampart, rampart-manager, rampart-cli}.rs │ ├── bin/{rampart, rampart-manager, rampart-cli, rampart-tui}.rs # тонкие, логика в lib
│ ├── engine/ # listener, tunnel (generic TCP proxy), challenge (PoW) │ ├── app/ # thin-main: runtime wiring, services (spawn-циклы)
│ ├── engine/ # listener, tunnel; challenge/ (pow+difficulty), subnet/ (tracker+monitor)
│ ├── filter/ # blacklist, rate_limit, geo — trait Filter │ ├── 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) │ ├── store/ # redis (+ trait StateStore)
│ ├── manager/ # api/, auth/, sync/ │ ├── manager/ # api/ (+ api/inventory/), auth/, sync/
│ ├── cli/ # команды CLI │ ├── cli/ # команды; commands/node/ (status,drain,emergency)
│ └── protocol/ # trait ProtocolHandler + registry (реализаций пока 0) │ ├── 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/ ├── 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-протокол-хуков │ └── hooks/hook_api.h # контракт подключаемых BPF-протокол-хуков
├── tests/ # интеграционные ├── tests/ # интеграционные (вне гейта 4-файлов: cargo требует 1 файл = 1 бинарь)
└── docs/ # kb/ (knowledge base) + research/ + ops-доки └── docs/ # kb/ (knowledge base) + research/ + ops-доки
``` ```
@ -78,46 +86,71 @@ guard/
## 2. Ближайшие задачи (v0.3) ## 2. Ближайшие задачи (v0.3)
### Subnet-level detection (ботнет с ротацией IP) ### Subnet-level detection (ботнет с ротацией IP) — ✅ готово (2026-08/09)
- [ ] **XDP**: карта `prefix_stats` (LRU_HASH, ключ /24 v4 | /64 v6) — счётчики SYN/pps - [x] **XDP**: карта `prefix_stats` (LRU_HASH, ключ /24 v4 | /64 v6) — счётчики SYN/pps
per-префикс рядом с per-IP (референс: caddy-mitigator CIDR promotion, lnvps_fw carpet-bomb). per-префикс (xdp/core/prefix_stats.h, трафик-слой: src/traffic/prefix/).
- [ ] **Detector**: префикс превышает порог при том что отдельные IP под лимитом - [x] **Detector**: превышение порога префиксом при IP под лимитом → флаг подсети
→ распределённая атака → флаг подсети. (src/traffic/intel/detector.rs, src/engine/subnet/).
- [ ] **Мягкая эскалация для подсетей**: monitor → strict limits → challenge → блок. - [x] **Мягкая эскалация**: monitor → strict_limit → challenge → block
Хард-бан /24 только через challenge (CGNAT: за одним /24 легитимно живут сотни людей). (SubnetVerdict-лестница; блок /24 идёт через ban_cidr, не слепой hard-ban).
- [ ] Блок самой подсети — уже умеем: `blacklist_map` это LPM trie (CIDR из коробки). - [x] Блок подсети через `blacklist_map` (LPM trie) — `XdpFilter::ban_cidr`.
### Движок без протоколов — сделать полезным ### Движок без протоколов — сделать полезным — ✅ готово
- [ ] **Первый протокол-плагин**: `protocol-http` (feature) — минимальный HTTP/1.1 - [x] **Первый протокол-плагин**: `protocol-http` — HTTP/1.1 request-head анализ
handshake-анализ (request line, заголовки, размер), чтобы edge-нода заработала (src/protocol/http/, tests/http_protocol.rs).
для веб-сервисов. - [x] **TCP-proxy режим**: tunnel.rs + ProtocolHandler интегрированы в listener/app wiring.
- [ ] **TCP-proxy режим**: generic upstream forwarding за ProtocolHandler - [x] **Fail-fast сообщение** при пустом registry (`ProtocolRegistry::primary()`).
(tunnel.rs уже generic — проверить интеграцию).
- [ ] **Fail-fast сообщение** при пустом registry — улучшить текст подсказки сборки.
### Подключение мёртвого интеллекта (правило: «мёртвый код = баг») ### Подключение мёртвого интеллекта (правило: «мёртвый код = баг») — ✅ готово
- [ ] Layer Traffic Intel подключить в hot path: AttackDetector/IpReputation → - [x] Traffic Intel в hot path: AttackDetector/IpReputation → метрики + auto-ban
метрики + auto-ban (сейчас не вызывается). (src/traffic/hook.rs, src/engine/tunnel.rs, src/app/services.rs).
- [ ] Blacklist: `clear_expired()` по таймеру. - [x] Blacklist: `clear_expired()` по таймеру (src/app/).
- [ ] RateLimiter: TTL-эвикция idle bucket'ов + cap карты. - [x] RateLimiter: TTL-эвикция idle bucket'ов.
### Безопасность (перенос из аудита v0.3, актуальное) ### Безопасность — ✅ кроме ролей
- [ ] Rate limiter на login endpoint manager API (5/60с). - [x] Rate limiter на login endpoint manager API (per-IP, tests в api/auth.rs).
- [ ] JWT: валидация ролей/audience, secret ≥ 32 байт. - [x] JWT: audience-валидация, secret ≥ 32 байт fail-fast (rampart-manager.rs).
- [ ] Redis: `KEYS` → `SCAN`, reconnect pubsub-подписчика. - [ ] JWT-роли: сейчас единственный hardcoded `admin`; RBAC-ролей нет — либо убрать поле,
либо делать роли (решение отложить до второго потребителя).
- [x] Redis: `KEYS` → `SCAN` (scan_options), pubsub-подписчик — честный reconnect
с экспоненциальным backoff (src/store/redis.rs).
### XDP ### XDP
- [ ] Verifier-проверка на реальном ядре (в контейнере нет CAP_BPF — компиляция OK, - [x] Verifier-проверка на реальном ядре (live-тест 2026-09: два бага RST-challenge найдены и
загрузка не проверялась). исправлены, d6bcae5).
- [ ] Rust loader (`src/xdp/`): пути к xdp/core/universal_filter.c, patch глобалов - [x] Rust loader (`src/xdp/filter/`): open/load/attach, patch глобалов G_* из config.toml
G_* из config.toml, ringbuf events → blacklist. (src/xdp/globals.rs), CString-safe if_nametoindex, реальный detach по fd прогрессы.
- [ ] Smoke-test attach в CI (VM runner с CAP_BPF). - [ ] 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 — переписать под новую структуру - [x] docs/deployment.md, configuration.md, runbook.md — старых крейтов/MC не упоминают (grep чисто).
(сейчас упоминают старые крейты/MC). - [x] docs/kb/README.md — индекс KB со ссылками.
- [ ] docs/kb/README.md — индекс KB со ссылками на все статьи. - [x] TUI (ratatui): live-метрики из Prometheus endpoint — готово (src/tui/, 564613b).
- [ ] TUI (ratatui): live-метрики из Prometheus endpoint (planned, v0.4).
## 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 ## 3. Backlog
@ -190,7 +223,8 @@ guard/
5. **По умолчанию безопасно**: нет дефолтных секретов; отсутствие обязательного env = fail-fast. 5. **По умолчанию безопасно**: нет дефолтных секретов; отсутствие обязательного env = fail-fast.
6. **Интеграционный тест на слой**: config parse, filter logic, registry 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`, 8. **CI guardrails**: `cargo clippy --all-targets -- -D warnings`, `cargo test`,
clang-build xdp/core/universal_filter.c, grep на `changeme`. clang-build xdp/core/universal_filter.c, grep на `changeme`.
9. **README/TODO не врут**: каждое число имеет ссылку на тест или отчёт. 9. **README/TODO не врут**: каждое число имеет ссылку на тест или отчёт.
@ -201,7 +235,7 @@ guard/
☐ cargo check / cargo test проходят ☐ cargo check / cargo test проходят
☐ cargo clippy --all-targets -- -D warnings — 0 warnings ☐ cargo clippy --all-targets -- -D warnings — 0 warnings
☐ cargo fmt --check проходит ☐ cargo fmt --check проходит
☐ Ни один модуль не превышает 300 строк ☐ Ни один файл не превышает 250 строк; ни в одной папке src/ больше 4 .rs-файлов
☐ Unit тесты покрывают happy path + 2+ error cases ☐ Unit тесты покрывают happy path + 2+ error cases
☐ Нет мёртвого кода: pub без вызовов, конфиг-поле без потребителя, метрика без writer ☐ Нет мёртвого кода: pub без вызовов, конфиг-поле без потребителя, метрика без writer
☐ Нет дефолтных секретов ☐ Нет дефолтных секретов
@ -214,7 +248,7 @@ guard/
``` ```
❌ Тесты после кода. Пиши вместе. ❌ Тесты после кода. Пиши вместе.
❌ Модуль > 300 строк — сигнал декомпозировать немедленно. ❌ Файл > 250 строк или > 4 .rs в папке — сигнал декомпозировать/сгруппировать немедленно.
❌ TODO в коде без issue. ❌ TODO в коде без issue.
❌ Мёртвый код: pub без вызовов, конфиг-поле без потребителя, метрика без writer. ❌ Мёртвый код: 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)*

View file

@ -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

161
scripts/check_module_size.sh Executable file
View file

@ -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

6
src/app/mod.rs Normal file
View file

@ -0,0 +1,6 @@
//! Wiring-логика edge-демона: сборка сервисов и главный цикл запуска.
pub mod runtime;
pub mod services;
pub use runtime::run;

229
src/app/runtime.rs Normal file
View file

@ -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<Arc<tokio::sync::Mutex<ClickHouseWriter>>> = 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<Arc<tokio::sync::Mutex<ClickHouseWriter>>>, 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() => {}
}
}

119
src/app/services.rs Normal file
View file

@ -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<Config>,
shutdown_rx: &watch::Receiver<bool>,
) -> anyhow::Result<Option<Arc<Mutex<XdpFilter>>>> {
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<Config>,
_shutdown_rx: &watch::Receiver<bool>,
) -> anyhow::Result<Option<Arc<Mutex<XdpFilter>>>> {
Ok(None)
}
/// Юзерспейс-агрегатор префиксов: активен, когда detect.prefix включён,
/// а XDP-путь не работает (не собран или выключен в конфиге).
pub(crate) fn start_subnet_tracker(config: &Config) -> Option<Arc<SubnetTracker>> {
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<Arc<HashSet<IpAddr>>> {
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))
}

View file

@ -55,7 +55,7 @@ async fn main() -> anyhow::Result<()> {
let cli = Cli::parse(); let cli = Cli::parse();
match cli.command { match cli.command {
Commands::Status => commands::status::run().await, Commands::Status => commands::node::status::run().await,
Commands::Doctor => commands::doctor::run().await, Commands::Doctor => commands::doctor::run().await,
Commands::Config { key, value } => commands::config::run(key, value).await, Commands::Config { key, value } => commands::config::run(key, value).await,
Commands::Blacklist { action } => match action { Commands::Blacklist { action } => match action {
@ -64,9 +64,9 @@ async fn main() -> anyhow::Result<()> {
BlacklistAction::List => commands::blacklist::list().await, BlacklistAction::List => commands::blacklist::list().await,
}, },
Commands::Emergency { mode } => match mode { Commands::Emergency { mode } => match mode {
EmergencyMode::Enable => commands::emergency::enable().await, EmergencyMode::Enable => commands::node::emergency::enable().await,
EmergencyMode::Disable => commands::emergency::disable().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,
} }
} }

View file

@ -46,16 +46,16 @@ async fn main() -> anyhow::Result<()> {
tokio::spawn(sync::heartbeat::start_heartbeat_check(state.clone())); tokio::spawn(sync::heartbeat::start_heartbeat_check(state.clone()));
let public = Router::new() 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)); .route("/api/v1/auth/login", post(api::auth::login));
let protected = Router::new() let protected = Router::new()
.route("/api/v1/servers", get(api::servers::list_servers)) .route("/api/v1/servers", get(api::inventory::servers::list_servers))
.route( .route(
"/api/v1/blacklist", "/api/v1/blacklist",
get(api::blacklist::list_blacklist).post(api::blacklist::add_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)); .route_layer(middleware::from_fn(auth::auth_middleware));
let cors = match std::env::var("CORS_ORIGIN") { let cors = match std::env::var("CORS_ORIGIN") {

View file

@ -1,330 +1,6 @@
use rampart::config::Config; use rampart::app;
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,
)));
}
#[tokio::main] #[tokio::main]
async fn main() -> anyhow::Result<()> { async fn main() -> anyhow::Result<()> {
tracing_subscriber::fmt() app::run().await
.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<Arc<tokio::sync::Mutex<ClickHouseWriter>>> = 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<Arc<tokio::sync::Mutex<ClickHouseWriter>>>, 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<Config>,
shutdown_rx: &watch::Receiver<bool>,
) -> anyhow::Result<Option<Arc<Mutex<XdpFilter>>>> {
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<Config>,
_shutdown_rx: &watch::Receiver<bool>,
) -> anyhow::Result<Option<Arc<Mutex<XdpFilter>>>> {
Ok(None)
}
/// Юзерспейс-агрегатор префиксов: активен, когда detect.prefix включён,
/// а XDP-путь не работает (не собран или выключен в конфиге).
fn start_subnet_tracker(config: &Config) -> Option<Arc<SubnetTracker>> {
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<Arc<HashSet<IpAddr>>> {
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() => {}
}
} }

View file

@ -1,6 +1,4 @@
pub mod blacklist; pub mod blacklist;
pub mod config; pub mod config;
pub mod doctor; pub mod doctor;
pub mod drain; pub mod node;
pub mod emergency;
pub mod status;

View file

@ -0,0 +1,4 @@
//! CLI-команды управления edge-узлами: статус, drain, аварийный режим.
pub mod drain;
pub mod emergency;
pub mod status;

View file

@ -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<String>,
}
impl Default for BackendConfig {
fn default() -> Self {
Self {
upstreams: default_upstreams(),
}
}
}
fn default_upstreams() -> Vec<String> {
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<String>,
#[serde(default = "default_blacklist_cache_ttl")]
pub blacklist_cache_ttl_secs: u64,
pub clickhouse_url: Option<String>,
}
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<String>,
#[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<String>,
}
/// Секция `[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<String>,
#[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
}

View file

@ -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<String>,
}
/// Секция `[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<String>,
#[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
}

135
src/config/sections/edge.rs Normal file
View file

@ -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<String>,
}
impl Default for BackendConfig {
fn default() -> Self {
Self {
upstreams: default_upstreams(),
}
}
}
fn default_upstreams() -> Vec<String> {
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
}

View file

@ -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};

View file

@ -0,0 +1,148 @@
use serde::Deserialize;
#[derive(Debug, Clone, Deserialize)]
pub struct StoreConfig {
pub redis_url: Option<String>,
#[serde(default = "default_blacklist_cache_ttl")]
pub blacklist_cache_ttl_secs: u64,
pub clickhouse_url: Option<String>,
}
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<String>,
#[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
}

View file

@ -0,0 +1,88 @@
use crate::metrics;
use std::collections::VecDeque;
use std::time::Instant;
/// Адаптирует сложность PoW к текущему темпу подключений.
pub struct DifficultyAdjuster {
window: VecDeque<Instant>,
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);
}
}

View file

@ -0,0 +1,8 @@
//! Универсальный SHA-256 hashcash: генерация challenge, решатель,
//! верификатор и адаптивная сложность.
mod difficulty;
mod pow;
pub use difficulty::DifficultyAdjuster;
pub use pow::{Challenge, enforce, solve};

View file

@ -1,10 +1,5 @@
//! Универсальный SHA-256 hashcash: генерация challenge, решатель,
//! верификатор и адаптивная сложность.
use crate::metrics;
use rand::RngCore; use rand::RngCore;
use sha2::{Digest, Sha256}; use sha2::{Digest, Sha256};
use std::collections::VecDeque;
use std::net::IpAddr; use std::net::IpAddr;
use std::time::{Duration, Instant}; use std::time::{Duration, Instant};
use subtle::ConstantTimeEq; use subtle::ConstantTimeEq;
@ -92,70 +87,6 @@ pub fn solve(challenge: &str, difficulty: u8) -> Option<String> {
None None
} }
/// Адаптирует сложность PoW к текущему темпу подключений.
pub struct DifficultyAdjuster {
window: VecDeque<Instant>,
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 и проверяет ответ. /// Проводит текстовый PoW-gate в потоке: выдаёт challenge и проверяет ответ.
/// ///
/// # Errors /// # Errors
@ -241,20 +172,4 @@ mod tests {
let big_nonce = "0".repeat(MAX_NONCE_LEN + 1); let big_nonce = "0".repeat(MAX_NONCE_LEN + 1);
assert!(!challenge.verify(&big_nonce)); 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);
}
} }

View file

@ -2,6 +2,9 @@
pub mod challenge; pub mod challenge;
pub mod listener; pub mod listener;
pub mod subnet_monitor; pub mod subnet;
pub mod subnet_tracker;
pub mod tunnel; pub mod tunnel;
// Совместимость путей после переезда в subnet/: `engine::subnet_tracker` /
// `engine::subnet_monitor` остаются валидными алиасами модулей.
pub use subnet::{monitor as subnet_monitor, tracker as subnet_tracker};

3
src/engine/subnet/mod.rs Normal file
View file

@ -0,0 +1,3 @@
//! Subnet-level детектор: юзерспейс-агрегатор префиксов и монитор вердиктов.
pub mod monitor;
pub mod tracker;

View file

@ -1,4 +1,5 @@
//! Rampart — универсальная платформа сетевой защиты (L3/L4/L7). //! Rampart — универсальная платформа сетевой защиты (L3/L4/L7).
pub mod app;
pub mod cli; pub mod cli;
pub mod config; pub mod config;
pub mod engine; pub mod engine;

View file

@ -0,0 +1,4 @@
//! Инвентарь manager-API: health, узлы и серверы.
pub mod health;
pub mod nodes;
pub mod servers;

View file

@ -1,5 +1,3 @@
pub mod auth; pub mod auth;
pub mod blacklist; pub mod blacklist;
pub mod health; pub mod inventory;
pub mod nodes;
pub mod servers;

View file

@ -1,8 +1,17 @@
use crate::manager::AppState; use crate::manager::AppState;
use futures::StreamExt;
use redis::AsyncCommands; use redis::AsyncCommands;
use redis::ScanOptions;
use std::sync::Arc; use std::sync::Arc;
use tokio::time::{Duration, interval}; 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<AppState>) { pub async fn start_heartbeat_check(state: Arc<AppState>) {
let mut ticker = interval(Duration::from_secs(30)); let mut ticker = interval(Duration::from_secs(30));
loop { loop {
@ -15,30 +24,99 @@ pub async fn start_heartbeat_check(state: Arc<AppState>) {
async fn check_nodes(state: &AppState) -> anyhow::Result<()> { async fn check_nodes(state: &AppState) -> anyhow::Result<()> {
let mut conn = state.redis_client.get_multiplexed_async_connection().await?; let mut conn = state.redis_client.get_multiplexed_async_connection().await?;
let keys: Vec<String> = 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(); let now = chrono::Utc::now().timestamp();
for key in &keys { for key in &keys {
let raw: Option<String> = conn.get(key).await?; let raw: Option<String> = conn.get(key).await?;
if let Some(json) = raw let Some(updated) = raw.and_then(|json| offline_update(&json, now)) else {
&& let Ok(mut node) = serde_json::from_str::<serde_json::Value>(&json) continue;
{ };
let hb = node["last_heartbeat"] tracing::warn!("Node {key} is offline (heartbeat expired)");
.as_str() let _: () = conn.set(key.as_str(), updated).await.unwrap_or_default();
.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)");
}
}
} }
Ok(()) Ok(())
} }
/// Incrementally iterates the keyspace via SCAN (cursor-based, non-blocking).
async fn collect_node_keys(conn: &mut redis::aio::MultiplexedConnection) -> anyhow::Result<Vec<String>> {
let opts = ScanOptions::default()
.with_pattern(NODES_KEY_PATTERN)
.with_count(SCAN_BATCH);
let mut iter = conn.scan_options::<String>(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<String> {
let mut node = serde_json::from_str::<serde_json::Value>(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:*");
}
}

View file

@ -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<bool>) {
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( pub async fn start_blacklist_sync(
client: &redis::Client, client: &redis::Client,
blacklist: Arc<Blacklist>, blacklist: Arc<Blacklist>,
mut shutdown: watch::Receiver<bool>, mut shutdown: watch::Receiver<bool>,
) { ) {
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<bool>,
backoff: &mut Duration,
) -> Flow {
#[allow(deprecated)] #[allow(deprecated)]
let conn = match client.get_async_connection().await { let conn = tokio::select! {
Ok(c) => c, connected = client.get_async_connection() => match connected {
Err(e) => { Ok(c) => c,
tracing::error!("failed to connect to Redis for blacklist sync: {e}"); Err(e) => {
return; 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(); let mut pubsub = conn.into_pubsub();
if let Err(e) = pubsub.subscribe("rampart:blacklist:events").await { tokio::select! {
tracing::error!("failed to subscribe to blacklist events: {e}"); subscribed = pubsub.subscribe(BLACKLIST_CHANNEL) => match subscribed {
return; 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 { loop {
let mut stream = pubsub.on_message();
let msg_fut = stream.next();
tokio::pin!(msg_fut);
tokio::select! { tokio::select! {
_ = shutdown.changed() => { maybe_msg = stream.next() => match maybe_msg {
if *shutdown.borrow() { Some(msg) => {
tracing::info!("shutting down blacklist subscriber"); if let Err(e) = handle_event(&msg, blacklist) {
return; tracing::error!("blacklist event error: {e}");
} } else {
} *backoff = BACKOFF_MIN;
result = &mut msg_fut => {
match result {
Some(msg) => {
if let Err(e) = handle_event(&msg, &blacklist) {
tracing::error!("blacklist event error: {e}");
}
} }
None => { },
tracing::error!("pubsub stream ended"); None => {
tokio::time::sleep(Duration::from_secs(1)).await; tracing::error!("blacklist pubsub stream ended; resubscribing");
return; return Flow::Retry;
} },
} },
} () = wait_for_shutdown(shutdown) => return Flow::Stop,
} }
} }
} }
@ -126,3 +189,23 @@ impl StateStore for RedisStore {
Ok(()) 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);
}
}

4
src/traffic/intel/mod.rs Normal file
View file

@ -0,0 +1,4 @@
//! Интеллект трафика: EWMA-сглаживание, детектор атак и репутация IP.
pub mod detector;
pub mod ewma;
pub mod reputation;

View file

@ -1,9 +1,11 @@
//! Профилирование трафика, детектор атак и репутация IP. //! Профилирование трафика, детектор атак и репутация IP.
pub mod alert;
pub mod detector;
pub mod ewma;
pub mod hook; pub mod hook;
pub mod intel;
pub mod prefix; pub mod prefix;
pub mod profiler; pub mod profile;
pub mod reputation;
// Совместимость путей после группировки: домены `intel/` и `profile/`
// остаются доступны по прежним путям `traffic::{detector, ewma, ...}`.
pub use intel::{detector, ewma, reputation};
pub use profile::{alert, profiler};

91
src/traffic/prefix/key.rs Normal file
View file

@ -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<IpAddr> {
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
}

14
src/traffic/prefix/mod.rs Normal file
View file

@ -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};

View file

@ -1,27 +1,6 @@
//! Subnet-level (prefix) детектор распределённых атак.
//!
//! Источник данных — трейт [`PrefixStatsSource`]: карта XDP `prefix_stats`
//! (feature `xdp`) либо юзерспейс-агрегатор [`crate::engine::subnet_tracker`].
//! Контракт с XDP-агентом: [`PrefixKey`] / [`PrefixStatsVal`] должны совпадать
//! побайтово со структурами ядра.
use anyhow::Result; use anyhow::Result;
use std::future::Future; use std::future::Future;
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; use std::net::IpAddr;
/// `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],
}
/// Значение карты `prefix_stats`. /// Значение карты `prefix_stats`.
#[repr(C)] #[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<IpAddr> {
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` означает, что источник не считает уникальные адреса /// `unique_sources == 0` означает, что источник не считает уникальные адреса
@ -129,30 +56,6 @@ pub trait PrefixStatsSource {
fn snapshot(&self) -> impl Future<Output = Result<Vec<(IpAddr, PrefixSnapshot)>>> + Send; fn snapshot(&self) -> impl Future<Output = Result<Vec<(IpAddr, PrefixSnapshot)>>> + 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 детектора. /// Эскалация вердиктов subnet-level детектора.
#[derive(Debug, Clone, Copy, PartialEq, Eq)] #[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SubnetVerdict { pub enum SubnetVerdict {

View file

@ -0,0 +1,3 @@
//! Профилирование трафика edge-ноды и алерты детектора атак.
pub mod alert;
pub mod profiler;

3
src/tui/metrics/mod.rs Normal file
View file

@ -0,0 +1,3 @@
//! Метрики TUI: парсер Prometheus-exposition и fetcher эндпоинтов.
pub mod fetch;
pub mod prometheus;

View file

@ -4,7 +4,8 @@
//! exposition endpoint (`/metrics`). //! exposition endpoint (`/metrics`).
pub mod app; pub mod app;
pub mod fetch; pub mod metrics;
pub mod prometheus;
pub mod state; pub mod state;
pub mod ui; pub mod ui;
pub use metrics::{fetch, prometheus};

View file

@ -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<Self> {
let mut nums = release
.split(|c: char| !c.is_ascii_digit())
.filter_map(|part| part.parse::<u16>().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<String>;
#[must_use]
fn path_exists(&self, path: &str) -> bool;
#[must_use]
fn symlink_target(&self, path: &str) -> Option<String>;
#[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<String> {
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<String> {
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<u32> {
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<String> {
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<KernelVersion>,
pub btf_available: bool,
pub privileged: bool,
pub driver_name: Option<String>,
pub warnings: Vec<String>,
}
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()
}

58
src/xdp/filter/attach.rs Normal file
View file

@ -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> {
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<RingBuffer<'static>> {
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());
}
}

43
src/xdp/filter/maps.rs Normal file
View file

@ -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);
}
}

View file

@ -1,13 +1,19 @@
use anyhow::{Context, Result, bail}; 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::IpAddr;
use std::net::Ipv4Addr; use std::net::Ipv4Addr;
use std::os::unix::io::AsFd; use std::os::unix::io::AsFd;
use super::XdpStats; use super::XdpStats;
use super::globals::{RODATA_MAP_NAME, XdpGlobals}; use super::globals::XdpGlobals;
use crate::traffic::prefix::{PrefixKey, PrefixStatsVal}; 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 { pub struct XdpFilter {
obj: Option<Object>, obj: Option<Object>,
ringbuf: Option<RingBuffer<'static>>, ringbuf: Option<RingBuffer<'static>>,
@ -34,7 +40,7 @@ impl XdpFilter {
pub fn load(&mut self) -> Result<()> { pub fn load(&mut self) -> Result<()> {
// Диагностика ДО загрузки: fail-fast на неподдерживаемом ядре/драйвере. // Диагностика ДО загрузки: 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 bpf_obj = include_bytes!(concat!(env!("OUT_DIR"), "/universal_filter.o"));
let mut open_obj = ObjectBuilder::default() let mut open_obj = ObjectBuilder::default()
@ -43,9 +49,15 @@ impl XdpFilter {
patch_rodata(&mut open_obj, &self.globals)?; patch_rodata(&mut open_obj, &self.globals)?;
let obj = open_obj.load().context("Failed to load XDP object (verifier error?)")?; 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 { 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 let prog = obj
@ -64,9 +76,8 @@ impl XdpFilter {
} }
pub fn unload(&mut self) -> Result<()> { pub fn unload(&mut self) -> Result<()> {
if self.ifindex != 0 { if let Some(obj) = self.obj.as_ref() {
let fd = unsafe { std::os::unix::io::BorrowedFd::borrow_raw(std::os::unix::io::RawFd::from(-1)) }; self.detach_program(obj)?;
let _ = Xdp::new(fd).detach(self.ifindex, XdpFlags::NONE);
} }
self.ringbuf = None; self.ringbuf = None;
self.obj = None; self.obj = None;
@ -75,6 +86,28 @@ impl XdpFilter {
Ok(()) 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) { pub fn drain_events(&self) {
if let Some(rb) = &self.ringbuf { if let Some(rb) = &self.ringbuf {
let _ = rb.consume(); let _ = rb.consume();
@ -107,22 +140,15 @@ impl XdpFilter {
fn blacklist_update(&self, prefix_len: u8, v4: &[u8; 4], duration_secs: u64) -> Result<()> { fn blacklist_update(&self, prefix_len: u8, v4: &[u8; 4], duration_secs: u64) -> Result<()> {
let map = self.find_map("blacklist_map")?; let map = self.find_map("blacklist_map")?;
let mut key = [0u8; 8]; let key = blacklist_key(prefix_len, v4);
key[0] = prefix_len; let expiry_ns = blacklist_expiry_ns(unix_ns(), duration_secs);
key[4..8].copy_from_slice(v4); map.update(&key, &expiry_ns.to_le_bytes(), MapFlags::ANY)?;
map.update(
&key,
&(unix_ns() + duration_secs * 1_000_000_000).to_le_bytes(),
MapFlags::ANY,
)?;
Ok(()) Ok(())
} }
pub fn unban_ip(&self, ip: Ipv4Addr) -> Result<()> { pub fn unban_ip(&self, ip: Ipv4Addr) -> Result<()> {
let map = self.find_map("blacklist_map")?; let map = self.find_map("blacklist_map")?;
let mut key = [0u8; 8]; let key = blacklist_key(32, &ip.octets());
key[0] = 32;
key[4..8].copy_from_slice(&ip.octets());
map.delete(&key)?; map.delete(&key)?;
Ok(()) Ok(())
} }
@ -175,6 +201,11 @@ impl XdpFilter {
} }
} }
// SAFETY: `XdpFilter` is only ever accessed from the single thread that owns it,
// held behind `Arc<Mutex<XdpFilter>>`; 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 {} unsafe impl Send for XdpFilter {}
impl Drop for XdpFilter { impl Drop for XdpFilter {
@ -182,39 +213,3 @@ impl Drop for XdpFilter {
let _ = self.unload(); 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<RingBuffer<'static>> {
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()?)
}

View file

@ -1,10 +1,9 @@
//! Загрузчик XDP-программы (L3/L4-уровень защиты). //! Загрузчик XDP-программы (L3/L4-уровень защиты).
mod stats;
pub use stats::XdpStats;
mod diagnostics; mod probe;
pub use diagnostics::{ pub use probe::{
AttachMode, EnvironmentReport, FilesystemProbe, KernelVersion, MIN_KERNEL, SystemProbe, driver_supports_native_xdp, AttachMode, EnvironmentReport, FilesystemProbe, KernelVersion, MIN_KERNEL, SystemProbe, XdpStats,
driver_supports_native_xdp,
}; };
#[cfg(feature = "xdp")] #[cfg(feature = "xdp")]
@ -18,9 +17,7 @@ mod filter;
pub use filter::XdpFilter; pub use filter::XdpFilter;
#[cfg(feature = "xdp")] #[cfg(feature = "xdp")]
mod metrics; pub use probe::XdpMetrics;
#[cfg(feature = "xdp")]
pub use metrics::XdpMetrics;
#[cfg(not(feature = "xdp"))] #[cfg(not(feature = "xdp"))]
mod noop; mod noop;

View file

@ -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};

View file

@ -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<KernelVersion>,
pub btf_available: bool,
pub privileged: bool,
pub driver_name: Option<String>,
pub warnings: Vec<String>,
}
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()
}

View file

@ -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<Self> {
let mut nums = release
.split(|c: char| !c.is_ascii_digit())
.filter_map(|part| part.parse::<u16>().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<String>;
#[must_use]
fn path_exists(&self, path: &str) -> bool;
#[must_use]
fn symlink_target(&self, path: &str) -> Option<String>;
#[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<String> {
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<String> {
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<u32> {
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<String> {
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)
}

16
src/xdp/probe/mod.rs Normal file
View file

@ -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;