Files
2026-08-31 10:55:44 +08:00

382 lines
15 KiB
Python
Raw Permalink 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.
"""斗鱼任务执行器:手册开通与积分(由 douyu_runner.py 按功能域拆分)。"""
from __future__ import annotations
import time
from datetime import UTC, datetime
from typing import TYPE_CHECKING, Any
from sqlalchemy.orm import Session
from core.douyu import DouyuActivityClient, DouyuActivityError
from ..models import Account, DouyuTask
from .douyu_runner_core import DOUYU_PAYMENT_POLL_INTERVAL, DOUYU_PAYMENT_POLL_SECONDS
from .douyu_service import account_uid, update_account_profile_from_cookie
class ManualMixin:
"""手册域:精英/电竞手册开通支付、积分查询与到账轮询。"""
if TYPE_CHECKING:
_stop: Any
@staticmethod
def _to_int(value: Any) -> int | None: ...
def _client(self, cookie: str) -> DouyuActivityClient: ...
def _task_payload(self, task: DouyuTask) -> dict: ...
def _push_log(self, level: str, message: str) -> None: ...
def _mark_task(self, *args: Any, **kwargs: Any) -> None: ...
def _update_task_progress(self, *args: Any, **kwargs: Any) -> None: ...
@classmethod
def _confirm_act_alias(cls, config: dict) -> str: ...
@classmethod
def _bind_qr_act_alias(cls, config: dict) -> str: ...
def _refresh_account_points(
self,
client: DouyuActivityClient,
account: Account,
cookie: str,
*,
ctn: str | None = None,
) -> dict:
"""刷新账号积分并写回账号表。"""
uid = account_uid(account, cookie)
if not uid:
raise DouyuActivityError("Cookie 中没有 acf_uid,无法查询积分")
ctn_value = ctn or client.acf_ccn(refresh_subscribe=False)
result = client.query_points(uid=uid, ctn=ctn_value)
points = self._to_int(result.get("points"))
account.uid = uid
account.points = points
update_account_profile_from_cookie(account, cookie)
account.updated_at = datetime.now(UTC)
return {"points": points, "points_query": result}
def _wait_points_after_payment(
self,
db: Session,
task: DouyuTask,
account: Account,
client: DouyuActivityClient,
cookie: str,
ctn: str,
result: dict,
) -> bool:
"""等待宝典支付到账;积分达到 300 视为开通成功。"""
deadline = time.monotonic() + DOUYU_PAYMENT_POLL_SECONDS
result["payment_polling"] = True
result["payment_target_points"] = 300
poll_count = 0
last_points = None
while not self._stop.is_set() and time.monotonic() <= deadline:
try:
points_result = self._refresh_account_points(
client, account, cookie, ctn=ctn
)
db.commit()
poll_count += 1
last_points = points_result["points"]
result.update(points_result)
result["payment_poll_count"] = poll_count
result["payment_polling"] = True
if last_points is not None and last_points >= 300:
result["payment_polling"] = False
result["elite_opened"] = True
return True
self._update_task_progress(
db,
task,
"running",
f"精英宝典支付码已生成,等待开通到账(当前积分 {last_points if last_points is not None else '-'}",
result,
)
except Exception as exc: # noqa: BLE001 外部接口与任务边界需要保留宽泛异常兜底
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["elite_opened"] = False
result["points"] = last_points
return False
def _refresh_esports_handbook(
self,
client: DouyuActivityClient,
account: Account,
*,
manual_id: str,
) -> dict:
"""刷新电竞手册开通状态并将积分写回账号。"""
result = client.esports_user_info(manual_id=manual_id)
manual_type = self._to_int(result.get("manual_type"))
manual_score = self._to_int(result.get("manual_score"))
account.esports_points = manual_score
account.updated_at = datetime.now(UTC)
return {
"esports_manual_type": manual_type,
"esports_manual_score": manual_score,
"esports_expire_time": result.get("expire_time"),
"esports_user_info": result,
"esports_points": manual_score,
"points": manual_score,
}
def _wait_esports_open_after_payment(
self,
db: Session,
task: DouyuTask,
account: Account,
client: DouyuActivityClient,
*,
manual_id: str,
result: dict,
baseline_manual_type: int | None,
baseline_manual_score: int | None,
) -> bool:
"""等待电竞手册支付到账,以 manualType=1 或积分变化作为成功条件。"""
deadline = time.monotonic() + DOUYU_PAYMENT_POLL_SECONDS
result["payment_polling"] = True
result["esports_manual_type_baseline"] = baseline_manual_type
result["esports_manual_score_baseline"] = baseline_manual_score
poll_count = 0
last_manual_type = baseline_manual_type
last_manual_score = baseline_manual_score
while not self._stop.is_set() and time.monotonic() <= deadline:
try:
handbook_result = self._refresh_esports_handbook(
client,
account,
manual_id=manual_id,
)
db.commit()
poll_count += 1
last_manual_type = handbook_result["esports_manual_type"]
last_manual_score = handbook_result["esports_manual_score"]
result.update(handbook_result)
result["payment_poll_count"] = poll_count
result["payment_polling"] = True
opened = (last_manual_type is not None and last_manual_type >= 1) or (
baseline_manual_score is not None
and last_manual_score is not None
and last_manual_score > baseline_manual_score
)
if opened:
result["payment_polling"] = False
result["esports_opened"] = True
return True
self._update_task_progress(
db,
task,
"running",
"电竞手册支付码已生成,等待开通到账"
f"(类型 {last_manual_type if last_manual_type is not None else '-'}"
f"积分 {last_manual_score if last_manual_score is not None else '-'}",
result,
)
except Exception as exc: # noqa: BLE001 外部接口与任务边界需要保留宽泛异常兜底
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["esports_opened"] = False
result["esports_manual_type"] = last_manual_type
result["esports_manual_score"] = last_manual_score
result["esports_points"] = last_manual_score
result["points"] = last_manual_score
return False
def _execute_create_elite_qr(
self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict
):
client = self._client(cookie)
ctn = str(self._task_payload(task).get("ctn") or "")
if not ctn:
ctn = client.acf_ccn(refresh_subscribe=True)
act_alias = self._confirm_act_alias(config) or self._bind_qr_act_alias(config)
if not act_alias:
self._mark_task(db, task, "failed", "请先配置开通宝典活动 actAlias")
return
result = client.create_elite_qr(
ctn=ctn,
act_alias=act_alias,
amount=int(config["elite_amount"]),
room_id=str(config["room_id"]),
)
account.bind_status = "elite_qr_created"
account.updated_at = datetime.now(UTC)
self._update_task_progress(
db, task, "running", "精英宝典支付码已生成,等待开通到账", result
)
opened = self._wait_points_after_payment(
db, task, account, client, cookie, ctn, result
)
if self._stop.is_set():
self._mark_task(db, task, "stopped", "任务已停止", result)
return
if opened:
account.bind_status = "elite_opened"
account.updated_at = datetime.now(UTC)
self._mark_task(
db, task, "success", f"精英宝典已开通,积分: {account.points}", result
)
return
self._mark_task(
db,
task,
"failed",
f"未检测到精英宝典开通到账,当前积分: {account.points if account.points is not None else '-'}",
result,
)
def _execute_create_esports_qr(
self,
db: Session,
task: DouyuTask,
account: Account,
cookie: str,
config: dict,
):
"""生成电竞手册支付二维码,并通过活动状态确认开通到账。"""
client = self._client(cookie)
ctn = str(self._task_payload(task).get("ctn") or "")
if not ctn:
ctn = client.acf_ccn(refresh_subscribe=True)
act_alias = str(config.get("esports_act_alias") or "").strip()
manual_id = str(config.get("esports_manual_id") or "").strip()
if not act_alias or not manual_id:
self._mark_task(
db, task, "failed", "请先配置电竞手册活动 actAlias 和 manualID"
)
return
baseline_manual_type = None
baseline_manual_score = None
baseline_result: dict = {}
try:
baseline_result = self._refresh_esports_handbook(
client, account, manual_id=manual_id
)
baseline_manual_type = baseline_result["esports_manual_type"]
baseline_manual_score = baseline_result["esports_manual_score"]
db.commit()
except Exception as exc: # noqa: BLE001 外部接口与任务边界需要保留宽泛异常兜底
self._push_log("warning", f"生成电竞手册支付码前查询活动状态失败: {exc}")
if baseline_manual_type is not None and baseline_manual_type >= 1:
account.esports_bind_status = "esports_opened"
account.updated_at = datetime.now(UTC)
self._mark_task(
db,
task,
"success",
f"电竞手册已开通,积分: {baseline_manual_score if baseline_manual_score is not None else '-'}",
{**baseline_result, "esports_opened": True, "payment_polling": False},
)
return
result = client.create_esports_qr(
ctn=ctn,
act_alias=act_alias,
amount=int(config["esports_amount"]),
room_id=str(config["room_id"]),
)
result.update(baseline_result)
account.esports_bind_status = "esports_qr_created"
account.updated_at = datetime.now(UTC)
self._update_task_progress(
db, task, "running", "电竞手册支付码已生成,等待开通到账", result
)
opened = self._wait_esports_open_after_payment(
db,
task,
account,
client,
manual_id=manual_id,
result=result,
baseline_manual_type=baseline_manual_type,
baseline_manual_score=baseline_manual_score,
)
if self._stop.is_set():
self._mark_task(db, task, "stopped", "任务已停止", result)
return
if opened:
account.esports_bind_status = "esports_opened"
account.updated_at = datetime.now(UTC)
self._mark_task(
db,
task,
"success",
f"电竞手册已开通,积分: {account.esports_points if account.esports_points is not None else '-'}",
result,
)
return
self._mark_task(
db,
task,
"failed",
"未检测到电竞手册开通到账"
f",类型: {result.get('esports_manual_type', '-')}"
f"积分: {result.get('esports_manual_score', '-')}",
result,
)
def _execute_query_esports_points(
self,
db: Session,
task: DouyuTask,
account: Account,
cookie: str,
config: dict,
):
"""查询电竞手册积分。"""
manual_id = str(config.get("esports_manual_id") or "").strip()
if not manual_id:
self._mark_task(db, task, "failed", "请先配置电竞手册 manualID")
return
client = self._client(cookie)
result = self._refresh_esports_handbook(client, account, manual_id=manual_id)
account.esports_bind_status = "esports_points_queried"
account.updated_at = datetime.now(UTC)
points = result["esports_points"]
self._mark_task(
db,
task,
"success",
f"电竞积分: {points if points is not None else '-'}",
result,
)
def _execute_query_points(
self, db: Session, task: DouyuTask, account: Account, cookie: str, config: dict
):
client = self._client(cookie)
ctn = client.acf_ccn(refresh_subscribe=False)
result = self._refresh_account_points(client, account, cookie, ctn=ctn)
account.bind_status = "points_queried"
account.updated_at = datetime.now(UTC)
points = result["points"]
self._mark_task(
db,
task,
"success",
f"积分: {points if points is not None else '-'}",
result,
)