#!/usr/bin/env python3
# =============================================================================
# proxyDrum
# =============================================================================
#
# FEATURES (reference / changelog)
# --------------------------------
#   1.  Three-thread architecture.
#         Thread 1 - Fetcher    : scans sources on a self-tuned schedule.
#         Thread 2 - Validator  : tests candidates, re-validates the pool,
#                                 and consumes the IPsum threat feed.
#         Thread 3 - HTTP proxy : serves clients via best-scoring proxy.
#         Thread 4 - Reporter   : emits a periodic structured status event.
#         Thread 5 - Adaptive   : recomputes pool target from client count.
#
#   2.  Fully environment-driven configuration. Every knob has a
#       PROXYDRUM_* env variable; CLI flags override the environment.
#
#   3.  Structured JSON logging (NDJSON) via structlog, with a colored
#       console renderer when PROXYDRUM_LOG_JSON=0. The periodic status
#       report is emitted through the same pipeline.
#
#   4.  Fluent client names derived deterministically from a fingerprint.
#
#   5.  Client fingerprinting: session-stable hash over (masked source IP,
#       User-Agent, Accept-Language, Accept-Encoding, Sec-CH-UA).
#
#   6.  Sticky session affinity: the same fingerprint is bound to the same
#       upstream proxy for the lifetime of the session.
#
#   7.  Scoring-based economy (no eviction, ever). Every proxy ever seen is
#       retained in memory. The composite score is a monotonic
#       "higher is better" value built from normalised components:
#         - speed_score     : 1 / (1 + response_time / ref)
#         - success_ratio   : successes / (successes + failures)
#         - streak_score    : min(1, success_streak / target)
#         - credit_score    : min(1, credit / credit_max)
#         - age_score       : min(1, age / age_cap)
#         - reputation_score: 1 for absent from IPsum, decaying to 0 as
#                             the number of listings rises
#       and a subtraction:
#         - failure_ratio   : failures / (successes + failures)
#       Default weights give reputation about 1/3 (specifically ~0.40) of
#       the positive contribution. All weights are configurable.
#
#       Long-lived proxies get progressively longer validation timeouts and
#       more retries, so transient network hiccups do not penalise them as
#       much as new proxies.
#
#   8.  Best-first pool ordering, continuous self-promotion. Proxies bubble
#       up as they accumulate streak, credit, age, and clean reputation;
#       they bubble down as they fail or appear on blacklists.
#
#   9.  Adaptive pool sizing with no upper bound.
#         target = max(TARGET_POOL_SIZE, active_clients * HEADROOM)
#       TARGET_POOL_SIZE is a floor (soft minimum) so proxies exist right
#       after startup and absorb initial connection bursts.
#
#   10. Self-tuning source scheduler. No per-source configured refresh
#       intervals. Each source's poll interval is learned from observed
#       change events (content hash differences between fetches).
#
#   11. Idle-aware pacing. When the pool is at target, the fetcher and
#       validator relax to their idle cadence. When the pool is below
#       target, they tighten.
#
#   12. Client wait queue. Bounded FIFO queue that stalls connecting clients
#       until a proxy becomes available.
#
#   13. Reserved proxies are never validated.
#
#   14. Dead-proxy handling on disconnect: penalise (not remove), so the
#       proxy sinks in the ranking and can recover later.
#
#   15. Periodic status report emitted through structlog at `info` level
#       with event name "report.status". The status is a compact overview
#       only: uptime, pool composition, proxy health, and client activity.
#       Per-proxy scores, top/bottom lists, and economy internals are not
#       emitted on stdout. The scoring system is opaque to the observer.
#
#   16. Content-hash based change detection for the scheduler.
#
#   17. No eviction. Every proxy ever visited is retained in memory
#       indefinitely. New candidates from source lists are appended to the
#       pool via the validator, which rewards or penalises them.
#
#   18. Continuous scanning. Sources are always scanned, at a lower
#       cadence when idle and a higher cadence under pressure.
#
#   19. Threat intelligence feed (IPsum). The validator periodically pulls
#       https://raw.githubusercontent.com/stamparm/ipsum/refs/heads/master/ipsum.txt
#       which maps IP addresses to the number of blocklists that IP appears
#       on. The number of listings contributes a reputation_score to the
#       composite score with a weight of at least one third of the positive
#       contributions by default.
#
#   20. Feed refresh policy. The feed is re-fetched at most once every
#       PROXYDRUM_IPSUM_REFRESH_INTERVAL seconds (default 24h), but the
#       decision is also driven by the "Last update:" header inside the
#       file: if that timestamp is older than the refresh interval, a new
#       copy is pulled even if the local interval has not elapsed.
#
# Copyright (c) 2006 Wizardry and Steamworks
# SPDX-License-Identifier: MIT
# =============================================================================

from __future__ import annotations

import argparse
import asyncio
import hashlib
import ipaddress
import logging
import os
import random
import re
import sys
import threading
import time
from collections import deque
from dataclasses import dataclass, field
from datetime import datetime, timezone
from email.utils import parsedate_to_datetime
from typing import Optional

import aiohttp
import structlog


# =============================================================================
# Configuration helpers (env-driven)
# =============================================================================

def _env_str(name: str, default: str) -> str:
    return os.environ.get(name, default)


def _env_int(name: str, default: int) -> int:
    try:
        return int(os.environ.get(name, default))
    except (TypeError, ValueError):
        return default


def _env_float(name: str, default: float) -> float:
    try:
        return float(os.environ.get(name, default))
    except (TypeError, ValueError):
        return default


def _env_bool(name: str, default: bool) -> bool:
    raw = os.environ.get(name)
    if raw is None:
        return default
    return raw.strip().lower() in ("1", "true", "yes", "on")


def _parse_sources(raw: str) -> list[str]:
    urls: list[str] = []
    for entry in raw.split(","):
        entry = entry.strip()
        if not entry:
            continue
        if "=" in entry:
            entry = entry.rpartition("=")[0].strip()
        if entry:
            urls.append(entry)
    return urls


# =============================================================================
# structlog configuration
# =============================================================================

def _configure_structlog(json_output: bool, level: int) -> None:
    logging.basicConfig(
        format="%(message)s",
        stream=sys.stdout,
        level=level,
    )

    timestamper = structlog.processors.TimeStamper(fmt="iso", utc=True)

    shared_processors = [
        structlog.contextvars.merge_contextvars,
        structlog.stdlib.add_log_level,
        structlog.stdlib.add_logger_name,
        timestamper,
        structlog.processors.StackInfoRenderer(),
        structlog.processors.format_exc_info,
        structlog.processors.UnicodeDecoder(),
    ]

    if json_output:
        renderer = structlog.processors.JSONRenderer(sort_keys=True)
    else:
        renderer = structlog.dev.ConsoleRenderer(colors=True)

    structlog.configure(
        processors=shared_processors + [
            structlog.stdlib.ProcessorFormatter.wrap_for_formatter,
        ],
        logger_factory=structlog.stdlib.LoggerFactory(),
        wrapper_class=structlog.stdlib.BoundLogger,
        cache_logger_on_first_use=True,
    )

    formatter = structlog.stdlib.ProcessorFormatter(
        foreign_pre_chain=shared_processors,
        processors=[
            structlog.stdlib.ProcessorFormatter.remove_processors_meta,
            renderer,
        ],
    )

    handler = logging.StreamHandler(sys.stdout)
    handler.setFormatter(formatter)

    root = logging.getLogger()
    root.handlers = [handler]
    root.setLevel(level)


def get_logger(name: str = "proxydrum") -> structlog.stdlib.BoundLogger:
    return structlog.get_logger(name)


# =============================================================================
# Fluent client name generator
# =============================================================================

_ADJECTIVES = [
    "calm", "brave", "swift", "quiet", "bright", "gentle", "clever",
    "nimble", "steady", "eager", "witty", "lively", "serene", "bold",
    "curious", "kind", "vivid", "mellow", "spry", "ardent", "sleek",
    "lucid", "candid", "fervid", "dapper", "jolly", "stout", "wistful",
    "zealous", "urbane", "lithe", "sable", "tawny", "russet", "amber",
    "azure", "coral", "ivory", "olive", "plum",
]

_ANIMALS = [
    "otter", "heron", "fox", "lynx", "badger", "falcon", "raven",
    "marten", "beaver", "hare", "stoat", "wren", "kestrel", "osprey",
    "tern", "puffin", "magpie", "jackdaw", "swift", "swallow",
    "ibis", "crane", "stork", "pelican", "cormorant", "curlew",
    "godwit", "sandpiper", "plover", "dunlin", "snipe", "woodcock",
    "quail", "partridge", "grouse", "ptarmigan", "capercaillie",
    "skylark", "pipit", "wagtail",
]


def _client_name_from_fingerprint(fingerprint: str) -> str:
    digest = hashlib.sha256(fingerprint.encode()).digest()
    adj_idx = digest[0] % len(_ADJECTIVES)
    animal_idx = digest[1] % len(_ANIMALS)
    return f"{_ADJECTIVES[adj_idx]}-{_ANIMALS[animal_idx]}"


# =============================================================================
# Default sources
# =============================================================================

_DEFAULT_SOURCES = [
    "https://raw.githubusercontent.com/proxygenerator1/ProxyGenerator/main/MostStable/http.txt",
]

_DEFAULT_IPSUM_URL = (
    "https://raw.githubusercontent.com/stamparm/ipsum/refs/heads/master/ipsum.txt"
)


@dataclass
class Config:
    # ---- proxy sources ------------------------------------------------- #
    proxy_sources: list[str] = field(
        default_factory=lambda: _parse_sources(
            _env_str("PROXYDRUM_SOURCES", ",".join(_DEFAULT_SOURCES))
        )
    )

    # ---- pool sizing (adaptive, no upper limit) ------------------------ #
    target_pool_size: int = field(
        default_factory=lambda: _env_int("PROXYDRUM_TARGET_POOL_SIZE", 10)
    )
    pool_headroom: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_POOL_HEADROOM", 3.0)
    )
    working_window_s: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_WORKING_WINDOW", 300.0)
    )

    # ---- scheduler (self-tuning) --------------------------------------- #
    source_min_interval: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_SOURCE_MIN_INTERVAL", 30.0)
    )
    source_max_interval: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_SOURCE_MAX_INTERVAL", 3600.0)
    )
    source_initial_interval: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_SOURCE_INITIAL_INTERVAL", 120.0)
    )
    source_safety_margin: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_SOURCE_SAFETY_MARGIN", 1.5)
    )
    source_history_size: int = field(
        default_factory=lambda: _env_int("PROXYDRUM_SOURCE_HISTORY", 8)
    )
    scheduler_tick: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_SCHEDULER_TICK", 5.0)
    )
    busy_interval_scale: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_BUSY_INTERVAL_SCALE", 0.5)
    )

    # ---- validation cadence ------------------------------------------- #
    validate_interval_idle: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_VALIDATE_INTERVAL_IDLE", 120.0)
    )
    validate_interval_busy: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_VALIDATE_INTERVAL_BUSY", 30.0)
    )
    validate_timeout_base: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_VALIDATE_TIMEOUT", 5.0)
    )
    validate_batch_size: int = field(
        default_factory=lambda: _env_int("PROXYDRUM_VALIDATE_BATCH_SIZE", 32)
    )
    adaptive_interval: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_ADAPTIVE_INTERVAL", 5.0)
    )

    # ---- economy / scoring system ------------------------------------- #
    score_weight_speed: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_SCORE_W_SPEED", 0.25)
    )
    score_weight_success: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_SCORE_W_SUCCESS", 0.10)
    )
    score_weight_streak: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_SCORE_W_STREAK", 0.08)
    )
    score_weight_credit: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_SCORE_W_CREDIT", 0.05)
    )
    score_weight_age: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_SCORE_W_AGE", 0.05)
    )
    score_weight_reputation: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_SCORE_W_REPUTATION", 0.35)
    )
    score_weight_failure: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_SCORE_W_FAILURE", 0.15)
    )
    score_speed_reference_ms: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_SCORE_SPEED_REF_MS", 1000.0)
    )
    credit_min: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_CREDIT_MIN", -100.0)
    )
    credit_max: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_CREDIT_MAX", 100.0)
    )
    credit_reward: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_CREDIT_REWARD", 1.0)
    )
    credit_penalty_base: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_CREDIT_PENALTY", 3.0)
    )
    longevity_streak_target: int = field(
        default_factory=lambda: _env_int("PROXYDRUM_LONGEVITY_STREAK", 10)
    )
    longevity_age_cap_s: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_LONGEVITY_AGE_CAP", 3600.0)
    )
    longevity_timeout_max: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_LONGEVITY_TIMEOUT_MAX", 3.0)
    )
    longevity_retries_max: int = field(
        default_factory=lambda: _env_int("PROXYDRUM_LONGEVITY_RETRIES_MAX", 2)
    )

    # ---- reputation / threat feed ------------------------------------- #
    ipsum_url: str = field(
        default_factory=lambda: _env_str("PROXYDRUM_IPSUM_URL", _DEFAULT_IPSUM_URL)
    )
    ipsum_enabled: bool = field(
        default_factory=lambda: _env_bool("PROXYDRUM_IPSUM_ENABLED", True)
    )
    ipsum_refresh_interval: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_IPSUM_REFRESH_INTERVAL",
                                           86400.0)
    )
    reputation_listings_soft: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_REPUTATION_LISTINGS_SOFT", 3.0)
    )
    reputation_listings_max: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_REPUTATION_LISTINGS_MAX", 10.0)
    )
    reputation_clean_score: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_REPUTATION_CLEAN_SCORE", 1.0)
    )

    # ---- test target --------------------------------------------------- #
    test_url: str = field(
        default_factory=lambda: _env_str(
            "PROXYDRUM_TEST_URL", "https://api.ipify.org")
    )

    # ---- HTTP proxy server -------------------------------------------- #
    listen_host: str = field(
        default_factory=lambda: _env_str("PROXYDRUM_LISTEN_HOST", "0.0.0.0")
    )
    listen_port: int = field(
        default_factory=lambda: _env_int("PROXYDRUM_LISTEN_PORT", 8080)
    )
    connect_timeout: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_CONNECT_TIMEOUT", 10.0)
    )
    upstream_timeout: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_UPSTREAM_TIMEOUT", 15.0)
    )

    # ---- client wait queue -------------------------------------------- #
    wait_queue_max: int = field(
        default_factory=lambda: _env_int("PROXYDRUM_WAIT_QUEUE_MAX", 200)
    )
    wait_timeout: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_WAIT_TIMEOUT", 30.0)
    )
    retry_after_hint: int = field(
        default_factory=lambda: _env_int("PROXYDRUM_RETRY_AFTER_HINT", 5)
    )

    # ---- affinity ------------------------------------------------------ #
    affinity_enabled: bool = field(
        default_factory=lambda: _env_bool("PROXYDRUM_AFFINITY_ENABLED", True)
    )
    affinity_ttl: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_AFFINITY_TTL", 1800.0)
    )
    affinity_max_entries: int = field(
        default_factory=lambda: _env_int("PROXYDRUM_AFFINITY_MAX_ENTRIES", 10000)
    )

    # ---- reporter ------------------------------------------------------ #
    report_interval: float = field(
        default_factory=lambda: _env_float("PROXYDRUM_REPORT_INTERVAL", 30.0)
    )

    # ---- logging ------------------------------------------------------- #
    log_level: int = field(
        default_factory=lambda: _env_int("PROXYDRUM_LOG_LEVEL", logging.INFO)
    )
    log_json: bool = field(
        default_factory=lambda: _env_bool("PROXYDRUM_LOG_JSON", True)
    )


# =============================================================================
# Runtime counters
# =============================================================================

class Metrics:
    def __init__(self) -> None:
        self._lock = threading.RLock()
        self.started_at = time.time()

        self.clients_connected = 0
        self.clients_total = 0
        self.clients_waiting = 0
        self.clients_rejected = 0
        self.clients_timed_out = 0
        self.requests_served = 0
        self.requests_failed = 0

        self.bytes_up = 0
        self.bytes_down = 0

        self.additions_total = 0
        self.rewards_total = 0
        self.penalties_total = 0

        self.ipsum_fetches = 0
        self.ipsum_entries = 0
        self.ipsum_last_fetch_at = 0.0
        self.ipsum_last_header_ts = 0.0

        self.last_validation_at = 0.0
        self.last_fetch_at = 0.0

        self._wait_samples: deque[float] = deque(maxlen=64)

    def client_connected(self) -> None:
        with self._lock:
            self.clients_connected += 1
            self.clients_total += 1

    def client_disconnected(self) -> None:
        with self._lock:
            self.clients_connected = max(0, self.clients_connected - 1)

    def client_waiting_enter(self) -> None:
        with self._lock:
            self.clients_waiting += 1

    def client_waiting_exit(self) -> None:
        with self._lock:
            self.clients_waiting = max(0, self.clients_waiting - 1)

    def client_rejected(self) -> None:
        with self._lock:
            self.clients_rejected += 1

    def client_timed_out(self) -> None:
        with self._lock:
            self.clients_timed_out += 1

    def record_wait(self, seconds: float) -> None:
        with self._lock:
            self._wait_samples.append(seconds)

    def average_wait(self) -> float:
        with self._lock:
            if not self._wait_samples:
                return 0.0
            return sum(self._wait_samples) / len(self._wait_samples)

    def request_served(self) -> None:
        with self._lock:
            self.requests_served += 1

    def request_failed(self) -> None:
        with self._lock:
            self.requests_failed += 1

    def add_bytes_up(self, n: int) -> None:
        with self._lock:
            self.bytes_up += n

    def add_bytes_down(self, n: int) -> None:
        with self._lock:
            self.bytes_down += n

    def proxy_added(self) -> None:
        with self._lock:
            self.additions_total += 1

    def proxy_rewarded(self) -> None:
        with self._lock:
            self.rewards_total += 1

    def proxy_penalized(self) -> None:
        with self._lock:
            self.penalties_total += 1

    def mark_validation(self) -> None:
        with self._lock:
            self.last_validation_at = time.time()

    def mark_fetch(self) -> None:
        with self._lock:
            self.last_fetch_at = time.time()

    def mark_ipsum(self, entries: int, header_ts: float) -> None:
        with self._lock:
            self.ipsum_fetches += 1
            self.ipsum_entries = entries
            self.ipsum_last_fetch_at = time.time()
            self.ipsum_last_header_ts = header_ts

    def snapshot(self) -> dict:
        with self._lock:
            now = time.time()
            return {
                "uptime_s": int(now - self.started_at),
                "clients_connected": self.clients_connected,
                "clients_total": self.clients_total,
                "clients_waiting": self.clients_waiting,
                "clients_rejected": self.clients_rejected,
                "clients_timed_out": self.clients_timed_out,
                "requests_served": self.requests_served,
                "requests_failed": self.requests_failed,
                "bytes_up": self.bytes_up,
                "bytes_down": self.bytes_down,
                "additions_total": self.additions_total,
                "rewards_total": self.rewards_total,
                "penalties_total": self.penalties_total,
                "ipsum_fetches": self.ipsum_fetches,
                "ipsum_entries": self.ipsum_entries,
                "ipsum_age_s": (
                    round(now - self.ipsum_last_fetch_at, 1)
                    if self.ipsum_last_fetch_at else None
                ),
                "avg_wait_s": round(self.average_wait(), 2),
                "last_validation_age": (
                    round(now - self.last_validation_at, 1)
                    if self.last_validation_at else None
                ),
                "last_fetch_age": (
                    round(now - self.last_fetch_at, 1)
                    if self.last_fetch_at else None
                ),
            }


# =============================================================================
# IPsum threat intelligence feed
# =============================================================================

class IpsumFeed:
    _LAST_UPDATE_RE = re.compile(r"#\s*Last update:\s*(.+?)\s*$", re.MULTILINE)
    _ENTRY_RE = re.compile(r"^(\d{1,3}(?:\.\d{1,3}){3})\s+(\d+)\s*$",
                           re.MULTILINE)

    def __init__(self, cfg: Config) -> None:
        self._cfg = cfg
        self._lock = threading.RLock()
        self._listings: dict[str, int] = {}
        self._header_ts: float = 0.0
        self._loaded: bool = False
        self.log = get_logger("proxydrum.ipsum")

    def is_loaded(self) -> bool:
        with self._lock:
            return self._loaded

    def header_age(self) -> Optional[float]:
        with self._lock:
            if not self._header_ts:
                return None
            return time.time() - self._header_ts

    def local_age(self, last_fetch_at: float) -> Optional[float]:
        if not last_fetch_at:
            return None
        return time.time() - last_fetch_at

    def size(self) -> int:
        with self._lock:
            return len(self._listings)

    def listings_for(self, ip: str) -> int:
        with self._lock:
            return self._listings.get(ip, 0)

    def parse(self, text: str) -> tuple[dict[str, int], float]:
        header_ts = 0.0
        m = self._LAST_UPDATE_RE.search(text)
        if m:
            try:
                dt = parsedate_to_datetime(m.group(1))
                if dt.tzinfo is None:
                    dt = dt.replace(tzinfo=timezone.utc)
                header_ts = dt.timestamp()
            except Exception:  # noqa: BLE001
                header_ts = 0.0

        mapping: dict[str, int] = {}
        for ip, count in self._ENTRY_RE.findall(text):
            try:
                mapping[ip] = int(count)
            except ValueError:
                continue
        return mapping, header_ts

    def install(self, mapping: dict[str, int], header_ts: float) -> None:
        with self._lock:
            self._listings = mapping
            self._header_ts = header_ts
            self._loaded = True

    async def maybe_refresh(self, session: aiohttp.ClientSession,
                            last_fetch_at: float) -> bool:
        if not self._cfg.ipsum_enabled:
            return False

        now = time.time()
        local_age = self.local_age(last_fetch_at)
        header_age = self.header_age()

        should_fetch = (
            not self.is_loaded()
            or local_age is None
            or local_age >= self._cfg.ipsum_refresh_interval
            or (header_age is not None
                and header_age >= self._cfg.ipsum_refresh_interval)
        )
        if not should_fetch:
            return False

        try:
            async with session.get(
                self._cfg.ipsum_url,
                timeout=aiohttp.ClientTimeout(total=30),
            ) as resp:
                if resp.status != 200:
                    self.log.warning("ipsum.http_error",
                                     url=self._cfg.ipsum_url,
                                     status=resp.status)
                    return False
                text = await resp.text()
        except Exception as exc:  # noqa: BLE001
            self.log.warning("ipsum.fetch_error",
                             url=self._cfg.ipsum_url,
                             error=str(exc),
                             error_type=type(exc).__name__)
            return False

        mapping, header_ts = self.parse(text)
        if not mapping:
            self.log.warning("ipsum.empty_or_unparsable",
                             url=self._cfg.ipsum_url,
                             bytes=len(text))
            return False

        self.install(mapping, header_ts)
        header_age_after = (time.time() - header_ts) if header_ts else None
        self.log.info("ipsum.loaded",
                      entries=len(mapping),
                      bytes=len(text),
                      header_age_s=round(header_age_after, 1)
                      if header_age_after is not None else None)
        return True

    def reputation_score(self, ip: str) -> float:
        n = self.listings_for(ip)
        clean = self._cfg.reputation_clean_score
        if n <= 0:
            return clean
        soft = max(0.0, self._cfg.reputation_listings_soft)
        hard = max(soft + 0.001, self._cfg.reputation_listings_max)

        if n <= soft:
            t = n / soft
            return clean * (1.0 - 0.40 * t)

        t = (n - soft) / (hard - soft)
        t = min(1.0, t)
        return clean * 0.60 * (1.0 - t)


# =============================================================================
# Scoring economy
# =============================================================================

_SCORE_CFG: Optional[Config] = None
_SCORE_IPSUM: Optional[IpsumFeed] = None


def _bind_score_context(cfg: Config, ipsum: IpsumFeed) -> None:
    global _SCORE_CFG, _SCORE_IPSUM
    _SCORE_CFG = cfg
    _SCORE_IPSUM = ipsum


@dataclass
class Proxy:
    address: str
    response_time: float
    first_seen: float
    last_seen: float
    credit: float = 0.0
    success_streak: int = 0
    failure_streak: int = 0
    total_successes: int = 0
    total_failures: int = 0
    reserved_by: Optional[str] = None
    reserved_at: Optional[float] = None

    @property
    def ip(self) -> str:
        return self.address.partition(":")[0]

    @property
    def age(self) -> float:
        return max(0.0, time.time() - self.first_seen)

    @property
    def total_observations(self) -> int:
        return self.total_successes + self.total_failures

    @property
    def success_ratio(self) -> float:
        n = self.total_observations
        if n == 0:
            return 0.0
        return self.total_successes / n

    @property
    def failure_ratio(self) -> float:
        n = self.total_observations
        if n == 0:
            return 0.0
        return self.total_failures / n

    @property
    def longevity_factor(self) -> float:
        cfg = _SCORE_CFG
        streak_target = cfg.longevity_streak_target if cfg else 10
        age_cap = cfg.longevity_age_cap_s if cfg else 3600.0
        streak = min(1.0, self.success_streak / max(1, streak_target))
        age = min(1.0, self.age / max(1.0, age_cap))
        return 0.6 * streak + 0.4 * age

    @property
    def reputation_score(self) -> float:
        ipsum = _SCORE_IPSUM
        if ipsum is None:
            return 1.0
        return ipsum.reputation_score(self.ip)

    @property
    def score(self) -> float:
        cfg = _SCORE_CFG
        if cfg is None:
            return 0.0

        ref = max(0.01, cfg.score_speed_reference_ms / 1000.0)
        speed_score = 1.0 / (1.0 + self.response_time / ref)

        streak_target = max(1, cfg.longevity_streak_target)
        streak_score = min(1.0, self.success_streak / streak_target)

        credit_range = max(1.0, cfg.credit_max)
        credit_score = max(0.0, min(1.0, self.credit / credit_range))

        age_cap = max(1.0, cfg.longevity_age_cap_s)
        age_score = min(1.0, self.age / age_cap)

        reputation = self.reputation_score

        return (
            cfg.score_weight_speed * speed_score
            + cfg.score_weight_success * self.success_ratio
            + cfg.score_weight_streak * streak_score
            + cfg.score_weight_credit * credit_score
            + cfg.score_weight_age * age_score
            + cfg.score_weight_reputation * reputation
            - cfg.score_weight_failure * self.failure_ratio
        )

    def __lt__(self, other: "Proxy") -> bool:
        return self.score > other.score


class ProxyPool:
    def __init__(self, cfg: Config, ipsum: IpsumFeed) -> None:
        _bind_score_context(cfg, ipsum)

        self._cfg = cfg
        self._lock = threading.RLock()
        self._proxies: dict[str, Proxy] = {}
        self.need_more = threading.Event()
        self.need_more.set()
        self.target: int = cfg.target_pool_size

    def update_target(self, active_clients: int) -> int:
        raw = max(0, active_clients) * self._cfg.pool_headroom
        target = int(max(self._cfg.target_pool_size, raw))
        with self._lock:
            self.target = target
            if len(self._proxies) >= target:
                self.need_more.clear()
            else:
                self.need_more.set()
        return target

    def size(self) -> int:
        with self._lock:
            return len(self._proxies)

    def reserved_count(self) -> int:
        with self._lock:
            return sum(1 for p in self._proxies.values() if p.reserved_by)

    def is_saturated(self) -> bool:
        with self._lock:
            return len(self._proxies) >= self.target

    def snapshot(self, include_reserved: bool = True) -> list[Proxy]:
        with self._lock:
            items = list(self._proxies.values())
        if not include_reserved:
            items = [p for p in items if not p.reserved_by]
        return sorted(items)

    def best(self, exclude: Optional[set[str]] = None) -> Optional[Proxy]:
        exclude = exclude or set()
        with self._lock:
            candidates = [p for p in self._proxies.values()
                          if p.address not in exclude and not p.reserved_by]
        if not candidates:
            return None
        return max(candidates, key=lambda p: p.score)

    def contains(self, address: str) -> bool:
        with self._lock:
            return address in self._proxies

    def get(self, address: str) -> Optional[Proxy]:
        with self._lock:
            return self._proxies.get(address)

    def top_n(self, n: int) -> list[Proxy]:
        with self._lock:
            items = sorted(self._proxies.values())
        return items[:n]

    def working_proxies(self, window_s: float) -> list[Proxy]:
        """
        Proxies whose last successful validation is within `window_s` and
        whose credit is non-negative.
        """
        now = time.time()
        with self._lock:
            items = list(self._proxies.values())
        return [p for p in items
                if (now - p.last_seen) <= window_s and p.credit >= 0.0]

    def dead_proxies(self, window_s: float) -> list[Proxy]:
        """
        Proxies whose last success is older than `window_s` or whose credit
        has gone negative.
        """
        now = time.time()
        with self._lock:
            items = list(self._proxies.values())
        return [p for p in items
                if (now - p.last_seen) > window_s or p.credit < 0.0]

    def reserve(self, address: str, client_name: str) -> bool:
        with self._lock:
            p = self._proxies.get(address)
            if p is None:
                return False
            if p.reserved_by is not None and p.reserved_by != client_name:
                return False
            p.reserved_by = client_name
            p.reserved_at = time.time()
            return True

    def release(self, address: str, client_name: str) -> None:
        with self._lock:
            p = self._proxies.get(address)
            if p is None:
                return
            if p.reserved_by == client_name:
                p.reserved_by = None
                p.reserved_at = None

    def reward(self, address: str, response_time: float) -> bool:
        now = time.time()
        with self._lock:
            p = self._proxies.get(address)
            if p is None:
                p = Proxy(
                    address=address,
                    response_time=response_time,
                    first_seen=now,
                    last_seen=now,
                    credit=self._cfg.credit_reward,
                    success_streak=1,
                    total_successes=1,
                )
                self._proxies[address] = p
                grew = True
            else:
                alpha = 0.4
                p.response_time = (
                    (1 - alpha) * p.response_time + alpha * response_time
                )
                p.last_seen = now
                p.credit = min(self._cfg.credit_max,
                               p.credit + self._cfg.credit_reward)
                p.success_streak += 1
                p.failure_streak = 0
                p.total_successes += 1
                grew = False

            if len(self._proxies) >= self.target:
                self.need_more.clear()
            return grew

    def penalize(self, address: str, reason: str = "validation") -> None:
        with self._lock:
            p = self._proxies.get(address)
            if p is None:
                return

            discount = 1.0 - 0.8 * p.longevity_factor
            cost = self._cfg.credit_penalty_base * discount
            p.credit = max(self._cfg.credit_min, p.credit - cost)
            p.failure_streak += 1
            p.success_streak = 0
            p.total_failures += 1

    def mark_dead(self, address: str) -> None:
        self.penalize(address, reason="client_failure")


# =============================================================================
# Client affinity tracker
# =============================================================================

class AffinityTracker:
    def __init__(self, cfg: Config, pool: ProxyPool) -> None:
        self._cfg = cfg
        self._pool = pool
        self._lock = threading.RLock()
        self._bindings: dict[str, tuple[str, float]] = {}

    @staticmethod
    def _mask_ip(peer_ip: str) -> str:
        try:
            ip = ipaddress.ip_address(peer_ip)
        except ValueError:
            return peer_ip
        if isinstance(ip, ipaddress.IPv4Address):
            return str(ipaddress.ip_network(f"{ip}/24", strict=False).network_address)
        return str(ipaddress.ip_network(f"{ip}/48", strict=False).network_address)

    @classmethod
    def fingerprint_from_request(cls, peer_ip: str,
                                 headers: dict[str, str]) -> str:
        lower = {k.lower(): v for k, v in headers.items()}
        masked_ip = cls._mask_ip(peer_ip)
        ua = lower.get("user-agent", "")
        al = lower.get("accept-language", "")
        ae = lower.get("accept-encoding", "")
        sec = lower.get("sec-ch-ua", "") or lower.get("sec-ch-ua-platform", "")
        material = "\x1f".join([masked_ip, ua, al, ae, sec])
        return hashlib.sha256(material.encode("utf-8")).hexdigest()

    @staticmethod
    def name_for(fingerprint: str) -> str:
        return _client_name_from_fingerprint(fingerprint)

    def bind(self, fingerprint: str,
             preferred: Optional[str] = None,
             exclude: Optional[set[str]] = None,
             ) -> Optional[str]:
        exclude = exclude or set()
        now = time.time()

        if not self._cfg.affinity_enabled:
            if preferred and self._pool.contains(preferred):
                return preferred
            best = self._pool.best(exclude=exclude)
            return best.address if best else None

        with self._lock:
            entry = self._bindings.get(fingerprint)
            if entry is not None:
                bound_addr, bound_at = entry
                if (now - bound_at) <= self._cfg.affinity_ttl \
                        and self._pool.contains(bound_addr):
                    self._bindings[fingerprint] = (bound_addr, now)
                    return bound_addr
                self._bindings.pop(fingerprint, None)

            if preferred and self._pool.contains(preferred):
                chosen = preferred
            else:
                best = self._pool.best(exclude=exclude)
                if best is None:
                    return None
                chosen = best.address

            if len(self._bindings) >= self._cfg.affinity_max_entries:
                oldest = min(self._bindings.items(), key=lambda kv: kv[1][1])
                self._bindings.pop(oldest[0], None)

            self._bindings[fingerprint] = (chosen, now)
            return chosen

    def release(self, fingerprint: str) -> None:
        with self._lock:
            self._bindings.pop(fingerprint, None)

    def size(self) -> int:
        with self._lock:
            return len(self._bindings)

    def prune(self) -> int:
        now = time.time()
        removed = 0
        with self._lock:
            for fp, (_, bound_at) in list(self._bindings.items()):
                if (now - bound_at) > self._cfg.affinity_ttl:
                    self._bindings.pop(fp, None)
                    removed += 1
        return removed


# =============================================================================
# Client wait queue
# =============================================================================

class ProxyWaitQueue:
    class QueueFull(Exception):
        pass

    def __init__(self, cfg: Config, metrics: Metrics) -> None:
        self._cfg = cfg
        self._metrics = metrics
        self._lock = threading.RLock()
        self._waiters: deque = deque()
        self.log = get_logger("proxydrum.waitqueue")

    def depth(self) -> int:
        with self._lock:
            return len(self._waiters)

    def request(self, client_name: str,
                fingerprint: str,
                loop: asyncio.AbstractEventLoop) -> asyncio.Future:
        with self._lock:
            if len(self._waiters) >= self._cfg.wait_queue_max:
                self._metrics.client_rejected()
                raise ProxyWaitQueue.QueueFull(
                    f"wait queue full ({len(self._waiters)})")

            fut = loop.create_future()
            self._waiters.append((client_name, fut, time.time(), fingerprint))
            self._metrics.client_waiting_enter()
            self.log.info("wait.enqueued",
                          client=client_name,
                          depth=len(self._waiters))
            return fut

    def cancel(self, fut: asyncio.Future) -> None:
        with self._lock:
            for i, (name, f, _, _) in enumerate(self._waiters):
                if f is fut:
                    del self._waiters[i]
                    self._metrics.client_waiting_exit()
                    return

    def grant_next(self, proxy_addr: str) -> bool:
        with self._lock:
            while self._waiters:
                name, fut, enqueued, fingerprint = self._waiters.popleft()
                waited = time.time() - enqueued
                if fut.done():
                    self._metrics.client_waiting_exit()
                    continue
                self._metrics.client_waiting_exit()
                self._metrics.record_wait(waited)
                try:
                    fut.get_loop().call_soon_threadsafe(
                        fut.set_result, (proxy_addr, fingerprint))
                except Exception:  # noqa: BLE001
                    self.log.exception("wait.grant_failed",
                                       client=name,
                                       proxy=proxy_addr)
                    continue
                self.log.info("wait.granted",
                              client=name,
                              proxy=proxy_addr,
                              waited_s=round(waited, 2),
                              depth=len(self._waiters))
                return True
            return False

    def expire_waiters(self, timeout: float) -> int:
        now = time.time()
        expired = 0
        with self._lock:
            surviving = deque()
            for name, fut, enqueued, fingerprint in self._waiters:
                if now - enqueued > timeout:
                    if not fut.done():
                        try:
                            fut.get_loop().call_soon_threadsafe(
                                fut.set_exception,
                                asyncio.TimeoutError("proxy wait timeout"))
                        except Exception:  # noqa: BLE001
                            pass
                    self._metrics.client_waiting_exit()
                    self._metrics.client_timed_out()
                    expired += 1
                    self.log.warning("wait.timeout",
                                     client=name,
                                     waited_s=round(now - enqueued, 2))
                else:
                    surviving.append((name, fut, enqueued, fingerprint))
            self._waiters = surviving
        return expired


# =============================================================================
# Proxy list parsing and testing helpers
# =============================================================================

_PROXY_RE = re.compile(r"(?:\d{1,3}\.){3}\d{1,3}:\d{1,5}")


def parse_proxy_list(text: str) -> list[str]:
    return [m for m in _PROXY_RE.findall(text)]


async def test_via_proxy(
    session: aiohttp.ClientSession,
    proxy: str,
    test_url: str,
    timeout: float,
    retries: int = 0,
) -> Optional[float]:
    proxy_url = f"http://{proxy}"
    attempts = 1 + max(0, retries)

    for attempt in range(attempts):
        start = time.monotonic()
        try:
            async with session.get(
                test_url,
                proxy=proxy_url,
                timeout=aiohttp.ClientTimeout(total=timeout),
                allow_redirects=False,
                headers={"User-Agent": "proxyDrum/0.1"},
            ) as resp:
                body = await resp.text()
                elapsed = time.monotonic() - start
                if 200 <= resp.status < 400 and body:
                    if "ipify" in test_url and "." not in body:
                        continue
                    return elapsed
        except Exception:  # noqa: BLE001
            pass
        if attempt + 1 < attempts:
            await asyncio.sleep(0.5)
    return None


# =============================================================================
# Source scheduler (self-tuning)
# =============================================================================

class SourceScheduler:
    @dataclass
    class State:
        url: str
        last_hash: Optional[str] = None
        last_fetched_at: float = 0.0
        change_ts: deque = field(default_factory=deque)
        no_change_streak: int = 0
        interval: float = 0.0

    def __init__(self, cfg: Config) -> None:
        self._cfg = cfg
        self._lock = threading.RLock()
        self._states: dict[str, SourceScheduler.State] = {
            url: SourceScheduler.State(
                url=url,
                interval=cfg.source_initial_interval,
            )
            for url in cfg.proxy_sources
        }
        self.log = get_logger("proxydrum.scheduler")

    def urls(self) -> list[str]:
        with self._lock:
            return list(self._states.keys())

    def state(self, url: str) -> "SourceScheduler.State":
        with self._lock:
            return self._states[url]

    def due_urls(self) -> list[str]:
        with self._lock:
            now = time.time()
            return [u for u, s in self._states.items()
                    if (now - s.last_fetched_at) >= s.interval]

    def _median(self, values: list[float]) -> float:
        if not values:
            return 0.0
        s = sorted(values)
        n = len(s)
        if n % 2 == 1:
            return s[n // 2]
        return 0.5 * (s[n // 2 - 1] + s[n // 2])

    def _recompute_interval_locked(self, s: "SourceScheduler.State") -> None:
        cfg = self._cfg
        if len(s.change_ts) >= 2:
            gaps = [s.change_ts[i] - s.change_ts[i - 1]
                    for i in range(1, len(s.change_ts))]
            base = self._median(gaps) * cfg.source_safety_margin
        else:
            base = cfg.source_initial_interval * (2 ** s.no_change_streak)
        s.interval = float(max(cfg.source_min_interval,
                               min(cfg.source_max_interval, base)))

    def observe(self, url: str, content: str) -> bool:
        h = hashlib.sha256(content.encode("utf-8", errors="replace")).hexdigest()
        now = time.time()
        with self._lock:
            s = self._states.get(url)
            if s is None:
                return False
            changed = (s.last_hash is not None and s.last_hash != h)
            s.last_hash = h
            s.last_fetched_at = now

            if changed:
                s.change_ts.append(now)
                while len(s.change_ts) > self._cfg.source_history_size + 1:
                    s.change_ts.popleft()
                s.no_change_streak = 0
            else:
                if s.change_ts:
                    s.no_change_streak += 1

            self._recompute_interval_locked(s)
            return changed

    def note_attempt(self, url: str) -> None:
        with self._lock:
            s = self._states.get(url)
            if s is None:
                return
            s.last_fetched_at = time.time()

    def schedule_report(self) -> list[dict]:
        with self._lock:
            now = time.time()
            out = []
            for u, s in self._states.items():
                out.append({
                    "name": u.rsplit("/", 1)[-1],
                    "interval_s": round(s.interval, 1),
                    "age_s": round(now - s.last_fetched_at, 1)
                    if s.last_fetched_at else None,
                    "changes_seen": len(s.change_ts),
                })
            return out


# =============================================================================
# Thread 1 - Fetcher
# =============================================================================

class Fetcher(threading.Thread):
    def __init__(self, cfg: Config, pool: ProxyPool,
                 scheduler: SourceScheduler,
                 candidate_queue: "asyncio.Queue[str]",
                 loop: asyncio.AbstractEventLoop,
                 metrics: Metrics) -> None:
        super().__init__(name="fetcher", daemon=True)
        self.cfg = cfg
        self.pool = pool
        self.scheduler = scheduler
        self.queue = candidate_queue
        self.loop = loop
        self.metrics = metrics
        self._stop = threading.Event()
        self.log = get_logger("proxydrum.fetcher")

    def stop(self) -> None:
        self._stop.set()

    async def _download_one(self, session: aiohttp.ClientSession,
                            url: str) -> Optional[str]:
        try:
            async with session.get(
                url,
                timeout=aiohttp.ClientTimeout(total=20),
            ) as resp:
                if resp.status != 200:
                    self.log.warning("source.http_error",
                                     url=url, status=resp.status)
                    self.scheduler.note_attempt(url)
                    return None
                text = await resp.text()
        except Exception as exc:  # noqa: BLE001
            self.log.warning("source.error",
                             url=url, error=str(exc),
                             error_type=type(exc).__name__)
            self.scheduler.note_attempt(url)
            return None

        changed = self.scheduler.observe(url, text)
        self.log.info("source.downloaded",
                      url=url, bytes=len(text), changed=changed,
                      learned_interval_s=round(
                          self.scheduler.state(url).interval, 1))
        return text

    async def _fetch_due(self) -> None:
        busy = self.pool.size() < self.pool.target
        due = self.scheduler.due_urls()
        if not due:
            return

        self.log.debug("fetch.dispatching", due=len(due), busy=busy)

        async with aiohttp.ClientSession(
                headers={"User-Agent": "proxyDrum/0.1"}) as session:
            tasks = [self._download_one(session, url) for url in due]
            results = await asyncio.gather(*tasks)

        candidates: list[str] = []
        for text in results:
            if text is None:
                continue
            for addr in parse_proxy_list(text):
                candidates.append(addr)

        if not candidates:
            self.log.info("fetch.no_candidates",
                          queue_size=self.queue.qsize())
            self.metrics.mark_fetch()
            return

        seen_this_round: set[str] = set()
        fresh: list[str] = []
        for addr in candidates:
            if addr in seen_this_round:
                continue
            seen_this_round.add(addr)
            fresh.append(addr)

        random.shuffle(fresh)
        self.log.info("fetch.queued",
                      count=len(fresh),
                      queue_size=self.queue.qsize() + len(fresh))
        for addr in fresh:
            await self.queue.put(addr)

        self.metrics.mark_fetch()

    def run(self) -> None:
        self.log.info("thread.started",
                      sources=self.scheduler.urls(),
                      tick=self.cfg.scheduler_tick,
                      min_interval=self.cfg.source_min_interval,
                      max_interval=self.cfg.source_max_interval)
        asyncio.set_event_loop(self.loop)
        while not self._stop.is_set():
            try:
                fut = asyncio.run_coroutine_threadsafe(
                    self._fetch_due(), self.loop)
                fut.result(timeout=120)
            except Exception as exc:  # noqa: BLE001
                self.log.exception("fetcher.unhandled_error",
                                   error=str(exc),
                                   error_type=type(exc).__name__)

            self._stop.wait(self.cfg.scheduler_tick)

        self.log.info("thread.stopped")


# =============================================================================
# Thread 2 - Validator
# =============================================================================

class Validator(threading.Thread):
    def __init__(self, cfg: Config, pool: ProxyPool,
                 candidate_queue: "asyncio.Queue[str]",
                 loop: asyncio.AbstractEventLoop,
                 metrics: Metrics,
                 ipsum: IpsumFeed) -> None:
        super().__init__(name="validator", daemon=True)
        self.cfg = cfg
        self.pool = pool
        self.queue = candidate_queue
        self.loop = loop
        self.metrics = metrics
        self.ipsum = ipsum
        self._stop = threading.Event()
        self.log = get_logger("proxydrum.validator")
        self._cycle = 0
        self._ipsum_last_fetch_at = 0.0

    def stop(self) -> None:
        self._stop.set()

    def _timeout_for(self, proxy: Optional[Proxy]) -> tuple[float, int]:
        if proxy is None:
            return self.cfg.validate_timeout_base, 0
        lf = proxy.longevity_factor
        timeout = self.cfg.validate_timeout_base * (
            1.0 + (self.cfg.longevity_timeout_max - 1.0) * lf)
        retries = int(round(self.cfg.longevity_retries_max * lf))
        return timeout, retries

    async def _refresh_ipsum_if_needed(self,
                                       session: aiohttp.ClientSession) -> None:
        if not self.cfg.ipsum_enabled:
            return
        updated = await self.ipsum.maybe_refresh(
            session, self._ipsum_last_fetch_at)
        if updated:
            self._ipsum_last_fetch_at = time.time()
            self.metrics.mark_ipsum(
                entries=self.ipsum.size(),
                header_ts=(time.time() - (self.ipsum.header_age() or 0.0)),
            )

    async def _validate_candidates(self,
                                   session: aiohttp.ClientSession) -> None:
        batch: list[tuple[str, Optional[Proxy]]] = []
        limit = self.cfg.validate_batch_size
        for _ in range(limit):
            try:
                addr = self.queue.get_nowait()
            except asyncio.QueueEmpty:
                break
            batch.append((addr, self.pool.get(addr)))

        if not batch:
            return

        self.log.debug("validate.candidates.begin",
                       batch_size=len(batch),
                       queue_size=self.queue.qsize(),
                       pool_size=self.pool.size())

        async def probe(addr: str, existing: Optional[Proxy]):
            timeout, retries = self._timeout_for(existing)
            elapsed = await test_via_proxy(
                session, addr, self.cfg.test_url, timeout, retries)
            return addr, elapsed

        results = await asyncio.gather(
            *(probe(addr, existing) for addr, existing in batch),
            return_exceptions=False)

        new_added = 0
        rewards = 0
        penalties = 0
        for addr, elapsed in results:
            if elapsed is None:
                if self.pool.contains(addr):
                    self.pool.penalize(addr)
                    self.metrics.proxy_penalized()
                    penalties += 1
            else:
                added = self.pool.reward(addr, elapsed)
                self.metrics.proxy_rewarded()
                rewards += 1
                if added:
                    new_added += 1
                    self.metrics.proxy_added()
                    ip = addr.partition(":")[0]
                    rep = self.ipsum.reputation_score(ip)
                    listings = self.ipsum.listings_for(ip)
                    self.log.info("proxy.added",
                                  proxy=addr,
                                  response_time_ms=round(elapsed * 1000, 1),
                                  pool_size=self.pool.size(),
                                  ipsum_listings=listings,
                                  reputation=round(rep, 3))

        self.metrics.mark_validation()
        self.log.debug("validate.candidates.complete",
                       batch_size=len(batch),
                       new=new_added,
                       rewards=rewards,
                       penalties=penalties,
                       pool_size=self.pool.size(),
                       queue_size=self.queue.qsize())

    async def _revalidate_pool(self,
                               session: aiohttp.ClientSession) -> None:
        snapshot = self.pool.snapshot(include_reserved=False)
        reserved = self.pool.reserved_count()
        if not snapshot:
            self.pool.need_more.set()
            self.log.info("validate.pool.empty", reserved=reserved)
            return

        self._cycle += 1
        cycle = self._cycle
        self.log.debug("validate.pool.begin",
                       cycle=cycle,
                       candidates=len(snapshot),
                       reserved_skipped=reserved,
                       pool_size=self.pool.size())

        async def probe(p: Proxy) -> tuple[Proxy, Optional[float]]:
            timeout, retries = self._timeout_for(p)
            elapsed = await test_via_proxy(
                session, p.address, self.cfg.test_url, timeout, retries)
            return p, elapsed

        results = await asyncio.gather(
            *(probe(p) for p in snapshot), return_exceptions=False)

        alive = 0
        penalized = 0
        for p, elapsed in results:
            if elapsed is None:
                self.pool.penalize(p.address)
                self.metrics.proxy_penalized()
                penalized += 1
            else:
                self.pool.reward(p.address, elapsed)
                alive += 1

        self.log.debug("validate.pool.complete",
                       cycle=cycle, alive=alive, penalized=penalized,
                       pool_size=self.pool.size(), reserved=reserved)

        if self.pool.size() < self.pool.target:
            self.pool.need_more.set()
            self.log.info("pool.below_target",
                          pool_size=self.pool.size(),
                          target=self.pool.target,
                          reserved=reserved)

    def _interval_for_pressure(self) -> float:
        if self.pool.size() < self.pool.target:
            return self.cfg.validate_interval_busy
        return self.cfg.validate_interval_idle

    def run(self) -> None:
        self.log.info("thread.started",
                      idle_interval=self.cfg.validate_interval_idle,
                      busy_interval=self.cfg.validate_interval_busy,
                      ipsum_enabled=self.cfg.ipsum_enabled,
                      ipsum_url=self.cfg.ipsum_url,
                      ipsum_refresh_interval=self.cfg.ipsum_refresh_interval)
        asyncio.set_event_loop(self.loop)

        async def _loop() -> None:
            async with aiohttp.ClientSession() as session:
                while not self._stop.is_set():
                    try:
                        await self._refresh_ipsum_if_needed(session)
                        await self._validate_candidates(session)
                        await self._revalidate_pool(session)
                    except Exception as exc:  # noqa: BLE001
                        self.log.exception("validator.unhandled_error",
                                           error=str(exc),
                                           error_type=type(exc).__name__)
                    await asyncio.sleep(self._interval_for_pressure())

        fut = asyncio.run_coroutine_threadsafe(_loop(), self.loop)
        while not self._stop.is_set():
            if fut.done():
                break
            self._stop.wait(1.0)
        if not fut.done():
            fut.cancel()

        self.log.info("thread.stopped")


# =============================================================================
# Thread 3 - HTTP forward proxy server
# =============================================================================

class HttpProxyServer(threading.Thread):
    def __init__(self, cfg: Config, pool: ProxyPool,
                 affinity: AffinityTracker,
                 wait_queue: ProxyWaitQueue,
                 loop: asyncio.AbstractEventLoop,
                 metrics: Metrics) -> None:
        super().__init__(name="httpproxy", daemon=True)
        self.cfg = cfg
        self.pool = pool
        self.affinity = affinity
        self.wait_queue = wait_queue
        self.loop = loop
        self.metrics = metrics
        self._stop = threading.Event()
        self._server: Optional[asyncio.AbstractServer] = None
        self.log = get_logger("proxydrum.httpproxy")

    def stop(self) -> None:
        self._stop.set()
        if self._server is not None:
            self._server.close()

    @staticmethod
    async def _read_headers(reader: asyncio.StreamReader,
                            limit: int = 65536) -> bytes:
        buf = bytearray()
        while b"\r\n\r\n" not in buf and b"\n\n" not in buf:
            chunk = await reader.read(4096)
            if not chunk:
                break
            buf.extend(chunk)
            if len(buf) > limit:
                raise ValueError("request headers too large")
        return bytes(buf)

    @staticmethod
    def _parse_headers(raw: bytes) -> dict[str, str]:
        headers: dict[str, str] = {}
        lines = raw.split(b"\r\n")
        if len(lines) == 1:
            lines = raw.split(b"\n")
        for line in lines[1:]:
            if not line.strip():
                continue
            if b":" not in line:
                continue
            key, _, value = line.partition(b":")
            headers[key.decode("latin-1").strip().lower()] = \
                value.decode("latin-1").strip()
        return headers

    async def _open_upstream(self, proxy: str,
                             target_host: str, target_port: int
                             ) -> Optional[tuple[asyncio.StreamReader,
                                                asyncio.StreamWriter]]:
        host, _, port_s = proxy.partition(":")
        try:
            up_reader, up_writer = await asyncio.wait_for(
                asyncio.open_connection(host, int(port_s)),
                timeout=self.cfg.connect_timeout,
            )
        except Exception as exc:  # noqa: BLE001
            self.log.debug("upstream.connect_failed",
                           proxy=proxy, error=str(exc),
                           error_type=type(exc).__name__)
            return None

        connect_req = (
            f"CONNECT {target_host}:{target_port} HTTP/1.1\r\n"
            f"Host: {target_host}:{target_port}\r\n"
            f"Proxy-Connection: keep-alive\r\n\r\n"
        ).encode()
        try:
            up_writer.write(connect_req)
            await up_writer.drain()
            status_line = await asyncio.wait_for(
                up_reader.readline(), timeout=self.cfg.connect_timeout)
            if b" 200 " not in status_line:
                raise RuntimeError(f"upstream refused: {status_line!r}")
            while True:
                line = await asyncio.wait_for(
                    up_reader.readline(), timeout=self.cfg.connect_timeout)
                if line in (b"\r\n", b"\n", b""):
                    break
        except Exception as exc:  # noqa: BLE001
            self.log.debug("upstream.connect_handshake_failed",
                           proxy=proxy,
                           target=f"{target_host}:{target_port}",
                           error=str(exc),
                           error_type=type(exc).__name__)
            try:
                up_writer.close()
            except Exception:
                pass
            return None

        return up_reader, up_writer

    async def _acquire_proxy(self, client_name: str,
                             fingerprint: str,
                             exclude: set[str],
                             ) -> Optional[str]:
        bound = self.affinity.bind(fingerprint, exclude=exclude)
        if bound and bound not in exclude:
            if self.pool.reserve(bound, client_name):
                return bound

        best = self.pool.best(exclude=exclude)
        if best is not None and self.pool.reserve(best.address, client_name):
            return best.address

        self.log.debug("wait.enter",
                       client=client_name,
                       pool_size=self.pool.size(),
                       reserved=self.pool.reserved_count(),
                       depth=self.wait_queue.depth())
        started = time.monotonic()
        try:
            fut = self.wait_queue.request(client_name, fingerprint,
                                          asyncio.get_running_loop())
        except ProxyWaitQueue.QueueFull:
            self.log.warning("wait.rejected_full",
                             client=client_name,
                             depth=self.wait_queue.depth())
            return None

        try:
            granted, _ = await asyncio.wait_for(
                fut, timeout=self.cfg.wait_timeout)
        except asyncio.TimeoutError:
            self.wait_queue.cancel(fut)
            self.metrics.client_timed_out()
            self.log.warning("wait.timed_out",
                             client=client_name,
                             waited_s=round(time.monotonic() - started, 2))
            return None
        except Exception as exc:  # noqa: BLE001
            self.wait_queue.cancel(fut)
            self.log.warning("wait.error",
                             client=client_name,
                             error=str(exc))
            return None

        if not self.pool.reserve(granted, client_name):
            self.affinity.release(fingerprint)
            return None
        waited = time.monotonic() - started
        self.log.info("wait.acquired",
                      client=client_name,
                      proxy=granted,
                      waited_s=round(waited, 2))
        return granted

    async def _handle_client(self, reader: asyncio.StreamReader,
                             writer: asyncio.StreamWriter) -> None:
        peer = writer.get_extra_info("peername")
        peer_ip = peer[0] if peer else "0.0.0.0"
        client_name = "unknown"
        try:
            header_block = await self._read_headers(reader)
            if not header_block:
                return

            head, _, rest = header_block.partition(b"\r\n\r\n")
            if not rest:
                head, _, rest = header_block.partition(b"\n\n")
            lines = head.splitlines()
            if not lines:
                return
            request_line = lines[0].decode("latin-1", errors="replace")
            parts = request_line.split()
            if len(parts) < 3:
                await self._reply(writer, 400, b"Bad Request")
                return
            method, target, _version = parts[0], parts[1], parts[2]

            headers = self._parse_headers(head)
            fingerprint = AffinityTracker.fingerprint_from_request(
                peer_ip, headers)
            client_name = AffinityTracker.name_for(fingerprint)

            self.metrics.client_connected()

            body_prefix = rest

            if method.upper() == "CONNECT":
                await self._handle_connect(peer, client_name, fingerprint,
                                           target, reader, writer)
                return

            await self._handle_http(method, target, head, body_prefix,
                                    client_name, fingerprint, reader, writer)

        except (asyncio.IncompleteReadError, ConnectionError):
            pass
        except Exception as exc:  # noqa: BLE001
            self.log.exception("client.unhandled_error",
                               client=client_name,
                               peer=str(peer),
                               error=str(exc),
                               error_type=type(exc).__name__)
        finally:
            self.metrics.client_disconnected()
            try:
                writer.close()
                await writer.wait_closed()
            except Exception:
                pass

    async def _handle_connect(self, peer, client_name: str,
                              fingerprint: str, target: str,
                              reader: asyncio.StreamReader,
                              writer: asyncio.StreamWriter) -> None:
        if ":" in target:
            target_host, _, port_s = target.rpartition(":")
            try:
                target_port = int(port_s)
            except ValueError:
                target_port = 443
        else:
            target_host, target_port = target, 443

        self.log.info("connect.begin",
                      client=client_name,
                      peer=str(peer),
                      target_host=target_host,
                      target_port=target_port,
                      pool_size=self.pool.size(),
                      reserved=self.pool.reserved_count())

        attempted: set[str] = set()
        while not self._stop.is_set():
            proxy_addr = await self._acquire_proxy(
                client_name, fingerprint, attempted)
            if proxy_addr is None:
                self.metrics.request_failed()
                await self._reply(
                    writer, 503, b"Service Unavailable",
                    retry_after=self.cfg.retry_after_hint)
                return
            attempted.add(proxy_addr)

            tunnel = await self._open_upstream(
                proxy_addr, target_host, target_port)
            if tunnel is None:
                self.pool.mark_dead(proxy_addr)
                self.metrics.proxy_penalized()
                self.affinity.release(fingerprint)
                self.pool.release(proxy_addr, client_name)
                self.wait_queue.grant_next(proxy_addr)
                continue

            up_reader, up_writer = tunnel

            writer.write(b"HTTP/1.1 200 Connection Established\r\n"
                         b"Proxy-Agent: proxyDrum\r\n\r\n")
            await writer.drain()

            self.log.info("connect.established",
                          client=client_name,
                          peer=str(peer),
                          proxy=proxy_addr,
                          target_host=target_host,
                          target_port=target_port)

            try:
                await asyncio.gather(
                    self._pump(reader, up_writer, "up"),
                    self._pump(up_reader, writer, "down"),
                    return_exceptions=True,
                )
            finally:
                self.pool.release(proxy_addr, client_name)
                self.metrics.request_served()
                self.wait_queue.grant_next(proxy_addr)

            self.log.info("connect.closed",
                          client=client_name,
                          peer=str(peer),
                          proxy=proxy_addr,
                          target_host=target_host)
            return

        self.log.warning("connect.no_upstream",
                         client=client_name,
                         peer=str(peer),
                         target_host=target_host)
        self.metrics.request_failed()
        self.pool.need_more.set()
        await self._reply(writer, 502, b"Bad Gateway")

    async def _handle_http(self, method: str, target: str,
                           head: bytes, body_prefix: bytes,
                           client_name: str, fingerprint: str,
                           reader: asyncio.StreamReader,
                           writer: asyncio.StreamWriter) -> None:
        if target.startswith("http://"):
            rest = target[len("http://"):]
            host_part, _, path = rest.partition("/")
            path = "/" + path
        else:
            host_part, path = target, "/"

        host, _, port_s = host_part.partition(":")
        port = int(port_s) if port_s else 80

        new_lines = [f"{method} {path} HTTP/1.1".encode("latin-1")]
        for line in head.splitlines()[1:]:
            lower = line.lower()
            if lower.startswith(b"proxy-connection:"):
                continue
            if lower.startswith(b"proxy-authorization:"):
                continue
            if lower.startswith(b"connection:"):
                continue
            new_lines.append(line)
        joined = b"\r\n".join(new_lines)
        if b"\r\nhost:" not in joined.lower():
            joined += f"\r\nHost: {host_part}".encode("latin-1")
        joined += b"\r\nConnection: keep-alive\r\n"
        origin_request = joined + b"\r\n\r\n" + body_prefix

        attempted: set[str] = set()
        while not self._stop.is_set():
            proxy_addr = await self._acquire_proxy(
                client_name, fingerprint, attempted)
            if proxy_addr is None:
                self.metrics.request_failed()
                await self._reply(
                    writer, 503, b"Service Unavailable",
                    retry_after=self.cfg.retry_after_hint)
                return
            attempted.add(proxy_addr)

            tunnel = await self._open_upstream(proxy_addr, host, port)
            if tunnel is None:
                self.pool.mark_dead(proxy_addr)
                self.metrics.proxy_penalized()
                self.affinity.release(fingerprint)
                self.pool.release(proxy_addr, client_name)
                self.wait_queue.grant_next(proxy_addr)
                continue

            up_reader, up_writer = tunnel
            try:
                up_writer.write(origin_request)
                await up_writer.drain()
                await self._pump(up_reader, writer, "down")
                self.metrics.request_served()
                return
            except Exception as exc:  # noqa: BLE001
                self.log.warning("http.forward_failed",
                                 client=client_name,
                                 proxy=proxy_addr, host=host,
                                 error=str(exc),
                                 error_type=type(exc).__name__)
                self.pool.mark_dead(proxy_addr)
                self.metrics.proxy_penalized()
                self.affinity.release(fingerprint)
                continue
            finally:
                self.pool.release(proxy_addr, client_name)
                self.wait_queue.grant_next(proxy_addr)
                try:
                    up_writer.close()
                except Exception:
                    pass

        self.log.warning("http.no_upstream",
                         client=client_name, host=host)
        self.metrics.request_failed()
        self.pool.need_more.set()
        await self._reply(writer, 502, b"Bad Gateway")

    async def _pump(self, reader: asyncio.StreamReader,
                    writer: asyncio.StreamWriter,
                    direction: str) -> None:
        try:
            while True:
                chunk = await reader.read(65536)
                if not chunk:
                    break
                writer.write(chunk)
                await writer.drain()
                if direction == "up":
                    self.metrics.add_bytes_up(len(chunk))
                else:
                    self.metrics.add_bytes_down(len(chunk))
        except Exception:
            pass
        finally:
            try:
                writer.close()
            except Exception:
                pass

    @staticmethod
    async def _reply(writer: asyncio.StreamWriter, status: int,
                     reason: bytes, retry_after: Optional[int] = None) -> None:
        try:
            headers = (b"HTTP/1.1 " + str(status).encode() + b" " + reason +
                       b"\r\nContent-Length: 0\r\nConnection: close\r\n")
            if retry_after is not None:
                headers += b"Retry-After: " + str(retry_after).encode() + b"\r\n"
            writer.write(headers + b"\r\n")
            await writer.drain()
        except Exception:
            pass

    def run(self) -> None:
        self.log.info("thread.started")
        asyncio.set_event_loop(self.loop)

        async def _serve() -> None:
            self._server = await asyncio.start_server(
                self._handle_client,
                host=self.cfg.listen_host,
                port=self.cfg.listen_port,
            )
            addrs = ", ".join(str(s.getsockname())
                              for s in self._server.sockets or [])
            self.log.info("server.listening", addresses=addrs)

            async def _housekeeping() -> None:
                while not self._stop.is_set():
                    evicted = self.affinity.prune()
                    if evicted:
                        self.log.debug("affinity.pruned",
                                       evicted=evicted,
                                       remaining=self.affinity.size())
                    self.wait_queue.expire_waiters(self.cfg.wait_timeout)
                    await asyncio.sleep(5)

            hk = asyncio.create_task(_housekeeping())
            try:
                async with self._server:
                    await self._server.serve_forever()
            finally:
                hk.cancel()

        fut = asyncio.run_coroutine_threadsafe(_serve(), self.loop)
        while not self._stop.is_set():
            if fut.done():
                break
            self._stop.wait(1.0)
        if not fut.done():
            fut.cancel()
        self.log.info("thread.stopped")


# =============================================================================
# Thread 4 - Reporter (compact status via structlog)
# =============================================================================

class Reporter(threading.Thread):
    """
    Emits a compact status event on a fixed interval. Only aggregate,
    human-relevant information is reported: uptime, pool composition,
    proxy health, and client activity. The scoring system is intentionally
    not surfaced.
    """

    def __init__(self, cfg: Config, pool: ProxyPool, metrics: Metrics,
                 affinity: AffinityTracker, wait_queue: ProxyWaitQueue,
                 scheduler: SourceScheduler, ipsum: IpsumFeed) -> None:
        super().__init__(name="reporter", daemon=True)
        self.cfg = cfg
        self.pool = pool
        self.metrics = metrics
        self.affinity = affinity
        self.wait_queue = wait_queue
        self.scheduler = scheduler
        self.ipsum = ipsum
        self._stop = threading.Event()
        self.log = get_logger("proxydrum.reporter")

    def stop(self) -> None:
        self._stop.set()

    @staticmethod
    def _fmt_bytes(n: int) -> str:
        for unit in ("B", "KiB", "MiB", "GiB"):
            if n < 1024:
                return f"{n:.1f}{unit}"
            n /= 1024.0
        return f"{n:.1f}TiB"

    @staticmethod
    def _fmt_duration(seconds: float) -> str:
        seconds = int(seconds)
        d, rem = divmod(seconds, 86400)
        h, rem = divmod(rem, 3600)
        m, s = divmod(rem, 60)
        if d:
            return f"{d}d{h:02d}h{m:02d}m"
        return f"{h:d}h{m:02d}m{s:02d}s"

    def build_status(self) -> dict:
        snap = self.metrics.snapshot()
        pool_size = self.pool.size()
        reserved = self.pool.reserved_count()
        target = self.pool.target
        queue_depth = self.wait_queue.depth()
        affinity_n = self.affinity.size()

        working = self.pool.working_proxies(self.cfg.working_window_s)
        dead = self.pool.dead_proxies(self.cfg.working_window_s)
        mean_lifetime_s = (
            sum(p.age for p in working) / len(working)
            if working else 0.0
        )

        sources = self.scheduler.schedule_report()

        return {
            "uptime_s": snap["uptime_s"],
            "uptime": self._fmt_duration(snap["uptime_s"]),
            "pool_size": pool_size,
            "pool_target": target,
            "pool_reserved": reserved,
            "proxies_working": len(working),
            "proxies_dead": len(dead),
            "proxies_mean_lifetime_s": round(mean_lifetime_s, 1),
            "proxies_mean_lifetime": self._fmt_duration(mean_lifetime_s)
            if mean_lifetime_s else "0h00m00s",
            "clients_connected": snap["clients_connected"],
            "clients_total": snap["clients_total"],
            "clients_waiting": queue_depth,
            "clients_rejected": snap["clients_rejected"],
            "clients_timed_out": snap["clients_timed_out"],
            "requests_served": snap["requests_served"],
            "requests_failed": snap["requests_failed"],
            "avg_wait_s": snap["avg_wait_s"],
            "additions_total": snap["additions_total"],
            "rewards_total": snap["rewards_total"],
            "penalties_total": snap["penalties_total"],
            "bytes_up": self._fmt_bytes(snap["bytes_up"]),
            "bytes_down": self._fmt_bytes(snap["bytes_down"]),
            "affinity_bindings": affinity_n,
            "ipsum_enabled": self.cfg.ipsum_enabled,
            "ipsum_entries": snap["ipsum_entries"],
            "ipsum_age_s": snap["ipsum_age_s"],
            "sources": sources,
        }

    def _emit(self) -> None:
        status = self.build_status()
        self.log.info(
            "report.status",
            uptime=status["uptime"],
            working=status["proxies_working"],
            dead=status["proxies_dead"],
            pool_size=status["pool_size"],
            pool_target=status["pool_target"],
            mean_lifetime=status["proxies_mean_lifetime"],
            clients=status["clients_connected"],
            waiting=status["clients_waiting"],
            served=status["requests_served"],
            failed=status["requests_failed"],
            avg_wait_s=status["avg_wait_s"],
            add=status["additions_total"],
            io=f"{status['bytes_up']}/{status['bytes_down']}",
            aff=status["affinity_bindings"],
        )

    def run(self) -> None:
        self.log.info("thread.started", interval=self.cfg.report_interval)
        while not self._stop.is_set():
            self._stop.wait(self.cfg.report_interval)
            if self._stop.is_set():
                break
            try:
                self._emit()
            except Exception as exc:  # noqa: BLE001
                self.log.exception("reporter.error", error=str(exc))
        self.log.info("thread.stopped")


# =============================================================================
# Thread 5 - Adaptive controller
# =============================================================================

class AdaptiveController(threading.Thread):
    def __init__(self, cfg: Config, pool: ProxyPool, metrics: Metrics) -> None:
        super().__init__(name="adaptive", daemon=True)
        self.cfg = cfg
        self.pool = pool
        self.metrics = metrics
        self._stop = threading.Event()
        self.log = get_logger("proxydrum.adaptive")

    def stop(self) -> None:
        self._stop.set()

    def run(self) -> None:
        self.log.info("thread.started", interval=self.cfg.adaptive_interval)
        last_target = -1
        while not self._stop.is_set():
            self._stop.wait(self.cfg.adaptive_interval)
            if self._stop.is_set():
                break
            try:
                snap = self.metrics.snapshot()
                active = snap["clients_connected"]
                target = self.pool.update_target(active)
                if target != last_target:
                    self.log.info("adaptive.target_changed",
                                  clients=active,
                                  old_target=last_target,
                                  new_target=target,
                                  pool_size=self.pool.size())
                    last_target = target
            except Exception as exc:  # noqa: BLE001
                self.log.exception("adaptive.error", error=str(exc))
        self.log.info("thread.stopped")


# =============================================================================
# Orchestrator
# =============================================================================

class Rotator:
    def __init__(self, cfg: Config) -> None:
        self.cfg = cfg
        self.loop = asyncio.new_event_loop()
        self.metrics = Metrics()
        self.ipsum = IpsumFeed(cfg)
        self.pool = ProxyPool(cfg, self.ipsum)
        self.affinity = AffinityTracker(cfg, self.pool)
        self.wait_queue = ProxyWaitQueue(cfg, self.metrics)
        self.scheduler = SourceScheduler(cfg)
        self.candidate_queue: asyncio.Queue[str] = asyncio.Queue()
        self.log = get_logger("proxydrum")

        self.fetcher = Fetcher(cfg, self.pool, self.scheduler,
                               self.candidate_queue, self.loop, self.metrics)
        self.validator = Validator(cfg, self.pool, self.candidate_queue,
                                   self.loop, self.metrics, self.ipsum)
        self.server = HttpProxyServer(cfg, self.pool, self.affinity,
                                      self.wait_queue, self.loop, self.metrics)
        self.reporter = Reporter(cfg, self.pool, self.metrics,
                                 self.affinity, self.wait_queue,
                                 self.scheduler, self.ipsum)
        self.adaptive = AdaptiveController(cfg, self.pool, self.metrics)

    def start(self) -> None:
        self.fetcher.start()
        self.validator.start()
        self.server.start()
        self.reporter.start()
        self.adaptive.start()

        self.log.info(
            "rotator.running",
            min_pool_size=self.cfg.target_pool_size,
            headroom=self.cfg.pool_headroom,
            listen_host=self.cfg.listen_host,
            listen_port=self.cfg.listen_port,
            test_url=self.cfg.test_url,
            sources=self.cfg.proxy_sources,
            ipsum_enabled=self.cfg.ipsum_enabled,
            ipsum_url=self.cfg.ipsum_url,
            ipsum_refresh_interval=self.cfg.ipsum_refresh_interval,
        )

        try:
            self.loop.run_forever()
        except KeyboardInterrupt:
            self.log.info("rotator.interrupt_received")
        finally:
            self.shutdown()

    def shutdown(self) -> None:
        for t in (self.fetcher, self.validator, self.server,
                  self.reporter, self.adaptive):
            t.stop()
        for t in (self.fetcher, self.validator, self.server,
                  self.reporter, self.adaptive):
            t.join(timeout=5)
        self.loop.call_soon_threadsafe(self.loop.stop)
        self.log.info("rotator.shutdown_complete")


# =============================================================================
# CLI
# =============================================================================

def build_argparser() -> argparse.ArgumentParser:
    p = argparse.ArgumentParser(
        prog="proxyDrum",
        description="Three-thread HTTP/HTTPS proxy rotator",
    )
    p.add_argument("--target", type=int, default=None,
                   help="minimum pool floor (override PROXYDRUM_TARGET_POOL_SIZE)")
    p.add_argument("--headroom", type=float, default=None,
                   help="pool headroom multiplier (override PROXYDRUM_POOL_HEADROOM)")
    p.add_argument("--test-url", default=None)
    p.add_argument("--listen", default=None)
    p.add_argument("--port", type=int, default=None)
    p.add_argument("--source", action="append", default=None,
                   help="add or override a source URL (repeatable)")
    p.add_argument("--no-affinity", dest="affinity",
                   action="store_false", default=None)
    p.add_argument("--affinity-ttl", type=float, default=None)
    p.add_argument("--wait-queue-max", type=int, default=None)
    p.add_argument("--wait-timeout", type=float, default=None)
    p.add_argument("--report-interval", type=float, default=None)
    p.add_argument("--no-ipsum", dest="ipsum",
                   action="store_false", default=None,
                   help="disable the IPsum threat feed")
    p.add_argument("--json-logs", dest="json_logs",
                   action="store_true", default=None)
    p.add_argument("--no-json-logs", dest="json_logs",
                   action="store_false")
    p.add_argument("--verbose", action="store_true")
    return p


def main() -> None:
    args = build_argparser().parse_args()

    cfg = Config()

    if args.target is not None:
        cfg.target_pool_size = args.target
    if args.headroom is not None:
        cfg.pool_headroom = args.headroom
    if args.test_url is not None:
        cfg.test_url = args.test_url
    if args.listen is not None:
        cfg.listen_host = args.listen
    if args.port is not None:
        cfg.listen_port = args.port
    if args.source:
        for entry in args.source:
            parsed = _parse_sources(entry)
            for url in parsed:
                if url not in cfg.proxy_sources:
                    cfg.proxy_sources.append(url)
    if args.affinity is not None:
        cfg.affinity_enabled = args.affinity
    if args.affinity_ttl is not None:
        cfg.affinity_ttl = args.affinity_ttl
    if args.wait_queue_max is not None:
        cfg.wait_queue_max = args.wait_queue_max
    if args.wait_timeout is not None:
        cfg.wait_timeout = args.wait_timeout
    if args.report_interval is not None:
        cfg.report_interval = args.report_interval
    if args.ipsum is not None:
        cfg.ipsum_enabled = args.ipsum
    if args.json_logs is not None:
        cfg.log_json = args.json_logs
    if args.verbose:
        cfg.log_level = logging.DEBUG

    _configure_structlog(cfg.log_json, cfg.log_level)

    Rotator(cfg).start()


if __name__ == "__main__":
    main()
