Skip to main content

moq_net/lite/
setup.rs

1//! The lite-05+ SETUP message: each endpoint advertises its capabilities once, as
2//! the sole message on a unidirectional Setup Stream, then closes it.
3
4use crate::coding::*;
5
6use super::{Message, Parameters, Version};
7
8/// Setup Parameter id for the Probe capability level.
9const PARAM_PROBE: u64 = 0x1;
10/// Setup Parameter id for the request Path (client-only, URI-less transports).
11const PARAM_PATH: u64 = 0x2;
12/// Setup Parameter id for the client's intended [`Role`] (client-only).
13const PARAM_ROLE: u64 = 0x3;
14/// Setup Parameter id for the link cost the dialer assigns to this connection.
15const PARAM_COST: u64 = 0x4;
16/// Setup Parameter id for the endpoint's Hop ID.
17const PARAM_HOP: u64 = 0x5;
18
19/// The cost of crossing a link that neither end priced.
20///
21/// One, so a mesh that configures no costs accumulates a route cost equal to the
22/// hop count and ranks routes exactly as pre-lite-06 shortest-path routing did. Pricing
23/// a link at 0 makes it free (a sibling in the same datacenter); pricing it higher
24/// makes it a last resort (a metered backbone).
25pub const DEFAULT_COST: u64 = 1;
26
27/// The probe capability an endpoint advertises in SETUP.
28///
29/// Monotonic: a higher level implies every lower one. An unknown (future) value
30/// decodes as the highest level we understand, so a peer that gains a new level is
31/// treated as at least [`Increase`](Self::Increase).
32#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Default)]
33pub enum ProbeLevel {
34	/// No probing. Equivalent to omitting the parameter.
35	#[default]
36	None,
37	/// The publisher can measure and periodically report at least one of the PROBE
38	/// metrics: its estimated bitrate, its round-trip time, or both. Either may be
39	/// unknown in any given report, since the two are independent on the wire.
40	Report,
41	/// The publisher can additionally pad the connection (or send redundant data).
42	Increase,
43}
44
45impl ProbeLevel {
46	/// The level to advertise for `session`, from what its transport actually exposes.
47	///
48	/// [`Report`](Self::Report) claims the publisher can measure and periodically
49	/// report. A transport that exposes neither a send-rate estimate nor an RTT can
50	/// honour neither, and the draft requires such a publisher to reset any Probe
51	/// Stream a subscriber opens. Advertising [`None`](Self::None) instead stops the
52	/// subscriber opening one at all.
53	///
54	/// Both metrics are sampled rather than declared, so this only works for a
55	/// transport whose figures exist by the time the session starts. QUIC and TCP
56	/// both qualify: their RTT comes from the handshake, which has already happened.
57	pub fn detect<S: crate::transport::poll::Session>(session: &S) -> Self {
58		use web_transport_trait::Stats as _;
59		let stats = session.stats();
60		match stats.estimated_send_rate().is_some() || stats.rtt().is_some() {
61			true => Self::Report,
62			false => Self::None,
63		}
64	}
65
66	/// Map the wire value to a level, saturating unknown values to [`Increase`](Self::Increase).
67	fn from_code(code: u64) -> Self {
68		match code {
69			0 => Self::None,
70			1 => Self::Report,
71			_ => Self::Increase,
72		}
73	}
74
75	/// The wire value for this level.
76	fn to_code(self) -> u64 {
77		match self {
78			Self::None => 0,
79			Self::Report => 1,
80			Self::Increase => 2,
81		}
82	}
83}
84
85/// The single direction a client intends to use the session for.
86///
87/// A client advertises this in its SETUP so the server can reject a token that lacks
88/// the matching scope during the handshake, instead of accepting a connection that
89/// then silently carries no media (a subscribe-only token used to publish, or vice
90/// versa). It only ever narrows what the server grants, so it is not a security
91/// boundary: the server still enforces the token's scope regardless.
92///
93/// A session is bidirectional by default, which the wire says by omitting the
94/// parameter. `Option<Role>` mirrors that: `None` is the default, and it's also what
95/// a client that predates the parameter decodes to.
96#[derive(Debug, Clone, Copy, PartialEq, Eq)]
97#[non_exhaustive]
98pub enum Role {
99	/// The client will publish tracks (ingest); the server must consume.
100	Publisher,
101	/// The client will subscribe to tracks (egress); the server must publish.
102	Subscriber,
103}
104
105impl Role {
106	/// Map the wire value to a role. `0` and any unrecognized future value are `None`
107	/// (bidirectional): the draft requires a receiver that does not recognize the value
108	/// to treat it as both directions, so a newer client can't break an older server (it
109	/// just loses the early reject and defers fully to the token's scope).
110	fn from_code(code: u64) -> Option<Self> {
111		match code {
112			1 => Some(Role::Publisher),
113			2 => Some(Role::Subscriber),
114			_ => None,
115		}
116	}
117
118	/// The wire value for this role.
119	fn to_code(self) -> u64 {
120		match self {
121			Role::Publisher => 1,
122			Role::Subscriber => 2,
123		}
124	}
125
126	/// Derive the advertised role from which origins a client wired up: publish-only is
127	/// a [`Publisher`](Role::Publisher), consume-only a [`Subscriber`](Role::Subscriber),
128	/// and both (or neither) advertises nothing. This keeps the advertised role from
129	/// drifting away from what the session actually does.
130	pub(crate) fn from_origins(publishes: bool, consumes: bool) -> Option<Self> {
131		match (publishes, consumes) {
132			(true, false) => Some(Role::Publisher),
133			(false, true) => Some(Role::Subscriber),
134			_ => None,
135		}
136	}
137
138	/// Lowercase label for this role (`"publisher"` / `"subscriber"`).
139	pub fn as_str(self) -> &'static str {
140		match self {
141			Role::Publisher => "publisher",
142			Role::Subscriber => "subscriber",
143		}
144	}
145}
146
147impl std::fmt::Display for Role {
148	fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
149		f.write_str(self.as_str())
150	}
151}
152
153/// The SETUP message, sent once per endpoint on the unidirectional Setup Stream.
154///
155/// lite-05+ only. The two endpoints' SETUP messages are independent: neither side
156/// blocks on the peer's before opening other streams, but a stream whose encoding
157/// depends on a negotiated capability (e.g. PROBE) must wait for it.
158#[derive(Debug, Clone, Default, PartialEq, Eq)]
159pub struct Setup {
160	/// The probe capability this endpoint supports. [`ProbeLevel::None`] when absent.
161	pub probe: ProbeLevel,
162	/// The request path, for transports that carry no request URI (native QUIC,
163	/// qmux over TCP/TLS, unix sockets), with `?` and the URI query appended when
164	/// there is one. Sent only by the client; a server never sends one and a relay
165	/// never forwards it. `None` on URI-carrying bindings, where it would be a
166	/// protocol violation. An empty path means the same thing as `None`; both are
167	/// on the wire so a client need not special-case the root.
168	pub path: Option<String>,
169	/// The single direction the client intends to use, or `None` for a bidirectional
170	/// session. `None` is sent as the absence of the parameter, which is also how a
171	/// client that predates the parameter decodes.
172	pub role: Option<Role>,
173	/// What subscribing from this endpoint costs (lite-06+), added by the peer to the
174	/// route cost of every announcement we forward it.
175	///
176	/// Directional: it prices the sender's own egress, so both ends declare their own
177	/// and the two need not match. `None` means the default cost of 1.
178	pub cost: Option<u64>,
179	/// This endpoint's Hop ID, the identity it stamps onto forwarded
180	/// announcements. The peer uses it to serve this endpoint's subscriptions from
181	/// a route that does not flow through it (the same split horizon the announce
182	/// filter applies). `None` when the endpoint has no meaningful identity (a
183	/// leaf that never forwards); a wire value of 0 decodes as `None`.
184	pub hop: Option<crate::Hop>,
185}
186
187impl Message for Setup {
188	const MAX_SIZE: usize = crate::setup::MAX_SETUP_SIZE;
189
190	fn decode_msg<R: bytes::Buf>(r: &mut R, version: Version) -> Result<Self, DecodeError> {
191		if !version.has_setup_stream() {
192			return Err(DecodeError::Version);
193		}
194
195		let params = Parameters::decode(r, version)?;
196		let probe = params
197			.get_varint(PARAM_PROBE, version)?
198			.map(ProbeLevel::from_code)
199			.unwrap_or_default();
200		let path = match params.get_bytes(PARAM_PATH) {
201			Some(bytes) => Some(
202				std::str::from_utf8(bytes)
203					.map_err(|_| DecodeError::InvalidValue)?
204					.to_string(),
205			),
206			None => None,
207		};
208		let role = params.get_varint(PARAM_ROLE, version)?.and_then(Role::from_code);
209		let cost = params.get_varint(PARAM_COST, version)?;
210		// 0 is legal on the wire but carries no identity (it can't be excluded),
211		// so it decodes as "not declared" rather than an error.
212		let hop = params
213			.get_varint(PARAM_HOP, version)?
214			.and_then(|id| crate::Hop::new(id).ok());
215
216		Ok(Self {
217			probe,
218			path,
219			role,
220			cost,
221			hop,
222		})
223	}
224
225	fn encode_msg<W: bytes::BufMut>(&self, w: &mut W, version: Version) -> Result<(), EncodeError> {
226		if !version.has_setup_stream() {
227			return Err(EncodeError::Version);
228		}
229
230		let mut params = Parameters::default();
231		// None is the wire default, so omit it to keep the message empty when nothing is set.
232		if self.probe != ProbeLevel::None {
233			params.set_varint(PARAM_PROBE, self.probe.to_code(), version)?;
234		}
235		if let Some(path) = &self.path {
236			params.set_bytes(PARAM_PATH, path.as_bytes().to_vec());
237		}
238		// Bidirectional is the wire default (absence of the parameter), so only a
239		// directional role is encoded.
240		if let Some(role) = self.role {
241			params.set_varint(PARAM_ROLE, role.to_code(), version)?;
242		}
243		if let Some(cost) = self.cost {
244			params.set_varint(PARAM_COST, cost, version)?;
245		}
246		if let Some(hop) = self.hop {
247			params.set_varint(PARAM_HOP, hop.id(), version)?;
248		}
249
250		params.encode(w, version)
251	}
252}
253
254/// Shared slot for the peer's SETUP, written once when its Setup stream is read.
255///
256/// Streams whose encoding depends on a negotiated capability (e.g. the PROBE
257/// stream) wait on this before deciding what to do. Cheap to clone: every handle
258/// shares the same slot.
259#[derive(Clone, Default)]
260pub(crate) struct PeerSetup(kio::Shared<Option<Setup>>);
261
262impl PeerSetup {
263	/// Record the peer's SETUP.
264	pub fn set(&self, setup: Setup) {
265		*self.0.lock() = Some(setup);
266	}
267
268	/// Poll for the peer's advertised probe level, waiting until its SETUP arrives.
269	pub fn poll_probe_level(&self, waiter: &kio::Waiter) -> std::task::Poll<ProbeLevel> {
270		self.poll_get(waiter, |setup| setup.probe)
271	}
272
273	/// Poll for the link cost the peer (the dialing side) declared in its SETUP.
274	/// `None` when it declared none, meaning the default cost of 1.
275	pub fn poll_cost(&self, waiter: &kio::Waiter) -> std::task::Poll<Option<u64>> {
276		self.poll_get(waiter, |setup| setup.cost)
277	}
278
279	/// Poll for the [`Hop`](crate::Hop) id the peer declared in its SETUP `Hop`
280	/// parameter. `None` when it declared none: a leaf with no identity worth excluding.
281	pub fn poll_hop(&self, waiter: &kio::Waiter) -> std::task::Poll<Option<crate::Hop>> {
282		self.poll_get(waiter, |setup| setup.hop)
283	}
284
285	/// Poll for a field of the peer's SETUP.
286	///
287	/// The peer MUST send exactly one SETUP, so this resolves once that stream is read.
288	/// Pends forever if it never does; the caller is session state, dropped with the
289	/// driver.
290	fn poll_get<T>(&self, waiter: &kio::Waiter, f: impl FnOnce(&Setup) -> T) -> std::task::Poll<T> {
291		let slot = std::task::ready!(self.0.poll(waiter, |setup| {
292			if setup.is_some() {
293				std::task::Poll::Ready(())
294			} else {
295				std::task::Poll::Pending
296			}
297		}));
298		std::task::Poll::Ready(f(slot.as_ref().expect("waited for Some")))
299	}
300}
301
302#[cfg(test)]
303mod tests {
304	use super::*;
305
306	fn round_trip(msg: &Setup) -> Setup {
307		let mut buf = bytes::BytesMut::new();
308		msg.encode(&mut buf, Version::Lite05).unwrap();
309		let mut slice = &buf[..];
310		let got = Setup::decode(&mut slice, Version::Lite05).unwrap();
311		assert!(bytes::Buf::remaining(&slice) == 0, "trailing bytes after decode");
312		got
313	}
314
315	#[test]
316	fn empty_round_trip() {
317		let msg = Setup::default();
318		assert_eq!(round_trip(&msg), msg);
319	}
320
321	/// A transport exposing neither metric can't honour a `Report` claim, and the
322	/// draft makes such a publisher reset any Probe Stream a subscriber opens. It
323	/// must advertise `None` so the subscriber never opens one.
324	#[test]
325	fn detect_reports_nothing_without_stats() {
326		use crate::lite::test_transport::{SinkSession, SinkStats};
327		let session = SinkSession::new(Default::default()).with_stats(SinkStats::default());
328		assert_eq!(ProbeLevel::detect(&session), ProbeLevel::None);
329	}
330
331	/// Either metric alone is enough to report, since the two PROBE fields are
332	/// independent on the wire.
333	#[test]
334	fn detect_reports_with_either_metric() {
335		use crate::lite::test_transport::{SinkSession, SinkStats};
336
337		let rtt_only = SinkStats::default().with_rtt(std::time::Duration::from_millis(40));
338		let session = SinkSession::new(Default::default()).with_stats(rtt_only);
339		assert_eq!(ProbeLevel::detect(&session), ProbeLevel::Report);
340
341		let rate_only = SinkStats::default().with_send_rate(1_000_000);
342		let session = SinkSession::new(Default::default()).with_stats(rate_only);
343		assert_eq!(ProbeLevel::detect(&session), ProbeLevel::Report);
344	}
345
346	#[test]
347	fn probe_levels_round_trip() {
348		for probe in [ProbeLevel::None, ProbeLevel::Report, ProbeLevel::Increase] {
349			let msg = Setup {
350				probe,
351				..Default::default()
352			};
353			assert_eq!(round_trip(&msg), msg);
354		}
355	}
356
357	#[test]
358	fn cost_round_trip() {
359		// Zero is a meaningful price (a free same-datacenter link), so it must survive
360		// the round trip as `Some(0)` rather than collapsing into "unpriced".
361		for cost in [None, Some(0), Some(1), Some(7)] {
362			let msg = Setup {
363				cost,
364				..Default::default()
365			};
366			assert_eq!(round_trip(&msg), msg);
367		}
368	}
369
370	/// Parameter values follow the session's varint codec, not a fixed one: a value that
371	/// takes the two-byte QUIC form on lite-06 fits one leading-ones byte on lite-07.
372	#[test]
373	fn parameter_values_use_the_version_codec() {
374		let msg = Setup {
375			cost: Some(100),
376			..Default::default()
377		};
378		for (version, wire) in [
379			(Version::Lite06, &[0x05, 0x01, 0x04, 0x02, 0x40, 0x64][..]),
380			(Version::Lite07, &[0x04, 0x01, 0x04, 0x01, 0x64][..]),
381		] {
382			let mut buf = Vec::new();
383			msg.encode(&mut buf, version).unwrap();
384			assert_eq!(buf, wire, "{version}");
385			assert_eq!(Setup::decode(&mut &buf[..], version).unwrap(), msg, "{version}");
386		}
387	}
388
389	/// Never emit a SETUP our own receiver would refuse.
390	#[test]
391	fn encode_enforces_the_setup_limit() {
392		// Count, id, and a 4-byte length varint precede the path.
393		let at_limit = crate::setup::MAX_SETUP_SIZE - 6;
394		let msg = Setup {
395			path: Some("a".repeat(at_limit)),
396			..Default::default()
397		};
398		assert_eq!(round_trip(&msg), msg);
399
400		let msg = Setup {
401			path: Some("a".repeat(at_limit + 1)),
402			..Default::default()
403		};
404		let mut buf = bytes::BytesMut::new();
405		assert!(matches!(
406			msg.encode(&mut buf, Version::Lite05),
407			Err(EncodeError::TooLarge)
408		));
409		assert!(buf.is_empty());
410	}
411
412	#[test]
413	fn path_round_trip() {
414		let msg = Setup {
415			probe: ProbeLevel::Report,
416			path: Some("/room/123".to_string()),
417			..Default::default()
418		};
419		assert_eq!(round_trip(&msg), msg);
420	}
421
422	#[test]
423	fn hop_round_trip() {
424		let msg = Setup {
425			hop: Some(crate::Hop::new(42).unwrap()),
426			..Default::default()
427		};
428		assert_eq!(round_trip(&msg), msg);
429	}
430
431	// A declared id of 0 carries no identity (it cannot be excluded), so it
432	// decodes as absent rather than erroring.
433	#[test]
434	fn hop_zero_decodes_as_none() {
435		use crate::coding::Encode;
436
437		let version = Version::Lite05;
438		let mut params = Parameters::default();
439		params.set_varint(super::PARAM_HOP, 0, version).unwrap();
440		let mut body = bytes::BytesMut::new();
441		params.encode(&mut body, version).unwrap();
442		// Frame the body with the Message Length prefix `Setup::decode` expects.
443		let mut buf = bytes::BytesMut::new();
444		(body.len() as u64).encode(&mut buf, version).unwrap();
445		buf.extend_from_slice(&body);
446		let mut slice = &buf[..];
447		let got = Setup::decode(&mut slice, version).unwrap();
448		assert_eq!(got.hop, None);
449	}
450
451	#[test]
452	fn empty_path_round_trips() {
453		// An empty path is valid and distinct from absent only on the wire; both mean
454		// the root, so a client doesn't have to special-case it.
455		let msg = Setup {
456			path: Some(String::new()),
457			..Default::default()
458		};
459		assert_eq!(round_trip(&msg), msg);
460	}
461
462	#[test]
463	fn roles_round_trip() {
464		for role in [Some(Role::Publisher), Some(Role::Subscriber), None] {
465			let msg = Setup {
466				path: Some("/room/123".to_string()),
467				role,
468				..Default::default()
469			};
470			assert_eq!(round_trip(&msg), msg);
471		}
472	}
473
474	#[test]
475	fn unknown_probe_level_saturates_to_increase() {
476		// Frame a SETUP message carrying an unknown probe level (99) by hand: the
477		// parameters body, prefixed with its length (the lite Message size prefix).
478		let mut params = Parameters::default();
479		params.set_varint(PARAM_PROBE, 99, Version::Lite05).unwrap();
480		let mut body = Vec::new();
481		params.encode(&mut body, Version::Lite05).unwrap();
482
483		let mut buf = bytes::BytesMut::new();
484		body.len().encode(&mut buf, Version::Lite05).unwrap();
485		buf.extend_from_slice(&body);
486
487		let mut slice = &buf[..];
488		let got = Setup::decode(&mut slice, Version::Lite05).unwrap();
489		assert_eq!(got.probe, ProbeLevel::Increase);
490	}
491
492	#[test]
493	fn role_wire_codes() {
494		// The draft pins Publisher=1 / Subscriber=2. A swap here would still round-trip
495		// against our own decoder, but break every other implementation.
496		for (role, code) in [(Role::Publisher, 1u64), (Role::Subscriber, 2)] {
497			assert_eq!(role.to_code(), code);
498			assert_eq!(Role::from_code(code), Some(role));
499		}
500	}
501
502	#[test]
503	fn unknown_role_decodes_as_bidirectional() {
504		// A role value the receiver doesn't recognize (a future extension, or an explicit
505		// 0) decodes to `None` rather than failing, so a newer client can't break an older
506		// server. The draft mandates this fallback.
507		for code in [0u64, 9, 250] {
508			let mut params = Parameters::default();
509			params.set_varint(PARAM_ROLE, code, Version::Lite05).unwrap();
510			let mut body = Vec::new();
511			params.encode(&mut body, Version::Lite05).unwrap();
512
513			let mut buf = bytes::BytesMut::new();
514			body.len().encode(&mut buf, Version::Lite05).unwrap();
515			buf.extend_from_slice(&body);
516
517			let mut slice = &buf[..];
518			let got = Setup::decode(&mut slice, Version::Lite05).unwrap();
519			assert_eq!(got.role, None, "role code {code} should decode as bidirectional");
520		}
521	}
522
523	#[test]
524	fn rejects_before_lite05() {
525		let msg = Setup::default();
526		let mut buf = bytes::BytesMut::new();
527		assert!(matches!(
528			msg.encode(&mut buf, Version::Lite04),
529			Err(EncodeError::Version)
530		));
531	}
532
533	#[test]
534	fn ignores_unknown_parameters() {
535		// Frame a SETUP carrying an unknown parameter ID alongside the path.
536		let mut params = Parameters::default();
537		params.set_bytes(PARAM_PATH, b"/foo".to_vec());
538		params.set_bytes(0xbeef, b"whatever".to_vec());
539
540		let mut body = Vec::new();
541		params.encode(&mut body, Version::Lite05).unwrap();
542
543		// Wrap with the message size prefix the Message impl expects.
544		let mut buf = bytes::BytesMut::new();
545		body.len().encode(&mut buf, Version::Lite05).unwrap();
546		buf.extend_from_slice(&body);
547
548		let mut slice = &buf[..];
549		let got = Setup::decode(&mut slice, Version::Lite05).unwrap();
550		assert_eq!(got.path.as_deref(), Some("/foo"));
551	}
552}