Disable proxy pool by default; add managed proxy lifecycle and Turnstile fixes.
Keep proxy_pool_enabled/managed off for direct registration, while shipping SSO OAuth CPA export, visible Turnstile click handling, and proxy harvest tooling.
This commit is contained in:
1 parent
1ed18236b9
commit
a4dbe948b2
118 files changed
+5351
-318
No files matched your search
+295
@@ -0,0 +1,295 @@
|
||||
#!/usr/bin/env python
|
||||
# -*- coding: utf-8 -*-
|
||||
"""可用代理持久化库存(JSON)。
|
||||
|
||||
只存「测通过 accounts.x.ai」的代理;失效/使用失败立即删。
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import random
|
||||
import threading
|
||||
import time
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from urllib.parse import urlparse
|
||||
|
||||
_ROOT = Path(__file__).resolve().parent
|
||||
_DEFAULT_DB = _ROOT / "proxy_alive.json"
|
||||
_lock = threading.RLock()
|
||||
|
||||
|
||||
def _now() -> float:
|
||||
return time.time()
|
||||
|
||||
|
||||
def normalize_proxy(raw: str | None) -> str:
|
||||
p = (raw or "").strip()
|
||||
if not p or p.lower() in ("none", "null", "direct", "off", "-"):
|
||||
return ""
|
||||
if "://" not in p:
|
||||
# host:port:user:pass
|
||||
if p.count(":") == 3:
|
||||
host, port, user, pwd = p.split(":", 3)
|
||||
p = f"http://{user}:{pwd}@{host}:{port}"
|
||||
else:
|
||||
p = f"http://{p}"
|
||||
return p
|
||||
|
||||
|
||||
def proxy_log_label(proxy: str | None) -> str:
|
||||
p = normalize_proxy(proxy)
|
||||
if not p:
|
||||
return "(direct)"
|
||||
try:
|
||||
u = urlparse(p)
|
||||
host = u.hostname or "?"
|
||||
port = f":{u.port}" if u.port else ""
|
||||
auth = "user:***@" if u.username else ""
|
||||
return f"{u.scheme or 'http'}://{auth}{host}{port}"
|
||||
except Exception:
|
||||
return "(proxy)"
|
||||
|
||||
|
||||
class ProxyStore:
|
||||
def __init__(self, path: str | Path | None = None):
|
||||
self.path = Path(path).expanduser() if path else _DEFAULT_DB
|
||||
if not self.path.is_absolute():
|
||||
self.path = (_ROOT / self.path).resolve()
|
||||
self._data: dict[str, Any] = {"version": 1, "proxies": {}, "sources": {}}
|
||||
self.load()
|
||||
|
||||
def load(self) -> None:
|
||||
with _lock:
|
||||
if self.path.is_file():
|
||||
try:
|
||||
raw = json.loads(self.path.read_text(encoding="utf-8"))
|
||||
if isinstance(raw, dict):
|
||||
self._data = raw
|
||||
self._data.setdefault("proxies", {})
|
||||
self._data.setdefault("sources", {})
|
||||
except Exception:
|
||||
self._data = {"version": 1, "proxies": {}, "sources": {}}
|
||||
else:
|
||||
self._data = {"version": 1, "proxies": {}, "sources": {}}
|
||||
|
||||
def save(self) -> None:
|
||||
with _lock:
|
||||
self.path.parent.mkdir(parents=True, exist_ok=True)
|
||||
tmp = self.path.with_suffix(self.path.suffix + ".tmp")
|
||||
tmp.write_text(
|
||||
json.dumps(self._data, ensure_ascii=False, indent=2),
|
||||
encoding="utf-8",
|
||||
)
|
||||
tmp.replace(self.path)
|
||||
|
||||
# ── proxies ──
|
||||
|
||||
def list_proxies(self, *, alive_only: bool = True) -> list[dict]:
|
||||
with _lock:
|
||||
out = []
|
||||
for p, meta in (self._data.get("proxies") or {}).items():
|
||||
if not isinstance(meta, dict):
|
||||
continue
|
||||
if alive_only and not meta.get("alive", True):
|
||||
continue
|
||||
item = dict(meta)
|
||||
item["proxy"] = p
|
||||
out.append(item)
|
||||
return out
|
||||
|
||||
def proxy_urls(self, *, alive_only: bool = True) -> list[str]:
|
||||
return [x["proxy"] for x in self.list_proxies(alive_only=alive_only)]
|
||||
|
||||
def count(self, *, alive_only: bool = True) -> int:
|
||||
return len(self.list_proxies(alive_only=alive_only))
|
||||
|
||||
def upsert(
|
||||
self,
|
||||
proxy: str,
|
||||
*,
|
||||
source_id: str = "",
|
||||
latency_ms: float | None = None,
|
||||
exit_ip: str = "",
|
||||
note: str = "",
|
||||
) -> str:
|
||||
p = normalize_proxy(proxy)
|
||||
if not p:
|
||||
return ""
|
||||
with _lock:
|
||||
cur = self._data["proxies"].get(p) or {}
|
||||
now = _now()
|
||||
meta = {
|
||||
"alive": True,
|
||||
"source_id": source_id or cur.get("source_id") or "",
|
||||
"first_seen": cur.get("first_seen") or now,
|
||||
"last_ok": now,
|
||||
"last_check": now,
|
||||
"ok_count": int(cur.get("ok_count") or 0) + 1,
|
||||
"fail_count": 0,
|
||||
"use_count": int(cur.get("use_count") or 0),
|
||||
"use_fail": int(cur.get("use_fail") or 0),
|
||||
"latency_ms": latency_ms if latency_ms is not None else cur.get("latency_ms"),
|
||||
"exit_ip": exit_ip or cur.get("exit_ip") or "",
|
||||
"note": note or cur.get("note") or "",
|
||||
}
|
||||
self._data["proxies"][p] = meta
|
||||
self.save()
|
||||
return p
|
||||
|
||||
def touch_use(self, proxy: str) -> None:
|
||||
p = normalize_proxy(proxy)
|
||||
if not p:
|
||||
return
|
||||
with _lock:
|
||||
meta = self._data["proxies"].get(p)
|
||||
if not meta:
|
||||
return
|
||||
meta["use_count"] = int(meta.get("use_count") or 0) + 1
|
||||
meta["last_use"] = _now()
|
||||
self.save()
|
||||
|
||||
def mark_ok(self, proxy: str, *, latency_ms: float | None = None, exit_ip: str = "") -> None:
|
||||
p = normalize_proxy(proxy)
|
||||
if not p:
|
||||
return
|
||||
with _lock:
|
||||
meta = self._data["proxies"].get(p)
|
||||
if not meta:
|
||||
return
|
||||
meta["alive"] = True
|
||||
meta["last_ok"] = _now()
|
||||
meta["last_check"] = _now()
|
||||
meta["ok_count"] = int(meta.get("ok_count") or 0) + 1
|
||||
meta["fail_count"] = 0
|
||||
if latency_ms is not None:
|
||||
meta["latency_ms"] = latency_ms
|
||||
if exit_ip:
|
||||
meta["exit_ip"] = exit_ip
|
||||
self.save()
|
||||
|
||||
def mark_fail(
|
||||
self,
|
||||
proxy: str,
|
||||
*,
|
||||
reason: str = "",
|
||||
remove_immediately: bool = True,
|
||||
max_fail: int = 1,
|
||||
) -> bool:
|
||||
"""返回 True 表示已删除。"""
|
||||
p = normalize_proxy(proxy)
|
||||
if not p:
|
||||
return False
|
||||
with _lock:
|
||||
meta = self._data["proxies"].get(p)
|
||||
if not meta:
|
||||
return False
|
||||
meta["fail_count"] = int(meta.get("fail_count") or 0) + 1
|
||||
meta["use_fail"] = int(meta.get("use_fail") or 0) + 1
|
||||
meta["last_check"] = _now()
|
||||
meta["last_fail_reason"] = (reason or "")[:200]
|
||||
if remove_immediately or meta["fail_count"] >= max_fail:
|
||||
self._data["proxies"].pop(p, None)
|
||||
self.save()
|
||||
return True
|
||||
meta["alive"] = False
|
||||
self.save()
|
||||
return False
|
||||
|
||||
def remove(self, proxy: str) -> bool:
|
||||
p = normalize_proxy(proxy)
|
||||
with _lock:
|
||||
if p in self._data["proxies"]:
|
||||
self._data["proxies"].pop(p, None)
|
||||
self.save()
|
||||
return True
|
||||
return False
|
||||
|
||||
def pick(self, mode: str = "weighted") -> str:
|
||||
"""从库存挑一个。weighted=按成功率/延迟加权;random/round_robin 简化为随机。"""
|
||||
items = self.list_proxies(alive_only=True)
|
||||
if not items:
|
||||
return ""
|
||||
if mode in ("random", "round_robin"):
|
||||
return random.choice(items)["proxy"]
|
||||
weights = []
|
||||
for it in items:
|
||||
ok = max(1, int(it.get("ok_count") or 1))
|
||||
fail = int(it.get("use_fail") or 0)
|
||||
lat = float(it.get("latency_ms") or 3000.0)
|
||||
# 越快越稳权重越高
|
||||
w = ok / (1 + fail) * (1.0 / max(0.2, lat / 1000.0))
|
||||
weights.append(max(0.01, w))
|
||||
return random.choices(items, weights=weights, k=1)[0]["proxy"]
|
||||
|
||||
# ── sources health ──
|
||||
|
||||
def source_meta(self, source_id: str) -> dict:
|
||||
with _lock:
|
||||
return dict((self._data.get("sources") or {}).get(source_id) or {})
|
||||
|
||||
def source_ok(self, source_id: str, *, harvested: int = 0, kept: int = 0) -> None:
|
||||
if not source_id:
|
||||
return
|
||||
with _lock:
|
||||
cur = self._data["sources"].setdefault(source_id, {})
|
||||
cur["enabled"] = True
|
||||
cur["last_ok"] = _now()
|
||||
cur["last_run"] = _now()
|
||||
cur["fail_streak"] = 0
|
||||
cur["ok_runs"] = int(cur.get("ok_runs") or 0) + 1
|
||||
cur["last_harvested"] = harvested
|
||||
cur["last_kept"] = kept
|
||||
self.save()
|
||||
|
||||
def source_fail(self, source_id: str, *, reason: str = "", disable_after: int = 3) -> bool:
|
||||
"""源连续失败 disable_after 次后禁用。返回是否已禁用。"""
|
||||
if not source_id:
|
||||
return False
|
||||
with _lock:
|
||||
cur = self._data["sources"].setdefault(source_id, {})
|
||||
cur["last_run"] = _now()
|
||||
cur["last_fail_reason"] = (reason or "")[:200]
|
||||
cur["fail_streak"] = int(cur.get("fail_streak") or 0) + 1
|
||||
cur["fail_runs"] = int(cur.get("fail_runs") or 0) + 1
|
||||
disabled = False
|
||||
if cur["fail_streak"] >= max(1, disable_after):
|
||||
cur["enabled"] = False
|
||||
cur["disabled_at"] = _now()
|
||||
disabled = True
|
||||
self.save()
|
||||
return disabled
|
||||
|
||||
def source_enabled(self, source_id: str, default: bool = True) -> bool:
|
||||
with _lock:
|
||||
cur = (self._data.get("sources") or {}).get(source_id) or {}
|
||||
if "enabled" not in cur:
|
||||
return default
|
||||
return bool(cur.get("enabled"))
|
||||
|
||||
def reenable_source(self, source_id: str) -> None:
|
||||
with _lock:
|
||||
cur = self._data["sources"].setdefault(source_id, {})
|
||||
cur["enabled"] = True
|
||||
cur["fail_streak"] = 0
|
||||
cur["reenabled_at"] = _now()
|
||||
self.save()
|
||||
|
||||
def stats(self) -> dict:
|
||||
items = self.list_proxies(alive_only=False)
|
||||
alive = [x for x in items if x.get("alive", True)]
|
||||
sources = self._data.get("sources") or {}
|
||||
return {
|
||||
"total": len(items),
|
||||
"alive": len(alive),
|
||||
"sources": {
|
||||
sid: {
|
||||
"enabled": meta.get("enabled", True),
|
||||
"fail_streak": meta.get("fail_streak", 0),
|
||||
"last_kept": meta.get("last_kept", 0),
|
||||
}
|
||||
for sid, meta in sources.items()
|
||||
},
|
||||
"db": str(self.path),
|
||||
}
|
||||
Reference in new issue
Block a user