Merge branch 'master' into next

This commit is contained in:
Johnathan Corgan
2026-09-18 20:42:26 +00:00
6 changed files with 1651 additions and 140 deletions
+22
View File
@@ -861,6 +861,28 @@ with v0.5.x or earlier peers.
now share a single batch, which the kernel applies as one transaction, so a
refused rebuild leaves the previous table in the packet path. The rules sent
are unchanged.
- The gateway's NAT rebuild no longer fails once the table holds more than
about 105 mappings. Each rebuild is one netlink batch. From about 105
mappings the default socket buffers could not hold its acknowledgements, so
rebuilds were logged as failed although they had taken effect. Past about
313 mappings the buffers could not hold the batch itself, and new `.fips`
names past that count got a virtual IP with no translation. In releases with
the gateway through 0.5.1, a rebuild past about 313 mappings also deleted the
whole `fips_gateway` table, which stopped every mapping, the `fips0`
masquerade and the port forwards. The rebuild now sizes its send buffer to
the batch and requests one acknowledgement per batch, and NAT errors now
name the kernel errno. A rebuild that still fails is logged, and the next
successful rebuild installs the mapping.
- The gateway's virtual-IP pool now limits how many mappings it holds and how
fast it creates them. Any host that can reach the LAN resolver could ask for
one new `.fips` name after another, and each got a mapping until the 65,535
addresses ran out, while every mapping made each NAT rebuild, each pool tick
and shutdown slower. The pool now refuses a new name once it holds 1000 live
mappings, and admits new names at 10 per second after a burst of 50. A
refused query gets SERVFAIL, and the gateway's "Pool allocation failed"
warning says which limit refused it. A name that already has a mapping is
answered before either limit is consulted, so names in use keep resolving
when the pool is full. The limits are compiled in, not configured.
- A new OpenWrt install no longer enables and starts `fips-gateway`. The
generated postinst turned it on unconditionally, contradicting the init
script's own header, the package README and the deployment tutorial, all of
+11
View File
@@ -53,6 +53,12 @@ fn main() {
std::process::exit(1);
}
/// Microseconds since `started`, saturating, for the timing fields on debug
/// lines.
fn elapsed_us(started: Instant) -> u64 {
u64::try_from(started.elapsed().as_micros()).unwrap_or(u64::MAX)
}
/// Take a conntrack snapshot off the runtime thread.
///
/// A failed read yields an empty snapshot, so every mapping reads zero
@@ -446,14 +452,19 @@ async fn main() {
// the pool lock: the runtime is current-thread, so a
// blocking read here would stall the DNS resolver, and the
// read must not happen under the lock the resolver needs.
let read_started = Instant::now();
let conntrack = read_conntrack(&mut conntrack_log).await;
let read_us = elapsed_us(read_started);
let mut pool_guard = tick_pool.lock().await;
let tick_started = Instant::now();
let events = pool_guard.tick(now, &conntrack);
let tick_us = elapsed_us(tick_started);
// Build snapshot for control socket
let pool_status = pool_guard.status();
let mappings = pool_guard.mapping_info(now);
drop(pool_guard);
debug!(mappings = mappings.len(), read_us, tick_us, "Pool tick");
let snapshot = control::build_snapshot(
pool_status,
+783 -137
View File
@@ -4,7 +4,10 @@
//! for translating between virtual IPs and FIPS mesh addresses.
use std::collections::HashMap;
use std::fmt;
use std::net::Ipv6Addr;
use std::os::fd::{AsRawFd, FromRawFd, OwnedFd};
use std::time::Instant;
use tracing::{debug, info};
use rustables::expr::{
@@ -23,6 +26,73 @@ const POSTROUTING_CHAIN: &str = "postrouting";
const DSTNAT_PRIORITY: i32 = -100;
const SRCNAT_PRIORITY: i32 = 100;
/// Largest value the kernel accepts for `SO_SNDBUFFORCE`.
///
/// The kernel clamps the requested value to `i32::MAX / 2` and then doubles
/// it, so the socket's send buffer never exceeds `2 * MAX_SNDBUF`.
const MAX_SNDBUF: libc::c_int = libc::c_int::MAX / 2;
/// Headroom added to half the batch length when sizing the send buffer.
const SNDBUF_HEADROOM: u64 = 64 * 1024;
/// The kernel refuses a netlink message longer than the send buffer less
/// this many bytes.
const SNDBUF_OVERHEAD: u64 = 32;
/// How long the rebuild waits for the kernel's acknowledgement. The rebuild
/// runs inside the gateway's event loop, so this bounds the stall there.
const ACK_TIMEOUT_SECS: libc::time_t = 5;
/// Length of a `struct nlmsghdr`.
const NLMSG_HDRLEN: usize = 16;
/// An errno value, displayed by name and number, e.g. `EMSGSIZE (90)`.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Errno(pub i32);
impl Errno {
/// The errno of the last failed libc call on this thread.
fn last() -> Self {
Errno(std::io::Error::last_os_error().raw_os_error().unwrap_or(0))
}
/// The symbolic name of the errno, for the values netlink can return.
fn name(self) -> &'static str {
match self.0 {
libc::EPERM => "EPERM",
libc::ENOENT => "ENOENT",
libc::EINTR => "EINTR",
libc::EBADF => "EBADF",
libc::EAGAIN => "EAGAIN",
libc::ENOMEM => "ENOMEM",
libc::EACCES => "EACCES",
libc::EFAULT => "EFAULT",
libc::EBUSY => "EBUSY",
libc::EEXIST => "EEXIST",
libc::ENODEV => "ENODEV",
libc::EINVAL => "EINVAL",
libc::ENFILE => "ENFILE",
libc::EMFILE => "EMFILE",
libc::ENOSPC => "ENOSPC",
libc::ERANGE => "ERANGE",
libc::ELOOP => "ELOOP",
libc::EMSGSIZE => "EMSGSIZE",
libc::EPROTONOSUPPORT => "EPROTONOSUPPORT",
libc::EOPNOTSUPP => "EOPNOTSUPP",
libc::EAFNOSUPPORT => "EAFNOSUPPORT",
libc::ENOBUFS => "ENOBUFS",
libc::ETIMEDOUT => "ETIMEDOUT",
_ => "errno",
}
}
}
impl fmt::Display for Errno {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{} ({})", self.name(), self.0)
}
}
/// Errors from NAT operations.
#[derive(Debug, thiserror::Error)]
pub enum NatError {
@@ -30,6 +100,18 @@ pub enum NatError {
Nftables(String),
#[error("rule not found for virtual IP {0}")]
RuleNotFound(Ipv6Addr),
/// The kernel rejected a message of the NAT batch; the batch was aborted.
#[error("kernel rejected netlink message {seq} of the NAT batch: {errno}")]
Kernel { errno: Errno, seq: u32 },
/// A netlink socket call failed.
#[error("netlink socket {op} failed: {errno}")]
Socket { op: &'static str, errno: Errno },
/// The batch is larger than any netlink send buffer the kernel allows.
#[error(
"NAT batch of {bytes} bytes exceeds the kernel's netlink limit of {} bytes",
admissible_limit()
)]
BatchTooLarge { bytes: usize },
}
impl From<rustables::error::QueryError> for NatError {
@@ -54,8 +136,8 @@ struct NatMapping {
/// One object a NAT rebuild sends, named rather than built.
///
/// `rebuild_batches` decides what a rebuild sends and in what order;
/// `send_batches` turns that decision into rustables objects and hands each
/// batch to the kernel. The split is what lets a test see the delete and the
/// `encode_batch` turns that decision into netlink bytes and `send_batch`
/// hands them to the kernel. The split is what lets a test see the delete and the
/// recreate share one transaction without a netlink socket, which is the
/// property that keeps the table in the packet path.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
@@ -161,14 +243,26 @@ impl NatManager {
mesh_addr,
},
);
self.rebuild()?;
debug!(
virtual_ip = %virtual_ip,
mesh_addr = %mesh_addr,
"Added DNAT/SNAT rules"
);
Ok(())
let (result, elapsed_us) = self.timed_rebuild();
let mappings = self.mappings.len();
match &result {
Ok(()) => debug!(
virtual_ip = %virtual_ip,
mesh_addr = %mesh_addr,
mappings,
elapsed_us,
"Added DNAT/SNAT rules"
),
Err(e) => debug!(
virtual_ip = %virtual_ip,
mesh_addr = %mesh_addr,
mappings,
elapsed_us,
error = %e,
"Added DNAT/SNAT rules"
),
}
result
}
/// Remove DNAT and SNAT rules for a virtual IP mapping.
@@ -176,10 +270,24 @@ impl NatManager {
if self.mappings.remove(&virtual_ip).is_none() {
return Err(NatError::RuleNotFound(virtual_ip));
}
self.rebuild()?;
debug!(virtual_ip = %virtual_ip, "Removed DNAT/SNAT rules");
Ok(())
let (result, elapsed_us) = self.timed_rebuild();
let mappings = self.mappings.len();
match &result {
Ok(()) => debug!(
virtual_ip = %virtual_ip,
mappings,
elapsed_us,
"Removed DNAT/SNAT rules"
),
Err(e) => debug!(
virtual_ip = %virtual_ip,
mappings,
elapsed_us,
error = %e,
"Removed DNAT/SNAT rules"
),
}
result
}
/// Flush all rules and delete the nftables table.
@@ -188,7 +296,7 @@ impl NatManager {
batch.add(&self.table, MsgType::Del);
batch
.send()
.map_err(|e| NatError::Nftables(e.to_string()))?;
.map_err(|e| NatError::Nftables(error_chain(&e)))?;
info!("Deleted nftables table '{TABLE_NAME}'");
Ok(())
@@ -238,133 +346,137 @@ impl NatManager {
vec![ops]
}
/// Build each op into its rustables object and send the batches in order.
fn send_batches(&self, batches: &[Vec<NatOp>]) -> Result<(), NatError> {
for ops in batches {
let mut batch = Batch::new();
for op in ops {
match *op {
NatOp::Table(msg_type) => batch.add(&self.table, msg_type),
NatOp::PreChain => batch.add(&self.pre_chain, MsgType::Add),
NatOp::PostChain => batch.add(&self.post_chain, MsgType::Add),
NatOp::FipsMasquerade => {
// Rewrite the source address of traffic leaving fips0.
// Without this, LAN clients' source addresses (e.g.
// fd02::20) are not routable on the mesh, so return
// traffic would be black-holed.
let rule = Rule::new(&self.post_chain)?
.with_expr(Meta::new(MetaType::OifName))
.with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec()))
.with_expr(Masquerade::default());
batch.add(&rule, MsgType::Add);
}
NatOp::Dnat(virtual_ip) => {
let mapping = self.mapping(virtual_ip)?;
let rule = Rule::new(&self.pre_chain)?
.with_expr(Meta::new(MetaType::NfProto))
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
.with_expr(
HighLevelPayload::Network(NetworkHeaderField::IPv6(
IPv6HeaderField::Daddr,
))
.build(),
)
.with_expr(Cmp::new(CmpOp::Eq, mapping.virtual_ip.octets()))
.with_expr(Immediate::new_data(
mapping.mesh_addr.octets().to_vec(),
Register::Reg1,
/// Build each op into its rustables object and encode the batch.
///
/// Only the last object before the batch end requests an
/// acknowledgement. rustables sets `NLM_F_ACK` on every message, and one
/// ack per message overflows the socket's receive buffer from about a
/// hundred mappings, after the kernel has already committed the batch.
/// The kernel reports a failing message whatever its flags, so errors
/// stay attributable.
fn encode_batch(&self, ops: &[NatOp]) -> Result<Vec<u8>, NatError> {
let mut batch = Batch::new();
for op in ops {
match *op {
NatOp::Table(msg_type) => batch.add(&self.table, msg_type),
NatOp::PreChain => batch.add(&self.pre_chain, MsgType::Add),
NatOp::PostChain => batch.add(&self.post_chain, MsgType::Add),
NatOp::FipsMasquerade => {
// Rewrite the source address of traffic leaving fips0.
// Without this, LAN clients' source addresses (e.g.
// fd02::20) are not routable on the mesh, so return
// traffic would be black-holed.
let rule = Rule::new(&self.post_chain)?
.with_expr(Meta::new(MetaType::OifName))
.with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec()))
.with_expr(Masquerade::default());
batch.add(&rule, MsgType::Add);
}
NatOp::Dnat(virtual_ip) => {
let mapping = self.mapping(virtual_ip)?;
let rule = Rule::new(&self.pre_chain)?
.with_expr(Meta::new(MetaType::NfProto))
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
.with_expr(
HighLevelPayload::Network(NetworkHeaderField::IPv6(
IPv6HeaderField::Daddr,
))
.with_expr(
Nat::default()
.with_nat_type(NatType::DNat)
.with_family(ProtocolFamily::Ipv6)
.with_ip_register(Register::Reg1),
);
batch.add(&rule, MsgType::Add);
}
NatOp::Snat(virtual_ip) => {
let mapping = self.mapping(virtual_ip)?;
let rule = Rule::new(&self.post_chain)?
.with_expr(Meta::new(MetaType::NfProto))
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
.with_expr(
HighLevelPayload::Network(NetworkHeaderField::IPv6(
IPv6HeaderField::Saddr,
))
.build(),
)
.with_expr(Cmp::new(CmpOp::Eq, mapping.mesh_addr.octets()))
.with_expr(Immediate::new_data(
mapping.virtual_ip.octets().to_vec(),
Register::Reg1,
))
.with_expr(
Nat::default()
.with_nat_type(NatType::SNat)
.with_family(ProtocolFamily::Ipv6)
.with_ip_register(Register::Reg1),
);
batch.add(&rule, MsgType::Add);
}
NatOp::PortForward(index) => {
let pf = self.port_forwards.get(index).expect(
"rebuild_batches only emits indices it read from port_forwards",
.build(),
)
.with_expr(Cmp::new(CmpOp::Eq, mapping.virtual_ip.octets()))
.with_expr(Immediate::new_data(
mapping.mesh_addr.octets().to_vec(),
Register::Reg1,
))
.with_expr(
Nat::default()
.with_nat_type(NatType::DNat)
.with_family(ProtocolFamily::Ipv6)
.with_ip_register(Register::Reg1),
);
let l4proto: u8 = match pf.proto {
Proto::Tcp => libc::IPPROTO_TCP as u8,
Proto::Udp => libc::IPPROTO_UDP as u8,
};
let dport_field = match pf.proto {
Proto::Tcp => TransportHeaderField::Tcp(TCPHeaderField::Dport),
Proto::Udp => TransportHeaderField::Udp(UDPHeaderField::Dport),
};
let target_ip = *pf.target.ip();
let target_port_be = pf.target.port().to_be_bytes();
let rule = Rule::new(&self.pre_chain)?
.with_expr(Meta::new(MetaType::IifName))
.with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec()))
.with_expr(Meta::new(MetaType::NfProto))
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
.with_expr(Meta::new(MetaType::L4Proto))
.with_expr(Cmp::new(CmpOp::Eq, [l4proto]))
.with_expr(HighLevelPayload::Transport(dport_field).build())
.with_expr(Cmp::new(CmpOp::Eq, pf.listen_port.to_be_bytes().to_vec()))
.with_expr(Immediate::new_data(
target_ip.octets().to_vec(),
Register::Reg1,
batch.add(&rule, MsgType::Add);
}
NatOp::Snat(virtual_ip) => {
let mapping = self.mapping(virtual_ip)?;
let rule = Rule::new(&self.post_chain)?
.with_expr(Meta::new(MetaType::NfProto))
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
.with_expr(
HighLevelPayload::Network(NetworkHeaderField::IPv6(
IPv6HeaderField::Saddr,
))
.with_expr(Immediate::new_data(target_port_be.to_vec(), Register::Reg2))
.with_expr(
Nat::default()
.with_nat_type(NatType::DNat)
.with_family(ProtocolFamily::Ipv6)
.with_ip_register(Register::Reg1)
.with_port_register(Register::Reg2),
);
batch.add(&rule, MsgType::Add);
}
NatOp::LanMasquerade => {
let mut lan_iface = self.lan_interface.clone().into_bytes();
lan_iface.push(0);
let rule = Rule::new(&self.post_chain)?
.with_expr(Meta::new(MetaType::IifName))
.with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec()))
.with_expr(Meta::new(MetaType::OifName))
.with_expr(Cmp::new(CmpOp::Eq, lan_iface))
.with_expr(Meta::new(MetaType::NfProto))
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
.with_expr(Masquerade::default());
batch.add(&rule, MsgType::Add);
}
.build(),
)
.with_expr(Cmp::new(CmpOp::Eq, mapping.mesh_addr.octets()))
.with_expr(Immediate::new_data(
mapping.virtual_ip.octets().to_vec(),
Register::Reg1,
))
.with_expr(
Nat::default()
.with_nat_type(NatType::SNat)
.with_family(ProtocolFamily::Ipv6)
.with_ip_register(Register::Reg1),
);
batch.add(&rule, MsgType::Add);
}
NatOp::PortForward(index) => {
let pf = self
.port_forwards
.get(index)
.expect("rebuild_batches only emits indices it read from port_forwards");
let l4proto: u8 = match pf.proto {
Proto::Tcp => libc::IPPROTO_TCP as u8,
Proto::Udp => libc::IPPROTO_UDP as u8,
};
let dport_field = match pf.proto {
Proto::Tcp => TransportHeaderField::Tcp(TCPHeaderField::Dport),
Proto::Udp => TransportHeaderField::Udp(UDPHeaderField::Dport),
};
let target_ip = *pf.target.ip();
let target_port_be = pf.target.port().to_be_bytes();
let rule = Rule::new(&self.pre_chain)?
.with_expr(Meta::new(MetaType::IifName))
.with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec()))
.with_expr(Meta::new(MetaType::NfProto))
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
.with_expr(Meta::new(MetaType::L4Proto))
.with_expr(Cmp::new(CmpOp::Eq, [l4proto]))
.with_expr(HighLevelPayload::Transport(dport_field).build())
.with_expr(Cmp::new(CmpOp::Eq, pf.listen_port.to_be_bytes().to_vec()))
.with_expr(Immediate::new_data(
target_ip.octets().to_vec(),
Register::Reg1,
))
.with_expr(Immediate::new_data(target_port_be.to_vec(), Register::Reg2))
.with_expr(
Nat::default()
.with_nat_type(NatType::DNat)
.with_family(ProtocolFamily::Ipv6)
.with_ip_register(Register::Reg1)
.with_port_register(Register::Reg2),
);
batch.add(&rule, MsgType::Add);
}
NatOp::LanMasquerade => {
let mut lan_iface = self.lan_interface.clone().into_bytes();
lan_iface.push(0);
let rule = Rule::new(&self.post_chain)?
.with_expr(Meta::new(MetaType::IifName))
.with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec()))
.with_expr(Meta::new(MetaType::OifName))
.with_expr(Cmp::new(CmpOp::Eq, lan_iface))
.with_expr(Meta::new(MetaType::NfProto))
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
.with_expr(Masquerade::default());
batch.add(&rule, MsgType::Add);
}
}
batch
.send()
.map_err(|e| NatError::Nftables(e.to_string()))?;
}
Ok(())
let mut bytes = batch.finalize();
keep_last_ack(&mut bytes)?;
Ok(bytes)
}
/// The mapping an op names, or the error a caller can report.
@@ -377,10 +489,372 @@ impl NatManager {
/// Rebuild the entire nftables table with all current rules, in one
/// netlink transaction.
fn rebuild(&self) -> Result<(), NatError> {
self.send_batches(&self.rebuild_batches())
for ops in self.rebuild_batches() {
send_batch(&self.encode_batch(&ops)?)?;
}
Ok(())
}
/// Rebuild, returning the outcome with the time the rebuild took in
/// microseconds, so a mapping change can log its cost on either path.
fn timed_rebuild(&self) -> (Result<(), NatError>, u64) {
let started = Instant::now();
let result = self.rebuild();
let elapsed_us = u64::try_from(started.elapsed().as_micros()).unwrap_or(u64::MAX);
(result, elapsed_us)
}
}
/// Largest batch, in bytes, that the kernel can admit in one send.
fn admissible_limit() -> u64 {
2 * MAX_SNDBUF as u64 - SNDBUF_OVERHEAD
}
/// The `SO_SNDBUFFORCE` value that lets a batch of `len` bytes through.
///
/// The kernel doubles the value it is given and refuses a message longer
/// than the result less 32 bytes, so half the length plus headroom is
/// enough. Saturates at the kernel's own clamp rather than wrapping.
fn sndbuf_for(len: usize) -> libc::c_int {
let want = (u64::try_from(len).unwrap_or(u64::MAX) / 2).saturating_add(SNDBUF_HEADROOM);
libc::c_int::try_from(want.min(MAX_SNDBUF as u64)).unwrap_or(MAX_SNDBUF)
}
/// Refuse a batch the kernel could not accept at any send-buffer size.
fn check_admissible(len: usize) -> Result<(), NatError> {
if u64::try_from(len).unwrap_or(u64::MAX) > admissible_limit() {
return Err(NatError::BatchTooLarge { bytes: len });
}
Ok(())
}
/// The fields of one `struct nlmsghdr` that the NAT batch code reads.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct NlHeader {
/// Offset of the header within the buffer.
offset: usize,
/// `nlmsg_len`: header plus payload, without alignment padding.
len: usize,
kind: u16,
flags: u16,
seq: u32,
}
impl NlHeader {
/// The message's payload, after the header.
fn payload<'a>(&self, buf: &'a [u8]) -> &'a [u8] {
&buf[self.offset + NLMSG_HDRLEN..self.offset + self.len]
}
}
/// Walk every netlink message header in `buf`.
///
/// Fails on a header shorter than `struct nlmsghdr` or a length that runs
/// past the end of the buffer.
fn nl_headers(buf: &[u8]) -> Result<Vec<NlHeader>, NatError> {
let mut headers = Vec::new();
let mut offset = 0;
while offset < buf.len() {
let rest = &buf[offset..];
if rest.len() < NLMSG_HDRLEN {
return Err(NatError::Nftables(format!(
"malformed netlink message at offset {offset}: {} bytes left, header needs {NLMSG_HDRLEN}",
rest.len()
)));
}
let field = |at: usize, width: usize| &rest[at..at + width];
let len = u32::from_ne_bytes(field(0, 4).try_into().expect("4-byte slice")) as usize;
if len < NLMSG_HDRLEN || len > rest.len() {
return Err(NatError::Nftables(format!(
"malformed netlink message at offset {offset}: length {len} with {} bytes left",
rest.len()
)));
}
headers.push(NlHeader {
offset,
len,
kind: u16::from_ne_bytes(field(4, 2).try_into().expect("2-byte slice")),
flags: u16::from_ne_bytes(field(6, 2).try_into().expect("2-byte slice")),
seq: u32::from_ne_bytes(field(8, 4).try_into().expect("4-byte slice")),
});
// Netlink messages are 4-byte aligned.
offset += (len + 3) & !3;
}
Ok(headers)
}
/// The finalized batch's objects: every message between the batch begin
/// and the batch end.
fn batch_objects(headers: &[NlHeader]) -> Result<&[NlHeader], NatError> {
match headers {
[_begin, objects @ .., _end] if !objects.is_empty() => Ok(objects),
_ => Err(NatError::Nftables(format!(
"NAT batch holds {} messages; it needs a begin, an object and an end",
headers.len()
))),
}
}
/// Clear `NLM_F_ACK` on every message of a finalized batch except the last
/// object before the batch end.
fn keep_last_ack(buf: &mut [u8]) -> Result<(), NatError> {
let headers = nl_headers(buf)?;
let last = batch_objects(&headers)?
.last()
.expect("batch_objects returns a non-empty slice")
.offset;
let ack = libc::NLM_F_ACK as u16;
for header in &headers {
let flags = if header.offset == last {
header.flags | ack
} else {
header.flags & !ack
};
buf[header.offset + 6..header.offset + 8].copy_from_slice(&flags.to_ne_bytes());
}
Ok(())
}
/// What the kernel's replies to a NAT batch have shown so far.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum AckState {
/// No verdict yet; read another datagram.
Pending,
/// The acknowledgement of the batch's last object arrived with no error
/// before it.
Done,
}
/// Reads the kernel's replies to a NAT batch, one datagram at a time.
///
/// The kernel aborts the whole batch when any message fails, yet it still
/// acknowledges the last message after the error. So the last ack alone does
/// not prove success: any error that arrives before it fails the batch.
struct AckReader {
/// Sequence number of the one message that requested an ack.
last_seq: u32,
}
impl AckReader {
/// Consume one received datagram, which may carry several messages.
fn feed(&self, datagram: &[u8]) -> Result<AckState, NatError> {
for header in nl_headers(datagram)? {
if i32::from(header.kind) != libc::NLMSG_ERROR {
continue;
}
let payload = header.payload(datagram);
if payload.len() < 4 {
return Err(NatError::Nftables(format!(
"malformed netlink error message: {} payload bytes, error field needs 4",
payload.len()
)));
}
let error = i32::from_ne_bytes(payload[..4].try_into().expect("4-byte slice"));
if error != 0 {
return Err(NatError::Kernel {
errno: Errno(error.saturating_neg()),
seq: header.seq,
});
}
if header.seq == self.last_seq {
return Ok(AckState::Done);
}
}
Ok(AckState::Pending)
}
}
/// Render an error with every source beneath it, so a wrapped errno is kept.
fn error_chain(error: &dyn std::error::Error) -> String {
let mut text = error.to_string();
let mut source = error.source();
while let Some(inner) = source {
text.push_str(": ");
text.push_str(&inner.to_string());
source = inner.source();
}
text
}
/// Set an integer socket option.
fn set_int_opt(
sock: &OwnedFd,
level: libc::c_int,
name: libc::c_int,
value: libc::c_int,
) -> Result<(), Errno> {
// SAFETY: the descriptor is open for the life of `sock`, and the pointer
// and length describe `value`, a c_int.
let rc = unsafe {
libc::setsockopt(
sock.as_raw_fd(),
level,
name,
(&value as *const libc::c_int).cast(),
std::mem::size_of::<libc::c_int>() as libc::socklen_t,
)
};
if rc < 0 { Err(Errno::last()) } else { Ok(()) }
}
/// Size of the buffer each reply datagram is read into.
///
/// The largest message nftables sends back, as rustables computes it
/// (`nft_nlmsg_maxsize`, which it does not export), and at least 64 KiB.
fn recv_buffer_len() -> usize {
// SAFETY: sysconf has no preconditions.
let page = unsafe { libc::sysconf(libc::_SC_PAGESIZE) };
(usize::from(u16::MAX) + usize::try_from(page).unwrap_or(0)).max(64 * 1024)
}
/// Open a netfilter netlink socket sized for a batch of `len` bytes.
fn open_batch_socket(len: usize) -> Result<OwnedFd, NatError> {
let socket_err = |op| move |errno| NatError::Socket { op, errno };
// SAFETY: socket has no memory preconditions.
let fd = unsafe {
libc::socket(
libc::AF_NETLINK,
libc::SOCK_RAW | libc::SOCK_CLOEXEC,
libc::NETLINK_NETFILTER,
)
};
if fd < 0 {
return Err(socket_err("open")(Errno::last()));
}
// SAFETY: socket returned a new descriptor that nothing else owns.
let sock = unsafe { OwnedFd::from_raw_fd(fd) };
// SAFETY: an all-zero sockaddr_nl is valid; the family is set below.
let mut addr: libc::sockaddr_nl = unsafe { std::mem::zeroed() };
addr.nl_family = libc::AF_NETLINK as libc::sa_family_t;
// SAFETY: the pointer and length describe `addr`, a sockaddr_nl.
let rc = unsafe {
libc::bind(
sock.as_raw_fd(),
(&addr as *const libc::sockaddr_nl).cast(),
std::mem::size_of::<libc::sockaddr_nl>() as libc::socklen_t,
)
};
if rc < 0 {
return Err(socket_err("bind")(Errno::last()));
}
// Without CAP_NET_ADMIN the forced size is refused; the plain option is
// then capped by wmem_max, and an oversized batch fails with EMSGSIZE.
let sndbuf = sndbuf_for(len);
match set_int_opt(&sock, libc::SOL_SOCKET, libc::SO_SNDBUFFORCE, sndbuf) {
Err(Errno(libc::EPERM)) => {
set_int_opt(&sock, libc::SOL_SOCKET, libc::SO_SNDBUF, sndbuf)
.map_err(socket_err("setsockopt SO_SNDBUF"))?;
}
other => other.map_err(socket_err("setsockopt SO_SNDBUFFORCE"))?,
}
// An error ack then carries only the failing header, not the message.
set_int_opt(&sock, libc::SOL_NETLINK, libc::NETLINK_CAP_ACK, 1)
.map_err(socket_err("setsockopt NETLINK_CAP_ACK"))?;
let timeout = libc::timeval {
tv_sec: ACK_TIMEOUT_SECS,
tv_usec: 0,
};
// SAFETY: the pointer and length describe `timeout`, a timeval.
let rc = unsafe {
libc::setsockopt(
sock.as_raw_fd(),
libc::SOL_SOCKET,
libc::SO_RCVTIMEO,
(&timeout as *const libc::timeval).cast(),
std::mem::size_of::<libc::timeval>() as libc::socklen_t,
)
};
if rc < 0 {
return Err(socket_err("setsockopt SO_RCVTIMEO")(Errno::last()));
}
Ok(sock)
}
/// Send one encoded NAT batch and wait for the kernel's verdict.
///
/// The batch goes out in a single send, so the kernel applies it as one
/// transaction. The send buffer is sized to the batch, because the default
/// one refuses a message past about 208 KiB, which is about 313 mappings.
fn send_batch(bytes: &[u8]) -> Result<(), NatError> {
check_admissible(bytes.len())?;
let headers = nl_headers(bytes)?;
let reader = AckReader {
last_seq: batch_objects(&headers)?
.last()
.expect("batch_objects returns a non-empty slice")
.seq,
};
let sock = open_batch_socket(bytes.len())?;
let sent = loop {
// SAFETY: the pointer and length describe `bytes`.
let rc = unsafe { libc::send(sock.as_raw_fd(), bytes.as_ptr().cast(), bytes.len(), 0) };
if rc >= 0 {
break rc as usize;
}
let errno = Errno::last();
if errno.0 != libc::EINTR {
return Err(NatError::Socket { op: "send", errno });
}
};
if sent != bytes.len() {
return Err(NatError::Nftables(format!(
"netlink send took {sent} of {} bytes",
bytes.len()
)));
}
let mut buf = vec![0u8; recv_buffer_len()];
loop {
// MSG_TRUNC makes a netlink recv return the datagram's full length,
// so a reply larger than the buffer is detected rather than cut.
// SAFETY: the pointer and length describe `buf`.
let rc = unsafe {
libc::recv(
sock.as_raw_fd(),
buf.as_mut_ptr().cast(),
buf.len(),
libc::MSG_TRUNC,
)
};
if rc < 0 {
let errno = Errno::last();
if errno.0 == libc::EINTR {
continue;
}
// EAGAIN is the receive timeout. ENOBUFS means error acks
// overflowed the receive buffer, since only one ack is requested.
return Err(NatError::Socket { op: "recv", errno });
}
let got = rc as usize;
if got == 0 {
return Err(NatError::Nftables(
"netlink socket returned no reply".into(),
));
}
if got > buf.len() {
return Err(NatError::Nftables(format!(
"netlink reply of {got} bytes truncated to {}",
buf.len()
)));
}
if reader.feed(&buf[..got])? == AckState::Done {
return Ok(());
}
}
}
// Coverage gap. These tests run unprivileged and open no netlink socket, so
// three failure paths in `open_socket` and `send_batch` go unexercised here.
// The receive timeout firing and a reply longer than the buffer (seen through
// `MSG_TRUNC`) need a kernel fault to provoke, so nothing runs them. The
// `SO_SNDBUF` fallback after `SO_SNDBUFFORCE` returns `EPERM` needs a process
// without CAP_NET_ADMIN, and the gateway suite's container is privileged, so
// nothing runs that either. The gateway suite covers only the success path.
#[cfg(test)]
mod tests {
use super::*;
@@ -497,4 +971,176 @@ mod tests {
1
);
}
/// The encoded rebuild of a manager holding `count` mappings.
fn encoded_rebuild(count: u16) -> Vec<u8> {
let mgr = manager_with_mappings(count);
let ops = mgr.rebuild_batches().remove(0);
mgr.encode_batch(&ops).expect("the rebuild encodes")
}
/// One netlink message, padded to 4 bytes.
fn nlmsg(kind: u16, seq: u32, payload: &[u8]) -> Vec<u8> {
let len = (NLMSG_HDRLEN + payload.len()) as u32;
let mut msg = Vec::new();
msg.extend_from_slice(&len.to_ne_bytes());
msg.extend_from_slice(&kind.to_ne_bytes());
msg.extend_from_slice(&0u16.to_ne_bytes());
msg.extend_from_slice(&seq.to_ne_bytes());
msg.extend_from_slice(&0u32.to_ne_bytes());
msg.extend_from_slice(payload);
msg.resize(msg.len().div_ceil(4) * 4, 0);
msg
}
/// The kernel's `NLMSG_ERROR` reply to message `seq`, as it sends it on a
/// socket with `NETLINK_CAP_ACK`: the error, then the request's header.
fn ack(seq: u32, error: i32) -> Vec<u8> {
let mut payload = error.to_ne_bytes().to_vec();
payload.extend_from_slice(&nlmsg(0x0a00, seq, &[])[..NLMSG_HDRLEN]);
nlmsg(libc::NLMSG_ERROR as u16, seq, &payload)
}
/// The largest batch the kernel admits: twice its send-buffer clamp,
/// less the 32 bytes netlink reserves.
const KERNEL_BATCH_LIMIT: usize = 2_147_483_614;
#[test]
fn rebuild_for_2000_mappings_requests_exactly_one_ack_on_the_last_message() {
let encoded = encoded_rebuild(2000);
assert!(
encoded.len() > 212_960,
"the 2000-mapping batch ({} bytes) must be past the default \
netlink send limit for this test to cover the large case",
encoded.len()
);
let headers = nl_headers(&encoded).expect("the batch parses");
assert_eq!(
headers.first().map(|h| h.kind),
Some(libc::NFNL_MSG_BATCH_BEGIN as u16)
);
assert_eq!(
headers.last().map(|h| h.kind),
Some(libc::NFNL_MSG_BATCH_END as u16)
);
let acked: Vec<usize> = headers
.iter()
.enumerate()
.filter(|(_, h)| h.flags & libc::NLM_F_ACK as u16 != 0)
.map(|(i, _)| i)
.collect();
assert_eq!(
acked,
vec![headers.len() - 2],
"only the last object before the batch end may request an ack; \
one ack per message overflows the receive buffer after the \
kernel has committed the batch"
);
}
#[test]
fn sndbuf_for_admits_the_2000_mapping_batch_and_small_batches_after_kernel_doubling() {
let large = encoded_rebuild(2000).len();
for len in [large, 0, 1, 212_961] {
let sndbuf = sndbuf_for(len);
assert!(
2 * sndbuf as u64 - 32 >= len as u64,
"a send buffer of {sndbuf}, doubled by the kernel, refuses a \
{len}-byte batch"
);
}
}
#[test]
fn sndbuf_for_saturates_at_the_kernel_clamp_for_huge_batches() {
for len in [2 * MAX_SNDBUF as usize, usize::MAX] {
assert_eq!(sndbuf_for(len), i32::MAX / 2, "sndbuf_for({len})");
}
}
#[test]
fn check_admissible_refuses_a_batch_larger_than_the_kernel_can_accept() {
assert!(check_admissible(KERNEL_BATCH_LIMIT).is_ok());
for len in [KERNEL_BATCH_LIMIT + 1, usize::MAX] {
match check_admissible(len) {
Err(e @ NatError::BatchTooLarge { bytes }) => {
assert_eq!(bytes, len);
assert!(
e.to_string().contains(&len.to_string()),
"the error names the batch size: {e}"
);
}
other => panic!("a {len}-byte batch was admitted: {other:?}"),
}
}
}
#[test]
fn ack_reader_fails_on_an_error_that_precedes_the_last_ack() {
// The kernel aborts the batch on a failing rule mid-batch, reports
// that rule's error, and still acknowledges the last message.
let reader = AckReader { last_seq: 4000 };
let error = ack(1234, -libc::ENOENT);
let last = ack(4000, 0);
let expect_error = |result: Result<AckState, NatError>| match result {
Err(NatError::Kernel { errno, seq }) => {
assert_eq!(errno, Errno(libc::ENOENT));
assert_eq!(seq, 1234);
}
other => panic!("the aborted batch was not reported: {other:?}"),
};
expect_error(reader.feed(&error));
expect_error(reader.feed(&[error.clone(), last.clone()].concat()));
}
#[test]
fn ack_reader_is_done_only_on_the_last_sequence_ack() {
let reader = AckReader { last_seq: 10 };
assert_eq!(reader.feed(&ack(10, 0)).expect("parses"), AckState::Done);
assert_eq!(reader.feed(&ack(5, 0)).expect("parses"), AckState::Pending);
assert_eq!(
reader
.feed(&nlmsg(libc::NLMSG_NOOP as u16, 10, &[]))
.expect("parses"),
AckState::Pending
);
assert_eq!(
reader
.feed(&[ack(5, 0), ack(10, 0)].concat())
.expect("parses"),
AckState::Done
);
let whole = ack(10, 0);
assert!(
reader.feed(&whole[..8]).is_err(),
"a header shorter than 16 bytes"
);
let mut overlong = whole.clone();
overlong[..4].copy_from_slice(&((whole.len() + 4) as u32).to_ne_bytes());
assert!(
reader.feed(&overlong).is_err(),
"a length past the end of the datagram"
);
let short = nlmsg(libc::NLMSG_ERROR as u16, 10, &[0, 0]);
assert!(
reader.feed(&short[..NLMSG_HDRLEN + 2]).is_err(),
"an NLMSG_ERROR payload shorter than its error field"
);
}
#[test]
fn kernel_error_display_names_the_errno() {
let text = NatError::Kernel {
errno: Errno(libc::EMSGSIZE),
seq: 7,
}
.to_string();
assert!(text.contains("EMSGSIZE"), "{text}");
assert!(text.contains(&format!("({})", libc::EMSGSIZE)), "{text}");
}
}
+229 -2
View File
@@ -10,6 +10,18 @@ use std::net::Ipv6Addr;
use std::time::Instant;
use tracing::{debug, info};
/// Most live mappings the pool holds before it refuses new names.
///
/// Every mapping adds rules to the NAT table, which is rebuilt whole on each
/// change, and work to every tick and to shutdown, so this bounds all three.
pub const MAPPING_CEILING: usize = 1000;
/// New mappings the pool admits in a burst, when idle long enough to refill.
pub const MAPPING_BURST: u32 = 50;
/// New mappings per second the pool admits once a burst is spent.
pub const MAPPING_RATE: u32 = 10;
/// Errors from pool operations.
#[derive(Debug, thiserror::Error)]
pub enum PoolError {
@@ -19,6 +31,10 @@ pub enum PoolError {
Exhausted(usize),
#[error("prefix length must be between 1 and 128")]
InvalidPrefix,
#[error("live-mapping ceiling reached ({0} mappings)")]
AtCeiling(usize),
#[error("new-mapping rate limit reached")]
RateLimited,
}
/// State of a virtual IP mapping.
@@ -217,6 +233,67 @@ impl ConntrackReadLog {
}
}
/// Token bucket for new mappings.
///
/// The level is kept in token-nanoseconds so refill is exact integer
/// arithmetic: one token is `NANOS` units, and each elapsed nanosecond adds
/// `rate` units.
#[derive(Debug)]
struct Bucket {
/// Current level, in units of `1 / NANOS` token.
level: u128,
/// Level when full.
capacity: u128,
/// Tokens added per second.
rate: u128,
/// When the level was last brought up to date; unset until first use.
last: Option<Instant>,
}
impl Bucket {
const NANOS: u128 = 1_000_000_000;
/// A full bucket of `capacity` tokens refilling at `rate` per second.
fn new(capacity: u32, rate: u32) -> Self {
let capacity = u128::from(capacity) * Self::NANOS;
Self {
level: capacity,
capacity,
rate: u128::from(rate),
last: None,
}
}
/// Add what has accrued since the last refill, up to capacity.
fn refill(&mut self, now: Instant) {
if let Some(last) = self.last {
let elapsed = now.saturating_duration_since(last).as_nanos();
self.level = self
.level
.saturating_add(elapsed.saturating_mul(self.rate))
.min(self.capacity);
}
// Never move backwards, so a stale `now` cannot credit time twice.
self.last = Some(self.last.map_or(now, |last| last.max(now)));
}
/// Whether at least one whole token is available.
fn has_token(&self) -> bool {
self.level >= Self::NANOS
}
/// Spend one token; the caller has checked `has_token`.
fn take(&mut self) {
self.level = self.level.saturating_sub(Self::NANOS);
}
/// Whole tokens available.
#[cfg(test)]
fn tokens(&self) -> u128 {
self.level / Self::NANOS
}
}
/// Virtual IP pool manager.
pub struct VirtualIpPool {
/// Available addresses (free pool).
@@ -231,11 +308,38 @@ pub struct VirtualIpPool {
grace_secs: u64,
/// Total pool size.
total: usize,
/// Most live mappings admitted before new names are refused.
ceiling: usize,
/// Rate limit on new mappings.
bucket: Bucket,
}
impl VirtualIpPool {
/// Create a new pool from a CIDR string (e.g., `fd01::/112`).
/// Create a new pool from a CIDR string (e.g., `fd01::/112`), with the
/// compiled-in admission limits.
pub fn new(cidr: &str, ttl_secs: u64, grace_secs: u64) -> Result<Self, PoolError> {
Self::with_limits(
cidr,
ttl_secs,
grace_secs,
MAPPING_CEILING,
MAPPING_BURST,
MAPPING_RATE,
)
}
/// Create a pool with explicit admission limits.
///
/// Production uses `new`; this exists so tests can set limits small
/// enough to reach without allocating the compiled-in counts.
pub fn with_limits(
cidr: &str,
ttl_secs: u64,
grace_secs: u64,
ceiling: usize,
burst: u32,
rate: u32,
) -> Result<Self, PoolError> {
let (base, prefix_len) = parse_ipv6_cidr(cidr)?;
if prefix_len == 0 || prefix_len > 128 {
return Err(PoolError::InvalidPrefix);
@@ -267,6 +371,8 @@ impl VirtualIpPool {
ttl_secs,
grace_secs,
total,
ceiling,
bucket: Bucket::new(burst, rate),
})
}
@@ -292,6 +398,21 @@ impl VirtualIpPool {
node_addr: NodeAddr,
mesh_addr: Ipv6Addr,
dns_name: &str,
) -> Result<(Ipv6Addr, bool), PoolError> {
self.allocate_at(node_addr, mesh_addr, dns_name, Instant::now())
}
/// `allocate` at a given instant, which drives the rate limit's refill
/// and stamps a new mapping.
///
/// An existing mapping is returned before either limit is consulted, so a
/// name already in use keeps resolving when new names are refused.
pub fn allocate_at(
&mut self,
node_addr: NodeAddr,
mesh_addr: Ipv6Addr,
dns_name: &str,
now: Instant,
) -> Result<(Ipv6Addr, bool), PoolError> {
// Idempotent: return existing mapping, refreshed.
if self.refresh_if_present(node_addr)
@@ -300,12 +421,21 @@ impl VirtualIpPool {
return Ok((mapping.virtual_ip, false));
}
// Ceiling first, so a refusal there costs no token and names the
// ceiling whatever the bucket holds.
if self.mappings.len() >= self.ceiling {
return Err(PoolError::AtCeiling(self.mappings.len()));
}
self.bucket.refill(now);
if !self.bucket.has_token() {
return Err(PoolError::RateLimited);
}
let virtual_ip = self
.available
.pop_front()
.ok_or(PoolError::Exhausted(self.mappings.len()))?;
self.bucket.take();
let now = Instant::now();
let mapping = VirtualIpMapping {
node_addr,
virtual_ip,
@@ -490,6 +620,7 @@ fn parse_ipv6_cidr(cidr: &str) -> Result<(Ipv6Addr, u32), PoolError> {
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
/// Session counts a test sets directly, handed to `tick` as the snapshot
/// the tick task would have read from conntrack.
@@ -589,6 +720,102 @@ mod tests {
);
}
/// A `/120` pool with the given limits, TTL and grace of 60 s.
fn limited_pool(ceiling: usize, burst: u32, rate: u32) -> VirtualIpPool {
VirtualIpPool::with_limits("fd01::/120", 60, 60, ceiling, burst, rate).unwrap()
}
/// Allocate node `i` at `now`.
fn alloc(pool: &mut VirtualIpPool, i: u8, now: Instant) -> Result<(Ipv6Addr, bool), PoolError> {
pool.allocate_at(make_node_addr(i), make_mesh_addr(i), "test.fips", now)
}
#[test]
fn ceiling_refuses_a_new_name_without_a_token_and_keeps_existing_names() {
let t0 = Instant::now();
let mut pool = limited_pool(3, 10, 1);
let mut vips = Vec::new();
for i in 1..=3u8 {
vips.push(alloc(&mut pool, i, t0).unwrap().0);
}
assert_eq!(pool.bucket.tokens(), 7);
assert!(
matches!(alloc(&mut pool, 4, t0), Err(PoolError::AtCeiling(3))),
"a fourth new name must be refused at a ceiling of 3"
);
assert_eq!(
pool.bucket.tokens(),
7,
"a ceiling refusal must not take a token"
);
assert_eq!(
alloc(&mut pool, 2, t0).unwrap(),
(vips[1], false),
"a name that already has a mapping must still resolve at the ceiling"
);
}
#[test]
fn ceiling_is_checked_before_the_rate_limit() {
let t0 = Instant::now();
// The bucket empties exactly as the ceiling is reached.
let mut pool = limited_pool(3, 3, 1);
for i in 1..=3u8 {
alloc(&mut pool, i, t0).unwrap();
}
assert_eq!(pool.bucket.tokens(), 0);
assert!(
matches!(alloc(&mut pool, 4, t0), Err(PoolError::AtCeiling(3))),
"a name refused at the ceiling must report the ceiling, not the rate"
);
}
#[test]
fn rate_limit_refuses_a_burst_keeps_existing_names_and_refills() {
let t0 = Instant::now();
let mut pool = limited_pool(100, 2, 1);
let (vip1, _) = alloc(&mut pool, 1, t0).unwrap();
alloc(&mut pool, 2, t0).unwrap();
assert!(
matches!(alloc(&mut pool, 3, t0), Err(PoolError::RateLimited)),
"a third new name at the same instant must be refused by a burst of 2"
);
assert_eq!(
alloc(&mut pool, 1, t0).unwrap(),
(vip1, false),
"an existing name must resolve with the bucket empty"
);
assert_eq!(pool.bucket.tokens(), 0);
assert!(
matches!(alloc(&mut pool, 3, t0), Err(PoolError::RateLimited)),
"resolving an existing name must not have freed a token"
);
let (_, is_new) = alloc(&mut pool, 3, t0 + Duration::from_secs(1)).unwrap();
assert!(is_new, "one refill interval later a new name must allocate");
}
#[test]
fn exhausted_pool_takes_no_token() {
let t0 = Instant::now();
// /126 = 3 usable addresses.
let mut pool = VirtualIpPool::with_limits("fd01::/126", 60, 60, 100, 10, 1).unwrap();
for i in 1..=3u8 {
alloc(&mut pool, i, t0).unwrap();
}
assert!(matches!(
alloc(&mut pool, 4, t0),
Err(PoolError::Exhausted(3))
));
assert_eq!(
pool.bucket.tokens(),
7,
"a refusal for an exhausted pool must not take a token"
);
}
#[test]
fn test_mapping_lifecycle_allocated_to_free() {
let mut pool = VirtualIpPool::new("fd01::/120", 1, 1).unwrap();
+3 -1
View File
@@ -423,7 +423,9 @@ services:
# attached after container start) and manage nftables NAT rules.
privileged: true
environment:
- RUST_LOG=info
# Debug for the NAT rebuild and pool tick timing lines only; the
# gateway's --log-level is ignored while RUST_LOG is set.
- RUST_LOG=info,fips::gateway::nat=debug,fips_gateway=debug
- FIPS_TEST_MODE=gateway
sysctls:
- net.ipv6.conf.all.disable_ipv6=0
+603
View File
@@ -477,6 +477,609 @@ else
check "Gateway shutdown (no completion message in logs)" 1
fi
# ── Long-lived gateway restart, shared by phases 11 and 12 ───────────────
# Start the stopped gateway with mappings that outlive the phase, and gate on
# its readiness. A failure is recorded through check under the label prefix
# $1 and returns 1, and the caller ends its phase. On success it sets:
# GW_STARTED the container's start time, for `docker logs --since`, since
# the log still holds every earlier phase
# GW_T0 epoch seconds at the start
# GW_BASELINE allocations after the readiness probe, which allocates
# GW_PROBE the address the readiness probe got for NPUB_B
gw_long_lived_start() {
local prefix="$1"
local config_file="$GENERATED_DIR/gateway/node-a.yaml"
local expect_rev
expect_rev=$(git -C "$SCRIPT_DIR" rev-parse --short=10 HEAD)
# Rewrite in place (same inode): the container sees the host file through
# a single-file bind mount, which a replace-by-rename would leave behind.
python3 - "$config_file" <<'PYEOF'
import sys, yaml
path = sys.argv[1]
with open(path, "r+") as f:
cfg = yaml.safe_load(f)
cfg["gateway"]["dns"]["ttl"] = 1800
cfg["gateway"]["pool_grace_period"] = 1800
f.seek(0)
yaml.dump(cfg, f, default_flow_style=False, sort_keys=False)
f.truncate()
PYEOF
docker start "$GATEWAY" >/dev/null
GW_STARTED=$(docker inspect -f '{{.State.StartedAt}}' "$GATEWAY")
GW_T0=$(date -u +%s)
echo " Gateway started at $GW_STARTED (expect rev $expect_rev)"
local seen_ttl seen_grace
seen_ttl=$(docker exec "$GATEWAY" grep -c "ttl: 1800" /etc/fips/fips.yaml || true)
seen_grace=$(docker exec "$GATEWAY" grep -c "pool_grace_period: 1800" /etc/fips/fips.yaml || true)
if [ "$seen_ttl" -ge 1 ] && [ "$seen_grace" -ge 1 ]; then
check "$prefix: container sees ttl 1800 and grace 1800" 0
else
check "$prefix: container config rewrite (ttl: $seen_ttl, grace: $seen_grace)" 1
return 1
fi
# Readiness is a hard gate here, unlike phases 1 and 2.
if wait_for_peers "$GATEWAY" 2 60; then
check "$prefix: gateway peers after restart" 0
else
check "$prefix: gateway peers after restart" 1
return 1
fi
local probe
GW_PROBE=""
for _ in $(seq 1 60); do
probe=$(docker exec "$CLIENT" dig +short AAAA "${NPUB_B}.fips" @${GW_DNS} 2>/dev/null || true)
GW_PROBE=$(grep -m1 "^fd01::" <<< "$probe" || true)
if [ -n "$GW_PROBE" ]; then
break
fi
sleep 1
done
if [ -n "$GW_PROBE" ]; then
check "$prefix: gateway DNS answers after restart" 0
else
check "$prefix: gateway DNS answers after restart" 1
return 1
fi
sleep 1
local started_log rev_lines
started_log=$(docker logs --timestamps --since "$GW_STARTED" "$GATEWAY" 2>&1)
# The co-resident daemon logs its own "(rev ...) starting" line, so
# anchor on the gateway's name.
rev_lines=$(grep -cE "fips-gateway [^ ]+ \(rev ${expect_rev}\) starting" <<< "$started_log" || true)
if [ "$rev_lines" -eq 1 ]; then
check "$prefix: startup line reads rev ${expect_rev}) with no -dirty" 0
else
check "$prefix: startup line for rev ${expect_rev} (found $rev_lines)" 1
return 1
fi
GW_BASELINE=$(grep -c "Allocated virtual IP" <<< "$started_log" || true)
return 0
}
# The pool's compiled-in admission limits. A new name is refused past
# MAPPING_CEILING live mappings, and past a burst of MAPPING_BURST new names
# are admitted at MAPPING_RATE per second, so phases that create many
# mappings retry the rate limit's refusals.
POOL_RS="$SCRIPT_DIR/../../../src/gateway/pool.rs"
# A compiled-in constant from pool.rs, or nothing if the line is not found.
pool_const() {
sed -nE "s/^pub const $1: [a-z0-9]+ = ([0-9]+);$/\1/p" "$POOL_RS"
}
POOL_CEILING=$(pool_const MAPPING_CEILING)
POOL_BURST=$(pool_const MAPPING_BURST)
POOL_RATE=$(pool_const MAPPING_RATE)
# Per-name retry bound for the driver, in seconds: a fixed 30 s. A bound
# sized from the rate (workers x refill interval x 5) is 2 s at 10/s, shorter
# than two of the driver's 1 to 1.5 s retry pauses, and never above 20 s for
# any rate of 1/s or more; 30 s lets a name wait through many pauses. No run
# has had a name reach it. Empty when the rate could not be read, which the
# phases report.
gw_retry_bound() {
[ -n "$POOL_RATE" ] && [ "$POOL_RATE" -gt 0 ] || return 0
echo 30
return 0
}
# Install the AAAA driver in the client: 4 closed-loop workers, one fresh
# socket per query, counts printed as key=value. Given a retry bound in
# seconds, a name answered SERVFAIL or not answered is asked again after a
# pause of 1 to 1.5 s until it is answered or the bound has passed since its
# first query; the exit status is 1 if any name was left unplaced. Without a
# bound each name is asked once and the exit status is 0.
gw_install_driver() {
docker exec -i "$CLIENT" sh -c 'cat > /tmp/gw_driver.py' <<'PYEOF'
import random, socket, struct, sys, threading, time
server = sys.argv[1]
bound = float(sys.argv[2]) if len(sys.argv) > 2 else None
names = [n.strip() for n in sys.stdin if n.strip()]
lock = threading.Lock()
counts = {"answered": 0, "servfail": 0, "timeout": 0, "other": 0,
"retries": 0, "unplaced": 0}
def query(name):
qid = random.getrandbits(16)
pkt = struct.pack(">HHHHHH", qid, 0x0100, 1, 0, 0, 0)
for label in (name + ".fips").split("."):
raw = label.encode()
pkt += bytes([len(raw)]) + raw
pkt += b"\x00" + struct.pack(">HH", 28, 1)
s = socket.socket(socket.AF_INET6, socket.SOCK_DGRAM)
s.settimeout(6)
try:
s.sendto(pkt, (server, 53))
while True:
data, _ = s.recvfrom(4096)
if len(data) >= 12 and struct.unpack(">H", data[:2])[0] == qid:
break
except socket.timeout:
return "timeout"
finally:
s.close()
flags, _, ancount = struct.unpack(">HHH", data[2:8])
rcode = flags & 0xF
if rcode == 2:
return "servfail"
if rcode == 0 and ancount > 0:
return "answered"
return "other"
def place(name):
first = time.monotonic()
while True:
outcome = query(name)
with lock:
counts[outcome] += 1
if bound is None or outcome == "answered":
return
if outcome == "other" or time.monotonic() - first > bound:
with lock:
counts["unplaced"] += 1
return
with lock:
counts["retries"] += 1
time.sleep(1 + random.random() * 0.5)
def worker():
while True:
with lock:
if not names:
return
name = names.pop()
place(name)
threads = [threading.Thread(target=worker) for _ in range(4)]
for t in threads:
t.start()
for t in threads:
t.join()
print(" ".join(f"{k}={v}" for k, v in counts.items()))
sys.exit(1 if counts["unplaced"] else 0)
PYEOF
}
# Phase 11: NAT rebuild past the default netlink socket limits
#
# Every change rebuilds the whole fips_gateway table in one netlink batch.
# With the default socket buffers that batch failed from about 105 mappings
# (the acks overflowed the receive buffer, after the commit) and past about
# 313 (the batch overflowed the send buffer, and nothing was committed).
# Drive 400 new names through a gateway whose mappings outlive the phase and
# judge the result on the kernel's own table, which is what a lost rebuild
# leaves wrong, not on the daemon's debug-level success line. The names go
# through the pool's rate limit, so the driver retries its refusals. Runs
# after phase 10, so the gateway container is stopped when it starts, and it
# leaves it stopped.
echo ""
echo "Phase 11: NAT rebuild past default socket limits"
NATBIG_NAMES=400
NATBIG_CAP=180
NATBIG_SETTLE=30
natbig_now() {
date -u +%s
}
natbig_slice() {
NATBIG_LOG=$(docker logs --timestamps --since "$NATBIG_STARTED" "$GATEWAY" 2>&1)
}
natbig_allocated() {
natbig_slice
NATBIG_ALLOCATED=$(grep -c "Allocated virtual IP" <<< "$NATBIG_LOG" || true)
}
# Rules the kernel holds right now. A failed listing is recorded through
# NATBIG_RC, never read as a table with no rules.
natbig_kernel() {
NATBIG_RC=0
NATBIG_NFT=$(docker exec "$GATEWAY" nft list table inet fips_gateway 2>&1) || NATBIG_RC=$?
NATBIG_DNAT=$(grep -cE "daddr fd01:[0-9a-f:]* .*dnat" <<< "$NATBIG_NFT" || true)
NATBIG_SNAT=$(grep -cE "saddr [0-9a-f:]+ .*snat" <<< "$NATBIG_NFT" || true)
NATBIG_MASQ=$(grep -c "masquerade" <<< "$NATBIG_NFT" || true)
}
natbig_phase() {
gw_long_lived_start "NAT batch" || return 0
NATBIG_STARTED="$GW_STARTED"
local t0="$GW_T0"
local baseline="$GW_BASELINE"
local target=$((baseline + NATBIG_NAMES))
echo " Baseline allocations after readiness: $baseline; target $target"
if [ "$baseline" -ge 1 ]; then
check "NAT batch: readiness probe allocated (baseline $baseline)" 0
else
check "NAT batch: readiness probe allocated (baseline $baseline)" 1
return 0
fi
# Names: real keys, since the daemon parses each one as a public key.
local names_file have_names
names_file=$(mktemp)
docker exec "$GATEWAY" bash -c \
"for i in \$(seq 1 $NATBIG_NAMES); do fipsctl keygen --stdout; done" \
| grep '^npub1' >"$names_file" || true
have_names=$(wc -l <"$names_file")
if [ "$have_names" -lt "$NATBIG_NAMES" ]; then
check "NAT batch: generated $NATBIG_NAMES names (got $have_names)" 1
rm -f "$names_file"
return 0
fi
local retry_bound
retry_bound=$(gw_retry_bound)
if [ -z "$retry_bound" ]; then
check "NAT batch: MAPPING_RATE read from pool.rs ('$POOL_RATE')" 1
rm -f "$names_file"
return 0
fi
gw_install_driver
local remaining out rc=0
remaining=$((NATBIG_CAP - ($(natbig_now) - t0)))
# timeout reads 0 as no limit and refuses a negative duration.
if [ "$remaining" -le 0 ]; then
check "NAT batch: setup exceeded ${NATBIG_CAP}s cap" 1
rm -f "$names_file"
return 0
fi
out=$(docker exec -i "$CLIENT" timeout "$remaining" \
python3 /tmp/gw_driver.py "$GW_DNS" "$retry_bound" <"$names_file" 2>&1) || rc=$?
rm -f "$names_file"
echo " [$(($(natbig_now) - t0))s] sent $NATBIG_NAMES names: $out (rc=$rc)"
natbig_allocated
if [ "$rc" -eq 0 ] && [ "$NATBIG_ALLOCATED" -eq "$target" ]; then
check "NAT batch: $NATBIG_ALLOCATED live mappings allocated" 0
else
check "NAT batch: live mappings allocated ($NATBIG_ALLOCATED of $target, rc $rc)" 1
fi
# Settle on the kernel: wait until it holds a DNAT rule per allocation,
# or the cap expires. Reaching the cap decides nothing by itself; the
# checks below do.
local settle=0
natbig_kernel
while [ "$NATBIG_DNAT" -ne "$NATBIG_ALLOCATED" ] && [ "$settle" -lt "$NATBIG_SETTLE" ]; do
sleep 1
settle=$((settle + 1))
natbig_kernel
done
# An error logged just after the last commit still counts.
sleep 2
natbig_kernel
natbig_allocated
local nat_fail
nat_fail=$(grep -c "Failed to add NAT rules" <<< "$NATBIG_LOG" || true)
echo " [$(($(natbig_now) - t0))s] allocated=$NATBIG_ALLOCATED nat_add_fail=$nat_fail" \
"table_rc=$NATBIG_RC dnat=$NATBIG_DNAT snat=$NATBIG_SNAT masquerade=$NATBIG_MASQ settle=${settle}s"
if [ "$nat_fail" -gt 0 ]; then
grep "Failed to add NAT rules" <<< "$NATBIG_LOG" | sed 's/\x1b\[[0-9;]*m//g' \
| sed -n '1p;$p' | sed 's/^/ /'
fi
if [ "$nat_fail" -eq 0 ]; then
check "NAT batch: no NAT rebuild failed" 0
else
check "NAT batch: NAT rebuilds failed ($nat_fail)" 1
fi
if [ "$NATBIG_RC" -eq 0 ]; then
check "NAT batch: nft lists the fips_gateway table" 0
else
check "NAT batch: nft list table failed (rc $NATBIG_RC)" 1
fi
if [ "$NATBIG_DNAT" -eq "$NATBIG_ALLOCATED" ] && [ "$NATBIG_SNAT" -eq "$NATBIG_ALLOCATED" ]; then
check "NAT batch: kernel holds a DNAT and SNAT rule per mapping ($NATBIG_DNAT)" 0
else
check "NAT batch: kernel rules (dnat $NATBIG_DNAT, snat $NATBIG_SNAT, allocated $NATBIG_ALLOCATED)" 1
fi
if [ "$NATBIG_MASQ" -eq 2 ]; then
check "NAT batch: fips0 and LAN masquerades present" 0
else
check "NAT batch: masquerade rules ($NATBIG_MASQ, expected 2)" 1
fi
docker stop --time=10 "$GATEWAY" >/dev/null 2>&1 || true
echo " Phase time: $(($(natbig_now) - t0))s"
}
natbig_phase
# Phase 12: Pool admission limits
#
# The pool refuses a new name once it holds MAPPING_CEILING live mappings,
# and past a burst of MAPPING_BURST admits new names at MAPPING_RATE per
# second; the constants are read from src/gateway/pool.rs. Restart the gateway
# with mappings that outlive the phase, fill the pool to the ceiling through
# the rate limit, then ask once each for 20 more new names: all 20 must be
# refused at the ceiling while an existing name still resolves. Without the
# ceiling those 20 are allocated, whatever the creation rate. Reports rebuild
# and tick durations on the way up and the shutdown duration at the ceiling.
# Runs after phase 11, so the gateway container is stopped when it starts,
# and it leaves it stopped.
echo ""
echo "Phase 12: Pool admission limits"
LIMITS_EXTRA=20
# Sized from both trees: with the limits the fill takes about
# (ceiling - burst) / rate seconds plus retry pauses, and without them a
# measurement run reached ceiling + 20 mappings far sooner.
LIMITS_CAP=300
# Shutdown took 2.4 s at 2000 mappings in a measurement run.
LIMITS_STOP_TIME=30
limits_now() {
date -u +%s
}
limits_slice() {
SLICE=$(docker logs --timestamps --since "$STARTED" "$GATEWAY" 2>&1)
}
# Figures below read lines a grep on the slice has already selected, so the
# log-string guard sees each daemon-log literal. The Python matches no log
# text itself: it takes the daemon's own timestamp and key=value fields.
LIMITS_PY_FIELDS='
import re, sys, statistics
from datetime import datetime, timezone
ANSI = re.compile(r"\x1b\[[0-9;]*m")
FIELD = re.compile(r"\b(\w+)=(\S+)")
def parse(line):
parts = ANSI.sub("", line.rstrip("\n")).split(" ", 2)
whole, _, frac = parts[1].rstrip("Z").partition(".")
t = datetime.strptime(whole, "%Y-%m-%dT%H:%M:%S").replace(tzinfo=timezone.utc)
t = t.timestamp() + float("0." + (frac or "0"))
return t, dict(FIELD.findall(parts[2] if len(parts) > 2 else ""))
'
limits_phase() {
local ceiling="$POOL_CEILING" burst="$POOL_BURST" rate="$POOL_RATE"
if [ -n "$ceiling" ] && [ -n "$burst" ] && [ -n "$rate" ] && [ "$rate" -gt 0 ] \
&& [ "$ceiling" -gt 1 ]; then
check "Limits: read ceiling $ceiling, burst $burst, rate $rate/s from pool.rs" 0
else
check "Limits: constants in pool.rs (ceiling '$ceiling', burst '$burst', rate '$rate')" 1
return 0
fi
local retry_bound
retry_bound=$(gw_retry_bound)
gw_long_lived_start "Limits" || return 0
STARTED="$GW_STARTED"
local t0="$GW_T0"
local baseline="$GW_BASELINE"
if [ "$baseline" -eq 1 ]; then
check "Limits: the readiness probe holds the only mapping" 0
else
check "Limits: mappings after readiness ($baseline, expected 1)" 1
return 0
fi
# Names: real keys, since the daemon parses each one as a public key.
local fill=$((ceiling - baseline))
local total=$((fill + LIMITS_EXTRA))
local names_file have_names
names_file=$(mktemp)
docker exec "$GATEWAY" bash -c \
"for i in \$(seq 1 $total); do fipsctl keygen --stdout; done" \
| grep '^npub1' >"$names_file" || true
have_names=$(wc -l <"$names_file")
if [ "$have_names" -lt "$total" ]; then
check "Limits: generated $total names (got $have_names)" 1
rm -f "$names_file"
return 0
fi
gw_install_driver
# Fill: exactly ceiling - 1 new names, each retried through the rate
# limit's refusals. A name left unplaced, or the cap firing, fails.
local remaining out rc=0
remaining=$((LIMITS_CAP - ($(limits_now) - t0)))
if [ "$remaining" -le 0 ]; then
check "Limits: setup exceeded ${LIMITS_CAP}s cap" 1
rm -f "$names_file"
return 0
fi
out=$(sed -n "1,${fill}p" "$names_file" | docker exec -i "$CLIENT" timeout "$remaining" \
python3 /tmp/gw_driver.py "$GW_DNS" "$retry_bound" 2>&1) || rc=$?
echo " [$(($(limits_now) - t0))s] fill of $fill names (retry bound ${retry_bound}s): $out (rc=$rc)"
if [ "$rc" -eq 0 ]; then
check "Limits: fill placed all $fill names" 0
else
check "Limits: fill driver failed (rc $rc)" 1
fi
sleep 1
limits_slice
local rate_refused_fill
rate_refused_fill=$(grep -c "new-mapping rate limit reached" <<< "$SLICE" || true)
# Past the ceiling: one query per name, never retried.
rc=0
remaining=$((LIMITS_CAP - ($(limits_now) - t0)))
if [ "$remaining" -le 0 ]; then
check "Limits: fill exceeded ${LIMITS_CAP}s cap" 1
rm -f "$names_file"
return 0
fi
out=$(sed -n "$((fill + 1)),${total}p" "$names_file" | docker exec -i "$CLIENT" \
timeout "$remaining" python3 /tmp/gw_driver.py "$GW_DNS" 2>&1) || rc=$?
rm -f "$names_file"
echo " [$(($(limits_now) - t0))s] $LIMITS_EXTRA names past the ceiling: $out (rc=$rc)"
local post_servfail
post_servfail=$(sed -nE 's/.*servfail=([0-9]+).*/\1/p' <<< "$out")
if [ "$rc" -eq 0 ] && [ "${post_servfail:-0}" -eq "$LIMITS_EXTRA" ]; then
check "Limits: all $LIMITS_EXTRA names past the ceiling got SERVFAIL" 0
else
check "Limits: names past the ceiling (servfail '${post_servfail}', rc $rc)" 1
fi
local probe
probe=$(docker exec "$CLIENT" dig +short AAAA "${NPUB_B}.fips" @${GW_DNS} 2>/dev/null || true)
if [ "$(grep -m1 "^fd01::" <<< "$probe" || true)" = "$GW_PROBE" ]; then
check "Limits: NPUB_B still resolves to $GW_PROBE at the ceiling" 0
else
check "Limits: NPUB_B at the ceiling (got '${probe:0:60}', had $GW_PROBE)" 1
fi
sleep 1
limits_slice
local allocated reclaimed ceiling_refused rate_refused nat_fail nat_rm_fail ndp_fail
allocated=$(grep -c "Allocated virtual IP" <<< "$SLICE" || true)
reclaimed=$(grep -c "Reclaimed virtual IP" <<< "$SLICE" || true)
ceiling_refused=$(grep -c "live-mapping ceiling reached" <<< "$SLICE" || true)
rate_refused=$(grep -c "new-mapping rate limit reached" <<< "$SLICE" || true)
nat_fail=$(grep -c "Failed to add NAT rules" <<< "$SLICE" || true)
nat_rm_fail=$(grep -c "Failed to remove NAT rules" <<< "$SLICE" || true)
ndp_fail=$(grep -c "Failed to add proxy NDP" <<< "$SLICE" || true)
echo " Slice counts: allocated=$allocated reclaimed=$reclaimed" \
"ceiling_refused=$ceiling_refused rate_refused=$rate_refused_fill/$rate_refused" \
"nat_add_fail=$nat_fail nat_remove_fail=$nat_rm_fail ndp_fail=$ndp_fail"
if [ "$allocated" -eq "$ceiling" ]; then
check "Limits: live mappings stop at the ceiling ($allocated)" 0
else
check "Limits: live mappings $allocated, ceiling $ceiling" 1
fi
if [ "$allocated" -gt 0 ]; then
if [ "$reclaimed" -eq 0 ]; then
check "Limits: no reclaims during the phase" 0
else
check "Limits: reclaims during the phase ($reclaimed)" 1
fi
else
check "Limits: allocations present in the slice (0)" 1
fi
if [ "$ceiling_refused" -ge "$LIMITS_EXTRA" ]; then
check "Limits: refusals logged as the ceiling ($ceiling_refused)" 0
else
check "Limits: ceiling refusals logged ($ceiling_refused, expected >= $LIMITS_EXTRA)" 1
fi
# The rate limit must have refused during the fill, so its count is a
# live signal, and must refuse nothing after it: past the ceiling the
# ceiling is checked first.
if [ "$rate_refused_fill" -ge 1 ]; then
check "Limits: the rate limit refused during the fill ($rate_refused_fill)" 0
else
check "Limits: rate refusals during the fill (0)" 1
fi
if [ "$rate_refused" -eq "$rate_refused_fill" ]; then
check "Limits: no rate refusal after the fill" 0
else
check "Limits: rate refusals after the fill ($((rate_refused - rate_refused_fill)))" 1
fi
if [ $((nat_fail + nat_rm_fail + ndp_fail)) -eq 0 ]; then
check "Limits: no NAT or proxy NDP failure lines" 0
else
check "Limits: failure lines (nat add $nat_fail, remove $nat_rm_fail, ndp $ndp_fail)" 1
fi
# Creations in any 10 s window of the fill can be at most a full burst
# plus 10 s of refill. The bound is exact, not approximate: the bucket
# holds at most burst whole tokens at the window's first creation and
# gains one per 1/rate s, so one more creation needs the daemon's log
# timestamps to lag its clock by a full token interval (100 ms at 10/s)
# more at one end of the window than the other. Do not loosen it for skew.
local bound=$((burst + 10 * rate)) most
most=$(grep "Allocated virtual IP" <<< "$SLICE" | python3 -c "$LIMITS_PY_FIELDS
b, fill = int(sys.argv[1]), int(sys.argv[2])
ts = sorted(parse(l)[0] for l in sys.stdin)[b:b + fill]
most = j = 0
for i in range(len(ts)):
while ts[i] - ts[j] > 10:
j += 1
most = max(most, i - j + 1)
print(most)
" "$baseline" "$fill" || true)
if [ -n "$most" ] && [ "$most" -le "$bound" ]; then
check "Limits: at most $most creations in any 10 s of the fill (bound $bound)" 0
else
check "Limits: creations in a 10 s window of the fill ('$most', bound $bound)" 1
fi
local added_timed tick_timed
added_timed=$(grep "Added DNAT/SNAT rules" <<< "$SLICE" | grep -c "elapsed_us" || true)
tick_timed=$(grep "Pool tick" <<< "$SLICE" | grep -c "tick_us" || true)
if [ "$added_timed" -ge 1 ] && [ "$tick_timed" -ge 1 ]; then
check "Limits: rebuild and tick timing lines present" 0
else
check "Limits: timing lines (rebuild $added_timed, tick $tick_timed)" 1
fi
echo " --- Rebuild duration (elapsed_us over the 20 adds ending at each count) ---"
local targets=("$ceiling") rebuild_report
[ "$ceiling" -gt 500 ] && targets=(500 "$ceiling")
rebuild_report=$(grep "Added DNAT/SNAT rules" <<< "$SLICE" | python3 -c "$LIMITS_PY_FIELDS
by_n = {}
for l in sys.stdin:
_, f = parse(l)
if 'error' in f or 'elapsed_us' not in f:
continue
by_n[int(f['mappings'])] = int(f['elapsed_us'])
missing = 0
for t in map(int, sys.argv[1:]):
if t not in by_n:
missing += 1
print(f' rebuild {t}: no successful add line with mappings={t}')
continue
xs = [by_n[n] for n in range(t - 19, t + 1) if n in by_n]
print(f' rebuild {t}: n={len(xs)} median={statistics.median(xs):.0f}us max={max(xs)}us')
print(f'REBUILD_MISSING={missing}')
" "${targets[@]}" || true)
echo "$rebuild_report" | grep -v '^REBUILD_MISSING='
if echo "$rebuild_report" | grep -q '^REBUILD_MISSING=0$'; then
check "Limits: a successful add line at ${targets[*]} mappings" 0
else
check "Limits: a successful add line at ${targets[*]} mappings" 1
fi
echo " --- Ticks as they fell ---"
grep "Pool tick" <<< "$SLICE" | python3 -c "$LIMITS_PY_FIELDS
for l in sys.stdin:
_, f = parse(l)
print(f\" tick: mappings={f.get('mappings')} read_us={f.get('read_us')} tick_us={f.get('tick_us')}\")
" || true
echo " --- Shutdown at the ceiling ---"
docker stop --time="$LIMITS_STOP_TIME" "$GATEWAY" >/dev/null 2>&1 || true
limits_slice
local shut_lines
shut_lines=$( { grep "fips-gateway shutting down" <<< "$SLICE"; \
grep "fips-gateway shutdown complete" <<< "$SLICE"; } || true)
if [ "$(echo "$shut_lines" | grep -c . || true)" -eq 2 ]; then
echo "$shut_lines" | python3 -c "$LIMITS_PY_FIELDS
ts = [parse(l)[0] for l in sys.stdin]
print(f' shutdown duration: {ts[1] - ts[0]:.3f}s')
" || true
check "Limits: gateway shutdown completed within ${LIMITS_STOP_TIME}s" 0
else
check "Limits: gateway shutdown completion line present" 1
fi
echo " Phase time: $(($(limits_now) - t0))s"
}
limits_phase
echo ""
echo "=== Results: $PASSED passed, $FAILED failed ==="
[ "$FAILED" -eq 0 ] && exit 0 || exit 1