1#![cfg_attr(docsrs, feature(doc_cfg))]
2#![doc = include_str!("../README.md")]
3#![allow(renamed_and_removed_lints)] #![allow(unknown_lints)] #![warn(missing_docs)]
7#![warn(noop_method_call)]
8#![warn(unreachable_pub)]
9#![warn(clippy::all)]
10#![deny(clippy::await_holding_lock)]
11#![deny(clippy::cargo_common_metadata)]
12#![deny(clippy::cast_lossless)]
13#![deny(clippy::checked_conversions)]
14#![allow(clippy::cognitive_complexity)] #![deny(clippy::debug_assert_with_mut_call)]
16#![deny(clippy::exhaustive_enums)]
17#![deny(clippy::exhaustive_structs)]
18#![deny(clippy::expl_impl_clone_on_copy)]
19#![deny(clippy::fallible_impl_from)]
20#![deny(clippy::implicit_clone)]
21#![deny(clippy::large_stack_arrays)]
22#![warn(clippy::manual_ok_or)]
23#![deny(clippy::missing_docs_in_private_items)]
24#![warn(clippy::needless_borrow)]
25#![warn(clippy::needless_pass_by_value)]
26#![warn(clippy::option_option)]
27#![deny(clippy::print_stderr)]
28#![deny(clippy::print_stdout)]
29#![warn(clippy::rc_buffer)]
30#![deny(clippy::ref_option_ref)]
31#![warn(clippy::semicolon_if_nothing_returned)]
32#![warn(clippy::trait_duplication_in_bounds)]
33#![deny(clippy::unchecked_time_subtraction)]
34#![deny(clippy::unnecessary_wraps)]
35#![warn(clippy::unseparated_literal_suffix)]
36#![deny(clippy::unwrap_used)]
37#![deny(clippy::mod_module_files)]
38#![allow(clippy::let_unit_value)] #![allow(clippy::uninlined_format_args)]
40#![allow(clippy::significant_drop_in_scrutinee)] #![allow(clippy::result_large_err)] #![allow(clippy::needless_raw_string_hashes)] #![allow(clippy::needless_lifetimes)] #![allow(mismatched_lifetime_syntaxes)] #![allow(clippy::collapsible_if)] #![deny(clippy::unused_async)]
47#![deny(clippy::string_slice)] #![allow(clippy::single_component_path_imports)]
54
55mod bootstrap;
56pub mod config;
57mod docid;
58mod docmeta;
59mod err;
60mod event;
61mod shared_ref;
62mod state;
63mod storage;
64
65#[cfg(feature = "dir-plugin")]
66mod as_plugin;
67#[cfg(feature = "bridge-client")]
68pub mod bridgedesc;
69#[cfg(feature = "dirfilter")]
70pub mod filter;
71
72use crate::docid::{CacheUsage, ClientRequest, DocQuery};
73use crate::err::BootstrapAction;
74#[cfg(not(feature = "experimental-api"))]
75use crate::shared_ref::SharedMutArc;
76#[cfg(feature = "experimental-api")]
77pub use crate::shared_ref::SharedMutArc;
78use crate::storage::{DynStore, Store};
79use bootstrap::AttemptId;
80use event::DirProgress;
81use postage::watch;
82use scopeguard::ScopeGuard;
83use tor_circmgr::CircMgr;
84use tor_dirclient::SourceInfo;
85use tor_dircommon::config::DirTolerance;
86use tor_error::{info_report, into_internal, warn_report};
87use tor_netdir::params::NetParameters;
88use tor_netdir::{DirEvent, MdReceiver, NetDir, NetDirProvider};
89
90use async_trait::async_trait;
91use futures::stream::BoxStream;
92use oneshot_fused_workaround as oneshot;
93use tor_netdoc::doc::netstatus::ProtoStatuses;
94use tor_rtcompat::scheduler::{TaskHandle, TaskSchedule};
95use tor_rtcompat::{Runtime, SpawnExt};
96use tracing::{debug, info, instrument, trace, warn};
97use web_time_compat::SystemTimeExt;
98
99use std::marker::PhantomData;
100use std::sync::atomic::{AtomicBool, Ordering};
101use std::sync::{Arc, Mutex};
102use std::time::Duration;
103use std::{collections::HashMap, sync::Weak};
104use std::{fmt::Debug, time::SystemTime};
105
106use crate::state::{DirState, NetDirChange};
107pub use config::DirMgrConfig;
108pub use docid::DocId;
109pub use err::Error;
110pub use event::{DirBlockage, DirBootstrapEvents, DirBootstrapStatus};
111pub use storage::DocumentText;
112pub use tor_dircommon::fallback::{FallbackDir, FallbackDirBuilder};
113pub use tor_netdir::Timeliness;
114
115#[cfg(feature = "dir-plugin")]
116pub use as_plugin::DirPlugin;
117
118use strum;
120
121pub type Result<T> = std::result::Result<T, Error>;
123
124#[derive(Clone)]
131pub struct DirMgrStore<R: Runtime> {
132 pub(crate) store: Arc<Mutex<crate::DynStore>>,
134
135 pub(crate) runtime: PhantomData<R>,
137}
138
139impl<R: Runtime> DirMgrStore<R> {
140 pub fn new(config: &DirMgrConfig, runtime: R, offline: bool) -> Result<Self> {
142 let store = Arc::new(Mutex::new(config.open_store(offline)?));
143 drop(runtime);
144 let runtime = PhantomData;
145 Ok(DirMgrStore { store, runtime })
146 }
147}
148
149#[async_trait]
151pub trait DirProvider: NetDirProvider {
152 fn reconfigure(
156 &self,
157 new_config: &DirMgrConfig,
158 how: tor_config::Reconfigure,
159 ) -> std::result::Result<(), tor_config::ReconfigureError>;
160
161 async fn bootstrap(&self) -> Result<()>;
163
164 fn bootstrap_events(&self) -> BoxStream<'static, DirBootstrapStatus>;
170
171 fn download_task_handle(&self) -> Option<TaskHandle> {
173 None
174 }
175}
176
177impl<R: Runtime> NetDirProvider for DirMgr<R> {
180 fn netdir(&self, timeliness: Timeliness) -> tor_netdir::Result<Arc<NetDir>> {
181 use tor_netdir::Error as NetDirError;
182 let netdir = self.netdir.get().ok_or(NetDirError::NoInfo)?;
183 let lifetime = match timeliness {
184 Timeliness::Strict => netdir.lifetime().clone(),
185 Timeliness::Timely => self
186 .config
187 .get()
188 .tolerance
189 .extend_lifetime(netdir.lifetime()),
190 Timeliness::Unchecked => return Ok(netdir),
191 };
192 let now = SystemTime::get();
194 if lifetime.valid_after() > now {
195 Err(NetDirError::DirNotYetValid)
196 } else if lifetime.valid_until() < now {
197 Err(NetDirError::DirExpired)
198 } else {
199 Ok(netdir)
200 }
201 }
202
203 fn events(&self) -> BoxStream<'static, DirEvent> {
204 Box::pin(self.events.subscribe())
205 }
206
207 fn params(&self) -> Arc<dyn AsRef<tor_netdir::params::NetParameters>> {
208 if let Some(netdir) = self.netdir.get() {
209 netdir
215 } else {
216 self.default_parameters
220 .lock()
221 .expect("Poisoned lock")
222 .clone()
223 }
224 }
230
231 fn protocol_statuses(&self) -> Option<(SystemTime, Arc<ProtoStatuses>)> {
232 self.protocols.lock().expect("Poisoned lock").clone()
233 }
234}
235
236#[async_trait]
237impl<R: Runtime> DirProvider for Arc<DirMgr<R>> {
238 fn reconfigure(
239 &self,
240 new_config: &DirMgrConfig,
241 how: tor_config::Reconfigure,
242 ) -> std::result::Result<(), tor_config::ReconfigureError> {
243 DirMgr::reconfigure(self, new_config, how)
244 }
245
246 #[instrument(level = "trace", skip_all)]
247 async fn bootstrap(&self) -> Result<()> {
248 DirMgr::bootstrap(self).await
249 }
250
251 fn bootstrap_events(&self) -> BoxStream<'static, DirBootstrapStatus> {
252 Box::pin(DirMgr::bootstrap_events(self))
253 }
254
255 fn download_task_handle(&self) -> Option<TaskHandle> {
256 Some(self.task_handle.clone())
257 }
258}
259
260pub struct DirMgr<R: Runtime> {
274 config: tor_config::MutCfg<DirMgrConfig>,
277 store: Arc<Mutex<DynStore>>,
282 netdir: Arc<SharedMutArc<NetDir>>,
289
290 protocols: Mutex<Option<(SystemTime, Arc<ProtoStatuses>)>>,
292
293 default_parameters: Mutex<Arc<NetParameters>>,
295
296 events: event::FlagPublisher<DirEvent>,
298
299 send_status: Mutex<watch::Sender<event::DirBootstrapStatus>>,
302
303 receive_status: DirBootstrapEvents,
309
310 circmgr: Option<Arc<CircMgr<R>>>,
312
313 runtime: R,
315
316 offline: bool,
318
319 bootstrap_started: AtomicBool,
326
327 #[cfg(feature = "dirfilter")]
329 filter: crate::filter::FilterConfig,
330
331 task_schedule: Mutex<Option<TaskSchedule<R>>>,
334
335 task_handle: TaskHandle,
337}
338
339#[derive(Debug, Clone)]
344#[non_exhaustive]
345pub enum DocSource {
346 LocalCache,
348 DirServer {
350 source: Option<SourceInfo>,
352 },
353}
354
355impl std::fmt::Display for DocSource {
356 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
357 match self {
358 DocSource::LocalCache => write!(f, "local cache"),
359 DocSource::DirServer { source: None } => write!(f, "directory server"),
360 DocSource::DirServer { source: Some(info) } => write!(f, "directory server {}", info),
361 }
362 }
363}
364
365impl<R: Runtime> DirMgr<R> {
366 pub fn load_once(runtime: R, config: DirMgrConfig) -> Result<Arc<NetDir>> {
376 let store = DirMgrStore::new(&config, runtime.clone(), true)?;
377 let dirmgr = Arc::new(Self::from_config(config, runtime, store, None, true)?);
378
379 let attempt = AttemptId::next();
381 trace!(%attempt, "Trying to load a full directory from cache");
382 let outcome = dirmgr.load_directory(attempt);
383 trace!(%attempt, "Load result: {outcome:?}");
384 let _success = outcome?;
385
386 dirmgr
387 .netdir(Timeliness::Timely)
388 .map_err(|_| Error::DirectoryNotPresent)
389 }
390
391 pub async fn load_or_bootstrap_once(
401 config: DirMgrConfig,
402 runtime: R,
403 store: DirMgrStore<R>,
404 circmgr: Arc<CircMgr<R>>,
405 ) -> Result<Arc<NetDir>> {
406 let dirmgr = DirMgr::bootstrap_from_config(config, runtime, store, circmgr).await?;
407 dirmgr
408 .timely_netdir()
409 .map_err(|_| Error::DirectoryNotPresent)
410 }
411
412 pub fn create_unbootstrapped(
416 config: DirMgrConfig,
417 runtime: R,
418 store: DirMgrStore<R>,
419 circmgr: Arc<CircMgr<R>>,
420 ) -> Result<Arc<Self>> {
421 Ok(Arc::new(DirMgr::from_config(
422 config,
423 runtime,
424 store,
425 Some(circmgr),
426 false,
427 )?))
428 }
429
430 #[instrument(level = "trace", skip_all)]
451 pub async fn bootstrap(self: &Arc<Self>) -> Result<()> {
452 if self.offline {
453 return Err(Error::OfflineMode);
454 }
455
456 if self
463 .bootstrap_started
464 .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
465 .is_err()
466 {
467 debug!("Attempted to bootstrap twice; ignoring.");
468 return Ok(());
469 }
470
471 let reset_bootstrap_started = scopeguard::guard(&self.bootstrap_started, |v| {
474 v.store(false, Ordering::SeqCst);
475 });
476
477 let schedule = {
478 let sched = self.task_schedule.lock().expect("poisoned lock").take();
479 match sched {
480 Some(sched) => sched,
481 None => {
482 debug!("Attempted to bootstrap twice; ignoring.");
483 return Ok(());
484 }
485 }
486 };
487
488 let attempt_id = AttemptId::next();
490 trace!(attempt=%attempt_id, "Starting to bootstrap directory");
491 let have_directory = self.load_directory(attempt_id)?;
492
493 let (mut sender, receiver) = if have_directory {
494 info!("Loaded a good directory from cache.");
495 (None, None)
496 } else {
497 info!("Didn't get usable directory from cache.");
498 let (sender, receiver) = oneshot::channel();
499 (Some(sender), Some(receiver))
500 };
501
502 let dirmgr_weak = Arc::downgrade(self);
504 self.runtime
505 .spawn(async move {
506 let mut schedule = scopeguard::guard(schedule, |schedule| {
513 if let Some(dm) = Weak::upgrade(&dirmgr_weak) {
514 *dm.task_schedule.lock().expect("poisoned lock") = Some(schedule);
515 }
516 });
517
518 if let Err(e) =
521 Self::reload_until_owner(&dirmgr_weak, &mut schedule, attempt_id, &mut sender)
522 .await
523 {
524 match e {
525 Error::ManagerDropped => {}
526 _ => warn_report!(e, "Unrecovered error while waiting for bootstrap",),
527 }
528 } else if let Err(e) =
529 Self::download_forever(dirmgr_weak.clone(), &mut schedule, attempt_id, sender)
530 .await
531 {
532 match e {
533 Error::ManagerDropped => {}
534 _ => warn_report!(e, "Unrecovered error while downloading"),
535 }
536 }
537 })
538 .map_err(|e| Error::from_spawn("directory updater task", e))?;
539
540 if let Some(receiver) = receiver {
541 match receiver.await {
542 Ok(()) => {
543 info!("We have enough information to build circuits.");
544 let _ = ScopeGuard::into_inner(reset_bootstrap_started);
546 }
547 Err(_) => {
548 warn!("Bootstrapping task exited before finishing.");
549 return Err(Error::CantAdvanceState);
550 }
551 }
552 }
553 Ok(())
554 }
555
556 pub fn bootstrap_started(&self) -> bool {
558 self.bootstrap_started.load(Ordering::SeqCst)
559 }
560
561 #[instrument(level = "trace", skip_all)]
564 pub async fn bootstrap_from_config(
565 config: DirMgrConfig,
566 runtime: R,
567 store: DirMgrStore<R>,
568 circmgr: Arc<CircMgr<R>>,
569 ) -> Result<Arc<Self>> {
570 let dirmgr = Self::create_unbootstrapped(config, runtime, store, circmgr)?;
571
572 dirmgr.bootstrap().await?;
573
574 Ok(dirmgr)
575 }
576
577 async fn reload_until_owner(
585 weak: &Weak<Self>,
586 schedule: &mut TaskSchedule<R>,
587 attempt_id: AttemptId,
588 on_complete: &mut Option<oneshot::Sender<()>>,
589 ) -> Result<()> {
590 let mut logged = false;
591 let mut bootstrapped;
592 {
593 let dirmgr = upgrade_weak_ref(weak)?;
594 bootstrapped = dirmgr.netdir.get().is_some();
595 }
596
597 loop {
598 {
599 let dirmgr = upgrade_weak_ref(weak)?;
600 trace!("Trying to take ownership of the directory cache lock");
601 if dirmgr.try_upgrade_to_readwrite()? {
602 if logged {
606 info!(
607 "The previous owning process has given up the lock. We are now in charge of managing the directory."
608 );
609 }
610 return Ok(());
611 }
612 }
613
614 if !logged {
615 logged = true;
616 if bootstrapped {
617 info!("Another process is managing the directory. We'll use its cache.");
618 } else {
619 info!(
620 "Another process is bootstrapping the directory. Waiting till it finishes or exits."
621 );
622 }
623 }
624
625 let pause = if bootstrapped {
628 std::time::Duration::new(120, 0)
629 } else {
630 std::time::Duration::new(5, 0)
631 };
632 schedule.sleep(pause).await?;
633 {
637 let dirmgr = upgrade_weak_ref(weak)?;
638 trace!("Trying to load from the directory cache");
639 if dirmgr.load_directory(attempt_id)? {
640 if let Some(send_done) = on_complete.take() {
642 let _ = send_done.send(());
643 }
644 if !bootstrapped {
645 info!("The directory is now bootstrapped.");
646 }
647 bootstrapped = true;
648 }
649 }
650 }
651 }
652
653 #[instrument(level = "trace", skip_all)]
658 async fn download_forever(
659 weak: Weak<Self>,
660 schedule: &mut TaskSchedule<R>,
661 mut attempt_id: AttemptId,
662 mut on_complete: Option<oneshot::Sender<()>>,
663 ) -> Result<()> {
664 let mut state: Box<dyn DirState> = {
665 let dirmgr = upgrade_weak_ref(&weak)?;
666 Box::new(state::GetConsensusState::new(
667 dirmgr.runtime.clone(),
668 dirmgr.config.get(),
669 CacheUsage::CacheOkay,
670 Some(dirmgr.netdir.clone()),
671 #[cfg(feature = "dirfilter")]
672 dirmgr
673 .filter
674 .clone()
675 .unwrap_or_else(|| Arc::new(crate::filter::NilFilter)),
676 ))
677 };
678
679 trace!("Entering download loop.");
680
681 loop {
682 let mut usable = false;
683
684 let retry_config = {
685 let dirmgr = upgrade_weak_ref(&weak)?;
686 dirmgr.config.get().schedule.retry_bootstrap()
690 };
691 let mut retry_delay = retry_config.schedule();
692
693 'retry_attempt: for try_num in retry_config.attempts() {
694 trace!(attempt=%attempt_id, ?try_num, "Trying to download a directory.");
695 let outcome = bootstrap::download(
696 Weak::clone(&weak),
697 &mut state,
698 schedule,
699 attempt_id,
700 &mut on_complete,
701 )
702 .await;
703 trace!(attempt=%attempt_id, ?try_num, ?outcome, "Download is over.");
704
705 if let Err(err) = outcome {
706 if state.is_ready(Readiness::Usable) {
707 usable = true;
708 info_report!(
709 err,
710 "Unable to completely download a directory. (Nevertheless, the directory is usable, so we'll pause for now)"
711 );
712 break 'retry_attempt;
713 }
714
715 match err.bootstrap_action() {
716 BootstrapAction::Nonfatal => {
717 return Err(into_internal!(
718 "Nonfatal error should not have propagated here"
719 )(err)
720 .into());
721 }
722 BootstrapAction::Reset => {}
723 BootstrapAction::Fatal => return Err(err),
724 }
725
726 let delay = retry_delay.next_delay(&mut rand::rng());
727 warn_report!(
728 err,
729 "Unable to download a usable directory. (We will restart in {})",
730 humantime::format_duration(delay),
731 );
732 {
733 let dirmgr = upgrade_weak_ref(&weak)?;
734 dirmgr.note_reset(attempt_id);
735 }
736 schedule.sleep(delay).await?;
737 state = state.reset();
738 } else {
739 info!(attempt=%attempt_id, "Directory is complete.");
740 usable = true;
741 break 'retry_attempt;
742 }
743 }
744
745 if !usable {
746 warn!(
748 "We failed {} times to bootstrap a directory. We're going to give up.",
749 retry_config.n_attempts()
750 );
751 return Err(Error::CantAdvanceState);
752 } else {
753 if let Some(send_done) = on_complete.take() {
755 let _ = send_done.send(());
756 }
757 }
758
759 let reset_at = state.reset_time();
760 match reset_at {
761 Some(t) => {
762 trace!("Sleeping until {}", time::OffsetDateTime::from(t));
763 schedule.sleep_until_wallclock(t).await?;
764 }
765 None => return Ok(()),
766 }
767 attempt_id = bootstrap::AttemptId::next();
768 trace!(attempt=%attempt_id, "Beginning new attempt to bootstrap directory");
769 state = state.reset();
770 }
771 }
772
773 fn circmgr(&self) -> Result<Arc<CircMgr<R>>> {
775 self.circmgr.clone().ok_or(Error::NoDownloadSupport)
776 }
777
778 pub fn reconfigure(
782 &self,
783 new_config: &DirMgrConfig,
784 how: tor_config::Reconfigure,
785 ) -> std::result::Result<(), tor_config::ReconfigureError> {
786 let config = self.config.get();
787 if new_config.cache_dir != config.cache_dir {
792 how.cannot_change("storage.cache_dir")?;
793 }
794 if new_config.cache_trust != config.cache_trust {
795 how.cannot_change("storage.permissions")?;
796 }
797 if new_config.authorities() != config.authorities() {
798 how.cannot_change("network.authorities")?;
799 }
800
801 if how == tor_config::Reconfigure::CheckAllOrNothing {
802 return Ok(());
803 }
804
805 let params_changed = new_config.override_net_params != config.override_net_params;
806
807 self.config
808 .map_and_replace(|cfg| cfg.update_from_config(new_config));
809
810 if params_changed {
811 let _ignore_err = self.netdir.mutate(|netdir| {
812 netdir.replace_overridden_parameters(&new_config.override_net_params);
813 Ok(())
814 });
815 {
816 let mut params = self.default_parameters.lock().expect("lock failed");
817 *params = Arc::new(NetParameters::from_map(&new_config.override_net_params));
818 }
819
820 self.events.publish(DirEvent::NewConsensus);
823 }
824
825 Ok(())
826 }
827
828 pub fn bootstrap_events(&self) -> event::DirBootstrapEvents {
834 self.receive_status.clone()
835 }
836
837 fn update_progress(&self, attempt_id: AttemptId, progress: DirProgress) {
840 let mut sender = self.send_status.lock().expect("poisoned lock");
842 let mut status = sender.borrow_mut();
843
844 status.update_progress(attempt_id, progress);
845 }
846
847 fn note_errors(&self, attempt_id: AttemptId, n_errors: usize) {
850 if n_errors == 0 {
851 return;
852 }
853 let mut sender = self.send_status.lock().expect("poisoned lock");
854 let mut status = sender.borrow_mut();
855
856 status.note_errors(attempt_id, n_errors);
857 }
858
859 fn note_reset(&self, attempt_id: AttemptId) {
861 let mut sender = self.send_status.lock().expect("poisoned lock");
862 let mut status = sender.borrow_mut();
863
864 status.note_reset(attempt_id);
865 }
866
867 fn try_upgrade_to_readwrite(&self) -> Result<bool> {
874 self.store
875 .lock()
876 .expect("Directory storage lock poisoned")
877 .upgrade_to_readwrite()
878 }
879
880 #[cfg(test)]
882 fn store_if_rw(&self) -> Option<&Mutex<DynStore>> {
883 let rw = !self
884 .store
885 .lock()
886 .expect("Directory storage lock poisoned")
887 .is_readonly();
888 if rw { Some(&self.store) } else { None }
890 }
891
892 #[allow(clippy::unnecessary_wraps)] fn from_config(
898 config: DirMgrConfig,
899 runtime: R,
900 store: DirMgrStore<R>,
901 circmgr: Option<Arc<CircMgr<R>>>,
902 offline: bool,
903 ) -> Result<Self> {
904 let netdir = Arc::new(SharedMutArc::new());
905 let events = event::FlagPublisher::new();
906 let default_parameters = NetParameters::from_map(&config.override_net_params);
907 let default_parameters = Mutex::new(Arc::new(default_parameters));
908
909 let (send_status, receive_status) = postage::watch::channel();
910 let send_status = Mutex::new(send_status);
911 let receive_status = DirBootstrapEvents {
912 inner: receive_status,
913 };
914 #[cfg(feature = "dirfilter")]
915 let filter = config.extensions.filter.clone();
916
917 let (task_schedule, task_handle) = TaskSchedule::new(runtime.clone());
919 let task_schedule = Mutex::new(Some(task_schedule));
920
921 let protocols = {
924 let store = store.store.lock().expect("lock poisoned");
925 store
926 .cached_protocol_recommendations()?
927 .map(|(t, p)| (t, Arc::new(p)))
928 };
929
930 Ok(DirMgr {
931 config: config.into(),
932 store: store.store,
933 netdir,
934 protocols: Mutex::new(protocols),
935 default_parameters,
936 events,
937 send_status,
938 receive_status,
939 circmgr,
940 runtime,
941 offline,
942 bootstrap_started: AtomicBool::new(false),
943 #[cfg(feature = "dirfilter")]
944 filter,
945 task_schedule,
946 task_handle,
947 })
948 }
949
950 fn load_directory(self: &Arc<Self>, attempt_id: AttemptId) -> Result<bool> {
955 let state = state::GetConsensusState::new(
956 self.runtime.clone(),
957 self.config.get(),
958 CacheUsage::CacheOnly,
959 None,
960 #[cfg(feature = "dirfilter")]
961 self.filter
962 .clone()
963 .unwrap_or_else(|| Arc::new(crate::filter::NilFilter)),
964 );
965 let _ = bootstrap::load(self, Box::new(state), attempt_id)?;
966
967 Ok(self.netdir.get().is_some())
968 }
969
970 pub fn events(&self) -> impl futures::Stream<Item = DirEvent> + use<R> {
977 self.events.subscribe()
978 }
979
980 pub fn text(&self, doc: &DocId) -> Result<Option<DocumentText>> {
983 use itertools::Itertools;
984 let mut result = HashMap::new();
985 let query: DocQuery = (*doc).into();
986 let store = self.store.lock().expect("store lock poisoned");
987 query.load_from_store_into(&mut result, &**store)?;
988 let item = result.into_iter().at_most_one().map_err(|_| {
989 Error::CacheCorruption("Found more than one entry in storage for given docid")
990 })?;
991 if let Some((docid, doctext)) = item {
992 if &docid != doc {
993 return Err(Error::CacheCorruption(
994 "Item from storage had incorrect docid.",
995 ));
996 }
997 Ok(Some(doctext))
998 } else {
999 Ok(None)
1000 }
1001 }
1002
1003 pub fn texts<T>(&self, docs: T) -> Result<HashMap<DocId, DocumentText>>
1008 where
1009 T: IntoIterator<Item = DocId>,
1010 {
1011 let partitioned = docid::partition_by_type(docs);
1012 let mut result = HashMap::new();
1013 let store = self.store.lock().expect("store lock poisoned");
1014 for (_, query) in partitioned.into_iter() {
1015 query.load_from_store_into(&mut result, &**store)?;
1016 }
1017 Ok(result)
1018 }
1019
1020 fn expand_response_text(&self, req: &ClientRequest, text: String) -> Result<String> {
1028 if let ClientRequest::Consensus(req) = req {
1029 if tor_consdiff::looks_like_diff(&text) {
1030 if let Some(old_d) = req.old_consensus_digests().next() {
1031 let db_val = {
1032 let s = self.store.lock().expect("Directory storage lock poisoned");
1033 s.consensus_by_sha3_digest_of_signed_part(old_d)?
1034 };
1035 if let Some((old_consensus, meta)) = db_val {
1036 info!("Applying a consensus diff");
1037 let new_consensus = tor_consdiff::apply_diff(
1038 old_consensus.as_str()?,
1039 &text,
1040 Some(*meta.sha3_256_of_signed()),
1041 tor_consdiff::DiffSizeStrictness::Apply,
1042 )?;
1043 new_consensus.check_digest()?;
1044 return Ok(new_consensus.to_string());
1045 }
1046 }
1047 return Err(Error::Unwanted(
1048 "Received a consensus diff we did not ask for",
1049 ));
1050 }
1051 }
1052 Ok(text)
1053 }
1054
1055 fn apply_netdir_changes(
1057 self: &Arc<Self>,
1058 state: &mut Box<dyn DirState>,
1059 store: &mut dyn Store,
1060 ) -> Result<()> {
1061 if let Some(change) = state.get_netdir_change() {
1062 match change {
1063 NetDirChange::AttemptReplace {
1064 netdir,
1065 consensus_meta,
1066 } => {
1067 if let Some(ref cm) = self.circmgr {
1070 if !cm
1071 .netdir_is_sufficient(netdir.as_ref().expect("AttemptReplace had None"))
1072 {
1073 debug!("Got a new NetDir, but it doesn't have enough guards yet.");
1074 return Ok(());
1075 }
1076 }
1077 let is_stale = {
1078 self.netdir
1080 .get()
1081 .map(|x| {
1082 x.lifetime().valid_after()
1083 > netdir
1084 .as_ref()
1085 .expect("AttemptReplace had None")
1086 .lifetime()
1087 .valid_after()
1088 })
1089 .unwrap_or(false)
1090 };
1091 if is_stale {
1092 warn!("Got a new NetDir, but it's older than the one we currently have!");
1093 return Err(Error::NetDirOlder);
1094 }
1095 let cfg = self.config.get();
1096 let mut netdir = netdir.take().expect("AttemptReplace had None");
1097 netdir.replace_overridden_parameters(&cfg.override_net_params);
1098 self.netdir.replace(netdir);
1099 self.events.publish(DirEvent::NewConsensus);
1100 self.events.publish(DirEvent::NewDescriptors);
1101
1102 info!("Marked consensus usable.");
1103 if !store.is_readonly() {
1104 store.mark_consensus_usable(consensus_meta)?;
1105 store.expire_all(&crate::storage::EXPIRATION_DEFAULTS)?;
1108 }
1109 Ok(())
1110 }
1111 NetDirChange::AddMicrodescs(mds) => {
1112 self.netdir.mutate(|netdir| {
1113 for md in mds.drain(..) {
1114 netdir.add_microdesc(md);
1115 }
1116 Ok(())
1117 })?;
1118 self.events.publish(DirEvent::NewDescriptors);
1119 Ok(())
1120 }
1121 NetDirChange::SetRequiredProtocol { timestamp, protos } => {
1122 if !store.is_readonly() {
1123 store.update_protocol_recommendations(timestamp, protos.as_ref())?;
1124 }
1125 let mut pr = self.protocols.lock().expect("Poisoned lock");
1126 *pr = Some((timestamp, protos));
1127 self.events.publish(DirEvent::NewProtocolRecommendation);
1128 Ok(())
1129 }
1130 }
1131 } else {
1132 Ok(())
1133 }
1134 }
1135
1136 #[cfg(feature = "dir-plugin")]
1139 pub fn get_plugin(&self) -> as_plugin::DirPlugin {
1140 as_plugin::DirPlugin {
1141 store: Arc::clone(&self.store),
1142 }
1143 }
1144}
1145
1146#[derive(Debug, Copy, Clone)]
1148enum Readiness {
1149 Complete,
1151 Usable,
1153}
1154
1155fn upgrade_weak_ref<T>(weak: &Weak<T>) -> Result<Arc<T>> {
1158 Weak::upgrade(weak).ok_or(Error::ManagerDropped)
1159}
1160
1161pub(crate) fn default_consensus_cutoff(
1164 now: SystemTime,
1165 tolerance: &DirTolerance,
1166) -> Result<SystemTime> {
1167 const MIN_AGE_TO_ALLOW: Duration = Duration::from_secs(3 * 3600);
1170 let allow_skew = std::cmp::max(MIN_AGE_TO_ALLOW, tolerance.post_valid_tolerance());
1171 let cutoff = time::OffsetDateTime::from(now - allow_skew);
1172 let (h, _m, _s) = cutoff.to_hms();
1179 let cutoff = cutoff.replace_time(
1180 time::Time::from_hms(h, 0, 0)
1181 .map_err(tor_error::into_internal!("Failed clock calculation"))?,
1182 );
1183 let cutoff = cutoff + Duration::from_secs(3600);
1184
1185 Ok(cutoff.into())
1186}
1187
1188pub fn supported_client_protocols() -> tor_protover::Protocols {
1191 use tor_protover::named::*;
1192 [
1195 DIRCACHE_CONSDIFF,
1197 ]
1198 .into_iter()
1199 .collect()
1200}
1201
1202#[cfg(test)]
1203mod test {
1204 #![allow(clippy::bool_assert_comparison)]
1206 #![allow(clippy::clone_on_copy)]
1207 #![allow(clippy::dbg_macro)]
1208 #![allow(clippy::mixed_attributes_style)]
1209 #![allow(clippy::print_stderr)]
1210 #![allow(clippy::print_stdout)]
1211 #![allow(clippy::single_char_pattern)]
1212 #![allow(clippy::unwrap_used)]
1213 #![allow(clippy::unchecked_time_subtraction)]
1214 #![allow(clippy::useless_vec)]
1215 #![allow(clippy::needless_pass_by_value)]
1216 #![allow(clippy::string_slice)] use super::*;
1219 use crate::docmeta::{AuthCertMeta, ConsensusMeta};
1220 use std::time::Duration;
1221 use tempfile::TempDir;
1222 use tor_basic_utils::test_rng::testing_rng;
1223 use tor_netdoc::doc::netstatus::ConsensusFlavor;
1224 use tor_netdoc::doc::{authcert::AuthCertKeyIds, netstatus::Lifetime};
1225 use tor_rtcompat::SleepProvider;
1226
1227 #[test]
1228 fn protocols() {
1229 let pr = supported_client_protocols();
1230 let expected = "DirCache=2".parse().unwrap();
1231 assert_eq!(pr, expected);
1232 }
1233
1234 pub(crate) fn new_mgr<R: Runtime>(runtime: R) -> (TempDir, DirMgr<R>) {
1235 let dir = TempDir::new().unwrap();
1236 let config = DirMgrConfig {
1237 cache_dir: dir.path().into(),
1238 ..Default::default()
1239 };
1240 let store = DirMgrStore::new(&config, runtime.clone(), false).unwrap();
1241 let dirmgr = DirMgr::from_config(config, runtime, store, None, false).unwrap();
1242
1243 (dir, dirmgr)
1244 }
1245
1246 #[test]
1247 fn failing_accessors() {
1248 tor_rtcompat::test_with_one_runtime!(|rt| async {
1249 let (_tempdir, mgr) = new_mgr(rt);
1250
1251 assert!(mgr.circmgr().is_err());
1252 assert!(mgr.netdir(Timeliness::Unchecked).is_err());
1253 });
1254 }
1255
1256 #[test]
1257 fn load_and_store_internals() {
1258 tor_rtcompat::test_with_one_runtime!(|rt| async {
1259 let now = rt.wallclock();
1260 let tomorrow = now + Duration::from_secs(86400);
1261 let later = tomorrow + Duration::from_secs(86400);
1262
1263 let (_tempdir, mgr) = new_mgr(rt);
1264
1265 let d1 = [5_u8; 32];
1267 let d2 = [7; 32];
1268 let d3 = [42; 32];
1269 let d4 = [99; 20];
1270 let d5 = [12; 20];
1271 let certid1 = AuthCertKeyIds {
1272 id_fingerprint: d4.into(),
1273 sk_fingerprint: d5.into(),
1274 };
1275 let certid2 = AuthCertKeyIds {
1276 id_fingerprint: d5.into(),
1277 sk_fingerprint: d4.into(),
1278 };
1279
1280 {
1281 let mut store = mgr.store.lock().unwrap();
1282
1283 store
1284 .store_microdescs(
1285 &[
1286 ("Fake micro 1", &d1),
1287 ("Fake micro 2", &d2),
1288 ("Fake micro 3", &d3),
1289 ],
1290 now,
1291 )
1292 .unwrap();
1293
1294 #[cfg(feature = "routerdesc")]
1295 store
1296 .store_routerdescs(&[("Fake rd1", now, &d4), ("Fake rd2", now, &d5)])
1297 .unwrap();
1298
1299 store
1300 .store_authcerts(&[
1301 (
1302 AuthCertMeta::new(certid1, now, tomorrow),
1303 "Fake certificate one",
1304 ),
1305 (
1306 AuthCertMeta::new(certid2, now, tomorrow),
1307 "Fake certificate two",
1308 ),
1309 ])
1310 .unwrap();
1311
1312 let cmeta = ConsensusMeta::new(
1313 Lifetime::new(now, tomorrow, later).unwrap(),
1314 [102; 32],
1315 [103; 32],
1316 );
1317 store
1318 .store_consensus(&cmeta, ConsensusFlavor::Microdesc, false, "Fake consensus!")
1319 .unwrap();
1320 }
1321
1322 let t1 = mgr.text(&DocId::Microdesc(d1)).unwrap().unwrap();
1324 assert_eq!(t1.as_str(), Ok("Fake micro 1"));
1325
1326 let t2 = mgr
1327 .text(&DocId::LatestConsensus {
1328 flavor: ConsensusFlavor::Microdesc,
1329 cache_usage: CacheUsage::CacheOkay,
1330 })
1331 .unwrap()
1332 .unwrap();
1333 assert_eq!(t2.as_str(), Ok("Fake consensus!"));
1334
1335 let t3 = mgr.text(&DocId::Microdesc([255; 32])).unwrap();
1336 assert!(t3.is_none());
1337
1338 let d_bogus = DocId::Microdesc([255; 32]);
1340 let res = mgr
1341 .texts(vec![
1342 DocId::Microdesc(d2),
1343 DocId::Microdesc(d3),
1344 d_bogus,
1345 DocId::AuthCert(certid2),
1346 #[cfg(feature = "routerdesc")]
1347 DocId::RouterDesc(d5),
1348 ])
1349 .unwrap();
1350 assert_eq!(
1351 res.get(&DocId::Microdesc(d2)).unwrap().as_str(),
1352 Ok("Fake micro 2")
1353 );
1354 assert_eq!(
1355 res.get(&DocId::Microdesc(d3)).unwrap().as_str(),
1356 Ok("Fake micro 3")
1357 );
1358 assert!(!res.contains_key(&d_bogus));
1359 assert_eq!(
1360 res.get(&DocId::AuthCert(certid2)).unwrap().as_str(),
1361 Ok("Fake certificate two")
1362 );
1363 #[cfg(feature = "routerdesc")]
1364 assert_eq!(
1365 res.get(&DocId::RouterDesc(d5)).unwrap().as_str(),
1366 Ok("Fake rd2")
1367 );
1368 });
1369 }
1370
1371 #[test]
1372 fn make_consensus_request() {
1373 tor_rtcompat::test_with_one_runtime!(|rt| async {
1374 let now = rt.wallclock();
1375 let tomorrow = now + Duration::from_secs(86400);
1376 let later = tomorrow + Duration::from_secs(86400);
1377
1378 let (_tempdir, mgr) = new_mgr(rt);
1379 let config = DirMgrConfig::default();
1380
1381 let req = {
1383 let store = mgr.store.lock().unwrap();
1384 bootstrap::make_consensus_request(
1385 now,
1386 ConsensusFlavor::Microdesc,
1387 &**store,
1388 &config,
1389 )
1390 .unwrap()
1391 };
1392 let tolerance = DirTolerance::default().post_valid_tolerance();
1393 match req {
1394 ClientRequest::Consensus(r) => {
1395 assert_eq!(r.old_consensus_digests().count(), 0);
1396 let date = r.last_consensus_date().unwrap();
1397 assert!(date >= now - tolerance);
1398 assert!(date <= now - tolerance + Duration::from_secs(3600));
1399 }
1400 _ => panic!("Wrong request type"),
1401 }
1402
1403 let d_prev = [42; 32];
1405 {
1406 let mut store = mgr.store.lock().unwrap();
1407
1408 let cmeta = ConsensusMeta::new(
1409 Lifetime::new(now, tomorrow, later).unwrap(),
1410 d_prev,
1411 [103; 32],
1412 );
1413 store
1414 .store_consensus(&cmeta, ConsensusFlavor::Microdesc, false, "Fake consensus!")
1415 .unwrap();
1416 }
1417
1418 let req = {
1420 let store = mgr.store.lock().unwrap();
1421 bootstrap::make_consensus_request(
1422 now,
1423 ConsensusFlavor::Microdesc,
1424 &**store,
1425 &config,
1426 )
1427 .unwrap()
1428 };
1429 match req {
1430 ClientRequest::Consensus(r) => {
1431 let ds: Vec<_> = r.old_consensus_digests().collect();
1432 assert_eq!(ds.len(), 1);
1433 assert_eq!(ds[0], &d_prev);
1434 assert_eq!(r.last_consensus_date(), Some(now));
1435 }
1436 _ => panic!("Wrong request type"),
1437 }
1438 });
1439 }
1440
1441 #[test]
1442 fn make_other_requests() {
1443 tor_rtcompat::test_with_one_runtime!(|rt| async {
1444 use rand::RngExt;
1445 let (_tempdir, mgr) = new_mgr(rt);
1446
1447 let certid1 = AuthCertKeyIds {
1448 id_fingerprint: [99; 20].into(),
1449 sk_fingerprint: [100; 20].into(),
1450 };
1451 let mut rng = testing_rng();
1452 #[cfg(feature = "routerdesc")]
1453 let rd_ids: Vec<DocId> = (0..1000).map(|_| DocId::RouterDesc(rng.random())).collect();
1454 let md_ids: Vec<DocId> = (0..1000).map(|_| DocId::Microdesc(rng.random())).collect();
1455 let config = DirMgrConfig::default();
1456
1457 let query = DocId::AuthCert(certid1);
1459 let store = mgr.store.lock().unwrap();
1460 let reqs =
1461 bootstrap::make_requests_for_documents(&mgr.runtime, &[query], &**store, &config)
1462 .unwrap();
1463 assert_eq!(reqs.len(), 1);
1464 let req = &reqs[0];
1465 if let ClientRequest::AuthCert(r) = req {
1466 assert_eq!(r.keys().next(), Some(&certid1));
1467 } else {
1468 panic!();
1469 }
1470
1471 let reqs =
1473 bootstrap::make_requests_for_documents(&mgr.runtime, &md_ids, &**store, &config)
1474 .unwrap();
1475 assert_eq!(reqs.len(), 2);
1476 assert!(matches!(reqs[0], ClientRequest::Microdescs(_)));
1477
1478 #[cfg(feature = "routerdesc")]
1480 {
1481 let reqs = bootstrap::make_requests_for_documents(
1482 &mgr.runtime,
1483 &rd_ids,
1484 &**store,
1485 &config,
1486 )
1487 .unwrap();
1488 assert_eq!(reqs.len(), 2);
1489 assert!(matches!(reqs[0], ClientRequest::RouterDescs(_)));
1490 }
1491 });
1492 }
1493
1494 #[test]
1495 fn expand_response() {
1496 tor_rtcompat::test_with_one_runtime!(|rt| async {
1497 let now = rt.wallclock();
1498 let day = Duration::from_secs(86400);
1499 let config = DirMgrConfig::default();
1500
1501 let (_tempdir, mgr) = new_mgr(rt);
1502
1503 let q = DocId::Microdesc([99; 32]);
1505 let r = {
1506 let store = mgr.store.lock().unwrap();
1507 bootstrap::make_requests_for_documents(&mgr.runtime, &[q], &**store, &config)
1508 .unwrap()
1509 };
1510 let expanded = mgr.expand_response_text(&r[0], "ABC".to_string());
1511 assert_eq!(&expanded.unwrap(), "ABC");
1512
1513 let latest_id = DocId::LatestConsensus {
1516 flavor: ConsensusFlavor::Microdesc,
1517 cache_usage: CacheUsage::CacheOkay,
1518 };
1519 let r = {
1520 let store = mgr.store.lock().unwrap();
1521 bootstrap::make_requests_for_documents(
1522 &mgr.runtime,
1523 &[latest_id],
1524 &**store,
1525 &config,
1526 )
1527 .unwrap()
1528 };
1529 let expanded = mgr.expand_response_text(&r[0], "DEF".to_string());
1530 assert_eq!(&expanded.unwrap(), "DEF");
1531
1532 {
1535 let mut store = mgr.store.lock().unwrap();
1536 let d_in = [0x99; 32]; let cmeta = ConsensusMeta::new(
1538 Lifetime::new(now, now + day, now + 2 * day).unwrap(),
1539 d_in,
1540 d_in,
1541 );
1542 store
1543 .store_consensus(
1544 &cmeta,
1545 ConsensusFlavor::Microdesc,
1546 false,
1547 "line 1\nline2\nline 3\n",
1548 )
1549 .unwrap();
1550 }
1551
1552 let r = {
1555 let store = mgr.store.lock().unwrap();
1556 bootstrap::make_requests_for_documents(
1557 &mgr.runtime,
1558 &[latest_id],
1559 &**store,
1560 &config,
1561 )
1562 .unwrap()
1563 };
1564 let expanded = mgr.expand_response_text(&r[0], "hello".to_string());
1565 assert_eq!(&expanded.unwrap(), "hello");
1566
1567 let diff = "network-status-diff-version 1
1569hash 9999999999999999999999999999999999999999999999999999999999999999 8382374ca766873eb0d2530643191c6eaa2c5e04afa554cbac349b5d0592d300
15702c
1571replacement line
1572.
1573".to_string();
1574 let expanded = mgr.expand_response_text(&r[0], diff);
1575
1576 assert_eq!(expanded.unwrap(), "line 1\nreplacement line\nline 3\n");
1577
1578 let diff = "network-status-diff-version 1
1580hash 9999999999999999999999999999999999999999999999999999999999999999 9999999999999999999999999999999999999999999999999999999999999999
15812c
1582replacement line
1583.
1584".to_string();
1585 let expanded = mgr.expand_response_text(&r[0], diff);
1586 assert!(expanded.is_err());
1587 });
1588 }
1589}