fix(site_finder): SAVEPOINT-isolate DB errors across cluster A (#2464) (#2465)
All checks were successful
Deploy / changes (push) Successful in 9s
Deploy / build-frontend (push) Has been skipped
Deploy / build-backend (push) Successful in 3m29s
Deploy / build-worker (push) Successful in 4m48s
Deploy / deploy (push) Successful in 1m43s

This commit is contained in:
bot-backend 2026-07-07 12:50:11 +00:00
parent 1fa99812af
commit 39dd63333e
14 changed files with 514 additions and 126 deletions

View file

@ -799,7 +799,12 @@ def get_competitors(
avg_price_map: dict[int, float] = {} avg_price_map: dict[int, float] = {}
price_source_map: dict[int, str] = {} price_source_map: dict[int, str] = {}
try: try:
price_rows = db.execute(_AVG_PRICE_SQL, {"obj_ids": obj_ids}).mappings().all() # SAVEPOINT — db shared с caller (POST .../competitors, Depends(get_db)). Голый
# db.rollback() orphan'ит outer SessionTransaction (velocity.py:170) — тем более
# здесь ЕЩЁ ДВА db.execute идут следом в этой же функции (sold-count,
# objective-fallback) на том же db; без savepoint сбой здесь отравил бы их тоже.
with db.begin_nested():
price_rows = db.execute(_AVG_PRICE_SQL, {"obj_ids": obj_ids}).mappings().all()
for r in price_rows: for r in price_rows:
oid = int(r["obj_id"]) oid = int(r["obj_id"])
price = _row_get(r, "avg_price_per_m2") price = _row_get(r, "avg_price_per_m2")
@ -818,11 +823,16 @@ def get_competitors(
# flats_sold → None → нейтральный stage. # flats_sold → None → нейтральный stage.
sold_count_map: dict[int, int] = {} sold_count_map: dict[int, int] = {}
try: try:
sold_rows = ( # SAVEPOINT — db shared, см. avg_price блок выше (иначе отравляет objective
db.execute(_SOLD_COUNT_SQL, {"obj_ids": obj_ids, "premise_kind": _SOLD_PREMISE_KIND}) # price fallback ниже на том же db).
.mappings() with db.begin_nested():
.all() sold_rows = (
) db.execute(
_SOLD_COUNT_SQL, {"obj_ids": obj_ids, "premise_kind": _SOLD_PREMISE_KIND}
)
.mappings()
.all()
)
for r in sold_rows: for r in sold_rows:
oid = int(r["obj_id"]) oid = int(r["obj_id"])
sold = _row_get(r, "flats_sold") sold = _row_get(r, "flats_sold")
@ -838,17 +848,21 @@ def get_competitors(
missing_price_ids = [oid for oid in obj_ids if oid not in avg_price_map] missing_price_ids = [oid for oid in obj_ids if oid not in avg_price_map]
if missing_price_ids: if missing_price_ids:
try: try:
obj_price_rows = ( # SAVEPOINT — db shared, см. avg_price блок выше. Последний из трёх
db.execute( # db.execute в get_competitors — тоже под savepoint, чтобы дальнейшие
_OBJECTIVE_PRICE_FALLBACK_SQL, # запросы analyze на этом же db (persist_analysis_run и т.п.) не отравились.
{ with db.begin_nested():
"obj_ids": missing_price_ids, obj_price_rows = (
"velocity_match_radius_m": _VELOCITY_MATCH_RADIUS_M, db.execute(
}, _OBJECTIVE_PRICE_FALLBACK_SQL,
{
"obj_ids": missing_price_ids,
"velocity_match_radius_m": _VELOCITY_MATCH_RADIUS_M,
},
)
.mappings()
.all()
) )
.mappings()
.all()
)
for r in obj_price_rows: for r in obj_price_rows:
oid = int(r["obj_id"]) oid = int(r["obj_id"])
price = _row_get(r, "median_price_per_m2") price = _row_get(r, "median_price_per_m2")

View file

@ -260,25 +260,29 @@ def _query_gas_city_grs(db: Session) -> dict:
""" """
outlet_counts = _query_gas_outlet_counts(db) outlet_counts = _query_gas_outlet_counts(db)
try: try:
rows = ( # SAVEPOINT — db shared с caller (GET .../connection-capacity, Depends(get_db)).
db.execute( # Голый db.rollback() orphan'ит outer SessionTransaction (см. velocity.py:170) —
text(""" # begin_nested откатывает только эту savepoint, outer tx остаётся usable.
SELECT grs_name, with db.begin_nested():
SUM(design_capacity_th_m3_h) AS design_sum, rows = (
SUM(free_capacity_th_m3_h) AS free_sum, db.execute(
AVG(free_capacity_pct) AS pct_avg, text("""
MAX(upgrade_due) AS upgrade_due, SELECT grs_name,
COUNT(*) AS outputs_count SUM(design_capacity_th_m3_h) AS design_sum,
FROM gas_grs_capacity SUM(free_capacity_th_m3_h) AS free_sum,
WHERE grs_name ILIKE '%свердловск%' AVG(free_capacity_pct) AS pct_avg,
OR grs_name ILIKE '%екатеринбург%' MAX(upgrade_due) AS upgrade_due,
GROUP BY grs_name COUNT(*) AS outputs_count
ORDER BY grs_name FROM gas_grs_capacity
""") WHERE grs_name ILIKE '%свердловск%'
OR grs_name ILIKE '%екатеринбург%'
GROUP BY grs_name
ORDER BY grs_name
""")
)
.mappings()
.all()
) )
.mappings()
.all()
)
except Exception as e: except Exception as e:
# Таблица может не существовать до применения миграции 181 (schema-first). # Таблица может не существовать до применения миграции 181 (schema-first).
# Блок аддитивный — деградируем в пустой, не роняя весь эндпоинт. Счётчики # Блок аддитивный — деградируем в пустой, не роняя весь эндпоинт. Счётчики
@ -333,25 +337,27 @@ def _query_gas_outlet_counts(db: Session) -> dict:
Таблицы может не быть (миграция 184 не применена) все нули (graceful, аддитивно). Таблицы может не быть (миграция 184 не применена) все нули (graceful, аддитивно).
""" """
try: try:
row = ( # SAVEPOINT — db shared с caller (см. _query_gas_city_grs выше).
db.execute( with db.begin_nested():
text(""" row = (
SELECT COUNT(*) AS outlets_total, db.execute(
COUNT(*) FILTER ( text("""
WHERE free_capacity_mln_m3 < 0 SELECT COUNT(*) AS outlets_total,
) AS outlets_deficit, COUNT(*) FILTER (
COUNT(*) FILTER ( WHERE free_capacity_mln_m3 < 0
WHERE free_capacity_needs_calc ) AS outlets_deficit,
) AS outlets_needs_calc COUNT(*) FILTER (
FROM gas_grs_outlet_points WHERE free_capacity_needs_calc
WHERE period_month = ( ) AS outlets_needs_calc
SELECT MAX(period_month) FROM gas_grs_outlet_points FROM gas_grs_outlet_points
) WHERE period_month = (
""") SELECT MAX(period_month) FROM gas_grs_outlet_points
)
""")
)
.mappings()
.first()
) )
.mappings()
.first()
)
except Exception as e: except Exception as e:
# Таблица может не существовать до применения миграции 184 (schema-first). # Таблица может не существовать до применения миграции 184 (schema-first).
# Счётчики аддитивны — деградируем в нули, не роняя весь эндпоинт. # Счётчики аддитивны — деградируем в нули, не роняя весь эндпоинт.
@ -384,39 +390,41 @@ def _query_gas_outlet_points(db: Session, parcel_wkt: str) -> list[dict]:
distance_m, lat, lon}] отсортировано по distance ASC. distance_m, lat, lon}] отсортировано по distance ASC.
""" """
try: try:
rows = ( # SAVEPOINT — db shared с caller (см. _query_gas_city_grs выше).
db.execute( with db.begin_nested():
text(""" rows = (
SELECT outlet_name, consumer_type, free_capacity_mln_m3, db.execute(
free_capacity_needs_calc AS needs_calc, text("""
ST_Distance( SELECT outlet_name, consumer_type, free_capacity_mln_m3,
geom::geography, free_capacity_needs_calc AS needs_calc,
ST_Centroid(ST_GeomFromText(:wkt, 4326))::geography ST_Distance(
) AS distance_m, geom::geography,
ST_Y(geom) AS lat, ST_Centroid(ST_GeomFromText(:wkt, 4326))::geography
ST_X(geom) AS lon ) AS distance_m,
FROM gas_grs_outlet_points ST_Y(geom) AS lat,
WHERE geom IS NOT NULL ST_X(geom) AS lon
AND period_month = ( FROM gas_grs_outlet_points
SELECT MAX(period_month) FROM gas_grs_outlet_points WHERE geom IS NOT NULL
) AND period_month = (
AND ST_DWithin( SELECT MAX(period_month) FROM gas_grs_outlet_points
geom::geography, )
ST_Centroid(ST_GeomFromText(:wkt, 4326))::geography, AND ST_DWithin(
CAST(:radius_m AS float) geom::geography,
) ST_Centroid(ST_GeomFromText(:wkt, 4326))::geography,
ORDER BY distance_m ASC CAST(:radius_m AS float)
LIMIT CAST(:limit AS int) )
"""), ORDER BY distance_m ASC
{ LIMIT CAST(:limit AS int)
"wkt": parcel_wkt, """),
"radius_m": _GAS_OUTLET_RADIUS_M, {
"limit": _GAS_OUTLET_LIMIT, "wkt": parcel_wkt,
}, "radius_m": _GAS_OUTLET_RADIUS_M,
"limit": _GAS_OUTLET_LIMIT,
},
)
.mappings()
.all()
) )
.mappings()
.all()
)
except Exception as e: except Exception as e:
# Таблица может не существовать до применения миграции 184 (schema-first). # Таблица может не существовать до применения миграции 184 (schema-first).
# Слой аддитивный — деградируем в пустой, не роняя весь эндпоинт. # Слой аддитивный — деградируем в пустой, не роняя весь эндпоинт.
@ -451,23 +459,25 @@ def _query_heat_latest(db: Session) -> dict:
"total_reserve_gcal_h": float|None}. "total_reserve_gcal_h": float|None}.
""" """
try: try:
rows = ( # SAVEPOINT — db shared с caller (см. _query_gas_city_grs выше).
db.execute( with db.begin_nested():
text(""" rows = (
SELECT h.org, h.system_name, h.reserve_gcal_h, h.period db.execute(
FROM heat_system_reserves h text("""
WHERE h.period IS NOT NULL SELECT h.org, h.system_name, h.reserve_gcal_h, h.period
AND h.period = ( FROM heat_system_reserves h
SELECT MAX(h2.period) FROM heat_system_reserves h2 WHERE h.period IS NOT NULL
WHERE h2.period IS NOT NULL AND h.period = (
AND h2.org = h.org SELECT MAX(h2.period) FROM heat_system_reserves h2
) WHERE h2.period IS NOT NULL
ORDER BY h.org, h.system_name AND h2.org = h.org
""") )
ORDER BY h.org, h.system_name
""")
)
.mappings()
.all()
) )
.mappings()
.all()
)
except Exception as e: except Exception as e:
# Таблица может не существовать до применения миграции 182 (schema-first). # Таблица может не существовать до применения миграции 182 (schema-first).
# Блок аддитивный — деградируем в пустой, не роняя весь эндпоинт. # Блок аддитивный — деградируем в пустой, не роняя весь эндпоинт.

View file

@ -142,14 +142,22 @@ def get_developer_attribution(
- БД-ошибке, включая ещё не задеплоенную миграцию 149 (analyze не падает). - БД-ошибке, включая ещё не задеплоенную миграцию 149 (analyze не падает).
""" """
try: try:
rows = ( # SAVEPOINT — db получена от caller (shared request-scoped Session,
db.execute( # analyze_parcel в parcels.py). Голый db.rollback() здесь НЕЛЬЗЯ — он
_DEVELOPER_FOR_PARCEL_SQL, # orphan'ит outer SessionTransaction и роняет persist_analysis_run сразу
{"cad_num": cad_num, "radius": radius_m}, # после (см. velocity.py:170 / connection_capacity_lookup._query_nearby_
# network_zones). Ошибка внутри begin_nested() откатывает ТОЛЬКО эту
# savepoint до propagate — outer-транзакция остаётся usable для
# последующих db.execute (#2464 cluster A finding 1).
with db.begin_nested():
rows = (
db.execute(
_DEVELOPER_FOR_PARCEL_SQL,
{"cad_num": cad_num, "radius": radius_m},
)
.mappings()
.all()
) )
.mappings()
.all()
)
except (OperationalError, ProgrammingError, DataError) as exc: except (OperationalError, ProgrammingError, DataError) as exc:
# Миграция 149 ещё не задеплоена / БД-ошибка — graceful degrade. # Миграция 149 ещё не задеплоена / БД-ошибка — graceful degrade.
logger.warning( logger.warning(

View file

@ -45,7 +45,12 @@ def parcel_pat_subzones(db: Session, parcel_wkt: str | None) -> list[dict[str, A
return [] return []
try: try:
rows = db.execute(_PAT_INTERSECTS_SQL, {"parcel_wkt": parcel_wkt}).mappings().all() # SAVEPOINT — db shared с caller (analyze_parcel, shared request-scoped Session).
# Голый db.rollback() orphan'ит outer SessionTransaction (см. velocity.py:170) —
# begin_nested откатывает только эту savepoint, не мешая persist_analysis_run
# и прочим db.execute дальше по /analyze (#2464 cluster A finding 7).
with db.begin_nested():
rows = db.execute(_PAT_INTERSECTS_SQL, {"parcel_wkt": parcel_wkt}).mappings().all()
except (OperationalError, ProgrammingError) as exc: except (OperationalError, ProgrammingError) as exc:
# Таблица ещё не задеплоена / структура изменилась — graceful degrade. # Таблица ещё не задеплоена / структура изменилась — graceful degrade.
logger.warning("parcel_pat_subzones: pat_subzones недоступна, skip: %s", exc) logger.warning("parcel_pat_subzones: pat_subzones недоступна, skip: %s", exc)

View file

@ -237,7 +237,15 @@ def compute_district_saturation(db: Session, district_name: str) -> dict[str, An
psycopg v3: bind через CAST(:dn AS text). Read-only. psycopg v3: bind через CAST(:dn AS text). Read-only.
""" """
try: try:
row = db.execute(_DISTRICT_SATURATION_SQL, {"dn": district_name}).mappings().first() # SAVEPOINT — db shared с caller (analyze_parcel, parcels.py). Caller УЖЕ
# оборачивает вызов этой функции в свой db.begin_nested(), но это НЕ спасает:
# если db.execute здесь падает и мы её же тут глотаем (return None, не re-raise),
# исключение никогда не долетает до caller-savepoint → её ROLLBACK TO SAVEPOINT
# никогда не срабатывает, а сам Postgres-connection уже aborted. Нужна СОБСТВЕННАЯ
# savepoint здесь — откатывается на выходе из `with` ДО того, как исключение
# поймано ниже, оставляя outer-транзакцию (caller'а) usable.
with db.begin_nested():
row = db.execute(_DISTRICT_SATURATION_SQL, {"dn": district_name}).mappings().first()
except Exception: except Exception:
logger.exception("saturation: query failed (district=%s)", district_name) logger.exception("saturation: query failed (district=%s)", district_name)
return None return None

View file

@ -570,9 +570,18 @@ def _safe_rows(
Зеркало market_metrics._query_* тонкие данные/сбой БД не должны валить расчёт Зеркало market_metrics._query_* тонкие данные/сбой БД не должны валить расчёт
остальных слоёв или worker'а целиком. остальных слоёв или worker'а целиком.
SAVEPOINT (``db.begin_nested()``): ``compute_all_layers`` вызывает эту функцию 3×
(L1/L2/L3) на ОДНОЙ Session без savepoint сбой L1 отравил бы транзакцию и L2/L3
молча получили бы [] тоже (плюс любой db.execute ПОСЛЕ compute_all_layers на этом же
db, включая caller'а — worker owns SessionLocal(), но при вызове из shared
request-scoped db, напр. forecasting/orchestrator.build_site_finder_report,
отравился бы весь оставшийся запрос). Голый db.rollback() здесь НЕЛЬЗЯ см.
velocity.py:170.
""" """
try: try:
return db.execute(sql, dict(params)).mappings().all() with db.begin_nested():
return db.execute(sql, dict(params)).mappings().all()
except Exception: except Exception:
# L1 фильтрует по resolved-микро (`districts`); L2/L3 — по админ `district`. # L1 фильтрует по resolved-микро (`districts`); L2/L3 — по админ `district`.
district_dbg = params.get("district", params.get("districts")) district_dbg = params.get("district", params.get("districts"))

View file

@ -470,9 +470,15 @@ def backfill_ekb_zone_regulations(
upserted += 1 upserted += 1
except Exception as exc: except Exception as exc:
# Per-zone изоляция: один сбойный фетч/апсёрт не должен ронять весь прогон. # Per-zone изоляция: один сбойный фетч/апсёрт не должен ронять весь прогон.
logger.warning( # rollback() (не begin_nested): db здесь — СОБСТВЕННАЯ сессия таска
"backfill_ekb_zone_regulations: зона idx=%s failed: %s", zone_index, exc # backfill_zone_regulations.py (SessionLocal() + try/finally close), не
) # shared request-scoped Session — orphan'ить нечего, следующая итерация
# цикла (get_cached_zone_regulation / upsert_zone_regulation) на этом же db
# иначе поймает "current transaction is aborted" после сбойного db.commit()
# выше (#2464 cluster A finding 6; сиблинг — write_db.rollback() в
# get_or_fetch_zone_regulation, тоже owned-сессия).
db.rollback()
logger.warning("backfill_ekb_zone_regulations: зона idx=%s failed: %s", zone_index, exc)
failed += 1 failed += 1
continue continue

View file

@ -545,6 +545,63 @@ def test_competitors_no_price_anywhere_is_none() -> None:
app.dependency_overrides.clear() app.dependency_overrides.clear()
# ── SAVEPOINT-регрессия (#2464 cluster A finding 3) ──────────────────────────
def test_competitors_avg_price_failure_does_not_poison_later_queries() -> None:
"""avg_price query падает → sold-count и objective-fallback (следом на ТОЙ ЖЕ
Session) всё равно отрабатывают, а не ловят "current transaction is aborted".
Раньше avg_price/sold-count/objective-fallback ловили db.execute-сбой без
SAVEPOINT (db.begin_nested()) сбой ЛЮБОГО из трёх отравлял shared db для
ОСТАЛЬНЫХ двух (и для persist_analysis_run дальше по /analyze). Фикс оборачивает
каждый в свою begin_nested(); здесь ловим именно avg_price и проверяем, что
objective-fallback (третий execute) всё равно наполняет цену.
"""
coord_result = MagicMock()
coord_result.mappings.return_value.first.return_value = _coord_row()
obj_result = MagicMock()
obj_result.mappings.return_value.all.return_value = [_obj_row(obj_id=1)]
sold_result = MagicMock()
sold_result.mappings.return_value.all.return_value = []
obj_price_result = MagicMock()
obj_price_result.mappings.return_value.all.return_value = [
_obj_price_row(obj_id=1, price=142_000.0)
]
db = MagicMock()
db.execute.side_effect = [
coord_result,
obj_result,
RuntimeError("simulated avg_price DB failure"), # avg_price query сбоит
sold_result,
obj_price_result, # objective-fallback ДОЛЖЕН всё равно выполниться
]
from app.core.db import get_db
app.dependency_overrides[get_db] = _override_db(db)
try:
client = TestClient(app)
resp = client.post(
"/api/v1/parcels/66:41:0303161:5/competitors",
json={"radius_km": 1.0, "time_window": "last_quarter"},
)
assert resp.status_code == 200, resp.text
comp = resp.json()["competitors"][0]
# avg_price сбойнул (без domrf-цены), но objective-fallback наполнил её —
# это возможно ТОЛЬКО если сбой avg_price не отравил сессию для fallback.
assert comp["avg_price_per_m2"] == pytest.approx(142_000.0)
assert comp["price_source"] == "objective"
# Все 5 запланированных execute реально дошли (sold-count и fallback не
# были пропущены/не упали каскадом из-за отравленной сессии).
assert db.execute.call_count == 5
# avg_price/sold-count/objective-fallback — каждый под своим SAVEPOINT.
assert db.begin_nested.call_count == 3
finally:
app.dependency_overrides.clear()
def test_competitors_multi_object_radius_yields_per_object_prices() -> None: def test_competitors_multi_object_radius_yields_per_object_prices() -> None:
"""Регрессия #2445 A5 (byte-identical к best_layouts #1956): несколько """Регрессия #2445 A5 (byte-identical к best_layouts #1956): несколько
конкурентов в радиусе КАЖДЫЙ получает свою domrf-цену, а не «1 объект на конкурентов в радиусе КАЖДЫЙ получает свою domrf-цену, а не «1 объект на

View file

@ -0,0 +1,138 @@
"""SAVEPOINT-регрессия для connection_capacity_lookup.py (#2464 cluster A finding 2).
gas/heat/gas-outlet query helpers (_query_gas_city_grs, _query_gas_outlet_counts,
_query_gas_outlet_points, _query_heat_latest) когда-то ловили db.execute-сбой в
bare ``except Exception`` без SAVEPOINT/rollback на реальном Postgres это отравляло
db shared с caller (GET .../connection-capacity, Depends(get_db)): следующий db.execute
в этом же запросе падал с "current transaction is aborted". Сиблинг
``_query_nearby_network_zones`` (в этом же файле) уже был правильным (db.begin_nested())
эти тесты проверяют, что остальные функции приведены в соответствие.
MagicMock-сессия не эмулирует реальный aborted-transaction Postgres тесты
проверяют (1) graceful fallback при сбое, (2) что begin_nested() реально вызван
вокруг db.execute (SAVEPOINT-обёртка присутствует в коде), (3) что db остаётся
usable для следующего вызова на том же mock-сессии (call_count растёт, ошибка не
проглатывает control flow).
"""
from __future__ import annotations
import os
from unittest.mock import MagicMock
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
from app.services.site_finder.connection_capacity_lookup import (
_query_gas_city_grs,
_query_gas_outlet_counts,
_query_gas_outlet_points,
_query_heat_latest,
)
_WKT = "POLYGON((60.8 56.7, 60.9 56.7, 60.9 56.8, 60.8 56.8, 60.8 56.7))"
def _mock_db_raising_once(final_rows: list[dict] | None = None) -> MagicMock:
"""db.execute кидает на 1-м вызове, на 2-м успешно отдаёт final_rows."""
db = MagicMock()
ok_result = MagicMock()
ok_result.mappings.return_value.all.return_value = final_rows or []
ok_result.mappings.return_value.first.return_value = (final_rows or [None])[0]
db.execute.side_effect = [RuntimeError("simulated DB failure"), ok_result]
return db
class TestGasCityGrsSavepoint:
def test_query_failure_degrades_gracefully(self) -> None:
db = MagicMock()
db.execute.side_effect = RuntimeError("gas_grs_capacity недоступна")
out = _query_gas_city_grs(db)
assert out["city_grs"] == []
assert out["total_free_th_m3_h"] is None
def test_uses_begin_nested_savepoint(self) -> None:
"""Фикс: query обёрнут в db.begin_nested(), не голый db.execute."""
db = MagicMock()
db.execute.return_value.mappings.return_value.all.return_value = []
db.execute.return_value.mappings.return_value.first.return_value = None
_query_gas_city_grs(db)
# 2 execute на функцию: _query_gas_outlet_counts (city-level счётчики) +
# сам gas_grs_capacity query — оба теперь под begin_nested.
assert db.begin_nested.call_count >= 1
def test_session_usable_after_failure(self) -> None:
"""db.execute сбоит один раз (city_grs), но db остаётся usable для
следующего (не связанного) вызова на этой же сессии."""
db = _mock_db_raising_once()
# outlet_counts вызывается ПЕРВЫМ внутри _query_gas_city_grs — заставим его
# тоже отдать пустой результат нормально, а основной query — кинуть.
db.execute.side_effect = [
MagicMock(
mappings=MagicMock(return_value=MagicMock(first=lambda: None))
), # outlet_counts
RuntimeError("main query failed"), # gas_grs_capacity query
]
out = _query_gas_city_grs(db)
assert out["city_grs"] == []
# Следующий вызов на том же db.execute отрабатывает (session не отравлена).
next_result = MagicMock()
next_result.mappings.return_value.all.return_value = [{"marker": "ok"}]
db.execute.side_effect = None
db.execute.return_value = next_result
assert db.execute("SELECT 1").mappings().all() == [{"marker": "ok"}]
class TestGasOutletCountsSavepoint:
def test_query_failure_degrades_to_zeros(self) -> None:
db = MagicMock()
db.execute.side_effect = RuntimeError("gas_grs_outlet_points недоступна")
out = _query_gas_outlet_counts(db)
assert out == {"outlets_total": 0, "outlets_deficit": 0, "outlets_needs_calc": 0}
def test_uses_begin_nested_savepoint(self) -> None:
db = MagicMock()
db.execute.return_value.mappings.return_value.first.return_value = None
_query_gas_outlet_counts(db)
assert db.begin_nested.call_count == 1
class TestGasOutletPointsSavepoint:
def test_query_failure_degrades_to_empty_list(self) -> None:
db = MagicMock()
db.execute.side_effect = RuntimeError("gas_grs_outlet_points (радиус) недоступна")
assert _query_gas_outlet_points(db, _WKT) == []
def test_uses_begin_nested_savepoint(self) -> None:
db = MagicMock()
db.execute.return_value.mappings.return_value.all.return_value = []
_query_gas_outlet_points(db, _WKT)
assert db.begin_nested.call_count == 1
class TestHeatLatestSavepoint:
def test_query_failure_degrades_gracefully(self) -> None:
db = MagicMock()
db.execute.side_effect = RuntimeError("heat_system_reserves недоступна")
out = _query_heat_latest(db)
assert out == {"systems": [], "total_reserve_gcal_h": None}
def test_uses_begin_nested_savepoint(self) -> None:
db = MagicMock()
db.execute.return_value.mappings.return_value.all.return_value = []
_query_heat_latest(db)
assert db.begin_nested.call_count == 1
def test_session_usable_after_failure(self) -> None:
"""execute сбоит один раз, db остаётся usable для следующего вызова —
имитирует heat gas_outlet_points nearby_network_zones в
get_connection_capacity() на одной shared Session."""
db = MagicMock()
db.execute.side_effect = RuntimeError("boom")
out = _query_heat_latest(db)
assert out["systems"] == []
db.execute.side_effect = None
ok_result = MagicMock()
ok_result.mappings.return_value.all.return_value = [{"marker": "next-query-ok"}]
db.execute.return_value = ok_result
assert db.execute("SELECT 1").mappings().all() == [{"marker": "next-query-ok"}]

View file

@ -289,3 +289,33 @@ def test_unexpected_error_also_degrades_to_none() -> None:
db.execute.side_effect = RuntimeError("boom") db.execute.side_effect = RuntimeError("boom")
result = get_developer_attribution(db, "66:41:0303161:1") result = get_developer_attribution(db, "66:41:0303161:1")
assert result is None assert result is None
# ── SAVEPOINT-регрессия (#2464 cluster A finding 1) ──────────────────────────
def test_db_error_uses_savepoint_and_session_stays_usable() -> None:
"""Фикс: db.execute обёрнут в db.begin_nested() (SAVEPOINT), НЕ голый rollback.
db здесь shared request-scoped Session (Depends(get_db) в parcels.py analyze_parcel);
без SAVEPOINT сбой fn_developer_for_parcel отравил бы persist_analysis_run сразу
после (parcels.py:4116). Мокаем db.execute: 1-й вызов (внутри
get_developer_attribution) кидает, 2-й (имитация persist_analysis_run следом на
том же db) отрабатывает нормально.
"""
db = MagicMock()
ok_result = MagicMock()
ok_result.mappings.return_value.all.return_value = [{"marker": "persist-ok"}]
db.execute.side_effect = [
ProgrammingError("stmt", {}, Exception("fn_developer_for_parcel does not exist")),
ok_result,
]
result = get_developer_attribution(db, "66:41:0303161:1")
assert result is None
assert db.begin_nested.call_count == 1 # SAVEPOINT вокруг резолвера
# Сессия осталась usable для следующего (не связанного) db.execute.
next_call = db.execute("SELECT 1")
assert next_call.mappings().all() == [{"marker": "persist-ok"}]
assert db.execute.call_count == 2

View file

@ -8,6 +8,7 @@ from __future__ import annotations
import importlib import importlib
import pathlib import pathlib
from contextlib import contextmanager
from typing import Any from typing import Any
import pytest import pytest
@ -211,6 +212,10 @@ class _Result:
class _FakeDB: class _FakeDB:
"""begin_nested — реальный (не-swallow) CM: зеркалит SAVEPOINT-фикс #2464
cluster A finding 7 (parcel_pat_subzones); исключение из db.execute должно
долетать до except в проде, не глушиться на выходе из ``with``."""
def __init__( def __init__(
self, self,
rows: list[dict[str, Any]] | None = None, rows: list[dict[str, Any]] | None = None,
@ -218,8 +223,14 @@ class _FakeDB:
) -> None: ) -> None:
self._rows = rows or [] self._rows = rows or []
self._raise = raise_exc self._raise = raise_exc
self.execute_calls = 0
@contextmanager
def begin_nested(self): # type: ignore[no-untyped-def]
yield
def execute(self, sql: Any, params: dict[str, Any] | None = None) -> _Result: def execute(self, sql: Any, params: dict[str, Any] | None = None) -> _Result:
self.execute_calls += 1
if self._raise is not None: if self._raise is not None:
raise self._raise raise self._raise
return _Result(self._rows) return _Result(self._rows)
@ -270,6 +281,27 @@ def test_pat_lookup_graceful_programming_error() -> None:
assert result == [] assert result == []
def test_pat_lookup_db_error_leaves_session_usable_for_next_query() -> None:
"""#2464 cluster A finding 7: db shared с caller (analyze_parcel) — сбой здесь
не должен отравлять session для persist_analysis_run и прочих db.execute дальше
по /analyze. Симулируем: execute кидает ОДИН раз, затем следующий (не связанный)
execute на ТОЙ ЖЕ session отрабатывает нормально."""
from sqlalchemy.exc import OperationalError
db = _FakeDB(raise_exc=OperationalError("stmt", {}, Exception("relation does not exist")))
result = parcel_pat_subzones(db, _WKT)
assert result == []
assert db.execute_calls == 1
# Сессия осталась usable: сбрасываем "raise" (имитируя, что savepoint откатился
# и не поломала outer-транзакцию) и проверяем, что db.execute снова отрабатывает.
db._raise = None
db._rows = [{"subzone_no": 3, "name": "N", "restriction": "R", "aerodrome": "A"}]
next_result = parcel_pat_subzones(db, _WKT)
assert next_result == [{"subzone_no": 3, "name": "N", "restriction": "R", "aerodrome": "A"}]
assert db.execute_calls == 2
# ── Резолв пути к JSON (#1150) ────────────────────────────────────────────── # ── Резолв пути к JSON (#1150) ──────────────────────────────────────────────

View file

@ -704,3 +704,25 @@ class TestComputeAllLayers:
db.execute.side_effect = [l1_fail, l2_ok, l3_ok] db.execute.side_effect = [l1_fail, l2_ok, l3_ok]
out = compute_all_layers(db, district="A") out = compute_all_layers(db, district="A")
assert out == [] assert out == []
def test_l1_failure_uses_savepoint_not_bare_rollback(self) -> None:
"""#2464 cluster A finding 5: _safe_rows оборачивает db.execute в begin_nested()
(не голый db.rollback()) L1 падает, но L2/L3 execute на ТОЙ ЖЕ Session должны
быть достижимы (session не отравлена/не orphan'ена outer-транзакция)."""
db = MagicMock()
l1_fail = RuntimeError("l1 down")
l2_ok = MagicMock()
l2_ok.mappings.return_value.all.return_value = [
{"district_name": "A", "n_objects_total": 1, "n_with_free_flats": 1, "hidden_units": 1}
]
l3_ok = MagicMock()
l3_ok.mappings.return_value.all.return_value = []
db.execute.side_effect = [l1_fail, l2_ok, l3_ok]
out = compute_all_layers(db, district="A")
# L1 деградировал в [], но L2 (следующий вызов на этом же db) успешно вернул строку —
# begin_nested() внутри _safe_rows не даёт сбою L1 отравить L2/L3.
assert db.execute.call_count == 3
assert db.begin_nested.call_count == 3 # SAVEPOINT на каждый из L1/L2/L3
assert any(r.layer == 2 for r in out)

View file

@ -63,14 +63,22 @@ class _FakeClient:
class _DB: class _DB:
"""Фейковая Session: считает commit'ы, ничего не пишет.""" """Фейковая Session: считает commit'ы/rollback'и, ничего не пишет.
rollback() owned-session (Celery-таск backfill_zone_regulations.py владеет
SessionLocal() целиком) fallback на per-zone сбое (#2464 cluster A finding 6).
"""
def __init__(self) -> None: def __init__(self) -> None:
self.commits = 0 self.commits = 0
self.rollbacks = 0
def commit(self) -> None: def commit(self) -> None:
self.commits += 1 self.commits += 1
def rollback(self) -> None:
self.rollbacks += 1
@pytest.fixture @pytest.fixture
def _no_sleep(monkeypatch: Any) -> None: def _no_sleep(monkeypatch: Any) -> None:
@ -146,12 +154,17 @@ def test_per_zone_error_isolation(monkeypatch: Any, _no_sleep: None) -> None:
zr, "upsert_zone_regulation", lambda db, reg, **kw: upserts.append(reg) or {} zr, "upsert_zone_regulation", lambda db, reg, **kw: upserts.append(reg) or {}
) )
result = zr.backfill_ekb_zone_regulations(_DB(), client=client, rate_delay_s=0.0) # type: ignore[arg-type] db = _DB()
result = zr.backfill_ekb_zone_regulations(db, client=client, rate_delay_s=0.0) # type: ignore[arg-type]
assert result["unique_zones"] == 2 assert result["unique_zones"] == 2
assert result["failed"] == 1 # первая зона упала assert result["failed"] == 1 # первая зона упала
assert result["upserted"] == 1 # вторая прошла → batch не остановился assert result["upserted"] == 1 # вторая прошла → batch не остановился
assert len(client.resolved) == 2 # дошли до обеих зон assert len(client.resolved) == 2 # дошли до обеих зон
# #2464 cluster A finding 6: сбойная зона откатывает СВОЮ (owned) сессию —
# иначе следующая итерация цикла поймала бы "current transaction is aborted".
assert db.rollbacks == 1
assert db.commits == 1 # только вторая (успешная) зона закоммитилась
def test_limit_caps_fetches(monkeypatch: Any, _no_sleep: None) -> None: def test_limit_caps_fetches(monkeypatch: Any, _no_sleep: None) -> None:

View file

@ -5,6 +5,7 @@ DB-слой compute_district_saturation — через минимальный mo
(тот же приём, что в test_poi_score.py). (тот же приём, что в test_poi_score.py).
""" """
from contextlib import contextmanager
from datetime import date from datetime import date
import pytest import pytest
@ -67,28 +68,21 @@ def test_multiplier_none_and_unknown_category_neutral():
def test_provision_unknown_category_none(): def test_provision_unknown_category_none():
assert ( assert (
compute_provision_ratio( compute_provision_ratio(poi_count=5, population=100000, age_share=0.1, category="park")
poi_count=5, population=100000, age_share=0.1, category="park"
)
is None is None
) )
def test_provision_zero_population_none(): def test_provision_zero_population_none():
assert ( assert (
compute_provision_ratio( compute_provision_ratio(poi_count=5, population=0, age_share=0.1, category="school") is None
poi_count=5, population=0, age_share=0.1, category="school"
)
is None
) )
def test_provision_missing_age_share_none(): def test_provision_missing_age_share_none():
"""Школа/детсад без age_share → None (когорту не оценить).""" """Школа/детсад без age_share → None (когорту не оценить)."""
assert ( assert (
compute_provision_ratio( compute_provision_ratio(poi_count=5, population=100000, age_share=None, category="school")
poi_count=5, population=100000, age_share=None, category="school"
)
is None is None
) )
@ -147,12 +141,21 @@ class _MockResult:
class _MockDb: class _MockDb:
"""Минимальный мок SQLAlchemy Session (как в test_poi_score.py).""" """Минимальный мок SQLAlchemy Session (как в test_poi_score.py).
``begin_nested`` реальный (не-swallow) context manager, зеркалит SAVEPOINT из
``compute_district_saturation`` (#2464 cluster A finding 4): исключение внутри
``with`` должно долетать до ``except`` в проде, не глушиться на выходе из CM.
"""
def __init__(self, row: dict | None, *, raise_on_execute: bool = False) -> None: def __init__(self, row: dict | None, *, raise_on_execute: bool = False) -> None:
self._row = row self._row = row
self._raise = raise_on_execute self._raise = raise_on_execute
@contextmanager
def begin_nested(self): # type: ignore[no-untyped-def]
yield
def execute(self, *_args: object, **_kwargs: object) -> _MockResult: def execute(self, *_args: object, **_kwargs: object) -> _MockResult:
if self._raise: if self._raise:
raise RuntimeError("simulated DB failure") raise RuntimeError("simulated DB failure")
@ -187,6 +190,39 @@ def test_saturation_db_error_returns_none():
assert compute_district_saturation(db, "Чкаловский") is None assert compute_district_saturation(db, "Чкаловский") is None
def test_saturation_db_error_leaves_session_usable_for_next_query():
"""#2464 cluster A finding 4: сбой query здесь не должен отравлять session.
Симулируем реальный сценарий: db.execute кидает ОДИН раз (внутри
compute_district_saturation), затем на ТОЙ ЖЕ session успешно отрабатывает
следующий (не связанный) запрос как это было бы с persist_analysis_run
сразу после saturation-блока в parcels.py analyze_parcel. Без SAVEPOINT
(begin_nested) второй execute на реальном Postgres упал бы с "current
transaction is aborted, commands ignored until end of transaction block".
"""
class _FlakyDb:
def __init__(self) -> None:
self.calls = 0
@contextmanager
def begin_nested(self): # type: ignore[no-untyped-def]
yield
def execute(self, *_args: object, **_kwargs: object) -> _MockResult:
self.calls += 1
if self.calls == 1:
raise RuntimeError("simulated DB failure")
return _MockResult({"marker": "next-query-succeeded"})
db = _FlakyDb()
assert compute_district_saturation(db, "Чкаловский") is None
# Сессия осталась usable: следующий (не связанный) db.execute отрабатывает.
result = db.execute("SELECT 1")
assert result.mappings().first() == {"marker": "next-query-succeeded"}
assert db.calls == 2
def test_saturation_shape_and_flags(): def test_saturation_shape_and_flags():
db = _MockDb(_chkalovsky_row()) db = _MockDb(_chkalovsky_row())
out = compute_district_saturation(db, "Чкаловский") out = compute_district_saturation(db, "Чкаловский")