- 每个账号独立获取代理,不再共用一个代理IP - 代理预检失败只标记该账号失败,不终止整个批次 - 代理验证改为并发:verify_proxies_concurrent 多IP同时验证 - resolve_working_proxy 多IP并发验证,找到可用即返回 - parse_proxy_response 改为返回列表,支持多行/逗号分隔 - 邮箱验证码匹配收件人地址,防止多账号并发取错验证码 - 修正 mail.bdhg.xyz 旧数据端口和SSL配置 - bdhg.xyz 默认配置改为 143端口非SSL Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
262 lines
9.9 KiB
Python
262 lines
9.9 KiB
Python
"""登录服务:复用 core/ 核心模块,在线程池中并发执行登录并推送日志。"""
|
|
|
|
import asyncio
|
|
import threading
|
|
import uuid
|
|
from concurrent.futures import ThreadPoolExecutor, as_completed
|
|
from datetime import datetime, timezone
|
|
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:
|
|
"""批量登录执行器,在线程中运行,通过 ThreadPoolExecutor 并发登录多个账号。"""
|
|
|
|
def __init__(
|
|
self,
|
|
db: Session,
|
|
account_ids: list[int],
|
|
created_by: int,
|
|
creator_role: str,
|
|
max_geetest_retries: int = 5,
|
|
max_proxy_retries: int = 10,
|
|
proxy_config: Optional[ProxyConfigModel] = None,
|
|
log_queue: Optional[asyncio.Queue] = None,
|
|
loop: Optional[asyncio.AbstractEventLoop] = None,
|
|
concurrency: int = 3,
|
|
):
|
|
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.max_proxy_retries = max_proxy_retries
|
|
self.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._stop = threading.Event()
|
|
self._counter_lock = threading.Lock()
|
|
self._completed = 0
|
|
|
|
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 _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']}")
|
|
return
|
|
|
|
worker_db = SessionLocal()
|
|
try:
|
|
task = worker_db.query(LoginTask).filter(LoginTask.id == task_id).first()
|
|
if not task:
|
|
return
|
|
|
|
task.status = "running"
|
|
worker_db.commit()
|
|
|
|
with self._counter_lock:
|
|
self._completed += 1
|
|
current = self._completed
|
|
|
|
self._push_log("info", f"[{current}/{total}] 开始登录: {acc_info['username']}")
|
|
|
|
# 每个账号独立获取代理
|
|
proxy_dict, proxy_msg = self._resolve_proxy()
|
|
if proxy_msg:
|
|
self._push_log("info", f"[{current}] {proxy_msg}")
|
|
|
|
# 代理预检失败时,该账号标记失败但不终止整个批次
|
|
if self.proxy_config and self.proxy_config.enabled and not proxy_dict:
|
|
task.status = "error"
|
|
task.message = f"代理不可用: {proxy_msg}"
|
|
task.finished_at = datetime.now(timezone.utc)
|
|
worker_db.commit()
|
|
self._push_log("error", f"[{current}] {acc_info['username']} 代理不可用: {proxy_msg}")
|
|
return
|
|
|
|
try:
|
|
account = Account(
|
|
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_geetest_retries=self.max_geetest_retries,
|
|
max_proxy_retries=self.max_proxy_retries,
|
|
)
|
|
result = loginer.login()
|
|
|
|
if result.success:
|
|
task.status = "success"
|
|
task.cookie = result.cookie
|
|
task.message = "登录成功"
|
|
self._push_log("success", f"[{current}] {acc_info['username']} 登录成功")
|
|
else:
|
|
task.status = "failed"
|
|
task.message = result.message
|
|
self._push_log("error", f"[{current}] {acc_info['username']} 登录失败: {result.message}")
|
|
|
|
except Exception as e:
|
|
task.status = "error"
|
|
task.message = str(e)
|
|
self._push_log("error", f"[{current}] {acc_info['username']} 登录异常: {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
|
|
self._push_log("info", f"批量登录任务 {batch_id} 开始,共 {len(self.account_ids)} 个账号,并发数: {concurrency}")
|
|
|
|
# 创建或复用任务记录(顺序执行,线程安全)
|
|
task_infos: list[dict] = [] # {task_id, acc_info}
|
|
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
|
|
|
|
# 复用该账号最近一条失败任务记录,避免重复产生多条
|
|
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.batch_id = batch_id
|
|
existing_task.status = "pending"
|
|
existing_task.cookie = ""
|
|
existing_task.message = ""
|
|
existing_task.finished_at = None
|
|
task = existing_task
|
|
else:
|
|
task = LoginTask(
|
|
batch_id=batch_id,
|
|
account_id=aid,
|
|
status="pending",
|
|
created_by=self.created_by,
|
|
)
|
|
self.db.add(task)
|
|
|
|
self.db.flush() # 获取 task.id
|
|
|
|
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,
|
|
},
|
|
})
|
|
|
|
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 in task_infos:
|
|
if self._stop.is_set():
|
|
self._push_log("warning", "任务已停止,跳过剩余账号")
|
|
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"批量登录任务 {batch_id} 完成")
|
|
self._push_log("result", "")
|
|
|
|
|
|
# 在模块末尾导入 SessionLocal(避免循环导入)
|
|
from ..database import SessionLocal
|