diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index e795e1f..037b116 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -25,6 +25,7 @@ jobs: - run: cargo fmt --all --check - run: cargo clippy --all-targets --all-features -- -D warnings - run: cargo check --all-features + - run: cargo check --no-default-features --all-targets - name: XDP — build binaries run: cargo build --features xdp --bins - name: XDP — smoke-check BPF compilation (clang) diff --git a/TODO.md b/TODO.md index 7ac7874..66e2eda 100644 --- a/TODO.md +++ b/TODO.md @@ -120,8 +120,11 @@ guard/ исправлены, d6bcae5). - [x] Rust loader (`src/xdp/filter/`): open/load/attach, patch глобалов G_* из config.toml (src/xdp/globals.rs), CString-safe if_nametoindex, реальный detach по fd прогрессы. -- [ ] Ringbuf events → blacklist: `drain_events()` сейчас только логирует debug; - нужен матчинг события → `ban_ip` (перенести из «loader», боится ошибок в data path). +- [x] Ringbuf events → blacklist: opt-in `auto_ban` ([xdp] config, default OFF) — + EVENT_RATE_LIMIT всегда, EVENT_CONN_DROP при fails>=порога; EVENT_BAN/POLICY_DROP не + банят (ядро уже отработало / статик-политика без аккумуляции). Парсер 24-байтного + `struct xdp_event` с фиксом endianness (IP — сетевой порядок, скаляры — host), + метрика `rampart_xdp_autobans_total`, tests/xdp_events.rs (10 тестов). - [ ] Smoke-test attach в CI (VM runner с CAP_BPF) + netns+veth харнесс — вывод №7 из конкурентного анализа, приоритет v0.4. @@ -146,10 +149,13 @@ guard/ 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`. +- [x] `cargo check --no-default-features` починен (2026-09-15, раунд 2): redis-код в + manager/sync/api/store/app закрыт `#[cfg(feature = "store-redis")]`, менеджер без фичи fail-fast'ит на старте; + CI теперь гоняет `cargo check --no-default-features --all-targets` гейтом. +- [x] `src/app/runtime.rs`: `#[allow(clippy::exit)]` убран — `app::run()` возвращает + `RunExit`, единственный `process::exit` теперь в `main` бина (clippy там не срабатывает). +- [x] Флакающий тест `pow_roundtrip::verifier_rejects_garbage_and_replay` (p≈1/16 из-за + случайного токена) — сидирован, 50 прогонов до (3 fail) / после (0 fail). - [ ] IPv6-банов в XDP-putи нет (`ban_cidr` bail'ит на v6) — IPv6-паритет (вывод №8). ## 3. Backlog diff --git a/src/app/mod.rs b/src/app/mod.rs index ba696c2..24b6750 100644 --- a/src/app/mod.rs +++ b/src/app/mod.rs @@ -2,5 +2,7 @@ pub mod runtime; pub mod services; +mod shutdown; pub use runtime::run; +pub use shutdown::RunExit; diff --git a/src/app/runtime.rs b/src/app/runtime.rs index 23a5b1a..1318c72 100644 --- a/src/app/runtime.rs +++ b/src/app/runtime.rs @@ -17,10 +17,12 @@ use std::sync::Arc; use std::sync::Mutex; use std::sync::atomic::{AtomicU64, Ordering}; use std::time::Duration; +use tokio::sync::oneshot; use tokio::sync::watch; use tracing_subscriber::EnvFilter; use super::services::{build_registry, build_whitelist, start_subnet_tracker, start_xdp}; +use super::shutdown::{RunExit, spawn_shutdown_drain}; fn attack_status_value(status: AttackStatus) -> i64 { match status { @@ -32,9 +34,11 @@ fn attack_status_value(status: AttackStatus) -> i64 { /// Точка входа edge-демона (`rampart`): то, что делал `main` бинаря. /// +/// Возвращает [`RunExit`] — выбор кода и самого завершения делает бинарь. +/// /// # Errors /// Ошибки старта: нечитаемый конфиг, пустой реестр протоколов, сбой listener. -pub async fn run() -> anyhow::Result<()> { +pub async fn run() -> anyhow::Result { tracing_subscriber::fmt() .with_env_filter(EnvFilter::from_default_env().add_directive("rampart=info".parse()?)) .init(); @@ -70,20 +74,9 @@ pub async fn run() -> anyhow::Result<()> { let allowed_1s = Arc::new(AtomicU64::new(0)); let (shutdown_tx, shutdown_rx) = watch::channel(false); + let (drain_tx, drain_rx) = oneshot::channel(); - 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); - }); + spawn_shutdown_drain(shutdown_tx.clone(), drain_tx); #[cfg(feature = "store-redis")] if let Some(redis_url) = &config.store.redis_url @@ -121,7 +114,7 @@ pub async fn run() -> anyhow::Result<()> { _ => None, }; - let xdp_filter = start_xdp(&config, &shutdown_rx)?; + let xdp_filter = start_xdp(&config, &shutdown_rx, &blacklist)?; let subnet_tracker = start_subnet_tracker(&config); monitor::spawn( &config.detect.prefix, @@ -197,7 +190,15 @@ pub async fn run() -> anyhow::Result<()> { registry, subnet_tracker, )); - listener::run(config, gateway, hook, shutdown_rx).await + tokio::select! { + result = listener::run(config, gateway, hook, shutdown_rx) => { + result?; + Ok(RunExit::Ok) + } + // Drain-таймаут сработал: возвращаем сигнал бинарю завершить процесс + // немедленно, как это делал прежний `process::exit(0)`. + _ = drain_rx => Ok(RunExit::DrainTimeout), + } } async fn push_attack_event(writer: &Option>>, pps: f64) { @@ -216,14 +217,3 @@ async fn push_attack_event(writer: &Option {} - _ = term.recv() => {} - } -} diff --git a/src/app/services.rs b/src/app/services.rs index c2b40bb..7553922 100644 --- a/src/app/services.rs +++ b/src/app/services.rs @@ -2,6 +2,7 @@ use crate::config::Config; use crate::engine::subnet::tracker::SubnetTracker; +use crate::filter::blacklist::Blacklist; use crate::protocol::ProtocolRegistry; use crate::xdp::XdpFilter; #[cfg(feature = "xdp")] @@ -10,7 +11,6 @@ use std::collections::HashSet; use std::net::IpAddr; use std::sync::Arc; use std::sync::Mutex; -use std::time::Duration; use tokio::sync::watch; /// Собирает реестр протоколов из скомпилированных реализаций. @@ -44,7 +44,10 @@ fn register_http(registry: &mut ProtocolRegistry, config: &Config) { pub(crate) fn start_xdp( config: &Arc, shutdown_rx: &watch::Receiver, + blacklist: &Arc, ) -> anyhow::Result>>> { + use std::time::Duration; + use crate::xdp::XdpMetrics; if !config.xdp.enabled { @@ -57,6 +60,10 @@ pub(crate) fn start_xdp( let mut guard = shared.lock().expect("xdp lock poisoned"); let globals = XdpGlobals::from_config(&config.xdp)?; guard.set_globals(globals); + // Порядок load-bearing: `load()` захватывает blacklist и политику + // авто-бана в ringbuf-колбэк (XdpFilter::load -> build_ringbuf), + // поэтому set_autoban обязан отработать ДО load(), а не после. + guard.set_autoban(blacklist.clone(), &config.xdp); guard.load()?; } let xdp_metrics = XdpMetrics::register()?; @@ -87,6 +94,7 @@ pub(crate) fn start_xdp( pub(crate) fn start_xdp( _config: &Arc, _shutdown_rx: &watch::Receiver, + _blacklist: &Arc, ) -> anyhow::Result>>> { Ok(None) } @@ -107,7 +115,23 @@ pub(crate) fn start_subnet_tracker(config: &Config) -> Option } } +/// Store-фичи проверяются на старте edge-демона: конфиг, запрашивающий redis, +/// в бинаре без `store-redis` должен упасть сразу, а не молча потерять +/// blacklist-sync (sync-задача в runtime.rs живёт под тем же cfg). +/// Точка вызова — `build_whitelist`: ближайший `?`-хук, принимающий конфиг +/// до бинда listener'а (сам runtime.rs вне правок этой задачи). +fn ensure_store_features(config: &Config) -> anyhow::Result<()> { + #[cfg(not(feature = "store-redis"))] + if config.store.redis_url.as_deref().is_some_and(|url| !url.is_empty()) { + anyhow::bail!("[store].redis_url requires a redis-backed store: built without store-redis feature"); + } + #[cfg(feature = "store-redis")] + let _ = config; + Ok(()) +} + pub(crate) fn build_whitelist(config: &Config) -> anyhow::Result>> { + ensure_store_features(config)?; let mut set = HashSet::with_capacity(config.whitelist.len()); for entry in &config.whitelist { let ip: IpAddr = entry @@ -117,3 +141,51 @@ pub(crate) fn build_whitelist(config: &Config) -> anyhow::Result Config { + Config { + store: crate::config::StoreConfig { + redis_url: Some("redis://127.0.0.1:6379/0".to_string()), + ..Default::default() + }, + ..Default::default() + } + } + + #[cfg(not(feature = "store-redis"))] + #[test] + fn redis_config_without_feature_fails_fast() { + let err = ensure_store_features(&config_with_redis()).expect_err("must fail without store-redis"); + assert!(err.to_string().contains("built without store-redis feature")); + } + + #[cfg(not(feature = "store-redis"))] + #[test] + fn empty_store_config_passes_without_feature() { + assert!(ensure_store_features(&Config::default()).is_ok()); + } + + #[cfg(feature = "store-redis")] + #[test] + fn guard_is_noop_with_feature() { + assert!(ensure_store_features(&config_with_redis()).is_ok()); + } + + /// Wiring-тасовка XDP auto-ban: `start_xdp` обязан принимать `Arc` + /// (т.е. существовать в сигнатуре с blacklist) и при выключенном `[xdp]` + /// не трогать ни ядро, ни blacklist. Порядок `set_autoban` -> `load()` + /// внутри enabled-ветки без ядра непроверяем — см. комментарий на месте вызова. + #[test] + fn start_xdp_accepts_blacklist_and_skips_kernel_when_disabled() { + let (_tx, rx) = watch::channel(false); + let blacklist = Arc::new(Blacklist::new()); + let started = start_xdp(&Arc::new(Config::default()), &rx, &blacklist) + .expect("xdp disabled: no kernel work and no error"); + assert!(started.is_none()); + assert!(blacklist.is_empty()); + } +} diff --git a/src/app/shutdown.rs b/src/app/shutdown.rs new file mode 100644 index 0000000..aaf94a8 --- /dev/null +++ b/src/app/shutdown.rs @@ -0,0 +1,67 @@ +//! Shutdown-drain машина edge-демона: сигнал, drain-таймер и код завершения. + +use std::time::Duration; +use tokio::sync::oneshot; +use tokio::sync::watch; + +/// Результат завершения демона. Код выхода вычисляется здесь, но сам +/// `process::exit` вызывается только в точке входа бинаря (`src/bin/rampart.rs`), +/// поэтому библиотечный код остаётся свободным от lint `exit`. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum RunExit { + /// Listener отработал штатно: drain уложился в таймаут. + Ok, + /// Истёк drain-таймаут: бинарь обязан завершить процесс немедленно + /// (замена прежнему `process::exit(0)` внутри shutdown-таймера). + DrainTimeout, +} + +impl RunExit { + /// Код завершения процесса для данного результата. + #[must_use] + pub const fn code(self) -> i32 { + match self { + Self::Ok => 0, + Self::DrainTimeout => 0, + } + } +} + +/// Ставит фоновую shutdown-задачу. +/// +/// shutdown-таймер обязан завершить процесс: сам `process::exit` вынесен в +/// бинарь, а здесь мы лишь уведомляем `run` об истечении drain-таймаута через +/// oneshot-канал. Логи и тайминг (5s после сигнала) сохранены без изменений. +pub(crate) fn spawn_shutdown_drain(shutdown_tx: watch::Sender, drain_tx: oneshot::Sender<()>) { + tokio::spawn(async move { + wait_for_signal().await; + tracing::info!("shutdown signal received, draining connections..."); + let _ = shutdown_tx.send(true); + tokio::time::sleep(Duration::from_secs(5)).await; + tracing::info!("shutdown timeout reached, exiting"); + let _ = drain_tx.send(()); + }); +} + +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() => {} + } +} + +#[cfg(test)] +mod tests { + use super::RunExit; + + #[test] + fn run_exit_codes_match_previous_process_exit() { + // Штатный путь и прежний drain-timeout (`process::exit(0)`) завершались с кодом 0. + assert_eq!(RunExit::Ok.code(), 0); + assert_eq!(RunExit::DrainTimeout.code(), 0); + } +} diff --git a/src/bin/rampart-manager.rs b/src/bin/rampart-manager.rs index b7768fa..5937bbc 100644 --- a/src/bin/rampart-manager.rs +++ b/src/bin/rampart-manager.rs @@ -5,7 +5,9 @@ use axum::{ routing::{get, post}, }; use dashmap::DashMap; -use rampart::manager::{AppState, api, auth, sync}; +#[cfg(feature = "store-redis")] +use rampart::manager::sync; +use rampart::manager::{AppState, api, auth}; use std::sync::Arc; use tower_http::cors::CorsLayer; use tracing_subscriber::EnvFilter; @@ -16,7 +18,15 @@ async fn main() -> anyhow::Result<()> { .with_env_filter(EnvFilter::from_default_env().add_directive("rampart=info".parse()?)) .init(); + // Manager живёт исключительно на redis-хранилище: без фичи стартовать + // нечего — fail-fast вместо полумёртвого API (features additive, TODO SOLID D). + if !cfg!(feature = "store-redis") { + anyhow::bail!("rampart-manager requires a redis state store: built without store-redis feature"); + } + + #[cfg(feature = "store-redis")] let redis_url = std::env::var("REDIS_URL").unwrap_or_else(|_| "redis://127.0.0.1:6379/0".to_string()); + #[cfg(feature = "store-redis")] let redis_client = redis::Client::open(redis_url)?; let jwt_secret = std::env::var("JWT_SECRET").map_err(|_| anyhow::anyhow!("JWT_SECRET must be set"))?; @@ -35,6 +45,7 @@ async fn main() -> anyhow::Result<()> { } let state = Arc::new(AppState { + #[cfg(feature = "store-redis")] redis_client, jwt_secret, jwt_audience, @@ -43,20 +54,25 @@ async fn main() -> anyhow::Result<()> { login_limiter: DashMap::new(), }); + #[cfg(feature = "store-redis")] tokio::spawn(sync::heartbeat::start_heartbeat_check(state.clone())); let public = Router::new() .route("/api/v1/health", get(api::inventory::health::health_check)) .route("/api/v1/auth/login", post(api::auth::login)); - let protected = Router::new() + // Redis-зависимые маршруты существуют только вместе с фичей; сам manager + // без неё bail'ит выше, поэтому 501-заглушки были бы недостижимым кодом. + let protected = Router::new(); + #[cfg(feature = "store-redis")] + let protected = protected .route("/api/v1/servers", get(api::inventory::servers::list_servers)) .route( "/api/v1/blacklist", get(api::blacklist::list_blacklist).post(api::blacklist::add_blacklist), ) - .route("/api/v1/nodes", get(api::inventory::nodes::list_nodes)) - .route_layer(middleware::from_fn(auth::auth_middleware)); + .route("/api/v1/nodes", get(api::inventory::nodes::list_nodes)); + let protected = protected.route_layer(middleware::from_fn(auth::auth_middleware)); let cors = match std::env::var("CORS_ORIGIN") { Ok(origin) if origin.is_empty() || origin == "*" => CorsLayer::new().allow_origin(tower_http::cors::Any), diff --git a/src/bin/rampart.rs b/src/bin/rampart.rs index d1c1011..0bbd751 100644 --- a/src/bin/rampart.rs +++ b/src/bin/rampart.rs @@ -1,6 +1,11 @@ -use rampart::app; +use rampart::app::{self, RunExit}; #[tokio::main] async fn main() -> anyhow::Result<()> { - app::run().await + match app::run().await? { + // Штатное завершение listener: управление возвращается tokio-рантайму. + RunExit::Ok => Ok(()), + // Drain-таймаут: бинарь — единственное место, где допустим process::exit. + exit => std::process::exit(exit.code()), + } } diff --git a/src/config/sections/platform.rs b/src/config/sections/platform.rs index eb66198..4169679 100644 --- a/src/config/sections/platform.rs +++ b/src/config/sections/platform.rs @@ -50,6 +50,16 @@ pub struct XdpConfig { pub throttle_enabled: bool, #[serde(default = "default_xdp_events_enabled")] pub events_enabled: bool, + /// OPT-IN: авто-бан по ringbuf-событиям XDP. По умолчанию выключен + /// («по умолчанию безопасно»): без него события только логируются. + #[serde(default)] + pub auto_ban_enabled: bool, + /// Длительность авто-бана в секундах (для RATE_LIMIT/CONN_DROP). + #[serde(default = "default_xdp_auto_ban_duration_secs")] + pub auto_ban_duration_secs: u64, + /// Порог `conn->fails` для авто-бана по EVENT_CONN_DROP (см. universal_filter.c:368). + #[serde(default = "default_xdp_auto_ban_conn_drop_fails")] + pub auto_ban_conn_drop_fails: u8, } impl Default for XdpConfig { @@ -67,6 +77,9 @@ impl Default for XdpConfig { challenge_timeout_ms: default_xdp_challenge_timeout_ms(), throttle_enabled: default_xdp_throttle_enabled(), events_enabled: default_xdp_events_enabled(), + auto_ban_enabled: false, + auto_ban_duration_secs: default_xdp_auto_ban_duration_secs(), + auto_ban_conn_drop_fails: default_xdp_auto_ban_conn_drop_fails(), } } } @@ -98,6 +111,12 @@ fn default_xdp_throttle_enabled() -> bool { fn default_xdp_events_enabled() -> bool { true } +fn default_xdp_auto_ban_duration_secs() -> u64 { + 60 +} +fn default_xdp_auto_ban_conn_drop_fails() -> u8 { + 3 +} #[derive(Debug, Clone, Deserialize)] pub struct LoggingConfig { diff --git a/src/engine/challenge/pow.rs b/src/engine/challenge/pow.rs index 0646a47..1e8350e 100644 --- a/src/engine/challenge/pow.rs +++ b/src/engine/challenge/pow.rs @@ -31,6 +31,18 @@ impl Challenge { } } + /// Challenge с явным токеном — для детерминированных тестов, где + /// случайный дайджест даёт flaky-исход. + #[must_use] + pub fn with_token(token: [u8; 32], difficulty: u8) -> Self { + Self { + token, + created_at: Instant::now(), + difficulty, + used: false, + } + } + #[must_use] pub fn is_expired(&self) -> bool { self.created_at.elapsed().as_secs() >= CHALLENGE_TTL_SECS diff --git a/src/manager/api/blacklist.rs b/src/manager/api/blacklist.rs index 2538df0..c22c708 100644 --- a/src/manager/api/blacklist.rs +++ b/src/manager/api/blacklist.rs @@ -1,6 +1,9 @@ +#[cfg(feature = "store-redis")] use crate::manager::AppState; +#[cfg(feature = "store-redis")] use axum::{Json, extract::State}; use serde::{Deserialize, Serialize}; +#[cfg(feature = "store-redis")] use std::sync::Arc; #[derive(Debug, Serialize, Deserialize)] @@ -29,6 +32,7 @@ pub struct AddBlacklistRequest { pub duration_secs: Option, } +#[cfg(feature = "store-redis")] pub async fn list_blacklist(State(state): State>) -> Json { let mut conn = match state.redis_client.get_multiplexed_async_connection().await { Ok(c) => c, @@ -69,6 +73,7 @@ pub async fn list_blacklist(State(state): State>) -> Json>, Json(req): Json, diff --git a/src/manager/api/inventory/nodes.rs b/src/manager/api/inventory/nodes.rs index 4f365df..bbc3ccb 100644 --- a/src/manager/api/inventory/nodes.rs +++ b/src/manager/api/inventory/nodes.rs @@ -1,7 +1,11 @@ +#[cfg(feature = "store-redis")] use crate::manager::AppState; +#[cfg(feature = "store-redis")] use axum::{Json, extract::State}; +#[cfg(feature = "store-redis")] use redis::AsyncCommands; use serde::{Deserialize, Serialize}; +#[cfg(feature = "store-redis")] use std::sync::Arc; #[derive(Debug, Serialize, Deserialize)] @@ -18,6 +22,7 @@ pub struct NodesResponse { pub nodes: Vec, } +#[cfg(feature = "store-redis")] pub async fn list_nodes(State(state): State>) -> Json { let mut conn = match state.redis_client.get_multiplexed_async_connection().await { Ok(c) => c, diff --git a/src/manager/api/inventory/servers.rs b/src/manager/api/inventory/servers.rs index bd8682c..6c2ccbc 100644 --- a/src/manager/api/inventory/servers.rs +++ b/src/manager/api/inventory/servers.rs @@ -1,7 +1,11 @@ +#[cfg(feature = "store-redis")] use crate::manager::AppState; +#[cfg(feature = "store-redis")] use axum::{Json, extract::State}; +#[cfg(feature = "store-redis")] use redis::AsyncCommands; use serde::{Deserialize, Serialize}; +#[cfg(feature = "store-redis")] use std::sync::Arc; #[derive(Debug, Serialize, Deserialize)] @@ -19,6 +23,7 @@ pub struct ServersResponse { pub servers: Vec, } +#[cfg(feature = "store-redis")] pub async fn list_servers(State(state): State>) -> Json { let mut conn = match state.redis_client.get_multiplexed_async_connection().await { Ok(c) => c, diff --git a/src/manager/mod.rs b/src/manager/mod.rs index e992ba6..2f89a26 100644 --- a/src/manager/mod.rs +++ b/src/manager/mod.rs @@ -10,6 +10,9 @@ use std::time::Instant; /// Общее состояние management-API. pub struct AppState { + /// Redis-клиент управления состоянием; присутствует только при собранной + /// фиче `store-redis` (manager требует её на старте — см. bin/rampart-manager). + #[cfg(feature = "store-redis")] pub redis_client: redis::Client, pub jwt_secret: String, pub jwt_audience: String, diff --git a/src/manager/sync/heartbeat.rs b/src/manager/sync/heartbeat.rs index 5ce0e54..773f689 100644 --- a/src/manager/sync/heartbeat.rs +++ b/src/manager/sync/heartbeat.rs @@ -1,17 +1,27 @@ +#[cfg(feature = "store-redis")] use crate::manager::AppState; +#[cfg(feature = "store-redis")] use futures::StreamExt; +#[cfg(feature = "store-redis")] use redis::AsyncCommands; +#[cfg(feature = "store-redis")] use redis::ScanOptions; +#[cfg(feature = "store-redis")] use std::sync::Arc; +#[cfg(feature = "store-redis")] use tokio::time::{Duration, interval}; /// SCAN pattern for node records (blocking KEYS is forbidden under load). +#[cfg(feature = "store-redis")] const NODES_KEY_PATTERN: &str = "rampart:nodes:*"; /// SCAN COUNT hint per round. +#[cfg(feature = "store-redis")] const SCAN_BATCH: usize = 100; /// Node is offline when the last heartbeat is older than this. +#[cfg(any(feature = "store-redis", test))] const NODE_TTL_SECS: i64 = 60; +#[cfg(feature = "store-redis")] pub async fn start_heartbeat_check(state: Arc) { let mut ticker = interval(Duration::from_secs(30)); loop { @@ -22,6 +32,7 @@ pub async fn start_heartbeat_check(state: Arc) { } } +#[cfg(feature = "store-redis")] async fn check_nodes(state: &AppState) -> anyhow::Result<()> { let mut conn = state.redis_client.get_multiplexed_async_connection().await?; let keys = collect_node_keys(&mut conn).await?; @@ -40,6 +51,7 @@ async fn check_nodes(state: &AppState) -> anyhow::Result<()> { } /// Incrementally iterates the keyspace via SCAN (cursor-based, non-blocking). +#[cfg(feature = "store-redis")] async fn collect_node_keys(conn: &mut redis::aio::MultiplexedConnection) -> anyhow::Result> { let opts = ScanOptions::default() .with_pattern(NODES_KEY_PATTERN) @@ -56,6 +68,7 @@ async fn collect_node_keys(conn: &mut redis::aio::MultiplexedConnection) -> anyh /// 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. +#[cfg(any(feature = "store-redis", test))] fn offline_update(raw: &str, now: i64) -> Option { let mut node = serde_json::from_str::(raw).ok()?; let hb = node["last_heartbeat"] @@ -115,6 +128,7 @@ mod tests { assert!(offline_update("not json", 2_000_000).is_none()); } + #[cfg(feature = "store-redis")] #[test] fn node_key_pattern_targets_rampart_nodes() { assert_eq!(NODES_KEY_PATTERN, "rampart:nodes:*"); diff --git a/src/xdp/filter/attach.rs b/src/xdp/filter/attach.rs index 3b076f8..2bf98cb 100644 --- a/src/xdp/filter/attach.rs +++ b/src/xdp/filter/attach.rs @@ -1,8 +1,15 @@ use anyhow::{Context, Result}; use libbpf_rs::{MapCore, Object, OpenObject, RingBuffer, RingBufferBuilder}; use std::ffi::CString; +use std::net::IpAddr; +use std::sync::Arc; +use std::time::Duration; +use crate::filter::blacklist::Blacklist; use crate::xdp::globals::{RODATA_MAP_NAME, XdpGlobals}; +use crate::xdp::metrics::XDP_AUTOBANS_TOTAL; + +use super::events::{AutoBanConfig, auto_ban_decision, parse_xdp_event}; /// Builds the NUL-terminated C string libc's `if_nametoindex` requires. /// Rejects interface names containing an interior NUL byte instead of @@ -22,18 +29,59 @@ pub(super) fn patch_rodata(open_obj: &mut OpenObject, globals: &XdpGlobals) -> R .with_context(|| format!("failed to set initial value of '{RODATA_MAP_NAME}'")) } -pub(super) fn build_ringbuf(obj: &Object) -> Result> { +/// Собирает ringbuf-колбэк для `events_map`. +/// +/// Колбэк обязан быть `'static`, поэтому всё разделяемое состояние (`Arc`, +/// копия `AutoBanConfig`) захватывается `move`-клонами. Когда `blacklist` — `None` +/// (фича не подключена) или политика выключена, события только парсятся и +/// логируются структурированно; баны не пишутся. В горячем пути нет `format!` — +/// `tracing`-поля и `&'static str` причины. +pub(super) fn build_ringbuf( + obj: &Object, + blacklist: Option>, + cfg: AutoBanConfig, +) -> Result> { let map = obj .maps() .find(|m| m.name() == "events_map") .context("events_map not found")?; let mut builder = RingBufferBuilder::new(); - builder.add(&map, |data: &[u8]| { - if data.len() >= 16 { - let ty = u32::from_ne_bytes(data[0..4].try_into().expect("4 bytes for type")); - let ip4 = u32::from_ne_bytes(data[4..8].try_into().expect("4 bytes for ip")); - let val = u64::from_ne_bytes(data[8..16].try_into().expect("8 bytes for val")); - tracing::debug!(event = ty, src_ip = ip4, data = val, "xdp event"); + builder.add(&map, move |data: &[u8]| { + let Some(event) = parse_xdp_event(data) else { + tracing::debug!( + len = data.len(), + "xdp ringbuf frame ignored (too short or unknown type)" + ); + return 0; + }; + tracing::debug!( + kind = ?event.kind, + src_ip = %event.src_ip, + metadata = event.metadata, + "xdp event" + ); + let Some(blacklist) = &blacklist else { + return 0; + }; + if let Some(ban) = auto_ban_decision(&event, &cfg) { + let addr = IpAddr::V4(ban.ip); + // `Blacklist::add` возвращает `()` (DashMap::insert игнорируется), а + // править blacklist.rs вне нашей зоны нельзя — поэтому «новый ли бан» + // определяем через публичный `is_blocked` ДО перезаписи. Это дедуп + // лога/метрики по активному бану, а не по факту вставки. + let already_banned = blacklist.is_blocked(addr); + blacklist.add(addr, Duration::from_secs(ban.duration_secs), ban.reason); + if already_banned { + tracing::debug!(src_ip = %ban.ip, reason = ban.reason, "xdp auto-ban refreshed"); + } else { + XDP_AUTOBANS_TOTAL.inc(); + tracing::info!( + src_ip = %ban.ip, + reason = ban.reason, + duration_secs = ban.duration_secs, + "xdp auto-ban issued" + ); + } } 0 })?; diff --git a/src/xdp/filter/events.rs b/src/xdp/filter/events.rs new file mode 100644 index 0000000..4c56788 --- /dev/null +++ b/src/xdp/filter/events.rs @@ -0,0 +1,166 @@ +//! Разбор и политика для событий XDP-ringbuf (`struct xdp_event`). +//! +//! Чистые (без I/O, без libbpf) функции, чтобы их можно было покрыть +//! юнит-тестами без ядра/Redis. Wiring колбэка — в [`super::attach`]. + +use std::net::Ipv4Addr; + +use crate::config::XdpConfig; + +/// Точный размер `struct xdp_event` (см. `xdp/core/common.h:61-66`): +/// `__u32 type` @0, `__u32 src_ip` @4, `__u32 metadata` @8, +/// 4 байта padding @12 (выравнивание `__u64`), `__u64 timestamp` @16 → 24. +pub const XDP_EVENT_LEN: usize = 24; + +/// Тип события из `enum event_type` (`common.h:54-59`). +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum XdpEventKind { + /// `EVENT_BAN` — пакет уже был в blacklist_map и сброшен ядром. + Ban, + /// `EVENT_RATE_LIMIT` — превышен порог троттлинга (SYN/UDP). + RateLimit, + /// `EVENT_POLICY_DROP` — статичная политика (UDP drop / вредоносные TCP-флаги). + PolicyDrop, + /// `EVENT_CONN_DROP` — flow сброшен (несоответствие seq или verdict хука). + ConnDrop, +} + +impl XdpEventKind { + /// Декодирует числовое значение из C-enum. Неизвестный тип → `None` + debug. + fn from_u32(raw: u32) -> Option { + match raw { + 0 => Some(Self::Ban), + 1 => Some(Self::RateLimit), + 2 => Some(Self::PolicyDrop), + 3 => Some(Self::ConnDrop), + other => { + tracing::debug!(event_type = other, "xdp event: unknown type, dropped"); + None + }, + } + } +} + +/// Разобранное событие ringbuf. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct XdpEvent { + pub kind: XdpEventKind, + /// IPv4 источника (для IPv4-mapped IPv6 ядро кладёт в `src_ip` младший u32). + pub src_ip: Ipv4Addr, + /// Поле `__u32 metadata`: для `RateLimit` — порог срабатывания троттлинга, + /// для `ConnDrop` — счётчик `conn->fails`, иначе `0`. + pub metadata: u32, + /// `bpf_ktime_get_ns()` на момент генерации (monotonic, не wall-clock). + pub timestamp_ns: u64, +} + +/// Разбирает 24-байтовый кадр `struct xdp_event`. +/// +/// # Byte order (проверено по C-источнику) +/// Скалярные поля (`type`, `metadata`, `timestamp`) пишутся ядром в ПОРЯДКЕ ХОСТА +/// без конверсии: `push_event` присваивает их как есть — `e->type = type`, +/// `e->metadata = meta`, `e->timestamp = bpf_ktime_get_ns()` +/// (`common.h:92-95`, `bpf_htonl`/`bpf_ntohl` к ним НЕ применяется). Userspace +/// читает те же байты на той же машине → `from_ne_bytes` корректно. +/// +/// А вот `src_ip` — сетевой порядок (big-endian octets): он приходит прямо из +/// заголовка пакета `src_ip = ip->saddr` (`universal_filter.c:149`) и пишется в +/// событие без `bpf_htonl` (`common.h:93`). Значит байты `data[4..8]` — это октеты +/// адреса в точковом порядке, и читать их надо как `Ipv4Addr::new(octets...)`, а НЕ +/// как `from_ne_bytes(u32)` (прежний баг: на LE-хосте давал перевернутый IP). +/// +/// Возвращает `None`, если кадр короче `XDP_EVENT_LEN` или тип неизвестен. +pub fn parse_xdp_event(bytes: &[u8]) -> Option { + let raw: &[u8; XDP_EVENT_LEN] = bytes.get(..XDP_EVENT_LEN)?.try_into().ok()?; + let kind = XdpEventKind::from_u32(u32::from_ne_bytes(raw[0..4].try_into().ok()?))?; + let src_ip = Ipv4Addr::new(raw[4], raw[5], raw[6], raw[7]); + let metadata = u32::from_ne_bytes(raw[8..12].try_into().ok()?); + // raw[12..16] — выравнивающий padding под __u64, не читаем (он не инициализирован). + let timestamp_ns = u64::from_ne_bytes(raw[16..24].try_into().ok()?); + Some(XdpEvent { + kind, + src_ip, + metadata, + timestamp_ns, + }) +} + +/// Параметры авто-бана, резолвинг из секции `[xdp]`. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct AutoBanConfig { + pub enabled: bool, + pub duration_secs: u64, + pub conn_drop_fails: u8, +} + +impl Default for AutoBanConfig { + fn default() -> Self { + Self { + enabled: false, + duration_secs: 60, + conn_drop_fails: 3, + } + } +} + +impl AutoBanConfig { + /// Берёт значения из `[xdp]`-секции конфига (единственный источник истины). + pub fn from_xdp(cfg: &XdpConfig) -> Self { + Self { + enabled: cfg.auto_ban_enabled, + duration_secs: cfg.auto_ban_duration_secs, + conn_drop_fails: cfg.auto_ban_conn_drop_fails, + } + } +} + +/// Готовый бан для применения к userspace-`Blacklist`. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct BanAction { + pub ip: Ipv4Addr, + pub duration_secs: u64, + /// Статичная причина (без аллокаций в горячем пути). + pub reason: &'static str, +} + +/// Чистое решение: банить ли источник события и на сколько. +/// +/// Возвращает `None`, если фича выключена или тип события не банится. +/// Разбор по типам обоснован C-кодом: +/// +/// * [`XdpEventKind::Ban`] — НЕ баним. Событие эмитится уже ПОСЛЕ того, как +/// ядро нашло IP в `blacklist_map` и дропнуло пакет (`universal_filter.c:220-225`); +/// повторный userspace-бан бессмысленен и мог бы скрыть штатный expiry. +/// * [`XdpEventKind::PolicyDrop`] — НЕ баним по умолчанию. Это статичная +/// политика на каждый пакет: UDP-drop-политика (`:234`) и вредоносные TCP-флаги +/// (`:259`). Бан по одному malformed-пакету перебанит легитимных, но глючных +/// клиентов — порог аккумулирования здесь отсутствует по смыслу. +/// * [`XdpEventKind::RateLimit`] — банیم. Порог троттлинга уже фактически превышен; +/// при этом UDP-путь (`:237-243`) только дропает и НЕ заносит в blacklist (в +/// отличие от SYN-пути `:295`) — userspace-бан добавляет реальную защиту. +/// * [`XdpEventKind::ConnDrop`] — баним, только если `metadata` (=`conn->fails`) +/// достиг порога (`:364-368`). Дроп по протокольному хуку приходит с `fails = 0` +/// (`:418`) и ниже порога, т.е. не банится. +#[must_use] +pub fn auto_ban_decision(event: &XdpEvent, config: &AutoBanConfig) -> Option { + if !config.enabled { + return None; + } + match event.kind { + XdpEventKind::Ban | XdpEventKind::PolicyDrop => None, + XdpEventKind::RateLimit => Some(BanAction { + ip: event.src_ip, + duration_secs: config.duration_secs, + reason: "xdp-rate-limit", + }), + XdpEventKind::ConnDrop if event.metadata >= u32::from(config.conn_drop_fails) => Some(BanAction { + ip: event.src_ip, + duration_secs: config.duration_secs, + reason: "xdp-conn-drop", + }), + XdpEventKind::ConnDrop => None, + } +} + +// Покрытие парсера и политики — в `tests/xdp_events.rs` (держим модуль под +// лимитом 250 строк; `src/xdp/filter/` уже на пределе в 4 файла на каталог). diff --git a/src/xdp/filter/mod.rs b/src/xdp/filter/mod.rs index 21de55f..0d8b4af 100644 --- a/src/xdp/filter/mod.rs +++ b/src/xdp/filter/mod.rs @@ -3,15 +3,20 @@ use libbpf_rs::{MapCore, MapFlags, Object, ObjectBuilder, RingBuffer, Xdp, XdpFl use std::net::IpAddr; use std::net::Ipv4Addr; use std::os::unix::io::AsFd; +use std::sync::Arc; use super::XdpStats; use super::globals::XdpGlobals; +use crate::config::XdpConfig; +use crate::filter::blacklist::Blacklist; use crate::traffic::prefix::{PrefixKey, PrefixStatsVal}; use attach::{build_ringbuf, interface_cstr, patch_rodata}; +use events::AutoBanConfig; use maps::{blacklist_expiry_ns, blacklist_key, unix_ns}; mod attach; +pub(crate) mod events; mod maps; pub struct XdpFilter { @@ -20,6 +25,8 @@ pub struct XdpFilter { ifindex: i32, interface: String, globals: XdpGlobals, + blacklist: Option>, + autoban: AutoBanConfig, } impl XdpFilter { @@ -30,6 +37,8 @@ impl XdpFilter { ifindex: 0, interface: interface.to_string(), globals: XdpGlobals::default(), + blacklist: None, + autoban: AutoBanConfig::default(), } } @@ -38,6 +47,16 @@ impl XdpFilter { self.globals = globals; } + /// Подключает OPT-IN авто-бан по ringbuf-событиям к userspace-`Blacklist`. + /// + /// До вызова `load()`. Политика читается из `[xdp]`-секции; если + /// `auto_ban_enabled = false` (по умолчанию), колбэк лишь логирует события и + /// не трогает blacklist. + pub fn set_autoban(&mut self, blacklist: Arc, config: &XdpConfig) { + self.blacklist = Some(blacklist); + self.autoban = AutoBanConfig::from_xdp(config); + } + pub fn load(&mut self) -> Result<()> { // Диагностика ДО загрузки: fail-fast на неподдерживаемом ядре/драйвере. super::probe::preflight(&self.interface)?; @@ -66,7 +85,8 @@ impl XdpFilter { .context("XDP program 'rampart_universal_filter' not found")?; Xdp::new(prog.as_fd()).attach(ifindex as i32, XdpFlags::NONE)?; - let rbuf = build_ringbuf(&obj)?; + super::metrics::init(); + let rbuf = build_ringbuf(&obj, self.blacklist.clone(), self.autoban)?; self.obj = Some(obj); self.ringbuf = Some(rbuf); diff --git a/src/xdp/metrics.rs b/src/xdp/metrics.rs new file mode 100644 index 0000000..fab4ce9 --- /dev/null +++ b/src/xdp/metrics.rs @@ -0,0 +1,20 @@ +//! Prometheus-метрики загрузчика XDP. +//! +//! Регистрируются в глобальный реестр `prometheus` (тот же, что собирает +//! `prometheus::gather()` в `crate::metrics::run_metrics_server`), поэтому +//! endpoint `/metrics` подхватывает их без правок вне `src/xdp/`. + +use prometheus::{IntCounter, register_int_counter}; +use std::sync::LazyLock; + +/// Сколько раз событие XDP-ringbuf привело к userspace-авто-бану +/// (счётчик инкрементится только на НОВЫЙ бан, не на продление). +pub static XDP_AUTOBANS_TOTAL: LazyLock = LazyLock::new(|| { + register_int_counter!("rampart_xdp_autobans_total", "Auto-bans issued from XDP ringbuf events") + .expect("XDP_AUTOBANS_TOTAL") +}); + +/// Форсирует регистрацию счётчика (для вызова на старте загрузки XDP). +pub fn init() { + LazyLock::force(&XDP_AUTOBANS_TOTAL); +} diff --git a/src/xdp/mod.rs b/src/xdp/mod.rs index 30ca8a8..15147c0 100644 --- a/src/xdp/mod.rs +++ b/src/xdp/mod.rs @@ -11,10 +11,17 @@ mod globals; #[cfg(feature = "xdp")] pub use globals::XdpGlobals; +#[cfg(feature = "xdp")] +mod metrics; + #[cfg(feature = "xdp")] mod filter; #[cfg(feature = "xdp")] pub use filter::XdpFilter; +#[cfg(feature = "xdp")] +pub use filter::events::{ + AutoBanConfig, BanAction, XDP_EVENT_LEN, XdpEvent, XdpEventKind, auto_ban_decision, parse_xdp_event, +}; #[cfg(feature = "xdp")] pub use probe::XdpMetrics; diff --git a/src/xdp/probe/diagnostics/report.rs b/src/xdp/probe/diagnostics/report.rs index 24f6f03..1e8a951 100644 --- a/src/xdp/probe/diagnostics/report.rs +++ b/src/xdp/probe/diagnostics/report.rs @@ -1,7 +1,7 @@ use anyhow::{Result, bail}; use super::types::{ - AttachMode, BTF_PATH, FilesystemProbe, KernelVersion, MIN_KERNEL, OSRELEASE_PATH, SystemProbe, driver_from_link, + AttachMode, BTF_PATH, KernelVersion, MIN_KERNEL, OSRELEASE_PATH, SystemProbe, driver_from_link, driver_supports_native_xdp, }; @@ -109,6 +109,8 @@ impl EnvironmentReport { /// См. [`EnvironmentReport::validate`]. #[cfg(feature = "xdp")] pub(crate) fn preflight(interface: &str) -> Result<()> { + use super::types::FilesystemProbe; + let report = EnvironmentReport::collect(&FilesystemProbe, interface); let kernel = report .kernel_version diff --git a/tests/pow_roundtrip.rs b/tests/pow_roundtrip.rs index 6d853d9..ad703a2 100644 --- a/tests/pow_roundtrip.rs +++ b/tests/pow_roundtrip.rs @@ -1,5 +1,11 @@ use rampart::engine::challenge::{Challenge, solve}; +/// Fixed 32-byte token: sha256(hex(token) || "garbage") starts with "ebc1", +/// so "garbage" provably fails the difficulty-2 check (first two hex chars +/// must be in "0123"). A random token made this flaky with probability 1/16 +/// per run, since any digest lands in the accepted prefix with p=(4/16)^2. +const FIXED_TOKEN: [u8; 32] = *b"rampart-pow-fixed-token-test-v1!"; + #[test] fn solver_output_passes_verifier() { let mut challenge = Challenge::generate(3); @@ -9,7 +15,7 @@ fn solver_output_passes_verifier() { #[test] fn verifier_rejects_garbage_and_replay() { - let mut challenge = Challenge::generate(2); + let mut challenge = Challenge::with_token(FIXED_TOKEN, 2); assert!(!challenge.verify("garbage")); let nonce = solve(&challenge.challenge_string(), 2).expect("solved"); diff --git a/tests/xdp_events.rs b/tests/xdp_events.rs new file mode 100644 index 0000000..88611a3 --- /dev/null +++ b/tests/xdp_events.rs @@ -0,0 +1,173 @@ +//! Интеграционные тесты wire-формата и политики авто-бана XDP-событий. +//! Ядро/Redis не нужны: проверяем чистые парсер и решение + дефолты конфига. +#![cfg(feature = "xdp")] + +use rampart::config::Config; +use rampart::xdp::{AutoBanConfig, XDP_EVENT_LEN, XdpEventKind, auto_ban_decision, parse_xdp_event}; +use std::net::Ipv4Addr; + +const IP: [u8; 4] = [203, 0, 113, 7]; + +/// Собирает кадр ровно как ядро (`push_event`, common.h:87-97): скаляры — native +/// (host) order, `src_ip` — сетевой порядок октетов (raw из заголовка пакета). +fn frame(kind: u32, octets: [u8; 4], metadata: u32, ts: u64) -> [u8; XDP_EVENT_LEN] { + let mut b = [0u8; XDP_EVENT_LEN]; + b[0..4].copy_from_slice(&kind.to_ne_bytes()); + b[4..8].copy_from_slice(&octets); + b[8..12].copy_from_slice(&metadata.to_ne_bytes()); + // b[12..16] — выравнивающий padding под __u64: ядро его не инициализирует. + b[16..24].copy_from_slice(&ts.to_ne_bytes()); + b +} + +#[test] +fn struct_size_is_24_bytes() { + // Раскладка common.h: __u32×3 + pad(4) + __u64 = 24, НЕ 16 как парсил код ранее. + assert_eq!(XDP_EVENT_LEN, 24); +} + +#[test] +fn round_trip_recovers_all_fields() { + let ev = parse_xdp_event(&frame(1, IP, 100, 0xDEAD_BEEF)).expect("valid frame parses"); + assert_eq!(ev.kind, XdpEventKind::RateLimit); + assert_eq!(ev.src_ip, Ipv4Addr::new(203, 0, 113, 7)); + assert_eq!(ev.metadata, 100); + assert_eq!(ev.timestamp_ns, 0xDEAD_BEEF); +} + +#[test] +fn dirty_padding_does_not_corrupt_parse() { + let mut f = frame(3, IP, 5, 0x1_0000_0000); + f[12..16].copy_from_slice(&[0xAA, 0xBB, 0xCC, 0xDD]); + let ev = parse_xdp_event(&f).expect("padding must be ignored"); + assert_eq!(ev.kind, XdpEventKind::ConnDrop); + assert_eq!(ev.metadata, 5); + assert_eq!(ev.timestamp_ns, 0x1_0000_0000); +} + +#[test] +fn src_ip_is_network_order_octets_not_native_u32() { + // КЛЮЧЕВОЙ вывод по endianness: октеты лежат в точковом порядке напрямую. + let f = frame(2, [198, 51, 100, 23], 0, 0); + let ev = parse_xdp_event(&f).expect("parses"); + assert_eq!(ev.src_ip, Ipv4Addr::new(198, 51, 100, 23)); + // Старый баг (`u32::from_ne_bytes` → Ipv4Addr::from(u32)) на little-endian + // давал бы переёрнутый адрес — фиксируем, что мы так НЕ делаем. + if cfg!(target_endian = "little") { + let naive = u32::from_ne_bytes(<[u8; 4]>::try_from(&f[4..8]).expect("4 octets")); + assert_eq!(Ipv4Addr::from(naive), Ipv4Addr::new(23, 100, 51, 198)); + assert_ne!(Ipv4Addr::from(naive), ev.src_ip); + } +} + +#[test] +fn truncated_frames_are_rejected() { + let full = frame(1, IP, 0, 0); + assert!(parse_xdp_event(&full[..16]).is_none(), "старая длина 16 — не валидна"); + assert!(parse_xdp_event(&full[..23]).is_none(), "на 1 байт короче — не валидна"); + assert!(parse_xdp_event(&full).is_some(), "ровно 24 — валидна"); + assert!(parse_xdp_event(&[]).is_none()); +} + +#[test] +fn unknown_event_type_is_rejected() { + for kind in [4u32, 9, 255, u32::MAX] { + assert!(parse_xdp_event(&frame(kind, IP, 0, 0)).is_none(), "type {kind} unknown"); + } +} + +#[test] +fn decision_matrix_by_type_and_enablement() { + let enabled = AutoBanConfig { + enabled: true, + ..AutoBanConfig::default() + }; + let disabled = AutoBanConfig { + enabled: false, + ..AutoBanConfig::default() + }; + // (kind, metadata=fails/хиты, банить ли при enabled) + let rows = [ + (0u32, 255u32, false), // EVENT_BAN — ядро уже забанило: никогда + (1, 100, true), // EVENT_RATE_LIMIT — всегда баним + (2, 0, false), // EVENT_POLICY_DROP — политика по дефолту не банит + (3, 10, true), // EVENT_CONN_DROP fails=10 >= 3 — баним + (3, 0, false), // EVENT_CONN_DROP fails=0 (hook-drop) — не баним + (3, 2, false), // EVENT_CONN_DROP fails=2 < 3 — не баним + ]; + for (kind, meta, should_ban_when_on) in rows { + let ev = parse_xdp_event(&frame(kind, IP, meta, 0)).expect("known type"); + assert!( + auto_ban_decision(&ev, &disabled).is_none(), + "disabled: kind {kind} meta {meta} must never ban" + ); + let decided = auto_ban_decision(&ev, &enabled).is_some(); + assert_eq!( + decided, should_ban_when_on, + "enabled: kind {kind} meta {meta} -> ban={decided}, expected {should_ban_when_on}" + ); + } +} + +#[test] +fn ban_carries_configured_duration_and_reason() { + let cfg = AutoBanConfig { + enabled: true, + duration_secs: 120, + conn_drop_fails: 3, + }; + let rl = parse_xdp_event(&frame(1, IP, 0, 0)).expect("rate limit frame"); + let ban = auto_ban_decision(&rl, &cfg).expect("rate limit bans"); + assert_eq!(ban.ip, Ipv4Addr::new(203, 0, 113, 7)); + assert_eq!(ban.duration_secs, 120); + assert_eq!(ban.reason, "xdp-rate-limit"); + + let cd = parse_xdp_event(&frame(3, IP, 4, 0)).expect("conn drop frame"); + assert_eq!( + auto_ban_decision(&cd, &cfg).expect("conn drop bans").reason, + "xdp-conn-drop" + ); +} + +#[test] +fn defaults_disable_autoban_without_xdp_keys() { + let cfg = Config::parse_str("").expect("empty config parses"); + assert!( + !cfg.xdp.auto_ban_enabled, + "авто-бан обязан быть выключен по умолчанию (safe-by-default)" + ); + assert_eq!(cfg.xdp.auto_ban_duration_secs, 60); + assert_eq!(cfg.xdp.auto_ban_conn_drop_fails, 3); + + let policy = AutoBanConfig::from_xdp(&cfg.xdp); + let rl = parse_xdp_event(&frame(1, IP, 0, 0)).expect("frame parses"); + assert!(auto_ban_decision(&rl, &policy).is_none(), "default policy bans nothing"); +} + +#[test] +fn xdp_section_enables_autoban_and_maps_fields() { + let cfg = Config::parse_str( + r#" +[xdp] +auto_ban_enabled = true +auto_ban_duration_secs = 45 +auto_ban_conn_drop_fails = 5 +"#, + ) + .expect("[xdp] auto-ban keys parse"); + let policy = AutoBanConfig::from_xdp(&cfg.xdp); + assert!(policy.enabled); + assert_eq!(policy.duration_secs, 45); + assert_eq!(policy.conn_drop_fails, 5); + + // fails=4 ниже нового порога 5 — не бан; fails=5 — бан. + let below = parse_xdp_event(&frame(3, IP, 4, 0)).expect("frame parses"); + assert!(auto_ban_decision(&below, &policy).is_none()); + let at = parse_xdp_event(&frame(3, IP, 5, 0)).expect("frame parses"); + assert_eq!( + auto_ban_decision(&at, &policy) + .expect("at-threshold bans") + .duration_secs, + 45 + ); +}