fetch: Prune remotes with sigrefs failures
Previously, trying to load `SignedRefs` for any given remote would result in a fetch failure. Teach the fetch to be more resilient by pruning remotes that result in an error when loading `SignedRefs`, and add the error to the validation failures.
This commit is contained in:
parent
6967bf8fca
commit
725ced09d5
|
|
@ -1,5 +1,5 @@
|
||||||
use std::collections::{BTreeMap, BTreeSet};
|
use std::collections::{BTreeMap, BTreeSet};
|
||||||
use std::ops::{Deref, Not as _};
|
use std::ops::Not as _;
|
||||||
|
|
||||||
use radicle::storage::git::Repository;
|
use radicle::storage::git::Repository;
|
||||||
pub use radicle::storage::refs::SignedRefsAt;
|
pub use radicle::storage::refs::SignedRefsAt;
|
||||||
|
|
@ -36,13 +36,6 @@ pub(crate) enum DelegateStatus<T = ()> {
|
||||||
NonDelegate { remote: PublicKey, data: T },
|
NonDelegate { remote: PublicKey, data: T },
|
||||||
}
|
}
|
||||||
|
|
||||||
impl DelegateStatus {
|
|
||||||
/// Construct a `DelegateStatus` without any data.
|
|
||||||
pub fn empty(remote: PublicKey, delegates: &BTreeSet<PublicKey>) -> Self {
|
|
||||||
Self::new((), remote, delegates)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
impl<T> DelegateStatus<T> {
|
impl<T> DelegateStatus<T> {
|
||||||
pub fn new(data: T, remote: PublicKey, delegates: &BTreeSet<PublicKey>) -> Self {
|
pub fn new(data: T, remote: PublicKey, delegates: &BTreeSet<PublicKey>) -> Self {
|
||||||
if delegates.contains(&remote) {
|
if delegates.contains(&remote) {
|
||||||
|
|
@ -51,42 +44,6 @@ impl<T> DelegateStatus<T> {
|
||||||
Self::NonDelegate { remote, data }
|
Self::NonDelegate { remote, data }
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Construct a `DelegateStatus` with [`SignedRefsAt`] signed reference
|
|
||||||
/// data, if it can be found in `repo`.
|
|
||||||
pub fn load<R, S>(
|
|
||||||
self,
|
|
||||||
cached: &Cached<R, S>,
|
|
||||||
) -> Result<
|
|
||||||
DelegateStatus<Option<SignedRefsAt>>,
|
|
||||||
radicle::storage::refs::sigrefs::read::error::Read,
|
|
||||||
>
|
|
||||||
where
|
|
||||||
R: AsRef<Repository>,
|
|
||||||
{
|
|
||||||
let remote = *self.remote();
|
|
||||||
self.traverse(|_| cached.load(&remote))
|
|
||||||
}
|
|
||||||
|
|
||||||
fn remote(&self) -> &PublicKey {
|
|
||||||
match self {
|
|
||||||
Self::Delegate { remote, .. } => remote,
|
|
||||||
Self::NonDelegate { remote, .. } => remote,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn traverse<U, E>(self, f: impl FnOnce(T) -> Result<U, E>) -> Result<DelegateStatus<U>, E> {
|
|
||||||
match self {
|
|
||||||
Self::Delegate { remote, data } => Ok(DelegateStatus::Delegate {
|
|
||||||
remote,
|
|
||||||
data: f(data)?,
|
|
||||||
}),
|
|
||||||
Self::NonDelegate { remote, data } => Ok(DelegateStatus::NonDelegate {
|
|
||||||
remote,
|
|
||||||
data: f(data)?,
|
|
||||||
}),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) fn validate(
|
pub(crate) fn validate(
|
||||||
|
|
@ -102,45 +59,52 @@ pub(crate) fn validate(
|
||||||
///
|
///
|
||||||
/// Construct using [`RemoteRefs::load`].
|
/// Construct using [`RemoteRefs::load`].
|
||||||
#[derive(Debug, Default)]
|
#[derive(Debug, Default)]
|
||||||
pub struct RemoteRefs(BTreeMap<PublicKey, SignedRefsAt>);
|
pub struct RemoteRefs(
|
||||||
|
pub(super) BTreeMap<
|
||||||
|
PublicKey,
|
||||||
|
Result<Option<SignedRefsAt>, radicle::storage::refs::sigrefs::read::error::Read>,
|
||||||
|
>,
|
||||||
|
);
|
||||||
|
|
||||||
impl RemoteRefs {
|
impl RemoteRefs {
|
||||||
/// Load the sigrefs for each remote in `remotes`.
|
/// Load the sigrefs for each remote in `remotes`.
|
||||||
///
|
|
||||||
/// If the sigrefs are missing for a given remote, regardless of delegate
|
|
||||||
/// status, then that remote is filtered out.
|
|
||||||
pub(crate) fn load<'a, R, S>(
|
pub(crate) fn load<'a, R, S>(
|
||||||
cached: &Cached<R, S>,
|
cached: &Cached<R, S>,
|
||||||
remotes: impl Iterator<Item = &'a PublicKey>,
|
remotes: impl Iterator<Item = &'a PublicKey>,
|
||||||
) -> Result<Self, error::RemoteRefs>
|
) -> Self
|
||||||
where
|
where
|
||||||
R: AsRef<Repository>,
|
R: AsRef<Repository>,
|
||||||
{
|
{
|
||||||
remotes
|
Self(
|
||||||
.filter_map(|id| match cached.load(id) {
|
remotes
|
||||||
Ok(None) => None,
|
.map(|remote| (*remote, cached.load(remote)))
|
||||||
Ok(Some(sr)) => Some(Ok((id, sr))),
|
.collect(),
|
||||||
Err(e) => Some(Err(e)),
|
)
|
||||||
})
|
|
||||||
.try_fold(RemoteRefs::default(), |mut acc, remote_refs| {
|
|
||||||
let (id, sigrefs) = remote_refs?;
|
|
||||||
acc.0.insert(*id, sigrefs);
|
|
||||||
Ok(acc)
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
impl Deref for RemoteRefs {
|
pub(crate) fn len(&self) -> usize {
|
||||||
type Target = BTreeMap<PublicKey, SignedRefsAt>;
|
self.0.len()
|
||||||
|
}
|
||||||
|
|
||||||
fn deref(&self) -> &Self::Target {
|
pub(crate) fn into_inner(
|
||||||
&self.0
|
self,
|
||||||
|
) -> BTreeMap<
|
||||||
|
PublicKey,
|
||||||
|
Result<Option<SignedRefsAt>, radicle::storage::refs::sigrefs::read::error::Read>,
|
||||||
|
> {
|
||||||
|
self.0
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<'a> IntoIterator for &'a RemoteRefs {
|
impl<'a> IntoIterator for &'a RemoteRefs {
|
||||||
type Item = <&'a BTreeMap<PublicKey, SignedRefsAt> as IntoIterator>::Item;
|
type Item = <&'a BTreeMap<
|
||||||
type IntoIter = <&'a BTreeMap<PublicKey, SignedRefsAt> as IntoIterator>::IntoIter;
|
PublicKey,
|
||||||
|
Result<Option<SignedRefsAt>, radicle::storage::refs::sigrefs::read::error::Read>,
|
||||||
|
> as IntoIterator>::Item;
|
||||||
|
type IntoIter = <&'a BTreeMap<
|
||||||
|
PublicKey,
|
||||||
|
Result<Option<SignedRefsAt>, radicle::storage::refs::sigrefs::read::error::Read>,
|
||||||
|
> as IntoIterator>::IntoIter;
|
||||||
|
|
||||||
fn into_iter(self) -> Self::IntoIter {
|
fn into_iter(self) -> Self::IntoIter {
|
||||||
self.0.iter()
|
self.0.iter()
|
||||||
|
|
|
||||||
|
|
@ -45,7 +45,7 @@ use radicle::storage::ReadRepository;
|
||||||
use crate::git::refs::{Policy, Update, Updates};
|
use crate::git::refs::{Policy, Update, Updates};
|
||||||
use crate::policy::BlockList;
|
use crate::policy::BlockList;
|
||||||
use crate::refs::{ReceivedRef, ReceivedRefname};
|
use crate::refs::{ReceivedRef, ReceivedRefname};
|
||||||
use crate::sigrefs;
|
use crate::sigrefs::RemoteRefs;
|
||||||
use crate::state::FetchState;
|
use crate::state::FetchState;
|
||||||
use crate::transport::WantsHaves;
|
use crate::transport::WantsHaves;
|
||||||
use crate::{policy, refs};
|
use crate::{policy, refs};
|
||||||
|
|
@ -482,12 +482,18 @@ pub struct DataRefs {
|
||||||
pub remote: PublicKey,
|
pub remote: PublicKey,
|
||||||
/// The set of signed references from each remote that was
|
/// The set of signed references from each remote that was
|
||||||
/// fetched.
|
/// fetched.
|
||||||
pub remotes: sigrefs::RemoteRefs,
|
pub remotes: RemoteRefs,
|
||||||
/// The data limit for this stage of fetching.
|
/// The data limit for this stage of fetching.
|
||||||
#[allow(dead_code)]
|
#[allow(dead_code)]
|
||||||
pub limit: u64,
|
pub limit: u64,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
impl DataRefs {
|
||||||
|
pub(crate) fn into_remote_refs(self) -> RemoteRefs {
|
||||||
|
self.remotes
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
impl ProtocolStage for DataRefs {
|
impl ProtocolStage for DataRefs {
|
||||||
// We don't need to ask for refs since we have all reference names
|
// We don't need to ask for refs since we have all reference names
|
||||||
// and `Oid`s in `rad/sigrefs`.
|
// and `Oid`s in `rad/sigrefs`.
|
||||||
|
|
@ -514,10 +520,13 @@ impl ProtocolStage for DataRefs {
|
||||||
) -> Result<WantsHaves, error::WantsHaves> {
|
) -> Result<WantsHaves, error::WantsHaves> {
|
||||||
let mut wants_haves = WantsHaves::default();
|
let mut wants_haves = WantsHaves::default();
|
||||||
|
|
||||||
for (remote, loaded) in &self.remotes {
|
for (remote, result) in self.remotes.into_iter() {
|
||||||
|
let Ok(Some(refs)) = result else {
|
||||||
|
continue;
|
||||||
|
};
|
||||||
wants_haves.add(
|
wants_haves.add(
|
||||||
refdb,
|
refdb,
|
||||||
loaded.refs().iter().filter_map(|(refname, tip)| {
|
refs.iter().filter_map(|(refname, tip)| {
|
||||||
let refname = Qualified::from_refstr(refname)
|
let refname = Qualified::from_refstr(refname)
|
||||||
.map(|refname| refname.with_namespace(Component::from(remote)))?;
|
.map(|refname| refname.with_namespace(Component::from(remote)))?;
|
||||||
Some((refname, *tip))
|
Some((refname, *tip))
|
||||||
|
|
@ -536,7 +545,10 @@ impl ProtocolStage for DataRefs {
|
||||||
) -> Result<Updates<'a>, error::Prepare> {
|
) -> Result<Updates<'a>, error::Prepare> {
|
||||||
let mut updates = Updates::default();
|
let mut updates = Updates::default();
|
||||||
|
|
||||||
for (remote, refs) in &self.remotes {
|
for (remote, result) in &self.remotes {
|
||||||
|
let Ok(Some(refs)) = result else {
|
||||||
|
continue;
|
||||||
|
};
|
||||||
let mut signed = HashSet::with_capacity(refs.refs().len());
|
let mut signed = HashSet::with_capacity(refs.refs().len());
|
||||||
for (name, tip) in refs.iter() {
|
for (name, tip) in refs.iter() {
|
||||||
let tracking: Namespaced<'_> = Qualified::from_refstr(name)
|
let tracking: Namespaced<'_> = Qualified::from_refstr(name)
|
||||||
|
|
|
||||||
|
|
@ -313,7 +313,7 @@ impl FetchState {
|
||||||
self.run_stage(handle, handshake, &sigrefs_at)?;
|
self.run_stage(handle, handshake, &sigrefs_at)?;
|
||||||
let remotes = refs_at.iter().map(|r| &r.remote);
|
let remotes = refs_at.iter().map(|r| &r.remote);
|
||||||
|
|
||||||
let signed_refs = sigrefs::RemoteRefs::load(&self.as_cached(handle), remotes)?;
|
let signed_refs = sigrefs::RemoteRefs::load(&self.as_cached(handle), remotes);
|
||||||
Ok(signed_refs)
|
Ok(signed_refs)
|
||||||
}
|
}
|
||||||
None => {
|
None => {
|
||||||
|
|
@ -333,7 +333,7 @@ impl FetchState {
|
||||||
let signed_refs = sigrefs::RemoteRefs::load(
|
let signed_refs = sigrefs::RemoteRefs::load(
|
||||||
&self.as_cached(handle),
|
&self.as_cached(handle),
|
||||||
fetched.iter().chain(delegates.iter()),
|
fetched.iter().chain(delegates.iter()),
|
||||||
)?;
|
);
|
||||||
Ok(signed_refs)
|
Ok(signed_refs)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -418,6 +418,7 @@ impl FetchState {
|
||||||
remote,
|
remote,
|
||||||
refs_at,
|
refs_at,
|
||||||
)?;
|
)?;
|
||||||
|
|
||||||
log::debug!(
|
log::debug!(
|
||||||
"Fetched data for {} remote(s) ({}ms)",
|
"Fetched data for {} remote(s) ({}ms)",
|
||||||
signed_refs.len(),
|
signed_refs.len(),
|
||||||
|
|
@ -429,10 +430,10 @@ impl FetchState {
|
||||||
remotes: signed_refs,
|
remotes: signed_refs,
|
||||||
limit: limit.refs,
|
limit: limit.refs,
|
||||||
};
|
};
|
||||||
self.run_stage(handle, handshake, &data_refs)?;
|
let fetched = self.run_stage(handle, handshake, &data_refs)?;
|
||||||
log::debug!(
|
log::debug!(
|
||||||
"Fetched data refs for {} remotes ({}ms)",
|
"Fetched data refs for {} remotes ({}ms)",
|
||||||
data_refs.remotes.len(),
|
fetched.len(),
|
||||||
start.elapsed().as_millis()
|
start.elapsed().as_millis()
|
||||||
);
|
);
|
||||||
|
|
||||||
|
|
@ -450,7 +451,8 @@ impl FetchState {
|
||||||
// remotes from the tips, thus not updating the production Git
|
// remotes from the tips, thus not updating the production Git
|
||||||
// repository.
|
// repository.
|
||||||
let mut failures = sigrefs::Validations::default();
|
let mut failures = sigrefs::Validations::default();
|
||||||
let signed_refs = data_refs.remotes;
|
|
||||||
|
let signed_refs = data_refs.into_remote_refs();
|
||||||
|
|
||||||
// We may prune fetched remotes, so we keep track of
|
// We may prune fetched remotes, so we keep track of
|
||||||
// non-pruned, fetched remotes here.
|
// non-pruned, fetched remotes here.
|
||||||
|
|
@ -469,21 +471,26 @@ impl FetchState {
|
||||||
|
|
||||||
// TODO(finto): this might read better if it got its own
|
// TODO(finto): this might read better if it got its own
|
||||||
// private function.
|
// private function.
|
||||||
for remote in signed_refs.keys() {
|
for (remote, refs) in signed_refs.into_inner() {
|
||||||
if handle.is_blocked(remote) {
|
if handle.is_blocked(&remote) {
|
||||||
log::trace!("Skipping blocked remote {remote}");
|
log::trace!("Skipping blocked remote {remote}");
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
let remote = sigrefs::DelegateStatus::empty(*remote, &delegates)
|
let remote = sigrefs::DelegateStatus::new(refs, remote, &delegates);
|
||||||
.load(&self.as_cached(handle))?;
|
|
||||||
match remote {
|
match remote {
|
||||||
sigrefs::DelegateStatus::NonDelegate { remote, data: None } => {
|
sigrefs::DelegateStatus::NonDelegate {
|
||||||
|
remote,
|
||||||
|
data: Ok(None),
|
||||||
|
} => {
|
||||||
log::debug!("Pruning non-delegate {remote} tips, missing 'rad/sigrefs'");
|
log::debug!("Pruning non-delegate {remote} tips, missing 'rad/sigrefs'");
|
||||||
failures.push(sigrefs::Validation::MissingRadSigRefs(remote));
|
failures.push(sigrefs::Validation::MissingRadSigRefs(remote));
|
||||||
self.prune(&remote);
|
self.prune(&remote);
|
||||||
}
|
}
|
||||||
sigrefs::DelegateStatus::Delegate { remote, data: None } => {
|
sigrefs::DelegateStatus::Delegate {
|
||||||
|
remote,
|
||||||
|
data: Ok(None),
|
||||||
|
} => {
|
||||||
log::debug!("Pruning delegate {remote} tips, missing 'rad/sigrefs'");
|
log::debug!("Pruning delegate {remote} tips, missing 'rad/sigrefs'");
|
||||||
failures.push(sigrefs::Validation::MissingRadSigRefs(remote));
|
failures.push(sigrefs::Validation::MissingRadSigRefs(remote));
|
||||||
self.prune(&remote);
|
self.prune(&remote);
|
||||||
|
|
@ -495,9 +502,26 @@ impl FetchState {
|
||||||
valid_delegates.remove(&remote);
|
valid_delegates.remove(&remote);
|
||||||
failed_delegates.insert(remote);
|
failed_delegates.insert(remote);
|
||||||
}
|
}
|
||||||
|
sigrefs::DelegateStatus::Delegate {
|
||||||
|
remote,
|
||||||
|
data: Err(err),
|
||||||
|
}
|
||||||
|
| sigrefs::DelegateStatus::NonDelegate {
|
||||||
|
remote,
|
||||||
|
data: Err(err),
|
||||||
|
} => {
|
||||||
|
log::debug!("Pruning {remote} tips due to: {err}");
|
||||||
|
self.prune(&remote);
|
||||||
|
valid_delegates.remove(&remote);
|
||||||
|
failed_delegates.insert(remote);
|
||||||
|
failures.push(sigrefs::Validation::Read {
|
||||||
|
remote,
|
||||||
|
source: err,
|
||||||
|
});
|
||||||
|
}
|
||||||
sigrefs::DelegateStatus::NonDelegate {
|
sigrefs::DelegateStatus::NonDelegate {
|
||||||
remote,
|
remote,
|
||||||
data: Some(sigrefs),
|
data: Ok(Some(sigrefs)),
|
||||||
} => {
|
} => {
|
||||||
if let Some(SignedRefsAt { at, .. }) =
|
if let Some(SignedRefsAt { at, .. }) =
|
||||||
SignedRefsAt::load(remote, handle.repository())?
|
SignedRefsAt::load(remote, handle.repository())?
|
||||||
|
|
@ -527,7 +551,7 @@ impl FetchState {
|
||||||
}
|
}
|
||||||
sigrefs::DelegateStatus::Delegate {
|
sigrefs::DelegateStatus::Delegate {
|
||||||
remote,
|
remote,
|
||||||
data: Some(sigrefs),
|
data: Ok(Some(sigrefs)),
|
||||||
} => {
|
} => {
|
||||||
if let Some(SignedRefsAt { at, .. }) =
|
if let Some(SignedRefsAt { at, .. }) =
|
||||||
SignedRefsAt::load(remote, handle.repository())?
|
SignedRefsAt::load(remote, handle.repository())?
|
||||||
|
|
|
||||||
|
|
@ -410,6 +410,12 @@ pub enum Validation {
|
||||||
},
|
},
|
||||||
#[error("missing `refs/namespaces/{0}/refs/rad/sigrefs`")]
|
#[error("missing `refs/namespaces/{0}/refs/rad/sigrefs`")]
|
||||||
MissingRadSigRefs(RemoteId),
|
MissingRadSigRefs(RemoteId),
|
||||||
|
#[error("failed to read `refs/namespaces/{remote}/refs/rad/sigrefs`: {source}")]
|
||||||
|
Read {
|
||||||
|
remote: RemoteId,
|
||||||
|
#[source]
|
||||||
|
source: crate::storage::refs::sigrefs::read::error::Read,
|
||||||
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Repository {
|
impl Repository {
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue