(db_type: str)
| 29 | |
| 30 | |
| 31 | async def create_database_if_not_exists(db_type: str): |
| 32 | if db_type == "mysql" or db_type == "db": |
| 33 | # Connect to the server without a database |
| 34 | server_url = f"mysql+asyncmy://{mysql_db_config['user']}:{mysql_db_config['password']}@{mysql_db_config['host']}:{mysql_db_config['port']}" |
| 35 | engine = create_async_engine(server_url, echo=False) |
| 36 | async with engine.connect() as conn: |
| 37 | await conn.execute(text(f"CREATE DATABASE IF NOT EXISTS {mysql_db_config['db_name']}")) |
| 38 | await engine.dispose() |
| 39 | elif db_type == "postgres": |
| 40 | # Connect to the default 'postgres' database |
| 41 | server_url = f"postgresql+asyncpg://{postgres_db_config['user']}:{postgres_db_config['password']}@{postgres_db_config['host']}:{postgres_db_config['port']}/postgres" |
| 42 | print(f"[init_db] Connecting to Postgres: host={postgres_db_config['host']}, port={postgres_db_config['port']}, user={postgres_db_config['user']}, dbname=postgres") |
| 43 | # Isolation level AUTOCOMMIT is required for CREATE DATABASE |
| 44 | engine = create_async_engine(server_url, echo=False, isolation_level="AUTOCOMMIT") |
| 45 | async with engine.connect() as conn: |
| 46 | # Check if database exists |
| 47 | result = await conn.execute(text(f"SELECT 1 FROM pg_database WHERE datname = '{postgres_db_config['db_name']}'")) |
| 48 | if not result.scalar(): |
| 49 | await conn.execute(text(f"CREATE DATABASE {postgres_db_config['db_name']}")) |
| 50 | await engine.dispose() |
| 51 | |
| 52 | |
| 53 | def get_async_engine(db_type: str = None): |
no test coverage detected