From 8c78e1f4b993d4df395e9b7c856e95f965ad4bb3 Mon Sep 17 00:00:00 2001 From: Arjen <18398758+Origami74@users.noreply.github.com> Date: Sun, 13 Sep 2026 09:22:46 -0300 Subject: [PATCH] feat(peer): a peer with a session is never dialled; a handshake creates no path state MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Three ways an address for a peer we already hold a session with used to reach the dialler — a beacon on a new transport, `update_peers` or `fipsctl connect`, a configured address whose transport came up later — and each was a second handshake, which the far side read as a rekey and this side resolved as a cross-connection, the two not composing. Two phones hearing the same Wi-Fi return at the same moment both dialled at once; one side swapped to its outbound session and freed the index it had just handed out in the rekey reply, the other kept its inbound session and the pre-rekey index, every frame between them was dropped, and the link-dead reap tore the peer down. About a minute dark on every Wi-Fi return. A peer that holds a session is never dialled now. An address on a transport it has no path over becomes a path candidate under that session; one on a transport whose path is not eligible re-points that path (the active path included: it is not answering, that is why we are here); one on a transport whose path is carrying acknowledged traffic changes nothing. The heartbeat tick probes the candidate under the existing session — one authenticated, replay-checked round trip — and the mandatory switch takes it if the current path stops answering. Nothing is lost against the dial: a session that is truly gone answers no probe either, is reaped by the link-dead timeout, and is dialled then; a peer that restarted dials us with a new epoch and wins promotion outright, as before. Applies to the control API's connect, to update_peers, to configured addresses (checked once a tick) and to transport discovery alike. The counterpart: a handshake creates no path state. A dial that does reach a peer with a session — a startup that lists two addresses dials both, a caller that still dials by hand, an older node dialling us — is classified and resolved exactly as before this work: rekey, duplicate or restart on the responder, whichever transport the msg1 arrived on; the cross-connection tie-break on the initiator. The address it ran to is left as a candidate for the probe exchange. Two reasons. Both ends must resolve a handshake on the same information, and "is this a new transport to a live peer" was a fact only one end could see. And the IK responder commits at msg1, which carries no freshness beyond the startup epoch: a captured msg1 replayed from any address would otherwise have planted a path, probed full-size for the life of the peering and counting as a transport the peer is on for the decrypt-failure gate. So that gate now counts garbage only on the active path or one the peer has acknowledged. On a connection-oriented transport the connection a dial opened is kept as the candidate's socket rather than closed as the losing leg: the probe rides it, and closing it would only have the first probe dial again — or, at the responder, find an ephemeral port that cannot be dialled at all. `api_disconnect` closes every path's connection, the standby's included; loopback records the closes it is asked for so a test can say so. Path heartbeats are gated and bounded. A peer with one live path is not path-heartbeated: selection has nothing to move to, the link heartbeat keeps liveness, and five probes a second on every single-path link was cost without a decision behind it. A standby the peer never acknowledges is given up after eight discovery probes, Dead and pruned after the grace; the active path is never given up. The active path's first probe is small, the handshake having proved it and seeded its MTU. And a Dead path is probed again when its transport returns: nothing on our side ever re-probed one, so after a NIC replug traffic stayed on the standby until the grace pruned the path and a beacon found it with no history. The presence edge now revives every Dead path on the transport as Probing, RTT window and ETX kept. Smaller: `add_path_candidate` re-points a known transport's path at a moved address (`refresh_path_addr`), for a Wi-Fi Aware data path that re-forms with a new link-local; `api_disconnect` closes every path's connection, not the active one alone; `path_show` is built from the `show_peers` path projection plus the three now-relative fields; `PathState` and `TransportRole` render through `as_str()`; `node.path.switch_margin` is validated finite and at least 1.0; `PathPolicy::PERMISSIVE` had no users; four doc comments an inserted function had split are put back on the function they describe. The dual-udp-flap scenario is config-driven: the dial owner lists udp/main and udp/, both dial at startup, and the second is proven as a path under the first's session by the probe exchange. --- CHANGELOG.md | 17 + docs/reference/control-socket.md | 30 +- src/config/mod.rs | 26 ++ src/config/transport.rs | 12 + src/control/queries.rs | 15 +- src/node/dataplane/encrypted.rs | 17 +- src/node/dataplane/peer_actions.rs | 5 +- src/node/handlers/handshake.rs | 19 +- src/node/handlers/path.rs | 127 +++---- src/node/lifecycle/mod.rs | 405 ++++++++++++-------- src/node/mod.rs | 4 +- src/node/tests/control.rs | 31 +- src/node/tests/multi_path.rs | 592 +++++++++++++++++++++++++---- src/node/tests/unit.rs | 16 +- src/peer/active.rs | 152 ++++++-- src/peer/mod.rs | 5 +- src/proto/fmp/core.rs | 7 + src/proto/fmp/tests/core.rs | 26 ++ src/proto/mmp/metrics.rs | 2 +- src/transport/loopback.rs | 16 + src/transport/mod.rs | 30 +- testing/chaos/README.md | 9 +- testing/chaos/sim/config_gen.py | 25 +- testing/chaos/sim/runner.py | 65 +--- 24 files changed, 1156 insertions(+), 497 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0bd4d5e6..c52e3431 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -44,6 +44,23 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 a path tells the peer with a `PathClose` on a surviving path, so the peer moves at once rather than after its own timeout. +- Connect semantics, for multi-path. A peer that holds a session is never + dialled again: an address for it on a transport it has no path over — + from a beacon, `update_peers`, `fipsctl connect`, a runtime peer lane, + or a configured address whose transport came up later — becomes a + candidate path, probed under the session by the next heartbeat tick; + one on a transport whose path has stopped answering re-points that + path; one on a transport whose path carries acknowledged traffic changes + nothing. A handshake never creates path state: a second handshake to a + peer with a session (two configured addresses dialled at startup) + resolves as it always did — rekey, duplicate or cross-connection + tie-break, whichever transport it ran over — and the address it ran to + is left as a candidate for the probe exchange. A standby the peer never + answers on is given up after eight probes, and until it is answered it + does not count as a transport the peer is on for the decrypt-failure + gate. On a connection-oriented transport the connection a dial opened + is kept as the candidate's socket rather than closed. + - Operator surface, for multi-path: `role: backup` on any transport (never carries a peer's traffic while a normal path is eligible); `fipsctl path show|pin|unpin` and the `path_show`, `path_pin`, diff --git a/docs/reference/control-socket.md b/docs/reference/control-socket.md index 0510b385..145d1fa5 100644 --- a/docs/reference/control-socket.md +++ b/docs/reference/control-socket.md @@ -116,7 +116,8 @@ table below lists every command currently registered. | ------- | ------ | ----------------------------- | | `show_status` | — | `version`, `npub`, `node_addr`, `ipv6_addr`, `state`, `is_leaf_only`, `is_root` (bool — this node is the spanning-tree root), `root` (hex node-addr of the current tree root), `persistent` (bool — identity is persisted, i.e. `persistent` set or an `nsec` configured), `peer_count`, `session_count`, `link_count`, `transport_count`, `connection_count`, `transport_peer_counts` (object mapping transport-type name to its connected-peer count; configured transports appear with `0`), `tun_state`, `tun_name`, `effective_ipv6_mtu`, `control_socket`, `pid`, `exe_path`, `uptime_secs`, `estimated_mesh_size`, `forwarding`, `sparklines`. | | `show_acl` | — | `allow_file`, `deny_file`, `enforcement_active`, `effective_mode`, `default_decision`, `allow_all`, `deny_all`, `allow_file_entries`, `deny_file_entries`, `allow_entries`, `deny_entries`. | -| `show_peers` | — | `peers[]` — per-peer object: `node_addr`, `npub`, `display_name`, `ipv6_addr`, `connectivity`, `link_id`, `direction`, `transport_addr`, `transport_type`, `is_parent`, `is_child`, `tree_depth`, `effective_depth` (`tree_depth + link_cost` — the metric `evaluate_parent` ranks on; `null` when the peer has no coords, or is unmeasured while another peer has an SRTT sample, per the cold-start gate), `stats`, `noise`, `current_k_bit`, `mmp`, `paths[]` (every path to the peer: `transport_id`, `transport`, `transport_type`, `addr`, `state`, `active`, `remote_active`, `role`, `pinned`, `last_rtt_ms`, `min_rtt_ms`, `rtt_samples`, `etx`, `score` — the `path_show` fields minus the now-relative ages), plus optional `nostr_traversal`, `rekey_in_progress`, `rekey_draining`. | +| `show_peers` | — | `peers[]` — per-peer object: `node_addr`, `npub`, `display_name`, `ipv6_addr`, `connectivity`, `link_id`, `direction`, `transport_addr`, `transport_type`, `is_parent`, `is_child`, `tree_depth`, `effective_depth` (`tree_depth + link_cost` — the metric `evaluate_parent` ranks on; `null` when the peer has no coords, or is unmeasured while another peer has an SRTT sample, per the cold-start gate), `stats`, `noise`, `current_k_bit`, `mmp`, `paths[]` (every path to the peer: `transport_id`, `transport` (instance name or null), `transport_type`, `addr`, `state` (`probing` / `live` / `suspect` / `dead`), `active`, `remote_active`, `role` (`normal` / `backup`), `pinned`, `last_rtt_ms`, `min_rtt_ms`, `rtt_samples`, `etx`, `score`), plus optional `nostr_traversal`, `rekey_in_progress`, `rekey_draining`. | +| `path_show` | `npub` (bech32) | Every path to one peer. `data`: `peer`, `link_cost`, `link_cost_held`, and `paths[]` — the `show_peers` per-path object plus `rx_live_ms_ago`, `tx_live_ms_ago` (ms since the last authentic frame heard there / the last ack proving the peer hears us there; `null` if never) and `acked_once`. Takes a parameter, so it is served on the daemon's main task like the mutating commands, not from the snapshot. | | `show_links` | — | `links[]` — `link_id`, `transport_id`, `remote_addr`, `direction`, `state`, `created_at_ms`, `stats`. | | `show_tree` | — | `my_node_addr`, `root`, `root_npub` (bech32 npub of the current tree root), `is_root`, `depth`, `my_coords[]`, `parent`, `parent_display_name`, `declaration_sequence`, `declaration_signed`, `peer_tree_count`, `peers[]`, `stats`. | | `show_sessions` | — | `sessions[]` — `remote_addr`, `npub`, `display_name`, `state` (`established`, `initiating`, `awaiting_msg3`, `unknown`), `is_initiator`, `last_activity_ms`, `stats`, optional `mmp`, `current_k_bit`, `is_draining`. | @@ -173,17 +174,24 @@ not reproduced here to avoid duplicating the source. | `probe_start` | `npub` (bech32) | Admits a diagnostic probe job and returns immediately. `data`: `probe_id`, `npub`, `node_addr`, `display_name`, `budget_ms`. | | `probe_poll` | `probe_id` (integer) | Reports a probe's progress. `data`: `state` (`running` / `done`) and `report`. A terminal job is removed on the poll that observes it, so the report is delivered once. | | `probe_cancel` | `probe_id` (integer) | Runs the probe's terminal actions immediately, without the teardown grace tick. | -| `path_show` | `npub` (bech32) | Every path to the peer. `data`: `peer`, `link_cost`, `link_cost_held`, and `paths[]` with `transport_id`, `transport` (instance name or null), `addr`, `state`, `active`, `remote_active`, `role`, `pinned`, `rx_live_ms_ago`, `tx_live_ms_ago`, `acked_once`, `last_rtt_ms`, `min_rtt_ms`, `rtt_samples`, `etx`, `score`. | -| `path_pin` | `npub` (bech32), `transport` (instance name or numeric id) | Pins this node's traffic to the peer to that transport's path. Applies on the next selection tick. Error if the peer has no path there. | -| `path_unpin` | `npub` (bech32) | Clears the pin. | +| `path_pin` | `npub` (bech32), `transport` (instance name or numeric id) | Pins this node's traffic to the peer to that transport's path. Applies on the next selection tick, and is suspended while that path is not eligible and re-applied when it is again. `data`: `{"pinned": }`. Error if the peer has no path there. | +| `path_unpin` | `npub` (bech32) | Clears the pin. `data`: `{"pinned": null}`. | -`connect` on a peer the node is **already connected to** neither tears the -live link down nor ignores the address: the address is tried as an alternate -path alongside the existing one, and the peer moves to it only if that -handshake authenticates. The response carries `refreshed` — `true` when such a -handshake was started, `false` when the peer is already on this exact path and -that path is fresh (a successful no-op). A `connect` that starts an ordinary -dial to a peer the node does not yet hold also reports `refreshed: false`. +`connect` has three outcomes, told apart by whether the node already holds +a session with the peer and by the response's `refreshed` field: + +- **No session:** an ordinary dial over the named transport. `refreshed: + false`; the peer appears in `show_peers` once the handshake completes. +- **Session, and the address is on a transport the peer has no path over, + or one whose path has stopped answering:** no handshake. The address + becomes a path candidate under the existing session (or re-points the + unanswering path), the next heartbeat tick probes it, and selection + moves traffic to it if it measures better or the current path stops + answering. `refreshed: true`. `path_show` lists it as `probing` until + the peer acknowledges, `live` after. +- **Session, and the peer is already reachable at exactly that address, + or that transport's path is carrying acknowledged traffic:** nothing + changes. `refreshed: false`. `connect` is ephemeral either way: the peer is not written to the config file and gets no auto-reconnect, so an attempt that fails leaves no residue. diff --git a/src/config/mod.rs b/src/config/mod.rs index 582dd628..41cbca47 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -1241,6 +1241,18 @@ impl Config { } } + // Path selection. The margin is the whole fail-back policy: a + // standby must beat the active path's score by this factor before + // traffic moves. Below 1.0 (or NaN, which compares false both ways) + // two paths of equal score would swap after every dwell, each swap + // re-seeding the path MTU and holding the tree-visible link cost. + let margin = self.node.path.switch_margin; + if !margin.is_finite() || margin < 1.0 { + return Err(ConfigError::Validation(format!( + "`node.path.switch_margin` = {margin} must be a finite number of at least 1.0: it is the factor a standby's score must beat the active path's by, and anything less makes two equal paths swap after every dwell" + ))); + } + let native = &self.node.native_api; // Both floors refuse a node that would start, answer every setup call // and then drop every datagram a peer sent. A zero `backlog` makes the @@ -2557,6 +2569,20 @@ node: assert!(config.node.discovery.is_none()); } + #[test] + fn a_switch_margin_below_one_is_refused() { + // Two paths of equal score would swap after every dwell. + for bad in [0.9, 0.0, -1.0, f64::NAN, f64::INFINITY] { + let mut config = Config::default(); + config.node.path.switch_margin = bad; + let err = config.validate().expect_err("validation should fail"); + assert!(err.to_string().contains("switch_margin"), "{bad}: {err}"); + } + let mut config = Config::default(); + config.node.path.switch_margin = 1.0; + config.validate().expect("1.0 means any better path wins"); + } + #[test] fn test_a_zero_netmon_poll_interval_is_refused() { // It was silently clamped to 1s, so a typo produced a node that polled diff --git a/src/config/transport.rs b/src/config/transport.rs index 311b597b..c2e47bc5 100644 --- a/src/config/transport.rs +++ b/src/config/transport.rs @@ -58,6 +58,18 @@ pub enum TransportRole { Backup, } +impl TransportRole { + /// The config and control-socket spelling: `normal`, `backup`. Same + /// strings serde reads and writes, fixed here so a variant rename cannot + /// silently change what `show_peers` and `path_show` emit. + pub fn as_str(self) -> &'static str { + match self { + Self::Normal => "normal", + Self::Backup => "backup", + } + } +} + /// UDP transport instance configuration. #[derive(Debug, Clone, Default, Serialize, Deserialize)] #[serde(deny_unknown_fields)] diff --git a/src/control/queries.rs b/src/control/queries.rs index 1af67fe7..f9c511d5 100644 --- a/src/control/queries.rs +++ b/src/control/queries.rs @@ -408,13 +408,9 @@ pub fn show_peers(node: &Node) -> Value { json!({ "peers": peers }) } -/// Render a snapshot [`EntityMmp`](super::snapshot::EntityMmp) into the inline -/// MMP JSON block, with the quality-index key named `quality_key` (`lqi` for -/// peers, `sqi` for sessions). Reproduces the on-loop key insertion order -/// exactly. `path_mtu` is emitted (inside the leading literal) only when -/// present (session-layer); for peers it is `None` and omitted. -/// Render a peer's path rows as the `paths` array of `show_peers`. -fn render_peer_paths(paths: &[super::snapshot::PeerPathRow]) -> Value { +/// Render a peer's path rows as the `paths` array of `show_peers`. Also the +/// base of `path_show`'s rows, which add the now-relative fields on top. +pub(crate) fn render_peer_paths(paths: &[super::snapshot::PeerPathRow]) -> Value { Value::Array( paths .iter() @@ -440,6 +436,11 @@ fn render_peer_paths(paths: &[super::snapshot::PeerPathRow]) -> Value { ) } +/// Render a snapshot [`EntityMmp`](super::snapshot::EntityMmp) into the inline +/// MMP JSON block, with the quality-index key named `quality_key` (`lqi` for +/// peers, `sqi` for sessions). Reproduces the on-loop key insertion order +/// exactly. `path_mtu` is emitted (inside the leading literal) only when +/// present (session-layer); for peers it is `None` and omitted. fn render_entity_mmp(mmp: &super::snapshot::EntityMmp, quality_key: &str) -> Value { // The on-loop `show_sessions` block places loss_rate/etx/goodput_bps/ // delivery ratios/path_mtu in the leading json! literal, while `show_peers` diff --git a/src/node/dataplane/encrypted.rs b/src/node/dataplane/encrypted.rs index f425bcce..c000f7d6 100644 --- a/src/node/dataplane/encrypted.rs +++ b/src/node/dataplane/encrypted.rs @@ -551,6 +551,12 @@ impl Node { /// landing on an index the allocator has already handed to a new owner: /// index reuse is immediate, with no quarantine. /// + /// Only a path the peer has proven counts: the active one, or one the + /// peer has acknowledged a probe on. A `Probing` path is an address we + /// were *told* about — a beacon, a config entry, a handshake source — + /// and until the peer answers there, garbage arriving on its transport + /// says nothing about the peer. + /// /// A peer with no transport bound yet is charged unconditionally, as /// before. pub(in crate::node) fn charge_decrypt_failure( @@ -558,10 +564,13 @@ impl Node { node_addr: &crate::NodeAddr, transport_id: crate::transport::TransportId, ) { - let on_path = self - .peers - .get(node_addr) - .is_some_and(|peer| peer.paths().is_empty() || peer.path_on(transport_id).is_some()); + let on_path = self.peers.get(node_addr).is_some_and(|peer| { + peer.paths().is_empty() + || peer.transport_id() == Some(transport_id) + || peer + .path_on(transport_id) + .is_some_and(|path| path.acked_once()) + }); if !on_path { trace!( peer = %self.peer_display_name(node_addr), diff --git a/src/node/dataplane/peer_actions.rs b/src/node/dataplane/peer_actions.rs index 03b56d55..285161cc 100644 --- a/src/node/dataplane/peer_actions.rs +++ b/src/node/dataplane/peer_actions.rs @@ -345,8 +345,9 @@ impl Node { "executor CrossConnectionLost is unreachable on \ driven net-new establish paths" ); - // Close this (losing) connection, drop its link, - // and restore `addr_to_link` to the winner. + // Close this connection, drop its link, and + // point `addr_to_link` for the new address at the + // winner, so a msg1 from it is recognised. if let Some(transport) = self.transports.get(&ambient.transport_id) { diff --git a/src/node/handlers/handshake.rs b/src/node/handlers/handshake.rs index 072a42b0..6fa37c7c 100644 --- a/src/node/handlers/handshake.rs +++ b/src/node/handlers/handshake.rs @@ -1573,13 +1573,24 @@ impl Node { // Clean up outbound connection state self.pending_outbound.remove(&key); - // Close the losing TCP connection (no-op for connectionless) + + // The handshake ran over some (transport, address). Whichever + // session won, that is where the peer answered just now — but a + // handshake creates no path state: the probe exchange, which is + // authenticated and replay-checked under the surviving session, + // is the one way a path is proven. Leave the address as a + // candidate for the heartbeat tick. The link record goes (the + // peer's link is the one its session rides), but the transport + // connection stays: on TCP, Tor or Nym that socket is what the + // probe will go out on, and closing it would only have the + // first probe dial it again — or, on the responder, find that + // our ephemeral port cannot be dialled at all. `api_disconnect` + // closes every path's connection. A connectionless close was a + // no-op either way. if let Some(link) = self.links.get(&link_id) { let tid = link.transport_id(); let addr = link.remote_addr().clone(); - if let Some(transport) = self.transports.get(&tid) { - transport.close_connection(&addr).await; - } + self.add_path_candidate(peer_node_addr, tid, addr); } self.remove_link(&link_id); diff --git a/src/node/handlers/path.rs b/src/node/handlers/path.rs index ea3fe45a..819cfa5d 100644 --- a/src/node/handlers/path.rs +++ b/src/node/handlers/path.rs @@ -121,39 +121,24 @@ impl Node { } /// `fipsctl path show `: every path to the peer, per direction. + /// + /// The per-path object is the `show_peers` one (`project_peer_paths`, + /// rendered by `render_peer_paths`) plus the three fields only a + /// now-relative read can give: the liveness ages and `acked_once`. + /// Built from the same projection so the two field lists cannot drift. pub(crate) fn api_path_show(&self, npub: &str) -> Result { let node_addr = self.resolve_peer_npub(npub)?; let peer = &self.peers[&node_addr]; let now_ms = crate::time::mono_ms(); - let active = peer.transport_id(); - let paths: Vec = peer - .paths() - .iter() - .map(|path| { - let ago = |at: Option| at.map(|t| now_ms.saturating_sub(t)); - serde_json::json!({ - "transport_id": path.transport_id().as_u32(), - "transport": self - .transports - .get(&path.transport_id()) - .and_then(|t| t.name().map(str::to_string)), - "addr": path.addr().to_string(), - "state": format!("{:?}", path.state()).to_lowercase(), - "active": Some(path.transport_id()) == active, - "remote_active": path.remote_active(), - "role": format!("{:?}", path.role()).to_lowercase(), - "pinned": path.pinned(), - "rx_live_ms_ago": ago(path.rx_live_at_ms()), - "tx_live_ms_ago": ago(path.tx_live_at_ms()), - "acked_once": path.acked_once(), - "last_rtt_ms": path.last_rtt_ms(), - "min_rtt_ms": path.min_rtt_ms(), - "rtt_samples": path.rtt_samples(), - "etx": path.etx(), - "score": path.score(), - }) - }) - .collect(); + let ago = |at: Option| at.map(|t| now_ms.saturating_sub(t)); + let mut paths = crate::control::queries::render_peer_paths(&self.project_peer_paths(peer)); + if let Some(rows) = paths.as_array_mut() { + for (row, path) in rows.iter_mut().zip(peer.paths()) { + row["rx_live_ms_ago"] = serde_json::json!(ago(path.rx_live_at_ms())); + row["tx_live_ms_ago"] = serde_json::json!(ago(path.tx_live_at_ms())); + row["acked_once"] = serde_json::json!(path.acked_once()); + } + } Ok(serde_json::json!({ "peer": npub, "link_cost": peer.link_cost(now_ms), @@ -204,13 +189,28 @@ impl Node { let Some(peer) = self.peers.get_mut(&node_addr) else { return; }; - if peer.transport_id() == Some(transport_id) { - // The active path: the handshake proved it. + if peer.transport_id() == Some(transport_id) + && peer.path_on(transport_id).is_some_and(|p| p.is_eligible()) + { + // The active path, and it is answering: nothing to add, and an + // address that has proven nothing does not displace it. An + // active path that has stopped answering falls through to the + // re-pointing below like any other. return; } let was_new = peer.path_on(transport_id).is_none(); peer.add_path(transport_id, remote_addr.clone()) .set_role(role); + // A known transport at a new address: the peer moved there (an + // Aware data path that re-formed, a DHCP lease that changed) and + // the old address answers nothing. Re-point the path; the heartbeat + // tick probes it from here. Only while the path is not eligible: an + // address that is carrying acknowledged traffic is not displaced by + // one that has proven nothing — a probe from the new address + // (`note_path_probe`) is what moves a working path. + let moved = !was_new + && peer.path_on(transport_id).is_some_and(|p| !p.is_eligible()) + && peer.refresh_path_addr(transport_id, remote_addr.clone()); if was_new { debug!( peer = %self.peer_display_name(&node_addr), @@ -218,47 +218,13 @@ impl Node { remote_addr = %remote_addr, "Peer beaconed on a new transport; path added, probing" ); - } - } - - /// Tests: add the path and probe it at once, as one heartbeat tick - /// would, without the tick's other sends. Uses the test-only - /// `take_probe`, which honours the per-path backoff. - #[cfg(test)] - pub(in crate::node) async fn maybe_probe_path( - &mut self, - node_addr: NodeAddr, - transport_id: TransportId, - remote_addr: TransportAddr, - ) { - self.add_path_candidate(node_addr, transport_id, remote_addr.clone()); - let now_ms = crate::time::mono_ms(); - let timing = self.heartbeat_timing(); - let Some(peer) = self.peers.get_mut(&node_addr) else { - return; - }; - if peer.transport_id() == Some(transport_id) { - return; - } - let Some((probe_id, remote_active, path_id)) = peer.take_probe( - transport_id, - now_ms, - timing.fast_ms, - timing.discovery_cap_ms, - ) else { - return; - }; - let probe = PathMessage { - probe_id, - remote_active, - path_id, - }; - let wire = self.pad_to_link_mtu(probe.encode_probe().to_vec(), transport_id, &remote_addr); - if let Err(e) = self - .send_encrypted_link_message_on_path(&node_addr, &wire, transport_id, remote_addr) - .await - { - debug!(peer = %self.peer_display_name(&node_addr), error = %e, "Path probe send failed"); + } else if moved { + debug!( + peer = %self.peer_display_name(&node_addr), + transport_id = %transport_id, + remote_addr = %remote_addr, + "Peer beaconed at a new address on a known transport; path re-addressed, probing" + ); } } @@ -766,10 +732,21 @@ impl Node { } /// A transport's presence came back: clear the probe backoff on every - /// path over it so the next discovery tick may probe at once. + /// path over it so the next heartbeat tick may probe at once, and + /// revive every path that went `Dead` with it — a replugged NIC, a + /// wifi interface that cycled — as `Probing`, history kept. Without + /// this nothing on our side ever probed a `Dead` path again: traffic + /// stayed on the standby until the grace pruned the path and a fresh + /// beacon found it with no history. pub(in crate::node) fn reset_probe_backoff_on_transport(&mut self, transport_id: TransportId) { - for peer in self.peers.values_mut() { - peer.reset_probe_backoff_on(transport_id); + for (node_addr, peer) in self.peers.iter_mut() { + if peer.reset_probe_backoff_on(transport_id) { + debug!( + peer = %node_addr, + %transport_id, + "Transport returned: dead path probing again" + ); + } } } } diff --git a/src/node/lifecycle/mod.rs b/src/node/lifecycle/mod.rs index deeb03d8..acb1372d 100644 --- a/src/node/lifecycle/mod.rs +++ b/src/node/lifecycle/mod.rs @@ -464,6 +464,21 @@ impl Node { .map(|t| t.transport_type().connection_oriented) .unwrap_or(false); + // A dial to a peer we already hold a session with (a startup that + // lists two addresses, a caller that still dials by hand) proves + // nothing by itself: the handshake creates no path state, and its + // outcome is the cross-connection tie-break as ever. The address + // is still a fact worth holding — leave it as a candidate for the + // heartbeat tick to probe under the session, on a datagram + // transport, so it becomes a path whichever way the dial goes. + if !is_connection_oriented && self.peers.contains_key(peer_identity.node_addr()) { + self.add_path_candidate( + *peer_identity.node_addr(), + transport_id, + remote_addr.clone(), + ); + } + // Allocate link ID and create link let link_id = self.allocate_link_id(); @@ -812,51 +827,38 @@ impl Node { let connected = self.peers.contains_key(&node_addr); if connected { - // Active peer: skip every candidate while the link we - // already hold is live — the current path *and* any - // alternate one. + // Active peer: never a dial, whatever the state of the + // link we hold. The address is a path to add or + // re-point, and the heartbeat tick probes it under the + // session we have — one round trip, authenticated, and + // selection moves traffic if the path proves better or + // the current one is not answering. // - // Only the same-path case used to be skipped, which left - // the stated intent ("avoid churning a healthy link") - // covering exactly the case that could not churn anything. - // A peer reachable twice — the ordinary result of two - // machines sharing a LAN and a cable, since each beacons on - // both — was therefore re-dialled on its alternate path - // every discovery tick, forever. Each dial that completed - // promoted and displaced the incumbent, so the peer's link - // migrated back and forth on a fixed cadence, tearing down - // and re-establishing its session each time. Measured on - // real hardware: seventeen dials to one peer in fifteen - // minutes, alternating wifi and cable, displacing a link - // reporting `etx = 1.0` and `loss = 0.0`. + // Dialling an active peer used to be the fallback for one + // that had gone quiet ("liveness is the gate"). Two + // reasons it is not any more. A handshake to a peer that + // already holds a session is read by the far side as a + // rekey, and when both ends do it at once — the ordinary + // case, since both hear the same medium come back — the + // rekey and the cross-connection resolution overlap and + // the two sides part on different session indices, dead + // to each other until the link timeout reaps them. And + // the dial buys nothing the probe does not: a session + // that is truly gone answers no probe either, is reaped + // by the link-dead timeout, and is dialled then; a peer + // that restarted dials us itself with a new epoch and + // wins the promotion outright. // - // When that peer is the parent — which the best path - // usually is — every migration also switched parents, - // invalidating the downstream coordinate cache and - // re-announcing to every peer. The cost of the churn was - // therefore mesh-wide while the benefit was nil: the link - // being replaced was already perfect. - // - // Failover is unaffected. Liveness is the gate, so a peer - // that stops answering goes stale within a heartbeat - // interval and every path, alternate included, is dialled - // again. What is given up is switching away from a link - // that is working, which is not a thing worth doing. - if self.active_peer_link_is_live(&node_addr) { - // A live peer beaconing on a transport we hold no - // path to it over is a path to add, not a link to - // replace: the heartbeat tick probes it under the - // existing session instead of dialling. - path_candidates.push((node_addr, candidate_transport_id, remote_addr)); - continue; - } - if self.is_connecting_to_peer_on_path( - &node_addr, - candidate_transport_id, - &remote_addr, - ) { - continue; - } + // (History: only the same-path case used to be skipped, + // so a peer reachable twice — two machines sharing a LAN + // and a cable — was re-dialled on its alternate path + // every discovery tick, each completed dial displacing + // the incumbent: seventeen dials to one peer in fifteen + // minutes, alternating wifi and cable, over a link + // reporting `etx = 1.0`. Liveness gating fixed that and + // left the quiet-peer dial; this removes the last of it.) + path_candidates.push((node_addr, candidate_transport_id, remote_addr)); + continue; } else if self.is_connecting_to_peer_on_path( &node_addr, candidate_transport_id, @@ -878,6 +880,7 @@ impl Node { for (node_addr, transport_id, remote_addr) in path_candidates { self.add_path_candidate(node_addr, transport_id, remote_addr); } + self.add_configured_path_candidates(); if transport_neighbors.is_empty() { return; @@ -2560,6 +2563,126 @@ impl Node { .collect() } + /// Configured addresses of live peers on transports they have no path + /// over become paths. Runs every discovery tick, idempotent and cheap: + /// a configured address whose transport was down at dial time (wifi + /// joined later, Tor came up) is otherwise never looked at again, since + /// a peer that is already active is not re-dialled. + fn add_configured_path_candidates(&mut self) { + let configs: Vec = self.config().auto_connect_peers().cloned().collect(); + for peer_config in configs { + let Ok(identity) = PeerIdentity::from_npub(&peer_config.npub) else { + continue; + }; + let node_addr = *identity.node_addr(); + if !self.peers.contains_key(&node_addr) || !self.active_peer_link_is_live(&node_addr) { + continue; + } + for addr in peer_config.addresses_by_priority() { + if addr.transport == "udp" && addr.addr.eq_ignore_ascii_case("nat") { + continue; + } + let Some((transport_id, remote_addr)) = self.resolve_peer_address(addr) else { + continue; + }; + let has_path = self + .peers + .get(&node_addr) + .is_some_and(|p| p.path_on(transport_id).is_some()); + if !has_path { + self.add_path_candidate(node_addr, transport_id, remote_addr); + } + } + } + } + + /// The transport and address a configured peer address dials to, or + /// `None` (logged at debug) if no operational transport can carry it. + /// + /// The transport field may name a specific instance (`"udp/aware"`): + /// the type half picks the resolver, the instance half is handed to + /// whichever resolver can honour it, and only the UDP one can. The + /// `"nat"` pseudo-address is not resolved here. + fn resolve_peer_address(&self, addr: &PeerAddress) -> Option<(TransportId, TransportAddr)> { + let spec = addr.spec(); + if addr.transport == "ethernet" { + return match self.resolve_ethernet_addr(&addr.addr) { + Ok(result) => Some(result), + Err(e) => { + debug!( + transport = %addr.transport, + addr = %addr.addr, + error = %e, + "Failed to resolve Ethernet address" + ); + None + } + }; + } + if addr.transport == "ble" { + #[cfg(ble_available)] + { + return match self.resolve_ble_addr(&addr.addr) { + Ok(result) => Some(result), + Err(e) => { + debug!( + transport = %addr.transport, + addr = %addr.addr, + error = %e, + "Failed to resolve BLE address" + ); + None + } + }; + } + #[cfg(not(ble_available))] + { + debug!(transport = %addr.transport, "BLE transport not available on this build"); + return None; + } + } + let tid = if spec.kind == "udp" + && let Ok(remote_socket_addr) = addr.addr.parse::() + { + match self.find_udp_transport_for_remote_addr(remote_socket_addr, spec.instance) { + Some((id, _)) => id, + None => { + debug!( + transport = %addr.transport, + addr = %addr.addr, + "No compatible operational UDP transport for address" + ); + return None; + } + } + } else if spec.instance.is_some() { + // Only the UDP resolver above can honour an instance name. + // Matching any instance of the type here would be the silent + // wrong-lane substitution this whole mechanism exists to + // prevent, so refuse instead. + debug!( + transport = %addr.transport, + addr = %addr.addr, + "Instance-qualified address for a transport type that \ + does not support instance selection" + ); + return None; + } else { + match self.find_transport_for_type(spec.kind) { + Some(id) => id, + None => { + debug!( + transport = %addr.transport, + addr = %addr.addr, + "No operational transport for address type" + ); + return None; + } + } + }; + Some((tid, TransportAddr::from_string(&addr.addr))) + } + async fn attempt_peer_address_list( &mut self, peer_config: &PeerConfig, @@ -2581,9 +2704,6 @@ impl Node { if attempted >= max_attempts { break; } - // The transport field may name a specific instance - // (`"udp/aware"`); everything below dispatches on the type half - // and hands the instance half to whichever resolver can honour it. let spec = addr.spec(); if spec.kind == "udp" && addr.addr.eq_ignore_ascii_case("nat") { @@ -2601,82 +2721,8 @@ impl Node { continue; } - let (transport_id, remote_addr) = if addr.transport == "ethernet" { - match self.resolve_ethernet_addr(&addr.addr) { - Ok(result) => result, - Err(e) => { - debug!( - transport = %addr.transport, - addr = %addr.addr, - error = %e, - "Failed to resolve Ethernet address" - ); - continue; - } - } - } else if addr.transport == "ble" { - #[cfg(ble_available)] - { - match self.resolve_ble_addr(&addr.addr) { - Ok(result) => result, - Err(e) => { - debug!( - transport = %addr.transport, - addr = %addr.addr, - error = %e, - "Failed to resolve BLE address" - ); - continue; - } - } - } - #[cfg(not(ble_available))] - { - debug!(transport = %addr.transport, "BLE transport not available on this build"); - continue; - } - } else { - let tid = if spec.kind == "udp" - && let Ok(remote_socket_addr) = addr.addr.parse::() - { - match self.find_udp_transport_for_remote_addr(remote_socket_addr, spec.instance) - { - Some((id, _)) => id, - None => { - debug!( - transport = %addr.transport, - addr = %addr.addr, - "No compatible operational UDP transport for address" - ); - continue; - } - } - } else if spec.instance.is_some() { - // Only the UDP resolver above can honour an instance name. - // Matching any instance of the type here would be the - // silent wrong-lane substitution this whole mechanism - // exists to prevent, so refuse instead. - debug!( - transport = %addr.transport, - addr = %addr.addr, - "Instance-qualified address for a transport type that \ - does not support instance selection" - ); - continue; - } else { - match self.find_transport_for_type(spec.kind) { - Some(id) => id, - None => { - debug!( - transport = %addr.transport, - addr = %addr.addr, - "No operational transport for address type" - ); - continue; - } - } - }; - (tid, TransportAddr::from_string(&addr.addr)) + let Some((transport_id, remote_addr)) = self.resolve_peer_address(addr) else { + continue; }; if self.is_connecting_to_peer_on_path(&peer_node_addr, transport_id, &remote_addr) { @@ -3243,27 +3289,55 @@ impl Node { .into_iter() .filter(|addr| !(addr.transport == "udp" && addr.addr.eq_ignore_ascii_case("nat"))) .collect(); - let has_alternative = concrete - .iter() - .any(|addr| !self.active_peer_matches_candidate(&peer_node_addr, addr)); - let attempt_candidates: Vec<_> = if has_alternative { - concrete - .into_iter() - .filter(|addr| !self.active_peer_matches_candidate(&peer_node_addr, addr)) - .collect() - } else if self.active_peer_needs_same_path_refresh(&peer_node_addr) { - concrete - } else { - Vec::new() - }; + // Every address the peer is not already on. (The same-path case — + // the one address it is on, gone quiet — was once re-dialled from + // here; the path's own heartbeat and the link-dead reap own that + // now, see below.) + let attempt_candidates: Vec<_> = concrete + .into_iter() + .filter(|addr| !self.active_peer_matches_candidate(&peer_node_addr, addr)) + .collect(); - if attempt_candidates.is_empty() { - return Ok(false); + // A peer we hold a session to gets a *path* at each address, under + // that session, never a second handshake: the probe exchange proves + // the path and selection moves traffic if it measures better or the + // current path stops answering. A handshake to a peer that already + // has one is read by the far side as a rekey, and two ends doing it + // at once — both hearing the same medium return — leave the rekey + // and the cross-connection resolution overlapping and the sides on + // different session indices. An address the peer is already + // reachable at is nothing to do; one on a transport that already + // has a path re-points that path (the peer moved); one on a new + // transport adds a path. A peer that has gone quiet is not dialled + // either: a dead session answers no probe, is reaped by the + // link-dead timeout, and is dialled then. See `poll_discovered_peers`. + let mut paths_added = false; + for addr in attempt_candidates { + let Some((transport_id, remote_addr)) = self.resolve_peer_address(&addr) else { + continue; + }; + let Some(peer) = self.peers.get(&peer_node_addr) else { + continue; + }; + if peer.is_reachable_at(transport_id, &remote_addr) { + continue; + } + let known_transport = peer.path_on(transport_id).is_some(); + info!( + peer = %self.peer_display_name(&peer_node_addr), + %transport_id, + addr = %remote_addr, + "{}", + if known_transport { + "Configured address moved on a known transport: path re-pointed, not dialled" + } else { + "Configured address on a new transport: added as a path, not dialled" + } + ); + self.add_path_candidate(peer_node_addr, transport_id, remote_addr); + paths_added = true; } - - self.attempt_peer_address_list(peer_config, peer_identity, false, &attempt_candidates) - .await?; - Ok(true) + Ok(paths_added) } async fn peer_address_candidates(&self, peer_config: &PeerConfig) -> Vec { @@ -3362,14 +3436,17 @@ impl Node { /// auto-reconnect). Reuses the same connection path as auto-connect /// peers. Returns JSON data on success or an error message. /// - /// For a peer the node is already connected to, the supplied address is - /// tried as an *alternate path* rather than ignored — the same treatment - /// [`Node::update_peers`] gives a refreshed runtime peer. The handshake - /// runs in parallel with the live link and promotion happens only once it - /// authenticates, so an address the caller got wrong cannot displace a - /// healthy path. The response's `refreshed` field reports whether such a - /// handshake was started; it is `false` when the peer is already on this - /// exact path and that path is fresh. + /// For a peer the node already holds a session with, the supplied + /// address is never dialled: it becomes a path candidate under that + /// session (or re-points a path that has stopped answering), the + /// heartbeat tick probes it, and selection moves traffic there only + /// once the peer has answered — the same treatment + /// [`Node::update_peers`] gives a refreshed runtime peer, so an address + /// the caller got wrong cannot displace a healthy path. The response's + /// `refreshed` field reports whether a path was added or re-pointed; + /// it is `false` when the peer is already reachable at that address or + /// that transport's path is carrying acknowledged traffic, and for an + /// ordinary dial to a peer the node does not yet hold. pub(crate) async fn api_connect( &mut self, npub: &str, @@ -3451,16 +3528,18 @@ impl Node { return Err(format!("peer not found: {npub}")); }; - // Read the transport path the peer is actually sending over BEFORE the - // teardown below drops the peer and its link — afterwards there is - // nothing left to derive it from. `current_addr` rather than the - // link's remote address, because roaming updates the former and it is - // the address the pool entry (and its inbound-slot accounting) is - // keyed by. - let transport_path = match (peer.transport_id(), peer.current_addr()) { - (Some(transport_id), Some(addr)) => Some((transport_id, addr.clone())), - _ => None, - }; + // Read every path the peer holds BEFORE the teardown below drops + // the peer and its link — afterwards there is nothing left to + // derive them from. The path's address rather than the link's + // remote address, because roaming updates the former and it is the + // address the pool entry (and its inbound-slot accounting) is keyed + // by. Every path, not the active one alone: a standby on a + // connection-oriented transport holds a pool entry of its own. + let transport_paths: Vec<(TransportId, TransportAddr)> = peer + .paths() + .iter() + .map(|path| (path.transport_id(), path.addr().clone())) + .collect(); // Notify the peer before we tear down the link, so it drops its own // session and re-handshakes symmetrically rather than holding a stale @@ -3483,10 +3562,10 @@ impl Node { // verbatim: closing twice is harmless, because every // `close_connection` implementation is `if let Some(conn) = // pool.remove(addr)` and the connectionless default is a no-op. - if let Some((transport_id, addr)) = transport_path - && let Some(transport) = self.transports.get(&transport_id) - { - transport.close_connection(&addr).await; + for (transport_id, addr) in transport_paths { + if let Some(transport) = self.transports.get(&transport_id) { + transport.close_connection(&addr).await; + } } // Suppress any pending auto-reconnect diff --git a/src/node/mod.rs b/src/node/mod.rs index bda9a946..6aca9de1 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -2697,10 +2697,10 @@ impl Node { transport: handle.and_then(|t| t.name().map(str::to_string)), transport_type: handle.map(|t| t.transport_type().name.to_string()), addr: path.addr().to_string(), - state: format!("{:?}", path.state()).to_lowercase(), + state: path.state().as_str().to_string(), active: Some(path.transport_id()) == active, remote_active: path.remote_active(), - role: format!("{:?}", path.role()).to_lowercase(), + role: path.role().as_str().to_string(), pinned: path.pinned(), last_rtt_ms: path.last_rtt_ms(), min_rtt_ms: path.min_rtt_ms(), diff --git a/src/node/tests/control.rs b/src/node/tests/control.rs index dbb08b99..06fb465b 100644 --- a/src/node/tests/control.rs +++ b/src/node/tests/control.rs @@ -150,13 +150,11 @@ async fn test_api_connect_on_current_fresh_path_is_a_no_op() { } /// `connect` naming a *different* address for a peer the node is already -/// connected to starts an alternate-path handshake instead of silently doing -/// nothing — the fix. -/// -/// The existing peer stays put while that handshake runs: promotion is the -/// handshake's job, not the command's. +/// connected to takes it as a path under the session it has — never a +/// second handshake, which the far side would read as a rekey. The peer and +/// its link stay put; the heartbeat tick probes the path from here. #[tokio::test] -async fn test_api_connect_starts_alternate_path_for_active_peer() { +async fn test_api_connect_takes_an_alternate_address_as_a_path() { let mut nodes = run_tree_test(2, &[(0, 1)], false).await; let node1_addr = *nodes[1].node.node_addr(); @@ -182,26 +180,25 @@ async fn test_api_connect_starts_alternate_path_for_active_peer() { assert_eq!( data["refreshed"], true, - "a new path for an active peer must start a refresh" + "a new address for an active peer is taken as a path" ); assert!( - nodes[0] + !nodes[0] .node .is_connecting_to_peer_on_path(&node1_addr, transport_id, &alternate), - "an outbound leg should exist on the alternate path" + "no handshake: a peer with a session is probed, not dialled" ); + assert_eq!(nodes[0].node.connection_count(), 0); let peer = nodes[0] .node .get_peer(&node1_addr) - .expect("the existing peer must survive the parallel handshake"); - assert_eq!( - peer.link_id(), - link_before, - "the alternate handshake must not tear the live link down before it authenticates" + .expect("the existing peer is untouched"); + assert_eq!(peer.link_id(), link_before, "the live link must not change"); + assert!( + peer.path_on(transport_id).is_some(), + "the transport still has its one path" ); - // Let the alternate handshake run to completion; the peer must still be - // there afterwards. for _ in 0..20 { if process_available_packets(&mut nodes).await == 0 { break; @@ -209,7 +206,7 @@ async fn test_api_connect_starts_alternate_path_for_active_peer() { } assert!( nodes[0].node.get_peer(&node1_addr).is_some(), - "node 1 should still be a peer after the alternate path resolves" + "node 1 is still a peer" ); cleanup_nodes(&mut nodes).await; diff --git a/src/node/tests/multi_path.rs b/src/node/tests/multi_path.rs index aca9b4f2..d29468e9 100644 --- a/src/node/tests/multi_path.rs +++ b/src/node/tests/multi_path.rs @@ -307,6 +307,22 @@ async fn dual_homed_pair() -> (Vec, TransportAddr, TransportAddr) { (nodes, wifi_0, wifi_1) } +/// Hand `nodes[i]` the address `addr` for `peer` on `transport` and run one +/// heartbeat tick: the production route by which a path gets probed +/// (`add_path_candidate` is what a beacon, a config entry or a completed +/// dial leaves behind; the tick is the one issuer of probes). Nothing is +/// delivered here — the caller drives `process_available_packets`. +async fn probe_candidate( + nodes: &mut [TestNode], + i: usize, + peer: NodeAddr, + transport: TransportId, + addr: TransportAddr, +) { + nodes[i].node.add_path_candidate(peer, transport, addr); + nodes[i].node.run_path_heartbeats().await; +} + #[test] fn path_message_round_trips_on_the_wire() { let probe = PathMessage { @@ -345,10 +361,7 @@ async fn a_probe_adds_a_path_at_both_ends_and_the_ack_makes_it_live() { let cable = nodes[0].transport_id; // Node 1 probes node 0 over the wifi. - nodes[1] - .node - .maybe_probe_path(addr_0, wifi(), wifi_0.clone()) - .await; + probe_candidate(&mut nodes, 1, addr_0, wifi(), wifi_0.clone()).await; { let peer = nodes[1].node.get_peer(&addr_0).unwrap(); let path = peer @@ -358,8 +371,9 @@ async fn a_probe_adds_a_path_at_both_ends_and_the_ack_makes_it_live() { assert!(path.tx_live_at_ms().is_none()); } - // Probe reaches node 0. - assert_eq!(process_available_packets(&mut nodes).await, 1); + // The probe reaches node 0 (the tick heartbeats the active path too, + // now that the peer has two). + assert!(process_available_packets(&mut nodes).await >= 1); { let peer = nodes[0].node.get_peer(&addr_1).unwrap(); let path = peer.path_on(wifi()).expect("the receiver adds the path"); @@ -374,8 +388,8 @@ async fn a_probe_adds_a_path_at_both_ends_and_the_ack_makes_it_live() { assert_eq!(peer.paths().len(), 2); } - // Ack reaches node 1. - assert_eq!(process_available_packets(&mut nodes).await, 1); + // The ack reaches node 1. + assert!(process_available_packets(&mut nodes).await >= 1); { let peer = nodes[1].node.get_peer(&addr_0).unwrap(); let path = peer.path_on(wifi()).unwrap(); @@ -446,6 +460,121 @@ async fn a_beacon_from_a_live_peer_on_a_new_transport_probes_instead_of_dialling ); } +#[tokio::test] +async fn a_second_handshake_creates_no_path_state_and_the_address_is_probed_instead() { + let (mut nodes, wifi_0, _wifi_1) = dual_homed_pair().await; + let addr_0 = *nodes[0].node.node_addr(); + let addr_1 = *nodes[1].node.node_addr(); + let cable = nodes[0].transport_id; + let session_before = ( + nodes[0].node.get_peer(&addr_1).unwrap().our_index(), + nodes[1].node.get_peer(&addr_0).unwrap().our_index(), + ); + + // Node 1 dials node 0 over the wifi, as a static config listing both + // addresses does at startup, while the cable session is live. The + // responder answers it as it would any msg1 from a peer it holds a + // session with — a duplicate here, the session being seconds old — + // and neither end derives a path from the handshake. + let identity_0 = PeerIdentity::from_pubkey_full(nodes[0].node.identity().pubkey_full()); + nodes[1] + .node + .initiate_connection(wifi(), wifi_0.clone(), identity_0) + .await + .expect("dial starts"); + for _ in 0..8 { + if process_available_packets(&mut nodes).await == 0 { + break; + } + } + + // Neither end re-peered: one peer each, the session untouched. + assert_eq!(nodes[0].node.peer_count(), 1); + assert_eq!(nodes[1].node.peer_count(), 1); + assert_eq!( + ( + nodes[0].node.get_peer(&addr_1).unwrap().our_index(), + nodes[1].node.get_peer(&addr_0).unwrap().our_index(), + ), + session_before, + "the handshake did not replace the session" + ); + // The dialler holds the address as an unproven candidate, traffic on + // the cable; the responder learned nothing from the handshake. + let p1 = nodes[1].node.get_peer(&addr_0).unwrap(); + assert_eq!(p1.transport_id(), Some(cable), "traffic stays on the cable"); + let wifi_path = p1 + .path_on(wifi()) + .expect("the dialled address is a candidate"); + assert_eq!(wifi_path.state(), PathState::Probing); + assert!(!wifi_path.acked_once(), "a handshake proves no path"); + assert!( + nodes[0] + .node + .get_peer(&addr_1) + .unwrap() + .path_on(wifi()) + .is_none(), + "the responder derives no path from a handshake" + ); + + // The connection the dial opened is left for the probe to ride: on a + // connection-oriented transport it is the path's socket, and closing + // it would have the first probe dial again (or, at the responder, find + // an ephemeral port that cannot be dialled). Loopback records the + // close it would have been asked for. + for node in &nodes { + let closed = match node.node.transports.get(&wifi()).expect("wifi transport") { + TransportHandle::Loopback(t) => t.closed(), + _ => unreachable!("tests run over loopback"), + }; + assert!( + closed.is_empty(), + "the dial's connection is kept for the candidate: {closed:?}" + ); + } + + // The heartbeat tick proves it, at both ends, under the shared session. + nodes[1].node.run_path_heartbeats().await; + for _ in 0..8 { + if process_available_packets(&mut nodes).await == 0 { + break; + } + } + let p1 = nodes[1].node.get_peer(&addr_0).unwrap(); + assert_eq!(p1.path_on(wifi()).unwrap().state(), PathState::Live); + assert_eq!(p1.transport_id(), Some(cable), "traffic still on the cable"); + let p0 = nodes[0].node.get_peer(&addr_1).unwrap(); + assert!( + p0.path_on(wifi()).is_some(), + "the probe taught the responder the path" + ); + assert_eq!(p0.transport_id(), Some(cable)); + + // An operator disconnect closes every path's connection, the standby's + // included: on a connection-oriented transport each holds a pool entry. + let npub_0 = nodes[0].node.identity().npub(); + nodes[1] + .node + .api_disconnect(&npub_0) + .await + .expect("disconnect"); + let closed_wifi = match nodes[1].node.transports.get(&wifi()).unwrap() { + TransportHandle::Loopback(t) => t.closed(), + _ => unreachable!(), + }; + assert_eq!( + closed_wifi, + vec![wifi_0.clone()], + "the standby's connection is closed too" + ); + let closed_cable = match nodes[1].node.transports.get(&cable).unwrap() { + TransportHandle::Loopback(t) => t.closed(), + _ => unreachable!(), + }; + assert_eq!(closed_cable.len(), 1, "and the active path's"); +} + #[tokio::test] async fn a_beaconed_path_is_probed_by_the_next_heartbeat_tick_and_once_only() { let (mut nodes, wifi_0, _wifi_1) = dual_homed_pair().await; @@ -509,24 +638,6 @@ async fn an_ack_for_no_outstanding_probe_changes_nothing() { ); } -#[tokio::test] -async fn the_active_path_is_not_probed() { - let (mut nodes, _wifi_0, _wifi_1) = dual_homed_pair().await; - let addr_0 = *nodes[0].node.node_addr(); - let cable = nodes[0].transport_id; - let cable_addr_0 = nodes[0].addr.clone(); - - nodes[1] - .node - .maybe_probe_path(addr_0, cable, cable_addr_0) - .await; - assert_eq!( - nodes[0].packet_rx.len(), - 0, - "the handshake proved the active path" - ); -} - // ============================================================================ // Presence loss withdraws a path (design §5–6) // ============================================================================ @@ -537,14 +648,8 @@ async fn pair_with_wifi_live() -> (Vec, TransportAddr, TransportAddr) let (mut nodes, wifi_0, wifi_1) = dual_homed_pair().await; let addr_0 = *nodes[0].node.node_addr(); let addr_1 = *nodes[1].node.node_addr(); - nodes[1] - .node - .maybe_probe_path(addr_0, wifi(), wifi_0.clone()) - .await; - nodes[0] - .node - .maybe_probe_path(addr_1, wifi(), wifi_1.clone()) - .await; + probe_candidate(&mut nodes, 1, addr_0, wifi(), wifi_0.clone()).await; + probe_candidate(&mut nodes, 0, addr_1, wifi(), wifi_1.clone()).await; for _ in 0..4 { if process_available_packets(&mut nodes).await == 0 { break; @@ -595,7 +700,9 @@ async fn losing_the_active_transport_moves_traffic_to_the_live_standby() { assert_eq!(seeded_by, Some(wifi())); // The same session carries on: a frame sent now goes out on the wifi - // and decrypts at the far end, which never saw a switch. + // and decrypts at the far end. (The withdrawal also told the far end, + // with a PathClose on the wifi, that the cable is gone: it moves its + // own traffic to the wifi on hearing it, and re-peers nowhere.) let before = nodes[0].packet_rx.len(); nodes[1] .node @@ -603,15 +710,21 @@ async fn losing_the_active_transport_moves_traffic_to_the_live_standby() { .await .expect("send over the standby"); assert_eq!(nodes[0].packet_rx.len(), before + 1); - let packet = nodes[0].packet_rx.try_recv().unwrap(); - assert_eq!(packet.transport_id, wifi()); - nodes[0].node.handle_encrypted_frame(packet).await; + while let Ok(packet) = nodes[0].packet_rx.try_recv() { + assert_eq!(packet.transport_id, wifi()); + nodes[0].node.handle_encrypted_frame(packet).await; + } let far = nodes[0] .node .get_peer(nodes[1].node.node_addr()) .expect("no re-peering"); assert_eq!(far.consecutive_decrypt_failures(), 0); - assert_eq!(far.transport_id(), Some(cable), "the far end did not move"); + assert_eq!(nodes[0].node.peer_count(), 1); + assert_eq!( + far.transport_id(), + Some(wifi()), + "told the cable is gone, the far end moved too" + ); } #[tokio::test] @@ -621,10 +734,7 @@ async fn losing_the_active_transport_with_only_a_probing_standby_reaps() { let cable = nodes[1].transport_id; // Probe sent, ack never processed: the wifi path is unproven. - nodes[1] - .node - .maybe_probe_path(addr_0, wifi(), wifi_0.clone()) - .await; + probe_candidate(&mut nodes, 1, addr_0, wifi(), wifi_0.clone()).await; assert_eq!( nodes[1] .node @@ -662,21 +772,39 @@ async fn a_dead_path_is_reprobed_when_its_transport_returns_and_forgotten_after_ let (mut nodes, wifi_0, _wifi_1) = pair_with_wifi_live().await; let addr_0 = *nodes[0].node.node_addr(); - nodes[1].node.withdraw_transport(wifi()).await; - - // Presence returns: the next discovery tick probes it again, and the - // ack brings it back Live with its history. (The withdrawal also sent - // node 0 a PathClose on the cable, which it processes here too.) - nodes[1].node.reset_probe_backoff_on_transport(wifi()); - nodes[1] + let samples_before = nodes[1] .node - .maybe_probe_path(addr_0, wifi(), wifi_0.clone()) - .await; + .get_peer(&addr_0) + .unwrap() + .path_on(wifi()) + .unwrap() + .rtt_samples(); + assert!(samples_before > 0); + nodes[1].node.withdraw_transport(wifi()).await; + // (The withdrawal also sent node 0 a PathClose on the cable.) for _ in 0..4 { if process_available_packets(&mut nodes).await == 0 { break; } } + let dead = nodes[1] + .node + .get_peer(&addr_0) + .unwrap() + .path_on(wifi()) + .unwrap(); + assert_eq!(dead.state(), PathState::Dead); + assert_eq!(dead.addr(), &wifi_0); + // Dead: the heartbeat tick leaves it alone. + nodes[1].node.run_path_heartbeats().await; + assert!( + nodes[0].packet_rx.try_recv().is_err(), + "a Dead path is not probed" + ); + + // Presence returns: the path is Probing again, the next heartbeat tick + // probes it, and the ack brings it back Live with its history. + nodes[1].node.reset_probe_backoff_on_transport(wifi()); assert_eq!( nodes[1] .node @@ -685,7 +813,24 @@ async fn a_dead_path_is_reprobed_when_its_transport_returns_and_forgotten_after_ .path_on(wifi()) .unwrap() .state(), - PathState::Live + PathState::Probing + ); + nodes[1].node.run_path_heartbeats().await; + for _ in 0..4 { + if process_available_packets(&mut nodes).await == 0 { + break; + } + } + let revived = nodes[1] + .node + .get_peer(&addr_0) + .unwrap() + .path_on(wifi()) + .unwrap(); + assert_eq!(revived.state(), PathState::Live); + assert!( + revived.rtt_samples() > samples_before, + "re-proved on top of its history, not from nothing" ); // Dead again, and this time the grace expires. @@ -751,6 +896,49 @@ async fn garbage_on_a_standby_path_counts_against_the_peer() { ); } +#[tokio::test] +async fn garbage_on_a_path_the_peer_never_acknowledged_is_not_counted() { + // A candidate is an address we were told about — a beacon, a config + // entry, a handshake source, possibly a replayed one. Until the peer + // answers a probe there, garbage on its transport says nothing about + // the peer, and cannot tear the peering down. + let (mut nodes, _wifi_0, wifi_1) = dual_homed_pair().await; + let addr_1 = *nodes[1].node.node_addr(); + nodes[0] + .node + .add_path_candidate(addr_1, wifi(), wifi_1.clone()); + assert!( + !nodes[0] + .node + .get_peer(&addr_1) + .unwrap() + .path_on(wifi()) + .unwrap() + .acked_once() + ); + let our_index = nodes[0] + .node + .get_peer(&addr_1) + .unwrap() + .our_index() + .unwrap(); + for counter in 0..(THRESHOLD * 2) as u64 { + nodes[0] + .node + .handle_encrypted_frame(ReceivedPacket::new( + wifi(), + wifi_1.clone(), + garbage_frame(our_index, counter), + )) + .await; + } + let peer = nodes[0] + .node + .get_peer(&addr_1) + .expect("an unproven path is not a transport the peer is on"); + assert_eq!(peer.consecutive_decrypt_failures(), 0); +} + // ============================================================================ // Selection (design §8) // ============================================================================ @@ -1039,24 +1227,124 @@ fn a_timed_out_echo_on_an_acknowledged_path_makes_it_suspect_and_selection_leave } #[test] -fn a_path_the_peer_never_acknowledged_backs_off_instead_of_going_suspect() { +fn a_peer_with_one_live_path_is_not_heartbeated() { let mut peer = ActivePeer::new(make_peer_identity(), LinkId::new(1), 0); peer.rebind_transport(tid(CABLE), TransportAddr::from_string("10.0.0.1:1")); - // Promotion-style: Live and tx_live from the handshake, never acked. + // Promotion-style: one Live path from the handshake. Nothing to decide, + // nothing sent — the link heartbeat keeps liveness. + let t0 = 1_000_000; + for k in 0..20 { + assert!( + peer.plan_heartbeats(t0 + k * FAST, &TIMING) + .sends + .is_empty(), + "a single-path peer gets no path probes" + ); + } + // A candidate makes it two: probes start, on the candidate and on the + // active path alike. + peer.add_path(tid(WIFI), TransportAddr::from_string("10.0.0.7:1")); + let plan = peer.plan_heartbeats(t0 + 21 * FAST, &TIMING); + assert_eq!(plan.sends.len(), 2); + assert!( + plan.sends + .iter() + .any(|s| s.transport_id == tid(WIFI) && s.full_size), + "the candidate's first probe is the full-size discovery probe" + ); + assert!( + plan.sends + .iter() + .any(|s| s.transport_id == tid(CABLE) && !s.full_size), + "the handshake proved the active path: its first probe is small" + ); +} + +#[test] +fn an_active_path_the_peer_never_acknowledged_backs_off_instead_of_going_suspect() { + let mut peer = ActivePeer::new(make_peer_identity(), LinkId::new(1), 0); + peer.rebind_transport(tid(CABLE), TransportAddr::from_string("10.0.0.1:1")); + // A candidate on the wifi makes the peer two-path, so the cable — Live + // from the handshake, never acked: an old node — is probed too. + peer.add_path(tid(WIFI), TransportAddr::from_string("10.0.0.7:1")); let t0 = 1_000_000; let plan = peer.plan_heartbeats(t0, &TIMING); - assert_eq!(plan.sends.len(), 1); + assert!(plan.sends.iter().any(|s| s.transport_id == tid(CABLE))); let plan = peer.plan_heartbeats(t0 + TIMEOUT, &TIMING); assert!(plan.suspects.is_empty(), "an old node is not a dead path"); let cable = peer.path_on(tid(CABLE)).unwrap(); assert_eq!(cable.state(), PathState::Live); - assert_eq!(plan.sends.len(), 1, "tried again, with the backoff doubled"); assert!( - peer.plan_heartbeats(t0 + TIMEOUT + 2 * FAST - 1, &TIMING) + plan.sends.iter().any(|s| s.transport_id == tid(CABLE)), + "tried again, with the backoff doubled" + ); + assert!( + !peer + .plan_heartbeats(t0 + TIMEOUT + 2 * FAST - 1, &TIMING) .sends - .is_empty(), + .iter() + .any(|s| s.transport_id == tid(CABLE)), "the unanswered probe pushed the next one out" ); + // However long it goes unanswered, the active path is never given up: + // the handshake proved it, and an old node answers no probe. + let mut t = t0 + TIMEOUT + 2 * FAST; + for _ in 0..200 { + peer.plan_heartbeats(t, &TIMING); + t += TIMEOUT; + } + assert_eq!(peer.path_on(tid(CABLE)).unwrap().state(), PathState::Live); + assert_eq!(peer.transport_id(), Some(tid(CABLE))); +} + +#[test] +fn a_standby_the_peer_never_acknowledges_is_given_up_after_the_discovery_budget() { + let mut peer = dual_path_peer(1, 5); + // A third address the peer never answers on: a NIC it no longer sends + // from, a replayed source, an old node's transport. + peer.add_path(tid(3), TransportAddr::from_string("10.0.0.9:1")); + let mut t = 1_000_000; + let mut probes = 0; + for _ in 0..400 { + let plan = peer.plan_heartbeats(t, &TIMING); + assert!( + plan.suspects.is_empty(), + "never acknowledged: never Suspect" + ); + probes += plan + .sends + .iter() + .filter(|s| s.transport_id == tid(3)) + .count(); + for s in plan.sends { + if s.transport_id != tid(3) { + peer.note_path_ack(s.transport_id, s.probe_id, false, 1, t + 1, u64::MAX); + } + } + if peer.path_on(tid(3)).unwrap().state() == PathState::Dead { + break; + } + t += TIMEOUT; + } + assert_eq!( + peer.path_on(tid(3)).unwrap().state(), + PathState::Dead, + "given up after the budget" + ); + assert_eq!(probes, crate::peer::MAX_DISCOVERY_PROBES as usize); + // Dead: not probed again, and forgotten after the grace. + assert!( + !peer + .plan_heartbeats(t + TIMEOUT, &TIMING) + .sends + .iter() + .any(|s| s.transport_id == tid(3)) + ); + peer.prune_dead_paths(t + 10 * 60_000, 5 * 60_000); + assert!(peer.path_on(tid(3)).is_none()); + // The proven paths are untouched. + assert_eq!(peer.path_on(tid(CABLE)).unwrap().state(), PathState::Live); + assert_eq!(peer.path_on(tid(WIFI)).unwrap().state(), PathState::Live); } #[test] @@ -1097,22 +1385,31 @@ fn withdrawing_our_only_path_leaves_it_suspect_and_still_probed() { fn presence_return_resets_the_discovery_backoff() { let mut peer = ActivePeer::new(make_peer_identity(), LinkId::new(1), 0); peer.rebind_transport(tid(CABLE), TransportAddr::from_string("10.0.0.1:1")); + // Two-path, so the never-acked cable is probed at all. + peer.add_path(tid(WIFI), TransportAddr::from_string("10.0.0.7:1")); let t0 = 1_000_000; + let on_cable = |plan: &crate::peer::HeartbeatPlan| { + plan.sends + .iter() + .filter(|s| s.transport_id == tid(CABLE)) + .count() + }; // Never answered: each timeout doubles the wait. - assert_eq!(peer.plan_heartbeats(t0, &TIMING).sends.len(), 1); - assert_eq!(peer.plan_heartbeats(t0 + TIMEOUT, &TIMING).sends.len(), 1); + assert_eq!(on_cable(&peer.plan_heartbeats(t0, &TIMING)), 1); + assert_eq!(on_cable(&peer.plan_heartbeats(t0 + TIMEOUT, &TIMING)), 1); assert_eq!( - peer.plan_heartbeats(t0 + 2 * TIMEOUT, &TIMING).sends.len(), + on_cable(&peer.plan_heartbeats(t0 + 2 * TIMEOUT, &TIMING)), 1 ); let t = t0 + 3 * TIMEOUT; - assert!( - peer.plan_heartbeats(t, &TIMING).sends.is_empty(), + assert_eq!( + on_cable(&peer.plan_heartbeats(t, &TIMING)), + 0, "the third timeout pushed the next probe past now" ); // The transport's presence cycles: probed at once. peer.reset_probe_backoff_on(tid(CABLE)); - assert_eq!(peer.plan_heartbeats(t, &TIMING).sends.len(), 1); + assert_eq!(on_cable(&peer.plan_heartbeats(t, &TIMING)), 1); } #[test] @@ -1127,7 +1424,7 @@ fn the_discovery_backoff_is_capped_and_never_goes_suspect() { discovery_cap_ms: 2_000, ..TIMING }; - // Drive the wifi standby through many unanswered probes. + // Drive the wifi standby through its unanswered probes. let mut t = 1_000_000; let mut sent_at = Vec::new(); for _ in 0..40 { @@ -1136,9 +1433,19 @@ fn the_discovery_backoff_is_capped_and_never_goes_suspect() { plan.suspects.is_empty(), "never acknowledged: not a dead path" ); + assert!( + plan.sends + .iter() + .filter(|s| s.transport_id == tid(WIFI)) + .all(|s| s.full_size), + "every probe on an unproven path is full-size" + ); if plan.sends.iter().any(|s| s.transport_id == tid(WIFI)) { sent_at.push(t); } + if peer.path_on(tid(WIFI)).unwrap().state() == PathState::Dead { + break; + } t += TIMEOUT; } let gaps: Vec = sent_at.windows(2).map(|w| w[1] - w[0]).collect(); @@ -1147,14 +1454,12 @@ fn the_discovery_backoff_is_capped_and_never_goes_suspect() { gaps.iter().all(|g| *g <= 2_000 + TIMEOUT), "the backoff is capped at discovery_cap_ms: {gaps:?}" ); - assert!( - peer.plan_heartbeats(t, &timing) - .sends - .iter() - .filter(|s| s.transport_id == tid(WIFI)) - .all(|s| s.full_size), - "every probe on an unproven path is full-size" + assert_eq!( + sent_at.len(), + crate::peer::MAX_DISCOVERY_PROBES as usize, + "and the budget bounds it: the path is given up, not probed forever" ); + assert_eq!(peer.path_on(tid(WIFI)).unwrap().state(), PathState::Dead); } #[test] @@ -1298,7 +1603,7 @@ fn unreachable_send_errors_are_classified() { #[tokio::test] async fn the_fast_tick_heartbeats_the_active_path_and_the_ack_measures_it() { - let (mut nodes, _wifi_0, _wifi_1) = pair_with_wifi_live().await; + let (mut nodes, wifi_0, _wifi_1) = dual_homed_pair().await; let addr_0 = *nodes[0].node.node_addr(); let cable = nodes[1].transport_id; assert!( @@ -1312,6 +1617,11 @@ async fn the_fast_tick_heartbeats_the_active_path_and_the_ack_measures_it() { "the handshake proved the cable; no probe has yet" ); + // A second path makes the peer worth heartbeating; the tick then + // probes the active cable as well as the wifi candidate. + nodes[1] + .node + .add_path_candidate(addr_0, wifi(), wifi_0.clone()); nodes[1].node.run_path_heartbeats().await; let queued = nodes[0].packet_rx.len(); assert!(queued >= 1, "a heartbeat probe went out on the cable"); @@ -1556,11 +1866,16 @@ fn a_probe_reaching_a_dead_path_revives_it() { async fn the_first_probe_on_a_path_is_full_size_and_so_is_its_ack() { let (mut nodes, wifi_0, _wifi_1) = dual_homed_pair().await; let addr_0 = *nodes[0].node.node_addr(); - nodes[1] - .node - .maybe_probe_path(addr_0, wifi(), wifi_0.clone()) - .await; - let probe = nodes[0].packet_rx.try_recv().expect("probe queued"); + probe_candidate(&mut nodes, 1, addr_0, wifi(), wifi_0.clone()).await; + // The tick also heartbeated the active cable, small; the wifi probe + // is the discovery one. + let mut probe = None; + while let Ok(packet) = nodes[0].packet_rx.try_recv() { + if packet.transport_id == wifi() { + probe = Some(packet); + } + } + let probe = probe.expect("probe queued"); let mtu = usize::from( nodes[1] .node @@ -1592,9 +1907,14 @@ fn a_proven_path_is_probed_full_size_once_a_minute() { let t0 = 1_000_000; let plan = peer.plan_heartbeats(t0, &TIMING); // dual_path_peer acked via take_probe, which sets no full-size stamp, - // so the first heartbeat is full-size; after that, small until a - // minute has passed. - assert!(plan.sends.iter().all(|s| s.full_size)); + // so the standby's first heartbeat is full-size — the active path's is + // not, the handshake having proved it; after that, small until a + // minute has passed, then full-size on both. + assert!( + plan.sends + .iter() + .all(|s| s.full_size == (s.transport_id == tid(WIFI))) + ); for s in plan.sends { peer.note_path_ack(s.transport_id, s.probe_id, false, 1, t0 + 1, u64::MAX); } @@ -1611,7 +1931,8 @@ fn a_proven_path_is_probed_full_size_once_a_minute() { ); } let plan = peer.plan_heartbeats(t0 + 61_000, &TIMING); - assert!(plan.sends.iter().any(|s| s.full_size)); + assert_eq!(plan.sends.len(), 2); + assert!(plan.sends.iter().all(|s| s.full_size)); } #[tokio::test] @@ -1773,3 +2094,108 @@ fn an_outage_does_not_keep_charging_the_path_s_etx() { "an outage is one event, not a lossy medium" ); } + +// --------------------------------------------------------------------------- +// A handshake proves a path +// --------------------------------------------------------------------------- + +#[test] +fn a_known_transport_at_a_new_address_is_re_pointed_there() { + let mut peer = ActivePeer::new(make_peer_identity(), LinkId::new(1), 0); + peer.rebind_transport(tid(CABLE), TransportAddr::from_string("10.0.0.1:1")); + peer.add_path(tid(WIFI), TransportAddr::from_string("10.0.0.7:1")); + + // `add_path` is "add or return": the address it carries is ignored for + // a path that exists. A move is a separate, explicit step. + peer.add_path(tid(WIFI), TransportAddr::from_string("10.0.0.8:1")); + assert_eq!( + peer.path_on(tid(WIFI)).unwrap().addr(), + &TransportAddr::from_string("10.0.0.7:1") + ); + assert!(peer.refresh_path_addr(tid(WIFI), TransportAddr::from_string("10.0.0.8:1"))); + assert_eq!( + peer.path_on(tid(WIFI)).unwrap().addr(), + &TransportAddr::from_string("10.0.0.8:1") + ); + assert!( + !peer.refresh_path_addr(tid(WIFI), TransportAddr::from_string("10.0.0.8:1")), + "the same address is not a move" + ); + assert!( + !peer.refresh_path_addr(tid(3), TransportAddr::from_string("10.0.0.9:1")), + "a transport with no path is not re-pointed: that is add_path's job" + ); + assert!(peer.path_on(tid(3)).is_none()); +} + +#[tokio::test] +async fn a_dial_to_a_stale_peer_over_another_transport_is_probed_and_taken_when_it_answers() { + let (mut nodes, wifi_0, _wifi_1) = dual_homed_pair().await; + let addr_0 = *nodes[0].node.node_addr(); + let cable = nodes[0].transport_id; + + // Node 1 has not heard node 0 in a long time: the cable is not live. + nodes[1].node.get_peer_mut(&addr_0).unwrap().touch(0); + assert!(!nodes[1].node.active_peer_link_is_live(&addr_0)); + + // An address on another transport arrives. Discovery never dials a + // peer with a session; a caller that does gets the same outcome: the + // handshake settles as a cross-connection, the address is a candidate. + let identity_0 = PeerIdentity::from_pubkey_full(nodes[0].node.identity().pubkey_full()); + nodes[1] + .node + .initiate_connection(wifi(), wifi_0.clone(), identity_0) + .await + .expect("dial starts"); + for _ in 0..8 { + if process_available_packets(&mut nodes).await == 0 { + break; + } + } + assert_eq!(nodes[0].node.peer_count(), 1); + assert_eq!(nodes[1].node.peer_count(), 1); + let p = nodes[1].node.get_peer(&addr_0).unwrap(); + assert_eq!( + p.transport_id(), + Some(cable), + "nothing moved on the handshake" + ); + assert_eq!(p.path_on(wifi()).unwrap().state(), PathState::Probing); + + // The heartbeat tick proves the wifi. The cable's own echoes then time + // out (the peer is silent there, which is why we are here): a hard + // signal, and the mandatory switch takes the path that answered. + nodes[1].node.run_path_heartbeats().await; + for _ in 0..8 { + if process_available_packets(&mut nodes).await == 0 { + break; + } + } + assert!( + nodes[1] + .node + .get_peer_mut(&addr_0) + .unwrap() + .mark_path_suspect(cable) + ); + nodes[1].node.run_path_selection(); + let p = nodes[1].node.get_peer(&addr_0).unwrap(); + assert_eq!(p.path_on(wifi()).unwrap().state(), PathState::Live); + assert_eq!( + p.transport_id(), + Some(wifi()), + "traffic moved to the path that answered" + ); + assert_eq!(p.current_addr(), Some(&wifi_0)); + // The link record followed the traffic. + let link = nodes[1].node.links.get(&p.link_id()).expect("peer's link"); + assert_eq!(link.transport_id(), wifi()); + assert_eq!(link.remote_addr(), &wifi_0); + assert!( + p.path_on(cable).is_some(), + "the cable stays for the heartbeat tick to probe" + ); + // And the session is still shared. + assert_eq!(nodes[0].node.peer_count(), 1); + assert_eq!(nodes[1].node.peer_count(), 1); +} diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index 9a872dc6..35e9eae7 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -1433,7 +1433,7 @@ async fn node_context_mirrors_config_and_immutable_facades() { } #[tokio::test] -async fn update_peers_races_new_alternative_without_dropping_active_peer() { +async fn update_peers_takes_a_new_alternative_as_a_path_without_dropping_active_peer() { // The node's *current* (pre-update) peer set must contain `old_peer`, so it // is baked into the Config at construction (immutable context = sole store). let peer_full = Identity::generate(); @@ -1497,16 +1497,18 @@ async fn update_peers_races_new_alternative_without_dropping_active_peer() { assert_eq!(outcome.updated, 1); assert_eq!(node.peer_count(), 1, "existing link must stay live"); - assert_eq!(node.connection_count(), 1); assert_eq!( - node.connections() - .next() - .and_then(|(_, machine)| machine.conn_source_addr()), - Some(&new_addr) + node.connection_count(), + 0, + "a peer with a session is not dialled on a new address: the address is a path" ); let active = node.get_peer(&peer_node_addr).unwrap(); assert_eq!(active.link_id(), old_link_id); - assert_eq!(active.current_addr(), Some(¤t_addr)); + // The path on this transport was never acknowledged (bound by + // `set_current_addr`, no ack), so it is not eligible and the new address + // re-points it; the heartbeat tick probes it there. An eligible path + // would have kept its address. + assert_eq!(active.current_addr(), Some(&new_addr)); for transport in node.transports.values_mut() { transport.stop().await.ok(); diff --git a/src/peer/active.rs b/src/peer/active.rs index aa5a93b8..5afff44e 100644 --- a/src/peer/active.rs +++ b/src/peer/active.rs @@ -22,6 +22,16 @@ use std::time::{Duration, Instant}; /// How often a full-size (MTU-padded) probe goes out on a proven path. const FULL_SIZE_PROBE_INTERVAL_MS: u64 = 60_000; +/// Discovery probes a never-acknowledged standby path gets before it is +/// given up as `Dead`. At the doubling backoff from the fast interval, +/// capped at the heartbeat interval, that is about half a minute. An +/// address we were told about but the peer never answers on — an old +/// node's transport, a replayed handshake source, a beacon from a NIC the +/// peer no longer sends on — is then no longer probed, and no longer +/// counts as a transport the peer is on. `prune_dead_paths` forgets it +/// after the grace; a fresh candidate starts the count over. +pub const MAX_DISCOVERY_PROBES: u32 = 8; + /// Fold one probe outcome into a path's ETX: the long EWMA (α = 1/32) of /// the delivery ratio, inverted and clamped like the link ETX. Per report a /// raw value is a flap generator on a lightly loaded link; the long average @@ -86,6 +96,20 @@ pub enum PathState { Dead, } +impl PathState { + /// The control-socket spelling: `probing`, `live`, `suspect`, `dead`. + /// A fixed string per variant, so a rename here cannot silently change + /// what `show_peers` and `path_show` emit and what fipstop matches. + pub fn as_str(self) -> &'static str { + match self { + Self::Probing => "probing", + Self::Live => "live", + Self::Suspect => "suspect", + Self::Dead => "dead", + } + } +} + /// Probe bookkeeping for one path: what is outstanding, and when the next /// one may go. #[derive(Clone, Copy, Debug, Default)] @@ -120,16 +144,6 @@ pub struct PathPolicy { pub rtt_window_ms: u64, } -impl PathPolicy { - /// Everything selectable at once; for tests. - pub const PERMISSIVE: Self = Self { - margin: 1.5, - dwell_ms: 0, - min_samples: 0, - rtt_window_ms: u64::MAX, - }; -} - /// A post-switch hold on the link cost the tree sees. #[derive(Clone, Copy, Debug)] struct CostHold { @@ -1111,6 +1125,34 @@ impl ActivePeer { &mut self.send.paths[idx] } + /// The peer's address on `transport_id` is `addr` now. Roams the path + /// there, if any, without touching its state; a `Dead` path brought back + /// at a new address is `Probing` again. Returns whether anything changed. + /// Nothing happens for a transport the peer has no path on — that is + /// [`add_path`](Self::add_path)'s job — nor for the same address. + /// + /// A peer's address on a transport does change under it: a Wi-Fi Aware + /// data path that re-forms comes up with a new link-local, and an + /// address that has moved is not one a probe can reach. Without this the + /// path kept the dead address for as long as it lived. + pub fn refresh_path_addr(&mut self, transport_id: TransportId, addr: TransportAddr) -> bool { + let Some(path) = self.path_on_mut(transport_id) else { + return false; + }; + if path.addr == addr { + return false; + } + path.addr = addr; + #[cfg(any(target_os = "linux", target_os = "macos"))] + path.clear_connected_udp(); + if path.state == PathState::Dead { + path.state = PathState::Probing; + path.dead_since_ms = None; + } + path.probe.next_at_ms = 0; + true + } + /// An authentic frame arrived on `transport_id`: the path there, if any, /// is `rx_live` as of `now_ms`. pub fn note_path_rx(&mut self, transport_id: TransportId, now_ms: u64) { @@ -1125,9 +1167,10 @@ impl ActivePeer { /// backoff has not expired. `backoff_cap_ms` bounds the retry interval, /// which doubles from `base_ms` per unanswered probe. /// - /// Tests only. In production [`plan_heartbeats`](Self::plan_heartbeats) - /// is the one issuer of probes, so no two writers race for - /// `probe.outstanding`. + /// Unit-test sampler for the selection tests, which need a path fed a + /// chosen round trip without driving a whole heartbeat schedule. Every + /// node-level test goes through [`plan_heartbeats`](Self::plan_heartbeats), + /// the one issuer of probes in production. #[cfg(test)] pub fn take_probe( &mut self, @@ -1163,12 +1206,11 @@ impl ActivePeer { remote_id: u32, now_ms: u64, ) { - let path = self.add_path(transport_id, addr.clone()); - if path.addr != addr { - path.addr = addr; - #[cfg(any(target_os = "linux", target_os = "macos"))] - path.clear_connected_udp(); - } + self.add_path(transport_id, addr.clone()); + self.refresh_path_addr(transport_id, addr); + let path = self + .path_on_mut(transport_id) + .expect("path added just above"); path.remote_id = Some(remote_id); if path.state == PathState::Dead { // The peer is probing a path we had given up on: it is back, @@ -1497,6 +1539,15 @@ impl ActivePeer { /// The silence hint: our active path silent for two of the peer's /// intervals on it while a standby hears the peer triggers a probe now, /// never `Suspect` (see §7 for the loop that would otherwise follow). + /// + /// A peer with one `Live` path is not heartbeated here at all. + /// Selection has nothing to move to, so a `Suspect` mark on it changes + /// nothing, and the link heartbeat already keeps its liveness; five + /// probes a second on every single-path link would be cost without a + /// decision behind it. Probes start the moment a second path — a + /// candidate included — exists, which is when a verdict can act; and a + /// lone path that is `Suspect` (the peer closed it, nowhere to go) is + /// probed so its ack can bring it back. pub fn plan_heartbeats(&mut self, now_ms: u64, timing: &HeartbeatTiming) -> HeartbeatPlan { let HeartbeatTiming { fast_ms, @@ -1506,6 +1557,23 @@ impl ActivePeer { } = *timing; let mut plan = HeartbeatPlan::default(); let active = self.send.active; + let alone_and_live = self + .send + .paths + .iter() + .filter(|p| p.state != PathState::Dead) + .count() + < 2 + && active + .and_then(|i| self.send.paths.get(i)) + .is_some_and(|p| p.state == PathState::Live); + if alone_and_live { + for path in self.send.paths.iter_mut() { + path.probe.outstanding = None; + path.probe.timed_out = None; + } + return plan; + } let newest_rx = self.send.paths.iter().filter_map(|p| p.rx_live_at_ms).max(); for (i, path) in self.send.paths.iter_mut().enumerate() { if path.state == PathState::Dead { @@ -1549,6 +1617,16 @@ impl ActivePeer { path.probe.next_at_ms = 0; } else { path.probe.unanswered = path.probe.unanswered.saturating_add(1); + // A standby the peer has never answered on is given up + // after the discovery budget. Never the active path: + // an old node answers no probe there either, and the + // handshake proved it. + if !ours && path.probe.unanswered >= MAX_DISCOVERY_PROBES { + path.state = PathState::Dead; + path.dead_since_ms = Some(now_ms); + path.probe.timed_out = None; + continue; + } } } @@ -1582,11 +1660,16 @@ impl ActivePeer { .min(discovery_cap_ms.max(interval)) }; path.probe.next_at_ms = now_ms.saturating_add(delay.max(1)); - let full_size = !path.acked_once - || path - .last_full_probe_ms - .is_none_or(|t| now_ms.saturating_sub(t) >= FULL_SIZE_PROBE_INTERVAL_MS); - if full_size { + // Discovery probes are full-size so a medium that passes small + // frames and drops large ones never proves itself — except on + // the active path, which the handshake proved and whose MTU it + // seeded; there the once-a-minute full-size probe is enough. + let full_size = (!path.acked_once && !ours) + || match path.last_full_probe_ms { + Some(t) => now_ms.saturating_sub(t) >= FULL_SIZE_PROBE_INTERVAL_MS, + None => !ours, + }; + if full_size || path.last_full_probe_ms.is_none() { path.last_full_probe_ms = Some(now_ms); } plan.sends.push(HeartbeatSend { @@ -1658,11 +1741,24 @@ impl ActivePeer { self.send.active = Some(new_active); } - /// Clear the probe backoff on every path over `transport_id`. - pub fn reset_probe_backoff_on(&mut self, transport_id: TransportId) { - if let Some(path) = self.path_on_mut(transport_id) { - path.reset_probe_backoff(); + /// The transport `transport_id` is back: clear the probe backoff on + /// its path, and a path given up as `Dead` while it was gone is + /// `Probing` again, its RTT window and ETX intact, so the next + /// heartbeat tick re-proves it rather than measuring it from nothing. + /// Returns whether a `Dead` path was revived. + pub fn reset_probe_backoff_on(&mut self, transport_id: TransportId) -> bool { + let Some(path) = self.path_on_mut(transport_id) else { + return false; + }; + path.reset_probe_backoff(); + if path.state == PathState::Dead { + path.state = PathState::Probing; + path.dead_since_ms = None; + path.probe.outstanding = None; + path.probe.timed_out = None; + return true; } + false } // === Handshake Resend === diff --git a/src/peer/mod.rs b/src/peer/mod.rs index 08212a85..58da6c12 100644 --- a/src/peer/mod.rs +++ b/src/peer/mod.rs @@ -9,8 +9,9 @@ mod active; pub(crate) mod machine; pub use active::{ - ActivePeer, ConnectivityState, HeartbeatPlan, HeartbeatSend, HeartbeatTiming, PathPolicy, - PathState, PathSwitch, PathWithdrawal, PeerPath, SwitchReason, + ActivePeer, ConnectivityState, HeartbeatPlan, HeartbeatSend, HeartbeatTiming, + MAX_DISCOVERY_PROBES, PathPolicy, PathState, PathSwitch, PathWithdrawal, PeerPath, + SwitchReason, }; use crate::NodeAddr; diff --git a/src/proto/fmp/core.rs b/src/proto/fmp/core.rs index 6df59200..7d7740dd 100644 --- a/src/proto/fmp/core.rs +++ b/src/proto/fmp/core.rs @@ -646,6 +646,13 @@ impl Fmp { if snap.has_existing_peer { let peer_addr = *wire.peer_identity.node_addr(); + // Which transport the msg1 arrived on plays no part here: a + // handshake never creates path state. A peer dialling us over + // a second transport gets the same answer as one dialling over + // the first — a rekey or a duplicate — and the transport becomes + // a path only through the authenticated, replay-checked probe + // exchange. Both ends then resolve on the same information; a + // rule that read this end's view of its own liveness split them. match (snap.existing_peer_epoch, wire.remote_epoch) { (Some(existing), Some(new)) if existing != new => { // Epoch mismatch → peer restart. diff --git a/src/proto/fmp/tests/core.rs b/src/proto/fmp/tests/core.rs index ad81608e..39c2ca00 100644 --- a/src/proto/fmp/tests/core.rs +++ b/src/proto/fmp/tests/core.rs @@ -707,3 +707,29 @@ fn retirements_are_grouped_after_drains_and_before_rekey_initiations() { matches!(actions[4], ConnAction::InitiateRekey { peer } if peer == make_node_addr(0x03)) ); } + +#[test] +fn a_live_peer_s_msg1_is_classified_the_same_on_every_transport() { + // Same epoch, session old enough to rekey. The transport the msg1 + // arrived on is not an input: a handshake never creates path state, so + // a live peer dialling over a second transport is the rekey it looks + // like, exactly as a dial over the first would be. Both ends then + // decide from the same facts. + let mut snap = establish_snapshot(); + snap.has_existing_peer = true; + snap.existing_peer_epoch = Some([1u8; 8]); + snap.has_session = true; + snap.existing_session_age_secs = 31; + let wire = wire_outcome(Some([1u8; 8])); + assert!(matches!( + Fmp::new().establish_inbound(&snap, &wire), + InboundDecision::RekeyRespond { .. } + )); + + // A restart is a restart, whatever transport it arrives on. + let restarted = wire_outcome(Some([2u8; 8])); + assert!(matches!( + Fmp::new().establish_inbound(&snap, &restarted), + InboundDecision::RestartThenPromote { .. } + )); +} diff --git a/src/proto/mmp/metrics.rs b/src/proto/mmp/metrics.rs index 2e18b65e..743bb759 100644 --- a/src/proto/mmp/metrics.rs +++ b/src/proto/mmp/metrics.rs @@ -341,7 +341,6 @@ impl MmpMetrics { } } - /// Smoothed ETX (long-term EWMA), or `None` if not yet initialized. /// The quality index, [`quality_index`](super::quality_index) of the /// smoothed ETX and the smoothed RTT. `None` until both are measured. pub fn quality_index(&self) -> Option { @@ -351,6 +350,7 @@ impl MmpMetrics { } } + /// Smoothed ETX (long-term EWMA), or `None` if not yet initialized. pub fn smoothed_etx(&self) -> Option { if self.etx_trend.initialized() { Some(self.etx_trend.long()) diff --git a/src/transport/loopback.rs b/src/transport/loopback.rs index 6cf21fb0..57ca7181 100644 --- a/src/transport/loopback.rs +++ b/src/transport/loopback.rs @@ -57,6 +57,11 @@ pub struct LoopbackTransport { /// transport's `IFF_RUNNING`. `None`: not interface-bound, no presence /// reported, which is the default. carrier: Mutex>, + /// Every `close_connection` the node asked for, in order. Loopback is + /// connectionless so the close itself does nothing; the record lets a + /// test say whether a socket a connection-oriented transport would own + /// was kept or closed. + closed: Mutex>, } impl LoopbackTransport { @@ -84,9 +89,20 @@ impl LoopbackTransport { registry, discovered: Mutex::new(Vec::new()), carrier: Mutex::new(None), + closed: Mutex::new(Vec::new()), } } + /// Record a `close_connection` call. + pub fn record_close(&self, addr: &TransportAddr) { + self.closed.lock().unwrap().push(addr.clone()); + } + + /// The addresses the node has asked to close on this transport. + pub fn closed(&self) -> Vec { + self.closed.lock().unwrap().clone() + } + /// Pretend this transport is bound to an interface with (or without) /// carrier; `None` returns it to reporting no presence at all. pub fn set_carrier(&self, carrier: Option) { diff --git a/src/transport/mod.rs b/src/transport/mod.rs index cc11123e..bf3e0e80 100644 --- a/src/transport/mod.rs +++ b/src/transport/mod.rs @@ -248,6 +248,20 @@ pub enum TransportError { } impl TransportError { + /// Whether the kernel refused the send for want of a route: the + /// interface is up but nothing is reachable through it. A hard signal + /// that the path is gone (`ENETUNREACH`, `EHOSTUNREACH`), distinct from + /// `is_transient`: the binder is not going to fix this. + pub fn is_unreachable(&self) -> bool { + match self { + Self::Io(e) => matches!( + e.kind(), + std::io::ErrorKind::NetworkUnreachable | std::io::ErrorKind::HostUnreachable + ), + _ => false, + } + } + /// Whether this failure is expected to clear on its own. /// /// The distinction callers need is not *what* went wrong but whether @@ -268,20 +282,6 @@ impl TransportError { /// which is a statement about the peer rather than about this node's /// ability to transmit, and the existing retry paths for them already sit /// at a different layer. - /// Whether the kernel refused the send for want of a route: the - /// interface is up but nothing is reachable through it. A hard signal - /// that the path is gone (`ENETUNREACH`, `EHOSTUNREACH`), distinct from - /// `is_transient`: the binder is not going to fix this. - pub fn is_unreachable(&self) -> bool { - match self { - Self::Io(e) => matches!( - e.kind(), - std::io::ErrorKind::NetworkUnreachable | std::io::ErrorKind::HostUnreachable - ), - _ => false, - } - } - pub fn is_transient(&self) -> bool { match self { // The interface is absent or mid-rebind. The binder is polling for @@ -1250,7 +1250,7 @@ impl TransportHandle { #[cfg(ble_available)] TransportHandle::Ble(t) => t.close_connection_async(addr).await, #[cfg(test)] - TransportHandle::Loopback(_) => {} // connectionless no-op + TransportHandle::Loopback(t) => t.record_close(addr), // connectionless; recorded for tests } } diff --git a/testing/chaos/README.md b/testing/chaos/README.md index 38c3f1f1..fe1d8d69 100644 --- a/testing/chaos/README.md +++ b/testing/chaos/README.md @@ -115,11 +115,10 @@ Explicit topologies exercising non-UDP transports. `node.path.*`; see the file header for what to read from a run. - **dual-udp-flap**: the all-IP twin (`[n01, n02, udp-veth+udp]`): a veth carrying IP with an interface-bound UDP instance at each end, plus UDP - over the bridge. The veth half is handed to the daemon by the runner - (`fipsctl connect ... udp/`) once the pair has peered, and becomes a - path under the existing session. Exercises `udp.interface` on the listen - and per-peer connected sockets, and a configured address on a new - transport becoming a path. + over the bridge. Both are static addresses on the dial owner, so two + handshakes run at startup and the second is kept as a path under the + first's session. Exercises `udp.interface` on the listen and per-peer + connected sockets, and a second handshake to a live peer becoming a path. - **tcp-mesh**: 6-node mesh with 4 UDP and 3 TCP edges. Both transports use static peer config. Netem mutation (30% fraction, every 20-40s) and link flaps (1 link max, 10-20s down). diff --git a/testing/chaos/sim/config_gen.py b/testing/chaos/sim/config_gen.py index 640ac8ef..2b01610e 100644 --- a/testing/chaos/sim/config_gen.py +++ b/testing/chaos/sim/config_gen.py @@ -58,24 +58,31 @@ def generate_peers_block( for peer_id in sorted(outbound_peers): peer = topology.nodes[peer_id] transport = topology.transport_for_edge(node_id, peer_id) - if topology.is_dual_udp_edge(node_id, peer_id): - # The veth half of a dual edge is found by beacon (Ethernet) or - # added as a path by the runner after the pair has peered - # (udp-veth); the bridge half is dialled from here. - transport = "udp" + dual = topology.is_dual_udp_edge(node_id, peer_id) + # (transport, addr, priority) per address. A dual edge's Ethernet + # half is found by beacon, so only its bridge half is listed; a dual + # udp-veth edge lists both halves, bridge first, and the daemon takes + # the second completed handshake as a path under the first's session. + addresses = [] if transport == UDP_VETH: link = next(l for l in topology.udp_veth_links(node_id) if l.peer_id == peer_id) - transport = f"udp/{link.instance}" - addr = link.peer_addr + if dual: + addresses.append((bridge, f"{peer.docker_ip}:{_TRANSPORT_PORTS['udp']}", 1)) + addresses.append((f"udp/{link.instance}", link.peer_addr, 10 if dual else 1)) else: + if dual: + transport = "udp" addr = f"{peer.docker_ip}:{_TRANSPORT_PORTS.get(transport, 2121)}" if transport == "udp": transport = bridge + addresses.append((transport, addr, 1)) lines.append(f' - npub: "{peer.npub}"') lines.append(f' alias: "{peer_id}"') lines.append(f" addresses:") - lines.append(f" - transport: {transport}") - lines.append(f' addr: "{addr}"') + for transport, addr, priority in addresses: + lines.append(f" - transport: {transport}") + lines.append(f' addr: "{addr}"') + lines.append(f" priority: {priority}") lines.append(f" connect_policy: auto_connect") return "\n".join(lines) diff --git a/testing/chaos/sim/runner.py b/testing/chaos/sim/runner.py index 4e3bcf93..a686b7f2 100644 --- a/testing/chaos/sim/runner.py +++ b/testing/chaos/sim/runner.py @@ -389,11 +389,6 @@ class SimRunner: self._sleep(wait) self._take_snapshot("warmup") - # The veth half of every udp-veth+udp edge: UDP has no beacon, so - # the runner hands the daemon the address once the pair has peered - # over the bridge, and it becomes a path under that session. - self._add_udp_veth_paths() - # Populate npub cache after convergence (nodes must be running) if self.peer_churn_mgr: self.peer_churn_mgr.refresh_all_npubs() @@ -404,67 +399,13 @@ class SimRunner: if self.link_swap_mgr: self.link_swap_mgr.setup_initial() - def _add_udp_veth_paths(self, only_node: str | None = None): - """Give each dual udp-veth edge its veth path. - - Sent from the edge's dial owner (the side whose static config holds - the bridge address) as a control-socket ``connect`` naming the - interface-bound instance: to a peer it already holds a session with, - the daemon adds that as a path rather than dialling. Waits for the - bridge session first, so the command cannot become the first dial. - """ - from .control import send_command - - outbound = self.topology.directed_outbound() - for node_id in sorted(self.topology.nodes): - if only_node is not None and node_id != only_node: - continue - for link in self.topology.udp_veth_links(node_id): - if not self.topology.is_dual_udp_edge(node_id, link.peer_id): - continue - if link.peer_id not in outbound.get(node_id, []): - continue - if node_id in self._down_nodes or link.peer_id in self._down_nodes: - continue - container = self.topology.container_name(node_id) - npub = self.topology.nodes[link.peer_id].npub - params = { - "npub": npub, - "address": link.peer_addr, - "transport": f"udp/{link.instance}", - } - added = None - for _ in range(30): - if send_command(container, "path_show", {"npub": npub}) is None: - time.sleep(1) # not peered over the bridge yet - continue - added = send_command(container, "connect", params) - break - if added is None: - log.warning( - "udp-veth path %s -> %s via %s not added", - node_id, link.peer_id, link.instance, - ) - else: - log.info( - "udp-veth path %s -> %s via %s (%s)", - node_id, link.peer_id, link.instance, link.peer_addr, - ) - def _handle_node_restart(self, node_id: str): """Called after a node container is restarted. - Re-adds the node's udp-veth paths once it has re-peered, and for - ephemeral identity nodes waits briefly for the daemon to start, - then queries its new npub and updates the peer churn manager's - cache. + For ephemeral identity nodes, waits briefly for the daemon to + start, then queries its new npub and updates the peer churn + manager's cache. """ - if any( - self.topology.is_dual_udp_edge(node_id, link.peer_id) - for link in self.topology.udp_veth_links(node_id) - ): - time.sleep(2) - self._add_udp_veth_paths(only_node=node_id) if not self.peer_churn_mgr: return if node_id not in self.peer_churn_mgr.ephemeral_nodes: