Files

114 lines
3.6 KiB
Python

"""斗鱼账号检测路由。"""
from __future__ import annotations
import threading
from fastapi import APIRouter, Depends, HTTPException
from fastapi.responses import StreamingResponse
from sqlalchemy.orm import Session
from ..database import get_db
from ..deps import get_current_user
from ..models import ProxyConfig as ProxyConfigModel
from ..models import User
from ..permissions import user_has_permission
from ..schemas import AccountCheckBatchOut, AccountCheckBatchRequest
from ..services.account_check_service import (
account_check_registry,
parse_account_check_lines,
)
router = APIRouter(prefix="/api/account-check", tags=["账号检测"])
def _get_runner_or_404(batch_id: str):
"""读取账号检测批次。"""
runner = account_check_registry.get(batch_id)
if not runner:
raise HTTPException(status_code=404, detail="批次不存在或服务已重启")
return runner
def _require_account_check_perm(user: User) -> None:
"""账号检测权限;兼容已有批量登录权限。"""
if not (
user_has_permission(user, "account:check")
or user_has_permission(user, "login:batch")
):
raise HTTPException(status_code=403, detail="权限不足")
@router.post("/batches", response_model=AccountCheckBatchOut)
def create_account_check_batch(
req: AccountCheckBatchRequest,
db: Session = Depends(get_db),
current: User = Depends(get_current_user),
):
"""启动斗鱼账号检测批次。"""
_require_account_check_perm(current)
try:
accounts = parse_account_check_lines(req.text)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
if not accounts:
raise HTTPException(status_code=400, detail="没有识别到有效账号")
proxy_config = db.query(ProxyConfigModel).first() if req.use_proxy else None
runner = account_check_registry.create(
accounts=accounts,
created_by=current.id,
concurrency=req.concurrency,
max_login_retries=req.max_login_retries,
max_total_time=req.max_total_time,
use_proxy=req.use_proxy,
proxy_config=proxy_config,
)
thread = threading.Thread(target=runner.run, daemon=True)
thread.start()
return runner.snapshot()
@router.get("/batches/{batch_id}", response_model=AccountCheckBatchOut)
def get_account_check_batch(
batch_id: str,
current: User = Depends(get_current_user),
):
"""查询斗鱼账号检测批次。"""
_require_account_check_perm(current)
runner = _get_runner_or_404(batch_id)
return runner.snapshot()
@router.post("/batches/{batch_id}/stop")
def stop_account_check_batch(
batch_id: str,
current: User = Depends(get_current_user),
):
"""停止斗鱼账号检测批次。"""
_require_account_check_perm(current)
runner = _get_runner_or_404(batch_id)
runner.stop()
return {"message": "已发送停止信号", "success": True}
@router.get("/batches/{batch_id}/download")
def download_account_check_batch(
batch_id: str,
current: User = Depends(get_current_user),
):
"""下载斗鱼账号检测分类结果 zip。"""
_require_account_check_perm(current)
runner = _get_runner_or_404(batch_id)
snapshot = runner.snapshot()
if snapshot["status"] in {"pending", "running"}:
raise HTTPException(status_code=400, detail="批次尚未完成")
content, filename = runner.build_zip()
return StreamingResponse(
iter([content]),
media_type="application/zip",
headers={"Content-Disposition": f'attachment; filename="{filename}"'},
)