diff --git a/protocols/compound/proposals.py b/protocols/compound/proposals.py index 76883f92..4e751ff4 100644 --- a/protocols/compound/proposals.py +++ b/protocols/compound/proposals.py @@ -5,6 +5,7 @@ from utils.alert import Alert, AlertSeverity, send_alert from utils.cache import get_last_queued_id_from_file, write_last_queued_id_to_file +from utils.http_client import request_with_retry from utils.logger import get_logger from utils.telegram import escape_markdown, send_error_message @@ -86,13 +87,13 @@ def get_proposals(): } try: - response = requests.post( + response = request_with_retry( + "post", TALLY_API_URL, json={"query": query, "variables": variables}, headers=headers, timeout=30, ) - response.raise_for_status() data = response.json() if "errors" in data: diff --git a/protocols/ethena/ethena.py b/protocols/ethena/ethena.py index c08651a8..bbb2d378 100644 --- a/protocols/ethena/ethena.py +++ b/protocols/ethena/ethena.py @@ -1,9 +1,8 @@ from datetime import datetime, timedelta, timezone -import requests - from utils.abi import load_abi from utils.alert import Alert, AlertSeverity, send_alert +from utils.http_client import fetch_json from utils.logger import get_logger from utils.telegram import send_error_message from utils.web3_wrapper import Chain, ChainManager @@ -29,19 +28,6 @@ ETHENA_SOURCE = "Ethena API" -def fetch_json(url: str) -> dict | None: - """Helper that fetches JSON with basic error handling.""" - try: - resp = requests.get(url, timeout=REQUEST_TIMEOUT) - if resp.status_code != 200: - logger.error("HTTP %s for %s\n%s", resp.status_code, url, resp.text) - return None - return resp.json() - except Exception as e: - logger.error("Failed to fetch %s: %s", url, e) - return None - - def _parse_timestamp(ts: str) -> datetime | None: """Parse the timestamp formats returned by Ethena's transparency API.""" formats = [ @@ -75,7 +61,7 @@ def is_stale_timestamp(ts: str, max_age_hours: int = 3) -> bool: def get_usde_supply() -> float | None: """Return total circulating USDe supply in USD terms (raw token amount / 1e18).""" - data = fetch_json(SUPPLY_URL) + data = fetch_json(SUPPLY_URL, timeout=REQUEST_TIMEOUT) if not data: return None @@ -98,7 +84,7 @@ def get_total_collateral_usd() -> float | None: e.g. $4.14B vs $4.03B) but is a stale snapshot (items lag several hours). We use the fresh net figure and add the reserve fund as the buffer. """ - data = fetch_json(COLLATERAL_URL) + data = fetch_json(COLLATERAL_URL, timeout=REQUEST_TIMEOUT) if not data: return None @@ -112,7 +98,7 @@ def get_reserve_fund() -> float | None: ``{timestamp, value}`` points; we take the most recent one and treat stale data (older than 3 hours) as unavailable. """ - data = fetch_json(RESERVE_FUND_URL) + data = fetch_json(RESERVE_FUND_URL, timeout=REQUEST_TIMEOUT) if not data: return None diff --git a/protocols/fluid/proposals.py b/protocols/fluid/proposals.py index 24382334..614327d3 100644 --- a/protocols/fluid/proposals.py +++ b/protocols/fluid/proposals.py @@ -4,6 +4,7 @@ from utils.alert import Alert, AlertSeverity, send_alert from utils.cache import get_last_queued_id_from_file, write_last_queued_id_to_file +from utils.http_client import request_with_retry from utils.logger import get_logger from utils.telegram import escape_markdown, send_error_message @@ -62,8 +63,7 @@ def get_proposals(): """Fetch and process Fluid governance proposals""" try: # Fetch queued proposals - response = requests.get(f"{FLUID_API_URL}?status=queued", timeout=30) - response.raise_for_status() + response = request_with_retry("get", f"{FLUID_API_URL}?status=queued", timeout=30) data = response.json() if "data" not in data or not data["data"]: diff --git a/protocols/infinifi/main.py b/protocols/infinifi/main.py index aede853d..44d31f94 100644 --- a/protocols/infinifi/main.py +++ b/protocols/infinifi/main.py @@ -1,6 +1,5 @@ from decimal import Decimal -import requests from web3 import Web3 from utils.abi import load_abi @@ -15,6 +14,7 @@ ) from utils.chains import Chain from utils.config import Config +from utils.http_client import fetch_json from utils.logger import get_logger from utils.telegram import send_error_message from utils.web3_wrapper import ChainManager @@ -57,20 +57,10 @@ def fetch_api_data(): """Fetches data from the Infinifi API.""" - url = f"{API_BASE_URL}{API_PROTOCOL_DATA}" - try: - headers = { - "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/91.0.4472.124 Safari/537.36" - } - response = requests.get(url, headers=headers, timeout=10) - if response.status_code == 200: - return response.json() - else: - logger.error("API Error %s for %s: %s", response.status_code, url, response.text[:200]) - return None - except Exception as e: - logger.error("API Request Failed for %s: %s", url, e) - return None + headers = { + "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/91.0.4472.124 Safari/537.36" + } + return fetch_json(f"{API_BASE_URL}{API_PROTOCOL_DATA}", headers=headers, timeout=10) def to_float(value, default=0.0): diff --git a/protocols/maker/proposals.py b/protocols/maker/proposals.py index 5ec998ec..9c360538 100644 --- a/protocols/maker/proposals.py +++ b/protocols/maker/proposals.py @@ -4,6 +4,7 @@ from utils.alert import Alert, AlertSeverity, send_alert from utils.cache import get_last_queued_id_from_file, write_last_queued_id_to_file +from utils.http_client import request_with_retry from utils.logger import get_logger from utils.telegram import escape_markdown, send_error_message @@ -16,8 +17,7 @@ def fetch_executive_proposals() -> list[dict]: """Fetch executive proposals from the Sky governance API.""" - response = requests.get(SKY_EXECUTIVE_API, timeout=30) - response.raise_for_status() + response = request_with_retry("get", SKY_EXECUTIVE_API, timeout=30) return response.json() diff --git a/protocols/spark/proposals.py b/protocols/spark/proposals.py index 9391b69b..3c39924b 100644 --- a/protocols/spark/proposals.py +++ b/protocols/spark/proposals.py @@ -10,6 +10,7 @@ from utils.alert import Alert, AlertSeverity, send_alert from utils.cache import get_last_queued_id_from_file, write_last_queued_id_to_file +from utils.http_client import request_with_retry from utils.logger import get_logger from utils.telegram import escape_markdown, send_error_message @@ -42,12 +43,12 @@ def fetch_proposals() -> list[dict]: """Fetch the most recent proposals from the Snapshot hub GraphQL API.""" - response = requests.post( + response = request_with_retry( + "post", SNAPSHOT_GRAPHQL_URL, json={"query": PROPOSALS_QUERY, "variables": {"space": SNAPSHOT_SPACE}}, timeout=30, ) - response.raise_for_status() payload = response.json() if payload.get("errors"): raise ValueError(f"Snapshot GraphQL errors: {payload['errors']}") diff --git a/protocols/yearn/kong.py b/protocols/yearn/kong.py index c71833f8..f2db8a0b 100644 --- a/protocols/yearn/kong.py +++ b/protocols/yearn/kong.py @@ -5,6 +5,7 @@ import requests from utils.chains import Chain +from utils.http_client import request_with_retry from utils.logger import get_logger logger = get_logger("yearn.kong") @@ -55,12 +56,15 @@ class KongRequestError(requests.RequestException): def _post_graphql(query: str, variables: Dict[str, object]) -> Dict[str, Any]: """Execute a Kong GraphQL query and return the response data.""" - response = requests.post( - KONG_GQL_URL, - json={"query": query, "variables": variables}, - timeout=30, - ) - response.raise_for_status() + try: + response = request_with_retry( + "post", + KONG_GQL_URL, + json={"query": query, "variables": variables}, + timeout=30, + ) + except requests.RequestException as e: + raise KongRequestError(f"Kong request failed: {e}") from e try: payload = response.json() diff --git a/tests/test_compound_proposals.py b/tests/test_compound_proposals.py index c5fcca37..63f2eb45 100644 --- a/tests/test_compound_proposals.py +++ b/tests/test_compound_proposals.py @@ -38,14 +38,14 @@ def test_compound_alert_uses_timeout_escapes_title_and_updates_reported_id(): } with ( - patch("protocols.compound.proposals.requests.post", return_value=_Response(payload)) as mock_post, + patch("protocols.compound.proposals.request_with_retry", return_value=_Response(payload)) as mock_request, patch("protocols.compound.proposals.get_last_queued_id_from_file", return_value=40), patch("protocols.compound.proposals.send_alert") as mock_send, patch("protocols.compound.proposals.write_last_queued_id_to_file") as mock_write, ): get_proposals() - assert mock_post.call_args.kwargs["timeout"] == 30 + assert mock_request.call_args.kwargs["timeout"] == 30 mock_send.assert_called_once() alert = mock_send.call_args.args[0] assert alert.protocol == "comp" @@ -57,7 +57,7 @@ def test_compound_alert_uses_timeout_escapes_title_and_updates_reported_id(): def test_compound_processing_error_alert_uses_plain_text(): with ( - patch("protocols.compound.proposals.requests.post", return_value=_Response({"data": {}})), + patch("protocols.compound.proposals.request_with_retry", return_value=_Response({"data": {}})), patch("protocols.compound.proposals.send_error_message") as mock_send, ): get_proposals() @@ -71,7 +71,7 @@ def test_compound_processing_error_alert_uses_plain_text(): def test_compound_fetch_error_alert_uses_plain_text(): with ( - patch("protocols.compound.proposals.requests.post", side_effect=requests.Timeout("timed out")), + patch("protocols.compound.proposals.request_with_retry", side_effect=requests.RequestException("timed out")), patch("protocols.compound.proposals.send_error_message") as mock_send, ): get_proposals() diff --git a/tests/test_fluid_proposals.py b/tests/test_fluid_proposals.py index 4d3e2c05..2d1f76d6 100644 --- a/tests/test_fluid_proposals.py +++ b/tests/test_fluid_proposals.py @@ -40,7 +40,7 @@ def test_fluid_proposal_alert_escapes_api_markdown_and_keeps_link(): } with ( - patch("protocols.fluid.proposals.requests.get", return_value=_Response(payload)), + patch("protocols.fluid.proposals.request_with_retry", return_value=_Response(payload)), patch("protocols.fluid.proposals.get_last_queued_id_from_file", return_value=130), patch("protocols.fluid.proposals.write_last_queued_id_to_file") as mock_write, patch("protocols.fluid.proposals.send_alert") as mock_send, @@ -59,7 +59,7 @@ def test_fluid_proposal_alert_escapes_api_markdown_and_keeps_link(): def test_fluid_proposal_fetch_error_routes_to_errors_channel(): with ( - patch("protocols.fluid.proposals.requests.get", side_effect=Exception("bad TYPE_1 payload")), + patch("protocols.fluid.proposals.request_with_retry", side_effect=Exception("bad TYPE_1 payload")), patch("protocols.fluid.proposals.send_error_message") as mock_send, ): get_proposals() diff --git a/tests/test_yearn_kong.py b/tests/test_yearn_kong.py index 66dd9e22..76501f5b 100644 --- a/tests/test_yearn_kong.py +++ b/tests/test_yearn_kong.py @@ -43,17 +43,18 @@ def _vault_payload() -> dict: def test_fetch_kong_vaults_uses_all_strategies_by_default(monkeypatch) -> None: calls = [] - def fake_post(url: str, json: dict, timeout: int) -> FakeResponse: - calls.append((url, json, timeout)) + def fake_request_with_retry(method: str, url: str, *, json: dict, timeout: int) -> FakeResponse: + calls.append((method, url, json, timeout)) return FakeResponse(_vault_payload()) - monkeypatch.setattr(kong.requests, "post", fake_post) + monkeypatch.setattr(kong, "request_with_retry", fake_request_with_retry) vaults = kong.fetch_kong_vaults(Chain.MAINNET) - assert calls[0][0] == kong.KONG_GQL_URL - assert calls[0][1]["variables"] == {"chainId": Chain.MAINNET.chain_id} - assert calls[0][2] == 30 + assert calls[0][0] == "post" + assert calls[0][1] == kong.KONG_GQL_URL + assert calls[0][2]["variables"] == {"chainId": Chain.MAINNET.chain_id} + assert calls[0][3] == 30 assert vaults == [ { "address": "0xabc", @@ -67,8 +68,8 @@ def fake_post(url: str, json: dict, timeout: int) -> FakeResponse: def test_fetch_kong_vaults_can_use_default_queue(monkeypatch) -> None: monkeypatch.setattr( - kong.requests, - "post", + kong, + "request_with_retry", lambda *_args, **_kwargs: FakeResponse(_vault_payload()), ) @@ -83,7 +84,7 @@ def test_fetch_kong_vaults_can_use_default_queue(monkeypatch) -> None: def test_fetch_kong_vaults_raises_on_graphql_errors(monkeypatch) -> None: payload = {"errors": [{"message": "bad query"}]} - monkeypatch.setattr(kong.requests, "post", lambda *_args, **_kwargs: FakeResponse(payload)) + monkeypatch.setattr(kong, "request_with_retry", lambda *_args, **_kwargs: FakeResponse(payload)) with pytest.raises(kong.KongRequestError): kong.fetch_kong_vaults(Chain.MAINNET) @@ -143,16 +144,16 @@ def test_fetch_kong_parent_vaults_filters_retired_and_malformed_vaults(monkeypat } calls = [] - def fake_post(url: str, json: dict, timeout: int) -> FakeResponse: - calls.append((url, json, timeout)) + def fake_request_with_retry(method: str, url: str, *, json: dict, timeout: int) -> FakeResponse: + calls.append((method, url, json, timeout)) return FakeResponse(payload) - monkeypatch.setattr(kong.requests, "post", fake_post) + monkeypatch.setattr(kong, "request_with_retry", fake_request_with_retry) vaults = kong.fetch_kong_parent_vaults(Chain.MAINNET) - assert "vaultType: 1" in calls[0][1]["query"] - assert calls[0][1]["variables"] == {"chainId": 1} + assert "vaultType: 1" in calls[0][2]["query"] + assert calls[0][2]["variables"] == {"chainId": 1} assert vaults == [ { "address": "0xParent", diff --git a/utils/http_client.py b/utils/http_client.py index b4637a98..bb05ce3f 100644 --- a/utils/http_client.py +++ b/utils/http_client.py @@ -54,6 +54,7 @@ def request_with_retry( except requests.exceptions.HTTPError as e: status_code = e.response.status_code if e.response is not None else None if status_code is not None and status_code < 500 and status_code != 429: + logger.error("HTTP %s for %s: %s", status_code, url, e.response.text[:200]) raise # Do not retry permanent client errors. last_exception = e except (requests.exceptions.ConnectionError, requests.exceptions.Timeout) as e: @@ -81,17 +82,12 @@ def fetch_json( timeout: int | None = None, **kwargs: Any, ) -> dict | None: - """Fetch JSON from a URL with error handling. + """Fetch JSON from a URL with retry and error handling. Returns the parsed JSON dict on success, or None on failure. """ - if timeout is None: - timeout = Config.get_request_timeout() try: - resp = requests.request(method, url, timeout=timeout, **kwargs) - if resp.status_code != 200: - logger.error("HTTP %s for %s: %s", resp.status_code, url, resp.text[:200]) - return None + resp = request_with_retry(method, url, timeout=timeout, **kwargs) return resp.json() except Exception as e: logger.error("Request failed for %s: %s", url, e) diff --git a/utils/tenderly/tenderly.py b/utils/tenderly/tenderly.py index b97d2d78..d1be1afa 100644 --- a/utils/tenderly/tenderly.py +++ b/utils/tenderly/tenderly.py @@ -10,9 +10,9 @@ import os from pathlib import Path -import requests from dotenv import load_dotenv +from utils.http_client import request_with_retry from utils.logger import get_logger load_dotenv() @@ -38,9 +38,7 @@ def get_response_hash(data: dict) -> str: def fetch_alerts() -> dict: """Fetch alerts from Tenderly API.""" headers = {"Accept": "application/json", "X-Access-Key": TENDERLY_API_KEY} - response = requests.get(TENDERLY_API_URL, headers=headers) - if response.status_code != 200: - raise Exception(f"Failed to get alerts: {response.status_code} - {response.text}") + response = request_with_retry("get", TENDERLY_API_URL, headers=headers) return response.json()