mirror of
https://github.com/jmcorgan/fips.git
synced 2026-10-06 11:38:24 +00:00
Merge branch 'maint'
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -53,6 +53,12 @@ fn main() {
|
||||
std::process::exit(1);
|
||||
}
|
||||
|
||||
/// Microseconds since `started`, saturating, for the timing fields on debug
|
||||
/// lines.
|
||||
fn elapsed_us(started: Instant) -> u64 {
|
||||
u64::try_from(started.elapsed().as_micros()).unwrap_or(u64::MAX)
|
||||
}
|
||||
|
||||
/// Take a conntrack snapshot off the runtime thread.
|
||||
///
|
||||
/// A failed read yields an empty snapshot, so every mapping reads zero
|
||||
@@ -446,14 +452,19 @@ async fn main() {
|
||||
// the pool lock: the runtime is current-thread, so a
|
||||
// blocking read here would stall the DNS resolver, and the
|
||||
// read must not happen under the lock the resolver needs.
|
||||
let read_started = Instant::now();
|
||||
let conntrack = read_conntrack(&mut conntrack_log).await;
|
||||
let read_us = elapsed_us(read_started);
|
||||
let mut pool_guard = tick_pool.lock().await;
|
||||
let tick_started = Instant::now();
|
||||
let events = pool_guard.tick(now, &conntrack);
|
||||
let tick_us = elapsed_us(tick_started);
|
||||
|
||||
// Build snapshot for control socket
|
||||
let pool_status = pool_guard.status();
|
||||
let mappings = pool_guard.mapping_info(now);
|
||||
drop(pool_guard);
|
||||
debug!(mappings = mappings.len(), read_us, tick_us, "Pool tick");
|
||||
|
||||
let snapshot = control::build_snapshot(
|
||||
pool_status,
|
||||
|
||||
+783
-137
@@ -4,7 +4,10 @@
|
||||
//! for translating between virtual IPs and FIPS mesh addresses.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::fmt;
|
||||
use std::net::Ipv6Addr;
|
||||
use std::os::fd::{AsRawFd, FromRawFd, OwnedFd};
|
||||
use std::time::Instant;
|
||||
use tracing::{debug, info};
|
||||
|
||||
use rustables::expr::{
|
||||
@@ -23,6 +26,73 @@ const POSTROUTING_CHAIN: &str = "postrouting";
|
||||
const DSTNAT_PRIORITY: i32 = -100;
|
||||
const SRCNAT_PRIORITY: i32 = 100;
|
||||
|
||||
/// Largest value the kernel accepts for `SO_SNDBUFFORCE`.
|
||||
///
|
||||
/// The kernel clamps the requested value to `i32::MAX / 2` and then doubles
|
||||
/// it, so the socket's send buffer never exceeds `2 * MAX_SNDBUF`.
|
||||
const MAX_SNDBUF: libc::c_int = libc::c_int::MAX / 2;
|
||||
|
||||
/// Headroom added to half the batch length when sizing the send buffer.
|
||||
const SNDBUF_HEADROOM: u64 = 64 * 1024;
|
||||
|
||||
/// The kernel refuses a netlink message longer than the send buffer less
|
||||
/// this many bytes.
|
||||
const SNDBUF_OVERHEAD: u64 = 32;
|
||||
|
||||
/// How long the rebuild waits for the kernel's acknowledgement. The rebuild
|
||||
/// runs inside the gateway's event loop, so this bounds the stall there.
|
||||
const ACK_TIMEOUT_SECS: libc::time_t = 5;
|
||||
|
||||
/// Length of a `struct nlmsghdr`.
|
||||
const NLMSG_HDRLEN: usize = 16;
|
||||
|
||||
/// An errno value, displayed by name and number, e.g. `EMSGSIZE (90)`.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub struct Errno(pub i32);
|
||||
|
||||
impl Errno {
|
||||
/// The errno of the last failed libc call on this thread.
|
||||
fn last() -> Self {
|
||||
Errno(std::io::Error::last_os_error().raw_os_error().unwrap_or(0))
|
||||
}
|
||||
|
||||
/// The symbolic name of the errno, for the values netlink can return.
|
||||
fn name(self) -> &'static str {
|
||||
match self.0 {
|
||||
libc::EPERM => "EPERM",
|
||||
libc::ENOENT => "ENOENT",
|
||||
libc::EINTR => "EINTR",
|
||||
libc::EBADF => "EBADF",
|
||||
libc::EAGAIN => "EAGAIN",
|
||||
libc::ENOMEM => "ENOMEM",
|
||||
libc::EACCES => "EACCES",
|
||||
libc::EFAULT => "EFAULT",
|
||||
libc::EBUSY => "EBUSY",
|
||||
libc::EEXIST => "EEXIST",
|
||||
libc::ENODEV => "ENODEV",
|
||||
libc::EINVAL => "EINVAL",
|
||||
libc::ENFILE => "ENFILE",
|
||||
libc::EMFILE => "EMFILE",
|
||||
libc::ENOSPC => "ENOSPC",
|
||||
libc::ERANGE => "ERANGE",
|
||||
libc::ELOOP => "ELOOP",
|
||||
libc::EMSGSIZE => "EMSGSIZE",
|
||||
libc::EPROTONOSUPPORT => "EPROTONOSUPPORT",
|
||||
libc::EOPNOTSUPP => "EOPNOTSUPP",
|
||||
libc::EAFNOSUPPORT => "EAFNOSUPPORT",
|
||||
libc::ENOBUFS => "ENOBUFS",
|
||||
libc::ETIMEDOUT => "ETIMEDOUT",
|
||||
_ => "errno",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Display for Errno {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
write!(f, "{} ({})", self.name(), self.0)
|
||||
}
|
||||
}
|
||||
|
||||
/// Errors from NAT operations.
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
pub enum NatError {
|
||||
@@ -30,6 +100,18 @@ pub enum NatError {
|
||||
Nftables(String),
|
||||
#[error("rule not found for virtual IP {0}")]
|
||||
RuleNotFound(Ipv6Addr),
|
||||
/// The kernel rejected a message of the NAT batch; the batch was aborted.
|
||||
#[error("kernel rejected netlink message {seq} of the NAT batch: {errno}")]
|
||||
Kernel { errno: Errno, seq: u32 },
|
||||
/// A netlink socket call failed.
|
||||
#[error("netlink socket {op} failed: {errno}")]
|
||||
Socket { op: &'static str, errno: Errno },
|
||||
/// The batch is larger than any netlink send buffer the kernel allows.
|
||||
#[error(
|
||||
"NAT batch of {bytes} bytes exceeds the kernel's netlink limit of {} bytes",
|
||||
admissible_limit()
|
||||
)]
|
||||
BatchTooLarge { bytes: usize },
|
||||
}
|
||||
|
||||
impl From<rustables::error::QueryError> for NatError {
|
||||
@@ -54,8 +136,8 @@ struct NatMapping {
|
||||
/// One object a NAT rebuild sends, named rather than built.
|
||||
///
|
||||
/// `rebuild_batches` decides what a rebuild sends and in what order;
|
||||
/// `send_batches` turns that decision into rustables objects and hands each
|
||||
/// batch to the kernel. The split is what lets a test see the delete and the
|
||||
/// `encode_batch` turns that decision into netlink bytes and `send_batch`
|
||||
/// hands them to the kernel. The split is what lets a test see the delete and the
|
||||
/// recreate share one transaction without a netlink socket, which is the
|
||||
/// property that keeps the table in the packet path.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
@@ -161,14 +243,26 @@ impl NatManager {
|
||||
mesh_addr,
|
||||
},
|
||||
);
|
||||
self.rebuild()?;
|
||||
|
||||
debug!(
|
||||
virtual_ip = %virtual_ip,
|
||||
mesh_addr = %mesh_addr,
|
||||
"Added DNAT/SNAT rules"
|
||||
);
|
||||
Ok(())
|
||||
let (result, elapsed_us) = self.timed_rebuild();
|
||||
let mappings = self.mappings.len();
|
||||
match &result {
|
||||
Ok(()) => debug!(
|
||||
virtual_ip = %virtual_ip,
|
||||
mesh_addr = %mesh_addr,
|
||||
mappings,
|
||||
elapsed_us,
|
||||
"Added DNAT/SNAT rules"
|
||||
),
|
||||
Err(e) => debug!(
|
||||
virtual_ip = %virtual_ip,
|
||||
mesh_addr = %mesh_addr,
|
||||
mappings,
|
||||
elapsed_us,
|
||||
error = %e,
|
||||
"Added DNAT/SNAT rules"
|
||||
),
|
||||
}
|
||||
result
|
||||
}
|
||||
|
||||
/// Remove DNAT and SNAT rules for a virtual IP mapping.
|
||||
@@ -176,10 +270,24 @@ impl NatManager {
|
||||
if self.mappings.remove(&virtual_ip).is_none() {
|
||||
return Err(NatError::RuleNotFound(virtual_ip));
|
||||
}
|
||||
self.rebuild()?;
|
||||
|
||||
debug!(virtual_ip = %virtual_ip, "Removed DNAT/SNAT rules");
|
||||
Ok(())
|
||||
let (result, elapsed_us) = self.timed_rebuild();
|
||||
let mappings = self.mappings.len();
|
||||
match &result {
|
||||
Ok(()) => debug!(
|
||||
virtual_ip = %virtual_ip,
|
||||
mappings,
|
||||
elapsed_us,
|
||||
"Removed DNAT/SNAT rules"
|
||||
),
|
||||
Err(e) => debug!(
|
||||
virtual_ip = %virtual_ip,
|
||||
mappings,
|
||||
elapsed_us,
|
||||
error = %e,
|
||||
"Removed DNAT/SNAT rules"
|
||||
),
|
||||
}
|
||||
result
|
||||
}
|
||||
|
||||
/// Flush all rules and delete the nftables table.
|
||||
@@ -188,7 +296,7 @@ impl NatManager {
|
||||
batch.add(&self.table, MsgType::Del);
|
||||
batch
|
||||
.send()
|
||||
.map_err(|e| NatError::Nftables(e.to_string()))?;
|
||||
.map_err(|e| NatError::Nftables(error_chain(&e)))?;
|
||||
|
||||
info!("Deleted nftables table '{TABLE_NAME}'");
|
||||
Ok(())
|
||||
@@ -238,133 +346,137 @@ impl NatManager {
|
||||
vec![ops]
|
||||
}
|
||||
|
||||
/// Build each op into its rustables object and send the batches in order.
|
||||
fn send_batches(&self, batches: &[Vec<NatOp>]) -> Result<(), NatError> {
|
||||
for ops in batches {
|
||||
let mut batch = Batch::new();
|
||||
for op in ops {
|
||||
match *op {
|
||||
NatOp::Table(msg_type) => batch.add(&self.table, msg_type),
|
||||
NatOp::PreChain => batch.add(&self.pre_chain, MsgType::Add),
|
||||
NatOp::PostChain => batch.add(&self.post_chain, MsgType::Add),
|
||||
NatOp::FipsMasquerade => {
|
||||
// Rewrite the source address of traffic leaving fips0.
|
||||
// Without this, LAN clients' source addresses (e.g.
|
||||
// fd02::20) are not routable on the mesh, so return
|
||||
// traffic would be black-holed.
|
||||
let rule = Rule::new(&self.post_chain)?
|
||||
.with_expr(Meta::new(MetaType::OifName))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec()))
|
||||
.with_expr(Masquerade::default());
|
||||
batch.add(&rule, MsgType::Add);
|
||||
}
|
||||
NatOp::Dnat(virtual_ip) => {
|
||||
let mapping = self.mapping(virtual_ip)?;
|
||||
let rule = Rule::new(&self.pre_chain)?
|
||||
.with_expr(Meta::new(MetaType::NfProto))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
|
||||
.with_expr(
|
||||
HighLevelPayload::Network(NetworkHeaderField::IPv6(
|
||||
IPv6HeaderField::Daddr,
|
||||
))
|
||||
.build(),
|
||||
)
|
||||
.with_expr(Cmp::new(CmpOp::Eq, mapping.virtual_ip.octets()))
|
||||
.with_expr(Immediate::new_data(
|
||||
mapping.mesh_addr.octets().to_vec(),
|
||||
Register::Reg1,
|
||||
/// Build each op into its rustables object and encode the batch.
|
||||
///
|
||||
/// Only the last object before the batch end requests an
|
||||
/// acknowledgement. rustables sets `NLM_F_ACK` on every message, and one
|
||||
/// ack per message overflows the socket's receive buffer from about a
|
||||
/// hundred mappings, after the kernel has already committed the batch.
|
||||
/// The kernel reports a failing message whatever its flags, so errors
|
||||
/// stay attributable.
|
||||
fn encode_batch(&self, ops: &[NatOp]) -> Result<Vec<u8>, NatError> {
|
||||
let mut batch = Batch::new();
|
||||
for op in ops {
|
||||
match *op {
|
||||
NatOp::Table(msg_type) => batch.add(&self.table, msg_type),
|
||||
NatOp::PreChain => batch.add(&self.pre_chain, MsgType::Add),
|
||||
NatOp::PostChain => batch.add(&self.post_chain, MsgType::Add),
|
||||
NatOp::FipsMasquerade => {
|
||||
// Rewrite the source address of traffic leaving fips0.
|
||||
// Without this, LAN clients' source addresses (e.g.
|
||||
// fd02::20) are not routable on the mesh, so return
|
||||
// traffic would be black-holed.
|
||||
let rule = Rule::new(&self.post_chain)?
|
||||
.with_expr(Meta::new(MetaType::OifName))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec()))
|
||||
.with_expr(Masquerade::default());
|
||||
batch.add(&rule, MsgType::Add);
|
||||
}
|
||||
NatOp::Dnat(virtual_ip) => {
|
||||
let mapping = self.mapping(virtual_ip)?;
|
||||
let rule = Rule::new(&self.pre_chain)?
|
||||
.with_expr(Meta::new(MetaType::NfProto))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
|
||||
.with_expr(
|
||||
HighLevelPayload::Network(NetworkHeaderField::IPv6(
|
||||
IPv6HeaderField::Daddr,
|
||||
))
|
||||
.with_expr(
|
||||
Nat::default()
|
||||
.with_nat_type(NatType::DNat)
|
||||
.with_family(ProtocolFamily::Ipv6)
|
||||
.with_ip_register(Register::Reg1),
|
||||
);
|
||||
batch.add(&rule, MsgType::Add);
|
||||
}
|
||||
NatOp::Snat(virtual_ip) => {
|
||||
let mapping = self.mapping(virtual_ip)?;
|
||||
let rule = Rule::new(&self.post_chain)?
|
||||
.with_expr(Meta::new(MetaType::NfProto))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
|
||||
.with_expr(
|
||||
HighLevelPayload::Network(NetworkHeaderField::IPv6(
|
||||
IPv6HeaderField::Saddr,
|
||||
))
|
||||
.build(),
|
||||
)
|
||||
.with_expr(Cmp::new(CmpOp::Eq, mapping.mesh_addr.octets()))
|
||||
.with_expr(Immediate::new_data(
|
||||
mapping.virtual_ip.octets().to_vec(),
|
||||
Register::Reg1,
|
||||
))
|
||||
.with_expr(
|
||||
Nat::default()
|
||||
.with_nat_type(NatType::SNat)
|
||||
.with_family(ProtocolFamily::Ipv6)
|
||||
.with_ip_register(Register::Reg1),
|
||||
);
|
||||
batch.add(&rule, MsgType::Add);
|
||||
}
|
||||
NatOp::PortForward(index) => {
|
||||
let pf = self.port_forwards.get(index).expect(
|
||||
"rebuild_batches only emits indices it read from port_forwards",
|
||||
.build(),
|
||||
)
|
||||
.with_expr(Cmp::new(CmpOp::Eq, mapping.virtual_ip.octets()))
|
||||
.with_expr(Immediate::new_data(
|
||||
mapping.mesh_addr.octets().to_vec(),
|
||||
Register::Reg1,
|
||||
))
|
||||
.with_expr(
|
||||
Nat::default()
|
||||
.with_nat_type(NatType::DNat)
|
||||
.with_family(ProtocolFamily::Ipv6)
|
||||
.with_ip_register(Register::Reg1),
|
||||
);
|
||||
let l4proto: u8 = match pf.proto {
|
||||
Proto::Tcp => libc::IPPROTO_TCP as u8,
|
||||
Proto::Udp => libc::IPPROTO_UDP as u8,
|
||||
};
|
||||
let dport_field = match pf.proto {
|
||||
Proto::Tcp => TransportHeaderField::Tcp(TCPHeaderField::Dport),
|
||||
Proto::Udp => TransportHeaderField::Udp(UDPHeaderField::Dport),
|
||||
};
|
||||
let target_ip = *pf.target.ip();
|
||||
let target_port_be = pf.target.port().to_be_bytes();
|
||||
|
||||
let rule = Rule::new(&self.pre_chain)?
|
||||
.with_expr(Meta::new(MetaType::IifName))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec()))
|
||||
.with_expr(Meta::new(MetaType::NfProto))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
|
||||
.with_expr(Meta::new(MetaType::L4Proto))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, [l4proto]))
|
||||
.with_expr(HighLevelPayload::Transport(dport_field).build())
|
||||
.with_expr(Cmp::new(CmpOp::Eq, pf.listen_port.to_be_bytes().to_vec()))
|
||||
.with_expr(Immediate::new_data(
|
||||
target_ip.octets().to_vec(),
|
||||
Register::Reg1,
|
||||
batch.add(&rule, MsgType::Add);
|
||||
}
|
||||
NatOp::Snat(virtual_ip) => {
|
||||
let mapping = self.mapping(virtual_ip)?;
|
||||
let rule = Rule::new(&self.post_chain)?
|
||||
.with_expr(Meta::new(MetaType::NfProto))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
|
||||
.with_expr(
|
||||
HighLevelPayload::Network(NetworkHeaderField::IPv6(
|
||||
IPv6HeaderField::Saddr,
|
||||
))
|
||||
.with_expr(Immediate::new_data(target_port_be.to_vec(), Register::Reg2))
|
||||
.with_expr(
|
||||
Nat::default()
|
||||
.with_nat_type(NatType::DNat)
|
||||
.with_family(ProtocolFamily::Ipv6)
|
||||
.with_ip_register(Register::Reg1)
|
||||
.with_port_register(Register::Reg2),
|
||||
);
|
||||
batch.add(&rule, MsgType::Add);
|
||||
}
|
||||
NatOp::LanMasquerade => {
|
||||
let mut lan_iface = self.lan_interface.clone().into_bytes();
|
||||
lan_iface.push(0);
|
||||
let rule = Rule::new(&self.post_chain)?
|
||||
.with_expr(Meta::new(MetaType::IifName))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec()))
|
||||
.with_expr(Meta::new(MetaType::OifName))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, lan_iface))
|
||||
.with_expr(Meta::new(MetaType::NfProto))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
|
||||
.with_expr(Masquerade::default());
|
||||
batch.add(&rule, MsgType::Add);
|
||||
}
|
||||
.build(),
|
||||
)
|
||||
.with_expr(Cmp::new(CmpOp::Eq, mapping.mesh_addr.octets()))
|
||||
.with_expr(Immediate::new_data(
|
||||
mapping.virtual_ip.octets().to_vec(),
|
||||
Register::Reg1,
|
||||
))
|
||||
.with_expr(
|
||||
Nat::default()
|
||||
.with_nat_type(NatType::SNat)
|
||||
.with_family(ProtocolFamily::Ipv6)
|
||||
.with_ip_register(Register::Reg1),
|
||||
);
|
||||
batch.add(&rule, MsgType::Add);
|
||||
}
|
||||
NatOp::PortForward(index) => {
|
||||
let pf = self
|
||||
.port_forwards
|
||||
.get(index)
|
||||
.expect("rebuild_batches only emits indices it read from port_forwards");
|
||||
let l4proto: u8 = match pf.proto {
|
||||
Proto::Tcp => libc::IPPROTO_TCP as u8,
|
||||
Proto::Udp => libc::IPPROTO_UDP as u8,
|
||||
};
|
||||
let dport_field = match pf.proto {
|
||||
Proto::Tcp => TransportHeaderField::Tcp(TCPHeaderField::Dport),
|
||||
Proto::Udp => TransportHeaderField::Udp(UDPHeaderField::Dport),
|
||||
};
|
||||
let target_ip = *pf.target.ip();
|
||||
let target_port_be = pf.target.port().to_be_bytes();
|
||||
|
||||
let rule = Rule::new(&self.pre_chain)?
|
||||
.with_expr(Meta::new(MetaType::IifName))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec()))
|
||||
.with_expr(Meta::new(MetaType::NfProto))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
|
||||
.with_expr(Meta::new(MetaType::L4Proto))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, [l4proto]))
|
||||
.with_expr(HighLevelPayload::Transport(dport_field).build())
|
||||
.with_expr(Cmp::new(CmpOp::Eq, pf.listen_port.to_be_bytes().to_vec()))
|
||||
.with_expr(Immediate::new_data(
|
||||
target_ip.octets().to_vec(),
|
||||
Register::Reg1,
|
||||
))
|
||||
.with_expr(Immediate::new_data(target_port_be.to_vec(), Register::Reg2))
|
||||
.with_expr(
|
||||
Nat::default()
|
||||
.with_nat_type(NatType::DNat)
|
||||
.with_family(ProtocolFamily::Ipv6)
|
||||
.with_ip_register(Register::Reg1)
|
||||
.with_port_register(Register::Reg2),
|
||||
);
|
||||
batch.add(&rule, MsgType::Add);
|
||||
}
|
||||
NatOp::LanMasquerade => {
|
||||
let mut lan_iface = self.lan_interface.clone().into_bytes();
|
||||
lan_iface.push(0);
|
||||
let rule = Rule::new(&self.post_chain)?
|
||||
.with_expr(Meta::new(MetaType::IifName))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec()))
|
||||
.with_expr(Meta::new(MetaType::OifName))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, lan_iface))
|
||||
.with_expr(Meta::new(MetaType::NfProto))
|
||||
.with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8]))
|
||||
.with_expr(Masquerade::default());
|
||||
batch.add(&rule, MsgType::Add);
|
||||
}
|
||||
}
|
||||
|
||||
batch
|
||||
.send()
|
||||
.map_err(|e| NatError::Nftables(e.to_string()))?;
|
||||
}
|
||||
Ok(())
|
||||
let mut bytes = batch.finalize();
|
||||
keep_last_ack(&mut bytes)?;
|
||||
Ok(bytes)
|
||||
}
|
||||
|
||||
/// The mapping an op names, or the error a caller can report.
|
||||
@@ -377,10 +489,372 @@ impl NatManager {
|
||||
/// Rebuild the entire nftables table with all current rules, in one
|
||||
/// netlink transaction.
|
||||
fn rebuild(&self) -> Result<(), NatError> {
|
||||
self.send_batches(&self.rebuild_batches())
|
||||
for ops in self.rebuild_batches() {
|
||||
send_batch(&self.encode_batch(&ops)?)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Rebuild, returning the outcome with the time the rebuild took in
|
||||
/// microseconds, so a mapping change can log its cost on either path.
|
||||
fn timed_rebuild(&self) -> (Result<(), NatError>, u64) {
|
||||
let started = Instant::now();
|
||||
let result = self.rebuild();
|
||||
let elapsed_us = u64::try_from(started.elapsed().as_micros()).unwrap_or(u64::MAX);
|
||||
(result, elapsed_us)
|
||||
}
|
||||
}
|
||||
|
||||
/// Largest batch, in bytes, that the kernel can admit in one send.
|
||||
fn admissible_limit() -> u64 {
|
||||
2 * MAX_SNDBUF as u64 - SNDBUF_OVERHEAD
|
||||
}
|
||||
|
||||
/// The `SO_SNDBUFFORCE` value that lets a batch of `len` bytes through.
|
||||
///
|
||||
/// The kernel doubles the value it is given and refuses a message longer
|
||||
/// than the result less 32 bytes, so half the length plus headroom is
|
||||
/// enough. Saturates at the kernel's own clamp rather than wrapping.
|
||||
fn sndbuf_for(len: usize) -> libc::c_int {
|
||||
let want = (u64::try_from(len).unwrap_or(u64::MAX) / 2).saturating_add(SNDBUF_HEADROOM);
|
||||
libc::c_int::try_from(want.min(MAX_SNDBUF as u64)).unwrap_or(MAX_SNDBUF)
|
||||
}
|
||||
|
||||
/// Refuse a batch the kernel could not accept at any send-buffer size.
|
||||
fn check_admissible(len: usize) -> Result<(), NatError> {
|
||||
if u64::try_from(len).unwrap_or(u64::MAX) > admissible_limit() {
|
||||
return Err(NatError::BatchTooLarge { bytes: len });
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// The fields of one `struct nlmsghdr` that the NAT batch code reads.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
struct NlHeader {
|
||||
/// Offset of the header within the buffer.
|
||||
offset: usize,
|
||||
/// `nlmsg_len`: header plus payload, without alignment padding.
|
||||
len: usize,
|
||||
kind: u16,
|
||||
flags: u16,
|
||||
seq: u32,
|
||||
}
|
||||
|
||||
impl NlHeader {
|
||||
/// The message's payload, after the header.
|
||||
fn payload<'a>(&self, buf: &'a [u8]) -> &'a [u8] {
|
||||
&buf[self.offset + NLMSG_HDRLEN..self.offset + self.len]
|
||||
}
|
||||
}
|
||||
|
||||
/// Walk every netlink message header in `buf`.
|
||||
///
|
||||
/// Fails on a header shorter than `struct nlmsghdr` or a length that runs
|
||||
/// past the end of the buffer.
|
||||
fn nl_headers(buf: &[u8]) -> Result<Vec<NlHeader>, NatError> {
|
||||
let mut headers = Vec::new();
|
||||
let mut offset = 0;
|
||||
while offset < buf.len() {
|
||||
let rest = &buf[offset..];
|
||||
if rest.len() < NLMSG_HDRLEN {
|
||||
return Err(NatError::Nftables(format!(
|
||||
"malformed netlink message at offset {offset}: {} bytes left, header needs {NLMSG_HDRLEN}",
|
||||
rest.len()
|
||||
)));
|
||||
}
|
||||
let field = |at: usize, width: usize| &rest[at..at + width];
|
||||
let len = u32::from_ne_bytes(field(0, 4).try_into().expect("4-byte slice")) as usize;
|
||||
if len < NLMSG_HDRLEN || len > rest.len() {
|
||||
return Err(NatError::Nftables(format!(
|
||||
"malformed netlink message at offset {offset}: length {len} with {} bytes left",
|
||||
rest.len()
|
||||
)));
|
||||
}
|
||||
headers.push(NlHeader {
|
||||
offset,
|
||||
len,
|
||||
kind: u16::from_ne_bytes(field(4, 2).try_into().expect("2-byte slice")),
|
||||
flags: u16::from_ne_bytes(field(6, 2).try_into().expect("2-byte slice")),
|
||||
seq: u32::from_ne_bytes(field(8, 4).try_into().expect("4-byte slice")),
|
||||
});
|
||||
// Netlink messages are 4-byte aligned.
|
||||
offset += (len + 3) & !3;
|
||||
}
|
||||
Ok(headers)
|
||||
}
|
||||
|
||||
/// The finalized batch's objects: every message between the batch begin
|
||||
/// and the batch end.
|
||||
fn batch_objects(headers: &[NlHeader]) -> Result<&[NlHeader], NatError> {
|
||||
match headers {
|
||||
[_begin, objects @ .., _end] if !objects.is_empty() => Ok(objects),
|
||||
_ => Err(NatError::Nftables(format!(
|
||||
"NAT batch holds {} messages; it needs a begin, an object and an end",
|
||||
headers.len()
|
||||
))),
|
||||
}
|
||||
}
|
||||
|
||||
/// Clear `NLM_F_ACK` on every message of a finalized batch except the last
|
||||
/// object before the batch end.
|
||||
fn keep_last_ack(buf: &mut [u8]) -> Result<(), NatError> {
|
||||
let headers = nl_headers(buf)?;
|
||||
let last = batch_objects(&headers)?
|
||||
.last()
|
||||
.expect("batch_objects returns a non-empty slice")
|
||||
.offset;
|
||||
let ack = libc::NLM_F_ACK as u16;
|
||||
for header in &headers {
|
||||
let flags = if header.offset == last {
|
||||
header.flags | ack
|
||||
} else {
|
||||
header.flags & !ack
|
||||
};
|
||||
buf[header.offset + 6..header.offset + 8].copy_from_slice(&flags.to_ne_bytes());
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// What the kernel's replies to a NAT batch have shown so far.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
enum AckState {
|
||||
/// No verdict yet; read another datagram.
|
||||
Pending,
|
||||
/// The acknowledgement of the batch's last object arrived with no error
|
||||
/// before it.
|
||||
Done,
|
||||
}
|
||||
|
||||
/// Reads the kernel's replies to a NAT batch, one datagram at a time.
|
||||
///
|
||||
/// The kernel aborts the whole batch when any message fails, yet it still
|
||||
/// acknowledges the last message after the error. So the last ack alone does
|
||||
/// not prove success: any error that arrives before it fails the batch.
|
||||
struct AckReader {
|
||||
/// Sequence number of the one message that requested an ack.
|
||||
last_seq: u32,
|
||||
}
|
||||
|
||||
impl AckReader {
|
||||
/// Consume one received datagram, which may carry several messages.
|
||||
fn feed(&self, datagram: &[u8]) -> Result<AckState, NatError> {
|
||||
for header in nl_headers(datagram)? {
|
||||
if i32::from(header.kind) != libc::NLMSG_ERROR {
|
||||
continue;
|
||||
}
|
||||
let payload = header.payload(datagram);
|
||||
if payload.len() < 4 {
|
||||
return Err(NatError::Nftables(format!(
|
||||
"malformed netlink error message: {} payload bytes, error field needs 4",
|
||||
payload.len()
|
||||
)));
|
||||
}
|
||||
let error = i32::from_ne_bytes(payload[..4].try_into().expect("4-byte slice"));
|
||||
if error != 0 {
|
||||
return Err(NatError::Kernel {
|
||||
errno: Errno(error.saturating_neg()),
|
||||
seq: header.seq,
|
||||
});
|
||||
}
|
||||
if header.seq == self.last_seq {
|
||||
return Ok(AckState::Done);
|
||||
}
|
||||
}
|
||||
Ok(AckState::Pending)
|
||||
}
|
||||
}
|
||||
|
||||
/// Render an error with every source beneath it, so a wrapped errno is kept.
|
||||
fn error_chain(error: &dyn std::error::Error) -> String {
|
||||
let mut text = error.to_string();
|
||||
let mut source = error.source();
|
||||
while let Some(inner) = source {
|
||||
text.push_str(": ");
|
||||
text.push_str(&inner.to_string());
|
||||
source = inner.source();
|
||||
}
|
||||
text
|
||||
}
|
||||
|
||||
/// Set an integer socket option.
|
||||
fn set_int_opt(
|
||||
sock: &OwnedFd,
|
||||
level: libc::c_int,
|
||||
name: libc::c_int,
|
||||
value: libc::c_int,
|
||||
) -> Result<(), Errno> {
|
||||
// SAFETY: the descriptor is open for the life of `sock`, and the pointer
|
||||
// and length describe `value`, a c_int.
|
||||
let rc = unsafe {
|
||||
libc::setsockopt(
|
||||
sock.as_raw_fd(),
|
||||
level,
|
||||
name,
|
||||
(&value as *const libc::c_int).cast(),
|
||||
std::mem::size_of::<libc::c_int>() as libc::socklen_t,
|
||||
)
|
||||
};
|
||||
if rc < 0 { Err(Errno::last()) } else { Ok(()) }
|
||||
}
|
||||
|
||||
/// Size of the buffer each reply datagram is read into.
|
||||
///
|
||||
/// The largest message nftables sends back, as rustables computes it
|
||||
/// (`nft_nlmsg_maxsize`, which it does not export), and at least 64 KiB.
|
||||
fn recv_buffer_len() -> usize {
|
||||
// SAFETY: sysconf has no preconditions.
|
||||
let page = unsafe { libc::sysconf(libc::_SC_PAGESIZE) };
|
||||
(usize::from(u16::MAX) + usize::try_from(page).unwrap_or(0)).max(64 * 1024)
|
||||
}
|
||||
|
||||
/// Open a netfilter netlink socket sized for a batch of `len` bytes.
|
||||
fn open_batch_socket(len: usize) -> Result<OwnedFd, NatError> {
|
||||
let socket_err = |op| move |errno| NatError::Socket { op, errno };
|
||||
|
||||
// SAFETY: socket has no memory preconditions.
|
||||
let fd = unsafe {
|
||||
libc::socket(
|
||||
libc::AF_NETLINK,
|
||||
libc::SOCK_RAW | libc::SOCK_CLOEXEC,
|
||||
libc::NETLINK_NETFILTER,
|
||||
)
|
||||
};
|
||||
if fd < 0 {
|
||||
return Err(socket_err("open")(Errno::last()));
|
||||
}
|
||||
// SAFETY: socket returned a new descriptor that nothing else owns.
|
||||
let sock = unsafe { OwnedFd::from_raw_fd(fd) };
|
||||
|
||||
// SAFETY: an all-zero sockaddr_nl is valid; the family is set below.
|
||||
let mut addr: libc::sockaddr_nl = unsafe { std::mem::zeroed() };
|
||||
addr.nl_family = libc::AF_NETLINK as libc::sa_family_t;
|
||||
// SAFETY: the pointer and length describe `addr`, a sockaddr_nl.
|
||||
let rc = unsafe {
|
||||
libc::bind(
|
||||
sock.as_raw_fd(),
|
||||
(&addr as *const libc::sockaddr_nl).cast(),
|
||||
std::mem::size_of::<libc::sockaddr_nl>() as libc::socklen_t,
|
||||
)
|
||||
};
|
||||
if rc < 0 {
|
||||
return Err(socket_err("bind")(Errno::last()));
|
||||
}
|
||||
|
||||
// Without CAP_NET_ADMIN the forced size is refused; the plain option is
|
||||
// then capped by wmem_max, and an oversized batch fails with EMSGSIZE.
|
||||
let sndbuf = sndbuf_for(len);
|
||||
match set_int_opt(&sock, libc::SOL_SOCKET, libc::SO_SNDBUFFORCE, sndbuf) {
|
||||
Err(Errno(libc::EPERM)) => {
|
||||
set_int_opt(&sock, libc::SOL_SOCKET, libc::SO_SNDBUF, sndbuf)
|
||||
.map_err(socket_err("setsockopt SO_SNDBUF"))?;
|
||||
}
|
||||
other => other.map_err(socket_err("setsockopt SO_SNDBUFFORCE"))?,
|
||||
}
|
||||
// An error ack then carries only the failing header, not the message.
|
||||
set_int_opt(&sock, libc::SOL_NETLINK, libc::NETLINK_CAP_ACK, 1)
|
||||
.map_err(socket_err("setsockopt NETLINK_CAP_ACK"))?;
|
||||
|
||||
let timeout = libc::timeval {
|
||||
tv_sec: ACK_TIMEOUT_SECS,
|
||||
tv_usec: 0,
|
||||
};
|
||||
// SAFETY: the pointer and length describe `timeout`, a timeval.
|
||||
let rc = unsafe {
|
||||
libc::setsockopt(
|
||||
sock.as_raw_fd(),
|
||||
libc::SOL_SOCKET,
|
||||
libc::SO_RCVTIMEO,
|
||||
(&timeout as *const libc::timeval).cast(),
|
||||
std::mem::size_of::<libc::timeval>() as libc::socklen_t,
|
||||
)
|
||||
};
|
||||
if rc < 0 {
|
||||
return Err(socket_err("setsockopt SO_RCVTIMEO")(Errno::last()));
|
||||
}
|
||||
Ok(sock)
|
||||
}
|
||||
|
||||
/// Send one encoded NAT batch and wait for the kernel's verdict.
|
||||
///
|
||||
/// The batch goes out in a single send, so the kernel applies it as one
|
||||
/// transaction. The send buffer is sized to the batch, because the default
|
||||
/// one refuses a message past about 208 KiB, which is about 313 mappings.
|
||||
fn send_batch(bytes: &[u8]) -> Result<(), NatError> {
|
||||
check_admissible(bytes.len())?;
|
||||
let headers = nl_headers(bytes)?;
|
||||
let reader = AckReader {
|
||||
last_seq: batch_objects(&headers)?
|
||||
.last()
|
||||
.expect("batch_objects returns a non-empty slice")
|
||||
.seq,
|
||||
};
|
||||
let sock = open_batch_socket(bytes.len())?;
|
||||
|
||||
let sent = loop {
|
||||
// SAFETY: the pointer and length describe `bytes`.
|
||||
let rc = unsafe { libc::send(sock.as_raw_fd(), bytes.as_ptr().cast(), bytes.len(), 0) };
|
||||
if rc >= 0 {
|
||||
break rc as usize;
|
||||
}
|
||||
let errno = Errno::last();
|
||||
if errno.0 != libc::EINTR {
|
||||
return Err(NatError::Socket { op: "send", errno });
|
||||
}
|
||||
};
|
||||
if sent != bytes.len() {
|
||||
return Err(NatError::Nftables(format!(
|
||||
"netlink send took {sent} of {} bytes",
|
||||
bytes.len()
|
||||
)));
|
||||
}
|
||||
|
||||
let mut buf = vec![0u8; recv_buffer_len()];
|
||||
loop {
|
||||
// MSG_TRUNC makes a netlink recv return the datagram's full length,
|
||||
// so a reply larger than the buffer is detected rather than cut.
|
||||
// SAFETY: the pointer and length describe `buf`.
|
||||
let rc = unsafe {
|
||||
libc::recv(
|
||||
sock.as_raw_fd(),
|
||||
buf.as_mut_ptr().cast(),
|
||||
buf.len(),
|
||||
libc::MSG_TRUNC,
|
||||
)
|
||||
};
|
||||
if rc < 0 {
|
||||
let errno = Errno::last();
|
||||
if errno.0 == libc::EINTR {
|
||||
continue;
|
||||
}
|
||||
// EAGAIN is the receive timeout. ENOBUFS means error acks
|
||||
// overflowed the receive buffer, since only one ack is requested.
|
||||
return Err(NatError::Socket { op: "recv", errno });
|
||||
}
|
||||
let got = rc as usize;
|
||||
if got == 0 {
|
||||
return Err(NatError::Nftables(
|
||||
"netlink socket returned no reply".into(),
|
||||
));
|
||||
}
|
||||
if got > buf.len() {
|
||||
return Err(NatError::Nftables(format!(
|
||||
"netlink reply of {got} bytes truncated to {}",
|
||||
buf.len()
|
||||
)));
|
||||
}
|
||||
if reader.feed(&buf[..got])? == AckState::Done {
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Coverage gap. These tests run unprivileged and open no netlink socket, so
|
||||
// three failure paths in `open_socket` and `send_batch` go unexercised here.
|
||||
// The receive timeout firing and a reply longer than the buffer (seen through
|
||||
// `MSG_TRUNC`) need a kernel fault to provoke, so nothing runs them. The
|
||||
// `SO_SNDBUF` fallback after `SO_SNDBUFFORCE` returns `EPERM` needs a process
|
||||
// without CAP_NET_ADMIN, and the gateway suite's container is privileged, so
|
||||
// nothing runs that either. The gateway suite covers only the success path.
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
@@ -497,4 +971,176 @@ mod tests {
|
||||
1
|
||||
);
|
||||
}
|
||||
|
||||
/// The encoded rebuild of a manager holding `count` mappings.
|
||||
fn encoded_rebuild(count: u16) -> Vec<u8> {
|
||||
let mgr = manager_with_mappings(count);
|
||||
let ops = mgr.rebuild_batches().remove(0);
|
||||
mgr.encode_batch(&ops).expect("the rebuild encodes")
|
||||
}
|
||||
|
||||
/// One netlink message, padded to 4 bytes.
|
||||
fn nlmsg(kind: u16, seq: u32, payload: &[u8]) -> Vec<u8> {
|
||||
let len = (NLMSG_HDRLEN + payload.len()) as u32;
|
||||
let mut msg = Vec::new();
|
||||
msg.extend_from_slice(&len.to_ne_bytes());
|
||||
msg.extend_from_slice(&kind.to_ne_bytes());
|
||||
msg.extend_from_slice(&0u16.to_ne_bytes());
|
||||
msg.extend_from_slice(&seq.to_ne_bytes());
|
||||
msg.extend_from_slice(&0u32.to_ne_bytes());
|
||||
msg.extend_from_slice(payload);
|
||||
msg.resize(msg.len().div_ceil(4) * 4, 0);
|
||||
msg
|
||||
}
|
||||
|
||||
/// The kernel's `NLMSG_ERROR` reply to message `seq`, as it sends it on a
|
||||
/// socket with `NETLINK_CAP_ACK`: the error, then the request's header.
|
||||
fn ack(seq: u32, error: i32) -> Vec<u8> {
|
||||
let mut payload = error.to_ne_bytes().to_vec();
|
||||
payload.extend_from_slice(&nlmsg(0x0a00, seq, &[])[..NLMSG_HDRLEN]);
|
||||
nlmsg(libc::NLMSG_ERROR as u16, seq, &payload)
|
||||
}
|
||||
|
||||
/// The largest batch the kernel admits: twice its send-buffer clamp,
|
||||
/// less the 32 bytes netlink reserves.
|
||||
const KERNEL_BATCH_LIMIT: usize = 2_147_483_614;
|
||||
|
||||
#[test]
|
||||
fn rebuild_for_2000_mappings_requests_exactly_one_ack_on_the_last_message() {
|
||||
let encoded = encoded_rebuild(2000);
|
||||
assert!(
|
||||
encoded.len() > 212_960,
|
||||
"the 2000-mapping batch ({} bytes) must be past the default \
|
||||
netlink send limit for this test to cover the large case",
|
||||
encoded.len()
|
||||
);
|
||||
|
||||
let headers = nl_headers(&encoded).expect("the batch parses");
|
||||
assert_eq!(
|
||||
headers.first().map(|h| h.kind),
|
||||
Some(libc::NFNL_MSG_BATCH_BEGIN as u16)
|
||||
);
|
||||
assert_eq!(
|
||||
headers.last().map(|h| h.kind),
|
||||
Some(libc::NFNL_MSG_BATCH_END as u16)
|
||||
);
|
||||
let acked: Vec<usize> = headers
|
||||
.iter()
|
||||
.enumerate()
|
||||
.filter(|(_, h)| h.flags & libc::NLM_F_ACK as u16 != 0)
|
||||
.map(|(i, _)| i)
|
||||
.collect();
|
||||
assert_eq!(
|
||||
acked,
|
||||
vec![headers.len() - 2],
|
||||
"only the last object before the batch end may request an ack; \
|
||||
one ack per message overflows the receive buffer after the \
|
||||
kernel has committed the batch"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sndbuf_for_admits_the_2000_mapping_batch_and_small_batches_after_kernel_doubling() {
|
||||
let large = encoded_rebuild(2000).len();
|
||||
for len in [large, 0, 1, 212_961] {
|
||||
let sndbuf = sndbuf_for(len);
|
||||
assert!(
|
||||
2 * sndbuf as u64 - 32 >= len as u64,
|
||||
"a send buffer of {sndbuf}, doubled by the kernel, refuses a \
|
||||
{len}-byte batch"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sndbuf_for_saturates_at_the_kernel_clamp_for_huge_batches() {
|
||||
for len in [2 * MAX_SNDBUF as usize, usize::MAX] {
|
||||
assert_eq!(sndbuf_for(len), i32::MAX / 2, "sndbuf_for({len})");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn check_admissible_refuses_a_batch_larger_than_the_kernel_can_accept() {
|
||||
assert!(check_admissible(KERNEL_BATCH_LIMIT).is_ok());
|
||||
for len in [KERNEL_BATCH_LIMIT + 1, usize::MAX] {
|
||||
match check_admissible(len) {
|
||||
Err(e @ NatError::BatchTooLarge { bytes }) => {
|
||||
assert_eq!(bytes, len);
|
||||
assert!(
|
||||
e.to_string().contains(&len.to_string()),
|
||||
"the error names the batch size: {e}"
|
||||
);
|
||||
}
|
||||
other => panic!("a {len}-byte batch was admitted: {other:?}"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn ack_reader_fails_on_an_error_that_precedes_the_last_ack() {
|
||||
// The kernel aborts the batch on a failing rule mid-batch, reports
|
||||
// that rule's error, and still acknowledges the last message.
|
||||
let reader = AckReader { last_seq: 4000 };
|
||||
let error = ack(1234, -libc::ENOENT);
|
||||
let last = ack(4000, 0);
|
||||
|
||||
let expect_error = |result: Result<AckState, NatError>| match result {
|
||||
Err(NatError::Kernel { errno, seq }) => {
|
||||
assert_eq!(errno, Errno(libc::ENOENT));
|
||||
assert_eq!(seq, 1234);
|
||||
}
|
||||
other => panic!("the aborted batch was not reported: {other:?}"),
|
||||
};
|
||||
|
||||
expect_error(reader.feed(&error));
|
||||
expect_error(reader.feed(&[error.clone(), last.clone()].concat()));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn ack_reader_is_done_only_on_the_last_sequence_ack() {
|
||||
let reader = AckReader { last_seq: 10 };
|
||||
|
||||
assert_eq!(reader.feed(&ack(10, 0)).expect("parses"), AckState::Done);
|
||||
assert_eq!(reader.feed(&ack(5, 0)).expect("parses"), AckState::Pending);
|
||||
assert_eq!(
|
||||
reader
|
||||
.feed(&nlmsg(libc::NLMSG_NOOP as u16, 10, &[]))
|
||||
.expect("parses"),
|
||||
AckState::Pending
|
||||
);
|
||||
assert_eq!(
|
||||
reader
|
||||
.feed(&[ack(5, 0), ack(10, 0)].concat())
|
||||
.expect("parses"),
|
||||
AckState::Done
|
||||
);
|
||||
|
||||
let whole = ack(10, 0);
|
||||
assert!(
|
||||
reader.feed(&whole[..8]).is_err(),
|
||||
"a header shorter than 16 bytes"
|
||||
);
|
||||
let mut overlong = whole.clone();
|
||||
overlong[..4].copy_from_slice(&((whole.len() + 4) as u32).to_ne_bytes());
|
||||
assert!(
|
||||
reader.feed(&overlong).is_err(),
|
||||
"a length past the end of the datagram"
|
||||
);
|
||||
let short = nlmsg(libc::NLMSG_ERROR as u16, 10, &[0, 0]);
|
||||
assert!(
|
||||
reader.feed(&short[..NLMSG_HDRLEN + 2]).is_err(),
|
||||
"an NLMSG_ERROR payload shorter than its error field"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn kernel_error_display_names_the_errno() {
|
||||
let text = NatError::Kernel {
|
||||
errno: Errno(libc::EMSGSIZE),
|
||||
seq: 7,
|
||||
}
|
||||
.to_string();
|
||||
assert!(text.contains("EMSGSIZE"), "{text}");
|
||||
assert!(text.contains(&format!("({})", libc::EMSGSIZE)), "{text}");
|
||||
}
|
||||
}
|
||||
|
||||
+229
-2
@@ -10,6 +10,18 @@ use std::net::Ipv6Addr;
|
||||
use std::time::Instant;
|
||||
use tracing::{debug, info};
|
||||
|
||||
/// Most live mappings the pool holds before it refuses new names.
|
||||
///
|
||||
/// Every mapping adds rules to the NAT table, which is rebuilt whole on each
|
||||
/// change, and work to every tick and to shutdown, so this bounds all three.
|
||||
pub const MAPPING_CEILING: usize = 1000;
|
||||
|
||||
/// New mappings the pool admits in a burst, when idle long enough to refill.
|
||||
pub const MAPPING_BURST: u32 = 50;
|
||||
|
||||
/// New mappings per second the pool admits once a burst is spent.
|
||||
pub const MAPPING_RATE: u32 = 10;
|
||||
|
||||
/// Errors from pool operations.
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
pub enum PoolError {
|
||||
@@ -19,6 +31,10 @@ pub enum PoolError {
|
||||
Exhausted(usize),
|
||||
#[error("prefix length must be between 1 and 128")]
|
||||
InvalidPrefix,
|
||||
#[error("live-mapping ceiling reached ({0} mappings)")]
|
||||
AtCeiling(usize),
|
||||
#[error("new-mapping rate limit reached")]
|
||||
RateLimited,
|
||||
}
|
||||
|
||||
/// State of a virtual IP mapping.
|
||||
@@ -217,6 +233,67 @@ impl ConntrackReadLog {
|
||||
}
|
||||
}
|
||||
|
||||
/// Token bucket for new mappings.
|
||||
///
|
||||
/// The level is kept in token-nanoseconds so refill is exact integer
|
||||
/// arithmetic: one token is `NANOS` units, and each elapsed nanosecond adds
|
||||
/// `rate` units.
|
||||
#[derive(Debug)]
|
||||
struct Bucket {
|
||||
/// Current level, in units of `1 / NANOS` token.
|
||||
level: u128,
|
||||
/// Level when full.
|
||||
capacity: u128,
|
||||
/// Tokens added per second.
|
||||
rate: u128,
|
||||
/// When the level was last brought up to date; unset until first use.
|
||||
last: Option<Instant>,
|
||||
}
|
||||
|
||||
impl Bucket {
|
||||
const NANOS: u128 = 1_000_000_000;
|
||||
|
||||
/// A full bucket of `capacity` tokens refilling at `rate` per second.
|
||||
fn new(capacity: u32, rate: u32) -> Self {
|
||||
let capacity = u128::from(capacity) * Self::NANOS;
|
||||
Self {
|
||||
level: capacity,
|
||||
capacity,
|
||||
rate: u128::from(rate),
|
||||
last: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Add what has accrued since the last refill, up to capacity.
|
||||
fn refill(&mut self, now: Instant) {
|
||||
if let Some(last) = self.last {
|
||||
let elapsed = now.saturating_duration_since(last).as_nanos();
|
||||
self.level = self
|
||||
.level
|
||||
.saturating_add(elapsed.saturating_mul(self.rate))
|
||||
.min(self.capacity);
|
||||
}
|
||||
// Never move backwards, so a stale `now` cannot credit time twice.
|
||||
self.last = Some(self.last.map_or(now, |last| last.max(now)));
|
||||
}
|
||||
|
||||
/// Whether at least one whole token is available.
|
||||
fn has_token(&self) -> bool {
|
||||
self.level >= Self::NANOS
|
||||
}
|
||||
|
||||
/// Spend one token; the caller has checked `has_token`.
|
||||
fn take(&mut self) {
|
||||
self.level = self.level.saturating_sub(Self::NANOS);
|
||||
}
|
||||
|
||||
/// Whole tokens available.
|
||||
#[cfg(test)]
|
||||
fn tokens(&self) -> u128 {
|
||||
self.level / Self::NANOS
|
||||
}
|
||||
}
|
||||
|
||||
/// Virtual IP pool manager.
|
||||
pub struct VirtualIpPool {
|
||||
/// Available addresses (free pool).
|
||||
@@ -231,11 +308,38 @@ pub struct VirtualIpPool {
|
||||
grace_secs: u64,
|
||||
/// Total pool size.
|
||||
total: usize,
|
||||
/// Most live mappings admitted before new names are refused.
|
||||
ceiling: usize,
|
||||
/// Rate limit on new mappings.
|
||||
bucket: Bucket,
|
||||
}
|
||||
|
||||
impl VirtualIpPool {
|
||||
/// Create a new pool from a CIDR string (e.g., `fd01::/112`).
|
||||
/// Create a new pool from a CIDR string (e.g., `fd01::/112`), with the
|
||||
/// compiled-in admission limits.
|
||||
pub fn new(cidr: &str, ttl_secs: u64, grace_secs: u64) -> Result<Self, PoolError> {
|
||||
Self::with_limits(
|
||||
cidr,
|
||||
ttl_secs,
|
||||
grace_secs,
|
||||
MAPPING_CEILING,
|
||||
MAPPING_BURST,
|
||||
MAPPING_RATE,
|
||||
)
|
||||
}
|
||||
|
||||
/// Create a pool with explicit admission limits.
|
||||
///
|
||||
/// Production uses `new`; this exists so tests can set limits small
|
||||
/// enough to reach without allocating the compiled-in counts.
|
||||
pub fn with_limits(
|
||||
cidr: &str,
|
||||
ttl_secs: u64,
|
||||
grace_secs: u64,
|
||||
ceiling: usize,
|
||||
burst: u32,
|
||||
rate: u32,
|
||||
) -> Result<Self, PoolError> {
|
||||
let (base, prefix_len) = parse_ipv6_cidr(cidr)?;
|
||||
if prefix_len == 0 || prefix_len > 128 {
|
||||
return Err(PoolError::InvalidPrefix);
|
||||
@@ -267,6 +371,8 @@ impl VirtualIpPool {
|
||||
ttl_secs,
|
||||
grace_secs,
|
||||
total,
|
||||
ceiling,
|
||||
bucket: Bucket::new(burst, rate),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -292,6 +398,21 @@ impl VirtualIpPool {
|
||||
node_addr: NodeAddr,
|
||||
mesh_addr: Ipv6Addr,
|
||||
dns_name: &str,
|
||||
) -> Result<(Ipv6Addr, bool), PoolError> {
|
||||
self.allocate_at(node_addr, mesh_addr, dns_name, Instant::now())
|
||||
}
|
||||
|
||||
/// `allocate` at a given instant, which drives the rate limit's refill
|
||||
/// and stamps a new mapping.
|
||||
///
|
||||
/// An existing mapping is returned before either limit is consulted, so a
|
||||
/// name already in use keeps resolving when new names are refused.
|
||||
pub fn allocate_at(
|
||||
&mut self,
|
||||
node_addr: NodeAddr,
|
||||
mesh_addr: Ipv6Addr,
|
||||
dns_name: &str,
|
||||
now: Instant,
|
||||
) -> Result<(Ipv6Addr, bool), PoolError> {
|
||||
// Idempotent: return existing mapping, refreshed.
|
||||
if self.refresh_if_present(node_addr)
|
||||
@@ -300,12 +421,21 @@ impl VirtualIpPool {
|
||||
return Ok((mapping.virtual_ip, false));
|
||||
}
|
||||
|
||||
// Ceiling first, so a refusal there costs no token and names the
|
||||
// ceiling whatever the bucket holds.
|
||||
if self.mappings.len() >= self.ceiling {
|
||||
return Err(PoolError::AtCeiling(self.mappings.len()));
|
||||
}
|
||||
self.bucket.refill(now);
|
||||
if !self.bucket.has_token() {
|
||||
return Err(PoolError::RateLimited);
|
||||
}
|
||||
let virtual_ip = self
|
||||
.available
|
||||
.pop_front()
|
||||
.ok_or(PoolError::Exhausted(self.mappings.len()))?;
|
||||
self.bucket.take();
|
||||
|
||||
let now = Instant::now();
|
||||
let mapping = VirtualIpMapping {
|
||||
node_addr,
|
||||
virtual_ip,
|
||||
@@ -490,6 +620,7 @@ fn parse_ipv6_cidr(cidr: &str) -> Result<(Ipv6Addr, u32), PoolError> {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use std::time::Duration;
|
||||
|
||||
/// Session counts a test sets directly, handed to `tick` as the snapshot
|
||||
/// the tick task would have read from conntrack.
|
||||
@@ -589,6 +720,102 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
/// A `/120` pool with the given limits, TTL and grace of 60 s.
|
||||
fn limited_pool(ceiling: usize, burst: u32, rate: u32) -> VirtualIpPool {
|
||||
VirtualIpPool::with_limits("fd01::/120", 60, 60, ceiling, burst, rate).unwrap()
|
||||
}
|
||||
|
||||
/// Allocate node `i` at `now`.
|
||||
fn alloc(pool: &mut VirtualIpPool, i: u8, now: Instant) -> Result<(Ipv6Addr, bool), PoolError> {
|
||||
pool.allocate_at(make_node_addr(i), make_mesh_addr(i), "test.fips", now)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn ceiling_refuses_a_new_name_without_a_token_and_keeps_existing_names() {
|
||||
let t0 = Instant::now();
|
||||
let mut pool = limited_pool(3, 10, 1);
|
||||
let mut vips = Vec::new();
|
||||
for i in 1..=3u8 {
|
||||
vips.push(alloc(&mut pool, i, t0).unwrap().0);
|
||||
}
|
||||
assert_eq!(pool.bucket.tokens(), 7);
|
||||
|
||||
assert!(
|
||||
matches!(alloc(&mut pool, 4, t0), Err(PoolError::AtCeiling(3))),
|
||||
"a fourth new name must be refused at a ceiling of 3"
|
||||
);
|
||||
assert_eq!(
|
||||
pool.bucket.tokens(),
|
||||
7,
|
||||
"a ceiling refusal must not take a token"
|
||||
);
|
||||
assert_eq!(
|
||||
alloc(&mut pool, 2, t0).unwrap(),
|
||||
(vips[1], false),
|
||||
"a name that already has a mapping must still resolve at the ceiling"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn ceiling_is_checked_before_the_rate_limit() {
|
||||
let t0 = Instant::now();
|
||||
// The bucket empties exactly as the ceiling is reached.
|
||||
let mut pool = limited_pool(3, 3, 1);
|
||||
for i in 1..=3u8 {
|
||||
alloc(&mut pool, i, t0).unwrap();
|
||||
}
|
||||
assert_eq!(pool.bucket.tokens(), 0);
|
||||
assert!(
|
||||
matches!(alloc(&mut pool, 4, t0), Err(PoolError::AtCeiling(3))),
|
||||
"a name refused at the ceiling must report the ceiling, not the rate"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rate_limit_refuses_a_burst_keeps_existing_names_and_refills() {
|
||||
let t0 = Instant::now();
|
||||
let mut pool = limited_pool(100, 2, 1);
|
||||
let (vip1, _) = alloc(&mut pool, 1, t0).unwrap();
|
||||
alloc(&mut pool, 2, t0).unwrap();
|
||||
assert!(
|
||||
matches!(alloc(&mut pool, 3, t0), Err(PoolError::RateLimited)),
|
||||
"a third new name at the same instant must be refused by a burst of 2"
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
alloc(&mut pool, 1, t0).unwrap(),
|
||||
(vip1, false),
|
||||
"an existing name must resolve with the bucket empty"
|
||||
);
|
||||
assert_eq!(pool.bucket.tokens(), 0);
|
||||
assert!(
|
||||
matches!(alloc(&mut pool, 3, t0), Err(PoolError::RateLimited)),
|
||||
"resolving an existing name must not have freed a token"
|
||||
);
|
||||
|
||||
let (_, is_new) = alloc(&mut pool, 3, t0 + Duration::from_secs(1)).unwrap();
|
||||
assert!(is_new, "one refill interval later a new name must allocate");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn exhausted_pool_takes_no_token() {
|
||||
let t0 = Instant::now();
|
||||
// /126 = 3 usable addresses.
|
||||
let mut pool = VirtualIpPool::with_limits("fd01::/126", 60, 60, 100, 10, 1).unwrap();
|
||||
for i in 1..=3u8 {
|
||||
alloc(&mut pool, i, t0).unwrap();
|
||||
}
|
||||
assert!(matches!(
|
||||
alloc(&mut pool, 4, t0),
|
||||
Err(PoolError::Exhausted(3))
|
||||
));
|
||||
assert_eq!(
|
||||
pool.bucket.tokens(),
|
||||
7,
|
||||
"a refusal for an exhausted pool must not take a token"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_mapping_lifecycle_allocated_to_free() {
|
||||
let mut pool = VirtualIpPool::new("fd01::/120", 1, 1).unwrap();
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user