安全修复: - WebSocket 端点添加认证(cookie/token),防止未授权窃听日志 - SPA serve_spa 添加路径遍历防护(resolve + relative_to 检查) - Token 改用 httpOnly Cookie 存储,移除前端 localStorage token(防 XSS 窃取) - 添加安全响应头中间件(X-Content-Type-Options/X-Frame-Options/Referrer-Policy) - HTTP 请求日志脱敏请求体中的 password/secret/token 等敏感字段 - 权限检查统一使用 user_has_permission(考虑自定义权限,修复 has_permission 忽略 custom_permissions 的缺陷) 性能与稳定性: - cookies.py 列表接口修复 N+1 查询(改为批量查询 Account) - login_service.py run() 结束时关闭 DB Session(防止连接泄漏) - _active_batches/_active_tests 全局字典添加 threading.Lock(防止并发竞态) 配置优化: - CORS 源支持环境变量 CORS_ORIGINS 配置 - Uvicorn reload 支持环境变量 UVICORN_RELOAD 控制(生产环境默认关闭) - Cookie 安全标志支持环境变量 COOKIE_SECURE 配置(HTTPS 部署时启用) - logs.py 权限不足返回 HTTP 403(原来返回 200 + message)
266 lines
11 KiB
Python
266 lines
11 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, get_proxy_manager
|
|
from ..models import Account as AccountModel, LoginTask, ProxyConfig as ProxyConfigModel
|
|
|
|
|
|
class LoginBatchRunner:
|
|
"""批量登录执行器,在线程中运行,通过 ThreadPoolExecutor 并发登录多个账号。"""
|
|
|
|
def __init__(
|
|
self,
|
|
db: Session,
|
|
account_ids: list[int],
|
|
created_by: int,
|
|
creator_permissions: list[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_permissions = creator_permissions
|
|
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
|
|
|
|
# 共享代理管理器(带锁,避免并发白名单限流;极验失败时可刷新代理)
|
|
self._shared_proxy_manager = None
|
|
if proxy_config and proxy_config.enabled and proxy_config.api_url:
|
|
wl_uid = proxy_config.whitelist_uid or "" if proxy_config.whitelist_enabled else ""
|
|
wl_ukey = proxy_config.whitelist_ukey or "" if proxy_config.whitelist_enabled else ""
|
|
self._shared_proxy_manager = get_proxy_manager(
|
|
proxy_config.api_url,
|
|
whitelist_uid=wl_uid,
|
|
whitelist_ukey=wl_ukey,
|
|
)
|
|
|
|
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代理:通过共享 ProxyManager(带已验证代理池缓存)获取,自动同步白名单
|
|
- 无代理:返回 (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._shared_proxy_manager:
|
|
proxy_url = self._shared_proxy_manager.get_proxy()
|
|
if proxy_url:
|
|
return {'http': proxy_url, 'https': proxy_url}, f'使用API代理: {proxy_url}'
|
|
return None, '代理不可用'
|
|
|
|
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,
|
|
proxy_manager=self._shared_proxy_manager,
|
|
)
|
|
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}")
|
|
|
|
try:
|
|
# 创建或复用任务记录(顺序执行,线程安全)
|
|
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 "login:view_all" not in self.creator_permissions:
|
|
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", "")
|
|
finally:
|
|
# 确保 DB Session 被关闭,避免连接泄漏
|
|
self.db.close()
|
|
|
|
|
|
# 在模块末尾导入 SessionLocal(避免循环导入)
|
|
from ..database import SessionLocal
|