use snarkvm::{
ledger::{
committee::Committee,
narwhal::{BatchCertificate, BatchHeader, Subdag},
},
prelude::{Address, Field, Network, cfg_iter},
};
use crate::helpers::now;
use indexmap::{IndexMap, IndexSet};
#[cfg(not(feature = "serial"))]
use rayon::prelude::*;
use std::{
collections::BTreeMap,
sync::{
Arc,
atomic::{AtomicI64, AtomicU64, Ordering},
},
};
use tokio::sync::{
mpsc::{self, error::TrySendError},
oneshot,
watch,
};
const TELEMETRY_QUEUE_CAPACITY: usize = 1024;
type ParticipationScores = (f64, f64, f64);
const DROPPED_WARNING_INTERVAL_IN_SECS: i64 = 10;
#[derive(Clone, Debug)]
struct ScoreSnapshot<N: Network> {
scores: IndexMap<Address<N>, ParticipationScores>,
gc_round: u64,
}
impl<N: Network> Default for ScoreSnapshot<N> {
fn default() -> Self {
Self { scores: Default::default(), gc_round: 0 }
}
}
#[derive(Clone, Debug)]
pub struct CertificateMetadata<N: Network> {
round: u64,
id: Field<N>,
author: Address<N>,
signers: Vec<Address<N>>,
}
impl<N: Network> CertificateMetadata<N> {
fn new(certificate: &BatchCertificate<N>) -> Self {
let author = certificate.author();
let signers = [author]
.into_iter()
.chain(certificate.signatures().map(|signature| signature.to_address()))
.collect::<Vec<_>>();
Self { round: certificate.round(), id: certificate.id(), author, signers }
}
}
#[derive(Debug)]
enum TelemetryUpdate<N: Network> {
Subdag { gc_round: u64, metadata: Vec<CertificateMetadata<N>> },
#[cfg_attr(not(test), allow(dead_code))]
Certificate(Box<CertificateMetadata<N>>),
Flush(oneshot::Sender<()>),
}
#[derive(Clone, Debug)]
pub struct Telemetry<N: Network> {
sender: mpsc::Sender<TelemetryUpdate<N>>,
scores: watch::Receiver<Arc<ScoreSnapshot<N>>>,
num_dropped: Arc<AtomicU64>,
max_dropped_round: Arc<AtomicU64>,
last_dropped_warning: Arc<AtomicI64>,
}
impl<N: Network> Telemetry<N> {
pub fn new() -> (Self, TelemetryWorker<N>) {
let (sender, receiver) = mpsc::channel(TELEMETRY_QUEUE_CAPACITY);
let (score_sender, score_receiver) = watch::channel::<Arc<ScoreSnapshot<N>>>(Default::default());
let telemetry = Self {
sender,
scores: score_receiver,
num_dropped: Default::default(),
max_dropped_round: Default::default(),
last_dropped_warning: Default::default(),
};
let worker = TelemetryWorker { receiver, scores: score_sender, state: TelemetryState::new(), gc_round: 0 };
(telemetry, worker)
}
pub fn insert_subdag(&self, subdag: &Subdag<N>) {
let anchor_round = subdag.anchor_round();
let Some(permit) = self.reserve() else {
self.max_dropped_round.fetch_max(anchor_round, Ordering::Relaxed);
return;
};
let gc_round = anchor_round.saturating_sub(BatchHeader::<N>::MAX_GC_ROUNDS as u64);
let certificates: Vec<_> = subdag.values().flatten().collect();
let metadata: Vec<_> =
cfg_iter!(certificates).map(|certificate| CertificateMetadata::new(certificate)).collect();
permit.send(TelemetryUpdate::Subdag { gc_round, metadata });
}
#[cfg(test)]
pub fn insert_certificate(&self, certificate: &BatchCertificate<N>) {
let Some(permit) = self.reserve() else {
self.max_dropped_round.fetch_max(certificate.round(), Ordering::Relaxed);
return;
};
permit.send(TelemetryUpdate::Certificate(Box::new(CertificateMetadata::new(certificate))));
}
pub fn get_participation_scores(&self, committee: &Committee<N>) -> IndexMap<Address<N>, (f64, f64)> {
let snapshot: Arc<ScoreSnapshot<N>> = self.scores.borrow().clone();
scores_for_committee(&snapshot.scores, committee)
}
#[cfg(test)]
pub fn is_complete(&self) -> bool {
let gc_round = self.scores.borrow().gc_round;
let max_dropped_round = self.max_dropped_round.load(Ordering::Relaxed);
if max_dropped_round == 0 {
return true;
}
gc_round >= max_dropped_round
}
pub fn published_gc_round(&self) -> u64 {
self.scores.borrow().gc_round
}
pub fn num_dropped(&self) -> u64 {
self.num_dropped.load(Ordering::Relaxed)
}
pub async fn flush(&self) -> Result<(), ()> {
let (sender, receiver) = oneshot::channel();
self.sender.send(TelemetryUpdate::Flush(sender)).await.map_err(|_| ())?;
receiver.await.map_err(|_| ())
}
fn reserve(&self) -> Option<mpsc::Permit<'_, TelemetryUpdate<N>>> {
match self.sender.try_reserve() {
Ok(permit) => Some(permit),
Err(TrySendError::Full(())) => {
let num_dropped = self.num_dropped.fetch_add(1, Ordering::Relaxed).saturating_add(1);
let now = now();
let last = self.last_dropped_warning.load(Ordering::Relaxed);
if now.saturating_sub(last) >= DROPPED_WARNING_INTERVAL_IN_SECS
&& self
.last_dropped_warning
.compare_exchange(last, now, Ordering::Relaxed, Ordering::Relaxed)
.is_ok()
{
warn!("Telemetry queue is full - dropping updates ({num_dropped} dropped in total)");
}
None
}
Err(TrySendError::Closed(())) => {
trace!("Telemetry worker is not running - dropping an update");
None
}
}
}
}
#[derive(Debug)]
pub struct TelemetryWorker<N: Network> {
receiver: mpsc::Receiver<TelemetryUpdate<N>>,
scores: watch::Sender<Arc<ScoreSnapshot<N>>>,
state: TelemetryState<N>,
gc_round: u64,
}
impl<N: Network> TelemetryWorker<N> {
pub async fn run(mut self) {
debug!("Starting the validator telemetry worker...");
while let Some(update) = self.receiver.recv().await {
let mut recompute = false;
let mut acks = Vec::new();
let mut next = Some(update);
while let Some(update) = next {
match update {
TelemetryUpdate::Subdag { gc_round, metadata } => {
self.state.garbage_collect_certificates(gc_round);
self.state.insert_certificate_metadata(&metadata);
self.gc_round = self.gc_round.max(gc_round);
recompute = true;
}
TelemetryUpdate::Certificate(metadata) => {
self.state.insert_certificate_metadata(std::slice::from_ref(&*metadata));
}
TelemetryUpdate::Flush(ack) => acks.push(ack),
}
next = self.receiver.try_recv().ok();
}
if recompute {
self.state.update_participation_scores();
self.scores.send_replace(Arc::new(ScoreSnapshot {
scores: self.state.participation_scores.clone(),
gc_round: self.gc_round,
}));
}
for ack in acks {
let _ = ack.send(());
}
}
debug!("The validator telemetry worker has stopped");
}
}
#[derive(Clone, Debug)]
pub struct TelemetryState<N: Network> {
tracked_certificates: BTreeMap<u64, IndexSet<Field<N>>>,
validator_signatures: IndexMap<Address<N>, IndexMap<u64, u32>>,
validator_certificates: IndexMap<Address<N>, IndexSet<u64>>,
participation_scores: IndexMap<Address<N>, ParticipationScores>,
}
impl<N: Network> Default for TelemetryState<N> {
fn default() -> Self {
Self::new()
}
}
impl<N: Network> TelemetryState<N> {
pub fn new() -> Self {
Self {
tracked_certificates: Default::default(),
validator_signatures: Default::default(),
validator_certificates: Default::default(),
participation_scores: Default::default(),
}
}
pub fn participation_scores(&self) -> &IndexMap<Address<N>, ParticipationScores> {
&self.participation_scores
}
pub fn insert_certificate(&mut self, certificate: &BatchCertificate<N>) {
self.insert_certificate_metadata(&[CertificateMetadata::new(certificate)]);
}
pub fn insert_certificate_metadata(&mut self, metadata: &[CertificateMetadata<N>]) {
for metadata in metadata {
if !self.tracked_certificates.entry(metadata.round).or_default().insert(metadata.id) {
continue;
}
for address in &metadata.signers {
self.validator_signatures
.entry(*address)
.or_default()
.entry(metadata.round)
.and_modify(|count| *count += 1)
.or_insert(1);
}
self.validator_certificates.entry(metadata.author).or_default().insert(metadata.round);
}
}
pub fn update_participation_scores(&mut self) {
fn weighted_score(certificate_score: f64, signature_score: f64) -> f64 {
let score = (0.9 * certificate_score) + (0.1 * signature_score);
(score * 100.0).round() / 100.0
}
let total_certificates = self.validator_certificates.values().map(|rounds| rounds.len()).sum::<usize>();
let signature_participation_scores: IndexMap<_, _> = self
.validator_signatures
.iter()
.map(|(address, signatures)| {
let total_signatures = signatures.values().sum::<u32>() as f64;
let score = total_signatures / total_certificates as f64 * 100.0;
(*address, score as u16)
})
.collect();
let tracked_rounds: Vec<_> = self.tracked_certificates.keys().skip_while(|r| *r % 2 == 0).copied().collect();
let certificate_participation_scores: IndexMap<_, _> = self
.validator_certificates
.iter()
.map(|(address, certificate_rounds)| {
let num_included_round_pairs = tracked_rounds
.chunks(2)
.filter(|chunk| chunk.iter().any(|r| certificate_rounds.contains(r)))
.count();
let num_round_pairs = (tracked_rounds.len().saturating_add(1)).saturating_div(2);
let score = num_included_round_pairs as f64 / num_round_pairs.max(1) as f64 * 100.0;
(*address, score as u16)
})
.collect();
let validator_addresses: IndexSet<_> =
signature_participation_scores.keys().chain(certificate_participation_scores.keys()).copied().collect();
let mut new_participation_scores = IndexMap::new();
for address in validator_addresses {
let signature_score = *signature_participation_scores.get(&address).unwrap_or(&0) as f64;
let certificate_score = *certificate_participation_scores.get(&address).unwrap_or(&0) as f64;
let combined_score = weighted_score(certificate_score, signature_score);
new_participation_scores.insert(address, (certificate_score, signature_score, combined_score));
}
self.participation_scores = new_participation_scores;
}
pub fn garbage_collect_certificates(&mut self, gc_round: u64) {
self.tracked_certificates.retain(|&round, _| round > gc_round);
self.validator_signatures.retain(|_, rounds| {
rounds.retain(|&round, _| round > gc_round);
!rounds.is_empty()
});
self.validator_certificates.retain(|_, rounds| {
rounds.retain(|&round| round > gc_round);
!rounds.is_empty()
});
}
}
fn scores_for_committee<N: Network>(
snapshot: &IndexMap<Address<N>, ParticipationScores>,
committee: &Committee<N>,
) -> IndexMap<Address<N>, (f64, f64)> {
committee
.members()
.iter()
.map(|(address, _)| {
let scores =
snapshot.get(address).map(|(cert_score, sig_score, _)| (*cert_score, *sig_score)).unwrap_or((0.0, 0.0));
(*address, scores)
})
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
use snarkvm::{
ledger::{
committee::test_helpers::sample_committee_for_round_and_members,
narwhal::batch_certificate::test_helpers::sample_batch_certificate_for_round,
},
prelude::MainnetV0,
utilities::TestRng,
};
use rand::RngExt;
type CurrentNetwork = MainnetV0;
#[test]
fn test_insert_certificates() {
let rng = &mut TestRng::default();
let mut state = TelemetryState::<CurrentNetwork>::new();
let current_round = 2;
let mut certificates = IndexSet::new();
for _ in 0..10 {
certificates.insert(sample_batch_certificate_for_round(current_round, rng));
}
assert!(state.tracked_certificates.is_empty());
for certificate in &certificates {
state.insert_certificate(certificate);
}
assert_eq!(state.tracked_certificates.get(¤t_round).unwrap().len(), certificates.len());
}
#[test]
fn test_insert_duplicate_certificate() {
let rng = &mut TestRng::default();
let mut state = TelemetryState::<CurrentNetwork>::new();
let current_round = 2;
let certificate = sample_batch_certificate_for_round(current_round, rng);
state.insert_certificate(&certificate);
let validator_signatures = state.validator_signatures.clone();
state.insert_certificate(&certificate);
assert_eq!(state.tracked_certificates.get(¤t_round).unwrap().len(), 1);
assert_eq!(state.validator_signatures, validator_signatures);
assert_eq!(state.validator_certificates.get(&certificate.author()).unwrap().len(), 1);
}
#[test]
fn test_participation_scores() {
let rng = &mut TestRng::default();
let mut state = TelemetryState::<CurrentNetwork>::new();
let current_round = 2;
let mut certificates = IndexSet::new();
certificates.insert(sample_batch_certificate_for_round(current_round, rng));
certificates.insert(sample_batch_certificate_for_round(current_round, rng));
certificates.insert(sample_batch_certificate_for_round(current_round, rng));
certificates.insert(sample_batch_certificate_for_round(current_round, rng));
let committee = sample_committee_for_round_and_members(
current_round,
vec![
certificates[0].author(),
certificates[1].author(),
certificates[2].author(),
certificates[3].author(),
],
rng,
);
assert!(state.tracked_certificates.is_empty());
for certificate in &certificates {
state.insert_certificate(certificate);
}
let participation_scores = scores_for_committee(state.participation_scores(), &committee);
assert_eq!(participation_scores.len(), committee.members().len());
for (address, _) in committee.members() {
assert_eq!(*participation_scores.get(address).unwrap(), (0.0, 0.0));
}
state.update_participation_scores();
let participation_scores = scores_for_committee(state.participation_scores(), &committee);
for (address, _) in committee.members() {
let (cert_score, sig_score) = *participation_scores.get(address).unwrap();
assert!(cert_score > 0.0 || sig_score > 0.0);
}
println!("{participation_scores:?}");
}
#[test]
fn test_garbage_collection() {
let rng = &mut TestRng::default();
let mut state = TelemetryState::<CurrentNetwork>::new();
let current_round = 2;
let next_round = current_round + 1;
let mut certificates = IndexSet::new();
let num_initial_certificates = rng.random_range(1..10);
for _ in 0..num_initial_certificates {
certificates.insert(sample_batch_certificate_for_round(current_round, rng));
}
let num_new_certificates = rng.random_range(1..10);
for _ in 0..num_new_certificates {
certificates.insert(sample_batch_certificate_for_round(next_round, rng));
}
for certificate in &certificates {
state.insert_certificate(certificate);
}
assert_eq!(state.tracked_certificates.get(¤t_round).unwrap().len(), num_initial_certificates);
assert_eq!(state.tracked_certificates.get(&next_round).unwrap().len(), num_new_certificates);
state.garbage_collect_certificates(current_round);
assert!(!state.tracked_certificates.contains_key(¤t_round));
assert_eq!(state.tracked_certificates.get(&next_round).unwrap().len(), num_new_certificates);
}
#[test]
fn test_scores_converge_after_dropped_updates() {
let rng = &mut TestRng::default();
const DROPPED_ROUNDS: std::ops::RangeInclusive<u64> = 5..=6;
let rounds: Vec<(u64, Vec<_>)> = (2..=21u64)
.map(|round| (round, (0..3).map(|_| sample_batch_certificate_for_round(round, rng)).collect()))
.collect();
let mut complete = TelemetryState::<CurrentNetwork>::new();
let mut lossy = TelemetryState::<CurrentNetwork>::new();
for (round, certificates) in &rounds {
for certificate in certificates {
complete.insert_certificate(certificate);
if !DROPPED_ROUNDS.contains(round) {
lossy.insert_certificate(certificate);
}
}
}
complete.update_participation_scores();
lossy.update_participation_scores();
assert_ne!(
complete.participation_scores(),
lossy.participation_scores(),
"the states should differ while the dropped rounds are still tracked"
);
complete.garbage_collect_certificates(*DROPPED_ROUNDS.end());
lossy.garbage_collect_certificates(*DROPPED_ROUNDS.end());
complete.update_participation_scores();
lossy.update_participation_scores();
assert_eq!(
complete.participation_scores(),
lossy.participation_scores(),
"the scores must converge once the dropped rounds are garbage collected"
);
assert!(!complete.participation_scores().is_empty(), "the test is vacuous if no scores survived");
}
#[tokio::test]
async fn test_worker_applies_updates() {
let rng = &mut TestRng::default();
let (telemetry, worker) = Telemetry::<CurrentNetwork>::new();
let handle = tokio::spawn(worker.run());
let current_round = 2;
let certificates: Vec<_> = (0..4).map(|_| sample_batch_certificate_for_round(current_round, rng)).collect();
let committee = sample_committee_for_round_and_members(
current_round,
certificates.iter().map(|certificate| certificate.author()).collect(),
rng,
);
let participation_scores = telemetry.get_participation_scores(&committee);
assert_eq!(participation_scores.len(), committee.members().len());
for (address, _) in committee.members() {
assert_eq!(*participation_scores.get(address).unwrap(), (0.0, 0.0));
}
for certificate in &certificates {
telemetry.insert_certificate(certificate);
}
telemetry.flush().await.unwrap();
assert_eq!(telemetry.num_dropped(), 0);
assert!(telemetry.is_complete());
assert_eq!(telemetry.get_participation_scores(&committee).len(), committee.members().len());
drop(telemetry);
handle.await.unwrap();
}
}