feat(adminpanel): 任务监控「回收退款」手动兜底

新增 POST /api/admin/tasks/<id>/reap/:仅限卡在「已预留」超 10 分钟的僵尸任务(SUBMITTED/
POLLING 故意不许杀,视频合法耗时 5-10 分钟+,强杀会与轮询结算竞态),标失败 + 退还冻结积分,
与自动回收 _reap_stale_standalone_image_tasks 同款账务路径,幂等。

自动回收只在同团队下次提交时顺手触发,团队弃用某功能后冻结积分会永久躺着——这是给超管的
手动出口。序列化器加 reapable 字段(与端点闸同口径),任务监控行级「回收退款」按钮。
4 条测试(成功退款+审计/窗口内拒绝/非 RESERVED 拒绝/非超管 403)。

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
zyc
2026-07-07 10:38:45 +08:00
co-authored by Claude Sonnet 5
parent 299b2c9ecd
commit d466bd60c6
7 changed files with 147 additions and 1 deletions
+14 -1
View File
@@ -118,18 +118,31 @@ class AdminTaskSerializer(serializers.ModelSerializer):
cost_anomaly = serializers.SerializerMethodField()
# 单任务毛利(¥):actual_cost(积分)÷汇率 base_cost。base_cost=0(成本未知)时 None,报表侧过滤
margin_yuan = serializers.SerializerMethodField()
# 可手动回收:卡在 RESERVED 超过 10 分钟(与 admin_task_reap 的服务端闸完全同口径,前端据此显示按钮)
reapable = serializers.SerializerMethodField()
class Meta:
model = AITask
fields = [
"id", "task_type", "status", "team", "team_name", "model_name",
"estimated_cost", "actual_cost", "base_cost", "margin_yuan", "cost_anomaly", "error_code", "created_at",
"estimated_cost", "actual_cost", "base_cost", "margin_yuan", "cost_anomaly", "error_code", "reapable", "created_at",
]
read_only_fields = fields
def get_cost_anomaly(self, obj) -> bool:
return is_cost_anomaly(obj.estimated_cost, obj.actual_cost)
def get_reapable(self, obj) -> bool:
from datetime import timedelta
from django.utils import timezone
return bool(
obj.status == AITask.Status.RESERVED
and obj.updated_at is not None
and obj.updated_at < timezone.now() - timedelta(minutes=10)
)
def get_margin_yuan(self, obj) -> str | None:
base = obj.base_cost or Decimal("0")
actual = obj.actual_cost or Decimal("0")
+61
View File
@@ -401,6 +401,67 @@ class AdminTaskMonitorTests(TestCase):
def test_retry_requires_admin(self):
self.assertEqual(self.nc.post(f"/api/admin/tasks/{self.t_failed.id}/retry/").status_code, 403)
# ── 手动回收僵尸任务(reap):标失败 + 退预留;10 分钟窗口内/非 RESERVED/非超管全拒 ──
def _mk_stuck_reserved(self, key: str, minutes_ago: int = 30):
"""卡死任务工厂:RESERVED + 冻结 20 积分 + 回拨 updated_at(auto_now 只能 queryset.update 绕)。"""
from datetime import timedelta
from django.utils import timezone
from apps.billing.models import CreditAccount
from apps.billing.services.ledger import reserve_credit
CreditAccount.objects.get_or_create(team=self.team, defaults={"balance": Decimal("1000")})
task = self.AITask.objects.create(
team=self.team, model_config=self.mc, task_type=self.AITask.Type.PRODUCT_IMAGE,
status=self.AITask.Status.RESERVED, estimated_cost="20", idempotency_key=key,
)
reserve_credit(team=self.team, user=self.normal, task=task, amount=Decimal("20"))
self.AITask.objects.filter(id=task.id).update(updated_at=timezone.now() - timedelta(minutes=minutes_ago))
task.refresh_from_db()
return task
def test_reap_stuck_reserved_refunds(self):
from apps.billing.models import CreditAccount, CreditReservation
task = self._mk_stuck_reserved("k-reap-stuck")
acct = CreditAccount.objects.get(team=self.team)
self.assertEqual(acct.reserved_balance, Decimal("20"))
# 列表侧 reapable 标记(前端据此显示按钮),与端点闸同口径
row = next(t for t in self.ac.get("/api/admin/tasks/?status=reserved").data["results"] if t["id"] == str(task.id))
self.assertTrue(row["reapable"])
r = self.ac.post(f"/api/admin/tasks/{task.id}/reap/")
self.assertEqual(r.status_code, 200)
self.assertEqual(r.data["status"], "failed")
self.assertFalse(r.data["reapable"])
task.refresh_from_db()
self.assertEqual(task.status, self.AITask.Status.FAILED)
acct.refresh_from_db()
self.assertEqual(acct.reserved_balance, Decimal("0"))
self.assertEqual(task.credit_reservation.status, CreditReservation.Status.RELEASED)
self.assertTrue(AdminAuditLog.objects.filter(action="task.reap").exists())
# 幂等收口:已终态再回收 → 400,不会双退
self.assertEqual(self.ac.post(f"/api/admin/tasks/{task.id}/reap/").status_code, 400)
def test_reap_fresh_reserved_rejected(self):
task = self._mk_stuck_reserved("k-reap-fresh", minutes_ago=0)
r = self.ac.post(f"/api/admin/tasks/{task.id}/reap/")
self.assertEqual(r.status_code, 400)
task.refresh_from_db()
self.assertEqual(task.status, self.AITask.Status.RESERVED)
# 新鲜任务不显示回收按钮
row = next(t for t in self.ac.get("/api/admin/tasks/?status=reserved").data["results"] if t["id"] == str(task.id))
self.assertFalse(row["reapable"])
def test_reap_non_reserved_rejected(self):
self.assertEqual(self.ac.post(f"/api/admin/tasks/{self.t_ok.id}/reap/").status_code, 400)
def test_reap_requires_admin(self):
task = self._mk_stuck_reserved("k-reap-perm")
self.assertEqual(self.nc.post(f"/api/admin/tasks/{task.id}/reap/").status_code, 403)
class AdminBillingTests(TestCase):
"""Phase 7:计费审计(流水浏览/手动调额)+ 4 层额度策略(CRUD + 拦截生效)+ 权限。"""
+2
View File
@@ -18,6 +18,7 @@ from .views import (
admin_quota_policies,
admin_quota_policy_detail,
admin_task_detail,
admin_task_reap,
admin_task_retry,
admin_tasks,
admin_prompt_template_detail,
@@ -55,6 +56,7 @@ urlpatterns = [
path("tasks/", admin_tasks, name="admin-tasks"),
path("tasks/<uuid:task_id>/", admin_task_detail, name="admin-task-detail"),
path("tasks/<uuid:task_id>/retry/", admin_task_retry, name="admin-task-retry"),
path("tasks/<uuid:task_id>/reap/", admin_task_reap, name="admin-task-reap"),
path("ledgers/", admin_ledgers, name="admin-ledgers"),
path("ledgers/adjust/", admin_ledger_adjust, name="admin-ledger-adjust"),
path("billing-config/", admin_billing_config, name="admin-billing-config"),
+46
View File
@@ -478,6 +478,52 @@ def admin_task_retry(request, task_id):
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)
# ─────────────────────────── 计费审计 + 额度策略 ───────────────────────────