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			local_close: protocol.local_close(),
281			runtime: runtime.clone(),
282			closed_watch: session.clone(),
283			session,
284			close: Some(close.consume()),
285			closed,
286			stats: stats.clone(),
287			send_bandwidth: send_producer,
288			mode: SamplerMode::Idle,
289			drain: Drain::Idle,
290		};
291
292		let session = Self {
293			close,
294			closed: closed_consumer,
295			stats,
296			version,
297			send_bandwidth,
298			recv_bandwidth,
299			goaway: Arc::new(goaway),
300		};
301		let driver = crate::Driver::new(
302			runtime.clone(),
303			crate::driver::State {
304				protocol,
305				supervisor: Some(supervisor),
306				result: None,
307			},
308		);
309
310		(session, driver)
311	}
312}
313
314/// The driver's transport-facing half of a [`Session`]: it executes the
315/// handles' close requests, publishes the transport's terminal error, and
316/// samples the connection stats (including the send-bandwidth estimate) while
317/// anyone is consuming them.
318///
319/// Finishes once the transport reports closed; everything else is moot then.
320pub(crate) struct Supervisor<S> {
321	local_close: Arc<std::sync::atomic::AtomicBool>,
322	runtime: crate::time::Clock,
323	session: S,
324	// A dedicated clone for the close watch, since each pending poll operation
325	// needs its own handle.
326	closed_watch: S,
327	/// Handle-side close requests; `None` once one was executed (only the
328	/// first close matters, and the channel closing is the last handle
329	/// dropping).
330	close: Option<kio::Consumer<Option<Close>>>,
331	/// Where the transport's end is published for [`Session::closed`].
332	closed: kio::Producer<Option<Ended>>,
333	stats: kio::Shared<StatsState>,
334	/// The send-rate estimate channel, when the backend reports one. `None`
335	/// also once every consumer is gone for good.
336	send_bandwidth: Option<bandwidth::Producer>,
337	mode: SamplerMode,
338	drain: Drain,
339}
340
341/// A [`Session::close`] waiting for the protocol to deliver what it queued.
342enum Drain {
343	/// Nobody asked for one.
344	Idle,
345	/// Requested: close once drained, or at the deadline.
346	Waiting(crate::runtime::Deadline<crate::time::Clock>),
347	/// The drain closed the transport, with this outcome.
348	Done(Result<(), Error>),
349}
350
351enum SamplerMode {
352	/// Nobody wants stats; sampling is paused.
353	Idle,
354	/// Someone does; sample when the deadline elapses.
355	Polling {
356		deadline: crate::runtime::Deadline<crate::time::Clock>,
357	},
358}
359
360impl<S: crate::transport::poll::Session> Supervisor<S> {
361	const POLL_INTERVAL: Duration = Duration::from_millis(100);
362
363	pub(crate) fn poll(&mut self, waiter: &kio::Waiter) -> Poll<()> {
364		let mut cx = std::task::Context::from_waker(waiter.waker());
365
366		// The transport's terminal error ends the supervisor.
367		if let Poll::Ready(err) = self.closed_watch.poll_closed(&mut cx) {
368			// Nothing samples once this returns, but `stats()` keeps serving
369			// this cell, so leave it holding the session's final counters
370			// rather than whichever sample the last demand happened to catch.
371			self.stats.lock().sample = snapshot(&self.session);
372			let drain = match std::mem::replace(&mut self.drain, Drain::Idle) {
373				Drain::Done(res) => Some(res),
374				_ => None,
375			};
376			if let Ok(mut closed) = self.closed.write() {
377				*closed = Some(Ended {
378					err: Error::from_transport(err),
379					drain,
380				});
381			}
382			return Poll::Ready(());
383		}
384
385		// Execute the handle-side close requests. The channel closing is the
386		// last handle dropping, with a request written just before winning over
387		// the implicit cancel. A drain loops back to watch for an abort, which
388		// cuts it short.
389		while let Some(close) = &self.close {
390			let draining = matches!(self.drain, Drain::Waiting(_));
391			let (request, last) = match close.poll(waiter, |state| match &**state {
392				Some(Close::Drain) if draining => Poll::Pending,
393				Some(request) => Poll::Ready(request.clone()),
394				None => Poll::Pending,
395			}) {
396				Poll::Ready(Ok(request)) => (request, false),
397				Poll::Ready(Err(last)) => (
398					last.clone().unwrap_or_else(|| {
399						self.local_close.store(true, std::sync::atomic::Ordering::Relaxed);
400						Close::Abort {
401							code: SessionError::Cancel.to_code(),
402							reason: "dropped".to_string(),
403						}
404					}),
405					true,
406				),
407				Poll::Pending => break,
408			};
409			match request {
410				Close::Abort { code, reason } => {
411					self.session.close(code, &reason);
412					self.drain = Drain::Idle;
413					self.close = None;
414				}
415				Close::Drain => {
416					if !draining {
417						self.drain = Drain::Waiting(crate::runtime::Deadline::after(&self.runtime, CLOSE_TIMEOUT));
418					}
419					// No handle is left to abort.
420					if last {
421						self.close = None;
422					}
423				}
424			}
425		}
426
427		self.poll_sampler(waiter);
428		Poll::Pending
429	}
430
431	/// Finish a requested drain once the protocol owes the peer nothing, or at
432	/// the deadline. Returns whether this closed the transport.
433	///
434	/// Called after the protocol ran this turn, since only then is `drained`
435	/// current.
436	pub(crate) fn poll_drain(&mut self, drained: bool, waiter: &kio::Waiter) -> bool {
437		let Drain::Waiting(deadline) = &mut self.drain else {
438			return false;
439		};
440		let res = match drained {
441			true => Ok(()),
442			false if deadline.poll(waiter).is_ready() => Err(Error::Timeout),
443			false => return false,
444		};
445		self.local_close.store(true, std::sync::atomic::Ordering::Relaxed);
446		self.session.close(SessionError::Cancel.to_code(), "");
447		self.drain = Drain::Done(res);
448		// The transport is closed, so no later request can change anything.
449		self.close = None;
450		true
451	}
452
453	/// Take one sample and arm the next deadline.
454	fn sample(&mut self) {
455		let sample = snapshot(&self.session);
456		if let Some(producer) = &self.send_bandwidth {
457			// An error means every consumer is gone for good; the stats cell
458			// still wants the sample.
459			if producer.set(sample.estimated_send_rate).is_err() {
460				self.send_bandwidth = None;
461			}
462		}
463		let mut stats = self.stats.lock();
464		stats.sample = sample;
465		stats.demanded = false;
466		drop(stats);
467		self.mode = SamplerMode::Polling {
468			deadline: crate::runtime::Deadline::after(&self.runtime, Self::POLL_INTERVAL),
469		};
470	}
471
472	fn poll_sampler(&mut self, waiter: &kio::Waiter) {
473		loop {
474			match &mut self.mode {
475				SamplerMode::Idle => {
476					// Demand is a bandwidth consumer appearing or a stats read.
477					let mut demanded = match &self.send_bandwidth {
478						Some(producer) => match producer.poll_used(waiter) {
479							Poll::Ready(Ok(())) => true,
480							Poll::Ready(Err(_)) => {
481								self.send_bandwidth = None;
482								false
483							}
484							Poll::Pending => false,
485						},
486						None => false,
487					};
488					demanded |= self
489						.stats
490						.poll(waiter, |state| match state.demanded {
491							true => Poll::Ready(()),
492							false => Poll::Pending,
493						})
494						.is_ready();
495					if !demanded {
496						return;
497					}
498					self.sample();
499				}
500				SamplerMode::Polling { deadline } => {
501					if deadline.poll(waiter).is_pending() {
502						return;
503					}
504					// The interval elapsed: pause unless someone still cares.
505					let used = self.send_bandwidth.as_ref().is_some_and(bandwidth::Producer::is_used);
506					if !used && !self.stats.read().demanded {
507						self.mode = SamplerMode::Idle;
508						continue;
509					}
510					self.sample();
511					// Loop so the fresh deadline registers the waiter.
512				}
513			}
514		}
515	}
516}
517
518/// A [`Stats`] snapshot of the transport's counters.
519///
520/// `estimated_recv_rate` is filled in at the [`Session`] level (it comes from
521/// MoQ PROBE, not the transport), so it stays `None` here.
522fn snapshot<S: crate::transport::poll::Session>(session: &S) -> Stats {
523	let stats = session.stats();
524	Stats {
525		rtt: stats.rtt(),
526		estimated_send_rate: stats.estimated_send_rate().map(bandwidth::Rate::from_bps),
527		bytes_sent: stats.bytes_sent(),
528		bytes_received: stats.bytes_received(),
529		bytes_lost: stats.bytes_lost(),
530		packets_sent: stats.packets_sent(),
531		packets_received: stats.packets_received(),
532		packets_lost: stats.packets_lost(),
533		..Default::default()
534	}
535}
536
537// The point of the sever: the handle's auto-traits no longer depend on which
538// transport the runtime drives, so every consumer (moq-ffi needs Send + Sync)
539// works over every transport, pinned `!Send` ones included.
540const _: () = {
541	const fn assert_send_sync<T: Send + Sync>() {}
542	assert_send_sync::<Session>();
543};