From d58d450017a437450eeee7f8ab8a412ff7eae325 Mon Sep 17 00:00:00 2001 From: "Tobias.Mikula" Date: Tue, 18 Aug 2026 16:08:45 +0200 Subject: [PATCH 1/3] feat: Adding a support of Flyway in the EventBus. --- .coverage | Bin 53248 -> 0 bytes .../{check_python.yml => quality_gates.yml} | 37 +++- database/README.md | 72 ++++++++ database/migrations/00_databases.ddl | 24 +++ .../migrations/V1.4.0.1__create_roles.ddl | 81 +++++++++ .../migrations/V1.4.0.2__initial_schema.ddl | 43 +++-- database/migrations/V1.4.0.3__grants.ddl | 47 +++++ flyway.toml | 33 ++++ src/utils/config_loader.py | 4 +- src/writers/writer_eventbridge.py | 3 +- tests/integration/conftest.py | 45 ++++- tests/integration/schemas/__init__.py | 15 -- tests/integration/test_baseline_migration.py | 163 ++++++++++++++++++ tests/integration/test_db_roles.py | 128 ++++++++++++++ 14 files changed, 641 insertions(+), 54 deletions(-) delete mode 100644 .coverage rename .github/workflows/{check_python.yml => quality_gates.yml} (79%) create mode 100644 database/README.md create mode 100644 database/migrations/00_databases.ddl create mode 100644 database/migrations/V1.4.0.1__create_roles.ddl rename tests/integration/schemas/postgres_schema.py => database/migrations/V1.4.0.2__initial_schema.ddl (72%) create mode 100644 database/migrations/V1.4.0.3__grants.ddl create mode 100644 flyway.toml delete mode 100644 tests/integration/schemas/__init__.py create mode 100644 tests/integration/test_baseline_migration.py create mode 100644 tests/integration/test_db_roles.py diff --git a/.coverage b/.coverage deleted file mode 100644 index 89085e7c3b7b780fd7111a4b73d480adf0977108..0000000000000000000000000000000000000000 GIT binary patch literal 0 HcmV?d00001 literal 53248 zcmeI4eQez19mnnS#}M_2lm{hr6? zdub9Br;|{_&&u6nKYpICf4|T3=HFd>TW{TLIJz=yT3O9eKE>5>JkPCH6prIWc*Wr5 zZDH66dIym575fYAid@&k*JI>%P6~aABX`FxA)V5Q^lo&Iv|9XP)ChNoU2p<7B!C2v z01{X>1PZrG;bbbs@7&>NgBjhiw6t!y@39+h>+9Xtr)=xpu(?lh^ORK)1-8ykrB|`c zNo8EOlwl*ID@Ja}NNbLf8&RB5-7B59^&z@uK_0AHv|tyjPxn+~2x>X{2uO@uMpmd;Vb( z7-mk@4QLHehX&l5#cNg2DOaASX0dp&@-P^eGd-;3iuBS@ znc1AnX9ofKSM%9%+Bi_dy|800RCA%aj#+ee(QPwNTbBK;YEBP1ZBe&YK4)yt>!rr; zP|7tsm>W5>L!nMU4_FZmC9jY39P~j+t%`hGx-$ zrkKToVChHN+o-W%gEE*g2dn!U>oK+t0|e#*uh?)tm!?C-raq+23Ek30^seGrcUt(U zI}&uA9x;^cYh9@8l)8Z0Q`(IbQlW6NrG?)qvTjsttzr_c>Ox*i&NGdAr?R?c=Php{ z`LIdSt5KdMy)jZe1MYN&L&@Hjh1@BYQ+tX#wS{mnoJ=J6f;%GUyhL9$RZh^H@(?(0 zAy~Zmo1N+mu1m%;;LqT^pdrCL6A4*b1{@?y9t|tYJ@6jgg5qW|OzL z&}n`f3_5Ro!A>@-<%Wvu64_H4(izLs?-~at zR5hy&xx<|{N`>O&s6I!|U5z$sj_w#)U7aRj@QlWN1AKvwg7EglA{de(U- z%bn~CFiXQNQf23&z$+NiL-z`FMs%Bf&_(B!D&8PP+G&8#lhYjh;f4f|01`j~NB{{S z0VIF~kN^@u0!RP}EI$GQAK=3@{tu8Dj=T+TxFG=~fCP{L5sF4y^-4@QSsN}4yb|~5xO0!5&epc zfO=V@KM5L%#cNDj@S_Em-R-98xM@2hmTrSJeNCW5maH_U4UbW$Zcczq%aUans6opZ zf@k>0r#Ci&SZwiPdB@1u-RxIB)QmYY0>8XLJ&XJE%vCAYS2JKzh638u05VCx#@O1h zuG+d~Xc^-k_$3oFNByc_ieK-ELYrs(WieDZF@7Jf#X_ZVZFO zg>(`&?NKZte^{#LDiG+c@#l;uP~l8}5F~t%ozvq-buGijPM-iuKIY3gO4_&qK5Ynq zjE_a2BvawkYF^+b7vIJeHR<>NDe-!a{7(Fucs+TLtd-xBS4f*AF824>_hWZN_eZB< ziO5muaI`b>xb%p;PtHZQMkMjj$7;DpXORFBKmter2_OL^z;PRcg($!HTdsNI|I`h^ zg6LbMI{t6o6f7isXn428bL0Qy=Ys{ZbR{qTZ`vFzv@BhwI{r`e1q(4BV)MuUjXphF zApVcv6f7it8#8bG->@XV&X51=dxM1<-v(92|8@PrLfnT!DflSvl;VGK;X-cwFRvAL z_%NtC{*NsQ>-q7&v{bC6_&>TnSP*=BQ5lx`9MxXbJ~#dsHv|h|AEtOp6~@$FXxn=6 zf7mCM>bVMH<&aNjJb?;lYJ8HN)8k9=f6(WwGDpdc{{^3mKuM;;slaMs(#J;#6*ckq z|L|r82_OL^fCP{L5u|!pT($j{9ZjKax!|Kb}3Q4jho0AoG$C$Nk_@O=u1!;7G3F8yX>93E}5GD-NzZ z$vu64e4tf|(~^xGcWC9o&(+p7z|j>Aa3s$BkUP2UlZ|2pstwDn zq3xuMG(w*A@HxACtf)mZm5Brl>!vY6#laRm>A|}_)rk?detrp0DA9qA2t1`G2yNBd5txa)i7}OfpLTNS-B6 zktfKbIr!WIHza@rkN^@u z0!RP}AOR$R1dsp{Kmtf$`4XVt|6}~WeA^d2Ljp(u2_OL^fCP{L5> "$GITHUB_OUTPUT" else echo "python_changed=false" >> "$GITHUB_OUTPUT" fi + if grep -Eq '^database/|^flyway\.toml$' <<< "$CHANGED_FILES"; then + echo "database_changed=true" >> "$GITHUB_OUTPUT" + else + echo "database_changed=false" >> "$GITHUB_OUTPUT" + fi + pylint-analysis: name: Pylint Static Code Analysis needs: detect @@ -139,7 +147,7 @@ jobs: integration-tests: name: Pytest Integration Tests needs: detect - if: needs.detect.outputs.python_changed == 'true' + if: needs.detect.outputs.python_changed == 'true' || needs.detect.outputs.database_changed == 'true' runs-on: ubuntu-latest timeout-minutes: 15 steps: @@ -152,13 +160,24 @@ jobs: - name: Set up dev Python environment uses: ./.github/actions/setup-dev-python-env + - name: Set up Java + uses: actions/setup-java@b6effb05e454b25005698d916606bdc6ffcbf961 + with: + distribution: temurin + java-version: '21' + + - name: Set up Flyway + uses: red-gate/setup-flyway@e024a17cd0890383f6996ed7edbded24c54ed86c + with: + version: '13.3.0' + - name: Run integration tests run: pytest tests/integration/ -v --tb=short --log-cli-level=INFO noop: name: No Operation needs: detect - if: needs.detect.outputs.python_changed != 'true' + if: needs.detect.outputs.python_changed != 'true' && needs.detect.outputs.database_changed != 'true' runs-on: ubuntu-latest steps: - run: echo "No changes in the *.py files — passing." diff --git a/database/README.md b/database/README.md new file mode 100644 index 0000000..9c062e6 --- /dev/null +++ b/database/README.md @@ -0,0 +1,72 @@ +# EventGate Database + +All database code lives here and is deployed with [Flyway](https://documentation.red-gate.com/flyway). +The migrations are the single source of truth for the schema, roles, and grants — the +same migrations build local, CI (integration tests), and real environments. + +## Layout + +``` +flyway.toml # Flyway configuration (locations, baseline, placeholders) — repo root +database/ +├── README.md +└── migrations/ + ├── 00_databases.ddl # One-off DB bootstrap (NOT a Flyway migration; no `V` prefix) + ├── V1.4.0.1__create_roles.ddl # owner / writer / reader roles + ├── V1.4.0.2__initial_schema.ddl # tables + └── V1.4.0.3__grants.ddl # ownership + least-privilege grants + ... +``` + +## Conventions + +- Versioned migrations follow Flyway's `V...__description.ext` format, + where `..` tracks the EventGate release the migration ships in and `` + increments per migration within that release. +- Extensions carry intent: `.ddl` for structural changes (tables, roles, constraints, indexes), + `.sql` for DML / data. + +## Roles + +| Role | Purpose | Used by | +|--------------------|-----------------------------------------------------|---------------------| +| master (superuser) | Runs the migrations | Flyway (deployment) | +| `eventgate_owner` | Owns the schema objects, may run DDL | Migrations | +| `eventgate_writer` | `SELECT` / `INSERT` / `UPDATE` on data tables | EventGate Lambda | +| `eventgate_reader` | `SELECT` only | EventStats Lambda | + +Role passwords are required Flyway placeholders (`eventgate_owner_password`, +`eventgate_writer_password`, `eventgate_reader_password`). Supply them from secrets in real +environments. + +## Local setup + +Requires the Flyway CLI (needs a JDK 17+) and Docker. + +```zsh +# 1. Start a local Postgres docker container +docker run --name=eventgate_db -e POSTGRES_PASSWORD=changeme -e POSTGRES_DB=eventgate_db -p 5432:5432 -d postgres:16 + +# 2. Apply the migrations (run from the repo root, where flyway.toml lives) +export FLYWAY_PLACEHOLDERS_EVENTGATE_OWNER_PASSWORD=changeme +export FLYWAY_PLACEHOLDERS_EVENTGATE_WRITER_PASSWORD=changeme +export FLYWAY_PLACEHOLDERS_EVENTGATE_READER_PASSWORD=changeme +flyway migrate + +# Inspect state / clean up +flyway info +docker kill eventgate_db && docker rm eventgate_db +``` + +## Adopting an existing database + +`flyway.toml` (repo root) sets `baselineOnMigrate = true` with `baselineVersion = 1.4.0.0`. On a +database that already contains the tables but has no Flyway history (i.e. production), the first +`flyway migrate` records the baseline and applies `V1.4.0.1+` on top. + +Before the first production migration: + +1. Compare the deployed schema with `V1.4.0.2__initial_schema.ddl`. +2. Back up the database and cluster roles. +3. Confirm the migration account can create roles and change ownership of every EventGate table. +4. Run `flyway info`, then `flyway migrate` with all role-password placeholders supplied from secrets. diff --git a/database/migrations/00_databases.ddl b/database/migrations/00_databases.ddl new file mode 100644 index 0000000..710ed5d --- /dev/null +++ b/database/migrations/00_databases.ddl @@ -0,0 +1,24 @@ +/* + * Copyright 2026 ABSA Group Limited + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +-- Database bootstrap (NOT a Flyway migration). +-- +-- Flyway connects to an existing database, so it cannot create the database it migrates. +-- This script is intentionally NOT prefixed with `V`, so Flyway ignores it. + +CREATE DATABASE eventgate_db + WITH + ENCODING = 'UTF8' + CONNECTION LIMIT = -1; diff --git a/database/migrations/V1.4.0.1__create_roles.ddl b/database/migrations/V1.4.0.1__create_roles.ddl new file mode 100644 index 0000000..ddc32ec --- /dev/null +++ b/database/migrations/V1.4.0.1__create_roles.ddl @@ -0,0 +1,81 @@ +/* + * Copyright 2026 ABSA Group Limited + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +-- Application database roles. +-- +-- eventgate_owner - owns the schema objects and may run DDL. +-- eventgate_writer - inserts/updates event data (main EventGate Lambda). +-- eventgate_reader - read-only access (EventStats Lambda). + +DO +$do$ + BEGIN + IF EXISTS ( + SELECT FROM pg_catalog.pg_roles + WHERE rolname = 'eventgate_owner') THEN + + RAISE NOTICE 'Role "eventgate_owner" already exists. Skipping.'; + ELSE + CREATE ROLE eventgate_owner WITH + LOGIN + NOSUPERUSER + INHERIT + NOCREATEDB + NOCREATEROLE + NOREPLICATION + PASSWORD '${eventgate_owner_password}'; + END IF; + END +$do$; + +DO +$do$ + BEGIN + IF EXISTS ( + SELECT FROM pg_catalog.pg_roles + WHERE rolname = 'eventgate_writer') THEN + RAISE NOTICE 'Role "eventgate_writer" already exists. Skipping.'; + ELSE + CREATE ROLE eventgate_writer WITH + LOGIN + NOSUPERUSER + INHERIT + NOCREATEDB + NOCREATEROLE + NOREPLICATION + PASSWORD '${eventgate_writer_password}'; + END IF; + END +$do$; + +DO +$do$ + BEGIN + IF EXISTS ( + SELECT FROM pg_catalog.pg_roles + WHERE rolname = 'eventgate_reader') THEN + RAISE NOTICE 'Role "eventgate_reader" already exists. Skipping.'; + ELSE + CREATE ROLE eventgate_reader WITH + LOGIN + NOSUPERUSER + INHERIT + NOCREATEDB + NOCREATEROLE + NOREPLICATION + PASSWORD '${eventgate_reader_password}'; + END IF; + END +$do$; diff --git a/tests/integration/schemas/postgres_schema.py b/database/migrations/V1.4.0.2__initial_schema.ddl similarity index 72% rename from tests/integration/schemas/postgres_schema.py rename to database/migrations/V1.4.0.2__initial_schema.ddl index 1a286e8..3740825 100644 --- a/tests/integration/schemas/postgres_schema.py +++ b/database/migrations/V1.4.0.2__initial_schema.ddl @@ -1,23 +1,21 @@ -# -# Copyright 2026 ABSA Group Limited -# -# Licensed under the Apache License, Version 2.0 (the "License"); -# you may not use this file except in compliance with the License. -# You may obtain a copy of the License at -# -# http://www.apache.org/licenses/LICENSE-2.0 -# -# Unless required by applicable law or agreed to in writing, software -# distributed under the License is distributed on an "AS IS" BASIS, -# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -# See the License for the specific language governing permissions and -# limitations under the License. -# +/* + * Copyright 2026 ABSA Group Limited + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ -"""PostgreSQL schema for integration tests.""" +-- Initial EventGate schema. -SCHEMA_SQL = """ --- Table matching WriterPostgres._postgres_run_write columns +-- Run header rows for the runs topic. CREATE TABLE IF NOT EXISTS public_cps_za_runs ( event_id VARCHAR(255) NOT NULL, job_ref VARCHAR(255) NOT NULL, @@ -29,7 +27,7 @@ timestamp_end BIGINT ); --- Table matching WriterPostgres._postgres_run_write job rows +-- Per-job rows belonging to a run. CREATE TABLE IF NOT EXISTS public_cps_za_runs_jobs ( internal_id SERIAL PRIMARY KEY, event_id VARCHAR(255) NOT NULL, @@ -42,7 +40,7 @@ additional_info JSONB ); --- Table matching WriterPostgres._postgres_edla_write columns +-- Data lake change events. CREATE TABLE IF NOT EXISTS public_cps_za_dlchange ( event_id VARCHAR(255) NOT NULL, tenant_id VARCHAR(255) NOT NULL, @@ -59,7 +57,7 @@ additional_info JSONB ); --- Table matching WriterPostgres._postgres_test_write columns +-- Test topic events. CREATE TABLE IF NOT EXISTS public_cps_za_test ( event_id VARCHAR(255) NOT NULL, tenant_id VARCHAR(255) NOT NULL, @@ -69,7 +67,7 @@ additional_info JSONB ); --- Table for test_status_change_writer +-- Aggregated latest status per job (see ADR 001). CREATE TABLE IF NOT EXISTS public_cps_za_status_change_aggregated_job ( job_id UUID PRIMARY KEY, job_group_id UUID, @@ -97,4 +95,3 @@ finished_at TIMESTAMPTZ, last_updated_at TIMESTAMPTZ NOT NULL ); -""" diff --git a/database/migrations/V1.4.0.3__grants.ddl b/database/migrations/V1.4.0.3__grants.ddl new file mode 100644 index 0000000..02401d7 --- /dev/null +++ b/database/migrations/V1.4.0.3__grants.ddl @@ -0,0 +1,47 @@ +/* + * Copyright 2026 ABSA Group Limited + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +-- Object ownership and least-privilege grants for the application roles. + +-- Owner: owns every table (and its sequences) in the public schema. +ALTER TABLE public_cps_za_runs OWNER TO eventgate_owner; +ALTER TABLE public_cps_za_runs_jobs OWNER TO eventgate_owner; +ALTER TABLE public_cps_za_dlchange OWNER TO eventgate_owner; +ALTER TABLE public_cps_za_test OWNER TO eventgate_owner; +ALTER TABLE public_cps_za_status_change_aggregated_job OWNER TO eventgate_owner; + +-- Both application roles (writer and reader) need to access the public schema. +GRANT USAGE ON SCHEMA public TO eventgate_writer, eventgate_reader; + +-- Reader: read-only access to EventGate data tables. +GRANT SELECT ON TABLE + public_cps_za_runs, + public_cps_za_runs_jobs, + public_cps_za_dlchange, + public_cps_za_test, + public_cps_za_status_change_aggregated_job +TO eventgate_reader; + +-- Writer: read and write EventGate data, but no DDL or migration metadata. +GRANT SELECT, INSERT, UPDATE ON TABLE + public_cps_za_runs, + public_cps_za_runs_jobs, + public_cps_za_dlchange, + public_cps_za_test, + public_cps_za_status_change_aggregated_job +TO eventgate_writer; + +-- Writer needs the SERIAL sequence (public_cps_za_runs_jobs.internal_id) to insert. +GRANT USAGE, SELECT ON SEQUENCE public_cps_za_runs_jobs_internal_id_seq TO eventgate_writer; diff --git a/flyway.toml b/flyway.toml new file mode 100644 index 0000000..2861a72 --- /dev/null +++ b/flyway.toml @@ -0,0 +1,33 @@ +# +# Copyright 2026 ABSA Group Limited +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +# Flyway configuration for EventGate database migrations. +# Lives at the project root alongside pyproject.toml, per project convention +# for top-level config. Paths below are relative to the repo root. + +[flyway] +locations = ["filesystem:database/migrations"] +sqlMigrationSuffixes = [".ddl", ".sql"] + +# Adopt the known legacy database: on a non-empty database without Flyway +# history, record a baseline at 1.4.0.0. +baselineOnMigrate = true +baselineVersion = "1.4.0.0" + +[environments.default] +url = "jdbc:postgresql://localhost:5432/eventgate_db" +user = "postgres" +password = "changeme" diff --git a/src/utils/config_loader.py b/src/utils/config_loader.py index 3093da4..c7565c2 100644 --- a/src/utils/config_loader.py +++ b/src/utils/config_loader.py @@ -59,7 +59,9 @@ def _load_json_from_path(path: str, aws_s3: ServiceResource) -> dict[str, Any]: name_parts = path.split("/") bucket_name = name_parts[2] bucket_object_key = "/".join(name_parts[3:]) - return json.loads(aws_s3.Bucket(bucket_name).Object(bucket_object_key).get()["Body"].read().decode("utf-8")) + bucket = aws_s3.Bucket(bucket_name) # type: ignore[attr-defined] + s3_object = bucket.Object(bucket_object_key) + return json.loads(s3_object.get()["Body"].read().decode("utf-8")) with open(path, "r", encoding="utf-8") as file: return json.load(file) diff --git a/src/writers/writer_eventbridge.py b/src/writers/writer_eventbridge.py index 076f007..d4908d4 100644 --- a/src/writers/writer_eventbridge.py +++ b/src/writers/writer_eventbridge.py @@ -36,7 +36,8 @@ class WriterEventBridge(Writer): def __init__(self, config: dict[str, Any]) -> None: super().__init__(config) - self._client: Optional["boto3.client"] = None + # boto3 clients are generated dynamically, so no precise static type exists. + self._client: Optional[Any] = None self._entries: list[dict[str, Any]] = [] self.event_bus_arn: str = config.get("event_bus_arn", "") logger.debug("Initialized EventBridge writer.") diff --git a/tests/integration/conftest.py b/tests/integration/conftest.py index 0ffe7f5..c62bf4e 100644 --- a/tests/integration/conftest.py +++ b/tests/integration/conftest.py @@ -20,6 +20,7 @@ import logging import os import shutil +import subprocess import time from concurrent.futures import ThreadPoolExecutor, as_completed from http.server import HTTPServer, BaseHTTPRequestHandler @@ -37,12 +38,13 @@ from testcontainers.kafka import KafkaContainer from testcontainers.postgres import PostgresContainer -from tests.integration.schemas.postgres_schema import SCHEMA_SQL from tests.integration.utils.jwt_helper import create_test_jwt_keypair, generate_token logger = logging.getLogger(__name__) PROJECT_ROOT = Path(__file__).parent.parent.parent +FLYWAY_CONFIG = PROJECT_ROOT / "flyway.toml" +TEST_ROLE_PASSWORD = "changeme" # Mock JWT Provider (runs in-process via threading) @@ -165,6 +167,40 @@ def _convert_dsn(dsn: str) -> str: return dsn.replace("postgresql+psycopg2://", "postgresql://") +def _run_flyway_migrate(dsn: str) -> None: + """Apply Flyway migrations from `database/migrations` to the given database. + Runs the same migrations used for real environments so integration tests + validate the migrations as the single source of truth for the schema. + Args: + dsn: psycopg2-style DSN of the target database. + Raises: + RuntimeError: If the `flyway migrate` command fails. + """ + parsed = urlparse(dsn) + jdbc_url = f"jdbc:postgresql://{parsed.hostname}:{parsed.port}{parsed.path}" + command = [ + "flyway", + f"-configFiles={FLYWAY_CONFIG}", + f"-workingDirectory={PROJECT_ROOT}", + f"-url={jdbc_url}", + f"-user={parsed.username}", + f"-password={parsed.password}", + "migrate", + ] + environment = os.environ.copy() + environment.update( + { + "FLYWAY_PLACEHOLDERS_EVENTGATE_OWNER_PASSWORD": TEST_ROLE_PASSWORD, + "FLYWAY_PLACEHOLDERS_EVENTGATE_WRITER_PASSWORD": TEST_ROLE_PASSWORD, + "FLYWAY_PLACEHOLDERS_EVENTGATE_READER_PASSWORD": TEST_ROLE_PASSWORD, + } + ) + flyway_process = subprocess.run(command, capture_output=True, text=True, check=False, env=environment) + if flyway_process.returncode != 0: + raise RuntimeError(f"Flyway migrate failed:\n{flyway_process.stdout}\n{flyway_process.stderr}") + logger.debug("Flyway migrate output:\n%s", flyway_process.stdout) + + @pytest.fixture(scope="session") def postgres_container() -> Generator[str, None, None]: """PostgreSQL container with initialized schema.""" @@ -192,11 +228,10 @@ def postgres_container() -> Generator[str, None, None]: if conn is None: raise TimeoutError(f"Timed out waiting for Postgres to become available after 5 attempts: {last_exc}") - conn.autocommit = True - with conn.cursor() as cursor: - cursor.execute(SCHEMA_SQL) conn.close() - logger.debug("PostgreSQL schema initialized.") + logger.debug("Postgres ready, applying Flyway migrations.") + _run_flyway_migrate(dsn) + logger.debug("PostgreSQL schema initialized via Flyway.") yield dsn diff --git a/tests/integration/schemas/__init__.py b/tests/integration/schemas/__init__.py deleted file mode 100644 index ebfbdd3..0000000 --- a/tests/integration/schemas/__init__.py +++ /dev/null @@ -1,15 +0,0 @@ -# -# Copyright 2026 ABSA Group Limited -# -# Licensed under the Apache License, Version 2.0 (the "License"); -# you may not use this file except in compliance with the License. -# You may obtain a copy of the License at -# -# http://www.apache.org/licenses/LICENSE-2.0 -# -# Unless required by applicable law or agreed to in writing, software -# distributed under the License is distributed on an "AS IS" BASIS, -# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -# See the License for the specific language governing permissions and -# limitations under the License. -# diff --git a/tests/integration/test_baseline_migration.py b/tests/integration/test_baseline_migration.py new file mode 100644 index 0000000..f269a88 --- /dev/null +++ b/tests/integration/test_baseline_migration.py @@ -0,0 +1,163 @@ +# +# Copyright 2026 ABSA Group Limited +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +import time +from typing import Generator + +import psycopg2 +import pytest +from testcontainers.postgres import PostgresContainer + +from tests.integration.conftest import _convert_dsn, _run_flyway_migrate + +# Mimics the complete hand-created production schema before Flyway is introduced. +LEGACY_SCHEMA_SQL = """ +CREATE TABLE public_cps_za_runs ( + event_id VARCHAR(255) NOT NULL, + job_ref VARCHAR(255) NOT NULL, + tenant_id VARCHAR(255) NOT NULL, + source_app VARCHAR(255) NOT NULL, + source_app_version VARCHAR(255) NOT NULL, + environment VARCHAR(255) NOT NULL, + timestamp_start BIGINT, + timestamp_end BIGINT +); + +CREATE TABLE public_cps_za_runs_jobs ( + internal_id SERIAL PRIMARY KEY, + event_id VARCHAR(255) NOT NULL, + country VARCHAR(255), + catalog_id VARCHAR(255) NOT NULL, + status VARCHAR(50) NOT NULL, + timestamp_start BIGINT, + timestamp_end BIGINT, + message TEXT, + additional_info JSONB +); + +CREATE TABLE public_cps_za_dlchange ( + event_id VARCHAR(255) NOT NULL, + tenant_id VARCHAR(255) NOT NULL, + source_app VARCHAR(255) NOT NULL, + source_app_version VARCHAR(255) NOT NULL, + environment VARCHAR(255) NOT NULL, + timestamp_event BIGINT, + country VARCHAR(255), + catalog_id VARCHAR(255) NOT NULL, + operation VARCHAR(255), + "location" TEXT, + "format" VARCHAR(255), + format_options JSONB, + additional_info JSONB +); + +CREATE TABLE public_cps_za_test ( + event_id VARCHAR(255) NOT NULL, + tenant_id VARCHAR(255) NOT NULL, + source_app VARCHAR(255) NOT NULL, + environment VARCHAR(255) NOT NULL, + timestamp_event BIGINT, + additional_info JSONB +); + +CREATE TABLE public_cps_za_status_change_aggregated_job ( + job_id UUID PRIMARY KEY, + job_group_id UUID, + parent_job_id UUID, + initial_job_id UUID, + job_ref TEXT, + job_name TEXT, + definition_id TEXT, + definition_version TEXT, + tenant_id TEXT, + country TEXT, + source_app TEXT, + source_app_version TEXT, + environment TEXT, + platform TEXT, + platform_metadata JSONB, + input_arguments JSONB, + additional_context JSONB, + attempt_number INTEGER NOT NULL DEFAULT 1 CHECK (attempt_number > 0), + status_type TEXT CHECK (status_type IN ('WAITING', 'RUNNING', 'SUCCEEDED', 'FAILED', 'KILLED')), + status_subtype TEXT, + status_detail TEXT, + created_at TIMESTAMPTZ, + started_at TIMESTAMPTZ, + finished_at TIMESTAMPTZ, + last_updated_at TIMESTAMPTZ NOT NULL +); + +INSERT INTO public_cps_za_test + (event_id, tenant_id, source_app, environment, timestamp_event) +VALUES ('preexisting-row', 'tenant', 'app', 'env', 42); +""" + + +@pytest.fixture(scope="module") +def preseeded_dsn() -> Generator[str, None, None]: + container = PostgresContainer("postgres:16", dbname="eventgate") + container.start() + dsn = _convert_dsn(container.get_connection_url()) + + conn = None + for attempt in range(1, 6): + try: + conn = psycopg2.connect(dsn) + break + except psycopg2.OperationalError: + if attempt < 5: + time.sleep(2) + if conn is None: + raise TimeoutError("Timed out waiting for Postgres to become available.") + + conn.autocommit = True + with conn.cursor() as cursor: + cursor.execute(LEGACY_SCHEMA_SQL) + conn.close() + + yield dsn + + container.stop() + + +def test_migrate_on_preseeded_db_preserves_existing_data(preseeded_dsn: str) -> None: + _run_flyway_migrate(preseeded_dsn) + + conn = psycopg2.connect(preseeded_dsn) + try: + with conn.cursor() as cursor: + cursor.execute("SELECT timestamp_event FROM public_cps_za_test WHERE event_id = 'preexisting-row'") + row = cursor.fetchone() + finally: + conn.close() + + assert row is not None + assert 42 == row[0] + + +def test_migrate_on_preseeded_db_records_baseline(preseeded_dsn: str) -> None: + _run_flyway_migrate(preseeded_dsn) + + conn = psycopg2.connect(preseeded_dsn) + try: + with conn.cursor() as cursor: + cursor.execute("SELECT version FROM flyway_schema_history WHERE type = 'BASELINE'") + baseline_versions = {row[0] for row in cursor.fetchall()} + finally: + conn.close() + + assert "1.4.0.0" in baseline_versions diff --git a/tests/integration/test_db_roles.py b/tests/integration/test_db_roles.py new file mode 100644 index 0000000..992ebf4 --- /dev/null +++ b/tests/integration/test_db_roles.py @@ -0,0 +1,128 @@ +# +# Copyright 2026 ABSA Group Limited +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +from urllib.parse import urlparse + +import psycopg2 +import pytest +from psycopg2 import errors + +from tests.integration.conftest import TEST_ROLE_PASSWORD + +SELECTABLE_TABLE = "public_cps_za_test" +OWNED_TABLES = ( + "public_cps_za_runs", + "public_cps_za_runs_jobs", + "public_cps_za_dlchange", + "public_cps_za_test", + "public_cps_za_status_change_aggregated_job", +) + + +def _connect_as(dsn: str, role: str) -> "psycopg2.extensions.connection": + parsed = urlparse(dsn) + return psycopg2.connect( + host=parsed.hostname, + port=parsed.port, + dbname=parsed.path.lstrip("/"), + user=role, + password=TEST_ROLE_PASSWORD, + ) + + +class TestReaderRole: + def test_reader_can_select(self, postgres_container: str) -> None: + conn = _connect_as(postgres_container, "eventgate_reader") + try: + with conn.cursor() as cursor: + cursor.execute(f"SELECT 1 FROM {SELECTABLE_TABLE} LIMIT 1") + finally: + conn.close() + + def test_reader_cannot_insert(self, postgres_container: str) -> None: + conn = _connect_as(postgres_container, "eventgate_reader") + try: + with conn.cursor() as cursor, pytest.raises(errors.InsufficientPrivilege): + cursor.execute( + f"INSERT INTO {SELECTABLE_TABLE} " + "(event_id, tenant_id, source_app, environment, timestamp_event) " + "VALUES ('e', 't', 'app', 'env', 1)" + ) + finally: + conn.close() + + def test_reader_cannot_read_flyway_history(self, postgres_container: str) -> None: + conn = _connect_as(postgres_container, "eventgate_reader") + try: + with conn.cursor() as cursor, pytest.raises(errors.InsufficientPrivilege): + cursor.execute("SELECT version FROM flyway_schema_history") + finally: + conn.close() + + +class TestWriterRole: + def test_writer_can_insert_and_select(self, postgres_container: str) -> None: + conn = _connect_as(postgres_container, "eventgate_writer") + try: + with conn.cursor() as cursor: + cursor.execute( + f"INSERT INTO {SELECTABLE_TABLE} " + "(event_id, tenant_id, source_app, environment, timestamp_event) " + "VALUES ('writer-e', 't', 'app', 'env', 1)" + ) + cursor.execute(f"SELECT 1 FROM {SELECTABLE_TABLE} LIMIT 1") + conn.commit() + finally: + conn.close() + + def test_writer_cannot_drop_table(self, postgres_container: str) -> None: + conn = _connect_as(postgres_container, "eventgate_writer") + try: + with conn.cursor() as cursor, pytest.raises(errors.InsufficientPrivilege): + cursor.execute(f"DROP TABLE {SELECTABLE_TABLE}") + finally: + conn.close() + + def test_writer_cannot_modify_flyway_history(self, postgres_container: str) -> None: + conn = _connect_as(postgres_container, "eventgate_writer") + try: + with conn.cursor() as cursor, pytest.raises(errors.InsufficientPrivilege): + cursor.execute("UPDATE flyway_schema_history SET description = description") + finally: + conn.close() + + +class TestOwnerRole: + def test_owner_owns_tables(self, postgres_container: str) -> None: + conn = _connect_as(postgres_container, "eventgate_owner") + try: + with conn.cursor() as cursor: + cursor.execute( + "SELECT tablename FROM pg_tables " "WHERE schemaname = 'public' AND tableowner = 'eventgate_owner'" + ) + owned = {row[0] for row in cursor.fetchall()} + finally: + conn.close() + assert set(OWNED_TABLES) == owned + + def test_owner_can_alter_table(self, postgres_container: str) -> None: + conn = _connect_as(postgres_container, "eventgate_owner") + try: + with conn.cursor() as cursor: + cursor.execute(f"ALTER TABLE {SELECTABLE_TABLE} ADD COLUMN tmp_col TEXT") + conn.rollback() + finally: + conn.close() From 8fdcb4c129117422db0bc20a39a93d9c2f34e350 Mon Sep 17 00:00:00 2001 From: "Tobias.Mikula" Date: Tue, 18 Aug 2026 16:20:15 +0200 Subject: [PATCH 2/3] Fixing the red-gate/setup-flyway setting. --- .github/workflows/quality_gates.yml | 2 ++ 1 file changed, 2 insertions(+) diff --git a/.github/workflows/quality_gates.yml b/.github/workflows/quality_gates.yml index 14de5c4..efd58b5 100644 --- a/.github/workflows/quality_gates.yml +++ b/.github/workflows/quality_gates.yml @@ -170,6 +170,8 @@ jobs: uses: red-gate/setup-flyway@e024a17cd0890383f6996ed7edbded24c54ed86c with: version: '13.3.0' + edition: community + i-agree-to-the-eula: true - name: Run integration tests run: pytest tests/integration/ -v --tb=short --log-cli-level=INFO From 8b73226ecf6539cc84747283a488c5f0d1272e5b Mon Sep 17 00:00:00 2001 From: "Tobias.Mikula" Date: Thu, 20 Aug 2026 16:14:47 +0200 Subject: [PATCH 3/3] Reaction for the review comments. --- database/README.md | 19 ++- database/migrations/V1.4.0.3__grants.ddl | 43 +++-- flyway.toml | 5 - tests/integration/test_baseline_migration.py | 163 ------------------- 4 files changed, 40 insertions(+), 190 deletions(-) delete mode 100644 tests/integration/test_baseline_migration.py diff --git a/database/README.md b/database/README.md index 9c062e6..a16f9c5 100644 --- a/database/README.md +++ b/database/README.md @@ -6,8 +6,8 @@ same migrations build local, CI (integration tests), and real environments. ## Layout -``` -flyway.toml # Flyway configuration (locations, baseline, placeholders) — repo root +```text +flyway.toml # Flyway configuration (locations, baseline, placeholders) — repo root database/ ├── README.md └── migrations/ @@ -60,13 +60,20 @@ docker kill eventgate_db && docker rm eventgate_db ## Adopting an existing database -`flyway.toml` (repo root) sets `baselineOnMigrate = true` with `baselineVersion = 1.4.0.0`. On a -database that already contains the tables but has no Flyway history (i.e. production), the first -`flyway migrate` records the baseline and applies `V1.4.0.1+` on top. +On a database that already contains the tables but has no Flyway history (i.e. production), a +plain `flyway migrate` fails because Flyway sees existing objects it didn't create. The first +migration against such a database must instead pass baseline flags explicitly, one time only: + +```zsh +flyway -baselineOnMigrate=true -baselineVersion=1.4.0.0 migrate +``` + +This records a baseline at `1.4.0.0` in `flyway_schema_history` and then applies `V1.4.0.1+` on +top. Before the first production migration: 1. Compare the deployed schema with `V1.4.0.2__initial_schema.ddl`. 2. Back up the database and cluster roles. 3. Confirm the migration account can create roles and change ownership of every EventGate table. -4. Run `flyway info`, then `flyway migrate` with all role-password placeholders supplied from secrets. +4. Run `flyway info`, then the baseline command above with all role-password placeholders supplied from secrets. diff --git a/database/migrations/V1.4.0.3__grants.ddl b/database/migrations/V1.4.0.3__grants.ddl index 02401d7..7282e9a 100644 --- a/database/migrations/V1.4.0.3__grants.ddl +++ b/database/migrations/V1.4.0.3__grants.ddl @@ -16,32 +16,43 @@ -- Object ownership and least-privilege grants for the application roles. -- Owner: owns every table (and its sequences) in the public schema. -ALTER TABLE public_cps_za_runs OWNER TO eventgate_owner; -ALTER TABLE public_cps_za_runs_jobs OWNER TO eventgate_owner; -ALTER TABLE public_cps_za_dlchange OWNER TO eventgate_owner; -ALTER TABLE public_cps_za_test OWNER TO eventgate_owner; -ALTER TABLE public_cps_za_status_change_aggregated_job OWNER TO eventgate_owner; +ALTER TABLE public.public_cps_za_runs OWNER TO eventgate_owner; +ALTER TABLE public.public_cps_za_runs_jobs OWNER TO eventgate_owner; +ALTER TABLE public.public_cps_za_dlchange OWNER TO eventgate_owner; +ALTER TABLE public.public_cps_za_test OWNER TO eventgate_owner; +ALTER TABLE public.public_cps_za_status_change_aggregated_job OWNER TO eventgate_owner; + +-- Owner also needs CREATE on the schema so it can create future tables directly +GRANT CREATE ON SCHEMA public TO eventgate_owner; -- Both application roles (writer and reader) need to access the public schema. GRANT USAGE ON SCHEMA public TO eventgate_writer, eventgate_reader; -- Reader: read-only access to EventGate data tables. GRANT SELECT ON TABLE - public_cps_za_runs, - public_cps_za_runs_jobs, - public_cps_za_dlchange, - public_cps_za_test, - public_cps_za_status_change_aggregated_job + public.public_cps_za_runs, + public.public_cps_za_runs_jobs, + public.public_cps_za_dlchange, + public.public_cps_za_test, + public.public_cps_za_status_change_aggregated_job TO eventgate_reader; -- Writer: read and write EventGate data, but no DDL or migration metadata. GRANT SELECT, INSERT, UPDATE ON TABLE - public_cps_za_runs, - public_cps_za_runs_jobs, - public_cps_za_dlchange, - public_cps_za_test, - public_cps_za_status_change_aggregated_job + public.public_cps_za_runs, + public.public_cps_za_runs_jobs, + public.public_cps_za_dlchange, + public.public_cps_za_test, + public.public_cps_za_status_change_aggregated_job TO eventgate_writer; -- Writer needs the SERIAL sequence (public_cps_za_runs_jobs.internal_id) to insert. -GRANT USAGE, SELECT ON SEQUENCE public_cps_za_runs_jobs_internal_id_seq TO eventgate_writer; +GRANT USAGE, SELECT ON SEQUENCE public.public_cps_za_runs_jobs_internal_id_seq TO eventgate_writer; + +-- Default privileges +ALTER DEFAULT PRIVILEGES FOR ROLE eventgate_owner IN SCHEMA public + GRANT SELECT ON TABLES TO eventgate_reader; +ALTER DEFAULT PRIVILEGES FOR ROLE eventgate_owner IN SCHEMA public + GRANT SELECT, INSERT, UPDATE ON TABLES TO eventgate_writer; +ALTER DEFAULT PRIVILEGES FOR ROLE eventgate_owner IN SCHEMA public + GRANT USAGE, SELECT ON SEQUENCES TO eventgate_writer; diff --git a/flyway.toml b/flyway.toml index 2861a72..e252d13 100644 --- a/flyway.toml +++ b/flyway.toml @@ -22,11 +22,6 @@ locations = ["filesystem:database/migrations"] sqlMigrationSuffixes = [".ddl", ".sql"] -# Adopt the known legacy database: on a non-empty database without Flyway -# history, record a baseline at 1.4.0.0. -baselineOnMigrate = true -baselineVersion = "1.4.0.0" - [environments.default] url = "jdbc:postgresql://localhost:5432/eventgate_db" user = "postgres" diff --git a/tests/integration/test_baseline_migration.py b/tests/integration/test_baseline_migration.py deleted file mode 100644 index f269a88..0000000 --- a/tests/integration/test_baseline_migration.py +++ /dev/null @@ -1,163 +0,0 @@ -# -# Copyright 2026 ABSA Group Limited -# -# Licensed under the Apache License, Version 2.0 (the "License"); -# you may not use this file except in compliance with the License. -# You may obtain a copy of the License at -# -# http://www.apache.org/licenses/LICENSE-2.0 -# -# Unless required by applicable law or agreed to in writing, software -# distributed under the License is distributed on an "AS IS" BASIS, -# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -# See the License for the specific language governing permissions and -# limitations under the License. -# - -import time -from typing import Generator - -import psycopg2 -import pytest -from testcontainers.postgres import PostgresContainer - -from tests.integration.conftest import _convert_dsn, _run_flyway_migrate - -# Mimics the complete hand-created production schema before Flyway is introduced. -LEGACY_SCHEMA_SQL = """ -CREATE TABLE public_cps_za_runs ( - event_id VARCHAR(255) NOT NULL, - job_ref VARCHAR(255) NOT NULL, - tenant_id VARCHAR(255) NOT NULL, - source_app VARCHAR(255) NOT NULL, - source_app_version VARCHAR(255) NOT NULL, - environment VARCHAR(255) NOT NULL, - timestamp_start BIGINT, - timestamp_end BIGINT -); - -CREATE TABLE public_cps_za_runs_jobs ( - internal_id SERIAL PRIMARY KEY, - event_id VARCHAR(255) NOT NULL, - country VARCHAR(255), - catalog_id VARCHAR(255) NOT NULL, - status VARCHAR(50) NOT NULL, - timestamp_start BIGINT, - timestamp_end BIGINT, - message TEXT, - additional_info JSONB -); - -CREATE TABLE public_cps_za_dlchange ( - event_id VARCHAR(255) NOT NULL, - tenant_id VARCHAR(255) NOT NULL, - source_app VARCHAR(255) NOT NULL, - source_app_version VARCHAR(255) NOT NULL, - environment VARCHAR(255) NOT NULL, - timestamp_event BIGINT, - country VARCHAR(255), - catalog_id VARCHAR(255) NOT NULL, - operation VARCHAR(255), - "location" TEXT, - "format" VARCHAR(255), - format_options JSONB, - additional_info JSONB -); - -CREATE TABLE public_cps_za_test ( - event_id VARCHAR(255) NOT NULL, - tenant_id VARCHAR(255) NOT NULL, - source_app VARCHAR(255) NOT NULL, - environment VARCHAR(255) NOT NULL, - timestamp_event BIGINT, - additional_info JSONB -); - -CREATE TABLE public_cps_za_status_change_aggregated_job ( - job_id UUID PRIMARY KEY, - job_group_id UUID, - parent_job_id UUID, - initial_job_id UUID, - job_ref TEXT, - job_name TEXT, - definition_id TEXT, - definition_version TEXT, - tenant_id TEXT, - country TEXT, - source_app TEXT, - source_app_version TEXT, - environment TEXT, - platform TEXT, - platform_metadata JSONB, - input_arguments JSONB, - additional_context JSONB, - attempt_number INTEGER NOT NULL DEFAULT 1 CHECK (attempt_number > 0), - status_type TEXT CHECK (status_type IN ('WAITING', 'RUNNING', 'SUCCEEDED', 'FAILED', 'KILLED')), - status_subtype TEXT, - status_detail TEXT, - created_at TIMESTAMPTZ, - started_at TIMESTAMPTZ, - finished_at TIMESTAMPTZ, - last_updated_at TIMESTAMPTZ NOT NULL -); - -INSERT INTO public_cps_za_test - (event_id, tenant_id, source_app, environment, timestamp_event) -VALUES ('preexisting-row', 'tenant', 'app', 'env', 42); -""" - - -@pytest.fixture(scope="module") -def preseeded_dsn() -> Generator[str, None, None]: - container = PostgresContainer("postgres:16", dbname="eventgate") - container.start() - dsn = _convert_dsn(container.get_connection_url()) - - conn = None - for attempt in range(1, 6): - try: - conn = psycopg2.connect(dsn) - break - except psycopg2.OperationalError: - if attempt < 5: - time.sleep(2) - if conn is None: - raise TimeoutError("Timed out waiting for Postgres to become available.") - - conn.autocommit = True - with conn.cursor() as cursor: - cursor.execute(LEGACY_SCHEMA_SQL) - conn.close() - - yield dsn - - container.stop() - - -def test_migrate_on_preseeded_db_preserves_existing_data(preseeded_dsn: str) -> None: - _run_flyway_migrate(preseeded_dsn) - - conn = psycopg2.connect(preseeded_dsn) - try: - with conn.cursor() as cursor: - cursor.execute("SELECT timestamp_event FROM public_cps_za_test WHERE event_id = 'preexisting-row'") - row = cursor.fetchone() - finally: - conn.close() - - assert row is not None - assert 42 == row[0] - - -def test_migrate_on_preseeded_db_records_baseline(preseeded_dsn: str) -> None: - _run_flyway_migrate(preseeded_dsn) - - conn = psycopg2.connect(preseeded_dsn) - try: - with conn.cursor() as cursor: - cursor.execute("SELECT version FROM flyway_schema_history WHERE type = 'BASELINE'") - baseline_versions = {row[0] for row in cursor.fetchall()} - finally: - conn.close() - - assert "1.4.0.0" in baseline_versions