From 8332511054a25c1c63873a2947d0d7c33b2c4082 Mon Sep 17 00:00:00 2001 From: baekyutae Date: Tue, 14 Jul 2026 12:29:48 +0900 Subject: [PATCH] =?UTF-8?q?feat(repo):=20DB=20=EC=97=B0=EA=B2=B0=20?= =?UTF-8?q?=EC=9E=AC=EC=82=AC=EC=9A=A9=20=EC=A0=84=20=EC=9C=A0=ED=9A=A8?= =?UTF-8?q?=EC=84=B1=20=EA=B2=80=EC=82=AC=EB=A5=BC=20=EC=B6=94=EA=B0=80(po?= =?UTF-8?q?ol=5Fpre=5Fping=3DTrue)=20(#100)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- services/core-api/src/core/dependencies.py | 6 +- .../tests/integration/test_db_pre_ping.py | 66 +++++++++++++++++++ .../feedback-loop-pipeline/src/bootstrap.py | 6 +- .../src/infra/db/session.py | 2 +- services/pipeline-worker/src/bootstrap.py | 6 +- .../search-service/src/infra/db/session.py | 2 +- 6 files changed, 81 insertions(+), 7 deletions(-) create mode 100644 services/core-api/tests/integration/test_db_pre_ping.py diff --git a/services/core-api/src/core/dependencies.py b/services/core-api/src/core/dependencies.py index ed9b117..b3d15f6 100644 --- a/services/core-api/src/core/dependencies.py +++ b/services/core-api/src/core/dependencies.py @@ -30,7 +30,11 @@ def build_dependency_container(settings: Settings | None = None) -> DependencyCo def _build_db_session_factory(settings: Settings) -> async_sessionmaker[AsyncSession]: - engine = create_async_engine(settings.database_url, future=True) + engine = create_async_engine( + settings.database_url, + future=True, + pool_pre_ping=True, + ) return async_sessionmaker(engine, expire_on_commit=False) diff --git a/services/core-api/tests/integration/test_db_pre_ping.py b/services/core-api/tests/integration/test_db_pre_ping.py new file mode 100644 index 0000000..e918d1b --- /dev/null +++ b/services/core-api/tests/integration/test_db_pre_ping.py @@ -0,0 +1,66 @@ +import asyncio + +import pytest +from sqlalchemy import text +from sqlalchemy.ext.asyncio import AsyncEngine, create_async_engine + +from src.core.config import Settings +from src.core.dependencies import _build_db_session_factory + + +@pytest.mark.asyncio +async def test_pre_ping_replaces_terminated_pooled_connection( + postgres_url: str, +) -> None: + settings = Settings( + _env_file=None, + GCP_PROJECT_ID="test-project", + GCS_VIDEO_BUCKET_NAME="test-bucket", + JWT_SECRET_KEY="test-jwt-key", + DATABASE_URL=postgres_url, + ) + session_factory = _build_db_session_factory(settings) + app_engine = session_factory.kw["bind"] + assert isinstance(app_engine, AsyncEngine) + + admin_engine = create_async_engine(postgres_url) + connection_terminated = asyncio.Event() + + try: + async with session_factory() as session: + terminated_backend_pid = await session.scalar( + text("SELECT pg_backend_pid()") + ) + + app_connection = await session.connection() + pooled_connection = await app_connection.get_raw_connection() + driver_connection = pooled_connection.driver_connection + driver_connection.add_termination_listener( + lambda _: connection_terminated.set() + ) + + assert terminated_backend_pid is not None + + async with admin_engine.begin() as connection: + terminated = await connection.scalar( + text("SELECT pg_terminate_backend(:pid)"), + {"pid": terminated_backend_pid}, + ) + + assert terminated is True + + await asyncio.wait_for( + connection_terminated.wait(), + timeout=1.0, + ) + + async with session_factory() as session: + replacement_backend_pid = await session.scalar( + text("SELECT pg_backend_pid()") + ) + + assert replacement_backend_pid is not None + assert replacement_backend_pid != terminated_backend_pid + finally: + await app_engine.dispose() + await admin_engine.dispose() diff --git a/services/feedback-loop-pipeline/src/bootstrap.py b/services/feedback-loop-pipeline/src/bootstrap.py index 17b80b2..286ac6a 100644 --- a/services/feedback-loop-pipeline/src/bootstrap.py +++ b/services/feedback-loop-pipeline/src/bootstrap.py @@ -9,7 +9,7 @@ from uuid import UUID, uuid4 import asyncpg -from sqlalchemy.ext.asyncio import AsyncEngine, create_async_engine +from sqlalchemy.ext.asyncio import AsyncEngine from src.config.settings import Settings from src.dataset.batch import DatasetBatchService @@ -23,7 +23,7 @@ from src.evaluation.evaluator import OfflineEvaluator from src.infra.db.legacy_reindex_lock import PostgresAdvisoryLegacyReindexLock from src.infra.db.legacy_reindex_store import LegacyReindexStore, VectorIndexCatalogStore -from src.infra.db.session import create_session_factory +from src.infra.db.session import create_db_engine, create_session_factory from src.infra.db.snapshot_restore import CatalogSnapshotIndexRestore from src.infra.db.stores import ( DbChunkTextSnapshot, @@ -453,7 +453,7 @@ async def _run_consumer( async def _build_runtime_context(settings: Settings): - engine = create_async_engine(settings.database_url, future=True) + engine = create_db_engine(settings.database_url) session_factory = create_session_factory(engine) pgmq_pool = None if settings.broker_type == "pgmq": diff --git a/services/feedback-loop-pipeline/src/infra/db/session.py b/services/feedback-loop-pipeline/src/infra/db/session.py index 38d3552..af143c9 100644 --- a/services/feedback-loop-pipeline/src/infra/db/session.py +++ b/services/feedback-loop-pipeline/src/infra/db/session.py @@ -2,7 +2,7 @@ def create_db_engine(database_url: str) -> AsyncEngine: - return create_async_engine(database_url, future=True) + return create_async_engine(database_url, future=True, pool_pre_ping=True) def create_session_factory(engine: AsyncEngine) -> async_sessionmaker[AsyncSession]: diff --git a/services/pipeline-worker/src/bootstrap.py b/services/pipeline-worker/src/bootstrap.py index 377f15f..094cb1a 100644 --- a/services/pipeline-worker/src/bootstrap.py +++ b/services/pipeline-worker/src/bootstrap.py @@ -98,7 +98,11 @@ async def create_production_bootstrap(settings: Settings) -> None: log = get_logger().bind(trace_id="-", video_id="-", user_id="-") # --- DB --- - engine = create_async_engine(settings.database_url, future=True) + engine = create_async_engine( + settings.database_url, + future=True, + pool_pre_ping=True, + ) session_factory = async_sessionmaker(engine, expire_on_commit=False) video_repo = VideoRepository( session_factory, diff --git a/services/search-service/src/infra/db/session.py b/services/search-service/src/infra/db/session.py index dc84689..87af6fd 100644 --- a/services/search-service/src/infra/db/session.py +++ b/services/search-service/src/infra/db/session.py @@ -2,7 +2,7 @@ def create_engine(database_url: str) -> AsyncEngine: - return create_async_engine(database_url, future=True) + return create_async_engine(database_url, future=True, pool_pre_ping=True) def create_session_factory(engine: AsyncEngine) -> async_sessionmaker[AsyncSession]: