Skip to main content

tor_hsclient/
connect.rs

1//! Main implementation of the connection functionality
2
3use std::collections::HashMap;
4use std::fmt::Debug;
5use std::marker::PhantomData;
6use std::sync::Arc;
7
8use async_trait::async_trait;
9use educe::Educe;
10use futures::{AsyncRead, AsyncWrite};
11use itertools::Itertools;
12use rand::RngExt;
13use tor_bytes::Writeable;
14use tor_cell::relaycell::hs::IntroduceAckStatus;
15use tor_cell::relaycell::hs::intro_payload::{self, IntroduceHandshakePayload};
16use tor_cell::relaycell::hs::pow::ProofOfWork;
17use tor_cell::relaycell::msg::{AnyRelayMsg, Introduce1, Rendezvous2};
18use tor_circmgr::build::onion_circparams_from_netparams;
19use tor_circmgr::{
20    ClientOnionServiceDataTunnel, ClientOnionServiceDirTunnel, ClientOnionServiceIntroTunnel,
21};
22use tor_dirclient::SourceInfo;
23use tor_error::{Bug, debug_report, warn_report};
24use tor_hscrypto::Subcredential;
25use tor_hscrypto::time::TimePeriod;
26use tor_netdir::params::NetParameters;
27use tor_proto::TargetHop;
28use tor_proto::client::circuit::handshake::hs_ntor::{self, HsNtorHkdfKeyGenerator};
29use tracing::{debug, instrument, trace, warn};
30use web_time_compat::{Duration, Instant, SystemTime};
31
32use retry_error::RetryError;
33use safelog::{DispRedacted, Sensitive};
34use tor_cell::relaycell::RelayMsg;
35use tor_cell::relaycell::hs::{
36    AuthKeyType, EstablishRendezvous, IntroduceAck, RendezvousEstablished,
37};
38use tor_checkable::{TimeBound, timed::TimeRangeBound};
39use tor_circmgr::hspool::HsCircPool;
40use tor_circmgr::timeouts::Action as TimeoutsAction;
41use tor_dirclient::request::Requestable as _;
42use tor_error::{HasRetryTime as _, RetryTime};
43use tor_error::{internal, into_internal};
44use tor_hscrypto::RendCookie;
45use tor_hscrypto::pk::{HsBlindId, HsId, HsIdKey};
46use tor_linkspec::{CircTarget, HasRelayIds, OwnedCircTarget, RelayId};
47use tor_llcrypto::pk::ed25519::Ed25519Identity;
48use tor_netdir::{NetDir, Relay};
49use tor_netdoc::doc::hsdesc::{HsDesc, IntroPointDesc};
50use tor_proto::client::circuit::{CircParameters, handshake};
51use tor_proto::{MetaCellDisposition, MsgHandler};
52use tor_rtcompat::{Runtime, SleepProviderExt as _, TimeoutError};
53
54use crate::Config;
55use crate::caps;
56use crate::err::RendPtIdentityForError;
57use crate::pow::HsPowClient;
58use crate::proto_oneshot;
59use crate::relay_info::ipt_to_circtarget;
60use crate::state::MockableConnectorData;
61use crate::{ConnError, DescriptorError, DescriptorErrorDetail};
62use crate::{FailedAttemptError, IntroPtIndex, rend_pt_identity_for_error};
63use crate::{HsClientConnector, HsClientSecretKeys};
64
65use ConnError as CE;
66use FailedAttemptError as FAE;
67
68/// Given `R, M` where `M: MocksForConnect<M>`, expand to the mockable `ClientCirc`
69// This is quite annoying.  But the alternative is to write out `<... as // ...>`
70// each time, since otherwise the compile complains about ambiguous associated types.
71macro_rules! DataTunnel{ { $R:ty, $M:ty } => {
72    <<$M as MocksForConnect<$R>>::HsCircPool as MockableCircPool<$R>>::DataTunnel
73} }
74
75/// Information about a hidden service, including our connection history
76#[derive(Default, Educe)]
77#[educe(Debug)]
78// This type is actually crate-private, since it isn't re-exported, but it must
79// be `pub` because it appears as a default for a type parameter in HsClientConnector.
80pub struct Data {
81    /// The latest known onion service descriptor for this service.
82    desc: DataHsDesc,
83
84    /// Information about the latest status of trying to connect to this service
85    /// through each of its introduction points.
86    ipts: DataIpts,
87    /// Information about the requery period of each HsDir we have recently queried.
88    ///
89    /// Each entry represents an HsDir that we cannot requery until
90    /// its specified timestamp elapses.
91    ///
92    /// Any HsDir that does not have an entry in this map can be requeried.
93    hsdirs: DataHsDirs,
94}
95
96/// An onion service descriptor and its associated HsBlindId.
97#[derive(Debug)]
98struct HsDescForTp {
99    /// The TP this descriptor is for.
100    ///
101    /// Used for determining whether a newly fetched descriptor
102    /// is for the same time period as this one.
103    time_period: TimePeriod,
104    /// The descriptor
105    desc: TimeRangeBound<HsDesc>,
106}
107
108/// Part of `Data` that relates to our information about the HsDir requery periods
109type DataHsDirs = HashMap<RelayIdForRequeryPeriod, SystemTime>;
110
111/// Marker type, to make typed HsDir [`RelayIdFor`] keys
112#[derive(Hash, Eq, PartialEq, Ord, PartialOrd, Copy, Clone, Debug)]
113struct RequeryPeriodMap;
114
115/// Lookup key for looking up and recording our IPT use experiences
116type RelayIdForRequeryPeriod = RelayIdFor<RequeryPeriodMap>;
117
118/// Part of `Data` that relates to the HS descriptor
119type DataHsDesc = Option<HsDescForTp>;
120
121/// Part of `Data` that relates to our information about introduction points
122type DataIpts = HashMap<RelayIdForExperience, IptExperience>;
123
124/// How things went last time we tried to use this introduction point
125///
126/// Neither this data structure, nor [`Data`], is responsible for arranging that we expire this
127/// information eventually.  If we keep reconnecting to the service, we'll retain information
128/// about each IPT indefinitely, at least so long as they remain listed in the descriptors we
129/// receive.
130///
131/// Expiry of unused data is handled by `state.rs`, according to `last_used` in `ServiceState`.
132///
133/// Choosing which IPT to prefer is done by obtaining an `IptSortKey`
134/// (from this and other information).
135//
136// Don't impl Ord for IptExperience.  We obtain `Option<&IptExperience>` from our
137// data structure, and if IptExperience were Ord then Option<&IptExperience> would be Ord
138// but it would be the wrong sort order: it would always prefer None, ie untried IPTs.
139#[derive(Debug)]
140struct IptExperience {
141    /// How long it took us to get whatever outcome occurred
142    ///
143    /// We prefer fast successes to slow ones.
144    /// Then, we prefer failures with earlier `RetryTime`,
145    /// and, lastly, faster failures to slower ones.
146    duration: Duration,
147
148    /// What happened and when we might try again
149    ///
150    /// Note that we don't actually *enforce* the `RetryTime` here, just sort by it
151    /// using `RetryTime::loose_cmp`.
152    ///
153    /// We *do* return an error that is itself `HasRetryTime` and expect our callers
154    /// to honour that.
155    outcome: Result<(), RetryTime>,
156}
157
158/// Actually make a HS connection, updating our recorded state as necessary
159///
160/// `connector` is provided only for obtaining the runtime and netdir (and `mock_for_state`).
161/// Obviously, `connect` is not supposed to go looking in `services`.
162///
163/// This function handles all necessary retrying of fallible operations,
164/// (and, therefore, must also limit the total work done for a particular call).
165///
166/// This function has a minimum of functionality, since it is the boundary
167/// between "mock connection, used for testing `state.rs`" and
168/// "mock circuit and netdir, used for testing `connect.rs`",
169/// so it is not, itself, unit-testable.
170#[instrument(level = "trace", skip_all)]
171pub(crate) async fn connect<R: Runtime>(
172    connector: &HsClientConnector<R>,
173    netdir: Arc<NetDir>,
174    config: Arc<Config>,
175    hsid: HsId,
176    data: &mut Data,
177    secret_keys: HsClientSecretKeys,
178) -> Result<ClientOnionServiceDataTunnel, ConnError> {
179    Context::new(
180        &connector.runtime,
181        &*connector.circpool,
182        netdir,
183        config,
184        hsid,
185        secret_keys,
186        (),
187    )?
188    .connect(data)
189    .await
190}
191
192/// Common context for a single request to connect to a hidden service
193///
194/// This saves on passing this same set of (immutable) values (or subsets thereof)
195/// to each method in the principal functional code, everywhere.
196/// It also provides a convenient type to be `Self`.
197///
198/// Its lifetime is one request to make a new client circuit to a hidden service,
199/// including all the retries and timeouts.
200struct Context<'c, R: Runtime, M: MocksForConnect<R>> {
201    /// Runtime
202    runtime: &'c R,
203    /// Circpool
204    circpool: &'c M::HsCircPool,
205    /// Netdir
206    //
207    // TODO holding onto the netdir for the duration of our attempts is not ideal
208    // but doing better is fairly complicated.  See discussions here:
209    //   https://gitlab.torproject.org/tpo/core/arti/-/merge_requests/1228#note_2910545
210    //   https://gitlab.torproject.org/tpo/core/arti/-/issues/884
211    netdir: Arc<NetDir>,
212    /// Configuration
213    config: Arc<Config>,
214    /// Secret keys to use
215    secret_keys: HsClientSecretKeys,
216    /// HS ID
217    hsid: DispRedacted<HsId>,
218    /// Blinded HS ID
219    hs_blind_id: HsBlindId,
220    /// The subcredential to use during this time period
221    subcredential: Subcredential,
222    /// Mock data
223    mocks: M,
224}
225
226/// Details of an established rendezvous point
227///
228/// Intermediate value for progress during a connection attempt.
229struct Rendezvous<'r, R: Runtime, M: MocksForConnect<R>> {
230    /// RPT as a `Relay`
231    rend_relay: Relay<'r>,
232    /// Rendezvous circuit
233    rend_tunnel: DataTunnel!(R, M),
234    /// Rendezvous cookie
235    rend_cookie: RendCookie,
236
237    /// Receiver that will give us the RENDEZVOUS2 message.
238    ///
239    /// The sending ended is owned by the handler
240    /// which receives control messages on the rendezvous circuit,
241    /// and which was installed when we sent `ESTABLISH_RENDEZVOUS`.
242    ///
243    /// (`RENDEZVOUS2` is the message containing the onion service's side of the handshake.)
244    rend2_rx: proto_oneshot::Receiver<Rendezvous2>,
245
246    /// Dummy, to placate compiler
247    ///
248    /// Covariant without dropck or interfering with Send/Sync will do fine.
249    marker: PhantomData<fn() -> (R, M)>,
250}
251
252/// Random value used as part of IPT selection
253type IptSortRand = u32;
254
255/// Details of an apparently-useable introduction point
256///
257/// Intermediate value for progress during a connection attempt.
258struct UsableIntroPt<'i> {
259    /// Index in HS descriptor
260    intro_index: IntroPtIndex,
261    /// IPT descriptor
262    intro_desc: &'i IntroPointDesc,
263    /// IPT `CircTarget`
264    intro_target: OwnedCircTarget,
265    /// Random value used as part of IPT selection
266    sort_rand: IptSortRand,
267}
268
269/// Lookup key for looking up and recording information about a relay
270///
271/// Used to identify a relay when looking to see what happened last time we used it,
272/// and storing that information after we tried it.
273///
274/// We store the experience information under an arbitrary one of the relay's identities,
275/// as returned by the `HasRelayIds::identities().next()`.
276/// When we do lookups, we check all the relay's identities to see if we find
277/// anything relevant.
278/// If relay identities permute in strange ways, whether we find our previous
279/// knowledge about them is not particularly well defined, but that's fine.
280///
281/// While this is, structurally, a relay identity, it is not suitable for other purposes.
282#[derive(Hash, Eq, PartialEq, Ord, PartialOrd, Debug)]
283struct RelayIdFor<K> {
284    /// The relay id
285    inner: RelayId,
286
287    /// Phantom data to allow parameterizing over `K`
288    ///
289    /// `K` is a marker type that represents the kind of map
290    /// this key will be used in.
291    marker: PhantomData<K>,
292}
293
294/// Marker type, to make typed Ipt exprience [`RelayIdFor`] keys
295#[derive(Hash, Eq, PartialEq, Ord, PartialOrd, Copy, Clone, Debug)]
296struct IptExperienceMap;
297
298/// Lookup key for looking up and recording our IPT use experiences
299type RelayIdForExperience = RelayIdFor<IptExperienceMap>;
300
301/// Details of an apparently-successful INTRODUCE exchange
302///
303/// Intermediate value for progress during a connection attempt.
304struct Introduced<R: Runtime, M: MocksForConnect<R>> {
305    /// End-to-end crypto NTORv3 handshake with the service
306    ///
307    /// Created as part of generating our `INTRODUCE1`,
308    /// and then used when processing `RENDEZVOUS2`.
309    handshake_state: hs_ntor::HsNtorClientState,
310
311    /// A set of peer extensions that we decided to negotiate and use.
312    peer_caps: caps::PeerCaps,
313
314    /// Dummy, to placate compiler
315    ///
316    /// `R` and `M` only used for getting to mocks.
317    /// Covariant without dropck or interfering with Send/Sync will do fine.
318    marker: PhantomData<fn() -> (R, M)>,
319}
320
321impl<K> RelayIdFor<K> {
322    /// Create a new key for use with `T`
323    fn new(inner: RelayId) -> Self {
324        Self {
325            inner,
326            marker: Default::default(),
327        }
328    }
329
330    /// Identities to use to try to find previous experience information about this IPT
331    fn for_lookup<T: HasRelayIds>(ids: &T) -> impl Iterator<Item = Self> + '_ {
332        ids.identities().map(|id| RelayIdFor::new(id.to_owned()))
333    }
334
335    /// Identity to use to store previous experience information about this IPT
336    fn for_store<T: HasRelayIds>(ids: &T) -> Result<Self, Bug> {
337        let id = ids
338            .identities()
339            .next()
340            .ok_or_else(|| internal!("introduction point relay with no identities"))?
341            .to_owned();
342        Ok(RelayIdFor::new(id))
343    }
344}
345
346/// Sort key for an introduction point, for selecting the best IPTs to try first
347///
348/// Ordering is most preferable first.
349///
350/// We use this to sort our `UsableIpt`s using `.sort_by_key`.
351/// (This implementation approach ensures that we obey all the usual ordering invariants.)
352#[derive(Ord, PartialOrd, Eq, PartialEq, Debug)]
353struct IptSortKey {
354    /// Sort by how preferable the experience was
355    outcome: IptSortKeyOutcome,
356    /// Failing that, choose randomly
357    sort_rand: IptSortRand,
358}
359
360/// Component of the [`IptSortKey`] representing outcome of our last attempt, if any
361///
362/// This is the main thing we use to decide which IPTs to try first.
363/// It is calculated for each IPT
364/// (via `.sort_by_key`, so repeatedly - it should therefore be cheap to make.)
365///
366/// Ordering is most preferable first.
367#[derive(Ord, PartialOrd, Eq, PartialEq, Debug)]
368enum IptSortKeyOutcome {
369    /// Prefer successes
370    Success {
371        /// Prefer quick ones
372        duration: Duration,
373    },
374    /// Failing that, try one we don't know to have failed
375    Untried,
376    /// Failing that, it'll have to be ones that didn't work last time
377    Failed {
378        /// Prefer failures with an earlier retry time
379        retry_time: tor_error::LooseCmpRetryTime,
380        /// Failing that, prefer quick failures (rather than slow ones eg timeouts)
381        duration: Duration,
382    },
383}
384
385impl From<Option<&IptExperience>> for IptSortKeyOutcome {
386    fn from(experience: Option<&IptExperience>) -> IptSortKeyOutcome {
387        use IptSortKeyOutcome as O;
388        match experience {
389            None => O::Untried,
390            Some(IptExperience { duration, outcome }) => match outcome {
391                Ok(()) => O::Success {
392                    duration: *duration,
393                },
394                Err(retry_time) => O::Failed {
395                    retry_time: (*retry_time).into(),
396                    duration: *duration,
397                },
398            },
399        }
400    }
401}
402
403/// Token indicating that a descriptor fetch is wanted
404#[derive(Clone, Copy, Eq, PartialEq, Debug)]
405struct RefetchDescriptor;
406
407impl<'c, R: Runtime, M: MocksForConnect<R>> Context<'c, R, M> {
408    /// Make a new `Context` from the input data
409    fn new(
410        runtime: &'c R,
411        circpool: &'c M::HsCircPool,
412        netdir: Arc<NetDir>,
413        config: Arc<Config>,
414        hsid: HsId,
415        secret_keys: HsClientSecretKeys,
416        mocks: M,
417    ) -> Result<Self, ConnError> {
418        let time_period = netdir.hs_time_period();
419        let (hs_blind_id_key, subcredential) = HsIdKey::try_from(hsid)
420            .map_err(|_| CE::InvalidHsId)?
421            .compute_blinded_key(time_period)
422            .map_err(
423                // TODO HS what on earth do these errors mean, in practical terms ?
424                // In particular, we'll want to convert them to a ConnError variant,
425                // but what ErrorKind should they have ?
426                into_internal!("key blinding error, don't know how to handle"),
427            )?;
428        let hs_blind_id = hs_blind_id_key.id();
429
430        Ok(Context {
431            netdir,
432            config,
433            hsid: DispRedacted(hsid),
434            hs_blind_id,
435            subcredential,
436            circpool,
437            runtime,
438            secret_keys,
439            mocks,
440        })
441    }
442
443    /// Actually make a HS connection, updating our recorded state as necessary
444    ///
445    /// Called by the `connect` function in this module.
446    ///
447    /// This function handles all necessary retrying of fallible operations,
448    /// (and, therefore, must also limit the total work done for a particular call).
449    #[instrument(level = "trace", skip_all)]
450    async fn connect(&self, data: &mut Data) -> Result<DataTunnel!(R, M), ConnError> {
451        // This function must do the following, retrying as appropriate.
452        //  - Look up the onion descriptor in the state.
453        //  - Download the onion descriptor if one isn't there.
454        //  - In parallel:
455        //    - Pick a rendezvous point from the netdirprovider and launch a
456        //      rendezvous circuit to it. Then send ESTABLISH_INTRO.
457        //    - Pick a number of introduction points (1 or more) and try to
458        //      launch circuits to them.
459        //  - On a circuit to an introduction point, send an INTRODUCE1 cell.
460        //  - Wait for a RENDEZVOUS2 cell on the rendezvous circuit
461        //  - Add a virtual hop to the rendezvous circuit.
462        //  - Return the rendezvous circuit.
463
464        let mocks = self.mocks.clone();
465
466        let desc = self
467            .descriptor_ensure(&mut data.desc, &mut data.hsdirs, None)
468            .await?;
469
470        mocks.test_got_desc(desc);
471
472        let tunnel = match self.intro_rend_connect(desc, &mut data.ipts).await {
473            Ok(tunnel) => tunnel,
474            Err(e) => {
475                let is_intro_nack = |e| {
476                    if let FAE::IntroductionFailed { status, .. } = e {
477                        status == IntroduceAckStatus::NOT_RECOGNIZED
478                    } else {
479                        false
480                    }
481                };
482
483                let retry = if let CE::Failed(ref errors) = e {
484                    // If any of the errors are an INTRODUCE_NACK,
485                    // then it's worth retrying one more time
486                    // with a fresh descriptor.
487                    errors
488                        .clone()
489                        .into_iter()
490                        .any(is_intro_nack)
491                        .then_some(RefetchDescriptor)
492                } else {
493                    None
494                };
495
496                if let Some(RefetchDescriptor) = retry {
497                    debug!(
498                        "Introduction to {} NACKed, refetching descriptor and retrying",
499                        &self.hsid,
500                    );
501                    // Refetch the descriptor and try one more time
502                    let desc = self
503                        .descriptor_ensure(&mut data.desc, &mut data.hsdirs, retry)
504                        .await?;
505                    mocks.test_got_desc(desc);
506                    self.intro_rend_connect(desc, &mut data.ipts).await?
507                } else {
508                    return Err(e);
509                }
510            }
511        };
512
513        mocks.test_got_tunnel(&tunnel);
514
515        Ok(tunnel)
516    }
517
518    /// Ensure that `Data.desc` contains the HS descriptor
519    ///
520    /// If we have a previously-downloaded descriptor, which is still valid,
521    /// just returns a reference to it.
522    ///
523    /// Otherwise, tries to obtain the descriptor by downloading it from hsdir(s).
524    ///
525    /// If `refetch` is `true`, a new descriptor will be refetched
526    /// from the hsdir(s) unconditionally.
527    ///
528    /// Does all necessary retries and timeouts.
529    /// Returns an error if no valid descriptor could be found.
530    #[instrument(level = "trace", skip_all)]
531    async fn descriptor_ensure<'d>(
532        &self,
533        data: &'d mut DataHsDesc,
534        recent_hsdirs: &'d mut DataHsDirs,
535        refetch: Option<RefetchDescriptor>,
536    ) -> Result<&'d HsDesc, CE> {
537        // Maximum number of hsdir connection and retrieval attempts we'll make
538        let max_total_attempts = self
539            .config
540            .retry
541            .hs_desc_fetch_attempts()
542            .try_into()
543            // User specified a very large u32.  We must be downcasting it to 16bit!
544            // let's give them as many retries as we can manage.
545            .unwrap_or(usize::MAX);
546
547        let now = self.runtime.wallclock();
548        let unwrap_valid_desc = |data: &'d mut DataHsDesc| -> &'d HsDesc {
549            data.as_ref()
550                .expect("Some but now None")
551                .desc
552                .as_ref()
553                .if_valid_at(&now)
554                .expect("Ok but now Err")
555        };
556
557        // We retain a previously obtained descriptor precisely until its lifetime expires,
558        // or until we refetch a more recent one
559        // as a result of an `intro_rend_connect()` failure caused by introduce NACK.
560        //
561        // When it expires, we discard it completely and try to obtain a new one.
562        //
563        // We only replace our cached descriptor if the new one has a higher revision counter.
564        //
565        // TODO SPEC: Discuss HS descriptor lifetime and expiry client behaviour
566        let now = self.runtime.wallclock();
567
568        let stored_revision = data.as_ref().and_then(|previously| {
569            if let Ok(desc) = previously.desc.as_ref().if_valid_at(&now) {
570                // Ideally we would just return desc but that confuses borrowck,
571                // so we have to use unwrap_valid_desc() each time
572                // we need the known-to-be-Some descriptor instead.
573                //
574                // https://github.com/rust-lang/rust/issues/51545
575                Some((desc.revision(), previously.time_period))
576            } else {
577                // Seems to be not valid now.  Try to fetch a fresh one.
578                None
579            }
580        });
581
582        match (stored_revision, refetch) {
583            (Some(_), None) => {
584                // Our cached descriptor is still timely,
585                // and we don't need to fetch a new one.
586                return Ok(unwrap_valid_desc(data));
587            }
588            (None, _) => {
589                // We don't have a timely descriptor,
590                // so ignore the requery_interval,
591                // and reach out to all HsDirs
592                recent_hsdirs.clear();
593            }
594            (_, Some(RefetchDescriptor)) => {
595                // We have been asked to try to fetch a new descriptor.
596                // We will only reach out to the HsDirs that are
597                // not within the `hs_dir_requery_interval`
598            }
599        }
600
601        // First, filter out any HsDirs that we *can* requery
602        recent_hsdirs.retain(|_hsdir, requery| *requery > now);
603
604        let working_tp = self.netdir.hs_time_period();
605        let hs_dirs = self.netdir.hs_dirs_download(
606            self.hs_blind_id,
607            working_tp,
608            &mut self.mocks.thread_rng(),
609        )?;
610
611        trace!(
612            "HS desc fetch for {}, for period {}, using {} hsdirs",
613            &self.hsid,
614            working_tp,
615            hs_dirs.len()
616        );
617
618        let hs_dirs = hs_dirs
619            .into_iter()
620            .filter(|hsdir| {
621                // Skip over any HsDirs that we are not allowed to requery right now
622                let should_skip = recent_hsdirs.keys().any(|recent| {
623                    RelayIdForRequeryPeriod::for_lookup(hsdir).any(|id| id == *recent)
624                });
625
626                !should_skip
627            })
628            .collect::<Vec<_>>();
629
630        if hs_dirs.is_empty() {
631            warn!(
632                "Tried to fetch HS desc for {}, for period {}, but all hsdirs are rate-limited",
633                &self.hsid, working_tp,
634            );
635
636            if stored_revision.is_none() {
637                // We can't fetch a new descriptor, and we don't have a cached one.
638                return Err(CE::NoUsableHsDirs);
639            } else {
640                // Return our cached descriptor
641                return Ok(unwrap_valid_desc(data));
642            }
643        }
644
645        let params = self.netdir.params();
646
647        // We might consider launching requests to multiple HsDirs in parallel.
648        //   https://gitlab.torproject.org/tpo/core/arti/-/merge_requests/1118#note_2894463
649        // But C Tor doesn't and our HS experts don't consider that important:
650        //   https://gitlab.torproject.org/tpo/core/arti/-/issues/913#note_2914436
651        // (Additionally, making multiple HSDir requests at once may make us
652        // more vulnerable to traffic analysis.)
653        let mut attempts = hs_dirs.iter().cycle().take(max_total_attempts);
654        let mut errors = RetryError::in_attempt_to("retrieve hidden service descriptor");
655        let desc = loop {
656            let relay = match attempts.next() {
657                Some(relay) => relay,
658                None => {
659                    return Err(if errors.is_empty() {
660                        CE::NoHsDirs
661                    } else {
662                        CE::DescriptorDownload(errors)
663                    });
664                }
665            };
666            let hsdir_for_error: Sensitive<Ed25519Identity> = (*relay.id()).into();
667
668            let hsdir = RelayIdForRequeryPeriod::for_store(relay)?;
669            // Ensure we wait at least hs_dir_requery_interval() until we try to
670            // fecth from this HsDir again
671            recent_hsdirs.insert(hsdir, now + self.config.retry.hs_dir_requery_interval());
672
673            match self.descriptor_fetch_attempt(relay, params).await {
674                Ok(desc) => break desc,
675                Err(error) => {
676                    if error.should_report_as_suspicious() {
677                        // Note that not every protocol violation is suspicious:
678                        // we only warn on the protocol violations that look like attempts
679                        // to do a traffic tagging attack via hsdir inflation.
680                        // (See proposal 360.)
681                        warn_report!(
682                            &error,
683                            "Suspicious failure while downloading hsdesc for {} from relay {}",
684                            &self.hsid,
685                            relay.display_relay_ids(),
686                        );
687                    } else {
688                        debug_report!(
689                            &error,
690                            "failed hsdir desc fetch for {} from {}/{}",
691                            &self.hsid,
692                            &relay.id(),
693                            &relay.rsa_id()
694                        );
695                    }
696                    errors.push_timed(
697                        tor_error::Report(DescriptorError {
698                            hsdir: hsdir_for_error,
699                            error,
700                        }),
701                        self.runtime.now(),
702                        Some(self.runtime.wallclock()),
703                    );
704                }
705            }
706        };
707
708        // If our existing descriptor is newer than the one we have just fetched,
709        // we should retain it.
710        if let Some(stored_revision) = stored_revision {
711            // It is safe to dangerously_assume_timely,
712            // as descriptor_fetch_attempt has already checked the timeliness of the descriptor.
713            let new_desc = desc.as_ref().dangerously_assume_timely();
714
715            // Revision counters are monotonically increasing within a given time period.
716            // If our newly fetched descriptor has the same HsBlindId as our cached one,
717            // it means they are both used for the same time period,
718            // and so we should only update our cache if the new descriptor is more recent
719            // (i.e. it has a higher revision counter).
720            if stored_revision >= (new_desc.revision(), working_tp) {
721                // Our cached descriptor is still timely, and has a higher revision counter
722                // than the one we've just fetched, so we retain it.
723                return Ok(unwrap_valid_desc(data));
724            }
725        }
726
727        // Store the bounded value in the cache for reuse,
728        // but return a reference to the unwrapped `HsDesc`.
729        //
730        // The `HsDesc` must be owned by `data.desc`,
731        // so first add it to `data.desc`,
732        // and then dangerously_assume_timely to get a reference out again.
733        //
734        // It is safe to dangerously_assume_timely,
735        // as descriptor_fetch_attempt has already checked the timeliness of the descriptor.
736        let desc = HsDescForTp {
737            time_period: working_tp,
738            desc,
739        };
740        let ret = data.insert(desc);
741        Ok(ret.desc.as_ref().dangerously_assume_timely())
742    }
743
744    /// Make one attempt to fetch the descriptor from a specific hsdir
745    ///
746    /// No timeout
747    ///
748    /// On success, returns the descriptor.
749    ///
750    /// While the returned descriptor is `TimeRangeBound`, its validity at the current time *has*
751    /// been checked.
752    #[instrument(level = "trace", skip_all)]
753    async fn descriptor_fetch_attempt(
754        &self,
755        hsdir: &Relay<'_>,
756        params: &NetParameters,
757    ) -> Result<TimeRangeBound<HsDesc>, DescriptorErrorDetail> {
758        let max_len: usize = self
759            .netdir
760            .params()
761            .hsdir_max_desc_size
762            .get()
763            .try_into()
764            .map_err(into_internal!("BoundedInt was not truly bounded!"))?;
765        let request = {
766            let mut r = tor_dirclient::request::HsDescDownloadRequest::new(self.hs_blind_id);
767            r.set_max_len(max_len);
768            r
769        };
770        trace!(
771            "hsdir for {}, trying {}/{}, request {:?} (http request {:?})",
772            &self.hsid,
773            &hsdir.id(),
774            &hsdir.rsa_id(),
775            &request,
776            request.debug_request()
777        );
778
779        let circuit = self
780            .circpool
781            .m_get_or_launch_dir(&self.netdir, OwnedCircTarget::from_circ_target(hsdir))
782            .await?;
783        let n_hops = circuit.m_num_hops()?;
784        let timeout_roundtrip =
785            self.estimate_timeout(&[(1, TimeoutsAction::RoundTrip { length: n_hops })]);
786
787        let source: Option<SourceInfo> = circuit
788            .m_source_info()
789            .map_err(into_internal!("Couldn't get SourceInfo for circuit"))?;
790
791        let mut stream = self
792            .runtime
793            // NOTE: In fact this timeout is overkill: this operation should succeed immediately,
794            // since we always send BEGINDIR messages optimistically (without waiting for a reply).
795            // But since our code is complex, and since it could become possible for this to block
796            // if the circuit is saturated or we implement proposal 367 or something,
797            // we may as well have _some_ timeout here.
798            .timeout(timeout_roundtrip, circuit.m_begin_dir_stream())
799            .await?
800            .map_err(DescriptorErrorDetail::Circuit)?;
801
802        let request_future =
803            tor_dirclient::send_request(self.runtime, &request, &mut stream, source);
804        let response = self
805            .runtime
806            .timeout(timeout_roundtrip, request_future)
807            .await?
808            .map_err(|dir_error| match dir_error {
809                tor_dirclient::Error::RequestFailed(rfe) => DescriptorErrorDetail::from(rfe.error),
810                tor_dirclient::Error::CircMgr(ce) => into_internal!(
811                    "tor-dirclient complains about circmgr going wrong but we gave it a stream"
812                )(ce)
813                .into(),
814                other => into_internal!(
815                    "tor-dirclient gave unexpected error, tor-hsclient code needs updating"
816                )(other)
817                .into(),
818            })?;
819
820        let desc_text = response.into_output_string().map_err(|rfe| rfe.error)?;
821        let hsc_desc_enc = self.secret_keys.keys.ks_hsc_desc_enc.as_ref();
822
823        let now = self.runtime.wallclock();
824
825        let desc = HsDesc::parse_decrypt_validate(
826            &desc_text,
827            &self.hs_blind_id,
828            &self.subcredential,
829            hsc_desc_enc,
830        )?;
831
832        // Check the validity time (relied on by descriptor_ensure)
833        desc.check_valid_at(&now)?;
834
835        // Validate cc_sendme_inc, if present, is within ±1 of params.cc_sendme_inc.
836        ensure_descriptor_compatible_with_params(desc.dangerously_peek(), params)?;
837
838        Ok(desc)
839    }
840
841    /// Given the descriptor, try to connect to service
842    ///
843    /// Does all necessary retries, timeouts, etc.
844    async fn intro_rend_connect(
845        &self,
846        desc: &HsDesc,
847        data: &mut DataIpts,
848    ) -> Result<DataTunnel!(R, M), CE> {
849        // Maximum number of rendezvous/introduction attempts we'll make
850        let max_total_attempts = self
851            .config
852            .retry
853            .hs_intro_rend_attempts()
854            .try_into()
855            // User specified a very large u32.  We must be downcasting it to 16bit!
856            // let's give them as many retries as we can manage.
857            .unwrap_or(usize::MAX);
858
859        // We can't reliably distinguish IPT failure from RPT failure, so we iterate over IPTs
860        // (best first) and each time use a random RPT.
861
862        // We limit the number of rendezvous establishment attempts, separately, since we don't
863        // try to talk to the intro pt until we've established the rendezvous circuit.
864        let mut rend_attempts = 0..max_total_attempts;
865
866        // But, we put all the errors into the same bucket, since we might have a mixture.
867        let mut errors = RetryError::in_attempt_to("make circuit to hidden service");
868
869        // Note that IntroPtIndex is *not* the index into this Vec.
870        // It is the index into the original list of introduction points in the descriptor.
871        let mut usable_intros: Vec<UsableIntroPt> = desc
872            .intro_points()
873            .iter()
874            .enumerate()
875            .map(|(intro_index, intro_desc)| {
876                let intro_index = intro_index.into();
877                let intro_target = ipt_to_circtarget(intro_desc, &self.netdir)
878                    .map_err(|error| FAE::UnusableIntro { error, intro_index })?;
879                // Lack of TAIT means this clone
880                let intro_target = OwnedCircTarget::from_circ_target(&intro_target);
881                Ok::<_, FailedAttemptError>(UsableIntroPt {
882                    intro_index,
883                    intro_desc,
884                    intro_target,
885                    sort_rand: self.mocks.thread_rng().random(),
886                })
887            })
888            .filter_map(|entry| match entry {
889                Ok(y) => Some(y),
890                Err(e) => {
891                    errors.push_timed(e, self.runtime.now(), Some(self.runtime.wallclock()));
892                    None
893                }
894            })
895            .collect_vec();
896
897        // Delete experience information for now-unlisted intro points
898        // Otherwise, as the IPTs change `Data` might grow without bound,
899        // if we keep reconnecting to the same HS.
900        data.retain(|k, _v| {
901            usable_intros
902                .iter()
903                .any(|ipt| RelayIdForExperience::for_lookup(&ipt.intro_target).any(|id| &id == k))
904        });
905
906        // Join with existing state recording our experiences,
907        // sort by descending goodness, and then randomly
908        // (so clients without any experience don't all pile onto the same, first, IPT)
909        usable_intros.sort_by_key(|ipt: &UsableIntroPt| {
910            let experience =
911                RelayIdForExperience::for_lookup(&ipt.intro_target).find_map(|id| data.get(&id));
912            IptSortKey {
913                outcome: experience.into(),
914                sort_rand: ipt.sort_rand,
915            }
916        });
917        self.mocks.test_got_ipts(&usable_intros);
918
919        let mut intro_attempts = usable_intros.iter().cycle().take(max_total_attempts);
920
921        // We retain a rendezvous we managed to set up in here.  That way if we created it, and
922        // then failed before we actually needed it, we can reuse it.
923        // If we exit with an error, we will waste it - but because we isolate things we do
924        // for different services, it wouldn't be reusable anyway.
925        let mut saved_rendezvous = None;
926
927        // If we are using proof-of-work DoS mitigation, this chooses an
928        // algorithm and initial effort, and adjusts that effort when we retry.
929        let mut pow_client = HsPowClient::new(&self.hs_blind_id, desc);
930
931        // We might consider making multiple INTRODUCE attempts to different
932        // IPTs in parallel, and somehow aggregating the errors and
933        // experiences.
934        // However our HS experts don't consider that important:
935        //   https://gitlab.torproject.org/tpo/core/arti/-/issues/913#note_2914438
936        // Parallelizing our HsCircPool circuit building would likely have
937        // greater impact. (See #1149.)
938        loop {
939            // When did we start doing things that depended on the IPT?
940            //
941            // Used for recording our experience with the selected IPT
942            let mut ipt_use_started = None::<Instant>;
943
944            // Error handling inner async block (analogous to an IEFE):
945            //  * Ok(Some()) means this attempt succeeded
946            //  * Ok(None) means all attempts exhausted
947            //  * Err(error) means this attempt failed
948            //
949            let outcome = async {
950                // We establish a rendezvous point first.  Although it appears from reading
951                // this code that this means we serialise establishment of the rendezvous and
952                // introduction circuits, this isn't actually the case.  The circmgr maintains
953                // a pool of circuits.  What actually happens in the "standing start" case is
954                // that we obtain a circuit for rendezvous from the circmgr's pool, expecting
955                // one to be available immediately; the circmgr will then start to build a new
956                // one to replenish its pool, and that happens in parallel with the work we do
957                // here - but in arrears.  If the circmgr pool is empty, then we must wait.
958                //
959                // Perhaps this should be parallelised here.  But that's really what the pool
960                // is for, since we expect building the rendezvous circuit and building the
961                // introduction circuit to take about the same length of time.
962                //
963                // We *do* serialise the ESTABLISH_RENDEZVOUS exchange, with the
964                // building of the introduction circuit.  That could be improved, at the cost
965                // of some additional complexity here.
966                //
967                // Our HS experts don't consider it important to increase the parallelism:
968                //   https://gitlab.torproject.org/tpo/core/arti/-/issues/913#note_2914444
969                //   https://gitlab.torproject.org/tpo/core/arti/-/issues/913#note_2914445
970                if saved_rendezvous.is_none() {
971                    debug!("hs conn to {}: setting up rendezvous point", &self.hsid);
972                    // Establish a rendezvous circuit.
973                    let Some(_): Option<usize> = rend_attempts.next() else {
974                        return Ok(None);
975                    };
976
977                    saved_rendezvous = Some(self.establish_rendezvous().await?);
978                }
979
980                let Some(ipt) = intro_attempts.next() else {
981                    return Ok(None);
982                };
983                let intro_index = ipt.intro_index;
984                let is_single_onion_service = desc.is_single_onion_service();
985
986                let proof_of_work = match pow_client.solve().await {
987                    Ok(solution) => solution,
988                    Err(e) => {
989                        debug!(
990                            "failing to compute proof-of-work, trying without. ({:?})",
991                            e
992                        );
993                        None
994                    }
995                };
996
997                // We record how long things take, starting from here, as
998                // as a statistic we'll use for the IPT in future.
999                // This is stored in a variable outside this async block,
1000                // so that the outcome handling can use it.
1001                ipt_use_started = Some(self.runtime.now());
1002
1003                // No `Option::get_or_try_insert_with`, or we'd avoid this expect()
1004                let rend_pt_for_error = rend_pt_identity_for_error(
1005                    &saved_rendezvous
1006                        .as_ref()
1007                        .expect("just made Some")
1008                        .rend_relay,
1009                );
1010                debug!(
1011                    "hs conn to {}: RPT {}",
1012                    &self.hsid,
1013                    rend_pt_for_error.as_inner()
1014                );
1015
1016                let (rendezvous, introduced) =
1017                    self.exchange_introduce(desc, ipt, &mut saved_rendezvous, proof_of_work)
1018                    .await
1019                    // TODO: Maybe try, once, to extend-and-reuse the intro circuit.
1020                    //
1021                    // If the introduction fails, the introduction circuit is in principle
1022                    // still usable.  We believe that in this case, C Tor extends the intro
1023                    // circuit by one hop to the next IPT to try.  That saves on building a
1024                    // whole new 3-hop intro circuit.  However, our HS experts tell us that
1025                    // if introduction fails at one IPT it is likely to fail at the others too,
1026                    // so that optimisation might reduce our network impact and time to failure,
1027                    // but isn't likely to improve our chances of success.
1028                    //
1029                    // However, it's not clear whether this approach risks contaminating
1030                    // the 2nd attempt with some fault relating to the introduction point.
1031                    // The 1st ipt might also gain more knowledge about which HS we're talking to.
1032                    //
1033                    // TODO SPEC: Discuss extend-and-reuse HS intro circuit after nack
1034                    ?;
1035                #[allow(unused_variables)] // it's *supposed* to be unused
1036                let saved_rendezvous = (); // don't use `saved_rendezvous` any more, use rendezvous
1037
1038                let rend_pt = rend_pt_identity_for_error(&rendezvous.rend_relay);
1039                let circ = self.complete_rendezvous(ipt, rendezvous, introduced, is_single_onion_service)
1040                    .await?;
1041
1042                debug!(
1043                    "hs conn to {}: RPT {} IPT {}: success",
1044                    &self.hsid,
1045                    rend_pt.as_inner(),
1046                    intro_index,
1047                );
1048                Ok::<_, FAE>(Some((intro_index, circ)))
1049            }
1050            .await;
1051
1052            // Store the experience `outcome` we had with IPT `intro_index`, in `data`
1053            #[allow(clippy::unused_unit)] // -> () is here for error handling clarity
1054            let mut store_experience = |intro_index, outcome| -> () {
1055                (|| {
1056                    let ipt = usable_intros
1057                        .iter()
1058                        .find(|ipt| ipt.intro_index == intro_index)
1059                        .ok_or_else(|| internal!("IPT not found by index"))?;
1060                    let id = RelayIdForExperience::for_store(&ipt.intro_target)?;
1061                    let started = ipt_use_started.ok_or_else(|| {
1062                        internal!("trying to record IPT use but no IPT start time noted")
1063                    })?;
1064                    let duration = self
1065                        .runtime
1066                        .now()
1067                        .checked_duration_since(started)
1068                        .ok_or_else(|| internal!("clock overflow calculating IPT use duration"))?;
1069                    data.insert(id, IptExperience { duration, outcome });
1070                    Ok::<_, Bug>(())
1071                })()
1072                .unwrap_or_else(|e| warn_report!(e, "error recording HS IPT use experience"));
1073            };
1074
1075            match outcome {
1076                Ok(Some((intro_index, y))) => {
1077                    // Record successful outcome in Data
1078                    store_experience(intro_index, Ok(()));
1079                    return Ok(y);
1080                }
1081                Ok(None) => return Err(CE::Failed(errors)),
1082                Err(error) => {
1083                    debug_report!(&error, "hs conn to {}: attempt failed", &self.hsid);
1084                    // Record error outcome in Data, if in fact we involved the IPT
1085                    // at all.  The IPT information is be retrieved from `error`,
1086                    // since only some of the errors implicate the introduction point.
1087                    if let Some(intro_index) = error.intro_index() {
1088                        store_experience(intro_index, Err(error.retry_time()));
1089                    }
1090                    errors.push_timed(error, self.runtime.now(), Some(self.runtime.wallclock()));
1091
1092                    // If we are using proof-of-work DoS mitigation, try harder next time
1093                    pow_client.increase_effort();
1094                }
1095            }
1096        }
1097    }
1098
1099    /// Make one attempt to establish a rendezvous circuit
1100    ///
1101    /// This doesn't really depend on anything,
1102    /// other than (obviously) the isolation implied by our circuit pool.
1103    /// In particular it doesn't depend on the introduction point.
1104    ///
1105    /// Applies timeouts as appropriate.
1106    #[instrument(level = "trace", skip_all)]
1107    async fn establish_rendezvous(&'c self) -> Result<Rendezvous<'c, R, M>, FAE> {
1108        let (rend_tunnel, rend_relay) = self
1109            .circpool
1110            .m_get_or_launch_client_rend(&self.netdir)
1111            .await
1112            .map_err(|error| FAE::RendezvousCircuitObtain { error })?;
1113
1114        let rend_pt = rend_pt_identity_for_error(&rend_relay);
1115
1116        let rend_cookie: RendCookie = self.mocks.thread_rng().random();
1117        let message = EstablishRendezvous::new(rend_cookie);
1118
1119        let (rend_established_tx, rend_established_rx) = proto_oneshot::channel();
1120        let (rend2_tx, rend2_rx) = proto_oneshot::channel();
1121
1122        /// Handler which expects `RENDEZVOUS_ESTABLISHED` and then
1123        /// `RENDEZVOUS2`.   Returns each message via the corresponding `oneshot`.
1124        struct Handler {
1125            /// Sender for a RENDEZVOUS_ESTABLISHED message.
1126            rend_established_tx: proto_oneshot::Sender<RendezvousEstablished>,
1127            /// Sender for a RENDEZVOUS2 message.
1128            rend2_tx: proto_oneshot::Sender<Rendezvous2>,
1129        }
1130        impl MsgHandler for Handler {
1131            fn handle_msg(
1132                &mut self,
1133                msg: AnyRelayMsg,
1134            ) -> Result<MetaCellDisposition, tor_proto::Error> {
1135                // The first message we expect is a RENDEZVOUS_ESTABALISHED.
1136                if self.rend_established_tx.still_expected() {
1137                    self.rend_established_tx
1138                        .deliver_expected_message(msg, MetaCellDisposition::Consumed)
1139                } else {
1140                    self.rend2_tx
1141                        .deliver_expected_message(msg, MetaCellDisposition::ConversationFinished)
1142                }
1143            }
1144        }
1145
1146        debug!(
1147            "hs conn to {}: RPT {}: sending ESTABLISH_RENDEZVOUS",
1148            &self.hsid,
1149            rend_pt.as_inner(),
1150        );
1151
1152        let failed_map_err = |error| FAE::RendezvousEstablish {
1153            error,
1154            rend_pt: rend_pt.clone(),
1155        };
1156        let handler = Handler {
1157            rend_established_tx,
1158            rend2_tx,
1159        };
1160
1161        let num_hops = rend_tunnel
1162            .m_num_own_hops()
1163            .map_err(|error| FAE::RendezvousCircuitObtain { error })?;
1164
1165        let timeout_roundtrip =
1166            self.estimate_timeout(&[(1, TimeoutsAction::RoundTrip { length: num_hops })]);
1167
1168        // TODO(conflux) This error handling is horrible. Problem is that this Mock system requires
1169        // to send back a tor_circmgr::Error while our reply handler requires a tor_proto::Error.
1170        // And unifying both is hard here considering it needs to be converted to yet another Error
1171        // type "FAE" so we have to do these hoops and jumps.
1172        rend_tunnel
1173            .m_start_conversation_last_hop(Some(message.into()), handler)
1174            .await
1175            .map_err(|e| {
1176                let proto_error = match e {
1177                    tor_circmgr::Error::Protocol { error, .. } => error,
1178                    _ => tor_proto::Error::CircuitClosed,
1179                };
1180                FAE::RendezvousEstablish {
1181                    error: proto_error,
1182                    rend_pt: rend_pt.clone(),
1183                }
1184            })?;
1185
1186        // `start_conversation` returns as soon as the control message has been sent.
1187        // We need to obtain the RENDEZVOUS_ESTABLISHED message, which is "returned" via the oneshot.
1188        let _: RendezvousEstablished = self
1189            .runtime
1190            .timeout(timeout_roundtrip, rend_established_rx.recv(failed_map_err))
1191            .await
1192            .map_err(
1193                |_timeout: tor_rtcompat::TimeoutError| FAE::RendezvousEstablishTimeout {
1194                    rend_pt: rend_pt.clone(),
1195                },
1196            )??;
1197
1198        debug!(
1199            "hs conn to {}: RPT {}: got RENDEZVOUS_ESTABLISHED",
1200            &self.hsid,
1201            rend_pt.as_inner(),
1202        );
1203
1204        Ok(Rendezvous {
1205            rend_tunnel,
1206            rend_cookie,
1207            rend_relay,
1208            rend2_rx,
1209            marker: PhantomData,
1210        })
1211    }
1212
1213    /// Attempt (once) to send an INTRODUCE1 and wait for the INTRODUCE_ACK
1214    ///
1215    /// `take`s the input `rendezvous` (but only takes it if it gets that far)
1216    /// and, if successful, returns it.
1217    /// (This arranges that the rendezvous is "used up" precisely if
1218    /// we sent its secret somewhere.)
1219    ///
1220    /// Although this function handles the `Rendezvous`,
1221    /// nothing in it actually involves the rendezvous point.
1222    /// So if there's a failure, it's purely to do with the introduction point.
1223    ///
1224    /// Applies timeouts as appropriate.
1225    #[allow(clippy::type_complexity)] // TODO: Refactor
1226    #[instrument(level = "trace", skip_all)]
1227    async fn exchange_introduce(
1228        &'c self,
1229        hsdesc: &HsDesc,
1230        ipt: &UsableIntroPt<'_>,
1231        rendezvous: &mut Option<Rendezvous<'c, R, M>>,
1232        proof_of_work: Option<ProofOfWork>,
1233    ) -> Result<(Rendezvous<'c, R, M>, Introduced<R, M>), FAE> {
1234        let intro_index = ipt.intro_index;
1235
1236        debug!(
1237            "hs conn to {}: IPT {}: obtaining intro circuit",
1238            &self.hsid, intro_index,
1239        );
1240
1241        let intro_circ = self
1242            .circpool
1243            .m_get_or_launch_intro(
1244                &self.netdir,
1245                ipt.intro_target.clone(), // &OwnedCircTarget isn't CircTarget apparently
1246            )
1247            .await
1248            .map_err(|error| FAE::IntroductionCircuitObtain { error, intro_index })?;
1249
1250        let rendezvous = rendezvous.take().ok_or_else(|| internal!("no rend"))?;
1251
1252        let rend_pt = rend_pt_identity_for_error(&rendezvous.rend_relay);
1253
1254        debug!(
1255            "hs conn to {}: RPT {} IPT {}: making introduction",
1256            &self.hsid,
1257            rend_pt.as_inner(),
1258            intro_index,
1259        );
1260
1261        // Now we construct an introduce1 message and perform the first part of the
1262        // rendezvous handshake.
1263        //
1264        // This process is tricky because the header of the INTRODUCE1 message
1265        // -- which depends on the IntroPt configuration -- is authenticated as
1266        // part of the HsDesc handshake.
1267
1268        // Construct the header, since we need it as input to our encryption.
1269        let intro_header = {
1270            let ipt_sid_key = ipt.intro_desc.ipt_sid_key();
1271            let intro1 = Introduce1::new(
1272                AuthKeyType::ED25519_SHA3_256,
1273                ipt_sid_key.as_bytes().to_vec(),
1274                vec![],
1275            );
1276            let mut header = vec![];
1277            intro1
1278                .encode_onto(&mut header)
1279                .map_err(into_internal!("couldn't encode intro1 header"))?;
1280            header
1281        };
1282
1283        let peer_caps = caps::PeerCaps::new(hsdesc);
1284
1285        // Construct the introduce payload, which tells the onion service how to find
1286        // our rendezvous point.  (We could do this earlier if we wanted.)
1287        let intro_payload = {
1288            let onion_key =
1289                intro_payload::OnionKey::NtorOnionKey(*rendezvous.rend_relay.ntor_onion_key());
1290            let linkspecs = rendezvous
1291                .rend_relay
1292                .linkspecs()
1293                .map_err(into_internal!("Couldn't encode link specifiers"))?;
1294            #[allow(unused_mut)] // TODO: Remove once negotiate-extensions is always-on.
1295            let mut payload = IntroduceHandshakePayload::new(
1296                rendezvous.rend_cookie,
1297                onion_key,
1298                linkspecs,
1299                proof_of_work,
1300            );
1301
1302            peer_caps.add_extensions(&mut payload);
1303
1304            let mut encoded = vec![];
1305            payload
1306                .write_onto(&mut encoded)
1307                .map_err(into_internal!("Couldn't encode introduce1 payload"))?;
1308            encoded
1309        };
1310
1311        // Perform the cryptographic handshake with the onion service.
1312        let service_info = hs_ntor::HsNtorServiceInfo::new(
1313            ipt.intro_desc.svc_ntor_key().clone(),
1314            ipt.intro_desc.ipt_sid_key().clone(),
1315            self.subcredential,
1316        );
1317        let handshake_state =
1318            hs_ntor::HsNtorClientState::new(&mut self.mocks.thread_rng(), service_info);
1319        let encrypted_body = handshake_state
1320            .client_send_intro(&intro_header, &intro_payload)
1321            .map_err(into_internal!("can't begin hs-ntor handshake"))?;
1322
1323        // Build our actual INTRODUCE1 message.
1324        let intro1_real = Introduce1::new(
1325            AuthKeyType::ED25519_SHA3_256,
1326            ipt.intro_desc.ipt_sid_key().as_bytes().to_vec(),
1327            encrypted_body,
1328        );
1329
1330        /// Handler which expects just `INTRODUCE_ACK`
1331        struct Handler {
1332            /// Sender for `INTRODUCE_ACK`
1333            intro_ack_tx: proto_oneshot::Sender<IntroduceAck>,
1334        }
1335        impl MsgHandler for Handler {
1336            fn handle_msg(
1337                &mut self,
1338                msg: AnyRelayMsg,
1339            ) -> Result<MetaCellDisposition, tor_proto::Error> {
1340                self.intro_ack_tx
1341                    .deliver_expected_message(msg, MetaCellDisposition::ConversationFinished)
1342            }
1343        }
1344        let failed_map_err = |error| FAE::IntroductionExchange { error, intro_index };
1345        let (intro_ack_tx, intro_ack_rx) = proto_oneshot::channel();
1346        let handler = Handler { intro_ack_tx };
1347
1348        let num_hops = intro_circ
1349            .m_num_hops()
1350            .map_err(|error| FAE::IntroductionCircuitObtain { error, intro_index })?;
1351        // NOTE: Should we allow this to be longer in case the introduction point is grievously
1352        // overloaded?
1353        let timeout_roundtrip =
1354            self.estimate_timeout(&[(1, TimeoutsAction::RoundTrip { length: num_hops })]);
1355
1356        debug!(
1357            "hs conn to {}: RPT {} IPT {}: making introduction - sending INTRODUCE1",
1358            &self.hsid,
1359            rend_pt.as_inner(),
1360            intro_index,
1361        );
1362
1363        // TODO(conflux) This error handling is horrible. Problem is that this Mock system requires
1364        // to send back a tor_circmgr::Error while our reply handler requires a tor_proto::Error.
1365        // And unifying both is hard here considering it needs to be converted to yet another Error
1366        // type "FAE" so we have to do these hoops and jumps.
1367        intro_circ
1368            .m_start_conversation_last_hop(Some(intro1_real.into()), handler)
1369            .await
1370            .map_err(|e| {
1371                let proto_error = match e {
1372                    tor_circmgr::Error::Protocol { error, .. } => error,
1373                    _ => tor_proto::Error::CircuitClosed,
1374                };
1375                FAE::IntroductionExchange {
1376                    error: proto_error,
1377                    intro_index,
1378                }
1379            })?;
1380
1381        // Status is checked by `.success()`, and we don't look at the extensions;
1382        // just discard the known-successful `IntroduceAck`
1383        let _: IntroduceAck = self
1384            .runtime
1385            .timeout(timeout_roundtrip, intro_ack_rx.recv(failed_map_err))
1386            .await
1387            .map_err(|_timeout: TimeoutError| FAE::IntroductionTimeout { intro_index })??
1388            .success()
1389            .map_err(|status| FAE::IntroductionFailed {
1390                status,
1391                intro_index,
1392            })?;
1393
1394        debug!(
1395            "hs conn to {}: RPT {} IPT {}: making introduction - success",
1396            &self.hsid,
1397            rend_pt.as_inner(),
1398            intro_index,
1399        );
1400
1401        // Having received INTRODUCE_ACK. we can forget about this circuit
1402        // (and potentially tear it down).
1403        drop(intro_circ);
1404
1405        Ok((
1406            rendezvous,
1407            Introduced {
1408                handshake_state,
1409                marker: PhantomData,
1410                peer_caps,
1411            },
1412        ))
1413    }
1414
1415    /// Attempt (once) to connect a rendezvous circuit using the given intro pt.
1416    ///
1417    /// That is to say, we simply wait for a RENDEZVOUS2 message,
1418    /// and if we get one, we add a virtual hop.
1419    ///
1420    /// Timeouts here might be due to the IPT, RPT, service,
1421    /// or any of the intermediate relays.
1422    ///
1423    /// If, rather than a timeout, we actually encounter some kind of error,
1424    /// we'll return the appropriate `FailedAttemptError`.
1425    /// (Who is responsible may vary, so the `FailedAttemptError` variant will reflect that.)
1426    async fn complete_rendezvous(
1427        &'c self,
1428        ipt: &UsableIntroPt<'_>,
1429        rendezvous: Rendezvous<'c, R, M>,
1430        introduced: Introduced<R, M>,
1431        is_single_onion_service: bool,
1432    ) -> Result<DataTunnel!(R, M), FAE> {
1433        /// Largest number of hops that the onion service must build for _its_
1434        /// circuits to our rendezvous points.
1435        ///
1436        /// This is 4 hops (assuming that it has full vanguards enabled) plus one for the
1437        /// renedezvous point itself.
1438        const MAX_PEER_REND_HOPS: usize = 5;
1439
1440        /// Largest number of retries that we think the peer might make if its
1441        /// circuits are failing.
1442        const MAX_PEER_CIRC_RETRIES: u32 = 3;
1443
1444        let rend_pt = rend_pt_identity_for_error(&rendezvous.rend_relay);
1445        let intro_index = ipt.intro_index;
1446        let failed_map_err = |error| FAE::RendezvousCompletionCircuitError {
1447            error,
1448            intro_index,
1449            rend_pt: rend_pt.clone(),
1450        };
1451
1452        debug!(
1453            "hs conn to {}: RPT {} IPT {}: awaiting rendezvous completion",
1454            &self.hsid,
1455            rend_pt.as_inner(),
1456            intro_index,
1457        );
1458
1459        let num_hops = rendezvous
1460            .rend_tunnel
1461            .m_num_own_hops()
1462            // This is not necessarily the best error, but it isn't totally wrong.
1463            // We can't wrap the tor_circuit error in anything else that makes sense.
1464            // See #2513.
1465            .map_err(|error| FAE::RendezvousCircuitObtain { error })?;
1466
1467        // Maximum length of the circuit that the peer will build to the rendezvous point.
1468        let peer_rend_circ_len = if is_single_onion_service {
1469            1
1470        } else {
1471            MAX_PEER_REND_HOPS
1472        };
1473
1474        // The total number of hops from the peer to us.
1475        //
1476        // We subtract 1 because both circuits terminate at the rendezvous point.
1477        let total_circ_len = peer_rend_circ_len + num_hops - 1;
1478
1479        // Limit on the duration of each attempt for activities involving both
1480        // RPT and IPT.
1481        let rpt_ipt_timeout = self.estimate_timeout(&[
1482            // The API requires us to specify a number of circuit builds and round trips.
1483            // So what we tell the estimator is a rather imprecise description.
1484            //
1485            // What we are timing here is:
1486            //
1487            //    INTRODUCE2 goes from IPT to HS.
1488            //    This happens in parallel with our waiting for the INTRODUCE_ACK,
1489            //    and we know that our own introduction circuit is always at least
1490            //    as long as the peer's (even if they are using full vanguards),
1491            //    so we don't need any additional delay here.
1492            //
1493            //    HS builds to our RPT
1494            (
1495                MAX_PEER_CIRC_RETRIES,
1496                TimeoutsAction::BuildCircuit {
1497                    length: peer_rend_circ_len,
1498                },
1499            ),
1500            //
1501            //    RENDEZVOUS1 goes from HS to RPT.  `peer_circ_len`, one-way.
1502            //    RENDEZVOUS2 goes from RPT to us.  `num_hops`, one-way.
1503            (
1504                1,
1505                TimeoutsAction::OneWay {
1506                    length: total_circ_len,
1507                },
1508            ),
1509        ]);
1510
1511        let rend2_msg: Rendezvous2 = self
1512            .runtime
1513            .timeout(rpt_ipt_timeout, rendezvous.rend2_rx.recv(failed_map_err))
1514            .await
1515            .map_err(|_: TimeoutError| FAE::RendezvousCompletionTimeout {
1516                intro_index,
1517                rend_pt: rend_pt.clone(),
1518            })??;
1519
1520        debug!(
1521            "hs conn to {}: RPT {} IPT {}: received RENDEZVOUS2",
1522            &self.hsid,
1523            rend_pt.as_inner(),
1524            intro_index,
1525        );
1526
1527        // In theory would be great if we could have multiple introduction attempts in parallel
1528        // with similar x,X values but different IPTs.  However, our HS experts don't
1529        // think increasing parallelism here is important:
1530        //   https://gitlab.torproject.org/tpo/core/arti/-/issues/913#note_2914438
1531        let handshake_state = introduced.handshake_state;
1532
1533        // Try to complete the cryptographic handshake.
1534        let keygen =
1535            self.mocks
1536                .rendezvous_handshake(handshake_state, rend2_msg, intro_index, &rend_pt)?;
1537
1538        let mut params = onion_circparams_from_netparams(self.netdir.params())
1539            .map_err(into_internal!("Failed to build CircParameters"))?;
1540
1541        if let Some(inc) = introduced.peer_caps.cc_sendme_inc() {
1542            params.ccontrol.override_sendme_inc(inc);
1543        }
1544
1545        let protocols = introduced.peer_caps.shared_protos();
1546
1547        rendezvous
1548            .rend_tunnel
1549            .m_extend_virtual(
1550                handshake::RelayProtocol::HsV3,
1551                handshake::HandshakeRole::Initiator,
1552                keygen,
1553                params,
1554                &protocols,
1555            )
1556            .await
1557            .map_err(into_internal!(
1558                "actually this is probably a 'circuit closed' error" // TODO HS
1559            ))?;
1560
1561        debug!(
1562            "hs conn to {}: RPT {} IPT {}: HS circuit established",
1563            &self.hsid,
1564            rend_pt.as_inner(),
1565            intro_index,
1566        );
1567
1568        Ok(rendezvous.rend_tunnel)
1569    }
1570
1571    /// Helper to estimate a timeout for a complicated operation
1572    ///
1573    /// `actions` is a list of `(count, action)`, where each entry
1574    /// represents doing `action`, `count` times sequentially.
1575    ///
1576    /// Combines the timeout estimates and returns an overall timeout.
1577    fn estimate_timeout(&self, actions: &[(u32, TimeoutsAction)]) -> Duration {
1578        // This algorithm is, perhaps, wrong.  For uncorrelated variables, a particular
1579        // percentile estimate for a sum of random variables, is not calculated by adding the
1580        // percentile estimates of the individual variables.
1581        //
1582        // But the actual lengths of times of the operations aren't uncorrelated.
1583        // If they were *perfectly* correlated, then this addition would be correct.
1584        // It will do for now; it just might be rather longer than it ought to be.
1585        actions
1586            .iter()
1587            .map(|(count, action)| {
1588                self.circpool
1589                    .m_estimate_timeout(action)
1590                    .saturating_mul(*count)
1591            })
1592            .fold(Duration::ZERO, Duration::saturating_add)
1593    }
1594}
1595
1596/// Mocks used for testing `connect.rs`
1597///
1598/// This is different to `MockableConnectorData`,
1599/// which is used to *replace* this file, when testing `state.rs`.
1600///
1601/// `MocksForConnect` provides mock facilities for *testing* this file.
1602//
1603// TODO this should probably live somewhere else, maybe tor-circmgr even?
1604// TODO this really ought to be made by macros or something
1605trait MocksForConnect<R>: Clone {
1606    /// HS circuit pool
1607    type HsCircPool: MockableCircPool<R>;
1608
1609    /// A random number generator
1610    type Rng: rand::Rng + rand::CryptoRng;
1611
1612    /// Key generator used for generating the keys for the virtual hop.
1613    type KeyGenerator: tor_proto::client::circuit::handshake::KeyGenerator + Send;
1614
1615    /// Tell tests we got this descriptor text
1616    fn test_got_desc(&self, _: &HsDesc) {}
1617    /// Tell tests we got this data tunnel.
1618    fn test_got_tunnel(&self, _: &DataTunnel!(R, Self)) {}
1619    /// Tell tests we have obtained and sorted the intros like this
1620    fn test_got_ipts(&self, _: &[UsableIntroPt]) {}
1621
1622    /// Return a random number generator
1623    fn thread_rng(&self) -> Self::Rng;
1624
1625    /// Complete the rendezvous handshake, returning the resulting keygen
1626    fn rendezvous_handshake(
1627        &self,
1628        handshake_state: hs_ntor::HsNtorClientState,
1629        rend2_msg: Rendezvous2,
1630        intro_index: IntroPtIndex,
1631        rend_pt: &RendPtIdentityForError,
1632    ) -> Result<Self::KeyGenerator, FAE>;
1633}
1634/// Mock for `HsCircPool`
1635///
1636/// Methods start with `m_` to avoid the following problem:
1637/// `ClientCirc::start_conversation` (say) means
1638/// to use the inherent method if one exists,
1639/// but will use a trait method if there isn't an inherent method.
1640///
1641/// So if the inherent method is renamed, the call in the impl here
1642/// turns into an always-recursive call.
1643/// This is not detected by the compiler due to the situation being
1644/// complicated by futures, `#[async_trait]` etc.
1645/// <https://github.com/rust-lang/rust/issues/111177>
1646#[async_trait]
1647trait MockableCircPool<R> {
1648    /// Directory tunnel.
1649    type DirTunnel: MockableClientDir;
1650    /// Data tunnel.
1651    type DataTunnel: MockableClientData;
1652    /// Intro tunnel.
1653    type IntroTunnel: MockableClientIntro;
1654
1655    async fn m_get_or_launch_dir(
1656        &self,
1657        netdir: &NetDir,
1658        target: impl CircTarget + Send + Sync + 'async_trait,
1659    ) -> tor_circmgr::Result<Self::DirTunnel>;
1660
1661    async fn m_get_or_launch_intro(
1662        &self,
1663        netdir: &NetDir,
1664        target: impl CircTarget + Send + Sync + 'async_trait,
1665    ) -> tor_circmgr::Result<Self::IntroTunnel>;
1666
1667    /// Client circuit
1668    async fn m_get_or_launch_client_rend<'a>(
1669        &self,
1670        netdir: &'a NetDir,
1671    ) -> tor_circmgr::Result<(Self::DataTunnel, Relay<'a>)>;
1672
1673    /// Estimate timeout
1674    fn m_estimate_timeout(&self, action: &TimeoutsAction) -> Duration;
1675}
1676
1677/// Mock for onion service client directory tunnel.
1678#[async_trait]
1679trait MockableClientDir: Debug {
1680    /// Client circuit
1681    type DirStream: AsyncRead + AsyncWrite + Send + Unpin;
1682    async fn m_begin_dir_stream(&self) -> tor_circmgr::Result<Self::DirStream>;
1683
1684    /// Get a tor_dirclient::SourceInfo for this circuit, if possible.
1685    fn m_source_info(&self) -> tor_proto::Result<Option<SourceInfo>>;
1686
1687    /// Return the length of this circuit.
1688    fn m_num_hops(&self) -> tor_circmgr::Result<usize>;
1689}
1690
1691/// Mock for onion service client data tunnel.
1692#[async_trait]
1693trait MockableClientData: Debug {
1694    /// Conversation
1695    type Conversation<'r>
1696    where
1697        Self: 'r;
1698    /// Converse
1699    async fn m_start_conversation_last_hop(
1700        &self,
1701        msg: Option<AnyRelayMsg>,
1702        reply_handler: impl MsgHandler + Send + 'static,
1703    ) -> tor_circmgr::Result<Self::Conversation<'_>>;
1704
1705    /// Add a virtual hop to the circuit.
1706    async fn m_extend_virtual(
1707        &self,
1708        protocol: handshake::RelayProtocol,
1709        role: handshake::HandshakeRole,
1710        handshake: impl handshake::KeyGenerator + Send,
1711        params: CircParameters,
1712        capabilities: &tor_protover::Protocols,
1713    ) -> tor_circmgr::Result<()>;
1714
1715    /// Return the number of our own hops in this circuit.
1716    ///
1717    /// This does not count any hops for the service's rendezvous circuit.
1718    /// It does count our virtual hop, if we have one.
1719    /// (That isn't a problem, since we only use this method to calculate
1720    /// timeouts, and we only calculate timeouts _before_ we establish
1721    /// the virtual hop.)
1722    fn m_num_own_hops(&self) -> tor_circmgr::Result<usize>;
1723}
1724
1725/// Mock for onion service client introduction tunnel.
1726#[async_trait]
1727trait MockableClientIntro: Debug {
1728    /// Conversation
1729    type Conversation<'r>
1730    where
1731        Self: 'r;
1732    /// Converse
1733    async fn m_start_conversation_last_hop(
1734        &self,
1735        msg: Option<AnyRelayMsg>,
1736        reply_handler: impl MsgHandler + Send + 'static,
1737    ) -> tor_circmgr::Result<Self::Conversation<'_>>;
1738
1739    /// Return the number of hops in this circuit.
1740    fn m_num_hops(&self) -> tor_circmgr::Result<usize>;
1741}
1742
1743impl<R: Runtime> MocksForConnect<R> for () {
1744    type HsCircPool = HsCircPool<R>;
1745    type Rng = rand::rngs::ThreadRng;
1746    type KeyGenerator = HsNtorHkdfKeyGenerator;
1747
1748    fn thread_rng(&self) -> Self::Rng {
1749        rand::rng()
1750    }
1751
1752    fn rendezvous_handshake(
1753        &self,
1754        handshake_state: hs_ntor::HsNtorClientState,
1755        rend2_msg: Rendezvous2,
1756        intro_index: IntroPtIndex,
1757        rend_pt: &RendPtIdentityForError,
1758    ) -> Result<Self::KeyGenerator, FAE> {
1759        // Try to complete the cryptographic handshake.
1760        handshake_state
1761            .client_receive_rend(rend2_msg.handshake_info())
1762            // If this goes wrong. either the onion service has mangled the crypto,
1763            // or the rendezvous point has misbehaved (that that is possible is a protocol bug),
1764            // or we have used the wrong handshake_state (let's assume that's not true).
1765            //
1766            // If this happens we'll go and try another RPT.
1767            .map_err(|error| FAE::RendezvousCompletionHandshake {
1768                error,
1769                intro_index,
1770                rend_pt: rend_pt.clone(),
1771            })
1772    }
1773}
1774#[async_trait]
1775impl<R: Runtime> MockableCircPool<R> for HsCircPool<R> {
1776    type DirTunnel = ClientOnionServiceDirTunnel;
1777    type DataTunnel = ClientOnionServiceDataTunnel;
1778    type IntroTunnel = ClientOnionServiceIntroTunnel;
1779
1780    #[instrument(level = "trace", skip_all)]
1781    async fn m_get_or_launch_dir(
1782        &self,
1783        netdir: &NetDir,
1784        target: impl CircTarget + Send + Sync + 'async_trait,
1785    ) -> tor_circmgr::Result<Self::DirTunnel> {
1786        Ok(HsCircPool::get_or_launch_client_dir(self, netdir, target).await?)
1787    }
1788    #[instrument(level = "trace", skip_all)]
1789    async fn m_get_or_launch_intro(
1790        &self,
1791        netdir: &NetDir,
1792        target: impl CircTarget + Send + Sync + 'async_trait,
1793    ) -> tor_circmgr::Result<Self::IntroTunnel> {
1794        Ok(HsCircPool::get_or_launch_client_intro(self, netdir, target).await?)
1795    }
1796    #[instrument(level = "trace", skip_all)]
1797    async fn m_get_or_launch_client_rend<'a>(
1798        &self,
1799        netdir: &'a NetDir,
1800    ) -> tor_circmgr::Result<(Self::DataTunnel, Relay<'a>)> {
1801        HsCircPool::get_or_launch_client_rend(self, netdir).await
1802    }
1803    fn m_estimate_timeout(&self, action: &TimeoutsAction) -> Duration {
1804        HsCircPool::estimate_timeout(self, action)
1805    }
1806}
1807#[async_trait]
1808impl MockableClientDir for ClientOnionServiceDirTunnel {
1809    /// Client circuit
1810    type DirStream = tor_proto::client::stream::DataStream;
1811    async fn m_begin_dir_stream(&self) -> tor_circmgr::Result<Self::DirStream> {
1812        Self::begin_dir_stream(self).await
1813    }
1814
1815    /// Get a tor_dirclient::SourceInfo for this circuit, if possible.
1816    fn m_source_info(&self) -> tor_proto::Result<Option<SourceInfo>> {
1817        SourceInfo::from_tunnel(self)
1818    }
1819
1820    fn m_num_hops(&self) -> tor_circmgr::Result<usize> {
1821        self.n_hops()
1822    }
1823}
1824
1825#[async_trait]
1826impl MockableClientData for ClientOnionServiceDataTunnel {
1827    type Conversation<'r> = tor_proto::Conversation<'r>;
1828
1829    async fn m_start_conversation_last_hop(
1830        &self,
1831        msg: Option<AnyRelayMsg>,
1832        reply_handler: impl MsgHandler + Send + 'static,
1833    ) -> tor_circmgr::Result<Self::Conversation<'_>> {
1834        Self::start_conversation(self, msg, reply_handler, TargetHop::LastHop).await
1835    }
1836
1837    async fn m_extend_virtual(
1838        &self,
1839        protocol: handshake::RelayProtocol,
1840        role: handshake::HandshakeRole,
1841        handshake: impl handshake::KeyGenerator + Send,
1842        params: CircParameters,
1843        capabilities: &tor_protover::Protocols,
1844    ) -> tor_circmgr::Result<()> {
1845        Self::extend_virtual(self, protocol, role, handshake, params, capabilities).await
1846    }
1847
1848    fn m_num_own_hops(&self) -> tor_circmgr::Result<usize> {
1849        self.n_hops()
1850    }
1851}
1852
1853#[async_trait]
1854impl MockableClientIntro for ClientOnionServiceIntroTunnel {
1855    type Conversation<'r> = tor_proto::Conversation<'r>;
1856
1857    async fn m_start_conversation_last_hop(
1858        &self,
1859        msg: Option<AnyRelayMsg>,
1860        reply_handler: impl MsgHandler + Send + 'static,
1861    ) -> tor_circmgr::Result<Self::Conversation<'_>> {
1862        Self::start_conversation(self, msg, reply_handler, TargetHop::LastHop).await
1863    }
1864
1865    fn m_num_hops(&self) -> tor_circmgr::Result<usize> {
1866        self.n_hops()
1867    }
1868}
1869
1870#[async_trait]
1871impl MockableConnectorData for Data {
1872    type DataTunnel = ClientOnionServiceDataTunnel;
1873    type MockGlobalState = ();
1874
1875    async fn connect<R: Runtime>(
1876        connector: &HsClientConnector<R>,
1877        netdir: Arc<NetDir>,
1878        config: Arc<Config>,
1879        hsid: HsId,
1880        data: &mut Self,
1881        secret_keys: HsClientSecretKeys,
1882    ) -> Result<Self::DataTunnel, ConnError> {
1883        connect(connector, netdir, config, hsid, data, secret_keys).await
1884    }
1885
1886    fn tunnel_is_ok(tunnel: &Self::DataTunnel) -> bool {
1887        !tunnel.is_closed()
1888    }
1889}
1890
1891/// Return an error if `desc` declares a set of parameters
1892/// that is too far away from the consensus.
1893fn ensure_descriptor_compatible_with_params(
1894    desc: &HsDesc,
1895    params: &NetParameters,
1896) -> Result<(), DescriptorErrorDetail> {
1897    if let Some((_, desc_inc)) = desc.flow_control() {
1898        let difference = u8::from(desc_inc).abs_diff(params.cc_sendme_inc.into());
1899        if difference > 1 {
1900            return Err(DescriptorErrorDetail::ParameterMismatch(
1901                "cc_sendme_inc too far from consensus".into(),
1902            ));
1903        }
1904    }
1905    Ok(())
1906}
1907
1908#[cfg(test)]
1909mod test {
1910    // @@ begin test lint list maintained by maint/add_warning @@
1911    #![allow(clippy::bool_assert_comparison)]
1912    #![allow(clippy::clone_on_copy)]
1913    #![allow(clippy::dbg_macro)]
1914    #![allow(clippy::mixed_attributes_style)]
1915    #![allow(clippy::print_stderr)]
1916    #![allow(clippy::print_stdout)]
1917    #![allow(clippy::single_char_pattern)]
1918    #![allow(clippy::unwrap_used)]
1919    #![allow(clippy::unchecked_time_subtraction)]
1920    #![allow(clippy::useless_vec)]
1921    #![allow(clippy::needless_pass_by_value)]
1922    #![allow(clippy::string_slice)] // See arti#2571
1923    //! <!-- @@ end test lint list maintained by maint/add_warning @@ -->
1924
1925    #![allow(dead_code, unused_variables)] // TODO HS TESTS delete, after tests are completed
1926
1927    use super::*;
1928    use crate::*;
1929    use itertools::chain;
1930    use std::iter;
1931    use tokio_crate as tokio;
1932    use tor_async_utils::JoinReadWrite;
1933    use tor_basic_utils::test_rng::{TestingRng, testing_rng};
1934    use tor_hscrypto::pk::{HsClientDescEncKey, HsClientDescEncKeypair};
1935    use tor_llcrypto::pk::curve25519;
1936    use tor_netdoc::doc::{hsdesc::test_data, netstatus::Lifetime};
1937    use tor_rtcompat::RuntimeSubstExt as _;
1938    use tor_rtcompat::tokio::TokioNativeTlsRuntime;
1939    use tor_rtmock::simple_time::SimpleMockTimeProvider;
1940    use tracing_test::traced_test;
1941
1942    #[derive(derive_more::Debug, Default)]
1943    struct MocksGlobal {
1944        hsdirs_asked: Vec<OwnedCircTarget>,
1945        got_desc: Option<HsDesc>,
1946        #[debug(skip)]
1947        rendezvous: Option<Box<dyn MsgHandler + Send + 'static>>,
1948        intro_acks: Vec<(IntroduceAck, MetaCellDisposition)>,
1949    }
1950
1951    #[derive(Clone, Debug)]
1952    struct Mocks<I> {
1953        mglobal: Arc<Mutex<MocksGlobal>>,
1954        id: I,
1955    }
1956
1957    struct MockKeyGenerator;
1958
1959    impl handshake::KeyGenerator for MockKeyGenerator {
1960        fn expand(self, _keylen: usize) -> tor_proto::Result<tor_bytes::SecretBuf> {
1961            todo!()
1962        }
1963    }
1964
1965    impl<R: Runtime> MocksForConnect<R> for Mocks<()> {
1966        type HsCircPool = Mocks<()>;
1967        type Rng = TestingRng;
1968        type KeyGenerator = MockKeyGenerator;
1969
1970        fn test_got_desc(&self, desc: &HsDesc) {
1971            self.mglobal.lock().unwrap().got_desc = Some(desc.clone());
1972        }
1973
1974        fn test_got_ipts(&self, desc: &[UsableIntroPt]) {}
1975
1976        fn thread_rng(&self) -> Self::Rng {
1977            testing_rng()
1978        }
1979
1980        fn rendezvous_handshake(
1981            &self,
1982            _handshake_state: hs_ntor::HsNtorClientState,
1983            _rend2_msg: Rendezvous2,
1984            _intro_index: IntroPtIndex,
1985            _rend_pt: &RendPtIdentityForError,
1986        ) -> Result<Self::KeyGenerator, FAE> {
1987            Ok(MockKeyGenerator)
1988        }
1989    }
1990    #[async_trait]
1991    impl<R: Runtime> MockableCircPool<R> for Mocks<()> {
1992        type DataTunnel = Mocks<()>;
1993        type DirTunnel = Mocks<()>;
1994        type IntroTunnel = Mocks<()>;
1995
1996        async fn m_get_or_launch_dir(
1997            &self,
1998            _netdir: &NetDir,
1999            target: impl CircTarget + Send + Sync + 'async_trait,
2000        ) -> tor_circmgr::Result<Self::DirTunnel> {
2001            let target = OwnedCircTarget::from_circ_target(&target);
2002            self.mglobal.lock().unwrap().hsdirs_asked.push(target);
2003            Ok(self.clone())
2004        }
2005        async fn m_get_or_launch_intro(
2006            &self,
2007            _netdir: &NetDir,
2008            target: impl CircTarget + Send + Sync + 'async_trait,
2009        ) -> tor_circmgr::Result<Self::IntroTunnel> {
2010            Ok(self.clone())
2011        }
2012        /// Client circuit
2013        async fn m_get_or_launch_client_rend<'a>(
2014            &self,
2015            netdir: &'a NetDir,
2016        ) -> tor_circmgr::Result<(Self::DataTunnel, Relay<'a>)> {
2017            // Pick one of the relays we know to be in the test net as our RPT
2018            let rpt = netdir.by_id(&Ed25519Identity::from([12; 32])).unwrap();
2019
2020            Ok((self.clone(), rpt))
2021        }
2022
2023        fn m_estimate_timeout(&self, action: &TimeoutsAction) -> Duration {
2024            Duration::from_secs(10)
2025        }
2026    }
2027    #[async_trait]
2028    impl MockableClientDir for Mocks<()> {
2029        type DirStream = JoinReadWrite<futures::io::Cursor<Box<[u8]>>, futures::io::Sink>;
2030        async fn m_begin_dir_stream(&self) -> tor_circmgr::Result<Self::DirStream> {
2031            let response = format!(
2032                r#"HTTP/1.1 200 OK
2033
2034{}"#,
2035                test_data::TEST_DATA_2
2036            )
2037            .into_bytes()
2038            .into_boxed_slice();
2039
2040            Ok(JoinReadWrite::new(
2041                futures::io::Cursor::new(response),
2042                futures::io::sink(),
2043            ))
2044        }
2045
2046        fn m_source_info(&self) -> tor_proto::Result<Option<SourceInfo>> {
2047            Ok(None)
2048        }
2049
2050        fn m_num_hops(&self) -> tor_circmgr::Result<usize> {
2051            Ok(4)
2052        }
2053    }
2054
2055    #[async_trait]
2056    impl MockableClientData for Mocks<()> {
2057        type Conversation<'r> = &'r ();
2058        async fn m_start_conversation_last_hop(
2059            &self,
2060            msg: Option<AnyRelayMsg>,
2061            mut reply_handler: impl MsgHandler + Send + 'static,
2062        ) -> tor_circmgr::Result<Self::Conversation<'_>> {
2063            match msg {
2064                Some(AnyRelayMsg::EstablishRendezvous(_)) => {
2065                    let reply = RendezvousEstablished::default();
2066                    let disp = reply_handler.handle_msg(reply.into()).unwrap();
2067                    assert_eq!(disp, MetaCellDisposition::Consumed);
2068                    // Save this, because we'll need to use it later,
2069                    // when handling the INTRODUCE1
2070                    let mut global = self.mglobal.lock().unwrap();
2071                    global.rendezvous = Some(Box::new(reply_handler));
2072                }
2073                _ => panic!("unexpected msg {msg:?}"),
2074            }
2075
2076            Ok(&())
2077        }
2078
2079        async fn m_extend_virtual(
2080            &self,
2081            protocol: handshake::RelayProtocol,
2082            role: handshake::HandshakeRole,
2083            handshake: impl handshake::KeyGenerator + Send,
2084            params: CircParameters,
2085            capabilities: &tor_protover::Protocols,
2086        ) -> tor_circmgr::Result<()> {
2087            Ok(())
2088        }
2089
2090        fn m_num_own_hops(&self) -> tor_circmgr::Result<usize> {
2091            Ok(4)
2092        }
2093    }
2094
2095    #[async_trait]
2096    impl MockableClientIntro for Mocks<()> {
2097        type Conversation<'r> = &'r ();
2098        async fn m_start_conversation_last_hop(
2099            &self,
2100            msg: Option<AnyRelayMsg>,
2101            mut reply_handler: impl MsgHandler + Send + 'static,
2102        ) -> tor_circmgr::Result<Self::Conversation<'_>> {
2103            match msg {
2104                Some(AnyRelayMsg::Introduce1(introduce1)) => {
2105                    let mut global = self.mglobal.lock().unwrap();
2106                    let (reply, expected_disp) = global.intro_acks.remove(0);
2107                    let disp = reply_handler.handle_msg(reply.into()).unwrap();
2108                    assert_eq!(disp, expected_disp);
2109
2110                    // Mock the service's response
2111                    let rendezvous = global
2112                        .rendezvous
2113                        .as_mut()
2114                        .expect("got INTRODUCE1 before ESTABLISH_RENDEZVOUS?!");
2115                    let reply = Rendezvous2::new(b"dummy handshake info, ignored");
2116                    let disp = rendezvous.handle_msg(reply.into()).unwrap();
2117                    assert_eq!(disp, MetaCellDisposition::ConversationFinished);
2118                }
2119                _ => panic!("unexpected msg {msg:?}"),
2120            }
2121
2122            Ok(&())
2123        }
2124
2125        fn m_num_hops(&self) -> tor_circmgr::Result<usize> {
2126            Ok(4)
2127        }
2128    }
2129
2130    fn ks_hsc_desc_enc() -> HsClientDescEncKeypair {
2131        let pk: HsClientDescEncKey = curve25519::PublicKey::from(test_data::TEST_PUBKEY_2).into();
2132        let sk = curve25519::StaticSecret::from(test_data::TEST_SECKEY_2).into();
2133        HsClientDescEncKeypair::new(pk, sk)
2134    }
2135
2136    fn expected_hsdesc(hsid: HsId, netdir: &NetDir, now: SystemTime) -> HsDesc {
2137        let time_period = netdir.hs_time_period();
2138        let (hs_blind_id_key, subcredential) = HsIdKey::try_from(hsid)
2139            .unwrap()
2140            .compute_blinded_key(time_period)
2141            .unwrap();
2142        let hs_blind_id = hs_blind_id_key.id();
2143
2144        HsDesc::parse_decrypt_validate(
2145            test_data::TEST_DATA_2,
2146            &hs_blind_id,
2147            &subcredential,
2148            Some(&ks_hsc_desc_enc()),
2149        )
2150        .unwrap()
2151        .if_valid_at(&now)
2152        .unwrap()
2153    }
2154
2155    fn build_test_netdir() -> Arc<NetDir> {
2156        let valid_after = humantime::parse_rfc3339("2023-02-09T12:00:00Z").unwrap();
2157        let fresh_until = valid_after + humantime::parse_duration("1 hours").unwrap();
2158        let valid_until = valid_after + humantime::parse_duration("24 hours").unwrap();
2159        let lifetime = Lifetime::new(valid_after, fresh_until, valid_until).unwrap();
2160
2161        let netdir = tor_netdir::testnet::construct_custom_netdir_with_params(
2162            tor_netdir::testnet::simple_net_func,
2163            iter::empty::<(&str, _)>(),
2164            Some(lifetime),
2165        )
2166        .expect("failed to build default testing netdir");
2167
2168        Arc::new(netdir.unwrap_if_sufficient().unwrap())
2169    }
2170
2171    #[traced_test]
2172    #[tokio::test]
2173    async fn test_connect() {
2174        use MetaCellDisposition::*;
2175        let netdir = build_test_netdir();
2176        let runtime = TokioNativeTlsRuntime::current().unwrap();
2177        let now = humantime::parse_rfc3339("2023-02-09T12:00:00Z").unwrap();
2178        let mock_sp = SimpleMockTimeProvider::from_wallclock(now);
2179        let runtime = runtime
2180            .with_sleep_provider(mock_sp.clone())
2181            .with_coarse_time_provider(mock_sp.clone());
2182
2183        let success = (
2184            IntroduceAck::new(IntroduceAckStatus::SUCCESS),
2185            ConversationFinished,
2186        );
2187
2188        let nack = (
2189            IntroduceAck::new(IntroduceAckStatus::NOT_RECOGNIZED),
2190            ConversationFinished,
2191        );
2192
2193        // The number of times to make Context:connect() fail due to intro NACK
2194        //
2195        // Set to 5 in order to trigger a rate-limit for all 6 HsDirs:
2196        //
2197        // there are 6 HsDirs in total, one of which is "used up" by the
2198        // first (successful) connect() attempt below.
2199        const INTRO_FAIL_COUNT: usize = 5;
2200
2201        /// The number of times we expect the client to retry the
2202        /// introduction per connect() call
2203        /// (it will essentially try two rounds of `intro_rend_connect()`,
2204        /// once with the cached descriptor, and once with the potentially
2205        /// new descriptor).
2206        const IPT_RETRY_COUNT: usize = 12;
2207
2208        // The first introduction will succeed
2209        let intro_acks = chain!(
2210            [&success],
2211            // But the next INTRO_FAIL_COUNT connect() will fail
2212            // (+1 because we want to fail *again*, in order to find
2213            // that there's now a limit on all our HsDirs)
2214            [&nack; IPT_RETRY_COUNT * (INTRO_FAIL_COUNT + 1)],
2215            // One more round of failures, to trigger a refecth after the rate-limit is lifted
2216            [&nack; IPT_RETRY_COUNT - 1],
2217            // After refetching the descriptor, the client will retry the introduction,
2218            // and succeed.
2219            [&success],
2220        )
2221        .cloned()
2222        .collect();
2223
2224        let mglobal = Arc::new(Mutex::new(MocksGlobal {
2225            intro_acks,
2226            ..Default::default()
2227        }));
2228
2229        let mocks = Mocks { mglobal, id: () };
2230        // From C Tor src/test/test_hs_common.c test_build_address
2231        let hsid = test_data::TEST_HSID_2.into();
2232        let mut data = Data::default();
2233        let mut expected_hsdirs_asked = 1;
2234
2235        let mut secret_keys_builder = HsClientSecretKeysBuilder::default();
2236        secret_keys_builder.ks_hsc_desc_enc(ks_hsc_desc_enc());
2237        let secret_keys = secret_keys_builder.build().unwrap();
2238
2239        let ctx = Context::new(
2240            &runtime,
2241            &mocks,
2242            Arc::clone(&netdir),
2243            Default::default(),
2244            hsid,
2245            secret_keys,
2246            mocks.clone(),
2247        )
2248        .unwrap();
2249
2250        let _got = ctx.connect(&mut data).await.unwrap();
2251
2252        // Our mock IPT hasn't sent any NACKs yet
2253        assert!(!logs_contain("NACKed, refetching descriptor and retrying"));
2254
2255        let hsdesc = expected_hsdesc(hsid, &netdir, now);
2256        {
2257            let mglobal = mocks.mglobal.lock().unwrap();
2258            assert_eq!(mglobal.hsdirs_asked.len(), expected_hsdirs_asked);
2259            // TODO hs: here and in other places, consider implementing PartialEq instead, or creating
2260            // an assert_dbg_eq macro (which would be part of a test_helpers crate or something)
2261            assert_eq!(
2262                format!("{:?}", mglobal.got_desc),
2263                format!("{:?}", Some(hsdesc.clone()))
2264            );
2265        }
2266
2267        // Check how long the descriptor is valid for
2268        let (start_time, end_time) = data.desc.as_ref().unwrap().desc.bounds_start_end();
2269        assert_eq!(start_time, None);
2270
2271        let desc_valid_until = humantime::parse_rfc3339("2023-02-11T20:00:00Z").unwrap();
2272        assert_eq!(end_time, Some(desc_valid_until));
2273
2274        // These attempts will all fail due to intro NACK,
2275        // and trigger a rate-limit for all 6 HsDirs
2276        for i in 1..=INTRO_FAIL_COUNT + 1 {
2277            let err = ctx.connect(&mut data).await.unwrap_err();
2278
2279            let is_intro_nack = |e| matches!(e, FAE::IntroductionFailed { status, .. });
2280
2281            // All attempts failed because of our repeated intro NACKs
2282            assert!(matches!(err, CE::Failed(e) if e.clone().into_iter().all(is_intro_nack)));
2283
2284            {
2285                assert!(logs_contain("NACKed, refetching descriptor and retrying"));
2286                let mglobal = mocks.mglobal.lock().unwrap();
2287                // Because all intro attempts failed with NACK (NOT_RECOGNIZED),
2288                // the client must've tried to refetch the descriptor
2289                if i <= INTRO_FAIL_COUNT {
2290                    // No rate limiting yet, so the client must've tried to fetch a new
2291                    // descriptor, before failing again.
2292                    expected_hsdirs_asked += 1;
2293                    assert!(!logs_contain("but all hsdirs are rate-limited"));
2294                    assert_eq!(mglobal.hsdirs_asked.len(), expected_hsdirs_asked);
2295                } else {
2296                    // The final failure won't lead to an HsDir fetch
2297                    // because all HsDirs will be rate-limited at that point
2298                    assert!(logs_contain("but all hsdirs are rate-limited"));
2299                    assert_eq!(mglobal.hsdirs_asked.len(), expected_hsdirs_asked);
2300                }
2301
2302                // Same descriptor each time
2303                // TODO hs: here and in other places, consider implementing PartialEq instead, or creating
2304                // an assert_dbg_eq macro (which would be part of a test_helpers crate or something)
2305                assert_eq!(
2306                    format!("{:?}", mglobal.got_desc),
2307                    format!("{:?}", Some(hsdesc.clone()))
2308                );
2309            }
2310
2311            let (start_time, end_time) = data.desc.as_ref().unwrap().desc.bounds_start_end();
2312            assert_eq!(start_time, None);
2313
2314            let desc_valid_until = humantime::parse_rfc3339("2023-02-11T20:00:00Z").unwrap();
2315            assert_eq!(end_time, Some(desc_valid_until));
2316        }
2317
2318        // By default, the HsDir fetches are rate-limited for 15min
2319        mock_sp.advance(Duration::from_secs(15 * 60));
2320        // Finally, we succeed.
2321        let _got = ctx.connect(&mut data).await.unwrap();
2322
2323        // And it turns out we did, in fact refetch the descriptor
2324
2325        // Finally, we try again, but find that all HsDirs are now rate-limited!
2326        // So now we advance the time to lift the rate limit, and hope that
2327        //
2328        // TODO HS TESTS: we could extend our mock infrastructure
2329        // to support returning a different hsdesc this time,
2330        // with various revision counters, to check that the client is indeed
2331        // keeping the newest one.
2332        {
2333            assert!(logs_contain("NACKed, refetching descriptor and retrying"));
2334            let mglobal = mocks.mglobal.lock().unwrap();
2335            // Because all intro attempts failed with NACK (NOT_RECOGNIZED),
2336            // the client must've tried to refetch the descriptor
2337            expected_hsdirs_asked += 1;
2338            assert_eq!(mglobal.hsdirs_asked.len(), expected_hsdirs_asked);
2339        }
2340
2341        // TODO HS TESTS: check the circuit in got is the one we gave out
2342
2343        // TODO HS TESTS: continue with this
2344    }
2345
2346    // TODO HS TESTS: Test IPT state management and expiry:
2347    //   - obtain a test descriptor with only a broken ipt
2348    //     (broken in the sense that intro can be attempted, but will fail somehow)
2349    //   - try to make a connection and expect it to fail
2350    //   - assert that the ipt data isn't empty
2351    //   - cause the descriptor to expire (advance clock)
2352    //   - start using a mocked RNG if we weren't already and pin its seed here
2353    //   - make a new descriptor with two IPTs: the broken one from earlier, and a new one
2354    //   - make a new connection
2355    //   - use test_got_ipts to check that the random numbers
2356    //     would sort the bad intro first, *and* that the good one is appears first
2357    //   - assert that connection succeeded
2358    //   - cause the circuit and descriptor to expire (advance clock)
2359    //   - go back to the previous descriptor contents, but with a new validity period
2360    //   - try to make a connection
2361    //   - use test_got_ipts to check that only the broken ipt is present
2362
2363    // TODO HS TESTS: test retries (of every retry loop we have here)
2364    // TODO HS TESTS: test error paths
2365}