Skip to main content

tor_proto/client/stream/
data.rs

1//! Declare DataStream, a type that wraps DataReader and DataWriter so as to be useful
2//! for byte-oriented communication.
3
4use crate::{Error, Result};
5use static_assertions::assert_impl_all;
6use tor_cell::relaycell::msg::EndReason;
7use tor_cell::relaycell::{RelayCellFormat, RelayCmd};
8
9use futures::io::{AsyncRead, AsyncWrite};
10use futures::stream::StreamExt;
11use futures::task::{Context, Poll};
12use futures::{Future, Stream};
13use pin_project::pin_project;
14use postage::watch;
15
16#[cfg(feature = "tokio")]
17use tokio_crate::io::ReadBuf;
18#[cfg(feature = "tokio")]
19use tokio_crate::io::{AsyncRead as TokioAsyncRead, AsyncWrite as TokioAsyncWrite};
20#[cfg(feature = "tokio")]
21use tokio_util::compat::{FuturesAsyncReadCompatExt, FuturesAsyncWriteCompatExt};
22use tor_cell::restricted_msg;
23
24use std::fmt::Debug;
25use std::io::Result as IoResult;
26use std::num::NonZero;
27use std::pin::Pin;
28#[cfg(any(feature = "stream-ctrl", feature = "experimental-api"))]
29use std::sync::Arc;
30#[cfg(feature = "stream-ctrl")]
31use std::sync::{Mutex, Weak};
32
33use educe::Educe;
34
35use crate::client::ClientTunnel;
36use crate::memquota::StreamAccount;
37use crate::stream::StreamReceiver;
38use crate::stream::StreamTarget;
39use crate::stream::Tunnel;
40use crate::stream::cmdcheck::{AnyCmdChecker, CmdChecker, StreamStatus};
41use crate::stream::flow_ctrl::state::StreamRateLimit;
42use crate::stream::flow_ctrl::xon_xoff::reader::{BufferIsEmpty, XonXoffReader, XonXoffReaderCtrl};
43use tor_async_utils::rate_limited_writer::{
44    DynamicRateLimitedWriter, RateLimitedWriter, RateLimitedWriterConfig,
45};
46use tor_basic_utils::onionperf_types::{OnionperfEvent, OnionperfStreamStatus};
47use tor_basic_utils::skip_fmt;
48use tor_cell::relaycell::msg::Data;
49use tor_error::internal;
50use tor_rtcompat::{CoarseTimeProvider, DynTimeProvider, SleepProvider};
51
52/// A stream of [`RateLimitedWriterConfig`] used to update a [`DynamicRateLimitedWriter`].
53///
54/// Unfortunately we need to store the result of a [`StreamExt::map`] and [`StreamExt::fuse`] in
55/// [`DataWriter`], which leaves us with this ugly type.
56/// We use a type alias to make `DataWriter` a little nicer.
57type RateConfigStream = futures::stream::Map<
58    futures::stream::Fuse<watch::Receiver<StreamRateLimit>>,
59    fn(StreamRateLimit) -> RateLimitedWriterConfig,
60>;
61
62/// An anonymized stream over the Tor network.
63///
64/// For most purposes, you can think of this type as an anonymized
65/// TCP stream: it can read and write data, and get closed when it's done.
66///
67/// [`DataStream`] implements [`futures::io::AsyncRead`] and
68/// [`futures::io::AsyncWrite`], so you can use it anywhere that those
69/// traits are expected.
70///
71/// # Examples
72///
73/// Connecting to an HTTP server and sending a request, using
74/// [`AsyncWriteExt::write_all`](futures::io::AsyncWriteExt::write_all):
75///
76/// ```ignore
77/// let mut stream = tor_client.connect(("icanhazip.com", 80), None).await?;
78///
79/// use futures::io::AsyncWriteExt;
80///
81/// stream
82///     .write_all(b"GET / HTTP/1.1\r\nHost: icanhazip.com\r\nConnection: close\r\n\r\n")
83///     .await?;
84///
85/// // Flushing the stream is important; see below!
86/// stream.flush().await?;
87/// ```
88///
89/// Reading the result, using [`AsyncReadExt::read_to_end`](futures::io::AsyncReadExt::read_to_end):
90///
91/// ```ignore
92/// use futures::io::AsyncReadExt;
93///
94/// let mut buf = Vec::new();
95/// stream.read_to_end(&mut buf).await?;
96///
97/// println!("{}", String::from_utf8_lossy(&buf));
98/// ```
99///
100/// # Usage with Tokio
101///
102/// If the `tokio` crate feature is enabled, this type also implements
103/// [`tokio::io::AsyncRead`](tokio_crate::io::AsyncRead) and
104/// [`tokio::io::AsyncWrite`](tokio_crate::io::AsyncWrite) for easier integration
105/// with code that expects those traits.
106///
107/// # Remember to call `flush`!
108///
109/// DataStream buffers data internally, in order to write as few cells
110/// as possible onto the network.  In order to make sure that your
111/// data has actually been sent, you need to make sure that
112/// [`AsyncWrite::poll_flush`] runs to completion: probably via
113/// [`AsyncWriteExt::flush`](futures::io::AsyncWriteExt::flush).
114///
115/// # Splitting the type
116///
117/// This type is internally composed of a [`DataReader`] and a [`DataWriter`]; the
118/// `DataStream::split` method can be used to split it into those two parts, for more
119/// convenient usage with e.g. stream combinators.
120///
121/// # How long does a stream live?
122///
123/// A `DataStream` will live until all references to it are dropped,
124/// or until it is closed explicitly.
125///
126/// If you split the stream into a `DataReader` and a `DataWriter`, it
127/// will survive until _both_ are dropped, or until it is closed
128/// explicitly.
129///
130/// A stream can also close because of a network error,
131/// or because the other side of the stream decided to close it.
132///
133// # Semver note
134//
135// Note that this type is re-exported as a part of the public API of
136// the `arti-client` crate.  Any changes to its API here in
137// `tor-proto` need to be reflected above.
138#[derive(Debug)]
139pub struct DataStream {
140    /// Underlying writer for this stream
141    w: DataWriter,
142    /// Underlying reader for this stream
143    r: DataReader,
144    /// A control object that can be used to monitor and control this stream
145    /// without needing to own it.
146    ///
147    /// Set to `None` if this is not a client stream.
148    #[cfg(feature = "stream-ctrl")]
149    ctrl: Option<Arc<ClientDataStreamCtrl>>,
150}
151assert_impl_all! { DataStream: Send, Sync }
152
153/// An object used to control and monitor a data stream.
154///
155/// # Notes
156///
157/// This is a separate type from [`DataStream`] because it's useful to have
158/// multiple references to this object, whereas a [`DataReader`] and [`DataWriter`]
159/// need to have a single owner for the `AsyncRead` and `AsyncWrite` APIs to
160/// work correctly.
161#[cfg(feature = "stream-ctrl")]
162#[cfg_attr(
163    feature = "rpc",
164    derive(derive_deftly::Deftly),
165    derive_deftly(tor_rpcbase::templates::Object)
166)]
167#[derive(Debug)]
168pub struct ClientDataStreamCtrl {
169    /// The circuit to which this stream is attached.
170    ///
171    /// Note that the stream's reader and writer halves each contain a `StreamTarget`,
172    /// which in turn has a strong reference to the `ClientCirc`.  So as long as any
173    /// one of those is alive, this reference will be present.
174    ///
175    /// We make this a Weak reference so that once the stream itself is closed,
176    /// we can't leak circuits.
177    tunnel: Weak<ClientTunnel>,
178
179    /// Shared user-visible information about the state of this stream.
180    ///
181    /// TODO RPC: This will probably want to be a `postage::Watch` or something
182    /// similar, if and when it stops moving around.
183    #[cfg(feature = "stream-ctrl")]
184    status: Arc<Mutex<DataStreamStatus>>,
185
186    /// The memory quota account that should be used for this stream's data
187    ///
188    /// Exists to keep the account alive
189    _memquota: StreamAccount,
190}
191
192/// The inner writer for [`DataWriter`].
193///
194/// This type is responsible for taking bytes and packaging them into cells.
195/// Rate limiting is implemented in [`DataWriter`] to avoid making this type more complex.
196#[derive(Debug)]
197struct DataWriterInner {
198    /// Internal state for this writer
199    ///
200    /// This is stored in an Option so that we can mutate it in the
201    /// AsyncWrite functions.  It might be possible to do better here,
202    /// and we should refactor if so.
203    state: Option<DataWriterState>,
204
205    /// The memory quota account that should be used for this stream's data
206    ///
207    /// Exists to keep the account alive
208    // If we liked, we could make this conditional; see DataReaderInner.memquota
209    _memquota: StreamAccount,
210
211    /// A control object that can be used to monitor and control this stream
212    /// without needing to own it.
213    ///
214    /// Set to `None` if this is not a client stream.
215    #[cfg(feature = "stream-ctrl")]
216    ctrl: Option<Arc<ClientDataStreamCtrl>>,
217}
218
219/// The write half of a [`DataStream`], implementing [`futures::io::AsyncWrite`].
220///
221/// See the [`DataStream`] docs for more information. In particular, note
222/// that this writer requires `poll_flush` to complete in order to guarantee that
223/// all data has been written.
224///
225/// # Usage with Tokio
226///
227/// If the `tokio` crate feature is enabled, this type also implements
228/// [`tokio::io::AsyncWrite`](tokio_crate::io::AsyncWrite) for easier integration
229/// with code that expects that trait.
230///
231/// # Drop and close
232///
233/// Note that dropping a `DataWriter` has no special effect on its own:
234/// if the `DataWriter` is dropped, the underlying stream will still remain open
235/// until the `DataReader` is also dropped.
236///
237/// If you want the stream to close earlier, use [`close`](futures::io::AsyncWriteExt::close)
238/// (or [`shutdown`](tokio_crate::io::AsyncWriteExt::shutdown) with `tokio`).
239///
240/// Remember that Tor does not support half-open streams:
241/// If you `close` or `shutdown` a stream,
242/// the other side will not see the stream as half-open,
243/// and so will (probably) not finish sending you any in-progress data.
244/// Do not use `close`/`shutdown` to communicate anything besides
245/// "I am done using this stream."
246///
247// # Semver note
248//
249// Note that this type is re-exported as a part of the public API of
250// the `arti-client` crate.  Any changes to its API here in
251// `tor-proto` need to be reflected above.
252#[derive(Debug)]
253pub struct DataWriter {
254    /// A wrapper around [`DataWriterInner`] that adds rate limiting.
255    writer: DynamicRateLimitedWriter<DataWriterInner, RateConfigStream, DynTimeProvider>,
256}
257
258impl DataWriter {
259    /// Create a new rate-limited [`DataWriter`] from a [`DataWriterInner`].
260    fn new(
261        inner: DataWriterInner,
262        rate_limit_updates: watch::Receiver<StreamRateLimit>,
263        time_provider: DynTimeProvider,
264    ) -> Self {
265        /// Converts a `rate` into a `RateLimitedWriterConfig`.
266        fn rate_to_config(rate: StreamRateLimit) -> RateLimitedWriterConfig {
267            let rate = rate.bytes_per_sec();
268            RateLimitedWriterConfig {
269                rate,        // bytes per second
270                burst: rate, // bytes
271                // This number is chosen arbitrarily, but the idea is that we want to balance
272                // between throughput and latency. Assume the user tries to write a large buffer
273                // (~600 bytes). If we set this too small (for example 1), we'll be waking up
274                // frequently and writing a small number of bytes each time to the
275                // `DataWriterInner`, even if this isn't enough bytes to send a cell. If we set this
276                // too large (for example 510), we'll be waking up infrequently to write a larger
277                // number of bytes each time. So even if the `DataWriterInner` has almost a full
278                // cell's worth of data queued (for example 490) and only needs 509-490=19 more
279                // bytes before a cell can be sent, it will block until the rate limiter allows 510
280                // more bytes.
281                //
282                // TODO(arti#2028): Is there an optimal value here?
283                wake_when_bytes_available: NonZero::new(200).expect("200 != 0"), // bytes
284            }
285        }
286
287        // get the current rate from the `watch::Receiver`, which we'll use as the initial rate
288        let initial_rate: StreamRateLimit = *rate_limit_updates.borrow();
289
290        // map the rate update stream to the type required by `DynamicRateLimitedWriter`
291        let rate_limit_updates = rate_limit_updates.fuse().map(rate_to_config as fn(_) -> _);
292
293        // build the rate limiter
294        let writer = RateLimitedWriter::new(inner, &rate_to_config(initial_rate), time_provider);
295        let writer = DynamicRateLimitedWriter::new(writer, rate_limit_updates);
296
297        Self { writer }
298    }
299
300    /// Return a [`ClientDataStreamCtrl`] object that can be used to monitor and
301    /// interact with this stream without holding the stream itself.
302    ///
303    /// Returns `None` if this is not a client stream.
304    #[cfg(feature = "stream-ctrl")]
305    pub fn client_stream_ctrl(&self) -> Option<&Arc<ClientDataStreamCtrl>> {
306        self.writer.inner().client_stream_ctrl()
307    }
308}
309
310impl AsyncWrite for DataWriter {
311    fn poll_write(
312        mut self: Pin<&mut Self>,
313        cx: &mut Context<'_>,
314        buf: &[u8],
315    ) -> Poll<IoResult<usize>> {
316        AsyncWrite::poll_write(Pin::new(&mut self.writer), cx, buf)
317    }
318
319    fn poll_flush(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<IoResult<()>> {
320        AsyncWrite::poll_flush(Pin::new(&mut self.writer), cx)
321    }
322
323    fn poll_close(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<IoResult<()>> {
324        AsyncWrite::poll_close(Pin::new(&mut self.writer), cx)
325    }
326}
327
328#[cfg(feature = "tokio")]
329impl TokioAsyncWrite for DataWriter {
330    fn poll_write(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &[u8]) -> Poll<IoResult<usize>> {
331        TokioAsyncWrite::poll_write(Pin::new(&mut self.compat_write()), cx, buf)
332    }
333
334    fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<IoResult<()>> {
335        TokioAsyncWrite::poll_flush(Pin::new(&mut self.compat_write()), cx)
336    }
337
338    fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<IoResult<()>> {
339        TokioAsyncWrite::poll_shutdown(Pin::new(&mut self.compat_write()), cx)
340    }
341}
342
343/// The read half of a [`DataStream`], implementing [`futures::io::AsyncRead`].
344///
345/// See the [`DataStream`] docs for more information.
346///
347/// # Usage with Tokio
348///
349/// If the `tokio` crate feature is enabled, this type also implements
350/// [`tokio::io::AsyncRead`](tokio_crate::io::AsyncRead) for easier integration
351/// with code that expects that trait.
352//
353// # Semver note
354//
355// Note that this type is re-exported as a part of the public API of
356// the `arti-client` crate.  Any changes to its API here in
357// `tor-proto` need to be reflected above.
358#[derive(Debug)]
359pub struct DataReader {
360    /// The [`DataReaderInner`] with a wrapper to support XON/XOFF flow control.
361    reader: XonXoffReader<DataReaderInner>,
362}
363
364impl DataReader {
365    /// Create a new [`DataReader`].
366    fn new(reader: DataReaderInner, xon_xoff_reader_ctrl: XonXoffReaderCtrl) -> Self {
367        Self {
368            reader: XonXoffReader::new(xon_xoff_reader_ctrl, reader),
369        }
370    }
371
372    /// Return a [`ClientDataStreamCtrl`] object that can be used to monitor and
373    /// interact with this stream without holding the stream itself.
374    ///
375    /// Returns `None` if this is not a client stream.
376    #[cfg(feature = "stream-ctrl")]
377    pub fn client_stream_ctrl(&self) -> Option<&Arc<ClientDataStreamCtrl>> {
378        self.reader.inner().client_stream_ctrl()
379    }
380}
381
382impl AsyncRead for DataReader {
383    fn poll_read(
384        mut self: Pin<&mut Self>,
385        cx: &mut Context<'_>,
386        buf: &mut [u8],
387    ) -> Poll<IoResult<usize>> {
388        AsyncRead::poll_read(Pin::new(&mut self.reader), cx, buf)
389    }
390
391    fn poll_read_vectored(
392        mut self: Pin<&mut Self>,
393        cx: &mut Context<'_>,
394        bufs: &mut [std::io::IoSliceMut<'_>],
395    ) -> Poll<IoResult<usize>> {
396        AsyncRead::poll_read_vectored(Pin::new(&mut self.reader), cx, bufs)
397    }
398}
399
400#[cfg(feature = "tokio")]
401impl TokioAsyncRead for DataReader {
402    fn poll_read(
403        self: Pin<&mut Self>,
404        cx: &mut Context<'_>,
405        buf: &mut ReadBuf<'_>,
406    ) -> Poll<IoResult<()>> {
407        TokioAsyncRead::poll_read(Pin::new(&mut self.compat()), cx, buf)
408    }
409}
410
411/// The inner reader for [`DataReader`].
412///
413/// This type is responsible for taking stream messages and extracting the stream data from them.
414/// Flow control logic is implemented in [`DataReader`] to avoid making this type more complex.
415#[derive(Debug)]
416pub(crate) struct DataReaderInner {
417    /// Internal state for this reader.
418    ///
419    /// This is stored in an Option so that we can mutate it in
420    /// poll_read().  It might be possible to do better here, and we
421    /// should refactor if so.
422    state: Option<DataReaderState>,
423
424    /// The memory quota account that should be used for this stream's data
425    ///
426    /// Exists to keep the account alive
427    // If we liked, we could make this conditional on not(cfg(feature = "stream-ctrl"))
428    // since, ClientDataStreamCtrl contains a StreamAccount clone too.  But that seems fragile.
429    _memquota: StreamAccount,
430
431    /// A control object that can be used to monitor and control this stream
432    /// without needing to own it.
433    ///
434    /// Set to `None` if this is not a client stream.
435    #[cfg(feature = "stream-ctrl")]
436    ctrl: Option<Arc<ClientDataStreamCtrl>>,
437}
438
439impl BufferIsEmpty for DataReaderInner {
440    /// The result will become stale,
441    /// so is most accurate immediately after a [`poll_read`](AsyncRead::poll_read).
442    fn is_empty(mut self: Pin<&mut Self>) -> bool {
443        match self
444            .state
445            .as_mut()
446            .expect("forgot to put `DataReaderState` back")
447        {
448            DataReaderState::Open(imp) => {
449                // check if the partial cell in `pending` is empty,
450                // and if the message stream is empty
451                imp.pending[imp.offset..].is_empty() && imp.s.is_empty()
452            }
453            // closed, so any data should have been discarded
454            DataReaderState::Closed => true,
455        }
456    }
457}
458
459/// Shared status flags for tracking the status of as `DataStream`.
460///
461/// We expect to refactor this a bit, so it's not exposed at all.
462//
463// TODO RPC: Possibly instead of manipulating the fields of DataStreamStatus
464// from various points in this module, we should instead construct
465// DataStreamStatus as needed from information available elsewhere.  In any
466// case, we should really  eliminate as much duplicate state here as we can.
467// (See discussions at !1198 for some challenges with this.)
468#[cfg(feature = "stream-ctrl")]
469#[derive(Clone, Debug, Default)]
470struct DataStreamStatus {
471    /// True if we've received a CONNECTED message.
472    //
473    // TODO: This is redundant with `connected` in DataReaderImpl.
474    received_connected: bool,
475    /// True if we have decided to send an END message.
476    //
477    // TODO RPC: There is not an easy way to set this from this module!  Really,
478    // the decision to send an "end" is made when the StreamTarget object is
479    // dropped, but we don't currently have any way to see when that happens.
480    // Perhaps we need a different shared StreamStatus object that the
481    // StreamTarget holds?
482    sent_end: bool,
483    /// True if we have received an END message telling us to close the stream.
484    received_end: bool,
485    /// True if we have received an error.
486    ///
487    /// (This is not a subset or superset of received_end; some errors are END
488    /// messages but some aren't; some END messages are errors but some aren't.)
489    received_err: bool,
490}
491
492#[cfg(feature = "stream-ctrl")]
493impl DataStreamStatus {
494    /// Remember that we've received a connected message.
495    fn record_connected(&mut self) {
496        self.received_connected = true;
497    }
498
499    /// Remember that we've received an error of some kind.
500    fn record_error(&mut self, e: &Error) {
501        // TODO: Probably we should remember the actual error in a box or
502        // something.  But that means making a redundant copy of the error
503        // even if nobody will want it.  Do we care?
504        match e {
505            Error::EndReceived(EndReason::DONE) => self.received_end = true,
506            Error::EndReceived(_) => {
507                self.received_end = true;
508                self.received_err = true;
509            }
510            _ => self.received_err = true,
511        }
512    }
513}
514
515restricted_msg! {
516    /// An allowable incoming message on a client data stream.
517    enum ClientDataStreamMsg:RelayMsg {
518        // SENDME is handled by the reactor.
519        Data, End, Connected,
520    }
521}
522
523// TODO RPC: Should we also implement this trait for everything that holds a
524// ClientDataStreamCtrl?
525#[cfg(feature = "stream-ctrl")]
526impl super::ctrl::ClientStreamCtrl for ClientDataStreamCtrl {
527    fn tunnel(&self) -> Option<Arc<ClientTunnel>> {
528        self.tunnel.upgrade()
529    }
530}
531
532#[cfg(feature = "stream-ctrl")]
533impl ClientDataStreamCtrl {
534    /// Return true if the underlying stream is connected. (That is, if it has
535    /// received a `CONNECTED` message, and has not been closed.)
536    pub fn is_connected(&self) -> bool {
537        let s = self.status.lock().expect("poisoned lock");
538        s.received_connected && !(s.sent_end || s.received_end || s.received_err)
539    }
540
541    // TODO RPC: Add more functions once we have the desired API more nailed
542    // down.
543}
544
545impl DataStream {
546    /// Wrap raw stream receiver and target parts as a DataStream.
547    ///
548    /// For non-optimistic stream, function `wait_for_connection`
549    /// must be called after to make sure CONNECTED is received.
550    pub(crate) fn new<P: SleepProvider + CoarseTimeProvider>(
551        time_provider: P,
552        receiver: StreamReceiver,
553        xon_xoff_reader_ctrl: XonXoffReaderCtrl,
554        target: StreamTarget,
555        memquota: StreamAccount,
556    ) -> Self {
557        Self::new_inner(
558            time_provider,
559            receiver,
560            xon_xoff_reader_ctrl,
561            target,
562            false,
563            memquota,
564        )
565    }
566
567    /// Wrap raw stream receiver and target parts as a connected DataStream.
568    ///
569    /// Unlike [`DataStream::new`], this creates a `DataStream` that does not expect to receive a
570    /// CONNECTED cell.
571    ///
572    /// This is used by hidden services, exit relays, and directory servers to accept streams.
573    #[cfg(any(feature = "hs-service", feature = "relay"))]
574    pub(crate) fn new_connected<P: SleepProvider + CoarseTimeProvider>(
575        time_provider: P,
576        receiver: StreamReceiver,
577        xon_xoff_reader_ctrl: XonXoffReaderCtrl,
578        target: StreamTarget,
579        memquota: StreamAccount,
580    ) -> Self {
581        Self::new_inner(
582            time_provider,
583            receiver,
584            xon_xoff_reader_ctrl,
585            target,
586            true,
587            memquota,
588        )
589    }
590
591    /// The shared implementation of the `new*()` functions.
592    fn new_inner<P: SleepProvider + CoarseTimeProvider>(
593        time_provider: P,
594        receiver: StreamReceiver,
595        xon_xoff_reader_ctrl: XonXoffReaderCtrl,
596        target: StreamTarget,
597        connected: bool,
598        memquota: StreamAccount,
599    ) -> Self {
600        let relay_cell_format = target.relay_cell_format();
601        let out_buf_len = Data::max_body_len(relay_cell_format);
602        let rate_limit_stream = target.rate_limit_stream().clone();
603
604        tracing::trace!(
605            onionperf = true,
606            event = ?OnionperfEvent::Stream(OnionperfStreamStatus::New),
607            stream_id = ?target.stream_id,
608            circ_id = ?match target.clone().tunnel {
609                Tunnel::Client(client) => Some(client.circ.unique_id()),
610                #[cfg(feature = "relay")]
611                Tunnel::Relay(_) => None, // TODO
612            },
613        );
614
615        #[cfg(feature = "stream-ctrl")]
616        let status = {
617            let mut data_stream_status = DataStreamStatus::default();
618            if connected {
619                data_stream_status.record_connected();
620            }
621            Arc::new(Mutex::new(data_stream_status))
622        };
623
624        #[cfg(feature = "stream-ctrl")]
625        let ctrl = {
626            let tunnel = match target.tunnel() {
627                crate::stream::Tunnel::Client(t) => Some(Arc::downgrade(t)),
628                #[cfg(feature = "relay")]
629                crate::stream::Tunnel::Relay(_) => None,
630            };
631
632            tunnel.map(|tunnel| {
633                Arc::new(ClientDataStreamCtrl {
634                    tunnel,
635                    status: status.clone(),
636                    _memquota: memquota.clone(),
637                })
638            })
639        };
640        let r = DataReaderInner {
641            state: Some(DataReaderState::Open(DataReaderImpl {
642                s: receiver,
643                pending: Vec::new(),
644                offset: 0,
645                connected,
646                #[cfg(feature = "stream-ctrl")]
647                status: status.clone(),
648            })),
649            _memquota: memquota.clone(),
650            #[cfg(feature = "stream-ctrl")]
651            ctrl: ctrl.clone(),
652        };
653        let w = DataWriterInner {
654            state: Some(DataWriterState::Ready(DataWriterImpl {
655                s: target,
656                buf: vec![0; out_buf_len].into_boxed_slice(),
657                n_pending: 0,
658                #[cfg(feature = "stream-ctrl")]
659                status,
660                relay_cell_format,
661            })),
662            _memquota: memquota,
663            #[cfg(feature = "stream-ctrl")]
664            ctrl: ctrl.clone(),
665        };
666
667        let time_provider = DynTimeProvider::new(time_provider);
668
669        DataStream {
670            w: DataWriter::new(w, rate_limit_stream, time_provider),
671            r: DataReader::new(r, xon_xoff_reader_ctrl),
672            #[cfg(feature = "stream-ctrl")]
673            ctrl,
674        }
675    }
676
677    /// Divide this DataStream into its constituent parts.
678    pub fn split(self) -> (DataReader, DataWriter) {
679        (self.r, self.w)
680    }
681
682    /// Wait until a CONNECTED cell is received, or some other cell
683    /// is received to indicate an error.
684    ///
685    /// Does nothing if this stream is already connected.
686    pub async fn wait_for_connection(&mut self) -> Result<()> {
687        // We must put state back before returning
688        let state = self
689            .r
690            .reader
691            .inner_mut()
692            .state
693            .take()
694            .expect("Missing state in DataReaderInner");
695
696        if let DataReaderState::Open(mut imp) = state {
697            let result = if imp.connected {
698                Ok(())
699            } else {
700                // This succeeds if the cell is CONNECTED, and fails otherwise.
701                std::future::poll_fn(|cx| Pin::new(&mut imp).read_cell(cx)).await
702            };
703            self.r.reader.inner_mut().state = Some(match result {
704                Err(_) => DataReaderState::Closed,
705                Ok(_) => DataReaderState::Open(imp),
706            });
707            result
708        } else {
709            Err(Error::from(internal!(
710                "Expected ready state, got {:?}",
711                state
712            )))
713        }
714    }
715
716    /// Return a [`ClientDataStreamCtrl`] object that can be used to monitor and
717    /// interact with this stream without holding the stream itself.
718    #[cfg(feature = "stream-ctrl")]
719    pub fn client_stream_ctrl(&self) -> Option<&Arc<ClientDataStreamCtrl>> {
720        self.ctrl.as_ref()
721    }
722}
723
724impl AsyncRead for DataStream {
725    fn poll_read(
726        mut self: Pin<&mut Self>,
727        cx: &mut Context<'_>,
728        buf: &mut [u8],
729    ) -> Poll<IoResult<usize>> {
730        AsyncRead::poll_read(Pin::new(&mut self.r), cx, buf)
731    }
732}
733
734#[cfg(feature = "tokio")]
735impl TokioAsyncRead for DataStream {
736    fn poll_read(
737        self: Pin<&mut Self>,
738        cx: &mut Context<'_>,
739        buf: &mut ReadBuf<'_>,
740    ) -> Poll<IoResult<()>> {
741        TokioAsyncRead::poll_read(Pin::new(&mut self.compat()), cx, buf)
742    }
743}
744
745impl AsyncWrite for DataStream {
746    fn poll_write(
747        mut self: Pin<&mut Self>,
748        cx: &mut Context<'_>,
749        buf: &[u8],
750    ) -> Poll<IoResult<usize>> {
751        AsyncWrite::poll_write(Pin::new(&mut self.w), cx, buf)
752    }
753    fn poll_flush(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<IoResult<()>> {
754        AsyncWrite::poll_flush(Pin::new(&mut self.w), cx)
755    }
756    fn poll_close(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<IoResult<()>> {
757        AsyncWrite::poll_close(Pin::new(&mut self.w), cx)
758    }
759}
760
761#[cfg(feature = "tokio")]
762impl TokioAsyncWrite for DataStream {
763    fn poll_write(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &[u8]) -> Poll<IoResult<usize>> {
764        TokioAsyncWrite::poll_write(Pin::new(&mut self.compat()), cx, buf)
765    }
766
767    fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<IoResult<()>> {
768        TokioAsyncWrite::poll_flush(Pin::new(&mut self.compat()), cx)
769    }
770
771    fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<IoResult<()>> {
772        TokioAsyncWrite::poll_shutdown(Pin::new(&mut self.compat()), cx)
773    }
774}
775
776/// Helper type: Like BoxFuture, but also requires that the future be Sync.
777type BoxSyncFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + Sync + 'a>>;
778
779/// An enumeration for the state of a DataWriter.
780///
781/// We have to use an enum here because, for as long as we're waiting
782/// for a flush operation to complete, the future returned by
783/// `flush_cell()` owns the DataWriterImpl.
784#[derive(Educe)]
785#[educe(Debug)]
786enum DataWriterState {
787    /// The writer has closed or gotten an error: nothing more to do.
788    Closed,
789    /// The writer is not currently flushing; more data can get queued
790    /// immediately.
791    Ready(DataWriterImpl),
792    /// The writer is flushing a cell.
793    Flushing(
794        #[educe(Debug(method = "skip_fmt"))] //
795        BoxSyncFuture<'static, (DataWriterImpl, Result<()>)>,
796    ),
797}
798
799/// Internal: the write part of a DataStream
800#[derive(Educe)]
801#[educe(Debug)]
802struct DataWriterImpl {
803    /// The underlying StreamTarget object.
804    s: StreamTarget,
805
806    /// Buffered data to send over the connection.
807    ///
808    /// This buffer is currently allocated using a number of bytes
809    /// equal to the maximum that we can package at a time.
810    //
811    // TODO: this buffer is probably smaller than we want, but it's good
812    // enough for now.  If we _do_ make it bigger, we'll have to change
813    // our use of Data::split_from to handle the case where we can't fit
814    // all the data.
815    #[educe(Debug(method = "skip_fmt"))]
816    buf: Box<[u8]>,
817
818    /// Number of unflushed bytes in buf.
819    n_pending: usize,
820
821    /// Relay cell format in use
822    relay_cell_format: RelayCellFormat,
823
824    /// Shared user-visible information about the state of this stream.
825    #[cfg(feature = "stream-ctrl")]
826    status: Arc<Mutex<DataStreamStatus>>,
827}
828
829impl DataWriterInner {
830    /// See [`DataWriter::client_stream_ctrl`].
831    #[cfg(feature = "stream-ctrl")]
832    fn client_stream_ctrl(&self) -> Option<&Arc<ClientDataStreamCtrl>> {
833        self.ctrl.as_ref()
834    }
835
836    /// Helper for poll_flush() and poll_close(): Performs a flush, then
837    /// closes the stream if should_close is true.
838    fn poll_flush_impl(
839        mut self: Pin<&mut Self>,
840        cx: &mut Context<'_>,
841        should_close: bool,
842    ) -> Poll<IoResult<()>> {
843        let state = self.state.take().expect("Missing state in DataWriter");
844
845        // TODO: this whole function is a bit copy-pasted.
846        let mut future: BoxSyncFuture<_> = match state {
847            DataWriterState::Ready(imp) => {
848                if imp.n_pending == 0 {
849                    // Nothing to flush!
850                    if should_close {
851                        // We need to actually continue with this function to do the closing.
852                        // Thus, make a future that does nothing and is ready immediately.
853                        Box::pin(futures::future::ready((imp, Ok(()))))
854                    } else {
855                        // There's nothing more to do; we can return.
856                        self.state = Some(DataWriterState::Ready(imp));
857                        return Poll::Ready(Ok(()));
858                    }
859                } else {
860                    // We need to flush the buffer's contents; Make a future for that.
861                    Box::pin(imp.flush_buf())
862                }
863            }
864            DataWriterState::Flushing(fut) => fut,
865            DataWriterState::Closed => {
866                self.state = Some(DataWriterState::Closed);
867                return Poll::Ready(Err(Error::NotConnected.into()));
868            }
869        };
870
871        match future.as_mut().poll(cx) {
872            Poll::Ready((imp, Err(e))) => {
873                match e {
874                    Error::NotConnected => (),
875                    _ => tracing::trace!(
876                        onionperf = true,
877                        event = ?OnionperfEvent::Stream(OnionperfStreamStatus::Failed),
878                        stream_id = ?imp.s.stream_id,
879                        reason = ?e,
880                    ),
881                }
882                self.state = Some(DataWriterState::Closed);
883                Poll::Ready(Err(e.into()))
884            }
885            Poll::Ready((mut imp, Ok(()))) => {
886                if should_close {
887                    // Tell the StreamTarget to close, so that the reactor
888                    // realizes that we are done sending. (Dropping `imp.s` does not
889                    // suffice, since there may be other clones of it.  In particular,
890                    // the StreamReceiver has one, which it uses to keep the stream
891                    // open, among other things.)
892                    imp.s.close();
893
894                    #[cfg(feature = "stream-ctrl")]
895                    {
896                        // TODO RPC:  This is not sufficient to track every case
897                        // where we might have sent an End.  See note on the
898                        // `sent_end` field.
899                        imp.status.lock().expect("lock poisoned").sent_end = true;
900                    }
901                    tracing::trace!(
902                        onionperf = true,
903                        event = ?OnionperfEvent::Stream(OnionperfStreamStatus::Closed),
904                        stream_id = ?imp.s.stream_id,
905                    );
906                    self.state = Some(DataWriterState::Closed);
907                } else {
908                    self.state = Some(DataWriterState::Ready(imp));
909                }
910                Poll::Ready(Ok(()))
911            }
912            Poll::Pending => {
913                self.state = Some(DataWriterState::Flushing(future));
914                Poll::Pending
915            }
916        }
917    }
918}
919
920impl AsyncWrite for DataWriterInner {
921    fn poll_write(
922        mut self: Pin<&mut Self>,
923        cx: &mut Context<'_>,
924        buf: &[u8],
925    ) -> Poll<IoResult<usize>> {
926        if buf.is_empty() {
927            return Poll::Ready(Ok(0));
928        }
929
930        let state = self.state.take().expect("Missing state in DataWriter");
931
932        let mut future = match state {
933            DataWriterState::Ready(mut imp) => {
934                let n_queued = imp.queue_bytes(buf);
935                if n_queued != 0 {
936                    self.state = Some(DataWriterState::Ready(imp));
937                    return Poll::Ready(Ok(n_queued));
938                }
939                // we couldn't queue anything, so the current cell must be full.
940                Box::pin(imp.flush_buf())
941            }
942            DataWriterState::Flushing(fut) => fut,
943            DataWriterState::Closed => {
944                self.state = Some(DataWriterState::Closed);
945                return Poll::Ready(Err(Error::NotConnected.into()));
946            }
947        };
948
949        match future.as_mut().poll(cx) {
950            Poll::Ready((_imp, Err(e))) => {
951                #[cfg(feature = "stream-ctrl")]
952                {
953                    _imp.status.lock().expect("lock poisoned").record_error(&e);
954                }
955                self.state = Some(DataWriterState::Closed);
956                Poll::Ready(Err(e.into()))
957            }
958            Poll::Ready((mut imp, Ok(()))) => {
959                // Great!  We're done flushing.  Queue as much as we can of this
960                // cell.
961                let n_queued = imp.queue_bytes(buf);
962                self.state = Some(DataWriterState::Ready(imp));
963                Poll::Ready(Ok(n_queued))
964            }
965            Poll::Pending => {
966                self.state = Some(DataWriterState::Flushing(future));
967                Poll::Pending
968            }
969        }
970    }
971
972    fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<IoResult<()>> {
973        self.poll_flush_impl(cx, false)
974    }
975
976    fn poll_close(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<IoResult<()>> {
977        self.poll_flush_impl(cx, true)
978    }
979}
980
981#[cfg(feature = "tokio")]
982impl TokioAsyncWrite for DataWriterInner {
983    fn poll_write(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &[u8]) -> Poll<IoResult<usize>> {
984        TokioAsyncWrite::poll_write(Pin::new(&mut self.compat_write()), cx, buf)
985    }
986
987    fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<IoResult<()>> {
988        TokioAsyncWrite::poll_flush(Pin::new(&mut self.compat_write()), cx)
989    }
990
991    fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<IoResult<()>> {
992        TokioAsyncWrite::poll_shutdown(Pin::new(&mut self.compat_write()), cx)
993    }
994}
995
996impl DataWriterImpl {
997    /// Try to flush the current buffer contents as a data cell.
998    async fn flush_buf(mut self) -> (Self, Result<()>) {
999        let result = if let Some((cell, remainder)) =
1000            Data::try_split_from(self.relay_cell_format, &self.buf[..self.n_pending])
1001        {
1002            // TODO: Eventually we may want a larger buffer; if we do,
1003            // this invariant will become false.
1004            assert!(remainder.is_empty());
1005            self.n_pending = 0;
1006            self.s.send(cell.into()).await
1007        } else {
1008            Ok(())
1009        };
1010
1011        (self, result)
1012    }
1013
1014    /// Add as many bytes as possible from `b` to our internal buffer;
1015    /// return the number we were able to add.
1016    fn queue_bytes(&mut self, b: &[u8]) -> usize {
1017        let empty_space = &mut self.buf[self.n_pending..];
1018        if empty_space.is_empty() {
1019            // that is, len == 0
1020            return 0;
1021        }
1022
1023        let n_to_copy = std::cmp::min(b.len(), empty_space.len());
1024        empty_space[..n_to_copy].copy_from_slice(&b[..n_to_copy]);
1025        self.n_pending += n_to_copy;
1026        n_to_copy
1027    }
1028}
1029
1030impl DataReaderInner {
1031    /// Return a [`ClientDataStreamCtrl`] object that can be used to monitor and
1032    /// interact with this stream without holding the stream itself.
1033    #[cfg(feature = "stream-ctrl")]
1034    pub(crate) fn client_stream_ctrl(&self) -> Option<&Arc<ClientDataStreamCtrl>> {
1035        self.ctrl.as_ref()
1036    }
1037}
1038
1039/// An enumeration for the state of a [`DataReaderInner`].
1040// TODO: We don't need to implement the state in this way anymore now that we've removed the saved
1041// future. There are a few ways we could simplify this. See:
1042// https://gitlab.torproject.org/tpo/core/arti/-/merge_requests/3076#note_3218210
1043#[derive(Educe)]
1044#[educe(Debug)]
1045// We allow this since it's expected that streams will spend most of their time in the `Open` state,
1046// and will be cleaned up shortly after closing.
1047#[allow(clippy::large_enum_variant)]
1048enum DataReaderState {
1049    /// In this state we have received an end cell or an error.
1050    Closed,
1051    /// In this state the reader is open.
1052    Open(DataReaderImpl),
1053}
1054
1055/// Wrapper for the read part of a [`DataStream`].
1056#[derive(Educe)]
1057#[educe(Debug)]
1058#[pin_project]
1059struct DataReaderImpl {
1060    /// The underlying StreamReceiver object.
1061    #[educe(Debug(method = "skip_fmt"))]
1062    #[pin]
1063    s: StreamReceiver,
1064
1065    /// If present, data that we received on this stream but have not
1066    /// been able to send to the caller yet.
1067    // TODO: This data structure is probably not what we want, but
1068    // it's good enough for now.
1069    #[educe(Debug(method = "skip_fmt"))]
1070    pending: Vec<u8>,
1071
1072    /// Index into pending to show what we've already read.
1073    offset: usize,
1074
1075    /// If true, we have received a CONNECTED cell on this stream.
1076    connected: bool,
1077
1078    /// Shared user-visible information about the state of this stream.
1079    #[cfg(feature = "stream-ctrl")]
1080    status: Arc<Mutex<DataStreamStatus>>,
1081}
1082
1083impl AsyncRead for DataReaderInner {
1084    fn poll_read(
1085        mut self: Pin<&mut Self>,
1086        cx: &mut Context<'_>,
1087        buf: &mut [u8],
1088    ) -> Poll<IoResult<usize>> {
1089        // We're pulling the state object out of the reader.  We MUST
1090        // put it back before this function returns.
1091        let mut state = self.state.take().expect("Missing state in DataReaderInner");
1092
1093        loop {
1094            let mut imp = match state {
1095                DataReaderState::Open(mut imp) => {
1096                    // There may be data to read already.
1097                    let n_copied = imp.extract_bytes(buf);
1098                    if n_copied != 0 || buf.is_empty() {
1099                        // We read data into the buffer, or the buffer was 0-len to begin with.
1100                        // Tell the caller.
1101                        self.state = Some(DataReaderState::Open(imp));
1102                        return Poll::Ready(Ok(n_copied));
1103                    }
1104
1105                    // No data available!  We have to try reading.
1106                    imp
1107                }
1108                DataReaderState::Closed => {
1109                    // TODO: Why are we returning an error rather than continuing to return EOF?
1110                    self.state = Some(DataReaderState::Closed);
1111                    return Poll::Ready(Err(Error::NotConnected.into()));
1112                }
1113            };
1114
1115            // See if a cell is ready.
1116            match Pin::new(&mut imp).read_cell(cx) {
1117                Poll::Ready(Err(e)) => {
1118                    // There aren't any survivable errors in the current
1119                    // design.
1120                    self.state = Some(DataReaderState::Closed);
1121                    #[cfg(feature = "stream-ctrl")]
1122                    {
1123                        imp.status.lock().expect("lock poisoned").record_error(&e);
1124                    }
1125                    let result = if matches!(e, Error::EndReceived(EndReason::DONE)) {
1126                        Ok(0)
1127                    } else {
1128                        Err(e.into())
1129                    };
1130                    return Poll::Ready(result);
1131                }
1132                Poll::Ready(Ok(())) => {
1133                    // It read a cell!  Continue the loop.
1134                    state = DataReaderState::Open(imp);
1135                }
1136                Poll::Pending => {
1137                    // No cells ready, so tell the
1138                    // caller to get back to us later.
1139                    self.state = Some(DataReaderState::Open(imp));
1140                    return Poll::Pending;
1141                }
1142            }
1143        }
1144    }
1145}
1146
1147#[cfg(feature = "tokio")]
1148impl TokioAsyncRead for DataReaderInner {
1149    fn poll_read(
1150        self: Pin<&mut Self>,
1151        cx: &mut Context<'_>,
1152        buf: &mut ReadBuf<'_>,
1153    ) -> Poll<IoResult<()>> {
1154        TokioAsyncRead::poll_read(Pin::new(&mut self.compat()), cx, buf)
1155    }
1156}
1157
1158impl DataReaderImpl {
1159    /// Pull as many bytes as we can off of self.pending, and return that
1160    /// number of bytes.
1161    fn extract_bytes(&mut self, buf: &mut [u8]) -> usize {
1162        let remainder = &self.pending[self.offset..];
1163        let n_to_copy = std::cmp::min(buf.len(), remainder.len());
1164        buf[..n_to_copy].copy_from_slice(&remainder[..n_to_copy]);
1165        self.offset += n_to_copy;
1166
1167        n_to_copy
1168    }
1169
1170    /// Return true iff there are no buffered bytes here to yield
1171    fn buf_is_empty(&self) -> bool {
1172        self.pending.len() == self.offset
1173    }
1174
1175    /// Load self.pending with the contents of a new data cell.
1176    fn read_cell(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<()>> {
1177        use ClientDataStreamMsg::*;
1178        let msg = match self.as_mut().project().s.poll_next(cx) {
1179            Poll::Pending => return Poll::Pending,
1180            Poll::Ready(Some(Ok(unparsed))) => match unparsed.decode::<ClientDataStreamMsg>() {
1181                Ok(cell) => cell.into_msg(),
1182                Err(e) => {
1183                    self.s.protocol_error();
1184                    return Poll::Ready(Err(Error::from_bytes_err(e, "message on a data stream")));
1185                }
1186            },
1187            Poll::Ready(Some(Err(e))) => return Poll::Ready(Err(e)),
1188            // TODO: This doesn't seem right to me, but seems to be the behaviour of the code before
1189            // the refactoring, so I've kept the same behaviour. I think if the cell stream is
1190            // terminated, we should be returning `None` here and not considering it as an error.
1191            // The `StreamReceiver` will have already returned an error if the cell stream was
1192            // terminated without an END message.
1193            Poll::Ready(None) => return Poll::Ready(Err(Error::NotConnected)),
1194        };
1195
1196        let result = match msg {
1197            Connected(_) if !self.connected => {
1198                self.connected = true;
1199                #[cfg(feature = "stream-ctrl")]
1200                {
1201                    self.status
1202                        .lock()
1203                        .expect("poisoned lock")
1204                        .record_connected();
1205                }
1206                Ok(())
1207            }
1208            Connected(_) => {
1209                self.s.protocol_error();
1210                Err(Error::StreamProto(
1211                    "Received a second connect cell on a data stream".to_string(),
1212                ))
1213            }
1214            Data(d) if self.connected => {
1215                self.add_data(d.into());
1216                Ok(())
1217            }
1218            Data(_) => {
1219                self.s.protocol_error();
1220                Err(Error::StreamProto(
1221                    "Received a data cell an unconnected stream".to_string(),
1222                ))
1223            }
1224            End(e) => Err(Error::EndReceived(e.reason())),
1225        };
1226
1227        Poll::Ready(result)
1228    }
1229
1230    /// Add the data from `d` to the end of our pending bytes.
1231    fn add_data(&mut self, mut d: Vec<u8>) {
1232        if self.buf_is_empty() {
1233            // No data pending?  Just take d as the new pending.
1234            self.pending = d;
1235            self.offset = 0;
1236        } else {
1237            // TODO(nickm) This has potential to grow `pending` without bound.
1238            // Fortunately, we don't currently read cells or call this
1239            // `add_data` method when pending is nonempty—but if we do in the
1240            // future, we'll have to be careful here.
1241            self.pending.append(&mut d);
1242        }
1243    }
1244}
1245
1246/// A `CmdChecker` that enforces invariants for outbound data streams.
1247#[derive(Debug)]
1248pub(crate) struct OutboundDataCmdChecker {
1249    /// True if we are expecting to receive a CONNECTED message on this stream.
1250    expecting_connected: bool,
1251}
1252
1253impl Default for OutboundDataCmdChecker {
1254    fn default() -> Self {
1255        Self {
1256            expecting_connected: true,
1257        }
1258    }
1259}
1260
1261impl CmdChecker for OutboundDataCmdChecker {
1262    fn check_msg(&mut self, msg: &tor_cell::relaycell::UnparsedRelayMsg) -> Result<StreamStatus> {
1263        use StreamStatus::*;
1264        match msg.cmd() {
1265            RelayCmd::CONNECTED => {
1266                if !self.expecting_connected {
1267                    Err(Error::StreamProto(
1268                        "Received CONNECTED twice on a stream.".into(),
1269                    ))
1270                } else {
1271                    self.expecting_connected = false;
1272                    Ok(Open)
1273                }
1274            }
1275            RelayCmd::DATA => {
1276                if !self.expecting_connected {
1277                    Ok(Open)
1278                } else {
1279                    Err(Error::StreamProto(
1280                        "Received DATA before CONNECTED on a stream".into(),
1281                    ))
1282                }
1283            }
1284            RelayCmd::END => Ok(Closed),
1285            _ => Err(Error::StreamProto(format!(
1286                "Unexpected {} on a data stream!",
1287                msg.cmd()
1288            ))),
1289        }
1290    }
1291
1292    fn consume_checked_msg(&mut self, msg: tor_cell::relaycell::UnparsedRelayMsg) -> Result<()> {
1293        let _ = msg
1294            .decode::<ClientDataStreamMsg>()
1295            .map_err(|err| Error::from_bytes_err(err, "cell on half-closed stream"))?;
1296        Ok(())
1297    }
1298}
1299
1300impl OutboundDataCmdChecker {
1301    /// Return a new boxed `DataCmdChecker` in a state suitable for a newly
1302    /// constructed connection.
1303    pub(crate) fn new_any() -> AnyCmdChecker {
1304        Box::<Self>::default()
1305    }
1306}