Merge pull request 'feat(tradein/estimate): потолок одновременных оценок — 4 слота, быстрый 429 (#3082)' (#3097) from fix/3082-estimate-concurrency into main
All checks were successful
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m42s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-backend (push) Successful in 1m16s
Deploy Trade-In / deploy (push) Successful in 1m25s
Deploy Trade-In / perimeter-smoke (push) Successful in 10s

This commit is contained in:
bot-backend 2026-08-26 07:36:37 +00:00
commit 929eff11d7
2 changed files with 179 additions and 0 deletions

View file

@ -70,6 +70,26 @@ _estimate_limiter = SlidingWindowLimiter(
limit=settings.estimate_rate_limit, window_s=settings.estimate_rate_limit_window_s
)
# ── Потолок одновременных оценок (#3082) ─────────────────────────────────────
#
# Рейт-лимит выше меряет ЧАСТОТУ (300/60с вправе стартовать в одну секунду), а
# квота — счётная и помесячная: ни один из них не ограничивает ПАРАЛЛЕЛИЗМ.
# Оценка 0.82.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.

View file

@ -0,0 +1,144 @@
"""#3082: потолок одновременных оценок на POST /estimate.
Рейт-лимит меряет частоту (300/60с вправе стартовать в одну секунду), квота
счётная и помесячная: параллелизм не ограничивал никто, и пила одновременных
оценок (0.82.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