Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 19 additions & 27 deletions scripts/reinvented_baseline.json
Original file line number Diff line number Diff line change
Expand Up @@ -2,64 +2,48 @@
"_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)",
"confidence": "med"
},
{
"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)",
Expand Down Expand Up @@ -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)",
Expand Down
89 changes: 63 additions & 26 deletions src/database/facade.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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]:
Expand Down
77 changes: 50 additions & 27 deletions src/services/telegram_command_dispatcher.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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('.', '_')}"
Expand Down
Loading
Loading