"""Отмена геокодинга по бюджету не оставляет ОСИРОТЕВШИЙ поток в сессии запроса (#3449). `geocode()` живёт под бюджетом 12 с (`estimator._with_budget` = `asyncio.wait_for`), а внутри ходит в БД (`_cache_get`/`_cache_put`/локальные тиры) по ТОЙ ЖЕ `Session`, с которой запрос идёт дальше. `asyncio.to_thread` отменить нельзя: по истечении бюджета снимается только ожидание со стороны loop'а, поток продолжает работать в чужой сессии. Ошибка геокодера при этом никого не трогает (её глушит `_with_budget`) — страдает СЛЕДУЮЩИЙ потребитель сессии, и у `_persist_estimate_and_commit` её не ловит никто: 500 и потерянная оценка клиента. Меряем значение, а не форму: «следующий шаг не вошёл в сессию, пока сирота не закончил». Следующий шаг здесь — ГОЛЫЙ `asyncio.to_thread(db...)`, как персист оценки, а не ещё один защищённый вызов: защита, которая живёт только внутри обёртки, ровно того пострадавшего и не закрывает. """ from __future__ import annotations import asyncio import os import pathlib import threading import time from typing import Any import pytest os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") from app.services import estimator as est from app.services import geocoder as geo # Шаг БД заметно длиннее бюджета — окно, в котором сирота ещё работает. _STEP_S = 0.3 _BUDGET_S = 0.05 class _ConcurrencyProbeSession: """Session-дублёр, который считает ОДНОВРЕМЕННЫЕ входы в сессию.""" def __init__(self) -> None: self._guard = threading.Lock() self._inside = 0 self.conflicts = 0 self.executed = 0 def execute(self, *_a: Any, **_kw: Any) -> Any: with self._guard: self._inside += 1 self.executed += 1 if self._inside > 1: self.conflicts += 1 try: time.sleep(_STEP_S) finally: with self._guard: self._inside -= 1 return self async def test_geocode_budget_cancel_does_not_leave_orphan_in_session( monkeypatch: pytest.MonkeyPatch, ) -> None: db = _ConcurrencyProbeSession() # Настоящий путь `geocode()` до первого шага БД; сам SQL подменён — меряем # владение сессией, а не содержимое кэша. def _slow_cache_get(session: Any, address_norm: str) -> None: session.execute("geocode cache lookup", {"addr": address_norm}) return None monkeypatch.setattr(geo, "_cache_get", _slow_cache_get) degraded = await est._with_budget( geo.geocode("ул. Пушкина, д. 1", db), # type: ignore[arg-type] _BUDGET_S, label="geocode", ) # Бюджет истёк — геокодер деградировал в None, вызывающий идёт дальше. assert degraded is None, "бюджет не сработал — тест ничего не проверил" # Следующий шаг ТОГО ЖЕ запроса по ТОЙ ЖЕ сессии (образец — персист оценки). await asyncio.to_thread(db.execute, "persist estimate") assert db.executed == 2, f"звучали не оба шага (executed={db.executed})" assert db.conflicts == 0, ( "персист вошёл в сессию, пока осиротевший поток геокодера ещё работал в ней: " "два потока в одной Session → «another operation is in progress» на персисте" ) def test_no_bare_to_thread_over_request_session() -> None: """Source-гейт: в геокодере и эстиматоре не осталось голых `asyncio.to_thread(`. Тест выше ловит ОДНУ проводку — ту, через которую идёт сценарий. Остальные 33 (`_cache_put`, `_fetch_anchor_comps`, персист, …) он не видит: возврат любой из них в голый вид прошёл бы мимо CI. Оба модуля сейчас на нуле по живым вызовам, поэтому гейт — ровно «ноль», без списка исключений. Понадобится вызов со СВОЕЙ сессией (как `user_events.record_event`) — заводить его в отдельном модуле или менять этот тест осознанно. Читаем через `module.__file__`: относительный путь зависел бы от cwd прогона. """ for module in (geo, est): src = pathlib.Path(module.__file__ or "").read_text(encoding="utf-8") bare = [ f"{i}: {line.strip()}" for i, line in enumerate(src.splitlines(), 1) if "asyncio.to_thread(" in line and not line.lstrip().startswith("#") ] assert not bare, ( f"{module.__name__}: голый asyncio.to_thread по сессии запроса — " f"отмена оставит сироту в чужой Session (#3449), нужен run_db_thread:\n" + "\n".join(bare) )