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}