"""登录服务:复用 core/ 核心模块,在线程池中执行登录并推送日志。""" import asyncio import threading import uuid from datetime import datetime from typing import Optional from sqlalchemy.orm import Session from core.douyu import DouyuLogin from core.models import Account, ProxyConfig as DouyuProxyConfig from core.douyu.proxy import resolve_working_proxy from ..models import Account as AccountModel, LoginTask, ProxyConfig as ProxyConfigModel from ..permissions import has_permission class LoginBatchRunner: """批量登录执行器,在线程中运行,通过 asyncio.Queue 推送日志。""" def __init__( self, db: Session, account_ids: list[int], created_by: int, creator_role: str, max_geetest_retries: int = 5, proxy_config: Optional[ProxyConfigModel] = None, log_queue: Optional[asyncio.Queue] = None, loop: Optional[asyncio.AbstractEventLoop] = None, ): self.db = db self.account_ids = account_ids self.created_by = created_by self.creator_role = creator_role self.max_geetest_retries = max_geetest_retries self.proxy_config = proxy_config self.log_queue = log_queue self.loop = loop self.batch_id = uuid.uuid4().hex[:12] self._stop = threading.Event() def stop(self): self._stop.set() def _push_log(self, level: str, message: str): if self.log_queue and self.loop: asyncio.run_coroutine_threadsafe( self.log_queue.put({"level": level, "message": message}), self.loop, ) def _resolve_proxy(self) -> tuple[Optional[dict], str]: """ 解析代理配置,返回 (proxy_dict, message)。 - 静态代理:直接返回 dict - API代理:调用 resolve_working_proxy 预检,自动同步白名单 - 无代理:返回 (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}, f'使用静态代理: {proxy_url}' # API代理:预检获取可用代理 if self.proxy_config.api_url: whitelist_uid = '' whitelist_ukey = '' if self.proxy_config.whitelist_enabled: whitelist_uid = self.proxy_config.whitelist_uid or '' whitelist_ukey = self.proxy_config.whitelist_ukey or '' proxy_url, msg = resolve_working_proxy( api_url=self.proxy_config.api_url, whitelist_uid=whitelist_uid, whitelist_ukey=whitelist_ukey, log_func=self._push_log, ) if proxy_url: return {'http': proxy_url, 'https': proxy_url}, msg return None, msg return None, '' def run(self): """在线程中执行批量登录。""" batch_id = self.batch_id self._push_log("info", f"批量登录任务 {batch_id} 开始,共 {len(self.account_ids)} 个账号") # 创建任务记录 tasks = [] for aid in self.account_ids: acc = self.db.query(AccountModel).filter(AccountModel.id == aid).first() if not acc: continue # 权限检查:客服只能跑分配给自己的 if not has_permission(self.creator_role, "login:view_all"): if acc.assigned_to != self.created_by: self._push_log("warning", f"跳过无权账号: {acc.username}") continue task = LoginTask( batch_id=batch_id, account_id=aid, status="pending", created_by=self.created_by, ) self.db.add(task) tasks.append((task, acc)) self.db.commit() # 代理预检 proxy_dict, proxy_msg = self._resolve_proxy() if proxy_msg: self._push_log("info", proxy_msg) # 如果启用了代理但预检失败,终止任务 if self.proxy_config and self.proxy_config.enabled and not proxy_dict: self._push_log("error", f"代理不可用,任务终止: {proxy_msg}") for task, _ in tasks: task.status = "error" task.message = f"代理不可用: {proxy_msg}" task.finished_at = datetime.utcnow() self.db.commit() return for i, (task, acc) in enumerate(tasks): if self._stop.is_set(): self._push_log("warning", "任务已停止") break task.status = "running" self.db.commit() self._push_log("info", f"[{i+1}/{len(tasks)}] 开始登录: {acc.username}") try: account = Account( 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, ) loginer = DouyuLogin( account, proxy=proxy_dict, max_geetest_retries=self.max_geetest_retries, ) result = loginer.login() if result.success: task.status = "success" task.cookie = result.cookie task.message = "登录成功" self._push_log("success", f"[{i+1}] {acc.username} 登录成功") else: task.status = "failed" task.message = result.message self._push_log("error", f"[{i+1}] {acc.username} 登录失败: {result.message}") except Exception as e: task.status = "error" task.message = str(e) self._push_log("error", f"[{i+1}] {acc.username} 登录异常: {e}") task.finished_at = datetime.utcnow() self.db.commit() self._push_log("info", f"批量登录任务 {batch_id} 完成") self._push_log("result", "")