diff --git a/CHANGELOG.md b/CHANGELOG.md index e843a567..516c3101 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -159,6 +159,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 the batch and requests one acknowledgement per batch, and NAT errors now name the kernel errno. A rebuild that still fails is logged, and the next successful rebuild installs the mapping. +- The gateway's virtual-IP pool now limits how many mappings it holds and how + fast it creates them. Any host that can reach the LAN resolver could ask for + one new `.fips` name after another, and each got a mapping until the 65,535 + addresses ran out, while every mapping made each NAT rebuild, each pool tick + and shutdown slower. The pool now refuses a new name once it holds 1000 live + mappings, and admits new names at 10 per second after a burst of 50. A + refused query gets SERVFAIL, and the gateway's "Pool allocation failed" + warning says which limit refused it. A name that already has a mapping is + answered before either limit is consulted, so names in use keep resolving + when the pool is full. The limits are compiled in, not configured. - A new OpenWrt install no longer enables and starts `fips-gateway`. The generated postinst turned it on unconditionally, contradicting the init script's own header, the package README and the deployment tutorial, all of diff --git a/src/bin/fips-gateway.rs b/src/bin/fips-gateway.rs index 73b8743a..def8029e 100644 --- a/src/bin/fips-gateway.rs +++ b/src/bin/fips-gateway.rs @@ -53,6 +53,12 @@ fn main() { std::process::exit(1); } +/// Microseconds since `started`, saturating, for the timing fields on debug +/// lines. +fn elapsed_us(started: Instant) -> u64 { + u64::try_from(started.elapsed().as_micros()).unwrap_or(u64::MAX) +} + /// Take a conntrack snapshot off the runtime thread. /// /// A failed read yields an empty snapshot, so every mapping reads zero @@ -446,14 +452,19 @@ async fn main() { // the pool lock: the runtime is current-thread, so a // blocking read here would stall the DNS resolver, and the // read must not happen under the lock the resolver needs. + let read_started = Instant::now(); let conntrack = read_conntrack(&mut conntrack_log).await; + let read_us = elapsed_us(read_started); let mut pool_guard = tick_pool.lock().await; + let tick_started = Instant::now(); let events = pool_guard.tick(now, &conntrack); + let tick_us = elapsed_us(tick_started); // Build snapshot for control socket let pool_status = pool_guard.status(); let mappings = pool_guard.mapping_info(now); drop(pool_guard); + debug!(mappings = mappings.len(), read_us, tick_us, "Pool tick"); let snapshot = control::build_snapshot( pool_status, diff --git a/src/gateway/nat.rs b/src/gateway/nat.rs index 1dade01b..7a0d2770 100644 --- a/src/gateway/nat.rs +++ b/src/gateway/nat.rs @@ -7,6 +7,7 @@ use std::collections::HashMap; use std::fmt; use std::net::Ipv6Addr; use std::os::fd::{AsRawFd, FromRawFd, OwnedFd}; +use std::time::Instant; use tracing::{debug, info}; use rustables::expr::{ @@ -242,14 +243,26 @@ impl NatManager { mesh_addr, }, ); - self.rebuild()?; - - debug!( - virtual_ip = %virtual_ip, - mesh_addr = %mesh_addr, - "Added DNAT/SNAT rules" - ); - Ok(()) + let (result, elapsed_us) = self.timed_rebuild(); + let mappings = self.mappings.len(); + match &result { + Ok(()) => debug!( + virtual_ip = %virtual_ip, + mesh_addr = %mesh_addr, + mappings, + elapsed_us, + "Added DNAT/SNAT rules" + ), + Err(e) => debug!( + virtual_ip = %virtual_ip, + mesh_addr = %mesh_addr, + mappings, + elapsed_us, + error = %e, + "Added DNAT/SNAT rules" + ), + } + result } /// Remove DNAT and SNAT rules for a virtual IP mapping. @@ -257,10 +270,24 @@ impl NatManager { if self.mappings.remove(&virtual_ip).is_none() { return Err(NatError::RuleNotFound(virtual_ip)); } - self.rebuild()?; - - debug!(virtual_ip = %virtual_ip, "Removed DNAT/SNAT rules"); - Ok(()) + let (result, elapsed_us) = self.timed_rebuild(); + let mappings = self.mappings.len(); + match &result { + Ok(()) => debug!( + virtual_ip = %virtual_ip, + mappings, + elapsed_us, + "Removed DNAT/SNAT rules" + ), + Err(e) => debug!( + virtual_ip = %virtual_ip, + mappings, + elapsed_us, + error = %e, + "Removed DNAT/SNAT rules" + ), + } + result } /// Flush all rules and delete the nftables table. @@ -467,6 +494,15 @@ impl NatManager { } Ok(()) } + + /// Rebuild, returning the outcome with the time the rebuild took in + /// microseconds, so a mapping change can log its cost on either path. + fn timed_rebuild(&self) -> (Result<(), NatError>, u64) { + let started = Instant::now(); + let result = self.rebuild(); + let elapsed_us = u64::try_from(started.elapsed().as_micros()).unwrap_or(u64::MAX); + (result, elapsed_us) + } } /// Largest batch, in bytes, that the kernel can admit in one send. diff --git a/src/gateway/pool.rs b/src/gateway/pool.rs index 900398d4..7776b160 100644 --- a/src/gateway/pool.rs +++ b/src/gateway/pool.rs @@ -10,6 +10,18 @@ use std::net::Ipv6Addr; use std::time::Instant; use tracing::{debug, info}; +/// Most live mappings the pool holds before it refuses new names. +/// +/// Every mapping adds rules to the NAT table, which is rebuilt whole on each +/// change, and work to every tick and to shutdown, so this bounds all three. +pub const MAPPING_CEILING: usize = 1000; + +/// New mappings the pool admits in a burst, when idle long enough to refill. +pub const MAPPING_BURST: u32 = 50; + +/// New mappings per second the pool admits once a burst is spent. +pub const MAPPING_RATE: u32 = 10; + /// Errors from pool operations. #[derive(Debug, thiserror::Error)] pub enum PoolError { @@ -19,6 +31,10 @@ pub enum PoolError { Exhausted(usize), #[error("prefix length must be between 1 and 128")] InvalidPrefix, + #[error("live-mapping ceiling reached ({0} mappings)")] + AtCeiling(usize), + #[error("new-mapping rate limit reached")] + RateLimited, } /// State of a virtual IP mapping. @@ -217,6 +233,67 @@ impl ConntrackReadLog { } } +/// Token bucket for new mappings. +/// +/// The level is kept in token-nanoseconds so refill is exact integer +/// arithmetic: one token is `NANOS` units, and each elapsed nanosecond adds +/// `rate` units. +#[derive(Debug)] +struct Bucket { + /// Current level, in units of `1 / NANOS` token. + level: u128, + /// Level when full. + capacity: u128, + /// Tokens added per second. + rate: u128, + /// When the level was last brought up to date; unset until first use. + last: Option, +} + +impl Bucket { + const NANOS: u128 = 1_000_000_000; + + /// A full bucket of `capacity` tokens refilling at `rate` per second. + fn new(capacity: u32, rate: u32) -> Self { + let capacity = u128::from(capacity) * Self::NANOS; + Self { + level: capacity, + capacity, + rate: u128::from(rate), + last: None, + } + } + + /// Add what has accrued since the last refill, up to capacity. + fn refill(&mut self, now: Instant) { + if let Some(last) = self.last { + let elapsed = now.saturating_duration_since(last).as_nanos(); + self.level = self + .level + .saturating_add(elapsed.saturating_mul(self.rate)) + .min(self.capacity); + } + // Never move backwards, so a stale `now` cannot credit time twice. + self.last = Some(self.last.map_or(now, |last| last.max(now))); + } + + /// Whether at least one whole token is available. + fn has_token(&self) -> bool { + self.level >= Self::NANOS + } + + /// Spend one token; the caller has checked `has_token`. + fn take(&mut self) { + self.level = self.level.saturating_sub(Self::NANOS); + } + + /// Whole tokens available. + #[cfg(test)] + fn tokens(&self) -> u128 { + self.level / Self::NANOS + } +} + /// Virtual IP pool manager. pub struct VirtualIpPool { /// Available addresses (free pool). @@ -231,11 +308,38 @@ pub struct VirtualIpPool { grace_secs: u64, /// Total pool size. total: usize, + /// Most live mappings admitted before new names are refused. + ceiling: usize, + /// Rate limit on new mappings. + bucket: Bucket, } impl VirtualIpPool { - /// Create a new pool from a CIDR string (e.g., `fd01::/112`). + /// Create a new pool from a CIDR string (e.g., `fd01::/112`), with the + /// compiled-in admission limits. pub fn new(cidr: &str, ttl_secs: u64, grace_secs: u64) -> Result { + Self::with_limits( + cidr, + ttl_secs, + grace_secs, + MAPPING_CEILING, + MAPPING_BURST, + MAPPING_RATE, + ) + } + + /// Create a pool with explicit admission limits. + /// + /// Production uses `new`; this exists so tests can set limits small + /// enough to reach without allocating the compiled-in counts. + pub fn with_limits( + cidr: &str, + ttl_secs: u64, + grace_secs: u64, + ceiling: usize, + burst: u32, + rate: u32, + ) -> Result { let (base, prefix_len) = parse_ipv6_cidr(cidr)?; if prefix_len == 0 || prefix_len > 128 { return Err(PoolError::InvalidPrefix); @@ -267,6 +371,8 @@ impl VirtualIpPool { ttl_secs, grace_secs, total, + ceiling, + bucket: Bucket::new(burst, rate), }) } @@ -292,6 +398,21 @@ impl VirtualIpPool { node_addr: NodeAddr, mesh_addr: Ipv6Addr, dns_name: &str, + ) -> Result<(Ipv6Addr, bool), PoolError> { + self.allocate_at(node_addr, mesh_addr, dns_name, Instant::now()) + } + + /// `allocate` at a given instant, which drives the rate limit's refill + /// and stamps a new mapping. + /// + /// An existing mapping is returned before either limit is consulted, so a + /// name already in use keeps resolving when new names are refused. + pub fn allocate_at( + &mut self, + node_addr: NodeAddr, + mesh_addr: Ipv6Addr, + dns_name: &str, + now: Instant, ) -> Result<(Ipv6Addr, bool), PoolError> { // Idempotent: return existing mapping, refreshed. if self.refresh_if_present(node_addr) @@ -300,12 +421,21 @@ impl VirtualIpPool { return Ok((mapping.virtual_ip, false)); } + // Ceiling first, so a refusal there costs no token and names the + // ceiling whatever the bucket holds. + if self.mappings.len() >= self.ceiling { + return Err(PoolError::AtCeiling(self.mappings.len())); + } + self.bucket.refill(now); + if !self.bucket.has_token() { + return Err(PoolError::RateLimited); + } let virtual_ip = self .available .pop_front() .ok_or(PoolError::Exhausted(self.mappings.len()))?; + self.bucket.take(); - let now = Instant::now(); let mapping = VirtualIpMapping { node_addr, virtual_ip, @@ -490,6 +620,7 @@ fn parse_ipv6_cidr(cidr: &str) -> Result<(Ipv6Addr, u32), PoolError> { #[cfg(test)] mod tests { use super::*; + use std::time::Duration; /// Session counts a test sets directly, handed to `tick` as the snapshot /// the tick task would have read from conntrack. @@ -589,6 +720,102 @@ mod tests { ); } + /// A `/120` pool with the given limits, TTL and grace of 60 s. + fn limited_pool(ceiling: usize, burst: u32, rate: u32) -> VirtualIpPool { + VirtualIpPool::with_limits("fd01::/120", 60, 60, ceiling, burst, rate).unwrap() + } + + /// Allocate node `i` at `now`. + fn alloc(pool: &mut VirtualIpPool, i: u8, now: Instant) -> Result<(Ipv6Addr, bool), PoolError> { + pool.allocate_at(make_node_addr(i), make_mesh_addr(i), "test.fips", now) + } + + #[test] + fn ceiling_refuses_a_new_name_without_a_token_and_keeps_existing_names() { + let t0 = Instant::now(); + let mut pool = limited_pool(3, 10, 1); + let mut vips = Vec::new(); + for i in 1..=3u8 { + vips.push(alloc(&mut pool, i, t0).unwrap().0); + } + assert_eq!(pool.bucket.tokens(), 7); + + assert!( + matches!(alloc(&mut pool, 4, t0), Err(PoolError::AtCeiling(3))), + "a fourth new name must be refused at a ceiling of 3" + ); + assert_eq!( + pool.bucket.tokens(), + 7, + "a ceiling refusal must not take a token" + ); + assert_eq!( + alloc(&mut pool, 2, t0).unwrap(), + (vips[1], false), + "a name that already has a mapping must still resolve at the ceiling" + ); + } + + #[test] + fn ceiling_is_checked_before_the_rate_limit() { + let t0 = Instant::now(); + // The bucket empties exactly as the ceiling is reached. + let mut pool = limited_pool(3, 3, 1); + for i in 1..=3u8 { + alloc(&mut pool, i, t0).unwrap(); + } + assert_eq!(pool.bucket.tokens(), 0); + assert!( + matches!(alloc(&mut pool, 4, t0), Err(PoolError::AtCeiling(3))), + "a name refused at the ceiling must report the ceiling, not the rate" + ); + } + + #[test] + fn rate_limit_refuses_a_burst_keeps_existing_names_and_refills() { + let t0 = Instant::now(); + let mut pool = limited_pool(100, 2, 1); + let (vip1, _) = alloc(&mut pool, 1, t0).unwrap(); + alloc(&mut pool, 2, t0).unwrap(); + assert!( + matches!(alloc(&mut pool, 3, t0), Err(PoolError::RateLimited)), + "a third new name at the same instant must be refused by a burst of 2" + ); + + assert_eq!( + alloc(&mut pool, 1, t0).unwrap(), + (vip1, false), + "an existing name must resolve with the bucket empty" + ); + assert_eq!(pool.bucket.tokens(), 0); + assert!( + matches!(alloc(&mut pool, 3, t0), Err(PoolError::RateLimited)), + "resolving an existing name must not have freed a token" + ); + + let (_, is_new) = alloc(&mut pool, 3, t0 + Duration::from_secs(1)).unwrap(); + assert!(is_new, "one refill interval later a new name must allocate"); + } + + #[test] + fn exhausted_pool_takes_no_token() { + let t0 = Instant::now(); + // /126 = 3 usable addresses. + let mut pool = VirtualIpPool::with_limits("fd01::/126", 60, 60, 100, 10, 1).unwrap(); + for i in 1..=3u8 { + alloc(&mut pool, i, t0).unwrap(); + } + assert!(matches!( + alloc(&mut pool, 4, t0), + Err(PoolError::Exhausted(3)) + )); + assert_eq!( + pool.bucket.tokens(), + 7, + "a refusal for an exhausted pool must not take a token" + ); + } + #[test] fn test_mapping_lifecycle_allocated_to_free() { let mut pool = VirtualIpPool::new("fd01::/120", 1, 1).unwrap(); diff --git a/testing/static/docker-compose.yml b/testing/static/docker-compose.yml index 77010007..eb025aee 100644 --- a/testing/static/docker-compose.yml +++ b/testing/static/docker-compose.yml @@ -375,7 +375,9 @@ services: # attached after container start) and manage nftables NAT rules. privileged: true environment: - - RUST_LOG=info + # Debug for the NAT rebuild and pool tick timing lines only; the + # gateway's --log-level is ignored while RUST_LOG is set. + - RUST_LOG=info,fips::gateway::nat=debug,fips_gateway=debug - FIPS_TEST_MODE=gateway sysctls: - net.ipv6.conf.all.disable_ipv6=0 diff --git a/testing/static/scripts/gateway-test.sh b/testing/static/scripts/gateway-test.sh index 331037ea..3153bf71 100755 --- a/testing/static/scripts/gateway-test.sh +++ b/testing/static/scripts/gateway-test.sh @@ -477,6 +477,191 @@ else check "Gateway shutdown (no completion message in logs)" 1 fi +# ── Long-lived gateway restart, shared by phases 11 and 12 ─────────────── + +# Start the stopped gateway with mappings that outlive the phase, and gate on +# its readiness. A failure is recorded through check under the label prefix +# $1 and returns 1, and the caller ends its phase. On success it sets: +# GW_STARTED the container's start time, for `docker logs --since`, since +# the log still holds every earlier phase +# GW_T0 epoch seconds at the start +# GW_BASELINE allocations after the readiness probe, which allocates +# GW_PROBE the address the readiness probe got for NPUB_B +gw_long_lived_start() { + local prefix="$1" + local config_file="$GENERATED_DIR/gateway/node-a.yaml" + local expect_rev + expect_rev=$(git -C "$SCRIPT_DIR" rev-parse --short=10 HEAD) + + # Rewrite in place (same inode): the container sees the host file through + # a single-file bind mount, which a replace-by-rename would leave behind. + python3 - "$config_file" <<'PYEOF' +import sys, yaml +path = sys.argv[1] +with open(path, "r+") as f: + cfg = yaml.safe_load(f) + cfg["gateway"]["dns"]["ttl"] = 1800 + cfg["gateway"]["pool_grace_period"] = 1800 + f.seek(0) + yaml.dump(cfg, f, default_flow_style=False, sort_keys=False) + f.truncate() +PYEOF + + docker start "$GATEWAY" >/dev/null + GW_STARTED=$(docker inspect -f '{{.State.StartedAt}}' "$GATEWAY") + GW_T0=$(date -u +%s) + echo " Gateway started at $GW_STARTED (expect rev $expect_rev)" + + local seen_ttl seen_grace + seen_ttl=$(docker exec "$GATEWAY" grep -c "ttl: 1800" /etc/fips/fips.yaml || true) + seen_grace=$(docker exec "$GATEWAY" grep -c "pool_grace_period: 1800" /etc/fips/fips.yaml || true) + if [ "$seen_ttl" -ge 1 ] && [ "$seen_grace" -ge 1 ]; then + check "$prefix: container sees ttl 1800 and grace 1800" 0 + else + check "$prefix: container config rewrite (ttl: $seen_ttl, grace: $seen_grace)" 1 + return 1 + fi + + # Readiness is a hard gate here, unlike phases 1 and 2. + if wait_for_peers "$GATEWAY" 2 60; then + check "$prefix: gateway peers after restart" 0 + else + check "$prefix: gateway peers after restart" 1 + return 1 + fi + local probe + GW_PROBE="" + for _ in $(seq 1 60); do + probe=$(docker exec "$CLIENT" dig +short AAAA "${NPUB_B}.fips" @${GW_DNS} 2>/dev/null || true) + GW_PROBE=$(grep -m1 "^fd01::" <<< "$probe" || true) + if [ -n "$GW_PROBE" ]; then + break + fi + sleep 1 + done + if [ -n "$GW_PROBE" ]; then + check "$prefix: gateway DNS answers after restart" 0 + else + check "$prefix: gateway DNS answers after restart" 1 + return 1 + fi + + sleep 1 + local started_log rev_lines + started_log=$(docker logs --timestamps --since "$GW_STARTED" "$GATEWAY" 2>&1) + # The co-resident daemon logs its own "(rev ...) starting" line, so + # anchor on the gateway's name. + rev_lines=$(grep -cE "fips-gateway [^ ]+ \(rev ${expect_rev}\) starting" <<< "$started_log" || true) + if [ "$rev_lines" -eq 1 ]; then + check "$prefix: startup line reads rev ${expect_rev}) with no -dirty" 0 + else + check "$prefix: startup line for rev ${expect_rev} (found $rev_lines)" 1 + return 1 + fi + GW_BASELINE=$(grep -c "Allocated virtual IP" <<< "$started_log" || true) + return 0 +} + +# The pool's compiled-in admission limits. A new name is refused past +# MAPPING_CEILING live mappings, and past a burst of MAPPING_BURST new names +# are admitted at MAPPING_RATE per second, so phases that create many +# mappings retry the rate limit's refusals. +POOL_RS="$SCRIPT_DIR/../../../src/gateway/pool.rs" + +# A compiled-in constant from pool.rs, or nothing if the line is not found. +pool_const() { + sed -nE "s/^pub const $1: [a-z0-9]+ = ([0-9]+);$/\1/p" "$POOL_RS" +} + +POOL_CEILING=$(pool_const MAPPING_CEILING) +POOL_BURST=$(pool_const MAPPING_BURST) +POOL_RATE=$(pool_const MAPPING_RATE) + +# Per-name retry bound for the driver, in seconds: a fixed 30 s. A bound +# sized from the rate (workers x refill interval x 5) is 2 s at 10/s, shorter +# than two of the driver's 1 to 1.5 s retry pauses, and never above 20 s for +# any rate of 1/s or more; 30 s lets a name wait through many pauses. No run +# has had a name reach it. Empty when the rate could not be read, which the +# phases report. +gw_retry_bound() { + [ -n "$POOL_RATE" ] && [ "$POOL_RATE" -gt 0 ] || return 0 + echo 30 + return 0 +} + +# Install the AAAA driver in the client: 4 closed-loop workers, one fresh +# socket per query, counts printed as key=value. Given a retry bound in +# seconds, a name answered SERVFAIL or not answered is asked again after a +# pause of 1 to 1.5 s until it is answered or the bound has passed since its +# first query; the exit status is 1 if any name was left unplaced. Without a +# bound each name is asked once and the exit status is 0. +gw_install_driver() { + docker exec -i "$CLIENT" sh -c 'cat > /tmp/gw_driver.py' <<'PYEOF' +import random, socket, struct, sys, threading, time +server = sys.argv[1] +bound = float(sys.argv[2]) if len(sys.argv) > 2 else None +names = [n.strip() for n in sys.stdin if n.strip()] +lock = threading.Lock() +counts = {"answered": 0, "servfail": 0, "timeout": 0, "other": 0, + "retries": 0, "unplaced": 0} +def query(name): + qid = random.getrandbits(16) + pkt = struct.pack(">HHHHHH", qid, 0x0100, 1, 0, 0, 0) + for label in (name + ".fips").split("."): + raw = label.encode() + pkt += bytes([len(raw)]) + raw + pkt += b"\x00" + struct.pack(">HH", 28, 1) + s = socket.socket(socket.AF_INET6, socket.SOCK_DGRAM) + s.settimeout(6) + try: + s.sendto(pkt, (server, 53)) + while True: + data, _ = s.recvfrom(4096) + if len(data) >= 12 and struct.unpack(">H", data[:2])[0] == qid: + break + except socket.timeout: + return "timeout" + finally: + s.close() + flags, _, ancount = struct.unpack(">HHH", data[2:8]) + rcode = flags & 0xF + if rcode == 2: + return "servfail" + if rcode == 0 and ancount > 0: + return "answered" + return "other" +def place(name): + first = time.monotonic() + while True: + outcome = query(name) + with lock: + counts[outcome] += 1 + if bound is None or outcome == "answered": + return + if outcome == "other" or time.monotonic() - first > bound: + with lock: + counts["unplaced"] += 1 + return + with lock: + counts["retries"] += 1 + time.sleep(1 + random.random() * 0.5) +def worker(): + while True: + with lock: + if not names: + return + name = names.pop() + place(name) +threads = [threading.Thread(target=worker) for _ in range(4)] +for t in threads: + t.start() +for t in threads: + t.join() +print(" ".join(f"{k}={v}" for k, v in counts.items())) +sys.exit(1 if counts["unplaced"] else 0) +PYEOF +} + # Phase 11: NAT rebuild past the default netlink socket limits # # Every change rebuilds the whole fips_gateway table in one netlink batch. @@ -484,10 +669,11 @@ fi # (the acks overflowed the receive buffer, after the commit) and past about # 313 (the batch overflowed the send buffer, and nothing was committed). # Drive 400 new names through a gateway whose mappings outlive the phase and -# judge the result on the kernel's own table, not on the daemon's debug-level -# success line, which the suite's log level does not show. Runs after phase -# 10, so the gateway container is stopped when it starts, and it leaves it -# stopped. +# judge the result on the kernel's own table, which is what a lost rebuild +# leaves wrong, not on the daemon's debug-level success line. The names go +# through the pool's rate limit, so the driver retries its refusals. Runs +# after phase 10, so the gateway container is stopped when it starts, and it +# leaves it stopped. echo "" echo "Phase 11: NAT rebuild past default socket limits" @@ -519,73 +705,10 @@ natbig_kernel() { } natbig_phase() { - local config_file="$GENERATED_DIR/gateway/node-a.yaml" - local expect_rev - expect_rev=$(git -C "$SCRIPT_DIR" rev-parse --short=10 HEAD) - - # Rewrite in place (same inode): the container sees the host file through - # a single-file bind mount, which a replace-by-rename would leave behind. - python3 - "$config_file" <<'PYEOF' -import sys, yaml -path = sys.argv[1] -with open(path, "r+") as f: - cfg = yaml.safe_load(f) - cfg["gateway"]["dns"]["ttl"] = 1800 - cfg["gateway"]["pool_grace_period"] = 1800 - f.seek(0) - yaml.dump(cfg, f, default_flow_style=False, sort_keys=False) - f.truncate() -PYEOF - - docker start "$GATEWAY" >/dev/null - NATBIG_STARTED=$(docker inspect -f '{{.State.StartedAt}}' "$GATEWAY") - local t0 - t0=$(natbig_now) - echo " Gateway started at $NATBIG_STARTED (expect rev $expect_rev)" - - local seen_ttl seen_grace - seen_ttl=$(docker exec "$GATEWAY" grep -c "ttl: 1800" /etc/fips/fips.yaml || true) - seen_grace=$(docker exec "$GATEWAY" grep -c "pool_grace_period: 1800" /etc/fips/fips.yaml || true) - if [ "$seen_ttl" -ge 1 ] && [ "$seen_grace" -ge 1 ]; then - check "NAT batch: container sees ttl 1800 and grace 1800" 0 - else - check "NAT batch: container config rewrite (ttl: $seen_ttl, grace: $seen_grace)" 1 - return 0 - fi - - if wait_for_peers "$GATEWAY" 2 60; then - check "NAT batch: gateway peers after restart" 0 - else - check "NAT batch: gateway peers after restart" 1 - return 0 - fi - local ready=false probe - for _ in $(seq 1 60); do - probe=$(docker exec "$CLIENT" dig +short AAAA "${NPUB_B}.fips" @${GW_DNS} 2>/dev/null || true) - if echo "$probe" | grep -q "^fd01::"; then - ready=true - break - fi - sleep 1 - done - if [ "$ready" = true ]; then - check "NAT batch: gateway DNS answers after restart" 0 - else - check "NAT batch: gateway DNS answers after restart" 1 - return 0 - fi - - sleep 1 - natbig_allocated - local rev_lines - rev_lines=$(grep -cE "fips-gateway [^ ]+ \(rev ${expect_rev}\) starting" <<< "$NATBIG_LOG" || true) - if [ "$rev_lines" -eq 1 ]; then - check "NAT batch: startup line reads rev ${expect_rev}) with no -dirty" 0 - else - check "NAT batch: startup line for rev ${expect_rev} (found $rev_lines)" 1 - return 0 - fi - local baseline="$NATBIG_ALLOCATED" + gw_long_lived_start "NAT batch" || return 0 + NATBIG_STARTED="$GW_STARTED" + local t0="$GW_T0" + local baseline="$GW_BASELINE" local target=$((baseline + NATBIG_NAMES)) echo " Baseline allocations after readiness: $baseline; target $target" if [ "$baseline" -ge 1 ]; then @@ -608,56 +731,14 @@ PYEOF return 0 fi - # Closed-loop AAAA driver, 4 workers, one fresh socket per query. - docker exec -i "$CLIENT" sh -c 'cat > /tmp/gw_natbig.py' <<'PYEOF' -import random, socket, struct, sys, threading -server = sys.argv[1] -names = [n.strip() for n in sys.stdin if n.strip()] -lock = threading.Lock() -counts = {"answered": 0, "servfail": 0, "timeout": 0, "other": 0} -def query(name): - qid = random.getrandbits(16) - pkt = struct.pack(">HHHHHH", qid, 0x0100, 1, 0, 0, 0) - for label in (name + ".fips").split("."): - raw = label.encode() - pkt += bytes([len(raw)]) + raw - pkt += b"\x00" + struct.pack(">HH", 28, 1) - s = socket.socket(socket.AF_INET6, socket.SOCK_DGRAM) - s.settimeout(6) - try: - s.sendto(pkt, (server, 53)) - while True: - data, _ = s.recvfrom(4096) - if len(data) >= 12 and struct.unpack(">H", data[:2])[0] == qid: - break - except socket.timeout: - return "timeout" - finally: - s.close() - flags, _, ancount = struct.unpack(">HHH", data[2:8]) - rcode = flags & 0xF - if rcode == 2: - return "servfail" - if rcode == 0 and ancount > 0: - return "answered" - return "other" -def worker(): - while True: - with lock: - if not names: - return - name = names.pop() - outcome = query(name) - with lock: - counts[outcome] += 1 -threads = [threading.Thread(target=worker) for _ in range(4)] -for t in threads: - t.start() -for t in threads: - t.join() -print(" ".join(f"{k}={v}" for k, v in counts.items())) -PYEOF - + local retry_bound + retry_bound=$(gw_retry_bound) + if [ -z "$retry_bound" ]; then + check "NAT batch: MAPPING_RATE read from pool.rs ('$POOL_RATE')" 1 + rm -f "$names_file" + return 0 + fi + gw_install_driver local remaining out rc=0 remaining=$((NATBIG_CAP - ($(natbig_now) - t0))) # timeout reads 0 as no limit and refuses a negative duration. @@ -667,7 +748,7 @@ PYEOF return 0 fi out=$(docker exec -i "$CLIENT" timeout "$remaining" \ - python3 /tmp/gw_natbig.py "$GW_DNS" <"$names_file" 2>&1) || rc=$? + python3 /tmp/gw_driver.py "$GW_DNS" "$retry_bound" <"$names_file" 2>&1) || rc=$? rm -f "$names_file" echo " [$(($(natbig_now) - t0))s] sent $NATBIG_NAMES names: $out (rc=$rc)" natbig_allocated @@ -727,6 +808,278 @@ PYEOF natbig_phase +# Phase 12: Pool admission limits +# +# The pool refuses a new name once it holds MAPPING_CEILING live mappings, +# and past a burst of MAPPING_BURST admits new names at MAPPING_RATE per +# second; the constants are read from src/gateway/pool.rs. Restart the gateway +# with mappings that outlive the phase, fill the pool to the ceiling through +# the rate limit, then ask once each for 20 more new names: all 20 must be +# refused at the ceiling while an existing name still resolves. Without the +# ceiling those 20 are allocated, whatever the creation rate. Reports rebuild +# and tick durations on the way up and the shutdown duration at the ceiling. +# Runs after phase 11, so the gateway container is stopped when it starts, +# and it leaves it stopped. +echo "" +echo "Phase 12: Pool admission limits" + +LIMITS_EXTRA=20 +# Sized from both trees: with the limits the fill takes about +# (ceiling - burst) / rate seconds plus retry pauses, and without them a +# measurement run reached ceiling + 20 mappings far sooner. +LIMITS_CAP=300 +# Shutdown took 2.4 s at 2000 mappings in a measurement run. +LIMITS_STOP_TIME=30 +limits_now() { + date -u +%s +} + +limits_slice() { + SLICE=$(docker logs --timestamps --since "$STARTED" "$GATEWAY" 2>&1) +} + +# Figures below read lines a grep on the slice has already selected, so the +# log-string guard sees each daemon-log literal. The Python matches no log +# text itself: it takes the daemon's own timestamp and key=value fields. +LIMITS_PY_FIELDS=' +import re, sys, statistics +from datetime import datetime, timezone +ANSI = re.compile(r"\x1b\[[0-9;]*m") +FIELD = re.compile(r"\b(\w+)=(\S+)") +def parse(line): + parts = ANSI.sub("", line.rstrip("\n")).split(" ", 2) + whole, _, frac = parts[1].rstrip("Z").partition(".") + t = datetime.strptime(whole, "%Y-%m-%dT%H:%M:%S").replace(tzinfo=timezone.utc) + t = t.timestamp() + float("0." + (frac or "0")) + return t, dict(FIELD.findall(parts[2] if len(parts) > 2 else "")) +' + +limits_phase() { + local ceiling="$POOL_CEILING" burst="$POOL_BURST" rate="$POOL_RATE" + if [ -n "$ceiling" ] && [ -n "$burst" ] && [ -n "$rate" ] && [ "$rate" -gt 0 ] \ + && [ "$ceiling" -gt 1 ]; then + check "Limits: read ceiling $ceiling, burst $burst, rate $rate/s from pool.rs" 0 + else + check "Limits: constants in pool.rs (ceiling '$ceiling', burst '$burst', rate '$rate')" 1 + return 0 + fi + local retry_bound + retry_bound=$(gw_retry_bound) + + gw_long_lived_start "Limits" || return 0 + STARTED="$GW_STARTED" + local t0="$GW_T0" + local baseline="$GW_BASELINE" + if [ "$baseline" -eq 1 ]; then + check "Limits: the readiness probe holds the only mapping" 0 + else + check "Limits: mappings after readiness ($baseline, expected 1)" 1 + return 0 + fi + + # Names: real keys, since the daemon parses each one as a public key. + local fill=$((ceiling - baseline)) + local total=$((fill + LIMITS_EXTRA)) + local names_file have_names + names_file=$(mktemp) + docker exec "$GATEWAY" bash -c \ + "for i in \$(seq 1 $total); do fipsctl keygen --stdout; done" \ + | grep '^npub1' >"$names_file" || true + have_names=$(wc -l <"$names_file") + if [ "$have_names" -lt "$total" ]; then + check "Limits: generated $total names (got $have_names)" 1 + rm -f "$names_file" + return 0 + fi + gw_install_driver + + # Fill: exactly ceiling - 1 new names, each retried through the rate + # limit's refusals. A name left unplaced, or the cap firing, fails. + local remaining out rc=0 + remaining=$((LIMITS_CAP - ($(limits_now) - t0))) + if [ "$remaining" -le 0 ]; then + check "Limits: setup exceeded ${LIMITS_CAP}s cap" 1 + rm -f "$names_file" + return 0 + fi + out=$(sed -n "1,${fill}p" "$names_file" | docker exec -i "$CLIENT" timeout "$remaining" \ + python3 /tmp/gw_driver.py "$GW_DNS" "$retry_bound" 2>&1) || rc=$? + echo " [$(($(limits_now) - t0))s] fill of $fill names (retry bound ${retry_bound}s): $out (rc=$rc)" + if [ "$rc" -eq 0 ]; then + check "Limits: fill placed all $fill names" 0 + else + check "Limits: fill driver failed (rc $rc)" 1 + fi + sleep 1 + limits_slice + local rate_refused_fill + rate_refused_fill=$(grep -c "new-mapping rate limit reached" <<< "$SLICE" || true) + + # Past the ceiling: one query per name, never retried. + rc=0 + remaining=$((LIMITS_CAP - ($(limits_now) - t0))) + if [ "$remaining" -le 0 ]; then + check "Limits: fill exceeded ${LIMITS_CAP}s cap" 1 + rm -f "$names_file" + return 0 + fi + out=$(sed -n "$((fill + 1)),${total}p" "$names_file" | docker exec -i "$CLIENT" \ + timeout "$remaining" python3 /tmp/gw_driver.py "$GW_DNS" 2>&1) || rc=$? + rm -f "$names_file" + echo " [$(($(limits_now) - t0))s] $LIMITS_EXTRA names past the ceiling: $out (rc=$rc)" + local post_servfail + post_servfail=$(sed -nE 's/.*servfail=([0-9]+).*/\1/p' <<< "$out") + if [ "$rc" -eq 0 ] && [ "${post_servfail:-0}" -eq "$LIMITS_EXTRA" ]; then + check "Limits: all $LIMITS_EXTRA names past the ceiling got SERVFAIL" 0 + else + check "Limits: names past the ceiling (servfail '${post_servfail}', rc $rc)" 1 + fi + + local probe + probe=$(docker exec "$CLIENT" dig +short AAAA "${NPUB_B}.fips" @${GW_DNS} 2>/dev/null || true) + if [ "$(grep -m1 "^fd01::" <<< "$probe" || true)" = "$GW_PROBE" ]; then + check "Limits: NPUB_B still resolves to $GW_PROBE at the ceiling" 0 + else + check "Limits: NPUB_B at the ceiling (got '${probe:0:60}', had $GW_PROBE)" 1 + fi + + sleep 1 + limits_slice + local allocated reclaimed ceiling_refused rate_refused nat_fail nat_rm_fail ndp_fail + allocated=$(grep -c "Allocated virtual IP" <<< "$SLICE" || true) + reclaimed=$(grep -c "Reclaimed virtual IP" <<< "$SLICE" || true) + ceiling_refused=$(grep -c "live-mapping ceiling reached" <<< "$SLICE" || true) + rate_refused=$(grep -c "new-mapping rate limit reached" <<< "$SLICE" || true) + nat_fail=$(grep -c "Failed to add NAT rules" <<< "$SLICE" || true) + nat_rm_fail=$(grep -c "Failed to remove NAT rules" <<< "$SLICE" || true) + ndp_fail=$(grep -c "Failed to add proxy NDP" <<< "$SLICE" || true) + echo " Slice counts: allocated=$allocated reclaimed=$reclaimed" \ + "ceiling_refused=$ceiling_refused rate_refused=$rate_refused_fill/$rate_refused" \ + "nat_add_fail=$nat_fail nat_remove_fail=$nat_rm_fail ndp_fail=$ndp_fail" + if [ "$allocated" -eq "$ceiling" ]; then + check "Limits: live mappings stop at the ceiling ($allocated)" 0 + else + check "Limits: live mappings $allocated, ceiling $ceiling" 1 + fi + if [ "$allocated" -gt 0 ]; then + if [ "$reclaimed" -eq 0 ]; then + check "Limits: no reclaims during the phase" 0 + else + check "Limits: reclaims during the phase ($reclaimed)" 1 + fi + else + check "Limits: allocations present in the slice (0)" 1 + fi + if [ "$ceiling_refused" -ge "$LIMITS_EXTRA" ]; then + check "Limits: refusals logged as the ceiling ($ceiling_refused)" 0 + else + check "Limits: ceiling refusals logged ($ceiling_refused, expected >= $LIMITS_EXTRA)" 1 + fi + # The rate limit must have refused during the fill, so its count is a + # live signal, and must refuse nothing after it: past the ceiling the + # ceiling is checked first. + if [ "$rate_refused_fill" -ge 1 ]; then + check "Limits: the rate limit refused during the fill ($rate_refused_fill)" 0 + else + check "Limits: rate refusals during the fill (0)" 1 + fi + if [ "$rate_refused" -eq "$rate_refused_fill" ]; then + check "Limits: no rate refusal after the fill" 0 + else + check "Limits: rate refusals after the fill ($((rate_refused - rate_refused_fill)))" 1 + fi + if [ $((nat_fail + nat_rm_fail + ndp_fail)) -eq 0 ]; then + check "Limits: no NAT or proxy NDP failure lines" 0 + else + check "Limits: failure lines (nat add $nat_fail, remove $nat_rm_fail, ndp $ndp_fail)" 1 + fi + + # Creations in any 10 s window of the fill can be at most a full burst + # plus 10 s of refill. The bound is exact, not approximate: the bucket + # holds at most burst whole tokens at the window's first creation and + # gains one per 1/rate s, so one more creation needs the daemon's log + # timestamps to lag its clock by a full token interval (100 ms at 10/s) + # more at one end of the window than the other. Do not loosen it for skew. + local bound=$((burst + 10 * rate)) most + most=$(grep "Allocated virtual IP" <<< "$SLICE" | python3 -c "$LIMITS_PY_FIELDS +b, fill = int(sys.argv[1]), int(sys.argv[2]) +ts = sorted(parse(l)[0] for l in sys.stdin)[b:b + fill] +most = j = 0 +for i in range(len(ts)): + while ts[i] - ts[j] > 10: + j += 1 + most = max(most, i - j + 1) +print(most) +" "$baseline" "$fill" || true) + if [ -n "$most" ] && [ "$most" -le "$bound" ]; then + check "Limits: at most $most creations in any 10 s of the fill (bound $bound)" 0 + else + check "Limits: creations in a 10 s window of the fill ('$most', bound $bound)" 1 + fi + + local added_timed tick_timed + added_timed=$(grep "Added DNAT/SNAT rules" <<< "$SLICE" | grep -c "elapsed_us" || true) + tick_timed=$(grep "Pool tick" <<< "$SLICE" | grep -c "tick_us" || true) + if [ "$added_timed" -ge 1 ] && [ "$tick_timed" -ge 1 ]; then + check "Limits: rebuild and tick timing lines present" 0 + else + check "Limits: timing lines (rebuild $added_timed, tick $tick_timed)" 1 + fi + + echo " --- Rebuild duration (elapsed_us over the 20 adds ending at each count) ---" + local targets=("$ceiling") rebuild_report + [ "$ceiling" -gt 500 ] && targets=(500 "$ceiling") + rebuild_report=$(grep "Added DNAT/SNAT rules" <<< "$SLICE" | python3 -c "$LIMITS_PY_FIELDS +by_n = {} +for l in sys.stdin: + _, f = parse(l) + if 'error' in f or 'elapsed_us' not in f: + continue + by_n[int(f['mappings'])] = int(f['elapsed_us']) +missing = 0 +for t in map(int, sys.argv[1:]): + if t not in by_n: + missing += 1 + print(f' rebuild {t}: no successful add line with mappings={t}') + continue + xs = [by_n[n] for n in range(t - 19, t + 1) if n in by_n] + print(f' rebuild {t}: n={len(xs)} median={statistics.median(xs):.0f}us max={max(xs)}us') +print(f'REBUILD_MISSING={missing}') +" "${targets[@]}" || true) + echo "$rebuild_report" | grep -v '^REBUILD_MISSING=' + if echo "$rebuild_report" | grep -q '^REBUILD_MISSING=0$'; then + check "Limits: a successful add line at ${targets[*]} mappings" 0 + else + check "Limits: a successful add line at ${targets[*]} mappings" 1 + fi + + echo " --- Ticks as they fell ---" + grep "Pool tick" <<< "$SLICE" | python3 -c "$LIMITS_PY_FIELDS +for l in sys.stdin: + _, f = parse(l) + print(f\" tick: mappings={f.get('mappings')} read_us={f.get('read_us')} tick_us={f.get('tick_us')}\") +" || true + + echo " --- Shutdown at the ceiling ---" + docker stop --time="$LIMITS_STOP_TIME" "$GATEWAY" >/dev/null 2>&1 || true + limits_slice + local shut_lines + shut_lines=$( { grep "fips-gateway shutting down" <<< "$SLICE"; \ + grep "fips-gateway shutdown complete" <<< "$SLICE"; } || true) + if [ "$(echo "$shut_lines" | grep -c . || true)" -eq 2 ]; then + echo "$shut_lines" | python3 -c "$LIMITS_PY_FIELDS +ts = [parse(l)[0] for l in sys.stdin] +print(f' shutdown duration: {ts[1] - ts[0]:.3f}s') +" || true + check "Limits: gateway shutdown completed within ${LIMITS_STOP_TIME}s" 0 + else + check "Limits: gateway shutdown completion line present" 1 + fi + echo " Phase time: $(($(limits_now) - t0))s" +} + +limits_phase + echo "" echo "=== Results: $PASSED passed, $FAILED failed ===" [ "$FAILED" -eq 0 ] && exit 0 || exit 1