Data Extraction & Scraping

Open Source Stack to Replace Bright Data (CAPTCHA + Residential)

By Anass Benameur Updated: 1 September 2026

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
LinkedIn 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)