From 14d6e64cdd01114b15467ca0c827eb92d4acea88 Mon Sep 17 00:00:00 2001 From: Artur Shiriev Date: Sun, 27 Sep 2026 23:17:00 +0300 Subject: [PATCH] feat: route GET/HEAD requests to an optional read replica --- AGENTS.md | 12 +++--- app/ioc.py | 28 ++++++++++++-- app/resources/db.py | 33 +++++++++++++++-- app/settings.py | 1 + pyproject.toml | 3 ++ tests/conftest.py | 6 +-- tests/test_db_routing.py | 80 ++++++++++++++++++++++++++++++++++++++++ 7 files changed, 148 insertions(+), 15 deletions(-) create mode 100644 tests/test_db_routing.py diff --git a/AGENTS.md b/AGENTS.md index 56210f3..6a45924 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -21,17 +21,19 @@ Python is 3.14, dependencies managed by `uv`. The API is exposed on `:8000`. **Request flow**: `app/__main__.py` → `granian` → `app.application:build_app` (factory) → `LitestarBootstrapper` from `lite-bootstrap` wraps a `litestar.Litestar` with OpenTelemetry (asyncpg + SQLAlchemy instrumentors), Sentry, CORS, Swagger, etc., based on `Settings.api_bootstrapper_config`. **Dependency injection** (`app/ioc.py`): `modern_di.Container` is created in `build_app` with the `Dependencies` group and attached via `modern_di_litestar.ModernDIPlugin`. Route handlers receive repositories as parameters; `application.py` declares them with `modern_di_litestar.FromDI(...)` so Litestar resolves them per-request. Provider scopes: -- `database_engine` — application-scoped factory, finalizer disposes the engine -- `session` — request-scoped, finalizer closes the session +- `database_engine` — application-scoped primary engine, finalizer disposes the engine +- `database_replica_engine` — application-scoped replica engine built from `DB_REPLICA_DSN`, or `None` when it is unset +- `dynamic_engine` — request-scoped; `choose_sa_engine` returns the replica for `GET`/`HEAD` requests when one is configured, the primary otherwise (including when there is no request). GET handlers must not write, and a read right after a write may see replica lag +- `session` — request-scoped, bound to `dynamic_engine`, finalizer closes the session - `*_repository` — request-scoped, depend on `session`, configured with `auto_commit=True` (advanced-alchemy commits on success / rolls back on exception) **Persistence**: Models inherit `advanced_alchemy.base.BigIntAuditBase` (gives `id: BigInt`, `created_at`, `updated_at`). Metadata is shared with `orm.DeclarativeBase.metadata` in `app/models.py` so Alembic autogen sees everything. Repositories are `SQLAlchemyAsyncRepositoryService[Model]` with a nested `BaseRepository(SQLAlchemyAsyncRepository[Model])`. Custom `CustomAsyncSession` in `app/resources/db.py` overrides `close()` so test transactions are not actually closed when the session is bound to an `AsyncConnection` — this is what makes the per-test rollback fixture work. -**Test isolation** (`tests/conftest.py`): `db_session` fixture opens a connection, starts a transaction, starts a SAVEPOINT, then **overrides** `Dependencies.database_engine` in the DI container to return that connection. All requests in the test reuse this connection; the outer transaction is rolled back in teardown, so DB state is clean between tests with no truncation needed. `app` and `client` fixtures build the real app and run it through `httpx.ASGITransport` + `asgi_lifespan.LifespanManager`. Polyfactory `SQLAlchemyFactory` is wired up via `set_async_session_in_base_sqlalchemy_factory`. +**Test isolation** (`tests/conftest.py`): `db_session` fixture opens a connection, starts a transaction, starts a SAVEPOINT, then **overrides** `Dependencies.dynamic_engine` in the DI container to return that connection. All requests in the test reuse this connection; the outer transaction is rolled back in teardown, so DB state is clean between tests with no truncation needed. `app` and `client` fixtures build the real app and run it through `httpx.ASGITransport` + `asgi_lifespan.LifespanManager`. Polyfactory `SQLAlchemyFactory` is wired up via `set_async_session_in_base_sqlalchemy_factory`. **Migrations**: `migrations/env.py` reads `app.models.METADATA` and rewrites the DSN driver from `postgresql+asyncpg` → `postgresql` (Alembic uses sync psycopg2). Always run autogen against an upgraded DB — `just migration` enforces this. -**Settings** (`app/settings.py`): `pydantic_settings.BaseSettings` reads from env vars (see `docker-compose.yml` for `SERVICE_DEBUG`, `SERVICE_ENVIRONMENT`, `DB_DSN`). `api_bootstrapper_config` builds the `LitestarConfig` consumed by `lite-bootstrap`. +**Settings** (`app/settings.py`): `pydantic_settings.BaseSettings` reads from env vars (see `docker-compose.yml` for `SERVICE_DEBUG`, `SERVICE_ENVIRONMENT`, `DB_DSN`; `DB_REPLICA_DSN` is optional). `api_bootstrapper_config` builds the `LitestarConfig` consumed by `lite-bootstrap`. ## Conventions @@ -39,7 +41,7 @@ Python is 3.14, dependencies managed by `uv`. The API is exposed on `:8000`. - Pydantic schemas in `app/schemas.py` use `from_attributes=True` (via `Base`) so they validate directly from ORM instances (`schemas.X.model_validate(orm_instance)`). Collection responses go through `Collection[T].from_models(...)` (e.g. `schemas.Decks`, `schemas.Cards`). - Deck responses are deliberately two-shaped: `list_decks`/`create_deck`/`update_deck` return the light `schemas.Deck` (no `cards`), while `get_deck` returns `schemas.DeckWithCards`. The split mirrors loading — lists use `noload` (no cards query), detail uses `selectinload` via `fetch_with_cards` — so the type states exactly what each endpoint loads. - Domain exceptions: register handlers in `application.build_app`'s `exception_handlers` dict (see `DuplicateKeyError` → `exceptions.duplicate_key_error_handler`). For per-handler 404s the code raises `litestar.exceptions.HTTPException` directly. -- `ruff` is configured with `select = ["ALL"]` and a line length of 120 — expect strict lint. Type-check with `ty`; suppressions already exist for `invalid-argument-type` around `LifespanManager` / `ASGITransport` / DTO list construction. +- `ruff` is configured with `select = ["ALL"]` and a line length of 120 — expect strict lint. `app/resources/` is exempt from the `TC` rules because modern-di resolves DI creator annotations at runtime. Type-check with `ty`; suppressions already exist for `invalid-argument-type` around `LifespanManager` / `ASGITransport` / DTO list construction. ## Agent skills diff --git a/app/ioc.py b/app/ioc.py index aa1df46..6190593 100644 --- a/app/ioc.py +++ b/app/ioc.py @@ -1,15 +1,37 @@ from modern_di import Group, Scope, providers from app.repositories import CardsRepository, DecksRepository -from app.resources.db import close_sa_engine, close_session, create_sa_engine, create_session +from app.resources.db import ( + choose_sa_engine, + close_sa_engine, + close_session, + create_primary_sa_engine, + create_replica_sa_engine, + create_session, +) class Dependencies(Group): database_engine = providers.Factory( - creator=create_sa_engine, cache=providers.CacheSettings(finalizer=close_sa_engine) + creator=create_primary_sa_engine, + cache=providers.CacheSettings(finalizer=close_sa_engine), + bound_type=None, + ) + database_replica_engine = providers.Factory( + creator=create_replica_sa_engine, + cache=providers.CacheSettings(finalizer=close_sa_engine), + bound_type=None, + ) + dynamic_engine = providers.Factory( + scope=Scope.REQUEST, + creator=choose_sa_engine, + kwargs={"primary_engine": database_engine, "replica_engine": database_replica_engine}, ) session = providers.Factory( - scope=Scope.REQUEST, creator=create_session, cache=providers.CacheSettings(finalizer=close_session) + scope=Scope.REQUEST, + creator=create_session, + cache=providers.CacheSettings(finalizer=close_session), + kwargs={"engine": dynamic_engine}, ) decks_repository = providers.Factory( diff --git a/app/resources/db.py b/app/resources/db.py index 99944bc..d54aacc 100644 --- a/app/resources/db.py +++ b/app/resources/db.py @@ -2,6 +2,8 @@ import logging import typing +import litestar +from sqlalchemy.engine.url import URL, make_url from sqlalchemy.ext import asyncio as sa from app.settings import settings @@ -10,9 +12,12 @@ logger = logging.getLogger(__name__) -def create_sa_engine() -> sa.AsyncEngine: +REPLICA_METHODS: typing.Final = frozenset({"GET", "HEAD"}) + + +def create_sa_engine(url: URL) -> sa.AsyncEngine: return sa.create_async_engine( - url=settings.db_dsn_parsed, + url=url, echo=settings.service_debug, echo_pool=settings.service_debug, pool_size=settings.db_pool_size, @@ -21,8 +26,28 @@ def create_sa_engine() -> sa.AsyncEngine: ) -async def close_sa_engine(engine: sa.AsyncEngine) -> None: - await engine.dispose() +def create_primary_sa_engine() -> sa.AsyncEngine: + return create_sa_engine(settings.db_dsn_parsed) + + +def create_replica_sa_engine() -> sa.AsyncEngine | None: + return create_sa_engine(make_url(settings.db_replica_dsn)) if settings.db_replica_dsn else None + + +async def close_sa_engine(engine: sa.AsyncEngine | None) -> None: + if engine: + await engine.dispose() + + +def choose_sa_engine( + *, + primary_engine: sa.AsyncEngine, + replica_engine: sa.AsyncEngine | None, + request: litestar.Request[typing.Any, typing.Any, typing.Any] | None = None, +) -> sa.AsyncEngine: + if replica_engine and request and request.method in REPLICA_METHODS: + return replica_engine + return primary_engine def create_session(engine: sa.AsyncEngine) -> sa.AsyncSession: diff --git a/app/settings.py b/app/settings.py index 1d74bdc..54ee2c0 100644 --- a/app/settings.py +++ b/app/settings.py @@ -11,6 +11,7 @@ class Settings(pydantic_settings.BaseSettings): log_level: str = "info" db_dsn: str = "postgresql+asyncpg://postgres:password@db/postgres" + db_replica_dsn: str = "" db_pool_size: int = 5 db_max_overflow: int = 0 db_pool_pre_ping: bool = True diff --git a/pyproject.toml b/pyproject.toml index b012d20..f2e42ec 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -67,6 +67,9 @@ isort.lines-after-imports = 2 isort.no-lines-before = ["standard-library", "local-folder"] [tool.ruff.lint.extend-per-file-ignores] +"app/resources/*.py" = [ + "TC", # modern-di reads DI creator annotations at runtime +] "tests/*.py" = [ "S101", # allow asserts ] diff --git a/tests/conftest.py b/tests/conftest.py index 1845ba7..64fd022 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -9,7 +9,7 @@ from app import ioc from app.application import build_app -from app.resources.db import create_sa_engine +from app.resources.db import create_primary_sa_engine if typing.TYPE_CHECKING: @@ -44,10 +44,10 @@ async def di_container(app: litestar.Litestar) -> typing.AsyncIterator[modern_di @pytest.fixture async def db_session(di_container: modern_di.Container) -> typing.AsyncIterator[AsyncSession]: - engine = create_sa_engine() + engine = create_primary_sa_engine() connection = await engine.connect() transaction = await connection.begin() - di_container.override(ioc.Dependencies.database_engine, connection) + di_container.override(ioc.Dependencies.dynamic_engine, connection) try: yield AsyncSession( diff --git a/tests/test_db_routing.py b/tests/test_db_routing.py new file mode 100644 index 0000000..8cd0457 --- /dev/null +++ b/tests/test_db_routing.py @@ -0,0 +1,80 @@ +import typing + +import litestar +import pytest +from litestar.enums import HttpMethod +from litestar.testing import RequestFactory +from modern_di import Scope + +from app import ioc +from app.resources.db import create_replica_sa_engine, create_sa_engine +from app.settings import settings + + +if typing.TYPE_CHECKING: + import modern_di + from sqlalchemy.ext.asyncio import AsyncEngine + + +@pytest.fixture +async def primary_engine(di_container: modern_di.Container) -> typing.AsyncIterator[AsyncEngine]: + engine = create_sa_engine(settings.db_dsn_parsed) + di_container.override(ioc.Dependencies.database_engine, engine) + yield engine + await engine.dispose() + + +@pytest.fixture +async def replica_engine(di_container: modern_di.Container) -> typing.AsyncIterator[AsyncEngine]: + engine = create_sa_engine(settings.db_dsn_parsed) + di_container.override(ioc.Dependencies.database_replica_engine, engine) + yield engine + await engine.dispose() + + +def build_request(method: HttpMethod) -> litestar.Request[typing.Any, typing.Any, typing.Any]: + request = RequestFactory().get("/") + request.scope["method"] = method + return request + + +def resolve_engine(di_container: modern_di.Container, method: HttpMethod | None) -> AsyncEngine: + context = {litestar.Request: build_request(method)} if method else None + with di_container.build_child_container(scope=Scope.REQUEST, context=context) as request_container: + return request_container.resolve_provider(ioc.Dependencies.dynamic_engine) + + +@pytest.mark.parametrize("method", [HttpMethod.GET, HttpMethod.HEAD]) +def test_safe_methods_use_replica( + di_container: modern_di.Container, primary_engine: AsyncEngine, replica_engine: AsyncEngine, method: HttpMethod +) -> None: + engine = resolve_engine(di_container, method) + assert engine is replica_engine + assert engine is not primary_engine + + +@pytest.mark.parametrize("method", [HttpMethod.POST, HttpMethod.PUT, HttpMethod.PATCH, HttpMethod.DELETE]) +def test_write_methods_use_primary( + di_container: modern_di.Container, primary_engine: AsyncEngine, replica_engine: AsyncEngine, method: HttpMethod +) -> None: + engine = resolve_engine(di_container, method) + assert engine is primary_engine + assert engine is not replica_engine + + +@pytest.mark.usefixtures("replica_engine") +def test_no_request_uses_primary(di_container: modern_di.Container, primary_engine: AsyncEngine) -> None: + assert resolve_engine(di_container, None) is primary_engine + + +def test_no_replica_configured_uses_primary(di_container: modern_di.Container, primary_engine: AsyncEngine) -> None: + assert di_container.resolve_provider(ioc.Dependencies.database_replica_engine) is None + assert resolve_engine(di_container, HttpMethod.GET) is primary_engine + + +async def test_replica_engine_built_from_replica_dsn(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setattr(settings, "db_replica_dsn", "postgresql+asyncpg://postgres:password@replica/postgres") + engine = create_replica_sa_engine() + assert engine + assert engine.url.host == "replica" + await engine.dispose()