Merge branch 'master' into next

Carries the platform work up: four commits, one per dependency group,
reorganized from ten before they landed on master. The path-MTU
never-loosen rule scoped to the link that measured it; the two control
socket commands that reported success without doing what they were asked;
the UDP listen socket descriptor handed to an embedder and labelled with
the transport instance it belongs to; and instance-qualified peer
addresses, with the validator that refuses a name no configured transport
answers to.

No conflicts, and no adaptation. `git diff --numstat` for what this merge
changes on next matches the master side file for file and count for
count, so the merge took master's diff verbatim rather than resolving
anything. Nothing here touches the wire.
This commit is contained in:
Johnathan Corgan
2026-08-20 22:03:24 +00:00
16 changed files with 1566 additions and 18 deletions
+95 -1
View File
@@ -257,6 +257,57 @@ with v0.4.x or earlier peers.
validation, since it refuses every inbound offer rather than disabling the
limit. Existing configurations parse unchanged, the key being optional.
- The UDP transport's listen socket descriptor can now be handed to an
embedder, for hosts that associate a socket with one interface or network
and steer inbound traffic by that association rather than routing by
destination address. On such a host a peer reachable only over a secondary
network fails in a way FIPS can neither see nor fix: the address is
well-formed, the send succeeds, the peer replies, and the host discards the
reply before it reaches our socket, so the link retries msg1 forever with no
error surfaced anywhere. The correction is a socket option chosen against
host state FIPS has no basis to reason about, so the descriptor goes to
whoever does. Call `Node::enable_app_owned_udp_fd()` after `Node::new` and
before `start()`, and read `AppOwnedUdpSocket { instance, fd }` off the
returned channel once the transport is up, following the existing
`enable_app_owned_tun` contract. One message is sent per UDP transport that
binds, so a multi-listener configuration yields all of them; nothing is sent
when no UDP transport is configured or one fails to bind, so an embedder
tells "no socket" from "here is the socket" by the receive timing out. The
`instance` field is the name the listener was configured under, `None` for a
single unnamed instance, and it is what makes more than one listener usable:
transports are created by iterating a map, so arrival order is luck, and an
embedder whose whole purpose is to bind one socket to one network would
otherwise have to guess which socket it just received. Guessing wrong pins
one lane's socket to the other lane's network, which is the failure the seam
exists to correct. FIPS keeps owning the socket, and
the descriptor carries no promise beyond "this is the transport's socket,
and it is open now". Two limits: the per-peer connected-UDP sockets that
Linux and macOS open after `start()` returns are not covered, and a
transport that adopts a socket handed in by the traversal bootstrap does not
fire the seam. Unix only, since the Windows UDP backend has no descriptor.
- A peer address may name which *instance* of a transport it belongs to, as
`transport: "udp/aware"` rather than `"udp"`, where the part after the slash
is the key the transport was configured under. A node running several
instances of one type could not be told them apart by a dialer: both bind
wildcard sockets, so the address-family test matches either, and selection
fell through to the lowest transport id. One socket carried every dial and
the other never carried traffic. A bare type is unqualified and matches any
instance, which is what every existing configuration and caller produces, so
nothing changes for a node that does not use the syntax. A qualified name is
never substituted with a different instance: that is the wrong-lane dial the
syntax exists to prevent, so an unmatched name fails the address instead, and
the same name is what an embedder binds by and what the dialer routes on. The
slash is already how FIPS qualifies an instance inside an address
(`eth0/aa:bb:...`) and cannot occur in a type name. Only UDP resolves an
instance name today; an address that qualifies any other transport type is
refused rather than matched loosely. **Because a qualified name never falls
back, the configuration validator rejects one that no configured transport
answers to**, naming the peer, the instance asked for and the instances that
exist. Otherwise the address would simply be skipped at every dial, which is
invisible for a peer that has a second address that works: the lane would
never carry traffic and nothing above debug logging would say so.
### Changed
- `node.rekey.enabled` now means "initiate rekeys" and nothing else. The
@@ -361,7 +412,6 @@ with v0.4.x or earlier peers.
folded into the new tables with a one-time deprecation warning; migrate your
`fips.yaml` to the new keys.
- Inbound traversal offers are now admitted against a per-sender allowance as
well as the global pool. The intake path previously took a permit from a
single semaphore before any identity check, with the sender's npub used only
@@ -818,6 +868,50 @@ with v0.4.x or earlier peers.
counter now charges at the node that makes the decision rather than at the
hop after it.
- `disconnect` on the control socket now closes the transport connection
rather than only the peer. It notified the peer and freed every node-side
structure, sessions, indices, links, address mapping, tree and bloom state,
and never touched the transport, so on a connection-oriented transport (TCP,
Tor, Nym, BLE) the pool entry, the socket and its inbound-slot accounting
survived the peer the node had just forgotten, until the far end closed or
the receive loop errored. An operator who disconnected a peer to free a slot
did not free the slot. No effect on UDP, Ethernet or loopback, whose
`close_connection` is the connectionless no-op. Still not addressed:
`disconnect` reports `peer not found` for an identity that is only
mid-handshake, so withdrawing a peer during its handshake leaves that leg
resending msg1 until the handshake timeout bounds it.
- `connect` on the control socket now tries the address it was given for a
peer the node is already connected to, instead of reporting success without
doing anything. The command built an ephemeral peer configuration and handed
it to the ordinary dial path, which returns success the moment the peer is
already held, so `fipsctl connect` printed success and the node never
attempted the path. An operator moving a peer onto a freshly provisioned
link, or a supervising process that has just seen a second path come up, had
no way to make the node use it: the peer stayed where it first authenticated
until that path died. The address is now tried as an alternate path
alongside the live one, through the same helper a runtime peer refresh uses,
so promotion happens only after the alternate handshake authenticates and a
wrong address cannot displace a healthy link. The response gains an additive
`refreshed` field distinguishing "started an alternate-path handshake" from
"already on this exact path and it is fresh". `connect` stays ephemeral: the
peer is not written to configuration and gets no auto-reconnect.
- A path MTU measured on one link no longer clamps a peer that has moved to
another. Every writer of the per-destination path-MTU cache keeps the
smaller of the existing and incoming value, which is right while a peer
stays put, but the entry is keyed by destination alone. So a peer first
reached over a narrow link stayed clamped to that link's ceiling for the
lifetime of the process: when it later became reachable over a wider
transport, promotion re-seeded, the seed saw a tighter existing value and
declined, and traffic kept running at the old link's ceiling on a link that
could carry far more, with nothing reporting it because the clamp was doing
exactly what it was told. The node now records which transport last seeded
each destination and treats a seed from a different one as authoritative
rather than as a loosening to refuse. Re-seeding the same transport still
keeps the tighter value, so repeated promotion does not reset discovery, and
a destination with no prior seed is unchanged.
### Security
- The peer static key is verified on both FMP handshake paths, not only at
+11
View File
@@ -164,6 +164,17 @@ not reproduced here to avoid duplicating the source.
| `connect` | `npub` (bech32), `address` (transport endpoint), `transport` (`udp`, `tcp`, `tor`, `nym`, `ethernet`) | Asks the node to dial the peer over the named transport. The named transport must be configured and running. Returns the API result on success or an error string on failure. |
| `disconnect` | `npub` (bech32) | Asks the node to drop the link to the named peer. |
`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` 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.
Both commands run on the daemon's main task and may block briefly
while the node mutates its state.
+155 -1
View File
@@ -41,7 +41,7 @@ pub use node::{
NodeConfig, NostrRendezvousConfig, NostrRendezvousPolicy, RateLimitConfig, RekeyConfig,
RendezvousConfig, RetryConfig, SessionConfig, SessionMmpConfig, TreeConfig,
};
pub use peer::{ConnectPolicy, PeerAddress, PeerConfig};
pub use peer::{ConnectPolicy, PeerAddress, PeerConfig, TransportSpec};
pub use transport::{
BleConfig, DirectoryServiceConfig, EthernetConfig, NymConfig, TcpConfig, TorConfig,
TransportInstances, TransportsConfig, UdpConfig,
@@ -1142,6 +1142,52 @@ impl Config {
}
}
// Reject a peer address naming a transport instance that no
// configured transport answers to. A qualified name deliberately never
// falls back: substituting a different instance is the wrong-lane dial
// the syntax exists to prevent, so the dialer refuses the address and
// says so at debug. Where the peer has a second address that does
// resolve, that refusal is invisible — the lane is simply never used,
// which is the failure the instance names were introduced to fix. A
// typo, a renamed transport, or a `Named` config collapsed back to
// `Single` all land here, and all of them are cheaper to find at
// startup than in a packet capture.
for peer in &self.peers {
for addr in &peer.addresses {
let spec = addr.spec();
let Some(want) = spec.instance else {
continue;
};
if spec.kind != "udp" {
return Err(ConfigError::Validation(format!(
"peer `{}` has address `{}` on transport `{}`, but only `udp` resolves an instance name; \
for any other type the dialer would have to pick an arbitrary instance, which is the wrong-lane dial the syntax exists to prevent. \
Drop the `/{want}` qualifier to match any instance of `{}`.",
peer.npub, addr.addr, addr.transport, spec.kind
)));
}
let configured: Vec<&str> = self
.transports
.udp
.iter()
.filter_map(|(name, _)| name)
.collect();
if !configured.contains(&want) {
let known = if configured.is_empty() {
"no named udp instances are configured (the udp transport is a single unnamed instance)".to_string()
} else {
format!("configured udp instances are: {}", configured.join(", "))
};
return Err(ConfigError::Validation(format!(
"peer `{}` has address `{}` on transport `{}`, but no udp transport is configured under the instance name `{want}`; \
a qualified name is never substituted, so this address would be skipped at every dial and the peer reached only over its other addresses, if it has any. \
{known}.",
peer.npub, addr.addr, addr.transport
)));
}
}
}
// Reject rekey triggers that fire immediately and forever. Both
// arms are checked regardless of `node.rekey.enabled` so that
// turning rekey on later cannot surface a config error at a
@@ -2355,6 +2401,114 @@ node:
assert!(!is_loopback_addr_str("example.com:443"));
}
#[test]
fn test_a_peer_address_naming_a_configured_udp_instance_passes_validation() {
let mut config = Config {
peers: vec![PeerConfig {
npub: "npub1peer".to_string(),
addresses: vec![PeerAddress::new("udp/aware", "203.0.113.1:2121")],
..Default::default()
}],
..Default::default()
};
config.transports.udp = TransportInstances::Named(HashMap::from([
("aware".to_string(), UdpConfig::default()),
("infra".to_string(), UdpConfig::default()),
]));
config
.validate()
.expect("an instance name that matches a configured transport must validate");
}
#[test]
fn test_a_peer_address_naming_an_unconfigured_udp_instance_is_rejected() {
let mut config = Config {
peers: vec![PeerConfig {
npub: "npub1peer".to_string(),
addresses: vec![PeerAddress::new("udp/awre", "203.0.113.1:2121")],
..Default::default()
}],
..Default::default()
};
config.transports.udp = TransportInstances::Named(HashMap::from([
("aware".to_string(), UdpConfig::default()),
("infra".to_string(), UdpConfig::default()),
]));
let err = config
.validate()
.expect_err("a typo in an instance name must not validate");
let text = err.to_string();
assert!(
text.contains("awre"),
"the error must name the instance asked for: {text}"
);
assert!(
text.contains("aware") && text.contains("infra"),
"the error must list the instances that do exist: {text}"
);
}
#[test]
fn test_an_unqualified_peer_address_still_validates_against_named_instances() {
let mut config = Config {
peers: vec![PeerConfig {
npub: "npub1peer".to_string(),
addresses: vec![PeerAddress::new("udp", "203.0.113.1:2121")],
..Default::default()
}],
..Default::default()
};
config.transports.udp =
TransportInstances::Named(HashMap::from([("aware".to_string(), UdpConfig::default())]));
config
.validate()
.expect("a bare type matches any instance and must stay valid");
}
#[test]
fn test_a_qualified_peer_address_is_rejected_when_the_udp_transport_is_unnamed() {
let mut config = Config {
peers: vec![PeerConfig {
npub: "npub1peer".to_string(),
addresses: vec![PeerAddress::new("udp/aware", "203.0.113.1:2121")],
..Default::default()
}],
..Default::default()
};
config.transports.udp = TransportInstances::Single(UdpConfig::default());
let err = config
.validate()
.expect_err("a Single config has no instance name to match and must not validate");
assert!(
err.to_string().contains("no named udp instances"),
"the error must say why nothing matched: {err}"
);
}
#[test]
fn test_a_qualified_peer_address_on_a_non_udp_transport_is_rejected() {
let config = Config {
peers: vec![PeerConfig {
npub: "npub1peer".to_string(),
addresses: vec![PeerAddress::new("ethernet/eth0", "eth0/aa:bb:cc:dd:ee:ff")],
..Default::default()
}],
..Default::default()
};
let err = config
.validate()
.expect_err("only udp resolves an instance name, so any other type must be refused");
assert!(
err.to_string().contains("only `udp` resolves"),
"the error must say which transport types support the syntax: {err}"
);
}
#[test]
fn test_validate_loopback_bind_with_external_peer_rejected() {
use crate::config::PeerAddress;
+137 -1
View File
@@ -22,6 +22,74 @@ pub enum ConnectPolicy {
Manual,
}
/// The separator between a transport type and an instance name in a
/// [`PeerAddress::transport`] field (`"udp/aware"`).
///
/// `/` rather than `:` or `.`: it is already how FIPS qualifies an instance
/// inside an *address* (`"eth0/aa:bb:cc:dd:ee:ff"` for Ethernet,
/// `"hci0/AA:BB:…"` for BLE), and it cannot occur in a transport type name,
/// so the split is unambiguous.
const INSTANCE_SEPARATOR: char = '/';
/// A [`PeerAddress::transport`] field, split into a transport *type* and an
/// optional *instance name*.
///
/// A node can run several instances of one transport type
/// ([`TransportInstances::Named`](crate::config::TransportInstances::Named)) —
/// two UDP sockets, say, one pinned to infrastructure Wi-Fi and one to a Wi-Fi
/// Aware data path. Both bind wildcard sockets, so the dialer's address-family
/// test cannot tell them apart: every dial would deterministically take the
/// same instance and the other socket would never carry traffic. Qualifying
/// the transport field with the configured instance name says which one an
/// address belongs to.
///
/// Syntax: `"<type>"` or `"<type>/<instance>"`, where `<instance>` is the key
/// the transport was configured under. A bare type is *unqualified* and
/// matches any instance of that type, which is what every existing config and
/// caller produces — so this is purely additive.
///
/// ```
/// use fips::config::TransportSpec;
///
/// assert_eq!(TransportSpec::parse("udp").instance, None);
/// assert_eq!(TransportSpec::parse("udp/aware").kind, "udp");
/// assert_eq!(TransportSpec::parse("udp/aware").instance, Some("aware"));
/// ```
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct TransportSpec<'a> {
/// The transport type name, as reported by `TransportType::name`.
pub kind: &'a str,
/// The required instance name, or `None` to match any instance.
pub instance: Option<&'a str>,
}
impl<'a> TransportSpec<'a> {
/// Split a transport field into type and optional instance name.
///
/// A field with no separator, an empty type, or an empty instance is
/// treated as an unqualified type — malformed input degrades to the
/// pre-existing behaviour rather than becoming an unmatchable name.
pub fn parse(field: &'a str) -> Self {
match field.split_once(INSTANCE_SEPARATOR) {
Some((kind, instance)) if !kind.is_empty() && !instance.is_empty() => Self {
kind,
instance: Some(instance),
},
_ => Self {
kind: field,
instance: None,
},
}
}
/// Whether a transport of type `kind` configured under `name` satisfies
/// this spec. An unqualified spec accepts any instance; a qualified one
/// accepts only an exact name match, and never falls back.
pub fn matches(&self, kind: &str, name: Option<&str>) -> bool {
self.kind == kind && self.instance.is_none_or(|want| name == Some(want))
}
}
/// A transport-specific address for reaching a peer.
///
/// Each peer can have multiple addresses across different transports,
@@ -29,7 +97,8 @@ pub enum ConnectPolicy {
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct PeerAddress {
/// Transport type (e.g., "udp", "tor", "ethernet").
/// Transport type (e.g., "udp", "tor", "ethernet"), optionally qualified
/// with a named instance (`"udp/aware"`) — see [`TransportSpec`].
pub transport: String,
/// Transport-specific address string.
@@ -106,6 +175,12 @@ impl PeerAddress {
self.seen_at_ms = Some(seen_at_ms);
self
}
/// The [`transport`](Self::transport) field split into type and optional
/// instance name.
pub fn spec(&self) -> TransportSpec<'_> {
TransportSpec::parse(&self.transport)
}
}
/// Configuration for a known peer.
@@ -202,3 +277,64 @@ impl PeerConfig {
matches!(self.connect_policy, ConnectPolicy::AutoConnect)
}
}
#[cfg(test)]
mod transport_spec_tests {
use super::*;
#[test]
fn a_bare_type_is_unqualified_and_matches_any_instance() {
let spec = TransportSpec::parse("udp");
assert_eq!(spec.kind, "udp");
assert_eq!(spec.instance, None);
assert!(spec.matches("udp", None), "an unnamed Single instance");
assert!(spec.matches("udp", Some("aware")), "a named instance");
assert!(!spec.matches("tcp", None), "a different transport type");
}
#[test]
fn a_qualified_type_matches_only_that_instance() {
let spec = TransportSpec::parse("udp/aware");
assert_eq!(spec.kind, "udp");
assert_eq!(spec.instance, Some("aware"));
assert!(spec.matches("udp", Some("aware")));
assert!(
!spec.matches("udp", Some("lan")),
"a different instance must not be substituted"
);
assert!(
!spec.matches("udp", None),
"an unnamed instance cannot satisfy a named request"
);
assert!(!spec.matches("tcp", Some("aware")));
}
#[test]
fn malformed_fields_degrade_to_an_unqualified_type() {
// Neither half may be empty; anything else keeps the whole string as
// the type, so a typo fails the type test rather than silently
// matching some instance.
for field in ["udp/", "/aware", "/"] {
let spec = TransportSpec::parse(field);
assert_eq!(spec.kind, field, "{field}");
assert_eq!(spec.instance, None, "{field}");
}
}
#[test]
fn only_the_first_separator_splits() {
let spec = TransportSpec::parse("udp/a/b");
assert_eq!(spec.kind, "udp");
assert_eq!(spec.instance, Some("a/b"));
}
#[test]
fn peer_address_exposes_its_spec() {
let addr = PeerAddress::new("udp/aware", "[fe80::1%7]:4872");
assert_eq!(addr.spec().instance, Some("aware"));
assert_eq!(
PeerAddress::new("udp", "1.2.3.4:2121").spec().instance,
None
);
}
}
+2
View File
@@ -103,4 +103,6 @@ pub use proto::fmp::{PromotionResult, cross_connection_winner};
pub use peer::{ActivePeer, ConnectivityState, PeerError};
// Re-export node types
#[cfg(unix)]
pub use node::AppOwnedUdpSocket;
pub use node::{Node, NodeError, NodeState, UpdatePeersOutcome};
+42 -2
View File
@@ -775,6 +775,19 @@ impl Node {
/// `path_mtu_lookup` empty for their FipsAddress, causing
/// `per_flow_max_mss` to fall back to the global ceiling and the
/// SYN-time TCP MSS clamp to over-estimate the effective path.
///
/// The never-loosen rule is scoped to a single link. A tighter value is
/// evidence about the path it was measured on, so re-seeding from a
/// *different* transport than the one that last seeded this destination
/// replaces it outright: the peer has moved, and the old measurement
/// describes a path it no longer uses. Without that, a peer once
/// reachable only over a low-MTU link stays clamped to it for the process
/// lifetime even after moving to a wider one.
///
/// A destination with no prior seed keeps the never-loosen rule unchanged
/// — nothing yet says which link its value describes, so a value learned
/// from discovery or from reactive `MtuExceeded` is assumed to be about
/// the link now being seeded and is not discarded.
pub(in crate::node) fn seed_path_mtu_for_link_peer(
&self,
peer_addr: &NodeAddr,
@@ -814,9 +827,33 @@ impl Node {
);
return;
};
// Taken while `path_mtu_lookup` is held. This is the only site that
// locks both, so no lock-order inversion is reachable.
let Ok(mut seeded_by) = self.path_mtu_seeded_by.write() else {
warn!(
peer = %self.peer_display_name(peer_addr),
"seed_path_mtu_for_link_peer: path_mtu_seeded_by write lock poisoned"
);
return;
};
// Only a *prior seed from another transport* proves the peer has
// moved. With no prior seed the existing value came from discovery or
// reactive learning about the path we are seeding now, so the
// never-loosen rule still applies to it.
let prior_seed = seeded_by.get(&fips_addr).copied();
let relinked = prior_seed.is_some_and(|prior| prior != transport_id);
// Recorded whether or not the value changes: the next seed needs to
// know which link this one described, otherwise a peer whose first
// seed was declined never registers a link at all and a later move
// cannot be detected.
seeded_by.insert(fips_addr, transport_id);
match map.get(&fips_addr).copied() {
Some(existing) if existing.mtu <= link_mtu => {
// Keep the tighter learned value; never loosen the clamp.
Some(existing) if !relinked && existing.mtu <= link_mtu => {
// Keep the tighter learned value; never loosen within a link.
// `relinked` is the case upstream's held/release lifecycle does
// not reach: two links to one peer can be up at once, so the
// old entry is never released and a wider seed from the new
// transport would otherwise be refused forever.
debug!(
peer = %self.peer_display_name(peer_addr),
fips_addr = %fips_addr,
@@ -834,6 +871,9 @@ impl Node {
fips_addr = %fips_addr,
link_mtu = link_mtu,
prior = ?other,
prior_transport = ?prior_seed,
transport_id = %transport_id,
relinked = relinked,
map_len = map.len(),
"seed_path_mtu_for_link_peer: wrote link MTU"
);
+133 -12
View File
@@ -409,14 +409,24 @@ impl Node {
/// service. A wildcard IPv4 socket cannot send to an IPv6 link-local
/// target, and vice versa, so callers must choose by socket family rather
/// than by transport type alone.
///
/// `instance` names a specific configured UDP instance
/// ([`TransportSpec`](crate::config::TransportSpec)) and is the only way to
/// discriminate two wildcard sockets: they are family-compatible with the
/// same addresses, so without a name the lowest `TransportId` always wins
/// and one socket carries everything. When it is `Some`, a non-matching
/// instance is never substituted — the caller gets `None` and can say so,
/// rather than dialing down the wrong lane.
fn find_udp_transport_for_remote_addr(
&self,
remote_addr: SocketAddr,
instance: Option<&str>,
) -> Option<(TransportId, SocketAddr)> {
self.transports
.iter()
.filter(|(id, handle)| {
handle.transport_type().name == "udp"
&& instance.is_none_or(|want| handle.name() == Some(want))
&& handle.is_operational()
&& !self.supervisor.nostr_rendezvous.is_bootstrap_transport(id)
})
@@ -1347,7 +1357,7 @@ impl Node {
for event in events {
let crate::mdns::LanEvent::Discovered(peer) = event;
let Some((transport_id, _local_addr)) =
self.find_udp_transport_for_remote_addr(peer.addr)
self.find_udp_transport_for_remote_addr(peer.addr, None)
else {
debug!(
addr = %peer.addr,
@@ -1721,6 +1731,21 @@ impl Node {
match handle.start().await {
Ok(()) => {
// Hand the freshly-bound socket to an embedder that
// armed `enable_app_owned_udp_fd`, labelled with the
// instance it belongs to so two UDP listeners can be
// told apart. Non-UDP handles report `None` for the
// fd by construction, so no transport-type test is
// needed here.
#[cfg(unix)]
if let (Some(tx), Some(fd)) =
(&self.supervisor.udp_fd_tx, handle.raw_fd())
{
let _ = tx.send(crate::node::AppOwnedUdpSocket {
instance: name.clone(),
fd,
});
}
self.transports.insert(id, handle);
Event::SubstrateUp { child }
}
@@ -2634,7 +2659,12 @@ impl Node {
if attempted >= max_attempts {
break;
}
if addr.transport == "udp" && addr.addr.eq_ignore_ascii_case("nat") {
// 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") {
if !allow_bootstrap_nat {
continue;
}
@@ -2684,10 +2714,11 @@ impl Node {
continue;
}
} else {
let tid = if addr.transport == "udp"
let tid = if spec.kind == "udp"
&& let Ok(remote_socket_addr) = addr.addr.parse::<SocketAddr>()
{
match self.find_udp_transport_for_remote_addr(remote_socket_addr) {
match self.find_udp_transport_for_remote_addr(remote_socket_addr, spec.instance)
{
Some((id, _)) => id,
None => {
debug!(
@@ -2698,8 +2729,20 @@ impl Node {
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(&addr.transport) {
match self.find_transport_for_type(spec.kind) {
Some(id) => id,
None => {
debug!(
@@ -3396,11 +3439,21 @@ impl Node {
let current_transport = peer
.transport_id()
.and_then(|id| self.transports.get(&id))
.map(|transport| transport.transport_type().name);
.map(|transport| (transport.transport_type().name, transport.name()));
// Compare against the candidate's *parsed* transport: a peer address
// may name an instance (`"udp/aware"`), while a handle reports its type
// and its instance name separately. Comparing the raw field to the type
// name would call every instance-qualified address an alternative path,
// so a platform lane that re-pushes its peers — Wi-Fi Aware does, on
// every NDP callback — would re-dial a peer it is already connected to,
// forever.
let spec = candidate.spec();
candidate.addr == current_addr
&& current_transport
.map(|transport| transport == candidate.transport)
.map(|(kind, instance)| {
kind == spec.kind && spec.instance.is_none_or(|want| instance == Some(want))
})
.unwrap_or(true)
}
@@ -3411,6 +3464,15 @@ impl Node {
/// Creates an ephemeral peer connection (not persisted to config, no
/// 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.
pub(crate) async fn api_connect(
&mut self,
npub: &str,
@@ -3426,13 +3488,43 @@ impl Node {
via_nostr: false,
};
// Pre-seed identity cache (same as initiate_peer_connections does)
if let Ok(identity) = PeerIdentity::from_npub(npub) {
// Pre-seed identity cache (same as initiate_peer_connections does).
// An unparseable npub is left to `initiate_peer_connection` below,
// which reports it as `InvalidPeerNpub`.
let peer_identity = PeerIdentity::from_npub(npub).ok();
if let Some(identity) = peer_identity.as_ref() {
self.peer_aliases
.insert(*identity.node_addr(), identity.short_npub());
self.register_identity(*identity.node_addr(), identity.pubkey_full());
}
// A peer we already hold a session to must not fall through to
// `initiate_peer_connection`: that returns Ok(()) the moment the peer
// is in `self.peers`, so the command would report success without ever
// trying the address it was handed. Route it through the same
// alternate-path helper `update_peers` uses instead.
if let Some(identity) = peer_identity
&& self.peers.contains_key(identity.node_addr())
{
let refreshed = self
.try_active_peer_alternative_addresses(&peer_config, identity)
.await
.map_err(|e| e.to_string())?;
info!(
npub = %npub,
address = %address,
transport = %transport,
refreshed = refreshed,
"API connect resolved against an already-connected peer"
);
return Ok(serde_json::json!({
"npub": npub,
"address": address,
"transport": transport,
"refreshed": refreshed,
}));
}
self.initiate_peer_connection(&peer_config)
.await
.map(|()| {
@@ -3446,6 +3538,7 @@ impl Node {
"npub": npub,
"address": address,
"transport": transport,
"refreshed": false,
})
})
.map_err(|e| e.to_string())
@@ -3453,15 +3546,27 @@ impl Node {
/// Disconnect a peer via the control API.
///
/// Notifies the peer, removes it locally, and suppresses auto-reconnect.
/// Notifies the peer, removes it locally, closes the transport connection
/// it was using, and suppresses auto-reconnect.
pub(crate) async fn api_disconnect(&mut self, npub: &str) -> Result<serde_json::Value, String> {
let peer_identity =
PeerIdentity::from_npub(npub).map_err(|e| format!("invalid npub '{npub}': {e}"))?;
let node_addr = *peer_identity.node_addr();
if !self.peers.contains_key(&node_addr) {
let Some(peer) = self.peers.get(&node_addr) else {
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,
};
// Notify the peer before we tear down the link, so it drops its own
// session and re-handshakes symmetrically rather than holding a stale
@@ -3474,6 +3579,22 @@ impl Node {
// Remove the peer (full cleanup: sessions, indices, links, tree, bloom)
self.remove_active_peer(&node_addr);
// Tear down the transport connection, not just the node-side state.
// `remove_active_peer` frees every node-side structure but never
// touches the transport, so on a connection-oriented transport the
// pool entry, the socket and its inbound-slot accounting would
// otherwise survive the peer the node has just forgotten — an operator
// who disconnects a peer to free a slot would not free the slot. This
// mirrors `cleanup_stale_connection`, and the reasoning there applies
// 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;
}
// Suppress any pending auto-reconnect
self.peering.reconciler.retry_pending.remove(&node_addr);
+10
View File
@@ -698,6 +698,14 @@ pub(crate) struct Supervisor {
/// [`Node::dns_local_addr`](crate::Node::dns_local_addr).
pub(in crate::node) dns_local_addr: Option<std::net::SocketAddr>,
/// Sender for each UDP listen socket the transport spawn binds — its raw
/// fd and the instance name it was configured under — armed by
/// [`Node::enable_app_owned_udp_fd`](crate::Node::enable_app_owned_udp_fd)
/// and fired from the transport-spawn arm of `start()` once the socket is
/// bound. `None` unless the embedder armed it.
#[cfg(unix)]
pub(in crate::node) udp_fd_tx: Option<std::sync::mpsc::Sender<crate::node::AppOwnedUdpSocket>>,
/// Node-side driver state for the Nostr overlay peer-rendezvous
/// subsystem: the engine handle, its startup timestamp, the one-shot
/// startup-sweep latch, and the per-peer bootstrap-transport bookkeeping
@@ -743,6 +751,8 @@ impl Supervisor {
dns_identity_rx: None,
dns_task: None,
dns_local_addr: None,
#[cfg(unix)]
udp_fd_tx: None,
nostr_rendezvous: crate::nostr::RendezvousDriver::default(),
lan_rendezvous: None,
#[cfg(unix)]
+116
View File
@@ -261,6 +261,37 @@ pub struct UpdatePeersOutcome {
pub unchanged: usize,
}
/// One bound UDP listen socket, handed to an embedder that armed
/// [`Node::enable_app_owned_udp_fd`].
///
/// A bare descriptor would be enough for the single-listener case and useless
/// for any other: a node configured with several named UDP instances
/// ([`TransportInstances::Named`](crate::config::TransportInstances::Named))
/// delivers one message per instance, and the whole point of the seam — the
/// embedder associating a socket with one host network — needs to know *which*
/// socket it is holding. Naming it here rather than making the embedder infer
/// it from arrival order is deliberate: transports are created from a
/// `HashMap`, so arrival order carries no meaning, and guessing wrong pins a
/// lane's socket to another lane's network, which is precisely the fault this
/// seam exists to correct.
///
/// A struct rather than a tuple so the receiving side reads as
/// `socket.instance` / `socket.fd`, and so a future addition (the bound local
/// address, say) does not break every embedder.
#[cfg(unix)]
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AppOwnedUdpSocket {
/// The configured instance name this listener was built from — the key in
/// a `Named` UDP config, and the same name a peer address qualifies its
/// transport field with (`"udp/aware"`, see
/// [`TransportSpec`](crate::config::TransportSpec)). `None` for a
/// `Single` config, which has no name to give.
pub instance: Option<String>,
/// The bound socket's raw descriptor. Borrowed, not owned: FIPS keeps the
/// socket, and the fd is valid only while the transport is running.
pub fd: std::os::unix::io::RawFd,
}
/// Key for addr_to_link reverse lookup.
type AddrKey = (TransportId, TransportAddr);
@@ -358,6 +389,18 @@ pub struct Node {
/// SYN/SYN-ACK clamp can use the smaller of the local-egress floor
/// and the learned per-destination path MTU.
path_mtu_lookup: crate::upper::tun::PathMtuLookup,
/// Which transport last supplied a *link seed* into `path_mtu_lookup`,
/// per destination.
///
/// A `PathMtuEntry` is released when the link that seeded it goes away,
/// but two links to one peer can be up at the same time — a phone on both
/// BLE and Wi-Fi Aware, say. Then nothing releases the first entry and a
/// wider seed from the second transport is refused by the never-loosen
/// rule forever. Recording the seeding transport is what distinguishes a
/// value that still describes the current path from one that describes a
/// path the peer has left. Absent for destinations reached over multiple
/// hops: those are never link-seeded, so never-loosen applies unchanged.
path_mtu_seeded_by: Arc<std::sync::RwLock<HashMap<crate::FipsAddress, TransportId>>>,
// === Transports & Links ===
/// Active transports (owned by Node).
@@ -780,6 +823,7 @@ impl Node {
peer_acl,
host_map,
path_mtu_lookup: Arc::new(std::sync::RwLock::new(HashMap::new())),
path_mtu_seeded_by: Arc::new(std::sync::RwLock::new(HashMap::new())),
#[cfg(unix)]
decrypt_registered_sessions: std::collections::HashSet::new(),
#[cfg(unix)]
@@ -930,6 +974,7 @@ impl Node {
peer_acl,
host_map,
path_mtu_lookup: Arc::new(std::sync::RwLock::new(HashMap::new())),
path_mtu_seeded_by: Arc::new(std::sync::RwLock::new(HashMap::new())),
#[cfg(unix)]
decrypt_registered_sessions: std::collections::HashSet::new(),
#[cfg(unix)]
@@ -3115,6 +3160,77 @@ impl Node {
(outbound_tx, tun_rx)
}
/// Set up an **app-owned UDP socket option**: FIPS keeps the socket, and
/// the embedder gets its raw fd so it can apply a host socket option FIPS
/// has no basis to choose. Call this after [`Node::new`] and **before**
/// [`Self::start`] — the fd does not exist until the transport binds.
///
/// The UDP transport binds one socket and selects the egress path per
/// destination address, which assumes the host routes by destination
/// alone. Not every host does. Where each socket is instead associated
/// with exactly one network interface (or "network") and inbound traffic
/// is steered by that association, a peer reachable only over a secondary,
/// non-default network is unreachable in a way FIPS can neither see nor
/// correct: the address is well-formed, the send succeeds, the peer
/// receives our handshake and replies, and the host discards the reply
/// before it reaches our socket. Handshake msg1 then retries forever with
/// no error surfaced anywhere. Correcting it is a `setsockopt`-class call
/// against host-specific network state, on a descriptor the transport
/// otherwise keeps entirely private.
///
/// ```no_run
/// # async fn f(node: &mut fips::Node) -> Result<(), Box<dyn std::error::Error>> {
/// let rx = node.enable_app_owned_udp_fd(); // after new(), before start()
/// node.start().await?;
/// let socket = rx.recv_timeout(std::time::Duration::from_secs(1))?;
/// # let _ = (socket.instance, socket.fd); Ok(())
/// # }
/// ```
///
/// One message is sent per UDP transport that successfully binds — the
/// usual single-listener configuration therefore yields exactly one, while
/// a config with several named UDP listeners yields one per listener, all
/// of which an embedder pinning sockets to a network needs. Each message
/// carries the instance name its listener was configured under
/// ([`AppOwnedUdpSocket::instance`]), which is the only thing that tells
/// two otherwise-identical descriptors apart: pinning the wrong socket to
/// the wrong network is exactly the fault this seam exists to let an
/// embedder fix, so an fd is never handed over unlabelled. A
/// [`Self::stop`] followed by another [`Self::start`] delivers the new
/// socket's fd on the same channel, since that is a genuinely different
/// descriptor. Nothing at all is sent when no UDP transport is configured
/// or a configured one fails to bind, so a receive that times out is how
/// an embedder tells "no socket" from "here is the socket". Calling this
/// twice replaces the first arming: the last receiver wins.
///
/// Scope, stated precisely so it is not read as more than it is:
///
/// - The fd is the transport's wildcard listen socket. On targets that
/// also run the per-peer connected-UDP fast path (Linux and macOS), the
/// additional `connect()`-ed sockets that path opens per established
/// peer, after `start()` has returned, are not covered by this seam. On
/// targets without that path the wildcard socket is the only UDP socket
/// the transport opens.
/// - Transports that adopt a socket supplied from outside (the
/// NAT-traversal bootstrap handoff) do not fire this, since whoever
/// supplied the socket already held its fd and could bind it before
/// handover.
/// - FIPS retains ownership. The fd is a borrow valid while the node is
/// running; using it after the transport stops can touch an unrelated
/// reused descriptor.
///
/// Unix-only: `RawFd` is a unix concept and the Windows UDP backend has no
/// descriptor. The channel is a [`std::sync::mpsc`] one because the
/// embedder is not necessarily on a tokio runtime; the sending end lives on
/// the supervisor as `udp_fd_tx` and [`Self::start`] fires it from the
/// transport-spawn arm, right after the handle reports a successful start.
#[cfg(unix)]
pub fn enable_app_owned_udp_fd(&mut self) -> std::sync::mpsc::Receiver<AppOwnedUdpSocket> {
let (udp_fd_tx, udp_fd_rx) = std::sync::mpsc::channel();
self.supervisor.udp_fd_tx = Some(udp_fd_tx);
udp_fd_rx
}
/// Address the built-in `.fips` DNS responder is listening on, or `None`
/// when it is not running (`dns.enabled = false`, the bind failed, or the
/// node is stopped).
+313
View File
@@ -0,0 +1,313 @@
//! Control API (`connect` / `disconnect`) behaviour tests.
//!
//! These drive `Node::api_connect` and `Node::api_disconnect` directly —
//! the same entry points the control socket's mutating commands dispatch to
//! (`src/control/commands.rs`) — so the assertions are about node state, not
//! socket framing.
use super::*;
use spanning_tree::{
TestNode, add_loopback_alias, cleanup_nodes, drain_all_packets, make_test_node,
process_available_packets, run_tree_test,
};
/// Count the in-flight handshake legs a node is running toward `peer`.
fn outbound_leg_count(node: &Node, peer: &NodeAddr) -> usize {
node.peer_machines
.values()
.filter(|machine| {
machine.leg().is_some()
&& machine
.conn_expected_identity()
.map(|id| id.node_addr() == peer)
.unwrap_or(false)
})
.count()
}
/// The loopback address string the control API would be handed for a node.
fn loopback_address(node: &TestNode) -> String {
node.addr.to_string()
}
/// `connect` for a peer the node does not know dials it and the handshake
/// completes: the baseline the other tests are contrasted against.
#[tokio::test]
async fn test_api_connect_dials_an_unknown_peer() {
let mut nodes = vec![make_test_node().await, make_test_node().await];
let node0_addr = *nodes[0].node.node_addr();
let node1_addr = *nodes[1].node.node_addr();
let node1_npub = nodes[1].node.npub();
let node1_address = loopback_address(&nodes[1]);
let data = nodes[0]
.node
.api_connect(&node1_npub, &node1_address, "loopback")
.await
.expect("api_connect should dial an unknown peer");
assert_eq!(
data["refreshed"], false,
"a first dial is not an alternate-path refresh"
);
let total = drain_all_packets(&mut nodes, false).await;
assert!(total > 0, "the dial should have produced packets");
assert!(
nodes[0].node.get_peer(&node1_addr).is_some(),
"node 0 should have node 1 as a peer after api_connect"
);
assert!(
nodes[1].node.get_peer(&node0_addr).is_some(),
"node 1 should have node 0 as a peer after api_connect"
);
cleanup_nodes(&mut nodes).await;
}
/// A second `connect` while the first handshake is still in flight must not
/// start a second leg — the `is_connecting_to_peer` guard.
#[tokio::test]
async fn test_api_connect_duplicate_while_connecting_starts_one_leg() {
let mut nodes = vec![make_test_node().await, make_test_node().await];
let node1_addr = *nodes[1].node.node_addr();
let node1_npub = nodes[1].node.npub();
let node1_address = loopback_address(&nodes[1]);
nodes[0]
.node
.api_connect(&node1_npub, &node1_address, "loopback")
.await
.expect("first api_connect should succeed");
assert_eq!(
outbound_leg_count(&nodes[0].node, &node1_addr),
1,
"the first connect should start exactly one handshake leg"
);
// Deliberately do not pump packets: the peer is still mid-handshake.
nodes[0]
.node
.api_connect(&node1_npub, &node1_address, "loopback")
.await
.expect("second api_connect should succeed");
assert_eq!(
outbound_leg_count(&nodes[0].node, &node1_addr),
1,
"a duplicate connect must not start a second handshake leg"
);
cleanup_nodes(&mut nodes).await;
}
/// `connect` naming the path an active peer is already on, while that path is
/// fresh, is a successful no-op — and says so.
///
/// This is the regression guard for the alternate-path fix: a caller that
/// re-announces the same peer at the same address every discovery cycle must
/// not churn a healthy link.
#[tokio::test]
async fn test_api_connect_on_current_fresh_path_is_a_no_op() {
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
let node1_addr = *nodes[1].node.node_addr();
let node1_npub = nodes[1].node.npub();
let node1_address = loopback_address(&nodes[1]);
let link_before = nodes[0]
.node
.get_peer(&node1_addr)
.expect("node 0 should have node 1")
.link_id();
let legs_before = outbound_leg_count(&nodes[0].node, &node1_addr);
let data = nodes[0]
.node
.api_connect(&node1_npub, &node1_address, "loopback")
.await
.expect("api_connect on the current path should succeed");
assert_eq!(
data["refreshed"], false,
"re-announcing the current fresh path is a no-op"
);
assert_eq!(
outbound_leg_count(&nodes[0].node, &node1_addr),
legs_before,
"no new handshake leg for the path the peer is already on"
);
let peer = nodes[0]
.node
.get_peer(&node1_addr)
.expect("the live peer must survive a no-op connect");
assert_eq!(peer.link_id(), link_before, "the live link must not change");
cleanup_nodes(&mut nodes).await;
}
/// `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.
#[tokio::test]
async fn test_api_connect_starts_alternate_path_for_active_peer() {
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
let node1_addr = *nodes[1].node.node_addr();
let node1_npub = nodes[1].node.npub();
let transport_id = nodes[0].transport_id;
// A second address that reaches node 1, standing in for a second path
// coming up.
let alternate = add_loopback_alias(&nodes[1].addr);
assert_ne!(alternate, nodes[1].addr);
let link_before = nodes[0]
.node
.get_peer(&node1_addr)
.expect("node 0 should have node 1")
.link_id();
let data = nodes[0]
.node
.api_connect(&node1_npub, &alternate.to_string(), "loopback")
.await
.expect("api_connect on an alternate path should succeed");
assert_eq!(
data["refreshed"], true,
"a new path for an active peer must start a refresh"
);
assert!(
nodes[0]
.node
.is_connecting_to_peer_on_path(&node1_addr, transport_id, &alternate),
"an outbound leg should exist on the alternate path"
);
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"
);
// 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;
}
}
assert!(
nodes[0].node.get_peer(&node1_addr).is_some(),
"node 1 should still be a peer after the alternate path resolves"
);
cleanup_nodes(&mut nodes).await;
}
/// An unparseable npub is rejected and changes nothing.
#[tokio::test]
async fn test_api_connect_rejects_invalid_npub() {
let mut node = make_node();
let err = node
.api_connect("notanpub", "loopback:0", "loopback")
.await
.expect_err("an invalid npub must be rejected");
assert!(
err.contains("notanpub"),
"the error should name the bad npub, got: {err}"
);
assert_eq!(node.peer_count(), 0);
assert!(node.peer_machines.is_empty());
}
/// `connect` naming a transport the node does not have fails cleanly rather
/// than half-registering a peer. This is also the pre-start case: a node with
/// no transports yet cannot dial anything.
#[tokio::test]
async fn test_api_connect_without_a_matching_transport_fails_cleanly() {
let mut node = make_node();
let peer = make_node();
let peer_npub = peer.npub();
let peer_addr = *peer.node_addr();
let err = node
.api_connect(&peer_npub, "127.0.0.1:1", "tor")
.await
.expect_err("no tor transport is configured");
assert!(
err.contains("no operational transport"),
"unexpected error: {err}"
);
assert_eq!(node.peer_count(), 0, "no peer may be registered");
assert!(
node.peer_machines.is_empty(),
"no handshake leg may be left behind"
);
assert_eq!(outbound_leg_count(&node, &peer_addr), 0);
}
/// `disconnect` on a connectionless transport removes the peer and the
/// transport close degrades to the no-op trait default — no error, no panic.
///
/// Disconnecting again reports `peer not found`, which is also the
/// double-close path: the first call already closed the connection.
#[tokio::test]
async fn test_api_disconnect_on_a_connectionless_transport() {
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
let node1_addr = *nodes[1].node.node_addr();
let node1_npub = nodes[1].node.npub();
nodes[0]
.node
.api_disconnect(&node1_npub)
.await
.expect("api_disconnect should succeed");
assert!(
nodes[0].node.get_peer(&node1_addr).is_none(),
"the peer must be gone"
);
let err = nodes[0]
.node
.api_disconnect(&node1_npub)
.await
.expect_err("a second disconnect has no peer to remove");
assert!(err.contains("peer not found"), "unexpected error: {err}");
cleanup_nodes(&mut nodes).await;
}
/// `disconnect` for a peer the node does not hold is rejected without any
/// partial teardown.
#[tokio::test]
async fn test_api_disconnect_unknown_peer_changes_nothing() {
let mut node = make_node();
let stranger = make_node();
let peers_before = node.peer_count();
let machines_before = node.peer_machines.len();
let links_before = node.links.len();
let err = node
.api_disconnect(&stranger.npub())
.await
.expect_err("an unknown peer cannot be disconnected");
assert!(err.contains("peer not found"), "unexpected error: {err}");
assert_eq!(node.peer_count(), peers_before);
assert_eq!(node.peer_machines.len(), machines_before);
assert_eq!(node.links.len(), links_before);
}
+1
View File
@@ -11,6 +11,7 @@ mod ble;
mod bloom;
mod bloom_poison;
mod bootstrap;
mod control;
mod decrypt_failure;
mod disconnect;
mod discovery;
+20
View File
@@ -56,6 +56,26 @@ pub(super) async fn lock_large_network_test() -> tokio::sync::MutexGuard<'static
LARGE_NETWORK_TEST_LOCK.lock().await
}
/// Register a second loopback address that delivers to the node already
/// listening on `existing`.
///
/// Gives a test node a second reachable path without a second transport:
/// packets sent to the returned address land in the same node's receive
/// channel. Used by the alternate-path tests, whose point is that the node
/// dialing it must treat the new address as a different path to the same
/// peer.
pub(super) fn add_loopback_alias(existing: &TransportAddr) -> TransportAddr {
let tx = LOOPBACK_REGISTRY
.lock()
.unwrap()
.get(existing)
.expect("existing loopback address must be registered")
.clone();
let alias = next_loopback_addr();
LOOPBACK_REGISTRY.lock().unwrap().insert(alias.clone(), tx);
alias
}
/// A test node bundling a Node with its transport and packet channel.
pub(super) struct TestNode {
pub(super) node: Node,
+64 -1
View File
@@ -8,7 +8,9 @@
use super::*;
use crate::config::{Config, TcpConfig};
use crate::transport::tcp::TcpTransport;
use crate::transport::{TransportAddr, TransportHandle, TransportId, packet_channel};
use crate::transport::{
ConnectionState, TransportAddr, TransportHandle, TransportId, packet_channel,
};
use spanning_tree::{
TestNode, cleanup_nodes, drain_all_packets, initiate_handshake, verify_tree_convergence,
};
@@ -396,3 +398,64 @@ async fn test_tcp_oriented_connect_failure_tears_down_link() {
cleanup_nodes(&mut nodes).await;
}
/// `api_disconnect` closes the pooled TCP connection, not just the peer.
///
/// Regression test: node-side teardown used to leave the socket, its pool
/// entry and its inbound-slot accounting alive after the node had forgotten
/// the peer they belonged to.
#[tokio::test]
async fn test_api_disconnect_closes_the_tcp_connection() {
let mut nodes = vec![make_test_node_tcp().await, make_test_node_tcp().await];
initiate_handshake(&mut nodes, 0, 1).await;
let total = drain_all_packets(&mut nodes, false).await;
assert!(total > 0, "should have processed packets");
let addr_1 = *nodes[1].node.node_addr();
let node1_npub = nodes[1].node.npub();
let peer = nodes[0]
.node
.get_peer(&addr_1)
.expect("node 0 should have node 1 as peer");
let peer_transport_id = peer.transport_id().expect("peer should have a transport");
let peer_addr = peer
.current_addr()
.expect("peer should have a current address")
.clone();
assert_eq!(
nodes[0]
.node
.transports
.get(&peer_transport_id)
.unwrap()
.connection_state(&peer_addr),
ConnectionState::Connected,
"the TCP connection should be pooled while the peer is up"
);
nodes[0]
.node
.api_disconnect(&node1_npub)
.await
.expect("api_disconnect should succeed");
assert!(
nodes[0].node.get_peer(&addr_1).is_none(),
"node 0 should have removed node 1"
);
assert_eq!(
nodes[0]
.node
.transports
.get(&peer_transport_id)
.unwrap()
.connection_state(&peer_addr),
ConnectionState::None,
"the pooled TCP connection must not outlive the peer"
);
cleanup_nodes(&mut nodes).await;
}
+433
View File
@@ -1256,6 +1256,69 @@ fn active_peer_same_path_discovery_refreshes_stale_peer() {
));
}
/// An instance-qualified candidate is the peer's *current* path only when it
/// names the instance the peer is actually on. Without this, every qualified
/// address looked like a different path from the `"udp"` a transport reports as
/// its type, so a platform lane that re-pushes its peers — Wi-Fi Aware does, on
/// every data-path callback — would re-dial a peer it is already connected to,
/// for as long as it stayed connected.
#[tokio::test]
async fn an_instance_qualified_candidate_matches_only_its_own_instance() {
let mut listeners = std::collections::HashMap::new();
for name in ["main", "backup"] {
listeners.insert(
name.to_string(),
crate::config::UdpConfig {
bind_addr: Some("127.0.0.1:0".to_string()),
..Default::default()
},
);
}
let mut config = crate::Config::new();
config.transports.udp = crate::config::TransportInstances::Named(listeners);
config.dns.enabled = false;
let mut node = make_node_with(config);
node.start().await.unwrap();
let main_id = *node
.transports
.iter()
.find(|(_, handle)| handle.name() == Some("main"))
.expect("the `main` listener came up")
.0;
let peer_full = Identity::generate();
let peer_identity = PeerIdentity::from_pubkey_full(peer_full.pubkey_full());
let peer_node_addr = *peer_identity.node_addr();
let mut active_peer = ActivePeer::new(peer_identity, LinkId::new(7), Node::now_ms());
active_peer.set_current_addr(main_id, TransportAddr::from_string("127.0.0.1:9"));
node.peers.insert(peer_node_addr, active_peer);
let matches = |transport: &str| {
let candidate = crate::config::PeerAddress::new(transport, "127.0.0.1:9");
node.active_peer_candidate_is_fresh_enough_to_skip(
&peer_node_addr,
std::slice::from_ref(&candidate),
)
};
assert!(
matches("udp"),
"an unqualified candidate still matches, as it always did",
);
assert!(
matches("udp/main"),
"the instance the peer is on is the same path, not an alternative",
);
assert!(
!matches("udp/backup"),
"a different instance is a genuinely different path",
);
node.stop().await.unwrap();
}
#[tokio::test]
async fn node_context_mirrors_config_and_immutable_facades() {
let mut node = make_node();
@@ -1945,6 +2008,142 @@ async fn test_seed_path_mtu_noop_for_unknown_transport() {
);
}
/// The upgrade case, and the reason the seeding transport is tracked.
///
/// A peer first reachable only over a narrow link, then moving to a wider
/// one, must not stay clamped to the narrow link's MTU. Every writer of
/// `path_mtu_lookup` keeps the tighter value, so without recording which link
/// a value described, the low MTU outlives the link it came from and pins the
/// peer for the process lifetime.
#[tokio::test]
async fn test_seed_path_mtu_reseeds_when_peer_moves_to_wider_transport() {
let mut node = make_node();
let (packet_tx, packet_rx) = packet_channel(64);
node.supervisor.packet_tx = Some(packet_tx);
node.packet_rx = Some(packet_rx);
let narrow = make_udp_transport_with_mtu(1, 1280).await;
let wide = make_udp_transport_with_mtu(2, 1452).await;
node.transports.insert(TransportId::new(1), narrow);
node.transports.insert(TransportId::new(2), wide);
let peer_addr = make_node_addr(0xE1);
let fips_addr = crate::FipsAddress::from_node_addr(&peer_addr);
let narrow_addr = TransportAddr::from_string("10.0.0.6:2121");
let wide_addr = TransportAddr::from_string("10.0.0.7:2121");
node.seed_path_mtu_for_link_peer(&peer_addr, TransportId::new(1), &narrow_addr);
assert_eq!(
node.path_mtu_lookup
.read()
.unwrap()
.get(&fips_addr)
.map(|e| e.mtu),
Some(1280),
"first seed takes the narrow link's MTU"
);
// The peer moves to the wider transport.
node.seed_path_mtu_for_link_peer(&peer_addr, TransportId::new(2), &wide_addr);
assert_eq!(
node.path_mtu_lookup
.read()
.unwrap()
.get(&fips_addr)
.map(|e| e.mtu),
Some(1452),
"a seed from a different transport must replace a value describing \
the link the peer has left"
);
for transport in node.transports.values_mut() {
transport.stop().await.ok();
}
}
/// A value learned *about the narrow link* is discarded on the move too — it
/// measured a path the peer no longer uses.
#[tokio::test]
async fn test_seed_path_mtu_discards_learned_value_from_abandoned_link() {
let mut node = make_node();
let (packet_tx, packet_rx) = packet_channel(64);
node.supervisor.packet_tx = Some(packet_tx);
node.packet_rx = Some(packet_rx);
let narrow = make_udp_transport_with_mtu(1, 1280).await;
let wide = make_udp_transport_with_mtu(2, 1452).await;
node.transports.insert(TransportId::new(1), narrow);
node.transports.insert(TransportId::new(2), wide);
let peer_addr = make_node_addr(0xE2);
let fips_addr = crate::FipsAddress::from_node_addr(&peer_addr);
let narrow_addr = TransportAddr::from_string("10.0.0.8:2121");
let wide_addr = TransportAddr::from_string("10.0.0.9:2121");
node.seed_path_mtu_for_link_peer(&peer_addr, TransportId::new(1), &narrow_addr);
// Reactive MtuExceeded tightens further, still on the narrow link.
node.path_mtu_lookup
.write()
.unwrap()
.insert(fips_addr, crate::upper::tun::PathMtuEntry::held(900));
node.seed_path_mtu_for_link_peer(&peer_addr, TransportId::new(2), &wide_addr);
assert_eq!(
node.path_mtu_lookup
.read()
.unwrap()
.get(&fips_addr)
.map(|e| e.mtu),
Some(1452),
"a tighter value measured on the abandoned link must not clamp the new one"
);
for transport in node.transports.values_mut() {
transport.stop().await.ok();
}
}
/// The guard against over-loosening. Promotion re-seeds on every handshake,
/// so discarding a tighter learned value on a *same-link* re-seed would reset
/// genuine PMTU discovery repeatedly and the estimate would never converge.
#[tokio::test]
async fn test_seed_path_mtu_keeps_tighter_value_when_reseeding_same_transport() {
let mut node = make_node();
let (packet_tx, packet_rx) = packet_channel(64);
node.supervisor.packet_tx = Some(packet_tx);
node.packet_rx = Some(packet_rx);
let udp = make_udp_transport_with_mtu(1, 1452).await;
node.transports.insert(TransportId::new(1), udp);
let peer_addr = make_node_addr(0xE3);
let fips_addr = crate::FipsAddress::from_node_addr(&peer_addr);
let transport_addr = TransportAddr::from_string("10.0.0.10:2121");
node.seed_path_mtu_for_link_peer(&peer_addr, TransportId::new(1), &transport_addr);
// Reactive learning tightens the same link.
node.path_mtu_lookup
.write()
.unwrap()
.insert(fips_addr, crate::upper::tun::PathMtuEntry::held(1200));
// Re-promotion on the same transport.
node.seed_path_mtu_for_link_peer(&peer_addr, TransportId::new(1), &transport_addr);
assert_eq!(
node.path_mtu_lookup
.read()
.unwrap()
.get(&fips_addr)
.map(|e| e.mtu),
Some(1200),
"re-seeding the same link must not undo reactive learning"
);
for transport in node.transports.values_mut() {
transport.stop().await.ok();
}
}
// === Outbound admission gate tests ===
/// Inject `count` synthetic active peers into `node.peers` so peer_count()
@@ -2524,6 +2723,240 @@ async fn start_skips_system_tun_when_app_owned() {
node.stop().await.unwrap();
}
/// Config for the app-owned-UDP-fd tests: one loopback UDP transport on an
/// ephemeral port and no DNS, mirroring `make_healthy_node`.
#[cfg(unix)]
fn udp_loopback_config() -> crate::Config {
let mut config = crate::Config::new();
config.transports.udp = crate::config::TransportInstances::Single(crate::config::UdpConfig {
bind_addr: Some("127.0.0.1:0".to_string()),
..Default::default()
});
config.dns.enabled = false;
config
}
/// App-owned UDP fd seam: the embedder gets the descriptor of the socket the
/// transport actually bound, and gets it only once the bind has happened —
/// there is no fd to hand out before `start()`.
#[cfg(unix)]
#[tokio::test]
async fn app_owned_udp_fd_seam_delivers_the_bound_socket() {
let mut node = make_node_with(udp_loopback_config());
let rx = node.enable_app_owned_udp_fd();
assert!(
rx.try_recv().is_err(),
"nothing is delivered at arm time — the socket is not bound until start()",
);
node.start().await.unwrap();
let socket = rx
.try_recv()
.expect("the seam fires once the UDP socket is bound");
let live_fd = node
.transports
.values()
.next()
.expect("the loopback UDP transport came up")
.raw_fd();
assert_eq!(
Some(socket.fd),
live_fd,
"the delivered fd must be the live transport's socket, not some other descriptor",
);
assert_eq!(
socket.instance, None,
"a `Single` UDP config has no instance name to report",
);
assert!(
rx.try_recv().is_err(),
"one UDP transport bound means exactly one delivery",
);
node.stop().await.unwrap();
}
/// One message per UDP transport that binds: an embedder pinning sockets to a
/// network needs every listener's fd, not just the first, so the seam does not
/// latch after the first send.
#[cfg(unix)]
#[tokio::test]
async fn app_owned_udp_fd_seam_delivers_every_udp_listener_that_binds() {
let mut listeners = std::collections::HashMap::new();
for name in ["main", "backup"] {
listeners.insert(
name.to_string(),
crate::config::UdpConfig {
bind_addr: Some("127.0.0.1:0".to_string()),
..Default::default()
},
);
}
let mut config = crate::Config::new();
config.transports.udp = crate::config::TransportInstances::Named(listeners);
config.dns.enabled = false;
let mut node = make_node_with(config);
let rx = node.enable_app_owned_udp_fd();
node.start().await.unwrap();
let mut delivered: Vec<_> = rx
.try_iter()
.map(|socket| (socket.instance, socket.fd))
.collect();
delivered.sort_unstable();
let mut live: Vec<_> = node
.transports
.values()
.filter_map(|handle| {
handle
.raw_fd()
.map(|fd| (handle.name().map(str::to_string), fd))
})
.collect();
live.sort_unstable();
assert_eq!(
delivered, live,
"every UDP listener that bound must be handed out, not just the first",
);
assert_eq!(delivered.len(), 2, "both named listeners bound");
// The label is what makes two descriptors usable: an embedder pins each
// socket to a different network, and arrival order — the transports come
// out of a `HashMap` — cannot tell it which is which.
let names: Vec<_> = delivered
.iter()
.map(|(instance, _)| instance.as_deref())
.collect();
assert!(
names.contains(&Some("main")) && names.contains(&Some("backup")),
"each fd names the configured instance it belongs to, got {names:?}",
);
assert_ne!(
delivered[0].1, delivered[1].1,
"two instances are two distinct sockets",
);
node.stop().await.unwrap();
}
/// No UDP transport means no fd: the channel stays silent rather than
/// delivering a sentinel, so a receive that times out is how an embedder tells
/// "there is no socket" from "here is the socket".
#[cfg(unix)]
#[tokio::test]
async fn app_owned_udp_fd_seam_stays_silent_without_a_udp_transport() {
let mut config = crate::Config::new();
config.dns.enabled = false;
let mut node = make_node_with(config);
let rx = node.enable_app_owned_udp_fd();
// A node with no transports configured fails bring-up with
// `NoOperationalTransports`; asserted so this test cannot silently stop
// exercising the no-UDP path if that outcome ever changes.
let started = node.start().await;
assert!(
started.is_err(),
"a transportless node has no operational transports",
);
assert!(
rx.try_recv().is_err(),
"no UDP transport means no fd is ever delivered",
);
}
/// A UDP transport that never bound has no fd to hand out. The bind address is
/// deliberately unparseable — a busy port would not do it, since
/// `UdpRawSocket::open` sets `SO_REUSEADDR`/`SO_REUSEPORT` before binding.
#[cfg(unix)]
#[tokio::test]
async fn app_owned_udp_fd_seam_stays_silent_when_the_udp_transport_fails_to_start() {
let mut config = crate::Config::new();
config.transports.udp = crate::config::TransportInstances::Single(crate::config::UdpConfig {
bind_addr: Some("not-a-socket-addr".to_string()),
..Default::default()
});
config.dns.enabled = false;
let mut node = make_node_with(config);
let rx = node.enable_app_owned_udp_fd();
// Transport-start failure is warn-and-continue; the node's overall start
// outcome is not what this test pins.
let _ = node.start().await;
assert!(
rx.try_recv().is_err(),
"a UDP transport that never bound has no fd to hand out",
);
}
/// The seam is per-`Node` state with a fresh channel per arming, so an embedder
/// that tears the mesh down and rebuilds the node — a radio off→on cycle — gets
/// the new socket on the new node's channel, with nothing shared between them.
#[cfg(unix)]
#[tokio::test]
async fn app_owned_udp_fd_seam_rearms_on_a_rebuilt_node() {
let mut node_a = make_node_with(udp_loopback_config());
let rx_a = node_a.enable_app_owned_udp_fd();
node_a.start().await.unwrap();
rx_a.try_recv()
.expect("node A's socket reached node A's rx");
node_a.stop().await.unwrap();
drop(node_a);
let mut node_b = make_node_with(udp_loopback_config());
let rx_b = node_b.enable_app_owned_udp_fd();
node_b.start().await.unwrap();
rx_b.try_recv()
.expect("node B's socket reached node B's rx");
// Nothing from B reaches A's channel. (A is dropped, so this is
// `Disconnected` rather than `Empty`; the specific variant is not part of
// the contract, only that no fd arrives.)
assert!(
rx_a.try_recv().is_err(),
"the channels are per-node — no shared or global arming state",
);
node_b.stop().await.unwrap();
}
/// Arming twice on the same node replaces the first arming: the last receiver
/// wins and the earlier one is disconnected. Asserted by firing through the
/// installed sender directly, so the test needs no real bind.
#[cfg(unix)]
#[test]
fn app_owned_udp_fd_seam_second_arm_replaces_the_first() {
let mut node = make_node_with(udp_loopback_config());
let rx1 = node.enable_app_owned_udp_fd();
let rx2 = node.enable_app_owned_udp_fd();
let sent = crate::node::AppOwnedUdpSocket {
instance: Some("aware".to_string()),
fd: 7,
};
node.supervisor
.udp_fd_tx
.as_ref()
.expect("the second arming installed a sender")
.send(sent.clone())
.expect("the surviving receiver is live");
assert_eq!(
rx2.try_recv().ok(),
Some(sent),
"the last receiver armed is the one the node feeds",
);
assert!(
rx1.try_recv().is_err(),
"the replaced receiver gets nothing",
);
}
/// The embedder-facing DNS contract, end to end.
///
/// An embedder that owns the TUN fd (Android `VpnService`) has no system DNS
+20
View File
@@ -844,6 +844,26 @@ impl TransportHandle {
}
}
/// Get the raw file descriptor of the bound socket (UDP only, returns None
/// for other transports and before the transport has started). Unix-only,
/// since `RawFd` is a unix concept and the Windows UDP backend has no
/// descriptor.
#[cfg(unix)]
pub fn raw_fd(&self) -> Option<std::os::unix::io::RawFd> {
match self {
TransportHandle::Udp(t) => t.raw_fd(),
#[cfg(any(target_os = "linux", target_os = "macos"))]
TransportHandle::Ethernet(_) => None,
TransportHandle::Tcp(_) => None,
TransportHandle::Tor(_) => None,
TransportHandle::Nym(_) => None,
#[cfg(target_os = "linux")]
TransportHandle::Ble(_) => None,
#[cfg(test)]
TransportHandle::Loopback(_) => None,
}
}
/// Get the interface name (Ethernet only, returns None for other transports).
pub fn interface_name(&self) -> Option<&str> {
match self {
+14
View File
@@ -85,6 +85,20 @@ impl UdpTransport {
self.local_addr
}
/// Raw file descriptor of the bound listen socket, or `None` before
/// `start_async` has bound one.
///
/// Unix-only: `RawFd` is a unix concept and the Windows backend is built
/// on `tokio::net::UdpSocket` with no descriptor to hand out. The socket
/// stays owned by the transport, so the descriptor is a borrow with no
/// lifetime promise beyond "this is the transport's socket, and it is
/// open right now".
#[cfg(unix)]
pub fn raw_fd(&self) -> Option<std::os::unix::io::RawFd> {
use std::os::unix::io::AsRawFd;
self.socket.as_ref().map(|socket| socket.as_raw_fd())
}
/// Configured recv buffer size — used when opening per-peer
/// `ConnectedPeerSocket`s so they get the same buffer ceiling as
/// the wildcard listen socket.