feat(douyu): 新增和平小店查询模块(绑定角色/商品列表/点券余额)

- activity_client: getIframeUrl 签发 code/sig,道聚城 getrole.out/recommend/balance 接口
  (2026-08 起需 isCode=1 + authType=delegate + sAnchorId 才能通过 Livelink 校验)
- 迁移 0016: accounts xpd_* 字段、douyu_config 小店配置(actAlias/actId/rid)、
  douyu_xpd_goods_snapshot 商品快照表
- runner/service/router/schema 注册 query_xpd_role/refresh_xpd_goods/query_xpd_balance
  三个任务,任务结果剥离 role.raw 防 openid 泄露
- 前端: 和平小店路由/菜单/DouyuTasksPage 三态/商品下拉/配置弹窗
This commit is contained in:
yml2213
2026-08-06 17:23:59 +08:00
parent f8de495029
commit fdc74eba8c
12 changed files with 702 additions and 25 deletions
@@ -0,0 +1,91 @@
"""新增斗鱼和平小店配置、账号字段与商品快照
Revision ID: 20260806_0016
Revises: 20260805_0015
Create Date: 2026-08-06
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
revision: str = "20260806_0016"
down_revision: Union[str, None] = "20260805_0015"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def _columns(bind, table_name: str) -> set[str]:
if not sa.inspect(bind).has_table(table_name):
return set()
return {column["name"] for column in sa.inspect(bind).get_columns(table_name)}
def upgrade() -> None:
bind = op.get_bind()
config_columns = _columns(bind, "douyu_config")
for column in [
sa.Column("xpd_act_alias", sa.String(length=64), nullable=True),
sa.Column("xpd_act_id", sa.String(length=64), nullable=True),
sa.Column("xpd_rid", sa.String(length=64), nullable=True),
]:
if column.name not in config_columns:
op.add_column("douyu_config", column)
account_columns = _columns(bind, "accounts")
for column in [
sa.Column("xpd_game_name", sa.String(length=128), nullable=True),
sa.Column("xpd_openid", sa.String(length=128), nullable=True),
sa.Column("xpd_role_id", sa.String(length=64), nullable=True),
sa.Column("xpd_plat_id", sa.Integer(), nullable=True),
sa.Column("xpd_area_id", sa.Integer(), nullable=True),
sa.Column("xpd_balance", sa.Integer(), nullable=True),
sa.Column("xpd_bind_status", sa.String(length=32), nullable=True),
]:
if column.name not in account_columns:
op.add_column("accounts", column)
if not sa.inspect(bind).has_table("douyu_xpd_goods_snapshot"):
op.create_table(
"douyu_xpd_goods_snapshot",
sa.Column("id", sa.Integer(), primary_key=True, autoincrement=True),
sa.Column("commodity_id", sa.String(length=64), nullable=False),
sa.Column("name", sa.String(length=256), nullable=True),
sa.Column("price", sa.Integer(), nullable=True),
sa.Column("org_price", sa.Integer(), nullable=True),
sa.Column("category", sa.String(length=32), nullable=True),
sa.Column("goods_left", sa.Integer(), nullable=True),
sa.Column("raw", sa.JSON(), nullable=True),
sa.Column("updated_at", sa.DateTime(), nullable=True),
)
op.create_index(
"ix_douyu_xpd_goods_snapshot_commodity_id",
"douyu_xpd_goods_snapshot",
["commodity_id"],
)
def downgrade() -> None:
bind = op.get_bind()
if sa.inspect(bind).has_table("douyu_xpd_goods_snapshot"):
op.drop_table("douyu_xpd_goods_snapshot")
account_columns = _columns(bind, "accounts")
for column_name in [
"xpd_bind_status",
"xpd_balance",
"xpd_area_id",
"xpd_plat_id",
"xpd_role_id",
"xpd_openid",
"xpd_game_name",
]:
if column_name in account_columns:
op.drop_column("accounts", column_name)
config_columns = _columns(bind, "douyu_config")
for column_name in ["xpd_rid", "xpd_act_id", "xpd_act_alias"]:
if column_name in config_columns:
op.drop_column("douyu_config", column_name)
+25
View File
@@ -70,6 +70,13 @@ class Account(Base):
esports_bind_status = Column(String(32), default="")
esports_change_role_wait_time = Column(Integer, nullable=True)
esports_can_change_time = Column(Integer, nullable=True)
xpd_game_name = Column(String(128), default="")
xpd_openid = Column(String(128), default="")
xpd_role_id = Column(String(64), default="")
xpd_plat_id = Column(Integer, nullable=True)
xpd_area_id = Column(Integer, nullable=True)
xpd_balance = Column(Integer, nullable=True)
xpd_bind_status = Column(String(32), default="")
created_at = Column(DateTime, default=_utcnow)
updated_at = Column(DateTime, default=_utcnow, onupdate=_utcnow)
@@ -132,6 +139,9 @@ class DouyuConfig(Base):
esports_chicken_skin_id = Column(String(64), default="0")
esports_firework_gift_id = Column(String(64), default="24767")
esports_firework_skin_id = Column(String(64), default="3850")
xpd_act_alias = Column(String(64), default="20260623KDQFH")
xpd_act_id = Column(String(64), default="46195")
xpd_rid = Column(String(64), default="9263298")
gold_pay_type = Column(Integer, default=1)
gift_id = Column(String(64), default="23643")
skin_id = Column(String(64), default="2942")
@@ -151,6 +161,21 @@ class DouyuGoodsSnapshot(Base):
updated_at = Column(DateTime, default=_utcnow, onupdate=_utcnow)
class DouyuXpdGoodsSnapshot(Base):
"""斗鱼和平小店商品快照"""
__tablename__ = "douyu_xpd_goods_snapshot"
id = Column(Integer, primary_key=True, autoincrement=True)
commodity_id = Column(String(64), nullable=False, index=True)
name = Column(String(256), default="")
price = Column(Integer, nullable=True)
org_price = Column(Integer, nullable=True)
category = Column(String(32), default="")
goods_left = Column(Integer, nullable=True)
raw = Column(JSON, nullable=True)
updated_at = Column(DateTime, default=_utcnow, onupdate=_utcnow)
class DouyuEsportsGoodsSnapshot(Base):
"""斗鱼电竞手册皮肤商城商品快照"""
__tablename__ = "douyu_esports_goods_snapshot"
+28 -1
View File
@@ -12,7 +12,7 @@ from sqlalchemy.orm import Session, joinedload
from ..database import SessionLocal, get_db
from ..deps import authenticate_websocket, get_current_user, require_permission
from ..models import Account, DouyuConfig, DouyuEsportsGoodsSnapshot, DouyuGoodsSnapshot, DouyuTask, User
from ..models import Account, DouyuConfig, DouyuEsportsGoodsSnapshot, DouyuGoodsSnapshot, DouyuTask, DouyuXpdGoodsSnapshot, User
from ..permissions import user_has_permission
from ..schemas import (
DouyuConfigOut,
@@ -21,6 +21,7 @@ from ..schemas import (
DouyuTaskAccountOut,
DouyuTaskBatchRequest,
DouyuTaskOut,
DouyuXpdGoodsOut,
)
from ..services.douyu_runner import DouyuBatchRunner, douyu_batch_registry
from ..services.douyu_service import (
@@ -113,6 +114,13 @@ def _account_out(account: Account) -> DouyuTaskAccountOut:
esports_bind_status=account.esports_bind_status or "",
esports_change_role_wait_time=account.esports_change_role_wait_time,
esports_can_change_time=account.esports_can_change_time,
xpd_game_name=account.xpd_game_name or "",
xpd_openid=account.xpd_openid or "",
xpd_role_id=account.xpd_role_id or "",
xpd_plat_id=account.xpd_plat_id,
xpd_area_id=account.xpd_area_id,
xpd_balance=account.xpd_balance,
xpd_bind_status=account.xpd_bind_status or "",
assigned_to=account.assigned_to,
assigned_username=account.assigned_user.username if account.assigned_user else None,
)
@@ -168,6 +176,11 @@ def _sanitize_task_result(result: dict | None, task_type: str, *, include_detail
):
data.pop(key, None)
# 和平小店任务 result 内嵌 role(含 gameOpenId 等),列表接口剥离其原始快照
role = data.get("role")
if isinstance(role, dict):
data["role"] = {key: value for key, value in role.items() if key != "raw"}
goods = data.get("goods")
if isinstance(goods, list):
data.pop("goods", None)
@@ -223,6 +236,9 @@ def _config_out(config: DouyuConfig) -> DouyuConfigOut:
esports_chicken_skin_id=douyu_config_value("esports_chicken_skin_id", config.esports_chicken_skin_id),
esports_firework_gift_id=douyu_config_value("esports_firework_gift_id", config.esports_firework_gift_id),
esports_firework_skin_id=douyu_config_value("esports_firework_skin_id", config.esports_firework_skin_id),
xpd_act_alias=douyu_config_value("xpd_act_alias", config.xpd_act_alias),
xpd_act_id=douyu_config_value("xpd_act_id", config.xpd_act_id),
xpd_rid=douyu_config_value("xpd_rid", config.xpd_rid),
gold_pay_type=douyu_config_value("gold_pay_type", config.gold_pay_type),
gift_id=douyu_config_value("gift_id", config.gift_id),
skin_id=douyu_config_value("skin_id", config.skin_id),
@@ -267,6 +283,7 @@ def list_task_accounts(
Account.tag.ilike(pattern),
Account.game_name.ilike(pattern),
Account.esports_game_name.ilike(pattern),
Account.xpd_game_name.ilike(pattern),
))
total = None
if page is not None:
@@ -331,6 +348,16 @@ def list_esports_goods(
return rows
@router.get("/xpd-goods", response_model=list[DouyuXpdGoodsOut])
def list_xpd_goods(
db: Session = Depends(get_db),
current: User = Depends(require_permission("douyu:task")),
):
"""查看已缓存的和平小店商品。"""
rows = db.query(DouyuXpdGoodsSnapshot).order_by(DouyuXpdGoodsSnapshot.id.asc()).all()
return rows
@router.post("/tasks/batch")
async def create_task_batch(
req: DouyuTaskBatchRequest,
+44
View File
@@ -554,6 +554,9 @@ class DouyuConfigOut(BaseModel):
esports_chicken_skin_id: str = "0"
esports_firework_gift_id: str = "24767"
esports_firework_skin_id: str = "3850"
xpd_act_alias: str = "20260623KDQFH"
xpd_act_id: str = "46195"
xpd_rid: str = "9263298"
gold_pay_type: int = 1
gift_id: str = "23643"
skin_id: str = "2942"
@@ -576,6 +579,9 @@ class DouyuConfigOut(BaseModel):
"esports_chicken_skin_id": self.esports_chicken_skin_id,
"esports_firework_gift_id": self.esports_firework_gift_id,
"esports_firework_skin_id": self.esports_firework_skin_id,
"xpd_act_alias": self.xpd_act_alias,
"xpd_act_id": self.xpd_act_id,
"xpd_rid": self.xpd_rid,
"gold_pay_type": self.gold_pay_type,
"gift_id": self.gift_id,
"skin_id": self.skin_id,
@@ -598,6 +604,9 @@ class DouyuConfigUpdate(BaseModel):
esports_chicken_skin_id: Optional[str] = None
esports_firework_gift_id: Optional[str] = None
esports_firework_skin_id: Optional[str] = None
xpd_act_alias: Optional[str] = None
xpd_act_id: Optional[str] = None
xpd_rid: Optional[str] = None
gold_pay_type: Optional[int] = Field(None, ge=1, le=9)
gift_id: Optional[str] = None
skin_id: Optional[str] = None
@@ -689,10 +698,45 @@ class DouyuTaskAccountOut(BaseModel):
esports_bind_status: str = ""
esports_change_role_wait_time: Optional[int] = None
esports_can_change_time: Optional[int] = None
xpd_game_name: str = ""
xpd_openid: str = ""
xpd_role_id: str = ""
xpd_plat_id: Optional[int] = None
xpd_area_id: Optional[int] = None
xpd_balance: Optional[int] = None
xpd_bind_status: str = ""
assigned_to: Optional[int] = None
assigned_username: Optional[str] = None
class DouyuXpdGoodsOut(BaseModel):
id: int
commodity_id: str
name: str = ""
price: Optional[int] = None
org_price: Optional[int] = None
category: str = ""
goods_left: Optional[int] = None
raw: Optional[dict[str, Any]] = None
updated_at: Optional[datetime] = None
model_config = ConfigDict(from_attributes=True)
@model_serializer
def _serialize(self) -> dict[str, Any]:
return {
"id": self.id,
"commodity_id": self.commodity_id,
"name": self.name,
"price": self.price,
"org_price": self.org_price,
"category": self.category,
"goods_left": self.goods_left,
"raw": self.raw,
"updated_at": _ensure_tz(self.updated_at).isoformat() if self.updated_at else None,
}
# ---- 代理配置 ----
class ProxyConfigOut(BaseModel):
enabled: bool = False
+163 -1
View File
@@ -15,7 +15,7 @@ from sqlalchemy.orm import Session, joinedload
from core.douyu import DouyuActivityClient, DouyuActivityError
from ..database import SessionLocal
from ..models import Account, DouyuEsportsGoodsSnapshot, DouyuGoodsSnapshot, DouyuTask
from ..models import Account, DouyuEsportsGoodsSnapshot, DouyuGoodsSnapshot, DouyuTask, DouyuXpdGoodsSnapshot
from .douyu_service import (
DOUYU_CONFIG_FIELDS,
account_uid,
@@ -415,6 +415,30 @@ class DouyuBatchRunner:
row.updated_at = now
db.commit()
def _upsert_xpd_goods(self, db: Session, goods: list[dict]) -> None:
"""写入和平小店商品快照。"""
now = datetime.now(timezone.utc)
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 "")
row.goods_left = self._to_int(raw.get("goods_left") or raw.get("iGoodsLeft"))
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}
@@ -993,6 +1017,141 @@ class DouyuBatchRunner:
{"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_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_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)
@@ -2284,6 +2443,9 @@ class DouyuBatchRunner:
"query_gold_balance": self._execute_query_gold_balance,
"query_exchange_records": self._execute_query_exchange_records,
"prefetch_csrf_token": self._execute_prefetch_csrf_token,
"query_xpd_role": self._execute_query_xpd_role,
"refresh_xpd_goods": self._execute_refresh_xpd_goods,
"query_xpd_balance": self._execute_query_xpd_balance,
}.get(task.task_type)
if handler is None:
self._mark_task(worker_db, task, "failed", "不支持的任务类型")
+7 -1
View File
@@ -38,6 +38,9 @@ SUPPORTED_DOUYU_TASK_TYPES = {
"refresh_goods": "刷新商品列表",
"query_exchange_records": "一键查询兑换记录",
"prefetch_csrf_token": "一键获取兑换 CSRF Token",
"query_xpd_role": "查询小店绑定角色",
"refresh_xpd_goods": "刷新小店商品列表",
"query_xpd_balance": "查询小店点券余额",
}
@@ -59,6 +62,9 @@ DOUYU_CONFIG_DEFAULTS = {
"gold_pay_type": 1,
"gift_id": "23643",
"skin_id": "2942",
"xpd_act_alias": "20260623KDQFH",
"xpd_act_id": "46195",
"xpd_rid": "9263298",
}
DOUYU_CONFIG_FIELDS = tuple(DOUYU_CONFIG_DEFAULTS.keys())
@@ -169,7 +175,7 @@ def create_douyu_planned_tasks(
raise ValueError("不支持的任务类型")
accounts = visible_douyu_task_accounts(db, account_ids)
if task_type in {"refresh_goods", "refresh_esports_goods"} and accounts:
if task_type in {"refresh_goods", "refresh_esports_goods", "refresh_xpd_goods"} and accounts:
# 商品快照是全局数据,一个可用 CK 足够;没有 CK 时前端无法选账号创建任务。
accounts = accounts[:1]