fix: serialize schema init for postgres
This commit is contained in:
+2
-1
@@ -3,9 +3,10 @@ from collections.abc import Generator
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.core.database import get_database_settings
|
||||
from packages.adapters.sqlalchemy_impl import build_session_factory, initialize_database
|
||||
from packages.adapters.sqlalchemy_impl import build_session_factory, ensure_database_exists, initialize_database
|
||||
|
||||
settings = get_database_settings()
|
||||
ensure_database_exists(settings.database_url)
|
||||
engine, SessionLocal = build_session_factory(
|
||||
settings.database_url,
|
||||
pool_size=settings.pool_size,
|
||||
|
||||
@@ -1,7 +1,8 @@
|
||||
from worker_app.core.config import get_settings
|
||||
from packages.adapters.sqlalchemy_impl import build_session_factory, initialize_database
|
||||
from packages.adapters.sqlalchemy_impl import build_session_factory, ensure_database_exists, initialize_database
|
||||
|
||||
settings = get_settings()
|
||||
ensure_database_exists(settings.database_url)
|
||||
engine, SessionLocal = build_session_factory(
|
||||
settings.database_url,
|
||||
pool_size=settings.database_pool_size,
|
||||
|
||||
@@ -4,7 +4,7 @@ from .asset_library_repository import SQLAlchemyAssetLibraryRepository
|
||||
from .asset_repository import SQLAlchemyAssetRepository
|
||||
from .ingest_job_repository import SQLAlchemyIngestJobRepository
|
||||
from .project_repository import SQLAlchemyProjectRepository
|
||||
from .session import Base, build_engine, build_session_factory, initialize_database
|
||||
from .session import Base, build_engine, build_session_factory, ensure_database_exists, initialize_database
|
||||
|
||||
__all__ = [
|
||||
"Base",
|
||||
@@ -14,5 +14,6 @@ __all__ = [
|
||||
"SQLAlchemyProjectRepository",
|
||||
"build_engine",
|
||||
"build_session_factory",
|
||||
"ensure_database_exists",
|
||||
"initialize_database",
|
||||
]
|
||||
|
||||
@@ -1,11 +1,15 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from sqlalchemy import create_engine
|
||||
from sqlalchemy.orm import Session, sessionmaker
|
||||
from sqlalchemy import create_engine, text
|
||||
from sqlalchemy.engine import URL, make_url
|
||||
from sqlalchemy.orm import sessionmaker
|
||||
|
||||
from packages.adapters.sqlalchemy_impl.models import Base
|
||||
|
||||
|
||||
SCHEMA_INIT_LOCK_ID = 2026061501
|
||||
|
||||
|
||||
def build_engine(
|
||||
database_url: str,
|
||||
*,
|
||||
@@ -42,5 +46,33 @@ def build_session_factory(
|
||||
return engine, session_factory
|
||||
|
||||
|
||||
def _build_admin_url(database_url: str) -> URL:
|
||||
url = make_url(database_url)
|
||||
return url.set(database="postgres")
|
||||
|
||||
|
||||
def ensure_database_exists(database_url: str) -> None:
|
||||
target_url = make_url(database_url)
|
||||
admin_engine = create_engine(_build_admin_url(database_url), isolation_level="AUTOCOMMIT")
|
||||
try:
|
||||
with admin_engine.connect() as connection:
|
||||
exists = connection.execute(
|
||||
text("SELECT 1 FROM pg_database WHERE datname = :database_name"),
|
||||
{"database_name": target_url.database},
|
||||
).scalar()
|
||||
if exists:
|
||||
return
|
||||
connection.execute(text(f'CREATE DATABASE "{target_url.database}"'))
|
||||
finally:
|
||||
admin_engine.dispose()
|
||||
|
||||
|
||||
def initialize_database(engine) -> None:
|
||||
Base.metadata.create_all(bind=engine)
|
||||
with engine.connect() as connection:
|
||||
connection.execute(text("SELECT pg_advisory_lock(:lock_id)"), {"lock_id": SCHEMA_INIT_LOCK_ID})
|
||||
try:
|
||||
Base.metadata.create_all(bind=connection)
|
||||
connection.commit()
|
||||
finally:
|
||||
connection.execute(text("SELECT pg_advisory_unlock(:lock_id)"), {"lock_id": SCHEMA_INIT_LOCK_ID})
|
||||
connection.commit()
|
||||
|
||||
Reference in New Issue
Block a user