From 9b04056b93ed2bb59607403c486fef5f4d697960 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 17 Sep 2026 19:59:24 +0000 Subject: [PATCH 1/6] fix(test/ci): compile the library tests in release mode, and keep them compiling `Node::debug_assert_peer_maps_coherent` exists only under `#[cfg(debug_assertions)]`, and the handshake-presence test called it without that gate, so `cargo test --release --lib` failed to build with E0599 while every sibling call site was gated correctly. Gate the call rather than the whole test, so the rest of its assertions still run under `--release`. Nothing noticed, because no runner ever compiled the test target with optimisations: local CI builds release binaries but tests in debug, and the GitHub unit-test job does the same. Add `cargo test --release --lib --no-run` to both runners so the next debug-only helper used without its gate fails at once instead of waiting for someone to run the suite in release. --- .github/workflows/ci.yml | 8 ++++++++ src/node/tests/unit.rs | 5 ++++- testing/ci-local.sh | 14 ++++++++++++++ 3 files changed, 26 insertions(+), 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 0fb4a7dc..f9bb2c6b 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -330,6 +330,14 @@ jobs: - name: Run library tests with the tick-body profiler enabled run: cargo test --lib --features profiling + # Debug-only helpers (anything behind #[cfg(debug_assertions)]) vanish in + # a release build, so a test calling one without the same gate breaks a + # build no other job performs: every run above compiles the test target + # in debug. Compile it in release too, without running it — the point is + # that it builds at all. Mirrored in testing/ci-local.sh. + - name: Compile the library tests in release mode + run: cargo test --release --lib --no-run + # ───────────────────────────────────────────────────────────────────────────── # Job 2b – Unit tests (macOS) # ───────────────────────────────────────────────────────────────────────────── diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index 12b0ff7b..5bc4ca87 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -3977,7 +3977,10 @@ fn handshake_presence_tracks_the_carrier_not_the_noise_handles() { "rekey-msg2 discriminator: {when}" ); // Fires the live-carrier coherence assertion; a machine that had gone - // invisible would panic here rather than fail an assert_eq above. + // invisible would panic here rather than fail an assert_eq above. The + // helper only exists in debug builds, so the rest of this test carries + // on without it under `--release`. + #[cfg(debug_assertions)] node.debug_assert_peer_maps_coherent(); }; diff --git a/testing/ci-local.sh b/testing/ci-local.sh index 9b097be5..5a4b6a4f 100755 --- a/testing/ci-local.sh +++ b/testing/ci-local.sh @@ -603,6 +603,20 @@ run_tests() { else record "unit-tests-profiling" 1 fi + + # Debug-only helpers (anything behind #[cfg(debug_assertions)]) vanish in a + # release build, so a test calling one without the same gate breaks a build + # nothing here ever performs: every run above compiles the test target in + # debug. Compile it in release too, without running it — the point is that + # it builds at all. Mirrored in .github/workflows/ci.yml; check-ci-parity.sh + # compares integration suites only and would not catch a stage added to one + # runner and not the other. + info "cargo test --release --lib --no-run" + if cargo test --release --lib --no-run 2>&1; then + record "release-test-compile" 0 + else + record "release-test-compile" 1 + fi } # ── Stage 3: Integration Tests ───────────────────────────────────────────── From b459ca029bfa36b6bc852194a18fed3dcae21141 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 17 Sep 2026 19:55:18 +0000 Subject: [PATCH 2/6] test(nat): dump diagnostics when a NAT lab ping fails A failed data-plane ping was the one step in the NAT lab that reported nothing. ping_peer sent ping6's output to /dev/null and the three scenarios called it bare, so under `set -euo pipefail` a failure killed the script with no message, without running the scenario's diagnostics dump, and without tearing down; the next scenario's opening cleanup then removed the failed containers and their logs. Two nat-symmetric failures on the builder logged only the peer counts and FAIL, and the cause had to be found by reproducing it locally. ping_peer now captures ping6's output, names the source container, the destination and the exit status on failure, and returns non-zero, matching the assertion helpers above it. The six call sites wrap it in the same `|| { dump_*_diagnostics; return 1; }` shape every other step uses. The healthy path stays silent and the pass/fail verdict is unchanged. --- testing/nat/scripts/nat-test.sh | 45 ++++++++++++++++++++++++++++----- 1 file changed, 38 insertions(+), 7 deletions(-) diff --git a/testing/nat/scripts/nat-test.sh b/testing/nat/scripts/nat-test.sh index 284edd57..cbf7c7dc 100755 --- a/testing/nat/scripts/nat-test.sh +++ b/testing/nat/scripts/nat-test.sh @@ -415,10 +415,23 @@ require_bootstrap_activity() { fi } +# Data-plane assertion, and the last step of every scenario. Like the two path +# assertions above it, it stays silent on success and names the container, the +# destination and what it read on failure; ping's own output is the only record +# of whether resolution, routing or the data path broke, so it is captured +# rather than discarded to /dev/null. Callers wrap it in the same +# `|| { dump_*_diagnostics; return 1; }` shape, because a bare call lets `set -e` +# tear the script down before any diagnostics run. ping_peer() { local container="$1" local npub="$2" - docker exec "$container" ping6 -c 3 -W 5 "${npub}.fips" >/dev/null + local output rc=0 + output="$(docker exec "$container" ping6 -c 3 -W 5 "${npub}.fips" 2>&1)" || rc=$? + if [ "$rc" != 0 ]; then + echo "PING FAIL: $container -> ${npub}.fips: ping6 exited ${rc}:" >&2 + printf '%s\n' "$output" >&2 + return 1 + fi } run_cone() { @@ -453,8 +466,14 @@ run_cone() { } # shellcheck disable=SC1090 source "$CONFIG_DIR/cone/npubs.env" - ping_peer fips-nat-cone-a${FIPS_CI_NAME_SUFFIX:-} "$NPUB_B" - ping_peer fips-nat-cone-b${FIPS_CI_NAME_SUFFIX:-} "$NPUB_A" + ping_peer fips-nat-cone-a${FIPS_CI_NAME_SUFFIX:-} "$NPUB_B" || { + dump_cone_diagnostics + return 1 + } + ping_peer fips-nat-cone-b${FIPS_CI_NAME_SUFFIX:-} "$NPUB_A" || { + dump_cone_diagnostics + return 1 + } cleanup } @@ -492,8 +511,14 @@ run_symmetric() { require_bootstrap_activity fips-nat-symmetric-b${FIPS_CI_NAME_SUFFIX:-} # shellcheck disable=SC1090 source "$CONFIG_DIR/symmetric/npubs.env" - ping_peer fips-nat-symmetric-a${FIPS_CI_NAME_SUFFIX:-} "$NPUB_B" - ping_peer fips-nat-symmetric-b${FIPS_CI_NAME_SUFFIX:-} "$NPUB_A" + ping_peer fips-nat-symmetric-a${FIPS_CI_NAME_SUFFIX:-} "$NPUB_B" || { + dump_symmetric_diagnostics + return 1 + } + ping_peer fips-nat-symmetric-b${FIPS_CI_NAME_SUFFIX:-} "$NPUB_A" || { + dump_symmetric_diagnostics + return 1 + } cleanup } @@ -528,8 +553,14 @@ run_lan() { } # shellcheck disable=SC1090 source "$CONFIG_DIR/lan/npubs.env" - ping_peer fips-nat-lan-a${FIPS_CI_NAME_SUFFIX:-} "$NPUB_B" - ping_peer fips-nat-lan-b${FIPS_CI_NAME_SUFFIX:-} "$NPUB_A" + ping_peer fips-nat-lan-a${FIPS_CI_NAME_SUFFIX:-} "$NPUB_B" || { + dump_lan_diagnostics + return 1 + } + ping_peer fips-nat-lan-b${FIPS_CI_NAME_SUFFIX:-} "$NPUB_A" || { + dump_lan_diagnostics + return 1 + } # Skip the final teardown when the mesh-lab harness wraps this # script: it needs to docker-logs the containers before teardown, # and will run its own cleanup after capture. Failure paths above From 51353faaaac93384be2b0e741f6b9acc29c5d62f Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 17 Sep 2026 20:00:34 +0000 Subject: [PATCH 3/6] fix(config): abort a persistent start when the identity key path cannot be examined `Path::exists` reports false both for a key file that is absent and for one whose metadata cannot be read, so a key symlinked onto a volume that did not mount, or one in a directory the daemon cannot search, read as a first boot. The node then generated a fresh identity, failed to store it, and carried on under an npub that every peer whose allowlist names the old one refuses. Presence is now decided by `symlink_metadata`, where only `NotFound` counts as an absence. A dangling symlink reads as present, so the read that follows aborts the start with the path named, and any other lookup failure aborts directly. The legacy system-directory fallback applies the same test. A key file that is genuinely absent still generates and persists a new identity on first boot. --- CHANGELOG.md | 13 ++++ src/config/mod.rs | 169 ++++++++++++++++++++++++++++++++++++++++++---- 2 files changed, 169 insertions(+), 13 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2555bb44..eee87df4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -79,6 +79,19 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 values the open-discovery tutorial described never occurred, and the tutorial no longer lists them. The response shape is unchanged. +#### Identity & config + +- A persistent node whose identity key path cannot be examined now refuses to + start instead of coming up under a new identity. `Path::exists` reports false + both for a key that is absent and for one whose metadata cannot be read, so a + key symlinked onto a volume that did not mount, or one in a directory the + daemon cannot search, read as a first boot: the node generated a fresh + identity, failed to store it, and carried on under an npub that every peer + whose allowlist names the old one refuses. Only a `NotFound` result is now + treated as an absence; any other failure to stat the path aborts the start and + names the path. A dangling symlink likewise aborts rather than being replaced. + The legacy `/etc/fips/fips.key` lookup follows the same rule. + ### Changed - The lockfile moves `chacha20` from 0.10.1 to 0.10.2, because 0.10.1 is yanked. diff --git a/src/config/mod.rs b/src/config/mod.rs index 67491aa6..cd26653e 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -120,15 +120,44 @@ const LEGACY_SYSTEM_CONFIG_DIR: &str = "/etc/fips"; /// - `key_path` sits in `system_dir`, so an operator using `./fips.yaml` or a /// user config is never redirected to a system key /// - a key does exist at `legacy_dir` -fn legacy_key_fallback(key_path: &Path, system_dir: &Path, legacy_dir: &Path) -> Option { - if system_dir == legacy_dir || key_path.exists() { - return None; +/// +/// Returns an error when either location cannot be examined, which the caller +/// aborts on: a lookup that failed is not evidence that no key is there. +fn legacy_key_fallback( + key_path: &Path, + system_dir: &Path, + legacy_dir: &Path, +) -> Result, ConfigError> { + if system_dir == legacy_dir || key_file_present(key_path)? { + return Ok(None); } if key_path.parent() != Some(system_dir) { - return None; + return Ok(None); } let legacy = legacy_dir.join(KEY_FILENAME); - legacy.exists().then_some(legacy) + Ok(key_file_present(&legacy)?.then_some(legacy)) +} + +/// Report whether an identity key file is present, distinguishing a genuine +/// absence from a lookup that could not be made. +/// +/// `Path::exists` answers false to both, which is what a persistent start must +/// not do: a key symlinked onto a volume that did not mount, or one in a +/// directory the daemon may not search, would read as a first boot and the +/// node would generate and run under a new identity that every peer +/// allowlisting its old npub refuses. `symlink_metadata` reports a symlink +/// itself as present, so the read that follows fails and aborts the start, +/// and any other lookup error is returned for the caller to abort on. Only +/// `NotFound` is an absence, which is the first-boot case. +fn key_file_present(path: &Path) -> Result { + match path.symlink_metadata() { + Ok(_) => Ok(true), + Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(false), + Err(e) => Err(ConfigError::KeyPathUnreadable { + path: path.to_path_buf(), + source: e, + }), + } } /// Derive the public key file path from a config file path. @@ -508,6 +537,10 @@ pub fn write_pub_file(path: &Path, npub: &str) -> Result<(), ConfigError> { /// 2. Persistent key file (`fips.key`) — reused across restarts /// 3. Generate new — creates keypair, writes `fips.key` and `fips.pub` /// +/// A key file that exists but cannot be read, including one whose metadata +/// the daemon cannot look up at all, aborts the start. Only a key file that +/// is genuinely absent reaches step 3. +/// /// - **`nsec` set explicitly**: always uses that, regardless of `persistent`. /// /// Returns the nsec string (bech32 or hex) to be used for identity creation. @@ -538,8 +571,10 @@ pub fn resolve_identity( let pub_path = pub_file_path(&config_ref); if config.node.identity.persistent { - // Persistent mode: load existing key file or generate-and-persist - if key_path.exists() { + // Persistent mode: load existing key file or generate-and-persist. + // A key path the daemon cannot examine aborts the start here rather + // than falling through to generation. + if key_file_present(&key_path)? { // Held in a guard, not a bare `String`: if the parse below fails, // the `?` returns and a bare local would be freed uncleared. let nsec = Zeroizing::new(read_key_file(&key_path)?); @@ -566,7 +601,7 @@ pub fn resolve_identity( &key_path, Path::new(SYSTEM_CONFIG_DIR), Path::new(LEGACY_SYSTEM_CONFIG_DIR), - ) { + )? { // Guarded for the same reason as the current-path read above. let nsec = Zeroizing::new(read_key_file(&legacy)?); let identity = Identity::from_secret_str(&nsec)?; @@ -746,6 +781,12 @@ pub enum ConfigError { #[error("refusing to write key file through a symlink: {path}")] KeyPathIsSymlink { path: PathBuf }, + #[error("cannot determine whether the identity key file {path} exists: {source}")] + KeyPathUnreadable { + path: PathBuf, + source: std::io::Error, + }, + #[error("identity error: {0}")] Identity(#[from] IdentityError), @@ -1629,7 +1670,7 @@ node: let key_path = system.join(KEY_FILENAME); assert_eq!( - legacy_key_fallback(&key_path, &system, &legacy), + legacy_key_fallback(&key_path, &system, &legacy).unwrap(), Some(legacy_key), "a key stranded at the legacy path must be adopted, not regenerated" ); @@ -1643,7 +1684,10 @@ node: write_stub_key(&legacy); let key_path = write_stub_key(&system); - assert_eq!(legacy_key_fallback(&key_path, &system, &legacy), None); + assert_eq!( + legacy_key_fallback(&key_path, &system, &legacy).unwrap(), + None + ); } #[test] @@ -1655,7 +1699,7 @@ node: write_stub_key(&dir); let absent = dir.join("nonexistent").join(KEY_FILENAME); - assert_eq!(legacy_key_fallback(&absent, &dir, &dir), None); + assert_eq!(legacy_key_fallback(&absent, &dir, &dir).unwrap(), None); } #[test] @@ -1670,7 +1714,10 @@ node: write_stub_key(&legacy); let key_path = elsewhere.join(KEY_FILENAME); - assert_eq!(legacy_key_fallback(&key_path, &system, &legacy), None); + assert_eq!( + legacy_key_fallback(&key_path, &system, &legacy).unwrap(), + None + ); } #[test] @@ -1682,7 +1729,36 @@ node: std::fs::create_dir_all(&system).unwrap(); let key_path = system.join(KEY_FILENAME); - assert_eq!(legacy_key_fallback(&key_path, &system, &legacy), None); + assert_eq!( + legacy_key_fallback(&key_path, &system, &legacy).unwrap(), + None + ); + } + + #[cfg(unix)] + #[test] + fn a_legacy_key_whose_metadata_cannot_be_read_is_present_not_absent() { + // A dangling symlink is the case that matters in the field: a key + // symlinked onto a volume that did not mount. Reading it as an + // absence sends a persistent node on to generate a new identity. + let root = TempDir::new().unwrap(); + let legacy = root.path().join("etc/fips"); + let system = root.path().join("usr/local/etc/fips"); + std::fs::create_dir_all(&legacy).unwrap(); + std::fs::create_dir_all(&system).unwrap(); + let legacy_key = legacy.join(KEY_FILENAME); + std::os::unix::fs::symlink( + root.path().join("unmounted").join(KEY_FILENAME), + &legacy_key, + ) + .unwrap(); + + let key_path = system.join(KEY_FILENAME); + assert_eq!( + legacy_key_fallback(&key_path, &system, &legacy).unwrap(), + Some(legacy_key), + "a legacy key the daemon cannot stat must be reported present, so the read aborts" + ); } #[test] @@ -2055,6 +2131,73 @@ node: assert_eq!(resolved.nsec, resolved2.nsec); } + #[cfg(unix)] + #[test] + fn persistent_start_aborts_when_the_key_path_is_a_dangling_symlink() { + // The key is symlinked onto a volume that did not mount. The node + // must not read that as a first boot and take a new identity, which + // every peer whose allowlist names the old npub would then refuse. + let temp_dir = TempDir::new().unwrap(); + let config_path = temp_dir.path().join("fips.yaml"); + let key_path = temp_dir.path().join("fips.key"); + let unmounted = temp_dir.path().join("unmounted").join("fips.key"); + + fs::write(&config_path, "node:\n identity:\n persistent: true\n").unwrap(); + std::os::unix::fs::symlink(&unmounted, &key_path).unwrap(); + + let config = Config::load_file(&config_path).unwrap(); + // `ResolvedIdentity` carries the secret and has no `Debug`, so the + // failure is matched rather than unwrapped. + let Err(err) = resolve_identity(&config, std::slice::from_ref(&config_path)) else { + panic!("a key path that cannot be read must abort the start, not generate a new key"); + }; + + assert!( + err.to_string().contains(&key_path.display().to_string()), + "the diagnostic must name the key path, got {err}" + ); + assert!( + key_path + .symlink_metadata() + .unwrap() + .file_type() + .is_symlink(), + "the symlink itself must be left in place" + ); + assert!( + !unmounted.exists(), + "nothing may be written through the symlink" + ); + assert!( + !temp_dir.path().join("fips.pub").exists(), + "an aborted start writes neither key file" + ); + } + + #[cfg(unix)] + #[test] + fn persistent_start_aborts_when_the_key_path_cannot_be_examined() { + // A key path whose parent is not a directory fails the lookup with an + // error that is not an absence, the same shape as a directory the + // daemon may not search, and unlike a permission case it behaves the + // same for root. + let temp_dir = TempDir::new().unwrap(); + let blocked = temp_dir.path().join("blocked"); + fs::write(&blocked, "not a directory\n").unwrap(); + let config_path = blocked.join("fips.yaml"); + + let mut config = Config::new(); + config.node.identity.persistent = true; + + let Err(err) = resolve_identity(&config, std::slice::from_ref(&config_path)) else { + panic!("a key path that cannot be examined must abort the start"); + }; + assert!( + err.to_string().contains(&blocked.display().to_string()), + "the diagnostic must name the key path, got {err}" + ); + } + #[test] fn test_to_yaml_empty_nsec_omitted() { let config = Config::new(); From 0f7ac05fe8c957419b4f592c883492a06cd4ec17 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 17 Sep 2026 20:07:46 +0000 Subject: [PATCH 4/6] fix(transport/tcp): key inbound pool entries by the connection four-tuple The kernel names a TCP connection by its four-tuple, but the connection pool keyed every entry by the remote address alone. A listener on a wildcard address, which is what the shipped configuration binds, accepts two connections with the same peer ip:port when they arrive on two different local addresses, and those two shared one pool entry: the second insert replaced the first entry and left its tasks running, the inbound-connection counter counted both, the first connection's teardown removed the second's entry, and the second's own teardown found nothing to remove. The counter therefore ended one above the connections it counts, and since it gates the inbound connection limit, a host repeating the collision could hold it at the limit and lock out further inbound connections until the daemon restarted. Inbound entries now carry the accepted socket's local address in their pool key, so two such connections get two entries and each receive loop tears down only the connection it was spawned for. An accepted socket whose local address cannot be read is dropped rather than pooled, since there is no key for it that is sure not to collide. Outbound entries keep the remote address alone: connect-on-send makes at most one connection per peer, so nothing would distinguish them, and the lookup stays a single hash probe. A caller answering a packet that arrived on an inbound connection knows only its remote address, so that lookup now resolves such an address to the matching entry rather than falling through to connect-on-send against the peer's ephemeral port. --- CHANGELOG.md | 13 ++ src/transport/tcp/mod.rs | 275 +++++++++++++++++++++++++++++++++----- src/transport/tcp/pool.rs | 65 ++++++++- 3 files changed, 320 insertions(+), 33 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index eee87df4..5ebedb74 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,19 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 #### Data-plane / transports +- Two inbound TCP connections that share a peer address but arrive on different + local addresses no longer share one pool entry. The kernel names a connection + by its four-tuple, so a listener on a wildcard address, which is what the + shipped configuration binds, can accept two connections whose peer `ip:port` + is the same on two different local addresses. The pool was keyed by the peer + address alone: the second connection's entry replaced the first's while the + inbound-connection counter counted both, the first connection's teardown then + removed the second's entry, and the second's own teardown found nothing to + remove, so the counter ended one above the connections it counts. That counter + gates the inbound connection limit, so a host repeating the collision could + hold it at the limit and lock out further inbound TCP connections until the + daemon restarted. Inbound entries now carry the accepted socket's local + address in their pool key as well as the remote one. - A peer that moves to a new address now loses the per-peer `connect(2)`-ed UDP socket pinned to the address it left. `set_current_addr` returns whether the address actually changed so the caller can drop the stale socket, and the diff --git a/src/transport/tcp/mod.rs b/src/transport/tcp/mod.rs index 25b29acc..e74f640a 100644 --- a/src/transport/tcp/mod.rs +++ b/src/transport/tcp/mod.rs @@ -12,9 +12,9 @@ //! ## Architecture //! //! Unlike UDP (one socket serves all peers), TCP requires one `TcpStream` -//! per peer. The transport maintains a connection pool mapping -//! `TransportAddr` to per-connection state, plus an optional `TcpListener` -//! for inbound connections. +//! per peer. The transport maintains a connection pool mapping each +//! connection's four-tuple to its per-connection state, plus an optional +//! `TcpListener` for inbound connections. //! //! ## Framing //! @@ -32,7 +32,10 @@ use super::{ }; use crate::config::TcpConfig; use crate::transport::framing::read_fmp_packet; -use pool::{ConnectingEntry, ConnectingPool, ConnectionPool, Direction, TcpConnection}; +use pool::{ + ConnectingEntry, ConnectingPool, ConnectionPool, Direction, PoolKey, TcpConnection, + key_for_remote, +}; use stats::TcpStats; use futures::FutureExt; @@ -57,7 +60,7 @@ use tracing::{debug, info, trace, warn}; /// /// Provides connection-oriented, reliable byte stream delivery over TCP/IP. /// Each peer has its own TCP connection; links are managed per-connection -/// with a connection pool keyed by `TransportAddr`. +/// with a connection pool keyed by `PoolKey`. pub struct TcpTransport { /// Unique transport identifier. transport_id: TransportId, @@ -269,7 +272,7 @@ impl TcpTransport { // aborting the task skips that path; decrement explicitly here // using the direction we stored on the connection record. let mut pool = self.pool.lock().await; - for (addr, conn) in pool.drain() { + for (key, conn) in pool.drain() { conn.recv_task.abort(); let _ = conn.recv_task.await; match conn.direction { @@ -278,7 +281,7 @@ impl TcpTransport { } debug!( transport_id = %self.transport_id, - remote_addr = %addr, + remote_addr = %key.remote, direction = ?conn.direction, "TCP connection closed (transport stopping)" ); @@ -326,7 +329,7 @@ impl TcpTransport { // Get or create connection let writer = { let pool = self.pool.lock().await; - pool.get(addr).map(|c| c.writer.clone()) + key_for_remote(&pool, addr).and_then(|key| pool.get(&key).map(|c| c.writer.clone())) }; let writer = match writer { @@ -355,7 +358,8 @@ impl TcpTransport { drop(w); // Remove failed connection from pool let mut pool = self.pool.lock().await; - if let Some(conn) = pool.remove(addr) { + let key = key_for_remote(&pool, addr); + if let Some(conn) = key.and_then(|key| pool.remove(&key)) { conn.recv_task.abort(); match conn.direction { Direction::Inbound => self.stats.record_pool_inbound_removed(), @@ -417,14 +421,15 @@ impl TcpTransport { let packet_tx = self.packet_tx.clone(); let pool = self.pool.clone(); let recv_stats = self.stats.clone(); - let remote_addr = addr.clone(); + let key = PoolKey::outbound(addr.clone()); + let recv_key = key.clone(); let mtu = mss_mtu; let recv_task = tokio::spawn(async move { tcp_receive_loop( read_half, transport_id, - remote_addr.clone(), + recv_key, packet_tx, pool, mtu, @@ -447,7 +452,7 @@ impl TcpTransport { }; let mut pool = self.pool.lock().await; - pool.insert(addr.clone(), conn); + pool.insert(key, conn); self.stats.record_connection_established(); self.stats.record_pool_outbound_added(); @@ -468,7 +473,8 @@ impl TcpTransport { /// and drops the write half (sends FIN to remote). pub async fn close_connection_async(&self, addr: &TransportAddr) { let mut pool = self.pool.lock().await; - if let Some(conn) = pool.remove(addr) { + let key = key_for_remote(&pool, addr); + if let Some(conn) = key.and_then(|key| pool.remove(&key)) { conn.recv_task.abort(); match conn.direction { Direction::Inbound => self.stats.record_pool_inbound_removed(), @@ -498,7 +504,7 @@ impl TcpTransport { // Already established? { let pool = self.pool.lock().await; - if pool.contains_key(addr) { + if key_for_remote(&pool, addr).is_some() { return Ok(()); } } @@ -610,7 +616,7 @@ impl TcpTransport { pub fn connection_state_sync(&self, addr: &TransportAddr) -> ConnectionState { // Check established pool first if let Ok(pool) = self.pool.try_lock() { - if pool.contains_key(addr) { + if key_for_remote(&pool, addr).is_some() { return ConnectionState::Connected; } } else { @@ -671,13 +677,14 @@ impl TcpTransport { let packet_tx = self.packet_tx.clone(); let pool = self.pool.clone(); let recv_stats = self.stats.clone(); - let remote_addr = addr.clone(); + let key = PoolKey::outbound(addr.clone()); + let recv_key = key.clone(); let recv_task = tokio::spawn(async move { tcp_receive_loop( read_half, transport_id, - remote_addr.clone(), + recv_key, packet_tx, pool, mss_mtu, @@ -702,7 +709,7 @@ impl TcpTransport { // Use try_lock since we're in a sync context and the pool // should be available (connection_state_sync already checked it) if let Ok(mut pool) = self.pool.try_lock() { - pool.insert(addr.clone(), conn); + pool.insert(key, conn); self.stats.record_connection_established(); self.stats.record_pool_outbound_added(); debug!( @@ -829,6 +836,24 @@ async fn accept_loop( loop { match listener.accept().await { Ok((stream, peer_addr)) => { + // The pool key is the four-tuple, so the local address is + // needed before anything else is done with the socket. A + // socket whose local address cannot be read is already + // broken; drop it rather than pool it under a key that + // could collide with another connection. + let local_addr = match stream.local_addr() { + Ok(a) => a, + Err(e) => { + warn!( + transport_id = %transport_id, + peer_addr = %peer_addr, + error = %e, + "Failed to read local address of accepted socket" + ); + continue; + } + }; + // Check inbound connection cap. Counts only inbound (accepted) // connections currently held in the pool; outbound (connect-on-send) // connections live in the same pool but are not subject to the @@ -889,6 +914,7 @@ async fn accept_loop( }; let remote_addr = TransportAddr::from_string(&peer_addr.to_string()); + let key = PoolKey::inbound(remote_addr.clone(), local_addr); // Split and spawn receive task let (read_half, write_half) = stream.into_split(); @@ -897,7 +923,7 @@ async fn accept_loop( let recv_pool = pool.clone(); let recv_packet_tx = packet_tx.clone(); let recv_stats = stats.clone(); - let recv_addr = remote_addr.clone(); + let recv_key = key.clone(); // Readiness barrier: the receive task must not reach its // cleanup path before the pool insert and counter bump below, @@ -909,7 +935,7 @@ async fn accept_loop( tcp_receive_loop( read_half, transport_id, - recv_addr, + recv_key, recv_packet_tx, recv_pool, conn_mtu, @@ -930,7 +956,7 @@ async fn accept_loop( }; let mut pool_guard = pool.lock().await; - pool_guard.insert(remote_addr.clone(), conn); + pool_guard.insert(key, conn); drop(pool_guard); stats.record_connection_accepted(); @@ -943,6 +969,7 @@ async fn accept_loop( debug!( transport_id = %transport_id, remote_addr = %remote_addr, + local_addr = %local_addr, mtu = conn_mtu, "Accepted inbound TCP connection" ); @@ -965,8 +992,11 @@ async fn accept_loop( /// Per-connection TCP receive loop. /// /// Reads complete FMP packets using the stream reader, delivers them to -/// the node via the packet channel. On error or EOF, removes the -/// connection from the pool and exits. `direction` is captured here so +/// the node via the packet channel. On error or EOF, removes its own +/// pool entry and exits. The entry is named by `key`, the connection's +/// four-tuple, so a loop tears down only the connection it was spawned +/// for even when another connection shares its peer address. +/// `direction` is captured here so /// the cleanup path can decrement the correct `pool_inbound` / /// `pool_outbound` counter regardless of whether the matching pool /// entry survived to be removed. @@ -980,7 +1010,7 @@ async fn accept_loop( async fn tcp_receive_loop( mut reader: tokio::net::tcp::OwnedReadHalf, transport_id: TransportId, - remote_addr: TransportAddr, + key: PoolKey, packet_tx: PacketTx, pool: ConnectionPool, mtu: u16, @@ -989,6 +1019,7 @@ async fn tcp_receive_loop( first_frame_timeout: Option, ready_rx: Option>, ) { + let remote_addr = &key.remote; debug!( transport_id = %transport_id, remote_addr = %remote_addr, @@ -1041,7 +1072,7 @@ async fn tcp_receive_loop( "TCP packet received" ); - let packet = ReceivedPacket::new(transport_id, remote_addr.clone(), data); + let packet = ReceivedPacket::new(transport_id, key.remote.clone(), data); if packet_tx.send(packet).await.is_err() { debug!( @@ -1071,7 +1102,7 @@ async fn tcp_receive_loop( // entry actually being removed so a double-cleanup never drives // the counter below zero. let mut pool_guard = pool.lock().await; - let removed = pool_guard.remove(&remote_addr).is_some(); + let removed = pool_guard.remove(&key).is_some(); drop(pool_guard); if removed { match direction { @@ -1233,6 +1264,16 @@ mod tests { } } + /// Listener on every local address, so that two connections can reach + /// it from one source port on two different local addresses. + fn wildcard_config() -> TcpConfig { + TcpConfig { + bind_addr: Some("0.0.0.0:0".to_string()), + mtu: Some(1400), + ..Default::default() + } + } + fn make_outbound_config() -> TcpConfig { TcpConfig { bind_addr: None, @@ -1473,7 +1514,7 @@ mod tests { // Connection should exist { let pool = t1.pool.lock().await; - assert!(pool.contains_key(&remote)); + assert!(pool.contains_key(&PoolKey::outbound(remote.clone()))); } // Close it @@ -1482,7 +1523,7 @@ mod tests { // Connection should be gone { let pool = t1.pool.lock().await; - assert!(!pool.contains_key(&remote)); + assert!(!pool.contains_key(&PoolKey::outbound(remote.clone()))); } t1.stop_async().await.unwrap(); @@ -2049,13 +2090,16 @@ mod tests { let listen = listener.local_addr().unwrap(); let client = TcpStream::connect(listen).await.unwrap(); let (server, peer_addr) = listener.accept().await.unwrap(); - let remote = TransportAddr::from_string(&peer_addr.to_string()); + let key = PoolKey::inbound( + TransportAddr::from_string(&peer_addr.to_string()), + server.local_addr().unwrap(), + ); let (read_half, write_half) = server.into_split(); let pool: ConnectionPool = Arc::new(Mutex::new(HashMap::new())); let stats = Arc::new(TcpStats::new()); pool.lock().await.insert( - remote.clone(), + key.clone(), TcpConnection { writer: Arc::new(Mutex::new(write_half)), recv_task: tokio::spawn(async {}), @@ -2073,7 +2117,7 @@ mod tests { tcp_receive_loop( read_half, TransportId::new(1), - remote.clone(), + key, tx, pool.clone(), 1400, @@ -2147,4 +2191,173 @@ mod tests { drop(client); transport.stop_async().await.unwrap(); } + + // ======================================================================== + // Inbound pool keying + // ======================================================================== + + /// Connect to `dst` from `src`, leaving the source port shareable. + /// + /// `SO_REUSEADDR` on both clients is what lets the second one bind the + /// source port the first is already using; the four-tuples still + /// differ, because the two connect to different local addresses. + #[cfg(target_os = "linux")] + fn client_from_port(src: SocketAddr, dst: SocketAddr) -> std::net::TcpStream { + let socket = socket2::Socket::new( + socket2::Domain::IPV4, + socket2::Type::STREAM, + Some(socket2::Protocol::TCP), + ) + .expect("create client socket"); + socket.set_reuse_address(true).expect("SO_REUSEADDR"); + socket.bind(&src.into()).expect("bind client socket"); + socket.connect(&dst.into()).expect("connect client socket"); + socket.into() + } + + /// A reply addressed to an inbound peer goes back over the connection + /// that peer opened, rather than dialing its ephemeral port. + /// + /// An inbound entry is keyed by the four-tuple, but a caller answering + /// a received packet knows only the remote address it came from. The + /// pool has to resolve that address to the entry; if it does not, the + /// send falls through to connect-on-send against the peer's ephemeral + /// port and fails. + #[tokio::test] + async fn a_reply_to_an_inbound_peer_uses_the_connection_it_arrived_on() { + let (tx, mut rx) = packet_channel(100); + let mut transport = TcpTransport::new(TransportId::new(1), None, make_config(), tx); + transport.start_async().await.unwrap(); + let listen = transport.local_addr().unwrap(); + + let mut peer = TcpStream::connect(listen).await.unwrap(); + peer.write_all(&build_msg1_frame()).await.unwrap(); + let packet = timeout(Duration::from_secs(2), rx.recv()) + .await + .expect("timeout waiting for the inbound frame") + .expect("packet channel closed"); + + let mut reply = vec![0xBB; 69]; + reply[0] = 0x02; + reply[1] = 0x00; + reply[2..4].copy_from_slice(&65u16.to_le_bytes()); + transport + .send_async(&packet.remote_addr, &reply) + .await + .expect("a reply to an inbound peer should use its connection"); + + let mut received = vec![0u8; reply.len()]; + timeout( + Duration::from_secs(2), + tokio::io::AsyncReadExt::read_exact(&mut peer, &mut received), + ) + .await + .expect("timeout waiting for the reply") + .expect("reply read failed"); + assert_eq!(received, reply); + + drop(peer); + transport.stop_async().await.unwrap(); + } + + /// Two inbound connections that share a peer address get two pool + /// entries, and each one's teardown releases only its own. + /// + /// The kernel names a connection by its four-tuple, so a wildcard + /// listener accepts two connections with the same peer address when + /// they arrive on different local addresses. Keyed by the peer address + /// alone, the second entry would replace the first while the inbound + /// counter counted both, the first connection's teardown would remove + /// the second's entry, and the second's teardown would find nothing to + /// remove and leave the counter one high for the life of the process. + /// The counter gates accepts, so repeating that locks the listener out. + /// + /// Break-check: with `PoolKey::inbound` ignoring its local address, + /// the pool holds one entry instead of two, no entry survives the + /// first client's close, and the inbound count never returns to zero. + /// + /// Linux only: the whole of 127.0.0.0/8 is local there without + /// configuration, which is what gives the listener two local addresses + /// to accept on. + #[cfg(target_os = "linux")] + #[tokio::test] + async fn two_inbound_connections_sharing_a_peer_address_get_separate_entries() { + let (tx, _rx) = packet_channel(100); + let mut transport = TcpTransport::new(TransportId::new(1), None, wildcard_config(), tx); + transport.start_async().await.unwrap(); + let port = transport.local_addr().unwrap().port(); + + // The first client picks the shared source port; the second binds + // the same one and reaches the listener on 127.0.0.2. + let first = client_from_port( + "127.0.0.1:0".parse().unwrap(), + SocketAddr::from(([127, 0, 0, 1], port)), + ); + let source = first.local_addr().unwrap(); + let second = client_from_port(source, SocketAddr::from(([127, 0, 0, 2], port))); + + assert!( + wait_until( + || transport.stats().pool_inbound_count() == 2, + Duration::from_secs(2) + ) + .await, + "both connections should have been accepted" + ); + { + let pool = transport.pool.lock().await; + assert_eq!( + pool.len(), + 2, + "two live connections must not share one pool entry" + ); + let mut locals: Vec<_> = pool + .keys() + .map(|key| key.local.expect("an inbound key carries a local address")) + .collect(); + locals.sort(); + assert_eq!( + locals, + vec![ + SocketAddr::from(([127, 0, 0, 1], port)), + SocketAddr::from(([127, 0, 0, 2], port)), + ], + "the two entries should be the two four-tuples" + ); + } + + // The first connection's teardown must leave the second alone. + drop(first); + assert!( + wait_until( + || transport.stats().pool_inbound_count() == 1, + Duration::from_secs(2) + ) + .await, + "closing one connection should release one inbound slot" + ); + { + let pool = transport.pool.lock().await; + let surviving: Vec<_> = pool.keys().map(|key| key.local.unwrap()).collect(); + assert_eq!( + surviving, + vec![SocketAddr::from(([127, 0, 0, 2], port))], + "the connection that was not closed should keep its entry" + ); + } + + // And the second connection's own teardown releases its slot. + drop(second); + assert!( + wait_until( + || transport.stats().pool_inbound_count() == 0, + Duration::from_secs(2) + ) + .await, + "the inbound count must return to zero, not strand a slot" + ); + assert!(transport.pool.lock().await.is_empty()); + + transport.stop_async().await.unwrap(); + } } diff --git a/src/transport/tcp/pool.rs b/src/transport/tcp/pool.rs index 8777c376..9c3124da 100644 --- a/src/transport/tcp/pool.rs +++ b/src/transport/tcp/pool.rs @@ -4,6 +4,7 @@ //! TCP transport. use std::collections::HashMap; +use std::net::SocketAddr; use std::sync::Arc; use tokio::net::TcpStream; use tokio::net::tcp::OwnedWriteHalf; @@ -34,14 +35,74 @@ pub(crate) struct TcpConnection { #[allow(dead_code)] pub(crate) mtu: u16, /// When the connection was established. - #[allow(dead_code)] pub(crate) established_at: Instant, /// Direction of the connection — drives pool-inbound/outbound accounting. pub(crate) direction: Direction, } +/// Key identifying one pooled connection. +/// +/// The kernel demultiplexes TCP by the connection four-tuple, so a remote +/// `ip:port` does not name a connection on its own: two accepted sockets +/// can carry the same peer address when they arrive on different local +/// addresses of a wildcard listener. Inbound entries therefore carry the +/// accepted socket's local address as well, and two such connections get +/// two entries instead of one overwriting the other. +/// +/// Outbound entries carry no local address. Nothing distinguishes two +/// outbound connections to one peer — the transport makes at most one — +/// and leaving the local address out keeps the connect-on-send lookup a +/// single hash probe. +#[derive(Clone, Debug, PartialEq, Eq, Hash)] +pub(crate) struct PoolKey { + /// Remote address, as the peer is named by callers and packets. + pub(crate) remote: TransportAddr, + /// Local address of an accepted socket; `None` for outbound. + pub(crate) local: Option, +} + +impl PoolKey { + /// Key for a connection this node initiated. + pub(crate) fn outbound(remote: TransportAddr) -> Self { + Self { + remote, + local: None, + } + } + + /// Key for a connection the listener accepted on `local`. + pub(crate) fn inbound(remote: TransportAddr, local: SocketAddr) -> Self { + Self { + remote, + local: Some(local), + } + } +} + +/// The pooled connections, keyed by [`PoolKey`]. +pub(crate) type PoolMap = HashMap; + /// Shared connection pool. -pub(crate) type ConnectionPool = Arc>>; +pub(crate) type ConnectionPool = Arc>; + +/// Resolve a remote address to the key of the connection that a caller +/// naming only that address should use. +/// +/// An outbound connection is keyed by the remote address alone, so the +/// common case is one hash probe. Inbound entries also carry a local +/// address, and are searched; where several share a remote address the +/// most recently established one wins, which is the entry an +/// address-keyed pool held before inbound keys became four-tuples. +pub(crate) fn key_for_remote(pool: &PoolMap, remote: &TransportAddr) -> Option { + let outbound = PoolKey::outbound(remote.clone()); + if pool.contains_key(&outbound) { + return Some(outbound); + } + pool.iter() + .filter(|(key, _)| &key.remote == remote) + .max_by_key(|(_, conn)| conn.established_at) + .map(|(key, _)| key.clone()) +} /// A pending background connection attempt. /// From 43b65125033d356d245ac8ee095362a6fb315167 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 17 Sep 2026 20:45:48 +0000 Subject: [PATCH 5/6] fix(gateway): harden the virtual-IP pool and ship it disabled on OpenWrt Any host able to query the LAN resolver could drain the gateway's 65,535-address virtual-IP pool one `.fips` name at a time, and each allocation rebuilt the whole nftables table in a way that could leave the host with no NAT at all. Four changes, each independently useful, close that off. Do not allocate for query types the gateway never answers with an address. handle_query minted a virtual IP for every query type and only then looked at what the client asked, answering an A or HTTPS query with NODATA after creating a mapping for it. The query type is now decided before the pool is touched, and only AAAA and ANY allocate. The refresh an existing mapping used to get from any query type is kept: it came from the reuse path in allocate, so a new pool method does that refresh alone and never creates anything, and the reuse path calls it. Rebuild the NAT table in one netlink transaction. rebuild() deleted the fips_gateway table in a batch of its own and discarded the result, then sent a second batch recreating the table, the chains, the fips0 masquerade and two rules per mapping. Between those sends the host had no NAT table, and a recreate the kernel refused left the table deleted, turning one failed mapping change into a total loss of forwarding until some later rebuild happened to succeed. The delete and the recreate now share one batch. A leading table add makes the delete legal on the first run, since rustables sends it with NLM_F_CREATE and no NLM_F_EXCL and the crate offers no flush. Deciding what to send is now separate from sending it, which is the seam the new unit tests use: they assert one batch, the add-delete-add prefix, and that every chain and rule follows the recreate, without a netlink socket or privileges. Read conntrack once per tick, off the runtime thread, and match by address. The session count searched each /proc/net/nf_conntrack line for `dst=` followed by the virtual IP's compressed Display form, while the kernel prints every tuple with `%pI6`, the full uncompressed form. That string cannot occur in that field, so the count was zero for every mapping on every kernel that has the file: nothing pinned an in-use mapping and one whose client did not re-query DNS was reclaimed about two minutes after its last DNS reference with traffic still flowing. Each `dst=` is now parsed and compared as an address. The read was also per mapping, under the pool lock, on the runtime thread that serves DNS; the tick now takes one snapshot in a blocking task before taking the lock. An unreadable source was silent, because read_to_string's error became zero through unwrap_or(0). Zero stays, since treating it as in-use would pin every mapping forever on a kernel without CONFIG_NF_CONNTRACK_PROCFS, but it is now reported at warn on the first failure and on each change of outcome, and at debug on a repeat. Ship the OpenWrt gateway disabled, and keep its state across upgrades. The generated postinst enabled and started fips-gateway on every install, against the init script's own header, the package README and the deployment tutorial, which all say the service ships disabled. A fresh install now leaves it alone. Upgrades are the awkward case: opkg runs the outgoing package's prerm first, and every released prerm disabled the gateway on its way out without recording whether it had been enabled. The new prerm stops the services on an upgrade but no longer disables them, and leaves a marker the incoming postinst reads. With the marker, enablement survived and the gateway starts only if it was enabled; without it, the outgoing package was a released one whose prerm destroyed that state, so the gateway is re-enabled rather than letting an upgrade turn off a working deployment. That re-enables a hand-disabled gateway once, which the CHANGELOG says. start_service now reads gateway.enabled from fips.yaml before touching anything, since starting a gateway the config disables used to take dnsmasq's `.fips` forwarding away from the daemon and hand it to a port whose daemon exits immediately. The four maintainer-script bodies move out of heredocs in the two build scripts into packaging/openwrt-ipk/scripts/, so the .ipk and the .apk install the same bodies and a test can run what ships. Coverage recorded rather than closed. The conntrack parser's first test builds its line from the kernel's own format string rather than a capture, because this host is built without CONFIG_NF_CONNTRACK_PROCFS and has no /proc/net/nf_conntrack, so the lab exercises only the unreadable path. Kernel acceptance of delete-then-recreate inside one transaction is not asserted by a unit test; the gateway suite is what proves it, since the manager rebuilds at startup and the daemon exits if that fails. The OpenWrt scenarios run the shipped script bodies under ash in a busybox container against stubbed init scripts, and assert their behaviour given opkg's call order, arguments and PKG_UPGRADE as read from opkg-lede's sources, not under a real opkg upgrade on a router image. Admission limits on the pool are deliberately not included here: they need a measurement run before their constants can be chosen. --- .github/workflows/ci.yml | 34 ++ CHANGELOG.md | 69 +++ packaging/openwrt-apk/build-apk.sh | 33 +- packaging/openwrt-ipk/build-ipk.sh | 29 +- .../openwrt-ipk/files/etc/init.d/fips-gateway | 15 + packaging/openwrt-ipk/scripts/postinst | 47 ++ packaging/openwrt-ipk/scripts/prerm | 26 + src/bin/fips-gateway.rs | 56 ++- src/gateway/dns.rs | 175 ++++++- src/gateway/nat.rs | 459 +++++++++++++----- src/gateway/pool.rs | 316 ++++++++++-- testing/ci-local.sh | 29 ++ testing/openwrt/fixtures/released-prerm | 6 + testing/openwrt/maintainer-scripts-test.sh | 46 ++ testing/openwrt/scenarios.sh | 342 +++++++++++++ 15 files changed, 1428 insertions(+), 254 deletions(-) create mode 100755 packaging/openwrt-ipk/scripts/postinst create mode 100755 packaging/openwrt-ipk/scripts/prerm create mode 100644 testing/openwrt/fixtures/released-prerm create mode 100755 testing/openwrt/maintainer-scripts-test.sh create mode 100755 testing/openwrt/scenarios.sh diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index f9bb2c6b..0a35dbe8 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -435,6 +435,40 @@ jobs: Write-Host "PSScriptAnalyzer: no issues" } +# ───────────────────────────────────────────────────────────────────────────── +# Job 2e – OpenWrt maintainer-script scenarios +# +# Runs the package's postinst/prerm and the fips-gateway init script under ash +# in a busybox container, against stubbed init scripts: a fresh install, an +# upgrade from a released package, an upgrade from a package carrying these +# scripts with the gateway enabled and with it disabled, a removal, and the +# init script's gateway.enabled guard. +# +# A job of its own rather than a leg of the integration matrix: it needs no +# FIPS binary and no shared test image, so as an integration leg it would wait +# on the build and then download and build both for nothing. +# +# The leg keeps `suite:` because testing/check-ci-parity.sh matches it against +# OPENWRT_SUITES in ci-local.sh; the step below does not read it. +# ───────────────────────────────────────────────────────────────────────────── + openwrt-scripts: + name: OpenWrt scripts (${{ matrix.suite }}) + runs-on: ubuntu-latest + if: ${{ !inputs.skip_integration }} + + strategy: + fail-fast: false + matrix: + include: + - suite: openwrt-scripts + + steps: + - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 # v6 + + - name: Run the OpenWrt maintainer-script scenarios + timeout-minutes: 5 + run: bash testing/openwrt/maintainer-scripts-test.sh + # ───────────────────────────────────────────────────────────────────────────── # Job 3 – Integration tests (static mesh + chaos simulation) # diff --git a/CHANGELOG.md b/CHANGELOG.md index 5ebedb74..4b09a88c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -105,6 +105,75 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 names the path. A dangling symlink likewise aborts rather than being replaced. The legacy `/etc/fips/fips.key` lookup follows the same rule. +#### Gateway + +- A `.fips` query the gateway answers without an address no longer takes an + address from the pool. Every query type was allocated a mapping before the + code looked at what the client had asked for, and an A or HTTPS query was + then answered with NODATA, so any host that can reach the LAN resolver could + consume the pool one name at a time with a query type it is never given an + address for. Only AAAA and ANY allocate now. A non-AAAA query for a name that + already has a mapping still refreshes that mapping's TTL clock, so a client + querying both types does not lose half of its refresh. +- Conntrack sessions are matched by address rather than by text, so live + traffic pins a gateway mapping again. The session count searched each + `/proc/net/nf_conntrack` line for `dst=` followed by the virtual IP in its + compressed form (`fd01::1`), while the kernel prints tuples in the full + uncompressed form (`dst=fd01:0000:0000:0000:0000:0000:0000:0001`), so the + count was zero for every mapping on every kernel. Nothing pinned an in-use + mapping, and one whose client did not re-query DNS was reclaimed about two + minutes after its last DNS reference while its traffic was still flowing. + Each `dst=` value is now parsed as an address and compared as one. +- The conntrack table is read once per tick instead of once per mapping, and + the read happens off the runtime thread. The whole file was read and scanned + for each mapping in turn, while the pool lock was held, on the same + single-threaded runtime that serves DNS. The tick now takes one snapshot with + a blocking task before it takes the lock, and the pool does a map lookup per + mapping. +- A conntrack source that cannot be read is reported. It still counts as zero + sessions for every mapping, as it always has, so reclamation keeps working + rather than pinning the whole pool; but the first failure and each change of + outcome after it are now logged, so an unreadable source is no longer + indistinguishable from an idle one. A kernel built without + `CONFIG_NF_CONNTRACK_PROCFS` has no `/proc/net/nf_conntrack` at all and fails + identically every tick, so a repeat is logged at debug rather than warn. +- The NAT table is rebuilt in one netlink transaction. A rebuild deleted the + `fips_gateway` table in a batch of its own, discarded that batch's result, + and only then sent the batch that recreated the table, the chains, the + `fips0` masquerade and every per-mapping rule. Between the two sends the + gateway had no NAT at all, and a recreate the kernel refused left the table + absent for good, taking down forwarding for every existing mapping rather + than failing the one change that was being made. The delete and the recreate + now share a single batch, which the kernel applies as one transaction, so a + refused rebuild leaves the previous table in the packet path. The rules sent + are unchanged. +- A new OpenWrt install no longer enables and starts `fips-gateway`. The + generated postinst turned it on unconditionally, contradicting the init + script's own header, the package README and the deployment tutorial, all of + which say the service ships disabled and is enabled deliberately. The + documented `service fips-gateway enable` / `service fips-gateway start` steps + are unchanged, and the shipped `fips.yaml` still carries `gateway.enabled: + true`, so enabling the service is all that is needed. +- **The first upgrade to this release re-enables and starts `fips-gateway` on + any router that has the package installed, including one where the gateway + was disabled by hand.** Every released package's prerm disabled the service + on its way out, leaving nothing behind that says whether the operator wanted + it on, so an upgrade cannot tell the two apart and keeps the gateway running + rather than silently turning off a working one. If you had disabled it, run + `service fips-gateway disable` once after upgrading. Later upgrades preserve + whatever state the service is in: the new prerm stops the services on an + upgrade but no longer disables them. +- `start_service` in the `fips-gateway` init script now reads `gateway.enabled` + from `/etc/fips/fips.yaml` before doing anything. Starting a gateway that the + config disables used to hand dnsmasq's `.fips` forwarding to the gateway's + port, add the LAN prefix and advertise the pool route, and only then start a + daemon that exits immediately because the gateway is disabled, leaving `.fips` + resolution pointed at a port nothing listens on. +- The four OpenWrt maintainer-script bodies now live in + `packaging/openwrt-ipk/scripts/` instead of inside heredocs in the two build + scripts, so the `.ipk` and `.apk` packages install the same bodies and the + scenarios in `testing/openwrt/` run what ships. + ### Changed - The lockfile moves `chacha20` from 0.10.1 to 0.10.2, because 0.10.1 is yanked. diff --git a/packaging/openwrt-apk/build-apk.sh b/packaging/openwrt-apk/build-apk.sh index 0788f7d0..cfe18a60 100755 --- a/packaging/openwrt-apk/build-apk.sh +++ b/packaging/openwrt-apk/build-apk.sh @@ -94,6 +94,7 @@ PROJECT_ROOT="$(cd "$SCRIPT_DIR/../.." && pwd)" # The installed-filesystem payload (init scripts, config, sysctl, etc.) is # shared with the .ipk package; there is one canonical copy in openwrt-ipk/. FILES_DIR="$PROJECT_ROOT/packaging/openwrt-ipk/files" +SCRIPTS_SRC="$PROJECT_ROOT/packaging/openwrt-ipk/scripts" DIST_DIR="$PROJECT_ROOT/dist" PKG_NAME="fips" @@ -228,33 +229,15 @@ EOF # ---- maintainer scripts ---- # Map our opkg maintainer scripts onto apk's lifecycle phases: -# opkg postinst -> apk post-install (enable + start services) +# opkg postinst -> apk post-install (enable + start the daemon) # opkg prerm -> apk pre-deinstall (stop + disable services) -cat > "$SCRIPTS_DIR/post-install" <<'EOF' -#!/bin/sh -# Run first-boot UCI setup (the script deletes itself when done). -if [ -x /etc/uci-defaults/90-fips-setup ]; then - /etc/uci-defaults/90-fips-setup && rm -f /etc/uci-defaults/90-fips-setup -fi - -/etc/init.d/fips enable -/etc/init.d/fips start -/etc/init.d/fips-gateway enable -/etc/init.d/fips-gateway start -exit 0 -EOF - -cat > "$SCRIPTS_DIR/pre-deinstall" <<'EOF' -#!/bin/sh -/etc/init.d/fips-gateway stop 2>/dev/null || true -/etc/init.d/fips-gateway disable 2>/dev/null || true -/etc/init.d/fips stop 2>/dev/null || true -/etc/init.d/fips disable 2>/dev/null || true -exit 0 -EOF - -chmod 0755 "$SCRIPTS_DIR/post-install" "$SCRIPTS_DIR/pre-deinstall" +# Both bodies come from packaging/openwrt-ipk/scripts/, the same files the +# .ipk ships, so the two packagers cannot drift apart and testing/openwrt/ +# exercises what both install. apk runs post-install only on a fresh install, +# so the postinst's upgrade branch is unreachable here. +install -m 0755 "$SCRIPTS_SRC/postinst" "$SCRIPTS_DIR/post-install" +install -m 0755 "$SCRIPTS_SRC/prerm" "$SCRIPTS_DIR/pre-deinstall" # --------------------------------------------------------------------------- # 3. Assemble the .apk via apk mkpkg diff --git a/packaging/openwrt-ipk/build-ipk.sh b/packaging/openwrt-ipk/build-ipk.sh index 14b9f7fd..1b2c200e 100755 --- a/packaging/openwrt-ipk/build-ipk.sh +++ b/packaging/openwrt-ipk/build-ipk.sh @@ -81,6 +81,7 @@ esac SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd)" PROJECT_ROOT="$(cd "$SCRIPT_DIR/../.." && pwd)" FILES_DIR="$SCRIPT_DIR/files" +SCRIPTS_SRC="$SCRIPT_DIR/scripts" # maintainer scripts (metadata, not payload) DIST_DIR="$PROJECT_ROOT/dist" PKG_NAME="fips" @@ -212,30 +213,10 @@ cat > "$CONTROL_DIR/conffiles" < "$CONTROL_DIR/postinst" <<'EOF' -#!/bin/sh -# Run first-boot UCI setup (the script deletes itself when done). -if [ -x /etc/uci-defaults/90-fips-setup ]; then - /etc/uci-defaults/90-fips-setup && rm -f /etc/uci-defaults/90-fips-setup -fi - -/etc/init.d/fips enable -/etc/init.d/fips start -/etc/init.d/fips-gateway enable -/etc/init.d/fips-gateway start -exit 0 -EOF -chmod 0755 "$CONTROL_DIR/postinst" - -cat > "$CONTROL_DIR/prerm" <<'EOF' -#!/bin/sh -/etc/init.d/fips-gateway stop 2>/dev/null || true -/etc/init.d/fips-gateway disable 2>/dev/null || true -/etc/init.d/fips stop 2>/dev/null || true -/etc/init.d/fips disable 2>/dev/null || true -exit 0 -EOF -chmod 0755 "$CONTROL_DIR/prerm" +# The maintainer scripts live in files of their own rather than in heredocs +# here, so testing/openwrt/ can run the same bodies the package ships. +install -m 0755 "$SCRIPTS_SRC/postinst" "$CONTROL_DIR/postinst" +install -m 0755 "$SCRIPTS_SRC/prerm" "$CONTROL_DIR/prerm" # ---- pack ---- diff --git a/packaging/openwrt-ipk/files/etc/init.d/fips-gateway b/packaging/openwrt-ipk/files/etc/init.d/fips-gateway index b9d2a9ec..4b9c7041 100755 --- a/packaging/openwrt-ipk/files/etc/init.d/fips-gateway +++ b/packaging/openwrt-ipk/files/etc/init.d/fips-gateway @@ -26,6 +26,14 @@ DAEMON_DNS_PORT=5354 GLOBAL_PREFIX="2001:2:f1b5::1/64" start_service() { + # The gateway daemon exits when gateway.enabled is not true, so without + # this check starting a disabled gateway would still take dnsmasq's .fips + # upstream away from the daemon and point it at a port nothing listens on. + if [ "$(gateway_config_enabled)" != "true" ]; then + logger -t fips-gateway "gateway.enabled is not true in $CONFIG; not starting" + return 1 + fi + # Apply gateway sysctls (proxy_ndp, IPv6 forwarding). sysctl -p /etc/sysctl.d/fips-gateway.conf 2>/dev/null || true @@ -72,6 +80,13 @@ reload_service() { restart } +# Extract the gateway "enabled" flag from fips.yaml. +# Prints the value indented under the top-level "gateway:" block, or nothing +# when there is no such block. +gateway_config_enabled() { + awk '/^gateway:/{found=1; next} found && /^[^ ]/{found=0} found && /enabled:/{gsub(/.*enabled:[[:space:]]*/, ""); gsub(/["'"'"']/, ""); print; exit}' "$CONFIG" +} + # Extract the gateway pool CIDR from fips.yaml. # Looks for "pool:" indented under the top-level "gateway:" block. gateway_pool_cidr() { diff --git a/packaging/openwrt-ipk/scripts/postinst b/packaging/openwrt-ipk/scripts/postinst new file mode 100755 index 00000000..3e505e6d --- /dev/null +++ b/packaging/openwrt-ipk/scripts/postinst @@ -0,0 +1,47 @@ +#!/bin/sh +# Maintainer script run after the FIPS package is unpacked. +# +# Installed as the .ipk CONTROL/postinst and registered as the .apk +# post-install script, so one body serves both packagers. +# +# The fips daemon is enabled and started on every install. The gateway is not: +# the package ships that service disabled, and the README and the deployment +# tutorial tell the operator to enable it deliberately. +# +# Upgrades are the awkward case, because opkg runs the OLD package's prerm +# before any script from the new one: +# +# marker present the old package was one of these, its prerm left +# enablement alone, and the gateway only needs starting +# again if it was enabled; +# no marker the old package's prerm disabled the gateway on its way +# out, so its former state is unrecoverable; the gateway is +# re-enabled, which also re-enables one an operator had +# disabled by hand. +# +# Under apk this script runs only on a fresh install, so it takes the first +# branch and the gateway stays off. + +UPGRADE_MARKER=/tmp/fips-prerm-upgrade + +# Run first-boot UCI setup (the script deletes itself when done). +if [ -x /etc/uci-defaults/90-fips-setup ]; then + /etc/uci-defaults/90-fips-setup && rm -f /etc/uci-defaults/90-fips-setup +fi + +/etc/init.d/fips enable +/etc/init.d/fips start + +if [ "${PKG_UPGRADE:-0}" = "1" ]; then + if [ -e "$UPGRADE_MARKER" ]; then + rm -f "$UPGRADE_MARKER" + else + /etc/init.d/fips-gateway enable + fi + + if /etc/init.d/fips-gateway enabled 2>/dev/null; then + /etc/init.d/fips-gateway start + fi +fi + +exit 0 diff --git a/packaging/openwrt-ipk/scripts/prerm b/packaging/openwrt-ipk/scripts/prerm new file mode 100755 index 00000000..f1db5153 --- /dev/null +++ b/packaging/openwrt-ipk/scripts/prerm @@ -0,0 +1,26 @@ +#!/bin/sh +# Maintainer script run before the FIPS package is removed or replaced. +# +# Installed as the .ipk CONTROL/prerm and registered as the .apk pre-deinstall +# script, so one body serves both packagers. +# +# opkg calls this with "upgrade " when the package is being +# replaced. Disabling the services there would erase the operator's choice, +# because nothing records it anywhere else, so an upgrade only stops them and +# leaves a marker telling the incoming postinst that enablement survived. +# A real removal stops and disables both, as before. + +UPGRADE_MARKER=/tmp/fips-prerm-upgrade + +if [ "$1" = "upgrade" ]; then + : > "$UPGRADE_MARKER" 2>/dev/null || true + /etc/init.d/fips-gateway stop 2>/dev/null || true + /etc/init.d/fips stop 2>/dev/null || true + exit 0 +fi + +/etc/init.d/fips-gateway stop 2>/dev/null || true +/etc/init.d/fips-gateway disable 2>/dev/null || true +/etc/init.d/fips stop 2>/dev/null || true +/etc/init.d/fips disable 2>/dev/null || true +exit 0 diff --git a/src/bin/fips-gateway.rs b/src/bin/fips-gateway.rs index ad4ea187..73b8743a 100644 --- a/src/bin/fips-gateway.rs +++ b/src/bin/fips-gateway.rs @@ -24,7 +24,7 @@ use tokio::signal::unix::{SignalKind, signal}; #[cfg(target_os = "linux")] use tokio::sync::{Mutex, mpsc, watch}; #[cfg(target_os = "linux")] -use tracing::{error, info, warn}; +use tracing::{debug, error, info, warn}; #[cfg(target_os = "linux")] use tracing_subscriber::{EnvFilter, fmt}; @@ -53,6 +53,53 @@ fn main() { std::process::exit(1); } +/// Take a conntrack snapshot off the runtime thread. +/// +/// A failed read yields an empty snapshot, so every mapping reads zero +/// sessions, which is what the pool did with an unreadable source before. The +/// alternative, treating "unknown" as "in use", would pin every mapping forever +/// on a kernel with no conntrack proc file and turn a read error into a pool +/// that never reclaims. The cost is the opposite error: a mapping carrying live +/// traffic can be reclaimed early while the source is unreadable. +#[cfg(target_os = "linux")] +async fn read_conntrack(log: &mut pool::ConntrackReadLog) -> pool::ConntrackSnapshot { + use fips::gateway::pool::ConntrackQuerier; + + match tokio::task::spawn_blocking(|| pool::ProcConntrack.snapshot()).await { + Ok(Ok(snapshot)) => { + log.observe(None); + snapshot + } + Ok(Err(e)) => { + report_unreadable_conntrack(log, e.kind(), &e.to_string()); + pool::ConntrackSnapshot::default() + } + Err(e) => { + report_unreadable_conntrack(log, std::io::ErrorKind::Other, &e.to_string()); + pool::ConntrackSnapshot::default() + } + } +} + +/// Log an unreadable conntrack source once per change of outcome. +#[cfg(target_os = "linux")] +fn report_unreadable_conntrack( + log: &mut pool::ConntrackReadLog, + kind: std::io::ErrorKind, + error: &str, +) { + match log.observe(Some(kind)) { + pool::ReadReport::Changed => warn!( + error, + "Conntrack unreadable; every mapping reads zero sessions" + ), + pool::ReadReport::Repeated => debug!( + error, + "Conntrack still unreadable; every mapping reads zero sessions" + ), + } +} + #[cfg(target_os = "linux")] #[tokio::main(flavor = "current_thread")] async fn main() { @@ -379,7 +426,7 @@ async fn main() { let tick_event_tx = event_tx; let tick_nat_count = Arc::clone(&nat_count); let mut tick_shutdown = shutdown_rx.clone(); - let conntrack = pool::ProcConntrack; + let mut conntrack_log = pool::ConntrackReadLog::default(); let snap_config = control::SnapshotConfig { pool_cidr: gw_config.pool.clone(), lan_interface: gw_config.lan_interface.clone(), @@ -395,6 +442,11 @@ async fn main() { tokio::select! { _ = interval.tick() => { let now = Instant::now(); + // Read conntrack once, off the runtime thread and before + // the pool lock: the runtime is current-thread, so a + // blocking read here would stall the DNS resolver, and the + // read must not happen under the lock the resolver needs. + let conntrack = read_conntrack(&mut conntrack_log).await; let mut pool_guard = tick_pool.lock().await; let events = pool_guard.tick(now, &conntrack); diff --git a/src/gateway/dns.rs b/src/gateway/dns.rs index cb71e367..4d9a069e 100644 --- a/src/gateway/dns.rs +++ b/src/gateway/dns.rs @@ -357,6 +357,32 @@ async fn handle_query( } }; + // What the client actually asked for. Only AAAA and ANY are answered with + // an address, and only those may mint a mapping: allocating for a query + // type the gateway answers with NODATA let any LAN host take a pool + // address per name without ever being given one. + let client_qtype = query + .questions + .first() + .map(|q| q.qtype) + .unwrap_or(QTYPE::TYPE(TYPE::AAAA)); + + if !matches!(client_qtype, QTYPE::TYPE(TYPE::AAAA) | QTYPE::ANY) { + // The client is still using the name, so an existing mapping's TTL + // clock is refreshed. A client that re-queries a mapped name with both + // A and AAAA should not lose half of its refresh, and with no + // conntrack sessions a DNS reference is all that keeps a mapping + // alive. Nothing is created. + let refreshed = pool.lock().await.refresh_if_present(node_addr); + debug!( + name = %fips_name, + mesh_addr = %mesh_addr, + refreshed, + "Non-AAAA .fips query, returning NODATA" + ); + return build_nodata(&query, ttl); + } + // Allocate virtual IP from pool let mut pool_guard = pool.lock().await; let (virtual_ip, is_new) = match pool_guard.allocate(node_addr, mesh_addr, &fips_name) { @@ -387,22 +413,7 @@ async fn handle_query( "Resolved .fips query" ); - // Check what the client originally asked for. - // Only return an AAAA record if the client asked for AAAA (or ANY). - // For A queries, return an empty NOERROR — the client's resolver will - // use the AAAA answer from its parallel AAAA query instead. - let client_qtype = query - .questions - .first() - .map(|q| q.qtype) - .unwrap_or(QTYPE::TYPE(TYPE::AAAA)); - - match client_qtype { - QTYPE::TYPE(TYPE::AAAA) | QTYPE::ANY => build_aaaa_response(&query, virtual_ip, ttl), - // All other types (A, HTTPS, etc.): return NODATA — the name exists - // but has no records of the requested type. - _ => build_nodata(&query, ttl), - } + build_aaaa_response(&query, virtual_ip, ttl) } #[cfg(test)] @@ -416,17 +427,28 @@ mod tests { /// Build a client-facing AAAA query. fn build_query(id: u16, qname: &str) -> Vec { + build_query_of_type(id, qname, QTYPE::TYPE(TYPE::AAAA)) + } + + /// Build a client-facing query of any type. + fn build_query_of_type(id: u16, qname: &str, qtype: QTYPE) -> Vec { let mut packet = Packet::new_query(id); - let question = Question::new( - Name::new_unchecked(qname), - QTYPE::TYPE(TYPE::AAAA), - CLASS::IN.into(), - false, - ); + let question = Question::new(Name::new_unchecked(qname), qtype, CLASS::IN.into(), false); packet.questions.push(question); packet.build_bytes_vec_compressed().unwrap() } + /// Assert the response is NODATA: NOERROR with no answer records. + fn assert_nodata(response: &[u8]) { + let packet = Packet::parse(response).unwrap(); + assert_eq!(packet.rcode(), RCODE::NoError); + assert!( + packet.answers.is_empty(), + "expected NODATA, got {} answer(s)", + packet.answers.len() + ); + } + /// Build an upstream NOERROR AAAA answer. fn build_answer(id: u16, qname: &str, addr: &str) -> Vec { let mut packet = Packet::new_reply(id); @@ -642,6 +664,115 @@ mod tests { )); } + #[tokio::test] + async fn an_a_query_returns_nodata_and_mints_no_mapping() { + let upstream_socket = UdpSocket::bind("[::1]:0").await.unwrap(); + let upstream = upstream_socket.local_addr().unwrap(); + let handle = spawn_upstream(upstream_socket, |id| { + vec![build_answer(id, "test.fips", "fd00::1")] + }); + + let pool = test_pool(); + let (event_tx, mut event_rx) = mpsc::channel(16); + let response = handle_query( + &build_query_of_type(0x1234, "test.fips", QTYPE::TYPE(TYPE::A)), + upstream, + TEST_TTL, + &pool, + &event_tx, + ) + .await + .unwrap(); + handle.await.unwrap(); + + assert_nodata(&response); + assert!( + matches!(event_rx.try_recv(), Err(mpsc::error::TryRecvError::Empty)), + "an A query minted a mapping, so any LAN host can take a pool \ + address per name with a query type it is never given one for" + ); + assert!( + pool.lock() + .await + .mapping_info(std::time::Instant::now()) + .is_empty(), + "an A query left a mapping in the pool" + ); + } + + #[tokio::test] + async fn an_a_query_refreshes_an_existing_mapping_without_creating_one() { + let pool = test_pool(); + let (event_tx, mut event_rx) = mpsc::channel(16); + + // An AAAA query mints the mapping. + let upstream_socket = UdpSocket::bind("[::1]:0").await.unwrap(); + let upstream = upstream_socket.local_addr().unwrap(); + let handle = spawn_upstream(upstream_socket, |id| { + vec![build_answer(id, "test.fips", "fd00::1")] + }); + let response = handle_query( + &build_query(0x1234, "test.fips"), + upstream, + TEST_TTL, + &pool, + &event_tx, + ) + .await + .unwrap(); + handle.await.unwrap(); + let virtual_ip = assert_pool_answer(&response); + assert!(matches!( + event_rx.try_recv().unwrap(), + PoolEvent::MappingCreated { .. } + )); + + let before = { + let guard = pool.lock().await; + guard + .lookup_virtual_ip(&virtual_ip) + .unwrap() + .last_referenced + }; + + // An A query for the same name refreshes it and creates nothing. A + // client that re-queries a mapped name with both types must not lose + // half of its refresh: with no conntrack sessions, the DNS reference + // is the only thing keeping the mapping alive. + let upstream_socket = UdpSocket::bind("[::1]:0").await.unwrap(); + let upstream = upstream_socket.local_addr().unwrap(); + let handle = spawn_upstream(upstream_socket, |id| { + vec![build_answer(id, "test.fips", "fd00::1")] + }); + let response = handle_query( + &build_query_of_type(0x1235, "test.fips", QTYPE::TYPE(TYPE::A)), + upstream, + TEST_TTL, + &pool, + &event_tx, + ) + .await + .unwrap(); + handle.await.unwrap(); + + assert_nodata(&response); + + let guard = pool.lock().await; + let mapping = guard + .lookup_virtual_ip(&virtual_ip) + .expect("the A query removed or replaced the mapping"); + assert!( + mapping.last_referenced > before, + "the A query did not refresh the mapping's TTL clock" + ); + drop(guard); + + assert!( + matches!(event_rx.try_recv(), Err(mpsc::error::TryRecvError::Empty)), + "the A query sent a second MappingCreated" + ); + } + #[tokio::test] async fn test_healthy_path_resolves() { let upstream_socket = UdpSocket::bind("[::1]:0").await.unwrap(); diff --git a/src/gateway/nat.rs b/src/gateway/nat.rs index b1859f61..f0c11d60 100644 --- a/src/gateway/nat.rs +++ b/src/gateway/nat.rs @@ -51,6 +51,30 @@ struct NatMapping { mesh_addr: Ipv6Addr, } +/// One object a NAT rebuild sends, named rather than built. +/// +/// `rebuild_batches` decides what a rebuild sends and in what order; +/// `send_batches` turns that decision into rustables objects and hands each +/// batch to the kernel. The split is what lets a test see the delete and the +/// recreate share one transaction without a netlink socket, which is the +/// property that keeps the table in the packet path. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum NatOp { + Table(MsgType), + PreChain, + PostChain, + /// Masquerade for traffic leaving through `fips0`. + FipsMasquerade, + /// DNAT for the mapping with this virtual IP. + Dnat(Ipv6Addr), + /// SNAT for the mapping with this virtual IP. + Snat(Ipv6Addr), + /// DNAT for the port forward at this index in `port_forwards`. + PortForward(usize), + /// LAN-side masquerade, emitted once when any port forward exists. + LanMasquerade, +} + /// NAT rule manager using nftables via rustables netlink API. /// /// Rebuilds the entire nftables table atomically on every change to @@ -71,15 +95,11 @@ pub struct NatManager { } impl NatManager { - /// Create the nftables table and NAT chains. + /// Build the manager's state without touching netlink. /// - /// Installs a masquerade rule for traffic exiting via `fips0` so that - /// LAN client source addresses are rewritten to the gateway's mesh - /// address, allowing return traffic to route back through the mesh. - /// - /// `lan_interface` is the gateway's LAN-facing interface name, - /// needed by the port-forward LAN-side masquerade rule. - pub fn new(lan_interface: String) -> Result { + /// Everything `new` does except sending the first rebuild, so a test can + /// exercise the batch builder with no socket and no privileges. + fn with_state(lan_interface: String) -> Self { let table = Table::new(ProtocolFamily::Inet).with_name(TABLE_NAME); let pre_chain = Chain::new(&table) .with_name(PREROUTING_CHAIN) @@ -90,14 +110,26 @@ impl NatManager { .with_type(ChainType::Nat) .with_hook(Hook::new(HookClass::PostRouting, SRCNAT_PRIORITY)); - let mgr = Self { + Self { table, pre_chain, post_chain, lan_interface, mappings: HashMap::new(), port_forwards: Vec::new(), - }; + } + } + + /// Create the nftables table and NAT chains. + /// + /// Installs a masquerade rule for traffic exiting via `fips0` so that + /// LAN client source addresses are rewritten to the gateway's mesh + /// address, allowing return traffic to route back through the mesh. + /// + /// `lan_interface` is the gateway's LAN-facing interface name, + /// needed by the port-forward LAN-side masquerade rule. + pub fn new(lan_interface: String) -> Result { + let mgr = Self::with_state(lan_interface); mgr.rebuild()?; info!("Created nftables table '{TABLE_NAME}' with NAT chains and fips0 masquerade"); @@ -167,133 +199,302 @@ impl NatManager { self.mappings.len() } - /// Atomically rebuild the entire nftables table with all current - /// rules. Deletes and recreates the table, chains, masquerade rule, - /// and all per-mapping DNAT/SNAT rules in a single netlink batch. - fn rebuild(&self) -> Result<(), NatError> { - // Delete existing table in a separate batch — ignore ENOENT on - // first call when the table doesn't exist yet. - let mut del_batch = Batch::new(); - del_batch.add(&self.table, MsgType::Del); - let _ = del_batch.send(); + /// The objects a rebuild sends, grouped into the batches that carry them. + /// + /// One batch, always. The kernel applies a batch as a single transaction, + /// so the table is deleted and recreated without ever leaving the packet + /// path, and a batch the kernel rejects leaves the previous table in + /// place. The leading `Add` is what makes the `Del` legal on a first run: + /// rustables sends a table `Add` with `NLM_F_CREATE` and no `NLM_F_EXCL`, + /// so it succeeds whether or not the table already exists and the `Del` + /// that follows always has a target. + fn rebuild_batches(&self) -> Vec> { + let mut ops = vec![ + NatOp::Table(MsgType::Add), + NatOp::Table(MsgType::Del), + NatOp::Table(MsgType::Add), + NatOp::PreChain, + NatOp::PostChain, + NatOp::FipsMasquerade, + ]; - // Recreate table, chains, and all rules atomically. - let mut batch = Batch::new(); - batch.add(&self.table, MsgType::Add); - batch.add(&self.pre_chain, MsgType::Add); - batch.add(&self.post_chain, MsgType::Add); - - // Masquerade rule: rewrite source address for traffic exiting fips0. - // Without this, LAN clients' source addresses (e.g. fd02::20) are - // not routable on the mesh, so return traffic would be black-holed. - let masq_rule = Rule::new(&self.post_chain)? - .with_expr(Meta::new(MetaType::OifName)) - .with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec())) - .with_expr(Masquerade::default()); - batch.add(&masq_rule, MsgType::Add); - - // Per-mapping DNAT/SNAT rules. for mapping in self.mappings.values() { - let dnat_rule = Rule::new(&self.pre_chain)? - .with_expr(Meta::new(MetaType::NfProto)) - .with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8])) - .with_expr( - HighLevelPayload::Network(NetworkHeaderField::IPv6(IPv6HeaderField::Daddr)) - .build(), - ) - .with_expr(Cmp::new(CmpOp::Eq, mapping.virtual_ip.octets())) - .with_expr(Immediate::new_data( - mapping.mesh_addr.octets().to_vec(), - Register::Reg1, - )) - .with_expr( - Nat::default() - .with_nat_type(NatType::DNat) - .with_family(ProtocolFamily::Ipv6) - .with_ip_register(Register::Reg1), - ); - batch.add(&dnat_rule, MsgType::Add); - - let snat_rule = Rule::new(&self.post_chain)? - .with_expr(Meta::new(MetaType::NfProto)) - .with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8])) - .with_expr( - HighLevelPayload::Network(NetworkHeaderField::IPv6(IPv6HeaderField::Saddr)) - .build(), - ) - .with_expr(Cmp::new(CmpOp::Eq, mapping.mesh_addr.octets())) - .with_expr(Immediate::new_data( - mapping.virtual_ip.octets().to_vec(), - Register::Reg1, - )) - .with_expr( - Nat::default() - .with_nat_type(NatType::SNat) - .with_family(ProtocolFamily::Ipv6) - .with_ip_register(Register::Reg1), - ); - batch.add(&snat_rule, MsgType::Add); + ops.push(NatOp::Dnat(mapping.virtual_ip)); + ops.push(NatOp::Snat(mapping.virtual_ip)); } - // Inbound port-forward rules. Each forward is - // one DNAT rule in prerouting keyed on (iif fips0, nfproto ipv6, - // l4proto, th dport). When any forwards are configured, emit a - // single LAN-side masquerade in postrouting so the LAN target - // host sees the gateway's LAN address as source and replies - // flow back through conntrack. - for pf in &self.port_forwards { - let l4proto: u8 = match pf.proto { - Proto::Tcp => libc::IPPROTO_TCP as u8, - Proto::Udp => libc::IPPROTO_UDP as u8, - }; - let dport_field = match pf.proto { - Proto::Tcp => TransportHeaderField::Tcp(TCPHeaderField::Dport), - Proto::Udp => TransportHeaderField::Udp(UDPHeaderField::Dport), - }; - let target_ip = *pf.target.ip(); - let target_port_be = pf.target.port().to_be_bytes(); - - let dnat_rule = Rule::new(&self.pre_chain)? - .with_expr(Meta::new(MetaType::IifName)) - .with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec())) - .with_expr(Meta::new(MetaType::NfProto)) - .with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8])) - .with_expr(Meta::new(MetaType::L4Proto)) - .with_expr(Cmp::new(CmpOp::Eq, [l4proto])) - .with_expr(HighLevelPayload::Transport(dport_field).build()) - .with_expr(Cmp::new(CmpOp::Eq, pf.listen_port.to_be_bytes().to_vec())) - .with_expr(Immediate::new_data( - target_ip.octets().to_vec(), - Register::Reg1, - )) - .with_expr(Immediate::new_data(target_port_be.to_vec(), Register::Reg2)) - .with_expr( - Nat::default() - .with_nat_type(NatType::DNat) - .with_family(ProtocolFamily::Ipv6) - .with_ip_register(Register::Reg1) - .with_port_register(Register::Reg2), - ); - batch.add(&dnat_rule, MsgType::Add); + // Inbound port-forward rules. Each forward is one DNAT rule in + // prerouting keyed on (iif fips0, nfproto ipv6, l4proto, th dport). + // When any forwards are configured, emit a single LAN-side masquerade + // in postrouting so the LAN target host sees the gateway's LAN address + // as source and replies flow back through conntrack. + for index in 0..self.port_forwards.len() { + ops.push(NatOp::PortForward(index)); } - if !self.port_forwards.is_empty() { - let mut lan_iface = self.lan_interface.clone().into_bytes(); - lan_iface.push(0); - let lan_masq = Rule::new(&self.post_chain)? - .with_expr(Meta::new(MetaType::IifName)) - .with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec())) - .with_expr(Meta::new(MetaType::OifName)) - .with_expr(Cmp::new(CmpOp::Eq, lan_iface)) - .with_expr(Meta::new(MetaType::NfProto)) - .with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8])) - .with_expr(Masquerade::default()); - batch.add(&lan_masq, MsgType::Add); + ops.push(NatOp::LanMasquerade); } - batch - .send() - .map_err(|e| NatError::Nftables(e.to_string()))?; + vec![ops] + } + + /// Build each op into its rustables object and send the batches in order. + fn send_batches(&self, batches: &[Vec]) -> Result<(), NatError> { + for ops in batches { + let mut batch = Batch::new(); + for op in ops { + match *op { + NatOp::Table(msg_type) => batch.add(&self.table, msg_type), + NatOp::PreChain => batch.add(&self.pre_chain, MsgType::Add), + NatOp::PostChain => batch.add(&self.post_chain, MsgType::Add), + NatOp::FipsMasquerade => { + // Rewrite the source address of traffic leaving fips0. + // Without this, LAN clients' source addresses (e.g. + // fd02::20) are not routable on the mesh, so return + // traffic would be black-holed. + let rule = Rule::new(&self.post_chain)? + .with_expr(Meta::new(MetaType::OifName)) + .with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec())) + .with_expr(Masquerade::default()); + batch.add(&rule, MsgType::Add); + } + NatOp::Dnat(virtual_ip) => { + let mapping = self.mapping(virtual_ip)?; + let rule = Rule::new(&self.pre_chain)? + .with_expr(Meta::new(MetaType::NfProto)) + .with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8])) + .with_expr( + HighLevelPayload::Network(NetworkHeaderField::IPv6( + IPv6HeaderField::Daddr, + )) + .build(), + ) + .with_expr(Cmp::new(CmpOp::Eq, mapping.virtual_ip.octets())) + .with_expr(Immediate::new_data( + mapping.mesh_addr.octets().to_vec(), + Register::Reg1, + )) + .with_expr( + Nat::default() + .with_nat_type(NatType::DNat) + .with_family(ProtocolFamily::Ipv6) + .with_ip_register(Register::Reg1), + ); + batch.add(&rule, MsgType::Add); + } + NatOp::Snat(virtual_ip) => { + let mapping = self.mapping(virtual_ip)?; + let rule = Rule::new(&self.post_chain)? + .with_expr(Meta::new(MetaType::NfProto)) + .with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8])) + .with_expr( + HighLevelPayload::Network(NetworkHeaderField::IPv6( + IPv6HeaderField::Saddr, + )) + .build(), + ) + .with_expr(Cmp::new(CmpOp::Eq, mapping.mesh_addr.octets())) + .with_expr(Immediate::new_data( + mapping.virtual_ip.octets().to_vec(), + Register::Reg1, + )) + .with_expr( + Nat::default() + .with_nat_type(NatType::SNat) + .with_family(ProtocolFamily::Ipv6) + .with_ip_register(Register::Reg1), + ); + batch.add(&rule, MsgType::Add); + } + NatOp::PortForward(index) => { + let pf = self.port_forwards.get(index).expect( + "rebuild_batches only emits indices it read from port_forwards", + ); + let l4proto: u8 = match pf.proto { + Proto::Tcp => libc::IPPROTO_TCP as u8, + Proto::Udp => libc::IPPROTO_UDP as u8, + }; + let dport_field = match pf.proto { + Proto::Tcp => TransportHeaderField::Tcp(TCPHeaderField::Dport), + Proto::Udp => TransportHeaderField::Udp(UDPHeaderField::Dport), + }; + let target_ip = *pf.target.ip(); + let target_port_be = pf.target.port().to_be_bytes(); + + let rule = Rule::new(&self.pre_chain)? + .with_expr(Meta::new(MetaType::IifName)) + .with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec())) + .with_expr(Meta::new(MetaType::NfProto)) + .with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8])) + .with_expr(Meta::new(MetaType::L4Proto)) + .with_expr(Cmp::new(CmpOp::Eq, [l4proto])) + .with_expr(HighLevelPayload::Transport(dport_field).build()) + .with_expr(Cmp::new(CmpOp::Eq, pf.listen_port.to_be_bytes().to_vec())) + .with_expr(Immediate::new_data( + target_ip.octets().to_vec(), + Register::Reg1, + )) + .with_expr(Immediate::new_data(target_port_be.to_vec(), Register::Reg2)) + .with_expr( + Nat::default() + .with_nat_type(NatType::DNat) + .with_family(ProtocolFamily::Ipv6) + .with_ip_register(Register::Reg1) + .with_port_register(Register::Reg2), + ); + batch.add(&rule, MsgType::Add); + } + NatOp::LanMasquerade => { + let mut lan_iface = self.lan_interface.clone().into_bytes(); + lan_iface.push(0); + let rule = Rule::new(&self.post_chain)? + .with_expr(Meta::new(MetaType::IifName)) + .with_expr(Cmp::new(CmpOp::Eq, b"fips0\0".to_vec())) + .with_expr(Meta::new(MetaType::OifName)) + .with_expr(Cmp::new(CmpOp::Eq, lan_iface)) + .with_expr(Meta::new(MetaType::NfProto)) + .with_expr(Cmp::new(CmpOp::Eq, [libc::NFPROTO_IPV6 as u8])) + .with_expr(Masquerade::default()); + batch.add(&rule, MsgType::Add); + } + } + } + + batch + .send() + .map_err(|e| NatError::Nftables(e.to_string()))?; + } Ok(()) } + + /// The mapping an op names, or the error a caller can report. + fn mapping(&self, virtual_ip: Ipv6Addr) -> Result<&NatMapping, NatError> { + self.mappings + .get(&virtual_ip) + .ok_or(NatError::RuleNotFound(virtual_ip)) + } + + /// Rebuild the entire nftables table with all current rules, in one + /// netlink transaction. + fn rebuild(&self) -> Result<(), NatError> { + self.send_batches(&self.rebuild_batches()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::net::SocketAddrV6; + + fn vip(last: u16) -> Ipv6Addr { + Ipv6Addr::new(0xfd01, 0, 0, 0, 0, 0, 0, last) + } + + fn mesh(last: u16) -> Ipv6Addr { + Ipv6Addr::new(0xfd02, 0, 0, 0, 0, 0, 0, last) + } + + /// A manager holding `count` mappings and no netlink socket. + fn manager_with_mappings(count: u16) -> NatManager { + let mut mgr = NatManager::with_state("br-lan".to_string()); + for i in 1..=count { + mgr.mappings.insert( + vip(i), + NatMapping { + virtual_ip: vip(i), + mesh_addr: mesh(i), + }, + ); + } + mgr + } + + #[test] + fn rebuild_deletes_and_recreates_the_table_inside_one_batch() { + let batches = manager_with_mappings(3).rebuild_batches(); + + assert_eq!( + batches.len(), + 1, + "a rebuild that sends the delete in a batch of its own leaves the \ + fips_gateway table absent between the two sends, so the gateway \ + has no NAT at all in that window: {batches:?}" + ); + assert_eq!( + batches[0][..3], + [ + NatOp::Table(MsgType::Add), + NatOp::Table(MsgType::Del), + NatOp::Table(MsgType::Add), + ], + "the delete needs a preceding add so it always has a target, and a \ + following add to recreate the table inside the same transaction" + ); + } + + #[test] + fn rebuild_deletes_the_table_exactly_once_and_before_every_rule() { + let batches = manager_with_mappings(2).rebuild_batches(); + let ops = &batches[0]; + + let deletes: Vec = ops + .iter() + .enumerate() + .filter(|(_, op)| matches!(op, NatOp::Table(MsgType::Del))) + .map(|(i, _)| i) + .collect(); + assert_eq!(deletes, vec![1], "the table is deleted once, at index 1"); + + // Everything that lives in the table has to be added after the delete + // and the recreate, or the delete would take it back out again. + for (index, op) in ops.iter().enumerate() { + if matches!(op, NatOp::Table(_)) { + continue; + } + assert!( + index > 2, + "{op:?} at index {index} would be removed by the table delete" + ); + } + } + + #[test] + fn rebuild_emits_a_dnat_and_an_snat_for_every_mapping() { + let ops = manager_with_mappings(3).rebuild_batches().remove(0); + + for i in 1..=3u16 { + assert!(ops.contains(&NatOp::Dnat(vip(i))), "no DNAT for {}", vip(i)); + assert!(ops.contains(&NatOp::Snat(vip(i))), "no SNAT for {}", vip(i)); + } + assert!(ops.contains(&NatOp::FipsMasquerade)); + assert!(!ops.contains(&NatOp::LanMasquerade), "no port forwards"); + } + + #[test] + fn rebuild_emits_the_lan_masquerade_once_when_port_forwards_exist() { + let mut mgr = manager_with_mappings(1); + mgr.port_forwards = vec![ + PortForward { + proto: Proto::Tcp, + listen_port: 8080, + target: SocketAddrV6::new(Ipv6Addr::LOCALHOST, 80, 0, 0), + }, + PortForward { + proto: Proto::Udp, + listen_port: 5353, + target: SocketAddrV6::new(Ipv6Addr::LOCALHOST, 53, 0, 0), + }, + ]; + + let ops = mgr.rebuild_batches().remove(0); + + assert!(ops.contains(&NatOp::PortForward(0))); + assert!(ops.contains(&NatOp::PortForward(1))); + assert_eq!( + ops.iter() + .filter(|op| matches!(op, NatOp::LanMasquerade)) + .count(), + 1 + ); + } } diff --git a/src/gateway/pool.rs b/src/gateway/pool.rs index 4b88bdff..900398d4 100644 --- a/src/gateway/pool.rs +++ b/src/gateway/pool.rs @@ -5,7 +5,7 @@ //! with conntrack to determine active sessions. use crate::NodeAddr; -use std::collections::{HashMap, VecDeque}; +use std::collections::{HashMap, HashSet, VecDeque}; use std::net::Ipv6Addr; use std::time::Instant; use tracing::{debug, info}; @@ -93,25 +93,127 @@ pub struct MappingInfo { pub last_ref_secs: u64, } -/// Trait for querying conntrack session counts. +/// Path the conntrack table is read from. +const CONNTRACK_PROC_PATH: &str = "/proc/net/nf_conntrack"; + +/// Active conntrack sessions counted by destination address. +/// +/// Taken once per tick, so the pool does a map lookup per mapping instead of +/// reading and scanning the whole conntrack table per mapping under its lock. +#[derive(Debug, Clone, Default)] +pub struct ConntrackSnapshot { + sessions: HashMap, +} + +impl ConntrackSnapshot { + /// Build a snapshot from counts already keyed by destination address. + pub fn from_counts(sessions: HashMap) -> Self { + Self { sessions } + } + + /// Sessions whose destination is `virtual_ip`, or zero if there are none. + pub fn sessions_for(&self, virtual_ip: Ipv6Addr) -> u32 { + self.sessions.get(&virtual_ip).copied().unwrap_or(0) + } + + /// Number of distinct destination addresses the snapshot saw. + pub fn len(&self) -> usize { + self.sessions.len() + } + + /// Whether the snapshot saw no sessions at all. + pub fn is_empty(&self) -> bool { + self.sessions.is_empty() + } +} + +/// Trait for taking a conntrack session snapshot. pub trait ConntrackQuerier: Send + Sync { - /// Returns the number of active conntrack entries whose original - /// destination matches the given virtual IP. - fn active_sessions(&self, virtual_ip: Ipv6Addr) -> Result; + /// Read the conntrack table once and count sessions by destination. + fn snapshot(&self) -> Result; } /// Conntrack querier that parses /proc/net/nf_conntrack. pub struct ProcConntrack; impl ConntrackQuerier for ProcConntrack { - fn active_sessions(&self, virtual_ip: Ipv6Addr) -> Result { - let content = std::fs::read_to_string("/proc/net/nf_conntrack")?; - let target = virtual_ip.to_string(); - let count = content - .lines() - .filter(|line| line.contains(&format!("dst={target}"))) - .count(); - Ok(count as u32) + fn snapshot(&self) -> Result { + let content = std::fs::read_to_string(CONNTRACK_PROC_PATH)?; + Ok(ConntrackSnapshot::from_counts(parse_conntrack(&content))) + } +} + +/// Count conntrack lines by the destination addresses they name. +/// +/// Every `dst=` value is parsed as an address and compared as an address. The +/// kernel prints tuples as `src=%pI6 dst=%pI6`, the full uncompressed form with +/// leading zeros, so a session to `fd01::1` is written +/// `dst=fd01:0000:0000:0000:0000:0000:0000:0001`; the previous code searched +/// each line for the address's compressed `Display` form, which cannot occur in +/// a fixed-width field, so it counted nothing on any kernel. +/// +/// A conntrack line carries the original and the reply tuple, each with its own +/// `dst=`, and the line is counted once per distinct address among them. That +/// keeps the meaning the count had before, which was "this line mentions the +/// address". A value that does not parse as an IPv6 address is skipped, which +/// is how IPv4 lines and any future field are ignored. +fn parse_conntrack(content: &str) -> HashMap { + let mut counts: HashMap = HashMap::new(); + let mut seen: HashSet = HashSet::new(); + + for line in content.lines() { + seen.clear(); + for token in line.split_whitespace() { + let Some(value) = token.strip_prefix("dst=") else { + continue; + }; + let Ok(addr) = value.parse::() else { + continue; + }; + seen.insert(addr); + } + for addr in &seen { + *counts.entry(*addr).or_insert(0) += 1; + } + } + + counts +} + +/// Whether a conntrack read outcome is new or a repeat of the last one. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ReadReport { + /// The outcome differs from the previous read, or is the first. + Changed, + /// The same outcome as the previous read. + Repeated, +} + +/// Remembers the last conntrack read outcome. +/// +/// A kernel built without `CONFIG_NF_CONNTRACK_PROCFS` has no +/// `/proc/net/nf_conntrack` at all, so every read fails the same way and a +/// per-tick warning would repeat for the life of the process. Warning on a +/// change of outcome still separates "the source is unreadable" from "there +/// are no sessions", which the pool could not distinguish before, without +/// filling the log. +#[derive(Debug, Default)] +pub struct ConntrackReadLog { + last: Option>, +} + +impl ConntrackReadLog { + /// Record a read outcome and say whether it is new. + /// + /// `None` is a successful read; `Some(kind)` is a failure of that kind. + pub fn observe(&mut self, outcome: Option) -> ReadReport { + let report = if self.last == Some(outcome) { + ReadReport::Repeated + } else { + ReadReport::Changed + }; + self.last = Some(outcome); + report } } @@ -168,6 +270,21 @@ impl VirtualIpPool { }) } + /// Refresh an existing mapping's TTL clock, never creating one. + /// + /// Returns whether a mapping for `node_addr` existed. A query the gateway + /// answers without an address still says the client is using the name, so + /// it must keep the mapping alive without minting one. + pub fn refresh_if_present(&mut self, node_addr: NodeAddr) -> bool { + match self.mappings.get_mut(&node_addr) { + Some(mapping) => { + mapping.last_referenced = Instant::now(); + true + } + None => false, + } + } + /// Allocate a virtual IP for the given node. Idempotent: returns /// existing mapping if one exists. pub fn allocate( @@ -176,9 +293,10 @@ impl VirtualIpPool { mesh_addr: Ipv6Addr, dns_name: &str, ) -> Result<(Ipv6Addr, bool), PoolError> { - // Idempotent: return existing mapping - if let Some(mapping) = self.mappings.get_mut(&node_addr) { - mapping.last_referenced = Instant::now(); + // Idempotent: return existing mapping, refreshed. + if self.refresh_if_present(node_addr) + && let Some(mapping) = self.mappings.get(&node_addr) + { return Ok((mapping.virtual_ip, false)); } @@ -215,15 +333,16 @@ impl VirtualIpPool { /// Periodic tick — drives state transitions. Returns events for /// the NAT and network modules. - pub fn tick(&mut self, now: Instant, conntrack: &dyn ConntrackQuerier) -> Vec { + pub fn tick(&mut self, now: Instant, conntrack: &ConntrackSnapshot) -> Vec { let mut events = Vec::new(); let mut to_free = Vec::new(); let ttl = std::time::Duration::from_secs(self.ttl_secs); let grace = std::time::Duration::from_secs(self.grace_secs); for (node_addr, mapping) in &mut self.mappings { - // Query conntrack for active sessions - let sessions = conntrack.active_sessions(mapping.virtual_ip).unwrap_or(0); + // One map lookup: the conntrack table was read once, before the + // pool lock was taken. + let sessions = conntrack.sessions_for(mapping.virtual_ip); mapping.session_count = sessions; // Live data-plane traffic pins the mapping: refresh the TTL @@ -372,26 +491,24 @@ fn parse_ipv6_cidr(cidr: &str) -> Result<(Ipv6Addr, u32), PoolError> { mod tests { use super::*; - /// Mock conntrack that returns a configurable session count. - struct MockConntrack { + /// Session counts a test sets directly, handed to `tick` as the snapshot + /// the tick task would have read from conntrack. + #[derive(Default)] + struct Sessions { counts: HashMap, } - impl MockConntrack { + impl Sessions { fn new() -> Self { - Self { - counts: HashMap::new(), - } + Self::default() } fn set(&mut self, addr: Ipv6Addr, count: u32) { self.counts.insert(addr, count); } - } - impl ConntrackQuerier for MockConntrack { - fn active_sessions(&self, virtual_ip: Ipv6Addr) -> Result { - Ok(*self.counts.get(&virtual_ip).unwrap_or(&0)) + fn snapshot(&self) -> ConntrackSnapshot { + ConntrackSnapshot::from_counts(self.counts.clone()) } } @@ -475,7 +592,7 @@ mod tests { #[test] fn test_mapping_lifecycle_allocated_to_free() { let mut pool = VirtualIpPool::new("fd01::/120", 1, 1).unwrap(); - let ct = MockConntrack::new(); + let ct = Sessions::new(); let node = make_node_addr(1); let mesh = make_mesh_addr(1); @@ -483,13 +600,13 @@ mod tests { // Tick before TTL — no change let now = Instant::now(); - let events = pool.tick(now, &ct); + let events = pool.tick(now, &ct.snapshot()); assert!(events.is_empty()); assert_eq!(pool.mappings.len(), 1); // Tick after TTL with no sessions — enters draining let later = now + std::time::Duration::from_secs(2); - let events = pool.tick(later, &ct); + let events = pool.tick(later, &ct.snapshot()); assert!(events.is_empty()); assert_eq!(pool.mappings.len(), 1); assert_eq!( @@ -499,7 +616,7 @@ mod tests { // Tick after grace period — freed let after_grace = later + std::time::Duration::from_secs(2); - let events = pool.tick(after_grace, &ct); + let events = pool.tick(after_grace, &ct.snapshot()); assert_eq!(events.len(), 1); assert!(matches!(events[0], PoolEvent::MappingRemoved { .. })); assert_eq!(pool.mappings.len(), 0); @@ -509,7 +626,7 @@ mod tests { #[test] fn test_mapping_lifecycle_active_draining_free() { let mut pool = VirtualIpPool::new("fd01::/120", 1, 1).unwrap(); - let mut ct = MockConntrack::new(); + let mut ct = Sessions::new(); let node = make_node_addr(1); let mesh = make_mesh_addr(1); @@ -518,25 +635,25 @@ mod tests { // Simulate active sessions ct.set(vip, 3); let now = Instant::now(); - let events = pool.tick(now, &ct); + let events = pool.tick(now, &ct.snapshot()); assert!(events.is_empty()); assert_eq!(pool.mappings[&node].state, MappingState::Active); // TTL expires after sessions drop to 0 → Draining let later = now + std::time::Duration::from_secs(2); ct.set(vip, 0); - let events = pool.tick(later, &ct); + let events = pool.tick(later, &ct.snapshot()); assert!(events.is_empty()); assert_eq!(pool.mappings[&node].state, MappingState::Draining); // Still draining, grace period not elapsed - let events = pool.tick(later, &ct); + let events = pool.tick(later, &ct.snapshot()); assert!(events.is_empty()); assert_eq!(pool.mappings[&node].state, MappingState::Draining); // Grace period elapsed → Free let much_later = later + std::time::Duration::from_secs(2); - let events = pool.tick(much_later, &ct); + let events = pool.tick(much_later, &ct.snapshot()); assert_eq!(events.len(), 1); assert!(matches!(events[0], PoolEvent::MappingRemoved { .. })); assert_eq!(pool.mappings.len(), 0); @@ -548,7 +665,7 @@ mod tests { // spanning well past the TTL must never be reclaimed and must // stay Active: live traffic refreshes last_referenced each tick. let mut pool = VirtualIpPool::new("fd01::/120", 1, 1).unwrap(); - let mut ct = MockConntrack::new(); + let mut ct = Sessions::new(); let node = make_node_addr(1); let mesh = make_mesh_addr(1); @@ -557,14 +674,14 @@ mod tests { let mut t = Instant::now(); // First tick activates the mapping. - let events = pool.tick(t, &ct); + let events = pool.tick(t, &ct.snapshot()); assert!(events.is_empty()); assert_eq!(pool.mappings[&node].state, MappingState::Active); // Advance many TTL-spans with continuous traffic. for _ in 0..10 { t += std::time::Duration::from_secs(5); // 5x the 1s TTL - let events = pool.tick(t, &ct); + let events = pool.tick(t, &ct.snapshot()); assert!(events.is_empty(), "mapping must not be reclaimed"); assert_eq!( pool.mappings[&node].state, @@ -580,7 +697,7 @@ mod tests { // Active -> drains when sessions hit 0 -> regains sessions before // grace elapses -> recovers to Active and is not freed. let mut pool = VirtualIpPool::new("fd01::/120", 1, 5).unwrap(); - let mut ct = MockConntrack::new(); + let mut ct = Sessions::new(); let node = make_node_addr(1); let mesh = make_mesh_addr(1); @@ -589,21 +706,21 @@ mod tests { // Activate with traffic. ct.set(vip, 1); let now = Instant::now(); - let events = pool.tick(now, &ct); + let events = pool.tick(now, &ct.snapshot()); assert!(events.is_empty()); assert_eq!(pool.mappings[&node].state, MappingState::Active); // TTL passes with sessions dropping to 0 -> Draining. let drained = now + std::time::Duration::from_secs(2); ct.set(vip, 0); - let events = pool.tick(drained, &ct); + let events = pool.tick(drained, &ct.snapshot()); assert!(events.is_empty()); assert_eq!(pool.mappings[&node].state, MappingState::Draining); // Traffic resumes before grace (5s) elapses -> recover to Active. let resumed = drained + std::time::Duration::from_secs(2); ct.set(vip, 3); - let events = pool.tick(resumed, &ct); + let events = pool.tick(resumed, &ct.snapshot()); assert!(events.is_empty()); assert_eq!(pool.mappings[&node].state, MappingState::Active); assert!(pool.mappings[&node].drain_start.is_none()); @@ -616,7 +733,7 @@ mod tests { // fresh drain_start so the full grace window is honored again, // not reclaimed immediately off a stale drain_start. let mut pool = VirtualIpPool::new("fd01::/120", 1, 5).unwrap(); - let mut ct = MockConntrack::new(); + let mut ct = Sessions::new(); let node = make_node_addr(1); let mesh = make_mesh_addr(1); @@ -625,36 +742,36 @@ mod tests { // Activate. ct.set(vip, 1); let now = Instant::now(); - pool.tick(now, &ct); + pool.tick(now, &ct.snapshot()); assert_eq!(pool.mappings[&node].state, MappingState::Active); // First drain. let first_drain = now + std::time::Duration::from_secs(2); ct.set(vip, 0); - pool.tick(first_drain, &ct); + pool.tick(first_drain, &ct.snapshot()); assert_eq!(pool.mappings[&node].state, MappingState::Draining); // Recover. let recover = first_drain + std::time::Duration::from_secs(2); ct.set(vip, 2); - pool.tick(recover, &ct); + pool.tick(recover, &ct.snapshot()); assert_eq!(pool.mappings[&node].state, MappingState::Active); // Second drain begins; drain_start must be re-stamped fresh. let second_drain = recover + std::time::Duration::from_secs(2); ct.set(vip, 0); - pool.tick(second_drain, &ct); + pool.tick(second_drain, &ct.snapshot()); assert_eq!(pool.mappings[&node].state, MappingState::Draining); // Just before the fresh grace window expires (5s): not reclaimed. let before_grace = second_drain + std::time::Duration::from_secs(4); - let events = pool.tick(before_grace, &ct); + let events = pool.tick(before_grace, &ct.snapshot()); assert!(events.is_empty(), "fresh grace window must be honored"); assert_eq!(pool.mappings.len(), 1); // After the fresh grace window: reclaimed. let after_grace = second_drain + std::time::Duration::from_secs(6); - let events = pool.tick(after_grace, &ct); + let events = pool.tick(after_grace, &ct.snapshot()); assert_eq!(events.len(), 1); assert!(matches!(events[0], PoolEvent::MappingRemoved { .. })); assert_eq!(pool.mappings.len(), 0); @@ -696,4 +813,99 @@ mod tests { let pool = VirtualIpPool::new("fd01::/96", 60, 60).unwrap(); assert_eq!(pool.total, 65535); // 2^16 - 1 (skip addr 0) } + + /// A conntrack line in the form the kernel prints. + /// + /// Built from the kernel's own format string, not captured from a running + /// kernel: `net/netfilter/nf_conntrack_standalone.c` prints each tuple with + /// `"src=%pI6 dst=%pI6 "`, and `%pI6` is the full uncompressed form with + /// leading zeros (`Documentation/core-api/printk-formats.rst`). Both were + /// read at v6.8. The host this was written on has no + /// `/proc/net/nf_conntrack` to capture from, because its kernel is built + /// without `CONFIG_NF_CONNTRACK_PROCFS`; OpenWrt's generic kernel config + /// sets it, which is the kernel this parser exists for. + const KERNEL_LINE: &str = "ipv6 10 tcp 6 431999 ESTABLISHED \ + src=fd02:0000:0000:0000:0000:0000:0000:0020 \ + dst=fd01:0000:0000:0000:0000:0000:0000:0001 sport=45678 dport=8000 \ + src=fd01:0000:0000:0000:0000:0000:0000:0001 \ + dst=fd02:0000:0000:0000:0000:0000:0000:0020 sport=8000 dport=45678 \ + [ASSURED] mark=0 use=1"; + + #[test] + fn conntrack_parse_counts_a_kernel_format_line_for_its_virtual_ip() { + let counts = parse_conntrack(KERNEL_LINE); + let virtual_ip: Ipv6Addr = "fd01::1".parse().unwrap(); + + assert_eq!( + counts.get(&virtual_ip).copied().unwrap_or(0), + 1, + "the kernel writes the uncompressed form, so matching on the \ + address's compressed Display form counts nothing" + ); + + // Healthy path: a different address in the same pool is not counted. + let other: Ipv6Addr = "fd01::10".parse().unwrap(); + assert_eq!(counts.get(&other).copied().unwrap_or(0), 0); + } + + #[test] + fn conntrack_parse_counts_a_line_once_however_many_tuples_name_the_address() { + // A hairpin flow: the address is the destination of both tuples. + let line = "ipv6 10 udp 17 29 \ + src=fd01:0000:0000:0000:0000:0000:0000:0001 \ + dst=fd01:0000:0000:0000:0000:0000:0000:0001 sport=1 dport=2 \ + src=fd01:0000:0000:0000:0000:0000:0000:0001 \ + dst=fd01:0000:0000:0000:0000:0000:0000:0001 sport=2 dport=1 \ + mark=0 use=1"; + let counts = parse_conntrack(line); + let virtual_ip: Ipv6Addr = "fd01::1".parse().unwrap(); + + assert_eq!(counts.get(&virtual_ip).copied().unwrap_or(0), 1); + } + + #[test] + fn conntrack_parse_counts_each_line_that_names_the_address() { + let content = format!("{KERNEL_LINE}\n{KERNEL_LINE}\n"); + let counts = parse_conntrack(&content); + let virtual_ip: Ipv6Addr = "fd01::1".parse().unwrap(); + + assert_eq!(counts.get(&virtual_ip).copied().unwrap_or(0), 2); + } + + #[test] + fn conntrack_parse_skips_a_value_that_is_not_an_ipv6_address() { + let content = "ipv4 2 tcp 6 431999 ESTABLISHED src=192.0.2.1 \ + dst=192.0.2.2 sport=1 dport=2 mark=0 use=1\n"; + + assert!(parse_conntrack(content).is_empty()); + } + + #[test] + fn conntrack_snapshot_reads_zero_for_an_address_it_did_not_see() { + let snapshot = ConntrackSnapshot::from_counts(parse_conntrack(KERNEL_LINE)); + + assert_eq!(snapshot.sessions_for("fd01::1".parse().unwrap()), 1); + assert_eq!(snapshot.sessions_for("fd01::99".parse().unwrap()), 0); + assert!(ConntrackSnapshot::default().is_empty()); + } + + #[test] + fn conntrack_read_log_warns_on_a_new_outcome_and_not_on_a_repeat() { + use std::io::ErrorKind; + + let mut log = ConntrackReadLog::default(); + + // The sequence a kernel without the proc file produces, then a source + // that comes back, then fails again. + assert_eq!(log.observe(Some(ErrorKind::NotFound)), ReadReport::Changed); + assert_eq!(log.observe(Some(ErrorKind::NotFound)), ReadReport::Repeated); + assert_eq!(log.observe(None), ReadReport::Changed); + assert_eq!(log.observe(None), ReadReport::Repeated); + assert_eq!(log.observe(Some(ErrorKind::NotFound)), ReadReport::Changed); + assert_eq!( + log.observe(Some(ErrorKind::PermissionDenied)), + ReadReport::Changed, + "a different failure is a different outcome and is worth a line" + ); + } } diff --git a/testing/ci-local.sh b/testing/ci-local.sh index 5a4b6a4f..76d5f4cb 100755 --- a/testing/ci-local.sh +++ b/testing/ci-local.sh @@ -192,6 +192,7 @@ CHAOS_SUITES=( # on disk and it remains runnable by hand via # testing/chaos/scripts/chaos.sh bloom-storm. GATEWAY_SUITES=(gateway) +OPENWRT_SUITES=(openwrt-scripts) SIDECAR_SUITES=(sidecar) FIREWALL_SUITES=(firewall) NAT_SUITES=(cone symmetric lan) @@ -229,6 +230,9 @@ list_suites() { echo " Gateway:" for s in "${GATEWAY_SUITES[@]}"; do echo " $s"; done echo "" + echo " OpenWrt packaging:" + for s in "${OPENWRT_SUITES[@]}"; do echo " $s"; done + echo "" echo " Firewall baseline:" for s in "${FIREWALL_SUITES[@]}"; do echo " $s"; done echo "" @@ -1098,6 +1102,16 @@ run_tor_directory() { run_integration() { stage "Stage 3: Integration Tests" + # First, and before the build context: the OpenWrt scenarios need no FIPS + # binary and no test image, so a packaging regression is reported in + # seconds rather than after the image build. + if [[ -z "$ONLY_SUITE" ]]; then + run_openwrt_scripts + elif [[ "$ONLY_SUITE" == "openwrt-scripts" ]]; then + run_openwrt_scripts + return + fi + # Populate THIS run's build context, then install the binaries into it. # Everything but the binaries is copied from the tracked context directory; # the binaries are installed fresh, and a previous run's are deliberately @@ -1256,6 +1270,8 @@ run_suite() { run_static "${suite#static-}" ;; gateway) run_gateway ;; + openwrt-scripts) + run_openwrt_scripts ;; firewall) run_firewall ;; nat-cone|nat-symmetric|nat-lan) @@ -1336,6 +1352,19 @@ print_summary() { # Verify the local default suite set and the GitHub matrix still cover the # same work. Runs first: it takes about a second, and a divergence should be # reported before a half-hour suite rather than after it. +# The OpenWrt maintainer scripts and the fips-gateway init script ship to +# routers and run there under ash, never under bash. This runs them under ash +# in a busybox container against stubbed init scripts, so an install, an +# upgrade from either generation of the package, and a removal each assert what +# the package left enabled and running. +run_openwrt_scripts() { + local rc=0 + info "[openwrt-scripts] Running the OpenWrt maintainer-script scenarios" + bash "$SCRIPT_DIR/openwrt/maintainer-scripts-test.sh" || rc=$? + record "openwrt-scripts" $rc + return $rc +} + run_ci_parity() { local rc=0 info "[ci-parity] Comparing the local suite set against the GitHub matrix" diff --git a/testing/openwrt/fixtures/released-prerm b/testing/openwrt/fixtures/released-prerm new file mode 100644 index 00000000..29c20988 --- /dev/null +++ b/testing/openwrt/fixtures/released-prerm @@ -0,0 +1,6 @@ +#!/bin/sh +/etc/init.d/fips-gateway stop 2>/dev/null || true +/etc/init.d/fips-gateway disable 2>/dev/null || true +/etc/init.d/fips stop 2>/dev/null || true +/etc/init.d/fips disable 2>/dev/null || true +exit 0 diff --git a/testing/openwrt/maintainer-scripts-test.sh b/testing/openwrt/maintainer-scripts-test.sh new file mode 100755 index 00000000..8a8d4eea --- /dev/null +++ b/testing/openwrt/maintainer-scripts-test.sh @@ -0,0 +1,46 @@ +#!/bin/bash +# ── OpenWrt maintainer-script scenarios ───────────────────────────────────── +# Runs testing/openwrt/scenarios.sh inside a busybox container, so the package +# scripts and the fips-gateway init script are interpreted by ash rather than +# by the host's bash or dash. The scripts ship to routers and are only ever run +# under ash there; a construct bash accepts and ash does not would otherwise +# surface on a router. +# +# The container is the only reason docker is needed: the scenarios touch no +# network and no FIPS binary, and they do not use the shared test image. +# +# Exit 0 = every scenario passed. Exit 1 = at least one failed. Exit 2 = the +# harness could not run; never treated as a pass. +# ───────────────────────────────────────────────────────────────────────────── +set -uo pipefail + +SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd)" +PROJECT_ROOT="$(cd "$SCRIPT_DIR/../.." && pwd)" + +# Pinned rather than :latest so the shell under test does not change under a +# run. Overridable for trying another ash build. +IMAGE="${OPENWRT_ASH_IMAGE:-busybox:1.37}" + +if ! command -v docker >/dev/null 2>&1; then + echo "openwrt-scripts: docker not found; cannot run the ash scenarios" >&2 + exit 2 +fi + +if [[ ! -f "$SCRIPT_DIR/scenarios.sh" ]]; then + echo "openwrt-scripts: missing $SCRIPT_DIR/scenarios.sh" >&2 + exit 2 +fi + +docker run --rm --network none \ + -v "$PROJECT_ROOT:/src:ro" \ + -e REPO=/src \ + -e "POSTINST=${POSTINST:-}" \ + -e "PRERM=${PRERM:-}" \ + "$IMAGE" sh /src/testing/openwrt/scenarios.sh +rc=$? + +if [[ $rc -ne 0 && $rc -ne 1 ]]; then + echo "openwrt-scripts: the container exited $rc, so the scenarios did not report" >&2 + exit 2 +fi +exit $rc diff --git a/testing/openwrt/scenarios.sh b/testing/openwrt/scenarios.sh new file mode 100755 index 00000000..b060d0b1 --- /dev/null +++ b/testing/openwrt/scenarios.sh @@ -0,0 +1,342 @@ +#!/bin/sh +# OpenWrt maintainer-script and init-guard scenarios, run under ash. +# +# Driven by testing/openwrt/maintainer-scripts-test.sh, which starts a busybox +# container so /bin/sh here is ash, the shell OpenWrt runs these scripts under. +# Nothing in this file needs opkg: the call order, the arguments and the +# PKG_UPGRADE environment are taken from opkg-lede's own sources, so what is +# exercised is the scripts' behaviour given that contract, not opkg itself. +# A real `opkg upgrade` on a router image stays uncovered. +# +# POSTINST and PRERM may be pointed at other files. That is the seam used to +# see a scenario red against the previously released scripts, and to re-break +# the fixed ones during a break-check. + +set -u + +REPO="${REPO:-/src}" +POSTINST="${POSTINST:-$REPO/packaging/openwrt-ipk/scripts/postinst}" +PRERM="${PRERM:-$REPO/packaging/openwrt-ipk/scripts/prerm}" +RELEASED_PRERM="$REPO/testing/openwrt/fixtures/released-prerm" +INIT_GATEWAY="$REPO/packaging/openwrt-ipk/files/etc/init.d/fips-gateway" +SHIPPED_YAML="$REPO/packaging/openwrt-ipk/files/etc/fips/fips.yaml" + +WORK=/tmp/fips-openwrt-scenarios +UPGRADE_MARKER=/tmp/fips-prerm-upgrade + +FAILURES=0 +CASES=0 + +note() { echo " $*"; } + +ok() { + CASES=$((CASES + 1)) + echo " ok $*" + return 0 +} + +bad() { + CASES=$((CASES + 1)) + FAILURES=$((FAILURES + 1)) + echo " FAIL $*" + return 0 +} + +# Stub init scripts that record every call and keep an enable state file, so a +# scenario can assert both what was invoked and what the package left behind. +install_stubs() { + mkdir -p /etc/init.d /etc/uci-defaults + + cat > /etc/init.d/fips-gateway <<'STUB' +#!/bin/sh +echo "fips-gateway $1" >> "$CALLS" +case "$1" in + enable) echo 1 > "$GW_STATE" ;; + disable) echo 0 > "$GW_STATE" ;; + enabled) [ "$(cat "$GW_STATE")" = 1 ] ;; +esac +STUB + + cat > /etc/init.d/fips <<'STUB' +#!/bin/sh +echo "fips $1" >> "$CALLS" +case "$1" in + enable) echo 1 > "$FIPS_STATE" ;; + disable) echo 0 > "$FIPS_STATE" ;; + enabled) [ "$(cat "$FIPS_STATE")" = 1 ] ;; +esac +STUB + + cat > /etc/uci-defaults/90-fips-setup <<'STUB' +#!/bin/sh +echo "uci-defaults" >> "$CALLS" +STUB + + chmod 0755 /etc/init.d/fips-gateway /etc/init.d/fips /etc/uci-defaults/90-fips-setup + return 0 +} + +reset_state() { + rm -rf "$WORK" + mkdir -p "$WORK" + CALLS="$WORK/calls" + GW_STATE="$WORK/gateway-enabled" + FIPS_STATE="$WORK/fips-enabled" + export CALLS GW_STATE FIPS_STATE + : > "$CALLS" + echo 0 > "$GW_STATE" + echo 0 > "$FIPS_STATE" + rm -f "$UPGRADE_MARKER" + unset PKG_UPGRADE + install_stubs + return 0 +} + +calls_oneline() { + tr '\n' ';' < "$CALLS" + return 0 +} + +assert_called() { + # assert_called + if grep -qxF "$1" "$CALLS"; then + ok "$2" + else + bad "$2 — '$1' is not among: $(calls_oneline)" + fi + return 0 +} + +assert_not_called() { + if grep -qxF "$1" "$CALLS"; then + bad "$2 — '$1' was called: $(calls_oneline)" + else + ok "$2" + fi + return 0 +} + +assert_file_is() { + # assert_file_is + got="$(cat "$1" 2>/dev/null)" + if [ "$got" = "$2" ]; then + ok "$3" + else + bad "$3 — expected '$2', got '$got'" + fi + return 0 +} + +assert_equals() { + # assert_equals + if [ "$1" = "$2" ]; then + ok "$3" + else + bad "$3 — expected '$2', got '$1'" + fi + return 0 +} + +assert_absent() { + if [ -e "$1" ]; then + bad "$2 — $1 still exists" + else + ok "$2" + fi + return 0 +} + +# ── 1. Fresh install ──────────────────────────────────────────────────────── +# opkg runs the postinst with "configure"; PKG_UPGRADE is set only on upgrades, +# so both its absence and an explicit 0 must leave the gateway alone. +scenario_fresh_install() { + for pkg_upgrade in unset 0; do + note "scenario 1: fresh install (PKG_UPGRADE $pkg_upgrade)" + reset_state + if [ "$pkg_upgrade" = "0" ]; then + PKG_UPGRADE=0 sh "$POSTINST" configure >/dev/null 2>&1 + else + sh "$POSTINST" configure >/dev/null 2>&1 + fi + + assert_called "fips enable" "the daemon is enabled on a fresh install" + assert_called "fips start" "the daemon is started on a fresh install" + assert_not_called "fips-gateway enable" "the gateway is not enabled on a fresh install" + assert_not_called "fips-gateway start" "the gateway is not started on a fresh install" + assert_file_is "$GW_STATE" "0" "the gateway is left disabled on a fresh install" + done + return 0 +} + +# ── 2. Upgrade from a released package ────────────────────────────────────── +# Its prerm disabled the gateway on its way out and left no marker, so the +# incoming postinst cannot tell an enabled gateway from a disabled one and +# re-enables it. +scenario_upgrade_from_released() { + note "scenario 2: upgrade from a released package" + reset_state + echo 1 > "$GW_STATE" + echo 1 > "$FIPS_STATE" + + sh "$RELEASED_PRERM" upgrade 0.5.1 >/dev/null 2>&1 + PKG_UPGRADE=1 sh "$POSTINST" configure >/dev/null 2>&1 + + assert_called "fips-gateway enable" "the gateway is re-enabled after a released prerm disabled it" + assert_called "fips-gateway start" "the gateway is started again" + assert_file_is "$GW_STATE" "1" "the gateway ends up enabled" + return 0 +} + +# ── 3. Upgrade from a package carrying these scripts, gateway enabled ─────── +scenario_upgrade_enabled() { + note "scenario 3: upgrade from these scripts, gateway enabled" + reset_state + echo 1 > "$GW_STATE" + echo 1 > "$FIPS_STATE" + + sh "$PRERM" upgrade 0.5.2 >/dev/null 2>&1 + assert_file_is "$GW_STATE" "1" "the outgoing prerm does not disable the gateway on an upgrade" + assert_not_called "fips-gateway disable" "the outgoing prerm does not call disable on an upgrade" + assert_called "fips-gateway stop" "the outgoing prerm still stops the gateway" + + PKG_UPGRADE=1 sh "$POSTINST" configure >/dev/null 2>&1 + assert_file_is "$GW_STATE" "1" "the gateway stays enabled across the upgrade" + assert_called "fips-gateway start" "an enabled gateway is started again" + assert_not_called "fips-gateway enable" "an enabled gateway does not need re-enabling" + assert_absent "$UPGRADE_MARKER" "the postinst removes the upgrade marker" + return 0 +} + +# ── 4. Upgrade from a package carrying these scripts, gateway disabled ────── +scenario_upgrade_disabled() { + note "scenario 4: upgrade from these scripts, gateway disabled" + reset_state + echo 1 > "$FIPS_STATE" + + sh "$PRERM" upgrade 0.5.2 >/dev/null 2>&1 + PKG_UPGRADE=1 sh "$POSTINST" configure >/dev/null 2>&1 + + assert_file_is "$GW_STATE" "0" "a disabled gateway stays disabled across the upgrade" + assert_not_called "fips-gateway enable" "a disabled gateway is not enabled by the upgrade" + assert_not_called "fips-gateway start" "a disabled gateway is not started by the upgrade" + assert_absent "$UPGRADE_MARKER" "the postinst removes the upgrade marker" + return 0 +} + +# ── 5. Removal ────────────────────────────────────────────────────────────── +scenario_removal() { + note "scenario 5: removal" + reset_state + echo 1 > "$GW_STATE" + echo 1 > "$FIPS_STATE" + + sh "$PRERM" remove >/dev/null 2>&1 + + assert_called "fips-gateway stop" "removal stops the gateway" + assert_called "fips-gateway disable" "removal disables the gateway" + assert_called "fips stop" "removal stops the daemon" + assert_called "fips disable" "removal disables the daemon" + assert_file_is "$GW_STATE" "0" "the gateway ends up disabled" + assert_absent "$UPGRADE_MARKER" "removal leaves no upgrade marker" + return 0 +} + +# ── 6. gateway_config_enabled reads the config ────────────────────────────── +scenario_config_reader() { + note "scenario 6: gateway_config_enabled" + reset_state + + # shellcheck source=/dev/null + . "$INIT_GATEWAY" + + CONFIG="$SHIPPED_YAML" + assert_equals "$(gateway_config_enabled)" "true" "the shipped fips.yaml reads as true" + + CONFIG="$WORK/disabled.yaml" + cat > "$CONFIG" <<'YAML' +identity: + key_file: "/etc/fips/node.key" + +gateway: + enabled: false + pool: "fd01::/112" + +peers: [] +YAML + assert_equals "$(gateway_config_enabled)" "false" "an explicitly disabled gateway reads as false" + + CONFIG="$WORK/no-gateway.yaml" + cat > "$CONFIG" <<'YAML' +identity: + key_file: "/etc/fips/node.key" + +dns: + enabled: true + +peers: [] +YAML + assert_equals "$(gateway_config_enabled)" "" "a config with no gateway block reads as empty" + return 0 +} + +# ── 7. start_service refuses to touch dnsmasq for a disabled gateway ──────── +# The init script's helpers are redefined after sourcing it, so start_service +# runs its own decision against recorded stubs instead of uci, procd and the +# network. +scenario_start_service_guard() { + note "scenario 7: start_service guard" + + # shellcheck source=/dev/null + . "$INIT_GATEWAY" + + sysctl() { return 0; } + modprobe() { return 0; } + logger() { return 0; } + sleep() { return 0; } + procd_set_param() { return 0; } + procd_close_instance() { return 0; } + dnsmasq_swap_fips_upstream() { echo "dnsmasq_swap $1" >> "$CALLS"; return 0; } + gateway_add_global_prefix() { echo "add_global_prefix" >> "$CALLS"; return 0; } + gateway_add_ra_route() { echo "add_ra_route" >> "$CALLS"; return 0; } + procd_open_instance() { echo "procd_open_instance" >> "$CALLS"; return 0; } + + reset_state + CONFIG="$SHIPPED_YAML" + start_service >/dev/null 2>&1 + assert_called "dnsmasq_swap 5353" "an enabled gateway still redirects dnsmasq" + assert_called "procd_open_instance" "an enabled gateway still starts the daemon" + + reset_state + CONFIG="$WORK/disabled.yaml" + cat > "$CONFIG" <<'YAML' +gateway: + enabled: false + pool: "fd01::/112" +YAML + start_service >/dev/null 2>&1 + assert_not_called "dnsmasq_swap 5353" "a disabled gateway does not redirect dnsmasq" + assert_not_called "add_global_prefix" "a disabled gateway does not add the LAN prefix" + assert_not_called "add_ra_route" "a disabled gateway does not advertise the pool route" + assert_not_called "procd_open_instance" "a disabled gateway does not start the daemon" + return 0 +} + +echo "OpenWrt maintainer-script scenarios (shell: $(readlink -f /proc/$$/exe 2>/dev/null || echo sh))" +echo " postinst: $POSTINST" +echo " prerm: $PRERM" + +scenario_fresh_install +scenario_upgrade_from_released +scenario_upgrade_enabled +scenario_upgrade_disabled +scenario_removal +scenario_config_reader +scenario_start_service_guard + +echo "" +if [ "$FAILURES" -eq 0 ]; then + echo "openwrt-scripts: all $CASES checks passed" + exit 0 +fi +echo "openwrt-scripts: $FAILURES of $CASES checks failed" +exit 1 From 9f09595dfa41b6e6b753df7badb788514be9a856 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 17 Sep 2026 21:24:16 +0000 Subject: [PATCH 6/6] fix(transport/tcp): key inbound pool entries by the connection four-tuple The kernel names a connection by its four-tuple, so a listener on a wildcard address, which is what the shipped configuration binds, can accept two connections whose peer `ip:port` is the same on two different local addresses. The pool was keyed by the peer address alone, so the second connection's entry replaced the first's while the inbound-connection counter counted both. That counter gates the inbound connection limit, so a host repeating the collision could hold it at the limit and lock out further inbound TCP connections until the daemon restarted. Inbound entries now carry the accepted socket's local address in their pool key as well as the remote one; outbound entries carry no local address, since nothing distinguishes two outbound connections to one peer and leaving it out keeps the connect-on-send lookup a single hash probe. Callers that send, close or query by peer address know only the remote, so `key_for_remote` recovers the key: an outbound entry is probed first, and inbound entries are found by scanning for the remote and taking the most recently established, which is the connection a peer that reconnected is using. `remove_own` is now generic over the pool's key type. It is shared with the SOCKS5 transport, which still pools by peer address and passes a `TransportAddr` unchanged. The identity check it performs stays necessary after the key is made more specific, because a peer can still reconnect on the same four-tuple, so a connection's own loops continue to remove an entry only when it carries their id. An accepted socket whose local address cannot be read has no four-tuple to be keyed by. It is dropped rather than pooled under a key that could collide, and the drop is counted as a rejected connection so it is visible in the transport's stats rather than only in the log. The scenario test for two inbound connections sharing a remote address recorded a known gap: that the two still shared one pool key, so the second accept replaced the first entry without stopping its tasks while counting a second inbound slot. That gap is what this closes, so the test now asserts two entries carrying the two four-tuples and an inbound counter that returns to zero, and its comment records the closure. A reply test covers the other half: a caller answering a received packet knows only the remote address, so the pool has to resolve it to the inbound entry rather than falling through to connect-on-send against the peer's ephemeral port. Break-checked by making the inbound key ignore its local address, which reds the scenario test on the entry count. --- src/transport/stream.rs | 28 +++-- src/transport/tcp/mod.rs | 236 ++++++++++++++++++++++++++++---------- src/transport/tcp/pool.rs | 66 ++++++++++- 3 files changed, 252 insertions(+), 78 deletions(-) diff --git a/src/transport/stream.rs b/src/transport/stream.rs index 64a99bb1..90cfd340 100644 --- a/src/transport/stream.rs +++ b/src/transport/stream.rs @@ -11,8 +11,6 @@ use std::time::Duration; use portable_atomic::{AtomicU64, Ordering}; use tokio::task::JoinHandle; -use crate::transport::TransportAddr; - /// Identity of one pooled stream connection. /// /// The pool is keyed by address, and a newer connection can take an address @@ -36,20 +34,26 @@ pub(crate) trait PooledConn { fn conn_id(&self) -> ConnId; } -/// Remove the entry at `addr`, but only if it is connection `id`. +/// Remove the entry at `key`, but only if it is connection `id`. /// /// This is the only way a connection's own writer or receive loop removes a -/// pool entry. An entry with another id belongs to a newer connection at the -/// same address, and is left alone. -pub(crate) fn remove_own( - pool: &mut HashMap, - addr: &TransportAddr, - id: ConnId, -) -> Option { - if pool.get(addr)?.conn_id() != id { +/// pool entry. An entry with another id belongs to a newer connection under +/// the same key, and is left alone. +/// +/// The key type is the pool's own: a transport that pools by peer address +/// passes a `TransportAddr`, and one that pools by four-tuple passes its own +/// key. The identity check is the same either way, and it stays necessary +/// after a key is made more specific, because a peer can still reconnect on +/// the same four-tuple. +pub(crate) fn remove_own(pool: &mut HashMap, key: &K, id: ConnId) -> Option +where + K: std::hash::Hash + Eq, + C: PooledConn, +{ + if pool.get(key)?.conn_id() != id { return None; } - pool.remove(addr) + pool.remove(key) } /// How long a deliberately closed connection's writer may keep writing the diff --git a/src/transport/tcp/mod.rs b/src/transport/tcp/mod.rs index 52b20e21..d4a8f54c 100644 --- a/src/transport/tcp/mod.rs +++ b/src/transport/tcp/mod.rs @@ -12,9 +12,9 @@ //! ## Architecture //! //! Unlike UDP (one socket serves all peers), TCP requires one `TcpStream` -//! per peer. The transport maintains a connection pool mapping -//! `TransportAddr` to per-connection state, plus an optional `TcpListener` -//! for inbound connections. +//! per peer. The transport maintains a connection pool mapping each +//! connection's four-tuple to its per-connection state, plus an optional +//! `TcpListener` for inbound connections. //! //! ## Framing //! @@ -35,7 +35,10 @@ use crate::transport::framing::read_fmp_packet; use crate::transport::stream::{ ConnId, WRITER_DRAIN_TIMEOUT, drain_writer, next_conn_id, remove_own, }; -use pool::{ConnectingEntry, ConnectingPool, ConnectionPool, Direction, TcpConnection}; +use pool::{ + ConnectingEntry, ConnectingPool, ConnectionPool, Direction, PoolKey, TcpConnection, + key_for_remote, +}; use stats::TcpStats; use futures::FutureExt; @@ -59,7 +62,7 @@ use tracing::{debug, info, trace, warn}; /// /// Provides connection-oriented, reliable byte stream delivery over TCP/IP. /// Each peer has its own TCP connection; links are managed per-connection -/// with a connection pool keyed by `TransportAddr`. +/// with a connection pool keyed by `PoolKey`, the connection's four-tuple. pub struct TcpTransport { /// Unique transport identifier. transport_id: TransportId, @@ -271,7 +274,7 @@ impl TcpTransport { // aborting the task skips that path; decrement explicitly here // using the direction we stored on the connection record. let mut pool = self.pool.lock().await; - for (addr, conn) in pool.drain() { + for (key, conn) in pool.drain() { conn.recv_task.abort(); conn.send_task.abort(); let _ = conn.recv_task.await; @@ -281,7 +284,7 @@ impl TcpTransport { } debug!( transport_id = %self.transport_id, - remote_addr = %addr, + remote_addr = %key.remote, direction = ?conn.direction, "TCP connection closed (transport stopping)" ); @@ -331,7 +334,7 @@ impl TcpTransport { // must not be able to await the wire (see `tcp_send_loop`). let send_tx = { let pool = self.pool.lock().await; - pool.get(addr).map(|c| c.send_tx.clone()) + key_for_remote(&pool, addr).and_then(|key| pool.get(&key).map(|c| c.send_tx.clone())) }; let send_tx = match send_tx { @@ -426,7 +429,9 @@ impl TcpTransport { let packet_tx = self.packet_tx.clone(); let pool = self.pool.clone(); let recv_stats = self.stats.clone(); - let remote_addr = addr.clone(); + let key = PoolKey::outbound(addr.clone()); + let recv_key = key.clone(); + let send_key = key.clone(); let mtu = mss_mtu; let id = next_conn_id(); @@ -434,7 +439,7 @@ impl TcpTransport { tcp_receive_loop( read_half, transport_id, - remote_addr.clone(), + recv_key, id, packet_tx, pool, @@ -454,7 +459,7 @@ impl TcpTransport { write_half, send_rx, transport_id, - addr.clone(), + send_key, id, self.pool.clone(), self.stats.clone(), @@ -471,7 +476,7 @@ impl TcpTransport { }; let mut pool = self.pool.lock().await; - pool.insert(addr.clone(), conn); + pool.insert(key, conn); self.stats.record_connection_established(); self.stats.record_pool_outbound_added(); @@ -498,7 +503,8 @@ impl TcpTransport { /// discard what it had queued. pub async fn close_connection_async(&self, addr: &TransportAddr) { let mut pool = self.pool.lock().await; - if let Some(conn) = pool.remove(addr) { + let key = key_for_remote(&pool, addr); + if let Some(conn) = key.and_then(|key| pool.remove(&key)) { let TcpConnection { send_tx, send_task, @@ -537,7 +543,7 @@ impl TcpTransport { // Already established? { let pool = self.pool.lock().await; - if pool.contains_key(addr) { + if key_for_remote(&pool, addr).is_some() { return Ok(()); } } @@ -649,7 +655,7 @@ impl TcpTransport { pub fn connection_state_sync(&self, addr: &TransportAddr) -> ConnectionState { // Check established pool first if let Ok(pool) = self.pool.try_lock() { - if pool.contains_key(addr) { + if key_for_remote(&pool, addr).is_some() { return ConnectionState::Connected; } } else { @@ -709,14 +715,16 @@ impl TcpTransport { let packet_tx = self.packet_tx.clone(); let pool = self.pool.clone(); let recv_stats = self.stats.clone(); - let remote_addr = addr.clone(); + let key = PoolKey::outbound(addr.clone()); + let recv_key = key.clone(); + let send_key = key.clone(); let id = next_conn_id(); let recv_task = tokio::spawn(async move { tcp_receive_loop( read_half, transport_id, - remote_addr.clone(), + recv_key, id, packet_tx, pool, @@ -736,7 +744,7 @@ impl TcpTransport { write_half, send_rx, transport_id, - addr.clone(), + send_key, id, self.pool.clone(), self.stats.clone(), @@ -755,7 +763,7 @@ impl TcpTransport { // Use try_lock since we're in a sync context and the pool // should be available (connection_state_sync already checked it) if let Ok(mut pool) = self.pool.try_lock() { - pool.insert(addr.clone(), conn); + pool.insert(key, conn); self.stats.record_connection_established(); self.stats.record_pool_outbound_added(); debug!( @@ -883,6 +891,25 @@ async fn accept_loop( loop { match listener.accept().await { Ok((stream, peer_addr)) => { + // The pool key is the four-tuple, so the local address is + // needed before anything else is done with the socket. A + // socket whose local address cannot be read is already + // broken; drop it rather than pool it under a key that could + // collide with another connection. + let local_addr = match stream.local_addr() { + Ok(a) => a, + Err(e) => { + warn!( + transport_id = %transport_id, + peer_addr = %peer_addr, + error = %e, + "Failed to read local address of accepted socket" + ); + stats.record_connection_rejected(); + continue; + } + }; + // Check inbound connection cap. Counts only inbound (accepted) // connections currently held in the pool; outbound (connect-on-send) // connections live in the same pool but are not subject to the @@ -943,6 +970,7 @@ async fn accept_loop( }; let remote_addr = TransportAddr::from_string(&peer_addr.to_string()); + let key = PoolKey::inbound(remote_addr.clone(), local_addr); // Split and spawn receive task let (read_half, write_half) = stream.into_split(); @@ -950,7 +978,8 @@ async fn accept_loop( let recv_pool = pool.clone(); let recv_packet_tx = packet_tx.clone(); let recv_stats = stats.clone(); - let recv_addr = remote_addr.clone(); + let recv_key = key.clone(); + let send_key = key.clone(); // Readiness barrier: the receive task must not reach its // cleanup path before the pool insert and counter bump below, @@ -963,7 +992,7 @@ async fn accept_loop( tcp_receive_loop( read_half, transport_id, - recv_addr, + recv_key, id, recv_packet_tx, recv_pool, @@ -982,7 +1011,7 @@ async fn accept_loop( write_half, send_rx, transport_id, - remote_addr.clone(), + send_key, id, pool.clone(), stats.clone(), @@ -999,7 +1028,7 @@ async fn accept_loop( }; let mut pool_guard = pool.lock().await; - pool_guard.insert(remote_addr.clone(), conn); + pool_guard.insert(key, conn); drop(pool_guard); stats.record_connection_accepted(); @@ -1012,6 +1041,7 @@ async fn accept_loop( debug!( transport_id = %transport_id, remote_addr = %remote_addr, + local_addr = %local_addr, mtu = conn_mtu, "Accepted inbound TCP connection" ); @@ -1049,7 +1079,7 @@ async fn accept_loop( /// /// The entry is removed only when it carries this connection's `id`. A writer /// can outlive its entry, and by the time its write fails a newer connection -/// may hold the address; that one is left alone. +/// may hold the same four-tuple; that one is left alone. /// /// Frames are written whole. A partial write followed by an error takes the /// connection down with it, so the peer never sees a frame it cannot @@ -1058,11 +1088,12 @@ async fn tcp_send_loop( mut writer: tokio::net::tcp::OwnedWriteHalf, mut frames: mpsc::Receiver>, transport_id: TransportId, - remote_addr: TransportAddr, + key: PoolKey, id: ConnId, pool: ConnectionPool, stats: Arc, ) { + let remote_addr = &key.remote; while let Some(frame) = frames.recv().await { match writer.write_all(&frame).await { Ok(()) => { @@ -1084,7 +1115,7 @@ async fn tcp_send_loop( ); let removed = { let mut pool = pool.lock().await; - remove_own(&mut pool, &remote_addr, id) + remove_own(&mut pool, &key, id) }; // The removed entry's `send_task` is this task, which returns // below, so only the receive task needs stopping. @@ -1134,7 +1165,7 @@ async fn tcp_send_loop( async fn tcp_receive_loop( mut reader: tokio::net::tcp::OwnedReadHalf, transport_id: TransportId, - remote_addr: TransportAddr, + key: PoolKey, id: ConnId, packet_tx: PacketTx, pool: ConnectionPool, @@ -1144,6 +1175,7 @@ async fn tcp_receive_loop( first_frame_timeout: Option, ready_rx: Option>, ) { + let remote_addr = &key.remote; debug!( transport_id = %transport_id, remote_addr = %remote_addr, @@ -1196,7 +1228,7 @@ async fn tcp_receive_loop( "TCP packet received" ); - let packet = ReceivedPacket::new(transport_id, remote_addr.clone(), data); + let packet = ReceivedPacket::new(transport_id, key.remote.clone(), data); if packet_tx.send(packet).await.is_err() { debug!( @@ -1226,7 +1258,7 @@ async fn tcp_receive_loop( // entry actually being removed so a double-cleanup never drives // the counter below zero. let mut pool_guard = pool.lock().await; - let removed = remove_own(&mut pool_guard, &remote_addr, id); + let removed = remove_own(&mut pool_guard, &key, id); drop(pool_guard); if let Some(conn) = removed { conn.send_task.abort(); @@ -1353,12 +1385,22 @@ fn read_mss_mtu(stream: &std::net::TcpStream, default_mtu: u16) -> u16 { #[cfg(test)] mod tests { + use super::pool::PoolMap; use super::*; use crate::transport::framing::build_msg1_frame; use crate::transport::packet_channel; use crate::transport::stream::park_writer; use tokio::time::{Duration, timeout}; + /// The pooled connection for `remote`, whatever key it sits under. + /// + /// The pool is keyed by four-tuple, so a test that knows only the peer + /// address resolves the key the same way the transport does. + fn conn_for<'a>(pool: &'a PoolMap, remote: &TransportAddr) -> Option<&'a TcpConnection> { + let key = key_for_remote(pool, remote)?; + pool.get(&key) + } + /// Poll `f` every 10ms until it holds or `limit` elapses. async fn wait_until bool>(mut f: F, limit: Duration) -> bool { let deadline = Instant::now() + limit; @@ -1630,7 +1672,7 @@ mod tests { // Connection should exist { let pool = t1.pool.lock().await; - assert!(pool.contains_key(&remote)); + assert!(conn_for(&pool, &remote).is_some()); } // Close it @@ -1639,7 +1681,7 @@ mod tests { // Connection should be gone { let pool = t1.pool.lock().await; - assert!(!pool.contains_key(&remote)); + assert!(conn_for(&pool, &remote).is_none()); } t1.stop_async().await.unwrap(); @@ -2275,7 +2317,7 @@ mod tests { let stats = Arc::new(TcpStats::new()); let id = next_conn_id(); pool.lock().await.insert( - remote.clone(), + PoolKey::outbound(remote.clone()), TcpConnection { send_tx: mpsc::channel(1).0, send_task: tokio::spawn(async {}), @@ -2295,7 +2337,7 @@ mod tests { tcp_receive_loop( read_half, TransportId::new(1), - remote.clone(), + PoolKey::outbound(remote.clone()), id, tx, pool.clone(), @@ -2397,13 +2439,7 @@ mod tests { Ok(Err(_)) => false, Err(_) => panic!("send blocked on a peer that stopped reading"), }, - async || { - t.pool - .lock() - .await - .get(remote) - .map(|c| c.send_tx.capacity()) - }, + async || conn_for(&*t.pool.lock().await, remote).map(|c| c.send_tx.capacity()), ) .await } @@ -2448,7 +2484,7 @@ mod tests { let pool: ConnectionPool = Arc::new(Mutex::new(HashMap::new())); let stats = Arc::new(TcpStats::new()); pool.lock().await.insert( - remote.clone(), + PoolKey::outbound(remote.clone()), TcpConnection { send_tx: mpsc::channel(1).0, send_task: tokio::spawn(async {}), @@ -2465,7 +2501,7 @@ mod tests { write_half, send_rx, TransportId::new(1), - remote.clone(), + PoolKey::outbound(remote.clone()), next_conn_id(), pool.clone(), stats.clone(), @@ -2485,7 +2521,7 @@ mod tests { ); assert_eq!( - pool.lock().await.get(&remote).map(|c| c.mtu), + conn_for(&*pool.lock().await, &remote).map(|c| c.mtu), Some(1234), "a failed writer removed the newer connection at its address" ); @@ -2560,7 +2596,7 @@ mod tests { { let pool = t1.pool.lock().await; assert_eq!(pool.len(), 1); - assert_eq!(pool.get(&remote).map(|c| c.mtu), Some(1300)); + assert_eq!(conn_for(&pool, &remote).map(|c| c.mtu), Some(1300)); } drop(sa); @@ -2574,7 +2610,7 @@ mod tests { ); tokio::time::sleep(Duration::from_millis(50)).await; assert_eq!( - t1.pool.lock().await.get(&remote).map(|c| c.mtu), + conn_for(&*t1.pool.lock().await, &remote).map(|c| c.mtu), Some(1300), "the displaced connection's teardown removed its successor" ); @@ -2618,7 +2654,7 @@ mod tests { let pool: ConnectionPool = Arc::new(Mutex::new(HashMap::new())); let stats = Arc::new(TcpStats::new()); pool.lock().await.insert( - remote.clone(), + PoolKey::outbound(remote.clone()), TcpConnection { send_tx: mpsc::channel(1).0, send_task: tokio::spawn(async {}), @@ -2636,7 +2672,7 @@ mod tests { tcp_receive_loop( read_half, TransportId::new(1), - remote.clone(), + PoolKey::outbound(remote.clone()), next_conn_id(), tx, pool.clone(), @@ -2654,7 +2690,7 @@ mod tests { ); assert_eq!( - pool.lock().await.get(&remote).map(|c| c.mtu), + conn_for(&*pool.lock().await, &remote).map(|c| c.mtu), Some(1234), "the teardown removed a newer connection at its address" ); @@ -2705,15 +2741,20 @@ mod tests { } /// Two live inbound connections can share one remote address when they - /// reach a wildcard listener on different local addresses. Closing the - /// older one must not remove the newer one's pool entry. + /// reach a wildcard listener on different local addresses. Each gets its + /// own pool entry, and closing the older one leaves the newer one's entry + /// and its connection alone. /// - /// Known gap, not covered here: two live inbound connections that share a - /// remote address still share one pool key. The second accept replaces - /// the first entry without stopping its tasks and counts a second inbound - /// slot, so the inbound counter ends one above the pool once both - /// connections close. This test checks only that the older connection's - /// teardown no longer removes the newer connection's entry. + /// The pool is keyed by four-tuple, so the two no longer share a key. That + /// closes the gap this test previously recorded: the second accept used to + /// replace the first entry without stopping its tasks while counting a + /// second inbound slot, so the inbound counter ended one above the pool + /// once both connections closed. The counter gates accepts, so repeating + /// that locked the listener out until the daemon restarted. + /// + /// Break-check: with `PoolKey::inbound` ignoring its local address, the + /// pool holds one entry rather than two and the inbound counter does not + /// return to zero. /// /// Linux only: it needs `127.0.0.2` on the loopback interface and Linux /// `SO_REUSEADDR` semantics to bind two client sockets to one port. @@ -2776,10 +2817,35 @@ mod tests { "B must arrive with the same remote address as A" ); assert_eq!(transport.stats().snapshot().connections_accepted, 2); + assert!( + wait_until( + || transport.stats().pool_inbound_count() == 2, + Duration::from_secs(2) + ) + .await, + "both connections should hold an inbound slot" + ); { let pool = transport.pool.lock().await; - assert_eq!(pool.len(), 1); - assert!(pool.contains_key(&remote)); + assert_eq!( + pool.len(), + 2, + "two live connections must not share one pool entry" + ); + let mut locals: Vec<_> = pool + .keys() + .map(|key| key.local.expect("an inbound key carries a local address")) + .collect(); + locals.sort(); + assert_eq!( + locals, + vec![ + SocketAddr::from(([127, 0, 0, 1], port)), + SocketAddr::from(([127, 0, 0, 2], port)), + ], + "the two entries should be the two four-tuples" + ); + assert!(conn_for(&pool, &remote).is_some()); } drop(a); @@ -2793,11 +2859,7 @@ mod tests { ); tokio::time::sleep(Duration::from_millis(50)).await; - let send_tx = transport - .pool - .lock() - .await - .get(&remote) + let send_tx = conn_for(&*transport.pool.lock().await, &remote) .map(|c| c.send_tx.clone()) .expect("closing the older connection removed the newer connection's entry"); send_tx.try_send(frame.clone()).unwrap(); @@ -2917,4 +2979,52 @@ mod tests { t1.stop_async().await.unwrap(); } + + // ======================================================================== + // Inbound pool keying + // ======================================================================== + + /// A reply addressed to an inbound peer goes back over the connection that + /// peer opened, rather than dialing its ephemeral port. + /// + /// An inbound entry is keyed by the four-tuple, but a caller answering a + /// received packet knows only the remote address it came from. The pool has + /// to resolve that address to the entry; if it does not, the send falls + /// through to connect-on-send against the peer's ephemeral port and fails. + #[tokio::test] + async fn a_reply_to_an_inbound_peer_uses_the_connection_it_arrived_on() { + let (tx, mut rx) = packet_channel(100); + let mut transport = TcpTransport::new(TransportId::new(1), None, make_config(), tx); + transport.start_async().await.unwrap(); + let listen = transport.local_addr().unwrap(); + + let mut peer = TcpStream::connect(listen).await.unwrap(); + peer.write_all(&build_msg1_frame()).await.unwrap(); + let packet = timeout(Duration::from_secs(2), rx.recv()) + .await + .expect("timeout waiting for the inbound frame") + .expect("packet channel closed"); + + let mut reply = vec![0xBB; 69]; + reply[0] = 0x02; + reply[1] = 0x00; + reply[2..4].copy_from_slice(&65u16.to_le_bytes()); + transport + .send_async(&packet.remote_addr, &reply) + .await + .expect("a reply to an inbound peer should use its connection"); + + let mut received = vec![0u8; reply.len()]; + timeout( + Duration::from_secs(2), + tokio::io::AsyncReadExt::read_exact(&mut peer, &mut received), + ) + .await + .expect("timeout waiting for the reply") + .expect("reply read failed"); + assert_eq!(received, reply); + + drop(peer); + transport.stop_async().await.unwrap(); + } } diff --git a/src/transport/tcp/pool.rs b/src/transport/tcp/pool.rs index ccacc4d7..26323151 100644 --- a/src/transport/tcp/pool.rs +++ b/src/transport/tcp/pool.rs @@ -4,6 +4,7 @@ //! TCP transport. use std::collections::HashMap; +use std::net::SocketAddr; use std::sync::Arc; use tokio::net::TcpStream; use tokio::sync::{Mutex, mpsc}; @@ -49,8 +50,8 @@ pub(crate) struct TcpConnection { /// MSS-derived MTU for this connection (used for dynamic MTU re-reading). #[allow(dead_code)] pub(crate) mtu: u16, - /// When the connection was established. - #[allow(dead_code)] + /// When the connection was established. Read by `key_for_remote` to pick + /// the newest of several inbound entries sharing a peer address. pub(crate) established_at: Instant, /// Direction of the connection — drives pool-inbound/outbound accounting. pub(crate) direction: Direction, @@ -67,8 +68,67 @@ impl PooledConn for TcpConnection { } } +/// Key identifying one pooled connection. +/// +/// The kernel names a TCP connection by its four-tuple, so a listener on a +/// wildcard address can accept two connections whose peer `ip:port` is the +/// same on two different local addresses. An inbound entry therefore carries +/// the accepted socket's local address as well, and two such connections get +/// two entries rather than displacing each other. +/// +/// Outbound entries carry no local address. Nothing distinguishes two +/// outbound connections to one peer, since the transport keeps at most one, +/// and leaving the local address out keeps the connect-on-send lookup a +/// single hash probe. +#[derive(Clone, Debug, PartialEq, Eq, Hash)] +pub(crate) struct PoolKey { + /// Remote address, as the peer is named by callers and packets. + pub(crate) remote: TransportAddr, + /// Local address of an accepted socket; `None` for outbound. + pub(crate) local: Option, +} + +impl PoolKey { + /// Key for a connection this node opened. + pub(crate) fn outbound(remote: TransportAddr) -> Self { + Self { + remote, + local: None, + } + } + + /// Key for a connection the listener accepted on `local`. + pub(crate) fn inbound(remote: TransportAddr, local: SocketAddr) -> Self { + Self { + remote, + local: Some(local), + } + } +} + +/// The pooled connections, keyed by [`PoolKey`]. +pub(crate) type PoolMap = HashMap; + /// Shared connection pool. -pub(crate) type ConnectionPool = Arc>>; +pub(crate) type ConnectionPool = Arc>; + +/// Resolve a bare remote address to the key of the connection to use for it. +/// +/// Callers that send, close or query by peer address know only the remote, so +/// the four-tuple has to be recovered. An outbound entry is tried first, so the +/// common case is one hash probe. Inbound entries also carry a local address, +/// so they are found by scanning for the remote and taking the most recently +/// established, which is the connection a peer that reconnected is using. +pub(crate) fn key_for_remote(pool: &PoolMap, remote: &TransportAddr) -> Option { + let outbound = PoolKey::outbound(remote.clone()); + if pool.contains_key(&outbound) { + return Some(outbound); + } + pool.iter() + .filter(|(key, _)| &key.remote == remote) + .max_by_key(|(_, conn)| conn.established_at) + .map(|(key, _)| key.clone()) +} /// A pending background connection attempt. ///