"""虎牙自动注册批次执行器。""" from __future__ import annotations import threading import uuid from concurrent.futures import ThreadPoolExecutor, as_completed from dataclasses import dataclass, field from datetime import datetime, timezone from typing import Optional from core.douyu.proxy_fetcher import ProxyFetcher from core.huya.auto_register import HuyaAutoRegisterResult, register_huya_with_sms_line from core.huya.cookie_utils import normalize_huya_cookie from core.sms_provider import SmsLine from ..database import SessionLocal from ..models import ProxyConfig as ProxyConfigModel from .huya_service import upsert_huya_cookie def _now() -> datetime: return datetime.now(timezone.utc) def _cookie_preview(cookie: str) -> str: normalized = normalize_huya_cookie(cookie or "") if not normalized: return "" return normalized[:50] + "..." if len(normalized) > 50 else normalized @dataclass class HuyaRegisterItemState: """单个手机号在批次中的状态。""" line: int phone: str provider: str sms_url: str = "" status: str = "pending" message: str = "等待开始" code: str = "" change_code: str = "" attempts: int = 0 change_attempts: int = 0 account_id: int | None = None username: str = "" uid: str = "" password: str = "" password_changed: bool = False cookie: str = "" cookie_preview: str = "" started_at: datetime | None = None finished_at: datetime | None = None def to_dict(self) -> dict: return { "line": self.line, "phone": self.phone, "provider": self.provider, "sms_url": self.sms_url, "status": self.status, "message": self.message, "code": self.code, "change_code": self.change_code, "attempts": self.attempts, "change_attempts": self.change_attempts, "account_id": self.account_id, "username": self.username, "uid": self.uid, "password": self.password, "password_changed": self.password_changed, "cookie": self.cookie, "cookie_preview": self.cookie_preview, "started_at": self.started_at, "finished_at": self.finished_at, } @dataclass class HuyaRegisterBatch: """自动注册批次内存快照。""" batch_id: str tag: str created_by: int concurrency: int wait_seconds: float poll_interval: float items: list[HuyaRegisterItemState] password_prefix: str = "hy" fixed_password: str = "" use_proxy: bool = False status: str = "pending" message: str = "等待开始" created_at: datetime = field(default_factory=_now) started_at: datetime | None = None finished_at: datetime | None = None class HuyaRegisterRunner: """在后台线程中批量执行虎牙手机号自动注册。""" def __init__( self, batch: HuyaRegisterBatch, sms_lines: list[SmsLine], proxy_config: Optional[ProxyConfigModel] = None, ): self.batch = batch self.sms_lines = sms_lines self.proxy_config = proxy_config self._lock = threading.Lock() self._stop = threading.Event() self._shared_proxy_fetcher = self._create_proxy_fetcher() def _create_proxy_fetcher(self) -> ProxyFetcher | None: """按需创建 API 代理获取器。""" if not self.batch.use_proxy or not self.proxy_config or not self.proxy_config.enabled: return None if not self.proxy_config.api_url: return None wl_platform = "xiequ" wl_credentials = None if self.proxy_config.whitelist_enabled: wl_platform = getattr(self.proxy_config, "whitelist_platform", None) or "xiequ" wl_credentials = getattr(self.proxy_config, "whitelist_credentials", None) if not wl_credentials and self.proxy_config.whitelist_uid and self.proxy_config.whitelist_ukey: wl_credentials = { "uid": self.proxy_config.whitelist_uid, "ukey": self.proxy_config.whitelist_ukey, } return ProxyFetcher( api_url=self.proxy_config.api_url, whitelist_platform=wl_platform, whitelist_credentials=wl_credentials, stop_event=self._stop, ) def stop(self): self._stop.set() with self._lock: if self.batch.status == "running": self.batch.message = "正在停止" def snapshot(self) -> dict: with self._lock: total = len(self.batch.items) success = sum(1 for item in self.batch.items if item.status == "success") failed = sum(1 for item in self.batch.items if item.status == "error") stopped = sum(1 for item in self.batch.items if item.status == "stopped") running = sum(1 for item in self.batch.items if item.status in {"sending", "waiting", "changing", "logging"}) return { "batch_id": self.batch.batch_id, "status": self.batch.status, "message": self.batch.message, "tag": self.batch.tag, "created_by": self.batch.created_by, "concurrency": self.batch.concurrency, "wait_seconds": self.batch.wait_seconds, "poll_interval": self.batch.poll_interval, "password_prefix": self.batch.password_prefix, "use_proxy": self.batch.use_proxy, "total": total, "success_count": success, "failed_count": failed, "stopped_count": stopped, "running_count": running, "created_at": self.batch.created_at, "started_at": self.batch.started_at, "finished_at": self.batch.finished_at, "items": [item.to_dict() for item in self.batch.items], } def _set_item(self, index: int, **updates): with self._lock: item = self.batch.items[index] for key, value in updates.items(): setattr(item, key, value) def _save_cookie(self, result: HuyaAutoRegisterResult) -> tuple[int | None, str, str]: db = SessionLocal() try: account = upsert_huya_cookie(db, result.cookie, tag=self.batch.tag, username_hint="") account.game_phone = result.phone if result.username: account.username = result.username if result.password: account.account_password = result.password if result.password_changed: account.status = "password_changed" account.updated_at = _now() db.commit() db.refresh(account) return account.id, account.username or "", account.uid or account.yyuid or "" finally: db.close() def _resolve_proxy(self) -> tuple[dict[str, str] | None, str]: """为单个手机号解析代理;返回代理字典和错误消息。""" if not self.batch.use_proxy: return None, "" if not self.proxy_config or not self.proxy_config.enabled: return None, "已开启代理,但代理配置未启用" if self.proxy_config.http or self.proxy_config.https: proxy_url = self.proxy_config.http or self.proxy_config.https return {"http": proxy_url, "https": proxy_url}, "" if self._shared_proxy_fetcher: proxy_url = self._shared_proxy_fetcher.fetch_new_proxy(max_attempts=3) if proxy_url: return {"http": proxy_url, "https": proxy_url}, "" return None, "获取代理失败" return None, "已开启代理,但未配置静态代理或代理 API" def _run_one(self, index: int, item: SmsLine): if self._stop.is_set(): self._set_item(index, status="stopped", message="已停止", finished_at=_now()) return proxies, proxy_error = self._resolve_proxy() if proxy_error: self._set_item(index, status="error", message=proxy_error, finished_at=_now()) return self._set_item(index, status="sending", message="注册并改密", started_at=_now(), finished_at=None) result = register_huya_with_sms_line( item, wait_seconds=self.batch.wait_seconds, poll_interval=self.batch.poll_interval, password_prefix=self.batch.password_prefix, fixed_password=self.batch.fixed_password, proxies=proxies, stop_event=self._stop, ) account_id = None username = result.username uid = result.uid message = result.message status = result.status cookie = result.cookie if result.success else "" if result.success: self._set_item( index, status="logging", message="保存账号密码", code=result.code, change_code=result.change_code, attempts=result.attempts, change_attempts=result.change_attempts, username=result.username, uid=result.uid, password=result.password, password_changed=result.password_changed, ) try: account_id, username, uid = self._save_cookie(result) except Exception as exc: status = "error" cookie = "" message = f"账号保存失败: {exc}" exposed_cookie = "" if result.password_changed else cookie self._set_item( index, status=status, message=message, code=result.code, change_code=result.change_code, attempts=result.attempts, change_attempts=result.change_attempts, account_id=account_id, username=username, uid=uid, password=result.password, password_changed=result.password_changed, cookie=normalize_huya_cookie(exposed_cookie), cookie_preview=_cookie_preview(exposed_cookie), finished_at=_now(), ) def run(self): """线程入口。""" with self._lock: self.batch.status = "running" self.batch.message = "批次运行中" self.batch.started_at = _now() try: if self._shared_proxy_fetcher: ok, msg = self._shared_proxy_fetcher.warmup_whitelist() if not ok: with self._lock: self.batch.message = f"代理白名单预热失败: {msg}" with ThreadPoolExecutor(max_workers=self.batch.concurrency) as executor: futures = [] for index, item in enumerate(self.sms_lines): if self._stop.is_set(): self._set_item(index, status="stopped", message="已停止", finished_at=_now()) continue futures.append(executor.submit(self._run_one, index, item)) for future in as_completed(futures): future.result() except Exception as exc: with self._lock: self.batch.status = "error" self.batch.message = f"批次执行异常: {exc}" self.batch.finished_at = _now() return with self._lock: if self._stop.is_set(): self.batch.status = "stopped" self.batch.message = "批次已停止" else: self.batch.status = "finished" self.batch.message = "批次已完成" self.batch.finished_at = _now() class HuyaRegisterRegistry: """管理自动注册批次。""" def __init__(self): self._lock = threading.Lock() self._runners: dict[str, HuyaRegisterRunner] = {} def create( self, sms_lines: list[SmsLine], tag: str, created_by: int, concurrency: int, wait_seconds: float, poll_interval: float, password_prefix: str = "hy", fixed_password: str = "", use_proxy: bool = False, proxy_config: Optional[ProxyConfigModel] = None, ) -> HuyaRegisterRunner: batch_id = uuid.uuid4().hex[:12] batch = HuyaRegisterBatch( batch_id=batch_id, tag=tag, created_by=created_by, concurrency=max(1, min(int(concurrency or 1), 5)), wait_seconds=max(15.0, float(wait_seconds or 180)), poll_interval=max(1.0, float(poll_interval or 5)), password_prefix=(password_prefix or "hy").strip()[:8] or "hy", fixed_password=(fixed_password or "").strip(), use_proxy=bool(use_proxy), items=[ HuyaRegisterItemState(line=index + 1, phone=item.phone, provider=item.provider, sms_url=item.url) for index, item in enumerate(sms_lines) ], ) runner = HuyaRegisterRunner(batch=batch, sms_lines=sms_lines, proxy_config=proxy_config) with self._lock: self._runners[batch_id] = runner return runner def get(self, batch_id: str) -> HuyaRegisterRunner | None: with self._lock: return self._runners.get(batch_id) huya_register_registry = HuyaRegisterRegistry()