tor_hsservice/publish/reactor.rs
1//! The onion service publisher reactor.
2//!
3//! Generates and publishes hidden service descriptors in response to various events.
4//!
5//! [`Reactor::run`] is the entry-point of the reactor. It starts the reactor,
6//! and runs until [`Reactor::run_once`] returns [`ShutdownStatus::Terminate`]
7//! or a fatal error occurs. `ShutdownStatus::Terminate` is returned if
8//! any of the channels the reactor is receiving events from is closed
9//! (i.e. when the senders are dropped).
10//!
11//! ## Publisher status
12//!
13//! The publisher has an internal [`PublishStatus`], distinct from its [`State`],
14//! which is used for onion service status reporting.
15//!
16//! The main loop of the reactor reads the current `PublishStatus` from `publish_status_rx`,
17//! and responds by generating and publishing a new descriptor if needed.
18//!
19//! See [`PublishStatus`] and [`Reactor::publish_status_rx`] for more details.
20//!
21//! ## When do we publish?
22//!
23//! We generate and publish a new descriptor if
24//! * the introduction points have changed
25//! * the onion service configuration has changed in a meaningful way (for example,
26//! if the `restricted_discovery` configuration or its [`Anonymity`](crate::Anonymity)
27//! has changed. See [`OnionServiceConfigPublisherView`]).
28//! * there is a new consensus
29//! * it is time to republish the descriptor (after we upload a descriptor,
30//! we schedule it for republishing at a random time between 60 minutes and 120 minutes
31//! in the future)
32//!
33//! ## Onion service status
34//!
35//! With respect to [`OnionServiceStatus`] reporting,
36//! the following state transitions are possible:
37//!
38//!
39//! ```ignore
40//!
41//! update_publish_status(UploadScheduled|AwaitingIpts|RateLimited)
42//! +---------------------------------------+
43//! | |
44//! | v
45//! | +---------------+
46//! | | Bootstrapping |
47//! | +---------------+
48//! | |
49//! | | uploaded to at least
50//! | not enough HsDir uploads succeeded | some HsDirs from each ring
51//! | +-----------------------------+-----------------------+
52//! | | | |
53//! | | all HsDir uploads succeeded |
54//! | | | |
55//! | v v v
56//! | +---------------------+ +---------+ +---------------------+
57//! | | DegradedUnreachable | | Running | | DegradedReachable |
58//! +----------+ | +---------------------+ +---------+ +---------------------+
59//! | Shutdown |-- | | | |
60//! +----------+ | | | |
61//! | | | |
62//! | | | |
63//! | +---------------------------+------------------------+
64//! | | invalid authorized_clients
65//! | | after handling config change
66//! | |
67//! | v
68//! | run_once() returns an error +--------+
69//! +-------------------------------->| Broken |
70//! +--------+
71//! ```
72//!
73//! We can also transition from `Broken`, `DegradedReachable`, or `DegradedUnreachable`
74//! back to `Bootstrapping` (those transitions were omitted for brevity).
75
76use tor_circmgr::ServiceOnionServiceDirTunnel;
77use tor_config::file_watcher::{
78 self, Event as FileEvent, FileEventReceiver, FileEventSender, FileWatcher, FileWatcherBuilder,
79};
80use tor_config_path::{CfgPath, CfgPathResolver};
81use tor_dirclient::SourceInfo;
82use tor_netdir::{DirEvent, NetDir};
83use tracing::instrument;
84
85use crate::config::OnionServiceConfigPublisherView;
86use crate::config::restricted_discovery::{
87 DirectoryKeyProviderList, RestrictedDiscoveryConfig, RestrictedDiscoveryKeys,
88};
89use crate::status::{DescUploadRetryError, Problem};
90
91use super::*;
92use derive_more::From;
93
94// TODO-CLIENT-AUTH: perhaps we should add a separate CONFIG_CHANGE_REPUBLISH_DEBOUNCE_INTERVAL
95// for rate-limiting the publish jobs triggered by a change in the config?
96//
97// Currently the descriptor publish tasks triggered by changes in the config
98// are rate-limited via the usual rate limiting mechanism
99// (which rate-limits the uploads for 1m).
100//
101// I think this is OK for now, but we might need to rethink this if it becomes problematic
102// (for example, we might want an even longer rate-limit, or to reset any existing rate-limits
103// each time the config is modified).
104
105/// The upload rate-limiting threshold.
106///
107/// Before initiating an upload, the reactor checks if the last upload was at least
108/// `UPLOAD_RATE_LIM_THRESHOLD` seconds ago. If so, it uploads the descriptor to all HsDirs that
109/// need it. If not, it schedules the upload to happen `UPLOAD_RATE_LIM_THRESHOLD` seconds from the
110/// current time.
111//
112// TODO: We may someday need to tune this value; it was chosen more or less arbitrarily.
113const UPLOAD_RATE_LIM_THRESHOLD: Duration = Duration::from_secs(60);
114
115/// The maximum number of concurrent upload tasks per time period.
116//
117// TODO: this value was arbitrarily chosen and may not be optimal. For now, it
118// will have no effect, since the current number of replicas is far less than
119// this value.
120//
121// The uploads for all TPs happen in parallel. As a result, the actual limit for the maximum
122// number of concurrent upload tasks is multiplied by a number which depends on the TP parameters
123// (currently 2, which means the concurrency limit will, in fact, be 32).
124//
125// We should try to decouple this value from the TP parameters.
126const MAX_CONCURRENT_UPLOADS: usize = 16;
127
128/// The maximum time allowed for uploading a descriptor to a single HSDir,
129/// across all attempts.
130pub(crate) const OVERALL_UPLOAD_TIMEOUT: Duration = Duration::from_secs(5 * 60);
131
132/// A reactor for the HsDir [`Publisher`]
133///
134/// The entrypoint is [`Reactor::run`].
135#[must_use = "If you don't call run() on the reactor, it won't publish any descriptors."]
136pub(super) struct Reactor<R: Runtime, M: Mockable> {
137 /// The immutable, shared inner state.
138 imm: Arc<Immutable<R, M>>,
139 /// A source for new network directories that we use to determine
140 /// our HsDirs.
141 dir_provider: Arc<dyn NetDirProvider>,
142 /// The mutable inner state,
143 inner: Arc<Mutex<Inner>>,
144 /// A channel for receiving IPT change notifications.
145 ipt_watcher: IptsPublisherView,
146 /// A channel for receiving onion service config change notifications.
147 config_rx: watch::Receiver<Arc<OnionServiceConfig>>,
148 /// A channel for receiving restricted discovery key_dirs change notifications.
149 key_dirs_rx: FileEventReceiver,
150 /// A channel for sending restricted discovery key_dirs change notifications.
151 ///
152 /// A copy of this sender is handed out to every `FileWatcher` created.
153 key_dirs_tx: FileEventSender,
154 /// A channel for receiving updates regarding our [`PublishStatus`].
155 ///
156 /// The main loop of the reactor watches for updates on this channel.
157 ///
158 /// When the [`PublishStatus`] changes to [`UploadScheduled`](PublishStatus::UploadScheduled),
159 /// we can start publishing descriptors.
160 ///
161 /// If the [`PublishStatus`] is [`AwaitingIpts`](PublishStatus::AwaitingIpts), publishing is
162 /// paused until we receive a notification on `ipt_watcher` telling us the IPT manager has
163 /// established some introduction points.
164 publish_status_rx: watch::Receiver<PublishStatus>,
165 /// A sender for updating our [`PublishStatus`].
166 ///
167 /// When our [`PublishStatus`] changes to [`UploadScheduled`](PublishStatus::UploadScheduled),
168 /// we can start publishing descriptors.
169 publish_status_tx: watch::Sender<PublishStatus>,
170 /// A channel for sending upload completion notifications.
171 ///
172 /// This channel is polled in the main loop of the reactor.
173 upload_task_complete_rx: mpsc::Receiver<TimePeriodUploadResult>,
174 /// A channel for receiving upload completion notifications.
175 ///
176 /// A copy of this sender is handed to each upload task.
177 upload_task_complete_tx: mpsc::Sender<TimePeriodUploadResult>,
178 /// A sender for notifying any pending upload tasks that the reactor is shutting down.
179 ///
180 /// Receivers can use this channel to find out when reactor is dropped.
181 ///
182 /// This is currently only used in [`upload_for_time_period`](Reactor::upload_for_time_period).
183 /// Any future background tasks can also use this channel to detect if the reactor is dropped.
184 ///
185 /// Closing this channel will cause any pending upload tasks to be dropped.
186 shutdown_tx: broadcast::Sender<Void>,
187 /// Path resolver for configuration files.
188 path_resolver: Arc<CfgPathResolver>,
189 /// Queue on which we receive messages from the [`PowManager`] telling us that a seed has
190 /// rotated and thus we need to republish the descriptor for a particular time period.
191 update_from_pow_manager_rx: mpsc::Receiver<TimePeriod>,
192}
193
194/// The immutable, shared state of the descriptor publisher reactor.
195#[derive(Clone)]
196struct Immutable<R: Runtime, M: Mockable> {
197 /// The runtime.
198 runtime: R,
199 /// Mockable state.
200 ///
201 /// This is used for launching circuits and for obtaining random number generators.
202 mockable: M,
203 /// The service for which we're publishing descriptors.
204 nickname: HsNickname,
205 /// The key manager,
206 keymgr: Arc<KeyMgr>,
207 /// A sender for updating the status of the onion service.
208 status_tx: PublisherStatusSender,
209 /// Proof-of-work state.
210 pow_manager: Arc<PowManager<R>>,
211}
212
213impl<R: Runtime, M: Mockable> Immutable<R, M> {
214 /// Create an [`AesOpeKey`] for generating revision counters for the descriptors associated
215 /// with the specified [`TimePeriod`].
216 ///
217 /// If the onion service is not running in offline mode, the key of the returned `AesOpeKey` is
218 /// the private part of the blinded identity key. Otherwise, the key is the private part of the
219 /// descriptor signing key.
220 ///
221 /// Returns an error if the service is running in offline mode and the descriptor signing
222 /// keypair of the specified `period` is not available.
223 //
224 // TODO (#1194): we don't support "offline" mode (yet), so this always returns an AesOpeKey
225 // built from the blinded id key
226 fn create_ope_key(&self, period: TimePeriod) -> Result<AesOpeKey, FatalError> {
227 let ope_key = match read_blind_id_keypair(&self.keymgr, &self.nickname, period)? {
228 Some(key) => {
229 let key: ed25519::ExpandedKeypair = key.into();
230 key.to_secret_key_bytes()[0..32]
231 .try_into()
232 .expect("Wrong length on slice")
233 }
234 None => {
235 // TODO (#1194): we don't support externally provisioned keys (yet), so this branch
236 // is unreachable (for now).
237 let desc_sign_key_spec =
238 DescSigningKeypairSpecifier::new(self.nickname.clone(), period);
239 let key: ed25519::Keypair = self
240 .keymgr
241 .get::<HsDescSigningKeypair>(&desc_sign_key_spec)?
242 // TODO (#1194): internal! is not the right type for this error (we need an
243 // error type for the case where a hidden service running in offline mode has
244 // run out of its pre-previsioned keys).
245 //
246 // This will be addressed when we add support for offline hs_id mode
247 .ok_or_else(|| {
248 internal!(
249 "identity keys are offline, but descriptor signing key is unavailable?!"
250 )
251 })?
252 .into();
253 key.to_bytes()
254 }
255 };
256
257 Ok(AesOpeKey::from_secret(&ope_key))
258 }
259
260 /// Generate a revision counter for a descriptor associated with the specified
261 /// [`TimePeriod`].
262 ///
263 /// Returns a revision counter generated according to the [encrypted time in period] scheme.
264 ///
265 /// [encrypted time in period]: https://spec.torproject.org/rend-spec/revision-counter-mgt.html#encrypted-time
266 fn generate_revision_counter(
267 &self,
268 params: &HsDirParams,
269 now: SystemTime,
270 ) -> Result<RevisionCounter, FatalError> {
271 // TODO: in the future, we might want to compute ope_key once per time period (as oppposed
272 // to each time we generate a new descriptor), for performance reasons.
273 let ope_key = self.create_ope_key(params.time_period())?;
274
275 // TODO: perhaps this should be moved to a new HsDirParams::offset_within_sr() function
276 let srv_start = params.start_of_shard_rand_period();
277 let offset = params.offset_within_srv_period(now).ok_or_else(|| {
278 internal!(
279 "current wallclock time not within SRV range?! (now={:?}, SRV_start={:?})",
280 now,
281 srv_start
282 )
283 })?;
284 let rev = ope_key.encrypt(offset);
285
286 Ok(RevisionCounter::from(rev))
287 }
288}
289
290/// Mockable state for the descriptor publisher reactor.
291///
292/// This enables us to mock parts of the [`Reactor`] for testing purposes.
293#[async_trait]
294pub(crate) trait Mockable: Clone + Send + Sync + Sized + 'static {
295 /// The type of random number generator.
296 type Rng: rand::Rng + rand::CryptoRng;
297
298 /// The type of client circuit.
299 type Tunnel: MockableDirTunnel;
300
301 /// Return a random number generator.
302 fn thread_rng(&self) -> Self::Rng;
303
304 /// Create a circuit of the specified `kind` to `target`.
305 async fn get_or_launch_hs_dir<T>(
306 &self,
307 netdir: &NetDir,
308 target: T,
309 ) -> Result<Self::Tunnel, tor_circmgr::Error>
310 where
311 T: CircTarget + Send + Sync;
312
313 /// Return an estimate-based value for how long we should allow a single
314 /// directory upload operation to complete.
315 ///
316 /// Includes circuit construction, stream opening, upload, and waiting for a
317 /// response.
318 fn estimate_upload_timeout(&self) -> Duration;
319}
320
321/// Mockable client circuit
322#[async_trait]
323pub(crate) trait MockableDirTunnel: Send + Sync {
324 /// The data stream type.
325 type DataStream: AsyncRead + AsyncWrite + Send + Unpin;
326
327 /// Start a new stream to the last relay in the circuit, using
328 /// a BEGIN_DIR cell.
329 async fn begin_dir_stream(&self) -> Result<Self::DataStream, tor_circmgr::Error>;
330
331 /// Try to get a SourceInfo for this circuit, for using it in a directory request.
332 fn source_info(&self) -> tor_proto::Result<Option<SourceInfo>>;
333}
334
335#[async_trait]
336impl MockableDirTunnel for ServiceOnionServiceDirTunnel {
337 type DataStream = tor_proto::client::stream::DataStream;
338
339 async fn begin_dir_stream(&self) -> Result<Self::DataStream, tor_circmgr::Error> {
340 Self::begin_dir_stream(self).await
341 }
342
343 fn source_info(&self) -> tor_proto::Result<Option<SourceInfo>> {
344 SourceInfo::from_tunnel(self)
345 }
346}
347
348/// The real version of the mockable state of the reactor.
349#[derive(Clone, From, Into)]
350pub(crate) struct Real<R: Runtime>(Arc<HsCircPool<R>>);
351
352#[async_trait]
353impl<R: Runtime> Mockable for Real<R> {
354 type Rng = rand::rngs::ThreadRng;
355 type Tunnel = ServiceOnionServiceDirTunnel;
356
357 fn thread_rng(&self) -> Self::Rng {
358 rand::rng()
359 }
360
361 #[instrument(level = "trace", skip_all)]
362 async fn get_or_launch_hs_dir<T>(
363 &self,
364 netdir: &NetDir,
365 target: T,
366 ) -> Result<Self::Tunnel, tor_circmgr::Error>
367 where
368 T: CircTarget + Send + Sync,
369 {
370 self.0.get_or_launch_svc_dir(netdir, target).await
371 }
372
373 fn estimate_upload_timeout(&self) -> Duration {
374 use tor_circmgr::timeouts::Action;
375 let est_build = self.0.estimate_timeout(&Action::BuildCircuit { length: 4 });
376 let est_roundtrip = self.0.estimate_timeout(&Action::RoundTrip { length: 4 });
377 // We assume that in the worst case we'll have to wait for an entire
378 // circuit construction and two round-trips to the hsdir.
379 let est_total = est_build + est_roundtrip * 2;
380 // We always allow _at least_ this much time, in case our estimate is
381 // ridiculously low.
382 let min_timeout = Duration::from_secs(30);
383 max(est_total, min_timeout)
384 }
385}
386
387/// The mutable state of a [`Reactor`].
388struct Inner {
389 /// The onion service config.
390 config: Arc<OnionServiceConfigPublisherView>,
391 /// Watcher for key_dirs.
392 ///
393 /// Set to `None` if the reactor is not running, or if `watch_configuration` is false.
394 ///
395 /// The watcher is recreated whenever the `restricted_discovery.key_dirs` change.
396 file_watcher: Option<FileWatcher>,
397 /// The relevant time periods.
398 ///
399 /// This includes the current time period, as well as any other time periods we need to be
400 /// publishing descriptors for.
401 ///
402 /// This is empty until we fetch our first netdir in [`Reactor::run`].
403 time_periods: Vec<TimePeriodContext>,
404 /// Our most up to date netdir.
405 ///
406 /// This is initialized in [`Reactor::run`].
407 netdir: Option<Arc<NetDir>>,
408 /// The timestamp of our last upload.
409 ///
410 /// This is the time when the last update was _initiated_ (rather than completed), to prevent
411 /// the publisher from spawning multiple upload tasks at once in response to multiple external
412 /// events happening in quick succession, such as the IPT manager sending multiple IPT change
413 /// notifications in a short time frame (#1142), or an IPT change notification that's
414 /// immediately followed by a consensus change. Starting two upload tasks at once is not only
415 /// inefficient, but it also causes the publisher to generate two different descriptors with
416 /// the same revision counter (the revision counter is derived from the current timestamp),
417 /// which ultimately causes the slower upload task to fail (see #1142).
418 ///
419 /// Note: This is only used for deciding when to reschedule a rate-limited upload. It is _not_
420 /// used for retrying failed uploads (these are handled internally by
421 /// [`Reactor::upload_descriptor_with_retries`]).
422 last_uploaded: Option<Instant>,
423 /// A max-heap containing the time periods for which we need to reupload the descriptor.
424 // TODO: we are currently reuploading more than nececessary.
425 // Ideally, this shouldn't contain contain duplicate TimePeriods,
426 // because we only need to retain the latest reupload time for each time period.
427 //
428 // Currently, if, for some reason, we upload the descriptor multiple times for the same TP,
429 // we will end up with multiple ReuploadTimer entries for that TP,
430 // each of which will (eventually) result in a reupload.
431 //
432 // TODO: maybe this should just be a HashMap<TimePeriod, Instant>
433 //
434 // See https://gitlab.torproject.org/tpo/core/arti/-/merge_requests/1971#note_2994950
435 reupload_timers: BinaryHeap<ReuploadTimer>,
436 /// The restricted discovery authorized clients.
437 ///
438 /// `None`, unless the service is running in restricted discovery mode.
439 authorized_clients: Option<Arc<RestrictedDiscoveryKeys>>,
440}
441
442/// The part of the reactor state that changes with every time period.
443struct TimePeriodContext {
444 /// The HsDir params.
445 params: HsDirParams,
446 /// The HsDirs to use in this time period.
447 ///
448 // We keep a list of `RelayIds` because we can't store a `Relay<'_>` inside the reactor
449 // (the lifetime of a relay is tied to the lifetime of its corresponding `NetDir`. To
450 // store `Relay<'_>`s in the reactor, we'd need a way of atomically swapping out both the
451 // `NetDir` and the cached relays, and to convince Rust what we're doing is sound)
452 hs_dirs: Vec<(RelayIds, DescriptorStatus)>,
453 /// The revision counter of the last successful upload, if any.
454 last_successful: Option<RevisionCounter>,
455 /// The outcome of the last upload, if any.
456 upload_results: Vec<HsDirUploadStatus>,
457}
458
459impl TimePeriodContext {
460 /// Create a new `TimePeriodContext`.
461 ///
462 /// Any of the specified `old_hsdirs` also present in the new list of HsDirs
463 /// (returned by `NetDir::hs_dirs_upload`) will have their `DescriptorStatus` preserved.
464 fn new<'r>(
465 params: HsDirParams,
466 blind_id: HsBlindId,
467 netdir: &Arc<NetDir>,
468 old_hsdirs: impl Iterator<Item = &'r (RelayIds, DescriptorStatus)>,
469 old_upload_results: Vec<HsDirUploadStatus>,
470 ) -> Result<Self, FatalError> {
471 let period = params.time_period();
472 let hs_dirs = Self::compute_hsdirs(period, blind_id, netdir, old_hsdirs)?;
473 let upload_results = old_upload_results
474 .into_iter()
475 .filter(|res|
476 // Check if the HsDir of this result still exists
477 hs_dirs
478 .iter()
479 .any(|(relay_ids, _status)| relay_ids == &res.relay_ids))
480 .collect();
481
482 Ok(Self {
483 params,
484 hs_dirs,
485 last_successful: None,
486 upload_results,
487 })
488 }
489
490 /// Recompute the HsDirs for this time period.
491 fn compute_hsdirs<'r>(
492 period: TimePeriod,
493 blind_id: HsBlindId,
494 netdir: &Arc<NetDir>,
495 mut old_hsdirs: impl Iterator<Item = &'r (RelayIds, DescriptorStatus)>,
496 ) -> Result<Vec<(RelayIds, DescriptorStatus)>, FatalError> {
497 let hs_dirs = netdir.hs_dirs_upload(blind_id, period)?;
498
499 Ok(hs_dirs
500 .map(|hs_dir| {
501 let mut builder = RelayIds::builder();
502 if let Some(ed_id) = hs_dir.ed_identity() {
503 builder.ed_identity(*ed_id);
504 }
505
506 if let Some(rsa_id) = hs_dir.rsa_identity() {
507 builder.rsa_identity(*rsa_id);
508 }
509
510 let relay_id = builder.build().unwrap_or_else(|_| RelayIds::empty());
511
512 // Have we uploaded the descriptor to thiw relay before? If so, we don't need to
513 // reupload it unless it was already dirty and due for a reupload.
514 let status = match old_hsdirs.find(|(id, _)| *id == relay_id) {
515 Some((_, status)) => *status,
516 None => DescriptorStatus::Dirty,
517 };
518
519 (relay_id, status)
520 })
521 .collect::<Vec<_>>())
522 }
523
524 /// Mark the descriptor dirty for all HSDirs of this time period.
525 fn mark_all_dirty(&mut self) {
526 self.hs_dirs
527 .iter_mut()
528 .for_each(|(_relay_id, status)| *status = DescriptorStatus::Dirty);
529 }
530
531 /// Update the upload result for this time period.
532 fn set_upload_results(&mut self, upload_results: Vec<HsDirUploadStatus>) {
533 self.upload_results = upload_results;
534 }
535}
536
537/// An error that occurs while trying to upload a descriptor.
538#[derive(Clone, Debug, thiserror::Error)]
539#[non_exhaustive]
540pub enum UploadError {
541 /// An error that has occurred after we have contacted a directory cache and made a circuit to it.
542 #[error("descriptor upload request failed: {}", _0.error)]
543 Request(#[from] RequestFailedError),
544
545 /// Failed to establish circuit to hidden service directory
546 #[error("could not build circuit to HsDir")]
547 Circuit(#[from] tor_circmgr::Error),
548
549 /// Failed to establish stream to hidden service directory
550 #[error("failed to establish directory stream to HsDir")]
551 Stream(#[source] tor_circmgr::Error),
552
553 /// An internal error.
554 #[error("Internal error")]
555 Bug(#[from] tor_error::Bug),
556}
557define_asref_dyn_std_error!(UploadError);
558
559impl UploadError {
560 /// Return true if this error is one that we should report as a suspicious event,
561 /// along with the dirserver, and description of the relevant document.
562 pub(crate) fn should_report_as_suspicious(&self) -> bool {
563 match self {
564 UploadError::Request(e) => e.error.should_report_as_suspicious_if_anon(),
565 UploadError::Circuit(_) => false, // TODO prop360
566 UploadError::Stream(_) => false, // TODO prop360
567 UploadError::Bug(_) => false,
568 }
569 }
570}
571
572impl<R: Runtime, M: Mockable> Reactor<R, M> {
573 /// Create a new `Reactor`.
574 #[allow(clippy::too_many_arguments)]
575 pub(super) fn new(
576 runtime: R,
577 nickname: HsNickname,
578 dir_provider: Arc<dyn NetDirProvider>,
579 mockable: M,
580 config: &OnionServiceConfig,
581 ipt_watcher: IptsPublisherView,
582 config_rx: watch::Receiver<Arc<OnionServiceConfig>>,
583 status_tx: PublisherStatusSender,
584 keymgr: Arc<KeyMgr>,
585 path_resolver: Arc<CfgPathResolver>,
586 pow_manager: Arc<PowManager<R>>,
587 update_from_pow_manager_rx: mpsc::Receiver<TimePeriod>,
588 ) -> Self {
589 /// The maximum size of the upload completion notifier channel.
590 ///
591 /// The channel we use this for is a futures::mpsc channel, which has a capacity of
592 /// `UPLOAD_CHAN_BUF_SIZE + num-senders`. We don't need the buffer size to be non-zero, as
593 /// each sender will send exactly one message.
594 const UPLOAD_CHAN_BUF_SIZE: usize = 0;
595
596 // Internally-generated instructions, no need for mq.
597 let (upload_task_complete_tx, upload_task_complete_rx) =
598 mpsc_channel_no_memquota(UPLOAD_CHAN_BUF_SIZE);
599
600 let (publish_status_tx, publish_status_rx) = watch::channel();
601 // Setting the buffer size to zero here is OK,
602 // since we never actually send anything on this channel.
603 let (shutdown_tx, _shutdown_rx) = broadcast::channel(0);
604
605 let authorized_clients =
606 Self::read_authorized_clients(&config.restricted_discovery, &path_resolver);
607
608 // Create a channel for watching for changes in the configured
609 // restricted_discovery.key_dirs.
610 let (key_dirs_tx, key_dirs_rx) = file_watcher::channel();
611
612 let imm = Immutable {
613 runtime,
614 mockable,
615 nickname,
616 keymgr,
617 status_tx,
618 pow_manager,
619 };
620
621 let inner = Inner {
622 time_periods: vec![],
623 config: Arc::new(config.into()),
624 file_watcher: None,
625 netdir: None,
626 last_uploaded: None,
627 reupload_timers: Default::default(),
628 authorized_clients,
629 };
630
631 Self {
632 imm: Arc::new(imm),
633 inner: Arc::new(Mutex::new(inner)),
634 dir_provider,
635 ipt_watcher,
636 config_rx,
637 key_dirs_rx,
638 key_dirs_tx,
639 publish_status_rx,
640 publish_status_tx,
641 upload_task_complete_rx,
642 upload_task_complete_tx,
643 shutdown_tx,
644 path_resolver,
645 update_from_pow_manager_rx,
646 }
647 }
648
649 /// Start the reactor.
650 ///
651 /// Under normal circumstances, this function runs indefinitely.
652 ///
653 /// Note: this also spawns the "reminder task" that we use to reschedule uploads whenever an
654 /// upload fails or is rate-limited.
655 pub(super) async fn run(mut self) -> Result<(), FatalError> {
656 debug!(nickname=%self.imm.nickname, "starting descriptor publisher reactor");
657
658 {
659 let netdir = self
660 .dir_provider
661 .wait_for_netdir(Timeliness::Timely)
662 .await?;
663 let time_periods = self.compute_time_periods(&netdir, &[])?;
664
665 let mut inner = self.inner.lock().expect("poisoned lock");
666
667 inner.netdir = Some(netdir);
668 inner.time_periods = time_periods;
669 }
670
671 // Create the initial key_dirs watcher.
672 self.update_file_watcher();
673
674 loop {
675 match self.run_once().await {
676 Ok(ShutdownStatus::Continue) => continue,
677 Ok(ShutdownStatus::Terminate) => {
678 debug!(nickname=%self.imm.nickname, "descriptor publisher is shutting down!");
679
680 self.imm.status_tx.send_shutdown();
681 return Ok(());
682 }
683 Err(e) => {
684 error_report!(
685 e,
686 "HS service {}: descriptor publisher crashed!",
687 self.imm.nickname
688 );
689
690 self.imm.status_tx.send_broken(e.clone());
691
692 return Err(e);
693 }
694 }
695 }
696 }
697
698 /// Run one iteration of the reactor loop.
699 async fn run_once(&mut self) -> Result<ShutdownStatus, FatalError> {
700 let mut netdir_events = self.dir_provider.events();
701
702 // Note: TrackingNow tracks the values it is compared with.
703 // This is equivalent to sleeping for (until - now) units of time,
704 let upload_rate_lim: TrackingNow = TrackingNow::now(&self.imm.runtime);
705 if let PublishStatus::RateLimited(until) = self.status() {
706 if upload_rate_lim > until {
707 // We are no longer rate-limited
708 self.expire_rate_limit().await?;
709 }
710 }
711
712 let reupload_tracking = TrackingNow::now(&self.imm.runtime);
713 let mut reupload_periods = vec![];
714 {
715 let mut inner = self.inner.lock().expect("poisoned lock");
716 let inner = &mut *inner;
717 while let Some(reupload) = inner.reupload_timers.peek().copied() {
718 // First, extract all the timeouts that already elapsed.
719 if reupload.when <= reupload_tracking {
720 inner.reupload_timers.pop();
721 reupload_periods.push(reupload.period);
722 } else {
723 // We are not ready to schedule any more reuploads.
724 //
725 // How much we need to sleep is implicitly
726 // tracked in reupload_tracking (through
727 // the TrackingNow implementation)
728 break;
729 }
730 }
731 }
732
733 // Check if it's time to schedule any reuploads.
734 for period in reupload_periods {
735 if self.mark_dirty(&period) {
736 debug!(
737 time_period=?period,
738 "descriptor reupload timer elapsed; scheduling reupload",
739 );
740 self.update_publish_status_unless_rate_lim(PublishStatus::UploadScheduled)
741 .await?;
742 }
743 }
744
745 select_biased! {
746 res = self.upload_task_complete_rx.next().fuse() => {
747 let Some(upload_res) = res else {
748 return Ok(ShutdownStatus::Terminate);
749 };
750
751 self.handle_upload_results(upload_res);
752 self.upload_result_to_svc_status()?;
753 },
754 () = upload_rate_lim.wait_for_earliest(&self.imm.runtime).fuse() => {
755 self.expire_rate_limit().await?;
756 },
757 () = reupload_tracking.wait_for_earliest(&self.imm.runtime).fuse() => {
758 // Run another iteration, executing run_once again. This time, we will remove the
759 // expired reupload from self.reupload_timers, mark the descriptor dirty for all
760 // relevant HsDirs, and schedule the upload by setting our status to
761 // UploadScheduled.
762 return Ok(ShutdownStatus::Continue);
763 },
764 netdir_event = netdir_events.next().fuse() => {
765 let Some(netdir_event) = netdir_event else {
766 debug!("netdir event stream ended");
767 return Ok(ShutdownStatus::Terminate);
768 };
769
770 if !matches!(netdir_event, DirEvent::NewConsensus) {
771 return Ok(ShutdownStatus::Continue);
772 };
773
774 // The consensus changed. Grab a new NetDir.
775 let netdir = match self.dir_provider.netdir(Timeliness::Timely) {
776 Ok(y) => y,
777 Err(e) => {
778 error_report!(e, "HS service {}: netdir unavailable. Retrying...", self.imm.nickname);
779 // Hopefully a netdir will appear in the future.
780 // in the meantime, suspend operations.
781 //
782 // TODO (#1218): there is a bug here: we stop reading on our inputs
783 // including eg publish_status_rx, but it is our job to log some of
784 // these things. While we are waiting for a netdir, all those messages
785 // are "stuck"; they'll appear later, with misleading timestamps.
786 //
787 // Probably this should be fixed by moving the logging
788 // out of the reactor, where it won't be blocked.
789 self.dir_provider.wait_for_netdir(Timeliness::Timely)
790 .await?
791 }
792 };
793 let relevant_periods = netdir.hs_all_time_periods();
794 self.handle_consensus_change(netdir).await?;
795 expire_publisher_keys(
796 &self.imm.keymgr,
797 &self.imm.nickname,
798 &relevant_periods,
799 ).unwrap_or_else(|e| {
800 error_report!(e, "failed to remove expired keys");
801 });
802 }
803 update = self.ipt_watcher.await_update().fuse() => {
804 if self.handle_ipt_change(update).await? == ShutdownStatus::Terminate {
805 return Ok(ShutdownStatus::Terminate);
806 }
807 },
808 config = self.config_rx.next().fuse() => {
809 let Some(config) = config else {
810 return Ok(ShutdownStatus::Terminate);
811 };
812
813 self.handle_svc_config_change(&config).await?;
814 },
815 res = self.key_dirs_rx.next().fuse() => {
816 let Some(event) = res else {
817 return Ok(ShutdownStatus::Terminate);
818 };
819
820 while let Some(_ignore) = self.key_dirs_rx.try_recv() {
821 // Discard other events, so that we only reload once.
822 }
823
824 self.handle_key_dirs_change(event).await?;
825 }
826 should_upload = self.publish_status_rx.next().fuse() => {
827 let Some(should_upload) = should_upload else {
828 return Ok(ShutdownStatus::Terminate);
829 };
830
831 // Our PublishStatus changed -- are we ready to publish?
832 if should_upload == PublishStatus::UploadScheduled {
833 self.update_publish_status_unless_waiting(PublishStatus::Idle).await?;
834 self.upload_all().await?;
835 }
836 }
837 update_tp_pow_seed = self.update_from_pow_manager_rx.next().fuse() => {
838 debug!("Update PoW seed for TP!");
839 let Some(time_period) = update_tp_pow_seed else {
840 return Ok(ShutdownStatus::Terminate);
841 };
842 self.mark_dirty(&time_period);
843 self.upload_all().await?;
844 }
845 }
846
847 Ok(ShutdownStatus::Continue)
848 }
849
850 /// Returns the current status of the publisher
851 fn status(&self) -> PublishStatus {
852 *self.publish_status_rx.borrow()
853 }
854
855 /// Handle a batch of upload outcomes,
856 /// possibly updating the status of the descriptor for the corresponding HSDirs.
857 fn handle_upload_results(&self, results: TimePeriodUploadResult) {
858 let mut inner = self.inner.lock().expect("poisoned lock");
859 let inner = &mut *inner;
860
861 // Check which time period these uploads pertain to.
862 let period = inner
863 .time_periods
864 .iter_mut()
865 .find(|ctx| ctx.params.time_period() == results.time_period);
866
867 let Some(period) = period else {
868 // The uploads were for a time period that is no longer relevant, so we
869 // can ignore the result.
870 return;
871 };
872
873 // We will need to reupload this descriptor at some point, so we pick
874 // a random time between 60 minutes and 120 minutes in the future.
875 //
876 // See https://spec.torproject.org/rend-spec/deriving-keys.html#WHEN-HSDESC
877 let mut rng = self.imm.mockable.thread_rng();
878 // TODO SPEC: Control republish period using a consensus parameter?
879 let minutes = rng.gen_range_checked(60..=120).expect("low > high?!");
880 let duration = Duration::from_secs(minutes * 60);
881 let reupload_when = self.imm.runtime.now() + duration;
882 let time_period = period.params.time_period();
883
884 info!(
885 time_period=?time_period,
886 "reuploading descriptor in {}",
887 humantime::format_duration(duration),
888 );
889
890 inner.reupload_timers.push(ReuploadTimer {
891 period: time_period,
892 when: reupload_when,
893 });
894
895 let mut upload_results = vec![];
896 for upload_res in results.hsdir_result {
897 let relay = period
898 .hs_dirs
899 .iter_mut()
900 .find(|(relay_ids, _status)| relay_ids == &upload_res.relay_ids);
901
902 let Some((_relay, status)): Option<&mut (RelayIds, _)> = relay else {
903 // This HSDir went away, so the result doesn't matter.
904 // Continue processing the rest of the results
905 continue;
906 };
907
908 if upload_res.upload_res.is_ok() {
909 let update_last_successful = match period.last_successful {
910 None => true,
911 Some(counter) => counter <= upload_res.revision_counter,
912 };
913
914 if update_last_successful {
915 period.last_successful = Some(upload_res.revision_counter);
916 // TODO (#1098): Is it possible that this won't update the statuses promptly
917 // enough. For example, it's possible for the reactor to see a Dirty descriptor
918 // and start an upload task for a descriptor has already been uploaded (or is
919 // being uploaded) in another task, but whose upload results have not yet been
920 // processed.
921 //
922 // This is probably made worse by the fact that the statuses are updated in
923 // batches (grouped by time period), rather than one by one as the upload tasks
924 // complete (updating the status involves locking the inner mutex, and I wanted
925 // to minimize the locking/unlocking overheads). I'm not sure handling the
926 // updates in batches was the correct decision here.
927 *status = DescriptorStatus::Clean;
928 }
929 }
930
931 upload_results.push(upload_res);
932 }
933
934 period.set_upload_results(upload_results);
935 }
936
937 /// Maybe update our list of HsDirs.
938 async fn handle_consensus_change(&mut self, netdir: Arc<NetDir>) -> Result<(), FatalError> {
939 trace!("the consensus has changed; recomputing HSDirs");
940
941 let _old: Option<Arc<NetDir>> = self.replace_netdir(netdir);
942
943 self.recompute_hs_dirs()?;
944 self.update_publish_status_unless_waiting(PublishStatus::UploadScheduled)
945 .await?;
946
947 // If the time period has changed, some of our upload results may now be irrelevant,
948 // so we might need to update our status (for example, if our uploads are
949 // for a no-longer-relevant time period, it means we might be able to update
950 // out status from "degraded" to "running")
951 self.upload_result_to_svc_status()?;
952
953 Ok(())
954 }
955
956 /// Recompute the HsDirs for all relevant time periods.
957 fn recompute_hs_dirs(&self) -> Result<(), FatalError> {
958 let mut inner = self.inner.lock().expect("poisoned lock");
959 let inner = &mut *inner;
960
961 let netdir = Arc::clone(
962 inner
963 .netdir
964 .as_ref()
965 .ok_or_else(|| internal!("started upload task without a netdir"))?,
966 );
967
968 // Update our list of relevant time periods.
969 let new_time_periods = self.compute_time_periods(&netdir, &inner.time_periods)?;
970 inner.time_periods = new_time_periods;
971
972 Ok(())
973 }
974
975 /// Compute the [`TimePeriodContext`]s for the time periods from the specified [`NetDir`].
976 ///
977 /// The specified `time_periods` are used to preserve the `DescriptorStatus` of the
978 /// HsDirs where possible.
979 fn compute_time_periods(
980 &self,
981 netdir: &Arc<NetDir>,
982 time_periods: &[TimePeriodContext],
983 ) -> Result<Vec<TimePeriodContext>, FatalError> {
984 netdir
985 .hs_all_time_periods()
986 .iter()
987 .map(|params| {
988 let period = params.time_period();
989 let blind_id_kp =
990 read_blind_id_keypair(&self.imm.keymgr, &self.imm.nickname, period)?
991 // Note: for now, read_blind_id_keypair cannot return Ok(None).
992 // It's supposed to return Ok(None) if we're in offline hsid mode,
993 // but that might change when we do #1194
994 .ok_or_else(|| internal!("offline hsid mode not supported"))?;
995
996 let blind_id: HsBlindIdKey = (&blind_id_kp).into();
997
998 // If our previous `TimePeriodContext`s also had an entry for `period`, we need to
999 // preserve the `DescriptorStatus` of its HsDirs. This helps prevent unnecessarily
1000 // publishing the descriptor to the HsDirs that already have it (the ones that are
1001 // marked with DescriptorStatus::Clean).
1002 //
1003 // In other words, we only want to publish to those HsDirs that
1004 // * are part of a new time period (which we have never published the descriptor
1005 // for), or
1006 // * have just been added to the ring of a time period we already knew about
1007 if let Some(ctx) = time_periods
1008 .iter()
1009 .find(|ctx| ctx.params.time_period() == period)
1010 {
1011 TimePeriodContext::new(
1012 params.clone(),
1013 blind_id.into(),
1014 netdir,
1015 ctx.hs_dirs.iter(),
1016 ctx.upload_results.clone(),
1017 )
1018 } else {
1019 // Passing an empty iterator here means all HsDirs in this TimePeriodContext
1020 // will be marked as dirty, meaning we will need to upload our descriptor to them.
1021 TimePeriodContext::new(
1022 params.clone(),
1023 blind_id.into(),
1024 netdir,
1025 iter::empty(),
1026 vec![],
1027 )
1028 }
1029 })
1030 .collect::<Result<Vec<TimePeriodContext>, FatalError>>()
1031 }
1032
1033 /// Replace the old netdir with the new, returning the old.
1034 fn replace_netdir(&self, new_netdir: Arc<NetDir>) -> Option<Arc<NetDir>> {
1035 self.inner
1036 .lock()
1037 .expect("poisoned lock")
1038 .netdir
1039 .replace(new_netdir)
1040 }
1041
1042 /// Replace our view of the service config with `new_config` if `new_config` contains changes
1043 /// that would cause us to generate a new descriptor.
1044 fn replace_config_if_changed(&self, new_config: Arc<OnionServiceConfigPublisherView>) -> bool {
1045 let mut inner = self.inner.lock().expect("poisoned lock");
1046 let old_config = &mut inner.config;
1047
1048 // The fields we're interested in haven't changed, so there's no need to update
1049 // `inner.config`.
1050 if *old_config == new_config {
1051 return false;
1052 }
1053
1054 let log_change = match (
1055 old_config.restricted_discovery.enabled,
1056 new_config.restricted_discovery.enabled,
1057 ) {
1058 (true, false) => Some("Disabling restricted discovery mode"),
1059 (false, true) => Some("Enabling restricted discovery mode"),
1060 _ => None,
1061 };
1062
1063 if let Some(msg) = log_change {
1064 info!(nickname=%self.imm.nickname, "{}", msg);
1065 }
1066
1067 let _old: Arc<OnionServiceConfigPublisherView> = std::mem::replace(old_config, new_config);
1068
1069 true
1070 }
1071
1072 /// Recreate the FileWatcher for watching the restricted discovery key_dirs.
1073 fn update_file_watcher(&self) {
1074 let mut inner = self.inner.lock().expect("poisoned lock");
1075 if inner.config.restricted_discovery.watch_configuration() {
1076 debug!("The restricted_discovery.key_dirs have changed, updating file watcher");
1077 let mut watcher = FileWatcher::builder(self.imm.runtime.clone());
1078
1079 let dirs = inner.config.restricted_discovery.key_dirs().clone();
1080
1081 watch_dirs(&mut watcher, &dirs, &self.path_resolver);
1082
1083 let watcher = watcher
1084 .start_watching(self.key_dirs_tx.clone())
1085 .map_err(|e| {
1086 // TODO: update the publish status (see also the module-level TODO about this).
1087 error_report!(e, "Cannot set file watcher");
1088 })
1089 .ok();
1090 inner.file_watcher = watcher;
1091 } else {
1092 if inner.file_watcher.is_some() {
1093 debug!("removing key_dirs watcher");
1094 }
1095 inner.file_watcher = None;
1096 }
1097 }
1098
1099 /// Read the intro points from `ipt_watcher`, and decide whether we're ready to start
1100 /// uploading.
1101 fn note_ipt_change(&self) -> PublishStatus {
1102 let mut ipts = self.ipt_watcher.borrow_for_publish();
1103 match ipts.ipts.as_mut() {
1104 Some(_ipts) => PublishStatus::UploadScheduled,
1105 None => PublishStatus::AwaitingIpts,
1106 }
1107 }
1108
1109 /// Update our list of introduction points.
1110 async fn handle_ipt_change(
1111 &mut self,
1112 update: Option<Result<(), crate::FatalError>>,
1113 ) -> Result<ShutdownStatus, FatalError> {
1114 trace!(nickname=%self.imm.nickname, "received IPT change notification from IPT manager");
1115 match update {
1116 Some(Ok(())) => {
1117 let should_upload = self.note_ipt_change();
1118 debug!(nickname=%self.imm.nickname, "the introduction points have changed");
1119
1120 self.mark_all_dirty();
1121 self.update_publish_status_unless_rate_lim(should_upload)
1122 .await?;
1123 Ok(ShutdownStatus::Continue)
1124 }
1125 Some(Err(e)) => Err(e),
1126 None => {
1127 debug!(nickname=%self.imm.nickname, "received shut down signal from IPT manager");
1128 Ok(ShutdownStatus::Terminate)
1129 }
1130 }
1131 }
1132
1133 /// Update the `PublishStatus` of the reactor with `new_state`,
1134 /// unless the current state is `AwaitingIpts`.
1135 async fn update_publish_status_unless_waiting(
1136 &mut self,
1137 new_state: PublishStatus,
1138 ) -> Result<(), FatalError> {
1139 // Only update the state if we're not waiting for intro points.
1140 if self.status() != PublishStatus::AwaitingIpts {
1141 self.update_publish_status(new_state).await?;
1142 }
1143
1144 Ok(())
1145 }
1146
1147 /// Update the `PublishStatus` of the reactor with `new_state`,
1148 /// unless the current state is `RateLimited`.
1149 async fn update_publish_status_unless_rate_lim(
1150 &mut self,
1151 new_state: PublishStatus,
1152 ) -> Result<(), FatalError> {
1153 // We can't exit this state until the rate-limit expires.
1154 if !matches!(self.status(), PublishStatus::RateLimited(_)) {
1155 self.update_publish_status(new_state).await?;
1156 }
1157
1158 Ok(())
1159 }
1160
1161 /// Unconditionally update the `PublishStatus` of the reactor with `new_state`.
1162 async fn update_publish_status(&mut self, new_state: PublishStatus) -> Result<(), Bug> {
1163 let onion_status = match new_state {
1164 PublishStatus::Idle => None,
1165 PublishStatus::UploadScheduled
1166 | PublishStatus::AwaitingIpts
1167 | PublishStatus::RateLimited(_) => Some(State::Bootstrapping),
1168 };
1169
1170 if let Some(onion_status) = onion_status {
1171 self.imm.status_tx.send(onion_status, None);
1172 }
1173
1174 trace!(
1175 "publisher reactor status change: {:?} -> {:?}",
1176 self.status(),
1177 new_state
1178 );
1179
1180 self.publish_status_tx.send(new_state).await.map_err(
1181 |_: postage::sink::SendError<_>| internal!("failed to send upload notification?!"),
1182 )?;
1183
1184 Ok(())
1185 }
1186
1187 /// Update the onion svc status based on the results of the last descriptor uploads.
1188 fn upload_result_to_svc_status(&self) -> Result<(), FatalError> {
1189 let inner = self.inner.lock().expect("poisoned lock");
1190 let netdir = inner
1191 .netdir
1192 .as_ref()
1193 .ok_or_else(|| internal!("handling upload results without netdir?!"))?;
1194
1195 let (state, err) = upload_result_state(netdir, &inner.time_periods);
1196 self.imm.status_tx.send(state, err);
1197
1198 Ok(())
1199 }
1200
1201 /// Update the descriptors based on the config change.
1202 async fn handle_svc_config_change(
1203 &mut self,
1204 config: &OnionServiceConfig,
1205 ) -> Result<(), FatalError> {
1206 let new_config = Arc::new(config.into());
1207 if self.replace_config_if_changed(Arc::clone(&new_config)) {
1208 self.update_file_watcher();
1209 self.update_authorized_clients_if_changed();
1210
1211 info!(nickname=%self.imm.nickname, "Config has changed, generating a new descriptor");
1212 self.mark_all_dirty();
1213
1214 // Schedule an upload, unless we're still waiting for IPTs.
1215 self.update_publish_status_unless_waiting(PublishStatus::UploadScheduled)
1216 .await?;
1217 }
1218
1219 Ok(())
1220 }
1221
1222 /// Update the descriptors based on a restricted discovery key_dirs change.
1223 ///
1224 /// If the authorized clients from the [`RestrictedDiscoveryConfig`] have changed,
1225 /// this marks the descriptor as dirty for all time periods,
1226 /// and schedules a reupload.
1227 async fn handle_key_dirs_change(&mut self, event: FileEvent) -> Result<(), FatalError> {
1228 debug!("The configured key_dirs have changed");
1229 match event {
1230 FileEvent::Rescan | FileEvent::FileChanged => {
1231 // These events are handled in the same way, by re-reading the keys from disk
1232 // and republishing the descriptor if necessary
1233 }
1234 _ => return Err(internal!("file watcher event {event:?}").into()),
1235 };
1236
1237 // Update the file watcher, in case the change was triggered by a key_dir move.
1238 self.update_file_watcher();
1239
1240 if self.update_authorized_clients_if_changed() {
1241 self.mark_all_dirty();
1242
1243 // Schedule an upload, unless we're still waiting for IPTs.
1244 self.update_publish_status_unless_waiting(PublishStatus::UploadScheduled)
1245 .await?;
1246 }
1247
1248 Ok(())
1249 }
1250
1251 /// Recreate the authorized_clients based on the current config.
1252 ///
1253 /// Returns `true` if the authorized clients have changed.
1254 fn update_authorized_clients_if_changed(&mut self) -> bool {
1255 let mut inner = self.inner.lock().expect("poisoned lock");
1256 let authorized_clients =
1257 Self::read_authorized_clients(&inner.config.restricted_discovery, &self.path_resolver);
1258
1259 let clients = &mut inner.authorized_clients;
1260 let changed = clients.as_ref() != authorized_clients.as_ref();
1261
1262 if changed {
1263 info!("The restricted discovery mode authorized clients have changed");
1264 *clients = authorized_clients;
1265 }
1266
1267 changed
1268 }
1269
1270 /// Read the authorized `RestrictedDiscoveryKeys` from `config`.
1271 fn read_authorized_clients(
1272 config: &RestrictedDiscoveryConfig,
1273 path_resolver: &CfgPathResolver,
1274 ) -> Option<Arc<RestrictedDiscoveryKeys>> {
1275 let authorized_clients = config.read_keys(path_resolver);
1276
1277 if matches!(authorized_clients.as_ref(), Some(c) if c.is_empty()) {
1278 warn!(
1279 "Running in restricted discovery mode, but we have no authorized clients. Service will be unreachable"
1280 );
1281 }
1282
1283 authorized_clients.map(Arc::new)
1284 }
1285
1286 /// Mark the descriptor dirty for all time periods.
1287 fn mark_all_dirty(&self) {
1288 trace!("marking the descriptor dirty for all time periods");
1289
1290 self.inner
1291 .lock()
1292 .expect("poisoned lock")
1293 .time_periods
1294 .iter_mut()
1295 .for_each(|tp| tp.mark_all_dirty());
1296 }
1297
1298 /// Mark the descriptor dirty for the specified time period.
1299 ///
1300 /// Returns `true` if the specified period is still relevant, and `false` otherwise.
1301 fn mark_dirty(&self, period: &TimePeriod) -> bool {
1302 let mut inner = self.inner.lock().expect("poisoned lock");
1303 let period_ctx = inner
1304 .time_periods
1305 .iter_mut()
1306 .find(|tp| tp.params.time_period() == *period);
1307
1308 match period_ctx {
1309 Some(ctx) => {
1310 trace!(time_period=?period, "marking the descriptor dirty");
1311 ctx.mark_all_dirty();
1312 true
1313 }
1314 None => false,
1315 }
1316 }
1317
1318 /// Try to upload our descriptor to the HsDirs that need it.
1319 ///
1320 /// If we've recently uploaded some descriptors, we return immediately and schedule the upload
1321 /// to happen after [`UPLOAD_RATE_LIM_THRESHOLD`].
1322 ///
1323 /// Failed uploads are retried
1324 /// (see [`upload_descriptor_with_retries`](Reactor::upload_descriptor_with_retries)).
1325 ///
1326 /// If restricted discovery mode is enabled and there are no authorized clients,
1327 /// we abort the upload and set our status to [`State::Broken`].
1328 //
1329 // Note: a broken restricted discovery config won't prevent future uploads from being scheduled
1330 // (for example if the IPTs change),
1331 // which can can cause the publisher's status to oscillate between `Bootstrapping` and `Broken`.
1332 // TODO: we might wish to refactor the publisher to be more sophisticated about this.
1333 //
1334 /// For each current time period, we spawn a task that uploads the descriptor to
1335 /// all the HsDirs on the HsDir ring of that time period.
1336 /// Each task shuts down on completion, or when the reactor is dropped.
1337 ///
1338 /// Each task reports its upload results (`TimePeriodUploadResult`)
1339 /// via the `upload_task_complete_tx` channel.
1340 /// The results are received and processed in the main loop of the reactor.
1341 ///
1342 /// Returns an error if it fails to spawn a task, or if an internal error occurs.
1343 async fn upload_all(&mut self) -> Result<(), FatalError> {
1344 trace!("starting descriptor upload task...");
1345
1346 // Abort the upload entirely if we have an empty list of authorized clients
1347 let authorized_clients = match self.authorized_clients() {
1348 Ok(authorized_clients) => authorized_clients,
1349 Err(e) => {
1350 error_report!(e, "aborting upload");
1351 self.imm.status_tx.send_broken(e.clone());
1352
1353 // Returning an error would shut down the reactor, so we have to return Ok here.
1354 return Ok(());
1355 }
1356 };
1357
1358 let last_uploaded = self.inner.lock().expect("poisoned lock").last_uploaded;
1359 let now = self.imm.runtime.now();
1360 // Check if we should rate-limit this upload.
1361 if let Some(ts) = last_uploaded {
1362 let duration_since_upload = now.duration_since(ts);
1363
1364 if duration_since_upload < UPLOAD_RATE_LIM_THRESHOLD {
1365 return Ok(self.start_rate_limit(UPLOAD_RATE_LIM_THRESHOLD).await?);
1366 }
1367 }
1368
1369 let mut inner = self.inner.lock().expect("poisoned lock");
1370 let inner = &mut *inner;
1371
1372 let _ = inner.last_uploaded.insert(now);
1373
1374 for period_ctx in inner.time_periods.iter_mut() {
1375 let upload_task_complete_tx = self.upload_task_complete_tx.clone();
1376
1377 // Figure out which HsDirs we need to upload the descriptor to (some of them might already
1378 // have our latest descriptor, so we filter them out).
1379 let hs_dirs = period_ctx
1380 .hs_dirs
1381 .iter()
1382 .filter_map(|(relay_id, status)| {
1383 if *status == DescriptorStatus::Dirty {
1384 Some(relay_id.clone())
1385 } else {
1386 None
1387 }
1388 })
1389 .collect::<Vec<_>>();
1390
1391 if hs_dirs.is_empty() {
1392 trace!("the descriptor is clean for all HSDirs. Nothing to do");
1393 return Ok(());
1394 }
1395
1396 let time_period = period_ctx.params.time_period();
1397 // This scope exists because rng is not Send, so it needs to fall out of scope before we
1398 // await anything.
1399 let netdir = Arc::clone(
1400 inner
1401 .netdir
1402 .as_ref()
1403 .ok_or_else(|| internal!("started upload task without a netdir"))?,
1404 );
1405
1406 let imm = Arc::clone(&self.imm);
1407 let ipt_upload_view = self.ipt_watcher.upload_view();
1408 let config = Arc::clone(&inner.config);
1409 let authorized_clients = authorized_clients.clone();
1410
1411 trace!(nickname=%self.imm.nickname, time_period=?time_period,
1412 "spawning upload task"
1413 );
1414
1415 let params = period_ctx.params.clone();
1416 let shutdown_rx = self.shutdown_tx.subscribe();
1417
1418 // Spawn a task to upload the descriptor to all HsDirs of this time period.
1419 //
1420 // This task will shut down when the reactor is dropped (i.e. when shutdown_rx is
1421 // dropped).
1422 let _handle: () = self
1423 .imm
1424 .runtime
1425 .spawn(async move {
1426 if let Err(e) = Self::upload_for_time_period(
1427 hs_dirs,
1428 &netdir,
1429 config,
1430 params,
1431 Arc::clone(&imm),
1432 ipt_upload_view.clone(),
1433 authorized_clients.clone(),
1434 upload_task_complete_tx,
1435 shutdown_rx,
1436 )
1437 .await
1438 {
1439 error_report!(
1440 e,
1441 "descriptor upload failed for HS service {} and time period {:?}",
1442 imm.nickname,
1443 time_period
1444 );
1445 }
1446 })
1447 .map_err(|e| FatalError::from_spawn("upload_for_time_period task", e))?;
1448 }
1449
1450 Ok(())
1451 }
1452
1453 /// Upload the descriptor for the time period specified in `params`.
1454 ///
1455 /// Failed uploads are retried
1456 /// (see [`upload_descriptor_with_retries`](Reactor::upload_descriptor_with_retries)).
1457 #[allow(clippy::too_many_arguments)] // TODO: refactor
1458 async fn upload_for_time_period(
1459 hs_dirs: Vec<RelayIds>,
1460 netdir: &Arc<NetDir>,
1461 config: Arc<OnionServiceConfigPublisherView>,
1462 params: HsDirParams,
1463 imm: Arc<Immutable<R, M>>,
1464 ipt_upload_view: IptsPublisherUploadView,
1465 authorized_clients: Option<Arc<RestrictedDiscoveryKeys>>,
1466 mut upload_task_complete_tx: mpsc::Sender<TimePeriodUploadResult>,
1467 shutdown_rx: broadcast::Receiver<Void>,
1468 ) -> Result<(), FatalError> {
1469 let time_period = params.time_period();
1470 trace!(time_period=?time_period, "uploading descriptor to all HSDirs for this time period");
1471
1472 let hsdir_count = hs_dirs.len();
1473
1474 /// An error returned from an upload future.
1475 //
1476 // Exhaustive, because this is a private type.
1477 #[derive(Clone, Debug, thiserror::Error)]
1478 enum PublishError {
1479 /// The upload was aborted because there are no IPTs.
1480 ///
1481 /// This happens because of an inevitable TOCTOU race, where after being notified by
1482 /// the IPT manager that the IPTs have changed (via `self.ipt_watcher.await_update`),
1483 /// we find out there actually are no IPTs, so we can't build the descriptor.
1484 ///
1485 /// This is a special kind of error that interrupts the current upload task, and is
1486 /// logged at `debug!` level rather than `warn!` or `error!`.
1487 ///
1488 /// Ideally, this shouldn't happen very often (if at all).
1489 #[error("No IPTs")]
1490 NoIpts,
1491
1492 /// The reactor has shut down
1493 #[error("The reactor has shut down")]
1494 Shutdown,
1495
1496 /// An fatal error.
1497 #[error("{0}")]
1498 Fatal(#[from] FatalError),
1499 }
1500
1501 let max_hsdesc_len: usize = netdir
1502 .params()
1503 .hsdir_max_desc_size
1504 .try_into()
1505 .expect("Unable to convert positive int32 to usize!?");
1506
1507 let upload_results = futures::stream::iter(hs_dirs)
1508 .map(|relay_ids| {
1509 let netdir = netdir.clone();
1510 let config = Arc::clone(&config);
1511 let imm = Arc::clone(&imm);
1512 let ipt_upload_view = ipt_upload_view.clone();
1513 let authorized_clients = authorized_clients.clone();
1514 let params = params.clone();
1515 let mut shutdown_rx = shutdown_rx.clone();
1516
1517 let ed_id = relay_ids
1518 .rsa_identity()
1519 .map(|id| id.to_string())
1520 .unwrap_or_else(|| "unknown".into());
1521 let rsa_id = relay_ids
1522 .rsa_identity()
1523 .map(|id| id.to_string())
1524 .unwrap_or_else(|| "unknown".into());
1525
1526 async move {
1527 let run_upload = |desc| async {
1528 let Some(hsdir) = netdir.by_ids(&relay_ids) else {
1529 // This should never happen (all of our relay_ids are from the stored
1530 // netdir).
1531 let err =
1532 "tried to upload descriptor to relay not found in consensus?!";
1533 warn!(
1534 nickname=%imm.nickname, hsdir_id=%ed_id, hsdir_rsa_id=%rsa_id,
1535 "{err}"
1536 );
1537 return Err(internal!("{err}").into());
1538 };
1539
1540 Self::upload_descriptor_with_retries(
1541 desc,
1542 &netdir,
1543 &hsdir,
1544 &ed_id,
1545 &rsa_id,
1546 Arc::clone(&imm),
1547 )
1548 .await
1549 };
1550
1551 // How long until we're supposed to time out?
1552 let worst_case_end = imm.runtime.now() + OVERALL_UPLOAD_TIMEOUT;
1553 // We generate a new descriptor before _each_ HsDir upload. This means each
1554 // HsDir could, in theory, receive a different descriptor (not just in terms of
1555 // revision-counters, but also with a different set of IPTs). It may seem like
1556 // this could lead to some HsDirs being left with an outdated descriptor, but
1557 // that's not the case: after the upload completes, the publisher will be
1558 // notified by the ipt_watcher of the IPT change event (if there was one to
1559 // begin with), which will trigger another upload job.
1560 let hsdesc = {
1561 // This scope is needed because the ipt_set MutexGuard is not Send, so it
1562 // needs to fall out of scope before the await point below
1563 let mut ipt_set = ipt_upload_view.borrow_for_publish();
1564
1565 // If there are no IPTs, we abort the upload. At this point, we might have
1566 // uploaded the descriptor to some, but not all, HSDirs from the specified
1567 // time period.
1568 //
1569 // Returning an error here means the upload completion task is never
1570 // notified of the outcome of any of these uploads (which means the
1571 // descriptor is not marked clean). This is OK, because if we suddenly find
1572 // out we have no IPTs, it means our built `hsdesc` has an outdated set of
1573 // IPTs, so we need to go back to the main loop to wait for IPT changes,
1574 // and generate a fresh descriptor anyway.
1575 //
1576 // Ideally, this shouldn't happen very often (if at all).
1577 let Some(ipts) = ipt_set.ipts.as_mut() else {
1578 return Err(PublishError::NoIpts);
1579 };
1580
1581 let hsdesc = {
1582 trace!(
1583 nickname=%imm.nickname, time_period=?time_period,
1584 "building descriptor"
1585 );
1586 let mut rng = imm.mockable.thread_rng();
1587 let mut key_rng = tor_llcrypto::rng::CautiousRng;
1588
1589 // We're about to generate a new version of the descriptor,
1590 // so let's generate a new revision counter.
1591 let now = imm.runtime.wallclock();
1592 let revision_counter = imm.generate_revision_counter(¶ms, now)?;
1593
1594 build_sign(
1595 &imm.keymgr,
1596 &imm.pow_manager,
1597 &config,
1598 netdir.params(),
1599 authorized_clients.as_deref(),
1600 ipts,
1601 time_period,
1602 revision_counter,
1603 &mut rng,
1604 &mut key_rng,
1605 imm.runtime.wallclock(),
1606 max_hsdesc_len,
1607 )?
1608 };
1609
1610 if let Err(e) =
1611 ipt_set.note_publication_attempt(&imm.runtime, worst_case_end)
1612 {
1613 let wait = e.log_retry_max(&imm.nickname)?;
1614 // TODO (#1226): retry instead of this
1615 return Err(FatalError::Bug(internal!(
1616 "ought to retry after {wait:?}, crashing instead"
1617 ))
1618 .into());
1619 }
1620
1621 hsdesc
1622 };
1623
1624 let VersionedDescriptor {
1625 desc,
1626 revision_counter,
1627 } = hsdesc;
1628 let desc: Arc<str> = desc.into();
1629
1630 trace!(
1631 nickname=%imm.nickname, time_period=?time_period,
1632 revision_counter=?revision_counter,
1633 "generated new descriptor for time period",
1634 );
1635
1636 // (Actually launch the upload attempt. No timeout is needed
1637 // here, since the backoff::Runner code will handle that for us.)
1638 let upload_res: UploadResult = select_biased! {
1639 shutdown = shutdown_rx.next().fuse() => {
1640 // This will always be None, since Void is uninhabited.
1641 let _: Option<Void> = shutdown;
1642
1643 // It looks like the reactor has shut down,
1644 // so there is no point in uploading the descriptor anymore.
1645 //
1646 // Let's shut down the upload task too.
1647 trace!(
1648 nickname=%imm.nickname, time_period=?time_period,
1649 "upload task received shutdown signal"
1650 );
1651
1652 return Err(PublishError::Shutdown);
1653 },
1654 res = run_upload(desc.clone()).fuse() => res,
1655 };
1656
1657 // Note: UploadResult::Failure is only returned when
1658 // upload_descriptor_with_retries fails, i.e. if all our retry
1659 // attempts have failed
1660 Ok(HsDirUploadStatus {
1661 relay_ids,
1662 upload_res,
1663 revision_counter,
1664 })
1665 }
1666 })
1667 // This fails to compile unless the stream is boxed. See https://github.com/rust-lang/rust/issues/104382
1668 .boxed()
1669 .buffer_unordered(MAX_CONCURRENT_UPLOADS)
1670 .try_collect::<Vec<_>>()
1671 .await;
1672
1673 let upload_results = match upload_results {
1674 Ok(v) => v,
1675 Err(PublishError::Fatal(e)) => return Err(e),
1676 Err(PublishError::NoIpts) => {
1677 debug!(
1678 nickname=%imm.nickname, time_period=?time_period,
1679 "no introduction points; skipping upload"
1680 );
1681
1682 return Ok(());
1683 }
1684 Err(PublishError::Shutdown) => {
1685 debug!(
1686 nickname=%imm.nickname, time_period=?time_period,
1687 "the reactor has shut down; aborting upload"
1688 );
1689
1690 return Ok(());
1691 }
1692 };
1693
1694 let (succeeded, _failed): (Vec<_>, Vec<_>) = upload_results
1695 .iter()
1696 .partition(|res| res.upload_res.is_ok());
1697
1698 debug!(
1699 nickname=%imm.nickname, time_period=?time_period,
1700 "descriptor uploaded successfully to {}/{} HSDirs",
1701 succeeded.len(), hsdir_count
1702 );
1703
1704 if upload_task_complete_tx
1705 .send(TimePeriodUploadResult {
1706 time_period,
1707 hsdir_result: upload_results,
1708 })
1709 .await
1710 .is_err()
1711 {
1712 return Err(internal!(
1713 "failed to notify reactor of upload completion (reactor shut down)"
1714 )
1715 .into());
1716 }
1717
1718 Ok(())
1719 }
1720
1721 /// Upload a descriptor to the specified HSDir.
1722 ///
1723 /// If an upload fails, this returns an `Err`. This function does not handle retries. It is up
1724 /// to the caller to retry on failure.
1725 ///
1726 /// This function does not handle timeouts.
1727 async fn upload_descriptor(
1728 hsdesc: Arc<str>,
1729 netdir: &Arc<NetDir>,
1730 hsdir: &Relay<'_>,
1731 imm: Arc<Immutable<R, M>>,
1732 ) -> Result<(), UploadError> {
1733 let request = HsDescUploadRequest::new(hsdesc);
1734
1735 trace!(nickname=%imm.nickname, hsdir_id=%hsdir.id(), hsdir_rsa_id=%hsdir.rsa_id(),
1736 "starting descriptor upload",
1737 );
1738
1739 let tunnel = imm
1740 .mockable
1741 .get_or_launch_hs_dir(netdir, OwnedCircTarget::from_circ_target(hsdir))
1742 .await?;
1743 let source: Option<SourceInfo> = tunnel
1744 .source_info()
1745 .map_err(into_internal!("Couldn't get SourceInfo for circuit"))?;
1746
1747 let mut stream = tunnel
1748 .begin_dir_stream()
1749 .await
1750 .map_err(UploadError::Stream)?;
1751
1752 let _response: String = send_request(&imm.runtime, &request, &mut stream, source)
1753 .await
1754 .map_err(|dir_error| -> UploadError {
1755 match dir_error {
1756 DirClientError::RequestFailed(e) => e.into(),
1757 DirClientError::CircMgr(e) => into_internal!(
1758 "tor-dirclient complains about circmgr going wrong but we gave it a stream"
1759 )(e)
1760 .into(),
1761 e => into_internal!("unexpected error")(e).into(),
1762 }
1763 })?
1764 .into_output_string()?; // This returns an error if we received an error response
1765
1766 Ok(())
1767 }
1768
1769 /// Upload a descriptor to the specified HSDir, retrying if appropriate.
1770 ///
1771 /// Any failed uploads are retried according to a [`PublisherBackoffSchedule`].
1772 /// Each failed upload is retried until it succeeds, or until the overall timeout specified
1773 /// by [`BackoffSchedule::overall_timeout`] elapses. Individual attempts are timed out
1774 /// according to the [`BackoffSchedule::single_attempt_timeout`].
1775 /// This function gives up after the overall timeout elapses,
1776 /// declaring the upload a failure, and never retrying it again.
1777 ///
1778 /// See also [`BackoffSchedule`].
1779 async fn upload_descriptor_with_retries(
1780 hsdesc: Arc<str>,
1781 netdir: &Arc<NetDir>,
1782 hsdir: &Relay<'_>,
1783 ed_id: &str,
1784 rsa_id: &str,
1785 imm: Arc<Immutable<R, M>>,
1786 ) -> UploadResult {
1787 /// The base delay to use for the backoff schedule.
1788 const BASE_DELAY_MSEC: u32 = 1000;
1789 let schedule = PublisherBackoffSchedule {
1790 retry_delay: RetryDelay::from_msec(BASE_DELAY_MSEC),
1791 mockable: imm.mockable.clone(),
1792 };
1793
1794 let runner = Runner::new(
1795 "upload a hidden service descriptor".into(),
1796 schedule.clone(),
1797 imm.runtime.clone(),
1798 );
1799
1800 let fallible_op = || async {
1801 let r = Self::upload_descriptor(hsdesc.clone(), netdir, hsdir, Arc::clone(&imm)).await;
1802
1803 if let Err(e) = &r {
1804 if e.should_report_as_suspicious() {
1805 // Note that not every protocol violation is suspicious:
1806 // we only warn on the protocol violations that look like attempts
1807 // to do a traffic tagging attack via hsdir inflation.
1808 // (See proposal 360.)
1809 warn_report!(
1810 e,
1811 "Suspicious error while uploading descriptor to {}/{}",
1812 ed_id,
1813 rsa_id
1814 );
1815 }
1816 }
1817 r
1818 };
1819
1820 let outcome: Result<(), BackoffError<UploadError>> = runner.run(fallible_op).await;
1821 match outcome {
1822 Ok(()) => {
1823 debug!(
1824 nickname=%imm.nickname, hsdir_id=%ed_id, hsdir_rsa_id=%rsa_id,
1825 "successfully uploaded descriptor to HSDir",
1826 );
1827
1828 Ok(())
1829 }
1830 Err(e) => {
1831 warn_report!(
1832 e,
1833 "failed to upload descriptor for service {} (hsdir_id={}, hsdir_rsa_id={})",
1834 imm.nickname,
1835 ed_id,
1836 rsa_id
1837 );
1838
1839 Err(e.into())
1840 }
1841 }
1842 }
1843
1844 /// Stop publishing descriptors until the specified delay elapses.
1845 async fn start_rate_limit(&mut self, delay: Duration) -> Result<(), Bug> {
1846 if !matches!(self.status(), PublishStatus::RateLimited(_)) {
1847 debug!(
1848 "We are rate-limited for {}; pausing descriptor publication",
1849 humantime::format_duration(delay)
1850 );
1851 let until = self.imm.runtime.now() + delay;
1852 self.update_publish_status(PublishStatus::RateLimited(until))
1853 .await?;
1854 }
1855
1856 Ok(())
1857 }
1858
1859 /// Handle the upload rate-limit being lifted.
1860 async fn expire_rate_limit(&mut self) -> Result<(), Bug> {
1861 debug!("We are no longer rate-limited; resuming descriptor publication");
1862 self.update_publish_status(PublishStatus::UploadScheduled)
1863 .await?;
1864 Ok(())
1865 }
1866
1867 /// Return the authorized clients, if restricted mode is enabled.
1868 ///
1869 /// Returns `Ok(None)` if restricted discovery mode is disabled.
1870 ///
1871 /// Returns an error if restricted discovery mode is enabled, but the client list is empty.
1872 #[cfg_attr(
1873 not(feature = "restricted-discovery"),
1874 allow(clippy::unnecessary_wraps)
1875 )]
1876 fn authorized_clients(&self) -> Result<Option<Arc<RestrictedDiscoveryKeys>>, FatalError> {
1877 cfg_if::cfg_if! {
1878 if #[cfg(feature = "restricted-discovery")] {
1879 let authorized_clients = self
1880 .inner
1881 .lock()
1882 .expect("poisoned lock")
1883 .authorized_clients
1884 .clone();
1885
1886 if authorized_clients.as_ref().as_ref().map(|v| v.is_empty()).unwrap_or_default() {
1887 return Err(FatalError::RestrictedDiscoveryNoClients);
1888 }
1889
1890 Ok(authorized_clients)
1891 } else {
1892 Ok(None)
1893 }
1894 }
1895 }
1896}
1897
1898/// Try to expand a path, logging a warning on failure.
1899fn maybe_expand_path(p: &CfgPath, r: &CfgPathResolver) -> Option<PathBuf> {
1900 // map_err returns unit for clarity
1901 #[allow(clippy::unused_unit, clippy::semicolon_if_nothing_returned)]
1902 p.path(r)
1903 .map_err(|e| {
1904 tor_error::warn_report!(e, "invalid path");
1905 ()
1906 })
1907 .ok()
1908}
1909
1910/// Add `path` to the specified `watcher`.
1911macro_rules! watch_path {
1912 ($watcher:expr, $path:expr, $watch_fn:ident, $($watch_fn_args:expr,)*) => {{
1913 if let Err(e) = $watcher.$watch_fn(&$path, $($watch_fn_args)*) {
1914 warn_report!(e, "failed to watch path {:?}", $path);
1915 } else {
1916 debug!("watching path {:?}", $path);
1917 }
1918 }}
1919}
1920
1921/// Add the specified directories to the watcher.
1922fn watch_dirs<R: Runtime>(
1923 watcher: &mut FileWatcherBuilder<R>,
1924 dirs: &DirectoryKeyProviderList,
1925 path_resolver: &CfgPathResolver,
1926) {
1927 for path in dirs {
1928 let path = path.path();
1929 let Some(path) = maybe_expand_path(path, path_resolver) else {
1930 warn!("failed to expand key_dir path {:?}", path);
1931 continue;
1932 };
1933
1934 // If the path doesn't exist, the notify watcher will return an error if we attempt to watch it,
1935 // so we skip over paths that don't exist at this time
1936 // (this obviously suffers from a TOCTOU race, but most of the time,
1937 // it is good enough at preventing the watcher from failing to watch.
1938 // If the race *does* happen it is not disastrous, i.e. the reactor won't crash,
1939 // but it will fail to set the watcher).
1940 if matches!(path.try_exists(), Ok(true)) {
1941 watch_path!(watcher, &path, watch_dir, "auth",);
1942 }
1943 // FileWatcher::watch_path causes the parent dir of the path to be watched.
1944 if matches!(path.parent().map(|p| p.try_exists()), Some(Ok(true))) {
1945 watch_path!(watcher, &path, watch_path,);
1946 }
1947 }
1948}
1949
1950/// Try to read the blinded identity key for a given `TimePeriod`.
1951///
1952/// Returns `None` if the service is running in "offline" mode.
1953///
1954// TODO (#1194): we don't currently have support for "offline" mode so this can never return
1955// `Ok(None)`.
1956pub(super) fn read_blind_id_keypair(
1957 keymgr: &Arc<KeyMgr>,
1958 nickname: &HsNickname,
1959 period: TimePeriod,
1960) -> Result<Option<HsBlindIdKeypair>, FatalError> {
1961 let svc_key_spec = HsIdKeypairSpecifier::new(nickname.clone());
1962 let hsid_kp = keymgr
1963 .get::<HsIdKeypair>(&svc_key_spec)?
1964 .ok_or_else(|| FatalError::MissingHsIdKeypair(nickname.clone()))?;
1965
1966 let blind_id_key_spec = BlindIdKeypairSpecifier::new(nickname.clone(), period);
1967
1968 // TODO: make the keystore selector configurable
1969 let keystore_selector = Default::default();
1970 match keymgr.get::<HsBlindIdKeypair>(&blind_id_key_spec)? {
1971 Some(kp) => Ok(Some(kp)),
1972 None => {
1973 let (_hs_blind_id_key, hs_blind_id_kp, _subcredential) = hsid_kp
1974 .compute_blinded_key(period)
1975 .map_err(|_| internal!("failed to compute blinded key"))?;
1976
1977 // Note: we can't use KeyMgr::generate because this key is derived from the HsId
1978 // (KeyMgr::generate uses the tor_keymgr::Keygen trait under the hood,
1979 // which assumes keys are randomly generated, rather than derived from existing keys).
1980
1981 keymgr.insert(hs_blind_id_kp, &blind_id_key_spec, keystore_selector, true)?;
1982
1983 let arti_path = |spec: &dyn KeySpecifier| {
1984 spec.arti_path()
1985 .map_err(into_internal!("invalid key specifier?!"))
1986 };
1987
1988 Ok(Some(
1989 keymgr.get::<HsBlindIdKeypair>(&blind_id_key_spec)?.ok_or(
1990 FatalError::KeystoreRace {
1991 action: "read",
1992 path: arti_path(&blind_id_key_spec)?,
1993 },
1994 )?,
1995 ))
1996 }
1997 }
1998}
1999
2000/// Determine the [`State`] of the publisher based on the upload results
2001/// from the current `time_periods`.
2002fn upload_result_state(
2003 netdir: &NetDir,
2004 time_periods: &[TimePeriodContext],
2005) -> (State, Option<Problem>) {
2006 let current_period = netdir.hs_time_period();
2007 let current_period_res = time_periods
2008 .iter()
2009 .find(|ctx| ctx.params.time_period() == current_period);
2010
2011 let succeeded_current_tp = current_period_res
2012 .iter()
2013 .flat_map(|res| &res.upload_results)
2014 .filter(|res| res.upload_res.is_ok())
2015 .collect_vec();
2016
2017 let secondary_tp_res = time_periods
2018 .iter()
2019 .filter(|ctx| ctx.params.time_period() != current_period)
2020 .collect_vec();
2021
2022 let succeeded_secondary_tp = secondary_tp_res
2023 .iter()
2024 .flat_map(|res| &res.upload_results)
2025 .filter(|res| res.upload_res.is_ok())
2026 .collect_vec();
2027
2028 // All of the failed uploads (for all TPs)
2029 let failed = time_periods
2030 .iter()
2031 .flat_map(|res| &res.upload_results)
2032 .filter(|res| res.upload_res.is_err())
2033 .collect_vec();
2034 let problems: Vec<DescUploadRetryError> = failed
2035 .iter()
2036 .flat_map(|e| e.upload_res.as_ref().map_err(|e| e.clone()).err())
2037 .collect();
2038
2039 let err = match problems.as_slice() {
2040 [_, ..] => Some(problems.into()),
2041 [] => None,
2042 };
2043
2044 if time_periods.len() < 2 {
2045 // We need at least TP contexts (one for the primary TP,
2046 // and another for the secondary one).
2047 //
2048 // If either is missing, we are unreachable for some or all clients.
2049 return (State::DegradedUnreachable, err);
2050 }
2051
2052 let state = match (
2053 succeeded_current_tp.as_slice(),
2054 succeeded_secondary_tp.as_slice(),
2055 ) {
2056 (&[], &[..]) | (&[..], &[]) if failed.is_empty() => {
2057 // We don't have any upload results for one or both TPs.
2058 // We are still bootstrapping.
2059 State::Bootstrapping
2060 }
2061 (&[_, ..], &[_, ..]) if failed.is_empty() => {
2062 // We have uploaded the descriptor to one or more HsDirs from both
2063 // HsDir rings (primary and secondary), and none of the uploads failed.
2064 // We are fully reachable.
2065 State::Running
2066 }
2067 (&[_, ..], &[_, ..]) => {
2068 // We have uploaded the descriptor to one or more HsDirs from both
2069 // HsDir rings (primary and secondary), but some of the uploads failed.
2070 // We are reachable, but we failed to upload the descriptor to all the HsDirs
2071 // that were supposed to have it.
2072 State::DegradedReachable
2073 }
2074 (&[..], &[]) | (&[], &[..]) => {
2075 // We have either
2076 // * uploaded the descriptor to some of the HsDirs from one of the rings,
2077 // but haven't managed to upload it to any of the HsDirs on the other ring, or
2078 // * all of the uploads failed
2079 //
2080 // Either way, we are definitely not reachable by all clients.
2081 State::DegradedUnreachable
2082 }
2083 };
2084
2085 (state, err)
2086}
2087
2088/// Whether the reactor should initiate an upload.
2089#[derive(Copy, Clone, Debug, Default, PartialEq)]
2090enum PublishStatus {
2091 /// We need to call upload_all.
2092 UploadScheduled,
2093 /// We are rate-limited until the specified [`Instant`].
2094 ///
2095 /// We have tried to schedule multiple uploads in a short time span,
2096 /// and we are rate-limited. We are waiting for a signal from the schedule_upload_tx
2097 /// channel to unblock us.
2098 RateLimited(Instant),
2099 /// We are idle and waiting for external events.
2100 ///
2101 /// We have enough information to build the descriptor, but since we have already called
2102 /// upload_all to upload it to all relevant HSDirs, there is nothing for us to do right nbow.
2103 Idle,
2104 /// We are waiting for the IPT manager to establish some introduction points.
2105 ///
2106 /// No descriptors will be published until the `PublishStatus` of the reactor is changed to
2107 /// `UploadScheduled`.
2108 #[default]
2109 AwaitingIpts,
2110}
2111
2112/// The backoff schedule for the task that publishes descriptors.
2113#[derive(Clone, Debug)]
2114struct PublisherBackoffSchedule<M: Mockable> {
2115 /// The delays
2116 retry_delay: RetryDelay,
2117 /// The mockable reactor state, needed for obtaining an rng.
2118 mockable: M,
2119}
2120
2121impl<M: Mockable> BackoffSchedule for PublisherBackoffSchedule<M> {
2122 fn max_retries(&self) -> Option<usize> {
2123 None
2124 }
2125
2126 fn overall_timeout(&self) -> Option<Duration> {
2127 Some(OVERALL_UPLOAD_TIMEOUT)
2128 }
2129
2130 fn single_attempt_timeout(&self) -> Option<Duration> {
2131 Some(self.mockable.estimate_upload_timeout())
2132 }
2133
2134 fn next_delay<E: RetriableError>(&mut self, _error: &E) -> Option<Duration> {
2135 Some(self.retry_delay.next_delay(&mut self.mockable.thread_rng()))
2136 }
2137}
2138
2139impl RetriableError for UploadError {
2140 fn should_retry(&self) -> bool {
2141 match self {
2142 UploadError::Request(_) | UploadError::Circuit(_) | UploadError::Stream(_) => true,
2143 UploadError::Bug(_) => false,
2144 }
2145 }
2146}
2147
2148/// The outcome of uploading a descriptor to the HSDirs from a particular time period.
2149#[derive(Debug, Clone)]
2150struct TimePeriodUploadResult {
2151 /// The time period.
2152 time_period: TimePeriod,
2153 /// The upload results.
2154 hsdir_result: Vec<HsDirUploadStatus>,
2155}
2156
2157/// The outcome of uploading a descriptor to a particular HsDir.
2158#[derive(Clone, Debug)]
2159struct HsDirUploadStatus {
2160 /// The identity of the HsDir we attempted to upload the descriptor to.
2161 relay_ids: RelayIds,
2162 /// The outcome of this attempt.
2163 upload_res: UploadResult,
2164 /// The revision counter of the descriptor we tried to upload.
2165 revision_counter: RevisionCounter,
2166}
2167
2168/// The outcome of uploading a descriptor.
2169type UploadResult = Result<(), DescUploadRetryError>;
2170
2171impl From<BackoffError<UploadError>> for DescUploadRetryError {
2172 fn from(e: BackoffError<UploadError>) -> Self {
2173 use BackoffError as BE;
2174 use DescUploadRetryError as DURE;
2175
2176 match e {
2177 BE::FatalError(e) => DURE::FatalError(e),
2178 BE::MaxRetryCountExceeded(e) => DURE::MaxRetryCountExceeded(e),
2179 BE::Timeout(e) => DURE::Timeout(e),
2180 BE::ExplicitStop(_) => {
2181 DURE::Bug(internal!("explicit stop in publisher backoff schedule?!"))
2182 }
2183 }
2184 }
2185}
2186
2187// NOTE: the rest of the publisher tests live in publish.rs
2188#[cfg(test)]
2189mod test {
2190 // @@ begin test lint list maintained by maint/add_warning @@
2191 #![allow(clippy::bool_assert_comparison)]
2192 #![allow(clippy::clone_on_copy)]
2193 #![allow(clippy::dbg_macro)]
2194 #![allow(clippy::mixed_attributes_style)]
2195 #![allow(clippy::print_stderr)]
2196 #![allow(clippy::print_stdout)]
2197 #![allow(clippy::single_char_pattern)]
2198 #![allow(clippy::unwrap_used)]
2199 #![allow(clippy::unchecked_time_subtraction)]
2200 #![allow(clippy::useless_vec)]
2201 #![allow(clippy::needless_pass_by_value)]
2202 #![allow(clippy::string_slice)] // See arti#2571
2203 //! <!-- @@ end test lint list maintained by maint/add_warning @@ -->
2204 use super::*;
2205 use tor_netdir::testnet;
2206
2207 /// Create a `TimePeriodContext` from the specified upload results.
2208 fn create_time_period_ctx(
2209 params: &HsDirParams,
2210 upload_results: Vec<HsDirUploadStatus>,
2211 ) -> TimePeriodContext {
2212 TimePeriodContext {
2213 params: params.clone(),
2214 hs_dirs: vec![],
2215 last_successful: None,
2216 upload_results,
2217 }
2218 }
2219
2220 /// Create a single `HsDirUploadStatus`
2221 fn create_upload_status(upload_res: UploadResult) -> HsDirUploadStatus {
2222 HsDirUploadStatus {
2223 relay_ids: RelayIds::empty(),
2224 upload_res,
2225 revision_counter: RevisionCounter::from(13),
2226 }
2227 }
2228
2229 /// Create a bunch of results, all with the specified `upload_res`.
2230 fn create_upload_results(upload_res: UploadResult) -> Vec<HsDirUploadStatus> {
2231 std::iter::repeat_with(|| create_upload_status(upload_res.clone()))
2232 .take(10)
2233 .collect()
2234 }
2235
2236 fn construct_netdir() -> NetDir {
2237 const SRV1: [u8; 32] = *b"The door refused to open. ";
2238 const SRV2: [u8; 32] = *b"It said, 'Five cents, please.' ";
2239
2240 let dir = testnet::construct_custom_netdir(|_, _, bld| {
2241 bld.shared_rand_cur(7, SRV1.into(), None)
2242 .shared_rand_prev(7, SRV2.into(), None);
2243 })
2244 .unwrap();
2245
2246 dir.unwrap_if_sufficient().unwrap()
2247 }
2248
2249 #[test]
2250 fn upload_result_status_bootstrapping() {
2251 let netdir = construct_netdir();
2252 let all_params = netdir.hs_all_time_periods();
2253 let current_period = netdir.hs_time_period();
2254 let primary_params = all_params
2255 .iter()
2256 .find(|param| param.time_period() == current_period)
2257 .unwrap();
2258 let results = [
2259 (vec![], vec![]),
2260 (vec![], create_upload_results(Ok(()))),
2261 (create_upload_results(Ok(())), vec![]),
2262 ];
2263
2264 for (primary_result, secondary_result) in results {
2265 let primary_ctx = create_time_period_ctx(primary_params, primary_result);
2266
2267 let secondary_params = all_params
2268 .iter()
2269 .find(|param| param.time_period() != current_period)
2270 .unwrap();
2271 let secondary_ctx = create_time_period_ctx(secondary_params, secondary_result.clone());
2272
2273 let (status, err) = upload_result_state(&netdir, &[primary_ctx, secondary_ctx]);
2274 assert_eq!(status, State::Bootstrapping);
2275 assert!(err.is_none());
2276 }
2277 }
2278
2279 #[test]
2280 fn upload_result_status_running() {
2281 let netdir = construct_netdir();
2282 let all_params = netdir.hs_all_time_periods();
2283 let current_period = netdir.hs_time_period();
2284 let primary_params = all_params
2285 .iter()
2286 .find(|param| param.time_period() == current_period)
2287 .unwrap();
2288
2289 let secondary_result = create_upload_results(Ok(()));
2290 let secondary_params = all_params
2291 .iter()
2292 .find(|param| param.time_period() != current_period)
2293 .unwrap();
2294 let secondary_ctx = create_time_period_ctx(secondary_params, secondary_result.clone());
2295
2296 let primary_result = create_upload_results(Ok(()));
2297 let primary_ctx = create_time_period_ctx(primary_params, primary_result);
2298 let (status, err) = upload_result_state(&netdir, &[primary_ctx, secondary_ctx]);
2299 assert_eq!(status, State::Running);
2300 assert!(err.is_none());
2301 }
2302
2303 #[test]
2304 fn upload_result_status_reachable() {
2305 let netdir = construct_netdir();
2306 let all_params = netdir.hs_all_time_periods();
2307 let current_period = netdir.hs_time_period();
2308 let primary_params = all_params
2309 .iter()
2310 .find(|param| param.time_period() == current_period)
2311 .unwrap();
2312
2313 let primary_result = create_upload_results(Ok(()));
2314 let primary_ctx = create_time_period_ctx(primary_params, primary_result.clone());
2315 let failed_res = create_upload_results(Err(DescUploadRetryError::Bug(internal!("test"))));
2316 let secondary_result = create_upload_results(Ok(()))
2317 .into_iter()
2318 .chain(failed_res.iter().cloned())
2319 .collect();
2320 let secondary_params = all_params
2321 .iter()
2322 .find(|param| param.time_period() != current_period)
2323 .unwrap();
2324 let secondary_ctx = create_time_period_ctx(secondary_params, secondary_result);
2325 let (status, err) = upload_result_state(&netdir, &[primary_ctx, secondary_ctx]);
2326
2327 // Degraded but reachable (because some of the secondary HsDir uploads failed).
2328 assert_eq!(status, State::DegradedReachable);
2329 assert!(matches!(err, Some(Problem::DescriptorUpload(_))));
2330 }
2331
2332 #[test]
2333 fn upload_result_status_unreachable() {
2334 let netdir = construct_netdir();
2335 let all_params = netdir.hs_all_time_periods();
2336 let current_period = netdir.hs_time_period();
2337 let primary_params = all_params
2338 .iter()
2339 .find(|param| param.time_period() == current_period)
2340 .unwrap();
2341 let mut primary_result =
2342 create_upload_results(Err(DescUploadRetryError::Bug(internal!("test"))));
2343 let primary_ctx = create_time_period_ctx(primary_params, primary_result.clone());
2344 // No secondary TP (we are unreachable).
2345 let (status, err) = upload_result_state(&netdir, &[primary_ctx]);
2346 assert_eq!(status, State::DegradedUnreachable);
2347 assert!(matches!(err, Some(Problem::DescriptorUpload(_))));
2348
2349 // Add a successful result
2350 primary_result.push(create_upload_status(Ok(())));
2351 let primary_ctx = create_time_period_ctx(primary_params, primary_result.clone());
2352 let (status, err) = upload_result_state(&netdir, &[primary_ctx]);
2353 // Still degraded, and unreachable (because we don't have a TimePeriodContext
2354 // for the secondary TP)
2355 assert_eq!(status, State::DegradedUnreachable);
2356 assert!(matches!(err, Some(Problem::DescriptorUpload(_))));
2357
2358 // If we add another time period where none of the uploads were successful,
2359 // we're *still* unreachable
2360 let secondary_result =
2361 create_upload_results(Err(DescUploadRetryError::Bug(internal!("test"))));
2362 let secondary_params = all_params
2363 .iter()
2364 .find(|param| param.time_period() != current_period)
2365 .unwrap();
2366 let secondary_ctx = create_time_period_ctx(secondary_params, secondary_result.clone());
2367 let primary_ctx = create_time_period_ctx(primary_params, primary_result.clone());
2368 let (status, err) = upload_result_state(&netdir, &[primary_ctx, secondary_ctx]);
2369 assert_eq!(status, State::DegradedUnreachable);
2370 assert!(matches!(err, Some(Problem::DescriptorUpload(_))));
2371 }
2372}