Files
proxy-god/proxy_chain_manager/service.py
Dr Jones b852fd264f
Some checks failed
CI / Test Python 3.10 (push) Has been cancelled
CI / Test Python 3.11 (push) Has been cancelled
CI / Test Python 3.12 (push) Has been cancelled
fix: audit round 3 - fail-closed leak, shared asyncio loop, GOST bundling, close-X UX, persona key lookup, +29 tests
2026-05-22 18:23:52 -07:00

950 lines
41 KiB
Python

from __future__ import annotations
import atexit
import asyncio
import logging
import random
import signal
import threading
import time
from collections.abc import Callable
from concurrent.futures import ThreadPoolExecutor
from typing import Any
from .config import (
Settings,
load_settings,
normalize_proxy_url,
redact_proxy_url,
save_settings,
)
from .dns_leak import flush_dns_cache
from .fetcher import fetch_proxy_json, normalize_entries
from .fingerprint import (
apply_webrtc_hardening,
disable_ipv6_on_adapters,
enable_ipv6_on_adapters,
get_computer_name,
random_hostname,
set_computer_name,
)
from .firewall import disengage as fw_disengage, emergency_disengage as fw_emergency_disengage, engage as fw_engage, is_admin
from .gost_util import build_gost_cmd, ensure_gost, popen_no_window, read_gost_log_tail, terminate_process
from .leak_detect import is_chain_leak, leak_reason
from .mac_spoof import restore_macs, spoof_all_physical
from .privacy_lan import (
LanSnapshot,
engage_lan_lockdown,
restore_lan,
)
from .telemetry_kill import (
TelemetrySnapshot,
engage_telemetry_kill,
restore_telemetry,
)
from .sysproxy import (
clear_system_proxy,
detect_policy_overrides,
diagnose_system_proxy,
is_system_proxy_set,
set_system_proxy,
)
from .validator import (
check_chain_exit_ip,
check_https_tunnel,
get_direct_ip,
validate_proxies,
)
from .vpn_detect import VpnStatus, detect_vpn
from .win_compat import probe as _win_probe
log = logging.getLogger(__name__)
Notify = Callable[[dict[str, Any]], None]
# Thread pool for blocking I/O (proxy list fetching) so we don't stall asyncio
_FETCH_POOL = ThreadPoolExecutor(max_workers=8, thread_name_prefix="fetcher")
# Back-compat for tests
from .leak_detect import is_same_subnet as _is_same_network # noqa: F401
async def _probe_tcp(host: str, port: int, timeout: float = 6.0) -> bool:
"""Open and close a TCP socket; True if the handshake completes within *timeout*."""
try:
reader, writer = await asyncio.wait_for(
asyncio.open_connection(host, port), timeout=timeout,
)
writer.close()
try:
await writer.wait_closed()
except Exception: # noqa: BLE001
pass
return True
except Exception: # noqa: BLE001
return False
class ChainService:
"""Background rotating proxy chain using GOST."""
def __init__(self, notify: Notify) -> None:
self._notify = notify
self._stop = threading.Event()
self._force_rotate = threading.Event()
self._thread: threading.Thread | None = None
self._proc = None
self._settings_lock = threading.Lock()
self._settings = load_settings()
self._current_chain: list[str] = []
# Sticky-exit: while time.monotonic() < self._sticky_until, leak-based
# auto-rotation is skipped so signup flows keep the same egress IP.
self._sticky_until: float = 0.0
# Shuffle-and-drain pool state — never repeat a proxy within a cycle
self._available: list[str] = [] # proxies not yet used this cycle
self._used: set[str] = set() # proxies used this cycle
# Per-session blacklist: proxies that crashed GOST immediately
self._blacklist: set[str] = set()
self._vpn: VpnStatus = VpnStatus()
self._mac_originals: dict[str, str] = {}
self._hostname_original: str | None = None
self._ipv6_adapters: list[str] = []
self._webrtc_was_applied: bool = False
self._lan_snap: LanSnapshot | None = None
self._telemetry_snap: TelemetrySnapshot | None = None
# Register emergency firewall disengage so a crash can't leave the
# kill-switch permanently engaged.
atexit.register(fw_emergency_disengage)
for _sig in (signal.SIGTERM, signal.SIGINT):
try:
signal.signal(_sig, self._signal_handler)
except (OSError, ValueError):
pass # not on main thread or unsupported on this platform
def _signal_handler(self, signum: int, frame: Any) -> None:
"""Received SIGTERM/SIGINT — stop cleanly and ensure firewall is disengaged."""
log.warning("Signal %s received — stopping ChainService.", signum)
self._stop.set()
terminate_process(self._proc)
fw_emergency_disengage()
def _manual_exit_url(self) -> str | None:
with self._settings_lock:
u = normalize_proxy_url(self._settings.manual_exit_proxy)
return u if u else None
@property
def settings(self) -> Settings:
with self._settings_lock:
return self._settings
@property
def current_chain(self) -> list[str]:
return list(self._current_chain)
def update_settings(self, s: Settings) -> None:
with self._settings_lock:
self._settings = s
save_settings(s)
def start(self) -> None:
if self._thread and self._thread.is_alive():
return
compat = _win_probe()
self._notify({"type": "log", "text": f"Host: {compat.summary()}"})
if not compat.has_net_cmdlets:
self._notify({
"type": "log",
"text": (
"Legacy PowerShell detected — MAC spoof / IPv6 disable / "
"hostname spoof unavailable on this host. Proxy chain + "
"kill-switch + WebRTC policy still work."
),
})
self._stop.clear()
self._thread = threading.Thread(
target=self._run_thread, name="ChainService", daemon=True
)
self._thread.start()
def stop(self) -> None:
self._stop.set()
terminate_process(self._proc)
self._proc = None
if self._thread:
self._thread.join(timeout=15)
if self._thread.is_alive():
log.warning("ChainService thread did not exit in 15s — forcing firewall disengage.")
fw_emergency_disengage()
self._teardown_network()
self._notify({"type": "state", "running": False})
def rotate_now(self) -> None:
self._force_rotate.set()
def set_sticky(self, seconds: float) -> None:
"""Hold current exit for ``seconds`` (suppresses leak-driven rotation).
Manual rotation still works. Pass 0 to clear sticky mode immediately.
"""
if seconds <= 0:
self._sticky_until = 0.0
return
self._sticky_until = time.monotonic() + float(seconds)
def sticky_remaining(self) -> float:
"""Seconds remaining on sticky-exit hold (0.0 if not sticky)."""
rem = self._sticky_until - time.monotonic()
return rem if rem > 0 else 0.0
# ─────────────────────────────────────────────────────────────────────────
def _teardown_network(self) -> None:
clear_system_proxy()
self._notify({"type": "log", "text": "System proxy cleared."})
self._restore_privacy()
with self._settings_lock:
ks_enabled = self._settings.kill_switch_enabled
if is_admin() and ks_enabled:
ok, msg = fw_disengage()
self._notify({"type": "log", "text": msg})
self._notify({"type": "firewall", "engaged": False})
def _apply_privacy(self) -> None:
s = self._settings
if s.mac_spoof_enabled and is_admin():
self._mac_originals, logs = spoof_all_physical(self._mac_originals or None)
for ln in logs:
self._notify({"type": "log", "text": f"MAC: {ln}"})
elif s.mac_spoof_enabled:
self._notify({"type": "log", "text": "MAC spoof enabled but not Admin — skipped."})
if s.spoof_hostname_enabled and is_admin():
self._hostname_original = get_computer_name()
new_name = random_hostname()
ok, msg = set_computer_name(new_name)
self._notify({"type": "log", "text": f"Hostname: {msg}"})
elif s.spoof_hostname_enabled:
self._notify({"type": "log", "text": "Hostname spoof enabled but not Admin — skipped."})
if s.disable_ipv6_while_active and is_admin():
self._ipv6_adapters, logs = disable_ipv6_on_adapters()
for ln in logs:
self._notify({"type": "log", "text": ln})
elif s.disable_ipv6_while_active:
self._notify({"type": "log", "text": "IPv6 disable enabled but not Admin — skipped."})
if s.harden_webrtc_enabled and is_admin():
ok, msg = apply_webrtc_hardening(True)
self._webrtc_was_applied = ok
self._notify({"type": "log", "text": msg})
elif s.harden_webrtc_enabled:
self._notify({"type": "log", "text": "WebRTC hardening enabled but not Admin — skipped."})
if s.lan_lockdown_enabled and is_admin():
snap, logs = engage_lan_lockdown()
self._lan_snap = snap
for ln in logs:
self._notify({"type": "log", "text": f"LAN: {ln}"})
elif s.lan_lockdown_enabled:
self._notify({"type": "log", "text": "LAN lockdown enabled but not Admin — skipped."})
if s.telemetry_kill_enabled and is_admin():
tsnap, tlogs = engage_telemetry_kill()
self._telemetry_snap = tsnap
for ln in tlogs[:8]:
self._notify({"type": "log", "text": f"Telemetry: {ln}"})
if len(tlogs) > 8:
self._notify({"type": "log", "text": f"Telemetry: … +{len(tlogs) - 8} more changes."})
elif s.telemetry_kill_enabled:
self._notify({"type": "log", "text": "Telemetry kill enabled but not Admin — skipped."})
def _restore_privacy(self) -> None:
if self._mac_originals:
for ln in restore_macs(self._mac_originals):
self._notify({"type": "log", "text": f"MAC restore: {ln}"})
self._mac_originals.clear()
if self._hostname_original and is_admin():
ok, msg = set_computer_name(self._hostname_original)
self._notify({"type": "log", "text": f"Hostname restore: {msg}"})
self._hostname_original = None
if self._ipv6_adapters:
for ln in enable_ipv6_on_adapters(self._ipv6_adapters):
self._notify({"type": "log", "text": ln})
self._ipv6_adapters.clear()
if self._webrtc_was_applied and is_admin():
apply_webrtc_hardening(False)
self._webrtc_was_applied = False
self._notify({"type": "log", "text": "WebRTC policy restored."})
if self._lan_snap is not None and is_admin():
for ln in restore_lan(self._lan_snap):
self._notify({"type": "log", "text": f"LAN: {ln}"})
self._lan_snap = None
if self._telemetry_snap is not None and is_admin():
for ln in restore_telemetry(self._telemetry_snap)[:8]:
self._notify({"type": "log", "text": f"Telemetry: {ln}"})
self._telemetry_snap = None
def _run_thread(self) -> None:
try:
asyncio.run(self._async_main())
except Exception:
log.exception("service thread failed")
self._notify({"type": "log", "text": "Fatal error in service thread (see log)."})
finally:
terminate_process(self._proc)
self._proc = None
self._teardown_network()
self._notify({"type": "state", "running": False})
# ─────────────────────────────────────────────────────────────────────────
# MAIN ASYNC LOOP
# ─────────────────────────────────────────────────────────────────────────
async def _async_main(self) -> None:
self._notify({"type": "state", "running": True})
log.info(
"Service start: chain_length=%d mode=%s sources=%d pinned=%s kill_switch=%s",
self._settings.chain_length,
self._settings.obfuscation_mode,
len(self._settings.sources),
self._settings.use_pinned_chain,
self._settings.kill_switch_enabled,
)
# ── GOST setup ───────────────────────────────────────────────────────
try:
gost = ensure_gost()
self._notify({"type": "log", "text": f"GOST ready: {gost}"})
except Exception as e:
self._notify({"type": "log", "text": f"GOST setup failed: {e}"})
return
if gost.stat().st_size < 10_000:
self._notify({"type": "log", "text": "GOST appears quarantined. Re-downloading..."})
gost.unlink(missing_ok=True)
try:
gost = ensure_gost()
except Exception as e:
self._notify({"type": "log", "text": f"GOST re-download failed: {e}"})
return
# ── VPN + direct IP ───────────────────────────────────────────────────
self._vpn = detect_vpn()
mode = "VPN-aware (/16)" if self._vpn.active else "strict (exact IP)"
self._notify({
"type": "vpn",
"active": self._vpn.active,
"label": self._vpn.label,
"adapter": self._vpn.adapter,
"leak_mode": mode,
})
self._notify({
"type": "log",
"text": (
f"VPN: {self._vpn.label}"
+ (f" ({self._vpn.adapter})" if self._vpn.adapter else "")
+ f" — leak check: {mode}"
),
})
# Warm-up: retry the direct-IP lookup a few times so a transient
# network blip during startup doesn't put us in permanent fail-closed.
real_ip = None
for attempt in range(3):
real_ip = await get_direct_ip(self._settings.ip_check_url)
if real_ip:
break
if attempt < 2:
await asyncio.sleep(2.0)
if real_ip:
self._notify({"type": "real_ip", "ip": real_ip})
label = "direct/VPN IP" if self._vpn.active else "your real IP"
self._notify({"type": "log", "text": f"{label.capitalize()}: {real_ip}"})
else:
# Leak detection now FAILS CLOSED on unknown real_ip — any chain
# whose exit IP can't be compared against a baseline will be
# treated as a leak and rotated. Make the operator aware up front.
self._notify({"type": "log", "text": (
"Could not determine direct IP after 3 attempts — leak detection "
"will FAIL CLOSED (rotate on any chain whose exit IP cannot be "
"verified). Fix internet / ip_check_url to restore."
)})
self._apply_privacy()
# ── Firewall kill-switch ──────────────────────────────────────────────
if self._settings.kill_switch_enabled:
if is_admin():
ok, msg = fw_engage(gost)
self._notify({"type": "log", "text": msg})
self._notify({"type": "firewall", "engaged": ok})
else:
self._notify({"type": "log", "text": "Kill-switch skipped (not Admin)."})
self._notify({"type": "firewall", "engaged": False})
else:
self._notify({"type": "log", "text": "Kill-switch disabled in settings."})
self._notify({"type": "firewall", "engaged": False})
log.debug(
"Listener %s | refresh=%ds health=%ds max_candidates=%d concurrency=%d",
self._settings.listen_addr(),
int(self._settings.full_refresh_seconds),
int(self._settings.health_check_seconds),
int(self._settings.max_candidates),
int(self._settings.validation_concurrency),
)
# ── Main rotation loop ────────────────────────────────────────────────
full_pool: list[str] = []
last_full = 0.0
rotation_num = 0
while not self._stop.is_set():
now = time.monotonic()
need_refresh = (
not full_pool
or now - last_full >= float(self._settings.full_refresh_seconds)
)
log.debug(
"Loop tick: full_pool=%d available=%d used=%d blacklist=%d need_refresh=%s age=%.0fs",
len(full_pool),
len(self._available),
len(self._used),
len(self._blacklist),
need_refresh,
(now - last_full) if full_pool else 0.0,
)
# Pinned (manual) chain mode — skip pool management
if self._settings.use_pinned_chain and not self._settings.pinned_chain:
self._notify({
"type": "log",
"text": (
"Pinned chain enabled but list is empty — add hops in Chain Builder "
"or disable 'Use pinned chain'."
),
})
if self._settings.use_pinned_chain and self._settings.pinned_chain:
chain = list(self._settings.pinned_chain)
# Optional per-hop TCP probe so a known-dead pinned chain
# doesn't loop forever on the 10 s retry below.
if self._settings.validate_pinned_on_start:
from urllib.parse import urlparse as _urlparse
dead = []
for hop in chain:
try:
p = _urlparse(hop if "://" in hop else "http://" + hop)
if p.hostname and p.port:
if not await _probe_tcp(p.hostname, int(p.port), 6.0):
dead.append(hop)
except Exception: # noqa: BLE001
dead.append(hop)
if dead:
self._notify({"type": "log", "text": (
f"Pinned chain: {len(dead)}/{len(chain)} hop(s) failed TCP "
f"probe — proceeding anyway (will retry on failure)."
)})
rotation_num += 1
self._notify({"type": "rotation", "n": rotation_num})
log.info("Pinned chain run: %d hops", len(chain))
ok = await self._run_chain(gost, chain, real_ip)
if not ok:
await self._sleep_interruptible(10)
if self._stop.is_set():
break
continue
# Fixed exit only: hop count = 1 → chain is only the manual exit (no pool)
mex = self._manual_exit_url()
if mex and not self._settings.use_pinned_chain and self._settings.chain_length == 1:
rotation_num += 1
self._notify({"type": "rotation", "n": rotation_num})
ok = await self._run_chain(gost, [mex], real_ip)
if not ok:
await self._sleep_interruptible(15)
if self._stop.is_set():
break
continue
# ── Pool refresh ─────────────────────────────────────────────────
if need_refresh:
self._notify({"type": "phase", "phase": "fetch"})
log.info("Pool refresh: fetching %d source(s)…", len(self._settings.sources))
full_pool = await self._build_pool()
last_full = time.monotonic()
self._available = list(full_pool)
random.shuffle(self._available)
self._used.clear()
self._notify({"type": "pool", "count": len(full_pool)})
if not full_pool:
self._notify({"type": "log", "text": "Empty pool — retrying in 60s..."})
await self._sleep_interruptible(60)
last_full = 0.0
continue
# ── Pick next chain (shuffle-and-drain, no repeats per cycle) ────
chain = self._pick_chain()
if not chain:
# Pool exhausted for this cycle — reshuffle and restart
self._notify({"type": "log", "text": "Pool cycle complete — reshuffling for next round."})
self._available = list(full_pool)
random.shuffle(self._available)
self._used.clear()
chain = self._pick_chain()
if not chain:
await self._sleep_interruptible(15)
continue
rotation_num += 1
self._notify({"type": "rotation", "n": rotation_num})
log.info(
"Auto chain #%d: %d hops | obfuscation=%s | pool_remain=%d",
rotation_num,
len(chain),
self._settings.obfuscation_mode,
len(self._available),
)
await self._run_chain(gost, chain, real_ip)
if self._stop.is_set():
break
log.info("Main loop exit (stop requested or fatal).")
terminate_process(self._proc)
self._proc = None
# ─────────────────────────────────────────────────────────────────────────
# CHAIN RUNNER
# ─────────────────────────────────────────────────────────────────────────
async def _run_chain(
self, gost: Any, chain: list[str], real_ip: str | None
) -> bool:
"""Start GOST with chain, verify exit IP, monitor until rotation/stop.
Returns True if chain ran successfully, False if it immediately failed."""
self._current_chain = list(chain)
self._notify({"type": "hops", "hops": chain, "status": "connecting"})
self._notify({"type": "phase", "phase": "gost_start"})
if self._settings.flush_dns_on_rotate:
ok, msg = flush_dns_cache()
if ok:
log.debug("DNS cache flushed before chain run")
# Per-rotation MAC re-randomization (only if mac_spoof is also on, since
# without spoof there's no original snapshot we own to mutate).
if (
self._settings.mac_rotate_on_chain_rotate
and self._settings.mac_spoof_enabled
and is_admin()
and self._mac_originals
):
_, logs = spoof_all_physical(self._mac_originals)
for ln in logs:
self._notify({"type": "log", "text": f"MAC rotate: {ln}"})
listen = self._settings.listen_addr()
cmd = build_gost_cmd(gost, listen, chain)
self._notify({
"type": "log",
"text": f"Chain #{len(self._used) // max(1, self._settings.chain_length)}: "
+ "".join(self._short(h) for h in chain),
})
red = " | ".join(redact_proxy_url(h) for h in chain)
log.debug("GOST listen=http://%s | forwards (redacted): %s", listen, red)
log.debug("GOST argv: %s … (%d args)", cmd[0], len(cmd))
terminate_process(self._proc)
self._proc = popen_no_window(cmd)
self._notify({"type": "log", "text": "GOST started — warming up (2s)…"})
for _ in range(4):
if self._stop.is_set():
terminate_process(self._proc)
self._proc = None
return False
await asyncio.sleep(0.5)
if self._proc.poll() is not None:
# GOST died immediately — blacklist pool proxies (never blacklist user fixed exit)
fixed = self._manual_exit_url()
for h in chain:
if fixed and h == fixed:
continue
self._blacklist.add(h)
self._available = [x for x in self._available if x != h]
self._notify({"type": "hops", "hops": chain, "status": "dead", "exit_ip": None})
self._notify({"type": "log", "text": "GOST exited immediately — proxies blacklisted."})
tail = read_gost_log_tail(20)
if tail.strip():
for ln in tail.splitlines()[-8:]:
if ln.strip():
self._notify({"type": "log", "text": f" GOST> {ln.rstrip()}"})
return False
local_proxy = f"http://{listen}"
timeout = min(30.0, self._settings.validation_timeout_seconds + 12.0)
self._notify({"type": "phase", "phase": "verify_chain"})
self._notify({
"type": "log",
"text": f"Exit IP check through local proxy (≤{int(timeout * 2 + 5)}s)…",
})
t0 = time.monotonic()
exit_ip = await check_chain_exit_ip(
local_proxy, self._settings.ip_check_url, timeout, chain_hops=len(chain)
)
log.debug("Initial exit IP check took %.2fs → %s", time.monotonic() - t0, exit_ip or "none")
if not exit_ip:
self._notify({"type": "hops", "hops": chain, "status": "dead", "exit_ip": None})
self._notify({"type": "log", "text": "Chain IP check failed. Rotating."})
terminate_process(self._proc)
self._proc = None
return False
if is_chain_leak(exit_ip, real_ip, self._vpn.active):
reason = leak_reason(exit_ip, real_ip, self._vpn.active)
self._notify({"type": "hops", "hops": chain, "status": "dead", "exit_ip": exit_ip})
self._notify({"type": "log", "text": f"Leak detected — {reason}. Rotating."})
log.warning("Chain leak: exit=%s real=%s vpn=%s%s", exit_ip, real_ip, self._vpn.active, reason)
# Blacklist the whole chain so we don't reuse broken proxies
fixed = self._manual_exit_url()
for h in chain:
if fixed and h == fixed:
continue
self._blacklist.add(h)
self._available = [x for x in self._available if x != h]
terminate_process(self._proc)
self._proc = None
return False
# HTTPS CONNECT must work before the chain is marked healthy.
https_ok, https_msg = await check_https_tunnel(local_proxy, timeout)
if not https_ok:
self._notify({"type": "hops", "hops": chain, "status": "dead", "exit_ip": exit_ip})
self._notify({"type": "log", "text": (
"HTTPS tunnel FAILED — chain forwards plain HTTP but refuses CONNECT. "
"Browsers will time out on HTTPS sites. Rotating. " + https_msg
)})
fixed = self._manual_exit_url()
for h in chain:
if fixed and h == fixed:
continue
self._blacklist.add(h)
self._available = [x for x in self._available if x != h]
terminate_process(self._proc)
self._proc = None
return False
self._notify({"type": "hops", "hops": chain, "status": "healthy", "exit_ip": exit_ip})
self._notify({"type": "log", "text": f"✓ Chain healthy — Exit IP: {exit_ip}"})
self._notify({"type": "log", "text": f"✓ HTTPS tunnel OK — {https_msg}"})
self._notify({"type": "phase", "phase": "running"})
# Pre-flight: surface Group Policy locks (they will override us).
pol_before = detect_policy_overrides()
if pol_before:
self._notify({"type": "log", "text": (
f"Group Policy proxy lock detected ({len(pol_before)} entries) — "
"these BEAT our settings. Browsers will keep the policy proxy "
"(or DIRECT) until those keys are removed."
)})
for p in pol_before[:3]:
self._notify({"type": "log", "text": f" ! {p}"})
set_system_proxy(
self._settings.local_host,
self._settings.local_port,
self._settings.proxy_bypass,
)
applied = is_system_proxy_set()
self._notify({
"type": "log",
"text": (
f"System proxy → {self._settings.listen_addr()} "
f"(HKCU + HKLM(adm) + Connections + WinHTTP) "
f"{'OK' if applied else 'FAILED — registry write rejected'}"
),
})
# Post-flight diagnostic — every layer's actual state.
for ln in diagnose_system_proxy():
self._notify({"type": "log", "text": f" proxy: {ln}"})
# ── Health monitor loop ───────────────────────────────────────────────
hc = int(self._settings.health_check_seconds)
while not self._stop.is_set():
log.debug("Health sleep: %ds until next exit check", hc)
result = await self._wait_health_interval()
if result in ("stop", "rotate"):
log.debug("Health loop break: %s", result)
break
if self._proc is None or self._proc.poll() is not None:
self._notify({"type": "log", "text": "GOST process died — rebuilding."})
break
self._notify({"type": "log", "text": "Health check..."})
t1 = time.monotonic()
exit_ip = await check_chain_exit_ip(
local_proxy, self._settings.ip_check_url, timeout, chain_hops=len(chain)
)
log.debug("Periodic exit IP check %.2fs → %s", time.monotonic() - t1, exit_ip or "none")
if is_chain_leak(exit_ip, real_ip, self._vpn.active):
reason = leak_reason(exit_ip, real_ip, self._vpn.active)
rem = self.sticky_remaining()
if rem > 0:
self._notify({
"type": "log",
"text": (
f"Health check flagged ({reason}) — STICKY EXIT active "
f"({int(rem)}s left), skipping rotation."
),
})
self._notify({"type": "hops", "hops": chain, "status": "healthy", "exit_ip": exit_ip})
continue
self._notify({"type": "hops", "hops": chain, "status": "dead", "exit_ip": exit_ip})
self._notify({"type": "log", "text": f"Health check failed ({reason}) — rotating."})
break
self._notify({"type": "hops", "hops": chain, "status": "healthy", "exit_ip": exit_ip})
self._notify({"type": "log", "text": f"✓ Still healthy — Exit IP: {exit_ip}"})
terminate_process(self._proc)
self._proc = None
return True
# ─────────────────────────────────────────────────────────────────────────
# POOL MANAGEMENT
# ─────────────────────────────────────────────────────────────────────────
@staticmethod
def _apply_mode_filter(candidates: list[str], mode: str) -> list[str]:
if mode == "http_only":
return [u for u in candidates if u.startswith("http://")]
if mode == "socks5_only":
return [u for u in candidates if u.startswith("socks5://")]
if mode == "random_mix":
out = list(candidates)
random.shuffle(out)
return out
return list(candidates)
def _pick_chain(self) -> list[str]:
"""Pick chain_length hops: optional fixed last hop + random prefix from pool."""
s = self._settings
chain = self._pick_chain_for_mode(s.obfuscation_mode)
if chain:
log.debug(
"Picked chain len=%d mode=%s available_left=%d",
len(chain),
s.obfuscation_mode,
len(self._available),
)
return chain
if s.obfuscation_mode != "auto" and not s.use_pinned_chain:
log.debug("Pick empty under mode=%s — retrying as auto", s.obfuscation_mode)
self._notify({
"type": "log",
"text": "Obfuscation filter left nothing usable — retrying this pick with ALL protocols (auto).",
})
ch2 = self._pick_chain_for_mode("auto")
if ch2:
log.debug("Auto retry picked len=%d", len(ch2))
return ch2
log.debug("Pick chain returned empty (pool exhausted?)")
return []
def _pick_chain_for_mode(self, mode: str) -> list[str]:
"""Build one chain using given mode (http_only / socks5_only / random_mix / auto)."""
s = self._settings
k = max(1, s.chain_length)
manual = self._manual_exit_url()
if manual and not s.use_pinned_chain:
mid_need = k - 1
base = [
u for u in self._available
if u not in self._blacklist and u != manual
]
candidates = self._apply_mode_filter(base, mode)
if mid_need == 0:
return [manual]
if len(candidates) < mid_need:
return []
prefix = candidates[:mid_need]
pset = set(prefix)
self._available = [u for u in self._available if u not in pset]
self._used.update(prefix)
self._used.add(manual)
return prefix + [manual]
base = [u for u in self._available if u not in self._blacklist]
candidates = self._apply_mode_filter(base, mode)
if len(candidates) < k:
return []
picked = candidates[:k]
picked_set = set(picked)
self._available = [u for u in self._available if u not in picked_set]
self._used.update(picked)
return picked
async def _build_pool(self) -> list[str]:
"""Fetch and validate proxies. Runs blocking I/O in thread pool."""
s = self._settings
loop = asyncio.get_running_loop()
log.info("_build_pool: %d source URL(s), prefer_elite=%s", len(s.sources), s.prefer_elite)
# Fetch all sources concurrently in thread pool (they are blocking)
async def _fetch_one(url: str) -> list[str]:
try:
rows = await asyncio.wait_for(
loop.run_in_executor(
_FETCH_POOL,
lambda u=url: fetch_proxy_json(u, timeout=50.0),
),
timeout=55.0,
)
entries = normalize_entries(rows, s.prefer_elite)
if s.prefer_elite and rows and not entries:
self._notify({
"type": "log",
"text": (
f"Source returned {len(rows)} proxies but none marked “elite” — all skipped. "
"Turn off “Elite proxies only” in Chain Builder if every fetch is empty."
),
})
self._notify({"type": "log", "text": f"Fetched {len(entries)} from source."})
return entries
except asyncio.TimeoutError:
self._notify({"type": "log", "text": f"Fetch timed out (55s): {url[:60]}..."})
return []
except Exception as e:
self._notify({"type": "log", "text": f"Fetch error: {e!s}"})
return []
src_urls = list(s.sources) if s.sources else list(Settings().sources)
if not s.sources:
self._notify({
"type": "log",
"text": "Settings had no proxy sources — using built-in defaults (save Settings to persist).",
})
if not src_urls:
self._notify({"type": "log", "text": "No proxy source URLs configured — cannot build a pool."})
return []
results = await asyncio.gather(*(_fetch_one(u) for u in src_urls))
raw_urls: list[str] = []
for chunk in results:
raw_urls.extend(chunk)
if not raw_urls:
self._notify({
"type": "log",
"text": "No proxies fetched from any source (check network, URLs, or turn off “Elite only” if every list is filtered empty).",
})
return []
# Deduplicate, remove blacklisted, shuffle before capping
seen: set[str] = set()
unique: list[str] = []
for u in raw_urls:
if u not in seen and u not in self._blacklist:
seen.add(u)
unique.append(u)
random.shuffle(unique)
if len(unique) > s.max_candidates:
unique = unique[: s.max_candidates]
log.debug("Unique candidates after dedupe/cap: %d (max_candidates=%d)", len(unique), s.max_candidates)
self._notify({"type": "phase", "phase": "validate"})
self._notify({"type": "log", "text": f"Validating {len(unique)} candidates..."})
def on_prog(done: int, total: int) -> None:
self._notify({"type": "validate_progress", "done": done, "total": total})
target = max(s.min_pool_size, s.chain_length * 3)
good = await validate_proxies(
unique,
s.ip_check_url,
s.validation_concurrency,
s.validation_timeout_seconds,
on_progress=on_prog,
target=target,
)
random.shuffle(good)
mex = normalize_proxy_url(s.manual_exit_proxy)
if mex and not s.use_pinned_chain:
self._notify({"type": "phase", "phase": "exit_check"})
self._notify({"type": "log", "text": "Validating fixed exit proxy…"})
v = await validate_proxies(
[mex],
s.ip_check_url,
1,
s.validation_timeout_seconds,
on_progress=None,
)
if v:
self._notify({"type": "log", "text": "Fixed exit proxy: OK (reachable)."})
else:
self._notify({
"type": "log",
"text": "Fixed exit proxy: validation failed — will still attempt; check URL/credentials.",
})
self._notify({"type": "log", "text": f"Valid proxies: {len(good)} / {len(unique)}"})
log.info("_build_pool done: valid=%d / tested=%d", len(good), len(unique))
return good
# ─────────────────────────────────────────────────────────────────────────
# UTILITIES
# ─────────────────────────────────────────────────────────────────────────
async def _wait_health_interval(self) -> str | None:
"""Sleep for health_check_seconds. Returns 'stop', 'rotate', or None."""
total = float(self._settings.health_check_seconds)
end = time.monotonic() + total
while time.monotonic() < end:
if self._stop.is_set():
return "stop"
if self._force_rotate.is_set():
self._force_rotate.clear()
self._notify({"type": "log", "text": "Manual rotate triggered."})
if self._settings.flush_dns_on_rotate:
ok, msg = flush_dns_cache()
self._notify({"type": "log", "text": f"DNS flush: {msg}" if ok else f"DNS flush failed: {msg}"})
return "rotate"
remaining = end - time.monotonic()
self._notify({"type": "countdown", "secs": max(0, int(remaining))})
await asyncio.sleep(1.0)
self._notify({"type": "countdown", "secs": 0})
return None
async def _sleep_interruptible(self, seconds: float) -> None:
end = time.monotonic() + seconds
while time.monotonic() < end:
if self._stop.is_set():
return
await asyncio.sleep(0.2)
@staticmethod
def _short(url: str) -> str:
url = redact_proxy_url(url)
url = (
url.replace("http://", "")
.replace("socks5://", "s5://")
.replace("socks4://", "s4://")
.replace("https://", "")
)
return url[:30] + "" if len(url) > 32 else url