Compare commits

..
Author SHA1 Message Date
enricobuehler 4bc7eecf05 feat(host/wire): MTU resilience for the video data plane
ci / docs-site (pull_request) Successful in 1m5s
ci / web (pull_request) Successful in 1m47s
apple / swift (pull_request) Successful in 1m22s
apple / screenshots (pull_request) Skipped
ci / rust-arm64 (pull_request) Successful in 2m46s
android / android (pull_request) Successful in 3m25s
windows / build (x86_64-pc-windows-msvc) (pull_request) Successful in 2m13s
windows / build (aarch64-pc-windows-msvc) (pull_request) Successful in 3m16s
ci / rust (pull_request) Successful in 22m23s
Video datagrams are sealed at a shard payload sized for a clean 1500-byte
MTU (1472-byte UDP payloads). A host whose route to the client crosses a
smaller-MTU hop (a VPN/overlay adapter claiming the LAN route, a lowered
NIC MTU) delivers every small flow — QUIC control, hole punch, input,
audio — while 100% of video datagrams die: the client sits on a black
screen reporting zero loss and the host streams into the void with every
gauge green. Field-reported as 'connects fine, black screen forever'.

Three legs, none of which changes a session on a healthy path:

- PUNKTFUNK_WIRE_MTU operator override: shard payload derived from a
  given on-wire IP MTU. Wire-compatible — Welcome::shard_payload is
  already negotiated per session (the v4/v6 split ships two values
  today) and every client follows the negotiated value.
- Detection: the QUIC MTU-discovery probe ceiling moves from quinn's
  stock 1452 to exactly the sealed video-datagram size (1472), so a
  control connection's settled MTU becomes a verdict on the path:
  settled at the ceiling proves it carries video, settled below proves
  it cannot. A per-session watcher samples after the search has settled
  (live-connection guard against mid-search false learns) and logs an
  actionable WARN naming the failure shape and the diagnosis commands.
- Healing: the measured budget is recorded per peer IP; the next
  handshake clamps shard_payload to fit, so a reconnect self-heals. A
  later session that reaches the ceiling erases the record.

Verified: core 286/286 --features quic + clippy -D warnings (macOS);
host clippy -D warnings + native:: tests 44/44 (pf-lxcheck container).
The regenerated C header picks up the new MIN_SHARD_PAYLOAD constant.
2026-08-04 18:30:33 +02:00
7 changed files with 366 additions and 223 deletions
+28 -219
View File
@@ -36,22 +36,6 @@ pub const LEGACY_STALE_MS: u64 = 1000;
/// engine's staleness zero lands at 1 s; this is the hardware-level net under an engine stall).
const BACKSTOP_LEGACY_MS: u32 = 2000;
/// The longest lease the engine honours, whatever the envelope claims — the receiver-side mirror of
/// the host's own `RUMBLE_TTL_CEIL_MS`.
///
/// No host built from this tree can exceed it (the `PUNKTFUNK_RUMBLE_TTL_MS` hatch is clamped to
/// `[150, 5000]` before it reaches the wire), so this is defence in depth against a third-party or
/// modified sender that stamps a long TTL and then wedges its renewal pump while the connection
/// stays up. It matters on exactly the platforms that sustain a level for the whole lease: Apple,
/// whose renderer deliberately keeps no staleness policy of its own, and a Deck slot, whose
/// keepalive re-kicks the actuator until the lease ends. Duration-parameterized embedders (SDL,
/// Android) already self-terminate at the clamped backstop.
///
/// Deliberately NOT `pub`: an embedder has no use for it, and every `pub` const in this crate is
/// emitted into `include/punktfunk_core.h` as an UNPREFIXED `#define` — a collision hazard the
/// header already has ~170 instances of, and one this has no reason to add to.
const MAX_LEASE_MS: u16 = 5_000;
/// One effective actuator command. `(0, 0)` means stop now. `backstop_ms` is a safety-net
/// duration for platform APIs that take one (SDL rumble, Android one-shots): the engine emits
/// explicit zeros at every policy stop, so the backstop only matters if the embedder thread itself
@@ -91,11 +75,8 @@ struct PadState {
/// A wire update landed since the last emit (level change OR renewal — renewals re-emit).
dirty: bool,
next_keepalive: Option<Instant>,
/// The exact value last handed to an embedder. `(0, 0)` ⇔ the engine believes this actuator is
/// silent. It replaces a free-running jitter phase because one field answers all three live
/// questions: would re-sending this be a no-op device write (the dedupe nudge), is a stop
/// redundant, and would the nudge synthesize the reserved stop.
last_emit: (u16, u16),
/// Current jitter phase (see [`ActuatorQuirks::dedup_jitter`]).
jitter: bool,
quirks: ActuatorQuirks,
}
@@ -107,7 +88,7 @@ impl PadState {
legacy_wire: None,
dirty: false,
next_keepalive: None,
last_emit: (0, 0),
jitter: false,
quirks: ActuatorQuirks {
keepalive_ms: 0,
min_pulse_ms: 0,
@@ -131,7 +112,6 @@ impl PadState {
self.legacy_wire = None;
self.next_keepalive = None;
self.dirty = false;
self.last_emit = (0, 0);
RumbleCommand {
pad,
low: 0,
@@ -139,40 +119,6 @@ impl PadState {
backstop_ms: 0,
}
}
/// Build the command for the pad's current level, and record what we handed out.
///
/// On a `dedup_jitter` actuator, re-emitting the value the device last took is a no-op write on
/// an SDL-class layer, so the low motor's LSB is nudged. Keying that on `last_emit` rather than
/// on a free-running phase is what makes it work on EVERY emit path. Previously the nudge lived
/// only in the keepalive branch, so a host renewal — which arrives every `ttl*3/10` ms, 120 ms
/// at the 400 ms default and 60 ms at the hatch floor — re-emitted the raw level, collided with
/// the last jittered write, was swallowed, AND re-anchored the keepalive. That stretched the
/// gap between *distinct* device writes to 80 ms at the default cadence and 100 ms at the
/// floor, on an actuator whose quirk declares 40.
///
/// The nudge is refused when it would synthesize the reserved `(0, 0)` stop. That is level
/// `(1, 0)` and only that: `high` must already be 0, and `low ^ 1 == 0` implies `low == 1`.
/// There the LSB steps up instead, so the phase still alternates (1 ↔ 3, two parts in 65535)
/// and the pad never receives a stop the policy did not order.
fn emit(&mut self, pad: u16) -> RumbleCommand {
let (mut low, high) = self.level;
if self.quirks.dedup_jitter && (low, high) == self.last_emit {
let alt = low ^ 1;
low = if (alt, high) == (0, 0) {
low | 0b10
} else {
alt
};
}
self.last_emit = (low, high);
RumbleCommand {
pad,
low,
high,
backstop_ms: self.backstop(),
}
}
}
/// The pure per-connection policy state machine. Time is always passed in (`now`) so the policy
@@ -210,8 +156,6 @@ impl RumbleEngine {
p.dirty = true;
match ttl_ms {
Some(t) => {
// Never honour a lease longer than [`MAX_LEASE_MS`], whatever the sender claims.
let t = t.min(MAX_LEASE_MS);
p.ttl_ms = t;
p.legacy_wire = None;
p.deadline = if (low, high) != (0, 0) {
@@ -270,25 +214,22 @@ impl RumbleEngine {
if p.dirty {
p.dirty = false;
if p.level == (0, 0) {
// Relay a stop only if the actuator is, as far as the engine knows, still
// buzzing. A zero on an already-silent pad heals nothing and costs every
// embedder a command — Android an unconditional log line plus a binder
// `cancel()`. Two senders produce them: the host's deliberate
// `RUMBLE_STOP_BURST` re-sends after the first stop already landed, and (behind
// `PUNKTFUNK_RUMBLE_ENVELOPE=0`) the legacy flat 500 ms refresh, which re-sends
// zeros for every latched pad for the rest of the session. The burst still
// heals the case it exists for: a LOST first stop leaves the pad buzzing, so
// `last_emit != (0, 0)` and the re-send does emit.
if p.last_emit != (0, 0) {
return (Some(p.silence(pad)), None);
}
continue;
return (Some(p.silence(pad)), None);
}
if p.quirks.keepalive_ms > 0 {
p.next_keepalive =
Some(now + Duration::from_millis(p.quirks.keepalive_ms as u64));
}
return (Some(p.emit(pad)), None);
let (low, high) = p.level;
return (
Some(RumbleCommand {
pad,
low,
high,
backstop_ms: p.backstop(),
}),
None,
);
}
// 4) actuator-decay keepalive, bounded by (1)/(2) above by construction: an expired
// or stale pad was silenced before reaching here, so a keepalive can never sustain a
@@ -298,7 +239,20 @@ impl RumbleEngine {
let due = *p.next_keepalive.get_or_insert(now + ka);
if now >= due {
p.next_keepalive = Some(now + ka);
return (Some(p.emit(pad)), None);
let (mut low, high) = p.level;
if p.quirks.dedup_jitter {
p.jitter = !p.jitter;
low ^= p.jitter as u16;
}
return (
Some(RumbleCommand {
pad,
low,
high,
backstop_ms: p.backstop(),
}),
None,
);
}
merge_wake(&mut wake, due);
}
@@ -403,22 +357,6 @@ pub(crate) struct Closed;
mod tests {
use super::*;
/// The Steam Deck's declared quirks — the only shipping actuator with `dedup_jitter`.
const DECK: ActuatorQuirks = ActuatorQuirks {
keepalive_ms: 40,
min_pulse_ms: 0,
dedup_jitter: true,
};
/// Drain the engine the way an embedder does: poll until nothing is due.
fn drain(e: &mut RumbleEngine, t: Instant) -> Vec<(u16, u16)> {
let mut out = Vec::new();
while let (Some(c), _) = e.poll(t) {
out.push((c.low, c.high));
}
out
}
fn ms(v: u64) -> Duration {
Duration::from_millis(v)
}
@@ -589,133 +527,4 @@ mod tests {
);
assert_eq!(shared.next_command(ms(10)), Err(Closed));
}
/// A host renewal must not repeat the value the device last took, or an SDL-class layer
/// swallows the write. Before the jitter moved onto every emit path it lived only in the
/// keepalive branch, so each renewal collided with the last jittered write and was deduped.
#[test]
fn renewal_keeps_the_dedupe_jitter_alternating() {
let mut e = RumbleEngine::new();
e.set_quirks(0, DECK);
let t0 = Instant::now();
e.wire_update(t0, 0, 100, 200, Some(400));
assert_eq!(drain(&mut e, t0), vec![(100, 200)]);
assert_eq!(drain(&mut e, t0 + ms(40)), vec![(101, 200)]);
assert_eq!(drain(&mut e, t0 + ms(80)), vec![(100, 200)]);
// The renewal at the 120 ms default cadence: same level, must still be a distinct write.
e.wire_update(t0 + ms(120), 0, 100, 200, Some(400));
assert_eq!(drain(&mut e, t0 + ms(120)), vec![(101, 200)]);
assert_eq!(drain(&mut e, t0 + ms(160)), vec![(100, 200)]);
}
/// Phase-robust version of the same property, at the TTL hatch's 60 ms renewal floor: no two
/// consecutive DISTINCT device writes may be further apart than the declared 40 ms cadence.
#[test]
fn renewal_never_gaps_distinct_writes_at_the_60ms_floor() {
let mut e = RumbleEngine::new();
e.set_quirks(0, DECK);
let t0 = Instant::now();
let (mut last, mut last_write, mut worst) = ((0u16, 0u16), 0u64, 0u64);
for tick in 0..=360u64 {
let t = t0 + ms(tick);
if tick % 60 == 0 {
e.wire_update(t, 0, 100, 200, Some(400));
}
for v in drain(&mut e, t) {
assert_ne!(v, (0, 0), "a live lease must never emit the stop sentinel");
if v != last {
worst = worst.max(tick - last_write);
last_write = tick;
last = v;
}
}
}
assert!(
worst <= 41,
"worst distinct-write gap {worst} ms exceeds the 40 ms declared cadence"
);
}
/// The nudge must stay behind `dedup_jitter`: an off-by-one amplitude on a default-quirks pad
/// would land in Apple's identical-target comparison and Android's one-shot amplitudes.
#[test]
fn default_quirks_pads_get_the_level_verbatim_on_every_renewal() {
let mut e = RumbleEngine::new(); // Apple / Android / plain SDL
let t0 = Instant::now();
e.wire_update(t0, 0, 100, 200, Some(400));
assert_eq!(e.poll(t0).0, Some(cmd(0, 100, 200, 800)));
e.wire_update(t0 + ms(120), 0, 100, 200, Some(400));
assert_eq!(e.poll(t0 + ms(120)).0, Some(cmd(0, 100, 200, 800)));
}
/// Level `(1, 0)` is the one value whose LSB flip is the reserved stop. The nudge steps up
/// instead, so the phase still alternates and no stop is invented under a live lease.
#[test]
fn jitter_never_synthesizes_the_stop_sentinel() {
let mut e = RumbleEngine::new();
e.set_quirks(0, DECK);
let t0 = Instant::now();
e.wire_update(t0, 0, 1, 0, Some(400));
assert_eq!(e.poll(t0).0, Some(cmd(0, 1, 0, 800)));
assert_eq!(e.poll(t0 + ms(40)).0, Some(cmd(0, 3, 0, 800)));
assert_eq!(e.poll(t0 + ms(80)).0, Some(cmd(0, 1, 0, 800)));
}
/// A zero for a pad the engine already believes is silent is dropped: it heals nothing and
/// costs every embedder a command. The deliberate stop-burst heal is unaffected, because a
/// LOST stop leaves the pad buzzing and the re-send therefore does emit.
#[test]
fn a_redundant_stop_is_dropped_but_the_burst_still_heals_a_lost_one() {
let mut e = RumbleEngine::new();
let t0 = Instant::now();
e.wire_update(t0, 0, 100, 200, Some(400));
assert_eq!(drain(&mut e, t0), vec![(100, 200)]);
// First stop reaches the embedder…
e.wire_update(t0 + ms(10), 0, 0, 0, Some(0));
assert_eq!(drain(&mut e, t0 + ms(10)), vec![(0, 0)]);
// …and the burst re-sends behind it are now silent.
e.wire_update(t0 + ms(20), 0, 0, 0, Some(0));
e.wire_update(t0 + ms(30), 0, 0, 0, Some(0));
assert_eq!(drain(&mut e, t0 + ms(30)), Vec::new());
// But if the pad is buzzing (the stop that mattered was lost), a re-send still emits.
e.wire_update(t0 + ms(40), 0, 100, 200, Some(400));
assert_eq!(drain(&mut e, t0 + ms(40)), vec![(100, 200)]);
e.wire_update(t0 + ms(50), 0, 0, 0, Some(0));
assert_eq!(drain(&mut e, t0 + ms(50)), vec![(0, 0)]);
}
/// The client bounds the host's lease. `RUMBLE_TTL_CEIL_MS` is sender-side only, so a modified
/// or third-party host could otherwise stamp a huge TTL and wedge its pump, leaving Apple and
/// the Deck buzzing for the whole of it.
#[test]
fn an_overlong_lease_is_clamped_to_the_ceiling() {
let mut e = RumbleEngine::new();
let t0 = Instant::now();
e.wire_update(t0, 0, 100, 200, Some(u16::MAX));
assert_eq!(e.poll(t0).0, Some(cmd(0, 100, 200, 5000)));
// Silenced at the ceiling, not at the 65 s the sender asked for.
assert!(e.poll(t0 + ms(MAX_LEASE_MS as u64 - 1)).0.is_none());
assert_eq!(
e.poll(t0 + ms(MAX_LEASE_MS as u64)).0,
Some(cmd(0, 0, 0, 0)),
"the lease must end at the ceiling"
);
}
/// A v2 envelope carrying `ttl_ms == 0` on a LIVE level. The audit suspected the zero would be
/// mistaken for the legacy sentinel in `backstop()`; it cannot, because the expiry check
/// preempts the relay branch — the pad silences on the same poll and never reaches a backstop.
/// Pinned so that ordering stays load-bearing rather than incidental.
#[test]
fn a_zero_ttl_envelope_silences_rather_than_taking_the_legacy_backstop() {
let mut e = RumbleEngine::new();
let t0 = Instant::now();
e.wire_update(t0, 0, 100, 200, Some(0));
assert_eq!(
e.poll(t0).0,
Some(cmd(0, 0, 0, 0)),
"a zero-length lease must expire immediately, not emit with a legacy backstop"
);
}
}
+112
View File
@@ -341,6 +341,50 @@ pub fn mtu1500_shard_payload_for(peer: core::net::IpAddr) -> usize {
}
}
/// Floor for a negotiated `shard_payload` (even, well under every real path). A path whose UDP
/// budget lands below this can't carry the QUIC control plane either (QUIC's own minimum is a
/// 1200-byte UDP payload), so shrinking video shards further buys nothing — the clamp helpers
/// bottom out here instead of producing degenerate confetti-sized shards.
pub const MIN_SHARD_PAYLOAD: usize = 512;
/// The sealed wire size of a video datagram carrying `shard_payload` bytes of shard — what
/// actually leaves the socket as UDP payload (punktfunk header + shard + crypto overhead).
pub const fn sealed_datagram_bytes(shard_payload: usize) -> usize {
HEADER_LEN + shard_payload + CRYPTO_OVERHEAD
}
/// The UDP-payload size a path must carry for full-size IPv4 video datagrams: the sealed size
/// of the [`mtu1500_shard_payload`] default (= 1472, the exact 1500-MTU IPv4 ceiling). Doubles
/// as the QUIC MTU-discovery probe ceiling (`quic/endpoint.rs`): with the ceiling set to
/// exactly this value, a control connection whose discovery settles AT the ceiling has proven
/// the path carries full-size video datagrams, and one that settles BELOW it has proven the
/// path cannot — a discrimination quinn's stock 1452 ceiling can't make in either direction.
pub const fn video_datagram_udp_ceiling() -> usize {
sealed_datagram_bytes(mtu1500_shard_payload())
}
/// Largest even shard payload whose sealed datagram fits in `udp_budget` bytes of UDP payload
/// (the quantity QUIC MTU discovery measures — [`video_datagram_udp_ceiling`] is its probe
/// ceiling). Clamped to the peer's family default ([`mtu1500_shard_payload_for`]) so a generous
/// budget never grows packets past today's wire, and floored at [`MIN_SHARD_PAYLOAD`].
pub fn shard_payload_for_udp_budget(udp_budget: usize, peer: core::net::IpAddr) -> usize {
let p = udp_budget.saturating_sub(HEADER_LEN + CRYPTO_OVERHEAD);
let p = p - p % 2; // FEC requires even shards
p.clamp(MIN_SHARD_PAYLOAD, mtu1500_shard_payload_for(peer))
}
/// [`shard_payload_for_udp_budget`] for an operator-supplied ON-WIRE IP MTU (the number
/// `netsh interface ipv4 show subinterfaces` / `ip link` shows): subtracts the family's IP+UDP
/// headers first — 28 for IPv4 (and IPv4-mapped), 48 for IPv6.
pub fn shard_payload_for_wire_mtu(wire_mtu: usize, peer: core::net::IpAddr) -> usize {
let ip_udp = match peer {
core::net::IpAddr::V4(_) => 28,
core::net::IpAddr::V6(v6) if v6.to_ipv4_mapped().is_some() => 28,
core::net::IpAddr::V6(_) => 48,
};
shard_payload_for_udp_budget(wire_mtu.saturating_sub(ip_udp), peer)
}
/// Everything needed to construct a [`Session`](crate::session::Session).
///
/// `Debug` is implemented by hand to redact `key`/`salt`, and `key`/`salt` are zeroized
@@ -514,6 +558,74 @@ mod tests {
assert!(HEADER_LEN + (p + 2) + CRYPTO_OVERHEAD > 1452, "not maximal");
}
/// The video-datagram ceiling IS the exact v4 sealed size — the QUIC MTU-discovery probe
/// ceiling (endpoint.rs) relies on this equality for its settled-at-vs-below verdict.
#[test]
fn video_datagram_ceiling_is_the_sealed_default() {
assert_eq!(
video_datagram_udp_ceiling(),
HEADER_LEN + mtu1500_shard_payload() + CRYPTO_OVERHEAD
);
assert_eq!(video_datagram_udp_ceiling(), 1472);
}
/// Budget-derived sizing: even, sealed-fits-the-budget, clamped to the family default
/// above and [`MIN_SHARD_PAYLOAD`] below.
#[test]
fn shard_payload_for_udp_budget_math() {
use core::net::IpAddr;
let v4: IpAddr = "192.168.1.50".parse().unwrap();
let v6: IpAddr = "fd00::50".parse().unwrap();
// The full ceiling reproduces the default exactly.
assert_eq!(
shard_payload_for_udp_budget(video_datagram_udp_ceiling(), v4),
mtu1500_shard_payload()
);
// A WARP/Tailscale-shaped 1280 budget: sealed result must fit the budget, stay even.
let p = shard_payload_for_udp_budget(1280, v4);
assert_eq!(p % 2, 0);
assert!(sealed_datagram_bytes(p) <= 1280);
assert!(sealed_datagram_bytes(p + 2) > 1280, "not maximal");
// Odd budgets round down to even shards.
assert_eq!(shard_payload_for_udp_budget(1281, v4) % 2, 0);
// A generous budget never grows past the family default (either family).
assert_eq!(
shard_payload_for_udp_budget(9000, v4),
mtu1500_shard_payload()
);
assert_eq!(
shard_payload_for_udp_budget(9000, v6),
mtu1500_shard_payload_v6()
);
// Degenerate budgets bottom out at the floor instead of confetti.
assert_eq!(shard_payload_for_udp_budget(100, v4), MIN_SHARD_PAYLOAD);
}
/// Operator-facing wire-MTU sizing subtracts the right IP+UDP header per family, and 1500
/// reproduces today's defaults exactly.
#[test]
fn shard_payload_for_wire_mtu_math() {
use core::net::IpAddr;
let v4: IpAddr = "192.168.1.50".parse().unwrap();
let v6: IpAddr = "fd00::50".parse().unwrap();
let mapped: IpAddr = "::ffff:192.168.1.50".parse().unwrap();
assert_eq!(
shard_payload_for_wire_mtu(1500, v4),
mtu1500_shard_payload()
);
assert_eq!(
shard_payload_for_wire_mtu(1500, mapped),
mtu1500_shard_payload()
);
assert_eq!(
shard_payload_for_wire_mtu(1500, v6),
mtu1500_shard_payload_v6()
);
// 1280 wire 28 64 = 1188 (v4); 48 64 = 1168 (v6).
assert_eq!(shard_payload_for_wire_mtu(1280, v4), 1188);
assert_eq!(shard_payload_for_wire_mtu(1280, v6), 1168);
}
/// Family selection: genuine v6 remotes get the v6 size; v4 — including the IPv4-mapped v6
/// form a dual-stack `[::]` socket reports for a v4 client — keeps the v4 size.
#[test]
@@ -47,6 +47,20 @@ fn stream_transport_idle(idle: std::time::Duration) -> Arc<quinn::TransportConfi
// plane latest-wins at the source — ~200 ms of stereo Opus (proportionally less at
// surround bitrates), so sustained congestion costs concealable drops, never lag.
t.datagram_send_buffer_size(4 * 1024);
// MTU discovery probes up to EXACTLY the sealed size of a full IPv4 video datagram (1472)
// instead of quinn's stock 1452. Two reasons: (a) on a clean 1500-MTU path QUIC gets the
// last 20 bytes per packet; (b) the ceiling turns discovery into a video-path verdict the
// host's wire-MTU watcher reads (`punktfunk-host` `native/wire_mtu.rs`) — settled == ceiling
// proves the path carries full-size video datagrams, settled BELOW it proves it cannot (a
// VPN/overlay adapter at MTU ~1280 blackholes every video packet while all the small flows
// pass: the "connects fine, black screen forever" field shape). With the stock 1452 ceiling
// a healthy path and a constrained one are indistinguishable at the top. This is the ONLY
// behavioral change on healthy paths, and it's confined to discovery: probes are padded
// PINGs quinn already expects to lose above a constrained hop — a lost probe settles the
// search lower, exactly as it did before.
let mut mtud = quinn::MtuDiscoveryConfig::default();
mtud.upper_bound(crate::config::video_datagram_udp_ceiling() as u16);
t.mtu_discovery_config(Some(mtud));
Arc::new(t)
}
+4 -3
View File
@@ -26,9 +26,7 @@
#![deny(clippy::undocumented_unsafe_blocks)]
use anyhow::{anyhow, Context, Result};
use punktfunk_core::config::{
mtu1500_shard_payload_for, CompositorPref, FecConfig, FecScheme, GamepadPref, Role,
};
use punktfunk_core::config::{CompositorPref, FecConfig, FecScheme, GamepadPref, Role};
use punktfunk_core::input::{InputEvent, InputKind};
use punktfunk_core::packet::{FLAG_PIC, FLAG_PROBE, FLAG_SOF};
use punktfunk_core::quic::{
@@ -72,6 +70,9 @@ use input::{input_thread, ClientInput};
/// The Hello→Welcome→Start negotiation (plan §W1); `serve_session` calls `handshake::negotiate`
/// after the pairing gate.
mod handshake;
/// MTU resilience for the video data plane: `PUNKTFUNK_WIRE_MTU` override, the per-session
/// path-MTU watch on the control connection, and the per-peer learned shard-payload clamp.
mod wire_mtu;
/// The mid-stream control task (plan §W1); `serve_session` spawns `control::run` after the
/// handshake to multiplex renegotiation / speed-test control messages onto the data-plane channels.
+10 -1
View File
@@ -491,7 +491,12 @@ pub(super) async fn negotiate(
// per-datagram loss on Wi-Fi — the "100 Mbps badly fails on the phone" root cause.
// Negotiated, so the client follows. Jumbo (≈8900) is a future negotiated bump (needs
// MAX_DATAGRAM_BYTES raised + end-to-end 9000 MTU).
shard_payload: mtu1500_shard_payload_for(peer.ip()) as u16,
// Resolution order (wire_mtu.rs): `PUNKTFUNK_WIRE_MTU` operator override, then a path
// budget learned from a prior session whose QUIC MTU discovery settled below the
// video-datagram ceiling (the "VPN on the host blackholes every video packet" field
// shape — small flows pass, the stream is an endless black screen), then this family
// default. Healthy paths take the default branch and are byte-identical to before.
shard_payload: wire_mtu::negotiated_shard_payload(peer.ip()) as u16,
encrypt: true,
key,
salt,
@@ -658,6 +663,10 @@ pub(super) async fn negotiate(
let start =
Start::decode(&io::read_msg(recv).await?).map_err(|e| anyhow!("Start decode: {e:?}"))?;
bringup.mark("start");
// The session is real: watch this connection's MTU discovery settle and turn it into a
// path verdict (WARN + learned clamp for the next session on a constrained path; clears a
// stale clamp on a healthy one). Bounded ~10 s task, ends by itself.
wire_mtu::spawn_watch(conn.clone(), welcome.shard_payload as usize);
Ok::<_, anyhow::Error>((
hello,
welcome,
@@ -0,0 +1,192 @@
//! MTU resilience for the video data plane (the "connects fine, black screen forever" field
//! shape).
//!
//! Video datagrams are sealed at a per-session `shard_payload` sized for a clean 1500-byte MTU
//! (1472-byte UDP payloads). A host whose route to the client runs through a smaller-MTU hop —
//! a VPN/overlay adapter (Tailscale/WARP/ZeroTier default to 1280) claiming the LAN route, or a
//! lowered NIC MTU — delivers every SMALL flow (QUIC control, hole punch, input, audio) while
//! 100 % of video datagrams die by fragmentation or local `WSAEMSGSIZE`: the client sits on a
//! black screen reporting `loss_ppm=0` (it can't see gaps in packets it never saw any of) and
//! the host streams into the void with every gauge green. Neither side observes the failure
//! directly — but the control connection CAN: its MTU discovery probes up to exactly the sealed
//! video-datagram size ([`video_datagram_udp_ceiling`], set in `quic/endpoint.rs`), so its
//! settled MTU is a verdict on the path.
//!
//! Three legs, none of which changes a session on a healthy path:
//! - **`PUNKTFUNK_WIRE_MTU=<bytes>`** — operator override; the shard payload is derived from
//! the given on-wire IP MTU. Wire-compatible with every deployed client:
//! `Welcome::shard_payload` is already negotiated per session (the v4/v6 split ships two
//! values today) and clients follow the negotiated value.
//! - **Watch** — a per-session task samples the control connection's discovered MTU once the
//! search has had time to finish. A connection still alive that settled BELOW the ceiling is
//! proof the path can't carry full-size video: log an actionable WARN and record the measured
//! budget for the peer.
//! - **Heal** — the next handshake from that peer clamps `shard_payload` to the recorded
//! budget, so a reconnect fixes the stream. A later session that reaches the ceiling erases
//! the record (the learn/heal loop is self-correcting in both directions).
use std::collections::HashMap;
use std::net::IpAddr;
use std::sync::{Mutex, OnceLock};
use punktfunk_core::config::{
mtu1500_shard_payload_for, sealed_datagram_bytes, shard_payload_for_udp_budget,
shard_payload_for_wire_mtu, video_datagram_udp_ceiling,
};
/// Measured UDP-payload budget per peer IP, learned from live control connections whose MTU
/// discovery settled below the video-datagram ceiling. In-memory only: a host restart
/// re-learns in one session, and entries self-correct (a later ceiling-hit erases, a lower
/// re-measure overwrites).
fn learned() -> &'static Mutex<HashMap<IpAddr, u16>> {
static LEARNED: OnceLock<Mutex<HashMap<IpAddr, u16>>> = OnceLock::new();
LEARNED.get_or_init(|| Mutex::new(HashMap::new()))
}
/// The shard payload for a new session to `peer`: `PUNKTFUNK_WIRE_MTU` override, else the
/// peer's learned path budget, else the family default (today's exact behavior). Logs whenever
/// the result differs from the default.
pub(super) fn negotiated_shard_payload(peer: IpAddr) -> usize {
let env = match std::env::var("PUNKTFUNK_WIRE_MTU") {
Ok(v) => match v.trim().parse::<usize>() {
Ok(mtu) => Some(mtu),
Err(_) => {
tracing::warn!(value = %v, "PUNKTFUNK_WIRE_MTU is not a number — ignoring it");
None
}
},
Err(_) => None,
};
let learned_budget = learned().lock().unwrap().get(&peer).copied();
resolve(env, learned_budget, peer)
}
/// Pure resolution (env override > learned budget > family default) — the tested core of
/// [`negotiated_shard_payload`].
fn resolve(env_wire_mtu: Option<usize>, learned_udp_budget: Option<u16>, peer: IpAddr) -> usize {
let default = mtu1500_shard_payload_for(peer);
if let Some(mtu) = env_wire_mtu {
let p = shard_payload_for_wire_mtu(mtu, peer);
if p != default {
tracing::info!(
wire_mtu = mtu,
shard_payload = p,
default,
"wire MTU: shard payload set from PUNKTFUNK_WIRE_MTU"
);
}
return p;
}
if let Some(budget) = learned_udp_budget {
let p = shard_payload_for_udp_budget(budget as usize, peer);
if p != default {
tracing::info!(
peer = %peer,
udp_budget = budget,
shard_payload = p,
default,
"wire MTU: shard payload clamped to this peer's measured path MTU (learned \
from a prior session's QUIC MTU discovery) video datagrams now fit the \
constrained hop"
);
return p;
}
}
default
}
/// Sample the control connection's discovered MTU after the search has settled and turn it
/// into a verdict. Spawned once per negotiated session; the task ends by itself after the
/// final sample (bounded ~10 s lifetime, holding only a cheap `Connection` handle).
pub(super) fn spawn_watch(conn: quinn::Connection, session_shard_payload: usize) {
tokio::spawn(async move {
let peer = conn.remote_address().ip();
let ceiling = video_datagram_udp_ceiling() as u16;
// Discovery finishes in a handful of RTTs on a LAN (well under the first sample) but
// needs a loss timeout per failed probe on a constrained path — the second sample
// covers that with margin. Max, because discovery only ever raises `current_mtu`.
let mut settled = 0u16;
for wait_s in [3u64, 7] {
tokio::time::sleep(std::time::Duration::from_secs(wait_s)).await;
settled = settled.max(conn.stats().path.current_mtu);
if settled >= ceiling {
break;
}
}
if settled >= ceiling {
// The path carries full-size video datagrams — erase any stale learned clamp so
// the next session returns to the default wire.
if learned().lock().unwrap().remove(&peer).is_some() {
tracing::info!(peer = %peer,
"wire MTU: path re-measured at full size — learned clamp cleared");
}
return;
}
// A closed connection stops discovering, so a session that ended before the final
// sample proves nothing (a healthy high-RTT path could still be mid-search): learn
// only from a connection that stayed alive through the whole window.
if conn.close_reason().is_some() {
return;
}
learned().lock().unwrap().insert(peer, settled);
if sealed_datagram_bytes(session_shard_payload) <= settled as usize {
// This session was already clamped small enough — the path is still constrained
// (keep the record fresh) but video fits, so no alarm.
tracing::info!(peer = %peer, discovered_udp_mtu = settled,
"wire MTU: constrained path re-measured; this session's video is sized to fit");
} else {
tracing::warn!(
peer = %peer,
discovered_udp_mtu = settled,
needed_udp_mtu = ceiling,
"wire MTU: this path CANNOT carry full-size video datagrams — the control \
plane works but every video packet is oversized for a hop, which streams as \
an endless black screen with zero reported loss. Typical cause: a VPN/overlay \
adapter (Tailscale / Cloudflare WARP / ZeroTier) claiming the LAN route, or a \
lowered NIC MTU compare `ping <client> -f -l 1450` vs `-l 1200` and check \
`netsh interface ipv4 show subinterfaces` (Windows) / `ip link` (Linux). The \
measured budget is recorded: the NEXT session from this client sizes video to \
fit automatically. To pin it for all sessions set PUNKTFUNK_WIRE_MTU."
);
}
});
}
#[cfg(test)]
mod tests {
use super::*;
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr};
const V4: IpAddr = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 2));
const V6: IpAddr = IpAddr::V6(Ipv6Addr::new(0x2001, 0xdb8, 0, 0, 0, 0, 0, 1));
#[test]
fn default_when_nothing_known() {
assert_eq!(resolve(None, None, V4), mtu1500_shard_payload_for(V4));
assert_eq!(resolve(None, None, V6), mtu1500_shard_payload_for(V6));
}
#[test]
fn env_override_beats_learned() {
// 1280 wire 28 IP/UDP 64 header/crypto = 1188.
assert_eq!(resolve(Some(1280), Some(1472), V4), 1188);
}
#[test]
fn learned_budget_clamps() {
// A WARP-shaped path: 1280-byte UDP budget → 1280 64 = 1216.
assert_eq!(resolve(None, Some(1280), V4), 1216);
}
#[test]
fn learned_at_or_above_ceiling_is_the_default_wire() {
assert_eq!(resolve(None, Some(1472), V4), mtu1500_shard_payload_for(V4));
assert_eq!(resolve(None, Some(2000), V4), mtu1500_shard_payload_for(V4));
}
#[test]
fn env_full_mtu_is_the_default_wire_both_families() {
assert_eq!(resolve(Some(1500), None, V4), mtu1500_shard_payload_for(V4));
assert_eq!(resolve(Some(1500), None, V6), mtu1500_shard_payload_for(V6));
}
}
+6
View File
@@ -333,6 +333,12 @@
#define INBOUND_REQ_FLAG 2147483648
#endif
// Floor for a negotiated `shard_payload` (even, well under every real path). A path whose UDP
// budget lands below this can't carry the QUIC control plane either (QUIC's own minimum is a
// 1200-byte UDP payload), so shrinking video shards further buys nothing — the clamp helpers
// bottom out here instead of producing degenerate confetti-sized shards.
#define MIN_SHARD_PAYLOAD 512
// 16-byte AEAD authentication tag appended by either session cipher.
#define TAG_LEN 16