feat: 登录流程整体重试机制 + 代理池耗尽等待恢复
- DouyuLogin.login() 添加整体重试循环,失败时换新代理从头重跑 - 新增 max_login_retries(默认3次)和 max_total_time(默认300秒) - 重试前重置 session cookies + 换新代理,避免残留状态 - 成功后归还代理到池;失败时标记代理坏 - 代理池耗尽时不再立即放弃账号,等待最多60秒恢复后再试 - API代理模式改为 DouyuLogin 内部通过 proxy_manager 自治管理 登录过程中的代理切换(极验、整体重试)不再由服务层预取 - Schema/LoginBatchRequest 添加 max_login_retries 和 max_total_time - 前端 API 调用改为对象参数形式,透传新参数
This commit is contained in:
+45
-6
@@ -73,6 +73,8 @@ class DouyuLogin:
|
|||||||
timeout: tuple[float, float] = REQUEST_TIMEOUT,
|
timeout: tuple[float, float] = REQUEST_TIMEOUT,
|
||||||
max_geetest_retries: int = 5,
|
max_geetest_retries: int = 5,
|
||||||
max_proxy_retries: int = 10,
|
max_proxy_retries: int = 10,
|
||||||
|
max_login_retries: int = 3,
|
||||||
|
max_total_time: float = 300,
|
||||||
whitelist_uid: str = "",
|
whitelist_uid: str = "",
|
||||||
whitelist_ukey: str = "",
|
whitelist_ukey: str = "",
|
||||||
proxy_manager: Optional[ProxyManager] = None,
|
proxy_manager: Optional[ProxyManager] = None,
|
||||||
@@ -82,6 +84,8 @@ class DouyuLogin:
|
|||||||
self.timeout = timeout
|
self.timeout = timeout
|
||||||
self.max_geetest_retries = max_geetest_retries
|
self.max_geetest_retries = max_geetest_retries
|
||||||
self.max_proxy_retries = max_proxy_retries # 0=无限切换直到成功
|
self.max_proxy_retries = max_proxy_retries # 0=无限切换直到成功
|
||||||
|
self.max_login_retries = max_login_retries # 登录整体重试次数(换代理从头重跑)
|
||||||
|
self.max_total_time = max_total_time # 单账号登录总时长上限(秒),超时则放弃
|
||||||
self.session = requests.Session()
|
self.session = requests.Session()
|
||||||
|
|
||||||
# 初始化代理管理器(优先使用外部传入的共享实例,避免并发刷新冲突)
|
# 初始化代理管理器(优先使用外部传入的共享实例,避免并发刷新冲突)
|
||||||
@@ -314,12 +318,27 @@ class DouyuLogin:
|
|||||||
|
|
||||||
def login(self) -> LoginResult:
|
def login(self) -> LoginResult:
|
||||||
"""
|
"""
|
||||||
完整登录流程
|
完整登录流程(带整体重试)。
|
||||||
|
|
||||||
|
任何步骤失败时,换新代理从头重跑,最多重试 max_login_retries 次。
|
||||||
|
整体超时 max_total_time 秒后放弃。
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
LoginResult: 登录结果,包含cookie
|
LoginResult: 登录结果,包含cookie
|
||||||
"""
|
"""
|
||||||
logger.info(f"开始登录账号: {self.account.username}")
|
logger.info(f"开始登录账号: {self.account.username}")
|
||||||
|
start_time = time.monotonic()
|
||||||
|
|
||||||
|
for attempt in range(1, self.max_login_retries + 1):
|
||||||
|
elapsed = time.monotonic() - start_time
|
||||||
|
if elapsed > self.max_total_time:
|
||||||
|
logger.warning(f"登录总耗时 {elapsed:.0f}s 超过上限 {self.max_total_time}s,放弃")
|
||||||
|
return LoginResult(success=False, message=f"登录超时({elapsed:.0f}s > {self.max_total_time}s)")
|
||||||
|
|
||||||
|
if attempt > 1:
|
||||||
|
logger.info(f"登录整体重试 {attempt}/{self.max_login_retries},换代理重新开始")
|
||||||
|
# 重试前:换新代理 + 重置 session(清 cookies)
|
||||||
|
self._prepare_retry()
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# 1️⃣ 第一次登录(获取极验参数)
|
# 1️⃣ 第一次登录(获取极验参数)
|
||||||
@@ -352,15 +371,35 @@ class DouyuLogin:
|
|||||||
cookie = self._complete_login(login_url)
|
cookie = self._complete_login(login_url)
|
||||||
|
|
||||||
logger.success(f"登录成功! Cookie长度: {len(cookie)}")
|
logger.success(f"登录成功! Cookie长度: {len(cookie)}")
|
||||||
|
# 成功后归还代理到池,让其他账号复用
|
||||||
|
if self.proxy_manager and self._current_proxy_url:
|
||||||
|
self.proxy_manager.release_proxy(self._current_proxy_url)
|
||||||
return LoginResult(success=True, cookie=cookie, message="登录成功")
|
return LoginResult(success=True, cookie=cookie, message="登录成功")
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"登录失败: {e}")
|
elapsed = time.monotonic() - start_time
|
||||||
return LoginResult(success=False, message=str(e))
|
logger.error(f"登录失败(尝试 {attempt}/{self.max_login_retries},已耗时 {elapsed:.0f}s): {e}")
|
||||||
finally:
|
if attempt < self.max_login_retries and elapsed < self.max_total_time:
|
||||||
# 归还代理到池,让其他账号可以复用
|
# 还有重试机会且未超时:标记当前代理坏,下一轮自动换新代理
|
||||||
if self.proxy_manager and self._current_proxy_url:
|
if self.proxy_manager and self._current_proxy_url:
|
||||||
self.proxy_manager.release_proxy(self._current_proxy_url)
|
self.proxy_manager.mark_bad(self._current_proxy_url)
|
||||||
|
time.sleep(2)
|
||||||
|
continue
|
||||||
|
# 所有重试耗尽或超时
|
||||||
|
return LoginResult(success=False, message=str(e))
|
||||||
|
|
||||||
|
return LoginResult(success=False, message=f"登录失败,已重试 {self.max_login_retries} 次")
|
||||||
|
|
||||||
|
def _prepare_retry(self) -> None:
|
||||||
|
"""重试前准备:换新代理、重置 session cookies。"""
|
||||||
|
# 重置 session(清掉旧 cookies,避免残留状态干扰)
|
||||||
|
self.session.cookies.clear()
|
||||||
|
# 刷新代理
|
||||||
|
if self.proxy_manager:
|
||||||
|
new_proxy = self.proxy_manager.get_proxy()
|
||||||
|
if new_proxy:
|
||||||
|
self._apply_proxy(new_proxy)
|
||||||
|
# 如果获取新代理失败,保留旧代理继续尝试
|
||||||
|
|
||||||
def _first_login(self) -> Tuple[str, str, str, dict]:
|
def _first_login(self) -> Tuple[str, str, str, dict]:
|
||||||
"""
|
"""
|
||||||
|
|||||||
@@ -56,6 +56,8 @@ async def create_batch(
|
|||||||
creator_permissions=get_user_permissions(current),
|
creator_permissions=get_user_permissions(current),
|
||||||
max_geetest_retries=req.max_geetest_retries,
|
max_geetest_retries=req.max_geetest_retries,
|
||||||
max_proxy_retries=req.max_proxy_retries,
|
max_proxy_retries=req.max_proxy_retries,
|
||||||
|
max_login_retries=req.max_login_retries,
|
||||||
|
max_total_time=req.max_total_time,
|
||||||
proxy_config=proxy,
|
proxy_config=proxy,
|
||||||
log_queue=log_queue,
|
log_queue=log_queue,
|
||||||
loop=loop,
|
loop=loop,
|
||||||
|
|||||||
@@ -126,6 +126,8 @@ class LoginBatchRequest(BaseModel):
|
|||||||
account_ids: list[int]
|
account_ids: list[int]
|
||||||
max_geetest_retries: int = 5
|
max_geetest_retries: int = 5
|
||||||
max_proxy_retries: int = 10 # 代理切换次数,0=无限切换直到成功
|
max_proxy_retries: int = 10 # 代理切换次数,0=无限切换直到成功
|
||||||
|
max_login_retries: int = 3 # 登录整体重试次数(换代理从头重跑)
|
||||||
|
max_total_time: float = 300 # 单账号登录总时长上限(秒),超时则放弃
|
||||||
concurrency: int = 3 # 并发数,1-10
|
concurrency: int = 3 # 并发数,1-10
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -2,6 +2,7 @@
|
|||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
import threading
|
import threading
|
||||||
|
import time
|
||||||
import uuid
|
import uuid
|
||||||
from types import SimpleNamespace
|
from types import SimpleNamespace
|
||||||
from concurrent.futures import ThreadPoolExecutor, as_completed
|
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||||
@@ -18,6 +19,9 @@ from ..models import Account as AccountModel, LoginTask, ProxyConfig as ProxyCon
|
|||||||
class LoginBatchRunner:
|
class LoginBatchRunner:
|
||||||
"""批量登录执行器,在线程中运行,通过 ThreadPoolExecutor 并发登录多个账号。"""
|
"""批量登录执行器,在线程中运行,通过 ThreadPoolExecutor 并发登录多个账号。"""
|
||||||
|
|
||||||
|
# 代理池耗尽时等待恢复的最大秒数
|
||||||
|
_PROXY_WAIT_MAX = 60
|
||||||
|
|
||||||
def __init__(
|
def __init__(
|
||||||
self,
|
self,
|
||||||
db: Session,
|
db: Session,
|
||||||
@@ -26,6 +30,8 @@ class LoginBatchRunner:
|
|||||||
creator_permissions: list[str],
|
creator_permissions: list[str],
|
||||||
max_geetest_retries: int = 5,
|
max_geetest_retries: int = 5,
|
||||||
max_proxy_retries: int = 10,
|
max_proxy_retries: int = 10,
|
||||||
|
max_login_retries: int = 3,
|
||||||
|
max_total_time: float = 300,
|
||||||
proxy_config: Optional[ProxyConfigModel] = None,
|
proxy_config: Optional[ProxyConfigModel] = None,
|
||||||
log_queue: Optional[asyncio.Queue] = None,
|
log_queue: Optional[asyncio.Queue] = None,
|
||||||
loop: Optional[asyncio.AbstractEventLoop] = None,
|
loop: Optional[asyncio.AbstractEventLoop] = None,
|
||||||
@@ -37,6 +43,8 @@ class LoginBatchRunner:
|
|||||||
self.creator_permissions = creator_permissions
|
self.creator_permissions = creator_permissions
|
||||||
self.max_geetest_retries = max_geetest_retries
|
self.max_geetest_retries = max_geetest_retries
|
||||||
self.max_proxy_retries = max_proxy_retries
|
self.max_proxy_retries = max_proxy_retries
|
||||||
|
self.max_login_retries = max_login_retries
|
||||||
|
self.max_total_time = max_total_time
|
||||||
self.proxy_config = proxy_config
|
self.proxy_config = proxy_config
|
||||||
self.log_queue = log_queue
|
self.log_queue = log_queue
|
||||||
self.loop = loop
|
self.loop = loop
|
||||||
@@ -67,13 +75,13 @@ class LoginBatchRunner:
|
|||||||
self.loop,
|
self.loop,
|
||||||
)
|
)
|
||||||
|
|
||||||
def _resolve_proxy(self) -> tuple[Optional[dict], str]:
|
def _resolve_static_proxy(self) -> tuple[Optional[dict], str]:
|
||||||
"""
|
"""
|
||||||
解析代理配置,返回 (proxy_dict, message)。
|
解析静态代理配置(仅处理无代理和静态代理场景)。
|
||||||
|
|
||||||
- 静态代理:直接返回 dict
|
API代理由 DouyuLogin 通过 proxy_manager 内部管理,
|
||||||
- API代理:通过共享 ProxyManager(带已验证代理池缓存)获取,自动同步白名单
|
不在此处预先获取——登录过程中的代理切换(极验失败、整体重试)
|
||||||
- 无代理:返回 (None, '')
|
都在 DouyuLogin 内部自治完成。
|
||||||
"""
|
"""
|
||||||
if not self.proxy_config or not self.proxy_config.enabled:
|
if not self.proxy_config or not self.proxy_config.enabled:
|
||||||
return None, ''
|
return None, ''
|
||||||
@@ -83,13 +91,7 @@ class LoginBatchRunner:
|
|||||||
proxy_url = 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}'
|
return {'http': proxy_url, 'https': proxy_url}, f'使用静态代理: {proxy_url}'
|
||||||
|
|
||||||
# API代理:通过共享代理管理器获取(优先从已验证代理池复用)
|
# API代理:不在此处获取,由 DouyuLogin 通过 proxy_manager 内部管理
|
||||||
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, ''
|
return None, ''
|
||||||
|
|
||||||
def _execute_one(self, task_id: int, acc_info: dict, total: int):
|
def _execute_one(self, task_id: int, acc_info: dict, total: int):
|
||||||
@@ -113,18 +115,61 @@ class LoginBatchRunner:
|
|||||||
|
|
||||||
self._push_log("info", f"[{current}/{total}] 开始登录: {acc_info['username']}")
|
self._push_log("info", f"[{current}/{total}] 开始登录: {acc_info['username']}")
|
||||||
|
|
||||||
# 每个账号独立获取代理
|
# 解析代理配置
|
||||||
proxy_dict, proxy_msg = self._resolve_proxy()
|
proxy_dict, proxy_msg = self._resolve_static_proxy()
|
||||||
if proxy_msg:
|
if proxy_msg:
|
||||||
self._push_log("info", f"[{current}] {proxy_msg}")
|
self._push_log("info", f"[{current}] {proxy_msg}")
|
||||||
|
|
||||||
# 代理预检失败时,该账号标记失败但不终止整个批次
|
# API代理模式下,验证代理池是否有可用代理;池空时等待恢复(最多 _PROXY_WAIT_MAX 秒)
|
||||||
if self.proxy_config and self.proxy_config.enabled and not proxy_dict:
|
# 注意:不预取代理传给 DouyuLogin,让 DouyuLogin 通过 proxy_manager 内部自治管理
|
||||||
|
is_api_proxy = (
|
||||||
|
self.proxy_config
|
||||||
|
and self.proxy_config.enabled
|
||||||
|
and not (self.proxy_config.http or self.proxy_config.https)
|
||||||
|
and self._shared_proxy_manager
|
||||||
|
)
|
||||||
|
if is_api_proxy:
|
||||||
|
# 先验证代理池是否有可用代理(获取后立即归还,不占用)
|
||||||
|
test_proxy = self._shared_proxy_manager.get_proxy()
|
||||||
|
if test_proxy:
|
||||||
|
self._shared_proxy_manager.release_proxy(test_proxy)
|
||||||
|
self._push_log("info", f"[{current}] 代理池可用,由 DouyuLogin 内部管理代理获取与切换")
|
||||||
|
else:
|
||||||
|
# 代理池暂时耗尽,等待冷却代理恢复或新代理入池
|
||||||
|
self._push_log("warning", f"[{current}] 代理池暂时耗尽,等待恢复...")
|
||||||
|
pool_available = False
|
||||||
|
for wait_sec in range(0, self._PROXY_WAIT_MAX, 10):
|
||||||
|
if self._stop.is_set():
|
||||||
task.status = "error"
|
task.status = "error"
|
||||||
task.message = f"代理不可用: {proxy_msg}"
|
task.message = "任务已停止"
|
||||||
task.finished_at = datetime.now(timezone.utc)
|
task.finished_at = datetime.now(timezone.utc)
|
||||||
worker_db.commit()
|
worker_db.commit()
|
||||||
self._push_log("error", f"[{current}] {acc_info['username']} 代理不可用: {proxy_msg}")
|
return
|
||||||
|
time.sleep(10)
|
||||||
|
test_proxy = self._shared_proxy_manager.get_proxy()
|
||||||
|
if test_proxy:
|
||||||
|
self._shared_proxy_manager.release_proxy(test_proxy)
|
||||||
|
pool_available = True
|
||||||
|
self._push_log("info", f"[{current}] 代理池恢复,由 DouyuLogin 内部管理代理")
|
||||||
|
break
|
||||||
|
remaining = self._PROXY_WAIT_MAX - wait_sec - 10
|
||||||
|
self._push_log("info", f"[{current}] 代理池仍为空,继续等待... (剩余 {remaining}s)")
|
||||||
|
|
||||||
|
if not pool_available:
|
||||||
|
task.status = "error"
|
||||||
|
task.message = f"代理池耗尽,等待 {self._PROXY_WAIT_MAX}s 后仍无可用代理"
|
||||||
|
task.finished_at = datetime.now(timezone.utc)
|
||||||
|
worker_db.commit()
|
||||||
|
self._push_log("error", f"[{current}] {acc_info['username']} 代理池耗尽,放弃")
|
||||||
|
return
|
||||||
|
|
||||||
|
# 静态代理启用但配置为空 → 不可用
|
||||||
|
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_manager and not proxy_dict:
|
||||||
|
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
|
return
|
||||||
|
|
||||||
try:
|
try:
|
||||||
@@ -143,6 +188,8 @@ class LoginBatchRunner:
|
|||||||
proxy=proxy_dict,
|
proxy=proxy_dict,
|
||||||
max_geetest_retries=self.max_geetest_retries,
|
max_geetest_retries=self.max_geetest_retries,
|
||||||
max_proxy_retries=self.max_proxy_retries,
|
max_proxy_retries=self.max_proxy_retries,
|
||||||
|
max_login_retries=self.max_login_retries,
|
||||||
|
max_total_time=self.max_total_time,
|
||||||
proxy_manager=self._shared_proxy_manager,
|
proxy_manager=self._shared_proxy_manager,
|
||||||
)
|
)
|
||||||
result = loginer.login()
|
result = loginer.login()
|
||||||
|
|||||||
@@ -1,9 +1,18 @@
|
|||||||
import api from './client';
|
import api from './client';
|
||||||
import type { BatchLoginResult, LoginTaskItem, MessageDeletedResponse, MessageResponse } from './types';
|
import type { BatchLoginResult, LoginTaskItem, MessageDeletedResponse, MessageResponse } from './types';
|
||||||
|
|
||||||
|
interface CreateBatchParams {
|
||||||
|
account_ids: number[];
|
||||||
|
max_geetest_retries?: number;
|
||||||
|
concurrency?: number;
|
||||||
|
max_proxy_retries?: number;
|
||||||
|
max_login_retries?: number;
|
||||||
|
max_total_time?: number;
|
||||||
|
}
|
||||||
|
|
||||||
export const loginApi = {
|
export const loginApi = {
|
||||||
createBatch: (account_ids: number[], max_geetest_retries?: number, concurrency?: number, max_proxy_retries?: number) =>
|
createBatch: (params: CreateBatchParams) =>
|
||||||
api.post<BatchLoginResult, BatchLoginResult>('/login/batch', { account_ids, max_geetest_retries, concurrency, max_proxy_retries }),
|
api.post<BatchLoginResult, BatchLoginResult>('/login/batch', params),
|
||||||
listTasks: (batch_id?: string) =>
|
listTasks: (batch_id?: string) =>
|
||||||
api.get<LoginTaskItem[], LoginTaskItem[]>('/login/tasks', { params: batch_id ? { batch_id } : {} }),
|
api.get<LoginTaskItem[], LoginTaskItem[]>('/login/tasks', { params: batch_id ? { batch_id } : {} }),
|
||||||
stop: (batch_id: string) => api.post<MessageResponse, MessageResponse>(`/login/stop/${batch_id}`),
|
stop: (batch_id: string) => api.post<MessageResponse, MessageResponse>(`/login/stop/${batch_id}`),
|
||||||
|
|||||||
@@ -143,7 +143,14 @@ export default function LoginTasksPage() {
|
|||||||
}
|
}
|
||||||
setLoading(true);
|
setLoading(true);
|
||||||
try {
|
try {
|
||||||
const result = await loginApi.createBatch(accountIds, 5, concurrency, maxProxyRetries);
|
const result = await loginApi.createBatch({
|
||||||
|
account_ids: accountIds,
|
||||||
|
max_geetest_retries: 5,
|
||||||
|
concurrency,
|
||||||
|
max_proxy_retries: maxProxyRetries,
|
||||||
|
max_login_retries: 3,
|
||||||
|
max_total_time: 300,
|
||||||
|
});
|
||||||
setBatchId(result.batch_id);
|
setBatchId(result.batch_id);
|
||||||
message.success(`已创建登录任务,共 ${result.count} 个账号`);
|
message.success(`已创建登录任务,共 ${result.count} 个账号`);
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user