newton-chain-watcher 0.5.2

newton chain watcher — smart event filter for direct on-chain tasks
//! Chain Watcher — smart event filter for direct on-chain tasks
//!
//! The chain watcher monitors blockchain events (via WebSocket) and distinguishes
//! between gateway-originated tasks and direct on-chain tasks. Only true direct
//! on-chain tasks are relayed to the gateway for processing.
//!
//! This crate provides:
//! - [`ChainWatcherListener`] — WebSocket event listener with reconnection
//! - [`ChainEvent`] / [`ChainEventType`] — typed event enums
//! - [`TaskFilter`] — Redis SET-backed filter for seen task IDs
//! - [`GatewayClient`] — HTTP client for pushing events to gateway's `/watcher` endpoint

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};

/// run the chain watcher event loop
///
/// Orchestrates the listener, filter, and gateway client:
/// 1. Starts WebSocket listener for on-chain events
/// 2. For each `NewTaskCreated` event, checks Redis SET (smart filter)
/// 3. Relays only direct on-chain tasks + all operator events to gateway
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");

    // Create the task filter (Redis-backed)
    let task_filter = filter::TaskFilter::new(&service_config.redis_url, chain_id).await?;

    // Create the gateway client. PSK is optional — when present, requests are
    // HMAC-signed so the gateway can reject forgeries (Octane #10).
    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(),
    );

    // Resolve contract addresses
    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"
    );

    // Start the WebSocket listener
    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;

    // Main event loop
    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, .. } => {
                        // Smart filter: check if gateway already handled this task
                        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"
                        );
                    }
                }

                // Relay 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"
                        );
                    }
                }
            }
        }
    }
}