From e786de7349a37bb09a50abcb008e1f62d8911301 Mon Sep 17 00:00:00 2001 From: Xiaoxia AI Date: Mon, 15 Jun 2026 19:41:57 +0800 Subject: [PATCH] fix: serialize schema init for postgres --- apps/api/app/db.py | 3 +- apps/worker/worker_app/db.py | 3 +- packages/adapters/sqlalchemy_impl/__init__.py | 3 +- packages/adapters/sqlalchemy_impl/session.py | 38 +++++++++++++++++-- 4 files changed, 41 insertions(+), 6 deletions(-) diff --git a/apps/api/app/db.py b/apps/api/app/db.py index 93980d99f..0bd772013 100644 --- a/apps/api/app/db.py +++ b/apps/api/app/db.py @@ -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, diff --git a/apps/worker/worker_app/db.py b/apps/worker/worker_app/db.py index ca58fe311..78aaee433 100644 --- a/apps/worker/worker_app/db.py +++ b/apps/worker/worker_app/db.py @@ -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, diff --git a/packages/adapters/sqlalchemy_impl/__init__.py b/packages/adapters/sqlalchemy_impl/__init__.py index 53a79cc54..fa3c87125 100644 --- a/packages/adapters/sqlalchemy_impl/__init__.py +++ b/packages/adapters/sqlalchemy_impl/__init__.py @@ -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", ] diff --git a/packages/adapters/sqlalchemy_impl/session.py b/packages/adapters/sqlalchemy_impl/session.py index 3a8614b8c..f7707a92f 100644 --- a/packages/adapters/sqlalchemy_impl/session.py +++ b/packages/adapters/sqlalchemy_impl/session.py @@ -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()