diff --git a/app/api/dependencies.py b/app/api/dependencies.py index 772f6c7..bed83c2 100644 --- a/app/api/dependencies.py +++ b/app/api/dependencies.py @@ -23,8 +23,18 @@ def extract_bearer_token(request: Request) -> str: return "" -def get_db(request: Request) -> Database: - return request.app.state.db +async def get_db(): + """每请求一个独立 Database/AsyncSession(绑定共享引擎)。 + + 之前直接返回 app.state.db(单例 AsyncSession)被并发请求复用,会间歇报 + ``sqlalchemy.exc.InvalidRequestError: concurrent operations are not permitted``。 + 改为每请求各建一个 Database 并在请求结束关闭,根除该竞态。 + """ + db = Database() + try: + yield db + finally: + await db.close() async def _resolve_identity(user: dict, db: Database) -> dict: diff --git a/app/infrastructure/repositories.py b/app/infrastructure/repositories.py index f18bdc3..ab48849 100644 --- a/app/infrastructure/repositories.py +++ b/app/infrastructure/repositories.py @@ -2102,17 +2102,33 @@ class StatsRepository: # --------------------------------------------------------------------------- # 数据库门面 # --------------------------------------------------------------------------- +# 共享异步引擎(按 url 缓存,进程内复用)—— 每请求一个独立 Database/AsyncSession, +# 引擎跨请求共享,避免「同一 AsyncSession 被并发请求复用」的竞态。 +_SHARED_ENGINES: dict[str, "object"] = {} + + +def _shared_engine(url: str): + eng = _SHARED_ENGINES.get(url) + if eng is None: + from .db import make_async_engine + eng = make_async_engine(url) + _SHARED_ENGINES[url] = eng + return eng + + class Database: - """持有全部 Repository,统一访问入口(基础设施层)。""" + """持有全部 Repository,统一访问入口(基础设施层)。 + + ``session is None`` 时新建一个独立 AsyncSession(绑定共享引擎)—— 供每请求 + 各持一个 Database,避免并发复用同一 session 抛 ``concurrent operations are not permitted``。 + """ def __init__(self, db_url: str | None = None, session: AsyncSession | None = None): if session is None: from sqlalchemy.ext.asyncio import async_sessionmaker - from .db import make_async_engine - url = db_url or config.DATABASE_URL - self._engine = make_async_engine(url) + self._engine = _shared_engine(url) self._session_factory = async_sessionmaker( bind=self._engine, class_=AsyncSession, autoflush=False, expire_on_commit=False, )