fix(database): object-level reset and optional DB_SCHEMA isolation
Safe one-shot reset, fail-fast on unknown revisions, object-level drops for non-owner corporate PG, rollback before advisory unlock, configurable DB_SCHEMA via search_path, e2e migration matrix tests.
This commit is contained in:
10
.env.example
10
.env.example
@@ -16,8 +16,14 @@ POSTGRES_DB=ss_tools
|
||||
POSTGRES_USER=postgres
|
||||
POSTGRES_PASSWORD=postgres
|
||||
|
||||
# Universal one-shot migration reset. Legacy Alembic revisions are replaced
|
||||
# by 0001_baseline; repeated starts preserve the baseline schema.
|
||||
# Схема, в которой живут все объекты приложения (по умолчанию public).
|
||||
# Для внешней корпоративной БД: DB_SCHEMA=ss_tools (схему один раз создаёт DBA).
|
||||
# DB_SCHEMA=public
|
||||
|
||||
# One-shot migration reset, only for databases NOT on the current chain (legacy/orphaned
|
||||
# Alembic revisions). true never wipes a chain-managed schema; false/unset fails fast on
|
||||
# an unknown revision instead of dropping anything. A dropping reset needs ownership of
|
||||
# the schema objects, not the schema itself.
|
||||
RESET_DATABASE_SCHEMA=false
|
||||
|
||||
# ── Application storage ─────────────────────────────────────────────────
|
||||
|
||||
@@ -21,6 +21,12 @@ ENCRYPTION_KEY=change-me-generate-a-fernet-key
|
||||
DATABASE_URL=postgresql+psycopg2://postgres:postgres@localhost:5432/ss_tools
|
||||
STORAGE_ROOT_PATH=/home/busya/dev/ss-tools-storage
|
||||
|
||||
# Схема, в которой живут все объекты приложения (по умолчанию public).
|
||||
# Для внешней корпоративной БД задайте DB_SCHEMA=ss_tools — тогда сброс и миграции
|
||||
# не требуют прав на public. Схему один раз создаёт DBA:
|
||||
# CREATE SCHEMA ss_tools AUTHORIZATION <app_user>;
|
||||
# DB_SCHEMA=public
|
||||
|
||||
# ── Admin bootstrap ────────────────────────────────────────────────────
|
||||
# INITIAL_ADMIN_CREATE=true
|
||||
# INITIAL_ADMIN_USERNAME=admin
|
||||
|
||||
@@ -82,9 +82,21 @@ from src.models import ( # noqa: F401, E402
|
||||
verification_run,
|
||||
)
|
||||
from src.models.mapping import Base # noqa: E402
|
||||
from src.core.env_settings import ( # noqa: E402
|
||||
database_connect_args,
|
||||
database_schema,
|
||||
)
|
||||
|
||||
target_metadata = Base.metadata
|
||||
|
||||
|
||||
# #region Alembic.Env.VersionTableSchema [C:1] [TYPE Function]
|
||||
# @BRIEF Explicit version-table schema for non-public deployments; None keeps the default schema.
|
||||
def _version_table_schema() -> str | None:
|
||||
schema = database_schema()
|
||||
return None if schema == "public" else schema
|
||||
# #endregion Alembic.Env.VersionTableSchema
|
||||
|
||||
# other values from the config, defined by the needs of env.py,
|
||||
# can be acquired:
|
||||
# my_important_option = config.get_main_option("my_important_option")
|
||||
@@ -109,6 +121,7 @@ def run_migrations_offline() -> None:
|
||||
target_metadata=target_metadata,
|
||||
literal_binds=True,
|
||||
dialect_opts={"paramstyle": "named"},
|
||||
version_table_schema=_version_table_schema(),
|
||||
)
|
||||
|
||||
with context.begin_transaction():
|
||||
@@ -124,7 +137,11 @@ def run_migrations_online() -> None:
|
||||
"""
|
||||
external_connection = config.attributes.get("connection")
|
||||
if external_connection is not None:
|
||||
context.configure(connection=external_connection, target_metadata=target_metadata)
|
||||
context.configure(
|
||||
connection=external_connection,
|
||||
target_metadata=target_metadata,
|
||||
version_table_schema=_version_table_schema(),
|
||||
)
|
||||
with context.begin_transaction():
|
||||
context.run_migrations()
|
||||
return
|
||||
@@ -133,10 +150,15 @@ def run_migrations_online() -> None:
|
||||
config.get_section(config.config_ini_section, {}),
|
||||
prefix="sqlalchemy.",
|
||||
poolclass=pool.NullPool,
|
||||
connect_args=database_connect_args(),
|
||||
)
|
||||
|
||||
with connectable.connect() as connection:
|
||||
context.configure(connection=connection, target_metadata=target_metadata)
|
||||
context.configure(
|
||||
connection=connection,
|
||||
target_metadata=target_metadata,
|
||||
version_table_schema=_version_table_schema(),
|
||||
)
|
||||
|
||||
with context.begin_transaction():
|
||||
context.run_migrations()
|
||||
|
||||
@@ -43,7 +43,7 @@ from src.models import ( # noqa: F401
|
||||
verification_run as _verification_run,
|
||||
)
|
||||
from src.models.mapping import Base # noqa: F401
|
||||
from .env_settings import database_url
|
||||
from .env_settings import database_connect_args, database_url
|
||||
|
||||
BASE_DIR = Path(__file__).resolve().parent.parent.parent
|
||||
DATABASE_URL = database_url()
|
||||
@@ -52,6 +52,9 @@ DATABASE_URL = database_url()
|
||||
def _build_engine(db_url: str):
|
||||
if db_url.startswith("sqlite"):
|
||||
return create_engine(db_url, connect_args={"check_same_thread": False})
|
||||
connect_args = database_connect_args()
|
||||
if connect_args:
|
||||
return create_engine(db_url, pool_pre_ping=True, connect_args=connect_args)
|
||||
return create_engine(db_url, pool_pre_ping=True)
|
||||
|
||||
|
||||
|
||||
@@ -5,6 +5,9 @@
|
||||
|
||||
import os
|
||||
from pathlib import Path
|
||||
import re
|
||||
|
||||
_SCHEMA_NAME_PATTERN = re.compile(r"^[A-Za-z_][A-Za-z0-9_]*$")
|
||||
|
||||
|
||||
# #region Core.EnvSettings.DatabaseUrl [C:2] [TYPE Function]
|
||||
@@ -18,6 +21,28 @@ def database_url() -> str:
|
||||
# #endregion Core.EnvSettings.DatabaseUrl
|
||||
|
||||
|
||||
# #region Core.EnvSettings.DatabaseSchema [C:2] [TYPE Function]
|
||||
# @ingroup Core
|
||||
# @BRIEF Return the PostgreSQL schema owning all application objects (default: public).
|
||||
def database_schema() -> str:
|
||||
value = os.getenv("DB_SCHEMA", "").strip() or "public"
|
||||
if not _SCHEMA_NAME_PATTERN.match(value):
|
||||
raise RuntimeError(f"DB_SCHEMA {value!r} is not a valid PostgreSQL schema name")
|
||||
return value
|
||||
# #endregion Core.EnvSettings.DatabaseSchema
|
||||
|
||||
|
||||
# #region Core.EnvSettings.DatabaseConnectArgs [C:2] [TYPE Function]
|
||||
# @ingroup Core
|
||||
# @BRIEF Pin the connection search_path to the application schema (no-op for public).
|
||||
def database_connect_args() -> dict:
|
||||
schema = database_schema()
|
||||
if schema == "public":
|
||||
return {}
|
||||
return {"options": f"-csearch_path={schema}"}
|
||||
# #endregion Core.EnvSettings.DatabaseConnectArgs
|
||||
|
||||
|
||||
# #region Core.EnvSettings.StorageRootPath [C:2] [TYPE Function]
|
||||
# @ingroup Core
|
||||
# @BRIEF Return the canonical filesystem storage root.
|
||||
|
||||
@@ -16,12 +16,21 @@ from alembic.script import ScriptDirectory
|
||||
from alembic.script.revision import ResolutionError
|
||||
from alembic.util.exc import CommandError
|
||||
from sqlalchemy import create_engine, inspect, text
|
||||
from sqlalchemy.exc import ProgrammingError
|
||||
|
||||
from alembic import command
|
||||
from src.core.env_settings import database_url
|
||||
from src.core.env_settings import database_connect_args, database_schema, database_url
|
||||
|
||||
_MIGRATION_LOCK_ID = 481920260824
|
||||
_BASELINE_REVISION = "0001_baseline"
|
||||
_RELATION_DDL_BY_KIND = {
|
||||
"r": "TABLE",
|
||||
"p": "TABLE",
|
||||
"f": "FOREIGN TABLE",
|
||||
"m": "MATERIALIZED VIEW",
|
||||
"v": "VIEW",
|
||||
"S": "SEQUENCE",
|
||||
}
|
||||
_RELATION_DROP_ORDER = {"r": 0, "p": 0, "f": 0, "m": 1, "v": 2, "S": 3}
|
||||
|
||||
|
||||
# #region Scripts.PrepareDatabase.Enabled [C:2] [TYPE Function]
|
||||
@@ -43,10 +52,71 @@ def _revision_is_known(config: Config, revision: str) -> bool:
|
||||
# #endregion Scripts.PrepareDatabase.RevisionIsKnown
|
||||
|
||||
|
||||
# #region Scripts.PrepareDatabase.ResetSchema [C:4] [TYPE Function]
|
||||
# @ingroup Scripts
|
||||
# @BRIEF Drop relations and enum/domain types inside the configured schema without requiring schema ownership.
|
||||
# @SIDE_EFFECT Issues transactional DROP ... CASCADE statements; foreign-owned objects fail atomically.
|
||||
def _reset_schema(connection, schema: str) -> None:
|
||||
preparer = connection.dialect.identifier_preparer
|
||||
relations = connection.execute(
|
||||
text(
|
||||
"SELECT c.relname, c.relkind FROM pg_class c "
|
||||
"JOIN pg_namespace n ON n.oid = c.relnamespace "
|
||||
"WHERE n.nspname = :schema AND c.relispartition = false "
|
||||
"AND c.relkind IN ('r', 'p', 'f', 'm', 'v', 'S')"
|
||||
),
|
||||
{"schema": schema},
|
||||
).fetchall()
|
||||
types_ = connection.execute(
|
||||
text(
|
||||
"SELECT t.typname, t.typtype FROM pg_type t "
|
||||
"JOIN pg_namespace n ON n.oid = t.typnamespace "
|
||||
"WHERE n.nspname = :schema AND t.typtype IN ('e', 'd')"
|
||||
),
|
||||
{"schema": schema},
|
||||
).fetchall()
|
||||
|
||||
for name, kind in sorted(relations, key=lambda row: _RELATION_DROP_ORDER[row[1]]):
|
||||
qualified = f"{preparer.quote_schema(schema)}.{preparer.quote(name)}"
|
||||
connection.execute(
|
||||
text(f"DROP {_RELATION_DDL_BY_KIND[kind]} IF EXISTS {qualified} CASCADE")
|
||||
)
|
||||
for name, typtype in types_:
|
||||
ddl_kind = "DOMAIN" if typtype == "d" else "TYPE"
|
||||
qualified = f"{preparer.quote_schema(schema)}.{preparer.quote(name)}"
|
||||
connection.execute(text(f"DROP {ddl_kind} IF EXISTS {qualified} CASCADE"))
|
||||
# #endregion Scripts.PrepareDatabase.ResetSchema
|
||||
|
||||
|
||||
# #region Scripts.PrepareDatabase.EnsureSchema [C:3] [TYPE Function]
|
||||
# @ingroup Scripts
|
||||
# @BRIEF Create the configured schema when missing; public is assumed present.
|
||||
# @SIDE_EFFECT Creates the schema when the connecting role may do so.
|
||||
def _ensure_schema(connection, schema: str) -> None:
|
||||
if schema == "public":
|
||||
return
|
||||
preparer = connection.dialect.identifier_preparer
|
||||
exists = connection.execute(
|
||||
text("SELECT 1 FROM pg_namespace WHERE nspname = :schema"), {"schema": schema}
|
||||
).scalar()
|
||||
if exists:
|
||||
return
|
||||
try:
|
||||
connection.execute(text(f"CREATE SCHEMA {preparer.quote(schema)}"))
|
||||
connection.commit()
|
||||
except ProgrammingError as exc:
|
||||
raise RuntimeError(
|
||||
f"Database schema {schema!r} does not exist and could not be created: {exc.orig}. "
|
||||
f"Ask the DBA to run: CREATE SCHEMA {preparer.quote(schema)} AUTHORIZATION <app_user>;"
|
||||
) from exc
|
||||
sys.stderr.write(f"[database] Created schema {schema!r}\n")
|
||||
# #endregion Scripts.PrepareDatabase.EnsureSchema
|
||||
|
||||
|
||||
# #region Scripts.PrepareDatabase.MaybeReset [C:4] [TYPE Function]
|
||||
# @ingroup Scripts
|
||||
# @BRIEF One-shot destructive schema reset gated by RESET_DATABASE_SCHEMA and orphaned-revision detection.
|
||||
# @SIDE_EFFECT Drops and recreates the public schema when reset conditions are met.
|
||||
# @BRIEF One-shot destructive reset reserved for databases not on the current chain; never wipes chain-managed schemas.
|
||||
# @SIDE_EFFECT Drops objects inside the configured schema (relations and enum/domain types) for orphaned revisions on explicit consent.
|
||||
def _maybe_reset(connection, config: Config) -> None:
|
||||
reset_enabled = _enabled(os.getenv("RESET_DATABASE_SCHEMA"))
|
||||
sys.stderr.write(
|
||||
@@ -57,19 +127,33 @@ def _maybe_reset(connection, config: Config) -> None:
|
||||
raise RuntimeError("RESET_DATABASE_SCHEMA is supported only for PostgreSQL")
|
||||
return
|
||||
|
||||
tables = set(inspect(connection).get_table_names())
|
||||
schema = database_schema()
|
||||
_ensure_schema(connection, schema)
|
||||
preparer = connection.dialect.identifier_preparer
|
||||
|
||||
tables = set(inspect(connection).get_table_names(schema=schema))
|
||||
revision = None
|
||||
if "alembic_version" in tables:
|
||||
revision = connection.execute(
|
||||
text("SELECT version_num FROM alembic_version LIMIT 1")
|
||||
text(f"SELECT version_num FROM {preparer.quote(schema)}.alembic_version LIMIT 1")
|
||||
).scalar()
|
||||
|
||||
if revision == _BASELINE_REVISION:
|
||||
sys.stderr.write("[database] Baseline already applied; one-shot reset skipped\n")
|
||||
orphaned_revision = bool(revision) and not _revision_is_known(config, revision)
|
||||
if revision is not None and not orphaned_revision:
|
||||
sys.stderr.write(
|
||||
f"[database] Alembic revision {revision} is on the current chain; "
|
||||
"one-shot reset skipped\n"
|
||||
)
|
||||
connection.commit()
|
||||
return
|
||||
orphaned_revision = bool(revision) and not _revision_is_known(config, revision)
|
||||
if not reset_enabled and not orphaned_revision:
|
||||
if not reset_enabled:
|
||||
if orphaned_revision:
|
||||
raise RuntimeError(
|
||||
f"Database Alembic revision {revision!r} is unknown to this image and "
|
||||
"RESET_DATABASE_SCHEMA is disabled; refusing destructive reset. Either "
|
||||
"(a) set RESET_DATABASE_SCHEMA=true to wipe and re-baseline "
|
||||
"(DESTROYS ALL DATA), or (b) migrate/stamp the database manually."
|
||||
)
|
||||
return
|
||||
if tables and revision is None:
|
||||
raise RuntimeError(
|
||||
@@ -77,12 +161,12 @@ def _maybe_reset(connection, config: Config) -> None:
|
||||
"refusing destructive reset"
|
||||
)
|
||||
|
||||
connection.execute(text("DROP SCHEMA IF EXISTS public CASCADE"))
|
||||
connection.execute(text("CREATE SCHEMA public"))
|
||||
if tables:
|
||||
_reset_schema(connection, schema)
|
||||
connection.commit()
|
||||
sys.stderr.write(
|
||||
f"[database] {'Orphaned' if orphaned_revision else 'Legacy'} Alembic revision "
|
||||
f"{revision or '<empty>'} reset successfully\n"
|
||||
f"{revision or '<empty>'} reset successfully in schema {schema!r}\n"
|
||||
)
|
||||
# #endregion Scripts.PrepareDatabase.MaybeReset
|
||||
|
||||
@@ -91,7 +175,9 @@ def _maybe_reset(connection, config: Config) -> None:
|
||||
# @ingroup Scripts
|
||||
# @BRIEF Run alembic upgrade head under a PostgreSQL advisory lock.
|
||||
def prepare_database() -> None:
|
||||
engine = create_engine(database_url(), pool_pre_ping=True)
|
||||
engine = create_engine(
|
||||
database_url(), pool_pre_ping=True, connect_args=database_connect_args()
|
||||
)
|
||||
try:
|
||||
with engine.connect() as connection:
|
||||
config = Config(str(Path(__file__).resolve().parents[2] / "alembic.ini"))
|
||||
@@ -106,6 +192,7 @@ def prepare_database() -> None:
|
||||
command.upgrade(config, "head")
|
||||
finally:
|
||||
if locked:
|
||||
connection.rollback()
|
||||
connection.execute(
|
||||
text("SELECT pg_advisory_unlock(:lock_id)"),
|
||||
{"lock_id": _MIGRATION_LOCK_ID},
|
||||
|
||||
@@ -34,6 +34,7 @@ from sqlalchemy import Column, String, Text, create_engine
|
||||
from sqlalchemy.orm import Session, declarative_base
|
||||
|
||||
from src.core.encryption import is_fernet_token
|
||||
from src.core.env_settings import database_connect_args
|
||||
|
||||
Base = declarative_base()
|
||||
|
||||
@@ -124,7 +125,7 @@ def main() -> None:
|
||||
# ── Load database URL ──────────────────────────────────────────
|
||||
db_url = os.getenv("DATABASE_URL", "")
|
||||
# Use psycopg2 for sync access in script
|
||||
engine = create_engine(db_url)
|
||||
engine = create_engine(db_url, connect_args=database_connect_args())
|
||||
|
||||
# ── Step 1: Environment passwords (AppConfigRecord.payload.environments) ──
|
||||
_r("── Environment passwords (ConfigManager) ──")
|
||||
|
||||
597
backend/tests/integration/test_alembic_prepare_database_e2e.py
Normal file
597
backend/tests/integration/test_alembic_prepare_database_e2e.py
Normal file
@@ -0,0 +1,597 @@
|
||||
# #region Test.Alembic.PrepareDatabaseE2E [C:5] [TYPE Module] [SEMANTICS test,alembic,migration,e2e,postgres,testcontainers]
|
||||
# @BRIEF E2E migration matrix for prepare_database against real PostgreSQL (Testcontainers).
|
||||
# @RELATION VERIFIES -> [Scripts.PrepareDatabase.PrepareDatabase]
|
||||
# @RELATION VERIFIES -> [Scripts.PrepareDatabase.MaybeReset]
|
||||
# @PRE Docker daemon is running and --run-integration is passed.
|
||||
# @TEST_FIXTURE: prepare_e2e_url -> INLINE_JSON
|
||||
# @TEST_EDGE: empty_schema -> A fresh database upgrades through the full chain to head.
|
||||
# @TEST_EDGE: repeated_run -> A second run on head is a no-op and preserves data.
|
||||
# @TEST_EDGE: intermediate_revision -> A known mid-chain revision upgrades forward without data loss.
|
||||
# @TEST_EDGE: chain_managed_reset -> RESET_DATABASE_SCHEMA=true must never wipe a chain-managed schema.
|
||||
# @TEST_EDGE: unversioned_disabled -> Unversioned non-empty schema without the flag fails at the baseline guard.
|
||||
# @TEST_EDGE: non_schema_owner -> A role owning objects but not schema public resets and migrates successfully.
|
||||
# @TEST_EDGE: foreign_owned -> Foreign-owned objects fail the reset atomically with a clean ownership error.
|
||||
# @TEST_EDGE: enum_types -> Stale native enum types sharing baseline names are dropped before create_all.
|
||||
# @TEST_EDGE: dedicated_schema -> DB_SCHEMA keeps app objects out of public and scopes the reset.
|
||||
# @TEST_EDGE: lock_blocked -> A held advisory lock blocks migration until released.
|
||||
# @TEST_EDGE: concurrent_runs -> Two parallel runs serialize on the advisory lock.
|
||||
# @TEST_INVARIANT: prepare_database_never_wipes_chain_schema -> VERIFIED_BY: [Test.Alembic.TestResetTrueNeverWipesChainManagedSchema]
|
||||
# @TEST_INVARIANT: prepare_database_surfaces_real_error -> VERIFIED_BY: [Test.Alembic.TestResetFailsCleanlyWhenObjectsAreForeignOwned]
|
||||
# @TEST_INVARIANT: dedicated_schema_isolation -> VERIFIED_BY: [Test.Alembic.TestDedicatedSchemaResetIsScoped]
|
||||
# @RATIONALE Every scenario runs prepare_database() (the entrypoint code path) against a pristine
|
||||
# database created by db_factory in the shared PostgreSQL container. Companion file
|
||||
# test_alembic_user_id_migration.py covers orphaned-revision refusal and legacy reset.
|
||||
# @REJECTED SQLite/metadata.create_all variants were rejected — they cannot exercise advisory
|
||||
# locks, schema ownership, or the real Alembic DDL path.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
from pathlib import Path
|
||||
import time
|
||||
import uuid
|
||||
|
||||
import psycopg2
|
||||
import pytest
|
||||
from alembic import command
|
||||
from alembic.config import Config
|
||||
from alembic.script import ScriptDirectory
|
||||
from sqlalchemy import create_engine, inspect, text
|
||||
from sqlalchemy.exc import ProgrammingError
|
||||
|
||||
pytestmark = pytest.mark.integration
|
||||
|
||||
_BACKEND_DIR = Path(__file__).parents[2]
|
||||
_ORPHANED_REVISION = "y4z5a6b7c8d9"
|
||||
|
||||
|
||||
# #region Test.Alembic.PrepareE2EUrl [C:2] [TYPE Fixture]
|
||||
# @ingroup Test.Alembic.PrepareDatabaseE2E
|
||||
# @BRIEF Pristine isolated database per test via db_factory in the shared PG container.
|
||||
# @POST Returns a host-accessible connection URL for an empty database.
|
||||
@pytest.fixture
|
||||
def prepare_e2e_url(db_factory) -> str:
|
||||
"""Provide a pristine empty PostgreSQL database for one E2E scenario."""
|
||||
return db_factory["create_db"]("_prepare_e2e")["host_url"]
|
||||
|
||||
|
||||
# #endregion Test.Alembic.PrepareE2EUrl
|
||||
|
||||
|
||||
def _alembic_config(url: str) -> Config:
|
||||
config = Config(str(_BACKEND_DIR / "alembic.ini"))
|
||||
config.set_main_option("script_location", str(_BACKEND_DIR / "alembic"))
|
||||
config.set_main_option("sqlalchemy.url", url)
|
||||
return config
|
||||
|
||||
|
||||
def _current_head() -> str:
|
||||
return ScriptDirectory.from_config(
|
||||
_alembic_config("postgresql+psycopg2://unused@invalid/invalid")
|
||||
).get_current_head()
|
||||
|
||||
|
||||
def _run_alembic_upgrade(url: str, revision: str) -> None:
|
||||
command.upgrade(_alembic_config(url), revision)
|
||||
|
||||
|
||||
def _table_names(url: str) -> set[str]:
|
||||
engine = create_engine(url)
|
||||
try:
|
||||
return set(inspect(engine).get_table_names())
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
|
||||
def _revision(url: str) -> str | None:
|
||||
engine = create_engine(url)
|
||||
try:
|
||||
with engine.connect() as connection:
|
||||
return connection.execute(text("SELECT version_num FROM alembic_version LIMIT 1")).scalar()
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
|
||||
def _table_names_in_schema(url: str, schema: str) -> set[str]:
|
||||
engine = create_engine(url)
|
||||
try:
|
||||
return set(inspect(engine).get_table_names(schema=schema))
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
|
||||
def _revision_in_schema(url: str, schema: str) -> str | None:
|
||||
engine = create_engine(url)
|
||||
try:
|
||||
with engine.connect() as connection:
|
||||
return connection.execute(
|
||||
text(f"SELECT version_num FROM {schema}.alembic_version LIMIT 1")
|
||||
).scalar()
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
|
||||
def _insert_sentinel(url: str, task_id: str) -> None:
|
||||
engine = create_engine(url)
|
||||
try:
|
||||
with engine.begin() as connection:
|
||||
connection.execute(
|
||||
text(
|
||||
"INSERT INTO task_records (id, type, status) "
|
||||
"VALUES (:id, 'e2e_migration', 'SUCCESS')"
|
||||
),
|
||||
{"id": task_id},
|
||||
)
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
|
||||
def _sentinel_exists(url: str, task_id: str) -> bool:
|
||||
engine = create_engine(url)
|
||||
try:
|
||||
with engine.connect() as connection:
|
||||
return (
|
||||
connection.execute(
|
||||
text("SELECT COUNT(*) FROM task_records WHERE id = :id"), {"id": task_id}
|
||||
).scalar_one()
|
||||
> 0
|
||||
)
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
|
||||
def _admin_connect(url: str):
|
||||
"""Connect to the scenario database as the container superuser via psycopg2."""
|
||||
from urllib.parse import urlparse
|
||||
|
||||
parsed = urlparse(url)
|
||||
return psycopg2.connect(
|
||||
host=parsed.hostname,
|
||||
port=parsed.port,
|
||||
user="test",
|
||||
password="test",
|
||||
dbname=parsed.path.lstrip("/"),
|
||||
)
|
||||
|
||||
|
||||
def _import_prepare_database():
|
||||
from src.scripts.prepare_database import prepare_database
|
||||
|
||||
return prepare_database
|
||||
|
||||
|
||||
# #region Test.Alembic.TestFreshDatabaseUpgradesToHead [C:3] [TYPE Function]
|
||||
# @ingroup Test.Alembic.PrepareDatabaseE2E
|
||||
# @BRIEF A fresh empty database upgrades through the full chain regardless of the reset flag.
|
||||
@pytest.mark.parametrize("reset_flag", ["false", "true"])
|
||||
def test_fresh_database_upgrades_to_head(
|
||||
monkeypatch: pytest.MonkeyPatch, prepare_e2e_url: str, reset_flag: str
|
||||
) -> None:
|
||||
"""The entrypoint migration path applies the whole chain on an empty database."""
|
||||
monkeypatch.setenv("DATABASE_URL", prepare_e2e_url)
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", reset_flag)
|
||||
|
||||
_import_prepare_database()()
|
||||
|
||||
assert _revision(prepare_e2e_url) == _current_head()
|
||||
assert "task_records" in _table_names(prepare_e2e_url)
|
||||
|
||||
|
||||
# #endregion Test.Alembic.TestFreshDatabaseUpgradesToHead
|
||||
|
||||
|
||||
# #region Test.Alembic.TestRepeatedRunsAreIdempotent [C:3] [TYPE Function]
|
||||
# @ingroup Test.Alembic.PrepareDatabaseE2E
|
||||
# @BRIEF A second run on an already-migrated database is a no-op that preserves data.
|
||||
def test_repeated_runs_are_idempotent_and_preserve_data(
|
||||
monkeypatch: pytest.MonkeyPatch, prepare_e2e_url: str
|
||||
) -> None:
|
||||
"""Re-running the entrypoint on head must not touch the schema or existing rows."""
|
||||
monkeypatch.setenv("DATABASE_URL", prepare_e2e_url)
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "false")
|
||||
prepare_database = _import_prepare_database()
|
||||
|
||||
prepare_database()
|
||||
sentinel = f"sentinel-{uuid.uuid4().hex[:8]}"
|
||||
_insert_sentinel(prepare_e2e_url, sentinel)
|
||||
|
||||
prepare_database()
|
||||
|
||||
assert _revision(prepare_e2e_url) == _current_head()
|
||||
assert _sentinel_exists(prepare_e2e_url, sentinel)
|
||||
|
||||
|
||||
# #endregion Test.Alembic.TestRepeatedRunsAreIdempotent
|
||||
|
||||
|
||||
# #region Test.Alembic.TestUpgradeFromIntermediateRevision [C:3] [TYPE Function]
|
||||
# @ingroup Test.Alembic.PrepareDatabaseE2E
|
||||
# @BRIEF A known mid-chain revision upgrades forward without data loss (prod upgrade path).
|
||||
def test_upgrade_from_intermediate_revision_preserves_data(
|
||||
monkeypatch: pytest.MonkeyPatch, prepare_e2e_url: str
|
||||
) -> None:
|
||||
"""Deploying a newer image over a partially migrated database keeps application data."""
|
||||
monkeypatch.setenv("DATABASE_URL", prepare_e2e_url)
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "false")
|
||||
_run_alembic_upgrade(prepare_e2e_url, "0010_authoring_contract")
|
||||
sentinel = f"sentinel-{uuid.uuid4().hex[:8]}"
|
||||
_insert_sentinel(prepare_e2e_url, sentinel)
|
||||
|
||||
_import_prepare_database()()
|
||||
|
||||
assert _revision(prepare_e2e_url) == _current_head()
|
||||
assert _sentinel_exists(prepare_e2e_url, sentinel)
|
||||
|
||||
|
||||
# #endregion Test.Alembic.TestUpgradeFromIntermediateRevision
|
||||
|
||||
|
||||
# #region Test.Alembic.TestResetTrueNeverWipesChainManagedSchema [C:3] [TYPE Function]
|
||||
# @ingroup Test.Alembic.PrepareDatabaseE2E
|
||||
# @BRIEF RESET_DATABASE_SCHEMA=true is one-shot: a chain-managed schema is never wiped.
|
||||
@pytest.mark.parametrize("revision", ["0001_baseline", "0010_authoring_contract"])
|
||||
def test_reset_true_never_wipes_chain_managed_schema(
|
||||
monkeypatch: pytest.MonkeyPatch, prepare_e2e_url: str, revision: str
|
||||
) -> None:
|
||||
"""Restarting with the bundled default RESET_DATABASE_SCHEMA=true must not drop data.
|
||||
|
||||
The reset is reserved for databases that are NOT on the current chain
|
||||
(legacy/orphaned revisions). Any known revision — baseline or mid-chain —
|
||||
must be forward-migrated by alembic upgrade, never dropped.
|
||||
"""
|
||||
monkeypatch.setenv("DATABASE_URL", prepare_e2e_url)
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "true")
|
||||
_run_alembic_upgrade(prepare_e2e_url, revision)
|
||||
sentinel = f"sentinel-{uuid.uuid4().hex[:8]}"
|
||||
_insert_sentinel(prepare_e2e_url, sentinel)
|
||||
|
||||
_import_prepare_database()()
|
||||
|
||||
assert _revision(prepare_e2e_url) == _current_head()
|
||||
assert _sentinel_exists(prepare_e2e_url, sentinel)
|
||||
|
||||
|
||||
# #endregion Test.Alembic.TestResetTrueNeverWipesChainManagedSchema
|
||||
|
||||
|
||||
# #region Test.Alembic.TestUnversionedSchemaFailsAtBaselineGuard [C:3] [TYPE Function]
|
||||
# @ingroup Test.Alembic.PrepareDatabaseE2E
|
||||
# @BRIEF An unversioned non-empty schema without the flag fails in the baseline guard.
|
||||
def test_unversioned_schema_without_flag_fails_at_baseline_guard(
|
||||
monkeypatch: pytest.MonkeyPatch, prepare_e2e_url: str
|
||||
) -> None:
|
||||
"""Disabled reset defers to the baseline guard; tables survive and the lock is released."""
|
||||
engine = create_engine(prepare_e2e_url)
|
||||
try:
|
||||
with engine.begin() as connection:
|
||||
connection.execute(text("CREATE TABLE unrelated_data (id INTEGER PRIMARY KEY)"))
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
monkeypatch.setenv("DATABASE_URL", prepare_e2e_url)
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "false")
|
||||
|
||||
from src.scripts.prepare_database import _MIGRATION_LOCK_ID
|
||||
|
||||
with pytest.raises(RuntimeError, match="Refusing to stamp"):
|
||||
_import_prepare_database()()
|
||||
|
||||
assert "unrelated_data" in _table_names(prepare_e2e_url)
|
||||
|
||||
engine = create_engine(prepare_e2e_url)
|
||||
try:
|
||||
with engine.connect() as connection:
|
||||
acquired = connection.execute(
|
||||
text("SELECT pg_try_advisory_lock(:lock_id)"), {"lock_id": _MIGRATION_LOCK_ID}
|
||||
).scalar_one()
|
||||
assert acquired is True
|
||||
connection.execute(
|
||||
text("SELECT pg_advisory_unlock(:lock_id)"), {"lock_id": _MIGRATION_LOCK_ID}
|
||||
)
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
|
||||
# #endregion Test.Alembic.TestUnversionedSchemaFailsAtBaselineGuard
|
||||
|
||||
|
||||
# #region Test.Alembic.TestNonSchemaOwnerResetDropsOwnedObjects [C:4] [TYPE Function]
|
||||
# @ingroup Test.Alembic.PrepareDatabaseE2E
|
||||
# @BRIEF A role owning schema objects but not schema public resets successfully and migrates.
|
||||
def test_non_schema_owner_reset_drops_owned_objects_and_migrates(
|
||||
monkeypatch: pytest.MonkeyPatch, prepare_e2e_url: str
|
||||
) -> None:
|
||||
"""Prod topology: the app role owns tables/types, the DBA owns schema public.
|
||||
|
||||
The object-level reset must succeed without schema ownership: drop the role-owned
|
||||
objects, then let the baseline recreate everything. Regression for the prod
|
||||
InsufficientPrivilege on DROP SCHEMA public.
|
||||
"""
|
||||
role = f"ss_e2e_{uuid.uuid4().hex[:10]}"
|
||||
password = f"pw_{uuid.uuid4().hex[:12]}"
|
||||
|
||||
admin = _admin_connect(prepare_e2e_url)
|
||||
try:
|
||||
admin.autocommit = True
|
||||
with admin.cursor() as cursor:
|
||||
cursor.execute(f'CREATE ROLE "{role}" LOGIN PASSWORD \'{password}\'')
|
||||
cursor.execute(f'GRANT CONNECT ON DATABASE "{admin.info.dbname}" TO "{role}"')
|
||||
cursor.execute(f'GRANT USAGE, CREATE ON SCHEMA public TO "{role}"')
|
||||
cursor.execute(
|
||||
"SELECT pg_get_userbyid(nspowner) FROM pg_namespace WHERE nspname = 'public'"
|
||||
)
|
||||
schema_owner = cursor.fetchone()[0]
|
||||
finally:
|
||||
admin.close()
|
||||
assert schema_owner != role, "test topology requires the role to NOT own schema public"
|
||||
|
||||
app_url = prepare_e2e_url.replace("://test:test@", f"://{role}:{password}@")
|
||||
app = create_engine(app_url)
|
||||
try:
|
||||
with app.begin() as connection:
|
||||
connection.execute(text("CREATE TABLE alembic_version (version_num VARCHAR(64) NOT NULL)"))
|
||||
connection.execute(text(f"INSERT INTO alembic_version VALUES ('{_ORPHANED_REVISION}')"))
|
||||
connection.execute(text("CREATE TABLE legacy_probe (id INTEGER PRIMARY KEY)"))
|
||||
connection.execute(text("CREATE TYPE legacy_only_enum AS ENUM ('A')"))
|
||||
connection.execute(
|
||||
text("CREATE TABLE legacy_typed (id INTEGER PRIMARY KEY, state legacy_only_enum)")
|
||||
)
|
||||
finally:
|
||||
app.dispose()
|
||||
|
||||
monkeypatch.setenv("DATABASE_URL", app_url)
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "true")
|
||||
|
||||
_import_prepare_database()()
|
||||
|
||||
assert _revision(prepare_e2e_url) == _current_head()
|
||||
tables = _table_names(prepare_e2e_url)
|
||||
assert "legacy_probe" not in tables
|
||||
assert "legacy_typed" not in tables
|
||||
assert "task_records" in tables
|
||||
|
||||
admin = _admin_connect(prepare_e2e_url)
|
||||
try:
|
||||
with admin.cursor() as cursor:
|
||||
cursor.execute(
|
||||
"SELECT COUNT(*) FROM pg_type t JOIN pg_namespace n ON n.oid = t.typnamespace "
|
||||
"WHERE n.nspname = 'public' AND t.typname = 'legacy_only_enum'"
|
||||
)
|
||||
assert cursor.fetchone()[0] == 0
|
||||
finally:
|
||||
admin.close()
|
||||
|
||||
|
||||
# #endregion Test.Alembic.TestNonSchemaOwnerResetDropsOwnedObjects
|
||||
|
||||
|
||||
# #region Test.Alembic.TestResetFailsCleanlyWhenObjectsAreForeignOwned [C:4] [TYPE Function]
|
||||
# @ingroup Test.Alembic.PrepareDatabaseE2E
|
||||
# @BRIEF Foreign-owned schema objects fail the reset atomically with a clean ownership error.
|
||||
def test_reset_fails_cleanly_when_objects_are_foreign_owned(
|
||||
monkeypatch: pytest.MonkeyPatch, prepare_e2e_url: str
|
||||
) -> None:
|
||||
"""The reset must not touch objects owned by other roles and must not lose data.
|
||||
|
||||
Nothing is dropped (transactional DDL rolls back), the ownership error surfaces
|
||||
unmasked, and the advisory lock is released.
|
||||
"""
|
||||
role = f"ss_e2e_{uuid.uuid4().hex[:10]}"
|
||||
password = f"pw_{uuid.uuid4().hex[:12]}"
|
||||
|
||||
admin = _admin_connect(prepare_e2e_url)
|
||||
try:
|
||||
admin.autocommit = True
|
||||
with admin.cursor() as cursor:
|
||||
cursor.execute(f'CREATE ROLE "{role}" LOGIN PASSWORD \'{password}\'')
|
||||
cursor.execute(f'GRANT CONNECT ON DATABASE "{admin.info.dbname}" TO "{role}"')
|
||||
cursor.execute(f'GRANT USAGE ON SCHEMA public TO "{role}"')
|
||||
cursor.execute("CREATE TABLE alembic_version (version_num VARCHAR(64) NOT NULL)")
|
||||
cursor.execute(f"INSERT INTO alembic_version VALUES ('{_ORPHANED_REVISION}')")
|
||||
cursor.execute("CREATE TABLE legacy_probe (id INTEGER PRIMARY KEY)")
|
||||
cursor.execute(f'GRANT SELECT ON public.alembic_version TO "{role}"')
|
||||
finally:
|
||||
admin.close()
|
||||
|
||||
app_url = prepare_e2e_url.replace("://test:test@", f"://{role}:{password}@")
|
||||
monkeypatch.setenv("DATABASE_URL", app_url)
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "true")
|
||||
|
||||
from src.scripts.prepare_database import _MIGRATION_LOCK_ID
|
||||
|
||||
with pytest.raises(ProgrammingError) as excinfo:
|
||||
_import_prepare_database()()
|
||||
|
||||
assert "must be owner" in str(excinfo.value)
|
||||
assert "InFailedSqlTransaction" not in str(excinfo.value)
|
||||
|
||||
tables = _table_names(prepare_e2e_url)
|
||||
assert "legacy_probe" in tables
|
||||
assert "alembic_version" in tables
|
||||
with create_engine(prepare_e2e_url).connect() as connection:
|
||||
revision = connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one()
|
||||
assert revision == _ORPHANED_REVISION
|
||||
|
||||
engine = create_engine(prepare_e2e_url)
|
||||
try:
|
||||
with engine.connect() as connection:
|
||||
acquired = connection.execute(
|
||||
text("SELECT pg_try_advisory_lock(:lock_id)"), {"lock_id": _MIGRATION_LOCK_ID}
|
||||
).scalar_one()
|
||||
assert acquired is True
|
||||
connection.execute(
|
||||
text("SELECT pg_advisory_unlock(:lock_id)"), {"lock_id": _MIGRATION_LOCK_ID}
|
||||
)
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
|
||||
# #endregion Test.Alembic.TestResetFailsCleanlyWhenObjectsAreForeignOwned
|
||||
|
||||
|
||||
# #region Test.Alembic.TestResetDropsEnumTypesThatWouldCollideWithBaseline [C:3] [TYPE Function]
|
||||
# @ingroup Test.Alembic.PrepareDatabaseE2E
|
||||
# @BRIEF A stale enum type sharing a baseline type name is dropped before create_all.
|
||||
def test_reset_drops_enum_types_that_would_collide_with_baseline(
|
||||
monkeypatch: pytest.MonkeyPatch, prepare_e2e_url: str
|
||||
) -> None:
|
||||
"""Legacy native enum types must be removed so the baseline can recreate them."""
|
||||
admin = _admin_connect(prepare_e2e_url)
|
||||
try:
|
||||
admin.autocommit = True
|
||||
with admin.cursor() as cursor:
|
||||
cursor.execute("CREATE TABLE alembic_version (version_num VARCHAR(64) NOT NULL)")
|
||||
cursor.execute(f"INSERT INTO alembic_version VALUES ('{_ORPHANED_REVISION}')")
|
||||
cursor.execute("CREATE TYPE gitprovider AS ENUM ('GITHUB')")
|
||||
cursor.execute(
|
||||
"CREATE TABLE stale_typed (id INTEGER PRIMARY KEY, provider gitprovider)"
|
||||
)
|
||||
finally:
|
||||
admin.close()
|
||||
|
||||
monkeypatch.setenv("DATABASE_URL", prepare_e2e_url)
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "true")
|
||||
|
||||
_import_prepare_database()()
|
||||
|
||||
assert _revision(prepare_e2e_url) == _current_head()
|
||||
tables = _table_names(prepare_e2e_url)
|
||||
assert "stale_typed" not in tables
|
||||
assert "git_repositories" in tables
|
||||
|
||||
|
||||
# #endregion Test.Alembic.TestResetDropsEnumTypesThatWouldCollideWithBaseline
|
||||
|
||||
|
||||
# #region Test.Alembic.TestDedicatedSchemaMigrates [C:3] [TYPE Function]
|
||||
# @ingroup Test.Alembic.PrepareDatabaseE2E
|
||||
# @BRIEF DB_SCHEMA places every object in a dedicated schema and leaves public empty.
|
||||
def test_dedicated_schema_migrates_without_touching_public(
|
||||
monkeypatch: pytest.MonkeyPatch, prepare_e2e_url: str
|
||||
) -> None:
|
||||
"""A non-public DB_SCHEMA keeps all app objects out of the public schema."""
|
||||
schema = "ss_tools"
|
||||
engine = create_engine(prepare_e2e_url)
|
||||
try:
|
||||
with engine.begin() as connection:
|
||||
connection.execute(text(f'CREATE SCHEMA "{schema}"'))
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
monkeypatch.setenv("DATABASE_URL", prepare_e2e_url)
|
||||
monkeypatch.setenv("DB_SCHEMA", schema)
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "true")
|
||||
|
||||
_import_prepare_database()()
|
||||
|
||||
assert _revision_in_schema(prepare_e2e_url, schema) == _current_head()
|
||||
assert "task_records" in _table_names_in_schema(prepare_e2e_url, schema)
|
||||
public_tables = _table_names_in_schema(prepare_e2e_url, "public")
|
||||
assert "task_records" not in public_tables
|
||||
assert "alembic_version" not in public_tables
|
||||
|
||||
|
||||
# #endregion Test.Alembic.TestDedicatedSchemaMigrates
|
||||
|
||||
|
||||
# #region Test.Alembic.TestDedicatedSchemaResetIsScoped [C:3] [TYPE Function]
|
||||
# @ingroup Test.Alembic.PrepareDatabaseE2E
|
||||
# @BRIEF The reset drops objects only inside DB_SCHEMA; public objects survive.
|
||||
def test_dedicated_schema_reset_is_scoped_to_that_schema(
|
||||
monkeypatch: pytest.MonkeyPatch, prepare_e2e_url: str
|
||||
) -> None:
|
||||
"""Legacy/orphaned objects in the dedicated schema are reset; public is untouched."""
|
||||
schema = "ss_tools"
|
||||
engine = create_engine(prepare_e2e_url)
|
||||
try:
|
||||
with engine.begin() as connection:
|
||||
connection.execute(text(f'CREATE SCHEMA "{schema}"'))
|
||||
connection.execute(
|
||||
text(f"CREATE TABLE {schema}.alembic_version (version_num VARCHAR(64) NOT NULL)")
|
||||
)
|
||||
connection.execute(
|
||||
text(f"INSERT INTO {schema}.alembic_version VALUES ('{_ORPHANED_REVISION}')")
|
||||
)
|
||||
connection.execute(text(f"CREATE TABLE {schema}.legacy_probe (id INTEGER PRIMARY KEY)"))
|
||||
connection.execute(text("CREATE TABLE public.keep_me (id INTEGER PRIMARY KEY)"))
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
monkeypatch.setenv("DATABASE_URL", prepare_e2e_url)
|
||||
monkeypatch.setenv("DB_SCHEMA", schema)
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "true")
|
||||
|
||||
_import_prepare_database()()
|
||||
|
||||
assert _revision_in_schema(prepare_e2e_url, schema) == _current_head()
|
||||
schema_tables = _table_names_in_schema(prepare_e2e_url, schema)
|
||||
assert "legacy_probe" not in schema_tables
|
||||
assert "task_records" in schema_tables
|
||||
assert "keep_me" in _table_names_in_schema(prepare_e2e_url, "public")
|
||||
|
||||
|
||||
# #endregion Test.Alembic.TestDedicatedSchemaResetIsScoped
|
||||
|
||||
|
||||
# #region Test.Alembic.TestMigrationLockBlocksUntilReleased [C:4] [TYPE Function]
|
||||
# @ingroup Test.Alembic.PrepareDatabaseE2E
|
||||
# @BRIEF A held migration advisory lock blocks prepare_database until released.
|
||||
def test_migration_lock_blocks_until_released(
|
||||
monkeypatch: pytest.MonkeyPatch, prepare_e2e_url: str
|
||||
) -> None:
|
||||
"""Another holder of the advisory lock postpones the migration run."""
|
||||
from src.scripts.prepare_database import _MIGRATION_LOCK_ID
|
||||
|
||||
monkeypatch.setenv("DATABASE_URL", prepare_e2e_url)
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "false")
|
||||
prepare_database = _import_prepare_database()
|
||||
|
||||
holder = create_engine(prepare_e2e_url)
|
||||
try:
|
||||
with holder.connect() as connection:
|
||||
connection.execute(
|
||||
text("SELECT pg_advisory_lock(:lock_id)"), {"lock_id": _MIGRATION_LOCK_ID}
|
||||
)
|
||||
connection.commit()
|
||||
|
||||
with ThreadPoolExecutor(max_workers=1) as pool:
|
||||
future = pool.submit(prepare_database)
|
||||
time.sleep(2.0)
|
||||
assert future.done() is False, "migration must block on the advisory lock"
|
||||
|
||||
connection.execute(
|
||||
text("SELECT pg_advisory_unlock(:lock_id)"), {"lock_id": _MIGRATION_LOCK_ID}
|
||||
)
|
||||
connection.commit()
|
||||
future.result(timeout=180)
|
||||
finally:
|
||||
holder.dispose()
|
||||
|
||||
assert _revision(prepare_e2e_url) == _current_head()
|
||||
|
||||
|
||||
# #endregion Test.Alembic.TestMigrationLockBlocksUntilReleased
|
||||
|
||||
|
||||
# #region Test.Alembic.TestConcurrentRunsSerialize [C:4] [TYPE Function]
|
||||
# @ingroup Test.Alembic.PrepareDatabaseE2E
|
||||
# @BRIEF Two concurrent prepare_database runs serialize on the advisory lock and both finish.
|
||||
def test_concurrent_runs_serialize_on_advisory_lock(
|
||||
monkeypatch: pytest.MonkeyPatch, prepare_e2e_url: str
|
||||
) -> None:
|
||||
"""Parallel entrypoint runs (e.g. backend replicas) must not corrupt the migration."""
|
||||
monkeypatch.setenv("DATABASE_URL", prepare_e2e_url)
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "false")
|
||||
prepare_database = _import_prepare_database()
|
||||
|
||||
with ThreadPoolExecutor(max_workers=2) as pool:
|
||||
futures = [pool.submit(prepare_database) for _ in range(2)]
|
||||
for future in futures:
|
||||
future.result(timeout=180)
|
||||
|
||||
assert _revision(prepare_e2e_url) == _current_head()
|
||||
assert "task_records" in _table_names(prepare_e2e_url)
|
||||
|
||||
|
||||
# #endregion Test.Alembic.TestConcurrentRunsSerialize
|
||||
# #endregion Test.Alembic.PrepareDatabaseE2E
|
||||
@@ -63,7 +63,7 @@ def test_prepare_database_resets_any_legacy_revision(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
migration_database_url: str,
|
||||
) -> None:
|
||||
"""An explicit reset replaces any old Alembic chain with the current baseline."""
|
||||
"""An explicit RESET_DATABASE_SCHEMA=true replaces any old Alembic chain with the baseline."""
|
||||
engine = create_engine(migration_database_url)
|
||||
try:
|
||||
with engine.begin() as connection:
|
||||
@@ -74,8 +74,8 @@ def test_prepare_database_resets_any_legacy_revision(
|
||||
connection.execute(text("CREATE TABLE legacy_probe (id INTEGER PRIMARY KEY)"))
|
||||
|
||||
monkeypatch.setenv("DATABASE_URL", migration_database_url)
|
||||
# Orphaned revisions are reset independently of stale Compose defaults.
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "false")
|
||||
# Destructive reset requires explicit consent via the flag.
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "true")
|
||||
monkeypatch.delenv("RESET_DATABASE_FROM_REVISION", raising=False)
|
||||
monkeypatch.delenv("RESET_DATABASE_NAME", raising=False)
|
||||
|
||||
@@ -88,7 +88,73 @@ def test_prepare_database_resets_any_legacy_revision(
|
||||
assert "task_records" in tables
|
||||
with engine.connect() as connection:
|
||||
revision = connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one()
|
||||
assert revision == "0001_baseline"
|
||||
assert revision == _current_head()
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
|
||||
def test_prepare_database_refuses_orphaned_revision_when_disabled(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
migration_database_url: str,
|
||||
) -> None:
|
||||
"""A disabled RESET_DATABASE_SCHEMA is authoritative: unknown revisions fail fast, no data loss."""
|
||||
engine = create_engine(migration_database_url)
|
||||
try:
|
||||
with engine.begin() as connection:
|
||||
connection.execute(text("CREATE TABLE alembic_version (version_num VARCHAR(64) NOT NULL)"))
|
||||
connection.execute(
|
||||
text("INSERT INTO alembic_version (version_num) VALUES ('y4z5a6b7c8d9')")
|
||||
)
|
||||
connection.execute(text("CREATE TABLE legacy_probe (id INTEGER PRIMARY KEY)"))
|
||||
|
||||
monkeypatch.setenv("DATABASE_URL", migration_database_url)
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "false")
|
||||
|
||||
from src.scripts.prepare_database import prepare_database
|
||||
|
||||
with pytest.raises(RuntimeError, match="unknown to this image"):
|
||||
prepare_database()
|
||||
|
||||
tables = set(inspect(engine).get_table_names())
|
||||
assert "legacy_probe" in tables
|
||||
with engine.connect() as connection:
|
||||
revision = connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one()
|
||||
assert revision == "y4z5a6b7c8d9"
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
|
||||
def test_prepare_database_unlock_survives_failure(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
migration_database_url: str,
|
||||
) -> None:
|
||||
"""A failed prepare run must surface the real error and release the migration advisory lock."""
|
||||
engine = create_engine(migration_database_url)
|
||||
try:
|
||||
with engine.begin() as connection:
|
||||
connection.execute(text("CREATE TABLE alembic_version (version_num VARCHAR(64) NOT NULL)"))
|
||||
connection.execute(
|
||||
text("INSERT INTO alembic_version (version_num) VALUES ('y4z5a6b7c8d9')")
|
||||
)
|
||||
|
||||
monkeypatch.setenv("DATABASE_URL", migration_database_url)
|
||||
monkeypatch.setenv("RESET_DATABASE_SCHEMA", "false")
|
||||
|
||||
from src.scripts.prepare_database import _MIGRATION_LOCK_ID, prepare_database
|
||||
|
||||
with pytest.raises(RuntimeError, match="unknown to this image"):
|
||||
prepare_database()
|
||||
|
||||
with engine.connect() as connection:
|
||||
acquired = connection.execute(
|
||||
text("SELECT pg_try_advisory_lock(:lock_id)"),
|
||||
{"lock_id": _MIGRATION_LOCK_ID},
|
||||
).scalar_one()
|
||||
assert acquired is True
|
||||
connection.execute(
|
||||
text("SELECT pg_advisory_unlock(:lock_id)"),
|
||||
{"lock_id": _MIGRATION_LOCK_ID},
|
||||
)
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
@@ -130,4 +196,15 @@ def _run_alembic_upgrade(revision: str) -> None:
|
||||
# #endregion _run_alembic_upgrade
|
||||
|
||||
|
||||
def _current_head() -> str:
|
||||
"""Return the head revision of the repository's Alembic chain."""
|
||||
from alembic.config import Config
|
||||
from alembic.script import ScriptDirectory
|
||||
|
||||
backend_dir = Path(__file__).parents[2]
|
||||
config = Config(str(backend_dir / "alembic.ini"))
|
||||
config.set_main_option("script_location", str(backend_dir / "alembic"))
|
||||
return ScriptDirectory.from_config(config).get_current_head()
|
||||
|
||||
|
||||
# #endregion Test.Alembic.TaskRecordsUserId
|
||||
|
||||
15
build.sh
15
build.sh
@@ -433,7 +433,10 @@ services:
|
||||
condition: service_healthy
|
||||
environment:
|
||||
DATABASE_URL: postgresql+psycopg2://\${POSTGRES_USER:-postgres}:\${POSTGRES_PASSWORD:?Set POSTGRES_PASSWORD}@db:5432/\${POSTGRES_DB:-ss_tools}
|
||||
# Safe one-shot default: legacy revisions are reset; 0001_baseline is preserved.
|
||||
DB_SCHEMA: \${DB_SCHEMA:-public}
|
||||
# One-shot destructive reset, only for databases NOT on the current chain
|
||||
# (legacy/orphaned Alembic revisions). Chain-managed schemas are never wiped.
|
||||
# The reset drops schema objects, so it needs object ownership, not schema ownership.
|
||||
RESET_DATABASE_SCHEMA: \${RESET_DATABASE_SCHEMA:-true}
|
||||
BACKEND_PORT: 8000
|
||||
AUTH_SECRET_KEY: \${AUTH_SECRET_KEY:?Set AUTH_SECRET_KEY in .env}
|
||||
@@ -570,8 +573,14 @@ POSTGRES_USER=postgres
|
||||
POSTGRES_PASSWORD=change-me
|
||||
POSTGRES_HOST_PORT=5432
|
||||
|
||||
# One-shot reset of any legacy Alembic revision before applying 0001_baseline.
|
||||
# Safe to leave true: once baseline is installed, subsequent starts skip reset.
|
||||
# Схема, в которой живут все объекты приложения (по умолчанию public).
|
||||
# Для внешней корпоративной БД: DB_SCHEMA=ss_tools (схему один раз создаёт DBA).
|
||||
DB_SCHEMA=public
|
||||
|
||||
# One-shot destructive reset, only for databases NOT on the current chain
|
||||
# (legacy/orphaned Alembic revisions). true never wipes a chain-managed schema;
|
||||
# false/unset fails fast on an unknown revision instead of dropping anything.
|
||||
# A reset that actually drops needs ownership of the schema objects, not the schema itself.
|
||||
RESET_DATABASE_SCHEMA=true
|
||||
|
||||
# ======================================================================
|
||||
|
||||
@@ -37,6 +37,7 @@ services:
|
||||
condition: service_healthy
|
||||
environment:
|
||||
DATABASE_URL: postgresql+psycopg2://postgres:postgres@db:5432/ss_tools
|
||||
DB_SCHEMA: ${DB_SCHEMA:-public}
|
||||
RESET_DATABASE_SCHEMA: ${RESET_DATABASE_SCHEMA:-false}
|
||||
BACKEND_PORT: 8000
|
||||
AUTH_SECRET_KEY: ${AUTH_SECRET_KEY:?Set AUTH_SECRET_KEY in .env.e2e}
|
||||
|
||||
@@ -61,7 +61,10 @@ services:
|
||||
environment:
|
||||
DATABASE_URL: postgresql+psycopg2://${POSTGRES_USER:-postgres}:${POSTGRES_PASSWORD:?Set POSTGRES_PASSWORD in .env}@${POSTGRES_HOST:-db}:${POSTGRES_PORT:-5432}/${POSTGRES_DB:-ss_tools}
|
||||
STORAGE_ROOT_PATH: /app/storage
|
||||
# Safe one-shot default: legacy revisions are reset; 0001_baseline is preserved.
|
||||
DB_SCHEMA: ${DB_SCHEMA:-public}
|
||||
# One-shot destructive reset, only for databases NOT on the current chain
|
||||
# (legacy/orphaned Alembic revisions). Chain-managed schemas are never wiped.
|
||||
# The reset drops schema objects, so it needs object ownership, not schema ownership.
|
||||
RESET_DATABASE_SCHEMA: ${RESET_DATABASE_SCHEMA:-true}
|
||||
BACKEND_PORT: 8000
|
||||
ALLOWED_ORIGINS: ${ALLOWED_ORIGINS:-}
|
||||
|
||||
@@ -32,6 +32,7 @@ services:
|
||||
environment:
|
||||
DATABASE_URL: postgresql+psycopg2://postgres:postgres@db:5432/ss_tools
|
||||
STORAGE_ROOT_PATH: /app/storage
|
||||
DB_SCHEMA: ${DB_SCHEMA:-public}
|
||||
RESET_DATABASE_SCHEMA: ${RESET_DATABASE_SCHEMA:-false}
|
||||
BACKEND_PORT: 8000
|
||||
ALLOWED_ORIGINS: ${ALLOWED_ORIGINS:-}
|
||||
|
||||
Reference in New Issue
Block a user