diff --git a/scripts/reinvented_baseline.json b/scripts/reinvented_baseline.json index c4591e37..9b8f7841 100644 --- a/scripts/reinvented_baseline.json +++ b/scripts/reinvented_baseline.json @@ -2,24 +2,16 @@ "_comment": "Baseline-снимок scripts/detect_reinvented.py (#1108). Фиксирует текущие находки-велосипеды, чтобы прогоны диффили и подсвечивали НОВЫЕ. Пересоздать: python scripts/detect_reinvented.py --update-baseline.", "findings": [ { - "file": "src/agent/manager.py", - "line": 429, + "file": "src/agent/backends/_stream.py", + "line": 463, "kind": "handrolled-retry-loop", "what": "самописный retry/backoff-цикл в _await_with_countdown()", "replacement": "tenacity.retry / aiolimiter (для rate-limit)", "confidence": "med" }, - { - "file": "src/database/facade.py", - "line": 236, - "kind": "handrolled-retry-loop", - "what": "самописный retry/backoff-цикл в _with_busy_retry()", - "replacement": "tenacity.retry / aiolimiter (для rate-limit)", - "confidence": "med" - }, { "file": "src/services/task_handlers/stats.py", - "line": 173, + "line": 174, "kind": "handrolled-retry-loop", "what": "самописный retry/backoff-цикл в handle_stats_all()", "replacement": "tenacity.retry / aiolimiter (для rate-limit)", @@ -27,39 +19,31 @@ }, { "file": "src/services/telegram_command_dispatcher.py", - "line": 144, + "line": 152, "kind": "handrolled-retry-loop", "what": "самописный retry/backoff-цикл в _run_loop()", "replacement": "tenacity.retry / aiolimiter (для rate-limit)", "confidence": "med" }, - { - "file": "src/services/telegram_command_dispatcher.py", - "line": 343, - "kind": "handrolled-retry-loop", - "what": "самописный retry/backoff-цикл в _update_command_safely()", - "replacement": "tenacity.retry / aiolimiter (для rate-limit)", - "confidence": "med" - }, { "file": "src/services/unified_dispatcher.py", - "line": 178, + "line": 182, "kind": "handrolled-retry-loop", "what": "самописный retry/backoff-цикл в _run_loop()", "replacement": "tenacity.retry / aiolimiter (для rate-limit)", "confidence": "med" }, { - "file": "src/telegram/collector.py", - "line": 674, + "file": "src/telegram/collector_mixins/collection.py", + "line": 322, "kind": "handrolled-retry-loop", "what": "самописный retry/backoff-цикл в collect_all_channels()", "replacement": "tenacity.retry / aiolimiter (для rate-limit)", "confidence": "med" }, { - "file": "src/telegram/collector.py", - "line": 2281, + "file": "src/telegram/collector_mixins/stats.py", + "line": 475, "kind": "handrolled-retry-loop", "what": "самописный retry/backoff-цикл в collect_all_stats()", "replacement": "tenacity.retry / aiolimiter (для rate-limit)", @@ -91,15 +75,23 @@ }, { "file": "src/web/bootstrap.py", - "line": 74, + "line": 76, "kind": "handrolled-retry-loop", "what": "самописный retry/backoff-цикл в _retry_telegram_pool_until_connected()", "replacement": "tenacity.retry / aiolimiter (для rate-limit)", "confidence": "med" }, + { + "file": "src/web/scheduler/context.py", + "line": 236, + "kind": "manual-counter", + "what": "самописный частотный счётчик в _dedupe_recent_unavailability_events()", + "replacement": "collections.Counter", + "confidence": "med" + }, { "file": "src/web/search/handlers.py", - "line": 84, + "line": 86, "kind": "handrolled-retry-loop", "what": "самописный retry/backoff-цикл в _telegram_search_via_worker()", "replacement": "tenacity.retry / aiolimiter (для rate-limit)", diff --git a/src/database/facade.py b/src/database/facade.py index 6752e4f4..aee7aacb 100644 --- a/src/database/facade.py +++ b/src/database/facade.py @@ -9,6 +9,8 @@ from typing import Any import aiosqlite +from tenacity import AsyncRetrying, RetryCallState, RetryError, retry_if_exception, stop_after_attempt +from tenacity.wait import wait_chain, wait_fixed from src.database.bundles import DatabaseRepositories from src.database.connection import ConnectionTuning, DBConnection @@ -70,6 +72,26 @@ def _is_sqlite_busy_error(exc: BaseException) -> bool: return any(part in message for part in _SQLITE_BUSY_MESSAGES) +def _is_retryable_busy_error(exc: BaseException) -> bool: + return isinstance(exc, sqlite3.OperationalError) and _is_sqlite_busy_error(exc) + + +def _wrap_async( + action: Callable[[], Awaitable[aiosqlite.Cursor | None]], +) -> Callable[[], Awaitable[aiosqlite.Cursor | None]]: + """Обернуть action в async-fn для tenacity.AsyncRetrying. + + tenacity await'ит результат ``fn()``, но ``action`` — синхронный ``Callable``, + возвращающий coroutine; без обёртки coroutine остался бы не-await'нутым и ни + одна DDL/DML внутри не выполнилась бы (#1132). + """ + + async def _runner() -> aiosqlite.Cursor | None: + return await action() + + return _runner + + class Database: def __init__( self, @@ -232,33 +254,48 @@ async def _with_busy_retry( operation: str, action: Callable[[], Awaitable[aiosqlite.Cursor | None]], ) -> aiosqlite.Cursor | None: + # tenacity вместо ручного цикла (#1132); лестница задержек — ровно + # self._busy_retry_delays_sec (wait_chain), попыток len(delays)+1. delays = self._busy_retry_delays_sec - for attempt in range(len(delays) + 1): - try: - return await action() - except sqlite3.OperationalError as exc: - if not _is_sqlite_busy_error(exc): - raise - if attempt >= len(delays): - logger.warning( - "Database stayed locked during %s after %d attempts", - operation, - attempt + 1, - ) - raise DatabaseBusyError( - "Database is busy. Retry the request in a few seconds." - ) from exc - delay = delays[attempt] - logger.warning( - "Database locked during %s; retrying in %.2fs (%d/%d)", - operation, - delay, - attempt + 1, - len(delays) + 1, - ) - if delay > 0: - await asyncio.sleep(delay) - return None + + def _log_retry(retry_state: RetryCallState) -> None: + # attempt_number — номер только что провалившейся попытки (1-based), + # после которой спим delays[attempt_number-1]. Берём задержку из нашей + # лестницы, а не из tenacity-internal retry_state.next_action (то поле + # недокументировано и может уйти при апгрейте tenacity). + idx = retry_state.attempt_number - 1 + delay = delays[idx] if 0 <= idx < len(delays) else 0.0 + logger.warning( + "Database locked during %s; retrying in %.2fs (%d/%d)", + operation, + delay, + retry_state.attempt_number, + len(delays) + 1, + ) + + retryer = AsyncRetrying( + retry=retry_if_exception(_is_retryable_busy_error), + stop=stop_after_attempt(len(delays) + 1), + wait=wait_chain(*(wait_fixed(delay) for delay in delays)), + sleep=asyncio.sleep, + before_sleep=_log_retry, + ) + try: + # Оборачиваем action в async-fn: tenacity AsyncRetrying await'ит + # результат fn(), а action — это синхронный Callable, возвращающий + # coroutine (lambda: begin_immediate(...)). Без обёртки coroutine + # остаётся не-await'нутым и BEGIN не выполняется (#1132). + return await retryer(_wrap_async(action)) + except RetryError as retry_error: + last_exc = retry_error.last_attempt.exception() + logger.warning( + "Database stayed locked during %s after %d attempts", + operation, + len(delays) + 1, + ) + raise DatabaseBusyError( + "Database is busy. Retry the request in a few seconds." + ) from last_exc @asynccontextmanager async def transaction(self) -> AsyncIterator[aiosqlite.Connection]: diff --git a/src/services/telegram_command_dispatcher.py b/src/services/telegram_command_dispatcher.py index 20bed211..2d2b2ae4 100644 --- a/src/services/telegram_command_dispatcher.py +++ b/src/services/telegram_command_dispatcher.py @@ -7,6 +7,9 @@ from pathlib import Path from typing import TYPE_CHECKING, Any +from tenacity import AsyncRetrying, RetryCallState, retry_if_exception_type +from tenacity.wait import wait_exponential + from src.config import AppConfig from src.database import Database, DatabaseBusyError from src.live_runtime_pause import LiveRuntimePauseGate @@ -345,34 +348,54 @@ async def _update_command_safely( **kwargs: Any, ) -> None: assert command_id is not None - delay = COMMAND_STATUS_UPDATE_BUSY_RETRY_INITIAL_SEC - while True: - try: - await self._db.repos.telegram_commands.update_command( - command_id, - status=status, - **kwargs, - ) - return - except DatabaseBusyError as exc: - logger.warning( - "telegram_command_dispatcher: DB busy while marking command %s %s: %s", - command_id, - log_action, - exc, - ) - if not retry_busy: - return - await asyncio.sleep(delay) - delay = min(delay * 2, COMMAND_STATUS_UPDATE_BUSY_RETRY_MAX_SEC) - except Exception as exc: - logger.warning( - "telegram_command_dispatcher: failed to mark command %s %s: %s", - command_id, - log_action, - exc, - ) + + async def _update() -> None: + await self._db.repos.telegram_commands.update_command( + command_id, + status=status, + **kwargs, + ) + + def _log_busy(retry_state: RetryCallState) -> None: + outcome = retry_state.outcome + logger.warning( + "telegram_command_dispatcher: DB busy while marking command %s %s: %s", + command_id, + log_action, + outcome.exception() if outcome is not None else None, + ) + + try: + if not retry_busy: + await _update() return + # tenacity вместо ручного цикла (#1132): бесконечный retry на busy с + # экспоненциальной лестницей INITIAL, 2×, 4×… и потолком MAX — как раньше. + retryer = AsyncRetrying( + retry=retry_if_exception_type(DatabaseBusyError), + wait=wait_exponential( + multiplier=COMMAND_STATUS_UPDATE_BUSY_RETRY_INITIAL_SEC, + max=COMMAND_STATUS_UPDATE_BUSY_RETRY_MAX_SEC, + ), + sleep=asyncio.sleep, + before_sleep=_log_busy, + ) + await retryer(_update) + except DatabaseBusyError as exc: + # Достижимо только при retry_busy=False: retryer повторяет busy бесконечно. + logger.warning( + "telegram_command_dispatcher: DB busy while marking command %s %s: %s", + command_id, + log_action, + exc, + ) + except Exception as exc: + logger.warning( + "telegram_command_dispatcher: failed to mark command %s %s: %s", + command_id, + log_action, + exc, + ) async def _dispatch(self, command_type: str, payload: dict[str, Any]) -> dict[str, Any]: handler_name = f"_handle_{command_type.replace('.', '_')}" diff --git a/tests/test_database.py b/tests/test_database.py index fb8cb5d8..c3aaee92 100644 --- a/tests/test_database.py +++ b/tests/test_database.py @@ -216,6 +216,113 @@ async def locked_execute(sql, params=()): assert await db.get_agent_messages(thread_id) == [] +# --- характеризующие тесты _with_busy_retry (#1132: ручной цикл → tenacity) --- +# Фиксируют точный контракт retry до/после замены реализации: лестницу задержек, +# число попыток, сквозной проброс не-busy ошибок и цепочку исключений. + + +@pytest.mark.anyio +async def test_busy_retry_sleeps_exact_configured_ladder(db, monkeypatch): + """При исчерпании попыток спим ровно _busy_retry_delays_sec, по порядку. + + Никакой экспоненты/джиттера: лестница фиксированная, попыток len(delays)+1. + """ + calls = 0 + + async def always_locked(): + nonlocal calls + calls += 1 + raise sqlite3.OperationalError("database is locked") + + slept: list[float] = [] + + async def fake_sleep(delay): + slept.append(delay) + + db._busy_retry_delays_sec = (0.05, 0.2, 0.7) + monkeypatch.setattr("asyncio.sleep", fake_sleep) + + with pytest.raises(DatabaseBusyError, match="Database is busy"): + await db._with_busy_retry("test op", always_locked) + + assert calls == 4 + assert slept == [0.05, 0.2, 0.7] + + +@pytest.mark.anyio +async def test_busy_retry_returns_action_result_after_transient_busy(db, monkeypatch): + """Успех после transient-busy: результат action возвращается, спали ровно delays[0]. + + Также неявно защищает ``_wrap_async``: action возвращает ``sentinel`` только если + coroutine реально await'нулась. Без адаптера tenacity вернул бы не-await'нутый + coroutine объекта ``flaky``, и ``result is sentinel`` упал бы — поэтому этот тест + краснеет, если убрать ``_wrap_async`` (#1132). См. ``test_busy_retry_sleeps_...`` + для тайминга и ``test_busy_retry_chains_last_busy_error_as_cause`` для цепочки. + """ + sentinel = object() + calls = 0 + + async def flaky(): + nonlocal calls + calls += 1 + if calls == 1: + raise sqlite3.OperationalError("database table is locked") + return sentinel + + slept: list[float] = [] + + async def fake_sleep(delay): + slept.append(delay) + + db._busy_retry_delays_sec = (0.05, 0.2) + monkeypatch.setattr("asyncio.sleep", fake_sleep) + + assert await db._with_busy_retry("test op", flaky) is sentinel + assert calls == 2 + assert slept == [0.05] + + +@pytest.mark.anyio +async def test_busy_retry_reraises_non_busy_operational_error_immediately(db, monkeypatch): + """Не-busy OperationalError пробрасывается сразу: без повторов и без sleep.""" + calls = 0 + + async def broken(): + nonlocal calls + calls += 1 + raise sqlite3.OperationalError("no such table: nope") + + async def fail_sleep(delay): # pragma: no cover - защита от ложного retry + raise AssertionError("non-busy error must not be retried") + + monkeypatch.setattr("asyncio.sleep", fail_sleep) + + with pytest.raises(sqlite3.OperationalError, match="no such table"): + await db._with_busy_retry("test op", broken) + + assert calls == 1 + + +@pytest.mark.anyio +async def test_busy_retry_chains_last_busy_error_as_cause(db, monkeypatch): + """DatabaseBusyError связан цепочкой (__cause__) с последней busy-ошибкой.""" + last_error = sqlite3.OperationalError("database is locked") + + async def always_locked(): + raise last_error + + async def fake_sleep(delay): + return None + + db._busy_retry_delays_sec = (0,) + monkeypatch.setattr("asyncio.sleep", fake_sleep) + + with pytest.raises(DatabaseBusyError) as exc_info: + await db._with_busy_retry("test op", always_locked) + + assert exc_info.value.__cause__ is last_error + + @pytest.mark.anyio async def test_account_session_encrypted_at_rest(tmp_path): db_path = str(tmp_path / "encrypted.db") diff --git a/tests/test_telegram_command_dispatcher.py b/tests/test_telegram_command_dispatcher.py index b96fdc1b..c6d393d6 100644 --- a/tests/test_telegram_command_dispatcher.py +++ b/tests/test_telegram_command_dispatcher.py @@ -1537,6 +1537,74 @@ async def test_run_loop_cancelled_reraises_when_update_busy(): db.repos.telegram_commands.update_command.assert_awaited_once() +# --- характеризующие тесты _update_command_safely (#1132: ручной цикл → tenacity) --- +# Фиксируют точный контракт retry до/после замены реализации: экспоненциальную +# лестницу задержек с потолком, retry_busy=False и заглатывание прочих ошибок. + + +def _busy_error(): + from src.database import DatabaseBusyError + + return DatabaseBusyError("Database is busy. Retry the request in a few seconds.") + + +async def test_update_command_safely_busy_backoff_doubles_and_caps(): + """Лестница задержек: INITIAL, 2×, 4×… с потолком MAX; попыток — до успеха.""" + db = _mock_db() + d = _dispatcher(db=db, pool=_mock_pool()) + failures = 6 + db.repos.telegram_commands.update_command = AsyncMock( + side_effect=[_busy_error()] * failures + [None] + ) + + with patch.object(mod.asyncio, "sleep", new_callable=AsyncMock) as mock_sleep: + await d._update_command_safely( + 7, status=TelegramCommandStatus.SUCCEEDED, log_action="succeeded" + ) + + expected = [] + delay = mod.COMMAND_STATUS_UPDATE_BUSY_RETRY_INITIAL_SEC + for _ in range(failures): + expected.append(delay) + delay = min(delay * 2, mod.COMMAND_STATUS_UPDATE_BUSY_RETRY_MAX_SEC) + + assert [c.args[0] for c in mock_sleep.await_args_list] == expected + assert db.repos.telegram_commands.update_command.await_count == failures + 1 + + +async def test_update_command_safely_busy_without_retry_is_single_shot(): + """retry_busy=False: одна попытка, без sleep, busy заглатывается.""" + db = _mock_db() + d = _dispatcher(db=db, pool=_mock_pool()) + db.repos.telegram_commands.update_command = AsyncMock(side_effect=_busy_error()) + + with patch.object(mod.asyncio, "sleep", new_callable=AsyncMock) as mock_sleep: + await d._update_command_safely( + 7, + status=TelegramCommandStatus.FAILED, + log_action="failed", + retry_busy=False, + ) + + db.repos.telegram_commands.update_command.assert_awaited_once() + mock_sleep.assert_not_awaited() + + +async def test_update_command_safely_swallows_generic_error_without_retry(): + """Не-busy ошибка: одна попытка, без sleep, наружу не пробрасывается.""" + db = _mock_db() + d = _dispatcher(db=db, pool=_mock_pool()) + db.repos.telegram_commands.update_command = AsyncMock(side_effect=RuntimeError("boom")) + + with patch.object(mod.asyncio, "sleep", new_callable=AsyncMock) as mock_sleep: + await d._update_command_safely( + 7, status=TelegramCommandStatus.SUCCEEDED, log_action="succeeded" + ) + + db.repos.telegram_commands.update_command.assert_awaited_once() + mock_sleep.assert_not_awaited() + + # ============================================================ # Additional tests for handler edge paths # ============================================================