use std::{collections::HashSet, sync::Arc};
use zksync_concurrency::{sync, time};
use zksync_consensus_roles::validator;
use crate::watch::Watch;
#[derive(Clone, Default, PartialEq, Eq)]
pub(crate) struct ValidatorAddrs(
pub(super) im::HashMap<validator::PublicKey, Arc<validator::Signed<validator::NetAddress>>>,
);
impl ValidatorAddrs {
pub(crate) fn get(
&self,
key: &validator::PublicKey,
) -> Option<&Arc<validator::Signed<validator::NetAddress>>> {
self.0.get(key)
}
pub(super) fn get_newer(&self, b: &Self) -> Vec<Arc<validator::Signed<validator::NetAddress>>> {
let mut newer = vec![];
for (k, v) in &self.0 {
if let Some(bv) = b.0.get(k) {
if !v.msg.is_newer(&bv.msg) {
continue;
}
}
newer.push(v.clone());
}
newer
}
pub(super) fn update(
&mut self,
validators: &validator::Schedule,
data: &[Arc<validator::Signed<validator::NetAddress>>],
) -> anyhow::Result<bool> {
let mut changed = false;
let mut done = HashSet::new();
for d in data {
if done.contains(&d.key) {
anyhow::bail!("duplicate entry for {:?}", d.key);
}
done.insert(d.key.clone());
if !validators.contains(&d.key) {
continue;
}
if let Some(x) = self.0.get(&d.key) {
if !d.msg.is_newer(&x.msg) {
continue;
}
}
d.verify()?;
self.0.insert(d.key.clone(), d.clone());
changed = true;
}
Ok(changed)
}
}
pub(crate) struct ValidatorAddrsWatch(Watch<ValidatorAddrs>);
impl Default for ValidatorAddrsWatch {
fn default() -> Self {
Self(Watch::new(ValidatorAddrs::default()))
}
}
impl ValidatorAddrsWatch {
pub(crate) fn subscribe(&self) -> sync::watch::Receiver<ValidatorAddrs> {
self.0.subscribe()
}
pub(crate) fn current(
&self,
) -> im::HashMap<validator::PublicKey, Arc<validator::Signed<validator::NetAddress>>> {
self.0.subscribe().borrow().0.clone()
}
pub(crate) async fn announce(
&self,
key: &validator::SecretKey,
addr: std::net::SocketAddr,
timestamp: time::Utc,
) {
let this = self.0.lock().await;
let mut validator_addrs = this.borrow().clone();
let version = validator_addrs
.get(&key.public())
.map(|x| x.msg.version + 1)
.unwrap_or(0);
let d = Arc::new(key.sign_msg(validator::NetAddress {
addr,
version,
timestamp,
}));
validator_addrs.0.insert(d.key.clone(), d);
this.send_replace(validator_addrs);
}
pub(crate) async fn update(
&self,
validators: &validator::Schedule,
data: &[Arc<validator::Signed<validator::NetAddress>>],
) -> anyhow::Result<()> {
let this = self.0.lock().await;
let mut validator_addrs = this.borrow().clone();
if validator_addrs.update(validators, data)? {
this.send_replace(validator_addrs);
}
Ok(())
}
}