Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
50 changes: 29 additions & 21 deletions protocols/yearn/alert_large_flows.py
Original file line number Diff line number Diff line change
@@ -1,21 +1,20 @@
#!/usr/bin/env python3
import argparse
import json
import logging
import os
import sys
import time
import urllib.error
import urllib.request
from decimal import Decimal, getcontext

import requests
from dotenv import load_dotenv

from utils.abi import load_abi
from utils.alert import Alert, AlertSeverity, send_alert
from utils.cache import cache_filename, get_last_value_for_key_from_file, write_last_value_to_file
from utils.chains import EXPLORER_URLS, Chain
from utils.defillama import fetch_prices
from utils.http_client import request_with_retry
from utils.telegram import send_envio_error_message
from utils.web3_wrapper import ChainManager

Expand All @@ -24,6 +23,8 @@
getcontext().prec = 40

ENVIO_GRAPHQL_URL = os.getenv("ENVIO_GRAPHQL_URL")
# Per-attempt timeout; request_with_retry retries 5xx/timeouts so a brief Envio stall does not skip the run.
ENVIO_REQUEST_TIMEOUT_SECONDS = 10
DEFAULT_LOG_LEVEL = os.getenv("ALERT_LARGE_FLOWS_LOG_LEVEL", "WARNING")
IGNORED_FROM_ADDRESS = "0x283132390ea87d6ecc20255b59ba94329ee17961"
PROTOCOL = "yearn"
Expand Down Expand Up @@ -180,20 +181,19 @@
_logger = logging.getLogger("alert_large_flows")


def http_json(url: str, method: str = "GET", body: dict | None = None, headers: dict | None = None):
_logger.info("http_json %s %s", method, url)
data = None
req_headers = {"Accept": "application/json"}
if headers:
req_headers.update(headers)
if body is not None:
data = json.dumps(body).encode("utf-8")
req_headers["Content-Type"] = "application/json"
req = urllib.request.Request(url, data=data, headers=req_headers, method=method)
with urllib.request.urlopen(req, timeout=30) as resp:
payload = json.loads(resp.read().decode("utf-8"))
_logger.info("http_json status=%s", resp.status)
return payload
def http_json(url: str, body: dict) -> dict:
"""POST a JSON body with retries on transient errors and return the decoded response."""
_logger.info("http_json POST %s", url)
response = request_with_retry(
"post",
url,
timeout=ENVIO_REQUEST_TIMEOUT_SECONDS,
json=body,
headers={"Accept": "application/json"},
)
_logger.info("http_json status=%s", response.status_code)
payload: dict = response.json()
return payload


def gql_request(query: str, variables: dict) -> dict | None:
Expand All @@ -206,13 +206,21 @@ def gql_request(query: str, variables: dict) -> dict | None:
payload = {"query": query, "variables": variables}

try:
return http_json(ENVIO_GRAPHQL_URL, method="POST", body=payload)
except urllib.error.HTTPError as exc:
return http_json(ENVIO_GRAPHQL_URL, body=payload)
except requests.HTTPError as exc:
status = exc.response.status_code if exc.response is not None else "unknown"
send_envio_error_message(
f"⚠️ Large Flow Alert: Envio GraphQL error (HTTP {status}). Skipping this run.",
PROTOCOL,
)
_logger.error("Envio request failed with HTTP %s", status)
return None
except (requests.Timeout, requests.ConnectionError) as exc:
send_envio_error_message(
f"⚠️ Large Flow Alert: Envio GraphQL error (HTTP {exc.code}). Skipping this run.",
f"⚠️ Large Flow Alert: Envio GraphQL request failed ({type(exc).__name__}). Skipping this run.",
PROTOCOL,
)
_logger.error("Envio request failed with HTTP %d", exc.code)
_logger.error("Envio request failed: %s", exc)
return None


Expand Down
40 changes: 27 additions & 13 deletions protocols/yearn/alert_small_parent_flows.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,18 +8,18 @@
import logging
import os
import time
import urllib.error
import urllib.request
from dataclasses import dataclass
from decimal import Decimal, getcontext
from typing import Callable

import requests
from dotenv import load_dotenv

from protocols.yearn.kong import fetch_kong_parent_vaults
from utils import store
from utils.alert import Alert, AlertSeverity, send_alert
from utils.chains import EXPLORER_URLS, Chain
from utils.http_client import request_with_retry
from utils.logger import get_logger
from utils.telegram import MAX_MESSAGE_LENGTH, YEARN_MAINTENANCE_CHANNEL, resolve_channel, send_envio_error_message

Expand All @@ -28,6 +28,8 @@
getcontext().prec = 60

ENVIO_GRAPHQL_URL = os.getenv("ENVIO_GRAPHQL_URL")
# Per-attempt timeout; request_with_retry retries 5xx/timeouts so a brief Envio stall does not skip the run.
ENVIO_REQUEST_TIMEOUT_SECONDS = 10
DEFAULT_LOG_LEVEL = os.getenv("SMALL_PARENT_FLOWS_LOG_LEVEL") or os.getenv("LOG_LEVEL", "INFO")
DEFAULT_THRESHOLD_RAW = 10_000
DEFAULT_LOOKBACK_SECONDS = 7200
Expand Down Expand Up @@ -249,20 +251,31 @@ class EventCursor:


def http_json(url: str, body: dict) -> dict:
"""POST a JSON body and return the decoded response."""
request = urllib.request.Request(
"""POST a JSON body with retries on transient errors and return the decoded response."""
response = request_with_retry(
"post",
url,
data=json.dumps(body).encode("utf-8"),
headers={"Accept": "application/json", "Content-Type": "application/json"},
method="POST",
timeout=ENVIO_REQUEST_TIMEOUT_SECONDS,
json=body,
headers={"Accept": "application/json"},
)
with urllib.request.urlopen(request, timeout=30) as response:
payload: object = json.loads(response.read().decode("utf-8"))
payload: object = response.json()
if not isinstance(payload, dict):
raise ValueError("Envio returned a non-object JSON response")
return payload


def describe_request_error(exc: Exception) -> str:
"""Return a short failure description that does not leak the Envio endpoint URL."""
if isinstance(exc, requests.HTTPError) and exc.response is not None:
return f"HTTP {exc.response.status_code} {exc.response.reason}".strip()
if isinstance(exc, requests.Timeout):
return f"timed out after {ENVIO_REQUEST_TIMEOUT_SECONDS}s"
if isinstance(exc, requests.ConnectionError):
return "connection error"
return str(exc)


def gql_request(query: str, variables: dict) -> dict:
"""Execute an Envio GraphQL query, routing failures to its ops channel.

Expand All @@ -274,15 +287,16 @@ def gql_request(query: str, variables: dict) -> dict:

try:
payload = http_json(ENVIO_GRAPHQL_URL, {"query": query, "variables": variables})
except (urllib.error.HTTPError, urllib.error.URLError, ConnectionError, OSError, ValueError) as exc:
except (requests.RequestException, OSError, ValueError) as exc:
detail = describe_request_error(exc)
send_envio_error_message(
f"Small parent flow monitor: Envio GraphQL request failed ({exc}). Skipping this run.",
f"Small parent flow monitor: Envio GraphQL request failed ({detail}). Skipping this run.",
PROTOCOL,
source="small_parent_flows",
alert_protocol=ALERT_PROTOCOL,
)
logger.error("Envio request failed: %s", exc)
raise EnvioUnavailableError(f"Envio request failed: {exc}") from exc
logger.error("Envio request failed: %s", detail)
raise EnvioUnavailableError(f"Envio request failed: {detail}") from exc

if payload.get("errors"):
send_envio_error_message(
Expand Down
42 changes: 42 additions & 0 deletions tests/test_small_parent_flows.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
import json
import sys
from decimal import Decimal

import pytest
import requests

from protocols.yearn import alert_small_parent_flows as monitor
from utils.alert import Alert, AlertSeverity
Expand Down Expand Up @@ -227,6 +229,46 @@ def failing_http(_url, _body):
assert kwargs["alert_protocol"] == "yearn-internal"


def _response(status: int, payload: dict | None = None) -> requests.Response:
response = requests.Response()
response.status_code = status
response.reason = "Gateway Timeout" if status == 504 else "OK"
response._content = b"" if payload is None else json.dumps(payload).encode()
response.url = "https://envio.example/graphql"
return response


def test_http_json_retries_gateway_timeout_with_short_timeout(monkeypatch) -> None:
calls = []
responses = [_response(504), _response(200, {"data": {"events": []}})]

def fake_request(method, url, timeout=None, **kwargs):
calls.append((method, timeout, kwargs["json"]))
return responses.pop(0)

monkeypatch.setattr("utils.http_client.requests.request", fake_request)
monkeypatch.setattr("utils.http_client.time.sleep", lambda _seconds: None)

assert monitor.http_json("https://envio.example/graphql", {"query": "q"}) == {"data": {"events": []}}
assert [call[1] for call in calls] == [10, 10]
assert calls[0][0] == "post"


def test_gql_request_reports_status_without_url_after_retries(monkeypatch) -> None:
reported = []

monkeypatch.setattr("utils.http_client.requests.request", lambda *args, **kwargs: _response(504))
monkeypatch.setattr("utils.http_client.time.sleep", lambda _seconds: None)
monkeypatch.setattr(monitor, "ENVIO_GRAPHQL_URL", "https://envio.example/graphql")
monkeypatch.setattr(monitor, "send_envio_error_message", lambda *args, **kwargs: reported.append(args[0]))

with pytest.raises(monitor.EnvioUnavailableError, match="HTTP 504 Gateway Timeout"):
monitor.gql_request("query {}", {})
assert len(reported) == 1
assert "HTTP 504 Gateway Timeout" in reported[0]
assert "envio.example" not in reported[0]


def test_first_run_lookback_floor_persists_without_events(monkeypatch) -> None:
since_values = []

Expand Down
Loading