- douyu_runner_manual/gold: 导入 core 的 DOUYU_PAYMENT_POLL_* 常量
(之前仅在 douyu_runner.py 单文件内可见, 拆分后到账轮询会 NameError)
- 新增 tests/test_douyu_runner_proxy.py: 静态/API/关闭三模式 + 降级直连 +
activity_client session 级 proxies, 兑现 8fe9fdf 提交说明中的行为断言
367 lines
17 KiB
Python
367 lines
17 KiB
Python
"""斗鱼任务执行器:鱼翅充值(由 douyu_runner.py 按功能域拆分)。"""
|
||
|
||
from __future__ import annotations
|
||
import re
|
||
import time
|
||
from decimal import Decimal
|
||
from datetime import datetime, timezone
|
||
from sqlalchemy.orm import Session
|
||
|
||
from core.douyu import DouyuActivityClient, FishFinRechargeClient, FishFinRechargeConfig, FishFinRechargeError
|
||
from ..models import Account, DouyuTask
|
||
from .douyu_service import update_account_profile_from_cookie
|
||
from .douyu_runner_core import DOUYU_PAYMENT_POLL_INTERVAL, DOUYU_PAYMENT_POLL_SECONDS
|
||
|
||
class GoldMixin:
|
||
"""鱼翅充值域:扫码充值、供应商直充与到账轮询。"""
|
||
def _refresh_account_gold_balance(self, client: DouyuActivityClient, account: Account) -> dict:
|
||
"""刷新鱼翅和钱包兑换余额并写回账号表。"""
|
||
gold = client.gold_account()
|
||
exchange = client.exchange_balance()
|
||
account.gold_balance = self._to_int(gold.get("gold"))
|
||
account.exchange_balance = self._to_int(exchange.get("count"))
|
||
account.updated_at = datetime.now(timezone.utc)
|
||
return {
|
||
"gold_balance": account.gold_balance,
|
||
"exchange_balance": account.exchange_balance,
|
||
"gold": gold,
|
||
"exchange_balance_query": exchange,
|
||
}
|
||
|
||
def _wait_gold_balance_after_payment(
|
||
self,
|
||
db: Session,
|
||
task: DouyuTask,
|
||
account: Account,
|
||
client: DouyuActivityClient,
|
||
result: dict,
|
||
baseline_gold: int | None,
|
||
) -> bool:
|
||
"""等待鱼翅充值到账;余额变化后写回账号表。"""
|
||
deadline = time.monotonic() + DOUYU_PAYMENT_POLL_SECONDS
|
||
result["payment_polling"] = True
|
||
result["baseline_gold_balance"] = baseline_gold
|
||
poll_count = 0
|
||
last_gold = baseline_gold
|
||
baseline_ready = baseline_gold is not None
|
||
while not self._stop.is_set() and time.monotonic() <= deadline:
|
||
try:
|
||
balance_result = self._refresh_account_gold_balance(client, account)
|
||
db.commit()
|
||
poll_count += 1
|
||
last_gold = balance_result["gold_balance"]
|
||
result.update(balance_result)
|
||
result["payment_poll_count"] = poll_count
|
||
result["payment_polling"] = True
|
||
if not baseline_ready and last_gold is not None:
|
||
baseline_gold = last_gold
|
||
result["baseline_gold_balance"] = baseline_gold
|
||
baseline_ready = True
|
||
self._update_task_progress(
|
||
db,
|
||
task,
|
||
"running",
|
||
f"鱼翅支付码已生成,已记录当前余额 {last_gold},等待到账",
|
||
result,
|
||
)
|
||
if self._stop.wait(DOUYU_PAYMENT_POLL_INTERVAL):
|
||
break
|
||
continue
|
||
changed = last_gold is not None and (baseline_gold is None or last_gold != baseline_gold)
|
||
if changed:
|
||
result["payment_polling"] = False
|
||
result["gold_recharged"] = True
|
||
return True
|
||
self._update_task_progress(
|
||
db,
|
||
task,
|
||
"running",
|
||
f"鱼翅支付码已生成,等待到账(当前鱼翅 {last_gold if last_gold is not None else '-'})",
|
||
result,
|
||
)
|
||
except Exception as exc:
|
||
poll_count += 1
|
||
result["payment_poll_count"] = poll_count
|
||
result["payment_poll_error"] = str(exc)
|
||
self._update_task_progress(db, task, "running", f"等待鱼翅到账: {exc}", result)
|
||
if self._stop.wait(DOUYU_PAYMENT_POLL_INTERVAL):
|
||
break
|
||
result["payment_polling"] = False
|
||
result["gold_recharged"] = False
|
||
result["gold_balance"] = last_gold
|
||
return False
|
||
|
||
@staticmethod
|
||
def _supplier_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
|
||
|
||
@classmethod
|
||
def _supplier_order_status(cls, payload: dict) -> int | None:
|
||
"""提取供应商订单状态,文档约定 0-4。"""
|
||
return cls._to_int(cls._supplier_value(payload, "order_status", "orderStatus", "supplier_order_status"))
|
||
|
||
@classmethod
|
||
def _supplier_message(cls, payload: dict) -> str:
|
||
"""提取供应商可展示的业务消息。"""
|
||
value = cls._supplier_value(payload, "msg", "message", "error_msg")
|
||
return str(value or "")[:256]
|
||
|
||
@staticmethod
|
||
def _supplier_result(payload: dict) -> dict:
|
||
"""保存必要订单状态,避免把完整供应商响应或签名暴露到任务结果。"""
|
||
data = payload.get("data") if isinstance(payload.get("data"), dict) else {}
|
||
response_result = payload.get("result") if isinstance(payload.get("result"), dict) else {}
|
||
result = {
|
||
key: value
|
||
for key, value in {**payload, **data, **response_result}.items()
|
||
if key not in {"sign", "cards", "card_no", "card_pwd", "recharge_arg"}
|
||
}
|
||
return result
|
||
|
||
@staticmethod
|
||
def _supplier_out_order_id(task: DouyuTask) -> str:
|
||
"""生成可追踪的供应商外部订单号;已有订单号必须在重试时复用。"""
|
||
existing = str(task.supplier_out_order_id or "").strip()
|
||
if existing:
|
||
return existing
|
||
batch_token = re.sub(r"[^A-Za-z0-9]", "", str(task.batch_id or "")).upper()[:16] or "LOCAL"
|
||
return f"DYGF{batch_token}T{task.id}"
|
||
|
||
def _wait_supplier_gold_order(
|
||
self,
|
||
db: Session,
|
||
task: DouyuTask,
|
||
client: FishFinRechargeClient,
|
||
result: dict,
|
||
) -> int | None:
|
||
"""轮询供应商直充订单至结束状态。"""
|
||
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:
|
||
# 回调可能已在另一个数据库会话中结束订单,刷新后直接使用其结果。
|
||
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
|
||
result.update({
|
||
"payment_poll_count": poll_count,
|
||
"supplier_code": code,
|
||
"supplier_order_status": status,
|
||
"supplier_order": self._supplier_result(payload),
|
||
})
|
||
if code != 200:
|
||
result["payment_polling"] = False
|
||
return status if status in {2, 3, 4} else 4
|
||
if status in {2, 3, 4}:
|
||
result["payment_polling"] = False
|
||
return status
|
||
self._update_task_progress(
|
||
db,
|
||
task,
|
||
"running",
|
||
f"供应商直充订单处理中(状态 {status if status is not None else '-'})",
|
||
result,
|
||
)
|
||
except FishFinRechargeError as exc:
|
||
poll_count += 1
|
||
result["payment_poll_count"] = poll_count
|
||
result["payment_poll_error"] = str(exc)
|
||
self._update_task_progress(db, task, "running", f"查询供应商订单失败: {exc}", result)
|
||
if self._stop.wait(DOUYU_PAYMENT_POLL_INTERVAL):
|
||
break
|
||
result["payment_polling"] = False
|
||
return None
|
||
|
||
def _execute_create_gold_qr(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
|
||
payload = self._task_payload(task)
|
||
amount = int(payload.get("amount") or payload.get("gold_amount") or 1)
|
||
channel = str(config.get("gold_recharge_channel") or "wechat_qr")
|
||
if channel == "supplier_api":
|
||
try:
|
||
self._execute_create_gold_supplier_order(db, task, account, cookie, config, amount)
|
||
except FishFinRechargeError as exc:
|
||
self._mark_task(db, task, "failed", str(exc), {"recharge_channel": "supplier_api"})
|
||
return
|
||
client = self._client(cookie)
|
||
baseline_gold = account.gold_balance
|
||
try:
|
||
baseline = self._refresh_account_gold_balance(client, account)
|
||
baseline_gold = baseline["gold_balance"]
|
||
db.commit()
|
||
except Exception as exc:
|
||
self._push_log("warning", f"生成鱼翅码前刷新余额失败: {exc}")
|
||
result = client.create_gold_qr(amount=amount, pay_type=int(config["gold_pay_type"]))
|
||
account.bind_status = "gold_qr_created"
|
||
account.updated_at = datetime.now(timezone.utc)
|
||
self._update_task_progress(db, task, "running", f"鱼翅 {amount} 元支付码已生成,等待到账", result)
|
||
recharged = self._wait_gold_balance_after_payment(db, task, account, client, result, baseline_gold)
|
||
if self._stop.is_set():
|
||
self._mark_task(db, task, "stopped", "任务已停止", result)
|
||
return
|
||
if recharged:
|
||
account.bind_status = "gold_recharged"
|
||
account.updated_at = datetime.now(timezone.utc)
|
||
self._mark_task(
|
||
db,
|
||
task,
|
||
"success",
|
||
f"鱼翅已到账,当前余额: {account.gold_balance if account.gold_balance is not None else '-'}",
|
||
result,
|
||
)
|
||
return
|
||
self._mark_task(
|
||
db,
|
||
task,
|
||
"failed",
|
||
f"未检测到鱼翅到账,当前余额: {account.gold_balance if account.gold_balance is not None else '-'}",
|
||
result,
|
||
)
|
||
|
||
def _execute_create_gold_supplier_order(
|
||
self,
|
||
db: Session,
|
||
task: DouyuTask,
|
||
account: Account,
|
||
cookie: str,
|
||
config: dict,
|
||
amount: int,
|
||
) -> None:
|
||
"""创建供应商鱼翅直充订单并轮询订单状态。"""
|
||
product_id = str(config.get("gold_api_product_id") or "").strip()
|
||
template_name = str(config.get("gold_api_account_template_name") or "斗鱼昵称").strip()
|
||
if not product_id:
|
||
raise FishFinRechargeError("请先在配置中填写供应商直充商品 ID")
|
||
# 充值商品按斗鱼昵称识别账号,UID 只能作为审计信息,不能作为充值值。
|
||
update_account_profile_from_cookie(account, cookie)
|
||
recharge_account = str(account.nickname or "").strip()
|
||
if not recharge_account:
|
||
raise FishFinRechargeError("账号缺少斗鱼昵称,无法发起供应商直充")
|
||
|
||
# 首次生成后持久化,网络重试或进程重启都继续查询同一笔订单。
|
||
order_no = self._supplier_out_order_id(task)
|
||
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")
|
||
if stage == "request":
|
||
params = event.get("params") or {}
|
||
self._push_log(
|
||
"info",
|
||
"供应商直充 | 下单 "
|
||
f"| 外部单号={params.get('out_order_id') or '-'} "
|
||
f"| 数量={params.get('buy_num') or '-'} "
|
||
f"| 金额={params.get('pay_amount') or '-'} "
|
||
f"| 商品={params.get('product_id') or '-'}",
|
||
)
|
||
if event.get("json_body"):
|
||
self._push_log(
|
||
"debug",
|
||
"供应商协议 | 请求 "
|
||
f"| {event.get('method')} {event.get('path')} "
|
||
f"| 签名摘要={event.get('sign_digest')} "
|
||
f"| 参数={FishFinRechargeClient._json_text(event['json_body'])}",
|
||
)
|
||
elif stage == "response":
|
||
status = self._to_int(event.get("order_status"))
|
||
status_labels = {0: "待处理", 1: "处理中", 2: "成功", 3: "失败", 4: "异常"}
|
||
status_text = status_labels.get(status, "-")
|
||
reason = str(event.get("fail_reason") or event.get("message") or "-")
|
||
self._push_log(
|
||
"info",
|
||
"供应商直充 | 响应 "
|
||
f"| HTTP={event.get('http_status') or '-'} "
|
||
f"| 业务码={event.get('code') or '-'} "
|
||
f"| 外部单号={event.get('out_order_id') or '-'} "
|
||
f"| 供应商单号={event.get('order_id') or '-'} "
|
||
f"| 状态={status_text} "
|
||
f"| 提示={reason}",
|
||
)
|
||
if event.get("response_body"):
|
||
self._push_log(
|
||
"debug",
|
||
"供应商协议 | 响应 "
|
||
f"| HTTP={event.get('http_status')} "
|
||
f"| 内容={FishFinRechargeClient._json_text(event['response_body'])}",
|
||
)
|
||
|
||
client = FishFinRechargeClient(FishFinRechargeConfig.from_env(), trace=trace)
|
||
order_payload = client.create_order(
|
||
buy_num=amount,
|
||
pay_amount=pay_amount,
|
||
out_order_id=order_no,
|
||
product_id=product_id,
|
||
recharge_arg=[{"templateName": template_name, "templateVal": recharge_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",
|
||
"out_order_id": order_no,
|
||
"order_id": self._supplier_value(order_payload, "order_id", "orderId"),
|
||
"recharge_account": recharge_account,
|
||
"douyu_uid": str(account.uid or "").strip(),
|
||
"buy_num": amount,
|
||
"product_id": product_id,
|
||
"pay_amount": format(pay_amount.normalize(), "f"),
|
||
"order_type": 0,
|
||
"supplier_code": code,
|
||
"supplier_order_status": status,
|
||
"supplier_order": self._supplier_result(order_payload),
|
||
}
|
||
if code != 200:
|
||
self._mark_task(db, task, "failed", self._supplier_message(order_payload) or "供应商创建直充订单失败", result)
|
||
return
|
||
account.bind_status = "gold_api_order_created"
|
||
account.updated_at = datetime.now(timezone.utc)
|
||
self._update_task_progress(db, task, "running", "供应商直充订单已创建,等待到账", result)
|
||
if status not in {2, 3, 4}:
|
||
status = self._wait_supplier_gold_order(db, task, client, result)
|
||
if self._stop.is_set():
|
||
self._mark_task(db, task, "stopped", "任务已停止", result)
|
||
return
|
||
if status == 2:
|
||
account.bind_status = "gold_recharged"
|
||
account.updated_at = datetime.now(timezone.utc)
|
||
self._mark_task(db, task, "success", "供应商直充成功", result)
|
||
return
|
||
if status in {3, 4}:
|
||
self._mark_task(db, task, "failed", "供应商直充失败", result)
|
||
return
|
||
self._mark_task(db, task, "failed", "供应商直充订单查询超时", result)
|
||
|
||
def _execute_query_gold_balance(self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict):
|
||
client = self._client(cookie)
|
||
result = self._refresh_account_gold_balance(client, account)
|
||
account.bind_status = "gold_balance_queried"
|
||
account.updated_at = datetime.now(timezone.utc)
|
||
self._mark_task(
|
||
db,
|
||
task,
|
||
"success",
|
||
f"鱼翅余额: {account.gold_balance if account.gold_balance is not None else '-'}",
|
||
result,
|
||
)
|
||
|