Skip to main content

tor_hsservice/
ipt_mgr.rs

1//! IPT Manager
2//!
3//! Maintains introduction points and publishes descriptors.
4//! Provides a stream of rendezvous requests.
5//!
6//! See [`IptManager::run_once`] for discussion of the implementation approach.
7
8use rand::RngExt;
9
10use crate::{internal_prelude::*, replay::OpenReplayLogError};
11
12use IptStatusStatus as ISS;
13use TrackedStatus as TS;
14use tor_relay_selection::{RelayExclusion, RelaySelector, RelayUsage};
15
16mod persist;
17pub(crate) use persist::IptStorageHandle;
18
19pub use crate::ipt_establish::IptError;
20
21/// Expiry time to put on an interim descriptor (IPT publication set Uncertain)
22///
23/// (Note that we use the same value in both cases, since it doesn't actually do
24/// much good to have a short expiration time. This expiration time only affects
25/// caches, and we can supersede an old descriptor just by publishing it. Thus,
26/// we pick a uniform publication time as done by the C tor implementation.)
27const IPT_PUBLISH_UNCERTAIN: Duration = Duration::from_secs(3 * 60 * 60); // 3 hours
28/// Expiry time to put on a final descriptor (IPT publication set Certain
29const IPT_PUBLISH_CERTAIN: Duration = IPT_PUBLISH_UNCERTAIN;
30
31//========== data structures ==========
32
33/// IPT Manager (for one hidden service)
34#[derive(Educe)]
35#[educe(Debug(bound))]
36pub(crate) struct IptManager<R, M> {
37    /// Immutable contents
38    imm: Immutable<R>,
39
40    /// Mutable state
41    state: State<R, M>,
42}
43
44/// Immutable contents of an IPT Manager
45///
46/// Contains things inherent to our identity, and
47/// handles to services that we'll be using.
48#[derive(Educe)]
49#[educe(Debug(bound))]
50pub(crate) struct Immutable<R> {
51    /// Runtime
52    #[educe(Debug(ignore))]
53    runtime: R,
54
55    /// Netdir provider
56    #[educe(Debug(ignore))]
57    dirprovider: Arc<dyn NetDirProvider>,
58
59    /// Nickname
60    nick: HsNickname,
61
62    /// Output MPSC for rendezvous requests
63    ///
64    /// Passed to IPT Establishers we create
65    output_rend_reqs: mpsc::Sender<RendRequest>,
66
67    /// Internal channel for updates from IPT Establishers (sender)
68    ///
69    /// When we make a new `IptEstablisher` we use this arrange for
70    /// its status updates to arrive, appropriately tagged, via `status_recv`
71    status_send: mpsc::Sender<(IptLocalId, IptStatus)>,
72
73    /// The key manager.
74    #[educe(Debug(ignore))]
75    keymgr: Arc<KeyMgr>,
76
77    /// Replay log directory
78    ///
79    /// Files are named after the (bare) IptLocalId
80    #[educe(Debug(ignore))]
81    replay_log_dir: tor_persist::state_dir::InstanceRawSubdir,
82
83    /// A sender for updating the status of the onion service.
84    #[educe(Debug(ignore))]
85    status_tx: IptMgrStatusSender,
86}
87
88/// State of an IPT Manager
89#[derive(Educe)]
90#[educe(Debug(bound))]
91pub(crate) struct State<R, M> {
92    /// Source of configuration updates
93    //
94    // TODO #1209 reject reconfigurations we can't cope with
95    // for example, state dir changes will go quite wrong
96    new_configs: watch::Receiver<Arc<OnionServiceConfig>>,
97
98    /// Last configuration update we received
99    ///
100    /// This is the snapshot of the config we are currently using.
101    /// (Doing it this way avoids running our algorithms
102    /// with a mixture of old and new config.)
103    current_config: Arc<OnionServiceConfig>,
104
105    /// Channel for updates from IPT Establishers (receiver)
106    ///
107    /// We arrange for all the updates to be multiplexed,
108    /// as that makes handling them easy in our event loop.
109    status_recv: mpsc::Receiver<(IptLocalId, IptStatus)>,
110
111    /// State: selected relays
112    ///
113    /// We append to this, and call `retain` on it,
114    /// so these are in chronological order of selection.
115    irelays: Vec<IptRelay>,
116
117    /// Did we fail to select a relay last time?
118    ///
119    /// This can only be caused (or triggered) by a busted netdir or config.
120    last_irelay_selection_outcome: Result<(), ()>,
121
122    /// Have we removed any IPTs but not yet cleaned up keys and logfiles?
123    #[educe(Debug(ignore))]
124    ipt_removal_cleanup_needed: bool,
125
126    /// Signal for us to shut down
127    shutdown: broadcast::Receiver<Void>,
128
129    /// The on-disk state storage handle.
130    #[educe(Debug(ignore))]
131    storage: IptStorageHandle,
132
133    /// Mockable state, normally [`Real`]
134    ///
135    /// This is in `State` so it can be passed mutably to tests,
136    /// even though the main code doesn't need `mut`
137    /// since `HsCircPool` is a service with interior mutability.
138    mockable: M,
139
140    /// Runtime (to placate compiler)
141    runtime: PhantomData<R>,
142}
143
144/// One selected relay, at which we are establishing (or relavantly advertised) IPTs
145struct IptRelay {
146    /// The actual relay
147    relay: RelayIds,
148
149    /// The retirement time we selected for this relay
150    planned_retirement: Instant,
151
152    /// IPTs at this relay
153    ///
154    /// At most one will have [`IsCurrent`].
155    ///
156    /// We append to this, and call `retain` on it,
157    /// so these are in chronological order of selection.
158    ipts: Vec<Ipt>,
159}
160
161/// One introduction point, representation in memory
162#[derive(Debug)]
163struct Ipt {
164    /// Local persistent identifier
165    lid: IptLocalId,
166
167    /// Handle for the establisher; we keep this here just for its `Drop` action
168    establisher: Box<ErasedIptEstablisher>,
169
170    /// `KS_hs_ipt_sid`, `KP_hs_ipt_sid`
171    ///
172    /// This is an `Arc` because:
173    ///  * The manager needs a copy so that it can save it to disk.
174    ///  * The establisher needs a copy to actually use.
175    ///  * The underlying secret key type is not `Clone`.
176    k_sid: Arc<HsIntroPtSessionIdKeypair>,
177
178    /// `KS_hss_ntor`, `KP_hss_ntor`
179    k_hss_ntor: Arc<HsSvcNtorKeypair>,
180
181    /// Last information about how it's doing including timing info
182    status_last: TrackedStatus,
183
184    /// Until when ought we to try to maintain it
185    ///
186    /// For introduction points we are publishing,
187    /// this is a copy of the value set by the publisher
188    /// in the `IptSet` we share with the publisher,
189    ///
190    /// (`None` means the IPT has not been advertised at all yet.)
191    ///
192    /// We must duplicate the information because:
193    ///
194    ///  * We can't have it just live in the shared `IptSet`
195    ///    because we need to retain it for no-longer-being published IPTs.
196    ///
197    ///  * We can't have it just live here because the publisher needs to update it.
198    ///
199    /// (An alternative would be to more seriously entangle the manager and publisher.)
200    last_descriptor_expiry_including_slop: Option<Instant>,
201
202    /// Is this IPT current - should we include it in descriptors ?
203    ///
204    /// `None` might mean:
205    ///  * WantsToRetire
206    ///  * We have >N IPTs and we have been using this IPT so long we want to rotate it out
207    ///    (the [`IptRelay`] has reached its `planned_retirement` time)
208    ///  * The IPT has wrong parameters of some kind, and needs to be replaced
209    ///    (Eg, we set it up with the wrong DOS_PARAMS extension)
210    is_current: Option<IsCurrent>,
211}
212
213/// Last information from establisher about an IPT, with timing info added by us
214#[derive(Debug)]
215enum TrackedStatus {
216    /// Corresponds to [`IptStatusStatus::Faulty`]
217    Faulty {
218        /// When we were first told this started to establish, if we know it
219        ///
220        /// This might be an early estimate, which would give an overestimate
221        /// of the establishment time, which is fine.
222        /// Or it might be `Err` meaning we don't know.
223        started: Result<Instant, ()>,
224
225        /// The error, if any.
226        error: Option<IptError>,
227    },
228
229    /// Corresponds to [`IptStatusStatus::Establishing`]
230    Establishing {
231        /// When we were told we started to establish, for calculating `time_to_establish`
232        started: Instant,
233    },
234
235    /// Corresponds to [`IptStatusStatus::Good`]
236    Good {
237        /// How long it took to establish (if we could determine that information)
238        ///
239        /// Can only be `Err` in strange situations.
240        time_to_establish: Result<Duration, ()>,
241
242        /// Details, from the Establisher
243        details: ipt_establish::GoodIptDetails,
244    },
245}
246
247/// Token indicating that this introduction point is current (not Retiring)
248#[derive(Copy, Clone, Debug, Eq, PartialEq, Hash, Ord, PartialOrd)]
249struct IsCurrent;
250
251//---------- related to mockability ----------
252
253/// Type-erased version of `Box<IptEstablisher>`
254///
255/// The real type is `M::IptEstablisher`.
256/// We use `Box<dyn Any>` to avoid propagating the `M` type parameter to `Ipt` etc.
257type ErasedIptEstablisher = dyn Any + Send + Sync + 'static;
258
259/// Mockable state in an IPT Manager - real version
260#[derive(Educe)]
261#[educe(Debug)]
262pub(crate) struct Real<R: Runtime> {
263    /// Circuit pool for circuits we need to make
264    ///
265    /// Passed to the each new Establisher
266    #[educe(Debug(ignore))]
267    pub(crate) circ_pool: Arc<HsCircPool<R>>,
268}
269
270//---------- errors ----------
271
272/// An error that happened while trying to select a relay
273///
274/// Used only within the IPT manager.
275/// Can only be caused by bad netdir or maybe bad config.
276#[derive(Debug, Error)]
277enum ChooseIptError {
278    /// Bad or insufficient netdir
279    #[error("bad or insufficient netdir")]
280    NetDir(#[from] tor_netdir::Error),
281    /// Too few suitable relays
282    #[error("too few suitable relays")]
283    TooFewUsableRelays,
284    /// Time overflow
285    #[error("time overflow (system clock set wrong?)")]
286    TimeOverflow,
287    /// Internal error
288    #[error("internal error")]
289    Bug(#[from] Bug),
290}
291
292/// An error that happened while trying to crate an IPT (at a selected relay)
293///
294/// Used only within the IPT manager.
295#[derive(Clone, Debug, Error)]
296pub(crate) enum CreateIptError {
297    /// Fatal error
298    #[error("fatal error")]
299    Fatal(#[from] FatalError),
300
301    /// Error accessing keystore
302    #[error("problems with keystores")]
303    Keystore(#[from] tor_keymgr::Error),
304
305    /// Error opening the intro request replay log
306    #[error(transparent)]
307    OpenReplayLog(#[from] OpenReplayLogError),
308}
309
310//========== Relays we've chosen, and IPTs ==========
311
312impl IptRelay {
313    /// Get a reference to this IPT relay's current intro point state (if any)
314    ///
315    /// `None` means this IPT has no current introduction points.
316    /// That might be, briefly, because a new intro point needs to be created;
317    /// or it might be because we are retiring the relay.
318    fn current_ipt(&self) -> Option<&Ipt> {
319        self.ipts
320            .iter()
321            .find(|ipt| ipt.is_current == Some(IsCurrent))
322    }
323
324    /// Get a mutable reference to this IPT relay's current intro point state (if any)
325    fn current_ipt_mut(&mut self) -> Option<&mut Ipt> {
326        self.ipts
327            .iter_mut()
328            .find(|ipt| ipt.is_current == Some(IsCurrent))
329    }
330
331    /// Should this IPT Relay be retired ?
332    ///
333    /// This is determined by our IPT relay rotation time.
334    fn should_retire(&self, now: &TrackingNow) -> bool {
335        now > &self.planned_retirement
336    }
337
338    /// Make a new introduction point at this relay
339    ///
340    /// It becomes the current IPT.
341    fn make_new_ipt<R: Runtime, M: Mockable<R>>(
342        &mut self,
343        imm: &Immutable<R>,
344        new_configs: &watch::Receiver<Arc<OnionServiceConfig>>,
345        mockable: &mut M,
346    ) -> Result<(), CreateIptError> {
347        let lid: IptLocalId = mockable.thread_rng().random();
348
349        let ipt = Ipt::start_establisher(
350            imm,
351            new_configs,
352            mockable,
353            &self.relay,
354            lid,
355            Some(IsCurrent),
356            None::<IptExpectExistingKeys>,
357            // None is precisely right: the descriptor hasn't been published.
358            PromiseLastDescriptorExpiryNoneIsGood {},
359        )?;
360
361        self.ipts.push(ipt);
362
363        Ok(())
364    }
365}
366
367/// Token, representing promise by caller of `start_establisher`
368///
369/// Caller who makes one of these structs promises that it is OK for `start_establisher`
370/// to set `last_descriptor_expiry_including_slop` to `None`.
371struct PromiseLastDescriptorExpiryNoneIsGood {}
372
373/// Token telling [`Ipt::start_establisher`] to expect existing keys in the keystore
374#[derive(Debug, Clone, Copy)]
375struct IptExpectExistingKeys;
376
377impl Ipt {
378    /// Start a new IPT establisher, and create and return an `Ipt`
379    #[allow(clippy::too_many_arguments)] // There's only two call sites
380    fn start_establisher<R: Runtime, M: Mockable<R>>(
381        imm: &Immutable<R>,
382        new_configs: &watch::Receiver<Arc<OnionServiceConfig>>,
383        mockable: &mut M,
384        relay: &RelayIds,
385        lid: IptLocalId,
386        is_current: Option<IsCurrent>,
387        expect_existing_keys: Option<IptExpectExistingKeys>,
388        _: PromiseLastDescriptorExpiryNoneIsGood,
389    ) -> Result<Ipt, CreateIptError> {
390        let mut rng = tor_llcrypto::rng::CautiousRng;
391
392        /// Load (from disk) or generate an IPT key with role IptKeyRole::$role
393        ///
394        /// Ideally this would be a closure, but it has to be generic over the
395        /// returned key type.  So it's a macro.  (A proper function would have
396        /// many type parameters and arguments and be quite annoying.)
397        macro_rules! get_or_gen_key { { $Keypair:ty, $role:ident } => { (||{
398            let spec = IptKeySpecifier {
399                nick: imm.nick.clone(),
400                role: IptKeyRole::$role,
401                lid,
402            };
403            // Our desired behaviour:
404            //  expect_existing_keys == None
405            //     The keys shouldn't exist.  Generate and insert.
406            //     If they do exist then things are badly messed up
407            //     (we're creating a new IPT with a fres lid).
408            //     So, then, crash.
409            //  expect_existing_keys == Some(IptExpectExistingKeys)
410            //     The key is supposed to exist.  Load them.
411            //     We ought to have stored them before storing in our on-disk records that
412            //     this IPT exists.  But this could happen due to file deletion or something.
413            //     And we could recover by creating fresh keys, although maybe some clients
414            //     would find the previous keys in old descriptors.
415            //     So if the keys are missing, make and store new ones, logging an error msg.
416            let k: Option<$Keypair> = imm.keymgr.get(&spec)?;
417            let arti_path = || {
418                spec
419                    .arti_path()
420                    .map_err(|e| {
421                        CreateIptError::Fatal(
422                            into_internal!("bad ArtiPath from IPT key spec")(e).into()
423                        )
424                    })
425            };
426            match (expect_existing_keys, k) {
427                (None, None) => { }
428                (Some(_), Some(k)) => return Ok(Arc::new(k)),
429                (None, Some(_)) => {
430                    return Err(FatalError::IptKeysFoundUnexpectedly(arti_path()?).into())
431                },
432                (Some(_), None) => {
433                    error!("bug: HS service {} missing previous key {:?}. Regenerating.",
434                           &imm.nick, arti_path()?);
435                }
436             }
437
438            let res = imm.keymgr.generate::<$Keypair>(
439                &spec,
440                tor_keymgr::KeystoreSelector::Primary,
441                &mut rng,
442                false, /* overwrite */
443            );
444
445            match res {
446                Ok(k) => Ok::<_, CreateIptError>(Arc::new(k)),
447                Err(tor_keymgr::Error::KeyAlreadyExists) => {
448                    Err(FatalError::KeystoreRace { action: "generate", path: arti_path()? }.into() )
449                },
450                Err(e) => Err(e.into()),
451            }
452        })() } }
453
454        let k_hss_ntor = get_or_gen_key!(HsSvcNtorKeypair, KHssNtor)?;
455        let k_sid = get_or_gen_key!(HsIntroPtSessionIdKeypair, KSid)?;
456
457        // we'll treat it as Establishing until we find otherwise
458        let status_last = TS::Establishing {
459            started: imm.runtime.now(),
460        };
461
462        // TODO #1186 Support ephemeral services (without persistent replay log)
463        let replay_log = IptReplayLog::new_logged(&imm.replay_log_dir, &lid)?;
464
465        let params = IptParameters {
466            replay_log,
467            config_rx: new_configs.clone(),
468            netdir_provider: imm.dirprovider.clone(),
469            introduce_tx: imm.output_rend_reqs.clone(),
470            lid,
471            target: relay.clone(),
472            k_sid: k_sid.clone(),
473            k_ntor: Arc::clone(&k_hss_ntor),
474            accepting_requests: ipt_establish::RequestDisposition::NotAdvertised,
475        };
476        let (establisher, mut watch_rx) = mockable.make_new_ipt(imm, params)?;
477
478        // This task will shut down when self.establisher is dropped, causing
479        // watch_tx to close.
480        imm.runtime
481            .spawn({
482                let mut status_send = imm.status_send.clone();
483                async move {
484                    loop {
485                        let Some(status) = watch_rx.next().await else {
486                            trace!("HS service IPT status task: establisher went away");
487                            break;
488                        };
489                        match status_send.send((lid, status)).await {
490                            Ok(()) => {}
491                            Err::<_, mpsc::SendError>(e) => {
492                                // Not using trace_report because SendError isn't HasKind
493                                trace!("HS service IPT status task: manager went away: {e}");
494                                break;
495                            }
496                        }
497                    }
498                }
499            })
500            .map_err(|cause| FatalError::Spawn {
501                spawning: "IPT establisher watch status task",
502                cause: cause.into(),
503            })?;
504
505        let ipt = Ipt {
506            lid,
507            establisher: Box::new(establisher),
508            k_hss_ntor,
509            k_sid,
510            status_last,
511            is_current,
512            last_descriptor_expiry_including_slop: None,
513        };
514
515        debug!(
516            "Hs service {}: {lid:?} establishing {} IPT at relay {}",
517            &imm.nick,
518            match expect_existing_keys {
519                None => "new",
520                Some(_) => "previous",
521            },
522            &relay,
523        );
524
525        Ok(ipt)
526    }
527
528    /// Returns `true` if this IPT has status Good (and should perhaps be published)
529    fn is_good(&self) -> bool {
530        match self.status_last {
531            TS::Good { .. } => true,
532            TS::Establishing { .. } | TS::Faulty { .. } => false,
533        }
534    }
535
536    /// Returns the error, if any, we are currently encountering at this IPT.
537    fn error(&self) -> Option<&IptError> {
538        match &self.status_last {
539            TS::Good { .. } | TS::Establishing { .. } => None,
540            TS::Faulty { error, .. } => error.as_ref(),
541        }
542    }
543
544    /// Construct the information needed by the publisher for this intro point
545    fn for_publish(&self, details: &ipt_establish::GoodIptDetails) -> Result<ipt_set::Ipt, Bug> {
546        let k_sid: &ed25519::Keypair = (*self.k_sid).as_ref();
547        tor_netdoc::doc::hsdesc::IntroPointDesc::builder()
548            .link_specifiers(details.link_specifiers.clone())
549            .ipt_kp_ntor(details.ipt_kp_ntor)
550            .kp_hs_ipt_sid(k_sid.verifying_key().into())
551            .kp_hss_ntor(self.k_hss_ntor.public().clone())
552            .build()
553            .map_err(into_internal!("failed to construct IntroPointDesc"))
554    }
555}
556
557impl HasKind for ChooseIptError {
558    fn kind(&self) -> ErrorKind {
559        use ChooseIptError as E;
560        use ErrorKind as EK;
561        match self {
562            E::NetDir(e) => e.kind(),
563            E::TooFewUsableRelays => EK::TorDirectoryUnusable,
564            E::TimeOverflow => EK::ClockSkew,
565            E::Bug(e) => e.kind(),
566        }
567    }
568}
569
570// This is somewhat abbreviated but it is legible and enough for most purposes.
571impl Debug for IptRelay {
572    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
573        writeln!(f, "IptRelay {}", self.relay)?;
574        write!(
575            f,
576            "          planned_retirement: {:?}",
577            self.planned_retirement
578        )?;
579        for ipt in &self.ipts {
580            write!(
581                f,
582                "\n          ipt {} {} {:?} ldeis={:?}",
583                match ipt.is_current {
584                    Some(IsCurrent) => "cur",
585                    None => "old",
586                },
587                ipt.lid,
588                ipt.status_last,
589                ipt.last_descriptor_expiry_including_slop,
590            )?;
591        }
592        Ok(())
593    }
594}
595
596//========== impls on IptManager and State ==========
597
598impl<R: Runtime, M: Mockable<R>> IptManager<R, M> {
599    //
600    //---------- constructor and setup ----------
601
602    /// Create a new IptManager
603    #[allow(clippy::too_many_arguments)] // this is an internal function with 1 call site
604    pub(crate) fn new(
605        runtime: R,
606        dirprovider: Arc<dyn NetDirProvider>,
607        nick: HsNickname,
608        config: watch::Receiver<Arc<OnionServiceConfig>>,
609        output_rend_reqs: mpsc::Sender<RendRequest>,
610        shutdown: broadcast::Receiver<Void>,
611        state_handle: &tor_persist::state_dir::InstanceStateHandle,
612        mockable: M,
613        keymgr: Arc<KeyMgr>,
614        status_tx: IptMgrStatusSender,
615    ) -> Result<Self, StartupError> {
616        let irelays = vec![]; // See TODO near persist::load call, in launch_background_tasks
617
618        // We don't need buffering; since this is written to by dedicated tasks which
619        // are reading watches.
620        //
621        // Internally-generated status updates (hopefully rate limited?), no need for mq.
622        let (status_send, status_recv) = mpsc_channel_no_memquota(0);
623
624        let storage = state_handle
625            .storage_handle("ipts")
626            .map_err(StartupError::StateDirectoryInaccessible)?;
627
628        let replay_log_dir = state_handle
629            .raw_subdir("iptreplay")
630            .map_err(StartupError::StateDirectoryInaccessible)?;
631
632        let imm = Immutable {
633            runtime,
634            dirprovider,
635            nick,
636            status_send,
637            output_rend_reqs,
638            keymgr,
639            replay_log_dir,
640            status_tx,
641        };
642        let current_config = config.borrow().clone();
643
644        let state = State {
645            current_config,
646            new_configs: config,
647            status_recv,
648            storage,
649            mockable,
650            shutdown,
651            irelays,
652            last_irelay_selection_outcome: Ok(()),
653            ipt_removal_cleanup_needed: false,
654            runtime: PhantomData,
655        };
656        let mgr = IptManager { imm, state };
657
658        Ok(mgr)
659    }
660
661    /// Send the IPT manager off to run and establish intro points
662    pub(crate) fn launch_background_tasks(
663        mut self,
664        mut publisher: IptsManagerView,
665    ) -> Result<(), StartupError> {
666        // TODO maybe this should be done in new(), so we don't have this dummy irelays
667        // but then new() would need the IptsManagerView
668        assert!(self.state.irelays.is_empty());
669        self.state.irelays = persist::load(
670            &self.imm,
671            &self.state.storage,
672            &self.state.new_configs,
673            &mut self.state.mockable,
674            &publisher.borrow_for_read(),
675        )?;
676
677        // Now that we've populated `irelays` and its `ipts` from the on-disk state,
678        // we should check any leftover disk files from previous runs.  Make a note.
679        self.state.ipt_removal_cleanup_needed = true;
680
681        let runtime = self.imm.runtime.clone();
682
683        self.imm.status_tx.send(IptMgrState::Bootstrapping, None);
684
685        // This task will shut down when the RunningOnionService is dropped, causing
686        // self.state.shutdown to become ready.
687        runtime
688            .spawn(self.main_loop_task(publisher))
689            .map_err(|cause| StartupError::Spawn {
690                spawning: "ipt manager",
691                cause: cause.into(),
692            })?;
693        Ok(())
694    }
695
696    //---------- internal utility and helper methods ----------
697
698    /// Iterate over *all* the IPTs we know about
699    ///
700    /// Yields each `IptRelay` at most once.
701    fn all_ipts(&self) -> impl Iterator<Item = (&IptRelay, &Ipt)> {
702        self.state
703            .irelays
704            .iter()
705            .flat_map(|ir| ir.ipts.iter().map(move |ipt| (ir, ipt)))
706    }
707
708    /// Iterate over the *current* IPTs
709    ///
710    /// Yields each `IptRelay` at most once.
711    fn current_ipts(&self) -> impl Iterator<Item = (&IptRelay, &Ipt)> {
712        self.state
713            .irelays
714            .iter()
715            .filter_map(|ir| Some((ir, ir.current_ipt()?)))
716    }
717
718    /// Iterate over the *current* IPTs in `Good` state
719    fn good_ipts(&self) -> impl Iterator<Item = (&IptRelay, &Ipt)> {
720        self.current_ipts().filter(|(_ir, ipt)| ipt.is_good())
721    }
722
723    /// Iterate over the current IPT errors.
724    ///
725    /// Used when reporting our state as [`Recovering`](crate::status::State::Recovering).
726    fn ipt_errors(&self) -> impl Iterator<Item = &IptError> {
727        self.all_ipts().filter_map(|(_ir, ipt)| ipt.error())
728    }
729
730    /// Target number of intro points
731    pub(crate) fn target_n_intro_points(&self) -> usize {
732        self.state.current_config.num_intro_points.into()
733    }
734
735    /// Maximum number of concurrent intro point relays
736    pub(crate) fn max_n_intro_relays(&self) -> usize {
737        let params = self.imm.dirprovider.params();
738        let num_extra = (*params).as_ref().hs_intro_num_extra_intropoints.get() as usize;
739        self.target_n_intro_points() + num_extra
740    }
741
742    //---------- main implementation logic ----------
743
744    /// Make some progress, if possible, and say when to wake up again
745    ///
746    /// Examines the current state and attempts to improve it.
747    ///
748    /// If `idempotently_progress_things_now` makes any changes,
749    /// it will return `None`.
750    /// It should then be called again immediately.
751    ///
752    /// Otherwise, it returns the time in the future when further work ought to be done:
753    /// i.e., the time of the earliest timeout or planned future state change -
754    /// as a [`TrackingNow`].
755    ///
756    /// In that case, the caller must call `compute_iptsetstatus_publish`,
757    /// since the IPT set etc. may have changed.
758    ///
759    /// ### Goals and algorithms
760    ///
761    /// We attempt to maintain a pool of N established and verified IPTs,
762    /// at N IPT Relays.
763    ///
764    /// When we have fewer than N IPT Relays
765    /// that have `Establishing` or `Good` IPTs (see below)
766    /// and fewer than k*N IPT Relays overall,
767    /// we choose a new IPT Relay at random from the consensus
768    /// and try to establish an IPT on it.
769    ///
770    /// (Rationale for the k*N limit:
771    /// we do want to try to replace faulty IPTs, but
772    /// we don't want an attacker to be able to provoke us into
773    /// rapidly churning through IPT candidates.)
774    ///
775    /// When we select a new IPT Relay, we randomly choose a planned replacement time,
776    /// after which it becomes `Retiring`.
777    ///
778    /// Additionally, any IPT becomes `Retiring`
779    /// after it has been used for a certain number of introductions
780    /// (c.f. C Tor `#define INTRO_POINT_MIN_LIFETIME_INTRODUCTIONS 16384`.)
781    /// When this happens we retain the IPT Relay,
782    /// and make new parameters to make a new IPT at the same Relay.
783    ///
784    /// An IPT is removed from our records, and we give up on it,
785    /// when it is no longer `Good` or `Establishing`
786    /// and all descriptors that mentioned it have expired.
787    ///
788    /// (Until all published descriptors mentioning an IPT expire,
789    /// we consider ourselves bound by those previously-published descriptors,
790    /// and try to maintain the IPT.
791    /// TODO: Allegedly this is unnecessary, but I don't see how it could be.)
792    ///
793    /// ### Performance
794    ///
795    /// This function is at worst O(N) where N is the number of IPTs.
796    /// When handling state changes relating to a particular IPT (or IPT relay)
797    /// it needs at most O(1) calls to progress that one IPT to its proper new state.
798    ///
799    /// See the performance note on [`run_once()`](Self::run_once).
800    #[allow(clippy::redundant_closure_call)]
801    fn idempotently_progress_things_now(&mut self) -> Result<Option<TrackingNow>, FatalError> {
802        /// Return value which means "we changed something, please run me again"
803        ///
804        /// In each case, if we make any changes which indicate we might
805        /// want to restart, we `return CONTINUE`, and
806        /// our caller will just call us again.
807        ///
808        /// This approach simplifies the logic: everything here is idempotent.
809        /// (It does mean the algorithm can be quadratic in the number of intro points,
810        /// but that number is reasonably small for a modern computer and the constant
811        /// factor is small too.)
812        const CONTINUE: Result<Option<TrackingNow>, FatalError> = Ok(None);
813
814        // This tracks everything we compare it to, using interior mutability,
815        // so that if there is no work to do and no timeouts have expired,
816        // we know when we will want to wake up.
817        let now = TrackingNow::now(&self.imm.runtime);
818
819        // ---------- collect garbage ----------
820
821        // Rotate out an old IPT(s)
822        for ir in &mut self.state.irelays {
823            if ir.should_retire(&now) {
824                if let Some(ipt) = ir.current_ipt_mut() {
825                    ipt.is_current = None;
826                    return CONTINUE;
827                }
828            }
829        }
830
831        // Forget old IPTs (after the last descriptor mentioning them has expired)
832        for ir in &mut self.state.irelays {
833            // When we drop the Ipt we drop the IptEstablisher, withdrawing the intro point
834            ir.ipts.retain(|ipt| {
835                let keep = ipt.is_current.is_some()
836                    || match ipt.last_descriptor_expiry_including_slop {
837                        None => false,
838                        Some(last) => now < last,
839                    };
840                // This is the only place in the manager where an IPT is dropped,
841                // other than when the whole service is dropped.
842                self.state.ipt_removal_cleanup_needed |= !keep;
843                keep
844            });
845            // No need to return CONTINUE, since there is no other future work implied
846            // by discarding a non-current IPT.
847        }
848
849        // Forget retired IPT relays (all their IPTs are gone)
850        self.state
851            .irelays
852            .retain(|ir| !(ir.should_retire(&now) && ir.ipts.is_empty()));
853        // If we deleted relays, we might want to select new ones.  That happens below.
854
855        // ---------- make progress ----------
856        //
857        // Consider selecting new relays and setting up new IPTs.
858
859        // Create new IPTs at already-chosen relays
860        for ir in &mut self.state.irelays {
861            if !ir.should_retire(&now) && ir.current_ipt_mut().is_none() {
862                // We don't have a current IPT at this relay, but we should.
863                match ir.make_new_ipt(&self.imm, &self.state.new_configs, &mut self.state.mockable)
864                {
865                    Ok(()) => return CONTINUE,
866                    Err(CreateIptError::Fatal(fatal)) => return Err(fatal),
867                    Err(
868                        e @ (CreateIptError::Keystore(_) | CreateIptError::OpenReplayLog { .. }),
869                    ) => {
870                        error_report!(e, "HS {}: failed to prepare new IPT", &self.imm.nick);
871                        // Let's not try any more of this.
872                        // We'll run the rest of our "make progress" algorithms,
873                        // presenting them with possibly-suboptimal state.  That's fine.
874                        // At some point we'll be poked to run again and then we'll retry.
875                        /// Retry no later than this:
876                        const STORAGE_RETRY: Duration = Duration::from_secs(60);
877                        now.update(STORAGE_RETRY);
878                        break;
879                    }
880                }
881            }
882        }
883
884        // Consider choosing a new IPT relay
885        {
886            // block {} prevents use of `n_good_ish_relays` for other (wrong) purposes
887
888            // We optimistically count an Establishing IPT as good-ish;
889            // specifically, for the purposes of deciding whether to select a new
890            // relay because we don't have enough good-looking ones.
891            let n_good_ish_relays = self
892                .current_ipts()
893                .filter(|(_ir, ipt)| match ipt.status_last {
894                    TS::Good { .. } | TS::Establishing { .. } => true,
895                    TS::Faulty { .. } => false,
896                })
897                .count();
898
899            #[allow(clippy::unused_unit, clippy::semicolon_if_nothing_returned)] // in map_err
900            if n_good_ish_relays < self.target_n_intro_points()
901                && self.state.irelays.len() < self.max_n_intro_relays()
902                && self.state.last_irelay_selection_outcome.is_ok()
903            {
904                self.state.last_irelay_selection_outcome = self
905                    .state
906                    .choose_new_ipt_relay(&self.imm, now.instant().get_now_untracked())
907                    .map_err(|error| {
908                        /// Call $report! with the message.
909                        // The macros are annoying and want a cost argument.
910                        macro_rules! report { { $report:ident } => {
911                            $report!(
912                                error,
913                                "HS service {} failed to select IPT relay",
914                                &self.imm.nick,
915                            )
916                        }}
917                        use ChooseIptError as E;
918                        match &error {
919                            E::NetDir(_) => report!(info_report),
920                            _ => report!(error_report),
921                        };
922                        ()
923                    });
924                return CONTINUE;
925            }
926        }
927
928        //---------- caller (run_once) will update publisher, and wait ----------
929
930        Ok(Some(now))
931    }
932
933    /// Import publisher's updates to latest descriptor expiry times
934    ///
935    /// Copies the `last_descriptor_expiry_including_slop` field
936    /// from each ipt in `publish_set` to the corresponding ipt in `self`.
937    ///
938    /// ### Performance
939    ///
940    /// This function is at worst O(N) where N is the number of IPTs.
941    /// See the performance note on [`run_once()`](Self::run_once).
942    fn import_new_expiry_times(irelays: &mut [IptRelay], publish_set: &PublishIptSet) {
943        // Every entry in the PublishIptSet ought to correspond to an ipt in self.
944        //
945        // If there are IPTs in publish_set.last_descriptor_expiry_including_slop
946        // that aren't in self, those are IPTs that we know were published,
947        // but can't establish since we have forgotten their details.
948        //
949        // We are not supposed to allow that to happen:
950        // we save IPTs to disk before we allow them to be published.
951        //
952        // (This invariant is across two data structures:
953        // `ipt_mgr::State` (specifically, `Ipt`) which is modified only here,
954        // and `ipt_set::PublishIptSet` which is shared with the publisher.
955        // See the comments in PublishIptSet.)
956
957        let all_ours = irelays.iter_mut().flat_map(|ir| ir.ipts.iter_mut());
958
959        for ours in all_ours {
960            if let Some(theirs) = publish_set
961                .last_descriptor_expiry_including_slop
962                .get(&ours.lid)
963            {
964                ours.last_descriptor_expiry_including_slop = Some(*theirs);
965            }
966        }
967    }
968
969    /// Expire old entries in publish_set.last_descriptor_expiry_including_slop
970    ///
971    /// Deletes entries where `now` > `last_descriptor_expiry_including_slop`,
972    /// ie, entries where the publication's validity time has expired,
973    /// meaning we don't need to maintain that IPT any more,
974    /// at least, not just because we've published it.
975    ///
976    /// We may expire even entries for IPTs that we, the manager, still want to maintain.
977    /// That's fine: this is (just) the information about what we have previously published.
978    ///
979    /// ### Performance
980    ///
981    /// This function is at worst O(N) where N is the number of IPTs.
982    /// See the performance note on [`run_once()`](Self::run_once).
983    fn expire_old_expiry_times(&self, publish_set: &mut PublishIptSet, now: &TrackingNow) {
984        // We don't want to bother waking up just to expire things,
985        // so use an untracked comparison.
986        let now = now.instant().get_now_untracked();
987
988        publish_set
989            .last_descriptor_expiry_including_slop
990            .retain(|_lid, expiry| *expiry <= now);
991    }
992
993    /// Compute the IPT set to publish, and update the data shared with the publisher
994    ///
995    /// `now` is current time and also the earliest wakeup,
996    /// which we are in the process of planning.
997    /// The noted earliest wakeup can be updated by this function,
998    /// for example, with a future time at which the IPT set ought to be published
999    /// (eg, the status goes from Unknown to Uncertain).
1000    ///
1001    /// ## IPT sets and lifetimes
1002    ///
1003    /// We remember every IPT we have published that is still valid.
1004    ///
1005    /// At each point in time we have an idea of set of IPTs we want to publish.
1006    /// The possibilities are:
1007    ///
1008    ///  * `Certain`:
1009    ///    We are sure of which IPTs we want to publish.
1010    ///    We try to do so, talking to hsdirs as necessary,
1011    ///    updating any existing information.
1012    ///    (We also republish to an hsdir if its descriptor will expire soon,
1013    ///    or we haven't published there since Arti was restarted.)
1014    ///
1015    ///  * `Unknown`:
1016    ///    We have no idea which IPTs to publish.
1017    ///    We leave whatever is on the hsdirs as-is.
1018    ///
1019    ///  * `Uncertain`:
1020    ///    We have some IPTs we could publish,
1021    ///    but we're not confident about them.
1022    ///    We publish these to a particular hsdir if:
1023    ///     - our last-published descriptor has expired
1024    ///     - or it will expire soon
1025    ///     - or if we haven't published since Arti was restarted.
1026    ///
1027    /// The idea of what to publish is calculated as follows:
1028    ///
1029    ///  * If we have at least N `Good` IPTs: `Certain`.
1030    ///    (We publish the "best" N IPTs for some definition of "best".
1031    ///    TODO: should we use the fault count?  recency?)
1032    ///
1033    ///  * Unless we have at least one `Good` IPT: `Unknown`.
1034    ///
1035    ///  * Otherwise: if there are IPTs in `Establishing`,
1036    ///    and they have been in `Establishing` only a short time \[1\]:
1037    ///    `Unknown`; otherwise `Uncertain`.
1038    ///
1039    /// The effect is that we delay publishing an initial descriptor
1040    /// by at most 1x the fastest IPT setup time,
1041    /// at most doubling the initial setup time.
1042    ///
1043    /// Each update to the IPT set that isn't `Unknown` comes with a
1044    /// proposed descriptor expiry time,
1045    /// which is used if the descriptor is to be actually published.
1046    /// The proposed descriptor lifetime for `Uncertain`
1047    /// is the minimum (30 minutes).
1048    /// Otherwise, we double the lifetime each time,
1049    /// unless any IPT in the previous descriptor was declared `Faulty`,
1050    /// in which case we reset it back to the minimum.
1051    /// TODO: Perhaps we should just pick fixed short and long lifetimes instead,
1052    /// to limit distinguishability.
1053    ///
1054    /// (Rationale: if IPTs are regularly misbehaving,
1055    /// we should be cautious and limit our exposure to the damage.)
1056    ///
1057    /// \[1\] NOTE: We wait a "short time" between establishing our first IPT,
1058    /// and publishing an incomplete (<N) descriptor -
1059    /// this is a compromise between
1060    /// availability (publishing as soon as we have any working IPT)
1061    /// and
1062    /// exposure and hsdir load
1063    /// (which would suggest publishing only when our IPT set is stable).
1064    /// One possible strategy is to wait as long again
1065    /// as the time it took to establish our first IPT.
1066    /// Another is to somehow use our circuit timing estimator.
1067    ///
1068    /// ### Performance
1069    ///
1070    /// This function is at worst O(N) where N is the number of IPTs.
1071    /// See the performance note on [`run_once()`](Self::run_once).
1072    #[allow(clippy::unnecessary_wraps)] // for regularity
1073    fn compute_iptsetstatus_publish(
1074        &mut self,
1075        now: &TrackingNow,
1076        publish_set: &mut PublishIptSet,
1077    ) -> Result<(), IptStoreError> {
1078        //---------- tell the publisher what to announce ----------
1079
1080        let very_recently: Option<(TrackingInstantOffsetNow, Duration)> = (|| {
1081            // on time overflow, don't treat any as started establishing very recently
1082
1083            let fastest_good_establish_time = self
1084                .current_ipts()
1085                .filter_map(|(_ir, ipt)| match ipt.status_last {
1086                    TS::Good {
1087                        time_to_establish, ..
1088                    } => Some(time_to_establish.ok()?),
1089                    TS::Establishing { .. } | TS::Faulty { .. } => None,
1090                })
1091                .min()?;
1092
1093            // Rationale:
1094            // we could use circuit timings etc., but arguably the actual time to establish
1095            // our fastest IPT is a better estimator here (and we want an optimistic,
1096            // rather than pessimistic estimate).
1097            //
1098            // This algorithm has potential to publish too early and frequently,
1099            // but our overall rate-limiting should keep it from getting out of hand.
1100            //
1101            // TODO: We might want to make this "1" tuneable, and/or tune the
1102            // algorithm as a whole based on experience.
1103            let wait_more = fastest_good_establish_time * 1;
1104            let very_recently = fastest_good_establish_time.checked_add(wait_more)?;
1105
1106            let very_recently = now.checked_sub(very_recently)?;
1107            Some((very_recently, wait_more))
1108        })();
1109
1110        let started_establishing_very_recently = || {
1111            let (very_recently, wait_more) = very_recently?;
1112            let lid = self
1113                .current_ipts()
1114                .filter_map(|(_ir, ipt)| {
1115                    let started = match ipt.status_last {
1116                        TS::Establishing { started } => Some(started),
1117                        TS::Good { .. } | TS::Faulty { .. } => None,
1118                    }?;
1119
1120                    (started > very_recently).then_some(ipt.lid)
1121                })
1122                .next()?;
1123            Some((lid, wait_more))
1124        };
1125
1126        let n_good_ipts = self.good_ipts().count();
1127        let publish_lifetime = if n_good_ipts >= self.target_n_intro_points() {
1128            // "Certain" - we are sure of which IPTs we want to publish
1129            debug!(
1130                "HS service {}: {} good IPTs, >= target {}, publishing",
1131                &self.imm.nick,
1132                n_good_ipts,
1133                self.target_n_intro_points()
1134            );
1135
1136            self.imm.status_tx.send(IptMgrState::Running, None);
1137
1138            Some(IPT_PUBLISH_CERTAIN)
1139        } else if self.good_ipts().next().is_none()
1140        /* !... .is_empty() */
1141        {
1142            // "Unknown" - we have no idea which IPTs to publish.
1143            debug!("HS service {}: no good IPTs", &self.imm.nick);
1144
1145            self.imm
1146                .status_tx
1147                .send_recovering(self.ipt_errors().cloned().collect_vec());
1148
1149            None
1150        } else if let Some((wait_for, wait_more)) = started_establishing_very_recently() {
1151            // "Unknown" - we say have no idea which IPTs to publish:
1152            // although we have *some* idea, we hold off a bit to see if things improve.
1153            // The wait_more period started counting when the fastest IPT became ready,
1154            // so the printed value isn't an offset from the message timestamp.
1155            debug!(
1156                "HS service {}: {} good IPTs, < target {}, waiting up to {}ms for {:?}",
1157                &self.imm.nick,
1158                n_good_ipts,
1159                self.target_n_intro_points(),
1160                wait_more.as_millis(),
1161                wait_for
1162            );
1163
1164            self.imm
1165                .status_tx
1166                .send_recovering(self.ipt_errors().cloned().collect_vec());
1167
1168            None
1169        } else {
1170            // "Uncertain" - we have some IPTs we could publish, but we're not confident
1171            debug!(
1172                "HS service {}: {} good IPTs, < target {}, publishing what we have",
1173                &self.imm.nick,
1174                n_good_ipts,
1175                self.target_n_intro_points()
1176            );
1177
1178            // We are close to being Running -- we just need more IPTs!
1179            let errors = self.ipt_errors().cloned().collect_vec();
1180            let errors = if errors.is_empty() {
1181                None
1182            } else {
1183                Some(errors)
1184            };
1185
1186            self.imm
1187                .status_tx
1188                .send(IptMgrState::DegradedReachable, errors.map(|e| e.into()));
1189
1190            Some(IPT_PUBLISH_UNCERTAIN)
1191        };
1192
1193        publish_set.ipts = if let Some(lifetime) = publish_lifetime {
1194            let selected = self.publish_set_select();
1195            for ipt in &selected {
1196                self.state.mockable.start_accepting(&*ipt.establisher);
1197            }
1198            Some(Self::make_publish_set(selected, lifetime)?)
1199        } else {
1200            None
1201        };
1202
1203        //---------- store persistent state ----------
1204
1205        persist::store(&self.imm, &mut self.state)?;
1206
1207        Ok(())
1208    }
1209
1210    /// Select IPTs to publish, given that we have decided to publish *something*
1211    ///
1212    /// Calculates set of ipts to publish, selecting up to the target `N`
1213    /// from the available good current IPTs.
1214    /// (Old, non-current IPTs, that we are trying to retire, are never published.)
1215    ///
1216    /// The returned list is in the same order as our data structure:
1217    /// firstly, by the ordering in `State.irelays`, and then within each relay,
1218    /// by the ordering in `IptRelay.ipts`.  Both of these are stable.
1219    ///
1220    /// ### Performance
1221    ///
1222    /// This function is at worst O(N) where N is the number of IPTs.
1223    /// See the performance note on [`run_once()`](Self::run_once).
1224    fn publish_set_select(&self) -> VecDeque<&Ipt> {
1225        /// Good candidate introduction point for publication
1226        type Candidate<'i> = &'i Ipt;
1227
1228        let target_n = self.target_n_intro_points();
1229
1230        let mut candidates: VecDeque<_> = self
1231            .state
1232            .irelays
1233            .iter()
1234            .filter_map(|ir: &_| -> Option<Candidate<'_>> {
1235                let current_ipt = ir.current_ipt()?;
1236                if !current_ipt.is_good() {
1237                    return None;
1238                }
1239                Some(current_ipt)
1240            })
1241            .collect();
1242
1243        // Take the last N good IPT relays
1244        //
1245        // The way we manage irelays means that this is always
1246        // the ones we selected most recently.
1247        //
1248        // TODO SPEC  Publication strategy when we have more than >N IPTs
1249        //
1250        // We could have a number of strategies here.  We could take some timing
1251        // measurements, or use the establishment time, or something; but we don't
1252        // want to add distinguishability.
1253        //
1254        // Another concern is manipulability, but
1255        // We can't be forced to churn because we don't remove relays
1256        // from our list of relays to try to use, other than on our own schedule.
1257        // But we probably won't want to be too reactive to the network environment.
1258        //
1259        // Since we only choose new relays when old ones are to retire, or are faulty,
1260        // choosing the most recently selected, rather than the least recently,
1261        // has the effect of preferring relays we don't know to be faulty,
1262        // to ones we have considered faulty least once.
1263        //
1264        // That's better than the opposite.  Also, choosing more recently selected relays
1265        // for publication may slightly bring forward the time at which all descriptors
1266        // mentioning that relay have expired, and then we can forget about it.
1267        while candidates.len() > target_n {
1268            // WTB: VecDeque::truncate_front
1269            let _: Candidate = candidates.pop_front().expect("empty?!");
1270        }
1271
1272        candidates
1273    }
1274
1275    /// Produce a `publish::IptSet`, from a list of IPT selected for publication
1276    ///
1277    /// Updates each chosen `Ipt`'s `last_descriptor_expiry_including_slop`
1278    ///
1279    /// The returned `IptSet` set is in the same order as `selected`.
1280    ///
1281    /// ### Performance
1282    ///
1283    /// This function is at worst O(N) where N is the number of IPTs.
1284    /// See the performance note on [`run_once()`](Self::run_once).
1285    fn make_publish_set<'i>(
1286        selected: impl IntoIterator<Item = &'i Ipt>,
1287        lifetime: Duration,
1288    ) -> Result<ipt_set::IptSet, FatalError> {
1289        let ipts = selected
1290            .into_iter()
1291            .map(|current_ipt| {
1292                let TS::Good { details, .. } = &current_ipt.status_last else {
1293                    return Err(internal!("was good but now isn't?!").into());
1294                };
1295
1296                let publish = current_ipt.for_publish(details)?;
1297
1298                // last_descriptor_expiry_including_slop was earlier merged in from
1299                // the previous IptSet, and here we copy it back
1300                let publish = ipt_set::IptInSet {
1301                    ipt: publish,
1302                    lid: current_ipt.lid,
1303                };
1304
1305                Ok::<_, FatalError>(publish)
1306            })
1307            .collect::<Result<_, _>>()?;
1308
1309        Ok(ipt_set::IptSet { ipts, lifetime })
1310    }
1311
1312    /// Delete persistent on-disk data (including keys) for old IPTs
1313    ///
1314    /// More precisely, scan places where per-IPT data files live,
1315    /// and delete anything that doesn't correspond to
1316    /// one of the IPTs in our main in-memory data structure.
1317    ///
1318    /// Does *not* deal with deletion of data handled via storage handles
1319    /// (`state_dir::StorageHandle`), `ipt_mgr/persist.rs` etc.;
1320    /// those are one file for each service, so old data is removed as we rewrite them.
1321    ///
1322    /// Does *not* deal with deletion of entire old hidden services.
1323    ///
1324    /// (This function works on the basis of the invariant that every IPT
1325    /// in [`ipt_set::PublishIptSet`] is also an [`Ipt`] in [`ipt_mgr::State`](State).
1326    /// See the comment in [`IptManager::import_new_expiry_times`].
1327    /// If that invariant is violated, we would delete on-disk files for the affected IPTs.
1328    /// That's fine since we couldn't re-establish them anyway.)
1329    fn expire_old_ipts_external_persistent_state(&self) -> Result<(), StateExpiryError> {
1330        self.state
1331            .mockable
1332            .expire_old_ipts_external_persistent_state_hook();
1333
1334        let all_ipts: HashSet<_> = self.all_ipts().map(|(_, ipt)| &ipt.lid).collect();
1335
1336        // Keys
1337
1338        let pat = IptKeySpecifierPattern {
1339            nick: Some(self.imm.nick.clone()),
1340            role: None,
1341            lid: None,
1342        }
1343        .arti_pattern()?;
1344
1345        let found = self.imm.keymgr.list_matching(&pat)?;
1346
1347        for entry in found {
1348            let path = entry.key_path();
1349            // Try to identify this key (including its IptLocalId)
1350            match IptKeySpecifier::try_from(path) {
1351                Ok(spec) if all_ipts.contains(&spec.lid) => continue,
1352                Ok(_) => trace!("deleting key for old IPT: {path}"),
1353                Err(bad) => info!("deleting unrecognised IPT key: {path} ({})", bad.report()),
1354            };
1355            // Not known, remove it
1356            self.imm.keymgr.remove_entry(&entry)?;
1357        }
1358
1359        // IPT replay logs
1360
1361        let handle_rl_err = |operation, path: &Path| {
1362            let path = path.to_owned();
1363            move |source| StateExpiryError::ReplayLog {
1364                operation,
1365                path,
1366                source: Arc::new(source),
1367            }
1368        };
1369
1370        // fs-mistrust doesn't offer CheckedDir::read_this_directory.
1371        // But, we probably don't mind that we're not doing many checks here.
1372        let replay_logs = self.imm.replay_log_dir.as_path();
1373        let replay_logs_dir =
1374            fs::read_dir(replay_logs).map_err(handle_rl_err("open dir", replay_logs))?;
1375
1376        for ent in replay_logs_dir {
1377            let ent = ent.map_err(handle_rl_err("read dir", replay_logs))?;
1378            let leaf = ent.file_name();
1379            // Try to identify this replay logfile (including its IptLocalId)
1380            match IptReplayLog::parse_log_leafname(&leaf) {
1381                Ok(lid) if all_ipts.contains(&lid) => continue,
1382                Ok(_) => trace!(
1383                    leaf = leaf.to_string_lossy().as_ref(),
1384                    "deleting replay log for old IPT"
1385                ),
1386                Err(bad) => info!(
1387                    "deleting garbage in IPT replay log dir: {} ({})",
1388                    leaf.to_string_lossy(),
1389                    bad
1390                ),
1391            }
1392            // Not known, remove it
1393            let path = ent.path();
1394            fs::remove_file(&path).map_err(handle_rl_err("remove", &path))?;
1395        }
1396
1397        Ok(())
1398    }
1399
1400    /// Run one iteration of the loop
1401    ///
1402    /// Either do some work, making changes to our state,
1403    /// or, if there's nothing to be done, wait until there *is* something to do.
1404    ///
1405    /// ### Implementation approach
1406    ///
1407    /// Every time we wake up we idempotently make progress
1408    /// by searching our whole state machine, looking for something to do.
1409    /// If we find something to do, we do that one thing, and search again.
1410    /// When we're done, we unconditionally recalculate the IPTs to publish, and sleep.
1411    ///
1412    /// This approach avoids the need for complicated reasoning about
1413    /// which state updates need to trigger other state updates,
1414    /// and thereby avoids several classes of potential bugs.
1415    /// However, it has some performance implications:
1416    ///
1417    /// ### Performance
1418    ///
1419    /// Events relating to an IPT occur, at worst,
1420    /// at a rate proportional to the current number of IPTs,
1421    /// times the maximum flap rate of any one IPT.
1422    ///
1423    /// [`idempotently_progress_things_now`](Self::idempotently_progress_things_now)
1424    /// can be called more than once for each such event,
1425    /// but only a finite number of times per IPT.
1426    ///
1427    /// Therefore, overall, our work rate is O(N^2) where N is the number of IPTs.
1428    /// We think this is tolerable,
1429    /// but it does mean that the principal functions should be written
1430    /// with an eye to avoiding "accidentally quadratic" algorithms,
1431    /// because that would make the whole manager cubic.
1432    /// Ideally we would avoid O(N.log(N)) algorithms.
1433    ///
1434    /// (Note that the number of IPTs can be significantly larger than
1435    /// the maximum target of 20, if the service is very busy so the intro points
1436    /// are cycling rapidly due to the need to replace the replay database.)
1437    async fn run_once(
1438        &mut self,
1439        // This is a separate argument for borrowck reasons
1440        publisher: &mut IptsManagerView,
1441    ) -> Result<ShutdownStatus, FatalError> {
1442        let now = {
1443            // Block to persuade borrow checker that publish_set isn't
1444            // held over an await point.
1445
1446            let mut publish_set = publisher.borrow_for_update(self.imm.runtime.clone());
1447
1448            Self::import_new_expiry_times(&mut self.state.irelays, &publish_set);
1449
1450            let mut loop_limit = 0..(
1451                // Work we do might be O(number of intro points),
1452                // but we might also have cycled the intro points due to many requests.
1453                // 10K is a guess at a stupid upper bound on the number of times we
1454                // might cycle ipts during a descriptor lifetime.
1455                // We don't need a tight bound; if we're going to crash. we can spin a bit first.
1456                (self.target_n_intro_points() + 1) * 10_000
1457            );
1458            let now = loop {
1459                let _: usize = loop_limit.next().expect("IPT manager is looping");
1460
1461                if let Some(now) = self.idempotently_progress_things_now()? {
1462                    break now;
1463                }
1464            };
1465
1466            // TODO #1214 Maybe something at level Error or Info, for example
1467            // Log an error if everything is terrilbe
1468            //   - we have >=N Faulty IPTs ?
1469            //    we have only Faulty IPTs and can't select another due to 2N limit ?
1470            // Log at info if and when we publish?  Maybe the publisher should do that?
1471
1472            if let Err(operr) = self.compute_iptsetstatus_publish(&now, &mut publish_set) {
1473                // This is not good, is it.
1474                publish_set.ipts = None;
1475                let wait = operr.log_retry_max(&self.imm.nick)?;
1476                now.update(wait);
1477            };
1478
1479            self.expire_old_expiry_times(&mut publish_set, &now);
1480
1481            drop(publish_set); // release lock, and notify publisher of any changes
1482
1483            if self.state.ipt_removal_cleanup_needed {
1484                let outcome = self.expire_old_ipts_external_persistent_state();
1485                log_ratelim!("removing state for old IPT(s)"; outcome);
1486                match outcome {
1487                    Ok(()) => self.state.ipt_removal_cleanup_needed = false,
1488                    Err(_already_logged) => {}
1489                }
1490            }
1491
1492            now
1493        };
1494
1495        assert_ne!(
1496            now.clone().shortest(),
1497            Some(Duration::ZERO),
1498            "IPT manager zero timeout, would loop"
1499        );
1500
1501        let mut new_configs = self.state.new_configs.next().fuse();
1502
1503        select_biased! {
1504            () = now.wait_for_earliest(&self.imm.runtime).fuse() => {},
1505            shutdown = self.state.shutdown.next().fuse() => {
1506                info!("HS service {}: terminating due to shutdown signal", &self.imm.nick);
1507                // We shouldn't be receiving anything on thisi channel.
1508                assert!(shutdown.is_none());
1509                return Ok(ShutdownStatus::Terminate)
1510            },
1511
1512            update = self.state.status_recv.next() => {
1513                let (lid, update) = update.ok_or_else(|| internal!("update mpsc ended!"))?;
1514                self.state.handle_ipt_status_update(&self.imm, lid, update);
1515            }
1516
1517            _dir_event = async {
1518                match self.state.last_irelay_selection_outcome {
1519                    Ok(()) => future::pending().await,
1520                    // This boxes needlessly but it shouldn't really happen
1521                    Err(()) => self.imm.dirprovider.events().next().await,
1522                }
1523            }.fuse() => {
1524                self.state.last_irelay_selection_outcome = Ok(());
1525            }
1526
1527            new_config = new_configs => {
1528                let Some(new_config) = new_config else {
1529                    trace!("HS service {}: terminating due to EOF on config updates stream",
1530                           &self.imm.nick);
1531                    return Ok(ShutdownStatus::Terminate);
1532                };
1533                if let Err(why) = (|| {
1534                    let dos = |config: &OnionServiceConfig| config.dos_extension()
1535                        .map_err(|e| e.report().to_string());
1536                    if dos(&self.state.current_config)? != dos(&new_config)? {
1537                        return Err("DOS parameters (rate limit) changed".to_string());
1538                    }
1539                    Ok(())
1540                })() {
1541                    // We need new IPTs with the new parameters.  (The previously-published
1542                    // IPTs will automatically be retained so long as needed, by the
1543                    // rest of our algorithm.)
1544                    info!("HS service {}: replacing IPTs: {}", &self.imm.nick, &why);
1545                    for ir in &mut self.state.irelays {
1546                        for ipt in &mut ir.ipts {
1547                            ipt.is_current = None;
1548                        }
1549                    }
1550                }
1551                self.state.current_config = new_config;
1552                self.state.last_irelay_selection_outcome = Ok(());
1553            }
1554        }
1555
1556        Ok(ShutdownStatus::Continue)
1557    }
1558
1559    /// IPT Manager main loop, runs as a task
1560    ///
1561    /// Contains the error handling, including catching panics.
1562    async fn main_loop_task(mut self, mut publisher: IptsManagerView) {
1563        loop {
1564            match async {
1565                AssertUnwindSafe(self.run_once(&mut publisher))
1566                    .catch_unwind()
1567                    .await
1568                    .map_err(|_: Box<dyn Any + Send>| internal!("IPT manager crashed"))?
1569            }
1570            .await
1571            {
1572                Err(crash) => {
1573                    error!("bug: HS service {} crashed! {}", &self.imm.nick, crash);
1574
1575                    self.imm.status_tx.send_broken(crash);
1576                    break;
1577                }
1578                Ok(ShutdownStatus::Continue) => continue,
1579                Ok(ShutdownStatus::Terminate) => {
1580                    self.imm.status_tx.send_shutdown();
1581
1582                    break;
1583                }
1584            }
1585        }
1586    }
1587}
1588
1589impl<R: Runtime, M: Mockable<R>> State<R, M> {
1590    /// Find the `Ipt` with persistent local id `lid`
1591    fn ipt_by_lid_mut(&mut self, needle: IptLocalId) -> Option<&mut Ipt> {
1592        self.irelays
1593            .iter_mut()
1594            .find_map(|ir| ir.ipts.iter_mut().find(|ipt| ipt.lid == needle))
1595    }
1596
1597    /// Choose a new relay to use for IPTs
1598    fn choose_new_ipt_relay(
1599        &mut self,
1600        imm: &Immutable<R>,
1601        now: Instant,
1602    ) -> Result<(), ChooseIptError> {
1603        let netdir = imm.dirprovider.timely_netdir()?;
1604
1605        let mut rng = self.mockable.thread_rng();
1606
1607        let relay = {
1608            let exclude_ids = self
1609                .irelays
1610                .iter()
1611                .flat_map(|e| e.relay.identities())
1612                .map(|id| id.to_owned())
1613                .collect();
1614            let selector = RelaySelector::new(
1615                RelayUsage::new_intro_point(),
1616                RelayExclusion::exclude_identities(exclude_ids),
1617            );
1618            selector
1619                .select_relay(&mut rng, &netdir)
1620                .0 // TODO: Someday we might want to report why we rejected everything on failure.
1621                .ok_or(ChooseIptError::TooFewUsableRelays)?
1622        };
1623
1624        let lifetime_low = netdir
1625            .params()
1626            .hs_intro_min_lifetime
1627            .try_into()
1628            .expect("Could not convert param to duration.");
1629        let lifetime_high = netdir
1630            .params()
1631            .hs_intro_max_lifetime
1632            .try_into()
1633            .expect("Could not convert param to duration.");
1634        let lifetime_range: std::ops::RangeInclusive<Duration> = lifetime_low..=lifetime_high;
1635        let retirement = rng
1636            .gen_range_checked(lifetime_range)
1637            // If the range from the consensus is invalid, just pick the high-bound.
1638            .unwrap_or(lifetime_high);
1639        let retirement = now
1640            .checked_add(retirement)
1641            .ok_or(ChooseIptError::TimeOverflow)?;
1642
1643        let new_irelay = IptRelay {
1644            relay: RelayIds::from_relay_ids(&relay),
1645            planned_retirement: retirement,
1646            ipts: vec![],
1647        };
1648        self.irelays.push(new_irelay);
1649
1650        debug!(
1651            "HS service {}: choosing new IPT relay {}",
1652            &imm.nick,
1653            relay.display_relay_ids()
1654        );
1655
1656        Ok(())
1657    }
1658
1659    /// Update `self`'s status tracking for one introduction point
1660    fn handle_ipt_status_update(&mut self, imm: &Immutable<R>, lid: IptLocalId, update: IptStatus) {
1661        let Some(ipt) = self.ipt_by_lid_mut(lid) else {
1662            // update from now-withdrawn IPT, ignore it (can happen due to the IPT being a task)
1663            return;
1664        };
1665
1666        debug!("HS service {}: {lid:?} status update {update:?}", &imm.nick);
1667
1668        let IptStatus {
1669            status: update,
1670            wants_to_retire,
1671            ..
1672        } = update;
1673
1674        #[allow(clippy::single_match)] // want to be explicit about the Ok type
1675        match wants_to_retire {
1676            Err(IptWantsToRetire) => ipt.is_current = None,
1677            Ok(()) => {}
1678        }
1679
1680        let now = || imm.runtime.now();
1681
1682        let started = match &ipt.status_last {
1683            TS::Establishing { started, .. } => Ok(*started),
1684            TS::Faulty { started, .. } => *started,
1685            TS::Good { .. } => Err(()),
1686        };
1687
1688        ipt.status_last = match update {
1689            ISS::Establishing => TS::Establishing {
1690                started: started.unwrap_or_else(|()| now()),
1691            },
1692            ISS::Good(details) => {
1693                let time_to_establish = started.and_then(|started| {
1694                    // return () at end of ok_or_else closure, for clarity
1695                    #[allow(clippy::unused_unit, clippy::semicolon_if_nothing_returned)]
1696                    now().checked_duration_since(started).ok_or_else(|| {
1697                        warn!("monotonic clock went backwards! (HS IPT)");
1698                        ()
1699                    })
1700                });
1701                TS::Good {
1702                    time_to_establish,
1703                    details,
1704                }
1705            }
1706            ISS::Faulty(error) => TS::Faulty { started, error },
1707        };
1708    }
1709}
1710
1711//========== mockability ==========
1712
1713/// Mockable state for the IPT Manager
1714///
1715/// This allows us to use a fake IPT Establisher and IPT Publisher,
1716/// so that we can unit test the Manager.
1717pub(crate) trait Mockable<R>: Debug + Send + Sync + Sized + 'static {
1718    /// IPT establisher type
1719    type IptEstablisher: Send + Sync + 'static;
1720
1721    /// A random number generator
1722    type Rng<'m>: rand::Rng + rand::CryptoRng + 'm;
1723
1724    /// Return a random number generator
1725    fn thread_rng(&mut self) -> Self::Rng<'_>;
1726
1727    /// Call `IptEstablisher::new`
1728    fn make_new_ipt(
1729        &mut self,
1730        imm: &Immutable<R>,
1731        params: IptParameters,
1732    ) -> Result<(Self::IptEstablisher, watch::Receiver<IptStatus>), FatalError>;
1733
1734    /// Call `IptEstablisher::start_accepting`
1735    fn start_accepting(&self, establisher: &ErasedIptEstablisher);
1736
1737    /// Allow tests to see when [`IptManager::expire_old_ipts_external_persistent_state`]
1738    /// is called.
1739    ///
1740    /// This lets tests see that it gets called at the right times,
1741    /// and not the wrong ones.
1742    fn expire_old_ipts_external_persistent_state_hook(&self);
1743}
1744
1745impl<R: Runtime> Mockable<R> for Real<R> {
1746    type IptEstablisher = IptEstablisher;
1747
1748    /// A random number generator
1749    type Rng<'m> = rand::rngs::ThreadRng;
1750
1751    /// Return a random number generator
1752    fn thread_rng(&mut self) -> Self::Rng<'_> {
1753        rand::rng()
1754    }
1755
1756    fn make_new_ipt(
1757        &mut self,
1758        imm: &Immutable<R>,
1759        params: IptParameters,
1760    ) -> Result<(Self::IptEstablisher, watch::Receiver<IptStatus>), FatalError> {
1761        IptEstablisher::launch(&imm.runtime, params, self.circ_pool.clone(), &imm.keymgr)
1762    }
1763
1764    fn start_accepting(&self, establisher: &ErasedIptEstablisher) {
1765        let establisher: &IptEstablisher = <dyn Any>::downcast_ref(establisher)
1766            .expect("upcast failure, ErasedIptEstablisher is not IptEstablisher!");
1767        establisher.start_accepting();
1768    }
1769
1770    fn expire_old_ipts_external_persistent_state_hook(&self) {}
1771}
1772
1773// TODO #1213 add more unit tests for IptManager
1774// Especially, we want to exercise all code paths in idempotently_progress_things_now
1775
1776#[cfg(test)]
1777mod test {
1778    // @@ begin test lint list maintained by maint/add_warning @@
1779    #![allow(clippy::bool_assert_comparison)]
1780    #![allow(clippy::clone_on_copy)]
1781    #![allow(clippy::dbg_macro)]
1782    #![allow(clippy::mixed_attributes_style)]
1783    #![allow(clippy::print_stderr)]
1784    #![allow(clippy::print_stdout)]
1785    #![allow(clippy::single_char_pattern)]
1786    #![allow(clippy::unwrap_used)]
1787    #![allow(clippy::unchecked_time_subtraction)]
1788    #![allow(clippy::useless_vec)]
1789    #![allow(clippy::needless_pass_by_value)]
1790    #![allow(clippy::string_slice)] // See arti#2571
1791    //! <!-- @@ end test lint list maintained by maint/add_warning @@ -->
1792    #![allow(clippy::match_single_binding)] // false positives, need the lifetime extension
1793    use super::*;
1794
1795    use crate::config::OnionServiceConfigBuilder;
1796    use crate::ipt_establish::GoodIptDetails;
1797    use crate::status::{OnionServiceStatus, StatusSender};
1798    use crate::test::{create_keymgr, create_storage_handles_from_state_dir};
1799    use rand::SeedableRng as _;
1800    use slotmap_careful::DenseSlotMap;
1801    use std::collections::BTreeMap;
1802    use std::sync::Mutex;
1803    use test_temp_dir::{TestTempDir, test_temp_dir};
1804    use tor_basic_utils::test_rng::TestingRng;
1805    use tor_netdir::testprovider::TestNetDirProvider;
1806    use tor_rtmock::MockRuntime;
1807    use tracing_test::traced_test;
1808    use walkdir::WalkDir;
1809
1810    slotmap_careful::new_key_type! {
1811        struct MockEstabId;
1812    }
1813
1814    type MockEstabs = Arc<Mutex<DenseSlotMap<MockEstabId, MockEstabState>>>;
1815
1816    fn ms(ms: u64) -> Duration {
1817        Duration::from_millis(ms)
1818    }
1819
1820    #[derive(Debug)]
1821    struct Mocks {
1822        rng: TestingRng,
1823        estabs: MockEstabs,
1824        expect_expire_ipts_calls: Arc<Mutex<usize>>,
1825    }
1826
1827    #[derive(Debug)]
1828    struct MockEstabState {
1829        st_tx: watch::Sender<IptStatus>,
1830        params: IptParameters,
1831    }
1832
1833    #[derive(Debug)]
1834    struct MockEstab {
1835        esid: MockEstabId,
1836        estabs: MockEstabs,
1837    }
1838
1839    impl Mockable<MockRuntime> for Mocks {
1840        type IptEstablisher = MockEstab;
1841        type Rng<'m> = &'m mut TestingRng;
1842
1843        fn thread_rng(&mut self) -> Self::Rng<'_> {
1844            &mut self.rng
1845        }
1846
1847        fn make_new_ipt(
1848            &mut self,
1849            _imm: &Immutable<MockRuntime>,
1850            params: IptParameters,
1851        ) -> Result<(Self::IptEstablisher, watch::Receiver<IptStatus>), FatalError> {
1852            let (st_tx, st_rx) = watch::channel();
1853            let estab = MockEstabState { st_tx, params };
1854            let esid = self.estabs.lock().unwrap().insert(estab);
1855            let estab = MockEstab {
1856                esid,
1857                estabs: self.estabs.clone(),
1858            };
1859            Ok((estab, st_rx))
1860        }
1861
1862        fn start_accepting(&self, _establisher: &ErasedIptEstablisher) {}
1863
1864        fn expire_old_ipts_external_persistent_state_hook(&self) {
1865            let mut expect = self.expect_expire_ipts_calls.lock().unwrap();
1866            eprintln!("expire_old_ipts_external_persistent_state_hook, expect={expect}");
1867            *expect = expect.checked_sub(1).expect("unexpected expiry");
1868        }
1869    }
1870
1871    impl Drop for MockEstab {
1872        fn drop(&mut self) {
1873            let mut estabs = self.estabs.lock().unwrap();
1874            let _: MockEstabState = estabs
1875                .remove(self.esid)
1876                .expect("dropping non-recorded MockEstab");
1877        }
1878    }
1879
1880    struct MockedIptManager<'d> {
1881        estabs: MockEstabs,
1882        pub_view: ipt_set::IptsPublisherView,
1883        shut_tx: broadcast::Sender<Void>,
1884        #[allow(dead_code)]
1885        cfg_tx: watch::Sender<Arc<OnionServiceConfig>>,
1886        #[allow(dead_code)] // ensures temp dir lifetime; paths stored in self
1887        temp_dir: &'d TestTempDir,
1888        expect_expire_ipts_calls: Arc<Mutex<usize>>, // use usize::MAX to not mind
1889    }
1890
1891    impl<'d> MockedIptManager<'d> {
1892        fn startup(
1893            runtime: MockRuntime,
1894            temp_dir: &'d TestTempDir,
1895            seed: u64,
1896            expect_expire_ipts_calls: usize,
1897        ) -> Self {
1898            let dir: TestNetDirProvider = tor_netdir::testnet::construct_netdir()
1899                .unwrap_if_sufficient()
1900                .unwrap()
1901                .into();
1902
1903            let nick: HsNickname = "nick".to_string().try_into().unwrap();
1904
1905            let cfg = OnionServiceConfigBuilder::default()
1906                .nickname(nick.clone())
1907                .build()
1908                .unwrap();
1909
1910            let (cfg_tx, cfg_rx) = watch::channel_with(Arc::new(cfg));
1911
1912            let (rend_tx, _rend_rx) = mpsc::channel(10);
1913            let (shut_tx, shut_rx) = broadcast::channel::<Void>(0);
1914
1915            let estabs: MockEstabs = Default::default();
1916            let expect_expire_ipts_calls = Arc::new(Mutex::new(expect_expire_ipts_calls));
1917
1918            let mocks = Mocks {
1919                rng: TestingRng::seed_from_u64(seed),
1920                estabs: estabs.clone(),
1921                expect_expire_ipts_calls: expect_expire_ipts_calls.clone(),
1922            };
1923
1924            // Don't provide a subdir; the ipt_mgr is supposed to add any needed subdirs
1925            let state_dir = temp_dir
1926                // untracked is OK because our return value captures 'd
1927                .subdir_untracked("state_dir");
1928
1929            let (state_handle, iptpub_state_handle) =
1930                create_storage_handles_from_state_dir(&state_dir, &nick);
1931
1932            let (mgr_view, pub_view) =
1933                ipt_set::ipts_channel(&runtime, iptpub_state_handle).unwrap();
1934
1935            let keymgr = create_keymgr(temp_dir);
1936            let keymgr = keymgr.into_untracked(); // OK because our return value captures 'd
1937            let status_tx = StatusSender::new(OnionServiceStatus::new_shutdown()).into();
1938            let mgr = IptManager::new(
1939                runtime.clone(),
1940                Arc::new(dir),
1941                nick,
1942                cfg_rx,
1943                rend_tx,
1944                shut_rx,
1945                &state_handle,
1946                mocks,
1947                keymgr,
1948                status_tx,
1949            )
1950            .unwrap();
1951
1952            mgr.launch_background_tasks(mgr_view).unwrap();
1953
1954            MockedIptManager {
1955                estabs,
1956                pub_view,
1957                shut_tx,
1958                cfg_tx,
1959                temp_dir,
1960                expect_expire_ipts_calls,
1961            }
1962        }
1963
1964        async fn shutdown_check_no_tasks(self, runtime: &MockRuntime) {
1965            drop(self.shut_tx);
1966            runtime.progress_until_stalled().await;
1967            assert_eq!(runtime.mock_task().n_tasks(), 1); // just us
1968        }
1969
1970        fn estabs_inventory(&self) -> impl Eq + Debug + 'static + use<> {
1971            let estabs = self.estabs.lock().unwrap();
1972            estabs
1973                .values()
1974                .map(|MockEstabState { params: p, .. }| {
1975                    (
1976                        p.lid,
1977                        (
1978                            p.target.clone(),
1979                            // We want to check the key values, but they're very hard to get at
1980                            // in a way we can compare.  Especially the private keys, for which
1981                            // we can't getting a clone or copy of the private key material out of the Arc.
1982                            // They're keypairs, we can use the debug rep which shows the public half.
1983                            // That will have to do.
1984                            format!("{:?}", p.k_sid),
1985                            format!("{:?}", p.k_ntor),
1986                        ),
1987                    )
1988                })
1989                .collect::<BTreeMap<_, _>>()
1990        }
1991    }
1992
1993    #[test]
1994    #[traced_test]
1995    fn test_mgr_lifecycle() {
1996        MockRuntime::test_with_various(|runtime| async move {
1997            let temp_dir = test_temp_dir!();
1998
1999            let m = MockedIptManager::startup(runtime.clone(), &temp_dir, 0, 1);
2000            runtime.progress_until_stalled().await;
2001
2002            assert_eq!(*m.expect_expire_ipts_calls.lock().unwrap(), 0);
2003
2004            // We expect it to try to establish 3 IPTs
2005            const EXPECT_N_IPTS: usize = 3;
2006            const EXPECT_MAX_IPTS: usize = EXPECT_N_IPTS + 2 /* num_extra */;
2007            assert_eq!(m.estabs.lock().unwrap().len(), EXPECT_N_IPTS);
2008            assert!(m.pub_view.borrow_for_publish().ipts.is_none());
2009
2010            // Advancing time a bit and it still shouldn't publish anything
2011            runtime.advance_by(ms(500)).await;
2012            runtime.progress_until_stalled().await;
2013            assert!(m.pub_view.borrow_for_publish().ipts.is_none());
2014
2015            let good = GoodIptDetails {
2016                link_specifiers: vec![],
2017                ipt_kp_ntor: [0x55; 32].into(),
2018            };
2019
2020            // Imagine that one of our IPTs becomes good
2021            m.estabs
2022                .lock()
2023                .unwrap()
2024                .values_mut()
2025                .next()
2026                .unwrap()
2027                .st_tx
2028                .borrow_mut()
2029                .status = IptStatusStatus::Good(good.clone());
2030
2031            // TODO #1213 test that we haven't called start_accepting
2032
2033            // It won't publish until a further fastest establish time
2034            // Ie, until a further 500ms = 1000ms
2035            runtime.progress_until_stalled().await;
2036            assert!(m.pub_view.borrow_for_publish().ipts.is_none());
2037            runtime.advance_by(ms(499)).await;
2038            assert!(m.pub_view.borrow_for_publish().ipts.is_none());
2039            runtime.advance_by(ms(1)).await;
2040            match m.pub_view.borrow_for_publish().ipts.as_mut().unwrap() {
2041                pub_view => {
2042                    assert_eq!(pub_view.ipts.len(), 1);
2043                    assert_eq!(pub_view.lifetime, IPT_PUBLISH_UNCERTAIN);
2044                }
2045            };
2046
2047            // TODO #1213 test that we have called start_accepting on the right IPTs
2048
2049            // Set the other IPTs to be Good too
2050            for e in m.estabs.lock().unwrap().values_mut().skip(1) {
2051                e.st_tx.borrow_mut().status = IptStatusStatus::Good(good.clone());
2052            }
2053            runtime.progress_until_stalled().await;
2054            match m.pub_view.borrow_for_publish().ipts.as_mut().unwrap() {
2055                pub_view => {
2056                    assert_eq!(pub_view.ipts.len(), EXPECT_N_IPTS);
2057                    assert_eq!(pub_view.lifetime, IPT_PUBLISH_CERTAIN);
2058                }
2059            };
2060
2061            // TODO #1213 test that we have called start_accepting on the right IPTs
2062
2063            let estabs_inventory = m.estabs_inventory();
2064
2065            // Shut down
2066            m.shutdown_check_no_tasks(&runtime).await;
2067
2068            // ---------- restart! ----------
2069            info!("*** Restarting ***");
2070
2071            let m = MockedIptManager::startup(runtime.clone(), &temp_dir, 1, 1);
2072            runtime.progress_until_stalled().await;
2073            assert_eq!(*m.expect_expire_ipts_calls.lock().unwrap(), 0);
2074
2075            assert_eq!(estabs_inventory, m.estabs_inventory());
2076
2077            // TODO #1213 test that we have called start_accepting on all the old IPTs
2078
2079            // ---------- New IPT relay selection ----------
2080
2081            let old_lids: Vec<String> = m
2082                .estabs
2083                .lock()
2084                .unwrap()
2085                .values()
2086                .map(|ess| ess.params.lid.to_string())
2087                .collect();
2088            eprintln!("IPTs to rotate out: {old_lids:?}");
2089
2090            let old_lid_files = || {
2091                WalkDir::new(temp_dir.as_path_untracked())
2092                    .into_iter()
2093                    .map(|ent| {
2094                        ent.unwrap()
2095                            .into_path()
2096                            .into_os_string()
2097                            .into_string()
2098                            .unwrap()
2099                    })
2100                    .filter(|path| old_lids.iter().any(|lid| path.contains(lid)))
2101                    .collect_vec()
2102            };
2103
2104            let no_files: [String; 0] = [];
2105
2106            assert_ne!(old_lid_files(), no_files);
2107
2108            // It might call the expiry function once, or once per IPT.
2109            // The latter is quadratic but this is quite rare, so that's fine.
2110            *m.expect_expire_ipts_calls.lock().unwrap() = EXPECT_MAX_IPTS;
2111
2112            // wait 2 days, > hs_intro_max_lifetime
2113            runtime.advance_by(ms(48 * 60 * 60 * 1_000)).await;
2114            runtime.progress_until_stalled().await;
2115
2116            // It must have called it at least once.
2117            assert_ne!(*m.expect_expire_ipts_calls.lock().unwrap(), EXPECT_MAX_IPTS);
2118
2119            // There should now be no files names after old IptLocalIds.
2120            assert_eq!(old_lid_files(), no_files);
2121
2122            // Shut down
2123            m.shutdown_check_no_tasks(&runtime).await;
2124        });
2125    }
2126}