"""虎牙任务批次执行器。""" import asyncio import threading 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.huya import HuyaHttpClient from ..database import SessionLocal from ..models import HuyaAccount, HuyaGoodsSnapshot, HuyaRechargeGoodsSnapshot, HuyaTask from .huya_service import HUYA_CONFIG_FIELDS, cookie_value, ensure_huya_config, huya_config_value HUYA_RECHARGE_ACT_ID = 25135 HUYA_RECHARGE_SOURCE_ID = "yellowcarlist" HUYA_RECHARGE_SCENE = 4 HUYA_RECHARGE_EXTRA_PRODUCTS = [ { "spu_id": "hy-5879340", "name": "精英宝典", "task_name": "开通精英宝典", "description": "得300积分丨解锁道具兑换权益", "sort": 0, }, ] class HuyaBatchRunner: """批量执行虎牙任务,通过队列推送实时日志。""" 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" and message: log_func = getattr(logger, level, logger.info) log_func(f"[huya] {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_info: dict) -> str: return ( account_info.get("nickname") or account_info.get("username") or account_info.get("uid") or f"#{account_info.get('account_id')}" ) @staticmethod def _to_int(value) -> int: text = str(value or "").strip() return int(text) if text.isdigit() else 0 def _resolve_uid(self, account_info: dict) -> int: cookie = account_info.get("cookie") or "" return ( self._to_int(account_info.get("yyuid")) or self._to_int(account_info.get("uid")) or self._to_int(cookie_value(cookie, "yyuid")) or self._to_int(cookie_value(cookie, "udb_uid")) ) @staticmethod def _role_name(bind_status) -> str: account_data = bind_status.accountData return account_data.gameRole.roleName or account_data.gameAccount.nick or "" @staticmethod def _bind_role_result(bind_status) -> dict: account_data = bind_status.accountData game_account = account_data.gameAccount game_role = account_data.gameRole return { "game_title": bind_status.gameName, "role_name": HuyaBatchRunner._role_name(bind_status), "change_bind_day": bind_status.changeBindDay, "is_bind_account": account_data.isBindAcount, "is_bind_role": account_data.isBindRole, "is_need_act_check": account_data.isNeedActCheck, "change_bind_time": account_data.changBindTime, "game_account": game_account.to_dict(), "game_role": game_role.to_dict(), } @staticmethod def _role_channel(bind_status) -> str: game_role = bind_status.accountData.gameRole parts = [bind_status.gameName, game_role.areaName, game_role.platName] return " / ".join(part for part in parts if part) @staticmethod def _format_local_time(timestamp: int) -> str: if not timestamp: return "" return datetime.fromtimestamp(timestamp).strftime("%Y-%m-%d %H:%M:%S") @classmethod def _bind_change_state(cls, bind_status) -> dict: account_data = bind_status.accountData is_bound = bool(account_data.isBindAcount and account_data.isBindRole) change_time = int(account_data.changBindTime or 0) now = int(datetime.now(timezone.utc).timestamp()) can_change = not is_bound or not change_time or change_time <= now return { "is_bound": is_bound, "can_change_bind": can_change, "change_bind_time": change_time, "change_available_at": cls._format_local_time(change_time), "change_bind_day": int(bind_status.changeBindDay or 0), } def _apply_role_to_account(self, account: HuyaAccount, bind_status, status: str): role_name = self._role_name(bind_status) account.status = status account.game_name = role_name or bind_status.gameName or account.game_name account.game_channel = self._role_channel(bind_status) or account.game_channel account.updated_at = datetime.now(timezone.utc) def _mark_task( self, worker_db: Session, task: HuyaTask, status: str, message: str, result: Optional[dict] = None, ): task.status = status task.message = message task.result = result task.finished_at = datetime.now(timezone.utc) worker_db.commit() def _execute_query_points( self, worker_db: Session, task: HuyaTask, account: HuyaAccount, account_info: dict, config_info: dict, ): sid = str(self.payload.get("sid") or config_info.get("sid") or "").strip() if not sid: self._mark_task(worker_db, task, "failed", "请先配置虎牙活动 SID") return sid_int = self._to_int(sid) if not sid_int: self._mark_task(worker_db, task, "failed", f"虎牙活动 SID 无效: {sid}") return uid = self._resolve_uid(account_info) if not uid: self._mark_task(worker_db, task, "failed", "无法从账号或 Cookie 解析 yyuid") return cookie = account_info.get("cookie") or "" if not cookie: self._mark_task(worker_db, task, "failed", "账号 Cookie 为空") return client = HuyaHttpClient(logger=lambda msg: self._push_log("info", f"[{uid}] {msg}")) response = client.query_user_score(uid=uid, cookie=cookie, sid=sid_int) if response is None: self._mark_task(worker_db, task, "error", "虎牙积分接口无响应") return result = response.to_dict() result["sid"] = sid_int if response.status != 200: self._mark_task( worker_db, task, "failed", response.msg or f"虎牙积分查询失败: {response.status}", result, ) return points = response.available_score account.points = points account.status = "points_queried" account.updated_at = datetime.now(timezone.utc) self._mark_task(worker_db, task, "success", f"积分: {points}", result) def _execute_refresh_goods( self, worker_db: Session, task: HuyaTask, account: HuyaAccount, account_info: dict, config_info: dict, ): sid = str(self.payload.get("sid") or config_info.get("sid") or "").strip() if not sid: self._mark_task(worker_db, task, "failed", "请先配置虎牙活动 SID") return sid_int = self._to_int(sid) if not sid_int: self._mark_task(worker_db, task, "failed", f"虎牙活动 SID 无效: {sid}") return uid = self._resolve_uid(account_info) if not uid: self._mark_task(worker_db, task, "failed", "无法从账号或 Cookie 解析 yyuid") return cookie = account_info.get("cookie") or "" if not cookie: self._mark_task(worker_db, task, "failed", "账号 Cookie 为空") return client = HuyaHttpClient(logger=lambda msg: self._push_log("info", f"[{uid}] {msg}")) response = client.get_act_prize_list(uid=uid, cookie=cookie, sid=sid_int) if response is None: self._mark_task(worker_db, task, "error", "虎牙商品列表接口无响应") return result = response.to_dict() result["sid"] = sid_int if response.status != 200: self._mark_task( worker_db, task, "failed", response.msg or f"虎牙商品列表刷新失败: {response.status}", result, ) return goods = [ item for item in result.get("goods", []) if item.get("product_id") and item.get("name") ] now = datetime.now(timezone.utc) worker_db.query(HuyaGoodsSnapshot).delete(synchronize_session=False) for item in goods: worker_db.add(HuyaGoodsSnapshot( product_id=item["product_id"], name=item["name"], price=item["price"], remain_text=item["remain_text"], raw=item, updated_at=now, )) account.status = "goods_refreshed" account.updated_at = now message = f"已刷新商品 {len(goods)} 个" self._mark_task(worker_db, task, "success", message, {**result, "goods": goods}) @staticmethod def _normalize_pay_channel(value) -> str: text = str(value or "").strip() lowered = text.lower() if lowered in {"weixin", "wx", "wechat", "微信"}: return "Weixin" return "Zfb" @staticmethod def _pay_channel_label(value: str) -> str: return "微信" if value == "Weixin" else "支付宝" @staticmethod def _recharge_price_text(price: int | None) -> str: if not price: return "" return f"{price / 100:.2f}元" def _execute_refresh_recharge_goods( self, worker_db: Session, task: HuyaTask, account: HuyaAccount, account_info: dict, config_info: dict, ): pid = self._to_int(config_info.get("room_pid")) if not pid: self._mark_task(worker_db, task, "failed", "请先配置虎牙直播间 ID") return uid = self._resolve_uid(account_info) if not uid: self._mark_task(worker_db, task, "failed", "无法从账号或 Cookie 解析 yyuid") return cookie = account_info.get("cookie") or "" if not cookie: self._mark_task(worker_db, task, "failed", "账号 Cookie 为空") return client = HuyaHttpClient(logger=lambda msg: self._push_log("info", f"[{uid}] {msg}")) task_resp = client.get_act_task_detail(uid=uid, cookie=cookie, act_id=HUYA_RECHARGE_ACT_ID) if task_resp is None: self._mark_task(worker_db, task, "error", "虎牙充值任务详情接口无响应") return task_result = task_resp.to_dict() if task_resp.status != 200: self._mark_task( worker_db, task, "failed", task_resp.msg or f"虎牙充值任务详情获取失败: {task_resp.status}", task_result, ) return candidates: list[dict] = [] seen: set[str] = set() def add_candidate(item: dict): spu_id = str(item.get("spu_id") or "").strip() if not spu_id or spu_id in seen: return seen.add(spu_id) candidates.append(item) for item in HUYA_RECHARGE_EXTRA_PRODUCTS: add_candidate(dict(item)) for index, item in enumerate(task_result.get("tasks", []), start=1): if int(item.get("task_type") or 0) != 67: continue add_candidate({ "spu_id": item.get("spu_id") or "", "name": item.get("name") or "", "task_id": str(item.get("task_id") or ""), "task_name": item.get("name") or "", "description": item.get("description") or "", "icon": item.get("icon") or "", "task_url": item.get("task_url") or "", "prizes": item.get("prizes") or [], "sort": index, }) if not candidates: self._mark_task(worker_db, task, "failed", "未从活动任务中发现充值商品", task_result) return now = datetime.now(timezone.utc) goods: list[dict] = [] failed: list[dict] = [] for candidate in candidates: spu_id = candidate["spu_id"] detail_resp = client.get_goods_info( uid=uid, guid="", cookie=cookie, pid=pid, spu_id=spu_id, sku_id=0, game_id="0", source_id=HUYA_RECHARGE_SOURCE_ID, scene=HUYA_RECHARGE_SCENE, ) if detail_resp is None: failed.append({"spu_id": spu_id, "message": "商品详情接口无响应"}) continue detail = detail_resp.to_dict() if detail_resp.code != 200 or not detail.get("sku_id"): failed.append({ "spu_id": spu_id, "message": detail_resp.message or f"商品详情获取失败: {detail_resp.code}", "detail": detail, }) continue item = { **candidate, **detail, "spu_id": detail.get("spu_id") or spu_id, "sku_id": str(detail.get("sku_id") or ""), "name": detail.get("name") or candidate.get("name") or spu_id, "description": detail.get("description") or candidate.get("description") or "", "icon": detail.get("icon") or candidate.get("icon") or "", "task_id": candidate.get("task_id") or "", "task_name": candidate.get("task_name") or candidate.get("name") or "", "raw_order": int(candidate.get("sort") or 0), } goods.append(item) worker_db.query(HuyaRechargeGoodsSnapshot).delete(synchronize_session=False) for item in goods: worker_db.add(HuyaRechargeGoodsSnapshot( spu_id=item["spu_id"], sku_id=item["sku_id"], name=item["name"], price=item.get("price") or None, stock=item.get("stock") or None, buy_limit=item.get("buy_limit") or None, icon=item.get("icon") or "", description=item.get("description") or "", task_id=item.get("task_id") or "", task_name=item.get("task_name") or "", raw=item, updated_at=now, )) account.status = "recharge_goods_refreshed" account.updated_at = now message = f"已刷新充值商品 {len(goods)} 个" if failed: message += f",失败 {len(failed)} 个" result = { "act_id": HUYA_RECHARGE_ACT_ID, "goods_count": len(goods), "failed_count": len(failed), "goods": goods, "failed": failed, "task_detail": task_result, } self._mark_task(worker_db, task, "success" if goods else "failed", message, result) def _execute_create_recharge_order( self, worker_db: Session, task: HuyaTask, account: HuyaAccount, account_info: dict, config_info: dict, ): pid = self._to_int(config_info.get("room_pid")) if not pid: self._mark_task(worker_db, task, "failed", "请先配置虎牙直播间 ID") return spu_id = str(self.payload.get("spu_id") or "").strip() if not spu_id: self._mark_task(worker_db, task, "failed", "请选择充值商品") return count = self._to_int(self.payload.get("count")) or 1 count = max(1, min(count, 999)) pay_channel = self._normalize_pay_channel(self.payload.get("pay_channel") or config_info.get("pay_channel")) uid = self._resolve_uid(account_info) if not uid: self._mark_task(worker_db, task, "failed", "无法从账号或 Cookie 解析 yyuid") return cookie = account_info.get("cookie") or "" if not cookie: self._mark_task(worker_db, task, "failed", "账号 Cookie 为空") return snapshot = worker_db.query(HuyaRechargeGoodsSnapshot).filter( HuyaRechargeGoodsSnapshot.spu_id == spu_id ).first() payload_sku_id = self._to_int(self.payload.get("sku_id")) sku_id = payload_sku_id or self._to_int(snapshot.sku_id if snapshot else "") product_name = str(self.payload.get("product_name") or (snapshot.name if snapshot else "") or spu_id) unit_price = int(snapshot.price or 0) if snapshot else 0 client = HuyaHttpClient(logger=lambda msg: self._push_log("info", f"[{uid}] {msg}")) detail_resp = client.get_goods_info( uid=uid, guid="", cookie=cookie, pid=pid, spu_id=spu_id, sku_id=sku_id or 0, game_id="0", source_id=HUYA_RECHARGE_SOURCE_ID, scene=HUYA_RECHARGE_SCENE, ) if detail_resp is None: self._mark_task(worker_db, task, "error", "虎牙充值商品详情接口无响应") return detail = detail_resp.to_dict() if detail_resp.code != 200: self._mark_task( worker_db, task, "failed", detail_resp.message or f"虎牙充值商品详情获取失败: {detail_resp.code}", detail, ) return sku_id = int(detail.get("sku_id") or sku_id or 0) product_name = detail.get("name") or product_name unit_price = int(detail.get("price") or unit_price or 0) if not sku_id: self._mark_task(worker_db, task, "failed", "充值商品缺少 SKU,请先刷新充值商品列表", detail) return order_resp = client.create_order( uid=uid, guid="", cookie=cookie, pid=pid, spu_id=spu_id, sku_id=sku_id, item_count=count, source_id=HUYA_RECHARGE_SOURCE_ID, game_id="0", scene=HUYA_RECHARGE_SCENE, order_type=6, ) if order_resp is None: self._mark_task(worker_db, task, "error", "虎牙下单接口无响应") return order_result = order_resp.to_dict() if order_resp.code != 200 or not order_resp.orderId: self._mark_task( worker_db, task, "failed", order_resp.message or f"虎牙下单失败: {order_resp.code}", {"goods": detail, "order": order_result}, ) return pay_resp = client.pay_order_submit( uid=uid, guid="", cookie=cookie, order_id=order_resp.orderId, pay_type=pay_channel, pid=pid, source_id=HUYA_RECHARGE_SOURCE_ID, scene=HUYA_RECHARGE_SCENE, item_count=count, ) if pay_resp is None: self._mark_task(worker_db, task, "error", "虎牙支付接口无响应", {"goods": detail, "order": order_result}) return pay_result = pay_resp.to_dict() if pay_resp.code != 200 or not pay_resp.payUrl: self._mark_task( worker_db, task, "failed", pay_resp.message or f"虎牙支付二维码生成失败: {pay_resp.code}", {"goods": detail, "order": order_result, "pay": pay_result}, ) return amount = int(pay_resp.amount or unit_price * count or 0) result = { "spu_id": spu_id, "sku_id": sku_id, "product_name": product_name, "count": count, "unit_price": unit_price, "amount": amount, "amount_text": self._recharge_price_text(amount), "pay_channel": pay_channel, "pay_channel_label": self._pay_channel_label(pay_channel), "order_id": order_resp.orderId, "app_order_id": pay_resp.appOrderId, "pay_order_id": pay_resp.payOrderId, "pay_url": pay_resp.payUrl, "goods": detail, "order": order_result, } account.status = "recharge_order_created" account.updated_at = datetime.now(timezone.utc) message = f"{product_name} x{count} {self._pay_channel_label(pay_channel)} {result['amount_text']}" self._mark_task(worker_db, task, "success", message, result) def _execute_get_bind_qr( self, worker_db: Session, task: HuyaTask, account: HuyaAccount, account_info: dict, config_info: dict, ): b_act_id = str(self.payload.get("bind_act_id") or config_info.get("bind_act_id") or "").strip() if not b_act_id: self._mark_task(worker_db, task, "failed", "请先配置虎牙绑定 bActId") return b_act_id_int = self._to_int(b_act_id) if not b_act_id_int: self._mark_task(worker_db, task, "failed", f"虎牙绑定 bActId 无效: {b_act_id}") return uid = self._resolve_uid(account_info) if not uid: self._mark_task(worker_db, task, "failed", "无法从账号或 Cookie 解析 yyuid") return cookie = account_info.get("cookie") or "" if not cookie: self._mark_task(worker_db, task, "failed", "账号 Cookie 为空") return client = HuyaHttpClient(logger=lambda msg: self._push_log("info", f"[{uid}] {msg}")) bind_status = client.check_user_bind_game_account( uid=uid, cookie=cookie, b_act_id=b_act_id_int, is_use_outer_act_id=1, ) if bind_status is None: self._mark_task(worker_db, task, "error", "虎牙绑定状态接口无响应") return if bind_status.status != 200: self._mark_task( worker_db, task, "failed", bind_status.msg or f"虎牙绑定状态查询失败: {bind_status.status}", bind_status.to_dict(), ) return role_info = self._bind_role_result(bind_status) change_state = self._bind_change_state(bind_status) if role_info["role_name"]: self._apply_role_to_account(account, bind_status, account.status or "imported") if not change_state["can_change_bind"]: role_name = role_info["role_name"] or "当前角色" available_at = change_state["change_available_at"] result = { "bind_act_id": b_act_id_int, **role_info, **change_state, "bind_status": bind_status.to_dict(), } self._mark_task( worker_db, task, "failed", f"{role_name} 暂不能更换,{available_at} 后可更换", result, ) return live_link = client.get_live_link_param( uid=uid, cookie=cookie, b_act_id=b_act_id_int, game_auth_scene=bind_status.gameAuthScene, ) if live_link is None: self._mark_task(worker_db, task, "error", "虎牙绑定二维码参数接口无响应") return if live_link.status != 200: self._mark_task( worker_db, task, "failed", live_link.msg or f"虎牙绑定二维码参数获取失败: {live_link.status}", live_link.to_log_dict(), ) return profile_nick = account_info.get("nickname") or account_info.get("username") or "" profile_avatar = "" profile_resp = client.get_user_profile_batch(uid=uid, cookie=cookie, target_uids=[uid]) if profile_resp is not None and profile_resp.profiles: profile = profile_resp.profiles[0] profile_nick = profile.nick or profile.passport or profile_nick profile_avatar = profile.avatar or "" urls = client.build_bind_urls( live_link.livelinkParam, b_act_id_int, game_auth_scene=bind_status.gameAuthScene, nick_name=profile_nick, face_url=profile_avatar, ) mini_qrcode = client.get_livelink_mini_qrcode(urls["qr_url"]) if not mini_qrcode: result = { "bind_act_id": b_act_id_int, "profile": { "nick": profile_nick, "avatar": profile_avatar, }, } self._mark_task(worker_db, task, "failed", "绑定小程序码获取失败", result) return result = { "bind_act_id": b_act_id_int, "mini_qrcode_image": mini_qrcode["mini_qrcode_image"], "qrcode_token": mini_qrcode["qrcode_token"], "bind_status": bind_status.to_dict(), **role_info, **change_state, "profile": { "nick": profile_nick, "avatar": profile_avatar, }, } account.status = "bind_qr_generated" account.game_name = role_info["role_name"] or bind_status.gameName or account.game_name account.nickname = profile_nick or account.nickname account.updated_at = datetime.now(timezone.utc) self._mark_task(worker_db, task, "success", "已生成绑定小程序码", result) def _execute_query_game_name( self, worker_db: Session, task: HuyaTask, account: HuyaAccount, account_info: dict, config_info: dict, ): b_act_id = str(self.payload.get("bind_act_id") or config_info.get("bind_act_id") or "").strip() if not b_act_id: self._mark_task(worker_db, task, "failed", "请先配置虎牙绑定 bActId") return b_act_id_int = self._to_int(b_act_id) if not b_act_id_int: self._mark_task(worker_db, task, "failed", f"虎牙绑定 bActId 无效: {b_act_id}") return uid = self._resolve_uid(account_info) if not uid: self._mark_task(worker_db, task, "failed", "无法从账号或 Cookie 解析 yyuid") return cookie = account_info.get("cookie") or "" if not cookie: self._mark_task(worker_db, task, "failed", "账号 Cookie 为空") return client = HuyaHttpClient(logger=lambda msg: self._push_log("info", f"[{uid}] {msg}")) bind_status = client.check_user_bind_game_account( uid=uid, cookie=cookie, b_act_id=b_act_id_int, is_use_outer_act_id=1, ) if bind_status is None: self._mark_task(worker_db, task, "error", "虎牙角色信息接口无响应") return if bind_status.status != 200: self._mark_task( worker_db, task, "failed", bind_status.msg or f"虎牙角色信息查询失败: {bind_status.status}", {"bind_act_id": b_act_id_int, "bind_status": bind_status.to_dict()}, ) return role_info = self._bind_role_result(bind_status) change_state = self._bind_change_state(bind_status) result = { "bind_act_id": b_act_id_int, **role_info, **change_state, "bind_status": bind_status.to_dict(), } role_name = role_info["role_name"] if role_name: self._apply_role_to_account(account, bind_status, "game_queried") self._mark_task(worker_db, task, "success", f"角色: {role_name}", result) return account.status = "game_not_bound" account.updated_at = datetime.now(timezone.utc) self._mark_task(worker_db, task, "success", "未绑定游戏角色", result) def _execute_confirm_bind( self, worker_db: Session, task: HuyaTask, account: HuyaAccount, account_info: dict, config_info: dict, ): b_act_id = str(self.payload.get("bind_act_id") or config_info.get("bind_act_id") or "").strip() if not b_act_id: self._mark_task(worker_db, task, "failed", "请先配置虎牙绑定 bActId") return b_act_id_int = self._to_int(b_act_id) if not b_act_id_int: self._mark_task(worker_db, task, "failed", f"虎牙绑定 bActId 无效: {b_act_id}") return uid = self._resolve_uid(account_info) if not uid: self._mark_task(worker_db, task, "failed", "无法从账号或 Cookie 解析 yyuid") return cookie = account_info.get("cookie") or "" if not cookie: self._mark_task(worker_db, task, "failed", "账号 Cookie 为空") return client = HuyaHttpClient(logger=lambda msg: self._push_log("info", f"[{uid}] {msg}")) role_status = client.check_user_bind_game_account( uid=uid, cookie=cookie, b_act_id=b_act_id_int, is_use_outer_act_id=0, ) if role_status is None: self._mark_task(worker_db, task, "error", "虎牙绑定角色查询接口无响应") return if role_status.status != 200: self._mark_task( worker_db, task, "failed", role_status.msg or f"虎牙绑定角色查询失败: {role_status.status}", {"bind_act_id": b_act_id_int, "bind_status": role_status.to_dict()}, ) return if not role_status.accountData.isBindAcount or not role_status.accountData.isBindRole: result = { "bind_act_id": b_act_id_int, "bind_confirmed": False, "bind_status": role_status.to_dict(), } self._mark_task(worker_db, task, "failed", "尚未绑定游戏角色,请先扫码完成绑定", result) return confirm_resp = client.confirm_bind_act_account( uid=uid, cookie=cookie, b_act_id=b_act_id_int, ) if confirm_resp is None: self._mark_task(worker_db, task, "error", "虎牙确认绑定接口无响应") return if confirm_resp.status != 200: result = { "bind_act_id": b_act_id_int, "bind_confirmed": False, "confirm_result": confirm_resp.to_dict(), "before_bind_status": role_status.to_dict(), } self._mark_task( worker_db, task, "failed", confirm_resp.msg or f"虎牙确认绑定失败: {confirm_resp.status}", result, ) return refreshed_status = client.check_user_bind_game_account( uid=uid, cookie=cookie, b_act_id=b_act_id_int, is_use_outer_act_id=1, ) if refreshed_status is None: result = { "bind_act_id": b_act_id_int, "bind_confirmed": False, "confirm_result": confirm_resp.to_dict(), "before_bind_status": role_status.to_dict(), } self._mark_task(worker_db, task, "error", "虎牙活动绑定状态刷新无响应", result) return if refreshed_status.status != 200: result = { "bind_act_id": b_act_id_int, "bind_confirmed": False, "confirm_result": confirm_resp.to_dict(), "before_bind_status": role_status.to_dict(), "bind_status": refreshed_status.to_dict(), } self._mark_task( worker_db, task, "failed", refreshed_status.msg or f"虎牙活动绑定状态刷新失败: {refreshed_status.status}", result, ) return bind_confirmed = bool( refreshed_status.accountData.isBindAcount and refreshed_status.accountData.isBindRole ) role_info = self._bind_role_result(refreshed_status) result = { "bind_act_id": b_act_id_int, "bind_confirmed": bind_confirmed, "confirm_result": confirm_resp.to_dict(), "before_bind_status": role_status.to_dict(), "bind_status": refreshed_status.to_dict(), **role_info, } if not bind_confirmed: self._mark_task(worker_db, task, "failed", "确认后仍未检测到活动绑定角色", result) return self._apply_role_to_account(account, refreshed_status, "bind_confirmed") role_name = self._role_name(refreshed_status) or "已绑定" self._mark_task(worker_db, task, "success", f"确认绑定: {role_name}", result) def _execute_one(self, task_id: int, account_info: dict, config_info: dict, total: int): worker_db = SessionLocal() try: task = worker_db.query(HuyaTask).filter(HuyaTask.id == task_id).first() account = worker_db.query(HuyaAccount).filter(HuyaAccount.id == account_info["account_id"]).first() if not task or not account: return if self._stop.is_set(): self._mark_task(worker_db, task, "failed", "任务已停止") return task.status = "running" task.message = "执行中" task.finished_at = None worker_db.commit() with self._counter_lock: self._started += 1 current = self._started name = self._account_name(account_info) self._push_log("info", f"[{current}/{total}] 开始虎牙任务: {name}") if self.task_type not in { "query_points", "get_bind_qr", "confirm_bind", "query_game_name", "refresh_goods", "refresh_recharge_goods", "create_recharge_order", }: self._mark_task(worker_db, task, "failed", "该虎牙任务执行器暂未实现") self._push_log("warning", f"[{current}] {name} 暂未实现: {self.task_type}") return try: if self.task_type == "query_points": self._execute_query_points(worker_db, task, account, account_info, config_info) elif self.task_type == "refresh_goods": self._execute_refresh_goods(worker_db, task, account, account_info, config_info) elif self.task_type == "refresh_recharge_goods": self._execute_refresh_recharge_goods(worker_db, task, account, account_info, config_info) elif self.task_type == "create_recharge_order": self._execute_create_recharge_order(worker_db, task, account, account_info, config_info) elif self.task_type == "get_bind_qr": self._execute_get_bind_qr(worker_db, task, account, account_info, config_info) elif self.task_type == "confirm_bind": self._execute_confirm_bind(worker_db, task, account, account_info, config_info) elif self.task_type == "query_game_name": self._execute_query_game_name(worker_db, task, account, account_info, config_info) worker_db.refresh(task) if task.status == "success": self._push_log("success", f"[{current}] {name} {task.message}") else: self._push_log("error", f"[{current}] {name} {task.message}") except Exception as exc: self._mark_task(worker_db, task, "error", f"执行异常: {exc}") self._push_log("error", f"[{current}] {name} 执行异常: {exc}") finally: worker_db.close() def run(self): """在线程中执行虎牙批次任务。""" self._push_log( "info", f"虎牙批次 {self.batch_id} 开始,共执行 {self.task_type},并发数: {self.concurrency}", ) try: config = ensure_huya_config(self.db) config_info = { field: huya_config_value(field, getattr(config, field, None)) for field in HUYA_CONFIG_FIELDS } tasks = ( self.db.query(HuyaTask) .options(joinedload(HuyaTask.account)) .filter(HuyaTask.batch_id == self.batch_id) .order_by(HuyaTask.id.asc()) .all() ) task_infos = [] for task in tasks: account = task.account if not account: task.status = "error" task.message = "账号不存在" task.finished_at = datetime.now(timezone.utc) continue task.status = "pending" task.message = "等待执行" task.finished_at = None task_infos.append({ "task_id": task.id, "account_info": { "account_id": account.id, "uid": account.uid or "", "yyuid": account.yyuid or "", "username": account.username or "", "nickname": account.nickname or "", "cookie": account.cookie or "", }, }) self.db.commit() total = len(task_infos) if total == 0: self._push_log("warning", "没有可执行的虎牙任务") self._push_log("result", "") return with ThreadPoolExecutor(max_workers=self.concurrency) as executor: futures = [] for item in task_infos: if self._stop.is_set(): self._push_log("warning", "任务已停止,跳过剩余账号") break futures.append(executor.submit( self._execute_one, item["task_id"], item["account_info"], config_info, total, )) for future in as_completed(futures): try: future.result() except Exception as exc: self._push_log("error", f"虎牙 Worker 异常: {exc}") self._push_log("info", f"虎牙批次 {self.batch_id} 完成") self._push_log("result", "") except Exception as exc: self._push_log("error", f"虎牙批次执行异常: {exc}") self._push_log("result", "") finally: self.db.close() class HuyaBatchRegistry: """管理运行中的虎牙批次。""" def __init__(self): self._batches: dict[str, dict] = {} def register( self, batch_id: str, log_queue: asyncio.Queue, loop: asyncio.AbstractEventLoop, runner: HuyaBatchRunner, ): self._batches[batch_id] = { "log_queue": log_queue, "loop": loop, "runner": runner, } 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) huya_batch_registry = HuyaBatchRegistry()