From 64b0dce57b3d081af3baadc9ce1917cea11eb85a Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Fri, 18 Sep 2026 17:33:14 +0000 Subject: [PATCH 1/2] fix(gateway): size the NAT rebuild's netlink socket to the batch Every NAT rebuild is one netlink batch, sent through the default socket that rustables opens, and every message in it requested an ack. From about 105 live mappings the acks overflowed the 208 KiB receive buffer and the read failed with ENOBUFS, although the kernel had already committed the batch, so those rebuilds were logged as failed while they had taken effect. Past about 313 mappings the batch itself exceeded the send buffer, the send failed with EMSGSIZE and nothing was committed, so new .fips names got a virtual IP with no translation. Releases with the gateway through 0.5.1 sent the table delete in a batch of its own, so there a rebuild past about 313 mappings also removed the whole table. Neither errno reached the log, because the error kept only the outer message. The rebuild now encodes the batch and sends it itself. SO_SNDBUFFORCE is sized to the batch, capped at the kernel's own clamp, and a batch no send buffer can hold is refused up front. Only the last object requests an ack; an ack reader fails on any error that arrives before that ack, since the kernel still acknowledges the last message of a batch it aborted, and SO_RCVTIMEO bounds the wait. NAT errors now name the errno. The batch is still one transaction, so the table never leaves the packet path, and a failed rebuild is still only logged: the next successful rebuild installs the mapping. Unit tests cover the encoding at 2000 mappings, the send-buffer sizing and its saturation, the admission limit and the ack reader. A new gateway suite phase floods 400 names, past both old thresholds, and judges the result on the kernel's table and on the absence of NAT failure lines. --- CHANGELOG.md | 12 + src/gateway/nat.rs | 860 +++++++++++++++++++++---- testing/static/scripts/gateway-test.sh | 250 +++++++ 3 files changed, 997 insertions(+), 125 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 4b09a88c..e843a567 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -147,6 +147,18 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 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. - 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 diff --git a/src/gateway/nat.rs b/src/gateway/nat.rs index f0c11d60..1dade01b 100644 --- a/src/gateway/nat.rs +++ b/src/gateway/nat.rs @@ -4,7 +4,9 @@ //! 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 tracing::{debug, info}; use rustables::expr::{ @@ -23,6 +25,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 +99,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 for NatError { @@ -54,8 +135,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)] @@ -188,7 +269,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 +319,137 @@ impl NatManager { vec![ops] } - /// Build each op into its rustables object and send the batches in order. - fn send_batches(&self, batches: &[Vec]) -> 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, 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 +462,363 @@ 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(()) } } +/// 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, 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 { + 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::() 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 { + 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::() 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::() 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 +935,176 @@ mod tests { 1 ); } + + /// The encoded rebuild of a manager holding `count` mappings. + fn encoded_rebuild(count: u16) -> Vec { + 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 { + 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 { + 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 = 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| 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}"); + } } diff --git a/testing/static/scripts/gateway-test.sh b/testing/static/scripts/gateway-test.sh index 9cb8a605..331037ea 100755 --- a/testing/static/scripts/gateway-test.sh +++ b/testing/static/scripts/gateway-test.sh @@ -477,6 +477,256 @@ else check "Gateway shutdown (no completion message in logs)" 1 fi +# 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, not on the daemon's debug-level +# success line, which the suite's log level does not show. 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() { + 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 + NATBIG_STARTED=$(docker inspect -f '{{.State.StartedAt}}' "$GATEWAY") + local t0 + t0=$(natbig_now) + echo " Gateway started at $NATBIG_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 "NAT batch: container sees ttl 1800 and grace 1800" 0 + else + check "NAT batch: container config rewrite (ttl: $seen_ttl, grace: $seen_grace)" 1 + return 0 + fi + + if wait_for_peers "$GATEWAY" 2 60; then + check "NAT batch: gateway peers after restart" 0 + else + check "NAT batch: gateway peers after restart" 1 + return 0 + fi + local ready=false probe + for _ in $(seq 1 60); do + probe=$(docker exec "$CLIENT" dig +short AAAA "${NPUB_B}.fips" @${GW_DNS} 2>/dev/null || true) + if echo "$probe" | grep -q "^fd01::"; then + ready=true + break + fi + sleep 1 + done + if [ "$ready" = true ]; then + check "NAT batch: gateway DNS answers after restart" 0 + else + check "NAT batch: gateway DNS answers after restart" 1 + return 0 + fi + + sleep 1 + natbig_allocated + local rev_lines + rev_lines=$(grep -cE "fips-gateway [^ ]+ \(rev ${expect_rev}\) starting" <<< "$NATBIG_LOG" || true) + if [ "$rev_lines" -eq 1 ]; then + check "NAT batch: startup line reads rev ${expect_rev}) with no -dirty" 0 + else + check "NAT batch: startup line for rev ${expect_rev} (found $rev_lines)" 1 + return 0 + fi + local baseline="$NATBIG_ALLOCATED" + 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 + + # Closed-loop AAAA driver, 4 workers, one fresh socket per query. + docker exec -i "$CLIENT" sh -c 'cat > /tmp/gw_natbig.py' <<'PYEOF' +import random, socket, struct, sys, threading +server = sys.argv[1] +names = [n.strip() for n in sys.stdin if n.strip()] +lock = threading.Lock() +counts = {"answered": 0, "servfail": 0, "timeout": 0, "other": 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 worker(): + while True: + with lock: + if not names: + return + name = names.pop() + outcome = query(name) + with lock: + counts[outcome] += 1 +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())) +PYEOF + + 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_natbig.py "$GW_DNS" <"$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 + echo "" echo "=== Results: $PASSED passed, $FAILED failed ===" [ "$FAILED" -eq 0 ] && exit 0 || exit 1 From 1e111fe044ea582bc56451ecb85e92c7a7754acf Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Fri, 18 Sep 2026 19:41:12 +0000 Subject: [PATCH 2/2] fix(gateway): bound live mappings and the new-mapping rate Any host that can reach the LAN resolver could ask for one new .fips name after another, and each got a virtual-IP mapping until the pool's 65,535 addresses ran out. Every mapping also adds to the cost of each NAT rebuild, each pool tick and shutdown. The pool now refuses a new name once it holds 1000 live mappings, and admits new names from a token bucket of 50 that refills at 10 per second. Both checks sit after the return for a name that already has a mapping, so names in use keep resolving when new ones are refused. The ceiling is checked first, so a refusal there takes no token and is always reported as the ceiling. A token is taken only once an address has been taken from the free list, so an exhausted pool costs none. Each limit has its own PoolError variant, and the "Pool allocation failed" warning names the one that refused. VirtualIpPool::new keeps its signature and uses the compiled-in limits; with_limits and allocate_at, which takes the instant that drives the refill, let the unit tests set small limits and a clock. The limits come from a flood of the gateway suite's gateway to 500, 1000 and 2000 live mappings with no limits in place. The measurement record is kept separately. Figures: - New mappings per second, over 100 creations: 215 for the first 100, 63 at 500, 47 at 1000, 32 at 2000. - NAT rebuild, median/max over the 20 adds up to the count: 17.5/26.5 ms at 500, 22.2/34.3 ms at 1000, 32.2/38.3 ms at 2000. - Pool tick: 24 to 32 us at 500, 48 to 53 us at 1000, 105 to 115 us at 2000, beside a conntrack read of 180 to 300 us. - Shutdown at 2000 mappings: 2.4 s. These are a best case for router hardware: they come from a container on a development host, not a router, and that host has no /proc/net/nf_conntrack, so the tick figures include no conntrack parsing. The ceiling is half the largest count measured. At 1000 a rebuild took about 22 ms and shutdown about 1.2 s. The rate is below the unthrottled creation rate at every count measured, so the bucket rather than the rebuild sets how fast a flood can fill the pool. From empty that now takes about 95 s; the unthrottled run passed 1000 in about 16 s. Ten rebuilds a second at the ceiling cost about a fifth of the gateway's single runtime thread on that host. The 2000 figures exist only in the measurement, since the committed phase stops at the ceiling. To measure and to keep measuring, the "Added DNAT/SNAT rules" and "Removed DNAT/SNAT rules" debug lines now carry the mapping count after the change and the rebuild's duration in microseconds, and are also emitted, with the error, when a rebuild fails, so a failure at some count leaves a record of that count. The pool tick logs a debug line with the mapping count and the durations of the conntrack read and of the tick. The suite's gateway logs at info because RUST_LOG=info overrides the entrypoint's --log-level debug, so the gw-gateway service now enables debug for the NAT manager and the gateway binary only. The gateway suite's last phase is now the regression. It restarts the gateway with mappings that outlive the phase and reads the limits from pool.rs. It fills the pool to the ceiling with the readiness probe's mapping plus ceiling - 1 new names, retrying rate refusals, then asks once each for 20 more. It asserts all 20 get SERVFAIL, the live count equals the ceiling, at least 20 ceiling refusals are logged, the rate limit refused during the fill and not after it, the probe's name keeps its address, nothing is reclaimed, no NAT or proxy NDP failure is logged, and no 10 s window of the fill holds more than burst + 10 x rate creations. It also reports rebuild durations at 500 and 1000, ticks as they fall, and the shutdown time at the ceiling. On a tree with these assertions but no checks in allocate the phase fails five of them: 1020 mappings, 0 SERVFAIL, no ceiling or rate refusals, and 713 creations in one 10 s window. With the limits it passes in about 140 s, with rebuilds of 17.8/24.5 ms at 500 and 20.9/33.9 ms at 1000 and a 1.2 s shutdown at the ceiling. The NAT batch phase now also goes through the rate limit, so its driver retries refusals; it shares the restart, readiness gate and DNS driver with the new phase instead of carrying its own copies. --- CHANGELOG.md | 10 + src/bin/fips-gateway.rs | 11 + src/gateway/nat.rs | 60 ++- src/gateway/pool.rs | 231 +++++++++- testing/static/docker-compose.yml | 4 +- testing/static/scripts/gateway-test.sh | 597 ++++++++++++++++++++----- 6 files changed, 776 insertions(+), 137 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index e843a567..516c3101 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -159,6 +159,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 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 diff --git a/src/bin/fips-gateway.rs b/src/bin/fips-gateway.rs index 73b8743a..def8029e 100644 --- a/src/bin/fips-gateway.rs +++ b/src/bin/fips-gateway.rs @@ -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, diff --git a/src/gateway/nat.rs b/src/gateway/nat.rs index 1dade01b..7a0d2770 100644 --- a/src/gateway/nat.rs +++ b/src/gateway/nat.rs @@ -7,6 +7,7 @@ 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::{ @@ -242,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. @@ -257,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. @@ -467,6 +494,15 @@ impl NatManager { } 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. diff --git a/src/gateway/pool.rs b/src/gateway/pool.rs index 900398d4..7776b160 100644 --- a/src/gateway/pool.rs +++ b/src/gateway/pool.rs @@ -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, +} + +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::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 { 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(); diff --git a/testing/static/docker-compose.yml b/testing/static/docker-compose.yml index 77010007..eb025aee 100644 --- a/testing/static/docker-compose.yml +++ b/testing/static/docker-compose.yml @@ -375,7 +375,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 diff --git a/testing/static/scripts/gateway-test.sh b/testing/static/scripts/gateway-test.sh index 331037ea..3153bf71 100755 --- a/testing/static/scripts/gateway-test.sh +++ b/testing/static/scripts/gateway-test.sh @@ -477,6 +477,191 @@ 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. @@ -484,10 +669,11 @@ fi # (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, not on the daemon's debug-level -# success line, which the suite's log level does not show. Runs after phase -# 10, so the gateway container is stopped when it starts, and it leaves it -# stopped. +# 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" @@ -519,73 +705,10 @@ natbig_kernel() { } natbig_phase() { - 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 - NATBIG_STARTED=$(docker inspect -f '{{.State.StartedAt}}' "$GATEWAY") - local t0 - t0=$(natbig_now) - echo " Gateway started at $NATBIG_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 "NAT batch: container sees ttl 1800 and grace 1800" 0 - else - check "NAT batch: container config rewrite (ttl: $seen_ttl, grace: $seen_grace)" 1 - return 0 - fi - - if wait_for_peers "$GATEWAY" 2 60; then - check "NAT batch: gateway peers after restart" 0 - else - check "NAT batch: gateway peers after restart" 1 - return 0 - fi - local ready=false probe - for _ in $(seq 1 60); do - probe=$(docker exec "$CLIENT" dig +short AAAA "${NPUB_B}.fips" @${GW_DNS} 2>/dev/null || true) - if echo "$probe" | grep -q "^fd01::"; then - ready=true - break - fi - sleep 1 - done - if [ "$ready" = true ]; then - check "NAT batch: gateway DNS answers after restart" 0 - else - check "NAT batch: gateway DNS answers after restart" 1 - return 0 - fi - - sleep 1 - natbig_allocated - local rev_lines - rev_lines=$(grep -cE "fips-gateway [^ ]+ \(rev ${expect_rev}\) starting" <<< "$NATBIG_LOG" || true) - if [ "$rev_lines" -eq 1 ]; then - check "NAT batch: startup line reads rev ${expect_rev}) with no -dirty" 0 - else - check "NAT batch: startup line for rev ${expect_rev} (found $rev_lines)" 1 - return 0 - fi - local baseline="$NATBIG_ALLOCATED" + 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 @@ -608,56 +731,14 @@ PYEOF return 0 fi - # Closed-loop AAAA driver, 4 workers, one fresh socket per query. - docker exec -i "$CLIENT" sh -c 'cat > /tmp/gw_natbig.py' <<'PYEOF' -import random, socket, struct, sys, threading -server = sys.argv[1] -names = [n.strip() for n in sys.stdin if n.strip()] -lock = threading.Lock() -counts = {"answered": 0, "servfail": 0, "timeout": 0, "other": 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 worker(): - while True: - with lock: - if not names: - return - name = names.pop() - outcome = query(name) - with lock: - counts[outcome] += 1 -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())) -PYEOF - + 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. @@ -667,7 +748,7 @@ PYEOF return 0 fi out=$(docker exec -i "$CLIENT" timeout "$remaining" \ - python3 /tmp/gw_natbig.py "$GW_DNS" <"$names_file" 2>&1) || rc=$? + 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 @@ -727,6 +808,278 @@ PYEOF 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