Hardened listener reload + config persistence + SYN Limit startup safety

This commit is contained in:
Alexey
2026-08-22 13:45:57 +03:00
parent 96d467a400
commit 189e10800a
28 changed files with 2749 additions and 1900 deletions
+78 -339
View File
@@ -1,11 +1,8 @@
use std::sync::{Arc, Mutex};
use std::sync::Mutex;
use tokio::sync::watch;
use tokio_util::sync::CancellationToken;
use tracing::warn;
use crate::config::ProxyConfig;
use crate::maestro::generation::RuntimeWatchState;
mod command;
mod iptables;
@@ -18,228 +15,91 @@ use self::model::{SynLimitNamespace, synlimit_namespace, synlimit_targets};
static ACTIVE_SYNLIMIT_NAMESPACE: Mutex<Option<SynLimitNamespace>> = Mutex::new(None);
/// Process-owned lifecycle handle for the SYN limiter reconciler.
pub(crate) struct SynlimitController {
shutdown: CancellationToken,
join: tokio::task::JoinHandle<()>,
}
impl SynlimitController {
/// Stops config observation after any in-flight reconcile completes.
pub(crate) async fn shutdown(self) {
self.shutdown.cancel();
let _ = self.join.await;
}
}
/// Spawns the process-scoped SYN limiter reconciler for active generations.
pub(crate) fn spawn_synlimit_controller(
runtime_watch_rx: watch::Receiver<Option<RuntimeWatchState>>,
) -> SynlimitController {
let shutdown = CancellationToken::new();
let join = tokio::spawn(watch_active_runtime_configs(
runtime_watch_rx,
shutdown.clone(),
|_generation_id, cfg| async move {
reconcile_synlimit_rules(&cfg).await;
},
));
SynlimitController { shutdown, join }
}
async fn watch_active_runtime_configs<F, Fut>(
mut runtime_watch_rx: watch::Receiver<Option<RuntimeWatchState>>,
shutdown: CancellationToken,
mut on_config: F,
) where
F: FnMut(u64, Arc<ProxyConfig>) -> Fut,
Fut: std::future::Future<Output = ()>,
{
let mut current = loop {
if let Some(state) = runtime_watch_rx.borrow().clone() {
break state;
}
tokio::select! {
biased;
_ = shutdown.cancelled() => return,
changed = runtime_watch_rx.changed() => {
if changed.is_err() {
return;
}
}
}
};
if shutdown.is_cancelled() {
return;
}
let initial_config = current.config_rx.borrow().clone();
on_config(current.generation_id, initial_config).await;
loop {
tokio::select! {
biased;
_ = shutdown.cancelled() => break,
changed = runtime_watch_rx.changed() => {
if changed.is_err() {
break;
}
let Some(next) = runtime_watch_rx.borrow().clone() else {
continue;
};
if next.generation_id != current.generation_id {
current = next;
let config = current.config_rx.borrow().clone();
on_config(current.generation_id, config).await;
}
}
changed = current.config_rx.changed() => {
if changed.is_err() {
let Some(next) = wait_for_new_runtime(
&mut runtime_watch_rx,
current.generation_id,
&shutdown,
).await else {
break;
};
current = next;
let config = current.config_rx.borrow().clone();
on_config(current.generation_id, config).await;
continue;
}
let active_generation_id = runtime_watch_rx
.borrow()
.as_ref()
.map(|state| state.generation_id);
if active_generation_id == Some(current.generation_id) {
let cfg = current.config_rx.borrow_and_update().clone();
on_config(current.generation_id, cfg).await;
}
}
}
}
}
async fn wait_for_new_runtime(
runtime_watch_rx: &mut watch::Receiver<Option<RuntimeWatchState>>,
previous_generation_id: u64,
shutdown: &CancellationToken,
) -> Option<RuntimeWatchState> {
loop {
if let Some(state) = runtime_watch_rx.borrow().clone()
&& state.generation_id != previous_generation_id
{
return Some(state);
}
tokio::select! {
biased;
_ = shutdown.cancelled() => return None,
changed = runtime_watch_rx.changed() => {
if changed.is_err() {
return None;
}
}
}
}
}
pub(crate) async fn reconcile_synlimit_rules(cfg: &ProxyConfig) {
/// Installs the complete startup SYN-limiter ruleset before accept loops start.
pub(crate) async fn reconcile_synlimit_rules(cfg: &ProxyConfig) -> Result<(), String> {
let targets = synlimit_targets(cfg);
let namespace = synlimit_namespace(&targets);
if let Some(previous_namespace) = set_active_synlimit_namespace(namespace.clone()) {
match clear_synlimit_rules_for_namespace(&previous_namespace).await {
Ok(true) => {
warn!("Removed previous SYN limiter namespace before reconcile");
}
Ok(false) => {}
Err(error) => {
warn!(error = %error, "Failed to clear previous SYN limiter namespace before reconcile");
}
}
}
if targets.is_empty() {
return;
return Ok(());
}
let Some(namespace) = namespace else {
return;
};
if !has_firewall_privileges() {
warn!("SYN limiter configured but firewall privileges are not available; rules not applied");
return;
return Err(
"SYN limiter requires root or CAP_NET_ADMIN for startup and shutdown".to_string(),
);
}
let namespace = synlimit_namespace(&targets)
.ok_or_else(|| "SYN limiter namespace could not be derived".to_string())?;
if clear_synlimit_rules_for_namespace(&namespace).await? {
warn!("Removed stale SYN limiter rules left by a previous run before startup");
}
match clear_synlimit_rules_for_namespace(&namespace).await {
Ok(true) => {
warn!("Removed stale SYN limiter rules left by a previous run before reconcile");
let apply_result = async {
if targets.has_iptables_targets() {
iptables::apply_synlimit_rules(&targets, &namespace).await?;
}
Ok(false) => {}
Err(error) => {
warn!(error = %error, "Failed to clear stale SYN limiter rules before reconcile");
if targets.has_nft_targets() {
nftables::apply_synlimit_rules(&targets, &namespace).await?;
}
if targets.has_pf_targets() {
pf::apply_synlimit_rules(&targets, &namespace).await?;
}
Ok::<(), String>(())
}
.await;
if let Err(apply_error) = apply_result {
return match clear_synlimit_rules_for_namespace(&namespace).await {
Ok(_) => Err(apply_error),
Err(cleanup_error) => Err(format!(
"{apply_error}; candidate cleanup failed: {cleanup_error}"
)),
};
}
if targets.has_iptables_targets() {
if let Err(error) = iptables::apply_synlimit_rules(&targets, &namespace).await {
warn!(error = %error, "Failed to apply iptables SYN limiter rules");
}
}
if targets.has_nft_targets() {
if let Err(error) = nftables::apply_synlimit_rules(&targets, &namespace).await {
warn!(error = %error, "Failed to apply nftables SYN limiter rules");
}
}
if targets.has_pf_targets() {
if let Err(error) = pf::apply_synlimit_rules(&targets, &namespace).await {
warn!(error = %error, "Failed to apply PF SYN limiter rules");
}
if let Err(error) = set_active_synlimit_namespace(namespace.clone()) {
return match clear_synlimit_rules_for_namespace(&namespace).await {
Ok(_) => Err(error),
Err(cleanup_error) => Err(format!(
"{error}; candidate cleanup failed: {cleanup_error}"
)),
};
}
Ok(())
}
/// Removes the ruleset installed by the current process, if any.
pub(crate) async fn clear_synlimit_rules_all_backends() -> Result<bool, String> {
let Some(namespace) = take_active_synlimit_namespace() else {
let Some(namespace) = active_synlimit_namespace()? else {
return Ok(false);
};
clear_synlimit_rules_for_namespace(&namespace).await
let removed = clear_synlimit_rules_for_namespace(&namespace).await?;
clear_active_synlimit_namespace(&namespace)?;
Ok(removed)
}
async fn clear_synlimit_rules_for_namespace(namespace: &SynLimitNamespace) -> Result<bool, String> {
if !has_firewall_privileges() {
return Ok(false);
return Err(
"SYN limiter cleanup requires root or CAP_NET_ADMIN privileges".to_string(),
);
}
let mut errors = Vec::new();
let mut removed = false;
match nftables::clear_rules_all_families(namespace).await {
Ok(value) => {
removed |= value;
}
Err(error) => {
errors.push(error);
}
Ok(value) => removed |= value,
Err(error) => errors.push(error),
}
match iptables::clear_rules_for_binary("iptables", namespace).await {
Ok(value) => {
removed |= value;
}
Err(error) => {
errors.push(error);
}
Ok(value) => removed |= value,
Err(error) => errors.push(error),
}
match iptables::clear_rules_for_binary("ip6tables", namespace).await {
Ok(value) => {
removed |= value;
}
Err(error) => {
errors.push(error);
}
Ok(value) => removed |= value,
Err(error) => errors.push(error),
}
match pf::clear_rules(namespace).await {
Ok(value) => {
removed |= value;
}
Err(error) => {
errors.push(error);
}
Ok(value) => removed |= value,
Err(error) => errors.push(error),
}
if errors.is_empty() {
@@ -249,161 +109,40 @@ async fn clear_synlimit_rules_for_namespace(namespace: &SynLimitNamespace) -> Re
}
}
fn set_active_synlimit_namespace(next: Option<SynLimitNamespace>) -> Option<SynLimitNamespace> {
fn set_active_synlimit_namespace(next: SynLimitNamespace) -> Result<(), String> {
match ACTIVE_SYNLIMIT_NAMESPACE.lock() {
Ok(mut active) => {
if *active == next {
None
} else {
std::mem::replace(&mut *active, next)
if active.is_some() {
return Err("SYN limiter namespace is already active".to_string());
}
*active = Some(next);
Ok(())
}
Err(error) => {
warn!(error = %error, "Failed to update active SYN limiter namespace");
None
}
Err(error) => Err(format!(
"failed to update active SYN limiter namespace: {error}"
)),
}
}
fn take_active_synlimit_namespace() -> Option<SynLimitNamespace> {
fn active_synlimit_namespace() -> Result<Option<SynLimitNamespace>, String> {
match ACTIVE_SYNLIMIT_NAMESPACE.lock() {
Ok(mut active) => active.take(),
Err(error) => {
warn!(error = %error, "Failed to read active SYN limiter namespace");
None
}
Ok(active) => Ok(active.clone()),
Err(error) => Err(format!(
"failed to read active SYN limiter namespace: {error}"
)),
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
use tokio::sync::{Notify, mpsc};
fn runtime_state(
generation_id: u64,
max_connections: u32,
) -> (
RuntimeWatchState,
watch::Sender<Arc<ProxyConfig>>,
watch::Sender<bool>,
) {
let mut config = ProxyConfig::default();
config.server.max_connections = max_connections;
let (config_tx, config_rx) = watch::channel(Arc::new(config));
let (admission_tx, admission_rx) = watch::channel(true);
(
RuntimeWatchState {
generation_id,
config_rx,
admission_rx,
},
config_tx,
admission_tx,
)
}
#[tokio::test]
async fn config_watcher_ignores_retired_generation_updates() {
let (initial, initial_config_tx, _initial_admission_tx) = runtime_state(1, 10);
let (runtime_tx, runtime_rx) = watch::channel(Some(initial));
let (observed_tx, mut observed_rx) = mpsc::unbounded_channel();
let watcher = tokio::spawn(watch_active_runtime_configs(
runtime_rx,
CancellationToken::new(),
move |generation_id, cfg| {
let observed_tx = observed_tx.clone();
async move {
let _ = observed_tx.send((generation_id, cfg.server.max_connections));
}
},
));
assert_eq!(observed_rx.recv().await, Some((1, 10)));
let (next, next_config_tx, _next_admission_tx) = runtime_state(2, 20);
runtime_tx.send_replace(Some(next));
assert_eq!(observed_rx.recv().await, Some((2, 20)));
let mut stale = ProxyConfig::default();
stale.server.max_connections = 30;
initial_config_tx.send_replace(Arc::new(stale));
assert!(
tokio::time::timeout(Duration::from_millis(50), observed_rx.recv())
.await
.is_err()
);
let mut active = ProxyConfig::default();
active.server.max_connections = 40;
next_config_tx.send_replace(Arc::new(active));
assert_eq!(observed_rx.recv().await, Some((2, 40)));
watcher.abort();
}
#[tokio::test]
async fn shutdown_waits_for_inflight_reconcile_and_stops_future_updates() {
let (initial, config_tx, _admission_tx) = runtime_state(1, 10);
let (_runtime_tx, runtime_rx) = watch::channel(Some(initial));
let shutdown = CancellationToken::new();
let started = Arc::new(Notify::new());
let release = Arc::new(Notify::new());
let calls = Arc::new(AtomicUsize::new(0));
let started_callback = started.clone();
let release_callback = release.clone();
let calls_callback = calls.clone();
let watcher_shutdown = shutdown.clone();
let watcher = tokio::spawn(watch_active_runtime_configs(
runtime_rx,
watcher_shutdown,
move |_generation_id, _cfg| {
let started = started_callback.clone();
let release = release_callback.clone();
let calls = calls_callback.clone();
async move {
calls.fetch_add(1, Ordering::AcqRel);
started.notify_one();
release.notified().await;
}
},
));
started.notified().await;
shutdown.cancel();
tokio::task::yield_now().await;
assert!(!watcher.is_finished());
release.notify_one();
tokio::time::timeout(Duration::from_secs(1), watcher)
.await
.unwrap()
.unwrap();
drop(_runtime_tx);
let mut updated = ProxyConfig::default();
updated.server.max_connections = 20;
assert!(config_tx.send(Arc::new(updated)).is_err());
assert_eq!(calls.load(Ordering::Acquire), 1);
}
#[tokio::test]
async fn shutdown_before_start_skips_initial_reconcile() {
let (initial, _config_tx, _admission_tx) = runtime_state(1, 10);
let (_runtime_tx, runtime_rx) = watch::channel(Some(initial));
let shutdown = CancellationToken::new();
shutdown.cancel();
let calls = Arc::new(AtomicUsize::new(0));
let calls_callback = calls.clone();
watch_active_runtime_configs(runtime_rx, shutdown, move |_generation_id, _cfg| {
let calls = calls_callback.clone();
async move {
calls.fetch_add(1, Ordering::AcqRel);
fn clear_active_synlimit_namespace(expected: &SynLimitNamespace) -> Result<(), String> {
match ACTIVE_SYNLIMIT_NAMESPACE.lock() {
Ok(mut active) => {
if active.as_ref() == Some(expected) {
*active = None;
}
})
.await;
assert_eq!(calls.load(Ordering::Acquire), 0);
Ok(())
}
Err(error) => Err(format!(
"failed to update active SYN limiter namespace: {error}"
)),
}
}
+25 -22
View File
@@ -5,6 +5,21 @@ use super::model::{SynLimitNamespace, SynLimitRule, SynLimitTargets};
const PF_ANCHOR_ROOT: &str = "telemt_synlimit";
#[derive(Clone, Copy)]
enum PfFamily {
Inet,
Inet6,
}
impl PfFamily {
fn as_str(self) -> &'static str {
match self {
Self::Inet => "inet",
Self::Inet6 => "inet6",
}
}
}
pub(super) async fn apply_synlimit_rules(
targets: &SynLimitTargets,
namespace: &SynLimitNamespace,
@@ -31,26 +46,23 @@ fn is_pf_anchor_hook_line(line: &str) -> bool {
fn pf_synlimit_script(targets: &SynLimitTargets) -> String {
let mut script = String::new();
for target in &targets.pf_v4 {
push_pf_rules(&mut script, target);
push_pf_rules(&mut script, PfFamily::Inet, target);
}
for target in &targets.pf_v6 {
push_pf_rules(&mut script, target);
push_pf_rules(&mut script, PfFamily::Inet6, target);
}
script
}
fn push_pf_rules(script: &mut String, target: &SynLimitRule) {
fn push_pf_rules(script: &mut String, family: PfFamily, target: &SynLimitRule) {
let destination = pf_destination(target.ip);
script.push_str(&format!(
"pass in quick proto tcp from any to {destination} port {port} flags S/SA keep state (max-src-conn-rate {rate}/{seconds})\n",
"pass in quick {family} proto tcp from any to {destination} port {port} flags S/SA keep state (max-src-conn-rate {rate}/{seconds})\n",
family = family.as_str(),
port = target.port,
rate = target.generic_hitcount,
seconds = target.generic_seconds,
));
script.push_str(&format!(
"block return-rst in quick proto tcp from any to {destination} port {port}\n",
port = target.port,
));
}
fn pf_destination(ip: Option<IpAddr>) -> String {
@@ -84,24 +96,15 @@ mod tests {
use crate::synlimit_control::model::test_rule;
#[test]
fn pf_script_uses_rate_limited_pass_before_reject() {
fn pf_script_uses_native_rate_limited_pass() {
let mut targets = SynLimitTargets::default();
targets.pf_v4 = vec![test_rule(Some(IpAddr::V4(Ipv4Addr::new(203, 0, 113, 7))), 443)];
let script = pf_synlimit_script(&targets);
assert!(script.contains(
"pass in quick proto tcp from any to 203.0.113.7 port 443 flags S/SA keep state (max-src-conn-rate 48/60)"
"pass in quick inet proto tcp from any to 203.0.113.7 port 443 flags S/SA keep state (max-src-conn-rate 48/60)"
));
assert!(script.contains(
"block return-rst in quick proto tcp from any to 203.0.113.7 port 443"
));
let pass_idx = script
.find("pass in quick proto tcp from any to 203.0.113.7 port 443")
.expect("rate-limited pass rule must be rendered");
let block_idx = script
.find("block return-rst in quick proto tcp from any to 203.0.113.7 port 443")
.expect("reject fallback rule must be rendered");
assert!(pass_idx < block_idx);
assert!(!script.contains("return-rst"));
}
#[test]
@@ -111,8 +114,8 @@ mod tests {
targets.pf_v6 = vec![test_rule(Some(IpAddr::V6(Ipv6Addr::LOCALHOST)), 8443)];
let script = pf_synlimit_script(&targets);
assert!(script.contains("to any port 443"));
assert!(script.contains("to ::1 port 8443"));
assert!(script.contains("pass in quick inet proto tcp from any to any port 443"));
assert!(script.contains("pass in quick inet6 proto tcp from any to ::1 port 8443"));
}
#[test]