diff --git a/CHANGELOG.md b/CHANGELOG.md index d12fdc93..d37c7fb1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -781,6 +781,33 @@ with v0.5.x or earlier peers. costs the initiator the msg2 key agreement until the cycle ends; the msg1 resend budget bounds that. The wire format is unchanged. +#### Session rekey + +- A SessionAck that fails to read no longer ends a session rekey this node + started. The handler took the rekey handshake off the session before reading + the ack's msg2 and abandoned the rekey when the read failed, although nothing + authenticates the ack before that read and the only tie to the rekey is the + datagram's source address. The handshake now goes back rolled back to its + state before the read, so the peer's genuine ack still completes the rekey, + and the refusal is counted as `ack_handshake_failed`, as it already was for a + first-contact session. The wire format is unchanged. + +#### Session coordinates + +- A node with no coordinates cached for a session's destination no longer + sends its own coordinates in their place. The lookup that supplies them falls + back to the node's own coordinates, which a first-contact SessionSetup needs + because its destination field cannot be empty, but the established data path, + the standalone CoordsWarmup and the rekey SessionSetup used the same + fallback. Every receiver files the destination coordinates it is sent under + the destination's address, so a destination reached this way cached its own + address under the sender's coordinates. On a cache miss a data frame now goes + out without coordinates and leaves the warmup budget for the first frames + after the cache is refilled, a standalone CoordsWarmup is not sent, and a + rekey SessionSetup, which can only miss for a direct peer, carries the + coordinates that peer announced. First-contact setup is unchanged. The wire + format is unchanged. + #### Control socket - `show_links` (`fipsctl show links`) now reports the traffic a link has @@ -851,6 +878,16 @@ with v0.5.x or earlier peers. indistinguishable from an idle one. A kernel built without `CONFIG_NF_CONNTRACK_PROCFS` has no `/proc/net/nf_conntrack` at all and fails identically every tick, so a repeat is logged at debug rather than warn. +- The gateway says at startup whether it can read conntrack sessions. It + reads the table once, as each tick does, and logs either the source it read + or that no source is readable and session pinning is off. An operator on a + kernel with no readable source learned this only from a warning at the first + failed tick. +- The gateway counts sessions on a kernel without `/proc/net/nf_conntrack`. + When the file is absent it dumps the conntrack table over netlink, as + `conntrack -L` does, so a mapping carrying traffic is pinned instead of + being reclaimed on its TTL and grace period alone. Kernels built without + `CONFIG_NF_CONNTRACK_PROCFS`, such as Ubuntu's, had session pinning off. - The NAT table is rebuilt in one netlink transaction. A rebuild deleted the `fips_gateway` table in a batch of its own, discarded that batch's result, and only then sent the batch that recreated the table, the chains, the diff --git a/Cargo.lock b/Cargo.lock index 836b80dc..f20236a8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -650,6 +650,12 @@ version = "0.10.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a6ef517f0926dd24a1582492c791b6a4818a4d94e789a334894aa15b0d12f55c" +[[package]] +name = "convert_case" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6245d59a3e82a7fc217c5828a6692dbc6dfb63a0c8c90495621f7b9d79704a0e" + [[package]] name = "convert_case" version = "0.10.0" @@ -760,7 +766,7 @@ checksum = "d8b9f2e4c67f833b660cdb0a3523065869fb35570177239812ed4c905aeff87b" dependencies = [ "bitflags 2.13.1", "crossterm_winapi", - "derive_more", + "derive_more 2.1.1", "document-features", "mio", "parking_lot", @@ -966,6 +972,19 @@ version = "0.5.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7cd812cc2bc1d69d4764bd80df88b4317eaef9e773c75226407d9bc0876b211c" +[[package]] +name = "derive_more" +version = "0.99.20" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6edb4b64a43d977b8e99788fe3a04d483834fba1215a7e02caa415b626497f7f" +dependencies = [ + "convert_case 0.4.0", + "proc-macro2", + "quote", + "rustc_version", + "syn 2.0.119", +] + [[package]] name = "derive_more" version = "2.1.1" @@ -981,7 +1000,7 @@ version = "2.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "799a97264921d8623a957f6c3b9011f3b5492f557bbb7a5a19b7fa6d06ba8dcb" dependencies = [ - "convert_case", + "convert_case 0.10.0", "proc-macro2", "quote", "rustc_version", @@ -1159,6 +1178,9 @@ dependencies = [ "libc", "libm", "mdns-sd", + "netlink-packet-core", + "netlink-packet-netfilter", + "netlink-sys", "nostr", "nostr-sdk", "portable-atomic", @@ -1954,6 +1976,19 @@ dependencies = [ "paste", ] +[[package]] +name = "netlink-packet-netfilter" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "27b511a24610c054dbbfea4c3e403fe4df2e811df6bc48fa2583b1b93333f020" +dependencies = [ + "bitflags 2.13.1", + "derive_more 0.99.20", + "libc", + "netlink-packet-core", + "zerocopy", +] + [[package]] name = "netlink-packet-route" version = "0.30.0" diff --git a/Cargo.toml b/Cargo.toml index 488aee6f..55dde798 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -61,6 +61,11 @@ libc = "0.2" rtnetlink = "0.21.0" rustables = "0.8.7" procfs = { version = "0.18", default-features = false } +# Conntrack dump for kernels without /proc/net/nf_conntrack. 0.3 is the +# release on netlink-packet-core 0.8, the version rtnetlink 0.21 uses. +netlink-packet-netfilter = "0.3" +netlink-packet-core = "0.8" +netlink-sys = "0.8" # bluer/BlueZ needs glibc — see build.rs `bluer_available` cfg gate. [target.'cfg(all(target_os = "linux", not(target_env = "musl")))'.dependencies] diff --git a/docs/design/fips-gateway.md b/docs/design/fips-gateway.md index a70b0afe..51ce4714 100644 --- a/docs/design/fips-gateway.md +++ b/docs/design/fips-gateway.md @@ -263,8 +263,13 @@ Timing: cached DNS responses. - **Tick interval**: the pool re-evaluates state every 10 s. -Active session counts come from `/proc/net/nf_conntrack`: an entry -counts as a session if its original destination is the virtual IP. +Active session counts come from `/proc/net/nf_conntrack`, or, on a +kernel without that file, from a dump of the IPv6 conntrack table over +`NETLINK_NETFILTER`, the request `conntrack -L` makes. The choice is +made on every tick, and the source is logged once at startup. Either +way an entry counts once toward each distinct IPv6 destination among +its original and reply tuples, so an entry counts as a session of a +virtual IP whose address is its original destination. If the pool is exhausted, new DNS queries return `SERVFAIL`. Existing mappings are never evicted prematurely — the correctness of diff --git a/docs/how-to/troubleshoot-gateway.md b/docs/how-to/troubleshoot-gateway.md index 7c570df5..6b532ded 100644 --- a/docs/how-to/troubleshoot-gateway.md +++ b/docs/how-to/troubleshoot-gateway.md @@ -159,6 +159,33 @@ failed to bind the socket and continued without it (the warning `Failed to bind gateway control socket — continuing without it` is in the journal in that case). +### Session pinning is off + +The gateway keeps a mapping while conntrack shows sessions to its +virtual IP. At startup it reads conntrack once, the same way each +10 s tick does, and logs which source answered: + +- `Conntrack source: proc; session pinning is on`: sessions are read + from `/proc/net/nf_conntrack`. +- `Conntrack source: netlink; session pinning is on`: the proc file + is absent, and sessions are read by dumping the conntrack table over + netlink, as `conntrack -L` does. The dump needs `CAP_NET_ADMIN`. +- `No conntrack source is readable; session pinning is off`: no + source could be read, and the line carries the error from each + source it tried. Every mapping then reads zero sessions, so a + mapping is reclaimed on its TTL and grace period alone, even while a + client that has not re-queried DNS still has traffic flowing through + it. + +The proc file exists only on a kernel built with +`CONFIG_NF_CONNTRACK_PROCFS`, and only once `nf_conntrack` is loaded. +Without it the gateway uses the netlink dump: + +```sh +ls /proc/net/nf_conntrack +grep NF_CONNTRACK_PROCFS /boot/config-$(uname -r) +``` + ## Outbound-half diagnostics Symptoms in this section all involve a LAN client trying to reach a diff --git a/packaging/debian/build-deb-container.sh b/packaging/debian/build-deb-container.sh index 1b147a54..ede5a92d 100755 --- a/packaging/debian/build-deb-container.sh +++ b/packaging/debian/build-deb-container.sh @@ -101,12 +101,22 @@ if [ -n "$FEATURES" ]; then # from the default build of the same commit. It refuses --features with # --no-build for that reason, so the two cases cannot share one command. # The version still comes from the host, because the image has no git. - BUILD_CMD="packaging/debian/build-deb.sh --features '$FEATURES' --version '$VERSION' --output-dir /out" + BUILD_CMD="packaging/debian/build-deb.sh --features '$FEATURES' --version '$VERSION' --output-dir /out --name-file /name/deb" else BUILD_CMD="cargo build --release --locked - packaging/debian/build-deb.sh --no-build --version '$VERSION' --output-dir /out" + packaging/debian/build-deb.sh --no-build --version '$VERSION' --output-dir /out --name-file /name/deb" fi +# The build names the package it produced rather than this script picking one +# out of the output directory. The output directory is the caller's and may +# already hold packages from earlier runs; a search there by name or by age +# could return one of those, and a package that sorts higher by name was +# returned in preference to the one just built. The name travels through a +# directory of its own, created fresh for this run, so a name left by an +# earlier run cannot be read and nothing extra is left in the output directory. +NAME_DIR=$(mktemp -d) +trap 'rm -rf "$NAME_DIR"' EXIT + # The source is mounted read-only so a build cannot leave artifacts in the tree. # CARGO_TARGET_DIR and the registry live in named volumes, which is what makes a # second run fast; they are per-base-image so a floor change does not reuse @@ -115,6 +125,7 @@ VOL_SUFFIX="${FIPS_BUILD_IMAGE//[:\/]/-}" docker run --rm \ -v "$REPO_ROOT":/src:ro \ -v "$DEST_ABS":/out \ + -v "$NAME_DIR":/name \ -v "fips-deb-target-${VOL_SUFFIX}":/target \ -v "fips-deb-registry-${VOL_SUFFIX}":/usr/local/cargo/registry \ -e CARGO_TARGET_DIR=/target \ @@ -123,8 +134,21 @@ docker run --rm \ "$IMAGE_TAG" \ bash -euo pipefail -c "$BUILD_CMD" >&2 -DEB=$(find "$DEST_ABS" -maxdepth 1 -name "fips_*_*.deb" -newermt '-10 minutes' -print | sort | tail -1) -[ -n "$DEB" ] || { echo "build-deb-container: no .deb was produced." >&2; exit 1; } +DEB_NAME="" +[ -f "$NAME_DIR/deb" ] && DEB_NAME=$(head -n 1 "$NAME_DIR/deb") +[ -n "$DEB_NAME" ] || { + echo "build-deb-container: the build did not name its package" >&2 + exit 1 +} +if [[ "$DEB_NAME" == */* || "$DEB_NAME" != fips_*_*.deb ]]; then + echo "build-deb-container: the build named '$DEB_NAME', which is not a package file name" >&2 + exit 1 +fi +DEB="$DEST_ABS/$DEB_NAME" +[ -f "$DEB" ] || { + echo "build-deb-container: the build named $DEB_NAME but $DEB does not exist" >&2 + exit 1 +} # Check the artifact here rather than in one workflow, so every producer is # gated: the release, the CI job, a local run and packaging/Makefile all reach diff --git a/packaging/debian/build-deb.sh b/packaging/debian/build-deb.sh index f1088f92..5712cbe1 100755 --- a/packaging/debian/build-deb.sh +++ b/packaging/debian/build-deb.sh @@ -2,7 +2,8 @@ # Build a .deb package for FIPS using cargo-deb. # # Usage: ./build-deb.sh [--target ] [--version ] [--no-build] -# [--features ] +# [--features ] [--output-dir ] +# [--name-file ] # # Prerequisites: cargo-deb (install with: cargo install cargo-deb) # Output: deploy/fips__.deb @@ -26,6 +27,10 @@ Options: --output-dir Where to put the finished .deb. Defaults to deploy/ under the project root. Exists so the container build can write to a mount and leave the source tree read-only. + --name-file Also write the finished package's file name (basename + only) to . The container build reads it so it + never has to guess which .deb in the output directory + this run produced. -h, --help Show this help EOF } @@ -35,6 +40,7 @@ VERSION_OVERRIDE="" NO_BUILD=0 FEATURES="" DEST_DIR="" +NAME_FILE="" while [[ $# -gt 0 ]]; do case "$1" in @@ -58,6 +64,10 @@ while [[ $# -gt 0 ]]; do DEST_DIR="${2:?missing value for --output-dir}" shift 2 ;; + --name-file) + NAME_FILE="${2:?missing value for --name-file}" + shift 2 + ;; -h|--help) usage exit 0 @@ -182,6 +192,9 @@ fi cp "${DEB_FILE}" "${DEST_DIR}/" BASENAME=$(basename "${DEB_FILE}") +if [[ -n "${NAME_FILE}" ]]; then + printf '%s\n' "${BASENAME}" > "${NAME_FILE}" +fi echo "Package built: ${DEST_DIR}/${BASENAME}" echo "" echo "Install with: sudo dpkg -i ${DEST_DIR}/${BASENAME}" diff --git a/packaging/debian/fips-firewall.service b/packaging/debian/fips-firewall.service index 7720d9bc..0c191694 100644 --- a/packaging/debian/fips-firewall.service +++ b/packaging/debian/fips-firewall.service @@ -9,6 +9,9 @@ ConditionPathExists=/etc/fips/fips.nft Type=oneshot RemainAfterExit=yes ExecStart=/usr/sbin/nft -f /etc/fips/fips.nft +# Reapplies the ruleset in place: the file adds then flushes the table, so one +# nft run replaces it in a single transaction, with no moment without it. +ExecReload=/usr/sbin/nft -f /etc/fips/fips.nft ExecStop=-/usr/sbin/nft delete table inet fips StandardOutput=journal StandardError=journal diff --git a/packaging/debian/postinst b/packaging/debian/postinst index 5c4bf4c8..b83869d8 100755 --- a/packaging/debian/postinst +++ b/packaging/debian/postinst @@ -2,6 +2,87 @@ # FIPS post-install script for Debian/Ubuntu set -e +# How long the upgrade path waits for each unit it starts. A bound rather than +# a blocking `systemctl start`: a unit that Requires= a daemon which never comes +# up has a start job that is never dispatched, and a blocking start on it never +# returns, which held apt, and every package operation queued behind it, for +# ever. +UNIT_START_LIMIT=60 +# fips-gateway.service waits up to 30s in ExecStartPre for fips0, behind a +# daemon that has only just been started, so it gets longer. +GATEWAY_START_LIMIT=90 + +# Set when the daemon or fips-dns does not come up on upgrade; checked at the +# end. +start_failed="" + +# Queue a start (or restart) of a unit and wait, up to a bound, for it to become +# active. Returns: +# 0 the unit is active and no job for it is still queued; +# 2 the unit was left inactive on purpose, because it is masked or because +# a Condition in it is not met, which is a skip, not a failure; +# 1 the job could not be queued, the unit failed, or it did not become +# active in time, after printing the unit's status. +# +# Active has to hold on two consecutive polls with no job pending: fips.service +# is Type=simple, so it reads active for an instant after the fork even when +# the exec then fails, and a unit being restarted reads active on its old +# process until the queued job runs. +unit_bounded() { + verb="$1" + unit="$2" + limit="$3" + + case "$(systemctl is-enabled "$unit" 2>/dev/null || true)" in + masked | masked-runtime) + echo "fips: $unit is masked; not starting it" + return 2 + ;; + esac + + # The condition result is only evidence about this start once the unit has + # evaluated its conditions again, so remember when it last did. + cond_before=$(systemctl show -p ConditionTimestampMonotonic --value "$unit" 2>/dev/null || true) + + if ! systemctl "$verb" --no-block "$unit"; then + echo "fips: could not queue $verb of $unit" >&2 + return 1 + fi + + seen=0 + waited=0 + while [ "$waited" -lt "$limit" ]; do + sleep 1 + waited=$((waited + 1)) + if systemctl is-active --quiet "$unit" && + [ -z "$(systemctl show -p Job --value "$unit" 2>/dev/null)" ]; then + seen=$((seen + 1)) + if [ "$seen" -ge 2 ]; then + return 0 + fi + continue + fi + seen=0 + job=$(systemctl show -p Job --value "$unit" 2>/dev/null || true) + state=$(systemctl show -p ActiveState --value "$unit" 2>/dev/null || true) + if [ -z "$job" ] && [ "$state" = "failed" ]; then + echo "fips: $unit failed to start" >&2 + systemctl status --no-pager --lines=15 "$unit" >&2 || true + return 1 + fi + cond_now=$(systemctl show -p ConditionTimestampMonotonic --value "$unit" 2>/dev/null || true) + if [ "$cond_now" != "$cond_before" ] && + [ "$(systemctl show -p ConditionResult --value "$unit" 2>/dev/null)" = "no" ]; then + echo "fips: $unit was skipped because a condition in the unit is not met" + return 2 + fi + done + + echo "fips: $unit did not become active within ${limit}s" >&2 + systemctl status --no-pager --lines=15 "$unit" >&2 || true + return 1 +} + case "$1" in configure) # Create fips system group for control socket access @@ -41,11 +122,55 @@ case "$1" in systemctl enable fips.service 2>/dev/null || true systemctl enable fips-dns.service 2>/dev/null || true - # On upgrade, restart services that were running before + # On upgrade, restart services that were running before. Each + # start is bounded, and a daemon or fips-dns that does not come up + # fails the install with its status printed, rather than holding + # apt. When the daemon does not come up the units that require it + # are not started: each would only wait out its own bound behind + # it. A daemon that was skipped (masked, or its condition not met) + # is not a failure, but the units that require it are not started + # either. if [ -n "$2" ]; then - systemctl start fips.service 2>/dev/null || true - if systemctl is-enabled --quiet fips-dns.service 2>/dev/null; then - systemctl start fips-dns.service 2>/dev/null || true + # Reapply the firewall ruleset in place, before the daemon + # starts, and only where the operator has it running: "try" + # leaves a unit that is not active alone, so this never opts + # a host in. The daemon-reload above has loaded the unit's + # ExecReload, so this reloads rather than restarts; a restart + # would run ExecStop, which deletes the table. A reload that + # fails leaves the previous ruleset loaded, so it is reported + # and the upgrade goes on. + if ! systemctl try-reload-or-restart fips-firewall.service; then + echo "fips: reloading fips-firewall.service failed; the ruleset loaded before the upgrade stays in force" >&2 + echo "fips: check /etc/fips/fips.nft and the rules in /etc/fips/fips.d/" >&2 + fi + + daemon_rc=0 + unit_bounded start fips.service "$UNIT_START_LIMIT" || daemon_rc=$? + if [ "$daemon_rc" -eq 1 ]; then + start_failed=1 + elif [ "$daemon_rc" -eq 2 ]; then + echo "fips: fips.service is not running, so the units that require it were not started" + elif [ "$daemon_rc" -eq 0 ] && + systemctl is-enabled --quiet fips-dns.service 2>/dev/null; then + dns_rc=0 + unit_bounded start fips-dns.service "$UNIT_START_LIMIT" || dns_rc=$? + [ "$dns_rc" -ne 1 ] || start_failed=1 + fi + + # prerm stopped the gateway; bring it back only where the + # operator enabled it, so an upgrade never turns it on. A + # restart rather than a start, so a gateway that an older + # prerm left running also moves onto the new binary. The + # gateway is an opt-in addition to a daemon that is running, + # so one that does not come up is reported, with its status, + # and does not fail the upgrade. + if [ "$daemon_rc" -eq 0 ] && + systemctl is-enabled --quiet fips-gateway.service 2>/dev/null; then + gw_rc=0 + unit_bounded restart fips-gateway.service "$GATEWAY_START_LIMIT" || gw_rc=$? + if [ "$gw_rc" -eq 1 ]; then + echo "fips: fips-gateway.service did not come back after the upgrade; the daemon did" >&2 + fi fi fi fi @@ -54,4 +179,12 @@ esac #DEBHELPER# +# Fail the configure step only here, after everything else has run, so a unit +# that did not come up leaves the package half-configured and apt non-zero. +if [ -n "$start_failed" ]; then + echo "fips: the upgrade is installed but its services did not all start;" >&2 + echo "fips: fix the cause above, then run: dpkg --configure -a" >&2 + exit 1 +fi + exit 0 diff --git a/packaging/debian/prerm b/packaging/debian/prerm index 399f5b6c..75ac046b 100755 --- a/packaging/debian/prerm +++ b/packaging/debian/prerm @@ -5,6 +5,9 @@ set -e case "$1" in remove|purge) if [ -d /run/systemd/system ]; then + # The gateway requires the daemon, so it goes first. + systemctl stop fips-gateway.service 2>/dev/null || true + systemctl disable fips-gateway.service 2>/dev/null || true systemctl stop fips-dns.service 2>/dev/null || true systemctl disable fips-dns.service 2>/dev/null || true systemctl stop fips.service 2>/dev/null || true @@ -17,6 +20,7 @@ case "$1" in upgrade) # Stop services before upgrade; postinst will restart them if [ -d /run/systemd/system ]; then + systemctl stop fips-gateway.service 2>/dev/null || true systemctl stop fips-dns.service 2>/dev/null || true systemctl stop fips.service 2>/dev/null || true fi diff --git a/packaging/systemd/fips-firewall.service b/packaging/systemd/fips-firewall.service index 76aaacd2..e078ee58 100644 --- a/packaging/systemd/fips-firewall.service +++ b/packaging/systemd/fips-firewall.service @@ -9,6 +9,9 @@ ConditionPathExists=/etc/fips/fips.nft Type=oneshot RemainAfterExit=yes ExecStart=/usr/sbin/nft -f /etc/fips/fips.nft +# Reapplies the ruleset in place: the file adds then flushes the table, so one +# nft run replaces it in a single transaction, with no moment without it. +ExecReload=/usr/sbin/nft -f /etc/fips/fips.nft ExecStop=-/usr/sbin/nft delete table inet fips StandardOutput=journal StandardError=journal diff --git a/src/bin/fips-gateway.rs b/src/bin/fips-gateway.rs index 6a6979fb..fb2c6164 100644 --- a/src/bin/fips-gateway.rs +++ b/src/bin/fips-gateway.rs @@ -65,14 +65,14 @@ fn elapsed_us(started: Instant) -> u64 { /// A failed read yields an empty snapshot, so every mapping reads zero /// sessions, which is what the pool did with an unreadable source before. The /// alternative, treating "unknown" as "in use", would pin every mapping forever -/// on a kernel with no conntrack proc file and turn a read error into a pool -/// that never reclaims. The cost is the opposite error: a mapping carrying live -/// traffic can be reclaimed early while the source is unreadable. +/// on a kernel with no readable conntrack source and turn a read error into a +/// pool that never reclaims. The cost is the opposite error: a mapping carrying +/// live traffic can be reclaimed early while the source is unreadable. #[cfg(target_os = "linux")] async fn read_conntrack(log: &mut pool::ConntrackReadLog) -> pool::ConntrackSnapshot { use fips::gateway::pool::ConntrackQuerier; - match tokio::task::spawn_blocking(|| pool::ProcConntrack.snapshot()).await { + match tokio::task::spawn_blocking(|| pool::SystemConntrack::default().snapshot()).await { Ok(Ok(snapshot)) => { log.observe(None); snapshot @@ -107,6 +107,42 @@ fn report_unreadable_conntrack( } } +/// Check once at startup which conntrack source the tick will read, and say so. +/// +/// Without this, an operator on a kernel with no readable source learns that +/// session pinning is off only from a warning at the first failed tick. +#[cfg(target_os = "linux")] +async fn report_conntrack_source() { + let probe = + tokio::task::spawn_blocking(|| pool::probe_conntrack(&pool::SystemConntrack::default())) + .await + .unwrap_or_else(|e| { + pool::ConntrackProbe::Missing(pool::ConntrackUnreadable { + proc: std::io::Error::other(e.to_string()), + netlink: None, + }) + }); + match probe { + pool::ConntrackProbe::Found(pool::ConntrackSource::Proc) => { + info!("Conntrack source: proc; session pinning is on") + } + pool::ConntrackProbe::Found(pool::ConntrackSource::Netlink) => { + info!("Conntrack source: netlink; session pinning is on") + } + pool::ConntrackProbe::Missing(e) => match e.netlink { + Some(netlink) => warn!( + proc_error = %e.proc, + netlink_error = %netlink, + "No conntrack source is readable; session pinning is off" + ), + None => warn!( + proc_error = %e.proc, + "No conntrack source is readable; session pinning is off" + ), + }, + } +} + #[cfg(target_os = "linux")] #[tokio::main(flavor = "current_thread")] async fn main() { @@ -371,6 +407,10 @@ async fn main() { std::process::exit(1); } + // The NAT table exists by now, so a kernel that provides the proc file + // has loaded nf_conntrack and the probe sees what the first tick will. + report_conntrack_source().await; + // --- Channels --- // Pool events (new/removed mappings) → NAT + net modules diff --git a/src/gateway/conntrack.rs b/src/gateway/conntrack.rs new file mode 100644 index 00000000..23fca85c --- /dev/null +++ b/src/gateway/conntrack.rs @@ -0,0 +1,330 @@ +//! Conntrack sessions read over netlink. +//! +//! A kernel built without `CONFIG_NF_CONNTRACK_PROCFS` has no +//! `/proc/net/nf_conntrack`, while `conntrack -L` still lists the table: it +//! asks the kernel for a dump over `NETLINK_NETFILTER`. This reader does the +//! same, so such a kernel can still pin mappings that carry traffic. + +use super::pool::{ConntrackQuerier, ConntrackSnapshot}; +use netlink_packet_core::{ + NLM_F_DUMP, NLM_F_REQUEST, NetlinkHeader, NetlinkMessage, NetlinkPayload, +}; +use netlink_packet_netfilter::conntrack::{ConntrackAttribute, ConntrackMessage, IPTuple, Tuple}; +use netlink_packet_netfilter::{ + NetfilterHeader, NetfilterMessage, NetfilterMessageInner, NetfilterProtoFamily, +}; +use netlink_sys::{Socket, SocketAddr, protocols::NETLINK_NETFILTER}; +use std::collections::{HashMap, HashSet}; +use std::io; +use std::net::{IpAddr, Ipv6Addr}; +use std::sync::atomic::{AtomicU32, Ordering}; +use std::time::Duration; + +/// Longest wait for each part of the kernel's reply. +/// +/// The read runs on a blocking thread once per tick, so a kernel that never +/// answers must not hold that thread for longer than a tick. +const READ_TIMEOUT: Duration = Duration::from_secs(2); + +/// Sequence number for the next dump request, so a reply to an earlier +/// request cannot be counted as part of this one. +static NEXT_SEQ: AtomicU32 = AtomicU32::new(1); + +/// Whether a dump has more to come after the buffer just counted. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum DumpState { + /// The kernel has more of the dump to send. + More, + /// The kernel sent the end-of-dump message. + Done, +} + +/// Conntrack querier that dumps the table over netlink. +/// +/// Needs `CAP_NET_ADMIN` in the gateway's network namespace, which the +/// gateway already needs for its NAT table. +pub struct NetlinkConntrack; + +impl ConntrackQuerier for NetlinkConntrack { + fn snapshot(&self) -> Result { + let mut socket = Socket::new(NETLINK_NETFILTER)?; + socket.bind_auto()?; + socket.connect(&SocketAddr::new(0, 0))?; + socket2::SockRef::from(&socket).set_read_timeout(Some(READ_TIMEOUT))?; + + let seq = NEXT_SEQ.fetch_add(1, Ordering::Relaxed); + socket.send(&dump_request(seq), 0)?; + + let mut counts = HashMap::new(); + loop { + // Sized by peeking first, so a large batch is not truncated. + let (buf, _) = socket.recv_from_full()?; + if count_dump(&buf, seq, &mut counts)? == DumpState::Done { + break; + } + } + Ok(ConntrackSnapshot::from_counts(counts)) + } +} + +/// A request for a dump of the IPv6 conntrack table. +/// +/// The family in the netfilter header makes the kernel leave out IPv4 +/// entries, which could never name a virtual IP. +fn dump_request(seq: u32) -> Vec { + let mut header = NetlinkHeader::default(); + header.flags = NLM_F_REQUEST | NLM_F_DUMP; + header.sequence_number = seq; + let mut message = NetlinkMessage::new( + header, + NetlinkPayload::from(NetfilterMessage::new( + NetfilterHeader::new(NetfilterProtoFamily::IPv6, 0, 0), + ConntrackMessage::Get(vec![]), + )), + ); + message.finalize(); + let mut buf = vec![0; message.buffer_len()]; + message.serialize(&mut buf); + buf +} + +/// Count the conntrack entries in one received buffer by destination. +/// +/// An entry counts once for each distinct IPv6 destination among its original +/// and reply tuples, the same rule the proc-file parser applies to a line. +/// A message carrying another sequence number is skipped. A dump the kernel +/// flags as interrupted is counted as received: reading the proc file is not +/// atomic across the table either, and failing the read would zero every +/// mapping for the tick. +pub fn count_dump( + buf: &[u8], + seq: u32, + counts: &mut HashMap, +) -> Result { + let mut offset = 0; + while offset < buf.len() { + let message = NetlinkMessage::::deserialize(&buf[offset..]) + .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e.to_string()))?; + // Messages are padded to four bytes. The length is at least a header, + // or the parse above would have failed, so the walk always advances. + let len = message.header.length as usize; + offset += (len + 3) & !3; + if message.header.sequence_number != seq { + continue; + } + match message.payload { + NetlinkPayload::Done(_) => return Ok(DumpState::Done), + NetlinkPayload::Error(e) if e.code.is_some() => return Err(e.to_io()), + NetlinkPayload::InnerMessage(NetfilterMessage { + inner: NetfilterMessageInner::Conntrack(ConntrackMessage::New(attrs)), + .. + }) => count_entry(&attrs, counts), + _ => {} + } + } + Ok(DumpState::More) +} + +/// Add one conntrack entry to the counts, once per distinct IPv6 destination +/// among its original and reply tuples. +fn count_entry(attrs: &[ConntrackAttribute], counts: &mut HashMap) { + let mut seen = HashSet::new(); + for attr in attrs { + let tuples = match attr { + ConntrackAttribute::CtaTupleOrig(t) | ConntrackAttribute::CtaTupleReply(t) => t, + _ => continue, + }; + for tuple in tuples { + let Tuple::Ip(ip) = tuple else { continue }; + for field in ip { + if let IPTuple::DestinationAddress(IpAddr::V6(dst)) = field { + seen.insert(*dst); + } + } + } + } + for dst in seen { + *counts.entry(dst).or_insert(0) += 1; + } +} + +#[cfg(test)] +mod tests { + use super::*; + use netlink_packet_core::{DoneMessage, ErrorMessage}; + use std::num::NonZeroI32; + + const SEQ: u32 = 7; + + fn v6(s: &str) -> Ipv6Addr { + s.parse().unwrap() + } + + /// One tuple naming a source and a destination. + fn tuple(src: Ipv6Addr, dst: Ipv6Addr) -> Vec { + vec![Tuple::Ip(vec![ + IPTuple::SourceAddress(IpAddr::V6(src)), + IPTuple::DestinationAddress(IpAddr::V6(dst)), + ])] + } + + /// A conntrack entry as a dump reply carries it. + fn entry(orig: Vec, reply: Vec) -> NetfilterMessage { + NetfilterMessage::new( + NetfilterHeader::new(NetfilterProtoFamily::IPv6, 0, 0), + ConntrackMessage::New(vec![ + ConntrackAttribute::CtaTupleOrig(orig), + ConntrackAttribute::CtaTupleReply(reply), + ]), + ) + } + + /// Serialise one netlink message with the given sequence number. + fn frame(payload: NetlinkPayload, seq: u32) -> Vec { + let mut header = NetlinkHeader::default(); + header.sequence_number = seq; + header.flags = netlink_packet_core::NLM_F_MULTIPART; + let mut message = NetlinkMessage::new(header, payload); + message.finalize(); + let mut buf = vec![0; message.buffer_len()]; + message.serialize(&mut buf); + buf + } + + fn done() -> NetlinkPayload { + NetlinkPayload::Done(DoneMessage::default()) + } + + /// The flow the gateway sees for a LAN client using a virtual IP: the + /// original tuple is client to virtual IP, and the reply, after DNAT and + /// masquerade, is the mesh address back to the gateway. + fn client_flow(virtual_ip: Ipv6Addr) -> NetfilterMessage { + entry( + tuple(v6("fd02::20"), virtual_ip), + tuple(v6("fd9a::1"), v6("fd9a::2")), + ) + } + + fn count(buf: &[u8]) -> (Result, HashMap) { + let mut counts = HashMap::new(); + let state = count_dump(buf, SEQ, &mut counts); + (state, counts) + } + + #[test] + fn netlink_dump_counts_an_entry_whose_original_destination_is_the_virtual_ip() { + let virtual_ip = v6("fd01::1"); + let buf = frame(NetlinkPayload::from(client_flow(virtual_ip)), SEQ); + + let (state, counts) = count(&buf); + + assert_eq!(state.unwrap(), DumpState::More); + assert_eq!(counts.get(&virtual_ip).copied(), Some(1)); + // The reply tuple's destination is counted too, as the proc parser + // counts every dst= on the line. + assert_eq!(counts.get(&v6("fd9a::2")).copied(), Some(1)); + } + + #[test] + fn netlink_dump_counts_an_entry_once_when_both_tuples_name_the_address() { + let addr = v6("fd01::1"); + let hairpin = entry(tuple(addr, addr), tuple(addr, addr)); + let buf = frame(NetlinkPayload::from(hairpin), SEQ); + + let (_, counts) = count(&buf); + + assert_eq!(counts.get(&addr).copied(), Some(1)); + } + + #[test] + fn netlink_dump_counts_each_entry_across_several_messages_in_one_buffer() { + let virtual_ip = v6("fd01::1"); + let other = v6("fd01::2"); + let mut buf = frame(NetlinkPayload::from(client_flow(virtual_ip)), SEQ); + buf.extend(frame(NetlinkPayload::from(client_flow(virtual_ip)), SEQ)); + buf.extend(frame(NetlinkPayload::from(client_flow(other)), SEQ)); + + let (state, counts) = count(&buf); + + assert_eq!(state.unwrap(), DumpState::More); + assert_eq!(counts.get(&virtual_ip).copied(), Some(2)); + assert_eq!(counts.get(&other).copied(), Some(1)); + } + + #[test] + fn netlink_dump_ignores_an_ipv4_entry() { + let v4 = |s: &str| IpAddr::V4(s.parse().unwrap()); + let ipv4 = NetfilterMessage::new( + NetfilterHeader::new(NetfilterProtoFamily::IPv4, 0, 0), + ConntrackMessage::New(vec![ConntrackAttribute::CtaTupleOrig(vec![Tuple::Ip( + vec![ + IPTuple::SourceAddress(v4("192.0.2.1")), + IPTuple::DestinationAddress(v4("192.0.2.2")), + ], + )])]), + ); + let mut buf = frame(NetlinkPayload::from(ipv4), SEQ); + buf.extend(frame(NetlinkPayload::from(client_flow(v6("fd01::1"))), SEQ)); + + let (state, counts) = count(&buf); + + assert_eq!(state.unwrap(), DumpState::More); + assert_eq!(counts.len(), 2, "only the IPv6 entry's two destinations"); + } + + #[test] + fn netlink_dump_reports_done_on_the_done_message() { + let virtual_ip = v6("fd01::1"); + let mut buf = frame(NetlinkPayload::from(client_flow(virtual_ip)), SEQ); + buf.extend(frame(done(), SEQ)); + + let (state, counts) = count(&buf); + + assert_eq!(state.unwrap(), DumpState::Done); + assert_eq!(counts.get(&virtual_ip).copied(), Some(1)); + } + + #[test] + fn netlink_dump_turns_an_eperm_error_message_into_permission_denied() { + let mut error = ErrorMessage::default(); + error.code = NonZeroI32::new(-libc::EPERM); + let buf = frame(NetlinkPayload::Error(error), SEQ); + + let (state, _) = count(&buf); + + assert_eq!( + state.expect_err("an error reply must fail the read").kind(), + io::ErrorKind::PermissionDenied + ); + } + + #[test] + fn netlink_dump_skips_a_message_with_another_sequence_number() { + let virtual_ip = v6("fd01::1"); + let mut buf = frame(NetlinkPayload::from(client_flow(virtual_ip)), SEQ + 1); + buf.extend(frame(done(), SEQ + 1)); + buf.extend(frame(NetlinkPayload::from(client_flow(virtual_ip)), SEQ)); + + let (state, counts) = count(&buf); + + assert_eq!( + state.unwrap(), + DumpState::More, + "another request's end of dump does not end this one" + ); + assert_eq!(counts.get(&virtual_ip).copied(), Some(1)); + } + + #[test] + fn netlink_dump_rejects_a_buffer_that_does_not_parse() { + let mut buf = frame(NetlinkPayload::from(client_flow(v6("fd01::1"))), SEQ); + buf.truncate(buf.len() - 4); + + let (state, _) = count(&buf); + + assert_eq!( + state.expect_err("a short buffer must fail the read").kind(), + io::ErrorKind::InvalidData + ); + } +} diff --git a/src/gateway/mod.rs b/src/gateway/mod.rs index db3f2dec..df3e4ed7 100644 --- a/src/gateway/mod.rs +++ b/src/gateway/mod.rs @@ -3,6 +3,7 @@ //! Allows unmodified LAN hosts to reach FIPS mesh destinations via //! DNS-allocated virtual IPs and kernel nftables NAT. +pub mod conntrack; pub mod control; pub mod dns; pub mod nat; diff --git a/src/gateway/pool.rs b/src/gateway/pool.rs index 7776b160..bac11373 100644 --- a/src/gateway/pool.rs +++ b/src/gateway/pool.rs @@ -109,7 +109,10 @@ pub struct MappingInfo { pub last_ref_secs: u64, } -/// Path the conntrack table is read from. +/// Path the conntrack table is read from when the kernel provides it. +/// +/// A kernel built without `CONFIG_NF_CONNTRACK_PROCFS` has no such file; +/// `SystemConntrack` then dumps the table over netlink instead. const CONNTRACK_PROC_PATH: &str = "/proc/net/nf_conntrack"; /// Active conntrack sessions counted by destination address. @@ -159,6 +162,132 @@ impl ConntrackQuerier for ProcConntrack { } } +/// Where a conntrack snapshot was read from. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ConntrackSource { + /// `/proc/net/nf_conntrack`. + Proc, + /// A conntrack table dump over `NETLINK_NETFILTER`. + Netlink, +} + +impl ConntrackSource { + /// Short name of the source. + pub fn name(self) -> &'static str { + match self { + Self::Proc => "proc", + Self::Netlink => "netlink", + } + } +} + +/// Why no conntrack source could be read. +#[derive(Debug)] +pub struct ConntrackUnreadable { + /// The error reading `/proc/net/nf_conntrack`. + pub proc: std::io::Error, + /// The error from the netlink dump, when the proc file was absent and the + /// dump was tried. + pub netlink: Option, +} + +impl ConntrackUnreadable { + /// The error that stands for the whole failed read. + /// + /// When the dump was tried, its error is the one that decided the read, so + /// it sets the kind; the absent proc file is kept in the message. Only a + /// proc error that stopped the read before the dump stands alone. + fn into_error(self) -> std::io::Error { + match self.netlink { + Some(netlink) => std::io::Error::new( + netlink.kind(), + format!("proc: {}; netlink: {netlink}", self.proc), + ), + None => self.proc, + } + } +} + +/// The conntrack reader the gateway uses, which also says which source +/// answered. +/// +/// The per-tick read and the startup probe both go through this type, so the +/// probe cannot report a source the tick would not use. The queriers are type +/// parameters so tests can substitute fakes. +/// +/// The proc file is read first. Only when it is absent is the table dumped +/// over netlink, and that is decided on every read: the file appears once +/// `nf_conntrack` is loaded in the namespace, so a choice fixed at startup +/// could keep using netlink on a kernel that has the file. +pub struct SystemConntrack

{ + proc: P, + netlink: N, +} + +impl SystemConntrack { + /// A reader over the given proc and netlink queriers. + pub fn new(proc: P, netlink: N) -> Self { + Self { proc, netlink } + } + + /// Read conntrack once and say which source the snapshot came from. + /// + /// A proc error other than an absent file, such as a permission error, is + /// returned without trying netlink. + pub fn read(&self) -> Result<(ConntrackSource, ConntrackSnapshot), ConntrackUnreadable> { + match self.proc.snapshot() { + Ok(snapshot) => Ok((ConntrackSource::Proc, snapshot)), + Err(proc) if proc.kind() == std::io::ErrorKind::NotFound => { + match self.netlink.snapshot() { + Ok(snapshot) => Ok((ConntrackSource::Netlink, snapshot)), + Err(netlink) => Err(ConntrackUnreadable { + proc, + netlink: Some(netlink), + }), + } + } + Err(proc) => Err(ConntrackUnreadable { + proc, + netlink: None, + }), + } + } +} + +impl Default for SystemConntrack { + fn default() -> Self { + Self::new(ProcConntrack, super::conntrack::NetlinkConntrack) + } +} + +impl ConntrackQuerier for SystemConntrack { + fn snapshot(&self) -> Result { + self.read() + .map(|(_, snapshot)| snapshot) + .map_err(ConntrackUnreadable::into_error) + } +} + +/// Outcome of the startup check for a readable conntrack source. +#[derive(Debug)] +pub enum ConntrackProbe { + /// Sessions can be read, from this source. + Found(ConntrackSource), + /// No source can be read, so every mapping reads zero sessions and session + /// pinning is off. + Missing(ConntrackUnreadable), +} + +/// Read conntrack once, as a tick would, and report which source answered. +pub fn probe_conntrack( + reader: &SystemConntrack, +) -> ConntrackProbe { + match reader.read() { + Ok((source, _)) => ConntrackProbe::Found(source), + Err(e) => ConntrackProbe::Missing(e), + } +} + /// Count conntrack lines by the destination addresses they name. /// /// Every `dst=` value is parsed as an address and compared as an address. The @@ -207,12 +336,12 @@ pub enum ReadReport { /// Remembers the last conntrack read outcome. /// -/// A kernel built without `CONFIG_NF_CONNTRACK_PROCFS` has no -/// `/proc/net/nf_conntrack` at all, so every read fails the same way and a -/// per-tick warning would repeat for the life of the process. Warning on a -/// change of outcome still separates "the source is unreadable" from "there -/// are no sessions", which the pool could not distinguish before, without -/// filling the log. +/// When no source is readable, for example a kernel with no +/// `/proc/net/nf_conntrack` whose netlink dump is refused, every read fails +/// the same way and a per-tick warning would repeat for the life of the +/// process. Warning on a change of outcome still separates "the source is +/// unreadable" from "there are no sessions", which the pool could not +/// distinguish before, without filling the log. #[derive(Debug, Default)] pub struct ConntrackReadLog { last: Option>, @@ -1135,4 +1264,130 @@ mod tests { "a different failure is a different outcome and is worth a line" ); } + + /// A conntrack querier that succeeds with an empty snapshot, or fails with + /// a fixed error kind. + struct FixedRead(Option); + + impl ConntrackQuerier for FixedRead { + fn snapshot(&self) -> Result { + match self.0 { + None => Ok(ConntrackSnapshot::default()), + Some(kind) => Err(kind.into()), + } + } + } + + #[test] + fn conntrack_probe_names_the_proc_source_when_the_proc_read_succeeds() { + let reader = SystemConntrack::new(FixedRead(None), NOT_CALLED); + + match probe_conntrack(&reader) { + ConntrackProbe::Found(source) => { + assert_eq!(source, ConntrackSource::Proc); + assert_eq!(source.name(), "proc"); + } + ConntrackProbe::Missing(e) => panic!("expected the proc source, got {e:?}"), + } + } + + #[test] + fn conntrack_probe_reports_missing_with_the_error_when_the_proc_read_fails() { + let reader = SystemConntrack::new( + FixedRead(Some(std::io::ErrorKind::PermissionDenied)), + NOT_CALLED, + ); + + match probe_conntrack(&reader) { + ConntrackProbe::Missing(e) => { + assert_eq!(e.proc.kind(), std::io::ErrorKind::PermissionDenied); + } + ConntrackProbe::Found(source) => panic!("expected no source, got {source:?}"), + } + } + + /// A netlink stand-in for tests where the dump must not be reached. It + /// fails with a kind no test expects, so reaching it shows in the result. + const NOT_CALLED: FixedRead = FixedRead(Some(std::io::ErrorKind::Unsupported)); + + /// A conntrack querier that reports one session to a fixed address. + struct OneSession(Ipv6Addr); + + impl ConntrackQuerier for OneSession { + fn snapshot(&self) -> Result { + Ok(ConntrackSnapshot::from_counts(HashMap::from([(self.0, 1)]))) + } + } + + #[test] + fn system_conntrack_falls_back_to_netlink_when_the_proc_file_is_absent() { + let addr: Ipv6Addr = "fd01::1".parse().unwrap(); + let reader = SystemConntrack::new( + FixedRead(Some(std::io::ErrorKind::NotFound)), + OneSession(addr), + ); + + let (source, snapshot) = reader.read().expect("the netlink dump answered"); + + assert_eq!(source, ConntrackSource::Netlink); + assert_eq!(source.name(), "netlink"); + assert_eq!(snapshot.sessions_for(addr), 1); + } + + #[test] + fn system_conntrack_does_not_fall_back_on_a_proc_error_other_than_not_found() { + let addr: Ipv6Addr = "fd01::1".parse().unwrap(); + let reader = SystemConntrack::new( + FixedRead(Some(std::io::ErrorKind::PermissionDenied)), + OneSession(addr), + ); + + let e = reader + .read() + .expect_err("a denied proc read is not a missing file"); + + assert_eq!(e.proc.kind(), std::io::ErrorKind::PermissionDenied); + assert!(e.netlink.is_none(), "netlink was not tried"); + assert_eq!( + reader.snapshot().unwrap_err().kind(), + std::io::ErrorKind::PermissionDenied + ); + } + + #[test] + fn system_conntrack_prefers_proc_when_it_reads() { + let proc_addr: Ipv6Addr = "fd01::1".parse().unwrap(); + let netlink_addr: Ipv6Addr = "fd01::2".parse().unwrap(); + let reader = SystemConntrack::new(OneSession(proc_addr), OneSession(netlink_addr)); + + let (source, snapshot) = reader.read().expect("the proc file answered"); + + assert_eq!(source, ConntrackSource::Proc); + assert_eq!(snapshot.sessions_for(proc_addr), 1); + assert_eq!(snapshot.sessions_for(netlink_addr), 0); + } + + #[test] + fn conntrack_probe_reports_both_errors_when_neither_source_reads() { + let reader = SystemConntrack::new( + FixedRead(Some(std::io::ErrorKind::NotFound)), + FixedRead(Some(std::io::ErrorKind::PermissionDenied)), + ); + + match probe_conntrack(&reader) { + ConntrackProbe::Missing(e) => { + assert_eq!(e.proc.kind(), std::io::ErrorKind::NotFound); + assert_eq!( + e.netlink.as_ref().map(std::io::Error::kind), + Some(std::io::ErrorKind::PermissionDenied) + ); + } + ConntrackProbe::Found(source) => panic!("expected no source, got {source:?}"), + } + // The per-tick read reports the error that decided it: the dump's. + assert_eq!( + reader.snapshot().unwrap_err().kind(), + std::io::ErrorKind::PermissionDenied + ); + } } diff --git a/src/node/handlers/rekey.rs b/src/node/handlers/rekey.rs index 96a970a1..bdb09a99 100644 --- a/src/node/handlers/rekey.rs +++ b/src/node/handlers/rekey.rs @@ -898,9 +898,19 @@ impl Node { } }; - // Build SessionSetup with coordinates + // Build SessionSetup with coordinates. The wire needs non-empty + // destination coordinates, and on a cache miss our own would be filed + // under the destination's address by every node on the path. A miss + // here means a direct peer, since we got past `find_next_hop`, and + // routing never refreshes a direct peer's cache entry, so fall back + // to the coordinates the peer announced to us. Only a peer that has + // not announced yet still gets ours, and that frame goes one hop, to + // the destination itself. let our_coords = self.tree_state.my_coords().clone(); - let dest_coords = self.get_dest_coords(dest_addr); + let dest_coords = self + .cached_dest_coords(dest_addr) + .or_else(|| self.tree_state.peer_coords(dest_addr).cloned()) + .unwrap_or_else(|| our_coords.clone()); let setup = SessionSetup::new(our_coords, dest_coords).with_handshake(msg1); let setup_payload = setup.encode(); diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index ccc3f7c6..1acaf018 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -915,7 +915,7 @@ impl Node { // never reaches for an entry holding a pending session; and // `set_pending_session` clears `rekey_state`, so a completed // initiator cycle leaves at most one of the two set. If that ever - // stops holding, these four sites become instances of the epoch + // stops holding, these three sites become instances of the epoch // discard the responder arm was fixed for. if entry.is_established() && entry.has_rekey_in_progress() && entry.is_rekey_initiator() { let mut handshake = match entry.take_rekey_state() { @@ -926,14 +926,28 @@ impl Node { } }; - // Process XX msg2 - if let Err(e) = handshake.read_message_2(&ack.handshake_payload) { - debug!(error = %e, "Failed to process rekey XX msg2"); - entry.abandon_rekey(); + // Process XX msg2, for the same reason and in the same way as the + // primary arm below. Nothing here has been authenticated: the + // only tie to our rekey is the datagram's source address, which + // the sender chooses. Abandoning would let anyone able to name + // the session end the cycle, so the handshake goes back, rolled + // back to its pre-read state so it can still read the genuine + // ack, and the refusal is counted. The rollback matters because + // `read_message_2` mixes the sender's ephemeral in before it + // authenticates. + if let Err(e) = handshake.try_read_message_2(&ack.handshake_payload) { + debug!(error = %e, "Failed to process rekey XX msg2, keeping the rekey"); + entry.set_rekey_state(handshake, true); self.sessions.insert(*src_addr, entry); + self.stats_mut() + .record_reject(RejectReason::Session(SessionReject::AckHandshakeFailed)); return; } + // The three abandons below stay abandons. Each is a local + // failure after msg2 has read (writing msg3, sending it, or + // completing the session), not a refusal of the ack. + // Generate XX msg3 let msg3 = match handshake.write_message_3() { Ok(m) => m, @@ -2489,11 +2503,22 @@ impl Node { let inner_plaintext = fsp_prepend_inner_header(timestamp, msg_type, inner_flags, &port_payload); + // With no coordinates cached for the destination, send without CP + // and leave the warmup budget unspent: our own coordinates in the + // destination's slot would be filed under its address by every + // receiver, and the budget is better spent on the first frames after + // discovery refills the cache. + let cached_dst = if wants_coords { + self.cached_dest_coords(dest_addr) + } else { + None + }; + let warming = cached_dst.is_some(); + // Determine whether coords fit within transport MTU. // If not, send standalone CoordsWarmup before the data packet. - let (include_coords, my_coords, dest_coords) = if wants_coords { + let (include_coords, my_coords, dest_coords) = if let Some(dst) = cached_dst { let src = self.tree_state.my_coords().clone(); - let dst = self.get_dest_coords(dest_addr); let coords_size = coords_wire_size(&src) + coords_wire_size(&dst); let total_wire = FIPS_OVERHEAD as usize + FSP_PORT_HEADER_SIZE + coords_size + payload.len(); @@ -2512,7 +2537,7 @@ impl Node { }; // Decrement warmup counter if we sent coords (piggybacked or standalone) - if wants_coords && let Some(entry) = self.sessions.get_mut(dest_addr) { + if warming && let Some(entry) = self.sessions.get_mut(dest_addr) { entry.set_coords_warmup_remaining(entry.coords_warmup_remaining() - 1); } @@ -3005,8 +3030,16 @@ impl Node { ) -> Result<(), NodeError> { let now_ms = Self::now_ms(); + // A warmup's only content is the two coordinates; with none cached + // for the destination, ours would stand in for its own. + let Some(dest_coords) = self.cached_dest_coords(dest_addr) else { + trace!( + dest = %self.peer_display_name(dest_addr), + "No cached coordinates for destination, skipping CoordsWarmup" + ); + return Ok(()); + }; let my_coords = self.tree_state.my_coords().clone(); - let dest_coords = self.get_dest_coords(dest_addr); // Read session metadata let entry = self @@ -3133,6 +3166,19 @@ impl Node { Ok(()) } + /// Look up destination coordinates in the coordinate cache, with no + /// fallback. + /// + /// Use this wherever the coordinates go on the wire as the destination's + /// own: a miss must not be filled with ours, which every receiver would + /// file under the destination's address. + pub(in crate::node) fn cached_dest_coords( + &self, + dest: &NodeAddr, + ) -> Option { + self.coord_cache.get(dest, Self::now_ms()).cloned() + } + /// Look up destination coordinates from available caches. /// /// Returns our own coordinates as a fallback (the SessionSetup will @@ -3142,9 +3188,8 @@ impl Node { &self, dest: &NodeAddr, ) -> crate::proto::stp::TreeCoordinate { - let now_ms = Self::now_ms(); - if let Some(coords) = self.coord_cache.get(dest, now_ms) { - return coords.clone(); + if let Some(coords) = self.cached_dest_coords(dest) { + return coords; } // Fallback: use our own coordinates. The SessionSetup dest_coords // field cannot be empty (wire format requires ≥1 entry). Using our diff --git a/src/node/rate_limit.rs b/src/node/rate_limit.rs index b456c002..726a1d84 100644 --- a/src/node/rate_limit.rs +++ b/src/node/rate_limit.rs @@ -93,11 +93,19 @@ impl TokenBucket { /// * `capacity` - Maximum number of tokens (burst capacity) /// * `refill_rate` - Tokens added per second pub fn with_params(capacity: u32, refill_rate: f64) -> Self { + Self::with_params_at(capacity, refill_rate, Instant::now()) + } + + /// Create a token bucket with custom parameters, full as of `now`. + /// + /// For a caller that keeps its own clock and passes the same clock's + /// readings to [`Self::try_acquire_at`]. + pub fn with_params_at(capacity: u32, refill_rate: f64, now: Instant) -> Self { Self { capacity, tokens: capacity as f64, refill_rate, - last_refill: Instant::now(), + last_refill: now, } } @@ -114,7 +122,20 @@ impl TokenBucket { /// Returns `true` if n tokens were available and consumed, `false` if /// rate limited (insufficient tokens). pub fn try_acquire_n(&mut self, n: u32) -> bool { - self.refill(); + self.try_acquire_n_at(n, Instant::now()) + } + + /// Try to consume one token, refilling as of `now`. + /// + /// `now` must come from the same clock as every earlier reading this + /// bucket was given. + pub fn try_acquire_at(&mut self, now: Instant) -> bool { + self.try_acquire_n_at(1, now) + } + + /// Try to consume n tokens, refilling as of `now`. + fn try_acquire_n_at(&mut self, n: u32, now: Instant) -> bool { + self.refill_at(now); if self.tokens >= n as f64 { self.tokens -= n as f64; @@ -127,14 +148,14 @@ impl TokenBucket { /// Check if tokens are available without consuming them. #[cfg(test)] pub fn available(&mut self) -> bool { - self.refill(); + self.refill_at(Instant::now()); self.tokens >= 1.0 } /// Get the current number of available tokens. #[cfg(test)] pub fn tokens(&mut self) -> f64 { - self.refill(); + self.refill_at(Instant::now()); self.tokens } @@ -144,9 +165,8 @@ impl TokenBucket { self.capacity } - /// Refill tokens based on elapsed time. - fn refill(&mut self) { - let now = Instant::now(); + /// Refill tokens based on the time elapsed up to `now`. + fn refill_at(&mut self, now: Instant) { let elapsed = now.duration_since(self.last_refill); let elapsed_secs = elapsed.as_secs_f64(); @@ -174,7 +194,7 @@ impl TokenBucket { /// estimated time until one token will be available. #[cfg(test)] pub fn time_until_available(&mut self) -> std::time::Duration { - self.refill(); + self.refill_at(Instant::now()); if self.tokens >= 1.0 { std::time::Duration::ZERO @@ -458,6 +478,8 @@ pub struct SessionSetupRateLimiter { stranger: (u32, f64), /// Burst and refill rate for a new link's established bucket. established: (u32, f64), + /// Where the limiter reads the time. `Instant::now` outside tests. + clock: fn() -> Instant, } impl SessionSetupRateLimiter { @@ -469,6 +491,7 @@ impl SessionSetupRateLimiter { buckets: HashMap::new(), stranger, established, + clock: Instant::now, } } @@ -477,22 +500,22 @@ impl SessionSetupRateLimiter { /// Returns `false` when the class's bucket for that link is empty, in /// which case the caller must drop the message before doing any work. pub fn try_admit(&mut self, link_peer: &NodeAddr, class: Msg1Class) -> bool { - let now = Instant::now(); + let now = (self.clock)(); let stranger = self.stranger; let established = self.established; let link = self .buckets .entry(*link_peer) .or_insert_with(|| LinkBuckets { - stranger: TokenBucket::with_params(stranger.0, stranger.1), - established: TokenBucket::with_params(established.0, established.1), + stranger: TokenBucket::with_params_at(stranger.0, stranger.1, now), + established: TokenBucket::with_params_at(established.0, established.1, now), seen: now, }); link.seen = now; let admitted = match class { - Msg1Class::Stranger => link.stranger.try_acquire(), - Msg1Class::EstablishedLink => link.established.try_acquire(), + Msg1Class::Stranger => link.stranger.try_acquire_at(now), + Msg1Class::EstablishedLink => link.established.try_acquire_at(now), }; if admitted { @@ -502,6 +525,13 @@ impl SessionSetupRateLimiter { admitted } + /// Read the time from `clock` instead of `Instant::now`, so a test can + /// drive the refill. + #[cfg(test)] + pub fn set_clock(&mut self, clock: fn() -> Instant) { + self.clock = clock; + } + /// Number of link peers currently holding buckets. #[cfg(test)] pub fn len(&self) -> usize { diff --git a/src/node/tests/session.rs b/src/node/tests/session.rs index b47ef144..d9883b50 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -1433,6 +1433,17 @@ async fn rekey_cutover_preserves_data_plane() { cleanup_nodes(&mut nodes).await; } +/// Deliver queued packets between the nodes until a round moves none, for at +/// most 50 rounds of 10 ms. +async fn pump_until_quiet(nodes: &mut [TestNode]) { + for _ in 0..50 { + tokio::time::sleep(Duration::from_millis(10)).await; + if process_available_packets(nodes).await == 0 { + break; + } + } +} + #[tokio::test] async fn test_tun_outbound_triggers_session_initiation() { // Two connected nodes, no session yet. @@ -4528,40 +4539,39 @@ async fn test_forged_setups_from_one_link_peer_stop_creating_session_entries_onc cleanup_nodes(&mut nodes).await; } +/// The limiter's clock for the refill test: tokio's paused clock, which the +/// test moves with `tokio::time::advance` and nothing else moves. +fn paused_now() -> std::time::Instant { + tokio::time::Instant::now().into_std() +} + #[tokio::test] async fn test_a_drained_setup_bucket_refills_and_admits_the_next_legitimate_setup() { - // The refill has to be slow enough that the draining loop below cannot be - // outrun by the refill it is draining against. At the 50/s this test used - // to run at, a token returned every 20 ms, so on a loaded runner the loop - // outlived its own window, the third setup was admitted, and the - // precondition failed on arrangement rather than on behaviour. At 2/s a - // delivery would have to take 500 ms to lose that race. - let mut nodes = make_setup_limited_pair(2, 2.0).await; + const BURST: u32 = 2; + const RATE: f64 = 2.0; + let mut nodes = make_setup_limited_pair(BURST, RATE).await; + + // From here the limiter reads a clock only the test moves, so the drain + // cannot race a refill however slowly each delivery runs, and the refill + // below is exactly the one the test grants. + tokio::time::pause(); + nodes[1].node.setup_rate_limiter.set_clock(paused_now); - // Deliver until one is actually refused, rather than assuming three is - // enough: a delivery the refill absorbs costs one more iteration and - // nothing else. The cap is what a runner slow enough to lose even this - // race trips, and it says so rather than reporting a drained bucket that - // was never drained. let before = nodes[1].node.stats().session.setup_rate_limited; - let mut delivered = 0; - while nodes[1].node.stats().session.setup_rate_limited == before { - assert!( - delivered < 50, - "the bucket must actually be drained before the refill is tested; \ - 50 forged setups drew no refusal, so each delivery is outlasting \ - the 500 ms refill interval" - ); + for _ in 0..=BURST { deliver_forged_setup_over_link(&mut nodes).await; - delivered += 1; } + assert_eq!( + nodes[1].node.stats().session.setup_rate_limited, + before + 1, + "the burst must be admitted and the one setup past it refused" + ); - // A full burst back from empty at 2/s, so the legitimate setup below meets - // the same bucket however many tokens the drain left behind. The point - // being made is that the denial is transient and clears on its own; the - // length of the window is a function of the configured rate, not of the - // claim. - tokio::time::sleep(Duration::from_millis(1200)).await; + // A full burst back from empty at the configured rate, so the legitimate + // setup below meets a bucket the refill alone has restored. The denial + // is transient and clears on its own; how long it lasts is a function of + // the configured rate. + tokio::time::advance(Duration::from_secs_f64(f64::from(BURST) / RATE)).await; establish_pair_session(&mut nodes).await; cleanup_nodes(&mut nodes).await; @@ -4866,6 +4876,24 @@ fn test_session_entry_size_stays_within_the_budget_the_cap_is_derived_from() { // Integration tests: a forged SessionAck against an in-flight initiation // ============================================================================ +/// A forged SessionAck of exactly the right length, carrying `from`'s tree +/// coordinates. +/// +/// The leading 33 bytes of its handshake payload are a valid compressed +/// point, which is what makes it discriminate a rollback: random bytes +/// usually fail `PublicKey::from_slice` before anything has been mixed into +/// the symmetric state. The bytes after it are zeroed, so the read fails only +/// once the point has been mixed in. +fn forged_session_ack(from: &TestNode) -> Vec { + let mut payload = Identity::generate().pubkey_full().serialize().to_vec(); + payload.resize(crate::noise::HANDSHAKE_MSG2_SIZE, 0); + assert_eq!(payload.len(), crate::noise::HANDSHAKE_MSG2_SIZE); + let coords = from.node.tree_state().my_coords().clone(); + SessionAck::new(coords.clone(), coords) + .with_handshake(payload) + .encode() +} + #[tokio::test] async fn test_forged_session_ack_leaves_the_initiation_able_to_complete_on_the_genuine_ack() { let mut nodes = make_rekey_disabled_pair().await; @@ -4887,17 +4915,7 @@ async fn test_forged_session_ack_leaves_the_initiation_able_to_complete_on_the_g .expect("initiating entry present") .last_activity(); - // A forged ack of exactly the right length. The leading 33 bytes are a - // valid compressed point, which is the point of the test: random bytes - // usually fail `PublicKey::from_slice` before anything has been mixed - // into the symmetric state, so they would not discriminate the rollback. - let mut payload = Identity::generate().pubkey_full().serialize().to_vec(); - payload.resize(crate::noise::HANDSHAKE_MSG2_SIZE, 0); - assert_eq!(payload.len(), crate::noise::HANDSHAKE_MSG2_SIZE); - let coords = nodes[1].node.tree_state().my_coords().clone(); - let forged = SessionAck::new(coords.clone(), coords) - .with_handshake(payload) - .encode(); + let forged = forged_session_ack(&nodes[1]); nodes[0] .node @@ -5168,6 +5186,617 @@ async fn test_a_session_ack_under_the_wrong_static_key_leaves_the_initiation_ali cleanup_nodes(&mut nodes).await; } +/// A SessionAck that fails to read must not end an FSP rekey the node +/// initiated. +/// +/// Nothing authenticates a SessionAck before its msg2 is read: the only tie +/// to the rekey is the datagram's source address, which the sender chooses. +/// So the rekey-initiator arm has to put its handshake back, rolled back to +/// its pre-read state, and let the genuine ack complete the cycle, as the +/// primary arm does for an initiation. +#[tokio::test] +async fn test_forged_session_ack_leaves_the_rekey_able_to_complete_on_the_genuine_ack() { + use crate::proto::fmp::wire::{CommonPrefix, PHASE_ESTABLISHED}; + use crate::transport::ReceivedPacket; + + // node 0 rekeys after one message; node 1 never initiates. + let mut cfg0 = Config::new(); + cfg0.node.rekey.after_messages = 1; + let mut cfg1 = Config::new(); + cfg1.node.rekey.after_messages = u64::MAX; + cfg1.node.rekey.after_secs = u64::MAX; + let mut nodes = run_tree_test_with_configs(vec![cfg0, cfg1], &[(0, 1)]).await; + verify_tree_convergence(&nodes); + populate_all_coord_caches(&mut nodes); + establish_pair_session(&mut nodes).await; + + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + + // One frame crosses node 0's rekey trigger. + nodes[0] + .node + .send_session_data(&node1_addr, 0, 0, b"before the rekey") + .await + .expect("send_session_data failed"); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + + // node 0 sends its rekey SessionSetup; only node 1 is pumped, so node 1 + // arms and its SessionAck waits in node 0's queue. + nodes[0].node.check_session_rekey().await; + assert!( + nodes[0] + .node + .get_session(&node1_addr) + .is_some_and(|e| e.has_rekey_in_progress() && e.is_rekey_initiator()), + "node 0 must have initiated a rekey" + ); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes[1..]).await; + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .is_some_and(|e| e.has_rekey_in_progress() && !e.is_rekey_initiator()), + "node 1 must have armed as the rekey responder" + ); + tokio::time::sleep(Duration::from_millis(20)).await; + let held: Vec = + std::iter::from_fn(|| nodes[0].packet_rx.try_recv().ok()).collect(); + assert!( + !held.is_empty(), + "node 1's SessionAck must be queued at node 0" + ); + for packet in &held { + assert_eq!( + CommonPrefix::parse(&packet.data).map(|p| p.phase), + Some(PHASE_ESTABLISHED), + "every held packet must be a link frame" + ); + } + + // The forgery arrives first, under node 1's address. + let forged = forged_session_ack(&nodes[1]); + nodes[0] + .node + .handle_session_payload(&node1_addr, &node1_addr, &forged, 1280, false) + .await; + let entry = nodes[0] + .node + .get_session(&node1_addr) + .expect("an unreadable ack must not remove the session"); + assert!( + entry.has_rekey_in_progress() && entry.is_rekey_initiator(), + "the rekey must still be in flight after an ack that did not read" + ); + assert_eq!( + nodes[0].node.stats().session.ack_handshake_failed, + 1, + "the refusal must be counted" + ); + + // Release the genuine ack. This is the assertion that tells the outcomes + // apart: an initiator that abandoned on the forgery meets the genuine ack + // with no rekey in flight and completes nothing. + for packet in held { + nodes[0].node.handle_encrypted_frame(packet).await; + } + for _ in 0..3 { + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + } + assert!( + nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .pending_new_session() + .is_some(), + "node 0 must complete the rekey on the genuine ack" + ); + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .pending_new_session() + .is_some(), + "node 1 must hold the new session after msg3" + ); + + // node 0 cuts over on its liveness timer, and data decodes both ways on + // the new epoch. + let now_ms = wall_clock_ms(); + nodes[0] + .node + .sessions + .get_mut(&node1_addr) + .unwrap() + .set_rekey_completed_ms(now_ms - 10_000); + nodes[0].node.check_session_rekey().await; + assert!( + nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .pending_new_session() + .is_none(), + "node 0 must have cut over" + ); + + let recv1_before = nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .traffic_counters() + .1; + nodes[0] + .node + .send_session_data(&node1_addr, 0, 0, b"after the rekey 0 to 1") + .await + .expect("send_session_data failed"); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + let entry1 = nodes[1].node.get_session(&node0_addr).unwrap(); + assert_eq!( + entry1.traffic_counters().1, + recv1_before + 1, + "node 0 to node 1 must decode on the new epoch" + ); + assert!( + entry1.pending_new_session().is_none(), + "node 0's first new-epoch frame must complete node 1's cutover" + ); + + let recv0_before = nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .traffic_counters() + .1; + nodes[1] + .node + .send_session_data(&node0_addr, 0, 0, b"after the rekey 1 to 0") + .await + .expect("send_session_data failed"); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + assert_eq!( + nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .traffic_counters() + .1, + recv0_before + 1, + "node 1 to node 0 must decode on the new epoch" + ); + + cleanup_nodes(&mut nodes).await; +} + +// ============================================================================ +// Integration tests: a destination with no cached coordinates +// ============================================================================ + +/// Build a two-node routable mesh with periodic rekey off and an established +/// FSP session from node 0 to node 1, with node 0's coordinate warmup budget +/// set to `warmup`. +/// +/// node 1's budget is 0, so none of its frames carry coordinates: each one +/// that did would re-warm node 0's entry for node 1, and these tests need +/// that entry to stay gone once they remove it. +async fn make_warmup_pair(warmup: u8) -> Vec { + let configs = (0..2) + .map(|i| { + let mut config = Config::new(); + config.node.rekey.enabled = false; + config.node.session.coords_warmup_packets = if i == 0 { warmup } else { 0 }; + config + }) + .collect(); + let mut nodes = run_tree_test_with_configs(configs, &[(0, 1)]).await; + verify_tree_convergence(&nodes); + populate_all_coord_caches(&mut nodes); + establish_pair_session(&mut nodes).await; + pump_until_quiet(&mut nodes).await; + nodes +} + +/// The address path `node` has cached for `addr`, if any. +/// +/// Compared by address because coordinates carried in a session frame arrive +/// without the declaration metadata a node's own copy holds. +fn cached_path(node: &TestNode, addr: &NodeAddr) -> Option> { + node.node + .coord_cache() + .get(addr, wall_clock_ms()) + .map(addr_path) +} + +/// The address path of a coordinate, self to root. +fn addr_path(coords: &crate::proto::stp::TreeCoordinate) -> Vec { + coords.node_addrs().copied().collect() +} + +/// Deliver a PathBroken naming `dest` to `nodes[at]`, reported by `reporter`. +/// The reporter must not be `dest`, or the signal is refused as forged before +/// it removes anything. +/// +/// Nothing is pumped afterwards: the handler starts a lookup, and in a small +/// mesh its answer would refill the entry the signal removed before the test +/// could send into the miss. +async fn deliver_path_broken( + nodes: &mut [TestNode], + at: usize, + dest: NodeAddr, + reporter: NodeAddr, +) { + use crate::proto::routing::PathBroken; + let encoded = PathBroken::new(dest, reporter).encode(); + nodes[at] + .node + .handle_path_broken(&reporter, &encoded[5..]) + .await; +} + +/// A data frame to a destination whose coordinates this node does not have +/// cached must not carry this node's own coordinates in their place. +/// +/// The shape reached in practice: a direct peer whose cache entry is gone, +/// here removed by a PathBroken from a third address. The destination warms +/// its cache from every coordinate-bearing frame, so a frame carrying the +/// sender's coordinates as the destination's leaves the destination holding +/// its own address under the sender's coordinates. Once the coordinates are +/// known again, the warmup budget the miss did not spend is spent on frames +/// that carry them. +#[tokio::test] +async fn test_a_data_frame_to_a_destination_with_no_cached_coordinates_carries_none_and_keeps_the_warmup() + { + let mut nodes = make_warmup_pair(1).await; + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + let node0_coords = nodes[0].node.tree_state().my_coords().clone(); + let node1_coords = nodes[1].node.tree_state().my_coords().clone(); + + // A transit router reports the path to node 1 broken. It is a third + // address: a report "from" node 1 about node 1 is refused as forged. + let reporter = NodeAddr::from_bytes([0xBB; 16]); + deliver_path_broken(&mut nodes, 0, node1_addr, reporter).await; + assert!( + nodes[0] + .node + .coord_cache() + .get(&node1_addr, wall_clock_ms()) + .is_none(), + "precondition: the PathBroken must have removed node 0's entry for node 1" + ); + let warmup = |nodes: &[TestNode]| { + nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .coords_warmup_remaining() + }; + assert_eq!(warmup(&nodes), 1, "PathBroken resets the warmup budget"); + + let mismatch_before = nodes[1] + .node + .metrics() + .forwarding + .coord_warm_key_mismatch + .get(); + let recv_before = nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .traffic_counters() + .1; + nodes[0] + .node + .send_session_data(&node1_addr, 0, 0, b"after the cache miss") + .await + .expect("send_session_data failed"); + // node 1 takes the frame, and node 0 then drops everything node 1 sent + // back. That includes the answer to the lookup the PathBroken started, + // which would both refill the cache and reset the warmup budget, and so + // hide whether the miss spent it. + process_available_packets(&mut nodes[1..]).await; + while nodes[0].packet_rx.try_recv().is_ok() {} + + assert_eq!( + nodes[1] + .node + .metrics() + .forwarding + .coord_warm_key_mismatch + .get(), + mismatch_before, + "node 1 must not be sent node 0's coordinates as its own" + ); + assert_ne!( + cached_path(&nodes[1], &node1_addr), + Some(addr_path(&node0_coords)), + "node 1 must not hold its own address under node 0's coordinates" + ); + assert_eq!( + nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .traffic_counters() + .1, + recv_before + 1, + "the frame must still be delivered" + ); + assert_eq!( + warmup(&nodes), + 1, + "a frame sent without coordinates must not spend the warmup budget" + ); + + // Node 0's cache is refilled without a discovery answer, so the only + // budget left to spend is the one the miss preserved. node 1's entry for + // node 0 is removed first, so the only way node 1 can learn node 0's + // coordinates again is from a frame that carries them. + let now_ms = wall_clock_ms(); + nodes[0] + .node + .insert_coord_hint(node1_addr, node1_coords, now_ms); + nodes[1].node.coord_cache.remove(&node0_addr); + nodes[0] + .node + .send_session_data(&node1_addr, 0, 0, b"after the refill") + .await + .expect("send_session_data failed"); + pump_until_quiet(&mut nodes).await; + + assert_eq!( + cached_path(&nodes[1], &node0_addr), + Some(addr_path(&node0_coords)), + "the first frame after the refill must carry node 0's coordinates" + ); + assert_eq!( + warmup(&nodes), + 0, + "and it spends the warmup budget the miss preserved" + ); + + cleanup_nodes(&mut nodes).await; +} + +/// A standalone CoordsWarmup to a destination with no cached coordinates +/// sends nothing: its only content would be this node's coordinates standing +/// in for the destination's. +#[tokio::test] +async fn test_a_coords_warmup_to_a_destination_with_no_cached_coordinates_sends_nothing() { + let mut nodes = make_warmup_pair(5).await; + let node1_addr = *nodes[1].node.node_addr(); + + nodes[0].node.coord_cache.remove(&node1_addr); + let warmup_before = nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .coords_warmup_remaining(); + assert_eq!( + nodes[1].packet_rx.len(), + 0, + "precondition: node 1's queue is empty" + ); + + nodes[0] + .node + .send_coords_warmup(&node1_addr) + .await + .expect("a skipped warmup is not an error"); + + assert_eq!( + nodes[1].packet_rx.len(), + 0, + "no CoordsWarmup may be sent without the destination's coordinates" + ); + assert_eq!( + nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .coords_warmup_remaining(), + warmup_before, + "the warmup budget must not move" + ); + + cleanup_nodes(&mut nodes).await; +} + +/// A rekey SessionSetup to a direct peer whose cache entry is gone carries +/// the peer's own announced coordinates, not this node's. +#[tokio::test] +async fn test_a_rekey_setup_to_a_peer_with_no_cached_coordinates_carries_the_peers_announced_ones() +{ + // node 0 rekeys after one message; node 1 never initiates. + let mut cfg0 = Config::new(); + cfg0.node.rekey.after_messages = 1; + let mut cfg1 = Config::new(); + cfg1.node.rekey.after_messages = u64::MAX; + cfg1.node.rekey.after_secs = u64::MAX; + let mut nodes = run_tree_test_with_configs(vec![cfg0, cfg1], &[(0, 1)]).await; + verify_tree_convergence(&nodes); + populate_all_coord_caches(&mut nodes); + establish_pair_session(&mut nodes).await; + + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + let node1_coords = nodes[1].node.tree_state().my_coords().clone(); + assert_eq!( + nodes[0].node.tree_state().peer_coords(&node1_addr), + Some(&node1_coords), + "precondition: node 0 knows node 1's announced coordinates" + ); + + // One frame crosses node 0's rekey trigger, sent while the cache still + // holds node 1, so only the rekey setup can meet the miss. + nodes[0] + .node + .send_session_data(&node1_addr, 0, 0, b"before the rekey") + .await + .expect("send_session_data failed"); + pump_until_quiet(&mut nodes).await; + + // node 1's entry for its own address goes too, so what it holds after + // the rekey can only have come from the setup. + nodes[0].node.coord_cache.remove(&node1_addr); + nodes[1].node.coord_cache.remove(&node1_addr); + let mismatch_before = nodes[1] + .node + .metrics() + .forwarding + .coord_warm_key_mismatch + .get(); + + nodes[0].node.check_session_rekey().await; + for _ in 0..3 { + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + } + + assert_eq!( + nodes[1] + .node + .metrics() + .forwarding + .coord_warm_key_mismatch + .get(), + mismatch_before, + "the rekey setup must not name node 0's coordinates as node 1's" + ); + assert_eq!( + cached_path(&nodes[1], &node1_addr), + Some(addr_path(&node1_coords)), + "the setup's destination coordinates must be node 1's own" + ); + assert!( + nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .pending_new_session() + .is_some(), + "the rekey must complete at node 0" + ); + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .pending_new_session() + .is_some(), + "and at node 1" + ); + + cleanup_nodes(&mut nodes).await; +} + +/// A transit router must not end up holding the destination under the +/// source's coordinates after the source's cache entry for it is gone. +/// +/// Not a red-first test: a source routes to a destination that is not a +/// direct peer only through its coordinate cache, so on a miss there is no +/// next hop and no frame reaches the transit router at all. This constructs +/// that rather than leaving it to a reading of the routing code. +#[tokio::test] +async fn test_a_transit_router_does_not_learn_the_source_coordinates_as_the_destinations() { + // Only A sends coordinates, so nothing but a frame from A can re-warm + // A's entry for B once it is removed. + let configs = (0..3) + .map(|i| { + let mut config = Config::new(); + config.node.rekey.enabled = false; + if i != 0 { + config.node.session.coords_warmup_packets = 0; + } + config + }) + .collect(); + let mut nodes = run_tree_test_with_configs(configs, &[(0, 1), (1, 2)]).await; + verify_tree_convergence(&nodes); + populate_all_coord_caches(&mut nodes); + + let a_addr = *nodes[0].node.node_addr(); + let t_addr = *nodes[1].node.node_addr(); + let b_addr = *nodes[2].node.node_addr(); + let a_coords = nodes[0].node.tree_state().my_coords().clone(); + let b_pubkey = nodes[2].node.identity().pubkey_full(); + + nodes[0] + .node + .initiate_session(b_addr, b_pubkey) + .await + .expect("initiate_session failed"); + pump_until_quiet(&mut nodes).await; + assert!( + nodes[0] + .node + .get_session(&b_addr) + .is_some_and(|e| e.is_established()), + "precondition: A and B must hold an established session" + ); + assert!( + nodes[2] + .node + .get_session(&a_addr) + .is_some_and(|e| e.is_established()), + "precondition: B must hold the session too" + ); + + // T reports the path to B broken, and A's entry for B goes. + deliver_path_broken(&mut nodes, 0, b_addr, t_addr).await; + assert!( + nodes[0] + .node + .coord_cache() + .get(&b_addr, wall_clock_ms()) + .is_none(), + "precondition: A's entry for B must be gone" + ); + + let mismatch_before = nodes[1] + .node + .metrics() + .forwarding + .coord_warm_key_mismatch + .get(); + let sent = nodes[0] + .node + .send_session_data(&b_addr, 0, 0, b"after the cache miss") + .await; + pump_until_quiet(&mut nodes).await; + + assert_ne!( + cached_path(&nodes[1], &b_addr), + Some(addr_path(&a_coords)), + "T must not hold B under A's coordinates" + ); + assert_eq!( + nodes[1] + .node + .metrics() + .forwarding + .coord_warm_key_mismatch + .get(), + mismatch_before, + "T must not be sent a coordinate filed under the wrong address" + ); + assert!( + sent.is_err(), + "with no coordinates cached for B, A has no route to it, which is why \ + no transit router can see the frame" + ); + + cleanup_nodes(&mut nodes).await; +} + // ============================================================================ // Tick-loop maintenance with periodic rekey disabled // ============================================================================ diff --git a/testing/deb-install/test.sh b/testing/deb-install/test.sh index ad299d6a..1cfa5ed3 100755 --- a/testing/deb-install/test.sh +++ b/testing/deb-install/test.sh @@ -58,6 +58,10 @@ SKIP=0 # however many scenarios run in one process. SUPPLIED_DEB="" DEB_PREPARED=0 +# The package the scenarios install, set by build_deb() on every path that +# succeeds. Scenarios use it rather than listing the cache directory, so which +# file they install never depends on what else happens to be in there. +DEB_PATH="" # ───────────────────────────────────────────────────────────────────── # Helpers @@ -202,6 +206,7 @@ build_deb() { fi rm -f "$DEB_CACHE_DIR"/*.deb cp "$SUPPLIED_DEB" "$DEB_CACHE_DIR/" + DEB_PATH="$DEB_CACHE_DIR/$(basename "$SUPPLIED_DEB")" DEB_PREPARED=1 log "Installing the supplied package $(basename "$SUPPLIED_DEB")" return 0 @@ -218,6 +223,7 @@ build_deb() { local cached_age cached_age=$(stat -c '%Y' "$cached_deb" 2>/dev/null || echo 0) if awk "BEGIN { exit !($cached_age >= $newest_src) }"; then + DEB_PATH="$cached_deb" log "Using cached .deb at $cached_deb" return 0 fi @@ -233,20 +239,60 @@ build_deb() { # not exhibit a defect that only the release environment produced. It stayed # green through five releases that could not start on two of the five # distributions in its own matrix. + # + # The cache holds one package at a time. Clearing it first is what keeps the + # reuse check above honest, since that check looks at whichever package it + # finds; and the package installed is the one the build names on the last + # line of its stdout, never one found by listing the directory. log "Building the .deb in the pinned build container (slow on first run)" - if ! bash "$REPO_ROOT/packaging/debian/build-deb-container.sh" \ - --output-dir "$DEB_CACHE_DIR" >&2; then + rm -f "$DEB_CACHE_DIR"/*.deb + local build_out + if ! build_out=$(bash "$REPO_ROOT/packaging/debian/build-deb-container.sh" \ + --output-dir "$DEB_CACHE_DIR"); then echo " ERROR: container build failed" >&2 return 1 fi - cached_deb=$(ls "$DEB_CACHE_DIR"/fips_*_amd64.deb 2>/dev/null | head -1) - if [ -n "$cached_deb" ]; then - log "Cached at $cached_deb ($(stat -c %s "$cached_deb") bytes)" - else - echo " ERROR: no .deb produced by the container build" >&2 + cached_deb=$(printf '%s\n' "$build_out" | tail -n 1) + if [ -z "$cached_deb" ] || [ ! -f "$cached_deb" ]; then + echo " ERROR: the container build did not report a package path: '$cached_deb'" >&2 return 1 fi + DEB_PATH="$cached_deb" + log "Cached at $cached_deb ($(stat -c %s "$cached_deb") bytes)" + return 0 +} + +# The packages a runtime image installs on top of the distro base image. +# Ubuntu 22.04 bundles systemd-resolved into systemd; other distros require it +# as a separate package. +runtime_packages() { + local base_image="$1" + if [ "$base_image" = "ubuntu:22.04" ]; then + echo "systemd iproute2 dbus dnsutils procps" + else + echo "systemd systemd-resolved iproute2 dbus dnsutils procps" + fi + return 0 +} + +# Patch a minimal gateway config into the container's fips.yaml, since the +# shipped one has the gateway disabled, and restart fips.service to load it. +# The caller checks that the daemon came back. +apply_gateway_config() { + local name="$1" + timeout "$CONFIG_RESTART_TIMEOUT" docker exec "$name" bash -c ' + systemctl unmask fips-gateway.service 2>/dev/null + cp /etc/fips/fips.yaml /etc/fips/fips.yaml.orig + cat >> /etc/fips/fips.yaml </dev/null 2>&1 + return } # ───────────────────────────────────────────────────────────────────── @@ -267,8 +313,7 @@ _run_deb_install_scenario() { build_deb || { fail ".deb build failed"; return; } - local cached_deb - cached_deb=$(ls "$DEB_CACHE_DIR"/fips_*_amd64.deb 2>/dev/null | head -1) + local cached_deb="$DEB_PATH" if [ -z "$cached_deb" ] || [ ! -f "$cached_deb" ]; then fail "no .deb available at $DEB_CACHE_DIR" return @@ -276,13 +321,8 @@ _run_deb_install_scenario() { local deb_basename deb_basename=$(basename "$cached_deb") - # Ubuntu 22.04 bundles systemd-resolved into systemd; other - # distros require it as a separate package. Compose the apt - # package list accordingly. - local apt_packages="systemd iproute2 dbus dnsutils procps" - if [ "$base_image" != "ubuntu:22.04" ]; then - apt_packages="systemd systemd-resolved iproute2 dbus dnsutils procps" - fi + local apt_packages + apt_packages=$(runtime_packages "$base_image") log "Building ${base_image} runtime image" cp "$cached_deb" "$CACHE_DIR/deb-for-image" @@ -515,19 +555,7 @@ DOCKERFILE # default preset) and ipv6 forwarding (gateway checks before # the DNS upstream check), which the container is started with; # see start_systemd_container_with_tun. - timeout "$CONFIG_RESTART_TIMEOUT" docker exec "$name" bash -c ' - systemctl unmask fips-gateway.service 2>/dev/null - # Patch in a minimal gateway config since the shipped fips.yaml - # has gateway disabled by default. - cp /etc/fips/fips.yaml /etc/fips/fips.yaml.orig - cat >> /etc/fips/fips.yaml </dev/null 2>&1 + apply_gateway_config "$name" sleep 3 if wait_for_service_active "$name" fips.service 5; then @@ -570,8 +598,474 @@ EOF cleanup_container "$name" } +# ───────────────────────────────────────────────────────────────────── +# Upgrade scenario +# +# Upgrades an installed package to a newer one and checks what the +# maintainer scripts do to the running services on the way. The newer +# package is made from the one under test inside the container: unpacked, +# given a higher Version and repacked, so the upgrade runs this tree's prerm +# and postinst without a second build. +# +# Its runtime image holds no package. The install scenario's image does, but +# under a tag every run shares and a file name every build of one version +# shares, so another run could retag it in the minutes between the two +# scenarios and this one would upgrade from that run's package. Instead the +# package this run built is copied into each container, and its checksum is +# compared there before anything is installed. +# ───────────────────────────────────────────────────────────────────── + +# Every apt run here is bounded: an upgrade that blocks in postinst is one of +# the defects this scenario exists to catch, and an unbounded one would hang +# the suite instead of failing it. +# The package's own worst case on a healthy daemon is its three bounded starts, +# 60s + 60s + 90s, so an upgrade bound above that reports the package's +# diagnosis rather than this one. +UPGRADE_APT_TIMEOUT=300 +# The dead-daemon reinstall skips the units behind the daemon, so its worst +# case is one 60s bound. +DEAD_DAEMON_APT_TIMEOUT=150 +# postinst waits up to 60s for a unit that does not start. When the daemon +# cannot start, apt has to return a failure well inside this. +DEAD_DAEMON_LIMIT=120 +# Bound on a single short command inside a container. +EXEC_TIMEOUT=60 + +# Run a short command in a container under EXEC_TIMEOUT. +cexec() { + local name="$1" + shift + timeout "$EXEC_TIMEOUT" docker exec "$name" "$@" + return +} + +# Run apt-get in /opt/fips-deb inside the container, bounded by the given +# number of seconds, keeping the existing configuration files. Sets APT_RC, +# APT_SECS and APT_OUT rather than returning a status, because every caller +# needs all three. +run_apt() { + local name="$1" limit="$2" + shift 2 + local start=$SECONDS + APT_RC=0 + APT_OUT=$(timeout "$limit" docker exec -w /opt/fips-deb "$name" \ + apt-get -o Dpkg::Options::=--force-confdef -o Dpkg::Options::=--force-confold \ + "$@" 2>&1) || APT_RC=$? + APT_SECS=$((SECONDS - start)) + return 0 +} + +# Boot an upgrade container, copy this run's package into it and install it. +# Returns 1, having recorded why, when any step fails. +upgrade_boot() { + local name="$1" image="$2" deb="$3" + if ! start_systemd_container_with_tun "$name" "$image"; then + fail "$name: container did not start" + return 1 + fi + if ! wait_for_systemd "$name"; then + fail "systemd did not boot in $name" + return 1 + fi + # /opt rather than /tmp: systemd mounts a fresh /tmp during boot. + if ! timeout "$EXEC_TIMEOUT" docker cp "$DEB_PATH" "$name:/opt/fips-deb/$deb"; then + fail "$name: could not copy $deb into the container" + return 1 + fi + local want have + want=$(sha256sum "$DEB_PATH" | cut -d' ' -f1) + have=$(cexec "$name" sha256sum "/opt/fips-deb/$deb" 2>/dev/null | cut -d' ' -f1) + if [ -z "$want" ] || [ "$want" != "$have" ]; then + fail "$name: the package in the container is not the one under test ('$have', want '$want')" + return 1 + fi + local start=$SECONDS rc=0 out + out=$(timeout "$UPGRADE_APT_TIMEOUT" docker exec -w /opt/fips-deb "$name" bash -c " + apt-get update >/dev/null 2>&1 + apt-get install -y --no-install-recommends ./${deb} 2>&1 + ") || rc=$? + echo " install took $((SECONDS - start))s" + if [ "$rc" -ne 0 ]; then + fail "$name: installing $deb exited $rc" + echo "$out" | tail -20 + return 1 + fi + return 0 +} + +# Make /opt/fips-deb/next.deb from the package under test, inside the +# container so the host needs no dpkg tooling. Its Version is the original's +# with "+upgrade1" appended, which must compare higher, and its fips.nft gains +# a named counter inside the fips table, so a check can tell whether the +# ruleset loaded after the upgrade is the new one. Any step failing is a +# failure of the scenario, never a skip. +make_next_package() { + local name="$1" deb="$2" out rc=0 + # shellcheck disable=SC2016 # the script expands inside the container + out=$(timeout "$EXEC_TIMEOUT" docker exec -w /opt/fips-deb -e DEB="$deb" "$name" \ + bash -euo pipefail -c ' + rm -rf /root/next + dpkg-deb -R "./$DEB" /root/next + old=$(dpkg-deb -f "./$DEB" Version) + new="${old}+upgrade1" + sed -i "s/^Version: .*/Version: ${new}/" /root/next/DEBIAN/control + dpkg --compare-versions "$new" gt "$old" + nft_file=/root/next/etc/fips/fips.nft + grep -q "^table inet fips {\$" "$nft_file" + sed -i "/^table inet fips {\$/a\\ counter fips_upgrade_probe { packets 0 bytes 0 }" "$nft_file" + grep -q "counter fips_upgrade_probe" "$nft_file" + nft -c -f "$nft_file" + if grep -q " etc/fips/fips.nft\$" /root/next/DEBIAN/md5sums 2>/dev/null; then + sum=$(md5sum "$nft_file" | cut -d" " -f1) + sed -i "s|^[0-9a-f]* etc/fips/fips.nft\$|${sum} etc/fips/fips.nft|" /root/next/DEBIAN/md5sums + grep -q "^${sum} etc/fips/fips.nft\$" /root/next/DEBIAN/md5sums + fi + dpkg-deb -b /root/next /opt/fips-deb/next.deb >/dev/null + echo "made next.deb at Version $new" + ' 2>&1) || rc=$? + if [ "$rc" -ne 0 ]; then + fail "$name: could not make the newer package (exit $rc)" + echo "$out" | tail -20 + return 1 + fi + echo " $out" + return 0 +} + +# Start fips.service and fips-dns.service the way the install scenario does +# and require both to be active. +start_daemon_units() { + local name="$1" + start_unit "$name" fips.service >/dev/null || true + start_unit_queued "$name" fips-dns.service >/dev/null || true + if wait_for_service_active "$name" fips.service && + wait_for_service_active "$name" fips-dns.service; then + return 0 + fi + cexec "$name" systemctl status --no-pager fips.service fips-dns.service 2>&1 | tail -20 + return 1 +} + +# Pass or fail on whether a unit is active. +check_active() { + local name="$1" unit="$2" what="$3" + if cexec "$name" systemctl is-active --quiet "$unit"; then + pass "$what: $unit active" + else + fail "$what: $unit not active" + cexec "$name" systemctl status --no-pager "$unit" 2>&1 | tail -15 + fi + return 0 +} + +# Pass or fail on whether a unit the host never enabled is still neither +# running nor enabled. +check_left_off() { + local name="$1" unit="$2" what="$3" state + state=$(cexec "$name" systemctl is-enabled "$unit" 2>/dev/null || true) + if ! cexec "$name" systemctl is-active --quiet "$unit" && [ "$state" = "disabled" ]; then + pass "$what: $unit inactive and disabled" + else + fail "$what: $unit is $(cexec "$name" systemctl is-active "$unit" 2>/dev/null) and '$state' (want inactive and disabled)" + fi + return 0 +} + +# Start `nft monitor tables` in the background, writing to +# /root/nft-monitor.log, and prove it is recording by adding and deleting a +# table of its own. An empty log from a monitor that never ran would otherwise +# read as a ruleset that was never removed. The probe table's name does not +# begin with "fips", so it cannot match a check on the fips table. +start_nft_monitor() { + local name="$1" + if ! timeout "$EXEC_TIMEOUT" docker exec -d "$name" \ + sh -c 'exec nft monitor tables > /root/nft-monitor.log 2>&1'; then + return 1 + fi + sleep 1 + cexec "$name" sh -c 'nft add table inet monprobe && nft delete table inet monprobe' || return 1 + local _i + for _i in 1 2 3 4 5; do + if cexec "$name" grep -Eq '^delete table inet monprobe( |$)' /root/nft-monitor.log; then + return 0 + fi + sleep 1 + done + cexec "$name" cat /root/nft-monitor.log 2>&1 | tail -10 + return 1 +} + +# Apply the gateway config and start fips-gateway.service, bounded, requiring +# it to be active. Leaves it enabled or not as the caller already set it. +start_gateway() { + local name="$1" + apply_gateway_config "$name" + if ! wait_for_service_active "$name" fips.service 10; then + return 1 + fi + start_unit "$name" fips-gateway.service "$GATEWAY_START_TIMEOUT" >/dev/null 2>&1 || true + if wait_for_service_active "$name" fips-gateway.service 10; then + return 0 + fi + cexec "$name" systemctl status --no-pager fips-gateway.service 2>&1 | tail -15 + return 1 +} + +# Print a unit's MainPID; 0 when it has no main process. +main_pid() { + local name="$1" unit="$2" + cexec "$name" systemctl show -p MainPID --value "$unit" 2>/dev/null || echo 0 + return 0 +} + +# Pass or fail on whether a unit runs a new process of the installed binary +# after the upgrade: active, a MainPID other than the one before, and an +# executable that is the installed file rather than one the upgrade replaced, +# which the kernel reports with a " (deleted)" suffix. +check_new_binary() { + local name="$1" unit="$2" binary="$3" before="$4" what="$5" pid exe + pid=$(main_pid "$name" "$unit") + exe=$(cexec "$name" readlink "/proc/$pid/exe" 2>/dev/null || true) + if cexec "$name" systemctl is-active --quiet "$unit" && [ "$pid" != 0 ] && + [ "$pid" != "$before" ] && [ "$exe" = "$binary" ]; then + pass "$what: $unit runs the upgraded $binary" + else + fail "$what: $unit is $(cexec "$name" systemctl is-active "$unit" 2>/dev/null), MainPID $before -> $pid, exe '$exe' (want active, a new process, $binary)" + fi + return 0 +} + +# Host that opted in to the firewall and enabled the gateway: the upgrade must +# apply the new ruleset in place, with no moment at which the fips table is +# absent, and bring the gateway back on the new binary. Purging the package +# must then leave no enablement behind for the gateway. +_upgrade_opted_in() { + local name="$1" image="$2" deb="$3" + log "upgrade on a host that opted in ($name)" + upgrade_boot "$name" "$image" "$deb" || { cleanup_container "$name"; return 0; } + make_next_package "$name" "$deb" || { cleanup_container "$name"; return 0; } + if ! cexec "$name" systemctl enable --now fips-firewall.service >/dev/null 2>&1 || + ! start_daemon_units "$name"; then + fail "opted in: the firewall, fips and fips-dns did not all start before the upgrade" + cexec "$name" systemctl status --no-pager fips-firewall.service 2>&1 | tail -15 + cleanup_container "$name" + return 0 + fi + cexec "$name" systemctl enable fips-gateway.service >/dev/null 2>&1 + if ! start_gateway "$name"; then + fail "opted in: fips-gateway did not start before the upgrade" + cleanup_container "$name" + return 0 + fi + if ! start_nft_monitor "$name"; then + fail "opted in: nft monitor is not observing table changes" + cleanup_container "$name" + return 0 + fi + local fips_pid gw_pid + fips_pid=$(main_pid "$name" fips.service) + gw_pid=$(main_pid "$name" fips-gateway.service) + + run_apt "$name" "$UPGRADE_APT_TIMEOUT" install -y ./next.deb + echo " upgrade took ${APT_SECS}s" + if [ "$APT_RC" -eq 0 ]; then + pass "opted in: upgrade exits 0" + else + fail "opted in: upgrade exited $APT_RC" + echo "$APT_OUT" | tail -20 + fi + check_active "$name" fips.service "opted in, after upgrade" + check_active "$name" fips-dns.service "opted in, after upgrade" + check_active "$name" fips-firewall.service "opted in, after upgrade" + if cexec "$name" nft list counter inet fips fips_upgrade_probe >/dev/null 2>&1; then + pass "opted in: the upgraded ruleset is loaded" + else + fail "opted in: the upgraded ruleset is not loaded (no fips_upgrade_probe counter)" + fi + if cexec "$name" grep -Eq '^delete table inet fips( |$)' /root/nft-monitor.log; then + fail "opted in: the fips table was deleted during the upgrade" + cexec "$name" cat /root/nft-monitor.log 2>&1 | tail -10 + else + pass "opted in: the fips table was never deleted during the upgrade" + fi + check_new_binary "$name" fips.service /usr/bin/fips "$fips_pid" "opted in, after upgrade" + check_new_binary "$name" fips-gateway.service /usr/bin/fips-gateway "$gw_pid" "opted in, after upgrade" + + run_apt "$name" "$UPGRADE_APT_TIMEOUT" purge -y fips + echo " purge took ${APT_SECS}s" + if [ "$APT_RC" -ne 0 ]; then + fail "opted in: purge exited $APT_RC" + echo "$APT_OUT" | tail -20 + fi + local link=/etc/systemd/system/multi-user.target.wants/fips-gateway.service state + state=$(cexec "$name" systemctl is-enabled fips-gateway.service 2>/dev/null || true) + if ! cexec "$name" test -e "$link" && ! cexec "$name" test -L "$link" && + [ "$state" != "enabled" ]; then + pass "opted in, after purge: no fips-gateway enablement left behind" + else + fail "opted in, after purge: fips-gateway still enabled ('$state', $(cexec "$name" ls -l "$link" 2>&1))" + fi + + cleanup_container "$name" + return 0 +} + +# Host that never opted in to the firewall and ran the gateway without enabling +# it: the upgrade must leave neither running nor enabled. Then the package is +# reinstalled three times: with the daemon masked, when apt must succeed and +# start nothing; with an enabled gateway that cannot start, when apt must +# succeed, say so, and leave the daemon running; and with a daemon that cannot +# start, when apt must fail, promptly, naming the unit, rather than wait for +# ever on a unit that requires a daemon which never comes up. +_upgrade_not_opted_in() { + local name="$1" image="$2" deb="$3" + log "upgrade on a host that never opted in ($name)" + upgrade_boot "$name" "$image" "$deb" || { cleanup_container "$name"; return 0; } + make_next_package "$name" "$deb" || { cleanup_container "$name"; return 0; } + if ! start_daemon_units "$name" || ! start_gateway "$name"; then + fail "not opted in: fips, fips-dns and fips-gateway did not start before the upgrade" + cleanup_container "$name" + return 0 + fi + + run_apt "$name" "$UPGRADE_APT_TIMEOUT" install -y ./next.deb + echo " upgrade took ${APT_SECS}s" + if [ "$APT_RC" -eq 0 ]; then + pass "not opted in: upgrade exits 0" + else + fail "not opted in: upgrade exited $APT_RC" + echo "$APT_OUT" | tail -20 + fi + check_active "$name" fips.service "not opted in, after upgrade" + check_active "$name" fips-dns.service "not opted in, after upgrade" + check_left_off "$name" fips-firewall.service "not opted in, after upgrade" + check_left_off "$name" fips-gateway.service "not opted in, after upgrade" + if cexec "$name" nft list table inet fips >/dev/null 2>&1; then + fail "not opted in, after upgrade: the fips firewall table is loaded" + else + pass "not opted in, after upgrade: no fips firewall table" + fi + + # A host that masked the daemon on purpose: the upgrade must skip it with a + # message, not fail. The package before this change printed nothing for a + # masked unit, so the message is what tells the two apart. + cexec "$name" bash -c 'systemctl stop fips-dns.service fips.service; systemctl mask fips.service' \ + >/dev/null 2>&1 + run_apt "$name" "$UPGRADE_APT_TIMEOUT" install --reinstall -y ./next.deb + echo " reinstall with the daemon masked took ${APT_SECS}s (exit $APT_RC)" + if [ "$APT_RC" -eq 0 ] && grep -q "fips.service is masked" <<<"$APT_OUT" && + ! cexec "$name" systemctl is-active --quiet fips.service && + ! cexec "$name" systemctl is-active --quiet fips-dns.service; then + pass "masked daemon: apt succeeds, says the unit was skipped and starts nothing" + else + fail "masked daemon: apt exited $APT_RC (want 0, a message that fips.service is masked, and fips and fips-dns inactive)" + echo "$APT_OUT" | tail -20 + fi + cexec "$name" bash -c 'systemctl unmask fips.service; systemctl daemon-reload' >/dev/null 2>&1 + if ! start_daemon_units "$name"; then + fail "dead daemon: fips and fips-dns did not start again after unmasking" + cleanup_container "$name" + return 0 + fi + + # An enabled gateway that fails on start, behind a daemon that is healthy: + # the gateway is an opt-in addition, so apt must report it and succeed. + # Restart=no so it reaches failed at once rather than looping. + cexec "$name" bash -c ' + mkdir -p /etc/systemd/system/fips-gateway.service.d + printf "[Service]\nExecStart=\nExecStart=/bin/false\nRestart=no\n" \ + > /etc/systemd/system/fips-gateway.service.d/broken.conf + systemctl daemon-reload + systemctl enable fips-gateway.service + ' >/dev/null 2>&1 + run_apt "$name" "$UPGRADE_APT_TIMEOUT" install --reinstall -y ./next.deb + echo " reinstall with a broken gateway took ${APT_SECS}s (exit $APT_RC)" + if [ "$APT_RC" -eq 0 ] && + grep -q "fips-gateway.service did not come back" <<<"$APT_OUT" && + cexec "$name" systemctl is-active --quiet fips.service; then + pass "broken gateway: apt succeeds, reports the gateway and leaves the daemon running" + else + fail "broken gateway: apt exited $APT_RC (want 0, a message that fips-gateway.service did not come back, and fips active)" + echo "$APT_OUT" | tail -20 + fi + cexec "$name" bash -c ' + systemctl disable fips-gateway.service + rm -rf /etc/systemd/system/fips-gateway.service.d + systemctl daemon-reload + systemctl reset-failed fips-gateway.service + ' >/dev/null 2>&1 + + # A daemon that fails on every start: exit 1, so Restart=on-failure loops, + # and the start job of fips-dns, which requires it, is never dispatched. + cexec "$name" bash -c ' + mkdir -p /etc/systemd/system/fips.service.d + printf "[Service]\nExecStart=\nExecStart=/bin/false\n" \ + > /etc/systemd/system/fips.service.d/broken.conf + systemctl daemon-reload + ' + run_apt "$name" "$DEAD_DAEMON_APT_TIMEOUT" install --reinstall -y ./next.deb + echo " reinstall with a dead daemon took ${APT_SECS}s (exit $APT_RC)" + if [ "$APT_RC" -ne 0 ] && [ "$APT_RC" -ne 124 ] && + [ "$APT_SECS" -lt "$DEAD_DAEMON_LIMIT" ] && + grep -q "fips.service did not become active" <<<"$APT_OUT"; then + pass "dead daemon: apt fails in ${APT_SECS}s and names fips.service" + else + fail "dead daemon: apt exited $APT_RC after ${APT_SECS}s (want a failure under ${DEAD_DAEMON_LIMIT}s naming fips.service; 124 is the harness bound)" + echo "$APT_OUT" | tail -20 + fi + + cleanup_container "$name" + return 0 +} + +_run_deb_upgrade_scenario() { + local distro_label="$1" + local base_image="$2" + # Scoped to the run like the container names, although it holds no + # package, so concurrent runs never rebuild an image under each other. + local image="fips-deb-upgrade:${distro_label}${FIPS_CI_NAME_SUFFIX:-}" + log ".deb upgrade: ${base_image}" + + if [ -z "$DEB_PATH" ] || [ ! -f "$DEB_PATH" ]; then + fail "no package to upgrade from" + return + fi + local deb + deb=$(basename "$DEB_PATH") + + log "Building $image (runtime packages and nftables, no fips package)" + build_image "$image" "$(cat </dev/null 2>&1 || true + return 0 +} + # Per-distro wrappers -test_debian12() { _run_deb_install_scenario debian12 debian:12; } +# debian12 also runs the upgrade scenario. One distro keeps the suite's cost +# down; this one because a oneshot start behind a daemon in its restart loop +# waits for ever on its systemd (252), while on Ubuntu 22.04's (249) the start +# returns with an error, so only here does the upgrade scenario see the hang. +test_debian12() { + _run_deb_install_scenario debian12 debian:12 + _run_deb_upgrade_scenario debian12 debian:12 +} test_debian13() { _run_deb_install_scenario debian13 debian:trixie; } test_ubuntu22() { _run_deb_install_scenario ubuntu22 ubuntu:22.04; } test_ubuntu24() { _run_deb_install_scenario ubuntu24 ubuntu:24.04; } diff --git a/testing/lib/wait-converge-test.sh b/testing/lib/wait-converge-test.sh index 48763531..a99974b1 100755 --- a/testing/lib/wait-converge-test.sh +++ b/testing/lib/wait-converge-test.sh @@ -7,7 +7,8 @@ # network are involved, so the suite is hermetic and safe to run in CI. # # It is not fast, though: it drives real timeouts against the real clock -# and takes about 45 seconds, measured 2026-07-23. The header claimed "a +# and took about 45 seconds when measured 2026-07-23, before case 7 added +# roughly another minute. The header claimed "a # few seconds" from the day it was written until then, which nothing had # contradicted because no runner had ever invoked it. # @@ -122,6 +123,59 @@ ping_quick_converge() { fi } +# Case 7 traces: the near-converged acceptance window. Each is keyed off +# PT like the traces above, and each is shaped so its expected verdict holds +# with at least a second of margin on either side of the accept window. + +# 7d: still progressing until t=4, so the hold that follows (armed at t=6 +# with stall_secs=2) has lasted only ~2s when the cap falls at 8s, short of +# a 4s accept window. +ping_late_progress() { + set_pt; local t=$PT + if (( t < 4 )); then + PASSED=$(( 15 + t )); FAILED=$(( 5 - t )) + else + PASSED=19; FAILED=1 + fi +} + +# 7f: holds at 19/1 from the start and converges at t=6, before a 10s cap. +# A gate that accepted as soon as the accept window elapsed (t~4) would +# report near_converged here instead of converged. +ping_converges_after_window() { + set_pt; local t=$PT + if (( t < 6 )); then + PASSED=19; FAILED=1 + else + PASSED=20; FAILED=0 + fi +} + +# 7g: enters the hold at 18/2 (armed at t=2), makes progress to 19/1 at +# t=4, then holds there long enough (re-armed at t=6, cap 11) to exceed a +# 3s accept window. Reds under a gate that arms the hold once and never +# re-arms it after progress. +ping_rearm_hold() { + set_pt; local t=$PT + if (( t < 4 )); then + PASSED=18; FAILED=2 + else + PASSED=19; FAILED=1 + fi +} + +# 7h: the same shape, but the progress to 19/1 comes at t=8, so the second +# hold (re-armed at t=10) is ~1s old when the cap falls at 11s. Reds under a +# gate that keeps the first hold's start across the progress. +ping_late_rearm() { + set_pt; local t=$PT + if (( t < 8 )); then + PASSED=18; FAILED=2 + else + PASSED=19; FAILED=1 + fi +} + HOLD_MSG="holding for full budget" STUCK_MSG="STUCK" NOCONV_MSG="tree did not converge" @@ -292,6 +346,111 @@ check "case6c: verdict carries a clean tree" "$c6c_cnt_ok" \ c6c_quiet_ok=0; grep -q "$NOCONV_MSG" "$VERDICT_OUT" && c6c_quiet_ok=1 check "case6c: no non-convergence message on a clean run" "$c6c_quiet_ok" +# --- Case 7: near-converged acceptance at the hard cap ----------------- +# +# The sixth argument lets a caller hand a mesh that has held within slack for +# at least that many seconds to its own strict assertion instead of failing +# the gate. It acts only at the hard cap, which is reachable only in states +# that are red without it, so it can never turn a run that would have +# converged by the cap into a red. With the argument absent or 0 the gate +# must behave exactly as before. +echo +echo "== Case 7: near-converged acceptance at the hard cap ==" + +echo "-- Case 7a: held within slack past the accept window, accepted at the cap --" +run_gate ping_never_converges 6 3 1 2 2; rc=$? +cat "$VERDICT_OUT" +c7a_rc_ok=1; [ "$rc" -eq 0 ] && c7a_rc_ok=0 +check "case7a: near-converged hold past the window returns 0" "$c7a_rc_ok" "rc=$rc" +c7a_out_ok=1; [ "$CONVERGE_OUTCOME" = "near_converged" ] && c7a_out_ok=0 +check "case7a: verdict is near_converged" "$c7a_out_ok" "CONVERGE_OUTCOME=$CONVERGE_OUTCOME" +c7a_cnt_ok=1 +[ "$CONVERGE_REACHED" -eq 19 ] && [ "$CONVERGE_PENDING" -eq 1 ] && c7a_cnt_ok=0 +check "case7a: verdict carries the shortfall" "$c7a_cnt_ok" \ + "reached=$CONVERGE_REACHED pending=$CONVERGE_PENDING" + +echo "-- Case 7b: the same trace with five arguments, then with an explicit 0 --" +run_gate ping_never_converges 6 3 1 2; rc=$? +cat "$VERDICT_OUT" +c7b_rc_ok=1; [ "$rc" -eq 1 ] && c7b_rc_ok=0 +check "case7b: five arguments still time out" "$c7b_rc_ok" "rc=$rc" +c7b_out_ok=1; [ "$CONVERGE_OUTCOME" = "timeout" ] && c7b_out_ok=0 +check "case7b: five arguments give verdict timeout" "$c7b_out_ok" "CONVERGE_OUTCOME=$CONVERGE_OUTCOME" +run_gate ping_never_converges 6 3 1 2 0; rc=$? +cat "$VERDICT_OUT" +c7b0_rc_ok=1; [ "$rc" -eq 1 ] && c7b0_rc_ok=0 +check "case7b: an explicit 0 still times out" "$c7b0_rc_ok" "rc=$rc" +c7b0_out_ok=1; [ "$CONVERGE_OUTCOME" = "timeout" ] && c7b0_out_ok=0 +check "case7b: an explicit 0 gives verdict timeout" "$c7b0_out_ok" "CONVERGE_OUTCOME=$CONVERGE_OUTCOME" + +echo "-- Case 7c: wedged far from convergence, acceptance does not rescue it --" +run_gate ping_far_stall 30 4 1 2 2; rc=$? +cat "$VERDICT_OUT" +c7c_rc_ok=1; [ "$rc" -eq 1 ] && c7c_rc_ok=0 +check "case7c: a stall beyond slack still reds with acceptance enabled" "$c7c_rc_ok" "rc=$rc" +c7c_out_ok=1; [ "$CONVERGE_OUTCOME" = "stalled" ] && c7c_out_ok=0 +check "case7c: verdict is stalled" "$c7c_out_ok" "CONVERGE_OUTCOME=$CONVERGE_OUTCOME" + +echo "-- Case 7d: hold shorter than the accept window at the cap --" +run_gate ping_late_progress 8 2 1 2 4; rc=$? +cat "$VERDICT_OUT" +c7d_rc_ok=1; [ "$rc" -eq 1 ] && c7d_rc_ok=0 +check "case7d: a hold shorter than the window still times out" "$c7d_rc_ok" "rc=$rc" +c7d_out_ok=1; [ "$CONVERGE_OUTCOME" = "timeout" ] && c7d_out_ok=0 +check "case7d: verdict is timeout" "$c7d_out_ok" "CONVERGE_OUTCOME=$CONVERGE_OUTCOME" + +echo "-- Case 7e: converged tree with acceptance enabled --" +run_gate ping_quick_converge 20 4 1 2 2; rc=$? +cat "$VERDICT_OUT" +c7e_rc_ok=1; [ "$rc" -eq 0 ] && c7e_rc_ok=0 +check "case7e: converged tree passes with acceptance enabled" "$c7e_rc_ok" "rc=$rc" +c7e_out_ok=1; [ "$CONVERGE_OUTCOME" = "converged" ] && c7e_out_ok=0 +check "case7e: verdict is converged" "$c7e_out_ok" "CONVERGE_OUTCOME=$CONVERGE_OUTCOME" + +echo "-- Case 7f: converges after the window would have elapsed, before the cap --" +run_gate ping_converges_after_window 10 2 1 2 2; rc=$? +cat "$VERDICT_OUT" +c7f_rc_ok=1; [ "$rc" -eq 0 ] && c7f_rc_ok=0 +check "case7f: late convergence before the cap passes" "$c7f_rc_ok" "rc=$rc" +c7f_out_ok=1; [ "$CONVERGE_OUTCOME" = "converged" ] && c7f_out_ok=0 +check "case7f: acceptance waits for the cap, verdict is converged" "$c7f_out_ok" \ + "CONVERGE_OUTCOME=$CONVERGE_OUTCOME" + +echo "-- Case 7g: hold, progress, then a second hold past the window --" +run_gate ping_rearm_hold 11 2 1 2 3; rc=$? +cat "$VERDICT_OUT" +c7g_rc_ok=1; [ "$rc" -eq 0 ] && c7g_rc_ok=0 +check "case7g: a re-armed hold past the window returns 0" "$c7g_rc_ok" "rc=$rc" +c7g_out_ok=1; [ "$CONVERGE_OUTCOME" = "near_converged" ] && c7g_out_ok=0 +check "case7g: verdict is near_converged" "$c7g_out_ok" "CONVERGE_OUTCOME=$CONVERGE_OUTCOME" +c7g_cnt_ok=1 +[ "$CONVERGE_REACHED" -eq 19 ] && [ "$CONVERGE_PENDING" -eq 1 ] && c7g_cnt_ok=0 +check "case7g: verdict carries the shortfall" "$c7g_cnt_ok" \ + "reached=$CONVERGE_REACHED pending=$CONVERGE_PENDING" + +echo "-- Case 7h: hold, late progress, second hold shorter than the window --" +run_gate ping_late_rearm 11 2 1 2 3; rc=$? +cat "$VERDICT_OUT" +c7h_rc_ok=1; [ "$rc" -eq 1 ] && c7h_rc_ok=0 +check "case7h: progress resets the hold, so a short second hold times out" "$c7h_rc_ok" "rc=$rc" +c7h_out_ok=1; [ "$CONVERGE_OUTCOME" = "timeout" ] && c7h_out_ok=0 +check "case7h: verdict is timeout" "$c7h_out_ok" "CONVERGE_OUTCOME=$CONVERGE_OUTCOME" + +echo "-- Case 7i: hold armed, but shorter than the window when the cap falls --" +# The hold arms at t~3 (stall_secs after the first reading) and the cap falls +# at 6, so it is ~3s old against a 10s window: seven seconds of margin, far +# above SECONDS' one-second granularity. This is the case that carries the +# duration condition; 7d and 7h pass even with it removed if their hold +# happens not to arm. +run_gate ping_never_converges 6 3 1 2 10; rc=$? +cat "$VERDICT_OUT" +c7i_rc_ok=1; [ "$rc" -eq 1 ] && c7i_rc_ok=0 +check "case7i: an armed hold shorter than the window still times out" "$c7i_rc_ok" "rc=$rc" +c7i_out_ok=1; [ "$CONVERGE_OUTCOME" = "timeout" ] && c7i_out_ok=0 +check "case7i: verdict is timeout" "$c7i_out_ok" "CONVERGE_OUTCOME=$CONVERGE_OUTCOME" +c7i_hold_ok=1; grep -q "$HOLD_MSG" "$VERDICT_OUT" && c7i_hold_ok=0 +check "case7i: the hold was armed" "$c7i_hold_ok" "expected '$HOLD_MSG' in output" + rm -f "$VERDICT_OUT" # --- Summary ---------------------------------------------------------- diff --git a/testing/lib/wait-converge.sh b/testing/lib/wait-converge.sh index a9588b15..dd7f42b7 100644 --- a/testing/lib/wait-converge.sh +++ b/testing/lib/wait-converge.sh @@ -7,7 +7,7 @@ # source "$(dirname "$0")/../../lib/wait-converge.sh" # wait_for_peers [timeout_secs] # wait_until_connected [poll_secs] \ -# [near_converged_slack] +# [near_converged_slack] [near_converged_accept_secs] # # wait_until_connected also sets CONVERGE_OUTCOME / CONVERGE_REACHED / # CONVERGE_PENDING; see the block above it. @@ -61,7 +61,7 @@ wait_for_peers() { # Verdict of the most recent wait_until_connected() call, so a caller can # report WHICH condition failed rather than only that one did: -# CONVERGE_OUTCOME converged | stalled | timeout +# CONVERGE_OUTCOME converged | near_converged | stalled | timeout # CONVERGE_REACHED reachable pairs at the moment of the verdict # CONVERGE_PENDING unreachable pairs at that moment # @@ -70,6 +70,10 @@ wait_for_peers() { # the strict all-pairs assertion 20/20. Without them the caller's summary # line reads "20 passed, 0 failed" on a non-convergence exit, which a # reader cannot tell from a connectivity failure. +# +# near_converged is a success return (0): the hard cap fell while the mesh +# had held within slack for at least near_converged_accept_secs, and the +# caller asked for that state to be handed to its own strict assertion. CONVERGE_OUTCOME="" CONVERGE_REACHED=0 CONVERGE_PENDING=0 @@ -87,7 +91,7 @@ _converge_verdict() { # progress-aware deadline instead of a fixed one. # # wait_until_connected [poll_secs] \ -# [near_converged_slack] +# [near_converged_slack] [near_converged_accept_secs] # # is the name of a function that runs the suite's own # connectivity check and sets two globals each call: @@ -112,19 +116,34 @@ _converge_verdict() { # rather than emitting a false RED with budget still unspent. A # genuinely never-converging single pair still hits the hard cap. # - hard cap: max_secs elapsed -> return 1 (never runs unbounded). +# - near-converged acceptance, off unless near_converged_accept_secs is +# non-zero: at the hard cap, if the current near-converged hold has +# lasted at least that long and the last poll was still within slack, +# return 0 with verdict near_converged instead. It acts only at the cap, +# a state that is red without it, so it can never turn a run that would +# have converged by the cap into a red one. Acceptance means "hand the +# mesh to the caller's strict assertion", never "skip that assertion". +# Progress disarms the hold, so only the hold that ends at the cap counts. # -# Returns 0 once fully connected, 1 on stall or timeout. +# Returns 0 once fully connected or accepted as near-converged, 1 on stall +# or timeout. wait_until_connected() { local ping_fn="$1" local max_secs="$2" local stall_secs="$3" local poll_secs="${4:-1}" local near_converged_slack="${5:-2}" + local near_converged_accept_secs="${6:-0}" local start_secs=$SECONDS local best=-1 local last_progress=$SECONDS local held_for_budget=0 + # Start of the current near-converged hold, or -1 when disarmed. Kept + # apart from held_for_budget, which is set once and never reset and so + # cannot measure the hold that ends at the cap. -1 rather than 0 because + # SECONDS can legitimately be 0. + local hold_start=-1 while (( SECONDS - start_secs < max_secs )); do "$ping_fn" @@ -136,6 +155,7 @@ wait_until_connected() { if (( PASSED > best )); then best=$PASSED last_progress=$SECONDS + hold_start=-1 echo " converge: $PASSED reachable, $FAILED pending (progressing) after $((SECONDS - start_secs))s" elif (( SECONDS - last_progress >= stall_secs )); then if (( FAILED > near_converged_slack )); then @@ -143,6 +163,9 @@ wait_until_connected() { echo " converge: STUCK — tree did not converge: $PASSED reachable / $FAILED pending, no progress for ${stall_secs}s (after $((SECONDS - start_secs))s)" return 1 fi + if (( hold_start < 0 )); then + hold_start=$SECONDS + fi if (( held_for_budget == 0 )); then held_for_budget=1 echo " converge: near-converged ($PASSED reachable / $FAILED pending <= slack=$near_converged_slack) — holding for full budget, not bailing (after $((SECONDS - start_secs))s)" @@ -151,6 +174,17 @@ wait_until_connected() { sleep "$poll_secs" done + # The slack test repeats what the loop already guarantees (once the hold + # is armed, any poll beyond slack exits as stalled), so that acceptance + # stays tied to slack if the loop's exits are ever reordered. + if (( near_converged_accept_secs > 0 && hold_start >= 0 \ + && SECONDS - hold_start >= near_converged_accept_secs \ + && FAILED <= near_converged_slack )); then + _converge_verdict near_converged + echo " converge: near-converged at the cap — $PASSED reachable / $FAILED pending, held $((SECONDS - hold_start))s >= ${near_converged_accept_secs}s; handing to the strict assertion (after ${max_secs}s)" + return 0 + fi + _converge_verdict timeout echo " converge: TIMEOUT — tree did not converge: $PASSED reachable / $FAILED pending after ${max_secs}s" return 1 diff --git a/testing/static/scripts/gateway-test.sh b/testing/static/scripts/gateway-test.sh index 3153bf71..dd441b91 100755 --- a/testing/static/scripts/gateway-test.sh +++ b/testing/static/scripts/gateway-test.sh @@ -153,6 +153,31 @@ if [ "$DNS_READY" != true ]; then echo " WARNING: Gateway DNS did not respond within 30s, continuing anyway" fi +# The gateway names its conntrack source once at startup, before the DNS +# resolver starts, so by now the line is in the log. Ask the gateway's own +# namespace which source it should have found: the proc file when it exists, +# and otherwise the netlink dump, which the container's NET_ADMIN allows. A +# failed `docker logs` reds the check rather than counting as zero lines. +if docker exec "$GATEWAY" test -e /proc/net/nf_conntrack; then + EXPECT_SRC=proc +else + EXPECT_SRC=netlink +fi +if GW_START_LOG=$(docker logs "$GATEWAY" 2>&1); then + SRC_PROC=$(grep -cF 'Conntrack source: proc; session pinning is on' <<< "$GW_START_LOG" || true) + SRC_NETLINK=$(grep -cF 'Conntrack source: netlink; session pinning is on' <<< "$GW_START_LOG" || true) + SRC_NONE=$(grep -cF 'No conntrack source is readable; session pinning is off' <<< "$GW_START_LOG" || true) + case "$EXPECT_SRC" in + proc) SRC_HIT=$SRC_PROC ;; + *) SRC_HIT=$SRC_NETLINK ;; + esac + SRC_ALL=$((SRC_PROC + SRC_NETLINK + SRC_NONE)) + SRC_OK=$([ "$SRC_HIT" -eq 1 ] && [ "$SRC_ALL" -eq 1 ] && echo 0 || echo 1) + check "Conntrack source line at startup (expect $EXPECT_SRC; proc lines $SRC_PROC, netlink lines $SRC_NETLINK, none lines $SRC_NONE)" "$SRC_OK" +else + check "Conntrack source line at startup (docker logs failed)" 1 +fi + # Phase 3: Client network setup — route virtual IP pool via gateway echo "" echo "Phase 3: Client network setup" @@ -278,6 +303,43 @@ else check "nftables DNAT rules" 1 fi +# Phase 5's GET left a TCP conntrack entry to the first virtual IP, which +# stays in the table in TIME_WAIT well past two ticks. The gateway must count +# it. Poll about once a second for 25 tries, which spans two 10s ticks, and +# read the count with a parser that cannot turn a failed query into a number: +# an error response has no `data`, and the parser exits non-zero on it. +if [ -n "$VIRTUAL_IP" ]; then + VIP_SESSIONS=error + for _ in $(seq 1 25); do + VIP_SESSIONS=$(docker exec "$GATEWAY" bash -c \ + 'echo "{\"command\":\"show_mappings\"}" | nc -U -w1 /run/fips/gateway.sock 2>/dev/null' \ + | VIP="$VIRTUAL_IP" python3 -c " +import os, sys, json +r = json.load(sys.stdin) +data = r.get('data') +if not isinstance(data, dict) or not isinstance(data.get('mappings'), list): + sys.exit(1) +hits = [m for m in data['mappings'] if m.get('virtual_ip') == os.environ['VIP']] +if len(hits) != 1 or not isinstance(hits[0].get('sessions'), int): + sys.exit(1) +print(hits[0]['sessions']) +" 2>/dev/null || echo "error") + if [ "$VIP_SESSIONS" != error ] && [ "$VIP_SESSIONS" -ge 1 ]; then + break + fi + sleep 1 + done + if [ "$VIP_SESSIONS" != error ] && [ "$VIP_SESSIONS" -ge 1 ]; then + check "Gateway counts a session to $VIRTUAL_IP (sessions $VIP_SESSIONS, source $EXPECT_SRC)" 0 + else + check "Gateway counts a session to $VIRTUAL_IP (sessions $VIP_SESSIONS, source $EXPECT_SRC)" 1 + echo " Kernel conntrack entries to $VIRTUAL_IP:" + docker exec "$GATEWAY" conntrack -L -f ipv6 -d "$VIRTUAL_IP" 2>&1 | sed 's/^/ /' || true + fi +else + check "Gateway counts a session (skipped — no virtual IP)" 1 +fi + # Phase 7: Inbound port forwarding — UDP and a second simultaneous TCP forward. # # Three forwards exercised: diff --git a/testing/static/scripts/rekey-test.sh b/testing/static/scripts/rekey-test.sh index 1e774c40..2ca46a75 100755 --- a/testing/static/scripts/rekey-test.sh +++ b/testing/static/scripts/rekey-test.sh @@ -154,6 +154,12 @@ trap 'echo ""; echo "Test interrupted"; exit 130' INT # wait early-exits on PASS, so successful reps are unaffected by the # extra headroom. BASELINE_CONVERGENCE_TIMEOUT=65 +# A mesh that has held within the gate's slack (two pairs) for this long when +# BASELINE_CONVERGENCE_TIMEOUT falls is handed to the strict all-pairs +# assertion instead of failing the gate. The strict assertion still decides. +# 10s accepts a straggler that held for most of the window and refuses a +# pair that only fell into the hold in the last few seconds. +BASELINE_NEAR_CONVERGED_ACCEPT=10 REKEY_SETTLE=12 # > DRAIN_WINDOW_SECS (10) so post-rekey samples are off the old session # First FMP rekey should follow shortly after the 35s interval once the mesh is # fully converged. Keep this bounded to preserve a meaningful scheduling check @@ -421,7 +427,13 @@ echo "" # it no longer false-times-out under concurrent CI load. The strict # ping_all below is the actual assertion, run only after convergence. echo "Phase 1: Pre-rekey connectivity (waiting for convergence)" -if wait_until_connected _baseline_ping "$BASELINE_CONVERGENCE_TIMEOUT" 20; then +if wait_until_connected _baseline_ping "$BASELINE_CONVERGENCE_TIMEOUT" 20 1 2 \ + "$BASELINE_NEAR_CONVERGED_ACCEPT"; then + if [ "$CONVERGE_OUTCOME" = "near_converged" ]; then + echo " Gate accepted a near-converged mesh" \ + "($CONVERGE_REACHED/$((CONVERGE_REACHED + CONVERGE_PENDING)) reachable);" \ + "the strict assertion below decides" + fi ping_all "" "$TIMEOUT" "$MAX_PING_ATTEMPTS" phase_result "Pre-rekey baseline (all 20 pairs)" if [ "$FAILED" -ne 0 ]; then