tor_proto/circuit/reactor/forward.rs
1//! A circuit's view of the forward state of the circuit.
2
3use crate::circuit::UniqId;
4use crate::circuit::reactor::backward::BackwardReactorCmd;
5use crate::circuit::reactor::hop_mgr::HopMgr;
6use crate::circuit::reactor::macros::derive_deftly_template_CircuitReactor;
7use crate::circuit::reactor::stream;
8use crate::circuit::reactor::{ControlHandler, ReactorResultChannel};
9use crate::congestion::sendme;
10use crate::stream::cmdcheck::AnyCmdChecker;
11use crate::stream::msg_streamid;
12use crate::util::err::ReactorError;
13use crate::{Error, HopNum, Result};
14
15#[cfg(any(feature = "hs-service", feature = "relay"))]
16use {
17 crate::stream::CloseStreamBehavior,
18 crate::stream::incoming::{
19 IncomingStreamRequestFilter, IncomingStreamRequestHandler, StreamReqSender,
20 },
21 tor_cell::relaycell::StreamId,
22};
23
24// TODO(circpad): once padding is stabilized, the padding module will be moved out of client.
25use crate::client::circuit::padding::PaddingController;
26
27use tor_cell::chancell::{CircId, msg::AnyChanMsg};
28use tor_cell::relaycell::msg::{Sendme, SendmeTag};
29use tor_cell::relaycell::{
30 AnyRelayMsgOuter, RelayCellDecoderResult, RelayCellFormat, RelayCmd, UnparsedRelayMsg,
31};
32use tor_error::internal;
33use tor_linkspec::HasRelayIds;
34use tor_rtcompat::Runtime;
35
36use derive_deftly::Deftly;
37use futures::SinkExt;
38use futures::channel::mpsc;
39use futures::{FutureExt as _, StreamExt, select_biased};
40use tracing::debug;
41
42use std::result::Result as StdResult;
43
44use crate::circuit::CircuitRxReceiver;
45
46/// The forward circuit reactor.
47///
48/// See the [`reactor`](crate::circuit::reactor) module-level docs.
49///
50/// Shuts downs down if an error occurs, or if either the [`Reactor`](super::Reactor)
51/// or the [`BackwardReactor`](super::BackwardReactor) shuts down:
52///
53/// * if the `Reactor` shuts down, we are alerted via the ctrl/command mpsc channels
54/// (their sending ends will close, which causes run_once() to return ReactorError::Shutdown)
55/// * if `BackwardReactor` shuts down, the `Reactor` will notice and will itself shut down,
56/// which, in turn, causes the `ForwardReactor` to shut down as described above
57#[derive(Deftly)]
58#[derive_deftly(CircuitReactor)]
59#[deftly(reactor_name = "forward reactor")]
60#[deftly(run_inner_fn = "Self::run_once")]
61#[must_use = "If you don't call run() on a reactor, the circuit won't work."]
62pub(super) struct ForwardReactor<R: Runtime, F: ForwardHandler> {
63 /// A handle to the runtime.
64 runtime: R,
65 /// An identifier for logging about this reactor's circuit.
66 unique_id: UniqId,
67 /// The circuit identifier on the inbound Tor channel.
68 circ_id: CircId,
69 /// Implementation-dependent part of the reactor.
70 ///
71 /// This enables us to customize the behavior of the reactor,
72 /// depending on whether we are a client or a relay.
73 inner: F,
74 /// Channel for receiving control commands.
75 command_rx: mpsc::UnboundedReceiver<CtrlCmd<F::CtrlCmd>>,
76 /// Channel for receiving control messages.
77 control_rx: mpsc::UnboundedReceiver<CtrlMsg<F::CtrlMsg>>,
78 /// The reading end of the inbound Tor channel.
79 ///
80 /// Yields cells moving from the client towards the exit, if we are a relay,
81 /// or cells moving towards *us*, if we are a client.
82 inbound_chan_rx: CircuitRxReceiver,
83 /// Sender for sending commands to the BackwardReactor.
84 ///
85 /// Used for sending:
86 ///
87 /// * circuit-level SENDMEs received from the other endpoint
88 /// (`[BackwardReactorCmd::HandleSendme]`)
89 /// * circuit-level SENDMEs that need to be delivered to the other endpoint
90 /// (using `[BackwardReactorCmd::SendRelayMsg]`)
91 ///
92 /// The receiver is in [`BackwardReactor`](super::BackwardReactor), which is responsible for
93 /// sending cell over the inbound channel.
94 backward_reactor_tx: mpsc::Sender<BackwardReactorCmd>,
95 /// Hop manager, storing per-hop state, and handles to the stream reactors.
96 ///
97 /// Contains the `CircHopList`.
98 hop_mgr: HopMgr<R>,
99 /// An implementation-specific event stream.
100 ///
101 /// Polled from the main loop of the reactor.
102 /// Each event is passed to [`ForwardHandler::handle_event`].
103 circ_events: mpsc::Receiver<F::CircEvent>,
104 /// A padding controller to which padding-related events should be reported.
105 padding_ctrl: PaddingController,
106}
107
108/// A control command aimed at the generic forward reactor.
109pub(crate) enum CtrlCmd<C> {
110 /// Begin accepting streams on this circuit.
111 //
112 // TODO(DEDUP): this is very similar to its client-side counterpart,
113 // except the hop is a Option<HopNum> instead of a TargetHop.
114 #[cfg(any(feature = "hs-service", feature = "relay"))]
115 #[allow(unused)] // TODO(dedup): this will be used by hs services
116 AwaitStreamRequests {
117 /// A channel for sending information about an incoming stream request.
118 incoming_sender: StreamReqSender,
119 /// A `CmdChecker` to keep track of which message types are acceptable.
120 cmd_checker: AnyCmdChecker,
121 /// Oneshot channel to notify on completion.
122 done: ReactorResultChannel<()>,
123 /// The hop that is allowed to create streams.
124 ///
125 /// Set to None if we are a relay wanting to accept stream requests.
126 hop: Option<HopNum>,
127 /// A filter used to check requests before passing them on.
128 filter: Box<dyn IncomingStreamRequestFilter>,
129 },
130
131 /// Close the specified pending incoming stream, sending the provided END message.
132 ///
133 /// A stream is said to be pending if the message for initiating the stream was received but
134 /// not has not been responded to yet.
135 ///
136 /// This should be used by responders for closing pending incoming streams initiated by the
137 /// other party on the circuit.
138 ///
139 /// TODO(dedup): this is almost identical to the ClosePendingStream control message
140 /// from the client-side. We can get rid of the duplication by rewriting
141 /// the client circuit reactor to use the new multi-reactor architecture
142 #[cfg(any(feature = "hs-service", feature = "relay"))]
143 ClosePendingStream {
144 /// The hop number the stream is on.
145 ///
146 /// Set to None if we are a relay.
147 hop: Option<HopNum>,
148 /// The stream ID to send the END for.
149 stream_id: StreamId,
150 /// The END message to send, if any.
151 message: CloseStreamBehavior,
152 /// Oneshot channel to notify on completion.
153 done: ReactorResultChannel<()>,
154 },
155 /// An implementation-dependent control command.
156 #[allow(unused)] // TODO(relay)
157 Custom(C),
158}
159
160/// A control message aimed at the generic forward reactor.
161pub(crate) enum CtrlMsg<M> {
162 /// An implementation-dependent control message.
163 #[allow(unused)] // TODO(relay)
164 Custom(M),
165}
166
167/// Trait for customizing the behavior of the forward reactor.
168///
169/// Used for plugging in the implementation-dependent (client vs relay)
170/// parts of the implementation into the generic one.
171pub(crate) trait ForwardHandler: ControlHandler {
172 /// Type that explains how to build an outgoing channel.
173 type BuildSpec: HasRelayIds;
174
175 /// The subclass of ChanMsg that can arrive on this type of circuit.
176 type CircChanMsg: TryFrom<AnyChanMsg, Error = crate::Error>;
177
178 /// An opaque event type.
179 ///
180 /// The [`ForwardReactor`] polls an MPSC stream yielding `CircEvent`s from the main loop.
181 /// Each event is passed to [`Self::handle_event`] for handling.
182 type CircEvent;
183
184 /// Handle a non-SENDME RELAY message on this circuit with stream ID 0.
185 async fn handle_meta_msg<R: Runtime>(
186 &mut self,
187 runtime: &R,
188 early: bool,
189 hopnum: Option<HopNum>,
190 msg: UnparsedRelayMsg,
191 relay_cell_format: RelayCellFormat,
192 ) -> StdResult<(), ReactorError>;
193
194 /// Handle a forward (TODO terminology) cell.
195 ///
196 /// The cell is
197 /// - moving from the client towards the exit, if we're a relay
198 /// - moving from the guard towards us, if we're a client
199 ///
200 /// Returns an error if the cell should cause the reactor to shut down,
201 /// or a [`ForwardCellDisposition`] specifying how it should be handled.
202 ///
203 /// Returns `None` if the cell was handled internally by this handler.
204 async fn handle_forward_cell<R: Runtime>(
205 &mut self,
206 hop_mgr: &mut HopMgr<R>,
207 cell: Self::CircChanMsg,
208 ) -> StdResult<Option<ForwardCellDisposition>, ReactorError>;
209
210 /// Handle an implementation-specific circuit event.
211 ///
212 /// Returns a command for the backward reactor.
213 fn handle_event(
214 &mut self,
215 event: Self::CircEvent,
216 ) -> StdResult<Option<BackwardReactorCmd>, ReactorError>;
217
218 /// Wait until the outbound channel, if there is one, is ready to accept more cells.
219 ///
220 /// Resolves immediately if there is no outbound channel.
221 /// Blocks if there is a pending outbound channel.
222 async fn outbound_chan_ready(&mut self) -> Result<()>;
223}
224
225/// What action to take in response to a cell arriving on our inbound Tor channel.
226pub(crate) enum ForwardCellDisposition {
227 /// Handle a decoded RELAY or RELAY_EARLY cell in the [`ForwardReactor`].
228 HandleRecognizedRelay {
229 /// The decoded cell.
230 cell: RelayCellDecoderResult,
231 /// Whether this was a RELAY_EARLY.
232 early: bool,
233 /// The hop this cell was for.
234 hopnum: Option<HopNum>,
235 /// The SENDME tag.
236 tag: SendmeTag,
237 },
238}
239
240impl<R: Runtime, F: ForwardHandler> ForwardReactor<R, F> {
241 /// Create a new [`ForwardReactor`].
242 #[allow(clippy::too_many_arguments)] // TODO
243 pub(super) fn new(
244 runtime: R,
245 unique_id: UniqId,
246 circ_id: CircId,
247 inner: F,
248 hop_mgr: HopMgr<R>,
249 inbound_chan_rx: CircuitRxReceiver,
250 control_rx: mpsc::UnboundedReceiver<CtrlMsg<F::CtrlMsg>>,
251 command_rx: mpsc::UnboundedReceiver<CtrlCmd<F::CtrlCmd>>,
252 backward_reactor_tx: mpsc::Sender<BackwardReactorCmd>,
253 circ_events: mpsc::Receiver<F::CircEvent>,
254 padding_ctrl: PaddingController,
255 ) -> Self {
256 Self {
257 runtime,
258 unique_id,
259 circ_id,
260 inbound_chan_rx,
261 control_rx,
262 command_rx,
263 inner,
264 backward_reactor_tx,
265 hop_mgr,
266 circ_events,
267 padding_ctrl,
268 }
269 }
270
271 /// Helper for [`run`](Self::run).
272 async fn run_once(&mut self) -> StdResult<(), ReactorError> {
273 let outbound_chan_ready = self.inner.outbound_chan_ready();
274
275 let inbound_chan_rx_fut = async {
276 // Avoid reading from the inbound_chan_rx Tor Channel if the outgoing sink is blocked
277 outbound_chan_ready.await?;
278 Ok(self.inbound_chan_rx.next().await)
279 };
280
281 select_biased! {
282 res = self.command_rx.next().fuse() => {
283 let cmd = res.ok_or_else(|| ReactorError::Shutdown)?;
284 self.handle_cmd(cmd).await
285 }
286 res = self.control_rx.next().fuse() => {
287 let msg = res.ok_or_else(|| ReactorError::Shutdown)?;
288 self.handle_msg(msg)
289 }
290 res = self.circ_events.next().fuse() => {
291 let ev = res.ok_or_else(|| ReactorError::Shutdown)?;
292 if let Some(cmd) = self.inner.handle_event(ev)? {
293 self.send_reactor_cmd(cmd).await?;
294 }
295
296 Ok(())
297 }
298 res = inbound_chan_rx_fut.fuse() => {
299 let cell = res.map_err(ReactorError::Err)?;
300 let Some(cell) = cell else {
301 debug!(
302 circ_uniq_id = %self.unique_id,
303 backward_circ_id = %self.circ_id,
304 "Backward channel has closed, shutting down forward relay reactor",
305 );
306
307 return Err(ReactorError::Shutdown);
308 };
309
310 let cell: F::CircChanMsg = cell.try_into()?;
311 let Some(disp) = self.inner.handle_forward_cell(&mut self.hop_mgr, cell).await? else {
312 return Ok(());
313 };
314
315 match disp {
316 ForwardCellDisposition::HandleRecognizedRelay { cell, early, hopnum, tag } => {
317 self.handle_relay_cell(cell, early, hopnum, tag).await
318 }
319 }
320 },
321 }
322 }
323
324 /// Handle a control command.
325 #[allow(clippy::unused_async)] // used if any(feature = "hs-service", feature = "relay")
326 async fn handle_cmd(&mut self, cmd: CtrlCmd<F::CtrlCmd>) -> StdResult<(), ReactorError> {
327 match cmd {
328 #[cfg(any(feature = "hs-service", feature = "relay"))]
329 CtrlCmd::AwaitStreamRequests {
330 incoming_sender,
331 cmd_checker,
332 done,
333 hop,
334 filter,
335 } => {
336 let handler = IncomingStreamRequestHandler {
337 incoming_sender,
338 cmd_checker,
339 hop_num: hop,
340 filter,
341 };
342
343 // Update the HopMgr with the
344 let ret = self.hop_mgr.set_incoming_handler(handler);
345 let _ = done.send(ret); // don't care if the corresponding receiver goes away.
346 Ok(())
347 }
348 #[cfg(any(feature = "hs-service", feature = "relay"))]
349 CtrlCmd::ClosePendingStream {
350 hop,
351 stream_id,
352 message,
353 done,
354 } => {
355 let ret = self.hop_mgr.close_pending(hop, stream_id, message).await;
356 let _ = done.send(ret); // don't care if the corresponding receiver goes away.
357
358 Ok(())
359 }
360 CtrlCmd::Custom(c) => self.inner.handle_cmd(c),
361 }
362 }
363
364 /// Handle a control message.
365 fn handle_msg(&mut self, msg: CtrlMsg<F::CtrlMsg>) -> StdResult<(), ReactorError> {
366 match msg {
367 CtrlMsg::Custom(c) => self.inner.handle_msg(c),
368 }
369 }
370
371 /// Note that we have received a RELAY cell.
372 ///
373 /// Updates the padding and CC state.
374 fn note_relay_cell_received(
375 &self,
376 hopnum: Option<HopNum>,
377 c_t_w: bool,
378 ) -> Result<(RelayCellFormat, bool)> {
379 let mut hops = self.hop_mgr.hops().write().expect("poisoned lock");
380 let hop = hops
381 .get_mut(hopnum)
382 .ok_or_else(|| internal!("msg from non-existent hop???"))?;
383
384 // Check whether we are allowed to receive more data for this circuit hop.
385 hop.inbound.decrement_cell_limit()?;
386
387 // Decrement the circuit sendme windows, and see if we need to
388 // send a sendme cell.
389 let send_circ_sendme = if c_t_w {
390 hop.ccontrol
391 .lock()
392 .expect("poisoned lock")
393 .note_data_received()?
394 } else {
395 false
396 };
397
398 let relay_cell_format = hop.settings.relay_crypt_protocol().relay_cell_format();
399
400 Ok((relay_cell_format, send_circ_sendme))
401 }
402
403 /// Handle a RELAY cell.
404 ///
405 // TODO(DEDUP): very similar to Client::handle_relay_cell()
406 async fn handle_relay_cell(
407 &mut self,
408 decode_res: RelayCellDecoderResult,
409 early: bool,
410 hopnum: Option<HopNum>,
411 tag: SendmeTag,
412 ) -> StdResult<(), ReactorError> {
413 // For padding purposes, if we are a relay, we set the hopnum to 0
414 // TODO(relay): is this right?
415 let hopnum_padding = hopnum.unwrap_or_else(|| HopNum::from(0));
416 if decode_res.is_padding() {
417 self.padding_ctrl.decrypted_padding(hopnum_padding)?;
418 } else {
419 self.padding_ctrl.decrypted_data(hopnum_padding);
420 }
421
422 let c_t_w = decode_res.cmds().any(sendme::cmd_counts_towards_windows);
423 let (relay_cell_format, send_circ_sendme) = self.note_relay_cell_received(hopnum, c_t_w)?;
424
425 // If we do need to send a circuit-level SENDME cell, do so.
426 if send_circ_sendme {
427 // This always sends a V1 (tagged) sendme cell, and thereby assumes
428 // that SendmeEmitMinVersion is no more than 1. If the authorities
429 // every increase that parameter to a higher number, this will
430 // become incorrect. (Higher numbers are not currently defined.)
431 let sendme = Sendme::from(tag);
432 let msg = AnyRelayMsgOuter::new(None, sendme.into());
433 let forward = BackwardReactorCmd::SendRelayMsg { hop: hopnum, msg };
434
435 // NOTE: sending the SENDME to the backward reactor for handling
436 // might seem counterintuitive, given that we have access to
437 // the congestion control object right here (via hop_mgr).
438 //
439 // However, the forward reactor does not have access to the
440 // outbound_chan_tx part of the inbound (towards the client) Tor channel,
441 // and so it cannot handle the SENDME on its own
442 // (because it cannot obtain the congestion signals),
443 // so the SENDME needs to be handled in the backward reactor.
444 //
445 // NOTE: this will block if the backward reactor is not ready
446 // to send any more cells.
447 self.send_reactor_cmd(forward).await?;
448 }
449
450 let (mut msgs, incomplete) = decode_res.into_parts();
451 while let Some(msg) = msgs.next() {
452 match self
453 .handle_relay_msg(early, hopnum, msg, relay_cell_format, c_t_w)
454 .await
455 {
456 Ok(()) => continue,
457 Err(e) => {
458 for m in msgs {
459 debug!(
460 circ_uniq_id = %self.unique_id,
461 backward_circ_id = %self.circ_id,
462 "Ignoring relay msg received after triggering shutdown: {m:?}",
463 );
464 }
465 if let Some(incomplete) = incomplete {
466 debug!(
467 circ_uniq_id = %self.unique_id,
468 backward_circ_id = %self.circ_id,
469 "Ignoring partial relay msg received after triggering shutdown: {:?}",
470 incomplete,
471 );
472 }
473
474 return Err(e);
475 }
476 }
477 }
478
479 Ok(())
480 }
481
482 /// Handle a single incoming RELAY message.
483 async fn handle_relay_msg(
484 &mut self,
485 early: bool,
486 hop: Option<HopNum>,
487 msg: UnparsedRelayMsg,
488 relay_cell_format: RelayCellFormat,
489 cell_counts_toward_windows: bool,
490 ) -> StdResult<(), ReactorError> {
491 // If this msg wants/refuses to have a Stream ID, does it
492 // have/not have one?
493 let streamid = msg_streamid(&msg)?;
494
495 // If this doesn't have a StreamId, it's a meta cell,
496 // not meant for a particular stream.
497 let Some(sid) = streamid else {
498 return self
499 .handle_meta_msg(early, hop, msg, relay_cell_format)
500 .await;
501 };
502
503 let msg = stream::CtrlMsg::DeliverStreamMsg {
504 sid,
505 msg,
506 cell_counts_toward_windows,
507 };
508
509 // All messages on streams are handled in the stream reactor
510 // (because that's where the stream map is)
511 //
512 // Internally, this will spawn a StreamReactor for the target hop,
513 // if not already spawned.
514 self.hop_mgr.send(hop, msg).await
515 }
516
517 /// Handle a RELAY or RELAY_EARLY message on this circuit with stream ID 0.
518 async fn handle_meta_msg(
519 &mut self,
520 early: bool,
521 hopnum: Option<HopNum>,
522 msg: UnparsedRelayMsg,
523 relay_cell_format: RelayCellFormat,
524 ) -> StdResult<(), ReactorError> {
525 match msg.cmd() {
526 RelayCmd::SENDME => {
527 let sendme = msg
528 .decode::<Sendme>()
529 .map_err(|e| Error::from_bytes_err(e, "sendme message"))?
530 .into_msg();
531
532 let cmd = BackwardReactorCmd::HandleSendme {
533 hop: hopnum,
534 sendme,
535 };
536
537 self.send_reactor_cmd(cmd).await
538 }
539 _ => {
540 self.inner
541 .handle_meta_msg(&self.runtime, early, hopnum, msg, relay_cell_format)
542 .await
543 }
544 }
545 }
546
547 /// Send a command to the backward reactor.
548 ///
549 /// Blocks if the `backward_reactor_tx` channel is full, i.e. if the backward reactor
550 /// is not ready to send any more cells.
551 ///
552 /// Returns an error if the backward reactor has shut down.
553 async fn send_reactor_cmd(
554 &mut self,
555 forward: BackwardReactorCmd,
556 ) -> StdResult<(), ReactorError> {
557 self.backward_reactor_tx.send(forward).await.map_err(|_| {
558 // The other reactor has shut down
559 ReactorError::Shutdown
560 })
561 }
562}