feat(tradein/estimate): потолок одновременных оценок — 4 слота, быстрый 429 (#3082) #3097
2 changed files with 179 additions and 0 deletions
|
|
@ -70,6 +70,26 @@ _estimate_limiter = SlidingWindowLimiter(
|
|||
limit=settings.estimate_rate_limit, window_s=settings.estimate_rate_limit_window_s
|
||||
)
|
||||
|
||||
# ── Потолок одновременных оценок (#3082) ─────────────────────────────────────
|
||||
#
|
||||
# Рейт-лимит выше меряет ЧАСТОТУ (300/60с вправе стартовать в одну секунду), а
|
||||
# квота — счётная и помесячная: ни один из них не ограничивает ПАРАЛЛЕЛИЗМ.
|
||||
# Оценка 0.8–2.4с держит соединение общего пула SQLAlchemy (дефолт 5+10) и
|
||||
# внешние тиры; пила одновременных оценок выедает пул и тормозит весь /api/v1/*.
|
||||
# Образец — public/mera.py::_suggest_slots (4 слота на секундное автодополнение).
|
||||
#
|
||||
# 4 слота: вместе с 4 слотами suggest — 8 одновременно удерживаемых соединений
|
||||
# из 15 возможных, остаток пула остаётся прочим ручкам. Ожидание слота 5с ≈ две
|
||||
# длительности оценки: если за это время слот не освободился, очередь глубока и
|
||||
# честный ответ — быстрый 429 с Retry-After, а не растущая очередь (очередь под
|
||||
# нагрузкой — те же занятые соединения плюс таймаут у клиента; mera.py:117-127).
|
||||
#
|
||||
# Семафор живёт в памяти процесса — при переходе на несколько воркеров uvicorn
|
||||
# (#3083) фактический лимит умножится на число воркеров; задачи согласовывать.
|
||||
_ESTIMATE_CONCURRENCY = 4
|
||||
_ESTIMATE_SLOT_WAIT_S = 5.0
|
||||
_estimate_slots = asyncio.Semaphore(_ESTIMATE_CONCURRENCY)
|
||||
|
||||
|
||||
def _resolve_quota_identity(
|
||||
request: Request,
|
||||
|
|
@ -490,6 +510,17 @@ async def estimate(
|
|||
account_quota.check_and_raise(db, quota_key, default_limit=quota_default_limit)
|
||||
from app.services.estimator import estimate_quality
|
||||
|
||||
# #3082: слот одновременности берём ПОСЛЕ дешёвых отказов (рейт-лимит, квота)
|
||||
# и ДО дорогой цепочки внешних вызовов; release — в finally ниже.
|
||||
try:
|
||||
await asyncio.wait_for(_estimate_slots.acquire(), timeout=_ESTIMATE_SLOT_WAIT_S)
|
||||
except TimeoutError:
|
||||
raise HTTPException(
|
||||
status_code=429,
|
||||
detail="Сервис оценки сейчас занят. Попробуйте ещё раз через несколько секунд.",
|
||||
headers={"Retry-After": "5"},
|
||||
) from None
|
||||
|
||||
# #654: ранее любое исключение estimate_quality всплывало необработанным и
|
||||
# маскировалось апстрим-прокси (Caddy) как непрозрачный 502. Ловим, логируем
|
||||
# через logger.exception (→ GlitchTip/Sentry получает stack trace) и отдаём
|
||||
|
|
@ -522,6 +553,10 @@ async def estimate(
|
|||
status_code=503,
|
||||
detail="estimate temporarily unavailable — try again shortly",
|
||||
) from None
|
||||
finally:
|
||||
# #3082: слот возвращаем сразу после дорогой части — инкремент квоты и
|
||||
# сериализация ответа ниже дёшевы и слот держать не должны.
|
||||
_estimate_slots.release()
|
||||
# #747: атомарно-условный инкремент — источник истины по лимиту. check_and_raise
|
||||
# выше остаётся быстрым pre-check (429 до дорогой оценки), но финальное решение
|
||||
# тут: при гонке двух /estimate на used=lim-1 второй получит False.
|
||||
|
|
|
|||
144
tradein-mvp/backend/tests/test_3082_estimate_concurrency_cap.py
Normal file
144
tradein-mvp/backend/tests/test_3082_estimate_concurrency_cap.py
Normal file
|
|
@ -0,0 +1,144 @@
|
|||
"""#3082: потолок одновременных оценок на POST /estimate.
|
||||
|
||||
Рейт-лимит меряет частоту (300/60с вправе стартовать в одну секунду), квота —
|
||||
счётная и помесячная: параллелизм не ограничивал никто, и пила одновременных
|
||||
оценок (0.8–2.4с каждая, соединение общего пула + внешние тиры) выедала пул и
|
||||
тормозила весь /api/v1/*. Фикс — asyncio.Semaphore по образцу public/mera.py.
|
||||
|
||||
Красный на origin/main по ЗНАЧЕНИЮ: без семафора все 5 конкурентных запросов
|
||||
проходят (0×429), с ним пятый получает быстрый 429 «занят». Никаких обращений
|
||||
к новым именам модуля напрямую — monkeypatch констант с raising=False, чтобы
|
||||
на main тест падал ассертом о поведении, а не AttributeError (см.
|
||||
red-must-mean-wrong-value).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
import asyncio
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from unittest.mock import patch
|
||||
from uuid import uuid4
|
||||
|
||||
import pytest
|
||||
from fastapi import FastAPI
|
||||
from httpx import ASGITransport, AsyncClient
|
||||
|
||||
from app.api.v1 import trade_in as trade_in_module
|
||||
from app.core.db import get_db
|
||||
from app.core.ratelimit import SlidingWindowLimiter
|
||||
from app.schemas.trade_in import AggregatedEstimate
|
||||
|
||||
_CONCURRENCY = 4 # прод-значение _ESTIMATE_CONCURRENCY; здесь литералом (см. док-стринг)
|
||||
|
||||
|
||||
def _canned_estimate() -> AggregatedEstimate:
|
||||
return AggregatedEstimate(
|
||||
estimate_id=uuid4(),
|
||||
median_price_rub=5_000_000,
|
||||
range_low_rub=4_500_000,
|
||||
range_high_rub=5_500_000,
|
||||
median_price_per_m2=100_000,
|
||||
confidence="medium",
|
||||
n_analogs=8,
|
||||
period_months=24,
|
||||
analogs=[],
|
||||
actual_deals=[],
|
||||
expires_at=datetime.now(tz=UTC) + timedelta(hours=24),
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture()
|
||||
def app() -> FastAPI:
|
||||
application = FastAPI()
|
||||
application.include_router(trade_in_module.router, prefix="/api/v1/trade-in")
|
||||
|
||||
def _override_db():
|
||||
yield None
|
||||
|
||||
application.dependency_overrides[get_db] = _override_db
|
||||
return application
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _wide_rate_limiter(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Рейт-лимит не должен мешать тесту параллелизма: 100/60с."""
|
||||
monkeypatch.setattr(
|
||||
trade_in_module, "_estimate_limiter", SlidingWindowLimiter(limit=100, window_s=60.0)
|
||||
)
|
||||
|
||||
|
||||
async def test_fifth_concurrent_estimate_gets_fast_429(
|
||||
app: FastAPI, monkeypatch: pytest.MonkeyPatch
|
||||
) -> None:
|
||||
"""4 оценки висят в работе → 5-я не ждёт в очереди, а быстро получает 429
|
||||
с текстом про занятость и Retry-After; после освобождения слотов те же 4
|
||||
завершаются 200 и следующий запрос снова проходит (release в finally)."""
|
||||
# raising=False: на main этих имён нет — тест должен упасть ассертом ниже,
|
||||
# а не AttributeError здесь.
|
||||
monkeypatch.setattr(trade_in_module, "_ESTIMATE_SLOT_WAIT_S", 0.1, raising=False)
|
||||
monkeypatch.setattr(
|
||||
trade_in_module, "_estimate_slots", asyncio.Semaphore(_CONCURRENCY), raising=False
|
||||
)
|
||||
|
||||
gate = asyncio.Event()
|
||||
|
||||
async def _slow_estimate(*args, **kwargs) -> AggregatedEstimate:
|
||||
await gate.wait()
|
||||
return _canned_estimate()
|
||||
|
||||
payload = {"address": "г. Екатеринбург, ул. Малышева, 1", "area_m2": 50.0, "rooms": 2}
|
||||
with (
|
||||
patch("app.services.account_quota.check_and_raise"),
|
||||
patch("app.services.account_quota.increment", return_value=True),
|
||||
patch("app.services.estimator.estimate_quality", new=_slow_estimate),
|
||||
):
|
||||
transport = ASGITransport(app=app)
|
||||
async with AsyncClient(transport=transport, base_url="http://test") as client:
|
||||
holders = [
|
||||
asyncio.create_task(client.post("/api/v1/trade-in/estimate", json=payload))
|
||||
for _ in range(_CONCURRENCY)
|
||||
]
|
||||
# Дать держателям дойти до acquire и занять все слоты.
|
||||
await asyncio.sleep(0.05)
|
||||
|
||||
fifth = await client.post("/api/v1/trade-in/estimate", json=payload)
|
||||
assert fifth.status_code == 429, (
|
||||
f"5-й конкурентный запрос прошёл ({fifth.status_code}) — "
|
||||
"потолка одновременности нет"
|
||||
)
|
||||
assert "занят" in fifth.json()["detail"]
|
||||
assert "Retry-After" in fifth.headers
|
||||
|
||||
gate.set()
|
||||
done = await asyncio.gather(*holders)
|
||||
assert [r.status_code for r in done] == [200] * _CONCURRENCY
|
||||
|
||||
# Слоты вернулись (release в finally) — новый запрос проходит.
|
||||
gate.set()
|
||||
again = await client.post("/api/v1/trade-in/estimate", json=payload)
|
||||
assert again.status_code == 200
|
||||
|
||||
|
||||
async def test_within_limit_behaviour_unchanged(
|
||||
app: FastAPI, monkeypatch: pytest.MonkeyPatch
|
||||
) -> None:
|
||||
"""В пределах лимита семафор прозрачен: одиночный запрос — 200, как раньше."""
|
||||
monkeypatch.setattr(trade_in_module, "_ESTIMATE_SLOT_WAIT_S", 0.1, raising=False)
|
||||
|
||||
async def _fast_estimate(*args, **kwargs) -> AggregatedEstimate:
|
||||
return _canned_estimate()
|
||||
|
||||
payload = {"address": "г. Екатеринбург, ул. Малышева, 1", "area_m2": 50.0, "rooms": 2}
|
||||
with (
|
||||
patch("app.services.account_quota.check_and_raise"),
|
||||
patch("app.services.account_quota.increment", return_value=True),
|
||||
patch("app.services.estimator.estimate_quality", new=_fast_estimate),
|
||||
):
|
||||
transport = ASGITransport(app=app)
|
||||
async with AsyncClient(transport=transport, base_url="http://test") as client:
|
||||
resp = await client.post("/api/v1/trade-in/estimate", json=payload)
|
||||
assert resp.status_code == 200
|
||||
Loading…
Add table
Reference in a new issue