From 550fb71df6c1c9ddc5839429350921b6ea4067cf Mon Sep 17 00:00:00 2001 From: Pietro Date: Thu, 13 Nov 2025 18:02:55 +0100 Subject: [PATCH] Fix: Ingester script retries + Add Gov. types --- .../node_rewards_ingester.py | 333 +++++++++++------- 1 file changed, 214 insertions(+), 119 deletions(-) diff --git a/tools/node-rewards-scheduler/node_rewards_ingester.py b/tools/node-rewards-scheduler/node_rewards_ingester.py index 9e149dd..fce60c5 100644 --- a/tools/node-rewards-scheduler/node_rewards_ingester.py +++ b/tools/node-rewards-scheduler/node_rewards_ingester.py @@ -7,6 +7,7 @@ import os import time from datetime import datetime, timedelta, timezone +from functools import wraps from typing import Any, Dict, Optional from urllib.parse import urljoin @@ -25,14 +26,72 @@ logger = logging.getLogger(__name__) +def retry_on_timeout(max_attempts=3, initial_delay=2, backoff_factor=2): + """ + Decorator to retry a function on timeout or connection errors with exponential backoff + + Args: + max_attempts: Maximum number of retry attempts + initial_delay: Initial delay in seconds before first retry + backoff_factor: Multiplier for delay between retries + """ + def decorator(func): + @wraps(func) + def wrapper(*args, **kwargs): + delay = initial_delay + last_exception = None + + for attempt in range(max_attempts): + try: + return func(*args, **kwargs) + except ( + requests.exceptions.Timeout, + requests.exceptions.ConnectionError, + TimeoutError, + ConnectionError, + Exception + ) as e: + last_exception = e + error_msg = str(e) + + # Check if it's a timeout-related error + is_timeout = ( + isinstance(e, (requests.exceptions.Timeout, TimeoutError)) or + "timed out" in error_msg.lower() or + "timeout" in error_msg.lower() + ) + + # Check if it's a connection error + is_connection_error = isinstance(e, (requests.exceptions.ConnectionError, ConnectionError)) + + # Only retry on timeout or connection errors + if not (is_timeout or is_connection_error): + # Re-raise if it's not a retryable error + raise + + if attempt < max_attempts - 1: + logger.warning( + f"Attempt {attempt + 1}/{max_attempts} failed with {type(e).__name__}: {error_msg}. " + f"Retrying in {delay} seconds..." + ) + time.sleep(delay) + delay *= backoff_factor + else: + logger.error( + f"All {max_attempts} attempts failed. Last error: {error_msg}" + ) + raise last_exception + + # This should never be reached, but just in case + raise last_exception + + return wrapper + return decorator + + IC_URL = "https://ic0.app" -# These have to be tuples -NODE_REWARDS_CANISTER_IDS = [ - ("uuew5-iiaaa-aaaaa-qbx4q-cai"), # Dev - # TODO: uncomment to enable prod - # ("sgymv-uiaaa-aaaaa-aaaia-cai"), # Prod -] +NODE_REWARDS_CANISTER_ID = "sgymv-uiaaa-aaaaa-aaaia-cai" GOVERNANCE_CANISTER_ID = "rrkah-fqaaa-aaaaa-aaaaq-cai" # NNS Governance canister ###################### TYPE DEFINITIONS ########################### @@ -124,7 +183,7 @@ } ) -RETURN_TYPE = Types.Variant( +GET_NODE_PROVIDERS_CALCULATION_RESPONSE_TYPE = Types.Variant( { "Ok": DAILY_RESULTS_TYPE, "Err": Types.Text, @@ -135,6 +194,27 @@ {"date_filter": Types.Opt(Types.Nat64)} ) +DATE_RANGE_FILTER_TYPE = Types.Record({ + "start_timestamp_seconds": Types.Opt(Types.Nat64), + "end_timestamp_seconds": Types.Opt(Types.Nat64), +}) + +MONTHLY_NODE_PROVIDER_REWARDS_TYPE = Types.Record({ + 'timestamp': Types.Nat64, + "start_date": Types.Opt(Types.Record({"year": Types.Nat32, "month": Types.Nat32, "day": Types.Nat32})), + "end_date": Types.Opt(Types.Record({"year": Types.Nat32, "month": Types.Nat32, "day": Types.Nat32})), + "rewards": Types.Vec(Types.Record({})), # Simplified, not parsing individual rewards + "xdr_conversion_rate": Types.Opt(Types.Record({})), + "minimum_xdr_permyriad_per_icp": Types.Opt(Types.Nat64), + "maximum_node_provider_rewards_e8s": Types.Opt(Types.Nat64), + "registry_version": Types.Opt(Types.Nat64), + "node_providers": Types.Vec(Types.Record({})), +}) + +LIST_NODE_PROVIDER_REWARDS_RESPONSE_TYPE = Types.Record({ + "rewards": Types.Vec(MONTHLY_NODE_PROVIDER_REWARDS_TYPE), +}) + class NodeRewardsClient: """Client for interacting with the node rewards canister""" @@ -146,6 +226,7 @@ def __init__(self, ic_url: str, canister_id: str): self.agent = Agent(self.identity, self.client) self.canister_id = canister_id + @retry_on_timeout(max_attempts=5, initial_delay=3, backoff_factor=2) 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") @@ -164,7 +245,7 @@ def get_rewards_daily(self, date: str) -> Dict[str, Any]: self.canister_id, "get_node_providers_rewards_calculation", arg_bytes, - RETURN_TYPE, + GET_NODE_PROVIDERS_CALCULATION_RESPONSE_TYPE, ) if not response or len(response) == 0: @@ -206,6 +287,7 @@ def get_rewards_daily(self, date: str) -> Dict[str, Any]: ) return result_dict + @retry_on_timeout(max_attempts=5, initial_delay=3, backoff_factor=2) def get_latest_governance_reward_event(self) -> Optional[float]: """Fetch latest governance reward event timestamp from governance canister""" @@ -224,6 +306,7 @@ def get_latest_governance_reward_event(self) -> Optional[float]: GOVERNANCE_CANISTER_ID, "list_node_provider_rewards", arg_bytes, + LIST_NODE_PROVIDER_REWARDS_RESPONSE_TYPE, ) if not response or len(response) == 0: @@ -248,9 +331,7 @@ class NodeRewardsPusher: def __init__(self, victoria_url: str): self.victoria_url = victoria_url - self.nrc_clients = [ - NodeRewardsClient(IC_URL, c) for c in NODE_REWARDS_CANISTER_IDS - ] + self.nrc_client = NodeRewardsClient(IC_URL, NODE_REWARDS_CANISTER_ID) @staticmethod def _unwrap_optional(value): @@ -261,7 +342,7 @@ def _unwrap_optional(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}" + 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""" @@ -280,6 +361,7 @@ def wait_for_victoria_metrics(self): logger.info(f" Waiting for VictoriaMetrics at {self.victoria_url}...") time.sleep(2) + @retry_on_timeout(max_attempts=3, initial_delay=5, backoff_factor=2) def push_metrics_for_date(self, date: str): """ Fetch node rewards data from IC canisters and push to VictoriaMetrics for a specific date @@ -289,130 +371,145 @@ def push_metrics_for_date(self, date: str): logger.info(f"Pushing node rewards data for {date}") target_date = datetime.strptime(date, "%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) + target_dt = target_date.replace(hour=0, minute=0, second=0, microsecond=0) + timestamp_ms = int(target_dt.replace(tzinfo=timezone.utc).timestamp() * 1000) metrics_lines = [] - 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, - ) + daily_results = self.nrc_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, + timestamp_ms, + canister_id=self.nrc_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) + # 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 - ) + 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", [])) - add_line_helper_with_provider("nodes_count", nodes_count) + # nodes_count + nodes_count = len(provider_rewards.get("daily_nodes_rewards", [])) + add_line_helper_with_provider("nodes_count", nodes_count) - # base_rewards - base_rewards = self._unwrap_optional( - provider_rewards.get("total_base_rewards_xdr_permyriad") + # 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 ) - if base_rewards is not None: + + 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 + ) + + # 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 "" + + # performance_multiplier + performance_multiplier = self._unwrap_optional( + node_result.get("performance_multiplier") + ) + if performance_multiplier is not None: add_line_helper_with_provider( - "total_base_rewards_xdr_permyriad", base_rewards + "performance_multiplier", + performance_multiplier, + node_id=node_id_str, ) - adjusted_rewards = self._unwrap_optional( - provider_rewards.get("total_adjusted_rewards_xdr_permyriad") + # 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: - add_line_helper_with_provider( - "total_adjusted_rewards_xdr_permyriad", adjusted_rewards + 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 "" + if not node_metrics: + continue - # daily_node_failure_rate is optional and contains a variant - failure_rate_data = self._unwrap_optional( - node_result.get("daily_node_failure_rate") + subnet_id = self._unwrap_optional( + node_metrics.get("subnet_assigned") ) - 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") - ) + subnet_id_str = str(subnet_id) if subnet_id else "" - if not node_metrics: - continue - - 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: + add_line_helper_with_provider( + "original_failure_rate", + original_fr, + node_id=node_id_str, + subnet_id=subnet_id_str, ) - 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, - ) - - # 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: - 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 - ) + # 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) + add_line_helper( + "subnets_failure_rate", failure_rate, subnet_id=subnet_id_str + ) + + # Governance timestamp + gov_timestamp = self.nrc_client.get_latest_governance_reward_event() + if gov_timestamp: + add_line_helper( + "governance_latest_reward_event_timestamp_seconds", gov_timestamp + ) + + # Calculate days since last governance distribution + now_utc = datetime.now(timezone.utc) + gov_datetime: datetime = datetime.fromtimestamp(gov_timestamp, tz=timezone.utc) + days_since = int((now_utc - gov_datetime).total_seconds() / 86400) # 86400 seconds in a day + add_line_helper( + "days_since_governance_distribution", days_since + ) if not metrics_lines: raise ValueError("After evaluation there were no metrics to upload") @@ -458,13 +555,11 @@ def backfill(self, days: int = 40): logger.info("✅ Backfill complete!") def wait_until_next_run(self): - """Wait until 00:05 UTC tomorrow""" + """Wait until 00:10 UTC tomorrow""" now = datetime.now(timezone.utc) - # Calculate next 00:05 UTC - next_run = now.replace(hour=0, minute=5, second=0, microsecond=0) + next_run = now.replace(hour=0, minute=10, 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()