diff --git a/core/backend/apps/ai/script_agent.py b/core/backend/apps/ai/script_agent.py index 48ec5fc..66e7f37 100644 --- a/core/backend/apps/ai/script_agent.py +++ b/core/backend/apps/ai/script_agent.py @@ -1670,7 +1670,7 @@ def stream_script_agent( ): """生成 SSE 帧字符串的同步生成器,供 StreamingHttpResponse 包裹。 target_index 非空 = 精准只改第 N 镜(读全脚本上下文,后端强制保留其余镜原样)。""" - from apps.ai.services import create_ai_task, stream_routed_text_request + from apps.ai.services import create_ai_task, reap_stale_script_tasks, stream_routed_text_request # 极速成片与专业创作都只产出口播。历史项目、模板或请求里的短剧/Vlog # 仅作兼容读取,绝不能重新进入实际生成链路。 @@ -1724,6 +1724,12 @@ def stream_script_agent( yield _sse({"type": "tool", "id": "analyze", "status": "done"}) task_type = AITask.Type.SCRIPT_OPTIMIZATION if mode == "revise" else AITask.Type.SCRIPT_GENERATION + # 顺手回收本团队被进程重启丢下的僵尸脚本任务(没有定时任务,沿用本项目「提交时顺手清」的做法): + # 存量卡「进行中」的行会在下一次生成时自愈,冻住的预扣积分也一并退回。 + try: + reap_stale_script_tasks(team=project.team) + except Exception: # noqa: BLE001 — 回收失败不该挡住本次生成 + logger.warning("reap stale script tasks failed", exc_info=True) try: task = create_ai_task( project=project, diff --git a/core/backend/apps/ai/services.py b/core/backend/apps/ai/services.py index 76b50a8..d1144cf 100644 --- a/core/backend/apps/ai/services.py +++ b/core/backend/apps/ai/services.py @@ -3409,6 +3409,60 @@ _STANDALONE_TASK_TYPE = { } +# 脚本 agent 的 SSE 流是**在进程内跑的**(专业创作走 web 进程,一键成片走 worker 线程), +# 进程一没(部署滚动更新、OOM、Pod 驱逐),那条流就没了,但 AITask 行还停在 SUBMITTED, +# 预扣的积分也一直冻着 —— 任务监控里就是「进行中」挂几个小时不动。 +# 图片有 _reap_stale_standalone_image_tasks、视频有 video_timeout_recovery,脚本此前是**裸奔**的。 +# 单次脚本生成实测在秒~分钟级,20 分钟仍没终态一律判僵尸(quick_create 的 SCRIPT_TIMEOUT 是 30 分钟, +# 这里留足余量,先于它把行收干净,免得编排端一直等一个永远不会动的任务)。 +SCRIPT_TASK_STALE_AFTER = timedelta(minutes=20) + +_SCRIPT_TASK_TYPES = (AITask.Type.SCRIPT_GENERATION, AITask.Type.SCRIPT_OPTIMIZATION) +_SCRIPT_ACTIVE_STATUSES = ( + AITask.Status.CREATED, + AITask.Status.RESERVED, + AITask.Status.SUBMITTED, + AITask.Status.POLLING, + AITask.Status.POSTPROCESSING, +) + + +def reap_stale_script_tasks(*, team=None, project=None) -> int: + """回收被进程重启丢下的脚本任务:置 FAILED + 退还预扣积分。返回回收条数。 + + 没有 celery beat,沿用本项目既有的「顺手回收」策略:每次新提交脚本、以及一键成片每轮推进时 + 各扫一次。这样存量僵尸行会在下一次生成时自愈,不需要人工进库改数据。""" + if team is None and project is None: + return 0 + cutoff = timezone.now() - SCRIPT_TASK_STALE_AFTER + qs = AITask.objects.filter( + task_type__in=_SCRIPT_TASK_TYPES, + status__in=_SCRIPT_ACTIVE_STATUSES, + updated_at__lt=cutoff, + ) + qs = qs.filter(project=project) if project is not None else qs.filter(team=team) + reaped = 0 + for task in qs: + try: + with transaction.atomic(): + 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="脚本任务僵尸回收") + reaped += 1 + except Exception: # noqa: BLE001 — 回收是兜底,失败不该挡住正常生成 + logger.warning("reap stale script task %s failed", task.id, exc_info=True) + if reaped: + logger.info("reaped %s stale script task(s)", reaped) + return reaped + + def _reap_stale_standalone_image_tasks(*, team) -> None: """兜底:worker 崩溃/重启(OOM、部署)可能留下卡在 RESERVED 的出图任务,额度被一直占住、 前端轮询也永远等不到结果。超过 10 分钟(远大于单张真实出图耗时 ~60s)仍 RESERVED 的判为僵尸: diff --git a/core/backend/apps/ai/test_script_task_reaper.py b/core/backend/apps/ai/test_script_task_reaper.py new file mode 100644 index 0000000..e2ca5ae --- /dev/null +++ b/core/backend/apps/ai/test_script_task_reaper.py @@ -0,0 +1,111 @@ +"""脚本任务僵尸回收。 + +复现的是线上真实故障:11:41 起了一条 script_generation,11:51 部署滚动更新把进程换掉, +到 14:18 任务监控里那一行还是「进行中」、10 积分预扣冻着不动。 + +成因两条,缺一不可: + 1. 脚本 SSE 流跑在进程内,进程没了任务行就再也没人推进 —— 而脚本任务当时没有任何回收器; + 2. quick_create._advance_script 的 30 分钟超时判定写在函数末尾,任务卡在 SUBMITTED 时 + 永远从「在途」分支提前 return,那个上限根本走不到。 +""" +from datetime import timedelta +from decimal import Decimal + +from django.test import TestCase +from django.utils import timezone + +from apps.accounts.models import Team, TeamMember, User +from apps.ai.models import AITask, ModelConfig, ModelProvider +from apps.ai.services import SCRIPT_TASK_STALE_AFTER, reap_stale_script_tasks +from apps.billing.models import CreditAccount +from apps.billing.services.ledger import reserve_credit +from apps.products.models import Product +from apps.projects.models import Project + + +class ScriptTaskReaperTests(TestCase): + def setUp(self): + self.user = User.objects.create_user(username="reaper", password="pass") + self.team = Team.objects.create(name="Reaper Team", owner=self.user) + TeamMember.objects.create(team=self.team, user=self.user, role=TeamMember.Role.OWNER) + self.account = CreditAccount.objects.create(team=self.team, balance="100.0000") + self.product = Product.objects.create(team=self.team, created_by=self.user, title="测试商品") + self.project = Project.objects.create( + team=self.team, created_by=self.user, product=self.product, name="P" + ) + # 迁移里已 seed 了 volcengine 等 provider,这里复用而不是重复创建(name 唯一) + provider, _ = ModelProvider.objects.get_or_create( + name="reaper-provider", defaults={"display_name": "回收测试"} + ) + self.model = ModelConfig.objects.create( + provider=provider, name="reaper-text-model", display_name="豆包", + capability=ModelConfig.Capability.TEXT, + ) + + def _task(self, *, age_minutes, status=AITask.Status.SUBMITTED, reserve="10.0000"): + task = AITask.objects.create( + team=self.team, created_by=self.user, project=self.project, model_config=self.model, + task_type=AITask.Type.SCRIPT_GENERATION, status=status, + idempotency_key=f"key-{age_minutes}-{status}", + ) + if reserve: + reserve_credit(team=self.team, user=self.user, task=task, amount=Decimal(reserve)) + stale_at = timezone.now() - timedelta(minutes=age_minutes) + AITask.objects.filter(id=task.id).update(created_at=stale_at, updated_at=stale_at) + task.refresh_from_db() + return task + + def test_orphaned_task_is_failed_and_credits_are_returned(self): + """这就是线上那一行:提交后进程被换掉,行永远停在 SUBMITTED,预扣积分冻着退不回来。""" + task = self._task(age_minutes=180) + self.assertEqual( + CreditAccount.objects.get(id=self.account.id).reserved_balance, + Decimal("10.0000"), + "预扣应先冻住额度", + ) + + self.assertEqual(reap_stale_script_tasks(team=self.team), 1) + + task.refresh_from_db() + self.assertEqual(task.status, AITask.Status.FAILED) + self.assertIn("僵尸任务自动回收", task.error_message) + self.assertEqual( + CreditAccount.objects.get(id=self.account.id).reserved_balance, + Decimal("0.0000"), + "冻住的预扣必须解冻退回", + ) + + def test_running_task_within_the_window_is_left_alone(self): + """正常在跑的脚本不能被误杀。""" + task = self._task(age_minutes=1) + + self.assertEqual(reap_stale_script_tasks(team=self.team), 0) + + task.refresh_from_db() + self.assertEqual(task.status, AITask.Status.SUBMITTED) + + def test_window_is_shorter_than_quick_create_timeout(self): + """回收窗口必须早于一键成片的 30 分钟超时,否则编排端会一直等一个永远不动的任务。""" + from apps.projects.services.quick_create import SCRIPT_TIMEOUT + + self.assertLess(SCRIPT_TASK_STALE_AFTER, SCRIPT_TIMEOUT) + + def test_finished_tasks_are_never_touched(self): + done = self._task(age_minutes=300, status=AITask.Status.SUCCEEDED, reserve="") + + self.assertEqual(reap_stale_script_tasks(team=self.team), 0) + + done.refresh_from_db() + self.assertEqual(done.status, AITask.Status.SUCCEEDED) + + def test_scoping_by_project_does_not_reap_other_projects(self): + other = Project.objects.create( + team=self.team, created_by=self.user, product=self.product, name="别的项目" + ) + mine = self._task(age_minutes=180) + + self.assertEqual(reap_stale_script_tasks(project=other), 0) + mine.refresh_from_db() + self.assertEqual(mine.status, AITask.Status.SUBMITTED) + + self.assertEqual(reap_stale_script_tasks(project=self.project), 1) diff --git a/core/backend/apps/projects/services/quick_create.py b/core/backend/apps/projects/services/quick_create.py index 70093fb..581afd0 100644 --- a/core/backend/apps/projects/services/quick_create.py +++ b/core/backend/apps/projects/services/quick_create.py @@ -27,6 +27,7 @@ from apps.ai.script_agent import ( ) from apps.ai.services import ( collect_video_review_blockers, + reap_stale_script_tasks, create_export_job, generate_base_asset, generate_person_triview, @@ -550,9 +551,17 @@ def _advance_script(job: QuickCreateJob) -> int | None: _enqueue_script(str(job.id)) return SCRIPT_POLL_SECONDS + # 先回收被进程重启丢下的僵尸脚本任务(部署滚动更新最常见),否则下面会一直命中「在途」分支: + # 任务行永远停在 SUBMITTED、预扣积分冻着,而本函数底部的 SCRIPT_TIMEOUT 判定根本走不到。 + reap_stale_script_tasks(project=job.project) latest = _latest_script_task(job.project) if latest is not None and latest.status in _ACTIVE_TASK_STATUSES: started_at = _parse_iso(metadata.get("script_started_at")) + # ★ 超时判定必须在这里也做一遍。原来它只写在函数末尾,而末尾只有「查不到在途任务」时才够得着 —— + # 任务卡在 SUBMITTED 时永远从这个分支 return,30 分钟上限形同虚设(实测挂了 2.5 小时还在转)。 + if started_at and timezone.now() - started_at > SCRIPT_TIMEOUT: + fail_quick_create(job, "脚本生成超时,请稍后重试或进入专业模式查看") + return None elapsed = int((timezone.now() - started_at).total_seconds()) if started_at else 0 _save_job(job, progress=min(44, 28 + elapsed // 8), message="正在生成分镜脚本…") return SCRIPT_POLL_SECONDS diff --git a/core/backend/apps/projects/test_quick_create.py b/core/backend/apps/projects/test_quick_create.py index 617904e..0a62b2e 100644 --- a/core/backend/apps/projects/test_quick_create.py +++ b/core/backend/apps/projects/test_quick_create.py @@ -584,6 +584,30 @@ class QuickCreateCoordinatorTests(TestCase): self.assertIsNone(delay) self.assertEqual(self.job.status, QuickCreateJob.Status.CANCELLED) + def test_script_stuck_in_flight_does_not_spin_forever(self): + """线上故障复现:11:41 起的脚本任务被部署打断,行永远停在 SUBMITTED, + 到 14:18 任务监控还是「进行中」。原因是「在途」分支提前 return, + 函数末尾那个 30 分钟 SCRIPT_TIMEOUT 判定根本走不到。""" + task = self._script_task(AITask.Status.SUBMITTED, key="stuck-by-deploy") + stale = timezone.now() - timedelta(hours=3) + AITask.objects.filter(id=task.id).update(created_at=stale, updated_at=stale) + self.job.status = QuickCreateJob.Status.RUNNING + self.job.phase = QuickCreateJob.Phase.SCRIPT + self.job.metadata = { + "script_started": True, + "script_started_at": stale.isoformat(), + } + self.job.save(update_fields=["status", "phase", "metadata", "updated_at"]) + + advance_quick_create(str(self.job.id)) + self.job.refresh_from_db() + task.refresh_from_db() + + # 僵尸行必须被收掉(积分解冻),整单落到可恢复的失败态,而不是永远转圈 + self.assertEqual(task.status, AITask.Status.FAILED) + self.assertNotEqual(self.job.status, QuickCreateJob.Status.RUNNING) + self.assertNotEqual(self.job.message, "正在生成分镜脚本…") + @patch("apps.projects.services.quick_create._enqueue_script") def test_recover_requeues_script_when_queue_drops_it(self, enqueue_script): self.job.status = QuickCreateJob.Status.RUNNING