Files
fips/src/gateway/pool.rs
T
Johnathan Corgan 43b6512503 fix(gateway): harden the virtual-IP pool and ship it disabled on OpenWrt
Any host able to query the LAN resolver could drain the gateway's
65,535-address virtual-IP pool one `.fips` name at a time, and each allocation
rebuilt the whole nftables table in a way that could leave the host with no
NAT at all. Four changes, each independently useful, close that off.

Do not allocate for query types the gateway never answers with an address.
handle_query minted a virtual IP for every query type and only then looked at
what the client asked, answering an A or HTTPS query with NODATA after
creating a mapping for it. The query type is now decided before the pool is
touched, and only AAAA and ANY allocate. The refresh an existing mapping used
to get from any query type is kept: it came from the reuse path in allocate,
so a new pool method does that refresh alone and never creates anything, and
the reuse path calls it.

Rebuild the NAT table in one netlink transaction. rebuild() deleted the
fips_gateway table in a batch of its own and discarded the result, then sent a
second batch recreating the table, the chains, the fips0 masquerade and two
rules per mapping. Between those sends the host had no NAT table, and a
recreate the kernel refused left the table deleted, turning one failed mapping
change into a total loss of forwarding until some later rebuild happened to
succeed. The delete and the recreate now share one batch. A leading table add
makes the delete legal on the first run, since rustables sends it with
NLM_F_CREATE and no NLM_F_EXCL and the crate offers no flush. Deciding what to
send is now separate from sending it, which is the seam the new unit tests
use: they assert one batch, the add-delete-add prefix, and that every chain
and rule follows the recreate, without a netlink socket or privileges.

Read conntrack once per tick, off the runtime thread, and match by address.
The session count searched each /proc/net/nf_conntrack line for `dst=`
followed by the virtual IP's compressed Display form, while the kernel prints
every tuple with `%pI6`, the full uncompressed form. That string cannot occur
in that field, so the count was zero for every mapping on every kernel that
has the file: nothing pinned an in-use mapping and one whose client did not
re-query DNS was reclaimed about two minutes after its last DNS reference with
traffic still flowing. Each `dst=` is now parsed and compared as an address.
The read was also per mapping, under the pool lock, on the runtime thread that
serves DNS; the tick now takes one snapshot in a blocking task before taking
the lock. An unreadable source was silent, because read_to_string's error
became zero through unwrap_or(0). Zero stays, since treating it as in-use
would pin every mapping forever on a kernel without
CONFIG_NF_CONNTRACK_PROCFS, but it is now reported at warn on the first
failure and on each change of outcome, and at debug on a repeat.

Ship the OpenWrt gateway disabled, and keep its state across upgrades. The
generated postinst enabled and started fips-gateway on every install, against
the init script's own header, the package README and the deployment tutorial,
which all say the service ships disabled. A fresh install now leaves it alone.
Upgrades are the awkward case: opkg runs the outgoing package's prerm first,
and every released prerm disabled the gateway on its way out without recording
whether it had been enabled. The new prerm stops the services on an upgrade
but no longer disables them, and leaves a marker the incoming postinst reads.
With the marker, enablement survived and the gateway starts only if it was
enabled; without it, the outgoing package was a released one whose prerm
destroyed that state, so the gateway is re-enabled rather than letting an
upgrade turn off a working deployment. That re-enables a hand-disabled gateway
once, which the CHANGELOG says. start_service now reads gateway.enabled from
fips.yaml before touching anything, since starting a gateway the config
disables used to take dnsmasq's `.fips` forwarding away from the daemon and
hand it to a port whose daemon exits immediately. The four maintainer-script
bodies move out of heredocs in the two build scripts into
packaging/openwrt-ipk/scripts/, so the .ipk and the .apk install the same
bodies and a test can run what ships.

Coverage recorded rather than closed. The conntrack parser's first test builds
its line from the kernel's own format string rather than a capture, because
this host is built without CONFIG_NF_CONNTRACK_PROCFS and has no
/proc/net/nf_conntrack, so the lab exercises only the unreadable path. Kernel
acceptance of delete-then-recreate inside one transaction is not asserted by a
unit test; the gateway suite is what proves it, since the manager rebuilds at
startup and the daemon exits if that fails. The OpenWrt scenarios run the
shipped script bodies under ash in a busybox container against stubbed init
scripts, and assert their behaviour given opkg's call order, arguments and
PKG_UPGRADE as read from opkg-lede's sources, not under a real opkg upgrade on
a router image.

Admission limits on the pool are deliberately not included here: they need a
measurement run before their constants can be chosen.
2026-09-17 20:46:12 +00:00

912 lines
33 KiB
Rust

//! Virtual IP pool manager.
//!
//! Manages allocation, TTL, and reclamation of virtual IPv6 addresses
//! from a configured CIDR range. Tracks mapping state and integrates
//! with conntrack to determine active sessions.
use crate::NodeAddr;
use std::collections::{HashMap, HashSet, VecDeque};
use std::net::Ipv6Addr;
use std::time::Instant;
use tracing::{debug, info};
/// Errors from pool operations.
#[derive(Debug, thiserror::Error)]
pub enum PoolError {
#[error("invalid CIDR: {0}")]
InvalidCidr(String),
#[error("pool exhausted ({0} addresses in use)")]
Exhausted(usize),
#[error("prefix length must be between 1 and 128")]
InvalidPrefix,
}
/// State of a virtual IP mapping.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MappingState {
/// Allocated via DNS query, no NAT sessions yet.
Allocated,
/// Active NAT sessions exist.
Active,
/// TTL expired but sessions remain.
Draining,
}
/// A single virtual IP ↔ FIPS mesh address mapping.
#[derive(Debug, Clone)]
pub struct VirtualIpMapping {
/// The FIPS node address this mapping is for.
pub node_addr: NodeAddr,
/// The virtual IP allocated from the pool.
pub virtual_ip: Ipv6Addr,
/// The FIPS mesh address (fd00::/8).
pub mesh_addr: Ipv6Addr,
/// The DNS name that was queried (e.g. "npub1abc...xyz.fips").
pub dns_name: String,
/// Current state.
pub state: MappingState,
/// When this mapping was created.
pub created: Instant,
/// When this mapping was last referenced (DNS query or session).
pub last_referenced: Instant,
/// When draining started (for grace period tracking).
pub drain_start: Option<Instant>,
/// Number of active conntrack sessions.
pub session_count: u32,
}
/// Events emitted by the pool on state transitions.
#[derive(Debug)]
pub enum PoolEvent {
/// A new mapping was allocated — NAT rules should be created.
MappingCreated {
virtual_ip: Ipv6Addr,
mesh_addr: Ipv6Addr,
},
/// A mapping was reclaimed — NAT rules should be removed.
MappingRemoved {
virtual_ip: Ipv6Addr,
mesh_addr: Ipv6Addr,
},
}
/// Pool utilization summary.
#[derive(Debug, Clone)]
pub struct PoolStatus {
pub total: usize,
pub allocated: usize,
pub active: usize,
pub draining: usize,
pub free: usize,
}
/// Summary of a single mapping for display.
#[derive(Debug, Clone)]
pub struct MappingInfo {
pub virtual_ip: Ipv6Addr,
pub mesh_addr: Ipv6Addr,
pub node_addr: NodeAddr,
pub dns_name: String,
pub state: MappingState,
pub session_count: u32,
pub age_secs: u64,
pub last_ref_secs: u64,
}
/// Path the conntrack table is read from.
const CONNTRACK_PROC_PATH: &str = "/proc/net/nf_conntrack";
/// Active conntrack sessions counted by destination address.
///
/// Taken once per tick, so the pool does a map lookup per mapping instead of
/// reading and scanning the whole conntrack table per mapping under its lock.
#[derive(Debug, Clone, Default)]
pub struct ConntrackSnapshot {
sessions: HashMap<Ipv6Addr, u32>,
}
impl ConntrackSnapshot {
/// Build a snapshot from counts already keyed by destination address.
pub fn from_counts(sessions: HashMap<Ipv6Addr, u32>) -> Self {
Self { sessions }
}
/// Sessions whose destination is `virtual_ip`, or zero if there are none.
pub fn sessions_for(&self, virtual_ip: Ipv6Addr) -> u32 {
self.sessions.get(&virtual_ip).copied().unwrap_or(0)
}
/// Number of distinct destination addresses the snapshot saw.
pub fn len(&self) -> usize {
self.sessions.len()
}
/// Whether the snapshot saw no sessions at all.
pub fn is_empty(&self) -> bool {
self.sessions.is_empty()
}
}
/// Trait for taking a conntrack session snapshot.
pub trait ConntrackQuerier: Send + Sync {
/// Read the conntrack table once and count sessions by destination.
fn snapshot(&self) -> Result<ConntrackSnapshot, std::io::Error>;
}
/// Conntrack querier that parses /proc/net/nf_conntrack.
pub struct ProcConntrack;
impl ConntrackQuerier for ProcConntrack {
fn snapshot(&self) -> Result<ConntrackSnapshot, std::io::Error> {
let content = std::fs::read_to_string(CONNTRACK_PROC_PATH)?;
Ok(ConntrackSnapshot::from_counts(parse_conntrack(&content)))
}
}
/// Count conntrack lines by the destination addresses they name.
///
/// Every `dst=` value is parsed as an address and compared as an address. The
/// kernel prints tuples as `src=%pI6 dst=%pI6`, the full uncompressed form with
/// leading zeros, so a session to `fd01::1` is written
/// `dst=fd01:0000:0000:0000:0000:0000:0000:0001`; the previous code searched
/// each line for the address's compressed `Display` form, which cannot occur in
/// a fixed-width field, so it counted nothing on any kernel.
///
/// A conntrack line carries the original and the reply tuple, each with its own
/// `dst=`, and the line is counted once per distinct address among them. That
/// keeps the meaning the count had before, which was "this line mentions the
/// address". A value that does not parse as an IPv6 address is skipped, which
/// is how IPv4 lines and any future field are ignored.
fn parse_conntrack(content: &str) -> HashMap<Ipv6Addr, u32> {
let mut counts: HashMap<Ipv6Addr, u32> = HashMap::new();
let mut seen: HashSet<Ipv6Addr> = HashSet::new();
for line in content.lines() {
seen.clear();
for token in line.split_whitespace() {
let Some(value) = token.strip_prefix("dst=") else {
continue;
};
let Ok(addr) = value.parse::<Ipv6Addr>() else {
continue;
};
seen.insert(addr);
}
for addr in &seen {
*counts.entry(*addr).or_insert(0) += 1;
}
}
counts
}
/// Whether a conntrack read outcome is new or a repeat of the last one.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ReadReport {
/// The outcome differs from the previous read, or is the first.
Changed,
/// The same outcome as the previous read.
Repeated,
}
/// Remembers the last conntrack read outcome.
///
/// A kernel built without `CONFIG_NF_CONNTRACK_PROCFS` has no
/// `/proc/net/nf_conntrack` at all, so every read fails the same way and a
/// per-tick warning would repeat for the life of the process. Warning on a
/// change of outcome still separates "the source is unreadable" from "there
/// are no sessions", which the pool could not distinguish before, without
/// filling the log.
#[derive(Debug, Default)]
pub struct ConntrackReadLog {
last: Option<Option<std::io::ErrorKind>>,
}
impl ConntrackReadLog {
/// Record a read outcome and say whether it is new.
///
/// `None` is a successful read; `Some(kind)` is a failure of that kind.
pub fn observe(&mut self, outcome: Option<std::io::ErrorKind>) -> ReadReport {
let report = if self.last == Some(outcome) {
ReadReport::Repeated
} else {
ReadReport::Changed
};
self.last = Some(outcome);
report
}
}
/// Virtual IP pool manager.
pub struct VirtualIpPool {
/// Available addresses (free pool).
available: VecDeque<Ipv6Addr>,
/// Active mappings keyed by NodeAddr.
mappings: HashMap<NodeAddr, VirtualIpMapping>,
/// Reverse map: virtual IP → NodeAddr.
reverse: HashMap<Ipv6Addr, NodeAddr>,
/// DNS TTL / mapping TTL in seconds.
ttl_secs: u64,
/// Grace period after last session before reclamation.
grace_secs: u64,
/// Total pool size.
total: usize,
}
impl VirtualIpPool {
/// Create a new pool from a CIDR string (e.g., `fd01::/112`).
pub fn new(cidr: &str, ttl_secs: u64, grace_secs: u64) -> Result<Self, PoolError> {
let (base, prefix_len) = parse_ipv6_cidr(cidr)?;
if prefix_len == 0 || prefix_len > 128 {
return Err(PoolError::InvalidPrefix);
}
let mut available = VecDeque::new();
let host_bits = 128 - prefix_len;
// Cap at 2^16 addresses to avoid massive allocations
let max_addrs: u128 = if host_bits > 16 {
1u128 << 16
} else {
1u128 << host_bits
};
let base_int = u128::from(base);
// Skip address 0 (network equivalent)
for i in 1..max_addrs {
available.push_back(Ipv6Addr::from(base_int + i));
}
let total = available.len();
info!(cidr = %cidr, addresses = total, "Virtual IP pool initialized");
Ok(Self {
available,
mappings: HashMap::new(),
reverse: HashMap::new(),
ttl_secs,
grace_secs,
total,
})
}
/// Refresh an existing mapping's TTL clock, never creating one.
///
/// Returns whether a mapping for `node_addr` existed. A query the gateway
/// answers without an address still says the client is using the name, so
/// it must keep the mapping alive without minting one.
pub fn refresh_if_present(&mut self, node_addr: NodeAddr) -> bool {
match self.mappings.get_mut(&node_addr) {
Some(mapping) => {
mapping.last_referenced = Instant::now();
true
}
None => false,
}
}
/// Allocate a virtual IP for the given node. Idempotent: returns
/// existing mapping if one exists.
pub fn allocate(
&mut self,
node_addr: NodeAddr,
mesh_addr: Ipv6Addr,
dns_name: &str,
) -> Result<(Ipv6Addr, bool), PoolError> {
// Idempotent: return existing mapping, refreshed.
if self.refresh_if_present(node_addr)
&& let Some(mapping) = self.mappings.get(&node_addr)
{
return Ok((mapping.virtual_ip, false));
}
let virtual_ip = self
.available
.pop_front()
.ok_or(PoolError::Exhausted(self.mappings.len()))?;
let now = Instant::now();
let mapping = VirtualIpMapping {
node_addr,
virtual_ip,
mesh_addr,
dns_name: dns_name.to_string(),
state: MappingState::Allocated,
created: now,
last_referenced: now,
drain_start: None,
session_count: 0,
};
self.mappings.insert(node_addr, mapping);
self.reverse.insert(virtual_ip, node_addr);
info!(
virtual_ip = %virtual_ip,
mesh_addr = %mesh_addr,
dns_name = %dns_name,
"Allocated virtual IP"
);
Ok((virtual_ip, true))
}
/// Periodic tick — drives state transitions. Returns events for
/// the NAT and network modules.
pub fn tick(&mut self, now: Instant, conntrack: &ConntrackSnapshot) -> Vec<PoolEvent> {
let mut events = Vec::new();
let mut to_free = Vec::new();
let ttl = std::time::Duration::from_secs(self.ttl_secs);
let grace = std::time::Duration::from_secs(self.grace_secs);
for (node_addr, mapping) in &mut self.mappings {
// One map lookup: the conntrack table was read once, before the
// pool lock was taken.
let sessions = conntrack.sessions_for(mapping.virtual_ip);
mapping.session_count = sessions;
// Live data-plane traffic pins the mapping: refresh the TTL
// clock whenever conntrack reports active sessions, so an
// in-use mapping never ages out from under the client.
if sessions > 0 {
mapping.last_referenced = now;
}
match mapping.state {
MappingState::Allocated => {
if sessions > 0 {
mapping.state = MappingState::Active;
debug!(
virtual_ip = %mapping.virtual_ip,
sessions,
"Mapping activated"
);
} else if now.duration_since(mapping.last_referenced) > ttl {
// TTL expired — enter draining with grace period so
// the mapping survives browser DNS cache, even if no
// conntrack sessions were observed (short HTTP requests
// may complete between ticks).
mapping.state = MappingState::Draining;
mapping.drain_start = Some(now);
debug!(
virtual_ip = %mapping.virtual_ip,
"Allocated mapping TTL expired, draining"
);
}
}
MappingState::Active => {
// The traffic refresh above keeps last_referenced == now
// while sessions > 0, so the TTL can only trip once the
// mapping is idle (no conntrack sessions). An actively used
// mapping never drains; an idle one enters the grace period.
if now.duration_since(mapping.last_referenced) > ttl {
mapping.state = MappingState::Draining;
mapping.drain_start = Some(now);
}
}
MappingState::Draining => {
if sessions > 0 {
// Traffic resumed before reclamation: recover to
// Active and clear drain_start so the next drain
// gets a fresh grace window rather than reusing a
// stale one.
mapping.state = MappingState::Active;
mapping.drain_start = None;
debug!(
virtual_ip = %mapping.virtual_ip,
sessions,
"Draining mapping recovered to active (traffic resumed)"
);
} else if let Some(drain_start) = mapping.drain_start
&& now.duration_since(drain_start) > grace
{
to_free.push(*node_addr);
}
}
}
}
// Free expired mappings
for node_addr in to_free {
if let Some(mapping) = self.mappings.remove(&node_addr) {
self.reverse.remove(&mapping.virtual_ip);
self.available.push_back(mapping.virtual_ip);
info!(
virtual_ip = %mapping.virtual_ip,
mesh_addr = %mapping.mesh_addr,
"Reclaimed virtual IP"
);
events.push(PoolEvent::MappingRemoved {
virtual_ip: mapping.virtual_ip,
mesh_addr: mapping.mesh_addr,
});
}
}
events
}
/// Pool utilization summary.
pub fn status(&self) -> PoolStatus {
let mut allocated = 0;
let mut active = 0;
let mut draining = 0;
for mapping in self.mappings.values() {
match mapping.state {
MappingState::Allocated => allocated += 1,
MappingState::Active => active += 1,
MappingState::Draining => draining += 1,
}
}
PoolStatus {
total: self.total,
allocated,
active,
draining,
free: self.available.len(),
}
}
/// Summary of all active mappings.
pub fn mapping_info(&self, now: Instant) -> Vec<MappingInfo> {
self.mappings
.values()
.map(|m| MappingInfo {
virtual_ip: m.virtual_ip,
mesh_addr: m.mesh_addr,
node_addr: m.node_addr,
dns_name: m.dns_name.clone(),
state: m.state,
session_count: m.session_count,
age_secs: now.duration_since(m.created).as_secs(),
last_ref_secs: now.duration_since(m.last_referenced).as_secs(),
})
.collect()
}
/// Look up which node a virtual IP maps to.
pub fn lookup_virtual_ip(&self, virtual_ip: &Ipv6Addr) -> Option<&VirtualIpMapping> {
self.reverse
.get(virtual_ip)
.and_then(|addr| self.mappings.get(addr))
}
}
/// Parse an IPv6 CIDR string into base address and prefix length.
fn parse_ipv6_cidr(cidr: &str) -> Result<(Ipv6Addr, u32), PoolError> {
let parts: Vec<&str> = cidr.split('/').collect();
if parts.len() != 2 {
return Err(PoolError::InvalidCidr(cidr.to_string()));
}
let addr: Ipv6Addr = parts[0]
.parse()
.map_err(|_| PoolError::InvalidCidr(cidr.to_string()))?;
let prefix: u32 = parts[1]
.parse()
.map_err(|_| PoolError::InvalidCidr(cidr.to_string()))?;
Ok((addr, prefix))
}
#[cfg(test)]
mod tests {
use super::*;
/// Session counts a test sets directly, handed to `tick` as the snapshot
/// the tick task would have read from conntrack.
#[derive(Default)]
struct Sessions {
counts: HashMap<Ipv6Addr, u32>,
}
impl Sessions {
fn new() -> Self {
Self::default()
}
fn set(&mut self, addr: Ipv6Addr, count: u32) {
self.counts.insert(addr, count);
}
fn snapshot(&self) -> ConntrackSnapshot {
ConntrackSnapshot::from_counts(self.counts.clone())
}
}
fn make_node_addr(byte: u8) -> NodeAddr {
let mut bytes = [0u8; 16];
bytes[0] = byte;
NodeAddr::from_bytes(bytes)
}
fn make_mesh_addr(byte: u8) -> Ipv6Addr {
let mut bytes = [0u8; 16];
bytes[0] = 0xfd;
bytes[15] = byte;
Ipv6Addr::from(bytes)
}
#[test]
fn test_parse_cidr() {
let (addr, prefix) = parse_ipv6_cidr("fd01::/112").unwrap();
assert_eq!(addr, "fd01::".parse::<Ipv6Addr>().unwrap());
assert_eq!(prefix, 112);
}
#[test]
fn test_parse_cidr_invalid() {
assert!(parse_ipv6_cidr("not-a-cidr").is_err());
assert!(parse_ipv6_cidr("fd01::").is_err());
assert!(parse_ipv6_cidr("fd01::/abc").is_err());
}
#[test]
fn test_pool_creation() {
let pool = VirtualIpPool::new("fd01::/120", 60, 60).unwrap();
// /120 = 8 host bits = 256 addresses, minus 1 (network) = 255
assert_eq!(pool.total, 255);
assert_eq!(pool.available.len(), 255);
}
#[test]
fn test_pool_allocation() {
let mut pool = VirtualIpPool::new("fd01::/120", 60, 60).unwrap();
let node = make_node_addr(1);
let mesh = make_mesh_addr(1);
let (vip, is_new) = pool.allocate(node, mesh, "test.fips").unwrap();
assert!(is_new);
assert_eq!(vip, "fd01::1".parse::<Ipv6Addr>().unwrap());
assert_eq!(pool.available.len(), 254);
}
#[test]
fn test_pool_idempotent() {
let mut pool = VirtualIpPool::new("fd01::/120", 60, 60).unwrap();
let node = make_node_addr(1);
let mesh = make_mesh_addr(1);
let (vip1, new1) = pool.allocate(node, mesh, "test.fips").unwrap();
let (vip2, new2) = pool.allocate(node, mesh, "test.fips").unwrap();
assert!(new1);
assert!(!new2);
assert_eq!(vip1, vip2);
assert_eq!(pool.available.len(), 254);
}
#[test]
fn test_pool_exhaustion() {
// /126 = 2 host bits = 4 addresses, minus 1 = 3
let mut pool = VirtualIpPool::new("fd01::/126", 60, 60).unwrap();
assert_eq!(pool.total, 3);
for i in 1..=3u8 {
pool.allocate(make_node_addr(i), make_mesh_addr(i), "test.fips")
.unwrap();
}
assert!(
pool.allocate(make_node_addr(4), make_mesh_addr(4), "test.fips")
.is_err()
);
}
#[test]
fn test_mapping_lifecycle_allocated_to_free() {
let mut pool = VirtualIpPool::new("fd01::/120", 1, 1).unwrap();
let ct = Sessions::new();
let node = make_node_addr(1);
let mesh = make_mesh_addr(1);
pool.allocate(node, mesh, "test.fips").unwrap();
// Tick before TTL — no change
let now = Instant::now();
let events = pool.tick(now, &ct.snapshot());
assert!(events.is_empty());
assert_eq!(pool.mappings.len(), 1);
// Tick after TTL with no sessions — enters draining
let later = now + std::time::Duration::from_secs(2);
let events = pool.tick(later, &ct.snapshot());
assert!(events.is_empty());
assert_eq!(pool.mappings.len(), 1);
assert_eq!(
pool.mappings.values().next().unwrap().state,
MappingState::Draining
);
// Tick after grace period — freed
let after_grace = later + std::time::Duration::from_secs(2);
let events = pool.tick(after_grace, &ct.snapshot());
assert_eq!(events.len(), 1);
assert!(matches!(events[0], PoolEvent::MappingRemoved { .. }));
assert_eq!(pool.mappings.len(), 0);
assert_eq!(pool.available.len(), 255); // returned to pool
}
#[test]
fn test_mapping_lifecycle_active_draining_free() {
let mut pool = VirtualIpPool::new("fd01::/120", 1, 1).unwrap();
let mut ct = Sessions::new();
let node = make_node_addr(1);
let mesh = make_mesh_addr(1);
let (vip, _) = pool.allocate(node, mesh, "test.fips").unwrap();
// Simulate active sessions
ct.set(vip, 3);
let now = Instant::now();
let events = pool.tick(now, &ct.snapshot());
assert!(events.is_empty());
assert_eq!(pool.mappings[&node].state, MappingState::Active);
// TTL expires after sessions drop to 0 → Draining
let later = now + std::time::Duration::from_secs(2);
ct.set(vip, 0);
let events = pool.tick(later, &ct.snapshot());
assert!(events.is_empty());
assert_eq!(pool.mappings[&node].state, MappingState::Draining);
// Still draining, grace period not elapsed
let events = pool.tick(later, &ct.snapshot());
assert!(events.is_empty());
assert_eq!(pool.mappings[&node].state, MappingState::Draining);
// Grace period elapsed → Free
let much_later = later + std::time::Duration::from_secs(2);
let events = pool.tick(much_later, &ct.snapshot());
assert_eq!(events.len(), 1);
assert!(matches!(events[0], PoolEvent::MappingRemoved { .. }));
assert_eq!(pool.mappings.len(), 0);
}
#[test]
fn test_active_traffic_never_reclaimed() {
// A mapping with continuous sessions > 0 across many ticks
// spanning well past the TTL must never be reclaimed and must
// stay Active: live traffic refreshes last_referenced each tick.
let mut pool = VirtualIpPool::new("fd01::/120", 1, 1).unwrap();
let mut ct = Sessions::new();
let node = make_node_addr(1);
let mesh = make_mesh_addr(1);
let (vip, _) = pool.allocate(node, mesh, "test.fips").unwrap();
ct.set(vip, 2);
let mut t = Instant::now();
// First tick activates the mapping.
let events = pool.tick(t, &ct.snapshot());
assert!(events.is_empty());
assert_eq!(pool.mappings[&node].state, MappingState::Active);
// Advance many TTL-spans with continuous traffic.
for _ in 0..10 {
t += std::time::Duration::from_secs(5); // 5x the 1s TTL
let events = pool.tick(t, &ct.snapshot());
assert!(events.is_empty(), "mapping must not be reclaimed");
assert_eq!(
pool.mappings[&node].state,
MappingState::Active,
"mapping must stay Active while traffic flows"
);
}
assert_eq!(pool.mappings.len(), 1);
}
#[test]
fn test_bursty_draining_recovers_to_active() {
// Active -> drains when sessions hit 0 -> regains sessions before
// grace elapses -> recovers to Active and is not freed.
let mut pool = VirtualIpPool::new("fd01::/120", 1, 5).unwrap();
let mut ct = Sessions::new();
let node = make_node_addr(1);
let mesh = make_mesh_addr(1);
let (vip, _) = pool.allocate(node, mesh, "test.fips").unwrap();
// Activate with traffic.
ct.set(vip, 1);
let now = Instant::now();
let events = pool.tick(now, &ct.snapshot());
assert!(events.is_empty());
assert_eq!(pool.mappings[&node].state, MappingState::Active);
// TTL passes with sessions dropping to 0 -> Draining.
let drained = now + std::time::Duration::from_secs(2);
ct.set(vip, 0);
let events = pool.tick(drained, &ct.snapshot());
assert!(events.is_empty());
assert_eq!(pool.mappings[&node].state, MappingState::Draining);
// Traffic resumes before grace (5s) elapses -> recover to Active.
let resumed = drained + std::time::Duration::from_secs(2);
ct.set(vip, 3);
let events = pool.tick(resumed, &ct.snapshot());
assert!(events.is_empty());
assert_eq!(pool.mappings[&node].state, MappingState::Active);
assert!(pool.mappings[&node].drain_start.is_none());
assert_eq!(pool.mappings.len(), 1);
}
#[test]
fn test_redrain_honors_fresh_grace_window() {
// After recovering from Draining, a subsequent drain must get a
// fresh drain_start so the full grace window is honored again,
// not reclaimed immediately off a stale drain_start.
let mut pool = VirtualIpPool::new("fd01::/120", 1, 5).unwrap();
let mut ct = Sessions::new();
let node = make_node_addr(1);
let mesh = make_mesh_addr(1);
let (vip, _) = pool.allocate(node, mesh, "test.fips").unwrap();
// Activate.
ct.set(vip, 1);
let now = Instant::now();
pool.tick(now, &ct.snapshot());
assert_eq!(pool.mappings[&node].state, MappingState::Active);
// First drain.
let first_drain = now + std::time::Duration::from_secs(2);
ct.set(vip, 0);
pool.tick(first_drain, &ct.snapshot());
assert_eq!(pool.mappings[&node].state, MappingState::Draining);
// Recover.
let recover = first_drain + std::time::Duration::from_secs(2);
ct.set(vip, 2);
pool.tick(recover, &ct.snapshot());
assert_eq!(pool.mappings[&node].state, MappingState::Active);
// Second drain begins; drain_start must be re-stamped fresh.
let second_drain = recover + std::time::Duration::from_secs(2);
ct.set(vip, 0);
pool.tick(second_drain, &ct.snapshot());
assert_eq!(pool.mappings[&node].state, MappingState::Draining);
// Just before the fresh grace window expires (5s): not reclaimed.
let before_grace = second_drain + std::time::Duration::from_secs(4);
let events = pool.tick(before_grace, &ct.snapshot());
assert!(events.is_empty(), "fresh grace window must be honored");
assert_eq!(pool.mappings.len(), 1);
// After the fresh grace window: reclaimed.
let after_grace = second_drain + std::time::Duration::from_secs(6);
let events = pool.tick(after_grace, &ct.snapshot());
assert_eq!(events.len(), 1);
assert!(matches!(events[0], PoolEvent::MappingRemoved { .. }));
assert_eq!(pool.mappings.len(), 0);
}
#[test]
fn test_pool_status() {
let mut pool = VirtualIpPool::new("fd01::/120", 60, 60).unwrap();
let status = pool.status();
assert_eq!(status.total, 255);
assert_eq!(status.free, 255);
assert_eq!(status.allocated, 0);
pool.allocate(make_node_addr(1), make_mesh_addr(1), "test.fips")
.unwrap();
let status = pool.status();
assert_eq!(status.allocated, 1);
assert_eq!(status.free, 254);
}
#[test]
fn test_lookup_virtual_ip() {
let mut pool = VirtualIpPool::new("fd01::/120", 60, 60).unwrap();
let node = make_node_addr(1);
let mesh = make_mesh_addr(1);
let (vip, _) = pool.allocate(node, mesh, "test.fips").unwrap();
let mapping = pool.lookup_virtual_ip(&vip).unwrap();
assert_eq!(mapping.node_addr, node);
assert_eq!(mapping.mesh_addr, mesh);
let unknown: Ipv6Addr = "fd01::ff".parse().unwrap();
assert!(pool.lookup_virtual_ip(&unknown).is_none());
}
#[test]
fn test_large_prefix_capped() {
// /96 = 32 host bits, but pool caps at 2^16
let pool = VirtualIpPool::new("fd01::/96", 60, 60).unwrap();
assert_eq!(pool.total, 65535); // 2^16 - 1 (skip addr 0)
}
/// A conntrack line in the form the kernel prints.
///
/// Built from the kernel's own format string, not captured from a running
/// kernel: `net/netfilter/nf_conntrack_standalone.c` prints each tuple with
/// `"src=%pI6 dst=%pI6 "`, and `%pI6` is the full uncompressed form with
/// leading zeros (`Documentation/core-api/printk-formats.rst`). Both were
/// read at v6.8. The host this was written on has no
/// `/proc/net/nf_conntrack` to capture from, because its kernel is built
/// without `CONFIG_NF_CONNTRACK_PROCFS`; OpenWrt's generic kernel config
/// sets it, which is the kernel this parser exists for.
const KERNEL_LINE: &str = "ipv6 10 tcp 6 431999 ESTABLISHED \
src=fd02:0000:0000:0000:0000:0000:0000:0020 \
dst=fd01:0000:0000:0000:0000:0000:0000:0001 sport=45678 dport=8000 \
src=fd01:0000:0000:0000:0000:0000:0000:0001 \
dst=fd02:0000:0000:0000:0000:0000:0000:0020 sport=8000 dport=45678 \
[ASSURED] mark=0 use=1";
#[test]
fn conntrack_parse_counts_a_kernel_format_line_for_its_virtual_ip() {
let counts = parse_conntrack(KERNEL_LINE);
let virtual_ip: Ipv6Addr = "fd01::1".parse().unwrap();
assert_eq!(
counts.get(&virtual_ip).copied().unwrap_or(0),
1,
"the kernel writes the uncompressed form, so matching on the \
address's compressed Display form counts nothing"
);
// Healthy path: a different address in the same pool is not counted.
let other: Ipv6Addr = "fd01::10".parse().unwrap();
assert_eq!(counts.get(&other).copied().unwrap_or(0), 0);
}
#[test]
fn conntrack_parse_counts_a_line_once_however_many_tuples_name_the_address() {
// A hairpin flow: the address is the destination of both tuples.
let line = "ipv6 10 udp 17 29 \
src=fd01:0000:0000:0000:0000:0000:0000:0001 \
dst=fd01:0000:0000:0000:0000:0000:0000:0001 sport=1 dport=2 \
src=fd01:0000:0000:0000:0000:0000:0000:0001 \
dst=fd01:0000:0000:0000:0000:0000:0000:0001 sport=2 dport=1 \
mark=0 use=1";
let counts = parse_conntrack(line);
let virtual_ip: Ipv6Addr = "fd01::1".parse().unwrap();
assert_eq!(counts.get(&virtual_ip).copied().unwrap_or(0), 1);
}
#[test]
fn conntrack_parse_counts_each_line_that_names_the_address() {
let content = format!("{KERNEL_LINE}\n{KERNEL_LINE}\n");
let counts = parse_conntrack(&content);
let virtual_ip: Ipv6Addr = "fd01::1".parse().unwrap();
assert_eq!(counts.get(&virtual_ip).copied().unwrap_or(0), 2);
}
#[test]
fn conntrack_parse_skips_a_value_that_is_not_an_ipv6_address() {
let content = "ipv4 2 tcp 6 431999 ESTABLISHED src=192.0.2.1 \
dst=192.0.2.2 sport=1 dport=2 mark=0 use=1\n";
assert!(parse_conntrack(content).is_empty());
}
#[test]
fn conntrack_snapshot_reads_zero_for_an_address_it_did_not_see() {
let snapshot = ConntrackSnapshot::from_counts(parse_conntrack(KERNEL_LINE));
assert_eq!(snapshot.sessions_for("fd01::1".parse().unwrap()), 1);
assert_eq!(snapshot.sessions_for("fd01::99".parse().unwrap()), 0);
assert!(ConntrackSnapshot::default().is_empty());
}
#[test]
fn conntrack_read_log_warns_on_a_new_outcome_and_not_on_a_repeat() {
use std::io::ErrorKind;
let mut log = ConntrackReadLog::default();
// The sequence a kernel without the proc file produces, then a source
// that comes back, then fails again.
assert_eq!(log.observe(Some(ErrorKind::NotFound)), ReadReport::Changed);
assert_eq!(log.observe(Some(ErrorKind::NotFound)), ReadReport::Repeated);
assert_eq!(log.observe(None), ReadReport::Changed);
assert_eq!(log.observe(None), ReadReport::Repeated);
assert_eq!(log.observe(Some(ErrorKind::NotFound)), ReadReport::Changed);
assert_eq!(
log.observe(Some(ErrorKind::PermissionDenied)),
ReadReport::Changed,
"a different failure is a different outcome and is worth a line"
);
}
}