from __future__ import annotations import asyncio import json import logging import time from typing import Callable import httpx log = logging.getLogger(__name__) # Secondary fallback endpoints for IP resolution _IP_FALLBACKS = [ "https://api.ipify.org?format=json", "https://httpbin.org/ip", "https://ifconfig.me/ip", ] def _parse_ip(text: str) -> str | None: """Parse an IP string from JSON or plain-text response body.""" body = text.strip() if not body: return None if body.startswith("{"): try: d = json.loads(body) ip = d.get("ip") or d.get("origin") or d.get("query") return str(ip).split(",")[0].strip() if ip else None except Exception: return None # Plain-text: must look like an IP (short, no HTML) if len(body) < 50 and "<" not in body: candidate = body.split()[0].strip() if candidate.count(".") >= 3 or ":" in candidate: return candidate return None async def validate_proxies( proxy_urls: list[str], check_url: str, concurrency: int, timeout_seconds: float, on_progress: Callable[[int, int], None] | None = None, ) -> list[str]: if not proxy_urls: return [] sem = asyncio.Semaphore(max(1, concurrency)) ok: list[str] = [] lock = asyncio.Lock() total = len(proxy_urls) done = 0 c = max(1, concurrency) waves = (total + c - 1) // c cap = min(600.0, float(waves) * float(timeout_seconds) + 45.0) cap = max(90.0, cap) log.info( "validate_proxies: n=%d concurrency=%d per_timeout=%.1fs wall_cap~=%.0fs check=%s", total, c, timeout_seconds, cap, (check_url[:70] + "…") if len(check_url) > 70 else check_url, ) t = httpx.Timeout(timeout_seconds, connect=min(8.0, timeout_seconds)) lim = max(32, min(256, c * 8)) limits = httpx.Limits(max_connections=lim, max_keepalive_connections=max(16, c * 4)) last_prog_t = 0.0 prog_step = max(1, total // 120) async def _run_all(client: httpx.AsyncClient) -> None: nonlocal done, last_prog_t async def one(u: str) -> None: nonlocal done, last_prog_t async with sem: try: r = await client.get(check_url, proxy=u) good = r.status_code == 200 and len(r.content) > 0 except Exception: good = False async with lock: done += 1 if good: ok.append(u) if on_progress: now = time.monotonic() if ( done == total or done == 1 or done % prog_step == 0 or (now - last_prog_t) >= 0.1 ): last_prog_t = now on_progress(done, total) await asyncio.gather(*(one(u) for u in proxy_urls)) try: async with httpx.AsyncClient( timeout=t, verify=False, follow_redirects=True, limits=limits, ) as client: await asyncio.wait_for(_run_all(client), timeout=cap) except asyncio.TimeoutError: log.warning( "validate_proxies wall-clock cap %.0fs exceeded (%d URLs) — returning partial results", cap, total, ) log.info("validate_proxies: finished ok=%d / %d", len(ok), total) return ok async def check_chain_exit_ip( listen_proxy: str, check_url: str, timeout_seconds: float, ) -> str | None: """Query the IP-check URL through the local chain proxy. Bounded total time — never stacks one slow request per fallback URL forever.""" per = max(5.0, min(20.0, float(timeout_seconds))) budget = max(15.0, min(60.0, float(timeout_seconds) * 2 + 5.0)) urls = [check_url] + [u for u in _IP_FALLBACKS if u != check_url][:2] log.debug( "check_chain_exit_ip: budget=%.1fs per_req=%.1fs fallback_count=%d", budget, per, len(urls), ) async def _run() -> str | None: t = httpx.Timeout(per, connect=min(8.0, per)) try: async with httpx.AsyncClient( proxy=listen_proxy, timeout=t, verify=False, follow_redirects=True, ) as c: for url in urls: try: r = await c.get(url) if r.status_code == 200: ip = _parse_ip(r.text) if ip: return ip except Exception: continue except Exception: pass return None try: return await asyncio.wait_for(_run(), timeout=budget) except asyncio.TimeoutError: log.debug("check_chain_exit_ip timed out after %.1fs", budget) return None async def get_direct_ip(check_url: str, timeout_seconds: float = 15.0) -> str | None: """Get our own exit IP without the proxy chain (so we can compare).""" per = max(5.0, min(15.0, float(timeout_seconds))) budget = max(12.0, min(45.0, float(timeout_seconds) * 2)) urls = [check_url] + [u for u in _IP_FALLBACKS if u != check_url][:2] async def _run() -> str | None: t = httpx.Timeout(per) try: async with httpx.AsyncClient( timeout=t, verify=False, follow_redirects=True, ) as c: for url in urls: try: r = await c.get(url) if r.status_code == 200: ip = _parse_ip(r.text) if ip: return ip except Exception: continue except Exception: pass return None try: return await asyncio.wait_for(_run(), timeout=budget) except asyncio.TimeoutError: log.debug("get_direct_ip timed out after %.1fs", budget) return None