fix(ptica): загрузчик теплоснабжения фиксирует по организациям, а не одной транзакцией (#2464) (#2973)
All checks were successful
Deploy / changes (push) Successful in 9s
Deploy / build-frontend (push) Has been skipped
Deploy / deploy-caddy (push) Has been skipped
Deploy / build-backend (push) Successful in 2m26s
Deploy / build-worker (push) Successful in 3m29s
Deploy / deploy (push) Successful in 1m29s
Deploy / deploy-status (push) Successful in 1s
Deploy / perimeter-smoke (push) Successful in 8s
All checks were successful
Deploy / changes (push) Successful in 9s
Deploy / build-frontend (push) Has been skipped
Deploy / deploy-caddy (push) Has been skipped
Deploy / build-backend (push) Successful in 2m26s
Deploy / build-worker (push) Successful in 3m29s
Deploy / deploy (push) Successful in 1m29s
Deploy / deploy-status (push) Successful in 1s
Deploy / perimeter-smoke (push) Successful in 8s
This commit is contained in:
parent
cfa0046b34
commit
9b18c23a5a
2 changed files with 122 additions and 0 deletions
|
|
@ -516,6 +516,23 @@ def load_heat_reserves(db: Session | None = None) -> dict[str, dict]:
|
|||
except Exception as e:
|
||||
logger.exception("load_heat_reserves: org %s failed: %s", org, e)
|
||||
out[org] = {"error": str(e)}
|
||||
if owns_session:
|
||||
# Сбойная организация не должна тащить свои частичные записи
|
||||
# в общий коммит следующих.
|
||||
db.rollback()
|
||||
continue
|
||||
if owns_session:
|
||||
# #2464: фиксируем ПОСЛЕ КАЖДОЙ организации, а не одним коммитом в
|
||||
# конце. Раньше одна транзакция оставалась открытой на весь батч —
|
||||
# восемь организаций, у каждой несколько HTTP-раундов к медленному
|
||||
# внешнему реестру с таймаутом _HTTP_TIMEOUT=60с. Открытая транзакция
|
||||
# столько времени держит соединение и тормозит vacuum, а падение в
|
||||
# конце обнуляло бы всё уже собранное.
|
||||
#
|
||||
# ТОЛЬКО на своей сессии: при db, переданном вызывающим, транзакцией
|
||||
# распоряжается он — коммитить её здесь значило бы зафиксировать
|
||||
# чужую работу (то же правило, что для плоского rollback).
|
||||
db.commit()
|
||||
db.commit()
|
||||
except Exception as e:
|
||||
db.rollback()
|
||||
|
|
|
|||
105
backend/tests/services/site_finder/test_2464_heat_loader_tx.py
Normal file
105
backend/tests/services/site_finder/test_2464_heat_loader_tx.py
Normal file
|
|
@ -0,0 +1,105 @@
|
|||
"""Загрузчик теплоснабжения фиксирует по организациям, а не одной транзакцией на всё (#2464).
|
||||
|
||||
`load_heat_reserves` открывал сессию, проходил по ВОСЬМИ организациям — у каждой несколько
|
||||
HTTP-раундов к медленному внешнему реестру с таймаутом 60 с — и коммитил один раз в самом
|
||||
конце. Одна транзакция оставалась открытой на всё это время: держала соединение, тормозила
|
||||
vacuum, а падение в конце обнуляло бы всё уже собранное.
|
||||
|
||||
Отдельная тонкость, из-за которой наивная правка была бы неверной: функция умеет принимать
|
||||
ЧУЖУЮ сессию (`db` аргументом). На ней коммитить по ходу нельзя — транзакцией распоряжается
|
||||
вызывающий, и промежуточный commit зафиксировал бы его работу. То же правило, что для
|
||||
плоского rollback.
|
||||
|
||||
Тесты считают коммиты на сессии-двойнике, то есть проверяют поведение. На origin/main
|
||||
коммит ровно один — красное по числу, а не по отсутствию символа.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
from typing import Any
|
||||
from unittest.mock import patch
|
||||
|
||||
|
||||
class _Session:
|
||||
def __init__(self) -> None:
|
||||
self.commits = 0
|
||||
self.rollbacks = 0
|
||||
self.closed = False
|
||||
|
||||
def commit(self) -> None:
|
||||
self.commits += 1
|
||||
|
||||
def rollback(self) -> None:
|
||||
self.rollbacks += 1
|
||||
|
||||
def close(self) -> None:
|
||||
self.closed = True
|
||||
|
||||
|
||||
def _run(*, own: bool, failing: set[str] | None = None) -> tuple[_Session, list[str], int]:
|
||||
from app.services.site_finder import eias_heat_loader as mod
|
||||
|
||||
failing = failing or set()
|
||||
db = _Session()
|
||||
visited: list[str] = []
|
||||
|
||||
def _fake_org(session: Any, org: str, _org_id: int) -> dict:
|
||||
visited.append(org)
|
||||
if org in failing:
|
||||
raise RuntimeError(f"внешний реестр не ответил по {org}")
|
||||
return {"rows": 1}
|
||||
|
||||
with (
|
||||
patch.object(mod, "SessionLocal", lambda: db),
|
||||
patch.object(mod, "load_org_reserves", _fake_org),
|
||||
):
|
||||
mod.load_heat_reserves() if own else mod.load_heat_reserves(db=db)
|
||||
|
||||
return db, visited, len(mod.ORGS)
|
||||
|
||||
|
||||
def test_own_session_commits_per_organization() -> None:
|
||||
"""На своей сессии — коммит после каждой организации.
|
||||
|
||||
На origin/main коммит ровно один на весь батч.
|
||||
"""
|
||||
db, visited, n_orgs = _run(own=True)
|
||||
|
||||
assert len(visited) == n_orgs, f"обошли {len(visited)} организаций из {n_orgs}"
|
||||
assert db.commits >= n_orgs, (
|
||||
f"коммитов {db.commits} при {n_orgs} организациях — транзакция остаётся открытой "
|
||||
"на весь батч, поверх десятков минут внешнего HTTP"
|
||||
)
|
||||
|
||||
|
||||
def test_borrowed_session_is_not_committed_per_organization() -> None:
|
||||
"""Контроль: на ЧУЖОЙ сессии промежуточных коммитов быть не должно.
|
||||
|
||||
Иначе правка фиксировала бы работу вызывающего — та же ошибка, что плоский
|
||||
rollback на общей сессии.
|
||||
"""
|
||||
db, visited, n_orgs = _run(own=False)
|
||||
|
||||
assert len(visited) == n_orgs
|
||||
assert (
|
||||
db.commits == 1
|
||||
), f"на чужой сессии {db.commits} коммитов — транзакцией распоряжается вызывающий"
|
||||
assert not db.closed, "чужая сессия закрыта — её закрывает вызывающий"
|
||||
|
||||
|
||||
def test_failing_organization_does_not_stop_the_rest() -> None:
|
||||
"""Контроль: сбой одной организации не рвёт обход остальных."""
|
||||
db, visited, n_orgs = _run(own=True, failing={mod_org()})
|
||||
|
||||
assert len(visited) == n_orgs, f"обход прервался: {visited}"
|
||||
assert db.rollbacks >= 1, "частичные записи сбойной организации не откачены"
|
||||
|
||||
|
||||
def mod_org() -> str:
|
||||
from app.services.site_finder import eias_heat_loader as mod
|
||||
|
||||
return mod.ORGS[1][0]
|
||||
Loading…
Add table
Reference in a new issue