readlink/realpath of /snap/bin/chromium becomes /usr/bin/snap, which is not a browser and breaks DrissionPage launch. Keep the snap wrapper path, reject snap host basenames, and prefer real deb chromium binaries.
840 lines
29 KiB
Python
840 lines
29 KiB
Python
"""CLI wrapper for grok_register_ttk — multi-thread register + async CPA mint pipeline.
|
||
|
||
Architecture:
|
||
Register workers (R) → accounts_cli + mint_queue
|
||
Mint workers (M) → cpa_auths/xai-*.json + optional hotload
|
||
|
||
Browser lifecycle:
|
||
- One Chromium per register worker, reused via TabPool.clear_session
|
||
- Full recycle every N accounts or on error
|
||
- Register browser released BEFORE mint (mint always standalone Chromium)
|
||
- Peak browsers ≈ R + M (not 2×R)
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import os
|
||
import queue
|
||
import sys
|
||
import threading
|
||
import time
|
||
import traceback
|
||
from typing import Any
|
||
|
||
# 强制走本目录的 grok_register_ttk
|
||
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
|
||
|
||
import grok_register_ttk as reg # noqa: E402
|
||
|
||
|
||
# Linux 适配: DrissionPage 默认找 'chrome', 我们装的是 chromium
|
||
# 保留原版 slim flags + proxy,再补 chromium 路径与 turnstilePatch。
|
||
_orig_create_browser_options = reg.create_browser_options
|
||
|
||
|
||
def _patched_create_browser_options():
|
||
# Prefer original factory (proxy + CHROMIUM_SLIM_FLAGS + extension + auto-headless)
|
||
used_orig = False
|
||
try:
|
||
opts = _orig_create_browser_options()
|
||
used_orig = True
|
||
except Exception:
|
||
from DrissionPage import ChromiumOptions
|
||
|
||
opts = ChromiumOptions()
|
||
opts.auto_port()
|
||
opts.set_timeouts(base=1)
|
||
for flag in getattr(reg, "CHROMIUM_SLIM_FLAGS", ()) or ():
|
||
try:
|
||
opts.set_argument(flag)
|
||
except Exception:
|
||
pass
|
||
|
||
try:
|
||
opts.auto_port()
|
||
except Exception:
|
||
pass
|
||
try:
|
||
opts.set_timeouts(base=1)
|
||
except Exception:
|
||
pass
|
||
|
||
# Prefer register resolve_browser_path (Chromium first when turnstilePatch exists).
|
||
# Do NOT force Google Chrome here — Chrome blocks --load-extension and CF often fails.
|
||
browser_path = None
|
||
resolve = getattr(reg, "resolve_browser_path", None)
|
||
if callable(resolve):
|
||
try:
|
||
browser_path = resolve()
|
||
except Exception:
|
||
browser_path = None
|
||
if not browser_path:
|
||
normalize = getattr(reg, "_normalize_browser_path", None)
|
||
for cand in (
|
||
"/usr/bin/chromium",
|
||
"/usr/bin/chromium-browser",
|
||
"/usr/lib/chromium/chromium",
|
||
"/snap/bin/chromium", # keep wrapper path; never /usr/bin/snap
|
||
"/usr/bin/google-chrome-stable",
|
||
"/usr/bin/google-chrome",
|
||
):
|
||
if callable(normalize):
|
||
norm = normalize(cand)
|
||
if norm:
|
||
browser_path = norm
|
||
break
|
||
elif os.path.isfile(cand) or os.path.islink(cand):
|
||
# refuse snap host binary
|
||
if os.path.basename(cand) == "snap":
|
||
continue
|
||
browser_path = cand
|
||
break
|
||
# Final sanitize: never pass /usr/bin/snap to DrissionPage
|
||
if browser_path and os.path.basename(str(browser_path)) in ("snap", "snap-confine", "snapd"):
|
||
print(f"[browser] reject invalid path={browser_path!r}, re-resolve", flush=True)
|
||
browser_path = None
|
||
resolve = getattr(reg, "resolve_browser_path", None)
|
||
if callable(resolve):
|
||
try:
|
||
browser_path = resolve()
|
||
except Exception:
|
||
browser_path = None
|
||
if browser_path:
|
||
try:
|
||
opts.set_browser_path(browser_path)
|
||
print(f"[browser] path={browser_path}", flush=True)
|
||
except Exception:
|
||
pass
|
||
|
||
# create_browser_options already handles turnstilePatch (skip on Google Chrome).
|
||
# Only re-add here for the fallback options path on Chromium.
|
||
ext_path = os.path.join(os.path.dirname(os.path.abspath(reg.__file__)), "turnstilePatch")
|
||
already = False
|
||
try:
|
||
already = any(
|
||
os.path.abspath(str(p)) == os.path.abspath(ext_path)
|
||
for p in (getattr(opts, "extensions", None) or [])
|
||
)
|
||
except Exception:
|
||
already = False
|
||
if os.path.isdir(ext_path) and not already:
|
||
is_chrome = bool(
|
||
browser_path
|
||
and "chrome" in os.path.basename(browser_path)
|
||
and "chromium" not in browser_path
|
||
)
|
||
if is_chrome:
|
||
print(
|
||
"[ext] skip turnstilePatch on Google Chrome (--load-extension blocked); "
|
||
"install chromium or set browser_prefer=chromium",
|
||
flush=True,
|
||
)
|
||
else:
|
||
try:
|
||
opts.add_extension(ext_path)
|
||
print(f"[ext] turnstilePatch loaded: {ext_path}", flush=True)
|
||
except Exception:
|
||
pass
|
||
|
||
# Fallback path only: original factory already applied headless when used_orig.
|
||
if not used_orig:
|
||
try:
|
||
apply = getattr(reg, "apply_headless_to_options", None)
|
||
should = getattr(reg, "should_use_headless", None)
|
||
if callable(apply) and callable(should):
|
||
use_hl = should(None, log_hint=True)
|
||
apply(opts, use_hl)
|
||
print(f"[browser] headless={use_hl} (fallback options)", flush=True)
|
||
else:
|
||
# Minimal fallback without helpers
|
||
has_disp = sys.platform == "win32" or bool(
|
||
(os.environ.get("DISPLAY") or "").strip()
|
||
or (os.environ.get("WAYLAND_DISPLAY") or "").strip()
|
||
)
|
||
if not has_disp:
|
||
try:
|
||
opts.headless(True)
|
||
except Exception:
|
||
opts.set_argument("--headless=new")
|
||
print("[browser] headless=True (no DISPLAY, minimal fallback)", flush=True)
|
||
except Exception as exc:
|
||
print(f"[browser] headless apply failed: {exc}", flush=True)
|
||
return opts
|
||
|
||
|
||
reg.create_browser_options = _patched_create_browser_options
|
||
|
||
|
||
# ── 线程安全日志 ──
|
||
|
||
_log_queue: queue.Queue = queue.Queue()
|
||
|
||
|
||
def _log_writer():
|
||
while True:
|
||
msg = _log_queue.get()
|
||
if msg is None:
|
||
break
|
||
print(msg, flush=True)
|
||
|
||
|
||
def log(worker_id: int | str, msg: str) -> None:
|
||
_log_queue.put(f"[{time.strftime('%H:%M:%S')}] [W{worker_id}] {msg}")
|
||
|
||
|
||
# ── 统计 ──
|
||
|
||
_stats_lock = threading.Lock()
|
||
_stats = {
|
||
"reg_success": 0,
|
||
"reg_fail": 0,
|
||
"mint_success": 0,
|
||
"mint_fail": 0,
|
||
"mint_skip": 0,
|
||
}
|
||
|
||
|
||
def _inc(key: str, n: int = 1) -> None:
|
||
with _stats_lock:
|
||
_stats[key] = _stats.get(key, 0) + n
|
||
|
||
|
||
# forever 任务索引
|
||
_next_idx_lock = threading.Lock()
|
||
_next_idx = [1]
|
||
|
||
# mint 队列结束哨兵
|
||
_MINT_STOP = object()
|
||
|
||
|
||
def resolve_mint_workers(
|
||
*,
|
||
cli_value: int,
|
||
threads: int,
|
||
config: dict,
|
||
inline_mint: bool,
|
||
) -> int:
|
||
"""Resolve mint worker count.
|
||
|
||
Priority: --inline-mint > CLI --mint-workers (>=0) > config cpa_mint_workers > auto.
|
||
auto (-1): min(threads, 4) when CPA export enabled, else 0.
|
||
0: inline mint on register threads.
|
||
"""
|
||
if inline_mint:
|
||
return 0
|
||
if cli_value >= 0:
|
||
return max(0, min(int(cli_value), 10))
|
||
cfg_v = config.get("cpa_mint_workers", -1)
|
||
try:
|
||
cfg_v = int(cfg_v)
|
||
except Exception:
|
||
cfg_v = -1
|
||
if cfg_v >= 0:
|
||
return max(0, min(cfg_v, 10))
|
||
# auto
|
||
if config.get("cpa_export_enabled", True):
|
||
return max(1, min(int(threads), 4))
|
||
return 0
|
||
|
||
|
||
def resolve_mint_queue_max(config: dict, mint_workers: int, cli_value: int | None = None) -> int:
|
||
if cli_value is not None and cli_value >= 0:
|
||
return int(cli_value)
|
||
try:
|
||
v = int(config.get("cpa_mint_queue_max", 0) or 0)
|
||
except Exception:
|
||
v = 0
|
||
if v > 0:
|
||
return v
|
||
# default backpressure: 2 × mint workers (0 if no mint pool)
|
||
return max(0, mint_workers * 2) if mint_workers > 0 else 0
|
||
|
||
|
||
class DummyStop:
|
||
def __call__(self) -> bool:
|
||
return False
|
||
|
||
|
||
def _is_hotmail_provider() -> bool:
|
||
try:
|
||
provider = reg.get_email_provider()
|
||
except Exception:
|
||
provider = (getattr(reg, "config", {}) or {}).get("email_provider", "")
|
||
return str(provider or "").strip().lower() in {"hotmail", "outlook", "outlookmail", "microsoft"}
|
||
|
||
|
||
def _mark_email_stage_error(email: str, reason: str) -> None:
|
||
"""Persist failed Hotmail/Outlook aliases so the next run does not reuse them."""
|
||
if not email or not _is_hotmail_provider():
|
||
return
|
||
try:
|
||
reg.mark_error(email, reason=str(reason)[:120])
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
def _ensure_browser(worker_id: int, force_recycle: bool = False):
|
||
"""Start browser if missing; optional full recycle."""
|
||
if force_recycle:
|
||
try:
|
||
reg.stop_browser()
|
||
except Exception:
|
||
pass
|
||
if reg.TabPool.get_browser() is None:
|
||
reg.start_browser(log_callback=lambda m: log(worker_id, m))
|
||
|
||
|
||
def register_one(
|
||
worker_id: int,
|
||
idx: int,
|
||
total: int,
|
||
accounts_file: str,
|
||
*,
|
||
do_mint_inline: bool = False,
|
||
mint_queue: queue.Queue | None = None,
|
||
) -> dict | None:
|
||
"""Run one registration. Enqueue CPA mint (default) instead of blocking.
|
||
|
||
Returns dict(email, sso, profile) or None.
|
||
"""
|
||
email = ""
|
||
dev_token = ""
|
||
try:
|
||
max_mail_retry = max(1, int((getattr(reg, "config", {}) or {}).get("mail_retry_count", 3) or 3))
|
||
except Exception:
|
||
max_mail_retry = 3
|
||
cancel = DummyStop()
|
||
|
||
try:
|
||
_ensure_browser(worker_id, force_recycle=False)
|
||
except Exception as exc:
|
||
log(worker_id, f"! 浏览器启动失败: {exc}")
|
||
return None
|
||
|
||
for mail_try in range(1, max_mail_retry + 1):
|
||
email = ""
|
||
dev_token = ""
|
||
try:
|
||
log(worker_id, f"--- 第 {idx}/{total} 个账号, 邮箱尝试 {mail_try}/{max_mail_retry} ---")
|
||
log(worker_id, "1. 打开注册页")
|
||
reg.open_signup_page(log_callback=lambda m: log(worker_id, m), cancel_callback=cancel)
|
||
log(worker_id, "2. 创建邮箱并提交")
|
||
email, dev_token = reg.fill_email_and_submit(
|
||
log_callback=lambda m: log(worker_id, m), cancel_callback=cancel
|
||
)
|
||
log(worker_id, f"邮箱: {email}")
|
||
log(worker_id, "3. 拉取验证码")
|
||
code = reg.fill_code_and_submit(
|
||
email,
|
||
dev_token,
|
||
log_callback=lambda m: log(worker_id, m),
|
||
cancel_callback=cancel,
|
||
)
|
||
log(worker_id, f"验证码: {code}")
|
||
break
|
||
except Exception as exc:
|
||
msg = str(exc)
|
||
# Retry whole email stage on: no OTP mail, submit never reached code page, CF/rate limit
|
||
retryable = any(
|
||
key in msg
|
||
for key in (
|
||
"未收到验证码",
|
||
"验证码",
|
||
"未进入验证码页",
|
||
"邮箱提交被拒绝",
|
||
"rate-limited",
|
||
"rate_limited",
|
||
"Turnstile",
|
||
)
|
||
)
|
||
if retryable and mail_try < max_mail_retry:
|
||
log(worker_id, f"! 邮箱阶段可重试,换邮箱/浏览器: {msg}")
|
||
_mark_email_stage_error(email, msg)
|
||
try:
|
||
reg.restart_browser(log_callback=lambda m: log(worker_id, m))
|
||
except Exception:
|
||
pass
|
||
reg.sleep_with_cancel(1, cancel)
|
||
continue
|
||
log(worker_id, f"! 邮箱阶段失败: {msg}")
|
||
_mark_email_stage_error(email, msg)
|
||
traceback.print_exc()
|
||
_inc("reg_fail")
|
||
try:
|
||
reg.restart_browser(log_callback=lambda m: log(worker_id, m))
|
||
except Exception:
|
||
pass
|
||
return None
|
||
|
||
try:
|
||
cfg = getattr(reg, "config", {}) or {}
|
||
try:
|
||
profile_timeout = int(cfg.get("profile_timeout", 120) or 120)
|
||
except Exception:
|
||
profile_timeout = 120
|
||
try:
|
||
# Prefer faster fail+retry over long dead waits; cap via config.
|
||
sso_timeout = int(cfg.get("sso_timeout_base", 120) or 120)
|
||
except Exception:
|
||
sso_timeout = 120
|
||
# Fast-fail band: keep CF stalls from dominating wall clock.
|
||
# batch10 lesson: success profiles finish in ~6–18s; 40–50s CF fail is enough.
|
||
profile_timeout = max(35, min(profile_timeout, 70))
|
||
sso_timeout = max(30, min(sso_timeout, 90))
|
||
|
||
log(worker_id, "4. 填写资料")
|
||
profile = reg.fill_profile_and_submit(
|
||
timeout=profile_timeout,
|
||
log_callback=lambda m: log(worker_id, m),
|
||
cancel_callback=cancel,
|
||
)
|
||
log(worker_id, f"资料已填: {profile.get('given_name')} {profile.get('family_name')}")
|
||
log(worker_id, "5. 等待 sso cookie")
|
||
sso = reg.wait_for_sso_cookie(
|
||
timeout=sso_timeout,
|
||
log_callback=lambda m: log(worker_id, m),
|
||
cancel_callback=cancel,
|
||
)
|
||
password = profile.get("password", "") or ""
|
||
line = f"{email}----{password}----{sso}\n"
|
||
with open(accounts_file, "a", encoding="utf-8") as f:
|
||
f.write(line)
|
||
log(worker_id, f"+ 注册成功: {email}")
|
||
reg.mark_used(email, password)
|
||
|
||
# Capture cookies BEFORE releasing browser (for mint cookie inject)
|
||
page = reg._get_page()
|
||
cookies = []
|
||
try:
|
||
import cpa_export as _cpa_exp
|
||
|
||
cookies = _cpa_exp.export_cookies_from_page(page) if page is not None else []
|
||
except Exception:
|
||
cookies = []
|
||
if cookies:
|
||
log(worker_id, f"[*] 导出 cookie {len(cookies)} 条供 mint 注入")
|
||
|
||
if page and reg.PERF_FLAGS.get("cookie_snapshot", True):
|
||
try:
|
||
reg.save_cookies_snapshot(page, "success", email)
|
||
except Exception:
|
||
pass
|
||
try:
|
||
reg.add_token_to_grok2api_pools(
|
||
sso, email=email, log_callback=lambda m: log(worker_id, m)
|
||
)
|
||
except Exception as exc:
|
||
log(worker_id, f"[Debug] grok2api: {exc}")
|
||
|
||
# Release / recycle register browser BEFORE mint so peak browsers ≈ R+M
|
||
try:
|
||
reg.prepare_browser_for_next_account(log_callback=lambda m: log(worker_id, m))
|
||
except Exception:
|
||
try:
|
||
reg.stop_browser()
|
||
except Exception:
|
||
pass
|
||
|
||
job = {
|
||
"email": email,
|
||
"password": password,
|
||
"sso": sso,
|
||
"profile": profile,
|
||
"idx": idx,
|
||
"cookies": cookies,
|
||
}
|
||
|
||
if do_mint_inline:
|
||
_run_mint_job(f"R{worker_id}", job, getattr(reg, "config", {}) or {})
|
||
elif mint_queue is not None:
|
||
# backpressure: wait while queue is saturated
|
||
qmax = int(getattr(mint_queue, "_reg_qmax", 0) or 0)
|
||
while qmax > 0 and mint_queue.qsize() >= qmax:
|
||
log(worker_id, f"[cpa] mint 队列背压 qsize={mint_queue.qsize()}≥{qmax},等待...")
|
||
time.sleep(1.0)
|
||
mint_queue.put(job)
|
||
log(worker_id, f"[cpa] enqueued mint for {email} (queue≈{mint_queue.qsize()})")
|
||
else:
|
||
log(worker_id, "[cpa] mint skipped (no queue / inline)")
|
||
|
||
_inc("reg_success")
|
||
return job
|
||
except Exception as exc:
|
||
log(worker_id, f"! 注册失败: {exc}")
|
||
reg.mark_error(email or "", reason=str(exc)[:120])
|
||
traceback.print_exc()
|
||
_inc("reg_fail")
|
||
try:
|
||
reg.restart_browser(log_callback=lambda m: log(worker_id, m))
|
||
except Exception:
|
||
pass
|
||
return None
|
||
|
||
|
||
def _run_mint_job(worker_id: int | str, job: dict[str, Any], config: dict) -> dict:
|
||
"""Standalone CPA mint (own Chromium). Never reuses register browser."""
|
||
email = job.get("email") or ""
|
||
password = job.get("password") or ""
|
||
if not email or not password:
|
||
_inc("mint_fail")
|
||
return {"ok": False, "error": "missing email/password", "email": email}
|
||
if not config.get("cpa_export_enabled", True):
|
||
_inc("mint_skip")
|
||
log(worker_id, f"[cpa] export disabled, skip {email}")
|
||
return {"ok": False, "skipped": True, "email": email}
|
||
try:
|
||
import cpa_export
|
||
|
||
# page=None always — force standalone path inside export
|
||
result = cpa_export.export_cpa_xai_for_account(
|
||
email,
|
||
password,
|
||
page=None,
|
||
cookies=job.get("cookies"),
|
||
sso=job.get("sso") or "",
|
||
config=config,
|
||
log_callback=lambda m: log(worker_id, m),
|
||
)
|
||
if result.get("ok"):
|
||
log(worker_id, f"+ CPA auth: {result.get('path')}")
|
||
_inc("mint_success")
|
||
elif result.get("skipped"):
|
||
_inc("mint_skip")
|
||
log(worker_id, f"[cpa] skipped: {result.get('reason')}")
|
||
else:
|
||
_inc("mint_fail")
|
||
log(worker_id, f"! CPA auth 未成功: {result.get('error') or result}")
|
||
return result
|
||
except Exception as exc:
|
||
_inc("mint_fail")
|
||
log(worker_id, f"! CPA export 异常: {exc}")
|
||
traceback.print_exc()
|
||
return {"ok": False, "error": str(exc), "email": email}
|
||
|
||
|
||
def _register_worker(
|
||
worker_id: int,
|
||
task_queue: queue.Queue,
|
||
total: int,
|
||
accounts_file: str,
|
||
mint_queue: queue.Queue | None,
|
||
forever: bool,
|
||
do_mint_inline: bool,
|
||
):
|
||
while True:
|
||
try:
|
||
idx = task_queue.get_nowait()
|
||
except queue.Empty:
|
||
if not forever:
|
||
break
|
||
with _next_idx_lock:
|
||
nxt = _next_idx[0]
|
||
_next_idx[0] = nxt + 5
|
||
for i in range(nxt, nxt + 5):
|
||
task_queue.put(i)
|
||
continue
|
||
|
||
retry = 0
|
||
while retry < 2:
|
||
try:
|
||
result = register_one(
|
||
worker_id,
|
||
idx,
|
||
total,
|
||
accounts_file,
|
||
do_mint_inline=do_mint_inline,
|
||
mint_queue=mint_queue,
|
||
)
|
||
if result:
|
||
break
|
||
retry += 1
|
||
if retry < 2:
|
||
log(worker_id, f"[retry] 账号 {idx} 失败,重试 {retry}/1")
|
||
try:
|
||
reg.restart_browser(log_callback=lambda m: log(worker_id, m))
|
||
except Exception:
|
||
pass
|
||
except Exception:
|
||
retry += 1
|
||
if retry < 2:
|
||
log(worker_id, f"[retry] 账号 {idx} 异常,重试 {retry}/1")
|
||
traceback.print_exc()
|
||
try:
|
||
reg.restart_browser(log_callback=lambda m: log(worker_id, m))
|
||
except Exception:
|
||
pass
|
||
|
||
if retry >= 2:
|
||
# register_one already counted fail on exception path; if both returned None, count once more only if needed
|
||
pass
|
||
|
||
# worker exit: free browser
|
||
try:
|
||
reg.stop_browser()
|
||
except Exception:
|
||
pass
|
||
log(worker_id, "register worker exit")
|
||
|
||
|
||
def _mint_worker(worker_id: str, mint_queue: queue.Queue, config: dict):
|
||
while True:
|
||
job = mint_queue.get()
|
||
try:
|
||
if job is _MINT_STOP:
|
||
break
|
||
if not isinstance(job, dict):
|
||
continue
|
||
_run_mint_job(worker_id, job, config)
|
||
finally:
|
||
mint_queue.task_done()
|
||
try:
|
||
from cpa_xai.browser_confirm import shutdown_mint_browsers
|
||
|
||
shutdown_mint_browsers()
|
||
except Exception:
|
||
pass
|
||
log(worker_id, "mint worker exit")
|
||
|
||
|
||
def main() -> int:
|
||
parser = argparse.ArgumentParser(description="CLI runner for grok_register_ttk (pipelined).")
|
||
parser.add_argument("--count", type=int, default=1, help="账号总数目标(0=不限;含已有)")
|
||
parser.add_argument(
|
||
"--extra",
|
||
type=int,
|
||
default=0,
|
||
help="在已有 accounts 基础上再新注册 N 个",
|
||
)
|
||
parser.add_argument("--threads", type=int, default=1, help="注册并发线程数(1-10)")
|
||
parser.add_argument(
|
||
"--mint-workers",
|
||
type=int,
|
||
default=-1,
|
||
help="CPA mint 并发:-1=用 config/auto;0=内联;1-10=固定。覆盖 config.cpa_mint_workers",
|
||
)
|
||
parser.add_argument(
|
||
"--mint-queue-max",
|
||
type=int,
|
||
default=-1,
|
||
help="mint 队列背压上限:-1=用 config/auto(2×workers);0=不限制",
|
||
)
|
||
parser.add_argument("--accounts-file", default=os.path.join(os.path.dirname(__file__), "accounts_cli.txt"))
|
||
parser.add_argument("--fast", action="store_true", default=True, help="快速模式(默认开):压缩 sleep、关截图")
|
||
parser.add_argument("--no-fast", action="store_true", help="关闭快速模式")
|
||
parser.add_argument("--no-browser-reuse", action="store_true", help="每号强制 quit 浏览器")
|
||
parser.add_argument("--browser-recycle-every", type=int, default=25, help="复用 N 次后完整回收")
|
||
parser.add_argument("--cookie-snapshot", action="store_true", help="注册成功写 cookie 快照(默认关,fast)")
|
||
parser.add_argument("--inline-mint", action="store_true", help="强制注册线程内联 mint(调试用)")
|
||
parser.add_argument(
|
||
"--headless",
|
||
action="store_true",
|
||
help="强制无头浏览器(无 DISPLAY 时默认也会自动 headless)",
|
||
)
|
||
parser.add_argument(
|
||
"--headed",
|
||
action="store_true",
|
||
help="强制有头浏览器(需 DISPLAY;无显示时仍会回退 headless 并提示 Xvfb)",
|
||
)
|
||
args = parser.parse_args()
|
||
|
||
reg.load_config()
|
||
cfg0 = getattr(reg, "config", {}) or {}
|
||
threads = max(1, min(args.threads, 10))
|
||
fast = bool(args.fast) and not bool(args.no_fast)
|
||
|
||
# Headless preference: CLI > config/env/auto (no DISPLAY → headless)
|
||
if args.headless and args.headed:
|
||
print("[!] --headless 与 --headed 互斥,忽略 --headed", flush=True)
|
||
args.headed = False
|
||
if args.headless:
|
||
cfg0["browser_headless"] = True
|
||
if isinstance(getattr(reg, "config", None), dict):
|
||
reg.config["browser_headless"] = True
|
||
elif args.headed:
|
||
cfg0["browser_headless"] = False
|
||
if isinstance(getattr(reg, "config", None), dict):
|
||
reg.config["browser_headless"] = False
|
||
|
||
mint_workers = resolve_mint_workers(
|
||
cli_value=args.mint_workers,
|
||
threads=threads,
|
||
config=cfg0,
|
||
inline_mint=bool(args.inline_mint),
|
||
)
|
||
do_mint_inline = mint_workers == 0
|
||
mint_qmax = resolve_mint_queue_max(
|
||
cfg0,
|
||
mint_workers,
|
||
cli_value=(None if args.mint_queue_max < 0 else args.mint_queue_max),
|
||
)
|
||
|
||
# perf knobs
|
||
reg.configure_perf(
|
||
fast=fast,
|
||
sleep_scale=0.15 if fast else 1.0,
|
||
skip_debug_io=fast,
|
||
cookie_snapshot=bool(args.cookie_snapshot) or not fast,
|
||
async_side_effects=True,
|
||
browser_reuse=not args.no_browser_reuse,
|
||
browser_recycle_every=max(1, int(args.browser_recycle_every)),
|
||
)
|
||
|
||
# 断点续跑
|
||
done_count = 0
|
||
if os.path.exists(args.accounts_file):
|
||
with open(args.accounts_file) as f:
|
||
done_count = sum(1 for line in f if line.strip())
|
||
|
||
if args.extra and args.extra > 0:
|
||
target_total = done_count + args.extra
|
||
remaining = args.extra
|
||
print(
|
||
f"[*] 配置加载完成,额外新注册 {args.extra} 个(当前已有 {done_count} → 目标 {target_total}),"
|
||
f"注册线程={threads} mint_workers={mint_workers} mint_queue_max={mint_qmax} fast={fast}",
|
||
flush=True,
|
||
)
|
||
args.count = target_total
|
||
elif args.count == 0:
|
||
remaining = None
|
||
print(
|
||
f"[*] 配置加载完成,不限数量,注册线程={threads} mint_workers={mint_workers} mint_queue_max={mint_qmax} fast={fast}",
|
||
flush=True,
|
||
)
|
||
else:
|
||
remaining = max(0, args.count - done_count)
|
||
print(
|
||
f"[*] 配置加载完成,目标 {args.count} 个账号,注册线程={threads} "
|
||
f"mint_workers={mint_workers} mint_queue_max={mint_qmax} fast={fast}",
|
||
flush=True,
|
||
)
|
||
print(f"[*] accounts_file = {args.accounts_file}", flush=True)
|
||
if done_count > 0:
|
||
print(f"[*] 断点续跑:已完成 {done_count}", flush=True)
|
||
if remaining is not None and remaining <= 0:
|
||
print("[*] 所有账号已完成,无需继续(可用 --extra N 再注册)", flush=True)
|
||
return 0
|
||
|
||
# Surface display/headless decision early (before TabPool cold-start)
|
||
try:
|
||
will_hl = reg.should_use_headless(None, log_hint=True)
|
||
has_disp = reg.has_display()
|
||
print(
|
||
f"[*] browser mode: {'headless' if will_hl else 'headed'} "
|
||
f"(DISPLAY={os.environ.get('DISPLAY', '')!r} has_display={has_disp})",
|
||
flush=True,
|
||
)
|
||
# Keep cpa browser fallback aligned when config left cpa_headless unset/false
|
||
# and we are on a headless host — only auto-enable if user did not force headed.
|
||
if will_hl and not args.headed:
|
||
if "cpa_headless" not in cfg0 or cfg0.get("cpa_headless") is False:
|
||
# Auto-headless host: enable cpa headless unless explicitly headed CLI
|
||
if not has_disp:
|
||
cfg0["cpa_headless"] = True
|
||
if isinstance(getattr(reg, "config", None), dict):
|
||
reg.config["cpa_headless"] = True
|
||
print("[*] cpa_headless 自动=true(无 DISPLAY;协议 mint 仍优先纯 HTTP)", flush=True)
|
||
except Exception as exc:
|
||
print(f"[!] headless probe failed: {exc}", flush=True)
|
||
|
||
log_thread = threading.Thread(target=_log_writer, daemon=True)
|
||
log_thread.start()
|
||
|
||
try:
|
||
reg.TabPool.init(reg.create_browser_options, log_callback=lambda m: log(0, m))
|
||
except Exception as exc:
|
||
print(f"[!] 浏览器初始化失败: {exc}", flush=True)
|
||
if not reg.has_display():
|
||
print(
|
||
"[!] 无 DISPLAY 时若仍失败,请检查 Chrome/Chromium 是否安装,并尝试:\n"
|
||
" xvfb-run -a uv run python -u register_cli.py --extra 1 --threads 1",
|
||
flush=True,
|
||
)
|
||
return 1
|
||
|
||
task_queue: queue.Queue = queue.Queue()
|
||
mint_queue: queue.Queue | None = queue.Queue() if not do_mint_inline else None
|
||
if mint_queue is not None:
|
||
mint_queue._reg_qmax = mint_qmax # type: ignore[attr-defined]
|
||
global _next_idx
|
||
_next_idx[0] = done_count + 1
|
||
if remaining is not None:
|
||
for i in range(done_count + 1, args.count + 1):
|
||
task_queue.put(i)
|
||
else:
|
||
for i in range(done_count + 1, done_count + threads * 5 + 1):
|
||
task_queue.put(i)
|
||
_next_idx[0] = done_count + threads * 5 + 1
|
||
|
||
forever = remaining is None
|
||
cfg = getattr(reg, "config", {}) or {}
|
||
|
||
# mint workers first (so queue consumers ready)
|
||
mint_threads: list[threading.Thread] = []
|
||
if mint_queue is not None and mint_workers > 0:
|
||
for i in range(1, mint_workers + 1):
|
||
wid = f"M{i}"
|
||
t = threading.Thread(
|
||
target=_mint_worker,
|
||
args=(wid, mint_queue, cfg),
|
||
daemon=True,
|
||
name=f"mint-{i}",
|
||
)
|
||
t.start()
|
||
mint_threads.append(t)
|
||
|
||
reg_threads: list[threading.Thread] = []
|
||
for wid in range(1, threads + 1):
|
||
t = threading.Thread(
|
||
target=_register_worker,
|
||
args=(wid, task_queue, args.count, args.accounts_file, mint_queue, forever, do_mint_inline),
|
||
daemon=True,
|
||
name=f"reg-{wid}",
|
||
)
|
||
t.start()
|
||
reg_threads.append(t)
|
||
|
||
try:
|
||
for t in reg_threads:
|
||
t.join()
|
||
except KeyboardInterrupt:
|
||
print("\n[!] 用户中断", flush=True)
|
||
|
||
# drain mint queue
|
||
if mint_queue is not None:
|
||
log(0, f"[cpa] 等待 mint 队列清空(qsize≈{mint_queue.qsize()})...")
|
||
mint_queue.join()
|
||
for _ in mint_threads:
|
||
mint_queue.put(_MINT_STOP)
|
||
for t in mint_threads:
|
||
t.join(timeout=600)
|
||
|
||
try:
|
||
reg.shutdown_browser()
|
||
except Exception:
|
||
pass
|
||
|
||
# stop side-effect pool
|
||
try:
|
||
pool = getattr(reg, "_side_effect_pool", None)
|
||
if pool is not None:
|
||
pool.shutdown(wait=False, cancel_futures=True)
|
||
except Exception:
|
||
pass
|
||
|
||
_log_queue.put(None)
|
||
log_thread.join(timeout=2)
|
||
|
||
with _stats_lock:
|
||
s = dict(_stats)
|
||
print(
|
||
f"=== 完成: 注册成功 {s.get('reg_success', 0)}, 注册失败 {s.get('reg_fail', 0)}, "
|
||
f"CPA成功 {s.get('mint_success', 0)}, CPA失败 {s.get('mint_fail', 0)}, "
|
||
f"CPA跳过 {s.get('mint_skip', 0)} ===",
|
||
flush=True,
|
||
)
|
||
return 0 if s.get("reg_success", 0) > 0 else 1
|
||
|
||
|
||
if __name__ == "__main__":
|
||
sys.exit(main())
|