From 4b245b66057d3462a43e3daa8ac93b9a7ef53be4 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Wed, 26 Aug 2026 12:29:31 +0500 Subject: [PATCH] =?UTF-8?q?feat(tradein/estimate):=20=D0=BF=D0=BE=D1=82?= =?UTF-8?q?=D0=BE=D0=BB=D0=BE=D0=BA=20=D0=BE=D0=B4=D0=BD=D0=BE=D0=B2=D1=80?= =?UTF-8?q?=D0=B5=D0=BC=D0=B5=D0=BD=D0=BD=D1=8B=D1=85=20=D0=BE=D1=86=D0=B5?= =?UTF-8?q?=D0=BD=D0=BE=D0=BA=20=E2=80=94=204=20=D1=81=D0=BB=D0=BE=D1=82?= =?UTF-8?q?=D0=B0,=20=D0=B1=D1=8B=D1=81=D1=82=D1=80=D1=8B=D0=B9=20429=20(#?= =?UTF-8?q?3082)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Рейт-лимит меряет частоту (300/60с вправе стартовать в одну секунду), квота — счётная и помесячная: параллелизм /estimate не ограничивал никто. Оценка 0.8–2.4с держит соединение общего пула SQLAlchemy (5+10) и внешние тиры — пила одновременных оценок выедала пул и тормозила весь /api/v1/*. Семафор по образцу public/mera.py::_suggest_slots: acquire после дешёвых отказов (рейт-лимит, квота) с ожиданием 5с ≈ две длительности оценки, timeout → 429 с Retry-After; release в finally сразу после дорогой части. 4+4 слота (estimate+suggest) = 8 удерживаемых соединений из 15 пула. Семафор в памяти процесса — при переходе на несколько воркеров (#3083) лимит умножится на их число; задачи согласовывать (о чём комментарий на месте). Co-Authored-By: Claude Opus 5 --- tradein-mvp/backend/app/api/v1/trade_in.py | 35 +++++ .../test_3082_estimate_concurrency_cap.py | 144 ++++++++++++++++++ 2 files changed, 179 insertions(+) create mode 100644 tradein-mvp/backend/tests/test_3082_estimate_concurrency_cap.py diff --git a/tradein-mvp/backend/app/api/v1/trade_in.py b/tradein-mvp/backend/app/api/v1/trade_in.py index 0c74f0c3..9c6688d0 100644 --- a/tradein-mvp/backend/app/api/v1/trade_in.py +++ b/tradein-mvp/backend/app/api/v1/trade_in.py @@ -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. diff --git a/tradein-mvp/backend/tests/test_3082_estimate_concurrency_cap.py b/tradein-mvp/backend/tests/test_3082_estimate_concurrency_cap.py new file mode 100644 index 00000000..3d389f24 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3082_estimate_concurrency_cap.py @@ -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