feat(xdp): auto-ban from ringbuf events; fix feature gating, exit path, pow flake
Some checks are pending
CI / Rust — check & clippy (push) Waiting to run
CI / Rust — test (push) Waiting to run
CI / Repo — module size & default secrets (push) Waiting to run
CI / Rust — cargo-deny (push) Waiting to run
CI / Docker — build edge image (push) Blocked by required conditions
Some checks are pending
CI / Rust — check & clippy (push) Waiting to run
CI / Rust — test (push) Waiting to run
CI / Repo — module size & default secrets (push) Waiting to run
CI / Rust — cargo-deny (push) Waiting to run
CI / Docker — build edge image (push) Blocked by required conditions
- xdp: real 24-byte xdp_event parsing (fixes LE byte-swap of src IP), opt-in [xdp] auto_ban (RATE_LIMIT always, CONN_DROP at fails>=threshold; EVENT_BAN/POLICY_DROP excluded by design), rampart_xdp_autobans_total; wired in app before load(); 10 tests in tests/xdp_events.rs - build: --no-default-features compiles — redis paths cfg-gated behind store-redis, manager fail-fasts without it; CI gates the config now - app: RunExit enum replaces process::exit in lib; single exit site in main - tests: seed PoW roundtrip token (was ~1/16 flaky; 50 pre-fix fails -> 0) - docs: TODO statuses refreshed (round 2)
This commit is contained in:
parent
aa615a1141
commit
bd40a44980
23 changed files with 714 additions and 50 deletions
1
.github/workflows/ci.yml
vendored
1
.github/workflows/ci.yml
vendored
|
|
@ -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)
|
||||
|
|
|
|||
18
TODO.md
18
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
|
||||
|
|
|
|||
|
|
@ -2,5 +2,7 @@
|
|||
|
||||
pub mod runtime;
|
||||
pub mod services;
|
||||
mod shutdown;
|
||||
|
||||
pub use runtime::run;
|
||||
pub use shutdown::RunExit;
|
||||
|
|
|
|||
|
|
@ -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<RunExit> {
|
||||
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<Arc<tokio::sync::Mutex<ClickHouseWriter>>>, pps: f64) {
|
||||
|
|
@ -216,14 +217,3 @@ async fn push_attack_event(writer: &Option<Arc<tokio::sync::Mutex<ClickHouseWrit
|
|||
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() => {}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<Config>,
|
||||
shutdown_rx: &watch::Receiver<bool>,
|
||||
blacklist: &Arc<Blacklist>,
|
||||
) -> anyhow::Result<Option<Arc<Mutex<XdpFilter>>>> {
|
||||
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<Config>,
|
||||
_shutdown_rx: &watch::Receiver<bool>,
|
||||
_blacklist: &Arc<Blacklist>,
|
||||
) -> anyhow::Result<Option<Arc<Mutex<XdpFilter>>>> {
|
||||
Ok(None)
|
||||
}
|
||||
|
|
@ -107,7 +115,23 @@ pub(crate) fn start_subnet_tracker(config: &Config) -> Option<Arc<SubnetTracker>
|
|||
}
|
||||
}
|
||||
|
||||
/// 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<Arc<HashSet<IpAddr>>> {
|
||||
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<Arc<HashSet<IpA
|
|||
}
|
||||
Ok(Arc::new(set))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn config_with_redis() -> 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>`
|
||||
/// (т.е. существовать в сигнатуре с 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());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
67
src/app/shutdown.rs
Normal file
67
src/app/shutdown.rs
Normal file
|
|
@ -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<bool>, 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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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),
|
||||
|
|
|
|||
|
|
@ -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()),
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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<u64>,
|
||||
}
|
||||
|
||||
#[cfg(feature = "store-redis")]
|
||||
pub async fn list_blacklist(State(state): State<Arc<AppState>>) -> Json<BlacklistResponse> {
|
||||
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<Arc<AppState>>) -> Json<Blacklis
|
|||
Json(BlacklistResponse { items, total })
|
||||
}
|
||||
|
||||
#[cfg(feature = "store-redis")]
|
||||
pub async fn add_blacklist(
|
||||
State(state): State<Arc<AppState>>,
|
||||
Json(req): Json<AddBlacklistRequest>,
|
||||
|
|
|
|||
|
|
@ -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<NodeInfo>,
|
||||
}
|
||||
|
||||
#[cfg(feature = "store-redis")]
|
||||
pub async fn list_nodes(State(state): State<Arc<AppState>>) -> Json<NodesResponse> {
|
||||
let mut conn = match state.redis_client.get_multiplexed_async_connection().await {
|
||||
Ok(c) => c,
|
||||
|
|
|
|||
|
|
@ -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<ServerEntry>,
|
||||
}
|
||||
|
||||
#[cfg(feature = "store-redis")]
|
||||
pub async fn list_servers(State(state): State<Arc<AppState>>) -> Json<ServersResponse> {
|
||||
let mut conn = match state.redis_client.get_multiplexed_async_connection().await {
|
||||
Ok(c) => c,
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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<AppState>) {
|
||||
let mut ticker = interval(Duration::from_secs(30));
|
||||
loop {
|
||||
|
|
@ -22,6 +32,7 @@ pub async fn start_heartbeat_check(state: Arc<AppState>) {
|
|||
}
|
||||
}
|
||||
|
||||
#[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<Vec<String>> {
|
||||
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<String> {
|
||||
let mut node = serde_json::from_str::<serde_json::Value>(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:*");
|
||||
|
|
|
|||
|
|
@ -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<RingBuffer<'static>> {
|
||||
/// Собирает ringbuf-колбэк для `events_map`.
|
||||
///
|
||||
/// Колбэк обязан быть `'static`, поэтому всё разделяемое состояние (`Arc<Blacklist>`,
|
||||
/// копия `AutoBanConfig`) захватывается `move`-клонами. Когда `blacklist` — `None`
|
||||
/// (фича не подключена) или политика выключена, события только парсятся и
|
||||
/// логируются структурированно; баны не пишутся. В горячем пути нет `format!` —
|
||||
/// `tracing`-поля и `&'static str` причины.
|
||||
pub(super) fn build_ringbuf(
|
||||
obj: &Object,
|
||||
blacklist: Option<Arc<Blacklist>>,
|
||||
cfg: AutoBanConfig,
|
||||
) -> 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");
|
||||
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
|
||||
})?;
|
||||
|
|
|
|||
166
src/xdp/filter/events.rs
Normal file
166
src/xdp/filter/events.rs
Normal file
|
|
@ -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<Self> {
|
||||
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<XdpEvent> {
|
||||
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<BanAction> {
|
||||
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 файла на каталог).
|
||||
|
|
@ -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<Arc<Blacklist>>,
|
||||
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<Blacklist>, 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);
|
||||
|
|
|
|||
20
src/xdp/metrics.rs
Normal file
20
src/xdp/metrics.rs
Normal file
|
|
@ -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<IntCounter> = 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);
|
||||
}
|
||||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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");
|
||||
|
|
|
|||
173
tests/xdp_events.rs
Normal file
173
tests/xdp_events.rs
Normal file
|
|
@ -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
|
||||
);
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue