Files
live-hub-py/web/backend/services/huya_runner.py
T
yml2213 a13dde57dd 恢复虎牙绑定二维码自动识别角色
生成二维码后后端轮询扫码与角色同步,扫码完成无需再手动点查询角色。
用 bind_polling 区分真实轮询与僵尸 running,避免再次锁死任务台。
2026-07-24 12:01:11 +08:00

1665 lines
62 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""虎牙任务批次执行器。"""
import asyncio
import time
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 core.huya.cookie_utils import normalize_huya_cookie
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_PAYMENT_POLL_SECONDS = 180
HUYA_PAYMENT_POLL_INTERVAL = 3
HUYA_BIND_ROLE_POLL_SECONDS = 180
HUYA_BIND_ROLE_POLL_INTERVAL = 3
HUYA_BIND_ZT_UUID = "b02faae1"
HUYA_BIND_ROOM_ID = "30596253"
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":
try:
huya_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"[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 ""
@staticmethod
def _has_bind_role(bind_status) -> bool:
return bool(bind_status and HuyaBatchRunner._role_name(bind_status))
@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 = [game_role.platName, game_role.areaName]
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")
@staticmethod
def _parse_scheduled_time(value) -> datetime | None:
text = str(value or "").strip()
if not text:
return None
try:
normalized = text.replace("Z", "+00:00")
dt = datetime.fromisoformat(normalized)
except ValueError:
return None
if dt.tzinfo is None:
return dt.astimezone()
return dt
def _wait_until(self, when: datetime, uid: int) -> bool:
target = when.timestamp()
local_text = self._format_local_time(int(target))
self._push_log("info", f"[{uid}] 定时兑换等待到 {local_text}")
while not self._stop.is_set():
remaining = target - time.time()
if remaining <= 0:
return True
time.sleep(min(0.2, max(0.02, remaining)))
return False
@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),
}
@classmethod
def _bind_ready_result(cls, bind_status) -> dict:
role_info = cls._bind_role_result(bind_status)
return {
**role_info,
**cls._bind_change_state(bind_status),
"bind_status": bind_status.to_dict(),
"bind_ready_for_confirm": bool(role_info["role_name"]),
"bind_phase": "role_ready" if role_info["role_name"] else "waiting_role",
}
def _resolve_bind_status(
self,
client: HuyaHttpClient,
uid: int,
cookie: str,
b_act_id_int: int,
):
"""按活动页逻辑解析绑定状态,优先返回含角色的状态。"""
outer_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 outer_status is None:
return None, {}
query_result = {
"bind_status_source": "outer",
"outer_bind_status": outer_status.to_dict(),
}
if outer_status.status != 200:
return outer_status, query_result
chosen_status = outer_status
account_data = outer_status.accountData
should_check_inner = (
not self._has_bind_role(outer_status)
and bool(account_data.isNeedActCheck or not account_data.isBindAcount or not account_data.isBindRole)
)
if should_check_inner:
inner_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 inner_status is not None:
query_result["inner_bind_status"] = inner_status.to_dict()
if inner_status.status == 200 and self._has_bind_role(inner_status):
chosen_status = inner_status
query_result["bind_status_source"] = "inner"
query_result["bind_status"] = chosen_status.to_dict()
return chosen_status, query_result
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 account.game_name
account.game_channel = self._role_channel(bind_status) or account.game_channel
account.updated_at = datetime.now(timezone.utc)
@staticmethod
def _bind_redirect_url(config_info: dict) -> str:
room_pid = str(config_info.get("room_pid") or "").strip()
if not room_pid:
return ""
return (
f"https://zt.huya.com/{HUYA_BIND_ZT_UUID}/pc/index.html"
f"?sourceId={HUYA_RECHARGE_SOURCE_ID}"
f"&pid={room_pid}"
f"&anchorUid={room_pid}"
f"&roomid={HUYA_BIND_ROOM_ID}"
)
def _wait_bind_role_result(
self,
client: HuyaHttpClient,
worker_db: Session,
task: HuyaTask,
account: HuyaAccount,
uid: int,
cookie: str,
b_act_id_int: int,
result: dict,
) -> tuple[str, dict]:
"""生成二维码后轮询扫码状态与角色同步,直到识别到角色、超时或停止。"""
deadline = time.monotonic() + HUYA_BIND_ROLE_POLL_SECONDS
qrcode_token = str(result.get("qrcode_token") or "")
qrcode_finished = not qrcode_token
result["bind_polling"] = True
self._update_task_progress(worker_db, task, "running", "已生成绑定小程序码,等待扫码绑定", result)
while not self._stop.is_set() and time.monotonic() < deadline:
if self._stop.wait(HUYA_BIND_ROLE_POLL_INTERVAL):
break
if qrcode_token and not qrcode_finished:
qrcode_status = client.get_livelink_qrcode_status(qrcode_token, timeout=10.0)
if qrcode_status is not None:
result["qrcode_status"] = qrcode_status
if qrcode_status["is_expired"] or qrcode_status["is_failure"]:
result.update({
"bind_phase": "qrcode_expired",
"bind_ready_for_confirm": False,
"bind_polling": False,
})
self._update_task_progress(
worker_db,
task,
"running",
"绑定小程序码已失效,请重新获取",
result,
)
return "", result
if qrcode_status["is_completed"]:
qrcode_finished = True
result["bind_phase"] = "qrcode_completed"
elif qrcode_status["is_scan"]:
result["bind_phase"] = "qrcode_scanned"
else:
result["bind_phase"] = "waiting_scan"
bind_status, bind_query_result = self._resolve_bind_status(
client=client,
uid=uid,
cookie=cookie,
b_act_id_int=b_act_id_int,
)
if bind_status is None:
continue
result.update(bind_query_result)
if bind_status.status != 200:
result.update({
"bind_phase": "role_check_failed",
"bind_status": bind_status.to_dict(),
})
self._update_task_progress(
worker_db,
task,
"running",
bind_status.msg or "等待绑定角色同步",
result,
)
continue
previous_phase = result.get("bind_phase")
ready = self._bind_ready_result(bind_status)
# 角色未就绪时保留扫码阶段文案,避免状态来回跳。
if not ready["role_name"] and previous_phase in {
"waiting_scan",
"qrcode_scanned",
"qrcode_completed",
}:
ready["bind_phase"] = previous_phase
result.update(ready)
result["bind_polling"] = True
role_name = ready["role_name"]
if role_name:
self._apply_role_to_account(account, bind_status, "game_queried")
result["bind_polling"] = False
self._update_task_progress(
worker_db,
task,
"running",
f"已识别角色: {role_name},待确认绑定",
result,
)
return role_name, result
if result.get("bind_phase") == "qrcode_completed":
message = "小程序绑定已完成,等待角色同步"
elif result.get("bind_phase") == "qrcode_scanned":
message = "已扫码,等待小程序绑定完成"
else:
message = "已生成绑定小程序码,等待扫码绑定"
self._update_task_progress(worker_db, task, "running", message, result)
result.update({
"bind_phase": "role_timeout" if not self._stop.is_set() else "stopped",
"bind_ready_for_confirm": False,
"bind_polling": False,
})
return "", result
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 _update_task_progress(
self,
worker_db: Session,
task: HuyaTask,
status: str,
message: str,
result: Optional[dict] = None,
):
task.status = status
task.message = message
if result is not None:
task.result = result
worker_db.commit()
@staticmethod
def _huya_order_status_label(status: int) -> str:
from core.huya.shop_structs import OrderStatus
labels = {
OrderStatus.DEPOSIT_WAIT_PAY: "待支付",
OrderStatus.DEPOSIT_PAID: "已支付",
OrderStatus.WAIT_DELIVER: "待发货",
OrderStatus.WAIT_RECEIVE: "待收货",
OrderStatus.FINISHED: "已完成",
OrderStatus.FINISHED_CLOSED: "已关闭",
OrderStatus.CANCELLED: "已取消",
OrderStatus.BALANCE_WAIT_PAY: "尾款待支付",
OrderStatus.CANCELLED_BALANCE_EXPIRED: "尾款超时取消",
}
return labels.get(int(status or 0), str(status or "未知"))
@classmethod
def _is_huya_order_paid(cls, order) -> bool:
from core.huya.shop_structs import OrderStatus
paid_statuses = {
OrderStatus.DEPOSIT_PAID,
OrderStatus.WAIT_DELIVER,
OrderStatus.WAIT_RECEIVE,
OrderStatus.FINISHED,
OrderStatus.FINISHED_CLOSED,
}
return int(getattr(order, "payTime", 0) or 0) > 0 or int(getattr(order, "orderStatus", 0) or 0) in paid_statuses
def _wait_recharge_payment(
self,
client: HuyaHttpClient,
uid: int,
guid: str,
cookie: str,
order_id: int,
result: dict,
) -> tuple[str, dict | None]:
deadline = time.time() + HUYA_PAYMENT_POLL_SECONDS
order_id_text = str(order_id)
last_order = None
while not self._stop.is_set() and time.time() < deadline:
resp = client.query_user_order_list(
uid=uid,
guid=guid,
cookie=cookie,
offset=0,
page_size=10,
order_type=1,
status=0,
timeout=10.0,
)
checked_at = datetime.now(timezone.utc).isoformat()
if resp is not None and getattr(resp, "orders", None):
for order in resp.orders:
if str(getattr(order, "orderId", "")) != order_id_text:
continue
last_order = order.to_dict()
status = int(getattr(order, "orderStatus", 0) or 0)
result.update({
"payment_checked_at": checked_at,
"payment_order": last_order,
"payment_order_status": status,
"payment_order_status_label": self._huya_order_status_label(status),
})
if self._is_huya_order_paid(order):
result.update({
"payment_status": "paid",
"payment_status_label": "已支付",
"payment_paid": True,
"payment_paid_at": checked_at,
})
return "paid", last_order
break
else:
result["payment_checked_at"] = checked_at
if self._stop.wait(HUYA_PAYMENT_POLL_INTERVAL):
break
result.update({
"payment_status": "timeout" if not self._stop.is_set() else "stopped",
"payment_status_label": "等待支付超时" if not self._stop.is_set() else "已停止监听",
"payment_paid": False,
"payment_timeout_seconds": HUYA_PAYMENT_POLL_SECONDS,
})
if last_order:
result["payment_order"] = last_order
return result["payment_status"], last_order
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_query_exchange_records(
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_user_prize_records(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
records = result.get("records", [])
for index, item in enumerate(records, start=1):
item["index"] = index
item["exchange_time_text"] = self._format_local_time(int(item.get("exchange_time") or 0))
if item.get("score") is not None:
item["score_text"] = f"{int(item.get('score') or 0)}积分"
account.status = "exchange_records_queried"
account.updated_at = datetime.now(timezone.utc)
count = len(records)
message = f"兑换记录 {count} 条" if count else "暂无兑换记录"
self._mark_task(worker_db, task, "success", message, 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})
def _execute_exchange_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
product_id = self._to_int(self.payload.get("product_id"))
if not product_id:
self._mark_task(worker_db, task, "failed", "请选择兑换商品")
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
snapshot = worker_db.query(HuyaGoodsSnapshot).filter(
HuyaGoodsSnapshot.product_id == str(product_id)
).first()
product_name = str(self.payload.get("product_name") or (snapshot.name if snapshot else "") or product_id)
scheduled_at = self._parse_scheduled_time(self.payload.get("scheduled_at"))
if self.payload.get("scheduled_at") and scheduled_at is None:
self._mark_task(worker_db, task, "failed", "定时兑换时间格式无效")
return
if scheduled_at and scheduled_at.timestamp() > time.time():
if not self._wait_until(scheduled_at, uid):
self._mark_task(worker_db, task, "stopped", "兑换任务已停止")
return
client = HuyaHttpClient(logger=lambda msg: self._push_log("info", f"[{uid}] {msg}"))
response = client.score_exchange_prize(uid=uid, cookie=cookie, sid=sid_int, pid=product_id)
if response is None:
self._mark_task(worker_db, task, "error", "虎牙兑换接口无响应")
return
result = response.to_dict()
result.update({
"sid": sid_int,
"product_id": str(product_id),
"product_name": product_name,
"scheduled_at": scheduled_at.isoformat() if scheduled_at else "",
"executed_at": datetime.now(timezone.utc).isoformat(),
"goods": snapshot.raw if snapshot else None,
})
if response.status != 200:
self._mark_task(
worker_db,
task,
"failed",
response.msg or f"虎牙兑换失败: {response.status}",
result,
)
return
account.status = "goods_exchanged"
account.updated_at = datetime.now(timezone.utc)
message = response.msg or f"兑换成功: {product_name}"
self._mark_task(worker_db, task, "success", message, result)
@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,
"payment_status": "pending",
"payment_status_label": "等待支付",
"payment_paid": False,
"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._update_task_progress(worker_db, task, "running", f"{message},等待扫码支付", result)
self._push_log("info", f"[{uid}] 已生成虎牙支付二维码,开始监听订单 {order_resp.orderId}")
payment_status, payment_order = self._wait_recharge_payment(
client=client,
uid=uid,
guid="",
cookie=cookie,
order_id=order_resp.orderId,
result=result,
)
account.updated_at = datetime.now(timezone.utc)
if payment_status == "paid":
account.status = "recharge_paid"
paid_message = f"支付成功: {product_name} x{count} {result['amount_text']}"
if payment_order and payment_order.get("pay_time"):
paid_message += f",支付时间 {self._format_local_time(int(payment_order['pay_time']) // 1000)}"
self._mark_task(worker_db, task, "success", paid_message, result)
return
if payment_status == "stopped":
account.status = "recharge_order_created"
self._mark_task(worker_db, task, "stopped", f"{message},已停止监听支付", result)
return
account.status = "recharge_order_created"
timeout_message = f"{message}{result['payment_status_label']}"
self._mark_task(worker_db, task, "success", timeout_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, bind_query_result = self._resolve_bind_status(
client=client,
uid=uid,
cookie=cookie,
b_act_id_int=b_act_id_int,
)
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_query_result or 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 ""
bind_redirect_url = self._bind_redirect_url(config_info)
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,
redirect_url=bind_redirect_url,
)
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.get("qrcode_token") or "",
**bind_query_result,
"bind_phase": "waiting_scan" if mini_qrcode.get("qrcode_token") else "waiting_role",
"bind_ready_for_confirm": False,
"bind_polling": True,
"bind_redirect_url": bind_redirect_url,
**role_info,
**change_state,
"profile": {
"nick": profile_nick,
"avatar": profile_avatar,
},
}
account.status = "bind_qr_generated"
account.game_name = role_info["role_name"] or account.game_name
account.game_channel = self._role_channel(bind_status) or account.game_channel
account.nickname = profile_nick or account.nickname
account.updated_at = datetime.now(timezone.utc)
# 生成二维码后自动轮询扫码/角色,避免用户必须手动点「查询角色」。
role_name, result = self._wait_bind_role_result(
client=client,
worker_db=worker_db,
task=task,
account=account,
uid=uid,
cookie=cookie,
b_act_id_int=b_act_id_int,
result=result,
)
if role_name:
self._mark_task(worker_db, task, "success", f"已识别角色: {role_name},待确认绑定", result)
return
if result.get("bind_phase") == "stopped":
self._mark_task(worker_db, task, "stopped", "任务已停止", result)
return
if result.get("bind_phase") == "qrcode_expired":
self._mark_task(worker_db, task, "failed", "绑定小程序码已失效,请重新获取", result)
return
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, bind_query_result = self._resolve_bind_status(
client=client,
uid=uid,
cookie=cookie,
b_act_id_int=b_act_id_int,
)
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_query_result},
)
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_query_result,
"bind_ready_for_confirm": bool(role_info["role_name"]),
"bind_phase": "role_ready" if role_info["role_name"] else "waiting_role",
}
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, role_query_result = self._resolve_bind_status(
client=client,
uid=uid,
cookie=cookie,
b_act_id_int=b_act_id_int,
)
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, **role_query_result},
)
return
role_info = self._bind_role_result(role_status)
role_name = role_info["role_name"]
if not role_name:
result = {
"bind_act_id": b_act_id_int,
"bind_confirmed": False,
**role_query_result,
}
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(),
**role_query_result,
}
self._mark_task(
worker_db,
task,
"failed",
confirm_resp.msg or f"虎牙确认绑定失败: {confirm_resp.status}",
result,
)
return
refreshed_status, refreshed_query_result = self._resolve_bind_status(
client=client,
uid=uid,
cookie=cookie,
b_act_id_int=b_act_id_int,
)
if refreshed_status is None:
result = {
"bind_act_id": b_act_id_int,
"bind_confirmed": True,
"confirm_result": confirm_resp.to_dict(),
"before_bind_status": role_status.to_dict(),
**role_info,
**role_query_result,
}
self._apply_role_to_account(account, role_status, "bind_confirmed")
self._mark_task(worker_db, task, "success", f"确认绑定: {role_name}", result)
return
if refreshed_status.status != 200:
result = {
"bind_act_id": b_act_id_int,
"bind_confirmed": True,
"confirm_result": confirm_resp.to_dict(),
"before_bind_status": role_status.to_dict(),
**role_info,
**role_query_result,
"refresh_error": refreshed_query_result,
}
self._apply_role_to_account(account, role_status, "bind_confirmed")
self._mark_task(worker_db, task, "success", f"确认绑定: {role_name}", result)
return
refreshed_confirmed = bool(
refreshed_status.accountData.isBindAcount
and refreshed_status.accountData.isBindRole
)
refreshed_role_info = self._bind_role_result(refreshed_status)
final_status = refreshed_status if refreshed_role_info["role_name"] else role_status
final_role_info = refreshed_role_info if refreshed_role_info["role_name"] else role_info
result = {
"bind_act_id": b_act_id_int,
"bind_confirmed": True,
"refreshed_is_bound": refreshed_confirmed,
"confirm_result": confirm_resp.to_dict(),
"before_bind_status": role_status.to_dict(),
**role_query_result,
"refresh_result": refreshed_query_result,
**final_role_info,
}
self._apply_role_to_account(account, final_status, "bind_confirmed")
role_name = final_role_info["role_name"] or role_name 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, "stopped", "任务已停止")
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",
"query_exchange_records",
"refresh_goods",
"refresh_recharge_goods",
"exchange_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 == "exchange_goods":
self._execute_exchange_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)
elif self.task_type == "query_exchange_records":
self._execute_query_exchange_records(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": normalize_huya_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] = {}
self._lock = threading.Lock()
def _cleanup_locked(self, ttl_seconds: int = 300):
now = time.time()
expired = [
batch_id
for batch_id, batch in self._batches.items()
if batch.get("finished") and now - float(batch.get("finished_at") or now) > ttl_seconds
]
for batch_id in expired:
self._batches.pop(batch_id, None)
def register(
self,
batch_id: str,
log_queue: asyncio.Queue,
loop: asyncio.AbstractEventLoop,
runner: HuyaBatchRunner,
):
with self._lock:
self._cleanup_locked()
self._batches[batch_id] = {
"log_queue": log_queue,
"loop": loop,
"runner": runner,
"finished": False,
"finished_at": None,
}
def get(self, batch_id: str):
with self._lock:
self._cleanup_locked()
return self._batches.get(batch_id)
def active_ids(self) -> set[str]:
with self._lock:
self._cleanup_locked()
return {
batch_id
for batch_id, batch in self._batches.items()
if not batch.get("finished")
}
def mark_finished(self, batch_id: str):
with self._lock:
batch = self._batches.get(batch_id)
if not batch:
return
batch["finished"] = True
batch["finished_at"] = time.time()
def pop(self, batch_id: str):
with self._lock:
return self._batches.pop(batch_id, None)
huya_batch_registry = HuyaBatchRegistry()