diff --git a/app/core/db/__init__.py b/app/core/db/__init__.py index c51e677..9924fdd 100644 --- a/app/core/db/__init__.py +++ b/app/core/db/__init__.py @@ -160,13 +160,23 @@ def get_worker_session_factory() -> async_sessionmaker[AsyncSession]: def get_migration_engine() -> AsyncEngine: """Get or create the migration engine (crm_migration role). - Used by Alembic for DDL operations. This engine connects as the table owner. - Falls back to the main engine if MIGRATION_DATABASE_URL is not set. + Used by Alembic and plugin migrations for DDL operations. + This engine connects as the table owner with BYPASSRLS. + + Raises: + RuntimeError: If MIGRATION_DATABASE_URL is not set. """ global _migration_engine if _migration_engine is None: settings = get_settings() - url = settings.migration_database_url or settings.database_url + url = settings.migration_database_url + if not url: + raise RuntimeError( + "MIGRATION_DATABASE_URL is not set. " + "Plugin migrations and Alembic require a dedicated migration " + "database connection (crm_migration role). " + "The application cannot start without it." + ) _migration_engine = create_async_engine( url, pool_size=2, diff --git a/app/main.py b/app/main.py index 27069af..3e64345 100644 --- a/app/main.py +++ b/app/main.py @@ -154,7 +154,8 @@ async def lifespan(app: FastAPI): # Initialize plugin registry and discover built-in plugins registry = get_registry() - registry.initialize(get_engine(), app) + from app.core.db import get_migration_engine + registry.initialize(get_migration_engine(), app) registry.discover_builtins() # Install discovered builtin plugins and activate only those marked active in DB @@ -201,12 +202,16 @@ async def lifespan(app: FastAPI): await db.flush() logger.info(f"Created plugin record: {name} (core={plugin.manifest.is_core})") - # Run migrations if not yet applied + # Run migrations if not yet applied — use MIGRATION engine (crm_migration) for DDL if plugin.manifest.migrations: try: - await registry.migration_runner.run_all_migrations( - db, name, plugin.manifest.migrations - ) + from app.core.db import get_migration_session_factory + mig_session_factory = get_migration_session_factory() + async with mig_session_factory() as mig_db: + await registry.migration_runner.run_all_migrations( + mig_db, name, plugin.manifest.migrations + ) + await mig_db.commit() except Exception as exc: logger.error(f"Migration FAILED for {name}: {exc}") if plugin_record.active: diff --git a/app/plugins/registry.py b/app/plugins/registry.py index 77f68b7..e9c524c 100644 --- a/app/plugins/registry.py +++ b/app/plugins/registry.py @@ -494,9 +494,13 @@ class PluginRegistry: f"Running migrations to update." ) - # Re-run migrations to apply any new migration files + # Re-run migrations to apply any new migration files — use migration engine (crm_migration) for DDL if plugin.manifest.migrations: - await self.migration_runner.run_all_migrations(db, name, plugin.manifest.migrations) + from app.core.db import get_migration_session_factory + mig_factory = get_migration_session_factory() + async with mig_factory() as mig_db: + await self.migration_runner.run_all_migrations(mig_db, name, plugin.manifest.migrations) + await mig_db.commit() # Update DB version to match manifest record.version = manifest_version @@ -548,9 +552,13 @@ class PluginRegistry: # Check dependencies are installed await self._check_dependencies_installed(db, name) - # Run migrations + # Run migrations — use migration engine (crm_migration) for DDL if plugin.manifest.migrations: - await self.migration_runner.run_all_migrations(db, name, plugin.manifest.migrations) + from app.core.db import get_migration_session_factory + mig_factory = get_migration_session_factory() + async with mig_factory() as mig_db: + await self.migration_runner.run_all_migrations(mig_db, name, plugin.manifest.migrations) + await mig_db.commit() # Call on_install hook await plugin.on_install(db, self._container) @@ -741,10 +749,14 @@ class PluginRegistry: # Call on_uninstall hook await plugin.on_uninstall(db, self._container) - # Optionally drop plugin tables + # Optionally drop plugin tables — use migration engine (crm_migration) for DDL dropped_tables: list[str] = [] if remove_data: - dropped_tables = await self.migration_runner.drop_plugin_tables(db, name) + from app.core.db import get_migration_session_factory + mig_factory = get_migration_session_factory() + async with mig_factory() as mig_db: + dropped_tables = await self.migration_runner.drop_plugin_tables(mig_db, name) + await mig_db.commit() # Remove DB record await db.delete(record)