From 4a5c905934ea5792100ab20a8bd704c387d5b3b2 Mon Sep 17 00:00:00 2001 From: Agent Zero Date: Fri, 31 Jul 2026 20:45:16 +0200 Subject: [PATCH] P0-fix: plugin migrations use migration engine (crm_migration) instead of API engine (crm_api) - main.py: registry.initialize(get_migration_engine()) instead of get_engine() - main.py: plugin migrations run via get_migration_session_factory() not async_session() - registry.py: upgrade_plugin, install_plugin, uninstall_plugin all use migration session for DDL - db/__init__.py: get_migration_engine() raises RuntimeError if MIGRATION_DATABASE_URL missing (no fallback) - Fixes fresh-install failure: crm_api has no DDL rights, plugin migrations need crm_migration --- app/core/db/__init__.py | 16 +++++++++++++--- app/main.py | 15 ++++++++++----- app/plugins/registry.py | 24 ++++++++++++++++++------ 3 files changed, 41 insertions(+), 14 deletions(-) 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)