1use std::task::Poll;
4
5use crate::Error;
6use crate::time::{Clock, Instant};
7
8#[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
26pub(crate) enum Protocol<S: crate::transport::poll::Session> {
35 Lite(Box<crate::lite::Driver<S>>),
37 Ietf(crate::ietf::Driver),
38}
39
40pub(crate) struct State<S: crate::transport::poll::Session> {
42 pub(crate) protocol: Protocol<S>,
43 pub(crate) supervisor: Option<crate::session::Supervisor<S>>,
51 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 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 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 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}