From 0f2ba939fd508fc41a870aa9ee1dc8f5fdd278b1 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 1 Oct 2026 19:29:04 +0000 Subject: [PATCH 01/23] Iterate the gateway's UDP socket inodes from a variable shellcheck reports SC2013 on the for loop over a command substitution in gateway_dns_held_by, and the OpenWrt Package workflow fails its lint step on it. Capture the awk output in a variable and iterate that instead; the words the loop sees are the same. cat stays in front of awk because busybox awk gives up on a missing /proc/net/udp6, and a while-read loop at the end of the pipe would run in a subshell under ash, where returning from the function is not possible. --- packaging/openwrt-ipk/files/etc/init.d/fips-gateway | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/packaging/openwrt-ipk/files/etc/init.d/fips-gateway b/packaging/openwrt-ipk/files/etc/init.d/fips-gateway index bc897b79..4d577455 100755 --- a/packaging/openwrt-ipk/files/etc/init.d/fips-gateway +++ b/packaging/openwrt-ipk/files/etc/init.d/fips-gateway @@ -185,11 +185,15 @@ dns_locked() { # is among the process's open descriptors. Another process holding the port, # such as an mDNS responder on 5353, does not count. gateway_dns_held_by() { - local hex inode fd + local hex inodes inode fd [ -n "$2" ] || return 1 hex="$(printf '%04X' "$1")" - for inode in $(cat /proc/net/udp /proc/net/udp6 2>/dev/null | - awk -v want=":$hex" 'substr($2, length($2) - 4) == want { print $10 }'); do + # cat rather than awk's own file arguments: busybox awk gives up on a + # missing /proc/net/udp6. A while-read loop on the pipe would run in a + # subshell under ash, where the return below could not leave this function. + inodes="$(cat /proc/net/udp /proc/net/udp6 2>/dev/null | + awk -v want=":$hex" 'substr($2, length($2) - 4) == want { print $10 }')" + for inode in $inodes; do for fd in /proc/"$2"/fd/*; do [ "$(readlink "$fd" 2>/dev/null)" = "socket:[$inode]" ] && return 0 done From a67aadbd49b761f3c8caf2804e6caf8dea70b60b Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 1 Oct 2026 22:40:40 +0000 Subject: [PATCH 02/23] Lint the OpenWrt package's shell scripts on every branch, not only on trunk pushes The OpenWrt Package workflow was the only place the package's shell scripts were linted, and it runs on trunk pushes, tags and pull requests but never on a branch push. A shellcheck finding introduced on a topic branch therefore first failed after the branch had landed, which is how all three trunks went red on the gateway init script. testing/check-shellcheck.sh now runs that pass: the scripts the package ships as POSIX sh with the workflow's exclusions, and the nak installer as bash. It also covers the maintainer scripts the workflow's list left out: the postinst and prerm that both the .ipk and the .apk package ship, and the preinst that the SDK feed Makefile ships. It fails when a shell script under the package's files or scripts directory is not on its list, so a new one cannot go unlinted. It exits 2 when shellcheck is missing or cannot read a file, 1 on a finding or a missing script, and 0 when clean. The OpenWrt Package workflow calls the guard in place of its two inline steps, keeping its install-if-missing step, so the runners cannot drift apart. ci-local.sh runs it as a static stage beside the other repository guards, and ci.yml's ci-parity job runs it on every branch push. --- .github/workflows/ci.yml | 11 ++ .github/workflows/package-openwrt.yml | 50 +-------- testing/check-shellcheck.sh | 144 ++++++++++++++++++++++++++ testing/ci-local.sh | 12 +++ 4 files changed, 172 insertions(+), 45 deletions(-) create mode 100755 testing/check-shellcheck.sh diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index f93d144c..daf8c116 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -79,6 +79,17 @@ jobs: run: bash testing/check-comment-refs.sh - name: Check no non-test code uses std 64-bit atomics run: python3 testing/check-portable-atomics.py + # The OpenWrt Package workflow runs this too, but only on trunk pushes, + # tags and pull requests; here a branch push sees a finding first. + # Kept in step with ci-local.sh's run_shellcheck by hand. + - name: Install shellcheck (if missing) + run: | + if ! command -v shellcheck >/dev/null 2>&1; then + sudo apt-get update && sudo apt-get install -y --no-install-recommends shellcheck + fi + shellcheck --version + - name: Check the OpenWrt package's shell scripts with shellcheck + run: bash testing/check-shellcheck.sh # Hermetic: synthetic ping functions, no containers, ~45s. Lives beside # the other two so both runners gate on it identically — putting it in # only one would create exactly the drift check-ci-parity.sh exists to diff --git a/.github/workflows/package-openwrt.yml b/.github/workflows/package-openwrt.yml index c9945dd6..880ca962 100644 --- a/.github/workflows/package-openwrt.yml +++ b/.github/workflows/package-openwrt.yml @@ -341,52 +341,12 @@ jobs: fi shellcheck --version - # Its own step, and its own shell dialect. The shipped-scripts lint below - # runs --shell=sh with the OpenWrt rc.common exclude set, which misfires - # on a bash script; install-nak.sh is also not shipped in the package. - - name: Lint install-nak.sh + # The scripts this package ships, as sh, and install-nak.sh, as bash. + # The guard is the one copy of the lint; ci.yml and ci-local.sh run it + # too, so a branch push sees a finding before it reaches a trunk. + - name: Lint shell scripts shell: bash - run: shellcheck --shell=bash .github/scripts/install-nak.sh - - - name: Lint shipped shell scripts - shell: bash - run: | - set -euo pipefail - FILES_DIR=packaging/openwrt-ipk/files - # Scripts shipped inside the .ipk. The init scripts use the OpenWrt - # `#!/bin/sh /etc/rc.common` shebang; tell shellcheck to treat them - # as POSIX sh and silence the unrecognized-shebang warning (SC1008). - # SC2317 (unreachable command) fires on rc.common's externally-invoked - # start_service/stop_service/reload_service hooks. - TARGETS=( - "$FILES_DIR/etc/init.d/fips" - "$FILES_DIR/etc/init.d/fips-gateway" - "$FILES_DIR/etc/fips/firewall.sh" - "$FILES_DIR/etc/hotplug.d/net/99-fips" - "$FILES_DIR/etc/uci-defaults/90-fips-setup" - "$FILES_DIR/usr/bin/fips-mesh-setup" - "$FILES_DIR/usr/bin/fips-ap-setup" - ) - fail=0 - for f in "${TARGETS[@]}"; do - if [ ! -f "$f" ]; then - echo "FAIL: missing $f" - fail=1 - continue - fi - echo "==> shellcheck $f" - if shellcheck --shell=sh --exclude=SC1008,SC2317,SC2034,SC3043,SC2086,SC2089,SC2090 "$f"; then - echo " PASS" - else - echo " FAIL" - fail=1 - fi - done - if [ "$fail" -ne 0 ]; then - echo "shellcheck FAILED" - exit 1 - fi - echo "shellcheck PASS (${#TARGETS[@]} scripts)" + run: bash testing/check-shellcheck.sh - name: Sysctl drop-in syntax check shell: bash diff --git a/testing/check-shellcheck.sh b/testing/check-shellcheck.sh new file mode 100755 index 00000000..2967a91b --- /dev/null +++ b/testing/check-shellcheck.sh @@ -0,0 +1,144 @@ +#!/bin/bash +# ── OpenWrt shell-script lint guard ───────────────────────────────────────── +# Runs shellcheck over the shell scripts the OpenWrt packages ship, and over +# .github/scripts/install-nak.sh, which the OpenWrt Package workflow runs to +# fetch its publishing tool. +# +# This is the one copy of that lint. The OpenWrt Package workflow calls it +# after building each .ipk, and ci.yml and ci-local.sh call it too, because +# that workflow runs only on trunk pushes, tags and pull requests: without the +# other two, a finding in a script edited on a topic branch first shows up +# after the branch has reached a trunk. +# +# What is checked, and how: +# * Every script under packaging/openwrt-ipk/files/ and the maintainer +# scripts under packaging/openwrt-ipk/scripts/, as POSIX sh. Both the .ipk +# and the .apk package take their payload and maintainer scripts from these +# two directories (build-apk.sh wraps the maintainer scripts with one +# header line, which is not linted separately); the SDK feed Makefile ships +# preinst as well. On a router they run under busybox ash. +# * install-nak.sh as bash, with no exclusions. It is a CI script, not a +# shipped one, and the sh exclusion set below misfires on bash. +# +# The sh exclusions, with reason: +# SC1008 the init scripts' `#!/bin/sh /etc/rc.common` shebang, an +# interpreter line the linter does not recognise. +# SC2317 rc.common's start_service/stop_service/reload_service hooks, +# which nothing in the file itself calls. +# SC2034 the init scripts' USE_PROCD, START, STOP, EXTRA_COMMANDS and +# EXTRA_HELP, which rc.common reads rather than the script. +# SC3043 `local`, which POSIX leaves undefined and ash supports. +# SC2086, SC2089, SC2090 firewall.sh builds an nft match, quotes included, +# in one variable and relies on word splitting to pass it as +# separate arguments; nft parses the quotes itself. +# +# Every sh-family script in those two directories must be on the list below. A +# new one that is not fails the guard, so a script added to the package is not +# silently left unlinted. +# +# Exit 0 = clean. Exit 1 = a finding, a listed script missing, or a shipped +# script not on the list. Exit 2 = the guard could not run (shellcheck or git +# missing, or shellcheck could not process a file); never treated as a pass. +# ───────────────────────────────────────────────────────────────────────────── +set -uo pipefail + +SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd)" +PROJECT_ROOT="$(cd "$SCRIPT_DIR/.." && pwd)" +cd "$PROJECT_ROOT" || { echo "check-shellcheck: cannot cd to $PROJECT_ROOT" >&2; exit 2; } + +IPK=packaging/openwrt-ipk +SH_EXCLUDE=SC1008,SC2317,SC2034,SC3043,SC2086,SC2089,SC2090 + +SH_TARGETS=( + "$IPK/files/etc/init.d/fips" + "$IPK/files/etc/init.d/fips-gateway" + "$IPK/files/etc/fips/firewall.sh" + "$IPK/files/etc/hotplug.d/net/99-fips" + "$IPK/files/etc/uci-defaults/90-fips-setup" + "$IPK/files/usr/bin/fips-mesh-setup" + "$IPK/files/usr/bin/fips-ap-setup" + "$IPK/scripts/preinst" + "$IPK/scripts/postinst" + "$IPK/scripts/prerm" +) +BASH_TARGETS=( + ".github/scripts/install-nak.sh" +) + +if ! command -v shellcheck >/dev/null 2>&1; then + echo "check-shellcheck: shellcheck not found; cannot lint the shell scripts" >&2 + echo "check-shellcheck: install it with 'apt-get install shellcheck'" >&2 + exit 2 +fi +if ! command -v git >/dev/null 2>&1; then + echo "check-shellcheck: git not found; cannot list the shipped scripts" >&2 + exit 2 +fi + +shellcheck --version | sed -n 's/^version: /check-shellcheck: shellcheck /p' + +findings=0 +broken=0 + +# ── Completeness: every shipped sh-family script is on the list ───────────── +shipped=$(git ls-files -- "$IPK/files" "$IPK/scripts") || { + echo "check-shellcheck: git ls-files failed, refusing to pass" >&2 + exit 2 +} +if [[ -z "$shipped" ]]; then + echo "check-shellcheck: git ls-files found nothing under $IPK, refusing to pass" >&2 + exit 2 +fi +while IFS= read -r f; do + # A tracked file deleted from the working tree: if listed, the lint below + # reports it missing; if not, there is nothing to ship. + [[ -f "$f" ]] || continue + head -n 1 "$f" | grep -qE '^#![[:space:]]*[^[:space:]]*/(env[[:space:]]+)?(ba|a|da)?sh([[:space:]]|$)' || continue + listed=0 + for t in "${SH_TARGETS[@]}"; do + [[ "$t" == "$f" ]] && { listed=1; break; } + done + if [[ $listed -eq 0 ]]; then + echo "FAIL: $f is a shipped shell script missing from SH_TARGETS in $0" + findings=1 + fi +done <<< "$shipped" + +# ── Lint ───────────────────────────────────────────────────────────────────── +lint() { + # lint : one file; sets findings or broken. + local f="$1" rc=0 + shift + if [[ ! -f "$f" ]]; then + echo "FAIL: missing $f" + findings=1 + return 0 + fi + echo "==> shellcheck $* $f" + shellcheck "$@" "$f" || rc=$? + case $rc in + 0) echo " PASS" ;; + 1) echo " FAIL"; findings=1 ;; + *) echo " shellcheck exited $rc: could not check $f"; broken=1 ;; + esac + return 0 +} + +for f in "${SH_TARGETS[@]}"; do + lint "$f" --shell=sh --exclude="$SH_EXCLUDE" +done +for f in "${BASH_TARGETS[@]}"; do + lint "$f" --shell=bash +done + +total=$(( ${#SH_TARGETS[@]} + ${#BASH_TARGETS[@]} )) +if [[ $broken -ne 0 ]]; then + echo "shellcheck could not check every script; refusing to pass" + exit 2 +fi +if [[ $findings -ne 0 ]]; then + echo "shellcheck FAILED" + exit 1 +fi +echo "shellcheck PASS ($total scripts: ${#SH_TARGETS[@]} as sh, ${#BASH_TARGETS[@]} as bash)" +exit 0 diff --git a/testing/ci-local.sh b/testing/ci-local.sh index 7a0aec69..f0831825 100755 --- a/testing/ci-local.sh +++ b/testing/ci-local.sh @@ -1564,6 +1564,17 @@ run_portable_atomics() { record "portable-atomics" $rc } +# The shell scripts the OpenWrt packages ship, and the nak installer. The +# OpenWrt Package workflow lints them on GitHub, but only for trunk pushes, +# tags and pull requests, so this is where a branch first sees a finding. +# Mirrored in ci.yml's ci-parity job by hand. Static, about a second. +run_shellcheck() { + local rc=0 + info "[shellcheck] Linting the OpenWrt package's shell scripts" + bash "$SCRIPT_DIR/check-shellcheck.sh" || rc=$? + record "shellcheck" $rc +} + # Every daemon log string a test matches on must still be emitted by src/. # A stale one does not fail — it stops observing, and an expect-zero assertion # built on it then passes for the wrong reason. @@ -1647,6 +1658,7 @@ main() { run_action_pins run_comment_refs run_portable_atomics + run_shellcheck run_wait_converge run_deb_version run_nextest_flaky From 4654e41f871409e47eb0342187741e67695c5302 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 1 Oct 2026 22:40:40 +0000 Subject: [PATCH 03/23] Keep verified coordinates on a single PathBroken and demote them only on reports from two links A PathBroken signal carries no end-to-end authentication, yet one admitted signal deleted the destination's cached coordinates even when a lookup had verified them. That removed the only thing refusing the next forged coordinate warm, so one forged PathBroken followed by one forged SessionSetup moved a node's route to any destination it had a session with. An admitted PathBroken now treats a verified entry differently from a hint: - A hint, or a verification that has aged out, is still removed. - A verified entry is kept while the lookup re-validates it. It is demoted to a hint, keeping its value, only when signals naming the destination arrive over two different links within 15 seconds. - The re-lookup now runs on every admitted PathBroken, not only when the destination's identity is cached. The vote is the authenticated link peer a signal arrived over, not the reporter it names: the reporter is plaintext the sender chooses, so every signal from one neighbour, forged or relayed, is one vote however many reporters it names. A node whose signals all arrive over one link, such as a leaf with a single peer, never reaches the quorum; a stale verified entry there lasts until the lookup each signal starts answers and replaces it, or at most until its verification ages out after five minutes, when the next signal removes it as it would a hint. The link quorum is a new sans-IO module with an injected clock. A lookup that verifies a destination again clears its quorum, so a report about the old path cannot combine with one about the new. When the path-MTU release fires, a kept entry also forgets the path MTU stored with it, as the removed entry used to; when the release is rate limited, a kept entry keeps it. New error-signal counters broken_below_quorum and broken_demoted count the two outcomes for a verified entry. Two advisory counters measure how often an admitted PathBroken looks implausible, without refusing anything: - broken_link_mismatch: the signal arrived over a link other than the one this node would forward to the destination on. The check uses the non-touching next-hop preview, so it does not refresh the cache entry. - broken_reporter_mismatch: the reporter is this node or the destination, or its known coordinates (a direct peer's from the tree, otherwise the coordinate cache) are not strictly closer to the destination than this node's. Both have a non-zero healthy floor: a genuine report can arrive off the forward link, and a reporter's view of the destination can differ from this node's. Most reporters' coordinates are not known here and are not counted. They size the forged-signal problem; they are not alarms, and the signal is acted on in full either way. The handler now receives the authenticated link peer the datagram arrived over, which the quorum and the link check need. The source-recovery steps in the mesh operation design described a PathBroken as removing the destination's cached coordinates and discovery as depending on a cached identity; they now describe the behaviour above. --- docs/design/fips-mesh-operation.md | 21 +- src/cache/coord_cache.rs | 69 +++ src/cache/entry.rs | 5 + src/control/snapshots/show_routing.json | 4 + src/node/handlers/lookup.rs | 3 + src/node/handlers/session.rs | 138 ++++- src/node/metrics.rs | 29 + src/node/mod.rs | 8 + src/node/stats.rs | 4 + src/node/tests/coord_forgery.rs | 760 ++++++++++++++++++++++++ src/node/tests/mod.rs | 1 + src/node/tests/session.rs | 23 +- src/proto/fsp/core.rs | 43 +- src/proto/fsp/mod.rs | 3 + src/proto/fsp/quorum.rs | 126 ++++ src/proto/fsp/tests/core.rs | 39 +- src/proto/fsp/tests/mod.rs | 1 + src/proto/fsp/tests/quorum.rs | 118 ++++ 18 files changed, 1342 insertions(+), 53 deletions(-) create mode 100644 src/node/tests/coord_forgery.rs create mode 100644 src/proto/fsp/quorum.rs create mode 100644 src/proto/fsp/tests/quorum.rs diff --git a/docs/design/fips-mesh-operation.md b/docs/design/fips-mesh-operation.md index 6ecbf9b0..91ccd88e 100644 --- a/docs/design/fips-mesh-operation.md +++ b/docs/design/fips-mesh-operation.md @@ -372,10 +372,27 @@ source. 1. Immediately send a standalone CoordsWarmup (0x14) message (rate-limited, same per-destination interval as CoordsRequired response) -2. Remove stale coordinates from cache -3. Initiate discovery for the destination +2. Handle the cached coordinates by where they came from. Unverified + coordinates (a hint copied off a passing packet, or a lookup result + whose verification has aged out) are removed. Coordinates a lookup + verified are kept while discovery re-validates them, because the signal + is unauthenticated and removing them would let the next forged warm + replace them. They are demoted to an unverified hint, keeping their value, + only when PathBroken signals naming the destination arrive over two + different links within 15 seconds. The vote is the authenticated link + peer, not the reporter the signal names, which the sender chooses. A node + whose signals all arrive over one link never demotes this way; its + verified coordinates last until discovery replaces them or their + verification ages out after 300 seconds. +3. Initiate discovery for the destination, whether or not its identity is + cached 4. Reset CP warmup counter +The source also counts, without refusing anything, a PathBroken that +arrives over a link other than its forward link to the destination, and one +whose reporter is not closer to the destination than the source is. Both +have a non-zero healthy floor. + ### MtuExceeded **Trigger**: A transit node receives a SessionDatagram but the total diff --git a/src/cache/coord_cache.rs b/src/cache/coord_cache.rs index d2acd8ae..9390803e 100644 --- a/src/cache/coord_cache.rs +++ b/src/cache/coord_cache.rs @@ -247,6 +247,27 @@ impl CoordCache { self.entries.remove(addr) } + /// Demote an entry to an unverified hint, keeping its coordinates and TTL. + /// + /// The entry goes on routing, but a hint may now replace it. Returns + /// whether an entry existed. + pub fn demote(&mut self, addr: &NodeAddr) -> bool { + match self.entries.get_mut(addr) { + Some(entry) => { + entry.mark_hint(); + true + } + None => false, + } + } + + /// Forget the path MTU stored with an entry, keeping the entry. + pub fn clear_path_mtu(&mut self, addr: &NodeAddr) { + if let Some(entry) = self.entries.get_mut(addr) { + entry.clear_path_mtu(); + } + } + /// Check if an address is cached (and not expired). pub fn contains(&self, addr: &NodeAddr, current_time_ms: u64) -> bool { self.get(addr, current_time_ms).is_some() @@ -844,4 +865,52 @@ mod tests { assert!(cache.contains(&make_node_addr(1), 10)); assert!(cache.contains(&make_node_addr(2), 10)); } + + #[test] + fn demote_keeps_the_value_and_lets_a_hint_replace_it() { + let mut cache = CoordCache::new(100, 1000); + let addr = make_node_addr(1); + let real = make_coords(&[1, 0]); + cache.insert_verified_with_path_mtu(addr, real.clone(), 10, 1400); + assert_eq!( + cache.insert(addr, make_coords(&[1, 2, 0]), 10), + HintOutcome::Rejected + ); + + assert!(cache.demote(&addr)); + + let entry = cache.get_entry(&addr).unwrap(); + assert_eq!(entry.coords(), &real); + assert_eq!(entry.source(), crate::cache::CoordSource::Hint); + assert!(!entry.is_verified(10)); + assert_eq!( + entry.path_mtu(), + Some(1400), + "demote leaves the path MTU to its caller" + ); + assert_eq!( + cache.insert(addr, make_coords(&[1, 2, 0]), 11), + HintOutcome::Changed + ); + } + + #[test] + fn demote_of_an_absent_entry_reports_none() { + let mut cache = CoordCache::new(100, 1000); + assert!(!cache.demote(&make_node_addr(1))); + assert!(cache.is_empty()); + } + + #[test] + fn clear_path_mtu_keeps_the_entry() { + let mut cache = CoordCache::new(100, 1000); + let addr = make_node_addr(1); + cache.insert_verified_with_path_mtu(addr, make_coords(&[1, 0]), 10, 1400); + + cache.clear_path_mtu(&addr); + + let entry = cache.get_entry(&addr).unwrap(); + assert_eq!(entry.path_mtu(), None); + assert!(entry.is_verified(10)); + } } diff --git a/src/cache/entry.rs b/src/cache/entry.rs index 8fe95ae0..6b62f358 100644 --- a/src/cache/entry.rs +++ b/src/cache/entry.rs @@ -133,6 +133,11 @@ impl CacheEntry { self.path_mtu = Some(mtu); } + /// Forget the path MTU, as when the path it described is released. + pub fn clear_path_mtu(&mut self) { + self.path_mtu = None; + } + /// Check if this entry has expired. pub fn is_expired(&self, current_time_ms: u64) -> bool { current_time_ms > self.expires_at diff --git a/src/control/snapshots/show_routing.json b/src/control/snapshots/show_routing.json index 0e0685da..f64fe2ee 100644 --- a/src/control/snapshots/show_routing.json +++ b/src/control/snapshots/show_routing.json @@ -36,6 +36,10 @@ "resp_unsolicited": 0 }, "error_signals": { + "broken_below_quorum": 0, + "broken_demoted": 0, + "broken_link_mismatch": 0, + "broken_reporter_mismatch": 0, "coords_required": 0, "emit_limiter_at_capacity": 0, "emit_over_dest_interval": 0, diff --git a/src/node/handlers/lookup.rs b/src/node/handlers/lookup.rs index 90a45208..3d3791e6 100644 --- a/src/node/handlers/lookup.rs +++ b/src/node/handlers/lookup.rs @@ -375,6 +375,9 @@ impl Node { now_ms, path_mtu, } => { + // Reports about the path this verified value replaces are + // not evidence against it. + self.broken_quorum.clear(&target); // The annotation is unsigned and accumulates hop by hop, so // any forwarder on the reverse path can lower it. A value // below the actionable floor cannot describe a usable path, diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index cb343bce..2b897aab 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -18,6 +18,7 @@ use crate::noise::{ use crate::proto::fmp::wire::{ ESTABLISHED_HEADER_SIZE, FLAG_KEY_EPOCH, FLAG_SP, build_established_header, }; +use crate::proto::fsp::quorum::QuorumVerdict; use crate::proto::fsp::wire::{ FSP_COMMON_PREFIX_SIZE, FSP_FLAG_CP, FSP_FLAG_K, FSP_HEADER_SIZE, FSP_PHASE_ESTABLISHED, FSP_PHASE_MSG1, FSP_PHASE_MSG2, FSP_PHASE_MSG3, FSP_PORT_HEADER_SIZE, FSP_PORT_IPV6_SHIM, @@ -193,7 +194,8 @@ impl Node { self.handle_coords_required(src_addr, error_body).await; } Some(RoutingSignalType::PathBroken) => { - self.handle_path_broken(src_addr, error_body).await; + self.handle_path_broken(src_addr, link_peer, error_body) + .await; } Some(RoutingSignalType::MtuExceeded) => { self.handle_mtu_exceeded(src_addr, error_body).await; @@ -1919,12 +1921,28 @@ impl Node { /// Handle a PathBroken error signal from a transit router. /// /// The router has coordinates but still can't route to the destination. - /// Send a standalone CoordsWarmup immediately (rate-limited), invalidate - /// cached coordinates, trigger re-discovery, and reset the warmup counter. + /// Send a standalone CoordsWarmup immediately (rate-limited), re-validate + /// the destination's coordinates by lookup, release its path MTU, and + /// reset the warmup counter. + /// + /// Cached coordinates that are only a hint are removed. Coordinates a + /// lookup verified are kept while the lookup re-validates them, and are + /// demoted to a hint only once reports arriving over distinct links reach + /// the quorum: the signal is unauthenticated, and deleting a verified + /// entry on one report is what let the next forged warm replace it. /// /// `src_addr` is the datagram's claimed source and is not - /// end-to-end authenticated; see `signal_verdict`. - pub(in crate::node) async fn handle_path_broken(&mut self, src_addr: &NodeAddr, inner: &[u8]) { + /// end-to-end authenticated; see `signal_verdict`. `link_peer` is the + /// authenticated peer the datagram arrived over. It is the report's vote + /// in the quorum, because the body's reporter is whatever the sender + /// wrote, and it is compared with the forward path to count mismatches; + /// it never refuses the signal. + pub(in crate::node) async fn handle_path_broken( + &mut self, + src_addr: &NodeAddr, + link_peer: &NodeAddr, + inner: &[u8], + ) { self.metrics().errors.path_broken.inc(); let msg = match PathBroken::decode(inner) { @@ -1936,11 +1954,10 @@ impl Node { }; // The premise: this signal carries no end-to-end authentication, so - // the body's `dest_addr` is attacker-chosen. `plan_path_broken` emits - // its coord-cache invalidation unconditionally, and the path-MTU - // release below is likewise unguarded, so both act on whatever address - // the body names unless the gate refuses it here, in the shell, which - // is the only layer that knows who sent the datagram. + // the body's `dest_addr` is attacker-chosen. The coord-cache action, + // the lookup and the path-MTU release below all act on whatever + // address the body names unless the gate refuses it here, in the + // shell, which is the only layer that knows who sent the datagram. let verdict = self.signal_verdict(src_addr, &msg.dest_addr); if verdict != SignalVerdict::Admit { debug!(src = %src_addr, dest = %msg.dest_addr, reporter = %msg.reporter, @@ -1960,6 +1977,7 @@ impl Node { reporter = %msg.reporter, "PathBroken: transit router reports routing failure" ); + self.count_path_broken_mismatches(link_peer, &msg); // Send standalone CoordsWarmup immediately (rate-limited) if self @@ -1978,46 +1996,68 @@ impl Node { "PathBroken response rate-limited, skipping standalone CoordsWarmup"); } - // Invalidate stale cached coordinates, then (only if the target's - // identity is cached — else the LookupResponse proof cannot be verified, - // e.g. when the XK responder receives PathBroken before msg3 completes) - // trigger re-discovery. The core emits invalidate-then-lookup in order. - let has_cached_identity = self.has_cached_identity(&msg.dest_addr); - let actions = self - .fsp - .plan_path_broken(msg.dest_addr, has_cached_identity); + // Only a live verified entry has anything for the quorum to protect, + // so only a report against one is recorded; that also bounds the + // quorum's keys by the destinations this node has looked up. The vote + // is the link the report arrived over, so a sender on one link counts + // once however many reporters it names. + let now = Self::now_ms(); + let verified = self + .coord_cache + .get_entry(&msg.dest_addr) + .is_some_and(|e| e.is_verified(now)); + let quorum = if verified { + self.broken_quorum.record(msg.dest_addr, *link_peer, now) + } else { + QuorumVerdict::Below { distinct: 0 } + }; + let actions = self.fsp.plan_path_broken(msg.dest_addr, verified, quorum); for action in actions { match action { FspAction::InvalidateCoords { addr } => { self.coord_cache.remove(&addr); } + FspAction::DemoteCoords { addr } => { + self.coord_cache.demote(&addr); + self.metrics().errors.broken_demoted.inc(); + debug!(dest = %addr, + "PathBroken quorum reached; demoted verified coordinates to a hint"); + } FspAction::InitiateLookup { dest } => { self.maybe_initiate_lookup(&dest).await; } _ => {} } } + if let QuorumVerdict::Below { distinct } = quorum + && verified + { + self.metrics().errors.broken_below_quorum.inc(); + debug!(dest = %msg.dest_addr, distinct, + "PathBroken below quorum; keeping verified coordinates pending the lookup"); + } + // The path this destination's stored MTU described is gone, so release // it rather than carrying it onto whatever path replaces it. Rate // limited per destination on its own budget: PathBroken is // unauthenticated, and an unlimited release discards a genuinely // learned bottleneck as fast as it is relearned. The budget is not // shared with any other signal, so nothing else can spend it. + // + // A cache entry kept above still carries the MTU the lookup stored + // with it, which describes the same path; clear it with the map so + // the two do not disagree. if self .path_mtu_release_limiter .should_send(&msg.dest_addr, Self::now_ms()) { self.path_mtu_lookup_release(&msg.dest_addr); + self.coord_cache.clear_path_mtu(&msg.dest_addr); } else { trace!(dest = %msg.dest_addr, "PathBroken path MTU release rate-limited, keeping the stored value"); } - if !has_cached_identity { - debug!(dest = %msg.dest_addr, - "Skipping discovery after PathBroken: no cached identity for target"); - } - // Reset coords warmup counter so the next N packets include // COORDS_PRESENT, re-warming transit caches along the new path. let n = self.config().node.session.coords_warmup_packets; @@ -2031,6 +2071,58 @@ impl Node { } } + /// Count, without acting on, the two ways an admitted PathBroken can + /// disagree with this node's own view of the path. + /// + /// Advisory only: both checks read state an attacker can influence and + /// both have a non-zero healthy floor, so they size the problem rather + /// than refuse anything. + /// + /// - Link: the signal arrived over a link other than the one this node + /// would forward to the destination on. A genuine report can also do + /// this when the reverse path differs from the forward one. + /// - Reporter: the reporter is this node or the destination, or its known + /// coordinates are not strictly closer to the destination than this + /// node's. Forwarding makes strict progress under each forwarder's own + /// view of the destination, and a PathBroken arises exactly where that + /// view may differ from this node's, so a genuine report can count here + /// too. A reporter's coordinates are rarely known unless it is a direct + /// peer, so this mostly reads as unknown and is not counted. + fn count_path_broken_mismatches(&self, link_peer: &NodeAddr, msg: &PathBroken) { + let now = Self::now_ms(); + let dest = &msg.dest_addr; + let reporter = &msg.reporter; + + if let (Some(hop), _) = self.preview_next_hop(dest, now) + && hop.node_addr != *link_peer + { + self.metrics().errors.broken_link_mismatch.inc(); + debug!(dest = %dest, link_peer = %link_peer, next_hop = %hop.node_addr, + "PathBroken arrived off the forward link"); + } + + let implausible = if reporter == self.node_addr() || reporter == dest { + true + } else { + let dest_coords = self.coord_cache.get(dest, now); + let reporter_coords = self + .tree_state + .peer_coords(reporter) + .or_else(|| self.coord_cache.get(reporter, now)); + match (dest_coords, reporter_coords) { + (Some(d), Some(r)) => { + r.distance_to(d) >= self.tree_state.my_coords().distance_to(d) + } + _ => false, + } + }; + if implausible { + self.metrics().errors.broken_reporter_mismatch.inc(); + debug!(dest = %dest, reporter = %reporter, + "PathBroken reporter is not closer to the destination than this node"); + } + } + /// Handle an MtuExceeded error signal from a transit router. /// /// A transit router couldn't forward our packet because it exceeded the diff --git a/src/node/metrics.rs b/src/node/metrics.rs index 77f9ca6d..984c19c8 100644 --- a/src/node/metrics.rs +++ b/src/node/metrics.rs @@ -536,6 +536,31 @@ pub struct ErrorMetrics { /// rising count is a forged or stale reactive signal. pub mtu_exceeded_uncorroborated: Counter, pub unbound: UnboundSignals, + /// Admitted `PathBroken` signals against coordinates a lookup verified + /// that left them in place, because reports over distinct links had not + /// yet reached the quorum. The lookup still ran. A genuine failure's + /// reports usually arrive over one link, so this is the ordinary outcome + /// of a real broken path as well as of signals forged or reflected + /// through one neighbour. + pub broken_below_quorum: Counter, + /// Verified coordinates demoted to a hint because reports arriving over + /// distinct links reached the quorum. A genuine failure reported from two + /// directions produces this. Forged reports produce it only when they + /// arrive over two different links. + pub broken_demoted: Counter, + /// Admitted `PathBroken` signals that arrived over a link other than the + /// one this node would forward to the destination on. Counted, never + /// refused. The healthy floor is not zero: 4 to 10 percent of genuine + /// reports arrived off the forward link in a loopback measurement, and a + /// real topology with asymmetric paths will see more. + pub broken_link_mismatch: Counter, + /// Admitted `PathBroken` signals whose reporter is this node, the + /// destination, or a node whose known coordinates are no closer to the + /// destination than this node's. Counted, never refused. The healthy + /// floor is not zero, since the reporter's view of the destination can + /// differ from this node's, and most reporters' coordinates are unknown + /// here and are not counted at all. + pub broken_reporter_mismatch: Counter, /// Routing errors this node declined to emit because the authenticated /// link peer that induced them had spent its budget. A rising count is /// either a peer flooding unroutable traffic or a hub relaying more @@ -565,6 +590,10 @@ impl ErrorMetrics { unbound_broken: self.unbound.broken.get(), unbound_mtu: self.unbound.mtu.get(), unbound_forged: self.unbound.forged.get(), + broken_below_quorum: self.broken_below_quorum.get(), + broken_demoted: self.broken_demoted.get(), + broken_link_mismatch: self.broken_link_mismatch.get(), + broken_reporter_mismatch: self.broken_reporter_mismatch.get(), emit_over_peer_budget: self.emit_over_peer_budget.get(), emit_over_dest_interval: self.emit_over_dest_interval.get(), emit_limiter_at_capacity: self.emit_limiter_at_capacity.get(), diff --git a/src/node/mod.rs b/src/node/mod.rs index 588c7c16..482eb57c 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -52,6 +52,7 @@ use crate::proto::fmp::wire::{ build_established_header, prepend_inner_header, }; use crate::proto::fsp::Fsp; +use crate::proto::fsp::quorum::LinkQuorum; use crate::proto::lookup::{Lookup, LookupBackoff, LookupForwardRateLimiter}; use crate::proto::mmp::Mmp; use crate::proto::routing::{self, Router, RoutingErrorRateLimiter}; @@ -600,6 +601,11 @@ pub struct Node { /// not a bound on this one, and one PathBroken drives both responses, so /// a shared limiter would let the coord-warmup arm pay for the release. path_mtu_release_limiter: RoutingErrorRateLimiter, + /// Distinct links that recently delivered a PathBroken naming a + /// destination whose coordinates a lookup verified. A verified entry is + /// demoted only when this reaches its quorum; any number of reports over + /// one link leave it in place while the lookup re-validates it. + broken_quorum: LinkQuorum, // === Peering Homeostasis === /// Owner of the peering-reconciler state relocated off `Node`: the sans-IO @@ -860,6 +866,7 @@ impl Node { path_mtu_release_limiter: RoutingErrorRateLimiter::with_interval_ms( handlers::session::PATH_MTU_RELEASE_MIN_INTERVAL.as_millis() as u64, ), + broken_quorum: LinkQuorum::new(), probes: handlers::probe::ProbeRegistry::new(), lookup: Lookup::new( LookupBackoff::with_params(backoff_base_secs, backoff_max_secs), @@ -1020,6 +1027,7 @@ impl Node { path_mtu_release_limiter: RoutingErrorRateLimiter::with_interval_ms( handlers::session::PATH_MTU_RELEASE_MIN_INTERVAL.as_millis() as u64, ), + broken_quorum: LinkQuorum::new(), probes: handlers::probe::ProbeRegistry::new(), lookup: Lookup::new(LookupBackoff::new(), LookupForwardRateLimiter::new()), discovery_sign_limiter: LookupSignRateLimiter::new(), diff --git a/src/node/stats.rs b/src/node/stats.rs index e4fd7d04..d85fed07 100644 --- a/src/node/stats.rs +++ b/src/node/stats.rs @@ -460,6 +460,10 @@ pub struct ErrorSignalStatsSnapshot { pub unbound_broken: u64, pub unbound_mtu: u64, pub unbound_forged: u64, + pub broken_below_quorum: u64, + pub broken_demoted: u64, + pub broken_link_mismatch: u64, + pub broken_reporter_mismatch: u64, pub emit_over_peer_budget: u64, pub emit_over_dest_interval: u64, pub emit_limiter_at_capacity: u64, diff --git a/src/node/tests/coord_forgery.rs b/src/node/tests/coord_forgery.rs new file mode 100644 index 00000000..bf772fd3 --- /dev/null +++ b/src/node/tests/coord_forgery.rs @@ -0,0 +1,760 @@ +//! Forged routing signals against a destination whose coordinates a lookup +//! verified. +//! +//! The attack these tests describe: a `PathBroken` naming a destination this +//! node has a session with, followed by a forged `SessionSetup` carrying a +//! different position for that destination under the same root. Every packet +//! is driven through `handle_session_datagram`, the entry point a real one +//! reaches, and the observable is the next hop `find_next_hop` picks. + +use super::*; +use crate::node::session::EndToEndState; +use crate::node::tests::spanning_tree::{TestNode, cleanup_nodes, run_tree_test}; +use crate::proto::fsp::SessionSetup; +use crate::proto::link::SessionDatagram; +use crate::proto::routing::PathBroken; +use crate::proto::stp::TreeCoordinate; + +/// Index of the victim in the fixture's node vector. +const V: usize = 0; +/// Index of the peer the destination genuinely sits under. +const P1: usize = 1; +/// Index of the peer the forged coordinates point at. +const P2: usize = 2; + +/// The victim, its two tree peers, and a destination it holds a session +/// with whose real and forged coordinates differ in the peer they hang off. +struct Fixture { + nodes: Vec, + dest: NodeAddr, + p1: NodeAddr, + p2: NodeAddr, + real: TreeCoordinate, + forged: TreeCoordinate, +} + +/// Wall-clock milliseconds, the clock `handle_path_broken` and +/// `find_next_hop` read. A verification stamped on any other clock reads as +/// aged out to the handler. +fn wall_ms() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_millis() as u64 +} + +/// `dest` hung directly under `parent`, sharing its root. +fn coords_under(dest: NodeAddr, parent: &TreeCoordinate) -> TreeCoordinate { + let mut addrs = vec![dest]; + addrs.extend(parent.node_addrs().copied()); + TreeCoordinate::from_addrs(addrs).unwrap() +} + +/// Install the entry `initiate_session` creates for `remote`, which is what +/// makes the admission gate accept a signal naming it. +fn install_initiating(node: &mut Node, remote: &Identity) { + use crate::noise::HandshakeState; + + let handshake = + HandshakeState::new_xk_initiator(node.identity().keypair(), remote.pubkey_full()); + let entry = crate::node::session::SessionEntry::new( + *remote.node_addr(), + remote.pubkey_full(), + EndToEndState::Initiating(handshake), + 1000, + true, + ); + node.sessions.insert(*remote.node_addr(), entry); +} + +/// Build the fixture and assert the two preconditions every test relies on. +async fn fixture() -> Fixture { + let mut nodes = run_tree_test(3, &[(V, P1), (V, P2)], false).await; + let p1 = *nodes[P1].node.node_addr(); + let p2 = *nodes[P2].node.node_addr(); + + let remote = Identity::generate(); + let dest = *remote.node_addr(); + install_initiating(&mut nodes[V].node, &remote); + + let real = coords_under(dest, nodes[P1].node.tree_state().my_coords()); + let forged = coords_under(dest, nodes[P2].node.tree_state().my_coords()); + nodes[V] + .node + .coord_cache_mut() + .insert_verified(dest, real.clone(), wall_ms()); + + let mut fx = Fixture { + nodes, + dest, + p1, + p2, + real, + forged, + }; + assert_eq!( + fx.next_hop(), + Some(fx.p1), + "precondition: the verified coordinates route via P1" + ); + let forged_hop = fx.nodes[V] + .node + .tree_state() + .find_next_hop(&fx.forged, &std::collections::BTreeSet::new()); + assert_eq!( + forged_hop, + Some(fx.p2), + "precondition: the forged coordinates would route via P2, so a \ + successful plant is visible as a flip" + ); + fx +} + +impl Fixture { + /// The next hop the victim picks toward the destination. + fn next_hop(&mut self) -> Option { + let dest = self.dest; + self.nodes[V] + .node + .find_next_hop(&dest) + .map(|p| *p.node_addr()) + } + + /// The victim's cache entry for the destination: its value and whether it + /// is still verified on the handler's clock. + fn entry(&self) -> Option<(TreeCoordinate, bool)> { + self.nodes[V] + .node + .coord_cache() + .get_entry(&self.dest) + .map(|e| (e.coords().clone(), e.is_verified(wall_ms()))) + } + + /// Deliver a `PathBroken` naming the destination, claiming `reporter` as + /// both the datagram source and the body's reporter, arriving over the + /// link to `link_peer`. + async fn path_broken(&mut self, reporter: NodeAddr, link_peer: NodeAddr) { + self.path_broken_from(reporter, reporter, link_peer).await; + } + + /// Deliver a `PathBroken` naming the destination from datagram source + /// `src`, whose body names `reporter`, arriving over the link to + /// `link_peer`. Both identities are the sender's to choose. + async fn path_broken_from(&mut self, src: NodeAddr, reporter: NodeAddr, link_peer: NodeAddr) { + let victim = *self.nodes[V].node.node_addr(); + let payload = PathBroken::new(self.dest, reporter).encode(); + let encoded = SessionDatagram::new(src, victim, payload).encode(); + self.nodes[V] + .node + .handle_session_datagram(&link_peer, &encoded[1..], false) + .await; + } + + /// Deliver a forged `SessionSetup` in transit, claiming to be from the + /// destination and carrying the forged coordinates for it, over the link + /// to P2. + async fn forged_setup(&mut self) { + let payload = SessionSetup::new(self.forged.clone(), self.forged.clone()).encode(); + let encoded = SessionDatagram::new(self.dest, self.dest, payload).encode(); + let p2 = self.p2; + self.nodes[V] + .node + .handle_session_datagram(&p2, &encoded[1..], false) + .await; + } +} + +/// One forged warm against a verified destination changes nothing: the +/// precedence rule refuses the hint, and the route stays on P1. +#[tokio::test] +async fn a_forged_same_root_warm_does_not_move_the_route_to_a_verified_destination() { + let mut fx = fixture().await; + let rejected = fx.nodes[V] + .node + .metrics() + .forwarding + .coord_hint_rejected + .get(); + + fx.forged_setup().await; + + assert_eq!(fx.next_hop(), Some(fx.p1), "a forged warm moved the route"); + assert!( + fx.nodes[V] + .node + .metrics() + .forwarding + .coord_hint_rejected + .get() + > rejected, + "the refused hint should be counted" + ); + cleanup_nodes(&mut fx.nodes).await; +} + +/// One forged `PathBroken` followed by one forged warm: the two-packet attack. +/// The signal must not strip the verification that refuses the warm. +#[tokio::test] +async fn a_forged_path_broken_then_a_forged_warm_does_not_move_the_route_to_a_verified_destination() +{ + let mut fx = fixture().await; + let p2 = fx.p2; + + fx.path_broken(make_node_addr(0xB1), p2).await; + fx.forged_setup().await; + + assert_eq!( + fx.next_hop(), + Some(fx.p1), + "the forged warm moved the route" + ); + assert_eq!( + fx.entry(), + Some((fx.real.clone(), true)), + "the verified entry must survive one PathBroken with its value and \ + its verification" + ); + let errors = &fx.nodes[V].node.metrics().errors; + assert_eq!(errors.broken_below_quorum.get(), 1); + assert_eq!(errors.broken_demoted.get(), 0); + cleanup_nodes(&mut fx.nodes).await; +} + +/// A direct forger on one link inventing two reporters is still one vote: the +/// entry stays verified and the forged warm that follows is refused. Keyed on +/// the body's reporter, the two invented names would have reached the quorum +/// and the warm would have moved the route. +#[tokio::test] +async fn a_forger_inventing_two_reporters_over_one_link_does_not_demote_a_verified_entry() { + let mut fx = fixture().await; + let p2 = fx.p2; + + fx.path_broken(make_node_addr(0xB1), p2).await; + fx.path_broken(make_node_addr(0xB2), p2).await; + fx.forged_setup().await; + + assert_eq!( + fx.entry(), + Some((fx.real.clone(), true)), + "two invented reporters over one link demoted the verified entry" + ); + assert_eq!( + fx.next_hop(), + Some(fx.p1), + "the forged warm moved the route" + ); + let errors = &fx.nodes[V].node.metrics().errors; + assert_eq!(errors.broken_below_quorum.get(), 2); + assert_eq!(errors.broken_demoted.get(), 0); + cleanup_nodes(&mut fx.nodes).await; +} + +/// Two genuinely different reporters whose signals both arrive over one link +/// are one vote too: what counts is the authenticated link, not the +/// reporter, nor the pair of the two. +#[tokio::test] +async fn two_reporters_over_the_same_link_do_not_reach_the_quorum() { + let mut fx = fixture().await; + let (p1, p2) = (fx.p1, fx.p2); + + fx.path_broken(p1, p1).await; + fx.path_broken(p2, p1).await; + + assert_eq!(fx.entry(), Some((fx.real.clone(), true))); + let errors = &fx.nodes[V].node.metrics().errors; + assert_eq!(errors.broken_below_quorum.get(), 2); + assert_eq!(errors.broken_demoted.get(), 0); + cleanup_nodes(&mut fx.nodes).await; +} + +/// Reports arriving over two different links reach the quorum and demote the +/// entry, and the forged warm that follows then moves the route. This is the +/// residual the quorum leaves: forged reports demote the entry once they +/// arrive over two different links, as a genuine failure reported from two +/// directions does. Asserted so that it stays visible. +#[tokio::test] +async fn reports_over_two_different_links_demote_a_verified_entry_and_a_following_warm_then_moves_the_route() + { + let mut fx = fixture().await; + let (p1, p2) = (fx.p1, fx.p2); + + fx.path_broken(make_node_addr(0xB1), p1).await; + fx.path_broken(make_node_addr(0xB2), p2).await; + + assert_eq!( + fx.entry(), + Some((fx.real.clone(), false)), + "a demoted entry keeps its value and loses its verification" + ); + assert_eq!( + fx.nodes[V] + .node + .coord_cache() + .get_entry(&fx.dest) + .map(|e| e.source()), + Some(crate::cache::CoordSource::Hint) + ); + assert_eq!( + fx.next_hop(), + Some(fx.p1), + "demotion alone does not move the route" + ); + let errors = &fx.nodes[V].node.metrics().errors; + assert_eq!(errors.broken_demoted.get(), 1); + assert_eq!(errors.broken_below_quorum.get(), 1); + + fx.forged_setup().await; + + assert_eq!( + fx.next_hop(), + Some(fx.p2), + "once demoted, the entry is a hint and a forged warm replaces it" + ); + cleanup_nodes(&mut fx.nodes).await; +} + +/// A node with a single peer receives every report over one link, so it never +/// demotes by quorum, however many reporters the reports name. What bounds a +/// stale verified entry there is its verification age: one still inside +/// `VERIFIED_TTL_MS` is kept, one past it is removed by the next report like +/// any hint. +#[tokio::test] +async fn on_a_single_link_node_only_the_verification_age_bounds_a_verified_entry() { + use crate::cache::VERIFIED_TTL_MS; + + // V - P1, and V holds a session with a destination beyond P1. + let mut nodes = run_tree_test(2, &[(V, P1)], false).await; + let p1 = *nodes[P1].node.node_addr(); + let victim = *nodes[V].node.node_addr(); + let remote = Identity::generate(); + let dest = *remote.node_addr(); + install_initiating(&mut nodes[V].node, &remote); + let coords = coords_under(dest, nodes[P1].node.tree_state().my_coords()); + + let report = |reporter: NodeAddr| { + let payload = PathBroken::new(dest, reporter).encode(); + SessionDatagram::new(reporter, victim, payload).encode() + }; + let verified = |node: &Node| { + node.coord_cache() + .get_entry(&dest) + .map(|e| e.is_verified(wall_ms())) + }; + + // A cache TTL long enough that expiry never removes the entry, so every + // removal below is the handler's. Ten seconds of margin inside the + // bound, so the test's own run time cannot age the entry out before the + // reports land. + let cache = nodes[V].node.coord_cache_mut(); + cache.set_default_ttl_ms(4 * VERIFIED_TTL_MS); + cache.insert_verified(dest, coords.clone(), wall_ms() - VERIFIED_TTL_MS + 10_000); + for r in 0..5u8 { + let encoded = report(make_node_addr(0xB0 + r)); + nodes[V] + .node + .handle_session_datagram(&p1, &encoded[1..], false) + .await; + } + assert_eq!( + verified(&nodes[V].node), + Some(true), + "reports over one link demoted or removed an entry still inside its \ + verification" + ); + let errors = &nodes[V].node.metrics().errors; + assert_eq!(errors.broken_below_quorum.get(), 5); + assert_eq!(errors.broken_demoted.get(), 0); + + // Past the bound the entry no longer refuses anything, and the next + // report removes it. + nodes[V].node.coord_cache_mut().insert_verified( + dest, + coords, + wall_ms() - VERIFIED_TTL_MS - 1_000, + ); + assert_eq!( + verified(&nodes[V].node), + Some(false), + "precondition: the entry is present and its verification has aged out" + ); + let encoded = report(make_node_addr(0xC0)); + nodes[V] + .node + .handle_session_datagram(&p1, &encoded[1..], false) + .await; + assert_eq!( + verified(&nodes[V].node), + None, + "an aged-out entry is removed" + ); + cleanup_nodes(&mut nodes).await; +} + +/// One reporter repeating itself over one link is one vote: the entry stays +/// verified and the forged warm is still refused. +#[tokio::test] +async fn one_reporter_repeating_a_path_broken_does_not_reach_the_quorum() { + let mut fx = fixture().await; + let p2 = fx.p2; + + fx.path_broken(make_node_addr(0xB1), p2).await; + fx.path_broken(make_node_addr(0xB1), p2).await; + fx.forged_setup().await; + + assert_eq!(fx.next_hop(), Some(fx.p1)); + assert_eq!(fx.entry(), Some((fx.real.clone(), true))); + let errors = &fx.nodes[V].node.metrics().errors; + assert_eq!(errors.broken_below_quorum.get(), 2); + assert_eq!(errors.broken_demoted.get(), 0); + cleanup_nodes(&mut fx.nodes).await; +} + +/// A verification that has aged out protects nothing, so the entry is +/// removed as a hint would be. +#[tokio::test] +async fn a_path_broken_removes_a_verified_entry_whose_verification_has_aged_out() { + let mut fx = fixture().await; + let dest = fx.dest; + let real = fx.real.clone(); + let p1 = fx.p1; + let cache = fx.nodes[V].node.coord_cache_mut(); + // A TTL long enough that the entry is still live when its verification + // is not, so the removal below is the handler's and not expiry's. + cache.set_default_ttl_ms(4 * crate::cache::VERIFIED_TTL_MS); + cache.insert_verified( + dest, + real, + wall_ms() - crate::cache::VERIFIED_TTL_MS - 1_000, + ); + assert_eq!( + fx.entry().map(|(_, verified)| verified), + Some(false), + "precondition: the entry is present and no longer verified" + ); + + fx.path_broken(make_node_addr(0xB1), p1).await; + + assert!(fx.entry().is_none(), "an unverified entry is removed"); + let errors = &fx.nodes[V].node.metrics().errors; + assert_eq!(errors.broken_below_quorum.get(), 0); + assert_eq!(errors.broken_demoted.get(), 0); + cleanup_nodes(&mut fx.nodes).await; +} + +/// When the path-MTU release fires, a kept entry forgets the MTU the lookup +/// stored with it, so it does not go on displaying a released value. +#[tokio::test] +async fn a_kept_entry_forgets_its_path_mtu_when_the_release_fires() { + let mut fx = fixture().await; + let dest = fx.dest; + let real = fx.real.clone(); + let p2 = fx.p2; + fx.nodes[V] + .node + .coord_cache_mut() + .insert_verified_with_path_mtu(dest, real, wall_ms(), 1200); + + fx.path_broken(make_node_addr(0xB1), p2).await; + + let entry = fx.nodes[V].node.coord_cache().get_entry(&dest).unwrap(); + assert!( + entry.is_verified(wall_ms()), + "precondition: the entry was kept" + ); + assert_eq!(entry.path_mtu(), None); + cleanup_nodes(&mut fx.nodes).await; +} + +/// When the path-MTU release is rate limited, a kept entry keeps its MTU, as +/// the path-MTU map does: the two describe the same path and must not +/// disagree. +#[tokio::test] +async fn a_kept_entry_keeps_its_path_mtu_when_the_release_is_rate_limited() { + let mut fx = fixture().await; + let dest = fx.dest; + let real = fx.real.clone(); + let p2 = fx.p2; + let fips = crate::FipsAddress::from_node_addr(&dest); + + // The first report spends the release budget for this destination. + fx.path_broken(make_node_addr(0xB1), p2).await; + + // A value learned again since, in both stores. + fx.nodes[V] + .node + .coord_cache_mut() + .insert_verified_with_path_mtu(dest, real, wall_ms(), 1200); + fx.nodes[V].node.path_mtu_lookup_insert(fips, 1200); + + // A second report inside the release interval (same link, so it stays + // below the quorum and the entry is kept). + fx.path_broken(make_node_addr(0xB2), p2).await; + + let node = &fx.nodes[V].node; + let entry = node.coord_cache().get_entry(&dest).unwrap(); + assert!( + entry.is_verified(wall_ms()), + "precondition: the entry was kept" + ); + assert_eq!( + node.path_mtu_lookup_get(&fips), + Some(1200), + "precondition: the release was rate limited" + ); + assert_eq!( + entry.path_mtu(), + Some(1200), + "a rate-limited release cleared the kept entry's MTU but not the map's" + ); + cleanup_nodes(&mut fx.nodes).await; +} + +/// The sum of the lookup counters `maybe_initiate_lookup` moves on every +/// outcome, so a test can see that it ran whatever it decided. +fn lookup_attempts(node: &Node) -> u64 { + let l = &node.metrics().lookup; + l.req_initiated.get() + + l.req_bloom_miss.get() + + l.req_deduplicated.get() + + l.req_backoff_suppressed.get() +} + +/// An admitted `PathBroken` re-validates the destination by lookup even when +/// its identity is not cached, which is the fixture's state: the session entry +/// alone does not register the identity. +#[tokio::test] +async fn an_admitted_path_broken_starts_a_lookup_without_a_cached_identity() { + let mut fx = fixture().await; + let dest = fx.dest; + assert!( + !fx.nodes[V].node.has_cached_identity(&dest), + "precondition: the destination's identity is not cached" + ); + let before = lookup_attempts(&fx.nodes[V].node); + let p1 = fx.p1; + + fx.path_broken(make_node_addr(0xB1), p1).await; + + assert!( + lookup_attempts(&fx.nodes[V].node) > before, + "the PathBroken should have run the lookup path" + ); + cleanup_nodes(&mut fx.nodes).await; +} + +/// The healthy path for keeping a verified entry below quorum: a genuine +/// failure reported once is still recovered, because the lookup the report +/// starts replaces the stale value with the destination's real position. +#[tokio::test] +async fn a_genuine_path_broken_still_replaces_stale_verified_coordinates_by_lookup() { + // V - P1 - D, with D a real node V has a session with, and a second peer + // P2 on its own, so a later report can arrive over a different link. + let mut nodes = run_tree_test(4, &[(0, 1), (1, 2), (0, 3)], false).await; + let p1 = *nodes[1].node.node_addr(); + let p2 = *nodes[3].node.node_addr(); + let dest = *nodes[2].node.node_addr(); + let dest_pubkey = nodes[2].node.identity().pubkey_full(); + let real: Vec = nodes[2] + .node + .tree_state() + .my_coords() + .node_addrs() + .copied() + .collect(); + nodes[0].node.register_identity(dest, dest_pubkey); + let handshake = crate::noise::HandshakeState::new_xk_initiator( + nodes[0].node.identity().keypair(), + dest_pubkey, + ); + nodes[0].node.sessions.insert( + dest, + crate::node::session::SessionEntry::new( + dest, + dest_pubkey, + EndToEndState::Initiating(handshake), + 1000, + true, + ), + ); + + // D "moved": V holds a verified position for it that is no longer true. + let stale = coords_under(dest, nodes[0].node.tree_state().my_coords()); + assert_ne!(stale.node_addrs().copied().collect::>(), real); + nodes[0] + .node + .coord_cache_mut() + .insert_verified(dest, stale, wall_ms()); + + let lookup = &nodes[0].node.metrics().lookup; + let (initiated, accepted) = (lookup.req_initiated.get(), lookup.resp_accepted.get()); + + let victim = *nodes[0].node.node_addr(); + let payload = PathBroken::new(dest, p1).encode(); + let encoded = SessionDatagram::new(p1, victim, payload).encode(); + nodes[0] + .node + .handle_session_datagram(&p1, &encoded[1..], false) + .await; + assert_eq!(nodes[0].node.metrics().errors.broken_below_quorum.get(), 1); + + for _ in 0..10 { + tokio::time::sleep(std::time::Duration::from_millis(50)).await; + crate::node::tests::spanning_tree::process_available_packets(&mut nodes).await; + } + + let lookup = &nodes[0].node.metrics().lookup; + assert!( + lookup.req_initiated.get() > initiated, + "the report must start a lookup" + ); + assert!( + lookup.resp_accepted.get() > accepted, + "the lookup must be answered" + ); + let entry = nodes[0].node.coord_cache().get_entry(&dest).unwrap(); + assert_eq!( + entry.coords().node_addrs().copied().collect::>(), + real, + "the lookup must replace the stale position with the real one" + ); + assert!(entry.is_verified(wall_ms())); + + // The lookup verified the destination again, so the report about the old + // path no longer counts: a report over a second link now starts a new + // quorum rather than completing the old one. + let payload = PathBroken::new(dest, make_node_addr(0xB2)).encode(); + let encoded = SessionDatagram::new(make_node_addr(0xB2), victim, payload).encode(); + nodes[0] + .node + .handle_session_datagram(&p2, &encoded[1..], false) + .await; + let errors = &nodes[0].node.metrics().errors; + assert_eq!( + errors.broken_demoted.get(), + 0, + "a report from before the re-verification combined with one after it" + ); + assert_eq!(errors.broken_below_quorum.get(), 2); + cleanup_nodes(&mut nodes).await; +} + +/// A PathBroken arriving off the forward link is counted and still acted on +/// in full: the advisory check refuses nothing. +#[tokio::test] +async fn a_path_broken_off_the_forward_link_is_counted_and_still_acted_on() { + let mut fx = fixture().await; + let dest = fx.dest; + let p2 = fx.p2; + fx.nodes[V] + .node + .sessions + .get_mut(&dest) + .unwrap() + .set_coords_warmup_remaining(0); + let lookups = lookup_attempts(&fx.nodes[V].node); + + fx.path_broken(make_node_addr(0xB1), p2).await; + + let node = &fx.nodes[V].node; + assert_eq!(node.metrics().errors.broken_link_mismatch.get(), 1); + assert_eq!( + node.sessions.get(&dest).unwrap().coords_warmup_remaining(), + node.config().node.session.coords_warmup_packets, + "the warmup counter is reset whatever the link" + ); + assert!( + node.config().node.session.coords_warmup_packets > 0, + "precondition: a reset is distinguishable from the zero set above" + ); + assert!( + lookup_attempts(node) > lookups, + "the lookup runs whatever the link" + ); + assert_eq!(node.metrics().errors.broken_below_quorum.get(), 1); + cleanup_nodes(&mut fx.nodes).await; +} + +/// The same report over the forward link is not counted. +#[tokio::test] +async fn a_path_broken_over_the_forward_link_is_not_counted_as_a_mismatch() { + let mut fx = fixture().await; + let p1 = fx.p1; + + fx.path_broken(make_node_addr(0xB1), p1).await; + + let errors = &fx.nodes[V].node.metrics().errors; + assert_eq!(errors.broken_link_mismatch.get(), 0); + assert_eq!( + errors.broken_below_quorum.get(), + 1, + "the signal was admitted and handled, so the zero above is a real zero" + ); + cleanup_nodes(&mut fx.nodes).await; +} + +/// The reporter-mismatch count after one PathBroken naming the fixture's +/// destination, reported by the address `pick` chooses, over the forward link. +async fn reporter_mismatches(pick: impl FnOnce(&Fixture) -> NodeAddr) -> u64 { + mismatches_from(|fx| { + let reporter = pick(fx); + (reporter, reporter) + }) + .await +} + +/// The reporter-mismatch count after one PathBroken naming the fixture's +/// destination over the forward link, with the datagram source and the +/// body's reporter that `pick` returns, in that order. +async fn mismatches_from(pick: impl FnOnce(&Fixture) -> (NodeAddr, NodeAddr)) -> u64 { + let mut fx = fixture().await; + let (src, reporter) = pick(&fx); + let p1 = fx.p1; + fx.path_broken_from(src, reporter, p1).await; + let errors = &fx.nodes[V].node.metrics().errors; + assert_eq!( + errors.broken_below_quorum.get(), + 1, + "precondition: the signal reached the handler's verified-entry path, \ + so the count below observed it" + ); + let n = errors.broken_reporter_mismatch.get(); + cleanup_nodes(&mut fx.nodes).await; + n +} + +#[tokio::test] +async fn a_reporter_farther_from_the_destination_than_this_node_is_counted() { + // P2 is a direct peer on the far side of this node from P1, under which + // the destination sits. + assert_eq!(reporter_mismatches(|fx| fx.p2).await, 1); +} + +#[tokio::test] +async fn a_reporter_closer_to_the_destination_than_this_node_is_not_counted() { + assert_eq!(reporter_mismatches(|fx| fx.p1).await, 0); +} + +#[tokio::test] +async fn a_reporter_with_unknown_coordinates_is_not_counted() { + assert_eq!(reporter_mismatches(|_| make_node_addr(0xB1)).await, 0); +} + +#[tokio::test] +async fn this_node_named_as_the_reporter_is_counted() { + assert_eq!( + reporter_mismatches(|fx| *fx.nodes[V].node.node_addr()).await, + 1 + ); +} + +/// The destination named as its own reporter cannot be reporting a broken +/// path to itself. The datagram source differs from the reporter here, +/// because a datagram claiming to come from the destination is refused by the +/// admission gate before this check is reached. +#[tokio::test] +async fn the_destination_named_as_the_reporter_is_counted() { + assert_eq!( + mismatches_from(|fx| (make_node_addr(0xB1), fx.dest)).await, + 1 + ); +} diff --git a/src/node/tests/mod.rs b/src/node/tests/mod.rs index b31c75aa..b57b0d33 100644 --- a/src/node/tests/mod.rs +++ b/src/node/tests/mod.rs @@ -13,6 +13,7 @@ mod bloom_poison; mod bootstrap; mod connected_udp; mod control; +mod coord_forgery; mod decrypt_failure; mod disconnect; mod discovery; diff --git a/src/node/tests/session.rs b/src/node/tests/session.rs index aba0d69a..8bb70718 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -4306,7 +4306,9 @@ async fn test_path_broken_releases_path_mtu_lookup_entry() { assertion below observes nothing" ); - tn.node.handle_path_broken(&reporter, inner).await; + tn.node + .handle_path_broken(&reporter, &reporter, inner) + .await; assert_eq!( tn.node.path_mtu_lookup_get(&dest_fips), @@ -4361,7 +4363,9 @@ async fn test_path_broken_resets_the_session_source_path_mtu() { assertion below observes nothing" ); - tn.node.handle_path_broken(&reporter, inner).await; + tn.node + .handle_path_broken(&reporter, &reporter, inner) + .await; assert_eq!( tn.node @@ -4699,7 +4703,8 @@ async fn test_path_broken_naming_a_dest_with_no_session_does_not_flush_cached_co let _ = node.coord_cache_mut().insert(dest, coords, 1000); let encoded = PathBroken::new(dest, reporter).encode(); - node.handle_path_broken(&reporter, &encoded[5..]).await; + node.handle_path_broken(&reporter, &reporter, &encoded[5..]) + .await; assert!( node.coord_cache().get(&dest, 1000).is_some(), @@ -4746,7 +4751,8 @@ async fn test_path_broken_naming_a_dest_whose_entry_is_an_unauthenticated_respon install_halfopen(&mut node, dest); let encoded = PathBroken::new(dest, reporter).encode(); - node.handle_path_broken(&reporter, &encoded[5..]).await; + node.handle_path_broken(&reporter, &reporter, &encoded[5..]) + .await; assert!( node.coord_cache().get(&dest, 1000).is_some(), @@ -4781,7 +4787,8 @@ async fn test_path_broken_for_a_session_we_initiated_still_flushes_cached_coords let _ = node.coord_cache_mut().insert(dest, coords, 1000); let encoded = PathBroken::new(dest, reporter).encode(); - node.handle_path_broken(&reporter, &encoded[5..]).await; + node.handle_path_broken(&reporter, &reporter, &encoded[5..]) + .await; assert!( node.coord_cache().get(&dest, 1000).is_none(), @@ -6993,7 +7000,7 @@ async fn deliver_path_broken( let encoded = PathBroken::new(dest, reporter).encode(); nodes[at] .node - .handle_path_broken(&reporter, &encoded[5..]) + .handle_path_broken(&reporter, &reporter, &encoded[5..]) .await; } @@ -8994,7 +9001,7 @@ async fn a_path_broken_flood_releases_the_stored_path_mtu_only_once_per_interval let inner = &encoded[5..]; node.path_mtu_lookup_insert(dest_fips, 700); - node.handle_path_broken(&reporter, inner).await; + node.handle_path_broken(&reporter, &reporter, inner).await; assert_eq!( node.path_mtu_lookup_get(&dest_fips), None, @@ -9002,7 +9009,7 @@ async fn a_path_broken_flood_releases_the_stored_path_mtu_only_once_per_interval ); node.path_mtu_lookup_insert(dest_fips, 700); - node.handle_path_broken(&reporter, inner).await; + node.handle_path_broken(&reporter, &reporter, inner).await; assert_eq!( node.path_mtu_lookup_get(&dest_fips), Some(700), diff --git a/src/proto/fsp/core.rs b/src/proto/fsp/core.rs index 33fc8b62..f45526a0 100644 --- a/src/proto/fsp/core.rs +++ b/src/proto/fsp/core.rs @@ -15,6 +15,7 @@ //! the plain-data [`DecryptSlot`] mirror before it reaches [`Fsp::classify_epoch`]. use super::limits::FSP_CUTOVER_DELAY_MS; +use super::quorum::QuorumVerdict; use crate::proto::stp::TreeCoordinate; use crate::{FipsAddress, NodeAddr}; @@ -95,12 +96,17 @@ pub(crate) enum FspAction { /// Invalidate the shared cached coordinates for `addr` /// (`coord_cache.remove`). InvalidateCoords { addr: NodeAddr }, + /// Mark the shared cached coordinates for `addr` as an unverified hint, + /// keeping the value (`coord_cache.demote`). The entry goes on routing, + /// but no longer refuses a hint that replaces it. + DemoteCoords { addr: NodeAddr }, /// Write `mtu` into the shared `FipsAddress`-keyed path-MTU lookup, keeping /// the tighter of existing-or-new (the shell applies the write under the /// `path_mtu_lookup` guard). TightenPathMtuLookup { fips_addr: FipsAddress, mtu: u16 }, - /// Trigger discovery toward `dest` (`maybe_initiate_lookup`); emitted only - /// when the target's identity is cached. + /// Trigger discovery toward `dest` (`maybe_initiate_lookup`). The + /// `CoordsRequired` decision emits it only when the target's identity is + /// cached; the `PathBroken` decision emits it always. InitiateLookup { dest: NodeAddr }, } @@ -445,20 +451,35 @@ impl Fsp { } } - /// Decide the reaction to a `PathBroken` signal: unconditionally invalidate - /// the stale cached coordinates for `dest`, then (only when the identity is - /// cached) trigger re-discovery. Order is invalidate-then-lookup, matching - /// the pre-refactor handler. The warmup send and counter reset stay shell. + /// Decide the reaction to an admitted `PathBroken` signal for `dest`. + /// + /// The signal is unauthenticated, so it may not discard coordinates a + /// lookup verified on its own say-so. `verified` is whether the cached + /// entry for `dest` is verified and still within its verification window; + /// `quorum` is what this report did to the destination's link quorum + /// (ignored when the entry is not verified). + /// + /// - Not verified: remove the entry, as a hint can be replaced by any warm + /// anyway, then look it up. + /// - Verified, quorum reached: demote the entry to a hint, keeping the + /// value, then look it up. + /// - Verified, below quorum: leave the entry alone and look it up. A + /// successful lookup replaces the value whatever it was. + /// + /// The lookup is emitted in every case. The warmup send, path-MTU release + /// and warmup-counter reset stay shell-side. pub(crate) fn plan_path_broken( &self, dest: NodeAddr, - has_cached_identity: bool, + verified: bool, + quorum: QuorumVerdict, ) -> Vec { - let mut actions = vec![FspAction::InvalidateCoords { addr: dest }]; - if has_cached_identity { - actions.push(FspAction::InitiateLookup { dest }); + let lookup = FspAction::InitiateLookup { dest }; + match (verified, quorum) { + (false, _) => vec![FspAction::InvalidateCoords { addr: dest }, lookup], + (true, QuorumVerdict::Reached) => vec![FspAction::DemoteCoords { addr: dest }, lookup], + (true, QuorumVerdict::Below { .. }) => vec![lookup], } - actions } /// Decide whether a path-MTU update should tighten the shared lookup: emit diff --git a/src/proto/fsp/mod.rs b/src/proto/fsp/mod.rs index 6ebb179c..7c9b37f5 100644 --- a/src/proto/fsp/mod.rs +++ b/src/proto/fsp/mod.rs @@ -17,11 +17,14 @@ //! `classify_epoch`, the initiation tie-break, and the pure MTU-clamp / //! bounded-queue / ECN transforms. No clock/crypto/I/O/tracing. //! - `limits.rs` — the session-rekey timing constants. +//! - `quorum.rs` — the distinct-link quorum that gates demoting verified +//! coordinates on `PathBroken`. Clock injected. //! - `wire.rs` — the FSP session wire codec and message types. Clock-free, //! crypto-free. pub(crate) mod core; pub(crate) mod limits; +pub(crate) mod quorum; pub(crate) mod wire; #[cfg(test)] diff --git a/src/proto/fsp/quorum.rs b/src/proto/fsp/quorum.rs new file mode 100644 index 00000000..b7a08e8c --- /dev/null +++ b/src/proto/fsp/quorum.rs @@ -0,0 +1,126 @@ +//! Distinct-link quorum for `PathBroken` signals. +//! +//! A `PathBroken` carries no end-to-end authentication, so one report is not +//! enough evidence to discard coordinates a lookup verified. This module +//! counts, per destination, the distinct links that delivered a report naming +//! it within a window, and says when enough of them have. +//! +//! The vote is keyed on the authenticated link peer the datagram arrived over, +//! not on the reporter the body names: the reporter is plaintext the sender +//! chooses, so a sender on one link could invent as many reporters as the +//! quorum needs. Every datagram from one neighbour, whether its own or relayed +//! through it, arrives over the same link and is one vote. +//! +//! Sans-IO: the caller passes the time in, and nothing here logs or counts. + +use std::collections::BTreeMap; + +use crate::NodeAddr; + +/// Distinct links that must deliver a report naming a destination within +/// [`QUORUM_WINDOW_MS`] before its verified coordinates are demoted. +/// +/// Two is the smallest value that one neighbour cannot reach alone, whether it +/// forges the reports or relays a reflection. No measurement has counted how +/// many links a genuine failure's reports arrive over; a larger value only +/// keeps a genuinely stale verified entry verified for longer, while the +/// lookup every report starts replaces it anyway. +/// +/// A node whose reports all arrive over one link, such as a leaf with a +/// single peer, never reaches this. A stale verified entry there lasts until +/// a lookup that a report started answers and replaces it, or at most until +/// its verification ages out after [`crate::cache::VERIFIED_TTL_MS`]; from +/// then on it no longer refuses a hint, and the next report removes it. +pub(crate) const QUORUM_LINKS: usize = 2; + +/// Window within which distinct links count toward one quorum, in +/// milliseconds. +/// +/// The sum of the default lookup attempt schedule (1 + 2 + 4 + 8 s), so the +/// reports that arrive during one re-validation cycle count together. An +/// anchor, not a measurement. +pub(crate) const QUORUM_WINDOW_MS: u64 = 15_000; + +/// What one report did to its destination's quorum. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum QuorumVerdict { + /// Not enough distinct links yet; `distinct` counts this one. + Below { distinct: usize }, + /// This report completed the quorum. The destination's record is cleared, + /// so the next quorum starts from nothing. + Reached, +} + +/// Distinct links seen per destination, each with the time it was first seen +/// inside the current window. +#[derive(Debug, Default)] +pub(crate) struct LinkQuorum { + seen: BTreeMap>, + last_sweep_ms: u64, +} + +impl LinkQuorum { + /// An empty quorum tracker. + pub(crate) fn new() -> Self { + Self::default() + } + + /// Record that a report naming `dest` arrived over the link to + /// `link_peer` at `now_ms`. + /// + /// A link already counted for `dest` inside the window does not count + /// again, and its first-seen time is not refreshed, so one neighbour + /// repeating itself can neither reach the quorum nor keep a record alive. + pub(crate) fn record( + &mut self, + dest: NodeAddr, + link_peer: NodeAddr, + now_ms: u64, + ) -> QuorumVerdict { + self.sweep(now_ms); + let seen = self.seen.entry(dest).or_default(); + seen.retain(|&(_, at)| live(at, now_ms)); + if !seen.iter().any(|&(l, _)| l == link_peer) { + seen.push((link_peer, now_ms)); + } + let distinct = seen.len(); + if distinct >= QUORUM_LINKS { + self.seen.remove(&dest); + QuorumVerdict::Reached + } else { + QuorumVerdict::Below { distinct } + } + } + + /// Forget every report about `dest`. + /// + /// Called when a lookup verifies `dest` again: reports about the path the + /// fresh value replaced are not evidence against it. + pub(crate) fn clear(&mut self, dest: &NodeAddr) { + self.seen.remove(dest); + } + + /// Number of destinations with a record. + #[cfg(test)] + pub(crate) fn len(&self) -> usize { + self.seen.len() + } + + /// Drop destinations with no live link, at most once per window, so + /// the map holds only destinations reported within the last window or so. + fn sweep(&mut self, now_ms: u64) { + if now_ms.saturating_sub(self.last_sweep_ms) < QUORUM_WINDOW_MS { + return; + } + self.last_sweep_ms = now_ms; + self.seen.retain(|_, seen| { + seen.retain(|&(_, at)| live(at, now_ms)); + !seen.is_empty() + }); + } +} + +/// Whether a report first seen at `at` still counts at `now_ms`. +fn live(at: u64, now_ms: u64) -> bool { + now_ms.saturating_sub(at) <= QUORUM_WINDOW_MS +} diff --git a/src/proto/fsp/tests/core.rs b/src/proto/fsp/tests/core.rs index 0e80f3af..a4b0150d 100644 --- a/src/proto/fsp/tests/core.rs +++ b/src/proto/fsp/tests/core.rs @@ -737,22 +737,43 @@ fn plan_coords_required_lookup_gated_on_identity() { } #[test] -fn plan_path_broken_invalidates_then_lookups() { +fn plan_path_broken_removes_an_unverified_entry_then_looks_it_up() { + use crate::proto::fsp::quorum::QuorumVerdict; + let fsp = Fsp::new(); + let dest = make_node_addr(6); + let expected = vec![ + FspAction::InvalidateCoords { addr: dest }, + FspAction::InitiateLookup { dest }, + ]; + // The quorum is irrelevant to a hint: removed either way. + for quorum in [QuorumVerdict::Below { distinct: 0 }, QuorumVerdict::Reached] { + assert_eq!(fsp.plan_path_broken(dest, false, quorum), expected); + } +} + +#[test] +fn plan_path_broken_keeps_a_verified_entry_below_quorum_and_looks_it_up() { + use crate::proto::fsp::quorum::QuorumVerdict; let fsp = Fsp::new(); let dest = make_node_addr(6); - // Cached identity: invalidate, then lookup — in that order. assert_eq!( - fsp.plan_path_broken(dest, true), + fsp.plan_path_broken(dest, true, QuorumVerdict::Below { distinct: 1 }), + vec![FspAction::InitiateLookup { dest }] + ); +} + +#[test] +fn plan_path_broken_demotes_a_verified_entry_at_quorum_then_looks_it_up() { + use crate::proto::fsp::quorum::QuorumVerdict; + let fsp = Fsp::new(); + let dest = make_node_addr(6); + assert_eq!( + fsp.plan_path_broken(dest, true, QuorumVerdict::Reached), vec![ - FspAction::InvalidateCoords { addr: dest }, + FspAction::DemoteCoords { addr: dest }, FspAction::InitiateLookup { dest }, ] ); - // No cached identity: invalidate only (still unconditional). - assert_eq!( - fsp.plan_path_broken(dest, false), - vec![FspAction::InvalidateCoords { addr: dest }] - ); } #[test] diff --git a/src/proto/fsp/tests/mod.rs b/src/proto/fsp/tests/mod.rs index e0f8f26e..51820b2a 100644 --- a/src/proto/fsp/tests/mod.rs +++ b/src/proto/fsp/tests/mod.rs @@ -1,2 +1,3 @@ mod core; +mod quorum; mod wire; diff --git a/src/proto/fsp/tests/quorum.rs b/src/proto/fsp/tests/quorum.rs new file mode 100644 index 00000000..2e153932 --- /dev/null +++ b/src/proto/fsp/tests/quorum.rs @@ -0,0 +1,118 @@ +//! Unit tests for the `PathBroken` link quorum. + +use crate::NodeAddr; +use crate::proto::fsp::quorum::{LinkQuorum, QUORUM_LINKS, QUORUM_WINDOW_MS, QuorumVerdict}; + +fn addr(v: u8) -> NodeAddr { + let mut bytes = [0u8; 16]; + bytes[0] = v; + NodeAddr::from_bytes(bytes) +} + +const T0: u64 = 1_000_000; + +#[test] +fn the_quorum_needs_two_links() { + // The tests below are written for this value; a change to it should be a + // deliberate edit here too. + assert_eq!(QUORUM_LINKS, 2); +} + +#[test] +fn one_report_stays_below_the_quorum() { + let mut q = LinkQuorum::new(); + assert_eq!( + q.record(addr(1), addr(0xA), T0), + QuorumVerdict::Below { distinct: 1 } + ); +} + +#[test] +fn two_distinct_links_inside_the_window_reach_the_quorum() { + let mut q = LinkQuorum::new(); + let _ = q.record(addr(1), addr(0xA), T0); + assert_eq!( + q.record(addr(1), addr(0xB), T0 + QUORUM_WINDOW_MS), + QuorumVerdict::Reached + ); +} + +#[test] +fn the_same_link_repeated_never_reaches_the_quorum() { + let mut q = LinkQuorum::new(); + for i in 0..10 { + assert_eq!( + q.record(addr(1), addr(0xA), T0 + i), + QuorumVerdict::Below { distinct: 1 }, + "repeat {i} counted again" + ); + } +} + +#[test] +fn links_further_apart_than_the_window_do_not_combine() { + let mut q = LinkQuorum::new(); + let _ = q.record(addr(1), addr(0xA), T0); + assert_eq!( + q.record(addr(1), addr(0xB), T0 + QUORUM_WINDOW_MS + 1), + QuorumVerdict::Below { distinct: 1 } + ); +} + +#[test] +fn a_repeat_does_not_refresh_its_links_first_seen_time() { + let mut q = LinkQuorum::new(); + let _ = q.record(addr(1), addr(0xA), T0); + let _ = q.record(addr(1), addr(0xA), T0 + QUORUM_WINDOW_MS); + // Had the repeat refreshed A's time, this would combine with it. + assert_eq!( + q.record(addr(1), addr(0xB), T0 + QUORUM_WINDOW_MS + 1), + QuorumVerdict::Below { distinct: 1 } + ); +} + +#[test] +fn reaching_the_quorum_clears_the_destination() { + let mut q = LinkQuorum::new(); + let _ = q.record(addr(1), addr(0xA), T0); + assert_eq!(q.record(addr(1), addr(0xB), T0 + 1), QuorumVerdict::Reached); + assert_eq!(q.len(), 0); + assert_eq!( + q.record(addr(1), addr(0xA), T0 + 2), + QuorumVerdict::Below { distinct: 1 }, + "the next demotion must need a fresh quorum" + ); +} + +#[test] +fn clear_forgets_the_reports_about_a_destination() { + let mut q = LinkQuorum::new(); + let _ = q.record(addr(1), addr(0xA), T0); + q.clear(&addr(1)); + assert_eq!( + q.record(addr(1), addr(0xB), T0 + 1), + QuorumVerdict::Below { distinct: 1 } + ); +} + +#[test] +fn destinations_are_counted_separately() { + let mut q = LinkQuorum::new(); + let _ = q.record(addr(1), addr(0xA), T0); + assert_eq!( + q.record(addr(2), addr(0xB), T0), + QuorumVerdict::Below { distinct: 1 } + ); +} + +#[test] +fn the_sweep_drops_destinations_with_no_live_link() { + let mut q = LinkQuorum::new(); + for d in 0..100 { + let _ = q.record(addr(d), addr(0xA), T0); + } + assert_eq!(q.len(), 100); + // One report a window and a millisecond later sweeps the stale records. + let _ = q.record(addr(200), addr(0xA), T0 + QUORUM_WINDOW_MS + 1); + assert_eq!(q.len(), 1); +} From 90446074f567bc31d26cabea84d91029fde12042 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 1 Oct 2026 19:42:19 +0000 Subject: [PATCH 04/23] Cache the destination's identity from its session before a PathBroken lookup The lookup a PathBroken starts now runs whether or not the destination's identity is cached, but its answer is verified against the cached key. With the identity evicted, the answer was discarded unverified, the lookup ran to its timeout, and packets queued for the destination were dropped as unreachable. An admitted PathBroken names a destination this node holds a session with, and the session holds its key: the one this node initiated to, or the one the responder handshake authenticated. Cache it from there before the lookup when the identity cache has none. --- docs/design/fips-mesh-operation.md | 4 +- src/node/handlers/session.rs | 18 ++++++ src/node/tests/coord_forgery.rs | 100 +++++++++++++++++++++++++++++ 3 files changed, 120 insertions(+), 2 deletions(-) diff --git a/docs/design/fips-mesh-operation.md b/docs/design/fips-mesh-operation.md index 91ccd88e..c38d1581 100644 --- a/docs/design/fips-mesh-operation.md +++ b/docs/design/fips-mesh-operation.md @@ -384,8 +384,8 @@ source. whose signals all arrive over one link never demotes this way; its verified coordinates last until discovery replaces them or their verification ages out after 300 seconds. -3. Initiate discovery for the destination, whether or not its identity is - cached +3. Initiate discovery for the destination. If its identity is not cached, + cache it first from the session's key, so the response can be verified 4. Reset CP warmup counter The source also counts, without refusing anything, a PathBroken that diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index 2b897aab..4d4c1f7b 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -2024,6 +2024,7 @@ impl Node { "PathBroken quorum reached; demoted verified coordinates to a hint"); } FspAction::InitiateLookup { dest } => { + self.cache_session_identity(&dest); self.maybe_initiate_lookup(&dest).await; } _ => {} @@ -2071,6 +2072,23 @@ impl Node { } } + /// Cache the identity of `dest` from its session entry if the identity + /// cache has none. + /// + /// A lookup's answer is verified against the target's cached key, so a + /// lookup for a destination whose identity has been evicted runs to its + /// timeout and drops the packets queued for it as unreachable. A session + /// already holds the key: the one this node initiated to, or the one the + /// responder handshake authenticated. + fn cache_session_identity(&mut self, dest: &NodeAddr) { + if self.has_cached_identity(dest) { + return; + } + if let Some(pubkey) = self.sessions.get(dest).map(|e| *e.remote_pubkey()) { + self.register_identity(*dest, pubkey); + } + } + /// Count, without acting on, the two ways an admitted PathBroken can /// disagree with this node's own view of the path. /// diff --git a/src/node/tests/coord_forgery.rs b/src/node/tests/coord_forgery.rs index bf772fd3..879eb20e 100644 --- a/src/node/tests/coord_forgery.rs +++ b/src/node/tests/coord_forgery.rs @@ -542,6 +542,106 @@ async fn an_admitted_path_broken_starts_a_lookup_without_a_cached_identity() { cleanup_nodes(&mut fx.nodes).await; } +/// The lookup a PathBroken starts must be answerable when the destination's +/// identity is not cached. The session already holds the destination's key, +/// so the lookup's answer can be verified with it: the position is learned, +/// the lookup does not run to its timeout, and packets queued for the +/// destination are not dropped as unreachable when it would have. +#[tokio::test] +async fn a_path_broken_without_a_cached_identity_starts_a_lookup_this_node_can_verify() { + // V - P1 - D, with D a real node V has a session with but whose identity + // V has not cached. + let mut nodes = run_tree_test(3, &[(0, 1), (1, 2)], false).await; + let p1 = *nodes[1].node.node_addr(); + let dest = *nodes[2].node.node_addr(); + let dest_pubkey = nodes[2].node.identity().pubkey_full(); + let real: Vec = nodes[2] + .node + .tree_state() + .my_coords() + .node_addrs() + .copied() + .collect(); + let handshake = crate::noise::HandshakeState::new_xk_initiator( + nodes[0].node.identity().keypair(), + dest_pubkey, + ); + nodes[0].node.sessions.insert( + dest, + crate::node::session::SessionEntry::new( + dest, + dest_pubkey, + EndToEndState::Initiating(handshake), + 1000, + true, + ), + ); + assert!( + !nodes[0].node.has_cached_identity(&dest), + "precondition: the destination's identity is not cached" + ); + nodes[0] + .node + .queue_pending_tun_packet_for_test(dest, vec![0x60; 40]); + + let lookup = &nodes[0].node.metrics().lookup; + let (accepted, miss, timed_out) = ( + lookup.resp_accepted.get(), + lookup.resp_identity_miss.get(), + lookup.resp_timed_out.get(), + ); + + let victim = *nodes[0].node.node_addr(); + let payload = PathBroken::new(dest, p1).encode(); + let encoded = SessionDatagram::new(p1, victim, payload).encode(); + let start = wall_ms(); + nodes[0] + .node + .handle_session_datagram(&p1, &encoded[1..], false) + .await; + for _ in 0..10 { + tokio::time::sleep(std::time::Duration::from_millis(50)).await; + crate::node::tests::spanning_tree::process_available_packets(&mut nodes).await; + } + + let lookup = &nodes[0].node.metrics().lookup; + assert_eq!( + lookup.resp_identity_miss.get(), + miss, + "the lookup's answer could not be verified for want of an identity" + ); + assert!( + lookup.resp_accepted.get() > accepted, + "the lookup's answer must be verified and accepted" + ); + let entry = nodes[0].node.coord_cache().get_entry(&dest).unwrap(); + assert_eq!( + entry.coords().node_addrs().copied().collect::>(), + real + ); + assert!(entry.is_verified(wall_ms())); + + // Run the lookup schedule well past its end: an answered lookup has + // nothing left to time out, so the queued packet survives. + for step in 1..=6 { + nodes[0] + .node + .check_pending_lookups(start + step * 20_000) + .await; + } + assert_eq!( + nodes[0].node.metrics().lookup.resp_timed_out.get(), + timed_out, + "the lookup ran to its timeout" + ); + assert_eq!( + nodes[0].node.pending_tun_total_packets(), + 1, + "the packet queued for the destination was dropped" + ); + cleanup_nodes(&mut nodes).await; +} + /// The healthy path for keeping a verified entry below quorum: a genuine /// failure reported once is still recovered, because the lookup the report /// starts replaces the stale value with the destination's real position. From 9e349a338e57309fb94e8526a63e24d9c1c7ec21 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 1 Oct 2026 14:27:42 +0000 Subject: [PATCH 05/23] Close the macOS encrypt worker queue when its worker thread exits Only the sender's Drop ever marked the macOS worker queue closed. If a worker thread exited, for example by panicking, its receiver dropped without telling anyone: the rx_loop kept pushing jobs for that worker, and once the queue reached its cap push_blocking waited on a condition variable nothing would ever signal. The rx_loop runs on the daemon's single-threaded runtime, so the whole node stopped. Give the receiver a Drop that marks the queue closed, frees the queued jobs (they hold key copies and sockets) outside the lock, and wakes any waiting sender. Dispatch then takes its existing path for a gone worker and drops the job, as the Linux channel already does. The lock is taken poison-tolerantly because this Drop can run while the worker thread unwinds, where a second panic would abort the process. The queue was hardwired to the send-job type and compiled only on macOS, so none of its blocking behaviour ran in the Linux test run. It is now generic over its item type and compiled for tests on every platform; on macOS it is monomorphised to the same job type, and non-test builds elsewhere still do not compile it. Tests cover the queue's push, drain and close behaviour, including a worker that panics while a push is blocked on its full queue. That worker panics only on a signal, after the test has checked that the push did not return while the worker was alive. The poisoned-lock test holds its sender in ManuallyDrop, so a regression shows as an assertion failure rather than aborting the whole test binary. --- src/node/encrypt_worker.rs | 242 ++++++++++++++++++++++++++++++++----- 1 file changed, 209 insertions(+), 33 deletions(-) diff --git a/src/node/encrypt_worker.rs b/src/node/encrypt_worker.rs index 1c94a268..198d7fae 100644 --- a/src/node/encrypt_worker.rs +++ b/src/node/encrypt_worker.rs @@ -56,15 +56,17 @@ use crate::transport::udp::io::AsyncUdpSocket; #[cfg(not(target_os = "macos"))] use crossbeam_channel::{Receiver, SendError, Sender, TrySendError, bounded}; use ring::aead::{Aad, LessSafeKey, Nonce}; +#[cfg(any(target_os = "macos", test))] +use std::collections::VecDeque; #[cfg(target_os = "macos")] -use std::collections::{BTreeMap, HashMap, VecDeque}; +use std::collections::{BTreeMap, HashMap}; use std::net::SocketAddr; #[cfg(unix)] use std::os::unix::io::AsRawFd; use std::sync::Arc; use std::sync::OnceLock; -#[cfg(target_os = "macos")] -use std::sync::{Condvar, Mutex}; +#[cfg(any(target_os = "macos", test))] +use std::sync::{Condvar, Mutex, PoisonError}; use tracing::{debug, trace, warn}; /// A pre-cooked FMP-encrypt-and-send job. All state-touching work @@ -218,43 +220,42 @@ impl QueuedFmpSendJob { /// same rationale as the bounded endpoint_commands channel upstream. const WORKER_CHANNEL_CAP: usize = 1024; -#[cfg(target_os = "macos")] -struct MacWorkerSender { - inner: Arc, +#[cfg(any(target_os = "macos", test))] +struct MacWorkerSender { + inner: Arc>, } -#[cfg(target_os = "macos")] -struct MacWorkerReceiver { - inner: Arc, +#[cfg(any(target_os = "macos", test))] +struct MacWorkerReceiver { + inner: Arc>, } -#[cfg(target_os = "macos")] -struct MacWorkerQueueInner { - state: Mutex, +#[cfg(any(target_os = "macos", test))] +struct MacWorkerQueueInner { + state: Mutex>, not_empty: Condvar, not_full: Condvar, cap: usize, } -#[cfg(target_os = "macos")] -#[derive(Default)] -struct MacWorkerQueueState { - queue: VecDeque, +#[cfg(any(target_os = "macos", test))] +struct MacWorkerQueueState { + queue: VecDeque, waiting: bool, closed: bool, } -#[cfg(target_os = "macos")] -enum MacWorkerTryPushError { - Full(Box), +#[cfg(any(target_os = "macos", test))] +enum MacWorkerTryPushError { + Full(Box), Closed, } -#[cfg(target_os = "macos")] +#[cfg(any(target_os = "macos", test))] struct MacWorkerPushError; -#[cfg(target_os = "macos")] -fn mac_worker_channel(cap: usize) -> (MacWorkerSender, MacWorkerReceiver) { +#[cfg(any(target_os = "macos", test))] +fn mac_worker_channel(cap: usize) -> (MacWorkerSender, MacWorkerReceiver) { let inner = Arc::new(MacWorkerQueueInner { state: Mutex::new(MacWorkerQueueState { queue: VecDeque::with_capacity(cap), @@ -273,9 +274,9 @@ fn mac_worker_channel(cap: usize) -> (MacWorkerSender, MacWorkerReceiver) { ) } -#[cfg(target_os = "macos")] -impl MacWorkerSender { - fn try_push(&self, job: QueuedFmpSendJob) -> Result<(), MacWorkerTryPushError> { +#[cfg(any(target_os = "macos", test))] +impl MacWorkerSender { + fn try_push(&self, job: T) -> Result<(), MacWorkerTryPushError> { let mut state = self .inner .state @@ -298,7 +299,7 @@ impl MacWorkerSender { Ok(()) } - fn push_blocking(&self, job: QueuedFmpSendJob) -> Result<(), MacWorkerPushError> { + fn push_blocking(&self, job: T) -> Result<(), MacWorkerPushError> { let mut state = self .inner .state @@ -328,8 +329,8 @@ impl MacWorkerSender { } } -#[cfg(target_os = "macos")] -impl Drop for MacWorkerSender { +#[cfg(any(target_os = "macos", test))] +impl Drop for MacWorkerSender { fn drop(&mut self) { let mut state = self .inner @@ -343,9 +344,34 @@ impl Drop for MacWorkerSender { } } -#[cfg(target_os = "macos")] -impl MacWorkerReceiver { - fn recv_batch(&self, batch: &mut Vec, max: usize) -> bool { +/// Closing from the receiver side matters because a sender waiting in +/// `push_blocking` on a full queue sleeps on `not_full`, and only the +/// receiver draining the queue wakes it. If the worker thread exits (a +/// panic unwinding included), nothing else would, and the rx_loop behind +/// that sender would block forever. +#[cfg(any(target_os = "macos", test))] +impl Drop for MacWorkerReceiver { + fn drop(&mut self) { + // This can run while the worker thread unwinds from a panic, where + // a second panic would abort the process, so tolerate poisoning. + let mut state = self + .inner + .state + .lock() + .unwrap_or_else(PoisonError::into_inner); + state.closed = true; + let queued = std::mem::take(&mut state.queue); + drop(state); + // Queued jobs hold key copies and sockets; free them outside the lock. + drop(queued); + self.inner.not_full.notify_all(); + self.inner.not_empty.notify_all(); + } +} + +#[cfg(any(target_os = "macos", test))] +impl MacWorkerReceiver { + fn recv_batch(&self, batch: &mut Vec, max: usize) -> bool { debug_assert!(batch.is_empty()); let mut state = self .inner @@ -378,7 +404,7 @@ impl MacWorkerReceiver { } #[cfg(target_os = "macos")] -type WorkerSender = MacWorkerSender; +type WorkerSender = MacWorkerSender; #[cfg(not(target_os = "macos"))] type WorkerSender = Sender; @@ -955,7 +981,7 @@ fn run_worker(idx: usize, rx: Receiver) { } #[cfg(target_os = "macos")] -fn run_worker_macos(idx: usize, rx: MacWorkerReceiver) { +fn run_worker_macos(idx: usize, rx: MacWorkerReceiver) { trace!(worker = idx, "FMP encrypt worker thread starting"); let batch_size = macos_worker_batch_size(); @@ -2402,3 +2428,153 @@ fn send_one_raw( Ok(r as usize) } } + +/// Tests for the bounded worker queue the macOS encrypt pool uses. The +/// queue is generic, so these run on every platform against small item +/// types. Every wait is bounded so a regression fails instead of hanging. +#[cfg(test)] +mod mac_queue_tests { + use super::*; + use std::panic::{AssertUnwindSafe, catch_unwind}; + use std::sync::mpsc; + use std::thread; + use std::time::Duration; + + const WAIT: Duration = Duration::from_secs(5); + + /// Run `push_blocking(item)` on a helper thread and return a channel + /// that yields its result. + fn spawn_pusher( + tx: MacWorkerSender, + item: T, + ) -> mpsc::Receiver> { + let (done_tx, done_rx) = mpsc::channel(); + thread::spawn(move || { + let result = tx.push_blocking(item); + let _ = done_tx.send(result); + }); + done_rx + } + + #[test] + fn push_blocking_returns_error_when_worker_thread_panics_with_full_queue() { + let (tx, rx) = mac_worker_channel::(2); + assert!(tx.try_push(1).is_ok()); + assert!(tx.try_push(2).is_ok()); + match tx.try_push(9) { + Err(MacWorkerTryPushError::Full(job)) => assert_eq!(*job, 9), + _ => panic!("try_push on a full queue should hand the job back"), + } + + // The worker owns the receiver and dies without draining, once the + // pusher below has had time to start waiting for space. + let (die_tx, die_rx) = mpsc::channel::<()>(); + let worker = thread::spawn(move || { + let _rx = rx; + let _ = die_rx.recv(); + panic!("simulated encrypt worker panic"); + }); + let done = spawn_pusher(tx, 3); + thread::sleep(Duration::from_millis(200)); + assert!( + matches!(done.try_recv(), Err(mpsc::TryRecvError::Empty)), + "push_blocking returned while the queue was full and the worker alive" + ); + die_tx + .send(()) + .expect("worker thread gone before its signal"); + + let result = done + .recv_timeout(WAIT) + .expect("push_blocking still blocked after the worker thread died"); + assert!(matches!(result, Err(MacWorkerPushError))); + assert!(worker.join().is_err(), "worker thread should have panicked"); + } + + #[test] + fn try_push_returns_closed_after_receiver_dropped() { + let (tx, rx) = mac_worker_channel::(2); + drop(rx); + assert!(matches!(tx.try_push(1), Err(MacWorkerTryPushError::Closed))); + let done = spawn_pusher(tx, 2); + let result = done + .recv_timeout(WAIT) + .expect("push_blocking blocked on a queue whose receiver is gone"); + assert!(matches!(result, Err(MacWorkerPushError))); + } + + #[test] + fn receiver_drop_releases_queued_items() { + let marker = Arc::new(()); + let (tx, rx) = mac_worker_channel::>(4); + assert!(tx.try_push(Arc::clone(&marker)).is_ok()); + assert!(tx.try_push(Arc::clone(&marker)).is_ok()); + assert_eq!(Arc::strong_count(&marker), 3); + drop(rx); + assert_eq!( + Arc::strong_count(&marker), + 1, + "queued items must be freed when the receiver goes away" + ); + drop(tx); + } + + #[test] + fn push_blocking_completes_when_worker_drains_full_queue() { + let (tx, rx) = mac_worker_channel::(2); + assert!(tx.try_push(1).is_ok()); + assert!(tx.try_push(2).is_ok()); + let done = spawn_pusher(tx, 3); + thread::sleep(Duration::from_millis(100)); + assert!( + matches!(done.try_recv(), Err(mpsc::TryRecvError::Empty)), + "push_blocking returned while the queue was still full" + ); + + let mut batch = Vec::new(); + assert!(rx.recv_batch(&mut batch, 16)); + assert_eq!(batch, vec![1, 2]); + let result = done + .recv_timeout(WAIT) + .expect("push_blocking not woken after the worker drained"); + assert!(result.is_ok()); + + batch.clear(); + assert!(rx.recv_batch(&mut batch, 16)); + assert_eq!(batch, vec![3]); + } + + #[test] + fn recv_batch_drains_then_reports_closed_after_sender_drop() { + let (tx, rx) = mac_worker_channel::(4); + for i in 1..=3 { + assert!(tx.try_push(i).is_ok()); + } + drop(tx); + let mut batch = Vec::new(); + assert!(rx.recv_batch(&mut batch, 16)); + assert_eq!(batch, vec![1, 2, 3]); + batch.clear(); + assert!(!rx.recv_batch(&mut batch, 16)); + assert!(batch.is_empty()); + } + + #[test] + fn receiver_drop_does_not_panic_on_poisoned_lock() { + let (tx, rx) = mac_worker_channel::(2); + // The sender's own Drop still expects an unpoisoned lock, so it must + // never run here: a panic there while a failed assertion unwinds + // would abort the whole test binary. + let _tx = std::mem::ManuallyDrop::new(tx); + let inner = Arc::clone(&rx.inner); + let poisoner = thread::spawn(move || { + let _guard = inner.state.lock().unwrap(); + panic!("poison the queue lock"); + }); + assert!(poisoner.join().is_err()); + assert!(rx.inner.state.is_poisoned()); + + let dropped = catch_unwind(AssertUnwindSafe(move || drop(rx))); + assert!(dropped.is_ok(), "receiver drop panicked on a poisoned lock"); + } +} From d4b36e4da19f3ec6e234e18a567672f8d368fac6 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 1 Oct 2026 14:45:29 +0000 Subject: [PATCH 06/23] Measure whether the kernel keeps an in-flight socket whose only reference is the message The native API hands a flow's descriptor to a listener's client inside an SCM_RIGHTS message and then closes its own copy. Reading xnu's descriptor collector says Darwin flushes such a socket if a collection runs before the client reads the message, because it roots only in-flight files and the listener's client half is not one. That would explain the intermittent macOS failures of two native API tests, where an accepted flow read as end of file with its held datagram gone. Add a probe that performs the same hand-off with the product's own pair type and hand-off code, provokes a collection by freeing an unrelated unix socket, and reads the received descriptor: Darwin is expected to return end of file, Linux the held bytes. On Linux 6.8, closing a unix socket while any descriptor is in flight runs the collector before the close returns, and the flow's half is in flight at that moment, so the Linux probe also exercises Linux's collector and checks that it keeps the socket. --- src/native/dgram_probe.rs | 99 +++++++++++++++++++++++++++++++++++++++ 1 file changed, 99 insertions(+) diff --git a/src/native/dgram_probe.rs b/src/native/dgram_probe.rs index 6e36952b..3550cc11 100644 --- a/src/native/dgram_probe.rs +++ b/src/native/dgram_probe.rs @@ -20,6 +20,17 @@ //! every unix so the platforms can be compared without the test itself being a //! variable. //! +//! **A second measurement lives here: what the kernel does to a socket that is +//! in flight.** When the daemon hands a flow's descriptor to a listener's +//! client, the descriptor sits inside an `SCM_RIGHTS` message until the client +//! reads it. xnu's descriptor garbage collector takes its roots only from files +//! that are in flight, and treats one whose only references are messages as +//! unreachable unless it is found in the receive buffer of another in-flight +//! socket. The listener's client half is not in flight, so a flow socket whose +//! daemon copy has been closed is flushed by any collection that runs before +//! the client reads it. Linux keeps such a socket. The tests at the end of this +//! file measure that difference directly. +//! //! `SOCK_CLOEXEC` is deliberately not passed in the type argument, though //! `super::seqpacket::pair` does pass it. Linux and FreeBSD accept it there and //! macOS does not, and that difference belongs to the port rather than to this @@ -464,3 +475,91 @@ fn freebsd_seqpacket_drops_a_zero_length_message_instead_of_delivering_it() { diagnosis of the stalled FreeBSD runs is wrong." ); } + +/// How long a collection is given to run after it has been queued. +/// +/// xnu runs its descriptor collector as an asynchronous thread call, so the +/// close that queues it returns before it has run. A fixed wait is the only +/// handle a test has on it; if the Darwin probe below ever misses, this is the +/// first number to raise. +const COLLECTION_WAIT: std::time::Duration = std::time::Duration::from_millis(100); + +/// Give the kernel's descriptor collector a reason to run, then time to run. +/// +/// Freeing any `AF_UNIX` socket queues xnu's collector, so a fresh pair closed +/// at once is enough. That is also why the collector can run at any moment on +/// a busy host: every unix socket any process closes queues it. Linux runs its +/// collector here too: in 6.8, the kernel this was read against, closing a unix +/// socket while any descriptor is in flight runs it before the close returns. +/// So the Linux probe below checks that Linux's collector keeps the socket, +/// rather than only that no collector ran. +pub(super) fn provoke_collection() { + let (a, b) = dgram_pair().expect("AF_UNIX SOCK_DGRAM socketpair"); + drop(a); + drop(b); + std::thread::sleep(COLLECTION_WAIT); +} + +/// Hand a flow's client half across a listener pair the way the daemon does, +/// close the sender's copy, provoke a collection, and read what reaches the +/// receiver. +/// +/// Built from the product's own pair type and hand-off code rather than from +/// [`dgram_pair`], so the measurement is of exactly the sockets the daemon +/// uses. Returns the first read on the received descriptor: the held bytes if +/// the socket survived, zero bytes if the collector flushed it. +#[cfg(any(target_os = "linux", target_os = "macos"))] +fn read_after_collection_in_flight() -> io::Result { + use super::{fdpass, seqpacket}; + use std::os::fd::AsFd; + + let (daemon, flow) = seqpacket::pair()?; + let (sender, receiver) = seqpacket::pair()?; + assert_eq!(send(daemon.as_raw_fd(), b"held")?, 4); + + fdpass::try_send(sender.as_raw_fd(), b"arrival", Some(flow.as_fd()))?; + // The message is now the only reference to the flow's client half. + drop(flow); + + provoke_collection(); + + let mut buf = [0u8; 64]; + let chunk = fdpass::recv(receiver.as_raw_fd(), &mut buf)?; + let flow = chunk + .fd + .ok_or_else(|| io::Error::other("the message arrived without its descriptor"))?; + recv(flow.as_raw_fd(), &mut buf) +} + +/// Darwin flushes an in-flight socket whose only reference is the message. +/// +/// This is the defect behind the native API's intermittent macOS failures, in +/// isolation: the listener's client receives a flow descriptor that reads as +/// end of file, with the datagram written to it before the hand-off gone. +/// A failure here means the collector did not run within +/// [`COLLECTION_WAIT`], or that Darwin no longer collects such a socket. In +/// the second case the daemon's hold on a handed-over descriptor is no longer +/// needed there. +#[cfg(target_os = "macos")] +#[test] +fn darwin_collects_an_in_flight_socket_whose_only_reference_is_the_message() { + let read = read_after_collection_in_flight(); + assert!( + matches!(read, Ok(0)), + "expected the collector to flush the in-flight flow socket, so its first \ + read returns end of file; got {read:?}" + ); +} + +/// Linux keeps the same socket through the collection the close provokes: a +/// queue held by a socket that is not in flight counts as a reference to what +/// it holds. +#[cfg(target_os = "linux")] +#[test] +fn linux_keeps_an_in_flight_socket_whose_only_reference_is_the_message() { + let read = read_after_collection_in_flight(); + assert!( + matches!(read, Ok(4)), + "expected the in-flight flow socket to survive with its datagram; got {read:?}" + ); +} From 8b396b662db5a8f60acb6a1b08736855fc390dbb Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 1 Oct 2026 22:40:40 +0000 Subject: [PATCH 07/23] Keep the daemon's copy of a native API descriptor until the client holds it A native API flow's descriptor reaches the client inside a message: an arrival for a flow a listener accepts, or a connect or listen reply. The daemon closed its own copy once the message was written, so until the client read it, the message was the only reference to the socket. xnu's descriptor collector flushes a socket in that state, and the client then receives a flow that reads as end of file with the datagrams the daemon held for it gone. That is the intermittent macOS failure of the two listener tests. A listener now keeps the daemon's copy of each flow it hands over until the client's first write on the flow, the listener's close, or the flow's end, on every platform the listener builds on. The flow is recorded before its reader starts, so a write already queued cannot race the record. The connection's serving loop, now a method on the connection so a test can run it over a real socket, keeps the copy sent in a connect or listen reply until the client's next command on that connection or the connection's end of file. A client sends its next command only after reading the reply, so either event means the descriptor has left the message. The shipped client closes the connection as soon as it has the reply, so it sees no change. The cost is accepted and documented: a flow a client accepts and closes without ever writing stays open, holding its port and a flow slot, until the listener closes, so a server that refuses flows by dropping them pays for each one until then; a client speaking the protocol directly that leaves its setup connection open sees a flow or listener it closes stay open until its next command or the connection's close. The reference and how-to pages, the client rustdoc on FipsStream, FipsListener and accept, and the design note say so, and the security reference records that a remote peer opening flows from many source ports to such a server can exhaust the node-wide max_flows ceiling. Tests cover an arrival surviving a provoked collection (deterministic on macOS), a held flow outliving its dropped descriptor until the listener closes, a client's write releasing it, a flow its client still holds working after the listener has closed and let its copy go, and a reply's copy kept until the next command and let go when the connection ends. The native API harness asserts the new lifetime of a refused flow. On macOS and FreeBSD the daemon notices a client's close only when a reader retries its read, up to a quarter second later, so the test helpers that wait for a close (forgotten, rebind, settle_closed) retry for up to five seconds by the clock rather than for a count of yields, and still_open waits two retry intervals there before asserting a flow is still open. --- docs/design/fips-native-api.md | 9 +- docs/how-to/use-the-native-datagram-api.md | 10 + docs/how-to/write-a-native-api-client.md | 12 +- docs/reference/native-api.md | 37 +- docs/reference/security.md | 26 +- src/native/client/mod.rs | 15 +- src/native/mod.rs | 581 +++++++++++++++++++-- src/native/seqpacket.rs | 13 + testing/native-api/client.py | 4 +- testing/native-api/test.sh | 26 +- 10 files changed, 661 insertions(+), 72 deletions(-) diff --git a/docs/design/fips-native-api.md b/docs/design/fips-native-api.md index 26ce5ea6..c16fcc55 100644 --- a/docs/design/fips-native-api.md +++ b/docs/design/fips-native-api.md @@ -133,9 +133,12 @@ and silent drops on its listeners. **Not a connection in the TCP sense.** A successful `connect` is a local registration and contacts no peer. There is no handshake, no keepalive and no notification that a peer went away. A flow ends when its descriptor closes, and -in no other way. In particular **a peer cannot end your flow: it has no close to -send.** That single fact shapes every program written against this interface, -and the consequences are drawn out in +in no other way. The one delay is an accepted flow never sent on, which ends +only once its listener has closed as well (see +[../reference/native-api.md](../reference/native-api.md#fipslistener)). In +particular **a peer cannot end your flow: it has no close to send.** That +single fact shapes every program written against this interface, and the +consequences are drawn out in [../how-to/use-the-native-datagram-api.md](../how-to/use-the-native-datagram-api.md#four-things-that-will-bite-you). ## See also diff --git a/docs/how-to/use-the-native-datagram-api.md b/docs/how-to/use-the-native-datagram-api.md index 355ba093..b7a19493 100644 --- a/docs/how-to/use-the-native-datagram-api.md +++ b/docs/how-to/use-the-native-datagram-api.md @@ -227,6 +227,16 @@ leaves the flows already accepted from it untouched. A program that parks streams in a `Vec` and never removes them holds ports and flow slots exactly as if it had leaked descriptors. +**An accepted flow you drop without ever sending on stays open until you +drop its listener.** Until your program has sent on an accepted flow, the +daemon keeps its own copy of the flow's descriptor, because on macOS the +kernel can otherwise destroy the flow while its descriptor is still on the +way to you. Your first `send` on the flow, or dropping the listener, lets +that copy go. So a server that refuses flows by dropping them unanswered +holds a port and a flow slot for each one until its listener goes, and a +long-lived listener that refuses many flows can walk the node into its flow +ceiling. + **Nothing peer-driven ever ends a flow, so your program has to.** The v1 wire carries no half-close. Nothing closes the daemon's half of a live accepted flow, so a loop written as "echo until the flow closes", or one diff --git a/docs/how-to/write-a-native-api-client.md b/docs/how-to/write-a-native-api-client.md index a128874a..a45c96ad 100644 --- a/docs/how-to/write-a-native-api-client.md +++ b/docs/how-to/write-a-native-api-client.md @@ -124,7 +124,7 @@ one producer on the surface, which is what lets a caller read it. ## Step 6: Keep descriptor hygiene -Five rules. Each one leaks a flow or loses one when broken. +Six rules. Each one leaks a flow or loses one when broken. **Request close-on-exec** with `MSG_CMSG_CLOEXEC` on the `recvmsg`, rather than setting it afterwards. Without it the descriptor survives an `exec` into a @@ -140,7 +140,15 @@ descriptors rather than dropping them on the floor. **Lift the descriptor out of an arrival you cannot parse** before discarding the message. Refusing a flow is closing its descriptor; discarding the message -without taking it leaks the flow instead. +without taking it leaks the flow instead. A refused flow ends only when the +listener closes, though, unless you wrote on it first: the daemon keeps its +own copy of an accepted flow's descriptor until your first write or the +listener's close. + +**Close the setup connection once you have the reply.** The daemon keeps its +own copy of the descriptor in its last reply until your next command on that +connection or the connection's close. A flow or listener you close while the +connection sits idle stays open until one of those happens. **Bound the partial line.** A daemon that stopped sending newlines would otherwise grow your buffer without end. The shipped client caps it at 64 KiB, diff --git a/docs/reference/native-api.md b/docs/reference/native-api.md index cea846d8..a6a21501 100644 --- a/docs/reference/native-api.md +++ b/docs/reference/native-api.md @@ -195,7 +195,10 @@ That is what lets the blanket reference implementation cover `&str` and One datagram flow, and the descriptor it rides on. A flow is an exact match of both ends and both ports. **The descriptor is the flow**: it lives while a -process holds that descriptor and ends when the last one closes it. +process holds that descriptor and ends when the last one closes it. The one +exception is an accepted flow that has never been sent on: the daemon keeps +its own copy of that descriptor until the first `send` or until the listener +is dropped. See `accept` below. `Send + Sync + 'static`, with no `Arc` and no borrow. There is **no `Clone` and no `try_clone`**. Because `send` and `recv` both take `&self`, a shared borrow @@ -340,6 +343,15 @@ no other way to refuse one. An unparseable arrival is therefore reported only after the descriptor it carried has been taken into ownership, so a parse failure refuses the flow rather than leaking it. +**A refused flow stays open until its listener is dropped**, unless the +program sent on it first. Until the first `send` on an accepted flow, the +daemon keeps its own copy of the flow's descriptor, because on macOS the +kernel can otherwise destroy a socket whose descriptor is still in an unread +arrival. That copy goes at the first `send`, when the listener is dropped, or +when the flow ends any other way, and the flow then ends with the program's +own close. Until then a dropped flow holds its port and its slot against +`max_flows`. + **`incoming()`** returns an `Incoming<'_>`, which borrows the listener for the iterator's lifetime, so the listener cannot be moved or dropped mid-iteration. @@ -360,7 +372,8 @@ none either. A bounded accept is `set_nonblocking` plus a wait of the caller's own on the descriptor. **Dropping** closes the descriptor and unbinds the port. Flows already accepted -from it are untouched; flows still pending on it go with it. +from it and still held are untouched; flows still pending on it go with it, and +so do flows accepted from it and dropped without ever being sent on. ### Incoming @@ -613,6 +626,20 @@ own half non-blocking and leaves the client's half blocking. `SOCK_SEQPACKET` is what preserves message boundaries in both directions, which is why the payload needs no framing. +**The daemon keeps a copy of the client's half after sending it.** While a +descriptor sits unread in a message, the message can be its only reference, +and the macOS kernel's descriptor collector destroys a socket in that state: +the client then receives a flow that reads as end of file with its datagrams +gone. So the daemon keeps its copy until one of these: + +- For a `connect` or `listen` reply, at the client's next command on the same + connection, or when that connection closes. +- For an arrival on a listener, at the client's first write on the flow, when + the listener closes, or when the flow ends any other way. + +A flow or listener the client closes before then ends when the daemon's copy +goes, not at the client's close. + A refused `connect` leaves the port free: the socket pair is built before the port is claimed, so a failure to build it needs no rollback. @@ -626,7 +653,11 @@ returning. **The connection owns nothing.** Closing it releases no flow and no listener, and a descriptor kept across the close keeps working. What owns the flow is the -descriptor. +descriptor. The connection does delay one thing: a descriptor from its last +reply that the client closes while the connection is still open, with no +further command sent, stays open until the next command or the connection's +close (see Passing the descriptor). The shipped client closes the connection +as soon as it has the reply, so it never meets this. ## Command reference diff --git a/docs/reference/security.md b/docs/reference/security.md index b69b5ea7..b140f893 100644 --- a/docs/reference/security.md +++ b/docs/reference/security.md @@ -245,12 +245,14 @@ machine. **The file descriptor carries the grant, not the connection.** A setup call hands the client a socket descriptor and the connection it was made on is then -closed; the flow or the held port lives until that descriptor is closed. A -descriptor is an ordinary kernel object, so it survives `fork`, survives -`exec` unless the client asked for it close-on-exec when it received it, and -can be handed to another process over `SCM_RIGHTS`. A process holding one can -send as this node on that flow, or receive on that port, without ever opening -the API socket and without being in the `fips` group. +closed; the flow or the held port lives until that descriptor is closed and +the daemon has let go of the copy it keeps while the descriptor is being +handed over. A descriptor is an ordinary kernel object, so it survives +`fork`, survives `exec` unless the client asked for it close-on-exec when it +received it, and can be handed to another process over `SCM_RIGHTS`. A +process holding one can send as this node on that flow, or receive on that +port, without ever opening the API socket and without being in the `fips` +group. Nothing revokes a descriptor already handed out. Restarting the daemon closes its own halves and ends every flow and listener at once, and that is the only revocation there is. @@ -269,6 +271,18 @@ peer had sent it, reaching any listener on this node under any peer identity the caller names. Leave it off outside a test harness; a packaged node does not enable it. +**A remote peer can fill the node's flow ceiling through a server that +refuses flows by dropping them.** Until a program first sends on a flow it +accepted, the daemon keeps its own copy of that flow's descriptor, so a flow +accepted and dropped unanswered keeps its slot against the node-wide +`node.native_api.max_flows` until its listener is dropped. A peer that opens +flows to such a listener from many source ports can therefore exhaust the +ceiling, and every other program on the node then gets `EMFILE` on `connect` +and silently loses arrivals on its listeners. This is the accepted cost of +keeping a flow alive while its descriptor is on the way to the program; see +[../how-to/use-the-native-datagram-api.md](../how-to/use-the-native-datagram-api.md) +for what releases the daemon's copy. + The socket is local only. It is not reachable over the network, and nothing about it changes the mesh's own authentication: a peer still verifies the node's signature, which is precisely why a local caller that can send through diff --git a/src/native/client/mod.rs b/src/native/client/mod.rs index ee4c8ade..f02ec7f1 100644 --- a/src/native/client/mod.rs +++ b/src/native/client/mod.rs @@ -373,7 +373,10 @@ fn expired(error: io::Error) -> io::Error { /// /// The protocol has no close command. Dropping the stream closes its /// descriptor, and that is what releases the flow and its local port at the -/// daemon. +/// daemon. An accepted flow that was never sent on is the exception: the daemon +/// keeps its own copy of its descriptor until the first [`FipsStream::send`] or +/// until the listener is dropped, so dropping the stream before either releases +/// nothing yet. #[derive(Debug)] pub struct FipsStream { fd: OwnedFd, @@ -627,8 +630,9 @@ impl AsFd for FipsStream { /// **The listener is a descriptor**, which is what makes it pollable: it joins /// an existing `poll`, `select` or `epoll` loop with no new mechanism, and /// [`FipsListener::accept`] is one `recvmsg` on it. Dropping the listener closes -/// that descriptor, which unbinds the port; flows already accepted from it are -/// untouched. +/// that descriptor, which unbinds the port; flows already accepted from it and +/// still held are untouched, and those accepted and dropped without ever being +/// sent on end with it. #[derive(Debug)] pub struct FipsListener { fd: OwnedFd, @@ -676,7 +680,10 @@ impl FipsListener { /// Refusing a flow is dropping the stream, which closes its descriptor. /// There is no other way to refuse one, which is why an unreadable arrival /// message is reported after the descriptor it carried has been taken: the - /// flow is then refused rather than leaked. + /// flow is then refused rather than leaked. A refused flow that was never + /// sent on still holds its port and its slot against the node's flow limit + /// until this listener is dropped, because the daemon keeps its own copy of + /// the descriptor until then. pub fn accept(&self) -> io::Result<(FipsStream, FipsAddr)> { let mut buf = [0u8; CHUNK]; let chunk = fdpass::recv(self.fd.as_raw_fd(), &mut buf)?; diff --git a/src/native/mod.rs b/src/native/mod.rs index 57a772f4..305c5760 100644 --- a/src/native/mod.rs +++ b/src/native/mod.rs @@ -18,7 +18,20 @@ //! replies and then has no further part in anything it opened. A flow lives //! until its own descriptor reaches end of file and a listener until its own //! does, whichever task holds them, which is what makes a descriptor this API -//! hands back behave like one a syscall would have. +//! hands back behave like one a syscall would have. The one exception: the +//! connection keeps a copy of the descriptor in its last reply until the +//! client's next command or its close, for the same reason a listener keeps a +//! copy of a flow it hands over (below). The shipped client closes the +//! connection as soon as it has the reply, so for it the copy is gone at once. +//! +//! **A listener keeps a copy of a flow it handed over.** While a flow's +//! descriptor sits unread in an arrival message, the message can be the only +//! reference to it, and xnu's descriptor collector flushes a socket in that +//! state. So `hand_over` keeps the daemon's copy until the client has written +//! on the flow or closed the listener, and a flow then reaches end of file once +//! both the client and that copy have let go. The cost is that a client which +//! accepts a flow and closes it without writing leaves it open until the +//! listener closes. See `dgram_probe.rs` for the measurement. //! //! **The wire is connected.** A datagram a client writes leaves this node over //! FSP, and one arriving on a held port reaches its flow. `max_payload` is the @@ -233,6 +246,15 @@ mod unix_impl { /// writes into. sock: Arc, counts: Arc, + /// The daemon's copy of the client's half, kept while that half may + /// still be in flight to a listener's client. + /// + /// Holding it keeps the socket reachable from outside the message that + /// carries it, which is what stops xnu's collector flushing it before + /// the client reads the arrival. `None` once released, and always for + /// a connected flow, whose descriptor went back in an RPC reply and is + /// held by the connection that sent it instead; see `Connection::run`. + pin: Option, } /// What `stats` reports about one flow. @@ -283,6 +305,31 @@ mod unix_impl { self.table().remove(&id); } + /// Let go of the daemon's copy of a flow's client half. + /// + /// The copy is dropped after the lock is released: if the client has + /// already closed its own, this is the last reference, and the close it + /// causes is the flow's end of file. + fn unpin(&self, id: u64) { + let pin = self.table().get_mut(&id).and_then(|flow| flow.pin.take()); + drop(pin); + } + + /// Let go of the daemon's copy of every flow on a local port. + /// + /// Called when a listener ends. Every pinned flow on its port came from + /// it: a connected flow carries no pin, and no later listener can hold + /// the port until this one's release has been served. + fn unpin_port(&self, port: u16) { + let pins: Vec = self + .table() + .values_mut() + .filter(|flow| flow.local == port) + .filter_map(|flow| flow.pin.take()) + .collect(); + drop(pins); + } + /// What the debug `stats` command reports, or `None` for a flow this /// node does not hold. fn stats(&self, id: u64) -> Option { @@ -503,6 +550,7 @@ mod unix_impl { key, peer, wiring, + None, &self.outbound, &self.node, &self.flows, @@ -700,6 +748,19 @@ mod unix_impl { #[cfg(test)] impl Connection { + /// Another connection to the same node, sharing its flow table, as a + /// second client of one daemon has. + pub(super) fn sibling(&self) -> Self { + Self::new( + self.node.clone(), + self.outbound.clone(), + Arc::clone(&self.flows), + self.limits, + Arc::clone(&self.npub), + self.debug, + ) + } + /// Build a connection wired to a channel a test serves. pub(super) fn for_test( node: mpsc::Sender, @@ -730,14 +791,24 @@ mod unix_impl { } } - /// Wait until a flow's reader task has observed the client's close. + /// Wait until a flow's reader task has observed the client's close, + /// failing by name if it never does. + /// + /// A flow the node no longer holds counts as closed: the reader sets + /// the flag and forgets the flow in the same turn, so the flag alone is + /// almost never there to be seen. Bounded by the clock rather than by + /// turns of the runtime, for the reason given at `CLOSE_WAIT`. pub(super) async fn settle_closed(&self, flow: u64) { - for _ in 0..1000 { - match self.flows.stats(flow) { - Some(stats) if stats.closed => return, - _ => tokio::task::yield_now().await, - } - } + let seen = super::tests::eventually(async || match self.flows.stats(flow) { + Some(stats) if !stats.closed => None, + _ => Some(()), + }) + .await; + assert!( + seen.is_some(), + "the reader never saw flow {flow} close within {:?}", + super::tests::CLOSE_WAIT + ); } } @@ -780,35 +851,48 @@ mod unix_impl { /// The reader cannot start any earlier: it stamps every datagram it /// forwards with the flow's key and the peer's address, and the local port /// is not known until the registry has answered. + /// + /// The flow is recorded before the reader is spawned. A handed-over client + /// can already have written, and on a multi-thread runtime the reader can + /// run at once, so recording afterwards would let its unpin find no flow + /// and leave the pin in place for the flow's whole life. + #[allow( + clippy::too_many_arguments, + reason = "one hand-off of plumbing to two tasks, with no state to group" + )] fn start( id: u64, key: FlowKey, peer: XOnlyPublicKey, wiring: Wiring, + pin: Option, outbound: &mpsc::Sender, node: &mpsc::Sender, flows: &Arc, ) { let counts = Arc::new(Counts::default()); + let pinned = pin.is_some(); + flows.record( + id, + Flow { + local: key.local, + sock: Arc::clone(&wiring.sock), + counts: Arc::clone(&counts), + pin, + }, + ); tokio::spawn(drain( id, key, peer, Arc::clone(&wiring.sock), - Arc::clone(&counts), + counts, + pinned, outbound.clone(), node.clone(), Arc::clone(flows), )); - tokio::spawn(feed(Arc::clone(&wiring.sock), wiring.inbound)); - flows.record( - id, - Flow { - local: key.local, - sock: wiring.sock, - counts, - }, - ); + tokio::spawn(feed(wiring.sock, wiring.inbound)); } /// Serve one listener until its client closes the descriptor. @@ -872,6 +956,15 @@ mod unix_impl { } debug!(port, "Native API listener closed by its client"); + // Before the release, so the port cannot have passed to another + // listener whose flows this would also let go of. When the client has + // closed the listener, an arrival it never read went with the + // listener's receive queue, so no descriptor from this port is still + // in flight to it. The other two ways out of the loop, the node going + // away and a failed read, can leave the client's half open with + // arrivals unread on it. Both are teardown, and the hold goes with the + // listener rather than outliving it. + flows.unpin_port(port); let _ = node .send(NativeMessage::Release { flows: Vec::new(), @@ -891,6 +984,15 @@ mod unix_impl { /// Both writes are try-sends. They go onto a socket pair whose client half /// has not been sent yet, so no process can read either one and a task that /// parked on one would stop serving this listener entirely. + /// + /// The daemon keeps its copy of the client's half after the arrival is + /// written, until the client writes on the flow or closes the listener. + /// Until the client reads the arrival, the message carrying the descriptor + /// is otherwise its only reference, and xnu's collector flushes a socket in + /// that state: the client then receives a flow that reads as end of file, + /// with its held datagrams gone. Nothing in the protocol says when the + /// client has read the arrival, so a client that closes the flow without + /// writing leaves it open until the listener closes. async fn hand_over( arrival: Arrival, listener: &Seqpacket, @@ -968,13 +1070,11 @@ mod unix_impl { accepted.key, accepted.peer, wiring, + Some(theirs), outbound, node, flows, ); - // Dropping our copy leaves the client holding the only reference to its - // half, so its close tears the flow down. - drop(theirs); } /// What a failed hand-off write says about the client, for the counter. @@ -1098,6 +1198,11 @@ mod unix_impl { /// Counting continues alongside the forwarding: `stats` is how a check /// observes that a datagram reached the daemon, independently of whether it /// then reached a peer. + /// + /// `pinned` says the daemon still holds a copy of the client's half. The + /// first datagram the client writes proves it holds the descriptor, so the + /// copy is let go then, once, and the client's close ends the flow from + /// there on. #[allow( clippy::too_many_arguments, reason = "one hand-off of plumbing to a task, with no state to group" @@ -1108,6 +1213,7 @@ mod unix_impl { peer: XOnlyPublicKey, sock: Arc, counts: Arc, + mut pinned: bool, outbound: mpsc::Sender, node: mpsc::Sender, flows: Arc, @@ -1116,6 +1222,10 @@ mod unix_impl { loop { match sock.recv(&mut buf).await { Ok(Received::Datagram(len)) => { + if pinned { + flows.unpin(id); + pinned = false; + } counts.datagrams.fetch_add(1, Ordering::Relaxed); counts.bytes.fetch_add(len as u64, Ordering::Relaxed); let sent = outbound @@ -1210,21 +1320,42 @@ mod unix_impl { npub: Arc, debug: bool, ) -> Result<(), std::io::Error> { - let mut connection = Connection::new(node, outbound, flows, limits, npub, debug); - let mut reader = BufReader::new(stream); - let mut line = Vec::new(); + Connection::new(node, outbound, flows, limits, npub, debug) + .run(stream) + .await + } - while read_command(&mut reader, &mut line).await? { - let (response, fd) = connection.answer(&line).await; - let mut json = serde_json::to_vec(&response)?; - json.push(b'\n'); - fdpass::reply(reader.get_ref(), &json, fd.as_ref().map(AsFd::as_fd)).await?; - // Dropping our copy leaves the client holding the only reference to - // its half, so its close tears the flow or the listener down. - drop(fd); + impl Connection { + /// Answer the commands on one client connection, in order, until it + /// closes or misbehaves. + /// + /// **The daemon keeps its copy of the descriptor in the last reply** + /// until the client's next command arrives or the connection ends. + /// Until the client reads the reply, the message carrying the + /// descriptor can be its only reference, and xnu's collector flushes a + /// socket in that state, as it does an unread arrival's. A client that + /// sends another command has read the reply first, unless it pipelined, + /// which the shipped client never does. Once the copy is gone the + /// client holds the only reference, so its close tears the flow or the + /// listener down; until then a close it makes waits for the copy. + pub(super) async fn run(mut self, stream: UnixStream) -> Result<(), std::io::Error> { + let mut reader = BufReader::new(stream); + let mut line = Vec::new(); + let mut kept: Option = None; + + while read_command(&mut reader, &mut line).await? { + drop(kept.take()); + let (response, fd) = self.answer(&line).await; + let mut json = serde_json::to_vec(&response)?; + json.push(b'\n'); + fdpass::reply(reader.get_ref(), &json, fd.as_ref().map(AsFd::as_fd)).await?; + kept = fd; + } + + // End of file, or an early return above: either way `kept` goes + // with this frame, and with it the last copy the daemon holds. + Ok(()) } - - Ok(()) } /// Read one newline-terminated command into `line`, refusing an oversized @@ -1365,6 +1496,47 @@ mod tests { sock } + /// How long a test waits for the daemon to notice that a client closed a + /// descriptor. + /// + /// On Linux the close wakes the task reading the daemon's half, so a wait + /// that will succeed does so within a few turns of the runtime. On macOS + /// and FreeBSD nothing wakes it: the reader sees the close only when its + /// bounded wait expires and it retries the read, as much as one + /// `CLOSE_RETRY` interval (`seqpacket.rs`) after the close, and a + /// listener's close that ends a flow's hold costs two of those in a row. + /// Counting turns of the runtime bounds nothing there: a thousand of them + /// ran out well inside one interval on a macOS runner. A wait that is + /// going to succeed still returns as soon as it does; this is how long one + /// that is not takes to say so. + pub(super) const CLOSE_WAIT: std::time::Duration = std::time::Duration::from_secs(5); + + // Four times the longest run of unnoticed closes a test waits through, so + // that lengthening `CLOSE_RETRY` cannot quietly use up the margin. + const _: () = assert!( + CLOSE_WAIT.as_millis() >= 8 * super::seqpacket::CLOSE_LATENCY.as_millis(), + "CLOSE_WAIT no longer covers two unnoticed closes four times over" + ); + + /// Retry `attempt` until it produces a value, for up to [`CLOSE_WAIT`]. + /// + /// `None` means the bound ran out. Attempts are a millisecond apart rather + /// than a yield apart, so a wait that lasts a quarter second on macOS is + /// not spent spinning, which for [`rebind`] would mean opening and closing + /// a socket pair on every turn. + pub(super) async fn eventually(mut attempt: impl AsyncFnMut() -> Option) -> Option { + tokio::time::timeout(CLOSE_WAIT, async { + loop { + if let Some(value) = attempt().await { + return value; + } + tokio::time::sleep(std::time::Duration::from_millis(1)).await; + } + }) + .await + .ok() + } + /// Send one command that opens a flow, returning the reply and descriptor. async fn open(connection: &mut Connection, line: &str) -> (serde_json::Value, StdUnixStream) { let (response, fd) = connection.answer(line.as_bytes()).await; @@ -1726,6 +1898,332 @@ mod tests { assert_eq!(&buf[..3], &[0x00, 0xff, 0x10]); } + /// Wait until a listener descriptor has an arrival to read. + /// + /// Polled without blocking and yielding in between, because the arrival is + /// written by the listener's task on this same runtime and a blocking wait + /// would stop the task it is waiting for. A readable listener means + /// `hand_over` has finished: it writes the arrival and settles what happens + /// to the daemon's copy of the descriptor in one synchronous step. + async fn readable(listener: &StdUnixStream) { + for _ in 0..1000 { + let mut poll = libc::pollfd { + fd: listener.as_raw_fd(), + events: libc::POLLIN, + revents: 0, + }; + // SAFETY: `poll` points at one live pollfd and the call cannot block. + let rc = unsafe { libc::poll(&mut poll, 1, 0) }; + if rc > 0 && (poll.revents & libc::POLLIN) != 0 { + return; + } + tokio::task::yield_now().await; + } + panic!("no arrival became readable on the listener"); + } + + /// Wait until the node no longer holds a flow, failing if it never lets go. + /// + /// The reader marks a flow closed before it gives the registry entry back + /// and forgets the flow, so `stats` can still find a closed flow for a few + /// turns of the runtime, and on macOS and FreeBSD the reader notices the + /// close itself only when it next retries. Bounded by [`CLOSE_WAIT`], so a + /// flow that is never released names itself. + async fn forgotten(connection: &mut Connection, flow: u64) { + let line = format!(r#"{{"command":"stats","params":{{"flow_id":{flow}}}}}"#); + let mut last = serde_json::Value::Null; + let gone = eventually(async || { + last = ask(connection, &line).await; + (last["data"]["errno"] == "ENOENT").then_some(()) + }) + .await; + assert!( + gone.is_some(), + "the node still holds flow {flow} after {CLOSE_WAIT:?}: {last}" + ); + } + + /// Bind a listener on a port a closed one held, retrying until the closed + /// one has given it back. + /// + /// The release reaches the registry from the closed listener's own task + /// once that task has noticed the close, which takes a few turns of the + /// runtime, or on macOS and FreeBSD up to a retry interval. Bounded by + /// [`CLOSE_WAIT`], so a port that is never given back fails here by name. + async fn rebind(connection: &mut Connection, port: u16) -> OwnedFd { + let line = format!(r#"{{"command":"listen","params":{{"local_port":{port}}}}}"#); + eventually(async || { + let (response, fd) = connection.answer(line.as_bytes()).await; + (serde_json::to_value(response).unwrap()["status"] == "ok") + .then(|| fd.expect("a rebound listener still gets a descriptor")) + }) + .await + .unwrap_or_else(|| { + panic!("the closed listener never gave port {port} back within {CLOSE_WAIT:?}") + }) + } + + /// Bind a listener on 4242, deliver one datagram to it from a new peer, + /// and accept the flow that announced, returning everything a test needs. + async fn arrive_and_accept( + connection: &mut Connection, + ) -> (StdUnixStream, u64, serde_json::Value, StdUnixStream) { + let (_value, listener) = listen(connection, 4242).await; + let value = ask(connection, &arrival(5000, 4242, "00ff10")).await; + assert_eq!(value["data"]["outcome"], "announced", "{value}"); + let flow = value["data"]["flow_id"].as_u64().unwrap(); + let (message, client) = accept(&listener); + (listener, flow, message, client) + } + + #[tokio::test] + async fn an_arrival_survives_a_kernel_collection_before_its_client_reads_it() { + // The macOS failure, made deterministic. Between the daemon writing the + // arrival and the client reading it, the flow's descriptor exists only + // inside the arrival message. xnu's descriptor collector flushes a + // socket in that state, so unless the daemon still holds its own copy, + // the client receives a flow that reads as end of file with its held + // datagram gone. On Linux the kernel keeps the socket either way, so + // this test can only fail on macOS. + let (mut connection, _outbound) = connect(); + let (_value, listener) = listen(&mut connection, 4242).await; + let value = ask(&mut connection, &arrival(5000, 4242, "00ff10")).await; + assert_eq!(value["data"]["outcome"], "announced", "{value}"); + + readable(&listener).await; + super::dgram_probe::provoke_collection(); + + let (message, mut client) = accept(&listener); + assert_eq!(message["held"], 1); + let mut buf = [0u8; 64]; + assert_eq!( + client.read(&mut buf).unwrap(), + 3, + "the accepted flow lost its held datagram to the kernel's collector" + ); + assert_eq!(&buf[..3], &[0x00, 0xff, 0x10]); + } + + #[tokio::test] + async fn a_handed_over_flow_outlives_its_dropped_descriptor_until_its_listener_closes() { + // The daemon keeps its copy of a handed-over descriptor until the + // client has shown it holds one, because dropping it is what exposes + // the descriptor to the macOS collector. On Linux the hold is + // observable this way: a client that closes the flow without ever + // writing does not end it while the listener that produced it is open. + let (mut connection, _outbound) = connect(); + let (listener, flow, _message, client) = arrive_and_accept(&mut connection).await; + + drop(client); + still_open( + &mut connection, + flow, + "the flow ended while the daemon should still hold its descriptor", + ) + .await; + + // Closing the listener ends the hold: whatever the client did with the + // arrival, the descriptor is no longer in flight in a live socket. + drop(listener); + forgotten(&mut connection, flow).await; + } + + #[tokio::test] + async fn a_handed_over_flow_closes_with_its_descriptor_once_its_client_has_written() { + // A datagram from the client proves it holds the descriptor, so the + // daemon lets its copy go and the client's close ends the flow at once, + // with the listener still open. + let (mut connection, _outbound) = connect(); + let (listener, flow, _message, mut client) = arrive_and_accept(&mut connection).await; + + client.write_all(b"x").unwrap(); + connection.settle(flow, 1).await; + drop(client); + forgotten(&mut connection, flow).await; + drop(listener); + } + + #[tokio::test] + async fn a_handed_over_flow_its_client_still_holds_outlives_its_listener_and_closes_with_its_descriptor() + { + // Closing the listener lets the daemon's copy go, and the client's own + // copy then carries the flow by itself: it keeps working after the + // listener has gone, and the client's close ends it with no write ever + // made. Letting the copy go must not end a flow the client still holds. + let (mut connection, _outbound) = connect(); + let (listener, flow, _message, mut client) = arrive_and_accept(&mut connection).await; + + drop(listener); + // The port comes back only after the listener's task has let its + // flows' copies go, so a rebound port means that has happened. + let _rebound = rebind(&mut connection, 4242).await; + + let mut buf = [0u8; 64]; + assert_eq!(client.read(&mut buf).unwrap(), 3); + let value = ask( + &mut connection, + &format!(r#"{{"command":"inject","params":{{"flow_id":{flow},"data":"ab"}}}}"#), + ) + .await; + assert_eq!(value["status"], "ok", "{value}"); + assert_eq!(client.read(&mut buf).unwrap(), 1); + assert_eq!(buf[0], 0xab); + + drop(client); + forgotten(&mut connection, flow).await; + } + + /// Serve a sibling of `connection` over a real socket, the way the daemon + /// serves a client, returning the client's end and the serving task. + /// + /// Through `run` rather than `answer`, because what these tests observe is + /// what the serving loop does with the descriptor in a reply it has sent. + fn serve_socket( + connection: &Connection, + ) -> (StdUnixStream, tokio::task::JoinHandle>) { + let (ours, theirs) = StdUnixStream::pair().expect("AF_UNIX socketpair"); + ours.set_nonblocking(true).expect("the socket is open"); + let ours = tokio::net::UnixStream::from_std(ours).expect("inside a runtime"); + let task = tokio::spawn(connection.sibling().run(ours)); + (bounded(theirs), task) + } + + /// Write one command on a client socket and read its reply line, with the + /// descriptor it carried. + async fn call(client: &StdUnixStream, line: &str) -> (serde_json::Value, Option) { + let mut writer = client; + writer.write_all(line.as_bytes()).unwrap(); + writer.write_all(b"\n").unwrap(); + readable(client).await; + let mut buf = [0u8; 4096]; + let chunk = super::fdpass::recv(client.as_raw_fd(), &mut buf) + .expect("a reply should be readable on the connection"); + let reply = buf[..chunk.len] + .strip_suffix(b"\n") + .expect("one whole reply line per read"); + (serde_json::from_slice(reply).unwrap(), chunk.fd) + } + + /// Open a flow through a served socket, returning its id and descriptor. + async fn connect_over(client: &StdUnixStream) -> (u64, OwnedFd) { + let line = format!( + r#"{{"command":"connect","params":{{"peer":"{PEER}","remote_port":4242,"local_port":4243}}}}"# + ); + let (value, fd) = call(client, &line).await; + assert_eq!(value["status"], "ok", "{value}"); + let flow = value["data"]["flow_id"].as_u64().unwrap(); + ( + flow, + fd.expect("a connect reply carries the flow's descriptor"), + ) + } + + /// Assert that the node still holds a flow open, after giving its reader + /// every chance to notice a close. + /// + /// Where a close wakes the reader, a few turns of the runtime are that + /// chance. On macOS and FreeBSD the reader notices a close only when its + /// bounded wait expires and it retries, so a check made sooner passes + /// whether or not the flow has closed. There this also waits out two of + /// those intervals: a reader that parked just before a close has retried + /// by then, with one interval to spare for a loaded runner. Elsewhere the + /// interval is zero and the wait costs nothing. + async fn still_open(connection: &mut Connection, flow: u64, why: &str) { + for _ in 0..1000 { + tokio::task::yield_now().await; + } + tokio::time::sleep(super::seqpacket::CLOSE_LATENCY * 2).await; + let value = ask( + connection, + &format!(r#"{{"command":"stats","params":{{"flow_id":{flow}}}}}"#), + ) + .await; + assert_eq!(value["status"], "ok", "{why}: {value}"); + assert_eq!(value["data"]["closed"], false, "{why}: {value}"); + } + + /// Any command at all, sent only so the serving loop reads one. + const NEXT: &str = r#"{"command":"stats","params":{"flow_id":0}}"#; + + #[tokio::test] + async fn a_descriptor_sent_in_a_reply_is_kept_until_the_clients_next_command() { + // Until the client reads a reply, the message carrying its descriptor + // can be the only reference to it, which is what xnu's collector + // flushes. So the daemon keeps its copy until the client sends another + // command, which it does only after reading the reply. On Linux the + // hold is observable as a flow that outlives the client's close. + let (mut probe, _outbound) = connect(); + let (client, _task) = serve_socket(&probe); + let (flow, fd) = connect_over(&client).await; + + drop(fd); + still_open( + &mut probe, + flow, + "the flow ended while the connection should still hold its descriptor", + ) + .await; + + call(&client, NEXT).await; + forgotten(&mut probe, flow).await; + } + + #[tokio::test] + async fn a_descriptor_sent_in_a_reply_is_let_go_when_the_connection_ends() { + // The connection's close ends its hold. The client's own copy then + // carries the flow alone, so the flow survives the connection and ends + // with the client's close. + let (mut probe, _outbound) = connect(); + let (client, task) = serve_socket(&probe); + let (flow, fd) = connect_over(&client).await; + + drop(client); + task.await + .unwrap() + .expect("a connection closed between commands ends cleanly"); + still_open( + &mut probe, + flow, + "the flow ended with the connection while the client still holds it", + ) + .await; + + drop(fd); + forgotten(&mut probe, flow).await; + } + + #[tokio::test] + async fn a_flow_its_client_closes_after_its_next_command_ends_at_once() { + // Once the client has sent another command the hold is over, so the + // client's close is the flow's end of file with the connection still + // open. + let (mut probe, _outbound) = connect(); + let (client, _task) = serve_socket(&probe); + let (flow, fd) = connect_over(&client).await; + + call(&client, NEXT).await; + drop(fd); + forgotten(&mut probe, flow).await; + drop(client); + } + + #[tokio::test] + async fn a_flow_whose_reply_is_never_followed_by_a_command_ends_with_the_connection() { + // A client that closes the flow and then the connection, sending + // nothing more, still ends the flow: the connection's end of file is + // the last point at which the daemon lets its copy go. + let (mut probe, _outbound) = connect(); + let (client, task) = serve_socket(&probe); + let (flow, fd) = connect_over(&client).await; + + drop(fd); + drop(client); + task.await + .unwrap() + .expect("a connection closed between commands ends cleanly"); + forgotten(&mut probe, flow).await; + } + #[test] fn a_listener_that_closed_and_one_that_stopped_reading_are_counted_apart() { // Both take the same cleanup, so the counter is the only place the @@ -1946,19 +2444,8 @@ mod tests { // The unbind is what `close(listen_fd)` means in Berkeley, and the // daemon can only observe it by reading its own half. Without the read - // arm the port is held for the node's lifetime and every attempt below - // fails. - for _ in 0..1000 { - let (response, fd) = connection - .answer(br#"{"command":"listen","params":{"local_port":4242}}"#) - .await; - if serde_json::to_value(response).unwrap()["status"] == "ok" { - assert!(fd.is_some(), "a rebound listener still gets a descriptor"); - return; - } - tokio::task::yield_now().await; - } - panic!("the closed listener never gave its port back"); + // arm the port is held for the node's lifetime and every attempt fails. + rebind(&mut connection, 4242).await; } #[tokio::test] diff --git a/src/native/seqpacket.rs b/src/native/seqpacket.rs index d2537a27..a9bd60b8 100644 --- a/src/native/seqpacket.rs +++ b/src/native/seqpacket.rs @@ -180,6 +180,19 @@ pub fn set_sndbuf(fd: &OwnedFd, bytes: usize) -> io::Result<()> { #[cfg(any(target_os = "macos", target_os = "freebsd"))] const CLOSE_RETRY: Duration = Duration::from_millis(250); +/// How long a reader can leave a closed peer unnoticed, for the native API +/// tests that assert a flow is still open and so must wait long enough to have +/// seen it close. +/// +/// `CLOSE_RETRY` where the reactor cannot see a close, because the reader then +/// notices one only when that bound expires and it retries the read. +#[cfg(all(test, any(target_os = "macos", target_os = "freebsd")))] +pub(super) const CLOSE_LATENCY: Duration = CLOSE_RETRY; + +/// Zero where a close wakes the reader itself. +#[cfg(all(test, not(any(target_os = "macos", target_os = "freebsd"))))] +pub(super) const CLOSE_LATENCY: Duration = Duration::ZERO; + /// The daemon's half of a flow's or a listener's socket pair, registered with /// the reactor. pub struct Seqpacket { diff --git a/testing/native-api/client.py b/testing/native-api/client.py index f5f41af8..221a0e0c 100755 --- a/testing/native-api/client.py +++ b/testing/native-api/client.py @@ -7,7 +7,9 @@ connection owns nothing — a flow lives until its own descriptor is closed, and listener until its own is — so the single connection is a convenience for the checks rather than a lifetime the daemon respects. Descriptors are what keep things alive, and this tool holds them until the step that closes them or until -it exits. +it exits. The daemon does keep its own copy of the descriptor in its last reply +until the next command arrives, so a check that closes one must send a command +before it expects the close to have taken effect. Kinds of step: diff --git a/testing/native-api/test.sh b/testing/native-api/test.sh index 2badf346..7470324c 100755 --- a/testing/native-api/test.sh +++ b/testing/native-api/test.sh @@ -695,12 +695,20 @@ check_backlog_is_not_the_clients_bound() { } check_refusing_a_flow_frees_it() { - log "Refusing an accepted flow is closing its descriptor, and that frees it" + log "Refusing an accepted flow is closing its descriptor, and its listener's close frees it" # There is no reject command: Berkeley has exactly one way to refuse a # connection and so does this. Keeping a second would let a client refuse a # flow two indistinguishable ways. # - # Two things have to follow the close, and the second is what makes this + # The daemon keeps its own copy of an accepted flow's descriptor until the + # client writes on the flow or closes the listener, because on macOS the + # kernel can destroy a socket whose descriptor is still in an unread + # arrival. So a flow refused without a write outlives its close, and goes + # when the listener does. Both halves are asserted: a flow freed at the + # close would mean the daemon let its copy go early, and one that outlived + # the listener would be a leak. + # + # Two things have to follow the release, and the second is what makes this # more than a repeat of the connected-flow close: the node forgets the flow, # and the registry entry and its key go with it, so the very same key # announces a new flow afterwards rather than delivering into the dead one. @@ -713,16 +721,22 @@ check_refusing_a_flow_frees_it() { {"command":"stats","params":{"flow_id":"@a"},"expect":{"status":"ok"}}, {"fd":"a","close":true}, {"sleep":1}, - {"command":"stats","params":{"flow_id":"@a"},"expect":{"status":"error"}}, + {"command":"stats","params":{"flow_id":"@a"}, + "expect":{"status":"ok","data.closed":false}}, + {"fd":"L","close":true}, + {"command":"stats","params":{"flow_id":"@a"},"settle":true, + "expect":{"status":"error"}}, + {"command":"listen","params":{"local_port":4304},"keep_listener":"M", + "expect":{"status":"ok"}}, {"command":"arrive","params":{"peer":"'"$PEER"'","src_port":5000,"dst_port":4304,"data":"bb"}, "expect":{"status":"ok","data.outcome":"announced"}}, - {"accept":"L","keep_fd":"b","expect":{"local_port":4304,"remote_port":5000}}, + {"accept":"M","keep_fd":"b","expect":{"local_port":4304,"remote_port":5000}}, {"fd":"b","read":1,"expect_bytes":"bb"} ]' if run_client "$script"; then - pass "a refused flow is gone and its key is free to arrive again" + pass "a refused flow lasts until its listener closes, then is gone and its key is free to arrive again" else - fail "closing a refused flow did not release it" + fail "a refused flow did not last until its listener closed, or was not released then" fi } From 888d393f8be25021ae687c71ba810ccfe23f0b36 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 1 Oct 2026 14:26:25 +0000 Subject: [PATCH 08/23] Report inputs the glibc floor check could not examine as unchecked check-glibc-floor.sh counted a missing path, a file that is not an ELF object and a .deb with no ELF executables as floor failures, so each exited 1 with "A binary above the floor installs cleanly and then fails to start" and the advice to rebuild in the container. Passing a release tarball, which is the usual mistake, therefore read as a broken release when nothing had been examined. A .deb that dpkg-deb could not unpack killed the script under set -e with no summary, leaving later arguments unexamined. Those inputs are now recorded as "could not check" with the reason, and listed at the end. A tarball's reason says to unpack it and pass its binaries. Exit status is 1 only when a binary is above the floor (still listing anything unchecked), 2 when anything could not be examined, and 0 only when every input was examined and passed. Neither caller branches on the status, and both stop on 1 or 2 as before. Two setup failures also exited 1: a missing packaging/build-floor.env ended the run through set -e when it was sourced, and an empty FIPS_GLIBC_FLOOR ended it through the ":?" expansion. The env file is now read only after a readable-file check, and an empty floor is caught by an explicit test; both exit 2 with a message saying there is no floor to check against. The file's shellcheck source directive now resolves from the script's directory, so shellcheck -x run from the repository root no longer reports SC1091. The floor check runs only at release time, so a regression in how it reports would surface on a release day. testing/glibc-floor/test.sh builds its inputs at run time from the host's own true executable (a text file, a tarball holding the binary, a missing path, an unreadable file, a corrupt .deb, a .deb holding only a script, and a .deb holding the binary), sets the floor explicitly in every case, and asserts the exit status, the "could not check" listing and the presence or absence of the rebuild advice. Scratch copies of the check cover a missing build-floor.env, an empty floor, a file that declares none, and one whose file supplies a floor. The test exits 2 rather than passing when readelf, dpkg-deb or a dynamic true is missing. ci-local.sh runs it beside the Debian version check, and ci.yml's static job runs the same script. --- .github/workflows/ci.yml | 5 + testing/check-glibc-floor.sh | 108 +++++++++++++++---- testing/ci-local.sh | 12 +++ testing/glibc-floor/test.sh | 195 +++++++++++++++++++++++++++++++++++ 4 files changed, 298 insertions(+), 22 deletions(-) create mode 100755 testing/glibc-floor/test.sh diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index daf8c116..c57cf625 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -106,6 +106,11 @@ jobs: # reports. Kept in step with ci-local.sh's run_nextest_flaky by hand. - name: Check the flaky-test reporter against its fixtures run: bash testing/nextest-flaky/test.sh + # The glibc floor check's cases, built from the host's own true + # executable. Not a matrix suite, so it is kept in step with + # ci-local.sh's run_glibc_floor by hand. + - name: Check the glibc floor check against its cases + run: bash testing/glibc-floor/test.sh fmt: name: Format check diff --git a/testing/check-glibc-floor.sh b/testing/check-glibc-floor.sh index 1ca81deb..cdce9506 100755 --- a/testing/check-glibc-floor.sh +++ b/testing/check-glibc-floor.sh @@ -16,17 +16,37 @@ # Anything else is treated as a single ELF binary. # # Reads the floor from packaging/build-floor.env unless FIPS_GLIBC_FLOOR is set. +# +# Exit 0 = every input was examined and none is above the floor. Exit 1 = a +# binary needs a newer glibc than the floor. Exit 2 = an input could not be +# examined (missing, not an ELF object, a .deb that would not unpack or holds +# no binaries), or the check could not run at all; never treated as a pass, +# and never reported as a binary above the floor. set -euo pipefail SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd)" REPO_ROOT="$(cd "$SCRIPT_DIR/.." && pwd)" +# A missing or empty floor means the check cannot run, so it exits 2 here +# rather than letting `set -e` or a `:?` expansion end the run with status 1, +# which callers would read as a binary above the floor. +FLOOR_ENV="$REPO_ROOT/packaging/build-floor.env" if [ -z "${FIPS_GLIBC_FLOOR:-}" ]; then - # shellcheck source=../packaging/build-floor.env - . "$REPO_ROOT/packaging/build-floor.env" + [ -r "$FLOOR_ENV" ] || { + echo "check-glibc-floor: cannot read $FLOOR_ENV and FIPS_GLIBC_FLOOR is not set;" >&2 + echo " there is no floor to check against." >&2 + exit 2 + } + # shellcheck source-path=SCRIPTDIR source=../packaging/build-floor.env + . "$FLOOR_ENV" fi -FLOOR="${FIPS_GLIBC_FLOOR:?no floor declared}" +if [ -z "${FIPS_GLIBC_FLOOR:-}" ]; then + echo "check-glibc-floor: $FLOOR_ENV declares no FIPS_GLIBC_FLOOR;" >&2 + echo " there is no floor to check against." >&2 + exit 2 +fi +FLOOR="$FIPS_GLIBC_FLOOR" for tool in readelf dpkg dpkg-deb; do command -v "$tool" >/dev/null 2>&1 || { @@ -65,6 +85,18 @@ max_glibc_need() { FAILED=0 CHECKED=0 +UNCHECKED=0 +UNCHECKED_LIST=() + +# unchecked