"""虎牙自动注册批次执行器。""" 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 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 .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 = "" 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], ): self.batch = batch self.sms_lines = sms_lines self._lock = threading.Lock() self._stop = threading.Event() 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, "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 _run_one(self, index: int, item: SmsLine): if self._stop.is_set(): self._set_item(index, status="stopped", message="已停止", 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, 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: 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 = "", ) -> 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(), 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) 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()