Skip to main content

tor_proto/relay/reactor/
forward.rs

1//! A relay's view of the forward (away from the client, towards the exit) state of a circuit.
2
3mod extend_handler;
4
5use extend_handler::ExtendRequestHandler;
6
7use crate::channel::{Channel, ChannelSender};
8use crate::circuit::CircuitRxReceiver;
9use crate::circuit::UniqId;
10use crate::circuit::celltypes::RelayMaybeEarlyChanMsg;
11use crate::circuit::reactor::ControlHandler;
12use crate::circuit::reactor::backward::BackwardReactorCmd;
13use crate::circuit::reactor::forward::{ForwardCellDisposition, ForwardHandler};
14use crate::circuit::reactor::hop_mgr::HopMgr;
15use crate::crypto::cell::OutboundRelayLayer;
16use crate::crypto::cell::RelayCellBody;
17use crate::relay::RelayCircChanMsg;
18use crate::util::err::ReactorError;
19use crate::{Error, HopNum, Result};
20
21// TODO(circpad): once padding is stabilized, the padding module will be moved out of client.
22use crate::client::circuit::padding::QueuedCellPaddingInfo;
23
24use crate::relay::channel_provider::ChannelProvider;
25use crate::relay::reactor::CircuitAccount;
26use tor_cell::chancell::msg::{AnyChanMsg, Destroy, PaddingNegotiate, Relay};
27use tor_cell::chancell::{AnyChanCell, BoxedCellBody, ChanMsg, CircId};
28use tor_cell::relaycell::msg::{Extended2, SendmeTag};
29use tor_cell::relaycell::{RelayCellDecoderResult, RelayCellFormat, RelayCmd, UnparsedRelayMsg};
30use tor_error::internal;
31use tor_linkspec::OwnedChanTarget;
32use tor_rtcompat::Runtime;
33
34use futures::channel::mpsc;
35use futures::{SinkExt as _, future};
36use tracing::{debug, trace};
37
38use std::result::Result as StdResult;
39use std::sync::Arc;
40use std::task::Poll;
41
42/// Placeholder for our custom control message type.
43type CtrlMsg = ();
44
45/// Placeholder for our custom control command type.
46type CtrlCmd = ();
47
48/// The maximum number of RELAY_EARLY cells allowed on a circuit.
49///
50// TODO(relay): should we come up with a consensus parameter for this? (arti#2349)
51const MAX_RELAY_EARLY_CELLS_PER_CIRCUIT: usize = 8;
52
53/// Relay-specific state for the forward reactor.
54pub(crate) struct Forward {
55    /// An identifier for logging about this reactor's circuit.
56    unique_id: UniqId,
57    /// The circuit identifier on the inbound Tor channel.
58    circ_id: CircId,
59    /// The outbound view of this circuit, if we are not the last hop.
60    ///
61    /// Delivers cells towards the exit.
62    ///
63    /// Only set for middle relays.
64    outbound: Option<Outbound>,
65    /// The cryptographic state for this circuit for inbound cells.
66    crypto_out: Box<dyn OutboundRelayLayer + Send>,
67    /// The number of RELAY_EARLY cells we have seen so far on this circuit.
68    ///
69    /// If we see more than [`MAX_RELAY_EARLY_CELLS_PER_CIRCUIT`] RELAY_EARLY cells, we tear down the circuit.
70    relay_early_count: usize,
71    /// Helper for handling circuit extension requests.
72    ///
73    /// Used for validating EXTEND2 cells.
74    extend_handler: ExtendRequestHandler,
75}
76
77/// A type of event issued by the relay forward reactor.
78pub(crate) enum CircEvent {
79    /// The outcome of an EXTEND2 request.
80    ExtendResult(StdResult<ExtendResult, ReactorError>),
81}
82
83/// A successful circuit extension result.
84pub(crate) struct ExtendResult {
85    /// The EXTENDED2 cell to send back to the client.
86    extended2: Extended2,
87    /// The outbound channel.
88    outbound: Outbound,
89    /// The reading end of the outbound Tor channel, if we are not the last hop.
90    ///
91    /// Yields cells moving from the exit towards the client, if we are a middle relay.
92    outbound_chan_rx: CircuitRxReceiver,
93}
94
95/// The outbound view of a relay circuit.
96struct Outbound {
97    /// The circuit identifier on the outbound Tor channel.
98    circ_id: CircId,
99    /// The outbound Tor channel.
100    channel: Arc<Channel>,
101    /// The sending end of the outbound Tor channel.
102    outbound_chan_tx: ChannelSender,
103}
104
105/// The outcome of `decode_relay_cell`.
106enum CellDecodeResult {
107    /// A decrypted cell.
108    Recognized(SendmeTag, RelayCellDecoderResult),
109    /// A cell we could not decrypt.
110    Unrecognizd(RelayCellBody),
111}
112
113impl Forward {
114    /// Create a new [`Forward`].
115    pub(crate) fn new(
116        inbound_chan: &Arc<Channel>,
117        circ_id: CircId,
118        unique_id: UniqId,
119        crypto_out: Box<dyn OutboundRelayLayer + Send>,
120        chan_provider: Arc<dyn ChannelProvider<BuildSpec = OwnedChanTarget> + Send + Sync>,
121        event_tx: mpsc::Sender<CircEvent>,
122        memquota: CircuitAccount,
123    ) -> Self {
124        let inbound_peer = Arc::clone(inbound_chan.peer_info());
125        let extend_handler = ExtendRequestHandler::new(
126            unique_id,
127            circ_id,
128            chan_provider,
129            inbound_peer,
130            event_tx,
131            memquota,
132        );
133
134        Self {
135            unique_id,
136            circ_id,
137            // Initially, we are the last hop in the circuit.
138            outbound: None,
139            crypto_out,
140            relay_early_count: 0,
141            extend_handler,
142        }
143    }
144
145    /// Decode `cell`, returning its corresponding hop number, tag and decoded body.
146    fn decode_relay_cell<R: Runtime>(
147        &mut self,
148        hop_mgr: &mut HopMgr<R>,
149        cell: RelayMaybeEarlyChanMsg,
150    ) -> Result<(Option<HopNum>, CellDecodeResult)> {
151        // Note: the client reactor will return the actual source hopnum
152        let hopnum = None;
153        let cmd = cell.cmd();
154        let mut body = cell.into_relay_body().into();
155        let Some(tag) = self.crypto_out.decrypt_outbound(cmd, &mut body) else {
156            return Ok((hopnum, CellDecodeResult::Unrecognizd(body)));
157        };
158
159        // The message is addressed to us! Now it's time to handle it...
160        let mut hops = hop_mgr.hops().write().expect("poisoned lock");
161        let decode_res = hops
162            .get_mut(hopnum)
163            .ok_or_else(|| internal!("msg from non-existent hop???"))?
164            .inbound
165            .decode(body.into())?;
166
167        Ok((hopnum, CellDecodeResult::Recognized(tag, decode_res)))
168    }
169
170    /// Handle a DROP message.
171    #[allow(clippy::unnecessary_wraps)] // Returns Err if circ-padding is enabled
172    fn handle_drop(&mut self) -> StdResult<(), ReactorError> {
173        cfg_if::cfg_if! {
174            if #[cfg(feature = "circ-padding")] {
175                Err(internal!("relay circuit padding not yet supported").into())
176            } else {
177                Ok(())
178            }
179        }
180    }
181
182    /// Handle the outcome of handling an EXTEND2.
183    fn handle_extend_result(
184        &mut self,
185        res: StdResult<ExtendResult, ReactorError>,
186    ) -> StdResult<Option<BackwardReactorCmd>, ReactorError> {
187        let ExtendResult {
188            extended2,
189            outbound,
190            outbound_chan_rx,
191        } = res?;
192
193        self.outbound = Some(outbound);
194
195        Ok(Some(BackwardReactorCmd::HandleCircuitExtended {
196            hop: None,
197            extended2,
198            outbound_chan_rx,
199        }))
200    }
201
202    /// Handle a RELAY or RELAY_EARLY cell.
203    fn handle_relay_cell<R: Runtime>(
204        &mut self,
205        hop_mgr: &mut HopMgr<R>,
206        cell: RelayMaybeEarlyChanMsg,
207    ) -> StdResult<Option<ForwardCellDisposition>, ReactorError> {
208        let early = matches!(cell, RelayMaybeEarlyChanMsg::RelayEarly(_));
209
210        if early {
211            self.relay_early_count += 1;
212
213            if self.relay_early_count > MAX_RELAY_EARLY_CELLS_PER_CIRCUIT {
214                return Err(
215                    Error::CircProto("Circuit received too many RELAY_EARLY cells".into()).into(),
216                );
217            }
218        }
219
220        let (hopnum, res) = self.decode_relay_cell(hop_mgr, cell)?;
221        let (tag, decode_res) = match res {
222            CellDecodeResult::Unrecognizd(body) => {
223                self.handle_unrecognized_cell(body, None, early)?;
224                return Ok(None);
225            }
226            CellDecodeResult::Recognized(tag, res) => (tag, res),
227        };
228
229        Ok(Some(ForwardCellDisposition::HandleRecognizedRelay {
230            cell: decode_res,
231            early,
232            hopnum,
233            tag,
234        }))
235    }
236
237    /// Handle a forward cell that we could not decrypt.
238    fn handle_unrecognized_cell(
239        &mut self,
240        body: RelayCellBody,
241        info: Option<QueuedCellPaddingInfo>,
242        early: bool,
243    ) -> StdResult<(), ReactorError> {
244        let Some(chan) = self.outbound.as_mut() else {
245            // The client shouldn't try to send us any cells before it gets
246            // an EXTENDED2 cell from us
247            return Err(Error::CircProto(
248                "Asked to forward cell before the circuit was extended?!".into(),
249            )
250            .into());
251        };
252
253        // TODO(relay): remove this log once we add some tests
254        // and confirm relaying cells works as expected
255        // (in practice it will be too noisy to be useful, even at trace level).
256        trace!(
257            circ_uniq_id = %self.unique_id,
258            forward_circ_id = %chan.circ_id,
259            "Forwarding unrecognized cell"
260        );
261
262        let msg = Relay::from(BoxedCellBody::from(body));
263        let relay = if early {
264            AnyChanMsg::RelayEarly(msg.into())
265        } else {
266            AnyChanMsg::Relay(msg)
267        };
268        let cell = AnyChanCell::new(Some(chan.circ_id), relay);
269
270        // Note: this future is always `Ready`, because we checked the sink for readiness
271        // before polling the input channel, so await won't block.
272        chan.outbound_chan_tx.start_send_unpin((cell, info))?;
273
274        Ok(())
275    }
276
277    /// Handle a TRUNCATE cell.
278    fn handle_truncate(&mut self) -> StdResult<(), ReactorError> {
279        // This is not strictly spec compliant,
280        // but since none of our implementations use TRUNCATE,
281        // we deem it a proto violation and shut down the circuit.
282        //
283        // TODO(spec): codify this in the spec
284        Err(Error::CircProto("TRUNCATE not allowed".into()).into())
285    }
286
287    /// Handle a DESTROY cell originating from the client.
288    fn handle_destroy_cell(&mut self, cell: &Destroy) -> StdResult<(), ReactorError> {
289        debug!(
290            circ_uniq_id = %self.unique_id,
291            backward_circ_id = %self.circ_id,
292            reason = %cell.reason(),
293            "Received outbound DESTROY, circuit shutting down",
294        );
295
296        // We don't need to send a DESTROY cell down the channel,
297        // because that's handled implicitly by our Drop implementation
298        Err(ReactorError::Shutdown)
299    }
300
301    /// Handle a PADDING_NEGOTIATE cell originating from the client.
302    #[allow(clippy::needless_pass_by_value)] // TODO(relay)
303    fn handle_padding_negotiate(&mut self, _cell: PaddingNegotiate) -> StdResult<(), ReactorError> {
304        Err(internal!("PADDING_NEGOTIATE is not implemented").into())
305    }
306}
307
308impl ForwardHandler for Forward {
309    type BuildSpec = OwnedChanTarget;
310    type CircChanMsg = RelayCircChanMsg;
311    type CircEvent = CircEvent;
312
313    async fn handle_meta_msg<R: Runtime>(
314        &mut self,
315        runtime: &R,
316        early: bool,
317        _hopnum: Option<HopNum>,
318        msg: UnparsedRelayMsg,
319        _relay_cell_format: RelayCellFormat,
320    ) -> StdResult<(), ReactorError> {
321        match msg.cmd() {
322            RelayCmd::DROP => self.handle_drop(),
323            RelayCmd::EXTEND2 => self.extend_handler.handle_extend2(runtime, early, msg),
324            RelayCmd::TRUNCATE => self.handle_truncate(),
325            cmd => Err(internal!("relay cmd {cmd} not supported").into()),
326        }
327    }
328
329    async fn handle_forward_cell<R: Runtime>(
330        &mut self,
331        hop_mgr: &mut HopMgr<R>,
332        cell: RelayCircChanMsg,
333    ) -> StdResult<Option<ForwardCellDisposition>, ReactorError> {
334        use RelayCircChanMsg::*;
335
336        match cell {
337            Relay(r) => self.handle_relay_cell(hop_mgr, r.into()),
338            RelayEarly(r) => self.handle_relay_cell(hop_mgr, r.into()),
339            Destroy(d) => {
340                self.handle_destroy_cell(&d)?;
341                Ok(None)
342            }
343            PaddingNegotiate(p) => {
344                self.handle_padding_negotiate(p)?;
345                Ok(None)
346            }
347        }
348    }
349
350    fn handle_event(
351        &mut self,
352        event: Self::CircEvent,
353    ) -> StdResult<Option<BackwardReactorCmd>, ReactorError> {
354        match event {
355            CircEvent::ExtendResult(res) => self.handle_extend_result(res),
356        }
357    }
358
359    async fn outbound_chan_ready(&mut self) -> Result<()> {
360        future::poll_fn(|cx| match &mut self.outbound {
361            Some(chan) => {
362                let _ = chan.outbound_chan_tx.poll_flush_unpin(cx);
363
364                chan.outbound_chan_tx.poll_ready_unpin(cx)
365            }
366            None => {
367                // Pedantically, if the channel doesn't exist, it can't be ready,
368                // but we have no choice here than to return Ready
369                // (returning Pending would cause the reactor to lock up).
370                //
371                // Returning ready here means the base reactor is allowed to read
372                // from its inbound channel. This is OK, because if we *do*
373                // read a cell from that channel and find ourselves needing to
374                // forward it to the next hop, we simply return a proto violation error,
375                // shutting down the reactor.
376                Poll::Ready(Ok(()))
377            }
378        })
379        .await
380    }
381}
382
383impl ControlHandler for Forward {
384    type CtrlMsg = CtrlMsg;
385    type CtrlCmd = CtrlCmd;
386
387    fn handle_cmd(&mut self, cmd: Self::CtrlCmd) -> StdResult<(), ReactorError> {
388        let () = cmd;
389        Ok(())
390    }
391
392    fn handle_msg(&mut self, msg: Self::CtrlMsg) -> StdResult<(), ReactorError> {
393        let () = msg;
394        Ok(())
395    }
396}
397
398impl Drop for Forward {
399    fn drop(&mut self) {
400        if let Some(outbound) = self.outbound.as_mut() {
401            // This will send a DESTROY down the outbound channel
402            let _ = outbound.channel.close_circuit(outbound.circ_id);
403        }
404    }
405}