1mod 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
21use 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
42type CtrlMsg = ();
44
45type CtrlCmd = ();
47
48const MAX_RELAY_EARLY_CELLS_PER_CIRCUIT: usize = 8;
52
53pub(crate) struct Forward {
55 unique_id: UniqId,
57 circ_id: CircId,
59 outbound: Option<Outbound>,
65 crypto_out: Box<dyn OutboundRelayLayer + Send>,
67 relay_early_count: usize,
71 extend_handler: ExtendRequestHandler,
75}
76
77pub(crate) enum CircEvent {
79 ExtendResult(StdResult<ExtendResult, ReactorError>),
81}
82
83pub(crate) struct ExtendResult {
85 extended2: Extended2,
87 outbound: Outbound,
89 outbound_chan_rx: CircuitRxReceiver,
93}
94
95struct Outbound {
97 circ_id: CircId,
99 channel: Arc<Channel>,
101 outbound_chan_tx: ChannelSender,
103}
104
105enum CellDecodeResult {
107 Recognized(SendmeTag, RelayCellDecoderResult),
109 Unrecognizd(RelayCellBody),
111}
112
113impl Forward {
114 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 outbound: None,
139 crypto_out,
140 relay_early_count: 0,
141 extend_handler,
142 }
143 }
144
145 fn decode_relay_cell<R: Runtime>(
147 &mut self,
148 hop_mgr: &mut HopMgr<R>,
149 cell: RelayMaybeEarlyChanMsg,
150 ) -> Result<(Option<HopNum>, CellDecodeResult)> {
151 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 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 #[allow(clippy::unnecessary_wraps)] 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 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 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 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 return Err(Error::CircProto(
248 "Asked to forward cell before the circuit was extended?!".into(),
249 )
250 .into());
251 };
252
253 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 chan.outbound_chan_tx.start_send_unpin((cell, info))?;
273
274 Ok(())
275 }
276
277 fn handle_truncate(&mut self) -> StdResult<(), ReactorError> {
279 Err(Error::CircProto("TRUNCATE not allowed".into()).into())
285 }
286
287 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 Err(ReactorError::Shutdown)
299 }
300
301 #[allow(clippy::needless_pass_by_value)] 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 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 let _ = outbound.channel.close_circuit(outbound.circ_id);
403 }
404 }
405}