fix(ptica): успехи батча каталога фиксируются по ходу, а не одним commit'ом в конце (#2464) #2960

Merged
bot-backend merged 1 commit from fix/2464-catalog-batch-commits into main 2026-08-20 08:54:27 +00:00
2 changed files with 156 additions and 3 deletions

View file

@ -462,12 +462,27 @@ async def scrape_catalog_objects(
ok = await scrape_catalog_object(db, session, obj_id, snapshot_date)
if ok:
stats["succeeded"] += 1
# Фиксируем сразу, а не одним commit'ом в конце (#2464). Раньше весь
# батч жил в одной незакоммиченной транзакции, и любой отказ ПОСЛЕ
# цикла — исключение в BrowserSession.__aexit__, снятие Celery-таски,
# перезапуск контейнера — обнулял все уже успешные UPDATE'ы.
# Это не теория: беговой режим здесь force=True («Загрузить все»),
# то есть SQL без LIMIT. На 20.08.2026 в очереди 13200 объектов из
# 13801 — многочасовой прогон, где отказ в конце стоил бы всего.
# SAVEPOINT внутри scrape_catalog_object к этому моменту уже снят,
# поэтому commit здесь корректен.
try:
db.commit()
except Exception:
db.rollback()
raise
else:
stats["failed"] += 1
# Commit outer transaction: SAVEPOINT (`begin_nested`) releases внутри loop,
# но outer tx остаётся autobegin'd — без commit() все UPDATE'ы откатятся
# при db.close() в Celery task.
# Финальный commit. Успешные строки зафиксированы по ходу цикла (см. выше), но
# этот вызов остаётся: он закрывает транзакцию, которую могли autobegin'ить
# неудачные итерации (их SAVEPOINT откатился, а внешняя транзакция открыта),
# и сохраняет прежнее поведение для вызывающих, которые на него полагались.
try:
db.commit()
except Exception:

View file

@ -0,0 +1,138 @@
"""Успешные объекты каталога фиксируются по ходу батча, а не одним commit'ом в конце (#2464).
`scrape_catalog_objects` держал весь батч в одной незакоммиченной транзакции: `db.commit()`
стоял ПОСЛЕ блока `async with BrowserSession(...)`. Любой отказ после цикла исключение в
`BrowserSession.__aexit__`, снятие Celery-таски, перезапуск контейнера обнулял все уже
успешные UPDATE'ы.
Масштаб не гипотетический. Беговой режим здесь `force=True` («Загрузить все»), то есть SQL
без LIMIT: на 20.08.2026 в очереди 13200 объектов из 13801 (последний успешный скрейп
19.05, до блокировки DOM.РФ, #2443). Многочасовой прогон, где отказ в конце стоил бы всего.
Тест не ходит в сеть и в БД: `BrowserSession` и `scrape_catalog_object` подменяются, сессия
БД счётчик вызовов. Проверяется поведение доживают ли успехи до фиксации, а не
наличие новых символов.
"""
from __future__ import annotations
import os
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
import asyncio
from datetime import date
from typing import Any
from unittest.mock import patch
import pytest
class _RecordingSession:
"""Сессия-счётчик: помнит порядок commit/rollback."""
def __init__(self) -> None:
self.events: list[str] = []
def commit(self) -> None:
self.events.append("commit")
def rollback(self) -> None:
self.events.append("rollback")
@property
def commits(self) -> int:
return self.events.count("commit")
class _FakeBrowserSession:
"""Подмена BrowserSession: ничего не делает, но умеет упасть на выходе."""
raise_on_exit: BaseException | None = None
def __init__(self, *_a: Any, **_kw: Any) -> None:
pass
async def __aenter__(self) -> _FakeBrowserSession:
return self
async def __aexit__(self, *_exc: Any) -> None:
if _FakeBrowserSession.raise_on_exit is not None:
exc = _FakeBrowserSession.raise_on_exit
_FakeBrowserSession.raise_on_exit = None
raise exc
async def warm_up(self) -> None:
return None
def _run(db: _RecordingSession, obj_ids: list[int], *, fail_on_exit: bool = False) -> Any:
from app.services.scrapers import domrf_catalog_object as mod
async def _fake_scrape(_db: Any, _s: Any, obj_id: int, _d: date) -> bool:
# Чётные — успех, нечётные — неудача: проверяем, что фиксируются именно успехи.
return obj_id % 2 == 0
_FakeBrowserSession.raise_on_exit = (
RuntimeError("падение на выходе из сессии") if fail_on_exit else None
)
with (
patch.object(mod, "BrowserSession", _FakeBrowserSession),
patch.object(mod, "scrape_catalog_object", _fake_scrape),
):
return asyncio.run(
mod.scrape_catalog_objects(
db=db, obj_ids=obj_ids, snapshot_date=date(2026, 8, 20), region_code=66
)
)
def test_successes_survive_a_late_failure() -> None:
"""Отказ ПОСЛЕ цикла не должен стоить уже собранных объектов.
На origin/main commit стоит после `async with`, поэтому исключение в __aexit__
случается ДО единственной фиксации ни один UPDATE не доживает.
"""
db = _RecordingSession()
with pytest.raises(RuntimeError, match="падение на выходе"):
_run(db, [2, 4, 6, 8], fail_on_exit=True)
assert db.commits >= 4, (
f"зафиксировано {db.commits} раз при 4 успешных объектах — успехи не пережили "
"отказ после цикла: весь батч висел в одной транзакции"
)
def test_commit_happens_per_success_not_per_object() -> None:
"""Фиксируются успехи, а не каждая итерация — неудачные строки коммитить нечего."""
db = _RecordingSession()
stats = _run(db, [1, 2, 3, 4, 5, 6])
assert stats["succeeded"] == 3
assert stats["failed"] == 3
# 3 успеха по ходу + финальная фиксация в конце.
assert db.commits == 4, f"ожидали 3 по ходу + 1 финальную, получили {db.commits}"
def test_empty_batch_touches_nothing() -> None:
"""Контроль: пустой список не открывает и не фиксирует ничего."""
db = _RecordingSession()
stats = _run(db, [])
assert stats["processed"] == 0
assert db.events == [], f"пустой батч тронул сессию: {db.events}"
def test_all_failed_still_closes_transaction() -> None:
"""Контроль: батч без единого успеха всё равно закрывает транзакцию.
Неудачная итерация откатывает свой SAVEPOINT, но внешняя транзакция остаётся
autobegin'нутой — финальный commit обязан остаться на месте.
"""
db = _RecordingSession()
stats = _run(db, [1, 3, 5])
assert stats["succeeded"] == 0
assert db.commits == 1, f"ожидали одну финальную фиксацию, получили {db.commits}"