diff --git a/src/config/mod.rs b/src/config/mod.rs index 7c311bc..bdf3391 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -2,8 +2,8 @@ mod sections; pub use sections::{ - BackendConfig, BanConfig, BindConfig, LimitsConfig, LoggingConfig, MetricsConfig, PowConfig, StoreConfig, - WorkerConfig, XdpConfig, + BackendConfig, BanConfig, BindConfig, DetectConfig, DetectPrefixConfig, LimitsConfig, LoggingConfig, MetricsConfig, + PowConfig, StoreConfig, WorkerConfig, XdpConfig, }; use serde::Deserialize; @@ -32,6 +32,8 @@ pub struct Config { #[serde(default)] pub pow: PowConfig, #[serde(default)] + pub detect: DetectConfig, + #[serde(default)] pub whitelist: Vec, } @@ -70,6 +72,9 @@ impl Config { anyhow::bail!("invalid whitelist entry: {entry}"); } } + if self.detect.prefix.enabled && self.detect.prefix.syn_threshold == 0 { + anyhow::bail!("detect.prefix.syn_threshold must be positive when enabled"); + } Ok(()) } } @@ -140,4 +145,43 @@ upstreams = ["not-an-addr"] let result = Config::parse_str("whitelist = [\"999.999.1.1\"]"); assert!(result.is_err()); } + + #[test] + fn prefix_detect_defaults_disabled() { + let config = Config::parse_str("").expect("empty config should parse"); + assert!(!config.detect.prefix.enabled); + assert_eq!(config.detect.prefix.syn_threshold, 500); + assert_eq!(config.detect.prefix.window_secs, 10); + assert_eq!(config.detect.prefix.min_unique_sources, 16); + } + + #[test] + fn prefix_detect_section_parses() { + let config = Config::parse_str( + r#" +[detect.prefix] +enabled = true +syn_threshold = 1000 +window_secs = 5 +min_unique_sources = 8 +"#, + ) + .expect("detect.prefix section should parse"); + assert!(config.detect.prefix.enabled); + assert_eq!(config.detect.prefix.syn_threshold, 1000); + assert_eq!(config.detect.prefix.window_secs, 5); + assert_eq!(config.detect.prefix.min_unique_sources, 8); + } + + #[test] + fn prefix_detect_zero_threshold_rejected_when_enabled() { + let result = Config::parse_str( + r#" +[detect.prefix] +enabled = true +syn_threshold = 0 +"#, + ); + assert!(result.is_err()); + } } diff --git a/src/config/sections.rs b/src/config/sections.rs index 029632e..97e8535 100644 --- a/src/config/sections.rs +++ b/src/config/sections.rs @@ -226,3 +226,44 @@ impl Default for PowConfig { fn default_pow_difficulty() -> u8 { 4 } + +/// Секция `[detect]`: детекторы распределённых атак. +#[derive(Debug, Clone, Default, Deserialize)] +pub struct DetectConfig { + #[serde(default)] + pub prefix: DetectPrefixConfig, +} + +/// Секция `[detect.prefix]`: subnet-level детектор распределённых атак. +#[derive(Debug, Clone, Deserialize)] +pub struct DetectPrefixConfig { + #[serde(default)] + pub enabled: bool, + #[serde(default = "default_prefix_syn_threshold")] + pub syn_threshold: u64, + #[serde(default = "default_prefix_window_secs")] + pub window_secs: u64, + #[serde(default = "default_prefix_min_unique_sources")] + pub min_unique_sources: u64, +} + +impl Default for DetectPrefixConfig { + fn default() -> Self { + Self { + enabled: false, + syn_threshold: default_prefix_syn_threshold(), + window_secs: default_prefix_window_secs(), + min_unique_sources: default_prefix_min_unique_sources(), + } + } +} + +fn default_prefix_syn_threshold() -> u64 { + 500 +} +fn default_prefix_window_secs() -> u64 { + 10 +} +fn default_prefix_min_unique_sources() -> u64 { + 16 +} diff --git a/src/engine/mod.rs b/src/engine/mod.rs index 86d7b50..36626ed 100644 --- a/src/engine/mod.rs +++ b/src/engine/mod.rs @@ -2,4 +2,6 @@ pub mod challenge; pub mod listener; +pub mod subnet_monitor; +pub mod subnet_tracker; pub mod tunnel; diff --git a/src/engine/subnet_monitor.rs b/src/engine/subnet_monitor.rs new file mode 100644 index 0000000..347cc04 --- /dev/null +++ b/src/engine/subnet_monitor.rs @@ -0,0 +1,143 @@ +//! Периодический subnet-монитор: источник снапшотов → детектор → эскалация. +//! +//! Источник выбирается по сборке и конфигу: XDP-карта `prefix_stats` +//! (feature `xdp` + загруженный фильтр) либо юзерспейс-агрегатор +//! [`SubnetTracker`]. + +use crate::config::DetectPrefixConfig; +use crate::engine::subnet_tracker::{SubnetTracker, TrackerSource}; +use crate::metrics; +#[cfg(feature = "xdp")] +use crate::traffic::prefix::XdpPrefixStats; +use crate::traffic::prefix::{PrefixSnapshot, PrefixStatsSource, PrefixVerdict, SubnetDetector, SubnetVerdict}; +use crate::xdp::XdpFilter; +use std::net::IpAddr; +use std::sync::{Arc, Mutex}; +use std::time::Duration; +use tokio::sync::watch; + +enum Source { + #[cfg(feature = "xdp")] + Xdp(XdpPrefixStats), + Tracker(TrackerSource), +} + +impl Source { + async fn snapshot(&self) -> anyhow::Result> { + match self { + #[cfg(feature = "xdp")] + Source::Xdp(xdp) => xdp.snapshot().await, + Source::Tracker(tracker) => tracker.snapshot().await, + } + } +} + +/// Запускает фоновый subnet-монитор, если `[detect.prefix] enabled`. +/// +/// Приоритет источника: XDP-карта при загруженном фильтре, иначе +/// юзерспейс-агрегатор. Без обоих источников монитор не запускается. +pub fn spawn( + detect: &DetectPrefixConfig, + tracker: Option>, + xdp: Option>>, + mut shutdown: watch::Receiver, +) { + if !detect.enabled { + return; + } + let Some(source) = build_source(tracker, xdp) else { + tracing::warn!("detect.prefix enabled but no stats source available; monitor disabled"); + return; + }; + let detector = SubnetDetector::from_config(detect); + let window = Duration::from_secs(detect.window_secs.max(1)); + + tokio::spawn(async move { + let mut tick = tokio::time::interval(window); + tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + loop { + tokio::select! { + biased; + _ = shutdown.changed() => { + if *shutdown.borrow() { + return; + } + } + _ = tick.tick() => { + match source.snapshot().await { + Ok(snaps) => apply_verdicts(&detector.evaluate(&snaps), &source).await, + Err(e) => tracing::debug!("prefix stats unavailable: {e}"), + } + } + } + } + }); +} + +fn build_source(tracker: Option>, xdp: Option>>) -> Option { + #[cfg(feature = "xdp")] + if let Some(filter) = xdp { + return Some(Source::Xdp(XdpPrefixStats::new(filter))); + } + #[cfg(not(feature = "xdp"))] + let _ = xdp; + tracker.map(|t| Source::Tracker(TrackerSource { tracker: t })) +} + +async fn apply_verdicts(verdicts: &[PrefixVerdict], source: &Source) { + for pv in verdicts { + metrics::SUBNET_VERDICTS + .with_label_values(&[verdict_name(pv.verdict)]) + .inc(); + tracing::info!( + prefix = %pv.prefix, + syn = pv.snapshot.syn_count, + unique = pv.snapshot.unique_sources, + verdict = verdict_name(pv.verdict), + "subnet verdict" + ); + if pv.verdict == SubnetVerdict::Block { + enforce_block(source, pv).await; + } + } +} + +async fn enforce_block(source: &Source, pv: &PrefixVerdict) { + match source { + #[cfg(feature = "xdp")] + Source::Xdp(xdp) => { + let result = { + let guard = xdp.filter_lock(); + guard.ban_cidr(pv.prefix, prefix_len(pv.prefix), 3600) + }; + if let Err(e) = result { + tracing::warn!("subnet block failed for {}: {e}", pv.prefix); + } + }, + Source::Tracker(_) => { + tracing::warn!( + prefix = %pv.prefix, + "block verdict without kernel enforcement; escalate at L7" + ); + }, + } +} + +/// Длина префикса, соответствующая гранулярности детекции (v4 /24, v6 /64). +#[cfg(feature = "xdp")] +fn prefix_len(prefix: IpAddr) -> u8 { + match prefix { + IpAddr::V4(_) => 24, + IpAddr::V6(_) => 64, + } +} + +#[must_use] +fn verdict_name(verdict: SubnetVerdict) -> &'static str { + match verdict { + SubnetVerdict::Monitor => "monitor", + SubnetVerdict::StrictLimit => "strict_limit", + SubnetVerdict::Challenge => "challenge", + SubnetVerdict::Block => "block", + } +} diff --git a/src/engine/subnet_tracker.rs b/src/engine/subnet_tracker.rs new file mode 100644 index 0000000..d14d60e --- /dev/null +++ b/src/engine/subnet_tracker.rs @@ -0,0 +1,131 @@ +//! Юзерспейс-агрегатор новых соединений по префиксам (/24, /64). +//! +//! Fallback-путь subnet-детектора, когда XDP выключен: питает +//! [`crate::traffic::prefix::SubnetDetector`] теми же снапшотами. + +use crate::traffic::prefix::{PrefixSnapshot, PrefixStatsSource, prefix_of}; +use dashmap::DashMap; +use std::net::IpAddr; +use std::sync::Arc; +use std::time::Instant; + +struct WindowState { + syn_count: u64, + sources: std::collections::HashSet, + started: Instant, +} + +/// Агрегатор новых соединений по префиксам. `take_snapshot` дренирует +/// накопленное: каждый вызов закрывает одно окно детекции. +pub struct SubnetTracker { + windows: DashMap, +} + +impl Default for SubnetTracker { + fn default() -> Self { + Self::new() + } +} + +impl SubnetTracker { + #[must_use] + pub fn new() -> Self { + Self { + windows: DashMap::new(), + } + } + + /// Регистрирует новое входящее соединение (юзерспейс-аналог SYN). + pub fn record_connection(&self, ip: IpAddr) { + let prefix = prefix_of(ip); + let mut state = self.windows.entry(prefix).or_insert_with(|| WindowState { + syn_count: 0, + sources: std::collections::HashSet::new(), + started: Instant::now(), + }); + state.syn_count = state.syn_count.saturating_add(1); + state.sources.insert(ip); + } + + /// Забирает и сбрасывает агрегаты текущего окна. + #[must_use] + pub fn take_snapshot(&self) -> Vec<(IpAddr, PrefixSnapshot)> { + let mut out = Vec::new(); + // Дренируем все ключи: каждое окно начинается с чистой карты, + // память не растёт при спуфинге множества префиксов. + self.windows.retain(|prefix, state| { + if state.syn_count == 0 { + return false; + } + out.push(( + *prefix, + PrefixSnapshot { + // Юзерспейс видит только установленные соединения; + // pkt_count отражает тот же поток. + pkt_count: state.syn_count, + syn_count: state.syn_count, + last_seen_ns: state.started.elapsed().as_nanos() as u64, + unique_sources: state.sources.len() as u64, + }, + )); + false + }); + out + } +} + +/// Адаптер [`SubnetTracker`] под трейт [`PrefixStatsSource`]. +pub struct TrackerSource { + pub tracker: Arc, +} + +impl PrefixStatsSource for TrackerSource { + async fn snapshot(&self) -> anyhow::Result> { + Ok(self.tracker.take_snapshot()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::net::{Ipv4Addr, Ipv6Addr}; + + fn v4(octets: [u8; 4]) -> IpAddr { + IpAddr::V4(Ipv4Addr::from(octets)) + } + + #[test] + fn groups_connections_by_prefix() { + let tracker = SubnetTracker::new(); + tracker.record_connection(v4([203, 0, 113, 1])); + tracker.record_connection(v4([203, 0, 113, 2])); + tracker.record_connection(v4([198, 51, 100, 9])); + let snaps = tracker.take_snapshot(); + assert_eq!(snaps.len(), 2); + let by_prefix: std::collections::HashMap = snaps.into_iter().collect(); + assert_eq!(by_prefix[&v4([203, 0, 113, 0])].syn_count, 2); + assert_eq!(by_prefix[&v4([203, 0, 113, 0])].unique_sources, 2); + assert_eq!(by_prefix[&v4([198, 51, 100, 0])].syn_count, 1); + } + + #[test] + fn snapshot_drains_window() { + let tracker = SubnetTracker::new(); + tracker.record_connection(v4([10, 0, 0, 1])); + assert_eq!(tracker.take_snapshot().len(), 1); + assert!(tracker.take_snapshot().is_empty()); + } + + #[test] + fn mapped_addresses_share_v4_prefix() { + let tracker = SubnetTracker::new(); + tracker.record_connection(v4([203, 0, 113, 5])); + // ::ffff:203.0.113.6 + let mapped = Ipv6Addr::new(0, 0, 0, 0, 0, 0xffff, 0xcb00, 0x7106); + tracker.record_connection(IpAddr::V6(mapped)); + let snaps = tracker.take_snapshot(); + assert_eq!(snaps.len(), 1); + assert_eq!(snaps[0].0, v4([203, 0, 113, 0])); + assert_eq!(snaps[0].1.unique_sources, 2); + } +} diff --git a/src/engine/tunnel.rs b/src/engine/tunnel.rs index f34bee5..8fd176c 100644 --- a/src/engine/tunnel.rs +++ b/src/engine/tunnel.rs @@ -2,6 +2,7 @@ use crate::config::Config; use crate::engine::challenge::{DifficultyAdjuster, enforce as enforce_pow}; +use crate::engine::subnet_tracker::SubnetTracker; use crate::filter::blacklist::Blacklist; use crate::filter::rate_limit::RateLimiter; use crate::metrics; @@ -30,6 +31,7 @@ pub struct Gateway { pub clickhouse: Option>>, pub allowed_1s: Arc, pub registry: Arc, + pub subnet_tracker: Option>, upstream_cursor: AtomicUsize, } @@ -46,6 +48,7 @@ impl Gateway { clickhouse: Option>>, allowed_1s: Arc, registry: Arc, + subnet_tracker: Option>, ) -> Self { Self { config, @@ -58,6 +61,7 @@ impl Gateway { clickhouse, allowed_1s, registry, + subnet_tracker, upstream_cursor: AtomicUsize::new(0), } } @@ -70,6 +74,10 @@ impl Gateway { pub async fn handle(&self, mut client: TcpStream, peer_addr: std::net::SocketAddr) -> anyhow::Result<()> { let peer_ip = peer_addr.ip(); + if let Some(tracker) = &self.subnet_tracker { + tracker.record_connection(peer_ip); + } + if self.blacklist.is_blocked(peer_ip) { metrics::CONNECTIONS_TOTAL.with_label_values(&["blocked"]).inc(); return Ok(()); diff --git a/src/metrics.rs b/src/metrics.rs index 74a835a..702dd00 100644 --- a/src/metrics.rs +++ b/src/metrics.rs @@ -29,6 +29,15 @@ pub static ATTACK_STATUS: LazyLock = LazyLock::new(|| { .expect("ATTACK_STATUS") }); +pub static SUBNET_VERDICTS: LazyLock = LazyLock::new(|| { + register_int_counter_vec!( + "rampart_subnet_verdicts_total", + "Subnet detector verdicts", + &["verdict"] + ) + .expect("SUBNET_VERDICTS") +}); + /// Отдаёт Prometheus-метрики по голому HTTP/0.9-совместимому ответу. pub async fn run_metrics_server(addr: &str) { let listener = match TcpListener::bind(addr).await { diff --git a/src/traffic/mod.rs b/src/traffic/mod.rs index f3daf29..1e7866b 100644 --- a/src/traffic/mod.rs +++ b/src/traffic/mod.rs @@ -3,5 +3,6 @@ pub mod alert; pub mod detector; pub mod ewma; +pub mod prefix; pub mod profiler; pub mod reputation; diff --git a/src/traffic/prefix.rs b/src/traffic/prefix.rs new file mode 100644 index 0000000..0a4a81b --- /dev/null +++ b/src/traffic/prefix.rs @@ -0,0 +1,296 @@ +//! Subnet-level (prefix) детектор распределённых атак. +//! +//! Источник данных — трейт [`PrefixStatsSource`]: карта XDP `prefix_stats` +//! (feature `xdp`) либо юзерспейс-агрегатор [`crate::engine::subnet_tracker`]. +//! Контракт с XDP-агентом: [`PrefixKey`] / [`PrefixStatsVal`] должны совпадать +//! побайтово со структурами ядра. + +use anyhow::Result; +use std::future::Future; +use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; + +/// `family` для IPv4-префиксов (/24). +pub const PREFIX_FAMILY_V4: u8 = 4; +/// `family` для IPv6-префиксов (/64). +pub const PREFIX_FAMILY_V6: u8 = 6; + +/// Ключ карты `prefix_stats` (`BPF_MAP_TYPE_LRU_HASH`, max_entries 65536). +#[repr(C)] +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct PrefixKey { + pub family: u8, + pub pad: [u8; 7], + pub addr: [u8; 16], +} + +/// Значение карты `prefix_stats`. +#[repr(C)] +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct PrefixStatsVal { + pub syn_count: u64, + pub pkt_count: u64, + pub last_seen_ns: u64, +} + +impl PrefixStatsVal { + /// Разбирает сырые 24 байта значения карты. + #[must_use] + pub fn from_bytes(bytes: [u8; 24]) -> Self { + let read_u64 = + |range: std::ops::Range| u64::from_le_bytes(bytes[range].try_into().expect("8-byte field")); + Self { + syn_count: read_u64(0..8), + pkt_count: read_u64(8..16), + last_seen_ns: read_u64(16..24), + } + } + + /// Сериализует значение в сырые 24 байта карты. + #[must_use] + pub fn to_bytes(self) -> [u8; 24] { + let mut out = [0u8; 24]; + out[0..8].copy_from_slice(&self.syn_count.to_le_bytes()); + out[8..16].copy_from_slice(&self.pkt_count.to_le_bytes()); + out[16..24].copy_from_slice(&self.last_seen_ns.to_le_bytes()); + out + } +} + +impl PrefixKey { + /// Строит ключ префикса для адреса (v4 → /24 в первых 4 байтах, v6 → /64). + #[must_use] + pub fn new(ip: IpAddr) -> Self { + let mut key = Self::default(); + match ip { + IpAddr::V4(v4) => { + key.family = PREFIX_FAMILY_V4; + key.addr[..4].copy_from_slice(&masked_v4(v4).octets()); + }, + IpAddr::V6(v6) => { + key.family = PREFIX_FAMILY_V6; + key.addr.copy_from_slice(&masked_v6_segments(v6)); + }, + } + key + } + + /// Префикс как IP-адрес (маскированный). + #[must_use] + pub fn to_prefix(self) -> Option { + match self.family { + PREFIX_FAMILY_V4 => Some(IpAddr::V4(Ipv4Addr::new( + self.addr[0], + self.addr[1], + self.addr[2], + self.addr[3], + ))), + PREFIX_FAMILY_V6 => Some(IpAddr::V6(Ipv6Addr::from(self.addr))), + _ => None, + } + } + + #[must_use] + pub fn from_bytes(bytes: [u8; 24]) -> Self { + Self { + family: bytes[0], + pad: bytes[1..8].try_into().expect("7 pad bytes"), + addr: bytes[8..24].try_into().expect("16 addr bytes"), + } + } + + #[must_use] + pub fn to_bytes(self) -> [u8; 24] { + let mut out = [0u8; 24]; + out[0] = self.family; + out[1..8].copy_from_slice(&self.pad); + out[8..24].copy_from_slice(&self.addr); + out + } +} + +/// Агрегат по одному префиксу за окно детекции. +/// +/// `unique_sources == 0` означает, что источник не считает уникальные адреса +/// (XDP-путь); spoof-gate в этом случае Block не подтверждает. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct PrefixSnapshot { + pub syn_count: u64, + pub pkt_count: u64, + pub last_seen_ns: u64, + pub unique_sources: u64, +} + +/// Источник агрегатов по префиксам. Позволяет тестам работать без ядра. +pub trait PrefixStatsSource { + /// Снимок агрегатов за текущее окно. + /// + /// # Errors + /// Ошибка чтения источника (карта не загружена, BPF-сбой и т.п.). + fn snapshot(&self) -> impl Future>> + Send; +} + +/// Возвращает префикс адреса: IPv4 и IPv4-mapped → /24, IPv6 → /64. +#[must_use] +pub fn prefix_of(ip: IpAddr) -> IpAddr { + match ip { + IpAddr::V4(v4) => IpAddr::V4(masked_v4(v4)), + IpAddr::V6(v6) => match v6.to_ipv4_mapped() { + // ::ffff:a.b.c.d живёт в том же /24-пространстве, что и a.b.c.d. + Some(v4) => IpAddr::V4(masked_v4(v4)), + None => IpAddr::V6(Ipv6Addr::from(masked_v6_segments(v6))), + }, + } +} + +fn masked_v4(ip: Ipv4Addr) -> Ipv4Addr { + let o = ip.octets(); + Ipv4Addr::new(o[0], o[1], o[2], 0) +} + +fn masked_v6_segments(ip: Ipv6Addr) -> [u8; 16] { + let mut segs = ip.octets(); + segs[8..].fill(0); + segs +} + +/// Эскалация вердиктов subnet-level детектора. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum SubnetVerdict { + Monitor, + StrictLimit, + Challenge, + Block, +} + +/// Вердикт по одному префиксу. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct PrefixVerdict { + pub prefix: IpAddr, + pub verdict: SubnetVerdict, + pub snapshot: PrefixSnapshot, +} + +/// Детектор распределённых атак по префиксам. +/// +/// Лестница эскалации по `syn_count` за окно: `< threshold` → Monitor, +/// `< 2×threshold` → StrictLimit, иначе Challenge; Block только при +/// `≥ 4×threshold` И прохождении spoof-gate (≥ `min_unique_sources` +/// уникальных IP — защита от ложного бана CGNAT/спуфинга). +#[derive(Debug, Clone)] +pub struct SubnetDetector { + syn_threshold: u64, + min_unique_sources: u64, +} + +impl SubnetDetector { + #[must_use] + pub fn new(syn_threshold: u64, min_unique_sources: u64) -> Self { + Self { + syn_threshold: syn_threshold.max(1), + min_unique_sources, + } + } + + #[must_use] + pub fn from_config(config: &crate::config::DetectPrefixConfig) -> Self { + Self::new(config.syn_threshold, config.min_unique_sources) + } + + /// Оценивает снимок и возвращает вердикты по активным префиксам. + #[must_use] + pub fn evaluate(&self, snapshots: &[(IpAddr, PrefixSnapshot)]) -> Vec { + snapshots + .iter() + .filter(|(_, snap)| snap.syn_count > 0) + .map(|(prefix, snap)| PrefixVerdict { + prefix: *prefix, + verdict: self.verdict_for(snap), + snapshot: *snap, + }) + .collect() + } + + #[must_use] + fn verdict_for(&self, snap: &PrefixSnapshot) -> SubnetVerdict { + let t = self.syn_threshold; + if snap.syn_count < t { + return SubnetVerdict::Monitor; + } + if snap.syn_count < t.saturating_mul(2) { + return SubnetVerdict::StrictLimit; + } + let extreme = snap.syn_count >= t.saturating_mul(4); + if extreme && self.spoof_gate_passed(snap) { + return SubnetVerdict::Block; + } + SubnetVerdict::Challenge + } + + #[must_use] + fn spoof_gate_passed(&self, snap: &PrefixSnapshot) -> bool { + // gate отключён (min_unique_sources == 0) либо подтверждено разнообразие. + self.min_unique_sources == 0 || snap.unique_sources >= self.min_unique_sources + } +} + +/// Реализация [`PrefixStatsSource`] поверх XDP-фильтра (feature `xdp`). +/// Читает карту `prefix_stats`; уникальные источники ядро не считает, +/// поэтому spoof-gate на этом пути Block не подтверждает. +#[cfg(feature = "xdp")] +pub struct XdpPrefixStats { + filter: std::sync::Arc>, +} + +#[cfg(feature = "xdp")] +impl XdpPrefixStats { + #[must_use] + pub fn new(filter: std::sync::Arc>) -> Self { + Self { filter } + } + + /// Блокирует мьютекс фильтра (для kernel-enforcement вердиктов). + /// + /// # Panics + /// Если мьютекс отравлен (паника в другом потоке во время удержания). + pub fn filter_lock(&self) -> std::sync::MutexGuard<'_, crate::xdp::XdpFilter> { + self.filter.lock().expect("xdp lock poisoned") + } +} + +#[cfg(feature = "xdp")] +impl PrefixStatsSource for XdpPrefixStats { + async fn snapshot(&self) -> Result> { + let raw = self.filter_lock().read_prefix_stats()?; + Ok(raw + .into_iter() + .filter_map(|(key, val)| { + let prefix = key.to_prefix()?; + Some(( + prefix, + PrefixSnapshot { + syn_count: val.syn_count, + pkt_count: val.pkt_count, + last_seen_ns: val.last_seen_ns, + unique_sources: 0, + }, + )) + }) + .collect()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn stats_val_roundtrip() { + let val = PrefixStatsVal { + syn_count: 42, + pkt_count: 100, + last_seen_ns: 1_234_567_890, + }; + assert_eq!(PrefixStatsVal::from_bytes(val.to_bytes()), val); + assert_eq!(size_of::(), 24); + } +} diff --git a/src/xdp/filter.rs b/src/xdp/filter.rs index cb06a16..907b57f 100644 --- a/src/xdp/filter.rs +++ b/src/xdp/filter.rs @@ -1,9 +1,11 @@ use anyhow::{Context, Result, bail}; use libbpf_rs::{MapCore, MapFlags, Object, ObjectBuilder, RingBuffer, RingBufferBuilder, Xdp, XdpFlags}; +use std::net::IpAddr; use std::net::Ipv4Addr; use std::os::unix::io::AsFd; use super::XdpStats; +use crate::traffic::prefix::{PrefixKey, PrefixStatsVal}; pub struct XdpFilter { obj: Option, @@ -78,17 +80,28 @@ impl XdpFilter { } pub fn ban_ip(&self, ip: Ipv4Addr, duration_secs: u64) -> Result<()> { + self.blacklist_update(32, &ip.octets(), duration_secs) + } + + /// Бан CIDR через `blacklist_map` (LPM trie): длина префикса в key[0]. + /// + /// # Errors + /// XDP не загружен, IPv6-префиксы (ядро их пока не банит) или ошибка карты. + pub fn ban_cidr(&self, prefix: IpAddr, prefix_len: u8, duration_secs: u64) -> Result<()> { + let IpAddr::V4(net) = prefix else { + bail!("ipv6 cidr bans not supported yet"); + }; + self.blacklist_update(prefix_len, &net.octets(), duration_secs) + } + + fn blacklist_update(&self, prefix_len: u8, v4: &[u8; 4], duration_secs: u64) -> Result<()> { let map = self.find_map("blacklist_map")?; let mut key = [0u8; 8]; - key[0] = 32; - key[4..8].copy_from_slice(&ip.octets()); - let now = std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .unwrap_or_default() - .as_nanos() as u64; + key[0] = prefix_len; + key[4..8].copy_from_slice(v4); map.update( &key, - &(now + duration_secs * 1_000_000_000).to_le_bytes(), + &(unix_ns() + duration_secs * 1_000_000_000).to_le_bytes(), MapFlags::ANY, )?; Ok(()) @@ -103,6 +116,28 @@ impl XdpFilter { Ok(()) } + /// Читает карту `prefix_stats` (агрегаты по префиксам /24 и /64). + /// + /// # Errors + /// XDP не загружен, карта отсутствует в объекте или ошибка BPF-чтения. + pub fn read_prefix_stats(&self) -> Result> { + let map = self.find_map("prefix_stats")?; + let mut out = Vec::new(); + for key in map.keys() { + let Ok(kb) = <[u8; 24]>::try_from(key.as_slice()) else { + continue; + }; + let Some(val) = map.lookup(&kb, MapFlags::ANY).context("prefix_stats lookup failed")? else { + continue; + }; + let Ok(vb) = <[u8; 24]>::try_from(val.as_slice()) else { + continue; + }; + out.push((PrefixKey::from_bytes(kb), PrefixStatsVal::from_bytes(vb))); + } + Ok(out) + } + pub fn get_stats(&self) -> Result { let map = self.find_map("stats_map")?; let sum = |idx: u32| -> u64 { @@ -137,6 +172,13 @@ impl Drop for XdpFilter { } } +fn unix_ns() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_nanos() as u64 +} + fn build_ringbuf(obj: &Object) -> Result> { let map = obj .maps() diff --git a/src/xdp/noop.rs b/src/xdp/noop.rs index bcf3636..9cb050f 100644 --- a/src/xdp/noop.rs +++ b/src/xdp/noop.rs @@ -1,5 +1,7 @@ -use anyhow::Result; -use std::net::Ipv4Addr; +use anyhow::{Ok, Result}; +use std::net::{IpAddr, Ipv4Addr}; + +use crate::traffic::prefix::{PrefixKey, PrefixStatsVal}; pub struct XdpFilter; @@ -17,9 +19,21 @@ impl XdpFilter { pub fn ban_ip(&self, _ip: Ipv4Addr, _duration_secs: u64) -> Result<()> { Ok(()) } + pub fn ban_cidr(&self, _prefix: IpAddr, _prefix_len: u8, _duration_secs: u64) -> Result<()> { + Ok(()) + } pub fn unban_ip(&self, _ip: Ipv4Addr) -> Result<()> { Ok(()) } + + /// Честная ошибка: чтение `prefix_stats` не подключено без feature `xdp`. + /// + /// # Errors + /// Всегда — сборка без XDP не имеет доступа к картам ядра. + pub fn read_prefix_stats(&self) -> Result> { + anyhow::bail!("prefix_stats reading not wired yet: built without xdp feature") + } + pub fn get_stats(&self) -> Result { Ok(super::XdpStats::default()) } diff --git a/tests/subnet_detection.rs b/tests/subnet_detection.rs new file mode 100644 index 0000000..9062de6 --- /dev/null +++ b/tests/subnet_detection.rs @@ -0,0 +1,215 @@ +//! Интеграционные тесты subnet-level детекции: трекер → детектор → вердикты. + +use rampart::config::Config; +use rampart::engine::subnet_tracker::SubnetTracker; +use rampart::traffic::prefix::{ + PrefixKey, PrefixSnapshot, PrefixStatsSource, PrefixStatsVal, SubnetDetector, SubnetVerdict, prefix_of, +}; +use std::collections::HashMap; +use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; + +fn v4(octets: [u8; 4]) -> IpAddr { + IpAddr::V4(Ipv4Addr::from(octets)) +} + +/// Фиксированный источник снапшотов — имитирует XDP-карту без ядра. +struct FixedSource(Vec<(IpAddr, PrefixSnapshot)>); + +impl PrefixStatsSource for FixedSource { + async fn snapshot(&self) -> anyhow::Result> { + Ok(self.0.clone()) + } +} + +#[tokio::test] +async fn tracker_feeds_detector_and_blocks_distributed_attack() { + let tracker = SubnetTracker::new(); + // 40 уникальных источников из одного /24 — распределённая атака + // (6 SYN от каждого: 240 >= 4 × threshold(50)). + for _ in 0..6 { + for last in 1..=40u8 { + tracker.record_connection(v4([203, 0, 113, last])); + } + } + let snaps = tracker.take_snapshot(); + + let detector = SubnetDetector::new(50, 16); + let verdicts = detector.evaluate(&snaps); + assert_eq!(verdicts.len(), 1); + assert_eq!(verdicts[0].prefix, v4([203, 0, 113, 0])); + assert_eq!(verdicts[0].verdict, SubnetVerdict::Block); + assert_eq!(verdicts[0].snapshot.syn_count, 240); + assert_eq!(verdicts[0].snapshot.unique_sources, 40); + + // Полный конвейер через трейт-источник (как в бинарнике). + let source_snaps = FixedSource(snaps).snapshot().await.expect("source snapshot"); + assert_eq!(detector.evaluate(&source_snaps)[0].verdict, SubnetVerdict::Block); +} + +#[tokio::test] +async fn spoof_gate_many_syn_few_sources_no_block() { + let tracker = SubnetTracker::new(); + // 500 SYN всего из 3 адресов — всплеск, но не распределённый. + for _ in 0..167 { + for ip in [v4([198, 51, 100, 1]), v4([198, 51, 100, 2]), v4([198, 51, 100, 3])] { + tracker.record_connection(ip); + } + } + let snaps = tracker.take_snapshot(); + let detector = SubnetDetector::new(50, 16); + let verdicts = detector.evaluate(&snaps); + assert_eq!(verdicts.len(), 1); + assert_eq!(verdicts[0].snapshot.syn_count, 501); + assert_eq!(verdicts[0].snapshot.unique_sources, 3); + assert_ne!(verdicts[0].verdict, SubnetVerdict::Block); + assert_eq!(verdicts[0].verdict, SubnetVerdict::Challenge); +} + +#[test] +fn ipv6_prefixes_aggregated_separately_from_ipv4() { + let tracker = SubnetTracker::new(); + let v6a = IpAddr::V6("2001:db8:1:2::1".parse::().expect("valid")); + let v6b = IpAddr::V6("2001:db8:1:2::ffff".parse::().expect("valid")); + tracker.record_connection(v6a); + tracker.record_connection(v6b); + tracker.record_connection(v4([10, 0, 0, 1])); + + let snaps: HashMap = tracker.take_snapshot().into_iter().collect(); + assert_eq!(snaps.len(), 2); + assert_eq!( + snaps[&IpAddr::V6("2001:db8:1:2::".parse::().expect("valid"))].unique_sources, + 2 + ); + assert_eq!(snaps[&v4([10, 0, 0, 0])].unique_sources, 1); +} + +#[test] +fn mapped_v6_joins_ipv4_prefix_bucket() { + // ::ffff:192.0.2.9 + assert_eq!( + prefix_of(IpAddr::V6(Ipv6Addr::new(0, 0, 0, 0, 0, 0xffff, 0xc000, 0x0209))), + v4([192, 0, 2, 0]) + ); +} + +#[test] +fn detect_prefix_config_section_parses_end_to_end() { + let config = Config::parse_str( + r#" +[detect.prefix] +enabled = true +syn_threshold = 250 +window_secs = 30 +min_unique_sources = 4 +"#, + ) + .expect("[detect.prefix] must parse"); + assert!(config.detect.prefix.enabled); + let detector = SubnetDetector::from_config(&config.detect.prefix); + let snaps = vec![( + v4([203, 0, 113, 0]), + PrefixSnapshot { + syn_count: 1200, + unique_sources: 5, + ..PrefixSnapshot::default() + }, + )]; + let verdicts = detector.evaluate(&snaps); + assert_eq!(verdicts[0].verdict, SubnetVerdict::Block); +} + +fn verdict_for(detector: &SubnetDetector, snap: &PrefixSnapshot) -> SubnetVerdict { + let verdicts = detector.evaluate(&[(v4([10, 9, 9, 0]), *snap)]); + verdicts[0].verdict +} + +#[test] +fn ipv4_prefix_masked_to_24() { + assert_eq!(prefix_of(v4([203, 0, 113, 77])), v4([203, 0, 113, 0])); + assert_eq!(prefix_of(v4([198, 51, 100, 1])), v4([198, 51, 100, 0])); +} + +#[test] +fn ipv6_prefix_masked_to_64() { + let ip = IpAddr::V6("2001:db8:1:2:3:4:5:6".parse::().expect("valid")); + assert_eq!( + prefix_of(ip), + IpAddr::V6("2001:db8:1:2::".parse::().expect("valid")) + ); +} + +#[test] +fn mapped_ipv6_treated_as_ipv4() { + // ::ffff:203.0.113.77 + let mapped = IpAddr::V6(Ipv6Addr::new(0, 0, 0, 0, 0, 0xffff, 0xcb00, 0x714d)); + assert_eq!(prefix_of(mapped), v4([203, 0, 113, 0])); +} + +#[test] +fn prefix_key_roundtrip_and_layout() { + assert_eq!(size_of::(), 24); + assert_eq!(size_of::(), 24); + let key = PrefixKey::new(v4([10, 1, 2, 3])); + assert_eq!(key.family, rampart::traffic::prefix::PREFIX_FAMILY_V4); + assert_eq!(key.to_prefix(), Some(v4([10, 1, 2, 0]))); + assert_eq!(PrefixKey::from_bytes(key.to_bytes()), key); + + let key6 = PrefixKey::new(IpAddr::V6("2001:db8::1".parse::().expect("valid"))); + assert_eq!(key6.family, rampart::traffic::prefix::PREFIX_FAMILY_V6); + assert_eq!( + key6.to_prefix(), + Some(IpAddr::V6("2001:db8::".parse::().expect("valid"))) + ); +} + +#[test] +fn escalation_ladder_monitor_to_block() { + let detector = SubnetDetector::new(100, 16); + let snap = |syn: u64| PrefixSnapshot { + syn_count: syn, + unique_sources: 32, + ..PrefixSnapshot::default() + }; + assert_eq!(verdict_for(&detector, &snap(99)), SubnetVerdict::Monitor); + assert_eq!(verdict_for(&detector, &snap(100)), SubnetVerdict::StrictLimit); + assert_eq!(verdict_for(&detector, &snap(150)), SubnetVerdict::StrictLimit); + assert_eq!(verdict_for(&detector, &snap(200)), SubnetVerdict::Challenge); + assert_eq!(verdict_for(&detector, &snap(400)), SubnetVerdict::Block); +} + +#[test] +fn spoof_gate_downgrades_block_to_challenge() { + let detector = SubnetDetector::new(100, 16); + // Много SYN из одного источника (или неизвестное разнообразие) → нет Block. + let few = PrefixSnapshot { + syn_count: 1000, + unique_sources: 3, + ..PrefixSnapshot::default() + }; + let unknown = PrefixSnapshot { + syn_count: 1000, + ..PrefixSnapshot::default() + }; + assert_eq!(verdict_for(&detector, &few), SubnetVerdict::Challenge); + assert_eq!(verdict_for(&detector, &unknown), SubnetVerdict::Challenge); +} + +#[test] +fn evaluate_skips_empty_prefixes() { + let detector = SubnetDetector::new(100, 0); + let snaps = vec![ + (v4([10, 0, 0, 0]), PrefixSnapshot::default()), + ( + v4([10, 0, 1, 0]), + PrefixSnapshot { + syn_count: 500, + unique_sources: 40, + ..PrefixSnapshot::default() + }, + ), + ]; + let verdicts = detector.evaluate(&snaps); + assert_eq!(verdicts.len(), 1); + assert_eq!(verdicts[0].prefix, v4([10, 0, 1, 0])); + assert_eq!(verdicts[0].verdict, SubnetVerdict::Block); +} diff --git a/xdp/core/prefix_stats.h b/xdp/core/prefix_stats.h new file mode 100644 index 0000000..7c4b7f9 --- /dev/null +++ b/xdp/core/prefix_stats.h @@ -0,0 +1,76 @@ +#ifndef RAMPART_PREFIX_STATS_H +#define RAMPART_PREFIX_STATS_H + +#include "common.h" + +// ── Per-prefix traffic statistics (shared ABI with Rust userspace) ── +// Layout is part of the userspace contract — do not reorder fields. + +// family: 4 = IPv4 /24 префикс, 6 = IPv6 /64 префикс +struct prefix_key { + __u8 family; + __u8 pad[7]; + __u8 addr[16]; // network-order байты префикса, zero-padded +}; + +struct prefix_stats_val { + __u64 syn_count; + __u64 pkt_count; + __u64 last_seen_ns; // bpf_ktime_get_ns() +}; + +// 🔗 Prefix statistics (LRU) — читается и сбрасывается юзерспейсом +struct { + __uint(type, BPF_MAP_TYPE_LRU_HASH); + __uint(max_entries, 65536); + __type(key, struct prefix_key); + __type(value, struct prefix_stats_val); +} prefix_stats SEC(".maps"); + +#define PREFIX_FAMILY_V4 4 +#define PREFIX_FAMILY_V6 6 + +// ── Build /24 IPv4 prefix key (network-order, low 8 bits masked out) ── +static __always_inline void prefix_key_from_v4(struct prefix_key *key, + __u32 src_ip_net_order) +{ + __builtin_memset(key, 0, sizeof(*key)); + key->family = PREFIX_FAMILY_V4; + __builtin_memcpy(key->addr, &src_ip_net_order, 4); + key->addr[3] = 0; // обнулить младшие 8 бит адреса +} + +// ── Build /64 IPv6 prefix key (network-order, low 64 bits masked out) ── +static __always_inline void prefix_key_from_v6(struct prefix_key *key, + const __u8 src_addr[16]) +{ + __builtin_memset(key, 0, sizeof(*key)); + key->family = PREFIX_FAMILY_V6; + __builtin_memcpy(key->addr, src_addr, 8); +} + +// ── Increment per-prefix counters (sliding window handled in userspace) ── +// is_syn: 1 → syn_count++, 0 → pkt_count++; last_seen_ns updated всегда. +static __always_inline void update_prefix_stats(struct prefix_key *key, + __u8 is_syn, __u64 now) +{ + struct prefix_stats_val *val = bpf_map_lookup_elem(&prefix_stats, key); + + if (!val) { + struct prefix_stats_val init = { + .syn_count = is_syn ? 1 : 0, + .pkt_count = is_syn ? 0 : 1, + .last_seen_ns = now, + }; + bpf_map_update_elem(&prefix_stats, key, &init, BPF_ANY); + return; + } + + if (is_syn) + __sync_fetch_and_add(&val->syn_count, 1); + else + __sync_fetch_and_add(&val->pkt_count, 1); + val->last_seen_ns = now; +} + +#endif /* RAMPART_PREFIX_STATS_H */ diff --git a/xdp/core/universal_filter.c b/xdp/core/universal_filter.c index 1e82730..019a6d0 100644 --- a/xdp/core/universal_filter.c +++ b/xdp/core/universal_filter.c @@ -26,6 +26,7 @@ #include "maps.h" #include "config.h" #include "stats.h" +#include "prefix_stats.h" #include "../hooks/hook_api.h" char __license[] SEC("license") = "GPL"; @@ -265,6 +266,10 @@ int rampart_universal_filter(struct xdp_md *ctx) if ((void *)data + sizeof(struct ethhdr) + ip_hdr_len + tcp_hdr_len > data_end) return XDP_DROP; + // ── Per-prefix stats key (/24; built once, after blacklist/bypass drops) ── + struct prefix_key pkey; + prefix_key_from_v4(&pkey, src_ip); + // Compute payload pointers __u8 *tcp_payload = (__u8 *)tcp + tcp_hdr_len; __u8 *tcp_payload_end = (__u8 *)data_end; @@ -294,6 +299,9 @@ int rampart_universal_filter(struct xdp_md *ctx) return XDP_DROP; } + // Count SYN only after throttle passed (don't count throttled SYNs) + update_prefix_stats(&pkey, 1, now); + // Create conntrack entry for new connection struct conntrack_entry ce = {}; ce.state = STATE_SYN_RECEIVED; @@ -313,6 +321,9 @@ int rampart_universal_filter(struct xdp_md *ctx) } // ── Connection tracking lookup ── + // pkt_count: all non-SYN TCP that survived blacklist/bypass/throttle + update_prefix_stats(&pkey, 0, now); + struct conntrack_entry *conn = bpf_map_lookup_elem(&conntrack_map, &flow); if (!conn) { // Unknown connection — drop