Skip to main content

tor_proto/client/reactor/
circuit.rs

1//! Module exposing types for representing circuits in the tunnel reactor.
2
3pub(crate) mod circhop;
4pub(super) mod extender;
5
6use crate::channel::Channel;
7use crate::circuit::cell_sender::CircuitCellSender;
8use crate::circuit::celltypes::CreateResponse;
9use crate::circuit::circhop::{HopSettings, ReactorStreamComponents};
10use crate::circuit::create::{Create2Wrap, CreateFastWrap, CreateHandshakeWrap};
11use crate::circuit::padding::CircPaddingDisposition;
12use crate::circuit::{CircuitRxReceiver, UniqId};
13use crate::client::circuit::handshake::{BoxedClientLayer, HandshakeRole};
14use crate::client::circuit::padding::{
15    self, PaddingController, PaddingEventStream, QueuedCellPaddingInfo,
16};
17use crate::client::circuit::{ClientCircChanMsg, MutableState, path};
18use crate::client::reactor::MetaCellDisposition;
19use crate::congestion::CongestionSignals;
20use crate::congestion::sendme;
21use crate::crypto::binding::CircuitBinding;
22use crate::crypto::cell::{
23    HopNum, InboundClientCrypt, InboundClientLayer, OutboundClientCrypt, OutboundClientLayer,
24    RelayCellBody,
25};
26use crate::crypto::handshake::fast::CreateFastClient;
27use crate::crypto::handshake::ntor::{NtorClient, NtorPublicKey};
28use crate::crypto::handshake::ntor_v3::{NtorV3Client, NtorV3PublicKey};
29use crate::crypto::handshake::{ClientHandshake, KeyGenerator};
30use crate::memquota::{CircuitAccount, SpecificAccount as _, StreamAccount};
31use crate::stream::cmdcheck::{AnyCmdChecker, StreamStatus};
32use crate::stream::flow_ctrl::state::WithSidechannelMitigations;
33use crate::stream::msg_streamid;
34use crate::streammap;
35use crate::tunnel::TunnelScopedCircId;
36use crate::util::err::ReactorError;
37use crate::util::timeout::TimeoutEstimator;
38use crate::{ClockSkew, Error, Result};
39
40use tor_async_utils::{SinkTrySend as _, SinkTrySendError as _};
41use tor_basic_utils::onionperf_types::{OnionperfCircuitStatus, OnionperfEvent};
42use tor_cell::chancell::msg::{AnyChanMsg, HandshakeType, Relay};
43use tor_cell::chancell::{AnyChanCell, ChanCmd, CircId};
44use tor_cell::chancell::{BoxedCellBody, ChanMsg};
45use tor_cell::relaycell::msg::{AnyRelayMsg, End, Sendme, SendmeTag, Truncated};
46use tor_cell::relaycell::{
47    AnyRelayMsgOuter, RelayCellDecoderResult, RelayCellFormat, RelayCmd, StreamId, UnparsedRelayMsg,
48};
49use tor_error::{Bug, internal};
50use tor_linkspec::RelayIds;
51use tor_llcrypto::pk;
52use web_time_compat::{Duration, Instant, SystemTime};
53
54use futures::SinkExt as _;
55use oneshot_fused_workaround as oneshot;
56use tor_rtcompat::{DynTimeProvider, SleepProvider as _};
57use tracing::{debug, instrument, trace, warn};
58
59use super::{
60    CellHandlers, CircuitHandshake, CloseStreamBehavior, ReactorResultChannel, SendRelayCell,
61};
62
63use crate::conflux::msghandler::ConfluxStatus;
64
65use std::borrow::Borrow;
66use std::pin::Pin;
67use std::result::Result as StdResult;
68use std::sync::Arc;
69
70use extender::HandshakeAuxDataHandler;
71
72#[cfg(feature = "hs-service")]
73use {
74    crate::circuit::CircHopSyncView,
75    crate::stream::{InboundDataCmdChecker, IncomingStreamRequest},
76    tor_cell::relaycell::msg::Begin,
77};
78
79#[cfg(feature = "conflux")]
80use {
81    crate::conflux::msghandler::{ConfluxAction, ConfluxCmd, ConfluxMsgHandler, OooRelayMsg},
82    crate::tunnel::TunnelId,
83};
84
85pub(super) use circhop::{CircHop, CircHopList};
86
87/// A circuit "leg" from a tunnel.
88///
89/// Regular (non-multipath) circuits have a single leg.
90/// Conflux (multipath) circuits have `N` (usually, `N = 2`).
91pub(crate) struct Circuit {
92    /// The time provider.
93    runtime: DynTimeProvider,
94    /// The channel this circuit is attached to.
95    channel: Arc<Channel>,
96    /// Sender object used to actually send cells.
97    ///
98    /// NOTE: Control messages could potentially add unboundedly to this, although that's
99    ///       not likely to happen (and isn't triggereable from the network, either).
100    pub(super) chan_sender: CircuitCellSender,
101    /// Input stream, on which we receive ChanMsg objects from this circuit's
102    /// channel.
103    ///
104    // TODO: could use a SPSC channel here instead.
105    pub(super) input: CircuitRxReceiver,
106    /// The cryptographic state for this circuit for inbound cells.
107    /// This object is divided into multiple layers, each of which is
108    /// shared with one hop of the circuit.
109    crypto_in: InboundClientCrypt,
110    /// The cryptographic state for this circuit for outbound cells.
111    crypto_out: OutboundClientCrypt,
112    /// List of hops state objects used by the reactor
113    pub(super) hops: CircHopList,
114    /// Mutable information about this circuit,
115    /// shared with the reactor's `ConfluxSet`.
116    mutable: Arc<MutableState>,
117    /// This circuit's identifier.
118    circ_id: CircId,
119    /// An identifier for logging about this reactor's circuit.
120    unique_id: TunnelScopedCircId,
121    /// A handler for conflux cells.
122    ///
123    /// Set once the conflux handshake is initiated by the reactor
124    /// using [`Reactor::handle_link_circuits`](super::Reactor::handle_link_circuits).
125    #[cfg(feature = "conflux")]
126    conflux_handler: Option<ConfluxMsgHandler>,
127    /// A padding controller to which padding-related events should be reported.
128    padding_ctrl: PaddingController,
129    /// An event stream telling us about padding-related events.
130    //
131    // TODO: it would be nice to have all of these streams wrapped in a single
132    // SelectAll, but we can't really do that, since we need the ability to move them
133    // from one conflux set to another, and a SelectAll doesn't let you actually
134    // remove one of its constituent streams.  This issue might get solved along
135    // with the rest of the next reactor refactoring.
136    pub(super) padding_event_stream: PaddingEventStream,
137    /// Current rules for blocking traffic, according to the padding controller.
138    #[cfg(feature = "circ-padding")]
139    padding_block: Option<padding::StartBlocking>,
140    /// The circuit timeout estimator.
141    ///
142    /// Used for computing half-stream expiration.
143    timeouts: Arc<dyn TimeoutEstimator>,
144    /// Memory quota account
145    #[allow(dead_code)] // Partly here to keep it alive as long as the circuit
146    memquota: CircuitAccount,
147}
148
149/// A command to run in response to a circuit event.
150///
151/// Unlike `RunOnceCmdInner`, doesn't know anything about `UniqId`s.
152/// The user of the `CircuitCmd`s is supposed to know the `UniqId`
153/// of the circuit the `CircuitCmd` came from.
154///
155/// This type gets mapped to a `RunOnceCmdInner` in the circuit reactor.
156#[derive(Debug, derive_more::From)]
157pub(super) enum CircuitCmd {
158    /// Send a RELAY cell on the circuit leg this command originates from.
159    Send(SendRelayCell),
160    /// Handle a SENDME message received on the circuit leg this command originates from.
161    HandleSendMe {
162        /// The hop number.
163        hop: HopNum,
164        /// The SENDME message to handle.
165        sendme: Sendme,
166    },
167    /// Close the specified stream on the circuit leg this command originates from.
168    CloseStream {
169        /// The hop number.
170        hop: HopNum,
171        /// The ID of the stream to close.
172        sid: StreamId,
173        /// The stream-closing behavior.
174        behav: CloseStreamBehavior,
175        /// The reason for closing the stream.
176        reason: streammap::TerminateReason,
177    },
178    /// Perform an action resulting from handling a conflux cell.
179    #[cfg(feature = "conflux")]
180    Conflux(ConfluxCmd),
181    /// Perform a clean shutdown on this circuit.
182    CleanShutdown,
183    /// Enqueue an out-of-order cell in the reactor.
184    #[cfg(feature = "conflux")]
185    Enqueue(OooRelayMsg),
186}
187
188/// Return a `CircProto` error for the specified unsupported cell.
189///
190/// This error will shut down the reactor.
191///
192/// Note: this is a macro to simplify usage (this way the caller doesn't
193/// need to .map() the result to the appropriate type)
194macro_rules! unsupported_client_cell {
195    ($msg:expr) => {{
196        unsupported_client_cell!(@ $msg, "")
197    }};
198
199    ($msg:expr, $hopnum:expr) => {{
200        let hop: HopNum = $hopnum;
201        let hop_display = format!(" from hop {}", hop.display());
202        unsupported_client_cell!(@ $msg, hop_display)
203    }};
204
205    (@ $msg:expr, $hopnum_display:expr) => {
206        Err(crate::Error::CircProto(format!(
207            "Unexpected {} cell{} on client circuit",
208            $msg.cmd(),
209            $hopnum_display,
210        )))
211    };
212}
213
214pub(super) use unsupported_client_cell;
215
216impl Circuit {
217    /// Create a new non-multipath circuit.
218    #[allow(clippy::too_many_arguments)]
219    pub(super) fn new(
220        runtime: DynTimeProvider,
221        channel: Arc<Channel>,
222        circ_id: CircId,
223        unique_id: TunnelScopedCircId,
224        input: CircuitRxReceiver,
225        memquota: CircuitAccount,
226        mutable: Arc<MutableState>,
227        padding_ctrl: PaddingController,
228        padding_event_stream: PaddingEventStream,
229        timeouts: Arc<dyn TimeoutEstimator>,
230    ) -> Self {
231        let chan_sender = CircuitCellSender::from_channel_sender(channel.sender());
232
233        let crypto_out = OutboundClientCrypt::new();
234        Circuit {
235            runtime,
236            channel,
237            chan_sender,
238            input,
239            crypto_in: InboundClientCrypt::new(),
240            hops: CircHopList::default(),
241            unique_id,
242            circ_id,
243            crypto_out,
244            mutable,
245            #[cfg(feature = "conflux")]
246            conflux_handler: None,
247            padding_ctrl,
248            padding_event_stream,
249            #[cfg(feature = "circ-padding")]
250            padding_block: None,
251            timeouts,
252            memquota,
253        }
254    }
255
256    /// Return the process-unique identifier of this circuit.
257    pub(super) fn unique_id(&self) -> UniqId {
258        self.unique_id.unique_id()
259    }
260
261    /// Return this circuit's identifier.
262    pub(super) fn circ_id(&self) -> CircId {
263        self.circ_id
264    }
265
266    /// Return the shared mutable state of this circuit.
267    pub(super) fn mutable(&self) -> &Arc<MutableState> {
268        &self.mutable
269    }
270
271    /// Add this circuit to a multipath tunnel, by associating it with a new [`TunnelId`],
272    /// and installing a [`ConfluxMsgHandler`] on this circuit.
273    ///
274    /// Once this is called, the circuit will be able to handle conflux cells.
275    #[cfg(feature = "conflux")]
276    pub(super) fn add_to_conflux_tunnel(
277        &mut self,
278        tunnel_id: TunnelId,
279        conflux_handler: ConfluxMsgHandler,
280    ) {
281        self.unique_id = TunnelScopedCircId::new(tunnel_id, self.unique_id.unique_id());
282        self.conflux_handler = Some(conflux_handler);
283    }
284
285    /// Send a LINK cell to the specified hop.
286    ///
287    /// This must be called *after* a [`ConfluxMsgHandler`] is installed
288    /// on the circuit with [`add_to_conflux_tunnel`](Self::add_to_conflux_tunnel).
289    #[cfg(feature = "conflux")]
290    pub(super) async fn begin_conflux_link(
291        &mut self,
292        hop: HopNum,
293        cell: AnyRelayMsgOuter,
294        runtime: &tor_rtcompat::DynTimeProvider,
295    ) -> Result<()> {
296        use tor_rtcompat::SleepProvider as _;
297
298        if self.conflux_handler.is_none() {
299            return Err(internal!(
300                "tried to send LINK cell before installing a ConfluxMsgHandler?!"
301            )
302            .into());
303        }
304
305        let cell = SendRelayCell {
306            hop: Some(hop),
307            early: false,
308            cell,
309        };
310        self.send_relay_cell(cell).await?;
311
312        let Some(conflux_handler) = self.conflux_handler.as_mut() else {
313            return Err(internal!("ConfluxMsgHandler disappeared?!").into());
314        };
315
316        Ok(conflux_handler.note_link_sent(runtime.wallclock())?)
317    }
318
319    /// Get the wallclock time when the handshake on this circuit is supposed to time out.
320    ///
321    /// Returns `None` if the handshake is not currently in progress.
322    pub(super) fn conflux_hs_timeout(&self) -> Option<SystemTime> {
323        cfg_if::cfg_if! {
324            if #[cfg(feature = "conflux")] {
325                self.conflux_handler.as_ref().map(|handler| handler.handshake_timeout())?
326            } else {
327                None
328            }
329        }
330    }
331
332    /// Handle a [`CtrlMsg::AddFakeHop`](super::CtrlMsg::AddFakeHop) message.
333    #[cfg(test)]
334    pub(super) fn handle_add_fake_hop(
335        &mut self,
336        format: RelayCellFormat,
337        fwd_lasthop: bool,
338        rev_lasthop: bool,
339        dummy_peer_id: path::HopDetail,
340        // TODO-CGO: Take HopSettings instead of CircParams.
341        // (Do this after we've got the virtual-hop refactorings done for
342        // virtual extending.)
343        params: &crate::client::circuit::CircParameters,
344        done: ReactorResultChannel<()>,
345    ) {
346        use tor_protover::{Protocols, named};
347
348        use crate::client::circuit::test::DummyCrypto;
349
350        assert!(matches!(format, RelayCellFormat::V0));
351        let _ = format; // TODO-CGO: remove this once we have CGO+hs implemented.
352
353        let fwd = Box::new(DummyCrypto::new(fwd_lasthop));
354        let rev = Box::new(DummyCrypto::new(rev_lasthop));
355        let binding = None;
356
357        let settings = HopSettings::from_params_and_caps(
358            // This is for testing only, so we'll assume full negotiation took place.
359            crate::circuit::circhop::HopNegotiationType::Full,
360            params,
361            &[named::FLOWCTRL_CC].into_iter().collect::<Protocols>(),
362        )
363        .expect("Can't construct HopSettings");
364        self.add_hop(dummy_peer_id, fwd, rev, binding, &settings)
365            .expect("could not add hop to circuit");
366        let _ = done.send(Ok(()));
367    }
368
369    /// Encode `msg` and encrypt it, returning the resulting cell
370    /// and tag that should be expected for an authenticated SENDME sent
371    /// in response to that cell.
372    fn encode_relay_cell(
373        crypto_out: &mut OutboundClientCrypt,
374        relay_format: RelayCellFormat,
375        hop: HopNum,
376        early: bool,
377        msg: AnyRelayMsgOuter,
378    ) -> Result<(AnyChanMsg, SendmeTag)> {
379        let mut body: RelayCellBody = msg
380            .encode(relay_format, &mut rand::rng())
381            .map_err(|e| Error::from_cell_enc(e, "relay cell body"))?
382            .into();
383        let cmd = if early {
384            ChanCmd::RELAY_EARLY
385        } else {
386            ChanCmd::RELAY
387        };
388        let tag = crypto_out.encrypt(cmd, &mut body, hop)?;
389        let msg = Relay::from(BoxedCellBody::from(body));
390        let msg = if early {
391            AnyChanMsg::RelayEarly(msg.into())
392        } else {
393            AnyChanMsg::Relay(msg)
394        };
395
396        Ok((msg, tag))
397    }
398
399    /// Encode `msg`, encrypt it, and send it to the 'hop'th hop.
400    ///
401    /// If there is insufficient outgoing *circuit-level* or *stream-level*
402    /// SENDME window, an error is returned instead.
403    ///
404    /// Does not check whether the cell is well-formed or reasonable.
405    ///
406    /// NOTE: the reactor should not call this function directly, only via
407    /// [`ConfluxSet::send_relay_cell_on_leg`](super::conflux::ConfluxSet::send_relay_cell_on_leg),
408    /// which will reroute the message, if necessary to the primary leg.
409    #[instrument(level = "trace", skip_all)]
410    pub(super) async fn send_relay_cell(&mut self, msg: SendRelayCell) -> Result<()> {
411        self.send_relay_cell_inner(msg, None).await
412    }
413
414    /// As [`send_relay_cell`](Self::send_relay_cell), but takes an optional
415    /// [`QueuedCellPaddingInfo`] in `padding_info`.
416    ///
417    /// If `padding_info` is None, `msg` must be non-padding: we report it as such to the
418    /// padding controller.
419    #[instrument(level = "trace", skip_all)]
420    async fn send_relay_cell_inner(
421        &mut self,
422        msg: SendRelayCell,
423        padding_info: Option<QueuedCellPaddingInfo>,
424    ) -> Result<()> {
425        let SendRelayCell {
426            hop,
427            early,
428            cell: msg,
429        } = msg;
430
431        let is_conflux_link = msg.cmd() == RelayCmd::CONFLUX_LINK;
432        if !is_conflux_link && self.is_conflux_pending() {
433            // Note: it is the responsibility of the reactor user to wait until
434            // at least one of the legs completes the handshake.
435            return Err(internal!("tried to send cell on unlinked circuit").into());
436        }
437
438        trace!(
439            circ_uniq_id = %self.unique_id,
440            forward_circ_id = %self.circ_id,
441            cell = ?msg,
442            "sending relay cell"
443        );
444
445        // Cloned, because we borrow mutably from self when we get the circhop.
446        let runtime = self.runtime.clone();
447        let c_t_w = sendme::cmd_counts_towards_windows(msg.cmd());
448        let stream_id = msg.stream_id();
449        let hop = hop.expect("missing hop in client SendRelayCell?!");
450        let circhop = self.hops.get_mut(hop).ok_or(Error::NoSuchHop)?;
451
452        // We might be out of capacity entirely; see if we are about to hit a limit.
453        //
454        // TODO: If we ever add a notion of _recoverable_ errors below, we'll
455        // need a way to restore this limit, and similarly for about_to_send().
456        circhop.decrement_outbound_cell_limit()?;
457
458        // We need to apply stream-level flow control *before* encoding the message.
459        if c_t_w {
460            if let Some(stream_id) = stream_id {
461                circhop.about_to_send(stream_id, msg.msg())?;
462            }
463        }
464
465        // Save the RelayCmd of the message before it gets consumed below.
466        // We need this to tell our ConfluxMsgHandler about the cell we've just sent,
467        // so that it can update its counters.
468        let relay_cmd = msg.cmd();
469
470        // NOTE(eta): Now that we've encrypted the cell, we *must* either send it or abort
471        //            the whole circuit (e.g. by returning an error).
472        let (msg, tag) = Self::encode_relay_cell(
473            &mut self.crypto_out,
474            circhop.relay_cell_format(),
475            hop,
476            early,
477            msg,
478        )?;
479        // The cell counted for congestion control, inform our algorithm of such and pass down the
480        // tag for authenticated SENDMEs.
481        if c_t_w {
482            circhop.ccontrol().note_data_sent(&runtime, &tag)?;
483        }
484
485        // Remember that we've enqueued this cell.
486        let padding_info = padding_info.or_else(|| self.padding_ctrl.queued_data(hop));
487
488        self.send_msg(msg, padding_info).await?;
489
490        #[cfg(feature = "conflux")]
491        if let Some(conflux) = self.conflux_handler.as_mut() {
492            conflux.note_cell_sent(relay_cmd);
493        }
494
495        Ok(())
496    }
497
498    /// Helper: process a cell on a channel.  Most cells get ignored
499    /// or rejected; a few get delivered to circuits.
500    ///
501    /// Return `CellStatus::CleanShutdown` if we should exit.
502    ///
503    // TODO: returning `Vec<CircuitCmd>` means we're unnecessarily
504    // allocating a `Vec` here. Generally, the number of commands is going to be small
505    // (usually 1, but > 1 when we start supporting packed cells).
506    //
507    // We should consider using smallvec instead. It might also be a good idea to have a
508    // separate higher-level type splitting this out into Single(CircuitCmd),
509    // and Multiple(SmallVec<[CircuitCmd; <capacity>]>).
510    pub(super) fn handle_cell(
511        &mut self,
512        handlers: &mut CellHandlers,
513        leg: UniqId,
514        cell: ClientCircChanMsg,
515    ) -> Result<Vec<CircuitCmd>> {
516        trace!(
517            circ_uniq_id = %self.unique_id,
518            forward_circ_id = %self.circ_id,
519            cell = ?cell,
520            "handling cell"
521        );
522        use ClientCircChanMsg::*;
523        match cell {
524            Relay(r) => self.handle_relay_cell(handlers, leg, r),
525            Destroy(d) => {
526                let reason = d.reason();
527                debug!(
528                    circ_uniq_id = %self.unique_id,
529                    forward_circ_id = %self.circ_id,
530                    "Received DESTROY cell. Reason: {} [{}]",
531                    reason.human_str(),
532                    reason
533                );
534
535                self.handle_destroy_cell().map(|c| vec![c])
536            }
537        }
538    }
539
540    /// Decode `cell`, returning its corresponding hop number, tag,
541    /// and decoded body.
542    fn decode_relay_cell(
543        &mut self,
544        cell: Relay,
545    ) -> Result<(HopNum, SendmeTag, RelayCellDecoderResult)> {
546        // This is always RELAY, not RELAY_EARLY, so long as this code is client-only.
547        let cmd = cell.cmd();
548        let mut body = cell.into_relay_body().into();
549
550        // Decrypt the cell. If it's recognized, then find the
551        // corresponding hop.
552        let (hopnum, tag) = self.crypto_in.decrypt(cmd, &mut body)?;
553
554        // Decode the cell.
555        let decode_res = self
556            .hop_mut(hopnum)
557            .ok_or_else(|| {
558                Error::from(internal!(
559                    "Trying to decode cell from nonexistent hop {:?}",
560                    hopnum
561                ))
562            })?
563            .decode(body.into())?;
564
565        Ok((hopnum, tag, decode_res))
566    }
567
568    /// React to a Relay or RelayEarly cell.
569    fn handle_relay_cell(
570        &mut self,
571        handlers: &mut CellHandlers,
572        leg: UniqId,
573        cell: Relay,
574    ) -> Result<Vec<CircuitCmd>> {
575        let (hopnum, tag, decode_res) = self.decode_relay_cell(cell)?;
576
577        if decode_res.is_padding() {
578            self.padding_ctrl.decrypted_padding(hopnum)?;
579        } else {
580            self.padding_ctrl.decrypted_data(hopnum);
581        }
582
583        // Check whether we are allowed to receive more data for this circuit hop.
584        self.hop_mut(hopnum)
585            .ok_or_else(|| internal!("nonexistent hop {:?}", hopnum))?
586            .decrement_inbound_cell_limit()?;
587
588        let c_t_w = decode_res.cmds().any(sendme::cmd_counts_towards_windows);
589
590        // Decrement the circuit sendme windows, and see if we need to
591        // send a sendme cell.
592        let send_circ_sendme = if c_t_w {
593            self.hop_mut(hopnum)
594                .ok_or_else(|| Error::CircProto("Sendme from nonexistent hop".into()))?
595                .ccontrol()
596                .note_data_received()?
597        } else {
598            false
599        };
600
601        let mut circ_cmds = vec![];
602        // If we do need to send a circuit-level SENDME cell, do so.
603        if send_circ_sendme {
604            // This always sends a V1 (tagged) sendme cell, and thereby assumes
605            // that SendmeEmitMinVersion is no more than 1.  If the authorities
606            // every increase that parameter to a higher number, this will
607            // become incorrect.  (Higher numbers are not currently defined.)
608            let sendme = Sendme::from(tag);
609            let cell = AnyRelayMsgOuter::new(None, sendme.into());
610            circ_cmds.push(CircuitCmd::Send(SendRelayCell {
611                hop: Some(hopnum),
612                early: false,
613                cell,
614            }));
615
616            // Inform congestion control of the SENDME we are sending. This is a circuit level one.
617            self.hop_mut(hopnum)
618                .ok_or_else(|| {
619                    Error::from(internal!(
620                        "Trying to send SENDME to nonexistent hop {:?}",
621                        hopnum
622                    ))
623                })?
624                .ccontrol()
625                .note_sendme_sent()?;
626        }
627
628        let (mut msgs, incomplete) = decode_res.into_parts();
629        while let Some(msg) = msgs.next() {
630            let msg_status = self.handle_relay_msg(handlers, hopnum, leg, c_t_w, msg)?;
631
632            match msg_status {
633                None => continue,
634                Some(msg @ CircuitCmd::CleanShutdown) => {
635                    for m in msgs {
636                        debug!(
637                            "{id}: Ignoring relay msg received after triggering shutdown: {m:?}",
638                            id = self.unique_id
639                        );
640                    }
641                    if let Some(incomplete) = incomplete {
642                        debug!(
643                            "{id}: Ignoring partial relay msg received after triggering shutdown: {:?}",
644                            incomplete,
645                            id = self.unique_id,
646                        );
647                    }
648                    circ_cmds.push(msg);
649                    return Ok(circ_cmds);
650                }
651                Some(msg) => {
652                    circ_cmds.push(msg);
653                }
654            }
655        }
656
657        Ok(circ_cmds)
658    }
659
660    /// Handle a single incoming relay message.
661    fn handle_relay_msg(
662        &mut self,
663        handlers: &mut CellHandlers,
664        hopnum: HopNum,
665        leg: UniqId,
666        cell_counts_toward_windows: bool,
667        msg: UnparsedRelayMsg,
668    ) -> Result<Option<CircuitCmd>> {
669        // If this msg wants/refuses to have a Stream ID, does it
670        // have/not have one?
671        let streamid = msg_streamid(&msg)?;
672
673        // If this doesn't have a StreamId, it's a meta cell,
674        // not meant for a particular stream.
675        let Some(streamid) = streamid else {
676            return self.handle_meta_cell(handlers, hopnum, msg);
677        };
678
679        #[cfg(feature = "conflux")]
680        let msg = if let Some(conflux) = self.conflux_handler.as_mut() {
681            match conflux.action_for_msg(hopnum, cell_counts_toward_windows, streamid, msg)? {
682                ConfluxAction::Deliver(msg) => {
683                    // The message either doesn't count towards the sequence numbers
684                    // or is already well-ordered, so we're ready to handle it.
685
686                    // It's possible that some of our buffered messages are now ready to be
687                    // handled. We don't check that here, however, because that's handled
688                    // by the reactor main loop.
689                    msg
690                }
691                ConfluxAction::Enqueue(msg) => {
692                    // Tell the reactor to enqueue this msg
693                    return Ok(Some(CircuitCmd::Enqueue(msg)));
694                }
695            }
696        } else {
697            // If we don't have a conflux_handler, it means this circuit is not part of
698            // a conflux tunnel, so we can just process the message.
699            msg
700        };
701
702        self.handle_in_order_relay_msg(
703            handlers,
704            hopnum,
705            leg,
706            cell_counts_toward_windows,
707            streamid,
708            msg,
709        )
710    }
711
712    /// Handle a single incoming relay message that is known to be in order.
713    pub(super) fn handle_in_order_relay_msg(
714        &mut self,
715        handlers: &mut CellHandlers,
716        hopnum: HopNum,
717        leg: UniqId,
718        cell_counts_toward_windows: bool,
719        streamid: StreamId,
720        msg: UnparsedRelayMsg,
721    ) -> Result<Option<CircuitCmd>> {
722        let now = self.runtime.now();
723
724        #[cfg(feature = "conflux")]
725        if let Some(conflux) = self.conflux_handler.as_mut() {
726            conflux.inc_last_seq_delivered(&msg);
727        }
728
729        let path = self.mutable.path();
730
731        let nonexistent_hop_err = || Error::CircProto("Cell from nonexistent hop!".into());
732        let hop = self.hop_mut(hopnum).ok_or_else(nonexistent_hop_err)?;
733
734        let hop_detail = path
735            .iter()
736            .nth(usize::from(hopnum))
737            .ok_or_else(nonexistent_hop_err)?;
738
739        // Returns the original message if it's an incoming stream request
740        // that we need to handle.
741        let res = hop.handle_msg(hop_detail, cell_counts_toward_windows, streamid, msg, now)?;
742
743        // If it was an incoming stream request, we don't need to worry about
744        // sending an XOFF as there's no stream data within this message.
745        if let Some(msg) = res {
746            cfg_if::cfg_if! {
747                if #[cfg(feature = "hs-service")] {
748                    return self.handle_incoming_stream_request(
749                        handlers,
750                        msg,
751                        streamid,
752                        hopnum,
753                        leg,
754                        // This is an onion service stream,
755                        // so we want sidechannel mitigations for flow control.
756                        WithSidechannelMitigations::Enabled,
757                    );
758                } else {
759                    return Err(
760                        Error::CircProto(format!("Cannot handle {} cells on this circuit", msg.cmd())),
761                    );
762                }
763            }
764        }
765
766        // We may want to send an XOFF if the incoming buffer is too large.
767        if let Some(cell) = hop.maybe_send_xoff(streamid)? {
768            let cell = AnyRelayMsgOuter::new(Some(streamid), cell.into());
769            let cell = SendRelayCell {
770                hop: Some(hopnum),
771                early: false,
772                cell,
773            };
774            return Ok(Some(CircuitCmd::Send(cell)));
775        }
776
777        Ok(None)
778    }
779
780    /// Handle a conflux message coming from the specified hop.
781    ///
782    /// Returns an error if
783    ///
784    ///   * this is not a conflux circuit (i.e. it doesn't have a [`ConfluxMsgHandler`])
785    ///   * this is a client circuit and the conflux message originated an unexpected hop
786    ///   * the cell was sent in violation of the handshake protocol
787    #[cfg(feature = "conflux")]
788    fn handle_conflux_msg(
789        &mut self,
790        hop: HopNum,
791        msg: UnparsedRelayMsg,
792    ) -> Result<Option<ConfluxCmd>> {
793        let Some(conflux_handler) = self.conflux_handler.as_mut() else {
794            // If conflux is not enabled, tear down the circuit
795            // (see 4.2.1. Cell Injection Side Channel Mitigations in prop329)
796            return Err(Error::CircProto(format!(
797                "Received {} cell from hop {} on non-conflux client circuit?!",
798                msg.cmd(),
799                hop.display(),
800            )));
801        };
802
803        Ok(conflux_handler.handle_conflux_msg(msg, hop))
804    }
805
806    /// For conflux: return the sequence number of the last cell sent on this leg.
807    ///
808    /// Returns an error if this circuit is not part of a conflux set.
809    #[cfg(feature = "conflux")]
810    pub(super) fn last_seq_sent(&self) -> Result<u64> {
811        let handler = self
812            .conflux_handler
813            .as_ref()
814            .ok_or_else(|| internal!("tried to get last_seq_sent of non-conflux circ"))?;
815
816        Ok(handler.last_seq_sent())
817    }
818
819    /// For conflux: set the sequence number of the last cell sent on this leg.
820    ///
821    /// Returns an error if this circuit is not part of a conflux set.
822    #[cfg(feature = "conflux")]
823    pub(super) fn set_last_seq_sent(&mut self, n: u64) -> Result<()> {
824        let handler = self
825            .conflux_handler
826            .as_mut()
827            .ok_or_else(|| internal!("tried to get last_seq_sent of non-conflux circ"))?;
828
829        handler.set_last_seq_sent(n);
830        Ok(())
831    }
832
833    /// For conflux: return the sequence number of the last cell received on this leg.
834    ///
835    /// Returns an error if this circuit is not part of a conflux set.
836    #[cfg(feature = "conflux")]
837    pub(super) fn last_seq_recv(&self) -> Result<u64> {
838        let handler = self
839            .conflux_handler
840            .as_ref()
841            .ok_or_else(|| internal!("tried to get last_seq_recv of non-conflux circ"))?;
842
843        Ok(handler.last_seq_recv())
844    }
845
846    /// A helper for handling incoming stream requests.
847    ///
848    // TODO: can we make this a method on CircHop to avoid the double HopNum lookup?
849    #[cfg(feature = "hs-service")]
850    fn handle_incoming_stream_request(
851        &mut self,
852        handlers: &mut CellHandlers,
853        msg: UnparsedRelayMsg,
854        stream_id: StreamId,
855        hop_num: HopNum,
856        leg: UniqId,
857        with_sidechannel_mitigations: WithSidechannelMitigations,
858    ) -> Result<Option<CircuitCmd>> {
859        use tor_cell::relaycell::msg::EndReason;
860        use tor_error::into_internal;
861        use tor_log_ratelim::log_ratelim;
862
863        use crate::stream::incoming::StreamReqInfo;
864
865        // We need to construct this early so that we don't double-borrow &mut self
866
867        let Some(handler) = handlers.incoming_stream_req_handler.as_mut() else {
868            return Err(Error::CircProto(
869                "Cannot handle BEGIN cells on this circuit".into(),
870            ));
871        };
872
873        // The handler's hop_num is only ever set to None for relays.
874        let expected_hop_num = handler
875            .hop_num
876            .ok_or_else(|| internal!("Handler HopNum is None in client impl?!"))?;
877
878        if hop_num != expected_hop_num {
879            return Err(Error::CircProto(format!(
880                "Expecting incoming streams from {}, but received {} cell from unexpected hop {}",
881                expected_hop_num.display(),
882                msg.cmd(),
883                hop_num.display()
884            )));
885        }
886
887        let message_closes_stream = handler.cmd_checker.check_msg(&msg)? == StreamStatus::Closed;
888
889        // TODO: we've already looked up the `hop` in handle_relay_cell, so we shouldn't
890        // have to look it up again! However, we can't pass the `&mut hop` reference from
891        // `handle_relay_cell` to this function, because that makes Rust angry (we'd be
892        // borrowing self as mutable more than once).
893        //
894        // TODO: we _could_ use self.hops.get_mut(..) instead self.hop_mut(..) inside
895        // handle_relay_cell to work around the problem described above
896        let hop = self.hops.get_mut(hop_num).ok_or(Error::CircuitClosed)?;
897
898        if message_closes_stream {
899            hop.ending_msg_received(stream_id)?;
900
901            return Ok(None);
902        }
903
904        let begin = msg
905            .decode::<Begin>()
906            .map_err(|e| Error::from_bytes_err(e, "Invalid Begin message"))?
907            .into_msg();
908
909        let req = IncomingStreamRequest::Begin(begin);
910
911        {
912            use crate::stream::IncomingStreamRequestDisposition::*;
913
914            let ctx = crate::stream::IncomingStreamRequestContext { request: &req };
915            // IMPORTANT: super::syncview::CircHopSyncView::n_open_streams() (called via disposition() below)
916            // accesses the stream map mutexes!
917            //
918            // This means it's very important not to call this function while any of the hop's
919            // stream map mutex is held.
920            let view = CircHopSyncView::new(hop.outbound());
921
922            match handler.filter.as_mut().disposition(&ctx, &view)? {
923                Accept => {}
924                CloseCircuit => return Ok(Some(CircuitCmd::CleanShutdown)),
925                RejectRequest(end) => {
926                    let end_msg = AnyRelayMsgOuter::new(Some(stream_id), end.into());
927                    let cell = SendRelayCell {
928                        hop: Some(hop_num),
929                        early: false,
930                        cell: end_msg,
931                    };
932                    return Ok(Some(CircuitCmd::Send(cell)));
933                }
934            }
935        }
936
937        // TODO: Sadly, we need to look up `&mut hop` yet again,
938        // since we needed to pass `&self.hops` by reference to our filter above. :(
939        let hop = self.hops.get_mut(hop_num).ok_or(Error::CircuitClosed)?;
940        let relay_cell_format = hop.relay_cell_format();
941
942        let memquota = StreamAccount::new(&self.memquota)?;
943
944        let cmd_checker = InboundDataCmdChecker::new_connected();
945        let stream_components = hop.add_ent_with_id(
946            self.chan_sender.time_provider(),
947            stream_id,
948            cmd_checker,
949            with_sidechannel_mitigations,
950            &memquota,
951        )?;
952
953        let outcome = Pin::new(&mut handler.incoming_sender).try_send(StreamReqInfo {
954            req,
955            stream_id,
956            hop: Some((leg, hop_num).into()),
957            stream_components,
958            memquota,
959            relay_cell_format,
960        });
961
962        log_ratelim!("Delivering message to incoming stream handler"; outcome);
963
964        if let Err(e) = outcome {
965            if e.is_full() {
966                // The IncomingStreamRequestHandler's stream is full; it isn't
967                // handling requests fast enough. So instead, we reply with an
968                // END cell.
969                let end_msg = AnyRelayMsgOuter::new(
970                    Some(stream_id),
971                    End::new_with_reason(EndReason::RESOURCELIMIT).into(),
972                );
973
974                let cell = SendRelayCell {
975                    hop: Some(hop_num),
976                    early: false,
977                    cell: end_msg,
978                };
979                return Ok(Some(CircuitCmd::Send(cell)));
980            } else if e.is_disconnected() {
981                // The IncomingStreamRequestHandler's stream has been dropped.
982                // In the Tor protocol as it stands, this always means that the
983                // circuit itself is out-of-use and should be closed. (See notes
984                // on `allow_stream_requests.`)
985                //
986                // Note that we will _not_ reach this point immediately after
987                // the IncomingStreamRequestHandler is dropped; we won't hit it
988                // until we next get an incoming request.  Thus, if we do later
989                // want to add early detection for a dropped
990                // IncomingStreamRequestHandler, we need to do it elsewhere, in
991                // a different way.
992                debug!(
993                    circ_uniq_id = %self.unique_id,
994                    forward_circ_id = %self.circ_id,
995                    "Incoming stream request receiver dropped",
996                );
997                // This will _cause_ the circuit to get closed.
998                return Err(Error::CircuitClosed);
999            } else {
1000                // There are no errors like this with the current design of
1001                // futures::mpsc, but we shouldn't just ignore the possibility
1002                // that they'll be added later.
1003                return Err(Error::from((into_internal!(
1004                    "try_send failed unexpectedly"
1005                ))(e)));
1006            }
1007        }
1008
1009        Ok(None)
1010    }
1011
1012    /// Helper: process a destroy cell.
1013    #[allow(clippy::unnecessary_wraps)]
1014    fn handle_destroy_cell(&mut self) -> Result<CircuitCmd> {
1015        // I think there is nothing more to do here.
1016        Ok(CircuitCmd::CleanShutdown)
1017    }
1018
1019    /// Handle a [`CtrlMsg::Create`](super::CtrlMsg::Create) message.
1020    pub(super) async fn handle_create(
1021        &mut self,
1022        recv_created: oneshot::Receiver<CreateResponse>,
1023        handshake: CircuitHandshake,
1024        settings: HopSettings,
1025        done: ReactorResultChannel<()>,
1026    ) -> StdResult<(), ReactorError> {
1027        let ret = match handshake {
1028            CircuitHandshake::CreateFast => self.create_firsthop_fast(recv_created, settings).await,
1029            CircuitHandshake::Ntor {
1030                public_key,
1031                ed_identity,
1032            } => {
1033                self.create_firsthop_ntor(recv_created, ed_identity, public_key, settings)
1034                    .await
1035            }
1036            CircuitHandshake::NtorV3 { public_key } => {
1037                self.create_firsthop_ntor_v3(recv_created, public_key, settings)
1038                    .await
1039            }
1040        };
1041        let _ = done.send(ret); // don't care if sender goes away
1042
1043        // TODO: maybe we don't need to flush here?
1044        // (we could let run_once() handle all the flushing)
1045        self.chan_sender.flush().await?;
1046
1047        Ok(())
1048    }
1049
1050    /// Helper: create the first hop of a circuit.
1051    ///
1052    /// This is parameterized not just on the RNG, but a wrapper object to
1053    /// build the right kind of create cell, and a handshake object to perform
1054    /// the cryptographic handshake.
1055    async fn create_impl<H, W, M>(
1056        &mut self,
1057        recvcreated: oneshot::Receiver<CreateResponse>,
1058        wrap: &W,
1059        key: &H::KeyType,
1060        mut settings: HopSettings,
1061        msg: &M,
1062    ) -> Result<()>
1063    where
1064        H: ClientHandshake + HandshakeAuxDataHandler,
1065        W: CreateHandshakeWrap,
1066        H::KeyGen: KeyGenerator,
1067        M: Borrow<H::ClientAuxData>,
1068    {
1069        // We don't need to shut down the circuit on failure here, since this
1070        // function consumes the PendingClientCirc and only returns
1071        // a ClientCirc on success.
1072
1073        let (state, msg) = H::client1(&mut rand::rng(), key, msg)?;
1074        let create_cell = wrap.to_chanmsg(msg);
1075        trace!(
1076            circ_uniq_id = %self.unique_id,
1077            forward_circ_id = %self.circ_id,
1078            create = %create_cell.cmd(),
1079            "Extending to hop 1",
1080        );
1081        self.send_msg(create_cell, None).await?;
1082
1083        let reply = recvcreated
1084            .await
1085            .map_err(|_| Error::CircProto("Circuit closed while waiting".into()))?;
1086
1087        let relay_handshake = wrap.decode_chanmsg(reply)?;
1088        let (server_msg, keygen) = H::client2(state, relay_handshake)?;
1089
1090        H::handle_server_aux_data(&mut settings, &server_msg)?;
1091
1092        let BoxedClientLayer { fwd, back, binding } = settings
1093            .relay_crypt_protocol()
1094            .construct_client_layers(HandshakeRole::Initiator, keygen)?;
1095
1096        trace!(
1097            circ_uniq_id = %self.unique_id,
1098            forward_circ_id = %self.circ_id,
1099            "Handshake complete; circuit created."
1100        );
1101
1102        trace!(
1103            onionperf = true,
1104            circ_uniq_id = %self.unique_id,
1105            forward_circ_id = %self.circ_id,
1106            event = ?OnionperfEvent::Circuit(OnionperfCircuitStatus::Extended),
1107        );
1108
1109        let peer_id = self.channel.target().clone();
1110
1111        self.add_hop(
1112            path::HopDetail::Relay(peer_id),
1113            fwd,
1114            back,
1115            binding,
1116            &settings,
1117        )?;
1118        Ok(())
1119    }
1120
1121    /// Use the (questionable!) CREATE_FAST handshake to connect to the
1122    /// first hop of this circuit.
1123    ///
1124    /// There's no authentication in CREATE_FAST,
1125    /// so we don't need to know whom we're connecting to: we're just
1126    /// connecting to whichever relay the channel is for.
1127    async fn create_firsthop_fast(
1128        &mut self,
1129        recvcreated: oneshot::Receiver<CreateResponse>,
1130        settings: HopSettings,
1131    ) -> Result<()> {
1132        // In a CREATE_FAST handshake, we can't negotiate a format other than this.
1133        let wrap = CreateFastWrap;
1134        self.create_impl::<CreateFastClient, _, _>(recvcreated, &wrap, &(), settings, &())
1135            .await
1136    }
1137
1138    /// Use the ntor handshake to connect to the first hop of this circuit.
1139    ///
1140    /// Note that the provided keys must match the channel's target,
1141    /// or the handshake will fail.
1142    async fn create_firsthop_ntor(
1143        &mut self,
1144        recvcreated: oneshot::Receiver<CreateResponse>,
1145        ed_identity: pk::ed25519::Ed25519Identity,
1146        pubkey: NtorPublicKey,
1147        settings: HopSettings,
1148    ) -> Result<()> {
1149        // Exit now if we have an Ed25519 or RSA identity mismatch.
1150        let target = RelayIds::builder()
1151            .ed_identity(ed_identity)
1152            .rsa_identity(pubkey.id)
1153            .build()
1154            .expect("Unable to build RelayIds");
1155        self.channel.check_match(&target)?;
1156
1157        let wrap = Create2Wrap {
1158            handshake_type: HandshakeType::NTOR,
1159        };
1160        self.create_impl::<NtorClient, _, _>(recvcreated, &wrap, &pubkey, settings, &())
1161            .await
1162    }
1163
1164    /// Use the ntor-v3 handshake to connect to the first hop of this circuit.
1165    ///
1166    /// Note that the provided key must match the channel's target,
1167    /// or the handshake will fail.
1168    async fn create_firsthop_ntor_v3(
1169        &mut self,
1170        recvcreated: oneshot::Receiver<CreateResponse>,
1171        pubkey: NtorV3PublicKey,
1172        settings: HopSettings,
1173    ) -> Result<()> {
1174        // Exit now if we have a mismatched key.
1175        let target = RelayIds::builder()
1176            .ed_identity(pubkey.id)
1177            .build()
1178            .expect("Unable to build RelayIds");
1179        self.channel.check_match(&target)?;
1180
1181        // Set the client extensions.
1182        let client_extensions = settings.circuit_request_extensions()?;
1183        let wrap = Create2Wrap {
1184            handshake_type: HandshakeType::NTOR_V3,
1185        };
1186
1187        self.create_impl::<NtorV3Client, _, _>(
1188            recvcreated,
1189            &wrap,
1190            &pubkey,
1191            settings,
1192            &client_extensions,
1193        )
1194        .await
1195    }
1196
1197    /// Add a hop to the end of this circuit.
1198    ///
1199    /// Will return an error if the circuit already has [`u8::MAX`] hops.
1200    pub(super) fn add_hop(
1201        &mut self,
1202        peer_id: path::HopDetail,
1203        fwd: Box<dyn OutboundClientLayer + 'static + Send>,
1204        rev: Box<dyn InboundClientLayer + 'static + Send>,
1205        binding: Option<CircuitBinding>,
1206        settings: &HopSettings,
1207    ) -> StdResult<(), Bug> {
1208        let hop_num = self.hops.len();
1209        debug_assert_eq!(hop_num, usize::from(self.num_hops()));
1210
1211        // There are several places in the code that assume that a `usize` hop number
1212        // can be cast or converted to a `u8` hop number,
1213        // so this check is important to prevent panics or incorrect behaviour.
1214        if hop_num == usize::from(u8::MAX) {
1215            return Err(internal!(
1216                "cannot add more hops to a circuit with `u8::MAX` hops"
1217            ));
1218        }
1219
1220        let hop_num = (hop_num as u8).into();
1221
1222        let hop = CircHop::new(self.unique_id, self.circ_id, hop_num, settings);
1223        self.hops.push(hop);
1224        self.crypto_in.add_layer(rev);
1225        self.crypto_out.add_layer(fwd);
1226        self.mutable.add_hop(peer_id, binding);
1227
1228        Ok(())
1229    }
1230
1231    /// Handle a RELAY cell on this circuit with stream ID 0.
1232    ///
1233    /// NOTE(prop349): this is part of Arti's "Base Circuit Hop Handler".
1234    /// This function returns a `CircProto` error if `msg` is an unsupported,
1235    /// unexpected, or otherwise invalid message:
1236    ///
1237    ///   * unexpected messages are rejected by returning an error using
1238    ///     [`unsupported_client_cell`]
1239    ///   * SENDME/TRUNCATED messages are rejected if they don't parse
1240    ///   * SENDME authentication tags are validated inside [`Circuit::handle_sendme`]
1241    ///   * conflux cells are handled in the client [`ConfluxMsgHandler`]
1242    ///
1243    /// The error is propagated all the way up to [`Circuit::handle_cell`],
1244    /// and eventually ends up being returned from the reactor's `run_once` function,
1245    /// causing it to shut down.
1246    fn handle_meta_cell(
1247        &mut self,
1248        handlers: &mut CellHandlers,
1249        hopnum: HopNum,
1250        msg: UnparsedRelayMsg,
1251    ) -> Result<Option<CircuitCmd>> {
1252        // SENDME cells and TRUNCATED get handled internally by the circuit.
1253
1254        // TODO: This pattern (Check command, try to decode, map error) occurs
1255        // several times, and would be good to extract simplify. Such
1256        // simplification is obstructed by a couple of factors: First, that
1257        // there is not currently a good way to get the RelayCmd from _type_ of
1258        // a RelayMsg.  Second, that decode() [correctly] consumes the
1259        // UnparsedRelayMsg.  I tried a macro-based approach, and didn't care
1260        // for it. -nickm
1261        if msg.cmd() == RelayCmd::SENDME {
1262            let sendme = msg
1263                .decode::<Sendme>()
1264                .map_err(|e| Error::from_bytes_err(e, "sendme message"))?
1265                .into_msg();
1266
1267            return Ok(Some(CircuitCmd::HandleSendMe {
1268                hop: hopnum,
1269                sendme,
1270            }));
1271        }
1272        if msg.cmd() == RelayCmd::TRUNCATED {
1273            let truncated = msg
1274                .decode::<Truncated>()
1275                .map_err(|e| Error::from_bytes_err(e, "truncated message"))?
1276                .into_msg();
1277            let reason = truncated.reason();
1278            debug!(
1279                circ_uniq_id = %self.unique_id,
1280                forward_circ_id = %self.circ_id,
1281                "Truncated from hop {}. Reason: {} [{}]",
1282                hopnum.display(),
1283                reason.human_str(),
1284                reason
1285            );
1286
1287            return Ok(Some(CircuitCmd::CleanShutdown));
1288        }
1289
1290        if msg.cmd() == RelayCmd::DROP {
1291            cfg_if::cfg_if! {
1292                if #[cfg(feature = "circ-padding")] {
1293                    return Ok(None);
1294                } else {
1295                    use crate::util::err::ExcessPadding;
1296                    return Err(Error::ExcessPadding(ExcessPadding::NoPaddingNegotiated, hopnum));
1297                }
1298            }
1299        }
1300
1301        trace!(
1302            circ_uniq_id = %self.unique_id,
1303            forward_circ_id = %self.circ_id,
1304            cell = ?msg,
1305            "Received meta-cell"
1306        );
1307
1308        #[cfg(feature = "conflux")]
1309        if matches!(
1310            msg.cmd(),
1311            RelayCmd::CONFLUX_LINK
1312                | RelayCmd::CONFLUX_LINKED
1313                | RelayCmd::CONFLUX_LINKED_ACK
1314                | RelayCmd::CONFLUX_SWITCH
1315        ) {
1316            let cmd = self.handle_conflux_msg(hopnum, msg)?;
1317            return Ok(cmd.map(CircuitCmd::from));
1318        }
1319
1320        if self.is_conflux_pending() {
1321            warn!(
1322                circ_uniq_id = %self.unique_id,
1323                forward_circ_id = %self.circ_id,
1324                "received unexpected cell {msg:?} on unlinked conflux circuit",
1325            );
1326            return Err(Error::CircProto(
1327                "Received unexpected cell on unlinked circuit".into(),
1328            ));
1329        }
1330
1331        // For all other command types, we'll only get them in response
1332        // to another command, which should have registered a responder.
1333        //
1334        // TODO: should the conflux state machine be a meta cell handler?
1335        // We'd need to add support for multiple meta handlers, and change the
1336        // MetaCellHandler API to support returning Option<RunOnceCmdInner>
1337        // (because some cells will require sending a response)
1338        if let Some(mut handler) = handlers.meta_handler.take() {
1339            // The handler has a TargetHop so we do a quick convert for equality check.
1340            if handler.expected_hop() == (self.unique_id(), hopnum).into() {
1341                // Somebody was waiting for a message -- maybe this message
1342                let ret = handler.handle_msg(msg, self);
1343                trace!(
1344                    circ_uniq_id = %self.unique_id,
1345                    forward_circ_id = %self.circ_id,
1346                    result = ?ret,
1347                    "meta handler completed",
1348                );
1349                match ret {
1350                    #[cfg(feature = "send-control-msg")]
1351                    Ok(MetaCellDisposition::Consumed) => {
1352                        handlers.meta_handler = Some(handler);
1353                        Ok(None)
1354                    }
1355                    Ok(MetaCellDisposition::ConversationFinished) => Ok(None),
1356                    #[cfg(feature = "send-control-msg")]
1357                    Ok(MetaCellDisposition::CloseCirc) => Ok(Some(CircuitCmd::CleanShutdown)),
1358                    Err(e) => Err(e),
1359                }
1360            } else {
1361                // Somebody wanted a message from a different hop!  Put this
1362                // one back.
1363                handlers.meta_handler = Some(handler);
1364
1365                unsupported_client_cell!(msg, hopnum)
1366            }
1367        } else {
1368            // No need to call shutdown here, since this error will
1369            // propagate to the reactor shut it down.
1370            unsupported_client_cell!(msg)
1371        }
1372    }
1373
1374    /// Handle a RELAY_SENDME cell on this circuit with stream ID 0.
1375    #[instrument(level = "trace", skip_all)]
1376    pub(super) fn handle_sendme(
1377        &mut self,
1378        hopnum: HopNum,
1379        msg: Sendme,
1380        signals: CongestionSignals,
1381    ) -> Result<Option<CircuitCmd>> {
1382        // Cloned, because we borrow mutably from self when we get the circhop.
1383        let runtime = self.runtime.clone();
1384
1385        // No need to call "shutdown" on errors in this function;
1386        // it's called from the reactor task and errors will propagate there.
1387        let hop = self
1388            .hop_mut(hopnum)
1389            .ok_or_else(|| Error::CircProto(format!("Couldn't find hop {}", hopnum.display())))?;
1390
1391        let tag = msg.into_sendme_tag().ok_or_else(||
1392                // Versions of Tor <=0.3.5 would omit a SENDME tag in this case;
1393                // but we don't support those any longer.
1394                 Error::CircProto("missing tag on circuit sendme".into()))?;
1395        // Update the CC object that we received a SENDME along with possible congestion signals.
1396        hop.ccontrol()
1397            .note_sendme_received(&runtime, tag, signals)?;
1398        Ok(None)
1399    }
1400
1401    /// Send a message onto the circuit's channel.
1402    ///
1403    /// If the channel is ready to accept messages, it will be sent immediately. If not, the message
1404    /// will be enqueued for sending at a later iteration of the reactor loop.
1405    ///
1406    /// `info` is the status returned from the padding controller when we told it we were queueing
1407    /// this data.  It should be provided whenever possible.
1408    ///
1409    /// # Note
1410    ///
1411    /// Making use of the enqueuing capabilities of this function is discouraged! You should first
1412    /// check whether the channel is ready to receive messages (`self.channel.poll_ready`), and
1413    /// ideally use this to implement backpressure (such that you do not read from other sources
1414    /// that would send here while you know you're unable to forward the messages on).
1415    #[instrument(level = "trace", skip_all)]
1416    async fn send_msg(
1417        &mut self,
1418        msg: AnyChanMsg,
1419        info: Option<QueuedCellPaddingInfo>,
1420    ) -> Result<()> {
1421        let cell = AnyChanCell::new(Some(self.circ_id), msg);
1422        // Note: this future is always `Ready`, so await won't block.
1423        Pin::new(&mut self.chan_sender)
1424            .send_unbounded((cell, info))
1425            .await?;
1426        Ok(())
1427    }
1428
1429    /// Remove all halfstreams that are expired at `now`.
1430    pub(super) fn remove_expired_halfstreams(&mut self, now: Instant) {
1431        self.hops.remove_expired_halfstreams(now);
1432    }
1433
1434    /// Return a reference to the hop corresponding to `hopnum`, if there is one.
1435    pub(super) fn hop(&self, hopnum: HopNum) -> Option<&CircHop> {
1436        self.hops.hop(hopnum)
1437    }
1438
1439    /// Return a mutable reference to the hop corresponding to `hopnum`, if there is one.
1440    pub(super) fn hop_mut(&mut self, hopnum: HopNum) -> Option<&mut CircHop> {
1441        self.hops.get_mut(hopnum)
1442    }
1443
1444    /// Begin a stream with the provided hop in this circuit.
1445    // TODO: see if there's a way that we can clean this up
1446    #[allow(clippy::too_many_arguments)]
1447    pub(super) fn begin_stream(
1448        &mut self,
1449        hop_num: HopNum,
1450        message: AnyRelayMsg,
1451        time_prov: &DynTimeProvider,
1452        cmd_checker: AnyCmdChecker,
1453        memquota: &StreamAccount,
1454    ) -> Result<(SendRelayCell, StreamId, ReactorStreamComponents)> {
1455        let Some(hop) = self.hop_mut(hop_num) else {
1456            return Err(internal!(
1457                "{}: Attempting to send a BEGIN cell to an unknown hop {hop_num:?}",
1458                self.unique_id,
1459            )
1460            .into());
1461        };
1462
1463        hop.begin_stream(message, time_prov, cmd_checker, memquota)
1464    }
1465
1466    /// Close the specified stream
1467    #[instrument(level = "trace", skip_all)]
1468    pub(super) async fn close_stream(
1469        &mut self,
1470        hop_num: HopNum,
1471        sid: StreamId,
1472        behav: CloseStreamBehavior,
1473        reason: streammap::TerminateReason,
1474        expiry: Instant,
1475    ) -> Result<()> {
1476        if let Some(hop) = self.hop_mut(hop_num) {
1477            let res = hop.close_stream(sid, behav, reason, expiry)?;
1478            if let Some(cell) = res {
1479                self.send_relay_cell(cell).await?;
1480            }
1481        }
1482        Ok(())
1483    }
1484
1485    /// Returns true if there are any streams on this circuit
1486    ///
1487    /// Important: this function locks the stream map of its each of the [`CircHop`]s
1488    /// in this circuit, so it must **not** be called from any function where the
1489    /// stream map lock is held.
1490    pub(super) fn has_streams(&self) -> bool {
1491        self.hops.has_streams()
1492    }
1493
1494    /// The number of hops in this circuit.
1495    pub(super) fn num_hops(&self) -> u8 {
1496        // `Circuit::add_hop` checks to make sure that we never have more than `u8::MAX` hops,
1497        // so `self.hops.len()` should be safe to cast to a `u8`.
1498        // If that assumption is violated,
1499        // we choose to panic rather than silently use the wrong hop due to an `as` cast.
1500        self.hops
1501            .len()
1502            .try_into()
1503            .expect("`hops.len()` has more than `u8::MAX` hops")
1504    }
1505
1506    /// Check whether this circuit has any hops.
1507    pub(super) fn has_hops(&self) -> bool {
1508        !self.hops.is_empty()
1509    }
1510
1511    /// Get the `HopNum` of the last hop, if this circuit is non-empty.
1512    ///
1513    /// Returns `None` if the circuit has no hops.
1514    pub(super) fn last_hop_num(&self) -> Option<HopNum> {
1515        let num_hops = self.num_hops();
1516        if num_hops == 0 {
1517            // asked for the last hop, but there are no hops
1518            return None;
1519        }
1520        Some(HopNum::from(num_hops - 1))
1521    }
1522
1523    /// Get the path of the circuit.
1524    ///
1525    /// **Warning:** Do not call while already holding the [`Self::mutable`] lock.
1526    pub(super) fn path(&self) -> Arc<path::Path> {
1527        self.mutable.path()
1528    }
1529
1530    /// Return a ClockSkew declaring how much clock skew the other side of this channel
1531    /// claimed that we had when we negotiated the connection.
1532    pub(super) fn clock_skew(&self) -> ClockSkew {
1533        self.channel.clock_skew()
1534    }
1535
1536    /// Does congestion control use stream SENDMEs for the given `hop`?
1537    ///
1538    /// Returns `None` if `hop` doesn't exist.
1539    pub(super) fn uses_stream_sendme(&self, hop: HopNum) -> Option<bool> {
1540        let hop = self.hop(hop)?;
1541        Some(hop.ccontrol().uses_stream_sendme())
1542    }
1543
1544    /// Returns whether this is a conflux circuit that is not linked yet.
1545    pub(super) fn is_conflux_pending(&self) -> bool {
1546        let Some(status) = self.conflux_status() else {
1547            return false;
1548        };
1549
1550        status != ConfluxStatus::Linked
1551    }
1552
1553    /// Returns the conflux status of this circuit.
1554    ///
1555    /// Returns `None` if this is not a conflux circuit.
1556    pub(super) fn conflux_status(&self) -> Option<ConfluxStatus> {
1557        cfg_if::cfg_if! {
1558            if #[cfg(feature = "conflux")] {
1559                self.conflux_handler
1560                    .as_ref()
1561                    .map(|handler| handler.status())
1562            } else {
1563                None
1564            }
1565        }
1566    }
1567
1568    /// Returns initial RTT on this leg, measured in the conflux handshake.
1569    #[cfg(feature = "conflux")]
1570    pub(super) fn init_rtt(&self) -> Option<Duration> {
1571        self.conflux_handler
1572            .as_ref()
1573            .map(|handler| handler.init_rtt())?
1574    }
1575
1576    /// Start or stop padding at the given hop.
1577    ///
1578    /// Replaces any previous padder at that hop.
1579    ///
1580    /// Return an error if that hop doesn't exist.
1581    #[cfg(feature = "circ-padding-manual")]
1582    pub(super) fn set_padding_at_hop(
1583        &self,
1584        hop: HopNum,
1585        padder: Option<padding::CircuitPadder>,
1586    ) -> Result<()> {
1587        if self.hop(hop).is_none() {
1588            return Err(Error::NoSuchHop);
1589        }
1590        self.padding_ctrl.install_padder_padding_at_hop(hop, padder);
1591        Ok(())
1592    }
1593
1594    /// Determine how exactly to handle a request to handle padding.
1595    ///
1596    /// This is fairly complicated; see the maybenot documentation for more information.
1597    ///
1598    /// ## Limitations
1599    ///
1600    /// In our current padding implementation, a circuit is either blocked or not blocked:
1601    /// we do not keep track of which hop is actually doing the blocking.
1602    #[cfg(feature = "circ-padding")]
1603    fn padding_disposition(&self, send_padding: &padding::SendPadding) -> CircPaddingDisposition {
1604        crate::circuit::padding::padding_disposition(
1605            send_padding,
1606            &self.chan_sender,
1607            self.padding_block.as_ref(),
1608        )
1609    }
1610
1611    /// Handle a request from our padding subsystem to send a padding packet.
1612    #[cfg(feature = "circ-padding")]
1613    pub(super) async fn send_padding(&mut self, send_padding: padding::SendPadding) -> Result<()> {
1614        use CircPaddingDisposition::*;
1615
1616        let target_hop = send_padding.hop;
1617
1618        match self.padding_disposition(&send_padding) {
1619            QueuePaddingNormally => {
1620                let queue_info = self.padding_ctrl.queued_padding(target_hop, send_padding);
1621                self.queue_padding_cell_for_hop(target_hop, queue_info)
1622                    .await?;
1623            }
1624            QueuePaddingAndBypass => {
1625                let queue_info = self.padding_ctrl.queued_padding(target_hop, send_padding);
1626                self.queue_padding_cell_for_hop(target_hop, queue_info)
1627                    .await?;
1628            }
1629            TreatQueuedCellAsPadding => {
1630                self.padding_ctrl
1631                    .replaceable_padding_already_queued(target_hop, send_padding);
1632            }
1633        }
1634        Ok(())
1635    }
1636
1637    /// Generate and encrypt a padding cell, and send it to a targeted hop.
1638    ///
1639    /// Ignores any padding-based blocking.
1640    #[cfg(feature = "circ-padding")]
1641    async fn queue_padding_cell_for_hop(
1642        &mut self,
1643        target_hop: HopNum,
1644        queue_info: Option<QueuedCellPaddingInfo>,
1645    ) -> Result<()> {
1646        use tor_cell::relaycell::msg::Drop as DropMsg;
1647        let msg = SendRelayCell {
1648            hop: Some(target_hop),
1649            // TODO circpad: we will probably want padding machines that can send EARLY cells.
1650            early: false,
1651            cell: AnyRelayMsgOuter::new(None, DropMsg::default().into()),
1652        };
1653        self.send_relay_cell_inner(msg, queue_info).await
1654    }
1655
1656    /// Enable padding-based blocking,
1657    /// or change the rule for padding-based blocking to the one in `block`.
1658    #[cfg(feature = "circ-padding")]
1659    pub(super) fn start_blocking_for_padding(&mut self, block: padding::StartBlocking) {
1660        self.chan_sender.start_blocking();
1661        self.padding_block = Some(block);
1662    }
1663
1664    /// Disable padding-based blocking.
1665    #[cfg(feature = "circ-padding")]
1666    pub(super) fn stop_blocking_for_padding(&mut self) {
1667        self.chan_sender.stop_blocking();
1668        self.padding_block = None;
1669    }
1670
1671    /// The estimated circuit build timeout for a circuit of the specified length.
1672    pub(super) fn estimate_cbt(&self, length: usize) -> Duration {
1673        self.timeouts.circuit_build_timeout(length)
1674    }
1675}
1676
1677impl Drop for Circuit {
1678    fn drop(&mut self) {
1679        let _ = self.channel.close_circuit(self.circ_id);
1680    }
1681}