1use std::num::NonZeroUsize;
5use std::ops::Deref;
6use std::result::Result as StdResult;
7use std::{
8 collections::HashMap,
9 sync::{Arc, Weak},
10 time::{Duration, SystemTime},
11};
12
13use crate::DirMgrConfig;
14use crate::DocSource;
15use crate::err::BootstrapAction;
16use crate::state::{DirState, PoisonedState};
17use crate::{
18 DirMgr, DocId, DocQuery, DocumentText, Error, Readiness, Result,
19 docid::{self, ClientRequest},
20 upgrade_weak_ref,
21};
22
23use futures::FutureExt;
24use futures::StreamExt;
25use oneshot_fused_workaround as oneshot;
26use tor_dirclient::DirResponse;
27use tor_error::{info_report, warn_report};
28use tor_rtcompat::Runtime;
29use tor_rtcompat::scheduler::TaskSchedule;
30use tracing::{debug, info, instrument, trace, warn};
31
32use crate::storage::Store;
33#[cfg(test)]
34use std::sync::LazyLock;
35#[cfg(test)]
36use std::sync::Mutex;
37use tor_circmgr::{CircMgr, DirInfo};
38use tor_netdir::{NetDir, NetDirProvider as _};
39use tor_netdoc::doc::netstatus::ConsensusFlavor;
40
41macro_rules! propagate_fatal_errors {
44 ( $e:expr ) => {
45 let v: Result<()> = $e;
46 if let Err(e) = v {
47 match e.bootstrap_action() {
48 BootstrapAction::Nonfatal => {}
49 _ => return Err(e),
50 }
51 }
52 };
53}
54
55#[derive(Copy, Clone, Debug, derive_more::Display, Eq, PartialEq, Ord, PartialOrd)]
62#[display("{0}", id)]
63pub(crate) struct AttemptId {
64 id: NonZeroUsize,
66}
67
68impl AttemptId {
69 pub(crate) fn next() -> Self {
76 use std::sync::atomic::{AtomicUsize, Ordering};
77 static NEXT: AtomicUsize = AtomicUsize::new(1);
79 let id = NEXT.fetch_add(1, Ordering::Relaxed);
80 let id = id.try_into().expect("Allocated too many AttemptIds");
81 Self { id }
82 }
83}
84
85fn note_request_outcome<R: Runtime>(
89 circmgr: &CircMgr<R>,
90 outcome: &tor_dirclient::Result<tor_dirclient::DirResponse>,
91) {
92 use tor_dirclient::{Error::RequestFailed, RequestFailedError};
93 let (err, source) = match outcome {
99 Ok(req) => {
100 if let (Some(e), Some(source)) = (req.error(), req.source()) {
101 (
102 RequestFailed(RequestFailedError {
103 error: e.clone(),
104 source: Some(source.clone()),
105 }),
106 source,
107 )
108 } else {
109 return;
110 }
111 }
112 Err(
113 error @ RequestFailed(RequestFailedError {
114 source: Some(source),
115 ..
116 }),
117 ) => (error.clone(), source),
118 _ => return,
119 };
120
121 note_cache_error(circmgr, source, &err.into());
122}
123
124fn note_cache_error<R: Runtime>(
126 circmgr: &CircMgr<R>,
127 source: &tor_dirclient::SourceInfo,
128 problem: &Error,
129) {
130 use tor_circmgr::ExternalActivity;
131
132 if !problem.indicates_cache_failure() {
133 return;
134 }
135
136 let real_source = match problem {
142 Error::NetDocError {
143 source: DocSource::DirServer { source: Some(info) },
144 ..
145 } => info,
146 _ => source,
147 };
148
149 info_report!(problem, "Marking {:?} as failed", real_source);
150 circmgr.note_external_failure(real_source.cache_id(), ExternalActivity::DirCache);
151 circmgr.retire_circ(source.unique_circ_id());
152}
153
154fn note_cache_success<R: Runtime>(circmgr: &CircMgr<R>, source: &tor_dirclient::SourceInfo) {
156 use tor_circmgr::ExternalActivity;
157
158 trace!("Marking {:?} as successful", source);
159 circmgr.note_external_success(source.cache_id(), ExternalActivity::DirCache);
160}
161
162fn load_and_apply_documents<R: Runtime>(
164 missing: &[DocId],
165 dirmgr: &Arc<DirMgr<R>>,
166 state: &mut Box<dyn DirState>,
167 changed: &mut bool,
168) -> Result<()> {
169 const CHUNK_SIZE: usize = 256;
174 for chunk in missing.chunks(CHUNK_SIZE) {
175 let documents = {
176 let store = dirmgr.store.lock().expect("store lock poisoned");
177 load_documents_from_store(chunk, &**store)?
178 };
179
180 state.add_from_cache(documents, changed)?;
181 }
182
183 Ok(())
184}
185
186fn load_documents_from_store(
189 missing: &[DocId],
190 store: &dyn Store,
191) -> Result<HashMap<DocId, DocumentText>> {
192 let mut loaded = HashMap::new();
193 for query in docid::partition_by_type(missing.iter().copied()).values() {
194 query.load_from_store_into(&mut loaded, store)?;
195 }
196 Ok(loaded)
197}
198
199fn update_state<R: Runtime>(
202 dirmgr: &Arc<DirMgr<R>>,
203 attempt_id: AttemptId,
204 state: &mut Box<dyn DirState>,
205) -> Result<SystemTime> {
206 let state_desc = state.describe();
207 let mut changed = false;
208 trace!(attempt=%attempt_id, state=%state_desc,"Attempting to load directory information from cache.");
209 let load_result = load_once(dirmgr, state, attempt_id, &mut changed);
210 trace!(attempt=%attempt_id, state=%state_desc, outcome=?load_result, "Load attempt complete.");
211 if let Err(e) = &load_result {
212 if let Some(source) = e.responsible_cache() {
215 dirmgr.note_errors(attempt_id, 1);
216 note_cache_error(dirmgr.circmgr()?.deref(), source, e);
217 }
218 }
219 propagate_fatal_errors!(load_result);
220 Ok(dirmgr.runtime.wallclock())
221}
222
223fn apply_state<R: Runtime>(
229 dirmgr: &Arc<DirMgr<R>>,
230 state: &mut Box<dyn DirState>,
231 attempt_id: AttemptId,
232) -> Result<()> {
233 let mut store = dirmgr.store.lock().expect("store lock poisoned");
234 dirmgr.apply_netdir_changes(state, &mut **store)?;
235 dirmgr.update_progress(attempt_id, state.bootstrap_progress());
236 Ok(())
237}
238
239enum AdvanceStateError {
241 AlreadyComplete,
243 CantAdvanceYet,
245}
246
247fn advance_state(
258 state: &mut Box<dyn DirState>,
259 attempt_id: AttemptId,
260) -> StdResult<(), AdvanceStateError> {
261 if state.is_ready(Readiness::Complete) {
262 trace!(attempt=%attempt_id, state=%state.describe(), "Directory is now Complete.");
263 return Err(AdvanceStateError::AlreadyComplete);
264 }
265
266 if state.can_advance() {
267 advance(state);
268 trace!(attempt=%attempt_id, state=%state.describe(), "State has advanced.");
269 return Ok(());
270 }
271
272 Err(AdvanceStateError::CantAdvanceYet)
273}
274
275enum DownloadOutcome {
279 DownloadFailed,
281 DirectoryOutdated,
283 Applied,
286}
287
288async fn perform_download<R: Runtime>(
305 attempt_id: AttemptId,
306 dirmgr: &Weak<DirMgr<R>>,
307 now: &mut SystemTime,
308 parallelism: u8,
309 schedule: &mut TaskSchedule<R>,
310 state: &mut Box<dyn DirState>,
311) -> Result<DownloadOutcome> {
312 let reset_time = no_more_than_a_week_from(*now, state.reset_time());
313
314 *now = {
315 let dirmgr = upgrade_weak_ref(dirmgr)?;
316 futures::select_biased! {
317 outcome = download_attempt(&dirmgr, state, parallelism.into(), attempt_id).fuse() => {
318 if let Err(e) = outcome {
319 warn_report!(e, attempt=%attempt_id, "Error while downloading.");
320 propagate_fatal_errors!(Err(e));
321 return Ok(DownloadOutcome::DownloadFailed);
322 } else {
323 trace!(attempt=%attempt_id, "Successfully downloaded some information.");
324 }
325 }
326 _ = schedule.sleep_until_wallclock(reset_time).fuse() => {
327 info!(attempt=%attempt_id, "Directory being fetched is now outdated; resetting download state.");
332 reset(state);
333 return Ok(DownloadOutcome::DirectoryOutdated);
334 },
335 };
336 dirmgr.runtime.wallclock()
337 };
338 Ok(DownloadOutcome::Applied)
339}
340
341pub(crate) fn make_consensus_request(
344 now: SystemTime,
345 flavor: ConsensusFlavor,
346 store: &dyn Store,
347 config: &DirMgrConfig,
348) -> Result<ClientRequest> {
349 let mut request = tor_dirclient::request::ConsensusRequest::new(flavor);
350
351 let default_cutoff = crate::default_consensus_cutoff(now, &config.tolerance)?;
352
353 match store.latest_consensus_meta(flavor) {
354 Ok(Some(meta)) => {
355 let valid_after = meta.lifetime().valid_after();
356 request.set_last_consensus_date(std::cmp::max(valid_after, default_cutoff));
357 request.push_old_consensus_digest(*meta.sha3_256_of_signed());
358 }
359 latest => {
360 if let Err(e) = latest {
361 warn_report!(e, "Error loading directory metadata");
362 }
363 request.set_last_consensus_date(default_cutoff);
367 }
368 }
369
370 request.set_skew_limit(
371 config.tolerance.post_valid_tolerance(),
374 config.tolerance.pre_valid_tolerance(),
377 );
378
379 Ok(ClientRequest::Consensus(request))
380}
381
382pub(crate) fn make_requests_for_documents<R: Runtime>(
384 rt: &R,
385 docs: &[DocId],
386 store: &dyn Store,
387 config: &DirMgrConfig,
388) -> Result<Vec<ClientRequest>> {
389 let mut res = Vec::new();
390 for q in docid::partition_by_type(docs.iter().copied())
391 .into_values()
392 .flat_map(|x| x.split_for_download().into_iter())
393 {
394 match q {
395 DocQuery::LatestConsensus { flavor, .. } => {
396 res.push(make_consensus_request(
397 rt.wallclock(),
398 flavor,
399 store,
400 config,
401 )?);
402 }
403 DocQuery::AuthCert(ids) => {
404 res.push(ClientRequest::AuthCert(ids.into_iter().collect()));
405 }
406 DocQuery::Microdesc(ids) => {
407 res.push(ClientRequest::Microdescs(ids.into_iter().collect()));
408 }
409 #[cfg(feature = "routerdesc")]
410 DocQuery::RouterDesc(ids) => {
411 res.push(ClientRequest::RouterDescs(ids.into_iter().collect()));
412 }
413 }
414 }
415 Ok(res)
416}
417
418#[instrument(level = "trace", skip_all)]
420async fn fetch_single<R: Runtime>(
421 rt: &R,
422 request: ClientRequest,
423 current_netdir: Option<&NetDir>,
424 circmgr: Arc<CircMgr<R>>,
425) -> Result<(ClientRequest, DirResponse)> {
426 let dirinfo: DirInfo = match current_netdir {
427 Some(netdir) => netdir.into(),
428 None => tor_circmgr::DirInfo::Nothing,
429 };
430 let outcome =
431 tor_dirclient::get_resource(request.as_requestable(), dirinfo, rt, circmgr.clone()).await;
432
433 note_request_outcome(&circmgr, &outcome);
434
435 let resource = outcome?;
436 Ok((request, resource))
437}
438
439#[cfg(test)]
445static CANNED_RESPONSE: LazyLock<Mutex<Vec<String>>> = LazyLock::new(|| Mutex::new(vec![]));
446
447#[instrument(level = "trace", skip_all)]
452async fn fetch_multiple<R: Runtime>(
453 dirmgr: Arc<DirMgr<R>>,
454 attempt_id: AttemptId,
455 missing: &[DocId],
456 parallelism: usize,
457) -> Result<Vec<(ClientRequest, DirResponse)>> {
458 let requests = {
459 let store = dirmgr.store.lock().expect("store lock poisoned");
460 make_requests_for_documents(&dirmgr.runtime, missing, &**store, &dirmgr.config.get())?
461 };
462
463 trace!(attempt=%attempt_id, "Launching {} requests for {} documents",
464 requests.len(), missing.len());
465
466 #[cfg(test)]
467 {
468 let m = CANNED_RESPONSE.lock().expect("Poisoned mutex");
469 if !m.is_empty() {
470 return Ok(requests
471 .into_iter()
472 .zip(m.iter().map(DirResponse::from_get_body))
473 .collect());
474 }
475 }
476
477 let circmgr = dirmgr.circmgr()?;
478 let netdir = dirmgr.netdir(tor_netdir::Timeliness::Timely).ok();
480
481 let responses: Vec<Result<(ClientRequest, DirResponse)>> = futures::stream::iter(requests)
484 .map(|query| fetch_single(&dirmgr.runtime, query, netdir.as_deref(), circmgr.clone()))
485 .buffer_unordered(parallelism)
486 .collect()
487 .await;
488
489 let mut useful_responses = Vec::new();
490 for r in responses {
491 match r {
493 Ok((request, response)) => {
494 if response.status_code() == 200 {
495 useful_responses.push((request, response));
496 } else {
497 trace!(
498 "cache declined request; reported status {:?}",
499 response.status_code()
500 );
501 }
502 }
503 Err(e) => warn_report!(e, "error while downloading"),
504 }
505 }
506
507 trace!(attempt=%attempt_id, "received {} useful responses from our requests.", useful_responses.len());
508
509 Ok(useful_responses)
510}
511
512fn load_once<R: Runtime>(
514 dirmgr: &Arc<DirMgr<R>>,
515 state: &mut Box<dyn DirState>,
516 attempt_id: AttemptId,
517 changed_out: &mut bool,
518) -> Result<()> {
519 let missing = state.missing_docs();
520 let mut changed = false;
521 let outcome: Result<()> = if missing.is_empty() {
522 trace!("Found no missing documents; can't advance current state");
523 Ok(())
524 } else {
525 trace!(
526 "Found {} missing documents; trying to load them",
527 missing.len()
528 );
529
530 load_and_apply_documents(&missing, dirmgr, state, &mut changed)
531 };
532
533 if changed {
537 dirmgr.update_progress(attempt_id, state.bootstrap_progress());
538 *changed_out = true;
539 }
540
541 outcome
542}
543
544pub(crate) fn load<R: Runtime>(
549 dirmgr: &Arc<DirMgr<R>>,
550 mut state: Box<dyn DirState>,
551 attempt_id: AttemptId,
552) -> Result<Box<dyn DirState>> {
553 let mut safety_counter = 0_usize;
554 loop {
555 trace!(attempt=%attempt_id, state=%state.describe(), "Loading from cache");
556 let mut changed = false;
557 let outcome = load_once(dirmgr, &mut state, attempt_id, &mut changed);
558 apply_state(dirmgr, &mut state, attempt_id)?;
559 trace!(attempt=%attempt_id, ?outcome, "Load operation completed.");
560
561 if let Err(e) = outcome {
562 match e.bootstrap_action() {
563 BootstrapAction::Nonfatal => {
564 debug!("Recoverable error loading from cache: {}", e);
565 }
566 BootstrapAction::Fatal | BootstrapAction::Reset => {
567 return Err(e);
568 }
569 }
570 }
571
572 if state.can_advance() {
573 state = state.advance();
574 trace!(attempt=%attempt_id, state=state.describe(), "State has advanced.");
575 safety_counter = 0;
576 } else {
577 if !changed {
578 trace!(attempt=%attempt_id, state=state.describe(), "No state advancement after load; nothing more to find in the cache.");
581 break;
582 }
583 safety_counter += 1;
584 assert!(
585 safety_counter < 100,
586 "Spent 100 iterations in the same state: this is a bug"
587 );
588 }
589 }
590
591 Ok(state)
592}
593
594#[instrument(level = "trace", skip_all)]
600async fn download_attempt<R: Runtime>(
601 dirmgr: &Arc<DirMgr<R>>,
602 state: &mut Box<dyn DirState>,
603 parallelism: usize,
604 attempt_id: AttemptId,
605) -> Result<()> {
606 let missing = state.missing_docs();
607 let fetched = fetch_multiple(Arc::clone(dirmgr), attempt_id, &missing, parallelism).await?;
608 let mut n_errors = 0;
609 for (client_req, dir_response) in fetched {
610 let source = dir_response.source().cloned();
611 let text = match String::from_utf8(dir_response.into_output_unchecked())
612 .map_err(Error::BadUtf8FromDirectory)
613 {
614 Ok(t) => t,
615 Err(e) => {
616 if let Some(source) = source {
617 n_errors += 1;
618 note_cache_error(dirmgr.circmgr()?.deref(), &source, &e);
619 }
620 continue;
621 }
622 };
623 match dirmgr.expand_response_text(&client_req, text) {
624 Ok(text) => {
625 let doc_source = DocSource::DirServer {
626 source: source.clone(),
627 };
628 let mut changed = false;
629 let outcome = state.add_from_download(
630 &text,
631 &client_req,
632 doc_source,
633 Some(&dirmgr.store),
634 &mut changed,
635 );
636
637 if !changed {
638 debug_assert!(outcome.is_err());
639 }
640
641 if let Some(source) = source {
642 if let Err(e) = &outcome {
643 n_errors += 1;
644 note_cache_error(dirmgr.circmgr()?.deref(), &source, e);
645 } else {
646 note_cache_success(dirmgr.circmgr()?.deref(), &source);
647 }
648 }
649
650 if let Err(e) = &outcome {
651 dirmgr.note_errors(attempt_id, 1);
652 warn_report!(e, "error while adding directory info");
653 }
654 propagate_fatal_errors!(outcome);
655 }
656 Err(e) => {
657 warn_report!(e, "Error when expanding directory text");
658 if let Some(source) = source {
659 n_errors += 1;
660 note_cache_error(dirmgr.circmgr()?.deref(), &source, &e);
661 }
662 propagate_fatal_errors!(Err(e));
663 }
664 }
665 }
666 if n_errors != 0 {
667 dirmgr.note_errors(attempt_id, n_errors);
668 }
669 dirmgr.update_progress(attempt_id, state.bootstrap_progress());
670
671 Ok(())
672}
673
674#[instrument(level = "trace", skip_all)]
684pub(crate) async fn download<R: Runtime>(
685 dirmgr: Weak<DirMgr<R>>,
686 state: &mut Box<dyn DirState>,
687 schedule: &mut TaskSchedule<R>,
688 attempt_id: AttemptId,
689 tx_usable: &mut Option<oneshot::Sender<()>>,
690) -> Result<()> {
691 let runtime = upgrade_weak_ref(&dirmgr)?.runtime.clone();
692
693 trace!(attempt=%attempt_id, state=%state.describe(), "Trying to download directory material.");
694
695 'next_state: loop {
696 let retry_config = state.dl_config();
697 let parallelism = retry_config.parallelism();
698
699 let mut now = update_state(&upgrade_weak_ref(&dirmgr)?, attempt_id, state)?;
703
704 apply_state(&upgrade_weak_ref(&dirmgr)?, state, attempt_id)?;
705
706 match advance_state(state, attempt_id) {
709 Ok(_) => continue 'next_state,
710 Err(AdvanceStateError::AlreadyComplete) => return Ok(()),
711 Err(AdvanceStateError::CantAdvanceYet) => {}
712 }
713
714 let reset_time = no_more_than_a_week_from(runtime.wallclock(), state.reset_time());
715
716 let mut retry = retry_config.schedule();
717 let mut delay = None;
718
719 'next_attempt: for attempt in retry_config.attempts() {
723 let next_delay = retry.next_delay(&mut rand::rng());
727 if let Some(delay) = delay.replace(next_delay) {
728 let time_until_reset = reset_time
729 .duration_since(now)
730 .unwrap_or(Duration::from_secs(0));
731 let real_delay = delay.min(time_until_reset);
732 debug!(attempt=%attempt_id, "Waiting {:?} for next download attempt...", real_delay);
733 schedule.sleep(real_delay).await?;
734
735 now = upgrade_weak_ref(&dirmgr)?.runtime.wallclock();
736 if now >= reset_time {
737 info!(attempt=%attempt_id, "Directory being fetched is now outdated; resetting download state.");
738 reset(state);
739 continue 'next_state;
740 }
741 }
742
743 info!(attempt=%attempt_id, "{}: {}", attempt + 1, state.describe());
744 match perform_download(attempt_id, &dirmgr, &mut now, parallelism, schedule, state)
745 .await?
746 {
747 DownloadOutcome::DownloadFailed => continue 'next_state,
748 DownloadOutcome::DirectoryOutdated => continue 'next_attempt,
749 DownloadOutcome::Applied => {}
750 }
751
752 propagate_fatal_errors!(apply_state(&upgrade_weak_ref(&dirmgr)?, state, attempt_id));
753
754 match advance_state(state, attempt_id) {
756 Err(AdvanceStateError::AlreadyComplete) => return Ok(()),
757 Ok(()) => continue 'next_state,
758 Err(AdvanceStateError::CantAdvanceYet) => {
759 if state.is_ready(Readiness::Usable) {
760 if let Some(tx) = tx_usable.take() {
761 trace!(attempt=%attempt_id, state=%state.describe(), "directory is now usable.");
762 let _ = tx.send(());
763 }
764 }
765 }
766 }
767 }
768
769 warn!(n_attempts=retry_config.n_attempts(),
771 state=%state.describe(),
772 "Unable to advance downloading state");
773 return Err(Error::CantAdvanceState);
774 }
775}
776
777fn reset(state: &mut Box<dyn DirState>) {
779 let cur_state = std::mem::replace(state, Box::new(PoisonedState));
780 *state = cur_state.reset();
781}
782
783fn advance(state: &mut Box<dyn DirState>) {
785 let cur_state = std::mem::replace(state, Box::new(PoisonedState));
786 *state = cur_state.advance();
787}
788
789fn no_more_than_a_week_from(now: SystemTime, v: Option<SystemTime>) -> SystemTime {
796 let one_week_later = now + Duration::new(86400 * 7, 0);
797 match v {
798 Some(t) => std::cmp::min(t, one_week_later),
799 None => one_week_later,
800 }
801}
802
803#[cfg(test)]
804mod test {
805 #![allow(clippy::bool_assert_comparison)]
807 #![allow(clippy::clone_on_copy)]
808 #![allow(clippy::dbg_macro)]
809 #![allow(clippy::mixed_attributes_style)]
810 #![allow(clippy::print_stderr)]
811 #![allow(clippy::print_stdout)]
812 #![allow(clippy::single_char_pattern)]
813 #![allow(clippy::unwrap_used)]
814 #![allow(clippy::unchecked_time_subtraction)]
815 #![allow(clippy::useless_vec)]
816 #![allow(clippy::needless_pass_by_value)]
817 #![allow(clippy::string_slice)] use super::*;
820 use crate::storage::DynStore;
821 use crate::test::new_mgr;
822 use std::sync::Mutex;
823 use tor_dircommon::retry::DownloadSchedule;
824 use tor_netdoc::doc::microdesc::MdDigest;
825 use tor_rtcompat::SleepProvider;
826 use web_time_compat::SystemTimeExt;
827
828 #[test]
829 fn week() {
830 let now = SystemTime::get();
831 let one_day = Duration::new(86400, 0);
832
833 assert_eq!(no_more_than_a_week_from(now, None), now + one_day * 7);
834 assert_eq!(
835 no_more_than_a_week_from(now, Some(now + one_day)),
836 now + one_day
837 );
838 assert_eq!(
839 no_more_than_a_week_from(now, Some(now - one_day)),
840 now - one_day
841 );
842 assert_eq!(
843 no_more_than_a_week_from(now, Some(now + 30 * one_day)),
844 now + one_day * 7
845 );
846 }
847
848 #[derive(Debug, Clone)]
852 struct DemoState {
853 second_time_around: bool,
854 got_items: HashMap<MdDigest, bool>,
855 }
856
857 const H1: MdDigest = *b"satellite's gone up to the skies";
859 const H2: MdDigest = *b"things like that drive me out of";
860 const H3: MdDigest = *b"my mind i watched it for a littl";
861 const H4: MdDigest = *b"while i like to watch things on ";
862 const H5: MdDigest = *b"TV Satellite of love Satellite--";
863
864 impl DemoState {
865 fn new1() -> Self {
866 DemoState {
867 second_time_around: false,
868 got_items: vec![(H1, false), (H2, false)].into_iter().collect(),
869 }
870 }
871 fn new2() -> Self {
872 DemoState {
873 second_time_around: true,
874 got_items: vec![(H3, false), (H4, false), (H5, false)]
875 .into_iter()
876 .collect(),
877 }
878 }
879 fn n_ready(&self) -> usize {
880 self.got_items.values().filter(|x| **x).count()
881 }
882 }
883
884 impl DirState for DemoState {
885 fn describe(&self) -> String {
886 format!("{:?}", self)
887 }
888 fn bootstrap_progress(&self) -> crate::event::DirProgress {
889 crate::event::DirProgress::default()
890 }
891 fn is_ready(&self, ready: Readiness) -> bool {
892 match (ready, self.second_time_around) {
893 (_, false) => false,
894 (Readiness::Complete, true) => self.n_ready() == self.got_items.len(),
895 (Readiness::Usable, true) => self.n_ready() >= self.got_items.len() - 1,
896 }
897 }
898 fn can_advance(&self) -> bool {
899 if self.second_time_around {
900 false
901 } else {
902 self.n_ready() == self.got_items.len()
903 }
904 }
905 fn missing_docs(&self) -> Vec<DocId> {
906 self.got_items
907 .iter()
908 .filter_map(|(id, have)| {
909 if *have {
910 None
911 } else {
912 Some(DocId::Microdesc(*id))
913 }
914 })
915 .collect()
916 }
917 fn add_from_cache(
918 &mut self,
919 docs: HashMap<DocId, DocumentText>,
920 changed: &mut bool,
921 ) -> Result<()> {
922 for id in docs.keys() {
923 if let DocId::Microdesc(id) = id {
924 if self.got_items.get(id) == Some(&false) {
925 self.got_items.insert(*id, true);
926 *changed = true;
927 }
928 }
929 }
930 Ok(())
931 }
932 fn add_from_download(
933 &mut self,
934 text: &str,
935 _request: &ClientRequest,
936 _source: DocSource,
937 _storage: Option<&Mutex<DynStore>>,
938 changed: &mut bool,
939 ) -> Result<()> {
940 for token in text.split_ascii_whitespace() {
941 if let Ok(v) = hex::decode(token) {
942 if let Ok(id) = v.try_into() {
943 if self.got_items.get(&id) == Some(&false) {
944 self.got_items.insert(id, true);
945 *changed = true;
946 }
947 }
948 }
949 }
950 Ok(())
951 }
952 fn dl_config(&self) -> DownloadSchedule {
953 DownloadSchedule::default()
954 }
955 fn advance(self: Box<Self>) -> Box<dyn DirState> {
956 if self.can_advance() {
957 Box::new(Self::new2())
958 } else {
959 self
960 }
961 }
962 fn reset_time(&self) -> Option<SystemTime> {
963 None
964 }
965 fn reset(self: Box<Self>) -> Box<dyn DirState> {
966 Box::new(Self::new1())
967 }
968 }
969
970 #[test]
971 fn all_in_cache() {
972 tor_rtcompat::test_with_one_runtime!(|rt| async {
974 let now = rt.wallclock();
975 let (_tempdir, mgr) = new_mgr(rt.clone());
976 let (mut schedule, _handle) = TaskSchedule::new(rt);
977
978 {
979 let mut store = mgr.store_if_rw().unwrap().lock().unwrap();
980 for h in [H1, H2, H3, H4, H5] {
981 store.store_microdescs(&[("ignore", &h)], now).unwrap();
982 }
983 }
984 let mgr = Arc::new(mgr);
985 let attempt_id = AttemptId::next();
986
987 let state = Box::new(DemoState::new1());
989 let result = super::load(&mgr, state, attempt_id).unwrap();
990 assert!(result.is_ready(Readiness::Complete));
991
992 let mut state: Box<dyn DirState> = Box::new(DemoState::new1());
994
995 let mut tx_usable = None;
996 super::download(
997 Arc::downgrade(&mgr),
998 &mut state,
999 &mut schedule,
1000 attempt_id,
1001 &mut tx_usable,
1002 )
1003 .await
1004 .unwrap();
1005 assert!(state.is_ready(Readiness::Complete));
1006 });
1007 }
1008
1009 #[test]
1010 fn partly_in_cache() {
1011 tor_rtcompat::test_with_one_runtime!(|rt| async {
1014 let now = rt.wallclock();
1015 let (_tempdir, mgr) = new_mgr(rt.clone());
1016 let (mut schedule, _handle) = TaskSchedule::new(rt);
1017
1018 {
1019 let mut store = mgr.store_if_rw().unwrap().lock().unwrap();
1020 for h in [H1, H2, H3] {
1021 store.store_microdescs(&[("ignore", &h)], now).unwrap();
1022 }
1023 }
1024 {
1025 let mut resp = CANNED_RESPONSE.lock().unwrap();
1026 *resp = vec![
1028 "7768696c652069206c696b6520746f207761746368207468696e6773206f6e20
1029 545620536174656c6c697465206f66206c6f766520536174656c6c6974652d2d"
1030 .to_owned(),
1031 ];
1032 }
1033 let mgr = Arc::new(mgr);
1034 let mut tx_usable = None;
1035 let attempt_id = AttemptId::next();
1036
1037 let mut state: Box<dyn DirState> = Box::new(DemoState::new1());
1038 super::download(
1039 Arc::downgrade(&mgr),
1040 &mut state,
1041 &mut schedule,
1042 attempt_id,
1043 &mut tx_usable,
1044 )
1045 .await
1046 .unwrap();
1047 assert!(state.is_ready(Readiness::Complete));
1048 });
1049 }
1050}