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}