certflow.utils.retry 源代码

"""SQLite 写操作退避重试工具

网络盘(坚果云等)上的 SQLite 文件在并发读写时极易出现
``database is locked`` / ``database table is locked`` 错误。本项目里
UI 主线程(保存订正、写打印日志)与队列工作线程(写 print_jobs)
会同时访问同一个 .db 文件,普通 sqlite 引擎无 ``busy_timeout``
时一旦撞车就直接抛锁异常。

本模块提供:
- :func:`retry_on_locked`:对可调用对象做指数退避重试,专门消化瞬时写锁
  (默认 12 次 / 上限 ~8s)。配合 :mod:`certflow.utils.database` 里的
  ``WAL + busy_timeout`` PRAGMA 使用,可彻底消除打印流程中的写保护锁。

Examples:
    >>> from certflow.utils.retry import retry_on_locked
    >>> retry_on_locked(session.commit)
"""

from __future__ import annotations

import time
from collections.abc import Callable
from typing import Any

from loguru import logger
from sqlalchemy.exc import OperationalError

# SQLite 锁相关错误关键字(含中文 "database is locked" 的英文原文)
_LOCK_ERROR_HINTS = (
    "database is locked",
    "database table is locked",
    "locked",
)

# 默认退避参数:12 次、base=0.05s、上限 0.8s → 累计约 8s 上限
_DEFAULT_MAX_ATTEMPTS = 12
_DEFAULT_BASE_DELAY = 0.05
_DEFAULT_MAX_DELAY = 0.8


def _is_lock_error(exc: Exception) -> bool:
    """判断异常是否为 SQLite 写锁相关(含 OperationalError / 包装异常)。"""
    msg = str(exc).lower()
    if isinstance(exc, OperationalError) and any(h in msg for h in _LOCK_ERROR_HINTS):
        return True
    # SQLAlchemy 有时会把底层异常包成更通用的类型,做一层兜底匹配
    return any(h in msg for h in _LOCK_ERROR_HINTS)


[文档] def retry_on_locked[ T, ]( fn: Callable[..., T], *args: Any, max_attempts: int = _DEFAULT_MAX_ATTEMPTS, base_delay: float = _DEFAULT_BASE_DELAY, max_delay: float = _DEFAULT_MAX_DELAY, silent: bool = False, **kwargs: Any, ) -> T: """执行可调用对象,遇 SQLite 写锁则指数退避重试。 仅对写锁类异常(``database is locked`` 等)重试;其他异常原样抛出。 适合包裹 ``session.commit()`` / ``session.flush()`` 这类幂等写库操作。 Args: fn: 待执行的可调用对象(如 ``session.commit``)。 *args / **kwargs: 透传给 ``fn`` 的位置/关键字参数。 max_attempts: 最大尝试次数(含首次),默认 12。 base_delay: 初始退避秒数,默认 0.05。 max_delay: 单次退避上限秒数,默认 0.8。 silent: 为 True 时连最后仍失败也吞掉异常返回 None(不推荐常开)。 Returns: ``fn`` 的返回值。 Examples: >>> retry_on_locked(session.commit) >>> retry_on_locked(session.flush) """ last_exc: Exception | None = None for attempt in range(1, max_attempts + 1): try: return fn(*args, **kwargs) except Exception as exc: # noqa: BLE001 - 仅对锁相关重试 if not _is_lock_error(exc): raise last_exc = exc if attempt >= max_attempts: break delay = min(base_delay * (2 ** (attempt - 1)), max_delay) if not silent: logger.warning( f"[DB-RETRY] 写库遇锁,第 {attempt}/{max_attempts} 次重试 " f"({delay:.2f}s 后),原因: {exc}" ) time.sleep(delay) if last_exc is not None: if silent: logger.error(f"[DB-RETRY] 重试耗尽,静默放弃: {last_exc}") return None # type: ignore[return-value] raise last_exc # 理论不可达:循环至少执行一次,last_exc 必然非 None 或已 return raise RuntimeError("[DB-RETRY] 重试逻辑异常:未触发任何分支") # pragma: no cover
[文档] def with_db_retry[ T, ]( max_attempts: int = _DEFAULT_MAX_ATTEMPTS, base_delay: float = _DEFAULT_BASE_DELAY, max_delay: float = _DEFAULT_MAX_DELAY, ) -> Callable[[Callable[..., T]], Callable[..., T]]: """装饰器版本:把含写库的函数包进 :func:`retry_on_locked`。 Examples: >>> @with_db_retry() ... def save(session): ... session.commit() """ def _deco(fn: Callable[..., T]) -> Callable[..., T]: def _wrapper(*args: Any, **kwargs: Any) -> T: return retry_on_locked( fn, *args, max_attempts=max_attempts, base_delay=base_delay, max_delay=max_delay, **kwargs, ) return _wrapper return _deco