"""斗鱼活动任务批次执行器。""" 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, DouyuXpdGoodsSnapshot from .douyu_service import ( DOUYU_CONFIG_FIELDS, account_uid, douyu_config_value, ensure_douyu_config, latest_success_cookie, douyu_task_payload, update_account_profile_from_cookie, ) DOUYU_LEGACY_BIND_ACT_ALIAS = "20250213NQCYX" DOUYU_BIND_ROLE_POLL_SECONDS = 65 DOUYU_BIND_ROLE_POLL_INTERVAL = 5 DOUYU_XPD_BIND_POLL_SECONDS = 300 DOUYU_XPD_BIND_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 DOUYU_CONFIRM_EFFECT_POLL_TIMES = 3 DOUYU_CONFIRM_EFFECT_POLL_INTERVAL = 3 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) @classmethod def _action_act_aliases(cls, config: dict) -> list[str]: """当前活动动作 alias;不包含只用于查询最新扫码态的 legacy/cjm。""" aliases: list[str] = [] for key in ("confirm_act_alias", "bind_act_alias"): alias = cls._action_act_alias(config, key) if alias and alias not in aliases: aliases.append(alias) return aliases @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() if not role_name or not cls._is_bound_act(bind_info): 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" can={snap['can_change_role']}" 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 _push_task_event(self, task: DouyuTask) -> None: """向批次 WS 推送任务状态事件(level=task),前端即时更新不依赖轮询。""" if not self.log_queue or not self.loop: return try: payload = douyu_task_payload(task) except Exception: logger.exception("[douyu] 推送任务状态失败: task_id={}", task.id) return event = { "level": "task", "message": "", "task": payload, } asyncio.run_coroutine_threadsafe(self.log_queue.put(event), self.loop) 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() self._push_task_event(task) 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() self._push_task_event(task) 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 _upsert_xpd_goods(self, db: Session, goods: list[dict]) -> None: """同步和平小店商品快照,移除上一次热门抢购等遗留商品。""" now = datetime.now(timezone.utc) commodity_ids = { str(raw.get("commodity_id") or raw.get("iGoodsId") or "") for raw in goods } commodity_ids.discard("") query = db.query(DouyuXpdGoodsSnapshot) if commodity_ids: query.filter(~DouyuXpdGoodsSnapshot.commodity_id.in_(commodity_ids)).delete( synchronize_session=False, ) else: query.delete(synchronize_session=False) for raw in goods: commodity_id = str(raw.get("commodity_id") or raw.get("iGoodsId") or "") if not commodity_id: continue row = ( db.query(DouyuXpdGoodsSnapshot) .filter(DouyuXpdGoodsSnapshot.commodity_id == commodity_id) .first() ) if row is None: row = DouyuXpdGoodsSnapshot(commodity_id=commodity_id) db.add(row) row.name = str(raw.get("name") or raw.get("sGoodsName") or "") row.price = self._to_int(raw.get("price") or raw.get("iPrice")) row.org_price = self._to_int(raw.get("org_price") or raw.get("iOrgPrice")) row.category = str(raw.get("category") or raw.get("iCategoryId") or "") goods_left = raw.get("goods_left") if goods_left is None: goods_left = raw.get("iGoodsLeft") row.goods_left = self._to_int(goods_left) 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} self._push_log( "info", f"绑定信息 act={alias} {self._format_bind_summary(info)}", ) 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_current_bound_info( self, candidates: list[dict], config: dict, *, extra_prefer_aliases: list[str] | None = None, ) -> dict | None: """选当前活动已生效绑定,避免把 legacy/cjm 的待确认态当成当前角色。""" if not candidates: return None prefer = [] for alias in [*(extra_prefer_aliases or []), *self._action_act_aliases(config)]: alias = str(alias or "").strip() if alias and alias not in prefer: prefer.append(alias) bound = [ info for info in candidates if self._is_bound_act(info) and str(info.get("role_name") or "").strip() and str(info.get("act_alias") or "") in prefer ] if not bound: return None return sorted(bound, key=lambda info: prefer.index(str(info.get("act_alias") or "")))[0] def _pick_baseline_bind_info( self, candidates: list[dict], config: dict, ) -> dict | None: """选“扫码前当前已绑定角色”作为 baseline。 活动 alias(20260120QYOOB)只反映当前已绑定; cjm 才是换绑最新态。baseline 应优先活动 alias 的已绑定结果, 避免把 cjm 上的待确认新角色误当成扫码前旧角色。 """ if not candidates: return None return self._pick_current_bound_info(candidates, config) 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"] # 待确认角色只回传前端展示,不写入账号表,避免“未换绑成功但角色信息已变新” account.bind_status = "game_queried" account.updated_at = datetime.now(timezone.utc) 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 _xpd_role_context(self, client: DouyuActivityClient, config: dict) -> dict: """获取小店 H5 参数 + 绑定角色信息,小店任务共用。""" embed = client.xpd_embed_query( act_alias=str(config["xpd_act_alias"]), rid=str(config["xpd_rid"]), ) role = client.xpd_get_role( embed_query=embed["query"], act_id=str(config["xpd_act_id"]), rid=str(config["xpd_rid"]), ) return {"embed": embed, "role": role} def _xpd_area_id(self, role: dict, account: Account) -> int: """角色大区: 微信=1, 手Q=2, 未知回退账号已存值或 1。""" role_type = str(role.get("type") or "") if role_type == "wx": return 1 if role_type == "qq": return 2 return account.xpd_area_id or 1 def _apply_xpd_role_to_account(self, account: Account, role: dict, area_id: int) -> None: account.xpd_game_name = str(role.get("role_name") or "") or account.xpd_game_name account.xpd_openid = str(role.get("game_open_id") or "") or account.xpd_openid account.xpd_role_id = str(role.get("role_id") or "") or account.xpd_role_id account.xpd_plat_id = self._to_int(role.get("plat_id")) account.xpd_area_id = area_id account.updated_at = datetime.now(timezone.utc) def _execute_query_xpd_role( self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict, ): """查询和平小店绑定角色。""" client = self._client(cookie) ctx = self._xpd_role_context(client, config) role = ctx["role"] if not role.get("role_id"): self._mark_task(db, task, "failed", "未获取到小店绑定角色") return area_id = self._xpd_area_id(role, account) self._apply_xpd_role_to_account(account, role, area_id) account.xpd_bind_status = "xpd_bound" db.commit() role_text = str(role.get("role_name") or "-") channel = "微信" if role.get("type") == "wx" else ("手Q" if role.get("type") == "qq" else str(role.get("type") or "-")) self._mark_task( db, task, "success", f"小店角色: {role_text}({channel})", {"role": role, "area_id": area_id}, ) def _execute_get_xpd_bind_qr( self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict, ): """生成和平小店绑定二维码并轮询等待微信扫码绑定/换绑完成。 识别到新角色后仅标记"待确认",不自动回写账号,由用户手动确认绑定。 """ client = self._client(cookie) act_alias = str(config.get("xpd_act_alias") or "").strip() if not act_alias: self._mark_task(db, task, "failed", "请先配置小店活动代号 actAlias") return result = client.xpd_bind_qr(act_alias=act_alias) account.xpd_bind_status = "xpd_bind_qr_ready" account.updated_at = datetime.now(timezone.utc) db.commit() # 记录绑定前状态:已绑定账号生成二维码后必须等扫码换绑,不能立即成功 try: before = client.xpd_bind_info(act_alias=act_alias) result["before_bound"] = bool(before.get("bind_role")) result["before_role_name"] = str(before.get("role_name") or "") result["before_area_name"] = str(before.get("area_name") or "") result["before_plat_name"] = str(before.get("plat_name") or "") except Exception: result["before_bound"] = False result["before_role_name"] = "" result["bind_polling"] = True self._update_task_progress( db, task, "running", "二维码已生成,请微信扫码在小程序中绑定角色", result, ) state = self._wait_xpd_bind(db, task, client, act_alias, result) result["bind_polling"] = False if state == "stopped": self._mark_task(db, task, "stopped", "任务已停止", result) return if state == "pending": role_text = str(result.get("role_name") or "-") self._mark_task(db, task, "success", f"已识别角色: {role_text},待确认绑定", result) return self._mark_task( db, task, "failed", "未检测到小店绑定(二维码仍有效,可再次生成后扫码)", result, ) def _wait_xpd_bind( self, db: Session, task: DouyuTask, client: DouyuActivityClient, act_alias: str, result: dict, ) -> str: """轮询 bindInfo 检测绑定/换绑角色,识别到后停在"待确认",不自动回写账号。 - 绑定前未绑定:检测到 bind_role=1 即识别到待确认角色 - 绑定前已绑定(换绑):检测到角色名变化才算换绑完成,角色不变继续等 返回 "pending"=已识别待确认角色, "stopped"=任务停止, "timeout"=超时未识别。 """ before_bound = bool(result.get("before_bound")) before_role_name = str(result.get("before_role_name") or "") deadline = time.monotonic() + DOUYU_XPD_BIND_POLL_SECONDS poll_count = 0 while not self._stop.is_set() and time.monotonic() <= deadline: try: info = client.xpd_bind_info(act_alias=act_alias) poll_count += 1 result["bind_poll_count"] = poll_count result["bind_polling"] = True role_name = str(info.get("role_name") or "") bound_now = bool(info.get("bind_role")) changed = before_bound and bool(role_name) and role_name != before_role_name if (not before_bound and bound_now and role_name) or changed: result.update({key: value for key, value in info.items() if key != "raw"}) result["bind_polling"] = False result["xpd_pending_confirm"] = True return "pending" self._update_task_progress( db, task, "running", f"等待扫码绑定(第 {poll_count} 次)", result, ) except Exception as exc: poll_count += 1 result["bind_poll_count"] = poll_count result["bind_poll_error"] = str(exc) self._update_task_progress( db, task, "running", f"等待扫码绑定: {exc}", result, ) if self._stop.wait(DOUYU_XPD_BIND_POLL_INTERVAL): break result["bind_polling"] = False return "stopped" if self._stop.is_set() else "timeout" def _execute_confirm_xpd_bind( self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict, ): """确认和平小店绑定:回查 bindInfo,确认绑定角色后将账号回写为已绑定。""" client = self._client(cookie) act_alias = str(config.get("xpd_act_alias") or "").strip() if not act_alias: self._mark_task(db, task, "failed", "请先配置小店活动代号 actAlias") return result = client.xpd_bind_info(act_alias=act_alias) role_name = str(result.get("role_name") or "") if not result.get("bind_role") or not role_name: account.xpd_bind_status = "xpd_not_bound" account.updated_at = datetime.now(timezone.utc) db.commit() self._mark_task(db, task, "failed", "尚未检测到小店绑定角色,请先扫码绑定", result) return # 优先用完整角色信息回写(与查询角色一致),失败时回退 bindInfo 角色名 try: ctx = self._xpd_role_context(client, config) role = ctx["role"] if role.get("role_id"): self._apply_xpd_role_to_account(account, role, self._xpd_area_id(role, account)) else: account.xpd_game_name = role_name account.updated_at = datetime.now(timezone.utc) except Exception: account.xpd_game_name = role_name account.updated_at = datetime.now(timezone.utc) account.xpd_bind_status = "xpd_bound" account.updated_at = datetime.now(timezone.utc) db.commit() result["xpd_pending_confirm"] = False result["xpd_bound"] = True self._mark_task(db, task, "success", f"小店绑定成功: {role_name}", result) def _execute_query_xpd_bind_info( self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict, ): """查询和平小店绑定信息(bindInfo)。 仅查询展示,不回写账号;确认绑定由 confirm_xpd_bind 任务完成。 """ client = self._client(cookie) act_alias = str(config.get("xpd_act_alias") or "").strip() if not act_alias: self._mark_task(db, task, "failed", "请先配置小店活动代号 actAlias") return result = client.xpd_bind_info(act_alias=act_alias) status = "已绑定" if result.get("bind_role") else "未绑定" text = str(result.get("role_name") or result.get("nick") or "-") self._mark_task( db, task, "success", f"小店绑定: {status} ({text})", result, ) def _execute_refresh_xpd_goods( self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict, ): """刷新和平小店商品列表快照(全局数据,任一可用 CK 即可)。""" client = self._client(cookie) ctx = self._xpd_role_context(client, config) role = ctx["role"] if not role.get("role_id"): self._mark_task(db, task, "failed", "未获取到小店绑定角色") return area_id = self._xpd_area_id(role, account) result = client.xpd_list_goods( embed_query=ctx["embed"]["query"], act_id=str(config["xpd_act_id"]), openid=str(role.get("game_open_id") or ""), roleid=str(role.get("role_id") or ""), areaid=str(area_id), ) goods = result["goods"] self._upsert_xpd_goods(db, goods) self._apply_xpd_role_to_account(account, role, area_id) account.xpd_bind_status = "xpd_goods_refreshed" db.commit() self._mark_task( db, task, "success", f"已刷新小店商品 {len(goods)} 个", {"goods_count": len(goods), "goods": goods}, ) def _execute_query_xpd_balance( self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict, ): """查询和平小店点券余额。""" client = self._client(cookie) ctx = self._xpd_role_context(client, config) role = ctx["role"] if not role.get("role_id"): self._mark_task(db, task, "failed", "未获取到小店绑定角色") return area_id = self._xpd_area_id(role, account) result = client.xpd_balance( embed_query=ctx["embed"]["query"], act_id=str(config["xpd_act_id"]), openid=str(role.get("game_open_id") or ""), roleid=str(role.get("role_id") or ""), plat=str(role.get("plat_id") or "1"), areaid=str(area_id), ) balance = result.get("balance") self._apply_xpd_role_to_account(account, role, area_id) account.xpd_balance = balance account.xpd_bind_status = "xpd_balance_queried" db.commit() if balance is None: self._mark_task(db, task, "failed", "未获取到小店点券余额") return self._mark_task( db, task, "success", f"小店点券余额: {balance}", {"balance": balance, "role": role, "area_id": area_id}, ) def _execute_query_xpd_fragments( self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict, ): """查询和平小店扭蛋碎片数量。 优先现查角色;getrole 受限(Livelink 风控/失效)时回退账号已存角色, 保证已绑定账号仍可查询。 """ client = self._client(cookie) act_id = str(config["xpd_act_id"]) embed_query: dict = {} openid = str(account.xpd_openid or "") roleid = str(account.xpd_role_id or "") plat = str(account.xpd_plat_id or "1") areaid = str(account.xpd_area_id or 1) role: dict = {} try: ctx = self._xpd_role_context(client, config) embed_query = ctx["embed"]["query"] role = ctx["role"] if isinstance(ctx.get("role"), dict) else {} if role.get("role_id"): role_area = self._xpd_area_id(role, account) openid = str(role.get("game_open_id") or "") or openid roleid = str(role.get("role_id") or "") or roleid plat = str(role.get("plat_id") or "1") or plat areaid = str(role_area) or areaid self._apply_xpd_role_to_account(account, role, role_area) except Exception: pass if not openid or not roleid: self._mark_task(db, task, "failed", "未获取到小店绑定角色,请先生成二维码扫码绑定") return result = client.xpd_fragments( embed_query=embed_query, act_id=act_id, openid=openid, roleid=roleid, plat=plat, areaid=areaid, ) fragments = result.get("fragments") account.xpd_fragments = fragments account.xpd_bind_status = "xpd_fragments_queried" account.updated_at = datetime.now(timezone.utc) db.commit() if fragments is None: self._mark_task(db, task, "failed", "未获取到小店扭蛋碎片数量") return self._mark_task( db, task, "success", f"小店扭蛋碎片: {fragments}", {"fragments": fragments, "role": role, "area_id": int(areaid)}, ) def _execute_exchange_xpd_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: pay_type = int(payload.get("pay_type") or 1) except (TypeError, ValueError): self._mark_task(db, task, "failed", "兑换货币参数无效") return if pay_type not in (1, 5): self._mark_task(db, task, "failed", "小店兑换仅支持点券或扭蛋碎片") return goods = ( db.query(DouyuXpdGoodsSnapshot) .filter(DouyuXpdGoodsSnapshot.commodity_id == commodity_id) .first() ) if not goods: self._mark_task(db, task, "failed", "未找到小店商品快照,请先刷新商品列表") return goods_snapshot = goods.raw if isinstance(goods.raw, dict) else {} goods_raw = goods_snapshot.get("raw") if isinstance(goods_snapshot.get("raw"), dict) else goods_snapshot price_key = "iPrice" if pay_type == 1 else "iJb2Price" price = self._to_int(goods_raw.get(price_key)) if price is None: price = goods.price if pay_type == 1 else None if price is None or price <= 0: currency = "点券" if pay_type == 1 else "扭蛋碎片" self._mark_task(db, task, "failed", f"该商品不支持使用{currency}兑换") return # iGoodsLeft=-1 表示活动未公开库存,不是售罄;只有 0 才阻止兑换。 if goods.goods_left == 0: self._mark_task(db, task, "failed", "该商品库存不足,请刷新商品列表后重试") return client = self._client(cookie) embed = client.xpd_embed_query( act_alias=str(config["xpd_act_alias"]), rid=str(config["xpd_rid"]), ) role: dict = {} try: role = client.xpd_get_role( embed_query=embed["query"], act_id=str(config["xpd_act_id"]), rid=str(config["xpd_rid"]), ) if role.get("role_id"): self._apply_xpd_role_to_account(account, role, self._xpd_area_id(role, account)) except DouyuActivityError as exc: self._push_log("warning", f"小店兑换前刷新角色失败,使用已保存角色: {exc}") if not role.get("role_id") and not account.xpd_role_id: self._mark_task(db, task, "failed", "未获取到小店绑定角色,请先生成二维码扫码绑定") return result = client.xpd_exchange_goods( embed_query=embed["query"], act_id=str(config["xpd_act_id"]), rid=str(config["xpd_rid"]), commodity_id=commodity_id, price=price, picture=str(goods_raw.get("sGoodsPic") or ""), pay_type=pay_type, action_id=str(goods_raw.get("iActionId") or ""), ) if pay_type == 1 and result.get("new_balance") is not None: account.xpd_balance = result["new_balance"] if pay_type == 5 and result.get("new_balance") is not None: account.xpd_fragments = result["new_balance"] account.xpd_bind_status = "xpd_goods_exchanged" account.updated_at = datetime.now(timezone.utc) db.commit() currency = "点券" if pay_type == 1 else "扭蛋碎片" display_role = role.get("role_name") or account.xpd_game_name or "" channel = "微信" if (role.get("type") == "wx" or account.xpd_area_id == 1) else "手Q" result.update({ "goods": goods_raw, "game_name": display_role, "game_channel": channel, "account_name": account.nickname or account.username or account.uid or f"#{account.id}", "currency": currency, }) self._mark_task(db, task, "success", f"兑换小店商品成功: {goods.name or commodity_id}({price}{currency})", result) 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 {} # 换绑冷却必须看活动当前绑定(QYOOB),不要用 cjm cooldown_info = self._pick_change_wait_bind_info(before_candidates, config) 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 or before_candidates[0] 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"] if cooldown_snapshot["is_bound_act"] else "" ) 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 confirm_alias: self._mark_task(db, task, "failed", "请先配置确认绑定活动 actAlias") return 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 "") # 确认前“已生效绑定”角色(bound_act=1),确认失败/回查失败时写库用,避免待确认角色污染 game_name before_bound = self._pick_current_bound_info( before_candidates, config, extra_prefer_aliases=[confirm_alias], ) 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 before_is_current_bound = ( before_bound is not None and str(before.get("act_alias") or "") == str(before_bound.get("act_alias") or "") and role_name == str(before_bound.get("role_name") or "") ) if before_is_current_bound: self._apply_bind_info_to_account(account, before_bound, "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_bound, **self._bind_snapshot(before_bound), "bind_ready_for_confirm": True, "bind_confirmed": True, "bind_phase": "confirmed", }, ) return # 确认接口优先用配置的确认 alias;没有则回退到命中查询的 alias use_confirm_alias = confirm_alias or query_alias try: confirm_result = client.confirm_bind(use_confirm_alias) except DouyuActivityError as exc: confirm_msg = str(exc) self._push_log("warning", f"确认绑定接口失败: {confirm_msg}") # 待绑定游戏账号侧换绑限制(未到换绑时间等)导致确认失败:保留原绑定并给出明确提示 if before_bound is not None: self._apply_bind_info_to_account(account, before_bound, "game_queried") else: account.bind_status = "game_queried" account.updated_at = datetime.now(timezone.utc) self._mark_task( db, task, "failed", f"待绑定游戏账号({role_name or '-'})未到换绑时间(不是斗鱼/虎牙账号),请重新换账号扫码绑定", { "act_alias": use_confirm_alias, "query_act_alias": query_alias, "query_act_aliases": query_aliases, "before_bind_info": before, **before_snapshot, "confirm_error": confirm_msg, "bind_ready_for_confirm": True, "bind_confirmed": False, "bind_phase": "confirm_failed", }, ) return confirm_raw = confirm_result.get("raw") or {} self._push_log( "info", f"确认绑定接口返回 act={use_confirm_alias} " f"error={confirm_raw.get('error')} msg={confirm_raw.get('msg') or '-'}", ) def _pick_bound_after(candidates: list[dict]) -> dict | None: """确认后回查:只认活动 alias 上的已生效绑定(bound_act=1)。""" return self._pick_current_bound_info( candidates, config, extra_prefer_aliases=[use_confirm_alias], ) try: after_candidates = self._fetch_bind_info_candidates(client, query_aliases) if not after_candidates: raise DouyuActivityError("确认后回查绑定信息失败") after = _pick_bound_after(after_candidates) # 已生效绑定存在同步延迟:确认接口已成功但未生效时短轮询等待 if after is None: for _ in range(DOUYU_CONFIRM_EFFECT_POLL_TIMES): if self._stop.wait(DOUYU_CONFIRM_EFFECT_POLL_INTERVAL): break after_candidates = self._fetch_bind_info_candidates(client, query_aliases) if not after_candidates: break after = _pick_bound_after(after_candidates) if after is not None: break if after is None: # 已生效绑定始终未出现:确认未生效(或同步延迟超时),保留原绑定 if before_bound is not None: self._apply_bind_info_to_account(account, before_bound, "game_queried") else: account.bind_status = "game_queried" account.updated_at = datetime.now(timezone.utc) self._mark_task( db, task, "failed", f"待绑定游戏账号({role_name or '-'})未到换绑时间(不是斗鱼/虎牙账号),请重新换账号扫码绑定", { "act_alias": use_confirm_alias, "query_act_alias": query_alias, "query_act_aliases": query_aliases, "before_bind_info": before, "confirm": confirm_result, "after_bind_info": None, **before_snapshot, "bind_ready_for_confirm": True, "bind_confirmed": False, "bind_phase": "confirm_failed", "confirm_wait_error": "确认后短轮询未等到已生效绑定", }, ) return self._push_log( "info", f"确认后回查 hit={after.get('act_alias') or '-'} " f"{self._format_bind_summary(after)}", ) except DouyuActivityError as exc: # 查询接口异常:确认接口已成功时按成功处理,但保留错误信息。 # 写库优先确认前已生效绑定,避免待确认角色被误写入。 if before_bound is not None: self._apply_bind_info_to_account(account, before_bound, "bind_confirmed") else: account.bind_status = "bind_confirmed" account.updated_at = datetime.now(timezone.utc) 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_role_name = after_snapshot["role_name"] or role_name self._apply_bind_info_to_account(account, after, "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, **after_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) # 兑换成功后自动刷新积分,更新账号最新积分信息(失败不阻断兑换成功) points_refresh = None try: points_refresh = self._refresh_account_points(client, account, cookie, ctn=ctn) except Exception as exc: self._push_log("warning", f"兑换后刷新积分失败: {exc}") 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, **(points_refresh or {}), }, ) 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 # 1) 优先“已生效绑定”角色(bound_act=1 且有角色名,通常是活动 alias)。 # 避免把扫码后未确认的新角色当成当前绑定结果。 bound_info = self._pick_current_bound_info(candidates, config) # 2) 待确认角色(cjm 扫码后未确认;无已绑定时也用于首次绑定展示) pending_info = self._pick_bind_info( candidates, prefer_pending=True, prefer_aliases=query_aliases, ) or candidates[0] pending_role = str(pending_info.get("role_name") or "").strip() bound_role = str(bound_info.get("role_name") or "") if bound_info else "" has_pending = bool(pending_role) and (pending_role != bound_role or not bound_info) # 查询结果角色 = 已生效绑定;首次绑定(无绑定)时回退待确认角色 bind_info = bound_info if bound_info else (pending_info if has_pending else None) snapshot = self._bind_snapshot(bind_info) if bind_info else {} role_name = snapshot.get("role_name") or "" is_bound = snapshot.get("is_bound_act", False) source_alias = str((bind_info or {}).get("act_alias") or "") summary = ( f"act={source_alias or '-'} {self._format_bind_summary(bind_info)}" if bind_info else "act=- role=-" ) self._push_log( "info", "查询角色 " f"aliases={','.join(query_aliases)} bound_hit={(bound_info or {}).get('act_alias') or '-'} " f"pending_hit={pending_info.get('act_alias') or '-'} " f"{self._format_bind_summary(bind_info) if bind_info else 'role=- bound_act=-'}", ) # 只有已生效绑定才写账号表;待确认角色不污染 game_name if bind_info and is_bound: self._apply_bind_info_to_account(account, bind_info, "game_queried") else: account.bind_status = "game_queried" if has_pending else "game_not_bound" account.updated_at = datetime.now(timezone.utc) result = { "act_alias": source_alias, "query_act_alias": source_alias, "query_act_aliases": query_aliases, **snapshot, "bind_ready_for_confirm": has_pending, "bind_confirmed": is_bound, "bind_phase": ( "confirmed" if is_bound 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 ], } # 已绑定角色之外另有待确认新角色时,单独字段展示,不覆盖 role_name if has_pending and bound_info and pending_role != bound_role: result["pending_role_name"] = pending_role result["pending_area_name"] = str(pending_info.get("area_name") or "") result["pending_plat_name"] = str(pending_info.get("plat_name") or "") if not bind_info: message = f"未获取到游戏名 | {summary}" elif is_bound: if has_pending and pending_role != bound_role: message = f"当前已绑定: {bound_role},待确认: {pending_role} | {summary}" else: message = f"当前已绑定: {bound_role} | {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 # 冷却是“当前已生效绑定”的属性,不能用 legacy/cjm 的待确认角色判断。 current_bound = self._pick_current_bound_info(candidates, config) if current_bound is not None: return current_bound action_aliases = self._action_act_aliases(config) return self._pick_bind_info( [info for info in candidates if str(info.get("act_alias") or "") in action_aliases], prefer_pending=False, prefer_aliases=action_aliases, ) 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 for item in candidates: self._push_log( "info", f"查询换绑时间候选 act={item.get('act_alias') or '-'} " f"{self._format_bind_summary(item)}", ) # 换绑时间看“当前已绑定”活动态,不要优先 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 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)}") 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_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, "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: 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: 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()