Skip to main content

tor_proto/circuit/reactor/
backward.rs

1//! A circuit's view of the backward state of the circuit.
2
3use crate::channel::Channel;
4use crate::circuit::UniqId;
5use crate::circuit::cell_sender::CircuitCellSender;
6use crate::circuit::reactor::ControlHandler;
7use crate::circuit::reactor::circhop::CircHopList;
8use crate::circuit::reactor::macros::derive_deftly_template_CircuitReactor;
9use crate::circuit::reactor::stream::ReadyStreamMsg;
10use crate::congestion::{CongestionControl, sendme};
11use crate::crypto::cell::RelayCellBody;
12use crate::util::err::ReactorError;
13use crate::util::poll_all::PollAll;
14use crate::{Error, HopNum, Result};
15
16// TODO(circpad): once padding is stabilized, the padding module will be moved out of client.
17use crate::client::circuit::padding::{
18    self, PaddingController, PaddingEvent, PaddingEventStream, QueuedCellPaddingInfo,
19};
20
21use tor_cell::chancell::msg::{AnyChanMsg, Relay};
22use tor_cell::chancell::{AnyChanCell, BoxedCellBody, ChanCmd, CircId};
23use tor_cell::relaycell::flow_ctrl::XonKBpsEwma;
24use tor_cell::relaycell::msg::{Sendme, SendmeTag};
25use tor_cell::relaycell::{AnyRelayMsgOuter, RelayCellFormat, RelayCmd, StreamId};
26use tor_error::internal;
27use tor_rtcompat::{DynTimeProvider, Runtime};
28
29use derive_deftly::Deftly;
30use futures::SinkExt;
31use futures::channel::mpsc;
32use futures::{FutureExt as _, StreamExt, future, select_biased};
33use tracing::{debug, trace};
34
35use std::pin::Pin;
36use std::result::Result as StdResult;
37use std::sync::{Arc, Mutex, RwLock};
38
39use crate::circuit::CircuitRxReceiver;
40
41#[cfg(feature = "circ-padding")]
42use crate::circuit::padding::{CircPaddingDisposition, padding_disposition};
43
44#[cfg(feature = "relay")]
45use tor_cell::relaycell::msg::Extended2;
46
47/// The "backward" circuit reactor of a relay.
48///
49/// See the [`reactor`](crate::circuit::reactor) module-level docs.
50///
51/// Shuts downs down if an error occurs, or if the [`Reactor`](super::Reactor),
52/// [`ForwardReactor`](super::ForwardReactor), or if one of the
53/// [`StreamReactor`](super::stream::StreamReactor)s of this circuit shuts down:
54///
55///   * if the `Reactor` shuts down, we are alerted via the ctrl/command mpsc channels
56///     (their sending ends will close, which causes run_once() to return ReactorError::Shutdown)
57///   * if `ForwardReactor` shuts down, the `Reactor` will notice and will itself shut down,
58///     which, in turn, causes the `BackwardReactor` to shut down as described above
59///   * if one of the `StreamReactor`s shuts down, the `ForwardReactor` will
60///     notice when it next tries to deliver a stream message to it, and shut down,
61///     causing the `BackwardReactor` and top-level `Reactor` to follow suit
62#[derive(Deftly)]
63#[derive_deftly(CircuitReactor)]
64#[deftly(reactor_name = "backward reactor")]
65#[deftly(run_inner_fn = "Self::run_once")]
66#[must_use = "If you don't call run() on a reactor, the circuit won't work."]
67pub(super) struct BackwardReactor<B: BackwardHandler> {
68    /// The time provider.
69    time_provider: DynTimeProvider,
70    /// An identifier for logging about this reactor's circuit.
71    unique_id: UniqId,
72    /// The circuit identifier on the backward Tor channel.
73    circ_id: CircId,
74    /// The inbound Tor channel.
75    channel: Arc<Channel>,
76    /// Implementation-dependent part of the reactor.
77    ///
78    /// This enables us to customize the behavior of the reactor,
79    /// depending on whether we are a client or a relay.
80    inner: B,
81    /// The reading end of the outbound Tor channel, if we are not the last hop.
82    ///
83    /// Yields cells moving from the exit towards the client, if we are a middle relay.
84    outbound_chan_rx: Option<CircuitRxReceiver>,
85    /// The per-hop state, shared with the forward reactor.
86    ///
87    /// The backward reactor acquires a read lock to this whenever it needs to
88    ///
89    ///   * send a circuit-level SENDME
90    ///   * handle a circuit-level SENDME
91    ///   * send a padding cell
92    ///
93    // Note: For the sending/handling of SENDMEs, we lock the hop list
94    // to extract the relay cell format and CC state of the hop.
95    // Technically, for the SENDME cases, we could've avoided locking
96    // the hop list from the BWD, by having the FWD share the relay cell format
97    // and CC state in the BackwardReactorCmd::{Send,Handle}Sendme command.
98    // But for the padding case, we *need* the hop list, because we need
99    // to work out what relay cell format to use when sending the padding cell.
100    // But for the sake of simplicity, I made the BWD consult the CircHopList in all cases.
101    //
102    // TODO: the backward reactor only ever reads from this.
103    // Conceptually, it is the forward reactor's HopMgr that owns this list:
104    // only HopMgr can add hops to the list.
105    //
106    // Perhaps we need a specialized abstraction that only allows reading here.
107    // This could be a wrapper over RwLock, providing a read-only API.
108    hops: Arc<RwLock<CircHopList>>,
109    /// The sending end of the backward Tor channel.
110    ///
111    /// Delivers cells towards the other endpoint: towards the client, if we are a relay,
112    /// or towards the exit, if we are a client.
113    inbound_chan_tx: CircuitCellSender,
114    /// Channel for receiving control commands.
115    command_rx: mpsc::UnboundedReceiver<CtrlCmd<B::CtrlCmd>>,
116    /// Channel for receiving control messages.
117    control_rx: mpsc::UnboundedReceiver<CtrlMsg<B::CtrlMsg>>,
118    /// Receiver for [`BackwardReactorCmd`]s coming from the forward reactor.
119    ///
120    /// The sender is in [`ForwardReactor`](super::ForwardReactor), which will forward all cells
121    /// carrying Tor stream data to us.
122    ///
123    /// This serves a dual purpose:
124    ///
125    ///   * it enables the `ForwardReactor` to deliver Tor stream data received
126    ///     from the other endpoint
127    ///   * it lets the `BackwardReactor` know if the `ForwardReactor` has shut down:
128    ///     we select! on this MPSC channel in the main loop, so if the `ForwardReactor`
129    ///     shuts down, we will get EOS upon calling `.next()`)
130    forward_reactor_rx: mpsc::Receiver<BackwardReactorCmd>,
131    /// A channel for receiving endpoint-bound stream messages from the StreamReactor(s)
132    /// (the stream messages are client-bound if we are a relay, or exit-bound if we are a client).
133    stream_rx: mpsc::Receiver<ReadyStreamMsg>,
134    /// A padding controller to which padding-related events should be reported.
135    padding_ctrl: PaddingController,
136    /// An event stream telling us about padding-related events.
137    padding_event_stream: PaddingEventStream,
138    /// Current rules for blocking traffic, according to the padding controller.
139    #[cfg(feature = "circ-padding")]
140    padding_block: Option<padding::StartBlocking>,
141}
142
143/// A control message aimed at the generic backward reactor.
144pub(crate) enum CtrlMsg<M> {
145    /// Inform the reactor that there's a flow control update for a given stream.
146    ///
147    /// The reactor will decide how to handle this update depending on the type of flow control and
148    /// the current state of the stream.
149    FlowCtrlUpdate {
150        /// The hop that the stream is on.
151        /// Relay circuits use `None`.
152        hop: Option<HopNum>,
153        /// The stream ID that the update is for.
154        stream_id: StreamId,
155        /// The type of flow control update, and any associated metadata.
156        msg: FlowCtrlMsg,
157    },
158    /// An implementation-dependent control message.
159    #[allow(unused)] // TODO(relay)
160    Custom(M),
161}
162
163/// A control command aimed at the generic backward reactor.
164pub(crate) enum CtrlCmd<C> {
165    /// An implementation-dependent control command.
166    #[allow(unused)] // TODO(relay)
167    Custom(C),
168}
169
170/// Trait for customizing the behavior of the backward reactor.
171///
172/// Used for plugging in the implementation-dependent (client vs relay)
173/// parts of the implementation into the generic one.
174pub(crate) trait BackwardHandler: ControlHandler {
175    /// The subclass of ChanMsg that can arrive on this type of circuit.
176    type CircChanMsg: TryFrom<AnyChanMsg, Error = crate::Error> + Send;
177
178    /// Encrypt a RelayCellBody that is moving in the backward direction.
179    fn encrypt_relay_cell(
180        &mut self,
181        cmd: ChanCmd,
182        body: &mut RelayCellBody,
183        hop: Option<HopNum>,
184    ) -> SendmeTag;
185
186    /// Handle a cell that was read from the Tor outbound channel.
187    ///
188    /// Returns an error if the cell should cause the reactor to shut down,
189    /// or a [`BackwardCellDisposition`] specifying how it should be handled.
190    fn handle_backward_cell(
191        &mut self,
192        circ_uniq_id: UniqId,
193        circ_id: CircId,
194        cell: Self::CircChanMsg,
195    ) -> StdResult<BackwardCellDisposition, ReactorError>;
196}
197
198/// What action to take in response to a cell arriving on our outbound Tor channel.
199pub(crate) enum BackwardCellDisposition {
200    /// Forward the cell, writing it to the inbound Tor channel.
201    Forward(AnyChanMsg),
202}
203
204#[allow(unused)] // TODO(relay)
205impl<B: BackwardHandler> BackwardReactor<B> {
206    /// Create a new [`BackwardReactor`].
207    #[allow(clippy::too_many_arguments)] // TODO
208    pub(super) fn new<R: Runtime>(
209        runtime: R,
210        channel: &Arc<Channel>,
211        circ_id: CircId,
212        unique_id: UniqId,
213        inner: B,
214        hops: Arc<RwLock<CircHopList>>,
215        forward_reactor_rx: mpsc::Receiver<BackwardReactorCmd>,
216        control_rx: mpsc::UnboundedReceiver<CtrlMsg<B::CtrlMsg>>,
217        command_rx: mpsc::UnboundedReceiver<CtrlCmd<B::CtrlCmd>>,
218        padding_ctrl: PaddingController,
219        padding_event_stream: PaddingEventStream,
220        stream_rx: mpsc::Receiver<ReadyStreamMsg>,
221    ) -> Self {
222        let channel = Arc::clone(channel);
223        let inbound_chan_tx = CircuitCellSender::from_channel_sender(channel.sender());
224
225        Self {
226            time_provider: DynTimeProvider::new(runtime),
227            outbound_chan_rx: None,
228            channel,
229            inner,
230            hops,
231            inbound_chan_tx,
232            unique_id,
233            circ_id,
234            forward_reactor_rx,
235            control_rx,
236            command_rx,
237            stream_rx,
238            padding_ctrl,
239            padding_event_stream,
240            #[cfg(feature = "circ-padding")]
241            padding_block: None,
242        }
243    }
244
245    /// Helper for [`run`](Self::run).
246    ///
247    /// Handles cells arriving on the outbound Tor channel,
248    /// and writes cells to the inbound Tor channel.
249    ///
250    /// Because the Tor application streams, the `forward_reactor_rx` MPSC streams,
251    /// and the outbound Tor channel MPSC stream are driven concurrently using [`PollAll`],
252    /// this function can send up to 3 cells per call over the inbound Tor channel:
253    ///
254    ///    * a cell carrying Tor stream data
255    ///    * a cell received from the outbound Tor channel, if we are a relay
256    ///      (moving from the exit towards the client)
257    ///    * a circuit-level SENDME
258    ///
259    /// However, in practice, leaky pipe is not really used,
260    /// and so relays that have application streams (i.e. the exits),
261    /// are not going to have an outbound Tor channel,
262    /// and so this will only really drive Tor stream data,
263    /// delivering at most 2 cells per call.
264    async fn run_once(&mut self) -> StdResult<(), ReactorError> {
265        use postage::prelude::{Sink as _, Stream as _};
266
267        /// The maximum number of events we expect to handle per reactor loop.
268        ///
269        /// This is bounded by the number of futures we push into the PollAll.
270        const PER_LOOP_EVENT_COUNT: usize = 3;
271
272        // A collection of futures we plan to drive concurrently.
273        let mut poll_all =
274            PollAll::<PER_LOOP_EVENT_COUNT, Option<CircuitEvent<B::CircChanMsg>>>::new();
275
276        // Flush the backward Tor channel sink, and check it for readiness
277        //
278        // TODO(flushing): here and everywhere else we need to flush:
279        //
280        // Currently, we try to flush every time we want to write to the sink,
281        // but may be suboptimal.
282        //
283        // However, we don't actually *wait* for the flush to complete
284        // (we just make a bit of progress by calling poll_flush),
285        // so it's possible that this is actually tolerable.
286        // We should run some tests, and if this turns out to be a performance bottleneck,
287        // we'll have to rethink our flushing approach.
288        let backward_chan_ready = future::poll_fn(|cx| {
289            // The flush outcome doesn't matter,
290            // so we simply move on to the readiness check.
291            // The reason we don't wait on the flush is because we don't
292            // want to flush on *every* reactor loop, but we do want to make
293            // a bit of progress each time.
294            //
295            // (TODO: do we want to handle errors here?)
296            let _ = self.inbound_chan_tx.poll_flush_unpin(cx);
297
298            self.inbound_chan_tx.poll_ready_unpin(cx)
299        });
300
301        // Concurrently, drive :
302        //  1. a future that reads from the StreamReactor, to see if there are
303        //  any application streams that have a message to send
304        //  (this resolves to a message that needs to be delivered to the peer)
305        poll_all.push(async {
306            // Internally, each stream reactor checks if we're allowed to send anything
307            // that counts towards SENDME windows (and ceases to send us stream data if not)
308            //
309            // The reason we don't check that here is because stream_rx multiplexes stream data
310            // from all hops, and we have no way of knowing which hop will want to send us stream
311            // data next, and therefore we can't know which hop's CC object to use
312            self.stream_rx.next().await.map(CircuitEvent::Send)
313        });
314
315        //  2. the stream of commands coming from the ForwardReactor
316        //  (this resolves to a BackwardReactorCmd)
317        poll_all.push(async {
318            let event = match self.forward_reactor_rx.next().await {
319                Some(cmd) => CircuitEvent::Forwarded(cmd),
320                None => {
321                    // The forward reactor has crashed, so we have to shut down.
322                    CircuitEvent::ForwardShutdown
323                }
324            };
325
326            Some(event)
327        });
328
329        // 3. Messages moving from the outbound channel towards the inbound Tor channel,
330        // if we have an outbound Tor channel.
331        //
332        // NOTE: in practice, clients and exits won't have an outbound Tor channel,
333        // so for them this will be a no-op.
334        poll_all.push(async {
335            let event = if let Some(outbound_chan_rx) = self.outbound_chan_rx.as_mut() {
336                // Forward channel unexpectedly closed, we should close too
337                match outbound_chan_rx.next().await {
338                    Some(msg) => match msg.try_into() {
339                        Err(e) => CircuitEvent::ProtoViolation(e),
340                        Ok(cell) => CircuitEvent::Cell(cell),
341                    },
342                    None => {
343                        // The forward reactor has crashed, so we have to shut down.
344                        CircuitEvent::ForwardShutdown
345                    }
346                }
347            } else {
348                future::pending().await
349            };
350
351            Some(event)
352        });
353
354        let poll_all = async move {
355            // Avoid polling **any** of the futures if the outgoing sink is blocked.
356            //
357            // This implements backpressure: we avoid reading from our input sources
358            // if we know we're unable to write to the inbound Tor channel sink.
359            //
360            // More specifically, if our inbound Tor channel sink is full and can no longer
361            // accept cells, we stop reading:
362            //
363            //   1. From the application streams (received from StreamReactor), if there are any.
364            //
365            //   2. From the forward_reactor_rx channel, used by the forward reactor to send us
366            //
367            //     - a circuit-level SENDME that we have received, or
368            //     - a circuit-level SENDME that we need to deliver to the client
369            //
370            //     Not reading from the forward_reactor_rx channel, in turn, causes the forward reactor
371            //     to block and therefore stop reading from **its** input sources,
372            //     propagating backpressure all the way to the other endpoint of the circuit.
373            //
374            //   3. From the outbound Tor channel, if there is one.
375            //
376            // This will delay any SENDMEs the client or exit might have sent along
377            // the way, and therefore count as a congestion signal.
378            //
379            // TODO: memquota setup to make sure this doesn't turn into a memory DOS vector
380            let _ = backward_chan_ready.await;
381
382            // TODO: it's important to not block reading from the forward_reactor_rx channel on the chan
383            // sender readiness (for instance, we should not block the sending of SENDMEs
384            // if the channel is blocked on a padding-induced block).
385            //
386            // This means we will need to move the forward_reactor_rx handling out of the PollAll
387            // to the select_biased! below.
388            poll_all.await
389        };
390
391        let events = select_biased! {
392            res = self.command_rx.next().fuse() => {
393                let cmd = res.ok_or_else(|| ReactorError::Shutdown)?;
394                self.handle_cmd(cmd)?;
395                return Ok(());
396            }
397            res = self.control_rx.next().fuse() => {
398                let msg = res.ok_or_else(|| ReactorError::Shutdown)?;
399                if let Some(new_msg) = self.handle_msg(msg)? {
400                    let mut events = <PollAll::<_, _> as Future>::Output::new();
401                    events.push(Some(CircuitEvent::Send(new_msg)));
402                    events
403                } else {
404                    return Ok(());
405                }
406            }
407            res = self.padding_event_stream.next().fuse() => {
408                // If there's a padding event, we need to handle it immediately,
409                // because it might tell us to start blocking the inbound_chan_tx sink,
410                // which, in turn, means we need to stop trying to read from
411                // the application streams.
412                let event = res.ok_or_else(|| ReactorError::Shutdown)?;
413
414                cfg_if::cfg_if! {
415                    if #[cfg(feature = "circ-padding")] {
416                        self.run_padding_event(event).await?;
417                    } else {
418                        // If padding isn't enabled, we never generate a padding event,
419                        // so we can be sure this case will never be called.
420                        void::unreachable(event.0);
421                    }
422                }
423                return Ok(())
424            }
425            res = poll_all.fuse() => res,
426        };
427
428        // Note: there shouldn't be more than N < PER_LOOP_EVENT_COUNT events to handle
429        // per reactor loop. We need to be careful here, because we must avoid blocking
430        // the reactor.
431        //
432        // If handling more than one event per loop turns out to be a problem, we may
433        // need to dispatch this to a background task instead.
434        //
435        // TODO(relay): this loop is actually a problem.
436        // As mentioned in the run_once() docs, this will attempt to send up
437        // to 3 cells on the inbound tor Channel (or 2 cells, assuming no leaky pipe).
438        //
439        // The problem is that the readiness check above (see backward_chan_ready)
440        // only checks that the queue has enough room for 1 cell, not *2 cells*.
441        // Trying to send more than 2 cell when there is only room for one
442        // will cause the reactor to block (and because there is nothing
443        // driving the flushing of this channel, this will be a hard block).
444        //
445        // We need to rethink the strategy here (e.g. by flushing in parallel
446        // with handle_event())
447        for event in events.into_iter().flatten() {
448            self.handle_event(event).await?;
449        }
450
451        Ok(())
452    }
453
454    /// Handle a control command.
455    fn handle_cmd(&mut self, cmd: CtrlCmd<B::CtrlCmd>) -> StdResult<(), ReactorError> {
456        match cmd {
457            CtrlCmd::Custom(c) => self.inner.handle_cmd(c),
458        }
459    }
460
461    /// Handle a control message.
462    ///
463    /// This may result in a new stream message that needs to be sent backward.
464    fn handle_msg(
465        &mut self,
466        msg: CtrlMsg<B::CtrlMsg>,
467    ) -> StdResult<Option<ReadyStreamMsg>, ReactorError> {
468        match msg {
469            CtrlMsg::Custom(c) => {
470                // In the future we may also want `inner.handle_msg(c)` to return an
471                // `Option<ReadyStreamMsg>`, and we can pass it through.
472                let () = self.inner.handle_msg(c)?;
473                Ok(None)
474            }
475            CtrlMsg::FlowCtrlUpdate {
476                hop,
477                stream_id,
478                msg,
479            } => match msg {
480                FlowCtrlMsg::Sendme => {
481                    // Congestion control decides if we can send stream level SENDMEs or not.
482                    let (cell_fmt, cc) = self.hop_info(hop)?;
483                    let uses_stream_sendme = cc.lock().expect("poisoned").uses_stream_sendme();
484
485                    if !uses_stream_sendme {
486                        // Nothing to do, so discard the SENDME.
487                        //
488                        // TODO(arti#2068): We should do something better here,
489                        // like ensure that nothing sends `FlowCtrlMsg::Sendme` when it shouldn't,
490                        // and making this an error instead.
491                        return Ok(None);
492                    }
493
494                    let sendme = Sendme::new_empty();
495                    let msg = AnyRelayMsgOuter::new(Some(stream_id), sendme.into());
496
497                    Ok(Some(ReadyStreamMsg {
498                        hop,
499                        msg,
500                        relay_cell_format: cell_fmt,
501                        ccontrol: Arc::clone(&cc),
502                    }))
503                }
504                FlowCtrlMsg::Xon(rate) => {
505                    todo!()
506                }
507            },
508        }
509    }
510
511    /// Perform some circuit-padding-based event on the specified circuit.
512    //
513    // TODO(DEDUP): this is almost identical to the client-side Conflux::run_padding_event()
514    #[cfg(feature = "circ-padding")]
515    async fn run_padding_event(
516        &mut self,
517        padding_event: PaddingEvent,
518    ) -> StdResult<(), ReactorError> {
519        use PaddingEvent as E;
520
521        match padding_event {
522            E::SendPadding(send_padding) => {
523                self.send_padding(send_padding).await?;
524            }
525            E::StartBlocking(start_blocking) => {
526                self.start_blocking_for_padding(start_blocking);
527            }
528            E::StopBlocking => {
529                self.stop_blocking_for_padding();
530            }
531        }
532        Ok(())
533    }
534
535    /// Handle a request from our padding subsystem to send a padding packet.
536    //
537    // TODO(DEDUP): this is almost identical to the client-side Client::send_padding()
538    #[cfg(feature = "circ-padding")]
539    async fn send_padding(&mut self, send_padding: padding::SendPadding) -> Result<()> {
540        use CircPaddingDisposition::*;
541
542        let target_hop = send_padding.hop;
543
544        match padding_disposition(
545            &send_padding,
546            &self.inbound_chan_tx,
547            self.padding_block.as_ref(),
548        ) {
549            QueuePaddingNormally => {
550                let queue_info = self.padding_ctrl.queued_padding(target_hop, send_padding);
551                self.queue_padding_cell_for_hop(target_hop, queue_info)
552                    .await?;
553            }
554            QueuePaddingAndBypass => {
555                let queue_info = self.padding_ctrl.queued_padding(target_hop, send_padding);
556                self.queue_padding_cell_for_hop(target_hop, queue_info)
557                    .await?;
558            }
559            TreatQueuedCellAsPadding => {
560                self.padding_ctrl
561                    .replaceable_padding_already_queued(target_hop, send_padding);
562            }
563        }
564        Ok(())
565    }
566
567    /// Enable padding-based blocking,
568    /// or change the rule for padding-based blocking to the one in `block`.
569    //
570    // TODO(DEDUP): copy of Client::start_blocking_for_padding()
571    #[cfg(feature = "circ-padding")]
572    pub(super) fn start_blocking_for_padding(&mut self, block: padding::StartBlocking) {
573        self.inbound_chan_tx.start_blocking();
574        self.padding_block = Some(block);
575    }
576
577    /// Disable padding-based blocking.
578    ///
579    // TODO(DEDUP): copy of Client::stop_blocking_for_padding()
580    #[cfg(feature = "circ-padding")]
581    pub(super) fn stop_blocking_for_padding(&mut self) {
582        self.inbound_chan_tx.stop_blocking();
583        self.padding_block = None;
584    }
585
586    /// Generate and encrypt a padding cell, and send it to a targeted hop.
587    ///
588    /// Ignores any padding-based blocking.
589    ///
590    // TODO(DEDUP): copy of Client::queue_padding_cell_for_hop()
591    #[cfg(feature = "circ-padding")]
592    async fn queue_padding_cell_for_hop(
593        &mut self,
594        target_hop: HopNum,
595        queue_info: Option<QueuedCellPaddingInfo>,
596    ) -> Result<()> {
597        use tor_cell::relaycell::msg::Drop as DropMsg;
598
599        let msg = AnyRelayMsgOuter::new(None, DropMsg::default().into());
600        let hopnum = Some(target_hop);
601
602        // TODO: the ccontrol state isn't actually needed here, because
603        // DROP cells don't count towards SENDME windows.
604        // Technically, we could avoid unnecessarily Arc::clone()ing the CC state
605        // here, and just extract the relay cell format.
606        // But for that we would need a specialized send_relay_cell_inner()-like function
607        // that doesn't take a CC object, or to make the CC object optional in
608        // send_relay_cell_inner().
609        let (relay_cell_format, ccontrol) = self.hop_info(hopnum)?;
610
611        self.send_relay_cell_inner(hopnum, relay_cell_format, msg, false, &ccontrol, queue_info)
612            .await
613    }
614
615    /// Determine how exactly to handle a request to handle padding.
616    #[cfg(feature = "circ-padding")]
617    fn padding_disposition(&self, send_padding: &padding::SendPadding) -> CircPaddingDisposition {
618        crate::circuit::padding::padding_disposition(
619            send_padding,
620            &self.inbound_chan_tx,
621            self.padding_block.as_ref(),
622        )
623    }
624
625    /// Handle a circuit event.
626    async fn handle_event(
627        &mut self,
628        event: CircuitEvent<B::CircChanMsg>,
629    ) -> StdResult<(), ReactorError> {
630        use CircuitEvent::*;
631
632        match event {
633            Cell(cell) => self.handle_backward_cell(cell).await,
634            Send(msg) => {
635                let ReadyStreamMsg {
636                    hop,
637                    relay_cell_format,
638                    msg,
639                    ccontrol,
640                } = msg;
641
642                self.send_relay_cell(hop, relay_cell_format, msg, false, &ccontrol)
643                    .await?;
644
645                Ok(())
646            }
647            Forwarded(cmd) => self.handle_reactor_cmd(cmd).await,
648            ForwardShutdown => {
649                // The forward reactor has crashed, so we have to shut down.
650                trace!(
651                    circ_uniq_id = %self.unique_id,
652                    backward_circ_id = %self.circ_id,
653                    "Backward relay reactor shutdown (forward reactor has closed)",
654                );
655
656                Err(ReactorError::Shutdown)
657            }
658            ProtoViolation(err) => Err(err.into()),
659        }
660    }
661
662    /// Return the RelayCellFormat and CC state of a given hop.
663    fn hop_info(
664        &self,
665        hopnum: Option<HopNum>,
666    ) -> Result<(RelayCellFormat, Arc<Mutex<CongestionControl>>)> {
667        let hops = self.hops.read().expect("poisoned lock");
668        let hop = hops
669            .get(hopnum)
670            .ok_or_else(|| internal!("tried to send padding to non-existent hop?!"))?;
671        let relay_cell_format = hop.settings.relay_crypt_protocol().relay_cell_format();
672        let ccontrol = Arc::clone(&hop.ccontrol);
673
674        Ok((relay_cell_format, ccontrol))
675    }
676
677    /// Handle a command sent to us by the forward reactor.
678    async fn handle_reactor_cmd(&mut self, msg: BackwardReactorCmd) -> StdResult<(), ReactorError> {
679        use BackwardReactorCmd::*;
680
681        match msg {
682            SendRelayMsg { hop, msg } => {
683                self.send_relay_msg(hop, msg).await?;
684            }
685            HandleSendme { hop, sendme } => {
686                self.handle_sendme(hop, sendme).await?;
687                return Ok(());
688            }
689            #[cfg(feature = "relay")]
690            HandleCircuitExtended {
691                hop,
692                extended2,
693                outbound_chan_rx,
694            } => {
695                self.outbound_chan_rx = Some(outbound_chan_rx);
696                let msg = AnyRelayMsgOuter::new(None, extended2.into());
697                self.send_relay_msg(hop, msg).await?;
698
699                debug!(
700                    circ_uniq_id = %self.unique_id,
701                    backward_circ_id = %self.circ_id,
702                    "Extended circuit to the next hop"
703                );
704            }
705        }
706
707        Ok(())
708    }
709
710    /// Send a relay message to the specified hop.
711    async fn send_relay_msg(
712        &mut self,
713        hopnum: Option<HopNum>,
714        msg: AnyRelayMsgOuter,
715    ) -> StdResult<(), ReactorError> {
716        let (relay_cell_format, ccontrol) = self.hop_info(hopnum)?;
717        let cmd = msg.cmd();
718
719        // TODO(relay): remove this log once we add some tests
720        // and confirm relaying cells works as expected
721        // (in practice it will be too noisy to be useful, even at trace level).
722        trace!(
723            circ_uniq_id = %self.unique_id,
724            backward_circ_id = %self.circ_id,
725            hopnum=?hopnum,
726            cmd = %cmd,
727            "Sending backward cell"
728        );
729
730        self.send_relay_cell(hopnum, relay_cell_format, msg, false, &ccontrol)
731            .await?;
732
733        if cmd == RelayCmd::SENDME {
734            ccontrol.lock().expect("poisoned lock").note_sendme_sent();
735        }
736
737        Ok(())
738    }
739
740    /// Handle a circuit-level SENDME (stream ID = 0).
741    ///
742    /// Returns an error if the SENDME does not have an authentication tag
743    /// (versions of Tor <=0.3.5 omit the SENDME tag, but we don't support
744    /// those any longer).
745    ///
746    /// Any error returned from this function will shut down the reactor.
747    ///
748    // TODO(DEDUP): duplicates the logic from the client-side Circuit::handle_sendme()
749    async fn handle_sendme(
750        &mut self,
751        hopnum: Option<HopNum>,
752        sendme: Sendme,
753    ) -> StdResult<(), ReactorError> {
754        let tag = sendme
755            .into_sendme_tag()
756            .ok_or_else(|| Error::CircProto("missing tag on circuit sendme".into()))?;
757
758        // NOTE: it's okay to await. We are only awaiting on the congestion_signals
759        // future which *should* resolve immediately
760        let signals = self.inbound_chan_tx.congestion_signals().await;
761
762        let hops = self.hops.read().expect("poisoned lock");
763        let hop = hops
764            .get(hopnum)
765            .ok_or_else(|| internal!("tried to send padding to non-existent hop?!"))?;
766
767        // Update the CC object that we received a SENDME along
768        // with possible congestion signals.
769        hop.ccontrol
770            .lock()
771            .expect("poisoned lock")
772            .note_sendme_received(&self.time_provider, tag, signals)?;
773
774        Ok(())
775    }
776
777    /// Encode `msg` and encrypt it, returning the resulting cell
778    /// and tag that should be expected for an authenticated SENDME sent
779    /// in response to that cell.
780    ///
781    // TODO(DEDUP): duplicates the logic from the client-side Circuit::encode_relay_cell()
782    fn encode_relay_cell(
783        &mut self,
784        relay_format: RelayCellFormat,
785        hop: Option<HopNum>,
786        early: bool,
787        msg: AnyRelayMsgOuter,
788    ) -> Result<(AnyChanMsg, SendmeTag)> {
789        let mut body: RelayCellBody = msg
790            .encode(relay_format, &mut rand::rng())
791            .map_err(|e| Error::from_cell_enc(e, "relay cell body"))?
792            .into();
793        let cmd = if early {
794            ChanCmd::RELAY_EARLY
795        } else {
796            ChanCmd::RELAY
797        };
798
799        // Use the implementation-dependent encryption logic
800        let tag = self.inner.encrypt_relay_cell(cmd, &mut body, hop);
801        let msg = Relay::from(BoxedCellBody::from(body));
802        let msg = if early {
803            AnyChanMsg::RelayEarly(msg.into())
804        } else {
805            AnyChanMsg::Relay(msg)
806        };
807
808        Ok((msg, tag))
809    }
810
811    /// Encode `msg`, encrypt it, and send it to the 'hop'th hop.
812    ///
813    /// If there is insufficient outgoing *circuit-level* or *stream-level*
814    /// SENDME window, an error is returned instead.
815    ///
816    /// Does not check whether the cell is well-formed or reasonable.
817    async fn send_relay_cell(
818        &mut self,
819        hop: Option<HopNum>,
820        relay_cell_format: RelayCellFormat,
821        msg: AnyRelayMsgOuter,
822        early: bool,
823        ccontrol: &Arc<Mutex<CongestionControl>>,
824    ) -> Result<()> {
825        self.send_relay_cell_inner(hop, relay_cell_format, msg, early, ccontrol, None)
826            .await
827    }
828
829    /// As [`send_relay_cell`](Self::send_relay_cell), but takes an optional
830    /// [`QueuedCellPaddingInfo`] in `padding_info`.
831    ///
832    /// If `padding_info` is None, `msg` must be non-padding: we report it as such to the
833    /// padding controller.
834    ///
835    // TODO(DEDUP): this contains parts of Circuit::send_relay_cell_inner()
836    async fn send_relay_cell_inner(
837        &mut self,
838        hop: Option<HopNum>,
839        relay_cell_format: RelayCellFormat,
840        msg: AnyRelayMsgOuter,
841        early: bool,
842        ccontrol: &Arc<Mutex<CongestionControl>>,
843        padding_info: Option<QueuedCellPaddingInfo>,
844    ) -> Result<()> {
845        let c_t_w = sendme::cmd_counts_towards_windows(msg.cmd());
846        let (msg, tag) = self.encode_relay_cell(relay_cell_format, hop, early, msg)?;
847        let cell = AnyChanCell::new(Some(self.circ_id), msg);
848
849        // TODO: we use HopNum(0) if we're a relay (i.e. if the hop is None).
850        // Is that ok?
851        let hop = hop.unwrap_or_else(|| HopNum::from(0));
852        // Remember that we've enqueued this cell.
853        let padding_info = padding_info.or_else(|| self.padding_ctrl.queued_data(hop));
854
855        // Note: this future is always `Ready`, because we checked the sink for readiness
856        // before polling the async streams, so await won't block.
857        Pin::new(&mut self.inbound_chan_tx)
858            .send_unbounded((cell, padding_info))
859            .await?;
860
861        if c_t_w {
862            ccontrol
863                .lock()
864                .expect("poisoned lock")
865                .note_data_sent(&self.time_provider, &tag)?;
866        }
867
868        Ok(())
869    }
870
871    /// Handle a backward cell (moving from the exit towards the client).
872    async fn handle_backward_cell(&mut self, cell: B::CircChanMsg) -> StdResult<(), ReactorError> {
873        match self
874            .inner
875            .handle_backward_cell(self.unique_id, self.circ_id, cell)?
876        {
877            BackwardCellDisposition::Forward(cell) => {
878                let cell = AnyChanCell::new(Some(self.circ_id), cell);
879                self.inbound_chan_tx
880                    .send((cell, None))
881                    .await
882                    .map_err(ReactorError::Err)
883            }
884        }
885    }
886}
887
888impl<B: BackwardHandler> Drop for BackwardReactor<B> {
889    fn drop(&mut self) {
890        // This will send a DESTROY down the inbound Tor channel
891        let _ = self.channel.close_circuit(self.circ_id);
892    }
893}
894
895/// A circuit event that must be handled by the [`BackwardReactor`].
896enum CircuitEvent<M> {
897    /// We received a cell that needs to be handled.
898    ///
899    /// (The cell is client-bound if we are a relay, or exit-bound if we are a client).
900    Cell(M),
901    /// A stream has a RELAY cell that needs
902    /// to be packaged and written to our Tor channel.
903    ///
904    /// (The message is client-bound if we are a relay, or exit-bound if we are a client).
905    Send(ReadyStreamMsg),
906    /// We received a cell from the ForwardReactor that we need to handle.
907    ///
908    /// This might be
909    ///
910    ///   * a circuit-level SENDME that we have received, or
911    ///   * a circuit-level SENDME that we need to deliver to the client
912    Forwarded(BackwardReactorCmd),
913    /// The forward reactor has shut down.
914    ///
915    /// We need to shut down too.
916    ForwardShutdown,
917    /// Protocol violation.
918    ///
919    /// This can happen if we receive a channel message that is not supported on the channel.
920    ProtoViolation(Error),
921}
922
923/// Instructions from the forward reactor.
924pub(crate) enum BackwardReactorCmd {
925    /// A circuit SENDME we received from the other endpoint.
926    HandleSendme {
927        /// The hop the SENDME came on.
928        hop: Option<HopNum>,
929        /// The SENDME.
930        sendme: Sendme,
931    },
932    /// A message we need to send back to the other endpoint.
933    SendRelayMsg {
934        /// The hop to encode the message for.
935        hop: Option<HopNum>,
936        /// The message to send.
937        msg: AnyRelayMsgOuter,
938    },
939    /// This relay circuit was extended by another hop.
940    ///
941    /// This causes the reactor send the `extended2` message on its inbound channel,
942    /// and start reading from `outbound_chan_rx` in the main loop.
943    //
944    ///
945    // TODO: I wish we didn't need to expose this relay-specific variant
946    // in the generic reactor but we have no choice: abstracting it away
947    // means either introducing a mutex between the relay-side forward/backward
948    // handlers, or yet another mpsc between them.
949    #[cfg(feature = "relay")]
950    HandleCircuitExtended {
951        /// The hop to encode the message for.
952        ///
953        /// In practice, this is always None, because only relays use this.
954        hop: Option<HopNum>,
955        /// The cell to send to the specified hop,
956        extended2: Extended2,
957        /// The reading end of the outbound Tor channel, if we are not the last hop.
958        ///
959        /// Yields cells moving from the exit towards the client, if we are a middle relay.
960        outbound_chan_rx: CircuitRxReceiver,
961    },
962}
963
964/// A flow control update message.
965///
966/// TODO(DEDUP): This is a duplicate of the client's
967/// `crate::client::reactor::control::FlowCtrlMsg`.
968#[derive(Debug)]
969pub(crate) enum FlowCtrlMsg {
970    /// Send a SENDME message on this stream.
971    Sendme,
972    /// Send an XON message on this stream with the given rate.
973    Xon(XonKBpsEwma),
974}