Files
live-hub-py/web/backend/services/douyu_runner.py
T
yml2213 592eef3582 feat(douyu): 兑换成功后生成结果图片并自动复制到剪贴板
- 账号表格新增兑换图片列,展示每账号最近兑换成功的结果卡片
- 新完成兑换任务自动生成图片并复制到剪贴板(批量时最后完成者覆盖)
- canvas 绘制兑换卡片:商品图/商品名/游戏角色/区服/兑换账号/时间
- 后端兑换成功结果补充 game_name/game_channel/account_name 字段
- 不支持剪贴板的环境降级为点击查看大图
2026-08-01 13:23:22 +08:00

2255 lines
93 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""斗鱼活动任务批次执行器。"""
from __future__ import annotations
import asyncio
import threading
import time
from concurrent.futures import ThreadPoolExecutor, as_completed
from datetime import datetime, timezone
from typing import Optional
from loguru import logger
from sqlalchemy.orm import Session, joinedload
from core.douyu import DouyuActivityClient, DouyuActivityError
from ..database import SessionLocal
from ..models import Account, DouyuEsportsGoodsSnapshot, DouyuGoodsSnapshot, DouyuTask
from .douyu_service import (
DOUYU_CONFIG_FIELDS,
account_uid,
douyu_config_value,
ensure_douyu_config,
latest_success_cookie,
update_account_profile_from_cookie,
)
DOUYU_LEGACY_BIND_ACT_ALIAS = "20250213NQCYX"
DOUYU_BIND_ROLE_POLL_SECONDS = 65
DOUYU_BIND_ROLE_POLL_INTERVAL = 5
DOUYU_PAYMENT_POLL_SECONDS = 600
DOUYU_PAYMENT_POLL_INTERVAL = 5
DOUYU_GIFT_POINTS_REFRESH_TIMES = 3
DOUYU_GIFT_POINTS_REFRESH_INTERVAL = 2
class DouyuBatchRunner:
"""批量执行斗鱼活动任务,通过队列推送实时日志。"""
def __init__(
self,
db: Session,
batch_id: str,
task_type: str,
payload: Optional[dict] = None,
log_queue: Optional[asyncio.Queue] = None,
loop: Optional[asyncio.AbstractEventLoop] = None,
concurrency: int = 3,
):
self.db = db
self.batch_id = batch_id
self.task_type = task_type
self.payload = payload or {}
self.log_queue = log_queue
self.loop = loop
self.concurrency = max(1, min(concurrency, 10))
self._stop = threading.Event()
self._counter_lock = threading.Lock()
self._started = 0
def stop(self):
self._stop.set()
def _push_log(self, level: str, message: str):
if level == "result":
try:
douyu_batch_registry.mark_finished(self.batch_id)
except NameError:
pass
if level != "result" and message:
log_func = getattr(logger, level, logger.info)
log_func(f"[douyu] {message}")
if self.log_queue and self.loop:
asyncio.run_coroutine_threadsafe(
self.log_queue.put({"level": level, "message": message}),
self.loop,
)
@staticmethod
def _account_name(account: Account) -> str:
return account.nickname or account.username or account.uid or f"#{account.id}"
@staticmethod
def _to_int(value) -> int | None:
if value is None:
return None
try:
return int(value)
except (TypeError, ValueError):
return None
@staticmethod
def _format_wait_time(seconds: int | None) -> str:
if seconds is None:
return ""
seconds = max(0, int(seconds))
days, rem = divmod(seconds, 86400)
hours, rem = divmod(rem, 3600)
minutes, sec = divmod(rem, 60)
if days:
return f"{days}{hours}小时{minutes}分"
if hours:
return f"{hours}小时{minutes}{sec}秒"
return f"{minutes}{sec}秒"
@staticmethod
def _action_act_alias(config: dict, key: str) -> str:
"""动作类接口用的活动 alias,排除只用于查询最新角色的 legacy alias。"""
alias = str(config.get(key) or "").strip()
query_only_alias = str(config.get("legacy_act_alias") or "").strip()
if not alias:
return ""
if alias in {query_only_alias, DOUYU_LEGACY_BIND_ACT_ALIAS}:
return ""
return alias
@classmethod
def _bind_qr_act_alias(cls, config: dict) -> str:
"""生成绑定二维码用的活动 alias。"""
return cls._action_act_alias(config, "bind_act_alias") or cls._action_act_alias(config, "confirm_act_alias")
@staticmethod
def _query_bind_act_aliases(config: dict) -> list[str]:
"""查询/轮询角色用的 alias 列表。
现网最新绑定信息在 legacy(cjm);活动 alias 可能仍用于扫码/确认,
所以按优先级去重返回多个,轮询时取“更像新扫码结果”的那个。
"""
ordered = [
str(config.get("legacy_act_alias") or "").strip(),
str(config.get("confirm_act_alias") or "").strip(),
str(config.get("bind_act_alias") or "").strip(),
]
aliases: list[str] = []
for alias in ordered:
if alias and alias not in aliases:
aliases.append(alias)
return aliases
@classmethod
def _confirm_act_alias(cls, config: dict) -> str:
"""确认绑定接口用的活动 alias。"""
return cls._action_act_alias(config, "confirm_act_alias") or cls._action_act_alias(config, "bind_act_alias")
# 兼容旧调用名
@classmethod
def _current_bind_act_alias(cls, config: dict) -> str:
return cls._confirm_act_alias(config) or cls._bind_qr_act_alias(config)
@staticmethod
def _role_channel(bind_info: dict) -> str:
return " / ".join(
part for part in [bind_info.get("area_name"), bind_info.get("plat_name")] if part
)
@staticmethod
def _is_truthy_flag(value) -> bool:
if value is True:
return True
if value is False or value is None:
return False
text = str(value).strip().lower()
return text in {"1", "true", "yes", "y"}
@classmethod
def _is_bound_act(cls, bind_info: dict | None) -> bool:
if not bind_info:
return False
return cls._is_truthy_flag(bind_info.get("is_bound_act"))
@classmethod
def _can_change_role(cls, bind_info: dict | None) -> bool:
"""综合 can_change_role 与换绑倒计时判断是否允许换绑。"""
if not bind_info:
return True
if str(bind_info.get("api_version") or "") == "esports":
can_change_time = cls._to_int(bind_info.get("can_change_time"))
if can_change_time is not None:
return can_change_time <= int(datetime.now(timezone.utc).timestamp())
wait_time = cls._to_int(bind_info.get("change_role_wait_time"))
if wait_time is not None and wait_time > 0:
return False
if bind_info.get("can_change_role") is not None:
return cls._is_truthy_flag(bind_info.get("can_change_role"))
# 已绑定但接口没给倒计时/开关时,默认允许(避免误杀首次绑定)
return True
@classmethod
def _is_change_cooling(cls, bind_info: dict | None) -> bool:
"""是否处于换绑冷却:已有绑定角色且当前不可换绑。"""
if not bind_info:
return False
role_name = str(bind_info.get("role_name") or "").strip()
is_bound = cls._is_bound_act(bind_info) or bool(role_name)
if not is_bound:
return False
return not cls._can_change_role(bind_info)
@classmethod
def _is_pending_role(
cls,
bind_info: dict | None,
*,
baseline_role_name: str = "",
baseline_is_bound_act: bool = False,
) -> bool:
"""判断是否出现了可确认的扫码角色。
规则:
1. 必须有角色名
2. 已绑定同一角色(或无 baseline 的已绑定)不算 pending
3. 角色名相对 baseline 变化,或从已绑定变成待确认,算 pending
4. need_bind_act / need_bind_role 且角色相对 baseline 有变化,算 pending
5. 无 baseline 且未绑定但有角色,视为待确认残留
"""
if not bind_info:
return False
role_name = str(bind_info.get("role_name") or "").strip()
if not role_name:
return False
is_bound_act = cls._is_bound_act(bind_info)
need_bind_act = cls._is_truthy_flag(bind_info.get("need_bind_act"))
need_bind_role = cls._is_truthy_flag(bind_info.get("need_bind_role"))
role_changed = bool(baseline_role_name) and role_name != baseline_role_name
if is_bound_act:
# 已绑定:仅当相对 baseline 角色发生变化时才视为新扫码结果
return role_changed
if need_bind_act or need_bind_role:
if not baseline_role_name:
return True
return role_changed or baseline_is_bound_act
if not baseline_role_name:
return True
if role_changed:
return True
# 同一角色从已绑定变为未绑定
return baseline_is_bound_act
def _bind_snapshot(self, bind_info: dict | None) -> dict:
info = bind_info or {}
role_name = str(info.get("role_name") or "")
return {
"role_name": role_name,
"area_name": info.get("area_name") or "",
"plat_name": info.get("plat_name") or "",
"nickname": info.get("nickname") or "",
"is_bound_act": self._is_bound_act(info),
"is_bound_role": self._is_truthy_flag(info.get("is_bound_role")),
"is_bound_account": self._is_truthy_flag(info.get("is_bound_account")),
"need_bind_act": self._is_truthy_flag(info.get("need_bind_act")),
"need_bind_role": self._is_truthy_flag(info.get("need_bind_role")),
"change_role_wait_time": self._to_int(info.get("change_role_wait_time")),
"can_change_time": self._to_int(info.get("can_change_time")),
"can_change_role": self._can_change_role(info),
"bind_info": info,
}
def _format_bind_summary(self, bind_info: dict | None, *, pending: bool | None = None) -> str:
snap = self._bind_snapshot(bind_info)
role = snap["role_name"] or "-"
area = snap["area_name"] or "-"
plat = snap["plat_name"] or "-"
pending_text = ""
if pending is not None:
pending_text = f" pending={1 if pending else 0}"
return (
f"role={role} area={area} plat={plat}"
f" bound_act={1 if snap['is_bound_act'] else 0}"
f" bound_role={1 if snap['is_bound_role'] else 0}"
f" need_act={1 if snap['need_bind_act'] else 0}"
f" need_role={1 if snap['need_bind_role'] else 0}"
f" wait={snap['change_role_wait_time'] if snap['change_role_wait_time'] is not None else '-'}"
f" can_change_time={snap['can_change_time'] if snap['can_change_time'] is not None else '-'}"
f"{pending_text}"
)
def _apply_bind_info_to_account(self, account: Account, bind_info: dict, status: str) -> None:
role_name = str(bind_info.get("role_name") or "")
account.game_name = role_name or account.game_name
account.game_channel = self._role_channel(bind_info) or account.game_channel
account.change_role_wait_time = self._to_int(bind_info.get("change_role_wait_time"))
account.bind_status = status
account.updated_at = datetime.now(timezone.utc)
def _apply_esports_bind_info_to_account(self, account: Account, bind_info: dict, status: str) -> None:
"""将电竞手册角色状态写入专属字段,避免覆盖精英宝典数据。"""
role_name = str(bind_info.get("role_name") or "")
account.esports_game_name = role_name or account.esports_game_name
account.esports_game_channel = self._role_channel(bind_info) or account.esports_game_channel
account.esports_change_role_wait_time = self._to_int(bind_info.get("change_role_wait_time"))
account.esports_can_change_time = self._to_int(bind_info.get("can_change_time"))
account.esports_bind_status = status
account.updated_at = datetime.now(timezone.utc)
def _esports_bind_state(
self,
client: DouyuActivityClient,
act_alias: str,
) -> dict:
"""查询电竞手册活动状态,接口同时返回当前角色和换绑冷却信息。"""
activity_info = client.esports_bind_info(act_alias)
activity_snapshot = self._bind_snapshot(activity_info)
esports_bound = activity_snapshot["is_bound_act"]
has_selected_role = bool(activity_snapshot["role_name"])
tx_act = activity_info.get("tx_act") or {}
return {
**activity_snapshot,
"is_bound_role": has_selected_role,
"act_alias": act_alias,
"game_id": tx_act.get("gameId") or "",
"esports_bound": esports_bound,
"has_selected_role": has_selected_role,
"bind_ready_for_confirm": has_selected_role and not esports_bound,
"bind_confirmed": esports_bound,
"bind_phase": "confirmed" if esports_bound else ("role_ready" if has_selected_role else "waiting_role"),
"activity_bind_info": activity_info,
"activity_bind_snapshot": activity_snapshot,
"role_source": "activity",
}
@staticmethod
def _esports_role_text(state: dict) -> str:
role_name = str(state.get("role_name") or "")
channel = " / ".join(
str(part)
for part in [state.get("plat_name"), state.get("area_name")]
if part
)
if not role_name:
return ""
return f"{role_name}{channel}" if channel else role_name
def _mark_task(
self,
db: Session,
task: DouyuTask,
status: str,
message: str,
result: dict | None = None,
) -> None:
task.status = status
task.message = message[:512]
if result is not None:
task.result = result
task.finished_at = datetime.now(timezone.utc)
db.commit()
def _update_task_progress(
self,
db: Session,
task: DouyuTask,
status: str,
message: str,
result: dict | None = None,
) -> None:
task.status = status
task.message = message[:512]
if result is not None:
task.result = result
db.commit()
def _upsert_goods(self, db: Session, goods: list[dict]) -> None:
now = datetime.now(timezone.utc)
for raw in goods:
commodity_id = str(raw.get("commodityId") or raw.get("commodity_id") or "")
if not commodity_id:
continue
row = (
db.query(DouyuGoodsSnapshot)
.filter(DouyuGoodsSnapshot.commodity_id == commodity_id)
.first()
)
score = self._to_int(raw.get("score"))
if row is None:
row = DouyuGoodsSnapshot(commodity_id=commodity_id)
db.add(row)
row.name = str(raw.get("commodityName") or raw.get("name") or "")
row.score = score
row.status = str(raw.get("status") or "")
row.raw = raw
row.updated_at = now
db.commit()
def _upsert_esports_goods(self, db: Session, goods: list[dict]) -> None:
now = datetime.now(timezone.utc)
for raw in goods:
commodity_id = str(raw.get("commodityId") or raw.get("commodity_id") or "")
if not commodity_id:
continue
row = (
db.query(DouyuEsportsGoodsSnapshot)
.filter(DouyuEsportsGoodsSnapshot.commodity_id == commodity_id)
.first()
)
if row is None:
row = DouyuEsportsGoodsSnapshot(commodity_id=commodity_id)
db.add(row)
row.name = str(raw.get("commodityName") or raw.get("name") or "")
row.score = self._to_int(raw.get("score"))
row.status = str(raw.get("status") or "")
row.raw = raw
row.updated_at = now
db.commit()
def _config_info(self, db: Session) -> dict:
config = ensure_douyu_config(db)
return {field: douyu_config_value(field, getattr(config, field, None)) for field in DOUYU_CONFIG_FIELDS}
def _refresh_account_points(
self,
client: DouyuActivityClient,
account: Account,
cookie: str,
*,
ctn: str | None = None,
) -> dict:
"""刷新账号积分并写回账号表。"""
uid = account_uid(account, cookie)
if not uid:
raise DouyuActivityError("Cookie 中没有 acf_uid,无法查询积分")
ctn_value = ctn or client.acf_ccn(refresh_subscribe=False)
result = client.query_points(uid=uid, ctn=ctn_value)
points = self._to_int(result.get("points"))
account.uid = uid
account.points = points
update_account_profile_from_cookie(account, cookie)
account.updated_at = datetime.now(timezone.utc)
return {"points": points, "points_query": result}
def _refresh_account_gold_balance(self, client: DouyuActivityClient, account: Account) -> dict:
"""刷新鱼翅和钱包兑换余额并写回账号表。"""
gold = client.gold_account()
exchange = client.exchange_balance()
account.gold_balance = self._to_int(gold.get("gold"))
account.exchange_balance = self._to_int(exchange.get("count"))
account.updated_at = datetime.now(timezone.utc)
return {
"gold_balance": account.gold_balance,
"exchange_balance": account.exchange_balance,
"gold": gold,
"exchange_balance_query": exchange,
}
def _wait_points_after_payment(
self,
db: Session,
task: DouyuTask,
account: Account,
client: DouyuActivityClient,
cookie: str,
ctn: str,
result: dict,
) -> bool:
"""等待宝典支付到账;积分达到 300 视为开通成功。"""
deadline = time.monotonic() + DOUYU_PAYMENT_POLL_SECONDS
result["payment_polling"] = True
result["payment_target_points"] = 300
poll_count = 0
last_points = None
while not self._stop.is_set() and time.monotonic() <= deadline:
try:
points_result = self._refresh_account_points(client, account, cookie, ctn=ctn)
db.commit()
poll_count += 1
last_points = points_result["points"]
result.update(points_result)
result["payment_poll_count"] = poll_count
result["payment_polling"] = True
if last_points is not None and last_points >= 300:
result["payment_polling"] = False
result["elite_opened"] = True
return True
self._update_task_progress(
db,
task,
"running",
f"精英宝典支付码已生成,等待开通到账(当前积分 {last_points if last_points is not None else '-'}",
result,
)
except Exception as exc:
poll_count += 1
result["payment_poll_count"] = poll_count
result["payment_poll_error"] = str(exc)
self._update_task_progress(db, task, "running", f"等待开通到账: {exc}", result)
if self._stop.wait(DOUYU_PAYMENT_POLL_INTERVAL):
break
result["payment_polling"] = False
result["elite_opened"] = False
result["points"] = last_points
return False
def _refresh_esports_handbook(
self,
client: DouyuActivityClient,
account: Account,
*,
manual_id: str,
) -> dict:
"""刷新电竞手册开通状态并将积分写回账号。"""
result = client.esports_user_info(manual_id=manual_id)
manual_type = self._to_int(result.get("manual_type"))
manual_score = self._to_int(result.get("manual_score"))
account.esports_points = manual_score
account.updated_at = datetime.now(timezone.utc)
return {
"esports_manual_type": manual_type,
"esports_manual_score": manual_score,
"esports_expire_time": result.get("expire_time"),
"esports_user_info": result,
"esports_points": manual_score,
"points": manual_score,
}
def _wait_esports_open_after_payment(
self,
db: Session,
task: DouyuTask,
account: Account,
client: DouyuActivityClient,
*,
manual_id: str,
result: dict,
baseline_manual_type: int | None,
baseline_manual_score: int | None,
) -> bool:
"""等待电竞手册支付到账,以 manualType=1 或积分变化作为成功条件。"""
deadline = time.monotonic() + DOUYU_PAYMENT_POLL_SECONDS
result["payment_polling"] = True
result["esports_manual_type_baseline"] = baseline_manual_type
result["esports_manual_score_baseline"] = baseline_manual_score
poll_count = 0
last_manual_type = baseline_manual_type
last_manual_score = baseline_manual_score
while not self._stop.is_set() and time.monotonic() <= deadline:
try:
handbook_result = self._refresh_esports_handbook(
client,
account,
manual_id=manual_id,
)
db.commit()
poll_count += 1
last_manual_type = handbook_result["esports_manual_type"]
last_manual_score = handbook_result["esports_manual_score"]
result.update(handbook_result)
result["payment_poll_count"] = poll_count
result["payment_polling"] = True
opened = (
last_manual_type is not None
and last_manual_type >= 1
) or (
baseline_manual_score is not None
and last_manual_score is not None
and last_manual_score > baseline_manual_score
)
if opened:
result["payment_polling"] = False
result["esports_opened"] = True
return True
self._update_task_progress(
db,
task,
"running",
"电竞手册支付码已生成,等待开通到账"
f"(类型 {last_manual_type if last_manual_type is not None else '-'}"
f"积分 {last_manual_score if last_manual_score is not None else '-'}",
result,
)
except Exception as exc:
poll_count += 1
result["payment_poll_count"] = poll_count
result["payment_poll_error"] = str(exc)
self._update_task_progress(db, task, "running", f"等待电竞手册到账: {exc}", result)
if self._stop.wait(DOUYU_PAYMENT_POLL_INTERVAL):
break
result["payment_polling"] = False
result["esports_opened"] = False
result["esports_manual_type"] = last_manual_type
result["esports_manual_score"] = last_manual_score
result["esports_points"] = last_manual_score
result["points"] = last_manual_score
return False
def _wait_gold_balance_after_payment(
self,
db: Session,
task: DouyuTask,
account: Account,
client: DouyuActivityClient,
result: dict,
baseline_gold: int | None,
) -> bool:
"""等待鱼翅充值到账;余额变化后写回账号表。"""
deadline = time.monotonic() + DOUYU_PAYMENT_POLL_SECONDS
result["payment_polling"] = True
result["baseline_gold_balance"] = baseline_gold
poll_count = 0
last_gold = baseline_gold
baseline_ready = baseline_gold is not None
while not self._stop.is_set() and time.monotonic() <= deadline:
try:
balance_result = self._refresh_account_gold_balance(client, account)
db.commit()
poll_count += 1
last_gold = balance_result["gold_balance"]
result.update(balance_result)
result["payment_poll_count"] = poll_count
result["payment_polling"] = True
if not baseline_ready and last_gold is not None:
baseline_gold = last_gold
result["baseline_gold_balance"] = baseline_gold
baseline_ready = True
self._update_task_progress(
db,
task,
"running",
f"鱼翅支付码已生成,已记录当前余额 {last_gold},等待到账",
result,
)
if self._stop.wait(DOUYU_PAYMENT_POLL_INTERVAL):
break
continue
changed = last_gold is not None and (baseline_gold is None or last_gold != baseline_gold)
if changed:
result["payment_polling"] = False
result["gold_recharged"] = True
return True
self._update_task_progress(
db,
task,
"running",
f"鱼翅支付码已生成,等待到账(当前鱼翅 {last_gold if last_gold is not None else '-'}",
result,
)
except Exception as exc:
poll_count += 1
result["payment_poll_count"] = poll_count
result["payment_poll_error"] = str(exc)
self._update_task_progress(db, task, "running", f"等待鱼翅到账: {exc}", result)
if self._stop.wait(DOUYU_PAYMENT_POLL_INTERVAL):
break
result["payment_polling"] = False
result["gold_recharged"] = False
result["gold_balance"] = last_gold
return False
def _refresh_points_after_elite_gift(
self,
db: Session,
task: DouyuTask,
account: Account,
client: DouyuActivityClient,
cookie: str,
ctn: str | None,
result: dict,
baseline_points: int | None,
gift_count: int,
) -> dict:
"""赠送精英令后短轮询积分;1 个精英令约等于 10 积分。"""
expected_gain = max(0, gift_count) * 10
target_points = baseline_points + expected_gain if baseline_points is not None else None
result["gift_points_baseline"] = baseline_points
result["gift_points_expected_gain"] = expected_gain
result["gift_points_target"] = target_points
last_points = None
refresh_result: dict = {}
for index in range(1, DOUYU_GIFT_POINTS_REFRESH_TIMES + 1):
refresh_result = self._refresh_account_points(client, account, cookie, ctn=ctn)
db.commit()
last_points = refresh_result["points"]
result.update(refresh_result)
result["gift_points_refresh_count"] = index
if target_points is None or (last_points is not None and last_points >= target_points):
result["gift_points_confirmed"] = target_points is None or last_points is not None
return refresh_result
if index < DOUYU_GIFT_POINTS_REFRESH_TIMES:
self._update_task_progress(
db,
task,
"running",
f"赠送精英令成功,等待积分同步(当前 {last_points if last_points is not None else '-'},预期 {target_points}",
result,
)
if self._stop.wait(DOUYU_GIFT_POINTS_REFRESH_INTERVAL):
break
result["gift_points_confirmed"] = False
result["points"] = last_points
return refresh_result
def _task_payload(self, task: DouyuTask) -> dict:
result = task.result if isinstance(task.result, dict) else {}
payload = result.get("payload") if isinstance(result.get("payload"), dict) else {}
return {**payload, **self.payload}
def _client(self, cookie: str) -> DouyuActivityClient:
return DouyuActivityClient(cookie, logger=lambda msg: self._push_log("debug", msg))
def _fetch_bind_info_candidates(
self,
client: DouyuActivityClient,
aliases: list[str],
) -> list[dict]:
"""按多个 actAlias 查询绑定信息,保留成功结果。"""
results: list[dict] = []
for alias in aliases:
if not alias:
continue
try:
info = client.bind_info(alias, v2=True)
except DouyuActivityError as exc:
self._push_log("warning", f"查询绑定信息失败 act={alias}: {exc}")
continue
info = {**info, "act_alias": alias}
results.append(info)
return results
def _pick_bind_info(
self,
candidates: list[dict],
*,
baseline_role_name: str = "",
baseline_is_bound_act: bool = False,
prefer_pending: bool = True,
prefer_aliases: list[str] | None = None,
) -> dict | None:
"""从多个 alias 结果里挑最有用的绑定信息。
- prefer_pending=True:优先选“待确认/新扫码角色”(通常来自 cjm)
- prefer_aliases:在同等条件下优先指定 alias(如活动当前绑定)
"""
if not candidates:
return None
def _alias_rank(info: dict) -> int:
alias = str(info.get("act_alias") or "")
if not prefer_aliases:
return 0
try:
return prefer_aliases.index(alias)
except ValueError:
return len(prefer_aliases) + 1
ranked = sorted(enumerate(candidates), key=lambda item: (_alias_rank(item[1]), item[0]))
ordered = [item[1] for item in ranked]
if prefer_pending:
for info in ordered:
if self._is_pending_role(
info,
baseline_role_name=baseline_role_name,
baseline_is_bound_act=baseline_is_bound_act,
):
return info
for info in ordered:
if str(info.get("role_name") or "").strip():
return info
return ordered[0]
def _pick_baseline_bind_info(
self,
candidates: list[dict],
config: dict,
) -> dict | None:
"""选“扫码前当前已绑定角色”作为 baseline。
活动 alias20260120QYOOB)只反映当前已绑定;
cjm 才是换绑最新态。baseline 应优先活动 alias 的已绑定结果,
避免把 cjm 上的待确认新角色误当成扫码前旧角色。
"""
if not candidates:
return None
prefer = []
for key in ("confirm_act_alias", "bind_act_alias", "legacy_act_alias"):
alias = str(config.get(key) or "").strip()
if alias and alias not in prefer:
prefer.append(alias)
# 先找“已绑定且有角色”的活动结果
for alias in prefer:
for info in candidates:
if str(info.get("act_alias") or "") != alias:
continue
if self._is_bound_act(info) and str(info.get("role_name") or "").strip():
return info
# 再退回任意有角色的结果(仍按活动 alias 优先)
return self._pick_bind_info(
candidates,
prefer_pending=False,
prefer_aliases=prefer,
)
def _wait_bind_role_result(
self,
client: DouyuActivityClient,
db: Session,
task: DouyuTask,
account: Account,
query_aliases: list[str],
result: dict,
baseline_role_name: str = "",
baseline_is_bound_act: bool = False,
) -> tuple[str, dict]:
"""生成二维码后轮询绑定信息,直到识别到待确认角色、超时或停止。"""
deadline = time.monotonic() + DOUYU_BIND_ROLE_POLL_SECONDS
result["bind_polling"] = True
result["bind_phase"] = result.get("bind_phase") or "waiting_scan"
result["bind_ready_for_confirm"] = False
result["bind_confirmed"] = False
result["baseline_role_name"] = baseline_role_name
result["baseline_is_bound_act"] = baseline_is_bound_act
result["query_act_aliases"] = query_aliases
self._push_log(
"info",
"开始轮询绑定角色 "
f"aliases={','.join(query_aliases) or '-'} "
f"baseline={baseline_role_name or '-'} bound={1 if baseline_is_bound_act else 0}",
)
self._update_task_progress(db, task, "running", "已生成绑定二维码,等待扫码绑定", result)
poll_count = 0
last_summary = ""
while not self._stop.is_set() and time.monotonic() < deadline:
if self._stop.wait(DOUYU_BIND_ROLE_POLL_INTERVAL):
break
candidates = self._fetch_bind_info_candidates(client, query_aliases)
if not candidates:
result["bind_poll_error"] = "所有 actAlias 查询绑定信息失败"
self._update_task_progress(db, task, "running", "等待绑定角色同步: 查询失败", result)
continue
poll_count += 1
bind_info = self._pick_bind_info(
candidates,
baseline_role_name=baseline_role_name,
baseline_is_bound_act=baseline_is_bound_act,
prefer_pending=True,
prefer_aliases=query_aliases,
) or candidates[0]
snapshot = self._bind_snapshot(bind_info)
is_pending_role = self._is_pending_role(
bind_info,
baseline_role_name=baseline_role_name,
baseline_is_bound_act=baseline_is_bound_act,
)
query_alias = str(bind_info.get("act_alias") or "")
summary = (
f"act={query_alias or '-'} "
f"{self._format_bind_summary(bind_info, pending=is_pending_role)}"
)
# 字段变化或每 3 次打印一次,避免刷屏但仍能看到过程
if summary != last_summary or poll_count == 1 or poll_count % 3 == 0:
self._push_log("info", f"轮询绑定#{poll_count}: {summary}")
last_summary = summary
# 注意:未识别到新角色时,不要把当前已绑定角色写进 role_name,
# 否则前端会把旧角色误当成“待确认角色/查询结果”。
if is_pending_role:
result.update({
**snapshot,
"act_alias": query_alias,
"query_act_alias": query_alias,
"bind_ready_for_confirm": True,
"bind_confirmed": False,
"bind_phase": "role_ready",
"bind_polling": False,
"poll_count": poll_count,
"bind_summary": summary,
"bind_candidates": [
{
"act_alias": item.get("act_alias"),
"role_name": item.get("role_name"),
"is_bound_act": self._is_bound_act(item),
}
for item in candidates
],
})
role_name = snapshot["role_name"]
self._apply_bind_info_to_account(account, bind_info, "game_queried")
self._push_log("success", f"识别到待确认角色: {role_name} (act={query_alias})")
self._update_task_progress(
db,
task,
"running",
f"已识别角色: {role_name},待确认绑定",
result,
)
return role_name, result
if snapshot["role_name"] and snapshot["is_bound_act"]:
result["current_role_name"] = snapshot["role_name"]
result["current_area_name"] = snapshot["area_name"]
result["current_plat_name"] = snapshot["plat_name"]
result.update({
"bind_info": bind_info,
"act_alias": query_alias,
"query_act_alias": query_alias,
"is_bound_act": snapshot["is_bound_act"],
"is_bound_role": snapshot["is_bound_role"],
"is_bound_account": snapshot["is_bound_account"],
"need_bind_act": snapshot["need_bind_act"],
"need_bind_role": snapshot["need_bind_role"],
"change_role_wait_time": snapshot["change_role_wait_time"],
"can_change_role": snapshot["can_change_role"],
"role_name": "",
"area_name": "",
"plat_name": "",
"bind_ready_for_confirm": False,
"bind_confirmed": False,
"bind_phase": "waiting_scan",
"bind_polling": True,
"poll_count": poll_count,
"bind_summary": summary,
})
self._update_task_progress(db, task, "running", f"等待扫码绑定 ({summary})", result)
result["bind_polling"] = False
result["bind_ready_for_confirm"] = False
result["bind_phase"] = "stopped" if self._stop.is_set() else "role_timeout"
result["poll_count"] = poll_count
if last_summary:
result["bind_summary"] = last_summary
self._push_log(
"warning",
f"轮询结束 phase={result['bind_phase']} polls={poll_count} last={last_summary or '-'}",
)
return "", result
def _execute_refresh_goods(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
client = self._client(cookie)
result = client.list_goods(manual_id=config["manual_id"], rid=config["rid"])
goods = result["goods"]
self._upsert_goods(db, goods)
account.bind_status = account.bind_status or "active"
account.updated_at = datetime.now(timezone.utc)
self._mark_task(db, task, "success", f"已刷新商品 {len(goods)} 个", {"goods_count": len(goods), "goods": goods})
def _execute_refresh_esports_goods(
self,
db: Session,
task: DouyuTask,
account: Account,
cookie: str,
config: dict,
):
"""刷新电竞手册皮肤商城快照。"""
client = self._client(cookie)
result = client.list_esports_goods(
manual_id=str(config["esports_manual_id"]),
rid=str(config["room_id"]),
)
goods = result["goods"]
self._upsert_esports_goods(db, goods)
account.esports_bind_status = "esports_goods_refreshed"
account.updated_at = datetime.now(timezone.utc)
self._mark_task(
db,
task,
"success",
f"已刷新电竞皮肤 {len(goods)} 个",
{"goods_count": len(goods), "esports_store_score": result["score"], "goods": goods},
)
def _execute_get_bind_qr(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
client = self._client(cookie)
qr_act_alias = self._bind_qr_act_alias(config)
query_aliases = self._query_bind_act_aliases(config)
if not qr_act_alias and not query_aliases:
self._mark_task(db, task, "failed", "请先配置绑定活动 actAlias")
return
if qr_act_alias and qr_act_alias not in query_aliases:
query_aliases = [qr_act_alias, *query_aliases]
before_candidates = self._fetch_bind_info_candidates(client, query_aliases)
if not before_candidates:
self._mark_task(db, task, "failed", "查询绑定信息失败,请检查 Cookie 或 actAlias")
return
# baseline 用活动 alias 的“当前已绑定”;pending 检测优先 cjm 的换绑最新态
before = self._pick_baseline_bind_info(before_candidates, config) or before_candidates[0]
# 换绑冷却必须看活动当前绑定(QYOOB),不要用 cjm
cooldown_info = self._pick_change_wait_bind_info(before_candidates, config) or before
pending_before = self._pick_bind_info(
before_candidates,
baseline_role_name=str(before.get("role_name") or ""),
baseline_is_bound_act=self._is_bound_act(before),
prefer_pending=True,
prefer_aliases=query_aliases,
) or before
before_snapshot = self._bind_snapshot(before)
cooldown_snapshot = self._bind_snapshot(cooldown_info)
current_role_name = before_snapshot["role_name"] or cooldown_snapshot["role_name"]
wait_time = cooldown_snapshot["change_role_wait_time"]
self._push_log(
"info",
"绑定前状态 "
f"qr_act={qr_act_alias or '-'} query={','.join(query_aliases)} "
f"baseline_hit={before.get('act_alias') or '-'} {self._format_bind_summary(before)} "
f"cooldown_hit={cooldown_info.get('act_alias') or '-'} "
f"{self._format_bind_summary(cooldown_info)} "
f"pending_hit={pending_before.get('act_alias') or '-'} "
f"{self._format_bind_summary(pending_before)}",
)
# 生成二维码前强制检查换绑冷却:冷却中直接失败,绝不发码
if self._is_change_cooling(cooldown_info):
role_label = current_role_name or cooldown_snapshot["role_name"] or "当前角色"
wait_text = self._format_wait_time(wait_time) or "冷却中"
self._push_log(
"warning",
f"换绑冷却中,跳过生成二维码 role={role_label} wait={wait_time if wait_time is not None else '-'} "
f"can={cooldown_info.get('can_change_role')}",
)
result = {
"act_alias": qr_act_alias or cooldown_info.get("act_alias") or before.get("act_alias"),
"query_act_alias": cooldown_info.get("act_alias") or before.get("act_alias"),
"query_act_aliases": query_aliases,
"before_bind_info": before,
"cooldown_bind_info": cooldown_info,
**cooldown_snapshot,
"current_role_name": role_label if role_label != "当前角色" else current_role_name,
"current_area_name": cooldown_snapshot["area_name"] or before_snapshot["area_name"],
"current_plat_name": cooldown_snapshot["plat_name"] or before_snapshot["plat_name"],
"change_role_wait_time": wait_time,
"change_role_wait_text": self._format_wait_time(wait_time),
"bind_ready_for_confirm": False,
"bind_confirmed": bool(cooldown_snapshot["is_bound_act"] or before_snapshot["is_bound_act"]),
"bind_phase": "change_waiting",
"bind_polling": False,
}
self._apply_bind_info_to_account(
account,
cooldown_info if cooldown_snapshot["role_name"] else before,
"bind_confirmed" if result["bind_confirmed"] else "game_queried",
)
self._mark_task(
db,
task,
"failed",
f"{role_label} 暂不能换绑,剩余 {wait_text}",
result,
)
return
if not qr_act_alias:
self._mark_task(db, task, "failed", "请先配置绑定二维码活动 actAlias")
return
self._push_log(
"info",
f"换绑校验通过,开始生成二维码 act={qr_act_alias} "
f"role={current_role_name or '-'} wait={wait_time if wait_time is not None else 0}",
)
qr_result = client.get_bind_qr(qr_act_alias)
result = {
**qr_result,
"act_alias": qr_act_alias,
"query_act_aliases": query_aliases,
"before_bind_info": before,
"current_role_name": current_role_name,
"current_area_name": before_snapshot["area_name"],
"current_plat_name": before_snapshot["plat_name"],
"role_name": "",
"area_name": "",
"plat_name": "",
"bind_ready_for_confirm": False,
"bind_confirmed": False,
"bind_phase": "waiting_scan",
"bind_polling": True,
}
account.bind_status = "bind_qr_generated"
account.updated_at = datetime.now(timezone.utc)
# 关键:先把二维码 progress 出去,前端 running 期间即可弹窗扫码。
self._update_task_progress(db, task, "running", "已生成绑定二维码,等待扫码绑定", result)
role_name, result = self._wait_bind_role_result(
client,
db,
task,
account,
query_aliases,
result,
baseline_role_name=current_role_name,
baseline_is_bound_act=before_snapshot["is_bound_act"],
)
if role_name:
self._mark_task(db, task, "success", f"已识别角色: {role_name},待确认绑定", result)
return
if result.get("bind_phase") == "stopped":
self._mark_task(db, task, "stopped", "任务已停止", result)
return
last_summary = str(result.get("bind_summary") or "")
timeout_msg = "已生成绑定二维码,未检测到新扫码角色"
if current_role_name:
timeout_msg = f"{timeout_msg}(当前仍是 {current_role_name}"
if last_summary:
timeout_msg = f"{timeout_msg} | {last_summary}"
self._mark_task(db, task, "success", timeout_msg, result)
def _execute_confirm_bind(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
client = self._client(cookie)
confirm_alias = self._confirm_act_alias(config)
query_aliases = self._query_bind_act_aliases(config)
if confirm_alias and confirm_alias not in query_aliases:
query_aliases = [confirm_alias, *query_aliases]
if not query_aliases:
self._mark_task(db, task, "failed", "请先配置绑定活动 actAlias")
return
before_candidates = self._fetch_bind_info_candidates(client, query_aliases)
if not before_candidates:
self._mark_task(db, task, "failed", "查询绑定信息失败,请检查 Cookie 或 actAlias")
return
before = self._pick_bind_info(
before_candidates,
baseline_role_name="",
baseline_is_bound_act=False,
prefer_pending=True,
prefer_aliases=query_aliases,
) or before_candidates[0]
before_snapshot = self._bind_snapshot(before)
role_name = before_snapshot["role_name"]
query_alias = str(before.get("act_alias") or "")
self._push_log(
"info",
f"确认前状态 confirm_act={confirm_alias or '-'} hit={query_alias or '-'} "
f"{self._format_bind_summary(before)}",
)
if not role_name:
self._mark_task(
db,
task,
"failed",
"尚未识别到待确认角色,请先扫码完成绑定",
{
"act_alias": confirm_alias or query_alias,
"query_act_alias": query_alias,
"query_act_aliases": query_aliases,
"before_bind_info": before,
**before_snapshot,
"bind_ready_for_confirm": False,
"bind_confirmed": False,
"bind_phase": "waiting_role",
},
)
return
if before_snapshot["is_bound_act"]:
self._apply_bind_info_to_account(account, before, "bind_confirmed")
self._mark_task(
db,
task,
"success",
f"已绑定: {role_name}",
{
"act_alias": confirm_alias or query_alias,
"query_act_alias": query_alias,
"query_act_aliases": query_aliases,
"before_bind_info": before,
**before_snapshot,
"bind_ready_for_confirm": True,
"bind_confirmed": True,
"bind_phase": "confirmed",
},
)
return
# 确认接口优先用配置的确认 alias;没有则回退到命中查询的 alias
use_confirm_alias = confirm_alias or query_alias
confirm_result = client.confirm_bind(use_confirm_alias)
try:
after_candidates = self._fetch_bind_info_candidates(client, query_aliases)
after = self._pick_bind_info(
after_candidates,
baseline_role_name="",
baseline_is_bound_act=False,
prefer_pending=False,
) if after_candidates else None
if after is None:
raise DouyuActivityError("确认后回查绑定信息失败")
except DouyuActivityError as exc:
# 确认接口已成功时,回查失败仍按确认成功处理,但保留错误信息。
self._apply_bind_info_to_account(account, before, "bind_confirmed")
self._mark_task(
db,
task,
"success",
f"绑定成功: {role_name}",
{
"act_alias": use_confirm_alias,
"query_act_alias": query_alias,
"query_act_aliases": query_aliases,
"before_bind_info": before,
"confirm": confirm_result,
**before_snapshot,
"bind_ready_for_confirm": True,
"bind_confirmed": True,
"bind_phase": "confirmed",
"refresh_error": str(exc),
},
)
return
after_snapshot = self._bind_snapshot(after)
final_info = after if after_snapshot["role_name"] else before
final_snapshot = after_snapshot if after_snapshot["role_name"] else before_snapshot
final_role_name = final_snapshot["role_name"] or role_name
if not after_snapshot["is_bound_act"]:
self._apply_bind_info_to_account(account, final_info, "game_queried")
self._mark_task(
db,
task,
"failed",
f"确认绑定未生效: {final_role_name}",
{
"act_alias": use_confirm_alias,
"query_act_alias": after.get("act_alias") or query_alias,
"query_act_aliases": query_aliases,
"before_bind_info": before,
"confirm": confirm_result,
**final_snapshot,
"bind_ready_for_confirm": True,
"bind_confirmed": False,
"bind_phase": "confirm_failed",
"after_bind_info": after,
},
)
return
self._apply_bind_info_to_account(account, final_info, "bind_confirmed")
self._mark_task(
db,
task,
"success",
f"绑定成功: {final_role_name}",
{
"act_alias": use_confirm_alias,
"query_act_alias": after.get("act_alias") or query_alias,
"query_act_aliases": query_aliases,
"before_bind_info": before,
"confirm": confirm_result,
"after_bind_info": after,
**final_snapshot,
"bind_ready_for_confirm": True,
"bind_confirmed": True,
"bind_phase": "confirmed",
},
)
def _execute_create_elite_qr(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
client = self._client(cookie)
ctn = str(self._task_payload(task).get("ctn") or "")
if not ctn:
ctn = client.acf_ccn(refresh_subscribe=True)
act_alias = self._confirm_act_alias(config) or self._bind_qr_act_alias(config)
if not act_alias:
self._mark_task(db, task, "failed", "请先配置开通宝典活动 actAlias")
return
result = client.create_elite_qr(
ctn=ctn,
act_alias=act_alias,
amount=int(config["elite_amount"]),
room_id=str(config["room_id"]),
)
account.bind_status = "elite_qr_created"
account.updated_at = datetime.now(timezone.utc)
self._update_task_progress(db, task, "running", "精英宝典支付码已生成,等待开通到账", result)
opened = self._wait_points_after_payment(db, task, account, client, cookie, ctn, result)
if self._stop.is_set():
self._mark_task(db, task, "stopped", "任务已停止", result)
return
if opened:
account.bind_status = "elite_opened"
account.updated_at = datetime.now(timezone.utc)
self._mark_task(db, task, "success", f"精英宝典已开通,积分: {account.points}", result)
return
self._mark_task(
db,
task,
"failed",
f"未检测到精英宝典开通到账,当前积分: {account.points if account.points is not None else '-'}",
result,
)
def _execute_prepare_esports_bind(
self,
db: Session,
task: DouyuTask,
account: Account,
cookie: str,
config: dict,
):
"""打开电竞手册绑定面板:查活动状态、当前角色和换绑冷却。"""
client = self._client(cookie)
act_alias = str(config.get("esports_act_alias") or "").strip()
if not act_alias:
self._mark_task(db, task, "failed", "请先配置电竞手册活动 actAlias")
return
state = self._esports_bind_state(client, act_alias)
esports_bound = state["esports_bound"]
role_text = self._esports_role_text(state)
self._apply_esports_bind_info_to_account(
account,
state,
"esports_bound" if esports_bound else ("esports_bind_ready" if role_text else "game_not_bound"),
)
# 二维码用于重新选择角色,不应因为已有角色或换绑冷却而隐藏。
# 冷却是否允许最终由 actBind 返回结果决定。
qr_result = client.get_esports_bind_qr(act_alias)
result = {
**state,
**qr_result,
"esports_bind_dialog": True,
"can_open_role_selector": True,
}
if esports_bound:
message = f"电竞手册已绑定: {role_text or '-'},可扫码换绑"
elif role_text:
message = f"当前角色: {role_text},可扫码切换角色或直接完成绑定"
else:
message = "请扫码选择游戏角色,完成后查询最新角色"
self._mark_task(db, task, "success", message, result)
def _execute_get_esports_bind_qr(
self,
db: Session,
task: DouyuTask,
account: Account,
cookie: str,
config: dict,
):
"""生成电竞手册切换角色用的腾讯入口。"""
client = self._client(cookie)
act_alias = str(config.get("esports_act_alias") or "").strip()
if not act_alias:
self._mark_task(db, task, "failed", "请先配置电竞手册活动 actAlias")
return
state = self._esports_bind_state(client, act_alias)
qr_result = client.get_esports_bind_qr(act_alias)
result = {
**state,
**qr_result,
"esports_bind_dialog": True,
"can_open_role_selector": True,
"bind_phase": "switching_role",
}
account.esports_bind_status = "esports_role_switching"
account.updated_at = datetime.now(timezone.utc)
self._mark_task(
db,
task,
"success",
"请在腾讯页面选择角色,返回后查询最新角色",
result,
)
def _execute_confirm_esports_bind(
self,
db: Session,
task: DouyuTask,
account: Account,
cookie: str,
config: dict,
):
"""通过 actBind 确认电竞手册绑定,并回读唯一状态接口确认结果。"""
client = self._client(cookie)
act_alias = str(config.get("esports_act_alias") or "").strip()
if not act_alias:
self._mark_task(db, task, "failed", "请先配置电竞手册活动 actAlias")
return
account.esports_bind_status = "esports_bind_confirming"
account.updated_at = datetime.now(timezone.utc)
self._update_task_progress(
db,
task,
"running",
"正在确认电竞手册绑定",
{
"bind_phase": "confirming",
"bind_confirmed": False,
},
)
confirm_result = client.confirm_esports_bind(
act_alias,
room_id=str(config.get("room_id") or "9263298"),
)
after_state = self._esports_bind_state(client, act_alias)
result = {
**after_state,
"confirm": confirm_result,
"bind_phase": "confirmed" if after_state["esports_bound"] else "confirm_failed",
}
if not after_state["esports_bound"]:
self._apply_esports_bind_info_to_account(account, after_state, "esports_bind_ready")
self._mark_task(
db,
task,
"failed",
f"电竞手册绑定未生效: {self._esports_role_text(after_state) or '-'}",
result,
)
return
self._apply_esports_bind_info_to_account(account, after_state, "esports_bound")
self._mark_task(
db,
task,
"success",
f"电竞手册绑定成功: {self._esports_role_text(after_state) or '-'}",
result,
)
def _execute_query_esports_game_name(
self,
db: Session,
task: DouyuTask,
account: Account,
cookie: str,
config: dict,
):
"""查询电竞手册活动返回的最新角色和换绑冷却状态。"""
client = self._client(cookie)
act_alias = str(config.get("esports_act_alias") or "").strip()
if not act_alias:
self._mark_task(db, task, "failed", "请先配置电竞手册活动 actAlias")
return
state = self._esports_bind_state(client, act_alias)
role_name = state["role_name"]
is_bound = state["esports_bound"]
self._apply_esports_bind_info_to_account(
account,
state,
"esports_bound" if is_bound else ("esports_bind_ready" if role_name else "game_not_bound"),
)
result = {
**state,
"esports_bind_dialog": True,
"can_open_role_selector": True,
}
if role_name:
message = f"最新角色: {self._esports_role_text(state)}"
else:
message = "未查询到游戏角色,请先切换角色"
self._mark_task(db, task, "success", message, result)
def _execute_create_esports_qr(
self,
db: Session,
task: DouyuTask,
account: Account,
cookie: str,
config: dict,
):
"""生成电竞手册支付二维码,并通过活动状态确认开通到账。"""
client = self._client(cookie)
ctn = str(self._task_payload(task).get("ctn") or "")
if not ctn:
ctn = client.acf_ccn(refresh_subscribe=True)
act_alias = str(config.get("esports_act_alias") or "").strip()
manual_id = str(config.get("esports_manual_id") or "").strip()
if not act_alias or not manual_id:
self._mark_task(db, task, "failed", "请先配置电竞手册活动 actAlias 和 manualID")
return
baseline_manual_type = None
baseline_manual_score = None
baseline_result: dict = {}
try:
baseline_result = self._refresh_esports_handbook(client, account, manual_id=manual_id)
baseline_manual_type = baseline_result["esports_manual_type"]
baseline_manual_score = baseline_result["esports_manual_score"]
db.commit()
except Exception as exc:
self._push_log("warning", f"生成电竞手册支付码前查询活动状态失败: {exc}")
if baseline_manual_type is not None and baseline_manual_type >= 1:
account.esports_bind_status = "esports_opened"
account.updated_at = datetime.now(timezone.utc)
self._mark_task(
db,
task,
"success",
f"电竞手册已开通,积分: {baseline_manual_score if baseline_manual_score is not None else '-'}",
{**baseline_result, "esports_opened": True, "payment_polling": False},
)
return
result = client.create_esports_qr(
ctn=ctn,
act_alias=act_alias,
amount=int(config["esports_amount"]),
room_id=str(config["room_id"]),
)
result.update(baseline_result)
account.esports_bind_status = "esports_qr_created"
account.updated_at = datetime.now(timezone.utc)
self._update_task_progress(db, task, "running", "电竞手册支付码已生成,等待开通到账", result)
opened = self._wait_esports_open_after_payment(
db,
task,
account,
client,
manual_id=manual_id,
result=result,
baseline_manual_type=baseline_manual_type,
baseline_manual_score=baseline_manual_score,
)
if self._stop.is_set():
self._mark_task(db, task, "stopped", "任务已停止", result)
return
if opened:
account.esports_bind_status = "esports_opened"
account.updated_at = datetime.now(timezone.utc)
self._mark_task(
db,
task,
"success",
f"电竞手册已开通,积分: {account.esports_points if account.esports_points is not None else '-'}",
result,
)
return
self._mark_task(
db,
task,
"failed",
"未检测到电竞手册开通到账"
f",类型: {result.get('esports_manual_type', '-')}"
f"积分: {result.get('esports_manual_score', '-')}",
result,
)
def _execute_query_esports_points(
self,
db: Session,
task: DouyuTask,
account: Account,
cookie: str,
config: dict,
):
"""查询电竞手册积分。"""
manual_id = str(config.get("esports_manual_id") or "").strip()
if not manual_id:
self._mark_task(db, task, "failed", "请先配置电竞手册 manualID")
return
client = self._client(cookie)
result = self._refresh_esports_handbook(client, account, manual_id=manual_id)
account.esports_bind_status = "esports_points_queried"
account.updated_at = datetime.now(timezone.utc)
points = result["esports_points"]
self._mark_task(db, task, "success", f"电竞积分: {points if points is not None else '-'}", result)
def _execute_donate_esports_gift(
self,
db: Session,
task: DouyuTask,
account: Account,
cookie: str,
config: dict,
*,
gift_name: str,
config_gift_id_key: str,
config_skin_id_key: str,
):
"""赠送电竞手册任务礼物并刷新独立积分。"""
payload = self._task_payload(task)
try:
gift_count = max(1, int(payload.get("gift_count") or payload.get("count") or 1))
except (TypeError, ValueError):
self._mark_task(db, task, "failed", "赠送数量必须是正整数")
return
manual_id = str(config.get("esports_manual_id") or "").strip()
gift_id = str(payload.get("gift_id") or config.get(config_gift_id_key) or "").strip()
skin_id = str(payload.get("skin_id") or config.get(config_skin_id_key) or "").strip()
room_id = str(payload.get("room_id") or config.get("room_id") or "").strip()
if not manual_id or not gift_id or not skin_id or not room_id:
self._mark_task(db, task, "failed", "请先完整配置电竞手册、房间和礼物参数")
return
client = self._client(cookie)
baseline_points = account.esports_points
try:
baseline = self._refresh_esports_handbook(client, account, manual_id=manual_id)
baseline_points = baseline["esports_points"]
db.commit()
except Exception as exc:
self._push_log("warning", f"赠送{gift_name}前刷新电竞积分失败: {exc}")
result = client.donate_esports_gift(
gift_name=gift_name,
gift_count=gift_count,
room_id=room_id,
gift_id=gift_id,
skin_id=skin_id,
)
result.update(
{
"gift_name": gift_name,
"gift_id": gift_id,
"skin_id": skin_id,
"gift_count": gift_count,
"esports_points_baseline": baseline_points,
}
)
refresh_errors = []
try:
result.update(self._refresh_account_gold_balance(client, account))
except Exception as exc:
refresh_errors.append(f"鱼翅余额: {exc}")
try:
points_result = self._refresh_esports_handbook(client, account, manual_id=manual_id)
result.update(points_result)
result["esports_points_after_gift"] = points_result["esports_points"]
result["esports_points_changed"] = (
baseline_points is not None
and points_result["esports_points"] is not None
and points_result["esports_points"] != baseline_points
)
except Exception as exc:
refresh_errors.append(f"电竞积分: {exc}")
if refresh_errors:
result["refresh_errors"] = refresh_errors
account.esports_bind_status = "esports_gift_donated"
account.updated_at = datetime.now(timezone.utc)
message = f"赠送{gift_name}成功: {gift_count}"
if account.gold_balance is not None:
message += f",鱼翅余额: {account.gold_balance}"
if account.esports_points is not None:
message += f",电竞积分: {account.esports_points}"
self._mark_task(db, task, "success", message, result)
def _execute_donate_esports_chicken_gift(
self,
db: Session,
task: DouyuTask,
account: Account,
cookie: str,
config: dict,
):
"""赠送冠军鸡腿。"""
self._execute_donate_esports_gift(
db,
task,
account,
cookie,
config,
gift_name="冠军鸡腿",
config_gift_id_key="esports_chicken_gift_id",
config_skin_id_key="esports_chicken_skin_id",
)
def _execute_donate_esports_firework_gift(
self,
db: Session,
task: DouyuTask,
account: Account,
cookie: str,
config: dict,
):
"""赠送冠军烟花。"""
self._execute_donate_esports_gift(
db,
task,
account,
cookie,
config,
gift_name="冠军烟花",
config_gift_id_key="esports_firework_gift_id",
config_skin_id_key="esports_firework_skin_id",
)
def _execute_create_gold_qr(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
payload = self._task_payload(task)
amount = int(payload.get("amount") or payload.get("gold_amount") or 1)
client = self._client(cookie)
baseline_gold = account.gold_balance
try:
baseline = self._refresh_account_gold_balance(client, account)
baseline_gold = baseline["gold_balance"]
db.commit()
except Exception as exc:
self._push_log("warning", f"生成鱼翅码前刷新余额失败: {exc}")
result = client.create_gold_qr(amount=amount, pay_type=int(config["gold_pay_type"]))
account.bind_status = "gold_qr_created"
account.updated_at = datetime.now(timezone.utc)
self._update_task_progress(db, task, "running", f"鱼翅 {amount} 元支付码已生成,等待到账", result)
recharged = self._wait_gold_balance_after_payment(db, task, account, client, result, baseline_gold)
if self._stop.is_set():
self._mark_task(db, task, "stopped", "任务已停止", result)
return
if recharged:
account.bind_status = "gold_recharged"
account.updated_at = datetime.now(timezone.utc)
self._mark_task(
db,
task,
"success",
f"鱼翅已到账,当前余额: {account.gold_balance if account.gold_balance is not None else '-'}",
result,
)
return
self._mark_task(
db,
task,
"failed",
f"未检测到鱼翅到账,当前余额: {account.gold_balance if account.gold_balance is not None else '-'}",
result,
)
def _execute_donate_elite_gift(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
payload = self._task_payload(task)
gift_count = int(payload.get("gift_count") or payload.get("count") or 1)
client = self._client(cookie)
ctn = None
baseline_points = account.points
try:
ctn = client.acf_ccn(refresh_subscribe=False)
baseline_result = self._refresh_account_points(client, account, cookie, ctn=ctn)
db.commit()
baseline_points = baseline_result["points"]
except Exception as exc:
self._push_log("warning", f"赠送精英令前刷新积分失败: {exc}")
result = client.donate_elite_gift(
gift_count=gift_count,
room_id=str(payload.get("room_id") or config["room_id"]),
gift_id=str(payload.get("gift_id") or config["gift_id"]),
skin_id=str(payload.get("skin_id") or config["skin_id"]),
)
result["gift_points_baseline"] = baseline_points
refresh_errors = []
try:
result.update(self._refresh_account_gold_balance(client, account))
except Exception as exc:
refresh_errors.append(f"鱼翅余额: {exc}")
try:
result.update(
self._refresh_points_after_elite_gift(
db,
task,
account,
client,
cookie,
ctn,
result,
baseline_points,
gift_count,
)
)
except Exception as exc:
refresh_errors.append(f"积分: {exc}")
if refresh_errors:
result["refresh_errors"] = refresh_errors
account.bind_status = "gift_donated"
account.updated_at = datetime.now(timezone.utc)
message = f"赠送精英令成功: {gift_count}"
if account.gold_balance is not None:
message += f",鱼翅余额: {account.gold_balance}"
if account.points is not None:
message += f",积分: {account.points}"
if result.get("gift_points_target") is not None and not result.get("gift_points_confirmed"):
message += f"(未确认涨到 {result['gift_points_target']}"
self._mark_task(db, task, "success", message, result)
def _execute_query_points(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
client = self._client(cookie)
ctn = client.acf_ccn(refresh_subscribe=False)
result = self._refresh_account_points(client, account, cookie, ctn=ctn)
account.bind_status = "points_queried"
account.updated_at = datetime.now(timezone.utc)
points = result["points"]
self._mark_task(db, task, "success", f"积分: {points if points is not None else '-'}", result)
def _execute_exchange_goods(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
import time as time_mod
payload = self._task_payload(task)
commodity_id = str(payload.get("commodity_id") or payload.get("commodityId") or "").strip()
if not commodity_id:
self._mark_task(db, task, "failed", "请选择兑换商品")
return
client = self._client(cookie)
ctn = client.acf_ccn(refresh_subscribe=False)
result = None
last_error = ""
for attempt in range(8 + 1):
if self._stop.is_set():
self._mark_task(db, task, "stopped", "任务已停止")
return
try:
result = client.exchange_goods(
manual_id=str(config["manual_id"]),
rid=str(config["rid"]),
commodity_id=commodity_id,
ctn=ctn,
)
break
except DouyuActivityError as exc:
last_error = str(exc)
if attempt >= 8:
self._mark_task(db, task, "failed", f"兑换失败(已重试{attempt}次): {last_error}")
return
error_lower = last_error.lower()
if any(kw in error_lower for kw in ("无效", "太快", "csrf")):
self._push_log("info", f" 重试 {attempt + 1}/8: {last_error},刷新 csrf_token...")
try:
token = client.csrf_token()
self._push_log("debug", f" csrf_token 已刷新: {token[:12]}...")
except Exception:
pass
else:
self._push_log("info", f" 重试 {attempt + 1}/8: {last_error}")
time_mod.sleep(0.3)
if result is None:
self._mark_task(db, task, "failed", f"兑换失败: {last_error}")
return
goods = (
db.query(DouyuGoodsSnapshot)
.filter(DouyuGoodsSnapshot.commodity_id == commodity_id)
.first()
)
account.bind_status = "goods_exchanged"
account.updated_at = datetime.now(timezone.utc)
self._mark_task(
db,
task,
"success",
f"兑换成功: {(goods.name if goods else '') or commodity_id}",
{
"goods": goods.raw if goods else None,
"game_name": account.game_name or "",
"game_channel": account.game_channel or "",
"account_name": account.nickname or account.username or account.uid or f"#{account.id}",
**result,
},
)
def _execute_exchange_esports_goods(
self,
db: Session,
task: DouyuTask,
account: Account,
cookie: str,
config: dict,
):
"""兑换电竞手册皮肤,并同步电竞积分。"""
payload = self._task_payload(task)
commodity_id = str(payload.get("commodity_id") or payload.get("commodityId") or "").strip()
if not commodity_id:
self._mark_task(db, task, "failed", "请选择电竞皮肤")
return
try:
quantity = max(1, int(payload.get("quantity") or payload.get("num") or 1))
except (TypeError, ValueError):
self._mark_task(db, task, "failed", "兑换数量必须是正整数")
return
manual_id = str(config.get("esports_manual_id") or "").strip()
room_id = str(config.get("room_id") or "").strip()
if not manual_id or not room_id:
self._mark_task(db, task, "failed", "请先配置电竞手册 manualID 和房间 ID")
return
client = self._client(cookie)
baseline_points = account.esports_points
try:
baseline = self._refresh_esports_handbook(client, account, manual_id=manual_id)
baseline_points = baseline["esports_points"]
db.commit()
except Exception as exc:
self._push_log("warning", f"兑换电竞皮肤前刷新积分失败: {exc}")
result = client.exchange_esports_goods(
manual_id=manual_id,
rid=room_id,
commodity_id=commodity_id,
quantity=quantity,
)
goods = (
db.query(DouyuEsportsGoodsSnapshot)
.filter(DouyuEsportsGoodsSnapshot.commodity_id == commodity_id)
.first()
)
result["goods"] = goods.raw if goods else None
result["esports_points_baseline"] = baseline_points
try:
points_result = self._refresh_esports_handbook(client, account, manual_id=manual_id)
result.update(points_result)
result["esports_points_after_exchange"] = points_result["esports_points"]
except Exception as exc:
result["esports_points_refresh_error"] = str(exc)
account.esports_bind_status = "esports_goods_exchanged"
account.updated_at = datetime.now(timezone.utc)
name = (goods.name if goods else "") or commodity_id
message = f"兑换电竞皮肤成功: {name}"
if quantity > 1:
message += f" x{quantity}"
if account.esports_points is not None:
message += f",电竞积分: {account.esports_points}"
result["game_name"] = account.esports_game_name or ""
result["game_channel"] = account.esports_game_channel or ""
result["account_name"] = account.nickname or account.username or account.uid or f"#{account.id}"
self._mark_task(db, task, "success", message, result)
def _execute_query_game_name(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
client = self._client(cookie)
query_aliases = self._query_bind_act_aliases(config)
if not query_aliases:
self._mark_task(db, task, "failed", "请先配置绑定活动 actAlias")
return
candidates = self._fetch_bind_info_candidates(client, query_aliases)
if not candidates:
self._mark_task(db, task, "failed", "查询绑定信息失败,请检查 Cookie 或 actAlias")
return
bind_info = self._pick_bind_info(
candidates,
prefer_pending=True,
prefer_aliases=query_aliases,
) or candidates[0]
snapshot = self._bind_snapshot(bind_info)
role_name = snapshot["role_name"]
source_alias = str(bind_info.get("act_alias") or "")
summary = f"act={source_alias or '-'} {self._format_bind_summary(bind_info)}"
self._push_log(
"info",
"查询角色 "
f"aliases={','.join(query_aliases)} hit={source_alias or '-'} "
f"{self._format_bind_summary(bind_info)}",
)
self._apply_bind_info_to_account(account, bind_info, "game_queried" if role_name else "game_not_bound")
pending = bool(role_name) and not snapshot["is_bound_act"]
result = {
"act_alias": source_alias,
"query_act_alias": source_alias,
"query_act_aliases": query_aliases,
**snapshot,
"bind_ready_for_confirm": pending,
"bind_confirmed": snapshot["is_bound_act"],
"bind_phase": (
"confirmed" if snapshot["is_bound_act"]
else ("role_ready" if role_name else "waiting_role")
),
"bind_summary": summary,
"bind_candidates": [
{
"act_alias": item.get("act_alias"),
"role_name": item.get("role_name"),
"is_bound_act": self._is_bound_act(item),
}
for item in candidates
],
}
if not role_name:
message = f"未获取到游戏名 | {summary}"
elif snapshot["is_bound_act"]:
message = f"当前已绑定: {role_name} | {summary}"
else:
message = f"待确认角色: {role_name} | {summary}"
self._mark_task(db, task, "success" if role_name else "failed", message, result)
def _pick_change_wait_bind_info(self, candidates: list[dict], config: dict) -> dict | None:
"""换绑倒计时优先看活动当前绑定(QYOOB),不是 cjm 换绑最新态。"""
if not candidates:
return None
# 1) 优先活动 alias 上有角色的结果
baseline = self._pick_baseline_bind_info(candidates, config)
if baseline is not None:
wait = self._to_int(baseline.get("change_role_wait_time"))
if wait is not None:
return baseline
# 2) 任意带 wait_time 的结果里取 wait 最大的
with_wait = []
for info in candidates:
wait = self._to_int(info.get("change_role_wait_time"))
if wait is not None:
with_wait.append((wait, info))
if with_wait:
with_wait.sort(key=lambda item: item[0], reverse=True)
return with_wait[0][1]
# 3) 回退 baseline / 首个有角色
return baseline or self._pick_bind_info(
candidates,
prefer_pending=False,
prefer_aliases=[
str(config.get("confirm_act_alias") or "").strip(),
str(config.get("bind_act_alias") or "").strip(),
str(config.get("legacy_act_alias") or "").strip(),
],
)
def _execute_query_change_bind_time(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
client = self._client(cookie)
query_aliases = self._query_bind_act_aliases(config)
if not query_aliases:
self._mark_task(db, task, "failed", "请先配置绑定活动 actAlias")
return
candidates = self._fetch_bind_info_candidates(client, query_aliases)
if not candidates:
self._mark_task(db, task, "failed", "查询绑定信息失败,请检查 Cookie 或 actAlias")
return
# 换绑时间看“当前已绑定”活动态,不要优先 cjm(cjm 常无 changeRoleWaitTime
bind_info = self._pick_change_wait_bind_info(candidates, config) or candidates[0]
snapshot = self._bind_snapshot(bind_info)
wait_time = snapshot["change_role_wait_time"]
account.change_role_wait_time = wait_time
account.bind_status = "change_time_queried"
account.updated_at = datetime.now(timezone.utc)
source_alias = str(bind_info.get("act_alias") or "")
wait_text = self._format_wait_time(wait_time)
can_change = snapshot["can_change_role"]
role_name = snapshot["role_name"] or "-"
if wait_time is None:
if can_change is False:
status_text = "不可换绑(接口未返回倒计时)"
elif can_change is True:
status_text = "可换绑"
else:
status_text = "未返回换绑倒计时"
elif wait_time <= 0:
status_text = "可换绑"
wait_text = "0"
else:
status_text = f"剩余 {wait_text}"
result = {
"act_alias": source_alias,
"query_act_alias": source_alias,
"query_act_aliases": query_aliases,
**snapshot,
"change_role_wait_text": wait_text,
"bind_ready_for_confirm": bool(snapshot["role_name"]) and not snapshot["is_bound_act"],
"bind_confirmed": snapshot["is_bound_act"],
"bind_summary": f"act={source_alias or '-'} {self._format_bind_summary(bind_info)}",
"bind_candidates": [
{
"act_alias": item.get("act_alias"),
"role_name": item.get("role_name"),
"is_bound_act": self._is_bound_act(item),
"change_role_wait_time": self._to_int(item.get("change_role_wait_time")),
"can_change_role": item.get("can_change_role"),
}
for item in candidates
],
}
self._push_log(
"info",
f"查询换绑时间 hit={source_alias or '-'} role={role_name} "
f"wait={wait_time if wait_time is not None else '-'} can={can_change}",
)
self._mark_task(
db,
task,
"success",
f"{role_name} {status_text} | act={source_alias or '-'}",
result,
)
def _execute_query_limited_goods(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
client = self._client(cookie)
result = client.query_limited_goods(manual_id=str(config["manual_id"]), rid=str(config["rid"]))
limited = result["limited_goods"]
names = [str(item.get("commodityName") or "") for item in limited if item.get("commodityName")]
message = "无限制商品" if not names else f"限兑 {len(names)} 个: {', '.join(names[:5])}"
account.bind_status = "limited_goods_queried"
account.updated_at = datetime.now(timezone.utc)
self._mark_task(db, task, "success", message, {"limited_count": len(limited), "limited_goods": limited})
def _execute_query_gold_balance(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
client = self._client(cookie)
result = self._refresh_account_gold_balance(client, account)
account.bind_status = "gold_balance_queried"
account.updated_at = datetime.now(timezone.utc)
self._mark_task(
db,
task,
"success",
f"鱼翅余额: {account.gold_balance if account.gold_balance is not None else '-'}",
result,
)
def _execute_query_exchange_records(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
client = self._client(cookie)
result = client.exchange_records(manual_id=str(config["manual_id"]))
records = result["records"]
account.bind_status = "exchange_records_queried"
account.updated_at = datetime.now(timezone.utc)
self._mark_task(db, task, "success", f"兑换记录 {len(records)} 条" if records else "暂无兑换记录", result)
def _execute_prefetch_csrf_token(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
client = self._client(cookie)
token = client.csrf_token()
account.bind_status = "csrf_token_ready"
account.updated_at = datetime.now(timezone.utc)
self._mark_task(db, task, "success", "获取 csrf_token 成功", {"csrf_token": token, "cookie": client.cookie})
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
task.status = "running"
task.message = "执行中"
worker_db.commit()
with self._counter_lock:
self._started += 1
current = self._started
self._push_log("info", f"[{current}/{total}] 开始: {self._account_name(account)}")
cookie = latest_success_cookie(worker_db, account.id)
if not cookie:
self._mark_task(worker_db, task, "failed", "账号没有成功登录 Cookie")
self._push_log("warning", f"[{current}] {self._account_name(account)} 无 Cookie")
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,
"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(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:
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()
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:
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()
class DouyuBatchRegistry:
"""管理运行中的斗鱼任务批次。"""
def __init__(self):
self._batches: dict[str, dict] = {}
def register(self, batch_id: str, log_queue: asyncio.Queue,
loop: asyncio.AbstractEventLoop, runner: DouyuBatchRunner):
self._batches[batch_id] = {
"log_queue": log_queue,
"loop": loop,
"runner": runner,
"finished": False,
"updated_at": time.time(),
}
def get(self, batch_id: str):
return self._batches.get(batch_id)
def pop(self, batch_id: str):
return self._batches.pop(batch_id, None)
def mark_finished(self, batch_id: str):
if batch_id in self._batches:
self._batches[batch_id]["finished"] = True
self._batches[batch_id]["updated_at"] = time.time()
def active_ids(self) -> set[str]:
return {
batch_id
for batch_id, info in self._batches.items()
if not info.get("finished")
}
douyu_batch_registry = DouyuBatchRegistry()