Merge branch 'maint'

Carries the five security commits up. Five of the nine behaviours merged as
text; four had to be re-implemented against this line refactored structure,
and two files this line had deleted needed their content relocated.

Relocated: the path-MTU source reset into the protocol MMP module, adapted
because this line stores the first-increase timestamp as milliseconds
rather than an instant; the lookup-carrier deadline into the renamed lookup
handler, where it routes through the write action rather than touching the
map directly; and the flap-dampening test into the spanning-tree limits
tests, rewritten against this line injected clock so it needs no real
sleep. The maint copies of both deleted files are removed.

The lookup write action gains a timestamp field so both stores it feeds are
stamped at the same instant, which is what the deadline reasoning depends
on. The alternative was a second clock read in the shell.

Config keys and the new validator error strings are expressed against this
line rendezvous naming rather than the maintenance line discovery naming.

Three defects survived the automatic merge and were caught by the compiler:
a serde default naming the old config struct, an import of a constant from
the old module path, and two tests missing an import this line requires per
test. The automatic merge also duplicated a test and a relocated block that
already existed here; both were dropped.

Green: fmt, build, clippy and test --lib, 1861 passed.
This commit is contained in:
Johnathan Corgan
2026-08-15 07:56:54 +00:00
35 changed files with 2579 additions and 120 deletions
+112
View File
@@ -102,6 +102,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
reporting channel at all, so someone with a finding had to guess at an
address or open a public issue.
- `node.rate_limit.session_setup_burst` (64) and
`node.rate_limit.session_setup_rate` (16.0), the parameters of the new
per-link-peer session-setup limiter. Setup messages naming a peer this node
is already established with are metered on a second per-link bucket derived
from `node.limits.max_peers`, `node.rekey.after_secs` and
`node.rate_limit.handshake_max_resends`, so raising the peer limit sizes it
automatically. A zero burst or a non-positive rate is rejected at config
validation rather than silently refusing every session.
- `node.rate_limit.established_handshake_burst` and
`node.rate_limit.established_handshake_rate`, the parameters of the new
established-link msg1 token bucket. Both are optional; omitting them (the
@@ -123,6 +132,13 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
refused together with `--no-build`, which would stamp the marker onto binaries
the features never reached.
- `node.rendezvous.nostr.max_concurrent_offers_per_npub`, defaulting to 4, which
bounds how many inbound traversal offers one sender npub may have in flight
at once. It sits inside `max_concurrent_incoming_offers`, which remains the
outer bound, so a value above that is inert; zero is rejected at config
validation, since it refuses every inbound offer rather than disabling the
limit. Existing configurations parse unchanged, the key being optional.
### Changed
- Node health is determined at start completion instead of unconditionally
@@ -175,6 +191,42 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
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
as a log field, so one sender could hold every slot and deny traversal
onboarding to every other peer for as long as it kept offering. Admission now
takes a per-npub permit and a global permit together. A sender over its own
allowance is refused at debug rather than warn, because the party tripping it
is by definition sending faster than the node wants and a record per
rejection would turn the spam into log volume; the global bound being reached
keeps its warn, which is the operator's signal that the node is genuinely
saturated. **This does not make the pool inexhaustible.** Nostr identities
cost nothing to generate and the signal subscription carries no author
restriction, so an attacker running four throwaway npubs still saturates the
shipped 16-slot pool at an unchanged total offer rate. What the change buys
is that one identity can no longer do it alone, and that the two refusals are
distinguishable in the log. The permit is still held across the whole
attempt; that duration remains inferred from the attempt timeout rather than
measured.
- Config validation now rejects a `node.rendezvous.nostr.signal_ttl_secs` that
is too large for the configured `replay_window_secs`. A traversal signal is
acceptable over its TTL plus 60s of clock-skew grace on each side, and that
span has to stay strictly inside the replay window, or a session id evicted
from the replay cache on expiry is still fresh enough to be accepted a second
time. The relation was documented but unenforced, so raising the TTL past
180s silently voided it. The bound is derived from the skew constant rather
than restated, and is checked whether or not nostr discovery is enabled, for
the same reason the rekey rules are. The shipped defaults (120s against 300s)
are unaffected, but a configuration that had widened the TTL or narrowed the
replay window now fails to load, with an error naming the concrete floor for
`replay_window_secs`. The NAT lab's config generator was one such
configuration and its generated `replay_window_secs` moves from 60 to 180.
Note that this covers eviction on expiry only: `seen_sessions_max_entries`
remains a separate capacity-eviction route that no config relation bounds.
- Config validation now rejects two `node.rekey` settings that appear to
disable the trigger and in fact fire it continuously. `after_messages` of
zero makes the message-count arm true on every poll, because the trigger
@@ -252,6 +304,66 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Fixed
- Inbound session-setup messages are now rate limited, keyed on the
authenticated link peer the datagram arrived over. The setup path allocated
a session entry and sent a routed SessionAck for every well-formed message
naming an address it had no entry for, and that address is an envelope field
the sender picks, so one neighbour could grow the session table at whatever
rate it could transmit and buy an ack per entry to a destination of its
choosing. The limiter sits ahead of every send and both handshake
constructions in the handler, so a refused message emits nothing and costs
no cryptography. The key is the link peer rather than the claimed source
address, which is what makes it a limit at all: keying on the source would
hand a single sender a fresh full bucket per forged message.
Two consequences worth stating rather than discovering. The limiter bounds
each neighbour's contribution and makes a flood attributable; it does not
give the node an absolute ceiling, which stays at roughly
`peers * rate * handshake_timeout_secs`. And a legitimate peer reaching this
node over the *same* link as an attacker shares that attacker's bucket, so
establishment behind a flooded neighbour is refused until it refills. Rekey
and restart traffic is deliberately not subject to that: setup messages
naming an already-established peer draw on a separate per-link bucket,
because suppressed key rotation is silent — nothing errors and no session
drops — and would have shown up only as a flat `rekey_armed`.
- A forged SessionAck no longer destroys an in-flight session initiation. The
handler removed the session entry to take ownership of the handshake state
and, when the XK msg2 read failed, returned without putting it back. Nothing
in that message is authenticated — the only thing tying it to the initiation
is the datagram's source address, which the sender chooses — so any node able
to reach the victim could cancel any initiation with 57 bytes of the right
length, and hold establishment down by repeating it. The entry is now kept.
Keeping it is not enough on its own, and the second half is the part worth
naming: the msg2 read mixes the sender's ephemeral into the symmetric state
before it authenticates anything, so an entry put back as the failed read
left it holds a handshake that can never read the genuine msg2, which trades
a one-round-trip denial for one lasting the full handshake timeout. The read
is therefore rolled back to its pre-read state before the entry goes back.
The three later failure paths in the same handler still drop the entry: each
is downstream of a msg2 that authenticated, so it is a local failure rather
than a possible forgery. The entry's activity stamp is deliberately not
refreshed on the failure path, so a spray cannot hold a dead initiation past
its original sweep deadline, and a new `ack_handshake_failed` counter makes
the refusals visible at the default log level. Only the XK handshake on this
branch is covered; the additional drop sites in the XX handshake on the
development branch are not.
- An unauthenticated session msg3 no longer discards a completed key epoch.
The four failure paths in the responder-side rekey arm of the msg3 handler
abandoned the whole rekey, which nulls a `pending` session sitting beside the
handshake, when only the handshake had failed. That `pending` session is the
epoch the real peer may already have cut over to, so discarding it kills the
reverse direction until the session idles out. Two unauthenticated messages
reached it: a forged setup message arms a handshake beside a completed rekey
once that rekey has waited a full idle timeout for a peer that never appeared
on the new epoch, and any garbage msg3 of the right length then finishes the
job. All four paths now abandon only the handshake. The neighbouring comment
in the setup handler, which claimed no path in that handler discards the
pending keys, has been corrected rather than left to be trusted: the
dual-initiation arm still calls the destructive form on an unauthenticated
msg1, and that is now named at the site as the one path that does.
- The packaged macOS daemon now recreates and binds its control socket at
`/var/run/fips/control.sock` instead of falling through to the shared
`/tmp/fips-control.sock` path after boot. The previous fallback also made the
+15 -5
View File
@@ -212,11 +212,20 @@ NIP-17 DM relay list (kind 10050), and falls back to `dm_relays` if
the inbox-relays fetch fails. Each side also publishes its own inbox
relay list on startup so dialers can discover it.
On the receiving side, an inbound semaphore bounds concurrent offer
processing at `max_concurrent_incoming_offers`. When the semaphore is
full, the offer is dropped with a warn log; this is the primary guard
against offer-spam from a misbehaving or compromised relay. A
`sessionId` replay cache (bounded by `seen_sessions_max_entries`, with
On the receiving side, admission is a pair of bounds taken together: a
per-sender allowance of `max_concurrent_offers_per_npub`, keyed on the
npub that signed the gift wrap, nested inside a global
`max_concurrent_incoming_offers`. A sender over its own allowance is
refused at debug, since by definition it is sending faster than the node
wants and a record per rejection would turn the spam into log volume; the
global bound being reached is the operator-visible warn, because that one
says the node is genuinely saturated. Together they keep one identity
from holding the whole pool. They do not make the pool inexhaustible:
nostr identities are free to generate, so an attacker running
`ceil(max_concurrent_incoming_offers / max_concurrent_offers_per_npub)`
throwaway npubs still saturates it at the same total offer rate. Raising
the attacker's cost beyond keypairs would mean pricing the offer itself.
A `sessionId` replay cache (bounded by `seen_sessions_max_entries`, with
entries valid for `replay_window_secs`) rejects duplicates.
The responder runs its own STUN query and replies with a
@@ -306,6 +315,7 @@ machinery:
| Mechanism | Default | What it prevents | Behavior at limit |
| --- | --- | --- | --- |
| Offer semaphore (`max_concurrent_incoming_offers`) | 16 | CPU and memory exhaustion from offer spam on DM relays. | Warn log, offer dropped. |
| Per-npub offer allowance (`max_concurrent_offers_per_npub`) | 4 | One sender identity holding every offer slot and denying traversal onboarding to everyone else. Does not prevent the same denial from several throwaway npubs. | Debug log, offer dropped. |
| Advert cache (`advert_cache_max_entries`) | 2048 | Memory growth from ambient advert traffic under `policy: open`. | LRU-by-expiry eviction. |
| Seen-sessions (`seen_sessions_max_entries`) | 2048 | Replay of stale `sessionId` values. | Oldest entry evicted. |
| Signal TTL (`signal_ttl_secs`) | 120 s | Indefinite in-flight offers on relays. | Expired offers rejected at validation. |
+24
View File
@@ -136,6 +136,8 @@ Handshake rate limiting protects against DoS on the Noise IK handshake path.
| `node.rate_limit.handshake_max_resends` | u32 | `5` | Max resends per handshake attempt |
| `node.rate_limit.established_handshake_burst` | u32 | derived | Burst capacity of the established-link bucket. Derived default is `node.limits.max_peers` (128) |
| `node.rate_limit.established_handshake_rate` | f64 | derived | Refill rate of that bucket. Derived default is `(max_peers / max(node.rekey.after_secs, 1)) * (1 + handshake_max_resends)`, floored at 1.0/s — 6.4/s at shipped defaults |
| `node.rate_limit.session_setup_burst` | u32 | `64` | Per-link-peer burst for inbound session-setup messages that would open a new session |
| `node.rate_limit.session_setup_rate` | f64 | `16.0` | Per-link-peer refill rate for those messages, in tokens per second |
Msg1 whose source matches an established link (rekey and restart
maintenance traffic) draws on a second bucket rather than competing with
@@ -150,6 +152,25 @@ burst and 16.4/s at shipped defaults, of which the established half is
reachable only by a source that already matches a live link. Size against
the sum when budgeting handshake crypto load for a host.
The `session_setup_*` pair is a separate limiter on the session layer, not
the link layer. It is keyed on the authenticated link peer a session datagram
arrived over, so each neighbour gets its own budget and a flood is
attributable. Setup messages naming a peer this node is already established
with (inbound rekey and restart traffic) draw on a second per-link bucket
derived from `max_peers`, `node.rekey.after_secs` and `handshake_max_resends`,
exactly as `established_handshake_*` is, so a stranger flood cannot suppress
rekey traffic sharing the link.
At the defaults one neighbour can force at most
`session_setup_rate * handshake_timeout_secs` half-open entries (480) and
`session_setup_rate * (1 + handshake_max_resends)` acks per second (96). The
node-wide ceiling is still that times the peer count, since the limiter bounds
each neighbour rather than the aggregate. A legitimate peer whose traffic
reaches this node over the *same* link as an attacker's shares that
attacker's stranger bucket, so establishment behind a flooded neighbour is
refused until the bucket refills; the initiator's own resend schedule (1s, 2s,
4s, 8s, 16s) covers a short drain.
### Retry / Backoff (`node.retry.*`)
Connection retry with exponential backoff.
@@ -207,6 +228,7 @@ inert otherwise.
| `node.discovery.nostr.policy` | string | `"configured_only"` | Advert discovery policy: `disabled`, `configured_only`, `open` |
| `node.discovery.nostr.open_discovery_max_pending` | usize | `64` | Max open-discovery peers queued in outbound retry/connection state at once |
| `node.discovery.nostr.max_concurrent_incoming_offers` | usize | `16` | Max concurrent inbound traversal offers processed at once (rate limit against offer spam) |
| `node.discovery.nostr.max_concurrent_offers_per_npub` | usize | `4` | Max concurrent inbound traversal offers accepted from any one sender npub, so a single identity cannot hold the whole pool. Sits inside `max_concurrent_incoming_offers`, which stays the outer bound; a larger value is inert. Zero is rejected, since it refuses every inbound offer rather than disabling the limit |
| `node.discovery.nostr.advert_cache_max_entries` | usize | `2048` | Max cached overlay adverts retained from relay traffic |
| `node.discovery.nostr.seen_sessions_max_entries` | usize | `2048` | Max seen-session IDs retained for replay detection |
| `node.discovery.nostr.advertise` | bool | `true` | Publish local endpoint adverts |
@@ -931,6 +953,8 @@ node:
handshake_resend_interval_ms: 1000
handshake_resend_backoff: 2.0
handshake_max_resends: 5
session_setup_burst: 64
session_setup_rate: 16.0
retry:
max_retries: 5
base_interval_secs: 5
+130
View File
@@ -26,6 +26,7 @@ mod peer;
mod transport;
use crate::node::REKEY_JITTER_SECS;
use crate::nostr::FRESHNESS_SKEW_TOLERANCE_MS;
use crate::upper::config::{DnsConfig, TunConfig};
use crate::{Identity, IdentityError};
use serde::{Deserialize, Serialize};
@@ -1113,6 +1114,67 @@ impl Config {
)));
}
// The per-link session-setup bucket. The same trap as above in a
// non-optional field: zero does not disable the limiter, it refuses
// every inbound setup message and so refuses every session.
if rl.session_setup_burst == 0 {
return Err(ConfigError::Validation(
"`node.rate_limit.session_setup_burst` is 0, which refuses every inbound SessionSetup rather than disabling the limit. \
Set a positive burst; a very large value effectively disables it."
.to_string(),
));
}
let setup_rate = rl.session_setup_rate;
if !(setup_rate.is_finite() && setup_rate > 0.0) {
return Err(ConfigError::Validation(format!(
"`node.rate_limit.session_setup_rate` is {setup_rate}, but must be a finite value greater than 0; \
a non-positive or non-finite refill rate never replenishes a link's setup bucket, so that link stops establishing sessions once its initial burst is spent."
)));
}
// The freshness window backstops session-id replay protection: an
// offer evicted from the replay cache must already be too old to pass
// the freshness check, or it can be accepted a second time. A signal
// is acceptable over `signal_ttl_secs` plus the skew tolerance on each
// side, so that span has to stay strictly inside `replay_window_secs`.
// Checked regardless of `nostr.enabled`, for the reason the rekey
// block above gives: enabling the feature later must not surface a
// config error at a surprising moment.
let skew_secs = FRESHNESS_SKEW_TOLERANCE_MS / 1000;
let freshness_window_secs = nostr.signal_ttl_secs.saturating_add(2 * skew_secs);
if freshness_window_secs >= nostr.replay_window_secs {
return Err(ConfigError::Validation(format!(
"`node.rendezvous.nostr.signal_ttl_secs` is {}, which with {skew_secs}s of clock-skew grace on each side makes a traversal signal acceptable over a {freshness_window_secs}s window, \
but `node.rendezvous.nostr.replay_window_secs` is {}. \
The freshness window must be strictly narrower than the replay window, or a session id evicted from the replay cache is still fresh enough to be accepted a second time. \
Raise `replay_window_secs` above {freshness_window_secs}, or lower `signal_ttl_secs`.",
nostr.signal_ttl_secs, nostr.replay_window_secs
)));
}
// Zero here is the same trap as `established_handshake_burst`: it
// reads as "no limit" and in fact refuses every inbound offer. The
// upper bound exists because the per-npub semaphore is built lazily
// inside the intake path rather than at startup, so an oversized
// value would panic there instead of failing loudly at load.
if nostr.max_concurrent_offers_per_npub == 0 {
return Err(ConfigError::Validation(
"`node.rendezvous.nostr.max_concurrent_offers_per_npub` is 0, which refuses every inbound traversal offer rather than disabling the per-sender limit. \
Omit the key for the default, or set a positive allowance; `max_concurrent_incoming_offers` remains the outer bound."
.to_string(),
));
}
if nostr.max_concurrent_offers_per_npub > tokio::sync::Semaphore::MAX_PERMITS {
return Err(ConfigError::Validation(format!(
"`node.rendezvous.nostr.max_concurrent_offers_per_npub` is {}, which exceeds the maximum {} permits a semaphore can hold. \
Use a value at or below `max_concurrent_incoming_offers`, which is the outer bound anything larger is inert against.",
nostr.max_concurrent_offers_per_npub,
tokio::sync::Semaphore::MAX_PERMITS
)));
}
Ok(())
}
@@ -2443,6 +2505,74 @@ node:
assert!(err.to_string().contains("after_secs"));
}
#[test]
fn test_validate_signal_ttl_at_or_above_the_replay_window_margin_rejected() {
// 180 is the boundary: 180 + 2 * 60 = 300, which is not strictly less
// than the default 300s replay window.
for signal_ttl_secs in [180, 181, 3600, u64::MAX] {
let mut config = Config::default();
config.node.rendezvous.nostr.signal_ttl_secs = signal_ttl_secs;
match config.validate() {
Err(e) => {
let msg = e.to_string();
assert!(msg.contains("signal_ttl_secs"), "got: {msg}");
assert!(msg.contains("replay_window_secs"), "got: {msg}");
}
Ok(()) => panic!("signal_ttl_secs = {signal_ttl_secs} should be rejected"),
}
}
}
#[test]
fn test_validate_signal_ttl_just_inside_the_replay_window_margin_accepted() {
let mut config = Config::default();
config.node.rendezvous.nostr.signal_ttl_secs = 179;
config
.validate()
.expect("179 + 2 * 60 = 299 leaves the freshness window inside the 300s replay window");
}
#[test]
fn test_validate_per_npub_offer_allowance_of_zero_rejected() {
let mut config = Config::default();
config.node.rendezvous.nostr.max_concurrent_offers_per_npub = 0;
let err = config.validate().expect_err("validation should fail");
assert!(
err.to_string().contains("max_concurrent_offers_per_npub"),
"got: {err}"
);
}
#[test]
fn test_validate_per_npub_offer_allowance_of_one_accepted() {
let mut config = Config::default();
config.node.rendezvous.nostr.max_concurrent_offers_per_npub = 1;
config
.validate()
.expect("an allowance of one offer per sender is restrictive but well defined");
}
#[test]
fn test_validate_shipped_defaults_satisfy_the_freshness_invariant() {
Config::default()
.validate()
.expect("the shipped defaults must satisfy every validation rule");
// Stated against the constant rather than a literal, so this reds if
// anyone changes FRESHNESS_SKEW_TOLERANCE_MS or either default without
// re-checking the relation they jointly have to satisfy.
let defaults = Config::default();
let nostr = &defaults.node.rendezvous.nostr;
assert!(
nostr.signal_ttl_secs + 2 * (FRESHNESS_SKEW_TOLERANCE_MS / 1000)
< nostr.replay_window_secs
);
}
#[test]
fn test_outbound_only_forces_ephemeral_bind() {
let cfg = UdpConfig {
+49
View File
@@ -91,6 +91,27 @@ pub struct RateLimitConfig {
/// `node.rekey.after_secs` and `handshake_max_resends`.
#[serde(default)]
pub established_handshake_rate: Option<f64>,
/// Per-link-peer burst capacity for inbound FSP SessionSetup messages
/// that would open a new session (`node.rate_limit.session_setup_burst`).
///
/// 64 absorbs a legitimate reconnect burst arriving behind one
/// neighbour. It bounds nothing on its own; `session_setup_rate` is what
/// bounds the sustained cost.
#[serde(default = "RateLimitConfig::default_session_setup_burst")]
pub session_setup_burst: u32,
/// Per-link-peer refill rate for those messages, in tokens per second
/// (`node.rate_limit.session_setup_rate`).
///
/// 16/s caps one neighbour's forced half-open occupancy at
/// `rate * handshake_timeout_secs` (480 entries at defaults) and its ack
/// amplification at `rate * (1 + handshake_max_resends)` (96 acks/s).
///
/// Setup messages naming a peer this node is already established with
/// are metered on a separate per-link bucket, derived from
/// `node.limits.max_peers` exactly as `established_handshake_*` is, so a
/// stranger flood cannot suppress rekey traffic sharing the link.
#[serde(default = "RateLimitConfig::default_session_setup_rate")]
pub session_setup_rate: f64,
}
impl Default for RateLimitConfig {
@@ -104,6 +125,8 @@ impl Default for RateLimitConfig {
handshake_max_resends: 5,
established_handshake_burst: None,
established_handshake_rate: None,
session_setup_burst: 64,
session_setup_rate: 16.0,
}
}
}
@@ -127,6 +150,12 @@ impl RateLimitConfig {
fn default_handshake_max_resends() -> u32 {
5
}
fn default_session_setup_burst() -> u32 {
64
}
fn default_session_setup_rate() -> f64 {
16.0
}
}
/// Retry/backoff configuration (`node.retry.*`).
@@ -374,6 +403,11 @@ pub struct NostrRendezvousConfig {
/// Acts as a rate limit against offer spam from relays.
#[serde(default = "NostrRendezvousConfig::default_max_concurrent_incoming_offers")]
pub max_concurrent_incoming_offers: usize,
/// Max concurrent inbound traversal offers accepted from any one sender
/// npub. Sits inside `max_concurrent_incoming_offers`, which remains the
/// outer bound.
#[serde(default = "NostrRendezvousConfig::default_max_concurrent_offers_per_npub")]
pub max_concurrent_offers_per_npub: usize,
/// Max cached overlay adverts retained from relay traffic.
/// Bounds memory under ambient advert volume.
#[serde(default = "NostrRendezvousConfig::default_advert_cache_max_entries")]
@@ -462,6 +496,7 @@ impl Default for NostrRendezvousConfig {
policy: NostrRendezvousPolicy::default(),
open_discovery_max_pending: Self::default_open_discovery_max_pending(),
max_concurrent_incoming_offers: Self::default_max_concurrent_incoming_offers(),
max_concurrent_offers_per_npub: Self::default_max_concurrent_offers_per_npub(),
advert_cache_max_entries: Self::default_advert_cache_max_entries(),
seen_sessions_max_entries: Self::default_seen_sessions_max_entries(),
attempt_timeout_secs: Self::default_attempt_timeout_secs(),
@@ -527,6 +562,20 @@ impl NostrRendezvousConfig {
16
}
/// Four, derived rather than picked. The initiator side already admits at
/// most one in-flight traversal per peer npub, so one concurrent offer per
/// peer is the honest steady state. An offer is published to and consumed
/// from the whole DM relay set, three URLs by default, and whether the
/// notification stream deduplicates one event delivered by three relays is
/// not established here — if it does not, one honest offer can present as
/// three near-simultaneous admissions before the replay check rejects the
/// duplicates. Four is that worst-case fan-out plus one, so a retry
/// overlapping a still-timing-out attempt is still admitted, and it is a
/// quarter of the default global bound.
fn default_max_concurrent_offers_per_npub() -> usize {
4
}
fn default_advert_cache_max_entries() -> usize {
2048
}
+3
View File
@@ -57,6 +57,7 @@ pub(crate) enum Step {
ResendPendingSessionHandshakes,
ResendPendingSessionMsg3,
PurgeIdleSessions,
PurgeExpiredPathMtu,
ProcessPendingRetries,
CheckTreeState,
CheckBloomState,
@@ -94,6 +95,7 @@ pub(crate) const STEPS: [Step; N_STEPS] = [
Step::ResendPendingSessionHandshakes,
Step::ResendPendingSessionMsg3,
Step::PurgeIdleSessions,
Step::PurgeExpiredPathMtu,
Step::ProcessPendingRetries,
Step::CheckTreeState,
Step::CheckBloomState,
@@ -126,6 +128,7 @@ impl Step {
Step::ResendPendingSessionHandshakes => "resend_pending_session_handshakes",
Step::ResendPendingSessionMsg3 => "resend_pending_session_msg3",
Step::PurgeIdleSessions => "purge_idle_sessions",
Step::PurgeExpiredPathMtu => "purge_expired_path_mtu",
Step::ProcessPendingRetries => "process_pending_retries",
Step::CheckTreeState => "check_tree_state",
Step::CheckBloomState => "check_bloom_state",
+12 -6
View File
@@ -27,7 +27,7 @@ impl Node {
/// has already had its msg_type byte stripped by dispatch.
pub(in crate::node) async fn handle_session_datagram(
&mut self,
_from: &NodeAddr,
from: &NodeAddr,
payload: &[u8],
incoming_ce: bool,
) {
@@ -51,7 +51,7 @@ impl Node {
// coords a peer put on the wire are equally valid whichever way those
// go, and the only arrivals this newly warms from are those with an
// exhausted TTL, whose every insert is already achievable at TTL 1.
self.try_warm_coord_cache_ref(&datagram_ref);
self.try_warm_coord_cache_ref(&datagram_ref, payload.len());
// Pre-resolve the next hop only for datagrams the core can actually
// forward: not locally destined, and carrying a TTL that survives the
@@ -107,6 +107,7 @@ impl Node {
self.metrics().forwarding.record_delivered(payload.len());
self.handle_session_payload(
&datagram_ref.src_addr,
from,
datagram_ref.payload,
datagram_ref.path_mtu,
incoming_ce,
@@ -219,7 +220,13 @@ impl Node {
///
/// Decode failures are logged and silently ignored — they don't block
/// forwarding.
fn try_warm_coord_cache_ref(&mut self, datagram: &SessionDatagramRef<'_>) {
///
/// `outer_len` is the length of the msg_type-stripped `SessionDatagram`
/// buffer this view was decoded from. It is carried in rather than
/// reconstructed from the header size so the malformed-frame byte counter
/// measures the same population as its siblings — which are charged the
/// outer slice — instead of the inner FSP payload.
fn try_warm_coord_cache_ref(&mut self, datagram: &SessionDatagramRef<'_>, outer_len: usize) {
let prefix = match FspCommonPrefix::parse(datagram.payload) {
Some(p) => p,
None => return,
@@ -276,11 +283,10 @@ impl Node {
// drill-down that separates a short frame from a bad version
// or a U-flagged one. The level stays at debug: any peer past
// the handshake can drive this at line rate.
self.metrics()
.forwarding
.record_warm_malformed(datagram.payload.len());
self.metrics().forwarding.record_warm_malformed(outer_len);
debug!(
len = datagram.payload.len(),
outer_len,
version = prefix.version,
flags = prefix.flags,
"Not a well-formed encrypted FSP message; not warming coords"
+6 -1
View File
@@ -355,7 +355,10 @@ impl Node {
// distinct resources; the `path_mtu_lookup` cache and the
// `nostr_rendezvous` subsystem are deliberately excluded
// from `Reloadable` since neither reloads from a backing
// file (see `node::reloadable`).
// file (see `node::reloadable`). The `path_mtu_lookup`
// cache is nevertheless swept on this tick, by
// `purge_expired_path_mtu` below: that is expiry, not
// reload.
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ReloadHostMap,
self.reload_host_map().await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PollPendingConnects,
@@ -374,6 +377,8 @@ impl Node {
self.resend_pending_session_msg3(now_ms).await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PurgeIdleSessions,
self.purge_idle_sessions(now_ms));
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PurgeExpiredPathMtu,
self.purge_expired_path_mtu(now_ms));
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ProcessPendingRetries,
self.process_pending_retries(now_ms).await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::CheckTreeState,
+28 -7
View File
@@ -296,7 +296,11 @@ impl Node {
.insert_with_path_mtu(target, coords, now_ms, path_mtu);
}
}
LookupAction::WritePathMtu { target, path_mtu } => {
LookupAction::WritePathMtu {
target,
now_ms,
path_mtu,
} => {
// Refused as absent on the CacheCoords arm the core always
// pairs with this one, so there is nothing to mirror; the
// warning and the counter are emitted there, once.
@@ -308,22 +312,37 @@ impl Node {
let fips_addr = crate::FipsAddress::from_node_addr(&target);
match self.path_mtu_lookup.write() {
Ok(mut map) => match map.get(&fips_addr).copied() {
Some(existing) if existing <= path_mtu => {
Some(existing) if existing.mtu <= path_mtu => {
// Keep the tighter learned value; never loosen
// the clamp. A reactive MtuExceeded or
// PathMtuNotification tighten takes precedence
// over a looser discovery estimate
// (cross-carrier keep-tighter).
//
// This arm deliberately leaves `learned_ms`
// alone. That is what bounds a replayed
// response: the replay of a value already
// stored takes this arm, so the entry still
// expires at first-write plus the TTL rather
// than being pushed out again on every
// injection. Refreshing the stamp here would
// read as a tidy-up and would silently restore
// indefinite pinning.
debug!(
target = %self.peer_display_name(&target),
fips_addr = %fips_addr,
path_mtu = path_mtu,
existing = existing,
existing = existing.mtu,
"LookupResponse: keeping tighter existing path_mtu_lookup value"
);
}
other => {
map.insert(fips_addr, path_mtu);
// The one carrier with no release path, so this
// is the one write that carries a deadline.
map.insert(
fips_addr,
crate::upper::tun::PathMtuEntry::learned(path_mtu, now_ms),
);
debug!(
target = %self.peer_display_name(&target),
fips_addr = %fips_addr,
@@ -753,18 +772,20 @@ impl Node {
return;
};
match map.get(&fips_addr).copied() {
Some(existing) if existing <= link_mtu => {
Some(existing) if existing.mtu <= link_mtu => {
// Keep the tighter learned value; never loosen the clamp.
debug!(
peer = %self.peer_display_name(peer_addr),
fips_addr = %fips_addr,
link_mtu = link_mtu,
existing = existing,
existing = existing.mtu,
"seed_path_mtu_for_link_peer: keeping tighter existing value"
);
}
other => {
map.insert(fips_addr, link_mtu);
// Held, not expiring: this describes a link this node can see
// for itself, and it is released when the link goes.
map.insert(fips_addr, crate::upper::tun::PathMtuEntry::held(link_mtu));
debug!(
peer = %self.peer_display_name(peer_addr),
fips_addr = %fips_addr,
+164 -18
View File
@@ -7,6 +7,7 @@
use crate::NodeAddr;
use crate::node::handlers::mmp::format_throughput;
use crate::node::rate_limit::Msg1Class;
use crate::node::reject::{RejectReason, SessionReject};
use crate::node::session::{EndToEndState, EpochSlot, SessionEntry};
use crate::node::{Node, NodeError};
@@ -100,9 +101,15 @@ impl Node {
/// - Phase 0x3 → SessionMsg3 (XK handshake msg3)
/// - Phase 0x0 + U flag → plaintext error signal (CoordsRequired/PathBroken)
/// - Phase 0x0 + !U → encrypted session message (data, reports, etc.)
///
/// `src_addr` is the datagram's claimed source, an envelope field the
/// sender chooses. `link_peer` is the authenticated FMP peer the datagram
/// arrived over, and is the only identity on this path worth keying a
/// limiter on.
pub(in crate::node) async fn handle_session_payload(
&mut self,
src_addr: &NodeAddr,
link_peer: &NodeAddr,
payload: &[u8],
path_mtu: u16,
ce_flag: bool,
@@ -122,7 +129,7 @@ impl Node {
match prefix.phase {
FSP_PHASE_MSG1 => {
self.handle_session_setup(src_addr, inner).await;
self.handle_session_setup(src_addr, link_peer, inner).await;
}
FSP_PHASE_MSG2 => {
self.handle_session_ack(src_addr, inner).await;
@@ -470,7 +477,12 @@ impl Node {
/// The remote node wants to establish an end-to-end session with us.
/// We create an XK responder handshake, process msg1, send SessionAck with msg2,
/// and transition to AwaitingMsg3.
async fn handle_session_setup(&mut self, src_addr: &NodeAddr, inner: &[u8]) {
async fn handle_session_setup(
&mut self,
src_addr: &NodeAddr,
link_peer: &NodeAddr,
inner: &[u8],
) {
let setup = match SessionSetup::decode(inner) {
Ok(s) => s,
Err(e) => {
@@ -488,6 +500,40 @@ impl Node {
return;
}
// Meter the setup before anything is spent on it. This sits ahead of
// every `send_session_datagram` call in the handler — the duplicate
// ack resend, the rekey ack and the fresh-setup ack alike — so a
// refused msg1 emits nothing at all, which is what bounds the ack
// amplification. It also precedes both responder handshake
// constructions, so a refusal costs no crypto. Moving it below the
// existing-entry lookup would leave the duplicate resend unmetered.
//
// The class is read from the session table but the *key* is the link
// peer: `src_addr` is chosen by the sender, so keying on it would be
// no limit at all. A setup naming an established peer cannot grow the
// table and is metered separately, so that a stranger flood over a
// shared link cannot stop that peer's rekey from arming.
let class = if self
.sessions
.get(src_addr)
.is_some_and(|e| e.is_established())
{
Msg1Class::EstablishedLink
} else {
Msg1Class::Stranger
};
if !self.setup_rate_limiter.try_admit(link_peer, class) {
debug!(
link_peer = %self.peer_display_name(link_peer),
src = %self.peer_display_name(src_addr),
?class,
"SessionSetup rate limited"
);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::SetupRateLimited));
return;
}
// Check for existing session with this remote
if let Some(existing) = self.sessions.get(src_addr) {
if existing.is_initiating() {
@@ -538,8 +584,13 @@ impl Node {
// on the new epoch, it stops vetoing: a peer that restarted,
// or one whose own cycle lapsed, would otherwise be refused
// for as long as our own sends kept the session from idling
// out. The pending keys are not discarded here either way —
// only an authenticated msg3 replaces them.
// out. The pending keys survive both outcomes of this test:
// the veto returns without touching them, and the
// fall-through only arms a handshake beside them. Adopting
// them is still an authenticated msg3's job alone. No arm
// below drops one either: the dual-initiation arm did until
// it was narrowed to abandon only the handshake, for the
// reason recorded at that site.
let pending_outranks = existing.pending_new_session().is_some()
&& !existing.pending_stale(
Self::now_ms(),
@@ -560,13 +611,28 @@ impl Node {
.record_reject(RejectReason::Session(SessionReject::RekeyTiebreak));
return;
}
// We lose — abandon our rekey, become responder below.
// We lose — abandon the armed handshake, become responder
// below.
//
// `abandon_handshake`, not `abandon_rekey`: the gate
// above is `has_rekey_in_progress`, which says only that
// *some* handshake is armed, not that we armed it. A
// handshake the peer armed carries `rekey_initiator ==
// false` and can sit beside a completed epoch that a
// stale `pending_outranks` no longer vetoes, so a
// stranger reaches this line with two unauthenticated
// setup messages: one to arm the handshake, one to lose
// the tie-break against it. Dropping the pending session
// there kills the epoch the peer may already have cut
// over to. Only the handshake is ours to discard, and
// discarding it costs nothing, since an armed handshake
// holds no key material either endpoint is using.
debug!(
src = %self.peer_display_name(src_addr),
"Dual FSP rekey initiation: we lose (larger addr), abandoning ours"
);
let entry = self.sessions.get_mut(src_addr).unwrap();
entry.abandon_rekey();
entry.abandon_handshake();
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::RekeyYielded));
} else if pending_outranks {
@@ -722,6 +788,18 @@ impl Node {
};
// Rekey path: entry is Established with rekey_state
//
// `abandon_rekey` below, not `abandon_handshake` as in the responder
// arm of `handle_session_msg3`, and the difference rests on an
// invariant rather than on a different judgement: an entry with
// `rekey_initiator` set holds no pending session, so the two calls
// are the same action here. `set_rekey_state(_, true)` has one
// caller, `initiate_session_rekey`, which `check_session_rekey`
// never reaches for an entry holding a pending session; and
// `set_pending_session` clears `rekey_state`, so a completed
// initiator cycle leaves at most one of the two set. If that ever
// stops holding, these four sites become instances of the epoch
// discard the responder arm was fixed for.
if entry.is_established() && entry.has_rekey_in_progress() && entry.is_rekey_initiator() {
let mut handshake = match entry.take_rekey_state() {
Some(hs) => hs,
@@ -811,11 +889,34 @@ impl Node {
};
// Process XK msg2: read_xk_message_2 (extracts responder's epoch)
if let Err(e) = handshake.read_xk_message_2(&ack.handshake_payload) {
//
// Nothing here has been authenticated: the only thing tying this
// message to the initiation is the datagram's source address, an
// envelope field the sender chooses. Dropping the entry would let
// anyone able to reach us cancel any initiation in flight, so the
// entry goes back with the handshake rolled back to its pre-read
// state and the stored msg1 still scheduled for resend. The rollback
// is what makes the reinsert worth anything: `read_xk_message_2`
// mixes the sender's ephemeral in before it authenticates, so a
// handshake put back as it was left could never read the genuine
// msg2. A real peer's corrupt ack is covered by the same path — the
// responder resends its stored msg2, and the handshake sweep reaps
// the entry on its original deadline if none arrives. `touch()` is
// deliberately not called, so a spray cannot push that deadline out.
if let Err(e) = handshake.try_read_xk_message_2(&ack.handshake_payload) {
debug!(error = %e, "Failed to process Noise XK msg2 in SessionAck");
return; // Entry was already removed, don't put back a broken session
entry.set_state(EndToEndState::Initiating(handshake));
self.sessions.insert(*src_addr, entry);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::AckHandshakeFailed));
return;
}
// The three drops below stay drops. Each is downstream of a msg2 that
// already authenticated, so they are local failures rather than
// possible forgeries, and a handshake left at `Message2Done` cannot
// be re-driven from a resent msg1 anyway.
// Generate XK msg3: write_xk_message_3 (sends encrypted static + epoch)
let msg3 = match handshake.write_xk_message_3() {
Ok(m) => m,
@@ -897,6 +998,24 @@ impl Node {
};
// Rekey path: entry is Established with rekey_state (responder side)
//
// Every failure below abandons only the handshake. Nothing in a msg3
// is authenticated until `read_xk_message_3` has both succeeded and
// produced a static key matching this session's peer, so a failure
// here proves nothing about the sender and must not cost the entry
// anything it would miss. A `pending_new_session` beside the
// handshake is the epoch the real peer may already have cut over to,
// and dropping it kills the reverse direction on two unauthenticated
// messages: a forged msg1 to arm the handshake, then any garbage
// msg3. `abandon_handshake` keeps it; `abandon_rekey` does not.
//
// What `abandon_handshake` leaves behind, and why each is safe here:
// `rekey_completed_ms` must survive, since `pending_stale` reads it
// and a zeroed stamp reads as freshly completed. A stranded
// `rekey_msg3_payload` belongs to an initiator cycle and clears
// itself once `resend_pending_session_msg3` exhausts its budget.
// `peer_new_epoch_confirmed` only stops that retransmission, and
// `rekey_initiator` is false throughout this arm by its own gate.
if entry.is_established() && entry.has_rekey_in_progress() && !entry.is_rekey_initiator() {
let mut handshake = match entry.take_rekey_state() {
Some(hs) => hs,
@@ -913,7 +1032,7 @@ impl Node {
error = %e,
"Failed to process rekey XK msg3"
);
entry.abandon_rekey();
entry.abandon_handshake();
self.sessions.insert(*src_addr, entry);
return;
}
@@ -925,8 +1044,14 @@ impl Node {
let rekey_pubkey = match handshake.remote_static() {
Some(pk) => *pk,
None => {
// Not independently exercised by any test: a successful
// `read_xk_message_3` always sets the remote static, so
// reaching this needs fault injection into
// `HandshakeState`. Changed with its three siblings so
// the arm has one rule rather than three plus an
// exception.
debug!("No remote static key after processing rekey XK msg3");
entry.abandon_rekey();
entry.abandon_handshake();
self.sessions.insert(*src_addr, entry);
return;
}
@@ -936,7 +1061,7 @@ impl Node {
src = %self.peer_display_name(src_addr),
"FSP rekey: initiator static key differs from the established peer key"
);
entry.abandon_rekey();
entry.abandon_handshake();
self.sessions.insert(*src_addr, entry);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::RekeyKeyMismatch));
@@ -947,8 +1072,11 @@ impl Node {
let session = match handshake.into_session() {
Ok(s) => s,
Err(e) => {
// Also not independently exercised, for the same reason
// as the missing-static arm above: a handshake that read
// msg3 successfully always converts.
debug!(error = %e, "Failed to create session from rekey XK msg3");
entry.abandon_rekey();
entry.abandon_handshake();
self.sessions.insert(*src_addr, entry);
return;
}
@@ -1441,19 +1569,29 @@ impl Node {
// Read existing, decide, and apply the write under one guard so
// the keep-tighter update stays atomic.
let prior = map.get(&fips_addr).copied();
let actions = self.fsp.plan_path_mtu_tighten(fips_addr, prior, new_mtu);
let actions =
self.fsp
.plan_path_mtu_tighten(fips_addr, prior.map(|e| e.mtu), new_mtu);
if actions.is_empty() {
debug!(
dest = %peer_name,
fips_addr = %fips_addr,
new_mtu,
existing = prior.unwrap_or(new_mtu),
existing = prior.map(|e| e.mtu).unwrap_or(new_mtu),
"PathMtuNotification: keeping tighter existing path_mtu_lookup value"
);
}
for action in actions {
if let FspAction::TightenPathMtuLookup { fips_addr, mtu } = action {
map.insert(fips_addr, mtu);
// Held, not expiring. This value arrives inside a
// session, and a session's teardown or a PathBroken
// naming it already releases the entry. A deadline here
// would instead recreate the gap this mirror exists to
// close: a peer repeating an identical value on a stable
// path takes the unchanged early-return above and never
// rewrites the entry, so an expiring one would vanish
// and stay gone.
map.insert(fips_addr, crate::upper::tun::PathMtuEntry::held(mtu));
debug!(
dest = %peer_name,
fips_addr = %fips_addr,
@@ -1814,19 +1952,27 @@ impl Node {
// Read existing, decide, and apply the write under one guard so
// the keep-tighter update stays atomic.
let prior = map.get(&fips_addr).copied();
let actions = self.fsp.plan_path_mtu_tighten(fips_addr, prior, msg.mtu);
let actions =
self.fsp
.plan_path_mtu_tighten(fips_addr, prior.map(|e| e.mtu), msg.mtu);
if actions.is_empty() {
debug!(
dest = %peer_name,
fips_addr = %fips_addr,
bottleneck_mtu = msg.mtu,
existing = prior.unwrap_or(msg.mtu),
existing = prior.map(|e| e.mtu).unwrap_or(msg.mtu),
"Reactive MtuExceeded: keeping tighter existing path_mtu_lookup value"
);
}
for action in actions {
if let FspAction::TightenPathMtuLookup { fips_addr, mtu } = action {
map.insert(fips_addr, mtu);
// Held, not expiring. The admission gate above requires
// a session for the named destination, and that
// session's teardown releases this entry. Nothing
// re-sends the signal once traffic is sized to fit, so a
// deadline would drop a genuine persistent bottleneck
// and start the next flow at the conservative ceiling.
map.insert(fips_addr, crate::upper::tun::PathMtuEntry::held(mtu));
debug!(
dest = %peer_name,
fips_addr = %fips_addr,
+89 -1
View File
@@ -7,7 +7,7 @@ use crate::proto::fmp::{
ConnAction, ConnSnapshot, LifecycleView, PeerSnapshot, RekeyResendSnapshot,
};
use crate::transport::LinkId;
use tracing::{debug, info};
use tracing::{debug, info, warn};
impl LifecycleView for Node {
fn stale_connections(&self, now_ms: u64, timeout_ms: u64) -> Vec<ConnSnapshot> {
@@ -515,4 +515,92 @@ impl Node {
);
}
}
/// Expire `path_mtu_lookup` entries that nothing else will ever release.
///
/// The three callers of `path_mtu_lookup_release` all fire on session
/// state, so an entry written by the discovery `LookupResponse` carrier
/// for a destination this node never opens a session with has no release
/// path at all. Keep-tighter then makes one such response permanent: a
/// `path_mtu` of 256 pins that destination's SYN-time MSS clamp at 119
/// bytes until the process restarts. Only those entries carry a
/// `learned_ms`, and only they are expired here.
///
/// The deadline is the coordinate cache's own TTL, because the same
/// `LookupResponse` writes both stores: the clamp cannot outlive the
/// route it was learned with, and shortening `node.cache.coord_ttl_secs`
/// shortens this with it. A TTL of zero disables the pass, matching
/// `purge_idle_sessions`.
pub(in crate::node) fn purge_expired_path_mtu(&mut self, now_ms: u64) {
use crate::upper::tun::PathMtuEntry;
let ttl_ms = self.config().node.cache.coord_ttl_secs * 1000;
if ttl_ms == 0 {
return; // disabled
}
let stale = |e: &PathMtuEntry| {
e.learned_ms
.is_some_and(|at| now_ms.saturating_sub(at) >= ttl_ms)
};
// Read-scan first: an ordinary tick expires nothing, and the TUN
// reader and writer take this lock on every packet.
let expired: Vec<crate::FipsAddress> = match self.path_mtu_lookup.read() {
Ok(map) => map
.iter()
.filter(|(_, e)| stale(e))
.map(|(a, _)| *a)
.collect(),
Err(e) => {
warn!(error = %e, "path_mtu_lookup read lock poisoned; expiry pass skipped");
return;
}
};
if expired.is_empty() {
return;
}
match self.path_mtu_lookup.write() {
Ok(mut map) => {
for addr in &expired {
// Re-test under the write lock. The read guard was dropped
// before this one was taken, so a fresh value may have
// landed in between; without this the pass would delete a
// value that was just learned.
if map.get(addr).is_some_and(stale) {
map.remove(addr);
}
}
}
Err(e) => {
warn!(error = %e, "path_mtu_lookup write lock poisoned; entries not expired");
return;
}
}
// Restore what local configuration knows, the same way
// `path_mtu_lookup_release` does. A tighter remote claim overwrites a
// direct peer's link MTU under keep-tighter, so an expired entry may
// be sitting on top of a seed, and a bare removal would drop that peer
// to the conservative ceiling until its link re-handshakes.
let gone: std::collections::HashSet<crate::FipsAddress> = expired.iter().copied().collect();
let seeds: Vec<(
crate::NodeAddr,
crate::transport::TransportId,
crate::transport::TransportAddr,
)> = self
.peers
.iter()
.filter(|(addr, _)| gone.contains(&crate::FipsAddress::from_node_addr(addr)))
.filter_map(|(addr, p)| Some((*addr, p.transport_id()?, p.current_addr()?.clone())))
.collect();
for (addr, tid, taddr) in seeds {
self.seed_path_mtu_for_link_peer(&addr, tid, &taddr);
}
debug!(
expired = expired.len(),
"Expired remote-learned path_mtu_lookup entries"
);
}
}
+10 -4
View File
@@ -118,10 +118,16 @@ impl ForwardingMetrics {
/// This is **not** a packet drop. The frame is still delivered or
/// forwarded by the normal path; only the opportunistic warm attempt was
/// abandoned, so this must never be folded into the rejection family or
/// rendered as dropped traffic. `bytes` is the payload size of the frame
/// whose warm attempt was abandoned, not volume dropped; it is carried so
/// the counter can be rendered as a packets-and-bytes pair like its
/// siblings.
/// rendered as dropped traffic. `bytes` is the outer `SessionDatagram`
/// payload of the frame whose warm attempt was abandoned — the buffer
/// after `dispatch_link_message` strips the msg_type byte, not the inner
/// FSP payload the warm path reads. That is the same basis as
/// [`Self::record_received`] and [`Self::record_reject_bytes`], and it has
/// to be: all three render through one `fwd_value` row on the fipstop
/// Routing tab, where a mixed basis reads as a smaller flood than the one
/// actually arriving. It is volume observed, not volume dropped, and is
/// carried so the counter can be rendered as a packets-and-bytes pair like
/// its siblings.
#[inline]
pub fn record_warm_malformed(&self, bytes: usize) {
self.warm_malformed_packets.inc();
+86 -5
View File
@@ -28,7 +28,7 @@ pub(crate) mod stats_history;
mod tests;
mod tree;
use self::rate_limit::HandshakeRateLimiter;
use self::rate_limit::{HandshakeRateLimiter, SessionSetupRateLimiter};
use self::reloadable::Reloadable;
/// Half-range of the symmetric jitter applied to the per-session rekey timer.
@@ -338,7 +338,7 @@ pub struct Node {
/// the TUN reader/writer threads at TCP MSS clamp time so the
/// SYN/SYN-ACK clamp can use the smaller of the local-egress floor
/// and the learned per-destination path MTU.
path_mtu_lookup: Arc<std::sync::RwLock<HashMap<crate::FipsAddress, u16>>>,
path_mtu_lookup: crate::upper::tun::PathMtuLookup,
// === Transports & Links ===
/// Active transports (owned by Node).
@@ -474,6 +474,8 @@ pub struct Node {
// === Rate Limiting ===
/// Rate limiter for msg1 processing (DoS protection).
msg1_rate_limiter: HandshakeRateLimiter,
/// Rate limiter for inbound FSP SessionSetup, keyed on the link peer.
setup_rate_limiter: SessionSetupRateLimiter,
/// Rate limiter for ICMP Packet Too Big messages.
icmp_rate_limiter: IcmpRateLimiter,
/// Routing-subsystem state (routing error-signal rate limiter).
@@ -585,6 +587,28 @@ fn build_msg1_rate_limiter(config: &Config) -> HandshakeRateLimiter {
)
}
/// Build the per-link session-setup limiter's two bucket sizes.
///
/// The stranger bucket is configured directly. The established bucket is
/// derived exactly as the FMP limiter's is, from `max_peers`, the rekey
/// period and the resend budget: a hub neighbour can legitimately carry the
/// rekey traffic of every session this node holds, so that is the population
/// the per-link bucket has to cover.
fn build_setup_rate_limiter(config: &Config) -> SessionSetupRateLimiter {
let rl = &config.node.rate_limit;
let established = rate_limit::derive_established_bucket(
config.node.limits.max_peers,
config.node.rekey.after_secs,
rl.handshake_max_resends,
rl.session_setup_burst,
rl.session_setup_rate,
);
SessionSetupRateLimiter::with_params(
(rl.session_setup_burst, rl.session_setup_rate),
established,
)
}
impl Node {
/// Create a new node from configuration.
pub fn new(config: Config) -> Result<Self, NodeError> {
@@ -630,6 +654,7 @@ impl Node {
config.node.cache.coord_ttl_secs * 1000,
);
let msg1_rate_limiter = build_msg1_rate_limiter(&config);
let setup_rate_limiter = build_setup_rate_limiter(&config);
let max_connections = config.node.limits.max_connections;
let max_peers = config.node.limits.max_peers;
@@ -705,6 +730,7 @@ impl Node {
peers_by_index: HashMap::new(),
pending_outbound: HashMap::new(),
msg1_rate_limiter,
setup_rate_limiter,
icmp_rate_limiter: IcmpRateLimiter::new(),
routing: Router::new(),
fmp: Fmp::new(),
@@ -777,6 +803,7 @@ impl Node {
config.node.cache.coord_ttl_secs * 1000,
);
let msg1_rate_limiter = build_msg1_rate_limiter(&config);
let setup_rate_limiter = build_setup_rate_limiter(&config);
let max_connections = config.node.limits.max_connections;
let max_peers = config.node.limits.max_peers;
@@ -849,6 +876,7 @@ impl Node {
peers_by_index: HashMap::new(),
pending_outbound: HashMap::new(),
msg1_rate_limiter,
setup_rate_limiter,
icmp_rate_limiter: IcmpRateLimiter::new(),
routing: Router::new(),
fmp: Fmp::new(),
@@ -2592,9 +2620,21 @@ impl Node {
self.sessions.remove(remote)
}
/// Read the path_mtu_lookup entry for a destination FipsAddress.
/// Read the path MTU stored for a destination FipsAddress.
#[cfg(test)]
pub(crate) fn path_mtu_lookup_get(&self, fips_addr: &crate::FipsAddress) -> Option<u16> {
self.path_mtu_lookup
.read()
.ok()
.and_then(|map| map.get(fips_addr).map(|e| e.mtu))
}
/// Read the whole path_mtu_lookup entry, including how it is released.
#[cfg(test)]
pub(crate) fn path_mtu_lookup_entry(
&self,
fips_addr: &crate::FipsAddress,
) -> Option<crate::upper::tun::PathMtuEntry> {
self.path_mtu_lookup
.read()
.ok()
@@ -2602,10 +2642,32 @@ impl Node {
}
/// Write a path_mtu_lookup entry directly (for tests that pre-seed the map).
///
/// Writes a held entry, which is what a locally derived seed or a
/// session-carried value stores, so pre-seeding does not put a test at
/// the mercy of the expiry pass. Use `path_mtu_lookup_learn` for the
/// discovery-carrier shape.
#[cfg(test)]
pub(crate) fn path_mtu_lookup_insert(&self, fips_addr: crate::FipsAddress, mtu: u16) {
if let Ok(mut map) = self.path_mtu_lookup.write() {
map.insert(fips_addr, mtu);
map.insert(fips_addr, crate::upper::tun::PathMtuEntry::held(mtu));
}
}
/// Write an expiring path_mtu_lookup entry directly, as the discovery
/// `LookupResponse` carrier does (for tests that drive the expiry pass).
#[cfg(test)]
pub(crate) fn path_mtu_lookup_learn(
&self,
fips_addr: crate::FipsAddress,
mtu: u16,
at_ms: u64,
) {
if let Ok(mut map) = self.path_mtu_lookup.write() {
map.insert(
fips_addr,
crate::upper::tun::PathMtuEntry::learned(mtu, at_ms),
);
}
}
@@ -2622,7 +2684,26 @@ impl Node {
/// link-peer seed keeps the second while discarding the first; a plain
/// removal would silently drop a direct peer back to the conservative
/// ceiling until its link re-handshakes.
fn path_mtu_lookup_release(&self, addr: &NodeAddr) {
///
/// Two stores describe the same dead path, so this releases both: the
/// `FipsAddress`-keyed map the TCP MSS clamp reads, and the session's own
/// source-side path MTU estimate.
fn path_mtu_lookup_release(&mut self, addr: &NodeAddr) {
// The session's own source-side estimate described the same dead path,
// and the increase ladder is the only thing that would ever raise it
// again. Reset it here so the two halves of "this path is gone" stay
// together. The two timeout callers remove the session before calling
// this, so this arm is reached only from the PathBroken route, where
// the session survives the event.
//
// It runs first so the `&mut self.sessions` borrow ends before the
// shared `self.peers` borrow the reseed below takes.
if let Some(entry) = self.sessions.get_mut(addr)
&& let Some(mmp) = entry.mmp_mut()
{
mmp.path_mtu.reset_source_mtu();
}
let fips_addr = crate::FipsAddress::from_node_addr(addr);
match self.path_mtu_lookup.write() {
Ok(mut map) => {
+150 -1
View File
@@ -36,10 +36,23 @@
//! The concurrency limb (`max_pending`) is deliberately *not* split, so no
//! equivalent inflation happens there: one counter bounds simultaneous
//! in-flight handshake state whoever holds the slot.
//!
//! ## Why the session-setup limiter *is* keyed, when the msg1 limiter is not
//!
//! "Not per-source, since UDP sources are spoofable" is about the FMP link
//! layer, where the source is a transport address on an unauthenticated
//! datagram. [`SessionSetupRateLimiter`] sits a layer up and keys on
//! something different: the FMP link peer the datagram arrived over, which
//! the hop-by-hop Noise AEAD authenticates and whose population is bounded by
//! the peer table. Keying on the FSP `src_addr` instead would be the mistake
//! that sentence warns about, since that field is chosen by the sender and a
//! single sender can mint an unbounded number of distinct values.
use crate::NodeAddr;
use std::collections::HashMap;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Instant;
use std::time::{Duration, Instant};
/// Default burst capacity (max tokens).
pub const DEFAULT_BURST_CAPACITY: u32 = 100;
@@ -398,6 +411,104 @@ impl HandshakeRateLimiter {
}
}
/// How long a link peer's buckets are kept after its last setup message.
const SETUP_BUCKET_IDLE: Duration = Duration::from_secs(300);
/// One link peer's pair of session-setup buckets.
struct LinkBuckets {
/// Setup messages that would create a new half-open session entry.
stranger: TokenBucket,
/// Setup messages naming a peer this node already has a session with:
/// rekey and restart traffic, which creates no new entry.
established: TokenBucket,
/// Last time this peer was charged, for idle pruning.
seen: Instant,
}
/// Rate limiter for inbound FSP SessionSetup messages, keyed on the link peer.
///
/// The setup path allocates a `SessionEntry` and sends a SessionAck for every
/// well-formed msg1 naming an address it has no entry for, and the address is
/// an envelope field the sender picks. Without a limiter one neighbour can
/// grow the session table at whatever rate it can transmit, and buy a routed
/// ack per entry to a destination it chooses.
///
/// **The key is the FMP link peer the datagram arrived over, never the FSP
/// `src_addr`.** The link peer is authenticated by the hop-by-hop Noise AEAD
/// and its population is bounded by the peer table; `src_addr` is chosen by
/// the sender, so keying on it would let one sender mint a fresh full bucket
/// per forged message. See the module doc for how this squares with the msg1
/// limiter being unkeyed.
///
/// Two buckets per link, for the same reason [`HandshakeRateLimiter`] has
/// two: a drained stranger bucket must not also stop an established peer's
/// rekey msg1 from arming. Suppressed rekey is quiet — nothing errors and no
/// session drops — so folding both classes into one bucket would let a
/// sprayer one hop away hold forward-secrecy rotation off for everything
/// behind that link with no signal but a flat `rekey_armed`.
///
/// Idle links are pruned lazily on the admit path. Pruning only ever relaxes
/// the limit, and it cannot be farmed: earning a fresh bucket costs a full
/// [`SETUP_BUCKET_IDLE`] of silence on that link, which at any sane sizing is
/// a far lower sustained rate than simply waiting for the bucket to refill.
pub struct SessionSetupRateLimiter {
/// Per-link-peer buckets, created on first use.
buckets: HashMap<NodeAddr, LinkBuckets>,
/// Burst and refill rate for a new link's stranger bucket.
stranger: (u32, f64),
/// Burst and refill rate for a new link's established bucket.
established: (u32, f64),
}
impl SessionSetupRateLimiter {
/// Create a limiter whose per-link buckets take the given parameters.
///
/// Each pair is `(burst, tokens per second)`.
pub fn with_params(stranger: (u32, f64), established: (u32, f64)) -> Self {
Self {
buckets: HashMap::new(),
stranger,
established,
}
}
/// Charge one setup message of `class` to `link_peer`.
///
/// Returns `false` when the class's bucket for that link is empty, in
/// which case the caller must drop the message before doing any work.
pub fn try_admit(&mut self, link_peer: &NodeAddr, class: Msg1Class) -> bool {
let now = Instant::now();
let stranger = self.stranger;
let established = self.established;
let link = self
.buckets
.entry(*link_peer)
.or_insert_with(|| LinkBuckets {
stranger: TokenBucket::with_params(stranger.0, stranger.1),
established: TokenBucket::with_params(established.0, established.1),
seen: now,
});
link.seen = now;
let admitted = match class {
Msg1Class::Stranger => link.stranger.try_acquire(),
Msg1Class::EstablishedLink => link.established.try_acquire(),
};
if admitted {
self.buckets
.retain(|_, link| now.duration_since(link.seen) < SETUP_BUCKET_IDLE);
}
admitted
}
/// Number of link peers currently holding buckets.
#[cfg(test)]
pub fn len(&self) -> usize {
self.buckets.len()
}
}
#[cfg(test)]
mod tests {
use super::*;
@@ -724,4 +835,42 @@ mod tests {
assert_eq!(burst, 100);
assert_eq!(rate, 10.0);
}
fn addr(byte: u8) -> NodeAddr {
NodeAddr::from_bytes([byte; 16])
}
#[test]
fn setup_limiter_draining_one_link_leaves_another_links_budget_untouched() {
let mut limiter = SessionSetupRateLimiter::with_params((2, 0.001), (2, 0.001));
let noisy = addr(0x01);
let quiet = addr(0x02);
assert!(limiter.try_admit(&noisy, Msg1Class::Stranger));
assert!(limiter.try_admit(&noisy, Msg1Class::Stranger));
assert!(
!limiter.try_admit(&noisy, Msg1Class::Stranger),
"the noisy link's own bucket must run out"
);
assert!(
limiter.try_admit(&quiet, Msg1Class::Stranger),
"a second link peer must not share the first one's budget"
);
assert_eq!(limiter.len(), 2);
}
#[test]
fn setup_limiter_draining_the_stranger_bucket_still_admits_established_peer_setups() {
let mut limiter = SessionSetupRateLimiter::with_params((1, 0.001), (1, 0.001));
let link = addr(0x01);
assert!(limiter.try_admit(&link, Msg1Class::Stranger));
assert!(!limiter.try_admit(&link, Msg1Class::Stranger));
assert!(
limiter.try_admit(&link, Msg1Class::EstablishedLink),
"rekey traffic must not be starved by a stranger flood on the \
same link; suppressed rotation is silent and would show only as \
a flat rekey_armed counter"
);
}
}
+12
View File
@@ -226,6 +226,18 @@ pub enum SessionReject {
/// rather than arming a second handshake. Tracked via
/// [`SessionStats::rekey_pending`](crate::node::stats::SessionStats).
RekeyPending,
/// An inbound SessionAck naming a session we are initiating failed the
/// XK msg2 read. The message carries no authenticator tying it to the
/// initiation, so the entry is kept and the handshake rolled back
/// rather than discarded; a sustained rate here is either a broken path
/// to the responder or somebody spraying forged acks to hold
/// establishment down. Tracked via
/// [`SessionStats::ack_handshake_failed`](crate::node::stats::SessionStats).
AckHandshakeFailed,
/// A setup message was refused by the per-link-peer setup limiter
/// before any handshake state was created or any ack sent. Tracked via
/// [`SessionStats::setup_rate_limited`](crate::node::stats::SessionStats).
SetupRateLimited,
}
/// MMP rejection reasons.
+12 -6
View File
@@ -38,13 +38,19 @@
//!
//! - `path_mtu_lookup` is an event-driven cache (`Arc<RwLock<HashMap>>`)
//! populated from observed path-MTU discovery traffic, not loaded from a
//! file. There is nothing to poll. Release is event-driven for the same
//! reason: an entry is dropped when the path it describes is declared
//! file. There is nothing to poll. Release is mostly event-driven for the
//! same reason: an entry is dropped when the path it describes is declared
//! invalid (a `PathBroken` report, session idle expiry, or handshake
//! timeout) and the locally derived link MTU is reseeded in its place, so
//! there is no expiry sweep either. (Its read side could adopt the same
//! lock-free `ArcSwap` shape in the future, but that is an optimization, not
//! a reload.)
//! timeout) and the locally derived link MTU is reseeded in its place. All
//! three of those events read session state, which leaves one carrier
//! uncovered: a lookup `LookupResponse` writes an entry for a destination
//! this node may never open a session with. Those entries, and only those,
//! carry a learn time and are expired at the coordinate cache's TTL by
//! `purge_expired_path_mtu` on the same tick, which then reseeds any direct
//! peer whose entry went. Locally derived link MTUs and values learned
//! inside a session carry no deadline. That sweep is expiry, not a reload.
//! (The read side could adopt the same lock-free `ArcSwap` shape in the
//! future, but that is an optimization, not a reload.)
//! - `nostr_rendezvous` is an async spawned subsystem, not a snapshot of disk
//! state.
//!
+15
View File
@@ -70,6 +70,15 @@ pub struct SessionStats {
/// already have adopted, so a sustained rate means one side keeps
/// rekeying while the other never appears on the new epoch.
pub pending_replaced: u64,
/// An inbound SessionAck failed the XK msg2 read against a session we
/// are initiating. The entry is kept and the handshake rolled back,
/// since the message authenticates nothing; a sustained rate is either
/// a broken path to the responder or forged acks holding establishment
/// down.
pub ack_handshake_failed: u64,
/// A setup message was refused by the per-link-peer setup limiter,
/// before any handshake state was created or any ack sent.
pub setup_rate_limited: u64,
}
impl SessionStats {
@@ -85,6 +94,8 @@ impl SessionStats {
rekey_pending: self.rekey_pending,
rekey_expired: self.rekey_expired,
pending_replaced: self.pending_replaced,
ack_handshake_failed: self.ack_handshake_failed,
setup_rate_limited: self.setup_rate_limited,
}
}
@@ -97,6 +108,8 @@ impl SessionStats {
SessionReject::RekeyTiebreak => self.rekey_tiebreak += 1,
SessionReject::RekeyYielded => self.rekey_yielded += 1,
SessionReject::RekeyPending => self.rekey_pending += 1,
SessionReject::AckHandshakeFailed => self.ack_handshake_failed += 1,
SessionReject::SetupRateLimited => self.setup_rate_limited += 1,
}
}
}
@@ -364,6 +377,8 @@ pub struct SessionStatsSnapshot {
pub rekey_pending: u64,
pub rekey_expired: u64,
pub pending_replaced: u64,
pub ack_handshake_failed: u64,
pub setup_rate_limited: u64,
}
#[derive(Clone, Debug, Default, Serialize)]
+113
View File
@@ -1047,6 +1047,119 @@ async fn test_originator_lookup_response_keeps_tighter_path_mtu_lookup() {
);
}
/// Build a verified LookupResponse for a fresh target and hand it to the
/// handler, returning the target and its FipsAddress. The identity is
/// registered so the proof verifies and the originator branch is taken.
fn make_verified_lookup_response(
node: &mut Node,
request_id: u64,
path_mtu: u16,
) -> (crate::NodeAddr, crate::FipsAddress, Vec<u8>) {
let target_identity = Identity::generate();
let target = *target_identity.node_addr();
let target_fips = crate::FipsAddress::from_node_addr(&target);
let root = make_node_addr(0xF0);
let coords = TreeCoordinate::from_addrs(vec![target, root]).unwrap();
node.register_identity(target, target_identity.pubkey_full());
let proof_data = LookupResponse::proof_bytes(request_id, &target, &coords);
let proof = target_identity.sign(&proof_data);
let mut response = LookupResponse::new(request_id, target, coords, proof);
response.path_mtu = path_mtu;
(target, target_fips, response.encode()[1..].to_vec())
}
#[tokio::test]
async fn test_lookup_response_path_mtu_expires_without_a_session() {
// The discovery carrier writes an entry for a destination this node may
// never open a session with, and all three release callers fire on
// session state. Without a deadline, one response carrying 256 pins that
// destination's SYN-time MSS clamp at 119 bytes until the process
// restarts.
let mut node = make_node();
let from = make_node_addr(0xAA);
let (_target, target_fips, body) = make_verified_lookup_response(&mut node, 803, 256);
node.handle_lookup_response(&from, &body).await;
let entry = node
.path_mtu_lookup_entry(&target_fips)
.expect("precondition: the response wrote an entry, or the rest observes nothing");
assert_eq!(entry.mtu, 256, "precondition: the annotation was stored");
let learned_ms = entry
.learned_ms
.expect("the discovery carrier has no release path, so its entry must carry a learn time");
assert_eq!(
node.session_count(),
0,
"precondition: no session exists, so nothing but the deadline would ever release this"
);
let ttl_ms = node.config().node.cache.coord_ttl_secs * 1000;
assert!(ttl_ms > 0, "precondition: the expiry pass is not disabled");
// The healthy half: an entry inside its lifetime must survive an
// ordinary tick, or the clamp loses a value it is entitled to.
node.purge_expired_path_mtu(learned_ms + ttl_ms - 1);
assert_eq!(
node.path_mtu_lookup_get(&target_fips),
Some(256),
"an entry inside its lifetime must survive the expiry pass"
);
node.purge_expired_path_mtu(learned_ms + ttl_ms + 1);
assert_eq!(
node.path_mtu_lookup_get(&target_fips),
None,
"past its deadline the entry must go, since nothing else will ever release it"
);
}
#[tokio::test]
async fn test_replayed_lookup_response_does_not_extend_the_path_mtu_deadline() {
// The response carries no replay dedupe, so a captured one can be
// re-injected indefinitely. What bounds the damage is that a replay of a
// value already stored takes the keep-tighter arm, which does not touch
// the learn time: each injection buys one TTL, not one per packet.
let mut node = make_node();
let from = make_node_addr(0xAA);
let (_target, target_fips, body) = make_verified_lookup_response(&mut node, 804, 256);
node.handle_lookup_response(&from, &body).await;
let first = node
.path_mtu_lookup_entry(&target_fips)
.expect("precondition: the first response wrote an entry");
let learned_ms = first
.learned_ms
.expect("precondition: the entry carries a learn time");
// Real elapsed wall-clock between replays. The handler stamps its own
// `Self::now_ms()`, so a version that refreshed the deadline would be
// indistinguishable from one that did not if all three landed in the
// same millisecond.
for _ in 0..2 {
std::thread::sleep(std::time::Duration::from_millis(5));
node.handle_lookup_response(&from, &body).await;
}
assert_eq!(
node.path_mtu_lookup_entry(&target_fips),
Some(first),
"a replay must leave the entry exactly as it was, deadline included"
);
let ttl_ms = node.config().node.cache.coord_ttl_secs * 1000;
node.purge_expired_path_mtu(learned_ms + ttl_ms + 1);
assert_eq!(
node.path_mtu_lookup_get(&target_fips),
None,
"replaying the same value must not push the deadline out"
);
}
// ============================================================================
// Open-Discovery Sweep — cache-injection unit test
// ============================================================================
+30 -3
View File
@@ -330,6 +330,13 @@ async fn test_coord_cache_warming_encrypted_msg_with_coords() {
node.coord_cache().get(&dest_addr, now_ms).is_some(),
"dest coords not cached from encrypted message"
);
// Changing what the malformed counter charges is close enough to changing
// when it fires that the well-formed case is pinned in the same place.
assert_eq!(
node.metrics().forwarding.warm_malformed_packets.get(),
0,
"a well-formed CP datagram must not be counted as an abandoned warm attempt"
);
}
#[tokio::test]
@@ -1191,8 +1198,28 @@ async fn test_coord_cache_warming_short_inner_payload_is_dropped_not_panic() {
24,
"every datagram in both loops must be counted as an abandoned warm attempt"
);
assert!(
node.metrics().forwarding.warm_malformed_bytes.get() > 0,
"the byte counter must move alongside the packet counter"
// The byte counter shares a fipstop row with `received_bytes` and
// `decode_error_bytes`, so it has to measure the same population: the
// outer SessionDatagram payload, not the inner FSP one. Two assertions
// produced two different ways, because a single one cannot tell "the
// basis matches" from "two counters are wrong in the same direction".
//
// Self-derived: every one of the 24 datagrams reaches the warm guard, as
// the two packet counts above already pin, and `record_received` charges
// the identical outer slice.
assert_eq!(
node.metrics().forwarding.warm_malformed_bytes.get(),
node.metrics().forwarding.received_bytes.get(),
"the byte counter must be charged the same outer payload as its \
siblings on the same row"
);
// Literal cross-check. The first loop sends inner lengths 4..=11, so
// outer 39..=46, summing to 340; the second sends inner 12..=27, so
// outer 47..=62, summing to 872. Charging the inner payload instead
// reads 60 + 312 = 372, about 15% of the wire volume that arrived.
assert_eq!(
node.metrics().forwarding.warm_malformed_bytes.get(),
1212,
"24 frames of 39..=46 and 47..=62 outer bytes sum to 1212"
);
}
+902 -16
View File
File diff suppressed because it is too large Load Diff
+6 -6
View File
@@ -1787,7 +1787,7 @@ async fn test_seed_path_mtu_inserts_when_empty() {
.read()
.unwrap()
.get(&fips_addr)
.copied();
.map(|e| e.mtu);
assert_eq!(
stored,
Some(1452),
@@ -1827,7 +1827,7 @@ async fn test_seeded_narrow_link_mtu_reaches_the_clamp_as_a_tight_ceiling() {
.read()
.unwrap()
.get(&fips_addr)
.copied(),
.map(|e| e.mtu),
Some(240),
"the seed stores a narrow link MTU unchanged"
);
@@ -1865,7 +1865,7 @@ async fn test_seed_path_mtu_keeps_tighter_existing_value() {
node.path_mtu_lookup
.write()
.unwrap()
.insert(fips_addr, 1280);
.insert(fips_addr, crate::upper::tun::PathMtuEntry::held(1280));
node.seed_path_mtu_for_link_peer(&peer_addr, TransportId::new(1), &transport_addr);
@@ -1874,7 +1874,7 @@ async fn test_seed_path_mtu_keeps_tighter_existing_value() {
.read()
.unwrap()
.get(&fips_addr)
.copied();
.map(|e| e.mtu);
assert_eq!(
stored,
Some(1280),
@@ -1904,7 +1904,7 @@ async fn test_seed_path_mtu_tightens_looser_existing_value() {
node.path_mtu_lookup
.write()
.unwrap()
.insert(fips_addr, 1452);
.insert(fips_addr, crate::upper::tun::PathMtuEntry::held(1452));
node.seed_path_mtu_for_link_peer(&peer_addr, TransportId::new(1), &transport_addr);
@@ -1913,7 +1913,7 @@ async fn test_seed_path_mtu_tightens_looser_existing_value() {
.read()
.unwrap()
.get(&fips_addr)
.copied();
.map(|e| e.mtu);
assert_eq!(
stored,
Some(1280),
+38
View File
@@ -13,6 +13,11 @@ use std::fmt;
/// Symmetric state during handshake.
///
/// Maintains the chaining key (ck), handshake hash (h), and current cipher.
///
/// `Clone` exists for [`HandshakeState::try_read_xk_message_2`], which has to
/// put the pre-read state back after a message that mixed material in before
/// failing to authenticate.
#[derive(Clone)]
struct SymmetricState {
/// Chaining key for key derivation.
ck: [u8; 32],
@@ -743,6 +748,39 @@ impl HandshakeState {
Ok(())
}
/// Read XK message 2, leaving the handshake untouched when the message
/// does not authenticate.
///
/// `read_xk_message_2` mixes the sender's ephemeral into the symmetric
/// state before it authenticates the encrypted epoch, so a message that
/// fails partway leaves a handshake that can never read the genuine msg2
/// afterwards. A caller that keeps its session entry across a failed read
/// — because the message may be a forgery rather than a real peer's
/// corrupt reply — needs the pre-read state back.
///
/// The saved set is exactly what `read_xk_message_2` writes:
/// `symmetric`, `remote_ephemeral`, `remote_epoch` and `progress`. **That
/// mirror is manual.** A later edit that adds a write to
/// `read_xk_message_2` without adding it here silently reintroduces the
/// poisoning, and no caller can detect it.
pub fn try_read_xk_message_2(&mut self, message: &[u8]) -> Result<(), NoiseError> {
let symmetric = self.symmetric.clone();
let remote_ephemeral = self.remote_ephemeral;
let remote_epoch = self.remote_epoch;
let progress = self.progress;
match self.read_xk_message_2(message) {
Ok(()) => Ok(()),
Err(e) => {
self.symmetric = symmetric;
self.remote_ephemeral = remote_ephemeral;
self.remote_epoch = remote_epoch;
self.progress = progress;
Err(e)
}
}
}
/// Write XK message 3 (initiator only).
///
/// XK msg3: `-> s, se` + encrypted epoch
+3
View File
@@ -2,6 +2,7 @@ mod advert;
mod driver;
mod failure_state;
mod handoff;
mod offer_admission;
mod runtime;
mod signal;
mod stun;
@@ -12,6 +13,8 @@ mod types;
#[cfg(test)]
mod tests;
pub(crate) use signal::FRESHNESS_SKEW_TOLERANCE_MS;
pub use driver::{AdvertTransportSnapshot, RendezvousDriver};
pub use handoff::{BootstrapHandoffResult, EstablishedTraversal, is_punch_packet};
pub use runtime::NostrRendezvous;
+202
View File
@@ -0,0 +1,202 @@
//! Per-npub admission control for inbound traversal offers.
//!
//! The intake path used to hold a single global semaphore sized by
//! `max_concurrent_incoming_offers`, acquired before any identity check, so
//! one sender could hold every slot and deny traversal onboarding to every
//! other peer. This type takes a per-npub permit and a global permit
//! together, so a single npub can occupy at most its own allowance.
//!
//! **What this does not fix.** The global pool stays exhaustible. Nostr
//! identities are free to generate and the signal subscription carries no
//! author restriction, so an attacker willing to run
//! `ceil(global_limit / per_sender_limit)` throwaway npubs still saturates
//! the pool at an unchanged total offer rate, and an honest peer still gets
//! the same drop. What the per-npub allowance buys is that one identity can
//! no longer do it alone, and that the refusal a spamming sender receives is
//! distinguishable in the log from genuine saturation. A defence that raises
//! the attacker's cost in more than keypairs would have to price the offer
//! itself, which is a larger change than this one.
//!
//! **Lock discipline, which the correctness of the map bound rests on.**
//! `try_admit` does the prune, the map lookup and both permit acquisitions
//! under one `std::sync::Mutex` and never awaits inside it. The prune uses
//! `Arc::strong_count` as an exact idle test: an `OwnedSemaphorePermit` owns
//! an `Arc<Semaphore>`, so a count of 1 means the map holds the only
//! reference and no permit for that npub is outstanding. Moving
//! `try_acquire_owned` outside the lock breaks this — a concurrent prune
//! could evict an entry whose permit was still live, and the next offer from
//! that npub would build a fresh semaphore and over-admit. No test in this
//! module catches that rearrangement, so it has to be preserved by reading.
use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use tokio::sync::{OwnedSemaphorePermit, Semaphore};
/// Why an inbound offer was refused a slot.
///
/// The two classes stay apart because they are different operator stories:
/// one says the node is saturated, the other says a single sender is over its
/// allowance and everyone else is unaffected.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum AdmissionReject {
/// The sender already holds its whole per-npub allowance.
SenderFull,
/// The node is at `max_concurrent_incoming_offers` across all senders.
GlobalFull,
}
/// A granted admission. Both permits release on drop, so the caller only has
/// to keep this alive for as long as the offer is being processed.
pub(super) struct OfferPermit {
_sender: OwnedSemaphorePermit,
_global: OwnedSemaphorePermit,
}
/// Admits inbound offers against a per-sender allowance nested inside a
/// global bound.
pub(super) struct OfferAdmission {
global: Arc<Semaphore>,
per_sender: Mutex<HashMap<String, Arc<Semaphore>>>,
per_sender_limit: usize,
}
impl OfferAdmission {
/// Build an admission gate bounded globally by `global_limit` and per
/// sender npub by `per_sender_limit`.
pub(super) fn new(global_limit: usize, per_sender_limit: usize) -> Self {
Self {
global: Arc::new(Semaphore::new(global_limit)),
per_sender: Mutex::new(HashMap::new()),
per_sender_limit,
}
}
/// Try to take one slot for `npub`, or say which bound refused it.
///
/// The per-sender check runs first so a spamming sender is turned away
/// without ever touching the global pool, and therefore causes no churn
/// there.
pub(super) fn try_admit(&self, npub: &str) -> Result<OfferPermit, AdmissionReject> {
let mut map = self
.per_sender
.lock()
.expect("offer-admission mutex poisoned");
// Drop entries with no outstanding permit. Every live per-sender
// permit implies a live global permit, so at most `global_limit`
// senders survive a prune and the map is bounded by
// `global_limit + 1` after the insert below. A benign race with a
// permit dropping concurrently can only retain an idle entry for one
// more round, or drop an entry all of whose permits are already free;
// neither over-admits.
map.retain(|_, sem| Arc::strong_count(sem) > 1);
let sem = map
.entry(npub.to_string())
.or_insert_with(|| Arc::new(Semaphore::new(self.per_sender_limit)))
.clone();
let sender = sem
.try_acquire_owned()
.map_err(|_| AdmissionReject::SenderFull)?;
let global = self
.global
.clone()
.try_acquire_owned()
.map_err(|_| AdmissionReject::GlobalFull)?;
Ok(OfferPermit {
_sender: sender,
_global: global,
})
}
/// How many sender npubs the map currently holds.
#[cfg(test)]
pub(super) fn tracked_senders(&self) -> usize {
self.per_sender
.lock()
.expect("offer-admission mutex poisoned")
.len()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn one_sender_cannot_take_more_than_its_allowance_or_deny_another_sender() {
let admission = OfferAdmission::new(16, 4);
// The attacker pushes for the whole global pool, which is what makes
// the victim's assertion below discriminating: without per-npub
// accounting it takes all sixteen and the victim is refused.
let mut held = Vec::new();
let mut refusals = Vec::new();
for _ in 0..16 {
match admission.try_admit("npub1attacker") {
Ok(permit) => held.push(permit),
Err(reject) => refusals.push(reject),
}
}
assert_eq!(
held.len(),
4,
"one npub must not hold more than its own allowance"
);
assert!(
refusals
.iter()
.all(|reject| *reject == AdmissionReject::SenderFull),
"a sender over its allowance is refused on its own account, not on the pool's: {refusals:?}"
);
assert!(
admission.try_admit("npub1victim").is_ok(),
"a second sender must still be admitted while the pool has slots"
);
}
#[test]
fn several_peers_bootstrapping_at_once_are_all_admitted() {
let admission = OfferAdmission::new(16, 4);
let mut held = Vec::new();
for peer in 0..8 {
held.push(
admission
.try_admit(&format!("npub1peer{peer}"))
.unwrap_or_else(|e| panic!("peer {peer} should be admitted, got {e:?}")),
);
}
}
#[test]
fn an_idle_sender_is_reclaimed_so_the_map_cannot_grow_without_bound() {
let admission = OfferAdmission::new(16, 4);
// Three senders stay busy for the whole test, so the map bound and
// not a collapse to a single entry is what the assertion measures.
let mut held = Vec::new();
for peer in 0..3 {
held.push(
admission
.try_admit(&format!("npub1busy{peer}"))
.expect("a busy peer is inside both bounds"),
);
}
for sender in 0..100 {
drop(
admission
.try_admit(&format!("npub1transient{sender}"))
.expect("a transient sender is inside both bounds"),
);
}
// Three still-held entries, plus the last transient sender, which was
// inserted after the prune that would have removed it.
assert_eq!(admission.tracked_senders(), 4);
}
}
+37 -13
View File
@@ -12,13 +12,14 @@ use nostr::prelude::{
};
use nostr_sdk::{Client, ClientOptions, prelude::RelayPoolNotification};
use serde::Serialize;
use tokio::sync::{Mutex, Notify, RwLock, Semaphore, broadcast, mpsc, oneshot};
use tokio::sync::{Mutex, Notify, RwLock, broadcast, mpsc, oneshot};
use tokio::task::JoinHandle;
use tracing::{debug, info, trace, warn};
use super::advert::{AdvertMachine, PublishPlan};
use super::failure_state::FailureState;
use super::handoff::EstablishedTraversal;
use super::offer_admission::{AdmissionReject, OfferAdmission};
use super::signal::{
FreshnessOutcome, SignalEnvelope, build_signal_event, create_traversal_answer,
create_traversal_offer, estimate_clock_skew, unwrap_signal_event, validate_offer_freshness,
@@ -204,7 +205,7 @@ pub struct NostrRendezvous {
advert: AdvertMachine,
traversal: TraversalMachine,
pending_answers: Mutex<HashMap<String, oneshot::Sender<SignalEnvelope<TraversalAnswer>>>>,
offer_slots: Arc<Semaphore>,
admission: OfferAdmission,
event_tx: mpsc::UnboundedSender<BootstrapEvent>,
event_rx: Mutex<mpsc::UnboundedReceiver<BootstrapEvent>>,
connect_task: Mutex<Option<JoinHandle<()>>>,
@@ -304,7 +305,10 @@ impl NostrRendezvous {
let pubkey = keys.public_key();
let npub = crate::encode_npub(&identity.pubkey());
let (event_tx, event_rx) = mpsc::unbounded_channel();
let offer_slots = Arc::new(Semaphore::new(config.max_concurrent_incoming_offers));
let admission = OfferAdmission::new(
config.max_concurrent_incoming_offers,
config.max_concurrent_offers_per_npub,
);
let failure_state = FailureState::new(
config.failure_streak_threshold,
@@ -332,7 +336,7 @@ impl NostrRendezvous {
advert,
traversal,
pending_answers: Mutex::new(HashMap::new()),
offer_slots,
admission,
event_tx,
event_rx: Mutex::new(event_rx),
connect_task: Mutex::new(None),
@@ -801,13 +805,30 @@ impl NostrRendezvous {
&& offer.message_type == "offer"
&& offer.recipient_npub == self.npub
{
let Ok(permit) = self.offer_slots.clone().try_acquire_owned() else {
warn!(
sender_npub = %sender_npub,
limit = self.config.max_concurrent_incoming_offers,
"rate-limited inbound traversal offer (max_concurrent_incoming_offers reached); offer dropped"
);
continue;
let permit = match self.admission.try_admit(&sender_npub) {
Ok(permit) => permit,
Err(AdmissionReject::GlobalFull) => {
warn!(
sender_npub = %sender_npub,
limit = self.config.max_concurrent_incoming_offers,
"rate-limited inbound traversal offer (max_concurrent_incoming_offers reached); offer dropped"
);
continue;
}
// Debug, not warn: the party that trips this is by
// definition sending faster than the node wants, so
// a record per rejection turns the spam into log
// volume. The global-full arm above stays at warn
// and remains the operator's signal that the node
// is actually saturated.
Err(AdmissionReject::SenderFull) => {
debug!(
sender_npub = %sender_npub,
limit = self.config.max_concurrent_offers_per_npub,
"inbound traversal offer refused: sender is at its per-npub offer allowance"
);
continue;
}
};
let runtime = Arc::clone(&self);
let peer_short = short_npub(&sender_npub);
@@ -1791,7 +1812,10 @@ impl NostrRendezvous {
.opts(ClientOptions::new().autoconnect(false))
.build();
let config = NostrRendezvousConfig::default();
let offer_slots = Arc::new(Semaphore::new(config.max_concurrent_incoming_offers));
let admission = OfferAdmission::new(
config.max_concurrent_incoming_offers,
config.max_concurrent_offers_per_npub,
);
let (event_tx, event_rx) = mpsc::unbounded_channel();
let failure_state = FailureState::new(
config.failure_streak_threshold,
@@ -1818,7 +1842,7 @@ impl NostrRendezvous {
advert,
traversal,
pending_answers: Mutex::new(HashMap::new()),
offer_slots,
admission,
event_tx,
event_rx: Mutex::new(event_rx),
connect_task: Mutex::new(None),
+5 -1
View File
@@ -11,7 +11,11 @@ use super::types::{BootstrapError, PunchHint, SIGNAL_KIND, TraversalAnswer, Trav
/// past ~minutes erodes the freshness guarantee that backstops session-id
/// replay protection. Tightening it below the size of a typical un-NTP'd
/// drift defeats the purpose. 60s sits comfortably between those.
pub(super) const FRESHNESS_SKEW_TOLERANCE_MS: u64 = 60_000;
///
/// `pub(crate)` because `Config::validate` derives the freshness window from
/// it rather than restating the number, so changing it here moves the
/// validation boundary with it.
pub(crate) const FRESHNESS_SKEW_TOLERANCE_MS: u64 = 60_000;
pub(super) struct SignalEnvelope<T> {
pub(super) payload: T,
+65 -4
View File
@@ -523,15 +523,76 @@ fn planned_remote_endpoints_cap_targets_from_an_oversized_candidate_list() {
)
.expect("endpoint planning should succeed");
// The exact figures are determined by the inputs: the plan is one
// reflexive-to-reflexive pair plus 300 reflexive-to-candidate pairs, all
// distinct, so the cap lands on exactly MAX_PUNCH_TARGETS. An inequality
// here could not tell the cap from a collapse to a single target.
assert!(tally.capped > 0, "the cap should have discarded targets");
assert!(tally.suspicious());
assert!(
endpoints.len() <= 8,
"expected at most 8 endpoints, got {}",
endpoints.len()
assert_eq!(tally.admitted, 8);
assert_eq!(endpoints.len(), 8);
}
/// The IPv6 arm of `is_never_punchable_ip` is exercised as a pure predicate
/// elsewhere; this drives it end to end through the planner, which is the
/// path the reflector attack actually uses.
///
/// The local address is a ULA so `lan_refs` is non-empty and `same_subnet_24`
/// is genuinely called with two IPv6 strings. It splits on `.` and requires
/// four parts, so no IPv6 candidate can ever satisfy the /24 gate and
/// `fd00::1` is refused off-subnet rather than admitted.
#[test]
fn planned_remote_endpoints_reject_never_punchable_ipv6_candidates() {
let (endpoints, _tally) = planned_remote_endpoints(
&[addr("fd00::2", 62000)],
Some(&addr("203.0.113.10", 62000)),
&[
addr("::1", 63000),
addr("::", 63000),
addr("ff02::1", 63000),
addr("fe80::1", 63000),
addr("fd00::1", 63000),
],
Some(&addr("198.51.100.20", 63000)),
)
.expect("endpoint planning should succeed");
assert_eq!(
endpoints,
vec!["198.51.100.20:63000".parse::<SocketAddr>().unwrap()]
);
}
/// The cap test above uses public candidates so the cap and not the filter is
/// what bounds the output. The attack shape is the opposite: several hundred
/// attacker-chosen unroutable addresses. This is the only test that pins the
/// filter and the cap acting together on that input, and its failure would
/// mean the reflector is back.
#[test]
fn planned_remote_endpoints_bound_an_oversized_list_of_unroutable_candidates() {
let mut remotes = Vec::new();
for host in 1..=150u8 {
remotes.push(addr(&format!("127.0.0.{host}"), 63000));
remotes.push(addr(&format!("224.0.0.{host}"), 63000));
}
let (endpoints, tally) = planned_remote_endpoints(
&[],
Some(&addr("203.0.113.10", 62000)),
&remotes,
Some(&addr("198.51.100.20", 63000)),
)
.expect("endpoint planning should succeed");
assert_eq!(
endpoints,
vec!["198.51.100.20:63000".parse::<SocketAddr>().unwrap()]
);
assert_eq!(tally.unroutable, 300);
assert_eq!(tally.admitted, 1);
assert!(tally.suspicious());
}
/// Guards the deployment whose STUN server sits inside the private network,
/// so the observed reflexive address is itself private. Applying the /24 gate
/// to a peer's reflexive address would drop it and remove the only branch
+20 -8
View File
@@ -469,14 +469,26 @@ pub(super) fn nonce() -> String {
/// step, which is what a resume produces, saturates the punch start delay to
/// zero so punching begins immediately; the attempt's own bounds are monotonic
/// `Instant` deadlines, so its length is unaffected. A backward step lengthens
/// that delay instead and can cost a single punch attempt, which retries. Early
/// eviction from the replay window cannot admit a replay under the shipped
/// defaults, because the freshness window a replayed offer would also have to
/// satisfy (`signal_ttl_secs` plus `FRESHNESS_SKEW_TOLERANCE_MS` on each side,
/// 240s under the shipped defaults) is strictly narrower than the replay window
/// itself (`replay_window_secs`, 300s). The margin holds while
/// `signal_ttl_secs + 120 < replay_window_secs`; nothing in config validation
/// enforces that relation today.
/// that delay instead and can cost a single punch attempt, which retries.
/// Expiry-driven eviction from the replay window cannot admit a replay, because
/// the freshness window a replayed offer would also have to satisfy
/// (`signal_ttl_secs` plus `FRESHNESS_SKEW_TOLERANCE_MS` on each side, 240s
/// under the shipped defaults) is strictly narrower than the replay window
/// itself (`replay_window_secs`, 300s). `Config::validate` enforces
/// `signal_ttl_secs + 2 * FRESHNESS_SKEW_TOLERANCE_MS < replay_window_secs`, so
/// that margin can no longer be configured away.
///
/// That covers the expiry route only. `mark_session_seen` also evicts on
/// capacity, dropping the entries nearest expiry once the cache exceeds
/// `seen_sessions_max_entries`, and a session id dropped that way stays
/// replayable for the rest of its freshness window whatever the time relation
/// is. Nothing bounds that route: whether it is reachable depends on the
/// admitted-offer rate against the cache size. The shipped defaults leave a
/// wide margin — filling 2048 entries inside 300s needs about 6.8 admitted
/// offers per second, against roughly 1.2/s from a 16-slot pool whose permits
/// are held for the length of an attempt — but raising
/// `max_concurrent_incoming_offers` narrows it. That admission rate is inferred
/// from the attempt timeout, not measured.
pub(super) fn now_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
+10 -1
View File
@@ -36,7 +36,15 @@ pub(crate) enum LookupAction {
path_mtu: u16,
},
/// Mirror path_mtu into the FipsAddress-keyed TUN-shared lookup map.
WritePathMtu { target: NodeAddr, path_mtu: u16 },
///
/// `now_ms` stamps the entry's expiry deadline. It is the same instant
/// [`Self::CacheCoords`] carries, because the clamp this write feeds must
/// not outlive the route it was learned with.
WritePathMtu {
target: NodeAddr,
now_ms: u64,
path_mtu: u16,
},
/// Reset the coords-warmup counter if an established session exists.
ResetWarmupIfEstablished { target: NodeAddr },
/// Retry queued TUN packets for the target if any are pending.
@@ -270,6 +278,7 @@ pub(crate) fn on_response_accepted(
},
LookupAction::WritePathMtu {
target: *target,
now_ms,
path_mtu,
},
LookupAction::ResetWarmupIfEstablished { target: *target },
+6
View File
@@ -233,9 +233,15 @@ fn on_response_accepted_clears_state_and_emits_effects() {
match &actions[1] {
LookupAction::WritePathMtu {
target: t,
now_ms: n,
path_mtu: p,
} => {
assert_eq!(*t, target);
assert_eq!(
*n, now_ms,
"the path-MTU entry must be stamped with the same instant the \
coordinates are cached at, or the clamp outlives its route"
);
assert_eq!(*p, path_mtu);
}
_ => panic!("action[1] must be WritePathMtu"),
+29
View File
@@ -179,6 +179,35 @@ impl PathMtuState {
// No change (equal or increase not yet confirmed)
false
}
/// Forget the source-side path MTU after the path it described is gone.
///
/// Called when a destination's path is declared broken. The tightened
/// value describes a path that no longer exists, and the increase ladder
/// in [`Self::apply_notification`] (three matching higher values spanning
/// two notification intervals) is far too slow to recover it on the
/// replacement path. Returning to the no-measurement state lets
/// [`Self::seed_source_mtu`] re-derive the value from the outbound
/// transport on the next send, exactly as a fresh session does.
///
/// Returning to `u16::MAX` is not a licence to send oversized packets:
/// the TUN outbound path caps every packet at `effective_ipv6_mtu()`
/// before it consults the per-destination gate, and that gate is simply
/// inert at `u16::MAX` — the state [`Self::new`] already starts in. If
/// that earlier cap is ever removed or made conditional, this reset stops
/// being safe.
///
/// Destination-side observation state (`last_observed_mtu`,
/// `observed_changed`, `last_notification_ms`) is deliberately left
/// alone: it describes the reverse direction, which this event says
/// nothing about, and clearing it would suppress our notifications to the
/// peer until a fresh observation arrived.
pub fn reset_source_mtu(&mut self) {
self.current_mtu = u16::MAX;
self.consecutive_increase_count = 0;
self.first_increase_ms = None;
self.pending_increase_mtu = 0;
}
}
impl Default for PathMtuState {
+51
View File
@@ -63,3 +63,54 @@ fn apply_notification_increase_counter_does_not_overflow_on_a_repeated_value() {
"the increase is not yet due, so the effective MTU must be unchanged"
);
}
#[test]
fn reset_source_mtu_returns_to_the_no_measurement_state_and_clears_the_increase_ladder() {
// Two halves. The value half is trivially observable; the ladder half is
// not, because immediately after a reset every reported value is a
// decrease from u16::MAX and the decrease branch reads none of the
// increase counters. Observing them takes a later decrease followed by a
// full three-notification increase sequence: a stale
// `pending_increase_mtu` makes the first of those three take the "same
// value as pending" arm, which never sets `first_increase_ms`, so the
// increase can never be accepted at all.
let t0 = 100_000u64;
let mut state = PathMtuState::new();
// Tighten, then part-build an increase sequence on top of it.
assert!(state.apply_notification(800, t0));
state.apply_notification(1400, t0);
state.apply_notification(1400, t0 + 11_000);
assert_eq!(
state.current_mtu(),
800,
"precondition: the tightened value is in place and the increase is pending"
);
state.reset_source_mtu();
assert_eq!(
state.current_mtu(),
u16::MAX,
"the reset must return the source side to the no-measurement state"
);
// A fresh decrease, which zeroes the count and the first-increase time but
// would leave a stale pending value behind if the reset had not cleared it.
assert!(state.apply_notification(1000, t0 + 20_000));
assert_eq!(state.current_mtu(), 1000);
// Three matching higher values spanning two notification intervals. This
// is accepted only if the sequence starts from a cleared ladder.
let t1 = t0 + 30_000;
state.apply_notification(1400, t1);
state.apply_notification(1400, t1 + 11_000);
state.apply_notification(1400, t1 + 21_000);
assert_eq!(
state.current_mtu(),
1400,
"a reset that left the increase ladder behind strands the first of \
the three notifications, so the increase is never accepted"
);
}
+74
View File
@@ -442,3 +442,77 @@ fn test_handle_parent_lost_keeps_its_published_bool_return() {
assert!(published(&mut state, &BTreeMap::new(), 2000, 2000));
assert!(state.is_root());
}
#[test]
fn test_flap_dampening_counter_resets_when_the_episode_lapses_not_when_it_engages() {
// Retirement clears the switch counter when a lapsed episode is retired,
// not when the episode is armed, so switches taken *during* an episode do
// not carry into the next window. Telling those two apart needs an episode
// that is genuinely live for at least one switch, which the injected clock
// supplies without any real elapsed time.
let my_node = make_node_addr(5);
let mut state = TreeState::new(my_node, 1000);
state.set_flap_dampening(3, 60, 2);
state.set_hold_down(0);
let peer_a = make_node_addr(1);
let peer_b = make_node_addr(2);
let root = make_node_addr(0);
state.update_peer(
ParentDeclaration::new(peer_a, root, 1, 1000),
make_coords(&[1, 0]),
);
state.update_peer(
ParentDeclaration::new(peer_b, root, 1, 1000),
make_coords(&[2, 0]),
);
assert!(!state.set_parent(peer_a, 1, 1000, 1000));
state.recompute_coords();
assert!(!state.set_parent(peer_b, 2, 2000, 2000));
state.recompute_coords();
assert!(
state.set_parent(peer_a, 3, 3000, 3000),
"the third switch inside the window must arm an episode"
);
state.recompute_coords();
assert!(
state.is_flap_dampened(3000),
"a two-second episode must be live immediately after it is armed"
);
// In-episode accumulation: in production a mandatory switch bypasses the
// veto and still feeds the counter. `set_parent` is the same entry point
// those switches reach, and it must not re-arm the running episode.
assert!(
!state.set_parent(peer_b, 4, 4000, 4000),
"a switch inside a live episode must not re-arm it"
);
state.recompute_coords();
assert!(state.is_flap_dampened(4000));
assert!(
!state.is_flap_dampened(6000),
"the episode must have lapsed after its two seconds"
);
// The lapsed episode is retired on the next switch, taking the counter
// with it: the in-episode switch must not count toward the next episode.
assert!(
!state.set_parent(peer_a, 5, 5000, 6000),
"retiring a lapsed episode must clear the switch counter"
);
state.recompute_coords();
assert!(
!state.set_parent(peer_b, 6, 6000, 7000),
"a second episode must cost a fresh threshold of switches"
);
state.recompute_coords();
assert!(
state.set_parent(peer_a, 7, 7000, 8000),
"a second episode must arm once the fresh threshold is reached"
);
state.recompute_coords();
assert!(state.is_flap_dampened(8000));
}
+67 -13
View File
@@ -35,12 +35,50 @@ use tracing::{error, warn};
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "freebsd"))]
use tun::Layer;
/// One `path_mtu_lookup` entry: the MTU the TCP MSS clamp reads, plus how
/// the entry is released.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct PathMtuEntry {
/// Path MTU in bytes.
pub mtu: u16,
/// Unix ms at which a discovery `LookupResponse` supplied this value, or
/// `None` for an entry that some event releases instead of a timer.
///
/// The discovery carrier is the one with no release path: it writes an
/// entry for a destination this node may never open a session with, and
/// all three callers of `path_mtu_lookup_release` fire on session state.
/// A link MTU this node derived from its own transport, and a value
/// learned inside a session, are both released by an event that says the
/// thing they describe is gone, so they carry no deadline.
pub learned_ms: Option<u64>,
}
impl PathMtuEntry {
/// An entry released by an event rather than a timer: a locally derived
/// link MTU, or a value learned inside a session.
pub fn held(mtu: u16) -> Self {
Self {
mtu,
learned_ms: None,
}
}
/// A remote party's claim stored at `at_ms` for a destination with no
/// other release path. Expires.
pub fn learned(mtu: u16, at_ms: u64) -> Self {
Self {
mtu,
learned_ms: Some(at_ms),
}
}
}
/// Read-only handle to the per-destination path MTU map. Populated by
/// the discovery handler on `LookupResponse`; read by the TUN reader
/// (outbound clamp) and writer (inbound clamp) at TCP MSS clamp time.
/// Keyed by [`FipsAddress`] (16 bytes, the IPv6 form of a fips peer
/// address).
pub type PathMtuLookup = Arc<RwLock<HashMap<FipsAddress, u16>>>;
pub type PathMtuLookup = Arc<RwLock<HashMap<FipsAddress, PathMtuEntry>>>;
/// Compute the effective TCP MSS ceiling for a packet given its peer
/// address bytes (a 16-byte IPv6 destination on outbound, source on
@@ -107,7 +145,7 @@ pub(crate) fn per_flow_max_mss(
);
return empty_lookup_ceiling;
};
let Some(&path_mtu) = map.get(&fips_addr) else {
let Some(entry) = map.get(&fips_addr).copied() else {
trace!(
fips_addr = %fips_addr,
global_max_mss,
@@ -117,6 +155,7 @@ pub(crate) fn per_flow_max_mss(
);
return empty_lookup_ceiling;
};
let path_mtu = entry.mtu;
let path_max_mss = mss_ceiling(path_mtu);
// The actionable floor deliberately does not apply here. Every value a
// remote party supplies is refused before it can reach this map, at the
@@ -1631,7 +1670,10 @@ mod tests {
// = min(1360, 1143) = 1143.
let lookup = empty_lookup();
let addr = fips_addr_with_node_byte(0x42);
lookup.write().unwrap().insert(addr, 1280);
lookup
.write()
.unwrap()
.insert(addr, PathMtuEntry::held(1280));
assert_eq!(per_flow_max_mss(&lookup, addr.as_bytes(), 1360), 1143);
}
@@ -1641,7 +1683,10 @@ mod tests {
// global 1143 (the smaller of the two).
let lookup = empty_lookup();
let addr = fips_addr_with_node_byte(0x42);
lookup.write().unwrap().insert(addr, 1452);
lookup
.write()
.unwrap()
.insert(addr, PathMtuEntry::held(1452));
// global=1143 (UDP-1280-derived); path_max = 1452-77-60 = 1315.
assert_eq!(per_flow_max_mss(&lookup, addr.as_bytes(), 1143), 1143);
}
@@ -1655,7 +1700,10 @@ mod tests {
// learned value governs.
let lookup = empty_lookup();
let addr = fips_addr_with_node_byte(0x42);
lookup.write().unwrap().insert(addr, 1452);
lookup
.write()
.unwrap()
.insert(addr, PathMtuEntry::held(1452));
// global=1360, path_max = 1452-77-60 = 1315; min(1360, 1315) = 1315.
// 1315 > 1143, so the conservative ceiling did NOT clamp here.
assert_eq!(per_flow_max_mss(&lookup, addr.as_bytes(), 1360), 1315);
@@ -1671,7 +1719,10 @@ mod tests {
for stored in [0u16, 1, 100, 137] {
let lookup = empty_lookup();
let addr = fips_addr_with_node_byte(0x42);
lookup.write().unwrap().insert(addr, stored);
lookup
.write()
.unwrap()
.insert(addr, PathMtuEntry::held(stored));
assert_eq!(
per_flow_max_mss(&lookup, addr.as_bytes(), 1360),
1143,
@@ -1705,7 +1756,10 @@ mod tests {
.map(|&(stored, _)| {
let lookup = empty_lookup();
let addr = fips_addr_with_node_byte(0x42);
lookup.write().unwrap().insert(addr, stored);
lookup
.write()
.unwrap()
.insert(addr, PathMtuEntry::held(stored));
(stored, per_flow_max_mss(&lookup, addr.as_bytes(), 1360))
})
.collect();
@@ -1725,10 +1779,10 @@ mod tests {
// sub-floor table above.
let lookup = empty_lookup();
let addr = fips_addr_with_node_byte(0x42);
lookup
.write()
.unwrap()
.insert(addr, super::super::icmp::MIN_ACTIONABLE_PATH_MTU);
lookup.write().unwrap().insert(
addr,
PathMtuEntry::held(super::super::icmp::MIN_ACTIONABLE_PATH_MTU),
);
assert_eq!(per_flow_max_mss(&lookup, addr.as_bytes(), 1360), 119);
}
@@ -1757,8 +1811,8 @@ mod tests {
let lookup = empty_lookup();
let a = fips_addr_with_node_byte(0x10);
let b = fips_addr_with_node_byte(0x20);
lookup.write().unwrap().insert(a, 1280);
lookup.write().unwrap().insert(b, 1452);
lookup.write().unwrap().insert(a, PathMtuEntry::held(1280));
lookup.write().unwrap().insert(b, PathMtuEntry::held(1452));
assert_eq!(per_flow_max_mss(&lookup, a.as_bytes(), 1360), 1143);
assert_eq!(per_flow_max_mss(&lookup, b.as_bytes(), 1360), 1315);
}
+4 -1
View File
@@ -102,7 +102,10 @@ node:
- "$stun_addr"
signal_ttl_secs: 30
attempt_timeout_secs: 6
replay_window_secs: 60
# Must stay above signal_ttl_secs plus 60s of clock-skew grace on each
# side, or config validation refuses to start the node. 180 keeps a 30s
# margin over the 150s freshness window the 30s TTL implies.
replay_window_secs: 180
punch_start_delay_ms: 500
punch_interval_ms: 100
punch_duration_ms: 2500