Merge master into next: package service and firewall upgrade handling, FSP rekey and session setup fixes, gateway conntrack over netlink

The session rekey fix takes next's XX form here: the rekey-initiator arm
reads msg2 with try_read_message_2, the rolling-back XX read the
first-contact arm already uses, and puts the handshake back on a failed
read. The forged-ack test builder builds an XX-sized msg2, and the
first-contact test uses it. The test for a rekey msg2 lost in transit is
not carried, as the forged-msg2 test it shares setup with was not; only
the pump helper the coordinate tests use comes across.
This commit is contained in:
Johnathan Corgan
2026-09-19 05:11:12 +00:00
24 changed files with 2515 additions and 125 deletions
+37
View File
@@ -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
Generated
+37 -2
View File
@@ -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"
+5
View File
@@ -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]
+7 -2
View File
@@ -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
+27
View File
@@ -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
+28 -4
View File
@@ -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
+14 -1
View File
@@ -2,7 +2,8 @@
# Build a .deb package for FIPS using cargo-deb.
#
# Usage: ./build-deb.sh [--target <triple>] [--version <version>] [--no-build]
# [--features <list>]
# [--features <list>] [--output-dir <dir>]
# [--name-file <path>]
#
# Prerequisites: cargo-deb (install with: cargo install cargo-deb)
# Output: deploy/fips_<version>_<arch>.deb
@@ -26,6 +27,10 @@ Options:
--output-dir <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 <path> Also write the finished package's file name (basename
only) to <path>. 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}"
+3
View File
@@ -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
+137 -4
View File
@@ -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
+4
View File
@@ -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
+3
View File
@@ -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
+44 -4
View File
@@ -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
+330
View File
@@ -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<ConntrackSnapshot, io::Error> {
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<u8> {
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<Ipv6Addr, u32>,
) -> Result<DumpState, io::Error> {
let mut offset = 0;
while offset < buf.len() {
let message = NetlinkMessage::<NetfilterMessage>::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<Ipv6Addr, u32>) {
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<Tuple> {
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<Tuple>, reply: Vec<Tuple>) -> 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<NetfilterMessage>, seq: u32) -> Vec<u8> {
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<NetfilterMessage> {
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<DumpState, io::Error>, HashMap<Ipv6Addr, u32>) {
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
);
}
}
+1
View File
@@ -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;
+262 -7
View File
@@ -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<std::io::Error>,
}
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<P = ProcConntrack, N = super::conntrack::NetlinkConntrack> {
proc: P,
netlink: N,
}
impl<P: ConntrackQuerier, N: ConntrackQuerier> SystemConntrack<P, N> {
/// 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<P: ConntrackQuerier, N: ConntrackQuerier> ConntrackQuerier for SystemConntrack<P, N> {
fn snapshot(&self) -> Result<ConntrackSnapshot, std::io::Error> {
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<P: ConntrackQuerier, N: ConntrackQuerier>(
reader: &SystemConntrack<P, N>,
) -> 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<Option<std::io::ErrorKind>>,
@@ -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<std::io::ErrorKind>);
impl ConntrackQuerier for FixedRead {
fn snapshot(&self) -> Result<ConntrackSnapshot, std::io::Error> {
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<ConntrackSnapshot, std::io::Error> {
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
);
}
}
+12 -2
View File
@@ -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();
+57 -12
View File
@@ -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<crate::proto::stp::TreeCoordinate> {
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
+43 -13
View File
@@ -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 {
+667 -38
View File
@@ -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<u8> {
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<ReceivedPacket> =
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<TestNode> {
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<Vec<NodeAddr>> {
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<NodeAddr> {
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
// ============================================================================
+524 -30
View File
@@ -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 <<EOF
gateway:
enabled: true
pool: "fd01::/112"
lan_interface: "eth0"
EOF
systemctl restart fips.service
' >/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 <<EOF
gateway:
enabled: true
pool: "fd01::/112"
lan_interface: "eth0"
EOF
systemctl restart fips.service
' >/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 <<DOCKERFILE
FROM ${base_image}
ENV DEBIAN_FRONTEND=noninteractive
RUN apt-get update && apt-get install -y --no-install-recommends \\
$(runtime_packages "$base_image") nftables && \\
apt-get clean && rm -rf /var/lib/apt/lists/* && \\
systemctl enable systemd-resolved && \\
mkdir -p /opt/fips-deb
CMD ["/lib/systemd/systemd"]
DOCKERFILE
)" || {
fail "upgrade image build failed"
return
}
# Named under the install scenario's prefix, so a CI step that collects
# that scenario's container logs on failure collects these too.
_upgrade_opted_in "fips-deb-test-${distro_label}-upg-a${FIPS_CI_NAME_SUFFIX:-}" "$image" "$deb"
_upgrade_not_opted_in "fips-deb-test-${distro_label}-upg-b${FIPS_CI_NAME_SUFFIX:-}" "$image" "$deb"
docker rmi "$image" >/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; }
+160 -1
View File
@@ -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 ----------------------------------------------------------
+38 -4
View File
@@ -7,7 +7,7 @@
# source "$(dirname "$0")/../../lib/wait-converge.sh"
# wait_for_peers <container> <min_peers> [timeout_secs]
# wait_until_connected <ping_fn> <max_secs> <stall_secs> [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 <ping_fn> <max_secs> <stall_secs> [poll_secs] \
# [near_converged_slack]
# [near_converged_slack] [near_converged_accept_secs]
#
# <ping_fn> 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
+62
View File
@@ -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:
+13 -1
View File
@@ -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