diff --git a/services/api/app/api/documents.py b/services/api/app/api/documents.py index 92aeadc..40dc41f 100644 --- a/services/api/app/api/documents.py +++ b/services/api/app/api/documents.py @@ -11,7 +11,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from app.config import settings from app.core.celery_app import celery_app 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.database import get_db from app.models.task import Task @@ -80,8 +80,7 @@ async def upload_for_plagiarism_check( headers={"Retry-After": "86400"}, ) - slot_acquired = await acquire_concurrent_slot(current_user.id, current_user.plan) - if not slot_acquired: + if not await check_concurrent_limit(db, current_user.id, current_user.plan): raise HTTPException( status_code=status.HTTP_429_TOO_MANY_REQUESTS, detail="Превышен лимит одновременных задач.", diff --git a/services/api/app/api/search.py b/services/api/app/api/search.py index 9d33c79..10c4a34 100644 --- a/services/api/app/api/search.py +++ b/services/api/app/api/search.py @@ -6,7 +6,7 @@ from fastapi import APIRouter, Depends, HTTPException, status from sqlalchemy.ext.asyncio import AsyncSession 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.database import get_db from app.models.task import Task @@ -46,9 +46,8 @@ async def create_search_task( headers={"Retry-After": "86400"}, ) - # Проверяем лимит одновременных задач (атомарно) - slot_acquired = await acquire_concurrent_slot(current_user.id, current_user.plan) - if not slot_acquired: + # Проверяем лимит одновременных задач (по факту в БД) + if not await check_concurrent_limit(db, current_user.id, current_user.plan): raise HTTPException( status_code=status.HTTP_429_TOO_MANY_REQUESTS, detail="Превышен лимит одновременных задач. Дождитесь завершения текущих.", diff --git a/services/api/app/core/rate_limiter.py b/services/api/app/core/rate_limiter.py index abd4449..d3e81c7 100644 --- a/services/api/app/core/rate_limiter.py +++ b/services/api/app/core/rate_limiter.py @@ -1,16 +1,27 @@ -"""Redis-based rate limiter с атомарными Lua-скриптами. +"""Rate limiter: дневные/месячные лимиты — в Redis, лимит одновременных задач — в БД. +Дневные/месячные лимиты (Redis, атомарные Lua-скрипты): Проблема наивного подхода (GET → проверка → INCR): - Race condition: 10 конкурентных запросов могут одновременно пройти GET, увидеть значение ниже лимита и все инкрементировать. - Решение: один Lua-скрипт выполняется атомарно на стороне Redis. Redis гарантирует, что между командами внутри скрипта нет других операций. + +Лимит одновременных задач — НЕ Redis-счётчик. Раньше был acquire/release-счётчик +(INCR при создании задачи, DECR при завершении) — но release_concurrent_slot() +никогда не вызывался ни из одного воркера, так что счётчик только рос и лимит +превышался навсегда (до часового TTL-автосброса), даже когда все задачи юзера +давно завершены. Вместо ручного счётчика, который может рассинхронизироваться +с реальностью, считаем активные задачи прямо по Task.status в Postgres — +рассинхронизации тогда не может быть в принципе. """ import logging from datetime import UTC, datetime +from sqlalchemy import func, select +from sqlalchemy.ext.asyncio import AsyncSession + from app.core.redis_client import get_redis logger = logging.getLogger(__name__) @@ -74,24 +85,6 @@ end 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: now = datetime.now(UTC) 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: - True — слот получен (задачу можно создавать). - False — все слоты заняты. + True — лимит не превышен, задачу можно создавать. """ + from app.models.task import Task # избегаем circular import на уровне модуля + limits = PLAN_LIMITS.get(plan, PLAN_LIMITS["free"]) max_concurrent = limits.get("concurrent", 1) - r = get_redis() - result = await r.eval( - _LUA_ACQUIRE_CONCURRENT, - 1, - f"concurrent:{user_id}", - max_concurrent, - 3_600, # TTL 1 час — автосброс если воркер упал не освободив слот + result = await db.execute( + select(func.count()) + .select_from(Task) + .where(Task.user_id == user_id, Task.status.in_(("queued", "processing"))) ) - return bool(result) - - -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) + current = result.scalar_one() + return current < max_concurrent