type: 收窄斗鱼商城与小店执行器类型

This commit is contained in:
yml2213
2026-08-30 19:57:37 +08:00
parent d400db987e
commit 2622e016b5
2 changed files with 301 additions and 75 deletions
+216 -52
View File
@@ -4,29 +4,75 @@ from __future__ import annotations
import random
import time
from datetime import datetime, timezone
from typing import Any, TYPE_CHECKING
from sqlalchemy.orm import Session
from core.douyu import DouyuActivityClient, DouyuActivityError
from ..models import Account, DouyuEsportsGoodsSnapshot, DouyuGoodsSnapshot, DouyuTask
if TYPE_CHECKING:
from .douyu_runner import DouyuBatchRunner
# 兑换节奏与重试(对齐 8.30 浏览器抓包:锁单->支付间隔约 2.2~4.2s;火爆类错误要长退避而不是秒级连打)
DOUYU_EXCHANGE_PRE_CREATE_JITTER = (0.3, 1.2)
DOUYU_EXCHANGE_LOCK_PAY_DELAY_RANGE = (2.0, 4.0)
DOUYU_EXCHANGE_LOCK_TTL_SECONDS = 300
DOUYU_EXCHANGE_RATE_LIMIT_HINTS = ("火爆", "频繁", "稍后", "太热", "限流", "繁忙", "人多", "排队", "手慢")
DOUYU_EXCHANGE_RATE_LIMIT_HINTS = (
"火爆",
"频繁",
"稍后",
"太热",
"限流",
"繁忙",
"人多",
"排队",
"手慢",
)
DOUYU_EXCHANGE_RATE_LIMIT_BACKOFFS = (30, 60, 90)
DOUYU_EXCHANGE_GENERIC_BACKOFFS = (2, 5, 10, 20)
class GoodsMixin:
"""商城兑换域:商品刷新、锁单/支付/兑换、兑换节奏与退避重试。"""
def _execute_refresh_goods(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
if TYPE_CHECKING:
_stop: Any
def _client(self, cookie: str) -> DouyuActivityClient: ...
def _task_payload(self, task: DouyuTask) -> dict: ...
def _upsert_goods(self, db: Session, goods: list[dict]) -> None: ...
def _upsert_esports_goods(self, db: Session, goods: list[dict]) -> None: ...
def _refresh_account_points(
self,
client: DouyuActivityClient,
account: Account,
cookie: str,
*,
ctn: str = "",
) -> dict: ...
def _refresh_esports_handbook(
self, client: DouyuActivityClient, account: Account, *, manual_id: str
) -> dict: ...
def _sleep_interruptible(self, seconds: float) -> bool: ...
def _push_log(self, level: str, message: str) -> None: ...
def _mark_task(self, *args: Any, **kwargs: Any) -> None: ...
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})
self._mark_task(
db,
task,
"success",
f"已刷新商品 {len(goods)}",
{"goods_count": len(goods), "goods": goods},
)
def _execute_refresh_esports_goods(
self,
@@ -51,12 +97,20 @@ class GoodsMixin:
task,
"success",
f"已刷新电竞皮肤 {len(goods)}",
{"goods_count": len(goods), "esports_store_score": result["score"], "goods": goods},
{
"goods_count": len(goods),
"esports_store_score": result["score"],
"goods": goods,
},
)
def _execute_lock_goods(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
def _execute_lock_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()
commodity_id = str(
payload.get("commodity_id") or payload.get("commodityId") or ""
).strip()
if not commodity_id:
self._mark_task(db, task, "failed", "请选择锁定商品")
return
@@ -89,7 +143,10 @@ class GoodsMixin:
"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}",
"account_name": account.nickname
or account.username
or account.uid
or f"#{account.id}",
**result,
},
)
@@ -104,7 +161,9 @@ class GoodsMixin:
):
payload = self._task_payload(task)
order_id = str(payload.get("order_id") or payload.get("orderId") or "").strip()
commodity_id = str(payload.get("commodity_id") or payload.get("commodityId") or "").strip()
commodity_id = str(
payload.get("commodity_id") or payload.get("commodityId") or ""
).strip()
if not order_id:
self._mark_task(db, task, "failed", "锁单订单号不能为空")
return
@@ -133,11 +192,16 @@ class GoodsMixin:
"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}",
"account_name": account.nickname
or account.username
or account.uid
or f"#{account.id}",
}
try:
ctn = client.acf_ccn(refresh_subscribe=False)
result.update(self._refresh_account_points(client, account, cookie, ctn=ctn))
result.update(
self._refresh_account_points(client, account, cookie, ctn=ctn)
)
except Exception as exc:
self._push_log("warning", f"支付锁单后刷新积分失败: {exc}")
self._mark_task(
@@ -206,9 +270,13 @@ class GoodsMixin:
return None
return None
def _execute_exchange_goods(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
def _execute_exchange_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()
commodity_id = str(
payload.get("commodity_id") or payload.get("commodityId") or ""
).strip()
if not commodity_id:
self._mark_task(db, task, "failed", "请选择兑换商品")
return
@@ -226,7 +294,9 @@ class GoodsMixin:
self._push_log("warning", f"查询手册状态失败,按经典流程继续: {exc}")
detail = None
try:
detail = client.query_goods_detail(manual_id=manual_id, commodity_id=commodity_id, rid=rid)["detail"]
detail = client.query_goods_detail(
manual_id=manual_id, commodity_id=commodity_id, rid=rid
)["detail"]
except DouyuActivityError as exc:
self._push_log("warning", f"查询商品详情失败,按经典流程继续: {exc}")
@@ -249,7 +319,10 @@ class GoodsMixin:
def finish_failed(message: str) -> None:
self._mark_task(
db, task, "failed", message,
db,
task,
"failed",
message,
{"commodity_id": commodity_id, "plan": plan, "detail": detail},
)
@@ -266,37 +339,63 @@ class GoodsMixin:
if action == "subscribe":
try:
client.subscribe_commodity(manual_id=manual_id, commodity_id=commodity_id)
client.subscribe_commodity(
manual_id=manual_id, commodity_id=commodity_id
)
except DouyuActivityError as exc:
finish_failed(f"预约到货失败: {exc}")
return
account.bind_status = "goods_subscribed"
account.updated_at = datetime.now(timezone.utc)
self._mark_task(
db, task, "success", f"已预约到货: {name} ({plan['text']})",
db,
task,
"success",
f"已预约到货: {name} ({plan['text']})",
{"commodity_id": commodity_id, "plan": plan, "detail": detail},
)
return
if action == "pre_exchange":
try:
check = client.pre_exchange_check(manual_id=manual_id, commodity_id=commodity_id)
check = client.pre_exchange_check(
manual_id=manual_id, commodity_id=commodity_id
)
user_status = int((check.get("data") or {}).get("userStatus") or 0)
if user_status == 1 and not plan["pre_exchange"]:
self._mark_task(
db, task, "success", f"已预兑: {name},等待开放后继续兑换",
{"commodity_id": commodity_id, "plan": plan, "detail": detail, "check": check},
db,
task,
"success",
f"已预兑: {name},等待开放后继续兑换",
{
"commodity_id": commodity_id,
"plan": plan,
"detail": detail,
"check": check,
},
)
return
confirm = client.confirm_pre_exchange(manual_id=manual_id, commodity_id=commodity_id, rid=rid)
confirm = client.confirm_pre_exchange(
manual_id=manual_id, commodity_id=commodity_id, rid=rid
)
except DouyuActivityError as exc:
finish_failed(f"预兑失败: {exc}")
return
account.bind_status = "goods_pre_exchanged"
account.updated_at = datetime.now(timezone.utc)
self._mark_task(
db, task, "success", f"预兑成功: {name},开放后可继续兑换",
{"commodity_id": commodity_id, "plan": plan, "detail": detail, "check": check, "confirm": confirm},
db,
task,
"success",
f"预兑成功: {name},开放后可继续兑换",
{
"commodity_id": commodity_id,
"plan": plan,
"detail": detail,
"check": check,
"confirm": confirm,
},
)
return
@@ -309,7 +408,9 @@ class GoodsMixin:
# 少量随机抖动(浏览器点击商品到锁单之间有人工间隔)
self._push_log("info", f" 兑换 {name}: {plan['text']}")
if not self._sleep_interruptible(random.uniform(*DOUYU_EXCHANGE_PRE_CREATE_JITTER)):
if not self._sleep_interruptible(
random.uniform(*DOUYU_EXCHANGE_PRE_CREATE_JITTER)
):
self._mark_task(db, task, "stopped", "任务已停止", result)
return
try:
@@ -321,8 +422,14 @@ class GoodsMixin:
)
except DouyuActivityError as exc:
# 锁单本身被频控时做一次长退避重试,避免无脑连打
backoff = DOUYU_EXCHANGE_RATE_LIMIT_BACKOFFS[0] if self._is_rate_limited_error(str(exc)) else 3
self._push_log("info", f" 锁定兑换商品失败({exc}){backoff:.0f}s 后重试一次")
backoff = (
DOUYU_EXCHANGE_RATE_LIMIT_BACKOFFS[0]
if self._is_rate_limited_error(str(exc))
else 3
)
self._push_log(
"info", f" 锁定兑换商品失败({exc}){backoff:.0f}s 后重试一次"
)
if not self._sleep_interruptible(backoff):
self._mark_task(db, task, "stopped", "任务已停止", result)
return
@@ -336,13 +443,19 @@ class GoodsMixin:
except DouyuActivityError as exc2:
finish_failed(f"锁定兑换商品失败: {exc2}")
return
result.update({**locked, "lock_order": locked, "goods": goods.raw if goods else None})
result.update(
{**locked, "lock_order": locked, "goods": goods.raw if goods else None}
)
# 锁单成功后按浏览器节奏(抓包实测 2.24s~4.2s)等待再支付
delay = random.uniform(*DOUYU_EXCHANGE_LOCK_PAY_DELAY_RANGE)
self._push_log("info", f" 锁单成功(订单 {locked['order_id']}){delay:.1f}s 后支付")
self._push_log(
"info", f" 锁单成功(订单 {locked['order_id']}){delay:.1f}s 后支付"
)
if not self._sleep_interruptible(delay):
self._mark_task(db, task, "stopped", "任务已停止,商品锁单仍可能有效", locked)
self._mark_task(
db, task, "stopped", "任务已停止,商品锁单仍可能有效", locked
)
return
try:
payment = self._pay_exchange_with_backoff(
@@ -356,23 +469,30 @@ class GoodsMixin:
finish_failed(str(exc))
return
if payment is None:
self._mark_task(db, task, "stopped", "任务已停止,商品锁单仍可能有效", locked)
self._mark_task(
db, task, "stopped", "任务已停止,商品锁单仍可能有效", locked
)
return
result.update({
"order_id": locked["order_id"],
"exchange_id": payment["exchange_id"],
"commodity_image": payment["commodity_image"] or locked["commodity_image"],
"exchange_num": payment["exchange_num"],
"payment": payment,
})
result.update(
{
"order_id": locked["order_id"],
"exchange_id": payment["exchange_id"],
"commodity_image": payment["commodity_image"]
or locked["commodity_image"],
"exchange_num": payment["exchange_num"],
"payment": payment,
}
)
account.bind_status = "goods_exchanged"
account.updated_at = datetime.now(timezone.utc)
# 兑换成功后自动刷新积分,更新账号最新积分信息(失败不阻断兑换成功)
points_refresh = None
try:
ctn = client.acf_ccn(refresh_subscribe=False)
points_refresh = self._refresh_account_points(client, account, cookie, ctn=ctn)
points_refresh = self._refresh_account_points(
client, account, cookie, ctn=ctn
)
except Exception as exc:
self._push_log("warning", f"兑换后刷新积分失败: {exc}")
self._mark_task(
@@ -384,7 +504,10 @@ class GoodsMixin:
"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}",
"account_name": account.nickname
or account.username
or account.uid
or f"#{account.id}",
**result,
**(points_refresh or {}),
},
@@ -400,7 +523,9 @@ class GoodsMixin:
):
"""兑换电竞手册皮肤,并同步电竞积分。"""
payload = self._task_payload(task)
commodity_id = str(payload.get("commodity_id") or payload.get("commodityId") or "").strip()
commodity_id = str(
payload.get("commodity_id") or payload.get("commodityId") or ""
).strip()
if not commodity_id:
self._mark_task(db, task, "failed", "请选择电竞皮肤")
return
@@ -419,7 +544,9 @@ class GoodsMixin:
client = self._client(cookie)
baseline_points = account.esports_points
try:
baseline = self._refresh_esports_handbook(client, account, manual_id=manual_id)
baseline = self._refresh_esports_handbook(
client, account, manual_id=manual_id
)
baseline_points = baseline["esports_points"]
db.commit()
except Exception as exc:
@@ -439,7 +566,9 @@ class GoodsMixin:
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)
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:
@@ -455,31 +584,66 @@ class GoodsMixin:
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}"
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_limited_goods(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
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"]))
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])}"
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})
self._mark_task(
db,
task,
"success",
message,
{"limited_count": len(limited), "limited_goods": limited},
)
def _execute_query_exchange_records(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
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)
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):
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})
self._mark_task(
db,
task,
"success",
"获取 csrf_token 成功",
{"csrf_token": token, "cookie": client.cookie},
)