Оценка не теряется, когда геокодер не уложился в бюджет (#3449) #3460

Merged
bot-backend merged 2 commits from fix/3449-geocoder-cancel-orphan into main 2026-09-12 10:11:20 +00:00
Collaborator

Оценка клиента больше не теряется из-за того, что геокодер не уложился в свои 12 секунд.

Closes #3449

Что было

asyncio.to_thread отменить нельзя. По истечении бюджета (estimator._with_budget = asyncio.wait_for) снимается только ожидание со стороны event loop'а — поток продолжает работать с ТОЙ ЖЕ Session, что и весь запрос (Depends(get_db)). Вызывающий тем временем идёт дальше: следующий источник, _fetch_anchor_comps, _persist_estimate_and_commit. Два потока в одной Session → «another operation is in progress» / InvalidRequestError на СЛЕДУЮЩЕМ шаге. Страдает не геокодер (его ошибку глушит _with_budget), а следующий потребитель сессии, и у персиста оценки её не ловит никто — 500 и потерянная оценка.

Что сделано

app/core/db.py: run_db_thread(fn, *args, **kwargs) — ТОЛЬКО защита от сироты: ensure_future + shield, в except asyncio.CancelledError дождаться потока (asyncio.wait([step])), прочитать step.exception() (иначе asyncio печатает «Task exception was never retrieved» без контекста) и пробросить отмену.

Commit/rollback в помощник НЕ вынесены умышленно: посреди геокодинга commit зафиксировал бы частичное состояние оценки. estimator._db_step (образец из #3444 — он коммитит ради возврата коннекта в пул перед внешним HTTP) переписан поверх run_db_thread и добавляет свои commit/rollback сам. Поведение прежнее, гейт tests/test_3408_db_step_cancel_orphan.py зелёный.

Заменено 34 вызова, работающих по сессии ЗАПРОСА:

файл сколько что
app/services/geocoder.py 12 _cache_get/_cache_put, геопортал, кадастр (house-match + forward), _local_houses_match, reverse, suggest
app/services/estimator.py 19 _empty_estimate, match_house_readonly, _backfill_house_fias, _lookup_house_facts, 5× _fetch_analogs, _fetch_dkp_corridor, _save_yandex_history_items, _fetch_anchor_comps, _fetch_house_imv_anchor, _lookup_target_quarter_by_coords, _price_from_inputs (ему инжектятся db-резолверы _ratio_resolver/_qi_lookup), _fetch_deals, _persist_estimate_and_commit, _fetch_price_trend, _is_premium_building
app/api/v1/geocode.py 2 _resolve_house_id_by_fias, _lookup_house_facts
app/api/v1/privacy_admin.py 1 erase_person_data

Оставлены голыми (сессия СВОЯ, чужую не держат):

  • app/services/user_events.py:101record_event открывает свой SessionLocal(), декаплён от транзакции запроса by design.
  • app/services/sber_index.py:429,457,518 — сессия СВОЯ на прогон (app/tasks/sber_index_pull.py), чужую транзакцию осиротевший поток не портит; upsert идемпотентный (ON CONFLICT) и повторится на следующем такте. Поправка к первой редакции этого описания: отменяющий там ЕСТЬ — run-задачи спавнятся детачед (spawn_tracked, scraper-kit scheduler), и на teardown asyncio.run() отменяет оставшиеся. Вывод «оставить голым» не меняется, обоснование было неверным.

Проверка

Гейт по значению — tests/test_3449_geocoder_cancel_orphan.py: _with_budget(geocode(...), 0.05) при шаге БД в 0.3 с, следом ГОЛЫЙ asyncio.to_thread(db.execute, "persist estimate") (образец персиста, а не ещё один защищённый вызов — защита, живущая только внутри обёртки, ровно того пострадавшего и не закрывает). Сессия-дублёр считает ОДНОВРЕМЕННЫЕ входы.

Фальсификация: в geocoder._geocode_resolve временно возвращён голый asyncio.to_thread(_cache_get, db, addr_norm) (без git stash — руками, рядом работают другие агенты), тест краснеет:

>       assert db.conflicts == 0, (
            "персист вошёл в сессию, пока осиротевший поток геокодера ещё работал в ней: "
            "два потока в одной Session → «another operation is in progress» на персисте"
        )
E       AssertionError: персист вошёл в сессию, пока осиротевший поток геокодера ещё работал в ней: два потока в одной Session → «another operation is in progress» на персисте
E       assert 1 == 0
E        +  where 1 = <tests.test_3449_geocoder_cancel_orphan._ConcurrencyProbeSession object at 0x1124fc830>.conflicts
tests/test_3449_geocoder_cancel_orphan.py:85: AssertionError
1 failed in 1.04s

Правка возвращена, тест зелёный.

Второй гейт (коммит be9aa2f9, по ревью): сценарный тест выше ловит ОДНУ проводку из 34 — ту, через которую сам и идёт; мутационный прогон ревьюера показал, что возврат голого to_thread в 5 из 6 других мест он не краснит. test_no_bare_to_thread_over_request_session читает исходники обоих модулей (через module.__file__) и требует нуля живых asyncio.to_thread(. Фальсификация — голый to_thread у _fetch_anchor_comps (estimator.py:4973, сценарным тестом НЕ покрыт):

E           AssertionError: app.services.estimator: голый asyncio.to_thread по сессии запроса — отмена оставит сироту в чужой Session (#3449), нужен run_db_thread:
E             4973: _anchor_comps_pre, _anchor_tier_pre = await asyncio.to_thread(
E           assert not ['4973: _anchor_comps_pre, _anchor_tier_pre = await asyncio.to_thread(']
tests/test_3449_geocoder_cancel_orphan.py:111: AssertionError
1 failed, 1 passed in 1.27s

Там же — предупреждение в докстринге _with_budget: защита run_db_thread одноразовая (except asyncio.CancelledError ловит ОДНУ отмену; вторая, прилетевшая во время asyncio.wait([step]), вылетает из самого ожидания, и поток остаётся сиротой). Живых путей нет — _with_budget нигде не вложен, Starlette не отменяет задачу на дисконнекте, uvicorn стартует без --timeout-graceful-shutdown — поэтому код не трогал, зафиксировал инвариант «не вкладывать бюджеты».

Прогоны в tradein-mvp/backend:

  • uv run python -m pytest tests/ -q5939 passed, 35 skipped за 148.72 с (rc=0; было 5938 до source-гейта).
  • uv run ruff check app tests → All checks passed.
  • uv run ruff format --check app tests → изменённые файлы отформатированы; 15 неформатных файлов в репозитории были такими ДО этой ветки и не тронуты.

Чего эта правка НЕ закрывает

services/house_metadata.py: get_house_metadata (он тоже под бюджетом, estimate_house_meta_timeout_s) ходит в БД СИНХРОННО прямо на loop'е, без to_thread вообще. Сироты там нет по построению, но есть блокировка loop'а на время чекаута коннекта — это класс #3408 п.1, отдельная работа, в эту ветку не тащу.

Приёмка из issue (ноль another operation is in progress / InvalidRequestError в логах tradein-backend за сутки под нагрузкой) проверяется ПОСЛЕ деплоя — здесь не заявляю.

Оценка клиента больше не теряется из-за того, что геокодер не уложился в свои 12 секунд. Closes #3449 ## Что было `asyncio.to_thread` отменить нельзя. По истечении бюджета (`estimator._with_budget` = `asyncio.wait_for`) снимается только ожидание со стороны event loop'а — поток продолжает работать с ТОЙ ЖЕ `Session`, что и весь запрос (`Depends(get_db)`). Вызывающий тем временем идёт дальше: следующий источник, `_fetch_anchor_comps`, `_persist_estimate_and_commit`. Два потока в одной `Session` → «another operation is in progress» / `InvalidRequestError` на СЛЕДУЮЩЕМ шаге. Страдает не геокодер (его ошибку глушит `_with_budget`), а следующий потребитель сессии, и у персиста оценки её не ловит никто — 500 и потерянная оценка. ## Что сделано `app/core/db.py: run_db_thread(fn, *args, **kwargs)` — ТОЛЬКО защита от сироты: `ensure_future` + `shield`, в `except asyncio.CancelledError` дождаться потока (`asyncio.wait([step])`), прочитать `step.exception()` (иначе asyncio печатает «Task exception was never retrieved» без контекста) и пробросить отмену. Commit/rollback в помощник НЕ вынесены умышленно: посреди геокодинга commit зафиксировал бы частичное состояние оценки. `estimator._db_step` (образец из #3444 — он коммитит ради возврата коннекта в пул перед внешним HTTP) переписан поверх `run_db_thread` и добавляет свои commit/rollback сам. Поведение прежнее, гейт `tests/test_3408_db_step_cancel_orphan.py` зелёный. Заменено 34 вызова, работающих по сессии ЗАПРОСА: | файл | сколько | что | |---|---|---| | `app/services/geocoder.py` | 12 | `_cache_get`/`_cache_put`, геопортал, кадастр (house-match + forward), `_local_houses_match`, reverse, suggest | | `app/services/estimator.py` | 19 | `_empty_estimate`, `match_house_readonly`, `_backfill_house_fias`, `_lookup_house_facts`, 5× `_fetch_analogs`, `_fetch_dkp_corridor`, `_save_yandex_history_items`, `_fetch_anchor_comps`, `_fetch_house_imv_anchor`, `_lookup_target_quarter_by_coords`, `_price_from_inputs` (ему инжектятся db-резолверы `_ratio_resolver`/`_qi_lookup`), `_fetch_deals`, `_persist_estimate_and_commit`, `_fetch_price_trend`, `_is_premium_building` | | `app/api/v1/geocode.py` | 2 | `_resolve_house_id_by_fias`, `_lookup_house_facts` | | `app/api/v1/privacy_admin.py` | 1 | `erase_person_data` | Оставлены голыми (сессия СВОЯ, чужую не держат): - `app/services/user_events.py:101` — `record_event` открывает свой `SessionLocal()`, декаплён от транзакции запроса by design. - `app/services/sber_index.py:429,457,518` — сессия СВОЯ на прогон (`app/tasks/sber_index_pull.py`), чужую транзакцию осиротевший поток не портит; upsert идемпотентный (`ON CONFLICT`) и повторится на следующем такте. Поправка к первой редакции этого описания: отменяющий там ЕСТЬ — run-задачи спавнятся детачед (`spawn_tracked`, scraper-kit scheduler), и на teardown `asyncio.run()` отменяет оставшиеся. Вывод «оставить голым» не меняется, обоснование было неверным. ## Проверка Гейт по значению — `tests/test_3449_geocoder_cancel_orphan.py`: `_with_budget(geocode(...), 0.05)` при шаге БД в 0.3 с, следом ГОЛЫЙ `asyncio.to_thread(db.execute, "persist estimate")` (образец персиста, а не ещё один защищённый вызов — защита, живущая только внутри обёртки, ровно того пострадавшего и не закрывает). Сессия-дублёр считает ОДНОВРЕМЕННЫЕ входы. Фальсификация: в `geocoder._geocode_resolve` временно возвращён голый `asyncio.to_thread(_cache_get, db, addr_norm)` (без `git stash` — руками, рядом работают другие агенты), тест краснеет: ``` > assert db.conflicts == 0, ( "персист вошёл в сессию, пока осиротевший поток геокодера ещё работал в ней: " "два потока в одной Session → «another operation is in progress» на персисте" ) E AssertionError: персист вошёл в сессию, пока осиротевший поток геокодера ещё работал в ней: два потока в одной Session → «another operation is in progress» на персисте E assert 1 == 0 E + where 1 = <tests.test_3449_geocoder_cancel_orphan._ConcurrencyProbeSession object at 0x1124fc830>.conflicts tests/test_3449_geocoder_cancel_orphan.py:85: AssertionError 1 failed in 1.04s ``` Правка возвращена, тест зелёный. Второй гейт (коммит `be9aa2f9`, по ревью): сценарный тест выше ловит ОДНУ проводку из 34 — ту, через которую сам и идёт; мутационный прогон ревьюера показал, что возврат голого `to_thread` в 5 из 6 других мест он не краснит. `test_no_bare_to_thread_over_request_session` читает исходники обоих модулей (через `module.__file__`) и требует нуля живых `asyncio.to_thread(`. Фальсификация — голый `to_thread` у `_fetch_anchor_comps` (`estimator.py:4973`, сценарным тестом НЕ покрыт): ``` E AssertionError: app.services.estimator: голый asyncio.to_thread по сессии запроса — отмена оставит сироту в чужой Session (#3449), нужен run_db_thread: E 4973: _anchor_comps_pre, _anchor_tier_pre = await asyncio.to_thread( E assert not ['4973: _anchor_comps_pre, _anchor_tier_pre = await asyncio.to_thread('] tests/test_3449_geocoder_cancel_orphan.py:111: AssertionError 1 failed, 1 passed in 1.27s ``` Там же — предупреждение в докстринге `_with_budget`: защита `run_db_thread` одноразовая (`except asyncio.CancelledError` ловит ОДНУ отмену; вторая, прилетевшая во время `asyncio.wait([step])`, вылетает из самого ожидания, и поток остаётся сиротой). Живых путей нет — `_with_budget` нигде не вложен, Starlette не отменяет задачу на дисконнекте, uvicorn стартует без `--timeout-graceful-shutdown` — поэтому код не трогал, зафиксировал инвариант «не вкладывать бюджеты». Прогоны в `tradein-mvp/backend`: - `uv run python -m pytest tests/ -q` → **5939 passed, 35 skipped** за 148.72 с (rc=0; было 5938 до source-гейта). - `uv run ruff check app tests` → All checks passed. - `uv run ruff format --check app tests` → изменённые файлы отформатированы; 15 неформатных файлов в репозитории были такими ДО этой ветки и не тронуты. ## Чего эта правка НЕ закрывает `services/house_metadata.py: get_house_metadata` (он тоже под бюджетом, `estimate_house_meta_timeout_s`) ходит в БД СИНХРОННО прямо на loop'е, без `to_thread` вообще. Сироты там нет по построению, но есть блокировка loop'а на время чекаута коннекта — это класс #3408 п.1, отдельная работа, в эту ветку не тащу. Приёмка из issue (ноль `another operation is in progress` / `InvalidRequestError` в логах `tradein-backend` за сутки под нагрузкой) проверяется ПОСЛЕ деплоя — здесь не заявляю.
bot-backend added 1 commit 2026-09-12 08:23:38 +00:00
Отмена по бюджету больше не оставляет сироту в сессии запроса (#3449)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / changes (pull_request) Successful in 10s
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 5m24s
c251c02f1e
`asyncio.to_thread` отменить нельзя: по истечении бюджета (`_with_budget` =
`asyncio.wait_for`, у геокодера 12 с) снимается только ожидание со стороны
loop'а — поток продолжает работать с ТОЙ ЖЕ `Session`, что и весь запрос.
Вызывающий тем временем идёт дальше: следующий источник, `_fetch_anchor_comps`,
`_persist_estimate_and_commit`. Два потока в одной `Session` дают «another
operation is in progress» / InvalidRequestError на СЛЕДУЮЩЕМ шаге. У источников
эту ошибку глушит `except` вокруг вызова, у персиста оценки не глушит никто —
500 и потерянная оценка клиента.

`app/core/db.py: run_db_thread` — ТОЛЬКО защита от сироты: `ensure_future` +
`shield`, на отмене дождаться потока (`asyncio.wait`), прочитать
`step.exception()` (иначе asyncio печатает «Task exception was never
retrieved» без контекста) и пробросить отмену. Commit/rollback туда НЕ вынесены:
посреди геокодинга commit зафиксировал бы частичное состояние оценки.
`estimator._db_step` переписан поверх и добавляет свои commit/rollback сам —
его поведение не меняется, гейт tests/test_3408_db_step_cancel_orphan.py
остаётся зелёным.

Заменено 34 вызова, работающих по сессии запроса: 12 в geocoder.py (кэш-чтение
и записи, геопортал, кадастр, houses, reverse, suggest), 19 в estimator.py
(в т.ч. `_backfill_house_fias`, `_save_yandex_history_items`,
`_fetch_anchor_comps`, `_price_from_inputs` с db-резолверами, персист оценки,
`_fetch_price_trend`, `_is_premium_building`), 2 в api/v1/geocode.py, 1 в
api/v1/privacy_admin.py. Не тронуты вызовы со СВОЕЙ сессией:
`user_events.schedule_event` (внутри `record_event` свой `SessionLocal`) и
`sber_index` (сессия задачи планировщика, отменять её некому).

Гейт по значению — tests/test_3449_geocoder_cancel_orphan.py: отмена по бюджету
во время шага БД геокодера, следом ГОЛЫЙ `to_thread(db.execute, ...)` (образец
персиста); проверяется, что он не вошёл в сессию, пока сирота ещё в ней.
На исходном коде тест краснеет: conflicts == 1.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Light1YT added 1 commit 2026-09-12 08:58:04 +00:00
Гейт на ВСЕ 34 проводки + запрет вложенных бюджетов (ревью #3460)
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
be9aa2f907
Сценарный тест ловил одну проводку из 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>
bot-backend merged commit 03d9745b4f into main 2026-09-12 10:11:20 +00:00
Sign in to join this conversation.
No reviewers
No milestone
No project
No assignees
2 participants
Notifications
Due date
The due date is invalid or out of range. Please use the format "yyyy-mm-dd".

No due date set.

Dependencies

No dependencies set.

Reference: lekss361/gendesign#3460
No description provided.