Files
2026-08-31 10:55:44 +08:00

190 lines
8.3 KiB
Python

"""斗鱼活动任务批次执行器(入口聚合;功能域已拆分到 douyu_runner_*.py)。"""
from __future__ import annotations
from concurrent.futures import ThreadPoolExecutor, as_completed
from sqlalchemy.orm import joinedload
from core.douyu import DouyuActivityError
from ..database import SessionLocal
from ..models import DouyuTask
from .cookie_check_service import check_douyu_cookie
from .douyu_runner_bind import BindMixin
from .douyu_runner_core import ( # noqa: F401 (douyu_batch_registry 供 routers 重导出)
DouyuBatchRunnerCore,
douyu_batch_registry,
)
from .douyu_runner_donate import DonateMixin
from .douyu_runner_gold import GoldMixin
from .douyu_runner_goods import GoodsMixin
from .douyu_runner_manual import ManualMixin
from .douyu_runner_xpd import XpdMixin
from .douyu_service import latest_success_login_task, update_account_profile_from_cookie
class DouyuBatchRunner(
DouyuBatchRunnerCore,
BindMixin,
ManualMixin,
GoldMixin,
DonateMixin,
GoodsMixin,
XpdMixin,
):
"""批量执行斗鱼活动任务(功能域 Mixin 聚合 + 批次调度)。"""
def _execute_one(self, task_id: int, config: dict, total: int):
worker_db = SessionLocal()
try:
task = (
worker_db.query(DouyuTask)
.options(joinedload(DouyuTask.account))
.filter(DouyuTask.id == task_id)
.first()
)
if not task or self._stop.is_set():
return
account = task.account
self._update_task_progress(worker_db, task, "running", "执行中")
with self._counter_lock:
self._started += 1
current = self._started
self._push_log(
"info", f"[{current}/{total}] 开始: {self._account_name(account)}"
)
login_task = latest_success_login_task(worker_db, account.id)
cookie = login_task.cookie if login_task else ""
if not cookie:
self._mark_task(worker_db, task, "failed", "账号没有成功登录 Cookie")
self._push_log(
"warning", f"[{current}] {self._account_name(account)} 无 Cookie"
)
return
assert login_task is not None
cookie_check = check_douyu_cookie(cookie)
login_task.ck_check_status = "valid" if cookie_check["valid"] else "invalid"
login_task.ck_check_result = {
key: cookie_check.get(key)
for key in ("fish_ball", "nickname", "level", "message")
}
login_task.ck_checked_at = cookie_check["checked_at"]
worker_db.commit()
if not cookie_check["valid"]:
message = f"Cookie 已失效,请重新登录:{cookie_check['message']}"
self._mark_task(worker_db, task, "failed", message)
self._push_log(
"warning", f"[{current}] {self._account_name(account)} {message}"
)
return
update_account_profile_from_cookie(account, cookie)
handler = {
"refresh_goods": self._execute_refresh_goods,
"refresh_esports_goods": self._execute_refresh_esports_goods,
"get_bind_qr": self._execute_get_bind_qr,
"confirm_bind": self._execute_confirm_bind,
"create_elite_qr": self._execute_create_elite_qr,
"prepare_esports_bind": self._execute_prepare_esports_bind,
"get_esports_bind_qr": self._execute_get_esports_bind_qr,
"query_esports_game_name": self._execute_query_esports_game_name,
"confirm_esports_bind": self._execute_confirm_esports_bind,
"create_esports_qr": self._execute_create_esports_qr,
"query_esports_points": self._execute_query_esports_points,
"donate_esports_chicken_gift": self._execute_donate_esports_chicken_gift,
"donate_esports_firework_gift": self._execute_donate_esports_firework_gift,
"create_gold_qr": self._execute_create_gold_qr,
"donate_elite_gift": self._execute_donate_elite_gift,
"query_points": self._execute_query_points,
"lock_goods": self._execute_lock_goods,
"pay_locked_order": self._execute_pay_locked_order,
"exchange_goods": self._execute_exchange_goods,
"exchange_esports_goods": self._execute_exchange_esports_goods,
"query_game_name": self._execute_query_game_name,
"query_change_bind_time": self._execute_query_change_bind_time,
"query_limited_goods": self._execute_query_limited_goods,
"query_gold_balance": self._execute_query_gold_balance,
"query_exchange_records": self._execute_query_exchange_records,
"prefetch_csrf_token": self._execute_prefetch_csrf_token,
"get_xpd_bind_qr": self._execute_get_xpd_bind_qr,
"query_xpd_bind_info": self._execute_query_xpd_bind_info,
"confirm_xpd_bind": self._execute_confirm_xpd_bind,
"query_xpd_role": self._execute_query_xpd_role,
"refresh_xpd_goods": self._execute_refresh_xpd_goods,
"query_xpd_balance": self._execute_query_xpd_balance,
"query_xpd_fragments": self._execute_query_xpd_fragments,
"query_xpd_purchase_records": self._execute_query_xpd_purchase_records,
"exchange_xpd_goods": self._execute_exchange_xpd_goods,
}.get(task.task_type)
if handler is None:
self._mark_task(worker_db, task, "failed", "不支持的任务类型")
return
handler(worker_db, task, account, cookie, config)
self._push_log(
"success", f"[{current}] {self._account_name(account)} {task.message}"
)
except DouyuActivityError as exc:
if "task" in locals() and task:
self._mark_task(worker_db, task, "failed", str(exc))
self._push_log("warning", f"斗鱼任务失败: {exc}")
except Exception as exc: # noqa: BLE001 外部接口与任务边界需要保留宽泛异常兜底
if "task" in locals() and task:
self._mark_task(worker_db, task, "error", str(exc))
self._push_log("error", f"斗鱼任务异常: {exc}")
finally:
worker_db.close()
def run(self):
"""执行批次任务。"""
self._push_log("info", f"斗鱼任务批次 {self.batch_id} 开始")
try:
config = self._config_info(self.db)
tasks = (
self.db.query(DouyuTask)
.filter(
DouyuTask.batch_id == self.batch_id, DouyuTask.status == "planned"
)
.order_by(DouyuTask.id.asc())
.all()
)
if not tasks:
self._push_log("warning", "没有可执行的斗鱼任务")
self._push_log("result", "")
return
for task in tasks:
task.status = "pending"
task.message = "等待执行"
self.db.commit()
for task in tasks:
self._push_task_event(task)
total = len(tasks)
with ThreadPoolExecutor(max_workers=self.concurrency) as executor:
futures = []
for task in tasks:
if self._stop.is_set():
break
futures.append(
executor.submit(self._execute_one, task.id, config, total)
)
for future in as_completed(futures):
try:
future.result()
except Exception as exc: # noqa: BLE001 外部接口与任务边界需要保留宽泛异常兜底
self._push_log("error", f"Worker 异常: {exc}")
if self._stop.is_set():
self._push_log("warning", f"斗鱼任务批次 {self.batch_id} 已停止")
else:
self._push_log("info", f"斗鱼任务批次 {self.batch_id} 完成")
self._push_log("result", "")
finally:
self.db.close()