- 碎片: 新增 query_xpd_fragments 任务, 复用 e语言 gamecoin 接口取扭蛋币 jb2, 账号已存角色回退(规避 Livelink 风控), accounts.xpd_fragments 字段与迁移 - 列表: 和平小店移除鱼翅/限制兑换/换绑时间列, 新增扭蛋碎片列; 操作区硬编码白底改用 theme token, 深色模式适配(三个手册共用) - 轮询: 任务状态经 WS(level=task)即时推送(二维码/绑定进度即时展示), pending/running/终态全覆盖, REST/WS 载荷统一 douyu_task_payload, raw 快照递归剥离, HTTP 轮询降频 5s/30s 兜底
447 lines
17 KiB
Python
447 lines
17 KiB
Python
"""斗鱼活动任务路由。"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import threading
|
|
from datetime import datetime, timezone
|
|
|
|
from fastapi import APIRouter, Depends, HTTPException, Query, WebSocket, WebSocketDisconnect
|
|
from sqlalchemy import or_
|
|
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, DouyuXpdGoodsSnapshot, User
|
|
from ..permissions import user_has_permission
|
|
from ..schemas import (
|
|
DouyuConfigOut,
|
|
DouyuConfigUpdate,
|
|
DouyuGoodsOut,
|
|
DouyuTaskAccountOut,
|
|
DouyuTaskBatchRequest,
|
|
DouyuTaskOut,
|
|
DouyuXpdGoodsOut,
|
|
)
|
|
from ..services.douyu_runner import DouyuBatchRunner, douyu_batch_registry
|
|
from ..services.douyu_service import (
|
|
DOUYU_CONFIG_FIELDS,
|
|
SUPPORTED_DOUYU_TASK_TYPES,
|
|
apply_douyu_config_defaults,
|
|
cleanup_orphan_douyu_tasks,
|
|
cookie_account_ids_query,
|
|
create_douyu_planned_tasks,
|
|
douyu_config_value,
|
|
douyu_task_payload,
|
|
ensure_douyu_config,
|
|
)
|
|
|
|
|
|
router = APIRouter(prefix="/api/douyu", tags=["斗鱼活动"])
|
|
|
|
|
|
def _can_view_all(user: User) -> bool:
|
|
return user_has_permission(user, "account:view_all")
|
|
|
|
|
|
def _visible_task_accounts_query(db: Session, current: User):
|
|
"""返回当前用户可用于斗鱼任务的账号查询。"""
|
|
cookie_ids = cookie_account_ids_query(db).subquery()
|
|
query = (
|
|
db.query(Account)
|
|
.options(joinedload(Account.assigned_user))
|
|
.filter(Account.id.in_(cookie_ids))
|
|
)
|
|
if _can_view_all(current):
|
|
return query
|
|
if user_has_permission(current, "account:view_assigned"):
|
|
return query.filter(Account.assigned_to == current.id)
|
|
raise HTTPException(status_code=403, detail="无权查看斗鱼账号")
|
|
|
|
|
|
def _visible_tasks_query(db: Session, current: User):
|
|
"""返回当前用户可查看的斗鱼任务查询。"""
|
|
query = db.query(DouyuTask).options(joinedload(DouyuTask.account))
|
|
if _can_view_all(current):
|
|
return query
|
|
if user_has_permission(current, "account:view_assigned"):
|
|
return query.join(DouyuTask.account).filter(Account.assigned_to == current.id)
|
|
raise HTTPException(status_code=403, detail="无权查看斗鱼任务")
|
|
|
|
|
|
def _require_task_account_access(db: Session, current: User, account_ids: list[int]) -> None:
|
|
"""确保任务只会提交到当前用户可操作的账号。"""
|
|
requested_ids = set(account_ids)
|
|
query = db.query(Account.id).filter(Account.id.in_(requested_ids))
|
|
if not _can_view_all(current):
|
|
if not user_has_permission(current, "account:view_assigned"):
|
|
raise HTTPException(status_code=403, detail="无权操作斗鱼账号")
|
|
query = query.filter(Account.assigned_to == current.id)
|
|
allowed_ids = {account_id for account_id, in query.all()}
|
|
if allowed_ids != requested_ids:
|
|
raise HTTPException(status_code=403, detail="包含无权操作的斗鱼账号")
|
|
|
|
|
|
def _require_batch_owner(db: Session, current: User, batch_id: str) -> None:
|
|
"""客服只能停止或订阅自己创建的任务批次。"""
|
|
if _can_view_all(current):
|
|
return
|
|
exists = (
|
|
db.query(DouyuTask.id)
|
|
.filter(DouyuTask.batch_id == batch_id, DouyuTask.created_by == current.id)
|
|
.first()
|
|
)
|
|
if not exists:
|
|
raise HTTPException(status_code=403, detail="无权操作该斗鱼任务批次")
|
|
|
|
|
|
def _account_out(account: Account) -> DouyuTaskAccountOut:
|
|
return DouyuTaskAccountOut(
|
|
id=account.id,
|
|
username=account.username,
|
|
uid=account.uid or "",
|
|
nickname=account.nickname or "",
|
|
tag=account.tag or "",
|
|
points=account.points,
|
|
game_name=account.game_name or "",
|
|
game_channel=account.game_channel or "",
|
|
gold_balance=account.gold_balance,
|
|
exchange_balance=account.exchange_balance,
|
|
bind_status=account.bind_status or "",
|
|
change_role_wait_time=account.change_role_wait_time,
|
|
esports_points=account.esports_points,
|
|
esports_game_name=account.esports_game_name or "",
|
|
esports_game_channel=account.esports_game_channel or "",
|
|
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_fragments=account.xpd_fragments,
|
|
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,
|
|
)
|
|
|
|
|
|
def _task_out(task: DouyuTask, *, include_detail: bool = False) -> DouyuTaskOut:
|
|
return DouyuTaskOut(**douyu_task_payload(task, include_detail=include_detail))
|
|
|
|
|
|
def _config_out(config: DouyuConfig) -> DouyuConfigOut:
|
|
return DouyuConfigOut(
|
|
manual_id=douyu_config_value("manual_id", config.manual_id),
|
|
rid=douyu_config_value("rid", config.rid),
|
|
bind_act_alias=douyu_config_value("bind_act_alias", config.bind_act_alias),
|
|
confirm_act_alias=douyu_config_value("confirm_act_alias", config.confirm_act_alias),
|
|
legacy_act_alias=douyu_config_value("legacy_act_alias", config.legacy_act_alias),
|
|
room_id=douyu_config_value("room_id", config.room_id),
|
|
elite_amount=douyu_config_value("elite_amount", config.elite_amount),
|
|
esports_manual_id=douyu_config_value("esports_manual_id", config.esports_manual_id),
|
|
esports_act_alias=douyu_config_value("esports_act_alias", config.esports_act_alias),
|
|
esports_amount=douyu_config_value("esports_amount", config.esports_amount),
|
|
esports_chicken_gift_id=douyu_config_value("esports_chicken_gift_id", config.esports_chicken_gift_id),
|
|
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),
|
|
updated_at=config.updated_at,
|
|
)
|
|
|
|
|
|
@router.get("/task-types")
|
|
def task_types(current: User = Depends(require_permission("douyu:task"))):
|
|
"""返回斗鱼任务类型。"""
|
|
return SUPPORTED_DOUYU_TASK_TYPES
|
|
|
|
|
|
@router.get("/accounts")
|
|
def list_task_accounts(
|
|
search: str = Query(""),
|
|
tag: str = Query(""),
|
|
ids: str = Query(""),
|
|
page: int | None = Query(None, ge=1),
|
|
page_size: int = Query(50, ge=1, le=200),
|
|
db: Session = Depends(get_db),
|
|
current: User = Depends(get_current_user),
|
|
):
|
|
"""查看可执行斗鱼任务的账号(必须有成功 Cookie)。
|
|
|
|
tag 按账号标签精确筛选;ids 为逗号分隔的账号 ID,用于工作台按已导入账号过滤;为空时返回全部。
|
|
"""
|
|
query = _visible_task_accounts_query(db, current)
|
|
id_list = [int(x) for x in ids.split(",") if x.strip().isdigit()]
|
|
if id_list:
|
|
query = query.filter(Account.id.in_(id_list))
|
|
tag_text = (tag or "").strip()
|
|
if tag_text:
|
|
query = query.filter(Account.tag == tag_text)
|
|
search_text = (search or "").strip()
|
|
if search_text:
|
|
pattern = f"%{search_text}%"
|
|
query = query.filter(or_(
|
|
Account.username.ilike(pattern),
|
|
Account.uid.ilike(pattern),
|
|
Account.nickname.ilike(pattern),
|
|
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:
|
|
total = query.order_by(None).count()
|
|
query = query.order_by(Account.id.desc())
|
|
if page is not None:
|
|
query = query.offset((page - 1) * page_size).limit(page_size)
|
|
else:
|
|
query = query.limit(500)
|
|
result = [_account_out(account) for account in query.all()]
|
|
if page is not None:
|
|
return {"items": result, "total": total or 0, "page": page, "page_size": page_size}
|
|
return result
|
|
|
|
|
|
@router.get("/config", response_model=DouyuConfigOut)
|
|
def get_config(
|
|
db: Session = Depends(get_db),
|
|
current: User = Depends(require_permission("douyu:config")),
|
|
):
|
|
"""获取斗鱼活动配置。"""
|
|
return _config_out(ensure_douyu_config(db))
|
|
|
|
|
|
@router.put("/config", response_model=DouyuConfigOut)
|
|
def update_config(
|
|
req: DouyuConfigUpdate,
|
|
db: Session = Depends(get_db),
|
|
current: User = Depends(require_permission("douyu:config")),
|
|
):
|
|
"""更新斗鱼活动配置。"""
|
|
config = ensure_douyu_config(db)
|
|
for field in DOUYU_CONFIG_FIELDS:
|
|
value = getattr(req, field)
|
|
if value is None:
|
|
continue
|
|
setattr(config, field, value.strip() if isinstance(value, str) else value)
|
|
apply_douyu_config_defaults(config)
|
|
config.updated_at = datetime.now(timezone.utc)
|
|
db.commit()
|
|
db.refresh(config)
|
|
return _config_out(config)
|
|
|
|
|
|
@router.get("/goods", response_model=list[DouyuGoodsOut])
|
|
def list_goods(
|
|
db: Session = Depends(get_db),
|
|
current: User = Depends(require_permission("douyu:task")),
|
|
):
|
|
"""查看已缓存的斗鱼商品快照。"""
|
|
rows = db.query(DouyuGoodsSnapshot).order_by(DouyuGoodsSnapshot.id.asc()).all()
|
|
return rows
|
|
|
|
|
|
@router.get("/esports-goods", response_model=list[DouyuGoodsOut])
|
|
def list_esports_goods(
|
|
db: Session = Depends(get_db),
|
|
current: User = Depends(require_permission("douyu:task")),
|
|
):
|
|
"""查看已缓存的电竞手册皮肤商城商品。"""
|
|
rows = db.query(DouyuEsportsGoodsSnapshot).order_by(DouyuEsportsGoodsSnapshot.id.asc()).all()
|
|
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,
|
|
db: Session = Depends(get_db),
|
|
current: User = Depends(require_permission("douyu:task")),
|
|
):
|
|
"""创建斗鱼任务记录并启动后台执行器。"""
|
|
if not req.account_ids:
|
|
raise HTTPException(status_code=400, detail="请选择斗鱼账号")
|
|
_require_task_account_access(db, current, req.account_ids)
|
|
|
|
cleanup_orphan_douyu_tasks(
|
|
db,
|
|
active_batch_ids=douyu_batch_registry.active_ids(),
|
|
statuses=("pending", "running"),
|
|
message="任务已中断(无执行器接管)",
|
|
)
|
|
|
|
try:
|
|
batch_id, count = create_douyu_planned_tasks(
|
|
db,
|
|
req.account_ids,
|
|
req.task_type,
|
|
current.id,
|
|
req.payload,
|
|
)
|
|
except ValueError as exc:
|
|
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
|
if count == 0:
|
|
raise HTTPException(status_code=400, detail="没有可执行的斗鱼账号,请先登录获取 Cookie")
|
|
|
|
log_queue = asyncio.Queue()
|
|
loop = asyncio.get_running_loop()
|
|
thread_db = SessionLocal()
|
|
runner = DouyuBatchRunner(
|
|
db=thread_db,
|
|
batch_id=batch_id,
|
|
task_type=req.task_type,
|
|
payload=req.payload,
|
|
log_queue=log_queue,
|
|
loop=loop,
|
|
concurrency=req.concurrency,
|
|
)
|
|
douyu_batch_registry.register(batch_id, log_queue, loop, runner)
|
|
|
|
thread = threading.Thread(target=runner.run, daemon=True)
|
|
thread.start()
|
|
|
|
return {"batch_id": batch_id, "count": count, "success": True}
|
|
|
|
|
|
@router.get("/tasks")
|
|
def list_tasks(
|
|
batch_id: str | None = None,
|
|
include_detail: bool = Query(False, description="是否返回完整任务结果(默认否,轮询请保持 false)"),
|
|
page: int | None = Query(None, ge=1),
|
|
page_size: int = Query(100, ge=1, le=500),
|
|
db: Session = Depends(get_db),
|
|
current: User = Depends(require_permission("douyu:task")),
|
|
):
|
|
"""查看斗鱼任务记录,默认只返回最近 100 条。"""
|
|
query = _visible_tasks_query(db, current)
|
|
if batch_id:
|
|
query = query.filter(DouyuTask.batch_id == batch_id)
|
|
total = None
|
|
if page is not None:
|
|
total = query.enable_eagerloads(False).order_by(None).count()
|
|
query = query.order_by(DouyuTask.id.desc())
|
|
if page is not None:
|
|
query = query.offset((page - 1) * page_size).limit(page_size)
|
|
else:
|
|
query = query.limit(page_size)
|
|
result = [_task_out(task, include_detail=include_detail) for task in query.all()]
|
|
if page is not None:
|
|
return {"items": result, "total": total or 0, "page": page, "page_size": page_size}
|
|
return result
|
|
|
|
|
|
@router.post("/tasks/cleanup-orphans")
|
|
def cleanup_orphan_tasks(
|
|
db: Session = Depends(get_db),
|
|
current: User = Depends(require_permission("douyu:task")),
|
|
):
|
|
"""手动清理没有内存执行器接管的斗鱼任务。"""
|
|
cleaned = cleanup_orphan_douyu_tasks(
|
|
db,
|
|
active_batch_ids=douyu_batch_registry.active_ids(),
|
|
statuses=("pending", "running"),
|
|
message="任务已中断(无执行器接管)",
|
|
)
|
|
return {"message": f"已清理 {cleaned} 个残留任务", "cleaned": cleaned, "success": True}
|
|
|
|
|
|
@router.get("/tasks/{task_id}", response_model=DouyuTaskOut)
|
|
def get_task(
|
|
task_id: int,
|
|
db: Session = Depends(get_db),
|
|
current: User = Depends(require_permission("douyu:task")),
|
|
):
|
|
"""获取单条斗鱼任务详情。"""
|
|
task = _visible_tasks_query(db, current).filter(DouyuTask.id == task_id).first()
|
|
if not task:
|
|
raise HTTPException(status_code=404, detail="任务不存在")
|
|
return _task_out(task, include_detail=True)
|
|
|
|
|
|
@router.post("/stop/{batch_id}")
|
|
def stop_batch(
|
|
batch_id: str,
|
|
db: Session = Depends(get_db),
|
|
current: User = Depends(require_permission("douyu:task")),
|
|
):
|
|
"""停止正在运行的斗鱼批次。"""
|
|
_require_batch_owner(db, current, batch_id)
|
|
batch = douyu_batch_registry.get(batch_id)
|
|
if batch:
|
|
if batch.get("finished"):
|
|
douyu_batch_registry.pop(batch_id)
|
|
cleaned = cleanup_orphan_douyu_tasks(db, batch_id=batch_id, message="批次已结束")
|
|
if cleaned:
|
|
return {"message": f"批次已结束,已清理 {cleaned} 个残留任务", "success": True}
|
|
raise HTTPException(status_code=404, detail="批次已结束")
|
|
batch["runner"].stop()
|
|
return {"message": "已发送停止信号", "success": True}
|
|
|
|
cleaned = cleanup_orphan_douyu_tasks(db, batch_id=batch_id, message="任务已停止(批次不存在)")
|
|
if cleaned:
|
|
return {"message": f"已清理 {cleaned} 个残留任务", "success": True}
|
|
raise HTTPException(status_code=404, detail="批次不存在或已结束")
|
|
|
|
|
|
@router.websocket("/ws/{batch_id}")
|
|
async def ws_douyu_logs(websocket: WebSocket, batch_id: str):
|
|
"""斗鱼实时日志推送通道。"""
|
|
user = authenticate_websocket(websocket)
|
|
if not user:
|
|
await websocket.close(code=1008, reason="未授权")
|
|
return
|
|
if not user_has_permission(user, "douyu:task"):
|
|
await websocket.close(code=1008, reason="无权限")
|
|
return
|
|
db = SessionLocal()
|
|
try:
|
|
_require_batch_owner(db, user, batch_id)
|
|
except HTTPException:
|
|
await websocket.close(code=1008, reason="无权访问该任务批次")
|
|
return
|
|
finally:
|
|
db.close()
|
|
await websocket.accept()
|
|
|
|
batch = douyu_batch_registry.get(batch_id)
|
|
if not batch:
|
|
await websocket.send_json({"level": "error", "message": "批次不存在或已结束"})
|
|
await websocket.close()
|
|
return
|
|
|
|
log_queue: asyncio.Queue = batch["log_queue"]
|
|
try:
|
|
while True:
|
|
try:
|
|
msg = await asyncio.wait_for(log_queue.get(), timeout=30)
|
|
await websocket.send_json(msg)
|
|
if msg.get("level") == "result":
|
|
await asyncio.sleep(0.1)
|
|
break
|
|
except asyncio.TimeoutError:
|
|
await websocket.send_json({"level": "heartbeat", "message": ""})
|
|
except WebSocketDisconnect:
|
|
pass
|
|
finally:
|
|
latest = douyu_batch_registry.get(batch_id)
|
|
if latest and latest.get("finished"):
|
|
douyu_batch_registry.pop(batch_id)
|