From 304dfdc8f97cac992ed7289c5fbcb517c2e48c98 Mon Sep 17 00:00:00 2001 From: nikolamilosa Date: Wed, 12 Nov 2025 20:45:33 +0100 Subject: [PATCH 1/5] refactor: cleaning up deps --- .../node-rewards-scheduler/scheduler.py | 532 ----------------- docker/docker-compose.tools.yaml | 2 +- docker/docker-compose.yaml | 12 +- .../node_rewards_ingester.py | 541 ++++++++++++++++++ .../prom-config-builder.Dockerfile | 5 - docker/tools/python.Dockerfile | 8 + 6 files changed, 557 insertions(+), 543 deletions(-) delete mode 100644 docker/config/node-rewards-scheduler/scheduler.py create mode 100644 docker/tools/node-rewards-scheduler/node_rewards_ingester.py delete mode 100644 docker/tools/prom-config-builder/prom-config-builder.Dockerfile create mode 100644 docker/tools/python.Dockerfile diff --git a/docker/config/node-rewards-scheduler/scheduler.py b/docker/config/node-rewards-scheduler/scheduler.py deleted file mode 100644 index 3e0e30e..0000000 --- a/docker/config/node-rewards-scheduler/scheduler.py +++ /dev/null @@ -1,532 +0,0 @@ -#!/usr/bin/env python3 -""" -Node Rewards Scheduler for VictoriaMetrics -Directly interacts with IC canisters to fetch node rewards data and push to VictoriaMetrics -""" - -import os -import sys -import time -import logging -import subprocess -from datetime import datetime, timedelta, timezone -from typing import Dict, List, Optional, Any -from urllib.parse import urljoin - -# Configure logging -logging.basicConfig( - level=logging.INFO, - format='%(asctime)s - %(levelname)s - %(message)s', - datefmt='%Y-%m-%d %H:%M:%S' -) -logger = logging.getLogger(__name__) - - -def install_dependencies(): - """Install required Python packages""" - logger.info("Installing dependencies...") - - # Install Python packages - packages = [ - 'requests', - 'ic-py', # IC Python agent library - ] - - try: - subprocess.check_call( - [sys.executable, '-m', 'pip', 'install', '--quiet', '--no-warn-script-location'] + packages, - stdout=subprocess.DEVNULL, - stderr=subprocess.DEVNULL - ) - logger.info("✅ Dependencies installed") - except subprocess.CalledProcessError as e: - logger.error(f"❌ Failed to install dependencies: {e}") - sys.exit(1) - - -class ICCanisterClient: - """Client for interacting with Internet Computer canisters""" - - # IC mainnet URL - IC_URL = "https://ic0.app" - - # Canister IDs - NODE_REWARDS_CANISTER_ID = "uuew5-iiaaa-aaaaa-qbx4q-cai" # Node rewards canister (DEV) - GOVERNANCE_CANISTER_ID = "rrkah-fqaaa-aaaaa-aaaaq-cai" # NNS Governance canister - - def __init__(self): - from ic.client import Client - from ic.identity import Identity - from ic.agent import Agent - - # Create anonymous identity - self.identity = Identity() - self.client = Client(url=self.IC_URL) - self.agent = Agent(self.identity, self.client) - - async def get_rewards_daily(self, date_str: str) -> Dict[str, Any]: - """Fetch daily rewards data from node rewards canister""" - from ic.candid import encode, decode, Types - from ic.principal import Principal - - try: - # Parse date - date_obj = datetime.strptime(date_str, '%Y-%m-%d') - - # Define Candid types for the argument - # DateUtc: record { year: nat16; month: nat8; day: nat8 } - date_utc_type = Types.Record({ - 'year': Types.Nat32, - 'month': Types.Nat32, - 'day': Types.Nat32, - }) - - # GetNodeProvidersRewardsCalculationRequest: record { day: DateUtc } - request_type = Types.Record({ - 'day': date_utc_type - }) - - # Define return type structure - # NodeMetricsDaily - node_metrics_daily_type = Types.Record({ - 'subnet_assigned': Types.Opt(Types.Principal), - 'subnet_assigned_failure_rate': Types.Opt(Types.Float64), - 'num_blocks_proposed': Types.Opt(Types.Nat64), - 'num_blocks_failed': Types.Opt(Types.Nat64), - 'original_failure_rate': Types.Opt(Types.Float64), - 'relative_failure_rate': Types.Opt(Types.Float64), - }) - - # DailyNodeFailureRate (variant) - daily_node_failure_rate_type = Types.Variant({ - 'SubnetMember': Types.Record({ - 'node_metrics': Types.Opt(node_metrics_daily_type), - }), - 'NonSubnetMember': Types.Record({ - 'extrapolated_failure_rate': Types.Opt(Types.Float64), - }), - }) - - # DailyNodeRewards - daily_node_rewards_type = Types.Record({ - 'node_id': Types.Opt(Types.Principal), - 'node_reward_type': Types.Opt(Types.Text), - 'region': Types.Opt(Types.Text), - 'dc_id': Types.Opt(Types.Text), - 'daily_node_failure_rate': Types.Opt(daily_node_failure_rate_type), - 'performance_multiplier': Types.Opt(Types.Float64), - 'rewards_reduction': Types.Opt(Types.Float64), - 'base_rewards_xdr_permyriad': Types.Opt(Types.Float64), - 'adjusted_rewards_xdr_permyriad': Types.Opt(Types.Float64), - }) - - # NodeTypeRegionBaseRewards - node_type_region_base_rewards_type = Types.Record({ - 'monthly_xdr_permyriad': Types.Opt(Types.Float64), - 'daily_xdr_permyriad': Types.Opt(Types.Float64), - 'node_reward_type': Types.Opt(Types.Text), - 'region': Types.Opt(Types.Text), - }) - - # Type3RegionBaseRewards - type3_region_base_rewards_type = Types.Record({ - 'region': Types.Opt(Types.Text), - 'nodes_count': Types.Opt(Types.Nat64), - 'avg_rewards_xdr_permyriad': Types.Opt(Types.Float64), - 'avg_coefficient': Types.Opt(Types.Float64), - 'daily_xdr_permyriad': Types.Opt(Types.Float64), - }) - - # DailyNodeProviderRewards - daily_node_provider_rewards_type = Types.Record({ - 'total_base_rewards_xdr_permyriad': Types.Opt(Types.Nat64), - 'total_adjusted_rewards_xdr_permyriad': Types.Opt(Types.Nat64), - 'base_rewards': Types.Vec(node_type_region_base_rewards_type), - 'base_rewards_type3': Types.Vec(type3_region_base_rewards_type), - 'daily_nodes_rewards': Types.Vec(daily_node_rewards_type), - }) - - # DailyResults - daily_results_type = Types.Record({ - 'subnets_failure_rate': Types.Vec(Types.Tuple(Types.Principal, Types.Float64)), - 'provider_results': Types.Vec(Types.Tuple(Types.Principal, daily_node_provider_rewards_type)), - }) - - # GetNodeProvidersRewardsCalculationResponse: Result - return_type = Types.Variant({ - 'Ok': daily_results_type, - 'Err': Types.Text, - }) - - # Build the argument value - arg_value = { - 'day': { - 'year': date_obj.year, - 'month': date_obj.month, - 'day': date_obj.day, - } - } - - # Encode with explicit type information - arg_bytes = encode([{ - 'type': request_type, - 'value': arg_value - }]) - - method_name = 'get_node_providers_rewards_calculation' - - # Make the query with return type - response = self.agent.query_raw( - self.NODE_REWARDS_CANISTER_ID, - method_name, - arg_bytes, - return_type - ) - - # Response is a list with dict containing 'type' and 'value' - if not response or len(response) == 0: - logger.warning(f"Empty response for {date_str}") - return {} - - result = response[0].get('value', {}) - - # Handle Result variant (Ok/Err) - if 'Ok' in result: - daily_results = result['Ok'] - - # Convert lists of tuples to dictionaries for easier processing - # provider_results: [[Principal, {...}], ...] -> {Principal: {...}} - provider_results_dict = {} - for principal, provider_data in daily_results.get('provider_results', []): - provider_results_dict[str(principal)] = provider_data - - # subnets_failure_rate: [[Principal, float], ...] -> {Principal: float} - subnets_failure_rate_dict = {} - for principal, failure_rate in daily_results.get('subnets_failure_rate', []): - subnets_failure_rate_dict[str(principal)] = failure_rate - - result_dict = { - 'provider_results': provider_results_dict, - 'subnets_failure_rate': subnets_failure_rate_dict, - } - - logger.info(f"Successfully fetched rewards for {date_str} ({len(provider_results_dict)} providers)") - return result_dict - - elif 'Err' in result: - error_msg = result['Err'] - logger.warning(f"Canister returned error for {date_str}: {error_msg}") - return {} - else: - logger.warning(f"Unexpected response format for {date_str}: {result}") - return {} - - except Exception as e: - logger.warning(f"Failed to fetch rewards for {date_str}: {e}", exc_info=True) - return {} - - async def get_latest_governance_reward_event(self) -> Optional[float]: - """Fetch latest governance reward event timestamp from governance canister""" - from ic.candid import encode, decode, Types - from ic.principal import Principal - - try: - request_type = Types.Record({ - 'date_filter': Types.Opt(Types.Nat64) - }) - - # Build request with no date filter to get all rewards - # In ic-py, optional None is represented as an empty list [] - arg_bytes = encode([{ - 'type': request_type, - 'value': { - 'date_filter': [] # Empty list = None for optional types - } - }]) - - # Make query - response = self.agent.query_raw( - self.GOVERNANCE_CANISTER_ID, - 'list_node_provider_rewards', - arg_bytes, - ) - - # Extract value from response - if not response or len(response) == 0: - logger.warning("Empty response from governance canister") - return None - - result = response[0].get('value', {}) - rewards_list = result.get('rewards', []) - - # Rewards are in descending order (latest first) - if rewards_list and len(rewards_list) > 0: - first_reward = rewards_list[0] - timestamp = first_reward.get('timestamp') - if timestamp: - logger.info(f"Latest governance reward timestamp: {timestamp}") - return float(timestamp) - - except Exception as e: - logger.warning(f"Failed to fetch governance rewards: {e}", exc_info=True) - - return None - - -class NodeRewardsPusher: - """Pushes node rewards metrics to VictoriaMetrics""" - - def __init__(self, victoria_url: str): - self.victoria_url = victoria_url - self.ic_client = ICCanisterClient() - self.scheduler_dir = '/scheduler' - - @staticmethod - def _unwrap_optional(value): - """Unwrap Candid optional values (represented as lists)""" - if isinstance(value, list): - return value[0] if len(value) > 0 else None - return value - - def setup_environment(self): - """Set up the environment""" - logger.info("Setting up environment...") - - # Create scheduler directory - os.makedirs(self.scheduler_dir, exist_ok=True) - os.chdir(self.scheduler_dir) - - logger.info("✅ Environment ready") - - def wait_for_victoria_metrics(self): - """Wait for VictoriaMetrics to be ready""" - import requests - - logger.info("Waiting for VictoriaMetrics to be ready...") - - while True: - try: - response = requests.get(f"{self.victoria_url}/-/ready", timeout=5) - if response.status_code == 200: - logger.info("✅ VictoriaMetrics is ready") - return - except requests.exceptions.RequestException: - pass - - logger.info(f" Waiting for VictoriaMetrics at {self.victoria_url}...") - time.sleep(2) - - async def push_metrics_for_date(self, date_str: str) -> bool: - """ - Fetch node rewards data from IC canisters and push to VictoriaMetrics for a specific date - Returns True if successful, False otherwise - """ - import requests - - try: - logger.info(f"Pushing node rewards data for {date_str}") - - # Parse target date - target_date = datetime.strptime(date_str, '%Y-%m-%d') - - # Calculate noon timestamp in milliseconds - noon_dt = target_date.replace(hour=12, minute=0, second=0, microsecond=0) - noon_timestamp_ms = int(noon_dt.replace(tzinfo=timezone.utc).timestamp() * 1000) - - # Fetch data from IC canister - daily_results = await self.ic_client.get_rewards_daily(date_str) - - if not daily_results: - logger.warning(f"⚠️ No data available for {date_str}") - return False - - # Build Prometheus metrics - metrics_lines = [] - - # Provider-level metrics - provider_results = daily_results.get('provider_results', {}) - for provider_id, provider_rewards in provider_results.items(): - provider_id_str = str(provider_id) - - # nodes_count - nodes_count = len(provider_rewards.get('daily_nodes_rewards', [])) - metrics_lines.append( - f'nodes_count{{provider_id="{provider_id_str}"}} {nodes_count} {noon_timestamp_ms}' - ) - - # base_rewards - base_rewards = self._unwrap_optional(provider_rewards.get('total_base_rewards_xdr_permyriad')) - if base_rewards is not None: - metrics_lines.append( - f'total_base_rewards_xdr_permyriad{{provider_id="{provider_id_str}"}} {base_rewards} {noon_timestamp_ms}' - ) - - # adjusted_rewards - adjusted_rewards = self._unwrap_optional(provider_rewards.get('total_adjusted_rewards_xdr_permyriad')) - if adjusted_rewards is not None: - metrics_lines.append( - f'total_adjusted_rewards_xdr_permyriad{{provider_id="{provider_id_str}"}} {adjusted_rewards} {noon_timestamp_ms}' - ) - - # Node-level metrics - for node_result in provider_rewards.get('daily_nodes_rewards', []): - node_id = self._unwrap_optional(node_result.get('node_id')) - node_id_str = str(node_id) if node_id else '' - - # daily_node_failure_rate is optional and contains a variant - failure_rate_data = self._unwrap_optional(node_result.get('daily_node_failure_rate')) - if isinstance(failure_rate_data, dict) and 'SubnetMember' in failure_rate_data: - # Extract node_metrics from SubnetMember variant - subnet_member = failure_rate_data['SubnetMember'] - node_metrics = self._unwrap_optional(subnet_member.get('node_metrics')) - - if node_metrics: - subnet_id = self._unwrap_optional(node_metrics.get('subnet_assigned')) - subnet_id_str = str(subnet_id) if subnet_id else '' - - # original_failure_rate - original_fr = self._unwrap_optional(node_metrics.get('original_failure_rate')) - if original_fr is not None: - metrics_lines.append( - f'original_failure_rate{{provider_id="{provider_id_str}",node_id="{node_id_str}",subnet_id="{subnet_id_str}"}} {original_fr} {noon_timestamp_ms}' - ) - - # relative_failure_rate - relative_fr = self._unwrap_optional(node_metrics.get('relative_failure_rate')) - if relative_fr is not None: - metrics_lines.append( - f'relative_failure_rate{{provider_id="{provider_id_str}",node_id="{node_id_str}",subnet_id="{subnet_id_str}"}} {relative_fr} {noon_timestamp_ms}' - ) - - # Subnet-level metrics - subnets_failure_rate = daily_results.get('subnets_failure_rate', {}) - for subnet_id, failure_rate in subnets_failure_rate.items(): - subnet_id_str = str(subnet_id) - metrics_lines.append( - f'subnet_failure_rate{{subnet_id="{subnet_id_str}"}} {failure_rate} {noon_timestamp_ms}' - ) - - # Governance timestamp - gov_timestamp = await self.ic_client.get_latest_governance_reward_event() - if gov_timestamp: - metrics_lines.append( - f'governance_latest_reward_event_timestamp_seconds {gov_timestamp} {noon_timestamp_ms}' - ) - - # Push to VictoriaMetrics - if metrics_lines: - metrics_payload = '\n'.join(metrics_lines) + '\n' - import_url = urljoin(self.victoria_url.rstrip('/') + '/', 'api/v1/import/prometheus') - - response = requests.post( - import_url, - data=metrics_payload.encode('utf-8'), - headers={'Content-Type': 'text/plain'}, - timeout=30 - ) - - if response.status_code in (200, 204): - logger.info(f"✅ Successfully pushed data for {date_str} ({len(metrics_lines)} metrics)") - return True - else: - logger.error(f"❌ Failed to push to VictoriaMetrics: {response.status_code} - {response.text}") - return False - else: - logger.warning(f"No metrics to push for {date_str}") - return False - - except Exception as e: - logger.error(f"❌ Error processing {date_str}: {e}", exc_info=True) - return False - - async def backfill(self, days: int = 40): - """Backfill historical data""" - - logger.info(f"Starting backfill of last {days} days...") - - for i in range(days, 0, -1): - date_obj = datetime.now(timezone.utc) - timedelta(days=i) - date_str = date_obj.strftime('%Y-%m-%d') - - logger.info(f"[{days - i + 1:2d}/{days}] Backfilling data for {date_str}...") - await self.push_metrics_for_date(date_str) - - logger.info("✅ Backfill complete!") - - def wait_until_next_run(self): - """Wait until 00:05 UTC tomorrow""" - now = datetime.now(timezone.utc) - - # Calculate next 00:05 UTC - next_run = now.replace(hour=0, minute=5, second=0, microsecond=0) - if now >= next_run: - # Already past 00:05 today, schedule for tomorrow - next_run += timedelta(days=1) - - wait_seconds = (next_run - now).total_seconds() - logger.info(f"Next run scheduled for {next_run.strftime('%Y-%m-%d %H:%M:%S')} UTC") - logger.info(f"Waiting {wait_seconds:.0f} seconds ({wait_seconds/3600:.1f} hours)...") - - time.sleep(wait_seconds) - - async def run_daily_scheduler(self): - """Run the daily scheduler loop""" - - while True: - try: - # Wait until 00:05 UTC - self.wait_until_next_run() - - # Push yesterday's data - yesterday = (datetime.now(timezone.utc) - timedelta(days=1)).strftime('%Y-%m-%d') - logger.info(f"Running scheduled push for {yesterday}") - await self.push_metrics_for_date(yesterday) - - except KeyboardInterrupt: - logger.info("Scheduler stopped by user") - break - except Exception as e: - logger.error(f"Error in scheduler loop: {e}", exc_info=True) - # Wait a bit before retrying - time.sleep(60) - - -async def main_async(): - """Async main entry point""" - # Get configuration from environment - victoria_url = os.environ.get('VICTORIA_METRICS_URL', 'http://localhost:9090') - - # Create pusher - pusher = NodeRewardsPusher(victoria_url) - - # Set up environment - pusher.setup_environment() - - # Wait for VictoriaMetrics - pusher.wait_for_victoria_metrics() - - # Backfill historical data - await pusher.backfill(days=40) - - # Run daily scheduler - await pusher.run_daily_scheduler() - - -def main(): - """Main entry point""" - logger.info("=" * 50) - logger.info("Node Rewards Scheduler for VictoriaMetrics") - logger.info("=" * 50) - logger.info("") - - # Install dependencies - install_dependencies() - - # Run async main loop - import asyncio - asyncio.run(main_async()) - - -if __name__ == '__main__': - main() - diff --git a/docker/docker-compose.tools.yaml b/docker/docker-compose.tools.yaml index bfc7ee3..f0f520b 100644 --- a/docker/docker-compose.tools.yaml +++ b/docker/docker-compose.tools.yaml @@ -2,7 +2,7 @@ services: prom-config-builder: build: context: . - dockerfile: ./tools/prom-config-builder/prom-config-builder.Dockerfile + dockerfile: ./tools/python.Dockerfile volumes: - ./config/prometheus:/config/prometheus - ./tools/prom-config-builder/:/tools/prom-config-builder diff --git a/docker/docker-compose.yaml b/docker/docker-compose.yaml index 58352a5..e6cf3b3 100644 --- a/docker/docker-compose.yaml +++ b/docker/docker-compose.yaml @@ -14,7 +14,7 @@ services: user: "${UID}:${GID}" victoriametrics: - image: victoriametrics/victoria-metrics:v1.97.1 + image: victoriametrics/victoria-metrics:v1.129.1 network_mode: host volumes: - ./config/prometheus:/config @@ -34,14 +34,16 @@ services: reservations: memory: 12G - node-rewards-scheduler: - image: python:3.11-slim + node-rewards-ingester: + build: + context: . + dockerfile: ./tools/python.Dockerfile network_mode: host environment: VICTORIA_METRICS_URL: http://localhost:9090 volumes: - - ./config/node-rewards-scheduler/scheduler.py:/app/scheduler.py:ro - command: python3 /app/scheduler.py + - ./tools/node-rewards-scheduler/node_rewards_ingester.py:/app/node_rewards_ingester.py + command: /app/node_rewards_ingester.py user: "${UID}:${GID}" depends_on: - victoriametrics diff --git a/docker/tools/node-rewards-scheduler/node_rewards_ingester.py b/docker/tools/node-rewards-scheduler/node_rewards_ingester.py new file mode 100644 index 0000000..65c75ef --- /dev/null +++ b/docker/tools/node-rewards-scheduler/node_rewards_ingester.py @@ -0,0 +1,541 @@ +""" +Node Rewards Scheduler for VictoriaMetrics +Directly interacts with IC canisters to fetch node rewards data and push to VictoriaMetrics +""" + +import logging +import os +import time +from datetime import datetime, timedelta, timezone +from typing import Any, Dict, Optional +from urllib.parse import urljoin + +import requests +from ic.agent import Agent +from ic.candid import Types, encode +from ic.client import Client +from ic.identity import Identity + +# Configure logging +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s - %(levelname)s - %(message)s", + datefmt="%Y-%m-%d %H:%M:%S", +) +logger = logging.getLogger(__name__) + + +class ICCanisterClient: + """Client for interacting with Internet Computer canisters""" + + # IC mainnet URL + IC_URL = "https://ic0.app" + + # Canister IDs + NODE_REWARDS_CANISTER_ID = ( + "uuew5-iiaaa-aaaaa-qbx4q-cai" # Node rewards canister (DEV) + ) + GOVERNANCE_CANISTER_ID = "rrkah-fqaaa-aaaaa-aaaaq-cai" # NNS Governance canister + + def __init__(self): + # Create anonymous identity + self.identity = Identity() + self.client = Client(url=self.IC_URL) + self.agent = Agent(self.identity, self.client) + + def get_rewards_daily(self, date_str: str) -> Dict[str, Any]: + """Fetch daily rewards data from node rewards canister""" + + try: + # Parse date + date_obj = datetime.strptime(date_str, "%Y-%m-%d") + + # Define Candid types for the argument + # DateUtc: record { year: nat16; month: nat8; day: nat8 } + date_utc_type = Types.Record( + { + "year": Types.Nat32, + "month": Types.Nat32, + "day": Types.Nat32, + } + ) + + # GetNodeProvidersRewardsCalculationRequest: record { day: DateUtc } + request_type = Types.Record({"day": date_utc_type}) + + # Define return type structure + # NodeMetricsDaily + node_metrics_daily_type = Types.Record( + { + "subnet_assigned": Types.Opt(Types.Principal), + "subnet_assigned_failure_rate": Types.Opt(Types.Float64), + "num_blocks_proposed": Types.Opt(Types.Nat64), + "num_blocks_failed": Types.Opt(Types.Nat64), + "original_failure_rate": Types.Opt(Types.Float64), + "relative_failure_rate": Types.Opt(Types.Float64), + } + ) + + # DailyNodeFailureRate (variant) + daily_node_failure_rate_type = Types.Variant( + { + "SubnetMember": Types.Record( + { + "node_metrics": Types.Opt(node_metrics_daily_type), + } + ), + "NonSubnetMember": Types.Record( + { + "extrapolated_failure_rate": Types.Opt(Types.Float64), + } + ), + } + ) + + # DailyNodeRewards + daily_node_rewards_type = Types.Record( + { + "node_id": Types.Opt(Types.Principal), + "node_reward_type": Types.Opt(Types.Text), + "region": Types.Opt(Types.Text), + "dc_id": Types.Opt(Types.Text), + "daily_node_failure_rate": Types.Opt(daily_node_failure_rate_type), + "performance_multiplier": Types.Opt(Types.Float64), + "rewards_reduction": Types.Opt(Types.Float64), + "base_rewards_xdr_permyriad": Types.Opt(Types.Float64), + "adjusted_rewards_xdr_permyriad": Types.Opt(Types.Float64), + } + ) + + # NodeTypeRegionBaseRewards + node_type_region_base_rewards_type = Types.Record( + { + "monthly_xdr_permyriad": Types.Opt(Types.Float64), + "daily_xdr_permyriad": Types.Opt(Types.Float64), + "node_reward_type": Types.Opt(Types.Text), + "region": Types.Opt(Types.Text), + } + ) + + # Type3RegionBaseRewards + type3_region_base_rewards_type = Types.Record( + { + "region": Types.Opt(Types.Text), + "nodes_count": Types.Opt(Types.Nat64), + "avg_rewards_xdr_permyriad": Types.Opt(Types.Float64), + "avg_coefficient": Types.Opt(Types.Float64), + "daily_xdr_permyriad": Types.Opt(Types.Float64), + } + ) + + # DailyNodeProviderRewards + daily_node_provider_rewards_type = Types.Record( + { + "total_base_rewards_xdr_permyriad": Types.Opt(Types.Nat64), + "total_adjusted_rewards_xdr_permyriad": Types.Opt(Types.Nat64), + "base_rewards": Types.Vec(node_type_region_base_rewards_type), + "base_rewards_type3": Types.Vec(type3_region_base_rewards_type), + "daily_nodes_rewards": Types.Vec(daily_node_rewards_type), + } + ) + + # DailyResults + daily_results_type = Types.Record( + { + "subnets_failure_rate": Types.Vec( + Types.Tuple(Types.Principal, Types.Float64) + ), + "provider_results": Types.Vec( + Types.Tuple(Types.Principal, daily_node_provider_rewards_type) + ), + } + ) + + # GetNodeProvidersRewardsCalculationResponse: Result + return_type = Types.Variant( + { + "Ok": daily_results_type, + "Err": Types.Text, + } + ) + + # Build the argument value + arg_value = { + "day": { + "year": date_obj.year, + "month": date_obj.month, + "day": date_obj.day, + } + } + + # Encode with explicit type information + arg_bytes = encode([{"type": request_type, "value": arg_value}]) + + method_name = "get_node_providers_rewards_calculation" + + # Make the query with return type + response = self.agent.query_raw( + self.NODE_REWARDS_CANISTER_ID, method_name, arg_bytes, return_type + ) + + # Response is a list with dict containing 'type' and 'value' + if not response or len(response) == 0: + logger.warning(f"Empty response for {date_str}") + return {} + + result = response[0].get("value", {}) + + # Handle Result variant (Ok/Err) + if "Ok" in result: + daily_results = result["Ok"] + + # Convert lists of tuples to dictionaries for easier processing + # provider_results: [[Principal, {...}], ...] -> {Principal: {...}} + provider_results_dict = {} + for principal, provider_data in daily_results.get( + "provider_results", [] + ): + provider_results_dict[str(principal)] = provider_data + + # subnets_failure_rate: [[Principal, float], ...] -> {Principal: float} + subnets_failure_rate_dict = {} + for principal, failure_rate in daily_results.get( + "subnets_failure_rate", [] + ): + subnets_failure_rate_dict[str(principal)] = failure_rate + + result_dict = { + "provider_results": provider_results_dict, + "subnets_failure_rate": subnets_failure_rate_dict, + } + + logger.info( + f"Successfully fetched rewards for {date_str} ({len(provider_results_dict)} providers)" + ) + return result_dict + + elif "Err" in result: + error_msg = result["Err"] + logger.warning(f"Canister returned error for {date_str}: {error_msg}") + return {} + else: + logger.warning(f"Unexpected response format for {date_str}: {result}") + return {} + + except Exception as e: + logger.warning( + f"Failed to fetch rewards for {date_str}: {e}", exc_info=True + ) + return {} + + def get_latest_governance_reward_event(self) -> Optional[float]: + """Fetch latest governance reward event timestamp from governance canister""" + + try: + request_type = Types.Record({"date_filter": Types.Opt(Types.Nat64)}) + + # Build request with no date filter to get all rewards + # In ic-py, optional None is represented as an empty list [] + arg_bytes = encode( + [ + { + "type": request_type, + "value": { + "date_filter": [] # Empty list = None for optional types + }, + } + ] + ) + + # Make query + response = self.agent.query_raw( + self.GOVERNANCE_CANISTER_ID, + "list_node_provider_rewards", + arg_bytes, + ) + + # Extract value from response + if not response or len(response) == 0: + logger.warning("Empty response from governance canister") + return None + + result = response[0].get("value", {}) + rewards_list = result.get("rewards", []) + + # Rewards are in descending order (latest first) + if rewards_list and len(rewards_list) > 0: + first_reward = rewards_list[0] + timestamp = first_reward.get("timestamp") + if timestamp: + logger.info(f"Latest governance reward timestamp: {timestamp}") + return float(timestamp) + + except Exception as e: + logger.warning(f"Failed to fetch governance rewards: {e}", exc_info=True) + + return None + + +class NodeRewardsPusher: + """Pushes node rewards metrics to VictoriaMetrics""" + + def __init__(self, victoria_url: str): + self.victoria_url = victoria_url + self.ic_client = ICCanisterClient() + + @staticmethod + def _unwrap_optional(value): + """Unwrap Candid optional values (represented as lists)""" + if isinstance(value, list): + return value[0] if len(value) > 0 else None + return value + + def wait_for_victoria_metrics(self): + """Wait for VictoriaMetrics to be ready""" + + logger.info("Waiting for VictoriaMetrics to be ready...") + + while True: + try: + response = requests.get(f"{self.victoria_url}/-/ready", timeout=5) + if response.status_code == 200: + logger.info("✅ VictoriaMetrics is ready") + return + except requests.exceptions.RequestException: + continue + + logger.info(f" Waiting for VictoriaMetrics at {self.victoria_url}...") + time.sleep(2) + + def push_metrics_for_date(self, date_str: str) -> bool: + """ + Fetch node rewards data from IC canisters and push to VictoriaMetrics for a specific date + Returns True if successful, False otherwise + """ + + try: + logger.info(f"Pushing node rewards data for {date_str}") + + # Parse target date + target_date = datetime.strptime(date_str, "%Y-%m-%d") + + # Calculate noon timestamp in milliseconds + noon_dt = target_date.replace(hour=12, minute=0, second=0, microsecond=0) + noon_timestamp_ms = int( + noon_dt.replace(tzinfo=timezone.utc).timestamp() * 1000 + ) + + # Fetch data from IC canister + daily_results = self.ic_client.get_rewards_daily(date_str) + + if not daily_results: + logger.warning(f"⚠️ No data available for {date_str}") + return False + + # Build Prometheus metrics + metrics_lines = [] + + # Provider-level metrics + provider_results = daily_results.get("provider_results", {}) + for provider_id, provider_rewards in provider_results.items(): + provider_id_str = str(provider_id) + + # nodes_count + nodes_count = len(provider_rewards.get("daily_nodes_rewards", [])) + metrics_lines.append( + f'nodes_count{{provider_id="{provider_id_str}"}} {nodes_count} {noon_timestamp_ms}' + ) + + # base_rewards + base_rewards = self._unwrap_optional( + provider_rewards.get("total_base_rewards_xdr_permyriad") + ) + if base_rewards is not None: + metrics_lines.append( + f'total_base_rewards_xdr_permyriad{{provider_id="{provider_id_str}"}} {base_rewards} {noon_timestamp_ms}' + ) + + # adjusted_rewards + adjusted_rewards = self._unwrap_optional( + provider_rewards.get("total_adjusted_rewards_xdr_permyriad") + ) + if adjusted_rewards is not None: + metrics_lines.append( + f'total_adjusted_rewards_xdr_permyriad{{provider_id="{provider_id_str}"}} {adjusted_rewards} {noon_timestamp_ms}' + ) + + # Node-level metrics + for node_result in provider_rewards.get("daily_nodes_rewards", []): + node_id = self._unwrap_optional(node_result.get("node_id")) + node_id_str = str(node_id) if node_id else "" + + # daily_node_failure_rate is optional and contains a variant + failure_rate_data = self._unwrap_optional( + node_result.get("daily_node_failure_rate") + ) + if ( + isinstance(failure_rate_data, dict) + and "SubnetMember" in failure_rate_data + ): + # Extract node_metrics from SubnetMember variant + subnet_member = failure_rate_data["SubnetMember"] + node_metrics = self._unwrap_optional( + subnet_member.get("node_metrics") + ) + + if node_metrics: + subnet_id = self._unwrap_optional( + node_metrics.get("subnet_assigned") + ) + subnet_id_str = str(subnet_id) if subnet_id else "" + + # original_failure_rate + original_fr = self._unwrap_optional( + node_metrics.get("original_failure_rate") + ) + if original_fr is not None: + metrics_lines.append( + f'original_failure_rate{{provider_id="{provider_id_str}",node_id="{node_id_str}",subnet_id="{subnet_id_str}"}} {original_fr} {noon_timestamp_ms}' + ) + + # relative_failure_rate + relative_fr = self._unwrap_optional( + node_metrics.get("relative_failure_rate") + ) + if relative_fr is not None: + metrics_lines.append( + f'relative_failure_rate{{provider_id="{provider_id_str}",node_id="{node_id_str}",subnet_id="{subnet_id_str}"}} {relative_fr} {noon_timestamp_ms}' + ) + + # Subnet-level metrics + subnets_failure_rate = daily_results.get("subnets_failure_rate", {}) + for subnet_id, failure_rate in subnets_failure_rate.items(): + subnet_id_str = str(subnet_id) + metrics_lines.append( + f'subnet_failure_rate{{subnet_id="{subnet_id_str}"}} {failure_rate} {noon_timestamp_ms}' + ) + + # Governance timestamp + gov_timestamp = self.ic_client.get_latest_governance_reward_event() + if gov_timestamp: + metrics_lines.append( + f"governance_latest_reward_event_timestamp_seconds {gov_timestamp} {noon_timestamp_ms}" + ) + + # Push to VictoriaMetrics + if metrics_lines: + metrics_payload = "\n".join(metrics_lines) + "\n" + import_url = urljoin( + self.victoria_url.rstrip("/") + "/", "api/v1/import/prometheus" + ) + + response = requests.post( + import_url, + data=metrics_payload.encode("utf-8"), + headers={"Content-Type": "text/plain"}, + timeout=30, + ) + + if response.status_code in (200, 204): + logger.info( + f"✅ Successfully pushed data for {date_str} ({len(metrics_lines)} metrics)" + ) + return True + else: + logger.error( + f"❌ Failed to push to VictoriaMetrics: {response.status_code} - {response.text}" + ) + return False + else: + logger.warning(f"No metrics to push for {date_str}") + return False + + except Exception as e: + logger.error(f"❌ Error processing {date_str}: {e}", exc_info=True) + return False + + def backfill(self, days: int = 40): + """Backfill historical data""" + + logger.info(f"Starting backfill of last {days} days...") + + for i in range(days, 0, -1): + date_obj = datetime.now(timezone.utc) - timedelta(days=i) + date_str = date_obj.strftime("%Y-%m-%d") + + logger.info( + f"[{days - i + 1:2d}/{days}] Backfilling data for {date_str}..." + ) + self.push_metrics_for_date(date_str) + + logger.info("✅ Backfill complete!") + + def wait_until_next_run(self): + """Wait until 00:05 UTC tomorrow""" + now = datetime.now(timezone.utc) + + # Calculate next 00:05 UTC + next_run = now.replace(hour=0, minute=5, second=0, microsecond=0) + if now >= next_run: + # Already past 00:05 today, schedule for tomorrow + next_run += timedelta(days=1) + + wait_seconds = (next_run - now).total_seconds() + logger.info( + f"Next run scheduled for {next_run.strftime('%Y-%m-%d %H:%M:%S')} UTC" + ) + logger.info( + f"Waiting {wait_seconds:.0f} seconds ({wait_seconds / 3600:.1f} hours)..." + ) + + time.sleep(wait_seconds) + + def run_daily_scheduler(self): + """Run the daily scheduler loop""" + + while True: + try: + # Wait until 00:05 UTC + self.wait_until_next_run() + + # Push yesterday's data + yesterday = (datetime.now(timezone.utc) - timedelta(days=1)).strftime( + "%Y-%m-%d" + ) + logger.info(f"Running scheduled push for {yesterday}") + self.push_metrics_for_date(yesterday) + + except KeyboardInterrupt: + logger.info("Scheduler stopped by user") + break + except Exception as e: + logger.error(f"Error in scheduler loop: {e}", exc_info=True) + # Wait a bit before retrying + time.sleep(60) + + +def main(): + """Main entry point""" + logger.info("=" * 50) + logger.info("Node Rewards Ingester for VictoriaMetrics") + logger.info("=" * 50) + logger.info("") + + # Get configuration from environment + victoria_url = os.environ.get("VICTORIA_METRICS_URL", "http://localhost:9090") + + # Create pusher + pusher = NodeRewardsPusher(victoria_url) + + # Wait for VictoriaMetrics + pusher.wait_for_victoria_metrics() + + # Backfill historical data + pusher.backfill(days=40) + + # Run daily scheduler + pusher.run_daily_scheduler() + + +if __name__ == "__main__": + main() diff --git a/docker/tools/prom-config-builder/prom-config-builder.Dockerfile b/docker/tools/prom-config-builder/prom-config-builder.Dockerfile deleted file mode 100644 index 5be884d..0000000 --- a/docker/tools/prom-config-builder/prom-config-builder.Dockerfile +++ /dev/null @@ -1,5 +0,0 @@ -FROM python:3.12-slim - -RUN pip install --no-cache-dir pyyaml - -ENTRYPOINT ["python3"] diff --git a/docker/tools/python.Dockerfile b/docker/tools/python.Dockerfile new file mode 100644 index 0000000..292ee89 --- /dev/null +++ b/docker/tools/python.Dockerfile @@ -0,0 +1,8 @@ +FROM python:3.12-slim + +RUN pip install --no-cache-dir \ + pyyaml \ + requests \ + ic-py + +ENTRYPOINT ["python3"] From f404ab50592843ec40c86c0eabc7bafd9ab64e2c Mon Sep 17 00:00:00 2001 From: nikolamilosa Date: Wed, 12 Nov 2025 21:13:45 +0100 Subject: [PATCH 2/5] further refactoring --- .../node_rewards_ingester.py | 341 ++++++++---------- 1 file changed, 157 insertions(+), 184 deletions(-) diff --git a/docker/tools/node-rewards-scheduler/node_rewards_ingester.py b/docker/tools/node-rewards-scheduler/node_rewards_ingester.py index 65c75ef..862ac37 100644 --- a/docker/tools/node-rewards-scheduler/node_rewards_ingester.py +++ b/docker/tools/node-rewards-scheduler/node_rewards_ingester.py @@ -25,208 +25,179 @@ logger = logging.getLogger(__name__) -class ICCanisterClient: - """Client for interacting with Internet Computer canisters""" - - # IC mainnet URL - IC_URL = "https://ic0.app" - - # Canister IDs - NODE_REWARDS_CANISTER_ID = ( - "uuew5-iiaaa-aaaaa-qbx4q-cai" # Node rewards canister (DEV) - ) - GOVERNANCE_CANISTER_ID = "rrkah-fqaaa-aaaaa-aaaaq-cai" # NNS Governance canister +IC_URL = "https://ic0.app" + +# These have to be tuples +NODE_REWARDS_CANISTER_IDS = [ + ("uuew5-iiaaa-aaaaa-qbx4q-cai"), # Dev + ("sgymv-uiaaa-aaaaa-aaaia-cai"), # Prod +] +GOVERNANCE_CANISTER_ID = "rrkah-fqaaa-aaaaa-aaaaq-cai" # NNS Governance canister + +###################### TYPE DEFINITIONS ########################### +DATE_UTC_TYPE = Types.Record( + { + "year": Types.Nat32, + "month": Types.Nat32, + "day": Types.Nat32, + } +) - def __init__(self): - # Create anonymous identity - self.identity = Identity() - self.client = Client(url=self.IC_URL) - self.agent = Agent(self.identity, self.client) +REQUEST_TYPE = Types.Record({"day": DATE_UTC_TYPE}) + +NODE_METRICS_DAILY_TYPE = Types.Record( + { + "subnet_assigned": Types.Opt(Types.Principal), + "subnet_assigned_failure_rate": Types.Opt(Types.Float64), + "num_blocks_proposed": Types.Opt(Types.Nat64), + "num_blocks_failed": Types.Opt(Types.Nat64), + "original_failure_rate": Types.Opt(Types.Float64), + "relative_failure_rate": Types.Opt(Types.Float64), + } +) - def get_rewards_daily(self, date_str: str) -> Dict[str, Any]: - """Fetch daily rewards data from node rewards canister""" +DAILY_NODE_FAILURE_RATE_TYPE = Types.Variant( + { + "SubnetMember": Types.Record( + { + "node_metrics": Types.Opt(NODE_METRICS_DAILY_TYPE), + } + ), + "NonSubnetMember": Types.Record( + { + "extrapolated_failure_rate": Types.Opt(Types.Float64), + } + ), + } +) - try: - # Parse date - date_obj = datetime.strptime(date_str, "%Y-%m-%d") - - # Define Candid types for the argument - # DateUtc: record { year: nat16; month: nat8; day: nat8 } - date_utc_type = Types.Record( - { - "year": Types.Nat32, - "month": Types.Nat32, - "day": Types.Nat32, - } - ) +DAILY_NODE_REWARDS_TYPE = Types.Record( + { + "node_id": Types.Opt(Types.Principal), + "node_reward_type": Types.Opt(Types.Text), + "region": Types.Opt(Types.Text), + "dc_id": Types.Opt(Types.Text), + "daily_node_failure_rate": Types.Opt(DAILY_NODE_FAILURE_RATE_TYPE), + "performance_multiplier": Types.Opt(Types.Float64), + "rewards_reduction": Types.Opt(Types.Float64), + "base_rewards_xdr_permyriad": Types.Opt(Types.Float64), + "adjusted_rewards_xdr_permyriad": Types.Opt(Types.Float64), + } +) - # GetNodeProvidersRewardsCalculationRequest: record { day: DateUtc } - request_type = Types.Record({"day": date_utc_type}) - - # Define return type structure - # NodeMetricsDaily - node_metrics_daily_type = Types.Record( - { - "subnet_assigned": Types.Opt(Types.Principal), - "subnet_assigned_failure_rate": Types.Opt(Types.Float64), - "num_blocks_proposed": Types.Opt(Types.Nat64), - "num_blocks_failed": Types.Opt(Types.Nat64), - "original_failure_rate": Types.Opt(Types.Float64), - "relative_failure_rate": Types.Opt(Types.Float64), - } - ) +NODE_TYPE_REGION_BASE_REWARDS_TYPE = Types.Record( + { + "monthly_xdr_permyriad": Types.Opt(Types.Float64), + "daily_xdr_permyriad": Types.Opt(Types.Float64), + "node_reward_type": Types.Opt(Types.Text), + "region": Types.Opt(Types.Text), + } +) - # DailyNodeFailureRate (variant) - daily_node_failure_rate_type = Types.Variant( - { - "SubnetMember": Types.Record( - { - "node_metrics": Types.Opt(node_metrics_daily_type), - } - ), - "NonSubnetMember": Types.Record( - { - "extrapolated_failure_rate": Types.Opt(Types.Float64), - } - ), - } - ) +TYPE3_REGION_BASE_REWARDS_TYPE = Types.Record( + { + "region": Types.Opt(Types.Text), + "nodes_count": Types.Opt(Types.Nat64), + "avg_rewards_xdr_permyriad": Types.Opt(Types.Float64), + "avg_coefficient": Types.Opt(Types.Float64), + "daily_xdr_permyriad": Types.Opt(Types.Float64), + } +) - # DailyNodeRewards - daily_node_rewards_type = Types.Record( - { - "node_id": Types.Opt(Types.Principal), - "node_reward_type": Types.Opt(Types.Text), - "region": Types.Opt(Types.Text), - "dc_id": Types.Opt(Types.Text), - "daily_node_failure_rate": Types.Opt(daily_node_failure_rate_type), - "performance_multiplier": Types.Opt(Types.Float64), - "rewards_reduction": Types.Opt(Types.Float64), - "base_rewards_xdr_permyriad": Types.Opt(Types.Float64), - "adjusted_rewards_xdr_permyriad": Types.Opt(Types.Float64), - } - ) +DAILY_NODE_PROVIDER_REWARDS_TYPE = Types.Record( + { + "total_base_rewards_xdr_permyriad": Types.Opt(Types.Nat64), + "total_adjusted_rewards_xdr_permyriad": Types.Opt(Types.Nat64), + "base_rewards": Types.Vec(NODE_TYPE_REGION_BASE_REWARDS_TYPE), + "base_rewards_type3": Types.Vec(TYPE3_REGION_BASE_REWARDS_TYPE), + "daily_nodes_rewards": Types.Vec(DAILY_NODE_REWARDS_TYPE), + } +) - # NodeTypeRegionBaseRewards - node_type_region_base_rewards_type = Types.Record( - { - "monthly_xdr_permyriad": Types.Opt(Types.Float64), - "daily_xdr_permyriad": Types.Opt(Types.Float64), - "node_reward_type": Types.Opt(Types.Text), - "region": Types.Opt(Types.Text), - } - ) +DAILY_RESULTS_TYPE = Types.Record( + { + "subnets_failure_rate": Types.Vec(Types.Tuple(Types.Principal, Types.Float64)), + "provider_results": Types.Vec( + Types.Tuple(Types.Principal, DAILY_NODE_PROVIDER_REWARDS_TYPE) + ), + } +) - # Type3RegionBaseRewards - type3_region_base_rewards_type = Types.Record( - { - "region": Types.Opt(Types.Text), - "nodes_count": Types.Opt(Types.Nat64), - "avg_rewards_xdr_permyriad": Types.Opt(Types.Float64), - "avg_coefficient": Types.Opt(Types.Float64), - "daily_xdr_permyriad": Types.Opt(Types.Float64), - } - ) +RETURN_TYPE = Types.Variant( + { + "Ok": DAILY_RESULTS_TYPE, + "Err": Types.Text, + } +) - # DailyNodeProviderRewards - daily_node_provider_rewards_type = Types.Record( - { - "total_base_rewards_xdr_permyriad": Types.Opt(Types.Nat64), - "total_adjusted_rewards_xdr_permyriad": Types.Opt(Types.Nat64), - "base_rewards": Types.Vec(node_type_region_base_rewards_type), - "base_rewards_type3": Types.Vec(type3_region_base_rewards_type), - "daily_nodes_rewards": Types.Vec(daily_node_rewards_type), - } - ) - # DailyResults - daily_results_type = Types.Record( - { - "subnets_failure_rate": Types.Vec( - Types.Tuple(Types.Principal, Types.Float64) - ), - "provider_results": Types.Vec( - Types.Tuple(Types.Principal, daily_node_provider_rewards_type) - ), - } - ) +class NodeRewardsClient: + """Client for interacting with the node rewards canister""" - # GetNodeProvidersRewardsCalculationResponse: Result - return_type = Types.Variant( - { - "Ok": daily_results_type, - "Err": Types.Text, - } - ) + def __init__(self, ic_url: str, canister_id: str): + # Create anonymous identity + self.identity = Identity() + self.client = Client(url=ic_url) + self.agent = Agent(self.identity, self.client) + self.canister_id = canister_id - # Build the argument value - arg_value = { - "day": { - "year": date_obj.year, - "month": date_obj.month, - "day": date_obj.day, - } + def get_rewards_daily(self, date: str) -> Dict[str, Any]: + """Fetch daily rewards data from node rewards canister""" + parsed_date = datetime.strptime(date, "%Y-%m-%d") + arg_value = { + "day": { + "year": parsed_date.year, + "month": parsed_date.month, + "day": parsed_date.day, } + } + arg_bytes = encode([{"type": REQUEST_TYPE, "value": arg_value}]) + + response = self.agent.query_raw( + self.canister_id, + "get_node_providers_rewards_calculation", + arg_bytes, + RETURN_TYPE, + ) - # Encode with explicit type information - arg_bytes = encode([{"type": request_type, "value": arg_value}]) + if not response or len(response) == 0: + logger.error(f"Empty response for {date}") + return {} - method_name = "get_node_providers_rewards_calculation" + result = response[0].get("value", {}) - # Make the query with return type - response = self.agent.query_raw( - self.NODE_REWARDS_CANISTER_ID, method_name, arg_bytes, return_type + if "Err" in result: + error_msg = result["Err"] + logger.error(f"Canister returned error for {date}: {error_msg}") + return {} + + if "Ok" not in result: + raise ValueError( + f"Unexpected result, contains neither `Err` nur `Ok`: {result}" ) - # Response is a list with dict containing 'type' and 'value' - if not response or len(response) == 0: - logger.warning(f"Empty response for {date_str}") - return {} + daily_results = result["Ok"] - result = response[0].get("value", {}) + # Convert lists of tuples to dictionaries for easier processing + # provider_results: [[Principal, {...}], ...] -> {Principal: {...}} + provider_results_dict = {} + for principal, provider_data in daily_results.get("provider_results", []): + provider_results_dict[str(principal)] = provider_data - # Handle Result variant (Ok/Err) - if "Ok" in result: - daily_results = result["Ok"] - - # Convert lists of tuples to dictionaries for easier processing - # provider_results: [[Principal, {...}], ...] -> {Principal: {...}} - provider_results_dict = {} - for principal, provider_data in daily_results.get( - "provider_results", [] - ): - provider_results_dict[str(principal)] = provider_data - - # subnets_failure_rate: [[Principal, float], ...] -> {Principal: float} - subnets_failure_rate_dict = {} - for principal, failure_rate in daily_results.get( - "subnets_failure_rate", [] - ): - subnets_failure_rate_dict[str(principal)] = failure_rate - - result_dict = { - "provider_results": provider_results_dict, - "subnets_failure_rate": subnets_failure_rate_dict, - } - - logger.info( - f"Successfully fetched rewards for {date_str} ({len(provider_results_dict)} providers)" - ) - return result_dict + # subnets_failure_rate: [[Principal, float], ...] -> {Principal: float} + subnets_failure_rate_dict = {} + for principal, failure_rate in daily_results.get("subnets_failure_rate", []): + subnets_failure_rate_dict[str(principal)] = failure_rate - elif "Err" in result: - error_msg = result["Err"] - logger.warning(f"Canister returned error for {date_str}: {error_msg}") - return {} - else: - logger.warning(f"Unexpected response format for {date_str}: {result}") - return {} + result_dict = { + "provider_results": provider_results_dict, + "subnets_failure_rate": subnets_failure_rate_dict, + } - except Exception as e: - logger.warning( - f"Failed to fetch rewards for {date_str}: {e}", exc_info=True - ) - return {} + logger.info( + f"Successfully fetched rewards for {date} ({len(provider_results_dict)} providers)" + ) + return result_dict def get_latest_governance_reward_event(self) -> Optional[float]: """Fetch latest governance reward event timestamp from governance canister""" @@ -249,7 +220,7 @@ def get_latest_governance_reward_event(self) -> Optional[float]: # Make query response = self.agent.query_raw( - self.GOVERNANCE_CANISTER_ID, + GOVERNANCE_CANISTER_ID, "list_node_provider_rewards", arg_bytes, ) @@ -281,7 +252,9 @@ class NodeRewardsPusher: def __init__(self, victoria_url: str): self.victoria_url = victoria_url - self.ic_client = ICCanisterClient() + self.nrc_clients = [ + NodeRewardsClient(IC_URL, c) for c in NODE_REWARDS_CANISTER_IDS + ] @staticmethod def _unwrap_optional(value): @@ -326,7 +299,7 @@ def push_metrics_for_date(self, date_str: str) -> bool: ) # Fetch data from IC canister - daily_results = self.ic_client.get_rewards_daily(date_str) + daily_results = self.nrc_clients[0].get_rewards_daily(date_str) if not daily_results: logger.warning(f"⚠️ No data available for {date_str}") @@ -416,7 +389,7 @@ def push_metrics_for_date(self, date_str: str) -> bool: ) # Governance timestamp - gov_timestamp = self.ic_client.get_latest_governance_reward_event() + gov_timestamp = self.nrc_clients[0].get_latest_governance_reward_event() if gov_timestamp: metrics_lines.append( f"governance_latest_reward_event_timestamp_seconds {gov_timestamp} {noon_timestamp_ms}" From ef41e8ade3e36edcd581645a475520e161229424 Mon Sep 17 00:00:00 2001 From: nikolamilosa Date: Wed, 12 Nov 2025 21:29:47 +0100 Subject: [PATCH 3/5] further refactoring --- .../node_rewards_ingester.py | 305 ++++++++---------- 1 file changed, 143 insertions(+), 162 deletions(-) diff --git a/docker/tools/node-rewards-scheduler/node_rewards_ingester.py b/docker/tools/node-rewards-scheduler/node_rewards_ingester.py index 862ac37..f4cd813 100644 --- a/docker/tools/node-rewards-scheduler/node_rewards_ingester.py +++ b/docker/tools/node-rewards-scheduler/node_rewards_ingester.py @@ -43,7 +43,7 @@ } ) -REQUEST_TYPE = Types.Record({"day": DATE_UTC_TYPE}) +GET_REWARDS_DAILY_REQUEST_TYPE = Types.Record({"day": DATE_UTC_TYPE}) NODE_METRICS_DAILY_TYPE = Types.Record( { @@ -130,6 +130,10 @@ } ) +LIST_NODE_PROVIDER_REWARDS_REQUEST_TYPE = Types.Record( + {"date_filter": Types.Opt(Types.Nat64)} +) + class NodeRewardsClient: """Client for interacting with the node rewards canister""" @@ -151,7 +155,9 @@ def get_rewards_daily(self, date: str) -> Dict[str, Any]: "day": parsed_date.day, } } - arg_bytes = encode([{"type": REQUEST_TYPE, "value": arg_value}]) + arg_bytes = encode( + [{"type": GET_REWARDS_DAILY_REQUEST_TYPE, "value": arg_value}] + ) response = self.agent.query_raw( self.canister_id, @@ -202,47 +208,36 @@ def get_rewards_daily(self, date: str) -> Dict[str, Any]: def get_latest_governance_reward_event(self) -> Optional[float]: """Fetch latest governance reward event timestamp from governance canister""" - try: - request_type = Types.Record({"date_filter": Types.Opt(Types.Nat64)}) - - # Build request with no date filter to get all rewards - # In ic-py, optional None is represented as an empty list [] - arg_bytes = encode( - [ - { - "type": request_type, - "value": { - "date_filter": [] # Empty list = None for optional types - }, - } - ] - ) - - # Make query - response = self.agent.query_raw( - GOVERNANCE_CANISTER_ID, - "list_node_provider_rewards", - arg_bytes, - ) + arg_bytes = encode( + [ + { + "type": LIST_NODE_PROVIDER_REWARDS_REQUEST_TYPE, + "value": { + "date_filter": [] # Empty list = None for optional types + }, + } + ] + ) - # Extract value from response - if not response or len(response) == 0: - logger.warning("Empty response from governance canister") - return None + response = self.agent.query_raw( + GOVERNANCE_CANISTER_ID, + "list_node_provider_rewards", + arg_bytes, + ) - result = response[0].get("value", {}) - rewards_list = result.get("rewards", []) + if not response or len(response) == 0: + logger.warning("Empty response from governance canister") + return None - # Rewards are in descending order (latest first) - if rewards_list and len(rewards_list) > 0: - first_reward = rewards_list[0] - timestamp = first_reward.get("timestamp") - if timestamp: - logger.info(f"Latest governance reward timestamp: {timestamp}") - return float(timestamp) + result = response[0].get("value", {}) + rewards_list = result.get("rewards", []) - except Exception as e: - logger.warning(f"Failed to fetch governance rewards: {e}", exc_info=True) + if rewards_list and len(rewards_list) > 0: + first_reward = rewards_list[0] + timestamp = first_reward.get("timestamp") + if timestamp: + logger.info(f"Latest governance reward timestamp: {timestamp}") + return float(timestamp) return None @@ -280,152 +275,138 @@ def wait_for_victoria_metrics(self): logger.info(f" Waiting for VictoriaMetrics at {self.victoria_url}...") time.sleep(2) - def push_metrics_for_date(self, date_str: str) -> bool: + def push_metrics_for_date(self, date: str): """ Fetch node rewards data from IC canisters and push to VictoriaMetrics for a specific date Returns True if successful, False otherwise """ - try: - logger.info(f"Pushing node rewards data for {date_str}") + logger.info(f"Pushing node rewards data for {date}") + target_date = datetime.strptime(date, "%Y-%m-%d") - # Parse target date - target_date = datetime.strptime(date_str, "%Y-%m-%d") + noon = target_date.replace(hour=12, minute=0, second=0, microsecond=0) + noon_timestamp_ms = int(noon.replace(tzinfo=timezone.utc).timestamp() * 1000) - # Calculate noon timestamp in milliseconds - noon_dt = target_date.replace(hour=12, minute=0, second=0, microsecond=0) - noon_timestamp_ms = int( - noon_dt.replace(tzinfo=timezone.utc).timestamp() * 1000 - ) + # Fetch data from IC canister + daily_results = self.nrc_clients[0].get_rewards_daily(date) - # Fetch data from IC canister - daily_results = self.nrc_clients[0].get_rewards_daily(date_str) + if not daily_results: + raise ValueError(f"⚠️ No data available for {date}") - if not daily_results: - logger.warning(f"⚠️ No data available for {date_str}") - return False + metrics_lines = [] - # Build Prometheus metrics - metrics_lines = [] + # Provider-level metrics + provider_results = daily_results.get("provider_results", {}) + for provider_id, provider_rewards in provider_results.items(): + provider_id_str = str(provider_id) - # Provider-level metrics - provider_results = daily_results.get("provider_results", {}) - for provider_id, provider_rewards in provider_results.items(): - provider_id_str = str(provider_id) + # nodes_count + nodes_count = len(provider_rewards.get("daily_nodes_rewards", [])) + metrics_lines.append( + f'nodes_count{{provider_id="{provider_id_str}"}} {nodes_count} {noon_timestamp_ms}' + ) - # nodes_count - nodes_count = len(provider_rewards.get("daily_nodes_rewards", [])) + # base_rewards + base_rewards = self._unwrap_optional( + provider_rewards.get("total_base_rewards_xdr_permyriad") + ) + if base_rewards is not None: metrics_lines.append( - f'nodes_count{{provider_id="{provider_id_str}"}} {nodes_count} {noon_timestamp_ms}' + f'total_base_rewards_xdr_permyriad{{provider_id="{provider_id_str}"}} {base_rewards} {noon_timestamp_ms}' ) - # base_rewards - base_rewards = self._unwrap_optional( - provider_rewards.get("total_base_rewards_xdr_permyriad") + # adjusted_rewards + adjusted_rewards = self._unwrap_optional( + provider_rewards.get("total_adjusted_rewards_xdr_permyriad") + ) + if adjusted_rewards is not None: + metrics_lines.append( + f'total_adjusted_rewards_xdr_permyriad{{provider_id="{provider_id_str}"}} {adjusted_rewards} {noon_timestamp_ms}' ) - if base_rewards is not None: - metrics_lines.append( - f'total_base_rewards_xdr_permyriad{{provider_id="{provider_id_str}"}} {base_rewards} {noon_timestamp_ms}' - ) - # adjusted_rewards - adjusted_rewards = self._unwrap_optional( - provider_rewards.get("total_adjusted_rewards_xdr_permyriad") + # Node-level metrics + for node_result in provider_rewards.get("daily_nodes_rewards", []): + node_id = self._unwrap_optional(node_result.get("node_id")) + node_id_str = str(node_id) if node_id else "" + + # daily_node_failure_rate is optional and contains a variant + failure_rate_data = self._unwrap_optional( + node_result.get("daily_node_failure_rate") ) - if adjusted_rewards is not None: - metrics_lines.append( - f'total_adjusted_rewards_xdr_permyriad{{provider_id="{provider_id_str}"}} {adjusted_rewards} {noon_timestamp_ms}' + if ( + isinstance(failure_rate_data, dict) + and "SubnetMember" in failure_rate_data + ): + # Extract node_metrics from SubnetMember variant + subnet_member = failure_rate_data["SubnetMember"] + node_metrics = self._unwrap_optional( + subnet_member.get("node_metrics") ) - # Node-level metrics - for node_result in provider_rewards.get("daily_nodes_rewards", []): - node_id = self._unwrap_optional(node_result.get("node_id")) - node_id_str = str(node_id) if node_id else "" - - # daily_node_failure_rate is optional and contains a variant - failure_rate_data = self._unwrap_optional( - node_result.get("daily_node_failure_rate") - ) - if ( - isinstance(failure_rate_data, dict) - and "SubnetMember" in failure_rate_data - ): - # Extract node_metrics from SubnetMember variant - subnet_member = failure_rate_data["SubnetMember"] - node_metrics = self._unwrap_optional( - subnet_member.get("node_metrics") + if node_metrics: + subnet_id = self._unwrap_optional( + node_metrics.get("subnet_assigned") ) + subnet_id_str = str(subnet_id) if subnet_id else "" - if node_metrics: - subnet_id = self._unwrap_optional( - node_metrics.get("subnet_assigned") + # original_failure_rate + original_fr = self._unwrap_optional( + node_metrics.get("original_failure_rate") + ) + if original_fr is not None: + metrics_lines.append( + f'original_failure_rate{{provider_id="{provider_id_str}",node_id="{node_id_str}",subnet_id="{subnet_id_str}"}} {original_fr} {noon_timestamp_ms}' ) - subnet_id_str = str(subnet_id) if subnet_id else "" - # original_failure_rate - original_fr = self._unwrap_optional( - node_metrics.get("original_failure_rate") - ) - if original_fr is not None: - metrics_lines.append( - f'original_failure_rate{{provider_id="{provider_id_str}",node_id="{node_id_str}",subnet_id="{subnet_id_str}"}} {original_fr} {noon_timestamp_ms}' - ) - - # relative_failure_rate - relative_fr = self._unwrap_optional( - node_metrics.get("relative_failure_rate") + # relative_failure_rate + relative_fr = self._unwrap_optional( + node_metrics.get("relative_failure_rate") + ) + if relative_fr is not None: + metrics_lines.append( + f'relative_failure_rate{{provider_id="{provider_id_str}",node_id="{node_id_str}",subnet_id="{subnet_id_str}"}} {relative_fr} {noon_timestamp_ms}' ) - if relative_fr is not None: - metrics_lines.append( - f'relative_failure_rate{{provider_id="{provider_id_str}",node_id="{node_id_str}",subnet_id="{subnet_id_str}"}} {relative_fr} {noon_timestamp_ms}' - ) - - # Subnet-level metrics - subnets_failure_rate = daily_results.get("subnets_failure_rate", {}) - for subnet_id, failure_rate in subnets_failure_rate.items(): - subnet_id_str = str(subnet_id) - metrics_lines.append( - f'subnet_failure_rate{{subnet_id="{subnet_id_str}"}} {failure_rate} {noon_timestamp_ms}' - ) - # Governance timestamp - gov_timestamp = self.nrc_clients[0].get_latest_governance_reward_event() - if gov_timestamp: - metrics_lines.append( - f"governance_latest_reward_event_timestamp_seconds {gov_timestamp} {noon_timestamp_ms}" - ) + # Subnet-level metrics + subnets_failure_rate = daily_results.get("subnets_failure_rate", {}) + for subnet_id, failure_rate in subnets_failure_rate.items(): + subnet_id_str = str(subnet_id) + metrics_lines.append( + f'subnet_failure_rate{{subnet_id="{subnet_id_str}"}} {failure_rate} {noon_timestamp_ms}' + ) - # Push to VictoriaMetrics - if metrics_lines: - metrics_payload = "\n".join(metrics_lines) + "\n" - import_url = urljoin( - self.victoria_url.rstrip("/") + "/", "api/v1/import/prometheus" - ) + # Governance timestamp + gov_timestamp = self.nrc_clients[0].get_latest_governance_reward_event() + if gov_timestamp: + metrics_lines.append( + f"governance_latest_reward_event_timestamp_seconds {gov_timestamp} {noon_timestamp_ms}" + ) - response = requests.post( - import_url, - data=metrics_payload.encode("utf-8"), - headers={"Content-Type": "text/plain"}, - timeout=30, - ) + if not metrics_lines: + raise ValueError("After evaluation there were no metrics to upload") - if response.status_code in (200, 204): - logger.info( - f"✅ Successfully pushed data for {date_str} ({len(metrics_lines)} metrics)" - ) - return True - else: - logger.error( - f"❌ Failed to push to VictoriaMetrics: {response.status_code} - {response.text}" - ) - return False - else: - logger.warning(f"No metrics to push for {date_str}") - return False + metrics_payload = "\n".join(metrics_lines) + "\n" + import_url = urljoin( + self.victoria_url.rstrip("/") + "/", "api/v1/import/prometheus" + ) - except Exception as e: - logger.error(f"❌ Error processing {date_str}: {e}", exc_info=True) - return False + response = requests.post( + import_url, + data=metrics_payload.encode("utf-8"), + headers={"Content-Type": "text/plain"}, + timeout=30, + ) + try: + response.raise_for_status() + except Exception: + logger.error( + f"❌ Failed to push to VictoriaMetrics: {response.status_code} - {response.text}" + ) + raise + + logger.info( + f"✅ Successfully pushed data for {date} ({len(metrics_lines)} metrics)" + ) def backfill(self, days: int = 40): """Backfill historical data""" @@ -433,13 +414,14 @@ def backfill(self, days: int = 40): logger.info(f"Starting backfill of last {days} days...") for i in range(days, 0, -1): - date_obj = datetime.now(timezone.utc) - timedelta(days=i) - date_str = date_obj.strftime("%Y-%m-%d") + date = datetime.now(timezone.utc) - timedelta(days=i) + date = date.strftime("%Y-%m-%d") - logger.info( - f"[{days - i + 1:2d}/{days}] Backfilling data for {date_str}..." - ) - self.push_metrics_for_date(date_str) + logger.info(f"[{days - i + 1:2d}/{days}] Backfilling data for {date}...") + try: + self.push_metrics_for_date(date) + except Exception as e: + logger.error(f"Failed to backfill data due to: {e}") logger.info("✅ Backfill complete!") @@ -468,7 +450,6 @@ def run_daily_scheduler(self): while True: try: - # Wait until 00:05 UTC self.wait_until_next_run() # Push yesterday's data From 2071ed14a81ba93510172363253ff96b05a0a5b4 Mon Sep 17 00:00:00 2001 From: nikolamilosa Date: Wed, 12 Nov 2025 21:54:36 +0100 Subject: [PATCH 4/5] further refactoring --- .../node_rewards_ingester.py | 90 ++++++++++++------- 1 file changed, 60 insertions(+), 30 deletions(-) diff --git a/docker/tools/node-rewards-scheduler/node_rewards_ingester.py b/docker/tools/node-rewards-scheduler/node_rewards_ingester.py index f4cd813..914ba23 100644 --- a/docker/tools/node-rewards-scheduler/node_rewards_ingester.py +++ b/docker/tools/node-rewards-scheduler/node_rewards_ingester.py @@ -258,6 +258,10 @@ def _unwrap_optional(value): return value[0] if len(value) > 0 else None return value + @staticmethod + def _make_line(metric_name: str, value: int, ts: int, **kwargs) -> str: + return f"{metric_name}{{ {' '.join([f'{key}="{value}"' for key, value in kwargs.items()])} }} {value} {ts}" + def wait_for_victoria_metrics(self): """Wait for VictoriaMetrics to be ready""" @@ -295,33 +299,48 @@ def push_metrics_for_date(self, date: str): metrics_lines = [] + # Helper function to not repeat the labels all the + # time and to single out the place for changing labels + def add_line_helper(metric_name: str, value, **kwargs): + metrics_lines.append( + self._make_line( + metric_name, + value, + noon_timestamp_ms, + canister_id=self.nrc_clients[0].canister_id, + **kwargs, + ) + ) + # Provider-level metrics provider_results = daily_results.get("provider_results", {}) for provider_id, provider_rewards in provider_results.items(): provider_id_str = str(provider_id) + def add_line_helper_with_provider(metric_name: str, value: int, **kwargs): + add_line_helper( + metric_name, value, provider_id=provider_id_str, **kwargs + ) + # nodes_count nodes_count = len(provider_rewards.get("daily_nodes_rewards", [])) - metrics_lines.append( - f'nodes_count{{provider_id="{provider_id_str}"}} {nodes_count} {noon_timestamp_ms}' - ) + add_line_helper_with_provider("nodes_count", nodes_count) # base_rewards base_rewards = self._unwrap_optional( provider_rewards.get("total_base_rewards_xdr_permyriad") ) if base_rewards is not None: - metrics_lines.append( - f'total_base_rewards_xdr_permyriad{{provider_id="{provider_id_str}"}} {base_rewards} {noon_timestamp_ms}' + add_line_helper_with_provider( + "total_base_rewards_xdr_permyriad", base_rewards ) - # adjusted_rewards adjusted_rewards = self._unwrap_optional( provider_rewards.get("total_adjusted_rewards_xdr_permyriad") ) if adjusted_rewards is not None: - metrics_lines.append( - f'total_adjusted_rewards_xdr_permyriad{{provider_id="{provider_id_str}"}} {adjusted_rewards} {noon_timestamp_ms}' + add_line_helper_with_provider( + "total_adjusted_rewards_xdr_permyriad", adjusted_rewards ) # Node-level metrics @@ -343,29 +362,37 @@ def push_metrics_for_date(self, date: str): subnet_member.get("node_metrics") ) - if node_metrics: - subnet_id = self._unwrap_optional( - node_metrics.get("subnet_assigned") - ) - subnet_id_str = str(subnet_id) if subnet_id else "" + if not node_metrics: + continue - # original_failure_rate - original_fr = self._unwrap_optional( - node_metrics.get("original_failure_rate") + subnet_id = self._unwrap_optional( + node_metrics.get("subnet_assigned") + ) + subnet_id_str = str(subnet_id) if subnet_id else "" + + # original_failure_rate + original_fr = self._unwrap_optional( + node_metrics.get("original_failure_rate") + ) + if original_fr is not None: + add_line_helper_with_provider( + "original_failure_rate", + original_fr, + node_id=node_id_str, + subnet_id=subnet_id_str, ) - if original_fr is not None: - metrics_lines.append( - f'original_failure_rate{{provider_id="{provider_id_str}",node_id="{node_id_str}",subnet_id="{subnet_id_str}"}} {original_fr} {noon_timestamp_ms}' - ) - - # relative_failure_rate - relative_fr = self._unwrap_optional( - node_metrics.get("relative_failure_rate") + + # relative_failure_rate + relative_fr = self._unwrap_optional( + node_metrics.get("relative_failure_rate") + ) + if relative_fr is not None: + add_line_helper_with_provider( + "relative_failure_rate", + relative_fr, + node_id=node_id_str, + subnet_id=subnet_id_str, ) - if relative_fr is not None: - metrics_lines.append( - f'relative_failure_rate{{provider_id="{provider_id_str}",node_id="{node_id_str}",subnet_id="{subnet_id_str}"}} {relative_fr} {noon_timestamp_ms}' - ) # Subnet-level metrics subnets_failure_rate = daily_results.get("subnets_failure_rate", {}) @@ -374,12 +401,15 @@ def push_metrics_for_date(self, date: str): metrics_lines.append( f'subnet_failure_rate{{subnet_id="{subnet_id_str}"}} {failure_rate} {noon_timestamp_ms}' ) + add_line_helper( + "subnets_failure_rate", failure_rate, subnet_id=subnet_id_str + ) # Governance timestamp gov_timestamp = self.nrc_clients[0].get_latest_governance_reward_event() if gov_timestamp: - metrics_lines.append( - f"governance_latest_reward_event_timestamp_seconds {gov_timestamp} {noon_timestamp_ms}" + add_line_helper( + "governance_latest_reward_event_timestamp_seconds", gov_timestamp ) if not metrics_lines: From e1acff95351d9a9a364fcf9c7ffcb565ab42228b Mon Sep 17 00:00:00 2001 From: nikolamilosa Date: Wed, 12 Nov 2025 22:02:15 +0100 Subject: [PATCH 5/5] further refactoring --- .../node_rewards_ingester.py | 208 +++++++++--------- 1 file changed, 105 insertions(+), 103 deletions(-) diff --git a/docker/tools/node-rewards-scheduler/node_rewards_ingester.py b/docker/tools/node-rewards-scheduler/node_rewards_ingester.py index 914ba23..bed42f4 100644 --- a/docker/tools/node-rewards-scheduler/node_rewards_ingester.py +++ b/docker/tools/node-rewards-scheduler/node_rewards_ingester.py @@ -30,7 +30,8 @@ # These have to be tuples NODE_REWARDS_CANISTER_IDS = [ ("uuew5-iiaaa-aaaaa-qbx4q-cai"), # Dev - ("sgymv-uiaaa-aaaaa-aaaia-cai"), # Prod + # TODO: uncomment to enable prod + # ("sgymv-uiaaa-aaaaa-aaaia-cai"), # Prod ] GOVERNANCE_CANISTER_ID = "rrkah-fqaaa-aaaaa-aaaaq-cai" # NNS Governance canister @@ -291,126 +292,127 @@ def push_metrics_for_date(self, date: str): noon = target_date.replace(hour=12, minute=0, second=0, microsecond=0) noon_timestamp_ms = int(noon.replace(tzinfo=timezone.utc).timestamp() * 1000) - # Fetch data from IC canister - daily_results = self.nrc_clients[0].get_rewards_daily(date) - - if not daily_results: - raise ValueError(f"⚠️ No data available for {date}") - metrics_lines = [] - - # Helper function to not repeat the labels all the - # time and to single out the place for changing labels - def add_line_helper(metric_name: str, value, **kwargs): - metrics_lines.append( - self._make_line( - metric_name, - value, - noon_timestamp_ms, - canister_id=self.nrc_clients[0].canister_id, - **kwargs, + for client in self.nrc_clients: + daily_results = client.get_rewards_daily(date) + + if not daily_results: + raise ValueError(f"⚠️ No data available for {date}") + + # Helper function to not repeat the labels all the + # time and to single out the place for changing labels + def add_line_helper(metric_name: str, value, **kwargs): + metrics_lines.append( + self._make_line( + metric_name, + value, + noon_timestamp_ms, + canister_id=client.canister_id, + **kwargs, + ) ) - ) - - # Provider-level metrics - provider_results = daily_results.get("provider_results", {}) - for provider_id, provider_rewards in provider_results.items(): - provider_id_str = str(provider_id) - def add_line_helper_with_provider(metric_name: str, value: int, **kwargs): - add_line_helper( - metric_name, value, provider_id=provider_id_str, **kwargs - ) + # Provider-level metrics + provider_results = daily_results.get("provider_results", {}) + for provider_id, provider_rewards in provider_results.items(): + provider_id_str = str(provider_id) - # nodes_count - nodes_count = len(provider_rewards.get("daily_nodes_rewards", [])) - add_line_helper_with_provider("nodes_count", nodes_count) + def add_line_helper_with_provider( + metric_name: str, value: int, **kwargs + ): + add_line_helper( + metric_name, value, provider_id=provider_id_str, **kwargs + ) - # base_rewards - base_rewards = self._unwrap_optional( - provider_rewards.get("total_base_rewards_xdr_permyriad") - ) - if base_rewards is not None: - add_line_helper_with_provider( - "total_base_rewards_xdr_permyriad", base_rewards - ) + # nodes_count + nodes_count = len(provider_rewards.get("daily_nodes_rewards", [])) + add_line_helper_with_provider("nodes_count", nodes_count) - adjusted_rewards = self._unwrap_optional( - provider_rewards.get("total_adjusted_rewards_xdr_permyriad") - ) - if adjusted_rewards is not None: - add_line_helper_with_provider( - "total_adjusted_rewards_xdr_permyriad", adjusted_rewards + # base_rewards + base_rewards = self._unwrap_optional( + provider_rewards.get("total_base_rewards_xdr_permyriad") ) + if base_rewards is not None: + add_line_helper_with_provider( + "total_base_rewards_xdr_permyriad", base_rewards + ) - # Node-level metrics - for node_result in provider_rewards.get("daily_nodes_rewards", []): - node_id = self._unwrap_optional(node_result.get("node_id")) - node_id_str = str(node_id) if node_id else "" - - # daily_node_failure_rate is optional and contains a variant - failure_rate_data = self._unwrap_optional( - node_result.get("daily_node_failure_rate") + adjusted_rewards = self._unwrap_optional( + provider_rewards.get("total_adjusted_rewards_xdr_permyriad") ) - if ( - isinstance(failure_rate_data, dict) - and "SubnetMember" in failure_rate_data - ): - # Extract node_metrics from SubnetMember variant - subnet_member = failure_rate_data["SubnetMember"] - node_metrics = self._unwrap_optional( - subnet_member.get("node_metrics") + if adjusted_rewards is not None: + add_line_helper_with_provider( + "total_adjusted_rewards_xdr_permyriad", adjusted_rewards ) - if not node_metrics: - continue - - subnet_id = self._unwrap_optional( - node_metrics.get("subnet_assigned") - ) - subnet_id_str = str(subnet_id) if subnet_id else "" + # Node-level metrics + for node_result in provider_rewards.get("daily_nodes_rewards", []): + node_id = self._unwrap_optional(node_result.get("node_id")) + node_id_str = str(node_id) if node_id else "" - # original_failure_rate - original_fr = self._unwrap_optional( - node_metrics.get("original_failure_rate") + # daily_node_failure_rate is optional and contains a variant + failure_rate_data = self._unwrap_optional( + node_result.get("daily_node_failure_rate") ) - if original_fr is not None: - add_line_helper_with_provider( - "original_failure_rate", - original_fr, - node_id=node_id_str, - subnet_id=subnet_id_str, + if ( + isinstance(failure_rate_data, dict) + and "SubnetMember" in failure_rate_data + ): + # Extract node_metrics from SubnetMember variant + subnet_member = failure_rate_data["SubnetMember"] + node_metrics = self._unwrap_optional( + subnet_member.get("node_metrics") ) - # relative_failure_rate - relative_fr = self._unwrap_optional( - node_metrics.get("relative_failure_rate") - ) - if relative_fr is not None: - add_line_helper_with_provider( - "relative_failure_rate", - relative_fr, - node_id=node_id_str, - subnet_id=subnet_id_str, + if not node_metrics: + continue + + subnet_id = self._unwrap_optional( + node_metrics.get("subnet_assigned") ) + subnet_id_str = str(subnet_id) if subnet_id else "" - # Subnet-level metrics - subnets_failure_rate = daily_results.get("subnets_failure_rate", {}) - for subnet_id, failure_rate in subnets_failure_rate.items(): - subnet_id_str = str(subnet_id) - metrics_lines.append( - f'subnet_failure_rate{{subnet_id="{subnet_id_str}"}} {failure_rate} {noon_timestamp_ms}' - ) - add_line_helper( - "subnets_failure_rate", failure_rate, subnet_id=subnet_id_str - ) + # original_failure_rate + original_fr = self._unwrap_optional( + node_metrics.get("original_failure_rate") + ) + if original_fr is not None: + add_line_helper_with_provider( + "original_failure_rate", + original_fr, + node_id=node_id_str, + subnet_id=subnet_id_str, + ) + + # relative_failure_rate + relative_fr = self._unwrap_optional( + node_metrics.get("relative_failure_rate") + ) + if relative_fr is not None: + add_line_helper_with_provider( + "relative_failure_rate", + relative_fr, + node_id=node_id_str, + subnet_id=subnet_id_str, + ) + + # Subnet-level metrics + subnets_failure_rate = daily_results.get("subnets_failure_rate", {}) + for subnet_id, failure_rate in subnets_failure_rate.items(): + subnet_id_str = str(subnet_id) + metrics_lines.append( + f'subnet_failure_rate{{subnet_id="{subnet_id_str}"}} {failure_rate} {noon_timestamp_ms}' + ) + add_line_helper( + "subnets_failure_rate", failure_rate, subnet_id=subnet_id_str + ) - # Governance timestamp - gov_timestamp = self.nrc_clients[0].get_latest_governance_reward_event() - if gov_timestamp: - add_line_helper( - "governance_latest_reward_event_timestamp_seconds", gov_timestamp - ) + # Governance timestamp + gov_timestamp = self.nrc_clients[0].get_latest_governance_reward_event() + if gov_timestamp: + add_line_helper( + "governance_latest_reward_event_timestamp_seconds", gov_timestamp + ) if not metrics_lines: raise ValueError("After evaluation there were no metrics to upload")