- shared/ = portable Android-origin sources vendored from deferred/desktop-server (app/build.gradle.kts srcDir repointed; PlaybackController.kt excluded as Android-only) - backend/ = bundled-lite engine (SQLite + inline queue); .venv symlinked from the old checkout, PYTHONPATH pins THIS backend's code over any editable install - repoRoot() resolves this project dir (env SHONAR_REPO still wins); desktop-dev.sh watches shared/ + backend/ - Verified: :app:compileKotlin + :app:test green (23 tests); engine boots on :8010, self-migrates, /healthz ok
76 lines
2.5 KiB
Python
76 lines
2.5 KiB
Python
"""Async SQLAlchemy engine/session management."""
|
|
|
|
from __future__ import annotations
|
|
|
|
from collections.abc import AsyncIterator
|
|
|
|
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine
|
|
|
|
from shonar.core.config import get_settings
|
|
|
|
_engine = None
|
|
_session_factory: async_sessionmaker[AsyncSession] | None = None
|
|
|
|
|
|
def get_engine():
|
|
global _engine, _session_factory
|
|
if _engine is None:
|
|
settings = get_settings()
|
|
kwargs: dict = {"pool_pre_ping": True}
|
|
url = settings.database_url
|
|
# SQLite (tests, desktop bundled engine) does not support the pg
|
|
# pool sizing kwargs.
|
|
if url.startswith("sqlite"):
|
|
kwargs = {}
|
|
else:
|
|
kwargs.update(pool_size=settings.db_pool_size, max_overflow=settings.db_max_overflow)
|
|
_engine = create_async_engine(url, **kwargs)
|
|
if url.startswith("sqlite"):
|
|
_configure_sqlite(_engine)
|
|
_session_factory = async_sessionmaker(_engine, expire_on_commit=False)
|
|
return _engine
|
|
|
|
|
|
def _configure_sqlite(engine) -> None: # noqa: ANN001
|
|
"""Per-connection SQLite pragmas.
|
|
|
|
foreign_keys is OFF by default in SQLite; the schema leans on
|
|
ON DELETE CASCADE, so every connection must enable it. WAL lets the
|
|
inline-queue writer and API readers coexist without SQLITE_BUSY
|
|
storms on the desktop box.
|
|
"""
|
|
from sqlalchemy import event
|
|
|
|
@event.listens_for(engine.sync_engine, "connect")
|
|
def _pragmas(dbapi_conn, _record): # noqa: ANN001
|
|
cur = dbapi_conn.cursor()
|
|
cur.execute("PRAGMA foreign_keys=ON")
|
|
cur.execute("PRAGMA journal_mode=WAL")
|
|
cur.execute("PRAGMA busy_timeout=5000")
|
|
cur.close()
|
|
|
|
|
|
async def dispose_engine() -> None:
|
|
global _engine, _session_factory
|
|
if _engine is not None:
|
|
await _engine.dispose()
|
|
_engine = None
|
|
_session_factory = None
|
|
|
|
|
|
async def get_session() -> AsyncIterator[AsyncSession]:
|
|
"""FastAPI dependency yielding a database session."""
|
|
assert _session_factory is not None, "engine not initialised"
|
|
async with _session_factory() as session:
|
|
try:
|
|
yield session
|
|
await session.commit()
|
|
except Exception:
|
|
await session.rollback()
|
|
raise
|
|
|
|
|
|
def session_factory() -> async_sessionmaker[AsyncSession]:
|
|
"""Shareable session factory for the worker (outside requests)."""
|
|
assert _session_factory is not None, "engine not initialised"
|
|
return _session_factory
|