"""登录服务:复用 core/ 核心模块,在线程池中并发执行登录并推送日志。""" import asyncio import os import threading import time import uuid from types import SimpleNamespace from concurrent.futures import ThreadPoolExecutor, as_completed from datetime import datetime, timezone from typing import Optional from sqlalchemy.orm import Session from loguru import logger from core.douyu import DouyuLogin, WgapiLoginAPI, IframeLoginAPI from core.douyu.proxy_fetcher import ProxyFetcher from ..models import Account as AccountModel, LoginTask, ProxyConfig as ProxyConfigModel from .cookie_check_service import check_douyu_cookie def _create_api_strategy(strategy_name: str): """根据名称创建登录接口策略实例。""" if strategy_name == "iframe": return IframeLoginAPI() return WgapiLoginAPI() CHECK_STATUS_MESSAGES = { "account_cancelled": "账号已注销", "password_wrong": "账号密码错误", "account_unverified": "账号未认证", "account_verified": "账号已认证", "account_auth_unknown": "账号认证状态未知", } CHECK_STATUS_LOG_LEVELS = { "account_cancelled": "warning", "password_wrong": "error", "account_unverified": "warning", "account_verified": "success", "account_auth_unknown": "warning", } def _positive_env_int(name: str, default: int) -> int: """读取正整数环境变量,非法值回退到默认值。""" try: value = int(os.getenv(name, str(default))) except ValueError: return default return value if value > 0 else default def get_relogin_limits() -> tuple[int, int]: """读取客服 CK 重登的重试与总时长限制。""" return ( _positive_env_int("COOKIE_RELOGIN_MAX_RETRIES", 5), _positive_env_int("COOKIE_RELOGIN_MAX_TOTAL_TIME", 600), ) def _snapshot_proxy_config(proxy_config: Optional[ProxyConfigModel]) -> Optional[SimpleNamespace]: """复制代理配置,避免后台线程访问已关闭会话中的 ORM 对象。""" if proxy_config is None: return None credentials = getattr(proxy_config, "whitelist_credentials", None) return SimpleNamespace( enabled=bool(getattr(proxy_config, "enabled", False)), http=getattr(proxy_config, "http", "") or "", https=getattr(proxy_config, "https", "") or "", api_url=getattr(proxy_config, "api_url", "") or "", whitelist_enabled=bool(getattr(proxy_config, "whitelist_enabled", False)), whitelist_platform=getattr(proxy_config, "whitelist_platform", None), whitelist_credentials=dict(credentials) if isinstance(credentials, dict) else credentials, whitelist_uid=getattr(proxy_config, "whitelist_uid", "") or "", whitelist_ukey=getattr(proxy_config, "whitelist_ukey", "") or "", ) class LoginBatchRunner: """批量登录执行器,在线程中运行,通过 ThreadPoolExecutor 并发登录多个账号。""" def __init__( self, db: Session, account_ids: list[int], created_by: int, creator_permissions: list[str], max_login_retries: int = 0, max_total_time: float = 0, proxy_config: Optional[ProxyConfigModel] = None, log_queue: Optional[asyncio.Queue] = None, loop: Optional[asyncio.AbstractEventLoop] = None, concurrency: int = 3, api_strategy: str = "wgapi", mode: str = "login", relogin_task_ids: Optional[list[int]] = None, ): self.db = db self.account_ids = account_ids self.created_by = created_by self.creator_permissions = creator_permissions # 限制最大重试次数,避免无限重试导致资源耗尽 # 0 表示使用默认值 20,其他值保持原样 self.max_login_retries = max_login_retries if max_login_retries > 0 else 20 self.max_total_time = max_total_time self.proxy_config = _snapshot_proxy_config(proxy_config) self.log_queue = log_queue self.loop = loop self.batch_id = uuid.uuid4().hex[:12] self.concurrency = max(1, min(concurrency, 10)) # 限制 1-10 self.api_strategy = _create_api_strategy(api_strategy) self.mode = mode if mode in ("login", "check", "relogin") else "login" self.relogin_task_ids = list(relogin_task_ids) if relogin_task_ids else [] self._stop = threading.Event() self._counter_lock = threading.Lock() self._completed = 0 # 共享代理获取器(无池,每次取新代理) self._shared_proxy_fetcher = None proxy_config = self.proxy_config if proxy_config and proxy_config.enabled and proxy_config.api_url: wl_platform = "xiequ" wl_credentials = None if proxy_config.whitelist_enabled: wl_platform = getattr(proxy_config, 'whitelist_platform', None) or "xiequ" wl_credentials = getattr(proxy_config, 'whitelist_credentials', None) # 向后兼容 if not wl_credentials and proxy_config.whitelist_uid and proxy_config.whitelist_ukey: wl_platform = "xiequ" wl_credentials = {"uid": proxy_config.whitelist_uid, "ukey": proxy_config.whitelist_ukey} self._shared_proxy_fetcher = ProxyFetcher( api_url=proxy_config.api_url, whitelist_platform=wl_platform, whitelist_credentials=wl_credentials, stop_event=self._stop, ) def stop(self): self._stop.set() def _sleep_or_stop(self, seconds: float) -> bool: """可中断等待;返回 True 表示收到停止信号。""" deadline = time.monotonic() + seconds while time.monotonic() < deadline: if self._stop.is_set(): return True time.sleep(min(0.2, deadline - time.monotonic())) return self._stop.is_set() def _mark_relogin_stopped( self, task_id: int, message: str = "重新登录已停止,旧 Cookie 已保留", ) -> None: """将未开始或已中断的重登任务收敛为保留旧 CK 的终态。""" worker_db = SessionLocal() try: task = worker_db.query(LoginTask).filter(LoginTask.id == task_id).first() if task and task.status in ("relogin_pending", "relogin_running"): task.status = "relogin_failed" task.message = message task.finished_at = datetime.now(timezone.utc) worker_db.commit() finally: worker_db.close() def _push_log(self, level: str, message: str): # 即使没有页面实时日志,也要保留批次进度到 app.log,便于排查卡点。 if message: log_level = level if level in {"debug", "info", "warning", "error", "success"} else "debug" getattr(logger, log_level)(f"[登录批次 {self.batch_id}] {message}") if self.log_queue and self.loop: asyncio.run_coroutine_threadsafe( self.log_queue.put({"level": level, "message": message}), self.loop, ) def _resolve_static_proxy(self) -> tuple[Optional[dict], str]: """解析静态代理配置。""" 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}, f'使用静态代理: {proxy_url}' # API代理:由 DouyuLogin 通过 proxy_fetcher 内部管理 return None, '' def _execute_one(self, task_id: int, acc_info: dict, total: int): """在独立线程中执行单个账号登录,使用独立的 DB 会话。""" if self._stop.is_set(): self._push_log("warning", f"任务已停止,跳过: {acc_info['username']}") if self.mode == "relogin": self._mark_relogin_stopped(task_id) return worker_db = SessionLocal() try: task = worker_db.query(LoginTask).filter(LoginTask.id == task_id).first() if not task: return if self.mode == "relogin": task.status = "relogin_running" task.message = "正在重新登录,旧 Cookie 保留中" else: task.status = "running" worker_db.commit() with self._counter_lock: self._completed += 1 current = self._completed action_name = "检测" if self.mode == "check" else "重新登录" if self.mode == "relogin" else "登录" self._push_log("info", f"[{current}/{total}] 开始{action_name}: {acc_info['username']}") try: # 代理配置也可能异常,必须由当前任务的失败处理收敛状态。 proxy_dict, proxy_msg = self._resolve_static_proxy() if proxy_msg: self._push_log("info", f"[{current}] {proxy_msg}") # 静态代理启用但配置为空 → 不可用 if self.proxy_config and self.proxy_config.enabled and not (self.proxy_config.http or self.proxy_config.https) and not self._shared_proxy_fetcher and not proxy_dict: if self.mode == "relogin": task.status = "relogin_failed" task.message = "重新登录失败: 代理不可用: 未配置代理(旧 Cookie 已保留)" else: task.status = "error" task.message = "代理不可用: 未配置代理" task.finished_at = datetime.now(timezone.utc) worker_db.commit() self._push_log("error", f"[{current}] {acc_info['username']} 代理不可用") return account = SimpleNamespace( username=acc_info["username"], password=acc_info["password"], email=acc_info["email"], email_password=acc_info["email_password"], email_imap_server=acc_info["email_imap_server"] or "", email_imap_port=acc_info["email_imap_port"] or 993, email_imap_ssl=acc_info["email_imap_ssl"], ) loginer = DouyuLogin( account, proxy=proxy_dict, max_login_retries=self.max_login_retries, max_total_time=self.max_total_time, proxy_fetcher=self._shared_proxy_fetcher, stop_event=self._stop, api_strategy=self.api_strategy, ) result = loginer.check_account() if self.mode == "check" else loginer.login() if self.mode == "check" and result.success: status = result.code if result.code in CHECK_STATUS_MESSAGES else "account_auth_unknown" task.status = status task.cookie = "" task.message = result.message or CHECK_STATUS_MESSAGES[status] level = CHECK_STATUS_LOG_LEVELS.get(status, "info") self._push_log(level, f"[{current}] {acc_info['username']} 检测结果: {task.message}") elif result.success: task.status = "success" task.cookie = result.cookie task.message = result.message or "登录成功" if self.mode == "relogin": check_result = check_douyu_cookie(result.cookie) task.ck_check_status = "valid" if check_result["valid"] else "invalid" task.ck_check_result = { "fish_ball": check_result["fish_ball"], "nickname": check_result["nickname"], "level": check_result["level"], "message": check_result["message"], } task.ck_checked_at = check_result["checked_at"] if check_result["valid"]: task.message = "重新登录成功,Cookie 有效" self._push_log("success", f"[{current}] {acc_info['username']} 重新登录成功,Cookie 已替换并验证有效") else: task.message = f"重新登录成功,但 Cookie 有效性检测失败: {check_result['message']}" self._push_log("warning", f"[{current}] {acc_info['username']} 重新登录成功,但 Cookie 有效性检测失败: {check_result['message']}") else: self._push_log("success", f"[{current}] {acc_info['username']} {task.message}") else: if self.mode == "relogin": # 重新登录失败时保留旧 Cookie 与成功状态,仅记录失败原因,行不消失 task.status = "relogin_failed" task.message = f"重新登录失败: {result.message}(旧 Cookie 已保留)" self._push_log("error", f"[{current}] {acc_info['username']} 重新登录失败: {result.message}") else: task.status = "failed" task.message = result.message self._push_log("error", f"[{current}] {acc_info['username']} {action_name}失败: {result.message}") except Exception as e: if self.mode == "relogin": task.status = "relogin_failed" task.message = f"重新登录异常: {e}(旧 Cookie 已保留)" self._push_log("error", f"[{current}] {acc_info['username']} 重新登录异常: {e}") else: task.status = "error" task.message = str(e) self._push_log("error", f"[{current}] {acc_info['username']} {action_name}异常: {e}") task.finished_at = datetime.now(timezone.utc) worker_db.commit() finally: worker_db.close() def run(self): """在线程中执行批量登录。""" batch_id = self.batch_id concurrency = self.concurrency action_name = "账号检测" if self.mode == "check" else "重新登录" if self.mode == "relogin" else "登录" self._push_log("info", f"批量{action_name}任务 {batch_id} 开始,共 {len(self.account_ids or self.relogin_task_ids)} 个账号,并发数: {concurrency}") # 批次开始前同步一次出口 IP 到白名单,后续 fetch_new_proxy 不再主动同步 if self._shared_proxy_fetcher: ok, msg = self._shared_proxy_fetcher.warmup_whitelist() if msg != "无白名单凭据,跳过": self._push_log("info" if ok else "warning", f"白名单预热: {msg}") def _append_task_info(task: LoginTask, acc: AccountModel): task.batch_id = batch_id if self.mode == "relogin": # 旧 Cookie 仅在 result.success 后由 _execute_one 原子替换。 task.status = "relogin_pending" task.message = "等待重新登录,旧 Cookie 已保留" else: task.status = "pending" task.message = "" task.finished_at = None self.db.flush() task_infos.append({ "task_id": task.id, "acc_info": { "username": acc.username, "password": acc.password, "email": acc.email, "email_password": acc.email_password, "email_imap_server": acc.email_imap_server or "", "email_imap_port": acc.email_imap_port or 993, "email_imap_ssl": acc.email_imap_ssl if acc.email_imap_ssl is not None else True, }, }) try: # 创建或复用任务记录(顺序执行,线程安全) task_infos: list[dict] = [] # {task_id, acc_info} if self.relogin_task_ids: # 重新登录模式:复用指定 Cookie 记录,登录成功后原地替换 Cookie for task_id in self.relogin_task_ids: task = self.db.query(LoginTask).filter(LoginTask.id == task_id).first() if not task: continue acc = self.db.query(AccountModel).filter(AccountModel.id == task.account_id).first() if not acc: self._push_log("warning", f"跳过无账号的任务 #{task_id}") continue if "login:view_all" not in self.creator_permissions: if acc.assigned_to != self.created_by: self._push_log("warning", f"跳过无权账号: {acc.username}") continue _append_task_info(task, acc) else: seen_account_ids = set() for aid in self.account_ids: if aid in seen_account_ids: continue seen_account_ids.add(aid) acc = ( self.db.query(AccountModel) .filter(AccountModel.id == aid) .with_for_update() .first() ) if not acc: continue # 权限检查:客服只能跑分配给自己的 if "login:view_all" not in self.creator_permissions: if acc.assigned_to != self.created_by: self._push_log("warning", f"跳过无权账号: {acc.username}") continue # 一个斗鱼账号只保留一条成功 CK:再次普通登录时更新最新成功记录。 latest_success_task = ( self.db.query(LoginTask) .filter(LoginTask.account_id == aid, LoginTask.status == "success") .order_by(LoginTask.finished_at.desc(), LoginTask.id.desc()) .first() ) if latest_success_task: duplicate_success_tasks = ( self.db.query(LoginTask) .filter( LoginTask.account_id == aid, LoginTask.status == "success", LoginTask.id != latest_success_task.id, ) .all() ) for duplicate_task in duplicate_success_tasks: self.db.delete(duplicate_task) task = latest_success_task else: # 复用该账号最近一条失败任务记录,避免重复产生多条失败历史。 existing_task = ( self.db.query(LoginTask) .filter(LoginTask.account_id == aid, LoginTask.status.in_(["failed", "error"])) .order_by(LoginTask.id.desc()) .first() ) if existing_task: existing_task.cookie = "" task = existing_task else: task = LoginTask( batch_id=batch_id, account_id=aid, status="pending", created_by=self.created_by, ) self.db.add(task) _append_task_info(task, acc) self.db.commit() total = len(task_infos) if total == 0: self._push_log("warning", "没有可执行的账号") self._push_log("result", "") return # 并发执行登录,每个账号独立获取代理 with ThreadPoolExecutor(max_workers=concurrency) as executor: futures = [] for item_index, item in enumerate(task_infos): if self._stop.is_set(): self._push_log("warning", "任务已停止,跳过剩余账号") if self.mode == "relogin": for pending_item in task_infos[item_index:]: self._mark_relogin_stopped(pending_item["task_id"]) break future = executor.submit( self._execute_one, item["task_id"], item["acc_info"], total, ) futures.append(future) # 等待所有任务完成 for future in as_completed(futures): try: future.result() except Exception as e: self._push_log("error", f"Worker 异常: {e}") self._push_log("info", f"批量{action_name}任务 {batch_id} 完成") self._push_log("result", "") finally: # 标记批次结束:WS 端据此决定何时清理注册表(断线重连可继续订阅日志)。 batch_registry.mark_finished(batch_id) # 确保 DB Session 被关闭,避免连接泄漏 self.db.close() class BatchRegistry: """管理运行中的登录批次状态。""" def __init__(self): self._batches: dict[str, dict] = {} def register(self, batch_id: str, log_queue: Optional[asyncio.Queue], loop: Optional[asyncio.AbstractEventLoop], runner: LoginBatchRunner, owner_id: Optional[int] = None): self._batches[batch_id] = { "log_queue": log_queue, "loop": loop, "runner": runner, "owner_id": owner_id, "finished": False, "finished_at": None, } def get(self, batch_id: str): return self._batches.get(batch_id) def mark_finished(self, batch_id: str): """标记批次已结束(幂等;不在本注册表的批次为无操作)。 与虎牙批次一致:WS 端只在 finished 后 pop,客户端断线重连仍可订阅 到运行中批次的实时日志。 """ batch = self._batches.get(batch_id) if batch: batch["finished"] = True batch["finished_at"] = time.time() def pop(self, batch_id: str): return self._batches.pop(batch_id, None) # 模块级单例 batch_registry = BatchRegistry() def cleanup_orphan_relogin_tasks( db: Session, message: str = "重新登录已中断(服务重启),旧 Cookie 已保留", ) -> int: """服务重启后收敛遗留重登状态,避免页面永久显示重登中。""" tasks = ( db.query(LoginTask) .filter(LoginTask.status.in_(("relogin_pending", "relogin_running"))) .all() ) if not tasks: return 0 finished_at = datetime.now(timezone.utc) for task in tasks: task.status = "relogin_failed" task.message = message task.finished_at = finished_at db.commit() return len(tasks) # 在模块末尾导入 SessionLocal(避免循环导入) from ..database import SessionLocal