Skip to main content

tor_proto/client/reactor/
conflux.rs

1//! Conflux-related functionality
2
3// TODO: replace Itertools::exactly_one() with a stdlib equivalent when there is one.
4//
5// See issue #48919 <https://github.com/rust-lang/rust/issues/48919>
6#![allow(unstable_name_collisions)]
7
8#[cfg(feature = "conflux")]
9pub(crate) mod msghandler;
10
11use std::pin::Pin;
12use std::sync::atomic::{self, AtomicU64};
13use std::sync::{Arc, Mutex};
14
15use futures::{FutureExt as _, StreamExt, select_biased};
16use itertools::Itertools;
17use itertools::structs::ExactlyOneError;
18use smallvec::{SmallVec, smallvec};
19use tor_rtcompat::{SleepProvider as _, SleepProviderExt as _};
20use tracing::{info, instrument, trace, warn};
21
22use tor_cell::relaycell::AnyRelayMsgOuter;
23use tor_error::{Bug, bad_api_usage, internal};
24use tor_linkspec::HasRelayIds as _;
25
26use crate::circuit::UniqId;
27use crate::circuit::circhop::SendRelayCell;
28use crate::client::circuit::TunnelMutableState;
29#[cfg(feature = "circ-padding")]
30use crate::client::circuit::padding::PaddingEvent;
31use crate::client::circuit::path::HopDetail;
32use crate::conflux::cmd_counts_towards_seqno;
33use crate::conflux::msghandler::{ConfluxStatus, RemoveLegReason};
34use crate::congestion::params::CongestionWindowParams;
35use crate::crypto::cell::HopNum;
36use crate::streammap;
37use crate::tunnel::TunnelId;
38use crate::util::err::ReactorError;
39use crate::util::poll_all::PollAll;
40use crate::util::tunnel_activity::TunnelActivity;
41
42use super::circuit::CircHop;
43use super::{Circuit, CircuitEvent};
44
45#[cfg(feature = "conflux")]
46use {
47    crate::conflux::msghandler::ConfluxMsgHandler,
48    msghandler::ClientConfluxMsgHandler,
49    tor_cell::relaycell::conflux::{V1DesiredUx, V1LinkPayload, V1Nonce},
50    tor_cell::relaycell::msg::{ConfluxLink, ConfluxSwitch},
51};
52
53/// The maximum number of conflux legs to store in the conflux set SmallVec.
54///
55/// Attempting to store more legs will cause the SmallVec to spill to the heap.
56///
57/// Note: this value was picked arbitrarily and may not be suitable.
58const MAX_CONFLUX_LEGS: usize = 16;
59
60/// The number of futures we add to the per-circuit [`PollAll`] future in
61/// [`ConfluxSet::next_circ_event`].
62///
63/// Used for the SmallVec size estimate;
64const NUM_CIRC_FUTURES: usize = 2;
65
66/// The expected number of circuit events to be returned from
67/// [`ConfluxSet::next_circ_event`]
68const CIRC_EVENT_COUNT: usize = MAX_CONFLUX_LEGS * NUM_CIRC_FUTURES;
69
70/// A set with one or more circuits.
71///
72/// ### Conflux set life cycle
73///
74/// Conflux sets are created by the reactor using [`ConfluxSet::new`].
75///
76/// Every `ConfluxSet` starts out as a single-path set consisting of a single 0-length circuit.
77///
78/// After constructing a `ConfluxSet`, the reactor will proceed to extend its (only) circuit.
79/// At this point, the `ConfluxSet` will be a single-path set with a single n-length circuit.
80///
81/// The reactor can then turn the `ConfluxSet` into a multi-path set
82/// (a multi-path set is a conflux set that contains more than 1 circuit).
83/// This is done using [`ConfluxSet::add_legs`], in response to a `CtrlMsg` sent
84/// by the reactor user (also referred to as the "conflux handshake initiator").
85/// After that, the conflux set is said to be a multi-path set with multiple N-length circuits.
86///
87/// Circuits can be removed from the set using [`ConfluxSet::remove`].
88///
89/// The lifetime of a `ConfluxSet` is tied to the lifetime of the reactor.
90/// When the reactor is dropped, its underlying `ConfluxSet` is dropped too.
91/// This can happen on an explicit shutdown request, or if a fatal error occurs.
92///
93/// Conversely, the `ConfluxSet` can also trigger a reactor shutdown.
94/// For example, if after being instructed to remove a circuit from the set
95/// using [`ConfluxSet::remove`], the set is completely depleted,
96/// the `ConfluxSet` will return a [`ReactorError::Shutdown`] error,
97/// which will cause the reactor to shut down.
98pub(super) struct ConfluxSet {
99    /// The unique identifier of the tunnel this conflux set belongs to.
100    ///
101    /// Used for setting the internal [`TunnelId`] of [`Circuit`]s
102    /// that gets used for logging purposes.
103    tunnel_id: TunnelId,
104    /// The circuits in this conflux set.
105    legs: SmallVec<[Circuit; MAX_CONFLUX_LEGS]>,
106    /// Tunnel state, shared with `ClientCirc`.
107    ///
108    /// Contains the [`MutableState`](super::MutableState) of each circuit in the set.
109    mutable: Arc<TunnelMutableState>,
110    /// The unique identifier of the primary leg
111    primary_id: UniqId,
112    /// The join point of the set, if this is a multi-path set.
113    ///
114    /// Initially the conflux set starts out as a single-path set with no join point.
115    /// When it is converted to a multipath set using [`add_legs`](Self::add_legs),
116    /// the join point is initialized to the last hop in the tunnel.
117    //
118    // TODO(#2017): for simplicity, we currently we force all legs to have the same length,
119    // to ensure the HopNum of the join point is the same for all of them.
120    //
121    // In the future we might want to relax this restriction.
122    join_point: Option<JoinPoint>,
123    /// The nonce associated with the circuits from this set.
124    #[cfg(feature = "conflux")]
125    nonce: V1Nonce,
126    /// The desired UX
127    #[cfg(feature = "conflux")]
128    desired_ux: V1DesiredUx,
129    /// The absolute sequence number of the last cell delivered to a stream.
130    ///
131    /// A clone of this is shared with each [`ConfluxMsgHandler`] created.
132    ///
133    /// When a message is received on a circuit leg, the `ConfluxMsgHandler`
134    /// of the leg compares the (leg-local) sequence number of the message
135    /// with this sequence number to determine whether the message is in-order.
136    ///
137    /// If the message is in-order, the `ConfluxMsgHandler` instructs the circuit
138    /// to deliver it to its corresponding stream.
139    ///
140    /// If the message is out-of-order, the `ConfluxMsgHandler` instructs the circuit
141    /// to instruct the reactor to buffer the message.
142    last_seq_delivered: Arc<AtomicU64>,
143    /// Whether we have selected our initial primary leg,
144    /// if this is a multipath conflux set.
145    selected_init_primary: bool,
146}
147
148/// The conflux join point.
149#[derive(Clone, derive_more::Debug)]
150struct JoinPoint {
151    /// The hop number.
152    hop: HopNum,
153    /// The [`HopDetail`] of the hop.
154    detail: HopDetail,
155    /// The stream map of the joint point, shared with each circuit leg.
156    #[debug(skip)]
157    streams: Arc<Mutex<streammap::StreamMap>>,
158}
159
160impl ConfluxSet {
161    /// Create a new conflux set, consisting of a single leg.
162    ///
163    /// Returns the newly created set and a reference to its [`TunnelMutableState`].
164    pub(super) fn new(
165        tunnel_id: TunnelId,
166        circuit_leg: Circuit,
167    ) -> (Self, Arc<TunnelMutableState>) {
168        let primary_id = circuit_leg.unique_id();
169        let circ_mutable = Arc::clone(circuit_leg.mutable());
170        let legs = smallvec![circuit_leg];
171        // Note: the join point is only set for multi-path tunnels
172        let join_point = None;
173
174        // TODO(#2035): read this from the consensus/config.
175        #[cfg(feature = "conflux")]
176        let desired_ux = V1DesiredUx::NO_OPINION;
177
178        let mutable = Arc::new(TunnelMutableState::default());
179        mutable.insert(primary_id, circ_mutable);
180
181        let set = Self {
182            tunnel_id,
183            legs,
184            primary_id,
185            join_point,
186            mutable: mutable.clone(),
187            #[cfg(feature = "conflux")]
188            nonce: V1Nonce::new(&mut rand::rng()),
189            #[cfg(feature = "conflux")]
190            desired_ux,
191            last_seq_delivered: Arc::new(AtomicU64::new(0)),
192            selected_init_primary: false,
193        };
194
195        (set, mutable)
196    }
197
198    /// Remove and return the only leg of this conflux set.
199    ///
200    /// Returns an error if there is more than one leg in the set,
201    /// or if called before any circuit legs are available.
202    ///
203    /// Calling this function will empty the [`ConfluxSet`].
204    pub(super) fn take_single_leg(&mut self) -> Result<Circuit, Bug> {
205        let circ = self
206            .legs
207            .iter()
208            .exactly_one()
209            .map_err(NotSingleLegError::from)?;
210        let circ_id = circ.unique_id();
211
212        debug_assert!(circ_id == self.primary_id);
213
214        self.remove_unchecked(circ_id)
215    }
216
217    /// Return a reference to the only leg of this conflux set,
218    /// along with the leg's ID.
219    ///
220    /// Returns an error if there is more than one leg in the set,
221    /// or if called before any circuit legs are available.
222    pub(super) fn single_leg(&self) -> Result<&Circuit, NotSingleLegError> {
223        Ok(self.legs.iter().exactly_one()?)
224    }
225
226    /// Return a mutable reference to the only leg of this conflux set,
227    /// along with the leg's ID.
228    ///
229    /// Returns an error if there is more than one leg in the set,
230    /// or if called before any circuit legs are available.
231    pub(super) fn single_leg_mut(&mut self) -> Result<&mut Circuit, NotSingleLegError> {
232        Ok(self.legs.iter_mut().exactly_one()?)
233    }
234
235    /// Return the primary leg of this conflux set.
236    ///
237    /// Returns an error if called before any circuit legs are available.
238    pub(super) fn primary_leg_mut(&mut self) -> Result<&mut Circuit, Bug> {
239        #[cfg(not(feature = "conflux"))]
240        if self.legs.len() > 1 {
241            return Err(internal!(
242                "got multipath tunnel, but conflux feature is disabled?!"
243            ));
244        }
245
246        if self.legs.is_empty() {
247            Err(bad_api_usage!(
248                "tried to get circuit leg before creating it?!"
249            ))
250        } else {
251            let circ = self
252                .leg_mut(self.primary_id)
253                .ok_or_else(|| internal!("conflux set is empty?!"))?;
254
255            Ok(circ)
256        }
257    }
258
259    /// Return a reference to the leg of this conflux set with the given id.
260    pub(super) fn leg(&self, leg_id: UniqId) -> Option<&Circuit> {
261        self.legs.iter().find(|circ| circ.unique_id() == leg_id)
262    }
263
264    /// Return a mutable reference to the leg of this conflux set with the given id.
265    pub(super) fn leg_mut(&mut self, leg_id: UniqId) -> Option<&mut Circuit> {
266        self.legs.iter_mut().find(|circ| circ.unique_id() == leg_id)
267    }
268
269    /// Return the number of legs in this conflux set.
270    pub(super) fn len(&self) -> usize {
271        self.legs.len()
272    }
273
274    /// Return whether this conflux set is empty.
275    pub(super) fn is_empty(&self) -> bool {
276        self.legs.len() == 0
277    }
278
279    /// Remove the specified leg from this conflux set.
280    ///
281    /// Returns an error if the given leg doesn't exist in the set.
282    ///
283    /// Returns an error instructing the reactor to perform a clean shutdown
284    /// ([`ReactorError::Shutdown`]), tearing down the entire [`ConfluxSet`], if
285    ///
286    ///   * the set is depleted (empty) after removing the specified leg
287    ///   * `leg` is currently the sending (primary) leg of this set
288    ///   * the closed leg had the highest non-zero last_seq_recv/sent
289    ///   * the closed leg had some in-progress data (inflight > cc_sendme_inc)
290    ///
291    /// We do not yet support resumption. See [2.4.3. Closing circuits] in prop329.
292    ///
293    /// [2.4.3. Closing circuits]: https://spec.torproject.org/proposals/329-traffic-splitting.html#243-closing-circuits
294    #[instrument(level = "trace", skip_all)]
295    pub(super) fn remove(&mut self, leg: UniqId) -> Result<Circuit, ReactorError> {
296        let circ = self.remove_unchecked(leg)?;
297
298        tracing::trace!(
299            circ_uniq_id = %circ.unique_id(),
300            forward_circ_id = %circ.circ_id(),
301            "Circuit removed from conflux set"
302        );
303
304        self.mutable.remove(circ.unique_id());
305
306        if self.legs.is_empty() {
307            // TODO: log the tunnel ID
308            tracing::debug!("Conflux set is now empty, tunnel reactor shutting down");
309
310            // The last circuit in the set has just died, so the reactor should exit.
311            return Err(ReactorError::Shutdown);
312        }
313
314        if leg == self.primary_id {
315            // We have just removed our sending leg,
316            // so it's time to close the entire conflux set.
317            return Err(ReactorError::Shutdown);
318        }
319
320        cfg_if::cfg_if! {
321            if #[cfg(feature = "conflux")] {
322                self.remove_conflux(circ)
323            } else {
324                // Conflux is disabled, so we can't possibly continue running if the only
325                // leg in the tunnel is gone.
326                //
327                // Technically this should be unreachable (because of the is_empty()
328                // check above)
329                Err(internal!("Multiple legs in single-path tunnel?!").into())
330            }
331        }
332    }
333
334    /// Handle the removal of a circuit,
335    /// returning an error if the reactor needs to shut down.
336    #[cfg(feature = "conflux")]
337    fn remove_conflux(&self, circ: Circuit) -> Result<Circuit, ReactorError> {
338        let Some(status) = circ.conflux_status() else {
339            return Err(internal!("Found non-conflux circuit in conflux set?!").into());
340        };
341
342        // TODO(conflux): should the circmgr be notified about the leg removal?
343        //
344        // "For circuits that are unlinked, the origin SHOULD immediately relaunch a new leg when it
345        // is closed, subject to the limits in [SIDE_CHANNELS]."
346
347        // If we've reached this point and the conflux set is non-empty,
348        // it means it's a multi-path set.
349        //
350        // Time to check if we need to tear down the entire set.
351        match status {
352            ConfluxStatus::Unlinked => {
353                // This circuit hasn't yet begun the conflux handshake,
354                // so we can safely remove it from the set
355                Ok(circ)
356            }
357            ConfluxStatus::Pending | ConfluxStatus::Linked => {
358                let (circ_last_seq_recv, circ_last_seq_sent) =
359                    (|| Ok::<_, ReactorError>((circ.last_seq_recv()?, circ.last_seq_sent()?)))()?;
360
361                // If the closed leg had the highest non-zero last_seq_recv/sent, close the set
362                if let Some(max_last_seq_recv) = self.max_last_seq_recv() {
363                    if circ_last_seq_recv > max_last_seq_recv {
364                        return Err(ReactorError::Shutdown);
365                    }
366                }
367
368                if let Some(max_last_seq_sent) = self.max_last_seq_sent() {
369                    if circ_last_seq_sent > max_last_seq_sent {
370                        return Err(ReactorError::Shutdown);
371                    }
372                }
373
374                let hop = self.join_point_hop(&circ)?;
375
376                let (inflight, cwnd) = (|| {
377                    let ccontrol = hop.ccontrol();
378                    let inflight = ccontrol.inflight()?;
379                    let cwnd = ccontrol.cwnd()?;
380
381                    Some((inflight, cwnd))
382                })()
383                .ok_or_else(|| {
384                    internal!("Congestion control algorithm doesn't track inflight cells or cwnd?!")
385                })?;
386
387                // If data is in progress on the leg (inflight > cc_sendme_inc),
388                // then all legs must be closed
389                if inflight >= u32::from(cwnd.params().sendme_inc()) {
390                    return Err(ReactorError::Shutdown);
391                }
392
393                Ok(circ)
394            }
395        }
396    }
397
398    /// Return the maximum relative last_seq_recv across all circuits.
399    #[cfg(feature = "conflux")]
400    fn max_last_seq_recv(&self) -> Option<u64> {
401        self.legs
402            .iter()
403            .filter_map(|leg| leg.last_seq_recv().ok())
404            .max()
405    }
406
407    /// Return the maximum relative last_seq_sent across all circuits.
408    #[cfg(feature = "conflux")]
409    fn max_last_seq_sent(&self) -> Option<u64> {
410        self.legs
411            .iter()
412            .filter_map(|leg| leg.last_seq_sent().ok())
413            .max()
414    }
415
416    /// Get the [`CircHop`] of the join point on the specified `circ`,
417    /// returning an error if this is a single path conflux set.
418    fn join_point_hop<'c>(&self, circ: &'c Circuit) -> Result<&'c CircHop, Bug> {
419        let Some(join_point) = self.join_point.as_ref().map(|p| p.hop) else {
420            return Err(internal!("No join point on conflux tunnel?!"));
421        };
422
423        circ.hop(join_point)
424            .ok_or_else(|| internal!("Conflux join point disappeared?!"))
425    }
426
427    /// Return an iterator of all circuits in the conflux set.
428    fn circuits(&self) -> impl Iterator<Item = &Circuit> {
429        self.legs.iter()
430    }
431
432    /// Return the most active [`TunnelActivity`] for any leg of this `ConfluxSet`.
433    pub(super) fn tunnel_activity(&self) -> TunnelActivity {
434        self.circuits()
435            .map(|c| c.hops.tunnel_activity())
436            .max()
437            .unwrap_or_else(TunnelActivity::never_used)
438    }
439
440    /// Add legs to the this conflux set.
441    ///
442    /// Returns an error if any of the legs is invalid.
443    ///
444    /// A leg is considered valid if
445    ///
446    ///   * the circuit has the same length as all the other circuits in the set
447    ///   * its last hop is equal to the designated join point
448    ///   * the circuit has no streams attached to any of its hops
449    ///   * the circuit is not already part of a conflux set
450    ///
451    /// Note: the circuits will not begin linking until
452    /// [`link_circuits`](Self::link_circuits) is called.
453    ///
454    /// IMPORTANT: this function does not prevent the construction of conflux sets
455    /// where the circuit legs share guard or middle relays. It is the responsibility
456    /// of the caller to enforce the following invariant from prop354:
457    ///
458    /// "If building a conflux leg: Reject any circuits that have the same Guard as the other conflux
459    /// "leg(s) in the current conflux set, EXCEPT when one of the primary Guards is also the chosen
460    /// "Exit of this conflux set (in which case, re-use the non-Exit Guard)."
461    ///
462    /// This is because at this level we don't actually know which relays are the guards,
463    /// so we can't know if the join point happens to be one of the Guard + Exit relays.
464    #[cfg(feature = "conflux")]
465    pub(super) fn add_legs(
466        &mut self,
467        legs: Vec<Circuit>,
468        runtime: &tor_rtcompat::DynTimeProvider,
469    ) -> Result<(), Bug> {
470        if legs.is_empty() {
471            return Err(bad_api_usage!("asked to add empty leg list to conflux set"));
472        }
473
474        let join_point = match self.join_point.take() {
475            Some(p) => {
476                // Preserve the existing join point, if there is one.
477                p
478            }
479            None => {
480                let (hop, detail, streams) = (|| {
481                    let first_leg = self.circuits().next()?;
482                    let first_leg_path = first_leg.path();
483                    let all_hops = first_leg_path.all_hops();
484                    let hop_num = first_leg.last_hop_num()?;
485                    let detail = all_hops.last()?;
486                    let hop = first_leg.hop(hop_num)?;
487                    let streams = Arc::clone(hop.stream_map());
488                    Some((hop_num, detail.clone(), streams))
489                })()
490                .ok_or_else(|| bad_api_usage!("asked to join circuit with no hops"))?;
491
492                JoinPoint {
493                    hop,
494                    detail,
495                    streams,
496                }
497            }
498        };
499
500        // Check two HopDetails for equality.
501        //
502        // Returns an error if one of the hops is virtual.
503        let hops_eq = |h1: &HopDetail, h2: &HopDetail| {
504            match (h1, h2) {
505                (HopDetail::Relay(t1), HopDetail::Relay(t2)) => Ok(t1.same_relay_ids(t2)),
506                #[cfg(feature = "hs-common")]
507                (HopDetail::Virtual, HopDetail::Virtual) => {
508                    // TODO(#2016): support onion service conflux
509                    Err(internal!("onion service conflux not supported"))
510                }
511                _ => Ok(false),
512            }
513        };
514
515        // A leg is considered valid if
516        //
517        //   * the circuit has the expected length
518        //     (the length of the first circuit we added to the set)
519        //   * its last hop is equal to the designated join point
520        //     (the last hop of the first circuit we added)
521        //   * the circuit has no streams attached to any of its hops
522        //   * the circuit is not already part of a conflux tunnel
523        //
524        // Returns an error if any hops are virtual.
525        let leg_is_valid = |leg: &Circuit| -> Result<bool, Bug> {
526            use crate::ccparams::Algorithm;
527
528            let path = leg.path();
529            let Some(last_hop) = path.all_hops().last() else {
530                // A circuit with no hops is invalid
531                return Ok(false);
532            };
533
534            // TODO: this sort of duplicates the check above.
535            // The difference is that above we read the hop detail
536            // information from the circuit Path, whereas here we get
537            // the actual last CircHop of the circuit.
538            let Some(last_hop_num) = leg.last_hop_num() else {
539                // A circuit with no hops is invalid
540                return Ok(false);
541            };
542
543            let circhop = leg
544                .hop(last_hop_num)
545                .ok_or_else(|| internal!("hop disappeared?!"))?;
546
547            // Ensure we negotiated a suitable cc algorithm
548            let is_cc_suitable = match circhop.ccontrol().algorithm() {
549                Algorithm::FixedWindow(_) => false,
550                Algorithm::Vegas(_) => true,
551            };
552
553            if !is_cc_suitable {
554                return Ok(false);
555            }
556
557            Ok(last_hop_num == join_point.hop
558                && hops_eq(last_hop, &join_point.detail)?
559                && !leg.has_streams()
560                && leg.conflux_status().is_none())
561        };
562
563        for leg in &legs {
564            if !leg_is_valid(leg)? {
565                return Err(bad_api_usage!("one more conflux circuits are invalid"));
566            }
567        }
568
569        // Select a join point, or put the existing one back into self.
570        self.join_point = Some(join_point.clone());
571
572        // The legs are valid, so add them to the set.
573        for circ in legs {
574            let mutable = Arc::clone(circ.mutable());
575            let unique_id = circ.unique_id();
576            self.legs.push(circ);
577            // Merge the mutable state of the circuit into our tunnel state.
578            self.mutable.insert(unique_id, mutable);
579        }
580
581        let cwnd_params = self.cwnd_params()?;
582        for circ in self.legs.iter_mut() {
583            // The circuits that have a None status don't know they're part of
584            // a multi-path tunnel yet. They need to be initialized with a
585            // conflux message handler, and have their join point fixed up
586            // to share a stream map with the join point on all the other circuits.
587            if circ.conflux_status().is_none() {
588                let handler = Box::new(ClientConfluxMsgHandler::new(
589                    join_point.hop,
590                    self.nonce,
591                    Arc::clone(&self.last_seq_delivered),
592                    cwnd_params,
593                    runtime.clone(),
594                ));
595                let conflux_handler =
596                    ConfluxMsgHandler::new(handler, Arc::clone(&self.last_seq_delivered));
597
598                circ.add_to_conflux_tunnel(self.tunnel_id, conflux_handler);
599
600                // Ensure the stream map of the last hop is shared by all the legs
601                let last_hop = circ
602                    .hop_mut(join_point.hop)
603                    .ok_or_else(|| bad_api_usage!("asked to join circuit with no hops"))?;
604                last_hop.set_stream_map(Arc::clone(&join_point.streams))?;
605            }
606        }
607
608        Ok(())
609    }
610
611    /// Get the [`CongestionWindowParams`] of the join point
612    /// on the first leg.
613    ///
614    /// Returns an error if the congestion control algorithm
615    /// doesn't have a congestion control window object,
616    /// or if the conflux set is empty, or the joint point hop
617    /// does not exist.
618    ///
619    // TODO: this function is a bit of a hack. In reality, we only
620    // need the cc_cwnd_init parameter (for SWITCH seqno validation).
621    // The fact that we obtain it from the cc params of the join point
622    // is an implementation detail (it's a workaround for the fact that
623    // at this point, these params can only obtained from a CircHop)
624    #[cfg(feature = "conflux")]
625    fn cwnd_params(&self) -> Result<CongestionWindowParams, Bug> {
626        let primary_leg = self
627            .leg(self.primary_id)
628            .ok_or_else(|| internal!("no primary leg?!"))?;
629        let join_point = self.join_point_hop(primary_leg)?;
630        let ccontrol = join_point.ccontrol();
631        let cwnd = ccontrol
632            .cwnd()
633            .ok_or_else(|| internal!("congestion control algorithm does not track the cwnd?!"))?;
634
635        Ok(*cwnd.params())
636    }
637
638    /// Try to update the primary leg based on the configured desired UX,
639    /// if needed.
640    ///
641    /// Returns the SWITCH cell to send on the primary leg,
642    /// if we switched primary leg.
643    #[cfg(feature = "conflux")]
644    pub(super) fn maybe_update_primary_leg(&mut self) -> crate::Result<Option<SendRelayCell>> {
645        use tor_error::into_internal;
646
647        let Some(join_point) = self.join_point.as_ref() else {
648            // Return early if this is not a multi-path tunnel
649            return Ok(None);
650        };
651
652        let join_point = join_point.hop;
653
654        if !self.should_update_primary_leg() {
655            // Nothing to do
656            return Ok(None);
657        }
658
659        let Some(new_primary_id) = self.select_primary_leg()? else {
660            // None of the legs satisfy our UX requirements, continue using the existing one.
661            return Ok(None);
662        };
663
664        // Check that the newly selected leg is actually different from the previous
665        if self.primary_id == new_primary_id {
666            // The primary leg stays the same, nothing to do.
667            return Ok(None);
668        }
669
670        let prev_last_seq_sent = self.primary_leg_mut()?.last_seq_sent()?;
671        self.primary_id = new_primary_id;
672        let new_last_seq_sent = self.primary_leg_mut()?.last_seq_sent()?;
673
674        // If this fails, it means we haven't updated our primary leg in a very long time.
675        //
676        // TODO(#2036): there are currently no safeguards to prevent us from staying
677        // on the same leg for "too long". Perhaps we should design should_update_primary_leg()
678        // such that it forces us to switch legs periodically, to prevent the seqno delta from
679        // getting too big?
680        let seqno_delta = u32::try_from(prev_last_seq_sent - new_last_seq_sent).map_err(
681            into_internal!("Seqno delta for switch does not fit in u32?!"),
682        )?;
683
684        // We need to carry the last_seq_sent over to the next leg
685        // (the next cell sent will have seqno = prev_last_seq_sent + 1)
686        self.primary_leg_mut()?
687            .set_last_seq_sent(prev_last_seq_sent)?;
688
689        let switch = ConfluxSwitch::new(seqno_delta);
690        let cell = AnyRelayMsgOuter::new(None, switch.into());
691        Ok(Some(SendRelayCell {
692            hop: Some(join_point),
693            early: false,
694            cell,
695        }))
696    }
697
698    /// Whether it's time to select a new primary leg.
699    #[cfg(feature = "conflux")]
700    fn should_update_primary_leg(&mut self) -> bool {
701        if !self.selected_init_primary {
702            self.maybe_select_init_primary();
703            return false;
704        }
705
706        // If we don't have at least 2 legs,
707        // we can't switch our primary leg.
708        if self.legs.len() < 2 {
709            return false;
710        }
711
712        // TODO(conflux-tuning): if it turns out we switch legs too frequently,
713        // we might want to implement some sort of rate-limiting here
714        // (see c-tor's conflux_can_switch).
715
716        true
717    }
718
719    /// Return the best leg according to the configured desired UX.
720    ///
721    /// Returns `None` if no suitable leg was found.
722    #[cfg(feature = "conflux")]
723    fn select_primary_leg(&self) -> Result<Option<UniqId>, Bug> {
724        match self.desired_ux {
725            V1DesiredUx::NO_OPINION | V1DesiredUx::MIN_LATENCY => {
726                self.select_primary_leg_min_rtt(false)
727            }
728            V1DesiredUx::HIGH_THROUGHPUT => self.select_primary_leg_min_rtt(true),
729            V1DesiredUx::LOW_MEM_LATENCY | V1DesiredUx::LOW_MEM_THROUGHPUT => {
730                // TODO(conflux-tuning): add support for low-memory algorithms
731                self.select_primary_leg_min_rtt(false)
732            }
733            _ => {
734                // Default to MIN_RTT if we don't recognize the desired UX value
735                warn!(
736                    tunnel_id = %self.tunnel_id,
737                    "Ignoring unrecognized conflux desired UX {}, using MIN_LATENCY",
738                    self.desired_ux
739                );
740                self.select_primary_leg_min_rtt(false)
741            }
742        }
743    }
744
745    /// Try to choose an initial primary leg, if we have an initial RTT measurement
746    /// for at least one of the legs.
747    #[cfg(feature = "conflux")]
748    fn maybe_select_init_primary(&mut self) {
749        let best = self
750            .legs
751            .iter()
752            .filter_map(|leg| leg.init_rtt().map(|rtt| (leg, rtt)))
753            .min_by_key(|(_leg, rtt)| *rtt)
754            .map(|(leg, _rtt)| leg.unique_id());
755
756        if let Some(best) = best {
757            self.primary_id = best;
758            self.selected_init_primary = true;
759        }
760    }
761
762    /// Return the leg with the best (lowest) RTT.
763    ///
764    /// If `check_can_send` is true, selects the lowest RTT leg that is ready to send.
765    ///
766    /// Returns `None` if no suitable leg was found.
767    #[cfg(feature = "conflux")]
768    fn select_primary_leg_min_rtt(&self, check_can_send: bool) -> Result<Option<UniqId>, Bug> {
769        let mut best: Option<(UniqId, std::time::Duration)> = None;
770
771        for circ in self.legs.iter() {
772            let leg_id = circ.unique_id();
773            let join_point = self.join_point_hop(circ)?;
774            let ccontrol = join_point.ccontrol();
775
776            if check_can_send && !ccontrol.can_send() {
777                continue;
778            }
779
780            let Some(ewma_rtt) = ccontrol.rtt().ewma_rtt().or_else(|| circ.init_rtt()) else {
781                return Err(internal!(
782                    "attempted to select primary leg before handshake completed?!"
783                ));
784            };
785
786            best = Some(match best.take() {
787                Some(best_so_far) if best_so_far.1 <= ewma_rtt => best_so_far,
788                None | Some(_) => (leg_id, ewma_rtt),
789            });
790        }
791
792        Ok(best.map(|(leg_id, _)| leg_id))
793    }
794
795    /// Returns `true` if our conflux join point is blocked on congestion control
796    /// on the specified `circuit`.
797    ///
798    /// Returns `false` if the join point is not blocked on cc,
799    /// or if this is a single-path set.
800    ///
801    /// Returns an error if this is a multipath tunnel,
802    /// but the joint point hop doesn't exist on the specified circuit.
803    #[cfg(feature = "conflux")]
804    fn is_join_point_blocked_on_cc(join_hop: HopNum, circuit: &Circuit) -> Result<bool, Bug> {
805        let join_circhop = circuit.hop(join_hop).ok_or_else(|| {
806            internal!(
807                "Join point hop {} not found on circuit {}?!",
808                join_hop.display(),
809                circuit.unique_id(),
810            )
811        })?;
812
813        Ok(!join_circhop.ccontrol().can_send())
814    }
815
816    /// Returns whether [`next_circ_event`](Self::next_circ_event)
817    /// should avoid polling the join point streams entirely.
818    #[cfg(feature = "conflux")]
819    fn should_skip_join_point(&self) -> Result<bool, Bug> {
820        let Some(primary_join_point) = self.primary_join_point() else {
821            // Single-path, there is no join point
822            return Ok(false);
823        };
824
825        let join_hop = primary_join_point.1;
826        let primary_blocked_on_cc = {
827            let primary = self
828                .leg(self.primary_id)
829                .ok_or_else(|| internal!("primary leg disappeared?!"))?;
830            Self::is_join_point_blocked_on_cc(join_hop, primary)?
831        };
832
833        if !primary_blocked_on_cc {
834            // Easy, we can just carry on
835            return Ok(false);
836        }
837
838        // Now, if the primary *is* blocked on cc, we may still be able to poll
839        // the join point streams (if we're using the right desired UX)
840        let should_skip = if self.desired_ux != V1DesiredUx::HIGH_THROUGHPUT {
841            // The primary leg is blocked on cc, and we can't switch because we're
842            // not using the high throughput algorithm, so we must stop reading
843            // the join point streams.
844            //
845            // Note: if the selected algorithm is HIGH_THROUGHPUT,
846            // it's okay to continue reading from the edge connection,
847            // because maybe_update_primary_leg() will select a new,
848            // non-blocked primary leg, just before sending.
849            trace!(
850                tunnel_id = %self.tunnel_id,
851                join_point = ?primary_join_point,
852                reason = "sending leg blocked on congestion control",
853                "Pausing join point stream reads"
854            );
855
856            true
857        } else {
858            // Ah-ha, the desired UX is HIGH_THROUGHPUT, which means we can switch
859            // to an unblocked leg before sending any cells over the join point,
860            // as long as there are some unblocked legs.
861
862            // TODO: figure out how to rewrite this with an idiomatic iterator combinator
863            let mut all_blocked_on_cc = true;
864            for leg in &self.legs {
865                all_blocked_on_cc = Self::is_join_point_blocked_on_cc(join_hop, leg)?;
866                if !all_blocked_on_cc {
867                    break;
868                }
869            }
870
871            if all_blocked_on_cc {
872                // All legs are blocked on cc, so we must stop reading from
873                // the join point streams for now.
874                trace!(
875                    tunnel_id = %self.tunnel_id,
876                    join_point = ?primary_join_point,
877                    reason = "all legs blocked on congestion control",
878                    "Pausing join point stream reads"
879                );
880
881                true
882            } else {
883                // At least one leg is not blocked, so we can continue reading
884                // from the join point streams
885                false
886            }
887        };
888
889        Ok(should_skip)
890    }
891
892    /// Returns the next ready [`CircuitEvent`],
893    /// obtained from processing the incoming/outgoing messages on all the circuits in this set.
894    ///
895    /// Will return an error if there are no circuits in this set,
896    /// or other internal errors occur.
897    ///
898    /// This is cancellation-safe.
899    #[allow(clippy::unnecessary_wraps)] // Can return Err if conflux is enabled
900    #[instrument(level = "trace", skip_all)]
901    pub(super) async fn next_circ_event(
902        &mut self,
903        runtime: &tor_rtcompat::DynTimeProvider,
904    ) -> Result<SmallVec<[CircuitEvent; CIRC_EVENT_COUNT]>, crate::Error> {
905        // Avoid polling the streams on the join point if our primary
906        // leg is blocked on cc
907        cfg_if::cfg_if! {
908            if #[cfg(feature = "conflux")] {
909                let mut should_poll_join_point = !self.should_skip_join_point()?;
910            } else {
911                let mut should_poll_join_point = true;
912            }
913        };
914        let join_point = self.primary_join_point().map(|join_point| join_point.1);
915
916        // Each circuit leg has a PollAll future (see poll_all_circ below)
917        // that drives two futures: one that reads from input channel,
918        // and another drives the application streams.
919        //
920        // *This* PollAll drives the PollAll futures of all circuit legs in lockstep,
921        // ensuring they all get a chance to make some progress on every reactor iteration.
922        //
923        // IMPORTANT: if you want to push additional futures into this,
924        // bear in mind that the ordering matters!
925        // If multiple futures resolve at the same time, their results will be processed
926        // in the order their corresponding futures were inserted into `PollAll`.
927        // So if futures A and B resolve at the same time, and future A was pushed
928        // into `PollAll` before future B, the result of future A will come
929        // before future B's result in the result list returned by poll_all.await.
930        //
931        // This means that the events corresponding to the first circuit in the tunnel
932        // will be executed first, followed by the events issued by the next circuit,
933        // and so on.
934        //
935        let mut poll_all =
936            PollAll::<MAX_CONFLUX_LEGS, SmallVec<[CircuitEvent; NUM_CIRC_FUTURES]>>::new();
937
938        for leg in &mut self.legs {
939            let unique_id = leg.unique_id();
940            let circ_id = leg.circ_id();
941            let tunnel_id = self.tunnel_id;
942            let runtime = runtime.clone();
943
944            // Garbage-collect all halfstreams that have expired.
945            //
946            // Note: this will iterate over the closed streams of all hops.
947            // If we think this will cause perf issues, one idea would be to make
948            // StreamMap::closed_streams into a min-heap, and add a branch to the
949            // select_biased! below to sleep until the first expiry is due
950            // (but my gut feeling is that iterating is cheaper)
951            leg.remove_expired_halfstreams(runtime.now());
952
953            // The client SHOULD abandon and close circuit if the LINKED message takes too long to
954            // arrive. This timeout MUST be no larger than the normal SOCKS/stream timeout in use for
955            // RELAY_BEGIN, but MAY be the Circuit Build Timeout value, instead. (The C-Tor
956            // implementation currently uses Circuit Build Timeout).
957            let conflux_hs_timeout = leg.conflux_hs_timeout();
958
959            let mut poll_all_circ = PollAll::<NUM_CIRC_FUTURES, CircuitEvent>::new();
960
961            let input = leg.input.next().map(move |res| match res {
962                Some(msg) => match msg.try_into() {
963                    Ok(cell) => CircuitEvent::HandleCell {
964                        leg: unique_id,
965                        cell,
966                    },
967                    // A message outside our restricted set is either a fatal internal error or
968                    // a protocol violation somehow so shutdown.
969                    //
970                    // TODO(relay): We have this spec ticket open about this behavior:
971                    // https://gitlab.torproject.org/tpo/core/torspec/-/issues/385. It is plausible
972                    // that we decide to either keep this circuit close behavior or close the
973                    // entire channel in this case. Resolution of the above ticket needs to fix
974                    // this part.
975                    Err(e) => CircuitEvent::ProtoViolation { err: e },
976                },
977                None => CircuitEvent::RemoveLeg {
978                    leg: unique_id,
979                    reason: RemoveLegReason::ChannelClosed,
980                },
981            });
982            poll_all_circ.push(input);
983
984            // This future resolves when the chan_sender sink (i.e. the outgoing TCP connection)
985            // becomes ready. We need it inside the next_ready_stream future below,
986            // to prevent reading from the application streams before we are ready to send.
987            let chan_ready_fut = futures::future::poll_fn(|cx| {
988                use futures::Sink as _;
989
990                // Ensure the chan sender sink is ready before polling the ready streams.
991                Pin::new(&mut leg.chan_sender).poll_ready(cx)
992            });
993
994            let exclude_hop = if should_poll_join_point {
995                // Avoid polling the join point more than once per reactor loop.
996                should_poll_join_point = false;
997                None
998            } else {
999                join_point
1000            };
1001
1002            let mut ready_streams = leg.hops.ready_streams_iterator(exclude_hop);
1003            let next_ready_stream = async move {
1004                // Avoid polling the application streams if the outgoing sink is blocked
1005                let _ = chan_ready_fut.await;
1006
1007                match ready_streams.next().await {
1008                    Some(x) => x,
1009                    None => {
1010                        info!(
1011                            circ_uniq_id = %unique_id,
1012                            forward_circ_id = %circ_id,
1013                            "no ready streams (maybe blocked on cc?)"
1014                        );
1015                        // There are no ready streams (for example, they may all be
1016                        // blocked due to congestion control), so there is nothing
1017                        // to do.
1018                        // We await an infinitely pending future so that we don't
1019                        // immediately return a `None` in the `select_biased!` below.
1020                        // We'd rather wait on `input.next()` than immediately return with
1021                        // no `CircuitEvent`, which could put the reactor into a spin loop.
1022                        let () = std::future::pending().await;
1023                        unreachable!();
1024                    }
1025                }
1026            };
1027
1028            poll_all_circ.push(next_ready_stream.map(move |cmd| CircuitEvent::RunCmd {
1029                leg: unique_id,
1030                cmd,
1031            }));
1032
1033            let mut next_padding_event_fut = leg.padding_event_stream.next();
1034
1035            // This selects between 3 events that cannot be handled concurrently.
1036            //
1037            // If the conflux handshake times out, we need to remove the circuit leg
1038            // (any pending padding events or application stream data should be discarded;
1039            // in fact, there shouldn't even be any open streams on circuits that are
1040            // in the conflux handshake phase).
1041            //
1042            // If there's a padding event, we need to handle it immediately,
1043            // because it might tell us to start blocking the chan_sender sink,
1044            // which, in turn, means we need to stop trying to read from the application streams.
1045            poll_all.push(
1046                async move {
1047                    let conflux_hs_timeout = if let Some(timeout) = conflux_hs_timeout {
1048                        // TODO: ask Diziet if we can have a sleep_until_instant() function
1049                        Box::pin(runtime.sleep_until_wallclock(timeout))
1050                            as Pin<Box<dyn Future<Output = ()> + Send>>
1051                    } else {
1052                        Box::pin(std::future::pending())
1053                    };
1054                    select_biased! {
1055                        () = conflux_hs_timeout.fuse() => {
1056                            warn!(
1057                                tunnel_id = %tunnel_id,
1058                                circ_uniq_id = %unique_id,
1059                                forward_circ_id = %circ_id,
1060                                "Conflux handshake timed out on circuit"
1061                            );
1062
1063                            // Conflux handshake has timed out, time to remove this circuit leg,
1064                            // and notify the handshake initiator.
1065                            smallvec![CircuitEvent::RemoveLeg {
1066                                leg: unique_id,
1067                                reason: RemoveLegReason::ConfluxHandshakeTimeout,
1068                            }]
1069                        }
1070                        padding_event = next_padding_event_fut => {
1071                            smallvec![CircuitEvent::PaddingAction {
1072                                leg: unique_id,
1073                                padding_event:
1074                                    padding_event.expect("PaddingEventStream, surprisingly, was terminated!"),
1075                            }]
1076                        }
1077                        ret = poll_all_circ.fuse() => ret,
1078                    }
1079                }
1080            );
1081        }
1082
1083        // Flatten the nested SmallVecs to simplify the calling code
1084        // (which will handle all the returned events sequentially).
1085        Ok(poll_all.await.into_iter().flatten().collect())
1086    }
1087
1088    /// The join point on the current primary leg.
1089    pub(super) fn primary_join_point(&self) -> Option<(UniqId, HopNum)> {
1090        self.join_point
1091            .as_ref()
1092            .map(|join_point| (self.primary_id, join_point.hop))
1093    }
1094
1095    /// Does congestion control use stream SENDMEs for the given hop?
1096    ///
1097    /// Returns `None` if either the `leg` or `hop` don't exist.
1098    pub(super) fn uses_stream_sendme(&self, leg: UniqId, hop: HopNum) -> Option<bool> {
1099        self.leg(leg)?.uses_stream_sendme(hop)
1100    }
1101
1102    /// Encode `msg`, encrypt it, and send it to the 'hop'th hop.
1103    ///
1104    /// See [`Circuit::send_relay_cell`].
1105    #[instrument(level = "trace", skip_all)]
1106    pub(super) async fn send_relay_cell_on_leg(
1107        &mut self,
1108        msg: SendRelayCell,
1109        leg: Option<UniqId>,
1110    ) -> crate::Result<()> {
1111        let conflux_join_point = self.join_point.as_ref().map(|join_point| join_point.hop);
1112        let leg = if let Some(join_point) = conflux_join_point {
1113            let hop = msg.hop.expect("missing hop in client SendRelayCell?!");
1114            // Conflux circuits always send multiplexed relay commands to
1115            // to the last hop (the join point).
1116            if cmd_counts_towards_seqno(msg.cell.cmd()) {
1117                if hop != join_point {
1118                    // For leaky pipe, we must continue using the original leg
1119                    leg
1120                } else {
1121                    let old_primary_leg = self.primary_id;
1122                    // Check if it's time to switch our primary leg.
1123                    #[cfg(feature = "conflux")]
1124                    if let Some(switch_cell) = self.maybe_update_primary_leg()? {
1125                        trace!(
1126                            old = ?old_primary_leg,
1127                            new = ?self.primary_id,
1128                            "Switching primary conflux leg..."
1129                        );
1130
1131                        self.primary_leg_mut()?.send_relay_cell(switch_cell).await?;
1132                    }
1133
1134                    // Use the possibly updated primary leg
1135                    Some(self.primary_id)
1136                }
1137            } else {
1138                // Non-multiplexed commands go on their original
1139                // circuit and hop
1140                leg
1141            }
1142        } else {
1143            // If there is no join point, it means this is not
1144            // a multi-path tunnel, so we continue using
1145            // the leg_id/hop the cmd came from.
1146            leg
1147        };
1148
1149        let leg = leg.unwrap_or(self.primary_id);
1150
1151        let circ = self
1152            .leg_mut(leg)
1153            .ok_or_else(|| internal!("leg disappeared?!"))?;
1154
1155        circ.send_relay_cell(msg).await
1156    }
1157
1158    /// Send a LINK cell down each unlinked leg.
1159    #[cfg(feature = "conflux")]
1160    pub(super) async fn link_circuits(
1161        &mut self,
1162        runtime: &tor_rtcompat::DynTimeProvider,
1163    ) -> crate::Result<()> {
1164        let (_leg_id, join_point) = self
1165            .primary_join_point()
1166            .ok_or_else(|| internal!("no join point when trying to send LINK"))?;
1167
1168        // Link all the circuits that haven't started the conflux handshake yet.
1169        for circ in self
1170            .legs
1171            .iter_mut()
1172            // TODO: it is an internal error if any of the legs don't have a conflux handler
1173            // (i.e. if conflux_status() returns None)
1174            .filter(|circ| circ.conflux_status() == Some(ConfluxStatus::Unlinked))
1175        {
1176            let v1_payload = V1LinkPayload::new(self.nonce, self.desired_ux);
1177            let link = ConfluxLink::new(v1_payload);
1178            let cell = AnyRelayMsgOuter::new(None, link.into());
1179
1180            circ.begin_conflux_link(join_point, cell, runtime).await?;
1181        }
1182
1183        // TODO(conflux): the caller should take care to not allow opening streams
1184        // until the conflux set is ready (i.e. until at least one of the legs completes
1185        // the handshake).
1186        //
1187        // We will probably need a channel for notifying the caller
1188        // of handshake completion/conflux set readiness
1189
1190        Ok(())
1191    }
1192
1193    /// Get the number of unlinked or non-conflux legs.
1194    #[cfg(feature = "conflux")]
1195    pub(super) fn num_unlinked(&self) -> usize {
1196        self.circuits()
1197            .filter(|circ| {
1198                let status = circ.conflux_status();
1199                status.is_none() || status == Some(ConfluxStatus::Unlinked)
1200            })
1201            .count()
1202    }
1203
1204    /// Check if the specified sequence number is the sequence number of the
1205    /// next message we're expecting to handle.
1206    pub(super) fn is_seqno_in_order(&self, seq_recv: u64) -> bool {
1207        let last_seq_delivered = self.last_seq_delivered.load(atomic::Ordering::Acquire);
1208        seq_recv == last_seq_delivered + 1
1209    }
1210
1211    /// Remove the circuit leg with the specified `UniqId` from this conflux set.
1212    ///
1213    /// Unlike [`ConfluxSet::remove`], this function does not check
1214    /// if the removal of the leg ought to trigger a reactor shutdown.
1215    ///
1216    /// Returns an error if the leg doesn't exit in the conflux set.
1217    fn remove_unchecked(&mut self, circ_uniq_id: UniqId) -> Result<Circuit, Bug> {
1218        let idx = self
1219            .legs
1220            .iter()
1221            .position(|circ| circ.unique_id() == circ_uniq_id)
1222            .ok_or_else(|| internal!("leg {circ_uniq_id:?} not found in conflux set"))?;
1223
1224        Ok(self.legs.remove(idx))
1225    }
1226
1227    /// Perform some circuit-padding-based event on the specified circuit.
1228    #[cfg(feature = "circ-padding")]
1229    pub(super) async fn run_padding_event(
1230        &mut self,
1231        circ_uniq_id: UniqId,
1232        padding_event: PaddingEvent,
1233    ) -> crate::Result<()> {
1234        use PaddingEvent as E;
1235        let Some(circ) = self.leg_mut(circ_uniq_id) else {
1236            // No such circuit; it must have gone away after generating this event.
1237            // Just ignore it.
1238            return Ok(());
1239        };
1240
1241        match padding_event {
1242            E::SendPadding(send_padding) => {
1243                circ.send_padding(send_padding).await?;
1244            }
1245            E::StartBlocking(start_blocking) => {
1246                circ.start_blocking_for_padding(start_blocking);
1247            }
1248            E::StopBlocking => {
1249                circ.stop_blocking_for_padding();
1250            }
1251        }
1252        Ok(())
1253    }
1254}
1255
1256/// An error returned when a method is expecting a single-leg conflux circuit,
1257/// but it is not single-leg.
1258#[derive(Clone, Debug, derive_more::Display, thiserror::Error)]
1259pub(super) struct NotSingleLegError(#[source] Bug);
1260
1261impl From<NotSingleLegError> for Bug {
1262    fn from(e: NotSingleLegError) -> Self {
1263        e.0
1264    }
1265}
1266
1267impl From<NotSingleLegError> for crate::Error {
1268    fn from(e: NotSingleLegError) -> Self {
1269        Self::from(e.0)
1270    }
1271}
1272
1273impl From<NotSingleLegError> for ReactorError {
1274    fn from(e: NotSingleLegError) -> Self {
1275        Self::from(e.0)
1276    }
1277}
1278
1279impl<I: Iterator> From<ExactlyOneError<I>> for NotSingleLegError {
1280    fn from(e: ExactlyOneError<I>) -> Self {
1281        // TODO: cannot wrap the ExactlyOneError with into_bad_api_usage
1282        // because it's not Send + Sync
1283        Self(bad_api_usage!("not a single leg conflux set ({e})"))
1284    }
1285}
1286
1287#[cfg(test)]
1288mod test {
1289    // Tested in [`crate::client::circuit::test`].
1290}