Files
Azmat@qq.com 0bd1db6bf9 模特库标签分页与全能创作收口:藏长视频、选择器分页
角色库导入打标签并支持筛选;模特库与全能创作角色/商品选择改为每页 20 条分页。临时限制成片 ≤60 秒,过滤对话里的超长时长选项,并收拢本地全能创作与后台用户相关修复。
2026-09-21 16:24:37 +08:00

1054 lines
47 KiB
Python

"""平台超管后台 · 跨团队端点。所有视图统一挂 IsPlatformAdmin,非超管一律 403,写操作记审计。"""
import logging
from decimal import Decimal, ROUND_HALF_UP
from django.core.exceptions import ValidationError as DjangoValidationError
from django.db import transaction
from django.db.models import Case, CharField, Count, F, Prefetch, Q, Value, When
from rest_framework import status
from rest_framework.authtoken.models import Token
from rest_framework.decorators import api_view, permission_classes
from rest_framework.response import Response
from apps.accounts.audit import log_admin_action
from apps.accounts.models import Invitation, Team, TeamMember, User
from apps.accounts.permissions import IsPlatformAdmin
from apps.accounts.serializers import InvitationSerializer
from apps.ai.model_catalog import invalidate_model_catalog_cache
from apps.ai.models import AITask, ModelConfig, ModelProvider, PromptTemplate, QualityWord
from apps.assets.models import Asset
from apps.assets.review import poll_asset_review, submit_asset_for_review
from apps.billing.models import BillingConfig, CreditAccount, CreditLedger, QuotaPolicy
from apps.billing.pricing import get_billing_config, invalidate_billing_config_cache
from apps.billing.services.ledger import adjust_credit
from apps.common.pagination import DefaultPagination
from apps.products.models import Product
from apps.projects.models import Project
from .serializers import (
COST_ANOMALY_RATIO,
AdminLedgerSerializer,
AdminModelConfigSerializer,
AdminModelProviderSerializer,
AdminProjectSerializer,
AdminQuotaPolicySerializer,
AdminReviewAssetSerializer,
AdminTaskDetailSerializer,
AdminTaskSerializer,
AdminTeamMemberSerializer,
AdminTeamSerializer,
AdminUserSerializer,
PromptTemplateSerializer,
QualityWordSerializer,
)
logger = logging.getLogger(__name__)
# 后台「刷新状态」只向供应商拉取已提交的异步视频任务;单次上限避免请求拖死。
_ADMIN_TASK_POLL_LIMIT = 20
_ADMIN_TASK_POLL_STATUSES = (AITask.Status.SUBMITTED, AITask.Status.POLLING)
# 任务监控「生成中」Tab:含已创建/已预留/已提交/轮询/后处理(与自由创作 IN_FLIGHT 对齐)
_ADMIN_TASK_INFLIGHT_STATUSES = (
AITask.Status.CREATED,
AITask.Status.RESERVED,
AITask.Status.SUBMITTED,
AITask.Status.POLLING,
AITask.Status.POSTPROCESSING,
)
def _team_qs():
return (
Team.objects.select_related("owner", "credit_account")
.annotate(member_count_anno=Count("members", distinct=True))
)
def _admin_user_qs():
"""用户列表:把成员关系连同团队钱包一次性 prefetch 出来,序列化器才能免查询地拍出余额。
team_member_count 用来判断钱包是不是多人共享池(决定发积分弹窗的提示文案)。"""
memberships = (
TeamMember.objects.select_related("team", "team__credit_account")
.annotate(team_member_count=Count("team__members", distinct=True))
.order_by("created_at")
)
return User.objects.prefetch_related(Prefetch("team_memberships", queryset=memberships))
def _wallet_team(user):
"""用户的钱包团队:第一个 active 成员关系(与 get_current_team 同口径)。无团队返回 None。"""
membership = (
user.team_memberships.filter(status=TeamMember.Status.ACTIVE)
.select_related("team")
.order_by("created_at")
.first()
)
return membership.team if membership else None
def _parse_points(raw, *, field: str, allow_negative: bool):
"""解析后台传来的积分数量。积分全站按整数流通,这里统一 HALF_UP 取整(与计价引擎同口径)。
返回 (Decimal, None) 或 (None, Response)。"""
from decimal import InvalidOperation
try:
value = Decimal(str(raw))
except (InvalidOperation, TypeError, ValueError):
return None, Response({field: ["数值格式不正确"]}, status=status.HTTP_400_BAD_REQUEST)
# Decimal("NaN")/("Infinity") 构造不抛,进库才炸 500 → 显式拦
if not value.is_finite():
return None, Response({field: ["数值格式不正确"]}, status=status.HTTP_400_BAD_REQUEST)
value = value.quantize(Decimal("1"), rounding=ROUND_HALF_UP)
if not allow_negative and value < 0:
return None, Response({field: ["积分不能为负"]}, status=status.HTTP_400_BAD_REQUEST)
if abs(value) > Decimal("10000000"):
return None, Response({field: ["单次不能超过 1,000 万积分"]}, status=status.HTTP_400_BAD_REQUEST)
return value, None
@api_view(["GET", "POST"])
@permission_classes([IsPlatformAdmin])
def admin_invitations(request):
"""GET 列所有邀请码(跨团队,可按 kind/status/search 过滤,分页);
POST 平台超管发「开团队码」(create_team,不绑团队,新用户凭码开新团队当 owner)。"""
if request.method == "GET":
qs = Invitation.objects.select_related("team", "used_by").order_by("-created_at")
kind = request.query_params.get("kind")
if kind in {Invitation.Kind.JOIN_TEAM, Invitation.Kind.CREATE_TEAM}:
qs = qs.filter(kind=kind)
st = request.query_params.get("status")
if st in dict(Invitation.Status.choices):
qs = qs.filter(status=st)
search = (request.query_params.get("search") or "").strip()
if search:
qs = qs.filter(Q(code__icontains=search) | Q(team__name__icontains=search))
paginator = DefaultPagination()
page = paginator.paginate_queryset(qs, request)
return paginator.get_paginated_response(InvitationSerializer(page, many=True).data)
email = str(request.data.get("email") or "").strip()
invite = Invitation.objects.create(
kind=Invitation.Kind.CREATE_TEAM,
team=None,
email=email,
created_by=request.user,
)
log_admin_action(
request,
"invite.issue_create_team",
target_type="invitation",
target_id=invite.id,
target_name=invite.code,
after={"kind": invite.kind, "email": email},
)
return Response(InvitationSerializer(invite).data, status=status.HTTP_201_CREATED)
@api_view(["POST"])
@permission_classes([IsPlatformAdmin])
def admin_revoke_invitation(request, invite_id):
"""撤销一个待用邀请码(任意团队)。已用/已撤销/已过期则原样返回(幂等)。"""
invite = Invitation.objects.select_related("team").filter(id=invite_id).first()
if invite is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
if invite.status == Invitation.Status.PENDING:
invite.status = Invitation.Status.REVOKED
invite.save(update_fields=["status", "updated_at"])
log_admin_action(
request,
"invite.revoke",
target_type="invitation",
target_id=invite.id,
target_name=invite.code,
)
return Response(InvitationSerializer(invite).data)
# ─────────────────────────── 团队管理 ───────────────────────────
@api_view(["GET"])
@permission_classes([IsPlatformAdmin])
def admin_teams(request):
"""列所有团队(跨团队,可按 status/search 过滤,分页)。
默认隐藏个人团队(后台直建用户的钱包容器,只在用户页以「余额」形态出现),?include_personal=1 可看全。"""
qs = _team_qs().order_by("-created_at")
if str(request.query_params.get("include_personal") or "") not in {"1", "true"}:
qs = qs.filter(is_personal=False)
st = request.query_params.get("status")
if st in dict(Team.Status.choices):
qs = qs.filter(status=st)
search = (request.query_params.get("search") or "").strip()
if search:
qs = qs.filter(name__icontains=search)
paginator = DefaultPagination()
page = paginator.paginate_queryset(qs, request)
return paginator.get_paginated_response(AdminTeamSerializer(page, many=True).data)
@api_view(["GET"])
@permission_classes([IsPlatformAdmin])
def admin_team_detail(request, team_id):
team = _team_qs().filter(id=team_id).first()
if team is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
members = team.members.select_related("user").order_by("created_at")
data = AdminTeamSerializer(team).data
data["members"] = AdminTeamMemberSerializer(members, many=True).data
return Response(data)
@api_view(["POST"])
@permission_classes([IsPlatformAdmin])
def admin_team_toggle(request, team_id):
"""启停团队。团队停用后其成员登录会被拒(login 校验团队状态)。"""
team = Team.objects.filter(id=team_id).first()
if team is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
before = team.status
team.status = Team.Status.ACTIVE if team.status == Team.Status.DISABLED else Team.Status.DISABLED
team.save(update_fields=["status", "updated_at"])
log_admin_action(
request,
"team.toggle_status",
target_type="team",
target_id=team.id,
target_name=team.name,
before={"status": before},
after={"status": team.status},
)
return Response(AdminTeamSerializer(_team_qs().get(id=team.id)).data)
# ─────────────────────────── 用户管理 ───────────────────────────
@api_view(["POST"])
@permission_classes([IsPlatformAdmin])
def admin_team_pricing(request, team_id):
"""团队差异化调价(jimeng 同款诉求):设置团队价格系数(0.10~10.00)。
最终积分价 = 挂牌价 × 系数;全部计费类型统一生效;视频在途任务用下单快照不受影响。"""
from decimal import InvalidOperation
team = Team.objects.filter(id=team_id).first()
if team is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
try:
value = Decimal(str(request.data.get("price_multiplier")))
except (InvalidOperation, TypeError, ValueError):
return Response({"price_multiplier": ["数值格式不正确"]}, status=status.HTTP_400_BAD_REQUEST)
if not value.is_finite() or value < Decimal("0.10") or value > Decimal("10"):
return Response({"price_multiplier": ["系数需在 0.10 ~ 10.00 之间"]}, status=status.HTTP_400_BAD_REQUEST)
before = str(team.price_multiplier)
# 显式 HALF_UP:quantize 默认银行家舍入(0.125→0.12),与全仓取整纪律不一致(review 确认)
team.price_multiplier = value.quantize(Decimal("0.01"), rounding=ROUND_HALF_UP)
team.save(update_fields=["price_multiplier", "updated_at"])
log_admin_action(
request,
"team.pricing_update",
target_type="team",
target_id=team.id,
target_name=team.name,
before={"price_multiplier": before},
after={"price_multiplier": str(team.price_multiplier)},
)
return Response(AdminTeamSerializer(_team_qs().get(id=team.id)).data)
@api_view(["GET", "POST"])
@permission_classes([IsPlatformAdmin])
def admin_users(request):
"""GET 列所有用户(可按 status/search 过滤,分页,每行带钱包余额);
POST 平台超管直接开户(不走邀请码),顺带配一个个人团队当钱包。"""
if request.method == "POST":
return _admin_create_user(request)
qs = _admin_user_qs().order_by("-date_joined")
st = request.query_params.get("status")
if st in dict(User.Status.choices):
qs = qs.filter(status=st)
search = (request.query_params.get("search") or "").strip()
if search:
qs = qs.filter(Q(username__icontains=search) | Q(first_name__icontains=search))
paginator = DefaultPagination()
page = paginator.paginate_queryset(qs, request)
return paginator.get_paginated_response(AdminUserSerializer(page, many=True).data)
def _admin_create_user(request):
"""后台直建用户:{username, display_name?, password, initial_credits?}。
username = 登录账号(仅英文+数字,固定 6 位);display_name = 展示用用户名(可中文)。
积分账户是 OneToOne 挂 Team 的,所以这里给新用户配一个 is_personal 的一人团队当专属钱包 ——
计费链路(预留/实扣/流水/额度)一行不用改,后台却能按「用户」视角管人和钱。
后续要恢复多人协作,把人拉进同一个团队即可,不需要迁数据。"""
import re
username = str(request.data.get("username") or "").strip()
display_name = str(
request.data.get("display_name")
or request.data.get("name")
or ""
).strip()
password = str(request.data.get("password") or "").strip()
if not display_name:
return Response({"display_name": ["请填写用户名"]}, status=status.HTTP_400_BAD_REQUEST)
if len(display_name) > 64:
return Response({"display_name": ["用户名不能超过 64 字"]}, status=status.HTTP_400_BAD_REQUEST)
if not username:
return Response({"username": ["请填写登录账号"]}, status=status.HTTP_400_BAD_REQUEST)
if not re.fullmatch(r"[A-Za-z0-9]{6}", username):
return Response(
{"username": ["登录账号须为 6 位英文或数字,不能含其他字符"]},
status=status.HTTP_400_BAD_REQUEST,
)
# 统一小写入库,避免 Abc123 / abc123 被当成两个账号
username = username.lower()
if User.objects.filter(username__iexact=username).exists():
return Response({"username": ["该登录账号已存在,不能重复"]}, status=status.HTTP_400_BAD_REQUEST)
if len(password) < 8:
return Response({"password": ["密码至少 8 位"]}, status=status.HTTP_400_BAD_REQUEST)
initial, err = _parse_points(request.data.get("initial_credits") or 0, field="initial_credits", allow_negative=False)
if err is not None:
return err
with transaction.atomic():
user = User.objects.create_user(
username=username,
password=password,
first_name=display_name,
)
# 个人团队名用展示名,侧栏第一行显示用户名、第二行显示登录账号
team = Team.objects.create(name=display_name, owner=user, is_personal=True)
TeamMember.objects.create(team=team, user=user, role=TeamMember.Role.OWNER)
CreditAccount.objects.create(team=team, balance=initial)
if initial > 0:
# 开户额度必须进流水,否则这笔钱在账单里查无凭证、无法对账(与注册赠送同口径)
CreditLedger.objects.create(
team=team,
user=request.user,
ledger_type=CreditLedger.Type.RECHARGE,
amount=initial,
balance_after=initial,
reason="后台开户初始积分",
metadata={"kind": "admin_create_user", "target_user_id": str(user.id)},
)
log_admin_action(
request,
"user.create",
target_type="user",
target_id=user.id,
target_name=user.username,
after={
"initial_credits": str(initial),
"wallet_team": str(team.id),
"display_name": display_name,
},
)
return Response(AdminUserSerializer(_admin_user_qs().get(id=user.id)).data, status=status.HTTP_201_CREATED)
@api_view(["POST"])
@permission_classes([IsPlatformAdmin])
def admin_user_credits(request, user_id):
"""给用户发/扣积分:{amount, reason?}。amount 可正可负,落 ADJUSTMENT 流水。
钱实际进的是这个用户的钱包团队 —— 个人团队即专属钱包;老用户若在多人团队里,
这笔会进团队共享池(前端弹窗已明确提示,后端流水的 metadata 记下目标用户备查)。"""
user = User.objects.filter(id=user_id).first()
if user is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
team = _wallet_team(user)
if team is None:
return Response({"detail": "该用户没有可用的积分钱包(无生效团队)"}, status=status.HTTP_400_BAD_REQUEST)
amount, err = _parse_points(request.data.get("amount"), field="amount", allow_negative=True)
if err is not None:
return err
if amount == 0:
return Response({"amount": ["调额金额不能为 0"]}, status=status.HTTP_400_BAD_REQUEST)
reason = str(request.data.get("reason") or "").strip()
try:
ledger = adjust_credit(
team=team,
amount=amount,
reason=reason or f"平台给 {user.username} 调额",
operator=request.user,
metadata={"target_user_id": str(user.id), "target_username": user.username},
)
except ValueError as exc:
return Response({"detail": str(exc)}, status=status.HTTP_400_BAD_REQUEST)
log_admin_action(
request,
"user.credit_adjust",
target_type="user",
target_id=user.id,
target_name=user.username,
after={"amount": str(amount), "balance_after": str(ledger.balance_after), "team": str(team.id), "reason": reason},
)
return Response(AdminUserSerializer(_admin_user_qs().get(id=user.id)).data, status=status.HTTP_201_CREATED)
@api_view(["POST"])
@permission_classes([IsPlatformAdmin])
def admin_user_toggle(request, user_id):
"""启停用户。停用即清 token 强制下线;不允许停用平台超管(防自锁)。"""
user = User.objects.filter(id=user_id).first()
if user is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
if user.is_platform_admin:
return Response({"detail": "不能停用平台超管"}, status=status.HTTP_400_BAD_REQUEST)
before = user.status
user.status = User.Status.ACTIVE if user.status == User.Status.DISABLED else User.Status.DISABLED
user.save(update_fields=["status"])
if user.status == User.Status.DISABLED:
Token.objects.filter(user=user).delete()
log_admin_action(
request,
"user.toggle_status",
target_type="user",
target_id=user.id,
target_name=user.username,
before={"status": before},
after={"status": user.status},
)
return Response(AdminUserSerializer(_admin_user_qs().get(id=user.id)).data)
@api_view(["POST"])
@permission_classes([IsPlatformAdmin])
def admin_user_reset_password(request, user_id):
"""平台超管强制改用户密码(改后清 token 强制重登)。"""
user = User.objects.filter(id=user_id).first()
if user is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
password = str(request.data.get("password") or "").strip()
if len(password) < 8:
return Response({"password": ["新密码至少 8 位"]}, status=status.HTTP_400_BAD_REQUEST)
user.set_password(password)
user.save(update_fields=["password"])
Token.objects.filter(user=user).delete()
log_admin_action(
request,
"user.reset_password",
target_type="user",
target_id=user.id,
target_name=user.username,
)
return Response(status=status.HTTP_204_NO_CONTENT)
# ─────────────────────────── 质量词(平台单层)───────────────────────────
@api_view(["GET", "POST"])
@permission_classes([IsPlatformAdmin])
def admin_quality_words(request):
"""GET 列全部质量词(配置量小,不分页,前端按 stage 分组);POST 新增一条。"""
if request.method == "GET":
qs = QualityWord.objects.all().order_by("stage", "sort", "created_at")
return Response(QualityWordSerializer(qs, many=True).data)
serializer = QualityWordSerializer(data=request.data)
serializer.is_valid(raise_exception=True)
obj = serializer.save()
log_admin_action(
request,
"quality_word.create",
target_type="quality_word",
target_id=obj.id,
target_name=f"{obj.stage}:{obj.text}",
after={"stage": obj.stage, "text": obj.text},
)
return Response(QualityWordSerializer(obj).data, status=status.HTTP_201_CREATED)
@api_view(["PATCH", "DELETE"])
@permission_classes([IsPlatformAdmin])
def admin_quality_word_detail(request, word_id):
obj = QualityWord.objects.filter(id=word_id).first()
if obj is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
if request.method == "DELETE":
log_admin_action(
request,
"quality_word.delete",
target_type="quality_word",
target_id=obj.id,
target_name=f"{obj.stage}:{obj.text}",
before={"stage": obj.stage, "text": obj.text},
)
obj.delete()
return Response(status=status.HTTP_204_NO_CONTENT)
before_text = obj.text
serializer = QualityWordSerializer(obj, data=request.data, partial=True)
serializer.is_valid(raise_exception=True)
serializer.save()
log_admin_action(
request,
"quality_word.update",
target_type="quality_word",
target_id=obj.id,
target_name=f"{obj.stage}:{obj.text}",
before={"text": before_text},
after={"text": obj.text},
)
return Response(QualityWordSerializer(obj).data)
# ─────────────────────────── 生图/视频 提示词模板 ───────────────────────────
@api_view(["GET"])
@permission_classes([IsPlatformAdmin])
def admin_prompt_templates(request):
"""列视频线 6 条提示词模板(固定 key,数量小不分页,按 key 排序)。"""
# 故事板已下线,分镜图提示词不再出现在后台可编辑列表里(历史行留在库里不动)
qs = PromptTemplate.objects.exclude(key=PromptTemplate.Key.STORYBOARD_FRAME).order_by("key")
return Response(PromptTemplateSerializer(qs, many=True).data)
@api_view(["PATCH"])
@permission_classes([IsPlatformAdmin])
def admin_prompt_template_detail(request, template_id):
"""改一条提示词模板的正文 / 比例 / 启用(key 固定不可改,不增不删)。"""
obj = PromptTemplate.objects.filter(id=template_id).first()
if obj is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
before = {"template": obj.template, "ratio": obj.ratio, "enabled": obj.enabled}
serializer = PromptTemplateSerializer(obj, data=request.data, partial=True)
serializer.is_valid(raise_exception=True)
serializer.save()
log_admin_action(
request,
"prompt_template.update",
target_type="prompt_template",
target_id=obj.id,
target_name=obj.key,
before=before,
after={"template": obj.template, "ratio": obj.ratio, "enabled": obj.enabled},
)
return Response(PromptTemplateSerializer(obj).data)
# ─────────────────────────── 火山人像审核队列 ───────────────────────────
_REVIEW_STATUSES = {"", "processing", "active", "failed"}
@api_view(["GET"])
@permission_classes([IsPlatformAdmin])
def admin_asset_reviews(request):
"""跨团队送审资产审核队列(角色定妆照/三视图/分镜图)。?review_status=none|processing|active|failed 过滤(none=未送审)。"""
qs = (
Asset.objects.filter(category__in=Asset.REVIEW_CATEGORIES, is_deleted=False)
.select_related("team")
.prefetch_related("files")
.order_by("-created_at")
)
rs = request.query_params.get("review_status")
if rs == "none":
qs = qs.filter(review_status="")
elif rs in {"processing", "active", "failed"}:
qs = qs.filter(review_status=rs)
search = (request.query_params.get("search") or "").strip()
if search:
qs = qs.filter(team__name__icontains=search)
paginator = DefaultPagination()
page = paginator.paginate_queryset(qs, request)
return paginator.get_paginated_response(AdminReviewAssetSerializer(page, many=True).data)
@api_view(["POST"])
@permission_classes([IsPlatformAdmin])
def admin_asset_reviews_submit(request):
"""批量送审(也用于失败重试):对给定 person 资产逐个 submit_asset_for_review。"""
ids = request.data.get("asset_ids") or []
assets = list(Asset.objects.filter(id__in=ids, category__in=Asset.REVIEW_CATEGORIES, is_deleted=False))
for asset in assets:
submit_asset_for_review(asset)
statuses = {str(a.id): a.review_status for a in Asset.objects.filter(id__in=ids)}
log_admin_action(
request,
"asset_review.submit",
target_type="asset",
target_name=f"{len(assets)} assets",
after={"count": len(assets)},
)
return Response({"submitted": len(assets), "statuses": statuses})
@api_view(["POST"])
@permission_classes([IsPlatformAdmin])
def admin_asset_reviews_poll(request):
"""自动送审未送审资产 + 轮询审核中状态。给 asset_ids 则只轮询这些,否则全平台兜底。"""
from apps.assets.review import drain_platform_reviews
ids = request.data.get("asset_ids")
if ids:
qs = Asset.objects.filter(
category__in=Asset.REVIEW_CATEGORIES, review_status="processing", is_deleted=False, id__in=ids
)
statuses = {}
for asset in qs:
statuses[str(asset.id)] = poll_asset_review(asset)
return Response({"polled": len(statuses), "submitted": 0, "statuses": statuses})
result = drain_platform_reviews()
return Response(result)
def _refresh_inflight_task(task: AITask, *, operator) -> AITask:
"""对单条在途任务向供应商拉一次最新态。无远端 ID / 非视频异步任务则原样返回。"""
if not task.provider_task_id:
return task
if task.task_type == AITask.Type.FREE_VIDEO:
from apps.ai.free_video import finalize_free_video
return finalize_free_video(task=task)
if task.task_type == AITask.Type.VIDEO_SEGMENT:
from apps.ai.services import poll_video_segment
from apps.projects.models import VideoSegment
segment_id = (task.request_payload or {}).get("video_segment_id")
if not segment_id:
return task
segment = VideoSegment.objects.filter(id=segment_id).first()
if segment is None:
return task
user = task.created_by or getattr(task.team, "owner", None) or operator
poll_video_segment(video_segment=segment, user=user)
task.refresh_from_db()
return task
return task
# ─────────────────────────── AI 任务监控 + 成本异常 ───────────────────────────
@api_view(["GET"])
@permission_classes([IsPlatformAdmin])
def admin_tasks(request):
"""全局 AITask 列表(?status= / ?task_type= / ?team= / ?anomaly=1 成本异常 筛 + 分页)。"""
# 列表不返回请求/响应/完整错误正文。部分图片、视频任务的 JSON 可达数 MB,若随列表页
# 从远程 MySQL 读取,会让只有 10 行的分页请求也长时间卡在“加载中”;详情接口仍完整读取。
qs = (
AITask.objects.select_related("team", "model_config")
.defer("request_payload", "response_payload", "error_message")
# task_type 仍用于调度;用请求来源区分「全能创作」和自由视频/图片,不读取整份 Prompt。
.annotate(
task_category=Case(
When(request_payload__feature="omni_create", then=Value("omni_create")),
default=Value("standard"),
output_field=CharField(),
)
)
.order_by("-created_at")
)
st = request.query_params.get("status")
if st in {"generating", "running", "in_flight"}:
qs = qs.filter(status__in=_ADMIN_TASK_INFLIGHT_STATUSES)
elif st in dict(AITask.Status.choices):
qs = qs.filter(status=st)
tt = request.query_params.get("task_type")
if tt in dict(AITask.Type.choices):
qs = qs.filter(task_type=tt)
category = request.query_params.get("category")
if category == "omni_create":
qs = qs.filter(request_payload__feature="omni_create")
team_id = request.query_params.get("team")
if team_id:
qs = qs.filter(team_id=team_id)
if request.query_params.get("anomaly") in {"1", "true"}:
qs = qs.filter(estimated_cost__gt=0, actual_cost__gt=F("estimated_cost") * COST_ANOMALY_RATIO)
paginator = DefaultPagination()
page = paginator.paginate_queryset(qs, request)
return paginator.get_paginated_response(AdminTaskSerializer(page, many=True).data)
@api_view(["POST"])
@permission_classes([IsPlatformAdmin])
def admin_tasks_poll(request):
"""向供应商刷新在途视频任务状态。给 task_ids 则只拉这些,否则拉全平台 submitted/polling。
图片/脚本等同步任务没有远端 task id,点刷新后前端会再拉一次列表,worker 已落库的终态会一并更新。
"""
ids = request.data.get("task_ids")
qs = (
AITask.objects.select_related("team", "model_config", "model_config__provider", "project")
.filter(status__in=_ADMIN_TASK_POLL_STATUSES)
.exclude(provider_task_id="")
.order_by("-updated_at")
)
if ids:
qs = qs.filter(id__in=ids)
statuses = {}
for task in qs[:_ADMIN_TASK_POLL_LIMIT]:
try:
refreshed = _refresh_inflight_task(task, operator=request.user)
statuses[str(refreshed.id)] = refreshed.status
except Exception: # noqa: BLE001 — 单条失败不阻断整页刷新
logger.warning("admin poll task %s failed", task.id, exc_info=True)
statuses[str(task.id)] = task.status
return Response({"polled": len(statuses), "statuses": statuses})
@api_view(["GET"])
@permission_classes([IsPlatformAdmin])
def admin_task_detail(request, task_id):
# 尝试链只在详情请求加载,列表仍保持原查询与一任务一行。
task = (
AITask.objects.select_related("team", "model_config")
.prefetch_related("model_attempts")
.filter(id=task_id)
.first()
)
if task is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
return Response(AdminTaskDetailSerializer(task).data)
@api_view(["POST"])
@permission_classes([IsPlatformAdmin])
def admin_task_retry(request, task_id):
"""失败任务重投(best-effort):仅 FAILED + 可重投的图像类(基础资产 / 独立生图)。
其余类型(脚本 / 故事板 / 视频 / 配音 / 导出)请到对应页面重跑,这里 400 不冒险误投。"""
task = AITask.objects.filter(id=task_id).first()
if task is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
if task.status != AITask.Status.FAILED:
return Response({"detail": "仅失败任务可重投"}, status=status.HTTP_400_BAD_REQUEST)
image_types = {AITask.Type.PRODUCT_IMAGE, AITask.Type.PERSON_IMAGE, AITask.Type.SCENE_IMAGE}
if task.task_type not in image_types:
return Response({"detail": "该任务类型暂不支持后台重投,请到对应页面重新生成"}, status=status.HTTP_400_BAD_REQUEST)
from apps.ai.tasks import generate_base_asset_task, generate_standalone_image_task
try:
if task.project_id:
generate_base_asset_task.delay(str(task.id))
else:
generate_standalone_image_task.delay(str(task.id))
except Exception: # noqa: BLE001 — 投递失败如实返回,不静默
return Response({"detail": "重投调度失败,请确认 worker 在线"}, status=status.HTTP_502_BAD_GATEWAY)
log_admin_action(
request,
"task.retry",
target_type="ai_task",
target_id=task.id,
target_name=task.task_type,
)
return Response({"retried": True, "task_id": str(task.id)})
@api_view(["POST"])
@permission_classes([IsPlatformAdmin])
def admin_task_reap(request, task_id):
"""手动回收僵尸任务:标失败 + 退还预留积分(与自动回收 _reap_stale_standalone_image_tasks 同款账务路径)。
自动回收只在「同团队下次提交」时顺手触发——团队从此不再用该功能,冻结积分就永远躺着;
这里给超管一个兜底出口。范围收紧到 RESERVED(worker 从未认领):SUBMITTED/POLLING 的视频
任务可能正在 ARK 生成(合法耗时 5-10 分钟+),后台强杀会与轮询结算竞态,不冒险。"""
from datetime import timedelta
from django.db import transaction
from django.core.exceptions import ObjectDoesNotExist
from django.utils import timezone
from apps.billing.services.ledger import release_credit
with transaction.atomic():
# 锁行后再判状态:worker 可能恰好并发认领,凭陈旧快照回收会把在途任务误杀
task = AITask.objects.select_for_update().filter(id=task_id).first()
if task is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
if task.status != AITask.Status.RESERVED:
return Response({"detail": "仅能回收卡在「已预留」状态的任务"}, status=status.HTTP_400_BAD_REQUEST)
if task.updated_at >= timezone.now() - timedelta(minutes=10):
return Response({"detail": "任务仍在 10 分钟活跃窗口内,可能只是在排队,稍后再试"}, status=status.HTTP_400_BAD_REQUEST)
task.status = AITask.Status.FAILED
task.error_message = "管理员手动回收(预留积分已退还)"
task.completed_at = timezone.now()
task.save(update_fields=["status", "error_message", "completed_at", "updated_at"])
try:
reservation = task.credit_reservation
except ObjectDoesNotExist:
reservation = None
if reservation is not None:
release_credit(reservation=reservation, reason="管理员手动回收僵尸任务")
log_admin_action(
request,
"task.reap",
target_type="ai_task",
target_id=task.id,
target_name=task.task_type,
)
return Response(AdminTaskSerializer(task).data)
# ─────────────────────────── 计费审计 + 额度策略 ───────────────────────────
@api_view(["GET"])
@permission_classes([IsPlatformAdmin])
def admin_ledgers(request):
"""全局信用流水浏览(?ledger_type= / ?team= 筛 + 分页)。"""
qs = CreditLedger.objects.select_related("team", "user").order_by("-created_at")
lt = request.query_params.get("ledger_type")
if lt in dict(CreditLedger.Type.choices):
qs = qs.filter(ledger_type=lt)
team_id = request.query_params.get("team")
if team_id:
qs = qs.filter(team_id=team_id)
paginator = DefaultPagination()
page = paginator.paginate_queryset(qs, request)
return paginator.get_paginated_response(AdminLedgerSerializer(page, many=True).data)
@api_view(["POST"])
@permission_classes([IsPlatformAdmin])
def admin_ledger_adjust(request):
"""手动调额(争议补偿):{team, amount, reason}。amount 可正可负,落 ADJUSTMENT 流水。"""
from decimal import InvalidOperation
team_id = request.data.get("team") or request.data.get("team_id")
team = Team.objects.filter(id=team_id).first()
if team is None:
return Response({"detail": "团队不存在"}, status=status.HTTP_404_NOT_FOUND)
try:
amount = Decimal(str(request.data.get("amount")))
except (InvalidOperation, TypeError, ValueError):
return Response({"amount": ["金额格式不正确"]}, status=status.HTTP_400_BAD_REQUEST)
if amount == 0:
return Response({"amount": ["调额金额不能为 0"]}, status=status.HTTP_400_BAD_REQUEST)
reason = str(request.data.get("reason") or "").strip()
try:
ledger = adjust_credit(team=team, amount=amount, reason=reason, operator=request.user)
except ValueError as exc:
return Response({"detail": str(exc)}, status=status.HTTP_400_BAD_REQUEST)
log_admin_action(
request,
"credit.adjust",
target_type="team",
target_id=team.id,
target_name=team.name,
after={"amount": str(amount), "balance_after": str(ledger.balance_after), "reason": reason},
)
return Response(AdminLedgerSerializer(ledger).data, status=status.HTTP_201_CREATED)
def _billing_config_payload(cfg: BillingConfig) -> dict:
return {
"points_per_yuan": str(cfg.points_per_yuan),
"video_margin_multiplier": str(cfg.video_margin_multiplier),
"video_reserve_buffer": str(cfg.video_reserve_buffer),
"updated_at": cfg.updated_at,
}
@api_view(["GET", "PATCH"])
@permission_classes([IsPlatformAdmin])
def admin_billing_config(request):
"""平台计费配置(积分汇率/视频毛利系数/预留 buffer)。PATCH 即刻生效于下一次估价/结算。"""
from decimal import InvalidOperation
cfg = get_billing_config()
if request.method == "GET":
return Response(_billing_config_payload(cfg))
before = _billing_config_payload(cfg)
updates: dict = {}
for field, minimum in (("points_per_yuan", Decimal("0.01")), ("video_margin_multiplier", Decimal("0.01")), ("video_reserve_buffer", Decimal("1"))):
if request.data.get(field) is None:
continue
try:
value = Decimal(str(request.data[field]))
except (InvalidOperation, TypeError, ValueError):
return Response({field: ["数值格式不正确"]}, status=status.HTTP_400_BAD_REQUEST)
# Decimal("NaN") 构造不抛,比较才抛 InvalidOperation → 500(review 确认),显式拦
if not value.is_finite():
return Response({field: ["数值格式不正确"]}, status=status.HTTP_400_BAD_REQUEST)
if value < minimum:
return Response({field: [f"不能小于 {minimum}"]}, status=status.HTTP_400_BAD_REQUEST)
updates[field] = value
if not updates:
return Response({"detail": "没有可更新的字段"}, status=status.HTTP_400_BAD_REQUEST)
for field, value in updates.items():
setattr(cfg, field, value)
cfg.save(update_fields=[*updates.keys(), "updated_at"])
invalidate_billing_config_cache()
after = _billing_config_payload(cfg)
log_admin_action(request, "billing_config.update", target_type="billing_config", target_id=cfg.id, before=before, after=after)
return Response(after)
@api_view(["GET", "POST"])
@permission_classes([IsPlatformAdmin])
def admin_quota_policies(request):
"""GET 列额度策略(?team= 筛);POST 新建。team 必填;user/project 留空=团队级。"""
if request.method == "GET":
qs = QuotaPolicy.objects.select_related("team").order_by("-created_at")
team_id = request.query_params.get("team")
if team_id:
qs = qs.filter(team_id=team_id)
paginator = DefaultPagination()
page = paginator.paginate_queryset(qs, request)
return paginator.get_paginated_response(AdminQuotaPolicySerializer(page, many=True).data)
serializer = AdminQuotaPolicySerializer(data=request.data)
serializer.is_valid(raise_exception=True)
obj = serializer.save()
log_admin_action(
request,
"quota_policy.create",
target_type="quota_policy",
target_id=obj.id,
target_name=str(obj.team_id),
after=serializer.data,
)
return Response(AdminQuotaPolicySerializer(obj).data, status=status.HTTP_201_CREATED)
@api_view(["PATCH", "DELETE"])
@permission_classes([IsPlatformAdmin])
def admin_quota_policy_detail(request, policy_id):
obj = QuotaPolicy.objects.filter(id=policy_id).first()
if obj is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
if request.method == "DELETE":
log_admin_action(request, "quota_policy.delete", target_type="quota_policy", target_id=obj.id, target_name=str(obj.team_id))
obj.delete()
return Response(status=status.HTTP_204_NO_CONTENT)
serializer = AdminQuotaPolicySerializer(obj, data=request.data, partial=True)
serializer.is_valid(raise_exception=True)
serializer.save()
log_admin_action(request, "quota_policy.update", target_type="quota_policy", target_id=obj.id, target_name=str(obj.team_id), after=serializer.data)
return Response(AdminQuotaPolicySerializer(obj).data)
# ─────────────────────────── 模型供应商 / 模型 ───────────────────────────
@api_view(["GET", "POST"])
@permission_classes([IsPlatformAdmin])
def admin_providers(request):
if request.method == "GET":
qs = ModelProvider.objects.annotate(model_count_anno=Count("models", distinct=True)).order_by("created_at")
return Response(AdminModelProviderSerializer(qs, many=True).data)
serializer = AdminModelProviderSerializer(data=request.data)
serializer.is_valid(raise_exception=True)
obj = serializer.save()
log_admin_action(request, "provider.create", target_type="model_provider", target_id=obj.id, target_name=obj.name)
return Response(AdminModelProviderSerializer(obj).data, status=status.HTTP_201_CREATED)
@api_view(["PATCH", "DELETE"])
@permission_classes([IsPlatformAdmin])
def admin_provider_detail(request, provider_id):
obj = ModelProvider.objects.filter(id=provider_id).first()
if obj is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
if request.method == "DELETE":
log_admin_action(request, "provider.delete", target_type="model_provider", target_id=obj.id, target_name=obj.name)
obj.delete()
return Response(status=status.HTTP_204_NO_CONTENT)
serializer = AdminModelProviderSerializer(obj, data=request.data, partial=True)
serializer.is_valid(raise_exception=True)
serializer.save()
log_admin_action(request, "provider.update", target_type="model_provider", target_id=obj.id, target_name=obj.name)
return Response(AdminModelProviderSerializer(ModelProvider.objects.annotate(model_count_anno=Count("models", distinct=True)).get(id=obj.id)).data)
@api_view(["GET", "POST"])
@permission_classes([IsPlatformAdmin])
def admin_models(request):
if request.method == "GET":
qs = ModelConfig.objects.select_related("provider").order_by("-is_default", "status", "capability", "provider__name", "created_at")
prov = request.query_params.get("provider")
if prov:
qs = qs.filter(provider_id=prov)
cap = request.query_params.get("capability")
if cap in dict(ModelConfig.Capability.choices):
qs = qs.filter(capability=cap)
return Response(AdminModelConfigSerializer(qs, many=True).data)
serializer = AdminModelConfigSerializer(data=request.data)
serializer.is_valid(raise_exception=True)
obj = serializer.save()
invalidate_model_catalog_cache()
log_admin_action(request, "model.create", target_type="model_config", target_id=obj.id, target_name=f"{obj.provider_id}:{obj.name}")
return Response(AdminModelConfigSerializer(obj).data, status=status.HTTP_201_CREATED)
@api_view(["PATCH", "DELETE"])
@permission_classes([IsPlatformAdmin])
def admin_model_detail(request, model_id):
obj = ModelConfig.objects.select_related("provider").filter(id=model_id).first()
if obj is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
if request.method == "DELETE":
log_admin_action(request, "model.delete", target_type="model_config", target_id=obj.id, target_name=obj.name)
obj.delete()
invalidate_model_catalog_cache()
return Response(status=status.HTTP_204_NO_CONTENT)
serializer = AdminModelConfigSerializer(obj, data=request.data, partial=True)
serializer.is_valid(raise_exception=True)
serializer.save()
invalidate_model_catalog_cache()
log_admin_action(request, "model.update", target_type="model_config", target_id=obj.id, target_name=obj.name)
return Response(AdminModelConfigSerializer(ModelConfig.objects.select_related("provider").get(id=obj.id)).data)
@api_view(["POST"])
@permission_classes([IsPlatformAdmin])
def admin_model_set_default(request, model_id):
"""把某模型设为其 capability 的默认模型(同 capability 其余清默认)。"""
from django.db import transaction
obj = ModelConfig.objects.filter(id=model_id).first()
if obj is None:
return Response({"detail": "not found"}, status=status.HTTP_404_NOT_FOUND)
with transaction.atomic():
ModelConfig.objects.filter(capability=obj.capability).exclude(id=obj.id).update(is_default=False)
obj.is_default = True
obj.save(update_fields=["is_default", "updated_at"])
invalidate_model_catalog_cache()
log_admin_action(request, "model.set_default", target_type="model_config", target_id=obj.id, target_name=f"{obj.capability}:{obj.name}")
return Response(AdminModelConfigSerializer(ModelConfig.objects.select_related("provider").get(id=obj.id)).data)
# ─────────────────────────── 收尾治理(项目监控 + 数据完整性)───────────────────────────
@api_view(["GET"])
@permission_classes([IsPlatformAdmin])
def admin_projects(request):
"""全局项目流水线监控(?status= / ?search= 筛 + 分页)。"""
qs = Project.objects.select_related("team", "product").order_by("-created_at")
st = request.query_params.get("status")
if st in dict(Project.Status.choices):
qs = qs.filter(status=st)
search = (request.query_params.get("search") or "").strip()
if search:
qs = qs.filter(Q(name__icontains=search) | Q(team__name__icontains=search))
paginator = DefaultPagination()
page = paginator.paginate_queryset(qs, request)
return paginator.get_paginated_response(AdminProjectSerializer(page, many=True).data)
@api_view(["GET"])
@permission_classes([IsPlatformAdmin])
def admin_integrity(request):
"""数据完整性 / 治理体检(只读):统计潜在孤儿 / 异常记录数,供平台超管巡检。"""
from django.utils import timezone
now = timezone.now()
checks = [
{"key": "teams_without_account", "label": "无信用账户的团队", "count": Team.objects.filter(credit_account__isnull=True).count()},
{"key": "products_without_image", "label": "无图商品(无主图且无图册)", "count": Product.objects.filter(cover_asset__isnull=True, images__isnull=True).distinct().count()},
{"key": "expired_pending_invites", "label": "已过期但未标记的邀请码", "count": Invitation.objects.filter(status=Invitation.Status.PENDING, expires_at__lt=now).count()},
{"key": "assets_without_files", "label": "无文件的资产(未删除)", "count": Asset.objects.filter(is_deleted=False, files__isnull=True).distinct().count()},
{"key": "failed_projects", "label": "失败的项目", "count": Project.objects.filter(status=Project.Status.FAILED).count()},
{"key": "stale_reserved_tasks", "label": "卡在 reserved 的任务", "count": AITask.objects.filter(status=AITask.Status.RESERVED).count()},
]
return Response({"checks": checks, "totals": {"teams": Team.objects.count(), "users": User.objects.count(), "projects": Project.objects.count(), "products": Product.objects.count()}})