use crate::{avs::errors::is_transient_rpc_error, provider::get_provider_with_retry};
use alloy::{
primitives::{Address, Bytes},
providers::Provider,
};
use eigensdk::{
client_avsregistry::reader::{AvsRegistryChainReader, AvsRegistryReader},
crypto_bls::{
alloy_registry_g1_point_to_g1_affine, alloy_registry_g2_point_to_g2_affine, BlsG1Point, BlsG2Point, OperatorId,
},
types::operator::{operator_id_from_g1_pub_key, OperatorPubKeys},
utils::slashing::middleware::{
bls_apk_registry::BLSApkRegistry, operator_state_retriever::OperatorStateRetriever,
socket_registry::SocketRegistry,
},
};
use futures::stream::{StreamExt, TryStreamExt};
use newton_core::operator_registry::OperatorRegistry;
use std::{
collections::HashMap,
sync::{atomic::Ordering, Arc},
};
use tokio::sync::{mpsc::UnboundedSender, oneshot::Sender, RwLock};
use tracing::{debug, error, info, instrument, warn};
use crate::registry::operator;
const MAX_CONCURRENT_OPERATOR_QUERIES: usize = 8;
fn skip_or_abort_operator(operator: Address, context: &str, err: eyre::Error) -> eyre::Result<bool> {
if is_transient_rpc_error(&err.to_string()) {
Err(err)
} else {
debug!(operator = %operator, "[{context}] skipping operator on non-transient error: {err}");
Ok(true)
}
}
#[derive(Debug, Clone)]
pub(crate) struct RegisteredOperator {
pub address: Address,
pub operator_id: OperatorId,
pub stakes: HashMap<u8, u128>,
}
#[derive(Debug)]
pub(crate) struct RegisteredSnapshot {
pub quorum_count: u8,
pub operators: Vec<RegisteredOperator>,
pub total_stakes: HashMap<u8, u128>,
pub block_number: u64,
}
async fn query_registered_operator_snapshot(
operator_registry_address: Address,
operator_state_retriever_address: Address,
http_rpc_url: String,
observed_quorums: &std::sync::atomic::AtomicBool,
) -> eyre::Result<Option<RegisteredSnapshot>> {
let provider = get_provider_with_retry(&http_rpc_url)?;
let operator_registry = OperatorRegistry::new(operator_registry_address, provider.clone());
let block_number = provider
.get_block_number()
.await
.map_err(|e| eyre::eyre!("failed to read current block number: {e}"))?;
let block_id = alloy::eips::BlockId::from(block_number);
let quorum_count = operator_registry
.quorumCount()
.block(block_id)
.call()
.await
.inspect_err(|e| error!("Failed to get quorum count: {e}"))?;
if quorum_count == 0 {
if observed_quorums.load(Ordering::Relaxed) {
warn!(
block_number,
"quorumCount read as 0 after previously observing a non-zero count; treating as no-signal (cache untouched)"
);
return Ok(None);
}
debug!("[query_registered_operator_snapshot] quorum count is 0 (bootstrap); no registered operators");
return Ok(Some(RegisteredSnapshot {
quorum_count,
operators: Vec::new(),
total_stakes: HashMap::new(),
block_number,
}));
}
observed_quorums.store(true, Ordering::Relaxed);
let reader = AvsRegistryChainReader::new(
operator_registry_address,
operator_state_retriever_address,
http_rpc_url.clone(),
)
.await
.map_err(|e| eyre::eyre!("failed to build AVS registry reader: {e}"))?;
let quorum_numbers: Vec<u8> = (0..quorum_count).collect();
let operators_by_quorum = reader
.get_operators_stake_in_quorums_at_block(block_number, Bytes::from(quorum_numbers))
.await
.map_err(|e| eyre::eyre!("failed to enumerate registered operators: {e}"))?;
let Some(snapshot) = assemble_registered_snapshot(quorum_count, &operators_by_quorum, block_number) else {
warn!(
block_number,
expected_quorums = quorum_count,
got_quorums = operators_by_quorum.len(),
"registered-operator snapshot is short; treating as no-signal (cache untouched)"
);
return Ok(None);
};
debug!(
"[query_registered_operator_snapshot] {} registered operators across {} quorums at block {}",
snapshot.operators.len(),
quorum_count,
block_number
);
Ok(Some(snapshot))
}
fn assemble_registered_snapshot(
quorum_count: u8,
operators_by_quorum: &[Vec<OperatorStateRetriever::Operator>],
block_number: u64,
) -> Option<RegisteredSnapshot> {
if operators_by_quorum.len() != quorum_count as usize {
return None;
}
let mut operators: Vec<RegisteredOperator> = Vec::new();
let mut index_by_address: HashMap<Address, usize> = HashMap::new();
let mut total_stakes: HashMap<u8, u128> = HashMap::new();
for (quorum_index, quorum_operators) in operators_by_quorum.iter().enumerate() {
let quorum_number = quorum_index as u8;
let mut quorum_total: u128 = 0;
for op in quorum_operators {
let stake = op.stake.to::<u128>();
quorum_total = quorum_total.saturating_add(stake);
let index = *index_by_address.entry(op.operator).or_insert_with(|| {
operators.push(RegisteredOperator {
address: op.operator,
operator_id: op.operatorId,
stakes: HashMap::new(),
});
operators.len() - 1
});
operators[index].stakes.insert(quorum_number, stake);
}
total_stakes.insert(quorum_number, quorum_total);
}
Some(RegisteredSnapshot {
quorum_count,
operators,
total_stakes,
block_number,
})
}
#[derive(strum_macros::Display, Debug, Clone, PartialEq)]
#[strum(serialize_all = "lowercase")]
pub enum StateSource {
Historic,
Event,
}
pub type VersionedStake = (u128, u64);
pub type OperatorStakeMap = HashMap<OperatorId, HashMap<u8, VersionedStake>>;
#[derive(Debug, Clone)]
pub struct OperatorStates {
pub operator_info_data: Arc<RwLock<HashMap<Address, (StateSource, OperatorPubKeys)>>>,
pub operator_addr_to_id: Arc<RwLock<HashMap<Address, (StateSource, OperatorId)>>>,
pub socket_dict: Arc<RwLock<HashMap<OperatorId, (StateSource, String)>>>,
pub operator_stake: Arc<RwLock<OperatorStakeMap>>,
pub total_stake: Arc<RwLock<HashMap<u8, VersionedStake>>>,
}
#[derive(Debug)]
pub struct OperatorSocket {
pub id: OperatorId,
pub socket: String,
}
#[derive(Debug, Clone, Copy)]
pub struct EpochAdvancedSignal {
pub chain_id: u64,
pub epoch: u32,
pub start_block: u64,
pub duration_blocks: u32,
}
#[derive(Debug)]
pub(crate) enum OperatorsInfoMessage {
InsertOperatorInfo(
Option<Address>,
Option<Box<OperatorPubKeys>>,
Option<OperatorSocket>,
StateSource,
),
InsertOperatorStake(OperatorId, u8, u128, u64),
InsertTotalStake(u8, u128, u64),
ReplaceOperatorStake(OperatorId, HashMap<u8, u128>, u64),
ReplaceTotalStakes(HashMap<u8, u128>, u64),
#[allow(dead_code)]
Remove(Address),
GetPubKeys(Address, Sender<Option<OperatorPubKeys>>),
GetSockets(Address, Sender<Option<String>>),
GetOperatorStake(OperatorId, u8, Sender<Option<u128>>),
GetTotalStake(u8, Sender<Option<u128>>),
}
pub(crate) fn apply_versioned_stake(map: &mut HashMap<u8, VersionedStake>, quorum: u8, stake: u128, block: u64) {
match map.get(&quorum) {
Some((_, existing_block)) if *existing_block > block => {}
_ => {
map.insert(quorum, (stake, block));
}
}
}
pub(crate) fn reconcile_snapshot(
map: &mut HashMap<u8, VersionedStake>,
snapshot: HashMap<u8, u128>,
snapshot_block: u64,
) {
map.retain(|quorum, (_, block)| *block > snapshot_block || snapshot.contains_key(quorum));
for (quorum, stake) in snapshot {
apply_versioned_stake(map, quorum, stake, snapshot_block);
}
}
pub(crate) async fn query_registered_operator_and_fill_db(
operator_registry_address: Address,
operator_state_retriever_address: Address,
http_rpc_url: String,
pub_keys: UnboundedSender<OperatorsInfoMessage>,
operator_states: &OperatorStates,
observed_quorums: &std::sync::atomic::AtomicBool,
) -> eyre::Result<()> {
let Some(snapshot) = query_registered_operator_snapshot(
operator_registry_address,
operator_state_retriever_address,
http_rpc_url.clone(),
observed_quorums,
)
.await?
else {
return Ok(());
};
let snapshot_block = snapshot.block_number;
let registered_addresses: Vec<Address> = snapshot.operators.iter().map(|o| o.address).collect();
let registered_ids: Vec<(Address, OperatorId)> =
snapshot.operators.iter().map(|o| (o.address, o.operator_id)).collect();
let stakes_by_addr: HashMap<Address, HashMap<u8, u128>> = snapshot
.operators
.iter()
.map(|o| (o.address, o.stakes.clone()))
.collect();
let total_stakes = snapshot.total_stakes;
let pubkey_handle =
query_registered_operator_pub_keys(operator_registry_address, http_rpc_url.clone(), registered_addresses);
let socket_handle =
query_registered_operator_sockets(operator_registry_address, http_rpc_url.clone(), registered_ids);
let mut operator_sockets: HashMap<OperatorId, String> = HashMap::new();
let mut operator_pub_keys: HashMap<Address, OperatorPubKeys> = HashMap::new();
let (pub_keys_res, sockets_res) = futures::join!(pubkey_handle, socket_handle);
let snapshot_complete = pub_keys_res.is_ok() && sockets_res.is_ok();
if let Ok(pub_keys) = pub_keys_res {
operator_pub_keys = pub_keys;
}
if let Ok(sockets) = sockets_res {
operator_sockets = sockets;
}
debug!(
"[query_registered_operator_and_fill_db] queried registered operator info. count of operators: {}",
operator_pub_keys.len()
);
if snapshot_complete {
let cached_addrs: std::collections::HashSet<Address> = {
let map = operator_states.operator_addr_to_id.read().await;
map.keys().copied().collect()
};
let current_addrs: std::collections::HashSet<Address> = operator_pub_keys.keys().copied().collect();
let stale: Vec<&Address> = cached_addrs.difference(¤t_addrs).collect();
if !stale.is_empty() {
debug!(
cached = cached_addrs.len(),
current = current_addrs.len(),
pruning = stale.len(),
"[query_registered_operator_and_fill_db] pruning stale operators from cache"
);
}
for stale_addr in stale {
debug!(
operator = %stale_addr,
"[query_registered_operator_and_fill_db] pruning stale operator"
);
let _ = pub_keys.send(OperatorsInfoMessage::Remove(*stale_addr));
}
} else {
warn!("[query_registered_operator_and_fill_db] incomplete operator snapshot (pubkey or socket RPC failed); skipping prune step (cache unchanged this tick)");
return Ok(());
}
for (address, operator_pub_key) in operator_pub_keys.iter() {
let operator_id = operator_id_from_g1_pub_key(operator_pub_key.g1_pub_key.clone())?;
if let Some(socket) = operator_sockets.get(&operator_id) {
let message = OperatorsInfoMessage::InsertOperatorInfo(
Some(*address),
Some(Box::new(operator_pub_key.clone())),
Some(OperatorSocket {
id: operator_id,
socket: socket.clone(),
}),
StateSource::Historic,
);
debug!("Inserting operator info for operator {address} from historic state: operator id {operator_id}");
let _ = pub_keys.send(message);
} else {
return Err(
alloy::contract::Error::TransportError(alloy::transports::TransportError::ErrorResp(
alloy::rpc::json_rpc::ErrorPayload {
code: -1,
message: "Socket not found".into(),
data: None,
},
))
.into(),
);
}
let stakes = stakes_by_addr.get(address).cloned().unwrap_or_default();
let _ = pub_keys.send(OperatorsInfoMessage::ReplaceOperatorStake(
operator_id,
stakes,
snapshot_block,
));
}
let _ = pub_keys.send(OperatorsInfoMessage::ReplaceTotalStakes(total_stakes, snapshot_block));
Ok(())
}
#[instrument(skip_all)]
pub async fn query_registered_operator_pub_keys(
operator_registry_address: Address,
http_rpc_url: String,
registered_operators: Vec<Address>,
) -> eyre::Result<(HashMap<Address, OperatorPubKeys>)> {
debug!("[query_registered_operator_pub_keys] operator_registry_address {operator_registry_address}");
let provider = get_provider_with_retry(&http_rpc_url)?;
let operator_registry = OperatorRegistry::new(operator_registry_address, provider.clone());
let bls_apk_registry_address = operator_registry
.blsApkRegistry()
.call()
.await
.inspect_err(|e| error!("Failed to get BLS APK registry address: {e}"))?;
let bls_apk_registry = BLSApkRegistry::new(bls_apk_registry_address, provider.clone());
debug!(
"[query_registered_operator_pub_keys] registered_operators count: {}",
registered_operators.len()
);
futures::stream::iter(registered_operators)
.map(|operator| {
let bls_apk_registry = bls_apk_registry.clone();
async move {
let fetch = async {
let operator_id = bls_apk_registry.getOperatorId(operator).call().await?;
let g1 = bls_apk_registry.getRegisteredPubkey(operator).call().await?._0;
let g2 = bls_apk_registry.getOperatorPubkeyG2(operator).call().await?;
Ok::<_, eyre::Error>((operator_id, g1, g2))
};
match fetch.await {
Ok((operator_id, g1, g2)) => {
let operator_pub_key = OperatorPubKeys {
g1_pub_key: BlsG1Point::new(alloy_registry_g1_point_to_g1_affine(g1)),
g2_pub_key: BlsG2Point::new(alloy_registry_g2_point_to_g2_affine(g2)),
};
debug!(
"[query_registered_operator_pub_keys] operator {}, operator_id {}",
operator, operator_id
);
debug!("[query_registered_operator_pub_keys] operator_pub_key {operator_pub_key:?}");
Ok(Some((operator, operator_pub_key)))
}
Err(e) => {
skip_or_abort_operator(operator, "query_registered_operator_pub_keys", e).map(|_skipped| None)
}
}
}
})
.buffer_unordered(MAX_CONCURRENT_OPERATOR_QUERIES)
.try_filter_map(|opt| async move { Ok(opt) })
.try_collect()
.await
.inspect_err(|e| {
error!("Transient failure fetching operator pubkeys; treating registered set as unobserved this tick: {e}")
})
}
#[instrument(skip_all)]
pub async fn query_registered_operator_sockets(
operator_registry_address: Address,
http_rpc_url: String,
registered_operators: Vec<(Address, OperatorId)>,
) -> eyre::Result<HashMap<OperatorId, String>> {
debug!("[query_registered_operator_sockets] operator_registry_address {operator_registry_address}");
let provider = get_provider_with_retry(&http_rpc_url)?;
let operator_registry = OperatorRegistry::new(operator_registry_address, provider.clone());
let socket_registry_address = operator_registry
.socketRegistry()
.call()
.await
.inspect_err(|e| error!("Failed to get socket_registry_address: {e}"))?;
debug!("query_registered_operator_sockets socket_registry_address {socket_registry_address}");
let socket_registry = SocketRegistry::new(socket_registry_address, provider);
debug!(
"[query_registered_operator_sockets] registered_operators count: {}",
registered_operators.len()
);
futures::stream::iter(registered_operators)
.map(|(operator, operator_id)| {
let socket_registry = socket_registry.clone();
async move {
match socket_registry.getOperatorSocket(operator_id).call().await {
Ok(socket) => {
debug!(
"[query_registered_operator_sockets] operator {}, operator_id {}, socket {}",
operator, operator_id, socket
);
Ok(Some((operator_id, socket)))
}
Err(e) => skip_or_abort_operator(operator, "query_registered_operator_sockets", e.into())
.map(|_skipped| None),
}
}
})
.buffer_unordered(MAX_CONCURRENT_OPERATOR_QUERIES)
.try_filter_map(|opt| async move { Ok(opt) })
.try_collect()
.await
.inspect_err(|e| {
error!("Transient failure fetching operator sockets; treating snapshot as incomplete this tick: {e}")
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn contract_revert_skips_operator() {
let err = eyre::eyre!("execution reverted: custom error 0x25ec6c1f, data: \"0x25ec6c1f\"");
let skipped = skip_or_abort_operator(Address::ZERO, "test", err).expect("revert should skip, not abort");
assert!(skipped);
}
#[test]
fn transient_error_aborts_tick() {
for msg in ["connection refused", "request timeout", "429 Too Many Requests"] {
let res = skip_or_abort_operator(Address::ZERO, "test", eyre::eyre!("{msg}"));
assert!(res.is_err(), "transient error {msg:?} should abort the tick");
}
}
use alloy::primitives::{aliases::U96, FixedBytes};
fn op(addr: u8, id: u8, stake: u128) -> OperatorStateRetriever::Operator {
OperatorStateRetriever::Operator {
operator: Address::repeat_byte(addr),
operatorId: FixedBytes::<32>::repeat_byte(id),
stake: U96::from(stake),
}
}
#[test]
fn snapshot_retains_registered_operator() {
let dewhitelisted = op(0xAA, 0x01, 500);
let snapshot = assemble_registered_snapshot(1, &[vec![dewhitelisted]], 100).expect("complete snapshot");
assert_eq!(snapshot.operators.len(), 1);
let kept = &snapshot.operators[0];
assert_eq!(kept.address, Address::repeat_byte(0xAA));
assert_eq!(kept.operator_id, FixedBytes::<32>::repeat_byte(0x01));
assert_eq!(kept.stakes.get(&0), Some(&500));
assert_eq!(snapshot.total_stakes.get(&0), Some(&500));
}
#[test]
fn snapshot_dedups_operator_across_quorums() {
let q0 = vec![op(0xAA, 0x01, 300), op(0xBB, 0x02, 700)];
let q1 = vec![op(0xAA, 0x01, 400)];
let snapshot = assemble_registered_snapshot(2, &[q0, q1], 100).expect("complete snapshot");
assert_eq!(snapshot.operators.len(), 2, "0xAA must not be double-counted");
let a = snapshot
.operators
.iter()
.find(|o| o.address == Address::repeat_byte(0xAA))
.expect("0xAA present");
assert_eq!(a.stakes.get(&0), Some(&300));
assert_eq!(a.stakes.get(&1), Some(&400));
assert_eq!(snapshot.total_stakes.get(&0), Some(&1000));
assert_eq!(snapshot.total_stakes.get(&1), Some(&400));
}
#[test]
fn short_snapshot_is_rejected() {
let only_one_quorum = vec![vec![op(0xAA, 0x01, 500)]];
assert!(
assemble_registered_snapshot(3, &only_one_quorum, 100).is_none(),
"short snapshot must be a no-signal None"
);
}
#[test]
fn empty_registered_set_is_a_valid_signal() {
let snapshot = assemble_registered_snapshot(2, &[vec![], vec![]], 100).expect("empty but complete");
assert!(snapshot.operators.is_empty());
assert_eq!(snapshot.total_stakes.get(&0), Some(&0));
assert_eq!(snapshot.total_stakes.get(&1), Some(&0));
}
#[test]
fn apply_versioned_stake_newer_block_wins() {
let mut map: HashMap<u8, VersionedStake> = HashMap::new();
apply_versioned_stake(&mut map, 0, 100, 10);
assert_eq!(map.get(&0), Some(&(100, 10)));
apply_versioned_stake(&mut map, 0, 200, 20);
assert_eq!(map.get(&0), Some(&(200, 20)));
apply_versioned_stake(&mut map, 0, 999, 15);
assert_eq!(map.get(&0), Some(&(200, 20)));
apply_versioned_stake(&mut map, 0, 250, 20);
assert_eq!(map.get(&0), Some(&(250, 20)));
}
#[test]
fn reconcile_snapshot_clears_dropped_quorum() {
let mut map: HashMap<u8, VersionedStake> = HashMap::new();
apply_versioned_stake(&mut map, 0, 100, 10);
apply_versioned_stake(&mut map, 1, 50, 10);
reconcile_snapshot(&mut map, HashMap::from([(0u8, 120u128)]), 20);
assert_eq!(map.get(&0), Some(&(120, 20)), "quorum 0 updated to snapshot value");
assert_eq!(
map.get(&1),
None,
"quorum 1 (dropped) must be cleared despite the earlier event"
);
}
#[test]
fn reconcile_snapshot_preserves_newer_event() {
let mut map: HashMap<u8, VersionedStake> = HashMap::new();
apply_versioned_stake(&mut map, 0, 500, 30);
reconcile_snapshot(&mut map, HashMap::new(), 20);
assert_eq!(
map.get(&0),
Some(&(500, 30)),
"post-snapshot event must not be cleared by an older snapshot"
);
}
#[test]
fn reconcile_snapshot_older_event_loses_to_snapshot() {
let mut map: HashMap<u8, VersionedStake> = HashMap::new();
apply_versioned_stake(&mut map, 0, 500, 10); reconcile_snapshot(&mut map, HashMap::from([(0u8, 700u128)]), 25);
assert_eq!(map.get(&0), Some(&(700, 25)), "newer snapshot supersedes older event");
}
}