Open Source Stack to Replace Bright Data (CAPTCHA + Residential)
Realistic answer upfront: no open-source tool bypasses every CAPTCHA reliably. But you can get 80-90% of what Web Unlocker does by stacking the right tools. Here's the actual stack.
CAPTCHA Bypass by Type
| CAPTCHA Type | Best Open-Source Tool | Success Rate |
|---|---|---|
| Cloudflare "Just a moment..." | FlareSolverr | 90-95% |
| Cloudflare Turnstile | CloudflareBypasser (Puppeteer + patchright) | 70-85% |
| reCAPTCHA v2 (image) | Buster extension (audio challenge) | 60-70% |
| reCAPTCHA v3 (score) | undetected-chromedriver + residential IP | 50-70% |
| hCaptcha | hcaptcha-challenger (ML-based) | 40-60% |
| Text/Image CAPTCHA | Tesseract OCR or ddddocr | 70-90% |
| Arkose/FunCaptcha | No reliable OSS solution | — |
| GeeTest | geetest-crack (unmaintained) | 20-40% |
The Core Tools
1. FlareSolverr — Cloudflare's biggest killer
docker run -d --name flaresolverr \
-p 8191:8191 \
--restart unless-stopped \
ghcr.io/flaresolverr/flaresolverr:latest
Your scraper hits http://localhost:8191/v1 with the target URL. It spins up a headless browser, solves the challenge, returns cookies + HTML. Handles cf_clearance cookies automatically.
2. undetected-chromedriver / patchright / camoufox — Stealth browsers
pip install undetected-chromedriver
# or better:
pip install patchright # patched Playwright, more actively maintained
pip install camoufox # patched Firefox, best fingerprint spoofing
Camoufox is the current gold standard for anti-detection. Better than any Playwright-stealth plugin.
3. ddddocr — Text CAPTCHA solver
pip install ddddocr
Feed it a CAPTCHA image, get the text. Works on 4-6 character alphanumeric CAPTCHAs with high accuracy. This alone would handle most Moroccan gov site CAPTCHAs if they use simple images.
4. hcaptcha-challenger — ML-based hCaptcha solver
pip install hcaptcha-challenger
Uses YOLO models to identify objects in hCaptcha challenges. Requires GPU for best speed but works on CPU too.
5. Buster — reCAPTCHA audio bypass (browser extension)
Loads as extension in Playwright/Selenium. Clicks the audio challenge, downloads audio, transcribes via speech-to-text API, submits answer. Use with speech_recognition + Google's free tier or Whisper.
Free Residential Proxy Alternatives
Since Bright Data is dead for you, here's the reality:
| Source | Cost | IPs | Reliability |
|---|---|---|---|
| Tor | Free | ~7000 exit nodes | Medium — many sites block Tor |
| ProxyBroker2 (scrapes free proxy lists) | Free | ~500-2000 working | Low — most die within hours |
| proxy-scraper-checker (github) | Free | Thousands | Low |
| Public SOCKS lists (spys.one, free-proxy-list.net) | Free | Varies | Low |
| Webshare free tier | Free | 10 datacenter | Medium — datacenter, not residential |
| VPN rotation (ProtonVPN free, Windscribe 10GB/mo) | Free | ~50 countries | High but slow |
The Honest Truth About "Free Residential"
There's no such thing as free residential proxies at scale. The IPs you find on scraper lists are:
- Compromised routers (illegal to use)
- Datacenter IPs mislabeled as residential
- Honeypots run by anti-fraud companies
- Dead within hours
The Actual Stack I'd Build
For your Morocco project, here's what I'd put together:
core/captcha_solver.py
import asyncio
import base64
from typing import Optional, Dict, Any
from pathlib import Path
import structlog
import httpx
import ddddocr
logger = structlog.get_logger()
class CaptchaSolver:
"""Unified CAPTCHA solver using open-source tools."""
def __init__(self, flaresolverr_url: str = "http://localhost:8191/v1"):
self.flaresolverr_url = flaresolverr_url
self.ocr = ddddocr.DdddOcr(show_ad=False)
self.slide_ocr = ddddocr.DdddOcr(det=False, ocr=False, show_ad=False)
self._client = httpx.AsyncClient(timeout=120.0)
async def solve_cloudflare(self, url: str, cookies: Optional[Dict] = None) -> Dict[str, Any]:
"""Solve Cloudflare challenge via FlareSolverr."""
payload = {
"cmd": "request.get",
"url": url,
"maxTimeout": 60000,
}
if cookies:
payload["cookies"] = [
{"name": k, "value": v} for k, v in cookies.items()
]
try:
response = await self._client.post(
self.flaresolverr_url,
json=payload
)
result = response.json()
if result.get("status") == "ok":
solution = result["solution"]
return {
"success": True,
"html": solution["response"],
"cookies": {c["name"]: c["value"] for c in solution.get("cookies", [])},
"user_agent": solution.get("userAgent"),
"status_code": solution.get("status"),
}
logger.warning("FlareSolverr failed", message=result.get("message"))
return {"success": False, "error": result.get("message")}
except Exception as e:
logger.error("FlareSolverr request failed", error=str(e))
return {"success": False, "error": str(e)}
def solve_image_captcha(self, image_bytes: bytes) -> Optional[str]:
"""Solve text-based image CAPTCHA using ddddocr."""
try:
result = self.ocr.classification(image_bytes)
logger.debug("Image CAPTCHA solved", text=result)
return result
except Exception as e:
logger.error("Image CAPTCHA solve failed", error=str(e))
return None
def solve_slider_captcha(self, bg_image: bytes, slider_image: bytes) -> Optional[int]:
"""Solve slider CAPTCHA - returns X coordinate of gap."""
try:
result = self.slide_ocr.slide_match(slider_image, bg_image, simple_target=True)
return result.get("target", [None])[0]
except Exception as e:
logger.error("Slider solve failed", error=str(e))
return None
async def close(self):
await self._client.aclose()
core/stealth_browser.py
import asyncio
from typing import Optional, Dict, Any
from pathlib import Path
import structlog
logger = structlog.get_logger()
class StealthBrowser:
"""Camoufox-based stealth browser for JS-heavy targets."""
def __init__(self, tor_proxy: bool = True):
self.tor_proxy = tor_proxy
self._browser = None
self._context = None
async def start(self):
from camoufox.async_api import AsyncCamoufox
proxy_config = None
if self.tor_proxy:
proxy_config = {
"server": "socks5://127.0.0.1:9050"
}
self._browser = await AsyncCamoufox(
headless=True,
humanize=True,
geoip=True,
proxy=proxy_config,
os=["windows", "macos", "linux"],
block_images=True,
block_webrtc=True,
).__aenter__()
logger.info("Camoufox browser started", tor=self.tor_proxy)
async def fetch(self, url: str, wait_for: Optional[str] = None) -> Dict[str, Any]:
"""Fetch URL with stealth browser."""
page = await self._browser.new_page()
try:
await page.goto(url, wait_until="domcontentloaded", timeout=60000)
if wait_for:
await page.wait_for_selector(wait_for, timeout=30000)
await asyncio.sleep(2)
html = await page.content()
cookies = await page.context.cookies()
return {
"success": True,
"html": html,
"cookies": {c["name"]: c["value"] for c in cookies},
"url": page.url,
}
except Exception as e:
logger.error("Stealth fetch failed", url=url, error=str(e))
return {"success": False, "error": str(e)}
finally:
await page.close()
async def close(self):
if self._browser:
await self._browser.__aexit__(None, None, None)
core/proxy_rotator.py — Free proxy scraping
import asyncio
import random
import time
from typing import List, Optional, Set
from dataclasses import dataclass, field
import httpx
import structlog
logger = structlog.get_logger()
@dataclass
class FreeProxy:
host: str
port: int
protocol: str = "http"
last_checked: float = 0
success_count: int = 0
failure_count: int = 0
latency_ms: int = 0
@property
def url(self) -> str:
return f"{self.protocol}://{self.host}:{self.port}"
@property
def score(self) -> float:
total = self.success_count + self.failure_count
if total == 0:
return 0.5
return self.success_count / total
class FreeProxyRotator:
"""Scrapes and validates free proxy lists."""
PROXY_SOURCES = [
"https://raw.githubusercontent.com/TheSpeedX/PROXY-List/master/http.txt",
"https://raw.githubusercontent.com/TheSpeedX/PROXY-List/master/socks5.txt",
"https://raw.githubusercontent.com/monosans/proxy-list/main/proxies/http.txt",
"https://raw.githubusercontent.com/monosans/proxy-list/main/proxies/socks5.txt",
"https://raw.githubusercontent.com/hookzof/socks5_list/master/proxy.txt",
"https://api.proxyscrape.com/v2/?request=getproxies&protocol=http&timeout=10000&country=all",
]
TEST_URL = "http://httpbin.org/ip"
def __init__(self, max_proxies: int = 500):
self.max_proxies = max_proxies
self.proxies: List[FreeProxy] = []
self._blacklist: Set[str] = set()
self._lock = asyncio.Lock()
async def refresh(self):
"""Fetch fresh proxy lists and validate them."""
logger.info("Refreshing free proxy list")
raw_proxies = set()
async with httpx.AsyncClient(timeout=15) as client:
for source in self.PROXY_SOURCES:
try:
response = await client.get(source)
if response.status_code == 200:
protocol = "socks5" if "socks5" in source else "http"
for line in response.text.strip().split("\n"):
line = line.strip()
if ":" in line and line not in self._blacklist:
raw_proxies.add((line, protocol))
except Exception as e:
logger.warning("Proxy source failed", source=source, error=str(e))
logger.info("Fetched raw proxies", count=len(raw_proxies))
candidates = list(raw_proxies)[:self.max_proxies * 3]
semaphore = asyncio.Semaphore(50)
async def validate(proxy_line: str, protocol: str) -> Optional[FreeProxy]:
async with semaphore:
try:
host, port = proxy_line.split(":")
proxy = FreeProxy(
host=host,
port=int(port),
protocol=protocol
)
start = time.time()
async with httpx.AsyncClient(
proxy=proxy.url,
timeout=10
) as client:
response = await client.get(self.TEST_URL)
if response.status_code == 200:
proxy.latency_ms = int((time.time() - start) * 1000)
proxy.success_count = 1
proxy.last_checked = time.time()
return proxy
except Exception:
return None
return None
results = await asyncio.gather(
*[validate(p, proto) for p, proto in candidates],
return_exceptions=True
)
working = [r for r in results if isinstance(r, FreeProxy)]
working.sort(key=lambda p: p.latency_ms)
async with self._lock:
self.proxies = working[:self.max_proxies]
logger.info("Proxy validation complete",
working=len(self.proxies),
tested=len(candidates))
async def get_proxy(self) -> Optional[FreeProxy]:
"""Get a random working proxy."""
async with self._lock:
if not self.proxies:
return None
top_proxies = sorted(
self.proxies,
key=lambda p: (p.score, -p.latency_ms),
reverse=True
)[:20]
return random.choice(top_proxies) if top_proxies else None
async def report_failure(self, proxy: FreeProxy):
async with self._lock:
proxy.failure_count += 1
if proxy.failure_count >= 3:
self._blacklist.add(f"{proxy.host}:{proxy.port}")
self.proxies = [p for p in self.proxies if p != proxy]
async def report_success(self, proxy: FreeProxy):
async with self._lock:
proxy.success_count += 1
proxy.last_checked = time.time()
The Honest Reality Check for Morocco
For your specific targets:
| Site | Protection Level | What to Use |
|---|---|---|
| directinfo.ma | Low — no CAPTCHA usually | Tor + httpx, no browser needed |
| ae.gov.ma | Low | Same |
| marchespublics.gov.ma | None to low | Same |
| charika.ma | Low — sometimes rate limits | Tor + delays |
| kerix.net | Low | Tor |
| Cloudflare-protected sites | Medium | FlareSolverr + Tor |
| Very High | Skip — you can't crack this without residential | |
| Google Maps | High | Skip — needs API key or residential |
You probably don't need CAPTCHA solvers for most of the Moroccan gov sites. They just don't use them heavily. The whole CAPTCHA stack is overkill for your actual targets.
Complete docker-compose.yml for the Stack
version: '3.8'
services:
tor:
image: dperson/torproxy
container_name: tor
restart: unless-stopped
ports:
- "9050:9050"
- "9051:9051"
environment:
- TOR_NewCircuitPeriod=30
- TOR_MaxCircuitDirtiness=60
flaresolverr:
image: ghcr.io/flaresolverr/flaresolverr:latest
container_name: flaresolverr
restart: unless-stopped
ports:
- "8191:8191"
environment:
- LOG_LEVEL=info
- TZ=Africa/Casablanca
mem_limit: 512m
scraper:
build: .
depends_on:
- tor
- flaresolverr
environment:
- TOR_SOCKS_HOST=tor
- TOR_SOCKS_PORT=9050
- FLARESOLVERR_URL=http://flaresolverr:8191/v1
volumes:
- ./data:/opt/morocco-seeds/data
Install Commands
# On your VPS
pip install ddddocr camoufox patchright httpx[socks] hcaptcha-challenger
# Install Camoufox browser binary
python -m camoufox fetch
# FlareSolverr via Docker
docker run -d --name flaresolverr \
-p 8191:8191 \
--restart unless-stopped \
-e LOG_LEVEL=info \
-e TZ=Africa/Casablanca \
--memory=512m \
ghcr.io/flaresolverr/flaresolverr:latest
# Tor via Docker (or install locally with apt)
docker run -d --name tor \
-p 9050:9050 \
-p 9051:9051 \
--restart unless-stopped \
-e TOR_NewCircuitPeriod=30 \
-e TOR_MaxCircuitDirtiness=60 \
dperson/torproxy
# Verify everything works
curl -x socks5://localhost:9050 https://check.torproject.org/api/ip
curl -X POST http://localhost:8191/v1 \
-H "Content-Type: application/json" \
-d '{"cmd":"request.get","url":"https://www.google.com","maxTimeout":60000}'
Integrating the Stack into Your Scraper
Now the actual wiring — how your existing TorHttpClient decides when to escalate from plain Tor → FlareSolverr → Camoufox. This is a tiered fetcher pattern.
core/tiered_fetcher.py
import asyncio
from enum import Enum
from typing import Optional, Dict, Any, Callable
from dataclasses import dataclass
import structlog
from core.tor_client import TorHttpClient
from core.captcha_solver import CaptchaSolver
from core.stealth_browser import StealthBrowser
logger = structlog.get_logger()
class FetchTier(str, Enum):
"""Escalation tiers, cheapest to most expensive."""
TOR_HTTPX = "tor_httpx" # ~50ms, ~0 cost
FLARESOLVERR = "flaresolverr" # ~5s, moderate cost
STEALTH_BROWSER = "stealth" # ~10s, high cost
@dataclass
class FetchResult:
success: bool
html: Optional[str] = None
status_code: Optional[int] = None
cookies: Optional[Dict[str, str]] = None
tier_used: Optional[FetchTier] = None
error: Optional[str] = None
latency_ms: int = 0
class TieredFetcher:
"""
Fetches URLs with automatic escalation:
Tier 1: Tor + httpx (fast, cheap)
Tier 2: FlareSolverr (Cloudflare-aware)
Tier 3: Camoufox stealth browser (heaviest weapon)
Learns per-domain which tier works and skips lower tiers on repeat hits.
"""
CHALLENGE_MARKERS = [
"Just a moment",
"Checking your browser",
"cf-browser-verification",
"cf_chl_",
"challenge-platform",
"captcha-delivery",
"hcaptcha",
"g-recaptcha",
"px-captcha",
]
def __init__(
self,
tor_client: TorHttpClient,
captcha_solver: CaptchaSolver,
stealth_browser: Optional[StealthBrowser] = None,
):
self.tor = tor_client
self.captcha = captcha_solver
self.stealth = stealth_browser
self._domain_tier: Dict[str, FetchTier] = {}
self._lock = asyncio.Lock()
def _detect_challenge(self, html: str, status: int) -> bool:
"""Return True if response looks like an anti-bot challenge."""
if status in (403, 503, 429):
return True
if not html:
return False
lowered = html[:5000].lower()
return any(marker.lower() in lowered for marker in self.CHALLENGE_MARKERS)
async def _remember_tier(self, domain: str, tier: FetchTier):
async with self._lock:
self._domain_tier[domain] = tier
async def _preferred_tier(self, domain: str) -> FetchTier:
async with self._lock:
return self._domain_tier.get(domain, FetchTier.TOR_HTTPX)
async def fetch(
self,
url: str,
domain: str,
method: str = "GET",
**kwargs
) -> FetchResult:
import time
start = time.time()
preferred = await self._preferred_tier(domain)
tier_order = {
FetchTier.TOR_HTTPX: [FetchTier.TOR_HTTPX, FetchTier.FLARESOLVERR, FetchTier.STEALTH_BROWSER],
FetchTier.FLARESOLVERR: [FetchTier.FLARESOLVERR, FetchTier.STEALTH_BROWSER, FetchTier.TOR_HTTPX],
FetchTier.STEALTH_BROWSER: [FetchTier.STEALTH_BROWSER, FetchTier.FLARESOLVERR, FetchTier.TOR_HTTPX],
}[preferred]
last_error = None
for tier in tier_order:
try:
result = await self._try_tier(tier, url, method, **kwargs)
if result.success and result.html:
if not self._detect_challenge(result.html, result.status_code or 200):
result.tier_used = tier
result.latency_ms = int((time.time() - start) * 1000)
await self._remember_tier(domain, tier)
logger.info("Fetch successful",
url=url,
tier=tier.value,
latency_ms=result.latency_ms)
return result
else:
logger.info("Challenge detected, escalating",
url=url,
current_tier=tier.value)
last_error = "challenge_detected"
continue
last_error = result.error
except Exception as e:
last_error = str(e)
logger.warning("Tier failed, escalating",
url=url,
tier=tier.value,
error=str(e))
continue
return FetchResult(
success=False,
error=last_error or "all_tiers_exhausted",
latency_ms=int((time.time() - start) * 1000)
)
async def _try_tier(
self,
tier: FetchTier,
url: str,
method: str,
**kwargs
) -> FetchResult:
if tier == FetchTier.TOR_HTTPX:
return await self._fetch_tor(url, method, **kwargs)
elif tier == FetchTier.FLARESOLVERR:
return await self._fetch_flaresolverr(url, **kwargs)
elif tier == FetchTier.STEALTH_BROWSER:
return await self._fetch_stealth(url, **kwargs)
raise ValueError(f"Unknown tier: {tier}")
async def _fetch_tor(self, url: str, method: str, **kwargs) -> FetchResult:
response = await self.tor.request(method, url, **kwargs)
return FetchResult(
success=response.status_code < 400,
html=response.text,
status_code=response.status_code,
cookies=dict(response.cookies),
)
async def _fetch_flaresolverr(self, url: str, **kwargs) -> FetchResult:
result = await self.captcha.solve_cloudflare(url, cookies=kwargs.get("cookies"))
return FetchResult(
success=result.get("success", False),
html=result.get("html"),
status_code=result.get("status_code"),
cookies=result.get("cookies"),
error=result.get("error"),
)
async def _fetch_stealth(self, url: str, **kwargs) -> FetchResult:
if not self.stealth:
return FetchResult(success=False, error="stealth_browser_not_configured")
result = await self.stealth.fetch(url, wait_for=kwargs.get("wait_for"))
return FetchResult(
success=result.get("success", False),
html=result.get("html"),
cookies=result.get("cookies"),
error=result.get("error"),
)
The learning behavior matters here. First hit to charika.ma might trickle through all three tiers because you don't know its defenses yet. Second hit skips straight to whichever tier worked. Over a run of 500K requests, this saves you hours of unnecessary browser spin-ups.
Session Warmup for Cookie/CSRF Sites
DirectInfo and marchespublics almost certainly need a warm-up GET before their POST endpoints work. Add this layer:
core/session_pool.py
import asyncio
import time
import random
from typing import Dict, Optional, List
from dataclasses import dataclass, field
from urllib.parse import urlparse
import structlog
from core.tiered_fetcher import TieredFetcher
logger = structlog.get_logger()
@dataclass
class SessionContext:
"""A warmed-up session with cookies + tokens ready to use."""
domain: str
cookies: Dict[str, str] = field(default_factory=dict)
csrf_token: Optional[str] = None
user_agent: Optional[str] = None
created_at: float = field(default_factory=time.time)
request_count: int = 0
max_requests: int = 200
max_age_seconds: int = 1800
@property
def is_expired(self) -> bool:
age = time.time() - self.created_at
return age > self.max_age_seconds or self.request_count >= self.max_requests
def use(self):
self.request_count += 1
class SessionPool:
"""
Pool of warmed-up sessions per domain.
Rotates sessions to distribute load and avoid single-session detection.
"""
def __init__(
self,
fetcher: TieredFetcher,
sessions_per_domain: int = 3,
):
self.fetcher = fetcher
self.sessions_per_domain = sessions_per_domain
self._pools: Dict[str, List[SessionContext]] = {}
self._warmup_configs: Dict[str, Dict] = {}
self._locks: Dict[str, asyncio.Lock] = {}
def register_warmup(
self,
domain: str,
warmup_url: str,
csrf_selector: Optional[str] = None,
csrf_cookie_name: Optional[str] = None,
):
"""Configure how a domain's sessions get warmed up."""
self._warmup_configs[domain] = {
"warmup_url": warmup_url,
"csrf_selector": csrf_selector,
"csrf_cookie_name": csrf_cookie_name,
}
self._locks[domain] = asyncio.Lock()
self._pools[domain] = []
async def _warm_session(self, domain: str) -> Optional[SessionContext]:
config = self._warmup_configs.get(domain)
if not config:
logger.warning("No warmup config for domain", domain=domain)
return SessionContext(domain=domain)
result = await self.fetcher.fetch(config["warmup_url"], domain=domain)
if not result.success:
logger.warning("Session warmup failed", domain=domain, error=result.error)
return None
session = SessionContext(
domain=domain,
cookies=result.cookies or {},
)
if config.get("csrf_cookie_name"):
session.csrf_token = session.cookies.get(config["csrf_cookie_name"])
elif config.get("csrf_selector") and result.html:
from selectolax.parser import HTMLParser
tree = HTMLParser(result.html)
el = tree.css_first(config["csrf_selector"])
if el:
session.csrf_token = (
el.attributes.get("content")
or el.attributes.get("value")
)
logger.info("Session warmed",
domain=domain,
has_csrf=bool(session.csrf_token),
cookie_count=len(session.cookies))
return session
async def acquire(self, domain: str) -> Optional[SessionContext]:
"""Get a fresh session for the given domain."""
lock = self._locks.get(domain)
if not lock:
self._locks[domain] = asyncio.Lock()
lock = self._locks[domain]
async with lock:
pool = self._pools.setdefault(domain, [])
pool[:] = [s for s in pool if not s.is_expired]
while len(pool) < self.sessions_per_domain:
session = await self._warm_session(domain)
if session:
pool.append(session)
else:
break
await asyncio.sleep(random.uniform(1.0, 3.0))
if not pool:
return None
session = random.choice(pool)
session.use()
return session
async def invalidate(self, session: SessionContext):
"""Mark a session as dead (e.g., after 403/captcha)."""
async with self._locks.get(session.domain, asyncio.Lock()):
pool = self._pools.get(session.domain, [])
if session in pool:
pool.remove(session)
logger.info("Session invalidated", domain=session.domain)
Usage in an OMPIC-style scraper:
# On init
self.sessions = SessionPool(self.fetcher, sessions_per_domain=3)
self.sessions.register_warmup(
domain="directinfo.ma",
warmup_url="https://directinfo.ma/",
csrf_selector='meta[name="csrf-token"]',
)
# Per request
session = await self.sessions.acquire("directinfo.ma")
if not session:
logger.error("No sessions available")
return None
headers = {"X-CSRF-Token": session.csrf_token} if session.csrf_token else {}
result = await self.fetcher.fetch(
url,
domain="directinfo.ma",
method="POST",
json=payload,
cookies=session.cookies,
headers=headers,
)
if result.status_code in (401, 403):
await self.sessions.invalidate(session)
One Last Piece — Token Bucket Rate Limiter Done Right
Since you're going to layer per-host limits over the concurrency semaphore from earlier, here's the actual working version:
core/token_bucket.py
import asyncio
import time
from typing import Dict
from aiolimiter import AsyncLimiter
import structlog
logger = structlog.get_logger()
class PerHostLimiter:
"""
Token bucket rate limiter, one bucket per host.
Configurable per-domain, with adaptive slowdown on 429/503.
"""
DEFAULT_RATES = {
"directinfo.ma": (5, 1.0),
"ae.gov.ma": (8, 1.0),
"marchespublics.gov.ma": (2, 1.0),
"charika.ma": (4, 1.0),
"kerix.net": (3, 1.0),
}
def __init__(self):
self._limiters: Dict[str, AsyncLimiter] = {}
self._current_rates: Dict[str, float] = {}
self._penalty_until: Dict[str, float] = {}
self._lock = asyncio.Lock()
async def _get_limiter(self, domain: str) -> AsyncLimiter:
async with self._lock:
if domain not in self._limiters:
rate, period = self.DEFAULT_RATES.get(domain, (3, 1.0))
self._limiters[domain] = AsyncLimiter(rate, period)
self._current_rates[domain] = rate
return self._limiters[domain]
async def acquire(self, domain: str):
"""Wait until a token is available for this domain."""
penalty_end = self._penalty_until.get(domain, 0)
if penalty_end > time.time():
wait = penalty_end - time.time()
logger.debug("In penalty box", domain=domain, wait_seconds=round(wait, 1))
await asyncio.sleep(wait)
limiter = await self._get_limiter(domain)
async with limiter:
pass
async def penalize(self, domain: str, status_code: int):
"""React to a rate-limit signal from the server."""
async with self._lock:
current = self._current_rates.get(domain, 3)
if status_code == 429:
new_rate = max(0.5, current / 2)
penalty_duration = 60
elif status_code == 503:
new_rate = max(0.5, current * 0.7)
penalty_duration = 30
else:
return
self._current_rates[domain] = new_rate
self._limiters[domain] = AsyncLimiter(new_rate, 1.0)
self._penalty_until[domain] = time.time() + penalty_duration
logger.warning("Rate limit penalty applied",
domain=domain,
old_rate=current,
new_rate=new_rate,
penalty_duration=penalty_duration)