pub mod config;
pub mod event;
pub mod filter;
pub mod gateway_client;
pub mod listener;
use event::ChainEventType;
use newton_core::config::NewtonAvsConfig;
use newton_metric::{
inc_chain_watcher_events_filtered, inc_chain_watcher_events_received, inc_chain_watcher_events_relayed,
inc_chain_watcher_gateway_relay_errors,
};
use tokio_util::sync::CancellationToken;
use tracing::{error, info};
pub async fn run(
avs_config: NewtonAvsConfig<config::ChainWatcherConfig>,
cancellation_token: CancellationToken,
) -> eyre::Result<()> {
let chain_id = avs_config.chain_id;
let service_config = &avs_config.service;
info!(chain_id, "starting chain watcher");
let task_filter = filter::TaskFilter::new(&service_config.redis_url, chain_id).await?;
if service_config.gateway_psk.is_none() {
tracing::warn!("chain watcher running without gateway_psk; /watcher requests will be unsigned");
}
let gateway_client = gateway_client::GatewayClient::with_psk(
service_config.gateway_watcher_url.clone(),
Some(service_config.max_retries),
service_config.gateway_psk.as_deref(),
);
let rpc = avs_config.rpc.get_or_err(chain_id)?;
let task_manager = avs_config.contracts.avs.newton_prover_task_manager;
let allocation_manager = avs_config.contracts.eigenlayer.allocation_manager;
let identity_registry = avs_config.contracts.avs.identity_registry;
info!(
chain_id,
%task_manager,
%allocation_manager,
%identity_registry,
"chain watcher configured"
);
let watcher_listener = listener::ChainWatcherListener::new(
chain_id,
rpc.ws.clone(),
task_manager,
allocation_manager,
identity_registry,
);
let mut trigger_rx = watcher_listener.start(cancellation_token.clone()).await;
loop {
tokio::select! {
_ = cancellation_token.cancelled() => {
info!(chain_id, "chain watcher received shutdown signal");
break Ok(());
}
Some(chain_event) = trigger_rx.recv() => {
let event_type_label = match &chain_event.event_type {
ChainEventType::DirectOnchainTask { .. } => "direct_task",
ChainEventType::OperatorAdded { .. } => "operator_added",
ChainEventType::OperatorRemoved { .. } => "operator_removed",
ChainEventType::IdentityDataBound { .. } => "identity_data_bound",
};
inc_chain_watcher_events_received(chain_id, event_type_label);
match &chain_event.event_type {
ChainEventType::DirectOnchainTask { task_id, .. } => {
if task_filter.is_seen_by_gateway(task_id).await {
inc_chain_watcher_events_filtered(chain_id);
info!(
chain_id,
task_id = %task_id,
"task already seen by gateway, skipping"
);
continue;
}
info!(
chain_id,
task_id = %task_id,
"direct on-chain task detected, relaying to gateway"
);
}
ChainEventType::OperatorAdded { operator, operator_set_id, .. } => {
info!(
chain_id,
%operator,
operator_set_id,
"operator added event, relaying to gateway"
);
}
ChainEventType::OperatorRemoved { operator, operator_set_id, .. } => {
info!(
chain_id,
%operator,
operator_set_id,
"operator removed event, relaying to gateway"
);
}
ChainEventType::IdentityDataBound { identity_owner, data_ref_id, .. } => {
info!(
chain_id,
%identity_owner,
%data_ref_id,
"identity data bound event, relaying to gateway"
);
}
}
match gateway_client.submit_event(&chain_event).await {
Ok(()) => {
inc_chain_watcher_events_relayed(chain_id, event_type_label);
}
Err(e) => {
inc_chain_watcher_gateway_relay_errors(chain_id);
error!(
chain_id,
error = %e,
"failed to relay event to gateway"
);
}
}
}
}
}
}