优化斗鱼充值

This commit is contained in:
yml2213
2026-08-13 21:43:27 +08:00
parent 902e669614
commit 0bf69c1240
19 changed files with 450 additions and 266 deletions
@@ -0,0 +1,44 @@
"""持久化斗鱼供应商直充外部订单号
Revision ID: 20260813_0027
Revises: 20260813_0026
Create Date: 2026-08-13
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
revision: str = "20260813_0027"
down_revision: Union[str, None] = "20260813_0026"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
bind = op.get_bind()
inspector = sa.inspect(bind)
if not inspector.has_table("douyu_tasks"):
return
columns = {column["name"] for column in inspector.get_columns("douyu_tasks")}
if "supplier_out_order_id" not in columns:
op.add_column("douyu_tasks", sa.Column("supplier_out_order_id", sa.String(length=64), nullable=True))
op.create_index(
"ix_douyu_tasks_supplier_out_order_id",
"douyu_tasks",
["supplier_out_order_id"],
unique=True,
)
def downgrade() -> None:
bind = op.get_bind()
inspector = sa.inspect(bind)
if not inspector.has_table("douyu_tasks"):
return
columns = {column["name"] for column in inspector.get_columns("douyu_tasks")}
if "supplier_out_order_id" in columns:
op.drop_index("ix_douyu_tasks_supplier_out_order_id", table_name="douyu_tasks")
op.drop_column("douyu_tasks", "supplier_out_order_id")
+2
View File
@@ -119,6 +119,8 @@ class DouyuTask(Base):
task_type = Column(String(64), nullable=False, index=True)
# 任务归属工作台;避免同一账号的精英/电竞/小店任务在前端串行展示或串弹二维码。
handbook_scope = Column(String(16), nullable=False, default="legacy", index=True)
# 供应商直充的来源订单号,用于轮询与异步回调的幂等关联。
supplier_out_order_id = Column(String(64), nullable=True, unique=True, index=True)
status = Column(String(32), default="pending")
message = Column(String(512), default="")
result = Column(JSON, nullable=True)
+76 -1
View File
@@ -6,7 +6,8 @@ import asyncio
import threading
from datetime import datetime, timezone
from fastapi import APIRouter, Depends, HTTPException, Query, WebSocket, WebSocketDisconnect
from fastapi import APIRouter, Depends, HTTPException, Query, Request, WebSocket, WebSocketDisconnect
from loguru import logger
from sqlalchemy import or_
from sqlalchemy.orm import Session, joinedload
@@ -37,11 +38,85 @@ from ..services.douyu_service import (
douyu_task_payload,
ensure_douyu_config,
)
from core.douyu import FishFinRechargeClient, FishFinRechargeConfig, FishFinRechargeConfigError
router = APIRouter(prefix="/api/douyu", tags=["斗鱼活动"])
def _supplier_response_value(payload: dict, *keys: str):
"""兼容供应商把订单字段放在根节点、data 或 result 中。"""
data = payload.get("data") if isinstance(payload.get("data"), dict) else {}
result = payload.get("result") if isinstance(payload.get("result"), dict) else {}
for source in (payload, data, result):
for key in keys:
if source.get(key) is not None:
return source[key]
return None
@router.post("/supplier-recharge/callback")
async def supplier_recharge_callback(request: Request, db: Session = Depends(get_db)):
"""接收供应商直充终态通知,验签后按 out_order_id 幂等更新任务。"""
try:
payload = await request.json()
except ValueError as exc:
raise HTTPException(status_code=400, detail="供应商回调不是 JSON") from exc
if not isinstance(payload, dict):
raise HTTPException(status_code=400, detail="供应商回调格式无效")
try:
client = FishFinRechargeClient(FishFinRechargeConfig.from_env())
except FishFinRechargeConfigError as exc:
logger.error("[douyu] 供应商直充回调配置无效: {}", exc)
raise HTTPException(status_code=503, detail="供应商直充回调未配置") from exc
if not client.verify_response_sign(payload, "POST"):
logger.warning("[douyu] 供应商直充回调验签失败: keys={}", sorted(payload))
raise HTTPException(status_code=401, detail="供应商回调签名无效")
out_order_id = str(_supplier_response_value(payload, "out_order_id") or "").strip()
if not out_order_id:
raise HTTPException(status_code=400, detail="供应商回调缺少 out_order_id")
task = (
db.query(DouyuTask)
.filter(DouyuTask.supplier_out_order_id == out_order_id, DouyuTask.task_type == "create_gold_qr")
.first()
)
if not task:
raise HTTPException(status_code=404, detail="供应商回调订单不存在")
status = DouyuBatchRunner._to_int(_supplier_response_value(payload, "order_status", "orderStatus"))
result = dict(task.result) if isinstance(task.result, dict) else {}
result.update({
"out_order_id": out_order_id,
"order_id": _supplier_response_value(payload, "order_id", "orderId") or result.get("order_id"),
"supplier_order_status": status,
"supplier_order": DouyuBatchRunner._supplier_result(payload),
"supplier_callback_received": True,
})
if task.status not in {"success", "failed", "stopped"}:
if status == 2:
task.status = "success"
task.message = "供应商直充成功(异步通知)"
task.finished_at = datetime.now(timezone.utc)
if task.account:
task.account.bind_status = "gold_recharged"
task.account.updated_at = datetime.now(timezone.utc)
elif status in {3, 4}:
reason = str(_supplier_response_value(payload, "fail_reason", "message", "msg") or "供应商直充失败")
task.status = "failed"
task.message = reason[:512]
task.finished_at = datetime.now(timezone.utc)
else:
task.status = "running"
task.message = f"供应商直充订单处理中(状态 {status if status is not None else '-'}"
task.result = result
db.commit()
logger.info("[douyu] 供应商直充回调已处理: out_order_id={} status={}", out_order_id, status)
acknowledgement = {"code": 200, "message": "success"}
acknowledgement["sign"] = client.sign(acknowledgement, "POST")
return acknowledgement
def _can_view_all(user: User) -> bool:
return user_has_permission(user, "account:view_all")
+25 -13
View File
@@ -740,7 +740,7 @@ class DouyuBatchRunner:
@classmethod
def _supplier_order_status(cls, payload: dict) -> int | None:
"""提取供应商订单状态,文档约定 0-4。"""
return cls._to_int(cls._supplier_value(payload, "order_status", "orderStatus"))
return cls._to_int(cls._supplier_value(payload, "order_status", "orderStatus", "supplier_order_status"))
@classmethod
def _supplier_message(cls, payload: dict) -> str:
@@ -768,13 +768,20 @@ class DouyuBatchRunner:
result: dict,
) -> int | None:
"""轮询供应商直充订单至结束状态。"""
order_no = str(result["customer_order_no"])
order_no = str(result["out_order_id"])
deadline = time.monotonic() + DOUYU_PAYMENT_POLL_SECONDS
poll_count = 0
result["payment_polling"] = True
while not self._stop.is_set() and time.monotonic() <= deadline:
try:
payload = client.query_order(customer_order_no=order_no)
# 回调可能已在另一个数据库会话中结束订单,刷新后直接使用其结果。
db.refresh(task)
if task.status in {"success", "failed"}:
callback_result = task.result if isinstance(task.result, dict) else result
result.update(callback_result)
result["payment_polling"] = False
return self._supplier_order_status(callback_result)
payload = client.query_order(order_no)
code = self._to_int(self._supplier_value(payload, "code"))
status = self._supplier_order_status(payload)
poll_count += 1
@@ -2505,10 +2512,12 @@ class DouyuBatchRunner:
if not charge_account:
raise FishFinRechargeError("账号缺少斗鱼 UID,无法发起供应商直充")
# task.id 是唯一且稳定的商户单号来源,重试时不会创建不同的供应商订单。
# task.id 是唯一且稳定的外部订单号来源,重试时不会创建不同的供应商订单。
order_no = f"DYGF{task.id}"
# customer_price 是用户选择的充值面值;goodsFaceValue=0.993 是供货成本,不能作为支付金额。
customer_price = Decimal(amount)
task.supplier_out_order_id = order_no
db.commit()
# pay_amount 是用户选择的充值面值;goodsFaceValue=0.993 是供货成本,不能作为支付金额。
pay_amount = Decimal(amount)
def trace(event: dict) -> None:
"""将脱敏供应商协议信息输出到任务日志,便于线上联调。"""
stage = event.get("stage")
@@ -2517,8 +2526,8 @@ class DouyuBatchRunner:
self._push_log(
"info",
"供应商直充请求 "
f"path={event.get('path')} uid={params.get('charge_account')} "
f"buy_num={params.get('buy_num')} customer_price={params.get('customer_price')} "
f"path={event.get('path')} out_order_id={params.get('out_order_id')} "
f"buy_num={params.get('buy_num')} pay_amount={params.get('pay_amount')} "
f"product_id={params.get('product_id')} types={event.get('parameter_types')} "
f"sign_digest={event.get('sign_digest')}",
)
@@ -2550,22 +2559,25 @@ class DouyuBatchRunner:
client = FishFinRechargeClient(FishFinRechargeConfig.from_env(), trace=trace)
order_payload = client.create_order(
charge_account=charge_account,
buy_num=amount,
customer_price=customer_price,
customer_order_no=order_no,
pay_amount=pay_amount,
out_order_id=order_no,
product_id=product_id,
recharge_arg=[{"templateName": template_name, "templateVal": charge_account}],
order_type=0,
notify_url=client.config.notify_url,
)
code = self._to_int(self._supplier_value(order_payload, "code"))
status = self._supplier_order_status(order_payload)
result = {
"recharge_channel": "supplier_api",
"customer_order_no": order_no,
"out_order_id": order_no,
"order_id": self._supplier_value(order_payload, "order_id", "orderId"),
"charge_account": charge_account,
"buy_num": amount,
"product_id": product_id,
"customer_price": format(customer_price.normalize(), "f"),
"pay_amount": format(pay_amount.normalize(), "f"),
"order_type": 0,
"supplier_code": code,
"supplier_order_status": status,
"supplier_order": self._supplier_result(order_payload),