Enhance PIA rate limit handling in gluetun WireGuard rotator. Introduce persistent cooldown management and improve retry logic for token requests. Update README to reflect new RATE_LIMIT_PATH variable and clarify cooldown behavior.
Build and Push Docker Images / build-and-push (push) Successful in 20s
Build and Push Docker Images / build-and-push (push) Successful in 20s
This commit is contained in:
@@ -50,7 +50,8 @@ docker run --rm --entrypoint /opt/venv/bin/python \
|
|||||||
| `TOKEN_CACHE_TTL` | `20h` | Hergebruik token i.p.v. opnieuw inloggen |
|
| `TOKEN_CACHE_TTL` | `20h` | Hergebruik token i.p.v. opnieuw inloggen |
|
||||||
| `PIA_CA_PATH` | `/config/cache/ca.rsa.4096.crt` | Gecachete PIA CA voor `addKey` TLS |
|
| `PIA_CA_PATH` | `/config/cache/ca.rsa.4096.crt` | Gecachete PIA CA voor `addKey` TLS |
|
||||||
| `FORCE_TOKEN_REFRESH` | `false` | `true` = token-cache negeren |
|
| `FORCE_TOKEN_REFRESH` | `false` | `true` = token-cache negeren |
|
||||||
| `RATE_LIMIT_WAIT_SECONDS` | `3600` | Wachttijd bij PIA rate-limit vóór retry |
|
| `RATE_LIMIT_WAIT_SECONDS` | `3600` | Cooldown bij PIA rate-limit (blijft gelden na container-restart) |
|
||||||
|
| `RATE_LIMIT_PATH` | `/config/cache/pia-rate-limit.json` | Persistente cooldown-timestamp |
|
||||||
| `TZ` | `Europe/Brussels` | Tijdzone voor scheduling |
|
| `TZ` | `Europe/Brussels` | Tijdzone voor scheduling |
|
||||||
|
|
||||||
`ROTATE_CRON` voorbeelden: `0 */6 * * *`, `0 3 * * 1-5`, `@hourly`. Quote in Compose: `'ROTATE_CRON=0 3 * * *'`.
|
`ROTATE_CRON` voorbeelden: `0 */6 * * *`, `0 3 * * 1-5`, `@hourly`. Quote in Compose: `'ROTATE_CRON=0 3 * * *'`.
|
||||||
@@ -113,6 +114,6 @@ Behoud minimaal:
|
|||||||
2. **Pin** de snelste server-IP (niet alleen regio)
|
2. **Pin** de snelste server-IP (niet alleen regio)
|
||||||
3. Altijd nieuw keypair + `addKey` (reauth), daarna containers herstarten
|
3. Altijd nieuw keypair + `addKey` (reauth), daarna containers herstarten
|
||||||
4. Caches: serverlist, token (~20u), CA-cert — alleen om overbodige API-calls te beperken
|
4. Caches: serverlist, token (~20u), CA-cert — alleen om overbodige API-calls te beperken
|
||||||
5. Bij rate-limit: `RATE_LIMIT_WAIT_SECONDS` (default 1 uur) wachten en opnieuw proberen
|
5. Bij rate-limit: cooldown van `RATE_LIMIT_WAIT_SECONDS` (default 1 uur), **persistent op disk** zodat restarts geen extra API-calls doen; daarna opnieuw proberen. Andere token-endpoints worden eerst nog geprobeerd.
|
||||||
|
|
||||||
**Let op:** zet gluetun **niet** in `RESTART_CONTAINERS`; gebruik `GLUETUN_CONTAINER`. Elke echte rotatie geeft korte downtime.
|
**Let op:** zet gluetun **niet** in `RESTART_CONTAINERS`; gebruik `GLUETUN_CONTAINER`. Elke echte rotatie geeft korte downtime.
|
||||||
|
|||||||
@@ -13,8 +13,9 @@ from zoneinfo import ZoneInfo
|
|||||||
|
|
||||||
from croniter import croniter
|
from croniter import croniter
|
||||||
|
|
||||||
# Matches rotate.py / pia.py EXIT_RATE_LIMITED (EX_TEMPFAIL)
|
import pia
|
||||||
EXIT_RATE_LIMITED = 75
|
|
||||||
|
EXIT_RATE_LIMITED = pia.EXIT_RATE_LIMITED
|
||||||
|
|
||||||
CRON_MACROS = {
|
CRON_MACROS = {
|
||||||
"@yearly": "0 0 1 1 *",
|
"@yearly": "0 0 1 1 *",
|
||||||
@@ -79,16 +80,25 @@ def sleep_until_next_rotate(expr: str) -> None:
|
|||||||
time.sleep(wait_s)
|
time.sleep(wait_s)
|
||||||
|
|
||||||
|
|
||||||
|
def wait_for_rate_limit_cooldown() -> None:
|
||||||
|
remaining = pia.rate_limit_remaining_seconds()
|
||||||
|
if remaining <= 0:
|
||||||
|
return
|
||||||
|
log(f"PIA rate-limit cooldown active; waiting {remaining}s before retry")
|
||||||
|
time.sleep(remaining)
|
||||||
|
|
||||||
|
|
||||||
def run_rotation(reason: str) -> None:
|
def run_rotation(reason: str) -> None:
|
||||||
wait_s = int(os.environ.get("RATE_LIMIT_WAIT_SECONDS", "3600"))
|
|
||||||
while True:
|
while True:
|
||||||
|
wait_for_rate_limit_cooldown()
|
||||||
log(reason)
|
log(reason)
|
||||||
result = subprocess.run(["/usr/local/bin/rotate.py"], check=False)
|
result = subprocess.run(["/usr/local/bin/rotate.py"], check=False)
|
||||||
if result.returncode == 0:
|
if result.returncode == 0:
|
||||||
return
|
return
|
||||||
if result.returncode == EXIT_RATE_LIMITED:
|
if result.returncode == EXIT_RATE_LIMITED:
|
||||||
log(f"PIA rate-limited (too many attempts); waiting {wait_s}s before retry")
|
# rotate.py / pia.py already marked the cooldown file.
|
||||||
time.sleep(wait_s)
|
if pia.rate_limit_remaining_seconds() <= 0:
|
||||||
|
pia.mark_rate_limited(reason="rotation exit code 75")
|
||||||
reason = "Retrying rotation after rate-limit wait"
|
reason = "Retrying rotation after rate-limit wait"
|
||||||
continue
|
continue
|
||||||
raise SystemExit(f"Rotation failed with exit code {result.returncode}")
|
raise SystemExit(f"Rotation failed with exit code {result.returncode}")
|
||||||
|
|||||||
@@ -226,6 +226,79 @@ def ensure_pia_ca() -> Path:
|
|||||||
return path
|
return path
|
||||||
|
|
||||||
|
|
||||||
|
def pia_ssl_context() -> ssl.SSLContext:
|
||||||
|
"""SSL context trusting PIA's CA.
|
||||||
|
|
||||||
|
OpenSSL 3.2+ / Python 3.13+ enable X509_STRICT by default, which rejects
|
||||||
|
PIA's ca.rsa.4096.crt because basicConstraints is not marked critical.
|
||||||
|
"""
|
||||||
|
context = ssl.create_default_context(cafile=str(ensure_pia_ca()))
|
||||||
|
if hasattr(ssl, "VERIFY_X509_STRICT"):
|
||||||
|
context.verify_flags &= ~ssl.VERIFY_X509_STRICT
|
||||||
|
return context
|
||||||
|
|
||||||
|
|
||||||
|
def rate_limit_path() -> Path:
|
||||||
|
return Path(os.environ.get("RATE_LIMIT_PATH", "/config/cache/pia-rate-limit.json"))
|
||||||
|
|
||||||
|
|
||||||
|
def rate_limit_wait_seconds() -> int:
|
||||||
|
return int(os.environ.get("RATE_LIMIT_WAIT_SECONDS", "3600"))
|
||||||
|
|
||||||
|
|
||||||
|
def rate_limit_remaining_seconds() -> int:
|
||||||
|
path = rate_limit_path()
|
||||||
|
if not path.is_file():
|
||||||
|
return 0
|
||||||
|
try:
|
||||||
|
data = json.loads(path.read_text(encoding="utf-8"))
|
||||||
|
until = float(data.get("until", 0))
|
||||||
|
except (OSError, json.JSONDecodeError, TypeError, ValueError):
|
||||||
|
return 0
|
||||||
|
return max(0, int(until - time.time()))
|
||||||
|
|
||||||
|
|
||||||
|
def mark_rate_limited(wait_s: int | None = None, reason: str = "") -> int:
|
||||||
|
wait = rate_limit_wait_seconds() if wait_s is None else wait_s
|
||||||
|
until = time.time() + wait
|
||||||
|
path = rate_limit_path()
|
||||||
|
path.parent.mkdir(parents=True, exist_ok=True)
|
||||||
|
path.write_text(
|
||||||
|
json.dumps(
|
||||||
|
{
|
||||||
|
"until": until,
|
||||||
|
"wait_seconds": wait,
|
||||||
|
"reason": reason,
|
||||||
|
"marked_at": datetime.now().astimezone().isoformat(timespec="seconds"),
|
||||||
|
},
|
||||||
|
indent=2,
|
||||||
|
)
|
||||||
|
+ "\n",
|
||||||
|
encoding="utf-8",
|
||||||
|
)
|
||||||
|
log(f"Marked PIA rate-limit cooldown for {wait}s ({reason or 'rate limited'})")
|
||||||
|
return wait
|
||||||
|
|
||||||
|
|
||||||
|
def clear_rate_limit() -> None:
|
||||||
|
path = rate_limit_path()
|
||||||
|
try:
|
||||||
|
path.unlink(missing_ok=True)
|
||||||
|
except OSError:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
def raise_if_rate_limit_cooldown() -> None:
|
||||||
|
remaining = rate_limit_remaining_seconds()
|
||||||
|
if remaining > 0:
|
||||||
|
log(f"PIA rate-limit cooldown active; {remaining}s remaining (no API calls)")
|
||||||
|
raise SystemExit(EXIT_RATE_LIMITED)
|
||||||
|
|
||||||
|
|
||||||
|
class RateLimited(RuntimeError):
|
||||||
|
"""One auth endpoint reported rate limiting; other methods may still work."""
|
||||||
|
|
||||||
|
|
||||||
def _rate_limited(body: str, status_code: int) -> bool:
|
def _rate_limited(body: str, status_code: int) -> bool:
|
||||||
return status_code == 429 or "too_many_attempts" in body
|
return status_code == 429 or "too_many_attempts" in body
|
||||||
|
|
||||||
@@ -279,15 +352,15 @@ def _token_via_central(username: str, password: str) -> str:
|
|||||||
except urllib.error.HTTPError as exc:
|
except urllib.error.HTTPError as exc:
|
||||||
body = exc.read().decode("utf-8", errors="replace")
|
body = exc.read().decode("utf-8", errors="replace")
|
||||||
if _rate_limited(body, exc.code):
|
if _rate_limited(body, exc.code):
|
||||||
raise SystemExit(EXIT_RATE_LIMITED) from exc
|
raise RateLimited(f"central token API: {body[:200]}") from exc
|
||||||
if _cloudflare_blocked(exc.code, body):
|
if _cloudflare_blocked(exc.code, body):
|
||||||
raise RuntimeError(f"central token API blocked by Cloudflare (1010)") from exc
|
raise RuntimeError("central token API blocked by Cloudflare (1010)") from exc
|
||||||
raise RuntimeError(f"central token API status {exc.code}: {body[:200]}") from exc
|
raise RuntimeError(f"central token API status {exc.code}: {body[:200]}") from exc
|
||||||
except urllib.error.URLError as exc:
|
except urllib.error.URLError as exc:
|
||||||
raise RuntimeError(f"central token API network error: {exc}") from exc
|
raise RuntimeError(f"central token API network error: {exc}") from exc
|
||||||
|
|
||||||
if _rate_limited(body, status):
|
if _rate_limited(body, status):
|
||||||
raise SystemExit(EXIT_RATE_LIMITED)
|
raise RateLimited(f"central token API: {body[:200]}")
|
||||||
return _parse_token_body(body, "central token API")
|
return _parse_token_body(body, "central token API")
|
||||||
|
|
||||||
|
|
||||||
@@ -307,7 +380,7 @@ def _token_via_gtoken(username: str, password: str) -> str:
|
|||||||
except urllib.error.HTTPError as exc:
|
except urllib.error.HTTPError as exc:
|
||||||
body = exc.read().decode("utf-8", errors="replace")
|
body = exc.read().decode("utf-8", errors="replace")
|
||||||
if _rate_limited(body, exc.code):
|
if _rate_limited(body, exc.code):
|
||||||
raise SystemExit(EXIT_RATE_LIMITED) from exc
|
raise RateLimited(f"gtoken API: {body[:200]}") from exc
|
||||||
if _cloudflare_blocked(exc.code, body):
|
if _cloudflare_blocked(exc.code, body):
|
||||||
raise RuntimeError("gtoken API blocked by Cloudflare (1010)") from exc
|
raise RuntimeError("gtoken API blocked by Cloudflare (1010)") from exc
|
||||||
raise RuntimeError(f"gtoken API status {exc.code}: {body[:200]}") from exc
|
raise RuntimeError(f"gtoken API status {exc.code}: {body[:200]}") from exc
|
||||||
@@ -315,13 +388,12 @@ def _token_via_gtoken(username: str, password: str) -> str:
|
|||||||
raise RuntimeError(f"gtoken API network error: {exc}") from exc
|
raise RuntimeError(f"gtoken API network error: {exc}") from exc
|
||||||
|
|
||||||
if _rate_limited(body, status):
|
if _rate_limited(body, status):
|
||||||
raise SystemExit(EXIT_RATE_LIMITED)
|
raise RateLimited(f"gtoken API: {body[:200]}")
|
||||||
return _parse_token_body(body, "gtoken API")
|
return _parse_token_body(body, "gtoken API")
|
||||||
|
|
||||||
|
|
||||||
def _token_via_meta(username: str, password: str, meta: WgServer) -> str:
|
def _token_via_meta(username: str, password: str, meta: WgServer) -> str:
|
||||||
ca_path = ensure_pia_ca()
|
context = pia_ssl_context()
|
||||||
context = ssl.create_default_context(cafile=str(ca_path))
|
|
||||||
conn = _HTTPSConnectionToIP(meta.ip, 443, meta.cn, context, timeout=10)
|
conn = _HTTPSConnectionToIP(meta.ip, 443, meta.cn, context, timeout=10)
|
||||||
try:
|
try:
|
||||||
conn.request(
|
conn.request(
|
||||||
@@ -339,13 +411,15 @@ def _token_via_meta(username: str, password: str, meta: WgServer) -> str:
|
|||||||
conn.close()
|
conn.close()
|
||||||
|
|
||||||
if _rate_limited(body, status):
|
if _rate_limited(body, status):
|
||||||
raise SystemExit(EXIT_RATE_LIMITED)
|
raise RateLimited(f"meta {meta.cn}: {body[:200]}")
|
||||||
if status != 200:
|
if status != 200:
|
||||||
raise RuntimeError(f"meta {meta.cn}/{meta.ip} status {status}: {body[:200]}")
|
raise RuntimeError(f"meta {meta.cn}/{meta.ip} status {status}: {body[:200]}")
|
||||||
return _parse_token_body(body, f"meta {meta.cn}")
|
return _parse_token_body(body, f"meta {meta.cn}")
|
||||||
|
|
||||||
|
|
||||||
def get_token(username: str, password: str, preferred_region: str | None = None) -> str:
|
def get_token(username: str, password: str, preferred_region: str | None = None) -> str:
|
||||||
|
raise_if_rate_limit_cooldown()
|
||||||
|
|
||||||
cache_path = token_cache_path()
|
cache_path = token_cache_path()
|
||||||
ttl = parse_duration_seconds(os.environ.get("TOKEN_CACHE_TTL", "20h"), 20 * 3600)
|
ttl = parse_duration_seconds(os.environ.get("TOKEN_CACHE_TTL", "20h"), 20 * 3600)
|
||||||
force = env_bool("FORCE_TOKEN_REFRESH", False)
|
force = env_bool("FORCE_TOKEN_REFRESH", False)
|
||||||
@@ -362,6 +436,7 @@ def get_token(username: str, password: str, preferred_region: str | None = None)
|
|||||||
pass
|
pass
|
||||||
|
|
||||||
errors: list[str] = []
|
errors: list[str] = []
|
||||||
|
saw_rate_limit = False
|
||||||
|
|
||||||
for label, getter in (
|
for label, getter in (
|
||||||
("central v2 token API", lambda: _token_via_central(username, password)),
|
("central v2 token API", lambda: _token_via_central(username, password)),
|
||||||
@@ -370,10 +445,13 @@ def get_token(username: str, password: str, preferred_region: str | None = None)
|
|||||||
try:
|
try:
|
||||||
log(f"Requesting PIA token via {label}")
|
log(f"Requesting PIA token via {label}")
|
||||||
token = getter()
|
token = getter()
|
||||||
|
clear_rate_limit()
|
||||||
log(f"Fetched and cached new PIA token via {label}")
|
log(f"Fetched and cached new PIA token via {label}")
|
||||||
return _store_token(token)
|
return _store_token(token)
|
||||||
except SystemExit:
|
except RateLimited as exc:
|
||||||
raise
|
saw_rate_limit = True
|
||||||
|
errors.append(f"{label}: {exc}")
|
||||||
|
log(f"Token via {label} rate-limited: {exc}")
|
||||||
except Exception as exc: # noqa: BLE001 - try next auth method
|
except Exception as exc: # noqa: BLE001 - try next auth method
|
||||||
errors.append(f"{label}: {exc}")
|
errors.append(f"{label}: {exc}")
|
||||||
log(f"Token via {label} failed: {exc}")
|
log(f"Token via {label} failed: {exc}")
|
||||||
@@ -405,14 +483,20 @@ def get_token(username: str, password: str, preferred_region: str | None = None)
|
|||||||
try:
|
try:
|
||||||
log(f"Requesting PIA token via {label}")
|
log(f"Requesting PIA token via {label}")
|
||||||
token = _token_via_meta(username, password, meta)
|
token = _token_via_meta(username, password, meta)
|
||||||
|
clear_rate_limit()
|
||||||
log(f"Fetched and cached new PIA token via {label}")
|
log(f"Fetched and cached new PIA token via {label}")
|
||||||
return _store_token(token)
|
return _store_token(token)
|
||||||
except SystemExit:
|
except RateLimited as exc:
|
||||||
raise
|
saw_rate_limit = True
|
||||||
|
errors.append(f"{label}: {exc}")
|
||||||
|
log(f"Token via {label} rate-limited: {exc}")
|
||||||
except Exception as exc: # noqa: BLE001
|
except Exception as exc: # noqa: BLE001
|
||||||
errors.append(f"{label}: {exc}")
|
errors.append(f"{label}: {exc}")
|
||||||
log(f"Token via {label} failed: {exc}")
|
log(f"Token via {label} failed: {exc}")
|
||||||
|
|
||||||
|
if saw_rate_limit:
|
||||||
|
mark_rate_limited(reason="token endpoints rate-limited")
|
||||||
|
raise SystemExit(EXIT_RATE_LIMITED)
|
||||||
raise SystemExit("PIA token request failed:\n- " + "\n- ".join(errors))
|
raise SystemExit("PIA token request failed:\n- " + "\n- ".join(errors))
|
||||||
|
|
||||||
|
|
||||||
@@ -437,21 +521,25 @@ class _HTTPSConnectionToIP(HTTPSConnection):
|
|||||||
|
|
||||||
|
|
||||||
def add_key(server: WgServer, token: str, public_key: str) -> dict[str, Any]:
|
def add_key(server: WgServer, token: str, public_key: str) -> dict[str, Any]:
|
||||||
ca_path = ensure_pia_ca()
|
raise_if_rate_limit_cooldown()
|
||||||
context = ssl.create_default_context(cafile=str(ca_path))
|
|
||||||
|
context = pia_ssl_context()
|
||||||
query = urllib.parse.urlencode({"pt": token, "pubkey": public_key})
|
query = urllib.parse.urlencode({"pt": token, "pubkey": public_key})
|
||||||
path = f"/addKey?{query}"
|
path = f"/addKey?{query}"
|
||||||
|
|
||||||
conn = _HTTPSConnectionToIP(server.ip, 1337, server.cn, context)
|
conn = _HTTPSConnectionToIP(server.ip, 1337, server.cn, context)
|
||||||
try:
|
try:
|
||||||
conn.request("GET", path, headers={"Content-Type": "application/json"})
|
conn.request("GET", path, headers={**HTTP_HEADERS, "Content-Type": "application/json"})
|
||||||
resp = conn.getresponse()
|
resp = conn.getresponse()
|
||||||
body = resp.read().decode("utf-8", errors="replace")
|
body = resp.read().decode("utf-8", errors="replace")
|
||||||
status = resp.status
|
status = resp.status
|
||||||
|
except ssl.SSLError as exc:
|
||||||
|
raise SystemExit(f"addKey TLS failed for {server.cn}/{server.ip}: {exc}") from exc
|
||||||
finally:
|
finally:
|
||||||
conn.close()
|
conn.close()
|
||||||
|
|
||||||
if _rate_limited(body, status):
|
if _rate_limited(body, status):
|
||||||
|
mark_rate_limited(reason=f"addKey {server.cn}")
|
||||||
raise SystemExit(EXIT_RATE_LIMITED)
|
raise SystemExit(EXIT_RATE_LIMITED)
|
||||||
if status != 200:
|
if status != 200:
|
||||||
raise SystemExit(f"addKey failed for {server.cn}/{server.ip}: status {status}: {body}")
|
raise SystemExit(f"addKey failed for {server.cn}/{server.ip}: status {status}: {body}")
|
||||||
|
|||||||
Reference in New Issue
Block a user