SPB Git forge

spb/ai-atlas

Public
41commits 1branches 0releases
4.6 MBsize
maindefault branch
12 days agolast push
HTML 77.2% TypeScript 10.5% Python 9.6% JavaScript 2.5%
2.4 KB · 76 lines python
Raw Blame History
1"""Database access: SQLAlchemy Core (async, asyncpg) with plain SQL. One engine per process."""2from __future__ import annotations34import json5from collections.abc import AsyncIterator, Mapping, Sequence6from contextlib import asynccontextmanager7from typing import Any89from sqlalchemy import text10from sqlalchemy.ext.asyncio import AsyncConnection, AsyncEngine, create_async_engine1112from aiatlas.config import settings1314_engine: AsyncEngine | None = None151617def engine() -> AsyncEngine:18    global _engine19    if _engine is None:20        _engine = create_async_engine(settings.database_url, pool_size=8, max_overflow=8, pool_pre_ping=True, pool_recycle=1800,21                                      connect_args={"server_settings": {"application_name": "aiatlas", "jit": "off"}})22    return _engine232425async def dispose() -> None:26    global _engine27    if _engine is not None:28        await _engine.dispose()29        _engine = None303132@asynccontextmanager33async def connection() -> AsyncIterator[AsyncConnection]:34    async with engine().connect() as conn:35        yield conn363738@asynccontextmanager39async def transaction() -> AsyncIterator[AsyncConnection]:40    async with engine().begin() as conn:41        yield conn424344def jsonb(value: Any) -> str:45    """Serialise a Python value for a `cast(:x as jsonb)` parameter."""46    return json.dumps(value, default=str, ensure_ascii=False)474849async def execute(conn: AsyncConnection, sql: str, /, **params: Any) -> None:50    await conn.execute(text(sql), params)515253async def fetch_all(conn: AsyncConnection, sql: str, /, **params: Any) -> list[dict[str, Any]]:54    result = await conn.execute(text(sql), params)55    return [dict(r._mapping) for r in result]565758async def fetch_one(conn: AsyncConnection, sql: str, /, **params: Any) -> dict[str, Any] | None:59    result = await conn.execute(text(sql), params)60    row = result.first()61    return dict(row._mapping) if row is not None else None626364async def fetch_val(conn: AsyncConnection, sql: str, /, **params: Any) -> Any:65    result = await conn.execute(text(sql), params)66    row = result.first()67    return row[0] if row is not None else None686970async def execute_many(conn: AsyncConnection, sql: str, rows: Sequence[Mapping[str, Any]]) -> None:71    if rows:72        await conn.execute(text(sql), list(rows))737475__all__ = ["connection", "dispose", "engine", "execute", "execute_many", "fetch_all", "fetch_one", "fetch_val", "jsonb", "transaction"]76