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};