Merge branch 'master' into next

This commit is contained in:
Johnathan Corgan
2026-05-02 02:14:36 +00:00
12 changed files with 440 additions and 27 deletions
+150 -4
View File
@@ -467,6 +467,8 @@ impl Node {
}
}
self.maybe_run_startup_open_discovery_sweep(&bootstrap)
.await;
self.queue_open_discovery_retries(&bootstrap).await;
}
@@ -617,6 +619,7 @@ impl Node {
warn!(error = %err, "Failed to publish initial Nostr overlay advert");
}
self.nostr_discovery = Some(runtime);
self.nostr_discovery_started_at_ms = Some(Self::now_ms());
info!("Nostr overlay discovery enabled");
}
Err(err) => {
@@ -1184,6 +1187,26 @@ impl Node {
}
async fn queue_open_discovery_retries(&mut self, bootstrap: &std::sync::Arc<NostrDiscovery>) {
self.run_open_discovery_sweep(bootstrap, None, "per-tick")
.await;
}
/// Open-discovery cache sweep. Iterates the cached overlay adverts and
/// queues retries for non-configured, not-yet-connected peers.
///
/// `max_age_secs`, if set, filters out adverts whose `created_at` is
/// older than `now - max_age_secs`. The per-tick sweep passes `None`
/// (relies on the cache's own `valid_until_ms` filter); the one-shot
/// startup sweep passes `Some(startup_sweep_max_age_secs)`.
///
/// `caller` is a short label included in log lines so per-tick and
/// startup sweeps are distinguishable in operator-facing logs.
async fn run_open_discovery_sweep(
&mut self,
bootstrap: &std::sync::Arc<NostrDiscovery>,
max_age_secs: Option<u64>,
caller: &'static str,
) {
if !self.config.node.discovery.nostr.enabled
|| self.config.node.discovery.nostr.policy != crate::config::NostrDiscoveryPolicy::Open
{
@@ -1197,28 +1220,63 @@ impl Node {
.map(|peer| peer.npub.clone())
.collect::<HashSet<_>>();
let now_ms = Self::now_ms();
let now_secs = now_ms / 1000;
let mut enqueue_budget = self.open_discovery_enqueue_budget(&configured_npubs);
if enqueue_budget == 0 {
debug!(
caller = %caller,
"open-discovery sweep: enqueue budget is 0, skipping"
);
return;
}
for (npub, endpoints) in bootstrap.cached_open_discovery_candidates(64).await {
let candidates = bootstrap.cached_open_discovery_candidates(64).await;
let cached_count = candidates.len();
let mut enqueued = 0usize;
let mut skipped_age = 0usize;
let mut skipped_configured = 0usize;
let mut skipped_self = 0usize;
let mut skipped_connected = 0usize;
let mut skipped_retry_pending = 0usize;
let mut skipped_connecting = 0usize;
let mut skipped_no_endpoints = 0usize;
let mut skipped_invalid_npub = 0usize;
for (npub, endpoints, created_at_secs) in candidates {
if enqueue_budget == 0 {
break;
}
if let Some(max_age) = max_age_secs
&& now_secs.saturating_sub(created_at_secs) > max_age
{
skipped_age = skipped_age.saturating_add(1);
continue;
}
if configured_npubs.contains(&npub) {
skipped_configured = skipped_configured.saturating_add(1);
continue;
}
let peer_identity = match PeerIdentity::from_npub(&npub) {
Ok(identity) => identity,
Err(_) => continue,
Err(_) => {
skipped_invalid_npub = skipped_invalid_npub.saturating_add(1);
continue;
}
};
let node_addr = *peer_identity.node_addr();
if node_addr == *self.identity.node_addr() || self.peers.contains_key(&node_addr) {
if node_addr == *self.identity.node_addr() {
skipped_self = skipped_self.saturating_add(1);
continue;
}
if self.peers.contains_key(&node_addr) {
skipped_connected = skipped_connected.saturating_add(1);
continue;
}
if self.retry_pending.contains_key(&node_addr) {
skipped_retry_pending = skipped_retry_pending.saturating_add(1);
continue;
}
let connecting = self.connections.values().any(|conn| {
@@ -1227,6 +1285,7 @@ impl Node {
.unwrap_or(false)
});
if connecting {
skipped_connecting = skipped_connecting.saturating_add(1);
continue;
}
@@ -1246,6 +1305,7 @@ impl Node {
priority = priority.saturating_add(1);
}
if addresses.is_empty() {
skipped_no_endpoints = skipped_no_endpoints.saturating_add(1);
continue;
}
@@ -1266,8 +1326,87 @@ impl Node {
state.retry_after_ms = now_ms;
state.expires_at_ms = Some(self.open_discovery_retry_expires_at_ms(now_ms));
self.retry_pending.insert(node_addr, state);
info!(
caller = %caller,
peer = %peer_identity.short_npub(),
advert_age_secs = now_secs.saturating_sub(created_at_secs),
"open-discovery sweep: queued retry for cached advert"
);
enqueue_budget = enqueue_budget.saturating_sub(1);
enqueued = enqueued.saturating_add(1);
}
// Always log a one-line summary on the startup sweep so operators
// can verify it ran. Per-tick sweeps are noisier; only summarize
// when something happened.
let total_skipped = skipped_age
+ skipped_configured
+ skipped_self
+ skipped_connected
+ skipped_retry_pending
+ skipped_connecting
+ skipped_no_endpoints
+ skipped_invalid_npub;
let should_summarize = caller == "startup" || enqueued > 0;
if should_summarize {
info!(
caller = %caller,
cached = cached_count,
queued = enqueued,
skipped_age = skipped_age,
skipped_configured = skipped_configured,
skipped_self = skipped_self,
skipped_connected = skipped_connected,
skipped_retry_pending = skipped_retry_pending,
skipped_connecting = skipped_connecting,
skipped_no_endpoints = skipped_no_endpoints,
skipped_invalid_npub = skipped_invalid_npub,
skipped_total = total_skipped,
"open-discovery sweep complete"
);
}
}
/// One-shot startup sweep: runs once after the configured settle
/// delay, iterating the cached overlay adverts and queueing retries
/// for any peer with a recent enough advert that we haven't already
/// configured statically or established a link to.
///
/// Gated identically to [`run_open_discovery_sweep`]: requires
/// `node.discovery.nostr.enabled` and `policy == open`.
async fn maybe_run_startup_open_discovery_sweep(
&mut self,
bootstrap: &std::sync::Arc<NostrDiscovery>,
) {
if self.startup_open_discovery_sweep_done {
return;
}
if !self.config.node.discovery.nostr.enabled
|| self.config.node.discovery.nostr.policy != crate::config::NostrDiscoveryPolicy::Open
{
// Mark done so we don't keep re-checking on every tick.
self.startup_open_discovery_sweep_done = true;
return;
}
let Some(started_at_ms) = self.nostr_discovery_started_at_ms else {
return;
};
let now_ms = Self::now_ms();
let delay_ms = self
.config
.node
.discovery
.nostr
.startup_sweep_delay_secs
.saturating_mul(1000);
if now_ms < started_at_ms.saturating_add(delay_ms) {
return;
}
let max_age_secs = self.config.node.discovery.nostr.startup_sweep_max_age_secs;
self.run_open_discovery_sweep(bootstrap, Some(max_age_secs), "startup")
.await;
self.startup_open_discovery_sweep_done = true;
}
fn available_outbound_slots(&self) -> usize {
@@ -1384,7 +1523,7 @@ impl Node {
if let Some(addr) = handle.onion_address() {
endpoints.push(OverlayEndpointAdvert {
transport: OverlayTransportKind::Tor,
addr: addr.to_string(),
addr: format!("{}:{}", addr, cfg.advertised_port()),
});
}
}
@@ -1593,6 +1732,13 @@ impl Node {
self.register_identity(peer_node_addr, peer_identity.pubkey_full());
let transport_id = self.allocate_transport_id();
// Adopted ephemeral UDP transports use UdpConfig::default() when the
// bootstrap runtime doesn't pass an override. Default MTU resolves to
// 1280 (IPv6 minimum), which is the only value guaranteed to survive
// arbitrary NAT-traversal middlebox paths. Inheriting from the named
// [transports.udp] config (Option 3 in ISSUE-2026-0013) would track
// operator config more closely but risks regressions on hostile paths;
// accepted as-is until a concrete use case justifies the change.
let mut transport = crate::transport::udp::UdpTransport::new(
transport_id,
traversal.transport_name.clone(),
+36 -10
View File
@@ -441,6 +441,15 @@ pub struct Node {
/// Optional Nostr/STUN overlay discovery coordinator for `udp:nat` peers.
nostr_discovery: Option<Arc<crate::discovery::nostr::NostrDiscovery>>,
/// Wall-clock ms when Nostr discovery successfully started, used to
/// schedule the one-shot startup advert sweep after a settle delay.
/// `None` until discovery comes up; remains `None` if discovery is
/// disabled or failed to start.
nostr_discovery_started_at_ms: Option<u64>,
/// Whether the one-shot startup advert sweep has run. Set to true
/// after the first sweep fires (under `policy: open`); thereafter
/// only the per-tick `queue_open_discovery_retries` continues.
startup_open_discovery_sweep_done: bool,
/// Per-peer UDP transports adopted from NAT traversal handoff.
bootstrap_transports: HashSet<TransportId>,
@@ -608,6 +617,8 @@ impl Node {
pending_connects: Vec::new(),
retry_pending: HashMap::new(),
nostr_discovery: None,
nostr_discovery_started_at_ms: None,
startup_open_discovery_sweep_done: false,
bootstrap_transports: HashSet::new(),
last_parent_reeval: None,
last_congestion_log: None,
@@ -737,6 +748,8 @@ impl Node {
pending_connects: Vec::new(),
retry_pending: HashMap::new(),
nostr_discovery: None,
nostr_discovery_started_at_ms: None,
startup_open_discovery_sweep_done: false,
bootstrap_transports: HashSet::new(),
last_parent_reeval: None,
last_congestion_log: None,
@@ -1023,18 +1036,31 @@ impl Node {
crate::upper::icmp::effective_ipv6_mtu(self.transport_mtu())
}
/// Get the transport MTU for a specific transport.
/// Get the transport MTU governing the global TUN-boundary MSS clamp.
///
/// When called without a specific transport context, returns the MTU
/// of the first operational transport, or 1280 (IPv6 minimum) as
/// fallback. This is used for initial TUN configuration where a
/// specific transport isn't yet known.
/// Returns the **minimum** MTU across all operational transports, or
/// 1280 (IPv6 minimum) as fallback. Used for initial TUN configuration
/// where a specific egress transport isn't yet known: the resulting
/// `effective_ipv6_mtu` (transport_mtu - 77) and `max_mss`
/// (effective_mtu - 60) form a conservative ceiling that fits ANY
/// configured-transport's egress, eliminating PMTU-D black holes that
/// would otherwise occur when a flow's actual egress is smaller than
/// the clamp ceiling assumed at TUN init.
///
/// Returning the smallest (rather than the first-iterated, which used
/// to vary across HashMap iteration order + async-startup race) makes
/// the clamp deterministic across daemon restarts.
///
/// See `ISSUE-2026-0011` for the empirical investigation.
pub fn transport_mtu(&self) -> u16 {
// Prefer the MTU from the first operational transport
for handle in self.transports.values() {
if handle.is_operational() {
return handle.mtu();
}
let min_operational = self
.transports
.values()
.filter(|h| h.is_operational())
.map(|h| h.mtu())
.min();
if let Some(mtu) = min_operational {
return mtu;
}
// Fallback to config: try UDP first, then Ethernet
if let Some((_, cfg)) = self.config.transports.udp.iter().next() {
+79
View File
@@ -954,3 +954,82 @@ fn test_promote_clears_retry_pending() {
"retry_pending should be cleared on successful promotion"
);
}
// ============================================================================
// transport_mtu() — ISSUE-2026-0011 regression coverage
// ============================================================================
/// Helper: spawn a UdpTransport with the given mtu, started and operational.
async fn make_udp_transport_with_mtu(id: u32, mtu: u16) -> TransportHandle {
let (packet_tx, _packet_rx) = packet_channel(64);
let transport_id = TransportId::new(id);
let mut udp = UdpTransport::new(
transport_id,
Some(format!("udp{}", id)),
crate::config::UdpConfig {
bind_addr: Some("127.0.0.1:0".to_string()),
mtu: Some(mtu),
..Default::default()
},
packet_tx,
);
udp.start_async().await.unwrap();
TransportHandle::Udp(udp)
}
#[tokio::test]
async fn test_transport_mtu_returns_min_across_operational() {
// Multiple operational transports with varied MTUs. The picker must
// return the smallest, deterministically, regardless of HashMap
// iteration order. This is the core ISSUE-2026-0011 regression test.
let mut node = make_node();
let (packet_tx, packet_rx) = packet_channel(64);
node.packet_tx = Some(packet_tx);
node.packet_rx = Some(packet_rx);
let udp1 = make_udp_transport_with_mtu(1, 1497).await;
let udp2 = make_udp_transport_with_mtu(2, 1280).await;
let udp3 = make_udp_transport_with_mtu(3, 1400).await;
node.transports.insert(TransportId::new(1), udp1);
node.transports.insert(TransportId::new(2), udp2);
node.transports.insert(TransportId::new(3), udp3);
// Expect the smallest (UDP-1280), not whichever HashMap iterates first.
assert_eq!(node.transport_mtu(), 1280);
// effective_ipv6_mtu = 1280 - 77 = 1203, max_mss = 1203 - 60 = 1143
// (verifies the downstream clamp value).
assert_eq!(node.effective_ipv6_mtu(), 1203);
for transport in node.transports.values_mut() {
transport.stop().await.ok();
}
}
#[tokio::test]
async fn test_transport_mtu_fallback_when_no_operational_transports() {
// No transports configured at all → falls back to 1280 (IPv6 minimum).
let node = make_node();
assert_eq!(node.transport_mtu(), 1280);
}
#[tokio::test]
async fn test_transport_mtu_min_with_single_operational() {
// Single transport: trivially returns its MTU. Pins the picker doesn't
// accidentally drop down to a smaller fallback when one transport is
// operational.
let mut node = make_node();
let (packet_tx, packet_rx) = packet_channel(64);
node.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);
assert_eq!(node.transport_mtu(), 1452);
for transport in node.transports.values_mut() {
transport.stop().await.ok();
}
}