tor_proto/stream/flow_ctrl/state.rs
1//! Code for implementing flow control (stream-level).
2
3use std::sync::Arc;
4
5use postage::watch;
6use tor_cell::relaycell::flow_ctrl::{Xoff, Xon, XonKBpsEwma};
7use tor_cell::relaycell::msg::AnyRelayMsg;
8use tor_cell::relaycell::{RelayMsg, UnparsedRelayMsg};
9
10use super::params::FlowCtrlParameters;
11use super::window::state::{HalfStreamWindowFlowCtrl, WindowFlowCtrl};
12use super::xon_xoff::reader::DrainRateRequest;
13use super::xon_xoff::state::{HalfStreamXonXoffFlowCtrl, XonXoffFlowCtrl};
14
15use crate::Result;
16use crate::congestion::sendme;
17use crate::util::notify::NotifySender;
18
19/// Private internals of [`StreamFlowCtrl`].
20#[enum_dispatch::enum_dispatch]
21#[derive(Debug)]
22enum StreamFlowCtrlInner {
23 /// "legacy" sendme-window-based flow control.
24 Window(WindowFlowCtrl),
25 /// XON/XOFF flow control.
26 XonXoff(XonXoffFlowCtrl),
27}
28
29/// Manages the circuit reactor's flow control for a stream.
30///
31/// Note that the flow control logic can be distributed across multiple parts of Arti.
32/// For example some flow control logic will exist in the circuit reactor,
33/// but other logic will exist in the stream's `DataStream`.
34/// So this doesn't include all flow control logic.
35#[derive(Debug)]
36pub(crate) struct StreamFlowCtrl {
37 /// Private internal enum.
38 inner: StreamFlowCtrlInner,
39}
40
41impl StreamFlowCtrl {
42 /// Returns a new sendme-window-based [`StreamFlowCtrl`].
43 pub(crate) fn new_window(window: sendme::StreamSendWindow) -> Self {
44 Self {
45 inner: StreamFlowCtrlInner::Window(WindowFlowCtrl::new(window)),
46 }
47 }
48
49 /// Returns a new xon/xoff-based [`StreamFlowCtrl`].
50 pub(crate) fn new_xon_xoff(
51 params: Arc<FlowCtrlParameters>,
52 with_sidechannel_mitigations: WithSidechannelMitigations,
53 rate_limit_updater: watch::Sender<StreamRateLimit>,
54 drain_rate_requester: NotifySender<DrainRateRequest>,
55 ) -> Self {
56 Self {
57 inner: StreamFlowCtrlInner::XonXoff(XonXoffFlowCtrl::new(
58 params,
59 with_sidechannel_mitigations,
60 rate_limit_updater,
61 drain_rate_requester,
62 )),
63 }
64 }
65
66 /// Once this end of the stream is closed and the stream becomes a
67 /// half-stream (`HalfStream`),
68 /// this method will turn the flow control object into a version
69 /// that is designed to be used for half-streams.
70 pub(crate) fn half_stream(self) -> HalfStreamFlowCtrl {
71 let inner = match self.inner {
72 StreamFlowCtrlInner::Window(x) => {
73 HalfStreamFlowCtrlInner::Window(HalfStreamWindowFlowCtrl::new(x))
74 }
75 StreamFlowCtrlInner::XonXoff(x) => {
76 HalfStreamFlowCtrlInner::XonXoff(HalfStreamXonXoffFlowCtrl::new(x))
77 }
78 };
79
80 HalfStreamFlowCtrl { inner }
81 }
82}
83
84// forward all trait methods to the inner enum
85impl FlowCtrlHooks for StreamFlowCtrl {
86 fn can_send<M: RelayMsg>(&self, msg: &M) -> bool {
87 self.inner.can_send(msg)
88 }
89
90 fn about_to_send(&mut self, msg: &AnyRelayMsg) -> Result<()> {
91 self.inner.about_to_send(msg)
92 }
93
94 fn put_for_incoming_sendme(&mut self, msg: UnparsedRelayMsg) -> Result<()> {
95 self.inner.put_for_incoming_sendme(msg)
96 }
97
98 fn handle_incoming_xon(&mut self, msg: UnparsedRelayMsg) -> Result<()> {
99 self.inner.handle_incoming_xon(msg)
100 }
101
102 fn handle_incoming_xoff(&mut self, msg: UnparsedRelayMsg) -> Result<()> {
103 self.inner.handle_incoming_xoff(msg)
104 }
105
106 fn maybe_send_xon(&mut self, rate: XonKBpsEwma, buffer_len: usize) -> Result<Option<Xon>> {
107 self.inner.maybe_send_xon(rate, buffer_len)
108 }
109
110 fn maybe_send_xoff(&mut self, buffer_len: usize) -> Result<Option<Xoff>> {
111 self.inner.maybe_send_xoff(buffer_len)
112 }
113
114 fn inbound_queue_max_len(&self) -> usize {
115 self.inner.inbound_queue_max_len()
116 }
117}
118
119/// Methods that can be called on a [`StreamFlowCtrl`].
120///
121/// We use a trait so that we can use `enum_dispatch` on the inner [`StreamFlowCtrlInner`] enum.
122#[enum_dispatch::enum_dispatch(StreamFlowCtrlInner)]
123pub(crate) trait FlowCtrlHooks {
124 /// Whether this stream is ready to send `msg`.
125 fn can_send<M: RelayMsg>(&self, msg: &M) -> bool;
126
127 /// Inform the flow control code that we're about to send `msg`.
128 /// Returns an error if the message should not be sent,
129 /// and the circuit should be closed.
130 // TODO: Consider having this method wrap the message in a type that
131 // "proves" we've applied flow control. This would make it easier to apply
132 // flow control earlier, e.g. in `OpenStreamEntStream`, without introducing
133 // ambiguity in the sending function as to whether flow control has already
134 // been applied or not.
135 fn about_to_send(&mut self, msg: &AnyRelayMsg) -> Result<()>;
136
137 /// Handle an incoming sendme.
138 ///
139 /// On success, return the number of cells left in the window.
140 ///
141 /// On failure, return an error: the caller should close the stream or
142 /// circuit with a protocol error.
143 ///
144 /// Takes the [`UnparsedRelayMsg`] so that we don't even try to decode it if we're not using the
145 /// correct type of flow control.
146 fn put_for_incoming_sendme(&mut self, msg: UnparsedRelayMsg) -> Result<()>;
147
148 /// Handle an incoming XON message.
149 ///
150 /// Takes the [`UnparsedRelayMsg`] so that we don't even try to decode it if we're not using the
151 /// correct type of flow control.
152 fn handle_incoming_xon(&mut self, msg: UnparsedRelayMsg) -> Result<()>;
153
154 /// Handle an incoming XOFF message.
155 ///
156 /// Takes the [`UnparsedRelayMsg`] so that we don't even try to decode it if we're not using the
157 /// correct type of flow control.
158 fn handle_incoming_xoff(&mut self, msg: UnparsedRelayMsg) -> Result<()>;
159
160 /// Check if we should send an XON message.
161 ///
162 /// If we should, then returns the XON message that should be sent.
163 /// Returns an error if XON/XOFF messages aren't supported for this type of flow control.
164 fn maybe_send_xon(&mut self, rate: XonKBpsEwma, buffer_len: usize) -> Result<Option<Xon>>;
165
166 /// Check if we should send an XOFF message.
167 ///
168 /// If we should, then returns the XOFF message that should be sent.
169 /// Returns an error if XON/XOFF messages aren't supported for this type of flow control.
170 fn maybe_send_xoff(&mut self, buffer_len: usize) -> Result<Option<Xoff>>;
171
172 /// The max queue length that should be used for stream messages incoming from the Tor network.
173 ///
174 /// This is the queue length between the user-facing stream reader (`DataReader`)
175 /// and the circuit reactor.
176 ///
177 /// If the queue would ever exceed this many messages, the stream should be closed.
178 fn inbound_queue_max_len(&self) -> usize;
179}
180
181/// Manages flow control for a half-stream (`HalfStream`).
182#[derive(Debug)]
183pub(crate) struct HalfStreamFlowCtrl {
184 /// Private internal enum.
185 inner: HalfStreamFlowCtrlInner,
186}
187
188/// Private internals of [`HalfStreamFlowCtrl`].
189#[enum_dispatch::enum_dispatch]
190#[derive(Debug)]
191enum HalfStreamFlowCtrlInner {
192 /// "legacy" sendme-window-based flow control.
193 Window(HalfStreamWindowFlowCtrl),
194 /// XON/XOFF flow control.
195 XonXoff(HalfStreamXonXoffFlowCtrl),
196}
197
198/// Methods that can be called on a [`HalfStreamFlowCtrl`].
199///
200/// We use a trait so that we can use `enum_dispatch` on the inner [`HalfStreamFlowCtrlInner`] enum.
201/// While this may seem unnecessary since this trait currently only has two methods,
202/// it's consistent with the [`FlowCtrlHooks`] trait above.
203#[enum_dispatch::enum_dispatch(HalfStreamFlowCtrlInner)]
204pub(crate) trait HalfStreamFlowCtrlHooks {
205 /// Handle some number of dropped stream messages.
206 ///
207 /// We don't know what kinds of stream messages were dropped, only the number of them.
208 ///
209 /// This method exists because currently the stream entry may drop some incoming stream
210 /// messages and they would never be processed by this flow control object otherwise.
211 fn handle_incoming_dropped(&mut self, msg_count: u16) -> Result<()>;
212
213 /// Handle an incoming message.
214 ///
215 /// If it's a flow control message, it will be consumed and `None` will be returned.
216 /// Otherwise the original message will be returned.
217 ///
218 /// Takes the [`UnparsedRelayMsg`] so that we don't even try to decode it if we're not using the
219 /// correct type of flow control.
220 fn handle_incoming_msg(&mut self, msg: UnparsedRelayMsg) -> Result<Option<UnparsedRelayMsg>>;
221}
222
223// forward all trait methods to the inner enum
224impl HalfStreamFlowCtrlHooks for HalfStreamFlowCtrl {
225 fn handle_incoming_dropped(&mut self, msg_count: u16) -> Result<()> {
226 self.inner.handle_incoming_dropped(msg_count)
227 }
228
229 fn handle_incoming_msg(&mut self, msg: UnparsedRelayMsg) -> Result<Option<UnparsedRelayMsg>> {
230 self.inner.handle_incoming_msg(msg)
231 }
232}
233
234/// A newtype wrapper for a tor stream rate limit that makes the units explicit.
235#[derive(Copy, Clone, Debug, PartialEq, Eq, PartialOrd, Ord)]
236pub(crate) struct StreamRateLimit {
237 /// The rate in bytes/s.
238 rate: u64,
239}
240
241impl StreamRateLimit {
242 /// A maximum rate limit.
243 pub(crate) const MAX: Self = Self::new_bytes_per_sec(u64::MAX);
244
245 /// A rate limit of 0.
246 pub(crate) const ZERO: Self = Self::new_bytes_per_sec(0);
247
248 /// A new [`StreamRateLimit`] with `rate` bytes/s.
249 pub(crate) const fn new_bytes_per_sec(rate: u64) -> Self {
250 Self { rate }
251 }
252
253 /// The rate in bytes/s.
254 pub(crate) const fn bytes_per_sec(&self) -> u64 {
255 self.rate
256 }
257}
258
259/// Whether sidechannel mitigations are enabled or not for flow control.
260#[derive(Copy, Clone, Debug, PartialEq, Eq)]
261pub(crate) enum WithSidechannelMitigations {
262 /// Flow control sidechannel mitigations are *enabled*.
263 ///
264 /// Should be enabled for clients (including onion services).
265 Enabled,
266 /// Flow control sidechannel mitigations are *disabled*.
267 ///
268 /// Should be disabled for exits.
269 Disabled,
270}
271
272impl std::fmt::Display for StreamRateLimit {
273 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
274 write!(f, "{} bytes/s", self.rate)
275 }
276}