Skip to main content

moq_net/
session.rs

1//! A MoQ session handle and a snapshot of its connection statistics.
2
3use std::{sync::Arc, task::Poll, time::Duration};
4
5use web_transport_trait::Stats as _;
6
7use crate::{Error, SessionError, Version, bandwidth, goaway};
8
9/// How long [`Session::close`] waits for queued data before closing anyway.
10const CLOSE_TIMEOUT: Duration = Duration::from_secs(1);
11
12/// A close requested by a session handle, executed by the driver.
13#[derive(Clone)]
14enum Close {
15	/// Close now with this code.
16	Abort { code: u32, reason: String },
17	/// Close once the protocol owes the peer nothing, or at [`CLOSE_TIMEOUT`].
18	Drain,
19}
20
21/// How the session ended, published by the driver.
22#[derive(Clone)]
23struct Ended {
24	/// The transport's terminal error.
25	err: Error,
26	/// The drain's outcome, when a drain is what closed the transport.
27	drain: Option<Result<(), Error>>,
28}
29
30/// The stats cell shared between the driver's sampler and the handles.
31struct StatsState {
32	/// The latest sample the driver took (or the construction-time snapshot).
33	sample: Stats,
34	/// A handle read the stats since the last sample: keep sampling.
35	demanded: bool,
36}
37
38/// A snapshot of connection statistics for a [`Session`].
39///
40/// Every field is optional: availability depends on the transport backend (native QUIC
41/// reports all of them, the browser WebTransport reports few or none) and on the
42/// connection state (e.g. `estimated_send_rate` is `None` until the congestion controller
43/// has a window). `None` means "not reported", not "zero".
44#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
45#[non_exhaustive]
46pub struct Stats {
47	/// Smoothed round-trip time estimate.
48	pub rtt: Option<Duration>,
49
50	/// Estimated send bandwidth from the congestion controller.
51	pub estimated_send_rate: Option<bandwidth::Rate>,
52
53	/// Estimated receive bandwidth from MoQ PROBE.
54	///
55	/// `None` unless the negotiated version supports PROBE (moq-lite-03+).
56	pub estimated_recv_rate: Option<bandwidth::Rate>,
57
58	/// Total bytes sent over the connection, including retransmissions and overhead.
59	pub bytes_sent: Option<u64>,
60
61	/// Total bytes received over the connection, including duplicates and overhead.
62	pub bytes_received: Option<u64>,
63
64	/// Total bytes lost (detected via retransmission or acknowledgement).
65	pub bytes_lost: Option<u64>,
66
67	/// Total datagrams sent.
68	pub packets_sent: Option<u64>,
69
70	/// Total datagrams received.
71	pub packets_received: Option<u64>,
72
73	/// Total datagrams detected as lost.
74	pub packets_lost: Option<u64>,
75}
76
77/// A MoQ transport session, wrapping a WebTransport connection.
78///
79/// Returned with a [`Driver`](crate::Driver) by [`crate::Client::connect`] and
80/// [`crate::Server::accept`]. The caller must poll or spawn that driver to run
81/// the session.
82///
83/// Like every handle in this library, the lifecycle is reference counted: clones
84/// share the connection, the transport closes when the last clone drops, and
85/// [`abort`](Self::abort) closes it explicitly with an error. The handle and the
86/// driver are severed in both directions: the driver holds no `Session` clone,
87/// so running it never keeps the session alive, and the `Session`
88/// holds no transport, so the handle is `Send + Sync` whatever transport the
89/// driver uses. Everything transport-shaped (the close, the close reason,
90/// the stats sample) is relayed through the driver.
91#[derive(Clone)]
92pub struct Session {
93	/// Handle side to driver: `Some` once [`abort`](Self::abort) or
94	/// [`close`](Self::close) ran; the channel closing (the last handle
95	/// dropping) is the implicit Cancel, unless a drain was already requested.
96	close: kio::Producer<Option<Close>>,
97	/// Driver to handle side: how the transport ended.
98	closed: kio::Consumer<Option<Ended>>,
99	stats: kio::Shared<StatsState>,
100	version: Version,
101	send_bandwidth: Option<bandwidth::Consumer>,
102	recv_bandwidth: Option<bandwidth::Consumer>,
103	goaway: Arc<goaway::Handle>,
104}
105
106impl Session {
107	/// Returns the negotiated protocol version.
108	pub fn version(&self) -> Version {
109		self.version
110	}
111
112	/// Returns a consumer for the estimated send bitrate (from the congestion controller).
113	///
114	/// Returns `None` if the QUIC backend doesn't support bandwidth estimation.
115	pub fn send_bandwidth(&self) -> Option<bandwidth::Consumer> {
116		self.send_bandwidth.clone()
117	}
118
119	/// Returns a consumer for the estimated receive bitrate (from PROBE).
120	///
121	/// Returns `None` if the MoQ version doesn't support PROBE (requires moq-lite-03+).
122	pub fn recv_bandwidth(&self) -> Option<bandwidth::Consumer> {
123		self.recv_bandwidth.clone()
124	}
125
126	/// Returns a snapshot of the current connection statistics.
127	///
128	/// Cheap and non-blocking: this reads the latest sample the session's
129	/// driver took, and schedules a refresh, so periodic polling observes
130	/// fresh counters (100ms cadence). See [`Stats`] for which
131	/// metrics each backend reports.
132	pub fn stats(&self) -> Stats {
133		let mut stats = {
134			let mut state = self.stats.lock();
135			// A read is demand: wake the sampler, but only mutate (and so wake)
136			// when the flag actually flips.
137			if !state.demanded {
138				state.demanded = true;
139			}
140			state.sample
141		};
142		stats.estimated_recv_rate = self.recv_bandwidth.as_ref().and_then(bandwidth::Consumer::peek);
143		stats
144	}
145
146	/// Close the transport with an explicit error, instead of waiting for the last
147	/// clone to drop. Idempotent: the first abort wins, and it cuts short a
148	/// [`close`](Self::close) still draining.
149	///
150	/// The close is executed by the session's driver, so it reaches the wire
151	/// once the runtime polls it (immediately on a live runtime).
152	pub fn abort(&self, err: Error) {
153		if let Ok(mut close) = self.close.write()
154			&& !matches!(*close, Some(Close::Abort { .. }))
155		{
156			*close = Some(Close::Abort {
157				code: SessionError::from(&err).to_code(),
158				reason: err.to_string(),
159			});
160		}
161	}
162
163	/// Close the session once the data it queued has been delivered.
164	///
165	/// Waits until every stream still serving the peer (a subscription whose track
166	/// finished, a fetch, a track info reply) has written its data and FIN and the
167	/// peer acknowledged them, then closes the transport. Returns
168	/// [`Error::Timeout`] if that takes longer than one second, closing anyway, or
169	/// the session's terminal error if it ended some other way first. A track that
170	/// is still live never finishes, so finish or abort tracks before closing.
171	///
172	/// moq-transport (IETF) sessions close without waiting.
173	pub async fn close(self) -> Result<(), Error> {
174		if let Ok(mut close) = self.close.write()
175			&& close.is_none()
176		{
177			*close = Some(Close::Drain);
178		}
179		// The drain outlives this handle, so dropping it cannot cut the drain short.
180		let closed = self.closed.clone();
181		drop(self);
182
183		match closed
184			.wait(|state| match &**state {
185				Some(ended) => Poll::Ready(ended.clone()),
186				None => Poll::Pending,
187			})
188			.await
189		{
190			Ok(ended) => ended.drain.unwrap_or(Err(ended.err)),
191			Err(kio::Closed) => Err(Error::Cancel),
192		}
193	}
194
195	/// Block until the transport session is closed, returning the reason.
196	///
197	/// A close code the peer sent is decoded through the session registry (so an auth
198	/// rejection arrives as `Error::Session(SessionError::Unauthorized)`); every peer code is
199	/// preserved as [`Error::Session`], and a close carrying no application code surfaces as
200	/// [`Error::Transport`]. See [`Error::from_transport`]. If the runtime drops
201	/// the driver instead of running it to completion, this resolves with
202	/// [`Error::Cancel`].
203	pub async fn closed(&self) -> Error {
204		match self
205			.closed
206			.wait(|state| match &**state {
207				Some(ended) => Poll::Ready(ended.err.clone()),
208				None => Poll::Pending,
209			})
210			.await
211		{
212			Ok(err) => err,
213			// The driver was dropped before it could observe the close.
214			Err(kio::Closed) => Error::Cancel,
215		}
216	}
217
218	/// Drain the peer gracefully: the handle for sending this session's single
219	/// GOAWAY.
220	///
221	/// The graceful counterpart to [`abort`](Self::abort). Send the message with
222	/// [`goaway::Producer::send`], then await [`closed`](Self::closed) to observe
223	/// the peer leaving.
224	///
225	/// Only a [`Goaway`](goaway::Goaway) carrying a [`timeout`](goaway::Goaway::timeout)
226	/// schedules a close of our own, so without one this waits for a peer that may
227	/// never leave. Set a deadline when the drain has to finish.
228	///
229	/// Available on every version. A version with no GOAWAY message (moq-lite-03
230	/// and earlier) simply carries no explanation to the peer; the deadline is the
231	/// sender's own timer either way, so the session still closes on schedule and
232	/// the caller does not branch on the negotiated version.
233	pub fn drain(&self) -> goaway::Producer {
234		self.goaway.producer()
235	}
236
237	/// Observe a GOAWAY from the peer, telling us to migrate elsewhere.
238	///
239	/// [`peek`](goaway::Consumer::peek) is the cheap synchronous check;
240	/// [`recv`](goaway::Consumer::recv) waits for one. Once a GOAWAY arrives, new
241	/// subscribe and announce-interest requests on this session are refused (both
242	/// drafts forbid opening new streams afterward); existing subscriptions keep
243	/// flowing until the session closes.
244	pub fn draining(&self) -> goaway::Consumer {
245		self.goaway.consumer()
246	}
247}
248
249impl Session {
250	pub(super) fn new<S>(
251		runtime: crate::time::Clock,
252		session: S,
253		version: Version,
254		recv_bandwidth: Option<bandwidth::Consumer>,
255		protocol: crate::driver::Protocol<S>,
256		goaway: goaway::Handle,
257	) -> (Self, crate::Driver<S>)
258	where
259		S: crate::transport::poll::Session,
260	{
261		let sample = snapshot(&session);
262
263		// Send bandwidth is version-agnostic: it depends on QUIC backend support.
264		let (send_bandwidth, send_producer) = if sample.estimated_send_rate.is_some() {
265			let producer = bandwidth::Producer::new();
266			(Some(producer.consume()), Some(producer))
267		} else {
268			(None, None)
269		};
270
271		let close = kio::Producer::new(None);
272		let closed = kio::Producer::new(None);
273		let closed_consumer = closed.consume();
274		let stats = kio::Shared::new(StatsState {
275			sample,
276			demanded: false,
277		});
278
279		let supervisor = Supervisor {
280			runtime: runtime.clone(),
281			closed_watch: session.clone(),
282			session,
283			close: Some(close.consume()),
284			closed,
285			stats: stats.clone(),
286			send_bandwidth: send_producer,
287			mode: SamplerMode::Idle,
288			drain: Drain::Idle,
289		};
290
291		let session = Self {
292			close,
293			closed: closed_consumer,
294			stats,
295			version,
296			send_bandwidth,
297			recv_bandwidth,
298			goaway: Arc::new(goaway),
299		};
300		let driver = crate::Driver::new(
301			runtime.clone(),
302			crate::driver::State {
303				protocol,
304				supervisor: Some(supervisor),
305				result: None,
306			},
307		);
308
309		(session, driver)
310	}
311}
312
313/// The driver's transport-facing half of a [`Session`]: it executes the
314/// handles' close requests, publishes the transport's terminal error, and
315/// samples the connection stats (including the send-bandwidth estimate) while
316/// anyone is consuming them.
317///
318/// Finishes once the transport reports closed; everything else is moot then.
319pub(crate) struct Supervisor<S> {
320	runtime: crate::time::Clock,
321	session: S,
322	// A dedicated clone for the close watch, since each pending poll operation
323	// needs its own handle.
324	closed_watch: S,
325	/// Handle-side close requests; `None` once one was executed (only the
326	/// first close matters, and the channel closing is the last handle
327	/// dropping).
328	close: Option<kio::Consumer<Option<Close>>>,
329	/// Where the transport's end is published for [`Session::closed`].
330	closed: kio::Producer<Option<Ended>>,
331	stats: kio::Shared<StatsState>,
332	/// The send-rate estimate channel, when the backend reports one. `None`
333	/// also once every consumer is gone for good.
334	send_bandwidth: Option<bandwidth::Producer>,
335	mode: SamplerMode,
336	drain: Drain,
337}
338
339/// A [`Session::close`] waiting for the protocol to deliver what it queued.
340enum Drain {
341	/// Nobody asked for one.
342	Idle,
343	/// Requested: close once drained, or at the deadline.
344	Waiting(crate::runtime::Deadline<crate::time::Clock>),
345	/// The drain closed the transport, with this outcome.
346	Done(Result<(), Error>),
347}
348
349enum SamplerMode {
350	/// Nobody wants stats; sampling is paused.
351	Idle,
352	/// Someone does; sample when the deadline elapses.
353	Polling {
354		deadline: crate::runtime::Deadline<crate::time::Clock>,
355	},
356}
357
358impl<S: crate::transport::poll::Session> Supervisor<S> {
359	const POLL_INTERVAL: Duration = Duration::from_millis(100);
360
361	pub(crate) fn poll(&mut self, waiter: &kio::Waiter) -> Poll<()> {
362		let mut cx = std::task::Context::from_waker(waiter.waker());
363
364		// The transport's terminal error ends the supervisor.
365		if let Poll::Ready(err) = self.closed_watch.poll_closed(&mut cx) {
366			// Nothing samples once this returns, but `stats()` keeps serving
367			// this cell, so leave it holding the session's final counters
368			// rather than whichever sample the last demand happened to catch.
369			self.stats.lock().sample = snapshot(&self.session);
370			let drain = match std::mem::replace(&mut self.drain, Drain::Idle) {
371				Drain::Done(res) => Some(res),
372				_ => None,
373			};
374			if let Ok(mut closed) = self.closed.write() {
375				*closed = Some(Ended {
376					err: Error::from_transport(err),
377					drain,
378				});
379			}
380			return Poll::Ready(());
381		}
382
383		// Execute the handle-side close requests. The channel closing is the
384		// last handle dropping, with a request written just before winning over
385		// the implicit cancel. A drain loops back to watch for an abort, which
386		// cuts it short.
387		while let Some(close) = &self.close {
388			let draining = matches!(self.drain, Drain::Waiting(_));
389			let (request, last) = match close.poll(waiter, |state| match &**state {
390				Some(Close::Drain) if draining => Poll::Pending,
391				Some(request) => Poll::Ready(request.clone()),
392				None => Poll::Pending,
393			}) {
394				Poll::Ready(Ok(request)) => (request, false),
395				Poll::Ready(Err(last)) => (
396					last.clone().unwrap_or_else(|| Close::Abort {
397						code: SessionError::Cancel.to_code(),
398						reason: "dropped".to_string(),
399					}),
400					true,
401				),
402				Poll::Pending => break,
403			};
404			match request {
405				Close::Abort { code, reason } => {
406					self.session.close(code, &reason);
407					self.drain = Drain::Idle;
408					self.close = None;
409				}
410				Close::Drain => {
411					if !draining {
412						self.drain = Drain::Waiting(crate::runtime::Deadline::after(&self.runtime, CLOSE_TIMEOUT));
413					}
414					// No handle is left to abort.
415					if last {
416						self.close = None;
417					}
418				}
419			}
420		}
421
422		self.poll_sampler(waiter);
423		Poll::Pending
424	}
425
426	/// Finish a requested drain once the protocol owes the peer nothing, or at
427	/// the deadline. Returns whether this closed the transport.
428	///
429	/// Called after the protocol ran this turn, since only then is `drained`
430	/// current.
431	pub(crate) fn poll_drain(&mut self, drained: bool, waiter: &kio::Waiter) -> bool {
432		let Drain::Waiting(deadline) = &mut self.drain else {
433			return false;
434		};
435		let res = match drained {
436			true => Ok(()),
437			false if deadline.poll(waiter).is_ready() => Err(Error::Timeout),
438			false => return false,
439		};
440		self.session.close(SessionError::Cancel.to_code(), "");
441		self.drain = Drain::Done(res);
442		// The transport is closed, so no later request can change anything.
443		self.close = None;
444		true
445	}
446
447	/// Take one sample and arm the next deadline.
448	fn sample(&mut self) {
449		let sample = snapshot(&self.session);
450		if let Some(producer) = &self.send_bandwidth {
451			// An error means every consumer is gone for good; the stats cell
452			// still wants the sample.
453			if producer.set(sample.estimated_send_rate).is_err() {
454				self.send_bandwidth = None;
455			}
456		}
457		let mut stats = self.stats.lock();
458		stats.sample = sample;
459		stats.demanded = false;
460		drop(stats);
461		self.mode = SamplerMode::Polling {
462			deadline: crate::runtime::Deadline::after(&self.runtime, Self::POLL_INTERVAL),
463		};
464	}
465
466	fn poll_sampler(&mut self, waiter: &kio::Waiter) {
467		loop {
468			match &mut self.mode {
469				SamplerMode::Idle => {
470					// Demand is a bandwidth consumer appearing or a stats read.
471					let mut demanded = match &self.send_bandwidth {
472						Some(producer) => match producer.poll_used(waiter) {
473							Poll::Ready(Ok(())) => true,
474							Poll::Ready(Err(_)) => {
475								self.send_bandwidth = None;
476								false
477							}
478							Poll::Pending => false,
479						},
480						None => false,
481					};
482					demanded |= self
483						.stats
484						.poll(waiter, |state| match state.demanded {
485							true => Poll::Ready(()),
486							false => Poll::Pending,
487						})
488						.is_ready();
489					if !demanded {
490						return;
491					}
492					self.sample();
493				}
494				SamplerMode::Polling { deadline } => {
495					if deadline.poll(waiter).is_pending() {
496						return;
497					}
498					// The interval elapsed: pause unless someone still cares.
499					let used = self.send_bandwidth.as_ref().is_some_and(bandwidth::Producer::is_used);
500					if !used && !self.stats.read().demanded {
501						self.mode = SamplerMode::Idle;
502						continue;
503					}
504					self.sample();
505					// Loop so the fresh deadline registers the waiter.
506				}
507			}
508		}
509	}
510}
511
512/// A [`Stats`] snapshot of the transport's counters.
513///
514/// `estimated_recv_rate` is filled in at the [`Session`] level (it comes from
515/// MoQ PROBE, not the transport), so it stays `None` here.
516fn snapshot<S: crate::transport::poll::Session>(session: &S) -> Stats {
517	let stats = session.stats();
518	Stats {
519		rtt: stats.rtt(),
520		estimated_send_rate: stats.estimated_send_rate().map(bandwidth::Rate::from_bps),
521		bytes_sent: stats.bytes_sent(),
522		bytes_received: stats.bytes_received(),
523		bytes_lost: stats.bytes_lost(),
524		packets_sent: stats.packets_sent(),
525		packets_received: stats.packets_received(),
526		packets_lost: stats.packets_lost(),
527		..Default::default()
528	}
529}
530
531// The point of the sever: the handle's auto-traits no longer depend on which
532// transport the runtime drives, so every consumer (moq-ffi needs Send + Sync)
533// works over every transport, pinned `!Send` ones included.
534const _: () = {
535	const fn assert_send_sync<T: Send + Sync>() {}
536	assert_send_sync::<Session>();
537};