2 Commits
3 changed files with 37 additions and 79 deletions
+1 -5
View File
@@ -4,12 +4,11 @@ build-backend = "hatchling.build"
[project] [project]
name = "aiowebhooks" name = "aiowebhooks"
version = "0.1.1" version = "0.1.0"
description = "async webhook sender (aiohttp) with round-robin urls, retry, and proxy rotation; optional discord.py embeds" description = "async webhook sender (aiohttp) with round-robin urls, retry, and proxy rotation; optional discord.py embeds"
requires-python = ">=3.10" requires-python = ">=3.10"
dependencies = [ dependencies = [
"aiohttp>=3.9", "aiohttp>=3.9",
"commons @ git+ssh://git@git.rethinkstudios.io/rethink-public/commons.git@v0.2.0",
] ]
[project.optional-dependencies] [project.optional-dependencies]
@@ -17,8 +16,5 @@ discord = [
"discord.py>=2.3", "discord.py>=2.3",
] ]
[tool.hatch.metadata]
allow-direct-references = true
[tool.hatch.build.targets.wheel] [tool.hatch.build.targets.wheel]
packages = ["src/aiowebhooks"] packages = ["src/aiowebhooks"]
+1 -1
View File
@@ -21,4 +21,4 @@ from .sender import Webhook
__all__ = ["Webhook", "WebhookResult", "WebhookError", "NoUrlsError"] __all__ = ["Webhook", "WebhookResult", "WebhookError", "NoUrlsError"]
__version__ = "0.1.1" __version__ = "0.1.0"
+31 -69
View File
@@ -13,7 +13,6 @@ from typing import Dict, List, Optional, Union
from urllib.parse import unquote, urlsplit from urllib.parse import unquote, urlsplit
import aiohttp import aiohttp
from commons import aretry
from .errors import NoUrlsError from .errors import NoUrlsError
from .result import WebhookResult from .result import WebhookResult
@@ -21,18 +20,6 @@ from .result import WebhookResult
log = logging.getLogger(__name__) log = logging.getLogger(__name__)
class _Retryable(Exception):
"""internal signal: a retryable HTTP status (429/5xx); carries the response
raised inside an attempt so commons.aretry drives the backoff + cap; the loop
catches the final one to return the REAL last response, not a synthetic result.
"""
def __init__(self, result: WebhookResult):
super().__init__(f"retryable status {result.status}")
self.result = result
def _proxy_string(proxies_dict: Optional[Dict[str, str]]) -> Optional[str]: def _proxy_string(proxies_dict: Optional[Dict[str, str]]) -> Optional[str]:
"""canonical host:port:user:pass (or host:port) from an aiohttp proxies dict """canonical host:port:user:pass (or host:port) from an aiohttp proxies dict
@@ -131,57 +118,29 @@ class Webhook:
async def _send_loop( async def _send_loop(
self, session: aiohttp.ClientSession, url: str, payload: Dict self, session: aiohttp.ClientSession, url: str, payload: Dict
) -> WebhookResult: ) -> WebhookResult:
"""status-retry (via commons.aretry) wrapping proxy rotation; never raises """retry/rotation loop for a single send"""
attempts = 0
commons.aretry owns the 429/5xx backoff schedule + retry cap (max_retries), retries = 0
retrying on the internal _Retryable signal. on exhaustion it re-raises the
last _Retryable, whose carried result is the REAL last response (not a
synthetic status-0). proxy rotation on connection errors lives inside the
attempt and is capped separately.
"""
self._attempt_no = 0
try:
return await aretry(
lambda: self._attempt(session, url, payload),
attempts=self.max_retries + 1,
on=(_Retryable,),
)
except _Retryable as exhausted:
return exhausted.result
async def _attempt(
self, session: aiohttp.ClientSession, url: str, payload: Dict
) -> WebhookResult:
"""one logical send: proxy rotation + a single POST; may raise _Retryable
raises _Retryable (carrying the real response) on a 429/5xx so the caller's
aretry applies backoff; honors an explicit 429 retry_after by sleeping it
before signalling. returns a final WebhookResult on success or a terminal
(non-retryable) failure — never lets a provider/connection error escape.
"""
timeout = aiohttp.ClientTimeout(total=self.timeout)
last_proxy: Optional[str] = None
proxy_tries = 0 proxy_tries = 0
last_proxy: Optional[str] = None
timeout = aiohttp.ClientTimeout(total=self.timeout)
while True: while True:
self._attempt_no += 1
attempts = self._attempt_no
proxy_url = None proxy_url = None
if self._proxies is not None: if self._proxies is not None:
try: try:
proxy_dict = self._proxies.get() proxy_dict = self._proxies.get()
except Exception: except Exception as error:
# duck-typed provider; any error from get() means no proxy is if type(error).__name__ == "ProxiesExhaustedError":
# available — fail cleanly rather than escaping send().
log.warning("webhook: proxy get() failed; no proxy available",
exc_info=True)
return WebhookResult( return WebhookResult(
ok=False, status=None, url=url, attempts=attempts, ok=False, status=None, url=url, attempts=attempts,
error="proxies unavailable", proxy=last_proxy, error="proxies exhausted", proxy=last_proxy,
) )
raise
last_proxy = _proxy_string(proxy_dict) last_proxy = _proxy_string(proxy_dict)
proxy_url = (proxy_dict or {}).get("http") or (proxy_dict or {}).get("https") proxy_url = (proxy_dict or {}).get("http") or (proxy_dict or {}).get("https")
attempts += 1
try: try:
async with session.post( async with session.post(
url, json=payload, proxy=proxy_url, timeout=timeout url, json=payload, proxy=proxy_url, timeout=timeout
@@ -195,20 +154,25 @@ class Webhook:
response=body, proxy=last_proxy, response=body, proxy=last_proxy,
) )
result = WebhookResult( wait = self._retry_after(status, resp.headers, body)
if wait is not None and retries < self.max_retries:
retries += 1
log.warning("webhook 429 on %s; waiting %.3fs (retry %d/%d)",
url, wait, retries, self.max_retries)
await asyncio.sleep(wait)
continue
if status >= 500 and retries < self.max_retries:
retries += 1
log.warning("webhook %d on %s; retry %d/%d",
status, url, retries, self.max_retries)
continue
return WebhookResult(
ok=False, status=status, url=url, attempts=attempts, ok=False, status=status, url=url, attempts=attempts,
error=f"http {status}", response=body, proxy=last_proxy, error=f"http {status}", response=body, proxy=last_proxy,
) )
wait = self._retry_after(status, resp.headers, body)
if wait is not None:
log.warning("webhook 429 on %s; honoring retry_after %.3fs", url, wait)
await asyncio.sleep(wait)
raise _Retryable(result)
if status >= 500:
raise _Retryable(result)
return result
except (aiohttp.ClientError, asyncio.TimeoutError) as error: except (aiohttp.ClientError, asyncio.TimeoutError) as error:
if self._proxies is not None and proxy_tries < self.max_proxy_retries: if self._proxies is not None and proxy_tries < self.max_proxy_retries:
if self._burn(last_proxy): if self._burn(last_proxy):
@@ -224,20 +188,18 @@ class Webhook:
) )
def _burn(self, proxy: Optional[str]) -> bool: def _burn(self, proxy: Optional[str]) -> bool:
"""burn the current proxy; return False if it can't be rotated """burn the current proxy; return False if the provider is exhausted
the provider is duck-typed and never imported, so we cannot catch its catches a provider ProxiesExhaustedError duck-typed by class name (the
exception types by class. ANY exception from burn (a ProxiesExhaustedError provider is never imported), so a dying pool ends the loop cleanly.
on a dead pool, a ValueError when the proxy isn't in the pool, etc.) means
we can't rotate — return False so the caller ends the loop with a failed
result rather than letting it escape send() (which must never raise).
""" """
try: try:
self._proxies.burn(proxy) self._proxies.burn(proxy)
return True return True
except Exception: except Exception as error:
log.warning("webhook: proxy burn failed; ending rotation", exc_info=True) if type(error).__name__ == "ProxiesExhaustedError":
return False return False
raise
@staticmethod @staticmethod
async def _read_body(resp) -> Optional[Union[str, Dict]]: async def _read_body(resp) -> Optional[Union[str, Dict]]: