Skip to main content

tor_proto/circuit/reactor/
stream.rs

1//! The stream reactor.
2
3use crate::circuit::circhop::CircHopOutbound;
4use crate::circuit::reactor::macros::derive_deftly_template_CircuitReactor;
5use crate::circuit::{CircHopSyncView, UniqId};
6use crate::congestion::{CongestionControl, sendme};
7use crate::memquota::{CircuitAccount, SpecificAccount as _, StreamAccount};
8use crate::stream::CloseStreamBehavior;
9use crate::stream::cmdcheck::StreamStatus;
10use crate::stream::flow_ctrl::state::WithSidechannelMitigations;
11use crate::streammap;
12use crate::util::err::ReactorError;
13use crate::{Error, HopNum};
14
15#[cfg(any(feature = "hs-service", feature = "relay"))]
16use crate::stream::incoming::{
17    InboundDataCmdChecker, IncomingStreamRequest, IncomingStreamRequestContext,
18    IncomingStreamRequestDisposition, IncomingStreamRequestHandler, StreamReqInfo,
19};
20
21use tor_async_utils::{SinkTrySend as _, SinkTrySendError as _};
22use tor_cell::chancell::CircId;
23use tor_cell::relaycell::msg::{AnyRelayMsg, Begin, BeginDir, End, EndReason, Resolve};
24use tor_cell::relaycell::{
25    AnyRelayMsgOuter, RelayCellFormat, RelayCmd, StreamId, UnparsedRelayMsg,
26};
27use tor_error::{internal, into_internal};
28use tor_log_ratelim::log_ratelim;
29use tor_rtcompat::{DynTimeProvider, Runtime, SleepProvider as _};
30
31use derive_deftly::Deftly;
32use futures::SinkExt;
33use futures::channel::mpsc;
34use futures::{FutureExt as _, StreamExt as _, future, select_biased};
35use tracing::debug;
36
37use std::pin::Pin;
38use std::result::Result as StdResult;
39use std::sync::{Arc, Mutex};
40use std::task::Poll;
41use std::time::Duration;
42
43/// Trait for customizing the behavior of the stream reactor.
44///
45/// Used for plugging in the implementation-dependent (client vs relay)
46/// parts of the implementation into the generic one.
47pub(crate) trait StreamHandler: Send + Sync + 'static {
48    /// Return the amount of time a newly closed stream
49    /// should be kept in the stream map for.
50    ///
51    /// This is the amount of time we are willing to wait for
52    /// an END ack before removing the half-stream from the map.
53    fn halfstream_expiry(&self, hop: &CircHopOutbound) -> Duration;
54
55    /// Whether sidechannel mitigations should be enabled for incoming streams.
56    fn flowctrl_sidechannel_mitigations(&self) -> WithSidechannelMitigations;
57}
58
59/// The stream reactor for a given hop.
60///
61/// Drives the application streams.
62///
63/// This reactor accepts [`CtrlMsg`]s from the forward reactor over its [`Self::cell_rx`]
64/// MPSC channel, and delivers them to the corresponding stream entries in the stream map.
65///
66/// The local streams are polled from the main loop, and any ready messages are sent
67/// to the backward reactor over the `bwd_tx` MPSC channel for packaging and delivery.
68///
69/// Shuts downs down if an error occurs, or if the sending end
70/// of the `cell_rx` MPSC channel, i.e. the forward reactor, closes.
71#[derive(Deftly)]
72#[derive_deftly(CircuitReactor)]
73#[deftly(reactor_name = "stream reactor")]
74#[deftly(run_inner_fn = "Self::run_once")]
75#[must_use = "If you don't call run() on a reactor, the circuit won't work."]
76pub(crate) struct StreamReactor {
77    /// The hop this stream reactor is for.
78    ///
79    /// This is `None` for relays.
80    hopnum: Option<HopNum>,
81    /// The state of this circuit hop.
82    hop: CircHopOutbound,
83    /// The time provider.
84    time_provider: DynTimeProvider,
85    /// An identifier for logging about this reactor's circuit.
86    unique_id: UniqId,
87    /// The circuit identifier on the inbound Tor channel.
88    circ_id: CircId,
89    /// Receiver for Tor stream data that need to be delivered to a Tor stream.
90    ///
91    /// The sender is in the [`HopMgr`](super::hop_mgr::HopMgr) of the
92    /// [`ForwardReactor`](super::ForwardReactor), which will forward all cells
93    /// carrying Tor stream data to us.
94    ///
95    /// This serves a dual purpose:
96    ///
97    ///   * it enables the `ForwardReactor` to deliver Tor stream data received from the client
98    ///   * it lets the `StreamReactor` know if the `ForwardReactor` has shut down:
99    ///     we select! on this MPSC channel in the main loop, so if the `ForwardReactor`
100    ///     shuts down, we will get EOS upon calling `.next()`)
101    cell_rx: mpsc::Receiver<CtrlMsg>,
102    /// Sender for sending Tor stream data to [`BackwardReactor`](super::BackwardReactor).
103    bwd_tx: mpsc::Sender<ReadyStreamMsg>,
104    /// A handler for incoming streams.
105    ///
106    /// Set to `None` if incoming streams are not allowed on this circuit.
107    ///
108    /// This handler is shared with the [`HopMgr`](super::hop_mgr::HopMgr) of this reactor,
109    /// which can install a new handler at runtime (for example, in response to a CtrlMsg).
110    /// The ability to update the handler after the reactor is launched is needed
111    /// for onion services, where the incoming stream request handler only gets installed
112    /// after the virtual hop is created.
113    #[cfg(any(feature = "hs-service", feature = "relay"))]
114    incoming: Arc<Mutex<Option<IncomingStreamRequestHandler>>>,
115    /// A handler for customizing the stream reactor behavior.
116    inner: Arc<dyn StreamHandler>,
117    /// Memory quota account
118    memquota: CircuitAccount,
119}
120
121#[allow(unused)] // TODO(relay)
122impl StreamReactor {
123    /// Create a new [`StreamReactor`].
124    #[allow(clippy::too_many_arguments)] // TODO
125    pub(crate) fn new<R: Runtime>(
126        runtime: R,
127        hopnum: Option<HopNum>,
128        hop: CircHopOutbound,
129        unique_id: UniqId,
130        circ_id: CircId,
131        cell_rx: mpsc::Receiver<CtrlMsg>,
132        bwd_tx: mpsc::Sender<ReadyStreamMsg>,
133        inner: Arc<dyn StreamHandler>,
134        #[cfg(any(feature = "hs-service", feature = "relay"))] //
135        incoming: Arc<Mutex<Option<IncomingStreamRequestHandler>>>,
136        memquota: CircuitAccount,
137    ) -> Self {
138        Self {
139            hopnum,
140            hop,
141            time_provider: DynTimeProvider::new(runtime),
142            unique_id,
143            circ_id,
144            #[cfg(any(feature = "hs-service", feature = "relay"))]
145            incoming,
146            cell_rx,
147            bwd_tx,
148            inner,
149            memquota,
150        }
151    }
152
153    /// Helper for [`run`](Self::run).
154    ///
155    /// Polls the stream map for messages
156    /// that need to be delivered to the other endpoint,
157    /// and the `cells_rx` MPSC stream for stream messages received
158    /// from the `ForwardReactor` that need to be delivered to the application streams.
159    async fn run_once(&mut self) -> StdResult<(), ReactorError> {
160        use postage::prelude::{Sink as _, Stream as _};
161
162        // Garbage-collect all halfstreams that have expired.
163        //
164        // Note: this will iterate over the closed streams of this hop.
165        // If we think this will cause perf issues, one idea would be to make
166        // StreamMap::closed_streams into a min-heap, and add a branch to the
167        // select_biased! below to sleep until the first expiry is due
168        // (but my gut feeling is that iterating is cheaper)
169        self.hop
170            .stream_map()
171            .lock()
172            .expect("poisoned lock")
173            .remove_expired_halfstreams(self.time_provider.now());
174
175        let mut streams = Arc::clone(self.hop.stream_map());
176        let can_send = self
177            .hop
178            .ccontrol()
179            .lock()
180            .expect("poisoned lock")
181            .can_send();
182        let mut ready_streams_fut = future::poll_fn(move |cx| {
183            if !can_send {
184                // We can't send anything on this hop that counts towards SENDME windows.
185                //
186                // Note: this does not block outgoing flow-control messages:
187                //
188                //   * circuit SENDMEs are initiated by the forward reactor,
189                //     by sending a BackwardReactorCmd::SendRelayMsg to BWD,
190                //   * stream SENDMEs will be initiated by StreamTarget::send_sendme(),
191                //     by sending a control message to the reactor
192                //     (TODO(relay): not yet implemented)
193                //   * XOFFs are sent in response to messages on streams
194                //     (i.e. RELAY messages with non-zero stream IDs).
195                //     These messages are delivered to us by the forward reactor
196                //     inside BackwardReactorCmd::HandleMsg
197                //   * XON will be initiated by StreamTarget::drain_rate_update(),
198                //     by sending a control message to the reactor
199                //     (TODO(relay): not yet implemented)\
200                return Poll::Pending;
201            }
202
203            let mut streams = streams.lock().expect("lock poisoned");
204            let Some((sid, msg)) = streams.poll_ready_streams_iter(cx).next() else {
205                // No ready streams
206                //
207                // TODO(flushing): if there are no ready Tor streams, we might want to defer
208                // flushing until stream data becomes available (or until a timeout elapses).
209                // The deferred flushing approach should enable us to send
210                // more than one message at a time to the channel reactor.
211                return Poll::Pending;
212            };
213
214            if msg.is_none() {
215                // This means the local sender has been dropped,
216                // which presumably can only happen if an error occurs,
217                // or if the Tor stream ends. In both cases, we're going to
218                // want to send an END to the client to let them know,
219                // and to remove the stream from the stream map.
220                //
221                // TODO(relay): the local sender part is not implemented yet
222                return Poll::Ready(StreamEvent::ApplicationStreamClosed(sid));
223            };
224
225            let msg = streams.take_ready_msg(sid).expect("msg disappeared");
226
227            Poll::Ready(StreamEvent::ReadyMsg { sid, msg })
228        });
229
230        select_biased! {
231            res = self.cell_rx.next().fuse() => {
232                let Some(cmd) = res else {
233                    // The forward reactor has shut down
234                    return Err(ReactorError::Shutdown);
235                };
236
237                self.handle_reactor_cmd(cmd).await?;
238            }
239            event = ready_streams_fut.fuse() => {
240                self.handle_stream_event(event).await?;
241            }
242        }
243
244        Ok(())
245    }
246
247    /// Handle a stream message sent to us by the forward reactor.
248    ///
249    /// Delivers the message to its corresponding application stream.
250    async fn handle_reactor_cmd(&mut self, msg: CtrlMsg) -> StdResult<(), ReactorError> {
251        match msg {
252            CtrlMsg::DeliverStreamMsg {
253                sid,
254                msg,
255                cell_counts_toward_windows,
256            } => {
257                self.deliver_message_to_stream(sid, msg, cell_counts_toward_windows)
258                    .await
259            }
260            #[cfg(any(feature = "hs-service", feature = "relay"))]
261            CtrlMsg::ClosePendingStream { stream_id, behav } => {
262                self.close_stream(stream_id, behav, streammap::TerminateReason::ExplicitEnd)
263                    .await
264            }
265        }
266    }
267
268    /// Deliver `msg` to the specified stream
269    async fn deliver_message_to_stream(
270        &mut self,
271        sid: StreamId,
272        msg: UnparsedRelayMsg,
273        cell_counts_toward_windows: bool,
274    ) -> StdResult<(), ReactorError> {
275        // We need to apply stream-level flow control *before* encoding the message.
276        // May optionally return a message that needs to be sent back to the client.
277        let bwd_msg = self.handle_msg(sid, msg, cell_counts_toward_windows)?;
278
279        if let Some(bwd_msg) = bwd_msg {
280            self.send_msg_to_bwd(bwd_msg).await?;
281        }
282
283        Ok(())
284    }
285
286    /// Handle a RELAY message that has a non-zero stream ID.
287    ///
288    /// A returned message is one that we need to send back to the client.
289    //
290    // TODO(relay): this is very similar to the client impl from
291    // Circuit::handle_in_order_relay_msg()
292    fn handle_msg(
293        &mut self,
294        streamid: StreamId,
295        msg: UnparsedRelayMsg,
296        cell_counts_toward_windows: bool,
297    ) -> StdResult<Option<AnyRelayMsgOuter>, ReactorError> {
298        let cmd = msg.cmd();
299        let possible_proto_violation_err = move |streamid: StreamId| {
300            Error::StreamProto(format!(
301                "Unexpected {cmd:?} message on unknown stream {streamid}"
302            ))
303        };
304        let now = self.time_provider.now();
305
306        // Check if any of our already-open streams want this message
307        let res = self.hop.handle_msg(
308            possible_proto_violation_err,
309            cell_counts_toward_windows,
310            streamid,
311            msg,
312            now,
313        )?;
314
315        // If it was an incoming stream request, we don't need to worry about
316        // sending an XOFF as there's no stream data within this message.
317        if let Some(msg) = res {
318            cfg_if::cfg_if! {
319                if #[cfg(any(feature = "hs-service", feature = "relay"))] {
320                    return self.handle_incoming_stream_request(streamid, msg);
321                } else {
322                    return Err(
323                        Error::CircProto(format!("Cannot handle {} cells on this circuit", msg.cmd())).into(),
324                    );
325                }
326            }
327        }
328
329        // We may want to send an XOFF if the incoming buffer is too large.
330        if let Some(cell) = self.hop.maybe_send_xoff(streamid)? {
331            let cell = AnyRelayMsgOuter::new(Some(streamid), cell.into());
332            return Ok(Some(cell));
333        }
334
335        Ok(None)
336    }
337
338    /// A helper for handling incoming stream requests.
339    ///
340    /// Accepts the specified incoming stream request,
341    /// by adding a new entry to our stream map.
342    ///
343    /// Returns the cell we need to send back to the client,
344    /// if an error occurred and the stream cannot be opened.
345    ///
346    /// Returns None if everything went well
347    /// (the CONNECTED response only comes if the external
348    /// consumer of our [Stream](futures::Stream) of incoming Tor streams
349    /// is able to actually establish the connection to the address
350    /// specified in the BEGIN).
351    ///
352    /// Any error returned from this function will shut down the reactor.
353    #[cfg(any(feature = "hs-service", feature = "relay"))]
354    fn handle_incoming_stream_request(
355        &mut self,
356        sid: StreamId,
357        msg: UnparsedRelayMsg,
358    ) -> StdResult<Option<AnyRelayMsgOuter>, ReactorError> {
359        let mut lock = self.incoming.lock().expect("poisoned lock");
360        let Some(handler) = lock.as_mut() else {
361            return Err(Error::CircProto(format!(
362                "Cannot handle {} cells on this circuit",
363                msg.cmd()
364            ))
365            .into());
366        };
367
368        if self.hopnum != handler.hop_num {
369            let expected_hopnum = match handler.hop_num {
370                Some(hopnum) => hopnum.display().to_string(),
371                None => "client".to_string(),
372            };
373
374            let actual_hopnum = match self.hopnum {
375                Some(hopnum) => hopnum.display().to_string(),
376                None => "None".to_string(),
377            };
378
379            return Err(Error::CircProto(format!(
380                "Expecting incoming streams from {}, but received {} cell from unexpected hop {}",
381                expected_hopnum,
382                msg.cmd(),
383                actual_hopnum,
384            ))
385            .into());
386        }
387
388        let message_closes_stream = handler.cmd_checker.check_msg(&msg)? == StreamStatus::Closed;
389
390        if message_closes_stream {
391            self.hop
392                .stream_map()
393                .lock()
394                .expect("poisoned lock")
395                .ending_msg_received(sid)?;
396
397            return Ok(None);
398        }
399
400        let req = parse_incoming_stream_req(msg)?;
401        let view = CircHopSyncView::new(&self.hop);
402
403        if let Some(reject) = Self::should_reject_incoming(handler, sid, &req, &view)? {
404            // We can't honor this request, so we bail by sending an END.
405            return Ok(Some(reject));
406        };
407
408        let memquota =
409            StreamAccount::new(&self.memquota).map_err(|e| ReactorError::Err(e.into()))?;
410
411        let cmd_checker = InboundDataCmdChecker::new_connected();
412        let stream_components = self.hop.add_ent_with_id(
413            &self.time_provider,
414            sid,
415            cmd_checker,
416            self.inner.flowctrl_sidechannel_mitigations(),
417            &memquota,
418        )?;
419
420        let outcome = Pin::new(&mut handler.incoming_sender).try_send(StreamReqInfo {
421            req,
422            stream_id: sid,
423            hop: None,
424            stream_components,
425            memquota,
426            relay_cell_format: self.hop.relay_cell_format(),
427        });
428
429        log_ratelim!("Delivering message to incoming stream handler"; outcome);
430
431        if let Err(e) = outcome {
432            if e.is_full() {
433                // The IncomingStreamRequestHandler's stream is full; it isn't
434                // handling requests fast enough. So instead, we reply with an
435                // END cell.
436                let end_msg = AnyRelayMsgOuter::new(
437                    Some(sid),
438                    End::new_with_reason(EndReason::RESOURCELIMIT).into(),
439                );
440
441                return Ok(Some(end_msg));
442            } else if e.is_disconnected() {
443                // The IncomingStreamRequestHandler's stream has been dropped.
444                // In the Tor protocol as it stands, this always means that the
445                // circuit itself is out-of-use and should be closed.
446                //
447                // Note that we will _not_ reach this point immediately after
448                // the IncomingStreamRequestHandler is dropped; we won't hit it
449                // until we next get an incoming request.  Thus, if we later
450                // want to add early detection for a dropped
451                // IncomingStreamRequestHandler, we need to do it elsewhere, in
452                // a different way.
453                debug!(
454                    circ_uniq_id = %self.unique_id,
455                    backward_circ_id = %self.circ_id,
456                    "Incoming stream request receiver dropped",
457                );
458                // This will _cause_ the circuit to get closed.
459                return Err(ReactorError::Err(Error::CircuitClosed));
460            } else {
461                // There are no errors like this with the current design of
462                // futures::mpsc, but we shouldn't just ignore the possibility
463                // that they'll be added later.
464                return Err(
465                    Error::from((into_internal!("try_send failed unexpectedly"))(e)).into(),
466                );
467            }
468        }
469
470        Ok(None)
471    }
472
473    /// Check if we should reject this incoming stream request or not.
474    ///
475    /// Returns a cell we need to send back to the client if we must reject the request,
476    /// or `None` if we are allowed to accept it.
477    ///`
478    /// Any error returned from this function will shut down the reactor.
479    #[cfg(any(feature = "hs-service", feature = "relay"))]
480    fn should_reject_incoming<'a>(
481        handler: &mut IncomingStreamRequestHandler,
482        sid: StreamId,
483        request: &IncomingStreamRequest,
484        view: &CircHopSyncView<'a>,
485    ) -> StdResult<Option<AnyRelayMsgOuter>, ReactorError> {
486        use IncomingStreamRequestDisposition::*;
487
488        let ctx = IncomingStreamRequestContext { request };
489
490        // Run the externally provided filter to check if we should
491        // open the stream or not.
492        match handler.filter.as_mut().disposition(&ctx, view)? {
493            Accept => {
494                // All is well, we can accept the stream request
495                Ok(None)
496            }
497            CloseCircuit => Err(ReactorError::Shutdown),
498            RejectRequest(end) => {
499                let end_msg = AnyRelayMsgOuter::new(Some(sid), end.into());
500
501                Ok(Some(end_msg))
502            }
503        }
504    }
505
506    /// Handle a [`StreamEvent`].
507    async fn handle_stream_event(&mut self, event: StreamEvent) -> StdResult<(), ReactorError> {
508        match event {
509            StreamEvent::ApplicationStreamClosed(sid) => {
510                self.close_stream(
511                    sid,
512                    CloseStreamBehavior::default(),
513                    streammap::TerminateReason::StreamTargetClosed,
514                )
515                .await
516            }
517            StreamEvent::ReadyMsg { sid, msg } => {
518                self.send_msg_to_bwd(AnyRelayMsgOuter::new(Some(sid), msg))
519                    .await
520            }
521        }
522    }
523
524    /// Close the stream that has the specified `sid`.
525    ///
526    /// The `behav` controls whether an `END` will be sent or not.
527    ///
528    /// This calls [`CircHopOutbound::close_stream`] under the hood,
529    /// which removes the stream from the stream map,
530    /// and returns an optional `END` cell to send back to the other party.
531    async fn close_stream(
532        &mut self,
533        sid: StreamId,
534        behav: CloseStreamBehavior,
535        reason: streammap::TerminateReason,
536    ) -> StdResult<(), ReactorError> {
537        let timeout = self.inner.halfstream_expiry(&self.hop);
538        let expire_at = self.time_provider.now() + timeout;
539        let res = self.hop.close_stream(
540            self.unique_id,
541            self.circ_id,
542            sid,
543            None,
544            behav,
545            reason,
546            expire_at,
547        )?;
548        let Some(msg) = res else {
549            // We may not need to send anything at all...
550            return Ok(());
551        };
552
553        self.send_msg_to_bwd(msg.cell).await
554    }
555
556    /// Wrap `msg` in [`ReadyStreamMsg`], and send it to the backward reactor.
557    async fn send_msg_to_bwd(&mut self, msg: AnyRelayMsgOuter) -> StdResult<(), ReactorError> {
558        // TODO(DEDUP): this contains parts of Circuit::send_relay_cell_inner()
559
560        // We might be out of capacity entirely; see if we are about to hit a limit.
561        //
562        // TODO: If we ever add a notion of _recoverable_ errors below, we'll
563        // need a way to restore this limit, and similarly for about_to_send().
564        self.hop.decrement_cell_limit()?;
565
566        // We need to apply stream-level flow control *before* encoding the message
567        // (the BWD handles the encoding)
568        if sendme::cmd_counts_towards_windows(msg.cmd()) {
569            if let Some(stream_id) = msg.stream_id() {
570                self.hop
571                    .about_to_send(self.unique_id, self.circ_id, stream_id, msg.msg())?;
572            }
573        }
574
575        // NOTE: on the client side, we call note_data_sent()
576        // just before writing the cell to the channel.
577        // We can't do that here, because we're not the ones
578        // encoding the cell, so we don't have the SENDME tag
579        // which is needed for note_data_sent().
580        //
581        // Instead, we notify the CC algorithm in the BWD,
582        // right after we've finished sending the cell.
583
584        let msg = ReadyStreamMsg {
585            hop: self.hopnum,
586            relay_cell_format: self.hop.relay_cell_format(),
587            ccontrol: Arc::clone(self.hop.ccontrol()),
588            msg,
589        };
590
591        self.bwd_tx
592            .send(msg)
593            .await
594            .map_err(|_| ReactorError::Shutdown)?;
595
596        Ok(())
597    }
598}
599
600/// A Tor stream-related event.
601enum StreamEvent {
602    /// An application stream was closed.
603    ///
604    /// The corresponding entry needs to be removed from the reactor's stream map.
605    ApplicationStreamClosed(StreamId),
606    /// A stream has a ready message.
607    ReadyMsg {
608        /// The ID of the stream to close.
609        sid: StreamId,
610        /// The message.
611        msg: AnyRelayMsg,
612    },
613}
614
615/// Convert an incoming stream request message (BEGIN, BEGIN_DIR, RESOLVE, etc.)
616/// to an [`IncomingStreamRequest`]
617///
618// TODO(dedup): when we rewrite the client reactor in the multi-reactor register,
619// we should rethink this part a bit: ideally, onion services shouldn't even
620// try to parse BEGIN_DIR, RESOLVE.
621//
622// We will likely need an implementation-specific hook for this,
623// similar to the `{Forward,Backward}Handler` implementation-specific handlers
624// we have for the FWD and BWD reactors.
625//
626// See https://gitlab.torproject.org/tpo/core/arti/-/merge_requests/4188#note_3432579
627#[cfg(any(feature = "hs-service", feature = "relay"))]
628fn parse_incoming_stream_req(msg: UnparsedRelayMsg) -> crate::Result<IncomingStreamRequest> {
629    /// Helper for parsing an incoming stream request
630    /// (BEGIN, BEGIN_DIR, or RESOLVE)
631    macro_rules! parse_stream_req {
632        ($msg:expr, $type:tt) => {{
633            let req = $msg
634                .decode::<$type>()
635                .map_err(|e| {
636                    Error::from_bytes_err(e, concat!("Invalid ", stringify!($type), " message"))
637                })?
638                .into_msg();
639
640            IncomingStreamRequest::$type(req)
641        }};
642    }
643
644    let req = match msg.cmd() {
645        RelayCmd::BEGIN => parse_stream_req!(msg, Begin),
646        RelayCmd::BEGIN_DIR => parse_stream_req!(msg, BeginDir),
647        RelayCmd::RESOLVE => parse_stream_req!(msg, Resolve),
648        cmd => {
649            // It's a bug if we reach this point, because CircHopOutbound::handle_msg()
650            // should have consumed the message (by forwarding it to the appropriate stream
651            // in its stream map)
652            return Err(internal!("{cmd} is not an incoming stream request").into());
653        }
654    };
655
656    Ok(req)
657}
658
659/// A stream message to be sent to the backward reactor for delivery.
660pub(crate) struct ReadyStreamMsg {
661    /// The hop number, or `None` if we are a relay.
662    pub(crate) hop: Option<HopNum>,
663    /// The message to send.
664    pub(crate) msg: AnyRelayMsgOuter,
665    /// The cell format used with the hop the message should be sent to.
666    pub(crate) relay_cell_format: RelayCellFormat,
667    /// The CC object to use.
668    pub(crate) ccontrol: Arc<Mutex<CongestionControl>>,
669}
670
671/// A control message
672/// that needs to be handled by [`StreamReactor`].
673pub(crate) enum CtrlMsg {
674    /// Stream data received from the other endpoint
675    /// that needs to be delivered to a Tor stream
676    DeliverStreamMsg {
677        /// The ID of the stream this message is for.
678        sid: StreamId,
679        /// The message.
680        msg: UnparsedRelayMsg,
681        /// Whether the cell this message came from counts towards flow-control windows.
682        cell_counts_toward_windows: bool,
683    },
684
685    /// Close the specified pending incoming stream, sending the provided END message.
686    #[cfg(any(feature = "hs-service", feature = "relay"))]
687    ClosePendingStream {
688        /// The stream ID to send the END for.
689        stream_id: StreamId,
690        /// The END message to send, if any.
691        behav: CloseStreamBehavior,
692    },
693}