use crate::error::ChainIoError;
use alloy::{
consensus::Transaction as _,
primitives::{Address, Bytes, FixedBytes, B256, U256},
providers::Provider,
sol_types::SolCall,
};
use ark_bn254::G1Projective;
use ark_ec::{short_weierstrass::Affine, AffineRepr, CurveGroup};
use async_trait::async_trait;
use dashmap::DashMap;
use eigensdk::{
client_avsregistry::{error::AvsRegistryError, reader::AvsRegistryReader},
crypto_bls::{BlsG1Point, PublicKey},
services_avsregistry::AvsRegistryService,
services_operatorsinfo::operator_info::OperatorInfoService,
types::{
avs_state::{OperatorAvsState, QuorumAvsState},
operator::{OperatorInfo, OperatorPubKeys},
},
utils::slashing::middleware::operator_state_retriever::OperatorStateRetriever::CheckSignaturesIndices,
};
use futures::future::join_all;
use newton_core::state_commit_registry::{
IStateRootCommittable::StateCommit, StateCommitRegistry::commitStateRootCall,
};
use std::{collections::HashMap, sync::Arc, time::Instant};
use tracing::{debug, error, info, instrument, trace, warn};
#[derive(Debug, Clone)]
struct QuorumOperatorState {
block_num: u64,
operators: HashMap<FixedBytes<32>, OperatorAvsState>,
inserted_at: Instant,
}
#[derive(Debug)]
struct QuorumAggregateState {
block_num: u64,
state: QuorumAvsState,
inserted_at: Instant,
}
type OperatorStateCache = Arc<DashMap<u8, QuorumOperatorState>>;
type QuorumStateCache = Arc<DashMap<u8, QuorumAggregateState>>;
type OperatorIdToAddrCache = Arc<DashMap<FixedBytes<32>, Address>>;
const CACHE_TTL_SECS: u64 = 300;
pub mod errors;
pub mod diagnostics;
pub mod writer;
pub mod nonce_allocator;
pub mod cancel;
#[async_trait]
pub trait AvsRegistryServiceCaller: AvsRegistryService + Send + Sync {
fn warm_operator_id_cache(&self, mappings: Vec<(FixedBytes<32>, Address)>);
fn cached_operator_count(&self) -> usize;
fn get_operator_address(&self, operator_id: &FixedBytes<32>) -> Option<Address>;
fn get_all_operator_addresses(&self) -> Vec<(FixedBytes<32>, Address)>;
fn invalidate_caches(&self);
}
#[derive(Debug, Clone)]
pub struct AvsRegistryServiceChainCaller<R: AvsRegistryReader + Send + Sync, S: OperatorInfoService + Send + Sync> {
avs_registry: R,
operators_info_service: S,
operator_state_cache: OperatorStateCache,
quorum_state_cache: QuorumStateCache,
operator_id_to_addr_cache: OperatorIdToAddrCache,
}
impl<R: AvsRegistryReader + Send + Sync, S: OperatorInfoService + Send + Sync> AvsRegistryServiceChainCaller<R, S> {
pub fn new(avs_registry: R, operators_info_service: S) -> Self {
Self {
avs_registry,
operators_info_service,
operator_state_cache: Arc::new(DashMap::new()),
quorum_state_cache: Arc::new(DashMap::new()),
operator_id_to_addr_cache: Arc::new(DashMap::new()),
}
}
}
#[derive(Clone)]
#[allow(missing_debug_implementations)]
pub struct AvsRegistryServiceArcCaller {
pub inner: Arc<dyn AvsRegistryServiceCaller>,
}
#[async_trait]
impl AvsRegistryService for AvsRegistryServiceArcCaller {
async fn get_operators_avs_state_at_block(
&self,
block_num: u64,
quorum_nums: &[u8],
) -> Result<HashMap<FixedBytes<32>, OperatorAvsState>, AvsRegistryError> {
self.inner
.get_operators_avs_state_at_block(block_num, quorum_nums)
.await
}
async fn get_quorums_avs_state_at_block(
&self,
quorum_nums: &[u8],
block_num: u64,
) -> Result<HashMap<u8, QuorumAvsState>, AvsRegistryError> {
self.inner.get_quorums_avs_state_at_block(quorum_nums, block_num).await
}
async fn get_check_signatures_indices(
&self,
reference_block_number: u64,
quorum_numbers: Vec<u8>,
non_signer_operator_ids: Vec<FixedBytes<32>>,
) -> Result<CheckSignaturesIndices, AvsRegistryError> {
self.inner
.get_check_signatures_indices(reference_block_number, quorum_numbers, non_signer_operator_ids)
.await
}
}
#[async_trait]
impl AvsRegistryServiceCaller for AvsRegistryServiceArcCaller {
fn warm_operator_id_cache(&self, mappings: Vec<(FixedBytes<32>, Address)>) {
self.inner.warm_operator_id_cache(mappings);
}
fn cached_operator_count(&self) -> usize {
self.inner.cached_operator_count()
}
fn get_operator_address(&self, operator_id: &FixedBytes<32>) -> Option<Address> {
self.inner.get_operator_address(operator_id)
}
fn get_all_operator_addresses(&self) -> Vec<(FixedBytes<32>, Address)> {
self.inner.get_all_operator_addresses()
}
fn invalidate_caches(&self) {
self.inner.invalidate_caches();
}
}
impl<R: AvsRegistryReader + Send + Sync, S: OperatorInfoService + Send + Sync> AvsRegistryServiceCaller
for AvsRegistryServiceChainCaller<R, S>
{
fn warm_operator_id_cache(&self, mappings: Vec<(FixedBytes<32>, Address)>) {
let mut count = 0;
for (operator_id, address) in mappings {
self.operator_id_to_addr_cache.insert(operator_id, address);
count += 1;
}
debug!(
"[AvsRegistryServiceChainCaller] warmed operator ID cache with {} mappings",
count
);
}
fn cached_operator_count(&self) -> usize {
self.operator_id_to_addr_cache.len()
}
fn get_operator_address(&self, operator_id: &FixedBytes<32>) -> Option<Address> {
self.operator_id_to_addr_cache.get(operator_id).map(|r| *r)
}
fn get_all_operator_addresses(&self) -> Vec<(FixedBytes<32>, Address)> {
self.operator_id_to_addr_cache
.iter()
.map(|entry| (*entry.key(), *entry.value()))
.collect()
}
fn invalidate_caches(&self) {
self.operator_state_cache.clear();
self.quorum_state_cache.clear();
self.operator_id_to_addr_cache.clear();
}
}
#[async_trait]
impl<R: AvsRegistryReader + Send + Sync, S: OperatorInfoService + Send + Sync> AvsRegistryService
for AvsRegistryServiceChainCaller<R, S>
{
#[instrument(skip(self), fields(block_num, quorum_count = quorum_nums.len()))]
async fn get_operators_avs_state_at_block(
&self,
block_num: u64,
quorum_nums: &[u8],
) -> Result<HashMap<FixedBytes<32>, OperatorAvsState>, AvsRegistryError> {
let start_time = std::time::Instant::now();
let mut operators_avs_state: HashMap<FixedBytes<32>, OperatorAvsState> = HashMap::new();
debug!(
"[AvsRegistryServiceChainCaller] fetching operators AVS state for block {} quorums {:?}",
block_num, quorum_nums
);
let mut quorums_to_fetch = Vec::new();
let mut cached_quorums = Vec::new();
for quorum_num in quorum_nums {
if let Some(entry) = self.operator_state_cache.get(quorum_num) {
let age_secs = entry.inserted_at.elapsed().as_secs();
if entry.block_num >= block_num && age_secs < CACHE_TTL_SECS {
cached_quorums.push(*quorum_num);
for (op_id, op_state) in &entry.operators {
operators_avs_state
.entry(*op_id)
.or_insert(OperatorAvsState {
operator_id: op_state.operator_id,
operator_info: op_state.operator_info.clone(),
stake_per_quorum: HashMap::new(),
block_num: op_state.block_num,
})
.stake_per_quorum
.extend(op_state.stake_per_quorum.clone());
}
debug!(
"[AvsRegistryServiceChainCaller] quorum {} cache HIT (cached block: {}, age: {}s)",
quorum_num, entry.block_num, age_secs
);
} else {
quorums_to_fetch.push(*quorum_num);
debug!(
"[AvsRegistryServiceChainCaller] quorum {} needs refresh (cached block: {}, age: {}s)",
quorum_num, entry.block_num, age_secs
);
}
} else {
quorums_to_fetch.push(*quorum_num);
debug!("[AvsRegistryServiceChainCaller] quorum {} not in cache", quorum_num);
}
}
if !cached_quorums.is_empty() {
debug!(
"[AvsRegistryServiceChainCaller] using cached data for {} quorums: {:?}",
cached_quorums.len(),
cached_quorums
);
}
if quorums_to_fetch.is_empty() {
let total_duration = start_time.elapsed();
debug!(
"[AvsRegistryServiceChainCaller] all quorums cached - returning {} operators in {} ms",
operators_avs_state.len(),
total_duration.as_millis()
);
return Ok(operators_avs_state);
}
debug!(
"[AvsRegistryServiceChainCaller] fetching from chain for {} quorums: {:?}",
quorums_to_fetch.len(),
quorums_to_fetch
);
debug!(
"[AvsRegistryServiceChainCaller] quorum_nums to fetch hex: {}",
hex!(&quorums_to_fetch)
);
debug!("[AvsRegistryServiceChainCaller] fetching operator stakes from AVS registry...");
let stakes_fetch_start = std::time::Instant::now();
let operators_stakes_in_quorums = self
.avs_registry
.get_operators_stake_in_quorums_at_block(block_num, Bytes::from(quorums_to_fetch.clone()))
.await
.inspect_err(|e| {
error!(
"[AvsRegistryServiceChainCaller] failed to get operator stakes: {}",
e.to_string()
);
})?;
let stakes_fetch_duration = stakes_fetch_start.elapsed();
debug!(
"[AvsRegistryServiceChainCaller] fetched operator stakes in {} ms",
stakes_fetch_duration.as_millis()
);
let total_operators: usize = operators_stakes_in_quorums.iter().map(|q| q.len()).sum();
debug!(
"[AvsRegistryServiceChainCaller] received stakes for {} quorums with {} total operators",
operators_stakes_in_quorums.len(),
total_operators
);
for (i, quorum_stakes) in operators_stakes_in_quorums.iter().enumerate() {
debug!(
"[AvsRegistryServiceChainCaller] quorum {} has {} operators with stakes",
i,
quorum_stakes.len()
);
}
trace!(
"[AvsRegistryServiceChainCaller] detailed stakes data: {:#?}",
operators_stakes_in_quorums
);
if operators_stakes_in_quorums.len() != quorums_to_fetch.len() {
error!(
"[AvsRegistryServiceChainCaller] quorum count mismatch - expected: {}, got: {}",
quorums_to_fetch.len(),
operators_stakes_in_quorums.len()
);
return Err(AvsRegistryError::InvalidQuorumNums);
}
let mut quorum_operator_maps: HashMap<u8, HashMap<FixedBytes<32>, OperatorAvsState>> = HashMap::new();
let parallel_fetch_start = std::time::Instant::now();
let mut unique_operators: HashMap<FixedBytes<32>, Vec<(u8, U256)>> = HashMap::new();
for (quorum_id, quorum_num) in quorums_to_fetch.iter().enumerate() {
for operator in &operators_stakes_in_quorums[quorum_id] {
let operator_key = FixedBytes(*operator.operatorId);
unique_operators
.entry(operator_key)
.or_default()
.push((*quorum_num, U256::from(operator.stake)));
}
}
debug!(
"[AvsRegistryServiceChainCaller] fetching info for {} unique operators in PARALLEL",
unique_operators.len()
);
let fetch_futures: Vec<_> = unique_operators
.keys()
.map(|operator_key| {
let operator_id: [u8; 32] = **operator_key;
async move {
let operator_id_hex = hex!(operator_id);
let fetch_start = std::time::Instant::now();
let (info_result, socket_result) = tokio::join!(
self.get_operator_info(operator_id),
self.get_operator_socket(operator_id)
);
let duration = fetch_start.elapsed();
trace!(
"[AvsRegistryServiceChainCaller] parallel fetch for operator {} completed in {} ms",
operator_id_hex,
duration.as_millis()
);
(FixedBytes(operator_id), info_result, socket_result)
}
})
.collect();
let fetch_results = join_all(fetch_futures).await;
let parallel_fetch_duration = parallel_fetch_start.elapsed();
debug!(
"[AvsRegistryServiceChainCaller] parallel fetch of {} operators completed in {} ms",
unique_operators.len(),
parallel_fetch_duration.as_millis()
);
let mut processed_operators = 0;
let mut failed_operators = 0;
for (operator_key, info_result, socket_result) in fetch_results {
let operator_id_hex = hex!(operator_key.as_slice());
let info = match info_result {
Ok(info) => {
debug!(
"[AvsRegistryServiceChainCaller] retrieved info for operator {}: g1_key={:?}",
operator_id_hex, info.g1_pub_key
);
info
}
Err(e) => {
error!(
"[AvsRegistryServiceChainCaller] failed to get info for operator {}: {}",
operator_id_hex,
e.to_string()
);
failed_operators += 1;
return Err(e);
}
};
let socket = match socket_result {
Ok(socket) => {
debug!(
"[AvsRegistryServiceChainCaller] retrieved socket for operator {}: {}",
operator_id_hex, socket
);
socket
}
Err(e) => {
error!(
"[AvsRegistryServiceChainCaller] failed to get socket for operator {}: {}",
operator_id_hex,
e.to_string()
);
failed_operators += 1;
return Err(e);
}
};
let quorum_stakes = unique_operators.get(&operator_key).unwrap();
let mut stake_per_quorum = HashMap::new();
for (quorum_num, stake) in quorum_stakes {
stake_per_quorum.insert(*quorum_num, U256::from(*stake));
}
operators_avs_state.insert(
operator_key,
OperatorAvsState {
operator_id: operator_key,
operator_info: OperatorInfo {
pub_keys: Some(info.clone()),
socket: Some(socket.clone()),
},
stake_per_quorum: stake_per_quorum.clone(),
block_num,
},
);
for (quorum_num, stake) in quorum_stakes {
quorum_operator_maps.entry(*quorum_num).or_default().insert(
operator_key,
OperatorAvsState {
operator_id: operator_key,
operator_info: OperatorInfo {
pub_keys: Some(info.clone()),
socket: Some(socket.clone()),
},
stake_per_quorum: {
let mut map = HashMap::new();
map.insert(*quorum_num, *stake);
map
},
block_num,
},
);
}
processed_operators += 1;
}
let total_duration = start_time.elapsed();
debug!(
"[AvsRegistryServiceChainCaller] completed get_operators_avs_state_at_block in {} ms",
total_duration.as_millis()
);
debug!(
"[AvsRegistryServiceChainCaller] final stats - unique operators: {}, processed entries: {}, failed: {}",
operators_avs_state.len(),
processed_operators,
failed_operators
);
for (operator_key, state) in &operators_avs_state {
let operator_id_hex = hex!(operator_key.as_slice());
debug!(
"[AvsRegistryServiceChainCaller] operator {} participates in {} quorums: {:?}",
operator_id_hex,
state.stake_per_quorum.len(),
state.stake_per_quorum.keys().collect::<Vec<_>>()
);
}
trace!(
"[AvsRegistryServiceChainCaller] detailed AVS state: {:#?}",
operators_avs_state
);
for (quorum_num, operators) in quorum_operator_maps {
self.operator_state_cache.insert(
quorum_num,
QuorumOperatorState {
block_num,
operators: operators.clone(),
inserted_at: Instant::now(),
},
);
debug!(
"[AvsRegistryServiceChainCaller] cached quorum {} with {} operators at block {}",
quorum_num,
operators.len(),
block_num
);
}
debug!(
"[AvsRegistryServiceChainCaller] stored {} quorums in cache (total cached: {}, refreshes after {}s or newer block)",
quorums_to_fetch.len(),
self.operator_state_cache.len(),
CACHE_TTL_SECS
);
Ok(operators_avs_state)
}
#[instrument(skip(self), fields(quorum_count = quorum_nums.len(), block_num))]
async fn get_quorums_avs_state_at_block(
&self,
quorum_nums: &[u8],
block_num: u64,
) -> Result<HashMap<u8, QuorumAvsState>, AvsRegistryError> {
let start_time = std::time::Instant::now();
debug!(
"[AvsRegistryServiceChainCaller] fetching quorums AVS state for block {} quorums {:?}",
block_num, quorum_nums
);
let mut result: HashMap<u8, QuorumAvsState> = HashMap::new();
let mut quorums_to_compute = Vec::new();
let mut cached_quorums = Vec::new();
for quorum_num in quorum_nums {
if let Some(entry) = self.quorum_state_cache.get(quorum_num) {
let age_secs = entry.inserted_at.elapsed().as_secs();
if entry.block_num >= block_num && age_secs < CACHE_TTL_SECS {
cached_quorums.push(*quorum_num);
result.insert(
*quorum_num,
QuorumAvsState {
quorum_num: entry.state.quorum_num,
total_stake: entry.state.total_stake,
agg_pub_key_g1: entry.state.agg_pub_key_g1.clone(),
block_num: entry.state.block_num,
},
);
debug!(
"[AvsRegistryServiceChainCaller] quorum {} aggregate cache HIT (cached block: {}, age: {}s)",
quorum_num, entry.block_num, age_secs
);
} else {
quorums_to_compute.push(*quorum_num);
debug!(
"[AvsRegistryServiceChainCaller] quorum {} aggregate needs recompute (cached block: {}, age: {}s)",
quorum_num,
entry.block_num,
age_secs
);
}
} else {
quorums_to_compute.push(*quorum_num);
debug!(
"[AvsRegistryServiceChainCaller] quorum {} aggregate not in cache",
quorum_num
);
}
}
if !cached_quorums.is_empty() {
debug!(
"[AvsRegistryServiceChainCaller] using cached aggregate for {} quorums: {:?}",
cached_quorums.len(),
cached_quorums
);
}
if quorums_to_compute.is_empty() {
let total_duration = start_time.elapsed();
debug!(
"[AvsRegistryServiceChainCaller] all quorum aggregates cached - returning {} quorums in {} ms",
result.len(),
total_duration.as_millis()
);
return Ok(result);
}
debug!(
"[AvsRegistryServiceChainCaller] computing fresh aggregates for {} quorums: {:?}",
quorums_to_compute.len(),
quorums_to_compute
);
debug!("[AvsRegistryServiceChainCaller] fetching operators AVS state first for quorums to compute...");
let operators_avs_state = self
.get_operators_avs_state_at_block(block_num, &quorums_to_compute)
.await?;
debug!(
"[AvsRegistryServiceChainCaller] got operators state for {} quorums, now computing aggregates in parallel for {} operators",
quorums_to_compute.len(),
operators_avs_state.len()
);
use futures::future::join_all;
let compute_tasks: Vec<_> = quorums_to_compute
.iter()
.map(|quorum_num| {
let operators = operators_avs_state.clone();
let qnum = *quorum_num;
async move {
let quorum_start = std::time::Instant::now();
debug!("[AvsRegistryServiceChainCaller] computing aggregate for quorum {}", qnum);
let mut pub_key_g1 = G1Projective::from(PublicKey::identity());
let mut total_stake: U256 = U256::from(0);
let mut participating_operators = 0;
let mut operators_with_keys = 0;
for (operator_key, operator) in operators.iter() {
let operator_stake = operator
.stake_per_quorum
.get(&qnum)
.unwrap_or(&U256::ZERO);
if !operator_stake.is_zero() {
participating_operators += 1;
let operator_id_hex = hex!(operator_key.as_slice());
trace!(
"[AvsRegistryServiceChainCaller] operator {} has stake {} in quorum {}",
operator_id_hex, operator_stake, qnum
);
if let Some(pub_keys) = &operator.operator_info.pub_keys {
operators_with_keys += 1;
pub_key_g1 += pub_keys.g1_pub_key.g1();
total_stake += operator_stake;
trace!(
"[AvsRegistryServiceChainCaller] added operator {} key to aggregate (stake: {})",
operator_id_hex, operator_stake
);
} else {
warn!(
"[AvsRegistryServiceChainCaller] operator {} has stake but no public keys",
operator_id_hex
);
}
}
}
let agg_pub_key_g1 = if pub_key_g1 == G1Projective::from(PublicKey::zero()) {
debug!("[AvsRegistryServiceChainCaller] quorum {} has zero aggregate public key", qnum);
BlsG1Point::new(Affine::zero())
} else {
debug!("[AvsRegistryServiceChainCaller] computed non-zero aggregate public key for quorum {}", qnum);
BlsG1Point::new(pub_key_g1.into_affine())
};
let quorum_duration = quorum_start.elapsed();
debug!(
"[AvsRegistryServiceChainCaller] quorum {} aggregate computed in {} ms: {} participating operators, {} with keys, total stake: {}",
qnum, quorum_duration.as_millis(), participating_operators, operators_with_keys, total_stake
);
(
qnum,
QuorumAvsState {
quorum_num: qnum,
total_stake,
agg_pub_key_g1,
block_num,
},
)
}
})
.collect();
let computed_results = join_all(compute_tasks).await;
let computed_states: HashMap<u8, QuorumAvsState> = computed_results.into_iter().collect();
let total_duration = start_time.elapsed();
debug!(
"[AvsRegistryServiceChainCaller] completed quorum aggregates computation in {} ms for {} quorums",
total_duration.as_millis(),
quorums_to_compute.len()
);
for (quorum_num, state) in &computed_states {
result.insert(
*quorum_num,
QuorumAvsState {
quorum_num: state.quorum_num,
total_stake: state.total_stake,
agg_pub_key_g1: state.agg_pub_key_g1.clone(),
block_num: state.block_num,
},
);
}
for (quorum_num, state) in computed_states {
self.quorum_state_cache.insert(
quorum_num,
QuorumAggregateState {
block_num,
state,
inserted_at: Instant::now(),
},
);
debug!(
"[AvsRegistryServiceChainCaller] cached quorum {} aggregate at block {}",
quorum_num, block_num
);
}
debug!(
"[AvsRegistryServiceChainCaller] stored {} quorum aggregates in cache (total cached: {}, refreshes after {}s or newer block)",
quorums_to_compute.len(),
self.quorum_state_cache.len(),
CACHE_TTL_SECS
);
Ok(result)
}
#[instrument(skip(self), fields(
reference_block_number,
quorum_count = quorum_numbers.len(),
non_signer_count = non_signer_operator_ids.len()
))]
async fn get_check_signatures_indices(
&self,
reference_block_number: u64,
quorum_numbers: Vec<u8>,
non_signer_operator_ids: Vec<FixedBytes<32>>,
) -> Result<CheckSignaturesIndices, AvsRegistryError> {
let start_time = std::time::Instant::now();
debug!(
"[AvsRegistryServiceChainCaller] getting check signatures indices - block: {}, quorums: {}, non-signers: {}",
reference_block_number,
String::from_utf8(quorum_numbers.clone()).unwrap_or_default(),
non_signer_operator_ids.len()
);
if !non_signer_operator_ids.is_empty() {
debug!(
"[AvsRegistryServiceChainCaller] non-signer operator IDs: {:?}",
non_signer_operator_ids
.iter()
.map(|id| hex!(id.as_slice()))
.collect::<Vec<_>>()
);
}
let result = self
.avs_registry
.get_check_signatures_indices(reference_block_number, quorum_numbers, non_signer_operator_ids)
.await
.inspect_err(|e| {
error!(
"[AvsRegistryServiceChainCaller] failed to get check signatures indices: {}",
e.to_string()
);
})?;
let duration = start_time.elapsed();
debug!(
"[AvsRegistryServiceChainCaller] retrieved check signatures indices in {} ms",
duration.as_millis()
);
Ok(result)
}
}
impl<R: AvsRegistryReader + Send + Sync, S: OperatorInfoService + Send + Sync> AvsRegistryServiceChainCaller<R, S> {
#[instrument(skip(self), fields(operator_id = %hex!(operator_id)))]
async fn get_operator_info(&self, operator_id: [u8; 32]) -> Result<OperatorPubKeys, AvsRegistryError> {
let start_time = std::time::Instant::now();
let operator_id_hex = hex!(operator_id);
let operator_key = FixedBytes(operator_id);
let operator_addr = if let Some(cached_addr) = self.operator_id_to_addr_cache.get(&operator_key) {
debug!(
"[AvsRegistryServiceChainCaller] operator ID {} cache HIT (address: {})",
operator_id_hex, *cached_addr
);
*cached_addr
} else {
debug!(
"[AvsRegistryServiceChainCaller] operator ID {} cache MISS, resolving via RPC",
operator_id_hex
);
let addr = self
.avs_registry
.get_operator_from_id(operator_id)
.await
.inspect_err(|e| {
error!(
"[AvsRegistryServiceChainCaller] failed to resolve operator ID {} to address: {}",
operator_id_hex,
e.to_string()
);
})?;
self.operator_id_to_addr_cache.insert(operator_key, addr);
debug!(
"[AvsRegistryServiceChainCaller] operator ID {} resolved to address: {} (cached)",
operator_id_hex, addr
);
addr
};
let info_result = self.operators_info_service.get_operator_info(operator_addr).await;
let info = match info_result {
Ok(Some(info)) => {
let duration = start_time.elapsed();
debug!(
"[AvsRegistryServiceChainCaller] retrieved operator info for {} in {} ms",
operator_id_hex,
duration.as_millis()
);
trace!(
"[AvsRegistryServiceChainCaller] operator {} info: g1_key={:?}, g2_key={:?}",
operator_id_hex,
info.g1_pub_key,
info.g2_pub_key
);
Ok(info)
}
Ok(None) => {
warn!(
"[AvsRegistryServiceChainCaller] no operator info found for ID {} (address: {})",
operator_id_hex, operator_addr
);
Err(AvsRegistryError::GetOperatorInfo)
}
Err(e) => {
error!(
"[AvsRegistryServiceChainCaller] error getting operator info for ID {} (address: {}): {}",
operator_id_hex,
operator_addr,
e.to_string()
);
Err(AvsRegistryError::GetOperatorInfo)
}
};
info
}
#[instrument(skip(self), fields(operator_id = %hex!(operator_id)))]
async fn get_operator_socket(&self, operator_id: [u8; 32]) -> Result<String, AvsRegistryError> {
let start_time = std::time::Instant::now();
let operator_id_hex = hex!(operator_id);
let operator_key = FixedBytes(operator_id);
let operator_addr = if let Some(cached_addr) = self.operator_id_to_addr_cache.get(&operator_key) {
debug!(
"[AvsRegistryServiceChainCaller] operator ID {} cache HIT for socket lookup (address: {})",
operator_id_hex, *cached_addr
);
*cached_addr
} else {
debug!(
"[AvsRegistryServiceChainCaller] operator ID {} cache MISS for socket lookup, resolving via RPC",
operator_id_hex
);
let addr = self
.avs_registry
.get_operator_from_id(operator_id)
.await
.inspect_err(|e| {
error!(
"[AvsRegistryServiceChainCaller] failed to resolve operator ID {} for socket: {}",
operator_id_hex,
e.to_string()
);
})?;
self.operator_id_to_addr_cache.insert(operator_key, addr);
debug!(
"[AvsRegistryServiceChainCaller] operator ID {} resolved to address {} for socket lookup (cached)",
operator_id_hex, addr
);
addr
};
let socket_result = self.operators_info_service.get_operator_socket(operator_addr).await;
let socket = match socket_result {
Ok(Some(socket)) => {
let duration = start_time.elapsed();
debug!(
"[AvsRegistryServiceChainCaller] retrieved socket for operator {} in {} ms: {}",
operator_id_hex,
duration.as_millis(),
socket
);
Ok(socket)
}
Ok(None) => {
warn!(
"[AvsRegistryServiceChainCaller] no socket found for operator ID {} (address: {})",
operator_id_hex, operator_addr
);
Err(AvsRegistryError::GetOperatorInfo)
}
Err(e) => {
error!(
"[AvsRegistryServiceChainCaller] error getting socket for operator ID {} (address: {}): {}",
operator_id_hex,
operator_addr,
e.to_string()
);
Err(AvsRegistryError::GetOperatorInfo)
}
};
socket
}
}
pub async fn fetch_commit_state_root_calldata<P: Provider>(
provider: &P,
tx_hash: B256,
registry_addr: Address,
) -> Result<(StateCommit, Bytes), ChainIoError> {
let tx = provider
.get_transaction_by_hash(tx_hash)
.await?
.ok_or_else(|| ChainIoError::CommitStateRootCallFail {
reason: format!("tx {tx_hash} not found"),
})?;
let tx_to = tx.inner.to();
if tx_to != Some(registry_addr) {
return Err(ChainIoError::CommitStateRootCallFail {
reason: format!("tx {tx_hash} sent to {tx_to:?}, expected StateCommitRegistry {registry_addr}"),
});
}
let calldata = tx.inner.input();
let decoded = commitStateRootCall::abi_decode(calldata).map_err(|e| ChainIoError::CommitStateRootCallFail {
reason: format!("commitStateRoot calldata decode failed: {e}"),
})?;
Ok((decoded.c, decoded.blsCertificate))
}