Skip to main content

tor_proto/client/reactor/
control.rs

1//! Module providing [`CtrlMsg`].
2
3use super::circuit::extender::CircuitExtender;
4use super::{
5    CircuitHandshake, CloseStreamBehavior, MetaCellHandler, Reactor, ReactorResultChannel,
6    RunOnceCmdInner, SendRelayCell,
7};
8use crate::Result;
9use crate::circuit::celltypes::CreateResponse;
10use crate::circuit::circhop::{HopSettings, ReactorStreamComponents};
11#[cfg(feature = "circ-padding-manual")]
12use crate::client::circuit::padding;
13use crate::client::circuit::path;
14use crate::client::reactor::{NoJoinPointError, NtorClient, ReactorError};
15use crate::client::{HopLocation, TargetHop};
16use crate::crypto::binding::CircuitBinding;
17use crate::crypto::cell::{InboundClientLayer, OutboundClientLayer};
18use crate::crypto::handshake::ntor_v3::{NtorV3Client, NtorV3PublicKey};
19use crate::memquota::StreamAccount;
20use crate::stream::cmdcheck::AnyCmdChecker;
21use crate::streammap;
22use crate::util::skew::ClockSkew;
23use crate::util::tunnel_activity::TunnelActivity;
24#[cfg(test)]
25use crate::{circuit::UniqId, client::circuit::CircParameters, crypto::cell::HopNum};
26use tor_cell::chancell::msg::HandshakeType;
27use tor_cell::relaycell::flow_ctrl::XonKBpsEwma;
28use tor_cell::relaycell::msg::{AnyRelayMsg, Sendme};
29use tor_cell::relaycell::{AnyRelayMsgOuter, RelayCellFormat, StreamId};
30use tor_error::{Bug, bad_api_usage, internal, into_bad_api_usage};
31use tracing::{debug, trace};
32#[cfg(feature = "hs-service")]
33use {
34    crate::client::reactor::IncomingStreamRequestHandler,
35    crate::stream::IncomingStreamRequestFilter, crate::stream::incoming::StreamReqSender,
36};
37
38#[cfg(test)]
39use tor_cell::relaycell::msg::SendmeTag;
40
41#[cfg(feature = "conflux")]
42use super::{Circuit, ConfluxLinkResultChannel};
43
44use oneshot_fused_workaround as oneshot;
45
46use crate::crypto::handshake::ntor::NtorPublicKey;
47use tor_linkspec::{EncodedLinkSpec, OwnedChanTarget};
48
49use std::result::Result as StdResult;
50
51/// A message telling the reactor to do something.
52///
53/// For each `CtrlMsg`, the reactor will send a cell on the underlying channel.
54///
55/// The difference between this and [`CtrlCmd`] is that `CtrlMsg`s
56/// cause the reactor to send cells on the reactor's `chan_sender`,
57/// whereas `CtrlCmd` do not.
58#[derive(educe::Educe)]
59#[educe(Debug)]
60pub(crate) enum CtrlMsg {
61    /// Create the first hop of this circuit.
62    Create {
63        /// A oneshot channel on which we'll receive the creation response.
64        recv_created: oneshot::Receiver<CreateResponse>,
65        /// The handshake type to use for the first hop.
66        handshake: CircuitHandshake,
67        /// Other parameters relevant for circuit creation.
68        settings: HopSettings,
69        /// Oneshot channel to notify on completion.
70        done: ReactorResultChannel<()>,
71    },
72    /// Extend a circuit by one hop, using the ntor handshake.
73    ExtendNtor {
74        /// The peer that we're extending to.
75        ///
76        /// Used to extend our record of the circuit's path.
77        peer_id: OwnedChanTarget,
78        /// The handshake type to use for this hop.
79        public_key: NtorPublicKey,
80        /// Information about how to connect to the relay we're extending to.
81        linkspecs: Vec<EncodedLinkSpec>,
82        /// Other parameters we are negotiating.
83        settings: HopSettings,
84        /// Oneshot channel to notify on completion.
85        done: ReactorResultChannel<()>,
86    },
87    /// Extend a circuit by one hop, using the ntorv3 handshake.
88    ExtendNtorV3 {
89        /// The peer that we're extending to.
90        ///
91        /// Used to extend our record of the circuit's path.
92        peer_id: OwnedChanTarget,
93        /// The handshake type to use for this hop.
94        public_key: NtorV3PublicKey,
95        /// Information about how to connect to the relay we're extending to.
96        linkspecs: Vec<EncodedLinkSpec>,
97        /// Other parameters we are negotiating.
98        settings: HopSettings,
99        /// Oneshot channel to notify on completion.
100        done: ReactorResultChannel<()>,
101    },
102    /// Begin a stream with the provided hop in this circuit.
103    ///
104    /// Allocates a stream ID, and sends the provided message to that hop.
105    BeginStream {
106        /// The hop number to begin the stream with.
107        hop: TargetHop,
108        /// The message to send.
109        message: AnyRelayMsg,
110        /// The stream account to use for anything we allocate for the purpose of this stream.
111        memquota: StreamAccount,
112        /// Oneshot channel to notify on completion, with the allocated stream ID.
113        done: ReactorResultChannel<(
114            StreamId,
115            HopLocation,
116            RelayCellFormat,
117            ReactorStreamComponents,
118        )>,
119        /// A `CmdChecker` to keep track of which message types are acceptable.
120        cmd_checker: AnyCmdChecker,
121    },
122    /// Close the specified pending incoming stream, sending the provided END message.
123    ///
124    /// A stream is said to be pending if the message for initiating the stream was received but
125    /// not has not been responded to yet.
126    ///
127    /// This should be used by responders for closing pending incoming streams initiated by the
128    /// other party on the circuit.
129    #[cfg(feature = "hs-service")]
130    ClosePendingStream {
131        /// The hop number the stream is on.
132        hop: HopLocation,
133        /// The stream ID to send the END for.
134        stream_id: StreamId,
135        /// The END message to send, if any.
136        message: CloseStreamBehavior,
137        /// Oneshot channel to notify on completion.
138        done: ReactorResultChannel<()>,
139    },
140    /// Send a given control message on this circuit.
141    #[cfg(feature = "send-control-msg")]
142    SendMsg {
143        /// The hop to receive this message.
144        hop: TargetHop,
145        /// The message to send.
146        msg: AnyRelayMsg,
147        /// A sender that we use to tell the caller that the message was sent
148        /// and the handler installed.
149        sender: oneshot::Sender<Result<()>>,
150    },
151    /// Send a given control message on this circuit, and install a control-message handler to
152    /// receive responses.
153    #[cfg(feature = "send-control-msg")]
154    SendMsgAndInstallHandler {
155        /// The message to send, if any
156        msg: Option<AnyRelayMsgOuter>,
157        /// A message handler to install.
158        ///
159        /// If this is `None`, there must already be a message handler installed
160        #[educe(Debug(ignore))]
161        handler: Option<Box<dyn MetaCellHandler + Send + 'static>>,
162        /// A sender that we use to tell the caller that the message was sent
163        /// and the handler installed.
164        sender: oneshot::Sender<Result<()>>,
165    },
166    /// Inform the reactor that there's a flow control update for a given stream.
167    ///
168    /// The reactor will decide how to handle this update depending on the type of flow control and
169    /// the current state of the stream.
170    FlowCtrlUpdate {
171        /// The type of flow control update, and any associated metadata.
172        msg: FlowCtrlMsg,
173        /// The stream ID that the update is for.
174        stream_id: StreamId,
175        /// The hop that the stream is on.
176        hop: HopLocation,
177    },
178    /// Get the clock skew claimed by the first hop of the circuit.
179    FirstHopClockSkew {
180        /// Oneshot channel to return the clock skew.
181        answer: oneshot::Sender<StdResult<ClockSkew, Bug>>,
182    },
183    /// Link the specified circuits into the current tunnel,
184    /// to form a multi-path tunnel.
185    #[cfg(feature = "conflux")]
186    #[allow(unused)] // TODO(conflux)
187    LinkCircuits {
188        /// The circuits to link into the tunnel,
189        #[educe(Debug(ignore))]
190        circuits: Vec<Circuit>,
191        /// Oneshot channel to notify sender when all the specified circuits have finished linking,
192        /// or have failed to link.
193        ///
194        /// A client circuit is said to be fully linked once the `RELAY_CONFLUX_LINKED_ACK` is sent
195        /// (see [set construction]).
196        ///
197        /// [set construction]: https://spec.torproject.org/proposals/329-traffic-splitting.html#set-construction
198        answer: ConfluxLinkResultChannel,
199    },
200}
201
202/// A message telling the reactor to do something.
203///
204/// The difference between this and [`CtrlMsg`] is that `CtrlCmd`s
205/// never cause cells to sent on the channel,
206/// while `CtrlMsg`s potentially do: `CtrlMsg`s are mapped to [`RunOnceCmdInner`] commands,
207/// some of which instruct the reactor to send cells down the channel.
208#[derive(educe::Educe)]
209#[educe(Debug)]
210pub(crate) enum CtrlCmd {
211    /// Shut down the reactor.
212    Shutdown,
213    /// Extend the circuit by one hop, in response to an out-of-band handshake.
214    ///
215    /// (This is used for onion services, where the negotiation takes place in
216    /// INTRODUCE and RENDEZVOUS messages.)
217    #[cfg(feature = "hs-common")]
218    ExtendVirtual {
219        /// The cryptographic algorithms and keys to use when communicating with
220        /// the newly added hop.
221        #[educe(Debug(ignore))]
222        cell_crypto: (
223            Box<dyn OutboundClientLayer + Send>,
224            Box<dyn InboundClientLayer + Send>,
225            Option<CircuitBinding>,
226        ),
227        /// A set of parameters to negotiate with this hop.
228        settings: HopSettings,
229        /// Oneshot channel to notify on completion.
230        done: ReactorResultChannel<()>,
231    },
232    /// Resolve a given [`TargetHop`] into a precise [`HopLocation`].
233    ResolveTargetHop {
234        /// The target hop to resolve.
235        hop: TargetHop,
236        /// Oneshot channel to notify on completion.
237        done: ReactorResultChannel<HopLocation>,
238    },
239    /// Begin accepting streams on this circuit.
240    #[cfg(feature = "hs-service")]
241    AwaitStreamRequest {
242        /// A channel for sending information about an incoming stream request.
243        incoming_sender: StreamReqSender,
244        /// A `CmdChecker` to keep track of which message types are acceptable.
245        cmd_checker: AnyCmdChecker,
246        /// Oneshot channel to notify on completion.
247        done: ReactorResultChannel<()>,
248        /// The hop that is allowed to create streams.
249        hop: TargetHop,
250        /// A filter used to check requests before passing them on.
251        #[educe(Debug(ignore))]
252        #[cfg(feature = "hs-service")]
253        filter: Box<dyn IncomingStreamRequestFilter>,
254    },
255    /// Request the binding key of a target hop.
256    #[cfg(feature = "hs-service")]
257    GetBindingKey {
258        /// The hop for which we want the key.
259        hop: TargetHop,
260        /// Oneshot channel to notify on completion.
261        done: ReactorResultChannel<Option<CircuitBinding>>,
262    },
263    /// (tests only) Add a hop to the list of hops on this circuit, with dummy cryptography.
264    #[cfg(test)]
265    AddFakeHop {
266        relay_cell_format: RelayCellFormat,
267        fwd_lasthop: bool,
268        rev_lasthop: bool,
269        peer_id: path::HopDetail,
270        // `CircParameters` is large and this command is test-only, so we box it.
271        params: Box<CircParameters>,
272        done: ReactorResultChannel<()>,
273    },
274    /// (tests only) Get the send window and expected tags for a given hop.
275    #[cfg(test)]
276    QuerySendWindow {
277        hop: HopNum,
278        leg: UniqId,
279        done: ReactorResultChannel<(u32, Vec<SendmeTag>)>,
280    },
281    /// Shut down the reactor, and return the underlying [`Circuit`],
282    /// if the tunnel is not multi-path.
283    ///
284    /// Returns an error if called on a multi-path reactor.
285    #[cfg(feature = "conflux")]
286    #[allow(unused)] // TODO(conflux)
287    ShutdownAndReturnCircuit {
288        /// Oneshot channel to return the underlying [`Circuit`],
289        /// or an error if the reactor's tunnel is multi-path.
290        answer: oneshot::Sender<StdResult<Circuit, Bug>>,
291    },
292
293    /// Install or remove a [`padding::CircuitPadder`] for a given hop.
294    ///
295    /// Any existing `CircuitPadder` at that hop is replaced.
296    #[cfg(feature = "circ-padding-manual")]
297    SetPadder {
298        /// The hop to modify.
299        hop: HopLocation,
300        /// The Padder to install, or None to remove any existing padder.
301        padder: Option<padding::CircuitPadder>,
302        /// A sender to alert after we've changed the padding.
303        sender: oneshot::Sender<Result<()>>,
304    },
305
306    /// Yield the most active [`TunnelActivity`] for any hop on any leg of this tunnel.
307    GetTunnelActivity {
308        /// A sender to receive the reply.
309        sender: oneshot::Sender<TunnelActivity>,
310    },
311}
312
313/// A flow control update message.
314#[derive(Debug)]
315pub(crate) enum FlowCtrlMsg {
316    /// Send a SENDME message on this stream.
317    Sendme,
318    /// Send an XON message on this stream with the given rate.
319    Xon(XonKBpsEwma),
320}
321
322/// A control message handler object. Keep a reference to the Reactor tying its lifetime to it.
323///
324/// Its `handle_msg` and `handle_cmd` handlers decide how messages and commands,
325/// respectively, are handled.
326pub(crate) struct ControlHandler<'a> {
327    /// Reference to the reactor of this
328    reactor: &'a mut Reactor,
329}
330
331impl<'a> ControlHandler<'a> {
332    /// Constructor.
333    pub(crate) fn new(reactor: &'a mut Reactor) -> Self {
334        Self { reactor }
335    }
336
337    /// Handle a control message.
338    pub(super) fn handle_msg(&mut self, msg: CtrlMsg) -> Result<Option<RunOnceCmdInner>> {
339        trace!(
340            tunnel_id = %self.reactor.tunnel_id,
341            msg = ?msg,
342            "reactor received control message"
343        );
344
345        match msg {
346            // This is handled earlier, since it requires blocking.
347            CtrlMsg::Create { done, .. } => {
348                if self.reactor.circuits.len() == 1 {
349                    // This should've been handled in Reactor::run_once()
350                    // (ControlHandler::handle_msg() is never called before wait_for_create()).
351                    debug_assert!(self.reactor.circuits.single_leg()?.has_hops());
352                    // Don't care if the receiver goes away
353                    let _ = done.send(Err(tor_error::bad_api_usage!(
354                        "cannot create first hop twice"
355                    )
356                    .into()));
357                } else {
358                    // Don't care if the receiver goes away
359                    let _ = done.send(Err(tor_error::bad_api_usage!(
360                        "cannot create first hop on multipath tunnel"
361                    )
362                    .into()));
363                }
364
365                Ok(None)
366            }
367            CtrlMsg::ExtendNtor {
368                peer_id,
369                public_key,
370                linkspecs,
371                settings,
372                done,
373            } => {
374                let Ok(circ) = self.reactor.circuits.single_leg_mut() else {
375                    // Don't care if the receiver goes away
376                    let _ = done.send(Err(tor_error::bad_api_usage!(
377                        "cannot extend multipath tunnel"
378                    )
379                    .into()));
380
381                    return Ok(None);
382                };
383
384                let (extender, cell) = CircuitExtender::<NtorClient>::begin(
385                    peer_id,
386                    HandshakeType::NTOR,
387                    &public_key,
388                    linkspecs,
389                    settings,
390                    &(),
391                    circ,
392                    done,
393                )?;
394                self.reactor
395                    .cell_handlers
396                    .set_meta_handler(Box::new(extender))?;
397
398                Ok(Some(RunOnceCmdInner::Send {
399                    leg: circ.unique_id(),
400                    cell,
401                    done: None,
402                }))
403            }
404            CtrlMsg::ExtendNtorV3 {
405                peer_id,
406                public_key,
407                linkspecs,
408                settings,
409                done,
410            } => {
411                let Ok(circ) = self.reactor.circuits.single_leg_mut() else {
412                    // Don't care if the receiver goes away
413                    let _ = done.send(Err(tor_error::bad_api_usage!(
414                        "cannot extend multipath tunnel"
415                    )
416                    .into()));
417
418                    return Ok(None);
419                };
420
421                let client_extensions = settings.circuit_request_extensions()?;
422
423                let (extender, cell) = CircuitExtender::<NtorV3Client>::begin(
424                    peer_id,
425                    HandshakeType::NTOR_V3,
426                    &public_key,
427                    linkspecs,
428                    settings,
429                    &client_extensions,
430                    circ,
431                    done,
432                )?;
433                self.reactor
434                    .cell_handlers
435                    .set_meta_handler(Box::new(extender))?;
436
437                Ok(Some(RunOnceCmdInner::Send {
438                    leg: circ.unique_id(),
439                    cell,
440                    done: None,
441                }))
442            }
443            CtrlMsg::BeginStream {
444                hop,
445                message,
446                memquota,
447                done,
448                cmd_checker,
449            } => {
450                // If resolving the hop fails,
451                // we want to report an error back to the initiator and not shut down the reactor.
452                let hop_location = match self.reactor.resolve_target_hop(hop) {
453                    Ok(x) => x,
454                    Err(e) => {
455                        let e = into_bad_api_usage!("Could not resolve {hop:?}")(e);
456                        // don't care if receiver goes away
457                        let _ = done.send(Err(e.into()));
458                        return Ok(None);
459                    }
460                };
461                let (leg_id, hop_num) = match self.reactor.resolve_hop_location(hop_location) {
462                    Ok(x) => x,
463                    Err(e) => {
464                        let e = into_bad_api_usage!("Could not resolve {hop_location:?}")(e);
465                        // don't care if receiver goes away
466                        let _ = done.send(Err(e.into()));
467                        return Ok(None);
468                    }
469                };
470                let circ = match self.reactor.circuits.leg_mut(leg_id) {
471                    Some(x) => x,
472                    None => {
473                        let e = bad_api_usage!("Circuit leg {leg_id:?} does not exist");
474                        // don't care if receiver goes away
475                        let _ = done.send(Err(e.into()));
476                        return Ok(None);
477                    }
478                };
479
480                let result = circ.begin_stream(
481                    hop_num,
482                    message,
483                    &self.reactor.runtime,
484                    cmd_checker,
485                    &memquota,
486                );
487
488                let (cell, stream_id, stream_components) = match result {
489                    Ok((cell, stream_id, receiver)) => (cell, stream_id, receiver),
490                    Err(e) => {
491                        // don't care if receiver goes away.
492                        let _ = done.send(Err(e.clone()));
493                        return Err(e);
494                    }
495                };
496
497                Ok(Some(RunOnceCmdInner::BeginStream {
498                    cell,
499                    stream_id,
500                    hop: hop_location,
501                    leg: leg_id,
502                    stream_components,
503                    done,
504                }))
505            }
506            #[cfg(feature = "hs-service")]
507            CtrlMsg::ClosePendingStream {
508                hop,
509                stream_id,
510                message,
511                done,
512            } => Ok(Some(RunOnceCmdInner::CloseStream {
513                hop,
514                sid: stream_id,
515                behav: message,
516                reason: streammap::TerminateReason::ExplicitEnd,
517                done: Some(done),
518            })),
519            CtrlMsg::FlowCtrlUpdate {
520                msg,
521                stream_id,
522                hop,
523            } => {
524                match msg {
525                    FlowCtrlMsg::Sendme => {
526                        let (leg_id, hop_num) = match self.reactor.resolve_hop_location(hop) {
527                            Ok(x) => x,
528                            Err(NoJoinPointError) => {
529                                // A stream tried to send a stream-level SENDME message to the join point of
530                                // a tunnel that has never had a join point. Currently in arti, only a
531                                // `StreamTarget` asks us to send a stream-level SENDME, and this tunnel
532                                // originally created the `StreamTarget` to begin with. So this is a
533                                // legitimate bug somewhere in the tunnel code.
534                                return Err(
535                                    internal!(
536                                        "Could not send a stream-level SENDME to a join point on a tunnel without a join point",
537                                    )
538                                    .into()
539                                );
540                            }
541                        };
542
543                        // Congestion control decides if we can send stream level SENDMEs or not.
544                        let sendme_required = match self.reactor.uses_stream_sendme(leg_id, hop_num)
545                        {
546                            Some(x) => x,
547                            None => {
548                                // The leg/hop has disappeared. This is fine since the stream may have ended
549                                // and been cleaned up while this `CtrlMsg::SendSendme` message was queued.
550                                // It is possible that is a bug and this is an incorrect leg/hop number, but
551                                // it's not currently possible to differentiate between an incorrect leg/hop
552                                // number and a circuit hop that has been closed.
553                                debug!(
554                                    "Could not send a stream-level SENDME on a hop that does not exist. Ignoring."
555                                );
556                                return Ok(None);
557                            }
558                        };
559
560                        if !sendme_required {
561                            // Nothing to do, so discard the SENDME.
562                            return Ok(None);
563                        }
564
565                        let sendme = Sendme::new_empty();
566                        let cell = AnyRelayMsgOuter::new(Some(stream_id), sendme.into());
567
568                        let cell = SendRelayCell {
569                            hop: Some(hop_num),
570                            early: false,
571                            cell,
572                        };
573
574                        Ok(Some(RunOnceCmdInner::Send {
575                            leg: leg_id,
576                            cell,
577                            done: None,
578                        }))
579                    }
580                    FlowCtrlMsg::Xon(rate) => Ok(Some(RunOnceCmdInner::MaybeSendXon {
581                        rate,
582                        hop,
583                        stream_id,
584                    })),
585                }
586            }
587            // TODO(conflux): this should specify which leg to send the msg on
588            // (currently we send it down the primary leg).
589            //
590            // This will involve updating ClientCIrc::send_raw_msg() to take a
591            // leg id argument (which is a breaking change.
592            #[cfg(feature = "send-control-msg")]
593            CtrlMsg::SendMsg { hop, msg, sender } => {
594                let Some((leg_id, hop_num)) = self.reactor.target_hop_to_hopnum_id(hop) else {
595                    // Don't care if receiver goes away
596                    let _ = sender.send(Err(bad_api_usage!("Unknown {hop:?}").into()));
597                    return Ok(None);
598                };
599
600                let cell = AnyRelayMsgOuter::new(None, msg);
601                let cell = SendRelayCell {
602                    hop: Some(hop_num),
603                    early: false,
604                    cell,
605                };
606
607                Ok(Some(RunOnceCmdInner::Send {
608                    leg: leg_id,
609                    cell,
610                    done: Some(sender),
611                }))
612            }
613            // TODO(conflux): this should specify which leg to send the msg on
614            // (currently we send it down the primary leg)
615            #[cfg(feature = "send-control-msg")]
616            CtrlMsg::SendMsgAndInstallHandler {
617                msg,
618                handler,
619                sender,
620            } => Ok(Some(RunOnceCmdInner::SendMsgAndInstallHandler {
621                msg,
622                handler,
623                done: sender,
624            })),
625            CtrlMsg::FirstHopClockSkew { answer } => {
626                Ok(Some(RunOnceCmdInner::FirstHopClockSkew { answer }))
627            }
628            #[cfg(feature = "conflux")]
629            CtrlMsg::LinkCircuits { circuits, answer } => {
630                Ok(Some(RunOnceCmdInner::Link { circuits, answer }))
631            }
632        }
633    }
634
635    /// Handle a control command.
636    #[allow(clippy::needless_pass_by_value)] // Needed when conflux is enabled
637    pub(super) fn handle_cmd(&mut self, msg: CtrlCmd) -> StdResult<(), ReactorError> {
638        trace!(
639            tunnel_id = %self.reactor.tunnel_id,
640            msg = ?msg,
641            "reactor received control command"
642        );
643
644        match msg {
645            CtrlCmd::Shutdown => self.reactor.handle_shutdown().map(|_| ()),
646            #[cfg(feature = "hs-common")]
647            #[allow(unreachable_code)]
648            CtrlCmd::ExtendVirtual {
649                cell_crypto,
650                settings,
651                done,
652            } => {
653                let (outbound, inbound, binding) = cell_crypto;
654
655                // TODO HS: Perhaps this should describe the onion service, or
656                // describe why the virtual hop was added, or something?
657                let peer_id = path::HopDetail::Virtual;
658
659                let Ok(leg) = self.reactor.circuits.single_leg_mut() else {
660                    // Don't care if the receiver goes away
661                    let _ = done.send(Err(tor_error::bad_api_usage!(
662                        "cannot extend multipath tunnel"
663                    )
664                    .into()));
665
666                    return Ok(());
667                };
668
669                trace!(circ=%self.reactor.tunnel_id, settings=?&settings,
670                    "Adding virtual hop to circuit");
671                leg.add_hop(peer_id, outbound, inbound, binding, &settings)?;
672                let _ = done.send(Ok(()));
673
674                Ok(())
675            }
676            CtrlCmd::ResolveTargetHop { hop, done } => {
677                let _ = done.send(
678                    self.reactor
679                        .resolve_target_hop(hop)
680                        .map_err(|_| crate::util::err::Error::NoSuchHop),
681                );
682                Ok(())
683            }
684            #[cfg(feature = "hs-service")]
685            CtrlCmd::AwaitStreamRequest {
686                cmd_checker,
687                incoming_sender,
688                hop,
689                done,
690                filter,
691            } => {
692                let Some((_, hop_num)) = self.reactor.target_hop_to_hopnum_id(hop) else {
693                    let _ = done.send(Err(crate::Error::NoSuchHop));
694                    return Ok(());
695                };
696                // TODO: At some point we might want to add a CtrlCmd for
697                // de-registering the handler.  See comments on `allow_stream_requests`.
698                let handler = IncomingStreamRequestHandler {
699                    incoming_sender,
700                    cmd_checker,
701                    hop_num: Some(hop_num),
702                    filter,
703                };
704
705                let ret = self
706                    .reactor
707                    .cell_handlers
708                    .set_incoming_stream_req_handler(handler);
709                let _ = done.send(ret); // don't care if the corresponding receiver goes away.
710
711                Ok(())
712            }
713            #[cfg(feature = "hs-service")]
714            CtrlCmd::GetBindingKey { hop, done } => {
715                let Some((leg_id, hop_num)) = self.reactor.target_hop_to_hopnum_id(hop) else {
716                    let _ = done.send(Err(tor_error::internal!(
717                        "Unknown TargetHop when getting binding key"
718                    )
719                    .into()));
720                    return Ok(());
721                };
722                let Some(circuit) = self.reactor.circuits.leg(leg_id) else {
723                    let _ = done.send(Err(tor_error::bad_api_usage!(
724                        "Unknown circuit id {leg_id} when getting binding key"
725                    )
726                    .into()));
727                    return Ok(());
728                };
729                // Get the binding key from the mutable state and send it back.
730                let key = circuit.mutable().binding_key(hop_num);
731                let _ = done.send(Ok(key));
732
733                Ok(())
734            }
735            #[cfg(test)]
736            CtrlCmd::AddFakeHop {
737                relay_cell_format,
738                fwd_lasthop,
739                rev_lasthop,
740                peer_id,
741                params,
742                done,
743            } => {
744                let Ok(leg) = self.reactor.circuits.single_leg_mut() else {
745                    // Don't care if the receiver goes away
746                    let _ = done.send(Err(tor_error::bad_api_usage!(
747                        "cannot add fake hop to multipath tunnel"
748                    )
749                    .into()));
750
751                    return Ok(());
752                };
753
754                leg.handle_add_fake_hop(
755                    relay_cell_format,
756                    fwd_lasthop,
757                    rev_lasthop,
758                    peer_id,
759                    &params,
760                    done,
761                );
762
763                Ok(())
764            }
765            #[cfg(test)]
766            CtrlCmd::QuerySendWindow { hop, leg, done } => {
767                // Immediately invoked function means that errors will be sent to the channel.
768                let _ = done.send((|| {
769                    let leg = self.reactor.circuits.leg_mut(leg).ok_or_else(|| {
770                        bad_api_usage!("cannot query send window of non-existent circuit")
771                    })?;
772
773                    let hop = leg.hop_mut(hop).ok_or(bad_api_usage!(
774                        "received QuerySendWindow for unknown hop {}",
775                        hop.display()
776                    ))?;
777
778                    Ok(hop.send_window_and_expected_tags())
779                })());
780
781                Ok(())
782            }
783            #[cfg(feature = "conflux")]
784            CtrlCmd::ShutdownAndReturnCircuit { answer } => {
785                self.reactor.handle_shutdown_and_return_circuit(answer)
786            }
787            #[cfg(feature = "circ-padding-manual")]
788            CtrlCmd::SetPadder {
789                hop,
790                padder,
791                sender,
792            } => {
793                let result = self.reactor.set_padding_at_hop(hop, padder);
794                let _ = sender.send(result);
795                Ok(())
796            }
797            CtrlCmd::GetTunnelActivity { sender } => {
798                let count = self.reactor.circuits.tunnel_activity();
799                let _ = sender.send(count);
800                Ok(())
801            }
802        }
803    }
804}