Skip to main content

moq_net/
driver.rs

1//! Drive a session on the caller's executor.
2
3use std::task::Poll;
4
5use crate::Error;
6use crate::time::{Clock, Instant};
7
8/// Drives a session with caller-supplied time.
9///
10/// Returned by [`crate::Client::connect`] and [`crate::Server::accept`]. Call
11/// [`poll`](Self::poll) when external activity wakes the waiter or when the
12/// returned deadline is reached, supplying nondecreasing instants; see
13/// [`crate::time::Driver`] for the contract. Completion is cached, so
14/// subsequent polls return the same error.
15///
16/// It holds no session handle: dropping the last [`crate::Session`] requests
17/// closure on the next poll. Dropping the driver cancels the session, and
18/// [`crate::Session::closed`] resolves with [`Error::Cancel`]. Its `Send`-ness
19/// follows its transport.
20#[must_use = "the session makes no progress unless its driver is polled"]
21pub struct Driver<S: crate::transport::poll::Session> {
22	state: State<S>,
23	clock: Clock,
24}
25
26/// The protocol half of a machine, one variant per negotiated wire protocol.
27///
28/// The lite driver is a named machine, so the machine's `Send`-ness follows
29/// the transport (a pinned `!Send` transport yields a `!Send` machine that
30/// stays on its thread). The ietf driver is still a boxed future; the box
31/// demands `Send` on native, which is why the ietf path requires a
32/// [`Boxable`](crate::transport::poll::Boxable) transport until it too becomes
33/// a named machine.
34pub(crate) enum Protocol<S: crate::transport::poll::Session> {
35	/// Boxed for size only: a concrete box, so `Send` stays inferred.
36	Lite(Box<crate::lite::Driver<S>>),
37	Ietf(crate::ietf::Driver),
38}
39
40/// Protocol and lifecycle work owned by the driver.
41pub(crate) struct State<S: crate::transport::poll::Session> {
42	pub(crate) protocol: Protocol<S>,
43	// The session supervisor, polled alongside the protocol: it executes the
44	// handles' close requests, publishes the transport's terminal error, and
45	// samples stats. It finishes once the transport reports closed, and the
46	// machine is not done until it has: the protocol's terminal transport close
47	// is what `Session::closed` observes, so resolving before it is published
48	// would leave waiters parked on a machine nobody polls again. `None` once
49	// finished, since a completed machine must not be polled again.
50	pub(crate) supervisor: Option<crate::session::Supervisor<S>>,
51	// Cached so a poll after completion doesn't re-poll a finished protocol.
52	pub(crate) result: Option<Result<(), Error>>,
53}
54
55impl<S: crate::transport::poll::Session> Driver<S> {
56	pub(crate) fn new(clock: Clock, state: State<S>) -> Self {
57		Self { state, clock }
58	}
59
60	/// Process ready work at `now`, registering for external activity.
61	///
62	/// Panics if `now` is earlier than the previous poll or construction time.
63	pub fn poll(&mut self, now: Instant, waiter: &kio::Waiter) -> Result<Option<Instant>, Error> {
64		self.clock.advance(now);
65		self.clock.register_driver(waiter);
66		match self.state.poll(waiter) {
67			Poll::Ready(Ok(())) => Err(Error::Closed),
68			Poll::Ready(Err(err)) => Err(err),
69			Poll::Pending => Ok(self.clock.timeout()),
70		}
71	}
72}
73
74impl<S: crate::transport::poll::Session> Protocol<S> {
75	pub(crate) fn local_close(&self) -> std::sync::Arc<std::sync::atomic::AtomicBool> {
76		match self {
77			Self::Lite(driver) => driver.local_close.clone(),
78			Self::Ietf(driver) => driver.local_close.clone(),
79		}
80	}
81
82	fn poll(&mut self, waiter: &kio::Waiter) -> Poll<Result<(), Error>> {
83		match self {
84			Self::Lite(driver) => driver.poll(waiter),
85			Self::Ietf(driver) => waiter.poll_future(std::pin::Pin::new(driver)),
86		}
87	}
88
89	/// Whether the protocol owes the peer no queued data, so a draining close can
90	/// proceed. The IETF driver does not track this and closes at once.
91	fn drained(&self) -> bool {
92		match self {
93			Self::Lite(driver) => driver.drained(),
94			Self::Ietf(_) => true,
95		}
96	}
97}
98
99impl<S: crate::transport::poll::Session> State<S> {
100	fn poll(&mut self, waiter: &kio::Waiter) -> Poll<Result<(), Error>> {
101		if let Some(supervisor) = &mut self.supervisor
102			&& supervisor.poll(waiter).is_ready()
103		{
104			self.supervisor = None;
105		}
106
107		if self.result.is_none() {
108			// The protocol's last act is closing the transport, and so is a drain
109			// once the protocol that just ran owes the peer nothing. Either wakes
110			// the supervisor's close watch; poll it now instead of waiting a turn.
111			let closed = match self.protocol.poll(waiter) {
112				Poll::Ready(result) => {
113					self.result = Some(result);
114					true
115				}
116				Poll::Pending => self
117					.supervisor
118					.as_mut()
119					.is_some_and(|supervisor| supervisor.poll_drain(self.protocol.drained(), waiter)),
120			};
121			if closed
122				&& let Some(supervisor) = &mut self.supervisor
123				&& supervisor.poll(waiter).is_ready()
124			{
125				self.supervisor = None;
126			}
127		}
128
129		match (&self.result, &self.supervisor) {
130			(Some(result), None) => Poll::Ready(result.clone()),
131			_ => Poll::Pending,
132		}
133	}
134}
135
136impl<S: crate::transport::poll::Session> crate::time::Driver for Driver<S> {
137	fn poll(&mut self, now: Instant, waiter: &kio::Waiter) -> Result<Option<Instant>, Error> {
138		self.poll(now, waiter)
139	}
140}
141
142impl<S: crate::transport::poll::Session> std::fmt::Debug for Driver<S> {
143	fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
144		f.debug_struct("Driver")
145			.field("done", &self.state.result.is_some())
146			.finish()
147	}
148}