f45baa2ce0
实现3层队列防护,防止批量提交导致队列爆炸: **用户级限流(核心)** - 每个用户同时 pending 的 generation 任务上限 3 个 - 超限返回 429:您的待处理任务过多,请等待完成后再提交 **全局限流(兜底)** - 系统 pending 任务超过 20 个一律拒绝 - 返回 503:系统繁忙,请稍后再试 **覆盖的入口** - 一键生成(generation_tasks 批量创建 + 重试) - 剪辑计划生成(edit_plans generate) - 任务中心重试(task_center 用户级 + 项目级) **实现细节** - 预检查 + 入队前检查双重保障 - 限流拒绝时任务标记为 failed,避免 pending 僵尸 - repository 不支持计数时自动降级跳过(兼容旧代码) - 全局检查始终生效,用户级检查需传 user_id **新增** - task_enqueue.py: UserPendingLimitExceeded / GlobalQueueFull 异常 - generation_task_repository: count_pending_by_user / count_pending_total - 13个单元测试,覆盖正常/超限/全局/降级等场景
28 lines
838 B
Python
Executable File
28 lines
838 B
Python
Executable File
from __future__ import annotations
|
|
|
|
from typing import Protocol
|
|
|
|
from packages.domain import GenerationTask
|
|
|
|
|
|
class GenerationTaskRepository(Protocol):
|
|
def create(self, task: GenerationTask) -> GenerationTask: ...
|
|
|
|
def get(self, task_id: str) -> GenerationTask | None: ...
|
|
|
|
def list_by_project(self, project_id: str) -> list[GenerationTask]: ...
|
|
|
|
def list_by_user(self, user_id: str) -> list[GenerationTask]: ...
|
|
|
|
def count_by_user(self, user_id: str) -> int: ...
|
|
|
|
def count_pending_by_user(self, user_id: str) -> int: ...
|
|
|
|
def count_pending_total(self) -> int: ...
|
|
|
|
def list_recent_by_user(self, user_id: str, limit: int = 5) -> list[GenerationTask]: ...
|
|
|
|
def list_by_source_edit_plan(self, plan_id: str) -> list[GenerationTask]: ...
|
|
|
|
def update(self, task: GenerationTask) -> GenerationTask: ...
|