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