Keep proxy_pool_enabled/managed off for direct registration, while shipping SSO OAuth CPA export, visible Turnstile click handling, and proxy harvest tooling.
502 lines
17 KiB
Python
502 lines
17 KiB
Python
#!/usr/bin/env python
|
|
# -*- coding: utf-8 -*-
|
|
"""代理生命周期管理:
|
|
|
|
1. 加权拉取公开源
|
|
2. 测通 https://accounts.x.ai/ 才入库
|
|
3. 库存定期巡检,失败删除
|
|
4. 使用中失败 / 打不开注册页 / 被 x.ai 屏蔽 → 立刻删除
|
|
5. 源连续无产出或拉失败 → 降权/禁用
|
|
|
|
CLI:
|
|
python proxy_manager.py harvest
|
|
python proxy_manager.py check
|
|
python proxy_manager.py loop
|
|
python proxy_manager.py status
|
|
python proxy_manager.py pick
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import json
|
|
import re
|
|
import sys
|
|
import threading
|
|
import time
|
|
from concurrent.futures import ThreadPoolExecutor, as_completed
|
|
from pathlib import Path
|
|
from typing import Callable
|
|
|
|
_ROOT = Path(__file__).resolve().parent
|
|
if str(_ROOT) not in sys.path:
|
|
sys.path.insert(0, str(_ROOT))
|
|
|
|
from proxy_sources import DEFAULT_SOURCES, harvest # noqa: E402
|
|
from proxy_store import ProxyStore, normalize_proxy, proxy_log_label # noqa: E402
|
|
|
|
try:
|
|
from curl_cffi import requests as _req
|
|
except Exception: # pragma: no cover
|
|
import requests as _req # type: ignore
|
|
|
|
|
|
LogFn = Callable[[str], None]
|
|
XAI_PROBE_URL = "https://accounts.x.ai/sign-up?redirect=grok-com"
|
|
XAI_HOST = "accounts.x.ai"
|
|
|
|
# 命中这些说明代理通了 x.ai 注册页(或至少没被墙/封到无法用)
|
|
_OK_MARKERS = (
|
|
"x.ai",
|
|
"sign-up",
|
|
"signup",
|
|
"sign up",
|
|
"email",
|
|
"cloudflare",
|
|
"turnstile",
|
|
"cf-turnstile",
|
|
"使用邮箱",
|
|
"create",
|
|
"accounts.x.ai",
|
|
)
|
|
# 明确坏标记
|
|
_BAD_MARKERS = (
|
|
"blocked due to abusive traffic",
|
|
"access denied",
|
|
"attention required",
|
|
"just a moment", # 纯 CF 墙页且无后续时也算可疑,但需配合 status
|
|
"error code: 5",
|
|
)
|
|
|
|
|
|
def _log_print(msg: str) -> None:
|
|
print(msg, flush=True)
|
|
|
|
|
|
def load_config() -> dict:
|
|
p = _ROOT / "config.json"
|
|
if not p.is_file():
|
|
return {}
|
|
try:
|
|
return json.loads(p.read_text(encoding="utf-8"))
|
|
except Exception:
|
|
return {}
|
|
|
|
|
|
def probe_accounts_xai(
|
|
proxy: str,
|
|
*,
|
|
timeout: float = 12.0,
|
|
) -> tuple[bool, str, float]:
|
|
"""测代理是否能访问 accounts.x.ai。
|
|
|
|
返回 (ok, detail, latency_ms)
|
|
"""
|
|
p = normalize_proxy(proxy)
|
|
if not p:
|
|
return False, "empty", 0.0
|
|
t0 = time.time()
|
|
try:
|
|
r = _req.get(
|
|
XAI_PROBE_URL,
|
|
proxies={"http": p, "https": p},
|
|
timeout=timeout,
|
|
impersonate="chrome120",
|
|
headers={
|
|
"User-Agent": (
|
|
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
|
|
"AppleWebKit/537.36 (KHTML, like Gecko) "
|
|
"Chrome/138.0.0.0 Safari/537.36"
|
|
),
|
|
"Accept": "text/html,application/xhtml+xml",
|
|
"Accept-Language": "en-US,en;q=0.9",
|
|
},
|
|
allow_redirects=True,
|
|
)
|
|
lat = (time.time() - t0) * 1000
|
|
status = int(getattr(r, "status_code", 0) or 0)
|
|
text = (r.text or "")[:8000]
|
|
low = text.lower()
|
|
final_url = str(getattr(r, "url", "") or "")
|
|
|
|
if status in (403, 429, 503):
|
|
if any(m in low for m in ("abusive traffic", "access denied", "blocked")):
|
|
return False, f"blocked status={status}", lat
|
|
if status >= 500:
|
|
return False, f"status={status}", lat
|
|
if status == 0:
|
|
return False, "no status", lat
|
|
|
|
# 明确坏
|
|
for m in _BAD_MARKERS:
|
|
if m in low and "sign" not in low and "email" not in low:
|
|
# CF interstitial 单独处理:有时还能过,不算立刻死
|
|
if m == "just a moment":
|
|
break
|
|
return False, f"bad marker:{m}", lat
|
|
|
|
# 成功条件:200/403 但页面像 x.ai 注册/登录相关
|
|
if status in (200, 301, 302, 303, 307, 308, 403, 401):
|
|
if XAI_HOST in final_url or any(m in low for m in _OK_MARKERS):
|
|
# 再加一刀:能拿到 HTML 长度
|
|
if len(text) < 80 and status != 200:
|
|
return False, f"short body status={status}", lat
|
|
return True, f"ok status={status} len={len(text)}", lat
|
|
return False, f"status={status} no markers", lat
|
|
except Exception as exc:
|
|
lat = (time.time() - t0) * 1000
|
|
return False, f"{type(exc).__name__}: {exc}"[:160], lat
|
|
|
|
|
|
def probe_batch(
|
|
proxies: list[str],
|
|
*,
|
|
workers: int = 40,
|
|
timeout: float = 12.0,
|
|
log: LogFn | None = None,
|
|
) -> list[tuple[str, bool, str, float]]:
|
|
log = log or _log_print
|
|
results: list[tuple[str, bool, str, float]] = []
|
|
total = len(proxies)
|
|
done = 0
|
|
lock = threading.Lock()
|
|
|
|
def one(p: str):
|
|
ok, detail, lat = probe_accounts_xai(p, timeout=timeout)
|
|
return p, ok, detail, lat
|
|
|
|
with ThreadPoolExecutor(max_workers=max(1, workers)) as ex:
|
|
futs = [ex.submit(one, p) for p in proxies]
|
|
for fu in as_completed(futs):
|
|
item = fu.result()
|
|
results.append(item)
|
|
with lock:
|
|
done += 1
|
|
if done % 25 == 0 or done == total:
|
|
ok_n = sum(1 for _, o, _, _ in results if o)
|
|
log(f"[proxy-probe] progress {done}/{total} ok={ok_n}")
|
|
return results
|
|
|
|
|
|
class ProxyManager:
|
|
def __init__(self, config: dict | None = None, log: LogFn | None = None):
|
|
self.cfg = config if config is not None else load_config()
|
|
self.log = log or _log_print
|
|
db = self.cfg.get("proxy_db") or self.cfg.get("proxy_alive_db") or "./proxy_alive.json"
|
|
self.store = ProxyStore(db)
|
|
self.sources = list(DEFAULT_SOURCES)
|
|
# 允许 config 覆盖源权重
|
|
overrides = self.cfg.get("proxy_source_weights") or {}
|
|
if isinstance(overrides, dict):
|
|
for s in self.sources:
|
|
if s.id in overrides:
|
|
try:
|
|
s.weight = int(overrides[s.id])
|
|
except Exception:
|
|
pass
|
|
|
|
# ── harvest + validate + insert ──
|
|
|
|
def harvest_and_store(self) -> dict:
|
|
pick_k = int(self.cfg.get("proxy_source_pick_k", 3) or 3)
|
|
per_limit = int(self.cfg.get("proxy_harvest_per_source", 150) or 150)
|
|
total_limit = int(self.cfg.get("proxy_harvest_total", 400) or 400)
|
|
workers = int(self.cfg.get("proxy_probe_workers", 40) or 40)
|
|
timeout = float(self.cfg.get("proxy_probe_timeout", 12) or 12)
|
|
max_keep = int(self.cfg.get("proxy_max_alive", 200) or 200)
|
|
|
|
self.log(
|
|
f"[proxy] harvest start pick_k={pick_k} per={per_limit} total={total_limit}"
|
|
)
|
|
candidates, by_source = harvest(
|
|
self.sources,
|
|
pick_k=pick_k,
|
|
per_source_limit=per_limit,
|
|
total_limit=total_limit,
|
|
store=self.store,
|
|
log=self.log,
|
|
)
|
|
# 去掉已在库的
|
|
existing = set(self.store.proxy_urls(alive_only=False))
|
|
fresh = [p for p in candidates if p not in existing]
|
|
self.log(f"[proxy] candidates={len(candidates)} fresh={len(fresh)} existing={len(existing)}")
|
|
if not fresh:
|
|
return {"harvested": len(candidates), "tested": 0, "kept": 0, "by_source": {}}
|
|
|
|
results = probe_batch(fresh, workers=workers, timeout=timeout, log=self.log)
|
|
kept = 0
|
|
kept_by_src: dict[str, int] = {sid: 0 for sid in by_source}
|
|
# reverse map proxy -> source
|
|
src_of: dict[str, str] = {}
|
|
for sid, items in by_source.items():
|
|
for p in items:
|
|
src_of.setdefault(p, sid)
|
|
|
|
for p, ok, detail, lat in results:
|
|
if not ok:
|
|
continue
|
|
sid = src_of.get(p, "")
|
|
self.store.upsert(p, source_id=sid, latency_ms=lat, note=detail)
|
|
kept += 1
|
|
if sid:
|
|
kept_by_src[sid] = kept_by_src.get(sid, 0) + 1
|
|
|
|
# 回写各源 kept,并惩罚 0 kept 的源
|
|
for sid, items in by_source.items():
|
|
k = kept_by_src.get(sid, 0)
|
|
if items:
|
|
if k > 0:
|
|
self.store.source_ok(sid, harvested=len(items), kept=k)
|
|
else:
|
|
# 有采集但测 accounts.x.ai 全挂 → 记失败(连续则禁用)
|
|
disabled = self.store.source_fail(
|
|
sid,
|
|
reason="harvested but 0 passed accounts.x.ai",
|
|
disable_after=int(self.cfg.get("proxy_source_disable_after", 3) or 3),
|
|
)
|
|
self.log(
|
|
f"[proxy] 源 {sid} 本轮 0 入库"
|
|
+ (" → 已禁用" if disabled else "")
|
|
)
|
|
|
|
# 超容量:按 last_ok 淘汰旧的
|
|
alive = self.store.list_proxies(alive_only=True)
|
|
if max_keep > 0 and len(alive) > max_keep:
|
|
alive.sort(key=lambda x: float(x.get("last_ok") or 0), reverse=True)
|
|
for item in alive[max_keep:]:
|
|
self.store.remove(item["proxy"])
|
|
self.log(f"[proxy] trim alive to {max_keep}")
|
|
|
|
# 同步 proxies.txt 方便人工看
|
|
self.export_txt()
|
|
st = self.store.stats()
|
|
self.log(f"[proxy] harvest done kept={kept} alive={st['alive']} db={st['db']}")
|
|
return {
|
|
"harvested": len(candidates),
|
|
"tested": len(fresh),
|
|
"kept": kept,
|
|
"alive": st["alive"],
|
|
"by_source": kept_by_src,
|
|
}
|
|
|
|
def export_txt(self) -> Path:
|
|
path = _ROOT / str(self.cfg.get("proxy_pool_file") or "proxies.txt")
|
|
if not path.is_absolute():
|
|
path = (_ROOT / path).resolve()
|
|
lines = [
|
|
"# auto-exported by proxy_manager — only accounts.x.ai-validated proxies",
|
|
f"# updated {time.strftime('%Y-%m-%d %H:%M:%S')}",
|
|
]
|
|
for item in sorted(
|
|
self.store.list_proxies(alive_only=True),
|
|
key=lambda x: float(x.get("latency_ms") or 99999),
|
|
):
|
|
lines.append(item["proxy"])
|
|
path.write_text("\n".join(lines) + "\n", encoding="utf-8")
|
|
return path
|
|
|
|
# ── recheck inventory ──
|
|
|
|
def recheck_inventory(self) -> dict:
|
|
items = self.store.list_proxies(alive_only=True)
|
|
if not items:
|
|
self.log("[proxy] recheck: empty inventory")
|
|
return {"checked": 0, "removed": 0, "ok": 0}
|
|
workers = int(self.cfg.get("proxy_probe_workers", 40) or 40)
|
|
timeout = float(self.cfg.get("proxy_probe_timeout", 12) or 12)
|
|
proxies = [x["proxy"] for x in items]
|
|
self.log(f"[proxy] recheck {len(proxies)} alive proxies")
|
|
results = probe_batch(proxies, workers=workers, timeout=timeout, log=self.log)
|
|
removed = 0
|
|
ok_n = 0
|
|
for p, ok, detail, lat in results:
|
|
if ok:
|
|
self.store.mark_ok(p, latency_ms=lat)
|
|
ok_n += 1
|
|
else:
|
|
if self.store.mark_fail(p, reason=detail, remove_immediately=True):
|
|
removed += 1
|
|
self.log(f"[proxy] drop {proxy_log_label(p)} ({detail})")
|
|
self.export_txt()
|
|
self.log(f"[proxy] recheck done ok={ok_n} removed={removed}")
|
|
return {"checked": len(proxies), "removed": removed, "ok": ok_n}
|
|
|
|
# ── runtime hooks for register ──
|
|
|
|
def acquire_for_use(self) -> str:
|
|
mode = str(self.cfg.get("proxy_pool_mode") or "weighted")
|
|
p = self.store.pick(mode=mode)
|
|
if p:
|
|
self.store.touch_use(p)
|
|
self.log(f"[proxy] acquire {proxy_log_label(p)} (alive={self.store.count()})")
|
|
return p
|
|
|
|
def report_success(self, proxy: str) -> None:
|
|
p = normalize_proxy(proxy)
|
|
if not p:
|
|
return
|
|
self.store.mark_ok(p)
|
|
|
|
def report_failure(self, proxy: str, reason: str = "") -> None:
|
|
p = normalize_proxy(proxy)
|
|
if not p:
|
|
return
|
|
reason = (reason or "").strip()
|
|
# 识别 x.ai 屏蔽 / 注册页缺失
|
|
low = reason.lower()
|
|
hard = any(
|
|
k in low
|
|
for k in (
|
|
"abusive traffic",
|
|
"access denied",
|
|
"blocked",
|
|
"403",
|
|
"未找到",
|
|
"使用邮箱注册",
|
|
"turnstile",
|
|
"timeout",
|
|
"连接已断开",
|
|
"page disconnected",
|
|
"err_proxy",
|
|
"proxy",
|
|
"tunnel",
|
|
)
|
|
)
|
|
removed = self.store.mark_fail(
|
|
p,
|
|
reason=reason,
|
|
remove_immediately=True if hard or True else False,
|
|
)
|
|
if removed:
|
|
self.log(f"[proxy] use-fail drop {proxy_log_label(p)} ({reason[:120]})")
|
|
# 若该源近期大量 use-fail,记源失败
|
|
# 简化:从 meta 取 source_id
|
|
# 已删除,无法取;跳过
|
|
self.export_txt()
|
|
|
|
def status(self) -> dict:
|
|
st = self.store.stats()
|
|
st["sources_cfg"] = [
|
|
{
|
|
"id": s.id,
|
|
"weight": s.weight,
|
|
"enabled_store": self.store.source_enabled(s.id, default=True),
|
|
"eff_weight": None,
|
|
}
|
|
for s in self.sources
|
|
]
|
|
from proxy_sources import effective_weight
|
|
|
|
for row in st["sources_cfg"]:
|
|
src = next(x for x in self.sources if x.id == row["id"])
|
|
row["eff_weight"] = round(
|
|
effective_weight(src, self.store.source_meta(src.id)), 2
|
|
)
|
|
return st
|
|
|
|
def loop_forever(self) -> None:
|
|
harvest_every = int(self.cfg.get("proxy_harvest_interval_sec", 600) or 600)
|
|
check_every = int(self.cfg.get("proxy_check_interval_sec", 300) or 300)
|
|
min_alive = int(self.cfg.get("proxy_min_alive", 20) or 20)
|
|
self.log(
|
|
f"[proxy] loop start harvest_every={harvest_every}s "
|
|
f"check_every={check_every}s min_alive={min_alive}"
|
|
)
|
|
last_harvest = 0.0
|
|
last_check = 0.0
|
|
# 启动先 harvest 一轮
|
|
try:
|
|
self.harvest_and_store()
|
|
last_harvest = time.time()
|
|
except Exception as exc:
|
|
self.log(f"[proxy] initial harvest error: {exc}")
|
|
while True:
|
|
now = time.time()
|
|
try:
|
|
if self.store.count() < min_alive or now - last_harvest >= harvest_every:
|
|
self.harvest_and_store()
|
|
last_harvest = time.time()
|
|
if now - last_check >= check_every:
|
|
self.recheck_inventory()
|
|
last_check = time.time()
|
|
except KeyboardInterrupt:
|
|
self.log("[proxy] loop stopped")
|
|
return
|
|
except Exception as exc:
|
|
self.log(f"[proxy] loop error: {exc}")
|
|
time.sleep(5)
|
|
|
|
|
|
# ── 给注册机用的薄封装 ──
|
|
|
|
_manager_singleton: ProxyManager | None = None
|
|
_manager_lock = threading.Lock()
|
|
|
|
|
|
def get_manager(config: dict | None = None, log: LogFn | None = None) -> ProxyManager:
|
|
global _manager_singleton
|
|
with _manager_lock:
|
|
if _manager_singleton is None:
|
|
_manager_singleton = ProxyManager(config=config, log=log)
|
|
elif config is not None:
|
|
_manager_singleton.cfg = config
|
|
if log is not None:
|
|
_manager_singleton.log = log
|
|
return _manager_singleton
|
|
|
|
|
|
def is_xai_block_reason(msg: str) -> bool:
|
|
low = (msg or "").lower()
|
|
keys = (
|
|
"abusive traffic",
|
|
"blocked due to",
|
|
"access denied",
|
|
"err_proxy",
|
|
"proxy connection",
|
|
"tunnel connection",
|
|
"未找到「使用邮箱注册」",
|
|
"未找到\"使用邮箱注册\"",
|
|
"与页面的连接已断开",
|
|
"page disconnected",
|
|
)
|
|
return any(k in low for k in keys)
|
|
|
|
|
|
def main(argv: list[str] | None = None) -> int:
|
|
parser = argparse.ArgumentParser(description="xAI 可用代理池管理")
|
|
parser.add_argument(
|
|
"cmd",
|
|
choices=["harvest", "check", "loop", "status", "pick", "export"],
|
|
help="harvest=采集入库 | check=巡检 | loop=常驻 | status | pick | export",
|
|
)
|
|
args = parser.parse_args(argv)
|
|
cfg = load_config()
|
|
mgr = ProxyManager(config=cfg)
|
|
if args.cmd == "harvest":
|
|
r = mgr.harvest_and_store()
|
|
print(json.dumps(r, ensure_ascii=False, indent=2))
|
|
return 0
|
|
if args.cmd == "check":
|
|
r = mgr.recheck_inventory()
|
|
print(json.dumps(r, ensure_ascii=False, indent=2))
|
|
return 0
|
|
if args.cmd == "loop":
|
|
mgr.loop_forever()
|
|
return 0
|
|
if args.cmd == "status":
|
|
print(json.dumps(mgr.status(), ensure_ascii=False, indent=2))
|
|
return 0
|
|
if args.cmd == "pick":
|
|
p = mgr.acquire_for_use()
|
|
print(p or "")
|
|
return 0 if p else 1
|
|
if args.cmd == "export":
|
|
path = mgr.export_txt()
|
|
print(path)
|
|
return 0
|
|
return 1
|
|
|
|
|
|
if __name__ == "__main__":
|
|
raise SystemExit(main())
|