217 lines
7.1 KiB
Python
217 lines
7.1 KiB
Python
"""斗鱼活动任务服务。"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import uuid
|
|
from datetime import datetime, timezone
|
|
|
|
from sqlalchemy.orm import Session
|
|
|
|
from core.douyu.activity_client import DouyuActivityClient
|
|
from core.douyu.cookie_utils import cookie_value
|
|
|
|
from ..models import Account, DouyuConfig, DouyuTask, LoginTask
|
|
|
|
|
|
SUPPORTED_DOUYU_TASK_TYPES = {
|
|
"get_bind_qr": "获取绑定二维码",
|
|
"confirm_bind": "确认绑定",
|
|
"create_elite_qr": "开通精英宝典30",
|
|
"prepare_esports_bind": "绑定电竞手册角色",
|
|
"get_esports_bind_qr": "切换电竞角色",
|
|
"query_esports_game_name": "查询最新游戏角色",
|
|
"confirm_esports_bind": "确认电竞手册绑定",
|
|
"create_esports_qr": "开通电竞手册30",
|
|
"query_esports_points": "查询电竞积分",
|
|
"donate_esports_chicken_gift": "赠送冠军鸡腿",
|
|
"donate_esports_firework_gift": "赠送冠军烟花",
|
|
"refresh_esports_goods": "刷新电竞皮肤商城",
|
|
"exchange_esports_goods": "兑换电竞皮肤",
|
|
"create_gold_qr": "充值鱼翅",
|
|
"donate_elite_gift": "赠送精英令",
|
|
"query_points": "一键查询积分",
|
|
"exchange_goods": "兑换商品",
|
|
"query_game_name": "一键获取游戏名",
|
|
"query_change_bind_time": "一键查询换绑时间",
|
|
"query_limited_goods": "一键查询限兑商品",
|
|
"query_gold_balance": "一键查询鱼翅余额",
|
|
"refresh_goods": "刷新商品列表",
|
|
"query_exchange_records": "一键查询兑换记录",
|
|
"prefetch_csrf_token": "一键获取兑换 CSRF Token",
|
|
}
|
|
|
|
|
|
DOUYU_CONFIG_DEFAULTS = {
|
|
"manual_id": "G4KA4Qnz4LDp7",
|
|
"rid": "9263298",
|
|
"bind_act_alias": "20260120QYOOB",
|
|
"confirm_act_alias": "20260120QYOOB",
|
|
"legacy_act_alias": "cjm",
|
|
"room_id": "9263298",
|
|
"elite_amount": 3000,
|
|
"esports_manual_id": "1W0yGkOdDbAXb",
|
|
"esports_act_alias": "20260226ZMGDG",
|
|
"esports_amount": 3000,
|
|
"esports_chicken_gift_id": "24766",
|
|
"esports_chicken_skin_id": "0",
|
|
"esports_firework_gift_id": "24767",
|
|
"esports_firework_skin_id": "3850",
|
|
"gold_pay_type": 1,
|
|
"gift_id": "23643",
|
|
"skin_id": "2942",
|
|
}
|
|
|
|
DOUYU_CONFIG_FIELDS = tuple(DOUYU_CONFIG_DEFAULTS.keys())
|
|
DOUYU_ACTIVE_TASK_STATUSES = ("planned", "pending", "running")
|
|
|
|
|
|
def douyu_config_value(field: str, value):
|
|
"""读取配置值;空值自动回退到默认值。"""
|
|
default = DOUYU_CONFIG_DEFAULTS[field]
|
|
if isinstance(default, int):
|
|
try:
|
|
return int(value if value is not None else default)
|
|
except (TypeError, ValueError):
|
|
return default
|
|
text = str(value or "").strip()
|
|
return text or str(default)
|
|
|
|
|
|
def apply_douyu_config_defaults(config: DouyuConfig) -> bool:
|
|
"""补齐斗鱼配置默认值,返回是否发生变更。"""
|
|
changed = False
|
|
for field in DOUYU_CONFIG_FIELDS:
|
|
if field == "bind_act_alias" and str(getattr(config, field, "") or "").strip() == "20250213NQCYX":
|
|
setattr(config, field, DOUYU_CONFIG_DEFAULTS[field])
|
|
changed = True
|
|
continue
|
|
normalized = douyu_config_value(field, getattr(config, field, None))
|
|
if getattr(config, field, None) != normalized:
|
|
setattr(config, field, normalized)
|
|
changed = True
|
|
return changed
|
|
|
|
|
|
def ensure_douyu_config(db: Session) -> DouyuConfig:
|
|
"""获取单条斗鱼配置,不存在则创建。"""
|
|
config = db.query(DouyuConfig).first()
|
|
if config:
|
|
if apply_douyu_config_defaults(config):
|
|
db.commit()
|
|
db.refresh(config)
|
|
return config
|
|
config = DouyuConfig(**DOUYU_CONFIG_DEFAULTS)
|
|
db.add(config)
|
|
db.commit()
|
|
db.refresh(config)
|
|
return config
|
|
|
|
|
|
def latest_success_cookie(db: Session, account_id: int) -> str:
|
|
"""读取账号最近一次成功登录 Cookie。"""
|
|
task = (
|
|
db.query(LoginTask)
|
|
.filter(
|
|
LoginTask.account_id == account_id,
|
|
LoginTask.status == "success",
|
|
LoginTask.cookie != "",
|
|
)
|
|
.order_by(LoginTask.id.desc())
|
|
.first()
|
|
)
|
|
return task.cookie if task else ""
|
|
|
|
|
|
def cookie_account_ids_query(db: Session):
|
|
"""返回拥有成功 Cookie 的斗鱼账号 ID 查询。"""
|
|
return (
|
|
db.query(LoginTask.account_id)
|
|
.filter(LoginTask.status == "success", LoginTask.cookie != "")
|
|
.distinct()
|
|
)
|
|
|
|
|
|
def visible_douyu_task_accounts(db: Session, account_ids: list[int]) -> list[Account]:
|
|
"""只保留存在成功 Cookie 的斗鱼账号。"""
|
|
if not account_ids:
|
|
return []
|
|
cookie_ids = cookie_account_ids_query(db).subquery()
|
|
return (
|
|
db.query(Account)
|
|
.filter(Account.id.in_(account_ids), Account.id.in_(cookie_ids))
|
|
.all()
|
|
)
|
|
|
|
|
|
def update_account_profile_from_cookie(account: Account, cookie: str) -> None:
|
|
"""从 Cookie 回填 uid/nickname。"""
|
|
profile = DouyuActivityClient.profile_from_cookie(cookie)
|
|
if profile["uid"]:
|
|
account.uid = profile["uid"]
|
|
if profile["nickname"]:
|
|
account.nickname = profile["nickname"]
|
|
|
|
|
|
def account_uid(account: Account, cookie: str) -> str:
|
|
"""优先从账号字段读取 uid,缺失时从 Cookie 中取。"""
|
|
return account.uid or cookie_value(cookie, "acf_uid")
|
|
|
|
|
|
def create_douyu_planned_tasks(
|
|
db: Session,
|
|
account_ids: list[int],
|
|
task_type: str,
|
|
created_by: int,
|
|
payload: dict | None = None,
|
|
) -> tuple[str, int]:
|
|
"""创建斗鱼任务记录,等待后台执行器消费。"""
|
|
if task_type not in SUPPORTED_DOUYU_TASK_TYPES:
|
|
raise ValueError("不支持的任务类型")
|
|
|
|
accounts = visible_douyu_task_accounts(db, account_ids)
|
|
if task_type in {"refresh_goods", "refresh_esports_goods"} and accounts:
|
|
# 商品快照是全局数据,一个可用 CK 足够;没有 CK 时前端无法选账号创建任务。
|
|
accounts = accounts[:1]
|
|
|
|
batch_id = uuid.uuid4().hex[:12]
|
|
payload = payload or {}
|
|
for account in accounts:
|
|
db.add(DouyuTask(
|
|
batch_id=batch_id,
|
|
account_id=account.id,
|
|
task_type=task_type,
|
|
status="planned",
|
|
message="任务已创建,等待执行",
|
|
result={"payload": payload} if payload else None,
|
|
created_by=created_by,
|
|
))
|
|
db.commit()
|
|
return batch_id, len(accounts)
|
|
|
|
|
|
def cleanup_orphan_douyu_tasks(
|
|
db: Session,
|
|
*,
|
|
active_batch_ids: set[str] | None = None,
|
|
batch_id: str | None = None,
|
|
statuses: tuple[str, ...] = DOUYU_ACTIVE_TASK_STATUSES,
|
|
message: str = "任务已中断(服务重启或批次丢失)",
|
|
) -> int:
|
|
"""清理没有执行器接管的斗鱼任务。"""
|
|
query = db.query(DouyuTask).filter(DouyuTask.status.in_(statuses))
|
|
if batch_id:
|
|
query = query.filter(DouyuTask.batch_id == batch_id)
|
|
elif active_batch_ids is not None and active_batch_ids:
|
|
query = query.filter(~DouyuTask.batch_id.in_(list(active_batch_ids)))
|
|
|
|
tasks = query.all()
|
|
if not tasks:
|
|
return 0
|
|
now = datetime.now(timezone.utc)
|
|
for task in tasks:
|
|
task.status = "stopped"
|
|
task.message = message
|
|
task.finished_at = now
|
|
db.commit()
|
|
return len(tasks)
|