Skip to main content

tor_circmgr/
mgr.rs

1//! Abstract code to manage a set of tunnels which has underlying circuit(s).
2//!
3//! This module implements the real logic for deciding when and how to
4//! launch tunnels, and for which tunnels to hand out in response to
5//! which requests.
6//!
7//! For testing and abstraction purposes, this module _does not_
8//! actually know anything about tunnels _per se_.  Instead,
9//! everything is handled using a set of traits that are internal to this
10//! crate:
11//!
12//!  * [`AbstractTunnel`] is a view of a tunnel.
13//!  * [`AbstractTunnelBuilder`] knows how to build an `AbstractCirc`.
14//!
15//! Using these traits, the [`AbstractTunnelMgr`] object manages a set of
16//! tunnels , launching them as necessary, and keeping track of the
17//! restrictions on their use.
18
19// TODO:
20// - Testing
21//    - Error from prepare_action()
22//    - Error reported by restrict_mut?
23
24use crate::config::CircuitTiming;
25use crate::usage::{SupportedTunnelUsage, TargetTunnelUsage};
26use crate::{DirInfo, Error, PathConfig, Result, timeouts};
27
28use retry_error::RetryError;
29use tor_async_utils::mpsc_channel_no_memquota;
30use tor_basic_utils::onionperf_types::{OnionperfCircuitStatus, OnionperfEvent};
31use tor_basic_utils::retry::RetryDelay;
32use tor_config::MutCfg;
33use tor_error::{AbsRetryTime, HasRetryTime, debug_report, info_report, internal, warn_report};
34#[cfg(feature = "vanguards")]
35use tor_guardmgr::vanguards::VanguardMgr;
36use tor_linkspec::CircTarget;
37use tor_proto::circuit::UniqId;
38use tor_proto::client::circuit::{CircParameters, Path};
39use tor_rtcompat::{Runtime, SleepProviderExt};
40
41use async_trait::async_trait;
42use futures::channel::mpsc;
43use futures::future::{FutureExt, Shared};
44use futures::stream::{FuturesUnordered, StreamExt};
45use oneshot_fused_workaround as oneshot;
46use std::collections::HashMap;
47use std::fmt::Debug;
48use std::hash::Hash;
49use std::panic::AssertUnwindSafe;
50use std::sync::{self, Arc, Weak};
51use tor_rtcompat::SpawnExt;
52use tracing::{debug, instrument, trace, warn};
53use web_time_compat::{Duration, Instant};
54mod streams;
55
56/// Alias to force use of RandomState, regardless of features enabled in `weak_tables`.
57///
58/// See <https://github.com/tov/weak-table-rs/issues/23> for discussion.
59///
60/// (We could probably get away with a weaker hash function in this case, since
61/// the attacker _probably_ doesn't have control over our pointers.)
62type PtrWeakHashSet<T> = weak_table::PtrWeakHashSet<T, std::hash::RandomState>;
63
64/// Description of how we got a tunnel.
65#[non_exhaustive]
66#[derive(Debug, Copy, Clone, Eq, PartialEq)]
67pub(crate) enum TunnelProvenance {
68    /// This channel was newly launched, or was in progress and finished while
69    /// we were waiting.
70    NewlyCreated,
71    /// This channel already existed when we asked for it.
72    Preexisting,
73}
74
75/// An error returned when we cannot apply circuit restriction.
76#[derive(Clone, Debug, thiserror::Error)]
77#[non_exhaustive]
78pub enum RestrictionFailed {
79    /// Tried to restrict a specification, but the tunnel didn't support the
80    /// requested usage.
81    #[error("Specification did not support desired usage")]
82    NotSupported,
83}
84
85/// Minimal abstract view of a tunnel.
86///
87/// From this module's point of view, tunnels are simply objects
88/// with unique identities, and a possible closed-state.
89#[async_trait]
90pub(crate) trait AbstractTunnel: Debug {
91    /// Type for a unique identifier for tunnels.
92    type Id: Clone + Debug + Hash + Eq + Send + Sync;
93    /// Return the unique identifier for this tunnel.
94    ///
95    /// # Requirements
96    ///
97    /// The values returned by this function are unique for distinct
98    /// tunnels.
99    fn id(&self) -> Self::Id;
100
101    /// Return true if this tunnel is usable for some purpose.
102    ///
103    /// Reasons a tunnel might be unusable include being closed.
104    fn usable(&self) -> bool;
105
106    /// Return a list of [`Path`] objects describing the only circuit in this tunnel.
107    ///
108    /// Returns an error if the tunnel has more than one tunnel.
109    fn single_path(&self) -> tor_proto::Result<Arc<Path>>;
110
111    /// Return the number of hops in this tunnel.
112    ///
113    /// Returns an error if the circuit is closed.
114    ///
115    /// NOTE: This function will currently return only the number of hops
116    /// _currently_ in the tunnel. If there is an extend operation in progress,
117    /// the currently pending hop may or may not be counted, depending on whether
118    /// the extend operation finishes before this call is done.
119    fn n_hops(&self) -> tor_proto::Result<usize>;
120
121    /// Return true if this tunnel is closed and therefore unusable.
122    fn is_closing(&self) -> bool;
123
124    /// Return a process-unique identifier for this tunnel.
125    fn unique_id(&self) -> UniqId;
126
127    /// Extend the tunnel via the most appropriate handshake to a new `target` hop.
128    async fn extend<T: CircTarget + Sync>(
129        &self,
130        target: &T,
131        params: CircParameters,
132    ) -> tor_proto::Result<()>;
133
134    /// Return a time at which this tunnel is last known to be used,
135    /// or None if it is in use right now (or has never been used).
136    async fn last_known_to_be_used_at(&self) -> tor_proto::Result<Option<Instant>>;
137}
138
139/// A plan for an `AbstractCircBuilder` that can maybe be mutated by tests.
140///
141/// You should implement this trait using all default methods for all code that isn't test code.
142pub(crate) trait MockablePlan {
143    /// Add a reason string that was passed to `SleepProvider::block_advance()` to this object
144    /// so that it knows what to pass to `::release_advance()`.
145    fn add_blocked_advance_reason(&mut self, _reason: String) {}
146}
147
148/// An object that knows how to build tunnels.
149///
150/// This creates tunnels in two phases. First, a plan is
151/// made for how to build the tunnel. This planning phase should be
152/// relatively fast, and must not suspend or block.  Its purpose is to
153/// get an early estimate of which operations the tunnel will be able
154/// to support when it's done.
155///
156/// Second, the tunnel is actually built, using the plan as input.
157
158#[async_trait]
159pub(crate) trait AbstractTunnelBuilder<R: Runtime>: Send + Sync {
160    /// The tunnel type that this builder knows how to build.
161    type Tunnel: AbstractTunnel + Send + Sync;
162    /// An opaque type describing how a given tunnel will be built.
163    /// It may represent some or all of a path-or it may not.
164    //
165    // TODO: It would be nice to have this parameterized on a lifetime,
166    // and have that lifetime depend on the lifetime of the directory.
167    // But I don't think that rust can do that.
168    //
169    // HACK(eta): I don't like the fact that `MockablePlan` is necessary here.
170    type Plan: Send + Debug + MockablePlan;
171
172    // TODO: I'd like to have a Dir type here to represent
173    // create::DirInfo, but that would need to be parameterized too,
174    // and would make everything complicated.
175
176    /// Form a plan for how to build a new tunnel that supports `usage`.
177    ///
178    /// Return an opaque Plan object, and a new spec describing what
179    /// the tunnel will actually support when it's built.  (For
180    /// example, if the input spec requests a tunnel that connect to
181    /// port 80, then "planning" the tunnel might involve picking an
182    /// exit that supports port 80, and the resulting spec might be
183    /// the exit's complete list of supported ports.)
184    ///
185    /// # Requirements
186    ///
187    /// The resulting Spec must support `usage`.
188    fn plan_tunnel(
189        &self,
190        usage: &TargetTunnelUsage,
191        dir: DirInfo<'_>,
192    ) -> Result<(Self::Plan, SupportedTunnelUsage)>;
193
194    /// Construct a tunnel according to a given plan.
195    ///
196    /// On success, return a spec describing what the tunnel can be used for,
197    /// and the tunnel that was just constructed.
198    ///
199    /// This function should implement some kind of a timeout for
200    /// tunnel that are taking too long.
201    ///
202    /// # Requirements
203    ///
204    /// The spec that this function returns _must_ support the usage
205    /// that was originally passed to `plan_tunnel`.  It _must_ also
206    /// contain the spec that was originally returned by
207    /// `plan_tunnel`.
208    async fn build_tunnel(&self, plan: Self::Plan) -> Result<(SupportedTunnelUsage, Self::Tunnel)>;
209
210    /// Return a "parallelism factor" with which tunnels should be
211    /// constructed for a given purpose.
212    ///
213    /// If this function returns N, then whenever we launch tunnels
214    /// for this purpose, then we launch N in parallel.
215    ///
216    /// The default implementation returns 1.  The value of 0 is
217    /// treated as if it were 1.
218    fn launch_parallelism(&self, usage: &TargetTunnelUsage) -> usize {
219        let _ = usage; // default implementation ignores this.
220        1
221    }
222
223    /// Return a "parallelism factor" for which tunnels should be
224    /// used for a given purpose.
225    ///
226    /// If this function returns N, then whenever we select among
227    /// open tunnels for this purpose, we choose at random from the
228    /// best N.
229    ///
230    /// The default implementation returns 1.  The value of 0 is
231    /// treated as if it were 1.
232    // TODO: Possibly this doesn't belong in this trait.
233    fn select_parallelism(&self, usage: &TargetTunnelUsage) -> usize {
234        let _ = usage; // default implementation ignores this.
235        1
236    }
237
238    /// Return true if we are currently attempting to learn tunnel
239    /// timeouts by building testing tunnels.
240    fn learning_timeouts(&self) -> bool;
241
242    /// Flush state to the state manager if we own the lock.
243    ///
244    /// Return `Ok(true)` if we saved, and `Ok(false)` if we didn't hold the lock.
245    fn save_state(&self) -> Result<bool>;
246
247    /// Return this builder's [`PathConfig`].
248    fn path_config(&self) -> Arc<PathConfig>;
249
250    /// Replace this builder's [`PathConfig`].
251    // TODO: This is dead_code because we only call this for the CircuitBuilder specialization of
252    // CircMgr, not from the generic version, because this trait doesn't provide guardmgr, which is
253    // needed by the [`CircMgr::reconfigure`] function that would be the only caller of this. We
254    // should add `guardmgr` to this trait, make [`CircMgr::reconfigure`] generic, and remove this
255    // dead_code marking.
256    #[allow(dead_code)]
257    fn set_path_config(&self, new_config: PathConfig);
258
259    /// Return a reference to this builder's timeout estimator.
260    fn estimator(&self) -> &timeouts::Estimator;
261
262    /// Return a reference to this builder's `VanguardMgr`.
263    #[cfg(feature = "vanguards")]
264    fn vanguardmgr(&self) -> &Arc<VanguardMgr<R>>;
265
266    /// Replace our state with a new owning state, assuming we have
267    /// storage permission.
268    fn upgrade_to_owned_state(&self) -> Result<()>;
269
270    /// Reload persistent state from disk, if we don't have storage permission.
271    fn reload_state(&self) -> Result<()>;
272
273    /// Return a reference to this builder's `GuardMgr`.
274    fn guardmgr(&self) -> &tor_guardmgr::GuardMgr<R>;
275
276    /// Reconfigure this builder using the latest set of network parameters.
277    ///
278    /// (NOTE: for now, this only affects tunnel timeout estimation.)
279    fn update_network_parameters(&self, p: &tor_netdir::params::NetParameters);
280}
281
282/// Enumeration to track the expiration state of a tunnel.
283///
284/// A tunnel an either be unused (at which point it should expire if it is
285/// _still unused_ by a certain time, or dirty (at which point it should
286/// expire after a certain duration).
287///
288/// All tunnels start out "unused" and become "dirty" when their spec
289/// is first restricted -- that is, when they are first handed out to be
290/// used for a request.
291#[derive(Debug, Clone, PartialEq, Eq)]
292enum ExpirationInfo {
293    /// The tunnel has never been used, and has never been restricted for use with a request.
294    Unused {
295        /// A time when the tunnel was created.
296        created: Instant,
297    },
298
299    /// The tunnel is not-long-lived; we will expire by waiting until a certain amount of time
300    /// after it was first used.
301    Dirty {
302        /// The time at which this tunnel's spec was first restricted.
303        dirty_since: Instant,
304    },
305
306    /// The tunnel is long-lived; we will expire by waiting until it has passed
307    /// a certain amount of time without having any streams attached to it.
308    LongLived {
309        /// Last time at which the tunnel was checked and found not to have any streams.
310        ///
311        /// (This is a bit complicated: We have to be vague here, since we need
312        /// an async check to find out that a tunnel is used, or when it actually
313        /// became disused.)
314        last_known_to_be_used_at: Instant,
315    },
316}
317
318impl ExpirationInfo {
319    /// Return an ExpirationInfo for a newly created tunnel.
320    fn new(now: Instant) -> Self {
321        ExpirationInfo::Unused { created: now }
322    }
323
324    /// Mark this ExpirationInfo as having been in-use at `now`.
325    ///
326    /// If `long_lived` is false, the associated tunnel should expire a certain amount of time
327    /// after it was _first_ used.
328    /// If `long_lived` is true, the associated tunnel should expire a certain amount of time
329    /// after it was _last_ used.
330    fn mark_used(&mut self, now: Instant, long_lived: bool) {
331        if long_lived {
332            *self = ExpirationInfo::LongLived {
333                last_known_to_be_used_at: now,
334            };
335        } else {
336            match self {
337                ExpirationInfo::Unused { .. } => {
338                    // This is our first time using this circuit; mark it dirty
339                    *self = ExpirationInfo::Dirty { dirty_since: now };
340                }
341                ExpirationInfo::Dirty { .. } => {
342                    // no need to update; we're tracking the time when the circuit _first_ became
343                    // dirty, so further uses don't matter.
344                }
345                ExpirationInfo::LongLived { .. } => {
346                    // shouldn't occur: we shouldn't be able to attach a stream with non-long-lived isolation
347                    // to a tunnel marked as long-lived.  In this case we leave the timestamp alone.
348                    // (If there were a bug here, it would be harmless, since we would
349                    // correct the timestamp the next time we tried to expire the circuit.)
350                }
351            }
352        }
353    }
354
355    /// Return an internal error if this ExpirationInfo is not marked as long-lived.
356    fn check_long_lived(&self) -> Result<()> {
357        match self {
358            ExpirationInfo::Unused { .. } | ExpirationInfo::Dirty { .. } => Err(internal!(
359                "Tunnel was not long-lived as expected. (Expiration status: {:?})",
360                self
361            )
362            .into()),
363            ExpirationInfo::LongLived { .. } => Ok(()),
364        }
365    }
366}
367
368/// Settings to determine when circuits are expired.
369#[derive(Clone, Debug)]
370pub(crate) struct ExpirationParameters {
371    /// Any unused circuit is expired this long after it was created.
372    expire_unused_after: Duration,
373    /// Any non long-lived dirty circuit is expired this long after it first becomes dirty.
374    expire_dirty_after: Duration,
375    /// Any long-lived circuit is expired after having been disused for this long.
376    expire_disused_after: Duration,
377}
378
379/// An entry for an open tunnel held by an `AbstractTunnelMgr`.
380#[derive(Debug, Clone)]
381pub(crate) struct OpenEntry<T> {
382    /// The supported usage for this tunnel.
383    spec: SupportedTunnelUsage,
384    /// The tunnel under management.
385    tunnel: Arc<T>,
386    /// When does this tunnel expire?
387    ///
388    /// (Note that expired tunnels are removed from the manager,
389    /// which does not actually close them until there are no more
390    /// references to them.)
391    expiration: ExpirationInfo,
392}
393
394impl<T: AbstractTunnel> OpenEntry<T> {
395    /// Make a new OpenEntry for a given tunnel and spec.
396    fn new(spec: SupportedTunnelUsage, tunnel: T, expiration: ExpirationInfo) -> Self {
397        OpenEntry {
398            spec,
399            tunnel: tunnel.into(),
400            expiration,
401        }
402    }
403
404    /// Return true if the underlying tunnel can be used for `usage`.
405    pub(crate) fn supports(&self, usage: &TargetTunnelUsage) -> bool {
406        self.tunnel.usable() && self.spec.supports(usage)
407    }
408
409    /// Change the underlying tunnel's permissible usage, based on its having
410    /// been used for `usage` at time `now`.
411    ///
412    /// Return an error if the tunnel may not be used for `usage`.
413    fn restrict_mut(&mut self, usage: &TargetTunnelUsage, now: Instant) -> Result<()> {
414        self.spec.restrict_mut(usage)?;
415        self.expiration.mark_used(now, self.spec.is_long_lived());
416        Ok(())
417    }
418
419    /// Find the "best" entry from a slice of OpenEntry for supporting
420    /// a given `usage`.
421    ///
422    /// If `parallelism` is some N greater than 1, we pick randomly
423    /// from the best `N` tunnels.
424    ///
425    /// # Requirements
426    ///
427    /// Requires that `ents` is nonempty, and that every element of `ents`
428    /// supports `spec`.
429    fn find_best<'a>(
430        // we do not mutate `ents`, but to return `&mut Self` we must have a mutable borrow
431        ents: &'a mut [&'a mut Self],
432        usage: &TargetTunnelUsage,
433        parallelism: usize,
434    ) -> &'a mut Self {
435        let _ = usage; // not yet used.
436        use rand::seq::IndexedMutRandom as _;
437        let parallelism = parallelism.clamp(1, ents.len());
438        // TODO: Actually look over the whole list to see which is better.
439        let slice = &mut ents[0..parallelism];
440        let mut rng = rand::rng();
441        slice.choose_mut(&mut rng).expect("Input list was empty")
442    }
443
444    /// Return true if this tunnel should be expired given that the current time is `now`,
445    /// and the current settings are `params`.
446    fn should_expire(&self, now: Instant, params: &ExpirationParameters) -> ShouldExpire {
447        match self.expiration {
448            ExpirationInfo::Unused { created } => {
449                ShouldExpire::certain(now, created + params.expire_unused_after)
450            }
451            ExpirationInfo::Dirty { dirty_since } => {
452                ShouldExpire::certain(now, dirty_since + params.expire_dirty_after)
453            }
454            ExpirationInfo::LongLived {
455                last_known_to_be_used_at,
456            } => {
457                ShouldExpire::uncertain(now, last_known_to_be_used_at + params.expire_disused_after)
458            }
459        }
460    }
461}
462
463/// When should a tunnel expire?
464///
465/// Reflects possible uncertainty.
466#[derive(Clone, Copy, Debug, Eq, PartialEq)]
467enum ShouldExpire {
468    /// The tunnel should expire now.
469    Now,
470    /// The circuit might expire now; we need to check.
471    ///
472    /// (This is the result we get when we know that this is a tunnel that should expire
473    /// if it has gone for some duration D without having any streams on it,
474    /// and that it definitely had a stream at time T.  It is now at least time T+D,
475    /// but we don't know whether the tunnel has any streams in the intervening time.
476    /// We need to call the async fn `last_known_to_be_used_at` to check.)
477    PossiblyNow,
478    /// The tunnel will not expire before the specified time.
479    NotBefore(Instant),
480}
481
482impl ShouldExpire {
483    /// Return a ShouldExpire reflecting an expiration that is known to be happening at `expiration`.
484    fn certain(now: Instant, expiration: Instant) -> Self {
485        if now >= expiration {
486            ShouldExpire::Now
487        } else {
488            ShouldExpire::NotBefore(expiration)
489        }
490    }
491
492    /// Return a ShouldExpire reflecting an expiration that is known to be no sooner than `expiration`,
493    /// but possibly later.
494    fn uncertain(now: Instant, expiration: Instant) -> Self {
495        if now >= expiration {
496            ShouldExpire::PossiblyNow
497        } else {
498            ShouldExpire::NotBefore(expiration)
499        }
500    }
501}
502
503/// A result type whose "Ok" value is the Id for a tunnel from B.
504type PendResult<B, R> = Result<<<B as AbstractTunnelBuilder<R>>::Tunnel as AbstractTunnel>::Id>;
505
506/// An in-progress tunnel request tracked by an `AbstractTunnelMgr`.
507///
508/// (In addition to tracking tunnels, `AbstractTunnelMgr` tracks
509/// _requests_ for tunnels.  The manager uses these entries if it
510/// finds that some tunnel created _after_ a request first launched
511/// might meet the request's requirements.)
512struct PendingRequest<B: AbstractTunnelBuilder<R>, R: Runtime> {
513    /// Usage for the operation requested by this request
514    usage: TargetTunnelUsage,
515    /// A channel to use for telling this request about tunnels that it
516    /// might like.
517    notify: mpsc::Sender<PendResult<B, R>>,
518}
519
520impl<B: AbstractTunnelBuilder<R>, R: Runtime> PendingRequest<B, R> {
521    /// Return true if this request would be supported by `spec`.
522    fn supported_by(&self, spec: &SupportedTunnelUsage) -> bool {
523        spec.supports(&self.usage)
524    }
525}
526
527/// An entry for an under-construction in-progress tunnel tracked by
528/// an `AbstractTunnelMgr`.
529#[derive(Debug)]
530struct PendingEntry<B: AbstractTunnelBuilder<R>, R: Runtime> {
531    /// Specification that this tunnel will support, if every pending
532    /// request that is waiting for it is attached to it.
533    ///
534    /// This spec becomes more and more restricted as more pending
535    /// requests are waiting for this tunnel.
536    ///
537    /// This spec is contained by circ_spec, and must support the usage
538    /// of every pending request that's waiting for this tunnel.
539    tentative_assignment: sync::Mutex<SupportedTunnelUsage>,
540    /// A shared future for requests to use when waiting for
541    /// notification of this tunnel's success.
542    receiver: Shared<oneshot::Receiver<PendResult<B, R>>>,
543}
544
545impl<B: AbstractTunnelBuilder<R>, R: Runtime> PendingEntry<B, R> {
546    /// Make a new PendingEntry that starts out supporting a given
547    /// spec.  Return that PendingEntry, along with a Sender to use to
548    /// report the result of building this tunnel.
549    fn new(spec: &SupportedTunnelUsage) -> (Self, oneshot::Sender<PendResult<B, R>>) {
550        let tentative_assignment = sync::Mutex::new(spec.clone());
551        let (sender, receiver) = oneshot::channel();
552        let receiver = receiver.shared();
553        let entry = PendingEntry {
554            tentative_assignment,
555            receiver,
556        };
557        (entry, sender)
558    }
559
560    /// Return true if this tunnel's current tentative assignment
561    /// supports `usage`.
562    fn supports(&self, usage: &TargetTunnelUsage) -> bool {
563        let assignment = self.tentative_assignment.lock().expect("poisoned lock");
564        assignment.supports(usage)
565    }
566
567    /// Try to change the tentative assignment of this tunnel by
568    /// restricting it for use with `usage`.
569    ///
570    /// Return an error if the current tentative assignment didn't
571    /// support `usage` in the first place.
572    fn tentative_restrict_mut(&self, usage: &TargetTunnelUsage) -> Result<()> {
573        if let Ok(mut assignment) = self.tentative_assignment.lock() {
574            assignment.restrict_mut(usage)?;
575        }
576        Ok(())
577    }
578
579    /// Find the best PendingEntry values from a slice for use with
580    /// `usage`.
581    ///
582    /// # Requirements
583    ///
584    /// The `ents` slice must not be empty.  Every element of `ents`
585    /// must support the given spec.
586    fn find_best(ents: &[Arc<Self>], usage: &TargetTunnelUsage) -> Vec<Arc<Self>> {
587        // TODO: Actually look over the whole list to see which is better.
588        let _ = usage; // currently unused
589        vec![Arc::clone(&ents[0])]
590    }
591}
592
593/// Wrapper type to represent the state between planning to build a
594/// tunnel and constructing it.
595#[derive(Debug)]
596struct TunnelBuildPlan<B: AbstractTunnelBuilder<R>, R: Runtime> {
597    /// The Plan object returned by [`AbstractTunnelBuilder::plan_tunnel`].
598    plan: B::Plan,
599    /// A sender to notify any pending requests when this tunnel is done.
600    sender: oneshot::Sender<PendResult<B, R>>,
601    /// A strong entry to the PendingEntry for this tunnel build attempt.
602    pending: Arc<PendingEntry<B, R>>,
603}
604
605/// The inner state of an [`AbstractTunnelMgr`].
606struct TunnelList<B: AbstractTunnelBuilder<R>, R: Runtime> {
607    /// A map from tunnel ID to [`OpenEntry`] values for all managed
608    /// open tunnels.
609    ///
610    /// A tunnel is added here from [`AbstractTunnelMgr::do_launch`] when we find
611    /// that it completes successfully, and has not been cancelled.
612    /// When we decide that such a tunnel should no longer be handed out for
613    /// any new requests, we "retire" the tunnel by removing it from this map.
614    #[allow(clippy::type_complexity)]
615    open_tunnels: HashMap<<B::Tunnel as AbstractTunnel>::Id, OpenEntry<B::Tunnel>>,
616    /// Weak-set of PendingEntry for tunnels that are being built.
617    ///
618    /// Because this set only holds weak references, and the only strong
619    /// reference to the PendingEntry is held by the task building the tunnel,
620    /// this set's members are lazily removed after the tunnel is either built
621    /// or fails to build.
622    ///
623    /// This set is used for two purposes:
624    ///
625    /// 1. When a tunnel request finds that there is no open tunnel for its
626    ///    purposes, it checks here to see if there is a pending tunnel that it
627    ///    could wait for.
628    /// 2. When a pending tunnel finishes building, it checks here to make sure
629    ///    that it has not been cancelled. (Removing an entry from this set marks
630    ///    it as cancelled.)
631    ///
632    /// An entry is added here in [`AbstractTunnelMgr::prepare_action`] when we
633    /// decide that a tunnel needs to be launched.
634    ///
635    /// Later, in [`AbstractTunnelMgr::do_launch`], once the tunnel has finished
636    /// (or failed), we remove the entry (by pointer identity).
637    /// If we cannot find the entry, we conclude that the request has been
638    /// _cancelled_, and so we discard any tunnel that was created.
639    pending_tunnels: PtrWeakHashSet<Weak<PendingEntry<B, R>>>,
640    /// Weak-set of PendingRequest for requests that are waiting for a
641    /// tunnel to be built.
642    ///
643    /// Because this set only holds weak references, and the only
644    /// strong reference to the PendingRequest is held by the task
645    /// waiting for the tunnel to be built, this set's members are
646    /// lazily removed after the request succeeds or fails.
647    pending_requests: PtrWeakHashSet<Weak<PendingRequest<B, R>>>,
648}
649
650impl<B: AbstractTunnelBuilder<R>, R: Runtime> TunnelList<B, R> {
651    /// Make a new empty `CircList`
652    fn new() -> Self {
653        TunnelList {
654            open_tunnels: HashMap::new(),
655            pending_tunnels: PtrWeakHashSet::new(),
656            pending_requests: PtrWeakHashSet::new(),
657        }
658    }
659
660    /// Add `e` to the list of open tunnels.
661    fn add_open(&mut self, e: OpenEntry<B::Tunnel>) {
662        let id = e.tunnel.id();
663        self.open_tunnels.insert(id, e);
664    }
665
666    /// Find all the usable open tunnels that support `usage`.
667    ///
668    /// Return None if there are no such tunnels.
669    fn find_open(&mut self, usage: &TargetTunnelUsage) -> Option<Vec<&mut OpenEntry<B::Tunnel>>> {
670        let list = self.open_tunnels.values_mut();
671        let v = SupportedTunnelUsage::find_supported(list, usage);
672        if v.is_empty() { None } else { Some(v) }
673    }
674
675    /// Find an open tunnel by ID.
676    ///
677    /// Return None if no such tunnels exists in this list.
678    fn get_open_mut(
679        &mut self,
680        id: &<B::Tunnel as AbstractTunnel>::Id,
681    ) -> Option<&mut OpenEntry<B::Tunnel>> {
682        self.open_tunnels.get_mut(id)
683    }
684
685    /// Extract an open tunnel by ID, removing it from this list.
686    ///
687    /// Return None if no such tunnel exists in this list.
688    fn take_open(
689        &mut self,
690        id: &<B::Tunnel as AbstractTunnel>::Id,
691    ) -> Option<OpenEntry<B::Tunnel>> {
692        self.open_tunnels.remove(id)
693    }
694
695    /// Remove tunnels based on expiration times.
696    ///
697    /// We remove every unused tunnel that is set to expire by
698    /// `unused_cutoff`, and every dirty tunnel that has been dirty
699    /// since before `dirty_cutoff`.
700    ///
701    /// Return the next time at which anything will definitely expire,
702    /// and a list of long-lived tunnels where we need to check their usage status
703    /// before we can be sure if they are expired.
704    #[must_use]
705    fn expire_tunnels(
706        &mut self,
707        now: Instant,
708        params: &ExpirationParameters,
709    ) -> (Option<Instant>, Vec<Weak<B::Tunnel>>) {
710        let mut need_check = Vec::new();
711        let mut earliest_expiration = None;
712        self.open_tunnels
713            .retain(|_k, v| match v.should_expire(now, params) {
714                // Expires now: Do not retain.
715                ShouldExpire::Now => false,
716
717                // Will expire at `when`: keep, but update `earliest_expiration`.
718                ShouldExpire::NotBefore(when) => {
719                    earliest_expiration = match earliest_expiration {
720                        Some(t) if t < when => Some(t),
721                        _ => Some(when),
722                    };
723                    true
724                }
725
726                // Need to check tunnel to see if/when it is disused.
727                ShouldExpire::PossiblyNow => {
728                    need_check.push(Arc::downgrade(&v.tunnel));
729                    true
730                }
731            });
732        (earliest_expiration, need_check)
733    }
734
735    /// Return the time when the tunnel with given `id`, should expire.
736    ///
737    /// Return None if no such tunnel exists.
738    fn tunnel_should_expire(
739        &mut self,
740        id: &<B::Tunnel as AbstractTunnel>::Id,
741        now: Instant,
742        params: &ExpirationParameters,
743    ) -> Option<ShouldExpire> {
744        self.open_tunnels
745            .get(id)
746            .map(|v| v.should_expire(now, params))
747    }
748
749    /// Update the "last known to be in use" time of a long-lived tunnel with ID `id`,
750    /// based on learning when it was last used.
751    ///
752    /// Expire the tunnel if appropriate.
753    ///
754    /// If the tunnel is still part of the map, return the next instant at which it might expire.
755    ///
756    /// Returns an error if the tunnel was present but was _not_ already marked as long-lived.
757    fn update_long_lived_tunnel_last_used(
758        &mut self,
759        id: &<B::Tunnel as AbstractTunnel>::Id,
760        now: Instant,
761        params: &ExpirationParameters,
762        disused_since: &tor_proto::Result<Option<Instant>>,
763    ) -> crate::Result<Option<Instant>> {
764        let Ok(disused_since) = disused_since else {
765            // got an error looking up disused time: discard the circuit.
766            let discard = self.take_open(id);
767            if let Some(ent) = discard {
768                ent.expiration.check_long_lived()?;
769            }
770            return Ok(None);
771        };
772        let Some(tun) = self.open_tunnels.get_mut(id) else {
773            // Circuit isn't there. Return.
774            return Ok(None);
775        };
776        tun.expiration.check_long_lived()?;
777        let last_known_in_use_at = disused_since.unwrap_or(now);
778
779        tun.expiration.mark_used(last_known_in_use_at, true);
780        match tun.should_expire(now, params) {
781            ShouldExpire::Now | ShouldExpire::PossiblyNow => {
782                let _discard = self.take_open(id);
783                Ok(None)
784            }
785            ShouldExpire::NotBefore(instant) => Ok(Some(instant)),
786        }
787    }
788
789    /// Add `pending` to the set of in-progress tunnels.
790    fn add_pending_tunnel(&mut self, pending: Arc<PendingEntry<B, R>>) {
791        self.pending_tunnels.insert(pending);
792    }
793
794    /// Find all pending tunnels that support `usage`.
795    ///
796    /// If no such tunnels are currently being built, return None.
797    fn find_pending_tunnels(
798        &self,
799        usage: &TargetTunnelUsage,
800    ) -> Option<Vec<Arc<PendingEntry<B, R>>>> {
801        let result: Vec<_> = self
802            .pending_tunnels
803            .iter()
804            .filter(|p| p.supports(usage))
805            .filter(|p| !matches!(p.receiver.peek(), Some(Err(_))))
806            .collect();
807
808        if result.is_empty() {
809            None
810        } else {
811            Some(result)
812        }
813    }
814
815    /// Return true if `circ` is still pending.
816    ///
817    /// A tunnel will become non-pending when finishes (successfully or not), or when it's
818    /// removed from this list via `clear_all_tunnels()`.
819    fn tunnel_is_pending(&self, circ: &Arc<PendingEntry<B, R>>) -> bool {
820        self.pending_tunnels.contains(circ)
821    }
822
823    /// Construct and add a new entry to the set of request waiting
824    /// for a tunnel.
825    ///
826    /// Return the request, and a new receiver stream that it should
827    /// use for notification of possible tunnels to use.
828    fn add_pending_request(&mut self, pending: &Arc<PendingRequest<B, R>>) {
829        self.pending_requests.insert(Arc::clone(pending));
830    }
831
832    /// Return all pending requests that would be satisfied by a tunnel
833    /// that supports `circ_spec`.
834    fn find_pending_requests(
835        &self,
836        circ_spec: &SupportedTunnelUsage,
837    ) -> Vec<Arc<PendingRequest<B, R>>> {
838        self.pending_requests
839            .iter()
840            .filter(|pend| pend.supported_by(circ_spec))
841            .collect()
842    }
843
844    /// Clear all pending and open tunnels.
845    ///
846    /// Calling `clear_all_tunnels` ensures that any request that is answered _after
847    /// this method runs_ will receive a tunnels that was launched _after this
848    /// method runs_.
849    fn clear_all_tunnels(&mut self) {
850        // Note that removing entries from pending_circs will also cause the
851        // tunnel tasks to realize that they are cancelled when they
852        // go to tell anybody about their results.
853        self.pending_tunnels.clear();
854        self.open_tunnels.clear();
855    }
856}
857
858/// Timing information for tunnels that have been built but never used.
859///
860/// Currently taken from the network parameters.
861struct UnusedTimings {
862    /// Minimum lifetime of a tunnel created while learning
863    /// tunnel timeouts.
864    learning: Duration,
865    /// Minimum lifetime of a tunnel created while not learning
866    /// tunnel timeouts.
867    not_learning: Duration,
868}
869
870// This isn't really fallible, given the definitions of the underlying
871// types.
872#[allow(clippy::fallible_impl_from)]
873impl From<&tor_netdir::params::NetParameters> for UnusedTimings {
874    fn from(v: &tor_netdir::params::NetParameters) -> Self {
875        // These try_into() calls can't fail, so unwrap() can't panic.
876        #[allow(clippy::unwrap_used)]
877        UnusedTimings {
878            learning: v
879                .unused_client_circ_timeout_while_learning_cbt
880                .try_into()
881                .unwrap(),
882            not_learning: v.unused_client_circ_timeout.try_into().unwrap(),
883        }
884    }
885}
886
887/// Abstract implementation for tunnel management.
888///
889/// The algorithm provided here is fairly simple. In its simplest form:
890///
891/// When somebody asks for a tunnel for a given operation: if we find
892/// one open already, we return it.  If we find in-progress tunnels
893/// that would meet our needs, we wait for one to finish (or for all
894/// to fail).  And otherwise, we launch one or more tunnels to meet the
895/// request's needs.
896///
897/// If this process fails, then we retry it, up to a timeout or a
898/// numerical limit.
899///
900/// If a tunnel not previously considered for a given request
901/// finishes before the request is satisfied, and if the tunnel would
902/// satisfy the request, we try to give that tunnel as an answer to
903/// that request even if it was not one of the tunnels that request
904/// was waiting for.
905pub(crate) struct AbstractTunnelMgr<B: AbstractTunnelBuilder<R>, R: Runtime> {
906    /// Builder used to construct tunnels.
907    builder: B,
908    /// An asynchronous runtime to use for launching tasks and
909    /// checking timeouts.
910    runtime: R,
911    /// A CircList to manage our list of tunnels, requests, and
912    /// pending tunnels.
913    tunnels: sync::Mutex<TunnelList<B, R>>,
914
915    /// Configured information about when to expire tunnels and requests.
916    circuit_timing: MutCfg<CircuitTiming>,
917
918    /// Minimum lifetime of an unused tunnel.
919    ///
920    /// Derived from the network parameters.
921    unused_timing: sync::Mutex<UnusedTimings>,
922}
923
924/// An action to take in order to satisfy a request for a tunnel.
925enum Action<B: AbstractTunnelBuilder<R>, R: Runtime> {
926    /// We found an open tunnel: return immediately.
927    Open(Arc<B::Tunnel>),
928    /// We found one or more pending tunnels: wait until one succeeds,
929    /// or all fail.
930    Wait(FuturesUnordered<Shared<oneshot::Receiver<PendResult<B, R>>>>),
931    /// We should launch tunnels: here are the instructions for how
932    /// to do so.
933    Build(Vec<TunnelBuildPlan<B, R>>),
934}
935
936impl<B: AbstractTunnelBuilder<R> + 'static, R: Runtime> AbstractTunnelMgr<B, R> {
937    /// Construct a new AbstractTunnelMgr.
938    pub(crate) fn new(builder: B, runtime: R, circuit_timing: CircuitTiming) -> Self {
939        let circs = sync::Mutex::new(TunnelList::new());
940        let dflt_params = tor_netdir::params::NetParameters::default();
941        let unused_timing = (&dflt_params).into();
942        AbstractTunnelMgr {
943            builder,
944            runtime,
945            tunnels: circs,
946            circuit_timing: circuit_timing.into(),
947            unused_timing: sync::Mutex::new(unused_timing),
948        }
949    }
950
951    /// Reconfigure this manager using the latest set of network parameters.
952    pub(crate) fn update_network_parameters(&self, p: &tor_netdir::params::NetParameters) {
953        let mut u = self
954            .unused_timing
955            .lock()
956            .expect("Poisoned lock for unused_timing");
957        *u = p.into();
958    }
959
960    /// Return this manager's [`CircuitTiming`].
961    pub(crate) fn circuit_timing(&self) -> Arc<CircuitTiming> {
962        self.circuit_timing.get()
963    }
964
965    /// Return this manager's [`CircuitTiming`].
966    pub(crate) fn set_circuit_timing(&self, new_config: CircuitTiming) {
967        self.circuit_timing.replace(new_config);
968    }
969    /// Return a circuit suitable for use with a given `usage`,
970    /// creating that circuit if necessary, and restricting it
971    /// under the assumption that it will be used for that spec.
972    ///
973    /// This is the primary entry point for AbstractTunnelMgr.
974    #[instrument(level = "trace", skip_all)]
975    pub(crate) async fn get_or_launch(
976        self: &Arc<Self>,
977        usage: &TargetTunnelUsage,
978        dir: DirInfo<'_>,
979    ) -> Result<(Arc<B::Tunnel>, TunnelProvenance)> {
980        /// Largest number of "resets" that we will accept in this attempt.
981        ///
982        /// A "reset" is an internally generated error that does not represent a
983        /// real problem; only a "whoops, got to try again" kind of a situation.
984        /// For example, if we reconfigure in the middle of an attempt and need
985        /// to re-launch the circuit, that counts as a "reset", since there was
986        /// nothing actually _wrong_ with the circuit we were building.
987        ///
988        /// We accept more resets than we do real failures. However,
989        /// we don't accept an unlimited number: we don't want to inadvertently
990        /// permit infinite loops here. If we ever bump against this limit, we
991        /// should not automatically increase it: we should instead figure out
992        /// why it is happening and try to make it not happen.
993        const MAX_RESETS: usize = 8;
994
995        let circuit_timing = self.circuit_timing();
996        let timeout_at = self.runtime.now() + circuit_timing.request_timeout;
997        let max_tries = circuit_timing.request_max_retries;
998        // We compute the maximum number of failures by dividing the maximum
999        // number of circuits to attempt by the number that will be launched in
1000        // parallel for each iteration.
1001        let max_failures = usize::div_ceil(
1002            max_tries as usize,
1003            std::cmp::max(1, self.builder.launch_parallelism(usage)),
1004        );
1005
1006        let mut retry_schedule = RetryDelay::from_msec(100);
1007        let mut retry_err = RetryError::<Box<Error>>::in_attempt_to("find or build a tunnel");
1008
1009        let mut n_failures = 0;
1010        let mut n_resets = 0;
1011
1012        for attempt_num in 1.. {
1013            // How much time is remaining?
1014            let remaining = match timeout_at.checked_duration_since(self.runtime.now()) {
1015                None => {
1016                    retry_err.push_timed(
1017                        Error::RequestTimeout,
1018                        self.runtime.now(),
1019                        Some(self.runtime.wallclock()),
1020                    );
1021                    break;
1022                }
1023                Some(t) => t,
1024            };
1025
1026            let error = match self.prepare_action(usage, dir, true) {
1027                Ok(action) => {
1028                    // We successfully found an action: Take that action.
1029                    let outcome = self
1030                        .runtime
1031                        .timeout(remaining, Arc::clone(self).take_action(action, usage))
1032                        .await;
1033
1034                    match outcome {
1035                        Ok(Ok(circ)) => {
1036                            // TODO: Give usage as Value, probably once tracing valuable feature is
1037                            // stable.
1038                            tracing::trace!(
1039                                onionperf = true,
1040                                usage = ?usage,
1041                                event = ?OnionperfEvent::Circuit(OnionperfCircuitStatus::Built),
1042                            );
1043                            return Ok(circ);
1044                        }
1045                        Ok(Err(e)) => {
1046                            debug!("Circuit attempt {} failed.", attempt_num);
1047                            Error::RequestFailed(e)
1048                        }
1049                        Err(_) => {
1050                            // We ran out of "remaining" time; there is nothing
1051                            // more to be done.
1052                            warn!("All tunnel attempts failed due to timeout");
1053                            retry_err.push_timed(
1054                                Error::RequestTimeout,
1055                                self.runtime.now(),
1056                                Some(self.runtime.wallclock()),
1057                            );
1058                            break;
1059                        }
1060                    }
1061                }
1062                Err(e) => {
1063                    // We couldn't pick the action!
1064                    debug_report!(
1065                        &e,
1066                        "Couldn't pick action for tunnel attempt {}",
1067                        attempt_num,
1068                    );
1069                    e
1070                }
1071            };
1072
1073            // There's been an error.  See how long we wait before we retry.
1074            let now = self.runtime.now();
1075            let retry_time =
1076                error.abs_retry_time(now, || retry_schedule.next_delay(&mut rand::rng()));
1077
1078            let (count, count_limit) = if error.is_internal_reset() {
1079                (&mut n_resets, MAX_RESETS)
1080            } else {
1081                (&mut n_failures, max_failures)
1082            };
1083            // Record the error, flattening it if needed.
1084            match error {
1085                // Flatten nested RetryError, using mockable time for each error
1086                Error::RequestFailed(e) => {
1087                    retry_err.extend_from_retry_error(e);
1088                }
1089                e => retry_err.push_timed(e, now, Some(self.runtime.wallclock())),
1090            }
1091
1092            *count += 1;
1093            // If we have reached our limit of this kind of problem, we're done.
1094            if *count >= count_limit {
1095                warn!("Reached circuit build retry limit, exiting...");
1096                break;
1097            }
1098
1099            // Wait, or not, as appropriate.
1100            match retry_time {
1101                AbsRetryTime::Immediate => {}
1102                AbsRetryTime::Never => break,
1103                AbsRetryTime::At(t) => {
1104                    let remaining = timeout_at.saturating_duration_since(now);
1105                    let delay = t.saturating_duration_since(now);
1106                    trace!(?delay, "Waiting to retry...");
1107                    self.runtime.sleep(std::cmp::min(delay, remaining)).await;
1108                }
1109            }
1110        }
1111
1112        warn!("Request failed");
1113        Err(Error::RequestFailed(retry_err))
1114    }
1115
1116    /// Make sure a circuit exists, without actually asking for it.
1117    ///
1118    /// Make sure that there is a circuit (built or in-progress) that could be
1119    /// used for `usage`, and launch one or more circuits in a background task
1120    /// if there is not.
1121    // TODO: This should probably take some kind of parallelism parameter.
1122    #[cfg(test)]
1123    pub(crate) fn ensure_tunnel(
1124        self: &Arc<Self>,
1125        usage: &TargetTunnelUsage,
1126        dir: DirInfo<'_>,
1127    ) -> Result<()> {
1128        let action = self.prepare_action(usage, dir, false)?;
1129        if let Action::Build(plans) = action {
1130            for plan in plans {
1131                let self_clone = Arc::clone(self);
1132                let _ignore_receiver = self_clone.spawn_launch(usage, plan);
1133            }
1134        }
1135
1136        Ok(())
1137    }
1138
1139    /// Choose which action we should take in order to provide a tunnel
1140    /// for a given `usage`.
1141    ///
1142    /// If `restrict_circ` is true, we restrict the spec of any
1143    /// circ we decide to use to mark that it _is_ being used for
1144    /// `usage`.
1145    #[instrument(level = "trace", skip_all)]
1146    fn prepare_action(
1147        &self,
1148        usage: &TargetTunnelUsage,
1149        dir: DirInfo<'_>,
1150        restrict_circ: bool,
1151    ) -> Result<Action<B, R>> {
1152        let mut list = self.tunnels.lock().expect("poisoned lock");
1153
1154        if let Some(mut open) = list.find_open(usage) {
1155            // We have open tunnels that meet the spec: return the best one.
1156            let parallelism = self.builder.select_parallelism(usage);
1157            let best = OpenEntry::find_best(&mut open, usage, parallelism);
1158            if restrict_circ {
1159                let now = self.runtime.now();
1160                best.restrict_mut(usage, now)?;
1161            }
1162            // TODO: If we have fewer tunnels here than our select
1163            // parallelism, perhaps we should launch more?
1164
1165            return Ok(Action::Open(best.tunnel.clone()));
1166        }
1167
1168        if let Some(pending) = list.find_pending_tunnels(usage) {
1169            // There are pending tunnels that could meet the spec.
1170            // Restrict them under the assumption that they could all
1171            // be used for this, and then wait until one is ready (or
1172            // all have failed)
1173            let best = PendingEntry::find_best(&pending, usage);
1174            if restrict_circ {
1175                for item in &best {
1176                    // TODO: Do we want to tentatively restrict _all_ of these?
1177                    // not clear to me.
1178                    item.tentative_restrict_mut(usage)?;
1179                }
1180            }
1181            let stream = best.iter().map(|item| item.receiver.clone()).collect();
1182            // TODO: if we have fewer tunnels here than our launch
1183            // parallelism, we might want to launch more.
1184
1185            return Ok(Action::Wait(stream));
1186        }
1187
1188        // Okay, we need to launch tunnels here.
1189        let parallelism = std::cmp::max(1, self.builder.launch_parallelism(usage));
1190        let mut plans = Vec::new();
1191        let mut last_err = None;
1192        for _ in 0..parallelism {
1193            match self.plan_by_usage(dir, usage) {
1194                Ok((pending, plan)) => {
1195                    list.add_pending_tunnel(pending);
1196                    plans.push(plan);
1197                }
1198                Err(e) => {
1199                    debug!("Unable to make a plan for {:?}: {}", usage, e);
1200                    last_err = Some(e);
1201                }
1202            }
1203        }
1204        if !plans.is_empty() {
1205            Ok(Action::Build(plans))
1206        } else if let Some(last_err) = last_err {
1207            Err(last_err)
1208        } else {
1209            // we didn't even try to plan anything!
1210            Err(internal!("no plans were built, but no errors were found").into())
1211        }
1212    }
1213
1214    /// Execute an action returned by pick-action, and return the
1215    /// resulting tunnel or error.
1216    #[allow(clippy::type_complexity)] // TODO #2010: Refactor
1217    #[instrument(level = "trace", skip_all)]
1218    async fn take_action(
1219        self: Arc<Self>,
1220        act: Action<B, R>,
1221        usage: &TargetTunnelUsage,
1222    ) -> std::result::Result<(Arc<B::Tunnel>, TunnelProvenance), RetryError<Box<Error>>> {
1223        /// Store the error `err` into `retry_err`, as appropriate.
1224        fn record_error<R: Runtime>(
1225            retry_err: &mut RetryError<Box<Error>>,
1226            source: streams::Source,
1227            building: bool,
1228            mut err: Error,
1229            runtime: &R,
1230        ) {
1231            if source == streams::Source::Right {
1232                // We don't care about this error, since it is from neither a tunnel we launched
1233                // nor one that we're waiting on.
1234                return;
1235            }
1236            if !building {
1237                // We aren't building our own tunnels, so our errors are
1238                // secondary reports of other tunnels' failures.
1239                err = Error::PendingFailed(Box::new(err));
1240            }
1241            retry_err.push_timed(err, runtime.now(), Some(runtime.wallclock()));
1242        }
1243        /// Return a string describing what it means, within the context of this
1244        /// function, to have gotten an answer from `source`.
1245        fn describe_source(building: bool, source: streams::Source) -> &'static str {
1246            match (building, source) {
1247                (_, streams::Source::Right) => "optimistic advice",
1248                (true, streams::Source::Left) => "tunnel we're building",
1249                (false, streams::Source::Left) => "pending tunnel",
1250            }
1251        }
1252
1253        // Get or make a stream of futures to wait on.
1254        let (building, wait_on_stream) = match act {
1255            Action::Open(c) => {
1256                // There's already a perfectly good open tunnel; we can return
1257                // it now.
1258                trace!("Returning existing tunnel.");
1259                return Ok((c, TunnelProvenance::Preexisting));
1260            }
1261            Action::Wait(f) => {
1262                // There is one or more pending tunnel that we're waiting for.
1263                // If any succeeds, we try to use it.  If they all fail, we
1264                // fail.
1265                trace!("Waiting for tunnel.");
1266                (false, f)
1267            }
1268            Action::Build(plans) => {
1269                // We're going to launch one or more tunnels in parallel.  We
1270                // report success if any succeeds, and failure of they all fail.
1271                trace!("Building new tunnel.");
1272                let futures = FuturesUnordered::new();
1273                for plan in plans {
1274                    let self_clone = Arc::clone(&self);
1275                    // (This is where we actually launch tunnels.)
1276                    futures.push(self_clone.spawn_launch(usage, plan));
1277                }
1278                (true, futures)
1279            }
1280        };
1281
1282        // Insert ourself into the list of pending requests, and make a
1283        // stream for us to listen on for notification from pending tunnels
1284        // other than those we are pending on.
1285        let (pending_request, additional_stream) = {
1286            // We don't want this queue to participate in memory quota tracking.
1287            // There isn't any tunnel yet, so there wouldn't be anything to account it to.
1288            // If this queue has the oldest data, probably the whole system is badly broken.
1289            // Tearing down the whole tunnel manager won't help.
1290            let (send, recv) = mpsc_channel_no_memquota(8);
1291            let pending = Arc::new(PendingRequest {
1292                usage: usage.clone(),
1293                notify: send,
1294            });
1295
1296            let mut list = self.tunnels.lock().expect("poisoned lock");
1297            list.add_pending_request(&pending);
1298
1299            (pending, recv)
1300        };
1301
1302        // We use our "select_biased" stream combiner here to ensure that:
1303        //   1) Circuits from wait_on_stream (the ones we're pending on) are
1304        //      preferred.
1305        //   2) We exit this function when those tunnels are exhausted.
1306        //   3) We still get notified about other tunnels that might meet our
1307        //      interests.
1308        //
1309        // The events from Left stream are the oes that we explicitly asked for,
1310        // so we'll treat errors there as real problems.  The events from the
1311        // Right stream are ones that we got opportunistically told about; it's
1312        // not a big deal if those fail.
1313        let mut incoming = streams::select_biased(wait_on_stream, additional_stream.map(Ok));
1314
1315        let mut retry_error = RetryError::in_attempt_to("wait for tunnels");
1316
1317        while let Some((src, id)) = incoming.next().await {
1318            match id {
1319                Ok(Ok(ref id)) => {
1320                    // Great, we have a tunnel . See if we can use it!
1321                    let mut list = self.tunnels.lock().expect("poisoned lock");
1322                    if let Some(ent) = list.get_open_mut(id) {
1323                        let now = self.runtime.now();
1324                        match ent.restrict_mut(usage, now) {
1325                            Ok(()) => {
1326                                // Great, this will work.  We drop the
1327                                // pending request now explicitly to remove
1328                                // it from the list.
1329                                drop(pending_request);
1330                                if matches!(ent.expiration, ExpirationInfo::Unused { .. }) {
1331                                    let try_to_expire_after = if ent.spec.is_long_lived() {
1332                                        self.circuit_timing().disused_circuit_timeout
1333                                    } else {
1334                                        self.circuit_timing().max_dirtiness
1335                                    };
1336                                    // Since this tunnel hasn't been used yet, schedule expiration
1337                                    // task after `max_dirtiness` from now.
1338                                    spawn_expiration_task(
1339                                        &self.runtime,
1340                                        Arc::downgrade(&self),
1341                                        ent.tunnel.id(),
1342                                        now + try_to_expire_after,
1343                                    );
1344                                }
1345                                return Ok((ent.tunnel.clone(), TunnelProvenance::NewlyCreated));
1346                            }
1347                            Err(e) => {
1348                                // In this case, a `UsageMismatched` error just means that we lost the race
1349                                // to restrict this tunnel.
1350                                let e = match e {
1351                                    Error::UsageMismatched(e) => Error::LostUsabilityRace(e),
1352                                    x => x,
1353                                };
1354                                if src == streams::Source::Left {
1355                                    info_report!(
1356                                        &e,
1357                                        "{} suggested we use {:?}, but restrictions failed",
1358                                        describe_source(building, src),
1359                                        id,
1360                                    );
1361                                } else {
1362                                    debug_report!(
1363                                        &e,
1364                                        "{} suggested we use {:?}, but restrictions failed",
1365                                        describe_source(building, src),
1366                                        id,
1367                                    );
1368                                }
1369                                record_error(&mut retry_error, src, building, e, &self.runtime);
1370                                continue;
1371                            }
1372                        }
1373                    }
1374                }
1375                Ok(Err(ref e)) => {
1376                    debug!("{} sent error {:?}", describe_source(building, src), e);
1377                    record_error(&mut retry_error, src, building, e.clone(), &self.runtime);
1378                }
1379                Err(oneshot::Canceled) => {
1380                    debug!(
1381                        "{} went away (Canceled), quitting take_action right away",
1382                        describe_source(building, src)
1383                    );
1384                    record_error(
1385                        &mut retry_error,
1386                        src,
1387                        building,
1388                        Error::PendingCanceled,
1389                        &self.runtime,
1390                    );
1391                    return Err(retry_error);
1392                }
1393            }
1394
1395            debug!(
1396                "While waiting on tunnel: {:?} from {}",
1397                id,
1398                describe_source(building, src)
1399            );
1400        }
1401
1402        // Nothing worked.  We drop the pending request now explicitly
1403        // to remove it from the list.  (We could just let it get dropped
1404        // implicitly, but that's a bit confusing.)
1405        drop(pending_request);
1406
1407        Err(retry_error)
1408    }
1409
1410    /// Given a directory and usage, compute the necessary objects to
1411    /// build a tunnel: A [`PendingEntry`] to keep track of the in-process
1412    /// tunnel, and a [`TunnelBuildPlan`] that we'll give to the thread
1413    /// that will build the tunnel.
1414    ///
1415    /// The caller should probably add the resulting `PendingEntry` to
1416    /// `self.circs`.
1417    ///
1418    /// This is an internal function that we call when we're pretty sure
1419    /// we want to build a tunnel.
1420    #[allow(clippy::type_complexity)]
1421    fn plan_by_usage(
1422        &self,
1423        dir: DirInfo<'_>,
1424        usage: &TargetTunnelUsage,
1425    ) -> Result<(Arc<PendingEntry<B, R>>, TunnelBuildPlan<B, R>)> {
1426        let (plan, bspec) = self.builder.plan_tunnel(usage, dir)?;
1427        let (pending, sender) = PendingEntry::new(&bspec);
1428        let pending = Arc::new(pending);
1429
1430        let plan = TunnelBuildPlan {
1431            plan,
1432            sender,
1433            pending: Arc::clone(&pending),
1434        };
1435
1436        Ok((pending, plan))
1437    }
1438
1439    /// Launch a managed tunnel for a target usage, without checking
1440    /// whether one already exists or is pending.
1441    ///
1442    /// Return a listener that will be informed when the tunnel is done.
1443    #[instrument(level = "trace", skip_all)]
1444    pub(crate) fn launch_by_usage(
1445        self: &Arc<Self>,
1446        usage: &TargetTunnelUsage,
1447        dir: DirInfo<'_>,
1448    ) -> Result<Shared<oneshot::Receiver<PendResult<B, R>>>> {
1449        let (pending, plan) = self.plan_by_usage(dir, usage)?;
1450
1451        self.tunnels
1452            .lock()
1453            .expect("Poisoned lock for tunnel list")
1454            .add_pending_tunnel(pending);
1455
1456        Ok(Arc::clone(self).spawn_launch(usage, plan))
1457    }
1458
1459    /// Spawn a background task to launch a tunnel, and report its status.
1460    ///
1461    /// The `usage` argument is the usage from the original request that made
1462    /// us build this tunnel.
1463    #[instrument(level = "trace", skip_all)]
1464    fn spawn_launch(
1465        self: Arc<Self>,
1466        usage: &TargetTunnelUsage,
1467        plan: TunnelBuildPlan<B, R>,
1468    ) -> Shared<oneshot::Receiver<PendResult<B, R>>> {
1469        let TunnelBuildPlan {
1470            mut plan,
1471            sender,
1472            pending,
1473        } = plan;
1474        let request_loyalty = self.circuit_timing().request_loyalty;
1475
1476        let wait_on_future = pending.receiver.clone();
1477        let runtime = self.runtime.clone();
1478        let runtime_copy = self.runtime.clone();
1479
1480        let tid = rand::random::<u64>();
1481        // We release this block when the tunnel builder task terminates.
1482        let reason = format!("tunnel builder task {}", tid);
1483        runtime.block_advance(reason.clone());
1484        // During tests, the `FakeBuilder` will need to release the block in order to fake a timeout
1485        // correctly.
1486        plan.add_blocked_advance_reason(reason);
1487
1488        let usage = usage.clone();
1489
1490        runtime
1491            .spawn(async move {
1492                let self_clone = Arc::clone(&self);
1493                let future = AssertUnwindSafe(self_clone.do_launch(plan, pending)).catch_unwind();
1494                let (new_spec, reply) = match future.await {
1495                    Ok(x) => x, // Success or regular failure
1496                    Err(e) => {
1497                        // Okay, this is a panic.  We have to tell the calling
1498                        // thread about it, then exit this tunnel builder task.
1499                        let _ = sender.send(Err(internal!("tunnel build task panicked").into()));
1500                        std::panic::panic_any(e);
1501                    }
1502                };
1503
1504                // Tell anybody who was listening about it that this
1505                // tunnel is now usable or failed.
1506                //
1507                // (We ignore any errors from `send`: That just means that nobody
1508                // was waiting for this tunnel.)
1509                let _ = sender.send(reply.clone());
1510
1511                // TODO: Give usage as Value, probably once tracing valuable feature is stable.
1512                if reply.is_ok() {
1513                    tracing::trace!(
1514                        onionperf = true,
1515                        usage = ?usage,
1516                        event = ?OnionperfEvent::Circuit(OnionperfCircuitStatus::Launched),
1517                    );
1518                } else {
1519                    tracing::trace!(
1520                        onionperf = true,
1521                        usage = ?usage,
1522                        event = ?OnionperfEvent::Circuit(OnionperfCircuitStatus::Failed),
1523                    );
1524                }
1525
1526                if let Some(new_spec) = new_spec {
1527                    // Wait briefly before we notify opportunistically.  This
1528                    // delay will give the tunnels that were originally
1529                    // specifically intended for a request a little more time
1530                    // to finish, before we offer it this tunnel instead.
1531                    let sl = runtime_copy.sleep(request_loyalty);
1532                    runtime_copy.allow_one_advance(request_loyalty);
1533                    sl.await;
1534
1535                    let pending = {
1536                        let list = self.tunnels.lock().expect("poisoned lock");
1537                        list.find_pending_requests(&new_spec)
1538                    };
1539                    for pending_request in pending {
1540                        let _ = pending_request.notify.clone().try_send(reply.clone());
1541                    }
1542                }
1543                runtime_copy.release_advance(format!("tunnel builder task {}", tid));
1544            })
1545            .expect("Couldn't spawn tunnel-building task");
1546
1547        wait_on_future
1548    }
1549
1550    /// Run in the background to launch a tunnel. Return a 2-tuple of the new
1551    /// tunnel spec and the outcome that should be sent to the initiator.
1552    #[instrument(level = "trace", skip_all)]
1553    async fn do_launch(
1554        self: Arc<Self>,
1555        plan: <B as AbstractTunnelBuilder<R>>::Plan,
1556        pending: Arc<PendingEntry<B, R>>,
1557    ) -> (Option<SupportedTunnelUsage>, PendResult<B, R>) {
1558        let outcome = self.builder.build_tunnel(plan).await;
1559
1560        match outcome {
1561            Err(e) => (None, Err(e)),
1562            Ok((new_spec, tunnel)) => {
1563                let id = tunnel.id();
1564
1565                let use_duration = self.pick_use_duration();
1566                let now = self.runtime.now();
1567                let exp_inst = now + use_duration;
1568                let runtime_copy = self.runtime.clone();
1569                spawn_expiration_task(&runtime_copy, Arc::downgrade(&self), tunnel.id(), exp_inst);
1570                // I used to call restrict_mut here, but now I'm not so
1571                // sure. Doing restrict_mut makes sure that this
1572                // tunnel will be suitable for the request that asked
1573                // for us in the first place, but that should be
1574                // ensured anyway by our tracking its tentative
1575                // assignment.
1576                //
1577                // new_spec.restrict_mut(&usage_copy).unwrap();
1578                let use_before = ExpirationInfo::new(now);
1579                let open_ent = OpenEntry::new(new_spec.clone(), tunnel, use_before);
1580                {
1581                    let mut list = self.tunnels.lock().expect("poisoned lock");
1582                    // Finally, before we return this tunnel, we need to make
1583                    // sure that this pending tunnel is still pending.  (If it
1584                    // is not pending, then it was cancelled through a call to
1585                    // `retire_all_tunnels`, and the configuration that we used
1586                    // to launch it is now sufficiently outdated that we should
1587                    // no longer give this tunnel to a client.)
1588                    if list.tunnel_is_pending(&pending) {
1589                        list.add_open(open_ent);
1590                        // We drop our reference to 'pending' here:
1591                        // this should make all the weak references to
1592                        // the `PendingEntry` become dangling.
1593                        drop(pending);
1594                        (Some(new_spec), Ok(id))
1595                    } else {
1596                        // This tunnel is no longer pending! It must have been cancelled, probably
1597                        // by a call to retire_all_tunnels()
1598                        drop(pending); // ibid
1599                        (None, Err(Error::CircCanceled))
1600                    }
1601                }
1602            }
1603        }
1604    }
1605
1606    /// Return the currently configured expiration parameters.
1607    fn expiration_params(&self) -> ExpirationParameters {
1608        let expire_unused_after = self.pick_use_duration();
1609        let expire_dirty_after = self.circuit_timing().max_dirtiness;
1610        let expire_disused_after = self.circuit_timing().disused_circuit_timeout;
1611
1612        ExpirationParameters {
1613            expire_unused_after,
1614            expire_dirty_after,
1615            expire_disused_after,
1616        }
1617    }
1618
1619    /// Plan and launch a new tunnel to a given target, bypassing our managed
1620    /// pool of tunnels.
1621    ///
1622    /// This method will always return a new tunnel, and never return a tunnel
1623    /// that this CircMgr gives out for anything else.
1624    ///
1625    /// The new tunnel will participate in the guard and timeout apparatus as
1626    /// appropriate, no retry attempt will be made if the tunnel fails.
1627    #[cfg(feature = "hs-common")]
1628    #[instrument(level = "trace", skip_all)]
1629    pub(crate) async fn launch_unmanaged(
1630        &self,
1631        usage: &TargetTunnelUsage,
1632        dir: DirInfo<'_>,
1633    ) -> Result<(SupportedTunnelUsage, B::Tunnel)> {
1634        let (_, plan) = self.plan_by_usage(dir, usage)?;
1635        self.builder.build_tunnel(plan.plan).await
1636    }
1637
1638    /// Remove the tunnel with a given `id` from this manager.
1639    ///
1640    /// After this function is called, that tunnel will no longer be handed
1641    /// out to any future requests.
1642    ///
1643    /// Return None if we have no tunnel with the given ID.
1644    pub(crate) fn take_tunnel(
1645        &self,
1646        id: &<B::Tunnel as AbstractTunnel>::Id,
1647    ) -> Option<Arc<B::Tunnel>> {
1648        let mut list = self.tunnels.lock().expect("poisoned lock");
1649        list.take_open(id).map(|e| e.tunnel)
1650    }
1651
1652    /// Remove all open and pending tunnels and from this manager, to ensure
1653    /// they can't be given out for any more requests.
1654    ///
1655    /// Calling `retire_all_tunnels` ensures that any tunnel request that gets
1656    /// an  answer _after this method runs_ will receive a tunnel that was
1657    /// launched _after this method runs_.
1658    ///
1659    /// We call this method this when our configuration changes in such a way
1660    /// that we want to make sure that any new (or pending) requests will
1661    /// receive tunnels that are built using the new configuration.
1662    //
1663    // For more information, see documentation on [`CircuitList::open_circs`],
1664    // [`CircuitList::pending_circs`], and comments in `do_launch`.
1665    pub(crate) fn retire_all_tunnels(&self) {
1666        let mut list = self.tunnels.lock().expect("poisoned lock");
1667        list.clear_all_tunnels();
1668    }
1669
1670    /// Expire tunnels according to the rules in `config` and the
1671    /// current time `now`.
1672    ///
1673    /// Expired tunnels will not be automatically closed, but they will
1674    /// no longer be given out for new tunnels.
1675    ///
1676    /// Return the earliest time at which any current tunnel will expire.
1677    pub(crate) async fn expire_tunnels(&self, now: Instant) -> Option<Instant> {
1678        let expiration_params = self.expiration_params();
1679
1680        // While holding the lock, we call TunnelList::expire_tunnels.
1681        // That function will expire what it can, and return a list of the tunnels for which
1682        // we need to call `disused_since`.
1683        let (mut earliest_expiration, need_to_check) = {
1684            let mut list = self.tunnels.lock().expect("poisoned lock");
1685            list.expire_tunnels(now, &expiration_params)
1686        };
1687
1688        // Now we've dropped the lock, and can do async checks.
1689        let mut last_known_usage = Vec::new();
1690        for tunnel in need_to_check {
1691            let Some(tunnel) = Weak::upgrade(&tunnel) else {
1692                continue; // The tunnel is already gone.
1693            };
1694            last_known_usage.push((tunnel.id(), tunnel.last_known_to_be_used_at().await));
1695        }
1696
1697        // Now get the lock again, and tell the list what we learned.
1698        //
1699        // Note that if this function is called twice simultaneously, in some corner cases, we might
1700        // decide to expire something twice.  That's okay.
1701        {
1702            let mut list = self.tunnels.lock().expect("poisoned lock");
1703            for (id, disused_since) in last_known_usage {
1704                match list.update_long_lived_tunnel_last_used(
1705                    &id,
1706                    now,
1707                    &expiration_params,
1708                    &disused_since,
1709                ) {
1710                    Ok(Some(may_expire)) => {
1711                        earliest_expiration = match earliest_expiration {
1712                            Some(exp) if exp < may_expire => Some(exp),
1713                            _ => Some(may_expire),
1714                        };
1715                    }
1716                    Ok(None) => {}
1717                    Err(e) => warn_report!(e, "Error while updating status on tunnel {:?}", id),
1718                }
1719            }
1720        }
1721
1722        earliest_expiration
1723    }
1724
1725    /// Consider expiring the tunnel with given tunnel `id`,
1726    /// according to the rules in `config` and the current time `now`.
1727    ///
1728    /// Returns None if the circuit is expired; otherwise returns the next time at which the circuit may expire.
1729    pub(crate) async fn consider_expiring_tunnel(
1730        &self,
1731        tun_id: &<B::Tunnel as AbstractTunnel>::Id,
1732        now: Instant,
1733    ) -> Result<Option<Instant>> {
1734        let expiration_params = self.expiration_params();
1735
1736        // With the lock, call TunneList::tunnel_should_expire, and expire it (or don't)
1737        // if the decision is obvious.
1738        let tunnel = {
1739            let mut list: sync::MutexGuard<'_, TunnelList<B, R>> =
1740                self.tunnels.lock().expect("poisoned lock");
1741            let Some(should_expire) = list.tunnel_should_expire(tun_id, now, &expiration_params)
1742            else {
1743                return Ok(None);
1744            };
1745            match should_expire {
1746                ShouldExpire::Now => {
1747                    let _discard = list.take_open(tun_id);
1748                    return Ok(None);
1749                }
1750                ShouldExpire::NotBefore(t) => return Ok(Some(t)),
1751                ShouldExpire::PossiblyNow => {
1752                    let Some(tunnel_ent) = list.get_open_mut(tun_id) else {
1753                        return Ok(None);
1754                    };
1755                    Arc::clone(&tunnel_ent.tunnel)
1756                }
1757            }
1758        };
1759
1760        // If we get here, then we have a long-lived tunnel for which we need to check `disused_since`
1761        let last_known_in_use_at = tunnel.last_known_to_be_used_at().await;
1762
1763        // Now we tell the TunnelList what we learned.
1764        {
1765            let mut list: sync::MutexGuard<'_, TunnelList<B, R>> =
1766                self.tunnels.lock().expect("poisoned lock");
1767            list.update_long_lived_tunnel_last_used(
1768                tun_id,
1769                now,
1770                &expiration_params,
1771                &last_known_in_use_at,
1772            )
1773        }
1774    }
1775
1776    /// Return the number of open tunnels held by this tunnel manager.
1777    pub(crate) fn n_tunnels(&self) -> usize {
1778        let list = self.tunnels.lock().expect("poisoned lock");
1779        list.open_tunnels.len()
1780    }
1781
1782    /// Return the number of pending tunnels tracked by this tunnel manager.
1783    #[cfg(test)]
1784    pub(crate) fn n_pending_tunnels(&self) -> usize {
1785        let list = self.tunnels.lock().expect("poisoned lock");
1786        list.pending_tunnels.len()
1787    }
1788
1789    /// Get a reference to this manager's runtime.
1790    pub(crate) fn peek_runtime(&self) -> &R {
1791        &self.runtime
1792    }
1793
1794    /// Get a reference to this manager's builder.
1795    pub(crate) fn peek_builder(&self) -> &B {
1796        &self.builder
1797    }
1798
1799    /// Pick a duration by when a new tunnel should expire from now
1800    /// if it has not yet been used
1801    fn pick_use_duration(&self) -> Duration {
1802        let timings = self
1803            .unused_timing
1804            .lock()
1805            .expect("Poisoned lock for unused_timing");
1806
1807        if self.builder.learning_timeouts() {
1808            timings.learning
1809        } else {
1810            // TODO: In Tor, this calculation also depends on
1811            // stuff related to predicted ports and channel
1812            // padding.
1813            use tor_basic_utils::RngExt as _;
1814            let mut rng = rand::rng();
1815            rng.gen_range_checked(timings.not_learning..=timings.not_learning * 2)
1816                .expect("T .. 2x T turned out to be an empty duration range?!")
1817        }
1818    }
1819}
1820
1821/// Spawn an expiration task that expires a tunnel at given instant.
1822///
1823/// When the timeout occurs, if the tunnel manager is still present,
1824/// the task will ask the manager to expire the tunnel, if the tunnel
1825/// is ready to expire.
1826//
1827// TODO: It would be good to do away with this function entirely, and have a smarter expiration
1828// function.  This one only exists because there is not an "expire some circuits" background task.
1829fn spawn_expiration_task<B, R>(
1830    runtime: &R,
1831    circmgr: Weak<AbstractTunnelMgr<B, R>>,
1832    circ_id: <<B as AbstractTunnelBuilder<R>>::Tunnel as AbstractTunnel>::Id,
1833    exp_inst: Instant,
1834) where
1835    R: Runtime,
1836    B: 'static + AbstractTunnelBuilder<R>,
1837{
1838    let now = runtime.now();
1839    let rt_copy = runtime.clone();
1840    let mut duration = exp_inst.saturating_duration_since(now);
1841
1842    // NOTE: Once there was an optimization here that ran the expiration immediately if
1843    // `duration` was zero.
1844    // I discarded that optimization when I made `consider_expiring_tunnel` async,
1845    // since we really want this function _not_ to be async,
1846    // because we run it in contexts where we hold a Mutex on the tunnel list.
1847
1848    // Spawn a timer expiration task with given expiration instant.
1849    if let Err(e) = runtime.spawn(async move {
1850        loop {
1851            rt_copy.sleep(duration).await;
1852            let cm = if let Some(cm) = Weak::upgrade(&circmgr) {
1853                cm
1854            } else {
1855                return;
1856            };
1857            match cm.consider_expiring_tunnel(&circ_id, exp_inst).await {
1858                Ok(None) => return,
1859                Ok(Some(when)) => {
1860                    duration = when.saturating_duration_since(rt_copy.now());
1861                }
1862                Err(e) => {
1863                    warn_report!(
1864                        e,
1865                        "Error while considering expiration for tunnel {:?}",
1866                        circ_id
1867                    );
1868                    return;
1869                }
1870            }
1871        }
1872    }) {
1873        warn_report!(e, "Unable to launch expiration task");
1874    }
1875}
1876
1877#[cfg(test)]
1878mod test {
1879    // @@ begin test lint list maintained by maint/add_warning @@
1880    #![allow(clippy::bool_assert_comparison)]
1881    #![allow(clippy::clone_on_copy)]
1882    #![allow(clippy::dbg_macro)]
1883    #![allow(clippy::mixed_attributes_style)]
1884    #![allow(clippy::print_stderr)]
1885    #![allow(clippy::print_stdout)]
1886    #![allow(clippy::single_char_pattern)]
1887    #![allow(clippy::unwrap_used)]
1888    #![allow(clippy::unchecked_time_subtraction)]
1889    #![allow(clippy::useless_vec)]
1890    #![allow(clippy::needless_pass_by_value)]
1891    #![allow(clippy::string_slice)] // See arti#2571
1892    //! <!-- @@ end test lint list maintained by maint/add_warning @@ -->
1893    use super::*;
1894    use crate::isolation::test::{IsolationTokenEq, assert_isoleq};
1895    use crate::mocks::{FakeBuilder, FakeCirc, FakeId, FakeOp};
1896    use crate::usage::{ExitPolicy, SupportedTunnelUsage};
1897    use crate::{
1898        Error, IsolationToken, StreamIsolation, TargetPort, TargetPorts, TargetTunnelUsage,
1899    };
1900    use std::sync::LazyLock;
1901    use tor_dircommon::fallback::FallbackList;
1902    use tor_guardmgr::TestConfig;
1903    use tor_llcrypto::pk::ed25519::Ed25519Identity;
1904    use tor_netdir::testnet;
1905    use tor_persist::TestingStateMgr;
1906    use tor_rtcompat::SleepProvider;
1907    use tor_rtmock::MockRuntime;
1908    use web_time_compat::InstantExt;
1909
1910    #[allow(deprecated)] // TODO #1885
1911    use tor_rtmock::MockSleepRuntime;
1912
1913    static FALLBACKS_EMPTY: LazyLock<FallbackList> = LazyLock::new(|| [].into());
1914
1915    fn di() -> DirInfo<'static> {
1916        (&*FALLBACKS_EMPTY).into()
1917    }
1918
1919    fn target_to_spec(target: &TargetTunnelUsage) -> SupportedTunnelUsage {
1920        match target {
1921            TargetTunnelUsage::Exit {
1922                ports,
1923                isolation,
1924                country_code,
1925                require_stability,
1926            } => SupportedTunnelUsage::Exit {
1927                policy: ExitPolicy::from_target_ports(&TargetPorts::from(&ports[..])),
1928                isolation: Some(isolation.clone()),
1929                country_code: country_code.clone(),
1930                all_relays_stable: *require_stability,
1931            },
1932            _ => unimplemented!(),
1933        }
1934    }
1935
1936    impl<U: PartialEq> IsolationTokenEq for OpenEntry<U> {
1937        fn isol_eq(&self, other: &Self) -> bool {
1938            self.spec.isol_eq(&other.spec)
1939                && self.tunnel == other.tunnel
1940                && self.expiration == other.expiration
1941        }
1942    }
1943
1944    impl<U: PartialEq> IsolationTokenEq for &mut OpenEntry<U> {
1945        fn isol_eq(&self, other: &Self) -> bool {
1946            self.spec.isol_eq(&other.spec)
1947                && self.tunnel == other.tunnel
1948                && self.expiration == other.expiration
1949        }
1950    }
1951
1952    fn make_builder<R: Runtime>(runtime: &R) -> FakeBuilder<R> {
1953        let state_mgr = TestingStateMgr::new();
1954        let guard_config = TestConfig::default();
1955        FakeBuilder::new(runtime, state_mgr, &guard_config)
1956    }
1957
1958    #[test]
1959    fn basic_tests() {
1960        MockRuntime::test_with_various(|rt| async move {
1961            #[allow(deprecated)] // TODO #1885
1962            let rt = MockSleepRuntime::new(rt);
1963
1964            let builder = make_builder(&rt);
1965
1966            let mgr = Arc::new(AbstractTunnelMgr::new(
1967                builder,
1968                rt.clone(),
1969                CircuitTiming::default(),
1970            ));
1971
1972            let webports = TargetTunnelUsage::new_from_ipv4_ports(&[80, 443]);
1973
1974            // Check initialization.
1975            assert_eq!(mgr.n_tunnels(), 0);
1976            assert!(mgr.peek_builder().script.lock().unwrap().is_empty());
1977
1978            // Launch a tunnel ; make sure we get it.
1979            let c1 = rt.wait_for(mgr.get_or_launch(&webports, di())).await;
1980            let c1 = c1.unwrap().0;
1981            assert_eq!(mgr.n_tunnels(), 1);
1982
1983            // Make sure we get the one we already made if we ask for it.
1984            let port80 = TargetTunnelUsage::new_from_ipv4_ports(&[80]);
1985            let c2 = mgr.get_or_launch(&port80, di()).await;
1986
1987            let c2 = c2.unwrap().0;
1988            assert!(FakeCirc::eq(&c1, &c2));
1989            assert_eq!(mgr.n_tunnels(), 1);
1990
1991            // Now try launching two tunnels "at once" to make sure that our
1992            // pending-tunnel code works.
1993
1994            let dnsport = TargetTunnelUsage::new_from_ipv4_ports(&[53]);
1995            let dnsport_restrict = TargetTunnelUsage::Exit {
1996                ports: vec![TargetPort::ipv4(53)],
1997                isolation: StreamIsolation::builder().build().unwrap(),
1998                country_code: None,
1999                require_stability: false,
2000            };
2001
2002            let (c3, c4) = rt
2003                .wait_for(futures::future::join(
2004                    mgr.get_or_launch(&dnsport, di()),
2005                    mgr.get_or_launch(&dnsport_restrict, di()),
2006                ))
2007                .await;
2008
2009            let c3 = c3.unwrap().0;
2010            let c4 = c4.unwrap().0;
2011            assert!(!FakeCirc::eq(&c1, &c3));
2012            assert!(FakeCirc::eq(&c3, &c4));
2013            assert_eq!(c3.id(), c4.id());
2014            assert_eq!(mgr.n_tunnels(), 2);
2015
2016            // Now we're going to remove c3 from consideration.  It's the
2017            // same as c4, so removing c4 will give us None.
2018            let c3_taken = mgr.take_tunnel(&c3.id()).unwrap();
2019            let now_its_gone = mgr.take_tunnel(&c4.id());
2020            assert!(FakeCirc::eq(&c3_taken, &c3));
2021            assert!(now_its_gone.is_none());
2022            assert_eq!(mgr.n_tunnels(), 1);
2023
2024            // Having removed them, let's launch another dnsport and make
2025            // sure we get a different tunnel.
2026            let c5 = rt.wait_for(mgr.get_or_launch(&dnsport, di())).await;
2027            let c5 = c5.unwrap().0;
2028            assert!(!FakeCirc::eq(&c3, &c5));
2029            assert!(!FakeCirc::eq(&c4, &c5));
2030            assert_eq!(mgr.n_tunnels(), 2);
2031
2032            // Now try launch_by_usage.
2033            let prev = mgr.n_pending_tunnels();
2034            assert!(mgr.launch_by_usage(&dnsport, di()).is_ok());
2035            assert_eq!(mgr.n_pending_tunnels(), prev + 1);
2036            // TODO: Actually make sure that launch_by_usage launched
2037            // the right thing.
2038        });
2039    }
2040
2041    #[test]
2042    fn request_timeout() {
2043        MockRuntime::test_with_various(|rt| async move {
2044            #[allow(deprecated)] // TODO #1885
2045            let rt = MockSleepRuntime::new(rt);
2046
2047            let ports = TargetTunnelUsage::new_from_ipv4_ports(&[80, 443]);
2048
2049            // This will fail once, and then completely time out.  The
2050            // result will be a failure.
2051            let builder = make_builder(&rt);
2052            builder.set(&ports, vec![FakeOp::Fail, FakeOp::Timeout]);
2053
2054            let mgr = Arc::new(AbstractTunnelMgr::new(
2055                builder,
2056                rt.clone(),
2057                CircuitTiming::default(),
2058            ));
2059            let c1 = mgr
2060                .peek_runtime()
2061                .wait_for(mgr.get_or_launch(&ports, di()))
2062                .await;
2063
2064            assert!(matches!(c1, Err(Error::RequestFailed(_))));
2065        });
2066    }
2067
2068    #[test]
2069    fn request_timeout2() {
2070        MockRuntime::test_with_various(|rt| async move {
2071            #[allow(deprecated)] // TODO #1885
2072            let rt = MockSleepRuntime::new(rt);
2073
2074            // Now try a more complicated case: we'll try to get things so
2075            // that we wait for a little over our predicted time because
2076            // of our wait-for-next-action logic.
2077            let ports = TargetTunnelUsage::new_from_ipv4_ports(&[80, 443]);
2078            let builder = make_builder(&rt);
2079            builder.set(
2080                &ports,
2081                vec![
2082                    FakeOp::Delay(Duration::from_millis(60_000 - 25)),
2083                    FakeOp::NoPlan,
2084                ],
2085            );
2086
2087            let mgr = Arc::new(AbstractTunnelMgr::new(
2088                builder,
2089                rt.clone(),
2090                CircuitTiming::default(),
2091            ));
2092            let c1 = mgr
2093                .peek_runtime()
2094                .wait_for(mgr.get_or_launch(&ports, di()))
2095                .await;
2096
2097            assert!(matches!(c1, Err(Error::RequestFailed(_))));
2098        });
2099    }
2100
2101    #[test]
2102    fn request_unplannable() {
2103        MockRuntime::test_with_various(|rt| async move {
2104            #[allow(deprecated)] // TODO #1885
2105            let rt = MockSleepRuntime::new(rt);
2106
2107            let ports = TargetTunnelUsage::new_from_ipv4_ports(&[80, 443]);
2108
2109            // This will fail a the planning stages, a lot.
2110            let builder = make_builder(&rt);
2111            builder.set(&ports, vec![FakeOp::NoPlan; 2000]);
2112
2113            let mgr = Arc::new(AbstractTunnelMgr::new(
2114                builder,
2115                rt.clone(),
2116                CircuitTiming::default(),
2117            ));
2118            let c1 = rt.wait_for(mgr.get_or_launch(&ports, di())).await;
2119
2120            assert!(matches!(c1, Err(Error::RequestFailed(_))));
2121        });
2122    }
2123
2124    #[test]
2125    fn request_fails_too_much() {
2126        MockRuntime::test_with_various(|rt| async move {
2127            #[allow(deprecated)] // TODO #1885
2128            let rt = MockSleepRuntime::new(rt);
2129            let ports = TargetTunnelUsage::new_from_ipv4_ports(&[80, 443]);
2130
2131            // This will fail 1000 times, which is above the retry limit.
2132            let builder = make_builder(&rt);
2133            builder.set(&ports, vec![FakeOp::Fail; 1000]);
2134
2135            let mgr = Arc::new(AbstractTunnelMgr::new(
2136                builder,
2137                rt.clone(),
2138                CircuitTiming::default(),
2139            ));
2140            let c1 = rt.wait_for(mgr.get_or_launch(&ports, di())).await;
2141
2142            assert!(matches!(c1, Err(Error::RequestFailed(_))));
2143        });
2144    }
2145
2146    #[test]
2147    fn request_wrong_spec() {
2148        MockRuntime::test_with_various(|rt| async move {
2149            #[allow(deprecated)] // TODO #1885
2150            let rt = MockSleepRuntime::new(rt);
2151            let ports = TargetTunnelUsage::new_from_ipv4_ports(&[80, 443]);
2152
2153            // The first time this is called, it will build a tunnel
2154            // with the wrong spec.  (A tunnel builder should never
2155            // actually _do_ that, but it's something we code for.)
2156            let builder = make_builder(&rt);
2157            builder.set(
2158                &ports,
2159                vec![FakeOp::WrongSpec(target_to_spec(
2160                    &TargetTunnelUsage::new_from_ipv4_ports(&[22]),
2161                ))],
2162            );
2163
2164            let mgr = Arc::new(AbstractTunnelMgr::new(
2165                builder,
2166                rt.clone(),
2167                CircuitTiming::default(),
2168            ));
2169            let c1 = rt.wait_for(mgr.get_or_launch(&ports, di())).await;
2170
2171            assert!(c1.is_ok());
2172        });
2173    }
2174
2175    #[test]
2176    fn request_retried() {
2177        MockRuntime::test_with_various(|rt| async move {
2178            #[allow(deprecated)] // TODO #1885
2179            let rt = MockSleepRuntime::new(rt);
2180            let ports = TargetTunnelUsage::new_from_ipv4_ports(&[80, 443]);
2181
2182            // This will fail twice, and then succeed. The result will be
2183            // a success.
2184            let builder = make_builder(&rt);
2185            builder.set(&ports, vec![FakeOp::Fail, FakeOp::Fail]);
2186
2187            let mgr = Arc::new(AbstractTunnelMgr::new(
2188                builder,
2189                rt.clone(),
2190                CircuitTiming::default(),
2191            ));
2192
2193            // This test doesn't exercise any timeout behaviour.
2194            rt.block_advance("test doesn't require advancing");
2195
2196            let (c1, c2) = rt
2197                .wait_for(futures::future::join(
2198                    mgr.get_or_launch(&ports, di()),
2199                    mgr.get_or_launch(&ports, di()),
2200                ))
2201                .await;
2202
2203            let c1 = c1.unwrap().0;
2204            let c2 = c2.unwrap().0;
2205
2206            assert!(FakeCirc::eq(&c1, &c2));
2207        });
2208    }
2209
2210    #[test]
2211    fn isolated() {
2212        MockRuntime::test_with_various(|rt| async move {
2213            #[allow(deprecated)] // TODO #1885
2214            let rt = MockSleepRuntime::new(rt);
2215            let builder = make_builder(&rt);
2216            let mgr = Arc::new(AbstractTunnelMgr::new(
2217                builder,
2218                rt.clone(),
2219                CircuitTiming::default(),
2220            ));
2221
2222            // Set our isolation so that iso1 and iso2 can't share a tunnel,
2223            // but no_iso can share a tunnel with either.
2224            let iso1 = TargetTunnelUsage::Exit {
2225                ports: vec![TargetPort::ipv4(443)],
2226                isolation: StreamIsolation::builder()
2227                    .owner_token(IsolationToken::new())
2228                    .build()
2229                    .unwrap(),
2230                country_code: None,
2231                require_stability: false,
2232            };
2233            let iso2 = TargetTunnelUsage::Exit {
2234                ports: vec![TargetPort::ipv4(443)],
2235                isolation: StreamIsolation::builder()
2236                    .owner_token(IsolationToken::new())
2237                    .build()
2238                    .unwrap(),
2239                country_code: None,
2240                require_stability: false,
2241            };
2242            let no_iso1 = TargetTunnelUsage::new_from_ipv4_ports(&[443]);
2243            let no_iso2 = no_iso1.clone();
2244
2245            // We're going to try launching these tunnels in 24 different
2246            // orders, to make sure that the outcome is correct each time.
2247            use itertools::Itertools;
2248            let timeouts: Vec<_> = [0_u64, 2, 4, 6]
2249                .iter()
2250                .map(|d| Duration::from_millis(*d))
2251                .collect();
2252
2253            for delays in timeouts.iter().permutations(4) {
2254                let d1 = delays[0];
2255                let d2 = delays[1];
2256                let d3 = delays[2];
2257                let d4 = delays[2];
2258                let (c_iso1, c_iso2, c_no_iso1, c_no_iso2) = rt
2259                    .wait_for(futures::future::join4(
2260                        async {
2261                            rt.sleep(*d1).await;
2262                            mgr.get_or_launch(&iso1, di()).await
2263                        },
2264                        async {
2265                            rt.sleep(*d2).await;
2266                            mgr.get_or_launch(&iso2, di()).await
2267                        },
2268                        async {
2269                            rt.sleep(*d3).await;
2270                            mgr.get_or_launch(&no_iso1, di()).await
2271                        },
2272                        async {
2273                            rt.sleep(*d4).await;
2274                            mgr.get_or_launch(&no_iso2, di()).await
2275                        },
2276                    ))
2277                    .await;
2278
2279                let c_iso1 = c_iso1.unwrap().0;
2280                let c_iso2 = c_iso2.unwrap().0;
2281                let c_no_iso1 = c_no_iso1.unwrap().0;
2282                let c_no_iso2 = c_no_iso2.unwrap().0;
2283
2284                assert!(!FakeCirc::eq(&c_iso1, &c_iso2));
2285                assert!(!FakeCirc::eq(&c_iso1, &c_no_iso1));
2286                assert!(!FakeCirc::eq(&c_iso1, &c_no_iso2));
2287                assert!(!FakeCirc::eq(&c_iso2, &c_no_iso1));
2288                assert!(!FakeCirc::eq(&c_iso2, &c_no_iso2));
2289                assert!(FakeCirc::eq(&c_no_iso1, &c_no_iso2));
2290            }
2291        });
2292    }
2293
2294    #[test]
2295    fn opportunistic() {
2296        MockRuntime::test_with_various(|rt| async move {
2297            #[allow(deprecated)] // TODO #1885
2298            let rt = MockSleepRuntime::new(rt);
2299
2300            // The first request will time out completely, but we're
2301            // making a second request after we launch it.  That
2302            // request should succeed, and notify the first request.
2303
2304            let ports1 = TargetTunnelUsage::new_from_ipv4_ports(&[80]);
2305            let ports2 = TargetTunnelUsage::new_from_ipv4_ports(&[80, 443]);
2306
2307            let builder = make_builder(&rt);
2308            builder.set(&ports1, vec![FakeOp::Timeout]);
2309
2310            let mgr = Arc::new(AbstractTunnelMgr::new(
2311                builder,
2312                rt.clone(),
2313                CircuitTiming::default(),
2314            ));
2315            // Note that ports2 will be wider than ports1, so the second
2316            // request will have to launch a new tunnel.
2317
2318            let (c1, c2) = rt
2319                .wait_for(futures::future::join(
2320                    mgr.get_or_launch(&ports1, di()),
2321                    async {
2322                        rt.sleep(Duration::from_millis(100)).await;
2323                        mgr.get_or_launch(&ports2, di()).await
2324                    },
2325                ))
2326                .await;
2327
2328            if let (Ok((c1, _)), Ok((c2, _))) = (c1, c2) {
2329                assert!(FakeCirc::eq(&c1, &c2));
2330            } else {
2331                panic!();
2332            };
2333        });
2334    }
2335
2336    #[test]
2337    fn prebuild() {
2338        MockRuntime::test_with_various(|rt| async move {
2339            // This time we're going to use ensure_tunnel() to make
2340            // sure that a tunnel gets built, and then launch two
2341            // other tunnels that will use it.
2342            #[allow(deprecated)] // TODO #1885
2343            let rt = MockSleepRuntime::new(rt);
2344            let builder = make_builder(&rt);
2345            let mgr = Arc::new(AbstractTunnelMgr::new(
2346                builder,
2347                rt.clone(),
2348                CircuitTiming::default(),
2349            ));
2350
2351            let ports1 = TargetTunnelUsage::new_from_ipv4_ports(&[80, 443]);
2352            let ports2 = TargetTunnelUsage::new_from_ipv4_ports(&[80]);
2353            let ports3 = TargetTunnelUsage::new_from_ipv4_ports(&[443]);
2354
2355            let ok = mgr.ensure_tunnel(&ports1, di());
2356            let (c1, c2) = rt
2357                .wait_for(futures::future::join(
2358                    async {
2359                        rt.sleep(Duration::from_millis(10)).await;
2360                        mgr.get_or_launch(&ports2, di()).await
2361                    },
2362                    async {
2363                        rt.sleep(Duration::from_millis(50)).await;
2364                        mgr.get_or_launch(&ports3, di()).await
2365                    },
2366                ))
2367                .await;
2368
2369            assert!(ok.is_ok());
2370
2371            let c1 = c1.unwrap().0;
2372            let c2 = c2.unwrap().0;
2373
2374            // If we had launched these separately, they wouldn't share
2375            // a tunnel.
2376            assert!(FakeCirc::eq(&c1, &c2));
2377        });
2378    }
2379
2380    #[test]
2381    fn expiration() {
2382        MockRuntime::test_with_various(|rt| async move {
2383            use crate::config::CircuitTimingBuilder;
2384            // Now let's make some tunnels -- one dirty, one clean, and
2385            // make sure that one expires and one doesn't.
2386            #[allow(deprecated)] // TODO #1885
2387            let rt = MockSleepRuntime::new(rt);
2388            let builder = make_builder(&rt);
2389
2390            let circuit_timing = CircuitTimingBuilder::default()
2391                .max_dirtiness(Duration::from_secs(15))
2392                .build()
2393                .unwrap();
2394
2395            let mgr = Arc::new(AbstractTunnelMgr::new(builder, rt.clone(), circuit_timing));
2396
2397            let imap = TargetTunnelUsage::new_from_ipv4_ports(&[993]);
2398            let pop = TargetTunnelUsage::new_from_ipv4_ports(&[995]);
2399
2400            let ok = mgr.ensure_tunnel(&imap, di());
2401            let pop1 = rt.wait_for(mgr.get_or_launch(&pop, di())).await;
2402
2403            assert!(ok.is_ok());
2404            let pop1 = pop1.unwrap().0;
2405
2406            rt.advance(Duration::from_secs(30)).await;
2407            rt.advance(Duration::from_secs(15)).await;
2408            let imap1 = rt.wait_for(mgr.get_or_launch(&imap, di())).await.unwrap().0;
2409
2410            // This should expire the pop tunnel, since it came from
2411            // get_or_launch() [which marks the tunnel as being
2412            // used].  It should not expire the imap tunnel, since
2413            // it was not dirty until 15 seconds after the cutoff.
2414            let now = rt.now();
2415
2416            mgr.expire_tunnels(now).await;
2417
2418            let (pop2, imap2) = rt
2419                .wait_for(futures::future::join(
2420                    mgr.get_or_launch(&pop, di()),
2421                    mgr.get_or_launch(&imap, di()),
2422                ))
2423                .await;
2424
2425            let pop2 = pop2.unwrap().0;
2426            let imap2 = imap2.unwrap().0;
2427
2428            assert!(!FakeCirc::eq(&pop2, &pop1));
2429            assert!(FakeCirc::eq(&imap2, &imap1));
2430        });
2431    }
2432
2433    /// Returns three exit policies; one that permits nothing, one that permits ports 80
2434    /// and 443 only, and one that permits all ports.
2435    fn get_exit_policies() -> (ExitPolicy, ExitPolicy, ExitPolicy) {
2436        // FIXME(eta): the below is copypasta; would be nice to have a better way of
2437        //             constructing ExitPolicy objects for testing maybe
2438        let network = testnet::construct_netdir().unwrap_if_sufficient().unwrap();
2439
2440        // Nodes with ID 0x0a through 0x13 and 0x1e through 0x27 are
2441        // exits.  Odd-numbered ones allow only ports 80 and 443;
2442        // even-numbered ones allow all ports.
2443        let id_noexit: Ed25519Identity = [0x05; 32].into();
2444        let id_webexit: Ed25519Identity = [0x11; 32].into();
2445        let id_fullexit: Ed25519Identity = [0x20; 32].into();
2446
2447        let not_exit = network.by_id(&id_noexit).unwrap();
2448        let web_exit = network.by_id(&id_webexit).unwrap();
2449        let full_exit = network.by_id(&id_fullexit).unwrap();
2450
2451        let ep_none = ExitPolicy::from_relay(&not_exit);
2452        let ep_web = ExitPolicy::from_relay(&web_exit);
2453        let ep_full = ExitPolicy::from_relay(&full_exit);
2454        (ep_none, ep_web, ep_full)
2455    }
2456
2457    #[test]
2458    fn test_find_supported() {
2459        let (ep_none, ep_web, ep_full) = get_exit_policies();
2460        let fake_circ = FakeCirc { id: FakeId::next() };
2461        let expiration = ExpirationInfo::Unused {
2462            created: Instant::get(),
2463        };
2464
2465        let mut entry_none = OpenEntry::new(
2466            SupportedTunnelUsage::Exit {
2467                policy: ep_none,
2468                isolation: None,
2469                country_code: None,
2470                all_relays_stable: true,
2471            },
2472            fake_circ.clone(),
2473            expiration.clone(),
2474        );
2475        let mut entry_none_c = entry_none.clone();
2476        let mut entry_web = OpenEntry::new(
2477            SupportedTunnelUsage::Exit {
2478                policy: ep_web,
2479                isolation: None,
2480                country_code: None,
2481                all_relays_stable: true,
2482            },
2483            fake_circ.clone(),
2484            expiration.clone(),
2485        );
2486        let mut entry_web_c = entry_web.clone();
2487        let mut entry_full = OpenEntry::new(
2488            SupportedTunnelUsage::Exit {
2489                policy: ep_full,
2490                isolation: None,
2491                country_code: None,
2492                all_relays_stable: true,
2493            },
2494            fake_circ,
2495            expiration,
2496        );
2497        let mut entry_full_c = entry_full.clone();
2498
2499        let usage_web = TargetTunnelUsage::new_from_ipv4_ports(&[80]);
2500        let empty: Vec<&mut OpenEntry<FakeCirc>> = vec![];
2501
2502        assert_isoleq!(
2503            SupportedTunnelUsage::find_supported(vec![&mut entry_none].into_iter(), &usage_web),
2504            empty
2505        );
2506
2507        // HACK(eta): We have to faff around with clones and such because
2508        //            `abstract_spec_find_supported` has a silly signature that involves `&mut`
2509        //            refs, which we can't have more than one of.
2510
2511        assert_isoleq!(
2512            SupportedTunnelUsage::find_supported(
2513                vec![&mut entry_none, &mut entry_web].into_iter(),
2514                &usage_web,
2515            ),
2516            vec![&mut entry_web_c]
2517        );
2518
2519        assert_isoleq!(
2520            SupportedTunnelUsage::find_supported(
2521                vec![&mut entry_none, &mut entry_web, &mut entry_full].into_iter(),
2522                &usage_web,
2523            ),
2524            vec![&mut entry_web_c, &mut entry_full_c]
2525        );
2526
2527        // Test preemptive tunnel usage:
2528
2529        let usage_preemptive_web = TargetTunnelUsage::Preemptive {
2530            port: Some(TargetPort::ipv4(80)),
2531            circs: 2,
2532            require_stability: false,
2533        };
2534        let usage_preemptive_dns = TargetTunnelUsage::Preemptive {
2535            port: None,
2536            circs: 2,
2537            require_stability: false,
2538        };
2539
2540        // shouldn't return anything unless there are >=2 tunnels
2541
2542        assert_isoleq!(
2543            SupportedTunnelUsage::find_supported(
2544                vec![&mut entry_none].into_iter(),
2545                &usage_preemptive_web
2546            ),
2547            empty
2548        );
2549
2550        assert_isoleq!(
2551            SupportedTunnelUsage::find_supported(
2552                vec![&mut entry_none].into_iter(),
2553                &usage_preemptive_dns
2554            ),
2555            empty
2556        );
2557
2558        assert_isoleq!(
2559            SupportedTunnelUsage::find_supported(
2560                vec![&mut entry_none, &mut entry_web].into_iter(),
2561                &usage_preemptive_web
2562            ),
2563            empty
2564        );
2565
2566        assert_isoleq!(
2567            SupportedTunnelUsage::find_supported(
2568                vec![&mut entry_none, &mut entry_web].into_iter(),
2569                &usage_preemptive_dns
2570            ),
2571            vec![&mut entry_none_c, &mut entry_web_c]
2572        );
2573
2574        assert_isoleq!(
2575            SupportedTunnelUsage::find_supported(
2576                vec![&mut entry_none, &mut entry_web, &mut entry_full].into_iter(),
2577                &usage_preemptive_web
2578            ),
2579            vec![&mut entry_web_c, &mut entry_full_c]
2580        );
2581    }
2582
2583    #[test]
2584    fn test_circlist_preemptive_target_circs() {
2585        MockRuntime::test_with_various(|rt| async move {
2586            #[allow(deprecated)] // TODO #1885
2587            let rt = MockSleepRuntime::new(rt);
2588            let netdir = testnet::construct_netdir().unwrap_if_sufficient().unwrap();
2589            let dirinfo = DirInfo::Directory(&netdir);
2590
2591            let builder = make_builder(&rt);
2592
2593            for circs in [2, 8].iter() {
2594                let mut circlist = TunnelList::<FakeBuilder<MockRuntime>, MockRuntime>::new();
2595
2596                let preemptive_target = TargetTunnelUsage::Preemptive {
2597                    port: Some(TargetPort::ipv4(80)),
2598                    circs: *circs,
2599                    require_stability: false,
2600                };
2601
2602                for _ in 0..*circs {
2603                    assert!(circlist.find_open(&preemptive_target).is_none());
2604
2605                    let usage = TargetTunnelUsage::new_from_ipv4_ports(&[80]);
2606                    let (plan, _) = builder.plan_tunnel(&usage, dirinfo).unwrap();
2607                    let (spec, circ) = rt.wait_for(builder.build_tunnel(plan)).await.unwrap();
2608                    let entry = OpenEntry::new(
2609                        spec,
2610                        circ,
2611                        ExpirationInfo::new(rt.now() + Duration::from_secs(60)),
2612                    );
2613                    circlist.add_open(entry);
2614                }
2615
2616                assert!(circlist.find_open(&preemptive_target).is_some());
2617            }
2618        });
2619    }
2620}