feat: subnet-level attack detection (prefix_stats)

- XDP: prefix_stats LRU map, per-/24 (v4) and /64 (v6) SYN/packet counters,
  incremented post-blacklist/throttle (xdp/core/prefix_stats.h)
- userspace: SubnetDetector escalation ladder Monitor->StrictLimit->Challenge->Block
  with spoof-gate (Block requires >= min_unique_sources, CGNAT-safe)
- config: [detect.prefix] section (enabled=false by default)
- fallback without XDP: engine SubnetTracker aggregates connections per-prefix
- fix: PrefixStatsVal::from_bytes for xdp feature build; prefix_len 24 vs 64

cargo build/clippy(-D warnings, all features)/test green: 83 tests
This commit is contained in:
loki5512344 2026-08-24 09:47:21 +02:00
parent 15f474486a
commit 40bfe956e2
Signed by: boba
GPG key ID: 253067914055423B
14 changed files with 1044 additions and 11 deletions

View file

@ -2,8 +2,8 @@
mod sections; mod sections;
pub use sections::{ pub use sections::{
BackendConfig, BanConfig, BindConfig, LimitsConfig, LoggingConfig, MetricsConfig, PowConfig, StoreConfig, BackendConfig, BanConfig, BindConfig, DetectConfig, DetectPrefixConfig, LimitsConfig, LoggingConfig, MetricsConfig,
WorkerConfig, XdpConfig, PowConfig, StoreConfig, WorkerConfig, XdpConfig,
}; };
use serde::Deserialize; use serde::Deserialize;
@ -32,6 +32,8 @@ pub struct Config {
#[serde(default)] #[serde(default)]
pub pow: PowConfig, pub pow: PowConfig,
#[serde(default)] #[serde(default)]
pub detect: DetectConfig,
#[serde(default)]
pub whitelist: Vec<String>, pub whitelist: Vec<String>,
} }
@ -70,6 +72,9 @@ impl Config {
anyhow::bail!("invalid whitelist entry: {entry}"); 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(()) Ok(())
} }
} }
@ -140,4 +145,43 @@ upstreams = ["not-an-addr"]
let result = Config::parse_str("whitelist = [\"999.999.1.1\"]"); let result = Config::parse_str("whitelist = [\"999.999.1.1\"]");
assert!(result.is_err()); 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());
}
} }

View file

@ -226,3 +226,44 @@ impl Default for PowConfig {
fn default_pow_difficulty() -> u8 { fn default_pow_difficulty() -> u8 {
4 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
}

View file

@ -2,4 +2,6 @@
pub mod challenge; pub mod challenge;
pub mod listener; pub mod listener;
pub mod subnet_monitor;
pub mod subnet_tracker;
pub mod tunnel; pub mod tunnel;

View file

@ -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<Vec<(IpAddr, PrefixSnapshot)>> {
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<Arc<SubnetTracker>>,
xdp: Option<Arc<Mutex<XdpFilter>>>,
mut shutdown: watch::Receiver<bool>,
) {
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<Arc<SubnetTracker>>, xdp: Option<Arc<Mutex<XdpFilter>>>) -> Option<Source> {
#[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",
}
}

View file

@ -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<IpAddr>,
started: Instant,
}
/// Агрегатор новых соединений по префиксам. `take_snapshot` дренирует
/// накопленное: каждый вызов закрывает одно окно детекции.
pub struct SubnetTracker {
windows: DashMap<IpAddr, WindowState>,
}
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<SubnetTracker>,
}
impl PrefixStatsSource for TrackerSource {
async fn snapshot(&self) -> anyhow::Result<Vec<(IpAddr, PrefixSnapshot)>> {
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<IpAddr, PrefixSnapshot> = 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);
}
}

View file

@ -2,6 +2,7 @@
use crate::config::Config; use crate::config::Config;
use crate::engine::challenge::{DifficultyAdjuster, enforce as enforce_pow}; use crate::engine::challenge::{DifficultyAdjuster, enforce as enforce_pow};
use crate::engine::subnet_tracker::SubnetTracker;
use crate::filter::blacklist::Blacklist; use crate::filter::blacklist::Blacklist;
use crate::filter::rate_limit::RateLimiter; use crate::filter::rate_limit::RateLimiter;
use crate::metrics; use crate::metrics;
@ -30,6 +31,7 @@ pub struct Gateway {
pub clickhouse: Option<Arc<TokioMutex<ClickHouseWriter>>>, pub clickhouse: Option<Arc<TokioMutex<ClickHouseWriter>>>,
pub allowed_1s: Arc<AtomicU64>, pub allowed_1s: Arc<AtomicU64>,
pub registry: Arc<ProtocolRegistry>, pub registry: Arc<ProtocolRegistry>,
pub subnet_tracker: Option<Arc<SubnetTracker>>,
upstream_cursor: AtomicUsize, upstream_cursor: AtomicUsize,
} }
@ -46,6 +48,7 @@ impl Gateway {
clickhouse: Option<Arc<TokioMutex<ClickHouseWriter>>>, clickhouse: Option<Arc<TokioMutex<ClickHouseWriter>>>,
allowed_1s: Arc<AtomicU64>, allowed_1s: Arc<AtomicU64>,
registry: Arc<ProtocolRegistry>, registry: Arc<ProtocolRegistry>,
subnet_tracker: Option<Arc<SubnetTracker>>,
) -> Self { ) -> Self {
Self { Self {
config, config,
@ -58,6 +61,7 @@ impl Gateway {
clickhouse, clickhouse,
allowed_1s, allowed_1s,
registry, registry,
subnet_tracker,
upstream_cursor: AtomicUsize::new(0), 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<()> { pub async fn handle(&self, mut client: TcpStream, peer_addr: std::net::SocketAddr) -> anyhow::Result<()> {
let peer_ip = peer_addr.ip(); 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) { if self.blacklist.is_blocked(peer_ip) {
metrics::CONNECTIONS_TOTAL.with_label_values(&["blocked"]).inc(); metrics::CONNECTIONS_TOTAL.with_label_values(&["blocked"]).inc();
return Ok(()); return Ok(());

View file

@ -29,6 +29,15 @@ pub static ATTACK_STATUS: LazyLock<IntGauge> = LazyLock::new(|| {
.expect("ATTACK_STATUS") .expect("ATTACK_STATUS")
}); });
pub static SUBNET_VERDICTS: LazyLock<IntCounterVec> = LazyLock::new(|| {
register_int_counter_vec!(
"rampart_subnet_verdicts_total",
"Subnet detector verdicts",
&["verdict"]
)
.expect("SUBNET_VERDICTS")
});
/// Отдаёт Prometheus-метрики по голому HTTP/0.9-совместимому ответу. /// Отдаёт Prometheus-метрики по голому HTTP/0.9-совместимому ответу.
pub async fn run_metrics_server(addr: &str) { pub async fn run_metrics_server(addr: &str) {
let listener = match TcpListener::bind(addr).await { let listener = match TcpListener::bind(addr).await {

View file

@ -3,5 +3,6 @@
pub mod alert; pub mod alert;
pub mod detector; pub mod detector;
pub mod ewma; pub mod ewma;
pub mod prefix;
pub mod profiler; pub mod profiler;
pub mod reputation; pub mod reputation;

296
src/traffic/prefix.rs Normal file
View file

@ -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<usize>| 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<IpAddr> {
match self.family {
PREFIX_FAMILY_V4 => Some(IpAddr::V4(Ipv4Addr::new(
self.addr[0],
self.addr[1],
self.addr[2],
self.addr[3],
))),
PREFIX_FAMILY_V6 => Some(IpAddr::V6(Ipv6Addr::from(self.addr))),
_ => None,
}
}
#[must_use]
pub fn from_bytes(bytes: [u8; 24]) -> Self {
Self {
family: bytes[0],
pad: bytes[1..8].try_into().expect("7 pad bytes"),
addr: bytes[8..24].try_into().expect("16 addr bytes"),
}
}
#[must_use]
pub fn to_bytes(self) -> [u8; 24] {
let mut out = [0u8; 24];
out[0] = self.family;
out[1..8].copy_from_slice(&self.pad);
out[8..24].copy_from_slice(&self.addr);
out
}
}
/// Агрегат по одному префиксу за окно детекции.
///
/// `unique_sources == 0` означает, что источник не считает уникальные адреса
/// (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<Output = Result<Vec<(IpAddr, PrefixSnapshot)>>> + Send;
}
/// Возвращает префикс адреса: IPv4 и IPv4-mapped → /24, IPv6 → /64.
#[must_use]
pub fn prefix_of(ip: IpAddr) -> IpAddr {
match ip {
IpAddr::V4(v4) => IpAddr::V4(masked_v4(v4)),
IpAddr::V6(v6) => match v6.to_ipv4_mapped() {
// ::ffff:a.b.c.d живёт в том же /24-пространстве, что и a.b.c.d.
Some(v4) => IpAddr::V4(masked_v4(v4)),
None => IpAddr::V6(Ipv6Addr::from(masked_v6_segments(v6))),
},
}
}
fn masked_v4(ip: Ipv4Addr) -> Ipv4Addr {
let o = ip.octets();
Ipv4Addr::new(o[0], o[1], o[2], 0)
}
fn masked_v6_segments(ip: Ipv6Addr) -> [u8; 16] {
let mut segs = ip.octets();
segs[8..].fill(0);
segs
}
/// Эскалация вердиктов subnet-level детектора.
#[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<PrefixVerdict> {
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<std::sync::Mutex<crate::xdp::XdpFilter>>,
}
#[cfg(feature = "xdp")]
impl XdpPrefixStats {
#[must_use]
pub fn new(filter: std::sync::Arc<std::sync::Mutex<crate::xdp::XdpFilter>>) -> 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<Vec<(IpAddr, PrefixSnapshot)>> {
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::<PrefixStatsVal>(), 24);
}
}

View file

@ -1,9 +1,11 @@
use anyhow::{Context, Result, bail}; use anyhow::{Context, Result, bail};
use libbpf_rs::{MapCore, MapFlags, Object, ObjectBuilder, RingBuffer, RingBufferBuilder, Xdp, XdpFlags}; use libbpf_rs::{MapCore, MapFlags, Object, ObjectBuilder, RingBuffer, RingBufferBuilder, Xdp, XdpFlags};
use std::net::IpAddr;
use std::net::Ipv4Addr; use std::net::Ipv4Addr;
use std::os::unix::io::AsFd; use std::os::unix::io::AsFd;
use super::XdpStats; use super::XdpStats;
use crate::traffic::prefix::{PrefixKey, PrefixStatsVal};
pub struct XdpFilter { pub struct XdpFilter {
obj: Option<Object>, obj: Option<Object>,
@ -78,17 +80,28 @@ impl XdpFilter {
} }
pub fn ban_ip(&self, ip: Ipv4Addr, duration_secs: u64) -> Result<()> { 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 map = self.find_map("blacklist_map")?;
let mut key = [0u8; 8]; let mut key = [0u8; 8];
key[0] = 32; key[0] = prefix_len;
key[4..8].copy_from_slice(&ip.octets()); key[4..8].copy_from_slice(v4);
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_nanos() as u64;
map.update( map.update(
&key, &key,
&(now + duration_secs * 1_000_000_000).to_le_bytes(), &(unix_ns() + duration_secs * 1_000_000_000).to_le_bytes(),
MapFlags::ANY, MapFlags::ANY,
)?; )?;
Ok(()) Ok(())
@ -103,6 +116,28 @@ impl XdpFilter {
Ok(()) Ok(())
} }
/// Читает карту `prefix_stats` (агрегаты по префиксам /24 и /64).
///
/// # Errors
/// XDP не загружен, карта отсутствует в объекте или ошибка BPF-чтения.
pub fn read_prefix_stats(&self) -> Result<Vec<(PrefixKey, PrefixStatsVal)>> {
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<XdpStats> { pub fn get_stats(&self) -> Result<XdpStats> {
let map = self.find_map("stats_map")?; let map = self.find_map("stats_map")?;
let sum = |idx: u32| -> u64 { 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<RingBuffer<'static>> { fn build_ringbuf(obj: &Object) -> Result<RingBuffer<'static>> {
let map = obj let map = obj
.maps() .maps()

View file

@ -1,5 +1,7 @@
use anyhow::Result; use anyhow::{Ok, Result};
use std::net::Ipv4Addr; use std::net::{IpAddr, Ipv4Addr};
use crate::traffic::prefix::{PrefixKey, PrefixStatsVal};
pub struct XdpFilter; pub struct XdpFilter;
@ -17,9 +19,21 @@ impl XdpFilter {
pub fn ban_ip(&self, _ip: Ipv4Addr, _duration_secs: u64) -> Result<()> { pub fn ban_ip(&self, _ip: Ipv4Addr, _duration_secs: u64) -> Result<()> {
Ok(()) Ok(())
} }
pub fn ban_cidr(&self, _prefix: IpAddr, _prefix_len: u8, _duration_secs: u64) -> Result<()> {
Ok(())
}
pub fn unban_ip(&self, _ip: Ipv4Addr) -> Result<()> { pub fn unban_ip(&self, _ip: Ipv4Addr) -> Result<()> {
Ok(()) Ok(())
} }
/// Честная ошибка: чтение `prefix_stats` не подключено без feature `xdp`.
///
/// # Errors
/// Всегда — сборка без XDP не имеет доступа к картам ядра.
pub fn read_prefix_stats(&self) -> Result<Vec<(PrefixKey, PrefixStatsVal)>> {
anyhow::bail!("prefix_stats reading not wired yet: built without xdp feature")
}
pub fn get_stats(&self) -> Result<super::XdpStats> { pub fn get_stats(&self) -> Result<super::XdpStats> {
Ok(super::XdpStats::default()) Ok(super::XdpStats::default())
} }

215
tests/subnet_detection.rs Normal file
View file

@ -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<Vec<(IpAddr, PrefixSnapshot)>> {
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::<Ipv6Addr>().expect("valid"));
let v6b = IpAddr::V6("2001:db8:1:2::ffff".parse::<Ipv6Addr>().expect("valid"));
tracker.record_connection(v6a);
tracker.record_connection(v6b);
tracker.record_connection(v4([10, 0, 0, 1]));
let snaps: HashMap<IpAddr, PrefixSnapshot> = tracker.take_snapshot().into_iter().collect();
assert_eq!(snaps.len(), 2);
assert_eq!(
snaps[&IpAddr::V6("2001:db8:1:2::".parse::<Ipv6Addr>().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::<Ipv6Addr>().expect("valid"));
assert_eq!(
prefix_of(ip),
IpAddr::V6("2001:db8:1:2::".parse::<Ipv6Addr>().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::<PrefixKey>(), 24);
assert_eq!(size_of::<PrefixStatsVal>(), 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::<Ipv6Addr>().expect("valid")));
assert_eq!(key6.family, rampart::traffic::prefix::PREFIX_FAMILY_V6);
assert_eq!(
key6.to_prefix(),
Some(IpAddr::V6("2001:db8::".parse::<Ipv6Addr>().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);
}

76
xdp/core/prefix_stats.h Normal file
View file

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

View file

@ -26,6 +26,7 @@
#include "maps.h" #include "maps.h"
#include "config.h" #include "config.h"
#include "stats.h" #include "stats.h"
#include "prefix_stats.h"
#include "../hooks/hook_api.h" #include "../hooks/hook_api.h"
char __license[] SEC("license") = "GPL"; 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) if ((void *)data + sizeof(struct ethhdr) + ip_hdr_len + tcp_hdr_len > data_end)
return XDP_DROP; 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 // Compute payload pointers
__u8 *tcp_payload = (__u8 *)tcp + tcp_hdr_len; __u8 *tcp_payload = (__u8 *)tcp + tcp_hdr_len;
__u8 *tcp_payload_end = (__u8 *)data_end; __u8 *tcp_payload_end = (__u8 *)data_end;
@ -294,6 +299,9 @@ int rampart_universal_filter(struct xdp_md *ctx)
return XDP_DROP; 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 // Create conntrack entry for new connection
struct conntrack_entry ce = {}; struct conntrack_entry ce = {};
ce.state = STATE_SYN_RECEIVED; ce.state = STATE_SYN_RECEIVED;
@ -313,6 +321,9 @@ int rampart_universal_filter(struct xdp_md *ctx)
} }
// ── Connection tracking lookup ── // ── 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); struct conntrack_entry *conn = bpf_map_lookup_elem(&conntrack_map, &flow);
if (!conn) { if (!conn) {
// Unknown connection — drop // Unknown connection — drop