use std::collections::HashSet;
use malachitebft_metrics::prometheus::encoding::EncodeLabelSet;
use malachitebft_metrics::prometheus::metrics::family::Family;
use malachitebft_metrics::prometheus::metrics::gauge::Gauge;
use malachitebft_metrics::Registry;
use tracing::{debug, warn};
use malachitebft_metrics::prometheus as prometheus_client;
use crate::state::{LocalNodeInfo, PeerInfo};
use crate::utils::Slots;
use crate::PeerType;
use libp2p::PeerId;
const MAX_PEER_SLOTS: usize = 100;
#[derive(Clone, Debug, Hash, PartialEq, Eq, EncodeLabelSet)]
pub(crate) struct PeerInfoLabels {
slot: String,
peer_moniker: String,
peer_id: String,
address: String,
peer_type: PeerType,
consensus_address: String, }
#[derive(Clone, Debug, Hash, PartialEq, Eq, EncodeLabelSet)]
pub(crate) struct MeshMembershipLabels {
peer_id: String,
peer_moniker: String,
topic: String, }
#[derive(Clone, Debug, Hash, PartialEq, Eq, EncodeLabelSet)]
pub(crate) struct ExplicitPeerLabels {
peer_id: String,
peer_moniker: String,
}
impl PeerInfo {
pub(crate) fn to_labels(&self, peer_id: &PeerId, slot: usize) -> PeerInfoLabels {
PeerInfoLabels {
slot: slot.to_string(),
peer_moniker: self.moniker.clone(),
peer_id: peer_id.to_string(),
address: self.address.to_string(),
peer_type: self.peer_type,
consensus_address: if self.peer_type.is_validator()
&& self.consensus_address != "unknown"
{
self.consensus_address.clone()
} else {
"none".to_string()
},
}
}
pub(crate) fn labels_match(&self, other: &PeerInfo) -> bool {
self.moniker == other.moniker
&& self.address == other.address
&& self.peer_type == other.peer_type
&& self.consensus_address == other.consensus_address
}
}
#[derive(Clone, Debug, Hash, PartialEq, Eq, EncodeLabelSet)]
pub(crate) struct LocalNodeLabels {
peer_id: String,
listen_addr: String,
consensus_address: String, }
pub(crate) struct Metrics {
local_node_info: Family<LocalNodeLabels, Gauge>,
discovered_peers: Family<PeerInfoLabels, Gauge>,
peer_mesh_membership: Family<MeshMembershipLabels, Gauge>,
explicit_peers: Family<ExplicitPeerLabels, Gauge>,
peer_slots: Slots<PeerId>,
}
impl std::fmt::Debug for Metrics {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Metrics")
.field("assigned_slots_count", &self.peer_slots.assigned())
.field("available_slots_count", &self.peer_slots.available())
.finish()
}
}
impl Metrics {
pub(crate) fn new(registry: &mut Registry) -> Self {
let local_node_info = Family::<LocalNodeLabels, Gauge>::default();
let peer_info = Family::<PeerInfoLabels, Gauge>::default();
let mesh_membership = Family::<MeshMembershipLabels, Gauge>::default();
let explicit_peers = Family::<ExplicitPeerLabels, Gauge>::default();
registry.register(
"local_node_info",
"Information about the local node (gauge value: 1 = validator, 0 = not validator)",
local_node_info.clone(),
);
registry.register(
"discovered_peers",
"Discovered/connected peers with basic info (gauge value = peer score)",
peer_info.clone(),
);
registry.register(
"peer_mesh_membership",
"Per-peer, per-topic gossipsub mesh membership (1 = in mesh, 0 = not in mesh)",
mesh_membership.clone(),
);
registry.register(
"explicit_peers",
"Peers added as explicit peers in gossipsub (1 = active, i64::MIN = disconnected)",
explicit_peers.clone(),
);
Self {
local_node_info,
discovered_peers: peer_info,
peer_mesh_membership: mesh_membership,
explicit_peers,
peer_slots: Slots::new(MAX_PEER_SLOTS),
}
}
pub(crate) fn set_local_node_info(&self, info: &LocalNodeInfo) {
let labels = LocalNodeLabels {
peer_id: info.peer_id.to_string(),
listen_addr: info.listen_addr.to_string(),
consensus_address: info
.consensus_address
.clone()
.unwrap_or_else(|| "none".to_string()),
};
let gauge_value = if info.is_validator { 1 } else { 0 };
self.local_node_info.get_or_create(&labels).set(gauge_value);
}
pub(crate) fn update_peer_metrics(
&mut self,
peer_id: &PeerId,
peer_info: &PeerInfo,
score: f64,
new_topics: Option<HashSet<String>>,
) -> Result<(), ()> {
let slot = match self.peer_slots.assign(*peer_id) {
Some(slot) => slot,
None => return Err(()), };
if let Some(ref new_topics) = new_topics {
let old_topics = &peer_info.topics;
for topic in old_topics.difference(new_topics) {
let mesh_labels = MeshMembershipLabels {
peer_id: peer_id.to_string(),
peer_moniker: peer_info.moniker.clone(),
topic: topic.clone(),
};
self.peer_mesh_membership.get_or_create(&mesh_labels).set(0);
}
for topic in new_topics.difference(old_topics) {
let mesh_labels = MeshMembershipLabels {
peer_id: peer_id.to_string(),
peer_moniker: peer_info.moniker.clone(),
topic: topic.clone(),
};
self.peer_mesh_membership.get_or_create(&mesh_labels).set(1);
}
}
let labels = peer_info.to_labels(peer_id, slot);
self.discovered_peers
.get_or_create(&labels)
.set(score as i64);
Ok(())
}
pub(crate) fn free_slot(&mut self, peer_id: &PeerId, peer_info: &PeerInfo) {
if let Some(slot) = self.peer_slots.release(peer_id) {
let labels = peer_info.to_labels(peer_id, slot);
self.discovered_peers.get_or_create(&labels).set(i64::MIN);
for topic in &peer_info.topics {
let mesh_labels = MeshMembershipLabels {
peer_id: peer_id.to_string(),
peer_moniker: peer_info.moniker.clone(),
topic: topic.clone(),
};
self.peer_mesh_membership.get_or_create(&mesh_labels).set(0);
}
debug!("Freed slot {slot} for peer {peer_id}");
}
}
pub(crate) fn record_explicit_peer(&self, peer_id: &PeerId, moniker: &str) {
let labels = ExplicitPeerLabels {
peer_id: peer_id.to_string(),
peer_moniker: moniker.to_string(),
};
self.explicit_peers.get_or_create(&labels).set(1);
}
pub(crate) fn mark_explicit_peer_stale(&self, peer_id: &PeerId, moniker: &str) {
let labels = ExplicitPeerLabels {
peer_id: peer_id.to_string(),
peer_moniker: moniker.to_string(),
};
self.explicit_peers.get_or_create(&labels).set(i64::MIN);
}
pub(crate) fn record_new_peer(&mut self, peer_id: &PeerId, peer_info: &PeerInfo) {
let slot = if let Some(existing_slot) = self.peer_slots.get(peer_id) {
existing_slot
} else {
let Some(new_slot) = self.peer_slots.assign(*peer_id) else {
warn!("No available metric slots for peer {peer_id}");
return;
};
new_slot
};
let labels = peer_info.to_labels(peer_id, slot);
self.discovered_peers
.get_or_create(&labels)
.set(peer_info.score as i64);
}
pub(crate) fn update_peer_labels(
&mut self,
peer_id: &PeerId,
old_peer_info: &PeerInfo,
new_peer_info: &PeerInfo,
) -> bool {
let Some(slot) = self.peer_slots.get(peer_id) else {
return false;
};
let labels_changed = !old_peer_info.labels_match(new_peer_info);
if labels_changed {
let old_labels = old_peer_info.to_labels(peer_id, slot);
tracing::debug!(%peer_id, ?old_labels, "Marking peer metric stale");
self.discovered_peers
.get_or_create(&old_labels)
.set(i64::MIN);
}
let new_labels = new_peer_info.to_labels(peer_id, slot);
self.discovered_peers
.get_or_create(&new_labels)
.set(new_peer_info.score as i64);
labels_changed
}
}