diff --git a/CHANGELOG.md b/CHANGELOG.md index fa595c8a..870d1060 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -570,6 +570,28 @@ 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. +- 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 f0c11d60..7a0d2770 100644 --- a/src/gateway/nat.rs +++ b/src/gateway/nat.rs @@ -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 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]) -> 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 +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, 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 +971,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/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 9cb8a605..3153bf71 100755 --- a/testing/static/scripts/gateway-test.sh +++ b/testing/static/scripts/gateway-test.sh @@ -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