All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m14s
Сценарный тест ловил одну проводку из 34 — ту, через которую сам и шёл (`geocoder._cache_get`). Мутационный прогон ревьюера: возврат голого `asyncio.to_thread` в 5 из 6 других мест тест НЕ краснит, то есть регресс «кто-то вернул вызов в голый вид» прошёл бы мимо CI в 33 случаях из 34. `test_no_bare_to_thread_over_request_session` читает исходники geocoder и estimator (через `module.__file__`, не по относительному пути — он зависел бы от cwd прогона) и требует нуля живых `asyncio.to_thread(`. Оба модуля сейчас на нуле, поэтому гейт без списка исключений. Фальсификация — голый `to_thread` у `_fetch_anchor_comps` (estimator:4973, сценарным тестом не покрыт): гейт краснеет с номером строки. Второе: защита `run_db_thread` одноразовая — `except asyncio.CancelledError` ловит ОДНУ отмену, вторая вылетает из самого `asyncio.wait([step])`, и поток остаётся сиротой. Живых путей нет (`_with_budget` нигде не вложен, Starlette не отменяет задачу на дисконнекте, uvicorn стартует без `--timeout-graceful-shutdown`), поэтому кода не трогаю — фиксирую инвариант «не вкладывать бюджеты» в докстринге `_with_budget`, чтобы вложение не завезли как безобидное. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
115 lines
5.7 KiB
Python
115 lines
5.7 KiB
Python
"""Отмена геокодинга по бюджету не оставляет ОСИРОТЕВШИЙ поток в сессии запроса (#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)
|
||
)
|