Loading crates/tor-dirmgr/src/bootstrap.rs +26 −31 Original line number Diff line number Diff line Loading @@ -303,11 +303,11 @@ async fn fetch_multiple<R: Runtime>( async fn load_once<R: Runtime>( dirmgr: &Arc<DirMgr<R>>, state: &mut Box<dyn DirState>, ) -> Result<(bool, Option<Error>)> { ) -> Result<Option<Error>> { let missing = state.missing_docs(); let outcome: Result<(bool, Option<Error>)> = if missing.is_empty() { let outcome: Result<Option<Error>> = if missing.is_empty() { trace!("Found no missing documents; can't advance current state"); Ok((false, None)) Ok(Some(Error::NoChange(DocSource::LocalCache))) } else { trace!( "Found {} missing documents; trying to load them", Loading @@ -323,13 +323,13 @@ async fn load_once<R: Runtime>( Err(Error::UntimelyObject(TimeValidityError::Expired(_))) => { // This is just an expired object from the cache; we don't need // to call that an error. Treat it as if it were absent. Ok((false, None)) Ok(None) } other => other, } }; if matches!(outcome, Ok((true, _))) { if outcome.is_ok() { dirmgr.update_status(state.bootstrap_status()); } Loading @@ -347,13 +347,13 @@ pub(crate) async fn load<R: Runtime>( let mut safety_counter = 0_usize; loop { trace!(state=%state.describe(), "Loading from cache"); let (changed, nonfatal_err) = load_once(&dirmgr, &mut state).await?; let nonfatal_err = load_once(&dirmgr, &mut state).await?; { let mut store = dirmgr.store.lock().expect("store lock poisoned"); dirmgr.apply_netdir_changes(&mut state, store.deref_mut())?; } if let Some(e) = nonfatal_err { if let Some(e) = &nonfatal_err { debug!("Recoverable loading from cache: {}", e); } if let Some(e) = state.blocking_error() { Loading @@ -364,7 +364,8 @@ pub(crate) async fn load<R: Runtime>( state = state.advance()?; safety_counter = 0; } else { if !changed { if matches!(nonfatal_err, Some(Error::NoChange(_))) { // XXXX refactor more. break; } safety_counter += 1; Loading @@ -383,14 +384,11 @@ pub(crate) async fn load<R: Runtime>( /// /// This can launch one or more download requests, but will not launch more /// than `parallelism` requests at a time. /// /// Return true if the state reports that it changed. async fn download_attempt<R: Runtime>( dirmgr: &Arc<DirMgr<R>>, state: &mut Box<dyn DirState>, parallelism: usize, ) -> Result<bool> { let mut changed = false; ) -> Result<()> { let missing = state.missing_docs(); let fetched = fetch_multiple(Arc::clone(dirmgr), &missing, parallelism).await?; for (client_req, dir_response) in fetched { Loading @@ -411,9 +409,8 @@ async fn download_attempt<R: Runtime>( let doc_source = DocSource::DirServer { source: source.clone(), }; let (b, opt_error) = let opt_error = state.add_from_download(&text, &client_req, doc_source, Some(&dirmgr.store))?; changed |= b; if let Some(e) = &opt_error { warn!("error while adding directory info: {}", e); } Loading @@ -434,11 +431,9 @@ async fn download_attempt<R: Runtime>( } } if changed { dirmgr.update_status(state.bootstrap_status()); } Ok(changed) Ok(()) } /// Download information into a DirState state machine until it is Loading Loading @@ -534,15 +529,10 @@ pub(crate) async fn download<R: Runtime>( let dirmgr = upgrade_weak_ref(&dirmgr)?; futures::select_biased! { outcome = download_attempt(&dirmgr, &mut state, parallelism.into()).fuse() => { match outcome { Err(e) => { if let Err(e) = outcome { warn!("Error while downloading: {}", e); continue 'next_attempt; } Ok(changed) => { changed } } } _ = runtime.sleep_until_wallclock(reset_time).fuse() => { // We need to reset. This can happen if (for Loading Loading @@ -716,10 +706,7 @@ mod test { }) .collect() } fn add_from_cache( &mut self, docs: HashMap<DocId, DocumentText>, ) -> Result<(bool, Option<Error>)> { fn add_from_cache(&mut self, docs: HashMap<DocId, DocumentText>) -> Result<Option<Error>> { let mut changed = false; for id in docs.keys() { if let DocId::Microdesc(id) = id { Loading @@ -729,7 +716,11 @@ mod test { } } } Ok((changed, None)) if changed { Ok(None) } else { Ok(Some(Error::NoChange(DocSource::LocalCache))) } } fn add_from_download( &mut self, Loading @@ -737,7 +728,7 @@ mod test { _request: &ClientRequest, _source: DocSource, _storage: Option<&Mutex<DynStore>>, ) -> Result<(bool, Option<Error>)> { ) -> Result<Option<Error>> { let mut changed = false; for token in text.split_ascii_whitespace() { if let Ok(v) = hex::decode(token) { Loading @@ -749,7 +740,11 @@ mod test { } } } Ok((changed, None)) if changed { Ok(None) } else { Ok(Some(Error::NoChange(DocSource::LocalCache))) } } fn dl_config(&self) -> DownloadSchedule { DownloadSchedule::default() Loading crates/tor-dirmgr/src/state.rs +43 −53 Original line number Diff line number Diff line Loading @@ -116,13 +116,9 @@ pub(crate) trait DirState: Send { /// Return an `Err` only on a fatal error that means we should stop /// downloading entirely; return recovered-from errors inside the `Ok()` /// variant. fn add_from_cache( &mut self, docs: HashMap<DocId, DocumentText>, ) -> Result<(bool, Option<Error>)>; fn add_from_cache(&mut self, docs: HashMap<DocId, DocumentText>) -> Result<Option<Error>>; /// Add information that we have just downloaded to this state; returns /// 'true' if there as any change in this state. /// Add information that we have just downloaded to this state. /// /// This method receives a copy of the original request, and should reject /// any documents that do not pertain to it. Loading @@ -139,7 +135,7 @@ pub(crate) trait DirState: Send { request: &ClientRequest, source: DocSource, storage: Option<&Mutex<DynStore>>, ) -> Result<(bool, Option<Error>)>; ) -> Result<Option<Error>>; /// Return a summary of this state as a [`DirStatus`]. fn bootstrap_status(&self) -> event::DirStatus; /// Return a configuration for attempting downloads. Loading Loading @@ -279,12 +275,9 @@ impl<R: Runtime> DirState for GetConsensusState<R> { fn dl_config(&self) -> DownloadSchedule { self.config.schedule.retry_consensus } fn add_from_cache( &mut self, docs: HashMap<DocId, DocumentText>, ) -> Result<(bool, Option<Error>)> { fn add_from_cache(&mut self, docs: HashMap<DocId, DocumentText>) -> Result<Option<Error>> { let text = match docs.into_iter().next() { None => return Ok((false, None)), None => return Ok(Some(Error::NoChange(DocSource::LocalCache))), Some(( DocId::LatestConsensus { flavor: ConsensusFlavor::Microdesc, Loading @@ -297,8 +290,8 @@ impl<R: Runtime> DirState for GetConsensusState<R> { let source = DocSource::LocalCache; match self.add_consensus_text(source, text.as_str().map_err(Error::BadUtf8InCache)?, None) { Ok(_) => Ok((true, None)), Err(e) => Ok((false, Some(e))), Ok(_) => Ok(None), Err(e) => Ok(Some(e)), } } fn add_from_download( Loading @@ -307,7 +300,7 @@ impl<R: Runtime> DirState for GetConsensusState<R> { request: &ClientRequest, source: DocSource, storage: Option<&Mutex<DynStore>>, ) -> Result<(bool, Option<Error>)> { ) -> Result<Option<Error>> { let requested_newer_than = match request { ClientRequest::Consensus(r) => r.last_consensus_date(), _ => None, Loading @@ -318,9 +311,9 @@ impl<R: Runtime> DirState for GetConsensusState<R> { let mut w = store.lock().expect("Directory storage lock poisoned"); w.store_consensus(meta, ConsensusFlavor::Microdesc, true, text)?; } Ok((true, None)) Ok(None) } Err(e) => Ok((false, Some(e))), Err(e) => Ok(Some(e)), } } fn advance(self: Box<Self>) -> Result<Box<dyn DirState>> { Loading Loading @@ -575,10 +568,7 @@ impl<R: Runtime> DirState for GetCertsState<R> { fn dl_config(&self) -> DownloadSchedule { self.config.schedule.retry_certs } fn add_from_cache( &mut self, docs: HashMap<DocId, DocumentText>, ) -> Result<(bool, Option<Error>)> { fn add_from_cache(&mut self, docs: HashMap<DocId, DocumentText>) -> Result<Option<Error>> { let mut changed = false; // Here we iterate over the documents we want, taking them from // our input and remembering them. Loading @@ -602,8 +592,10 @@ impl<R: Runtime> DirState for GetCertsState<R> { } if changed { self.try_checking_sigs(); Ok(nonfatal_error) } else { Ok(Some(Error::NoChange(DocSource::LocalCache))) } Ok((changed, nonfatal_error)) } fn add_from_download( &mut self, Loading @@ -611,7 +603,7 @@ impl<R: Runtime> DirState for GetCertsState<R> { request: &ClientRequest, source: DocSource, storage: Option<&Mutex<DynStore>>, ) -> Result<(bool, Option<Error>)> { ) -> Result<Option<Error>> { let asked_for: HashSet<_> = match request { ClientRequest::AuthCert(a) => a.keys().collect(), _ => return Err(internal!("expected an AuthCert request").into()), Loading Loading @@ -644,7 +636,7 @@ impl<R: Runtime> DirState for GetCertsState<R> { // We want to exit early if we aren't saving any certificates. if newcerts.is_empty() { return Ok((false, nonfatal_error)); return Ok(nonfatal_error.or(Some(Error::NoChange(source)))); } if let Some(store) = storage { Loading @@ -671,9 +663,10 @@ impl<R: Runtime> DirState for GetCertsState<R> { if changed { self.try_checking_sigs(); Ok(nonfatal_error) } else { Ok(Some(Error::NoChange(source))) } Ok((changed, nonfatal_error)) } fn advance(self: Box<Self>) -> Result<Box<dyn DirState>> { Loading Loading @@ -997,10 +990,7 @@ impl<R: Runtime> DirState for GetMicrodescsState<R> { fn dl_config(&self) -> DownloadSchedule { self.config.schedule.retry_microdescs } fn add_from_cache( &mut self, docs: HashMap<DocId, DocumentText>, ) -> Result<(bool, Option<Error>)> { fn add_from_cache(&mut self, docs: HashMap<DocId, DocumentText>) -> Result<Option<Error>> { let mut microdescs = Vec::new(); for (id, text) in docs { if let DocId::Microdesc(digest) = id { Loading @@ -1014,10 +1004,12 @@ impl<R: Runtime> DirState for GetMicrodescsState<R> { } } let changed = !microdescs.is_empty(); if microdescs.is_empty() { Ok(Some(Error::NoChange(DocSource::LocalCache))) } else { self.register_microdescs(microdescs); Ok((changed, None)) Ok(None) } } fn add_from_download( Loading @@ -1026,7 +1018,7 @@ impl<R: Runtime> DirState for GetMicrodescsState<R> { request: &ClientRequest, source: DocSource, storage: Option<&Mutex<DynStore>>, ) -> Result<(bool, Option<Error>)> { ) -> Result<Option<Error>> { let requested: HashSet<_> = if let ClientRequest::Microdescs(req) = request { req.digests().collect() } else { Loading Loading @@ -1079,8 +1071,12 @@ impl<R: Runtime> DirState for GetMicrodescsState<R> { )?; } } if new_mds.is_empty() { Ok(Some(Error::NoChange(source))) } else { self.register_microdescs(new_mds.into_iter().map(|(_, md)| md)); Ok((true, nonfatal_err)) Ok(nonfatal_err) } } fn advance(self: Box<Self>) -> Result<Box<dyn DirState>> { Ok(self) Loading Loading @@ -1358,10 +1354,7 @@ mod test { source.clone(), Some(&store), ); assert!(matches!( outcome, Ok((false, Some(Error::NetDocError { .. }))) )); assert!(matches!(outcome, Ok(Some(Error::NetDocError { .. })))); // make sure it wasn't stored... assert!(store .lock() Loading @@ -1372,10 +1365,7 @@ mod test { // Now try again, with a real consensus... but the wrong authorities. let outcome = state.add_from_download(CONSENSUS, &req, source.clone(), Some(&store)); assert!(matches!( outcome, Ok((false, Some(Error::UnrecognizedAuthorities))) )); assert!(matches!(outcome, Ok(Some(Error::UnrecognizedAuthorities)))); assert!(store .lock() .unwrap() Loading @@ -1396,7 +1386,7 @@ mod test { Arc::new(crate::filter::NilFilter), ); let outcome = state.add_from_download(CONSENSUS, &req, source, Some(&store)); assert!(outcome.unwrap().0); assert!(outcome.unwrap().is_none()); assert!(store .lock() .unwrap() Loading Loading @@ -1427,7 +1417,7 @@ mod test { let text: crate::storage::InputString = CONSENSUS.to_owned().into(); let map = vec![(docid, text.into())].into_iter().collect(); let outcome = state.add_from_cache(map); assert!(outcome.unwrap().0); assert!(outcome.unwrap().is_none()); assert!(state.can_advance()); }); } Loading @@ -1451,7 +1441,7 @@ mod test { let req = tor_dirclient::request::ConsensusRequest::new(ConsensusFlavor::Microdesc); let req = crate::docid::ClientRequest::Consensus(req); let outcome = state.add_from_download(CONSENSUS, &req, source, None); assert!(outcome.unwrap().0); assert!(outcome.unwrap().is_none()); Box::new(state).advance().unwrap() } Loading Loading @@ -1495,7 +1485,7 @@ mod test { .into_iter() .collect(); let outcome = state.add_from_cache(docs); assert!(outcome.unwrap().0); // no error, and something changed. assert!(outcome.unwrap().is_none()); // no error, and something changed. assert!(!state.can_advance()); // But we aren't done yet. let missing = state.missing_docs(); assert_eq!(missing.len(), 1); // Now we're only missing one! Loading @@ -1513,7 +1503,7 @@ mod test { let req = ClientRequest::AuthCert(req); let outcome = state.add_from_download(AUTHCERT_5A23, &req, source.clone(), Some(&store)); assert!(!outcome.unwrap().0); // no error, but nothing changed. assert!(matches!(outcome.unwrap(), Some(Error::Unwanted(_)))); let missing2 = state.missing_docs(); assert_eq!(missing, missing2); // No change. assert!(store Loading @@ -1528,7 +1518,7 @@ mod test { req.push(authcert_id_5a23()); // Right idea this time! let req = ClientRequest::AuthCert(req); let outcome = state.add_from_download(AUTHCERT_5A23, &req, source, Some(&store)); assert!(outcome.unwrap().0); // No error, _and_ something changed! assert!(outcome.unwrap().is_none()); // No error, _and_ something changed! let missing3 = state.missing_docs(); assert!(missing3.is_empty()); assert!(state.can_advance()); Loading Loading @@ -1631,7 +1621,7 @@ mod test { .into_iter() .collect(); let outcome = state.add_from_cache(docs); assert!(outcome.unwrap().0); // successfully loaded one MD. assert!(outcome.unwrap().is_none()); // successfully loaded one MD. assert!(!state.can_advance()); assert!(!state.is_ready(Readiness::Complete)); assert!(!state.is_ready(Readiness::Usable)); Loading @@ -1655,7 +1645,7 @@ mod test { let req = ClientRequest::Microdescs(req); let source = DocSource::DirServer { source: None }; let outcome = state.add_from_download(response.as_str(), &req, source, Some(&store)); assert!(outcome.unwrap().0); // successfully loaded MDs assert!(outcome.unwrap().is_none()); // successfully loaded MDs match state.get_netdir_change().unwrap() { NetDirChange::AttemptReplace { netdir, .. } => { assert!(netdir.take().is_some()); Loading Loading
crates/tor-dirmgr/src/bootstrap.rs +26 −31 Original line number Diff line number Diff line Loading @@ -303,11 +303,11 @@ async fn fetch_multiple<R: Runtime>( async fn load_once<R: Runtime>( dirmgr: &Arc<DirMgr<R>>, state: &mut Box<dyn DirState>, ) -> Result<(bool, Option<Error>)> { ) -> Result<Option<Error>> { let missing = state.missing_docs(); let outcome: Result<(bool, Option<Error>)> = if missing.is_empty() { let outcome: Result<Option<Error>> = if missing.is_empty() { trace!("Found no missing documents; can't advance current state"); Ok((false, None)) Ok(Some(Error::NoChange(DocSource::LocalCache))) } else { trace!( "Found {} missing documents; trying to load them", Loading @@ -323,13 +323,13 @@ async fn load_once<R: Runtime>( Err(Error::UntimelyObject(TimeValidityError::Expired(_))) => { // This is just an expired object from the cache; we don't need // to call that an error. Treat it as if it were absent. Ok((false, None)) Ok(None) } other => other, } }; if matches!(outcome, Ok((true, _))) { if outcome.is_ok() { dirmgr.update_status(state.bootstrap_status()); } Loading @@ -347,13 +347,13 @@ pub(crate) async fn load<R: Runtime>( let mut safety_counter = 0_usize; loop { trace!(state=%state.describe(), "Loading from cache"); let (changed, nonfatal_err) = load_once(&dirmgr, &mut state).await?; let nonfatal_err = load_once(&dirmgr, &mut state).await?; { let mut store = dirmgr.store.lock().expect("store lock poisoned"); dirmgr.apply_netdir_changes(&mut state, store.deref_mut())?; } if let Some(e) = nonfatal_err { if let Some(e) = &nonfatal_err { debug!("Recoverable loading from cache: {}", e); } if let Some(e) = state.blocking_error() { Loading @@ -364,7 +364,8 @@ pub(crate) async fn load<R: Runtime>( state = state.advance()?; safety_counter = 0; } else { if !changed { if matches!(nonfatal_err, Some(Error::NoChange(_))) { // XXXX refactor more. break; } safety_counter += 1; Loading @@ -383,14 +384,11 @@ pub(crate) async fn load<R: Runtime>( /// /// This can launch one or more download requests, but will not launch more /// than `parallelism` requests at a time. /// /// Return true if the state reports that it changed. async fn download_attempt<R: Runtime>( dirmgr: &Arc<DirMgr<R>>, state: &mut Box<dyn DirState>, parallelism: usize, ) -> Result<bool> { let mut changed = false; ) -> Result<()> { let missing = state.missing_docs(); let fetched = fetch_multiple(Arc::clone(dirmgr), &missing, parallelism).await?; for (client_req, dir_response) in fetched { Loading @@ -411,9 +409,8 @@ async fn download_attempt<R: Runtime>( let doc_source = DocSource::DirServer { source: source.clone(), }; let (b, opt_error) = let opt_error = state.add_from_download(&text, &client_req, doc_source, Some(&dirmgr.store))?; changed |= b; if let Some(e) = &opt_error { warn!("error while adding directory info: {}", e); } Loading @@ -434,11 +431,9 @@ async fn download_attempt<R: Runtime>( } } if changed { dirmgr.update_status(state.bootstrap_status()); } Ok(changed) Ok(()) } /// Download information into a DirState state machine until it is Loading Loading @@ -534,15 +529,10 @@ pub(crate) async fn download<R: Runtime>( let dirmgr = upgrade_weak_ref(&dirmgr)?; futures::select_biased! { outcome = download_attempt(&dirmgr, &mut state, parallelism.into()).fuse() => { match outcome { Err(e) => { if let Err(e) = outcome { warn!("Error while downloading: {}", e); continue 'next_attempt; } Ok(changed) => { changed } } } _ = runtime.sleep_until_wallclock(reset_time).fuse() => { // We need to reset. This can happen if (for Loading Loading @@ -716,10 +706,7 @@ mod test { }) .collect() } fn add_from_cache( &mut self, docs: HashMap<DocId, DocumentText>, ) -> Result<(bool, Option<Error>)> { fn add_from_cache(&mut self, docs: HashMap<DocId, DocumentText>) -> Result<Option<Error>> { let mut changed = false; for id in docs.keys() { if let DocId::Microdesc(id) = id { Loading @@ -729,7 +716,11 @@ mod test { } } } Ok((changed, None)) if changed { Ok(None) } else { Ok(Some(Error::NoChange(DocSource::LocalCache))) } } fn add_from_download( &mut self, Loading @@ -737,7 +728,7 @@ mod test { _request: &ClientRequest, _source: DocSource, _storage: Option<&Mutex<DynStore>>, ) -> Result<(bool, Option<Error>)> { ) -> Result<Option<Error>> { let mut changed = false; for token in text.split_ascii_whitespace() { if let Ok(v) = hex::decode(token) { Loading @@ -749,7 +740,11 @@ mod test { } } } Ok((changed, None)) if changed { Ok(None) } else { Ok(Some(Error::NoChange(DocSource::LocalCache))) } } fn dl_config(&self) -> DownloadSchedule { DownloadSchedule::default() Loading
crates/tor-dirmgr/src/state.rs +43 −53 Original line number Diff line number Diff line Loading @@ -116,13 +116,9 @@ pub(crate) trait DirState: Send { /// Return an `Err` only on a fatal error that means we should stop /// downloading entirely; return recovered-from errors inside the `Ok()` /// variant. fn add_from_cache( &mut self, docs: HashMap<DocId, DocumentText>, ) -> Result<(bool, Option<Error>)>; fn add_from_cache(&mut self, docs: HashMap<DocId, DocumentText>) -> Result<Option<Error>>; /// Add information that we have just downloaded to this state; returns /// 'true' if there as any change in this state. /// Add information that we have just downloaded to this state. /// /// This method receives a copy of the original request, and should reject /// any documents that do not pertain to it. Loading @@ -139,7 +135,7 @@ pub(crate) trait DirState: Send { request: &ClientRequest, source: DocSource, storage: Option<&Mutex<DynStore>>, ) -> Result<(bool, Option<Error>)>; ) -> Result<Option<Error>>; /// Return a summary of this state as a [`DirStatus`]. fn bootstrap_status(&self) -> event::DirStatus; /// Return a configuration for attempting downloads. Loading Loading @@ -279,12 +275,9 @@ impl<R: Runtime> DirState for GetConsensusState<R> { fn dl_config(&self) -> DownloadSchedule { self.config.schedule.retry_consensus } fn add_from_cache( &mut self, docs: HashMap<DocId, DocumentText>, ) -> Result<(bool, Option<Error>)> { fn add_from_cache(&mut self, docs: HashMap<DocId, DocumentText>) -> Result<Option<Error>> { let text = match docs.into_iter().next() { None => return Ok((false, None)), None => return Ok(Some(Error::NoChange(DocSource::LocalCache))), Some(( DocId::LatestConsensus { flavor: ConsensusFlavor::Microdesc, Loading @@ -297,8 +290,8 @@ impl<R: Runtime> DirState for GetConsensusState<R> { let source = DocSource::LocalCache; match self.add_consensus_text(source, text.as_str().map_err(Error::BadUtf8InCache)?, None) { Ok(_) => Ok((true, None)), Err(e) => Ok((false, Some(e))), Ok(_) => Ok(None), Err(e) => Ok(Some(e)), } } fn add_from_download( Loading @@ -307,7 +300,7 @@ impl<R: Runtime> DirState for GetConsensusState<R> { request: &ClientRequest, source: DocSource, storage: Option<&Mutex<DynStore>>, ) -> Result<(bool, Option<Error>)> { ) -> Result<Option<Error>> { let requested_newer_than = match request { ClientRequest::Consensus(r) => r.last_consensus_date(), _ => None, Loading @@ -318,9 +311,9 @@ impl<R: Runtime> DirState for GetConsensusState<R> { let mut w = store.lock().expect("Directory storage lock poisoned"); w.store_consensus(meta, ConsensusFlavor::Microdesc, true, text)?; } Ok((true, None)) Ok(None) } Err(e) => Ok((false, Some(e))), Err(e) => Ok(Some(e)), } } fn advance(self: Box<Self>) -> Result<Box<dyn DirState>> { Loading Loading @@ -575,10 +568,7 @@ impl<R: Runtime> DirState for GetCertsState<R> { fn dl_config(&self) -> DownloadSchedule { self.config.schedule.retry_certs } fn add_from_cache( &mut self, docs: HashMap<DocId, DocumentText>, ) -> Result<(bool, Option<Error>)> { fn add_from_cache(&mut self, docs: HashMap<DocId, DocumentText>) -> Result<Option<Error>> { let mut changed = false; // Here we iterate over the documents we want, taking them from // our input and remembering them. Loading @@ -602,8 +592,10 @@ impl<R: Runtime> DirState for GetCertsState<R> { } if changed { self.try_checking_sigs(); Ok(nonfatal_error) } else { Ok(Some(Error::NoChange(DocSource::LocalCache))) } Ok((changed, nonfatal_error)) } fn add_from_download( &mut self, Loading @@ -611,7 +603,7 @@ impl<R: Runtime> DirState for GetCertsState<R> { request: &ClientRequest, source: DocSource, storage: Option<&Mutex<DynStore>>, ) -> Result<(bool, Option<Error>)> { ) -> Result<Option<Error>> { let asked_for: HashSet<_> = match request { ClientRequest::AuthCert(a) => a.keys().collect(), _ => return Err(internal!("expected an AuthCert request").into()), Loading Loading @@ -644,7 +636,7 @@ impl<R: Runtime> DirState for GetCertsState<R> { // We want to exit early if we aren't saving any certificates. if newcerts.is_empty() { return Ok((false, nonfatal_error)); return Ok(nonfatal_error.or(Some(Error::NoChange(source)))); } if let Some(store) = storage { Loading @@ -671,9 +663,10 @@ impl<R: Runtime> DirState for GetCertsState<R> { if changed { self.try_checking_sigs(); Ok(nonfatal_error) } else { Ok(Some(Error::NoChange(source))) } Ok((changed, nonfatal_error)) } fn advance(self: Box<Self>) -> Result<Box<dyn DirState>> { Loading Loading @@ -997,10 +990,7 @@ impl<R: Runtime> DirState for GetMicrodescsState<R> { fn dl_config(&self) -> DownloadSchedule { self.config.schedule.retry_microdescs } fn add_from_cache( &mut self, docs: HashMap<DocId, DocumentText>, ) -> Result<(bool, Option<Error>)> { fn add_from_cache(&mut self, docs: HashMap<DocId, DocumentText>) -> Result<Option<Error>> { let mut microdescs = Vec::new(); for (id, text) in docs { if let DocId::Microdesc(digest) = id { Loading @@ -1014,10 +1004,12 @@ impl<R: Runtime> DirState for GetMicrodescsState<R> { } } let changed = !microdescs.is_empty(); if microdescs.is_empty() { Ok(Some(Error::NoChange(DocSource::LocalCache))) } else { self.register_microdescs(microdescs); Ok((changed, None)) Ok(None) } } fn add_from_download( Loading @@ -1026,7 +1018,7 @@ impl<R: Runtime> DirState for GetMicrodescsState<R> { request: &ClientRequest, source: DocSource, storage: Option<&Mutex<DynStore>>, ) -> Result<(bool, Option<Error>)> { ) -> Result<Option<Error>> { let requested: HashSet<_> = if let ClientRequest::Microdescs(req) = request { req.digests().collect() } else { Loading Loading @@ -1079,8 +1071,12 @@ impl<R: Runtime> DirState for GetMicrodescsState<R> { )?; } } if new_mds.is_empty() { Ok(Some(Error::NoChange(source))) } else { self.register_microdescs(new_mds.into_iter().map(|(_, md)| md)); Ok((true, nonfatal_err)) Ok(nonfatal_err) } } fn advance(self: Box<Self>) -> Result<Box<dyn DirState>> { Ok(self) Loading Loading @@ -1358,10 +1354,7 @@ mod test { source.clone(), Some(&store), ); assert!(matches!( outcome, Ok((false, Some(Error::NetDocError { .. }))) )); assert!(matches!(outcome, Ok(Some(Error::NetDocError { .. })))); // make sure it wasn't stored... assert!(store .lock() Loading @@ -1372,10 +1365,7 @@ mod test { // Now try again, with a real consensus... but the wrong authorities. let outcome = state.add_from_download(CONSENSUS, &req, source.clone(), Some(&store)); assert!(matches!( outcome, Ok((false, Some(Error::UnrecognizedAuthorities))) )); assert!(matches!(outcome, Ok(Some(Error::UnrecognizedAuthorities)))); assert!(store .lock() .unwrap() Loading @@ -1396,7 +1386,7 @@ mod test { Arc::new(crate::filter::NilFilter), ); let outcome = state.add_from_download(CONSENSUS, &req, source, Some(&store)); assert!(outcome.unwrap().0); assert!(outcome.unwrap().is_none()); assert!(store .lock() .unwrap() Loading Loading @@ -1427,7 +1417,7 @@ mod test { let text: crate::storage::InputString = CONSENSUS.to_owned().into(); let map = vec![(docid, text.into())].into_iter().collect(); let outcome = state.add_from_cache(map); assert!(outcome.unwrap().0); assert!(outcome.unwrap().is_none()); assert!(state.can_advance()); }); } Loading @@ -1451,7 +1441,7 @@ mod test { let req = tor_dirclient::request::ConsensusRequest::new(ConsensusFlavor::Microdesc); let req = crate::docid::ClientRequest::Consensus(req); let outcome = state.add_from_download(CONSENSUS, &req, source, None); assert!(outcome.unwrap().0); assert!(outcome.unwrap().is_none()); Box::new(state).advance().unwrap() } Loading Loading @@ -1495,7 +1485,7 @@ mod test { .into_iter() .collect(); let outcome = state.add_from_cache(docs); assert!(outcome.unwrap().0); // no error, and something changed. assert!(outcome.unwrap().is_none()); // no error, and something changed. assert!(!state.can_advance()); // But we aren't done yet. let missing = state.missing_docs(); assert_eq!(missing.len(), 1); // Now we're only missing one! Loading @@ -1513,7 +1503,7 @@ mod test { let req = ClientRequest::AuthCert(req); let outcome = state.add_from_download(AUTHCERT_5A23, &req, source.clone(), Some(&store)); assert!(!outcome.unwrap().0); // no error, but nothing changed. assert!(matches!(outcome.unwrap(), Some(Error::Unwanted(_)))); let missing2 = state.missing_docs(); assert_eq!(missing, missing2); // No change. assert!(store Loading @@ -1528,7 +1518,7 @@ mod test { req.push(authcert_id_5a23()); // Right idea this time! let req = ClientRequest::AuthCert(req); let outcome = state.add_from_download(AUTHCERT_5A23, &req, source, Some(&store)); assert!(outcome.unwrap().0); // No error, _and_ something changed! assert!(outcome.unwrap().is_none()); // No error, _and_ something changed! let missing3 = state.missing_docs(); assert!(missing3.is_empty()); assert!(state.can_advance()); Loading Loading @@ -1631,7 +1621,7 @@ mod test { .into_iter() .collect(); let outcome = state.add_from_cache(docs); assert!(outcome.unwrap().0); // successfully loaded one MD. assert!(outcome.unwrap().is_none()); // successfully loaded one MD. assert!(!state.can_advance()); assert!(!state.is_ready(Readiness::Complete)); assert!(!state.is_ready(Readiness::Usable)); Loading @@ -1655,7 +1645,7 @@ mod test { let req = ClientRequest::Microdescs(req); let source = DocSource::DirServer { source: None }; let outcome = state.add_from_download(response.as_str(), &req, source, Some(&store)); assert!(outcome.unwrap().0); // successfully loaded MDs assert!(outcome.unwrap().is_none()); // successfully loaded MDs match state.get_netdir_change().unwrap() { NetDirChange::AttemptReplace { netdir, .. } => { assert!(netdir.take().is_some()); Loading