diff --git a/deployments/charts/openg2p-master-data/values.yaml b/deployments/charts/openg2p-master-data/values.yaml index b0e502d..b19b7cd 100644 --- a/deployments/charts/openg2p-master-data/values.yaml +++ b/deployments/charts/openg2p-master-data/values.yaml @@ -11,6 +11,11 @@ global: masterDataDBUser: '{{ if eq .Release.Name "master-data" }}master_data_user{{ else }}{{ printf "%s_master_data_user" .Release.Name | replace "-" "_" }}{{ end }}' masterDataDBSecret: '{{ if eq .Release.Name "master-data" }}master-data{{ else }}{{ .Release.Name }}-master-data{{ end }}' masterDataDBUserPasswordKey: '{{ if eq .Release.Name "master-data" }}master-data-db-user{{ else }}{{ .Release.Name }}-master-data-db-user{{ end }}' + # SQLAlchemy pool per gunicorn worker process. Defaults match fastapi-common. + dbPoolSize: 5 + dbPoolMaxOverflow: 10 + dbPoolPrePing: true + dbPoolRecycle: 1800 # IAM (JWT validation + DP_ role -> permission resolution). Consumed by # master-data Settings (MASTER_DATA_API_*); override per environment. @@ -183,6 +188,10 @@ masterDataAPI: MASTER_DATA_API_DB_DBNAME: '{{ tpl .Values.global.masterDataDB $ }}' MASTER_DATA_API_DB_PORT: '{{ .Values.global.masterDataDBPort }}' MASTER_DATA_API_DB_USERNAME: '{{ tpl .Values.global.masterDataDBUser $ }}' + MASTER_DATA_API_DB_POOL_SIZE: '{{ .Values.global.dbPoolSize }}' + MASTER_DATA_API_DB_POOL_MAX_OVERFLOW: '{{ .Values.global.dbPoolMaxOverflow }}' + MASTER_DATA_API_DB_POOL_PRE_PING: '{{ .Values.global.dbPoolPrePing }}' + MASTER_DATA_API_DB_POOL_RECYCLE: '{{ .Values.global.dbPoolRecycle }}' MASTER_DATA_API_PORT: '{{ .Values.containerPort }}' MASTER_DATA_API_OPENAPI_ROOT_PATH: '{{ .Values.openapiRootPath }}' # Auth / CSRF (iam-core fields inherited on master-data Settings). diff --git a/master-data-api/.env.example b/master-data-api/.env.example index 7143ca9..bd8cca7 100644 --- a/master-data-api/.env.example +++ b/master-data-api/.env.example @@ -7,6 +7,11 @@ MASTER_DATA_API_DB_HOSTNAME=localhost MASTER_DATA_API_DB_PORT=5432 MASTER_DATA_API_DB_USERNAME=postgres MASTER_DATA_API_DB_PASSWORD=password +# SQLAlchemy pool per process. Defaults are 5 / 10. +MASTER_DATA_API_DB_POOL_SIZE=5 +MASTER_DATA_API_DB_POOL_MAX_OVERFLOW=10 +MASTER_DATA_API_DB_POOL_PRE_PING=true +MASTER_DATA_API_DB_POOL_RECYCLE=1800 MASTER_DATA_API_OPENAPI_ROOT_PATH=/ MASTER_DATA_API_PORT=8040 diff --git a/master-data-api/src/openg2p_gen2_master_data/engine.py b/master-data-api/src/openg2p_gen2_master_data/engine.py deleted file mode 100644 index cfd3477..0000000 --- a/master-data-api/src/openg2p_gen2_master_data/engine.py +++ /dev/null @@ -1,42 +0,0 @@ -"""Database engine and session management.""" - -import logging - -from sqlalchemy.ext.asyncio import AsyncEngine, AsyncSession, async_sessionmaker, create_async_engine -from openg2p_fastapi_common.context import dbengine - -_logger = logging.getLogger("master-data-engine") - -_engine: AsyncEngine | None = None - - -def get_engine() -> AsyncEngine: - """ - Get the master-data database engine. - - First tries the framework's context variable, then falls back to - a module-level cached engine created from config. - """ - global _engine - - engine = dbengine.get() - if engine is not None: - return engine - - if _engine is not None: - return _engine - - from .config import Settings - - config = Settings.get_config() - - if config.db_datasource: - _engine = create_async_engine(config.db_datasource, echo=config.db_logging) - return _engine - - raise RuntimeError("Database not configured. Check db_datasource in settings.") - - -def get_session_maker() -> async_sessionmaker[AsyncSession]: - """Get an async session maker bound to the master-data database engine.""" - return async_sessionmaker(get_engine(), expire_on_commit=False) diff --git a/master-data-api/src/openg2p_gen2_master_data/services/g2p_attribute_service.py b/master-data-api/src/openg2p_gen2_master_data/services/g2p_attribute_service.py index c1fc590..6010211 100644 --- a/master-data-api/src/openg2p_gen2_master_data/services/g2p_attribute_service.py +++ b/master-data-api/src/openg2p_gen2_master_data/services/g2p_attribute_service.py @@ -2,10 +2,10 @@ import uuid from typing import List, Optional +from openg2p_fastapi_common.context import get_async_session_maker from openg2p_fastapi_common.service import BaseService from sqlalchemy import delete, func, select -from ..engine import get_session_maker from ..helpers.data_policy_helper import DataPolicyHelper from ..models import G2PAttribute, G2PAttributeValue from ..repositories import AttributeValueRepository @@ -63,7 +63,8 @@ def _to_value_data(row: G2PAttributeValue) -> AttributeValueData: ) async def get_attributes(self) -> List[AttributeData]: - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: stmt = select(G2PAttribute).order_by(G2PAttribute.attribute_id) rows = (await session.execute(stmt)).scalars().all() return [self._to_attribute_data(r) for r in rows] @@ -75,7 +76,8 @@ async def get_attribute_values( page_number: int = 1, data_policies: Optional[List[dict]] = None, ) -> tuple[List[AttributeValueData], int]: - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: policy_condition = await self._build_attribute_value_policy_condition( data_policies, session, @@ -213,7 +215,8 @@ async def add_attribute( "attribute_code and attribute_display are required", ) - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: if await self._attribute_code_exists(session, attribute_code): raise AttributeServiceError( "G2P-ATTR-409", @@ -232,7 +235,8 @@ async def add_attribute( return self._to_attribute_data(attribute) async def update_attribute(self, payload) -> AttributeData: - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: attribute = await session.get(G2PAttribute, payload.attribute_id) if not attribute: raise AttributeServiceError( @@ -303,7 +307,8 @@ async def delete_attribute( *, cascade: bool = False, ) -> str: - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: attribute = await session.get(G2PAttribute, attribute_id) if not attribute: raise AttributeServiceError( @@ -353,7 +358,8 @@ async def add_attribute_value( parent_value_id = self._empty_to_none(parent_value_id) - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: attribute = await session.get(G2PAttribute, attribute_id) if not attribute: raise AttributeServiceError( @@ -388,7 +394,8 @@ async def add_attribute_value( return self._to_value_data(value) async def update_attribute_value(self, payload) -> AttributeValueData: - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: value = await self._get_value_by_id( session, payload.value_id, @@ -464,7 +471,8 @@ async def delete_attribute_value( value_id: str, attribute_id: Optional[str] = None, ) -> tuple[str, str]: - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: value = await self._get_value_by_id(session, value_id, attribute_id) if not value: raise AttributeServiceError( diff --git a/master-data-api/src/openg2p_gen2_master_data/services/g2p_geo_service.py b/master-data-api/src/openg2p_gen2_master_data/services/g2p_geo_service.py index 25474bb..ed8f1fe 100644 --- a/master-data-api/src/openg2p_gen2_master_data/services/g2p_geo_service.py +++ b/master-data-api/src/openg2p_gen2_master_data/services/g2p_geo_service.py @@ -2,10 +2,10 @@ import uuid from typing import List, Optional -from sqlalchemy import delete, func, or_, select +from openg2p_fastapi_common.context import get_async_session_maker from openg2p_fastapi_common.service import BaseService +from sqlalchemy import delete, func, or_, select -from ..engine import get_session_maker from ..helpers.data_policy_helper import DataPolicyHelper from ..models import G2PGeoLevel, G2PGeoLevelValue from ..repositories import GeoLevelValueRepository @@ -67,7 +67,8 @@ async def get_all_geo_levels(self) -> List[GeoLevelData]: Returns: List of GeoLevelData """ - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: query = select(G2PGeoLevel) levels = (await session.execute(query)).scalars().all() return [self._to_level_data(level) for level in levels] @@ -89,7 +90,8 @@ async def get_geo_level_values( Returns: List of GeoLevelValueData """ - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: level = await session.get(G2PGeoLevel, level_id) if not level: # Fall back to the level's NAME. Callers hand-configure this — @@ -215,7 +217,8 @@ async def add_geo_level( parent_level_id = self._empty_to_none(parent_level_id) - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: if await self._mnemonic_exists(session, level_mnemonic, parent_level_id): raise GeoServiceError( "G2P-GEO-409", @@ -241,7 +244,8 @@ async def add_geo_level( return self._to_level_data(level) async def update_geo_level(self, payload) -> GeoLevelData: - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: level = await session.get(G2PGeoLevel, payload.level_id) if not level: raise GeoServiceError("G2P-GEO-404", f"level_id not found: {payload.level_id}") @@ -286,7 +290,8 @@ async def update_geo_level(self, payload) -> GeoLevelData: return self._to_level_data(level) async def delete_geo_level(self, level_id: str) -> str: - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: level = await session.get(G2PGeoLevel, level_id) if not level: raise GeoServiceError("G2P-GEO-404", f"level_id not found: {level_id}") @@ -336,7 +341,8 @@ async def add_geo_level_value( parent_level_value_id = self._empty_to_none(parent_level_value_id) - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: level = await session.get(G2PGeoLevel, level_id) if not level: raise GeoServiceError("G2P-GEO-404", f"level_id not found: {level_id}") @@ -369,7 +375,8 @@ async def add_geo_level_value( return self._to_value_data(value) async def update_geo_level_value(self, payload) -> GeoLevelValueData: - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: value = await session.get(G2PGeoLevelValue, payload.level_value_id) if not value: raise GeoServiceError( @@ -435,7 +442,8 @@ async def delete_geo_level_value( *, cascade: bool = False, ) -> str: - async with get_session_maker()() as session: + session_maker = get_async_session_maker() + async with session_maker() as session: value = await session.get(G2PGeoLevelValue, level_value_id) if not value: raise GeoServiceError(