fix(api): лимит одновременных задач считать по БД, не по Redis-счётчику
Живая проверка: юзер сделал один поиск (давно завершился, status='done'),
второй поиск сразу упёрся в "Превышен лимит одновременных задач" — на
free-тарифе лимит 1.
Причина: acquire_concurrent_slot() инкрементирует Redis-счётчик
concurrent:{user_id} при создании КАЖДОЙ задачи (search.py, documents.py),
а release_concurrent_slot() — которая должна его декрементировать по
завершении — НЕ ВЫЗЫВАЛАСЬ НИГДЕ В КОДЕ (grep подтвердил: только
определение, ни одного вызова). Счётчик только рос, лимит превышался
навсегда для практически любого юзера после первой же задачи — до
часового TTL-автосброса.
Фикс — не "доставить забытый release()" (это лечит симптом, но оставляет
класс бага: счётчик и реальность могут разойтись любым другим путём), а
убрать сам отдельный счётчик. check_concurrent_limit() считает активные
задачи (status IN queued/processing) напрямую в Postgres — Task.status уже
корректно обновляется во всех воркерах (проверено многократно в этой
сессии), рассинхронизация невозможна по конструкции. Redis-лимиты
(дневные/месячные, Lua-скрипт) не тронуты — там свой, рабочий, механизм.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -11,7 +11,7 @@ from sqlalchemy.ext.asyncio import AsyncSession
|
|||||||
from app.config import settings
|
from app.config import settings
|
||||||
from app.core.celery_app import celery_app
|
from app.core.celery_app import celery_app
|
||||||
from app.core.minio_client import get_minio
|
from app.core.minio_client import get_minio
|
||||||
from app.core.rate_limiter import acquire_concurrent_slot, check_and_increment_limit
|
from app.core.rate_limiter import check_and_increment_limit, check_concurrent_limit
|
||||||
from app.core.security import get_current_verified_user
|
from app.core.security import get_current_verified_user
|
||||||
from app.database import get_db
|
from app.database import get_db
|
||||||
from app.models.task import Task
|
from app.models.task import Task
|
||||||
@@ -80,8 +80,7 @@ async def upload_for_plagiarism_check(
|
|||||||
headers={"Retry-After": "86400"},
|
headers={"Retry-After": "86400"},
|
||||||
)
|
)
|
||||||
|
|
||||||
slot_acquired = await acquire_concurrent_slot(current_user.id, current_user.plan)
|
if not await check_concurrent_limit(db, current_user.id, current_user.plan):
|
||||||
if not slot_acquired:
|
|
||||||
raise HTTPException(
|
raise HTTPException(
|
||||||
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
|
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
|
||||||
detail="Превышен лимит одновременных задач.",
|
detail="Превышен лимит одновременных задач.",
|
||||||
|
|||||||
@@ -6,7 +6,7 @@ from fastapi import APIRouter, Depends, HTTPException, status
|
|||||||
from sqlalchemy.ext.asyncio import AsyncSession
|
from sqlalchemy.ext.asyncio import AsyncSession
|
||||||
|
|
||||||
from app.core.celery_app import celery_app
|
from app.core.celery_app import celery_app
|
||||||
from app.core.rate_limiter import acquire_concurrent_slot, check_and_increment_limit
|
from app.core.rate_limiter import check_and_increment_limit, check_concurrent_limit
|
||||||
from app.core.security import get_current_verified_user
|
from app.core.security import get_current_verified_user
|
||||||
from app.database import get_db
|
from app.database import get_db
|
||||||
from app.models.task import Task
|
from app.models.task import Task
|
||||||
@@ -46,9 +46,8 @@ async def create_search_task(
|
|||||||
headers={"Retry-After": "86400"},
|
headers={"Retry-After": "86400"},
|
||||||
)
|
)
|
||||||
|
|
||||||
# Проверяем лимит одновременных задач (атомарно)
|
# Проверяем лимит одновременных задач (по факту в БД)
|
||||||
slot_acquired = await acquire_concurrent_slot(current_user.id, current_user.plan)
|
if not await check_concurrent_limit(db, current_user.id, current_user.plan):
|
||||||
if not slot_acquired:
|
|
||||||
raise HTTPException(
|
raise HTTPException(
|
||||||
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
|
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
|
||||||
detail="Превышен лимит одновременных задач. Дождитесь завершения текущих.",
|
detail="Превышен лимит одновременных задач. Дождитесь завершения текущих.",
|
||||||
|
|||||||
@@ -1,16 +1,27 @@
|
|||||||
"""Redis-based rate limiter с атомарными Lua-скриптами.
|
"""Rate limiter: дневные/месячные лимиты — в Redis, лимит одновременных задач — в БД.
|
||||||
|
|
||||||
|
Дневные/месячные лимиты (Redis, атомарные Lua-скрипты):
|
||||||
Проблема наивного подхода (GET → проверка → INCR):
|
Проблема наивного подхода (GET → проверка → INCR):
|
||||||
- Race condition: 10 конкурентных запросов могут одновременно пройти GET,
|
- Race condition: 10 конкурентных запросов могут одновременно пройти GET,
|
||||||
увидеть значение ниже лимита и все инкрементировать.
|
увидеть значение ниже лимита и все инкрементировать.
|
||||||
|
|
||||||
Решение: один Lua-скрипт выполняется атомарно на стороне Redis.
|
Решение: один Lua-скрипт выполняется атомарно на стороне Redis.
|
||||||
Redis гарантирует, что между командами внутри скрипта нет других операций.
|
Redis гарантирует, что между командами внутри скрипта нет других операций.
|
||||||
|
|
||||||
|
Лимит одновременных задач — НЕ Redis-счётчик. Раньше был acquire/release-счётчик
|
||||||
|
(INCR при создании задачи, DECR при завершении) — но release_concurrent_slot()
|
||||||
|
никогда не вызывался ни из одного воркера, так что счётчик только рос и лимит
|
||||||
|
превышался навсегда (до часового TTL-автосброса), даже когда все задачи юзера
|
||||||
|
давно завершены. Вместо ручного счётчика, который может рассинхронизироваться
|
||||||
|
с реальностью, считаем активные задачи прямо по Task.status в Postgres —
|
||||||
|
рассинхронизации тогда не может быть в принципе.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
from datetime import UTC, datetime
|
from datetime import UTC, datetime
|
||||||
|
|
||||||
|
from sqlalchemy import func, select
|
||||||
|
from sqlalchemy.ext.asyncio import AsyncSession
|
||||||
|
|
||||||
from app.core.redis_client import get_redis
|
from app.core.redis_client import get_redis
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
@@ -74,24 +85,6 @@ end
|
|||||||
return {new_val, 1}
|
return {new_val, 1}
|
||||||
"""
|
"""
|
||||||
|
|
||||||
# Атомарная проверка + инкремент счётчика одновременных задач.
|
|
||||||
# Возвращает 1 если слот получен, 0 если превышен лимит.
|
|
||||||
_LUA_ACQUIRE_CONCURRENT = """
|
|
||||||
local key = KEYS[1]
|
|
||||||
local limit = tonumber(ARGV[1])
|
|
||||||
local ttl = tonumber(ARGV[2])
|
|
||||||
|
|
||||||
local current = tonumber(redis.call('GET', key) or '0')
|
|
||||||
if current >= limit then
|
|
||||||
return 0
|
|
||||||
end
|
|
||||||
|
|
||||||
redis.call('INCR', key)
|
|
||||||
redis.call('EXPIRE', key, ttl)
|
|
||||||
return 1
|
|
||||||
"""
|
|
||||||
|
|
||||||
|
|
||||||
def _period_suffix(period: str) -> str:
|
def _period_suffix(period: str) -> str:
|
||||||
now = datetime.now(UTC)
|
now = datetime.now(UTC)
|
||||||
return now.strftime("%Y-%m-%d") if period == "day" else now.strftime("%Y-%m")
|
return now.strftime("%Y-%m-%d") if period == "day" else now.strftime("%Y-%m")
|
||||||
@@ -136,33 +129,27 @@ async def check_and_increment_limit(
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
async def acquire_concurrent_slot(user_id: int, plan: str) -> bool:
|
async def check_concurrent_limit(db: AsyncSession, user_id: int, plan: str) -> bool:
|
||||||
"""
|
"""
|
||||||
Атомарно захватить слот одновременной задачи.
|
Проверить лимит одновременных задач по факту в БД (queued/processing).
|
||||||
|
|
||||||
|
Возможен редкий race (два запроса одновременно оба видят N-1 активных и
|
||||||
|
оба проходят) — на практике не критично для этого лимита (защита от
|
||||||
|
злоупотребления, не от превышения на единицу), а взамен исключён класс
|
||||||
|
багов "счётчик разошёлся с реальностью и завис навсегда".
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
True — слот получен (задачу можно создавать).
|
True — лимит не превышен, задачу можно создавать.
|
||||||
False — все слоты заняты.
|
|
||||||
"""
|
"""
|
||||||
|
from app.models.task import Task # избегаем circular import на уровне модуля
|
||||||
|
|
||||||
limits = PLAN_LIMITS.get(plan, PLAN_LIMITS["free"])
|
limits = PLAN_LIMITS.get(plan, PLAN_LIMITS["free"])
|
||||||
max_concurrent = limits.get("concurrent", 1)
|
max_concurrent = limits.get("concurrent", 1)
|
||||||
|
|
||||||
r = get_redis()
|
result = await db.execute(
|
||||||
result = await r.eval(
|
select(func.count())
|
||||||
_LUA_ACQUIRE_CONCURRENT,
|
.select_from(Task)
|
||||||
1,
|
.where(Task.user_id == user_id, Task.status.in_(("queued", "processing")))
|
||||||
f"concurrent:{user_id}",
|
|
||||||
max_concurrent,
|
|
||||||
3_600, # TTL 1 час — автосброс если воркер упал не освободив слот
|
|
||||||
)
|
)
|
||||||
return bool(result)
|
current = result.scalar_one()
|
||||||
|
return current < max_concurrent
|
||||||
|
|
||||||
async def release_concurrent_slot(user_id: int) -> None:
|
|
||||||
"""Освободить слот одновременной задачи после завершения."""
|
|
||||||
r = get_redis()
|
|
||||||
key = f"concurrent:{user_id}"
|
|
||||||
# DECR безопасен: Redis не уходит в отрицательные значения если мы контролируем acquire
|
|
||||||
current = await r.get(key)
|
|
||||||
if current and int(current) > 0:
|
|
||||||
await r.decr(key)
|
|
||||||
|
|||||||
Reference in New Issue
Block a user