Skip to main content

tor_proto/stream/flow_ctrl/xon_xoff/
state.rs

1//! Circuit reactor's stream XON/XOFF flow control.
2//!
3//! ## Notes on consensus parameters
4//!
5//! ### `cc_xoff_client`
6//!
7//! This is the number of bytes that we buffer within a [`DataStream`]. The actual total number of
8//! bytes buffered can be *much* larger. For example there will be additional buffering:
9//!
10//! - Within the arti socks/http proxy: Arti's proxy code needs to read some bytes from the stream, store
11//!   it in a temporary buffer, then write the buffer to the socket. If the socket would block, the
12//!   data would remain in that temporary buffer. In practice arti uses only a small byte buffer (APP_STREAM_BUF_LEN) at
13//!   the time of writing, which is hopefully negligible. See `arti::socks::copy_interactive()`.
14//! - Within the kernel: There are two additional buffers that will store stream data before the
15//!   application connected over socks will see the data: Arti's socket send buffer and the
16//!   application's socket receive buffer. If the application were to stop reading from its socket,
17//!   stream data would accumulate first in the socket's receive buffer. Once full, stream data
18//!   would accumulate in arti's socket's send buffer. This can become relatively large, especially
19//!   with buffer autotuning enabled. On a Linux 6.15 system with curl downloading a large file and
20//!   stopping mid-download, the receive buffer was 6,116,738 bytes and the send buffer was
21//!   2,631,062 bytes. This sums to around 8.7 MB of stream data buffered in the kernel, which is
22//!   significantly higher than the current consensus value of `cc_xoff_client`.
23//!   NOTE: Arti's proxy sockets now use fixed-size `DEFAULT_{SEND,RECV}_BUF_SIZE` kernel buffers.
24//!
25//! This means that the total number of bytes buffered before an XOFF is sent can be much larger
26//! than `cc_xoff_client`.
27//!
28//! While we should take into account the kernel and arti socks buffering above, we also need to
29//! keep in mind that arti-client is a library that can be used by others. These library users might
30//! not do any kernel or socks buffering, for example if they write a rust program that handles the
31//! stream data entirely within their program. We don't want to set `cc_xoff_client` too low that it
32//! harms the performance for these users, even if it's fine for the arti socks proxy case.
33
34use std::num::Saturating;
35use std::sync::Arc;
36
37use postage::watch;
38use tor_cell::relaycell::flow_ctrl::{FlowCtrlVersion, Xoff, Xon, XonKBpsEwma};
39use tor_cell::relaycell::msg::AnyRelayMsg;
40use tor_cell::relaycell::{RelayCmd, RelayMsg, UnparsedRelayMsg};
41use tracing::trace;
42
43use super::reader::DrainRateRequest;
44
45use crate::stream::flow_ctrl::params::{CellCount, FlowCtrlParameters};
46use crate::stream::flow_ctrl::state::{
47    FlowCtrlHooks, HalfStreamFlowCtrlHooks, StreamRateLimit, WithSidechannelMitigations,
48};
49use crate::util::notify::NotifySender;
50use crate::{Error, Result};
51
52#[cfg(doc)]
53use {crate::client::stream::DataStream, crate::stream::flow_ctrl::state::StreamFlowCtrl};
54
55/// State for XON/XOFF flow control.
56#[derive(Debug)]
57pub(crate) struct XonXoffFlowCtrl {
58    /// Consensus parameters.
59    params: Arc<FlowCtrlParameters>,
60    /// How we communicate rate limit updates to the
61    /// [`DataWriter`](crate::client::stream::DataWriter).
62    rate_limit_updater: watch::Sender<StreamRateLimit>,
63    /// How we communicate requests for new drain rate updates to the
64    /// [`XonXoffReader`](crate::stream::flow_ctrl::xon_xoff::reader::XonXoffReader).
65    drain_rate_requester: NotifySender<DrainRateRequest>,
66    /// The last rate limit we sent.
67    last_sent_xon_xoff: Option<XonXoffMsg>,
68    /// The buffer limit at which we should send an XOFF.
69    ///
70    /// In prop324 it says that this will be either `cc_xoff_client` or `cc_xoff_exit` depending on
71    /// whether we're a client/hs or exit, but we deviate from the spec here (see how it is set
72    /// below).
73    xoff_limit: CellCount<{ tor_cell::relaycell::PAYLOAD_MAX_SIZE_ALL as u32 }>,
74    /// DropMark sidechannel mitigations.
75    ///
76    /// This is only enabled if we are a client (including an onion service).
77    //
78    // We could use a `Box` here so that this only takes up space if sidechannel mitigations are
79    // enabled. But `SidechannelMitigation` is (at the time of writing) only 16 bytes. We could
80    // reconsider in the future if we add more functionality to `SidechannelMitigation`.
81    sidechannel_mitigation: Option<SidechannelMitigation>,
82}
83
84impl XonXoffFlowCtrl {
85    /// Returns a new xon/xoff-based state.
86    pub(crate) fn new(
87        params: Arc<FlowCtrlParameters>,
88        with_sidechannel_mitigations: WithSidechannelMitigations,
89        rate_limit_updater: watch::Sender<StreamRateLimit>,
90        drain_rate_requester: NotifySender<DrainRateRequest>,
91    ) -> Self {
92        let sidechannel_mitigation = match with_sidechannel_mitigations {
93            WithSidechannelMitigations::Enabled => Some(SidechannelMitigation::new()),
94            WithSidechannelMitigations::Disabled => None,
95        };
96
97        // We use the same XOFF limit regardless of if we're a client or exit.
98        // See https://gitlab.torproject.org/tpo/core/torspec/-/issues/371#note_3260658
99        let xoff_limit = std::cmp::max(params.cc_xoff_client, params.cc_xoff_exit);
100
101        Self {
102            params,
103            rate_limit_updater,
104            drain_rate_requester,
105            last_sent_xon_xoff: None,
106            xoff_limit,
107            sidechannel_mitigation,
108        }
109    }
110}
111
112impl FlowCtrlHooks for XonXoffFlowCtrl {
113    fn can_send<M: RelayMsg>(&self, _msg: &M) -> bool {
114        // we perform rate-limiting in the `DataWriter`,
115        // so we send any messages that made it past the `DataWriter`
116        true
117    }
118
119    fn about_to_send(&mut self, msg: &AnyRelayMsg) -> Result<()> {
120        // if sidechannel mitigations are enabled and this is a RELAY_DATA message,
121        // notify that we sent a data message
122        if let Some(ref mut sidechannel_mitigation) = self.sidechannel_mitigation {
123            if let AnyRelayMsg::Data(data_msg) = msg {
124                sidechannel_mitigation.sent_stream_data(data_msg.as_ref().len());
125            }
126        }
127
128        Ok(())
129    }
130
131    fn put_for_incoming_sendme(&mut self, _msg: UnparsedRelayMsg) -> Result<()> {
132        let msg = "Stream level SENDME not allowed due to congestion control";
133        Err(Error::CircProto(msg.into()))
134    }
135
136    fn handle_incoming_xon(&mut self, msg: UnparsedRelayMsg) -> Result<()> {
137        let xon = msg
138            .decode::<Xon>()
139            .map_err(|e| Error::from_bytes_err(e, "failed to decode XON message"))?
140            .into_msg();
141
142        // > Parties SHOULD treat XON or XOFF cells with unrecognized versions as a protocol
143        // > violation.
144        if *xon.version() != 0 {
145            return Err(Error::CircProto("Unrecognized XON version".into()));
146        }
147
148        // if sidechannel mitigations are enabled, notify that an XON was received
149        if let Some(ref mut sidechannel_mitigation) = self.sidechannel_mitigation {
150            sidechannel_mitigation.received_xon(&self.params)?;
151        }
152
153        trace!("Received an XON with rate {}", xon.kbytes_per_sec_ewma());
154
155        let rate = match xon.kbytes_per_sec_ewma() {
156            XonKBpsEwma::Limited(rate_kbytes_per_sec) => {
157                let rate_kbytes_per_sec = u64::from(rate_kbytes_per_sec.get());
158                // convert from kilobytes/s to bytes/s
159                StreamRateLimit::new_bytes_per_sec(rate_kbytes_per_sec * 1000)
160            }
161            XonKBpsEwma::Unlimited => StreamRateLimit::MAX,
162        };
163
164        *self.rate_limit_updater.borrow_mut() = rate;
165        Ok(())
166    }
167
168    fn handle_incoming_xoff(&mut self, msg: UnparsedRelayMsg) -> Result<()> {
169        let xoff = msg
170            .decode::<Xoff>()
171            .map_err(|e| Error::from_bytes_err(e, "failed to decode XOFF message"))?
172            .into_msg();
173
174        // > Parties SHOULD treat XON or XOFF cells with unrecognized versions as a protocol
175        // > violation.
176        if *xoff.version() != 0 {
177            return Err(Error::CircProto("Unrecognized XOFF version".into()));
178        }
179
180        // if sidechannel mitigations are enabled, notify that an XOFF was received
181        if let Some(ref mut sidechannel_mitigation) = self.sidechannel_mitigation {
182            sidechannel_mitigation.received_xoff(&self.params)?;
183        }
184
185        trace!("Received an XOFF");
186
187        // update the rate limit and notify the `DataWriter`
188        *self.rate_limit_updater.borrow_mut() = StreamRateLimit::ZERO;
189
190        Ok(())
191    }
192
193    fn maybe_send_xon(&mut self, rate: XonKBpsEwma, buffer_len: usize) -> Result<Option<Xon>> {
194        if buffer_len as u64 > self.xoff_limit.as_bytes() {
195            // we can't send an XON, and we should have already sent an XOFF when the queue first
196            // exceeded the limit (see `maybe_send_xoff()`)
197            debug_assert!(matches!(self.last_sent_xon_xoff, Some(XonXoffMsg::Xoff)));
198
199            // inform the stream reader that we need a new drain rate
200            self.drain_rate_requester.notify();
201            return Ok(None);
202        }
203
204        self.last_sent_xon_xoff = Some(XonXoffMsg::Xon(rate));
205
206        trace!("Want to send an XON with rate {rate}");
207
208        Ok(Some(Xon::new(FlowCtrlVersion::V0, rate)))
209    }
210
211    fn maybe_send_xoff(&mut self, buffer_len: usize) -> Result<Option<Xoff>> {
212        // if the last XON/XOFF we sent was an XOFF, no need to send another
213        if matches!(self.last_sent_xon_xoff, Some(XonXoffMsg::Xoff)) {
214            return Ok(None);
215        }
216
217        if buffer_len as u64 <= self.xoff_limit.as_bytes() {
218            return Ok(None);
219        }
220
221        // either we have never sent an XOFF or XON, or we last sent an XON
222
223        // remember that we last sent an XOFF
224        self.last_sent_xon_xoff = Some(XonXoffMsg::Xoff);
225
226        // inform the stream reader that we need a new drain rate
227        self.drain_rate_requester.notify();
228
229        trace!("Want to send an XOFF");
230
231        Ok(Some(Xoff::new(FlowCtrlVersion::V0)))
232    }
233
234    fn inbound_queue_max_len(&self) -> usize {
235        // Congestion control doesn't have an upper limit for the number of in-flight
236        // cells that the other end might send,
237        // so we need to expect any number of cells on this stream.
238        //
239        // Since dealing with mpsc queues that may be bounded or unbounded is a pain (requires a
240        // bunch of enum wrappers), we'll set a very high bound.
241        // This bound should be high enough that we'll never reach it in practice
242        // (and if we do, it's surely a bug or an attack),
243        // but not too high as to cause `futures_channel::mpsc::channel()` to panic.
244        //
245        // Here we choose a max of 2_000_000 messages,
246        // which is approx 1000 MB of stream data (assuming packed cells).
247        //
248        // TODO(arti#2540): We should use an unbounded queue for XON/XOFF flow control,
249        // and should return `None` here.
250        2_000_000
251    }
252}
253
254/// State for XON/XOFF flow control on a half-stream.
255#[derive(Debug)]
256pub(crate) struct HalfStreamXonXoffFlowCtrl {
257    /// The original [`XonXoffFlowCtrl`] from the full stream.
258    ///
259    /// We keep this since we need to continue validating any incoming messages
260    /// and continue applying the sidechannel mitigations.
261    inner: XonXoffFlowCtrl,
262}
263
264impl HalfStreamXonXoffFlowCtrl {
265    /// Returns a new xon/xoff-based state for a half-stream.
266    pub(crate) fn new(flow_ctrl: XonXoffFlowCtrl) -> Self {
267        Self { inner: flow_ctrl }
268    }
269}
270
271impl HalfStreamFlowCtrlHooks for HalfStreamXonXoffFlowCtrl {
272    fn handle_incoming_dropped(&mut self, _msg_count: u16) -> Result<()> {
273        // Nothing to do here.
274        Ok(())
275    }
276
277    fn handle_incoming_msg(&mut self, msg: UnparsedRelayMsg) -> Result<Option<UnparsedRelayMsg>> {
278        match msg.cmd() {
279            RelayCmd::SENDME => {
280                self.inner.put_for_incoming_sendme(msg)?;
281                Ok(None)
282            }
283            RelayCmd::XON => {
284                self.inner.handle_incoming_xon(msg)?;
285                Ok(None)
286            }
287            RelayCmd::XOFF => {
288                self.inner.handle_incoming_xoff(msg)?;
289                Ok(None)
290            }
291            // Nothing to do here.
292            _ => Ok(Some(msg)),
293        }
294    }
295}
296
297/// An XON or XOFF message with no associated data.
298#[derive(Debug, PartialEq, Eq)]
299enum XonXoff {
300    /// XON message.
301    Xon,
302    /// XOFF message.
303    Xoff,
304}
305
306/// An XON or XOFF message with associated data.
307#[derive(Debug)]
308enum XonXoffMsg {
309    /// XON message with a rate.
310    // TODO: I'm expecting that we'll want the `XonKBpsEwma` in the future.
311    // If that doesn't end up being the case, then we should remove it.
312    #[expect(dead_code)]
313    Xon(XonKBpsEwma),
314    /// XOFF message.
315    Xoff,
316}
317
318/// Sidechannel mitigations for DropMark attacks.
319///
320/// > In order to mitigate DropMark attacks, both XOFF and advisory XON transmission must be
321/// > restricted.
322///
323/// These restrictions should be implemented for clients (OPs and onion services).
324#[derive(Debug)]
325struct SidechannelMitigation {
326    /// The last rate limit update we received.
327    last_recvd_xon_xoff: Option<XonXoff>,
328    /// Number of sent stream bytes.
329    ///
330    /// C-tor has some logic to try to fit this into a 32-bit integer,
331    /// but lets not do that unless we need to as it will make bugs more likely.
332    bytes_sent_total: Saturating<u64>,
333    /// The number of advisory XON messages we've received.
334    ///
335    /// Note: Advisory XONs are XON->XON messages, and not XOFF->XON messages.
336    num_advisory_xon_recvd: Saturating<u64>,
337    /// The number of XOFF messages we've received.
338    num_xoff_recvd: Saturating<u64>,
339}
340
341impl SidechannelMitigation {
342    /// A new [`SidechannelMitigation`].
343    fn new() -> Self {
344        Self {
345            last_recvd_xon_xoff: None,
346            bytes_sent_total: Saturating(0),
347            num_advisory_xon_recvd: Saturating(0),
348            num_xoff_recvd: Saturating(0),
349        }
350    }
351
352    /// A (likely underestimated) guess of the XOFF limit that the other endpoint is using.
353    fn peer_xoff_limit_bytes(params: &FlowCtrlParameters) -> u64 {
354        // We need to consider that `xoff_client` and `xoff_exit` may be different, we don't know
355        // here exactly what kind of peer we're connected to, and that we may have a different view
356        // of the consensus than the peer.
357        // We deviate from prop324 here and use a more relaxed threshold.
358        // See https://gitlab.torproject.org/tpo/core/torspec/-/issues/371#note_3260658
359        let min = std::cmp::min(
360            params.cc_xoff_client.as_bytes(),
361            params.cc_xoff_exit.as_bytes(),
362        );
363        min / 2
364    }
365
366    /// A (likely underestimated) guess of the advisory XON limit that the other endpoint is using.
367    fn peer_xon_limit_bytes(params: &FlowCtrlParameters) -> u64 {
368        // We need to consider that we may have a different view of the consensus than the peer.
369        // We deviate from prop324 here and use a more relaxed threshold.
370        // See https://gitlab.torproject.org/tpo/core/torspec/-/issues/371#note_3260658
371        params.cc_xon_rate.as_bytes() / 2
372    }
373
374    /// Notify that we have sent stream data.
375    fn sent_stream_data(&mut self, stream_bytes: usize) {
376        // perform a saturating conversion to u64
377        let stream_bytes: u64 = stream_bytes.try_into().unwrap_or(u64::MAX);
378        self.bytes_sent_total += stream_bytes;
379    }
380
381    /// Notify that we have received an XON message.
382    fn received_xon(&mut self, params: &FlowCtrlParameters) -> Result<()> {
383        // Check to make sure that XON is not sent too early, for dropmark attacks. The main
384        // sidechannel risk is early cells, but we also check to see that we did not get more XONs
385        // than make sense for the number of bytes we sent.
386        //
387        // The ordering is important here. For example we first want to check if we received an
388        // advisory XON that was too early, before we check if we received the advisory XON too
389        // frequently.
390
391        // Ensure that we have sent some bytes. This might be covered by other checks below, but this
392        // is the most important check so we do it explicitly here first.
393        if self.bytes_sent_total.0 == 0 {
394            const MSG: &str = "Received XON before sending any data";
395            return Err(Error::CircProto(MSG.into()));
396        }
397
398        // is this an advisory XON?
399        let is_advisory = match self.last_recvd_xon_xoff {
400            // if we last received an XON, then this is advisory since we are already sending data
401            Some(XonXoff::Xon) => true,
402            // if we last received an XOFF, then this isn't advisory since we're being asked to
403            // resume sending data
404            Some(XonXoff::Xoff) => false,
405            // if we never received an XON nor XOFF, then this is advisory since we are already
406            // sending data
407            None => true,
408        };
409
410        // set this before we possibly return early below, since this must be set regardless of if
411        // it's an advisory XON or not
412        self.last_recvd_xon_xoff = Some(XonXoff::Xon);
413
414        // we only restrict advisory XON messages
415        if !is_advisory {
416            return Ok(());
417        }
418
419        self.num_advisory_xon_recvd += 1;
420
421        // > Clients also SHOULD ensure that advisory XONs do not arrive before the minimum of the
422        // > XOFF limit and 'cc_xon_rate' full cells worth of bytes have been transmitted.
423        //
424        // NOTE: We use a more relaxed threshold for the XON and XOFF limits than in prop324.
425        let advisory_not_expected_before = std::cmp::min(
426            Self::peer_xoff_limit_bytes(params),
427            Self::peer_xon_limit_bytes(params),
428        );
429        if self.bytes_sent_total.0 < advisory_not_expected_before {
430            const MSG: &str = "Received advisory XON too early";
431            return Err(Error::CircProto(MSG.into()));
432        }
433
434        // > Clients SHOULD ensure that advisory XONs do not arrive more frequently than every
435        // > 'cc_xon_rate' cells worth of sent data.
436        //
437        // It should be an error if
438        //   XON frequency > 1/peer_xon_limit_bytes
439        // where
440        //   XON frequency = num_advisory_xon_recvd/bytes_sent_total
441        //
442        // so
443        //   num_advisory_xon_recvd/bytes_sent_total > 1/peer_xon_limit_bytes
444        //
445        // or to better work with integers
446        //   num_advisory_xon_recvd > bytes_sent_total/peer_xon_limit_bytes
447        //
448        // NOTE: We use a more relaxed threshold for the XON limit than in prop324.
449        let peer_xon_limit_bytes = Self::peer_xon_limit_bytes(params);
450        if peer_xon_limit_bytes != 0
451            && self.num_advisory_xon_recvd.0 > self.bytes_sent_total.0 / peer_xon_limit_bytes
452        {
453            const MSG: &str = "Received advisory XON too frequently";
454            return Err(Error::CircProto(MSG.into()));
455        }
456
457        Ok(())
458    }
459
460    /// Notify that we have received an XOFF message.
461    fn received_xoff(&mut self, params: &FlowCtrlParameters) -> Result<()> {
462        // Check to make sure that XOFF is not sent too early, for dropmark attacks. The
463        // main sidechannel risk is early cells, but we also check to make sure that we have not
464        // received more XOFFs than could have been generated by the bytes we sent.
465        //
466        // The ordering is important here. For example we first want to disallow consecutive XOFFs,
467        // then check if we received an XOFF that was too early, and finally check if we received
468        // the XOFF too frequently.
469
470        self.num_xoff_recvd += 1;
471
472        // Ensure that we have sent some bytes. This might be covered by other checks below, but this
473        // is the most important check so we do it explicitly here first.
474        if self.bytes_sent_total.0 == 0 {
475            const MSG: &str = "Received XOFF before sending any data";
476            return Err(Error::CircProto(MSG.into()));
477        }
478
479        // disallow consecutive XOFF messages
480        if self.last_recvd_xon_xoff == Some(XonXoff::Xoff) {
481            const MSG: &str = "Received consecutive XOFF messages";
482            return Err(Error::CircProto(MSG.into()));
483        }
484
485        // > clients MUST ensure that an XOFF does not arrive before it has sent the appropriate
486        // > XOFF limit of bytes on a stream ('cc_xoff_exit' for exits, 'cc_xoff_client' for
487        // > onions).
488        //
489        // NOTE: We use a more relaxed threshold for the XOFF limit than in prop324.
490        if self.bytes_sent_total.0 < Self::peer_xoff_limit_bytes(params) {
491            const MSG: &str = "Received XOFF too early";
492            return Err(Error::CircProto(MSG.into()));
493        }
494
495        // > Clients also SHOULD ensure than XOFFs do not arrive more frequently than every XOFF
496        // > limit worth of sent data.
497        //
498        // It should be an error if
499        //   XOFF frequency > 1/peer_xoff_limit_bytes
500        // where
501        //   XOFF frequency = num_xoff_recvd/bytes_sent_total
502        //
503        // so
504        //   num_xoff_recvd/bytes_sent_total > 1/peer_xoff_limit_bytes
505        //
506        // or to better work with integers
507        //   num_xoff_recvd > bytes_sent_total/peer_xoff_limit_bytes
508        //
509        // NOTE: We use a more relaxed threshold for the XOFF limit than in prop324.
510        let peer_xoff_limit_bytes = Self::peer_xoff_limit_bytes(params);
511        if peer_xoff_limit_bytes != 0
512            && self.num_xoff_recvd.0 > self.bytes_sent_total.0 / peer_xoff_limit_bytes
513        {
514            return Err(Error::CircProto("Received XOFF too frequently".into()));
515        }
516
517        self.last_recvd_xon_xoff = Some(XonXoff::Xoff);
518
519        Ok(())
520    }
521}
522
523#[cfg(test)]
524mod test {
525    use super::*;
526
527    use crate::stream::flow_ctrl::params::CellCount;
528
529    #[test]
530    fn sidechannel_mitigation() {
531        let params = [
532            FlowCtrlParameters {
533                cc_xoff_client: CellCount::new(2),
534                cc_xoff_exit: CellCount::new(4),
535                cc_xon_rate: CellCount::new(8),
536                cc_xon_change_pct: 1,
537                cc_xon_ewma_cnt: 1,
538            },
539            FlowCtrlParameters {
540                cc_xoff_client: CellCount::new(8),
541                cc_xoff_exit: CellCount::new(4),
542                cc_xon_rate: CellCount::new(2),
543                cc_xon_change_pct: 1,
544                cc_xon_ewma_cnt: 1,
545            },
546        ];
547
548        for params in params {
549            let xon_limit = SidechannelMitigation::peer_xon_limit_bytes(&params);
550            let xoff_limit = SidechannelMitigation::peer_xoff_limit_bytes(&params);
551
552            let mut x = SidechannelMitigation::new();
553            // cannot receive XON as first message
554            assert!(x.received_xon(&params).is_err());
555
556            let mut x = SidechannelMitigation::new();
557            // cannot receive XOFF as first message
558            assert!(x.received_xoff(&params).is_err());
559
560            let mut x = SidechannelMitigation::new();
561            // cannot receive XOFF after sending fewer than `xoff_limit` bytes
562            x.sent_stream_data(xoff_limit as usize - 1);
563            assert!(x.received_xoff(&params).is_err());
564
565            let mut x = SidechannelMitigation::new();
566            // can receive XOFF after sending `xoff_limit` bytes
567            x.sent_stream_data(xoff_limit as usize);
568            assert!(x.received_xoff(&params).is_ok());
569            // but cannot receive another XOFF immediately after
570            assert!(x.received_xoff(&params).is_err());
571
572            let mut x = SidechannelMitigation::new();
573            // can receive XOFF after sending `xoff_limit` bytes
574            x.sent_stream_data(xoff_limit as usize);
575            assert!(x.received_xoff(&params).is_ok());
576            // but cannot receive another XOFF even after sending another `xoff_limit` bytes
577            x.sent_stream_data(xoff_limit as usize);
578            assert!(x.received_xoff(&params).is_err());
579
580            let mut x = SidechannelMitigation::new();
581            // can receive XOFF after sending `xoff_limit` bytes
582            x.sent_stream_data(xoff_limit as usize);
583            assert!(x.received_xoff(&params).is_ok());
584            // and can immediately receive an XON
585            assert!(x.received_xon(&params).is_ok());
586            // and can receive another XOFF after sending another `xoff_limit` bytes
587            x.sent_stream_data(xoff_limit as usize);
588            assert!(x.received_xoff(&params).is_ok());
589
590            let mut x = SidechannelMitigation::new();
591            // cannot receive XON after sending fewer than `xon_limit` bytes
592            x.sent_stream_data(xon_limit as usize - 1);
593            assert!(x.received_xon(&params).is_err());
594
595            let mut x = SidechannelMitigation::new();
596            // can receive XON after sending a large number of bytes
597            x.sent_stream_data(xon_limit as usize * 3);
598            assert!(x.received_xon(&params).is_ok());
599            // and can immediately receive another XON
600            assert!(x.received_xon(&params).is_ok());
601            // and can immediately receive another XON
602            assert!(x.received_xon(&params).is_ok());
603            // but cannot receive another XON immediately after
604            assert!(x.received_xon(&params).is_err());
605
606            let mut x = SidechannelMitigation::new();
607            // can receive XOFF after sending a large number of bytes
608            x.sent_stream_data(xoff_limit as usize * 3);
609            assert!(x.received_xoff(&params).is_ok());
610            // and can immediately receive an XON
611            assert!(x.received_xon(&params).is_ok());
612            // and can immediately receive an XOFF
613            assert!(x.received_xoff(&params).is_ok());
614            // and can immediately receive an XON
615            assert!(x.received_xon(&params).is_ok());
616            // and can immediately receive an XOFF
617            assert!(x.received_xoff(&params).is_ok());
618            // and can immediately receive an XON
619            assert!(x.received_xon(&params).is_ok());
620            // but cannot immediately receive an XOFF
621            assert!(x.received_xoff(&params).is_err());
622        }
623    }
624}