初步增加, 扫码登录成功
This commit is contained in:
@@ -154,6 +154,8 @@ _SENSITIVE_COLUMNS: tuple[tuple[str, str], ...] = (
|
||||
("huya_register_items", "password"),
|
||||
("huya_register_items", "cookie"),
|
||||
("huya_register_success_logs", "password"),
|
||||
("yyb_recharge_tasks", "login_qr_data"),
|
||||
("yyb_recharge_tasks", "payment_qr_data"),
|
||||
("proxy_config", "api_url"),
|
||||
("proxy_config", "http"),
|
||||
("proxy_config", "https"),
|
||||
|
||||
+6
-1
@@ -12,7 +12,7 @@ from fastapi.responses import FileResponse
|
||||
from starlette.middleware.base import BaseHTTPMiddleware
|
||||
|
||||
from .database import init_db
|
||||
from .routers import auth, users, accounts, account_check, dashboard, login, proxy, cookies, huya, douyu
|
||||
from .routers import auth, users, accounts, account_check, dashboard, login, proxy, cookies, huya, douyu, yyb
|
||||
from .schemas import AppInfo
|
||||
from .version import get_app_version
|
||||
from utils import setup_logger
|
||||
@@ -41,6 +41,10 @@ async def lifespan(app: FastAPI):
|
||||
cleaned_douyu = cleanup_orphan_douyu_tasks(db, message="任务已中断(服务重启)")
|
||||
if cleaned_douyu:
|
||||
logger.info(f"启动清理斗鱼残留任务: {cleaned_douyu} 条")
|
||||
from .services.yyb_service import cleanup_orphan_yyb_tasks
|
||||
cleaned_yyb = cleanup_orphan_yyb_tasks(db, message="任务已中断(服务重启)")
|
||||
if cleaned_yyb:
|
||||
logger.info(f"启动清理应用宝残留任务: {cleaned_yyb} 条")
|
||||
finally:
|
||||
db.close()
|
||||
yield
|
||||
@@ -97,6 +101,7 @@ app.include_router(proxy.router)
|
||||
app.include_router(cookies.router)
|
||||
app.include_router(huya.router)
|
||||
app.include_router(douyu.router)
|
||||
app.include_router(yyb.router)
|
||||
|
||||
|
||||
@app.get("/api/health")
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
"""增加应用宝和平精英充值任务"""
|
||||
|
||||
from typing import Sequence, Union
|
||||
from alembic import op
|
||||
import sqlalchemy as sa
|
||||
|
||||
revision: str = "20260812_0020"
|
||||
down_revision: Union[str, None] = "20260808_0019"
|
||||
branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
bind = op.get_bind()
|
||||
if sa.inspect(bind).has_table("yyb_recharge_tasks"):
|
||||
return
|
||||
op.create_table(
|
||||
"yyb_recharge_tasks",
|
||||
sa.Column("id", sa.Integer(), primary_key=True, autoincrement=True),
|
||||
sa.Column("task_id", sa.String(64), nullable=False),
|
||||
sa.Column("created_by", sa.Integer(), sa.ForeignKey("users.id"), nullable=False),
|
||||
sa.Column("worker_job_id", sa.String(64), nullable=False),
|
||||
sa.Column("provider", sa.String(16), nullable=False, server_default=""),
|
||||
sa.Column("platform", sa.String(16), nullable=False, server_default="android"),
|
||||
sa.Column("points", sa.Integer(), nullable=True),
|
||||
sa.Column("product_id", sa.String(128), nullable=True, server_default=""),
|
||||
sa.Column("zone_id", sa.String(64), nullable=True, server_default=""),
|
||||
sa.Column("zone_name", sa.String(128), nullable=True, server_default=""),
|
||||
sa.Column("role_id", sa.String(64), nullable=True, server_default=""),
|
||||
sa.Column("role_name", sa.String(128), nullable=True, server_default=""),
|
||||
sa.Column("status", sa.String(32), nullable=False, server_default="created"),
|
||||
sa.Column("phase", sa.String(32), nullable=False, server_default="login"),
|
||||
sa.Column("message", sa.String(512), nullable=True, server_default=""),
|
||||
sa.Column("result", sa.JSON(), nullable=True),
|
||||
sa.Column("login_qr_data", sa.Text(), nullable=True),
|
||||
sa.Column("payment_qr_data", sa.Text(), nullable=True),
|
||||
sa.Column("created_at", sa.DateTime(), nullable=True),
|
||||
sa.Column("finished_at", sa.DateTime(), nullable=True),
|
||||
)
|
||||
op.create_index("uq_yyb_recharge_tasks_task_id", "yyb_recharge_tasks", ["task_id"], unique=True)
|
||||
op.create_index("ix_yyb_recharge_tasks_created_by", "yyb_recharge_tasks", ["created_by"])
|
||||
op.create_index("ix_yyb_recharge_tasks_worker_job_id", "yyb_recharge_tasks", ["worker_job_id"])
|
||||
op.create_index("ix_yyb_recharge_tasks_status", "yyb_recharge_tasks", ["status"])
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
bind = op.get_bind()
|
||||
if sa.inspect(bind).has_table("yyb_recharge_tasks"):
|
||||
op.drop_table("yyb_recharge_tasks")
|
||||
@@ -0,0 +1,28 @@
|
||||
"""扩展应用宝二维码密文字段容量"""
|
||||
|
||||
from typing import Sequence, Union
|
||||
|
||||
from alembic import op
|
||||
import sqlalchemy as sa
|
||||
|
||||
|
||||
revision: str = "20260812_0021"
|
||||
down_revision: Union[str, None] = "20260812_0020"
|
||||
branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def _is_mysql() -> bool:
|
||||
return op.get_bind().dialect.name in {"mysql", "mariadb"}
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
if _is_mysql():
|
||||
op.execute("ALTER TABLE yyb_recharge_tasks MODIFY login_qr_data MEDIUMTEXT NULL")
|
||||
op.execute("ALTER TABLE yyb_recharge_tasks MODIFY payment_qr_data MEDIUMTEXT NULL")
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
if _is_mysql():
|
||||
op.execute("ALTER TABLE yyb_recharge_tasks MODIFY login_qr_data TEXT NULL")
|
||||
op.execute("ALTER TABLE yyb_recharge_tasks MODIFY payment_qr_data TEXT NULL")
|
||||
@@ -125,6 +125,35 @@ class DouyuTask(Base):
|
||||
account = relationship("Account", back_populates="douyu_tasks")
|
||||
|
||||
|
||||
class YybRechargeTask(Base):
|
||||
"""应用宝和平精英点券充值任务。
|
||||
|
||||
YYB 登录身份与斗鱼账号无关,因此任务只关联创建者,不复用 Account。
|
||||
"""
|
||||
__tablename__ = "yyb_recharge_tasks"
|
||||
|
||||
id = Column(Integer, primary_key=True, autoincrement=True)
|
||||
task_id = Column(String(64), unique=True, nullable=False, index=True)
|
||||
created_by = Column(Integer, ForeignKey("users.id"), nullable=False, index=True)
|
||||
worker_job_id = Column(String(64), nullable=False, index=True)
|
||||
provider = Column(String(16), default="", nullable=False)
|
||||
platform = Column(String(16), default="android", nullable=False)
|
||||
points = Column(Integer, nullable=True)
|
||||
product_id = Column(String(128), default="")
|
||||
zone_id = Column(String(64), default="")
|
||||
zone_name = Column(String(128), default="")
|
||||
role_id = Column(String(64), default="")
|
||||
role_name = Column(String(128), default="")
|
||||
status = Column(String(32), default="created", nullable=False, index=True)
|
||||
phase = Column(String(32), default="login", nullable=False)
|
||||
message = Column(String(512), default="")
|
||||
result = Column(JSON, nullable=True)
|
||||
login_qr_data = Column(EncryptedText(), default="")
|
||||
payment_qr_data = Column(EncryptedText(), default="")
|
||||
created_at = Column(DateTime, default=_utcnow)
|
||||
finished_at = Column(DateTime, nullable=True)
|
||||
|
||||
|
||||
class DouyuConfig(Base):
|
||||
"""斗鱼业务配置"""
|
||||
__tablename__ = "douyu_config"
|
||||
|
||||
@@ -26,6 +26,9 @@ PERMISSIONS = {
|
||||
# 斗鱼活动
|
||||
"douyu:task": "斗鱼任务管理",
|
||||
"douyu:config": "斗鱼配置管理",
|
||||
"yyb:session": "应用宝扫码登录与角色查询",
|
||||
"yyb:recharge": "应用宝充值与付款",
|
||||
"yyb:history": "查看应用宝充值历史",
|
||||
# 虎牙
|
||||
"huya:account": "虎牙账号管理(兼容旧权限)",
|
||||
"huya:view_all": "查看所有虎牙账号",
|
||||
@@ -62,6 +65,9 @@ ROLE_PERMISSIONS = {
|
||||
"cookie:export",
|
||||
"douyu:task",
|
||||
"douyu:config",
|
||||
"yyb:session",
|
||||
"yyb:recharge",
|
||||
"yyb:history",
|
||||
"huya:account",
|
||||
"huya:view_all",
|
||||
"huya:import",
|
||||
@@ -80,6 +86,8 @@ ROLE_PERMISSIONS = {
|
||||
"support": [
|
||||
"account:view_assigned",
|
||||
"douyu:task",
|
||||
"yyb:session",
|
||||
"yyb:history",
|
||||
"huya:view_assigned",
|
||||
"huya:task",
|
||||
],
|
||||
|
||||
@@ -0,0 +1,124 @@
|
||||
"""应用宝和平精英充值工作台 API。"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import uuid
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException, Query
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from ..database import get_db
|
||||
from ..deps import get_current_user, require_permission
|
||||
from ..models import User, YybRechargeTask
|
||||
from ..permissions import user_has_permission
|
||||
from ..schemas import YybLoginRequest, YybSelectionRequest, YybTaskCreateRequest
|
||||
from ..services.yyb_service import public_task, sync_task
|
||||
from ..services.yyb_worker_client import YybWorkerClient, YybWorkerError
|
||||
|
||||
|
||||
router = APIRouter(prefix="/api/yyb", tags=["应用宝充值"])
|
||||
|
||||
|
||||
def _get_task(db: Session, task_id: int, current: User) -> YybRechargeTask:
|
||||
task = db.query(YybRechargeTask).filter(YybRechargeTask.id == task_id).first()
|
||||
if not task:
|
||||
raise HTTPException(404, "充值任务不存在")
|
||||
if task.created_by != current.id and not user_has_permission(current, "yyb:history"):
|
||||
raise HTTPException(403, "无权查看该充值任务")
|
||||
return task
|
||||
|
||||
|
||||
def _worker_call(call):
|
||||
try:
|
||||
return call()
|
||||
except YybWorkerError as exc:
|
||||
raise HTTPException(502, str(exc)) from exc
|
||||
|
||||
|
||||
@router.post("/tasks")
|
||||
def create_task(payload: YybTaskCreateRequest, db: Session = Depends(get_db), current: User = Depends(require_permission("yyb:session"))):
|
||||
data = _worker_call(YybWorkerClient().create_job)
|
||||
task = YybRechargeTask(task_id=uuid.uuid4().hex[:16], created_by=current.id,
|
||||
worker_job_id=str(data["job_id"]), status=str(data.get("status", "created")),
|
||||
phase="login", message="请选择登录方式")
|
||||
db.add(task)
|
||||
db.commit()
|
||||
db.refresh(task)
|
||||
return public_task(task, include_qr=False)
|
||||
|
||||
|
||||
@router.post("/tasks/{task_id}/login")
|
||||
def login(task_id: int, payload: YybLoginRequest, db: Session = Depends(get_db), current: User = Depends(require_permission("yyb:session"))):
|
||||
task = _get_task(db, task_id, current)
|
||||
data = _worker_call(lambda: YybWorkerClient().login(task.worker_job_id, payload.provider, payload.timeout))
|
||||
task.provider = payload.provider
|
||||
task.status = str(data.get("status", "waiting_login"))
|
||||
task.phase = "login"
|
||||
task.message = "请扫码登录并在手机确认"
|
||||
if data.get("qr_data"):
|
||||
task.login_qr_data = data["qr_data"]
|
||||
db.commit()
|
||||
return public_task(task)
|
||||
|
||||
|
||||
@router.get("/tasks/{task_id}")
|
||||
def get_task(task_id: int, db: Session = Depends(get_db), current: User = Depends(get_current_user)):
|
||||
task = _get_task(db, task_id, current)
|
||||
try:
|
||||
sync_task(db, task, YybWorkerClient())
|
||||
except YybWorkerError:
|
||||
# Worker 暂时重启时仍返回最近一次持久化状态。
|
||||
pass
|
||||
return public_task(task, include_qr=user_has_permission(current, "yyb:session"),
|
||||
include_payment_qr=user_has_permission(current, "yyb:recharge"))
|
||||
|
||||
|
||||
@router.get("/tasks")
|
||||
def list_tasks(limit: int = Query(50, ge=1, le=200), db: Session = Depends(get_db), current: User = Depends(require_permission("yyb:history"))):
|
||||
tasks = db.query(YybRechargeTask).order_by(YybRechargeTask.id.desc()).limit(limit).all()
|
||||
return [public_task(task, include_qr=False) for task in tasks]
|
||||
|
||||
|
||||
@router.get("/tasks/{task_id}/selection-options")
|
||||
def selection_options(task_id: int, platform: str = Query("android"), points: int | None = Query(None), zone_id: str | None = Query(None), db: Session = Depends(get_db), current: User = Depends(require_permission("yyb:session"))):
|
||||
task = _get_task(db, task_id, current)
|
||||
data = _worker_call(lambda: YybWorkerClient().selection_options(task.worker_job_id, platform, points, zone_id))
|
||||
task.platform = platform
|
||||
db.commit()
|
||||
return data
|
||||
|
||||
|
||||
@router.post("/tasks/{task_id}/selection")
|
||||
def selection(task_id: int, payload: YybSelectionRequest, db: Session = Depends(get_db), current: User = Depends(require_permission("yyb:session"))):
|
||||
task = _get_task(db, task_id, current)
|
||||
data = _worker_call(lambda: YybWorkerClient().selection(task.worker_job_id, payload.model_dump()))
|
||||
selected = data.get("selection", payload.model_dump())
|
||||
for field in ("platform", "points", "product_id", "zone_id", "zone_name", "role_id", "role_name"):
|
||||
setattr(task, field, selected[field])
|
||||
task.phase, task.status, task.message = "payment", "ready", "选择已保存,可以创建付款码"
|
||||
db.commit()
|
||||
return public_task(task)
|
||||
|
||||
|
||||
@router.post("/tasks/{task_id}/payment")
|
||||
def payment(task_id: int, db: Session = Depends(get_db), current: User = Depends(require_permission("yyb:recharge"))):
|
||||
task = _get_task(db, task_id, current)
|
||||
if not task.product_id or not task.role_id:
|
||||
raise HTTPException(400, "请先完成商品、区服和角色选择")
|
||||
data = _worker_call(lambda: YybWorkerClient().payment(task.worker_job_id))
|
||||
task.status = str(data.get("status", "running"))
|
||||
task.phase = "payment"
|
||||
task.message = "正在创建订单和付款码"
|
||||
db.commit()
|
||||
return public_task(task)
|
||||
|
||||
|
||||
@router.post("/tasks/{task_id}/stop")
|
||||
def stop(task_id: int, db: Session = Depends(get_db), current: User = Depends(require_permission("yyb:session"))):
|
||||
task = _get_task(db, task_id, current)
|
||||
data = _worker_call(lambda: YybWorkerClient().stop(task.worker_job_id))
|
||||
task.status = "failed"
|
||||
task.phase = "stopped"
|
||||
task.message = "任务已停止"
|
||||
db.commit()
|
||||
return public_task(task)
|
||||
@@ -13,6 +13,25 @@ from .huya_defaults import (
|
||||
)
|
||||
|
||||
|
||||
class YybTaskCreateRequest(BaseModel):
|
||||
"""创建 YYB 任务;登录方式在下一步选择。"""
|
||||
|
||||
|
||||
class YybLoginRequest(BaseModel):
|
||||
provider: str = Field(..., pattern="^(qq|wechat)$")
|
||||
timeout: int = Field(600, ge=60, le=1800)
|
||||
|
||||
|
||||
class YybSelectionRequest(BaseModel):
|
||||
platform: str = Field(..., pattern="^(android|ios)$")
|
||||
points: int = Field(..., gt=0)
|
||||
product_id: str = Field(..., min_length=1, max_length=128)
|
||||
zone_id: str = Field(..., min_length=1, max_length=64)
|
||||
zone_name: str = Field("", max_length=128)
|
||||
role_id: str = Field(..., min_length=1, max_length=64)
|
||||
role_name: str = Field(..., min_length=1, max_length=128)
|
||||
|
||||
|
||||
def _ensure_tz(dt: Optional[datetime]) -> Optional[datetime]:
|
||||
"""确保 datetime 带有 UTC 时区信息,无时区的视为 UTC。"""
|
||||
if dt is None:
|
||||
|
||||
@@ -0,0 +1,74 @@
|
||||
"""应用宝充值任务服务。"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any
|
||||
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from ..models import YybRechargeTask
|
||||
from .yyb_worker_client import YybWorkerClient
|
||||
|
||||
|
||||
def _utcnow():
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
|
||||
def cleanup_orphan_yyb_tasks(db: Session, message: str) -> int:
|
||||
rows = db.query(YybRechargeTask).filter(YybRechargeTask.status.in_(["created", "waiting_login", "running", "ready"])).all()
|
||||
for task in rows:
|
||||
task.status = "failed"
|
||||
task.message = message
|
||||
task.finished_at = _utcnow()
|
||||
if rows:
|
||||
db.commit()
|
||||
return len(rows)
|
||||
|
||||
|
||||
def sync_task(db: Session, task: YybRechargeTask, worker: YybWorkerClient) -> YybRechargeTask:
|
||||
data = worker.get_job(task.worker_job_id)
|
||||
task.status = str(data.get("status", task.status))
|
||||
task.phase = str(data.get("phase", task.phase))
|
||||
task.message = str(data.get("message", task.message))
|
||||
if data.get("provider"):
|
||||
task.provider = str(data["provider"])
|
||||
if data.get("qr_data"):
|
||||
task.login_qr_data = str(data["qr_data"])
|
||||
task.result = {**(task.result or {}), "login_qr_mime_type": data.get("qr_mime_type", "image/jpeg")}
|
||||
if data.get("payment_qr_data"):
|
||||
task.payment_qr_data = str(data["payment_qr_data"])
|
||||
task.result = {**(task.result or {}), "payment_qr_mime_type": data.get("payment_qr_mime_type", "image/png")}
|
||||
task.result = {
|
||||
**(task.result or {}),
|
||||
"logs": data.get("logs", []),
|
||||
"payment_status": data.get("payment_status"),
|
||||
}
|
||||
if task.status in {"success", "failed"} and task.finished_at is None:
|
||||
task.finished_at = _utcnow()
|
||||
db.commit()
|
||||
db.refresh(task)
|
||||
return task
|
||||
|
||||
|
||||
def public_task(task: YybRechargeTask, include_qr: bool = True,
|
||||
include_payment_qr: bool | None = None) -> dict[str, Any]:
|
||||
result: dict[str, Any] = {
|
||||
"id": task.id, "task_id": task.task_id,
|
||||
"provider": task.provider, "platform": task.platform, "points": task.points,
|
||||
"product_id": task.product_id, "zone_id": task.zone_id, "zone_name": task.zone_name,
|
||||
"role_id": task.role_id, "role_name": task.role_name, "status": task.status,
|
||||
"phase": task.phase, "message": task.message, "result": task.result,
|
||||
"created_by": task.created_by, "created_at": task.created_at,
|
||||
"finished_at": task.finished_at,
|
||||
}
|
||||
if task.result:
|
||||
result["login_qr_mime_type"] = task.result.get("login_qr_mime_type", "image/jpeg")
|
||||
result["payment_qr_mime_type"] = task.result.get("payment_qr_mime_type", "image/png")
|
||||
if include_payment_qr is None:
|
||||
include_payment_qr = include_qr
|
||||
if include_qr:
|
||||
result["login_qr_data"] = task.login_qr_data or ""
|
||||
if include_payment_qr:
|
||||
result["payment_qr_data"] = task.payment_qr_data or ""
|
||||
return result
|
||||
@@ -0,0 +1,57 @@
|
||||
"""应用宝 Worker HTTP 客户端。"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
from typing import Any
|
||||
|
||||
import requests
|
||||
|
||||
|
||||
class YybWorkerError(RuntimeError):
|
||||
"""Worker 返回业务错误。"""
|
||||
|
||||
|
||||
class YybWorkerClient:
|
||||
def __init__(self) -> None:
|
||||
self.base_url = os.getenv("YYB_WORKER_URL", "http://127.0.0.1:8810").rstrip("/")
|
||||
self.key = os.getenv("YYB_WORKER_KEY", "")
|
||||
self.timeout = float(os.getenv("YYB_WORKER_TIMEOUT", "30"))
|
||||
|
||||
def _request(self, method: str, path: str, payload: dict[str, Any] | None = None) -> dict[str, Any]:
|
||||
headers = {"Accept": "application/json"}
|
||||
if self.key:
|
||||
headers["Authorization"] = f"Bearer {self.key}"
|
||||
try:
|
||||
response = requests.request(method, self.base_url + path, json=payload,
|
||||
headers=headers, timeout=self.timeout)
|
||||
data = response.json()
|
||||
except (requests.RequestException, ValueError) as exc:
|
||||
raise YybWorkerError(f"应用宝 Worker 不可用: {exc}") from exc
|
||||
if response.status_code >= 400:
|
||||
raise YybWorkerError(str(data.get("detail", "Worker 请求失败")))
|
||||
return data
|
||||
|
||||
def create_job(self) -> dict[str, Any]:
|
||||
return self._request("POST", "/v1/jobs")
|
||||
|
||||
def login(self, worker_job_id: str, provider: str, timeout: int = 600) -> dict[str, Any]:
|
||||
return self._request("POST", f"/v1/jobs/{worker_job_id}/login",
|
||||
{"provider": provider, "timeout": timeout})
|
||||
|
||||
def get_job(self, worker_job_id: str) -> dict[str, Any]:
|
||||
return self._request("GET", f"/v1/jobs/{worker_job_id}")
|
||||
|
||||
def selection_options(self, worker_job_id: str, platform: str,
|
||||
points: int | None = None, zone_id: str | None = None) -> dict[str, Any]:
|
||||
return self._request("POST", f"/v1/jobs/{worker_job_id}/selection-options",
|
||||
{"platform": platform, "points": points, "zone_id": zone_id})
|
||||
|
||||
def selection(self, worker_job_id: str, selection: dict[str, Any]) -> dict[str, Any]:
|
||||
return self._request("POST", f"/v1/jobs/{worker_job_id}/selection", selection)
|
||||
|
||||
def payment(self, worker_job_id: str) -> dict[str, Any]:
|
||||
return self._request("POST", f"/v1/jobs/{worker_job_id}/payment")
|
||||
|
||||
def stop(self, worker_job_id: str) -> dict[str, Any]:
|
||||
return self._request("POST", f"/v1/jobs/{worker_job_id}/stop")
|
||||
Reference in New Issue
Block a user